|
|
|
|
@@ -36,6 +36,31 @@ type ToolResultMessage = {
|
|
|
|
|
content?: unknown;
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
type SubagentDeliveryPath = "queued" | "steered" | "direct" | "none";
|
|
|
|
|
|
|
|
|
|
type SubagentAnnounceDeliveryResult = {
|
|
|
|
|
delivered: boolean;
|
|
|
|
|
path: SubagentDeliveryPath;
|
|
|
|
|
error?: string;
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
function summarizeDeliveryError(error: unknown): string {
|
|
|
|
|
if (error instanceof Error) {
|
|
|
|
|
return error.message || "error";
|
|
|
|
|
}
|
|
|
|
|
if (typeof error === "string") {
|
|
|
|
|
return error;
|
|
|
|
|
}
|
|
|
|
|
if (error === undefined || error === null) {
|
|
|
|
|
return "unknown error";
|
|
|
|
|
}
|
|
|
|
|
try {
|
|
|
|
|
return JSON.stringify(error);
|
|
|
|
|
} catch {
|
|
|
|
|
return "error";
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
function extractToolResultText(content: unknown): string {
|
|
|
|
|
if (typeof content === "string") {
|
|
|
|
|
return sanitizeTextContent(content);
|
|
|
|
|
@@ -45,6 +70,9 @@ function extractToolResultText(content: unknown): string {
|
|
|
|
|
text?: unknown;
|
|
|
|
|
output?: unknown;
|
|
|
|
|
content?: unknown;
|
|
|
|
|
result?: unknown;
|
|
|
|
|
error?: unknown;
|
|
|
|
|
summary?: unknown;
|
|
|
|
|
};
|
|
|
|
|
if (typeof obj.text === "string") {
|
|
|
|
|
return sanitizeTextContent(obj.text);
|
|
|
|
|
@@ -55,6 +83,15 @@ function extractToolResultText(content: unknown): string {
|
|
|
|
|
if (typeof obj.content === "string") {
|
|
|
|
|
return sanitizeTextContent(obj.content);
|
|
|
|
|
}
|
|
|
|
|
if (typeof obj.result === "string") {
|
|
|
|
|
return sanitizeTextContent(obj.result);
|
|
|
|
|
}
|
|
|
|
|
if (typeof obj.error === "string") {
|
|
|
|
|
return sanitizeTextContent(obj.error);
|
|
|
|
|
}
|
|
|
|
|
if (typeof obj.summary === "string") {
|
|
|
|
|
return sanitizeTextContent(obj.summary);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
if (!Array.isArray(content)) {
|
|
|
|
|
return "";
|
|
|
|
|
@@ -72,12 +109,41 @@ function extractSubagentOutputText(message: unknown): string {
|
|
|
|
|
return "";
|
|
|
|
|
}
|
|
|
|
|
const role = (message as { role?: unknown }).role;
|
|
|
|
|
const content = (message as { content?: unknown }).content;
|
|
|
|
|
if (role === "assistant") {
|
|
|
|
|
return extractAssistantText(message) ?? "";
|
|
|
|
|
const assistantText = extractAssistantText(message);
|
|
|
|
|
if (assistantText) {
|
|
|
|
|
return assistantText;
|
|
|
|
|
}
|
|
|
|
|
if (typeof content === "string") {
|
|
|
|
|
return sanitizeTextContent(content);
|
|
|
|
|
}
|
|
|
|
|
if (Array.isArray(content)) {
|
|
|
|
|
return (
|
|
|
|
|
extractTextFromChatContent(content, {
|
|
|
|
|
sanitizeText: sanitizeTextContent,
|
|
|
|
|
normalizeText: (text) => text.trim(),
|
|
|
|
|
joinWith: "",
|
|
|
|
|
}) ?? ""
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
return "";
|
|
|
|
|
}
|
|
|
|
|
if (role === "toolResult" || role === "tool") {
|
|
|
|
|
return extractToolResultText((message as ToolResultMessage).content);
|
|
|
|
|
}
|
|
|
|
|
if (typeof content === "string") {
|
|
|
|
|
return sanitizeTextContent(content);
|
|
|
|
|
}
|
|
|
|
|
if (Array.isArray(content)) {
|
|
|
|
|
return (
|
|
|
|
|
extractTextFromChatContent(content, {
|
|
|
|
|
sanitizeText: sanitizeTextContent,
|
|
|
|
|
normalizeText: (text) => text.trim(),
|
|
|
|
|
joinWith: "",
|
|
|
|
|
}) ?? ""
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
return "";
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@@ -314,6 +380,126 @@ async function maybeQueueSubagentAnnounce(params: {
|
|
|
|
|
return "none";
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
function queueOutcomeToDeliveryResult(
|
|
|
|
|
outcome: "steered" | "queued" | "none",
|
|
|
|
|
): SubagentAnnounceDeliveryResult {
|
|
|
|
|
if (outcome === "steered") {
|
|
|
|
|
return {
|
|
|
|
|
delivered: true,
|
|
|
|
|
path: "steered",
|
|
|
|
|
};
|
|
|
|
|
}
|
|
|
|
|
if (outcome === "queued") {
|
|
|
|
|
return {
|
|
|
|
|
delivered: true,
|
|
|
|
|
path: "queued",
|
|
|
|
|
};
|
|
|
|
|
}
|
|
|
|
|
return {
|
|
|
|
|
delivered: false,
|
|
|
|
|
path: "none",
|
|
|
|
|
};
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
async function sendSubagentAnnounceDirectly(params: {
|
|
|
|
|
targetRequesterSessionKey: string;
|
|
|
|
|
triggerMessage: string;
|
|
|
|
|
directIdempotencyKey: string;
|
|
|
|
|
directOrigin?: DeliveryContext;
|
|
|
|
|
requesterIsSubagent: boolean;
|
|
|
|
|
}): Promise<SubagentAnnounceDeliveryResult> {
|
|
|
|
|
try {
|
|
|
|
|
await callGateway({
|
|
|
|
|
method: "agent",
|
|
|
|
|
params: {
|
|
|
|
|
sessionKey: params.targetRequesterSessionKey,
|
|
|
|
|
message: params.triggerMessage,
|
|
|
|
|
deliver: !params.requesterIsSubagent,
|
|
|
|
|
channel: params.requesterIsSubagent ? undefined : params.directOrigin?.channel,
|
|
|
|
|
accountId: params.requesterIsSubagent ? undefined : params.directOrigin?.accountId,
|
|
|
|
|
to: params.requesterIsSubagent ? undefined : params.directOrigin?.to,
|
|
|
|
|
threadId:
|
|
|
|
|
!params.requesterIsSubagent &&
|
|
|
|
|
params.directOrigin?.threadId != null &&
|
|
|
|
|
params.directOrigin.threadId !== ""
|
|
|
|
|
? String(params.directOrigin.threadId)
|
|
|
|
|
: undefined,
|
|
|
|
|
idempotencyKey: params.directIdempotencyKey,
|
|
|
|
|
},
|
|
|
|
|
expectFinal: true,
|
|
|
|
|
timeoutMs: 15_000,
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
return {
|
|
|
|
|
delivered: true,
|
|
|
|
|
path: "direct",
|
|
|
|
|
};
|
|
|
|
|
} catch (err) {
|
|
|
|
|
return {
|
|
|
|
|
delivered: false,
|
|
|
|
|
path: "direct",
|
|
|
|
|
error: summarizeDeliveryError(err),
|
|
|
|
|
};
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
async function deliverSubagentCompletionAnnouncement(params: {
|
|
|
|
|
requesterSessionKey: string;
|
|
|
|
|
announceId?: string;
|
|
|
|
|
triggerMessage: string;
|
|
|
|
|
summaryLine?: string;
|
|
|
|
|
requesterOrigin?: DeliveryContext;
|
|
|
|
|
directOrigin?: DeliveryContext;
|
|
|
|
|
targetRequesterSessionKey: string;
|
|
|
|
|
requesterIsSubagent: boolean;
|
|
|
|
|
expectsCompletionMessage: boolean;
|
|
|
|
|
directIdempotencyKey: string;
|
|
|
|
|
}): Promise<SubagentAnnounceDeliveryResult> {
|
|
|
|
|
// Non-completion mode mirrors historical behavior: try queued/steered delivery first,
|
|
|
|
|
// then (only if not queued) attempt direct delivery.
|
|
|
|
|
if (!params.expectsCompletionMessage) {
|
|
|
|
|
const queueOutcome = await maybeQueueSubagentAnnounce({
|
|
|
|
|
requesterSessionKey: params.requesterSessionKey,
|
|
|
|
|
announceId: params.announceId,
|
|
|
|
|
triggerMessage: params.triggerMessage,
|
|
|
|
|
summaryLine: params.summaryLine,
|
|
|
|
|
requesterOrigin: params.requesterOrigin,
|
|
|
|
|
});
|
|
|
|
|
const queued = queueOutcomeToDeliveryResult(queueOutcome);
|
|
|
|
|
if (queued.delivered) {
|
|
|
|
|
return queued;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Completion-mode uses direct send first so manual spawns can return immediately
|
|
|
|
|
// in the common ready-to-deliver case.
|
|
|
|
|
const direct = await sendSubagentAnnounceDirectly({
|
|
|
|
|
targetRequesterSessionKey: params.targetRequesterSessionKey,
|
|
|
|
|
triggerMessage: params.triggerMessage,
|
|
|
|
|
directIdempotencyKey: params.directIdempotencyKey,
|
|
|
|
|
directOrigin: params.directOrigin,
|
|
|
|
|
requesterIsSubagent: params.requesterIsSubagent,
|
|
|
|
|
});
|
|
|
|
|
if (direct.delivered || !params.expectsCompletionMessage) {
|
|
|
|
|
return direct;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// If completion path failed direct delivery, try queueing as a fallback so the
|
|
|
|
|
// report can still be delivered once the requester session is idle.
|
|
|
|
|
const queueOutcome = await maybeQueueSubagentAnnounce({
|
|
|
|
|
requesterSessionKey: params.requesterSessionKey,
|
|
|
|
|
announceId: params.announceId,
|
|
|
|
|
triggerMessage: params.triggerMessage,
|
|
|
|
|
summaryLine: params.summaryLine,
|
|
|
|
|
requesterOrigin: params.requesterOrigin,
|
|
|
|
|
});
|
|
|
|
|
if (queueOutcome === "steered" || queueOutcome === "queued") {
|
|
|
|
|
return queueOutcomeToDeliveryResult(queueOutcome);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return direct;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
function loadSessionEntryByKey(sessionKey: string) {
|
|
|
|
|
const cfg = loadConfig();
|
|
|
|
|
const agentId = resolveAgentIdFromSessionKey(sessionKey);
|
|
|
|
|
@@ -472,7 +658,7 @@ export async function runSubagentAnnounceFlow(params: {
|
|
|
|
|
let outcome: SubagentRunOutcome | undefined = params.outcome;
|
|
|
|
|
// Lifecycle "end" can arrive before auto-compaction retries finish. If the
|
|
|
|
|
// subagent is still active, wait for the embedded run to fully settle.
|
|
|
|
|
if (childSessionId && isEmbeddedPiRunActive(childSessionId)) {
|
|
|
|
|
if (!expectsCompletionMessage && childSessionId && isEmbeddedPiRunActive(childSessionId)) {
|
|
|
|
|
const settled = await waitForEmbeddedPiRunEnd(childSessionId, settleTimeoutMs);
|
|
|
|
|
if (!settled && isEmbeddedPiRunActive(childSessionId)) {
|
|
|
|
|
// The child run is still active (e.g., compaction retry still in progress).
|
|
|
|
|
@@ -531,7 +717,12 @@ export async function runSubagentAnnounceFlow(params: {
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if (!reply?.trim() && childSessionId && isEmbeddedPiRunActive(childSessionId)) {
|
|
|
|
|
if (
|
|
|
|
|
!expectsCompletionMessage &&
|
|
|
|
|
!reply?.trim() &&
|
|
|
|
|
childSessionId &&
|
|
|
|
|
isEmbeddedPiRunActive(childSessionId)
|
|
|
|
|
) {
|
|
|
|
|
// Avoid announcing "(no output)" while the child run is still producing output.
|
|
|
|
|
shouldDeleteChildSession = false;
|
|
|
|
|
return false;
|
|
|
|
|
@@ -548,7 +739,7 @@ export async function runSubagentAnnounceFlow(params: {
|
|
|
|
|
} catch {
|
|
|
|
|
// Best-effort only; fall back to direct announce behavior when unavailable.
|
|
|
|
|
}
|
|
|
|
|
if (activeChildDescendantRuns > 0) {
|
|
|
|
|
if (!expectsCompletionMessage && activeChildDescendantRuns > 0) {
|
|
|
|
|
// The finished run still has active descendant subagents. Defer announcing
|
|
|
|
|
// this run until descendants settle so we avoid posting in-progress updates.
|
|
|
|
|
shouldDeleteChildSession = false;
|
|
|
|
|
@@ -669,34 +860,32 @@ export async function runSubagentAnnounceFlow(params: {
|
|
|
|
|
// Send to the requester session. For nested subagents this is an internal
|
|
|
|
|
// follow-up injection (deliver=false) so the orchestrator receives it.
|
|
|
|
|
let directOrigin = targetRequesterOrigin;
|
|
|
|
|
if (!requesterIsSubagent && !directOrigin) {
|
|
|
|
|
if (!requesterIsSubagent) {
|
|
|
|
|
const { entry } = loadRequesterSessionEntry(targetRequesterSessionKey);
|
|
|
|
|
directOrigin = deliveryContextFromSession(entry);
|
|
|
|
|
directOrigin = resolveAnnounceOrigin(entry, targetRequesterOrigin);
|
|
|
|
|
}
|
|
|
|
|
// Use a deterministic idempotency key so the gateway dedup cache
|
|
|
|
|
// catches duplicates if this announce is also queued by the gateway-
|
|
|
|
|
// level message queue while the main session is busy (#17122).
|
|
|
|
|
const directIdempotencyKey = buildAnnounceIdempotencyKey(announceId);
|
|
|
|
|
await callGateway({
|
|
|
|
|
method: "agent",
|
|
|
|
|
params: {
|
|
|
|
|
sessionKey: targetRequesterSessionKey,
|
|
|
|
|
message: triggerMessage,
|
|
|
|
|
deliver: !requesterIsSubagent,
|
|
|
|
|
channel: requesterIsSubagent ? undefined : directOrigin?.channel,
|
|
|
|
|
accountId: requesterIsSubagent ? undefined : directOrigin?.accountId,
|
|
|
|
|
to: requesterIsSubagent ? undefined : directOrigin?.to,
|
|
|
|
|
threadId:
|
|
|
|
|
!requesterIsSubagent && directOrigin?.threadId != null && directOrigin.threadId !== ""
|
|
|
|
|
? String(directOrigin.threadId)
|
|
|
|
|
: undefined,
|
|
|
|
|
idempotencyKey: directIdempotencyKey,
|
|
|
|
|
},
|
|
|
|
|
expectFinal: true,
|
|
|
|
|
timeoutMs: 15_000,
|
|
|
|
|
const delivery = await deliverSubagentCompletionAnnouncement({
|
|
|
|
|
requesterSessionKey: targetRequesterSessionKey,
|
|
|
|
|
announceId,
|
|
|
|
|
triggerMessage,
|
|
|
|
|
summaryLine: taskLabel,
|
|
|
|
|
requesterOrigin: targetRequesterOrigin,
|
|
|
|
|
directOrigin,
|
|
|
|
|
targetRequesterSessionKey,
|
|
|
|
|
requesterIsSubagent,
|
|
|
|
|
expectsCompletionMessage: expectsCompletionMessage,
|
|
|
|
|
directIdempotencyKey,
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
didAnnounce = true;
|
|
|
|
|
didAnnounce = delivery.delivered;
|
|
|
|
|
if (!delivery.delivered && delivery.path === "direct" && delivery.error) {
|
|
|
|
|
defaultRuntime.error?.(
|
|
|
|
|
`Subagent completion direct announce failed for run ${params.childRunId}: ${delivery.error}`,
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
} catch (err) {
|
|
|
|
|
defaultRuntime.error?.(`Subagent announce failed: ${String(err)}`);
|
|
|
|
|
// Best-effort follow-ups; ignore failures to avoid breaking the caller response.
|
|
|
|
|
|