Skip to content

Explicit websocket handshake buffering and safer transport request queue - #4

Open
jasonLaster wants to merge 3 commits into
statemachinefrom
codex/implement-channel-specific-safety-model
Open

jasonLaster wants to merge 3 commits into
statemachinefrom
codex/implement-channel-specific-safety-model

Conversation

@jasonLaster

Copy link
Copy Markdown
Owner

Motivation

  • Implement an explicit, per-socket handshake safety model so server.welcome is always first for a new socket while still queuing and delivering request-triggered pushes that happen during the handshake.
  • Prevent timed-out or disposed client requests from being flushed after reconnect by making outbound queue entries keyed by request id and skipping cancelled entries on flush.
  • Keep latestPushByChannel for late in-page subscribers but add explicit reconnect/replay semantics for orchestration event gaps.
  • Make keybindings bootstrap failures non-fatal so the server remains operational while still surfacing issues.

Description

  • Reworked the server push bus in apps/server/src/wsServer/pushBus.ts to allocate global sequence numbers, register clients in pre_welcome phase, buffer publishAll pushes per-client with a bounded backlog, and expose registerClient, activateClient, and unregisterClient APIs that atomically flip a client to active and flush the backlog in sequence order.
  • Simplified WebSocket connection handling in apps/server/src/wsServer.ts to call pushBus.registerClient(ws), send server.welcome via pushBus.publishClient(...), then pushBus.activateClient(ws) and to use pushBus.unregisterClient(ws) on close/error; also replaced the previous Ref-based broadcast set with the push-bus registration model.
  • Changed client transport semantics in apps/web/src/wsTransport.ts by converting the outbound queue from string[] to an array of entries keyed by request id ({id, encoded}), removing queued entries when requests time out or the transport is disposed, and skipping queued entries whose pending request no longer exists when flushing.
  • Added reconnect recovery logic in apps/web/src/routes/__root.tsx that detects a sequence gap and calls api.orchestration.replayEvents(...) (falling back to snapshot sync) instead of relying on latestPushByChannel for log-like event recovery.
  • Made keybindings startup non-fatal by catching start errors and logging a warning instead of failing server bootstrap in apps/server/src/wsServer.ts (keeps existing issues behavior intact).
  • Added/updated tests: apps/server/src/wsServer/pushBus.test.ts (handshake buffering), apps/web/src/wsTransport.test.ts (timeout/dispose/queue semantics), and apps/server/src/wsServer.test.ts (request-triggered pushes during handshake and non-fatal keybindings).

Testing

  • Ran lint: mise exec -- bun lint succeeded with no errors.
  • Ran typecheck: mise exec -- bun typecheck failed in this environment because pinned Effect platform packages could not be fetched (HTTP 403 from the private pkg.pr.new registry), so workspace type dependencies are missing and full typecheck could not complete.
  • Ran unit tests with Vitest: mise exec -- bunx vitest run ... could not run the modified suites due to the same missing effect/platform packages (import errors), so the new tests exist but could not be executed end-to-end in this environment.

Codex Task

Add typed websocket push envelopes with ordered sequence numbers,
structured decode diagnostics, and a client transport state machine.

On the server, route pushes through a shared push bus, gate welcome
delivery on explicit readiness barriers, and make keybindings startup
an explicit runtime with start/ready lifecycle.

Fix the flaky keybindings watcher test on Linux CI:
- Subscribe eagerly via PubSub.subscribe instead of lazy
  Stream.fromPubSub to avoid missing events before the fiber
  is scheduled.
- Replace the hand-rolled Effect.sleep polling loop with
  fs.watchFile (Node's built-in stat-based poller on libuv
  timers) so file change detection is reliable regardless of
  Effect fiber scheduling pressure.

Add orchestration runtime receipts for checkpoint baseline capture,
checkpoint diff finalization, and turn quiescence, then update the
integration harness/tests to wait on those receipts instead of
timing-sensitive polling.

Use DrainableWorker for queue-based reactors so unit tests can
deterministically drain() instead of sleeping.

Adopt idiomatic Effect patterns: Effect.ignoreCause({ log: true })
for log-and-ignore error handling, Effect.tapCause for side-effecting
on errors before ignoring.
@jasonLaster
jasonLaster force-pushed the statemachine branch 6 times, most recently from 332973c to 0711302 Compare March 10, 2026 05:53
@jasonLaster
jasonLaster force-pushed the statemachine branch 5 times, most recently from 47a82c8 to 52410bd Compare March 10, 2026 18:59
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant