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.
svelte-kit-vice/src/arts/connection/DESIGN_CONN.md

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; $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:

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:

  • 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:

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:
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:

  • 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:

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 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:

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:

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 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:

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: 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:

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/conn owns 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:

  • .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.

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:

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 create EngineConnections.

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 call disconnect() 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:

  • open
  • send
  • close
  • bufferedAmount
  • onOpen
  • onMessage
  • onClose
  • onError

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 replyTo arrives, 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() 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

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 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:

{
    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:

  • 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

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:

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():

  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:

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 = 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

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

Powered by TurnKey Linux.