You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
1988 lines
43 KiB
1988 lines
43 KiB
|
3 months ago
|
---
|
||
|
|
title: arts/connection — design & implementation record
|
||
|
|
type: decision-log
|
||
|
|
audience: human + agent
|
||
|
|
authority: historical design record — the realtime connection registry (transports, reconnect, heartbeat, request/reply, channels, session bridge); superseded by src/arts/connection/README.md as the current reference
|
||
|
|
status: historical
|
||
|
|
source: moved from src/arts/connection/DESIGN_CONN.md (2026-07-03, arts-docs-reconciliation B2; kept verbatim)
|
||
|
|
---
|
||
|
|
|
||
|
|
# `arts/connection` — Design & Implementation Document
|
||
|
|
|
||
|
|
> **Historical note 2026-05-14.** This document predates the directory rename
|
||
|
|
> sweep. The canonical artifact name is now `connection`; `$connection` points
|
||
|
|
> to `src/arts/connection`. Mentions of `conn`, `sess`, `timr`, `logr`,
|
||
|
|
> `cach` and `aapp` below are historical design context unless they refer to
|
||
|
|
> literal wire/event keys.
|
||
|
|
>
|
||
|
|
> **Purpose.** Implement a zero-dependency runtime artifact for managing multiple realtime connections at App level.
|
||
|
|
>
|
||
|
|
> **Audience.** Claude Code / implementer familiar with TypeScript, Svelte 5 runes, SvelteKit, and the existing `src/arts/*` artifact architecture.
|
||
|
|
>
|
||
|
|
> **Status.** Design target for first implementation.
|
||
|
|
>
|
||
|
|
> **Core decision.** The artifact is called `conn` / `connections`, not `sock` / `sockets`, because the public model must support WebSocket, SSE, polling, WebTransport, mocks, and future transports through adapters.
|
||
|
|
>
|
||
|
|
> **Naming correction applied.** `EngineConnections` / `ActiveConnections` are the root registry names. An individual runtime unit is `Connection`. Any older references below to `EngineConnection` / `ActiveConnection` describe that unit conceptually, but are not implementation names.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 1. Problem statement
|
||
|
|
|
||
|
|
The application needs first-class support for realtime connections at App level.
|
||
|
|
|
||
|
|
The system must support:
|
||
|
|
|
||
|
|
- Multiple independent connections.
|
||
|
|
- Each connection with its own lifecycle, transport, auth, reconnect policy, heartbeat, buffer policy, serializer, channels and event listeners.
|
||
|
|
- Different transport types: WebSocket now, SSE/polling/mock soon, WebTransport later.
|
||
|
|
- Integration with `ActiveApp`, `Session`, `Http`, `Logger`, and browser lifecycle when configured.
|
||
|
|
- Strict TypeScript types.
|
||
|
|
- Zero dependencies.
|
||
|
|
- No domain state inside `arts/conn`.
|
||
|
|
|
||
|
|
The correct mental model is:
|
||
|
|
|
||
|
|
```txt
|
||
|
|
App
|
||
|
|
└── Connections
|
||
|
|
├── connection('main')
|
||
|
|
│ ├── channel('orders')
|
||
|
|
│ └── channel('cart')
|
||
|
|
│
|
||
|
|
├── connection('chat')
|
||
|
|
│ ├── channel('room:123')
|
||
|
|
│ └── channel('presence:room:123')
|
||
|
|
│
|
||
|
|
└── connection('market')
|
||
|
|
└── channel('prices')
|
||
|
|
```
|
||
|
|
|
||
|
|
A connection is not necessarily a WebSocket. A connection is an abstract realtime transport endpoint.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 1.1 Mandatory v1 scope
|
||
|
|
|
||
|
|
The first implementation of `arts/conn` MUST include the following capabilities.
|
||
|
|
|
||
|
|
This list is normative. If any item is missing, the implementation is incomplete.
|
||
|
|
|
||
|
|
### 1.1.1 Multi-connection registry
|
||
|
|
|
||
|
|
`arts/conn` must provide an App-level registry of named connections.
|
||
|
|
|
||
|
|
Required API:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
const Connections = App.createActiveConnections<AppConnections>();
|
||
|
|
|
||
|
|
const Main = Connections.createConnection('main', options);
|
||
|
|
const Chat = Connections.createConnection('chat', options);
|
||
|
|
const Market = Connections.createConnection('market', options);
|
||
|
|
|
||
|
|
Connections.connection('main');
|
||
|
|
Connections.hasConnection('main');
|
||
|
|
Connections.connectionNames();
|
||
|
|
```
|
||
|
|
|
||
|
|
Required behavior:
|
||
|
|
|
||
|
|
- More than one connection may exist at the same time.
|
||
|
|
- Each connection has independent lifecycle/state/config.
|
||
|
|
- Connections are addressed by stable string names.
|
||
|
|
- Duplicate names are rejected.
|
||
|
|
- Missing names throw a typed error.
|
||
|
|
- The registry can open, close and reconnect connections by name.
|
||
|
|
|
||
|
|
Required hub methods:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
openConnection(name): Promise<ConnectionConnectResult>;
|
||
|
|
closeConnection(name, reason?: string): void;
|
||
|
|
reconnectConnection(name, reason?: string): Promise<ConnectionConnectResult>;
|
||
|
|
|
||
|
|
openAll(): Promise<readonly ConnectionConnectResult[]>;
|
||
|
|
closeAll(reason?: string): void;
|
||
|
|
reconnectAll(reason?: string): Promise<readonly ConnectionConnectResult[]>;
|
||
|
|
```
|
||
|
|
|
||
|
|
### 1.1.2 Transport adapter
|
||
|
|
|
||
|
|
The connection engine MUST be transport-agnostic.
|
||
|
|
|
||
|
|
The engine may not hardcode WebSocket behavior directly.
|
||
|
|
|
||
|
|
Every physical transport must implement:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export 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;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Required behavior:
|
||
|
|
|
||
|
|
- `EngineConnection` only talks to `ConnectionTransport`.
|
||
|
|
- Reconnect, heartbeat, ack, channels and buffering work above the transport.
|
||
|
|
- Future transports must be addable without changing the engine core.
|
||
|
|
|
||
|
|
### 1.1.3 WebSocket transport
|
||
|
|
|
||
|
|
v1 MUST include:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
createWebSocketTransport(options);
|
||
|
|
```
|
||
|
|
|
||
|
|
Required behavior:
|
||
|
|
|
||
|
|
- Opens a native `WebSocket`.
|
||
|
|
- Sends string or binary payloads.
|
||
|
|
- Receives messages and emits them through the adapter interface.
|
||
|
|
- Exposes `bufferedAmount`.
|
||
|
|
- Maps browser `open`, `message`, `close`, `error` events.
|
||
|
|
- Supports dynamic URL:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
url: string | (() => string);
|
||
|
|
```
|
||
|
|
|
||
|
|
Required usage:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
const Main = Connections.createConnection('main', {
|
||
|
|
transport: createWebSocketTransport({
|
||
|
|
url: () => '/realtime'
|
||
|
|
})
|
||
|
|
});
|
||
|
|
```
|
||
|
|
|
||
|
|
### 1.1.4 Mock transport
|
||
|
|
|
||
|
|
v1 MUST include:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
createMockTransport();
|
||
|
|
```
|
||
|
|
|
||
|
|
Required behavior:
|
||
|
|
|
||
|
|
- Implements `ConnectionTransport`.
|
||
|
|
- Allows deterministic tests without browser/network.
|
||
|
|
- Can manually emit open/message/close/error events.
|
||
|
|
- Captures sent messages for assertions.
|
||
|
|
|
||
|
|
Required testing helpers:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
transport.emitOpen();
|
||
|
|
transport.emitMessage(raw);
|
||
|
|
transport.emitClose(event);
|
||
|
|
transport.emitError(error);
|
||
|
|
transport.sentMessages();
|
||
|
|
```
|
||
|
|
|
||
|
|
The engine test suite should primarily use the mock transport.
|
||
|
|
|
||
|
|
### 1.1.5 State per connection
|
||
|
|
|
||
|
|
Each connection MUST expose its own state.
|
||
|
|
|
||
|
|
Required state:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export type ConnectionState =
|
||
|
|
| 'idle'
|
||
|
|
| 'connecting'
|
||
|
|
| 'open'
|
||
|
|
| 'reconnecting'
|
||
|
|
| 'closing'
|
||
|
|
| 'closed'
|
||
|
|
| 'failed';
|
||
|
|
```
|
||
|
|
|
||
|
|
Required `EngineConnection`/`ActiveConnection` properties:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
readonly state: ConnectionState;
|
||
|
|
readonly connected: boolean;
|
||
|
|
readonly generation: number;
|
||
|
|
readonly error: unknown | null;
|
||
|
|
readonly openedAt: number | null;
|
||
|
|
readonly closedAt: number | null;
|
||
|
|
readonly lastMessageAt: number | null;
|
||
|
|
readonly reconnectAttempt: number;
|
||
|
|
```
|
||
|
|
|
||
|
|
Required event:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
onState(listener: (change: ConnectionStateChange) => void): () => void;
|
||
|
|
```
|
||
|
|
|
||
|
|
Required change shape:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export interface ConnectionStateChange {
|
||
|
|
readonly connection: string;
|
||
|
|
readonly from: ConnectionState;
|
||
|
|
readonly to: ConnectionState;
|
||
|
|
readonly generation: number;
|
||
|
|
readonly error?: unknown;
|
||
|
|
readonly at: number;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
### 1.1.6 Aggregated reactive state
|
||
|
|
|
||
|
|
`ActiveConnections` MUST expose an aggregated reactive view of all registered connections.
|
||
|
|
|
||
|
|
Required properties:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
readonly size: number;
|
||
|
|
readonly activeNames: readonly string[];
|
||
|
|
|
||
|
|
readonly states: Readonly<Record<string, ConnectionState>>;
|
||
|
|
|
||
|
|
readonly connectedNames: readonly string[];
|
||
|
|
readonly connectingNames: readonly string[];
|
||
|
|
readonly reconnectingNames: readonly string[];
|
||
|
|
readonly failedNames: readonly string[];
|
||
|
|
readonly closedNames: readonly string[];
|
||
|
|
|
||
|
|
readonly allConnected: boolean;
|
||
|
|
readonly anyConnected: boolean;
|
||
|
|
readonly anyConnecting: boolean;
|
||
|
|
readonly anyReconnecting: boolean;
|
||
|
|
readonly anyFailed: boolean;
|
||
|
|
```
|
||
|
|
|
||
|
|
Required behavior:
|
||
|
|
|
||
|
|
- `Connection` owns the real state.
|
||
|
|
- `Connections` observes registered connections and derives aggregate state.
|
||
|
|
- There must not be a single ambiguous `Connections.state`.
|
||
|
|
|
||
|
|
### 1.1.7 Send / request / ack
|
||
|
|
|
||
|
|
v1 MUST support fire-and-forget sends and request/reply acks.
|
||
|
|
|
||
|
|
Required APIs:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
send<TPayload>(
|
||
|
|
type: string,
|
||
|
|
payload: TPayload,
|
||
|
|
options?: ConnectionSendOptions
|
||
|
|
): Promise<ConnectionSendResult>;
|
||
|
|
|
||
|
|
request<TPayload, TResult = unknown>(
|
||
|
|
type: string,
|
||
|
|
payload: TPayload,
|
||
|
|
options?: ConnectionRequestOptions
|
||
|
|
): Promise<ConnectionAckResult<TResult>>;
|
||
|
|
```
|
||
|
|
|
||
|
|
Required frame fields:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
id?: string;
|
||
|
|
type: string;
|
||
|
|
payload: unknown;
|
||
|
|
ack?: boolean;
|
||
|
|
replyTo?: string;
|
||
|
|
error?: ConnectionFrameError;
|
||
|
|
```
|
||
|
|
|
||
|
|
Required behavior:
|
||
|
|
|
||
|
|
- `request()` generates a message id.
|
||
|
|
- `request()` sends a frame with `ack: true`.
|
||
|
|
- Incoming frames with `replyTo` resolve the matching pending request.
|
||
|
|
- Requests time out using `timeoutMs`.
|
||
|
|
- Disconnect resolves all pending requests as closed.
|
||
|
|
- Normal runtime failures return tagged results, not uncaught exceptions.
|
||
|
|
|
||
|
|
### 1.1.8 Reconnect
|
||
|
|
|
||
|
|
v1 MUST support configurable reconnect.
|
||
|
|
|
||
|
|
Required options:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export interface ConnectionReconnectOptions {
|
||
|
|
readonly enabled?: boolean;
|
||
|
|
readonly minDelayMs?: number;
|
||
|
|
readonly maxDelayMs?: number;
|
||
|
|
readonly factor?: number;
|
||
|
|
readonly jitterMs?: number;
|
||
|
|
readonly maxAttempts?: number;
|
||
|
|
readonly reconnectOnVisible?: boolean;
|
||
|
|
readonly reconnectOnOnline?: boolean;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Required behavior:
|
||
|
|
|
||
|
|
- Unexpected close triggers reconnect when enabled.
|
||
|
|
- Intentional `disconnect()` does not reconnect.
|
||
|
|
- Backoff uses factor + jitter.
|
||
|
|
- Successful reconnect resets attempts.
|
||
|
|
- Reconnect increments connection generation.
|
||
|
|
- After reconnect, configured channels rejoin.
|
||
|
|
- After reconnect, auth re-runs if configured.
|
||
|
|
- Reconnect delay must be computed through the shared timer/backoff primitives
|
||
|
|
(`$libs/timers.computeBackoffDelay`) and scheduled through injected
|
||
|
|
`TimerScheduler`/`App.timers`, not raw `setTimeout` inside the connection
|
||
|
|
engine. A native fallback is allowed only in the standalone engine path.
|
||
|
|
|
||
|
|
### 1.1.9 Heartbeat
|
||
|
|
|
||
|
|
v1 MUST support heartbeat.
|
||
|
|
|
||
|
|
Required options:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export interface ConnectionHeartbeatOptions {
|
||
|
|
readonly enabled?: boolean;
|
||
|
|
readonly intervalMs?: number;
|
||
|
|
readonly timeoutMs?: number;
|
||
|
|
readonly pingType?: string;
|
||
|
|
readonly pongType?: string;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Required behavior:
|
||
|
|
|
||
|
|
- When connection is open, heartbeat sends ping frames if transport can send.
|
||
|
|
- Receiving pong or any message updates liveness.
|
||
|
|
- Timeout closes/reconnects according to reconnect policy.
|
||
|
|
- Heartbeat can be disabled per connection.
|
||
|
|
|
||
|
|
### 1.1.10 Channels
|
||
|
|
|
||
|
|
v1 MUST support logical channels/topics inside each connection.
|
||
|
|
|
||
|
|
Required API:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
const Orders = Main.channel('orders');
|
||
|
|
const Cart = Main.channel('cart');
|
||
|
|
```
|
||
|
|
|
||
|
|
Required behavior:
|
||
|
|
|
||
|
|
- `connection.channel(name)` caches by name.
|
||
|
|
- Same channel name returns same channel instance.
|
||
|
|
- Frames with `topic` are routed to the matching channel.
|
||
|
|
- Frames without `topic` go to connection-level listeners.
|
||
|
|
- Channels expose `join`, `leave`, `send`, `request`, `on`, `onAny`.
|
||
|
|
- Channels may auto-join and rejoin after reconnect.
|
||
|
|
|
||
|
|
Required channel methods:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
join(params?: unknown): Promise<ConnectionChannelJoinResult>;
|
||
|
|
leave(): Promise<void>;
|
||
|
|
|
||
|
|
send(type, payload, options?): Promise<ConnectionSendResult>;
|
||
|
|
request(type, payload, options?): Promise<ConnectionAckResult>;
|
||
|
|
|
||
|
|
on(type, listener): () => void;
|
||
|
|
onAny(listener): () => void;
|
||
|
|
```
|
||
|
|
|
||
|
|
### 1.1.11 Session integration opt-in
|
||
|
|
|
||
|
|
Session integration MUST be per connection and opt-in.
|
||
|
|
|
||
|
|
Required option:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
session:
|
||
|
|
| false
|
||
|
|
| {
|
||
|
|
readonly enabled?: boolean;
|
||
|
|
readonly reauthOnRefresh?: boolean;
|
||
|
|
readonly disconnectOnExpire?: boolean;
|
||
|
|
};
|
||
|
|
```
|
||
|
|
|
||
|
|
Required behavior:
|
||
|
|
|
||
|
|
- Connections with `session: false` never react to session events.
|
||
|
|
- Connections with session integration may reauth on `REFRESHED`.
|
||
|
|
- Connections with session integration may disconnect on `EXPIRED`/`REVOKED`.
|
||
|
|
- No connection should assume JWT/Bearer.
|
||
|
|
- Auth payload is provided by user code.
|
||
|
|
|
||
|
|
Required auth shape:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
auth?: {
|
||
|
|
getAuth?: () =>
|
||
|
|
| Record<string, unknown>
|
||
|
|
| null
|
||
|
|
| Promise<Record<string, unknown> | null>;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
or shorthand:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
auth: () => Record<string, unknown> | null | Promise<Record<string, unknown> | null>;
|
||
|
|
```
|
||
|
|
|
||
|
|
### 1.1.12 Minimum v1 deliverables
|
||
|
|
|
||
|
|
The first implementation is not accepted unless these files exist and are tested:
|
||
|
|
|
||
|
|
```txt
|
||
|
|
engine-connections.ts
|
||
|
|
active-connections.svelte.ts
|
||
|
|
engine-connection.ts
|
||
|
|
active-connection.svelte.ts
|
||
|
|
transport.ts
|
||
|
|
transports/websocket.ts
|
||
|
|
transports/mock.ts
|
||
|
|
serializer.ts
|
||
|
|
serializers/json.ts
|
||
|
|
channel.ts
|
||
|
|
reconnect.ts
|
||
|
|
heartbeat.ts
|
||
|
|
ack.ts
|
||
|
|
backpressure.ts
|
||
|
|
types.ts
|
||
|
|
errors.ts
|
||
|
|
consts.ts
|
||
|
|
```
|
||
|
|
|
||
|
|
Presence, SSE and polling may be deferred, but WebSocket + mock are mandatory.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 2. Explicit non-goals
|
||
|
|
|
||
|
|
`arts/conn` must NOT implement:
|
||
|
|
|
||
|
|
- Chat UI.
|
||
|
|
- Notification UI.
|
||
|
|
- Business/domain event models.
|
||
|
|
- Authorization policies.
|
||
|
|
- User/session storage.
|
||
|
|
- Webhook receiving/signing.
|
||
|
|
- Durable queues.
|
||
|
|
- CRDTs/collaboration protocols.
|
||
|
|
- Server-side room management.
|
||
|
|
- Provider-specific SDK logic for Ably, Pusher, Supabase, etc.
|
||
|
|
- Socket.IO protocol compatibility.
|
||
|
|
|
||
|
|
Webhooks are a separate concern and should become a separate artifact, for example:
|
||
|
|
|
||
|
|
```txt
|
||
|
|
arts/webh
|
||
|
|
```
|
||
|
|
|
||
|
|
Do not put webhooks inside `arts/conn`.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 3. Relationship with existing artifacts
|
||
|
|
|
||
|
|
### 3.1 With `aapp`
|
||
|
|
|
||
|
|
`ActiveApp` is the composition root.
|
||
|
|
|
||
|
|
`connections` should live on App, not on Session:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
const Connections = App.createActiveConnections<AppConnections>();
|
||
|
|
```
|
||
|
|
|
||
|
|
or, once integrated:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
App.connections.connection('main');
|
||
|
|
```
|
||
|
|
|
||
|
|
Do **not** put this under:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
Sess.connections;
|
||
|
|
Sess.socket;
|
||
|
|
Sess.realtime;
|
||
|
|
```
|
||
|
|
|
||
|
|
Session is only a source of auth/lifecycle events.
|
||
|
|
|
||
|
|
### 3.2 With `sess`
|
||
|
|
|
||
|
|
`arts/sess` owns session lifecycle:
|
||
|
|
|
||
|
|
```txt
|
||
|
|
ADOPTED
|
||
|
|
REFRESHED
|
||
|
|
EXPIRED
|
||
|
|
REVOKED
|
||
|
|
EXTERNAL_CHANGED
|
||
|
|
```
|
||
|
|
|
||
|
|
`arts/conn` may subscribe to those events through App auto-wiring, but only per connection and only if configured.
|
||
|
|
|
||
|
|
Example:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
connections: {
|
||
|
|
main: {
|
||
|
|
session: {
|
||
|
|
enabled: true,
|
||
|
|
reauthOnRefresh: true,
|
||
|
|
disconnectOnExpire: true
|
||
|
|
}
|
||
|
|
},
|
||
|
|
|
||
|
|
market: {
|
||
|
|
session: false
|
||
|
|
}
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Public connections may ignore session entirely.
|
||
|
|
|
||
|
|
### 3.3 With `http`
|
||
|
|
|
||
|
|
`arts/conn` should not depend on `arts/http`, but transport adapters may use an HTTP client if injected.
|
||
|
|
|
||
|
|
Example SSE adapter:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
createSseTransport({
|
||
|
|
url: () => '/events',
|
||
|
|
send: async (frame) => {
|
||
|
|
await App.http.post('/events/send', { body: frame });
|
||
|
|
}
|
||
|
|
});
|
||
|
|
```
|
||
|
|
|
||
|
|
### 3.4 With `logr`
|
||
|
|
|
||
|
|
Logger injection is allowed.
|
||
|
|
|
||
|
|
All logs should use a stable scope, for example:
|
||
|
|
|
||
|
|
```txt
|
||
|
|
conn
|
||
|
|
conn:main
|
||
|
|
conn:main:orders
|
||
|
|
```
|
||
|
|
|
||
|
|
### 3.5 With `stor`
|
||
|
|
|
||
|
|
No persistence by default.
|
||
|
|
|
||
|
|
Connection state is runtime state. Do not persist connections in storage.
|
||
|
|
|
||
|
|
### 3.6 With `timr`
|
||
|
|
|
||
|
|
`arts/conn` should use `App.timers` when built from `aapp`. Reconnect,
|
||
|
|
heartbeat and pending-ack timeouts must be keyed timers so app-level debug
|
||
|
|
panels and `App.dispose()` can see and cancel them uniformly.
|
||
|
|
|
||
|
|
Recommended key shape:
|
||
|
|
|
||
|
|
```txt
|
||
|
|
conn:<connection>:reconnect
|
||
|
|
conn:<connection>:heartbeat
|
||
|
|
conn:<connection>:ack:<id>
|
||
|
|
```
|
||
|
|
|
||
|
|
The engine accepts a structural scheduler option compatible with
|
||
|
|
`TimerScheduler` instead of importing Svelte-specific active wrappers.
|
||
|
|
|
||
|
|
### 3.7 With `libs/http` and `libs/timers`
|
||
|
|
|
||
|
|
Protocol vocabulary that is not connection-specific belongs in `libs/*`.
|
||
|
|
|
||
|
|
- HTTP method/header/status names used by polling/SSE transports must come
|
||
|
|
from `$libs/http`.
|
||
|
|
- Reconnect backoff must reuse `$libs/timers`.
|
||
|
|
- `arts/conn` owns only connection-specific constants: states, events,
|
||
|
|
transport kinds, frame keys, close reasons, logger messages and error names.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 4. Naming
|
||
|
|
|
||
|
|
Artifact:
|
||
|
|
|
||
|
|
```txt
|
||
|
|
src/arts/connection/
|
||
|
|
```
|
||
|
|
|
||
|
|
Public names:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
EngineConnections;
|
||
|
|
ActiveConnections;
|
||
|
|
|
||
|
|
EngineConnection;
|
||
|
|
ActiveConnection;
|
||
|
|
|
||
|
|
ConnectionChannel;
|
||
|
|
ConnectionPresence;
|
||
|
|
|
||
|
|
ConnectionTransport;
|
||
|
|
ConnectionSerializer;
|
||
|
|
```
|
||
|
|
|
||
|
|
Factory names:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
createEngineConnections;
|
||
|
|
createActiveConnections;
|
||
|
|
|
||
|
|
createWebSocketTransport;
|
||
|
|
createSseTransport;
|
||
|
|
createPollingTransport;
|
||
|
|
createMockTransport;
|
||
|
|
|
||
|
|
jsonConnectionSerializer;
|
||
|
|
```
|
||
|
|
|
||
|
|
Avoid names such as:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
SocketService;
|
||
|
|
RealtimeManager;
|
||
|
|
ConnectionManager;
|
||
|
|
```
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 5. Directory structure
|
||
|
|
|
||
|
|
Target structure:
|
||
|
|
|
||
|
|
```txt
|
||
|
|
src/arts/connection/
|
||
|
|
index.ts
|
||
|
|
|
||
|
|
consts.ts
|
||
|
|
types.ts
|
||
|
|
errors.ts
|
||
|
|
|
||
|
|
engine-connections.ts
|
||
|
|
active-connections.svelte.ts
|
||
|
|
|
||
|
|
engine-connection.ts
|
||
|
|
active-connection.svelte.ts
|
||
|
|
|
||
|
|
channel.ts
|
||
|
|
presence.ts
|
||
|
|
|
||
|
|
transport.ts
|
||
|
|
transports/
|
||
|
|
websocket.ts
|
||
|
|
sse.ts
|
||
|
|
polling.ts
|
||
|
|
mock.ts
|
||
|
|
|
||
|
|
serializer.ts
|
||
|
|
serializers/
|
||
|
|
json.ts
|
||
|
|
|
||
|
|
reconnect.ts
|
||
|
|
heartbeat.ts
|
||
|
|
backpressure.ts
|
||
|
|
ack.ts
|
||
|
|
|
||
|
|
app-integration.ts
|
||
|
|
|
||
|
|
engine-connections.test.ts
|
||
|
|
engine-connection.test.ts
|
||
|
|
channel.test.ts
|
||
|
|
transports.test.ts
|
||
|
|
```
|
||
|
|
|
||
|
|
Rules:
|
||
|
|
|
||
|
|
- `.ts` files must be runes-free.
|
||
|
|
- `.svelte.ts` files may use `$state`.
|
||
|
|
- `index.ts` must not accidentally force server-side imports of runes files where inappropriate.
|
||
|
|
- Prefer explicit file imports in App integration if necessary.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 6. Core abstractions
|
||
|
|
|
||
|
|
### 6.1 Connection hub
|
||
|
|
|
||
|
|
The hub owns a registry of named connections.
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export interface EngineConnections<TConnections extends ConnectionMap = ConnectionMap> {
|
||
|
|
createConnection<K extends keyof TConnections & string>(
|
||
|
|
name: K,
|
||
|
|
options: ConnectionOptions<TConnections[K]>
|
||
|
|
): EngineConnection<TConnections[K]>;
|
||
|
|
|
||
|
|
createConnection<TChannels extends ConnectionChannelMap = ConnectionChannelMap>(
|
||
|
|
name: string,
|
||
|
|
options: ConnectionOptions<TChannels>
|
||
|
|
): EngineConnection<TChannels>;
|
||
|
|
|
||
|
|
connection<K extends keyof TConnections & string>(name: K): EngineConnection<TConnections[K]>;
|
||
|
|
|
||
|
|
connection<TChannels extends ConnectionChannelMap = ConnectionChannelMap>(
|
||
|
|
name: string
|
||
|
|
): EngineConnection<TChannels>;
|
||
|
|
|
||
|
|
has(name: string): boolean;
|
||
|
|
|
||
|
|
names(): readonly string[];
|
||
|
|
|
||
|
|
close(name: string, reason?: string): void;
|
||
|
|
|
||
|
|
closeAll(reason?: string): void;
|
||
|
|
|
||
|
|
dispose(): void;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Behavior:
|
||
|
|
|
||
|
|
- `createConnection(name, options)` creates and registers a connection.
|
||
|
|
- Creating the same name twice should throw `ConnectionAlreadyExistsError`.
|
||
|
|
- `connection(name)` returns an existing connection or throws `ConnectionNotFoundError`.
|
||
|
|
- `closeAll()` disconnects all registered connections.
|
||
|
|
- `dispose()` disconnects everything and clears listeners.
|
||
|
|
|
||
|
|
### 6.1.1 Hub lifecycle methods
|
||
|
|
|
||
|
|
In addition to create/lookup methods, the hub MUST expose lifecycle methods by connection name:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
openConnection(name: string): Promise<ConnectionConnectResult>;
|
||
|
|
closeConnection(name: string, reason?: string): void;
|
||
|
|
reconnectConnection(name: string, reason?: string): Promise<ConnectionConnectResult>;
|
||
|
|
|
||
|
|
openAll(): Promise<readonly ConnectionConnectResult[]>;
|
||
|
|
closeAll(reason?: string): void;
|
||
|
|
reconnectAll(reason?: string): Promise<readonly ConnectionConnectResult[]>;
|
||
|
|
```
|
||
|
|
|
||
|
|
These are convenience methods over:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
Connections.connection(name).connect();
|
||
|
|
Connections.connection(name).disconnect();
|
||
|
|
Connections.connection(name).reconnect();
|
||
|
|
```
|
||
|
|
|
||
|
|
Both levels are intentionally supported.
|
||
|
|
|
||
|
|
### 6.2 Active hub
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export interface ActiveConnections<
|
||
|
|
TConnections extends ConnectionMap = ConnectionMap
|
||
|
|
> extends EngineConnections<TConnections> {
|
||
|
|
readonly size: number;
|
||
|
|
readonly activeNames: readonly string[];
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
`ActiveConnections` wraps `EngineConnections` but exposes reactive derived state.
|
||
|
|
|
||
|
|
Important:
|
||
|
|
|
||
|
|
- `createEngineConnections()` must never create active/runes objects.
|
||
|
|
- `createActiveConnections()` may internally create `EngineConnections`.
|
||
|
|
|
||
|
|
Correct:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
const Connections = createActiveConnections();
|
||
|
|
const Main = Connections.createConnection('main', options);
|
||
|
|
```
|
||
|
|
|
||
|
|
Avoid:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
createEngineConnections().createActiveConnection(...)
|
||
|
|
```
|
||
|
|
|
||
|
|
That mixes layers.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 7. Connection abstraction
|
||
|
|
|
||
|
|
Each connection has its own:
|
||
|
|
|
||
|
|
- name
|
||
|
|
- transport
|
||
|
|
- serializer
|
||
|
|
- state
|
||
|
|
- reconnect policy
|
||
|
|
- heartbeat policy
|
||
|
|
- buffer policy
|
||
|
|
- auth provider
|
||
|
|
- channel registry
|
||
|
|
- event listeners
|
||
|
|
- pending acks/requests
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export interface EngineConnection<TChannels extends ConnectionChannelMap = ConnectionChannelMap> {
|
||
|
|
readonly name: string;
|
||
|
|
readonly state: ConnectionState;
|
||
|
|
readonly connected: boolean;
|
||
|
|
readonly generation: number;
|
||
|
|
|
||
|
|
connect(): Promise<ConnectionConnectResult>;
|
||
|
|
|
||
|
|
disconnect(reason?: string): void;
|
||
|
|
|
||
|
|
reconnect(reason?: string): Promise<ConnectionConnectResult>;
|
||
|
|
|
||
|
|
reauthenticate(): Promise<ConnectionAuthResult>;
|
||
|
|
|
||
|
|
send<TPayload>(
|
||
|
|
type: string,
|
||
|
|
payload: TPayload,
|
||
|
|
options?: ConnectionSendOptions
|
||
|
|
): Promise<ConnectionSendResult>;
|
||
|
|
|
||
|
|
request<TPayload, TResult = unknown>(
|
||
|
|
type: string,
|
||
|
|
payload: TPayload,
|
||
|
|
options?: ConnectionRequestOptions
|
||
|
|
): Promise<ConnectionAckResult<TResult>>;
|
||
|
|
|
||
|
|
channel<K extends keyof TChannels & string>(
|
||
|
|
name: K,
|
||
|
|
options?: ConnectionChannelOptions
|
||
|
|
): ConnectionChannel<TChannels[K]>;
|
||
|
|
|
||
|
|
channel<TEvents extends ConnectionEventMap = ConnectionEventMap>(
|
||
|
|
name: string,
|
||
|
|
options?: ConnectionChannelOptions
|
||
|
|
): ConnectionChannel<TEvents>;
|
||
|
|
|
||
|
|
channels(): readonly ConnectionChannel[];
|
||
|
|
|
||
|
|
hasChannel(name: string): boolean;
|
||
|
|
|
||
|
|
leaveChannel(name: string): Promise<void>;
|
||
|
|
|
||
|
|
onState(listener: (change: ConnectionStateChange) => void): () => void;
|
||
|
|
|
||
|
|
onAny(listener: (frame: ConnectionFrame, meta: ConnectionMessageMeta) => void): () => void;
|
||
|
|
|
||
|
|
dispose(): void;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
`ActiveConnection` extends `EngineConnection` and exposes reactive properties:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export interface ActiveConnection<
|
||
|
|
TChannels extends ConnectionChannelMap = ConnectionChannelMap
|
||
|
|
> extends EngineConnection<TChannels> {
|
||
|
|
readonly state: ConnectionState;
|
||
|
|
readonly connected: boolean;
|
||
|
|
readonly generation: number;
|
||
|
|
readonly channelNames: readonly string[];
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 8. State model
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export type ConnectionState =
|
||
|
|
| 'idle'
|
||
|
|
| 'connecting'
|
||
|
|
| 'open'
|
||
|
|
| 'reconnecting'
|
||
|
|
| 'closing'
|
||
|
|
| 'closed'
|
||
|
|
| 'failed';
|
||
|
|
```
|
||
|
|
|
||
|
|
State transitions:
|
||
|
|
|
||
|
|
```txt
|
||
|
|
idle -> connecting -> open
|
||
|
|
idle -> connecting -> failed
|
||
|
|
|
||
|
|
open -> reconnecting -> open
|
||
|
|
open -> closing -> closed
|
||
|
|
open -> closed
|
||
|
|
|
||
|
|
reconnecting -> open
|
||
|
|
reconnecting -> failed
|
||
|
|
reconnecting -> closed
|
||
|
|
|
||
|
|
failed -> connecting
|
||
|
|
closed -> connecting
|
||
|
|
```
|
||
|
|
|
||
|
|
Rules:
|
||
|
|
|
||
|
|
- `connect()` when already open should return `{ ok: true, state: 'open', reused: true }`.
|
||
|
|
- `connect()` while connecting should share the existing promise.
|
||
|
|
- `reconnect()` increments generation.
|
||
|
|
- `disconnect()` must cancel reconnect timers, heartbeat timers and pending acks.
|
||
|
|
- `dispose()` must call `disconnect()` and remove all listeners.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 9. Transport adapter
|
||
|
|
|
||
|
|
The connection engine must not know whether the transport is WebSocket, SSE or polling.
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export type ConnectionTransportState =
|
||
|
|
| 'idle'
|
||
|
|
| 'opening'
|
||
|
|
| 'open'
|
||
|
|
| 'closing'
|
||
|
|
| 'closed'
|
||
|
|
| 'failed';
|
||
|
|
|
||
|
|
export interface ConnectionCloseEvent {
|
||
|
|
readonly code?: number;
|
||
|
|
readonly reason?: string;
|
||
|
|
readonly clean?: boolean;
|
||
|
|
}
|
||
|
|
|
||
|
|
export 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;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
### 9.1 WebSocket transport
|
||
|
|
|
||
|
|
```ts
|
||
|
|
createWebSocketTransport({
|
||
|
|
url: () => '/realtime',
|
||
|
|
protocols: ['app-v1']
|
||
|
|
});
|
||
|
|
```
|
||
|
|
|
||
|
|
Must support:
|
||
|
|
|
||
|
|
- `open`
|
||
|
|
- `send`
|
||
|
|
- `close`
|
||
|
|
- `bufferedAmount`
|
||
|
|
- `onOpen`
|
||
|
|
- `onMessage`
|
||
|
|
- `onClose`
|
||
|
|
- `onError`
|
||
|
|
|
||
|
|
### 9.2 SSE transport
|
||
|
|
|
||
|
|
SSE is receive-first. Sending is optional.
|
||
|
|
|
||
|
|
```ts
|
||
|
|
createSseTransport({
|
||
|
|
url: () => '/events',
|
||
|
|
|
||
|
|
send: async (frame) => {
|
||
|
|
await App.http.post('/events/send', { body: frame });
|
||
|
|
}
|
||
|
|
});
|
||
|
|
```
|
||
|
|
|
||
|
|
If no `send` function is provided:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
canSend === false;
|
||
|
|
```
|
||
|
|
|
||
|
|
Calling `send()` should return:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
{ ok: false, reason: 'send_not_supported' }
|
||
|
|
```
|
||
|
|
|
||
|
|
at the connection layer.
|
||
|
|
|
||
|
|
### 9.3 Polling transport
|
||
|
|
|
||
|
|
```ts
|
||
|
|
createPollingTransport({
|
||
|
|
poll: async (cursor) => {
|
||
|
|
return App.http.get('/events', { query: { cursor } });
|
||
|
|
},
|
||
|
|
|
||
|
|
send: async (frame) => {
|
||
|
|
await App.http.post('/events', { body: frame });
|
||
|
|
},
|
||
|
|
|
||
|
|
intervalMs: 5_000
|
||
|
|
});
|
||
|
|
```
|
||
|
|
|
||
|
|
### 9.4 Mock transport
|
||
|
|
|
||
|
|
Required for tests.
|
||
|
|
|
||
|
|
```ts
|
||
|
|
const transport = createMockTransport();
|
||
|
|
|
||
|
|
transport.emitMessage(raw);
|
||
|
|
transport.emitClose({ code: 1006, clean: false });
|
||
|
|
transport.emitError(new Error('boom'));
|
||
|
|
```
|
||
|
|
|
||
|
|
The mock transport must expose deterministic test helpers but should still implement the same `ConnectionTransport` interface.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 10. Serializer adapter
|
||
|
|
|
||
|
|
Do not hardcode JSON deeply into the engine.
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export interface ConnectionSerializer<TFrame = ConnectionFrame> {
|
||
|
|
encode(frame: TFrame): string | ArrayBuffer;
|
||
|
|
decode(raw: string | ArrayBuffer): TFrame;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Default:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
jsonConnectionSerializer();
|
||
|
|
```
|
||
|
|
|
||
|
|
`jsonConnectionSerializer` must:
|
||
|
|
|
||
|
|
- encode frames with `JSON.stringify`
|
||
|
|
- decode string frames with `JSON.parse`
|
||
|
|
- reject binary frames unless explicitly configured
|
||
|
|
- return tagged errors or throw only internally and have the engine convert errors to `ConnectionError`
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 11. Frame model
|
||
|
|
|
||
|
|
Base frame:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export 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;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Meanings:
|
||
|
|
|
||
|
|
- `id`: unique message id generated by the client when needed.
|
||
|
|
- `topic`: optional channel name.
|
||
|
|
- `type`: event type.
|
||
|
|
- `payload`: event payload.
|
||
|
|
- `ts`: epoch ms.
|
||
|
|
- `ack`: when true, sender expects reply.
|
||
|
|
- `replyTo`: id of original request when this is an ack/reply.
|
||
|
|
- `error`: reply error payload.
|
||
|
|
|
||
|
|
Global events have no `topic`.
|
||
|
|
|
||
|
|
Channel events include `topic`.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 12. Event maps and typing
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export type ConnectionEventMap = Record<string, unknown>;
|
||
|
|
|
||
|
|
export type ConnectionChannelMap = Record<string, ConnectionEventMap>;
|
||
|
|
|
||
|
|
export type ConnectionMap = Record<string, ConnectionChannelMap>;
|
||
|
|
```
|
||
|
|
|
||
|
|
Example:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
type AppConnections = {
|
||
|
|
main: {
|
||
|
|
orders: {
|
||
|
|
'order.created': OrderCreated;
|
||
|
|
'order.updated': OrderUpdated;
|
||
|
|
};
|
||
|
|
|
||
|
|
cart: {
|
||
|
|
'cart.changed': CartChanged;
|
||
|
|
};
|
||
|
|
};
|
||
|
|
|
||
|
|
chat: {
|
||
|
|
rooms: {
|
||
|
|
'message.created': MessageCreated;
|
||
|
|
'typing.started': TypingStarted;
|
||
|
|
};
|
||
|
|
};
|
||
|
|
|
||
|
|
market: {
|
||
|
|
prices: {
|
||
|
|
'price.tick': PriceTick;
|
||
|
|
};
|
||
|
|
};
|
||
|
|
};
|
||
|
|
```
|
||
|
|
|
||
|
|
Usage:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
const Connections = App.createActiveConnections<AppConnections>();
|
||
|
|
|
||
|
|
const Main = Connections.createConnection('main', { ... });
|
||
|
|
const Orders = Main.channel('orders');
|
||
|
|
|
||
|
|
Orders.on('order.updated', (payload) => {
|
||
|
|
payload.orderId;
|
||
|
|
});
|
||
|
|
```
|
||
|
|
|
||
|
|
The generic typing should help at call-sites, but runtime validation is optional and schema-driven.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 13. Channels
|
||
|
|
|
||
|
|
A channel is a logical topic within one connection.
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export type ConnectionChannelState = 'idle' | 'joining' | 'joined' | 'leaving' | 'left' | 'failed';
|
||
|
|
|
||
|
|
export interface ConnectionChannel<TEvents extends ConnectionEventMap = ConnectionEventMap> {
|
||
|
|
readonly name: string;
|
||
|
|
readonly state: ConnectionChannelState;
|
||
|
|
|
||
|
|
join(params?: unknown): Promise<ConnectionChannelJoinResult>;
|
||
|
|
|
||
|
|
leave(): Promise<void>;
|
||
|
|
|
||
|
|
send<K extends keyof TEvents & string>(
|
||
|
|
type: K,
|
||
|
|
payload: TEvents[K],
|
||
|
|
options?: ConnectionSendOptions
|
||
|
|
): Promise<ConnectionSendResult>;
|
||
|
|
|
||
|
|
request<K extends keyof TEvents & string, TResult = unknown>(
|
||
|
|
type: K,
|
||
|
|
payload: TEvents[K],
|
||
|
|
options?: ConnectionRequestOptions
|
||
|
|
): Promise<ConnectionAckResult<TResult>>;
|
||
|
|
|
||
|
|
on<K extends keyof TEvents & string>(
|
||
|
|
type: K,
|
||
|
|
listener: (payload: TEvents[K], meta: ConnectionMessageMeta) => void
|
||
|
|
): () => void;
|
||
|
|
|
||
|
|
onAny(listener: (frame: ConnectionFrame, meta: ConnectionMessageMeta) => void): () => void;
|
||
|
|
|
||
|
|
dispose(): void;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Rules:
|
||
|
|
|
||
|
|
- `connection.channel(name)` should cache and return the same instance for the same name.
|
||
|
|
- Channels do not create physical connections.
|
||
|
|
- Channels route frames by `frame.topic === channel.name`.
|
||
|
|
- `leave()` sends a leave frame if connected, then marks local state as left.
|
||
|
|
- Rejoin after reconnect is configurable.
|
||
|
|
|
||
|
|
Channel options:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export interface ConnectionChannelOptions {
|
||
|
|
readonly autoJoin?: boolean;
|
||
|
|
readonly rejoinOnReconnect?: boolean;
|
||
|
|
readonly params?: () => unknown | Promise<unknown>;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 14. Presence
|
||
|
|
|
||
|
|
Presence is optional and channel-scoped.
|
||
|
|
|
||
|
|
Do not implement a complex distributed presence model in v1. Provide a minimal abstraction that can be driven by frames.
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export interface ConnectionPresence<TPresence = unknown> {
|
||
|
|
readonly members: readonly ConnectionPresenceMember<TPresence>[];
|
||
|
|
|
||
|
|
onJoin(listener: (member: ConnectionPresenceMember<TPresence>) => void): () => void;
|
||
|
|
|
||
|
|
onLeave(listener: (member: ConnectionPresenceMember<TPresence>) => void): () => void;
|
||
|
|
|
||
|
|
onSync(listener: (members: readonly ConnectionPresenceMember<TPresence>[]) => void): () => void;
|
||
|
|
}
|
||
|
|
|
||
|
|
export interface ConnectionPresenceMember<TPresence = unknown> {
|
||
|
|
readonly id: string;
|
||
|
|
readonly payload: TPresence;
|
||
|
|
readonly joinedAt?: number;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Presence may be implemented later if v1 scope is too large. Do not block the base connection engine on presence.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 15. Send/request/ack semantics
|
||
|
|
|
||
|
|
### 15.1 Send
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export type ConnectionSendResult =
|
||
|
|
| { readonly ok: true; readonly id?: string }
|
||
|
|
| {
|
||
|
|
readonly ok: false;
|
||
|
|
readonly reason:
|
||
|
|
| 'closed'
|
||
|
|
| 'send_not_supported'
|
||
|
|
| 'buffer_full'
|
||
|
|
| 'invalid_frame'
|
||
|
|
| 'serialize_failed'
|
||
|
|
| 'transport_error';
|
||
|
|
readonly error?: unknown;
|
||
|
|
};
|
||
|
|
```
|
||
|
|
|
||
|
|
`send()` should not throw for normal runtime failures.
|
||
|
|
|
||
|
|
### 15.2 Request/ack
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export type ConnectionAckResult<T = unknown> =
|
||
|
|
| { readonly ok: true; readonly payload: T }
|
||
|
|
| {
|
||
|
|
readonly ok: false;
|
||
|
|
readonly reason: 'timeout' | 'closed' | 'rejected' | 'transport_error' | 'invalid_reply';
|
||
|
|
readonly error?: unknown;
|
||
|
|
};
|
||
|
|
```
|
||
|
|
|
||
|
|
Request behavior:
|
||
|
|
|
||
|
|
- Generate id.
|
||
|
|
- Send frame with `ack: true`.
|
||
|
|
- Store pending resolver by id.
|
||
|
|
- When frame with `replyTo` arrives, resolve matching promise.
|
||
|
|
- Timeout rejects as tagged result.
|
||
|
|
- Disconnect resolves all pending acks with `{ ok: false, reason: 'closed' }`.
|
||
|
|
|
||
|
|
Options:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export interface ConnectionRequestOptions extends ConnectionSendOptions {
|
||
|
|
readonly timeoutMs?: number;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 16. Reconnect
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export interface ConnectionReconnectOptions {
|
||
|
|
readonly enabled?: boolean;
|
||
|
|
readonly minDelayMs?: number;
|
||
|
|
readonly maxDelayMs?: number;
|
||
|
|
readonly factor?: number;
|
||
|
|
readonly jitterMs?: number;
|
||
|
|
readonly maxAttempts?: number;
|
||
|
|
readonly reconnectOnVisible?: boolean;
|
||
|
|
readonly reconnectOnOnline?: boolean;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Defaults:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
{
|
||
|
|
enabled: true,
|
||
|
|
minDelayMs: 500,
|
||
|
|
maxDelayMs: 15_000,
|
||
|
|
factor: 1.8,
|
||
|
|
jitterMs: 500,
|
||
|
|
maxAttempts: undefined,
|
||
|
|
reconnectOnVisible: true,
|
||
|
|
reconnectOnOnline: true
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Rules:
|
||
|
|
|
||
|
|
- Use exponential backoff with jitter.
|
||
|
|
- Reset attempt count after successful open.
|
||
|
|
- Do not reconnect after intentional `disconnect()` unless `reconnect()` is called.
|
||
|
|
- Reconnect increments generation.
|
||
|
|
- After reconnect:
|
||
|
|
1. reauthenticate if configured
|
||
|
|
2. rejoin channels that had `rejoinOnReconnect !== false`
|
||
|
|
3. resume message flow
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 17. Heartbeat
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export interface ConnectionHeartbeatOptions {
|
||
|
|
readonly enabled?: boolean;
|
||
|
|
readonly intervalMs?: number;
|
||
|
|
readonly timeoutMs?: number;
|
||
|
|
readonly pingType?: string;
|
||
|
|
readonly pongType?: string;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Defaults:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
{
|
||
|
|
enabled: true,
|
||
|
|
intervalMs: 25_000,
|
||
|
|
timeoutMs: 10_000,
|
||
|
|
pingType: 'connection.ping',
|
||
|
|
pongType: 'connection.pong'
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Behavior:
|
||
|
|
|
||
|
|
- When open, periodically send ping if transport can send.
|
||
|
|
- If no pong/any message is received within timeout, close/reconnect.
|
||
|
|
- For SSE receive-only transports, heartbeat may be receive-based only.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 18. Backpressure and buffering
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export interface ConnectionBufferOptions {
|
||
|
|
readonly policy?: 'buffer' | 'drop' | 'fail';
|
||
|
|
readonly maxMessages?: number;
|
||
|
|
readonly maxBytes?: number;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Rules:
|
||
|
|
|
||
|
|
- If disconnected and policy is `buffer`, enqueue messages up to `maxMessages`.
|
||
|
|
- If disconnected and policy is `drop`, return ok false or drop according to options.
|
||
|
|
- If disconnected and policy is `fail`, return `{ ok: false, reason: 'closed' }`.
|
||
|
|
- If `transport.bufferedAmount > maxBytes`, return `{ ok: false, reason: 'buffer_full' }`.
|
||
|
|
- Buffered messages flush after reconnect, before normal sends.
|
||
|
|
- Request/ack messages should generally not be buffered unless explicitly allowed.
|
||
|
|
|
||
|
|
Default:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
{
|
||
|
|
policy: 'fail',
|
||
|
|
maxMessages: 100,
|
||
|
|
maxBytes: 1_000_000
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 19. Auth and reauthentication
|
||
|
|
|
||
|
|
Auth is connection-specific and opaque.
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export type ConnectionAuthPayload = Record<string, unknown>;
|
||
|
|
|
||
|
|
export interface ConnectionAuthOptions {
|
||
|
|
readonly getAuth?: () => ConnectionAuthPayload | null | Promise<ConnectionAuthPayload | null>;
|
||
|
|
|
||
|
|
readonly authType?: string;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Reauth behavior:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
reauthenticate(): Promise<ConnectionAuthResult>
|
||
|
|
```
|
||
|
|
|
||
|
|
Implementation may send a frame:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
{
|
||
|
|
type: 'connection.auth',
|
||
|
|
payload: await getAuth(),
|
||
|
|
ack: true
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Result:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export type ConnectionAuthResult =
|
||
|
|
| { readonly ok: true }
|
||
|
|
| {
|
||
|
|
readonly ok: false;
|
||
|
|
readonly reason: 'no_auth_provider' | 'closed' | 'rejected' | 'timeout' | 'transport_error';
|
||
|
|
readonly error?: unknown;
|
||
|
|
};
|
||
|
|
```
|
||
|
|
|
||
|
|
Do not hardcode Bearer/JWT.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 20. App integration
|
||
|
|
|
||
|
|
### 20.1 App options
|
||
|
|
|
||
|
|
Suggested App-level options:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export interface ActiveAppConnectionsOptions {
|
||
|
|
readonly defaults?: Partial<ConnectionOptions>;
|
||
|
|
|
||
|
|
readonly connections?: Record<string, Partial<ConnectionOptions>>;
|
||
|
|
|
||
|
|
readonly autoCreate?: boolean;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
But do not force all connections to be declared in App config. Support both patterns:
|
||
|
|
|
||
|
|
Declarative:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
const App = createActiveApp({
|
||
|
|
connections: {
|
||
|
|
connections: {
|
||
|
|
main: {
|
||
|
|
reconnect: { enabled: true }
|
||
|
|
},
|
||
|
|
market: {
|
||
|
|
buffer: { policy: 'drop' }
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
});
|
||
|
|
```
|
||
|
|
|
||
|
|
Typed creation:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
const Connections = App.createActiveConnections<AppConnections>();
|
||
|
|
|
||
|
|
const Main = Connections.createConnection('main', {
|
||
|
|
transport: createWebSocketTransport({
|
||
|
|
url: () => '/realtime'
|
||
|
|
})
|
||
|
|
});
|
||
|
|
```
|
||
|
|
|
||
|
|
### 20.2 Why factory on App
|
||
|
|
|
||
|
|
Use an App factory, not a generic `App.connections` baked into `ActiveApp<T>` by default.
|
||
|
|
|
||
|
|
Recommended:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
const Connections = App.createActiveConnections<AppConnections>();
|
||
|
|
```
|
||
|
|
|
||
|
|
Reasons:
|
||
|
|
|
||
|
|
- Type flows at call site.
|
||
|
|
- App itself avoids becoming generic on every connection map.
|
||
|
|
- Mirrors the already accepted factory pattern for typed runtime helpers.
|
||
|
|
- Avoids forcing connections in apps that do not need realtime.
|
||
|
|
|
||
|
|
### 20.3 Auto-wiring
|
||
|
|
|
||
|
|
Auto-wire only obvious infrastructure:
|
||
|
|
|
||
|
|
- `Logger` injected into connections.
|
||
|
|
- App defaults merged into connection options.
|
||
|
|
- Session events wired only when connection opts into session integration.
|
||
|
|
- Frontend visibility/online events used only when reconnect options enable them.
|
||
|
|
|
||
|
|
Do not auto-wire domain behavior.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 21. Options
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export interface ConnectionOptions<TChannels extends ConnectionChannelMap = ConnectionChannelMap> {
|
||
|
|
readonly transport: ConnectionTransport | (() => ConnectionTransport);
|
||
|
|
|
||
|
|
readonly serializer?: ConnectionSerializer;
|
||
|
|
|
||
|
|
readonly reconnect?: ConnectionReconnectOptions | false;
|
||
|
|
|
||
|
|
readonly heartbeat?: ConnectionHeartbeatOptions | false;
|
||
|
|
|
||
|
|
readonly buffer?: ConnectionBufferOptions;
|
||
|
|
|
||
|
|
readonly auth?:
|
||
|
|
| ConnectionAuthOptions
|
||
|
|
| (() => ConnectionAuthPayload | null | Promise<ConnectionAuthPayload | null>);
|
||
|
|
|
||
|
|
readonly session?:
|
||
|
|
| false
|
||
|
|
| {
|
||
|
|
readonly enabled?: boolean;
|
||
|
|
readonly reauthOnRefresh?: boolean;
|
||
|
|
readonly disconnectOnExpire?: boolean;
|
||
|
|
};
|
||
|
|
|
||
|
|
readonly channels?: {
|
||
|
|
readonly [K in keyof TChannels & string]?: ConnectionChannelOptions;
|
||
|
|
};
|
||
|
|
|
||
|
|
readonly loggerScope?: string;
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Note:
|
||
|
|
|
||
|
|
- If `transport` is a function, call it on each fresh physical connect/reconnect if the previous transport cannot be reused.
|
||
|
|
- If a transport instance can reopen safely, reuse is allowed, but keep behavior deterministic.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 22. Errors
|
||
|
|
|
||
|
|
Use custom errors for programmer errors, not normal runtime failures.
|
||
|
|
|
||
|
|
Examples:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
ConnectionAlreadyExistsError;
|
||
|
|
ConnectionNotFoundError;
|
||
|
|
ConnectionInvalidNameError;
|
||
|
|
ConnectionInvalidFrameError;
|
||
|
|
ConnectionDisposedError;
|
||
|
|
ConnectionChannelAlreadyExistsError;
|
||
|
|
ConnectionChannelNotFoundError;
|
||
|
|
```
|
||
|
|
|
||
|
|
Runtime network failures should generally return tagged results or emit state/error events, not throw unexpectedly.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 23. Public constants
|
||
|
|
|
||
|
|
```ts
|
||
|
|
export const CONNECTION_STATE_IDLE = 'idle';
|
||
|
|
export const CONNECTION_STATE_CONNECTING = 'connecting';
|
||
|
|
export const CONNECTION_STATE_OPEN = 'open';
|
||
|
|
export const CONNECTION_STATE_RECONNECTING = 'reconnecting';
|
||
|
|
export const CONNECTION_STATE_CLOSING = 'closing';
|
||
|
|
export const CONNECTION_STATE_CLOSED = 'closed';
|
||
|
|
export const CONNECTION_STATE_FAILED = 'failed';
|
||
|
|
|
||
|
|
export const DEFAULT_RECONNECT_MIN_DELAY_MS = 500;
|
||
|
|
export const DEFAULT_RECONNECT_MAX_DELAY_MS = 15_000;
|
||
|
|
export const DEFAULT_RECONNECT_FACTOR = 1.8;
|
||
|
|
export const DEFAULT_RECONNECT_JITTER_MS = 500;
|
||
|
|
|
||
|
|
export const DEFAULT_HEARTBEAT_INTERVAL_MS = 25_000;
|
||
|
|
export const DEFAULT_HEARTBEAT_TIMEOUT_MS = 10_000;
|
||
|
|
|
||
|
|
export const DEFAULT_ACK_TIMEOUT_MS = 10_000;
|
||
|
|
|
||
|
|
export const TRANSPORT_KIND_WEBSOCKET = 'websocket';
|
||
|
|
export const TRANSPORT_KIND_MOCK = 'mock';
|
||
|
|
|
||
|
|
export const FRAME_KEY_ID = 'id';
|
||
|
|
export const FRAME_KEY_TYPE = 'type';
|
||
|
|
export const FRAME_KEY_TOPIC = 'topic';
|
||
|
|
export const FRAME_KEY_PAYLOAD = 'payload';
|
||
|
|
export const FRAME_KEY_REPLY_TO = 'replyTo';
|
||
|
|
|
||
|
|
export const LOGGER_CATEGORY = 'conn';
|
||
|
|
export const LOG_MSG_LISTENER_THREW_PREFIX = 'Listener threw on ';
|
||
|
|
|
||
|
|
export const ERROR_NAME_DISPOSED = 'ConnectionDisposedError';
|
||
|
|
export const ERROR_NAME_ALREADY_EXISTS = 'ConnectionAlreadyExistsError';
|
||
|
|
export const ERROR_NAME_NOT_FOUND = 'ConnectionNotFoundError';
|
||
|
|
export const ERROR_NAME_INVALID_NAME = 'ConnectionInvalidNameError';
|
||
|
|
export const ERROR_NAME_INVALID_FRAME = 'ConnectionInvalidFrameError';
|
||
|
|
```
|
||
|
|
|
||
|
|
Prefer exported constants over magic strings/numbers. Engine bodies should not
|
||
|
|
inline connection states, event names, transport kinds, close reasons, frame
|
||
|
|
field names, logger categories or error names.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 24. Implementation algorithm
|
||
|
|
|
||
|
|
### 24.1 createEngineConnections
|
||
|
|
|
||
|
|
- Keep `Map<string, EngineConnection>`.
|
||
|
|
- Validate connection names.
|
||
|
|
- Merge hub defaults with per-connection options.
|
||
|
|
- Create `EngineConnection`.
|
||
|
|
- Register it.
|
||
|
|
- Return it.
|
||
|
|
- Dispose closes all.
|
||
|
|
|
||
|
|
### 24.2 createEngineConnection
|
||
|
|
|
||
|
|
Internal state:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
let state: ConnectionState = 'idle';
|
||
|
|
let generation = 0;
|
||
|
|
let disposed = false;
|
||
|
|
let connectPromise: Promise<ConnectionConnectResult> | null = null;
|
||
|
|
|
||
|
|
const channels = new Map<string, ConnectionChannel>();
|
||
|
|
const listeners = new Set<...>();
|
||
|
|
const pendingAcks = new Map<string, PendingAck>();
|
||
|
|
const outboundBuffer: ConnectionFrame[] = [];
|
||
|
|
```
|
||
|
|
|
||
|
|
On `connect()`:
|
||
|
|
|
||
|
|
1. If disposed, throw `ConnectionDisposedError`.
|
||
|
|
2. If open, return reused success.
|
||
|
|
3. If connecting, return existing promise.
|
||
|
|
4. Create/open transport.
|
||
|
|
5. Attach transport listeners.
|
||
|
|
6. Transition to `connecting`.
|
||
|
|
7. Await transport open.
|
||
|
|
8. Transition to `open`.
|
||
|
|
9. Authenticate if configured.
|
||
|
|
10. Start heartbeat.
|
||
|
|
11. Flush buffer.
|
||
|
|
12. Auto-join configured channels.
|
||
|
|
13. Return success.
|
||
|
|
|
||
|
|
On transport message:
|
||
|
|
|
||
|
|
1. Decode with serializer.
|
||
|
|
2. Validate frame shape.
|
||
|
|
3. If `replyTo`, resolve pending ack.
|
||
|
|
4. If heartbeat pong, mark heartbeat alive.
|
||
|
|
5. If `topic`, route to channel.
|
||
|
|
6. Emit global listeners.
|
||
|
|
|
||
|
|
On transport close:
|
||
|
|
|
||
|
|
1. Stop heartbeat.
|
||
|
|
2. Resolve pending acks as closed.
|
||
|
|
3. If intentional, transition closed.
|
||
|
|
4. If not intentional and reconnect enabled, transition reconnecting and schedule reconnect.
|
||
|
|
5. Else transition failed/closed.
|
||
|
|
|
||
|
|
On `disconnect()`:
|
||
|
|
|
||
|
|
1. Mark intentional close.
|
||
|
|
2. Stop reconnect timer.
|
||
|
|
3. Stop heartbeat.
|
||
|
|
4. Close transport.
|
||
|
|
5. Resolve pending acks.
|
||
|
|
6. Transition closed.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 25. Validation
|
||
|
|
|
||
|
|
Runtime frame validation should be minimal by default:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
function isConnectionFrame(value: unknown): value is ConnectionFrame {
|
||
|
|
return (
|
||
|
|
typeof value === 'object' &&
|
||
|
|
value !== null &&
|
||
|
|
typeof (value as { type?: unknown }).type === 'string' &&
|
||
|
|
'payload' in value
|
||
|
|
);
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
Optional Standard Schema support may be added for event payloads, but do not make it required for v1.
|
||
|
|
|
||
|
|
If added later:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
schemas: {
|
||
|
|
orders: {
|
||
|
|
'order.updated': OrderUpdatedSchema
|
||
|
|
}
|
||
|
|
}
|
||
|
|
```
|
||
|
|
|
||
|
|
No vendor-specific validation library.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 26. Testing checklist
|
||
|
|
|
||
|
|
### Engine hub
|
||
|
|
|
||
|
|
- creates connection
|
||
|
|
- rejects duplicate connection
|
||
|
|
- retrieves connection
|
||
|
|
- throws on missing connection
|
||
|
|
- closes one connection
|
||
|
|
- closes all connections
|
||
|
|
- disposes cleanly
|
||
|
|
|
||
|
|
### Connection lifecycle
|
||
|
|
|
||
|
|
- connects successfully
|
||
|
|
- connect dedups in-flight promise
|
||
|
|
- disconnect closes transport
|
||
|
|
- reconnect increments generation
|
||
|
|
- no reconnect after intentional disconnect
|
||
|
|
- reconnect after unexpected close
|
||
|
|
- failed state after max attempts
|
||
|
|
|
||
|
|
### Transport
|
||
|
|
|
||
|
|
- WebSocket transport maps events correctly
|
||
|
|
- mock transport works deterministically
|
||
|
|
- SSE transport reports `canSend = false` when no send provided
|
||
|
|
- polling transport receives messages repeatedly
|
||
|
|
|
||
|
|
### Channels
|
||
|
|
|
||
|
|
- channel is cached by name
|
||
|
|
- routes messages by topic
|
||
|
|
- global messages do not hit channel listeners
|
||
|
|
- channel leave updates state
|
||
|
|
- rejoin after reconnect when configured
|
||
|
|
|
||
|
|
### Send/request
|
||
|
|
|
||
|
|
- send serializes frame
|
||
|
|
- send fails when closed and buffer policy fail
|
||
|
|
- send buffers when policy buffer
|
||
|
|
- send drops/fails on buffer full
|
||
|
|
- request resolves on replyTo
|
||
|
|
- request times out
|
||
|
|
- disconnect resolves pending requests as closed
|
||
|
|
|
||
|
|
### Heartbeat
|
||
|
|
|
||
|
|
- sends ping
|
||
|
|
- receives pong
|
||
|
|
- reconnects/closes on heartbeat timeout
|
||
|
|
- disabled heartbeat does nothing
|
||
|
|
|
||
|
|
### Session integration
|
||
|
|
|
||
|
|
- session refresh triggers reauth only for configured connections
|
||
|
|
- session revoke disconnects only configured connections
|
||
|
|
- public connections remain open
|
||
|
|
|
||
|
|
### Active wrappers
|
||
|
|
|
||
|
|
- `state` reactive
|
||
|
|
- `connected` reactive
|
||
|
|
- `generation` reactive
|
||
|
|
- channel names reactive if exposed
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 27. Implementation order for Claude Code
|
||
|
|
|
||
|
|
Implement in this order:
|
||
|
|
|
||
|
|
1. `consts.ts`
|
||
|
|
2. `types.ts`
|
||
|
|
3. `errors.ts`
|
||
|
|
4. `serializer.ts` + `serializers/json.ts`
|
||
|
|
5. `transport.ts`
|
||
|
|
6. `transports/mock.ts`
|
||
|
|
7. `engine-connection.ts`
|
||
|
|
8. `channel.ts`
|
||
|
|
9. `engine-connections.ts`
|
||
|
|
10. Tests for engine with mock transport
|
||
|
|
11. `transports/websocket.ts`
|
||
|
|
12. `reconnect.ts`
|
||
|
|
13. `heartbeat.ts`
|
||
|
|
14. `backpressure.ts`
|
||
|
|
15. `ack.ts`
|
||
|
|
16. `active-connection.svelte.ts`
|
||
|
|
17. `active-connections.svelte.ts`
|
||
|
|
18. `app-integration.ts`
|
||
|
|
19. Optional SSE/polling transports
|
||
|
|
20. Documentation examples
|
||
|
|
|
||
|
|
Do not implement presence in the first pass unless everything above is stable.
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 28. Example usage
|
||
|
|
|
||
|
|
```ts
|
||
|
|
type AppConnections = {
|
||
|
|
main: {
|
||
|
|
orders: {
|
||
|
|
'order.updated': {
|
||
|
|
orderId: string;
|
||
|
|
status: string;
|
||
|
|
};
|
||
|
|
};
|
||
|
|
|
||
|
|
cart: {
|
||
|
|
'cart.changed': {
|
||
|
|
cartId: string;
|
||
|
|
total: number;
|
||
|
|
};
|
||
|
|
};
|
||
|
|
};
|
||
|
|
|
||
|
|
market: {
|
||
|
|
prices: {
|
||
|
|
'price.tick': {
|
||
|
|
symbol: string;
|
||
|
|
price: number;
|
||
|
|
};
|
||
|
|
};
|
||
|
|
};
|
||
|
|
};
|
||
|
|
|
||
|
|
const Connections = App.createActiveConnections<AppConnections>();
|
||
|
|
|
||
|
|
const Main = Connections.createConnection('main', {
|
||
|
|
transport: createWebSocketTransport({
|
||
|
|
url: () => '/realtime'
|
||
|
|
}),
|
||
|
|
|
||
|
|
auth: () => {
|
||
|
|
const session = Sess.current;
|
||
|
|
|
||
|
|
return session?.credential ? { token: session.credential.accessToken } : null;
|
||
|
|
},
|
||
|
|
|
||
|
|
session: {
|
||
|
|
enabled: true,
|
||
|
|
reauthOnRefresh: true,
|
||
|
|
disconnectOnExpire: true
|
||
|
|
}
|
||
|
|
});
|
||
|
|
|
||
|
|
const Orders = Main.channel('orders');
|
||
|
|
|
||
|
|
Orders.on('order.updated', (payload) => {
|
||
|
|
console.log(payload.orderId, payload.status);
|
||
|
|
});
|
||
|
|
|
||
|
|
await Main.connect();
|
||
|
|
```
|
||
|
|
|
||
|
|
Second connection:
|
||
|
|
|
||
|
|
```ts
|
||
|
|
const Market = Connections.createConnection('market', {
|
||
|
|
transport: createWebSocketTransport({
|
||
|
|
url: () => 'wss://prices.example.com/stream'
|
||
|
|
}),
|
||
|
|
|
||
|
|
reconnect: {
|
||
|
|
minDelayMs: 250,
|
||
|
|
maxDelayMs: 3_000
|
||
|
|
},
|
||
|
|
|
||
|
|
buffer: {
|
||
|
|
policy: 'drop',
|
||
|
|
maxMessages: 50
|
||
|
|
},
|
||
|
|
|
||
|
|
session: false
|
||
|
|
});
|
||
|
|
|
||
|
|
const Prices = Market.channel('prices');
|
||
|
|
|
||
|
|
Prices.on('price.tick', (tick) => {
|
||
|
|
console.log(tick.symbol, tick.price);
|
||
|
|
});
|
||
|
|
|
||
|
|
await Market.connect();
|
||
|
|
```
|
||
|
|
|
||
|
|
---
|
||
|
|
|
||
|
|
## 29. Final architectural rule
|
||
|
|
|
||
|
|
`arts/conn` owns realtime connection infrastructure.
|
||
|
|
|
||
|
|
It does not own auth, session, permissions, user state, cache, or webhooks.
|
||
|
|
|
||
|
|
The final layering should remain:
|
||
|
|
|
||
|
|
```txt
|
||
|
|
arts/sess
|
||
|
|
session lifecycle
|
||
|
|
|
||
|
|
arts/http
|
||
|
|
request/response
|
||
|
|
|
||
|
|
arts/conn
|
||
|
|
realtime connections
|
||
|
|
|
||
|
|
arts/webh
|
||
|
|
signed webhook receive/send infrastructure
|
||
|
|
|
||
|
|
arts/cache
|
||
|
|
future derived state/query cache
|
||
|
|
|
||
|
|
ActiveApp
|
||
|
|
composition root that wires them when configured
|
||
|
|
```
|