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
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
- [Cloud credential sources](cloud-credential-sources.md) — Pear's two auth stores and the resolveCloudAuth() helper that unifies them with ../workforce and ../cloud
- [Cloud agent semantics](cloud-agent-semantics.md) — a Pear "cloud agent" is a provider credential; delete fixed, attach/warm still unbuilt
- [Relayfile workspace model](relayfile-workspace-model.md) — relayWorkspaceId is a fake random UUID; integrations Connect must use the SDK; account-wide workspace decided; catalog logos/filter done
- [Duplicate event hardening](duplicate-event-hardening.md) — broker starts, event streams, PTY chunks, spawned personas, and integration notifications must be idempotent and replay-tolerant
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
# Duplicate Event Hardening

Treat duplicate broker/event delivery as expected behavior in Pear. Renderer startup flows, dashboard reconnects, broker stream refreshes, relay replay, and spawned persona sessions can all surface the same logical event more than once.

When changing broker start, event streaming, PTY output, spawned personas, or integration notifications:

- Make side effects idempotent and return whether state actually changed before notifying agents or publishing metadata.
- Coalesce same-project start/attach calls with keyed in-flight promises.
- Dedupe live events by stable identity (`event_id`, `id`, or `seq`) before falling back to short content-based suppression.
- Use generation tokens for refreshed listeners so stale event streams cannot publish IPC output.
- Preserve PTY-specific duplicate guards in main and tolerate repeated chunk metadata in renderer buffers.
- Add regression tests for replay, reconnect, repeated `ensureBroker()`, and repeated terminal output cases.
14 changes: 14 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
# Agent Instructions

## Duplicate Event Hardening

Pear broker work must treat duplicate delivery as a normal failure mode. Renderer effects, dashboard reconnects, broker event-stream refreshes, and relay replay can all cause the same logical event to be observed more than once.

- 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 that identity first and use short content-based windows only as a fallback.
- 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.
2 changes: 1 addition & 1 deletion electron.vite.config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ import tailwindcss from '@tailwindcss/vite'
export default defineConfig({
main: {
plugins: [externalizeDepsPlugin({
exclude: ['@agent-relay/sdk', '@agent-relay/cloud', '@agent-relay/harness-driver']
exclude: ['@agent-relay/sdk', '@agent-relay/cloud', '@agent-relay/harness-driver', 'zod']
})],
build: {
rollupOptions: {
Expand Down
91 changes: 91 additions & 0 deletions src/main/broker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -270,6 +270,22 @@ describe('BrokerManager local + cloud coexistence', () => {
expect(cloud.shutdown).toHaveBeenCalled()
})

it('reports an existing local broker start as reused', async () => {
const manager = new BrokerManager()
mock.state.nextLocalAgents = []

const firstStart = await manager.start(PROJECT_ID, '/tmp/project-1', 'pear-project-1', undefined as never, [])
const local = lastSpawned()
const secondStart = await manager.start(PROJECT_ID, '/tmp/project-1', 'pear-project-1', undefined as never, [])

expect(firstStart).toBe(true)
expect(secondStart).toBe(false)
expect(mock.HarnessDriverClient.spawn).toHaveBeenCalledTimes(1)
expect(local.shutdown).not.toHaveBeenCalled()

await manager.shutdown()
})

it('keeps the cloud session alive when a local broker starts afterwards', async () => {
const manager = new BrokerManager()
const cloud = await attachCloud(manager, ['cloud-agent'])
Expand Down Expand Up @@ -426,6 +442,81 @@ describe('BrokerManager local + cloud coexistence', () => {
await manager.shutdown()
})

it('dedupes repeated PTY chunks from overlapping event streams', 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')

const chunkEvent = {
kind: 'worker_stream',
name: 'claude-1',
chunk: 'pong\n',
seq: 22
}
listener?.(chunkEvent)
listener?.(chunkEvent)

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', 'pong\n']])

await manager.shutdown()
})

it('keeps legitimate repeated PTY chunks with different broker sequences', 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')

listener?.({
kind: 'worker_stream',
name: 'claude-1',
chunk: 'pong\n',
seq: 23
})
listener?.({
kind: 'worker_stream',
name: 'claude-1',
chunk: 'pong\n',
seq: 24
})

const ptyCalls = (win.webContents.send as ReturnType<typeof vi.fn>).mock.calls
.filter(([channel]) => channel === 'broker:pty-chunk')
expect(ptyCalls).toHaveLength(2)

await manager.shutdown()
})

it('keeps repeated PTY chunks when broker events have no identity', 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')

listener?.({
kind: 'worker_stream',
name: 'claude-1',
chunk: 'pong\n'
})
listener?.({
kind: 'worker_stream',
name: 'claude-1',
chunk: 'pong\n'
})

const ptyCalls = (win.webContents.send as ReturnType<typeof vi.fn>).mock.calls
.filter(([channel]) => channel === 'broker:pty-chunk')
expect(ptyCalls).toHaveLength(2)

await manager.shutdown()
})

it('refreshEventStream rebinds the harness stream from the last seen sequence', async () => {
const manager = new BrokerManager()
const local = await startLocal(manager)
Expand Down
54 changes: 47 additions & 7 deletions src/main/broker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -343,6 +343,8 @@ const BROKER_EVENT_HISTORY_TTL_MS = 12 * 60 * 60 * 1_000
const DEFAULT_RECONCILE_MESSAGE_LIMIT = 50
const MAX_RECONCILE_MESSAGE_LIMIT = 100
const EVENT_STREAM_REBIND_COOLDOWN_MS = 5_000
const PTY_CHUNK_IDENTITY_DEDUPE_TTL_MS = 60_000
const MAX_PTY_CHUNK_DEDUPE_ENTRIES = 2_000
// After this many consecutive failures to open a PTY input stream, give up on
// the WS fast path for that agent briefly and send over HTTP while it cools down.
const MAX_INPUT_STREAM_OPEN_FAILURES = 3
Expand Down Expand Up @@ -1075,7 +1077,7 @@ function sessionKeyFor(session: BrokerSession): string {

export class BrokerManager {
private sessions = new Map<string, BrokerSession>()
private startPromises = new Map<string, Promise<void>>()
private startPromises = new Map<string, Promise<boolean | void>>()
private revivePromises = new Map<string, Promise<boolean>>()
// Which broker sessions (by session key) an agent name is registered on.
// Both a project's local and cloud brokers join the same relay workspace,
Expand All @@ -1097,6 +1099,7 @@ export class BrokerManager {
private brokerTimeoutCounts = new Map<string, number>()
private eventObservers = new Set<BrokerEventObserver>()
private eventHistory: BrokerEventRecord[] = []
private recentPtyChunks = new Map<string, number>()
private eventSerial = 0

get cwd(): string | null {
Expand Down Expand Up @@ -1124,7 +1127,7 @@ export class BrokerManager {
name: string,
win: BrowserWindow,
channels: string[] = []
): Promise<void> {
): Promise<boolean> {
const normalizedProjectId = projectId.trim()
if (!normalizedProjectId) {
throw new Error('Project id is required')
Expand All @@ -1140,7 +1143,7 @@ export class BrokerManager {
await this.syncChannels(normalizedProjectId, nextChannels)
await this.refreshEventStream(normalizedProjectId, 'existing-session-start', win)
this.sendStatus(normalizedProjectId, 'connected')
return
return false
}

const inFlight = this.startPromises.get(normalizedProjectId)
Expand All @@ -1156,14 +1159,14 @@ export class BrokerManager {
started.name = name
await this.syncChannels(normalizedProjectId, nextChannels)
this.sendStatus(normalizedProjectId, 'connected')
return
return false
} catch (err) {
this.sendStatusToWindow(win, normalizedProjectId, 'error', String(err))
throw err
}
}

const startBroker = async (): Promise<void> => {
const startBroker = async (): Promise<boolean> => {
const existingClient = await this.connectExistingBroker(normalizedProjectId, cwd)
if (existingClient) {
const unsubEvent = this.attachClient(normalizedProjectId, existingClient, win)
Expand All @@ -1190,7 +1193,7 @@ export class BrokerManager {
source: 'local'
})
this.sendStatus(normalizedProjectId, 'connected')
return
return true
}

const agentRelayMcpCommand = resolveAgentRelayMcpCommand()
Expand Down Expand Up @@ -1241,12 +1244,13 @@ export class BrokerManager {
source: 'local'
})
this.sendStatus(normalizedProjectId, 'connected')
return true
}

const startPromise = startBroker()
this.startPromises.set(normalizedProjectId, startPromise)
try {
await startPromise
return await startPromise
} catch (err) {
console.error(`[broker] Failed to start for project ${normalizedProjectId}:`, err)
this.sendStatusToWindow(win, normalizedProjectId, 'error', String(err))
Expand Down Expand Up @@ -1880,6 +1884,9 @@ export class BrokerManager {
'name' in event && typeof event.name === 'string' &&
'chunk' in event && typeof event.chunk === 'string'
) {
if (this.isDuplicatePtyChunk(sessionKey, event.name, event)) {
return
}
const targetWindow = this.windowForSession(sessionKey, win)
if (targetWindow && !targetWindow.isDestroyed()) {
targetWindow.webContents.send('broker:pty-chunk', projectId, event.name, event.chunk)
Expand Down Expand Up @@ -1932,6 +1939,39 @@ export class BrokerManager {
}
}

private isDuplicatePtyChunk(sessionKey: string, name: string, event: BrokerEvent): boolean {
const now = Date.now()
for (const [key, seenAt] of this.recentPtyChunks) {
if (
now - seenAt > PTY_CHUNK_IDENTITY_DEDUPE_TTL_MS ||
this.recentPtyChunks.size > MAX_PTY_CHUNK_DEDUPE_ENTRIES
) {
this.recentPtyChunks.delete(key)
}
}

const eventRecord = event as Record<string, unknown>
const seq = typeof eventRecord.seq === 'number' || typeof eventRecord.seq === 'string'
? String(eventRecord.seq)
: ''
const eventId = typeof eventRecord.event_id === 'string'
? eventRecord.event_id
: typeof eventRecord.id === 'string'
? eventRecord.id
: ''
const identity = eventId || (seq ? `seq:${seq}` : '')
if (!identity) return false

const key = `${sessionKey}:${name}:${identity}`
const previous = this.recentPtyChunks.get(key)
if (previous !== undefined && now - previous <= PTY_CHUNK_IDENTITY_DEDUPE_TTL_MS) {
return true
}

this.recentPtyChunks.set(key, now)
return false
}

private publishBrokerEvent(
sessionKey: string,
projectId: string,
Expand Down
10 changes: 6 additions & 4 deletions src/main/ipc-handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -179,10 +179,12 @@ export function registerIpcHandlers(): void {
console.warn(`[broker] Project path no longer exists; skipping broker start: ${cwd}`)
return false
}
await brokerManager.start(projectId, cwd, name, win, channels)
void integrationsManager.notifyAgentState(projectId).catch((error) => {
console.warn('[integrations] Failed to notify agents after broker start:', error instanceof Error ? error.message : String(error))
})
const started = await brokerManager.start(projectId, cwd, name, win, channels)
if (started) {
void integrationsManager.notifyAgentState(projectId).catch((error) => {
console.warn('[integrations] Failed to notify agents after broker start:', error instanceof Error ? error.message : String(error))
})
}
return true
})

Expand Down
Loading