Skip to content
Draft
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
5 changes: 1 addition & 4 deletions packages/junior/src/chat/agent/handoff.ts
Original file line number Diff line number Diff line change
Expand Up @@ -145,7 +145,6 @@ export async function commitHandoff(args: {
metadata?: CompactContextArgs["metadata"];
onStatus?: (status: { text: string }) => void | Promise<void>;
profile: ModelProfile;
runtimeContextSourceMessages?: PiMessage[];
signal?: AbortSignal;
sourceMessages: PiMessage[];
triggeringToolCallId?: string;
Expand All @@ -154,9 +153,7 @@ export async function commitHandoff(args: {
if (args.profile === args.activeModelProfile) {
return undefined;
}
const runtimeContext = retainRuntimeTurnContext(
args.runtimeContextSourceMessages ?? args.sourceMessages,
);
const runtimeContext = retainRuntimeTurnContext(args.sourceMessages);
const phaseUsageSummary = extractGenAiUsageSummary(
...args.sourceMessages
.slice(args.beforeMessageCount)
Expand Down
103 changes: 16 additions & 87 deletions packages/junior/src/chat/agent/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -323,34 +323,11 @@ async function executeAgentRunInPrivacyContext(
};
const signal = policy.signal;
const state = run.state ?? {};
const observers = {
onStatus: run.onEvent
? async (status: { text: string }) => {
await run.onEvent?.({ type: "status", text: status.text });
}
: undefined,
onToolInvocation: run.onEvent
? async (invocation: {
params: Record<string, unknown>;
toolCallId: string;
toolName: string;
}) => {
await run.onEvent?.({
type: "tool_started",
params: invocation.params,
toolCallId: invocation.toolCallId,
toolName: invocation.toolName,
});
}
: undefined,
onToolResult: run.onEvent
? async (
report: import("@/chat/tool-support/tool-execution-report").ToolExecutionReport,
) => {
await run.onEvent?.({ type: "tool_finished", report });
}
: undefined,
};
const onStatus = run.onEvent
? async (status: { text: string }) => {
await run.onEvent?.({ type: "status", text: status.text });
}
: undefined;
const delivery = run.delivery;
const authorization = run.authorization;
const durability = run.durability ?? {};
Expand Down Expand Up @@ -733,25 +710,21 @@ async function executeAgentRunInPrivacyContext(
actorId: slackActor?.userId,
runId,
};
const scheduleHandoff = async (args: {
profile: ModelProfile;
runtimeContextSourceMessages?: PiMessage[];
signal?: AbortSignal;
sourceMessages: PiMessage[];
triggeringToolCallId?: string;
}) => {
const handoffExecute = async (
profile: ModelProfile,
options: { signal?: AbortSignal; toolCallId: string },
) => {
const pending = await commitHandoff({
activeModelProfile,
beforeMessageCount: runResume.beforeMessageCount,
conversationContext: input.conversationContext,
conversationId,
metadata: handoffMetadata,
onStatus: observers.onStatus,
profile: args.profile,
runtimeContextSourceMessages: args.runtimeContextSourceMessages,
signal: args.signal,
sourceMessages: args.sourceMessages,
triggeringToolCallId: args.triggeringToolCallId,
onStatus,
profile,
signal: options.signal,
sourceMessages: [...agent!.state.messages],
triggeringToolCallId: options.toolCallId,
turnRoute: turnRoute!,
});
if (!pending) {
Expand All @@ -763,16 +736,6 @@ async function executeAgentRunInPrivacyContext(
activeModelId = pending.modelId;
turnRoute = pending.turnRoute;
};
const handoffExecute = async (
profile: ModelProfile,
options: { signal?: AbortSignal; toolCallId: string },
) =>
await scheduleHandoff({
profile,
signal: options.signal,
sourceMessages: [...agent!.state.messages],
triggeringToolCallId: options.toolCallId,
});
const requestHandoff = handoffControl({
activeProfile: activeModelProfile,
enabled: handoffEnabled,
Expand Down Expand Up @@ -946,8 +909,7 @@ async function executeAgentRunInPrivacyContext(
},
modelId: activeModelId,
modelProfile: activeModelProfile,
onCompactionStart: () =>
observers.onStatus?.({ text: "Compacting context" }),
onCompactionStart: () => onStatus?.({ text: "Compacting context" }),
pendingMessages,
piMessages: messages,
...(pairPendingRuntimeContext && pendingMessages
Expand Down Expand Up @@ -1373,33 +1335,6 @@ async function executeAgentRunInPrivacyContext(
return result;
};

let run: Promise<unknown>;
let handoffApplied = false;
const requestedProfile =
activeModelProfile === botConfig.defaultProfile
? turnRoute!.profile
: undefined;
if (
requestedProfile &&
requestedProfile !== botConfig.defaultProfile
) {
const handoffAbortController = new AbortController();
await runAgentStep(
scheduleHandoff({
profile: requestedProfile,
runtimeContextSourceMessages: shouldPromptAgent
? [
...(contextMessage ? [contextMessage] : []),
freshPromptMessage,
]
: undefined,
signal: handoffAbortController.signal,
sourceMessages: [...agent!.state.messages],
}),
() => handoffAbortController.abort(),
);
handoffApplied = Boolean(applyPendingHandoff());
}
const compactionAbortController = new AbortController();
const capacityUpdate = await runAgentStep(
applyActiveContextCompaction(
Expand All @@ -1418,13 +1353,7 @@ async function executeAgentRunInPrivacyContext(
),
() => compactionAbortController.abort(),
);
if (shouldPromptAgent && handoffApplied && !capacityUpdate) {
await runResume.requireDurableInputCheckpoint([
...agent!.state.messages,
freshPromptMessage,
]);
}
run =
let run =
shouldPromptAgent && !capacityUpdate
? agent!.prompt(freshPromptMessage)
: agent!.continue();
Expand Down
7 changes: 1 addition & 6 deletions packages/junior/src/chat/agent/tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,10 +25,7 @@ import {
} from "@/chat/agent/sandbox";
import { SkillSandbox } from "@/chat/sandbox/skill-sandbox";
import type { Skill, SkillMetadata } from "@/chat/skills";
import {
createPluginHookRunner,
type PluginHookRunner,
} from "@/chat/plugins/agent-hooks";
import { createPluginHookRunner } from "@/chat/plugins/agent-hooks";
import { pluginCatalogRuntime } from "@/chat/plugins/catalog-runtime";
import { McpToolManager } from "@/chat/mcp/tool-manager";
import {
Expand Down Expand Up @@ -216,7 +213,6 @@ export interface ToolWiring {
agentTools: AgentTool[];
getPendingAuthPause: () => AuthorizationPauseError | undefined;
mcpToolManager: McpToolManager;
pluginHooks: PluginHookRunner;
/** Project core-owned review state into one durable Pi tool result. */
projectActionReviewResult<
TResult extends { details?: unknown; isError?: boolean },
Expand Down Expand Up @@ -565,7 +561,6 @@ export async function wireAgentTools(
agentTools,
getPendingAuthPause,
mcpToolManager,
pluginHooks,
projectActionReviewResult(toolCallId, result) {
return actionReview.projectToolResult(toolCallId, result);
},
Expand Down
10 changes: 0 additions & 10 deletions packages/junior/src/chat/pi/messages.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,13 +12,3 @@ export const piMessageSchema = z

/** Durable Pi transcript message stored across turns. */
export type PiMessage = z.output<typeof piMessageSchema>;

/** Reporting transcript entries only render messages with structured content parts. */
export const piContentMessageSchema = z
.object({
content: z.array(z.unknown()),
role: z.string().min(1),
})
.passthrough()
// @ts-expect-error non-overlapping boundary cast; rule forbids as-unknown-as chains
.transform((value) => value as PiMessage);
Loading