import { createHash } from "node:crypto"; import { constants, createWriteStream } from "node:fs"; import fs from "node:fs/promises"; import path from "node:path"; import { spawn } from "node:child_process"; import { z } from "zod"; import type { JazzThoughtStore } from "../jazz/store.js"; import { readVerifiedBlob } from "./blob.js"; import { requestArtifactStorage } from "./append.js"; import { ARTIFACT_KINDS, ARTIFACT_MEDIA_TYPES, ARTIFACT_VISIBILITY } from "./types.js"; const id = z.string().min(1).max(200).regex(/^[A-Za-z0-9][A-Za-z0-9._:-]*$/); const text = z.string().min(1).max(2_000); const MAX_SSH_OUTPUT_BYTES = 12 * 1024 * 1024; export const exeArtifactRequestSchema = z.object({ requestId: id, leaseId: id, relativeOutputFilePath: z.string().min(1).max(500).refine(normalizedRelative, "not-normalized-relative"), artifactId: id, artifactVersion: z.number().int().positive().max(1_000_000), kind: z.enum(ARTIFACT_KINDS), title: text, summary: text, mediaType: z.enum(ARTIFACT_MEDIA_TYPES), provenanceLabel: text, visibility: z.enum(ARTIFACT_VISIBILITY).default("private"), requestingAgentId: id.optional(), runId: id.optional(), sourceEventId: id.optional(), supersedes: id.optional(), }).strict(); export type ExeArtifactRequest = z.infer; export interface ExeTransportConfig { host: "cameron.exe.xyz"; workspaceRoot: "/srv/thoughtstream/workspaces/stream-v1"; workspaceId: "stream-v1"; leaseId: string; stagingRoot: string; } export interface ExeTransportReceipt { requestId: string; requestEventId: string; artifactEventId: string; artifactId: string; sha256: string; bytes: number; status: "stored"; timestamp: string } const REMOTE_HELPER = String.raw`import base64,hashlib,json,os,stat,sys cfg=json.loads(base64.urlsafe_b64decode(sys.argv[1]+"==")) root=cfg["root"]; rel=cfg["relative"] def fail(code): print(json.dumps({"ok":False,"reasonCode":code},separators=(",",":")));sys.exit(0) try: mpath=os.path.join(root,".thoughtstream-workspace.json") rst=os.lstat(root) if stat.S_ISLNK(rst.st_mode) or not stat.S_ISDIR(rst.st_mode) or os.path.realpath(root)!=root: fail("workspace-invalid") with open(mpath,"rb") as f: m=json.load(f) if m.get("workspaceId")!=cfg["workspaceId"] or m.get("leaseId")!=cfg["leaseId"] or m.get("vmIdentity")!=cfg["host"]: fail("workspace-mismatch") parts=rel.split("/") if not rel or rel.startswith("/") or "\\" in rel or any(x in ("",".","..") for x in parts): fail("path-invalid") cur=os.path.join(root,"output") for part in parts: cur=os.path.join(cur,part); s=os.lstat(cur) if stat.S_ISLNK(s.st_mode): fail("symlink") if os.path.realpath(cur)!=cur or not stat.S_ISREG(os.lstat(cur).st_mode): fail("not-regular") ext=os.path.splitext(cur)[1].lower(); pairs={".md":["text/markdown"],".markdown":["text/markdown"],".txt":["text/plain"],".json":["application/json"],".yaml":["application/yaml","text/yaml"],".yml":["application/yaml","text/yaml"],".png":["image/png"],".jpg":["image/jpeg"],".jpeg":["image/jpeg"]} if cfg["mediaType"] not in pairs.get(ext,[]): fail("media-mismatch") maximum=10485760 if cfg["mediaType"].startswith("image/") else 1048576 before=os.lstat(cur) if before.st_size<=0 or before.st_size>maximum: fail("size-invalid") fd=os.open(cur,os.O_RDONLY|os.O_NOFOLLOW) try: opened=os.fstat(fd); data=b"" while len(data)<=maximum: chunk=os.read(fd,min(65536,maximum+1-len(data))) if not chunk: break data+=chunk after=os.fstat(fd) finally: os.close(fd) ident=lambda s:(s.st_dev,s.st_ino,s.st_size,s.st_mtime_ns,s.st_ctime_ns) if ident(before)!=ident(opened) or ident(opened)!=ident(after) or len(data)!=after.st_size: fail("unstable-snapshot") mt=cfg["mediaType"] if mt=="image/png" and not data.startswith(b"\x89PNG\r\n\x1a\n"): fail("magic-invalid") if mt=="image/jpeg" and not (data.startswith(b"\xff\xd8") and data.endswith(b"\xff\xd9")): fail("magic-invalid") if not mt.startswith("image/"): try: value=data.decode("utf-8") except UnicodeDecodeError: fail("encoding-invalid") if "\x00" in value: fail("encoding-invalid") if mt=="application/json": json.loads(value) header=json.dumps({"ok":True,"bytes":len(data),"sha256":hashlib.sha256(data).hexdigest(),"mediaType":mt,"relativePath":rel},separators=(",",":" )).encode() sys.stdout.buffer.write(str(len(header)).encode()+b"\n"+header+data) except (OSError,ValueError,json.JSONDecodeError): fail("remote-validation-failed")`; const REMOTE_RECEIPT_HELPER = String.raw`import base64,json,os,sys p=sys.argv[1] r=json.loads(base64.urlsafe_b64decode(sys.argv[2]+"==")) n=r["requestId"]+".json" assert all(c.isalnum() or c in "._:-" for c in r["requestId"]) d=(json.dumps(r,separators=(",",":"),sort_keys=True)+"\n").encode() q=os.path.join(p,n) try: f=os.open(q,os.O_WRONLY|os.O_CREAT|os.O_EXCL,0o600) os.write(f,d) os.close(f) except FileExistsError: assert open(q,"rb").read()==d`; const REMOTE_PROCESSED_HELPER = String.raw`import os,sys s=os.path.join(sys.argv[1],sys.argv[3]) d=os.path.join(sys.argv[2],sys.argv[3]) if os.path.exists(d): assert open(s,"rb").read()==open(d,"rb").read() os.unlink(s) else: os.rename(s,d)`; export async function processExeArtifactRequest(store: JazzThoughtStore, config: ExeTransportConfig, raw: unknown): Promise { const request = exeArtifactRequestSchema.parse(raw); if (request.leaseId !== config.leaseId) throw reason("lease-mismatch"); await ensurePrivateStaging(config.stagingRoot); const framed = await remoteFetch(config, request); const hash = createHash("sha256").update(framed.bytes).digest("hex"); if (hash !== framed.header.sha256 || framed.bytes.byteLength !== framed.header.bytes || framed.header.mediaType !== request.mediaType || framed.header.relativePath !== request.relativeOutputFilePath) throw reason("content-integrity"); const requestStage = await fs.mkdtemp(path.join(config.stagingRoot, "request-")); await fs.chmod(requestStage, 0o700); const staged = path.join(requestStage, "artifact" + path.extname(request.relativeOutputFilePath)); try { await fs.writeFile(staged, framed.bytes, { mode: 0o600, flag: "wx" }); const result = await requestArtifactStorage({ store, workspaceRoot: requestStage, workspaceLeaseId: config.leaseId, workspaceRootLabel: `exe-workspace:${config.workspaceId}`, requestId: request.requestId, relativeFilePath: path.basename(staged), artifactId: request.artifactId, artifactVersion: request.artifactVersion, kind: request.kind, title: request.title, summary: request.summary, mediaType: request.mediaType, provenance: { source: "exe-dev-private-workspace", label: request.provenanceLabel }, visibility: request.visibility, ...(request.requestingAgentId ? { requestingAgentId: request.requestingAgentId } : {}), ...(request.runId ? { requestingRunId: request.runId } : {}), ...(request.sourceEventId ? { sourceEventId: request.sourceEventId } : {}), ...(request.supersedes ? { supersedesArtifactEventId: request.supersedes } : {}), }); const persisted = await store.getEvent(result.event.id); if (!persisted || persisted.id !== result.event.id) throw reason("durable-readback-failed"); const blob = persisted.payload.blob as { algorithm: "sha256"; sha256: string; relativePath: string; byteCount: number }; const readback = await readVerifiedBlob(store.getArtifactRoot(), blob, request.mediaType); if (!readback.equals(framed.bytes)) throw reason("durable-readback-failed"); const receipt: ExeTransportReceipt = { requestId: request.requestId, requestEventId: result.requestEvent!.id, artifactEventId: result.event.id, artifactId: request.artifactId, sha256: hash, bytes: framed.bytes.byteLength, status: "stored", timestamp: persisted.observedAt }; await remoteReceipt(config, receipt); return receipt; } finally { await fs.rm(requestStage, { recursive: true, force: true }); } } export async function runExeArtifactTransportOnce(store: JazzThoughtStore, config: ExeTransportConfig): Promise<{ receipts: ExeTransportReceipt[]; failures: Array<{ requestId: string; reasonCode: string }> }> { const lock = await acquireLock(path.join(config.stagingRoot, "transport.lock")); try { const names = await remoteList(config); const receipts: ExeTransportReceipt[] = []; const failures: Array<{ requestId: string; reasonCode: string }> = []; for (const name of names) { let raw: unknown; let requestId = name.replace(/\.json$/, ""); try { raw = JSON.parse(await remoteReadRequest(config, name)); requestId = typeof (raw as { requestId?: unknown }).requestId === "string" ? (raw as { requestId: string }).requestId : requestId; receipts.push(await processExeArtifactRequest(store, config, raw)); await remoteProcessed(config, name); } catch (error) { failures.push({ requestId: safeId(requestId), reasonCode: reasonCode(error) }); } } return { receipts, failures }; } finally { await lock.close(); await fs.rm(path.join(config.stagingRoot, "transport.lock"), { force: true }); } } async function remoteFetch(config: ExeTransportConfig, request: ExeArtifactRequest): Promise<{ header: { bytes: number; sha256: string; mediaType: string; relativePath: string }; bytes: Buffer }> { const payload = b64({ root: config.workspaceRoot, workspaceId: config.workspaceId, leaseId: config.leaseId, host: config.host, relative: request.relativeOutputFilePath, mediaType: request.mediaType }); const output = await ssh(config.host, ["python3", "-c", b64(REMOTE_HELPER), payload]); const newline = output.indexOf(10); if (newline < 1) throw reason("remote-frame-invalid"); const length = Number(output.subarray(0, newline).toString("ascii")); if (!Number.isSafeInteger(length) || length < 2 || length > 2048) { const failure = parseFailure(output); throw reason(failure); } const header = JSON.parse(output.subarray(newline + 1, newline + 1 + length).toString("utf8")); const bytes = output.subarray(newline + 1 + length); if (!header.ok) throw reason(header.reasonCode ?? "remote-validation-failed"); return { header, bytes }; } async function remoteList(config: ExeTransportConfig): Promise { const output = await ssh(config.host, ["python3", "-c", b64("import json,os,sys;p=sys.argv[1];v=sorted(x for x in os.listdir(p) if x.endswith('.json') and '/' not in x);assert len(v)<=256;print(json.dumps(v))"), `${config.workspaceRoot}/requests`]); return JSON.parse(output.toString("utf8")); } async function remoteReadRequest(config: ExeTransportConfig, name: string): Promise { if (!/^[A-Za-z0-9][A-Za-z0-9._-]{0,199}\.json$/.test(name)) throw reason("request-name-invalid"); return (await ssh(config.host, ["python3", "-c", b64("import os,sys; p=os.path.join(sys.argv[1],sys.argv[2]); s=os.lstat(p); assert not os.path.islink(p) and os.path.isfile(p) and s.st_size<=16384; sys.stdout.buffer.write(open(p,'rb').read())"), `${config.workspaceRoot}/requests`, name])).toString("utf8"); } async function remoteReceipt(config: ExeTransportConfig, receipt: ExeTransportReceipt): Promise { await ssh(config.host, ["python3", "-c", b64(REMOTE_RECEIPT_HELPER), `${config.workspaceRoot}/receipts`, b64(receipt)]); } async function remoteProcessed(config: ExeTransportConfig, name: string): Promise { await ssh(config.host, ["python3", "-c", b64(REMOTE_PROCESSED_HELPER), `${config.workspaceRoot}/requests`, `${config.workspaceRoot}/requests/processed`, name]); } function ssh(host: string, args: string[], input?: Buffer): Promise { const quote=(arg:string)=>`'${arg.replaceAll("'", "'\\''")}'`; return new Promise((resolve, reject) => { const remote = args[1] === "-c" && args.length >= 3 ? `python3 -c 'import base64,sys;exec(base64.urlsafe_b64decode(\"${args[2]}==\"))' ${quote(args[3] ?? "")}${args.slice(4).map((arg)=>` ${quote(arg)}`).join("")}` : args.map(quote).join(" "); const child = spawn("ssh", ["-o", "BatchMode=yes", "-o", "ConnectTimeout=10", "-o", "StrictHostKeyChecking=yes", "-o", "ClearAllForwardings=yes", "--", host, remote], { stdio: ["pipe", "pipe", "pipe"] }); const out: Buffer[]=[]; const err: Buffer[]=[]; let outputBytes=0; let outputOverflow=false; child.stdout.on("data", x=>{ outputBytes+=x.length; if(outputBytes>MAX_SSH_OUTPUT_BYTES){ outputOverflow=true; child.kill("SIGKILL"); return; } out.push(x); }); child.stderr.on("data", x=>{ if(Buffer.concat(err).byteLength<4096) err.push(x); }); child.on("error", reject); child.on("close", code=>outputOverflow?reject(reason("ssh-output-too-large")):code===0?resolve(Buffer.concat(out)):reject(Object.assign(reason("ssh-failed"), { detail: Buffer.concat(err).toString("utf8").slice(0, 500) }))); if (input) child.stdin.end(input); else child.stdin.end(); }); } async function ensurePrivateStaging(root: string): Promise { if (!path.isAbsolute(root)) throw reason("staging-invalid"); await fs.mkdir(root,{recursive:true,mode:0o700}); await fs.chmod(root,0o700); const s=await fs.lstat(root); if (!s.isDirectory()||s.isSymbolicLink()||await fs.realpath(root)!==root) throw reason("staging-invalid"); } async function acquireLock(file: string): Promise { await ensurePrivateStaging(path.dirname(file)); try { return await fs.open(file, constants.O_CREAT|constants.O_EXCL|constants.O_WRONLY,0o600); } catch { throw reason("already-running"); } } function b64(value: unknown): string { return Buffer.from(typeof value === "string" ? value : JSON.stringify(value)).toString("base64url"); } function normalizedRelative(value: string): boolean { return !path.isAbsolute(value) && !value.includes("\\") && !value.endsWith("/") && value.split("/").every(part=>part!==""&&part!=="."&&part!==".."); } function reason(code: string): Error { return Object.assign(new Error(code), { reasonCode: code }); } function reasonCode(error: unknown): string { const code=(error as {reasonCode?:unknown})?.reasonCode; return typeof code==="string"&&/^[a-z0-9-]{1,64}$/.test(code)?code:"request-rejected"; } function safeId(value: string): string { return /^[A-Za-z0-9][A-Za-z0-9._:-]{0,199}$/.test(value)?value:"invalid-request"; } function parseFailure(output: Buffer): string { try { const value=JSON.parse(output.toString("utf8")); return typeof value.reasonCode==="string"?value.reasonCode:"remote-frame-invalid"; } catch { return "remote-frame-invalid"; } }