From 39a546fd87e32186774e4e514c3a70ec8ec4ca23 Mon Sep 17 00:00:00 2001 From: Repository verification Date: Sat, 19 Sep 2026 04:14:06 +0000 Subject: [PATCH] refactor(improvement): consolidate training validation and execution lifecycle --- CHANGELOG.md | 10 ++ api-surface.json | 6 +- bench/CHANGELOG.md | 4 + bench/package.json | 2 +- docs/api/index.md | 16 +- docs/api/primitive-catalog.md | 5 +- docs/canonical-api.md | 32 +++- package.json | 2 +- scripts/verify-package-exports.mjs | 9 ++ src/improvement/candidate-validation.ts | 25 +++ src/improvement/improve-types.ts | 1 + src/improvement/improve.test.ts | 18 +++ src/improvement/improve.ts | 22 +-- src/improvement/meta-harness.test.ts | 25 +++ src/improvement/method-execution.ts | 7 +- .../profile-improvement-harness.ts | 26 ++- src/improvement/training.ts | 150 +++++++++++------- src/mcp/local-harness.ts | 68 +------- src/runtime/process-tree.ts | 69 ++++++++ .../fixtures/agent-improvement-proposal.json | 10 +- .../agent-profile-improvement-proposal.json | 6 +- tests/profile-training-boundaries.test.ts | 149 +++++++++++++++++ tests/profile-training.test.ts | 80 +++++++++- 23 files changed, 567 insertions(+), 175 deletions(-) create mode 100644 src/improvement/candidate-validation.ts create mode 100644 src/runtime/process-tree.ts diff --git a/CHANGELOG.md b/CHANGELOG.md index 054278449..aeba09aaa 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,15 @@ # Changelog +## 0.243.0 + +Training, optimization, and bound harnesses now share candidate-validator admission and invocation. **Migration:** validators must return `undefined` synchronously or throw. Promises and other return values are rejected instead of silently accepting an unchecked candidate; use a block body for side effects. Composed optimizer leaves obey the same rule, while original callbacks remain bound into execution identity. + +`improve` and `createProfileImprovementHarness` accept their own frozen profile results directly. Retraining updates small-model, subagent, and mode references to the prior receipted checkpoint, preserving unrelated model choices. Serving evidence is snapshotted before asynchronous checkpoint revalidation. + +Training's `timeoutMs` is optional and uses the existing chunked deadline timer when supplied; there is no implicit deadline or seven-day ceiling. Cancellation still stops admission and waiting, without claiming remote cleanup. The command trainer snapshots direct-call requests before asynchronous work, honors caller-selected output budgets, accepts empty pinned configuration files, and streams declared input hashes without a one-gigabyte ceiling. Checkpoints remain bounded and nonempty; dataset bounds, partition checks, receipt identity, and publication ordering are unchanged. + +Command training and coding harnesses share one confirmed process-group teardown implementation, allowing graceful final writes before escalation. Existing coding-harness timing policy and import paths are preserved. No new dependency, executor, training algorithm, dataset exporter, scheduler, or deployment client is added. + ## 0.242.0 `improve(profile, { mode: 'training', ... })` and the bound harness's `train` method now produce a checkpoint-backed candidate with an Interface training receipt. The command trainer pins its executable and inputs, supplies only an explicit public environment, and cancels its POSIX process group. Managed trainers and verified serving adapters use the same typed boundary; they own remote job cleanup. diff --git a/api-surface.json b/api-surface.json index da8d5ba8f..dffaf00ad 100644 --- a/api-surface.json +++ b/api-surface.json @@ -48,7 +48,7 @@ "ConversationStreamEvent": "type a39182b3dbfc", "ConversationTurn": "type d8280ca3c636", "CreateKnowledgeImprovementActivationExecutorOptions": "type 4b3e5fe02df7", - "CreateProfileImprovementHarnessOptions": "type 36de01aba28e", + "CreateProfileImprovementHarnessOptions": "type ac085504184e", "D1DatabaseLike": "type ffce9ec5de30", "D1StmtLike": "type cd8b3c46cbcd", "DEFAULT_ROUTER_BASE_URL": "value 4db78e51a917", @@ -94,7 +94,7 @@ "ImproveScenarioPartitions": "type 37a3508406b1", "ImproveSkillsOptions": "type c1f5a69faefc", "ImproveSurface": "type b711b683b151", - "ImproveTrainingOptions": "type 8bfc42270e58", + "ImproveTrainingOptions": "type 0d216f385a91", "ImproveTrainingResult": "type 81c28cf9d5ab", "ImprovementCandidate": "type 0c22a91c6396", "ImprovementCodeCandidate": "type 588fa6d3b2f5", @@ -253,7 +253,7 @@ "formatSupervisedKnowledgeTask": "value bdcf6b28157d", "generateSpanId": "value 2f8329045cac", "getModels": "value 95cb7c012c48", - "improve": "value 45f11e9feb52", + "improve": "value 4c0da1dd4fbf", "isDelegatedLoopMode": "value d0f2042750ec", "knowledgeReadinessDeliverable": "value f6f33b24a926", "loopEventToOtelSpan": "value 63bec8b09ae0", diff --git a/bench/CHANGELOG.md b/bench/CHANGELOG.md index 0f13076a1..6bc25daa1 100644 --- a/bench/CHANGELOG.md +++ b/bench/CHANGELOG.md @@ -1,5 +1,9 @@ # Changelog +## 0.13.5 + +Consume Runtime 0.243.0 through the existing workspace dependency. Benchmark APIs and grading behavior are unchanged. + ## 0.13.4 Require Interface `^2.10.0` and consume Runtime 0.242.0 through the published dependency ranges, keeping benchmark consumers on the checkpoint-training receipt contract. diff --git a/bench/package.json b/bench/package.json index d521260f1..fdfcae30d 100644 --- a/bench/package.json +++ b/bench/package.json @@ -1,6 +1,6 @@ { "name": "@tangle-network/agent-bench", - "version": "0.13.4", + "version": "0.13.5", "type": "module", "description": "Benchmark adapters and execution for agent-runtime across coding, tool-use, RAG, memory, browser, and terminal tasks.", "repository": { diff --git a/docs/api/index.md b/docs/api/index.md index c19534784..350260c5d 100644 --- a/docs/api/index.md +++ b/docs/api/index.md @@ -3528,7 +3528,7 @@ Exact materialized profile presented for validation before any candidate run. ##### profile -> **profile**: `AgentProfile` +> **profile**: `object` Exact baseline profile. It is parsed, detached, and frozen at construction. @@ -3917,9 +3917,11 @@ Pins trainer, serving adapter and their private dependencies, just like the boun > **outputDirectory**: `string` -##### timeoutMs +##### timeoutMs? + +> `optional` **timeoutMs?**: `number` -> **timeoutMs**: `number` +Optional overall deadline. Omit to rely on caller cancellation; long durations are supported. ##### maxCheckpointBytes @@ -3983,6 +3985,8 @@ Explicit public environment only. Ambient credentials are never inherited. > **maxOutputBytes**: `number` +Total stdout + stderr byte budget. Output is drained, not retained in memory. + *** ### CreateKnowledgeImprovementActivationExecutorOptions @@ -7503,6 +7507,8 @@ Runs one exact materialized profile on one scenario. > **ImproveCandidateValidator** = (`input`) => `void` +Accept by returning void synchronously; reject by throwing. Async callbacks are refused. + #### Parameters ##### input @@ -9002,8 +9008,6 @@ Train and serve a checkpoint without implying that it improved held-out quality. ###### profile -`AgentProfile` - ###### opts [`ImproveTrainingOptions`](#improvetrainingoptions) @@ -9032,8 +9036,6 @@ Optimize one exact profile surface with a complete method. ###### profile -`AgentProfile` - ###### opts [`ImproveMethodOptions`](#improvemethodoptions)\<`TScenario`, `TArtifact`\> diff --git a/docs/api/primitive-catalog.md b/docs/api/primitive-catalog.md index 651bafb8e..d1d6b2add 100644 --- a/docs/api/primitive-catalog.md +++ b/docs/api/primitive-catalog.md @@ -7,7 +7,7 @@ # Primitive catalog — the never-stale anti-reinvention inventory -> **GENERATED** from `@tangle-network/agent-runtime@0.242.0` and `@tangle-network/agent-eval@0.182.0` by `scripts/gen-primitive-catalog.mjs`. Do NOT hand-edit — run `pnpm run docs:api`. This is the mechanical companion to the JUDGMENT in `canonical-api.md` (§2 decision table + §1.5 AgentProfile law): that doc says WHICH primitive to reach for and what NOT to build; this catalog proves WHAT exists. Per-symbol signatures + `file:line` live in the per-module pages under `docs/api/`. +> **GENERATED** from `@tangle-network/agent-runtime@0.243.0` and `@tangle-network/agent-eval@0.182.0` by `scripts/gen-primitive-catalog.mjs`. Do NOT hand-edit — run `pnpm run docs:api`. This is the mechanical companion to the JUDGMENT in `canonical-api.md` (§2 decision table + §1.5 AgentProfile law): that doc says WHICH primitive to reach for and what NOT to build; this catalog proves WHAT exists. Per-symbol signatures + `file:line` live in the per-module pages under `docs/api/`. ## 1. agent-runtime — own public surface @@ -156,6 +156,7 @@ Import from `@tangle-network/agent-runtime` — 298 exports. | `AgentEvalErrorCode` | type | Error taxonomy for `@tangle-network/agent-eval`. | | `AgenticGeneratorShotDisposition` | type | Worktree decision emitted before a completed shot is retried, accepted, or | | `AgenticGeneratorShotExecution` | type | Runtime's exact terminal turn plus its complete normalized event stream. | +| `ImproveCandidateValidator` | type | Accept by returning void synchronously; reject by throwing. Async callbacks are refused. | | `ImproveCodeRunOptions` | type | Runtime-owned code search in isolated git worktrees. | | `ImproveMethodFactory` | type | Build a complete method after trace findings are available. | | `ImproveMethodOptions` | type | Complete-method configuration for every non-code profile surface. | @@ -175,7 +176,7 @@ Import from `@tangle-network/agent-runtime` — 298 exports. | `Verifier` | type | Verifies the edited worktree. Sync or async; throws only on a setup fault | | `WorktreeCheckRunner` | type | The single shell-command-in-worktree runner seam (replaces the per-executor copies). | -**Undocumented supporting types** (add a TSDoc line at the declaration to earn a table row): `AgentAdapter`, `AgentBackendContext`, `AgentBackendInput`, `AgentExecutionBackend`, `AgenticGeneratorOptions`, `AgenticGeneratorShotReceipt`, `AgentKnowledgeProvider`, `AgentKnowledgeReadinessCheckOptions`, `AgentTaskContext`, `AgentTaskRunResult`, `AgentTaskSpec`, `BackendCallPolicy`, `ChatModelCandidate`, `CheckpointServingPort`, `ControlBudget`, `ControlEvalResult`, `ControlledTrainingCommand`, `ControlRunResult`, `ControlStep`, `Conversation`, `ConversationDriveState`, `ConversationJournal`, `ConversationJournalEntry`, `ConversationParticipant`, `ConversationPolicy`, `ConversationResult`, `ConversationTurn`, `CreateKnowledgeImprovementActivationExecutorOptions`, `CreateProfileImprovementHarnessOptions`, `D1StmtLike`, `DataAcquisitionPlan`, `DelegatedLoopResult`, `EvalRunEvent`, `EvalRunGeneration`, `EvalRunsExportConfig`, `EvalRunsExportResult`, `HaltContext`, `HaltSignal`, `ImproveCodeBaseOptions`, `ImproveCodeResult`, `ImproveCustomCodeGeneratorOptions`, `ImprovementCodeCandidate`, `ImprovementProfileCandidate`, `ImproveMethodContext`, `ImproveMethodResult`, `ImproveRuntimeCodeGeneratorOptions`, `ImproveSkillsOptions`, `ImproveTrainingOptions`, `KnowledgeImprovementActivationExecutor`, `KnowledgeImprovementCandidatePair`, `KnowledgeImprovementExperimentBundles`, `KnowledgeImprovementJobMeasurement`, `KnowledgeImprovementJobResult`, `KnowledgeReadinessCheckInput`, `KnowledgeReadinessDecision`, `KnowledgeReadinessReport`, `KnowledgeRequirement`, `LoopRunnerCliArgs`, `LoopRunnerCliResult`, `McpServeSpec`, `OfficialSensitiveCandidateInput`, `OtelAttribute`, `OtelExportConfig`, `OtelExporter`, `OtelSpan`, `PersonaConversationResult`, `ProfileTrainerRequest`, `RawTraceDistillerOptions`, `ReflectiveGeneratorOptions`, `ResearchLoopResult`, `ResearchLoopRunnerOptions`, `ResolvedChatModel`, `RunAgentTaskOptions`, `RunAgentTaskStreamOptions`, `RunConversationOptions`, `RunDelegatedLoopOptions`, `RunKnowledgeImprovementJobOptions`, `RunPersonaConfig`, `RunPersonaConversationOptions`, `RuntimeDecisionEvidenceRef`, `RuntimeDecisionPoint`, `RuntimeEventCollector`, `RuntimeEventOtelOptions`, `RuntimeHookContext`, `RuntimeHookErrorContext`, `RuntimeHookEvent`, `RuntimeRunCompleteInput`, `RuntimeRunCost`, `RuntimeRunHandle`, `RuntimeRunOptions`, `RuntimeRunPersistenceAdapter`, `RuntimeRunRow`, `RuntimeSession`, `RuntimeSessionStore`, `RuntimeStreamEventCollector`, `RuntimeStreamEventSummary`, `RuntimeTelemetryOptions`, `SanitizedKnowledgeReadinessReport`, `SanitizedKnowledgeRequirement`, `ServerSentEventOptions`, `SupervisedKnowledgeUpdateInput`, `SupervisedKnowledgeUpdateOptions`, `SupervisedKnowledgeUpdateResult`, `TrainingDatasetDocument`, `VetoedFact`, `WorktreeLoopRunnerOptions`, `AgenticGeneratorExecutorForWorktree`, `AgentRuntimeEvent`, `AgentRuntimeEventSink`, `AgentTaskStatus`, `AuthSource`, `ChatModelValidation`, `ControlDecision`, `ConversationStreamEvent`, `DeepReadonly`, `DelegatedLoopMode`, `DelegatedLoopRegistry`, `DelegatedLoopRunner`, `HaltPredicate`, `HaltReason`, `ImproveCandidateValidator`, `ImproveCodeOptions`, `ImprovementCandidate`, `ImprovementProfileCandidatePopulation`, `ImprovementProfilePopulationCandidate`, `ImprovementProfilePopulationLineage`, `ImproveMethodSource`, `ImproveOptimizationRunOptions`, `ImproveProfileSurface`, `ImproveResult`, `ImproveTrainingResult`, `KnowledgeReadinessCheck`, `KnowledgeReadinessCheckResult`, `ProfileImprovementHarnessRunOptions`, `ProfileImprovementHarnessTrainOptions`, `RuntimeDecisionKind`, `RuntimeHookTarget`, `RuntimeRunStatus`, `RuntimeStreamEvent`, `RuntimeStreamEventSink`, `SupervisedKnowledgeUpdater`, `TrainingBoundaryResult`, `TurnOrder`. +**Undocumented supporting types** (add a TSDoc line at the declaration to earn a table row): `AgentAdapter`, `AgentBackendContext`, `AgentBackendInput`, `AgentExecutionBackend`, `AgenticGeneratorOptions`, `AgenticGeneratorShotReceipt`, `AgentKnowledgeProvider`, `AgentKnowledgeReadinessCheckOptions`, `AgentTaskContext`, `AgentTaskRunResult`, `AgentTaskSpec`, `BackendCallPolicy`, `ChatModelCandidate`, `CheckpointServingPort`, `ControlBudget`, `ControlEvalResult`, `ControlledTrainingCommand`, `ControlRunResult`, `ControlStep`, `Conversation`, `ConversationDriveState`, `ConversationJournal`, `ConversationJournalEntry`, `ConversationParticipant`, `ConversationPolicy`, `ConversationResult`, `ConversationTurn`, `CreateKnowledgeImprovementActivationExecutorOptions`, `CreateProfileImprovementHarnessOptions`, `D1StmtLike`, `DataAcquisitionPlan`, `DelegatedLoopResult`, `EvalRunEvent`, `EvalRunGeneration`, `EvalRunsExportConfig`, `EvalRunsExportResult`, `HaltContext`, `HaltSignal`, `ImproveCodeBaseOptions`, `ImproveCodeResult`, `ImproveCustomCodeGeneratorOptions`, `ImprovementCodeCandidate`, `ImprovementProfileCandidate`, `ImproveMethodContext`, `ImproveMethodResult`, `ImproveRuntimeCodeGeneratorOptions`, `ImproveSkillsOptions`, `ImproveTrainingOptions`, `KnowledgeImprovementActivationExecutor`, `KnowledgeImprovementCandidatePair`, `KnowledgeImprovementExperimentBundles`, `KnowledgeImprovementJobMeasurement`, `KnowledgeImprovementJobResult`, `KnowledgeReadinessCheckInput`, `KnowledgeReadinessDecision`, `KnowledgeReadinessReport`, `KnowledgeRequirement`, `LoopRunnerCliArgs`, `LoopRunnerCliResult`, `McpServeSpec`, `OfficialSensitiveCandidateInput`, `OtelAttribute`, `OtelExportConfig`, `OtelExporter`, `OtelSpan`, `PersonaConversationResult`, `ProfileTrainerRequest`, `RawTraceDistillerOptions`, `ReflectiveGeneratorOptions`, `ResearchLoopResult`, `ResearchLoopRunnerOptions`, `ResolvedChatModel`, `RunAgentTaskOptions`, `RunAgentTaskStreamOptions`, `RunConversationOptions`, `RunDelegatedLoopOptions`, `RunKnowledgeImprovementJobOptions`, `RunPersonaConfig`, `RunPersonaConversationOptions`, `RuntimeDecisionEvidenceRef`, `RuntimeDecisionPoint`, `RuntimeEventCollector`, `RuntimeEventOtelOptions`, `RuntimeHookContext`, `RuntimeHookErrorContext`, `RuntimeHookEvent`, `RuntimeRunCompleteInput`, `RuntimeRunCost`, `RuntimeRunHandle`, `RuntimeRunOptions`, `RuntimeRunPersistenceAdapter`, `RuntimeRunRow`, `RuntimeSession`, `RuntimeSessionStore`, `RuntimeStreamEventCollector`, `RuntimeStreamEventSummary`, `RuntimeTelemetryOptions`, `SanitizedKnowledgeReadinessReport`, `SanitizedKnowledgeRequirement`, `ServerSentEventOptions`, `SupervisedKnowledgeUpdateInput`, `SupervisedKnowledgeUpdateOptions`, `SupervisedKnowledgeUpdateResult`, `TrainingDatasetDocument`, `VetoedFact`, `WorktreeLoopRunnerOptions`, `AgenticGeneratorExecutorForWorktree`, `AgentRuntimeEvent`, `AgentRuntimeEventSink`, `AgentTaskStatus`, `AuthSource`, `ChatModelValidation`, `ControlDecision`, `ConversationStreamEvent`, `DeepReadonly`, `DelegatedLoopMode`, `DelegatedLoopRegistry`, `DelegatedLoopRunner`, `HaltPredicate`, `HaltReason`, `ImproveCodeOptions`, `ImprovementCandidate`, `ImprovementProfileCandidatePopulation`, `ImprovementProfilePopulationCandidate`, `ImprovementProfilePopulationLineage`, `ImproveMethodSource`, `ImproveOptimizationRunOptions`, `ImproveProfileSurface`, `ImproveResult`, `ImproveTrainingResult`, `KnowledgeReadinessCheck`, `KnowledgeReadinessCheckResult`, `ProfileImprovementHarnessRunOptions`, `ProfileImprovementHarnessTrainOptions`, `RuntimeDecisionKind`, `RuntimeHookTarget`, `RuntimeRunStatus`, `RuntimeStreamEvent`, `RuntimeStreamEventSink`, `SupervisedKnowledgeUpdater`, `TrainingBoundaryResult`, `TurnOrder`. ### Vertical agent — manifest + surface proposal source diff --git a/docs/canonical-api.md b/docs/canonical-api.md index 55c6b655d..3b412a2ce 100644 --- a/docs/canonical-api.md +++ b/docs/canonical-api.md @@ -4,7 +4,7 @@ Generated signatures and the complete export list live in docs/api/. Run pnpm docs:freshness after editing this file. --> -> **Version 0.242.0.** +> **Version 0.243.0.** > [`docs/api/primitive-catalog.md`](./api/primitive-catalog.md) lists every export and import path. > `agent-eval` must satisfy `>=0.182.0 <0.183.0`. > `sandbox` must satisfy `>=0.36.4 <0.42.0`. @@ -292,9 +292,35 @@ artifact-addressed Router model identity. Runtime retains complete receipt ances checkpoint bytes after serving and candidate validation, and durably publishes the receipt before the runnable profile. Use the existing benchmark and held-out gates to assess that profile. -The controlled command trainer runs without a shell or inherited credentials and cancels its -POSIX process group. Managed training and serving adapters own their remote jobs and cleanup; +Pass a returned frozen profile directly into `improve` or a new bound harness; no mutable cast +or reconstruction is needed. When retraining a checkpoint-backed profile, references to that +same checkpoint in the small model, subagents, and modes follow the new receipt. Unrelated +model choices remain unchanged. + +Training has no implicit deadline. Omit `timeoutMs` to rely on caller cancellation, or specify a +positive safe duration; the shared deadline timer supports long jobs without native timer overflow. +A managed adapter that ignores cancellation may continue working after Runtime stops awaiting it. + +The controlled command trainer runs without a shell or inherited credentials. It snapshots the +request before asynchronous work, streams pinned input hashes (empty configuration files are +valid), and drains stdout/stderr within the caller's positive `maxOutputBytes` budget. Checkpoints +must still be nonempty and fit `maxCheckpointBytes`; dataset snapshot bounds remain in force. +It uses the same confirmed POSIX process-group teardown as the coding harnesses: permit graceful +shutdown, then escalate if necessary. This is trusted host execution, not an OS sandbox; separately +detached sessions are outside the owned process group. +Managed training and serving adapters own their remote jobs and cleanup; inspect the failure stage and the training/serving uncertainty flags rather than assuming a timeout removed external resources. A cancellation before adapter dispatch starts no job. Once profile publication commits, later cancellation does not retract the committed result. This is a local execution primitive, not a durable remote-job scheduler or a Router deployment API. + + +### Candidate validation has one acceptance rule + +Across training, optimization, composed method leaves, and bound harnesses, `validateCandidate` +accepts by returning `undefined` synchronously and rejects by throwing. A promise or another +return value is an error, not an accepted candidate. Use a block body for side effects, rather +than returning the result of an assertion or array operation. Do asynchronous preparation before +calling the improvement API; keep measurement and held-out decisions in the existing evaluators. +The original callback remains part of execution identity, so centralizing its invocation does +not collapse distinct validator implementations into one cache identity. diff --git a/package.json b/package.json index c3a20456d..a25167109 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@tangle-network/agent-runtime", - "version": "0.242.0", + "version": "0.243.0", "description": "Shared task-lifecycle skeleton for agents: a recursive loop kernel for chat turns, one-shot tasks, and multi-attempt loops, with trace capture and eval-gated self-improvement. Domain behavior lives in adapters; scoring and ship-gates in @tangle-network/agent-eval.", "homepage": "https://github.com/tangle-network/agent-runtime#readme", "repository": { diff --git a/scripts/verify-package-exports.mjs b/scripts/verify-package-exports.mjs index bf4df135f..1fa2ef105 100644 --- a/scripts/verify-package-exports.mjs +++ b/scripts/verify-package-exports.mjs @@ -307,6 +307,15 @@ try { const harnessTraining: Promise = profileHarness.train(trainingOptions) if (trainingResult.succeeded) { const trainingReceipt: AgentTrainingReceipt = trainingResult.receipt + const { timeoutMs: _timeout, ...withoutDeadline } = trainingOptions + const retrained: Promise = improve(trainingResult.profile, withoutDeadline) + const rebound = createProfileImprovementHarness({ + profile: trainingResult.profile, + executionRef: trainingOptions.executionRef, + agent: async () => 'fixture', + }) + void retrained + void rebound void trainingReceipt } void training diff --git a/src/improvement/candidate-validation.ts b/src/improvement/candidate-validation.ts new file mode 100644 index 000000000..ec055c601 --- /dev/null +++ b/src/improvement/candidate-validation.ts @@ -0,0 +1,25 @@ +import { ConfigError } from '../errors' +import type { ImproveCandidateValidationInput, ImproveCandidateValidator } from './improve-types' + +/** Validate the callback itself before any optimizer, trainer, or candidate work starts. */ +export function assertCandidateValidator( + validator: unknown, +): asserts validator is ImproveCandidateValidator | undefined { + if (validator !== undefined && typeof validator !== 'function') { + throw new ConfigError('validateCandidate must be a function when present') + } +} + +/** A validator accepts synchronously by returning void, or rejects by throwing. */ +export function validateProfileCandidate( + validator: ImproveCandidateValidator | undefined, + input: ImproveCandidateValidationInput, +): void { + assertCandidateValidator(validator) + const result: unknown = validator?.(Object.freeze(input)) + if (result !== undefined) { + // Observe a rejected async callback without treating its eventual result as admission. + void Promise.resolve(result).catch(() => {}) + throw new ConfigError('candidate validators must return void synchronously or throw') + } +} diff --git a/src/improvement/improve-types.ts b/src/improvement/improve-types.ts index e82e361e4..6d8489ba7 100644 --- a/src/improvement/improve-types.ts +++ b/src/improvement/improve-types.ts @@ -77,6 +77,7 @@ export interface ImproveCandidateValidationInput { isBaseline: boolean } +/** Accept by returning void synchronously; reject by throwing. Async callbacks are refused. */ export type ImproveCandidateValidator = (input: ImproveCandidateValidationInput) => void export type ImproveOptimizationRunOptions = Omit< diff --git a/src/improvement/improve.test.ts b/src/improvement/improve.test.ts index 2ee0e5c6b..c91a18cc6 100644 --- a/src/improvement/improve.test.ts +++ b/src/improvement/improve.test.ts @@ -409,6 +409,24 @@ describe('improve method execution', () => { expect(forbiddenCalls).toBe(0) }) + it('does not ignore an async validator on a composed optimization leaf', async () => { + let executions = 0 + const leaf = withMethodRuntimeControls(fixedMethod('improved'), { + costAttribution: 'optimizer-run', + validateCandidate: async () => {}, + }) + await expect( + improve(promptProfile(), { + ...methodOptions(sequentialOptimizationMethod({ name: 'async-leaf', methods: [leaf] })), + agent: async (profile, scenario, context) => { + executions++ + return paidProfile(profile, scenario, context) + }, + }), + ).rejects.toThrow(/synchronous/) + expect(executions).toBe(0) + }) + it('runs a complete method without exposing final-test cases and materializes its prompt', async () => { let observed: OptimizationMethodInput | undefined let observedEvaluationRef = '' diff --git a/src/improvement/improve.ts b/src/improvement/improve.ts index 2dea7eb7b..d0fd0eefc 100644 --- a/src/improvement/improve.ts +++ b/src/improvement/improve.ts @@ -10,7 +10,7 @@ */ import type { Scenario } from '@tangle-network/agent-eval/contract' -import { type AgentProfile, agentProfileSchema } from '@tangle-network/agent-interface' +import { agentProfileSchema } from '@tangle-network/agent-interface' import { immutableCandidateValue } from '../candidate-execution/digest' import { ConfigError } from '../errors' import { runCodeImprovement } from './code-execution' @@ -22,6 +22,7 @@ import type { ImproveResult, } from './improve-types' import { runMethodImprovement } from './method-execution' +import type { ReadonlyAgentProfile } from './profile-types' import { type ImproveTrainingOptions, type ImproveTrainingResult, @@ -83,14 +84,14 @@ export { createCommandProfileTrainer } from './training' /** Train and serve a checkpoint without implying that it improved held-out quality. */ export function improve( - profile: AgentProfile, + profile: ReadonlyAgentProfile, opts: ImproveTrainingOptions, ): Promise /** * Optimize one exact profile surface with a complete method. */ export function improve( - profile: AgentProfile, + profile: ReadonlyAgentProfile, opts: ImproveMethodOptions, ): Promise /** @@ -100,7 +101,7 @@ export function improve( opts: ImproveCodeRunOptions, ): Promise> export async function improve( - profileOrCode: AgentProfile | ImproveCodeRunOptions, + profileOrCode: ReadonlyAgentProfile | ImproveCodeRunOptions, opts?: ImproveMethodOptions | ImproveTrainingOptions, ): Promise | ImproveTrainingResult> { if (opts === undefined) { @@ -110,8 +111,12 @@ export async function improve( } return runCodeImprovement(code) } - if ('mode' in opts && opts.mode === 'training') { - return runProfileTraining(profileOrCode as AgentProfile, opts) + if (opts === null || typeof opts !== 'object') { + throw new ConfigError('improve(): options must be an object') + } + if ('mode' in opts) { + if (opts.mode !== 'training') throw new ConfigError('improve(): unsupported mode') + return runProfileTraining(profileOrCode as ReadonlyAgentProfile, opts) } if ((opts as { surface?: string }).surface === 'code') { throw new ConfigError("improve(): code takes one argument: improve({ surface: 'code', ... })") @@ -122,8 +127,5 @@ export async function improve( `improve(): input is not a valid AgentProfile: ${parsedProfile.error.message}`, ) } - return runMethodImprovement( - immutableCandidateValue(parsedProfile.data), - opts as ImproveMethodOptions, - ) + return runMethodImprovement(immutableCandidateValue(parsedProfile.data), opts) } diff --git a/src/improvement/meta-harness.test.ts b/src/improvement/meta-harness.test.ts index 372649694..627032df0 100644 --- a/src/improvement/meta-harness.test.ts +++ b/src/improvement/meta-harness.test.ts @@ -275,4 +275,29 @@ describe('createProfileImprovementHarness', () => { }) expect(() => harness.run({ validateCandidate: null } as never)).toThrow(ConfigError) }) + it('rejects invalid validator overrides consistently on training and optimization', () => { + const harness = createProfileImprovementHarness({ + profile: baselineProfile(), + executionRef: canonicalCandidateDigest({ fixture: 'validator-overrides' }), + agent: paidProfile, + validateCandidate: () => {}, + }) + expect(() => harness.run({ validateCandidate: null } as never)).toThrow(ConfigError) + expect(() => harness.train({ validateCandidate: null } as never)).toThrow(ConfigError) + }) + + it('refuses an async validator before optimization or agent execution', async () => { + let executions = 0 + const harness = createProfileImprovementHarness({ + profile: baselineProfile(), + executionRef: canonicalCandidateDigest({ fixture: 'async-validator' }), + agent: async (...args: Parameters) => { + executions++ + return paidProfile(...args) + }, + validateCandidate: async () => {}, + }) + await expect(harness.run(runOptions())).rejects.toThrow(/synchronous/) + expect(executions).toBe(0) + }) }) diff --git a/src/improvement/method-execution.ts b/src/improvement/method-execution.ts index a459e83bf..7c3da67d6 100644 --- a/src/improvement/method-execution.ts +++ b/src/improvement/method-execution.ts @@ -24,6 +24,7 @@ import { } from '@tangle-network/agent-interface' import { canonicalCandidateDigest, immutableCandidateValue } from '../candidate-execution/digest' import { ConfigError } from '../errors' +import { assertCandidateValidator, validateProfileCandidate } from './candidate-validation' import { copyImproveCost } from './improve-result' import type { ImproveCandidateValidationInput, @@ -309,6 +310,7 @@ export async function runMethodImprovement { /** Exact baseline profile. It is parsed, detached, and frozen at construction. */ - profile: AgentProfile + profile: ReadonlyAgentProfile /** * Immutable identity of the bound executor, models, tools, component mapping, * and every closure or external setting that can change measured behavior. @@ -88,23 +88,23 @@ export function createProfileImprovementHarness { + assertCandidateValidator(validator) + return validator ?? defaultValidator + } return Object.freeze({ profile, profileDigest: canonicalAgentProfileDigest(profile), executionRef, train(trainOptions: ProfileImprovementHarnessTrainOptions) { - const validateCandidate = trainOptions.validateCandidate ?? defaultValidator + const validateCandidate = validatorFor(trainOptions.validateCandidate) return improve(profile, { ...trainOptions, mode: 'training', @@ -113,15 +113,7 @@ export function createProfileImprovementHarness) { - if ( - runOptions.validateCandidate !== undefined && - typeof runOptions.validateCandidate !== 'function' - ) { - throw new ConfigError( - 'ProfileImprovementHarness.run: validateCandidate must be a function when present', - ) - } - const validateCandidate = runOptions.validateCandidate ?? defaultValidator + const validateCandidate = validatorFor(runOptions.validateCandidate) return improve(profile, { ...runOptions, executionRef, diff --git a/src/improvement/training.ts b/src/improvement/training.ts index 7127e4be6..12c2700ac 100644 --- a/src/improvement/training.ts +++ b/src/improvement/training.ts @@ -2,14 +2,14 @@ import { spawn } from 'node:child_process' import { createHash } from 'node:crypto' import { constants } from 'node:fs' import { chmod, mkdir, mkdtemp, open, realpath, rm, writeFile } from 'node:fs/promises' -import { isAbsolute, join } from 'node:path' +import { dirname, isAbsolute, join } from 'node:path' import { type AgentProfile, - type AgentProfileTraining, type AgentTrainingDatasetIdentity, type AgentTrainingReceipt, type AgentTrainingTask, agentProfileEnvironmentSchema, + agentProfileTrainingSchema, agentTrainingDatasetIdentitySchema, agentTrainingParametersSchema, agentTrainingReceiptSchema, @@ -27,11 +27,14 @@ import { immutableCandidateValue, sha256Bytes, } from '../candidate-execution/digest' -import { runAbortable } from '../runtime/supervise/abortable' +import { terminateProcessTreeAndConfirm } from '../runtime/process-tree' +import { linkAbort, runAbortable } from '../runtime/supervise/abortable' +import { armDeadlineTimer } from '../runtime/supervise/deadline' import { publishExclusiveDurableFile, syncDurableDirectory, } from '../runtime/supervise/durable-file' +import { assertCandidateValidator, validateProfileCandidate } from './candidate-validation' import type { ImproveCandidateValidator } from './improve-types' import type { ReadonlyAgentProfile } from './profile-types' @@ -92,7 +95,8 @@ export interface ImproveTrainingOptions { executionRef: Sha256Digest serving: CheckpointServingPort outputDirectory: string - timeoutMs: number + /** Optional overall deadline. Omit to rely on caller cancellation; long durations are supported. */ + timeoutMs?: number maxCheckpointBytes: number signal?: AbortSignal validateCandidate?: ImproveCandidateValidator @@ -138,12 +142,12 @@ export interface ControlledTrainingCommand { inputs: Array<{ path: string; digest: Sha256Digest }> /** Explicit public environment only. Ambient credentials are never inherited. */ environment: Record + /** Total stdout + stderr byte budget. Output is drained, not retained in memory. */ maxOutputBytes: number } -function positiveLimit(value: number, maximum: number, label: string): void { - if (!Number.isSafeInteger(value) || value <= 0 || value > maximum) - throw new Error(`invalid ${label}`) +function positiveLimit(value: number, label: string): void { + if (!Number.isSafeInteger(value) || value <= 0) throw new Error(`invalid ${label}`) } async function hashFile( @@ -160,14 +164,15 @@ async function hashFile( const file = await open(path, constants.O_RDONLY | constants.O_NOFOLLOW | constants.O_NONBLOCK) try { const before = await file.stat() - if (!before.isFile() || before.size <= 0 || before.size > maximum) - throw new Error('artifact must be a bounded nonempty regular file') + if (!before.isFile() || !Number.isSafeInteger(before.size) || before.size > maximum) + throw new Error('artifact must be a bounded regular file') const hash = createHash('sha256') const chunks: Buffer[] = [] let bytes = 0 for await (const chunk of file.createReadStream({ autoClose: false, signal })) { bytes += chunk.length if (bytes > maximum) throw new Error('artifact exceeded its byte limit') + if (bytes > before.size) throw new Error('artifact changed while being hashed') hash.update(chunk) if (capture) chunks.push(Buffer.from(chunk)) } @@ -193,7 +198,7 @@ async function hashFile( /** Execute one pinned command without a shell, in the runtime-owned job directory. POSIX only. */ export function createCommandProfileTrainer(input: ControlledTrainingCommand): ProfileTrainer { const command = immutableCandidateValue(input) - positiveLimit(command.maxOutputBytes, 16 * 1024 * 1024, 'trainer output limit') + positiveLimit(command.maxOutputBytes, 'trainer output limit') if ( !command.id || command.id.trim() !== command.id || @@ -229,9 +234,15 @@ export function createCommandProfileTrainer(input: ControlledTrainingCommand): P try { if (process.platform === 'win32') throw new Error('controlled trainers require POSIX process-group cancellation') + // Serialize before spawning and before any await: caller mutation cannot redirect the job, + // and a malformed request cannot leave a live child behind an early rejected promise. + const requestSnapshot = immutableCandidateValue(request) + const inputBytes = Buffer.from(canonicalCandidateBytes(requestSnapshot)) const verifyInputs = async () => { for (const file of [command.executable, ...command.inputs]) { - if ((await hashFile(file.path, 1024 * 1024 * 1024, signal)).digest !== file.digest) { + if ( + (await hashFile(file.path, Number.MAX_SAFE_INTEGER, signal)).digest !== file.digest + ) { throw new Error('trainer executable or input digest mismatch') } } @@ -240,7 +251,7 @@ export function createCommandProfileTrainer(input: ControlledTrainingCommand): P signal.throwIfAborted() await new Promise((resolve, reject) => { const child = spawn(command.executable.path, command.args, { - cwd: join(request.checkpointPath, '..'), + cwd: dirname(requestSnapshot.checkpointPath), env: command.environment, shell: false, detached: true, @@ -248,13 +259,21 @@ export function createCommandProfileTrainer(input: ControlledTrainingCommand): P }) let failure: Error | undefined let outputBytes = 0 + let leaderClosed = false + let termination: Promise | undefined const stop = () => { - if (!child.pid) return - try { - process.kill(-child.pid, 'SIGKILL') - } catch (error) { - if ((error as NodeJS.ErrnoException).code !== 'ESRCH') failure ??= error as Error + if (!termination) { + termination = terminateProcessTreeAndConfirm( + child, + () => leaderClosed, + 'createCommandProfileTrainer', + ) + // Teardown can finish before close. Observe its failure in both event orders. + void termination.catch((error) => { + failure ??= error + }) } + return termination } const abort = () => { failure ??= new Error('trainer cancelled') @@ -281,14 +300,16 @@ export function createCommandProfileTrainer(input: ControlledTrainingCommand): P // A successful parent may not leave descendants mutating the checkpoint. child.on('exit', stop) child.on('close', (code, exitSignal) => { + leaderClosed = true signal.removeEventListener('abort', abort) - stop() - if (failure) reject(failure) - else if (code !== 0 || exitSignal !== null) - reject(new Error(`trainer exited unsuccessfully (${code ?? exitSignal})`)) - else resolve() + void stop().then(() => { + if (failure) reject(failure) + else if (code !== 0 || exitSignal !== null) + reject(new Error(`trainer exited unsuccessfully (${code ?? exitSignal})`)) + else resolve() + }, reject) }) - child.stdin.end(Buffer.from(canonicalCandidateBytes(request))) + child.stdin.end(inputBytes) }) await verifyInputs() return { succeeded: true, value: undefined } @@ -352,19 +373,18 @@ function datasetIdentity(bytes: Uint8Array): AgentTrainingDatasetIdentity { /** Training materializes a candidate; it never emits a ship verdict or changes a live agent. */ export async function runProfileTraining( - profile: AgentProfile, + profile: ReadonlyAgentProfile, options: ImproveTrainingOptions, ): Promise { let stage: Extract['stage'] = 'admission' let jobDirectory: string | undefined - const controller = new AbortController() - const inputSignal = options.signal + const cancellation = linkAbort(...(options.signal ? [options.signal] : [])) + const signal = cancellation.signal const outputDirectory = options.outputDirectory const timeoutMs = options.timeoutMs let servingMayExist = false let trainingMayExist = false - const abort = () => controller.abort(inputSignal?.reason ?? new Error('training cancelled')) - let timer: ReturnType | undefined + let clearTimer: (() => void) | undefined try { if (options.mode !== 'training') throw new Error('training mode is required') const parent = snapshotAgentProfile(profile) @@ -375,8 +395,8 @@ export async function runProfileTraining( const parentBytes = Buffer.from(JSON.stringify(parent), 'utf8') sha256DigestSchema.parse(options.executionRef) sha256DigestSchema.parse(options.dataset.digest) - positiveLimit(options.timeoutMs, 7 * 24 * 60 * 60 * 1000, 'training timeout') - positiveLimit(options.maxCheckpointBytes, Number.MAX_SAFE_INTEGER, 'checkpoint byte limit') + if (timeoutMs !== undefined) positiveLimit(timeoutMs, 'training timeout') + positiveLimit(options.maxCheckpointBytes, 'checkpoint byte limit') if ( typeof options.trainer?.execute !== 'function' || typeof options.serving?.serve !== 'function' @@ -387,8 +407,7 @@ export async function runProfileTraining( agentTrainingParametersSchema.parse(options.parameters), ) agentTrainingReceiptSchema.shape.trainer.parse({ ...trainer, parameters }) - if (options.validateCandidate !== undefined && typeof options.validateCandidate !== 'function') - throw new Error('invalid candidate validator') + assertCandidateValidator(options.validateCandidate) const execute = options.trainer.execute.bind(options.trainer) const serve = options.serving.serve.bind(options.serving) const executionRef = options.executionRef @@ -396,29 +415,29 @@ export async function runProfileTraining( const sourceDataset = options.dataset.path const maxCheckpointBytes = options.maxCheckpointBytes const validateCandidate = options.validateCandidate - const ancestry: AgentProfileTraining['ancestors'] = parent.metadata?.training - ? [parent.metadata.training.receipt, ...parent.metadata.training.ancestors] - : [] - if (ancestry.length > 8) throw new Error('training receipt ancestry limit exceeded') - const validate = (candidate: AgentProfile, isBaseline: boolean) => { - const result: unknown = validateCandidate?.({ + // Interface owns both the ancestry bound and the portable receipt shape. + const ancestry = agentProfileTrainingSchema.shape.ancestors.parse( + parent.metadata?.training + ? [parent.metadata.training.receipt, ...parent.metadata.training.ancestors] + : [], + ) + const validate = (candidate: AgentProfile, isBaseline: boolean) => + validateProfileCandidate(validateCandidate, { profile: candidate, surface: 'agent-profile', candidateSurface: JSON.stringify(candidate), value: candidate, isBaseline, }) - if (result !== undefined) { - void Promise.resolve(result).catch(() => {}) - throw new Error('candidate validators must return void synchronously or throw') - } + if (timeoutMs !== undefined) { + clearTimer = armDeadlineTimer( + timeoutMs, + () => cancellation.abort(new Error('training deadline exceeded')), + true, + ) } - validate(parent, true) - inputSignal?.addEventListener('abort', abort, { once: true }) - if (inputSignal?.aborted) abort() - timer = setTimeout(() => controller.abort(new Error('training deadline exceeded')), timeoutMs) - const signal = controller.signal signal.throwIfAborted() + validate(parent, true) await mkdir(outputDirectory, { recursive: true }) const outputRoot = await realpath(outputDirectory) jobDirectory = await mkdtemp(join(outputRoot, 'training-')) @@ -467,6 +486,7 @@ export async function runProfileTraining( throw new Error('trainer changed its pinned inputs') } const artifact = await hashFile(artifactPath, maxCheckpointBytes, signal) + if (artifact.bytes === 0) throw new Error('checkpoint must be nonempty') await chmod(artifactPath, 0o400) const checkpoint = await open(artifactPath, constants.O_RDONLY | constants.O_NOFOLLOW) try { @@ -494,8 +514,7 @@ export async function runProfileTraining( throw new Error(served?.reason ?? 'checkpoint serving is unverified') if (served.value.artifactDigest !== artifact.digest) throw new Error('Router serving evidence names a different checkpoint') - if ((await hashFile(artifactPath, maxCheckpointBytes, signal)).digest !== artifact.digest) - throw new Error('checkpoint changed during serving') + // Snapshot caller-owned evidence before awaiting another read. const receipt = immutableCandidateValue( agentTrainingReceiptSchema.parse({ version: 1, @@ -512,10 +531,30 @@ export async function runProfileTraining( }, }), ) + if ((await hashFile(artifactPath, maxCheckpointBytes, signal)).digest !== artifact.digest) + throw new Error('checkpoint changed during serving') stage = 'profile' + const previousModel = parent.metadata?.training?.receipt.checkpoint.routerModelId + const routerModelId = receipt.checkpoint.routerModelId + const rebindModels = (entries: Record) => + Object.fromEntries( + Object.entries(entries).map(([key, value]) => [ + key, + previousModel && value.model === previousModel + ? { ...value, model: routerModelId } + : value, + ]), + ) const candidate = snapshotAgentProfile({ ...parent, - model: { ...parent.model, default: receipt.checkpoint.routerModelId }, + model: { + ...parent.model, + default: routerModelId, + ...(previousModel && parent.model.small === previousModel ? { small: routerModelId } : {}), + }, + // Only references to the old trained artifact follow the new receipt; other models stay put. + ...(parent.subagents && { subagents: rebindModels(parent.subagents) }), + ...(parent.modes && { modes: rebindModels(parent.modes) }), metadata: { ...parent.metadata, training: { receipt, ancestors: ancestry } }, }) validate(candidate, false) @@ -560,14 +599,19 @@ export async function runProfileTraining( mode: 'training', succeeded: false, stage, - reason: error instanceof Error ? error.message : 'training failed', + reason: + error instanceof Error + ? error.message + : typeof error === 'string' + ? error + : 'training failed', servingMayExist, trainingMayExist, ...(cleanupError ? { cleanupError } : {}), ...(jobDirectory ? { outputDirectory: jobDirectory } : {}), } } finally { - if (timer) clearTimeout(timer) - inputSignal?.removeEventListener('abort', abort) + clearTimer?.() + cancellation.release() } } diff --git a/src/mcp/local-harness.ts b/src/mcp/local-harness.ts index 52e64958a..07f2ea351 100644 --- a/src/mcp/local-harness.ts +++ b/src/mcp/local-harness.ts @@ -38,6 +38,7 @@ import { basename, delimiter, dirname, isAbsolute, join, resolve, sep } from 'no import type { AgentProfile, HarnessType, ReasoningEffort } from '@tangle-network/agent-interface' import { ValidationError } from '../errors' import { parseCodexUsageRecord } from '../runtime/harness-usage' +import { processGroupExists, terminateProcessTreeAndConfirm } from '../runtime/process-tree' import { concreteProfileModel } from '../runtime/supervise/model-policy' import { codexSensitiveEnvironmentName, @@ -47,6 +48,7 @@ import { redactCodexHome, } from './codex-diagnostics' +export { terminateProcessTreeAndConfirm } from '../runtime/process-tree' export type { CodexExecutionFailureDiagnostic } from './codex-diagnostics' export { CodexExecutionDiagnosticError } from './codex-diagnostics' @@ -541,9 +543,6 @@ export interface LocalHarnessResult { } const DEFAULT_MAX_OUTPUT_BYTES = 64 * 1024 * 1024 -const processKillGraceMs = 250 -const processGroupExitConfirmMs = 1_000 -const processGroupExitPollMs = 10 class RollingByteCapture { private readonly chunks: Buffer[] = [] @@ -1683,69 +1682,6 @@ function sha256File(path: string): Promise { }) } -function signalProcessTree(child: ChildProcess, signal: NodeJS.Signals): void { - if (process.platform !== 'win32' && typeof child.pid === 'number') { - try { - process.kill(-child.pid, signal) - return - } catch (err) { - if ((err as NodeJS.ErrnoException).code === 'ESRCH') return - } - } - try { - child.kill(signal) - } catch { - // The process may have exited between the timer and signal delivery. - } -} - -export async function terminateProcessTreeAndConfirm( - child: ChildProcess, - leaderClosed: () => boolean, - context = 'runLocalHarness', -): Promise { - signalProcessTree(child, 'SIGTERM') - if (process.platform === 'win32' || typeof child.pid !== 'number') { - const graceDeadline = Date.now() + processKillGraceMs - while (!leaderClosed() && Date.now() < graceDeadline) { - await delayUntilNextProcessCheck(graceDeadline) - } - if (!leaderClosed()) signalProcessTree(child, 'SIGKILL') - return - } - const processGroupId = child.pid - const graceDeadline = Date.now() + processKillGraceMs - if (await waitForProcessGroupExit(processGroupId, graceDeadline)) return - - signalProcessTree(child, 'SIGKILL') - const killDeadline = Date.now() + processGroupExitConfirmMs - if (await waitForProcessGroupExit(processGroupId, killDeadline)) return - throw new Error(`${context}: process group ${processGroupId} survived SIGKILL`) -} - -async function waitForProcessGroupExit(processGroupId: number, deadline: number): Promise { - while (processGroupExists(processGroupId)) { - if (Date.now() >= deadline) return false - await delayUntilNextProcessCheck(deadline) - } - return true -} - -async function delayUntilNextProcessCheck(deadline: number): Promise { - const delayMs = Math.min(processGroupExitPollMs, Math.max(0, deadline - Date.now())) - if (delayMs > 0) await new Promise((resolve) => setTimeout(resolve, delayMs)) -} - -function processGroupExists(processGroupId: number): boolean { - try { - process.kill(-processGroupId, 0) - return true - } catch (err) { - if ((err as NodeJS.ErrnoException).code === 'ESRCH') return false - return true - } -} - /** * Parse and validate the one terminal usage event emitted by `codex exec --json`. * diff --git a/src/runtime/process-tree.ts b/src/runtime/process-tree.ts new file mode 100644 index 000000000..cc84069b6 --- /dev/null +++ b/src/runtime/process-tree.ts @@ -0,0 +1,69 @@ +import type { ChildProcess } from 'node:child_process' + +const processKillGraceMs = 250 +const processGroupExitConfirmMs = 1_000 +const processGroupExitPollMs = 10 + +function signalProcessTree(child: ChildProcess, signal: NodeJS.Signals): void { + if (process.platform !== 'win32' && typeof child.pid === 'number') { + try { + process.kill(-child.pid, signal) + return + } catch (err) { + if ((err as NodeJS.ErrnoException).code === 'ESRCH') return + } + } + try { + child.kill(signal) + } catch { + // The process may have exited between the timer and signal delivery. + } +} + +/** Terminate one owned process group and confirm its disappearance; detached sessions are outside it. */ +export async function terminateProcessTreeAndConfirm( + child: ChildProcess, + leaderClosed: () => boolean, + context = 'runLocalHarness', +): Promise { + signalProcessTree(child, 'SIGTERM') + if (process.platform === 'win32' || typeof child.pid !== 'number') { + const graceDeadline = Date.now() + processKillGraceMs + while (!leaderClosed() && Date.now() < graceDeadline) { + await delayUntilNextProcessCheck(graceDeadline) + } + if (!leaderClosed()) signalProcessTree(child, 'SIGKILL') + return + } + const processGroupId = child.pid + const graceDeadline = Date.now() + processKillGraceMs + if (await waitForProcessGroupExit(processGroupId, graceDeadline)) return + + signalProcessTree(child, 'SIGKILL') + const killDeadline = Date.now() + processGroupExitConfirmMs + if (await waitForProcessGroupExit(processGroupId, killDeadline)) return + throw new Error(`${context}: process group ${processGroupId} survived SIGKILL`) +} + +async function waitForProcessGroupExit(processGroupId: number, deadline: number): Promise { + while (processGroupExists(processGroupId)) { + if (Date.now() >= deadline) return false + await delayUntilNextProcessCheck(deadline) + } + return true +} + +async function delayUntilNextProcessCheck(deadline: number): Promise { + const delayMs = Math.min(processGroupExitPollMs, Math.max(0, deadline - Date.now())) + if (delayMs > 0) await new Promise((resolve) => setTimeout(resolve, delayMs)) +} + +export function processGroupExists(processGroupId: number): boolean { + try { + process.kill(-processGroupId, 0) + return true + } catch (err) { + if ((err as NodeJS.ErrnoException).code === 'ESRCH') return false + return true + } +} diff --git a/src/testing/fixtures/agent-improvement-proposal.json b/src/testing/fixtures/agent-improvement-proposal.json index 5727a5faf..486c72d97 100644 --- a/src/testing/fixtures/agent-improvement-proposal.json +++ b/src/testing/fixtures/agent-improvement-proposal.json @@ -1,6 +1,6 @@ { "changedSurfaces": ["prompt"], - "digest": "sha256:224843e066e78674681eb5334f0cf291d30037c92727a5bf6d5a219109922f3e", + "digest": "sha256:510321d21135976eadd7168da0bf49b896ae72f66bab6776b414ccb0a03e7379", "evaluation": { "decision": { "contributingChecks": [ @@ -4882,7 +4882,7 @@ ], "metadata": { "fixture": "agent-improvement-proposal", - "runtimeVersion": "0.242.0" + "runtimeVersion": "0.243.0" }, "objectives": [ { @@ -4993,8 +4993,8 @@ "baselineContentHash": "sha256:5c21ee53e513fc604cb09754e21c392b24a424da0ef37dbf8f1ee4a8a0b08f09", "candidateContentHash": "sha256:60fcbb1c728194bd51d7d19cb732d1c3f1881dce7e0a6266b41c8b98cfd65693", "kind": "agent-eval-loop", - "recordDigest": "sha256:91b24d43247fc2688b49735c25914a44a3e4ce2421337433837796cea8f01542", - "runId": "agent-runtime-0.242.0-proposal-fixture", + "recordDigest": "sha256:dfb35e9df09682957107869a1452590e0d8873d8082144649d49e6fe050e9bba", + "runId": "agent-runtime-0.243.0-proposal-fixture", "schema": "agent-candidate-experiment" } }, @@ -5021,5 +5021,5 @@ ], "kind": "agent-improvement-proposal", "proposedAt": "2026-07-10T01:00:00.000Z", - "runId": "agent-runtime-0.242.0-proposal-fixture" + "runId": "agent-runtime-0.243.0-proposal-fixture" } diff --git a/src/testing/fixtures/agent-profile-improvement-proposal.json b/src/testing/fixtures/agent-profile-improvement-proposal.json index 8dd9fd415..5d9be3354 100644 --- a/src/testing/fixtures/agent-profile-improvement-proposal.json +++ b/src/testing/fixtures/agent-profile-improvement-proposal.json @@ -1,6 +1,6 @@ { "changedSurfaces": ["prompt", "skills"], - "digest": "sha256:97be4c92577f9497ce4d2897353c43cd07f29e94126b2e928e50cc66fba33113", + "digest": "sha256:8e8d52f54702a5946c776cff15a834336b8a04850228eb73f849fdafced8ded7", "evaluation": { "decision": { "contributingChecks": [ @@ -1715,7 +1715,7 @@ ], "metadata": { "fixture": "agent-profile-improvement-proposal", - "runtimeVersion": "0.242.0" + "runtimeVersion": "0.243.0" }, "objectives": [ { @@ -1826,7 +1826,7 @@ "baselineContentHash": "sha256:21c495a37c418c10bde64fbaa188beddeed31f1f051ea60a6a6582a9ee0db704", "candidateContentHash": "sha256:103f77bc8481601eef1ad5fe6ba84a40dffabc3a44f421f8c8559121edab84e9", "kind": "agent-eval-loop", - "recordDigest": "sha256:951ee09c461d8007f03a42a3d1f58510e877a5e097843b8c4138c241af121cdb", + "recordDigest": "sha256:38bd277479b222145ef0068a8f8f62b17eb866b4afc5f70e0922155d314213d0", "runId": "profile-improvement-1", "schema": "agent-profile-improvement-experiment" } diff --git a/tests/profile-training-boundaries.test.ts b/tests/profile-training-boundaries.test.ts index bf25595c2..123ba1f31 100644 --- a/tests/profile-training-boundaries.test.ts +++ b/tests/profile-training-boundaries.test.ts @@ -14,6 +14,7 @@ import { } from '../src/improvement/training' const faults = vi.hoisted(() => ({ + afterOpen: undefined as ((path: string) => void) | undefined, afterWrite: undefined as ((path: string) => void) | undefined, afterSync: undefined as ((path: string) => void) | undefined, afterPublication: undefined as ((path: string) => void) | undefined, @@ -32,6 +33,7 @@ vi.mock('node:fs/promises', async (importOriginal) => { }, async open(...args: Parameters) { const handle = await actual.open(...args) + faults.afterOpen?.(String(args[0])) const sync = handle.sync.bind(handle) handle.sync = async () => { await sync() @@ -62,6 +64,8 @@ vi.mock('../src/runtime/supervise/durable-file', async (importOriginal) => { let root: string | undefined afterEach(async () => { + vi.useRealTimers() + faults.afterOpen = undefined faults.afterWrite = undefined faults.afterSync = undefined faults.afterPublication = undefined @@ -275,3 +279,148 @@ describe('training cancellation, checkpoint integrity and publication', () => { await assertNoProfile(result.outputDirectory) }) }) + +describe('training composition and caller controls', () => { + it('has no implicit deadline when timeoutMs is omitted', async () => { + const { profile, options } = await fixture() + const { timeoutMs: _timeout, ...withoutTimeout } = options + const result = await runProfileTraining(profile, withoutTimeout) + assert(result.succeeded, JSON.stringify(result)) + }) + + it('honors a deadline beyond native timer range without expiring early', async () => { + const { profile, options, execute } = await fixture() + const day = 24 * 60 * 60 * 1000 + let reportStarted: () => void = () => {} + const started = new Promise((resolve) => { + reportStarted = resolve + }) + execute.mockImplementation(() => { + reportStarted() + return new Promise(() => {}) + }) + vi.useFakeTimers({ toFake: ['Date', 'setTimeout', 'clearTimeout'] }) + let settled = false + const work = runProfileTraining(profile, { ...options, timeoutMs: 35 * day }) + void work.then(() => { + settled = true + }) + // Admission failure must surface instead of leaving the test waiting for a trainer. + await Promise.race([ + started, + work.then((result) => { + assert(result.succeeded, JSON.stringify(result)) + }), + ]) + await vi.advanceTimersByTimeAsync(35 * day - 1) + assert.equal(settled, false) + await vi.advanceTimersByTimeAsync(1) + const result = await work + assert(!result.succeeded) + assert.equal(result.stage, 'training') + assert.equal(result.trainingMayExist, true) + assert.match(result.reason, /deadline/) + assert.equal(vi.getTimerCount(), 0) + }) + + it('keeps all references to the retrained checkpoint aligned without changing other models', async () => { + const { profile, options } = await fixture() + const first = await runProfileTraining(profile, options) + assert(first.succeeded, JSON.stringify(first)) + const oldModel = first.profile.model!.default! + const nextParent: AgentProfile = { + ...first.profile, + model: { ...first.profile.model, small: oldModel }, + subagents: { + trained: { model: oldModel, description: 'trained specialist' }, + other: { model: 'other-base', description: 'independent specialist' }, + }, + modes: { trained: { model: oldModel }, other: { model: 'other-base' } }, + } + options.trainer.execute = async (request) => { + await writeFile(request.checkpointPath, 'new weights') + return { succeeded: true, value: undefined } + } + const result = await runProfileTraining(nextParent, options) + assert(result.succeeded, JSON.stringify(result)) + const model = result.profile.model!.default! + assert.notEqual(model, oldModel) + assert.equal(result.profile.model!.small, model) + assert.equal(result.profile.subagents!.trained!.model, model) + assert.equal(result.profile.modes!.trained!.model, model) + assert.equal(result.profile.subagents!.other!.model, 'other-base') + assert.equal(result.profile.modes!.other!.model, 'other-base') + assert.equal(nextParent.model!.small, oldModel) + assert.deepEqual(result.profile.metadata!.training!.ancestors, [first.receipt]) + }) + + it('snapshots serving evidence before awaiting more filesystem work', async () => { + const { profile, options, serve } = await fixture() + const original = sha256Utf8('original evidence') + serve.mockImplementation(async (input) => { + const value = { + artifactDigest: input.artifactDigest, + routerModelId: input.routerModelId, + evidenceDigest: original, + } + faults.afterOpen = (path) => { + if (path.endsWith('checkpoint.bin')) + value.evidenceDigest = sha256Utf8('changed after serving') + } + return { succeeded: true, value } + }) + const result = await runProfileTraining(profile, options) + assert(result.succeeded, JSON.stringify(result)) + assert.equal(result.receipt.checkpoint.servingDigest, original) + }) +}) + +describe('training admission controls', () => { + it.each([0, -1, Number.NaN, Number.POSITIVE_INFINITY, Number.MAX_SAFE_INTEGER])( + 'rejects an invalid or unrepresentable deadline before dispatch: %s', + async (timeoutMs) => { + const { profile, options, execute, serve } = await fixture() + const result = await runProfileTraining(profile, { ...options, timeoutMs }) + assert(!result.succeeded) + assert.equal(result.stage, 'admission') + assert.equal(execute.mock.calls.length, 0) + assert.equal(serve.mock.calls.length, 0) + assert.equal(result.outputDirectory, undefined) + }, + ) + + it('still cancels a managed job when no deadline is configured', async () => { + const { profile, options, controller, execute } = await fixture() + delete options.timeoutMs + execute.mockImplementation(() => { + controller.abort(new Error('caller stopped the job')) + return new Promise(() => {}) + }) + const result = await runProfileTraining(profile, options) + assert(!result.succeeded) + assert.equal(result.stage, 'training') + assert.equal(result.trainingMayExist, true) + assert.match(result.reason, /caller stopped the job/) + }) + + it.each([ + ['fulfilled promise', async () => {}], + [ + 'rejected promise', + async () => { + throw new Error('late rejection') + }, + ], + ['false instead of throwing', () => false], + ] as const)( + 'refuses a validator returning %s before trainer dispatch', + async (_name, validator) => { + const { profile, options, execute } = await fixture() + const result = await runProfileTraining(profile, { ...options, validateCandidate: validator }) + assert(!result.succeeded) + assert.equal(result.stage, 'admission') + assert.match(result.reason, /synchronous/) + assert.equal(execute.mock.calls.length, 0) + }, + ) +}) diff --git a/tests/profile-training.test.ts b/tests/profile-training.test.ts index e25f8b42a..a3ad945e6 100644 --- a/tests/profile-training.test.ts +++ b/tests/profile-training.test.ts @@ -67,6 +67,7 @@ const serve: CheckpointServingPort = { async function withFixture( run: (options: ImproveTrainingOptions, dir: string) => Promise, script = TRAIN, + maxOutputBytes = 4096, ): Promise { const dir = await mkdtemp(join(tmpdir(), 'profile-training-test-')) try { @@ -102,7 +103,7 @@ async function withFixture( args: [scriptPath], inputs: [{ path: scriptPath, digest: sha256Utf8(script) }], environment: {}, - maxOutputBytes: 4096, + maxOutputBytes, }) await run( { @@ -355,7 +356,7 @@ process.stdin.on('end', () => { await withFixture(async (options) => { const first = await improve(parent(), options) assert(first.succeeded) - const second = await improve(first.profile as AgentProfile, { + const second = await improve(first.profile, { ...options, parameters: { epochs: 30, learningRate: 0.1 }, }) @@ -393,4 +394,79 @@ process.stdin.on('end', () => { await assertNoProfile(options.outputDirectory) }) }) + it('snapshots a direct command request before asynchronous input verification', async () => { + await withFixture(async (options, dir) => { + const originalPath = join(dir, 'original-checkpoint') + const request = { + version: 1 as const, + invocationId: 'direct-call', + datasetPath: options.dataset.path, + checkpointPath: originalPath, + parentProfilePath: options.dataset.path, + parentProfileDigest: canonicalAgentProfileDigest(parent()), + parameters: options.parameters, + executionRef: options.executionRef, + } + const pending = options.trainer.execute(request, new AbortController().signal) + request.checkpointPath = join(dir, 'redirected-checkpoint') + const result = await pending + assert(result.succeeded, JSON.stringify(result)) + assert((await readFile(originalPath)).length > 0) + await assert.rejects(readFile(request.checkpointPath), { code: 'ENOENT' }) + }) + }) + + it('honors a caller output bound larger than sixteen MiB', async () => { + await withFixture( + async (options) => { + const result = await improve(parent(), options) + assert(result.succeeded, JSON.stringify(result)) + }, + TRAIN.replace( + 'let weight = 0;', + "process.stdout.write('x'.repeat(17 * 1024 * 1024)); let weight = 0;", + ), + 18 * 1024 * 1024, + ) + }) + + it('allows a byte-pinned empty configuration input while still requiring a nonempty checkpoint', async () => { + await withFixture(async (options, dir) => { + const executable = await realpath(process.execPath) + const config = join(dir, 'empty-config') + await writeFile(config, '') + const trainer = createCommandProfileTrainer({ + id: 'empty-config-trainer', + executable: { path: executable, digest: sha256Bytes(await readFile(executable)) }, + args: [join(dir, 'trainer.cjs')], + inputs: [ + { path: join(dir, 'trainer.cjs'), digest: sha256Utf8(TRAIN) }, + { path: config, digest: sha256Utf8('') }, + ], + environment: {}, + maxOutputBytes: 4096, + }) + const result = await improve(parent(), { ...options, trainer }) + assert(result.succeeded, JSON.stringify(result)) + }) + }) + it('waits for a descendant to flush during shared process-group cleanup', async () => { + await withFixture( + async (options) => { + const result = await improve(parent(), options) + assert(result.succeeded, JSON.stringify(result)) + assert.equal(await readFile(result.artifactPath, 'utf8'), 'final child checkpoint') + }, + ` +const { spawn } = require('node:child_process'); +let input = ''; process.stdin.on('data', b => input += b); +process.stdin.on('end', () => { + const r = JSON.parse(input); + const code = "process.on('SIGTERM', () => { require('node:fs').writeFileSync(process.argv[1], 'final child checkpoint'); process.exit(0); }); process.send('ready'); setInterval(() => {}, 1000);"; + const child = spawn(process.execPath, ['-e', code, r.checkpointPath], { stdio: ['ignore', 'ignore', 'ignore', 'ipc'] }); + child.on('message', () => { child.disconnect(); process.exit(0); }); +}); +`, + ) + }) })