mirror of
https://github.com/openclaw/openclaw.git
synced 2026-04-18 17:37:27 +00:00
* Agents: add subagent orchestration controls
* Agents: add subagent orchestration controls (WIP uncommitted changes)
* feat(subagents): add depth-based spawn gating for sub-sub-agents
* feat(subagents): tool policy, registry, and announce chain for nested agents
* feat(subagents): system prompt, docs, changelog for nested sub-agents
* fix(subagents): prevent model fallback override, show model during active runs, and block context overflow fallback
Bug 1: When a session has an explicit model override (e.g., gpt/openai-codex),
the fallback candidate logic in resolveFallbackCandidates silently appended the
global primary model (opus) as a backstop. On reinjection/steer with a transient
error, the session could fall back to opus which has a smaller context window
and crash. Fix: when storedModelOverride is set, pass fallbacksOverride ?? []
instead of undefined, preventing the implicit primary backstop.
Bug 2: Active subagents showed 'model n/a' in /subagents list because
resolveModelDisplay only read entry.model/modelProvider (populated after run
completes). Fix: fall back to modelOverride/providerOverride fields which are
populated at spawn time via sessions.patch.
Bug 3: Context overflow errors (prompt too long, context_length_exceeded) could
theoretically escape runEmbeddedPiAgent and be treated as failover candidates
in runWithModelFallback, causing a switch to a model with a smaller context
window. Fix: in runWithModelFallback, detect context overflow errors via
isLikelyContextOverflowError and rethrow them immediately instead of trying the
next model candidate.
* fix(subagents): track spawn depth in session store and fix announce routing for nested agents
* Fix compaction status tracking and dedupe overflow compaction triggers
* fix(subagents): enforce depth block via session store and implement cascade kill
* fix: inject group chat context into system prompt
* fix(subagents): always write model to session store at spawn time
* Preserve spawnDepth when agent handler rewrites session entry
* fix(subagents): suppress announce on steer-restart
* fix(subagents): fallback spawned session model to runtime default
* fix(subagents): enforce spawn depth when caller key resolves by sessionId
* feat(subagents): implement active-first ordering for numeric targets and enhance task display
- Added a test to verify that subagents with numeric targets follow an active-first list ordering.
- Updated `resolveSubagentTarget` to sort subagent runs based on active status and recent activity.
- Enhanced task display in command responses to prevent truncation of long task descriptions.
- Introduced new utility functions for compacting task text and managing subagent run states.
* fix(subagents): show model for active runs via run record fallback
When the spawned model matches the agent's default model, the session
store's override fields are intentionally cleared (isDefault: true).
The model/modelProvider fields are only populated after the run
completes. This left active subagents showing 'model n/a'.
Fix: store the resolved model on SubagentRunRecord at registration
time, and use it as a fallback in both display paths (subagents tool
and /subagents command) when the session store entry has no model info.
Changes:
- SubagentRunRecord: add optional model field
- registerSubagentRun: accept and persist model param
- sessions-spawn-tool: pass resolvedModel to registerSubagentRun
- subagents-tool: pass run record model as fallback to resolveModelDisplay
- commands-subagents: pass run record model as fallback to resolveModelDisplay
* feat(chat): implement session key resolution and reset on sidebar navigation
- Added functions to resolve the main session key and reset chat state when switching sessions from the sidebar.
- Updated the `renderTab` function to handle session key changes when navigating to the chat tab.
- Introduced a test to verify that the session resets to "main" when opening chat from the sidebar navigation.
* fix: subagent timeout=0 passthrough and fallback prompt duplication
Bug 1: runTimeoutSeconds=0 now means 'no timeout' instead of applying 600s default
- sessions-spawn-tool: default to undefined (not 0) when neither timeout param
is provided; use != null check so explicit 0 passes through to gateway
- agent.ts: accept 0 as valid timeout (resolveAgentTimeoutMs already handles
0 → MAX_SAFE_TIMEOUT_MS)
Bug 2: model fallback no longer re-injects the original prompt as a duplicate
- agent.ts: track fallback attempt index; on retries use a short continuation
message instead of the full original prompt since the session file already
contains it from the first attempt
- Also skip re-sending images on fallback retries (already in session)
* feat(subagents): truncate long task descriptions in subagents command output
- Introduced a new utility function to format task previews, limiting their length to improve readability.
- Updated the command handler to use the new formatting function, ensuring task descriptions are truncated appropriately.
- Adjusted related tests to verify that long task descriptions are now truncated in the output.
* refactor(subagents): update subagent registry path resolution and improve command output formatting
- Replaced direct import of STATE_DIR with a utility function to resolve the state directory dynamically.
- Enhanced the formatting of command output for active and recent subagents, adding separators for better readability.
- Updated related tests to reflect changes in command output structure.
* fix(subagent): default sessions_spawn to no timeout when runTimeoutSeconds omitted
The previous fix (75a791106) correctly handled the case where
runTimeoutSeconds was explicitly set to 0 ("no timeout"). However,
when models omit the parameter entirely (which is common since the
schema marks it as optional), runTimeoutSeconds resolved to undefined.
undefined flowed through the chain as:
sessions_spawn → timeout: undefined (since undefined != null is false)
→ gateway agent handler → agentCommand opts.timeout: undefined
→ resolveAgentTimeoutMs({ overrideSeconds: undefined })
→ DEFAULT_AGENT_TIMEOUT_SECONDS (600s = 10 minutes)
This caused subagents to be killed at exactly 10 minutes even though
the user's intent (via TOOLS.md) was for subagents to run without a
timeout.
Fix: default runTimeoutSeconds to 0 (no timeout) when neither
runTimeoutSeconds nor timeoutSeconds is provided by the caller.
Subagent spawns are long-running by design and should not inherit the
600s agent-command default timeout.
* fix(subagent): accept timeout=0 in agent-via-gateway path (second 600s default)
* fix: thread timeout override through getReplyFromConfig dispatch path
getReplyFromConfig called resolveAgentTimeoutMs({ cfg }) with no override,
always falling back to the config default (600s). Add timeoutOverrideSeconds
to GetReplyOptions and pass it through as overrideSeconds so callers of the
dispatch chain can specify a custom timeout (0 = no timeout).
This complements the existing timeout threading in agentCommand and the
cron isolated-agent runner, which already pass overrideSeconds correctly.
* feat(model-fallback): normalize OpenAI Codex model references and enhance fallback handling
- Added normalization for OpenAI Codex model references, specifically converting "gpt-5.3-codex" to "openai-codex" before execution.
- Updated the `resolveFallbackCandidates` function to utilize the new normalization logic.
- Enhanced tests to verify the correct behavior of model normalization and fallback mechanisms.
- Introduced a new test case to ensure that the normalization process works as expected for various input formats.
* feat(tests): add unit tests for steer failure behavior in openclaw-tools
- Introduced a new test file to validate the behavior of subagents when steer replacement dispatch fails.
- Implemented tests to ensure that the announce behavior is restored correctly and that the suppression reason is cleared as expected.
- Enhanced the subagent registry with a new function to clear steer restart suppression.
- Updated related components to support the new test scenarios.
* fix(subagents): replace stop command with kill in slash commands and documentation
- Updated the `/subagents` command to replace `stop` with `kill` for consistency in controlling sub-agent runs.
- Modified related documentation to reflect the change in command usage.
- Removed legacy timeoutSeconds references from the sessions-spawn-tool schema and tests to streamline timeout handling.
- Enhanced tests to ensure correct behavior of the updated commands and their interactions.
* feat(tests): add unit tests for readLatestAssistantReply function
- Introduced a new test file for the `readLatestAssistantReply` function to validate its behavior with various message scenarios.
- Implemented tests to ensure the function correctly retrieves the latest assistant message and handles cases where the latest message has no text.
- Mocked the gateway call to simulate different message histories for comprehensive testing.
* feat(tests): enhance subagent kill-all cascade tests and announce formatting
- Added a new test to verify that the `kill-all` command cascades through ended parents to active descendants in subagents.
- Updated the subagent announce formatting tests to reflect changes in message structure, including the replacement of "Findings:" with "Result:" and the addition of new expectations for message content.
- Improved the handling of long findings and stats in the announce formatting logic to ensure concise output.
- Refactored related functions to enhance clarity and maintainability in the subagent registry and tools.
* refactor(subagent): update announce formatting and remove unused constants
- Modified the subagent announce formatting to replace "Findings:" with "Result:" and adjusted related expectations in tests.
- Removed constants for maximum announce findings characters and summary words, simplifying the announcement logic.
- Updated the handling of findings to retain full content instead of truncating, ensuring more informative outputs.
- Cleaned up unused imports in the commands-subagents file to enhance code clarity.
* feat(tests): enhance billing error handling in user-facing text
- Added tests to ensure that normal text mentioning billing plans is not rewritten, preserving user context.
- Updated the `isBillingErrorMessage` and `sanitizeUserFacingText` functions to improve handling of billing-related messages.
- Introduced new test cases for various scenarios involving billing messages to ensure accurate processing and output.
- Enhanced the subagent announce flow to correctly manage active descendant runs, preventing premature announcements.
* feat(subagent): enhance workflow guidance and auto-announcement clarity
- Added a new guideline in the subagent system prompt to emphasize trust in push-based completion, discouraging busy polling for status updates.
- Updated documentation to clarify that sub-agents will automatically announce their results, improving user understanding of the workflow.
- Enhanced tests to verify the new guidance on avoiding polling loops and to ensure the accuracy of the updated prompts.
* fix(cron): avoid announcing interim subagent spawn acks
* chore: clean post-rebase imports
* fix(cron): fall back to child replies when parent stays interim
* fix(subagents): make active-run guidance advisory
* fix(subagents): update announce flow to handle active descendants and enhance test coverage
- Modified the announce flow to defer announcements when active descendant runs are present, ensuring accurate status reporting.
- Updated tests to verify the new behavior, including scenarios where no fallback requester is available and ensuring proper handling of finished subagents.
- Enhanced the announce formatting to include an `expectFinal` flag for better clarity in the announcement process.
* fix(subagents): enhance announce flow and formatting for user updates
- Updated the announce flow to provide clearer instructions for user updates based on active subagent runs and requester context.
- Refactored the announcement logic to improve clarity and ensure internal context remains private.
- Enhanced tests to verify the new message expectations and formatting, including updated prompts for user-facing updates.
- Introduced a new function to build reply instructions based on session context, improving the overall announcement process.
* fix: resolve prep blockers and changelog placement (#14447) (thanks @tyler6204)
* fix: restore cron delivery-plan import after rebase (#14447) (thanks @tyler6204)
* fix: resolve test failures from rebase conflicts (#14447) (thanks @tyler6204)
* fix: apply formatting after rebase (#14447) (thanks @tyler6204)
772 lines
26 KiB
TypeScript
772 lines
26 KiB
TypeScript
import crypto from "node:crypto";
|
||
import type { SubagentRunRecord } from "../../agents/subagent-registry.js";
|
||
import type { CommandHandler } from "./commands-types.js";
|
||
import { AGENT_LANE_SUBAGENT } from "../../agents/lanes.js";
|
||
import { abortEmbeddedPiRun } from "../../agents/pi-embedded.js";
|
||
import {
|
||
clearSubagentRunSteerRestart,
|
||
listSubagentRunsForRequester,
|
||
markSubagentRunTerminated,
|
||
markSubagentRunForSteerRestart,
|
||
replaceSubagentRunAfterSteer,
|
||
} from "../../agents/subagent-registry.js";
|
||
import {
|
||
extractAssistantText,
|
||
resolveInternalSessionKey,
|
||
resolveMainSessionAlias,
|
||
sanitizeTextContent,
|
||
stripToolMessages,
|
||
} from "../../agents/tools/sessions-helpers.js";
|
||
import {
|
||
type SessionEntry,
|
||
loadSessionStore,
|
||
resolveStorePath,
|
||
updateSessionStore,
|
||
} from "../../config/sessions.js";
|
||
import { callGateway } from "../../gateway/call.js";
|
||
import { logVerbose } from "../../globals.js";
|
||
import { formatTimeAgo } from "../../infra/format-time/format-relative.ts";
|
||
import { parseAgentSessionKey } from "../../routing/session-key.js";
|
||
import { INTERNAL_MESSAGE_CHANNEL } from "../../utils/message-channel.js";
|
||
import { stopSubagentsForRequester } from "./abort.js";
|
||
import { clearSessionQueues } from "./queue.js";
|
||
import { formatRunLabel, formatRunStatus, sortSubagentRuns } from "./subagents-utils.js";
|
||
|
||
type SubagentTargetResolution = {
|
||
entry?: SubagentRunRecord;
|
||
error?: string;
|
||
};
|
||
|
||
const COMMAND = "/subagents";
|
||
const COMMAND_KILL = "/kill";
|
||
const COMMAND_STEER = "/steer";
|
||
const COMMAND_TELL = "/tell";
|
||
const ACTIONS = new Set(["list", "kill", "log", "send", "steer", "info", "help"]);
|
||
const RECENT_WINDOW_MINUTES = 30;
|
||
const SUBAGENT_TASK_PREVIEW_MAX = 110;
|
||
const STEER_ABORT_SETTLE_TIMEOUT_MS = 5_000;
|
||
|
||
function formatDurationCompact(valueMs?: number) {
|
||
if (!valueMs || !Number.isFinite(valueMs) || valueMs <= 0) {
|
||
return "n/a";
|
||
}
|
||
const minutes = Math.max(1, Math.round(valueMs / 60_000));
|
||
if (minutes < 60) {
|
||
return `${minutes}m`;
|
||
}
|
||
const hours = Math.floor(minutes / 60);
|
||
const minutesRemainder = minutes % 60;
|
||
if (hours < 24) {
|
||
return minutesRemainder > 0 ? `${hours}h${minutesRemainder}m` : `${hours}h`;
|
||
}
|
||
const days = Math.floor(hours / 24);
|
||
const hoursRemainder = hours % 24;
|
||
return hoursRemainder > 0 ? `${days}d${hoursRemainder}h` : `${days}d`;
|
||
}
|
||
|
||
function formatTokenShort(value?: number) {
|
||
if (!value || !Number.isFinite(value) || value <= 0) {
|
||
return undefined;
|
||
}
|
||
const n = Math.floor(value);
|
||
if (n < 1_000) {
|
||
return `${n}`;
|
||
}
|
||
if (n < 10_000) {
|
||
return `${(n / 1_000).toFixed(1).replace(/\.0$/, "")}k`;
|
||
}
|
||
if (n < 1_000_000) {
|
||
return `${Math.round(n / 1_000)}k`;
|
||
}
|
||
return `${(n / 1_000_000).toFixed(1).replace(/\.0$/, "")}m`;
|
||
}
|
||
|
||
function truncateLine(value: string, maxLength: number) {
|
||
if (value.length <= maxLength) {
|
||
return value;
|
||
}
|
||
return `${value.slice(0, maxLength).trimEnd()}...`;
|
||
}
|
||
|
||
function compactLine(value: string) {
|
||
return value.replace(/\s+/g, " ").trim();
|
||
}
|
||
|
||
function formatTaskPreview(value: string) {
|
||
return truncateLine(compactLine(value), SUBAGENT_TASK_PREVIEW_MAX);
|
||
}
|
||
|
||
function resolveModelDisplay(
|
||
entry?: {
|
||
model?: unknown;
|
||
modelProvider?: unknown;
|
||
modelOverride?: unknown;
|
||
providerOverride?: unknown;
|
||
},
|
||
fallbackModel?: string,
|
||
) {
|
||
const model = typeof entry?.model === "string" ? entry.model.trim() : "";
|
||
const provider = typeof entry?.modelProvider === "string" ? entry.modelProvider.trim() : "";
|
||
let combined = model.includes("/") ? model : model && provider ? `${provider}/${model}` : model;
|
||
if (!combined) {
|
||
// Fall back to override fields which are populated at spawn time,
|
||
// before the first run completes and writes model/modelProvider.
|
||
const overrideModel =
|
||
typeof entry?.modelOverride === "string" ? entry.modelOverride.trim() : "";
|
||
const overrideProvider =
|
||
typeof entry?.providerOverride === "string" ? entry.providerOverride.trim() : "";
|
||
combined = overrideModel.includes("/")
|
||
? overrideModel
|
||
: overrideModel && overrideProvider
|
||
? `${overrideProvider}/${overrideModel}`
|
||
: overrideModel;
|
||
}
|
||
if (!combined) {
|
||
combined = fallbackModel?.trim() || "";
|
||
}
|
||
if (!combined) {
|
||
return "model n/a";
|
||
}
|
||
const slash = combined.lastIndexOf("/");
|
||
if (slash >= 0 && slash < combined.length - 1) {
|
||
return combined.slice(slash + 1);
|
||
}
|
||
return combined;
|
||
}
|
||
|
||
function resolveTotalTokens(entry?: {
|
||
totalTokens?: unknown;
|
||
inputTokens?: unknown;
|
||
outputTokens?: unknown;
|
||
}) {
|
||
if (!entry || typeof entry !== "object") {
|
||
return undefined;
|
||
}
|
||
if (typeof entry.totalTokens === "number" && Number.isFinite(entry.totalTokens)) {
|
||
return entry.totalTokens;
|
||
}
|
||
const input = typeof entry.inputTokens === "number" ? entry.inputTokens : 0;
|
||
const output = typeof entry.outputTokens === "number" ? entry.outputTokens : 0;
|
||
const total = input + output;
|
||
return total > 0 ? total : undefined;
|
||
}
|
||
|
||
function resolveIoTokens(entry?: { inputTokens?: unknown; outputTokens?: unknown }) {
|
||
if (!entry || typeof entry !== "object") {
|
||
return undefined;
|
||
}
|
||
const input =
|
||
typeof entry.inputTokens === "number" && Number.isFinite(entry.inputTokens)
|
||
? entry.inputTokens
|
||
: 0;
|
||
const output =
|
||
typeof entry.outputTokens === "number" && Number.isFinite(entry.outputTokens)
|
||
? entry.outputTokens
|
||
: 0;
|
||
const total = input + output;
|
||
if (total <= 0) {
|
||
return undefined;
|
||
}
|
||
return { input, output, total };
|
||
}
|
||
|
||
function resolveUsageDisplay(entry?: {
|
||
totalTokens?: unknown;
|
||
inputTokens?: unknown;
|
||
outputTokens?: unknown;
|
||
}) {
|
||
const io = resolveIoTokens(entry);
|
||
const promptCache = resolveTotalTokens(entry);
|
||
const parts: string[] = [];
|
||
if (io) {
|
||
const input = formatTokenShort(io.input) ?? "0";
|
||
const output = formatTokenShort(io.output) ?? "0";
|
||
parts.push(`tokens ${formatTokenShort(io.total)} (in ${input} / out ${output})`);
|
||
} else if (typeof promptCache === "number" && promptCache > 0) {
|
||
parts.push(`tokens ${formatTokenShort(promptCache)} prompt/cache`);
|
||
}
|
||
if (typeof promptCache === "number" && io && promptCache > io.total) {
|
||
parts.push(`prompt/cache ${formatTokenShort(promptCache)}`);
|
||
}
|
||
return parts.join(", ");
|
||
}
|
||
|
||
function resolveDisplayStatus(entry: SubagentRunRecord) {
|
||
const status = formatRunStatus(entry);
|
||
return status === "error" ? "failed" : status;
|
||
}
|
||
|
||
function formatTimestamp(valueMs?: number) {
|
||
if (!valueMs || !Number.isFinite(valueMs) || valueMs <= 0) {
|
||
return "n/a";
|
||
}
|
||
return new Date(valueMs).toISOString();
|
||
}
|
||
|
||
function formatTimestampWithAge(valueMs?: number) {
|
||
if (!valueMs || !Number.isFinite(valueMs) || valueMs <= 0) {
|
||
return "n/a";
|
||
}
|
||
return `${formatTimestamp(valueMs)} (${formatTimeAgo(Date.now() - valueMs, { fallback: "n/a" })})`;
|
||
}
|
||
|
||
function resolveRequesterSessionKey(params: Parameters<CommandHandler>[0]): string | undefined {
|
||
const raw = params.sessionKey?.trim() || params.ctx.CommandTargetSessionKey?.trim();
|
||
if (!raw) {
|
||
return undefined;
|
||
}
|
||
const { mainKey, alias } = resolveMainSessionAlias(params.cfg);
|
||
return resolveInternalSessionKey({ key: raw, alias, mainKey });
|
||
}
|
||
|
||
function resolveSubagentTarget(
|
||
runs: SubagentRunRecord[],
|
||
token: string | undefined,
|
||
): SubagentTargetResolution {
|
||
const trimmed = token?.trim();
|
||
if (!trimmed) {
|
||
return { error: "Missing subagent id." };
|
||
}
|
||
if (trimmed === "last") {
|
||
const sorted = sortSubagentRuns(runs);
|
||
return { entry: sorted[0] };
|
||
}
|
||
const sorted = sortSubagentRuns(runs);
|
||
const recentCutoff = Date.now() - RECENT_WINDOW_MINUTES * 60_000;
|
||
const numericOrder = [
|
||
...sorted.filter((entry) => !entry.endedAt),
|
||
...sorted.filter((entry) => !!entry.endedAt && (entry.endedAt ?? 0) >= recentCutoff),
|
||
];
|
||
if (/^\d+$/.test(trimmed)) {
|
||
const idx = Number.parseInt(trimmed, 10);
|
||
if (!Number.isFinite(idx) || idx <= 0 || idx > numericOrder.length) {
|
||
return { error: `Invalid subagent index: ${trimmed}` };
|
||
}
|
||
return { entry: numericOrder[idx - 1] };
|
||
}
|
||
if (trimmed.includes(":")) {
|
||
const match = runs.find((entry) => entry.childSessionKey === trimmed);
|
||
return match ? { entry: match } : { error: `Unknown subagent session: ${trimmed}` };
|
||
}
|
||
const lowered = trimmed.toLowerCase();
|
||
const byLabel = runs.filter((entry) => formatRunLabel(entry).toLowerCase() === lowered);
|
||
if (byLabel.length === 1) {
|
||
return { entry: byLabel[0] };
|
||
}
|
||
if (byLabel.length > 1) {
|
||
return { error: `Ambiguous subagent label: ${trimmed}` };
|
||
}
|
||
const byLabelPrefix = runs.filter((entry) =>
|
||
formatRunLabel(entry).toLowerCase().startsWith(lowered),
|
||
);
|
||
if (byLabelPrefix.length === 1) {
|
||
return { entry: byLabelPrefix[0] };
|
||
}
|
||
if (byLabelPrefix.length > 1) {
|
||
return { error: `Ambiguous subagent label prefix: ${trimmed}` };
|
||
}
|
||
const byRunId = runs.filter((entry) => entry.runId.startsWith(trimmed));
|
||
if (byRunId.length === 1) {
|
||
return { entry: byRunId[0] };
|
||
}
|
||
if (byRunId.length > 1) {
|
||
return { error: `Ambiguous run id prefix: ${trimmed}` };
|
||
}
|
||
return { error: `Unknown subagent id: ${trimmed}` };
|
||
}
|
||
|
||
function buildSubagentsHelp() {
|
||
return [
|
||
"Subagents",
|
||
"Usage:",
|
||
"- /subagents list",
|
||
"- /subagents kill <id|#|all>",
|
||
"- /subagents log <id|#> [limit] [tools]",
|
||
"- /subagents info <id|#>",
|
||
"- /subagents send <id|#> <message>",
|
||
"- /subagents steer <id|#> <message>",
|
||
"- /kill <id|#|all>",
|
||
"- /steer <id|#> <message>",
|
||
"- /tell <id|#> <message>",
|
||
"",
|
||
"Ids: use the list index (#), runId/session prefix, label, or full session key.",
|
||
].join("\n");
|
||
}
|
||
|
||
type ChatMessage = {
|
||
role?: unknown;
|
||
content?: unknown;
|
||
name?: unknown;
|
||
toolName?: unknown;
|
||
};
|
||
|
||
function normalizeMessageText(text: string) {
|
||
return text.replace(/\s+/g, " ").trim();
|
||
}
|
||
|
||
export function extractMessageText(message: ChatMessage): { role: string; text: string } | null {
|
||
const role = typeof message.role === "string" ? message.role : "";
|
||
const shouldSanitize = role === "assistant";
|
||
const content = message.content;
|
||
if (typeof content === "string") {
|
||
const normalized = normalizeMessageText(
|
||
shouldSanitize ? sanitizeTextContent(content) : content,
|
||
);
|
||
return normalized ? { role, text: normalized } : null;
|
||
}
|
||
if (!Array.isArray(content)) {
|
||
return null;
|
||
}
|
||
const chunks: string[] = [];
|
||
for (const block of content) {
|
||
if (!block || typeof block !== "object") {
|
||
continue;
|
||
}
|
||
if ((block as { type?: unknown }).type !== "text") {
|
||
continue;
|
||
}
|
||
const text = (block as { text?: unknown }).text;
|
||
if (typeof text === "string") {
|
||
const value = shouldSanitize ? sanitizeTextContent(text) : text;
|
||
if (value.trim()) {
|
||
chunks.push(value);
|
||
}
|
||
}
|
||
}
|
||
const joined = normalizeMessageText(chunks.join(" "));
|
||
return joined ? { role, text: joined } : null;
|
||
}
|
||
|
||
function formatLogLines(messages: ChatMessage[]) {
|
||
const lines: string[] = [];
|
||
for (const msg of messages) {
|
||
const extracted = extractMessageText(msg);
|
||
if (!extracted) {
|
||
continue;
|
||
}
|
||
const label = extracted.role === "assistant" ? "Assistant" : "User";
|
||
lines.push(`${label}: ${extracted.text}`);
|
||
}
|
||
return lines;
|
||
}
|
||
|
||
type SessionStoreCache = Map<string, Record<string, SessionEntry>>;
|
||
|
||
function loadSubagentSessionEntry(
|
||
params: Parameters<CommandHandler>[0],
|
||
childKey: string,
|
||
storeCache?: SessionStoreCache,
|
||
) {
|
||
const parsed = parseAgentSessionKey(childKey);
|
||
const storePath = resolveStorePath(params.cfg.session?.store, { agentId: parsed?.agentId });
|
||
let store = storeCache?.get(storePath);
|
||
if (!store) {
|
||
store = loadSessionStore(storePath);
|
||
storeCache?.set(storePath, store);
|
||
}
|
||
return { storePath, store, entry: store[childKey] };
|
||
}
|
||
|
||
export const handleSubagentsCommand: CommandHandler = async (params, allowTextCommands) => {
|
||
if (!allowTextCommands) {
|
||
return null;
|
||
}
|
||
const normalized = params.command.commandBodyNormalized;
|
||
const handledPrefix = normalized.startsWith(COMMAND)
|
||
? COMMAND
|
||
: normalized.startsWith(COMMAND_KILL)
|
||
? COMMAND_KILL
|
||
: normalized.startsWith(COMMAND_STEER)
|
||
? COMMAND_STEER
|
||
: normalized.startsWith(COMMAND_TELL)
|
||
? COMMAND_TELL
|
||
: null;
|
||
if (!handledPrefix) {
|
||
return null;
|
||
}
|
||
if (!params.command.isAuthorizedSender) {
|
||
logVerbose(
|
||
`Ignoring ${handledPrefix} from unauthorized sender: ${params.command.senderId || "<unknown>"}`,
|
||
);
|
||
return { shouldContinue: false };
|
||
}
|
||
|
||
const rest = normalized.slice(handledPrefix.length).trim();
|
||
const restTokens = rest.split(/\s+/).filter(Boolean);
|
||
let action = "list";
|
||
if (handledPrefix === COMMAND) {
|
||
const [actionRaw] = restTokens;
|
||
action = actionRaw?.toLowerCase() || "list";
|
||
if (!ACTIONS.has(action)) {
|
||
return { shouldContinue: false, reply: { text: buildSubagentsHelp() } };
|
||
}
|
||
restTokens.splice(0, 1);
|
||
} else if (handledPrefix === COMMAND_KILL) {
|
||
action = "kill";
|
||
} else {
|
||
action = "steer";
|
||
}
|
||
|
||
const requesterKey = resolveRequesterSessionKey(params);
|
||
if (!requesterKey) {
|
||
return { shouldContinue: false, reply: { text: "⚠️ Missing session key." } };
|
||
}
|
||
const runs = listSubagentRunsForRequester(requesterKey);
|
||
|
||
if (action === "help") {
|
||
return { shouldContinue: false, reply: { text: buildSubagentsHelp() } };
|
||
}
|
||
|
||
if (action === "list") {
|
||
const sorted = sortSubagentRuns(runs);
|
||
const now = Date.now();
|
||
const recentCutoff = now - RECENT_WINDOW_MINUTES * 60_000;
|
||
const storeCache: SessionStoreCache = new Map();
|
||
let index = 1;
|
||
const activeLines = sorted
|
||
.filter((entry) => !entry.endedAt)
|
||
.map((entry) => {
|
||
const { entry: sessionEntry } = loadSubagentSessionEntry(
|
||
params,
|
||
entry.childSessionKey,
|
||
storeCache,
|
||
);
|
||
const usageText = resolveUsageDisplay(sessionEntry);
|
||
const label = truncateLine(formatRunLabel(entry, { maxLength: 48 }), 48);
|
||
const task = formatTaskPreview(entry.task);
|
||
const runtime = formatDurationCompact(now - (entry.startedAt ?? entry.createdAt));
|
||
const status = resolveDisplayStatus(entry);
|
||
const line = `${index}. ${label} (${resolveModelDisplay(sessionEntry, entry.model)}, ${runtime}${usageText ? `, ${usageText}` : ""}) ${status}${task.toLowerCase() !== label.toLowerCase() ? ` - ${task}` : ""}`;
|
||
index += 1;
|
||
return line;
|
||
});
|
||
const recentLines = sorted
|
||
.filter((entry) => !!entry.endedAt && (entry.endedAt ?? 0) >= recentCutoff)
|
||
.map((entry) => {
|
||
const { entry: sessionEntry } = loadSubagentSessionEntry(
|
||
params,
|
||
entry.childSessionKey,
|
||
storeCache,
|
||
);
|
||
const usageText = resolveUsageDisplay(sessionEntry);
|
||
const label = truncateLine(formatRunLabel(entry, { maxLength: 48 }), 48);
|
||
const task = formatTaskPreview(entry.task);
|
||
const runtime = formatDurationCompact(
|
||
(entry.endedAt ?? now) - (entry.startedAt ?? entry.createdAt),
|
||
);
|
||
const status = resolveDisplayStatus(entry);
|
||
const line = `${index}. ${label} (${resolveModelDisplay(sessionEntry, entry.model)}, ${runtime}${usageText ? `, ${usageText}` : ""}) ${status}${task.toLowerCase() !== label.toLowerCase() ? ` - ${task}` : ""}`;
|
||
index += 1;
|
||
return line;
|
||
});
|
||
|
||
const lines = ["active subagents:", "-----"];
|
||
if (activeLines.length === 0) {
|
||
lines.push("(none)");
|
||
} else {
|
||
lines.push(activeLines.join("\n"));
|
||
}
|
||
lines.push("", `recent subagents (last ${RECENT_WINDOW_MINUTES}m):`, "-----");
|
||
if (recentLines.length === 0) {
|
||
lines.push("(none)");
|
||
} else {
|
||
lines.push(recentLines.join("\n"));
|
||
}
|
||
return { shouldContinue: false, reply: { text: lines.join("\n") } };
|
||
}
|
||
|
||
if (action === "kill") {
|
||
const target = restTokens[0];
|
||
if (!target) {
|
||
return {
|
||
shouldContinue: false,
|
||
reply: {
|
||
text:
|
||
handledPrefix === COMMAND
|
||
? "Usage: /subagents kill <id|#|all>"
|
||
: "Usage: /kill <id|#|all>",
|
||
},
|
||
};
|
||
}
|
||
if (target === "all" || target === "*") {
|
||
stopSubagentsForRequester({
|
||
cfg: params.cfg,
|
||
requesterSessionKey: requesterKey,
|
||
});
|
||
return { shouldContinue: false };
|
||
}
|
||
const resolved = resolveSubagentTarget(runs, target);
|
||
if (!resolved.entry) {
|
||
return {
|
||
shouldContinue: false,
|
||
reply: { text: `⚠️ ${resolved.error ?? "Unknown subagent."}` },
|
||
};
|
||
}
|
||
if (resolved.entry.endedAt) {
|
||
return {
|
||
shouldContinue: false,
|
||
reply: { text: `${formatRunLabel(resolved.entry)} is already finished.` },
|
||
};
|
||
}
|
||
|
||
const childKey = resolved.entry.childSessionKey;
|
||
const { storePath, store, entry } = loadSubagentSessionEntry(params, childKey);
|
||
const sessionId = entry?.sessionId;
|
||
if (sessionId) {
|
||
abortEmbeddedPiRun(sessionId);
|
||
}
|
||
const cleared = clearSessionQueues([childKey, sessionId]);
|
||
if (cleared.followupCleared > 0 || cleared.laneCleared > 0) {
|
||
logVerbose(
|
||
`subagents kill: cleared followups=${cleared.followupCleared} lane=${cleared.laneCleared} keys=${cleared.keys.join(",")}`,
|
||
);
|
||
}
|
||
if (entry) {
|
||
entry.abortedLastRun = true;
|
||
entry.updatedAt = Date.now();
|
||
store[childKey] = entry;
|
||
await updateSessionStore(storePath, (nextStore) => {
|
||
nextStore[childKey] = entry;
|
||
});
|
||
}
|
||
markSubagentRunTerminated({
|
||
runId: resolved.entry.runId,
|
||
childSessionKey: childKey,
|
||
reason: "killed",
|
||
});
|
||
// Cascade: also stop any sub-sub-agents spawned by this child.
|
||
stopSubagentsForRequester({
|
||
cfg: params.cfg,
|
||
requesterSessionKey: childKey,
|
||
});
|
||
return { shouldContinue: false };
|
||
}
|
||
|
||
if (action === "info") {
|
||
const target = restTokens[0];
|
||
if (!target) {
|
||
return { shouldContinue: false, reply: { text: "ℹ️ Usage: /subagents info <id|#>" } };
|
||
}
|
||
const resolved = resolveSubagentTarget(runs, target);
|
||
if (!resolved.entry) {
|
||
return {
|
||
shouldContinue: false,
|
||
reply: { text: `⚠️ ${resolved.error ?? "Unknown subagent."}` },
|
||
};
|
||
}
|
||
const run = resolved.entry;
|
||
const { entry: sessionEntry } = loadSubagentSessionEntry(params, run.childSessionKey);
|
||
const runtime =
|
||
run.startedAt && Number.isFinite(run.startedAt)
|
||
? (formatDurationCompact((run.endedAt ?? Date.now()) - run.startedAt) ?? "n/a")
|
||
: "n/a";
|
||
const outcome = run.outcome
|
||
? `${run.outcome.status}${run.outcome.error ? ` (${run.outcome.error})` : ""}`
|
||
: "n/a";
|
||
const lines = [
|
||
"ℹ️ Subagent info",
|
||
`Status: ${resolveDisplayStatus(run)}`,
|
||
`Label: ${formatRunLabel(run)}`,
|
||
`Task: ${run.task}`,
|
||
`Run: ${run.runId}`,
|
||
`Session: ${run.childSessionKey}`,
|
||
`SessionId: ${sessionEntry?.sessionId ?? "n/a"}`,
|
||
`Transcript: ${sessionEntry?.sessionFile ?? "n/a"}`,
|
||
`Runtime: ${runtime}`,
|
||
`Created: ${formatTimestampWithAge(run.createdAt)}`,
|
||
`Started: ${formatTimestampWithAge(run.startedAt)}`,
|
||
`Ended: ${formatTimestampWithAge(run.endedAt)}`,
|
||
`Cleanup: ${run.cleanup}`,
|
||
run.archiveAtMs ? `Archive: ${formatTimestampWithAge(run.archiveAtMs)}` : undefined,
|
||
run.cleanupHandled ? "Cleanup handled: yes" : undefined,
|
||
`Outcome: ${outcome}`,
|
||
].filter(Boolean);
|
||
return { shouldContinue: false, reply: { text: lines.join("\n") } };
|
||
}
|
||
|
||
if (action === "log") {
|
||
const target = restTokens[0];
|
||
if (!target) {
|
||
return { shouldContinue: false, reply: { text: "📜 Usage: /subagents log <id|#> [limit]" } };
|
||
}
|
||
const includeTools = restTokens.some((token) => token.toLowerCase() === "tools");
|
||
const limitToken = restTokens.find((token) => /^\d+$/.test(token));
|
||
const limit = limitToken ? Math.min(200, Math.max(1, Number.parseInt(limitToken, 10))) : 20;
|
||
const resolved = resolveSubagentTarget(runs, target);
|
||
if (!resolved.entry) {
|
||
return {
|
||
shouldContinue: false,
|
||
reply: { text: `⚠️ ${resolved.error ?? "Unknown subagent."}` },
|
||
};
|
||
}
|
||
const history = await callGateway<{ messages: Array<unknown> }>({
|
||
method: "chat.history",
|
||
params: { sessionKey: resolved.entry.childSessionKey, limit },
|
||
});
|
||
const rawMessages = Array.isArray(history?.messages) ? history.messages : [];
|
||
const filtered = includeTools ? rawMessages : stripToolMessages(rawMessages);
|
||
const lines = formatLogLines(filtered as ChatMessage[]);
|
||
const header = `📜 Subagent log: ${formatRunLabel(resolved.entry)}`;
|
||
if (lines.length === 0) {
|
||
return { shouldContinue: false, reply: { text: `${header}\n(no messages)` } };
|
||
}
|
||
return { shouldContinue: false, reply: { text: [header, ...lines].join("\n") } };
|
||
}
|
||
|
||
if (action === "send" || action === "steer") {
|
||
const steerRequested = action === "steer";
|
||
const target = restTokens[0];
|
||
const message = restTokens.slice(1).join(" ").trim();
|
||
if (!target || !message) {
|
||
return {
|
||
shouldContinue: false,
|
||
reply: {
|
||
text: steerRequested
|
||
? handledPrefix === COMMAND
|
||
? "Usage: /subagents steer <id|#> <message>"
|
||
: `Usage: ${handledPrefix} <id|#> <message>`
|
||
: "Usage: /subagents send <id|#> <message>",
|
||
},
|
||
};
|
||
}
|
||
const resolved = resolveSubagentTarget(runs, target);
|
||
if (!resolved.entry) {
|
||
return {
|
||
shouldContinue: false,
|
||
reply: { text: `⚠️ ${resolved.error ?? "Unknown subagent."}` },
|
||
};
|
||
}
|
||
if (steerRequested && resolved.entry.endedAt) {
|
||
return {
|
||
shouldContinue: false,
|
||
reply: { text: `${formatRunLabel(resolved.entry)} is already finished.` },
|
||
};
|
||
}
|
||
const { entry: targetSessionEntry } = loadSubagentSessionEntry(
|
||
params,
|
||
resolved.entry.childSessionKey,
|
||
);
|
||
const targetSessionId =
|
||
typeof targetSessionEntry?.sessionId === "string" && targetSessionEntry.sessionId.trim()
|
||
? targetSessionEntry.sessionId.trim()
|
||
: undefined;
|
||
|
||
if (steerRequested) {
|
||
// Suppress stale announce before interrupting the in-flight run.
|
||
markSubagentRunForSteerRestart(resolved.entry.runId);
|
||
|
||
// Force an immediate interruption and make steer the next run.
|
||
if (targetSessionId) {
|
||
abortEmbeddedPiRun(targetSessionId);
|
||
}
|
||
const cleared = clearSessionQueues([resolved.entry.childSessionKey, targetSessionId]);
|
||
if (cleared.followupCleared > 0 || cleared.laneCleared > 0) {
|
||
logVerbose(
|
||
`subagents steer: cleared followups=${cleared.followupCleared} lane=${cleared.laneCleared} keys=${cleared.keys.join(",")}`,
|
||
);
|
||
}
|
||
|
||
// Best effort: wait for the interrupted run to settle so the steer
|
||
// message is appended on the existing conversation state.
|
||
try {
|
||
await callGateway({
|
||
method: "agent.wait",
|
||
params: {
|
||
runId: resolved.entry.runId,
|
||
timeoutMs: STEER_ABORT_SETTLE_TIMEOUT_MS,
|
||
},
|
||
timeoutMs: STEER_ABORT_SETTLE_TIMEOUT_MS + 2_000,
|
||
});
|
||
} catch {
|
||
// Continue even if wait fails; steer should still be attempted.
|
||
}
|
||
}
|
||
|
||
const idempotencyKey = crypto.randomUUID();
|
||
let runId: string = idempotencyKey;
|
||
try {
|
||
const response = await callGateway<{ runId: string }>({
|
||
method: "agent",
|
||
params: {
|
||
message,
|
||
sessionKey: resolved.entry.childSessionKey,
|
||
sessionId: targetSessionId,
|
||
idempotencyKey,
|
||
deliver: false,
|
||
channel: INTERNAL_MESSAGE_CHANNEL,
|
||
lane: AGENT_LANE_SUBAGENT,
|
||
timeout: 0,
|
||
},
|
||
timeoutMs: 10_000,
|
||
});
|
||
const responseRunId = typeof response?.runId === "string" ? response.runId : undefined;
|
||
if (responseRunId) {
|
||
runId = responseRunId;
|
||
}
|
||
} catch (err) {
|
||
if (steerRequested) {
|
||
// Replacement launch failed; restore announce behavior for the
|
||
// original run so completion is not silently suppressed.
|
||
clearSubagentRunSteerRestart(resolved.entry.runId);
|
||
}
|
||
const messageText =
|
||
err instanceof Error ? err.message : typeof err === "string" ? err : "error";
|
||
return { shouldContinue: false, reply: { text: `send failed: ${messageText}` } };
|
||
}
|
||
|
||
if (steerRequested) {
|
||
replaceSubagentRunAfterSteer({
|
||
previousRunId: resolved.entry.runId,
|
||
nextRunId: runId,
|
||
fallback: resolved.entry,
|
||
runTimeoutSeconds: resolved.entry.runTimeoutSeconds ?? 0,
|
||
});
|
||
return {
|
||
shouldContinue: false,
|
||
reply: {
|
||
text: `steered ${formatRunLabel(resolved.entry)} (run ${runId.slice(0, 8)}).`,
|
||
},
|
||
};
|
||
}
|
||
|
||
const waitMs = 30_000;
|
||
const wait = await callGateway<{ status?: string; error?: string }>({
|
||
method: "agent.wait",
|
||
params: { runId, timeoutMs: waitMs },
|
||
timeoutMs: waitMs + 2000,
|
||
});
|
||
if (wait?.status === "timeout") {
|
||
return {
|
||
shouldContinue: false,
|
||
reply: { text: `⏳ Subagent still running (run ${runId.slice(0, 8)}).` },
|
||
};
|
||
}
|
||
if (wait?.status === "error") {
|
||
const waitError = typeof wait.error === "string" ? wait.error : "unknown error";
|
||
return {
|
||
shouldContinue: false,
|
||
reply: {
|
||
text: `⚠️ Subagent error: ${waitError} (run ${runId.slice(0, 8)}).`,
|
||
},
|
||
};
|
||
}
|
||
|
||
const history = await callGateway<{ messages: Array<unknown> }>({
|
||
method: "chat.history",
|
||
params: { sessionKey: resolved.entry.childSessionKey, limit: 50 },
|
||
});
|
||
const filtered = stripToolMessages(Array.isArray(history?.messages) ? history.messages : []);
|
||
const last = filtered.length > 0 ? filtered[filtered.length - 1] : undefined;
|
||
const replyText = last ? extractAssistantText(last) : undefined;
|
||
return {
|
||
shouldContinue: false,
|
||
reply: {
|
||
text:
|
||
replyText ?? `✅ Sent to ${formatRunLabel(resolved.entry)} (run ${runId.slice(0, 8)}).`,
|
||
},
|
||
};
|
||
}
|
||
|
||
return { shouldContinue: false, reply: { text: buildSubagentsHelp() } };
|
||
};
|