Skip to content

API reference > @kontsedal/olas-realtime

olas-realtime package ​

Functions ​

Function

Description

createConnectionState(ctx)

Reactive connection-state signal — 'connected' | 'reconnecting' | 'offline' | 'unknown'.

Backed by RealtimeService.onConnectionChange?(...). If the consumer's transport doesn't implement that method, the returned signal stays at'unknown' for the lifetime of the controller — the hook has no way to observe connection state, so it reports that honestly rather than lying with 'connected' (T6.7). With a reporter, it starts optimistically at'connected' until the first change corrects it.

Every createConnectionState and onReconnect on one RealtimeService shares one onConnectionChange listener. The first live one subscribes, and the last one to dispose (or suspend) unsubscribes. One that starts while the listener is live begins at the transport's latest report, since the transport's own "current state" call came once, on that subscribe. With no report since the listener opened, it begins at 'connected'. A resume counts as a start: the state from before the suspend is dropped, because nothing listened while the controller was suspended.

Useful for "stale-during-disconnect" UIs and as a refetch trigger when the connection comes back up:

ts
const conn = createConnectionState(ctx)
const orders = bindQuery(ctx, ordersQuery)
ctx.effect(() => {
  if (conn.value === 'connected') {
    orders.invalidateAll()
  }
})

createLiveStream(ctx, channel, options)

Subscribe to channel, buffer events into a ReadSignal<readonly TEvent[]> with capacity oldest-drop semantics and flushMs coalescing. The subscription lives inside ctx.effect so pause/resume re-runs it (we readisPaused.value as a tracked dep).

Naming: the create* prefix is the convention for ctx-taking composables (createPersisted, createRealtimePatcher). The define* prefix is reserved for module-scope factories (defineQuery, defineController).

Buffer semantics (SPEC §16.5): - capacity caps memory; oldest entries drop when exceeded. - flushMs coalesces N events into one signal write — prevents thrashing under 1000-events/sec bursts. flushMs <= 0 flushes synchronously. - pause() tears down the subscription; already-buffered events survive, but events arriving DURING the pause are LOST (not received). Recover a gap with onReconnect(...) + query invalidate, not the buffer. - clear() resets the buffer (and any unflushed pending events) without touching the subscription. - channel is a name or a signal of one. A new name moves the subscription to the new channel and empties the buffer, as clear() does: a tail holds one channel's events, and nothing in an event says which channel it came from. A change during a pause empties it too, and resume() subscribes to the new name.

createRealtimePatcher(ctx, channel, handlers)

Subscribe to channel for the lifetime of the surrounding controller and dispatch each event to the matching handler by event.type. Wrapper around the recurring SPEC §16.5 "realtime → cache patches" pattern.

channel is a name or a signal of one. With a signal, a new name unsubscribes from the old channel and subscribes to the new one, so a per-route room is computed(() => 'room:' + params.value.roomId).

Handlers run inside untracked(...) so accidental signal reads (e.g.query.setData((prev) => prev.value)) don't add deps to the enclosing effect, which would otherwise re-subscribe whenever those signals change.

A '*' wildcard handler is invoked for every event, regardless of whether a specific handler matched. Specific handlers run first; the wildcard sees the same event afterwards (in the same untracked scope).

onReconnect(ctx, fn)

Trigger fn() when the realtime connection transitions back to'connected' from a non-connected state. Typical use: invalidate queries that may have missed updates during the disconnect window.

The handler fires AFTER the transition is observed; it does NOT fire on the initial 'connected' value (no transition happened yet). Wrapped in untracked so cache writes don't accidentally hook the effect.

A resume that moves the state from 'offline' or 'reconnecting' back to'connected' counts: createConnectionState restarts at the transport's latest report, or at 'connected' when there is none. The controller's own channel subscriptions stopped while it was suspended, so the refetch is due either way.

Type Aliases ​

Type Alias

Description

ConnectionState

A realtime transport's connection state, as createConnectionState reports it. 'unknown' means the transport has no onConnectionChange, so the state cannot be observed. RealtimeService.onConnectionChange describes the other three.

LiveStream

Live-streaming buffer over a single realtime channel. Pause/resume controls the subscription (not the buffer — already-buffered events are preserved across a pause). **Events that arrive DURING a pause are lost:**pause() tears down the underlying subscription, so nothing is received (let alone buffered) until resume(). To recover a gap, pair withonReconnect(...) + a query invalidate (refetch authoritative state) rather than relying on the buffer. clear() empties the buffer without touching the subscription. A channel signal's change empties it too, since the buffered events came from the old channel. SPEC §16.5 tail-buffer pattern.

LiveStreamOptions

Live-stream options. capacity is the maximum buffer length (oldest events are dropped); flushMs coalesces bursty writes into a single signal update.flushMs <= 0 flushes synchronously per event.

PatcherHandlers

Map of event.type literal → handler. Each key is one type of the discriminated union, and its handler receives that variant only:Extract<TEvent, { type: K }>, so 'comment-added': (ev) => ev.comment needs no narrowing.

The '*' wildcard key receives every event the dispatcher saw, typed as the whole union, including those a specific handler already consumed. Use it for logging, instrumentation, or "I don't know all the types yet" diagnostics.

RealtimeDeps

Slice of ctx.deps consumed by this package.

RealtimeHandler

Per-event callback handed to RealtimeService.subscribe.

RealtimeService

Consumer-implemented realtime transport. The package ships no default — apps wire their own (WebSocket, Pusher, Supabase Realtime, etc.) and pass the implementation through ctx.deps.realtime after augmentingAmbientDeps:

ts
declare module '@kontsedal/olas-core' {
  interface AmbientDeps {
    realtime: RealtimeService
  }
}

RealtimeSubscription

A handle returned by RealtimeService.subscribe(...). Matches the shape used by most WebSocket / SSE / Pusher / Ably / Supabase clients in the wild. SPEC §16.5.

Released under the MIT License.