feat(pi): enforce session continuity and cleanup states (#208)

This commit is contained in:
Yao
2026-09-20 13:02:17 +08:00
committed by GitHub
parent 23920a5b83
commit 3a6605594d
29 changed files with 975 additions and 88 deletions
+3
View File
@@ -57,6 +57,9 @@ pressure:
- One PR, one independently verifiable behavior. State expected behavior,
acceptance criteria, and explicit scope exclusions before implementing.
- Public topic branches use functional slugs; a GitHub Issue remains linked in
the Issue/PR metadata, and internal split identifiers never enter branch or
commit names.
- No direct commits to `main`. Every change goes through a worktree branch and
a PR.
- Never weaken sandbox path checks, API authentication, credential injection,
+6 -1
View File
@@ -19,7 +19,12 @@
log used by the builtin engine. Final text, native tool trajectory, spans,
usage, retry/error state, and terminal completion are durable; text deltas
remain transient and bounded stderr is redacted diagnostics only. Pi native
tools do not receive Harness approval or path-confinement authority.
- Adds Pi session-file lease and SQLite-backed continuity proof. Live owners
receive retryable busy responses; stale owners are recovered explicitly;
mismatched/corrupt/refused resumes are visible errors. Pi cancellation,
timeout, and unknown process-tree cleanup are distinct `cancelled`,
`timed_out`, and `cleanup_pending` states, and retained workspaces are never
released on uncertain cleanup.
### Fixes
+6 -2
View File
@@ -73,8 +73,12 @@ git worktree add .worktrees/<feature-name> -b feat/<feature-name> origin/main
- Use one branch and, when applicable, one worktree per topic. Do not mix
unrelated product work in a branch.
- Branch names use `feat/<slug>`, `fix/<slug>`, `docs/<slug>`, `test/<slug>`,
`build/<slug>`, or `chore/<slug>`. Issue-driven work uses
`fix/issue-<number>-<slug>` or `feat/issue-<number>-<slug>`.
`build/<slug>`, or `chore/<slug>`. Issue-driven work may use a public
`fix/issue-<number>-<slug>` or `feat/issue-<number>-<slug>` name, or a
functional slug-only name such as `feat/pi-session-continuity` when the
Issue number should remain in Issue/PR metadata rather than the branch.
Never put internal split identifiers such as `PR-07` in public names or
commit subjects.
- Push the branch and open a PR against `main`; do not bypass review through a
local fast-forward or merge.
- Before review, update from `origin/main`, resolve conflicts in the topic
+1 -1
View File
@@ -24,7 +24,7 @@ that every Claude hosted capability exists locally.
| Area | Endpoint group | Status | Notes |
| --- | --- | --- | --- |
| Agents | `/v1/agents` | Supported | Create, list, retrieve, update, archive, and list versions. |
| Sessions | `/v1/sessions` | Supported | Create, list, retrieve, stop, delete/archive, durable event ingestion/listing, sequence-resumable SSE, and message convenience endpoint. Session status includes `requires_action` while an approval group is pending. Event envelopes include `seq`, immutable metadata, and optional model/token/stop-reason/duration metadata. Pi stdout final replies and native tool trajectories use the same event history; transient deltas are not replayed. |
| Sessions | `/v1/sessions` | Supported | Create, list, retrieve, stop, delete/archive, durable event ingestion/listing, sequence-resumable SSE, and message convenience endpoint. Session status includes `requires_action` while an approval group is pending. Event envelopes include `seq`, immutable metadata, and optional model/token/stop-reason/duration metadata. Pi stdout final replies and native tool trajectories use the same event history; transient deltas are not replayed. Pi cleanup may surface `cancelled`, `timed_out`, or `cleanup_pending` rather than fabricated success. |
| Session artifacts | `/v1/sessions/{id}/artifacts` | Supported | Create/list artifact records and fetch content. |
| Files | `/v1/files` | Supported | Upload/list/retrieve/delete workspace files and fetch content. |
| Environments | `/v1/environments` | Supported | Create/list/retrieve/update/archive environment templates. |
+17
View File
@@ -1021,3 +1021,20 @@ Embedded test servers or custom hosts that do not provide a restart hook return
`501 unsupported`. When available, restart stops accepting requests, drains the
session manager, closes SQLite, and starts a new process with the same command
line arguments.
## Pi lifecycle status boundary
When `loop_engine.provider` is `pi`, the API may expose `cancelled`,
`timed_out`, or `cleanup_pending` rather than collapsing every child-process
outcome into `failed` or `terminated`. `cleanup_pending` is fail-closed: the
runtime has not proved that the process tree released the workspace, so it does
not clean up or accept a new turn. A live cross-runtime Pi session-file owner
returns a retryable `pi_session_busy` error. Resume refusal, corrupt headers,
path/schema mismatch, and missing SQLite continuity proof are visible errors;
they never silently start a second Pi history.
Pi-native tool events are trajectory records only. They do not carry Harness
`requires_confirmation`, do not run through the builtin `ToolResolver`, and do
not receive Harness local path confinement or an Allow/Deny card. Docker and
Kubernetes Pi transport, RPC, and a Pi-to-Harness approval bridge are not part
of this local demo.
+26 -17
View File
@@ -22,13 +22,13 @@ written to `AGENTS.md` in the sandbox work directory before launch.
The process runs with the session sandbox's host work directory as its current
working directory. Therefore the Pi foundation currently requires the **local**
sandbox provider; Docker, Kubernetes, and self-hosted sandbox pairings are
rejected by Settings rather than failing their first turn. The user prompt is
written on stdin and stdin is then closed; it is never added to command
arguments. Pi's `models.json` is written under the session's private Pi config
directory with restrictive file permissions where supported. It defines a
managed `sandbase` provider with the selected model and matching OpenAI or
Anthropic API kind, and contains only the `$SANDBASE_PI_API_KEY` credential
reference. Pi stdout is consumed as bounded LF-delimited JSONL without `readline`. The
rejected by Settings. The user prompt is written on stdin and stdin is then
closed; it is never added to command arguments. Explicit agent skills are passed
as one managed `--skill` flag per directory. Pi's `models.json` is written under
the session's private Pi config directory with restrictive file permissions and
contains only the `$SANDBASE_PI_API_KEY` credential reference.
Pi stdout is consumed as bounded LF-delimited JSONL without `readline`. The
adapter validates documented event shapes, strips structured tool markup across
chunk boundaries, and appends final text, thinking, native tool use/result,
model spans, and terminal events to the same SQLite EventLogger used by the
@@ -38,27 +38,36 @@ authority. Each Pi model request records usage exactly once. Unknown events are
inert, malformed authority-bearing events fail the turn, and stderr is limited
to a redacted 64 KiB diagnostic tail.
Pi children receive the session abort signal. An interrupt, stop, delete, or
runtime shutdown waits for the launched Pi process tree to exit before the
session work directory can be released; on POSIX the child runs in its own
process group, and on Windows the session waits for `taskkill` to finish
terminating the process tree. If tree ownership cannot be confirmed before the
cleanup deadline, the session becomes `cleanup_pending` and the workspace is
retained. Timeouts become `timed_out`; Pi cancellation becomes `cancelled`.
Pi must be installed and discoverable as `pi`; the Settings test reports a
missing CLI and a turn fails explicitly if it cannot be launched. On Windows,
the npm `pi.cmd` shim is invoked through its neighbouring `pi.ps1` script with
a fixed PowerShell argument forwarder, rather than a shell command string.
Pi children receive the session abort signal. An interrupt, stop, delete, or
runtime shutdown waits for the launched Pi process tree to exit before the
session work directory can be released; on POSIX the child runs in its own
process group, and on Windows the session waits for `taskkill` to finish
terminating the process tree.
The adapter now holds a cross-runtime lease beside the managed session file for
this entire child lifetime. A live owner returns a retryable `pi_session_busy`
error; an expired owner is recovered with an observable continuity notice. A
non-empty file must match the SQLite `pi_session_state` header id/schema/path,
otherwise resume is refused rather than silently forking. A Pi resume refusal
from stderr is persisted as `pi_resume_refused` and remains visible.
## Current scope and boundaries
The current print-mode adapter produces durable CMA events and visible Pi-native
tool trajectory. Native Pi tools are not Harness `ToolResolver` tools: they do
not receive Harness `always_ask` approval, local file path confinement, or a
fake Allow/Deny card. Foundation policy still rejects agents declaring
`always_ask`, `never_allow`, or disabled tools. Pi session-file leases,
resume/recovery continuity, turn deadlines, and stronger cleanup states remain
the next continuity work package. Docker/Kubernetes Pi transport, RPC, and a
Pi→Harness approval bridge remain excluded. It also adds no OpenAI API surface.
fake Allow/Deny card. Foundation policy still rejects agents declaring `always_ask`, `never_allow`,
or disabled tools. Continuity is guarded by the managed lease and SQLite
header state; a failed proof remains visible and cannot silently fork history.
Docker/Kubernetes Pi transport, RPC, and a Pi→Harness approval bridge remain
excluded. It also adds no OpenAI API surface.
Pi recognizes `models.json` provider settings and resolves `$ENV_VAR` values
at request time; this is why the per-session config references
+33
View File
@@ -0,0 +1,33 @@
# Pi conformance evidence
Verification date: 2026-09-15. This record separates deterministic harness
proof from real-provider proof; a pinned CLI version alone is not evidence that
a model turn used the expected trust, skill, or resume behavior.
## Local observations
- `pi --version` on the Windows host returned `0.84.4`.
- The runtime invokes Pi through stdin and the controlled launcher tests verify
multiline/Unicode input is not placed in argv.
- The launcher tests verify private `models.json` contains only
`$SANDBASE_PI_API_KEY`, source API-key/base-URL names are not inherited, and
the composed `AGENTS.md` is written in the managed work directory.
- Explicit skill directories are forwarded as one managed `--skill` argument
per directory; `model_config.speed` maps to Pi `--thinking` (`fast`→`off`,
`standard`→`medium`, `extended`→`high`).
- Three repeated controlled turns use one managed session file; a concurrent
owner gets `pi_session_busy`, and a changed header/schema is rejected.
## Real Pi provider gate
The host had no `OPENAI_API_KEY`, `ANTHROPIC_API_KEY`, `MINIMAX_API_KEY`, or
`SANDBASE_PI_API_KEY` configured at verification time. Therefore a real
non-production model turn proving Pi project trust, `AGENTS.md` loading,
explicit `--skill` behavior, provider usage, and three-turn Pi resume was not
run. This is an external credential blocker, not a passing claim. The exact
next safe command, after a temporary non-production credential is supplied, is
to run the deterministic controlled test plus a fresh workspace walkthrough
with Pi `0.84.4`, then record only redacted observations here.
No credentials, personal paths, or provider diagnostics are stored in this
repository.
+7 -1
View File
@@ -43,7 +43,7 @@ export interface ApiSession {
title: string | null;
agent: ApiAgent | { id: string; type: 'agent'; name: string };
environment_id: string;
status: 'idle' | 'running' | 'requires_action' | 'terminated' | 'failed';
status: 'idle' | 'running' | 'requires_action' | 'terminated' | 'failed' | 'cancelled' | 'timed_out' | 'cleanup_pending';
resources: ApiSessionResource[];
vault_ids: string[];
usage: {
@@ -184,6 +184,12 @@ export function toApiSessionStatus(status: string): ApiSession['status'] {
return 'terminated';
case 'failed':
return 'failed';
case 'cancelled':
return 'cancelled';
case 'timed_out':
return 'timed_out';
case 'cleanup_pending':
return 'cleanup_pending';
default:
return 'idle';
}
+14
View File
@@ -736,6 +736,19 @@ const M031_SESSION_LOOP_ENGINE = `
ALTER TABLE sessions ADD COLUMN loop_engine TEXT NOT NULL DEFAULT 'builtin';
`;
const M032_PI_SESSION_STATE = `
CREATE TABLE pi_session_state (
session_id TEXT PRIMARY KEY,
session_file TEXT NOT NULL,
pi_session_id TEXT NOT NULL,
schema_version TEXT NOT NULL DEFAULT '',
status TEXT NOT NULL DEFAULT 'active',
continuity_notice TEXT,
last_turn_at TEXT,
FOREIGN KEY (session_id) REFERENCES sessions(id)
);
`;
export const MIGRATIONS: Migration[] = [
{ version: 1, name: '001_initial', sql: M001_INITIAL },
{ version: 2, name: '002_memory', sql: M002_MEMORY },
@@ -768,4 +781,5 @@ export const MIGRATIONS: Migration[] = [
{ version: 29, name: '029_webhook_retries', sql: M029_WEBHOOK_RETRIES },
{ version: 30, name: '030_event_metadata', sql: M030_EVENT_METADATA },
{ version: 31, name: '031_session_loop_engine', sql: M031_SESSION_LOOP_ENGINE },
{ version: 32, name: '032_pi_session_state', sql: M032_PI_SESSION_STATE },
];
+12 -1
View File
@@ -1,3 +1,4 @@
import type { Database } from '@/core/db/database.js';
import { DefaultStrategy } from '@/strategy/default-strategy.js';
import { PiStrategy } from '@/strategy/pi-strategy.js';
import { PiLauncher } from '@/strategy/pi-launcher.js';
@@ -17,6 +18,7 @@ export interface RuntimeLoopEngine {
export interface RuntimeLoopEngineBootstrapOptions {
dataDir?: string;
database?: Database;
}
export function bootstrapRuntimeLoopEngine(
@@ -26,7 +28,16 @@ export function bootstrapRuntimeLoopEngine(
const builtin = new DefaultStrategy();
const strategies: Partial<Record<SessionLoopEngine, AgentStrategy>> = { builtin };
if (options.dataDir) {
strategies.pi = new PiStrategy(new PiLauncher({ dataDir: options.dataDir }));
const timeoutSeconds = settings.loop_engine.options.timeout_seconds;
const timeoutMs = typeof timeoutSeconds === 'number' && Number.isFinite(timeoutSeconds) && timeoutSeconds > 0
? Math.trunc(timeoutSeconds * 1000)
: undefined;
const launcher = new PiLauncher({
dataDir: options.dataDir,
database: options.database,
...(timeoutMs ? { timeoutMs } : {}),
});
strategies.pi = new PiStrategy(launcher, options.database);
}
const provider = settings.loop_engine.provider;
+2
View File
@@ -30,6 +30,7 @@ export interface RuntimeSessionServicesOptions {
/** Resolve the strategy matching a persisted session engine. */
resolveStrategy?: (loopEngine: SessionLoopEngine) => AgentStrategy;
skills: Skill[];
skillsDir?: string;
memory?: MemoryProvider;
artifactStore: Pick<ArtifactStore, 'path'>;
defaultMaxSteps: number;
@@ -66,6 +67,7 @@ export function createRuntimeSessionServices(options: RuntimeSessionServicesOpti
eventLogger,
compactor: new ContextCompactor(),
skills: options.skills,
skillsDir: options.skillsDir,
memory: options.memory,
snapshots,
defaultMaxSteps: options.defaultMaxSteps,
+3 -1
View File
@@ -34,7 +34,7 @@ export interface DelegationServiceDeps {
provisionSandbox: (session: Session, sandboxId: string) => Promise<SandboxInstance>;
composeSystemPrompt: (agent: AgentDefinition) => string;
buildSandboxTools: (agent: AgentDefinition, sandbox: SandboxInstance) => Record<string, any>;
}
resolveSkillDirs?: (agent: AgentDefinition) => string[];}
export class DelegationService {
constructor(private readonly deps: DelegationServiceDeps) {}
@@ -165,6 +165,7 @@ export class DelegationService {
id: subSessionId,
agentId: target.name,
agentName: target.name,
agentDefinition: target,
loopEngine: session.loopEngine ?? 'builtin',
// Inherit the parent's environment so anything downstream that
// resolves configuration from it sees the same backend the
@@ -179,6 +180,7 @@ export class DelegationService {
messages: [{ role: 'user', content: [{ type: 'text', text: task }] }] as any,
modelConfig,
model,
skillDirs: this.deps.resolveSkillDirs?.(target) ?? [],
tools,
sandbox,
eventLog: memLog,
+20 -1
View File
@@ -11,6 +11,7 @@
* terminal state (stop/delete/failed).
*/
import { dirname, resolve, sep } from 'node:path';
import type { SessionExecutor, ExecuteOptions } from './session-manager.js';
import type { Session, SessionEvent, SessionLoopEngine } from '@/types/session.js';
import type { UserEvent } from '@/types/cma-protocol.js';
@@ -55,6 +56,8 @@ export interface ExecutorDeps {
compactor?: ContextCompactor;
/** Loaded skills, injected into agent system prompts by name (R4). */
skills?: Skill[];
/** Root directory containing explicit skill packages for Pi --skill flags. */
skillsDir?: string;
/** Optional long-term memory provider, scoped by context_id (R9.16–18). */
memory?: MemoryProvider;
/** Optional workspace snapshot manager (R9.11). */
@@ -89,6 +92,7 @@ export class DefaultSessionExecutor implements SessionExecutor {
provisionSandbox: (session, sandboxId) => this.sandboxLifecycle.provisionDetached(session, sandboxId),
composeSystemPrompt: (agent) => this.contextBuilder.composeSystemPrompt(agent),
buildSandboxTools: (agent, sandbox) => this.toolResolver.buildSandboxTools(agent, sandbox),
resolveSkillDirs: (agent) => this.skillDirsFor(agent),
});
this.toolResolver = new ToolResolver({ delegationService: this.delegationService });
}
@@ -169,12 +173,13 @@ export class DefaultSessionExecutor implements SessionExecutor {
// 6. Execute strategy
const context: StrategyContext = {
session,
session: { ...session, agentDefinition: agent },
userEvent: event,
systemPrompt,
messages: messages as any,
modelConfig,
model,
...(this.deps.skillsDir ? { skillDirs: this.skillDirsFor(agent) } : {}),
tools,
sandbox,
eventLog: eventLogger,
@@ -201,6 +206,20 @@ export class DefaultSessionExecutor implements SessionExecutor {
// lifetime and are destroyed via cleanupSession() on terminal states.
}
private skillDirsFor(agent: AgentDefinition): string[] {
const root = this.deps.skillsDir;
if (!root) return [];
const resolvedRoot = resolve(root);
const byId = new Map((this.deps.skills ?? []).map((skill) => [skill.id, skill]));
return (agent.skills ?? []).flatMap((reference) => {
const skill = byId.get(reference.skill_id);
if (!skill?.file) return [];
const file = resolve(resolvedRoot, skill.file);
if (file !== resolvedRoot && !file.startsWith(`${resolvedRoot}${sep}`)) return [];
return [dirname(file)];
});
}
/**
* Destroy the sandbox + MCP connections bound to a session. Called by
* SessionManager when the session reaches a terminal state.
+3
View File
@@ -5,6 +5,9 @@ const STATUS_TO_EVENT: Partial<Record<SessionStatus, SessionEvent['type']>> = {
paused: 'session.status_idle',
requires_action: 'session.status_idle',
completed: 'session.status_terminated',
cancelled: 'session.status_terminated',
timed_out: 'session.status_terminated',
cleanup_pending: 'session.status_terminated',
// 'failed' is terminal → status_terminated. The detailed session.error event
// is appended separately by SessionManager.runTurn's catch block.
failed: 'session.status_terminated',
+32 -10
View File
@@ -426,7 +426,9 @@ export class SessionManager {
if (!isTerminal(this.get(sessionId)?.status ?? session.status)) {
this.updateStatus(sessionId, 'completed');
}
await this.releaseSandbox(sessionId);
if (this.get(sessionId)?.status !== 'cleanup_pending') {
await this.releaseSandbox(sessionId);
}
}
/**
@@ -444,7 +446,9 @@ export class SessionManager {
if (!isTerminal(this.get(sessionId)?.status ?? session.status)) {
this.updateStatus(sessionId, 'completed');
}
await this.releaseSandbox(sessionId);
if (this.get(sessionId)?.status !== 'cleanup_pending') {
await this.releaseSandbox(sessionId);
}
const deletedEvent = this.eventLogger.append(sessionId, { type: 'session.deleted' });
this.broadcast(sessionId, deletedEvent);
}
@@ -546,7 +550,9 @@ export class SessionManager {
}
}
const completedAt = isTerminal(newStatus) ? new Date().toISOString() : null;
const completedAt = new Set<SessionStatus>(['completed', 'cancelled', 'timed_out']).has(newStatus)
? new Date().toISOString()
: null;
this.db.prepare(
`UPDATE sessions SET status = ?, updated_at = datetime('now'), completed_at = ? WHERE id = ?`,
).run(newStatus, completedAt, sessionId);
@@ -601,15 +607,20 @@ export class SessionManager {
this.updateStatus(sessionId, requiresAction ? 'requires_action' : 'paused');
}
} catch (err) {
// An abort (user.interrupt) is normal control flow, not a failure:
// the session returns to idle so the user can send a follow-up.
if (abortController.signal.aborted || isAbortError(err)) {
const errorCode = errorCodeOf(err);
if (errorCode === 'pi_cleanup_pending') {
const errorEvent = this.eventLogger.append(sessionId, {
type: 'session.error',
content: [{ type: 'text', text: err instanceof Error ? err.message : String(err) }],
});
this.broadcast(sessionId, errorEvent);
if (this.get(sessionId)?.status === 'running') this.updateStatus(sessionId, 'cleanup_pending');
} else if (abortController.signal.aborted || isAbortError(err)) {
const current = this.get(sessionId);
if (current && current.status === 'running') {
this.updateStatus(sessionId, 'paused');
this.updateStatus(sessionId, current.loopEngine === 'pi' ? 'cancelled' : 'paused');
}
} else {
// Turn failed unrecoverably — terminal. Log error + release sandbox.
const errorEvent = this.eventLogger.append(sessionId, {
type: 'session.error',
content: [{ type: 'text', text: err instanceof Error ? err.message : String(err) }],
@@ -617,9 +628,14 @@ export class SessionManager {
this.broadcast(sessionId, errorEvent);
const current = this.get(sessionId);
if (current && !isTerminal(current.status)) {
this.updateStatus(sessionId, 'failed');
if (errorCode === 'pi_cleanup_pending') this.updateStatus(sessionId, 'cleanup_pending');
else if (errorCode === 'pi_timed_out') this.updateStatus(sessionId, 'timed_out');
else if (errorCode === 'pi_session_busy') this.updateStatus(sessionId, 'paused');
else this.updateStatus(sessionId, 'failed');
}
await this.releaseSandbox(sessionId);
// A cleanup_pending child may still own the workspace. Never release it
// based on parent close or a failed taskkill result.
if (errorCode !== 'pi_cleanup_pending') await this.releaseSandbox(sessionId);
}
} finally {
this.abortControllers.delete(sessionId);
@@ -735,3 +751,9 @@ function lifecycleMetadataFor(status: SessionStatus, events: SessionEvent[]): Re
},
};
}
function errorCodeOf(error: unknown): string | undefined {
if (!error || typeof error !== 'object' || !('code' in error)) return undefined;
const code = (error as { code?: unknown }).code;
return typeof code === 'string' ? code : undefined;
}
+3 -1
View File
@@ -9,7 +9,9 @@
* requires_action → running
* failed → running (resume) | completed (stopped/deleted)
*
* Terminal state: completed only (no transitions out). `failed` is recoverable.
* Terminal states are completed, cancelled, timed_out, and cleanup_pending.
* cleanup_pending deliberately has no outbound transition: the workspace
* remains retained until an operator can prove child-tree cleanup.
*/
import { type SessionStatus, SESSION_TRANSITIONS } from '@/types/session.js';
+1
View File
@@ -57,6 +57,7 @@ export function describeSettingsAdapters(installedSandboxes: string[] = ['local'
// the external executable and a turn still fails explicitly if it is absent.
descriptor('pi', 'Pi CLI', true, 'runtime', objectSchema({
default_max_steps: { type: 'integer', minimum: 1, maximum: 1000, default: 25 },
timeout_seconds: { type: 'integer', minimum: 1, maximum: 86400, default: 300 },
})),
descriptor('harness', 'Harness', false, 'runtime', objectSchema()),
descriptor('codex', 'Codex', false, 'runtime', objectSchema()),
+2 -1
View File
@@ -94,7 +94,7 @@ async function startServer(opts: StartServerOptions) {
const effectiveSettings = runtimeComposition.settings.effective_config;
const memory = runtimeComposition.memory;
const loopEngine = bootstrapRuntimeLoopEngine(effectiveSettings, { dataDir });
const loopEngine = bootstrapRuntimeLoopEngine(effectiveSettings, { dataDir, database: db });
const artifactStore = runtimeComposition.artifactStore;
const {
@@ -116,6 +116,7 @@ async function startServer(opts: StartServerOptions) {
return strategy;
},
skills,
skillsDir,
memory,
artifactStore,
defaultMaxSteps: loopEngine.defaultMaxSteps,
+169 -36
View File
@@ -4,7 +4,14 @@ import { chmodSync, existsSync, lstatSync, mkdirSync, renameSync, rmSync, writeF
import { randomUUID } from 'node:crypto';
import { delimiter, dirname, extname, join, resolve, sep, win32 } from 'node:path';
import { referencedEnvVars, resolveEnvVarsFrom } from '@/core/config/env-resolver.js';
import type { Database } from '@/core/db/database.js';
import type { ModelConfig } from '@/types/model.js';
import { acquirePiSessionFileLease, type PiSessionFileLease } from './pi/session-lease.js';
import {
assertPiSessionContinuity,
markPiSessionContinuityFailure,
type PiContinuityError,
} from './pi/session-continuity.js';
const INHERITED_ENVIRONMENT_KEYS = [
'HOME', 'LANG', 'LC_ALL', 'LOGNAME', 'PATH', 'SHELL', 'TERM', 'TMPDIR',
@@ -23,10 +30,31 @@ export interface PiLaunchRequest {
systemPrompt: string;
/** Concrete configuration selected for this agent turn. */
model: PiModelConfig;
/** Explicit Pi skill directories, one `--skill` flag per directory. */
skillDirs?: string[];
thinkingLevel?: 'off' | 'minimal' | 'low' | 'medium' | 'high' | 'xhigh' | 'max';
/** Cancels the in-flight child only after it has exited and released its workdir. */
abortSignal?: AbortSignal;
}
export class PiCleanupPendingError extends Error {
readonly code = 'pi_cleanup_pending';
constructor(message = 'Pi process tree cleanup is pending; the workspace remains retained') {
super(message);
this.name = 'PiCleanupPendingError';
}
}
export class PiTimeoutError extends Error {
readonly code = 'pi_timed_out';
constructor(readonly timeoutMs: number) {
super(`Pi turn timed out after ${timeoutMs}ms`);
this.name = 'PiTimeoutError';
}
}
export interface PiProcessExit {
code: number | null;
signal: NodeJS.Signals | null;
@@ -37,6 +65,8 @@ export interface PiProcessHandle {
readonly child: ChildProcess;
readonly stdout: Readable | null;
readonly stderr: Readable | null;
readonly sessionFile?: string;
readonly leaseRecovered?: boolean;
wait(): Promise<PiProcessExit>;
terminate(force?: boolean): Promise<void>;
}
@@ -59,6 +89,14 @@ export interface PiLauncherOptions {
/** Host Settings environment used for resolving model placeholders. */
environment?: NodeJS.ProcessEnv;
fileExists?: (path: string) => boolean;
/** Host database used to validate the Pi session header against continuity state. */
database?: Database;
/** Stale lease expiry; the default is deliberately short and observable. */
leaseStaleAfterMs?: number;
/** Maximum time to wait for a child tree after abort before cleanup_pending. */
cleanupTimeoutMs?: number;
/** Per-turn timeout; defaults to five minutes. */
timeoutMs?: number;
/** Test seam for terminating an in-flight Pi child process tree. */
terminateProcess?: PiProcessTerminator;
/** Bounded grace period before a POSIX process-group kill is escalated. */
@@ -137,6 +175,10 @@ export class PiLauncher {
private readonly terminateProcess: PiProcessTerminator;
private readonly terminationGraceMs: number;
private readonly processGroupAlive: PiProcessGroupInspector;
private readonly database?: Database;
private readonly leaseStaleAfterMs: number;
private readonly cleanupTimeoutMs: number;
private readonly timeoutMs: number;
constructor(private readonly options: PiLauncherOptions) {
this.command = options.command ?? 'pi';
@@ -148,6 +190,10 @@ export class PiLauncher {
this.terminateProcess = options.terminateProcess ?? terminatePiProcess;
this.terminationGraceMs = options.terminationGraceMs ?? 1_000;
this.processGroupAlive = options.processGroupAlive ?? isProcessGroupAlive;
this.database = options.database;
this.leaseStaleAfterMs = options.leaseStaleAfterMs ?? 30_000;
this.cleanupTimeoutMs = options.cleanupTimeoutMs ?? 10_000;
this.timeoutMs = options.timeoutMs ?? 300_000;
}
/** Compatibility wrapper used by the foundation: drain output and wait. */
@@ -163,46 +209,105 @@ export class PiLauncher {
/** Start a Pi child for a protocol adapter to consume incrementally. */
async start(request: PiLaunchRequest): Promise<PiProcessHandle> {
// Resolve against the host Settings environment before materializing either
// the per-session config or Pi's credential alias. The original host names
// are then explicitly removed from the restricted child environment.
const model = this.resolveModel(request.model);
const modelEnvironmentKeys = new Set([
...referencedEnvVars(request.model.api_key),
...referencedEnvVars(request.model.base_url),
]);
const paths = this.prepareSessionPaths(request.sessionId);
this.materializeModelsConfig(paths.configDir, model);
this.materializeAgentsPrompt(request.workDir, request.systemPrompt);
const invocation = piInvocationFor([
'-p', '--mode', 'json', '--model', `sandbase/${model.model}`, '--session', paths.sessionFile,
], {
command: this.command,
commandArgs: this.commandArgs,
platform: this.platform,
environment: this.environment,
fileExists: this.fileExists,
const lease = await acquirePiSessionFileLease(paths.sessionFile, {
staleAfterMs: this.leaseStaleAfterMs,
});
const env = restrictedPiEnvironment(this.environment, {
PI_CODING_AGENT_DIR: paths.configDir,
PI_TELEMETRY: '0',
SANDBASE_PI_API_KEY: model.api_key,
}, modelEnvironmentKeys);
return spawnPiProcess(this.spawnImpl, invocation.file, invocation.args, {
cwd: request.workDir,
env,
stdio: ['pipe', 'pipe', 'pipe'],
windowsHide: true,
// Detached POSIX children let an interrupt address the entire Pi process
// group, including any CLI descendants that retain the workspace.
detached: this.platform !== 'win32',
}, request.prompt, request.abortSignal, {
platform: this.platform,
terminateProcess: this.terminateProcess,
terminationGraceMs: this.terminationGraceMs,
processGroupAlive: this.processGroupAlive,
});
try {
if (this.database) assertPiSessionContinuity(this.database, request.sessionId, paths.sessionFile);
this.materializeModelsConfig(paths.configDir, model);
this.materializeAgentsPrompt(request.workDir, request.systemPrompt);
const invocation = piInvocationFor([
'-p', '--mode', 'json', '--model', `sandbase/${model.model}`, '--session', paths.sessionFile,
...(request.thinkingLevel ? ['--thinking', request.thinkingLevel] : []),
...(request.skillDirs ?? []).flatMap((directory) => ['--skill', directory]),
], {
command: this.command,
commandArgs: this.commandArgs,
platform: this.platform,
environment: this.environment,
fileExists: this.fileExists,
});
const env = restrictedPiEnvironment(this.environment, {
PI_CODING_AGENT_DIR: paths.configDir,
PI_TELEMETRY: '0',
SANDBASE_PI_API_KEY: model.api_key,
}, modelEnvironmentKeys);
const abortController = new AbortController();
let timedOut = false;
const onRequestAbort = () => abortController.abort();
if (request.abortSignal) {
if (request.abortSignal.aborted) abortController.abort();
else request.abortSignal.addEventListener('abort', onRequestAbort, { once: true });
}
const timeout = setTimeout(() => {
timedOut = true;
abortController.abort();
}, this.timeoutMs);
let raw: PiProcessHandle;
try {
raw = await spawnPiProcess(this.spawnImpl, invocation.file, invocation.args, {
cwd: request.workDir,
env,
stdio: ['pipe', 'pipe', 'pipe'],
windowsHide: true,
detached: this.platform !== 'win32',
}, request.prompt, abortController.signal, {
platform: this.platform,
terminateProcess: this.terminateProcess,
terminationGraceMs: this.terminationGraceMs,
cleanupTimeoutMs: this.cleanupTimeoutMs,
processGroupAlive: this.processGroupAlive,
});
} catch (error) {
clearTimeout(timeout);
request.abortSignal?.removeEventListener('abort', onRequestAbort);
throw error;
}
const wait = async (): Promise<PiProcessExit> => {
let failure: unknown;
try {
const exit = await raw.wait();
if (timedOut) throw new PiTimeoutError(this.timeoutMs);
return exit;
} catch (error) {
failure = error;
if (timedOut && !(error instanceof PiCleanupPendingError)) {
throw new PiTimeoutError(this.timeoutMs);
}
throw error;
} finally {
clearTimeout(timeout);
request.abortSignal?.removeEventListener('abort', onRequestAbort);
if (isCleanupPendingError(failure)) lease.suspendHeartbeat();
else await lease.release().catch(() => {});
}
};
return {
...raw,
sessionFile: paths.sessionFile,
leaseRecovered: lease.recoveredStale,
wait,
};
} catch (error) {
if (this.database && isPiContinuityError(error)) {
// Persist the reason so a later retry cannot silently fork a new file.
const message = error instanceof Error ? error.message : String(error);
markPiSessionContinuityFailure(this.database, request.sessionId, paths.sessionFile, error.code, message);
}
await lease.release().catch(() => {});
throw error;
}
}
private prepareSessionPaths(sessionId: string): { sessionFile: string; configDir: string } {
@@ -312,9 +417,21 @@ type PiProcessAbortOptions = {
platform: NodeJS.Platform;
terminateProcess: PiProcessTerminator;
terminationGraceMs: number;
cleanupTimeoutMs: number;
processGroupAlive: PiProcessGroupInspector;
};
function isCleanupPendingError(error: unknown): error is PiCleanupPendingError {
return error instanceof PiCleanupPendingError || (
error instanceof Error && (error as Error & { code?: unknown }).code === 'pi_cleanup_pending'
);
}
function isPiContinuityError(error: unknown): error is PiContinuityError {
return error instanceof Error && error.name === 'PiContinuityError'
&& typeof (error as Error & { code?: unknown }).code === 'string';
}
async function spawnPiProcess(
spawnImpl: typeof spawn,
file: string,
@@ -326,6 +443,7 @@ async function spawnPiProcess(
platform: process.platform,
terminateProcess: terminatePiProcess,
terminationGraceMs: 1_000,
cleanupTimeoutMs: 10_000,
processGroupAlive: isProcessGroupAlive,
},
): Promise<PiProcessHandle> {
@@ -344,6 +462,7 @@ async function spawnPiProcess(
let childClosed = false;
let windowsTreeTerminationComplete = false;
let forceTimer: ReturnType<typeof setTimeout> | undefined;
let cleanupTimer: ReturnType<typeof setTimeout> | undefined;
let settled = false;
let resolveWait: (exit: PiProcessExit) => void = () => {};
let rejectWait: (error: unknown) => void = () => {};
@@ -388,6 +507,11 @@ async function spawnPiProcess(
const onAbort = () => {
if (abortRequested) return;
abortRequested = true;
cleanupTimer = setTimeout(() => {
rejectWait(new PiCleanupPendingError(
`Pi process tree cleanup did not complete within ${abortOptions.cleanupTimeoutMs}ms; workspace retained`,
));
}, abortOptions.cleanupTimeoutMs);
const termination = requestTermination(false);
if (abortOptions.platform === 'win32') {
void termination.then(
@@ -408,15 +532,18 @@ async function spawnPiProcess(
rejectAbortIfSafe();
return;
}
void waitForProcessGroupExit(child.pid, abortOptions.processGroupAlive).then(() => {
rejectAbortIfSafe();
});
void waitForProcessGroupExit(child.pid, abortOptions.processGroupAlive, abortOptions.cleanupTimeoutMs)
.then(() => {
rejectAbortIfSafe();
})
.catch(() => {});
}, abortOptions.terminationGraceMs);
}
rejectAbortIfSafe();
};
const cleanup = () => {
if (forceTimer) clearTimeout(forceTimer);
if (cleanupTimer) clearTimeout(cleanupTimer);
abortSignal?.removeEventListener('abort', onAbort);
};
@@ -467,7 +594,7 @@ async function spawnPiProcess(
terminate: async (force = false) => {
await requestTermination(force);
if (force && abortOptions.platform !== 'win32' && child.pid !== undefined) {
await waitForProcessGroupExit(child.pid, abortOptions.processGroupAlive);
await waitForProcessGroupExit(child.pid, abortOptions.processGroupAlive, abortOptions.cleanupTimeoutMs);
}
},
};
@@ -548,8 +675,13 @@ function isProcessGroupAlive(pid: number): boolean {
async function waitForProcessGroupExit(
pid: number,
processGroupAlive: PiProcessGroupInspector,
timeoutMs = Number.POSITIVE_INFINITY,
): Promise<void> {
const deadline = Date.now() + timeoutMs;
while (processGroupAlive(pid)) {
if (Date.now() >= deadline) {
throw new PiCleanupPendingError(`Pi process group did not exit within ${timeoutMs}ms`);
}
await new Promise((resolvePromise) => setTimeout(resolvePromise, 10));
}
}
@@ -561,6 +693,7 @@ function abortError(): Error {
}
function piLaunchError(error: unknown): Error {
if (error instanceof PiCleanupPendingError || error instanceof PiTimeoutError) return error;
if (error instanceof Error && error.name === 'AbortError') return error;
const code = typeof error === 'object' && error !== null && 'code' in error
? (error as { code?: unknown }).code
+54 -1
View File
@@ -1,6 +1,16 @@
import type { TextBlock } from '@/types/cma-protocol.js';
import type { AgentStrategy, StrategyContext } from '@/types/strategy.js';
import type { PiLaunchRequest, PiProcessHandle } from './pi-launcher.js';
import type { Database } from '@/core/db/database.js';
import {
getPiSessionState,
inspectPiSessionFile,
markPiSessionContinuityFailure,
recordPiSessionState,
PiContinuityError,
PI_RESUME_REFUSED_MARKER,
} from './pi/session-continuity.js';
import { PiCleanupPendingError, PiTimeoutError } from './pi-launcher.js';
import { PiStderrTail } from './pi/stderr-tail.js';
import { PiTranslator } from './pi/translator.js';
@@ -19,8 +29,11 @@ export class PiStrategy implements AgentStrategy {
readonly name = 'pi';
/** Pi owns model transport, so the executor must not construct an AI SDK model. */
readonly requiresModel = false;
private readonly database?: Database;
constructor(private readonly launcher: PiTurnLauncher) {}
constructor(private readonly launcher: PiTurnLauncher, database?: Database) {
this.database = database;
}
async *execute(context: StrategyContext) {
if (!context.sandbox.hostWorkDir) {
@@ -46,6 +59,10 @@ export class PiStrategy implements AgentStrategy {
prompt: context.userEvent.content.map((block) => block.text).join('\n'),
systemPrompt: context.systemPrompt,
model: context.modelConfig,
...(context.skillDirs?.length ? { skillDirs: context.skillDirs } : {}),
...(thinkingLevelForSpeed(context.session.agentDefinition?.model_config?.speed)
? { thinkingLevel: thinkingLevelForSpeed(context.session.agentDefinition?.model_config?.speed) }
: {}),
...(context.abortSignal ? { abortSignal: context.abortSignal } : {}),
};
@@ -93,6 +110,25 @@ export class PiStrategy implements AgentStrategy {
throw new Error(withStderr(summary.lastTurnError, stderr.text()));
}
if (this.database && handle.sessionFile) {
const header = inspectPiSessionFile(handle.sessionFile);
if (!header) {
throw new PiContinuityError('pi_session_discontinuous', 'Pi completed without writing a session header');
}
const previous = getPiSessionState(this.database, context.session.id);
if (previous && (previous.piSessionId !== header.id || previous.schemaVersion !== header.schemaVersion)) {
throw new PiContinuityError('pi_session_discontinuous', 'Pi changed its session identity or schema during the turn');
}
recordPiSessionState(this.database, context.session.id, handle.sessionFile, header);
if (handle.leaseRecovered) {
const notice = context.eventLog.append(context.session.id, {
type: 'agent.message',
content: [{ type: 'text', text: 'Pi continuity notice: recovered a stale session-file lease before this turn.' }],
});
context.broadcast(notice);
}
}
// `turn_complete` is the durable adapter terminal marker. It is appended
// before strategy completion and never yielded separately.
const terminal = context.eventLog.append(context.session.id, { type: 'turn_complete' });
@@ -104,12 +140,29 @@ export class PiStrategy implements AgentStrategy {
}
await stderrDrain.catch(() => {});
if (error instanceof Error && error.name === 'AbortError') throw error;
if (error instanceof PiCleanupPendingError || error instanceof PiTimeoutError || error instanceof PiContinuityError) {
throw error;
}
const message = error instanceof Error ? error.message : String(error);
if (this.database && handle.sessionFile && stderr.text().includes(PI_RESUME_REFUSED_MARKER)) {
const refusal = new PiContinuityError('pi_resume_refused', `Pi refused to resume its managed session: ${stderr.text()}`);
markPiSessionContinuityFailure(this.database, context.session.id, handle.sessionFile, refusal.code, refusal.message);
throw refusal;
}
throw new Error(withStderr(message, stderr.text()));
}
}
}
function thinkingLevelForSpeed(speed: string | undefined): PiLaunchRequest['thinkingLevel'] {
switch (speed) {
case 'fast': return 'off';
case 'extended': return 'high';
case 'standard': return 'medium';
default: return undefined;
}
}
async function drainStderr(stream: AsyncIterable<Uint8Array | string>, tail: PiStderrTail): Promise<void> {
for await (const chunk of stream) tail.append(chunk);
}
+141
View File
@@ -0,0 +1,141 @@
import { existsSync, readFileSync } from 'node:fs';
import type { Database } from '@/core/db/database.js';
export const PI_RESUME_REFUSED_MARKER = 'Stored session working directory does not exist';
export interface PiSessionHeader {
id: string;
schemaVersion: string;
}
export interface PiSessionState {
sessionId: string;
sessionFile: string;
piSessionId: string;
schemaVersion: string;
status: string;
continuityNotice?: string;
}
export class PiContinuityError extends Error {
readonly code: string;
constructor(code: string, message: string) {
super(message);
this.name = 'PiContinuityError';
this.code = code;
}
}
/** Read and validate only the Pi session header; the file is not event authority. */
export function inspectPiSessionFile(sessionFile: string): PiSessionHeader | undefined {
if (!existsSync(sessionFile)) return undefined;
const source = readFileSync(sessionFile, 'utf8');
const firstLine = source.split(/\r?\n/, 1)[0]?.trim();
if (!firstLine) return undefined;
let parsed: unknown;
try {
parsed = JSON.parse(firstLine);
} catch {
throw new PiContinuityError('pi_session_file_corrupt', `Pi session file has an invalid JSON header: ${sessionFile}`);
}
if (!parsed || typeof parsed !== 'object' || Array.isArray(parsed)) {
throw new PiContinuityError('pi_session_file_schema_invalid', 'Pi session header must be a JSON object');
}
const record = parsed as Record<string, unknown>;
if (record.type !== 'session' || typeof record.id !== 'string' || !record.id) {
throw new PiContinuityError('pi_session_file_schema_invalid', 'Pi session header is missing type=session or id');
}
const schema = record.version ?? record.schemaVersion ?? record.schema_version ?? 'unknown';
return { id: record.id, schemaVersion: String(schema) };
}
/** Refuse a non-empty file unless SQLite can prove it is the same Pi session. */
export function assertPiSessionContinuity(
db: Database,
sessionId: string,
sessionFile: string,
): { header?: PiSessionHeader; state?: PiSessionState } {
const header = inspectPiSessionFile(sessionFile);
const state = getPiSessionState(db, sessionId);
if (!header && !state) return {};
if (!header && state) {
throw new PiContinuityError('pi_session_discontinuous', 'Pi session state exists but its managed session file is empty');
}
if (header && !state) {
throw new PiContinuityError('pi_session_discontinuous', 'Pi session file exists without matching SandBase continuity state');
}
if (!header || !state) return {};
if (state.status !== 'active') {
throw new PiContinuityError('pi_session_discontinuous', `Pi continuity state is ${state.status}; resume is refused until it is repaired`);
}
if (state.sessionFile !== sessionFile || state.piSessionId !== header.id || state.schemaVersion !== header.schemaVersion) {
throw new PiContinuityError('pi_session_discontinuous', 'Pi session file identity or schema does not match SandBase continuity state');
}
return { header, state };
}
export function getPiSessionState(db: Database, sessionId: string): PiSessionState | undefined {
const row = db.prepare(`
SELECT session_id, session_file, pi_session_id, schema_version, status, continuity_notice
FROM pi_session_state WHERE session_id = ?
`).get(sessionId) as {
session_id: string;
session_file: string;
pi_session_id: string;
schema_version: string;
status: string;
continuity_notice: string | null;
} | undefined;
if (!row) return undefined;
return {
sessionId: row.session_id,
sessionFile: row.session_file,
piSessionId: row.pi_session_id,
schemaVersion: row.schema_version,
status: row.status,
...(row.continuity_notice ? { continuityNotice: row.continuity_notice } : {}),
};
}
export function recordPiSessionState(
db: Database,
sessionId: string,
sessionFile: string,
header: PiSessionHeader,
status = 'active',
continuityNotice?: string,
): void {
db.prepare(`
INSERT INTO pi_session_state (
session_id, session_file, pi_session_id, schema_version, status, continuity_notice, last_turn_at
) VALUES (?, ?, ?, ?, ?, ?, datetime('now'))
ON CONFLICT(session_id) DO UPDATE SET
session_file = excluded.session_file,
pi_session_id = excluded.pi_session_id,
schema_version = excluded.schema_version,
status = excluded.status,
continuity_notice = excluded.continuity_notice,
last_turn_at = excluded.last_turn_at
`).run(sessionId, sessionFile, header.id, header.schemaVersion, status, continuityNotice ?? null);
}
export function markPiSessionContinuityFailure(
db: Database,
sessionId: string,
sessionFile: string,
code: string,
message: string,
): void {
const existing = getPiSessionState(db, sessionId);
if (existing) {
db.prepare(`UPDATE pi_session_state SET status = ?, continuity_notice = ? WHERE session_id = ?`)
.run(code, message, sessionId);
return;
}
db.prepare(`
INSERT OR IGNORE INTO pi_session_state (
session_id, session_file, pi_session_id, schema_version, status, continuity_notice
) VALUES (?, ?, '', '', ?, ?)
`).run(sessionId, sessionFile, code, message);
}
+174
View File
@@ -0,0 +1,174 @@
import { randomUUID } from 'node:crypto';
import { hostname } from 'node:os';
import { open, readFile, rename, rm, stat, writeFile } from 'node:fs/promises';
import { join, dirname } from 'node:path';
const DEFAULT_STALE_AFTER_MS = 30_000;
export class PiSessionBusyError extends Error {
readonly code = 'pi_session_busy';
readonly retryable = true;
constructor(readonly sessionFile: string, readonly owner?: PiLeaseRecord) {
super(`Pi session file is busy: ${sessionFile}`);
this.name = 'PiSessionBusyError';
}
}
export interface PiLeaseRecord {
version: 1;
ownerId: string;
pid: number;
host: string;
acquiredAt: string;
heartbeatAt: string;
expiresAt: string;
}
export interface PiSessionFileLease {
readonly leasePath: string;
readonly ownerId: string;
readonly recoveredStale: boolean;
renew(): Promise<void>;
/** Stop heartbeats without deleting the lease when cleanup ownership is unknown. */
suspendHeartbeat(): void;
release(): Promise<void>;
}
export interface PiSessionLeaseOptions {
staleAfterMs?: number;
now?: () => number;
ownerId?: string;
pid?: number;
host?: string;
}
/**
* Acquire a cross-runtime lease by exclusive creation and expiry heartbeat.
* A stale owner is recovered by an atomic rename before the next attempt.
*/
export async function acquirePiSessionFileLease(
sessionFile: string,
options: PiSessionLeaseOptions = {},
): Promise<PiSessionFileLease> {
const staleAfterMs = options.staleAfterMs ?? DEFAULT_STALE_AFTER_MS;
if (!Number.isSafeInteger(staleAfterMs) || staleAfterMs <= 0) {
throw new RangeError('Pi session lease expiry must be a positive safe integer');
}
const now = options.now ?? Date.now;
const ownerId = options.ownerId ?? randomUUID();
const pid = options.pid ?? process.pid;
const host = options.host ?? hostname();
const leasePath = `${sessionFile}.lease`;
let recoveredStale = false;
await import('node:fs/promises').then(({ mkdir }) => mkdir(dirname(leasePath), { recursive: true }));
for (let attempt = 0; attempt < 3; attempt += 1) {
const acquiredAt = now();
const record = makeRecord(ownerId, pid, host, acquiredAt, staleAfterMs);
try {
const handle = await open(leasePath, 'wx', 0o600);
try {
await handle.writeFile(JSON.stringify(record));
} finally {
await handle.close();
}
let released = false;
let renewal: Promise<void> = Promise.resolve();
const heartbeat = setInterval(() => {
renewal = renewal.then(async () => {
if (released) return;
const next = makeRecord(ownerId, pid, host, now(), staleAfterMs);
await writeFile(leasePath, JSON.stringify(next), { encoding: 'utf8', mode: 0o600 });
}).catch(() => {
// The owner remains conservative: a failed heartbeat makes the lease
// appear stale only after the expiry window, never immediately safe.
});
}, Math.max(1_000, Math.floor(staleAfterMs / 3)));
return {
leasePath,
ownerId,
recoveredStale,
renew: async () => {
await renewal;
if (released) return;
const next = makeRecord(ownerId, pid, host, now(), staleAfterMs);
await writeFile(leasePath, JSON.stringify(next), { encoding: 'utf8', mode: 0o600 });
},
suspendHeartbeat: () => {
clearInterval(heartbeat);
},
release: async () => {
if (released) return;
released = true;
clearInterval(heartbeat);
await renewal;
try {
const current = JSON.parse(await readFile(leasePath, 'utf8')) as Partial<PiLeaseRecord>;
if (current.ownerId === ownerId) await rm(leasePath, { force: true });
} catch {
// Never remove an unreadable lease: ownership cannot be proven.
}
},
};
} catch (error) {
if (!isAlreadyExists(error)) throw error;
const existing = await readExistingLease(leasePath);
if (!isStale(existing, now(), staleAfterMs)) {
throw new PiSessionBusyError(sessionFile, existing?.record);
}
const recoveryPath = join(dirname(leasePath), `.${ownerId}.stale`);
try {
await rename(leasePath, recoveryPath);
await rm(recoveryPath, { force: true });
recoveredStale = true;
} catch {
// A concurrent owner won the race. Re-check on the next attempt.
}
}
}
throw new PiSessionBusyError(sessionFile);
}
function makeRecord(ownerId: string, pid: number, host: string, now: number, staleAfterMs: number): PiLeaseRecord {
const timestamp = new Date(now).toISOString();
return {
version: 1,
ownerId,
pid,
host,
acquiredAt: timestamp,
heartbeatAt: timestamp,
expiresAt: new Date(now + staleAfterMs).toISOString(),
};
}
async function readExistingLease(path: string): Promise<{ record?: PiLeaseRecord; mtimeMs?: number } | undefined> {
try {
const [contents, metadata] = await Promise.all([readFile(path, 'utf8'), stat(path)]);
try {
const record = JSON.parse(contents) as PiLeaseRecord;
if (record && record.version === 1 && typeof record.expiresAt === 'string') return { record, mtimeMs: metadata.mtimeMs };
} catch {
// Treat a partially written/crashed lease as stale only by mtime.
}
return { mtimeMs: metadata.mtimeMs };
} catch {
return undefined;
}
}
function isStale(value: { record?: PiLeaseRecord; mtimeMs?: number } | undefined, now: number, staleAfterMs: number): boolean {
if (!value) return true;
if (value.record) return Date.parse(value.record.expiresAt) <= now;
return (value.mtimeMs ?? now) + staleAfterMs <= now;
}
function isAlreadyExists(error: unknown): boolean {
return typeof error === 'object' && error !== null && 'code' in error
&& (error as { code?: unknown }).code === 'EEXIST';
}
+12 -10
View File
@@ -17,7 +17,10 @@ export type SessionStatus =
| 'paused'
| 'requires_action'
| 'completed'
| 'failed';
| 'failed'
| 'cancelled'
| 'timed_out'
| 'cleanup_pending';
/**
* Valid state transitions for the Session state machine.
@@ -28,16 +31,15 @@ export type SessionStatus =
* at any point in its life — including while queued or idle (paused).
*/
export const SESSION_TRANSITIONS: Record<SessionStatus, SessionStatus[]> = {
queued: ['running', 'completed', 'failed'],
running: ['paused', 'requires_action', 'completed', 'failed'],
paused: ['running', 'completed', 'failed'],
requires_action: ['running', 'completed', 'failed'],
queued: ['running', 'completed', 'failed', 'cancelled', 'timed_out', 'cleanup_pending'],
running: ['paused', 'requires_action', 'completed', 'failed', 'cancelled', 'timed_out', 'cleanup_pending'],
paused: ['running', 'completed', 'failed', 'cancelled'],
requires_action: ['running', 'completed', 'failed', 'cancelled'],
completed: [],
// A failed session is recoverable: a new user message resumes it (failed →
// running), preserving the full event log. It can also be stopped/deleted,
// which drives it to the completed terminal state. Only completed is truly
// terminal (no outbound edges).
failed: ['running', 'completed'],
failed: ['running', 'completed', 'cancelled'],
cancelled: [],
timed_out: [],
cleanup_pending: [],
};
// ============================================================
+2
View File
@@ -40,6 +40,8 @@ export interface StrategyContext {
messages: CoreMessage[];
/** Resolved selected model configuration, including no AI SDK construction. */
modelConfig?: ModelConfig;
/** Explicit skill directories for engines that support managed skill loading. */
skillDirs?: string[];
/** Constructed only for strategies that require the AI SDK model transport. */
model?: LanguageModel;
tools: Record<string, CoreTool>;
+4 -4
View File
@@ -91,6 +91,7 @@ function createWindowsPiTerminationHarness(treeTermination: Promise<void>): Wind
queueMicrotask(() => piChild.emit('close', 0));
return treeTermination;
},
cleanupTimeoutMs: 25,
});
const manager = new SessionManager(db, undefined, 'pi');
const cleanupCalls: string[] = [];
@@ -155,10 +156,9 @@ describe('Pi Windows termination cleanup', () => {
void stop.then(() => { stopSettled = true; }, () => { stopSettled = true; });
await parentClosed;
rejectTreeTermination?.(new Error('taskkill failed'));
await Promise.resolve();
await Promise.resolve();
expect(stopSettled).toBe(false);
await waitFor(() => stopSettled);
expect(stopSettled).toBe(true);
expect(manager.get(sessionId)?.status).toBe('cleanup_pending');
expect(cleanupCalls).toEqual([]);
db.close();
});
+66
View File
@@ -0,0 +1,66 @@
import { Readable } from 'node:stream';
import { mkdtempSync, mkdirSync, rmSync, writeFileSync } from 'node:fs';
import { afterEach, describe, expect, it } from 'vitest';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { Database } from '@/core/db/database.js';
import { PiStrategy } from '@/strategy/pi-strategy.js';
import { recordPiSessionState, getPiSessionState } from '@/strategy/pi/session-continuity.js';
import type { StrategyContext } from '@/types/strategy.js';
const directories: string[] = [];
afterEach(() => {
for (const directory of directories.splice(0)) rmSync(directory, { recursive: true, force: true });
});
describe('Pi resume refusal', () => {
it('makes a Pi resume refusal visible and records continuity failure state', async () => {
const directory = mkdtempSync(join(tmpdir(), 'ma-pi-resume-refused-'));
directories.push(directory);
const db = new Database(join(directory, 'data.db'));
db.runMigrations();
db.exec(`INSERT INTO environments (id, name, config) VALUES ('env_default', 'local', '{}')`);
db.exec(`INSERT INTO agents (id, name, definition) VALUES ('agent_pi', 'pi-agent', '{}')`);
const sessionId = 'sess_pi_resume';
db.exec(`INSERT INTO sessions (id, agent_id, agent_name, environment_id, status, resources, vault_ids, loop_engine) VALUES ('${sessionId}', 'agent_pi', 'pi-agent', 'env_default', 'paused', '[]', '[]', 'pi')`);
const sessionFile = join(directory, 'pi-sessions', `${sessionId}.jsonl`);
mkdirSync(join(directory, 'pi-sessions'), { recursive: true });
writeFileSync(sessionFile, '{"type":"session","id":"pi-resume","version":1}\n');
recordPiSessionState(db, sessionId, sessionFile, { id: 'pi-resume', schemaVersion: '1' });
const strategy = new PiStrategy({
async launch() {},
async start() {
return {
child: {} as any,
sessionFile,
stdout: Readable.from(['{"type":"session","id":"pi-resume"}\n']),
stderr: Readable.from(['Stored session working directory does not exist\n']),
wait: async () => ({ code: 1, signal: null }),
terminate: async () => {},
};
},
}, db);
const context = {
session: {
id: sessionId, loopEngine: 'pi', agentId: 'agent_pi', agentName: 'pi-agent',
environmentId: 'env_default', status: 'running', createdAt: new Date(), updatedAt: new Date(),
},
userEvent: { type: 'user.message', content: [{ type: 'text', text: 'resume' }] },
systemPrompt: 'system', messages: [],
modelConfig: { name: 'fixture', provider: 'openai', model: 'fixture-model', api_key: 'fixture-key' },
tools: {}, sandbox: { sessionId, hostWorkDir: directory }, eventLog: {
append: (_id: string, event: any) => ({ id: 'sevt_1', sessionId, seq: 1, type: event.type, content: event.content, createdAt: new Date() } as any),
getLatestSeq: () => 0,
recordUsage: () => {},
}, broadcast: () => {}, config: {},
} as unknown as StrategyContext;
const run = async () => {
for await (const _event of strategy.execute(context)) {}
};
await expect(run()).rejects.toMatchObject({ code: 'pi_resume_refused' });
expect(getPiSessionState(db, sessionId)).toMatchObject({ status: 'pi_resume_refused' });
db.close();
});
});
+109
View File
@@ -0,0 +1,109 @@
import { afterEach, describe, expect, it } from 'vitest';
import { mkdtempSync, rmSync, writeFileSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { Database } from '@/core/db/database.js';
import {
assertPiSessionContinuity,
getPiSessionState,
inspectPiSessionFile,
recordPiSessionState,
} from '@/strategy/pi/session-continuity.js';
import { PiLauncher } from '@/strategy/pi-launcher.js';
const directories: string[] = [];
afterEach(() => {
for (const directory of directories.splice(0)) rmSync(directory, { recursive: true, force: true });
});
function setup() {
const directory = mkdtempSync(join(tmpdir(), 'ma-pi-continuity-'));
directories.push(directory);
const db = new Database(join(directory, 'data.db'));
db.runMigrations();
db.exec(`INSERT INTO environments (id, name, config) VALUES ('env_default', 'local', '{}')`);
db.exec(`INSERT INTO agents (id, name, definition) VALUES ('agent_pi', 'pi-agent', '{}')`);
const session = 'sess_pi_continuity';
db.exec(`INSERT INTO sessions (id, agent_id, agent_name, environment_id, status, resources, vault_ids, loop_engine) VALUES ('${session}', 'agent_pi', 'pi-agent', 'env_default', 'paused', '[]', '[]', 'pi')`);
return { db, directory, session, file: join(directory, 'session.jsonl') };
}
describe('Pi session continuity state', () => {
it('accepts a new empty file, then requires matching SQLite state for resume', () => {
const value = setup();
writeFileSync(value.file, '');
expect(assertPiSessionContinuity(value.db, value.session, value.file)).toEqual({});
writeFileSync(value.file, '{"type":"session","id":"pi-1","version":1}\n');
expect(() => assertPiSessionContinuity(value.db, value.session, value.file)).toThrow('without matching SandBase continuity state');
recordPiSessionState(value.db, value.session, value.file, { id: 'pi-1', schemaVersion: '1' });
expect(assertPiSessionContinuity(value.db, value.session, value.file).state).toMatchObject({ piSessionId: 'pi-1', schemaVersion: '1', status: 'active' });
expect(getPiSessionState(value.db, value.session)?.sessionFile).toBe(value.file);
value.db.close();
});
it('rejects a changed header identity, schema, and malformed header', () => {
const value = setup();
writeFileSync(value.file, '{"type":"session","id":"pi-1","version":1}\n');
recordPiSessionState(value.db, value.session, value.file, { id: 'pi-1', schemaVersion: '1' });
writeFileSync(value.file, '{"type":"session","id":"pi-2","version":1}\n');
expect(() => assertPiSessionContinuity(value.db, value.session, value.file)).toThrow('identity or schema');
writeFileSync(value.file, '{"type":"session","id":"pi-1","version":2}\n');
expect(() => assertPiSessionContinuity(value.db, value.session, value.file)).toThrow('identity or schema');
writeFileSync(value.file, '{"type":"not-session","id":"pi-1"}\n');
expect(() => inspectPiSessionFile(value.file)).toThrow('missing type=session');
value.db.close();
});
it('proves repeated launcher turns share one managed file and reject a concurrent writer', async () => {
const value = setup();
const script = join(value.directory, 'controlled-pi.mjs');
writeFileSync(script, `
import { writeFileSync } from 'node:fs';
const args = process.argv.slice(2);
const sessionFile = args[args.indexOf('--session') + 1];
if (args[args.indexOf('--thinking') + 1] !== 'medium' || !args.includes('--skill')) process.exit(2);
let input = '';
process.stdin.setEncoding('utf8');
process.stdin.on('data', (chunk) => { input += chunk; });
process.stdin.on('end', () => {
writeFileSync(sessionFile, JSON.stringify({ type: 'session', id: 'pi-real-fixture', version: 1 }) + '\\n');
process.stdout.write(JSON.stringify({ type: 'session', id: 'pi-real-fixture', version: 1 }) + '\\n');
process.stdout.write(JSON.stringify({ type: 'agent_end' }) + '\\n');
setTimeout(() => process.exit(0), 100);
});
`);
const launcher = new PiLauncher({
dataDir: value.directory,
database: value.db,
command: process.execPath,
commandArgs: [script],
timeoutMs: 2_000,
cleanupTimeoutMs: 100,
});
const request = {
sessionId: value.session,
workDir: value.directory,
prompt: '多行\\nUnicode ✓',
systemPrompt: 'fixture system',
model: { provider: 'openai', model: 'fixture-model', api_key: 'fixture-key' },
skillDirs: [value.directory] as string[],
thinkingLevel: 'medium' as const,
};
const first = await launcher.start(request);
await first.wait();
const sessionFile = join(value.directory, 'pi-sessions', `${value.session}.jsonl`);
recordPiSessionState(value.db, value.session, sessionFile, { id: 'pi-real-fixture', schemaVersion: '1' });
const second = await launcher.start(request);
await second.wait();
expect(getPiSessionState(value.db, value.session)).toMatchObject({ piSessionId: 'pi-real-fixture', status: 'active' });
const held = await launcher.start(request);
await expect(launcher.start(request)).rejects.toMatchObject({ code: 'pi_session_busy' });
await held.wait();
value.db.close();
});
});
+38
View File
@@ -0,0 +1,38 @@
import { afterEach, describe, expect, it } from 'vitest';
import { existsSync, mkdtempSync, rmSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { acquirePiSessionFileLease, PiSessionBusyError } from '@/strategy/pi/session-lease.js';
const directories: string[] = [];
afterEach(() => {
for (const directory of directories.splice(0)) rmSync(directory, { recursive: true, force: true });
});
describe('Pi session-file lease', () => {
it('rejects a live concurrent owner and removes its lease on release', async () => {
const directory = mkdtempSync(join(tmpdir(), 'ma-pi-lease-'));
directories.push(directory);
const sessionFile = join(directory, 'session.jsonl');
const first = await acquirePiSessionFileLease(sessionFile, { ownerId: 'owner-a', now: () => 1_000, staleAfterMs: 10_000 });
await expect(acquirePiSessionFileLease(sessionFile, { ownerId: 'owner-b', now: () => 2_000, staleAfterMs: 10_000 }))
.rejects.toMatchObject({ code: 'pi_session_busy' } satisfies Partial<PiSessionBusyError>);
expect(existsSync(`${sessionFile}.lease`)).toBe(true);
await first.release();
expect(existsSync(`${sessionFile}.lease`)).toBe(false);
});
it('recovers an expired owner through an atomic stale rename', async () => {
const directory = mkdtempSync(join(tmpdir(), 'ma-pi-lease-stale-'));
directories.push(directory);
const sessionFile = join(directory, 'session.jsonl');
const first = await acquirePiSessionFileLease(sessionFile, { ownerId: 'owner-a', now: () => 1_000, staleAfterMs: 10 });
const recovered = await acquirePiSessionFileLease(sessionFile, { ownerId: 'owner-b', now: () => 2_000, staleAfterMs: 10 });
expect(recovered.recoveredStale).toBe(true);
await first.release();
await recovered.release();
expect(existsSync(`${sessionFile}.lease`)).toBe(false);
});
});
+15
View File
@@ -0,0 +1,15 @@
import { describe, expect, it } from 'vitest';
import { isTerminal, transition } from '@/core/session/state-machine.js';
import type { SessionStatus } from '@/types/session.js';
describe('Pi terminal lifecycle states', () => {
it.each(['cancelled', 'timed_out', 'cleanup_pending'] as SessionStatus[])('treats %s as terminal and distinct from completed', (status) => {
expect(isTerminal(status)).toBe(true);
expect(status).not.toBe('completed');
expect(transition('running', status)).toBe(status);
});
it('keeps cleanup_pending from becoming a normal completed release', () => {
expect(isTerminal('cleanup_pending')).toBe(true);
});
});