diff --git a/src/actor.ts b/src/actor.ts index 09775d6..53f4022 100644 --- a/src/actor.ts +++ b/src/actor.ts @@ -8,6 +8,7 @@ import { TOOL_OUTPUT_MAX, } from "./config"; import { + type ArtifactRef, type AttachmentRef, Event, type EventData, @@ -673,6 +674,11 @@ export class ConversationActor { input: step.input, }); } else if (step.kind === "tool-result") { + // Documents the tool produced come out of its result and become + // first-class: refcounted so GC can reclaim them, versioned so a + // rerun under the same name builds history, and lifted onto the event + // itself so a client never has to dig them out of a tool's payload. + const artifacts = this.promoteArtifacts(step.output, messageId); this.persist(Event.ToolResult, { runId, threadId: this.conversationId, @@ -680,6 +686,7 @@ export class ConversationActor { toolCallId: step.toolCallId, toolName: step.toolName, output: capToolOutput(step.output), + ...(artifacts.length ? { artifacts } : {}), isError: step.isError, }); } else if (step.kind === "usage") { @@ -725,6 +732,36 @@ export class ConversationActor { this.currentRunId = null; } + /** + * Record a tool's `artifacts[]` and hand back the refs, version stamped. + * + * Bytes are already in the blob store by the time this runs — the tool put + * them there. This is the bookkeeping: a `blob_refs` row so the orphan sweep + * knows the blob is live, and an `artifacts` row so the document can be listed + * and read back without scanning the log. + */ + private promoteArtifacts(output: unknown, messageId: string): ArtifactRef[] { + if (!output || typeof output !== "object") return []; + const raw = (output as { artifacts?: unknown }).artifacts; + if (!Array.isArray(raw)) return []; + const out: ArtifactRef[] = []; + for (const a of raw as ArtifactRef[]) { + if (!a?.sha256 || !a.name) continue; + this.store.addBlobRef(a.sha256, this.conversationId); + const version = this.store.recordArtifact({ + conversationId: this.conversationId, + name: a.name, + sha256: a.sha256, + title: a.title, + mime: a.mime || "application/octet-stream", + size: a.size ?? 0, + messageId, + }); + out.push({ ...a, version }); + } + return out; + } + /** * Log a running tool's progress (see ToolProgressData). * diff --git a/src/client/app.css b/src/client/app.css index eca2234..c13430b 100644 --- a/src/client/app.css +++ b/src/client/app.css @@ -908,6 +908,37 @@ body.selecting #selectBtn { width: 16px; height: 16px; } +/* Documents this chat produced. Sits at the right of the header, and only when + there is something to list. */ +.docsbtn { + display: inline-flex; + align-items: center; + gap: 5px; + margin-left: auto; + padding: 4px 8px; + border: 1px solid var(--rule); + border-radius: var(--radius-md); + background: none; + color: var(--ink-muted); + cursor: pointer; +} +.docsbtn[hidden] { + display: none; +} +.docsbtn:hover { + color: var(--ink); + background: var(--bg-sunk); + border-color: var(--rule-strong); +} +.docsbtn svg { + width: 15px; + height: 15px; +} +.docscount { + font-family: var(--mono); + font-size: 11px; +} + .scroll { flex: 1; overflow-y: auto; diff --git a/src/client/app.js b/src/client/app.js index dae37a2..3380516 100644 --- a/src/client/app.js +++ b/src/client/app.js @@ -1318,7 +1318,7 @@ import { mountSidebar } from "./sidebar.js"; p.order.push(host); } p.domains[host]++; - rsrchLabel(p.gather, "Gathering " + p.read + " sources and counting\u2026"); + rsrchLabel(p.gather, "Gathering " + p.read + (p.read === 1 ? " source" : " sources")); rsrchTally(p); break; } @@ -1335,17 +1335,13 @@ import { mountSidebar } from "./sidebar.js"; case "citing": rsrchLabel(p.write, "Attaching citations"); break; - // The document arrives here rather than in the tool result: the result is - // permanent conversation context, so it carries only a summary, and the - // report rides the progress channel instead — durable, rendered, and never - // sent to a model. Replay rebuilds the card from this same event. case "done": rsrchState(p.write, "done"); rsrchLabel(p.write, d.title ? "Created \u201c" + d.title + "\u201d" : "Report written"); - if (d.report && !t.artifacted) { - t.artifacted = true; - addArtifact(t.block.rec, researchArtifact(t, d)); - } + // The card comes from the tool-result's artifacts[]; this copy of the + // text is kept only so its preview can render without waiting on a fetch + // during the run that produced it. + if (d.report) t.liveReport = d.report; break; } autoScroll(); @@ -1429,15 +1425,32 @@ import { mountSidebar } from "./sidebar.js"; // with copy and download. One pane, reused — opening a second artifact // replaces the first rather than stacking. var paneDoc = null; + // Artifacts are blob references, so the pane fetches the bytes rather than the + // event carrying them: content-addressed and immutable, which is exactly the + // cache the browser is good at (one GET per sha, then never again). function openPane(doc) { paneDoc = doc; - $("paneTitle").textContent = doc.title; + $("paneTitle").textContent = doc.title || doc.name; var body = $("paneBody"); body.innerHTML = ""; - renderStaticMd(body, doc.content); - body.scrollTop = 0; $("pane").hidden = false; $("app").classList.add("pane-open"); + var want = doc.sha256; + fetch("/api/blobs/" + encodeURIComponent(doc.sha256)) + .then(function (r) { + return r.ok ? r.text() : Promise.reject(new Error("HTTP " + r.status)); + }) + .then(function (text) { + if (!paneDoc || paneDoc.sha256 !== want) return; // opened something else meanwhile + paneDoc.content = text; + body.innerHTML = ""; + renderStaticMd(body, text); + body.scrollTop = 0; + }) + .catch(function () { + if (!paneDoc || paneDoc.sha256 !== want) return; + body.textContent = "This document could not be loaded."; + }); } function closePane() { paneDoc = null; @@ -1447,14 +1460,18 @@ import { mountSidebar } from "./sidebar.js"; } // The report is already markdown text, so a blob URL saves it without a round // trip to the server; it's revoked as soon as the click is dispatched. + // The blob endpoint sets Content-Disposition from `?name`, so the browser + // saves it under its real filename without the app re-encoding the bytes. function downloadMd() { if (!paneDoc) return; - var url = URL.createObjectURL(new Blob([paneDoc.content], { type: "text/markdown" })); var a = document.createElement("a"); - a.href = url; - a.download = paneDoc.filename || "document.md"; + a.href = + "/api/blobs/" + + encodeURIComponent(paneDoc.sha256) + + "?name=" + + encodeURIComponent(paneDoc.name); + a.download = paneDoc.name; a.click(); - URL.revokeObjectURL(url); } // PDF via the browser's own print pipeline rather than a bundled generator: // the pane is already the rendered document, so print CSS hides everything @@ -1472,7 +1489,53 @@ import { mountSidebar } from "./sidebar.js"; window.addEventListener("afterprint", restore); window.print(); } + // ---- the documents button ---------------------------------------------- + // What this chat has produced, listed from the artifacts projection rather + // than by walking the thread. It appears only once there IS a document, so a + // chat that never made one carries no chrome for it. + var artifactList = []; + function refreshArtifacts(id) { + var btn = $("docsBtn"); + if (!id) { + artifactList = []; + btn.hidden = true; + return; + } + fetch("/api/conversations/" + encodeURIComponent(id) + "/artifacts") + .then(function (r) { + return r.ok ? r.json() : { artifacts: [] }; + }) + .then(function (d) { + if (convId !== id) return; // switched chats while it was in flight + artifactList = d.artifacts || []; + btn.hidden = artifactList.length === 0; + $("docsCount").textContent = artifactList.length || ""; + }) + .catch(function () {}); + } + function mountDocsButton() { + var btn = $("docsBtn"); + btn.addEventListener("click", function (e) { + if (!artifactList.length) return; + e.stopPropagation(); + var r = btn.getBoundingClientRect(); + showContextMenu( + r.right, + r.bottom + 6, + artifactList.map(function (a) { + return { + label: (a.title || a.name) + (a.versions > 1 ? " · v" + a.version : ""), + onClick: function () { + openPane(a); + }, + }; + }), + { align: "right", trigger: btn }, + ); + }); + } function mountPane() { + mountDocsButton(); $("paneClose").addEventListener("click", closePane); $("paneCopy").addEventListener("click", function () { if (!paneDoc) return; @@ -1508,27 +1571,13 @@ import { mountSidebar } from "./sidebar.js"; // live panel (timeline, tally, streaming text) is scaffolding for work in // progress; once the work is done it's replaced by a card that names the thing // and opens it in the pane. - function researchArtifact(t, output) { - var report = output.report || ""; - // The subagent names its own document (write_report). The report's own H1 and - // then the question are fallbacks for a run that ended before it filed. - var heading = /^#{1,3}\s+(.+?)\s*$/m.exec(report); - var title = - output.title || - (heading && heading[1]) || - (t.input && t.input.question) || - "Research findings"; - var doc = { - title: title, - filename: output.filename || "research-findings.md", - content: report, - }; - + function artifactCard(ref, peekText) { + var title = ref.title || ref.name; var wrap = document.createElement("button"); wrap.type = "button"; wrap.className = "artifact"; wrap.addEventListener("click", function () { - openPane(doc); + openPane(ref); }); var text = document.createElement("span"); @@ -1538,11 +1587,9 @@ import { mountSidebar } from "./sidebar.js"; h.textContent = title; var kind = document.createElement("span"); kind.className = "artifactkind"; - kind.textContent = - "Document" + - (output.sources && output.sources.length - ? " \u00b7 " + output.sources.length + " sources" - : ""); + // Version only once there IS history — "v1" on a document written once is + // noise about a thing that hasn't happened. + kind.textContent = "Document" + (ref.version > 1 ? " \u00b7 v" + ref.version : ""); text.appendChild(h); text.appendChild(kind); @@ -1552,50 +1599,44 @@ import { mountSidebar } from "./sidebar.js"; peek.className = "artifactpeek"; var paper = document.createElement("span"); paper.className = "artifactpaper"; - paper.textContent = report + peek.appendChild(paper); + if (peekText) paper.textContent = peekPreview(peekText); + // No live copy (a replayed log): pull just enough of the blob to preview. + else + fetch("/api/blobs/" + encodeURIComponent(ref.sha256)) + .then(function (r) { + return r.ok ? r.text() : ""; + }) + .then(function (body) { + paper.textContent = peekPreview(body); + }) + .catch(function () {}); + + wrap.appendChild(text); + wrap.appendChild(peek); + return wrap; + } + function peekPreview(body) { + return String(body || "") .split("\n") .filter(function (l) { return l.trim(); }) .slice(0, 14) .join("\n"); - peek.appendChild(paper); - - wrap.appendChild(text); - wrap.appendChild(peek); - return wrap; } - // Artifacts belong to the message, not to the step that happened to make them: - // a document produced halfway through a turn shouldn't be buried in a tool - // gutter above three more paragraphs of reply. - // - // They hang off the TURN, after `.body`, rather than inside it. Inside, the - // best we could do is move the container to the end each time one lands — and - // the reply is still streaming, so the next prose block appends after it and - // the document ends up in the middle again. Outside the body it has nothing to - // race: `.body` is the turn's last child, so anything after it is last, always. - function addArtifact(rec, el) { - if (!rec.artifacts) { - rec.artifacts = document.createElement("div"); - rec.artifacts.className = "artifacts"; - rec.turn.appendChild(rec.artifacts); - } - rec.artifacts.appendChild(el); - } - function renderResearchResult(t, output) { - // The `done` progress event already built the card (it carries the document; - // this result deliberately doesn't). This is the fallback for a run whose - // progress never landed, and for events logged before the split. + // Any tool's documents, off the event itself. `artifacts[]` is on tool-result + // for every tool (see spec, "Artifacts — the promotion path"), so this is not a + // research feature — a sandbox step that promotes an output file renders the + // same card by the same path. + function renderArtifacts(rec, t, data) { + if (!Array.isArray(data.artifacts) || !data.artifacts.length) return; if (t.artifacted) return; t.artifacted = true; - if (t.panel) { - rsrchState(t.panel.write, "done"); - rsrchLabel( - t.panel.write, - output.title ? "Created \u201c" + output.title + "\u201d" : "Report written", - ); + for (const ref of data.artifacts) { + addArtifact(rec, artifactCard(ref, t.liveReport)); } - if (output.report) addArtifact(t.block.rec, researchArtifact(t, output)); + refreshArtifacts(convId); // the header list just gained one } // Upgrade a fetched page's favicon from the default service .ico using what the // page actually declares (fetch.ts pageFavicons): prefer its own SVG favicon @@ -1784,10 +1825,9 @@ import { mountSidebar } from "./sidebar.js"; var read = s.read + (s.read === 1 ? " page" : " pages"); return "Researched · " + read + " in " + Math.round(s.ms / 1000) + "s"; }, - result: function (t, output) { - if (output && (output.report || output.sources)) renderResearchResult(t, output); - else defaultResult(t, output); - }, + // The document renders from the event's artifacts[], the same path any + // tool's output files take — nothing tool-specific left to draw here. + result: function () {}, }, read_artifact: { icon: ICON_RESEARCH, @@ -2078,6 +2118,7 @@ import { mountSidebar } from "./sidebar.js"; retryCount = 0; hat.set("live"); convId = id; + refreshArtifacts(id); // Clear the previous conversation's thread synchronously, before the async // history load. Otherwise, when the chat shell is re-revealed after a detour // through a satellite view, it briefly shows the old conversation (a full diff --git a/src/client/index.html b/src/client/index.html index 98d954d..671c552 100644 --- a/src/client/index.html +++ b/src/client/index.html @@ -61,6 +61,11 @@ New conversation + +
diff --git a/src/drive.ts b/src/drive.ts index a92877d..379d99c 100644 --- a/src/drive.ts +++ b/src/drive.ts @@ -244,6 +244,7 @@ export class JobDriver { model: spec.model, abortSignal: signal, store: this.store, + blobs: this.blobs, owner, conversationId: actor.conversationId, project, diff --git a/src/events.ts b/src/events.ts index e5eabf6..ed9762c 100644 --- a/src/events.ts +++ b/src/events.ts @@ -24,6 +24,27 @@ export const Event = { } as const; export type EventName = (typeof Event)[keyof typeof Event]; +/** + * A document a tool produced: a reference, never bytes (spec, "Artifacts — the + * promotion path"). Agent-made files use the same content-addressed blob store + * as user uploads — one mechanism — so an artifact can be downloaded, fed back + * into a later tool, or materialized into the sandbox by its sha256. + * + * `name` is the document's identity within a conversation; writing it again + * makes a new version rather than a new document. + */ +export interface ArtifactRef { + sha256: string; + /** Filename, e.g. "hack-club-funding.md". */ + name: string; + /** Human title for the card; falls back to the name. */ + title?: string; + mime: string; + size: number; + /** Assigned by the store when the reference is recorded. */ + version?: number; +} + /** A blob referenced by a message (bytes in the BlobStore, keyed by sha256). */ export interface AttachmentRef { sha256: string; @@ -139,6 +160,8 @@ export interface ToolResultData { toolCallId: string; toolName: string; output: unknown; + /** Documents this tool produced. References only — bytes live in the blob store. */ + artifacts?: ArtifactRef[]; /** True when the tool threw / errored rather than returning a result. */ isError?: boolean; } diff --git a/src/http.ts b/src/http.ts index a87f8fd..e4522dd 100644 --- a/src/http.ts +++ b/src/http.ts @@ -670,6 +670,16 @@ export function apiRoutes(deps: { store: Store; blobs: BlobStore; kick?: () => v }, // Content-addressed blobs: upload (raw body) → sha256; fetch by sha256. + // Documents this conversation's tools produced — newest version of each, for + // the header list and the artifact pane. Bytes come from /api/blobs/:sha256. + "/api/conversations/:id/artifacts": { + GET: (req: Bun.BunRequest<"/api/conversations/:id/artifacts">) => { + const denied = guardConv(req, req.params.id); + if (denied) return denied; + return Response.json({ artifacts: store.listArtifacts(req.params.id) }); + }, + }, + "/api/blobs": { POST: (req: Bun.BunRequest<"/api/blobs">) => uploadBlob(req, store, blobs), }, diff --git a/src/inference.ts b/src/inference.ts index 8dc547d..bb9c073 100644 --- a/src/inference.ts +++ b/src/inference.ts @@ -1,5 +1,6 @@ import { type JSONValue, type LanguageModel, type ModelMessage, stepCountIs, streamText } from "ai"; import type { RunStep } from "./actor"; +import type { BlobStore } from "./blobs"; import { type LoadCatalogOptions, loadCatalog } from "./catalog"; import type { TokenUsage } from "./events"; import { contextToText, getContext, lardConnected, lardEnabled } from "./lard"; @@ -235,6 +236,8 @@ export interface RunOptions { temperature?: number; /** For per-user lard: the run's store + the conversation owner's `sub`. */ store?: Store; + /** Where a tool's output files land (agent artifacts share the upload store). */ + blobs?: BlobStore; owner?: string; /** The conversation, so the sandbox binds to its persistent per-chat container. */ conversationId?: string; @@ -263,6 +266,7 @@ export async function* run(messages: ModelMessage[], opts: RunOptions): AsyncGen store: opts.store, owner: opts.owner, conversationId: opts.conversationId, + blobs: opts.blobs, model, // deep_research runs its subagent on the same model as the run onProgress: opts.onProgress, }); diff --git a/src/store.ts b/src/store.ts index 10637d4..5b8d00d 100644 --- a/src/store.ts +++ b/src/store.ts @@ -135,6 +135,23 @@ export interface ContextFileMeta { createdAt: number; } +/** One stored version of a document a tool produced. */ +export interface ArtifactVersion { + name: string; + version: number; + sha256: string; + title: string | null; + mime: string; + size: number; + messageId: string | null; + createdAt: number; +} + +/** A document at its newest version, with a count of how many exist. */ +export interface ArtifactSummary extends ArtifactVersion { + versions: number; +} + /** A search hit: a conversation plus an excerpt of the matching message. */ export interface ConversationSearchResult extends ConversationSummary { /** Text around the match (title match → excerpt of the first message). */ @@ -296,6 +313,29 @@ CREATE TABLE IF NOT EXISTS blob_refs ( ); CREATE INDEX IF NOT EXISTS idx_blob_refs_conv ON blob_refs (conversation_id); +-- Documents a conversation's tools produced (see spec "Artifacts — the promotion +-- path"). A PROJECTION of the log, like blob_refs: the tool-result event stays +-- authoritative, this is the index that makes "list this chat's documents" and +-- "read the newest report.md" cheap instead of a scan. +-- +-- (conversation_id, name) is the document; version is its revision, so a rerun +-- into the same filename builds history rather than shadowing what came before. +-- The bytes are the blob, addressed by sha256 and shared with every other file +-- in the system. +CREATE TABLE IF NOT EXISTS artifacts ( + conversation_id TEXT NOT NULL, + name TEXT NOT NULL, + version INTEGER NOT NULL, + sha256 TEXT NOT NULL, + title TEXT, + mime TEXT NOT NULL, + size INTEGER NOT NULL, + message_id TEXT, + created_at INTEGER NOT NULL, + PRIMARY KEY (conversation_id, name, version) +); +CREATE INDEX IF NOT EXISTS idx_artifacts_conv ON artifacts (conversation_id, created_at DESC); + -- Auth sessions (indiko OAuth). The cookie holds an opaque high-entropy id; the -- row carries the user's identity (sub) and cached profile JSON. Expired rows -- are swept periodically. Only used when auth is enabled. @@ -719,56 +759,100 @@ export class Store { } /** - * Documents produced by this conversation's tools, newest first. + * Record a document a tool produced, as a new version of its name. * - * Derived from the event log rather than stored in a table of its own, because - * the log already holds them: `deep_research` emits its finished report on the - * progress channel (see ConversationActor.toolProgress), deliberately keeping - * it out of the tool result and so out of every later model prompt. That makes - * the log the one copy, and this the read side of it. + * Content addressing makes the no-op case free: bytes identical to the current + * version aren't a new version, they're the same document written twice. That + * matters because a rerun with the same answer shouldn't manufacture history. + * + * Returns the version it landed on, so the caller can say "v3" without a + * second query. */ - listArtifacts(conversationId: string): Array<{ filename: string; title: string; at: number }> { - const rows = this.db + recordArtifact(a: { + conversationId: string; + name: string; + sha256: string; + title?: string; + mime: string; + size: number; + messageId?: string; + }): number { + const latest = this.db .prepare( - `SELECT json_extract(data, '$.data.filename') AS filename, - json_extract(data, '$.data.title') AS title, - created_at - FROM events - WHERE conversation_id = ? - AND event = 'tool-progress' - AND json_extract(data, '$.phase') = 'done' - AND json_extract(data, '$.data.filename') IS NOT NULL - ORDER BY seq DESC`, + "SELECT version, sha256 FROM artifacts WHERE conversation_id = ? AND name = ? ORDER BY version DESC LIMIT 1", ) - .all(conversationId) as Array<{ filename: string; title: string; created_at: number }>; - // A document rewritten under the same name is one document at its newest. - const seen = new Set(); - const out: Array<{ filename: string; title: string; at: number }> = []; - for (const r of rows) { - if (seen.has(r.filename)) continue; - seen.add(r.filename); - out.push({ filename: r.filename, title: r.title ?? r.filename, at: r.created_at }); - } - return out; + .get(a.conversationId, a.name) as { version: number; sha256: string } | null; + if (latest?.sha256 === a.sha256) return latest.version; // same bytes, same document + const version = (latest?.version ?? 0) + 1; + this.db + .prepare( + `INSERT INTO artifacts + (conversation_id, name, version, sha256, title, mime, size, message_id, created_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`, + ) + .run( + a.conversationId, + a.name, + version, + a.sha256, + a.title ?? null, + a.mime, + a.size, + a.messageId ?? null, + Date.now(), + ); + return version; } - /** The newest version of one document's markdown, or null if there isn't one. */ - readArtifact(conversationId: string, filename: string): string | null { - const row = this.db + /** Every version of every document in a conversation, newest first. */ + artifactVersions(conversationId: string, name: string): ArtifactVersion[] { + return this.db .prepare( - `SELECT json_extract(data, '$.data.report') AS report - FROM events - WHERE conversation_id = ? - AND event = 'tool-progress' - AND json_extract(data, '$.phase') = 'done' - AND json_extract(data, '$.data.filename') = ? - ORDER BY seq DESC LIMIT 1`, + `SELECT name, version, sha256, title, mime, size, message_id AS messageId, created_at AS createdAt + FROM artifacts WHERE conversation_id = ? AND name = ? ORDER BY version DESC`, ) - .get(conversationId, filename) as { report: string | null } | null; - return row?.report ?? null; + .all(conversationId, name) as ArtifactVersion[]; + } + + /** + * One row per document — its newest version, plus how many there are. This is + * what the header list and the artifact pane read. + */ + listArtifacts(conversationId: string): ArtifactSummary[] { + return this.db + .prepare( + `SELECT a.name, a.version, a.sha256, a.title, a.mime, a.size, + a.message_id AS messageId, a.created_at AS createdAt, + (SELECT COUNT(*) FROM artifacts v + WHERE v.conversation_id = a.conversation_id AND v.name = a.name) AS versions + FROM artifacts a + WHERE a.conversation_id = ? + AND a.version = (SELECT MAX(v.version) FROM artifacts v + WHERE v.conversation_id = a.conversation_id AND v.name = a.name) + ORDER BY a.created_at DESC`, + ) + .all(conversationId) as ArtifactSummary[]; + } + + /** One document by name — its newest version, or a specific one. */ + getArtifact(conversationId: string, name: string, version?: number): ArtifactVersion | null { + const row = version + ? this.db + .prepare( + `SELECT name, version, sha256, title, mime, size, message_id AS messageId, created_at AS createdAt + FROM artifacts WHERE conversation_id = ? AND name = ? AND version = ?`, + ) + .get(conversationId, name, version) + : this.db + .prepare( + `SELECT name, version, sha256, title, mime, size, message_id AS messageId, created_at AS createdAt + FROM artifacts WHERE conversation_id = ? AND name = ? ORDER BY version DESC LIMIT 1`, + ) + .get(conversationId, name); + return (row as ArtifactVersion | null) ?? null; } - /** Set the title only if none is set yet (auto-title never clobbers a rename). + /** Set the title only if none is set yet /** Set the title only if none is set yet (auto-title never clobbers a rename). * Returns whether it actually set one. */ setTitleIfEmpty(id: string, title: string): boolean { const t = title.trim(); diff --git a/src/tools.ts b/src/tools.ts index 8f7d8fc..175923d 100644 --- a/src/tools.ts +++ b/src/tools.ts @@ -1,4 +1,6 @@ import { jsonSchema, type LanguageModel, type Tool, type ToolSet, tool } from "ai"; +import type { BlobStore } from "./blobs"; +import type { ArtifactRef } from "./events"; import { type Executor, formatExecResult, getExecutor } from "./executor"; import { createFetchProvider, type FetchProvider } from "./fetch"; import { @@ -87,8 +89,9 @@ function deepResearch( model: LanguageModel, search: SearchProvider, fetcher: FetchProvider, - onProgress?: ToolContext["onProgress"], + ctx: ToolContext, ) { + const onProgress = ctx.onProgress; return tool({ description: "Hand off a question that needs real research — several searches, several " + @@ -131,19 +134,41 @@ function deepResearch( ? (phase, data) => onProgress({ toolCallId, toolName: "deep_research", phase, data }) : undefined, }); - // The summary, not the report. A finished document runs to thousands of - // tokens and a tool result is permanent context — it would be re-sent on - // every later turn of the conversation, which is precisely the cost the - // subagent exists to avoid. The document itself reached the UI over the - // progress channel, which is durable and rendered but never shown to a - // model, so nothing is lost by leaving it out here. - return { + // The report becomes a blob and the result carries a REFERENCE — never the + // bytes. A finished document is thousands of tokens and a tool result is + // permanent context, re-sent on every later turn; the reference is a + // handful of tokens and is also the handle everything else works from, + // since agent output shares the content-addressed store with user uploads + // (read it back, download it, feed it to a later tool by sha256). + const result: { + summary: string; + document: string; + title: string; + sources: number; + stats: unknown; + artifacts?: ArtifactRef[]; + } = { summary: out.summary, document: out.filename, title: out.title, sources: out.sources.length, stats: out.stats, }; + if (ctx.blobs) { + const bytes = new TextEncoder().encode(out.report); + const ref = await ctx.blobs.put(bytes); + ctx.store?.recordBlob(ref.sha256, "text/markdown", ref.size); + result.artifacts = [ + { + sha256: ref.sha256, + name: out.filename, + title: out.title, + mime: "text/markdown", + size: ref.size, + }, + ]; + } + return result; }, }); } @@ -160,34 +185,42 @@ function deepResearch( * So the full text is available on request rather than by default. The cost is * paid once, in the turn that needs it, by the model that asked. */ -function readArtifact(store: Store, conversationId: string) { +function readArtifact(store: Store, blobs: BlobStore, conversationId: string) { return tool({ description: "Read the full text of a document produced earlier in this conversation " + "(e.g. a report from deep_research, by its filename). Use it when the user " + "asks about something a document covers in more detail than the summary you " + "were given. The user can already see the document, so answer from it rather " + - "than reproducing it wholesale.", - inputSchema: jsonSchema<{ filename: string }>({ + "than reproducing it wholesale. Pass a version to read an older revision; " + + "the newest is used by default.", + inputSchema: jsonSchema<{ name: string; version?: number }>({ type: "object", properties: { - filename: { - type: "string", - description: "The document's name, e.g. 'hack-club-funding.md'.", + name: { type: "string", description: "The document's name, e.g. 'hack-club-funding.md'." }, + version: { + type: "number", + description: "An earlier revision to read. Omit for the newest.", }, }, - required: ["filename"], + required: ["name"], additionalProperties: false, }), - execute: async ({ filename }) => { - const doc = store.readArtifact(conversationId, filename); - if (doc) return doc; - // A wrong name is worth answering with the right ones — the model named it - // from memory of a tool result several turns back. - const have = store.listArtifacts(conversationId); - return have.length - ? `No document named "${filename}". This conversation has: ${have.map((a) => a.filename).join(", ")}.` - : "No documents have been produced in this conversation yet."; + execute: async ({ name, version }) => { + const doc = store.getArtifact(conversationId, name, version); + if (!doc) { + // A wrong name is worth answering with the right ones — the model named + // it from memory of a tool result several turns back. + const have = store.listArtifacts(conversationId); + return have.length + ? `No document named "${name}"${version ? ` at version ${version}` : ""}. This conversation has: ${have + .map((a) => `${a.name} (v${a.version})`) + .join(", ")}.` + : "No documents have been produced in this conversation yet."; + } + const blob = await blobs.get(doc.sha256); + if (!blob) return `The bytes for "${name}" are missing from the blob store.`; + return await blob.text(); }, }); } @@ -260,7 +293,7 @@ const REGISTRY: Array<{ if (!ctx.model || !getConfig().research.enabled) return null; const search = createSearchProvider(); const fetcher = createFetchProvider(); - return search && fetcher ? deepResearch(ctx.model, search, fetcher, ctx.onProgress) : null; + return search && fetcher ? deepResearch(ctx.model, search, fetcher, ctx) : null; }, }, { @@ -270,8 +303,11 @@ const REGISTRY: Array<{ // answer "there aren't any" is prompt weight in every chat that never // researched anything. create: (ctx) => - ctx.store && ctx.conversationId && ctx.store.listArtifacts(ctx.conversationId).length - ? readArtifact(ctx.store, ctx.conversationId) + ctx.store && + ctx.blobs && + ctx.conversationId && + ctx.store.listArtifacts(ctx.conversationId).length + ? readArtifact(ctx.store, ctx.blobs, ctx.conversationId) : null, }, { @@ -375,6 +411,8 @@ export interface ToolContext { * model in hand. */ model?: LanguageModel; + /** Where a tool's output files go: agent artifacts share the user-upload store. */ + blobs?: BlobStore; /** * Report from inside a long-running tool, so the UI can show the work as it * happens rather than a spinner that lasts minutes. Absent when nothing is diff --git a/tests/actor.test.ts b/tests/actor.test.ts index 6e7b3b3..39ca337 100644 --- a/tests/actor.test.ts +++ b/tests/actor.test.ts @@ -93,6 +93,42 @@ test("tool progress reaches subscribers while the tool is still running", async expect(store.replay("t-prog", 0).some((e) => e.event === Event.ToolProgress)).toBe(true); }); +test("a tool's artifacts are promoted: refcounted, versioned, lifted onto the event", async () => { + const a = new ConversationActor("t-art", store); + const events: WireEvent[] = []; + a.follow({ push: (e) => events.push(e), closed: false }); + const ref = { + sha256: "c".repeat(64), + name: "report.md", + title: "A report", + mime: "text/markdown", + size: 42, + }; + + await a.runText("r-a", "m-a", async function* (_signal) { + yield { + kind: "tool-result", + toolCallId: "call-1", + toolName: "deep_research", + output: { summary: "short", artifacts: [ref] }, + }; + }); + + // Lifted onto the event, so a client reads documents off the tool-result + // rather than digging through one tool's payload shape. + const result = events.find((e) => e.event === Event.ToolResult)!.data as { + artifacts?: Array<{ name: string; version: number }>; + }; + expect(result.artifacts).toEqual([{ ...ref, version: 1 }]); + // Listable and readable without scanning the log. + expect(store.listArtifacts("t-art")).toHaveLength(1); + expect(store.getArtifact("t-art", "report.md")?.sha256).toBe(ref.sha256); + // And refcounted, so the orphan sweep knows the blob is live rather than + // reclaiming a document someone is still reading. + store.recordBlob(ref.sha256, ref.mime, ref.size); + expect(store.findOrphanBlobs(Date.now() + 1000)).not.toContain(ref.sha256); +}); + test("message-end omits usage when no usage step is yielded", async () => { const a = new ConversationActor("t-nousage", store); const events: WireEvent[] = []; diff --git a/tests/artifacts.test.ts b/tests/artifacts.test.ts index 09889de..eafaa25 100644 --- a/tests/artifacts.test.ts +++ b/tests/artifacts.test.ts @@ -2,87 +2,90 @@ import { expect, test } from "bun:test"; import { Store } from "../src/store"; /** - * A finished `deep_research` run, as the actor logs it: the document rides the - * progress channel (durable, rendered, never sent to a model) rather than the - * tool result. These are the read side of that. + * Documents a tool produced (spec, "Artifacts — the promotion path"). Bytes live + * in the content-addressed blob store; this table is the projection that makes + * "list this chat's documents" and "read the newest report.md" cheap. */ -function seedDocument( - store: Store, - conv: string, - seq: number, - filename: string, - title: string, - report: string, -): void { - const db = ( - store as unknown as { db: { query: (s: string) => { run: (...a: unknown[]) => void } } } - ).db; - db.query("INSERT OR IGNORE INTO conversations (id, created_at, last_seq) VALUES (?, ?, 0)").run( - conv, - Date.now(), - ); - db.query( - "INSERT INTO events (id, conversation_id, seq, event, data, created_at) VALUES (?, ?, ?, 'tool-progress', ?, ?)", - ).run( - `${conv}:${seq}`, - conv, - seq, - JSON.stringify({ - threadId: conv, - toolName: "deep_research", - phase: "done", - data: { filename, title, report }, - }), - Date.now() + seq, - ); +function fresh(): Store { + return new Store(":memory:"); +} +function record(store: Store, conv: string, name: string, sha256: string, title = "T"): number { + return store.recordArtifact({ + conversationId: conv, + name, + sha256, + title, + mime: "text/markdown", + size: 100, + messageId: "m1", + }); } -test("a document written to the log can be read back by name", () => { - const store = new Store(":memory:"); - seedDocument(store, "c1", 1, "funding.md", "How it's funded", "# Funding\n\nThe body."); - expect(store.readArtifact("c1", "funding.md")).toBe("# Funding\n\nThe body."); - expect(store.listArtifacts("c1")).toEqual([ - { filename: "funding.md", title: "How it's funded", at: expect.any(Number) }, - ]); +test("a recorded document is listed at version 1", () => { + const store = fresh(); + expect(record(store, "c1", "funding.md", "a".repeat(64), "How it's funded")).toBe(1); + const [doc] = store.listArtifacts("c1"); + expect(doc).toMatchObject({ + name: "funding.md", + version: 1, + versions: 1, + title: "How it's funded", + mime: "text/markdown", + }); }); -test("documents are scoped to their conversation", () => { - const store = new Store(":memory:"); - seedDocument(store, "c1", 1, "a.md", "A", "body a"); - seedDocument(store, "c2", 1, "b.md", "B", "body b"); - expect(store.readArtifact("c1", "b.md")).toBeNull(); - expect(store.listArtifacts("c2").map((a) => a.filename)).toEqual(["b.md"]); +test("writing the same name again is a new version, and the newest wins", () => { + const store = fresh(); + record(store, "c1", "report.md", "a".repeat(64), "Draft"); + expect(record(store, "c1", "report.md", "b".repeat(64), "Final")).toBe(2); + // One document, two revisions — not two documents. + const list = store.listArtifacts("c1"); + expect(list).toHaveLength(1); + expect(list[0]).toMatchObject({ version: 2, versions: 2, title: "Final" }); + expect(store.getArtifact("c1", "report.md")?.sha256).toBe("b".repeat(64)); }); -test("a name rewritten later reads as one document, at its newest", () => { - const store = new Store(":memory:"); - seedDocument(store, "c1", 1, "report.md", "Draft", "first pass"); - seedDocument(store, "c1", 2, "report.md", "Final", "second pass"); - expect(store.readArtifact("c1", "report.md")).toBe("second pass"); - expect(store.listArtifacts("c1")).toHaveLength(1); - expect(store.listArtifacts("c1")[0]?.title).toBe("Final"); +test("history is readable by version", () => { + const store = fresh(); + record(store, "c1", "report.md", "a".repeat(64), "Draft"); + record(store, "c1", "report.md", "b".repeat(64), "Final"); + expect(store.getArtifact("c1", "report.md", 1)?.title).toBe("Draft"); + expect(store.artifactVersions("c1", "report.md").map((v) => v.version)).toEqual([2, 1]); + expect(store.getArtifact("c1", "report.md", 9)).toBeNull(); }); -test("an unknown name reads as absent, not as an error", () => { - const store = new Store(":memory:"); - expect(store.readArtifact("c1", "nope.md")).toBeNull(); - expect(store.listArtifacts("c1")).toEqual([]); +test("identical bytes are the same document written twice, not a new version", () => { + // Content addressing makes this free, and it matters: a rerun that reaches the + // same answer shouldn't manufacture history. + const store = fresh(); + const sha = "a".repeat(64); + expect(record(store, "c1", "report.md", sha)).toBe(1); + expect(record(store, "c1", "report.md", sha)).toBe(1); + expect(store.artifactVersions("c1", "report.md")).toHaveLength(1); +}); + +test("documents are scoped to their conversation", () => { + const store = fresh(); + record(store, "c1", "a.md", "a".repeat(64)); + record(store, "c2", "b.md", "b".repeat(64)); + expect(store.getArtifact("c1", "b.md")).toBeNull(); + expect(store.listArtifacts("c2").map((a) => a.name)).toEqual(["b.md"]); }); -test("progress events that aren't documents are ignored", () => { - const store = new Store(":memory:"); - const db = ( - store as unknown as { db: { query: (s: string) => { run: (...a: unknown[]) => void } } } - ).db; - db.query("INSERT INTO conversations (id, created_at, last_seq) VALUES ('c1', ?, 0)").run( - Date.now(), - ); - // A `read` phase — hundreds of these per run, none of them a document. - db.query( - "INSERT INTO events (id, conversation_id, seq, event, data, created_at) VALUES ('c1:1','c1',1,'tool-progress',?,?)", - ).run( - JSON.stringify({ threadId: "c1", phase: "read", data: { url: "https://a.test" } }), - Date.now(), - ); +test("an empty conversation has no documents", () => { + const store = fresh(); expect(store.listArtifacts("c1")).toEqual([]); + expect(store.getArtifact("c1", "nope.md")).toBeNull(); +}); + +test("different names in one conversation are different documents", () => { + const store = fresh(); + record(store, "c1", "a.md", "a".repeat(64)); + record(store, "c1", "b.md", "b".repeat(64)); + expect( + store + .listArtifacts("c1") + .map((a) => a.name) + .sort(), + ).toEqual(["a.md", "b.md"]); });