Files
xarmian e7b1c3b5ae feat(collab): per-item Room manager with op-log replay + grace TTL (TASK-1255) (#453)
* feat(collab): per-item Room manager with op-log replay + grace TTL (TASK-1255)

Wires the OpBus + op-log + WS handler from prior phase-1 PRs into a
working dumb-relay collab server. Per-item Room created lazily on
first Join, kept alive across transient disconnects via a 60s grace
TTL, reclaimed when the grace expires with no fresh subscribers.

Components:

- internal/collab/room.go — Room struct + lifecycle
  · roomConn pairs (id, conn, bus channel, write mutex). The id is
    server-assigned per WS so writeLoop can suppress own-event echoes
    without decoding the Y.Doc to read the Yjs ClientID.
  · readLoop discriminates yMessageSync vs yMessageAwareness on
    byte 0. Sync frames are persisted to the op-log AND broadcast;
    awareness frames are broadcast only (presence is ephemeral).
    Persistence happens BEFORE broadcast so a crash mid-publish loses
    at most a live keystroke that the originator will replay on
    reconnect anyway.
  · writeLoop drains the bus subscription and writes non-self events
    to the WS, gated by a per-conn write mutex (gorilla's "one writer
    at a time" rule).
  · removeConn arms a 60s graceTimer when the last conn drops; a
    fresh addConn cancels the timer. onGraceExpired re-checks
    len(conns) == 0 under the room mutex and only THEN sets
    closing=true + calls back to the manager. The race between
    "manager.getOrCreate found us" and "grace timer fired" is
    handled by addConn returning errRoomClosing; the manager retries
    via getOrCreate which mints a fresh Room.

- internal/collab/manager.go — RoomManager + RoomManagerConfig
  · NewRoomManager wires production defaults (DefaultGraceTTL = 60s,
    DefaultSchemaVersion = "1"). NewRoomManagerWithConfig accepts an
    explicit config so tests can drop graceTTL to a few ms without
    sleeping a minute. graceTTL is per-manager, not a package var,
    so parallel tests with different TTLs don't trip the race
    detector.
  · Join is the public entry point: getOrCreate → addConn (with
    retry on errRoomClosing) → replayTo → spawn writeLoop goroutine
    → run readLoop inline → wait for writeLoop drain → return. The
    inline read keeps the HTTP handler in scope so its
    `defer conn.Close()` doesn't fire until both loops exit.
  · Close is for graceful server shutdown — closes every active
    conn under the room mutex, then drains the manager's room map.

- internal/collab/manager_test.go — 7 tests covering: lazy create,
  op-log replay-on-connect (two seed rows arrive in order), sync
  broadcast + persist (peer B sees A's frame, originator does not
  echo, op-log gains a row), awareness broadcast WITHOUT persist,
  cross-item isolation (item-a frames don't leak to item-b
  subscribers), grace-TTL reclaim with a 50ms config TTL, grace
  cancel on reconnect within window, manager.Close shuts down
  every active conn. All tests run with -race; the bus's
  concurrent-publish test was already covered by TASK-1253.

- internal/server/handlers_collab.go — wire to RoomManager
  · Returns 503 when s.collab is nil (matches the SSE handler's
    "events bus not configured" 503 — fail loud rather than silently
    accept the upgrade).
  · Otherwise hands the upgraded conn to s.collab.Join, which
    blocks until the WS closes. Unexpected close codes get the same
    warn-log as before; normal closures stay quiet.

- internal/server/server.go — adds *collab.RoomManager field +
  SetCollabRoomManager setter (nil-safe optional, like SetEventBus).

- cmd/pad/main.go — wires NewMemoryOpBus + NewRoomManager into
  the running server alongside the event-bus wiring. Single-instance
  only today; multi-replica fanout via Redis is a deferred IDEA per
  the Plan body.

- internal/server/handlers_collab_test.go — adds
  testServerWithCollab helper (so existing collab tests get a real
  RoomManager) plus TestCollabUpgradeUnavailableWithoutRoomManager
  which asserts the 503 path for unwired servers.

Parent: PLAN-1248. Phase 1 — Backend foundation.

* fix(collab): per-room appendMu + Server.Stop closes RoomManager per Codex review (round 1)

P1 — concurrent peers raced AppendYjsUpdate, violating the
single-writer-per-item contract documented on the store call. Each
peer's readLoop runs in its own goroutine, so two peers in the same
room could call AppendYjsUpdate concurrently. On Postgres that
risks the BIGSERIAL allocation-vs-commit-order cursor gap that
TASK-1252's contract was specifically guarding against. Add an
appendMu on Room held across the persist+publish sequence; reads,
awareness frames, and OTHER rooms remain unserialised.

Regression test (TestRoomManagerSerializesSyncAppends) drives 4
peers × 10 writes concurrently and asserts the op-log gains exactly
40 rows. Without appendMu this would intermittently surface fewer
rows or out-of-order ids on Postgres; with it the count is
deterministic and the race detector stays clean.

P2 — Server.Stop did not close s.collab. Active collab WS goroutines
+ grace timers could keep using s.store after the server's other
cleanup paths winding down. Add s.collab.Close() before
rateLimiters.Stop so any Join goroutines holding rate-limiter
handles can wind down cleanly. nil-safe via the existing collab
optional-attachment pattern.

* fix(collab): start writer before replay to avoid bus-overflow drops per Codex review (round 2)

P2: a joining peer subscribed to live events BEFORE its writer
goroutine started. During a long replay, live sync events would pile
up in the 64-event bus channel; once full, MemoryOpBus.Publish
silently drops them, leaving the new peer connected but permanently
missing those updates.

Restructure runConn to spawn the writer goroutine FIRST so it drains
the bus subscription concurrently with the replay. Both replay and
writer go through rc.writeMessage, which holds the per-conn write
mutex, so we never violate gorilla's one-writer-at-a-time rule.

Yjs CRDTs are commutative — applying live op 100 before replay op 50
yields the same final Y.Doc as the reverse order — so interleaving
is correct. The trade-off is a brief "out of causal order" UX wobble
during replay, which is acceptable: the alternative would require
either an unbounded queue or losing updates the way the original
order did.

* fix(pad): call srv.Stop() in serveCmd shutdown so collab sessions close per Codex review (round 3)

P2: serveCmd's SIGINT/SIGTERM path called srv.Shutdown but never
srv.Stop. http.Server.Shutdown does NOT terminate hijacked
connections (WebSockets), so active collab sessions kept running
until process exit and could race the deferred store close. The
RoomManager.Close path added in round 1 only fires inside Stop, so
without this call the production shutdown was effectively bypassing
the new cleanup.

Add srv.Stop() after srv.Shutdown in the serveCmd shutdown
sequence. Stop also runs the existing background-loop teardowns
(orphan GC, MCP audit writer, MCP session tracker) which were
previously already part of Stop's contract — those will continue to
fire as they always have, so this commit's only behavioural change
is "now also closes the collab room manager".

* fix(collab): WaitGroup drain barrier + bigger bus buffer per Codex review (round 3)

P1 — RoomManager.Close was not a true drain barrier. closeAll
closed the WebSockets but did NOT wait for the corresponding Join
goroutines (running runConn) to exit. Server.Stop returned before
in-flight collab work finished, racing the deferred store close
on process exit. Fix: track every Join in m.activeJoins
(sync.WaitGroup); Close iterates closeAll first (waking up every
reader by closing the conn), then activeJoins.Wait — guaranteeing
no collab goroutine is still running by the time Close returns.

P2 — replay-time bus overflow could still drop sync events on a
slow drain (writeLoop blocks on the same writeMu replayTo holds,
so a long replay starves the bus drain even with the writer
goroutine started before replay). Two-part response:

(a) Bump the per-subscriber bus channel buffer from 64 to 256.
Sized for a 5x safety margin on a 1k-row replay against a
chatty 5-peer room (~50 events/sec during a ~1s replay).

(b) The architectural fix — force-close subscribers on overflow,
honoring the bus's documented slow-peer recovery contract — is
filed as TASK-1273 follow-up. That requires extending the OpBus
interface (per-subscriber drop callback or counter) and an active
health-check tick in the room manager; both are out of scope for
TASK-1255's "lazy room + grace TTL" deliverable.

For PLAN-1248's single-instance scope and typical editor load,
256 covers realistic workloads. Pathological / load-test scenarios
exposing overflow can recover via Yjs's state-vector negotiation
on reconnect, and TASK-1273 will tighten that to an active kick.

* fix(collab): closed flag gates Join + Close idempotency per Codex review (round 4)

P2: http.Server.Shutdown does NOT wait for hijacked WebSocket
handlers, so a Join() call from a freshly-upgraded conn could fire
AFTER Close() returned. The previous Add-then-Wait pattern was
correct for already-started Joins but couldn't catch a Join that
hadn't yet hit Add when Close fired. Race: Close iterates the (empty)
rooms map, Wait sees zero waiters, Close returns; THEN Join hits
Add and proceeds against a torn-down store.

Add a `closed` flag gated by the same mutex that wraps
activeJoins.Add. Three orderings, all safe:

  1. Add before Close.closed=true → Wait blocks until Done.
  2. Close.closed=true before Add → Join sees closed=true under
     the same lock and returns errManagerClosed without ever
     incrementing the WaitGroup.
  3. Close called twice → second call short-circuits (idempotent).

getOrCreate also gets a closed-flag short-circuit so a future
caller can't bypass the gate by skipping Join.

Test: TestRoomManagerJoinAfterCloseFailsFast asserts post-Close
Join returns errManagerClosed, plus a second Close() is a no-op.

All 15 collab tests pass under -race.
2026-05-08 15:38:45 -04:00

188 lines
6.4 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package collab
import (
"log/slog"
"sync"
"time"
)
// subscriber wraps a delivery channel with the item filter it cares
// about. We map *chan to *subscriber (rather than chan→itemID) so
// MemoryOpBus.Unsubscribe is O(1) and so future enhancements (per-sub
// metrics, last-activity timestamps) have a place to live.
type memSubscriber struct {
ch chan OpEvent
itemID string
}
// MemoryOpBus is the in-process OpBus implementation used in every
// shipping target today (single-binary self-hosted, single-replica
// pad-cloud). Mirrors the shape of internal/events.MemoryBus so a
// future RedisOpBus is a drop-in: same Subscribe/Publish/Close
// surface, same slow-subscriber drop semantics, same cleanup order.
type MemoryOpBus struct {
mu sync.RWMutex
subscribers map[chan OpEvent]*memSubscriber
}
// NewMemoryOpBus returns a ready-to-use in-process bus.
func NewMemoryOpBus() *MemoryOpBus {
return &MemoryOpBus{
subscribers: make(map[chan OpEvent]*memSubscriber),
}
}
// subscriberBufSize is the per-subscriber bus channel buffer. Sized
// generously so a slow drain (replay holding the conn's writeMu, OS
// write buffer momentarily backed up, …) doesn't spill over into the
// drop path under normal editor load.
//
// Sizing rationale: a 1k-row op-log replay at ~1ms/row takes ~1
// second; during that window a chatty 5-peer room produces O(50)
// live events, so 256 events covers a 5× safety margin. Larger
// documents (10k+ rows) under sustained write load can still
// overflow — that's the documented "force-close slow peers" path
// in the bus's Publish doc comment, deferred to a follow-up task.
const subscriberBufSize = 256
// Subscribe registers a subscriber for itemID and returns a buffered
// channel of size subscriberBufSize.
func (b *MemoryOpBus) Subscribe(itemID string) chan OpEvent {
b.mu.Lock()
defer b.mu.Unlock()
ch := make(chan OpEvent, subscriberBufSize)
b.subscribers[ch] = &memSubscriber{
ch: ch,
itemID: itemID,
}
return ch
}
// Unsubscribe removes a subscriber and closes its channel. No-op when
// the channel is unknown — a double-Unsubscribe (e.g. handler defer
// firing after Close has already cleaned up) is safe.
func (b *MemoryOpBus) Unsubscribe(ch chan OpEvent) {
b.mu.Lock()
defer b.mu.Unlock()
if _, ok := b.subscribers[ch]; ok {
delete(b.subscribers, ch)
close(ch)
}
}
// Publish fans an event out to every subscriber whose itemID matches.
//
// Non-blocking: a slow subscriber whose buffer is full has new events
// dropped (logged at warn level so operators can spot a stuck client)
// rather than back-pressuring the broadcast loop. One stuck consumer
// must NEVER poison every other peer in the same room.
//
// Recovery contract for dropped sync events:
//
// The room manager (TASK-1255) is the bus's only consumer in
// production. It reads from each subscriber channel and writes to the
// owning peer's WebSocket. A full channel means that peer's
// WebSocket write is backed up — slow network, blocked client,
// half-broken socket, etc. The room manager is responsible for
// detecting that condition (e.g. via the channel-len threshold
// exposed by the room health check, TASK-1255 + TASK-1256) and
// force-closing the slow peer's WebSocket. A fresh reconnect then
// replays everything the peer missed by loading op-log rows since
// the peer's last known id (TASK-1252 + Yjs state-vector
// negotiation). Awareness drops are unrecoverable, but presence is
// ephemeral so that's fine — sync drops are the only case that
// matters for correctness, and they're recoverable via the
// op-log + reconnect path.
//
// The bus deliberately does NOT take corrective action itself: it
// has no concept of WHICH peer owns which channel and no way to
// signal a force-close. That's the room manager's domain.
//
// Mutation contract for OpEvent.Data:
//
// Data is cloned at the publish boundary so subscribers (and any
// other publisher) cannot affect each other through a shared backing
// array. Subscribers MUST treat their received Data as read-only —
// mutating one subscriber's view would still mutate every other
// subscriber's, since they share the same clone. Cloning per
// subscriber would push that cost onto every fan-out; the chosen
// trade-off is "one allocation per Publish, immutability by
// convention on the read side".
func (b *MemoryOpBus) Publish(event OpEvent) {
if event.Timestamp == 0 {
event.Timestamp = time.Now().UnixMilli()
}
// Clone Data so a publisher's later buffer reuse cannot mutate
// what subscribers observe. One allocation per Publish is cheaper
// than a corrupted Yjs decode somewhere downstream — and the
// gorilla/websocket ReadMessage caller IS allowed to reuse its
// read buffer between messages, so we must not assume the input
// slice is owned by us.
if event.Data != nil {
cloned := make([]byte, len(event.Data))
copy(cloned, event.Data)
event.Data = cloned
}
b.mu.RLock()
defer b.mu.RUnlock()
for _, sub := range b.subscribers {
if sub.itemID != event.ItemID {
continue
}
select {
case sub.ch <- event:
default:
slog.Warn(
"collab: dropping op for slow subscriber",
"type", event.Type,
"item_id", event.ItemID,
"client_id", event.ClientID,
)
}
}
}
// SubscriberCount returns the number of active subscribers whose
// itemID matches. Used by the room manager to decide when the active
// peer count has dropped to zero (room enters its 60s grace before
// teardown) and when it climbs from zero (room comes back live).
func (b *MemoryOpBus) SubscriberCount(itemID string) int {
b.mu.RLock()
defer b.mu.RUnlock()
n := 0
for _, sub := range b.subscribers {
if sub.itemID == itemID {
n++
}
}
return n
}
// Close shuts down the bus. All subscriber channels are closed under
// the same write-lock that gates Publish/Subscribe so a final inflight
// Publish can't race a Close into delivering on a closed channel.
//
// After Close returns, the bus must not be used. Subscribe will leak
// goroutines blocked on the closed channel; Publish becomes a no-op
// over an empty subscriber map but is otherwise undefined.
func (b *MemoryOpBus) Close() {
b.mu.Lock()
defer b.mu.Unlock()
for ch := range b.subscribers {
delete(b.subscribers, ch)
close(ch)
}
}
// Compile-time assertion that MemoryOpBus satisfies OpBus. Catches a
// drifted interface signature at build time rather than at the first
// caller site.
var _ OpBus = (*MemoryOpBus)(nil)