fix: bind manual production graph and snapshots

This commit is contained in:
2026-08-10 00:38:46 +02:00
parent c10857405a
commit e97f335b88
4 changed files with 58 additions and 7 deletions
+16 -6
View File
@@ -67,7 +67,12 @@ const watchdog=setTimeout(()=>{if(state!=="STARTING")return;console.error("manua
const PRELOAD=`data:text/javascript;base64,${Buffer.from(PRELOAD_SOURCE,"utf8").toString("base64")}`;
async function controlRequest(control,payload){if(control?.host!==HOST||!Number.isSafeInteger(control?.port)||control.port<1||control.port>65535)throw new Error("backend control identity mismatch");return await new Promise((resolvePromise,reject)=>{const socket=net.createConnection({host:control.host,port:control.port}),timer=setTimeout(()=>socket.destroy(new Error("backend control timeout")),2000);let bytes="";socket.setEncoding("utf8");socket.on("connect",()=>socket.end(JSON.stringify(payload)));socket.on("data",chunk=>{bytes+=chunk;if(bytes.length>2048)socket.destroy(new Error("backend control response too large"));});socket.on("error",reject);socket.on("close",()=>{clearTimeout(timer);let value;try{value=JSON.parse(bytes);}catch{return reject(new Error("backend control response is malformed"));}resolvePromise(value);});});}
function ownedValue(repo,root,nonce,{backendLog=null,entrypoint,stage="PREPARING",createdAt=new Date().toISOString()}={}){return{schemaVersion:1,kind:"p1-manual-acceptance",nonce,repositoryRoot:repo,root,status:"PENDING",stage,createdAt,listener:{host:HOST,port:PORT,state:"stopped"},backendLog,entrypoint,resources:[root,{kind:"fastify",host:HOST,port:PORT}]};}
async function distManifest(repo){
const base=join(repo,"backend","dist"), files=[];
async function walk(dir){ for(const entry of await readdir(dir,{withFileTypes:true})){ const path=join(dir,entry.name); if(entry.isSymbolicLink())throw new Error("production distribution contains a symlink"); if(entry.isDirectory())await walk(path); else if(entry.isFile()){ const bytes=await readFile(path); if(bytes.length>33554432)throw new Error("production distribution file is unbounded"); files.push({path:relative(base,path).split(sep).join("/"),size:bytes.length,sha256:createHash("sha256").update(bytes).digest("hex")}); } else throw new Error("production distribution entry is unsupported"); if(files.length>20000)throw new Error("production distribution is unbounded"); }}
await walk(base); files.sort((a,b)=>a.path.localeCompare(b.path)); const bytes=Buffer.from(JSON.stringify({root:base,files})); return {root:base,files,sha256:createHash("sha256").update(bytes).digest("hex")};
}
function ownedValue(repo,root,nonce,{backendLog=null,entrypoint,distManifest=null,stage="PREPARING",createdAt=new Date().toISOString()}={}){return{schemaVersion:1,kind:"p1-manual-acceptance",nonce,repositoryRoot:repo,root,status:"PENDING",stage,createdAt,listener:{host:HOST,port:PORT,state:"stopped"},backendLog,entrypoint,distManifest,resources:[root,{kind:"fastify",host:HOST,port:PORT}]};}
function validEntrypoint(value,repo){return value?.path===join(repo,"backend/dist/server.js")&&Number.isSafeInteger(value.dev)&&Number.isSafeInteger(value.ino)&&Number.isSafeInteger(value.size)&&value.size>0&&HEX64.test(value.sha256??"");}
async function readBoundEntrypoint(repo){
if(!Number.isInteger(constants.O_NOFOLLOW))throw new Error("production entrypoint no-follow protection is unavailable");const path=join(repo,"backend/dist/server.js");let handle;
@@ -413,9 +418,9 @@ export async function prepareManual(options={}){
const{repositoryRoot=defaultRepositoryRoot,skipBuild=false}=options,repo=realpathSync(repositoryRoot),root=fixedManualRoot(repo),lifecycle=await acquireLifecycle(repo,"prepare");let entryBinding,ownershipCreated=false;
try{
await requireLifecycleContext(lifecycle);await checkPrerequisites(repo);if(!skipBuild)await run("npm",["--prefix",join(repo,"backend"),"run","build"]);await requireLifecycleContext(lifecycle);
entryBinding=await readBoundEntrypoint(repo);const entrypoint=entryBinding.identity;await entryBinding.handle.close();entryBinding=undefined;
entryBinding=await readBoundEntrypoint(repo);const entrypoint=entryBinding.identity;await entryBinding.handle.close();entryBinding=undefined;const distribution=await distManifest(repo);
noSymlinkExisting(repo,root);try{await mkdir(root,{recursive:false,mode:0o700});}catch(error){if(error.code==="EEXIST")throw new Error("manual acceptance root already exists; stop/cleanup it explicitly");throw error;}await bindLifecycleRoot(lifecycle,root);
const nonce=randomBytes(32).toString("hex"),createdAt=new Date().toISOString();await exclusiveRecord(join(root,"ownership.json"),ownedValue(repo,root,nonce,{entrypoint,createdAt}),"manual ownership");ownershipCreated=true;await requireLifecycleContext(lifecycle,{root:true});
const nonce=randomBytes(32).toString("hex"),createdAt=new Date().toISOString();await exclusiveRecord(join(root,"ownership.json"),ownedValue(repo,root,nonce,{entrypoint,distManifest:distribution,createdAt}),"manual ownership");ownershipCreated=true;await requireLifecycleContext(lifecycle,{root:true});
for(const path of ["installation/registry","installation/data","installation/runtime","fixture-secrets","fixtures/descriptors","requests","responses","exports/raw","exports/extracted","rendered","logs","commands"]){await mkdir(join(root,path),{recursive:true,mode:path==="fixture-secrets"?0o700:0o755});await requireLifecycleContext(lifecycle,{root:true});}
const backendLogPath=join(root,"logs/backend.log"),backendLogHandle=await open(backendLogPath,"wx",0o600);let backendLogEntry;try{await backendLogHandle.chmod(0o600);await backendLogHandle.sync();backendLogEntry=await backendLogHandle.stat();}finally{await backendLogHandle.close();}directorySync(dirname(backendLogPath));const backendLog={path:backendLogPath,dev:backendLogEntry.dev,ino:backendLogEntry.ino};
await requireLifecycleContext(lifecycle,{root:true});
@@ -423,7 +428,7 @@ export async function prepareManual(options={}){
const items=descriptors();for(const workspace of items)await atomicWrite(join(root,"fixtures/descriptors",`${workspace.workspace.id}.json`),`${JSON.stringify(workspace,null,2)}\n`);
const secrets={"dwh-password":`DWH-${randomBytes(16).toString("hex")}`,"evidence-signed-urls.json":JSON.stringify([`https://evidence.example.test/guide.md?token=SIGNED-${randomBytes(16).toString("hex")}`]),"evidence-access":`ACCESS-${randomBytes(16).toString("hex")}`,"evidence-secret":`SECRET-${randomBytes(16).toString("hex")}`,"evidence-session":`SESSION-${randomBytes(16).toString("hex")}`};for(const[name,value]of Object.entries(secrets))await atomicWrite(join(root,"fixture-secrets",name),value,0o600);
const env={};for(const workspace of items){const ns=workspace.workspace.id.toUpperCase().replaceAll("-","_"),prefix=`THT_WS_${ns}`;Object.assign(env,{[`${prefix}_DWH_TRANSPORT`]:"postgres_direct",[`${prefix}_DWH_HOST`]:"dwh.invalid",[`${prefix}_DWH_PORT`]:"5432",[`${prefix}_DWH_USER`]:"reader",[`${prefix}_DWH_PASSWORD_FILE`]:join(root,"fixture-secrets/dwh-password")});}Object.assign(env,{THT_WORKSPACE_SECRET_ROOTS:join(root,"fixture-secrets"),THT_WS_P1_HTTP_EVIDENCE_SIGNED_URLS_FILE:join(root,"fixture-secrets/evidence-signed-urls.json"),THT_WS_P1_S3_EVIDENCE_ACCESS_KEY_FILE:join(root,"fixture-secrets/evidence-access"),THT_WS_P1_S3_EVIDENCE_SECRET_KEY_FILE:join(root,"fixture-secrets/evidence-secret"),THT_WS_P1_S3_EVIDENCE_SESSION_TOKEN_FILE:join(root,"fixture-secrets/evidence-session")});
await atomicWrite(join(root,"installation/bindings.env"),Object.entries(env).map(([k,v])=>`${k}=${quote(v)}`).join("\n")+"\n");await atomicWrite(join(root,"installation/base.yaml"),"{}\n");for(const[name,value]of Object.entries(requestFixtures(items)))await atomicWrite(join(root,"requests",name),`${JSON.stringify(value,null,2)}\n`);await writeCommands(repo,root);await atomicWrite(join(root,"GUIDE.md"),guide(repo,root),0o600);await requireLifecycleContext(lifecycle,{root:true});await atomicWrite(join(root,"ownership.json"),`${JSON.stringify(ownedValue(repo,root,nonce,{backendLog,entrypoint,stage:"READY",createdAt}),null,2)}\n`);await requireLifecycleContext(lifecycle,{root:true});return{repositoryRoot:repo,root,nonce};
await atomicWrite(join(root,"installation/bindings.env"),Object.entries(env).map(([k,v])=>`${k}=${quote(v)}`).join("\n")+"\n");await atomicWrite(join(root,"installation/base.yaml"),"{}\n");for(const[name,value]of Object.entries(requestFixtures(items)))await atomicWrite(join(root,"requests",name),`${JSON.stringify(value,null,2)}\n`);await writeCommands(repo,root);await atomicWrite(join(root,"GUIDE.md"),guide(repo,root),0o600);await requireLifecycleContext(lifecycle,{root:true});await atomicWrite(join(root,"ownership.json"),`${JSON.stringify(ownedValue(repo,root,nonce,{backendLog,entrypoint,distManifest:distribution,stage:"READY",createdAt}),null,2)}\n`);await requireLifecycleContext(lifecycle,{root:true});return{repositoryRoot:repo,root,nonce};
}catch(error){
if(entryBinding)await entryBinding.handle.close().catch(()=>{});
if(ownershipCreated){try{await requireLifecycleContext(lifecycle,{root:true});}catch{await cleanupFailedPrepare(repo,lifecycle).catch(()=>{});throw new Error("manual acceptance parent or root identity changed during prepare");}}
@@ -433,6 +438,11 @@ export async function prepareManual(options={}){
function portAvailable(port,label=`${HOST}:${port}`){return new Promise((resolvePromise,reject)=>{const server=net.createServer();server.once("error",error=>error.code==="EADDRINUSE"?reject(new Error(`${label} is occupied`)):reject(error));server.listen({host:HOST,port,exclusive:true},()=>server.close(()=>resolvePromise()));});}
async function requireCanonicalDirectory(path,label){const entry=await lstat(path);if(!entry.isDirectory()||entry.isSymbolicLink()||await realpath(path)!==path)throw new Error(`${label} directory identity is unsafe`);return entry;}
async function requireAbsent(path,label){try{await lstat(path);throw new Error(`${label} is legacy or unsafe`);}catch(error){if(error.code!=="ENOENT")throw error;}}
async function validateDistIntegrity(repo,owned){
if(!owned.distManifest?.sha256||owned.distManifest.root!==join(repo,"backend","dist"))throw new Error("production distribution manifest is missing");
const current=await distManifest(repo); if(current.sha256!==owned.distManifest.sha256||JSON.stringify(current.files)!==JSON.stringify(owned.distManifest.files))throw new Error("production distribution identity changed");
return current;
}
async function validateServeFilesystem(repo,root,owned){
if(root!==fixedManualRoot(repo))throw new Error("owned root identity is unsafe");
for(const [path,label] of [
@@ -443,7 +453,7 @@ async function validateServeFilesystem(repo,root,owned){
])await requireCanonicalDirectory(path,label);
await requireAbsent(legacySupervisorPath(root),"legacy supervisor");
if(owned.stage!=="READY")throw new Error("manual acceptance preparation is incomplete");
const script=join(repo,"backend/dist/server.js");await requireEntrypointPathIdentity(owned.entrypoint);
const script=join(repo,"backend/dist/server.js");await requireEntrypointPathIdentity(owned.entrypoint);await validateDistIntegrity(repo,owned);
const logPath=join(root,"logs/backend.log");
if(owned.backendLog?.path!==logPath)throw new Error("backend log ownership identity is unsafe");
return{script,logPath};
@@ -468,7 +478,7 @@ async function readPid(root){const path=join(root,"backend.pid"),entry=await lst
async function validateProcess(repo,root,owned,pidRecord){
const script=join(repo,"backend/dist/server.js"),entrypoint=owned.entrypoint;
if(pidRecord.schemaVersion!==1||pidRecord.kind!=="p1-manual-backend"||pidRecord.status!=="RUNNING"||!Number.isSafeInteger(pidRecord.pid)||pidRecord.pid<2||!HEX64.test(pidRecord.reservationNonce??"")||pidRecord.nonce!==owned.nonce||pidRecord.root!==root||pidRecord.repositoryRoot!==repo||pidRecord.executable!==process.execPath||pidRecord.preload!==PRELOAD||pidRecord.script!==script||JSON.stringify(pidRecord.entrypoint)!==JSON.stringify(entrypoint)||!pidRecord.startIdentity||pidRecord.control?.host!==HOST||pidRecord.control?.port!==CONTROL_PORT)throw new Error("backend process identity mismatch; refusing cooperative control");
await requireEntrypointPathIdentity(entrypoint);if(!alive(pidRecord.pid))throw new Error("backend PID is stale; operator inspection required");
await requireEntrypointPathIdentity(entrypoint);await validateDistIntegrity(repo,owned);if(!alive(pidRecord.pid))throw new Error("backend PID is stale; operator inspection required");
const[start,args,cwd,executable]=await Promise.all([processStart(pidRecord.pid),processArgs(pidRecord.pid),processCwd(pidRecord.pid),processExecutable(pidRecord.pid)]);
const expectedArgs=[pidRecord.executable,"--import",pidRecord.preload,pidRecord.script,`--p1-manual-nonce=${owned.nonce}`,`--p1-root=${root}`,`--p1-control-nonce=${pidRecord.reservationNonce}`,`--p1-entry-sha256=${entrypoint.sha256}`,`--p1-entry-dev=${entrypoint.dev}`,`--p1-entry-ino=${entrypoint.ino}`].join(" ");
if(start!==pidRecord.startIdentity||cwd!==repo||executable!==realpathSync(pidRecord.executable)||args!==expectedArgs)throw new Error("backend process identity mismatch; refusing cooperative control");return true;
@@ -601,6 +601,22 @@ test("prepare binds production entry bytes and serve rejects a regular replaceme
await assert.rejects(lstat(malicious)); await assert.rejects(lstat(join(run.root,"backend.pid")));
});
test("serve rejects replacement of an imported production dependency", { concurrency: false }, async () => {
const repo=await fakeRepo(),malicious=join(repo,"dependency-executed"); await installFakeServer(repo);
const server=join(repo,"backend/dist/server.js"), original=await readFile(server); await writeFile(join(repo,"backend/dist/dep.js"),"export const dependency = true;\n"); await writeFile(server,`import \"./dep.js\";\n${original}`);
const run=await prepareManual({repositoryRoot:repo,skipBuild:true}); const before=await readFile(server); await writeFile(join(repo,"backend/dist/dep.js"),`import {writeFileSync} from \"node:fs\"; writeFileSync(${JSON.stringify(malicious)},\"bad\");\n`);
assert.deepEqual(await readFile(server),before); await assert.rejects(serveManual({repositoryRoot:repo}),/distribution|identity|manifest/i); await assert.rejects(lstat(malicious)); await assert.rejects(lstat(join(run.root,"backend.pid")));
});
test("serve refuses a replaced backend dist dependency after prepare", { concurrency: false }, async () => {
const repo=await fakeRepo(); await installFakeServer(repo);
const dependency=join(repo,"backend/dist/dependency.js"); await writeFile(dependency,"export const value = 1;\n");
const run=await prepareManual({repositoryRoot:repo,skipBuild:true}),marker=join(repo,"dependency-replaced-executed");
await writeFile(dependency,`import { writeFileSync } from "node:fs";writeFileSync(${JSON.stringify(marker)},"bad");export const value = 2;\n`);
await assert.rejects(serveManual({repositoryRoot:repo}),/distribution|manifest|identity/i);
await assert.rejects(lstat(marker)); await assert.rejects(lstat(join(run.root,"backend.pid")));
});
test("opened production FD prevents deterministic check-spawn replacement execution", { concurrency: false }, async () => {
const repo=await fakeRepo(),safe=join(repo,"safe-executed"),malicious=join(repo,"malicious-executed"); await installFakeServer(repo,{marker:safe});
const run=await prepareManual({repositoryRoot:repo,skipBuild:true}),script=join(repo,"backend/dist/server.js"),replacement=join(repo,"replacement-server.js");
+22 -1
View File
@@ -1,5 +1,6 @@
#!/usr/bin/env node
import { spawnSync } from "node:child_process";
import { createHash } from "node:crypto";
import { lstatSync, realpathSync } from "node:fs";
import { lstat, mkdir, readFile, realpath } from "node:fs/promises";
import { basename, dirname, isAbsolute, join, relative, resolve, sep } from "node:path";
@@ -90,6 +91,16 @@ export async function renderOwnedSnapshot({ repositoryRoot = defaultRepositoryRo
if (!match || !HEX40.test(match[1])) throw new Error("snapshot is not commit addressed");
assertNoSymlinks(root, snapshot); const snapshotEntry = await lstat(snapshot);
if (!snapshotEntry.isFile() || snapshotEntry.isSymbolicLink() || await realpath(snapshot) !== snapshot) throw new Error("snapshot is unsafe");
const snapshotBytes = await readFile(snapshot);
const snapshotDigest = createHash("sha256").update(snapshotBytes).digest("hex");
const manifestPath = join(dirname(snapshot), "snapshot.json");
assertNoSymlinks(root, manifestPath);
const manifestEntry = await lstat(manifestPath);
if (!manifestEntry.isFile() || manifestEntry.isSymbolicLink() || await realpath(manifestPath) !== manifestPath) throw new Error("snapshot manifest is unsafe");
let manifest;
try { manifest = JSON.parse(await readFile(manifestPath, "utf8")); } catch { throw new Error("snapshot manifest is malformed"); }
const yamlName = `${match[2]}.yaml`;
if (manifest?.head !== match[1] || manifest?.files?.[yamlName] !== snapshotDigest) throw new Error("snapshot content does not match its immutable manifest");
if (!below(renderedRoot, output) || dirname(output) !== renderedRoot || !output.endsWith(".yaml")) throw new Error("output is not an owned rendered path");
assertNoSymlinks(root, dirname(output));
try { if ((await lstat(output)).isSymbolicLink()) throw new Error("output is unsafe"); } catch (error) { if (error.code !== "ENOENT") throw error; }
@@ -103,7 +114,17 @@ export async function renderOwnedSnapshot({ repositoryRoot = defaultRepositoryRo
semanticRuntime: { internalQdrantUrl: "http://qdrant:6333", internalEmbeddingUrl: "http://embedding:11434", internalEmbeddingModel: "qwen3-embedding:0.6b", internalEmbeddingDimensions: 1024 },
});
let lease;
try { lease = runner.acquireWorkspaceRuntime(snapshot); if(beforePublish)await beforePublish({output,renderedRoot}); await atomicCopy(lease.path, output); }
try {
lease = runner.acquireWorkspaceRuntime(snapshot);
const verifySnapshot = async () => {
const current = await readFile(snapshot);
if (createHash("sha256").update(current).digest("hex") !== snapshotDigest) throw new Error("snapshot content changed during rendering");
};
await verifySnapshot();
if(beforePublish)await beforePublish({output,renderedRoot});
await verifySnapshot();
await atomicCopy(lease.path, output);
}
finally {
if (lease) lease.release();
for (const key of Object.keys(env)) { if (prior[key] === undefined) delete process.env[key]; else process.env[key] = prior[key]; }
@@ -1,4 +1,5 @@
import assert from "node:assert/strict";
import { createHash } from "node:crypto";
import { chmod, lstat, mkdir, mkdtemp, readFile, realpath, rename, rm, symlink, writeFile } from "node:fs/promises";
import { tmpdir } from "node:os";
import { dirname, join } from "node:path";
@@ -32,6 +33,7 @@ evidence:
source: {type: filesystem, uri: workspace-content/p1-filesystem/evidence, patterns: ["**/*.md"], max_bytes: 10485760}
policy: {max_chunk_chars: 4000, retain_published_generations: 3}
`);
const snapshotBytes=await readFile(snapshot); await writeFile(join(dirname(snapshot),"snapshot.json"),JSON.stringify({head:commit,revisions:[{id:"p1-filesystem",commit,blob:"0".repeat(40),snapshotPath:snapshot,state:"operational"}],files:{"p1-filesystem.yaml":createHash("sha256").update(snapshotBytes).digest("hex")}}));
const env={THT_WS_P1_FILESYSTEM_DWH_TRANSPORT:"postgres_direct",THT_WS_P1_FILESYSTEM_DWH_HOST:"dwh.invalid",THT_WS_P1_FILESYSTEM_DWH_PORT:"5432",THT_WS_P1_FILESYSTEM_DWH_USER:"reader",THT_WS_P1_FILESYSTEM_DWH_PASSWORD_FILE:secret};
return {repo,root,snapshot,env};
}
@@ -44,6 +46,8 @@ test("renderer rejects unowned, symlink, and out-of-root paths",async()=>{ const
test("renderer releases its acquired lease when atomic output copy fails",async()=>{ const f=await fixture(); const output=join(f.root,"rendered/existing.yaml"); await mkdir(output); await assert.rejects(renderOwnedSnapshot({repositoryRoot:f.repo,ownershipPath:join(f.root,"ownership.json"),snapshotPath:f.snapshot,outputPath:output,env:{...process.env,...f.env}}),/anchored|publication|unsafe/); assert.deepEqual(await (await import("node:fs/promises")).readdir(join(f.root,"installation/registry/snapshots/runtime")),[]); });
test("renderer rejects a regular snapshot replacement against its manifest",async()=>{ const f=await fixture(); const output=join(f.root,"rendered/replaced.yaml"); await assert.rejects(renderOwnedSnapshot({repositoryRoot:f.repo,ownershipPath:join(f.root,"ownership.json"),snapshotPath:f.snapshot,outputPath:output,env:{...process.env,...f.env},beforePublish:async()=>{await writeFile(f.snapshot,"workspace:\n schema_version: 3\n id: p1-filesystem\n name: replaced\n")}}),/snapshot content changed/); await assert.rejects(lstat(output)); });
test("renderer anchors publication when rendered parent is concurrently swapped", async()=>{
const f=await fixture(),output=join(f.root,"rendered/raced.yaml"),moved=join(f.root,"rendered-moved"),outside=join(f.repo,"outside-rendered"); await mkdir(outside);
await assert.rejects(renderOwnedSnapshot({repositoryRoot:f.repo,ownershipPath:join(f.root,"ownership.json"),snapshotPath:f.snapshot,outputPath:output,env:{...process.env,...f.env},beforePublish:async()=>{await rename(join(f.root,"rendered"),moved);await symlink(outside,join(f.root,"rendered"));}}),/identity|changed|unsafe|publication/i);