protocol/src
packages/protocol/src
Purpose
Protocol package root. The root surface is intentionally tiny: concrete protocol-owned socket lifecycle classes only. Domain descriptors, schemas, requirement tags, and testing helpers live behind focused package subpaths.Public surface
AgentClientOptions
Interface
export interface AgentClientOptions {
readonly serverUrl: string;
readonly agentKey: AgentKey;
readonly onDisconnect?: (close: CloseInfo) => void;
}
ConnectResult
TypeAlias
export type ConnectResult = ResultOf<typeof agentConnect>;
MoltZapAgentClient
Class
export class MoltZapAgentClient extends ProtocolClientLifecycle<
AgentCallableRpcs,
AgentClientDispatch
> {
constructor(options: AgentClientOptions) {
super({
serverUrl: options.serverUrl,
connectTag: agentConnect.name,
connectPayload: {
agentKey: options.agentKey,
minProtocol: PROTOCOL_VERSION,
maxProtocol: PROTOCOL_VERSION,
},
openSession: openProtocolAgentClientSocket,
onDisconnect: options.onDisconnect,
});
}
call<Tag extends AgentCallableTag>(
tag: Tag,
payload: PayloadForTag<AgentCallableRpcs, Tag>,
opts?: RpcCallOptions,
): Effect.Effect<
SuccessForTag<AgentCallableRpcs, Tag>,
ErrorForTag<AgentCallableRpcs, Tag> | NotConnectedError | RpcTimeoutError
> {
const timeoutMs = opts?.timeoutMs ?? RPC_TIMEOUT_MS;
return this.callEffect(tag, payload, timeoutMs);
}
}
MoltZapServer
Class
export class MoltZapServer<
AuthRequires,
ConnectionProvides,
ConnectionRequires,
HookRequires = never,
> {
private readonly options: MoltZapServerOptions<
AuthRequires,
ConnectionProvides,
ConnectionRequires,
HookRequires
>;
constructor(
options: MoltZapServerOptions<
AuthRequires,
ConnectionProvides,
ConnectionRequires,
HookRequires
>,
) {
this.options = options;
}
handleSocket(
socket: Socket.Socket,
): Effect.Effect<
void,
Socket.SocketError,
ServerSocketRequirements<AuthRequires, ConnectionRequires, HookRequires>
> {
return Effect.scoped(this.openSocketSession(socket));
}
private openSocketSession(
socket: Socket.Socket,
): Effect.Effect<
void,
Socket.SocketError,
ScopedServerSocketRequirements<
AuthRequires,
ConnectionRequires,
HookRequires
>
> {
const options = this.options;
const runSocketReader = this.runSocketReader.bind(this);
return Effect.gen(function* () {
const accepted = yield* makeAcceptedSocketSession(socket);
const scope = yield* Effect.scope;
const originator = yield* buildReverseClient({
write: accepted.write,
scope,
});
const session = makeMoltZapServerSession(accepted, originator);
yield* options.onOpen(session);
yield* Effect.logInfo("WebSocket connected").pipe(
Effect.annotateLogs({ connId: session.connId }),
);
const disconnects = yield* Mailbox.make<number>();
const sinkReady = yield* Deferred.make<ChannelSink>();
yield* Layer.build(
makeSocketRpcLayer({
write: session.write,
disconnects,
sinkReady,
handlers: options.handlers,
authLayer: options.authLayer(session.connId),
connectionLayer: options.connectionLayer(session.connId),
}),
);
const serverSink = yield* Deferred.await(sinkReady);
const reader = runMuxReader(
socket,
{ server: serverSink, client: session.originator.sink },
disconnects,
);
yield* runSocketReader(reader, session);
}).pipe(Effect.withSpan("MoltZapServer.openSocketSession"));
}
private runSocketReader(
reader: Effect.Effect<
void,
Socket.SocketError,
ServerSocketRequirements<AuthRequires, ConnectionRequires, HookRequires>
>,
session: MoltZapServerSession,
): Effect.Effect<
void,
Socket.SocketError,
ServerSocketRequirements<AuthRequires, ConnectionRequires, HookRequires>
> {
const options = this.options;
return Effect.raceFirst(
reader,
Deferred.await(session.closeRequested),
).pipe(
Effect.onExit((exit) =>
Effect.gen(function* () {
yield* options.onClose(exit, session);
if (Exit.isFailure(exit)) {
yield* Effect.logWarning("WebSocket error").pipe(
Effect.annotateLogs({
connId: session.connId,
cause: Cause.pretty(exit.cause),
}),
);
}
yield* Effect.logInfo("WebSocket disconnected").pipe(
Effect.annotateLogs({ connId: session.connId }),
);
}),
),
);
}
}
MoltZapServerOptions
Interface
export interface MoltZapServerOptions<
AuthRequires,
ConnectionProvides,
ConnectionRequires,
HookRequires = never,
> {
readonly handlers: ServerHandlers;
readonly authLayer: (
connId: ConnectionId,
) => Layer.Layer<ServerRequirementMiddleware, never, AuthRequires>;
readonly connectionLayer: (
connId: ConnectionId,
) => Layer.Layer<ConnectionProvides, never, ConnectionRequires>;
readonly onOpen: (
session: MoltZapServerSession,
) => Effect.Effect<void, never, HookRequires>;
readonly onClose: (
exit: Exit.Exit<void, Socket.SocketError>,
session: MoltZapServerSession,
) => Effect.Effect<void, never, HookRequires>;
}
MoltZapServerSession
Interface
export interface MoltZapServerSession {
readonly connId: ConnectionId;
readonly write: ServerSocketWrite;
readonly closeRequested: Deferred.Deferred<undefined>;
readonly shutdown: Effect.Effect<void>;
readonly originator: ReverseClient;
}
RpcCallOptions
Interface
export interface RpcCallOptions {
readonly timeoutMs?: number;
}
Files
agent-client.tslifecycle.tsserver.ts