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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 14 additions & 0 deletions packages/server/src/agents/AcpAdapter.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import { spawn as cpSpawn, type ChildProcess } from 'node:child_process';
import { Readable, Writable } from 'node:stream';
import { randomUUID } from 'node:crypto';
import { dirname, resolve, join } from 'node:path';

Check warning on line 4 in packages/server/src/agents/AcpAdapter.ts

View workflow job for this annotation

GitHub Actions / lint

'join' is defined but never used. Allowed unused vars must match /^_/u
import { fileURLToPath } from 'node:url';
import { existsSync } from 'node:fs';

Expand Down Expand Up @@ -233,7 +233,7 @@
session.onOutputChunk?.(update);
self.onAnySessionOutput?.(session.agentId, update);
if (self.onToolCall) {
let toolName = (update as any).title ?? (update as any).name ?? (update as any).toolName ?? 'unknown';

Check warning on line 236 in packages/server/src/agents/AcpAdapter.ts

View workflow job for this annotation

GitHub Actions / lint

Unexpected any. Specify a different type

Check warning on line 236 in packages/server/src/agents/AcpAdapter.ts

View workflow job for this annotation

GitHub Actions / lint

Unexpected any. Specify a different type

Check failure on line 236 in packages/server/src/agents/AcpAdapter.ts

View workflow job for this annotation

GitHub Actions / lint

'toolName' is never reassigned. Use 'const' instead
self.onToolCall(session.agentId, { toolName, status: 'completed' });
}
break;
Expand Down Expand Up @@ -421,6 +421,20 @@
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)}`;
Expand Down
4 changes: 2 additions & 2 deletions packages/server/src/agents/AgentManager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand All @@ -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);
}
Expand Down
113 changes: 102 additions & 11 deletions packages/server/src/agents/CopilotSdkAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,14 @@
cwd: string;
model?: string;
pendingToolCalls?: Map<string, string>;
/** 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;
sawReasoningDelta?: boolean;
}

export class CopilotSdkAdapter extends AgentAdapter {
Expand Down Expand Up @@ -988,6 +996,73 @@
/**
* 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 what to forward
* to onOutput (agent:stream broadcast).
*
* 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 }): { type: string; data?: any } | null {
let update: Record<string, unknown> | null = null;
let forward: { type: string; data?: any } | null = event;
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 = 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':
// 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':
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<AgentMetadata> {
const client = await this.ensureClient();
const aid = (opts.agentId ?? makeAgentId(opts.role, Date.now().toString())) as AgentId;
Expand Down Expand Up @@ -1040,10 +1115,13 @@
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
Expand Down Expand Up @@ -1125,8 +1203,9 @@
}
}

if (this.onOutput) {
try { this.onOutput(aid, event); } catch { /* */ }
const forward = this.processStreamEvent(agentSession, event as { type: string; data?: any });
if (forward && this.onOutput) {
try { this.onOutput(aid, forward as SessionEvent); } catch { /* */ }
}
});

Expand Down Expand Up @@ -1162,16 +1241,20 @@

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<void>((resolve) => {
const handler = (event: SessionEvent) => {
let timer: ReturnType<typeof setTimeout> | undefined;

Check failure on line 1248 in packages/server/src/agents/CopilotSdkAdapter.ts

View workflow job for this annotation

GitHub Actions / lint

'timer' is never reassigned. Use 'const' instead
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);
Expand All @@ -1186,6 +1269,8 @@
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);
Expand Down Expand Up @@ -1255,10 +1340,13 @@
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;
Expand Down Expand Up @@ -1294,8 +1382,6 @@
});
} 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;
Expand All @@ -1315,6 +1401,11 @@
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, forward as SessionEvent); } catch { /* */ }
}
});

return { agentId: aid, sessionId: sessionId, status: 'running' as const };
Expand Down
10 changes: 10 additions & 0 deletions packages/server/src/agents/ModelConfig.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*/
Expand Down
23 changes: 23 additions & 0 deletions packages/server/src/agents/MultiAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
7 changes: 7 additions & 0 deletions packages/server/src/agents/copilotSdkEventMapper.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 */
Expand All @@ -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;
Expand Down
17 changes: 13 additions & 4 deletions packages/server/src/api/routes/messages.ts
Original file line number Diff line number Diff line change
Expand Up @@ -149,9 +149,18 @@ 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;
const allMsgs = fd.messages?.listMessages({ taskId, limit: limit + 50, authorTypes }) ?? [];
const mainChatMsgs = allMsgs.filter(m => !m.channel?.startsWith('dm:'));
json(200, mainChatMsgs.slice(-limit).reverse());
// 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';
// 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'));
// 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;
}

Expand Down Expand Up @@ -227,7 +236,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,
});
}
Expand Down
Loading
Loading