From 86cebc9e8b0ef9b241f41f83052371045cf2ade7 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 4 Oct 2026 03:23:18 +0000 Subject: [PATCH] Wire knowledge records into repo memory behind SHIP_KNOWLEDGE_PROVENANCE (S21) Sidecar provenance records (new ship_knowledge_records table / knowledge-records dir) written alongside notes; shadow logs what retrieval would hide or label without changing context; on hides invalidated/unverifiable notes and labels fresh|stale|unknown. Run-publish notes record the run as source; condenser summaries record as derivation:summary with derivedFrom; Knowledge delete redacts derived records. Default off: no store built, nothing written. Co-Authored-By: Claude Sonnet 5.5 Claude-Session: https://claude.ai/code/session_01VqsBNqvaWezf1DwQrAnVgX --- src/durable.ts | 24 ++ src/knowledge-provenance.test.ts | 450 ++++++++++++++++++++++++++++++ src/knowledge-provenance.ts | 238 ++++++++++++++++ src/knowledge-store.ts | 192 +++++++++++++ src/memory.ts | 13 + src/repo-memory.ts | 48 +++- src/runtime.ts | 6 +- src/scoped-repo-memory.ts | 40 ++- web/src/views/knowledge.server.ts | 4 +- 9 files changed, 997 insertions(+), 18 deletions(-) create mode 100644 src/knowledge-provenance.test.ts create mode 100644 src/knowledge-provenance.ts create mode 100644 src/knowledge-store.ts diff --git a/src/durable.ts b/src/durable.ts index 7770b65..8a307c4 100644 --- a/src/durable.ts +++ b/src/durable.ts @@ -1173,11 +1173,17 @@ export function durableAgent( // Playbook + recent-run notes, recorded so replay never re-reads // a tree or memory that has since changed. let repoContext = ""; + // S21: ids of the notes injected into context, for summary derivedFrom. + // Filled on first execution only (the step's recorded output is a + // string); empty on replay, which recordSummary tolerates (putIfAbsent). + const contextNoteIds: string[] = []; + CONTEXT_NOTES.set(ctx, contextNoteIds); if (checkout !== null && repoKey !== null) { repoContext = await ctx.step("repo-context", () => loadRepoContext(executor, { repo: repoKey, ...(config.repoMemory !== undefined ? { memory: config.repoMemory } : {}), + ...(config.repoMemory?.recordSummary !== undefined ? { onNotes: (ids: string[]) => contextNoteIds.push(...ids) } : {}), }), ); } @@ -1687,6 +1693,9 @@ export function durableAgent( } +/** S21: note ids injected into a run's context, by workflow ctx (see the repo-context step). */ +const CONTEXT_NOTES = new WeakMap(); + /** * The native CodeAct loop as a harness adapter — the current code, * re-entry-pointed. Every recorded step name is unchanged when `stepPrefix` @@ -1910,6 +1919,16 @@ export function nativeAdapter(config: DurableAgentConfig): HarnessAdapter { return summaryStep.text; }, condense, + config.repoMemory?.recordSummary !== undefined && input.repo !== undefined + ? ({ summary }) => + config.repoMemory!.recordSummary!({ + id: `summary-${ws.ctx.runId}-${p}turn-${turn}`, + repo: repoKeyOf(input.repo!, input.repositoryScopeVersion), + runId: ws.ctx.runId, + summary, + derivedFrom: [...(CONTEXT_NOTES.get(ws.ctx) ?? [])], + }) + : undefined, ); } @@ -2680,6 +2699,11 @@ async function publishIfRepoRun( repo: repoKeyOf(repoUrl, input.repositoryScopeVersion), note: runNote({ task: input.task, summary, ...(pr !== undefined ? { pr } : {}) }), runId: ctx.runId, + // S21: where the claim came from — the run, at the revision it pushed. + // Only passed with provenance on, so a plain store sees the old call. + ...(config.repoMemory!.provenanceMode?.() !== undefined && config.repoMemory!.provenanceMode() !== "off" + ? { provenance: { revision: push.kind === "pushed" ? push.sha : "", sourceKind: "run" as const } } + : {}), }) .catch(() => {}); // memory is advisory — never fail a publish over it return true; diff --git a/src/knowledge-provenance.test.ts b/src/knowledge-provenance.test.ts new file mode 100644 index 0000000..c1d4e7d --- /dev/null +++ b/src/knowledge-provenance.test.ts @@ -0,0 +1,450 @@ +import assert from "node:assert/strict"; +import { mkdtemp, readFile, readdir } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { test } from "node:test"; + +import type { AdapterGenerateResult, ModelAdapter, Message } from "@neutron-build/ai"; +import { LocalExecutor } from "@neutron-build/agents"; +import type { AgentExecutor } from "@neutron-build/agents"; +import { MemoryEventStore, executeRun } from "@neutron-build/workflow"; + +import { durableAgent } from "./durable.js"; +import type { ExecutorProvider } from "./durable.js"; +import { KnowledgeProvenance, knowledgeProvenanceMode, knowledgeScope } from "./knowledge-provenance.js"; +import type { KnowledgeMode } from "./knowledge-provenance.js"; +import { createRecord, verify } from "./knowledge-record.js"; +import type { KnowledgeRecord } from "./knowledge-record.js"; +import { FileKnowledgeStore, InMemoryKnowledgeStore, NucleusKnowledgeStore } from "./knowledge-store.js"; +import type { KnowledgeStore } from "./knowledge-store.js"; +import { condenseIfNeeded } from "./memory.js"; +import type { NucleusPgwire } from "./nucleus-pgwire.js"; +import { FileRepoMemory, loadRepoContext } from "./repo-memory.js"; +import { ScopedRepoMemory } from "./scoped-repo-memory.js"; + +const REPO = "https://github.com/tyler/app"; +const REPO_URL = "https://github.com/tyler/app"; + +/** A fake Nucleus that understands exactly the statements NucleusKnowledgeStore issues. */ +function fakeNucleus(options: { failList?: boolean } = {}): { db: NucleusPgwire; rows: Record[]; sql: string[] } { + const rows: Record[] = []; + const sql: string[] = []; + const db = { + async query(text: string, params: unknown[] = []): Promise[]> { + sql.push(text); + if (/^CREATE TABLE IF NOT EXISTS ship_knowledge_records/.test(text)) return []; + if (/^INSERT INTO ship_knowledge_records/.test(text)) { + rows.push({ record_id: params[0], repo: params[1], project: params[2], record: params[3] }); + return []; + } + if (/^DELETE FROM ship_knowledge_records WHERE record_id/.test(text)) { + for (let i = rows.length - 1; i >= 0; i--) if (rows[i]!.record_id === params[0]) rows.splice(i, 1); + return []; + } + if (/^SELECT record FROM ship_knowledge_records WHERE record_id/.test(text)) return rows.filter((r) => r.record_id === params[0]).map((r) => ({ record: r.record })); + if (/^SELECT record FROM ship_knowledge_records WHERE repo/.test(text)) { + if (options.failList) throw new Error("store down"); + return rows.filter((r) => r.repo === params[0]).map((r) => ({ record: r.record })); + } + throw new Error(`fakeNucleus cannot parse: ${text}`); + }, + } as unknown as NucleusPgwire; + return { db, rows, sql }; +} + +type StoreFactory = { name: string; make: () => Promise }; +const factories: StoreFactory[] = [ + { name: "in-memory", make: async () => new InMemoryKnowledgeStore() }, + { name: "file", make: async () => new FileKnowledgeStore(await mkdtemp(join(tmpdir(), "kn-store-"))) }, + { name: "fake nucleus", make: async () => new NucleusKnowledgeStore(fakeNucleus().db) }, +]; + +const sample = (id: string, extra: Partial = {}): KnowledgeRecord => ({ + ...createRecord({ + id, + kind: "hypothesis", + statement: `statement ${id}`, + source: { kind: "run", repo: REPO, revision: "abc", runId: "run-1" }, + scope: knowledgeScope(REPO), + createdAt: "2026-10-04T00:00:00.000Z", + }), + ...extra, +}); + +for (const f of factories) { + test(`${f.name} knowledge store: put/get/list/remove, putIfAbsent never overwrites, repos stay apart`, async () => { + const store = await f.make(); + await store.put(sample("a")); + assert.equal(await store.putIfAbsent(sample("a", { statement: "OVERWRITE" })), false); + assert.equal(await store.putIfAbsent(sample("b")), true); + assert.equal((await store.get("a"))?.statement, "statement a"); + await store.put(sample("c", { scope: knowledgeScope("github.com/other/repo") })); + assert.deepEqual((await store.list(REPO)).map((r) => r.id).sort(), ["a", "b"]); + await store.put(sample("a", { statement: "replaced" })); + assert.equal((await store.list(REPO)).filter((r) => r.id === "a").length, 1, "put replaces, never duplicates"); + assert.equal((await store.get("a"))?.statement, "replaced"); + await store.remove("a"); + assert.equal(await store.get("a"), undefined); + assert.deepEqual((await store.list(REPO)).map((r) => r.id), ["b"]); + }); +} + +test("the nucleus store creates only a NEW table and never touches ship_memory", async () => { + const { db, sql } = fakeNucleus(); + const store = new NucleusKnowledgeStore(db); + await store.put(sample("a")); + await store.list(REPO); + assert.ok(sql.every((s) => /ship_knowledge_records/.test(s)), "every statement targets the new table"); + assert.ok(!sql.some((s) => /ALTER|DROP|ship_memory/i.test(s))); +}); + +test("a failed table ensure is retried, not cached", async () => { + let creates = 0; + const db = { + async query(text: string): Promise[]> { + if (text.startsWith("CREATE TABLE")) { + creates++; + if (creates === 1) throw new Error("transient"); + } + return []; + }, + } as unknown as NucleusPgwire; + const store = new NucleusKnowledgeStore(db); + await assert.rejects(store.list(REPO), /transient/); + assert.deepEqual(await store.list(REPO), []); + assert.equal(creates, 2); +}); + +test("knowledgeProvenanceMode: off unless explicitly shadow or on", () => { + assert.equal(knowledgeProvenanceMode({}), "off"); + assert.equal(knowledgeProvenanceMode({ SHIP_KNOWLEDGE_PROVENANCE: "" }), "off"); + assert.equal(knowledgeProvenanceMode({ SHIP_KNOWLEDGE_PROVENANCE: "bogus" }), "off"); + assert.equal(knowledgeProvenanceMode({ SHIP_KNOWLEDGE_PROVENANCE: "shadow" }), "shadow"); + assert.equal(knowledgeProvenanceMode({ SHIP_KNOWLEDGE_PROVENANCE: " ON " }), "on"); +}); + +async function rig(mode: KnowledgeMode | "none", opts: { store?: KnowledgeStore; plain?: boolean } = {}) { + const dir = await mkdtemp(join(tmpdir(), "kn-rig-")); + const raw = new FileRepoMemory(join(dir, "memory")); + const store = opts.store ?? new InMemoryKnowledgeStore(); + const logs: string[] = []; + const provenance = mode === "none" ? undefined : new KnowledgeProvenance({ mode, store, log: (l) => logs.push(l) }); + const memory = new ScopedRepoMemory(raw, new MemoryEventStore(), provenance); + return { dir, raw, store, logs, memory }; +} + +test("default off: the stored note is byte-identical, no provenance is written, retrieval is untouched", async () => { + const bare = await rig("none"); + const off = await rig("off"); + for (const r of [bare, off]) { + await r.memory.record({ repo: REPO_URL, note: "n1", runId: "run-1", provenance: { revision: "abc" } }); + } + const file = async (r: { dir: string }) => { + const names = await readdir(join(r.dir, "memory")); + return (await readFile(join(r.dir, "memory", names[0]!), "utf8")).replace(/"noteId":"[^"]+"/, "").replace(/"createdAt":"[^"]+"/, ""); + }; + assert.equal(await file(off), await file(bare)); + assert.ok(!(await file(off)).includes("provenance") && !(await file(off)).includes("abc"), "the hint never reaches the note"); + assert.deepEqual(await off.store.list(REPO), [], "mode off writes nothing"); + assert.equal(off.memory.provenanceMode(), "off"); + const ctx = { head: "zzz" }; + const a = await bare.memory.recent(REPO_URL, 5, { context: ctx }); + const b = await off.memory.recent(REPO_URL, 5, { context: ctx }); + assert.deepEqual(a.map((n) => [n.note, n.freshness]), b.map((n) => [n.note, n.freshness])); + assert.ok(b.every((n) => n.freshness === undefined)); + assert.deepEqual(off.logs, []); +}); + +test("shadow: writes a hypothesis record alongside the note (run as source), note bytes unchanged", async () => { + const s = await rig("shadow"); + const note = await s.memory.record({ repo: REPO_URL, note: "task → PR 1. done", runId: "run-9", provenance: { revision: "sha1" } }); + const rec = await s.store.get(note.noteId); + assert.ok(rec); + assert.equal(rec.kind, "hypothesis", "a model-written note is never born a fact"); + assert.deepEqual(rec.source, { kind: "run", repo: REPO, revision: "sha1", runId: "run-9" }); + assert.deepEqual(rec.scope, { repo: REPO, project: "" }, "project is optional: repo-only scope"); + assert.deepEqual(rec.derivedFrom, []); + const stored = (await s.raw.recent(REPO, 5))[0]!; + assert.deepEqual(Object.keys(stored).sort(), ["createdAt", "note", "noteId", "repo", "runId"]); + // a dashboard note has no run: its source is a human, not a run + const manual = await s.memory.record({ repo: REPO_URL, note: "manual" }); + assert.equal((await s.store.get(manual.noteId))?.source.kind, "human"); +}); + +test("a provenance write failure never fails recording the note", async () => { + const broken: KnowledgeStore = { + put: async () => { throw new Error("boom"); }, + putIfAbsent: async () => { throw new Error("boom"); }, + get: async () => undefined, + list: async () => { throw new Error("boom"); }, + remove: async () => undefined, + }; + const s = await rig("on", { store: broken }); + const note = await s.memory.record({ repo: REPO_URL, note: "kept", runId: "r" }); + assert.equal((await s.raw.recent(REPO, 5))[0]?.noteId, note.noteId); +}); + +async function seeded(mode: KnowledgeMode) { + const s = await rig(mode); + const ids: Record = {}; + for (const [name, revision] of [["current", "HEAD1"], ["old", "OLD"], ["norev", ""], ["poison", "HEAD1"]] as const) { + ids[name] = (await s.memory.record({ repo: REPO_URL, note: name, runId: `run-${name}`, provenance: { revision } })).noteId; + } + // a note recorded before the flag: no provenance record at all + ids.legacy = (await s.raw.record({ repo: REPO, note: "legacy" })).noteId; + const poisoned = (await s.store.get(ids.poison!))!; + await s.store.put({ ...poisoned, invalidated: { reason: "deleted source" } }); + return { ...s, ids }; +} + +test("on: invalidated notes are hidden, the rest labelled fresh|stale|unknown, legacy shown as unknown", async () => { + const s = await seeded("on"); + const notes = await s.memory.recent(REPO_URL, 10, { context: { head: "HEAD1" } }); + const label = Object.fromEntries(notes.map((n) => [n.note, n.freshness])); + assert.deepEqual(label, { current: "fresh", old: "stale", norev: "unknown", legacy: "unknown" }); + assert.ok(!("poison" in label), "invalidated provenance is never served"); +}); + +test("on: freshness without a resolvable head is unknown, never fresh", async () => { + const s = await seeded("on"); + const notes = await s.memory.recent(REPO_URL, 10, { context: {} }); + assert.ok(notes.every((n) => n.freshness === "unknown")); +}); + +test("on: the dashboard listing (no context) is never filtered, so a hidden note stays deletable", async () => { + const s = await seeded("on"); + const all = await s.memory.recent(REPO_URL, 10); + assert.equal(all.length, 5); + assert.ok(all.every((n) => n.freshness === undefined)); +}); + +test("on: a record outside the viewer's project is hidden; a derived record whose parent is missing is hidden", async () => { + const s = await rig("on"); + const n = await s.memory.record({ repo: REPO_URL, note: "scoped elsewhere", runId: "r", provenance: { revision: "H", project: "proj-a" } }); + const d = await s.memory.record({ repo: REPO_URL, note: "orphan", runId: "r2", provenance: { revision: "H" } }); + await s.store.put({ ...(await s.store.get(d.noteId))!, derivedFrom: ["missing-parent"] }); + assert.ok(n.noteId); + assert.deepEqual(await s.memory.recent(REPO_URL, 10, { context: { head: "H" } }), [], "repo-only viewer sees neither"); +}); + +test("on: an unreadable provenance store fails closed (no notes); shadow fails open and logs", async () => { + const { db } = fakeNucleus({ failList: true }); + for (const [mode, expected] of [["on", 0], ["shadow", 1]] as const) { + const s = await rig(mode, { store: new NucleusKnowledgeStore(db) }); + await s.raw.record({ repo: REPO, note: "x" }); + const notes = await s.memory.recent(REPO_URL, 5, { context: { head: "H" } }); + assert.equal(notes.length, expected, mode); + assert.ok(s.logs.some((l) => l.includes("knowledge-provenance-error")), `${mode} logs the failure`); + } +}); + +test("shadow: retrieval is IDENTICAL to off, and the log says what on would have done", async () => { + const shadow = await seeded("shadow"); + const off = await rig("none"); + for (const n of ["current", "old", "norev", "poison"]) await off.raw.record({ repo: REPO, note: n }); + await off.raw.record({ repo: REPO, note: "legacy" }); + const a = (await shadow.memory.recent(REPO_URL, 10, { context: { head: "HEAD1" } })).map((n) => [n.note, n.freshness]); + const b = (await off.memory.recent(REPO_URL, 10, { context: { head: "HEAD1" } })).map((n) => [n.note, n.freshness]); + assert.deepEqual(a.sort(), b.sort()); + assert.ok(a.some(([note]) => note === "poison"), "the invalidated note is still served in shadow"); + const decision = JSON.parse(shadow.logs.at(-1)!); + assert.equal(decision.event, "knowledge-provenance"); + assert.equal(decision.mode, "shadow"); + assert.deepEqual(decision.wouldHide.map((h: { noteId: string }) => h.noteId), [shadow.ids.poison]); + assert.match(decision.wouldHide[0].reason, /invalidated: deleted source/); + assert.equal(decision.wouldLabel[shadow.ids.current!], "fresh"); + assert.equal(decision.wouldLabel[shadow.ids.old!], "stale"); + assert.deepEqual(decision.unrecorded, [shadow.ids.legacy]); +}); + +function execStub(head: string | null): AgentExecutor { + return { + exec: async (cmd: string) => { + if (cmd.startsWith("git rev-parse")) return head === null ? { exitCode: 128, stdout: "", stderr: "x" } : { exitCode: 0, stdout: `${head}\n`, stderr: "" }; + return { exitCode: 1, stdout: "", stderr: "" }; + }, + } as unknown as AgentExecutor; +} + +test("loadRepoContext: labels only with on; shadow and off produce the same bytes as before", async () => { + const contexts: Record = {}; + for (const mode of ["none", "off", "shadow", "on"] as const) { + const s = await seeded(mode === "none" ? "off" : mode); + const memory = mode === "none" ? new ScopedRepoMemory(s.raw, new MemoryEventStore()) : s.memory; + contexts[mode] = (await loadRepoContext(execStub("HEAD1"), { repo: REPO_URL, memory })).replace(/\[\d{4}-\d\d-\d\d\]/g, "[d]"); + } + assert.equal(contexts.off, contexts.none); + assert.equal(contexts.shadow, contexts.none.replace(/^/, ""), "shadow does not change what the model reads"); + assert.ok(!/\((fresh|stale|unknown)\)/.test(contexts.shadow!)); + assert.match(contexts.on!, /\[d\] \(fresh\) current/); + assert.match(contexts.on!, /\[d\] \(stale\) old/); + assert.ok(!contexts.on!.includes("poison")); +}); + +test("loadRepoContext in on mode: no git head means every note reads unknown", async () => { + const s = await seeded("on"); + const text = await loadRepoContext(execStub(null), { repo: REPO_URL, memory: s.memory }); + assert.ok(!/\((fresh|stale)\)/.test(text)); + assert.match(text, /\(unknown\) current/); +}); + +test("redaction: deleting a note removes its record, invalidates summaries, deletes embeddings; they are never served", async () => { + const s = await rig("on"); + const n = await s.memory.record({ repo: REPO_URL, note: "secret", runId: "r1", provenance: { revision: "H" } }); + const keep = await s.memory.record({ repo: REPO_URL, note: "other", runId: "r2", provenance: { revision: "H" } }); + await s.memory.recordSummary({ id: "summary-1", repo: REPO_URL, runId: "r3", summary: "recap of secret", derivedFrom: [n.noteId] }); + await s.store.put({ ...createRecord({ id: "emb-1", kind: "hypothesis", statement: "vec", source: { kind: "agent-claim", repo: REPO, revision: "" }, scope: knowledgeScope(REPO), createdAt: "t", derivedFrom: [n.noteId], derivation: "embedding" }) }); + + await s.memory.remove(n.noteId, REPO_URL); + + assert.deepEqual((await s.raw.recent(REPO, 10)).map((x) => x.noteId), [keep.noteId], "the note itself is gone"); + assert.equal(await s.store.get(n.noteId), undefined); + assert.equal(await s.store.get("emb-1"), undefined, "embeddings cannot be edited down: deleted"); + const summary = await s.store.get("summary-1"); + assert.ok(summary?.invalidated, "summary kept for audit but flagged"); + assert.ok((await s.store.get(keep.noteId)) !== undefined, "unrelated records untouched"); + const served = await s.memory.recent(REPO_URL, 10, { context: { head: "H" } }); + assert.deepEqual(served.map((x) => x.note), ["other"]); +}); + +test("redaction cascades from a note with NO provenance record when the repo is given", async () => { + const s = await rig("on"); + const legacy = await s.raw.record({ repo: REPO, note: "predates the flag" }); + await s.memory.recordSummary({ id: "summary-L", repo: REPO_URL, runId: "r", summary: "recap", derivedFrom: [legacy.noteId] }); + await s.memory.remove(legacy.noteId, REPO_URL); + assert.ok((await s.store.get("summary-L"))?.invalidated); +}); + +test("redaction runs in shadow too (a shadow record holds the deleted text); a failure is logged, not thrown", async () => { + const s = await rig("shadow"); + const n = await s.memory.record({ repo: REPO_URL, note: "gone", runId: "r" }); + await s.memory.remove(n.noteId, REPO_URL); + assert.equal(await s.store.get(n.noteId), undefined); + + const failing: KnowledgeStore = { ...new InMemoryKnowledgeStore(), put: async () => undefined, putIfAbsent: async () => true, get: async () => { throw new Error("down"); }, list: async () => [], remove: async () => undefined }; + const sh = await rig("shadow", { store: failing }); + const m = await sh.raw.record({ repo: REPO, note: "k" }); + await sh.memory.remove(m.noteId, REPO_URL); + assert.deepEqual(await sh.raw.recent(REPO, 5), [], "the note is still deleted"); + const on = await rig("on", { store: failing }); + const m2 = await on.raw.record({ repo: REPO, note: "k2" }); + await assert.rejects(on.memory.remove(m2.noteId, REPO_URL), /down/, "on does not swallow a redaction failure"); +}); + +test("redaction does nothing with provenance off", async () => { + const s = await rig("off"); + const n = await s.memory.record({ repo: REPO_URL, note: "x" }); + await s.memory.remove(n.noteId, REPO_URL); + assert.deepEqual(await s.raw.recent(REPO, 5), []); +}); + +test("a condenser summary is a derived agent claim: a hypothesis, idempotent, and unusable as evidence", async () => { + const s = await rig("shadow"); + const n = await s.memory.record({ repo: REPO_URL, note: "src", runId: "r1" }); + await s.memory.recordSummary({ id: "summary-r-1", repo: REPO_URL, runId: "r", summary: "first", derivedFrom: [n.noteId] }); + await s.memory.recordSummary({ id: "summary-r-1", repo: REPO_URL, runId: "r", summary: "replay", derivedFrom: [] }); + const sum = (await s.store.get("summary-r-1"))!; + assert.equal(sum.derivation, "summary"); + assert.equal(sum.kind, "hypothesis"); + assert.equal(sum.source.kind, "agent-claim"); + assert.deepEqual(sum.derivedFrom, [n.noteId]); + assert.equal(sum.statement, "first", "a replayed run does not overwrite the first record"); + // it cannot verify another claim: evidence that cites a record is refused + const target = (await s.store.get(n.noteId))!; + const res = verify(target, [{ kind: "run", repo: REPO, revision: "x", runId: "other", recordId: sum.id }]); + assert.equal(res.verified, false); + // nor can it be promoted without evidence + assert.equal(verify(sum, []).verified, false); +}); + +test("condenseIfNeeded: the observer sees the summary; messages are identical with or without it; a throw is swallowed", async () => { + const messages: Message[] = [ + { role: "system", content: "sys" }, + { role: "user", content: "task" }, + ...Array.from({ length: 30 }, (_, i) => ({ role: i % 2 ? "assistant" : "user", content: `turn ${i} ${"x".repeat(200)}` }) as Message), + ]; + const cfg = { maxTokens: 100, keepRecent: 4, maxSummaryLayers: 3 }; + const summarize = async () => "RECAP"; + const plain = await condenseIfNeeded(messages, summarize, cfg); + const seen: { summary: string; condensed: number }[] = []; + const observed = await condenseIfNeeded(messages, summarize, cfg, (i) => { seen.push(i); }); + assert.deepEqual(observed, plain); + assert.deepEqual(seen, [{ summary: "RECAP", condensed: 26 }]); + const thrown = await condenseIfNeeded(messages, summarize, cfg, () => { throw new Error("nope"); }); + assert.deepEqual(thrown, plain); + const under = await condenseIfNeeded(messages, summarize, { ...cfg, maxTokens: 10_000_000 }, (i) => { seen.push(i); }); + assert.equal(under, messages); + assert.equal(seen.length, 1, "no summary, no callback"); +}); + +// ---- the real call path: a durable repo run --------------------------------- + +async function repoRun(mode: KnowledgeMode | "none") { + const bareDir = await mkdtemp(join(tmpdir(), "kn-bare-")); + const seedDir = await mkdtemp(join(tmpdir(), "kn-seed-")); + const seeder = new LocalExecutor({ root: seedDir }); + await seeder.exec( + `git init -q -b main . && git config user.email t@t && git config user.name t && printf 'Run tests with make check-42.\\n' > SHIP.md && git add -A && git commit -qm seed && git clone -q --bare . ${bareDir}/owner/repo.git`, + ); + const r = await rig(mode); + await r.memory.record({ repo: "file:///owner/repo", note: "previously fixed the parser", runId: "old-run", provenance: { revision: "0000000" } }); + const prompts: string[] = []; + const model: ModelAdapter = { + provider: "scripted", + modelId: "s1", + async doGenerate(options): Promise { + prompts.push(String(options.messages[0]?.content ?? "")); + const finishes = options.messages.filter((m) => typeof m.content === "string" && m.content.includes("Before finishing")).length; + const asst = options.messages.filter((m) => m.role === "assistant").length; + const text = asst === 0 ? "```bash\ntrue\n```" : finishes > 0 ? "```finish\nnothing to change\n```" : "```bash\ntrue\n```"; + return { content: [{ type: "text", text }], finishReason: "stop", usage: { inputTokens: 1, outputTokens: 1, totalTokens: 2 }, raw: null }; + }, + async *doStream() { throw new Error("unused"); }, + }; + const work = await mkdtemp(join(tmpdir(), "kn-work-")); + const provider: ExecutorProvider = { async create() { return { handle: work }; }, attach: (handle: string) => new LocalExecutor({ root: handle }) }; + const wf = durableAgent({ model, executor: provider, workdir: ".", repoMemory: r.memory }); + const store = new MemoryEventStore(); + const outcome = await executeRun({ workflow: wf, runId: "run-kn", store, input: { task: "check the build", repo: `file://${bareDir}/owner/repo.git` } }); + assert.equal(outcome.status, "completed"); + const steps = (await store.load("run-kn")).filter((e) => e.type === "step-completed").map((e) => (e as { name: string }).name); + return { ...r, prompts, steps }; +} + +test("durable repo run: off, shadow and on share one step sequence; shadow's prompt equals off's; publish records the run as source", async () => { + const off = await repoRun("none"); + const shadow = await repoRun("shadow"); + const on = await repoRun("on"); + assert.deepEqual(shadow.steps, off.steps, "durable step names and order are unchanged"); + assert.deepEqual(on.steps, off.steps); + assert.equal(shadow.prompts[0], off.prompts[0], "shadow never changes what the model reads"); + assert.match(on.prompts[0]!, /previously fixed the parser/); + assert.match(on.prompts[0]!, /\((fresh|stale|unknown)\) previously fixed the parser/); + assert.ok(!/\((fresh|stale|unknown)\) previously/.test(shadow.prompts[0]!)); + + for (const r of [shadow, on]) { + const recs = await r.store.list("file:///owner/repo"); + const published = recs.find((x) => x.source.runId === "run-kn"); + assert.ok(published, "the publish note has a provenance record"); + assert.equal(published.kind, "hypothesis"); + assert.equal(published.source.kind, "run"); + } + assert.deepEqual(await off.store.list("file:///owner/repo"), [], "off writes no records"); + assert.ok(shadow.logs.some((l) => l.includes('"event":"knowledge-provenance"')), "shadow logged its would-be decision"); +}); + +test("screen() itself never alters the list in shadow, even for invalidated or unrecorded notes", async () => { + const store = new InMemoryKnowledgeStore(); + const logs: string[] = []; + const shadow = new KnowledgeProvenance({ mode: "shadow", store, log: (l) => logs.push(l) }); + const base = { repo: REPO, createdAt: "t" }; + await store.put({ ...sample("bad"), invalidated: { reason: "r" } }); + const notes = [ + { ...base, noteId: "bad", note: "bad" }, + { ...base, noteId: "unrecorded", note: "u" }, + ]; + const out = await shadow.screen(REPO, notes, { head: "abc" }); + assert.deepEqual(out, notes); + assert.ok(out.every((n) => !("freshness" in n)), "no label leaks into shadow output"); + assert.equal(logs.length, 1); +}); diff --git a/src/knowledge-provenance.ts b/src/knowledge-provenance.ts new file mode 100644 index 0000000..950f228 --- /dev/null +++ b/src/knowledge-provenance.ts @@ -0,0 +1,238 @@ +import { createRecord, freshness, redactionPlan, applyRedactionPlan, visibleTo } from "./knowledge-record.js"; +import type { Freshness, KnowledgeRecord, KnowledgeScope, KnowledgeSource, CurrentRevisionOf } from "./knowledge-record.js"; +import type { KnowledgeStore } from "./knowledge-store.js"; +import { canonicalRepositoryURL } from "./repository-reference.js"; +import type { RepoNote } from "./repo-memory.js"; + +/** + * S21 wiring: knowledge-record.ts applied to repo memory, behind + * SHIP_KNOWLEDGE_PROVENANCE. + * + * off (default) nothing here runs. No store is built, nothing is written. + * shadow provenance records are written ALONGSIDE the existing notes, and + * retrieval logs what it WOULD have hidden or labelled. What a run + * retrieves is never changed. + * on context retrieval hides notes whose provenance is invalidated or + * unverifiable and labels the rest fresh|stale|unknown. + * + * PROJECT: Ship's repo-memory scope key (`repoKeyOf`) is repo-only; a project + * is a separate record that can own several repos and has no hand in a note's + * key. KnowledgeScope.project is required by the record module, so it is + * OPTIONAL here: absent means REPO_ONLY_PROJECT (""), which matches exactly + * itself, as visibleTo requires. Passing a real project later partitions + * records without a schema change. Nothing infers a project from a repo. + * + * Redaction is not mode-gated beyond "off": a shadow record holds the note's + * text, so deleting the note must delete (or invalidate) its record and + * anything derived from it, or "shadow" would quietly retain deleted content. + */ +export type KnowledgeMode = "off" | "shadow" | "on"; + +export function knowledgeProvenanceMode(env: NodeJS.ProcessEnv = process.env): KnowledgeMode { + const raw = (env.SHIP_KNOWLEDGE_PROVENANCE ?? "").trim().toLowerCase(); + return raw === "shadow" ? "shadow" : raw === "on" || raw === "1" || raw === "true" ? "on" : "off"; +} + +/** The project component for a note recorded without a project (see header). */ +export const REPO_ONLY_PROJECT = ""; + +export const knowledgeScope = (repo: string, project: string = REPO_ONLY_PROJECT): KnowledgeScope => ({ + repo: canonicalRepositoryURL(repo) ?? repo, + project, +}); + +/** What a caller knows about where a note's claim came from. All optional. */ +export interface NoteProvenanceHint { + /** The revision the claim is about (a run records the sha it pushed). */ + revision?: string; + branch?: string; + /** Overrides the inferred source kind (a run => 'run', no run => 'human'). */ + sourceKind?: KnowledgeSource["kind"]; + by?: string; + project?: string; +} + +/** What retrieval knows about the repo right now, so freshness can be judged. */ +export interface RetrievalContext { + head?: string; + branch?: string; +} + +export interface ShadowDecision { + event: "knowledge-provenance"; + mode: KnowledgeMode; + repo: string; + /** Notes that `on` would have hidden, with the reason. */ + wouldHide: { noteId: string; reason: string }[]; + /** Label `on` would have attached, by note id. */ + wouldLabel: Record; + /** Notes with no provenance record (written before the flag): shown, labelled unknown. */ + unrecorded: string[]; +} + +export interface KnowledgeProvenanceOptions { + mode: KnowledgeMode; + store: KnowledgeStore; + /** One JSON line per retrieval that had something to report. */ + log?: (line: string) => void; + now?: () => string; +} + +export class KnowledgeProvenance { + readonly mode: KnowledgeMode; + #store: KnowledgeStore; + #log: (line: string) => void; + #now: () => string; + + constructor(options: KnowledgeProvenanceOptions) { + this.mode = options.mode; + this.#store = options.store; + this.#log = options.log ?? ((line) => console.error(line)); + this.#now = options.now ?? (() => new Date().toISOString()); + } + + /** + * Write the provenance record for a note that was just recorded. + * + * A note is a HYPOTHESIS: its text is a model-written summary of a run, so + * "the run is the source" says where it came from, not that it is true. It + * becomes a fact only through verify() with independent evidence. + */ + async recordNote(note: RepoNote, hint: NoteProvenanceHint = {}): Promise { + if (this.mode === "off") return; + const scope = knowledgeScope(note.repo, hint.project); + const kind = hint.sourceKind ?? (note.runId !== undefined ? "run" : "human"); + const source: KnowledgeSource = { + kind, + repo: scope.repo, + revision: hint.revision ?? "", + ...(hint.branch !== undefined ? { branch: hint.branch } : {}), + ...(note.runId !== undefined ? { runId: note.runId } : {}), + ...(kind === "human" ? { by: hint.by ?? "dashboard" } : hint.by !== undefined ? { by: hint.by } : {}), + }; + await this.#store.put( + createRecord({ id: note.noteId, kind: "hypothesis", statement: note.note, source, scope, createdAt: note.createdAt }), + ); + } + + /** + * A condenser summary: derivation 'summary', an agent claim, never evidence. + * Idempotent by id — a replayed run reaches the same call again and must not + * overwrite the first record (which may know more of derivedFrom). + */ + async recordSummary(input: { + id: string; + repo: string; + runId: string; + summary: string; + derivedFrom: string[]; + project?: string; + }): Promise { + if (this.mode === "off") return; + const scope = knowledgeScope(input.repo, input.project); + await this.#store.putIfAbsent( + createRecord({ + id: input.id, + kind: "hypothesis", + statement: input.summary, + source: { kind: "agent-claim", repo: scope.repo, revision: "", runId: input.runId, by: "condenser" }, + scope, + createdAt: this.#now(), + derivedFrom: input.derivedFrom, + derivation: "summary", + }), + ); + } + + /** + * Decide what retrieval would do with `notes`. `on` returns the filtered, + * labelled notes; `shadow` returns them untouched and logs the decision. + * A provenance store that cannot be read: `on` fails closed (no notes), + * `shadow` fails open and says so. + */ + async screen(repo: string, notes: RepoNote[], ctx: RetrievalContext = {}, project?: string): Promise { + if (this.mode === "off" || notes.length === 0) return notes; + const scope = knowledgeScope(repo, project); + let records: KnowledgeRecord[]; + try { + records = await this.#store.list(scope.repo); + } catch (error) { + this.#log(JSON.stringify({ event: "knowledge-provenance-error", mode: this.mode, repo: scope.repo, error: String(error) })); + return this.mode === "on" ? [] : notes; + } + const visible = new Set(visibleTo(records, scope).map((r) => r.id)); + const byId = new Map(records.map((r) => [r.id, r])); + // Without a resolvable head freshness is `unknown`, never `fresh`. + const current: CurrentRevisionOf = () => + ctx.head === undefined || ctx.head === "" + ? null + : { head: ctx.head, ...(ctx.branch !== undefined ? { branch: ctx.branch } : {}), exists: true }; + + const kept: RepoNote[] = []; + const decision: ShadowDecision = { event: "knowledge-provenance", mode: this.mode, repo: scope.repo, wouldHide: [], wouldLabel: {}, unrecorded: [] }; + for (const note of notes) { + const record = byId.get(note.noteId); + if (record === undefined) { + decision.unrecorded.push(note.noteId); + decision.wouldLabel[note.noteId] = "unknown"; + kept.push({ ...note, freshness: "unknown" }); + continue; + } + if (!visible.has(record.id)) { + decision.wouldHide.push({ noteId: note.noteId, reason: record.invalidated ? `invalidated: ${record.invalidated.reason}` : "outside scope or provenance unverifiable" }); + continue; + } + const label = freshness(record, current).freshness; + decision.wouldLabel[note.noteId] = label; + kept.push({ ...note, freshness: label }); + } + if (this.mode === "shadow") { + this.#log(JSON.stringify(decision)); + return notes; + } + return kept; + } + + /** + * Delete notes' provenance: the records themselves, and everything derived + * from them (summaries invalidated, embeddings/exports deleted). Returns the + * ids the plan DELETED (other than the requested ones), so a caller that + * also holds those as notes can remove them too. + */ + async redact(noteIds: readonly string[], repoHint?: string): Promise<{ deleted: string[]; invalidated: string[] }> { + if (this.mode === "off") return { deleted: [], invalidated: [] }; + const found: KnowledgeRecord[] = []; + const repos = new Set(repoHint !== undefined && repoHint !== "" ? [knowledgeScope(repoHint).repo] : []); + for (const id of noteIds) { + const r = await this.#store.get(id); + if (r !== undefined) repos.add(r.scope.repo); + } + for (const repo of repos) found.push(...(await this.#store.list(repo))); + // A requested id with no record (a note from before the flag) can still be + // the parent of a summary; a stub makes it a root so the cascade runs. + const have = new Set(found.map((r) => r.id)); + const stubs: KnowledgeRecord[] = []; + for (const id of noteIds) { + if (have.has(id)) continue; + for (const repo of repos) { + stubs.push({ ...createRecord({ id, kind: "hypothesis", statement: "", source: { kind: "document", repo, revision: "" }, scope: { repo, project: REPO_ONLY_PROJECT }, createdAt: this.#now() }) }); + } + } + const plan = redactionPlan([...found, ...stubs], { deleteIds: noteIds }); + const real = new Set(found.map((r) => r.id)); + const after = new Map(applyRedactionPlan(found, plan).map((r) => [r.id, r])); + const deleted: string[] = []; + const invalidated: string[] = []; + for (const step of plan.steps) { + if (!real.has(step.id)) continue; + if (step.action === "delete") { + await this.#store.remove(step.id); + if (!noteIds.includes(step.id)) deleted.push(step.id); + } else { + await this.#store.put(after.get(step.id)!); + invalidated.push(step.id); + } + } + return { deleted, invalidated }; + } +} diff --git a/src/knowledge-store.ts b/src/knowledge-store.ts new file mode 100644 index 0000000..875921d --- /dev/null +++ b/src/knowledge-store.ts @@ -0,0 +1,192 @@ +import { readFile, readdir } from "node:fs/promises"; +import { join } from "node:path"; + +import { withFileLock, writeTextFile } from "./file-store.js"; +import type { KnowledgeRecord } from "./knowledge-record.js"; +import type { NucleusPgwire } from "./nucleus-pgwire.js"; +import { stateDir } from "./run-store.js"; + +/** + * Where S21 knowledge records live: a SIDECAR to ship_memory / the repo-memory + * JSONL, never a change to them. + * + * A note keeps its existing row/line byte for byte. Its provenance (source, + * derivedFrom, scope.project, corrections, verification, invalidated) is a + * KnowledgeRecord stored under the same id in a NEW table (Nucleus) or a NEW + * directory (file store). Nucleus cannot safely ALTER a populated table + * (migrations.ts), so a new table created with CREATE TABLE IF NOT EXISTS is + * the additive route: it needs no migration, holds no pre-existing rows, and + * dropping it loses only provenance, never a note. + * + * Records are keyed by `id`. For a note that is the note's noteId; a condensed + * summary has its own deterministic id. `repo` is the scope key the record was + * written under (canonical when the repo has a canonical form). + */ +export interface KnowledgeStore { + /** Insert or replace by id. */ + put(record: KnowledgeRecord): Promise; + /** Insert only when no record has this id. Returns whether it inserted. */ + putIfAbsent(record: KnowledgeRecord): Promise; + get(id: string): Promise; + /** Every record written under this repo key. */ + list(repo: string): Promise; + remove(id: string): Promise; +} + +export class InMemoryKnowledgeStore implements KnowledgeStore { + #records = new Map(); + async put(record: KnowledgeRecord): Promise { + this.#records.set(record.id, structuredClone(record)); + } + async putIfAbsent(record: KnowledgeRecord): Promise { + if (this.#records.has(record.id)) return false; + this.#records.set(record.id, structuredClone(record)); + return true; + } + async get(id: string): Promise { + const r = this.#records.get(id); + return r === undefined ? undefined : structuredClone(r); + } + async list(repo: string): Promise { + return [...this.#records.values()].filter((r) => r.scope.repo === repo).map((r) => structuredClone(r)); + } + async remove(id: string): Promise { + this.#records.delete(id); + } +} + +/** + * File-backed: one JSONL per repo key under `/knowledge-records`. A + * different directory from repo-memory, so FileRepoMemory.repos() (which counts + * every .jsonl it finds) never sees these. Writes are lock + atomic rewrite, + * like FileRepoMemory.remove. + */ +export class FileKnowledgeStore implements KnowledgeStore { + #dir: string; + constructor(dir = join(stateDir(), "knowledge-records")) { + this.#dir = dir; + } + #file(repo: string): string { + return join(this.#dir, `${repo.replace(/[^a-zA-Z0-9._-]/g, "_")}.jsonl`); + } + #parse(raw: string): KnowledgeRecord[] { + const out: KnowledgeRecord[] = []; + for (const line of raw.split("\n")) { + if (line.trim() === "") continue; + try { + out.push(JSON.parse(line) as KnowledgeRecord); + } catch { + // torn tail line — skip + } + } + return out; + } + async #mutate(repo: string, fn: (records: KnowledgeRecord[]) => KnowledgeRecord[] | null): Promise { + const path = this.#file(repo); + await withFileLock(path, async () => { + const current = this.#parse(await readFile(path, "utf8").catch(() => "")); + const next = fn(current); + if (next === null) return; + await writeTextFile(path, next.map((r) => JSON.stringify(r)).join("\n") + (next.length > 0 ? "\n" : "")); + }); + } + async put(record: KnowledgeRecord): Promise { + await this.#mutate(record.scope.repo, (rs) => [...rs.filter((r) => r.id !== record.id), record]); + } + async putIfAbsent(record: KnowledgeRecord): Promise { + let inserted = false; + await this.#mutate(record.scope.repo, (rs) => { + if (rs.some((r) => r.id === record.id)) return null; + inserted = true; + return [...rs, record]; + }); + return inserted; + } + async #all(): Promise { + const files = await readdir(this.#dir).catch(() => [] as string[]); + const out: KnowledgeRecord[] = []; + for (const f of files) { + if (f.endsWith(".jsonl")) out.push(...this.#parse(await readFile(join(this.#dir, f), "utf8").catch(() => ""))); + } + return out; + } + async get(id: string): Promise { + return (await this.#all()).find((r) => r.id === id); + } + async list(repo: string): Promise { + // Filter by scope.repo: the filename is sanitized, so repos can collide on one file. + return this.#parse(await readFile(this.#file(repo), "utf8").catch(() => "")).filter((r) => r.scope.repo === repo); + } + async remove(id: string): Promise { + const found = await this.get(id); + if (found === undefined) return; + await this.#mutate(found.scope.repo, (rs) => rs.filter((r) => r.id !== id)); + } +} + +/** Nucleus-backed: a NEW table, created on first use. */ +export class NucleusKnowledgeStore implements KnowledgeStore { + #db: NucleusPgwire; + #ready: Promise | null = null; + constructor(db: NucleusPgwire) { + this.#db = db; + } + #ensure(): Promise { + this.#ready ??= this.#db + .query( + `CREATE TABLE IF NOT EXISTS ship_knowledge_records ( + record_id TEXT, + repo TEXT, + project TEXT, + record TEXT, + updated_at TEXT + )`, + ) + .then(() => undefined) + // Same rule as ship_memory: a failed ensure must not be cached. + .catch((error: unknown) => { + this.#ready = null; + throw error; + }); + return this.#ready; + } + #parse(row: Record): KnowledgeRecord | undefined { + try { + return JSON.parse(String(row.record)) as KnowledgeRecord; + } catch { + return undefined; + } + } + async put(record: KnowledgeRecord): Promise { + await this.#ensure(); + // No upsert: delete-then-insert, the pattern the other Nucleus stores use. + await this.#db.query("DELETE FROM ship_knowledge_records WHERE record_id = $1", [record.id]); + await this.#insert(record); + } + async #insert(record: KnowledgeRecord): Promise { + await this.#db.query( + "INSERT INTO ship_knowledge_records (record_id, repo, project, record, updated_at) VALUES ($1, $2, $3, $4, $5)", + [record.id, record.scope.repo, record.scope.project, JSON.stringify(record), new Date().toISOString()], + ); + } + async putIfAbsent(record: KnowledgeRecord): Promise { + await this.#ensure(); + if ((await this.get(record.id)) !== undefined) return false; + await this.#insert(record); + return true; + } + async get(id: string): Promise { + await this.#ensure(); + const rows = await this.#db.query("SELECT record FROM ship_knowledge_records WHERE record_id = $1", [id]); + return rows.length === 0 ? undefined : this.#parse(rows[0]!); + } + async list(repo: string): Promise { + await this.#ensure(); + const rows = await this.#db.query("SELECT record FROM ship_knowledge_records WHERE repo = $1", [repo]); + return rows.map((r) => this.#parse(r)).filter((r): r is KnowledgeRecord => r !== undefined); + } + async remove(id: string): Promise { + await this.#ensure(); + await this.#db.query("DELETE FROM ship_knowledge_records WHERE record_id = $1", [id]); + } +} diff --git a/src/memory.ts b/src/memory.ts index 01e3b38..0fc22e4 100644 --- a/src/memory.ts +++ b/src/memory.ts @@ -91,6 +91,12 @@ export async function condenseIfNeeded( messages: Message[], summarize: Summarizer, config: CondenseConfig = defaultCondenseConfig, + /** + * S21: told about each summary this call generates, so it can be recorded as + * a derived (never evidence) knowledge record. Observe-only: it cannot change + * the messages, and a throw is swallowed. Omitted => identical to before. + */ + onSummary?: (info: { summary: string; condensed: number }) => Promise | void, ): Promise { if (historyTokens(messages) <= config.maxTokens) { return messages; @@ -127,6 +133,13 @@ export async function condenseIfNeeded( .map((m) => `${m.role.toUpperCase()}: ${typeof m.content === "string" ? m.content : JSON.stringify(m.content)}`) .join("\n\n"); const summary = await summarize(transcript); + if (onSummary !== undefined) { + try { + await onSummary({ summary, condensed: middle.length }); + } catch { + // advisory + } + } const summaryMessage: Message = { role: "user", diff --git a/src/repo-memory.ts b/src/repo-memory.ts index 1b86e80..020423d 100644 --- a/src/repo-memory.ts +++ b/src/repo-memory.ts @@ -7,6 +7,7 @@ import type { AgentExecutor } from "@neutron-build/agents"; import { appendLineSync, withFileLock, writeTextFile } from "./file-store.js"; import { frameUntrusted } from "./guard.js"; import type { NucleusPgwire } from "./nucleus-pgwire.js"; +import type { KnowledgeMode, NoteProvenanceHint, RetrievalContext } from "./knowledge-provenance.js"; import { stateDir } from "./run-store.js"; /** @@ -38,16 +39,32 @@ export interface RepoNote { note: string; runId?: string; createdAt: string; + /** + * S21: set ONLY on notes returned for model context with + * SHIP_KNOWLEDGE_PROVENANCE=on. Never stored. Absent means "not judged". + */ + freshness?: "fresh" | "stale" | "unknown"; } +/** What record() accepts: a note, plus (optionally) where its claim came from. Provenance is never stored on the note itself. */ +export type RecordInput = Omit & { provenance?: NoteProvenanceHint }; + export interface RepoMemoryStore { - record(note: Omit): Promise; - /** Most recent first, bounded in the query rather than after the fact. */ - recent(repo: string, limit: number): Promise; + record(note: RecordInput): Promise; + /** + * Most recent first, bounded in the query rather than after the fact. + * `context` marks a read for model context (S21 screens only those; the + * dashboard listing is never filtered, so a hidden note stays deletable). + */ + recent(repo: string, limit: number, options?: { context?: RetrievalContext }): Promise; + /** S21: the provenance mode this store applies to context reads. Absent = off. */ + provenanceMode?(): KnowledgeMode; + /** S21: record a condenser summary as a derived record. Absent = nothing recorded. */ + recordSummary?(input: { id: string; repo: string; runId: string; summary: string; derivedFrom: string[] }): Promise; /** Every repo that has notes, with its note count — for the dashboard. */ repos(): Promise<{ repo: string; count: number }[]>; - /** Delete exactly one note. */ - remove(noteId: string): Promise; + /** Delete exactly one note. `repo` (optional) lets S21 redaction find summaries derived from a note that has no provenance record. */ + remove(noteId: string, repo?: string): Promise; } /** File-backed memory: one JSONL per repo under the state dir. */ @@ -62,7 +79,8 @@ export class FileRepoMemory implements RepoMemoryStore { return join(this.#dir, `${repo.replace(/[^a-zA-Z0-9._-]/g, "_")}.jsonl`); } - async record(note: Omit): Promise { + async record(input: RecordInput): Promise { + const { provenance: _provenance, ...note } = input; const full: RepoNote = { ...note, noteId: `note-${randomUUID().slice(0, 12)}`, createdAt: new Date().toISOString() }; // Append, not read-modify-write — two concurrent records must not clobber // each other's note. Durably, so a note is not lost to a crash. @@ -162,7 +180,8 @@ export class NucleusRepoMemory implements RepoMemoryStore { return this.#ready; } - async record(note: Omit): Promise { + async record(input: RecordInput): Promise { + const { provenance: _provenance, ...note } = input; await this.#ensure(); const full: RepoNote = { ...note, noteId: `note-${randomUUID().slice(0, 12)}`, createdAt: new Date().toISOString() }; await this.#db.query( @@ -249,7 +268,7 @@ const RECENT_NOTES = 5; */ export async function loadRepoContext( executor: AgentExecutor, - options: { repo: string; memory?: RepoMemoryStore }, + options: { repo: string; memory?: RepoMemoryStore; onNotes?: (noteIds: string[]) => void }, ): Promise { const parts: string[] = []; for (const file of PLAYBOOK_FILES) { @@ -276,9 +295,18 @@ export async function loadRepoContext( } } if (options.memory !== undefined) { - const notes = await options.memory.recent(options.repo, RECENT_NOTES); + // S21: with provenance off (the default) this is the same call as before. + let notes: RepoNote[]; + if ((options.memory.provenanceMode?.() ?? "off") === "off") { + notes = await options.memory.recent(options.repo, RECENT_NOTES); + } else { + const head = await executor.exec("git rev-parse HEAD 2>/dev/null").catch(() => undefined); + const sha = head !== undefined && head.exitCode === 0 ? head.stdout.trim() : ""; + notes = await options.memory.recent(options.repo, RECENT_NOTES, { context: sha !== "" ? { head: sha } : {} }); + } if (notes.length > 0) { - const lines = notes.map((n) => `- [${n.createdAt.slice(0, 10)}] ${n.note}`); + options.onNotes?.(notes.map((n) => n.noteId)); + const lines = notes.map((n) => `- [${n.createdAt.slice(0, 10)}]${n.freshness !== undefined ? ` (${n.freshness})` : ""} ${n.note}`); // Ship's own notes about past runs — but their text came from task and // summary strings that originated outside, so they are framed too. A // poisoned note would otherwise be an instruction channel into every diff --git a/src/runtime.ts b/src/runtime.ts index 3434796..22641d6 100644 --- a/src/runtime.ts +++ b/src/runtime.ts @@ -1,4 +1,6 @@ import { ScopedRepoMemory } from "./scoped-repo-memory.js"; +import { KnowledgeProvenance, knowledgeProvenanceMode } from "./knowledge-provenance.js"; +import { FileKnowledgeStore, NucleusKnowledgeStore } from "./knowledge-store.js"; import { canonicalRepositoryURL, repoSlug } from "./repository-reference.js"; import { ScopedRepoStatsStore } from "./scoped-repo-stats.js"; import { finishReviewReplacement } from "./revision-launch.js"; @@ -466,7 +468,7 @@ export function fileRuntime(): ShipRuntime { governance: new FileGovernanceStore(), fleet: new FileFleetStore(), placement: new FilePlacementStore(), - memory: new ScopedRepoMemory(new FileRepoMemory(),store), + memory: new ScopedRepoMemory(new FileRepoMemory(),store,knowledgeProvenanceMode()==='off'?undefined:new KnowledgeProvenance({mode:knowledgeProvenanceMode(),store:new FileKnowledgeStore()})), steer: new FileSteerStore(), live: new FileLiveStore(), users: new FileUserStore(), @@ -589,7 +591,7 @@ export async function nucleusRuntime( governance: new NucleusGovernanceStore(db), fleet: new NucleusFleetStore(db), placement: new NucleusPlacementStore(db), - memory: new ScopedRepoMemory(new NucleusRepoMemory(db),store), + memory: new ScopedRepoMemory(new NucleusRepoMemory(db),store,knowledgeProvenanceMode()==='off'?undefined:new KnowledgeProvenance({mode:knowledgeProvenanceMode(),store:new NucleusKnowledgeStore(db)})), steer: new NucleusSteerStore(db), live: new NucleusLiveStore(db), users: new NucleusUserStore(db), diff --git a/src/scoped-repo-memory.ts b/src/scoped-repo-memory.ts index 402af5a..a8febbc 100644 --- a/src/scoped-repo-memory.ts +++ b/src/scoped-repo-memory.ts @@ -1,5 +1,6 @@ import type {EventStore} from '@neutron-build/workflow'; -import type {RepoMemoryStore,RepoNote} from './repo-memory.js'; +import type {RecordInput,RepoMemoryStore,RepoNote} from './repo-memory.js'; +import type {KnowledgeMode,KnowledgeProvenance,RetrievalContext} from './knowledge-provenance.js'; import {canonicalRepositoryURL} from './repository-reference.js'; import {repoKeyOf} from './repository-scope.js'; @@ -15,11 +16,40 @@ import {repoKeyOf} from './repository-scope.js'; * legacy notes still pass the recorded-run proof before entering context. */ export class ScopedRepoMemory implements RepoMemoryStore { - constructor(private inner:RepoMemoryStore,private events:EventStore){} - record(note:Omit){return this.inner.record({...note,repo:canonicalRepositoryURL(note.repo)??note.repo})} + /** `provenance` is S21 (SHIP_KNOWLEDGE_PROVENANCE); absent or mode off => every method below is the pre-S21 behaviour. */ + constructor(private inner:RepoMemoryStore,private events:EventStore,private provenance?:KnowledgeProvenance){} + provenanceMode():KnowledgeMode{return this.provenance?.mode??'off'} + async record(input:RecordInput){ + const {provenance:hint,...note}=input; + const stored=await this.inner.record({...note,repo:canonicalRepositoryURL(note.repo)??note.repo}); + // Advisory, like the note itself: a provenance failure never fails a record. + await this.provenance?.recordNote(stored,hint).catch(()=>{}); + return stored; + } repos(){return this.inner.repos()} - remove(id:string){return this.inner.remove(id)} - async recent(repo:string,limit:number):Promise { + async remove(id:string,repo?:string){ + await this.inner.remove(id,repo); + if(this.provenance===undefined||this.provenance.mode==='off')return; + // Redaction must not be skipped silently: `on` surfaces a failure to the + // caller, shadow logs it (a shadow record holds the deleted note's text). + try{ + const {deleted}=await this.provenance.redact([id],repo); + for(const d of deleted)await this.inner.remove(d,repo); + }catch(error){ + if(this.provenance.mode==='on')throw error; + console.error(JSON.stringify({event:'knowledge-provenance-error',op:'redact',noteId:id,error:String(error)})); + } + } + recordSummary(input:{id:string;repo:string;runId:string;summary:string;derivedFrom:string[]}){return this.provenance?.recordSummary(input)??Promise.resolve()} + async recent(repo:string,limit:number,options?:{context?:RetrievalContext}):Promise { + const p=this.provenance; + if(p===undefined||p.mode==='off'||options?.context===undefined)return this.#gather(repo,limit); + // `on` filters after the limit, so over-fetch to keep the block full; shadow never changes the result. + const notes=await this.#gather(repo,p.mode==='on'?limit*3:limit); + const screened=await p.screen(repo,notes,options.context); + return p.mode==='on'?screened.slice(0,limit):notes; + } + async #gather(repo:string,limit:number):Promise { const identity=canonicalRepositoryURL(repo); if(!identity)return this.inner.recent(repo,limit); const v1=repoKeyOf(identity); diff --git a/web/src/views/knowledge.server.ts b/web/src/views/knowledge.server.ts index 7a09b62..e8b41e1 100644 --- a/web/src/views/knowledge.server.ts +++ b/web/src/views/knowledge.server.ts @@ -34,7 +34,9 @@ export async function action({ request }: { request: Request }): Promise