diff --git a/packages/server/src/server/session.ts b/packages/server/src/server/session.ts index 1f977f788..5e1269d68 100644 --- a/packages/server/src/server/session.ts +++ b/packages/server/src/server/session.ts @@ -86,7 +86,6 @@ import { } from "./agent/lifecycle-command.js"; import { buildStoredAgentPayload, - resolveEffectiveThinkingOptionId, resolveStoredAgentPayloadUpdatedAt, toAgentPayload, } from "./agent/agent-projections.js"; @@ -167,6 +166,11 @@ import { createWorkspaceProvisioningService, type WorkspaceProvisioningService, } from "./session/workspace-provisioning/workspace-provisioning-service.js"; +import { + createAgentUpdatesService, + matchesAgentUpdatesFilter, + type AgentUpdatesService, +} from "./session/agent-updates/agent-updates-service.js"; import { expandTilde } from "../utils/path.js"; import { searchHomeDirectories, searchWorkspaceEntries } from "../utils/directory-suggestions.js"; import type { CheckoutDiffManager } from "./checkout-diff-manager.js"; @@ -345,14 +349,7 @@ type FetchAgentsResponsePayload = Extract< >["payload"]; type FetchAgentsResponseEntry = FetchAgentsResponsePayload["entries"][number]; type FetchAgentsResponsePageInfo = FetchAgentsResponsePayload["pageInfo"]; -type AgentUpdatePayload = Extract["payload"]; type AgentUpdatesFilter = FetchAgentsRequestFilter; -interface AgentUpdatesSubscriptionState { - subscriptionId: string; - filter?: AgentUpdatesFilter; - isBootstrapping: boolean; - pendingUpdatesByAgentId: Map; -} type FetchWorkspacesRequestMessage = Extract< SessionInboundMessage, { type: "fetch_workspaces_request" } @@ -560,7 +557,7 @@ export class Session { private readonly pushTokenStore: PushTokenStore; private unsubscribeAgentEvents: (() => void) | null = null; private unsubscribeTerminalWorkspaceContributionEvents: (() => void) | null = null; - private agentUpdatesSubscription: AgentUpdatesSubscriptionState | null = null; + private readonly agentUpdates: AgentUpdatesService; private workspaceUpdatesSubscription: WorkspaceUpdatesSubscriptionState | null = null; private clientActivity: { deviceType: "web" | "mobile"; @@ -808,6 +805,17 @@ export class Session { this.clientCapabilities.has(CLIENT_CAPS.terminalReflowableSnapshot), getClientBufferedAmount: () => this.getTransportBufferedAmount(), }); + this.agentUpdates = createAgentUpdatesService({ + emit: (message) => this.emit(message), + buildAgentPayload: (agent) => this.buildAgentPayload(agent), + buildStoredAgentPayload: (record) => this.buildStoredAgentPayload(record), + isProviderVisibleToClient: (provider) => this.isProviderVisibleToClient(provider), + buildProjectPlacementForWorkspaceId: (workspaceId) => + this.buildProjectPlacementForWorkspaceId(workspaceId), + emitWorkspaceUpdateForWorkspaceId: (workspaceId) => + this.emitWorkspaceUpdateForWorkspaceId(workspaceId), + logger: this.sessionLogger, + }); this.createAgentLifecycleDispatch = new CreateAgentLifecycleDispatch({ paseoHome: this.paseoHome, worktreesRoot: this.worktreesRoot, @@ -822,14 +830,7 @@ export class Session { listActiveWorkspaces: () => this.listActiveWorkspaceRefs(), archiveWorkspaceRecord: (workspaceId) => this.archiveWorkspaceRecord(workspaceId), emit: (message) => this.emit(message), - emitAgentRemove: (agentId) => { - if (this.agentUpdatesSubscription) { - this.bufferOrEmitAgentUpdate(this.agentUpdatesSubscription, { - kind: "remove", - agentId, - }); - } - }, + emitAgentRemove: (agentId) => this.agentUpdates.removeAgent(agentId), emitWorkspaceUpdatesForWorkspaceIds: (workspaceIds) => this.emitWorkspaceUpdatesForWorkspaceIds(workspaceIds), markWorkspaceArchiving: (workspaceIds, archivingAt) => @@ -1171,7 +1172,7 @@ export class Session { }, "agent.session.forward_update", ); - void this.forwardAgentUpdate(event.agent); + void this.agentUpdates.forwardLiveAgent(event.agent); return; } @@ -1286,176 +1287,6 @@ export class Session { return LEGACY_PROVIDER_IDS.has(provider); } - private agentThinkingOptionMatchesFilter( - agent: AgentSnapshotPayload, - filter: AgentUpdatesFilter, - ): boolean { - if (filter.thinkingOptionId === undefined) { - return true; - } - const expectedThinkingOptionId = resolveEffectiveThinkingOptionId({ - configuredThinkingOptionId: filter.thinkingOptionId ?? null, - }); - const resolvedThinkingOptionId = - agent.effectiveThinkingOptionId ?? - resolveEffectiveThinkingOptionId({ - runtimeInfo: agent.runtimeInfo, - configuredThinkingOptionId: agent.thinkingOptionId ?? null, - }); - return resolvedThinkingOptionId === expectedThinkingOptionId; - } - - private matchesAgentStructuralFilter( - agent: AgentSnapshotPayload, - project: ProjectPlacementPayload, - filter: AgentUpdatesFilter, - ): boolean { - if (filter.statuses && filter.statuses.length > 0) { - const statuses = new Set(filter.statuses); - if (!statuses.has(agent.status)) { - return false; - } - } - - if (typeof filter.requiresAttention === "boolean") { - const requiresAttention = agent.requiresAttention ?? false; - if (requiresAttention !== filter.requiresAttention) { - return false; - } - } - - if (filter.projectKeys && filter.projectKeys.length > 0) { - const projectKeys = new Set(filter.projectKeys.filter((item) => item.trim().length > 0)); - if (projectKeys.size > 0 && !projectKeys.has(project.projectKey)) { - return false; - } - } - return true; - } - - private matchesAgentFilter(options: { - agent: AgentSnapshotPayload; - project: ProjectPlacementPayload; - filter?: AgentUpdatesFilter; - }): boolean { - const { agent, project, filter } = options; - - if (filter?.labels) { - const matchesLabels = Object.entries(filter.labels).every( - ([key, value]) => agent.labels[key] === value, - ); - if (!matchesLabels) { - return false; - } - } - - const includeArchived = filter?.includeArchived ?? false; - if (!includeArchived && agent.archivedAt) { - return false; - } - - if (filter && !this.agentThinkingOptionMatchesFilter(agent, filter)) { - return false; - } - - if (filter && !this.matchesAgentStructuralFilter(agent, project, filter)) { - return false; - } - - return true; - } - - private getAgentUpdateTargetId(update: AgentUpdatePayload): string { - return update.kind === "remove" ? update.agentId : update.agent.id; - } - - private bufferOrEmitAgentUpdate( - subscription: AgentUpdatesSubscriptionState, - payload: AgentUpdatePayload, - ): void { - if (payload.kind === "upsert" && !this.isProviderVisibleToClient(payload.agent.provider)) { - return; - } - if (subscription.isBootstrapping) { - subscription.pendingUpdatesByAgentId.set(this.getAgentUpdateTargetId(payload), payload); - return; - } - - this.emit({ - type: "agent_update", - payload, - }); - } - - private async emitStoredAgentUpdate(record: StoredAgentRecord): Promise { - const payload = this.buildStoredAgentPayload(record); - const subscription = this.agentUpdatesSubscription; - if (!subscription) { - return payload; - } - - const project = payload.workspaceId - ? await this.buildProjectPlacementForWorkspaceId(payload.workspaceId) - : null; - if (!project) { - this.bufferOrEmitAgentUpdate(subscription, { - kind: "remove", - agentId: payload.id, - }); - return payload; - } - - const matches = this.matchesAgentFilter({ - agent: payload, - project, - filter: subscription.filter, - }); - this.bufferOrEmitAgentUpdate( - subscription, - matches - ? { - kind: "upsert", - agent: payload, - project, - } - : { - kind: "remove", - agentId: payload.id, - }, - ); - return payload; - } - - private flushBootstrappedAgentUpdates(options?: { - snapshotUpdatedAtByAgentId?: Map; - }): void { - const subscription = this.agentUpdatesSubscription; - if (!subscription || !subscription.isBootstrapping) { - return; - } - - subscription.isBootstrapping = false; - const pending = Array.from(subscription.pendingUpdatesByAgentId.values()); - subscription.pendingUpdatesByAgentId.clear(); - - for (const payload of pending) { - if (payload.kind === "upsert") { - const snapshotUpdatedAt = options?.snapshotUpdatedAtByAgentId?.get(payload.agent.id); - if (typeof snapshotUpdatedAt === "number") { - const updateUpdatedAt = Date.parse(payload.agent.updatedAt); - if (!Number.isNaN(updateUpdatedAt) && updateUpdatedAt <= snapshotUpdatedAt) { - continue; - } - } - } - - this.emit({ - type: "agent_update", - payload, - }); - } - } - private async buildProjectPlacementForWorkspace( workspace: PersistedWorkspaceRecord, projectRecord?: PersistedProjectRecord | null, @@ -1494,51 +1325,6 @@ export class Session { return this.buildProjectPlacementForWorkspace(workspace, project); } - private async forwardAgentUpdate(agent: ManagedAgent): Promise { - try { - const subscription = this.agentUpdatesSubscription; - const payload = await this.buildAgentPayload(agent); - if (subscription) { - const project = payload.workspaceId - ? await this.buildProjectPlacementForWorkspaceId(payload.workspaceId) - : null; - if (!project) { - this.bufferOrEmitAgentUpdate(subscription, { - kind: "remove", - agentId: payload.id, - }); - } else { - const matches = this.matchesAgentFilter({ - agent: payload, - project, - filter: subscription.filter, - }); - - if (matches) { - this.bufferOrEmitAgentUpdate(subscription, { - kind: "upsert", - agent: payload, - project, - }); - } else { - this.bufferOrEmitAgentUpdate(subscription, { - kind: "remove", - agentId: payload.id, - }); - } - } - } - - // A lifecycle change updates exactly the agent's owning workspace, never - // every workspace sharing its cwd. Ownership is the agent's workspaceId. - if (payload.workspaceId) { - await this.emitWorkspaceUpdateForWorkspaceId(payload.workspaceId); - } - } catch (error) { - this.sessionLogger.error({ err: error }, "Failed to emit agent update"); - } - } - /** * Main entry point for processing session messages */ @@ -2060,12 +1846,7 @@ export class Session { }, }); - if (this.agentUpdatesSubscription) { - this.bufferOrEmitAgentUpdate(this.agentUpdatesSubscription, { - kind: "remove", - agentId, - }); - } + this.agentUpdates.removeAgent(agentId); if (knownWorkspaceId) { await this.emitWorkspaceUpdateForWorkspaceId(knownWorkspaceId); @@ -2099,8 +1880,8 @@ export class Session { agentId, ); - if (this.agentUpdatesSubscription) { - const payload = await this.emitStoredAgentUpdate(archivedRecord); + if (this.agentUpdates.hasSubscription()) { + const payload = await this.agentUpdates.emitStoredRecord(archivedRecord); if (payload.workspaceId) { await this.emitWorkspaceUpdateForWorkspaceId(payload.workspaceId); } @@ -2117,7 +1898,7 @@ export class Session { const affectedWorkspaceIds = new Set(); if (!result.live) { - const payload = await this.emitStoredAgentUpdate(result.record); + const payload = await this.agentUpdates.emitStoredRecord(result.record); if (payload.workspaceId) { affectedWorkspaceIds.add(payload.workspaceId); } @@ -2678,7 +2459,7 @@ export class Session { if (!createdWorktree && msg.workspaceId) { await this.writeInitialWorkspaceTitleIfUntitled(workspaceId, workspacePromptTitle); } - await this.forwardAgentUpdate(snapshot); + await this.agentUpdates.forwardLiveAgent(snapshot); if (!createdWorktree && trimmedPrompt) { await this.scheduleAutoNameLocalWorkspaceTitleForFirstAgent({ workspaceId, @@ -2775,7 +2556,7 @@ export class Session { const snapshot = await this.agentManager.resumeAgentFromPersistence(handle, overrides); await unarchiveAgentState(this.agentStorage, this.agentManager, snapshot.id); await this.agentManager.hydrateTimelineFromProvider(snapshot.id); - await this.forwardAgentUpdate(snapshot); + await this.agentUpdates.forwardLiveAgent(snapshot); const timelineSize = this.agentManager.getTimeline(snapshot.id).length; if (requestId) { const agentPayload = await this.buildAgentPayload(snapshot); @@ -2925,7 +2706,7 @@ export class Session { ); } await this.agentManager.hydrateTimelineFromProvider(agentId); - await this.forwardAgentUpdate(snapshot); + await this.agentUpdates.forwardLiveAgent(snapshot); const timelineSize = this.agentManager.getTimeline(agentId).length; if (requestId) { this.emit({ @@ -3735,7 +3516,7 @@ export class Session { continue; } if ( - !this.matchesAgentFilter({ + !matchesAgentUpdatesFilter({ agent: entry.agent, project: entry.project, filter, @@ -4509,12 +4290,10 @@ export class Session { try { if (subscriptionId) { - this.agentUpdatesSubscription = { + this.agentUpdates.beginSubscription({ subscriptionId, filter: request.filter, - isBootstrapping: true, - pendingUpdatesByAgentId: new Map(), - }; + }); } const payload = await this.listFetchAgentsEntries(request); @@ -4535,12 +4314,12 @@ export class Session { }, }); - if (subscriptionId && this.agentUpdatesSubscription?.subscriptionId === subscriptionId) { - this.flushBootstrappedAgentUpdates({ snapshotUpdatedAtByAgentId }); + if (subscriptionId) { + this.agentUpdates.flushBootstrapped(subscriptionId, { snapshotUpdatedAtByAgentId }); } } catch (error) { - if (subscriptionId && this.agentUpdatesSubscription?.subscriptionId === subscriptionId) { - this.agentUpdatesSubscription = null; + if (subscriptionId) { + this.agentUpdates.clearSubscription(subscriptionId); } const code = error instanceof SessionRequestError ? error.code : "fetch_agents_failed"; const message = error instanceof Error ? error.message : "Failed to fetch agents"; @@ -5977,6 +5756,7 @@ export class Session { this.unsubscribeAgentEvents(); this.unsubscribeAgentEvents = null; } + this.agentUpdates.dispose(); if (this.unsubscribeTerminalWorkspaceContributionEvents) { this.unsubscribeTerminalWorkspaceContributionEvents(); this.unsubscribeTerminalWorkspaceContributionEvents = null; diff --git a/packages/server/src/server/session.workspaces.test.ts b/packages/server/src/server/session.workspaces.test.ts index 5ce8b27ea..b5873873d 100644 --- a/packages/server/src/server/session.workspaces.test.ts +++ b/packages/server/src/server/session.workspaces.test.ts @@ -15,6 +15,7 @@ import { z } from "zod"; import { Session } from "./session.js"; import type { SessionOptions } from "./session.js"; +import type { AgentUpdatesService } from "./session/agent-updates/agent-updates-service.js"; import type { AgentSnapshotPayload, SessionOutboundMessage } from "@getpaseo/protocol/messages"; import type { TerminalManager } from "../terminal/terminal-manager.js"; import { createTerminalManager } from "../terminal/terminal-manager.js"; @@ -113,7 +114,7 @@ interface SessionTestAccess { get(workspaceId: string): Promise; upsert(record: unknown): Promise; }; - agentUpdatesSubscription: unknown; + agentUpdates: AgentUpdatesService; workspaceUpdatesSubscription: unknown; interruptAgentIfRunning(agentId: string): unknown; recreateOwningWorktreeForRestore( @@ -128,7 +129,6 @@ interface SessionTestAccess { [key: string]: unknown; }>; reconcileAndEmitWorkspaceUpdates(...args: unknown[]): Promise; - forwardAgentUpdate(...args: unknown[]): Promise; handleArchiveAgentRequest(agentId: string, requestId: string): Promise; handleMessage(message: unknown): Promise; handleCreatePaseoWorktreeRequest(params: unknown): Promise; @@ -182,6 +182,22 @@ function asTestSession(session: Session | TestSession): TestSession { return asSessionInternals(session); } +type AgentUpdatesSubscriptionFilter = Parameters< + AgentUpdatesService["beginSubscription"] +>[0]["filter"]; + +// Drives the agent-updates module to a live (non-bootstrapping) subscription — +// the post-extraction equivalent of assigning a subscription with +// `isBootstrapping: false`. begin → flush leaves an empty buffer and emits nothing. +function activateAgentUpdatesSubscription( + session: TestSession, + subscriptionId: string, + filter?: AgentUpdatesSubscriptionFilter, +): void { + session.agentUpdates.beginSubscription({ subscriptionId, filter }); + session.agentUpdates.flushBootstrapped(subscriptionId); +} + const AgentIdEntrySchema = z.object({ agent: z.object({ id: z.string() }) }); function makeAgent(input: { @@ -1084,14 +1100,9 @@ test("agent_update placement does not refresh git snapshots", async () => { session.workspaceRegistry.list = async () => [workspace]; session.workspaceRegistry.get = async (id: string) => id === workspace.workspaceId ? workspace : null; - session.agentUpdatesSubscription = { - subscriptionId: "sub-agents", - filter: {}, - isBootstrapping: false, - pendingUpdatesByAgentId: new Map(), - }; + activateAgentUpdatesSubscription(session, "sub-agents", {}); - await session.forwardAgentUpdate( + await session.agentUpdates.forwardLiveAgent( makeManagedAgent({ id: "agent-1", cwd: REPO_CWD, @@ -1131,14 +1142,9 @@ test("agent_update emits remove when the agent has no workspaceId", async () => }), ); - session.agentUpdatesSubscription = { - subscriptionId: "sub-agents", - filter: {}, - isBootstrapping: false, - pendingUpdatesByAgentId: new Map(), - }; + activateAgentUpdatesSubscription(session, "sub-agents", {}); - await session.forwardAgentUpdate( + await session.agentUpdates.forwardLiveAgent( makeManagedAgent({ id: "agent-1", cwd: UNREGISTERED_CWD, @@ -1295,12 +1301,7 @@ test("archive emits an authoritative agent_update upsert for subscribed clients" }), ); - session.agentUpdatesSubscription = { - subscriptionId: "sub-agents", - filter: { includeArchived: true }, - isBootstrapping: false, - pendingUpdatesByAgentId: new Map(), - }; + activateAgentUpdatesSubscription(session, "sub-agents", { includeArchived: true }); await session.handleArchiveAgentRequest("agent-1", "req-archive"); @@ -1663,12 +1664,7 @@ test("close_items_request archives agents and kills terminals in one batch", asy }), ); - session.agentUpdatesSubscription = { - subscriptionId: "sub-agents", - filter: { includeArchived: true }, - isBootstrapping: false, - pendingUpdatesByAgentId: new Map(), - }; + activateAgentUpdatesSubscription(session, "sub-agents", { includeArchived: true }); await session.handleMessage({ type: "close_items_request", @@ -1850,12 +1846,7 @@ test("close_items_request archives stored agents that are not currently loaded", }), ); - session.agentUpdatesSubscription = { - subscriptionId: "sub-agents", - filter: { includeArchived: true }, - isBootstrapping: false, - pendingUpdatesByAgentId: new Map(), - }; + activateAgentUpdatesSubscription(session, "sub-agents", { includeArchived: true }); await session.handleMessage({ type: "close_items_request", @@ -2005,12 +1996,7 @@ test("close_items_request continues after an archive failure", async () => { }), ); - session.agentUpdatesSubscription = { - subscriptionId: "sub-agents", - filter: { includeArchived: true }, - isBootstrapping: false, - pendingUpdatesByAgentId: new Map(), - }; + activateAgentUpdatesSubscription(session, "sub-agents", { includeArchived: true }); await session.handleMessage({ type: "close_items_request", @@ -2393,7 +2379,7 @@ test("fetch_agent_history_request pages archived historical rows separately", as }), }, ]); - expect(session.agentUpdatesSubscription).toBeNull(); + expect(session.agentUpdates.hasSubscription()).toBe(false); }); test("fetch_agent_history_request skips rows whose workspace project record is missing", async () => { @@ -3601,7 +3587,7 @@ test("import_agent_request registers a workspace for a never-seen cwd", async () session.agentManager.setTitle = async () => undefined; session.agentStorage.list = async () => []; session.agentStorage.get = async () => null; - session.forwardAgentUpdate = async () => undefined; + session.agentUpdates.forwardLiveAgent = async () => undefined; session.workspaceUpdatesSubscription = { subscriptionId: "sub-import", @@ -4465,7 +4451,7 @@ test("refresh_agent_request unarchives the owning workspace when its directory e session.agentManager.reloadAgentSession = async () => managed; session.agentManager.hydrateTimelineFromProvider = async () => undefined; session.agentManager.getTimeline = () => []; - session.forwardAgentUpdate = async () => undefined; + session.agentUpdates.forwardLiveAgent = async () => undefined; const unarchivedWorkspaceIds: string[][] = []; const realEmit = session.emitWorkspaceUpdatesForWorkspaceIds.bind(session); @@ -4563,7 +4549,7 @@ test("refresh_agent_request leaves the owning workspace archived when its direct session.agentManager.reloadAgentSession = async () => managed; session.agentManager.hydrateTimelineFromProvider = async () => undefined; session.agentManager.getTimeline = () => []; - session.forwardAgentUpdate = async () => undefined; + session.agentUpdates.forwardLiveAgent = async () => undefined; await session.handleMessage({ type: "refresh_agent_request", @@ -4660,7 +4646,7 @@ test("refresh_agent_request recreates a deleted worktree directory and unarchive session.agentManager.reloadAgentSession = async () => managed; session.agentManager.hydrateTimelineFromProvider = async () => undefined; session.agentManager.getTimeline = () => []; - session.forwardAgentUpdate = async () => undefined; + session.agentUpdates.forwardLiveAgent = async () => undefined; const unarchivedWorkspaceIds: string[][] = []; const realEmit = session.emitWorkspaceUpdatesForWorkspaceIds.bind(session); @@ -4764,7 +4750,7 @@ test("refresh_agent_request leaves the worktree archived and surfaces a typed er session.agentManager.reloadAgentSession = async () => managed; session.agentManager.hydrateTimelineFromProvider = async () => undefined; session.agentManager.getTimeline = () => []; - session.forwardAgentUpdate = async () => undefined; + session.agentUpdates.forwardLiveAgent = async () => undefined; await session.handleMessage({ type: "refresh_agent_request", @@ -4894,7 +4880,7 @@ test("refresh_agent_request recreates a real deleted worktree against a temp git session.agentManager.reloadAgentSession = async () => managed; session.agentManager.hydrateTimelineFromProvider = async () => undefined; session.agentManager.getTimeline = () => []; - session.forwardAgentUpdate = async () => undefined; + session.agentUpdates.forwardLiveAgent = async () => undefined; await session.handleMessage({ type: "refresh_agent_request", diff --git a/packages/server/src/server/session/agent-updates/agent-updates-service.test.ts b/packages/server/src/server/session/agent-updates/agent-updates-service.test.ts new file mode 100644 index 000000000..3a5fa2e3c --- /dev/null +++ b/packages/server/src/server/session/agent-updates/agent-updates-service.test.ts @@ -0,0 +1,497 @@ +import { describe, expect, test } from "vitest"; +import type pino from "pino"; +import { createAgentUpdatesService, matchesAgentUpdatesFilter } from "./agent-updates-service.js"; +import type { + AgentSnapshotPayload, + ProjectPlacementPayload, + SessionOutboundMessage, +} from "../../messages.js"; +import type { ManagedAgent } from "../../agent/agent-manager.js"; +import type { StoredAgentRecord } from "../../agent/agent-storage.js"; + +// No mocks — every dependency is an injected in-memory fake. The agent payloads +// are supplied through the fake builders, so each test fully controls the +// (agent, project, filter) triple the service reasons about and asserts the +// emitted `agent_update` payloads. + +type AgentUpdatePayload = Extract["payload"]; + +function makeAgentPayload(input: { + id: string; + workspaceId?: string; + provider?: string; + status?: AgentSnapshotPayload["status"]; + updatedAt?: string; + archivedAt?: string | null; + labels?: Record; + requiresAttention?: boolean; + effectiveThinkingOptionId?: string | null; + thinkingOptionId?: string | null; +}): AgentSnapshotPayload { + const updatedAt = input.updatedAt ?? "2026-03-01T12:00:00.000Z"; + const provider = input.provider ?? "codex"; + return { + id: input.id, + provider, + cwd: "/tmp/repo", + ...(input.workspaceId ? { workspaceId: input.workspaceId } : {}), + model: null, + thinkingOptionId: input.thinkingOptionId ?? null, + effectiveThinkingOptionId: input.effectiveThinkingOptionId ?? null, + createdAt: updatedAt, + updatedAt, + lastUserMessageAt: null, + status: input.status ?? "running", + capabilities: { + supportsStreaming: true, + supportsSessionPersistence: true, + supportsDynamicModes: true, + supportsMcpServers: true, + supportsReasoningStream: true, + supportsToolInvocations: true, + }, + currentModeId: null, + availableModes: [], + pendingPermissions: [], + persistence: null, + runtimeInfo: { provider, sessionId: null }, + title: null, + labels: input.labels ?? {}, + requiresAttention: input.requiresAttention ?? false, + attentionReason: null, + attentionTimestamp: null, + archivedAt: input.archivedAt ?? null, + }; +} + +function makeProject(overrides?: Partial): ProjectPlacementPayload { + return { + projectKey: "proj-1", + projectName: "repo", + workspaceName: "main", + checkout: { + cwd: "/tmp/repo", + isGit: false, + currentBranch: null, + remoteUrl: null, + worktreeRoot: null, + isPaseoOwnedWorktree: false, + mainRepoRoot: null, + }, + ...overrides, + }; +} + +function buildHarness() { + const emitted: SessionOutboundMessage[] = []; + const workspaceUpdates: string[] = []; + const loggedErrors: unknown[][] = []; + const payloadById = new Map(); + const projectByWorkspaceId = new Map(); + let providerVisible: (provider: string) => boolean = () => true; + let buildAgentPayloadError: Error | null = null; + + const service = createAgentUpdatesService({ + emit: (message) => emitted.push(message), + buildAgentPayload: async (agent) => { + if (buildAgentPayloadError) { + throw buildAgentPayloadError; + } + const payload = payloadById.get(agent.id); + if (!payload) { + throw new Error(`no payload registered for ${agent.id}`); + } + return payload; + }, + buildStoredAgentPayload: (record) => { + const payload = payloadById.get(record.id); + if (!payload) { + throw new Error(`no payload registered for ${record.id}`); + } + return payload; + }, + isProviderVisibleToClient: (provider) => providerVisible(provider), + buildProjectPlacementForWorkspaceId: async (workspaceId) => + projectByWorkspaceId.get(workspaceId) ?? null, + emitWorkspaceUpdateForWorkspaceId: async (workspaceId) => { + workspaceUpdates.push(workspaceId); + }, + logger: { error: (...args: unknown[]) => loggedErrors.push(args) } as unknown as pino.Logger, + }); + + return { + service, + emitted, + workspaceUpdates, + loggedErrors, + // Register the payload a builder returns for an agent id, plus the project + // its workspaceId resolves to (null = no placement). + register( + payload: AgentSnapshotPayload, + project: ProjectPlacementPayload | null = makeProject(), + ) { + payloadById.set(payload.id, payload); + if (payload.workspaceId) { + projectByWorkspaceId.set(payload.workspaceId, project); + } + return payload; + }, + setProviderVisible(fn: (provider: string) => boolean) { + providerVisible = fn; + }, + failBuildAgentPayload(error: Error) { + buildAgentPayloadError = error; + }, + agentUpdates(): AgentUpdatePayload[] { + return emitted + .filter((message) => message.type === "agent_update") + .map( + (message) => + (message as Extract).payload, + ); + }, + managed(id: string): ManagedAgent { + return { id } as unknown as ManagedAgent; + }, + stored(id: string): StoredAgentRecord { + return { id } as unknown as StoredAgentRecord; + }, + }; +} + +describe("matchesAgentUpdatesFilter", () => { + const project = makeProject(); + + test("no filter matches", () => { + expect(matchesAgentUpdatesFilter({ agent: makeAgentPayload({ id: "a" }), project })).toBe(true); + }); + + test("label match vs mismatch", () => { + const agent = makeAgentPayload({ id: "a", labels: { surface: "voice" } }); + expect( + matchesAgentUpdatesFilter({ agent, project, filter: { labels: { surface: "voice" } } }), + ).toBe(true); + expect( + matchesAgentUpdatesFilter({ agent, project, filter: { labels: { surface: "cli" } } }), + ).toBe(false); + }); + + test("archived agents are excluded unless includeArchived", () => { + const agent = makeAgentPayload({ id: "a", archivedAt: "2026-03-02T00:00:00.000Z" }); + expect(matchesAgentUpdatesFilter({ agent, project, filter: {} })).toBe(false); + expect(matchesAgentUpdatesFilter({ agent, project, filter: { includeArchived: true } })).toBe( + true, + ); + }); + + test("thinking-option filter compares the resolved option", () => { + const high = makeAgentPayload({ id: "a", effectiveThinkingOptionId: "high" }); + expect( + matchesAgentUpdatesFilter({ agent: high, project, filter: { thinkingOptionId: "high" } }), + ).toBe(true); + expect( + matchesAgentUpdatesFilter({ agent: high, project, filter: { thinkingOptionId: "low" } }), + ).toBe(false); + // undefined means "don't filter on thinking option". + expect(matchesAgentUpdatesFilter({ agent: high, project, filter: {} })).toBe(true); + }); + + test("status filter", () => { + const agent = makeAgentPayload({ id: "a", status: "running" }); + expect(matchesAgentUpdatesFilter({ agent, project, filter: { statuses: ["running"] } })).toBe( + true, + ); + expect(matchesAgentUpdatesFilter({ agent, project, filter: { statuses: ["closed"] } })).toBe( + false, + ); + }); + + test("requiresAttention filter", () => { + const agent = makeAgentPayload({ id: "a", requiresAttention: true }); + expect(matchesAgentUpdatesFilter({ agent, project, filter: { requiresAttention: true } })).toBe( + true, + ); + expect( + matchesAgentUpdatesFilter({ agent, project, filter: { requiresAttention: false } }), + ).toBe(false); + }); + + test("projectKeys filter, ignoring blank entries", () => { + const agent = makeAgentPayload({ id: "a" }); + const inProject = makeProject({ projectKey: "proj-1" }); + const otherProject = makeProject({ projectKey: "proj-2" }); + expect( + matchesAgentUpdatesFilter({ agent, project: inProject, filter: { projectKeys: ["proj-1"] } }), + ).toBe(true); + expect( + matchesAgentUpdatesFilter({ + agent, + project: otherProject, + filter: { projectKeys: ["proj-1"] }, + }), + ).toBe(false); + // A whitespace-only key is trimmed away, leaving no constraint. + expect( + matchesAgentUpdatesFilter({ agent, project: otherProject, filter: { projectKeys: [" "] } }), + ).toBe(true); + }); +}); + +describe("forwardLiveAgent", () => { + test("emits an upsert for a matching agent and updates its workspace", async () => { + const h = buildHarness(); + h.service.beginSubscription({ subscriptionId: "sub", filter: {} }); + h.service.flushBootstrapped("sub"); + h.register(makeAgentPayload({ id: "a", workspaceId: "ws-1" })); + + await h.service.forwardLiveAgent(h.managed("a")); + + expect(h.agentUpdates()).toEqual([ + { kind: "upsert", agent: expect.objectContaining({ id: "a" }), project: makeProject() }, + ]); + expect(h.workspaceUpdates).toEqual(["ws-1"]); + }); + + test("emits a remove when the agent's workspace resolves to no project", async () => { + const h = buildHarness(); + h.service.beginSubscription({ subscriptionId: "sub", filter: {} }); + h.service.flushBootstrapped("sub"); + h.register(makeAgentPayload({ id: "a", workspaceId: "ws-1" }), null); + + await h.service.forwardLiveAgent(h.managed("a")); + + expect(h.agentUpdates()).toEqual([{ kind: "remove", agentId: "a" }]); + }); + + test("emits a remove when the agent does not match the filter", async () => { + const h = buildHarness(); + h.service.beginSubscription({ subscriptionId: "sub", filter: { statuses: ["closed"] } }); + h.service.flushBootstrapped("sub"); + h.register(makeAgentPayload({ id: "a", workspaceId: "ws-1", status: "running" })); + + await h.service.forwardLiveAgent(h.managed("a")); + + expect(h.agentUpdates()).toEqual([{ kind: "remove", agentId: "a" }]); + }); + + test("with no subscription, emits no agent_update but still updates the workspace", async () => { + const h = buildHarness(); + h.register(makeAgentPayload({ id: "a", workspaceId: "ws-1" })); + + await h.service.forwardLiveAgent(h.managed("a")); + + expect(h.agentUpdates()).toEqual([]); + expect(h.workspaceUpdates).toEqual(["ws-1"]); + }); + + test("drops an upsert whose provider is not visible to the client", async () => { + const h = buildHarness(); + h.service.beginSubscription({ subscriptionId: "sub", filter: {} }); + h.service.flushBootstrapped("sub"); + h.setProviderVisible(() => false); + h.register(makeAgentPayload({ id: "a", workspaceId: "ws-1", provider: "pi" })); + + await h.service.forwardLiveAgent(h.managed("a")); + + expect(h.agentUpdates()).toEqual([]); + // The workspace-update tail still fires regardless of visibility. + expect(h.workspaceUpdates).toEqual(["ws-1"]); + }); + + test("swallows and logs a build error without throwing", async () => { + const h = buildHarness(); + h.service.beginSubscription({ subscriptionId: "sub", filter: {} }); + h.service.flushBootstrapped("sub"); + h.failBuildAgentPayload(new Error("boom")); + + await expect(h.service.forwardLiveAgent(h.managed("a"))).resolves.toBeUndefined(); + expect(h.loggedErrors).toHaveLength(1); + expect(h.agentUpdates()).toEqual([]); + }); +}); + +describe("emitStoredRecord", () => { + test("returns the built payload and emits an upsert when matching", async () => { + const h = buildHarness(); + h.service.beginSubscription({ subscriptionId: "sub", filter: {} }); + h.service.flushBootstrapped("sub"); + h.register(makeAgentPayload({ id: "a", workspaceId: "ws-1" })); + + const payload = await h.service.emitStoredRecord(h.stored("a")); + + expect(payload.id).toBe("a"); + expect(h.agentUpdates()).toEqual([ + { kind: "upsert", agent: expect.objectContaining({ id: "a" }), project: makeProject() }, + ]); + }); + + test("emits a remove when no project resolves", async () => { + const h = buildHarness(); + h.service.beginSubscription({ subscriptionId: "sub", filter: {} }); + h.service.flushBootstrapped("sub"); + h.register(makeAgentPayload({ id: "a", workspaceId: "ws-1" }), null); + + await h.service.emitStoredRecord(h.stored("a")); + + expect(h.agentUpdates()).toEqual([{ kind: "remove", agentId: "a" }]); + }); + + test("emits a remove when the record no longer matches the filter", async () => { + const h = buildHarness(); + h.service.beginSubscription({ subscriptionId: "sub", filter: { includeArchived: false } }); + h.service.flushBootstrapped("sub"); + h.register( + makeAgentPayload({ id: "a", workspaceId: "ws-1", archivedAt: "2026-03-02T00:00:00.000Z" }), + ); + + await h.service.emitStoredRecord(h.stored("a")); + + expect(h.agentUpdates()).toEqual([{ kind: "remove", agentId: "a" }]); + }); + + test("returns the payload but emits nothing when there is no subscription", async () => { + const h = buildHarness(); + h.register(makeAgentPayload({ id: "a", workspaceId: "ws-1" })); + + const payload = await h.service.emitStoredRecord(h.stored("a")); + + expect(payload.id).toBe("a"); + expect(h.agentUpdates()).toEqual([]); + }); +}); + +describe("bootstrap buffering", () => { + test("buffers updates while bootstrapping and replays them on flush", async () => { + const h = buildHarness(); + h.service.beginSubscription({ subscriptionId: "sub", filter: {} }); + h.register(makeAgentPayload({ id: "a", workspaceId: "ws-1" })); + + await h.service.forwardLiveAgent(h.managed("a")); + expect(h.agentUpdates()).toEqual([]); // buffered, not emitted yet + + h.service.flushBootstrapped("sub"); + expect(h.agentUpdates()).toEqual([ + { kind: "upsert", agent: expect.objectContaining({ id: "a" }), project: makeProject() }, + ]); + }); + + test("keeps only the latest buffered update per agent", async () => { + const h = buildHarness(); + h.service.beginSubscription({ subscriptionId: "sub", filter: {} }); + + h.register(makeAgentPayload({ id: "a", workspaceId: "ws-1", status: "running" })); + await h.service.forwardLiveAgent(h.managed("a")); + h.register(makeAgentPayload({ id: "a", workspaceId: "ws-1", status: "closed" })); + await h.service.forwardLiveAgent(h.managed("a")); + + h.service.flushBootstrapped("sub"); + + const updates = h.agentUpdates(); + expect(updates).toHaveLength(1); + expect(updates[0]).toMatchObject({ + kind: "upsert", + agent: { id: "a", status: "closed" }, + }); + }); + + test("skips a buffered upsert that is not newer than the snapshot", async () => { + const h = buildHarness(); + h.service.beginSubscription({ subscriptionId: "sub", filter: {} }); + + h.register( + makeAgentPayload({ id: "stale", workspaceId: "ws-1", updatedAt: "2026-03-01T00:00:00.000Z" }), + ); + await h.service.forwardLiveAgent(h.managed("stale")); + h.register( + makeAgentPayload({ id: "fresh", workspaceId: "ws-1", updatedAt: "2026-03-05T00:00:00.000Z" }), + ); + await h.service.forwardLiveAgent(h.managed("fresh")); + + h.service.flushBootstrapped("sub", { + snapshotUpdatedAtByAgentId: new Map([ + ["stale", Date.parse("2026-03-02T00:00:00.000Z")], // snapshot newer → drop + ["fresh", Date.parse("2026-03-02T00:00:00.000Z")], // update newer → keep + ]), + }); + + expect(h.agentUpdates()).toEqual([ + { kind: "upsert", agent: expect.objectContaining({ id: "fresh" }), project: makeProject() }, + ]); + }); + + test("a removed agent is always replayed, even against a snapshot", async () => { + const h = buildHarness(); + h.service.beginSubscription({ subscriptionId: "sub", filter: {} }); + + h.service.removeAgent("a"); + expect(h.agentUpdates()).toEqual([]); // buffered + + h.service.flushBootstrapped("sub", { + snapshotUpdatedAtByAgentId: new Map([["a", Date.parse("2030-01-01T00:00:00.000Z")]]), + }); + expect(h.agentUpdates()).toEqual([{ kind: "remove", agentId: "a" }]); + }); + + test("does not buffer an upsert for a provider that is not visible", async () => { + const h = buildHarness(); + h.service.beginSubscription({ subscriptionId: "sub", filter: {} }); + h.setProviderVisible(() => false); + h.register(makeAgentPayload({ id: "a", workspaceId: "ws-1", provider: "pi" })); + + await h.service.forwardLiveAgent(h.managed("a")); + h.service.flushBootstrapped("sub"); + + expect(h.agentUpdates()).toEqual([]); + }); +}); + +describe("subscription lifecycle", () => { + test("flushBootstrapped is a no-op for a stale subscription id", async () => { + const h = buildHarness(); + h.service.beginSubscription({ subscriptionId: "sub", filter: {} }); + h.register(makeAgentPayload({ id: "a", workspaceId: "ws-1" })); + await h.service.forwardLiveAgent(h.managed("a")); + + h.service.flushBootstrapped("other-sub"); + expect(h.agentUpdates()).toEqual([]); // still buffering + + h.service.flushBootstrapped("sub"); + expect(h.agentUpdates()).toHaveLength(1); + }); + + test("clearSubscription only clears the current subscription", () => { + const h = buildHarness(); + h.service.beginSubscription({ subscriptionId: "sub", filter: {} }); + + h.service.clearSubscription("other-sub"); + expect(h.service.hasSubscription()).toBe(true); + + h.service.clearSubscription("sub"); + expect(h.service.hasSubscription()).toBe(false); + }); + + test("dispose drops the subscription", () => { + const h = buildHarness(); + h.service.beginSubscription({ subscriptionId: "sub", filter: {} }); + expect(h.service.hasSubscription()).toBe(true); + + h.service.dispose(); + expect(h.service.hasSubscription()).toBe(false); + }); + + test("removeAgent is a no-op without a subscription", () => { + const h = buildHarness(); + h.service.removeAgent("a"); + expect(h.agentUpdates()).toEqual([]); + }); + + test("removeAgent emits a remove for a live subscription", () => { + const h = buildHarness(); + h.service.beginSubscription({ subscriptionId: "sub", filter: {} }); + h.service.flushBootstrapped("sub"); + + h.service.removeAgent("a"); + + expect(h.agentUpdates()).toEqual([{ kind: "remove", agentId: "a" }]); + }); +}); diff --git a/packages/server/src/server/session/agent-updates/agent-updates-service.ts b/packages/server/src/server/session/agent-updates/agent-updates-service.ts new file mode 100644 index 000000000..4630823cc --- /dev/null +++ b/packages/server/src/server/session/agent-updates/agent-updates-service.ts @@ -0,0 +1,326 @@ +import type pino from "pino"; +import type { + AgentSnapshotPayload, + ProjectPlacementPayload, + SessionInboundMessage, + SessionOutboundMessage, +} from "../../messages.js"; +import type { ManagedAgent } from "../../agent/agent-manager.js"; +import type { StoredAgentRecord } from "../../agent/agent-storage.js"; +import { resolveEffectiveThinkingOptionId } from "../../agent/agent-projections.js"; + +type AgentUpdatePayload = Extract["payload"]; +type AgentUpdatesFilter = NonNullable< + Extract["filter"] +>; + +interface AgentUpdatesSubscriptionState { + subscriptionId: string; + filter?: AgentUpdatesFilter; + isBootstrapping: boolean; + pendingUpdatesByAgentId: Map; +} + +/** + * Owns the single per-client `agent_update` subscription: when a client subscribes + * via `fetch_agents_request`, every later agent lifecycle change (live forward, + * stored-record archive/detach, delete) is filtered against the subscription's + * filter and either emitted or — while the initial snapshot is still being built — + * buffered and replayed on flush. Keeping the mutable subscription state, the + * bootstrap buffer, the provider-visibility gate, and the filter predicate behind + * one interface stops the rest of session.ts from poking the subscription shape or + * hand-rolling `agent_update` payloads, and the (previously untested) filter/buffer/ + * flush branches become exercisable through injected fakes. + * + * The snapshot listing path applies the SAME filter via the pure + * `matchesAgentUpdatesFilter` so a subscription's initial page and its live updates + * stay consistent. + */ +export interface AgentUpdatesService { + beginSubscription(input: { subscriptionId: string; filter?: AgentUpdatesFilter }): void; + flushBootstrapped( + subscriptionId: string, + options?: { snapshotUpdatedAtByAgentId?: Map }, + ): void; + clearSubscription(subscriptionId: string): void; + hasSubscription(): boolean; + forwardLiveAgent(agent: ManagedAgent): Promise; + emitStoredRecord(record: StoredAgentRecord): Promise; + removeAgent(agentId: string): void; + dispose(): void; +} + +export interface AgentUpdatesServiceDeps { + emit(message: SessionOutboundMessage): void; + buildAgentPayload(agent: ManagedAgent): Promise; + buildStoredAgentPayload(record: StoredAgentRecord): AgentSnapshotPayload; + isProviderVisibleToClient(provider: string): boolean; + buildProjectPlacementForWorkspaceId(workspaceId: string): Promise; + emitWorkspaceUpdateForWorkspaceId(workspaceId: string): Promise; + logger: pino.Logger; +} + +function agentThinkingOptionMatchesFilter( + agent: AgentSnapshotPayload, + filter: AgentUpdatesFilter, +): boolean { + if (filter.thinkingOptionId === undefined) { + return true; + } + const expectedThinkingOptionId = resolveEffectiveThinkingOptionId({ + configuredThinkingOptionId: filter.thinkingOptionId ?? null, + }); + const resolvedThinkingOptionId = + agent.effectiveThinkingOptionId ?? + resolveEffectiveThinkingOptionId({ + runtimeInfo: agent.runtimeInfo, + configuredThinkingOptionId: agent.thinkingOptionId ?? null, + }); + return resolvedThinkingOptionId === expectedThinkingOptionId; +} + +function matchesAgentStructuralFilter( + agent: AgentSnapshotPayload, + project: ProjectPlacementPayload, + filter: AgentUpdatesFilter, +): boolean { + if (filter.statuses && filter.statuses.length > 0) { + const statuses = new Set(filter.statuses); + if (!statuses.has(agent.status)) { + return false; + } + } + + if (typeof filter.requiresAttention === "boolean") { + const requiresAttention = agent.requiresAttention ?? false; + if (requiresAttention !== filter.requiresAttention) { + return false; + } + } + + if (filter.projectKeys && filter.projectKeys.length > 0) { + const projectKeys = new Set(filter.projectKeys.filter((item) => item.trim().length > 0)); + if (projectKeys.size > 0 && !projectKeys.has(project.projectKey)) { + return false; + } + } + return true; +} + +/** + * Pure predicate shared by the live subscription stream and the snapshot listing + * pager: does an agent (with its resolved project placement) satisfy a + * `fetch_agents` filter? + */ +export function matchesAgentUpdatesFilter(input: { + agent: AgentSnapshotPayload; + project: ProjectPlacementPayload; + filter?: AgentUpdatesFilter; +}): boolean { + const { agent, project, filter } = input; + + if (filter?.labels) { + const matchesLabels = Object.entries(filter.labels).every( + ([key, value]) => agent.labels[key] === value, + ); + if (!matchesLabels) { + return false; + } + } + + const includeArchived = filter?.includeArchived ?? false; + if (!includeArchived && agent.archivedAt) { + return false; + } + + if (filter && !agentThinkingOptionMatchesFilter(agent, filter)) { + return false; + } + + if (filter && !matchesAgentStructuralFilter(agent, project, filter)) { + return false; + } + + return true; +} + +function agentUpdateTargetId(update: AgentUpdatePayload): string { + return update.kind === "remove" ? update.agentId : update.agent.id; +} + +export function createAgentUpdatesService(deps: AgentUpdatesServiceDeps): AgentUpdatesService { + let subscription: AgentUpdatesSubscriptionState | null = null; + + function bufferOrEmit(sub: AgentUpdatesSubscriptionState, payload: AgentUpdatePayload): void { + if (payload.kind === "upsert" && !deps.isProviderVisibleToClient(payload.agent.provider)) { + return; + } + if (sub.isBootstrapping) { + sub.pendingUpdatesByAgentId.set(agentUpdateTargetId(payload), payload); + return; + } + + deps.emit({ + type: "agent_update", + payload, + }); + } + + function beginSubscription(input: { subscriptionId: string; filter?: AgentUpdatesFilter }): void { + subscription = { + subscriptionId: input.subscriptionId, + filter: input.filter, + isBootstrapping: true, + pendingUpdatesByAgentId: new Map(), + }; + } + + function flushBootstrapped( + subscriptionId: string, + options?: { snapshotUpdatedAtByAgentId?: Map }, + ): void { + if (!subscription || subscription.subscriptionId !== subscriptionId) { + return; + } + if (!subscription.isBootstrapping) { + return; + } + + subscription.isBootstrapping = false; + const pending = Array.from(subscription.pendingUpdatesByAgentId.values()); + subscription.pendingUpdatesByAgentId.clear(); + + for (const payload of pending) { + if (payload.kind === "upsert") { + const snapshotUpdatedAt = options?.snapshotUpdatedAtByAgentId?.get(payload.agent.id); + if (typeof snapshotUpdatedAt === "number") { + const updateUpdatedAt = Date.parse(payload.agent.updatedAt); + if (!Number.isNaN(updateUpdatedAt) && updateUpdatedAt <= snapshotUpdatedAt) { + continue; + } + } + } + + deps.emit({ + type: "agent_update", + payload, + }); + } + } + + function clearSubscription(subscriptionId: string): void { + if (subscription && subscription.subscriptionId === subscriptionId) { + subscription = null; + } + } + + function hasSubscription(): boolean { + return subscription !== null; + } + + function removeAgent(agentId: string): void { + if (!subscription) { + return; + } + bufferOrEmit(subscription, { kind: "remove", agentId }); + } + + async function emitStoredRecord(record: StoredAgentRecord): Promise { + const payload = deps.buildStoredAgentPayload(record); + const sub = subscription; + if (!sub) { + return payload; + } + + const project = payload.workspaceId + ? await deps.buildProjectPlacementForWorkspaceId(payload.workspaceId) + : null; + if (!project) { + bufferOrEmit(sub, { + kind: "remove", + agentId: payload.id, + }); + return payload; + } + + const matches = matchesAgentUpdatesFilter({ + agent: payload, + project, + filter: sub.filter, + }); + bufferOrEmit( + sub, + matches + ? { + kind: "upsert", + agent: payload, + project, + } + : { + kind: "remove", + agentId: payload.id, + }, + ); + return payload; + } + + async function forwardLiveAgent(agent: ManagedAgent): Promise { + try { + const sub = subscription; + const payload = await deps.buildAgentPayload(agent); + if (sub) { + const project = payload.workspaceId + ? await deps.buildProjectPlacementForWorkspaceId(payload.workspaceId) + : null; + if (!project) { + bufferOrEmit(sub, { + kind: "remove", + agentId: payload.id, + }); + } else { + const matches = matchesAgentUpdatesFilter({ + agent: payload, + project, + filter: sub.filter, + }); + + if (matches) { + bufferOrEmit(sub, { + kind: "upsert", + agent: payload, + project, + }); + } else { + bufferOrEmit(sub, { + kind: "remove", + agentId: payload.id, + }); + } + } + } + + // A lifecycle change updates exactly the agent's owning workspace, never + // every workspace sharing its cwd. Ownership is the agent's workspaceId. + if (payload.workspaceId) { + await deps.emitWorkspaceUpdateForWorkspaceId(payload.workspaceId); + } + } catch (error) { + deps.logger.error({ err: error }, "Failed to emit agent update"); + } + } + + function dispose(): void { + subscription = null; + } + + return { + beginSubscription, + flushBootstrapped, + clearSubscription, + hasSubscription, + forwardLiveAgent, + emitStoredRecord, + removeAgent, + dispose, + }; +}