|
|
|
|
# conn
|
|
|
|
|
|
|
|
|
|
`conn` es el artefacto de conexiones realtime. No es "un wrapper de
|
|
|
|
|
WebSocket": es un registro de conexiones con transporte intercambiable,
|
|
|
|
|
reconexión, heartbeat, request/reply, canales, integración opcional con sesión
|
|
|
|
|
y diagnósticos estructurados.
|
|
|
|
|
|
|
|
|
|
La regla de nombres es importante:
|
|
|
|
|
|
|
|
|
|
| Concepto | Nombre |
|
|
|
|
|
| -------- | ------ |
|
|
|
|
|
| Raíz imperativa | `createEngineConnections()` / `EngineConnections` |
|
|
|
|
|
| Raíz reactiva | `createActiveConnections()` / `ActiveConnections` |
|
|
|
|
|
| Unidad individual | `Connection` |
|
|
|
|
|
| Topic lógico dentro de una conexión | `ConnectionChannel` |
|
|
|
|
|
|
|
|
|
|
No existe `EngineConnection` ni `ActiveConnection`. `Engine*` y `Active*`
|
|
|
|
|
quedan reservados para raíces de artefacto; una conexión individual no es una
|
|
|
|
|
raíz, es una entidad gestionada por `EngineConnections`.
|
|
|
|
|
|
|
|
|
|
## Uso Mínimo
|
|
|
|
|
|
|
|
|
|
```ts
|
|
|
|
|
import { createEngineConnections, createWebSocketTransport } from '$conn';
|
|
|
|
|
|
|
|
|
|
const Connections = createEngineConnections();
|
|
|
|
|
|
|
|
|
|
const Main = Connections.createConnection('main', {
|
|
|
|
|
transport: createWebSocketTransport({ url: () => '/realtime' }),
|
|
|
|
|
heartbeat: false
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
await Main.connect();
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
La raíz mantiene el mapa de conexiones:
|
|
|
|
|
|
|
|
|
|
```ts
|
|
|
|
|
Connections.names();
|
|
|
|
|
Connections.connection('main');
|
|
|
|
|
await Connections.openConnection('main');
|
|
|
|
|
Connections.closeConnection('main', 'manual');
|
|
|
|
|
await Connections.reconnectAll('network-restored');
|
|
|
|
|
Connections.dispose();
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
`dispose()` de la raíz cierra conexiones, cancela timers propios y limpia los
|
|
|
|
|
listeners. Después de `dispose()`, las operaciones públicas lanzan errores
|
|
|
|
|
tipados `Conn*`.
|
|
|
|
|
|
|
|
|
|
## ActiveConnections
|
|
|
|
|
|
|
|
|
|
La capa activa añade estado derivado para UI:
|
|
|
|
|
|
|
|
|
|
```ts
|
|
|
|
|
const Connections = createActiveConnections();
|
|
|
|
|
|
|
|
|
|
Connections.size;
|
|
|
|
|
Connections.activeNames;
|
|
|
|
|
Connections.states;
|
|
|
|
|
Connections.connectedNames;
|
|
|
|
|
Connections.failedNames;
|
|
|
|
|
Connections.anyConnected;
|
|
|
|
|
Connections.anyFailed;
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
`ActiveConnections` conserva la misma API de creación/consulta que el engine,
|
|
|
|
|
pero sus colecciones reflejan los cambios de estado de cada conexión.
|
|
|
|
|
|
|
|
|
|
## Integración Con App
|
|
|
|
|
|
|
|
|
|
`aapp` expone una factory porque las conexiones son app-scoped y suelen tener
|
|
|
|
|
tipos específicos del proyecto:
|
|
|
|
|
|
|
|
|
|
```ts
|
|
|
|
|
const Connections = App.createActiveConnections<AppConnections>();
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
App inyecta:
|
|
|
|
|
|
|
|
|
|
- `App.Logger`, como `Logger` común de `$libs/logr`.
|
|
|
|
|
- `App.Timers`, para reconexión, heartbeat y timeouts de ack.
|
|
|
|
|
- Un bridge estructural de sesión, si `App.Sess` existe.
|
|
|
|
|
|
|
|
|
|
Cada conexión decide si usa la sesión:
|
|
|
|
|
|
|
|
|
|
```ts
|
|
|
|
|
const Main = Connections.createConnection('main', {
|
|
|
|
|
transport: createWebSocketTransport({ url: '/realtime' }),
|
|
|
|
|
auth: () => ({ token: App.Sess?.current?.credential }),
|
|
|
|
|
session: {
|
|
|
|
|
enabled: true,
|
|
|
|
|
reauthOnRefresh: true,
|
|
|
|
|
disconnectOnExpire: true
|
|
|
|
|
}
|
|
|
|
|
});
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
Cuando `Timers` no se inyecta, `createEngineConnections()` crea un scheduler
|
|
|
|
|
privado con el mismo logger. Los diagnósticos del scheduler salen bajo la
|
|
|
|
|
categoría `timr`; los de conexiones salen bajo `conn` o `conn:<name>`.
|
|
|
|
|
|
|
|
|
|
## Estados
|
|
|
|
|
|
|
|
|
|
Estados de conexión:
|
|
|
|
|
|
|
|
|
|
```txt
|
|
|
|
|
idle -> connecting -> open
|
|
|
|
|
open -> reconnecting -> open
|
|
|
|
|
open -> closing -> closed
|
|
|
|
|
connecting/reconnecting -> failed
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
Campos útiles:
|
|
|
|
|
|
|
|
|
|
```ts
|
|
|
|
|
Main.state;
|
|
|
|
|
Main.connected;
|
|
|
|
|
Main.generation;
|
|
|
|
|
Main.error;
|
|
|
|
|
Main.openedAt;
|
|
|
|
|
Main.closedAt;
|
|
|
|
|
Main.lastMessageAt;
|
|
|
|
|
Main.reconnectAttempt;
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
`generation` cambia cuando se abre una conexión nueva. Los timeouts de ack,
|
|
|
|
|
heartbeat y reconexión usan esa generación para no resolver trabajo viejo sobre
|
|
|
|
|
una conexión nueva.
|
|
|
|
|
|
|
|
|
|
## Transports
|
|
|
|
|
|
|
|
|
|
El contrato mínimo es `ConnectionTransport`:
|
|
|
|
|
|
|
|
|
|
```ts
|
|
|
|
|
interface ConnectionTransport {
|
|
|
|
|
readonly kind: string;
|
|
|
|
|
readonly state: ConnectionTransportState;
|
|
|
|
|
readonly bufferedAmount: number;
|
|
|
|
|
readonly canSend: boolean;
|
|
|
|
|
open(): Promise<void>;
|
|
|
|
|
send(data: string | ArrayBuffer): Promise<void> | void;
|
|
|
|
|
close(code?: number, reason?: string): void;
|
|
|
|
|
onOpen(listener: () => void): () => void;
|
|
|
|
|
onMessage(listener: (message: string | ArrayBuffer) => void): () => void;
|
|
|
|
|
onClose(listener: (event: ConnectionCloseEvent) => void): () => void;
|
|
|
|
|
onError(listener: (error: unknown) => void): () => void;
|
|
|
|
|
}
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
Incluidos:
|
|
|
|
|
|
|
|
|
|
- `createWebSocketTransport()` para navegador/runtime con WebSocket.
|
|
|
|
|
- `createMockTransport()` para tests, loopback y páginas de diagnóstico.
|
|
|
|
|
|
|
|
|
|
El transporte no decide reconexión, auth, heartbeat ni canales. Solo abre,
|
|
|
|
|
envía, cierra y emite eventos.
|
|
|
|
|
|
|
|
|
|
## Frames
|
|
|
|
|
|
|
|
|
|
El frame canónico:
|
|
|
|
|
|
|
|
|
|
```ts
|
|
|
|
|
interface ConnectionFrame<TType extends string = string, TPayload = unknown> {
|
|
|
|
|
readonly id?: string;
|
|
|
|
|
readonly topic?: string;
|
|
|
|
|
readonly type: TType;
|
|
|
|
|
readonly payload: TPayload;
|
|
|
|
|
readonly ts?: number;
|
|
|
|
|
readonly ack?: boolean;
|
|
|
|
|
readonly replyTo?: string;
|
|
|
|
|
readonly error?: ConnectionFrameError;
|
|
|
|
|
}
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
El serializer por defecto es JSON y valida la forma mínima del frame. Errores
|
|
|
|
|
de encode/decode no se lanzan como strings dispersos: devuelven resultados
|
|
|
|
|
tagged o errores `ConnInvalidFrameError` según el punto de entrada.
|
|
|
|
|
|
|
|
|
|
## Send Y Request/Reply
|
|
|
|
|
|
|
|
|
|
`send()` devuelve un resultado tagged:
|
|
|
|
|
|
|
|
|
|
```ts
|
|
|
|
|
const result = await Main.send('project.updated', { id: 'p1' });
|
|
|
|
|
|
|
|
|
|
if (!result.ok) {
|
|
|
|
|
console.log(result.reason);
|
|
|
|
|
}
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
`request()` usa `ack: true`, genera un `id` y espera un frame entrante con
|
|
|
|
|
`replyTo` igual a ese id:
|
|
|
|
|
|
|
|
|
|
```ts
|
|
|
|
|
const reply = await Main.request<{ id: string }, { ok: boolean }>('project.sync', {
|
|
|
|
|
id: 'p1'
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
if (reply.ok) {
|
|
|
|
|
reply.payload.ok;
|
|
|
|
|
}
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
Razones de fallo principales:
|
|
|
|
|
|
|
|
|
|
- `timeout`
|
|
|
|
|
- `closed`
|
|
|
|
|
- `rejected`
|
|
|
|
|
- `transport_error`
|
|
|
|
|
- `invalid_reply`
|
|
|
|
|
|
|
|
|
|
## Canales
|
|
|
|
|
|
|
|
|
|
Los canales son topics nombrados dentro de una conexión. Se cachean por nombre:
|
|
|
|
|
|
|
|
|
|
```ts
|
|
|
|
|
type ProjectEvents = {
|
|
|
|
|
'project.updated': { id: string; version: number };
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
const Projects = Main.channel<ProjectEvents>('tenant:projects');
|
|
|
|
|
|
|
|
|
|
Projects.on('project.updated', (payload, meta) => {
|
|
|
|
|
console.log(payload.id, meta.receivedAt);
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
await Projects.join({ tenantId: 'acme' });
|
|
|
|
|
await Projects.send('project.updated', { id: 'p1', version: 2 });
|
|
|
|
|
await Projects.leave();
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
Estados de canal:
|
|
|
|
|
|
|
|
|
|
```txt
|
|
|
|
|
idle -> joining -> joined -> leaving -> left
|
|
|
|
|
joining -> failed
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
`dispose()` del canal limpia listeners y deja el canal en estado terminal
|
|
|
|
|
`left`.
|
|
|
|
|
|
|
|
|
|
## Reconexión
|
|
|
|
|
|
|
|
|
|
La reconexión usa `timr` y backoff configurable:
|
|
|
|
|
|
|
|
|
|
```ts
|
|
|
|
|
Connections.createConnection('main', {
|
|
|
|
|
transport,
|
|
|
|
|
reconnect: {
|
|
|
|
|
enabled: true,
|
|
|
|
|
minDelayMs: 500,
|
|
|
|
|
maxDelayMs: 15_000,
|
|
|
|
|
factor: 1.8,
|
|
|
|
|
jitterMs: 500,
|
|
|
|
|
maxAttempts: 8,
|
|
|
|
|
reconnectOnVisible: true,
|
|
|
|
|
reconnectOnOnline: true
|
|
|
|
|
}
|
|
|
|
|
});
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
Si `reconnectOnVisible` o `reconnectOnOnline` están activos, el módulo escucha
|
|
|
|
|
eventos del navegador y pide reconexión cuando la conexión está cerrada o
|
|
|
|
|
fallida. Esa capa no recibe una función `logDebug`; recibe `Diagnostics`, que
|
|
|
|
|
incluye el `Logger` completo y emite eventos catalogados.
|
|
|
|
|
|
|
|
|
|
## Heartbeat
|
|
|
|
|
|
|
|
|
|
```ts
|
|
|
|
|
Connections.createConnection('main', {
|
|
|
|
|
transport,
|
|
|
|
|
heartbeat: {
|
|
|
|
|
enabled: true,
|
|
|
|
|
intervalMs: 25_000,
|
|
|
|
|
timeoutMs: 10_000,
|
|
|
|
|
pingType: 'conn.ping',
|
|
|
|
|
pongType: 'conn.pong'
|
|
|
|
|
}
|
|
|
|
|
});
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
El heartbeat envía `pingType` periódicamente y espera `pongType`. Si vence el
|
|
|
|
|
timeout, cierra la conexión con `heartbeat_timeout` y deja que la política de
|
|
|
|
|
reconexión decida el siguiente paso.
|
|
|
|
|
|
|
|
|
|
## Auth Y Sesión
|
|
|
|
|
|
|
|
|
|
Auth de conexión:
|
|
|
|
|
|
|
|
|
|
```ts
|
|
|
|
|
Connections.createConnection('main', {
|
|
|
|
|
transport,
|
|
|
|
|
auth: {
|
|
|
|
|
getAuth: () => ({ token }),
|
|
|
|
|
authType: 'conn.auth',
|
|
|
|
|
timeoutMs: 10_000
|
|
|
|
|
}
|
|
|
|
|
});
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
El resultado de auth es tagged:
|
|
|
|
|
|
|
|
|
|
```ts
|
|
|
|
|
await Main.reauthenticate(); // { ok: true } | { ok: false, reason, error? }
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
Con `session.enabled`, el bridge de sesión puede:
|
|
|
|
|
|
|
|
|
|
- reautenticar cuando `sess` emite refresh;
|
|
|
|
|
- desconectar cuando la sesión expira;
|
|
|
|
|
- desconectar cuando la sesión se revoca.
|
|
|
|
|
|
|
|
|
|
`conn` no crea sesiones ni decide permisos. En servidor, los joins/sends de un
|
|
|
|
|
canal deben validarse con `auth/sess/perm`.
|
|
|
|
|
|
|
|
|
|
## Buffer
|
|
|
|
|
|
|
|
|
|
Cuando la conexión no está abierta, `send()` puede comportarse según policy:
|
|
|
|
|
|
|
|
|
|
```ts
|
|
|
|
|
buffer: {
|
|
|
|
|
policy: 'buffer', // 'buffer' | 'drop' | 'fail'
|
|
|
|
|
maxMessages: 100,
|
|
|
|
|
maxBytes: 1_000_000
|
|
|
|
|
}
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
- `buffer`: encola y drena al abrir.
|
|
|
|
|
- `drop`: acepta la llamada pero descarta.
|
|
|
|
|
- `fail`: devuelve `{ ok: false, reason: 'closed' }`.
|
|
|
|
|
|
|
|
|
|
## Diagnostics Y Logger
|
|
|
|
|
|
|
|
|
|
Las opciones públicas aceptan `logger?: Logger` desde `$libs/logr`. No existe
|
|
|
|
|
un `ConnectionLogger` propio.
|
|
|
|
|
|
|
|
|
|
```ts
|
|
|
|
|
import type { Logger } from '$libs/logr';
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
Internamente `conn` usa `createConnectionDiagnostics(logger)` y eventos
|
|
|
|
|
catalogados en `CONNECTION_DIAGNOSTIC_EVENTS`:
|
|
|
|
|
|
|
|
|
|
- `connect_failed`
|
|
|
|
|
- `transport_error`
|
|
|
|
|
- `send_failed`
|
|
|
|
|
- `frame_decode_failed`
|
|
|
|
|
- `frame_encode_failed`
|
|
|
|
|
- `auth_failed`
|
|
|
|
|
- `reauth_failed`
|
|
|
|
|
- `heartbeat_timeout`
|
|
|
|
|
- `reconnect_exhausted`
|
|
|
|
|
- `browser_reconnect`
|
|
|
|
|
- `session_refreshed`
|
|
|
|
|
- `session_expired`
|
|
|
|
|
- `session_revoked`
|
|
|
|
|
- `listener_threw`
|
|
|
|
|
|
|
|
|
|
Los mensajes y categorías viven en `consts.ts`. Si quieres enviar eventos a
|
|
|
|
|
Sentry, Loki o Datadog, inyecta un `EngineLogger` con el transporte adecuado;
|
|
|
|
|
`conn` solo emite al contrato común.
|
|
|
|
|
|
|
|
|
|
## Errores
|
|
|
|
|
|
|
|
|
|
Programmer errors lanzan clases tipadas:
|
|
|
|
|
|
|
|
|
|
- `ConnDisposedError`
|
|
|
|
|
- `ConnConnectionAlreadyExistsError`
|
|
|
|
|
- `ConnConnectionNotFoundError`
|
|
|
|
|
- `ConnInvalidConnectionNameError`
|
|
|
|
|
- `ConnInvalidFrameError`
|
|
|
|
|
- `ConnChannelAlreadyExistsError`
|
|
|
|
|
- `ConnChannelNotFoundError`
|
|
|
|
|
- `ConnWebSocketUnavailableError`
|
|
|
|
|
|
|
|
|
|
Fallos runtime de transporte/envío/auth/request devuelven resultados tagged
|
|
|
|
|
para que el consumidor pueda decidir sin `try/catch` obligatorio.
|
|
|
|
|
|
|
|
|
|
## Página De Prueba
|
|
|
|
|
|
|
|
|
|
La página interactiva está en:
|
|
|
|
|
|
|
|
|
|
```txt
|
|
|
|
|
/test/conn
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
Incluye chat WebSocket real usando `scripts/conn-chat-server.mjs`:
|
|
|
|
|
|
|
|
|
|
```txt
|
|
|
|
|
npm run dev:conn-chat
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
También `/test/ecosystem` usa `conn` dentro de la demo total con transporte
|
|
|
|
|
mock loopback para probar integración con `aapp`, `timr`, `logr`, `perm` y
|
|
|
|
|
`cach`.
|
|
|
|
|
|
|
|
|
|
## Testing
|
|
|
|
|
|
|
|
|
|
Tests principales:
|
|
|
|
|
|
|
|
|
|
```txt
|
|
|
|
|
src/arts/conn/test/engine-connections.test.ts
|
|
|
|
|
src/arts/conn/test/active-connections.test.ts
|
|
|
|
|
src/arts/conn/test/connection.test.ts
|
|
|
|
|
src/arts/conn/test/connection-state.test.ts
|
|
|
|
|
src/arts/conn/test/websocket.test.ts
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
Casos que deben mantenerse cubiertos:
|
|
|
|
|
|
|
|
|
|
- creación duplicada y lookup inexistente;
|
|
|
|
|
- `this` no requerido al desestructurar métodos del engine;
|
|
|
|
|
- reloj inyectado para timestamps deterministas;
|
|
|
|
|
- reconnect con fake timers;
|
|
|
|
|
- heartbeat timeout;
|
|
|
|
|
- request/reply con timeout y rechazo;
|
|
|
|
|
- channel join/leave/dispose;
|
|
|
|
|
- bridge de sesión refresh/expire/revoke;
|
|
|
|
|
- serializer inválido y transport errors.
|