Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
node_modules/
node_modules
out/
dist/
bin/relayfile-mount
Expand Down
4 changes: 2 additions & 2 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,12 +6,12 @@ Pear broker work must treat duplicate delivery as a normal failure mode. Rendere

- Make lifecycle operations idempotent. `broker:start`, agent registration, integration notifications, and mount/link setup should return or record whether they actually changed state before triggering side effects.
- Coalesce concurrent starts or attaches with keyed in-flight promises. Repeated UI calls for the same project/root/channels should wait on the existing operation instead of starting another broker or event stream.
- Prefer stable event identity over content matching. PTY and broker events should carry `event_id`, `id`, or `seq`; dedupe by identity AND content hash together. Never drop a PTY chunk on content alone (identical consecutive chunks are normal terminal traffic — repeated keystroke echoes, byte-identical TUI repaint frames) and never on seq alone (a daemon restart resets seq; dropping fresh low-seq chunks corrupts the escape-sequence stream — the stacked-duplicate-lines rendering corruption). When in doubt, deliver: a rare duplicate repaints once; a dropped chunk mangles the screen until the next full repaint.
- Prefer stable event identity over content matching, and use the identity fields each event *actually* carries. Chat/broker events carry `event_id`, `id`, or `seq`; dedupe those by identity AND content hash together. `worker_stream` PTY chunks do **not** carry `seq`/`event_id` — those are ephemeral and excluded from the daemon replay buffer, so any dedup keyed on them is structurally inert (it can never match). PTY chunks instead carry a cumulative per-worker byte `offset`; dedupe them by `(generation, offset)` identity AND content hash, where `generation` is the broker event-stream generation the delivering listener was attached under (it scopes the offset so a fresh worker stream's low offset can't collide with a stale remembered one). Never drop a PTY chunk on content alone (identical consecutive chunks are normal terminal traffic — repeated keystroke echoes, byte-identical TUI repaint frames) and never on identity alone (a repeated offset with different bytes is fresh output, not a replay — dropping it corrupts the escape-sequence stream into the stacked-duplicate-lines rendering corruption). When correlation metadata is absent (no `offset`), deliver every chunk and log the blind spot loudly — do not silently claim a protection that isn't running. When in doubt, deliver: a rare duplicate repaints once; a dropped chunk mangles the screen until the next full repaint.
- Scope live event listeners with a generation token when reconnecting or refreshing streams. Stale callbacks from an older listener must bail before publishing IPC events.
- Keep PTY stream delivery separately guarded in main and renderer code. Main should suppress duplicate `worker_stream` chunks before `broker:pty-chunk`; renderer buffers should tolerate repeated chunk metadata as a final guardrail.
- Do not post integration or launch metadata on reused broker sessions. Notify agents only after a real broker start, reconnect, or state transition, and make repeated payloads no-ops when possible.
- Add regression tests when touching broker start, event streaming, PTY buffering, spawned personas, or integration notifications. Include duplicate/replay cases, not just the happy path.
- Add low-noise telemetry for suppressed duplicates and missing event identity so replay issues are visible without flooding the terminal.
- Add low-noise telemetry for suppressed duplicates and a rate-limited *loud* warning (with a running count) when correlation metadata is missing, so both real replays and dedup blind spots are visible without flooding the terminal.

## Terminal Screen Convergence

Expand Down
119 changes: 47 additions & 72 deletions src/main/broker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2039,24 +2039,27 @@ exit 2
await manager.shutdown()
})

it('keeps legitimate repeated PTY chunks with different broker sequences', async () => {
it('keeps legitimate repeated PTY chunks with different offsets', async () => {
const manager = new BrokerManager()
const win = createMockWindow()
const local = await startLocalWithWindow(manager, win)
const listener = local.onEvent.mock.calls.at(-1)?.[0]
expect(listener).toBeTypeOf('function')

// Byte-identical output at advancing offsets is normal terminal traffic
// (a repeated keystroke echo). Distinct offsets are distinct identities;
// both must render — never drop on content alone.
listener?.({
kind: 'worker_stream',
name: 'claude-1',
chunk: 'pong\n',
seq: 23
offset: 23
})
listener?.({
kind: 'worker_stream',
name: 'claude-1',
chunk: 'pong\n',
seq: 24
offset: 24
})

const ptyCalls = (win.webContents.send as ReturnType<typeof vi.fn>).mock.calls
Expand All @@ -2066,10 +2069,10 @@ exit 2
await manager.shutdown()
})

it('delivers identical PTY chunks when broker events have no identity', async () => {
it('delivers identical PTY chunks when broker events have no offset', async () => {
// Identical consecutive chunks are NORMAL terminal traffic (the same
// keystroke echoed twice, byte-identical TUI repaint frames). Without an
// identity there is no way to tell a transport replay from real output,
// offset there is no way to tell a transport replay from real output,
// and dropping real bytes mangles escape sequences — the stacked
// "duplicate line" rendering corruption. Never drop on content alone.
const manager = new BrokerManager()
Expand Down Expand Up @@ -2099,11 +2102,11 @@ exit 2
await manager.shutdown()
})

it('logs the identity-less PTY stream blind spot once per stream while delivering', async () => {
// AGENTS.md: low-noise telemetry for missing event identity. One line per
// stream, not per chunk — an identity-less stream hits the branch on
// every chunk and per-chunk logging would flood.
const infoSpy = vi.spyOn(console, 'info').mockImplementation(() => {})
it('loudly warns about the offset-less PTY stream blind spot per stream while delivering', async () => {
// AGENTS.md: loud, rate-limited telemetry when correlation metadata is
// absent. Without an offset the deduper is blind — it warns (at most once
// per interval per stream, carrying a running count) but always delivers.
const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {})
try {
const manager = new BrokerManager()
const win = createMockWindow()
Expand All @@ -2115,10 +2118,12 @@ exit 2
listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'two\n' })
listener?.({ kind: 'worker_stream', name: 'codex-1', chunk: 'three\n' })

const blindSpotLogs = infoSpy.mock.calls.filter(([first]) =>
typeof first === 'string' && first.includes('no seq/event_id')
const blindSpotLogs = warnSpy.mock.calls.filter(([first]) =>
typeof first === 'string' && first.includes('no offset')
)
expect(blindSpotLogs).toHaveLength(2) // once for claude-1, once for codex-1
// claude-1's two chunks land in one rate-limit interval → one warn;
// codex-1 is a distinct stream → its own warn.
expect(blindSpotLogs).toHaveLength(2)
expect(blindSpotLogs[0][0]).toContain('claude-1')
expect(blindSpotLogs[1][0]).toContain('codex-1')

Expand All @@ -2129,108 +2134,78 @@ exit 2

await manager.shutdown()
} finally {
infoSpy.mockRestore()
warnSpy.mockRestore()
}
})

it('delivers chunks whose seq repeats with different bytes (daemon seq reset)', async () => {
// A daemon restart resets its event seq counter. The replacement stream
// reuses seq numbers we have already seen — but the bytes are new output
// and must render. Only an (identity AND content) match is a replay.
it('delivers chunks whose offset repeats with different bytes (never drop on identity alone)', async () => {
// An offset that repeats with different bytes is an anomaly, not a replay
// (offset is meant to be unique within a generation+worker). The bytes are
// real output and must render — only an (identity AND content) match drops.
const manager = new BrokerManager()
const win = createMockWindow()
const local = await startLocalWithWindow(manager, win)
const listener = local.onEvent.mock.calls.at(-1)?.[0]
expect(listener).toBeTypeOf('function')

listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'before restart\n', seq: 90 })
listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'after restart\n', seq: 90 })
listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'before\n', offset: 90 })
listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'after\n', offset: 90 })

const ptyCalls = (win.webContents.send as ReturnType<typeof vi.fn>).mock.calls
.filter(([channel]) => channel === 'broker:pty-chunk')
expect(ptyCalls).toEqual([
['broker:pty-chunk', PROJECT_ID, 'claude-1', 'before restart\n'],
['broker:pty-chunk', PROJECT_ID, 'claude-1', 'after restart\n']
])
expect(ptyCalls.map((call) => call[3])).toEqual(['before\n', 'after\n'])

await manager.shutdown()
})

it('delivers low-seq chunks after a watermark reset instead of dropping a window', async () => {
// After the daemon restarts, fresh chunks arrive with seqs far below the
// old watermark. They must all be delivered immediately — the previous
// TTL-based dedup dropped them for up to 60s, losing repaint bytes.
it('drops a delayed replay of an earlier offset with identical bytes', async () => {
// A double-emit re-delivers recent chunks: same offset, same bytes,
// possibly long after first delivery. These are the one provable-duplicate
// case — full identity (generation + offset) AND content match.
const manager = new BrokerManager()
const win = createMockWindow()
const local = await startLocalWithWindow(manager, win)
const listener = local.onEvent.mock.calls.at(-1)?.[0]
expect(listener).toBeTypeOf('function')

listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'old-high\n', seq: 5_000 })
listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'fresh-1\n', seq: 1 })
listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'fresh-2\n', seq: 2 })
listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'a\n', offset: 10 })
listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'b\n', offset: 11 })
// The tail of the stream is re-emitted.
listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'a\n', offset: 10 })
listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'b\n', offset: 11 })

const ptyCalls = (win.webContents.send as ReturnType<typeof vi.fn>).mock.calls
.filter(([channel]) => channel === 'broker:pty-chunk')
expect(ptyCalls).toEqual([
['broker:pty-chunk', PROJECT_ID, 'claude-1', 'old-high\n'],
['broker:pty-chunk', PROJECT_ID, 'claude-1', 'fresh-1\n'],
['broker:pty-chunk', PROJECT_ID, 'claude-1', 'fresh-2\n']
])
expect(ptyCalls.map((call) => call[3])).toEqual(['a\n', 'b\n'])

await manager.shutdown()
})

it('drops a delayed replay of an older seq with identical bytes', async () => {
// Rebind replay re-delivers recent events: same seq, same bytes, possibly
// long after first delivery. These are the one provable-duplicate case.
it('tracks PTY dedup identity per agent stream', async () => {
// Two agents share the session event stream; one agent's offset window must
// not swallow the other's chunks even at a colliding offset.
const manager = new BrokerManager()
const win = createMockWindow()
const local = await startLocalWithWindow(manager, win)
const listener = local.onEvent.mock.calls.at(-1)?.[0]
expect(listener).toBeTypeOf('function')

listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'a\n', seq: 10 })
listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'b\n', seq: 11 })
// Rebind replays the tail of the stream.
listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'a\n', seq: 10 })
listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'b\n', seq: 11 })
listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'one\n', offset: 100 })
listener?.({ kind: 'worker_stream', name: 'codex-1', chunk: 'one\n', offset: 100 })
listener?.({ kind: 'worker_stream', name: 'codex-1', chunk: 'two\n', offset: 101 })

const ptyCalls = (win.webContents.send as ReturnType<typeof vi.fn>).mock.calls
.filter(([channel]) => channel === 'broker:pty-chunk')
expect(ptyCalls).toEqual([
['broker:pty-chunk', PROJECT_ID, 'claude-1', 'a\n'],
['broker:pty-chunk', PROJECT_ID, 'claude-1', 'b\n']
])

await manager.shutdown()
})

it('tracks PTY dedup watermarks per agent stream', async () => {
// Two agents share the session event-stream seq space; one agent's
// watermark must not swallow the other's chunks.
const manager = new BrokerManager()
const win = createMockWindow()
const local = await startLocalWithWindow(manager, win)
const listener = local.onEvent.mock.calls.at(-1)?.[0]
expect(listener).toBeTypeOf('function')

listener?.({ kind: 'worker_stream', name: 'claude-1', chunk: 'one\n', seq: 100 })
listener?.({ kind: 'worker_stream', name: 'codex-1', chunk: 'one\n', seq: 100 })
listener?.({ kind: 'worker_stream', name: 'codex-1', chunk: 'two\n', seq: 101 })

const ptyCalls = (win.webContents.send as ReturnType<typeof vi.fn>).mock.calls
.filter(([channel]) => channel === 'broker:pty-chunk')
expect(ptyCalls).toEqual([
['broker:pty-chunk', PROJECT_ID, 'claude-1', 'one\n'],
['broker:pty-chunk', PROJECT_ID, 'codex-1', 'one\n'],
['broker:pty-chunk', PROJECT_ID, 'codex-1', 'two\n']
expect(ptyCalls.map((call) => [call[2], call[3]])).toEqual([
['claude-1', 'one\n'],
['codex-1', 'one\n'],
['codex-1', 'two\n']
])

await manager.shutdown()
})

it('keeps distinct PTY chunks when broker events have no identity', async () => {
it('keeps distinct PTY chunks when broker events have no offset', async () => {
const manager = new BrokerManager()
const win = createMockWindow()
const local = await startLocalWithWindow(manager, win)
Expand Down
2 changes: 1 addition & 1 deletion src/main/broker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2510,7 +2510,7 @@ export class BrokerManager {
'name' in event && typeof event.name === 'string' &&
'chunk' in event && typeof event.chunk === 'string'
) {
if (this.ptyDeduper.isDuplicatePtyChunk(sessionKey, event.name, event)) {
if (this.ptyDeduper.isDuplicatePtyChunk(sessionKey, event.name, event, eventStreamGeneration)) {
return
}
const targetWindow = this.windowForSession(sessionKey, win)
Expand Down
Loading
Loading