From 9565fed532a89b3e95638ccd6fc2a974d2d28c25 Mon Sep 17 00:00:00 2001
From: Justin Chu
Date: Thu, 11 Jun 2026 16:19:10 -0700
Subject: [PATCH 1/6] fix: director runtime inheritance + copilot streaming
(chat silence, word dupes)
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
Director silently dead (lead → director messages dropped):
- An unconfigured director role fell back to the global default runtime
(codex), which may not be installed. Director now inherits the Lead's
runtime/model unless explicitly configured (new ModelConfig.hasRoleConfig,
applied in both gateway registration paths and re-read in spawnDirector).
- AcpAdapter.spawn now fails fast when the runtime binary is missing —
previously the ENOENT arrived async after spawn() "succeeded", leaving a
dead session that dropped every steer.
- Director spawn failures are surfaced to the user's chat as a system
message (LeadManager.onSystemNotice, wired to persist + ws broadcast in
wireWsToLead) instead of being swallowed by empty catch blocks.
- Lead/Director spawns now actually pass the configured model to the
adapter (previously model was only persisted for display).
Copilot SDK streaming:
- Main chat showed no streaming: copilot sessions never invoked the
ACP-style onOutputChunk that LeadManager wires for chat:stream. The
adapter now synthesizes SessionUpdate chunks (message/thought/tool_call)
and emits them on the session, so the existing lead pipeline works.
- Agents page repeated words: the SDK emits both incremental
assistant.message_delta events AND a full assistant.message at turn end;
both were broadcast as deltas. The full-content event is now dropped
when deltas were streamed this turn (tracked per session, reset at
session.idle), while non-streaming turns still pass it through.
Co-Authored-By: Claude Fable 5
---
packages/server/src/agents/AcpAdapter.ts | 14 ++++
.../server/src/agents/CopilotSdkAdapter.ts | 74 ++++++++++++++++++-
packages/server/src/agents/ModelConfig.ts | 10 +++
packages/server/src/cli/gateway.ts | 31 +++++++-
packages/server/src/lead/LeadManager.ts | 74 +++++++++++++++++--
packages/server/tests/agents/adapters.test.ts | 11 +++
.../tests/agents/copilot-sdk-adapter.test.ts | 56 ++++++++++++++
.../server/tests/agents/model-config.test.ts | 12 +++
.../server/tests/lead/lead-manager.test.ts | 51 +++++++++++++
9 files changed, 321 insertions(+), 12 deletions(-)
diff --git a/packages/server/src/agents/AcpAdapter.ts b/packages/server/src/agents/AcpAdapter.ts
index c4e33114..1a577966 100644
--- a/packages/server/src/agents/AcpAdapter.ts
+++ b/packages/server/src/agents/AcpAdapter.ts
@@ -421,6 +421,20 @@ export class AcpAdapter extends AgentAdapter {
throw new Error(`Unknown runtime "${runtimeName}". Available: ${Object.keys(this.runtimes).join(', ')}`);
}
+ // Fail fast when the runtime binary is missing. Without this check the
+ // ENOENT arrives async after spawn() already "succeeded", leaving a dead
+ // session that silently drops every message.
+ const resolvedCmd = resolveCommand(runtime.command);
+ if (resolvedCmd === runtime.command) {
+ const { commandExists } = await import('../utils/platform.js');
+ if (!commandExists(runtime.command)) {
+ throw new Error(
+ `Runtime "${runtimeName}" is not available: command "${runtime.command}" not found on PATH. ` +
+ `Install it, or configure a different runtime/model for this role.`
+ );
+ }
+ }
+
const prompt = opts.systemPrompt ?? `You are a ${opts.role} agent. Complete your assigned tasks.`;
const args = interpolateArgs(runtime.args, { prompt, cwd: opts.cwd });
const sessionLocalId = `acp-${randomUUID().slice(0, 8)}`;
diff --git a/packages/server/src/agents/CopilotSdkAdapter.ts b/packages/server/src/agents/CopilotSdkAdapter.ts
index 44edece8..a7da54f9 100644
--- a/packages/server/src/agents/CopilotSdkAdapter.ts
+++ b/packages/server/src/agents/CopilotSdkAdapter.ts
@@ -50,6 +50,12 @@ export interface CopilotAgentSession {
cwd: string;
model?: string;
pendingToolCalls?: Map;
+ /** ACP-style per-chunk callback (set by LeadManager for chat streaming). */
+ onOutputChunk?: (update: unknown) => void;
+ /** True once a text/reasoning delta was streamed this turn — used to drop
+ * the redundant full-content event that would duplicate streamed words. */
+ sawTextDelta?: boolean;
+ sawReasoningDelta?: boolean;
}
export class CopilotSdkAdapter extends AgentAdapter {
@@ -988,6 +994,64 @@ export class CopilotSdkAdapter extends AgentAdapter {
/**
* Spawn a new Copilot agent session with flightdeck tools injected.
*/
+ /**
+ * Normalize a Copilot SDK event into an ACP-style SessionUpdate for the
+ * session's onOutputChunk (lead chat streaming), and decide whether the raw
+ * event should still be forwarded to onOutput (agent:stream broadcast).
+ *
+ * Returns false for full-content events whose text was already streamed as
+ * deltas this turn — forwarding those would repeat every word in the UI.
+ */
+ private processStreamEvent(agentSession: CopilotAgentSession, event: { type: string; data?: any }): boolean {
+ let update: Record | null = null;
+ let forward = true;
+ switch (event.type) {
+ case 'assistant.message_delta':
+ agentSession.sawTextDelta = true;
+ update = { sessionUpdate: 'agent_message_chunk', content: { type: 'text', text: event.data?.deltaContent ?? '' } };
+ break;
+ case 'assistant.reasoning_delta':
+ agentSession.sawReasoningDelta = true;
+ update = { sessionUpdate: 'agent_thought_chunk', content: { type: 'text', text: event.data?.deltaContent ?? '' } };
+ break;
+ case 'assistant.message':
+ if (agentSession.sawTextDelta) forward = false;
+ else if (event.data?.content) update = { sessionUpdate: 'agent_message_chunk', content: { type: 'text', text: event.data.content } };
+ break;
+ case 'assistant.reasoning':
+ if (agentSession.sawReasoningDelta) forward = false;
+ else if (event.data?.content) update = { sessionUpdate: 'agent_thought_chunk', content: { type: 'text', text: event.data.content } };
+ break;
+ case 'tool.execution_start':
+ update = {
+ sessionUpdate: 'tool_call',
+ toolCallId: event.data?.toolCallId ?? '',
+ title: event.data?.name ?? event.data?.toolName ?? '',
+ rawInput: event.data?.arguments,
+ status: 'pending',
+ };
+ break;
+ case 'tool.execution_complete':
+ update = {
+ sessionUpdate: 'tool_call_update',
+ toolCallId: event.data?.toolCallId ?? '',
+ title: event.data?.name ?? event.data?.toolName ?? '',
+ content: event.data?.content ? [{ type: 'text', text: String(event.data.content) }] : [],
+ status: 'completed',
+ };
+ break;
+ case 'session.idle':
+ // Turn boundary — next turn may or may not stream deltas
+ agentSession.sawTextDelta = false;
+ agentSession.sawReasoningDelta = false;
+ break;
+ }
+ if (update && agentSession.onOutputChunk) {
+ try { agentSession.onOutputChunk(update); } catch { /* consumer errors must not kill the event stream */ }
+ }
+ return forward;
+ }
+
async spawn(opts: BaseSpawnOptions): Promise {
const client = await this.ensureClient();
const aid = (opts.agentId ?? makeAgentId(opts.role, Date.now().toString())) as AgentId;
@@ -1125,7 +1189,8 @@ export class CopilotSdkAdapter extends AgentAdapter {
}
}
- if (this.onOutput) {
+ const forward = this.processStreamEvent(agentSession, event as { type: string; data?: any });
+ if (forward && this.onOutput) {
try { this.onOutput(aid, event); } catch { /* */ }
}
});
@@ -1294,8 +1359,6 @@ export class CopilotSdkAdapter extends AgentAdapter {
});
} catch { /* */ }
}
- if (this.onOutput) { try { this.onOutput(aid, event); } catch { /* */ } }
-
// Track tool calls (same as create path)
if ((event.type as string) === 'tool.execution_start') {
const data = (event as any).data;
@@ -1315,6 +1378,11 @@ export class CopilotSdkAdapter extends AgentAdapter {
try { this.onToolCall(aid, { toolName, status: 'completed' }); } catch { /* */ }
}
}
+
+ const forward = this.processStreamEvent(agentSession, event as { type: string; data?: any });
+ if (forward && this.onOutput) {
+ try { this.onOutput(aid, event); } catch { /* */ }
+ }
});
return { agentId: aid, sessionId: sessionId, status: 'running' as const };
diff --git a/packages/server/src/agents/ModelConfig.ts b/packages/server/src/agents/ModelConfig.ts
index 454ca693..3bff3667 100644
--- a/packages/server/src/agents/ModelConfig.ts
+++ b/packages/server/src/agents/ModelConfig.ts
@@ -103,6 +103,16 @@ export class ModelConfig {
});
}
+ /**
+ * Whether the user explicitly configured this role (runtime, model, or
+ * model pool). When false, callers may apply their own fallback — e.g.
+ * Director inherits the Lead's runtime/model instead of a global default.
+ */
+ hasRoleConfig(role: string): boolean {
+ const rc = this.getAgentsConfig().roles?.[role];
+ return !!rc && !!(rc.runtime || rc.model || (rc.enabledModels && rc.enabledModels.length > 0));
+ }
+
/**
* Get config for a single role.
*/
diff --git a/packages/server/src/cli/gateway.ts b/packages/server/src/cli/gateway.ts
index 07814b15..da12151f 100644
--- a/packages/server/src/cli/gateway.ts
+++ b/packages/server/src/cli/gateway.ts
@@ -448,7 +448,11 @@ export async function startGateway(deps: GatewayDeps): Promise {
const { ModelConfig } = await import('../agents/ModelConfig.js');
const modelConfig = new ModelConfig(projectCwd);
const leadRoleConfig = modelConfig.getRoleConfig('lead');
- const directorRoleConfig = modelConfig.getRoleConfig('director');
+ // Director inherits the Lead's runtime/model unless explicitly configured —
+ // the global default (e.g. codex) may not even be installed
+ const directorRoleConfig = modelConfig.hasRoleConfig('director')
+ ? modelConfig.getRoleConfig('director')
+ : leadRoleConfig;
// Create LeadManager
const projectConfig = fd.project.getConfig();
@@ -460,7 +464,9 @@ export async function startGateway(deps: GatewayDeps): Promise {
projectName: name,
cwd: projectCwd,
leadRuntime: leadRoleConfig.runtime as import('@flightdeck-ai/shared').AgentRuntime,
+ leadModel: leadRoleConfig.model || undefined,
directorRuntime: directorRoleConfig.runtime as import('@flightdeck-ai/shared').AgentRuntime,
+ directorModel: directorRoleConfig.model || undefined,
heartbeat: {
enabled: projectConfig.heartbeatEnabled === true,
interval: 30 * 60 * 1000,
@@ -599,7 +605,10 @@ export async function startGateway(deps: GatewayDeps): Promise {
const { ModelConfig } = await import('../agents/ModelConfig.js');
const modelConfig = new ModelConfig(projectCwd);
const leadRoleConfig = modelConfig.getRoleConfig('lead');
- const directorRoleConfig = modelConfig.getRoleConfig('director');
+ // Director inherits the Lead's runtime/model unless explicitly configured
+ const directorRoleConfig = modelConfig.hasRoleConfig('director')
+ ? modelConfig.getRoleConfig('director')
+ : leadRoleConfig;
// LeadManager
const projectConfig = fd.project.getConfig();
@@ -611,7 +620,9 @@ export async function startGateway(deps: GatewayDeps): Promise {
projectName: name,
cwd: projectCwd,
leadRuntime: leadRoleConfig.runtime as import('@flightdeck-ai/shared').AgentRuntime,
+ leadModel: leadRoleConfig.model || undefined,
directorRuntime: directorRoleConfig.runtime as import('@flightdeck-ai/shared').AgentRuntime,
+ directorModel: directorRoleConfig.model || undefined,
heartbeat: {
enabled: projectConfig.heartbeatEnabled === true,
interval: 30 * 60 * 1000,
@@ -1003,10 +1014,24 @@ async function spawnAgents(
}
// eslint-disable-next-line @typescript-eslint/no-explicit-any
-function wireWsToLead(wsServer: any, leadManager: { steerLead(event: any): Promise; getLastMergedSourceIds?(): string[]; setStreamHandler?(handler: (update: any) => void): void; cancelLead?(): Promise }, fd: Flightdeck, projectName: string, notifier?: InstanceType | null): void {
+function wireWsToLead(wsServer: any, leadManager: { steerLead(event: any): Promise; getLastMergedSourceIds?(): string[]; setStreamHandler?(handler: (update: any) => void): void; cancelLead?(): Promise; onSystemNotice?: ((content: string) => void) | null }, fd: Flightdeck, projectName: string, notifier?: InstanceType | null): void {
// Pre-generate message ID for stream↔final message consistency
const msgIdRef = { current: makeMessageId('lead', Date.now().toString()) };
+ // Surface operational notices (e.g. Director spawn failures) in the chat:
+ // persist so they survive reloads, broadcast so the user sees them live
+ if ('onSystemNotice' in leadManager) {
+ leadManager.onSystemNotice = (content: string) => {
+ try {
+ const sysMsg = fd.messages?.createMessage({
+ parentId: null, taskId: null, authorType: 'system', authorId: null,
+ content, metadata: null,
+ });
+ if (sysMsg) wsServer.broadcast({ type: 'chat:message', project: projectName, message: sysMsg });
+ } catch { /* best effort */ }
+ };
+ }
+
// Wire streaming updates (tool calls, thoughts) from Lead to WebSocket
if (leadManager.setStreamHandler && wsServer.streamChunk) {
leadManager.setStreamHandler((update: SessionUpdate) => {
diff --git a/packages/server/src/lead/LeadManager.ts b/packages/server/src/lead/LeadManager.ts
index 13198a4d..8688869e 100644
--- a/packages/server/src/lead/LeadManager.ts
+++ b/packages/server/src/lead/LeadManager.ts
@@ -101,8 +101,12 @@ export interface LeadManagerOptions {
cwd?: string;
/** Runtime name for Lead (e.g. 'copilot', 'opencode'). Falls back to adapter default. */
leadRuntime?: AgentRuntime;
+ /** Concrete model ID for Lead. Falls back to the runtime's default model. */
+ leadModel?: string;
/** Runtime name for Director. Falls back to leadRuntime, then adapter default. */
directorRuntime?: AgentRuntime;
+ /** Concrete model ID for Director. Falls back to leadModel semantics in the gateway. */
+ directorModel?: string;
}
export class LeadManager {
@@ -122,9 +126,13 @@ export class LeadManager {
private projectName: string | undefined;
private agentCwd: string;
private leadRuntime: AgentRuntime | undefined;
+ private leadModel: string | undefined;
private directorRuntime: AgentRuntime | undefined;
+ private directorModel: string | undefined;
/** Optional callback invoked during heartbeat when scout should run */
public onScoutHeartbeat: (() => Promise) | null = null;
+ /** Optional callback to surface operational notices to the user's chat (wired by the gateway). */
+ public onSystemNotice: ((content: string) => void) | null = null;
private directorSessionId: string | null = null;
private directorAgentId: string | null = null;
@@ -144,7 +152,22 @@ export class LeadManager {
this.agentCwd = opts.cwd ?? process.cwd();
this.sessionStore = new SessionStore(opts.projectName ?? 'default', opts.sqlite.db);
this.leadRuntime = opts.leadRuntime;
+ this.leadModel = opts.leadModel;
this.directorRuntime = opts.directorRuntime;
+ this.directorModel = opts.directorModel;
+ }
+
+ /** Surface an operational problem to the user: chat notice if wired, else persisted system message. */
+ private notifySystem(content: string): void {
+ if (this.onSystemNotice) {
+ try { this.onSystemNotice(content); return; } catch { /* fall through to store */ }
+ }
+ try {
+ this.messageStore?.createMessage({
+ parentId: null, taskId: null, authorType: 'system', authorId: null,
+ content, metadata: null,
+ });
+ } catch { /* best effort — never break the caller */ }
}
/** Spawn a new Lead agent session */
@@ -190,7 +213,9 @@ export class LeadManager {
console.error(` Lead ${lead.id} woken (session: ${meta.sessionId})`);
// Still spawn director alongside
if (!this.directorSessionId) {
- try { await this.spawnDirector(); } catch { /* non-fatal */ }
+ try { await this.spawnDirector(); } catch (err) {
+ this.reportDirectorSpawnFailure(err);
+ }
}
this.retireOtherAgents('lead', lead.id);
return meta.sessionId;
@@ -201,7 +226,7 @@ export class LeadManager {
}
}
- // Re-read model config to pick up runtime changes (e.g. user switched from copilot to claude)
+ // Re-read model config to pick up runtime/model changes (e.g. user switched from copilot to claude)
try {
const { ModelConfig } = await import('../agents/ModelConfig.js');
const mc = new ModelConfig(this.agentCwd);
@@ -210,6 +235,7 @@ export class LeadManager {
console.error(` Lead runtime changed: ${this.leadRuntime} → ${leadConfig.runtime}`);
this.leadRuntime = leadConfig.runtime as AgentRuntime;
}
+ if (leadConfig.model) this.leadModel = leadConfig.model;
} catch { /* fallback to existing runtime */ }
// Purge stale offline agents before spawning
@@ -274,6 +300,7 @@ export class LeadManager {
cwd: this.agentCwd,
projectName: this.projectName,
runtime: this.leadRuntime,
+ ...(this.leadModel ? { model: this.leadModel } : {}),
...(systemPrompt ? { systemPrompt } : {}),
});
this.leadSessionId = meta.sessionId;
@@ -305,7 +332,11 @@ export class LeadManager {
try {
await this.spawnDirector();
log('Lead', 'Director auto-spawned alongside Lead');
- } catch { /* Director spawn failure is non-fatal */ }
+ } catch (err) {
+ // Non-fatal for the Lead, but the user must hear about it —
+ // a silently missing Director looks like "Director ignores messages"
+ this.reportDirectorSpawnFailure(err);
+ }
}
// Retire all other leads (one project = one active lead)
@@ -672,6 +703,17 @@ export class LeadManager {
return count;
}
+ /** Log a Director spawn failure and surface it to the user's chat. */
+ private reportDirectorSpawnFailure(err: unknown): void {
+ const reason = err instanceof Error ? err.message : String(err);
+ log('Director', `Spawn FAILED (runtime: ${this.directorRuntime}): ${reason}`);
+ this.notifySystem(
+ `⚠️ Director failed to start (runtime: ${this.directorRuntime ?? 'default'}): ${reason}\n` +
+ `Messages to the Director will be dropped until this is fixed. ` +
+ `Configure the director role's runtime/model in Settings → Roles, or install the missing runtime.`
+ );
+ }
+
/** Spawn Director as a persistent ACP session */
async spawnDirector(): Promise {
// Check project-level runtime restrictions
@@ -685,7 +727,19 @@ export class LeadManager {
if (e instanceof Error && e.message.includes('not allowed')) throw e;
/* project.getConfig() may not exist in tests — skip check */
}
- log('Director', `Spawning (runtime: ${this.directorRuntime})...`);
+ // Re-read model config — Director inherits the Lead's runtime/model
+ // unless the user explicitly configured the director role
+ try {
+ const { ModelConfig } = await import('../agents/ModelConfig.js');
+ const mc = new ModelConfig(this.agentCwd);
+ const cfg = mc.hasRoleConfig('director') ? mc.getRoleConfig('director') : mc.getRoleConfig('lead');
+ if (cfg.runtime) this.directorRuntime = cfg.runtime as AgentRuntime;
+ if (cfg.model) this.directorModel = cfg.model;
+ } catch { /* keep constructor-provided values */ }
+ if (!this.directorRuntime && this.leadRuntime) this.directorRuntime = this.leadRuntime;
+ if (!this.directorModel && this.leadModel) this.directorModel = this.leadModel;
+
+ log('Director', `Spawning (runtime: ${this.directorRuntime}, model: ${this.directorModel ?? 'default'})...`);
// Try to wake a hibernated director first
const hibernatedDirectors = this.sqlite.listAgents().filter(a => a.role === 'director' && a.status === 'hibernated' && a.acpSessionId);
if (hibernatedDirectors.length > 0) {
@@ -763,6 +817,7 @@ export class LeadManager {
cwd: this.agentCwd,
projectName: this.projectName,
runtime: this.directorRuntime,
+ ...(this.directorModel ? { model: this.directorModel } : {}),
...(fullSystemPrompt ? { systemPrompt: fullSystemPrompt } : {}),
});
this.directorSessionId = meta.sessionId;
@@ -774,13 +829,17 @@ export class LeadManager {
id: meta.agentId,
role: 'director',
runtime: this.directorRuntime ?? this.leadRuntime ?? 'acp',
- runtimeName: this.directorRuntime ?? this.leadRuntime ?? 'codex',
+ runtimeName: this.directorRuntime ?? this.leadRuntime ?? null,
acpSessionId: meta.sessionId,
status: 'idle',
currentSpecId: null,
costAccumulated: 0,
lastHeartbeat: null,
});
+ const directorDisplayModel = this.directorModel ?? meta.model;
+ if (directorDisplayModel) {
+ try { this.sqlite.updateAgentModel(meta.agentId as any, directorDisplayModel); } catch { /* display only */ }
+ }
// Notify Lead about new Director (only if replacing an old one, not on first boot)
if (this.leadSessionId && this.directorAgentId) {
@@ -923,7 +982,10 @@ export class LeadManager {
return '';
}
}
- if (!this.directorSessionId) return '';
+ if (!this.directorSessionId) {
+ log('Director', 'steer dropped — no Director session (spawn failed or not started)');
+ return '';
+ }
const response = await this.acpAdapter.steer(this.directorSessionId, { content: message });
log('Director', `→ response (${Date.now() - directorStart}ms): "${truncate(response)}"`);
return response;
diff --git a/packages/server/tests/agents/adapters.test.ts b/packages/server/tests/agents/adapters.test.ts
index c182b053..2efe0c3e 100644
--- a/packages/server/tests/agents/adapters.test.ts
+++ b/packages/server/tests/agents/adapters.test.ts
@@ -39,6 +39,17 @@ describe('AcpAdapter', () => {
expect(meta).toBeNull();
});
+ it('spawn fails fast when the runtime binary is missing', async () => {
+ // Regression: the ENOENT used to arrive async after spawn() "succeeded",
+ // leaving a dead session that silently dropped every message.
+ const runtimes: Record = {
+ ghost: { command: 'definitely-not-a-real-binary-xyz', args: [], adapter: 'acp' },
+ };
+ adapter = new AcpAdapter(runtimes, 'ghost');
+ await expect(adapter.spawn({ role: 'director', cwd: '/tmp' }))
+ .rejects.toThrow(/not found on PATH/);
+ });
+
it('kill terminates the session', async () => {
const runtimes: Record = {
codex: { command: 'sleep', args: ['30'], adapter: 'acp' },
diff --git a/packages/server/tests/agents/copilot-sdk-adapter.test.ts b/packages/server/tests/agents/copilot-sdk-adapter.test.ts
index 95952876..ee5548f1 100644
--- a/packages/server/tests/agents/copilot-sdk-adapter.test.ts
+++ b/packages/server/tests/agents/copilot-sdk-adapter.test.ts
@@ -237,4 +237,60 @@ describe('CopilotSdkAdapter', () => {
await adapter.shutdown();
});
});
+
+ // ─── Stream event processing (dedup + ACP-style chunk synthesis) ─────
+
+ describe('processStreamEvent', () => {
+ function makeSession() {
+ const chunks: any[] = [];
+ const session: any = { onOutputChunk: (u: any) => chunks.push(u) };
+ return { session, chunks };
+ }
+ const process = (session: any, event: any) =>
+ (adapter as any).processStreamEvent(session, event);
+
+ it('drops the full assistant.message after deltas were streamed (no word duplication)', () => {
+ const { session, chunks } = makeSession();
+ expect(process(session, { type: 'assistant.message_delta', data: { deltaContent: 'hello ' } })).toBe(true);
+ expect(process(session, { type: 'assistant.message_delta', data: { deltaContent: 'world' } })).toBe(true);
+ // Full message repeats the streamed text — must NOT be forwarded
+ expect(process(session, { type: 'assistant.message', data: { content: 'hello world' } })).toBe(false);
+ const texts = chunks.filter(c => c.sessionUpdate === 'agent_message_chunk').map(c => c.content.text);
+ expect(texts).toEqual(['hello ', 'world']);
+ });
+
+ it('forwards the full assistant.message when nothing was streamed', () => {
+ const { session, chunks } = makeSession();
+ expect(process(session, { type: 'assistant.message', data: { content: 'hello world' } })).toBe(true);
+ expect(chunks).toEqual([{ sessionUpdate: 'agent_message_chunk', content: { type: 'text', text: 'hello world' } }]);
+ });
+
+ it('resets delta tracking at turn boundary (session.idle)', () => {
+ const { session } = makeSession();
+ process(session, { type: 'assistant.message_delta', data: { deltaContent: 'a' } });
+ process(session, { type: 'session.idle' });
+ // Next turn answers without streaming — full message must pass through
+ expect(process(session, { type: 'assistant.message', data: { content: 'next turn' } })).toBe(true);
+ });
+
+ it('dedupes reasoning the same way as text', () => {
+ const { session } = makeSession();
+ process(session, { type: 'assistant.reasoning_delta', data: { deltaContent: 'thinking…' } });
+ expect(process(session, { type: 'assistant.reasoning', data: { content: 'thinking…' } })).toBe(false);
+ });
+
+ it('synthesizes ACP-style tool call updates for onOutputChunk', () => {
+ const { session, chunks } = makeSession();
+ process(session, { type: 'tool.execution_start', data: { toolCallId: 't1', name: 'grep', arguments: { q: 'x' } } });
+ process(session, { type: 'tool.execution_complete', data: { toolCallId: 't1', name: 'grep', content: 'match' } });
+ expect(chunks[0]).toMatchObject({ sessionUpdate: 'tool_call', toolCallId: 't1', title: 'grep', status: 'pending' });
+ expect(chunks[1]).toMatchObject({ sessionUpdate: 'tool_call_update', toolCallId: 't1', title: 'grep', status: 'completed', content: [{ type: 'text', text: 'match' }] });
+ });
+
+ it('works without onOutputChunk wired (workers)', () => {
+ const session: any = {};
+ expect(process(session, { type: 'assistant.message_delta', data: { deltaContent: 'x' } })).toBe(true);
+ expect(process(session, { type: 'assistant.message', data: { content: 'x' } })).toBe(false);
+ });
+ });
});
diff --git a/packages/server/tests/agents/model-config.test.ts b/packages/server/tests/agents/model-config.test.ts
index e6d96832..5beaebb2 100644
--- a/packages/server/tests/agents/model-config.test.ts
+++ b/packages/server/tests/agents/model-config.test.ts
@@ -153,6 +153,18 @@ agents:
expect(cfg.model).toBe('high');
});
+ it('hasRoleConfig distinguishes explicit config from fallbacks', () => {
+ // lead/worker are configured in the fixture; director is not
+ expect(mc.hasRoleConfig('lead')).toBe(true);
+ expect(mc.hasRoleConfig('worker')).toBe(true);
+ expect(mc.hasRoleConfig('director')).toBe(false);
+ // Regression: an unconfigured director used to resolve to the global
+ // default runtime (codex), which may not even be installed — callers
+ // need this signal to inherit the Lead's runtime/model instead.
+ mc.setRole('director', 'claude:opus');
+ expect(mc.hasRoleConfig('director')).toBe(true);
+ });
+
it('handles missing config file gracefully', () => {
const emptyDir = join(tmpdir(), `fd-empty-${randomUUID().slice(0, 8)}`);
mkdirSync(join(emptyDir, '.flightdeck'), { recursive: true });
diff --git a/packages/server/tests/lead/lead-manager.test.ts b/packages/server/tests/lead/lead-manager.test.ts
index a8929008..23fb4c6b 100644
--- a/packages/server/tests/lead/lead-manager.test.ts
+++ b/packages/server/tests/lead/lead-manager.test.ts
@@ -144,4 +144,55 @@ describe('LeadManager', () => {
expect(lm.getLeadSessionId()).toBeNull();
});
});
+
+ describe('director spawn', () => {
+ function fakeAdapter(spawns: any[], opts?: { failDirector?: boolean }) {
+ return {
+ spawn: async (o: any) => {
+ if (opts?.failDirector && o.role === 'director') {
+ throw new Error('Runtime "codex" is not available: command "codex" not found on PATH.');
+ }
+ spawns.push(o);
+ return { agentId: `${o.role}-fake`, sessionId: `sess-${o.role}`, status: 'running' };
+ },
+ getSession: () => undefined,
+ steer: async () => '',
+ kill: async () => { /* noop */ },
+ getMetadata: async () => null,
+ };
+ }
+
+ it('director inherits lead runtime/model when not explicitly configured', async () => {
+ // Regression: an unconfigured director used to fall back to the global
+ // default runtime (codex) — which may not be installed — and then
+ // silently dropped every message from the Lead.
+ const cwd = project.subpath('.');
+ const { mkdirSync } = await import('node:fs');
+ mkdirSync(join(cwd, '.flightdeck'), { recursive: true });
+ writeFileSync(join(cwd, '.flightdeck', 'config.yaml'), [
+ 'agents:', ' roles:', ' lead:', ' runtime: copilot', ' model: claude-sonnet-4.6', '',
+ ].join('\n'));
+ const spawns: any[] = [];
+ const lm = new LeadManager({
+ sqlite, project, acpAdapter: fakeAdapter(spawns), cwd,
+ leadRuntime: 'copilot' as any, leadModel: 'claude-sonnet-4.6',
+ });
+ await lm.spawnDirector();
+ const d = spawns.find(s => s.role === 'director');
+ expect(d.runtime).toBe('copilot');
+ expect(d.model).toBe('claude-sonnet-4.6');
+ });
+
+ it('surfaces director spawn failure via onSystemNotice instead of failing silently', async () => {
+ const spawns: any[] = [];
+ const notices: string[] = [];
+ const lm = new LeadManager({
+ sqlite, project, acpAdapter: fakeAdapter(spawns, { failDirector: true }), cwd: project.subpath('.'),
+ });
+ lm.onSystemNotice = c => notices.push(c);
+ await lm.spawnLead();
+ expect(notices.some(n => n.includes('Director failed to start'))).toBe(true);
+ expect(notices.some(n => n.includes('not found on PATH'))).toBe(true);
+ });
+ });
});
From a814b49249887d5e3b4cfeb8ccd64570ce94a869 Mon Sep 17 00:00:00 2001
From: Justin Chu
Date: Thu, 11 Jun 2026 20:48:19 -0700
Subject: [PATCH 2/6] =?UTF-8?q?feat:=20agent=20DM=20sender=E2=86=92recipie?=
=?UTF-8?q?nt=20display=20with=20collapsed-by-default=20UX?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
Agent↔agent messages were effectively invisible and unattributed in the
web UI, three stacked causes:
- GET /messages unconditionally filtered out every dm: channel message,
so agent DMs never reached the main chat (new include_agent_dms param,
only admits authorType=agent DMs to keep user steers out)
- the recipient indicator checked msg.channelId, but DM recipients are
stored in msg.channel ('dm:') — wrong field, never rendered
- live dm:message ws events were never ingested by the chat hook
New display config: agentMessages ('off' | 'summary' | 'detail',
default 'summary') in shared DisplayConfig, presets, merge + validation.
- summary: consecutive agent DMs collapse into one expandable line
("worker-x → director · N agent messages") so the main chat stays
readable by default
- detail: full bubbles, each showing sender → recipient
- off: hidden entirely (and not fetched)
Settings popover gets an "Agent messages" visibility selector.
Also: HTTP-level regression test pinning the model-selection roundtrip
(PUT /agents/:id/model visible in next GET /agents).
Co-Authored-By: Claude Fable 5
---
packages/server/src/api/routes/messages.ts | 6 +-
.../server/tests/api/model-roundtrip.test.ts | 86 +++++++++++++++++
packages/shared/src/display.ts | 13 +++
.../web/src/components/DisplaySettings.tsx | 12 ++-
packages/web/src/hooks/useChat.tsx | 26 +++++-
packages/web/src/lib/api.ts | 3 +-
packages/web/src/lib/types.ts | 2 +
packages/web/src/pages/Chat.tsx | 92 +++++++++++++++++--
8 files changed, 227 insertions(+), 13 deletions(-)
create mode 100644 packages/server/tests/api/model-roundtrip.test.ts
diff --git a/packages/server/src/api/routes/messages.ts b/packages/server/src/api/routes/messages.ts
index 7f857971..34101ddf 100644
--- a/packages/server/src/api/routes/messages.ts
+++ b/packages/server/src/api/routes/messages.ts
@@ -149,8 +149,12 @@ export async function handleMessageRoutes(
const authorTypesParam = url.searchParams.get('author_types');
const authorTypes = authorTypesParam ? authorTypesParam.split(',') : undefined;
const limit = parseInt(url.searchParams.get('limit') ?? '50', 10) || 50;
+ // Agent↔agent DM traffic is excluded from the main chat by default;
+ // include_agent_dms=true lets the UI show it (collapsed by default)
+ const includeAgentDms = url.searchParams.get('include_agent_dms') === 'true';
const allMsgs = fd.messages?.listMessages({ taskId, limit: limit + 50, authorTypes }) ?? [];
- const mainChatMsgs = allMsgs.filter(m => !m.channel?.startsWith('dm:'));
+ const mainChatMsgs = allMsgs.filter(m =>
+ !m.channel?.startsWith('dm:') || (includeAgentDms && m.authorType === 'agent'));
json(200, mainChatMsgs.slice(-limit).reverse());
return true;
}
diff --git a/packages/server/tests/api/model-roundtrip.test.ts b/packages/server/tests/api/model-roundtrip.test.ts
new file mode 100644
index 00000000..a98374b4
--- /dev/null
+++ b/packages/server/tests/api/model-roundtrip.test.ts
@@ -0,0 +1,86 @@
+import { describe, it, expect, beforeEach, afterEach } from 'vitest';
+import { createHttpServer } from '../../src/api/HttpServer.js';
+import { Flightdeck } from '../../src/facade.js';
+import type { ProjectManager } from '../../src/projects/ProjectManager.js';
+import type { AgentId } from '@flightdeck-ai/shared';
+import { existsSync, rmSync } from 'node:fs';
+import { join } from 'node:path';
+import { homedir } from 'node:os';
+import http from 'node:http';
+
+/**
+ * Full HTTP-level regression for the "model selection doesn't display" bug:
+ * PUT /agents/:id/model must be visible in the next GET /agents response.
+ * (The model column used to be written via raw SQL but missing from the
+ * drizzle schema, so the GET never returned it and the UI looked unsaved.)
+ */
+describe('Agent model HTTP roundtrip', () => {
+ const projectName = `test-model-rt-${Date.now()}`;
+ let fd: Flightdeck;
+ let server: http.Server;
+ let port: number;
+
+ beforeEach(async () => {
+ fd = new Flightdeck(projectName);
+ const projectManager = {
+ list: () => [projectName],
+ get: (name: string) => (name === projectName ? fd : null),
+ create: () => {},
+ delete: () => true,
+ closeAll: () => {},
+ } as unknown as ProjectManager;
+ server = createHttpServer({
+ projectManager,
+ leadManagers: new Map(),
+ agentManagers: new Map([[projectName, fd.agentManager]]),
+ wsServers: new Map(),
+ webhookNotifiers: new Map(),
+ cronStores: new Map(),
+ port: 0,
+ corsOrigin: '*',
+ } as any);
+ port = await new Promise(resolve => {
+ server.listen(0, '127.0.0.1', () => resolve((server.address() as any).port));
+ });
+ });
+
+ afterEach(() => {
+ server.close();
+ fd.close();
+ const projDir = join(homedir(), '.flightdeck', 'v2', 'projects', projectName);
+ if (existsSync(projDir)) rmSync(projDir, { recursive: true, force: true });
+ });
+
+ async function req(method: string, path: string, body?: unknown) {
+ const res = await fetch(`http://127.0.0.1:${port}/api/projects/${projectName}${path}`, {
+ method,
+ headers: body ? { 'Content-Type': 'application/json' } : undefined,
+ body: body ? JSON.stringify(body) : undefined,
+ });
+ return { status: res.status, data: await res.json().catch(() => null) };
+ }
+
+ it('PUT /agents/:id/model is reflected in GET /agents', async () => {
+ fd.registerAgent({
+ id: 'worker-rt-1' as AgentId,
+ role: 'worker',
+ runtime: 'acp',
+ runtimeName: 'copilot',
+ acpSessionId: null,
+ status: 'idle',
+ currentSpecId: null,
+ costAccumulated: 0,
+ lastHeartbeat: null,
+ });
+
+ const put = await req('PUT', '/agents/worker-rt-1/model', { model: 'claude-sonnet-4.6', runtime: 'copilot' });
+ expect(put.status).toBe(200);
+
+ const get = await req('GET', '/agents?include_retired=true');
+ expect(get.status).toBe(200);
+ const agent = (get.data as any[]).find(a => a.id === 'worker-rt-1');
+ expect(agent).toBeDefined();
+ expect(agent.model).toBe('claude-sonnet-4.6');
+ expect(agent.runtimeName).toBe('copilot');
+ });
+});
diff --git a/packages/shared/src/display.ts b/packages/shared/src/display.ts
index e7558119..d7abf28e 100644
--- a/packages/shared/src/display.ts
+++ b/packages/shared/src/display.ts
@@ -11,6 +11,12 @@ export interface DisplayConfig {
toolCalls: ToolVisibility;
/** Flightdeck internal tool calls (flightdeck_* prefix) visibility */
flightdeckTools: ToolVisibility;
+ /**
+ * Agent↔agent DM visibility in the main chat.
+ * 'off' hides them, 'summary' collapses consecutive DMs into an expandable
+ * one-liner (default), 'detail' shows full bubbles with sender → recipient.
+ */
+ agentMessages?: ToolVisibility;
/** Per-tool overrides (tool name → visibility) */
toolOverrides?: Record;
}
@@ -21,21 +27,25 @@ export const DISPLAY_PRESETS = {
thinking: false,
toolCalls: 'off' as const,
flightdeckTools: 'off' as const,
+ agentMessages: 'off' as const,
},
summary: {
thinking: false,
toolCalls: 'summary' as const,
flightdeckTools: 'off' as const,
+ agentMessages: 'summary' as const,
},
detail: {
thinking: true,
toolCalls: 'detail' as const,
flightdeckTools: 'summary' as const,
+ agentMessages: 'detail' as const,
},
debug: {
thinking: true,
toolCalls: 'detail' as const,
flightdeckTools: 'detail' as const,
+ agentMessages: 'detail' as const,
},
} as const;
@@ -116,6 +126,8 @@ export function mergeDisplayConfig(
thinking: partial.thinking ?? base.thinking,
toolCalls: partial.toolCalls ?? base.toolCalls,
flightdeckTools: partial.flightdeckTools ?? base.flightdeckTools,
+ // Persisted configs may predate this field — default to collapsed
+ agentMessages: partial.agentMessages ?? base.agentMessages ?? 'summary',
toolOverrides,
};
}
@@ -132,6 +144,7 @@ export function isValidDisplayConfig(v: unknown): v is PartialDisplayConfig {
if (obj.thinking !== undefined && typeof obj.thinking !== 'boolean') return false;
if (obj.toolCalls !== undefined && !isValidToolVisibility(obj.toolCalls)) return false;
if (obj.flightdeckTools !== undefined && !isValidToolVisibility(obj.flightdeckTools)) return false;
+ if (obj.agentMessages !== undefined && !isValidToolVisibility(obj.agentMessages)) return false;
if (obj.toolOverrides !== undefined) {
if (typeof obj.toolOverrides !== 'object' || obj.toolOverrides === null || Array.isArray(obj.toolOverrides)) return false;
for (const val of Object.values(obj.toolOverrides as Record)) {
diff --git a/packages/web/src/components/DisplaySettings.tsx b/packages/web/src/components/DisplaySettings.tsx
index 105871da..b1c25c14 100644
--- a/packages/web/src/components/DisplaySettings.tsx
+++ b/packages/web/src/components/DisplaySettings.tsx
@@ -17,7 +17,8 @@ export function DisplaySettings({ onClose }: { onClose: () => void }) {
const preset = DISPLAY_PRESETS[p];
return preset.thinking === displayConfig.thinking
&& preset.toolCalls === displayConfig.toolCalls
- && preset.flightdeckTools === displayConfig.flightdeckTools;
+ && preset.flightdeckTools === displayConfig.flightdeckTools
+ && preset.agentMessages === (displayConfig.agentMessages ?? 'summary');
}) ?? 'custom';
// M9: Basic focus trap — keep keyboard focus inside the dialog
@@ -105,6 +106,15 @@ export function DisplaySettings({ onClose }: { onClose: () => void }) {
onChange={v => setDisplayConfig({ flightdeckTools: v })}
/>
+
+ {/* Agent ↔ agent messages */}
+
+ Agent messages
+ setDisplayConfig({ agentMessages: v })}
+ />
+
{/* Advanced: per-tool overrides */}
diff --git a/packages/web/src/hooks/useChat.tsx b/packages/web/src/hooks/useChat.tsx
index 136f69e5..8818334d 100644
--- a/packages/web/src/hooks/useChat.tsx
+++ b/packages/web/src/hooks/useChat.tsx
@@ -70,9 +70,16 @@ export function ChatProvider({ children }: { children: ReactNode }) {
};
// SWR for initial message load
+ const agentMsgVisibility = displayConfig.agentMessages ?? 'summary';
const { data: initialMessages } = useSWR(
- projectName ? ['messages', projectName, displayConfig.flightdeckTools] : null,
- () => api.getMessages(projectName!, { limit: 100, author_types: displayConfig.flightdeckTools === 'detail' ? undefined : 'user,lead,system' })
+ projectName ? ['messages', projectName, displayConfig.flightdeckTools, agentMsgVisibility] : null,
+ () => api.getMessages(projectName!, {
+ limit: 100,
+ author_types: displayConfig.flightdeckTools === 'detail'
+ ? undefined
+ : agentMsgVisibility !== 'off' ? 'user,lead,system,agent' : 'user,lead,system',
+ include_agent_dms: agentMsgVisibility !== 'off',
+ })
);
// Clear WS messages when project changes
@@ -101,7 +108,9 @@ export function ChatProvider({ children }: { children: ReactNode }) {
case 'chat:message': {
const msg = event.message;
const isDebugMode = displayConfigRef.current.flightdeckTools === 'detail';
- if (!isDebugMode && msg.authorType && msg.authorType !== 'user' && msg.authorType !== 'lead' && msg.authorType !== 'system') break;
+ const showAgentMsgs = (displayConfigRef.current.agentMessages ?? 'summary') !== 'off';
+ if (!isDebugMode && msg.authorType && msg.authorType !== 'user' && msg.authorType !== 'lead' && msg.authorType !== 'system'
+ && !(showAgentMsgs && msg.authorType === 'agent')) break;
// Browser notification when Lead replies and tab is not focused
if (msg.authorType === 'lead' && !document.hasFocus() && 'Notification' in window && Notification.permission === 'granted') {
new Notification('Flightdeck — Lead replied', { body: (msg.content || '').slice(0, 100) });
@@ -159,6 +168,17 @@ export function ChatProvider({ children }: { children: ReactNode }) {
return [...prev.slice(-(MAX_MESSAGES - 1)), event.message];
});
break;
+ case 'dm:message': {
+ // Live agent↔agent DMs — shown in main chat unless turned off
+ if ((displayConfigRef.current.agentMessages ?? 'summary') === 'off') break;
+ const dm = (event as any).message;
+ if (!dm?.id) break;
+ setWsMessages(prev => {
+ if (prev.some(m => m.id === dm.id)) return prev;
+ return [...prev.slice(-(MAX_MESSAGES - 1)), dm];
+ });
+ break;
+ }
}
});
}, [subscribe]);
diff --git a/packages/web/src/lib/api.ts b/packages/web/src/lib/api.ts
index e3a61b56..c5cb149a 100644
--- a/packages/web/src/lib/api.ts
+++ b/packages/web/src/lib/api.ts
@@ -65,13 +65,14 @@ export const api = {
if (opts?.limit) params.set('limit', String(opts.limit));
return get(projectPath(project, `/activity?${params}`));
},
- getMessages: (project: string, opts?: { thread_id?: string; task_id?: string; limit?: number; author_types?: string; channel?: string }) => {
+ getMessages: (project: string, opts?: { thread_id?: string; task_id?: string; limit?: number; author_types?: string; channel?: string; include_agent_dms?: boolean }) => {
const params = new URLSearchParams();
if (opts?.thread_id) params.set('thread_id', opts.thread_id);
if (opts?.task_id) params.set('task_id', opts.task_id);
if (opts?.limit) params.set('limit', String(opts.limit));
if (opts?.author_types) params.set('author_types', opts.author_types);
if (opts?.channel) params.set('channel', opts.channel);
+ if (opts?.include_agent_dms) params.set('include_agent_dms', 'true');
return get(projectPath(project, `/messages?${params}`));
},
getReport: async (project: string): Promise => {
diff --git a/packages/web/src/lib/types.ts b/packages/web/src/lib/types.ts
index 547a00d2..494169aa 100644
--- a/packages/web/src/lib/types.ts
+++ b/packages/web/src/lib/types.ts
@@ -69,6 +69,8 @@ export interface ChatMessage {
senderName?: string | null;
replyToId?: string | null;
attachments?: Array<{ url: string; filename: string; mimeType: string; size: number }> | null;
+ /** DM/channel routing: 'dm:' identifies the recipient of an agent DM */
+ channel?: string | null;
channelId?: string | null;
createdAt: string;
updatedAt: string | null;
diff --git a/packages/web/src/pages/Chat.tsx b/packages/web/src/pages/Chat.tsx
index 2b27cb55..ec12288c 100644
--- a/packages/web/src/pages/Chat.tsx
+++ b/packages/web/src/pages/Chat.tsx
@@ -150,11 +150,18 @@ const MessageBubble = memo(function MessageBubble({ msg, messages, replyCountMap
{msg.authorType === 'agent' && msg.authorId && (
{msg.authorId}
)}
- {msg.channelId?.startsWith('dm:') && (
-
- → {msg.channelId.replace('dm:', '').replace(/-[a-z0-9]+$/, '')}
-
- )}
+ {(() => {
+ // Agent DMs store the recipient in `channel` ('dm:');
+ // `channelId` is only set by external bridges
+ const dmChannel = msg.channel?.startsWith('dm:') ? msg.channel : msg.channelId?.startsWith('dm:') ? msg.channelId : null;
+ if (!dmChannel) return null;
+ const recipient = dmChannel.slice(3);
+ return (
+
+ → {recipient}
+
+ );
+ })()}
{new Date(msg.createdAt).toLocaleTimeString()}
@@ -276,6 +283,55 @@ export function ToolCallCard({ tc, level }: { tc: ToolCallState; level: 'summary
);
}
+/** True for agent↔agent DM traffic (sender in authorId, recipient in channel). */
+function isAgentDm(m: ChatMessage): boolean {
+ return m.authorType === 'agent' && !!(m.channel?.startsWith('dm:') || m.channelId?.startsWith('dm:'));
+}
+
+function dmRecipient(m: ChatMessage): string {
+ const ch = m.channel?.startsWith('dm:') ? m.channel : m.channelId ?? '';
+ return ch.startsWith('dm:') ? ch.slice(3) : '?';
+}
+
+/** Collapsed run of consecutive agent↔agent DMs — one line, expandable. */
+function AgentDmGroup({ msgs, replyCountMap, onReply, agents }: {
+ msgs: ChatMessage[];
+ replyCountMap?: Map;
+ onReply: (m: ChatMessage) => void;
+ agents?: Array<{ id: string; role: string; runtime?: string; runtimeName?: string; model?: string; status?: string }>;
+}) {
+ const [open, setOpen] = useState(false);
+ const pairs = useMemo(
+ () => [...new Set(msgs.map(m => `${m.authorId ?? '?'} → ${dmRecipient(m)}`))],
+ [msgs]
+ );
+ return (
+
+
+ {open && (
+
+ {msgs.map(m => (
+
+ ))}
+
+ )}
+
+ );
+}
+
function StreamingBubble({ content, chunks, toolCallMap, displayConfig }: {
content: string;
chunks?: StreamChunk[];
@@ -591,6 +647,26 @@ export default function Chat() {
return map;
}, [filteredMessages]);
+ // Group consecutive agent↔agent DMs per display config:
+ // off → dropped, summary → collapsed expandable run, detail → full bubbles
+ const renderItems = useMemo(() => {
+ const visibility = displayConfig.agentMessages ?? 'summary';
+ const items: Array<{ kind: 'msg'; msg: ChatMessage } | { kind: 'dmGroup'; key: string; msgs: ChatMessage[] }> = [];
+ for (const m of filteredMessages) {
+ if (isAgentDm(m)) {
+ if (visibility === 'off') continue;
+ if (visibility === 'summary') {
+ const last = items[items.length - 1];
+ if (last?.kind === 'dmGroup') last.msgs.push(m);
+ else items.push({ kind: 'dmGroup', key: `dmg-${m.id}`, msgs: [m] });
+ continue;
+ }
+ }
+ items.push({ kind: 'msg', msg: m });
+ }
+ return items;
+ }, [filteredMessages, displayConfig.agentMessages]);
+
const handleReply = useCallback((m: ChatMessage) => setReplyTo(m), []);
useEffect(() => {
@@ -746,8 +822,10 @@ export default function Chat() {
)}
- {filteredMessages.map(msg => (
-
+ {renderItems.map(item => item.kind === 'dmGroup' ? (
+
+ ) : (
+
))}
{streamEntries.map(([id, content]) => (
Date: Thu, 11 Jun 2026 22:03:55 -0700
Subject: [PATCH 3/6] fix: apply verified findings from multi-agent repo audit
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
From a 6-dimension audit (40 agents, adversarially verified — 15 confirmed
of 34 raw findings), fixing the low-risk high-value ones:
- CopilotSdkAdapter.steer(): every call leaked a session.on handler (and
the 5-min timeout left another) — unsubscribe on idle/timeout. N steers
used to mean N live handlers firing on every subsequent event.
- LeadManager.resumeLead/resumeDirector: insertAgent on an existing agent
row hit the primary key, so resumes of persisted leads/directors always
failed into "marking offline". Upsert instead. resumeDirector also never
passed runtime, routing copilot-sdk directors to the ACP adapter.
- DecisionLog.readAll(): one corrupted JSONL line crashed every route that
lists decisions — skip bad lines (matches SuggestionStore behavior).
- useAgents: dmMessagesRef survived project switches — unbounded memory
growth across projects; clear it with the other per-project state.
- HTTP DM path now stores recipient (was only derivable from channel).
- drizzle schema: read_state composite PK (matches sql/schema.sql).
- web ws types: declare the project field the server actually sends.
Co-Authored-By: Claude Fable 5
---
.../server/src/agents/CopilotSdkAdapter.ts | 14 ++--
packages/server/src/api/routes/messages.ts | 2 +-
packages/server/src/db/schema.ts | 7 +-
packages/server/src/lead/LeadManager.ts | 66 ++++++++++++-------
packages/server/src/storage/DecisionLog.ts | 8 ++-
packages/web/src/hooks/useAgents.tsx | 2 +
packages/web/src/lib/ws.ts | 4 +-
7 files changed, 67 insertions(+), 36 deletions(-)
diff --git a/packages/server/src/agents/CopilotSdkAdapter.ts b/packages/server/src/agents/CopilotSdkAdapter.ts
index a7da54f9..76c0e5b4 100644
--- a/packages/server/src/agents/CopilotSdkAdapter.ts
+++ b/packages/server/src/agents/CopilotSdkAdapter.ts
@@ -1227,16 +1227,20 @@ export class CopilotSdkAdapter extends AgentAdapter {
await agentSession.session.send({ prompt: message.content });
- // Wait for idle (turn complete)
+ // Wait for idle (turn complete). session.on returns an unsubscribe fn —
+ // without calling it, every steer() leaks a handler that keeps firing
+ // on all future session events.
await new Promise((resolve) => {
- const handler = (event: SessionEvent) => {
+ let timer: ReturnType | undefined;
+ const unsubscribe = agentSession.session.on((event: SessionEvent) => {
if (event.type === 'session.idle') {
+ if (timer) clearTimeout(timer);
+ unsubscribe();
resolve();
}
- };
- agentSession.session.on(handler);
+ });
// Timeout after 5 minutes
- setTimeout(() => resolve(), 5 * 60 * 1000);
+ timer = setTimeout(() => { unsubscribe(); resolve(); }, 5 * 60 * 1000);
});
return agentSession.output.slice(outputBefore);
diff --git a/packages/server/src/api/routes/messages.ts b/packages/server/src/api/routes/messages.ts
index 34101ddf..c4d764e9 100644
--- a/packages/server/src/api/routes/messages.ts
+++ b/packages/server/src/api/routes/messages.ts
@@ -231,7 +231,7 @@ export async function handleMessageRoutes(
parentId: body.parentId ?? null, taskId: null,
authorType: 'agent', authorId: agentId,
content: (body.content as string).length > 4000 ? (body.content as string).slice(0, 4000) + '\n\u2026[truncated]' : body.content,
- metadata: null, channel: `dm:${body.to}`,
+ metadata: null, channel: `dm:${body.to}`, recipient: body.to,
replyToId: body.parentId ?? body.replyToId ?? null,
});
}
diff --git a/packages/server/src/db/schema.ts b/packages/server/src/db/schema.ts
index b33520a2..1bcb85a2 100644
--- a/packages/server/src/db/schema.ts
+++ b/packages/server/src/db/schema.ts
@@ -1,4 +1,4 @@
-import { sqliteTable, text, integer, real, index } from 'drizzle-orm/sqlite-core';
+import { sqliteTable, text, integer, real, index, primaryKey } from 'drizzle-orm/sqlite-core';
import { sql } from 'drizzle-orm';
/** ISO 8601 UTC timestamp with Z suffix — use instead of datetime('now') to avoid timezone ambiguity */
@@ -104,7 +104,10 @@ export const readState = sqliteTable('read_state', {
agentId: text('agent_id').notNull(),
channel: text('channel').notNull().default('dm'),
lastReadAt: text('last_read_at').notNull(),
-}, () => []);
+}, (table) => [
+ // Matches sql/schema.sql — markRead's ON CONFLICT(agent_id, channel) relies on it
+ primaryKey({ columns: [table.agentId, table.channel] }),
+]);
// ── Channel Subscriptions ────────────────────────────────────────────
diff --git a/packages/server/src/lead/LeadManager.ts b/packages/server/src/lead/LeadManager.ts
index 8688869e..9ad9f33b 100644
--- a/packages/server/src/lead/LeadManager.ts
+++ b/packages/server/src/lead/LeadManager.ts
@@ -1049,20 +1049,27 @@ export class LeadManager {
runtime: this.leadRuntime,
});
this.leadSessionId = meta.sessionId;
- this.leadAgentId = meta.agentId;
+ this.leadAgentId = meta.agentId;
this.wireStreamHandler();
- this.sqlite.insertAgent({
- id: meta.agentId,
- role: 'lead',
- runtime: this.leadRuntime ?? 'acp',
- runtimeName: this.leadRuntime ?? 'codex',
- acpSessionId: meta.sessionId,
- status: 'idle',
- currentSpecId: null,
- costAccumulated: 0,
- lastHeartbeat: null,
- });
+ // The agent row usually still exists when resuming — an INSERT would
+ // hit the primary key and make every resume "fail". Upsert instead.
+ if (this.sqlite.getAgent(meta.agentId as any)) {
+ this.sqlite.updateAgentAcpSession(meta.agentId as any, meta.sessionId);
+ this.sqlite.updateAgentStatus(meta.agentId as any, 'idle');
+ } else {
+ this.sqlite.insertAgent({
+ id: meta.agentId,
+ role: 'lead',
+ runtime: this.leadRuntime ?? 'acp',
+ runtimeName: this.leadRuntime ?? null,
+ acpSessionId: meta.sessionId,
+ status: 'idle',
+ currentSpecId: null,
+ costAccumulated: 0,
+ lastHeartbeat: null,
+ });
+ }
return meta.sessionId;
} catch (err) {
@@ -1082,21 +1089,30 @@ export class LeadManager {
cwd,
role: 'director',
model,
+ // Without the runtime, MultiAdapter routes the resume to the
+ // default (ACP) adapter even for copilot-sdk directors
+ runtime: this.directorRuntime,
});
this.directorSessionId = meta.sessionId;
- this.directorAgentId = meta.agentId;
-
- this.sqlite.insertAgent({
- id: meta.agentId,
- role: 'director',
- runtime: this.directorRuntime ?? this.leadRuntime ?? 'acp',
- runtimeName: this.directorRuntime ?? this.leadRuntime ?? 'codex',
- acpSessionId: meta.sessionId,
- status: 'idle',
- currentSpecId: null,
- costAccumulated: 0,
- lastHeartbeat: null,
- });
+ this.directorAgentId = meta.agentId;
+
+ // Upsert — the agent row usually still exists when resuming
+ if (this.sqlite.getAgent(meta.agentId as any)) {
+ this.sqlite.updateAgentAcpSession(meta.agentId as any, meta.sessionId);
+ this.sqlite.updateAgentStatus(meta.agentId as any, 'idle');
+ } else {
+ this.sqlite.insertAgent({
+ id: meta.agentId,
+ role: 'director',
+ runtime: this.directorRuntime ?? this.leadRuntime ?? 'acp',
+ runtimeName: this.directorRuntime ?? this.leadRuntime ?? null,
+ acpSessionId: meta.sessionId,
+ status: 'idle',
+ currentSpecId: null,
+ costAccumulated: 0,
+ lastHeartbeat: null,
+ });
+ }
return meta.sessionId;
} catch (err) {
diff --git a/packages/server/src/storage/DecisionLog.ts b/packages/server/src/storage/DecisionLog.ts
index ac534cc3..e835a21d 100644
--- a/packages/server/src/storage/DecisionLog.ts
+++ b/packages/server/src/storage/DecisionLog.ts
@@ -23,7 +23,13 @@ export class DecisionLog {
const filepath = join(this.decisionsDir, filename);
if (!existsSync(filepath)) return [];
const lines = readFileSync(filepath, 'utf-8').trim().split('\n').filter(Boolean);
- return lines.map(l => JSON.parse(l) as Decision);
+ // A single corrupted line must not take down every API route that lists
+ // decisions — skip it instead
+ const decisions: Decision[] = [];
+ for (const l of lines) {
+ try { decisions.push(JSON.parse(l) as Decision); } catch { /* skip corrupted line */ }
+ }
+ return decisions;
}
list(opts?: DecisionListOptions, filename: string = 'decisions.jsonl'): Decision[] {
diff --git a/packages/web/src/hooks/useAgents.tsx b/packages/web/src/hooks/useAgents.tsx
index 95832832..3ab1e7e9 100644
--- a/packages/web/src/hooks/useAgents.tsx
+++ b/packages/web/src/hooks/useAgents.tsx
@@ -56,8 +56,10 @@ export function AgentProvider({ children }: { children: ReactNode }) {
useEffect(() => {
agentOutputsRef.current.clear();
agentStreamChunksRef.current.clear();
+ dmMessagesRef.current.clear();
setAgentOutputs(new Map());
setAgentStreamChunks(new Map());
+ setDmMessages(new Map());
}, [projectName]);
useEffect(() => {
diff --git a/packages/web/src/lib/ws.ts b/packages/web/src/lib/ws.ts
index 39d83406..c063c770 100644
--- a/packages/web/src/lib/ws.ts
+++ b/packages/web/src/lib/ws.ts
@@ -3,10 +3,10 @@ import type { DisplayConfig, ContentType } from '@flightdeck-ai/shared/display';
import { WS_INITIAL_BACKOFF_MS, WS_MAX_BACKOFF_MS } from './constants.ts';
export type WsEvent =
- | { type: 'chat:message'; message: ChatMessage }
+ | { type: 'chat:message'; message: ChatMessage; project?: string }
| { type: 'chat:stream'; message_id: string; delta: string; done: boolean; content_type?: ContentType; tool_name?: string }
| { type: 'thread:created'; thread: Thread }
- | { type: 'task:comment'; task_id: string; message: ChatMessage }
+ | { type: 'task:comment'; task_id: string; message: ChatMessage; project?: string }
| { type: 'display:config'; config: DisplayConfig }
| { type: 'state:update'; stats: Record }
| { type: 'agent:stream'; agentId: string; delta: string; contentType: 'text' | 'thinking' | 'tool_call' | 'tool_result'; toolName?: string }
From 31612b437c252c1c5234301b84fa1b481bfead15 Mon Sep 17 00:00:00 2001
From: Justin Chu
Date: Thu, 11 Jun 2026 22:23:08 -0700
Subject: [PATCH 4/6] fix: DM format unification, safe agent purge, adapter
session lifecycle
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
Remaining verified audit findings (items 2-5):
- purgeOfflineAgents now skips hibernated agents that are still resumable
(saved acp_session_id) or have in-flight tasks assigned — purging those
orphaned tasks and lost resumable sessions on every Lead spawn.
- DM channel naming unified on 'dm:' + recipient column.
Three formats coexisted (bare 'dm'+recipient via appendDM, 'dm:'
without recipient via AgentManager steers, 'dm:'+recipient via
HTTP), so no single query could see all DMs — and bare-'dm' rows leaked
past the main-chat dm: filter unattributed. Idempotent migration in
SqliteStore normalizes existing rows; getUnreadDMs tolerates legacy
shapes; web recipient arrow reads channel/recipient/channelId.
- MultiAdapter sessionAdapterMap is now cleaned when sessions end
naturally (crash/exit), not just on kill() — transparent onSessionEnd
property hook that preserves gateway-assigned handlers.
- CopilotSdkAdapter spawn/resume no longer stack session.on handlers when
the same session is re-registered; disposer stored per session and
released on kill.
- Deflake messagelog markRead test (ms-granularity boundary).
Co-Authored-By: Claude Fable 5
---
packages/server/src/agents/AgentManager.ts | 4 +-
.../server/src/agents/CopilotSdkAdapter.ts | 14 ++++-
packages/server/src/agents/MultiAdapter.ts | 23 ++++++++
packages/server/src/comms/MessageStore.ts | 10 ++--
packages/server/src/storage/SqliteStore.ts | 22 +++++++-
.../server/tests/agents/multi-adapter.test.ts | 54 +++++++++++++++++++
.../server/tests/comms/message-store.test.ts | 22 ++++++++
.../tests/orchestrator/suspended.test.ts | 36 ++++++++++++-
.../server/tests/storage/messagelog.test.ts | 7 ++-
packages/server/tests/storage/sqlite.test.ts | 18 +++++++
packages/web/src/lib/types.ts | 2 +
packages/web/src/pages/Chat.tsx | 18 +++----
12 files changed, 209 insertions(+), 21 deletions(-)
create mode 100644 packages/server/tests/agents/multi-adapter.test.ts
diff --git a/packages/server/src/agents/AgentManager.ts b/packages/server/src/agents/AgentManager.ts
index 2ed307e3..a4dcb290 100644
--- a/packages/server/src/agents/AgentManager.ts
+++ b/packages/server/src/agents/AgentManager.ts
@@ -517,7 +517,7 @@ export class AgentManager {
parentId: null, taskId: null,
authorType: 'system', authorId: null,
content: message.length > 4000 ? message.slice(0, 4000) + '\n…[truncated]' : message,
- metadata: null, channel: `dm:${agentId}`,
+ metadata: null, channel: `dm:${agentId}`, recipient: agentId,
});
if (this.onDmMessage) this.onDmMessage(this.projectName, dmMsg);
}
@@ -542,7 +542,7 @@ export class AgentManager {
parentId: null, taskId: null,
authorType: 'system', authorId: null,
content: message.length > 4000 ? message.slice(0, 4000) + '\n…[truncated]' : message,
- metadata: null, channel: `dm:${agentId}`,
+ metadata: null, channel: `dm:${agentId}`, recipient: agentId,
});
if (this.onDmMessage) this.onDmMessage(this.projectName, dmMsg);
}
diff --git a/packages/server/src/agents/CopilotSdkAdapter.ts b/packages/server/src/agents/CopilotSdkAdapter.ts
index 76c0e5b4..88cbc7bf 100644
--- a/packages/server/src/agents/CopilotSdkAdapter.ts
+++ b/packages/server/src/agents/CopilotSdkAdapter.ts
@@ -52,6 +52,8 @@ export interface CopilotAgentSession {
pendingToolCalls?: Map;
/** ACP-style per-chunk callback (set by LeadManager for chat streaming). */
onOutputChunk?: (update: unknown) => void;
+ /** Disposer for the session.on event handler (prevents stacking on resume). */
+ eventUnsubscribe?: () => void;
/** True once a text/reasoning delta was streamed this turn — used to drop
* the redundant full-content event that would duplicate streamed words. */
sawTextDelta?: boolean;
@@ -1104,10 +1106,13 @@ export class CopilotSdkAdapter extends AgentAdapter {
model: opts.model,
};
+ // Replacing an existing session entry? Detach its handler first so
+ // events don't fan out to stale handlers from previous registrations.
+ this.sessions.get(sessionId)?.eventUnsubscribe?.();
this.sessions.set(sessionId, agentSession);
// Wire up event handlers
- session.on((event: SessionEvent) => {
+ agentSession.eventUnsubscribe = session.on((event: SessionEvent) => {
agentSession.lastActivityAt = new Date();
// Capture resolved model from session.created event
@@ -1255,6 +1260,8 @@ export class CopilotSdkAdapter extends AgentAdapter {
try {
await agentSession.session.disconnect();
} catch { /* */ }
+ agentSession.eventUnsubscribe?.();
+ agentSession.eventUnsubscribe = undefined;
agentSession.status = 'ended';
// Clean up after grace period
setTimeout(() => this.sessions.delete(sessionId), 60_000);
@@ -1324,10 +1331,13 @@ export class CopilotSdkAdapter extends AgentAdapter {
model: opts.model,
};
+ // Resuming the same session id again must not stack handlers — detach
+ // the previous registration first.
+ this.sessions.get(sessionId)?.eventUnsubscribe?.();
this.sessions.set(sessionId, agentSession);
// Wire same event handlers as spawn
- session.on((event: SessionEvent) => {
+ agentSession.eventUnsubscribe = session.on((event: SessionEvent) => {
agentSession.lastActivityAt = new Date();
if (event.type === 'assistant.message') {
agentSession.output += event.data.content;
diff --git a/packages/server/src/agents/MultiAdapter.ts b/packages/server/src/agents/MultiAdapter.ts
index d4650f80..4af269e0 100644
--- a/packages/server/src/agents/MultiAdapter.ts
+++ b/packages/server/src/agents/MultiAdapter.ts
@@ -17,6 +17,29 @@ export class MultiAdapter extends AgentAdapter {
super();
this.acpAdapter = acpAdapter;
this.copilotSdkAdapter = copilotSdkAdapter ?? null;
+ this.hookSessionEnd(this.acpAdapter);
+ if (this.copilotSdkAdapter) this.hookSessionEnd(this.copilotSdkAdapter);
+ }
+
+ /**
+ * Drop sessionAdapterMap entries when sessions end naturally (crash, exit) —
+ * kill() is not the only way sessions die, and without this the map grows
+ * for the lifetime of the gateway. Implemented as a transparent property
+ * hook so consumers (gateway, facade) can still assign onSessionEnd later.
+ */
+ private hookSessionEnd(adapter: AgentAdapter): void {
+ let external: ((sessionId: string, session: unknown) => void) | null =
+ (adapter as { onSessionEnd?: ((sessionId: string, session: unknown) => void) | null }).onSessionEnd ?? null;
+ const wrapper = (sessionId: string, session: unknown) => {
+ this.sessionAdapterMap.delete(sessionId);
+ if (external) { try { external(sessionId, session); } catch { /* consumer errors are not ours */ } }
+ };
+ Object.defineProperty(adapter, 'onSessionEnd', {
+ get: () => wrapper,
+ set: (fn) => { external = fn; },
+ configurable: true,
+ enumerable: true,
+ });
}
private pickAdapter(runtime?: string): AgentAdapter {
diff --git a/packages/server/src/comms/MessageStore.ts b/packages/server/src/comms/MessageStore.ts
index a11dd8f6..93e55d40 100644
--- a/packages/server/src/comms/MessageStore.ts
+++ b/packages/server/src/comms/MessageStore.ts
@@ -1,4 +1,4 @@
-import { eq, desc, and, lt, gt, isNotNull, sql } from 'drizzle-orm';
+import { eq, desc, and, or, lt, gt, isNotNull, sql } from 'drizzle-orm';
import { messages, readState, channelSubscriptions, channels } from '../db/schema.js';
import type { FlightdeckDatabase } from '../db/database.js';
import { messageId } from '@flightdeck-ai/shared';
@@ -331,6 +331,8 @@ export class MessageStore {
// ── DM ────────────────────────────────────────────────────────────────
appendDM(from: string, to: string, content: string): ChatMessage {
+ // Canonical DM format: channel 'dm:' + recipient column.
+ // (Bare channel 'dm' predates this and is migrated in SqliteStore.)
return this.createMessage({
parentId: null,
taskId: null,
@@ -338,7 +340,7 @@ export class MessageStore {
authorId: from,
content,
metadata: null,
- channel: 'dm',
+ channel: `dm:${to}`,
recipient: to,
});
}
@@ -346,8 +348,8 @@ export class MessageStore {
getUnreadDMs(agentId: string): ChatMessage[] {
const lastRead = this.getLastRead(agentId, 'dm');
const conditions = [
- eq(messages.recipient, agentId),
- eq(messages.channel, 'dm'),
+ // Match canonical 'dm:' rows; recipient covers legacy bare-'dm' rows
+ or(eq(messages.channel, `dm:${agentId}`), and(eq(messages.channel, 'dm'), eq(messages.recipient, agentId)))!,
];
if (lastRead) conditions.push(gt(messages.createdAt, lastRead));
const rows = this.db
diff --git a/packages/server/src/storage/SqliteStore.ts b/packages/server/src/storage/SqliteStore.ts
index 8913ef1e..87534137 100644
--- a/packages/server/src/storage/SqliteStore.ts
+++ b/packages/server/src/storage/SqliteStore.ts
@@ -67,6 +67,11 @@ export class SqliteStore extends EventEmitter {
// Re-run index creation after columns are ensured
try { this._db.run(sql.raw('CREATE INDEX IF NOT EXISTS `idx_messages_channel` ON `messages` (`channel`)')); } catch {}
try { this._db.run(sql.raw('CREATE INDEX IF NOT EXISTS `idx_messages_recipient` ON `messages` (`recipient`)')); } catch {}
+ // One-time DM format migration (idempotent): unify on channel 'dm:'
+ // + recipient column. Legacy rows used bare channel 'dm' (recipient-only)
+ // or 'dm:' without recipient, so no single query could find all DMs.
+ try { this._db.run(sql.raw(`UPDATE messages SET channel = 'dm:' || recipient WHERE channel = 'dm' AND recipient IS NOT NULL AND recipient != ''`)); } catch {}
+ try { this._db.run(sql.raw(`UPDATE messages SET recipient = substr(channel, 4) WHERE channel LIKE 'dm:%' AND (recipient IS NULL OR recipient = '')`)); } catch {}
}
private addColumnIfMissing(table: string, column: string, type: string): void {
@@ -340,9 +345,22 @@ export class SqliteStore extends EventEmitter {
return result.changes > 0;
}
- /** Remove all agents with status 'hibernated'. Returns count deleted. */
+ /**
+ * Remove dead hibernated agents. Skips agents that are still resumable
+ * (have a saved session) or have in-flight tasks assigned — purging those
+ * would orphan the tasks and lose resumable sessions.
+ */
purgeOfflineAgents(): number {
- const result = this._db.delete(agents).where(eq(agents.status, 'hibernated')).run();
+ const result = this._db.run(sql`
+ DELETE FROM agents
+ WHERE status = 'hibernated'
+ AND (acp_session_id IS NULL OR acp_session_id = '')
+ AND id NOT IN (
+ SELECT assigned_agent FROM tasks
+ WHERE assigned_agent IS NOT NULL
+ AND state IN ('running', 'in_review', 'claimed')
+ )
+ `);
return result.changes;
}
diff --git a/packages/server/tests/agents/multi-adapter.test.ts b/packages/server/tests/agents/multi-adapter.test.ts
new file mode 100644
index 00000000..9b7bea0e
--- /dev/null
+++ b/packages/server/tests/agents/multi-adapter.test.ts
@@ -0,0 +1,54 @@
+import { describe, it, expect } from 'vitest';
+import { MultiAdapter } from '../../src/agents/MultiAdapter.js';
+import { AgentAdapter, type SpawnOptions, type SteerMessage, type AgentMetadata } from '../../src/agents/AgentAdapter.js';
+import type { AgentId, AgentRuntime } from '@flightdeck-ai/shared';
+
+class FakeAdapter extends AgentAdapter {
+ readonly runtime: AgentRuntime = 'acp';
+ onSessionEnd: ((sessionId: string, session: unknown) => void) | null = null;
+ steerCalls: string[] = [];
+ private counter = 0;
+
+ async spawn(opts: SpawnOptions): Promise {
+ const sessionId = `fake-${++this.counter}`;
+ return { agentId: `${opts.role}-fake` as AgentId, sessionId, status: 'running' };
+ }
+ async steer(sessionId: string, _message: SteerMessage): Promise {
+ this.steerCalls.push(sessionId);
+ return 'ok';
+ }
+ async kill(_sessionId: string): Promise { /* noop */ }
+ async getMetadata(sessionId: string): Promise {
+ return { agentId: 'a' as AgentId, sessionId, status: 'running' };
+ }
+}
+
+describe('MultiAdapter session lifecycle', () => {
+ it('cleans sessionAdapterMap when a session ends naturally', async () => {
+ // Regression: entries were only removed on explicit kill(); sessions
+ // that crashed or exited on their own leaked map entries forever.
+ const acp = new FakeAdapter();
+ const multi = new MultiAdapter(acp);
+ const meta = await multi.spawn({ role: 'worker', cwd: '/tmp' });
+ const map = (multi as any).sessionAdapterMap as Map;
+ expect(map.has(meta.sessionId)).toBe(true);
+
+ // Simulate natural session end (what AcpAdapter fires on process exit)
+ acp.onSessionEnd?.(meta.sessionId, {});
+ expect(map.has(meta.sessionId)).toBe(false);
+ });
+
+ it('preserves externally-assigned onSessionEnd handlers', async () => {
+ const acp = new FakeAdapter();
+ const multi = new MultiAdapter(acp);
+ const meta = await multi.spawn({ role: 'worker', cwd: '/tmp' });
+
+ // Gateway-style assignment AFTER MultiAdapter construction
+ const seen: string[] = [];
+ acp.onSessionEnd = (sessionId) => { seen.push(sessionId); };
+
+ acp.onSessionEnd?.(meta.sessionId, {});
+ expect(seen).toEqual([meta.sessionId]);
+ expect(((multi as any).sessionAdapterMap as Map).has(meta.sessionId)).toBe(false);
+ });
+});
diff --git a/packages/server/tests/comms/message-store.test.ts b/packages/server/tests/comms/message-store.test.ts
index 0b92e87a..22528de3 100644
--- a/packages/server/tests/comms/message-store.test.ts
+++ b/packages/server/tests/comms/message-store.test.ts
@@ -159,4 +159,26 @@ describe('MessageStore', () => {
const afterRead = store.getUnreadDMs('worker-1');
expect(afterRead.length).toBe(0);
});
+
+ it('appendDM writes canonical dm: channel with recipient set', () => {
+ // Regression: bare channel 'dm' leaked into the main chat feed
+ // (the dm: prefix filter missed it) and was invisible to the
+ // per-agent DM panel (which queries channel 'dm:')
+ const msg = store.appendDM('lead', 'worker-9', 'hello');
+ expect(msg.channel).toBe('dm:worker-9');
+ expect(msg.recipient).toBe('worker-9');
+ });
+
+ it('getUnreadDMs finds both canonical and legacy bare-dm rows', () => {
+ // Legacy row shape (pre-migration): channel 'dm' + recipient
+ store.createMessage({
+ parentId: null, taskId: null, authorType: 'agent', authorId: 'lead',
+ content: 'legacy format', metadata: null, channel: 'dm', recipient: 'worker-2',
+ });
+ // Canonical row shape
+ store.appendDM('director', 'worker-2', 'canonical format');
+
+ const unread = store.getUnreadDMs('worker-2');
+ expect(unread.map(m => m.content).sort()).toEqual(['canonical format', 'legacy format']);
+ });
});
diff --git a/packages/server/tests/orchestrator/suspended.test.ts b/packages/server/tests/orchestrator/suspended.test.ts
index 65e52b47..1e173656 100644
--- a/packages/server/tests/orchestrator/suspended.test.ts
+++ b/packages/server/tests/orchestrator/suspended.test.ts
@@ -93,9 +93,43 @@ describe('Orchestrator suspended agents', () => {
const purged = store.purgeOfflineAgents();
expect(purged).toBe(1);
-
+
const remaining = store.listAgents(true);
expect(remaining).toHaveLength(1);
expect(remaining[0].status).toBe('retired');
});
+
+ it('purgeOfflineAgents keeps resumable agents and agents with in-flight tasks', () => {
+ // Resumable: has a saved session — purging it would lose the resume
+ store.insertAgent({
+ id: 'agent-resumable' as AgentId,
+ role: 'worker', runtime: 'acp', acpSessionId: 'sess-keep',
+ status: 'hibernated', currentSpecId: null, costAccumulated: 0, lastHeartbeat: null,
+ });
+ // Has a running task assigned — purging would orphan the task
+ store.insertAgent({
+ id: 'agent-busy-task' as AgentId,
+ role: 'worker', runtime: 'acp', acpSessionId: null,
+ status: 'hibernated', currentSpecId: null, costAccumulated: 0, lastHeartbeat: null,
+ });
+ store.insertTask({
+ id: 'task-inflight' as any, specId: null, title: 'wip', description: '',
+ state: 'running', role: 'worker', dependsOn: [], priority: 0,
+ assignedAgent: 'agent-busy-task' as AgentId, acpSessionId: null,
+ createdAt: new Date().toISOString(), updatedAt: new Date().toISOString(),
+ });
+ // Truly dead: no session, no tasks
+ store.insertAgent({
+ id: 'agent-dead' as AgentId,
+ role: 'worker', runtime: 'acp', acpSessionId: null,
+ status: 'hibernated', currentSpecId: null, costAccumulated: 0, lastHeartbeat: null,
+ });
+
+ const purged = store.purgeOfflineAgents();
+ expect(purged).toBe(1);
+ const ids = store.listAgents(true).map(a => a.id);
+ expect(ids).toContain('agent-resumable');
+ expect(ids).toContain('agent-busy-task');
+ expect(ids).not.toContain('agent-dead');
+ });
});
diff --git a/packages/server/tests/storage/messagelog.test.ts b/packages/server/tests/storage/messagelog.test.ts
index a12e486c..b89212a9 100644
--- a/packages/server/tests/storage/messagelog.test.ts
+++ b/packages/server/tests/storage/messagelog.test.ts
@@ -88,13 +88,18 @@ describe('MessageStore (channel & DM)', () => {
expect(ms.getUnreadDMs('lead-1')).toEqual([]);
});
- it('excludes already-read DMs after markRead', () => {
+ it('excludes already-read DMs after markRead', async () => {
ms.appendDM('worker-1', 'lead-1', 'first');
ms.markRead('lead-1');
const unread1 = ms.getUnreadDMs('lead-1');
expect(unread1).toHaveLength(0);
+ // Timestamps have millisecond granularity — a message created in the
+ // same ms as markRead is indistinguishable from an already-read one.
+ // Step past the boundary so the test is deterministic.
+ await new Promise(r => setTimeout(r, 2));
+
// New message after markRead should show up
ms.appendDM('worker-1', 'lead-1', 'second');
const unread2 = ms.getUnreadDMs('lead-1');
diff --git a/packages/server/tests/storage/sqlite.test.ts b/packages/server/tests/storage/sqlite.test.ts
index 0bf08cc5..23b05621 100644
--- a/packages/server/tests/storage/sqlite.test.ts
+++ b/packages/server/tests/storage/sqlite.test.ts
@@ -95,6 +95,24 @@ describe('SqliteStore', () => {
expect(store.getAgent('agent-test' as AgentId)!.status).toBe('busy');
});
+ it('migrates legacy DM message formats to canonical dm:', () => {
+ // Legacy shape A: bare channel 'dm' with recipient
+ store.rawClient.prepare(
+ `INSERT INTO messages (id, author_type, author_id, content, channel, recipient, created_at) VALUES (?, 'agent', 'lead', 'legacy A', 'dm', 'worker-7', ?)`
+ ).run('msg-legacy-a', new Date().toISOString());
+ // Legacy shape B: channel 'dm:' without recipient
+ store.rawClient.prepare(
+ `INSERT INTO messages (id, author_type, author_id, content, channel, created_at) VALUES (?, 'system', NULL, 'legacy B', 'dm:worker-8', ?)`
+ ).run('msg-legacy-b', new Date().toISOString());
+ store.close();
+
+ // Reopen — constructor migration normalizes both shapes
+ store = new SqliteStore(join(tmpDir, 'test.sqlite'));
+ const rows = store.rawClient.prepare(`SELECT id, channel, recipient FROM messages ORDER BY id`).all() as Array<{ id: string; channel: string; recipient: string | null }>;
+ expect(rows.find(r => r.id === 'msg-legacy-a')).toMatchObject({ channel: 'dm:worker-7', recipient: 'worker-7' });
+ expect(rows.find(r => r.id === 'msg-legacy-b')).toMatchObject({ channel: 'dm:worker-8', recipient: 'worker-8' });
+ });
+
it('persists and reads back agent model and runtime name', () => {
// Regression: agents.model was written via raw SQL but missing from the
// drizzle schema, so select() never returned it and the UI showed the
diff --git a/packages/web/src/lib/types.ts b/packages/web/src/lib/types.ts
index 494169aa..3208955c 100644
--- a/packages/web/src/lib/types.ts
+++ b/packages/web/src/lib/types.ts
@@ -71,6 +71,8 @@ export interface ChatMessage {
attachments?: Array<{ url: string; filename: string; mimeType: string; size: number }> | null;
/** DM/channel routing: 'dm:' identifies the recipient of an agent DM */
channel?: string | null;
+ /** DM recipient agent ID */
+ recipient?: string | null;
channelId?: string | null;
createdAt: string;
updatedAt: string | null;
diff --git a/packages/web/src/pages/Chat.tsx b/packages/web/src/pages/Chat.tsx
index ec12288c..319be5e1 100644
--- a/packages/web/src/pages/Chat.tsx
+++ b/packages/web/src/pages/Chat.tsx
@@ -151,11 +151,8 @@ const MessageBubble = memo(function MessageBubble({ msg, messages, replyCountMap
{msg.authorId}
)}
{(() => {
- // Agent DMs store the recipient in `channel` ('dm:');
- // `channelId` is only set by external bridges
- const dmChannel = msg.channel?.startsWith('dm:') ? msg.channel : msg.channelId?.startsWith('dm:') ? msg.channelId : null;
- if (!dmChannel) return null;
- const recipient = dmChannel.slice(3);
+ const recipient = dmRecipient(msg);
+ if (recipient === '?' || (!msg.channel?.startsWith('dm:') && msg.channel !== 'dm' && !msg.channelId?.startsWith('dm:') && !msg.recipient)) return null;
return (
→ {recipient}
@@ -283,14 +280,17 @@ export function ToolCallCard({ tc, level }: { tc: ToolCallState; level: 'summary
);
}
-/** True for agent↔agent DM traffic (sender in authorId, recipient in channel). */
+/** True for agent↔agent DM traffic (sender in authorId, recipient in channel/recipient). */
function isAgentDm(m: ChatMessage): boolean {
- return m.authorType === 'agent' && !!(m.channel?.startsWith('dm:') || m.channelId?.startsWith('dm:'));
+ return m.authorType === 'agent'
+ && !!(m.channel?.startsWith('dm:') || m.channel === 'dm' || m.channelId?.startsWith('dm:') || m.recipient);
}
function dmRecipient(m: ChatMessage): string {
- const ch = m.channel?.startsWith('dm:') ? m.channel : m.channelId ?? '';
- return ch.startsWith('dm:') ? ch.slice(3) : '?';
+ if (m.channel?.startsWith('dm:')) return m.channel.slice(3);
+ if (m.recipient) return m.recipient;
+ if (m.channelId?.startsWith('dm:')) return m.channelId.slice(3);
+ return '?';
}
/** Collapsed run of consecutive agent↔agent DMs — one line, expandable. */
From dc6de7faebc4b5791d50f04595d91e71f144399c Mon Sep 17 00:00:00 2001
From: Justin Chu
Date: Fri, 12 Jun 2026 08:07:52 -0700
Subject: [PATCH 5/6] feat: replace streamed text with the final message at
turn end
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
Streaming dedup upgraded from drop-the-duplicate to replace-on-final:
deltas stream live as before, but when the authoritative assistant.message
arrives the UI swaps the accumulated text for it. This self-corrects any
delta-level glitches — duplicated fragments inside the stream (seen with
stacked event handlers before 31612b4, and possible upstream), dropped ws
frames, out-of-order delivery — instead of trusting the accumulation.
- CopilotSdkAdapter: after a streamed turn, assistant.message is forwarded
as synthetic 'assistant.message_final' (previously dropped)
- mapper/gateway: message_final broadcasts agent:stream with replace: true
- web useAgents: streamed text chunks are marked ephemeral and the per-turn
output offset is tracked; a replace event swaps them for the final text
while preserving interleaved tool/thinking chunks
- ACP runtimes are unaffected (never send replace); main chat already had
these semantics via the persisted chat:message
Co-Authored-By: Claude Fable 5
---
.../server/src/agents/CopilotSdkAdapter.ts | 31 +++++++++-----
.../src/agents/copilotSdkEventMapper.ts | 7 ++++
packages/server/src/cli/gateway.ts | 2 +-
.../tests/agents/copilot-sdk-adapter.test.ts | 42 +++++++++++++------
packages/web/src/hooks/useAgents.tsx | 27 +++++++++++-
packages/web/src/lib/ws.ts | 2 +-
6 files changed, 83 insertions(+), 28 deletions(-)
diff --git a/packages/server/src/agents/CopilotSdkAdapter.ts b/packages/server/src/agents/CopilotSdkAdapter.ts
index 88cbc7bf..2c6069c6 100644
--- a/packages/server/src/agents/CopilotSdkAdapter.ts
+++ b/packages/server/src/agents/CopilotSdkAdapter.ts
@@ -998,15 +998,18 @@ export class CopilotSdkAdapter extends AgentAdapter {
*/
/**
* Normalize a Copilot SDK event into an ACP-style SessionUpdate for the
- * session's onOutputChunk (lead chat streaming), and decide whether the raw
- * event should still be forwarded to onOutput (agent:stream broadcast).
+ * session's onOutputChunk (lead chat streaming), and decide what to forward
+ * to onOutput (agent:stream broadcast).
*
- * Returns false for full-content events whose text was already streamed as
- * deltas this turn — forwarding those would repeat every word in the UI.
+ * Returns the event to forward (possibly rewritten), or null to drop it.
+ * When the full assistant.message arrives after deltas were streamed, it is
+ * rewritten to the synthetic 'assistant.message_final' — the UI replaces the
+ * accumulated stream with it, so any delta-level glitches (duplicated or
+ * dropped fragments) self-correct at turn end.
*/
- private processStreamEvent(agentSession: CopilotAgentSession, event: { type: string; data?: any }): boolean {
+ private processStreamEvent(agentSession: CopilotAgentSession, event: { type: string; data?: any }): { type: string; data?: any } | null {
let update: Record | null = null;
- let forward = true;
+ let forward: { type: string; data?: any } | null = event;
switch (event.type) {
case 'assistant.message_delta':
agentSession.sawTextDelta = true;
@@ -1017,11 +1020,17 @@ export class CopilotSdkAdapter extends AgentAdapter {
update = { sessionUpdate: 'agent_thought_chunk', content: { type: 'text', text: event.data?.deltaContent ?? '' } };
break;
case 'assistant.message':
- if (agentSession.sawTextDelta) forward = false;
- else if (event.data?.content) update = { sessionUpdate: 'agent_message_chunk', content: { type: 'text', text: event.data.content } };
+ if (agentSession.sawTextDelta) {
+ forward = event.data?.content
+ ? { type: 'assistant.message_final', data: { content: event.data.content } }
+ : null;
+ } else if (event.data?.content) {
+ update = { sessionUpdate: 'agent_message_chunk', content: { type: 'text', text: event.data.content } };
+ }
break;
case 'assistant.reasoning':
- if (agentSession.sawReasoningDelta) forward = false;
+ // Thinking blocks don't need replace fidelity — drop the duplicate
+ if (agentSession.sawReasoningDelta) forward = null;
else if (event.data?.content) update = { sessionUpdate: 'agent_thought_chunk', content: { type: 'text', text: event.data.content } };
break;
case 'tool.execution_start':
@@ -1196,7 +1205,7 @@ export class CopilotSdkAdapter extends AgentAdapter {
const forward = this.processStreamEvent(agentSession, event as { type: string; data?: any });
if (forward && this.onOutput) {
- try { this.onOutput(aid, event); } catch { /* */ }
+ try { this.onOutput(aid, forward as SessionEvent); } catch { /* */ }
}
});
@@ -1395,7 +1404,7 @@ export class CopilotSdkAdapter extends AgentAdapter {
const forward = this.processStreamEvent(agentSession, event as { type: string; data?: any });
if (forward && this.onOutput) {
- try { this.onOutput(aid, event); } catch { /* */ }
+ try { this.onOutput(aid, forward as SessionEvent); } catch { /* */ }
}
});
diff --git a/packages/server/src/agents/copilotSdkEventMapper.ts b/packages/server/src/agents/copilotSdkEventMapper.ts
index f4a453ef..b98e4631 100644
--- a/packages/server/src/agents/copilotSdkEventMapper.ts
+++ b/packages/server/src/agents/copilotSdkEventMapper.ts
@@ -2,6 +2,8 @@ export interface StreamBroadcast {
delta: string;
contentType: 'text' | 'thinking' | 'tool_call' | 'tool_result';
toolName?: string;
+ /** Replace the text streamed this turn instead of appending (final message). */
+ replace?: boolean;
}
/** Map a Copilot SDK session event to a WebSocket stream broadcast payload */
@@ -16,6 +18,11 @@ export function mapCopilotSdkEvent(event: { type: string; data?: any }): StreamB
case 'assistant.message':
// Complete message — use as text if no delta events were streamed
return event.data?.content ? { delta: event.data.content, contentType: 'text' } : null;
+ case 'assistant.message_final':
+ // Synthetic (CopilotSdkAdapter): authoritative full text after a streamed
+ // turn — the UI replaces the accumulated deltas with it, so delta-level
+ // glitches (duplicated/dropped fragments) self-correct at turn end
+ return event.data?.content ? { delta: event.data.content, contentType: 'text', replace: true } : null;
case 'assistant.reasoning':
// Complete reasoning block
return event.data?.content ? { delta: event.data.content, contentType: 'thinking' } : null;
diff --git a/packages/server/src/cli/gateway.ts b/packages/server/src/cli/gateway.ts
index da12151f..02d145be 100644
--- a/packages/server/src/cli/gateway.ts
+++ b/packages/server/src/cli/gateway.ts
@@ -350,7 +350,7 @@ export async function startGateway(deps: GatewayDeps): Promise {
if (!mapped) return;
for (const wsServer of wsServers.values()) {
if (mapped.delta) {
- wsServer.broadcast({ type: 'agent:stream', agentId, delta: mapped.delta, contentType: mapped.contentType, toolName: mapped.toolName });
+ wsServer.broadcast({ type: 'agent:stream', agentId, delta: mapped.delta, contentType: mapped.contentType, toolName: mapped.toolName, ...(mapped.replace ? { replace: true } : {}) });
}
}
};
diff --git a/packages/server/tests/agents/copilot-sdk-adapter.test.ts b/packages/server/tests/agents/copilot-sdk-adapter.test.ts
index ee5548f1..af631c30 100644
--- a/packages/server/tests/agents/copilot-sdk-adapter.test.ts
+++ b/packages/server/tests/agents/copilot-sdk-adapter.test.ts
@@ -249,19 +249,22 @@ describe('CopilotSdkAdapter', () => {
const process = (session: any, event: any) =>
(adapter as any).processStreamEvent(session, event);
- it('drops the full assistant.message after deltas were streamed (no word duplication)', () => {
+ it('rewrites the full assistant.message to message_final after deltas were streamed', () => {
const { session, chunks } = makeSession();
- expect(process(session, { type: 'assistant.message_delta', data: { deltaContent: 'hello ' } })).toBe(true);
- expect(process(session, { type: 'assistant.message_delta', data: { deltaContent: 'world' } })).toBe(true);
- // Full message repeats the streamed text — must NOT be forwarded
- expect(process(session, { type: 'assistant.message', data: { content: 'hello world' } })).toBe(false);
+ expect(process(session, { type: 'assistant.message_delta', data: { deltaContent: 'hello ' } })).toBeTruthy();
+ expect(process(session, { type: 'assistant.message_delta', data: { deltaContent: 'world' } })).toBeTruthy();
+ // Full message repeats the streamed text — forwarded as a REPLACE event
+ // so the UI swaps the accumulated stream for the authoritative text
+ const fwd = process(session, { type: 'assistant.message', data: { content: 'hello world' } });
+ expect(fwd).toEqual({ type: 'assistant.message_final', data: { content: 'hello world' } });
const texts = chunks.filter(c => c.sessionUpdate === 'agent_message_chunk').map(c => c.content.text);
expect(texts).toEqual(['hello ', 'world']);
});
- it('forwards the full assistant.message when nothing was streamed', () => {
+ it('forwards the full assistant.message unchanged when nothing was streamed', () => {
const { session, chunks } = makeSession();
- expect(process(session, { type: 'assistant.message', data: { content: 'hello world' } })).toBe(true);
+ const event = { type: 'assistant.message', data: { content: 'hello world' } };
+ expect(process(session, event)).toBe(event);
expect(chunks).toEqual([{ sessionUpdate: 'agent_message_chunk', content: { type: 'text', text: 'hello world' } }]);
});
@@ -269,14 +272,15 @@ describe('CopilotSdkAdapter', () => {
const { session } = makeSession();
process(session, { type: 'assistant.message_delta', data: { deltaContent: 'a' } });
process(session, { type: 'session.idle' });
- // Next turn answers without streaming — full message must pass through
- expect(process(session, { type: 'assistant.message', data: { content: 'next turn' } })).toBe(true);
+ // Next turn answers without streaming — full message passes through unchanged
+ const event = { type: 'assistant.message', data: { content: 'next turn' } };
+ expect(process(session, event)).toBe(event);
});
- it('dedupes reasoning the same way as text', () => {
+ it('drops duplicate reasoning (thinking does not need replace fidelity)', () => {
const { session } = makeSession();
process(session, { type: 'assistant.reasoning_delta', data: { deltaContent: 'thinking…' } });
- expect(process(session, { type: 'assistant.reasoning', data: { content: 'thinking…' } })).toBe(false);
+ expect(process(session, { type: 'assistant.reasoning', data: { content: 'thinking…' } })).toBeNull();
});
it('synthesizes ACP-style tool call updates for onOutputChunk', () => {
@@ -289,8 +293,20 @@ describe('CopilotSdkAdapter', () => {
it('works without onOutputChunk wired (workers)', () => {
const session: any = {};
- expect(process(session, { type: 'assistant.message_delta', data: { deltaContent: 'x' } })).toBe(true);
- expect(process(session, { type: 'assistant.message', data: { content: 'x' } })).toBe(false);
+ expect(process(session, { type: 'assistant.message_delta', data: { deltaContent: 'x' } })).toBeTruthy();
+ expect(process(session, { type: 'assistant.message', data: { content: 'x' } }))
+ .toEqual({ type: 'assistant.message_final', data: { content: 'x' } });
+ });
+
+ it('mapper marks message_final as a replace broadcast', async () => {
+ const { mapCopilotSdkEvent } = await import('../../src/agents/copilotSdkEventMapper.js');
+ expect(mapCopilotSdkEvent({ type: 'assistant.message_final', data: { content: 'final text' } }))
+ .toEqual({ delta: 'final text', contentType: 'text', replace: true });
+ // Plain deltas and unstreamed full messages stay append-mode
+ expect(mapCopilotSdkEvent({ type: 'assistant.message_delta', data: { deltaContent: 'x' } }))
+ .toEqual({ delta: 'x', contentType: 'text' });
+ expect(mapCopilotSdkEvent({ type: 'assistant.message', data: { content: 'y' } }))
+ .toEqual({ delta: 'y', contentType: 'text' });
});
});
});
diff --git a/packages/web/src/hooks/useAgents.tsx b/packages/web/src/hooks/useAgents.tsx
index 3ab1e7e9..66b2dc11 100644
--- a/packages/web/src/hooks/useAgents.tsx
+++ b/packages/web/src/hooks/useAgents.tsx
@@ -12,6 +12,8 @@ export interface StreamChunk {
content: string;
contentType?: ContentType;
toolName?: string;
+ /** Streamed text of the in-flight turn — replaced by the final message. */
+ ephemeral?: boolean;
}
export interface AgentContextValue {
@@ -30,6 +32,8 @@ export function AgentProvider({ children }: { children: ReactNode }) {
const agentOutputsRef = useRef(new Map());
const agentStreamChunksRef = useRef(new Map());
const dmMessagesRef = useRef(new Map());
+ /** Per-agent output offset where the current turn's streamed text began. */
+ const turnStartRef = useRef(new Map());
const dirtyRef = useRef(false);
const rafRef = useRef(null);
const { subscribe } = useWsEventBus();
@@ -57,6 +61,7 @@ export function AgentProvider({ children }: { children: ReactNode }) {
agentOutputsRef.current.clear();
agentStreamChunksRef.current.clear();
dmMessagesRef.current.clear();
+ turnStartRef.current.clear();
setAgentOutputs(new Map());
setAgentStreamChunks(new Map());
setDmMessages(new Map());
@@ -68,9 +73,27 @@ export function AgentProvider({ children }: { children: ReactNode }) {
mutateAgents();
} else if (event.type === 'agent:stream') {
const prev = agentOutputsRef.current.get(event.agentId) ?? '';
- agentOutputsRef.current.set(event.agentId, prev + event.delta);
const prevChunks = agentStreamChunksRef.current.get(event.agentId) ?? [];
- agentStreamChunksRef.current.set(event.agentId, [...prevChunks, { content: event.delta, contentType: event.contentType ?? 'text', toolName: (event as any).toolName }].slice(-MAX_CHUNKS));
+ const isText = !event.contentType || event.contentType === 'text';
+ if (event.replace && isText) {
+ // Final authoritative text for the turn — replace the streamed
+ // deltas with it so any delta-level glitches self-correct
+ const start = turnStartRef.current.get(event.agentId);
+ const base = start !== undefined ? prev.slice(0, start) : prev;
+ agentOutputsRef.current.set(event.agentId, base + event.delta);
+ agentStreamChunksRef.current.set(event.agentId,
+ [...prevChunks.filter(c => !c.ephemeral), { content: event.delta, contentType: 'text' as const }].slice(-MAX_CHUNKS));
+ turnStartRef.current.delete(event.agentId);
+ } else {
+ // Mark where this turn's text began so a later replace knows
+ // how much of the output to swap out
+ if (isText && !turnStartRef.current.has(event.agentId)) {
+ turnStartRef.current.set(event.agentId, prev.length);
+ }
+ agentOutputsRef.current.set(event.agentId, prev + event.delta);
+ agentStreamChunksRef.current.set(event.agentId,
+ [...prevChunks, { content: event.delta, contentType: event.contentType ?? 'text', toolName: (event as any).toolName, ...(isText ? { ephemeral: true } : {}) }].slice(-MAX_CHUNKS));
+ }
scheduleFlush();
} else if (event.type === 'tool:event') {
const agentId = event.agentId;
diff --git a/packages/web/src/lib/ws.ts b/packages/web/src/lib/ws.ts
index c063c770..65fd025f 100644
--- a/packages/web/src/lib/ws.ts
+++ b/packages/web/src/lib/ws.ts
@@ -9,7 +9,7 @@ export type WsEvent =
| { type: 'task:comment'; task_id: string; message: ChatMessage; project?: string }
| { type: 'display:config'; config: DisplayConfig }
| { type: 'state:update'; stats: Record }
- | { type: 'agent:stream'; agentId: string; delta: string; contentType: 'text' | 'thinking' | 'tool_call' | 'tool_result'; toolName?: string }
+ | { type: 'agent:stream'; agentId: string; delta: string; contentType: 'text' | 'thinking' | 'tool_call' | 'tool_result'; toolName?: string; replace?: boolean }
| { type: 'tool:event'; toolName: string; agentId: string; input: unknown; output: unknown; status: string; durationMs?: number; error?: string }
| { type: 'dm:message'; project: string; message: { id: string; channel: string; content: string; authorType: string; authorId: string | null; createdAt: string; [k: string]: unknown } };
From 8e2fa6e00f2a68d35a695ff5b086c5ee58d771cb Mon Sep 17 00:00:00 2001
From: Justin Chu
Date: Fri, 12 Jun 2026 08:11:53 -0700
Subject: [PATCH 6/6] =?UTF-8?q?fix:=20address=20PR=20#4=20review=20?=
=?UTF-8?q?=E2=80=94=20message=20window=20order,=20DM=20group=20reply=20co?=
=?UTF-8?q?ntext?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
- GET /messages: listMessages returns newest-first, but slice(-limit)
kept the OLDEST entries of the fetched window — chats longer than
`limit` never showed their most recent messages. Keep the newest
`limit` and reverse to ascending. Also over-fetch (limit + 200) so DM
filtering can't underfill the page. Route-level regression test pins
newest-window + ascending order.
- AgentDmGroup: pass the full visible message list to MessageBubble so
reply previews resolve parents outside the collapsed group.
- useChat: drop the `as any` on dm:message ingestion — use the typed
event payload.
Co-Authored-By: Claude Fable 5
---
packages/server/src/api/routes/messages.ts | 9 ++-
.../server/tests/api/messages-window.test.ts | 78 +++++++++++++++++++
packages/web/src/hooks/useChat.tsx | 4 +-
packages/web/src/pages/Chat.tsx | 8 +-
4 files changed, 93 insertions(+), 6 deletions(-)
create mode 100644 packages/server/tests/api/messages-window.test.ts
diff --git a/packages/server/src/api/routes/messages.ts b/packages/server/src/api/routes/messages.ts
index c4d764e9..6d0bcdb2 100644
--- a/packages/server/src/api/routes/messages.ts
+++ b/packages/server/src/api/routes/messages.ts
@@ -152,10 +152,15 @@ export async function handleMessageRoutes(
// Agent↔agent DM traffic is excluded from the main chat by default;
// include_agent_dms=true lets the UI show it (collapsed by default)
const includeAgentDms = url.searchParams.get('include_agent_dms') === 'true';
- const allMsgs = fd.messages?.listMessages({ taskId, limit: limit + 50, authorTypes }) ?? [];
+ // Over-fetch: DM filtering below can discard a large share of the window
+ const allMsgs = fd.messages?.listMessages({ taskId, limit: limit + 200, authorTypes }) ?? [];
const mainChatMsgs = allMsgs.filter(m =>
!m.channel?.startsWith('dm:') || (includeAgentDms && m.authorType === 'agent'));
- json(200, mainChatMsgs.slice(-limit).reverse());
+ // listMessages returns newest-first: keep the NEWEST `limit` entries,
+ // then reverse to the ascending order the chat renders in.
+ // (slice(-limit) used to keep the oldest entries — chats longer than
+ // `limit` never showed their most recent messages.)
+ json(200, mainChatMsgs.slice(0, limit).reverse());
return true;
}
diff --git a/packages/server/tests/api/messages-window.test.ts b/packages/server/tests/api/messages-window.test.ts
new file mode 100644
index 00000000..b92c395b
--- /dev/null
+++ b/packages/server/tests/api/messages-window.test.ts
@@ -0,0 +1,78 @@
+import { describe, it, expect, beforeEach, afterEach } from 'vitest';
+import { createHttpServer } from '../../src/api/HttpServer.js';
+import { Flightdeck } from '../../src/facade.js';
+import type { ProjectManager } from '../../src/projects/ProjectManager.js';
+import { existsSync, rmSync } from 'node:fs';
+import { join } from 'node:path';
+import { homedir } from 'node:os';
+import http from 'node:http';
+
+/**
+ * Regression: GET /messages used slice(-limit) on a newest-first list, so
+ * chats longer than `limit` returned the OLDEST window and never showed
+ * the most recent messages.
+ */
+describe('GET /messages window', () => {
+ const projectName = `test-msg-window-${Date.now()}`;
+ let fd: Flightdeck;
+ let server: http.Server;
+ let port: number;
+
+ beforeEach(async () => {
+ fd = new Flightdeck(projectName);
+ const projectManager = {
+ list: () => [projectName],
+ get: (name: string) => (name === projectName ? fd : null),
+ create: () => {},
+ delete: () => true,
+ closeAll: () => {},
+ } as unknown as ProjectManager;
+ server = createHttpServer({
+ projectManager,
+ leadManagers: new Map(),
+ agentManagers: new Map(),
+ wsServers: new Map(),
+ webhookNotifiers: new Map(),
+ cronStores: new Map(),
+ port: 0,
+ corsOrigin: '*',
+ } as any);
+ port = await new Promise(resolve => {
+ server.listen(0, '127.0.0.1', () => resolve((server.address() as any).port));
+ });
+ });
+
+ afterEach(() => {
+ server.close();
+ fd.close();
+ const projDir = join(homedir(), '.flightdeck', 'v2', 'projects', projectName);
+ if (existsSync(projDir)) rmSync(projDir, { recursive: true, force: true });
+ });
+
+ it('returns the NEWEST `limit` messages in ascending order', async () => {
+ // 120 messages with strictly increasing ids/timestamps
+ const base = Date.now();
+ for (let i = 0; i < 120; i++) {
+ fd.messages!.createMessage({
+ id: `msg-${String(i).padStart(3, '0')}`,
+ parentId: null, taskId: null, authorType: 'user', authorId: 'u',
+ content: `message ${i}`, metadata: null,
+ });
+ // Distinct createdAt per row (createMessage stamps now): nudge clock
+ // by overwriting created_at deterministically
+ fd.sqlite.rawClient.prepare(`UPDATE messages SET created_at = ? WHERE id = ?`)
+ .run(new Date(base + i * 1000).toISOString(), `msg-${String(i).padStart(3, '0')}`);
+ }
+
+ const res = await fetch(`http://127.0.0.1:${port}/api/projects/${projectName}/messages?limit=100`);
+ expect(res.status).toBe(200);
+ const msgs = await res.json() as Array<{ id: string; content: string }>;
+ expect(msgs).toHaveLength(100);
+ // Must be the newest 100 (20..119), not the oldest
+ expect(msgs[0].content).toBe('message 20');
+ expect(msgs[msgs.length - 1].content).toBe('message 119');
+ // Ascending render order
+ const indices = msgs.map(m => parseInt(m.content.replace('message ', ''), 10));
+ expect([...indices].sort((a, b) => a - b)).toEqual(indices);
+ });
+});
diff --git a/packages/web/src/hooks/useChat.tsx b/packages/web/src/hooks/useChat.tsx
index 8818334d..76c125f0 100644
--- a/packages/web/src/hooks/useChat.tsx
+++ b/packages/web/src/hooks/useChat.tsx
@@ -171,7 +171,9 @@ export function ChatProvider({ children }: { children: ReactNode }) {
case 'dm:message': {
// Live agent↔agent DMs — shown in main chat unless turned off
if ((displayConfigRef.current.agentMessages ?? 'summary') === 'off') break;
- const dm = (event as any).message;
+ // The ws payload is the persisted message row — same shape the
+ // REST endpoint returns, just typed loosely on the event union
+ const dm = event.message as unknown as ChatMessage;
if (!dm?.id) break;
setWsMessages(prev => {
if (prev.some(m => m.id === dm.id)) return prev;
diff --git a/packages/web/src/pages/Chat.tsx b/packages/web/src/pages/Chat.tsx
index 319be5e1..27c2440b 100644
--- a/packages/web/src/pages/Chat.tsx
+++ b/packages/web/src/pages/Chat.tsx
@@ -294,8 +294,10 @@ function dmRecipient(m: ChatMessage): string {
}
/** Collapsed run of consecutive agent↔agent DMs — one line, expandable. */
-function AgentDmGroup({ msgs, replyCountMap, onReply, agents }: {
+function AgentDmGroup({ msgs, allMessages, replyCountMap, onReply, agents }: {
msgs: ChatMessage[];
+ /** Full visible message list — reply previews may reference messages outside the group */
+ allMessages: ChatMessage[];
replyCountMap?: Map;
onReply: (m: ChatMessage) => void;
agents?: Array<{ id: string; role: string; runtime?: string; runtimeName?: string; model?: string; status?: string }>;
@@ -324,7 +326,7 @@ function AgentDmGroup({ msgs, replyCountMap, onReply, agents }: {
{open && (
{msgs.map(m => (
-
+
))}
)}
@@ -823,7 +825,7 @@ export default function Chat() {
)}
{renderItems.map(item => item.kind === 'dmGroup' ? (
-
+
) : (
))}