diff --git a/AGENTS.md b/AGENTS.md index e1afb66..df6d8cf 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -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, diff --git a/CHANGELOG.md b/CHANGELOG.md index dd46934..28b8010 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 05af7b8..2f51fc6 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -73,8 +73,12 @@ git worktree add .worktrees/ -b feat/ 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/`, `fix/`, `docs/`, `test/`, - `build/`, or `chore/`. Issue-driven work uses - `fix/issue--` or `feat/issue--`. + `build/`, or `chore/`. Issue-driven work may use a public + `fix/issue--` or `feat/issue--` 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 diff --git a/docs/api-matrix.md b/docs/api-matrix.md index 3f74b0e..4feaebf 100644 --- a/docs/api-matrix.md +++ b/docs/api-matrix.md @@ -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. | diff --git a/docs/api.md b/docs/api.md index a3a7c68..a7365bf 100644 --- a/docs/api.md +++ b/docs/api.md @@ -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. diff --git a/docs/pi-loop-engine.md b/docs/pi-loop-engine.md index 4bab009..be9c810 100644 --- a/docs/pi-loop-engine.md +++ b/docs/pi-loop-engine.md @@ -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 diff --git a/docs/test-evidence/pi-conformance.md b/docs/test-evidence/pi-conformance.md new file mode 100644 index 0000000..fa00133 --- /dev/null +++ b/docs/test-evidence/pi-conformance.md @@ -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. diff --git a/src/api/standard.ts b/src/api/standard.ts index 9bd05e7..6e3d35d 100644 --- a/src/api/standard.ts +++ b/src/api/standard.ts @@ -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'; } diff --git a/src/core/db/migrations.ts b/src/core/db/migrations.ts index f3d0ccd..70b23d8 100644 --- a/src/core/db/migrations.ts +++ b/src/core/db/migrations.ts @@ -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 }, ]; diff --git a/src/core/runtime/loop-engine-bootstrap.ts b/src/core/runtime/loop-engine-bootstrap.ts index 2d51784..de06698 100644 --- a/src/core/runtime/loop-engine-bootstrap.ts +++ b/src/core/runtime/loop-engine-bootstrap.ts @@ -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> = { 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; diff --git a/src/core/runtime/session-runtime.ts b/src/core/runtime/session-runtime.ts index 4922e6f..089a232 100644 --- a/src/core/runtime/session-runtime.ts +++ b/src/core/runtime/session-runtime.ts @@ -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; 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, diff --git a/src/core/session/delegation-service.ts b/src/core/session/delegation-service.ts index 1a48ae7..9848341 100644 --- a/src/core/session/delegation-service.ts +++ b/src/core/session/delegation-service.ts @@ -34,7 +34,7 @@ export interface DelegationServiceDeps { provisionSandbox: (session: Session, sandboxId: string) => Promise; composeSystemPrompt: (agent: AgentDefinition) => string; buildSandboxTools: (agent: AgentDefinition, sandbox: SandboxInstance) => Record; -} + 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, diff --git a/src/core/session/executor.ts b/src/core/session/executor.ts index 62092d5..4659d0d 100644 --- a/src/core/session/executor.ts +++ b/src/core/session/executor.ts @@ -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. diff --git a/src/core/session/session-lifecycle.ts b/src/core/session/session-lifecycle.ts index 8019768..e4ec622 100644 --- a/src/core/session/session-lifecycle.ts +++ b/src/core/session/session-lifecycle.ts @@ -5,6 +5,9 @@ const STATUS_TO_EVENT: Partial> = { 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', diff --git a/src/core/session/session-manager.ts b/src/core/session/session-manager.ts index dd9ed57..1515dd8 100644 --- a/src/core/session/session-manager.ts +++ b/src/core/session/session-manager.ts @@ -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(['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; +} diff --git a/src/core/session/state-machine.ts b/src/core/session/state-machine.ts index 6ae1823..57b3f51 100644 --- a/src/core/session/state-machine.ts +++ b/src/core/session/state-machine.ts @@ -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'; diff --git a/src/core/settings/adapters.ts b/src/core/settings/adapters.ts index 3a84364..fd8b3f7 100644 --- a/src/core/settings/adapters.ts +++ b/src/core/settings/adapters.ts @@ -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()), diff --git a/src/index.ts b/src/index.ts index b2ab5ce..e2c66be 100644 --- a/src/index.ts +++ b/src/index.ts @@ -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, diff --git a/src/strategy/pi-launcher.ts b/src/strategy/pi-launcher.ts index 418f349..5c56243 100644 --- a/src/strategy/pi-launcher.ts +++ b/src/strategy/pi-launcher.ts @@ -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; terminate(force?: boolean): Promise; } @@ -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 { - // 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 => { + 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 { @@ -344,6 +462,7 @@ async function spawnPiProcess( let childClosed = false; let windowsTreeTerminationComplete = false; let forceTimer: ReturnType | undefined; + let cleanupTimer: ReturnType | 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 { + 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 diff --git a/src/strategy/pi-strategy.ts b/src/strategy/pi-strategy.ts index d784a98..7c44f99 100644 --- a/src/strategy/pi-strategy.ts +++ b/src/strategy/pi-strategy.ts @@ -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, tail: PiStderrTail): Promise { for await (const chunk of stream) tail.append(chunk); } diff --git a/src/strategy/pi/session-continuity.ts b/src/strategy/pi/session-continuity.ts new file mode 100644 index 0000000..1465822 --- /dev/null +++ b/src/strategy/pi/session-continuity.ts @@ -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; + 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); +} diff --git a/src/strategy/pi/session-lease.ts b/src/strategy/pi/session-lease.ts new file mode 100644 index 0000000..a7f6d70 --- /dev/null +++ b/src/strategy/pi/session-lease.ts @@ -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; + /** Stop heartbeats without deleting the lease when cleanup ownership is unknown. */ + suspendHeartbeat(): void; + release(): Promise; +} + +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 { + 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 = 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; + 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'; +} diff --git a/src/types/session.ts b/src/types/session.ts index 63bc626..dd89909 100644 --- a/src/types/session.ts +++ b/src/types/session.ts @@ -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 = { - 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: [], }; // ============================================================ diff --git a/src/types/strategy.ts b/src/types/strategy.ts index 6a9fc2d..0acf8d0 100644 --- a/src/types/strategy.ts +++ b/src/types/strategy.ts @@ -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; diff --git a/tests/unit/pi-engine-session.test.ts b/tests/unit/pi-engine-session.test.ts index ae7d073..ef7499a 100644 --- a/tests/unit/pi-engine-session.test.ts +++ b/tests/unit/pi-engine-session.test.ts @@ -91,6 +91,7 @@ function createWindowsPiTerminationHarness(treeTermination: Promise): 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(); }); diff --git a/tests/unit/pi-resume-refusal.test.ts b/tests/unit/pi-resume-refusal.test.ts new file mode 100644 index 0000000..e7460aa --- /dev/null +++ b/tests/unit/pi-resume-refusal.test.ts @@ -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(); + }); +}); diff --git a/tests/unit/pi-session-continuity.test.ts b/tests/unit/pi-session-continuity.test.ts new file mode 100644 index 0000000..34a680d --- /dev/null +++ b/tests/unit/pi-session-continuity.test.ts @@ -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(); + }); +}); diff --git a/tests/unit/pi-session-lease.test.ts b/tests/unit/pi-session-lease.test.ts new file mode 100644 index 0000000..5c867b5 --- /dev/null +++ b/tests/unit/pi-session-lease.test.ts @@ -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); + 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); + }); +}); diff --git a/tests/unit/pi-status-lifecycle.test.ts b/tests/unit/pi-status-lifecycle.test.ts new file mode 100644 index 0000000..15cdd2e --- /dev/null +++ b/tests/unit/pi-status-lifecycle.test.ts @@ -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); + }); +});