@latkit/port
v0.8.0
Published
Where Latkit crosses a boundary: ports, frames, typed protocols, and a model and its recordings served across them.
Readme
@latkit/port
Where latkit crosses a boundary: a two-method port over workers, webviews, sockets, and one thread;
one binary frame that carries typed arrays intact; typed request, reply, and stream protocols with
the checks their served side runs; and @latkit/model models, engines, and recordings served and
connected across a port.
Install
npm install @latkit/portA port
import { messagePort, socketPort } from '@latkit/port';
const worker = messagePort(new Worker(new URL('./worker.ts', import.meta.url), { type: 'module' }));
const server = socketPort(new WebSocket('wss://example.org/model'));A Port has post, subscribe, and an optional drain that resolves once the transport has room
for more. messagePort wraps anything with the DOM message-target shape and carries what
structured clone carries, transfer list included. bytePort wraps any channel that carries bytes
faithfully and rides each message on one binary frame, so typed arrays view the received buffer in
place even where structured clone does not survive. socketPort is bytePort over a browser
WebSocket or a node ws socket, queueing posts until it opens. loopback() is two ports wired
to each other in one realm, every message crossing as a frame on a microtask: a client and a server
on one thread, or a test's two ends, where fail(reason) delivers a transport failure. Each
constructor takes its target structurally, so a Worker, a webview API, or a socket passes as it
is; only Port is a named type.
Every message is JSON values plus typed arrays (Uint8Array through Float64Array), anywhere in
the value, on every transport. A service written against a worker runs unchanged against a socket.
messagePort does not refuse what structured clone would carry beyond that; loopback does.
Models, engines, and recordings across a port
// the worker: cases parse here, and the engine records here
import { messagePort, serveEngine, serveModel, serveRecording } from '@latkit/port';
serveEngine(messagePort(self), new GridkitEngine(server)); // records any model a peer gives it
serveRecording(messagePort(self), recording); // a recording the worker keeps
self.addEventListener('message', ({ data }) => {
// one case per channel, its packs served as they are asked for
if (data.open) serveModel(messagePort(data.open.port), new GridkitCase(data.open.bytes));
});
// the page
import { connectEngine, connectModel, connectRecording, messagePort } from '@latkit/port';
const port = messagePort(worker);
const engine = await connectEngine(port);
const { port1, port2 } = new MessageChannel();
worker.postMessage({ open: { port: port2, bytes } }, [port2]);
const model = await connectModel(messagePort(port1), { progress });
const recording = engine.record(model, input); // held by the worker, followed here
const kept = await connectRecording(port, model, 'fault-4'); // or opens the worker's ownOnly sources, changes, and the windows read cross. A model's core crosses at once and each class
shard as it is first asked for, with the case's bytes on request; a model opened from packs serves
them as they came. An engine records any model it is given: a model its own realm serves is
recorded where it lives, and any other is lent by its source, which the engine reads only as it
needs, for as long as the recording lasts; a file an input gives is lent the same way, its bytes
crossing only as the engine reads them. Each recording is held where the engine runs, its frames
in the engine's store: the far side follows its changes, its clock, ranges, state, and log, reads
its frames a window of at most 4 MiB at a time, and lets it go by closing it, as the port's close
lets every one go. The served engine checks every input, queues what it cannot take at once, and
stops when the far recording stops. The studies an engine offers cross with it:
connectEngine resolves once they are in, and the connected engine follows each change and checks
a study's form where it is. A kept recording opens with Recording.from against
the model it records, its clock at hand and its samples read in windows of at most 4 MiB. A
connected side is a Remote<T>: the model, engine, or recording, plus close.
Documents across a port
import { connectDocument, serveDocument } from '@latkit/port';
// Server or worker: format is a Document.Format; retain its document across connections.
const document = await format.open(nativeBytes);
serveDocument(serverPort, document);
// Browser: mutations are asynchronous; the view and schematic lookups are local.
const session = await connectDocument(clientPort);
session.on('change', () => render(session.view.schematic));
await session.apply({ kind: 'set', element, column: 'kv', value: 138 });
await session.undo();
const bytes = await session.bytes(); // Current native case; the host owns saving.
const snapshot = await session.model();
// Keep it until every reader or recording using it finishes.
snapshot.close();
// On a new transport, reconcile a possibly lost edit acknowledgment.
await connectDocument(newPort, { resume: session });
session.close();serveDocument(port, () => openNativeDocument()) defers loading until the first document open
request. Each service invokes its factory at most once and shares its result or failure. Closing
an unused service never invokes it; closing during loading prevents attachment without cancelling
host-owned work. Supplied documents and promises remain supported. For reconnects or multiple
clients, have the factory return the same host-owned document; see
Lazy document loading.
Opening a session, editing, and exporting native bytes do not build a model. session.model()
requests an immutable snapshot only when needed; its bytes remain frozen across later edits.
Each document has one serialized owner, retained across connections. Commands carry the
document's epoch, base revision, and client sequence; stale indexed edits throw
DocumentConflict. Refusals retain Refusal.at. Updates carry the schematic parts a change
replaced, compared by identity; layout changes retain the netlist and model. A gap refreshes the
cached view. Acknowledgments mean accepted in memory.
session.inspect(elementOrKey, signal?) returns { version, inspection }: editable values and
complete wiring read together without materializing a model. Retain the revision with a form and
submit through session.apply(version, ...operations); the owner rejects stale drafts. Inspections
copy only public fields and are bounded to 1 MiB; the service and frame format remain unchanged.
Retry state, queues, snapshots, and slow-peer event buffers are bounded. Model snapshots reuse
scoped model services (serveModel / connectModel accept an optional id) and the existing
engine reference path. See Document sessions for lifetime,
reconnect, limits, and the document wire contract.
A protocol
Both ends import one value: the name on the port, the request, reply, and event types, and the check the served side runs on every request.
import { check, protocol } from '@latkit/port';
type Request =
{ readonly op: 'greet'; readonly name: string } | { readonly op: 'count'; readonly upTo: number };
export const HELLO = protocol<Request, string, { readonly tick: number }>(
'hello',
check.requests<Request>({ greet: { name: check.string }, count: { upTo: check.index } }),
);check holds the checks: string, boolean, finite, index, bounded, bytes, oneOf,
nullable, optional, object, array, stringMap, record, and requests, whose shape map
the compiler keeps exhaustive over the request union's op and exact in every field's type. A
check returns when a value is what it claims and throws a TypeError naming what is wrong, so a
refused request is answered with that reason and never reaches the handler.
Serve and connect
// the worker
import { serve } from '@latkit/port';
const hello = serve(port, HELLO, async (request) => {
if (request.op === 'greet') return `Hello, ${request.name}.`;
return String(request.upTo);
});
hello.emit({ tick: 1 });
// the page
import { connect } from '@latkit/port';
const hello = connect(port, HELLO);
hello.on((event) => console.log(event.tick));
const greeting = await hello.call({ op: 'greet', name: 'Ada' }, { signal });
hello.close();Several protocols share one port; each serve and connect sees only its own. A handler failure
rejects that one call with the handler's message. A cancelled call aborts the handler's signal.
Either side may close; a transport failure closes every connection on the port with its reason. A
call to a protocol no peer serves settles only when the transport closes, which is why one protocol
value is imported at both ends.
Stream
A handler that returns an async iterable streams, one yield per item, and the service awaits the
port's drain between items so backpressure reaches the producer.
serve(port, FRAMES, async function* (request, signal) {
for await (const frame of frames(request, signal)) yield frame;
});
for await (const frame of connect(port, FRAMES).stream(request, { signal })) paint(frame);Leaving the loop early, or aborting signal, cancels the handler and ends the iteration quietly. A
handler failure ends it with that error. A reply whose buffers the handler relinquishes is wrapped
with transferred(value, buffers), for a reply or for a streamed item alike. Call options
(signal, progress, transfer) and a handler's shape are stated inline on call, stream, and
serve.
Series
serveSeries / connectSeries expose a standalone Series without transferring its sample store. Reads are bounded and borrowed sample buffers are copied before transfer. Closing the connection leaves the source owned by its host.
const stop = serveSeries(port, series, { id: 'voltage', snapshot: true });
const remote = await connectSeries(port, { id: 'voltage', signal });
const block = await remote.read(0, window, signal);
remote.close();
stop();snapshot: true pins the committed prefix at serve time and reports a sealed remote history. Without it, the connection follows appended frames and sealing. A failed live connection stops appends, preserves its last valid state, and rejects subsequent reads and lookups with the original failure. IDs allow multiple histories to share a port.
