fix(coding-agent): implement session stop hook semantics

This commit is contained in:
ben
2026-06-17 15:53:08 +08:00
committed by can1357
parent b73325b103
commit c93774f892
11 changed files with 274 additions and 86 deletions
+2 -2
View File
@@ -216,8 +216,8 @@ Cancelable pre-events:
- `before_provider_request` (may replace provider request payload)
- `after_provider_response`
- `context`
- `agent_start` / `agent_end` — agent loop lifecycle; `agent_end` is the main-agent stop-style hook and can queue hidden continuation context with `pi.sendMessage(..., { deliverAs: "nextTurn", triggerTurn: true })`
- `session_stop` — task/subagent session completion lifecycle; use this instead of `agent_end` for subagent-only cleanup or mission status tracking
- `agent_start` / `agent_end` — agent loop lifecycle notification; `agent_end` remains notification-only
- `session_stop` — main-session stop hook, awaited before settle; may continue with `{ continue: true, additionalContext }` or `{ decision: "block", reason }`; capped at 8 consecutive continuations and never fires for task/subagent sessions
- `turn_start` / `turn_end`
- `message_start` / `message_update` / `message_end`
+3 -3
View File
@@ -202,9 +202,9 @@ pi.on("turn_end", async (_event, ctx) => {
ctx.ui.setStatus("tokens", `~${ctx.getContextUsage()?.tokens ?? "?"} tokens`);
});
pi.on("session_stop", async (event, ctx) => {
// Fires for task/subagent session completion; main-agent stop-style continuation stays on agent_end.
ctx.ui.setStatus("session", `${event.messages.length} completion messages`);
pi.on("session_stop", async (event) => {
if (event.stop_hook_active) return;
return { continue: true, additionalContext: `Review final status after turn ${event.turn_id}.` };
});
```
+1 -1
View File
@@ -36,7 +36,7 @@
- Fixed hashline visible-line validation for ACP editor reads so `INS.POST` anchors displayed by bridge-backed range and multi-range `read` output are merged into the session snapshot before `edit` validates them ([#2773](https://github.com/can1357/oh-my-pi/issues/2773)).
### Added
- Added a `session_stop` extension event for task/subagent completion, leaving `agent_end` scoped to main-agent stop-style continuation ([#2834](https://github.com/can1357/oh-my-pi/issues/2834)).
- Added a main-session `session_stop` extension event with continuation feedback and an 8-continuation loop cap ([#2834](https://github.com/can1357/oh-my-pi/issues/2834)).
### Fixed
@@ -46,6 +46,8 @@ import type {
SessionBeforeSwitchResult,
SessionBeforeTreeResult,
SessionCompactingResult,
SessionStopEvent,
SessionStopEventResult,
ToolCallEvent,
ToolCallEventResult,
ToolResultEvent,
@@ -135,7 +137,9 @@ type RunnerEmitResult<TEvent extends RunnerEmitEvent> = TEvent extends { type: "
? SessionBeforeTreeResult | undefined
: TEvent extends { type: "session.compacting" }
? SessionCompactingResult | undefined
: undefined;
: TEvent extends { type: "session_stop" }
? SessionStopEventResult | undefined
: undefined;
export type NewSessionHandler = (options?: {
parentSession?: string;
@@ -322,8 +326,8 @@ export class ExtensionRunner {
await this.emit({ type: "credential_disabled", ...event });
}
async emitSessionStop(messages: AgentMessage[]): Promise<void> {
await this.emit({ type: "session_stop", messages });
async emitSessionStop(event: Omit<SessionStopEvent, "type">): Promise<SessionStopEventResult | undefined> {
return await this.emit({ type: "session_stop", ...event });
}
getUIContext(): ExtensionUIContext {
@@ -592,7 +596,7 @@ export class ExtensionRunner {
async emit<TEvent extends RunnerEmitEvent>(event: TEvent): Promise<RunnerEmitResult<TEvent>> {
const ctx = this.createContext();
let result: SessionBeforeEventResult | SessionCompactingResult | undefined;
let result: SessionBeforeEventResult | SessionCompactingResult | SessionStopEventResult | undefined;
if (this.#isSessionShutdownEvent(event)) {
const timeoutMs = handlerTimeoutForEvent(event.type);
@@ -631,6 +635,13 @@ export class ExtensionRunner {
if (event.type === "session.compacting" && handlerResult) {
result = handlerResult as SessionCompactingResult;
}
if (event.type === "session_stop" && handlerResult) {
result = handlerResult as SessionStopEventResult;
if (result.continue === true || result.decision === "block") {
return result as RunnerEmitResult<TEvent>;
}
}
}
}
@@ -83,6 +83,7 @@ import type {
SessionShutdownEvent,
SessionStartEvent,
SessionStopEvent,
SessionStopEventResult,
SessionSwitchEvent,
SessionTreeEvent,
TodoReminderEvent,
@@ -526,7 +527,14 @@ export interface BeforeAgentStartEvent {
systemPrompt: string[];
}
export type { AgentEndEvent, AgentStartEvent, SessionStopEvent, TurnEndEvent, TurnStartEvent } from "../shared-events";
export type {
AgentEndEvent,
AgentStartEvent,
SessionStopEvent,
SessionStopEventResult,
TurnEndEvent,
TurnStartEvent,
} from "../shared-events";
/** Fired when a message starts (user, assistant, or toolResult) */
export interface MessageStartEvent {
@@ -980,7 +988,7 @@ export interface ExtensionAPI {
on(event: "before_agent_start", handler: ExtensionHandler<BeforeAgentStartEvent, BeforeAgentStartEventResult>): void;
on(event: "agent_start", handler: ExtensionHandler<AgentStartEvent>): void;
on(event: "agent_end", handler: ExtensionHandler<AgentEndEvent>): void;
on(event: "session_stop", handler: ExtensionHandler<SessionStopEvent>): void;
on(event: "session_stop", handler: ExtensionHandler<SessionStopEvent, SessionStopEventResult>): void;
on(event: "turn_start", handler: ExtensionHandler<TurnStartEvent>): void;
on(event: "turn_end", handler: ExtensionHandler<TurnEndEvent>): void;
on(event: "message_start", handler: ExtensionHandler<MessageStartEvent>): void;
@@ -93,6 +93,17 @@ export interface SessionShutdownEvent {
type: "session_shutdown";
}
/** Fired when a main-agent turn is about to settle; handlers may request one continuation turn. */
export interface SessionStopEvent {
type: "session_stop";
messages: AgentMessage[];
turn_id: number;
last_assistant_message?: AgentMessage;
session_id: string;
session_file?: string;
stop_hook_active: boolean;
}
/** Preparation data for tree navigation (used by session_before_tree event) */
export interface TreePreparation {
/** Node being switched to */
@@ -145,6 +156,7 @@ export type SessionEvent =
| SessionBeforeCompactEvent
| SessionCompactingEvent
| SessionCompactEvent
| SessionStopEvent
| SessionShutdownEvent
| SessionBeforeTreeEvent
| SessionTreeEvent
@@ -181,12 +193,6 @@ export interface AgentEndEvent {
messages: AgentMessage[];
}
/** Fired when a task/subagent session loop ends */
export interface SessionStopEvent {
type: "session_stop";
messages: AgentMessage[];
}
/** Fired at the start of each turn */
export interface TurnStartEvent {
type: "turn_start";
@@ -334,6 +340,18 @@ export interface SessionCompactingResult {
preserveData?: Record<string, unknown>;
}
/** Return type for `session_stop` handlers */
export interface SessionStopEventResult {
/** Continue the main session with additional context before settling */
continue?: boolean;
/** OMP-native model-visible context for the continuation */
additionalContext?: string;
/** Claude/Codex-compatible block decision; maps to a continuation */
decision?: "block";
/** Claude/Codex-compatible model-visible continuation reason */
reason?: string;
}
/** Return type for `session_before_tree` handlers */
export interface SessionBeforeTreeResult {
/** If true, cancel the navigation entirely */
@@ -18,6 +18,7 @@ import * as os from "node:os";
import * as path from "node:path";
import { scheduler } from "node:timers/promises";
import { isPromise } from "node:util/types";
import type { InMemorySnapshotStore } from "@oh-my-pi/hashline";
import {
type AfterToolCallContext,
@@ -178,6 +179,7 @@ import type {
SessionBeforeCompactResult,
SessionBeforeSwitchResult,
SessionBeforeTreeResult,
SessionStopEventResult,
ToolExecutionEndEvent,
ToolExecutionStartEvent,
ToolExecutionUpdateEvent,
@@ -296,6 +298,8 @@ import { ToolChoiceQueue } from "./tool-choice-queue";
import { classifyUnexpectedStop, isUnexpectedStopCandidate } from "./unexpected-stop-classifier";
import { YieldQueue } from "./yield-queue";
const SESSION_STOP_CONTINUATION_CAP = 8;
/** Session-specific events that extend the core AgentEvent */
export type AgentSessionEvent =
| AgentEvent
@@ -1247,6 +1251,16 @@ export class AgentSession {
cutoffCount: number;
}
| undefined = undefined;
#sessionStopContinuationCount = 0;
#sessionStopHookActive = false;
#pendingProviderRequestNonMessageTokens: number | undefined = undefined;
#lastProviderUsageNonMessage:
| {
promptTokens: number;
nonMessageTokens: number;
cutoffCount: number;
}
| undefined = undefined;
// Bumped whenever the pending in-flight snapshot is set/cleared. The
// status-line context memo includes this so clearing the snapshot on
// turn-end/abort invalidates the cache even though the message list is
@@ -3526,6 +3540,55 @@ export class AgentSession {
}
}
#sessionStopContinuationContext(result: SessionStopEventResult | undefined): string | undefined {
if (!result) return undefined;
if (result.continue === true) {
return result.additionalContext ?? result.reason;
}
if (result.decision === "block") {
return result.reason ?? result.additionalContext;
}
return undefined;
}
async #emitSessionStopEvent(messages: AgentMessage[]): Promise<void> {
if (this.#agentKind === "sub" || !this.#extensionRunner?.hasHandlers("session_stop")) return;
const result = await this.#extensionRunner.emitSessionStop({
messages,
turn_id: this.#turnIndex,
last_assistant_message: this.getLastAssistantMessage(),
session_id: this.sessionId,
session_file: this.sessionFile,
stop_hook_active: this.#sessionStopHookActive,
});
const additionalContext = this.#sessionStopContinuationContext(result);
if (!additionalContext) {
this.#sessionStopContinuationCount = 0;
this.#sessionStopHookActive = false;
return;
}
if (this.#sessionStopContinuationCount >= SESSION_STOP_CONTINUATION_CAP) {
logger.warn("session_stop continuation cap reached", {
sessionId: this.sessionId,
cap: SESSION_STOP_CONTINUATION_CAP,
});
this.#sessionStopContinuationCount = 0;
this.#sessionStopHookActive = false;
return;
}
this.#sessionStopContinuationCount++;
this.#sessionStopHookActive = true;
await this.sendCustomMessage(
{
customType: "session-stop-continuation",
content: additionalContext,
display: false,
attribution: "agent",
},
{ deliverAs: "nextTurn", triggerTurn: true },
);
}
/** Emit extension events based on session events */
async #emitExtensionEvent(event: AgentSessionEvent): Promise<void> {
if (!this.#extensionRunner) return;
@@ -3534,6 +3597,7 @@ export class AgentSession {
await this.#extensionRunner.emit({ type: "agent_start" });
} else if (event.type === "agent_end") {
await this.#extensionRunner.emit({ type: "agent_end", messages: event.messages });
await this.#emitSessionStopEvent(event.messages);
} else if (event.type === "turn_start") {
const hookEvent: TurnStartEvent = {
type: "turn_start",
+1 -14
View File
@@ -1841,7 +1841,6 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
let error: string | undefined;
let aborted = false;
let abortReasonText: string | undefined;
const pendingExtensionMessages: Promise<unknown>[] = [];
const checkAbort = () => {
if (abortSignal.aborted) {
throw new ToolAbortError();
@@ -2118,6 +2117,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
void session.abort();
}
const pendingExtensionMessages: Array<Promise<unknown>> = [];
const extensionRunner = session.extensionRunner;
if (extensionRunner) {
extensionRunner.initialize(
@@ -2226,19 +2226,6 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
const session = monitor.takeActiveSession();
if (session) {
monitor.captureSalvage(session);
const extensionRunner = session.extensionRunner;
if (!aborted && extensionRunner?.hasHandlers("session_stop")) {
try {
await extensionRunner.emitSessionStop(session.messages);
while (pendingExtensionMessages.length > 0) {
await Promise.all(pendingExtensionMessages.splice(0));
}
} catch (err) {
logger.warn("Failed to emit session_stop event", {
error: err instanceof Error ? err.message : String(err),
});
}
}
const registry = AgentRegistry.global();
if (aborted) {
// Hard abort (caller signal / wall-clock / budget): terminal teardown.
@@ -7,7 +7,7 @@ import * as fs from "node:fs";
import * as os from "node:os";
import * as path from "node:path";
import { scheduler } from "node:timers/promises";
import { Agent, AgentBusyError, type AgentTool } from "@oh-my-pi/pi-agent-core";
import { Agent, AgentBusyError, type AgentMessage, type AgentTool } from "@oh-my-pi/pi-agent-core";
import type { AssistantMessage, Message, ToolCall } from "@oh-my-pi/pi-ai";
import { createMockModel } from "@oh-my-pi/pi-ai/providers/mock";
import { AssistantMessageEventStream } from "@oh-my-pi/pi-ai/utils/event-stream";
@@ -17,6 +17,7 @@ import type { Rule } from "@oh-my-pi/pi-coding-agent/capability/rule";
import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry";
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
import { TtsrManager } from "@oh-my-pi/pi-coding-agent/export/ttsr";
import type { ExtensionRunner } from "@oh-my-pi/pi-coding-agent/extensibility/extensions";
import { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session";
import { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage";
import { convertToLlm } from "@oh-my-pi/pi-coding-agent/session/messages";
@@ -251,6 +252,136 @@ describe("AgentSession concurrent prompt guard", () => {
).toBe(true);
});
it("continues a main session from session_stop feedback before settling", async () => {
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
const mock = createMockModel({
handler: () => ({ content: ["Done"] }),
});
const agent = new Agent({
getApiKey: () => "test-key",
initialState: { model, systemPrompt: ["Test"], tools: [] },
streamFn: mock.stream,
convertToLlm,
});
const stopEvents: Array<{
stop_hook_active: boolean;
session_id: string;
last_assistant_message?: AgentMessage;
}> = [];
const extensionRunner = {
emit: vi.fn().mockResolvedValue(undefined),
emitBeforeAgentStart: vi.fn().mockResolvedValue(undefined),
hasHandlers: vi.fn((eventType: string) => eventType === "session_stop"),
emitSessionStop: vi.fn(event => {
stopEvents.push(event);
if (stopEvents.length === 1) {
return Promise.resolve({ continue: true, additionalContext: "Mission incomplete; continue." });
}
return Promise.resolve(undefined);
}),
} as unknown as ExtensionRunner;
const sessionManager = SessionManager.inMemory();
const settings = Settings.isolated();
const authStorage = await AuthStorage.create(path.join(tempDir, "testauth.db"));
authStorages.push(authStorage);
const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models.yml"));
authStorage.setRuntimeApiKey("anthropic", "test-key");
session = new AgentSession({ agent, sessionManager, settings, modelRegistry, extensionRunner });
await session.prompt("First message");
await session.waitForIdle();
const callMessages = mock.calls.map(call => call.context.messages);
expect(callMessages).toHaveLength(2);
expect(
callMessages[1]?.some(message =>
typeof message.content === "string"
? message.content.includes("Mission incomplete; continue.")
: message.content.some(
content => content.type === "text" && content.text.includes("Mission incomplete; continue."),
),
),
).toBe(true);
expect(stopEvents.map(event => event.stop_hook_active)).toEqual([false, true]);
expect(stopEvents[0]?.session_id).toBe(session.sessionId);
expect(stopEvents[0]?.last_assistant_message?.role).toBe("assistant");
});
it("caps consecutive session_stop continuations at eight", async () => {
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
const mock = createMockModel({
handler: () => ({ content: ["Pass"] }),
});
const agent = new Agent({
getApiKey: () => "test-key",
initialState: { model, systemPrompt: ["Test"], tools: [] },
streamFn: mock.stream,
convertToLlm,
});
const extensionRunner = {
emit: vi.fn().mockResolvedValue(undefined),
emitBeforeAgentStart: vi.fn().mockResolvedValue(undefined),
hasHandlers: vi.fn((eventType: string) => eventType === "session_stop"),
emitSessionStop: vi.fn(() => Promise.resolve({ decision: "block" as const, reason: "Run another pass." })),
} as unknown as ExtensionRunner;
const sessionManager = SessionManager.inMemory();
const settings = Settings.isolated();
const authStorage = await AuthStorage.create(path.join(tempDir, "testauth.db"));
authStorages.push(authStorage);
const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models.yml"));
authStorage.setRuntimeApiKey("anthropic", "test-key");
session = new AgentSession({ agent, sessionManager, settings, modelRegistry, extensionRunner });
await session.prompt("First message");
await session.waitForIdle();
expect(mock.calls).toHaveLength(9);
expect(extensionRunner.emitSessionStop).toHaveBeenCalledTimes(9);
});
it("does not emit session_stop for subagent sessions", async () => {
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
const mock = createMockModel({
handler: () => ({ content: ["Subagent done"] }),
});
const agent = new Agent({
getApiKey: () => "test-key",
initialState: { model, systemPrompt: ["Test"], tools: [] },
streamFn: mock.stream,
convertToLlm,
});
const extensionRunner = {
emit: vi.fn().mockResolvedValue(undefined),
emitBeforeAgentStart: vi.fn().mockResolvedValue(undefined),
hasHandlers: vi.fn((eventType: string) => eventType === "session_stop"),
emitSessionStop: vi.fn().mockResolvedValue(undefined),
} as unknown as ExtensionRunner;
const sessionManager = SessionManager.inMemory();
const settings = Settings.isolated();
const authStorage = await AuthStorage.create(path.join(tempDir, "testauth.db"));
authStorages.push(authStorage);
const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models.yml"));
authStorage.setRuntimeApiKey("anthropic", "test-key");
session = new AgentSession({
agent,
sessionManager,
settings,
modelRegistry,
extensionRunner,
agentKind: "sub",
});
await session.prompt("Subagent message");
await session.waitForIdle();
expect(mock.calls).toHaveLength(1);
expect(extensionRunner.emit).toHaveBeenCalledWith({ type: "agent_end", messages: expect.any(Array) });
expect(extensionRunner.emitSessionStop).not.toHaveBeenCalled();
});
it("should allow prompt() after previous completes", async () => {
// Create session with a stream that completes immediately
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
@@ -605,7 +605,7 @@ describe("ExtensionRunner", () => {
});
describe("session_stop", () => {
it("invokes handlers with completed stopped session messages", async () => {
it("invokes handlers with completed main-session messages and returns continuation feedback", async () => {
const eventsPath = path.join(tempDir.path(), "session-stop-events.jsonl");
const extCode = `
import * as fs from "node:fs";
@@ -617,8 +617,14 @@ describe("ExtensionRunner", () => {
JSON.stringify({
type: event.type,
messages: event.messages,
turn_id: event.turn_id,
last_assistant_message: event.last_assistant_message,
session_id: event.session_id,
session_file: event.session_file,
stop_hook_active: event.stop_hook_active,
}) + "\\n",
);
return { continue: true, additionalContext: "Run one more pass." };
});
}
`;
@@ -634,7 +640,7 @@ describe("ExtensionRunner", () => {
);
const completedMessage: AgentMessage = {
role: "assistant",
content: [{ type: "text", text: "subagent finished" }],
content: [{ type: "text", text: "main session finished" }],
api: "anthropic-messages",
provider: "anthropic",
model: "claude-sonnet-4-5",
@@ -650,7 +656,14 @@ describe("ExtensionRunner", () => {
timestamp: 123,
};
await runner.emitSessionStop([completedMessage]);
const stopResult = await runner.emitSessionStop({
messages: [completedMessage],
turn_id: 2,
last_assistant_message: completedMessage,
session_id: "session-123",
session_file: "/tmp/session.jsonl",
stop_hook_active: false,
});
const events = fs
.readFileSync(eventsPath, "utf8")
@@ -661,8 +674,14 @@ describe("ExtensionRunner", () => {
{
type: "session_stop",
messages: [completedMessage],
turn_id: 2,
last_assistant_message: completedMessage,
session_id: "session-123",
session_file: "/tmp/session.jsonl",
stop_hook_active: false,
},
]);
expect(stopResult).toEqual({ continue: true, additionalContext: "Run one more pass." });
});
});
@@ -51,9 +51,6 @@ function createMockSession(
const session = {
state,
get messages() {
return state.messages;
},
agent: { state: { systemPrompt: ["test"] } },
model: undefined,
extensionRunner: undefined,
@@ -161,7 +158,6 @@ describe("runSubprocess yield reminders", () => {
}
return undefined;
},
hasHandlers: () => false,
} as unknown as NonNullable<AgentSession["extensionRunner"]>;
mockCreateAgentSession(session);
@@ -275,52 +271,6 @@ describe("runSubprocess yield reminders", () => {
expect(result.output.includes("SYSTEM WARNING")).toBe(false);
});
it("emits session_stop once after yield reminders finish", async () => {
const prompts: string[] = [];
const sessionStopPromptCounts: number[] = [];
const sessionStopMessages: Array<AgentSession["messages"]> = [];
const session = createMockSession(({ text, promptIndex, emit, state }) => {
prompts.push(text);
if (promptIndex === 1) {
const assistant = createAssistantStopMessage("did some work");
state.messages.push(assistant);
emit({ type: "message_end", message: assistant });
return;
}
emit({
type: "tool_execution_end",
toolCallId: "tool-session-stop",
toolName: "yield",
result: {
content: [{ type: "text", text: "Result submitted." }],
details: { status: "success", data: { done: true } },
},
isError: false,
});
});
const mutableSession = session as unknown as {
extensionRunner: NonNullable<AgentSession["extensionRunner"]>;
};
mutableSession.extensionRunner = {
initialize: () => {},
onError: () => {},
emit: async () => undefined,
hasHandlers: (eventType: string) => eventType === "session_stop",
emitSessionStop: async (messages: AgentSession["messages"]) => {
sessionStopPromptCounts.push(prompts.length);
sessionStopMessages.push(messages);
},
} as unknown as NonNullable<AgentSession["extensionRunner"]>;
mockCreateAgentSession(session);
const result = await runSubprocess({ ...baseOptions, id: "subagent-session-stop" });
expect(result.output).toContain('"done": true');
expect(sessionStopPromptCounts).toEqual([2]);
expect(sessionStopMessages).toHaveLength(1);
expect(sessionStopMessages[0]?.some(message => message.role === "assistant")).toBe(true);
});
it("keeps null yield warning when subagent submits success without data", async () => {
const session = createMockSession(({ promptIndex, emit, state }) => {
if (promptIndex === 1) {