43 KiB
arts/connection — Design & Implementation Document
Historical note 2026-05-14. This document predates the directory rename sweep. The canonical artifact name is now
connection;$connectionpoints tosrc/arts/connection. Mentions ofconn,sess,timr,logr,cachandaappbelow 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, notsock/sockets, because the public model must support WebSocket, SSE, polling, WebTransport, mocks, and future transports through adapters.Naming correction applied.
EngineConnections/ActiveConnectionsare the root registry names. An individual runtime unit isConnection. Any older references below toEngineConnection/ActiveConnectiondescribe 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:
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:
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:
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:
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:
EngineConnectiononly talks toConnectionTransport.- 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:
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,errorevents. - Supports dynamic URL:
url: string | (() => string);
Required usage:
const Main = Connections.createConnection('main', {
transport: createWebSocketTransport({
url: () => '/realtime'
})
});
1.1.4 Mock transport
v1 MUST include:
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:
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:
export type ConnectionState =
| 'idle'
| 'connecting'
| 'open'
| 'reconnecting'
| 'closing'
| 'closed'
| 'failed';
Required EngineConnection/ActiveConnection properties:
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:
onState(listener: (change: ConnectionStateChange) => void): () => void;
Required change shape:
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:
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:
Connectionowns the real state.Connectionsobserves 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:
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:
id?: string;
type: string;
payload: unknown;
ack?: boolean;
replyTo?: string;
error?: ConnectionFrameError;
Required behavior:
request()generates a message id.request()sends a frame withack: true.- Incoming frames with
replyToresolve 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:
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 injectedTimerScheduler/App.timers, not rawsetTimeoutinside the connection engine. A native fallback is allowed only in the standalone engine path.
1.1.9 Heartbeat
v1 MUST support heartbeat.
Required options:
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:
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
topicare routed to the matching channel. - Frames without
topicgo to connection-level listeners. - Channels expose
join,leave,send,request,on,onAny. - Channels may auto-join and rejoin after reconnect.
Required channel methods:
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:
session:
| false
| {
readonly enabled?: boolean;
readonly reauthOnRefresh?: boolean;
readonly disconnectOnExpire?: boolean;
};
Required behavior:
- Connections with
session: falsenever 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:
auth?: {
getAuth?: () =>
| Record<string, unknown>
| null
| Promise<Record<string, unknown> | null>;
}
or shorthand:
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:
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:
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:
const Connections = App.createActiveConnections<AppConnections>();
or, once integrated:
App.connections.connection('main');
Do not put this under:
Sess.connections;
Sess.socket;
Sess.realtime;
Session is only a source of auth/lifecycle events.
3.2 With sess
arts/sess owns session lifecycle:
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:
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:
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:
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:
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/connowns only connection-specific constants: states, events, transport kinds, frame keys, close reasons, logger messages and error names.
4. Naming
Artifact:
src/arts/connection/
Public names:
EngineConnections;
ActiveConnections;
EngineConnection;
ActiveConnection;
ConnectionChannel;
ConnectionPresence;
ConnectionTransport;
ConnectionSerializer;
Factory names:
createEngineConnections;
createActiveConnections;
createWebSocketTransport;
createSseTransport;
createPollingTransport;
createMockTransport;
jsonConnectionSerializer;
Avoid names such as:
SocketService;
RealtimeManager;
ConnectionManager;
5. Directory structure
Target structure:
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:
.tsfiles must be runes-free..svelte.tsfiles may use$state.index.tsmust 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.
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 throwsConnectionNotFoundError.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:
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:
Connections.connection(name).connect();
Connections.connection(name).disconnect();
Connections.connection(name).reconnect();
Both levels are intentionally supported.
6.2 Active hub
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 createEngineConnections.
Correct:
const Connections = createActiveConnections();
const Main = Connections.createConnection('main', options);
Avoid:
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
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:
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
export type ConnectionState =
| 'idle'
| 'connecting'
| 'open'
| 'reconnecting'
| 'closing'
| 'closed'
| 'failed';
State transitions:
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 calldisconnect()and remove all listeners.
9. Transport adapter
The connection engine must not know whether the transport is WebSocket, SSE or polling.
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
createWebSocketTransport({
url: () => '/realtime',
protocols: ['app-v1']
});
Must support:
opensendclosebufferedAmountonOpenonMessageonCloseonError
9.2 SSE transport
SSE is receive-first. Sending is optional.
createSseTransport({
url: () => '/events',
send: async (frame) => {
await App.http.post('/events/send', { body: frame });
}
});
If no send function is provided:
canSend === false;
Calling send() should return:
{ ok: false, reason: 'send_not_supported' }
at the connection layer.
9.3 Polling transport
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.
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.
export interface ConnectionSerializer<TFrame = ConnectionFrame> {
encode(frame: TFrame): string | ArrayBuffer;
decode(raw: string | ArrayBuffer): TFrame;
}
Default:
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:
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
export type ConnectionEventMap = Record<string, unknown>;
export type ConnectionChannelMap = Record<string, ConnectionEventMap>;
export type ConnectionMap = Record<string, ConnectionChannelMap>;
Example:
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:
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.
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:
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.
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
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
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
replyToarrives, resolve matching promise. - Timeout rejects as tagged result.
- Disconnect resolves all pending acks with
{ ok: false, reason: 'closed' }.
Options:
export interface ConnectionRequestOptions extends ConnectionSendOptions {
readonly timeoutMs?: number;
}
16. Reconnect
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:
{
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()unlessreconnect()is called. - Reconnect increments generation.
- After reconnect:
- reauthenticate if configured
- rejoin channels that had
rejoinOnReconnect !== false - resume message flow
17. Heartbeat
export interface ConnectionHeartbeatOptions {
readonly enabled?: boolean;
readonly intervalMs?: number;
readonly timeoutMs?: number;
readonly pingType?: string;
readonly pongType?: string;
}
Defaults:
{
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
export interface ConnectionBufferOptions {
readonly policy?: 'buffer' | 'drop' | 'fail';
readonly maxMessages?: number;
readonly maxBytes?: number;
}
Rules:
- If disconnected and policy is
buffer, enqueue messages up tomaxMessages. - 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:
{
policy: 'fail',
maxMessages: 100,
maxBytes: 1_000_000
}
19. Auth and reauthentication
Auth is connection-specific and opaque.
export type ConnectionAuthPayload = Record<string, unknown>;
export interface ConnectionAuthOptions {
readonly getAuth?: () => ConnectionAuthPayload | null | Promise<ConnectionAuthPayload | null>;
readonly authType?: string;
}
Reauth behavior:
reauthenticate(): Promise<ConnectionAuthResult>
Implementation may send a frame:
{
type: 'connection.auth',
payload: await getAuth(),
ack: true
}
Result:
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:
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:
const App = createActiveApp({
connections: {
connections: {
main: {
reconnect: { enabled: true }
},
market: {
buffer: { policy: 'drop' }
}
}
}
});
Typed creation:
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:
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:
Loggerinjected 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
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
transportis 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:
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
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:
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():
- If disposed, throw
ConnectionDisposedError. - If open, return reused success.
- If connecting, return existing promise.
- Create/open transport.
- Attach transport listeners.
- Transition to
connecting. - Await transport open.
- Transition to
open. - Authenticate if configured.
- Start heartbeat.
- Flush buffer.
- Auto-join configured channels.
- Return success.
On transport message:
- Decode with serializer.
- Validate frame shape.
- If
replyTo, resolve pending ack. - If heartbeat pong, mark heartbeat alive.
- If
topic, route to channel. - Emit global listeners.
On transport close:
- Stop heartbeat.
- Resolve pending acks as closed.
- If intentional, transition closed.
- If not intentional and reconnect enabled, transition reconnecting and schedule reconnect.
- Else transition failed/closed.
On disconnect():
- Mark intentional close.
- Stop reconnect timer.
- Stop heartbeat.
- Close transport.
- Resolve pending acks.
- Transition closed.
25. Validation
Runtime frame validation should be minimal by default:
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:
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 = falsewhen 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
statereactiveconnectedreactivegenerationreactive- channel names reactive if exposed
27. Implementation order for Claude Code
Implement in this order:
consts.tstypes.tserrors.tsserializer.ts+serializers/json.tstransport.tstransports/mock.tsengine-connection.tschannel.tsengine-connections.ts- Tests for engine with mock transport
transports/websocket.tsreconnect.tsheartbeat.tsbackpressure.tsack.tsactive-connection.svelte.tsactive-connections.svelte.tsapp-integration.ts- Optional SSE/polling transports
- Documentation examples
Do not implement presence in the first pass unless everything above is stable.
28. Example usage
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:
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:
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