feat: archive Codex sessions as private cognition evidence

This commit is contained in:
冰朔 2026-08-05 12:42:16 +08:00
commit 7d85d79e61
8 changed files with 798 additions and 0 deletions

View file

@ -0,0 +1,503 @@
#!/usr/bin/env node
import crypto from "node:crypto";
import fs from "node:fs";
import path from "node:path";
import process from "node:process";
import readline from "node:readline";
import { fileURLToPath } from "node:url";
import { pipeline } from "node:stream/promises";
import { Transform } from "node:stream";
const SCHEMA = "guanghu.codex-session-cold-archive/v1";
const INDEX_SCHEMA = "guanghu.development-cognition-event/v1";
const DEFAULT_STABLE_SECONDS = 300;
const MAX_VISIBLE_TEXT = 4000;
function parseArgs(argv) {
const args = { command: argv[2] || "", execute: false };
for (let i = 3; i < argv.length; i += 1) {
const key = argv[i];
if (key === "--execute") {
args.execute = true;
continue;
}
if (!key.startsWith("--") || i + 1 >= argv.length) {
throw new Error(`INVALID_ARGUMENT: ${key}`);
}
args[key.slice(2).replaceAll("-", "_")] = argv[++i];
}
return args;
}
function required(value, name) {
if (!value) throw new Error(`REQUIRED_ARGUMENT: --${name.replaceAll("_", "-")}`);
return path.resolve(value);
}
function isoForPath(date = new Date()) {
return date.toISOString().replaceAll(":", "-").replace(/\.\d{3}Z$/, "Z");
}
function safeRelative(root, candidate) {
const relative = path.relative(root, candidate);
if (!relative || relative.startsWith("..") || path.isAbsolute(relative)) {
throw new Error(`UNSAFE_RELATIVE_PATH: ${candidate}`);
}
return relative.split(path.sep).join("/");
}
async function listSessionFiles(root) {
const found = [];
async function walk(current) {
const entries = await fs.promises.readdir(current, { withFileTypes: true });
for (const entry of entries) {
if (entry.name.startsWith("._") || entry.name === ".DS_Store") continue;
const absolute = path.join(current, entry.name);
if (entry.isDirectory()) await walk(absolute);
else if (entry.isFile() && entry.name.endsWith(".jsonl")) found.push(absolute);
}
}
await walk(root);
return found.sort();
}
export async function removeAppleDouble(root) {
let removed = 0;
async function walk(current) {
const entries = await fs.promises.readdir(current, { withFileTypes: true });
for (const entry of entries) {
const absolute = path.join(current, entry.name);
if (entry.isDirectory()) {
await walk(absolute);
} else if (entry.isFile() && entry.name.startsWith("._")) {
await fs.promises.unlink(absolute);
removed += 1;
}
}
}
await walk(path.resolve(root));
return removed;
}
export function redactVisibleText(input, home = process.env.HOME || "") {
let text = String(input ?? "");
const counts = {
private_key: 0,
authorization: 0,
named_secret: 0,
email: 0,
ip_address: 0,
credential_url: 0,
high_entropy: 0,
};
const replace = (pattern, label, key) => {
text = text.replace(pattern, () => {
counts[key] += 1;
return `[REDACTED:${label}]`;
});
};
replace(/-----BEGIN [^-]+ PRIVATE KEY-----[\s\S]*?-----END [^-]+ PRIVATE KEY-----/gi, "PRIVATE_KEY", "private_key");
replace(/\b(?:Bearer|Basic)\s+[A-Za-z0-9._~+/=-]{8,}/gi, "AUTHORIZATION", "authorization");
replace(/\b(?:api[_ -]?key|token|secret|password|passwd|authorization)\b\s*[:=]\s*["']?[^"',\s}]{4,}/gi, "NAMED_SECRET", "named_secret");
replace(/\b[A-Z0-9._%+-]+@[A-Z0-9.-]+\.[A-Z]{2,}\b/gi, "EMAIL", "email");
replace(/\b(?:(?:25[0-5]|2[0-4]\d|1?\d?\d)\.){3}(?:25[0-5]|2[0-4]\d|1?\d?\d)\b/g, "IP_ADDRESS", "ip_address");
replace(/\bhttps?:\/\/[^/\s:@]+:[^@\s/]+@[^\s]+/gi, "CREDENTIAL_URL", "credential_url");
text = text.replace(
/\b(?=[A-Za-z0-9_+/.=-]{32,}\b)(?=[A-Za-z0-9_+/.=-]*[A-Za-z])(?=[A-Za-z0-9_+/.=-]*\d)[A-Za-z0-9_+/.=-]+\b/g,
(match) => {
// Git object IDs are essential engineering evidence. Named-secret rules above
// still redact a hex credential when it is attached to a secret-bearing key.
if (/^(?:[0-9a-f]{40}|[0-9a-f]{64})$/i.test(match)) return match;
counts.high_entropy += 1;
return "[REDACTED:HIGH_ENTROPY]";
},
);
if (home) text = text.split(home).join("$HOME");
const originalLength = text.length;
if (text.length > MAX_VISIBLE_TEXT) text = `${text.slice(0, MAX_VISIBLE_TEXT)}…[TRUNCATED]`;
return { text, counts, truncated: originalLength > MAX_VISIBLE_TEXT };
}
function mergeCounts(target, incoming) {
for (const [key, value] of Object.entries(incoming)) target[key] = (target[key] || 0) + value;
}
function textFromResponseContent(content) {
if (typeof content === "string") return content;
if (!Array.isArray(content)) return "";
return content
.filter((item) => item && typeof item === "object" && typeof item.text === "string")
.map((item) => item.text)
.join("\n");
}
function correctionSignals(text) {
const markers = [
"不是", "别搞错", "方向", "纠正", "你忘了", "走偏", "不应该",
"应该是", "我的意思", "我说的是", "不要", "边界",
];
return markers.filter((marker) => text.includes(marker));
}
function visibleEvent(record, securityCounts) {
const timestamp = record.timestamp || record.payload?.timestamp || null;
const payload = record.payload || {};
let role = "";
let kind = "";
let rawText = "";
if (record.type === "event_msg" && payload.type === "user_message") {
role = "user";
kind = "visible_message";
rawText = payload.message || "";
} else if (record.type === "event_msg" && payload.type === "agent_message") {
role = "assistant";
kind = "visible_message";
rawText = payload.message || "";
} else if (record.type === "response_item" && payload.type === "function_call") {
role = "system";
kind = "tool_call";
rawText = payload.name || "unknown_tool";
} else if (record.type === "response_item" && payload.type === "custom_tool_call") {
role = "system";
kind = "tool_call";
rawText = payload.name || "unknown_tool";
} else if (record.type === "event_msg" && payload.type === "task_started") {
return { schema: INDEX_SCHEMA, timestamp, role: "system", kind: "task_started" };
} else if (record.type === "event_msg" && payload.type === "task_complete") {
return { schema: INDEX_SCHEMA, timestamp, role: "system", kind: "task_complete" };
} else if (
record.type === "response_item"
&& payload.type === "message"
&& ["user", "assistant"].includes(payload.role)
) {
role = payload.role;
kind = "visible_message";
rawText = textFromResponseContent(payload.content);
} else {
return null;
}
if (!rawText) return null;
const digest = crypto.createHash("sha256").update(String(rawText)).digest("hex");
const redacted = redactVisibleText(rawText);
mergeCounts(securityCounts, redacted.counts);
const event = {
schema: INDEX_SCHEMA,
timestamp,
role,
kind,
text_sha256: digest,
text_length: String(rawText).length,
text: redacted.text,
truncated: redacted.truncated,
};
if (role === "user") event.correction_signals = correctionSignals(redacted.text);
return event;
}
async function hashFile(file) {
const hash = crypto.createHash("sha256");
const stream = fs.createReadStream(file);
for await (const chunk of stream) hash.update(chunk);
return hash.digest("hex");
}
function sameStat(before, after) {
return before.size === after.size && before.mtimeMs === after.mtimeMs;
}
async function copyAndIndex(sourceFile, stagingFile, stagingIndex) {
const sourceHash = crypto.createHash("sha256");
const securityCounts = {};
const metadata = {
record_count: 0,
invalid_json_count: 0,
type_counts: {},
session_id: null,
cwd_fingerprint: null,
first_timestamp: null,
last_timestamp: null,
visible_event_count: 0,
};
let lineBuffer = "";
const indexHandle = await fs.promises.open(stagingIndex, "wx", 0o600);
const processLine = async (line) => {
if (!line.trim()) return;
metadata.record_count += 1;
try {
const record = JSON.parse(line);
const subtype = record.payload?.type ? `${record.type}:${record.payload.type}` : record.type;
metadata.type_counts[subtype] = (metadata.type_counts[subtype] || 0) + 1;
const timestamp = record.timestamp || record.payload?.timestamp || null;
if (timestamp && !metadata.first_timestamp) metadata.first_timestamp = timestamp;
if (timestamp) metadata.last_timestamp = timestamp;
if (record.type === "session_meta") {
metadata.session_id ||= record.payload?.id || record.payload?.session_id || null;
if (record.payload?.cwd) {
metadata.cwd_fingerprint = crypto.createHash("sha256").update(record.payload.cwd).digest("hex");
}
}
// Hidden reasoning, encrypted content, tool arguments and tool output are never indexed.
if (record.type === "response_item" && record.payload?.type === "reasoning") return;
const event = visibleEvent(record, securityCounts);
if (event) {
await indexHandle.write(`${JSON.stringify(event)}\n`);
metadata.visible_event_count += 1;
}
} catch {
metadata.invalid_json_count += 1;
}
};
const transform = new Transform({
transform(chunk, _encoding, callback) {
sourceHash.update(chunk);
lineBuffer += chunk.toString("utf8");
const lines = lineBuffer.split("\n");
lineBuffer = lines.pop() || "";
(async () => {
for (const line of lines) await processLine(line);
})().then(() => callback(null, chunk), callback);
},
flush(callback) {
(async () => {
if (lineBuffer) await processLine(lineBuffer);
})().then(() => callback(), callback);
},
});
try {
await pipeline(
fs.createReadStream(sourceFile),
transform,
fs.createWriteStream(stagingFile, { flags: "wx", mode: 0o600 }),
);
} finally {
await indexHandle.close();
}
return {
source_sha256: sourceHash.digest("hex"),
security_counts: securityCounts,
metadata,
};
}
async function ensureObject(staging, objectPath, expectedHash) {
await fs.promises.mkdir(path.dirname(objectPath), { recursive: true, mode: 0o700 });
if (fs.existsSync(objectPath)) {
const existingHash = await hashFile(objectPath);
if (existingHash !== expectedHash) throw new Error(`OBJECT_HASH_COLLISION: ${objectPath}`);
await fs.promises.unlink(staging);
return "deduplicated";
}
await fs.promises.rename(staging, objectPath);
const readback = await hashFile(objectPath);
if (readback !== expectedHash) throw new Error(`ARCHIVE_READBACK_FAILED: ${objectPath}`);
return "created";
}
async function archiveOne({ sourceRoot, destination, sourceFile, stableSeconds, stagingRoot }) {
const relative = safeRelative(sourceRoot, sourceFile);
const before = await fs.promises.stat(sourceFile);
const ageSeconds = (Date.now() - before.mtimeMs) / 1000;
if (ageSeconds < stableSeconds) {
return { relative_path: relative, state: "pending_live", size: before.size, mtime: before.mtime.toISOString() };
}
const nonce = crypto.randomUUID();
const staging = path.join(stagingRoot, `${nonce}.jsonl`);
const stagingIndex = path.join(stagingRoot, `${nonce}.index.jsonl`);
const copied = await copyAndIndex(sourceFile, staging, stagingIndex);
const after = await fs.promises.stat(sourceFile);
if (!sameStat(before, after)) {
await fs.promises.rm(staging, { force: true });
await fs.promises.rm(stagingIndex, { force: true });
return { relative_path: relative, state: "pending_changed_during_copy", size: after.size, mtime: after.mtime.toISOString() };
}
const archiveReadbackHash = await hashFile(staging);
if (archiveReadbackHash !== copied.source_sha256) {
throw new Error(`STAGING_READBACK_FAILED: ${relative}`);
}
const shard = copied.source_sha256.slice(0, 2);
const objectRelative = `objects/sha256/${shard}/${copied.source_sha256}.jsonl`;
const indexRelative = `private-index/sha256/${shard}/${copied.source_sha256}.jsonl`;
const objectState = await ensureObject(staging, path.join(destination, objectRelative), copied.source_sha256);
const indexHash = await hashFile(stagingIndex);
const indexState = await ensureObject(stagingIndex, path.join(destination, indexRelative), indexHash);
return {
relative_path: relative,
state: "archived_verified",
size: before.size,
mtime: before.mtime.toISOString(),
sha256: copied.source_sha256,
object: objectRelative,
object_state: objectState,
private_index: indexRelative,
private_index_sha256: indexHash,
private_index_state: indexState,
...copied.metadata,
security_counts: copied.security_counts,
};
}
export async function archiveSessions({
source,
destination,
stableSeconds = DEFAULT_STABLE_SECONDS,
now = new Date(),
onProgress = null,
}) {
const sourceRoot = path.resolve(source);
const destinationRoot = path.resolve(destination);
if (!fs.statSync(sourceRoot).isDirectory()) throw new Error("SOURCE_NOT_DIRECTORY");
await fs.promises.mkdir(destinationRoot, { recursive: true, mode: 0o700 });
const stagingRoot = path.join(destinationRoot, ".staging");
await fs.promises.mkdir(stagingRoot, { recursive: true, mode: 0o700 });
const snapshotId = isoForPath(now);
const files = await listSessionFiles(sourceRoot);
const entries = [];
for (let index = 0; index < files.length; index += 1) {
const sourceFile = files[index];
const entry = await archiveOne({
sourceRoot,
destination: destinationRoot,
sourceFile,
stableSeconds: Number(stableSeconds),
stagingRoot,
});
entries.push(entry);
onProgress?.({ current: index + 1, total: files.length, relative_path: entry.relative_path, state: entry.state });
}
const archived = entries.filter((entry) => entry.state === "archived_verified");
const pending = entries.filter((entry) => entry.state !== "archived_verified");
const manifest = {
schema: SCHEMA,
snapshot_id: snapshotId,
created_at: now.toISOString(),
source_kind: "codex_session_jsonl",
source_root_fingerprint: crypto.createHash("sha256").update(sourceRoot).digest("hex"),
privacy_boundary: {
raw_objects: "PRIVATE_JZAO_ONLY",
private_index: "PRIVATE_JZAO_ONLY_REDACTED_VISIBLE_EVENTS",
repository_publication: "SCHEMA_CODE_AND_REVIEWED_DISTILLATION_ONLY",
hidden_reasoning_indexed: false,
tool_arguments_or_outputs_indexed: false,
},
totals: {
discovered_files: entries.length,
archived_verified: archived.length,
pending_live: pending.length,
archived_bytes: archived.reduce((sum, entry) => sum + entry.size, 0),
},
entries,
};
const snapshotDir = path.join(destinationRoot, "snapshots", snapshotId);
await fs.promises.mkdir(snapshotDir, { recursive: true, mode: 0o700 });
const manifestPath = path.join(snapshotDir, "manifest.json");
await fs.promises.writeFile(manifestPath, `${JSON.stringify(manifest, null, 2)}\n`, { flag: "wx", mode: 0o600 });
const manifestSha256 = await hashFile(manifestPath);
const latest = {
schema: "guanghu.codex-session-cold-archive-pointer/v1",
snapshot_id: snapshotId,
manifest: path.relative(destinationRoot, manifestPath).split(path.sep).join("/"),
manifest_sha256: manifestSha256,
archived_verified: archived.length,
pending_live: pending.length,
updated_at: now.toISOString(),
};
const latestTemp = path.join(destinationRoot, `.LATEST-${crypto.randomUUID()}.json`);
await fs.promises.writeFile(latestTemp, `${JSON.stringify(latest, null, 2)}\n`, { mode: 0o600 });
await fs.promises.rename(latestTemp, path.join(destinationRoot, "LATEST.json"));
await fs.promises.rm(stagingRoot, { recursive: true, force: true });
const appleDoubleRemoved = await removeAppleDouble(destinationRoot);
return { manifest, manifestPath, manifestSha256, latest, appleDoubleRemoved };
}
export async function verifyManifest({ destination, manifest }) {
const destinationRoot = path.resolve(destination);
const manifestPath = path.resolve(manifest);
const parsed = JSON.parse(await fs.promises.readFile(manifestPath, "utf8"));
if (parsed.schema !== SCHEMA) throw new Error("MANIFEST_SCHEMA_MISMATCH");
let checked = 0;
for (const entry of parsed.entries) {
if (entry.state !== "archived_verified") continue;
const objectPath = path.join(destinationRoot, entry.object);
const indexPath = path.join(destinationRoot, entry.private_index);
if (safeRelative(destinationRoot, objectPath) !== entry.object) throw new Error("OBJECT_PATH_MISMATCH");
if (safeRelative(destinationRoot, indexPath) !== entry.private_index) throw new Error("INDEX_PATH_MISMATCH");
if (await hashFile(objectPath) !== entry.sha256) throw new Error(`OBJECT_VERIFY_FAILED: ${entry.relative_path}`);
if (await hashFile(indexPath) !== entry.private_index_sha256) throw new Error(`INDEX_VERIFY_FAILED: ${entry.relative_path}`);
checked += 1;
}
const appleDouble = [];
async function scan(current) {
for (const entry of await fs.promises.readdir(current, { withFileTypes: true })) {
if (entry.name.startsWith("._")) appleDouble.push(path.join(current, entry.name));
if (entry.isDirectory()) await scan(path.join(current, entry.name));
}
}
await scan(destinationRoot);
if (appleDouble.length) throw new Error(`APPLEDOUBLE_PRESENT: ${appleDouble.length}`);
return { ok: true, checked_objects: checked, pending_live: parsed.totals.pending_live, manifest_sha256: await hashFile(manifestPath) };
}
export async function restoreManifest({ destination, manifest, restoreRoot, execute = false }) {
const destinationRoot = path.resolve(destination);
const parsed = JSON.parse(await fs.promises.readFile(path.resolve(manifest), "utf8"));
if (parsed.schema !== SCHEMA) throw new Error("MANIFEST_SCHEMA_MISMATCH");
const plan = [];
for (const entry of parsed.entries) {
if (entry.state !== "archived_verified") continue;
const output = path.join(path.resolve(restoreRoot), entry.relative_path);
safeRelative(path.resolve(restoreRoot), output);
plan.push({ relative_path: entry.relative_path, sha256: entry.sha256, output });
if (!execute) continue;
await fs.promises.mkdir(path.dirname(output), { recursive: true, mode: 0o700 });
if (fs.existsSync(output)) {
if (await hashFile(output) === entry.sha256) continue;
throw new Error(`RESTORE_TARGET_EXISTS: ${output}`);
}
await fs.promises.copyFile(path.join(destinationRoot, entry.object), output, fs.constants.COPYFILE_EXCL);
if (await hashFile(output) !== entry.sha256) throw new Error(`RESTORE_READBACK_FAILED: ${output}`);
}
return { execute, planned_files: plan.length, plan };
}
async function main() {
const args = parseArgs(process.argv);
let result;
if (args.command === "archive") {
result = await archiveSessions({
source: required(args.source, "source"),
destination: required(args.destination, "destination"),
stableSeconds: Number(args.stable_seconds || DEFAULT_STABLE_SECONDS),
onProgress: (progress) => process.stderr.write(
`ARCHIVE_PROGRESS ${progress.current}/${progress.total} ${progress.state} ${progress.relative_path}\n`,
),
});
result = {
manifest: result.manifestPath,
manifest_sha256: result.manifestSha256,
totals: result.manifest.totals,
latest: result.latest,
apple_double_removed: result.appleDoubleRemoved,
};
} else if (args.command === "verify") {
result = await verifyManifest({
destination: required(args.destination, "destination"),
manifest: required(args.manifest, "manifest"),
});
} else if (args.command === "restore") {
result = await restoreManifest({
destination: required(args.destination, "destination"),
manifest: required(args.manifest, "manifest"),
restoreRoot: required(args.restore_root, "restore_root"),
execute: args.execute,
});
} else {
throw new Error("COMMAND_REQUIRED: archive | verify | restore");
}
process.stdout.write(`${JSON.stringify(result, null, 2)}\n`);
}
const isCli = process.argv[1] && path.resolve(process.argv[1]) === fileURLToPath(import.meta.url);
if (isCli) main().catch((error) => {
process.stderr.write(`${error.message}\n`);
process.exitCode = 1;
});