From 8788f0a2a7a7c1bac822c55840953007a4b90c56 Mon Sep 17 00:00:00 2001 From: Drew Stone Date: Fri, 17 Jul 2026 18:42:20 -0600 Subject: [PATCH] feat(improvement): bind activations to knowledge changes --- README.md | 25 +- package.json | 6 +- pnpm-lock.yaml | 22 +- src/file-transaction.ts | 18 +- src/kb-improvement.ts | 697 ++++++++++++++++++++++++++++++----- src/mutation-lock.ts | 126 ++++++- tests/kb-improvement.test.ts | 516 +++++++++++++++++++++++++- tests/mutation-lock.test.ts | 9 +- 8 files changed, 1280 insertions(+), 139 deletions(-) diff --git a/README.md b/README.md index 2f31388..5a51a47 100644 --- a/README.md +++ b/README.md @@ -308,7 +308,8 @@ import { improveKnowledgeBase, knowledgeImprovementCandidateRef, promoteKnowledgeCandidate, - withKnowledgeImprovementCandidate, + restoreKnowledgeCandidateBaseline, + withKnowledgeImprovementComparison, } from '@tangle-network/agent-knowledge' const staged = await improveKnowledgeBase({ @@ -330,19 +331,31 @@ const staged = await improveKnowledgeBase({ const candidate = knowledgeImprovementCandidateRef(staged) console.log(staged.evaluation, candidate) -await withKnowledgeImprovementCandidate({ root: './kb', candidate }, async (snapshot) => { - await inspectCandidateFiles(snapshot.root, snapshot.evaluation) +await withKnowledgeImprovementComparison({ root: './kb', candidate }, async (comparison) => { + await compareKnowledgeFiles( + comparison.baseline.root, + comparison.candidate.root, + comparison.evaluation, + ) }) // Call this only after your product records approval for this exact candidate. const promoted = await promoteKnowledgeCandidate({ root: './kb', candidate }) -console.log(promoted.promoted) +console.log(promoted.promoted, promoted.mutation) + +// The same candidate reference can restore its exact frozen baseline. +await restoreKnowledgeCandidateBaseline({ root: './kb', candidate }) ``` `improveKnowledgeBase` stages a measured candidate by default and does not change the live knowledge base. -Calling it again with the same `runId` resumes interrupted work. -`withKnowledgeImprovementCandidate` materializes the measured bytes in an isolated temporary directory for the callback, checks them again afterward, and removes the directory. +Calling it again with the same `runId` resumes interrupted candidate generation. +Resume an interrupted promotion or restore through the same transition function with its exact candidate and activation. +`withKnowledgeImprovementComparison` materializes the exact measured baseline and candidate in isolated temporary directories for one trusted callback, checks both again afterward, and removes them. +`withKnowledgeImprovementCandidate` remains the focused candidate-only read. `promoteKnowledgeCandidate` applies only the frozen bytes identified by the approved candidate reference, and refuses if the live base changed. +Promotion and restore results require `mutation` with the logical before/after hashes and transaction identity observed under the knowledge write lock; resumed transactions report `recovered: true`, while a later call that finds no pending work reports `changed: false`. +Passing `activation` makes the shared `AgentImprovementActivationResult` durable before the file transaction closes; `loadKnowledgeImprovementActivationResult` provides the read-only retry path. +Runtime supplies the result builder and product-owned result store, while this package owns knowledge files and their co-located result. The current release accepts one strict run-state format; incomplete runs created before 3.0 must be completed or restarted before upgrading. The exact candidate workflow requires Linux; other knowledge, retrieval, and evaluation APIs remain cross-platform. diff --git a/package.json b/package.json index 280523a..82b9972 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@tangle-network/agent-knowledge", - "version": "3.1.0", + "version": "3.2.0", "description": "Source-grounded, eval-gated knowledge growth primitives for agents.", "homepage": "https://github.com/tangle-network/agent-knowledge#readme", "repository": { @@ -69,8 +69,8 @@ "verify:package": "node scripts/verify-package.mjs" }, "dependencies": { - "@tangle-network/agent-eval": "^0.122.1", - "@tangle-network/agent-interface": "^0.30.0", + "@tangle-network/agent-eval": "^0.122.7", + "@tangle-network/agent-interface": "^0.31.0", "proper-lockfile": "4.1.2", "zod": "^4.3.6" }, diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index c9cf349..4a0c904 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -13,11 +13,11 @@ importers: .: dependencies: '@tangle-network/agent-eval': - specifier: ^0.122.1 - version: 0.122.1(typescript@5.9.3) + specifier: ^0.122.7 + version: 0.122.7(typescript@5.9.3) '@tangle-network/agent-interface': - specifier: ^0.30.0 - version: 0.30.0 + specifier: ^0.31.0 + version: 0.31.0 proper-lockfile: specifier: 4.1.2 version: 4.1.2 @@ -633,8 +633,8 @@ packages: '@tangle-network/agent-core@0.4.11': resolution: {integrity: sha512-5B1IjrJ8xDR7w8Hv/MSk2ixul6NEJQ5Ftzo+z6l7ipeFpc0yaf1+Kkml4HDUEmGW/rqdrNpBf9/MpT7i/SY0PA==} - '@tangle-network/agent-eval@0.122.1': - resolution: {integrity: sha512-nBjoslpjtxDxBTeBKgxgSHb3CJH+nKmXX8u5X8fepB9OllmaqhnF4B/DoNp95b/KaezVnd4i2hT2Zip/WS2Kxg==} + '@tangle-network/agent-eval@0.122.7': + resolution: {integrity: sha512-uOYVynC/+BnlNqZsKRG8NXz/wqeGB/AUImOPb/juUDQMIDzbkmxw9CiZGDmWH7kjCX7Nb8Z5zfpzO3VCTpO74Q==} engines: {node: '>=20'} hasBin: true @@ -647,8 +647,8 @@ packages: '@tangle-network/agent-interface@0.26.0': resolution: {integrity: sha512-/z4HavFr/9AbaHxi/13bFP9iSUt8oitZbw4wS9g9KAt2TeN2ec5/A2C106udWkKIETZLnu465TfU+K4KDVNkKw==} - '@tangle-network/agent-interface@0.30.0': - resolution: {integrity: sha512-GmwPwamrzOFPn+Qx7W0nt9QRUF9VUiZrMNLmF6QpL23R+7TD21HllPblVOYB5MVqA728FpDLXit3HfqY8mStwg==} + '@tangle-network/agent-interface@0.31.0': + resolution: {integrity: sha512-OvP8OebhbFd4d/Mxt1QDPAckJdBIA9omxoXU17dBbyEwhZKsiM6rbuOPaPDcv0Sf4koNLI25CfOYuzqlDR6ypQ==} '@tangle-network/sandbox@0.9.7': resolution: {integrity: sha512-9pCwJ5MlF7RUpp0AQKQDFyR0yu+E0udEhWkqhrlb/RuoJxlt72zVPuzO4FnMb1MZTkfjStmomC3k5xQyqi1YSA==} @@ -1528,13 +1528,13 @@ snapshots: '@tangle-network/agent-interface': 0.26.0 zod: 4.4.3 - '@tangle-network/agent-eval@0.122.1(typescript@5.9.3)': + '@tangle-network/agent-eval@0.122.7(typescript@5.9.3)': dependencies: '@asteasolutions/zod-to-openapi': 8.5.0(zod@4.4.3) '@ax-llm/ax': 23.0.0(zod@4.4.3) '@hono/node-server': 2.0.1(hono@4.12.30) '@tangle-network/agent-core': 0.4.11 - '@tangle-network/agent-interface': 0.30.0 + '@tangle-network/agent-interface': 0.31.0 '@tangle-network/tcloud': 0.4.14(typescript@5.9.3)(zod@4.4.3) hono: 4.12.30 zod: 4.4.3 @@ -1560,7 +1560,7 @@ snapshots: '@noble/hashes': 1.8.0 zod: 4.4.3 - '@tangle-network/agent-interface@0.30.0': + '@tangle-network/agent-interface@0.31.0': dependencies: '@noble/hashes': 1.8.0 zod: 4.4.3 diff --git a/src/file-transaction.ts b/src/file-transaction.ts index f1bb67a..fdb1b52 100644 --- a/src/file-transaction.ts +++ b/src/file-transaction.ts @@ -34,6 +34,7 @@ const transactionSchema = z kind: z.literal('knowledge-file-transaction'), transactionId: z.string().uuid(), purpose: z.string().min(1), + recoveryOwner: z.string().min(1).max(256).optional(), createdAt: z.string().min(1), entries: z.array(transactionEntrySchema).min(1), }) @@ -82,6 +83,7 @@ export async function prepareKnowledgeFileTransaction(input: { root: string transactionRoot: string purpose: string + recoveryOwner?: string mutations: readonly KnowledgeFileMutation[] includeUnchanged?: boolean now?: () => Date @@ -159,6 +161,7 @@ export async function prepareKnowledgeFileTransaction(input: { kind: 'knowledge-file-transaction', transactionId: randomUUID(), purpose: input.purpose, + ...(input.recoveryOwner ? { recoveryOwner: input.recoveryOwner } : {}), createdAt: (input.now ?? (() => new Date()))().toISOString(), entries: changed.map((item) => item.entry), }) @@ -227,6 +230,7 @@ export async function recoverKnowledgeFileTransaction(input: { transactionRoot: string expectedPurpose: string direction?: 'apply' | 'rollback' + finish?: boolean validate?: (transaction: KnowledgeFileTransaction) => void assertOwned?: () => void }): Promise { @@ -260,12 +264,14 @@ export async function recoverKnowledgeFileTransaction(input: { beforeCommit: input.assertOwned, }) } - await finishKnowledgeFileTransaction({ - root: input.root, - transactionRoot: input.transactionRoot, - transaction, - assertOwned: input.assertOwned, - }) + if (input.finish !== false) { + await finishKnowledgeFileTransaction({ + root: input.root, + transactionRoot: input.transactionRoot, + transaction, + assertOwned: input.assertOwned, + }) + } return true } diff --git a/src/kb-improvement.ts b/src/kb-improvement.ts index 2592d29..2d64453 100644 --- a/src/kb-improvement.ts +++ b/src/kb-improvement.ts @@ -8,6 +8,16 @@ import { type RunRecord, validateRunRecord, } from '@tangle-network/agent-eval' +import { + type AgentImprovementActivation, + type AgentImprovementActivationResult, + agentImprovementActivationResultSchema, + agentImprovementActivationSchema, + canonicalCandidateDigest, + omitTopLevelDigest, + type Sha256Digest, + sha256DigestSchema, +} from '@tangle-network/agent-interface' import { z } from 'zod' import { isMissingFile, @@ -154,6 +164,33 @@ export interface KnowledgeImprovementResult { blocked: boolean } +export type KnowledgeImprovementTarget = 'candidate' | 'baseline' + +export interface KnowledgeImprovementMutationReceipt { + target: KnowledgeImprovementTarget + beforeHash: string + afterHash: string + changed: boolean + transactionId: string | null + recovered: boolean +} + +export interface KnowledgeImprovementMutationResult extends KnowledgeImprovementResult { + candidate: KnowledgeImprovementCandidateRecord + mutation: KnowledgeImprovementMutationReceipt + activationResult?: AgentImprovementActivationResult +} + +export interface KnowledgeImprovementActivationPersistence { + activation: AgentImprovementActivation + attemptedAt: string + identity: string + /** May run again after interruption; keep this deterministic and free of external side effects. */ + createResult( + mutation: KnowledgeImprovementMutationReceipt, + ): Promise | AgentImprovementActivationResult +} + const digestSchema = z.string().regex(/^[a-f0-9]{64}$/) const runIdSchema = z.string().min(1).max(2_048) const safePathSegmentSchema = z @@ -161,6 +198,29 @@ const safePathSegmentSchema = z .min(1) .max(128) .regex(/^[A-Za-z0-9][A-Za-z0-9._-]*$/) + +const knowledgeImprovementMutationReceiptSchema = z + .object({ + target: z.enum(['candidate', 'baseline']), + beforeHash: digestSchema, + afterHash: digestSchema, + changed: z.boolean(), + transactionId: z.string().uuid().nullable(), + recovered: z.boolean(), + }) + .strict() + +const knowledgeImprovementActivationRecordSchema = z + .object({ + kind: z.literal('knowledge-improvement-activation-result'), + candidateId: safePathSegmentSchema, + mutation: knowledgeImprovementMutationReceiptSchema, + result: agentImprovementActivationResultSchema, + }) + .strict() +type KnowledgeImprovementActivationRecord = z.infer< + typeof knowledgeImprovementActivationRecordSchema +> const improvementStatusSchema = z.enum([ 'running', 'candidate-ready', @@ -388,6 +448,7 @@ export type KnowledgeImprovementCandidateRef = z.infer< export interface PromoteKnowledgeCandidateOptions { root: string candidate: KnowledgeImprovementCandidateRef + activation?: KnowledgeImprovementActivationPersistence ownerId?: string leaseTtlMs?: number now?: () => Date @@ -396,11 +457,30 @@ export interface PromoteKnowledgeCandidateOptions { export type RestoreKnowledgeCandidateBaselineOptions = PromoteKnowledgeCandidateOptions +export interface LoadKnowledgeImprovementActivationResultOptions { + root: string + candidate: KnowledgeImprovementCandidateRef + activation: AgentImprovementActivation + identity: string +} + export interface UseKnowledgeImprovementCandidateOptions { root: string candidate: KnowledgeImprovementCandidateRef } +export interface ResolvedKnowledgeImprovementComparisonSnapshot { + root: string + hash: string +} + +export interface ResolvedKnowledgeImprovementComparison { + reference: KnowledgeImprovementCandidateRef + evaluation: KnowledgeImprovementMetric + baseline: ResolvedKnowledgeImprovementComparisonSnapshot + candidate: ResolvedKnowledgeImprovementComparisonSnapshot +} + export interface ResolvedKnowledgeImprovementCandidate { root: string candidate: KnowledgeImprovementCandidateRef @@ -531,6 +611,27 @@ export async function loadKnowledgeImprovementState( } } +/** Load the durable result for one exact activation without changing knowledge or run state. */ +export async function loadKnowledgeImprovementActivationResult( + options: LoadKnowledgeImprovementActivationResultOptions, +): Promise { + assertExactCandidatePlatform() + const candidate = Object.freeze(KnowledgeImprovementCandidateRefSchema.parse(options.candidate)) + const activation = verifyCanonicalKnowledgeActivation(options.activation) + const target = targetForKnowledgeActivation(activation) + assertKnowledgeActivationAuthority(activation, candidate, target, options.identity) + return withKnowledgeImprovementRun(options.root, candidate.runId, false, async (runDir) => { + const record = await loadKnowledgeActivationRecord( + runDir, + candidate, + activation, + target, + options.identity, + ) + return record?.result ?? null + }) +} + async function loadKnowledgeImprovementStateFromRun( root: string, runId: string, @@ -559,13 +660,47 @@ export function knowledgeImprovementCandidateRef( return candidateRefFor(result.runId, result.state, result.candidate) } -/** Use one measured snapshot while its directory identity remains open and stable. */ +/** Use both frozen sides of one measured comparison in isolated, integrity-checked copies. */ +export async function withKnowledgeImprovementComparison( + options: UseKnowledgeImprovementCandidateOptions, + use: (comparison: ResolvedKnowledgeImprovementComparison) => Promise | T, +): Promise { + assertExactCandidatePlatform() + const reference = Object.freeze(KnowledgeImprovementCandidateRefSchema.parse(options.candidate)) + return withKnowledgeImprovementRun(options.root, reference.runId, false, async (runDir) => { + const state = await loadKnowledgeImprovementStateFromRun(options.root, reference.runId, runDir) + return withMeasuredCandidateSnapshot(options.root, runDir, state, reference, (resolved) => + withBaselineSnapshot(runDir, reference.baseHash, (baselineRoot) => + withIsolatedKnowledgeCopy(baselineRoot, reference.baseHash, 'baseline', (baseline) => + withIsolatedKnowledgeCopy( + resolved.root, + reference.candidateHash, + 'candidate', + (candidate) => + use( + Object.freeze({ + reference, + evaluation: immutableJsonValue(structuredClone(resolved.evidence.evaluation)), + baseline: Object.freeze({ root: baseline, hash: reference.baseHash }), + candidate: Object.freeze({ root: candidate, hash: reference.candidateHash }), + }), + ), + ), + ), + ), + ) + }) +} + +/** Use the frozen candidate side of one measured comparison. */ export async function withKnowledgeImprovementCandidate( options: UseKnowledgeImprovementCandidateOptions, use: (candidate: ResolvedKnowledgeImprovementCandidate) => Promise | T, ): Promise { assertExactCandidatePlatform() - const candidateRef = KnowledgeImprovementCandidateRefSchema.parse(options.candidate) + const candidateRef = Object.freeze( + KnowledgeImprovementCandidateRefSchema.parse(options.candidate), + ) return withKnowledgeImprovementRun(options.root, candidateRef.runId, false, async (runDir) => { const state = await loadKnowledgeImprovementStateFromRun( options.root, @@ -573,12 +708,14 @@ export async function withKnowledgeImprovementCandidate( runDir, ) return withMeasuredCandidateSnapshot(options.root, runDir, state, candidateRef, (resolved) => - withIsolatedCandidateCopy(resolved.root, candidateRef.candidateHash, (root) => - use({ - root, - candidate: candidateRef, - evaluation: resolved.evidence.evaluation, - }), + withIsolatedKnowledgeCopy(resolved.root, candidateRef.candidateHash, 'candidate', (root) => + use( + Object.freeze({ + root, + candidate: candidateRef, + evaluation: immutableJsonValue(structuredClone(resolved.evidence.evaluation)), + }), + ), ), ) }) @@ -587,28 +724,29 @@ export async function withKnowledgeImprovementCandidate( /** Promote one previously measured candidate without rerunning research or evaluation. */ export async function promoteKnowledgeCandidate( options: PromoteKnowledgeCandidateOptions, -): Promise { +): Promise { return transitionKnowledgeCandidate(options, 'candidate') } /** Restore the frozen baseline paired with one previously measured candidate. */ export async function restoreKnowledgeCandidateBaseline( options: RestoreKnowledgeCandidateBaselineOptions, -): Promise { +): Promise { return transitionKnowledgeCandidate(options, 'baseline') } -type KnowledgeCandidateTarget = 'candidate' | 'baseline' - async function transitionKnowledgeCandidate( options: PromoteKnowledgeCandidateOptions, - target: KnowledgeCandidateTarget, -): Promise { + target: KnowledgeImprovementTarget, +): Promise { assertExactCandidatePlatform() const candidateRef = Object.freeze( KnowledgeImprovementCandidateRefSchema.parse(options.candidate), ) const now = options.now ?? (() => new Date()) + const activation = options.activation + ? resolveKnowledgeActivationPersistence(options.activation, candidateRef, target) + : undefined return withKnowledgeImprovementRun(options.root, candidateRef.runId, false, async (runDir) => { const lease = await acquireRunLease(runDir, { ownerId: options.ownerId ?? `pid-${process.pid}`, @@ -631,6 +769,7 @@ async function transitionKnowledgeCandidate( assertRunOwned: lease.assertOwned, now, onState: options.onState, + activation, }, target, ) @@ -702,24 +841,7 @@ async function improveKnowledgeBaseInRun( if (state.status === 'promoted' && !promotedCandidate) { throw new Error('promoted knowledge state has no promoted candidate') } - const resumablePromotion = state.candidates.find( - (candidate) => - (candidate.status === 'candidate-ready' || candidate.status === 'promoted') && - candidate.evidenceHash !== undefined && - candidate.promotionPlanHash !== undefined, - ) - const resumableCandidateRef = resumablePromotion - ? candidateRefFor(runId, state, resumablePromotion) - : undefined - await withKnowledgeMutation(options.root, () => undefined, { - resumeTransaction: resumableCandidateRef - ? { - purpose: knowledgeCandidateTransitionPurpose(resumableCandidateRef, 'candidate'), - validate: (transaction) => - assertCandidateTransitionTransaction(transaction, resumableCandidateRef, 'candidate'), - } - : undefined, - }) + await withKnowledgeMutation(options.root, () => undefined) await ensureBaselineSnapshot(runDir, options.root, state.baseHash) if (state.status === 'promoted') { @@ -845,6 +967,7 @@ interface KnowledgeCandidateTransitionInput { runDir: string state: KnowledgeImprovementRunState candidateRef: KnowledgeImprovementCandidateRef + activation?: ResolvedKnowledgeImprovementActivationPersistence leaseTtlMs: number assertRunOwned(): void now: () => Date @@ -852,16 +975,21 @@ interface KnowledgeCandidateTransitionInput { lifecycle?: RunRagKnowledgeImprovementLoopResult } +interface ResolvedKnowledgeImprovementActivationPersistence + extends KnowledgeImprovementActivationPersistence { + activation: AgentImprovementActivation +} + async function applyKnowledgeCandidateTarget( input: KnowledgeCandidateTransitionInput, - target: KnowledgeCandidateTarget, -): Promise { + target: KnowledgeImprovementTarget, +): Promise { const { candidateRef, runDir, state } = input assertStateIdentity(input.root, candidateRef, state) const candidate = state.candidates.find((entry) => entry.candidateId === candidateRef.candidateId) if ( !candidate || - canonicalJson(candidateRefFor(candidateRef.runId, state, candidate)) !== + canonicalJson(candidateIdentityFor(candidateRef.runId, state, candidate)) !== canonicalJson(candidateRef) ) { throw new Error('knowledge candidate approval does not match the measured candidate') @@ -870,13 +998,85 @@ async function applyKnowledgeCandidateTarget( const action = target === 'candidate' ? 'promotion' : 'restore' const desiredHash = target === 'candidate' ? candidateRef.candidateHash : candidateRef.baseHash const purpose = knowledgeCandidateTransitionPurpose(candidateRef, target) + const recoveryOwner = input.activation + ? `knowledge-improvement-activation:${input.activation.activation.digest}` + : 'knowledge-improvement-candidate-transition' return withKnowledgeMutation( input.root, async (mutationLock) => { - input.assertRunOwned() + const assertOwned = () => { + mutationLock.assertOwned() + input.assertRunOwned() + } + assertOwned() const transactionRoot = mutationLock.transactionRoot - let pending: KnowledgeFileTransaction | null = null + const recovery = + mutationLock.recovery?.purpose === purpose ? mutationLock.recovery : undefined + const existingActivation = input.activation + ? await loadKnowledgeActivationRecord( + runDir, + candidateRef, + input.activation.activation, + target, + input.activation.identity, + ) + : null + if (existingActivation) { + if (recovery?.direction === 'rollback') { + throw new Error('stored knowledge activation result conflicts with a pending rollback') + } + if (recovery) { + const recoveredHash = await hashKnowledgeBase(input.root) + if (recoveredHash !== existingActivation.mutation.afterHash) { + throw new Error('stored knowledge activation result does not match the recovered files') + } + await finishKnowledgeFileTransaction({ + root: input.root, + transactionRoot, + transaction: recovery.transaction, + assertOwned, + }) + } + if (existingActivation.result.outcome.status === 'conflict') { + const reason = candidateTransitionConflictReason( + target, + candidateRef, + existingActivation.mutation.afterHash, + ) + const blocked = await blockCandidateTransition( + input, + candidate, + target, + reason, + existingActivation.mutation.afterHash, + ) + return { ...blocked, activationResult: existingActivation.result } + } + return candidateTransitionResult( + input, + candidate, + target === 'candidate' && state.status === 'promoted', + false, + existingActivation.mutation, + existingActivation.result, + ) + } + if (recovery?.direction === 'rollback') { + await finishKnowledgeFileTransaction({ + root: input.root, + transactionRoot, + transaction: recovery.transaction, + assertOwned, + }) + } + const recovered = recovery?.direction === 'apply' ? recovery : undefined + let pending: KnowledgeFileTransaction | null = + input.activation && recovered ? recovered.transaction : null const currentHash = await hashKnowledgeBase(input.root) + let transactionId = recovered?.transactionId ?? null + if (recovered && currentHash !== desiredHash) { + throw new Error(`recovered knowledge ${action} does not match the approved target`) + } if (state.status === 'promoted' && state.promotedCandidateId !== candidate.candidateId) { throw new Error( `knowledge run already promoted '${state.promotedCandidateId ?? 'unknown'}'`, @@ -901,10 +1101,34 @@ async function applyKnowledgeCandidateTarget( target === 'candidate' ? `base changed before promotion: expected ${state.baseHash}, got ${currentHash}` : `knowledge changed before restore: expected ${candidateRef.candidateHash}, got ${currentHash}` - return await blockCandidateTransition(input, candidate, target, reason) + const mutation = Object.freeze({ + target, + beforeHash: currentHash, + afterHash: currentHash, + changed: false, + transactionId: null, + recovered: false, + }) satisfies KnowledgeImprovementMutationReceipt + const activationResult = input.activation + ? await persistKnowledgeActivationResult( + runDir, + candidateRef, + input.activation, + target, + mutation, + ) + : undefined + const blocked = await blockCandidateTransition( + input, + candidate, + target, + reason, + currentHash, + ) + return activationResult ? { ...blocked, activationResult } : blocked } - if (currentHash !== desiredHash) { + if (currentHash !== desiredHash && !pending) { pending = await withMeasuredCandidateSnapshot( input.root, runDir, @@ -920,6 +1144,7 @@ async function applyKnowledgeCandidateTarget( root: input.root, transactionRoot, purpose, + recoveryOwner, mutations: await knowledgePlanMutations(targetRoot, plan), includeUnchanged: true, now: input.now, @@ -929,6 +1154,7 @@ async function applyKnowledgeCandidateTarget( if (!pending) { throw new Error(`knowledge ${action} plan unexpectedly contained no file changes`) } + transactionId = pending.transactionId try { assertCandidateTransitionTransaction(pending, candidateRef, target) } catch (error) { @@ -937,19 +1163,13 @@ async function applyKnowledgeCandidateTarget( root: input.root, transactionRoot, transaction: pending, - beforeCommit() { - mutationLock.assertOwned() - input.assertRunOwned() - }, + beforeCommit: assertOwned, }) await finishKnowledgeFileTransaction({ root: input.root, transactionRoot, transaction: pending, - assertOwned() { - mutationLock.assertOwned() - input.assertRunOwned() - }, + assertOwned, }) } catch (cleanupError) { throw new AggregateError( @@ -966,23 +1186,19 @@ async function applyKnowledgeCandidateTarget( root: input.root, transactionRoot, transaction: pending, - beforeCommit() { - mutationLock.assertOwned() - input.assertRunOwned() - }, + beforeCommit: assertOwned, }) } - mutationLock.assertOwned() - input.assertRunOwned() + assertOwned() if ((await hashKnowledgeBase(input.root)) !== desiredHash) { throw new Error(`knowledge ${action} content does not match the approved target`) } await writeKnowledgeIndex(input.root) } catch (error) { if (!pending) throw error + if (recovered && input.activation) throw error try { - mutationLock.assertOwned() - input.assertRunOwned() + assertOwned() } catch (ownershipError) { throw new AggregateError( [error, ownershipError], @@ -994,19 +1210,13 @@ async function applyKnowledgeCandidateTarget( root: input.root, transactionRoot, transaction: pending, - beforeCommit() { - mutationLock.assertOwned() - input.assertRunOwned() - }, + beforeCommit: assertOwned, }) await finishKnowledgeFileTransaction({ root: input.root, transactionRoot, transaction: pending, - assertOwned() { - mutationLock.assertOwned() - input.assertRunOwned() - }, + assertOwned, }) await writeKnowledgeIndex(input.root) } catch (rollbackError) { @@ -1026,36 +1236,81 @@ async function applyKnowledgeCandidateTarget( state.updatedAt = input.now().toISOString() await saveState(runDir, state, input.onState) await ensureCandidateTransitionEvent(runDir, candidateRef, target) + let finalHash = await hashKnowledgeBase(input.root) + if (finalHash !== desiredHash) { + throw new Error(`knowledge ${action} changed before its result was returned`) + } + const mutation = Object.freeze({ + target, + beforeHash: transactionId ? sourceKnowledgeHash(candidateRef, target) : currentHash, + afterHash: finalHash, + changed: transactionId !== null, + transactionId, + recovered: recovered !== undefined, + }) satisfies KnowledgeImprovementMutationReceipt + const activationResult = input.activation + ? await persistKnowledgeActivationResult( + runDir, + candidateRef, + input.activation, + target, + mutation, + ) + : undefined + if ((await hashKnowledgeBase(input.root)) !== finalHash) { + throw new Error(`knowledge ${action} changed while its result was persisted`) + } if (pending) { await finishKnowledgeFileTransaction({ root: input.root, transactionRoot, transaction: pending, - assertOwned() { - mutationLock.assertOwned() - input.assertRunOwned() - }, + assertOwned, }) } - return candidateTransitionResult(input, candidate, target === 'candidate', false) + finalHash = await hashKnowledgeBase(input.root) + if (finalHash !== desiredHash) { + throw new Error(`knowledge ${action} changed before its result was returned`) + } + return candidateTransitionResult( + input, + candidate, + target === 'candidate', + false, + mutation, + activationResult, + ) }, { staleMs: input.leaseTtlMs, resumeTransaction: { purpose, + recoveryOwner, validate: (transaction) => assertCandidateTransitionTransaction(transaction, candidateRef, target), + deferFinish: input.activation !== undefined, }, }, ) } +function candidateTransitionConflictReason( + target: KnowledgeImprovementTarget, + candidate: KnowledgeImprovementCandidateRef, + currentHash: string, +): string { + return target === 'candidate' + ? `base changed before promotion: expected ${candidate.baseHash}, got ${currentHash}` + : `knowledge changed before restore: expected ${candidate.candidateHash}, got ${currentHash}` +} + async function blockCandidateTransition( input: KnowledgeCandidateTransitionInput, candidate: KnowledgeImprovementCandidateRecord, - target: KnowledgeCandidateTarget, + target: KnowledgeImprovementTarget, reason: string, -): Promise { + currentHash: string, +): Promise { if (input.state.status === 'promoted') { candidate.status = 'blocked' candidate.updatedAt = input.now().toISOString() @@ -1068,7 +1323,14 @@ async function blockCandidateTransition( candidateId: candidate.candidateId, reason, }) - return candidateTransitionResult(input, candidate, false, true) + return candidateTransitionResult(input, candidate, false, true, { + target, + beforeHash: currentHash, + afterHash: currentHash, + changed: false, + transactionId: null, + recovered: false, + }) } function candidateTransitionResult( @@ -1076,21 +1338,280 @@ function candidateTransitionResult( candidate: KnowledgeImprovementCandidateRecord, promoted: boolean, blocked: boolean, -): KnowledgeImprovementResult { + mutation: KnowledgeImprovementMutationReceipt, + activationResult?: AgentImprovementActivationResult, +): KnowledgeImprovementMutationResult { return { runId: input.candidateRef.runId, state: input.state, candidate, ...(input.lifecycle ? { lifecycle: input.lifecycle } : {}), + mutation, + ...(activationResult ? { activationResult } : {}), promoted, blocked, } } +function sourceKnowledgeHash( + candidate: KnowledgeImprovementCandidateRef, + target: KnowledgeImprovementTarget, +): string { + return target === 'candidate' ? candidate.baseHash : candidate.candidateHash +} + +function targetForKnowledgeActivation( + activation: AgentImprovementActivation, +): KnowledgeImprovementTarget { + return activation.intent === 'activate-candidate' ? 'candidate' : 'baseline' +} + +function resolveKnowledgeActivationPersistence( + input: KnowledgeImprovementActivationPersistence, + candidate: KnowledgeImprovementCandidateRef, + target: KnowledgeImprovementTarget, +): ResolvedKnowledgeImprovementActivationPersistence { + const activation = verifyCanonicalKnowledgeActivation(input.activation) + assertKnowledgeActivationAuthority(activation, candidate, target, input.identity) + const attemptedAt = z.iso.datetime().parse(input.attemptedAt) + if ( + Date.parse(attemptedAt) < Date.parse(activation.authorizedAt) || + Date.parse(attemptedAt) >= Date.parse(activation.expiresAt) + ) { + throw new Error('knowledge activation attempt is outside its authorization window') + } + return Object.freeze({ + activation, + attemptedAt, + identity: input.identity, + createResult: input.createResult, + }) +} + +function assertKnowledgeActivationAuthority( + activation: AgentImprovementActivation, + candidate: KnowledgeImprovementCandidateRef, + target: KnowledgeImprovementTarget, + identity: string, +): void { + if (!identity.trim()) throw new Error('knowledge activation identity is required') + if (targetForKnowledgeActivation(activation) !== target) { + throw new Error('knowledge activation intent does not match the requested transition') + } + if (activation.targets.length !== 1) { + throw new Error('knowledge activation requires exactly one target') + } + const authorizedTarget = activation.targets[0] + if (authorizedTarget.surface !== 'knowledge' || authorizedTarget.identity !== identity) { + throw new Error('knowledge activation target does not match this knowledge base') + } + if ( + authorizedTarget.expectedBaseDigest !== + prefixedKnowledgeDigest(sourceKnowledgeHash(candidate, target)) + ) { + throw new Error('knowledge activation does not authorize the measured source state') + } +} + +function verifyCanonicalKnowledgeActivation( + value: AgentImprovementActivation, +): AgentImprovementActivation { + const activation = agentImprovementActivationSchema.parse(value) + if (canonicalCandidateDigest(omitTopLevelDigest(activation)) !== activation.digest) { + throw new Error('knowledge activation digest does not match its canonical content') + } + return immutableJsonValue(structuredClone(activation)) +} + +function verifyCanonicalKnowledgeActivationResult( + value: AgentImprovementActivationResult, +): AgentImprovementActivationResult { + const result = agentImprovementActivationResultSchema.parse(value) + if (canonicalCandidateDigest(omitTopLevelDigest(result)) !== result.digest) { + throw new Error('knowledge activation result digest does not match its canonical content') + } + return immutableJsonValue(structuredClone(result)) +} + +async function persistKnowledgeActivationResult( + runDir: string, + candidate: KnowledgeImprovementCandidateRef, + persistence: ResolvedKnowledgeImprovementActivationPersistence, + target: KnowledgeImprovementTarget, + mutation: KnowledgeImprovementMutationReceipt, +): Promise { + const result = assertKnowledgeActivationResult( + persistence.activation, + candidate, + target, + persistence.identity, + mutation, + await persistence.createResult(Object.freeze({ ...mutation })), + persistence.attemptedAt, + ) + const record = knowledgeImprovementActivationRecordSchema.parse({ + kind: 'knowledge-improvement-activation-result', + candidateId: candidate.candidateId, + mutation, + result, + }) + const existing = await loadKnowledgeActivationRecord( + runDir, + candidate, + persistence.activation, + target, + persistence.identity, + ) + if (existing) { + if (canonicalJson(existing) !== canonicalJson(record)) { + throw new Error('knowledge activation result identity conflicts with durable content') + } + return existing.result + } + await writeJsonDurableWithinRoot( + runDir, + knowledgeActivationResultPath(persistence.activation.digest), + record, + ) + const stored = await loadKnowledgeActivationRecord( + runDir, + candidate, + persistence.activation, + target, + persistence.identity, + ) + if (!stored || canonicalJson(stored) !== canonicalJson(record)) { + throw new Error('knowledge activation result was not durably persisted') + } + return stored.result +} + +async function loadKnowledgeActivationRecord( + runDir: string, + candidate: KnowledgeImprovementCandidateRef, + activation: AgentImprovementActivation, + target: KnowledgeImprovementTarget, + identity: string, +): Promise { + let raw: unknown + try { + const file = await readRegularFileWithinRoot( + runDir, + knowledgeActivationResultPath(activation.digest), + ) + raw = JSON.parse(file.bytes.toString('utf8')) as unknown + } catch (error) { + if (isMissingFile(error)) return null + throw error + } + const record = knowledgeImprovementActivationRecordSchema.parse(raw) + if (record.candidateId !== candidate.candidateId) { + throw new Error('knowledge activation result belongs to another candidate') + } + const result = assertKnowledgeActivationResult( + activation, + candidate, + target, + identity, + record.mutation, + record.result, + ) + return immutableJsonValue({ ...record, result }) +} + +function assertKnowledgeActivationResult( + activation: AgentImprovementActivation, + candidate: KnowledgeImprovementCandidateRef, + target: KnowledgeImprovementTarget, + identity: string, + mutationInput: KnowledgeImprovementMutationReceipt, + resultInput: AgentImprovementActivationResult, + attemptedAt?: string, +): AgentImprovementActivationResult { + assertKnowledgeActivationAuthority(activation, candidate, target, identity) + const mutation = knowledgeImprovementMutationReceiptSchema.parse(mutationInput) + const result = verifyCanonicalKnowledgeActivationResult(resultInput) + if ( + result.idempotencyKey !== activation.digest || + (attemptedAt !== undefined && result.attemptedAt !== attemptedAt) || + Date.parse(result.attemptedAt) < Date.parse(activation.authorizedAt) || + Date.parse(result.attemptedAt) >= Date.parse(activation.expiresAt) || + mutation.target !== target || + mutation.changed !== (mutation.beforeHash !== mutation.afterHash) || + (!mutation.changed && (mutation.transactionId !== null || mutation.recovered)) + ) { + throw new Error('knowledge activation result does not bind its authorized mutation') + } + + const sourceDigest = prefixedKnowledgeDigest(sourceKnowledgeHash(candidate, target)) + const desiredDigest = prefixedKnowledgeDigest( + target === 'candidate' ? candidate.candidateHash : candidate.baseHash, + ) + const beforeDigest = prefixedKnowledgeDigest(mutation.beforeHash) + const afterDigest = prefixedKnowledgeDigest(mutation.afterHash) + const outcome = result.outcome + if (mutation.changed) { + if ( + mutation.transactionId === null || + beforeDigest !== sourceDigest || + afterDigest !== desiredDigest || + outcome.status !== 'applied' || + outcome.transactionId !== mutation.transactionId || + outcome.targets.length !== 1 || + outcome.targets[0]?.surface !== 'knowledge' || + outcome.targets[0]?.identity !== identity || + outcome.targets[0]?.beforeDigest !== beforeDigest || + outcome.targets[0]?.afterDigest !== afterDigest + ) { + throw new Error('knowledge activation result does not prove the applied transaction') + } + return result + } + + const expectedStatus = afterDigest === desiredDigest ? 'already-applied' : 'conflict' + if ( + afterDigest === sourceDigest || + outcome.status !== expectedStatus || + outcome.targets.length !== 1 || + outcome.targets[0]?.surface !== 'knowledge' || + outcome.targets[0]?.identity !== identity || + outcome.targets[0]?.currentDigest !== afterDigest + ) { + throw new Error('knowledge activation result does not prove the observed target state') + } + return result +} + +function knowledgeActivationResultPath(digest: Sha256Digest): string { + const parsed = sha256DigestSchema.parse(digest) + return `activation-results/${parsed.slice('sha256:'.length)}.json` +} + +function prefixedKnowledgeDigest(hash: string): Sha256Digest { + return sha256DigestSchema.parse(`sha256:${digestSchema.parse(hash)}`) +} + +function immutableJsonValue(value: T): T { + if (value === null || typeof value !== 'object') return value + for (const child of Object.values(value)) immutableJsonValue(child) + return Object.freeze(value) +} + function candidateRefFor( runId: string, state: KnowledgeImprovementRunState, candidate: KnowledgeImprovementCandidateRecord, +): KnowledgeImprovementCandidateRef { + if (candidate.status !== 'candidate-ready' && candidate.status !== 'promoted') { + throw new Error(`knowledge candidate '${candidate.candidateId}' is not ready`) + } + return candidateIdentityFor(runId, state, candidate) +} + +function candidateIdentityFor( + runId: string, + state: KnowledgeImprovementRunState, + candidate: KnowledgeImprovementCandidateRecord, ): KnowledgeImprovementCandidateRef { if (!candidate.candidateHash) { throw new Error(`knowledge candidate '${candidate.candidateId}' has no content hash`) @@ -1101,9 +1622,6 @@ function candidateRefFor( if (!candidate.promotionPlanHash) { throw new Error(`knowledge candidate '${candidate.candidateId}' has no promotion plan hash`) } - if (candidate.status !== 'candidate-ready' && candidate.status !== 'promoted') { - throw new Error(`knowledge candidate '${candidate.candidateId}' is not ready`) - } return Object.freeze({ kind: 'knowledge-improvement-candidate', runId, @@ -1155,21 +1673,22 @@ async function withMeasuredCandidateSnapshot( }) } -async function withIsolatedCandidateCopy( +async function withIsolatedKnowledgeCopy( sourceRoot: string, expectedHash: string, + target: KnowledgeImprovementTarget, use: (root: string) => Promise | T, ): Promise { - const isolationRoot = await mkdtemp(join(tmpdir(), 'agent-knowledge-candidate-')) - const candidateRoot = join(isolationRoot, 'candidate') + const isolationRoot = await mkdtemp(join(tmpdir(), 'agent-knowledge-snapshot-')) + const snapshotRoot = join(isolationRoot, 'snapshot') try { - await copyKnowledgeWorkspace(sourceRoot, candidateRoot) - if ((await hashKnowledgeBase(candidateRoot)) !== expectedHash) { - throw new Error('isolated knowledge candidate does not match its approved content') + await copyKnowledgeWorkspace(sourceRoot, snapshotRoot) + if ((await hashKnowledgeBase(snapshotRoot)) !== expectedHash) { + throw new Error(`isolated knowledge ${target} does not match its measured content`) } - const result = await use(candidateRoot) - if ((await hashKnowledgeBase(candidateRoot)) !== expectedHash) { - throw new Error('knowledge candidate snapshot changed during use') + const result = await use(snapshotRoot) + if ((await hashKnowledgeBase(snapshotRoot)) !== expectedHash) { + throw new Error(`knowledge ${target} snapshot changed during use`) } return result } finally { @@ -1812,7 +2331,7 @@ async function copyKnowledgeWorkspace(sourceRoot: string, targetRoot: string): P function knowledgeCandidateTransitionPurpose( candidate: KnowledgeImprovementCandidateRef, - target: KnowledgeCandidateTarget, + target: KnowledgeImprovementTarget, ): string { const action = target === 'candidate' ? 'promotion' : 'restore' return `knowledge-${action}:${contentHash(candidate)}` @@ -1865,7 +2384,7 @@ async function knowledgePlanMutations( function assertCandidateTransitionPlan( plan: readonly KnowledgeFileTransactionPlanEntry[], candidate: KnowledgeImprovementCandidateRef, - target: KnowledgeCandidateTarget, + target: KnowledgeImprovementTarget, ): void { const approvedDirection = target === 'candidate' ? plan : reverseKnowledgeFilePlan(plan) const actualPlanHash = knowledgeFileTransactionPlanHash(approvedDirection) @@ -1879,7 +2398,7 @@ function assertCandidateTransitionPlan( function assertCandidateTransitionTransaction( transaction: KnowledgeFileTransaction, candidate: KnowledgeImprovementCandidateRef, - target: KnowledgeCandidateTarget, + target: KnowledgeImprovementTarget, ): void { assertCandidateTransitionPlan(transaction.entries, candidate, target) } @@ -1899,7 +2418,7 @@ function reverseKnowledgeFilePlan( async function ensureCandidateTransitionEvent( runDir: string, candidateRef: KnowledgeImprovementCandidateRef, - target: KnowledgeCandidateTarget, + target: KnowledgeImprovementTarget, ): Promise { if (await hasCandidateTransitionEvent(runDir, candidateRef, target)) return await appendLedger(runDir, { @@ -1915,7 +2434,7 @@ async function ensureCandidateTransitionEvent( async function hasCandidateTransitionEvent( runDir: string, candidateRef: KnowledgeImprovementCandidateRef, - target: KnowledgeCandidateTarget, + target: KnowledgeImprovementTarget, ): Promise { const eventType = target === 'candidate' ? 'candidate.promoted' : 'candidate.restored' let matched = false diff --git a/src/mutation-lock.ts b/src/mutation-lock.ts index 115e037..59532a2 100644 --- a/src/mutation-lock.ts +++ b/src/mutation-lock.ts @@ -29,9 +29,17 @@ const activeReadRoots = new AsyncLocalStorage>() export interface KnowledgeMutationLock { readonly transactionRoot: string + readonly recovery?: KnowledgeMutationRecovery assertOwned(): void } +export interface KnowledgeMutationRecovery { + transactionId: string + purpose: string + direction: 'apply' | 'rollback' + transaction: KnowledgeFileTransaction +} + export class KnowledgeLockLostError extends Error { constructor(message: string, options?: { cause?: unknown }) { super(message, options) @@ -48,8 +56,10 @@ export interface KnowledgeMutationOptions { staleMs?: number resumeTransaction?: { purpose: string + recoveryOwner?: string validate?: (transaction: KnowledgeFileTransaction) => void direction?: 'apply' | 'rollback' + deferFinish?: boolean } retries?: LockOptions['retries'] } @@ -57,6 +67,7 @@ export interface KnowledgeMutationOptions { export interface PendingKnowledgeMutation { transactionId: string purpose: string + recoveryOwner?: string createdAt: string direction: 'apply' | 'rollback' paths: string[] @@ -96,30 +107,41 @@ export async function withKnowledgeMutation( options.retries ?? ({ retries: 100, factor: 1.1, minTimeout: 10, maxTimeout: 200, randomize: true } as const), }) + let recovery: KnowledgeMutationRecovery | undefined const mutationLock: KnowledgeMutationLock = { transactionRoot: join(cacheDir, 'file-transactions'), + get recovery() { + return recovery + }, assertOwned: acquired.assertOwned, } const scope: KnowledgeMutationScope = { active: true, lock: mutationLock } try { - const pending = await loadKnowledgeFileTransaction({ + const pendingState = await inspectKnowledgeFileTransaction({ root: resolvedRoot, transactionRoot: mutationLock.transactionRoot, }) - if (pending) { - const resume = options.resumeTransaction - if (!resume || pending.purpose !== resume.purpose) { - throw new Error(`knowledge transaction '${pending.purpose}' requires its owner to resume`) - } - } + const pending = pendingState?.transaction ?? null + const resume = pending + ? assertKnowledgeTransactionResume(pending, options.resumeTransaction) + : undefined const locks = new Map(active) locks.set(resolvedRoot, scope) return await activeRoots.run(locks, async () => { const epoch = await beginMutationEpoch(cacheDir) - let completed = false + const finishEpoch = async (completed: boolean) => { + mutationLock.assertOwned() + const stillPending = await loadKnowledgeFileTransaction({ + root: resolvedRoot, + transactionRoot: mutationLock.transactionRoot, + }) + if (completed && stillPending) { + throw new Error('knowledge mutation returned with an unfinished transaction') + } + if (!stillPending) await finishMutationEpoch(cacheDir, epoch) + } try { if (pending) { - const resume = options.resumeTransaction if (!resume) throw new Error( `knowledge transaction '${pending.purpose}' requires its owner to resume`, @@ -129,22 +151,32 @@ export async function withKnowledgeMutation( transactionRoot: mutationLock.transactionRoot, expectedPurpose: resume.purpose, direction: resume.direction, + finish: resume.deferFinish !== true, validate: resume.validate, assertOwned: acquired.assertOwned, }) + recovery = Object.freeze({ + transactionId: pending.transactionId, + purpose: pending.purpose, + direction: resume.direction ?? pendingState?.direction ?? 'apply', + transaction: pending, + }) } mutationLock.assertOwned() const result = await mutate(mutationLock) mutationLock.assertOwned() - completed = true + await finishEpoch(true) return result - } finally { - mutationLock.assertOwned() - const stillPending = await loadKnowledgeFileTransaction({ - root: resolvedRoot, - transactionRoot: mutationLock.transactionRoot, - }) - if (completed || !stillPending) await finishMutationEpoch(cacheDir, epoch) + } catch (error) { + try { + await finishEpoch(false) + } catch (finishError) { + throw new AggregateError( + [error, finishError], + 'knowledge mutation failed and its durable state could not be inspected', + ) + } + throw error } }) } finally { @@ -167,12 +199,31 @@ export async function inspectPendingKnowledgeMutation( return { transactionId: transaction.transactionId, purpose: transaction.purpose, + ...(transaction.recoveryOwner ? { recoveryOwner: transaction.recoveryOwner } : {}), createdAt: transaction.createdAt, direction, paths: transaction.entries.map((entry) => entry.path), } } +function assertKnowledgeTransactionResume( + pending: KnowledgeFileTransaction, + resume: KnowledgeMutationOptions['resumeTransaction'], +): NonNullable { + if (!resume || pending.purpose !== resume.purpose) { + throw new Error(`knowledge transaction '${pending.purpose}' requires its owner to resume`) + } + if (pending.recoveryOwner !== resume.recoveryOwner) { + if (pending.recoveryOwner) { + throw new Error( + `knowledge transaction '${pending.purpose}' must be resumed by '${pending.recoveryOwner}'`, + ) + } + throw new Error(`knowledge transaction '${pending.purpose}' has no recovery owner`) + } + return resume +} + export async function recoverPendingKnowledgeMutation( root: string, options: RecoverPendingKnowledgeMutationOptions, @@ -332,12 +383,47 @@ async function waitForActiveMutationEpoch( epoch: number, options: KnowledgeReadOptions, ): Promise { - if (!(await hasActiveMutationLock(root, options))) { - throw new Error(`knowledge mutation epoch ${epoch} is odd with no active writer`) - } + if (!(await hasActiveMutationLock(root, options))) + await healAbandonedMutationEpoch(root, epoch, options) await new Promise((resolve) => setTimeout(resolve, options.waitMs ?? DEFAULT_READ_WAIT_MS)) } +async function healAbandonedMutationEpoch( + root: string, + expectedEpoch: number, + options: KnowledgeReadOptions, +): Promise { + await withSafeDirectory(root, '.agent-knowledge', false, async (cacheDir) => { + let acquired: DurableFileLock + try { + acquired = await acquireDurableFileLock(root, { + lockfilePath: join(cacheDir, 'mutation.lock.durable'), + staleMs: options.staleMs, + retries: 0, + }) + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ELOCKED') return + throw error + } + try { + const currentEpoch = await readMutationEpochFromCache(cacheDir) + if (currentEpoch !== expectedEpoch || !isOdd(currentEpoch)) return + const pending = await inspectKnowledgeFileTransaction({ + root, + transactionRoot: join(cacheDir, 'file-transactions'), + }) + if (pending) { + throw new Error( + `knowledge transaction '${pending.transaction.purpose}' requires its owner to resume`, + ) + } + await finishMutationEpoch(cacheDir, currentEpoch) + } finally { + await acquired.release() + } + }) +} + async function hasActiveMutationLock( root: string, options: Pick, diff --git a/tests/kb-improvement.test.ts b/tests/kb-improvement.test.ts index c8e7846..dbb8cd4 100644 --- a/tests/kb-improvement.test.ts +++ b/tests/kb-improvement.test.ts @@ -13,6 +13,12 @@ import { } from 'node:fs/promises' import { tmpdir } from 'node:os' import { dirname, join, relative } from 'node:path' +import { + type AgentImprovementActivation, + type AgentImprovementActivationResult, + canonicalCandidateDigest, + type Sha256Digest, +} from '@tangle-network/agent-interface' import { describe, expect, it } from 'vitest' import { applyKnowledgeWriteBlocks, @@ -23,15 +29,21 @@ import { hashKnowledgeBase, improveKnowledgeBase, initKnowledgeBase, + inspectPendingKnowledgeMutation, + type KnowledgeImprovementCandidateRef, + type KnowledgeImprovementMutationReceipt, knowledgeImprovementCandidateRef, knowledgeImprovementRunDir, + loadKnowledgeImprovementActivationResult, loadKnowledgeImprovementEvents, loadKnowledgeImprovementState, promoteKnowledgeCandidate, + recoverPendingKnowledgeMutation, restoreKnowledgeCandidateBaseline, sha256, stableId, withKnowledgeImprovementCandidate, + withKnowledgeImprovementComparison, } from '../src/index' import { withKnowledgeMutation } from '../src/mutation-lock' @@ -116,6 +128,96 @@ function passingMetric() { } } +function candidateDigest(seed: string): Sha256Digest { + return canonicalCandidateDigest({ seed }) +} + +function canonicalDocument>( + material: T, +): T & { digest: Sha256Digest } { + return { ...material, digest: canonicalCandidateDigest(material) } +} + +function knowledgeActivation( + candidate: KnowledgeImprovementCandidateRef, + intent: AgentImprovementActivation['intent'], + identity = 'knowledge:test', +): AgentImprovementActivation { + const expectedBaseHash = + intent === 'activate-candidate' ? candidate.baseHash : candidate.candidateHash + return canonicalDocument({ + kind: 'agent-improvement-activation' as const, + proposalDigest: candidateDigest('proposal'), + reviewDigest: candidateDigest('review'), + experimentDigest: candidateDigest('experiment'), + candidateBundleDigest: candidateDigest('candidate-bundle'), + intent, + targets: [ + { + surface: 'knowledge' as const, + identity, + expectedBaseDigest: `sha256:${expectedBaseHash}` as Sha256Digest, + }, + ] as AgentImprovementActivation['targets'], + fundingOwner: 'tenant:test', + authorizedBy: 'reviewer:test', + authorizedAt: '2026-07-17T00:00:00.000Z', + expiresAt: '2026-07-18T00:00:00.000Z', + }) +} + +function knowledgeActivationResult( + activation: AgentImprovementActivation, + candidate: KnowledgeImprovementCandidateRef, + mutation: KnowledgeImprovementMutationReceipt, + attemptedAt: string, + identity = 'knowledge:test', +): AgentImprovementActivationResult { + const desiredHash = + activation.intent === 'activate-candidate' ? candidate.candidateHash : candidate.baseHash + const outcome: AgentImprovementActivationResult['outcome'] = mutation.changed + ? { + status: 'applied', + transactionId: mutation.transactionId!, + targets: [ + { + surface: 'knowledge', + identity, + beforeDigest: `sha256:${mutation.beforeHash}`, + afterDigest: `sha256:${mutation.afterHash}`, + }, + ], + } + : mutation.afterHash === desiredHash + ? { + status: 'already-applied', + targets: [ + { + surface: 'knowledge', + identity, + currentDigest: `sha256:${mutation.afterHash}`, + }, + ], + } + : { + status: 'conflict', + targets: [ + { + surface: 'knowledge', + identity, + currentDigest: `sha256:${mutation.afterHash}`, + }, + ], + } + return canonicalDocument({ + kind: 'agent-improvement-activation-result' as const, + idempotencyKey: activation.digest, + attemptedAt, + completedAt: attemptedAt, + outcome, + }) +} + async function improveAndPromote(options: Parameters[0]) { const staged = await improveKnowledgeBase(options) const promoted = await promoteKnowledgeCandidate({ @@ -400,6 +502,14 @@ describe('improveKnowledgeBase', () => { const repeated = await promoteKnowledgeCandidate({ root, candidate }) expect(repeated.promoted).toBe(true) expect(repeated.state.promotedCandidateId).toBe(candidate.candidateId) + expect(repeated.mutation).toEqual({ + target: 'candidate', + beforeHash: candidate.candidateHash, + afterHash: candidate.candidateHash, + changed: false, + transactionId: null, + recovered: false, + }) expect(updateCalls).toBe(1) }) }) @@ -425,10 +535,19 @@ describe('improveKnowledgeBase', () => { }) const candidate = knowledgeImprovementCandidateRef(staged) - await expect(promoteKnowledgeCandidate({ root, candidate })).resolves.toMatchObject({ + const promoted = await promoteKnowledgeCandidate({ root, candidate }) + expect(promoted).toMatchObject({ promoted: true, blocked: false, }) + expect(promoted.mutation).toEqual({ + target: 'candidate', + beforeHash: baseHash, + afterHash: candidate.candidateHash, + changed: true, + transactionId: expect.any(String), + recovered: false, + }) expect(await hashKnowledgeBase(root)).toBe(candidate.candidateHash) const restored = await restoreKnowledgeCandidateBaseline({ root, candidate }) @@ -438,6 +557,14 @@ describe('improveKnowledgeBase', () => { state: { status: 'candidate-ready' }, candidate: { status: 'candidate-ready' }, }) + expect(restored.mutation).toEqual({ + target: 'baseline', + beforeHash: candidate.candidateHash, + afterHash: baseHash, + changed: true, + transactionId: expect.any(String), + recovered: false, + }) expect(restored.state.promotedCandidateId).toBeUndefined() expect(await hashKnowledgeBase(root)).toBe(baseHash) await expect(readFile(originalPage, 'utf8')).resolves.toBe('# Original\n') @@ -489,6 +616,14 @@ describe('improveKnowledgeBase', () => { promoted: false, blocked: false, state: { status: 'candidate-ready' }, + mutation: { + target: 'baseline', + beforeHash: candidate.candidateHash, + afterHash: candidate.baseHash, + changed: true, + transactionId: expect.any(String), + recovered: true, + }, }) expect(await hashKnowledgeBase(root)).toBe(candidate.baseHash) const events = await loadKnowledgeImprovementEvents(root, candidate.runId) @@ -583,6 +718,14 @@ describe('improveKnowledgeBase', () => { const recovered = await promoteKnowledgeCandidate({ root, candidate }) expect(recovered.promoted).toBe(true) + expect(recovered.mutation).toEqual({ + target: 'candidate', + beforeHash: candidate.baseHash, + afterHash: candidate.candidateHash, + changed: true, + transactionId: expect.any(String), + recovered: true, + }) expect(updateCalls).toBe(1) const events = await loadKnowledgeImprovementEvents(root, candidate.runId) expect(events.filter((event) => event.type === 'candidate.promoted')).toHaveLength(1) @@ -591,6 +734,296 @@ describe('improveKnowledgeBase', () => { }) }) + it('persists an approved result before closing its file transaction and never reapplies it', async () => { + await withKb(async (root) => { + const staged = await improveKnowledgeBase({ + root, + goal: 'Keep approval results recoverable with their knowledge mutation', + runId: 'protected-activation-recovery', + updateKnowledge: async ({ candidateRoot }) => { + await writeFile(join(candidateRoot, 'knowledge', 'candidate.md'), '# Candidate\n') + return { applied: true, summary: 'created measured knowledge' } + }, + evaluate: passingMetric, + }) + const candidate = knowledgeImprovementCandidateRef(staged) + const activation = knowledgeActivation(candidate, 'activate-candidate') + const attemptedAt = '2026-07-17T12:00:00.000Z' + const transactionRoot = join(root, '.agent-knowledge', 'file-transactions') + let interruptedMutation: KnowledgeImprovementMutationReceipt | undefined + + await expect( + promoteKnowledgeCandidate({ + root, + candidate, + activation: { + activation, + attemptedAt, + identity: 'knowledge:test', + createResult(mutation) { + interruptedMutation = mutation + throw new Error('activation result store unavailable') + }, + }, + }), + ).rejects.toThrow(/activation result store unavailable/) + expect(interruptedMutation).toMatchObject({ + target: 'candidate', + beforeHash: candidate.baseHash, + afterHash: candidate.candidateHash, + changed: true, + transactionId: expect.any(String), + recovered: false, + }) + await expect(hashKnowledgeBase(root)).rejects.toThrow(/requires its owner to resume/) + await expect(readFile(join(root, 'knowledge', 'candidate.md'), 'utf8')).resolves.toBe( + '# Candidate\n', + ) + expect( + (await readdir(transactionRoot)).filter((entry) => entry.startsWith('active-')), + ).toHaveLength(1) + const pending = await inspectPendingKnowledgeMutation(root) + expect(pending).toMatchObject({ + transactionId: interruptedMutation?.transactionId, + recoveryOwner: `knowledge-improvement-activation:${activation.digest}`, + }) + await expect( + recoverPendingKnowledgeMutation(root, { + transactionId: pending!.transactionId, + action: 'apply', + }), + ).rejects.toThrow(/must be resumed by 'knowledge-improvement-activation:sha256:/) + await expect(promoteKnowledgeCandidate({ root, candidate })).rejects.toThrow( + /must be resumed by 'knowledge-improvement-activation:sha256:/, + ) + await expect( + loadKnowledgeImprovementActivationResult({ + root, + candidate, + activation, + identity: 'knowledge:test', + }), + ).resolves.toBeNull() + await expect( + improveKnowledgeBase({ + root, + goal: 'Keep approval results recoverable with their knowledge mutation', + runId: 'protected-activation-recovery', + }), + ).rejects.toThrow(/requires its owner to resume/) + + const recovered = await promoteKnowledgeCandidate({ + root, + candidate, + activation: { + activation, + attemptedAt, + identity: 'knowledge:test', + createResult: (mutation) => + knowledgeActivationResult(activation, candidate, mutation, attemptedAt), + }, + }) + expect(recovered.mutation).toEqual({ + ...interruptedMutation, + recovered: true, + }) + expect(recovered.activationResult?.outcome.status).toBe('applied') + expect( + (await readdir(transactionRoot)).filter((entry) => entry.startsWith('active-')), + ).toEqual([]) + await expect( + loadKnowledgeImprovementActivationResult({ + root, + candidate, + activation, + identity: 'knowledge:test', + }), + ).resolves.toEqual(recovered.activationResult) + + await restoreKnowledgeCandidateBaseline({ root, candidate }) + await expect(hashKnowledgeBase(root)).resolves.toBe(candidate.baseHash) + const retried = await promoteKnowledgeCandidate({ + root, + candidate, + activation: { + activation, + attemptedAt, + identity: 'knowledge:test', + createResult() { + throw new Error('durable activation must not be recomputed') + }, + }, + }) + expect(retried.activationResult).toEqual(recovered.activationResult) + await expect(hashKnowledgeBase(root)).resolves.toBe(candidate.baseHash) + }) + }) + + it('repairs blocked run state and its event from a durable conflict result', async () => { + await withKb(async (root) => { + const staged = await improveKnowledgeBase({ + root, + goal: 'Keep conflict results and run state consistent', + runId: 'activation-conflict-recovery', + updateKnowledge: async ({ candidateRoot }) => { + await writeFile(join(candidateRoot, 'knowledge', 'candidate.md'), '# Candidate\n') + return { applied: true, summary: 'created measured knowledge' } + }, + evaluate: passingMetric, + }) + const candidate = knowledgeImprovementCandidateRef(staged) + await promoteKnowledgeCandidate({ root, candidate }) + await writeFile(join(root, 'knowledge', 'candidate.md'), '# Concurrent change\n') + const activation = knowledgeActivation(candidate, 'restore-baseline') + const attemptedAt = '2026-07-17T12:00:00.000Z' + + await expect( + restoreKnowledgeCandidateBaseline({ + root, + candidate, + activation: { + activation, + attemptedAt, + identity: 'knowledge:test', + createResult: (mutation) => + knowledgeActivationResult(activation, candidate, mutation, attemptedAt), + }, + onState() { + throw new Error('operator stopped before conflict event persistence') + }, + }), + ).rejects.toThrow(/operator stopped before conflict event persistence/) + const stored = await loadKnowledgeImprovementActivationResult({ + root, + candidate, + activation, + identity: 'knowledge:test', + }) + expect(stored?.outcome.status).toBe('conflict') + expect( + (await loadKnowledgeImprovementEvents(root, candidate.runId)).filter( + (event) => event.type === 'restore.blocked', + ), + ).toHaveLength(0) + + const recovered = await restoreKnowledgeCandidateBaseline({ + root, + candidate, + activation: { + activation, + attemptedAt, + identity: 'knowledge:test', + createResult() { + throw new Error('durable conflict result must not be recomputed') + }, + }, + }) + expect(recovered).toMatchObject({ + promoted: false, + blocked: true, + state: { status: 'blocked' }, + candidate: { status: 'blocked' }, + }) + expect(recovered.activationResult).toEqual(stored) + await restoreKnowledgeCandidateBaseline({ + root, + candidate, + activation: { + activation, + attemptedAt, + identity: 'knowledge:test', + createResult() { + throw new Error('durable conflict result must not be recomputed') + }, + }, + }) + expect( + (await loadKnowledgeImprovementEvents(root, candidate.runId)).filter( + (event) => event.type === 'restore.blocked', + ), + ).toHaveLength(1) + }) + }) + + it('binds already-applied promotion and baseline restore to their exact activation', async () => { + await withKb(async (root) => { + const staged = await improveKnowledgeBase({ + root, + goal: 'Record exact no-op and restore activation outcomes', + runId: 'activation-already-and-restore', + updateKnowledge: async ({ candidateRoot }) => { + await writeFile(join(candidateRoot, 'knowledge', 'candidate.md'), '# Candidate\n') + return { applied: true, summary: 'created measured knowledge' } + }, + evaluate: passingMetric, + }) + const candidate = knowledgeImprovementCandidateRef(staged) + await promoteKnowledgeCandidate({ root, candidate }) + const attemptedAt = '2026-07-17T12:00:00.000Z' + const promoteActivation = knowledgeActivation(candidate, 'activate-candidate') + + const already = await promoteKnowledgeCandidate({ + root, + candidate, + activation: { + activation: promoteActivation, + attemptedAt, + identity: 'knowledge:test', + createResult: (mutation) => + knowledgeActivationResult(promoteActivation, candidate, mutation, attemptedAt), + }, + }) + expect(already.mutation).toEqual({ + target: 'candidate', + beforeHash: candidate.candidateHash, + afterHash: candidate.candidateHash, + changed: false, + transactionId: null, + recovered: false, + }) + expect(already.activationResult?.outcome.status).toBe('already-applied') + await expect( + loadKnowledgeImprovementActivationResult({ + root, + candidate, + activation: promoteActivation, + identity: 'knowledge:test', + }), + ).resolves.toEqual(already.activationResult) + + const restoreActivation = knowledgeActivation(candidate, 'restore-baseline') + const restored = await restoreKnowledgeCandidateBaseline({ + root, + candidate, + activation: { + activation: restoreActivation, + attemptedAt, + identity: 'knowledge:test', + createResult: (mutation) => + knowledgeActivationResult(restoreActivation, candidate, mutation, attemptedAt), + }, + }) + expect(restored.mutation).toMatchObject({ + target: 'baseline', + beforeHash: candidate.candidateHash, + afterHash: candidate.baseHash, + changed: true, + transactionId: expect.any(String), + recovered: false, + }) + expect(restored.activationResult?.outcome.status).toBe('applied') + await expect( + loadKnowledgeImprovementActivationResult({ + root, + candidate, + activation: restoreActivation, + identity: 'knowledge:test', + }), + ).resolves.toEqual(restored.activationResult) + await expect(hashKnowledgeBase(root)).resolves.toBe(candidate.baseHash) + }) + }) + it('rejects a forged promotion journal entry outside the measured knowledge files', async () => { await withKb(async (root) => { const packagePath = join(root, 'package.json') @@ -719,6 +1152,57 @@ describe('improveKnowledgeBase', () => { }) }) + it('resolves the exact frozen baseline and candidate from one measured comparison', async () => { + await withKb(async (root) => { + const liveBefore = await hashKnowledgeBase(root) + const staged = await improveKnowledgeBase({ + root, + goal: 'Compare the frozen baseline and candidate bytes', + runId: 'paired-frozen-snapshots', + updateKnowledge: async ({ candidateRoot }) => { + await writeFile(join(candidateRoot, 'knowledge', 'measured.md'), '# Measured\n') + return { applied: true, summary: 'created measured knowledge' } + }, + evaluate: () => ({ ...passingMetric(), dimensions: { quality: 1 } }), + }) + const candidate = knowledgeImprovementCandidateRef(staged) + const isolatedRoots: string[] = [] + + await withKnowledgeImprovementComparison({ root, candidate }, async (comparison) => { + isolatedRoots.push(comparison.baseline.root, comparison.candidate.root) + expect(Object.isFrozen(comparison)).toBe(true) + expect(Object.isFrozen(comparison.baseline)).toBe(true) + expect(Object.isFrozen(comparison.candidate)).toBe(true) + expect(Object.isFrozen(comparison.evaluation)).toBe(true) + expect(Object.isFrozen(comparison.evaluation.provenance)).toBe(true) + expect(Object.isFrozen(comparison.evaluation.dimensions)).toBe(true) + expect(comparison.reference).toEqual(candidate) + expect(comparison.baseline.hash).toBe(candidate.baseHash) + expect(comparison.candidate.hash).toBe(candidate.candidateHash) + await expect(hashKnowledgeBase(comparison.baseline.root)).resolves.toBe(candidate.baseHash) + await expect(hashKnowledgeBase(comparison.candidate.root)).resolves.toBe( + candidate.candidateHash, + ) + await expect( + readFile(join(comparison.baseline.root, 'knowledge', 'measured.md'), 'utf8'), + ).rejects.toMatchObject({ code: 'ENOENT' }) + await expect( + readFile(join(comparison.candidate.root, 'knowledge', 'measured.md'), 'utf8'), + ).resolves.toBe('# Measured\n') + }) + await expect( + withKnowledgeImprovementComparison({ root, candidate }, ({ baseline }) => + writeFile(join(baseline.root, 'knowledge', 'changed.md'), '# Changed\n'), + ), + ).rejects.toThrow(/baseline snapshot changed during use/) + + await expect(hashKnowledgeBase(root)).resolves.toBe(liveBefore) + for (const isolatedRoot of isolatedRoots) { + await expect(stat(isolatedRoot)).rejects.toMatchObject({ code: 'ENOENT' }) + } + }) + }) + it('rejects a measured snapshot changed while an approved candidate is in use', async () => { await withKb(async (root) => { const staged = await improveKnowledgeBase({ @@ -741,6 +1225,36 @@ describe('improveKnowledgeBase', () => { }) }) + it('does not return a successful mutation receipt when a state callback changes knowledge', async () => { + await withKb(async (root) => { + const staged = await improveKnowledgeBase({ + root, + goal: 'Reject post-commit knowledge changes', + runId: 'post-commit-change', + updateKnowledge: async ({ candidateRoot }) => { + await writeFile(join(candidateRoot, 'knowledge', 'candidate.md'), '# Candidate\n') + return { applied: true, summary: 'created measured knowledge' } + }, + evaluate: passingMetric, + }) + const candidate = knowledgeImprovementCandidateRef(staged) + + await expect( + promoteKnowledgeCandidate({ + root, + candidate, + async onState() { + await writeFile(join(root, 'knowledge', 'unmeasured.md'), '# Unmeasured\n') + }, + }), + ).rejects.toThrow(/changed before its result was returned/) + await expect(hashKnowledgeBase(root)).rejects.toThrow(/requires its owner to resume/) + await expect(readFile(join(root, 'knowledge', 'unmeasured.md'), 'utf8')).resolves.toBe( + '# Unmeasured\n', + ) + }) + }) + it('rejects a frozen measured copy changed after approval', async () => { await withKb(async (root) => { const source = refundSource() diff --git a/tests/mutation-lock.test.ts b/tests/mutation-lock.test.ts index 277846b..657192c 100644 --- a/tests/mutation-lock.test.ts +++ b/tests/mutation-lock.test.ts @@ -88,7 +88,7 @@ describe('knowledge read epochs', () => { }) }) - it('fails loudly on an abandoned odd mutation epoch', async () => { + it('repairs an abandoned odd mutation epoch when no transaction remains', async () => { await withRoot(async (root) => { await mkdir(join(root, '.agent-knowledge'), { recursive: true }) await writeFile( @@ -96,7 +96,10 @@ describe('knowledge read epochs', () => { '{\n "epoch": 1,\n "updatedAt": "2026-07-13T00:00:00.000Z"\n}\n', ) - await expect(buildKnowledgeIndex(root)).rejects.toThrow(/odd with no active writer/) + await expect(buildKnowledgeIndex(root)).resolves.toMatchObject({ pages: [] }) + await expect( + readFile(join(root, '.agent-knowledge', 'mutation-epoch.json'), 'utf8'), + ).resolves.toContain('"epoch": 2') }) }) @@ -128,7 +131,7 @@ describe('knowledge read epochs', () => { }), ).rejects.toThrow(/simulated interruption/) - await expect(buildKnowledgeIndex(root)).rejects.toThrow(/odd with no active writer/) + await expect(buildKnowledgeIndex(root)).rejects.toThrow(/requires its owner to resume/) }) })