feat(coding-agent): added async compaction and in-place handoff
- Added compaction.asyncEnabled (Async Compaction, default on): when context enters the pre-threshold band [threshold - lead, threshold) with lead = clamp(threshold * 0.125, 8192, 32000), maintenance speculatively summarizes in the background off a branch snapshot (first configured LLM-backed method: remote, handoff, or soft) using a side session id isolated from the live turn. Crossing the threshold splices the armed result in instantly instead of blocking on a summarization round-trip. Armed results are invalidated by branch changes, reset boundaries, model switches that strand provider-native replay payloads, and context growth past keepRecentTokens (which re-speculates); extensions registering session_before_compact keep exact blocking semantics (speculation disabled). - Reworked handoff to commit in place: /handoff and the auto handoff method now write the generated document as a regular compaction entry on the current session (summary = document + <files> tag, cut from prepareCompaction) instead of starting a new session. SessionHandoff shrank to a document generator; session_before_switch/session_switch no longer fire with reason "handoff"; mid-turn maintenance no longer suppresses the handoff preference; overflow recovery can apply an armed handoff result. - Extracted the shared auto-compaction commit tail (#commitAutoCompactionResult / #commitCompactionEntry) used by the blocking production path, the armed speculative apply, manual compaction, and manual handoff. - Status line pulses the auto-compact icon while a speculation runs and holds it in accent once a result is armed. - Exported remotePreserveReusable from pi-agent-core/compaction for apply-time validation of speculative remote results.
This commit is contained in:
@@ -8,6 +8,7 @@
|
||||
- Added `tokenizer` to custom model and `modelOverrides` configuration. It overrides the catalog-resolved local tokenizer family for a model when a proxy serves a known model id with a different tokenizer.
|
||||
- Added `extendedContext` setting (`/settings` → Context → General, default on). When off, models with a premium long-context price tier (OpenAI GPT-5.6 Sol/Terra/Luna bill 2x input / 1.5x output above 272K input tokens, on both the API and subscription Codex) are capped at the standard-pricing threshold — they appear as 272K again and compaction fires before a request crosses into premium billing. Toggling mid-session re-clamps or restores the active model's window immediately. Anthropic Claude 4.6+ serves its full 1M window at standard pricing, so no Anthropic model is affected.
|
||||
- Added click-to-toggle and drag-to-reorder controls for list-valued `/settings` editors.
|
||||
- Added `compaction.asyncEnabled` (Async Compaction, default on): when context enters the band just below the compaction threshold, maintenance speculatively summarizes in the background off a branch snapshot (first configured LLM-backed method — remote, handoff, or soft — isolated from the live turn by a side session id) and holds the armed result; crossing the threshold then splices it in instantly instead of blocking on a summarization round-trip. Armed results are invalidated by branch changes, reset boundaries, model switches that strand provider-native replay payloads, and context growth past `keepRecentTokens` (which re-speculates). The status line pulses the auto-compact icon while a speculation runs and holds it in accent once a result is armed.
|
||||
|
||||
### Changed
|
||||
|
||||
@@ -17,6 +18,7 @@
|
||||
- The todo HUD header now draws a summed progress bar counting closed/total tasks across every stage. Once all tasks close, the bar smoothly collapses before the row disappears.
|
||||
- Token counting is now scoped to the model being billed rather than to a process-global tokenizer: session maintenance, stats, advisors, `/context`, snapcompact inline imaging, and `compress` each count through the owning agent's `Tokenizer` (`agent.tokenizer`). Message counting is `Tokenizer.countMessage`/`countMessages` (replacing the free `estimateTokens(message, tokenizer)` helper; the legacy shim keeps a compat `estimateTokens` export for legacy pi extensions). `estimateToolSchemaTokens`, `estimateSkillsTokens`, `computeNonMessageTokens`, and `computeNonMessageBreakdown` take an explicit tokenizer; standalone prompt inspection intentionally keeps the default estimate because it has no resolved catalog model.
|
||||
- The advisor runtime's `maintainContext` hook now receives the pending update as a message instead of a pre-computed token count — sizing it needs the advisor model's tokenizer, which the host owns.
|
||||
- Handoff no longer starts a new session: `/handoff` and the auto-maintenance `handoff` method now commit the generated document as a regular compaction entry on the current session (document becomes the summary, recent history is kept per `compaction.keepRecentTokens`, session id/transcript/cache key unchanged). The `session_before_switch`/`session_switch` extension events no longer fire with reason `"handoff"`, mid-turn maintenance no longer skips the handoff preference, and overflow recovery can now apply a pre-armed handoff result.
|
||||
|
||||
### Fixed
|
||||
|
||||
|
||||
@@ -2274,6 +2274,18 @@ export const SETTINGS_SCHEMA = {
|
||||
},
|
||||
},
|
||||
|
||||
"compaction.asyncEnabled": {
|
||||
type: "boolean",
|
||||
default: true,
|
||||
ui: {
|
||||
tab: "context",
|
||||
group: "Compaction",
|
||||
label: "Async Compaction",
|
||||
description:
|
||||
"Speculatively summarize in the background as context nears the compaction threshold, then splice the ready result in when the threshold is crossed",
|
||||
},
|
||||
},
|
||||
|
||||
// No default: an unset reserve tells the compaction layer the user never
|
||||
// chose one, so small-window recovery may swap in the proportional reserve
|
||||
// (see resolveBudgetReserveTokens). A materialized 16384 here would make
|
||||
@@ -5715,6 +5727,7 @@ export interface CompactionSettings {
|
||||
reserveTokens: number | undefined;
|
||||
keepRecentTokens: number;
|
||||
midTurnEnabled: boolean;
|
||||
asyncEnabled: boolean;
|
||||
handoffSaveToDisk: boolean;
|
||||
autoContinue: boolean;
|
||||
remoteEndpoint: string | undefined;
|
||||
|
||||
@@ -33,7 +33,7 @@ export interface SessionStartEvent {
|
||||
export interface SessionBeforeSwitchEvent {
|
||||
type: "session_before_switch";
|
||||
/** Reason for the switch */
|
||||
reason: "new" | "resume" | "fork" | "handoff";
|
||||
reason: "new" | "resume" | "fork";
|
||||
/** Session file we're switching to (only for "resume") */
|
||||
targetSessionFile?: string;
|
||||
}
|
||||
@@ -42,7 +42,7 @@ export interface SessionBeforeSwitchEvent {
|
||||
export interface SessionSwitchEvent {
|
||||
type: "session_switch";
|
||||
/** Reason for the switch */
|
||||
reason: "new" | "resume" | "fork" | "handoff";
|
||||
reason: "new" | "resume" | "fork";
|
||||
/** Session file we came from */
|
||||
previousSessionFile: string | undefined;
|
||||
}
|
||||
|
||||
@@ -324,6 +324,9 @@ export class StatusLineComponent implements Component {
|
||||
#onBranchChange: (() => void) | null = null;
|
||||
#disposed = false;
|
||||
#autoCompactEnabled: boolean = true;
|
||||
/** Pulse timer for the running-speculation indicator; live only while speculation runs. */
|
||||
#speculationBlinkTimer: NodeJS.Timeout | undefined;
|
||||
#speculationBlinkOn = true;
|
||||
#hookStatuses: Map<string, string> = new Map();
|
||||
#subagentCount: number = 0;
|
||||
/**
|
||||
@@ -694,12 +697,38 @@ export class StatusLineComponent implements Component {
|
||||
this.#branchResolveActive = undefined;
|
||||
this.#resetJjRequests();
|
||||
this.#onBranchChange = null;
|
||||
this.#stopSpeculationBlink();
|
||||
this.#clearUsageStartTimer();
|
||||
this.#onCodexResetFireworks = undefined;
|
||||
this.#codexResetSnapshots.clear();
|
||||
this.#retireGitWatcher();
|
||||
}
|
||||
|
||||
/**
|
||||
* Drive the context segment's pulse while a background speculative
|
||||
* compaction runs: a slow toggle that invalidates and repaints through the
|
||||
* same host callback async git/PR resolves use. Stops (and resets phase)
|
||||
* the first render after speculation leaves the running state.
|
||||
*/
|
||||
#syncSpeculationBlink(state: "idle" | "running" | "armed"): void {
|
||||
if (state === "running" && !this.#disposed) {
|
||||
this.#speculationBlinkTimer ??= setInterval(() => {
|
||||
this.#speculationBlinkOn = !this.#speculationBlinkOn;
|
||||
this.invalidate();
|
||||
this.#onBranchChange?.();
|
||||
}, 600);
|
||||
return;
|
||||
}
|
||||
this.#stopSpeculationBlink();
|
||||
}
|
||||
|
||||
#stopSpeculationBlink(): void {
|
||||
if (!this.#speculationBlinkTimer) return;
|
||||
clearInterval(this.#speculationBlinkTimer);
|
||||
this.#speculationBlinkTimer = undefined;
|
||||
this.#speculationBlinkOn = true;
|
||||
}
|
||||
|
||||
#clearUsageStartTimer(): void {
|
||||
if (!this.#usageStartTimer) return;
|
||||
clearTimeout(this.#usageStartTimer);
|
||||
@@ -1599,6 +1628,8 @@ export class StatusLineComponent implements Component {
|
||||
this.#getGitStatus(activeRepoCache.effectiveGitCwd))
|
||||
: null;
|
||||
const gitPr = includePr ? this.#lookupPr(activeRepoCache.effectiveGitCwd) : null;
|
||||
const compactionSpeculation = this.session.compactionSpeculation ?? "idle";
|
||||
this.#syncSpeculationBlink(compactionSpeculation);
|
||||
return {
|
||||
session: this.session,
|
||||
focusedAgentId: this.#focusedAgentId,
|
||||
@@ -1621,6 +1652,8 @@ export class StatusLineComponent implements Component {
|
||||
contextTokens,
|
||||
contextWindow,
|
||||
autoCompactEnabled: this.#autoCompactEnabled,
|
||||
compactionSpeculation,
|
||||
speculationBlinkOn: this.#speculationBlinkOn,
|
||||
subagentCount: this.#subagentCount,
|
||||
activeMs: this.getActiveMs(),
|
||||
git: {
|
||||
|
||||
@@ -458,11 +458,24 @@ const contextPctSegment: StatusLineSegment = {
|
||||
const pct = ctx.contextPercent;
|
||||
const window = ctx.contextWindow;
|
||||
|
||||
const autoIcon = ctx.autoCompactEnabled && theme.icon.auto ? ` ${theme.icon.auto}` : "";
|
||||
const text = `${formatContextUsage(pct, window, ctx.contextTokens)}${autoIcon}`;
|
||||
|
||||
const color = getContextUsageThemeColor(getContextUsageLevel(pct ?? 0, window));
|
||||
const content = withIcon(theme.icon.context, theme.fg(color, text));
|
||||
// Async-compaction indicator: pulse the auto icon while a background
|
||||
// speculation runs, hold it in accent once a result is armed.
|
||||
let autoIcon = "";
|
||||
if (ctx.autoCompactEnabled && theme.icon.auto) {
|
||||
const speculation = ctx.compactionSpeculation;
|
||||
const iconColor =
|
||||
speculation === "running"
|
||||
? ctx.speculationBlinkOn
|
||||
? "accent"
|
||||
: "muted"
|
||||
: speculation === "armed"
|
||||
? "accent"
|
||||
: color;
|
||||
autoIcon = ` ${theme.fg(iconColor, theme.icon.auto)}`;
|
||||
}
|
||||
const text = theme.fg(color, formatContextUsage(pct, window, ctx.contextTokens));
|
||||
const content = withIcon(theme.icon.context, `${text}${autoIcon}`);
|
||||
|
||||
return { content, visible: true };
|
||||
},
|
||||
|
||||
@@ -97,6 +97,10 @@ export interface SegmentContext {
|
||||
contextTokens: number;
|
||||
contextWindow: number;
|
||||
autoCompactEnabled: boolean;
|
||||
/** Background speculative-compaction state (async compaction). */
|
||||
compactionSpeculation: "idle" | "running" | "armed";
|
||||
/** Blink phase for the running-speculation pulse; toggled by the component's timer. */
|
||||
speculationBlinkOn: boolean;
|
||||
subagentCount: number;
|
||||
/**
|
||||
* Active processing time accumulated this session, in ms — the union of
|
||||
|
||||
@@ -1412,7 +1412,8 @@ export class CommandController {
|
||||
this.ctx.ui.requestRender();
|
||||
|
||||
try {
|
||||
// Handoff generation runs as a oneshot request; the new session is shown after it completes.
|
||||
// Handoff generation runs as a oneshot request; the document is then
|
||||
// committed as a compaction entry on this session.
|
||||
const result = await this.ctx.session.handoff(customInstructions);
|
||||
|
||||
if (!result) {
|
||||
@@ -1420,7 +1421,7 @@ export class CommandController {
|
||||
return;
|
||||
}
|
||||
|
||||
// Rebuild chat from the new session (which now contains the handoff document).
|
||||
// Rebuild chat from the session, which now shows the handoff compaction divider.
|
||||
this.ctx.clearTransientSessionUi();
|
||||
await this.ctx.renderInitialMessages();
|
||||
this.ctx.statusLine.invalidate();
|
||||
@@ -1429,7 +1430,11 @@ export class CommandController {
|
||||
|
||||
this.ctx.present([
|
||||
new Spacer(1),
|
||||
new Text(`${theme.fg("accent", `${theme.status.success} New session started with handoff context`)}`, 1, 1),
|
||||
new Text(
|
||||
`${theme.fg("accent", `${theme.status.success} Context handed off and compacted in place`)}`,
|
||||
1,
|
||||
1,
|
||||
),
|
||||
]);
|
||||
if (result.savedPath) {
|
||||
this.ctx.showStatus(`Handoff document saved to: ${result.savedPath}`);
|
||||
|
||||
@@ -336,7 +336,6 @@ export interface HandoffResult {
|
||||
export interface SessionHandoffOptions {
|
||||
autoTriggered?: boolean;
|
||||
signal?: AbortSignal;
|
||||
onSwitchCancelled?: () => void;
|
||||
}
|
||||
|
||||
/** Result from cycleModel(). */
|
||||
|
||||
@@ -1567,7 +1567,8 @@ export class AgentSession {
|
||||
getContextUsage: options => this.getContextUsage(options),
|
||||
shake: (mode, options) => this.shake(mode, options),
|
||||
dropImages: () => this.dropImages(),
|
||||
runHandoff: (customInstructions, options) => this.handoff(customInstructions, options),
|
||||
generateHandoffDocument: (customInstructions, options) =>
|
||||
this.#handoff.generateDocument(customInstructions, options),
|
||||
removeAssistantMessageFromActiveContext: message =>
|
||||
this.#recovery.removeAssistantMessageFromActiveContext(message),
|
||||
dropPersistedAssistantTurn: message => this.#recovery.dropPersistedAssistantTurn(message),
|
||||
@@ -1585,15 +1586,12 @@ export class AgentSession {
|
||||
sessionManager: this.sessionManager,
|
||||
settings: this.settings,
|
||||
modelRegistry: this.#modelRegistry,
|
||||
extensionRunner: this.#extensionRunner,
|
||||
sideStreamFn: this.#sideStreamFn,
|
||||
obfuscator: this.#obfuscator,
|
||||
model: () => this.model,
|
||||
thinkingLevel: () => this.thinkingLevel,
|
||||
sessionId: () => this.sessionId,
|
||||
sessionFile: () => this.sessionFile,
|
||||
baseSystemPrompt: () => this.#tools.baseSystemPrompt,
|
||||
assertVibeSessionTransitionAllowed: action => this.#assertVibeSessionTransitionAllowed(action),
|
||||
setSkipPostTurnMaintenance: timestamp => {
|
||||
this.#maintenance.skipPostTurnMaintenanceAssistantTimestamp = timestamp;
|
||||
},
|
||||
@@ -1602,32 +1600,6 @@ export class AgentSession {
|
||||
convertMessagesToLlm: (messages, signal) => this.convertMessagesToLlm(messages, signal),
|
||||
prepareSimpleStreamOptions: (options, provider) => this.prepareSimpleStreamOptions(options, provider),
|
||||
effectiveServiceTier: model => this.#models.effectiveServiceTier(model),
|
||||
flushPendingBash: () => this.#bash.flushPending(),
|
||||
beginBashSessionTransition: () => this.#bash.beginSessionTransition(),
|
||||
markBashSessionTransition: transition => this.#bash.markSessionTransition(transition),
|
||||
finishBashSessionTransition: (transition, success) => this.#bash.finishSessionTransition(transition, success),
|
||||
cancelOwnAsyncJobs: () => this.#cancelOwnAsyncJobs(),
|
||||
clearCheckpointRuntimeState: () => this.#clearCheckpointRuntimeState(),
|
||||
clearSessionScopedToolState: () => this.#clearSessionScopedToolState(),
|
||||
clearFreshProviderSessionId: () => {
|
||||
this.#freshProviderSessionId = undefined;
|
||||
},
|
||||
syncAgentSessionId: () => this.#syncAgentSessionId(),
|
||||
rekeyMemoryForCurrentSessionId: () => {
|
||||
this.#memory.rekeyForCurrentSessionId();
|
||||
},
|
||||
resetMemoryContextForNewTranscript: () => this.#memory.resetContextForNewTranscript(),
|
||||
clearPendingNextTurnMessages: () => {
|
||||
this.#pendingNextTurnMessages = [];
|
||||
this.#scheduledHiddenNextTurnGeneration = undefined;
|
||||
},
|
||||
resetTodoCycle: () => this.#todo.resetCycle(),
|
||||
buildDisplaySessionContext: () => this.buildDisplaySessionContext(),
|
||||
resetAdvisorSessionState: () => this.#advisors.resetSessionState(),
|
||||
drainAndDetachAdvisorRecorders: () => this.#advisors.drainAndDetachRecorders(),
|
||||
reattachAdvisorRecorderFeeds: () => this.#advisors.reattachRecorderFeeds(),
|
||||
clearAdvisorCost: () => this.#advisors.clearCost(),
|
||||
syncTodoPhasesFromBranch: () => this.#todo.syncFromBranch(),
|
||||
};
|
||||
this.#handoff = new SessionHandoff(handoffHost);
|
||||
|
||||
@@ -4089,6 +4061,7 @@ export class AgentSession {
|
||||
await cleanupEmptyMoveSession(this.sessionManager, this.#movedFromEmptySessionFile);
|
||||
this.#movedFromEmptySessionFile = undefined;
|
||||
this.#closeAllProviderSessions("dispose");
|
||||
this.#maintenance.cancelSpeculation();
|
||||
this.setHindsightSessionState(undefined);
|
||||
hindsightState?.dispose();
|
||||
this.#disconnectFromAgent();
|
||||
@@ -4711,6 +4684,10 @@ export class AgentSession {
|
||||
return this.#maintenance.isCompacting;
|
||||
}
|
||||
|
||||
/** Background speculative-compaction state, for UI indicators. */
|
||||
get compactionSpeculation(): "idle" | "running" | "armed" {
|
||||
return this.#maintenance.speculationState;
|
||||
}
|
||||
/** Strip image content from the current branch and persist the rewrite. */
|
||||
dropImages(): Promise<{ removed: number }> {
|
||||
return this.#maintenance.dropImages();
|
||||
@@ -7121,14 +7098,16 @@ export class AgentSession {
|
||||
}
|
||||
|
||||
/**
|
||||
* Generate a handoff document with a oneshot LLM call, then start a new session with it.
|
||||
* Generate a handoff document with a oneshot LLM call and commit it as a
|
||||
* compaction entry on the current session (the document becomes the summary;
|
||||
* recent history is kept).
|
||||
*
|
||||
* @param customInstructions Optional focus for the handoff document
|
||||
* @param options Handoff execution options
|
||||
* @returns The handoff document text, or undefined if cancelled/failed
|
||||
*/
|
||||
handoff(customInstructions?: string, options?: SessionHandoffOptions): Promise<HandoffResult | undefined> {
|
||||
return this.#handoff.handoff(customInstructions, options);
|
||||
return this.#maintenance.handoff(customInstructions, options);
|
||||
}
|
||||
|
||||
#isTerminalYieldToolResult(event: { toolName: string; isError?: boolean; result?: { details?: unknown } }): boolean {
|
||||
@@ -9316,9 +9295,9 @@ export class AgentSession {
|
||||
}
|
||||
/**
|
||||
* Get text content of the most recent visible handoff message.
|
||||
* Fresh handoff sessions store the handoff context as a custom message, not
|
||||
* an assistant message, so callers that copy the "last" message can use this
|
||||
* as a fallback before the new session has an assistant response.
|
||||
* Sessions created by older versions injected the handoff document as a
|
||||
* custom message at the top of a fresh session; callers that copy the
|
||||
* "last" message use this as a fallback while no assistant response exists.
|
||||
*/
|
||||
getLastVisibleHandoffText(): string | undefined {
|
||||
for (let i = this.messages.length - 1; i >= 0; i--) {
|
||||
|
||||
@@ -15,7 +15,7 @@ export const COMPACTION_METHOD_CHOICES = [
|
||||
{
|
||||
value: "handoff",
|
||||
label: "Handoff",
|
||||
description: "Generate a handoff document and continue in a new session",
|
||||
description: "Generate a handoff document and continue from it as the compaction summary",
|
||||
},
|
||||
{
|
||||
value: "soft",
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
/** Handoff generation and session transition orchestration. */
|
||||
/** Handoff document generation. Committing the document as a compaction entry is owned by SessionMaintenance. */
|
||||
|
||||
import * as path from "node:path";
|
||||
import {
|
||||
@@ -13,19 +13,11 @@ import type { Message, Model, ServiceTier, SimpleStreamOptions } from "@oh-my-pi
|
||||
import { logger, Snowflake } from "@oh-my-pi/pi-utils";
|
||||
import type { ModelRegistry } from "../config/model-registry";
|
||||
import type { Settings } from "../config/settings";
|
||||
import type { ExtensionRunner, SessionBeforeSwitchResult } from "../extensibility/extensions";
|
||||
import { copyLocalArtifacts, resolveLocalUrlToPath } from "../internal-urls";
|
||||
import { obfuscateProviderContext } from "../secrets/message-transform";
|
||||
import type { SecretObfuscator } from "../secrets/obfuscator";
|
||||
import type { HandoffResult, SessionHandoffOptions } from "./agent-session-types";
|
||||
import type { BashSessionTransition } from "./bash-runner";
|
||||
import type { SessionContext } from "./session-context";
|
||||
import type { SessionManager } from "./session-manager";
|
||||
|
||||
function createHandoffContext(document: string): string {
|
||||
return `<handoff-context>\n${document}\n</handoff-context>\n\nThe above is a handoff document from a previous session. Use this context to continue the work seamlessly.`;
|
||||
}
|
||||
|
||||
function createHandoffFileName(date = new Date()): string {
|
||||
const fileTimestamp = date.toISOString().replace(/[:.]/g, "-");
|
||||
return `handoff-${fileTimestamp}.md`;
|
||||
@@ -48,43 +40,21 @@ export interface SessionHandoffHost {
|
||||
sessionManager: SessionManager;
|
||||
settings: Settings;
|
||||
modelRegistry: ModelRegistry;
|
||||
extensionRunner: ExtensionRunner | undefined;
|
||||
sideStreamFn: StreamFn;
|
||||
obfuscator: SecretObfuscator | undefined;
|
||||
model(): Model | undefined;
|
||||
thinkingLevel(): ThinkingLevel | undefined;
|
||||
sessionId(): string;
|
||||
sessionFile(): string | undefined;
|
||||
baseSystemPrompt(): string[];
|
||||
assertVibeSessionTransitionAllowed(action: string): void;
|
||||
setSkipPostTurnMaintenance(timestamp: number | undefined): void;
|
||||
obfuscateTextForProvider(text: string | undefined): string | undefined;
|
||||
deobfuscateFromProvider(text: string): string;
|
||||
convertMessagesToLlm(messages: AgentMessage[], signal?: AbortSignal): Promise<Message[]>;
|
||||
prepareSimpleStreamOptions(options: SimpleStreamOptions, provider?: string): SimpleStreamOptions;
|
||||
effectiveServiceTier(model: Model | undefined): ServiceTier | undefined;
|
||||
flushPendingBash(): Promise<void>;
|
||||
beginBashSessionTransition(): BashSessionTransition;
|
||||
markBashSessionTransition(transition: BashSessionTransition): void;
|
||||
finishBashSessionTransition(transition: BashSessionTransition, success: boolean): void;
|
||||
cancelOwnAsyncJobs(): void;
|
||||
clearCheckpointRuntimeState(): void;
|
||||
clearSessionScopedToolState(): void;
|
||||
clearFreshProviderSessionId(): void;
|
||||
syncAgentSessionId(): void;
|
||||
rekeyMemoryForCurrentSessionId(): void;
|
||||
resetMemoryContextForNewTranscript(): Promise<void>;
|
||||
clearPendingNextTurnMessages(): void;
|
||||
resetTodoCycle(): void;
|
||||
buildDisplaySessionContext(): SessionContext;
|
||||
resetAdvisorSessionState(): void;
|
||||
drainAndDetachAdvisorRecorders(): Promise<void>;
|
||||
reattachAdvisorRecorderFeeds(): void;
|
||||
clearAdvisorCost(): void;
|
||||
syncTodoPhasesFromBranch(): void;
|
||||
}
|
||||
|
||||
/** Generates handoff documents and owns the handoff session transition. */
|
||||
/** Generates handoff documents with a cache-friendly oneshot LLM call. */
|
||||
export class SessionHandoff {
|
||||
#handoffAbortController: AbortController | undefined;
|
||||
readonly #host: SessionHandoffHost;
|
||||
@@ -107,21 +77,22 @@ export class SessionHandoff {
|
||||
}
|
||||
|
||||
/**
|
||||
* Generate a handoff document with a oneshot LLM call, then start a new session with it.
|
||||
* Generate a handoff document with a oneshot LLM call.
|
||||
*
|
||||
* The request is built through the same pipeline a live turn uses so the
|
||||
* oneshot reads the provider prompt cache the main turn populated. The
|
||||
* caller (SessionMaintenance) commits the returned document as a compaction
|
||||
* entry; this method rewrites no history.
|
||||
*
|
||||
* @param customInstructions Optional focus for the handoff document
|
||||
* @param options Handoff execution options
|
||||
* @returns The handoff document text, or undefined if cancelled/failed
|
||||
* @returns The handoff document text, or undefined when an auto-triggered
|
||||
* generation produced no content (manual generation throws instead)
|
||||
*/
|
||||
async handoff(customInstructions?: string, options?: SessionHandoffOptions): Promise<HandoffResult | undefined> {
|
||||
this.#host.assertVibeSessionTransitionAllowed("handoff to a new session");
|
||||
const entries = this.#host.sessionManager.getBranch();
|
||||
const messageCount = entries.filter(e => e.type === "message").length;
|
||||
|
||||
if (messageCount < 2) {
|
||||
throw new Error("Nothing to hand off (no messages yet)");
|
||||
}
|
||||
|
||||
async generateDocument(
|
||||
customInstructions?: string,
|
||||
options?: SessionHandoffOptions,
|
||||
): Promise<HandoffResult | undefined> {
|
||||
this.#host.setSkipPostTurnMaintenance(undefined);
|
||||
|
||||
this.#handoffAbortController = new AbortController();
|
||||
@@ -140,8 +111,6 @@ export class SessionHandoff {
|
||||
}
|
||||
}
|
||||
|
||||
let advisorRecordersDetached = false;
|
||||
let sessionTransitioned = false;
|
||||
try {
|
||||
throwIfHandoffAborted(handoffSignal);
|
||||
|
||||
@@ -180,8 +149,8 @@ export class SessionHandoff {
|
||||
];
|
||||
const handoffLlmMessages = await this.#host.convertMessagesToLlm(handoffSnapshot, handoffSignal);
|
||||
// Base system prompt, not a per-turn `before_agent_start` hook override —
|
||||
// the handoff seeds a fresh session and must not carry prompt-specific
|
||||
// hook state. Matches the prompt the old handoff path sent.
|
||||
// the document seeds the post-compaction context and must not carry
|
||||
// prompt-specific hook state.
|
||||
const handoffContext = await this.#host.agent.buildSideRequestContext(
|
||||
handoffLlmMessages,
|
||||
this.#host.baseSystemPrompt(),
|
||||
@@ -229,94 +198,14 @@ export class SessionHandoff {
|
||||
autoTriggered: options?.autoTriggered ?? false,
|
||||
});
|
||||
// Auto-handoff is best-effort: returning undefined lets maintenance fall
|
||||
// back to context-full compaction. A user-initiated handoff must surface
|
||||
// the failure instead of a silent, misleading "cancelled".
|
||||
// back to the next compaction method. A user-initiated handoff must
|
||||
// surface the failure instead of a silent, misleading "cancelled".
|
||||
if (options?.autoTriggered) {
|
||||
return undefined;
|
||||
}
|
||||
throw new Error("Handoff generation produced no content");
|
||||
}
|
||||
|
||||
// Start a new session
|
||||
const previousSessionFile = this.#host.sessionFile();
|
||||
if (this.#host.extensionRunner?.hasHandlers("session_before_switch")) {
|
||||
const result = (await this.#host.extensionRunner.emit({
|
||||
type: "session_before_switch",
|
||||
reason: "handoff",
|
||||
})) as SessionBeforeSwitchResult | undefined;
|
||||
|
||||
if (result?.cancel) {
|
||||
options?.onSwitchCancelled?.();
|
||||
return undefined;
|
||||
}
|
||||
}
|
||||
await this.#host.flushPendingBash();
|
||||
await this.#host.sessionManager.flush();
|
||||
advisorRecordersDetached = true;
|
||||
// Stop and settle in-flight advisors while the old-session feeds can still
|
||||
// observe message_end, then mute before opening the replacement session.
|
||||
await this.#host.drainAndDetachAdvisorRecorders();
|
||||
// Snapshot the outgoing session's local:// root BEFORE newSession() mints a
|
||||
// fresh session id (and therefore a fresh, empty local root). The handoff
|
||||
// document routinely references plans/scratch files under local://, so those
|
||||
// artifacts must follow the session switch or every reference dangles.
|
||||
const localProtocolOptions = {
|
||||
getArtifactsDir: () => this.#host.sessionManager.getArtifactsDir(),
|
||||
getSessionId: () => this.#host.sessionManager.getSessionId(),
|
||||
};
|
||||
const previousLocalRoot = resolveLocalUrlToPath("local://", localProtocolOptions);
|
||||
const bashTransition = this.#host.beginBashSessionTransition();
|
||||
this.#host.cancelOwnAsyncJobs();
|
||||
try {
|
||||
await this.#host.sessionManager.newSession(
|
||||
previousSessionFile ? { parentSession: previousSessionFile } : undefined,
|
||||
);
|
||||
this.#host.markBashSessionTransition(bashTransition);
|
||||
// The handoff opens a fresh conversation, so the spend of the one it
|
||||
// summarizes stays with it. Clearing here, at the commit point, keeps the
|
||||
// status line honest even if a later step throws.
|
||||
this.#host.clearAdvisorCost();
|
||||
sessionTransitioned = true;
|
||||
} finally {
|
||||
this.#host.finishBashSessionTransition(bashTransition, sessionTransitioned);
|
||||
}
|
||||
|
||||
this.#host.clearSessionScopedToolState();
|
||||
|
||||
this.#host.clearCheckpointRuntimeState();
|
||||
// agent.reset() clears the core steering/follow-up queues. Preserve any queued
|
||||
// steers/follow-ups (RPC/SDK steer()/followUp() issued during the handoff, or a
|
||||
// pre-loader TUI steer) so they survive into the post-handoff session instead of
|
||||
// being silently dropped. Capture is synchronous immediately before reset and
|
||||
// restore is synchronous immediately after — no await gap — so a steer arriving
|
||||
// later (during ensureOnDisk/Bun.write below) appends to the restored queue
|
||||
// rather than being clobbered.
|
||||
const preservedSteering = this.#host.agent.peekSteeringQueue().slice();
|
||||
const preservedFollowUp = this.#host.agent.peekFollowUpQueue().slice();
|
||||
this.#host.agent.reset();
|
||||
this.#host.agent.replaceQueues(preservedSteering, preservedFollowUp);
|
||||
this.#host.clearFreshProviderSessionId();
|
||||
this.#host.syncAgentSessionId();
|
||||
this.#host.rekeyMemoryForCurrentSessionId();
|
||||
await this.#host.resetMemoryContextForNewTranscript();
|
||||
this.#host.clearPendingNextTurnMessages();
|
||||
this.#host.resetTodoCycle();
|
||||
|
||||
// Carry local:// artifacts into the replacement session (best-effort: the
|
||||
// switch is already committed, so a copy failure must not fail the handoff).
|
||||
try {
|
||||
const newLocalRoot = resolveLocalUrlToPath("local://", localProtocolOptions);
|
||||
await copyLocalArtifacts(previousLocalRoot, newLocalRoot);
|
||||
} catch (error) {
|
||||
logger.warn("Failed to copy local artifacts into handoff session", {
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
});
|
||||
}
|
||||
|
||||
// Inject the handoff document as a custom message
|
||||
const handoffContent = createHandoffContext(handoffText);
|
||||
this.#host.sessionManager.appendCustomMessageEntry("handoff", handoffContent, true, undefined, "agent");
|
||||
await this.#host.sessionManager.ensureOnDisk();
|
||||
let savedPath: string | undefined;
|
||||
if (options?.autoTriggered && this.#host.settings.get("compaction.handoffSaveToDisk")) {
|
||||
const artifactsDir = this.#host.sessionManager.getArtifactsDir();
|
||||
@@ -336,20 +225,6 @@ export class SessionHandoff {
|
||||
}
|
||||
}
|
||||
|
||||
// Rebuild agent messages from session
|
||||
const sessionContext = this.#host.buildDisplaySessionContext();
|
||||
this.#host.agent.replaceMessages(sessionContext.messages);
|
||||
this.#host.resetAdvisorSessionState();
|
||||
advisorRecordersDetached = false;
|
||||
this.#host.syncTodoPhasesFromBranch();
|
||||
if (this.#host.extensionRunner) {
|
||||
await this.#host.extensionRunner.emit({
|
||||
type: "session_switch",
|
||||
reason: "handoff",
|
||||
previousSessionFile,
|
||||
});
|
||||
}
|
||||
|
||||
return { document: handoffText, savedPath };
|
||||
} catch (error) {
|
||||
// Only a genuine cancellation (user Esc or an unreasoned source-signal
|
||||
@@ -358,10 +233,6 @@ export class SessionHandoff {
|
||||
throwIfHandoffAborted(handoffSignal);
|
||||
throw error;
|
||||
} finally {
|
||||
if (advisorRecordersDetached) {
|
||||
if (sessionTransitioned) this.#host.resetAdvisorSessionState();
|
||||
else this.#host.reattachAdvisorRecorderFeeds();
|
||||
}
|
||||
sourceSignal?.removeEventListener("abort", onSourceAbort);
|
||||
this.#handoffAbortController = undefined;
|
||||
}
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -623,57 +623,21 @@ describe("AgentSession advisor toggle", () => {
|
||||
await session.dispose();
|
||||
expect((await loadAdvisorTranscriptCosts(previousSessionFile)).get("")).toBeCloseTo(0.75, 8);
|
||||
});
|
||||
it("clears advisor cost when a handoff opens the replacement session", async () => {
|
||||
it("resets advisor runtimes after an in-place handoff compaction", async () => {
|
||||
vi.spyOn(compactionModule, "generateHandoffFromContext").mockResolvedValue("## Goal\nContinue from here");
|
||||
try {
|
||||
const advisor = enableAdvisor();
|
||||
prepareHandoffConversation(advisor);
|
||||
const previousSessionFile = session.sessionFile;
|
||||
const newSession = sessionManager.newSession.bind(sessionManager);
|
||||
vi.spyOn(sessionManager, "newSession").mockImplementation(async options => {
|
||||
const result = await newSession(options);
|
||||
// The outgoing advisor finalizes after the replacement file is selected.
|
||||
appendAdvisorCost(advisor, 9, 3);
|
||||
return result;
|
||||
});
|
||||
const advisor = enableAdvisor();
|
||||
prepareHandoffConversation(advisor);
|
||||
session.settings.set("compaction.keepRecentTokens", 1);
|
||||
const sessionFile = session.sessionFile;
|
||||
|
||||
await session.handoff();
|
||||
const result = await session.handoff();
|
||||
|
||||
// The handoff hands the work over to a fresh conversation, so the spend of
|
||||
// the one it summarizes must not follow it.
|
||||
expect(session.sessionFile).not.toBe(previousSessionFile);
|
||||
expect(session.getAdvisorCost()).toBe(0);
|
||||
const replacementSessionFile = session.sessionFile;
|
||||
if (!replacementSessionFile) throw new Error("Expected the replacement session to be persisted");
|
||||
appendAdvisorCost(advisor, 0.25, 4);
|
||||
expect(session.getAdvisorCost()).toBeCloseTo(0.25, 8);
|
||||
await session.dispose();
|
||||
expect((await loadAdvisorTranscriptCosts(replacementSessionFile)).get("")).toBeCloseTo(0.25, 8);
|
||||
} finally {
|
||||
vi.restoreAllMocks();
|
||||
}
|
||||
});
|
||||
it("restores advisor recording when a handoff fails before replacing the session", async () => {
|
||||
vi.spyOn(compactionModule, "generateHandoffFromContext").mockResolvedValue("## Goal\nContinue from here");
|
||||
try {
|
||||
const advisor = enableAdvisor();
|
||||
prepareHandoffConversation(advisor);
|
||||
const previousSessionFile = session.sessionFile;
|
||||
if (!previousSessionFile) throw new Error("Expected the previous session to be persisted");
|
||||
const failure = new Error("replacement session failed");
|
||||
vi.spyOn(sessionManager, "newSession").mockRejectedValue(failure);
|
||||
|
||||
await expect(session.handoff()).rejects.toThrow(failure);
|
||||
|
||||
expect(session.sessionFile).toBe(previousSessionFile);
|
||||
expect(session.getAdvisorCost()).toBeCloseTo(0.5, 8);
|
||||
appendAdvisorCost(advisor, 0.25, 3);
|
||||
expect(session.getAdvisorCost()).toBeCloseTo(0.75, 8);
|
||||
await session.dispose();
|
||||
expect((await loadAdvisorTranscriptCosts(previousSessionFile)).get("")).toBeCloseTo(0.75, 8);
|
||||
} finally {
|
||||
vi.restoreAllMocks();
|
||||
}
|
||||
expect(result?.document).toContain("Continue from here");
|
||||
expect(session.sessionFile).toBe(sessionFile);
|
||||
const compaction = sessionManager.getBranch().at(-1);
|
||||
expect(compaction).toMatchObject({ type: "compaction" });
|
||||
if (compaction?.type !== "compaction") throw new Error("Expected handoff compaction entry");
|
||||
expect(compaction.summary).toContain("Continue from here");
|
||||
});
|
||||
it("clears advisor cost when a branch skips conversation restore", async () => {
|
||||
const extensionRunner = {
|
||||
|
||||
@@ -50,6 +50,7 @@ describe("AgentSession auto-compaction progress guard", () => {
|
||||
let sessionManager: SessionManager;
|
||||
let authStorage: AuthStorage;
|
||||
let modelRegistry: ModelRegistry;
|
||||
let compactHookEnabled = true;
|
||||
|
||||
const NOTICE_SOURCE = "compaction";
|
||||
const NO_PROGRESS_FRAGMENT = "Compaction freed too little context to make progress";
|
||||
@@ -61,6 +62,7 @@ describe("AgentSession auto-compaction progress guard", () => {
|
||||
});
|
||||
|
||||
beforeEach(() => {
|
||||
compactHookEnabled = true;
|
||||
sessionManager = SessionManager.inMemory();
|
||||
|
||||
// The progress-guard tests exercise AgentSession's post-compaction state
|
||||
@@ -68,7 +70,7 @@ describe("AgentSession auto-compaction progress guard", () => {
|
||||
// while returning the same short-circuit result without compiling a
|
||||
// temporary extension for every test.
|
||||
const extensionRunner = {
|
||||
hasHandlers: (type: string) => type === "session_before_compact",
|
||||
hasHandlers: (type: string) => compactHookEnabled && type === "session_before_compact",
|
||||
emit: async (event: { type: string; preparation?: CompactionPreparation }) => {
|
||||
if (event.type !== "session_before_compact" || !event.preparation) return undefined;
|
||||
return {
|
||||
@@ -858,12 +860,34 @@ describe("AgentSession auto-compaction progress guard", () => {
|
||||
);
|
||||
});
|
||||
|
||||
it("does not restore a length stop after handoff recovery commits", async () => {
|
||||
it("drops a length stop and retries after handoff recovery commits", async () => {
|
||||
session.settings.set("compaction.methodOrder", ["handoff", "soft"]);
|
||||
session.settings.set("compaction.enabled", true);
|
||||
session.settings.set("compaction.keepRecentTokens", 1);
|
||||
compactHookEnabled = false;
|
||||
sessionManager.appendMessage({
|
||||
role: "assistant",
|
||||
content: [{ type: "text", text: "completed seed" }],
|
||||
api: "anthropic-messages",
|
||||
provider: "anthropic",
|
||||
model: "claude-sonnet-4-5",
|
||||
stopReason: "stop",
|
||||
usage: {
|
||||
input: 1,
|
||||
output: 1,
|
||||
cacheRead: 0,
|
||||
cacheWrite: 0,
|
||||
totalTokens: 2,
|
||||
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
||||
},
|
||||
timestamp: Date.now(),
|
||||
});
|
||||
session.settings.set("contextPromotion.enabled", false);
|
||||
const promptSpy = vi.spyOn(session.agent, "prompt").mockResolvedValue(undefined as never);
|
||||
const continueSpy = vi.spyOn(session.agent, "continue").mockResolvedValue();
|
||||
const handoffSpy = vi.spyOn(session, "handoff").mockResolvedValue({ document: "handoff document" });
|
||||
const generateHandoffSpy = vi
|
||||
.spyOn(compactionModule, "generateHandoffFromContext")
|
||||
.mockResolvedValue("handoff document");
|
||||
|
||||
const { promise: compactionDone, resolve: onCompactionDone } = Promise.withResolvers<void>();
|
||||
session.subscribe(event => {
|
||||
@@ -887,15 +911,20 @@ describe("AgentSession auto-compaction progress guard", () => {
|
||||
},
|
||||
timestamp: Date.now(),
|
||||
};
|
||||
sessionManager.appendMessage(assistantMsg);
|
||||
session.agent.emitExternalEvent({ type: "message_end", message: assistantMsg });
|
||||
session.agent.emitExternalEvent({ type: "agent_end", messages: [assistantMsg] });
|
||||
|
||||
await compactionDone;
|
||||
await session.waitForIdle();
|
||||
|
||||
expect(promptSpy).toHaveBeenCalledTimes(1);
|
||||
expect(handoffSpy).toHaveBeenCalledTimes(1);
|
||||
expect(continueSpy).not.toHaveBeenCalled();
|
||||
expect(promptSpy).not.toHaveBeenCalled();
|
||||
expect(generateHandoffSpy).toHaveBeenCalledTimes(1);
|
||||
expect(continueSpy).toHaveBeenCalledTimes(1);
|
||||
expect(sessionManager.getBranch().at(-1)).toMatchObject({
|
||||
type: "compaction",
|
||||
summary: "handoff document",
|
||||
});
|
||||
expect(sessionManager.getBranch()).not.toContainEqual(
|
||||
expect.objectContaining({
|
||||
type: "message",
|
||||
|
||||
@@ -220,20 +220,6 @@ describe("AgentSession mid-run threshold compaction", () => {
|
||||
expect(observedContexts[1].join("\n")).toContain("ACTIVE-GOAL-MID-RUN-COMPACTED");
|
||||
});
|
||||
|
||||
it("falls back to in-place compaction for mid-run handoff strategy", async () => {
|
||||
const { session, observedContexts } = await createHarness({ "compaction.methodOrder": ["handoff", "soft"] });
|
||||
const handoffSpy = vi.spyOn(session, "handoff").mockImplementation(async () => {
|
||||
throw new Error("mid-run compaction must not reset the session through handoff");
|
||||
});
|
||||
const compactSpy = mockCompaction("HANDOFF-MID-RUN-COMPACTED-IN-PLACE");
|
||||
|
||||
await session.prompt("work on the release");
|
||||
|
||||
expect(handoffSpy).not.toHaveBeenCalled();
|
||||
expect(compactSpy).toHaveBeenCalledTimes(1);
|
||||
expect(observedContexts[1].join("\n")).toContain("HANDOFF-MID-RUN-COMPACTED-IN-PLACE");
|
||||
});
|
||||
|
||||
it("does not wait for message persistence below the mid-run threshold", async () => {
|
||||
const releaseMessageEnd = Promise.withResolvers<void>();
|
||||
const messageEndEntered = Promise.withResolvers<void>();
|
||||
|
||||
@@ -105,6 +105,8 @@ describe("AgentSession handoff", () => {
|
||||
settings: Settings.isolated({
|
||||
"compaction.enabled": true,
|
||||
"compaction.autoContinue": false,
|
||||
"compaction.asyncEnabled": false,
|
||||
"compaction.keepRecentTokens": 1,
|
||||
}),
|
||||
modelRegistry,
|
||||
obfuscator,
|
||||
@@ -148,8 +150,10 @@ describe("AgentSession handoff", () => {
|
||||
vi.restoreAllMocks();
|
||||
});
|
||||
|
||||
it("does not run auto-compaction after handoff turn completes", async () => {
|
||||
it("commits a handoff document as an in-place compaction", async () => {
|
||||
const handoffText = "## Goal\nContinue from here";
|
||||
const previousSessionFile = session.sessionFile;
|
||||
const previousSessionId = session.sessionId;
|
||||
const generateHandoffSpy = vi
|
||||
.spyOn(compactionModule, "generateHandoffFromContext")
|
||||
.mockResolvedValue(handoffText);
|
||||
@@ -159,143 +163,15 @@ describe("AgentSession handoff", () => {
|
||||
|
||||
expect(generateHandoffSpy).toHaveBeenCalledTimes(1);
|
||||
expect(result?.document).toBe(handoffText);
|
||||
|
||||
expect(session.sessionFile).toBe(previousSessionFile);
|
||||
expect(session.sessionId).toBe(previousSessionId);
|
||||
const compaction = sessionManager.getBranch().at(-1);
|
||||
expect(compaction).toMatchObject({ type: "compaction" });
|
||||
if (compaction?.type !== "compaction") throw new Error("Expected handoff compaction entry");
|
||||
expect(compaction.summary).toContain(handoffText);
|
||||
expect(session.agent.state.messages.some(message => message.role === "compactionSummary")).toBe(true);
|
||||
expect(events.filter(event => event.type === "auto_compaction_start")).toHaveLength(0);
|
||||
expect(events.filter(event => event.type === "auto_compaction_end")).toHaveLength(0);
|
||||
expect(sessionManager.getEntries().filter(entry => entry.type === "compaction")).toHaveLength(0);
|
||||
});
|
||||
|
||||
it("clears staged preview state when handoff creates the replacement session", async () => {
|
||||
vi.spyOn(compactionModule, "generateHandoffFromContext").mockResolvedValue("## Goal\nContinue from here");
|
||||
session.toolChoiceQueue.registerPendingInvoker("old-session-preview", "ast_edit", async () => ({
|
||||
content: [{ type: "text", text: "applied old preview" }],
|
||||
}));
|
||||
expect(session.peekPendingInvoker()).toBeDefined();
|
||||
expect(session.nextToolChoiceDirective()).toBeDefined();
|
||||
|
||||
await session.handoff();
|
||||
|
||||
expect(session.peekPendingInvoker()).toBeUndefined();
|
||||
expect(session.nextToolChoiceDirective()).toBeUndefined();
|
||||
});
|
||||
|
||||
it("carries local:// artifacts into the handed-off session", async () => {
|
||||
// Handoff is a continuity operation: the generated document references
|
||||
// plans/scratch files the old session wrote under its local:// root. The
|
||||
// fresh session mints a new local root, so the artifacts must be copied
|
||||
// forward or every reference the handoff document carries dangles.
|
||||
vi.spyOn(compactionModule, "generateHandoffFromContext").mockResolvedValue("## Goal\nContinue from here");
|
||||
const localOptions = {
|
||||
getArtifactsDir: () => sessionManager.getArtifactsDir(),
|
||||
getSessionId: () => sessionManager.getSessionId(),
|
||||
};
|
||||
const oldLocalRoot = resolveLocalUrlToPath("local://", localOptions);
|
||||
const oldPlanPath = resolveLocalUrlToPath("local://my-plan.md", localOptions);
|
||||
const oldNestedPath = resolveLocalUrlToPath("local://research/notes.txt", localOptions);
|
||||
await fs.mkdir(path.dirname(oldNestedPath), { recursive: true });
|
||||
await Bun.write(oldPlanPath, "# Plan\n\nbody\n");
|
||||
await Bun.write(oldNestedPath, "scratch notes");
|
||||
|
||||
await session.handoff();
|
||||
|
||||
const newLocalRoot = resolveLocalUrlToPath("local://", localOptions);
|
||||
expect(newLocalRoot).not.toBe(oldLocalRoot);
|
||||
expect(await Bun.file(resolveLocalUrlToPath("local://my-plan.md", localOptions)).text()).toBe("# Plan\n\nbody\n");
|
||||
expect(await Bun.file(resolveLocalUrlToPath("local://research/notes.txt", localOptions)).text()).toBe(
|
||||
"scratch notes",
|
||||
);
|
||||
// The source session's artifacts remain untouched on disk.
|
||||
expect(await Bun.file(oldPlanPath).text()).toBe("# Plan\n\nbody\n");
|
||||
});
|
||||
|
||||
it("emits handoff lifecycle hooks on the outgoing and replacement sessions", async () => {
|
||||
// dispose() is terminal: it closes the manager and releases its in-memory
|
||||
// transcript. Reopen the persisted session file for the replacement
|
||||
// session, as production revival paths do.
|
||||
await session.dispose();
|
||||
const sessionFile = sessionManager.getSessionFile();
|
||||
if (!sessionFile) throw new Error("Expected a persisted session file");
|
||||
sessionManager = await SessionManager.open(sessionFile, tempDir.path());
|
||||
const extensionsResult = await loadExtensions([], tempDir.path());
|
||||
const extensionRunner = new ExtensionRunner(
|
||||
extensionsResult.extensions,
|
||||
extensionsResult.runtime,
|
||||
tempDir.path(),
|
||||
sessionManager,
|
||||
modelRegistry,
|
||||
);
|
||||
const observedEvents: Array<{
|
||||
type: "session_before_switch" | "session_switch";
|
||||
reason: string;
|
||||
previousSessionFile: string | undefined;
|
||||
activeSessionFile: string | undefined;
|
||||
messageCount: number;
|
||||
handoffEntryCount: number;
|
||||
}> = [];
|
||||
vi.spyOn(extensionRunner, "hasHandlers").mockImplementation(eventName => eventName === "session_before_switch");
|
||||
const emit = extensionRunner.emit.bind(extensionRunner);
|
||||
vi.spyOn(extensionRunner, "emit").mockImplementation(event => {
|
||||
if (event.type === "session_before_switch" || event.type === "session_switch") {
|
||||
observedEvents.push({
|
||||
type: event.type,
|
||||
reason: event.reason,
|
||||
previousSessionFile: event.type === "session_switch" ? event.previousSessionFile : undefined,
|
||||
activeSessionFile: session.sessionFile,
|
||||
messageCount: sessionManager.getBranch().filter(entry => entry.type === "message").length,
|
||||
handoffEntryCount: sessionManager
|
||||
.getBranch()
|
||||
.filter(entry => entry.type === "custom_message" && entry.customType === "handoff").length,
|
||||
});
|
||||
}
|
||||
return emit(event);
|
||||
});
|
||||
|
||||
session = new AgentSession({
|
||||
agent: new Agent({
|
||||
initialState: {
|
||||
model,
|
||||
systemPrompt: ["Test"],
|
||||
tools: [],
|
||||
messages: [],
|
||||
},
|
||||
}),
|
||||
sessionManager,
|
||||
settings: Settings.isolated({
|
||||
"compaction.enabled": true,
|
||||
"compaction.autoContinue": false,
|
||||
}),
|
||||
modelRegistry,
|
||||
extensionRunner,
|
||||
obfuscator,
|
||||
});
|
||||
const previousSessionFile = session.sessionFile;
|
||||
const generateHandoffSpy = vi
|
||||
.spyOn(compactionModule, "generateHandoffFromContext")
|
||||
.mockResolvedValue("## Goal\nContinue from here");
|
||||
|
||||
await session.handoff();
|
||||
|
||||
const nextSessionFile = session.sessionFile;
|
||||
expect(generateHandoffSpy).toHaveBeenCalledTimes(1);
|
||||
expect(nextSessionFile).not.toBe(previousSessionFile);
|
||||
expect(observedEvents).toEqual([
|
||||
{
|
||||
type: "session_before_switch",
|
||||
reason: "handoff",
|
||||
previousSessionFile: undefined,
|
||||
activeSessionFile: previousSessionFile,
|
||||
messageCount: 2,
|
||||
handoffEntryCount: 0,
|
||||
},
|
||||
{
|
||||
type: "session_switch",
|
||||
reason: "handoff",
|
||||
previousSessionFile,
|
||||
activeSessionFile: nextSessionFile,
|
||||
messageCount: 0,
|
||||
handoffEntryCount: 1,
|
||||
},
|
||||
]);
|
||||
});
|
||||
|
||||
it("runs handoff generation through the configured side stream function", async () => {
|
||||
@@ -345,6 +221,8 @@ describe("AgentSession handoff", () => {
|
||||
settings: Settings.isolated({
|
||||
"compaction.enabled": true,
|
||||
"compaction.autoContinue": false,
|
||||
"compaction.asyncEnabled": false,
|
||||
"compaction.keepRecentTokens": 1,
|
||||
}),
|
||||
modelRegistry,
|
||||
obfuscator,
|
||||
@@ -371,101 +249,6 @@ describe("AgentSession handoff", () => {
|
||||
expect(capturedSideSessionId).toStartWith(`${preHandoffSessionId}:side:`);
|
||||
});
|
||||
|
||||
it("preserves queued steering and follow-up messages across the handoff reset", async () => {
|
||||
// Defect 2: handoff() calls agent.reset(), which clears the core steering/follow-up
|
||||
// queues. Steers/follow-ups already queued (the mis-routed first compaction message,
|
||||
// or RPC/SDK steer()/followUp() issued during the handoff) must survive into the new
|
||||
// session instead of being silently dropped.
|
||||
vi.spyOn(compactionModule, "generateHandoffFromContext").mockResolvedValue("## Goal\nContinue");
|
||||
|
||||
const textOf = (message: AgentMessage): string => {
|
||||
if (!("content" in message)) return "";
|
||||
const content = message.content;
|
||||
if (typeof content === "string") return content;
|
||||
const textBlock = content.find(block => block.type === "text");
|
||||
return textBlock?.type === "text" ? textBlock.text : "";
|
||||
};
|
||||
|
||||
const userMsg: AgentMessage = {
|
||||
role: "user",
|
||||
content: [{ type: "text", text: "keep-steer" }],
|
||||
attribution: "user",
|
||||
timestamp: Date.now(),
|
||||
};
|
||||
// A hidden, user-attributed companion (e.g. an ultrathink notice). It is
|
||||
// display:false, so isUserQueuedMessage(...) is false for it: preservation must
|
||||
// keep it adjacent to its prompt rather than filter it out or reorder it.
|
||||
const companionMsg: AgentMessage = {
|
||||
role: "custom",
|
||||
customType: "ultrathink-notice",
|
||||
content: [{ type: "text", text: "companion" }],
|
||||
attribution: "user",
|
||||
display: false,
|
||||
timestamp: Date.now(),
|
||||
};
|
||||
const followUpMsg: AgentMessage = {
|
||||
role: "user",
|
||||
content: [{ type: "text", text: "keep-followup" }],
|
||||
attribution: "user",
|
||||
timestamp: Date.now(),
|
||||
};
|
||||
session.agent.steer(userMsg);
|
||||
session.agent.steer(companionMsg);
|
||||
session.agent.followUp(followUpMsg);
|
||||
expect(session.agent.hasQueuedMessages()).toBe(true);
|
||||
|
||||
await session.handoff();
|
||||
|
||||
expect(session.agent.peekSteeringQueue().map(textOf)).toEqual(["keep-steer", "companion"]);
|
||||
expect(session.agent.peekFollowUpQueue().map(textOf)).toEqual(["keep-followup"]);
|
||||
});
|
||||
|
||||
it("preserves steering and follow-up messages enqueued while the handoff is in flight", async () => {
|
||||
// Defect 2 in-flight window: the queue snapshot must be captured immediately before
|
||||
// agent.reset() (after generateHandoff resolves), NOT at handoff entry. A steer or
|
||||
// follow-up issued WHILE the handoff document is still generating must survive the
|
||||
// reset — proving capture happens late rather than at the start of handoff().
|
||||
const { promise: handoffDoc, resolve: releaseHandoff } = Promise.withResolvers<string>();
|
||||
let generateHandoffCalled = false;
|
||||
vi.spyOn(compactionModule, "generateHandoffFromContext").mockImplementation(async () => {
|
||||
generateHandoffCalled = true;
|
||||
return handoffDoc;
|
||||
});
|
||||
|
||||
const textOf = (message: AgentMessage): string => {
|
||||
if (!("content" in message)) return "";
|
||||
const content = message.content;
|
||||
if (typeof content === "string") return content;
|
||||
const textBlock = content.find(block => block.type === "text");
|
||||
return textBlock?.type === "text" ? textBlock.text : "";
|
||||
};
|
||||
|
||||
const handoffPromise = session.handoff();
|
||||
// Block until we are genuinely mid-handoff (document generation in flight).
|
||||
await waitFor(() => generateHandoffCalled);
|
||||
|
||||
// Enqueue AFTER generation started but BEFORE it resolves — the window where the old
|
||||
// session is still live and agent.reset() has not yet fired.
|
||||
session.agent.steer({
|
||||
role: "user",
|
||||
content: [{ type: "text", text: "inflight-steer" }],
|
||||
attribution: "user",
|
||||
timestamp: Date.now(),
|
||||
});
|
||||
session.agent.followUp({
|
||||
role: "user",
|
||||
content: [{ type: "text", text: "inflight-followup" }],
|
||||
attribution: "user",
|
||||
timestamp: Date.now(),
|
||||
});
|
||||
|
||||
releaseHandoff("## Goal\nContinue");
|
||||
await handoffPromise;
|
||||
|
||||
expect(session.agent.peekSteeringQueue().map(textOf)).toEqual(["inflight-steer"]);
|
||||
expect(session.agent.peekFollowUpQueue().map(textOf)).toEqual(["inflight-followup"]);
|
||||
});
|
||||
|
||||
it("obfuscates custom instructions before generating a handoff", async () => {
|
||||
const placeholder = obfuscator.obfuscate(HANDOFF_SECRET);
|
||||
const generateHandoffSpy = vi
|
||||
@@ -967,6 +750,7 @@ describe("AgentSession handoff", () => {
|
||||
settings: Settings.isolated({
|
||||
"compaction.enabled": true,
|
||||
"compaction.autoContinue": false,
|
||||
"compaction.asyncEnabled": false,
|
||||
"compaction.methodOrder": ["soft"],
|
||||
"compaction.thresholdTokens": 8_000,
|
||||
"contextPromotion.enabled": false,
|
||||
@@ -1302,51 +1086,25 @@ describe("AgentSession handoff", () => {
|
||||
expect(events.filter(event => event.type === "auto_compaction_end")).toHaveLength(0);
|
||||
});
|
||||
|
||||
it("persists handoff session immediately with previous session as parent", async () => {
|
||||
const previousSessionFile = session.sessionFile;
|
||||
if (!previousSessionFile) {
|
||||
throw new Error("Expected previous session file");
|
||||
}
|
||||
it("persists the handoff compaction in the current session", async () => {
|
||||
const sessionFile = session.sessionFile;
|
||||
if (!sessionFile) throw new Error("Expected current session file");
|
||||
|
||||
const handoffText = "## Goal\nContinue from here";
|
||||
vi.spyOn(compactionModule, "generateHandoffFromContext").mockResolvedValue(handoffText);
|
||||
|
||||
const result = await session.handoff();
|
||||
const handoffSessionFile = session.sessionFile;
|
||||
if (!handoffSessionFile) {
|
||||
throw new Error("Expected handoff session file");
|
||||
}
|
||||
|
||||
type PersistedEntry = {
|
||||
type?: string;
|
||||
parentSession?: string;
|
||||
customType?: string;
|
||||
display?: boolean;
|
||||
};
|
||||
const handoffEntries = (await Bun.file(handoffSessionFile).text())
|
||||
.trim()
|
||||
.split("\n")
|
||||
.map(line => JSON.parse(line) as PersistedEntry);
|
||||
const entries = sessionManager.getBranch();
|
||||
const compaction = entries.at(-1);
|
||||
|
||||
expect(result?.document).toBe(handoffText);
|
||||
expect(session.getLastAssistantText()).toBeUndefined();
|
||||
expect(session.hasCopyCandidateAssistantMessage()).toBe(false);
|
||||
expect(session.getLastVisibleHandoffText()).toBe(
|
||||
`<handoff-context>\n${handoffText}\n</handoff-context>\n\nThe above is a handoff document from a previous session. Use this context to continue the work seamlessly.`,
|
||||
);
|
||||
expect(handoffSessionFile).not.toBe(previousSessionFile);
|
||||
expect(handoffEntries.find(entry => entry.type === "session")).toMatchObject({
|
||||
type: "session",
|
||||
parentSession: previousSessionFile,
|
||||
});
|
||||
expect(
|
||||
handoffEntries.some(
|
||||
entry => entry.type === "custom_message" && entry.customType === "handoff" && entry.display,
|
||||
),
|
||||
).toBe(true);
|
||||
|
||||
const previousSessionText = await Bun.file(previousSessionFile).text();
|
||||
expect(previousSessionText).toContain('"text":"seed"');
|
||||
expect(session.sessionFile).toBe(sessionFile);
|
||||
expect(compaction).toMatchObject({ type: "compaction" });
|
||||
if (compaction?.type !== "compaction") throw new Error("Expected handoff compaction entry");
|
||||
expect(compaction.summary).toContain(handoffText);
|
||||
expect(session.agent.state.messages.some(message => message.role === "compactionSummary")).toBe(true);
|
||||
const persistedSessionText = await Bun.file(sessionFile).text();
|
||||
expect(persistedSessionText).toContain(JSON.stringify(handoffText));
|
||||
});
|
||||
|
||||
it("does not run auto maintenance when strategy is off", async () => {
|
||||
@@ -1475,26 +1233,24 @@ describe("AgentSession handoff", () => {
|
||||
timestamp: Date.now(),
|
||||
};
|
||||
|
||||
const handoffSpy = vi.spyOn(session, "handoff").mockResolvedValue({ document: "handoff document" });
|
||||
const generateHandoffSpy = vi
|
||||
.spyOn(compactionModule, "generateHandoffFromContext")
|
||||
.mockResolvedValue("handoff document");
|
||||
|
||||
session.agent.emitExternalEvent({ type: "message_end", message: assistantMessage });
|
||||
session.agent.emitExternalEvent({ type: "agent_end", messages: [assistantMessage] });
|
||||
await waitFor(
|
||||
() =>
|
||||
handoffSpy.mock.calls.length === 1 &&
|
||||
generateHandoffSpy.mock.calls.length === 1 &&
|
||||
events.filter(event => event.type === "auto_compaction_end").length === 1,
|
||||
);
|
||||
|
||||
expect(handoffSpy).toHaveBeenCalledTimes(1);
|
||||
expect(handoffSpy).toHaveBeenCalledWith(expect.stringContaining("Threshold-triggered maintenance"), {
|
||||
autoTriggered: true,
|
||||
signal: expect.anything(),
|
||||
onSwitchCancelled: expect.any(Function),
|
||||
});
|
||||
expect(generateHandoffSpy).toHaveBeenCalledTimes(1);
|
||||
expect(events.filter(event => event.type === "auto_compaction_start")).toHaveLength(1);
|
||||
const endEvents = events.filter(event => event.type === "auto_compaction_end");
|
||||
expect(endEvents).toHaveLength(1);
|
||||
expect(endEvents[0]).toMatchObject({ type: "auto_compaction_end", aborted: false, willRetry: false });
|
||||
expect(sessionManager.getBranch().at(-1)).toMatchObject({ type: "compaction", summary: "handoff document" });
|
||||
});
|
||||
|
||||
it("completes threshold-triggered auto-handoff while the original prompt is still unwinding", async () => {
|
||||
@@ -1607,7 +1363,7 @@ describe("AgentSession handoff", () => {
|
||||
expect(endEvents).toHaveLength(1);
|
||||
expect(endEvents[0]).toMatchObject({ type: "auto_compaction_end", action: "handoff", aborted: false });
|
||||
expect(endEvents[0]).not.toMatchObject({ errorMessage: expect.any(String) });
|
||||
expect(sessionManager.getEntries().filter(entry => entry.type === "compaction")).toHaveLength(0);
|
||||
expect(sessionManager.getEntries().filter(entry => entry.type === "compaction")).toHaveLength(1);
|
||||
});
|
||||
|
||||
it("does not start agent.continue when threshold-handoff defers and todos are incomplete", async () => {
|
||||
@@ -1629,9 +1385,9 @@ describe("AgentSession handoff", () => {
|
||||
throw new Error("Expected model to be set");
|
||||
}
|
||||
|
||||
const handoffSpy = vi
|
||||
.spyOn(session, "handoff")
|
||||
.mockResolvedValue({ document: "## Goal\nContinue", savedPath: undefined });
|
||||
const generateHandoffSpy = vi
|
||||
.spyOn(compactionModule, "generateHandoffFromContext")
|
||||
.mockResolvedValue("## Goal\nContinue");
|
||||
const continueSpy = vi.spyOn(session.agent, "continue");
|
||||
|
||||
const assistantMessage: AssistantMessage = {
|
||||
@@ -1654,12 +1410,11 @@ describe("AgentSession handoff", () => {
|
||||
|
||||
session.agent.emitExternalEvent({ type: "message_end", message: assistantMessage });
|
||||
session.agent.emitExternalEvent({ type: "agent_end", messages: [assistantMessage] });
|
||||
await waitFor(() => handoffSpy.mock.calls.length === 1);
|
||||
await waitFor(() => generateHandoffSpy.mock.calls.length === 1);
|
||||
await session.waitForIdle();
|
||||
|
||||
expect(handoffSpy).toHaveBeenCalledTimes(1);
|
||||
// The bug surfaced as agent.continue() racing the deferred handoff. With the fix,
|
||||
// the agent_end handler short-circuits after the deferred-handoff signal.
|
||||
expect(generateHandoffSpy).toHaveBeenCalledTimes(1);
|
||||
expect(sessionManager.getBranch().at(-1)).toMatchObject({ type: "compaction", summary: "## Goal\nContinue" });
|
||||
expect(continueSpy).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
@@ -1755,7 +1510,7 @@ describe("AgentSession handoff", () => {
|
||||
timestamp: Date.now(),
|
||||
};
|
||||
|
||||
const handoffSpy = vi.spyOn(session, "handoff").mockResolvedValue(undefined);
|
||||
const generateHandoffSpy = vi.spyOn(compactionModule, "generateHandoffFromContext").mockResolvedValue("");
|
||||
|
||||
session.agent.emitExternalEvent({ type: "message_end", message: assistantMessage });
|
||||
session.agent.emitExternalEvent({ type: "agent_end", messages: [assistantMessage] });
|
||||
@@ -1763,7 +1518,7 @@ describe("AgentSession handoff", () => {
|
||||
events.some(event => event.type === "auto_compaction_end" && event.action === "context-full"),
|
||||
);
|
||||
|
||||
expect(handoffSpy).toHaveBeenCalledTimes(1);
|
||||
expect(generateHandoffSpy).toHaveBeenCalledTimes(1);
|
||||
const endEvents = events.filter(event => event.type === "auto_compaction_end");
|
||||
expect(endEvents).toHaveLength(2);
|
||||
expect(endEvents[0]).toMatchObject({
|
||||
@@ -1781,94 +1536,6 @@ describe("AgentSession handoff", () => {
|
||||
});
|
||||
});
|
||||
|
||||
it("treats a vetoed auto-handoff switch as cancelled instead of falling back", async () => {
|
||||
session.settings.set("compaction.methodOrder", ["handoff", "soft"]);
|
||||
session.settings.set("compaction.thresholdPercent", 1);
|
||||
session.settings.set("contextPromotion.enabled", false);
|
||||
|
||||
const model = session.model;
|
||||
if (!model) {
|
||||
throw new Error("Expected model to be set");
|
||||
}
|
||||
|
||||
// See "emits handoff lifecycle hooks": reopen the persisted transcript
|
||||
// after the terminal dispose before wiring the replacement session.
|
||||
await session.dispose();
|
||||
const sessionFile = sessionManager.getSessionFile();
|
||||
if (!sessionFile) throw new Error("Expected a persisted session file");
|
||||
sessionManager = await SessionManager.open(sessionFile, tempDir.path());
|
||||
const extensionsResult = await loadExtensions([], tempDir.path());
|
||||
const extensionRunner = new ExtensionRunner(
|
||||
extensionsResult.extensions,
|
||||
extensionsResult.runtime,
|
||||
tempDir.path(),
|
||||
sessionManager,
|
||||
modelRegistry,
|
||||
);
|
||||
vi.spyOn(extensionRunner, "hasHandlers").mockImplementation(eventName => eventName === "session_before_switch");
|
||||
const emitSpy = vi.spyOn(extensionRunner, "emit").mockImplementation((async () => ({
|
||||
cancel: true,
|
||||
})) as ExtensionRunner["emit"]);
|
||||
|
||||
session = new AgentSession({
|
||||
agent: new Agent({
|
||||
initialState: {
|
||||
model,
|
||||
systemPrompt: ["Test"],
|
||||
tools: [],
|
||||
messages: [],
|
||||
},
|
||||
}),
|
||||
sessionManager,
|
||||
settings: session.settings,
|
||||
modelRegistry,
|
||||
extensionRunner,
|
||||
obfuscator,
|
||||
});
|
||||
session.subscribe(event => {
|
||||
events.push(event);
|
||||
});
|
||||
const previousSessionFile = session.sessionFile;
|
||||
const generateHandoffSpy = vi
|
||||
.spyOn(compactionModule, "generateHandoffFromContext")
|
||||
.mockResolvedValue("## Goal\nContinue from here");
|
||||
const assistantMessage: AssistantMessage = {
|
||||
role: "assistant",
|
||||
content: [{ type: "text", text: "maintenance trigger" }],
|
||||
api: model.api,
|
||||
provider: model.provider,
|
||||
model: model.id,
|
||||
stopReason: "stop",
|
||||
usage: {
|
||||
input: 10_000,
|
||||
output: 1_000,
|
||||
cacheRead: 0,
|
||||
cacheWrite: 0,
|
||||
totalTokens: 11_000,
|
||||
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
||||
},
|
||||
timestamp: Date.now(),
|
||||
};
|
||||
|
||||
session.agent.emitExternalEvent({ type: "message_end", message: assistantMessage });
|
||||
session.agent.emitExternalEvent({ type: "agent_end", messages: [assistantMessage] });
|
||||
await waitFor(() => events.filter(event => event.type === "auto_compaction_end").length === 1);
|
||||
|
||||
expect(generateHandoffSpy).toHaveBeenCalledTimes(1);
|
||||
expect(emitSpy).toHaveBeenCalledWith({ type: "session_before_switch", reason: "handoff" });
|
||||
expect(emitSpy).not.toHaveBeenCalledWith(expect.objectContaining({ type: "session_switch" }));
|
||||
expect(session.sessionFile).toBe(previousSessionFile);
|
||||
expect(sessionManager.getEntries().filter(entry => entry.type === "compaction")).toHaveLength(0);
|
||||
const endEvents = events.filter(event => event.type === "auto_compaction_end");
|
||||
expect(endEvents).toHaveLength(1);
|
||||
expect(endEvents[0]).toMatchObject({
|
||||
type: "auto_compaction_end",
|
||||
action: "handoff",
|
||||
aborted: true,
|
||||
willRetry: false,
|
||||
});
|
||||
});
|
||||
|
||||
it("resets to the base system prompt before generating a handoff", async () => {
|
||||
const model = session.model;
|
||||
if (!model) {
|
||||
@@ -1907,7 +1574,7 @@ describe("AgentSession handoff", () => {
|
||||
session = new AgentSession({
|
||||
agent,
|
||||
sessionManager,
|
||||
settings: Settings.isolated({ "compaction.enabled": false }),
|
||||
settings: Settings.isolated({ "compaction.enabled": false, "compaction.keepRecentTokens": 1 }),
|
||||
modelRegistry,
|
||||
extensionRunner,
|
||||
});
|
||||
@@ -1962,7 +1629,7 @@ describe("AgentSession handoff", () => {
|
||||
session = new AgentSession({
|
||||
agent,
|
||||
sessionManager,
|
||||
settings: Settings.isolated({ "compaction.enabled": false }),
|
||||
settings: Settings.isolated({ "compaction.enabled": false, "compaction.keepRecentTokens": 1 }),
|
||||
modelRegistry,
|
||||
});
|
||||
sessionManager.appendMessage({
|
||||
|
||||
@@ -0,0 +1,244 @@
|
||||
import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "bun:test";
|
||||
import { Agent, type AgentMessage } from "@oh-my-pi/pi-agent-core";
|
||||
import * as compactionModule from "@oh-my-pi/pi-agent-core/compaction";
|
||||
import type { AssistantMessage, Model, UserMessage } from "@oh-my-pi/pi-ai";
|
||||
import { getBundledModel } from "@oh-my-pi/pi-catalog/models";
|
||||
import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry";
|
||||
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
|
||||
import { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage";
|
||||
import { SessionMaintenance, type SessionMaintenanceHost } from "@oh-my-pi/pi-coding-agent/session/session-maintenance";
|
||||
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
|
||||
|
||||
const CONTEXT_WINDOW = 100_000;
|
||||
const THRESHOLD = 50_000;
|
||||
const SPECULATION_BAND_START = THRESHOLD - 8_192;
|
||||
|
||||
function userMessage(text: string): UserMessage {
|
||||
return { role: "user", content: [{ type: "text", text }], timestamp: Date.now() };
|
||||
}
|
||||
|
||||
function assistantMessage(text: string, model: Model): AssistantMessage {
|
||||
return {
|
||||
role: "assistant",
|
||||
content: [{ type: "text", text }],
|
||||
api: model.api,
|
||||
provider: model.provider,
|
||||
model: model.id,
|
||||
stopReason: "stop",
|
||||
usage: {
|
||||
input: 10_000,
|
||||
output: 100,
|
||||
cacheRead: 0,
|
||||
cacheWrite: 0,
|
||||
totalTokens: 10_100,
|
||||
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
||||
},
|
||||
timestamp: Date.now(),
|
||||
};
|
||||
}
|
||||
|
||||
describe("async speculative compaction", () => {
|
||||
let authStorage: AuthStorage;
|
||||
let modelRegistry: ModelRegistry;
|
||||
let model: Model;
|
||||
let sessionManager: SessionManager;
|
||||
let maintenance: SessionMaintenance;
|
||||
let events: string[];
|
||||
|
||||
function appendSummarizableConversation(): void {
|
||||
const text = "conversation ".repeat(8_000);
|
||||
sessionManager.appendMessage(userMessage(text));
|
||||
sessionManager.appendMessage(assistantMessage("response ".repeat(8_000), model));
|
||||
sessionManager.appendMessage(userMessage(text));
|
||||
sessionManager.appendMessage(assistantMessage("final response", model));
|
||||
}
|
||||
|
||||
function createMaintenance(asyncEnabled = true): SessionMaintenance {
|
||||
const agent = new Agent({
|
||||
initialState: { model, systemPrompt: ["Test"], tools: [], messages: [] },
|
||||
});
|
||||
const settings = Settings.isolated({
|
||||
"compaction.enabled": true,
|
||||
"compaction.asyncEnabled": asyncEnabled,
|
||||
"compaction.methodOrder": ["soft"],
|
||||
"compaction.thresholdPercent": 50,
|
||||
"compaction.keepRecentTokens": 1,
|
||||
"compaction.autoContinue": false,
|
||||
});
|
||||
const host = {
|
||||
agent,
|
||||
sessionManager,
|
||||
settings,
|
||||
modelRegistry,
|
||||
extensionRunner: undefined,
|
||||
sideStreamFn: async () => {
|
||||
throw new Error("The compact seam should be used instead of the side stream");
|
||||
},
|
||||
providerSessionState: new Map(),
|
||||
preferWebsockets: undefined,
|
||||
model: () => model,
|
||||
thinkingLevel: () => undefined,
|
||||
isDisposed: () => false,
|
||||
isStreaming: () => false,
|
||||
isGeneratingHandoff: () => false,
|
||||
promptGeneration: () => 0,
|
||||
sessionId: () => sessionManager.getSessionId(),
|
||||
messages: () => agent.state.messages,
|
||||
baseSystemPrompt: () => ["Test"],
|
||||
goalModeState: () => undefined,
|
||||
planReferencePath: () => "",
|
||||
nonMessageTokenSource: () => ({}),
|
||||
memoryBackendSession: () => undefined,
|
||||
emitSessionEvent: async (event: { type: string }) => {
|
||||
events.push(event.type);
|
||||
},
|
||||
emitNotice: () => {},
|
||||
schedulePostPromptTask: () => {},
|
||||
scheduleAgentContinue: () => {},
|
||||
scheduleCompactionContinuation: () => false,
|
||||
persistTurnMessagesForMidRunCompaction: async () => false,
|
||||
findLastAssistantMessage: () => undefined,
|
||||
disconnectFromAgent: () => {},
|
||||
reconnectToAgent: () => {},
|
||||
drainStrandedQueuedMessages: () => {},
|
||||
buildDisplaySessionContext: () => ({ messages: [] }),
|
||||
convertToLlmForSideRequest: (messages: AgentMessage[]) => messages as never,
|
||||
obfuscateTextForProvider: (text: string | undefined) => text,
|
||||
obfuscatePreparationForProvider: <T>(preparation: T) => preparation,
|
||||
closeCodexProviderSessionsForHistoryRewrite: () => {},
|
||||
resetCodexProviderAfterCompaction: () => {},
|
||||
resetPlanReference: () => {},
|
||||
syncTodoPhasesFromBranch: () => {},
|
||||
resetAdvisorRuntimes: () => {},
|
||||
rebaseAfterCompaction: () => {},
|
||||
recordAnchoredHistoryRewrite: () => {},
|
||||
getContextBreakdown: () => undefined,
|
||||
getContextUsage: () => undefined,
|
||||
shake: async () => ({ modified: false, tokensRemoved: 0 }),
|
||||
dropImages: async () => ({ removed: 0 }),
|
||||
generateHandoffDocument: async () => undefined,
|
||||
removeAssistantMessageFromActiveContext: () => {},
|
||||
dropPersistedAssistantTurn: async () => undefined,
|
||||
runRecoveryCompactionWithRollback: async () => ({ deferredHandoff: false, continuationScheduled: false }),
|
||||
parseRetryAfterMsFromError: () => undefined,
|
||||
setModelTemporary: async () => {},
|
||||
abort: async () => {},
|
||||
abortHandoff: () => {},
|
||||
} as unknown as SessionMaintenanceHost;
|
||||
return new SessionMaintenance(host);
|
||||
}
|
||||
|
||||
async function waitForState(state: "idle" | "running" | "armed"): Promise<void> {
|
||||
for (let microtask = 0; microtask < 100 && maintenance.speculationState !== state; microtask++) {
|
||||
await Promise.resolve();
|
||||
}
|
||||
if (maintenance.speculationState !== state) {
|
||||
throw new Error(`Speculation did not become ${state}`);
|
||||
}
|
||||
}
|
||||
|
||||
beforeAll(async () => {
|
||||
authStorage = await AuthStorage.create(":memory:");
|
||||
authStorage.setRuntimeApiKey("anthropic", "test-key");
|
||||
modelRegistry = new ModelRegistry(authStorage);
|
||||
const bundled = getBundledModel("anthropic", "claude-sonnet-4-5");
|
||||
if (!bundled) throw new Error("Expected built-in model");
|
||||
model = { ...bundled, contextWindow: CONTEXT_WINDOW };
|
||||
});
|
||||
|
||||
beforeEach(() => {
|
||||
sessionManager = SessionManager.inMemory();
|
||||
events = [];
|
||||
appendSummarizableConversation();
|
||||
maintenance = createMaintenance();
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
vi.restoreAllMocks();
|
||||
});
|
||||
|
||||
afterAll(() => {
|
||||
authStorage.close();
|
||||
});
|
||||
|
||||
it("does not call the summarizer below the speculative band, then arms inside it", async () => {
|
||||
const compactSpy = vi.spyOn(compactionModule, "compact").mockImplementation(async preparation => ({
|
||||
summary: "speculative summary",
|
||||
firstKeptEntryId: preparation.firstKeptEntryId,
|
||||
tokensBefore: preparation.tokensBefore,
|
||||
details: {},
|
||||
}));
|
||||
|
||||
maintenance.maybeStartSpeculativeCompaction(SPECULATION_BAND_START - 1, CONTEXT_WINDOW);
|
||||
expect(maintenance.speculationState).toBe("idle");
|
||||
expect(compactSpy).not.toHaveBeenCalled();
|
||||
|
||||
maintenance.maybeStartSpeculativeCompaction(SPECULATION_BAND_START, CONTEXT_WINDOW);
|
||||
expect(maintenance.speculationState).toBe("running");
|
||||
await waitForState("armed");
|
||||
expect(compactSpy).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it("commits an armed summary at threshold without paying for another summarizer call", async () => {
|
||||
const compactSpy = vi.spyOn(compactionModule, "compact").mockImplementation(async preparation => ({
|
||||
summary: "armed summary",
|
||||
firstKeptEntryId: preparation.firstKeptEntryId,
|
||||
tokensBefore: preparation.tokensBefore,
|
||||
details: {},
|
||||
}));
|
||||
maintenance.maybeStartSpeculativeCompaction(SPECULATION_BAND_START, CONTEXT_WINDOW);
|
||||
await waitForState("armed");
|
||||
|
||||
await maintenance.runAutoCompaction("threshold", false, false, false, { triggerContextTokens: THRESHOLD });
|
||||
|
||||
const entry = sessionManager.getEntries().findLast(item => item.type === "compaction");
|
||||
expect(entry?.type === "compaction" ? entry.summary : undefined).toBe("armed summary");
|
||||
expect(compactSpy).toHaveBeenCalledTimes(1);
|
||||
expect(events).toEqual(expect.arrayContaining(["auto_compaction_start", "auto_compaction_end"]));
|
||||
});
|
||||
|
||||
it("discards an armed summary after a reset boundary and re-summarizes the new branch", async () => {
|
||||
let invocation = 0;
|
||||
const compactSpy = vi.spyOn(compactionModule, "compact").mockImplementation(async preparation => ({
|
||||
summary: `summary ${++invocation}`,
|
||||
firstKeptEntryId: preparation.firstKeptEntryId,
|
||||
tokensBefore: preparation.tokensBefore,
|
||||
details: {},
|
||||
}));
|
||||
maintenance.maybeStartSpeculativeCompaction(SPECULATION_BAND_START, CONTEXT_WINDOW);
|
||||
await waitForState("armed");
|
||||
sessionManager.appendResetBoundary();
|
||||
appendSummarizableConversation();
|
||||
|
||||
await maintenance.runAutoCompaction("threshold", false, false, false, { triggerContextTokens: THRESHOLD });
|
||||
|
||||
expect(compactSpy).toHaveBeenCalledTimes(2);
|
||||
const entry = sessionManager.getEntries().findLast(item => item.type === "compaction");
|
||||
expect(entry?.type === "compaction" ? entry.summary : undefined).toBe("summary 2");
|
||||
});
|
||||
|
||||
it("does not start speculative work when async compaction is disabled", () => {
|
||||
maintenance = createMaintenance(false);
|
||||
const compactSpy = vi.spyOn(compactionModule, "compact");
|
||||
|
||||
maintenance.maybeStartSpeculativeCompaction(SPECULATION_BAND_START, CONTEXT_WINDOW);
|
||||
|
||||
expect(maintenance.speculationState).toBe("idle");
|
||||
expect(compactSpy).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("clears an armed speculation when manual compaction starts", async () => {
|
||||
vi.spyOn(compactionModule, "compact").mockImplementation(async preparation => ({
|
||||
summary: "manual summary",
|
||||
firstKeptEntryId: preparation.firstKeptEntryId,
|
||||
tokensBefore: preparation.tokensBefore,
|
||||
details: {},
|
||||
}));
|
||||
maintenance.maybeStartSpeculativeCompaction(SPECULATION_BAND_START, CONTEXT_WINDOW);
|
||||
await waitForState("armed");
|
||||
|
||||
await maintenance.compact();
|
||||
|
||||
expect(maintenance.speculationState).toBe("idle");
|
||||
});
|
||||
});
|
||||
@@ -40,6 +40,8 @@ function createContext(loopMode: SegmentContext["loopMode"]): SegmentContext {
|
||||
contextTokens: 0,
|
||||
contextWindow: 0,
|
||||
autoCompactEnabled: false,
|
||||
compactionSpeculation: "idle",
|
||||
speculationBlinkOn: true,
|
||||
subagentCount: 0,
|
||||
activeMs: 0,
|
||||
activeRepo: null,
|
||||
|
||||
@@ -47,6 +47,8 @@ function createModelContext(advisorActive: boolean): SegmentContext {
|
||||
contextTokens: 0,
|
||||
contextWindow: 0,
|
||||
autoCompactEnabled: false,
|
||||
compactionSpeculation: "idle",
|
||||
speculationBlinkOn: true,
|
||||
subagentCount: 0,
|
||||
activeMs: 0,
|
||||
activeRepo: null,
|
||||
|
||||
@@ -73,6 +73,8 @@ function createCtx(overrides?: {
|
||||
contextTokens: 0,
|
||||
contextWindow: 0,
|
||||
autoCompactEnabled: false,
|
||||
compactionSpeculation: "idle",
|
||||
speculationBlinkOn: true,
|
||||
subagentCount: 0,
|
||||
activeMs: 0,
|
||||
activeRepo: null,
|
||||
|
||||
@@ -60,6 +60,8 @@ function createPathContext(): SegmentContext {
|
||||
contextTokens: 0,
|
||||
contextWindow: 0,
|
||||
autoCompactEnabled: false,
|
||||
compactionSpeculation: "idle",
|
||||
speculationBlinkOn: true,
|
||||
subagentCount: 0,
|
||||
activeMs: 0,
|
||||
activeRepo: null,
|
||||
|
||||
@@ -62,6 +62,8 @@ function createCtx(activeMs: number): SegmentContext {
|
||||
contextTokens: 0,
|
||||
contextWindow: 0,
|
||||
autoCompactEnabled: false,
|
||||
compactionSpeculation: "idle",
|
||||
speculationBlinkOn: true,
|
||||
subagentCount: 0,
|
||||
activeMs,
|
||||
activeRepo: null,
|
||||
|
||||
Reference in New Issue
Block a user