diff --git a/docs/event-loop-contract-audit.md b/docs/event-loop-contract-audit.md index 6c3f801e0e..0374040e8e 100644 --- a/docs/event-loop-contract-audit.md +++ b/docs/event-loop-contract-audit.md @@ -155,11 +155,11 @@ The classifier follows execution home, not API spelling. A direct accessor callb - `pipeline/document-worker.ts:529` (withReadDbAsync) - `pipeline/document-worker.ts:557` (withReadDbAsync) - `pipeline/document-worker.ts:571` (withReadDbAsync) -- `pipeline/dreaming-attention.ts:138` (withReadDb) -- `pipeline/dreaming-attention.ts:150` (withReadDb) -- `pipeline/dreaming-attention.ts:220` (withReadDb) -- `pipeline/dreaming-attention.ts:251` (withReadDb) -- `pipeline/dreaming-attention.ts:285` (withReadDb) +- `pipeline/dreaming-attention.ts:140` (withReadDb) +- `pipeline/dreaming-attention.ts:152` (withReadDb) +- `pipeline/dreaming-attention.ts:226` (withReadDb) +- `pipeline/dreaming-attention.ts:257` (withReadDb) +- `pipeline/dreaming-attention.ts:291` (withReadDb) - `db:dreaming.operations.cite-evidence.read` (withReadDb) - `db:dreaming.operations.duplicate-group.read` (withReadDb) - `db:dreaming.operations.entity-name.read` (withReadDb) diff --git a/platform/core/src/migrations/169-dreaming-evidence-review-cursor.ts b/platform/core/src/migrations/169-dreaming-evidence-review-cursor.ts new file mode 100644 index 0000000000..b9076d2607 --- /dev/null +++ b/platform/core/src/migrations/169-dreaming-evidence-review-cursor.ts @@ -0,0 +1,20 @@ +import type { MigrationDb } from "./contract"; + +export function up(db: MigrationDb): void { + if ( + db.prepare("SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = 'dreaming_evidence_consumption'").get() == + null + ) + return; + const columns = db.prepare("PRAGMA table_info(dreaming_evidence_consumption)").all() as Array<{ name: string }>; + if (!columns.some((column) => column.name === "cursor_basis")) { + db.exec( + "ALTER TABLE dreaming_evidence_consumption ADD COLUMN cursor_basis TEXT NOT NULL DEFAULT 'delivery' CHECK (cursor_basis IN ('delivery', 'review'))", + ); + } + if (!columns.some((column) => column.name === "stalled_passes")) { + db.exec( + "ALTER TABLE dreaming_evidence_consumption ADD COLUMN stalled_passes INTEGER NOT NULL DEFAULT 0 CHECK (stalled_passes >= 0)", + ); + } +} diff --git a/platform/core/src/migrations/index.ts b/platform/core/src/migrations/index.ts index ddb9710fc8..429f4aaf00 100644 --- a/platform/core/src/migrations/index.ts +++ b/platform/core/src/migrations/index.ts @@ -169,6 +169,7 @@ import { up as dreamingHistory } from "./165-dreaming-history"; import { up as dreamingPassPeakContext } from "./166-dreaming-pass-peak-context"; import { up as memoriesFtsPorter } from "./167-memories-fts-porter"; import { up as retireSourceParagraphClaims } from "./168-retire-source-paragraph-claims"; +import { up as dreamingEvidenceReviewCursor } from "./169-dreaming-evidence-review-cursor"; export type { Migration, MigrationArtifacts, MigrationDb } from "./contract"; export const MIGRATIONS: readonly Migration[] = [ @@ -1527,6 +1528,17 @@ export const MIGRATIONS: readonly Migration[] = [ name: "retire-source-paragraph-claims", up: retireSourceParagraphClaims, }, + { + version: 169, + name: "dreaming-evidence-review-cursor", + up: dreamingEvidenceReviewCursor, + artifacts: { + columns: [ + { table: "dreaming_evidence_consumption", column: "cursor_basis" }, + { table: "dreaming_evidence_consumption", column: "stalled_passes" }, + ], + }, + }, ]; function checksum(m: Migration): string { let h = 0; diff --git a/platform/core/src/migrations/migrations.test.ts b/platform/core/src/migrations/migrations.test.ts index 24db138c0b..ab7df601f9 100644 --- a/platform/core/src/migrations/migrations.test.ts +++ b/platform/core/src/migrations/migrations.test.ts @@ -415,7 +415,7 @@ describe("migration framework", () => { runMigrations(db); const applied = db.query("SELECT MAX(version) AS version FROM schema_migrations").get() as { version: number }; - expect(applied.version).toBe(168); + expect(applied.version).toBe(169); expect( db.query("SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'vector_repair_checkpoints'").get(), ).toEqual({ name: "vector_repair_checkpoints" }); @@ -2326,6 +2326,30 @@ describe("migration framework", () => { ).not.toBeNull(); }); + test("migration 169 marks existing evidence cursors as delivery-based", () => { + db = createFreshDb(); + runMigrations(db); + db.exec(` + ALTER TABLE dreaming_evidence_consumption DROP COLUMN cursor_basis; + ALTER TABLE dreaming_evidence_consumption DROP COLUMN stalled_passes; + INSERT INTO dreaming_evidence_consumption + (agent_id, source_kind, source_id, source_captured_at, source_entry_id, source_revision, + delivered_offset, source_length, pass_id, updated_at) + VALUES ('default', 'transcript', 'legacy', '2026-01-01 00:00:00', '', 'r1', 40, 40, 'pass-1', datetime('now')); + `); + db.prepare("DELETE FROM schema_migrations WHERE version = 169").run(); + runMigrations(db); + + expect( + db + .query("SELECT cursor_basis, stalled_passes FROM dreaming_evidence_consumption WHERE source_id = 'legacy'") + .get(), + ).toEqual({ cursor_basis: "delivery", stalled_passes: 0 }); + expect(() => + db.exec("UPDATE dreaming_evidence_consumption SET cursor_basis = 'guess' WHERE source_id = 'legacy'"), + ).toThrow(); + }); + test("migration 063 limits memories_fts updates to content changes", () => { db = createFreshDb(); runMigrations(db); diff --git a/platform/daemon/src/daemon.ts b/platform/daemon/src/daemon.ts index e996e55be2..f77cb843e4 100644 --- a/platform/daemon/src/daemon.ts +++ b/platform/daemon/src/daemon.ts @@ -165,6 +165,7 @@ import { } from "./pipeline"; import { randomUUID } from "node:crypto"; import { recordDreamingPassTelemetry } from "./pipeline/dreaming"; +import { importedSourceAttentionUpsert } from "./pipeline/dreaming-evidence-consumption"; import { dbOwnerTransaction } from "./db-owner-runtime"; import { startDeferredRuntimeAfterDreaming } from "./dreaming-startup"; import { type DreamingWorkerHandle, startDreamingWorker } from "./pipeline/dreaming-worker"; @@ -2842,27 +2843,11 @@ async function main() { agentId: resolveDaemonAgentId(), workspaceRoot: AGENTS_DIR, onBatch: async (_jobId, sourceId) => { - const agentId = resolveDaemonAgentId(); - const subjectRef = `source:${sourceId}`; - const details = JSON.stringify({ sourceId, reason: "transcript-import-committed" }); - await dbOwnerTransaction( - [ - { - sql: `INSERT INTO dreaming_attention - (id, agent_id, kind, subject_ref, details_json, priority) - VALUES (?, ?, 'evidence_requeue', ?, ?, 50) - ON CONFLICT(agent_id, kind, subject_ref) DO UPDATE SET - details_json = excluded.details_json, - priority = MAX(dreaming_attention.priority, excluded.priority), - generation = dreaming_attention.generation + 1, - resolved_at = NULL, - resolved_by_pass_id = NULL`, - params: [randomUUID(), agentId, subjectRef, details], - result: "run", - }, - ], - { operation: "sources.import.dreaming-attention", lane: "write" }, - ); + const nudge = importedSourceAttentionUpsert(resolveDaemonAgentId(), sourceId); + await dbOwnerTransaction([{ sql: nudge.sql, params: nudge.params, result: "run" }], { + operation: "sources.import.dreaming-attention", + lane: "write", + }); }, }); } diff --git a/platform/daemon/src/db-owner-protocol.ts b/platform/daemon/src/db-owner-protocol.ts index 46394cd782..b7f842d0ca 100644 --- a/platform/daemon/src/db-owner-protocol.ts +++ b/platform/daemon/src/db-owner-protocol.ts @@ -256,6 +256,7 @@ export type DbOwnerRequest = | { readonly kind: "dreaming_episodic_backlog_exists"; readonly input: DbOwnerDreamingEpisodicBacklogExists } | { readonly kind: "dreaming_evidence_search"; readonly input: DbOwnerDreamingEvidenceSearch } | { readonly kind: "dreaming_evidence_source"; readonly input: DbOwnerDreamingEvidenceSource } + | { readonly kind: "dreaming_evidence_review"; readonly input: DbOwnerDreamingEvidenceReview } | { readonly kind: "dreaming_pass_finalize"; readonly input: DbOwnerDreamingPassFinalize } | { readonly kind: "dreaming_review_due"; readonly input: DbOwnerDreamingReviewDue } | { readonly kind: "dreaming_evidence_classify"; readonly input: DbOwnerDreamingEvidenceClassify } @@ -387,6 +388,16 @@ export interface DbOwnerDreamingEvidenceSource { readonly sourceRef: string; } +export interface DbOwnerDreamingEvidenceReview { + readonly agentId: string; + readonly passId: string; + readonly items: ReadonlyArray<{ + readonly sourceRef: string; + readonly contentOffset: number; + readonly through?: string; + }>; +} + export interface DbOwnerDreamingPassFinalize { readonly passId: string; readonly mode: string; diff --git a/platform/daemon/src/db-owner-runtime.ts b/platform/daemon/src/db-owner-runtime.ts index 75bc4dba96..bca2bdcb97 100644 --- a/platform/daemon/src/db-owner-runtime.ts +++ b/platform/daemon/src/db-owner-runtime.ts @@ -207,6 +207,7 @@ async function executeInlineOwnerRequest(accessor: DbAccessor, request: DbOwnerR case "dreaming_episodic_backlog_exists": case "dreaming_evidence_search": case "dreaming_evidence_source": + case "dreaming_evidence_review": case "memory_head": case "dreaming_pass_finalize": case "dreaming_review_due": diff --git a/platform/daemon/src/db-owner-worker.ts b/platform/daemon/src/db-owner-worker.ts index fb4c07187e..9435441ad9 100644 --- a/platform/daemon/src/db-owner-worker.ts +++ b/platform/daemon/src/db-owner-worker.ts @@ -560,6 +560,13 @@ export function runDbOwnerWorker(): void { return readDreamingEvidenceSourceInDb(db as never, request.input); } + async function executeDreamingEvidenceReview( + request: Extract, + ): Promise { + const { reviewDreamingEvidenceInDb } = await import("./pipeline/dreaming-evidence-consumption"); + return reviewDreamingEvidenceInDb(db as never, request.input); + } + async function executeDreamingPassFinalize( request: Extract, context: JobExecutionContext, @@ -1082,6 +1089,7 @@ export function runDbOwnerWorker(): void { return await executeDreamingEpisodicBacklogExists(job.request); if (job.request.kind === "dreaming_evidence_search") return await executeDreamingEvidenceSearch(job.request); if (job.request.kind === "dreaming_evidence_source") return await executeDreamingEvidenceSource(job.request); + if (job.request.kind === "dreaming_evidence_review") return await executeDreamingEvidenceReview(job.request); if (job.request.kind === "dreaming_pass_finalize") return executeDreamingPassFinalize(job.request, context); if (job.request.kind === "dreaming_review_due") return await executeDreamingReviewDue(job.request); if (job.request.kind === "dreaming_evidence_classify") return await executeDreamingEvidenceClassify(job.request); diff --git a/platform/daemon/src/episodic-sources.ts b/platform/daemon/src/episodic-sources.ts index 8e3bf4a9bb..e86b5ef302 100644 --- a/platform/daemon/src/episodic-sources.ts +++ b/platform/daemon/src/episodic-sources.ts @@ -883,6 +883,17 @@ function selectEpisodicSourceRefs( : transcriptHasUpdatedAt ? "COALESCE(session_transcripts.updated_at, session_transcripts.created_at)" : "session_transcripts.created_at"; + const transcriptEntryId = tableHasColumn(db, "session_transcripts", "source_id") + ? "COALESCE(session_transcripts.source_id, '')" + : "''"; + const transcriptRevision = + transcriptEntryId === "''" + ? transcriptSearchTime + : `CASE WHEN ${transcriptEntryId} = '' THEN ${transcriptSearchTime} ELSE ${ + tableHasColumn(db, "session_transcripts", "content_hash") + ? `COALESCE(session_transcripts.content_hash, ${transcriptSearchTime})` + : transcriptSearchTime + } END`; const transcriptCompleted = transcriptHasCompletedAt ? "session_transcripts.completed_at IS NOT NULL" : "0"; const commonArgs = [...contentArgs, params.agentId, ...sinceArgs, ...beforeArgs, ...deliveredArgs, ...reviewedArgs]; const candidateKinds = @@ -956,8 +967,8 @@ function selectEpisodicSourceRefs( WHERE agent_id = ? AND ${transcriptCompleted} ${params.since ? `AND (julianday(${transcriptSearchTime}) >= julianday(?) OR julianday(${transcriptSearchTime}) < julianday(?))` : ""} ${params.before ? `AND julianday(${transcriptSearchTime}) <= julianday(?)` : ""} - ${deliveredPredicate("transcript", "session_key", transcriptSearchTime, "''", transcriptSearchTime)} - ${reviewedPredicate("transcript", "session_key", transcriptSearchTime, "''", transcriptSearchTime)} + ${deliveredPredicate("transcript", "session_key", transcriptSearchTime, transcriptEntryId, transcriptRevision)} + ${reviewedPredicate("transcript", "session_key", transcriptSearchTime, transcriptEntryId, transcriptRevision)} ${transcriptCandidate.sql} ${transcriptExcluded.sql}`, args: [...commonArgs, ...transcriptCandidate.args, ...transcriptExcluded.args], diff --git a/platform/daemon/src/pipeline/dreaming-attention.ts b/platform/daemon/src/pipeline/dreaming-attention.ts index 46495ab83f..bfca4669fe 100644 --- a/platform/daemon/src/pipeline/dreaming-attention.ts +++ b/platform/daemon/src/pipeline/dreaming-attention.ts @@ -19,6 +19,8 @@ export const DREAMING_CONTENT_ATTENTION_KINDS = [ "surprisal", ] as const; +const AGENT_FACING_ATTENTION_FILTER = "AND NOT (kind = 'evidence_requeue' AND subject_ref LIKE 'source:%')"; + export interface DreamingAttention { readonly id: string; readonly kind: DreamingAttentionKind; @@ -137,7 +139,7 @@ export function getDreamingAttention( // @ts-expect-error LEGACY_SYNC_DB_ACCESS: withReadDb migration site return accessor.withReadDb( (db: import("../db-accessor").ReadDb) => getDreamingAttentionInDb(db, agentId, limit), - "pipeline/dreaming-attention.ts:138", + "pipeline/dreaming-attention.ts:140", ); } @@ -163,7 +165,7 @@ export function getDreamingAttentionWorkloadDiagnostics( pending: row.pending, oldestAgeMs: oldestMs > 0 ? Math.max(0, nowMs - oldestMs) : null, }; - }, "pipeline/dreaming-attention.ts:150"); + }, "pipeline/dreaming-attention.ts:152"); } export async function getDreamingAttentionScoped( accessor: DbAccessor, @@ -172,11 +174,13 @@ export async function getDreamingAttentionScoped( readonly kind?: string; readonly status?: "pending" | "resolved"; readonly limit?: number; + readonly agentFacing?: boolean; }, ): Promise { const boundedLimit = Math.max(1, Math.min(Math.floor(options.limit ?? 20), 100)); const kindFilter = typeof options.kind === "string" && options.kind.length > 0 ? "AND kind = ?" : ""; const statusFilter = options.status === "resolved" ? "AND resolved_at IS NOT NULL" : "AND resolved_at IS NULL"; + const audienceFilter = options.agentFacing === true ? AGENT_FACING_ATTENTION_FILTER : ""; const params: string[] = [agentId]; if (kindFilter && options.kind !== undefined) params.push(options.kind); const rows = await ownerQueryAll<{ @@ -191,7 +195,7 @@ export async function getDreamingAttentionScoped( "dreaming.attention.scoped", `SELECT id, kind, subject_ref AS subjectRef, details_json AS detailsJson, priority, created_at AS createdAt FROM dreaming_attention - WHERE agent_id = ? ${kindFilter} ${statusFilter} + WHERE agent_id = ? ${kindFilter} ${statusFilter} ${audienceFilter} ORDER BY priority DESC, created_at ASC, id ASC LIMIT ?`, [...params, boundedLimit], @@ -208,11 +212,13 @@ export function getDreamingAttentionAcrossScopes( readonly kind?: string; readonly status?: "pending" | "resolved"; readonly limit?: number; + readonly agentFacing?: boolean; }, ): readonly (DreamingAttention & { readonly agentId: string })[] { const boundedLimit = Math.max(1, Math.min(Math.floor(options.limit ?? 50), 200)); const kindFilter = typeof options.kind === "string" && options.kind.length > 0 ? "AND kind = ?" : ""; const statusFilter = options.status === "resolved" ? "AND resolved_at IS NOT NULL" : "AND resolved_at IS NULL"; + const audienceFilter = options.agentFacing === true ? AGENT_FACING_ATTENTION_FILTER : ""; const params: unknown[] = []; if (kindFilter) params.push(options.kind); params.push(boundedLimit); @@ -223,7 +229,7 @@ export function getDreamingAttentionAcrossScopes( `SELECT agent_id AS agentId, id, kind, subject_ref AS subjectRef, details_json AS detailsJson, priority, created_at AS createdAt FROM dreaming_attention - WHERE 1=1 ${kindFilter} ${statusFilter} + WHERE 1=1 ${kindFilter} ${statusFilter} ${audienceFilter} ORDER BY priority DESC, created_at ASC, id ASC LIMIT ?`, ) @@ -240,7 +246,7 @@ export function getDreamingAttentionAcrossScopes( ...attention, details: parseDetails(detailsJson), })); - }, "pipeline/dreaming-attention.ts:220"); + }, "pipeline/dreaming-attention.ts:226"); } export function getDreamingAttentionById( @@ -273,7 +279,7 @@ export function getDreamingAttentionById( priority: row.priority, createdAt: row.createdAt, }; - }, "pipeline/dreaming-attention.ts:251"); + }, "pipeline/dreaming-attention.ts:257"); } export function getDreamingAttentionSnapshots( @@ -284,7 +290,7 @@ export function getDreamingAttentionSnapshots( // @ts-expect-error LEGACY_SYNC_DB_ACCESS: withReadDb migration site return accessor.withReadDb( (db: import("../db-accessor").ReadDb) => getDreamingAttentionSnapshotsInDb(db, agentId, limit), - "pipeline/dreaming-attention.ts:285", + "pipeline/dreaming-attention.ts:291", ); } diff --git a/platform/daemon/src/pipeline/dreaming-capabilities.ts b/platform/daemon/src/pipeline/dreaming-capabilities.ts index 234020914b..29b31e6e4e 100644 --- a/platform/daemon/src/pipeline/dreaming-capabilities.ts +++ b/platform/daemon/src/pipeline/dreaming-capabilities.ts @@ -24,13 +24,17 @@ import { getOntologyLinkEvidence } from "../ontology-link-evidence"; import { type GraphWriteCaps, findDuplicateEntityMerges } from "../ontology-proposals"; import { detectProspectiveContradictionRisk } from "./antonyms"; import { getDreamingAttentionAcrossScopes, getDreamingAttentionScoped } from "./dreaming-attention"; -import { nextDreamingEvidenceFragment, renderDreamingEvidence } from "./dreaming-evidence"; +import { dreamingEvidenceResumeStart, nextDreamingEvidenceFragment, renderDreamingEvidence } from "./dreaming-evidence"; import { - deliveredOffsetForSource, + DREAMING_EVIDENCE_STALL_PASSES, + type DreamingEvidenceCursor, + type DreamingEvidenceReviewRequest, + evidenceCursorForSource, extendDeliveredOffset, passDeliveredRanges, passFullyServedSourceRefs, pendingDreamingEvidenceContinuations, + reviewDreamingEvidenceInDb, } from "./dreaming-evidence-consumption"; import { DREAMING_ONTOLOGY_OPERATION_SCHEMA } from "./dreaming-operation-contract"; import { @@ -48,6 +52,7 @@ const bounded = (value: number | undefined, fallback: number, max: number): numb Math.min(Math.max(Math.floor(value ?? fallback), 1), max); const MAX_EVIDENCE_EXCERPT_CHARS = 2_000; +const EVIDENCE_RESUME_OVERLAP_CHARS = Math.floor(MAX_EVIDENCE_EXCERPT_CHARS / 10); const MAX_EVIDENCE_RESULT_CHARS = 16_000; const MAX_EVIDENCE_PAGE_CHARS = 250_000; @@ -161,6 +166,7 @@ export const DREAMING_CAPABILITY_IDS = [ "get_entity", "list_aspect_claims", "search_evidence", + "review_evidence", "validate_proposal", "zoom_history", "runbook_write", @@ -313,7 +319,16 @@ function drainDreamingEvidenceQueueInDb(db: ReadDb, input: DbOwnerDreamingEviden const servedInPass = input.passId ? passDeliveredRanges(db, input.passId, scopeId) : new Map(); const pageChars = evidencePageChars(input.evidenceChars); let budgetExhausted = false; - const fresh = searchEpisodicSources(db, { + const cursors = new Map(); + const cursorFor = (source: EpisodicSourceRecord): DreamingEvidenceCursor => { + const ref = `${source.kind}:${source.id}`; + const cursor = cursors.get(ref) ?? evidenceCursorForSource(db, scopeId, source); + cursors.set(ref, cursor); + return cursor; + }; + const stalled = (source: EpisodicSourceRecord): boolean => + cursorFor(source).stalledPasses >= DREAMING_EVIDENCE_STALL_PASSES; + const scanned = searchEpisodicSources(db, { agentId: scopeId, query: "", kind: input.kind, @@ -321,6 +336,7 @@ function drainDreamingEvidenceQueueInDb(db: ReadDb, input: DbOwnerDreamingEviden excludeSourceRefs: input.passId ? passFullyServedSourceRefs(db, input.passId, scopeId) : [], limit: DELIVERY_QUEUE_SCAN_LIMIT, }); + const fresh = [...scanned.filter((source) => !stalled(source)), ...scanned.filter(stalled)]; const page = (sources: readonly EpisodicSourceRecord[], max: number, skip = new Set()) => { const items: Record[] = []; let remaining = pageChars; @@ -332,10 +348,17 @@ function drainDreamingEvidenceQueueInDb(db: ReadDb, input: DbOwnerDreamingEviden break; } skip.add(ref); - const offset = extendDeliveredOffset(deliveredOffsetForSource(db, scopeId, source), servedInPass.get(ref)); + const reviewed = cursorFor(source).offset; + const served = extendDeliveredOffset(reviewed, servedInPass.get(ref)); + const offset = + served > reviewed ? served : dreamingEvidenceResumeStart(source, reviewed, EVIDENCE_RESUME_OVERLAP_CHARS); const fragment = projectEvidenceFragment(source, offset, Math.max(remaining, MAX_EVIDENCE_EXCERPT_CHARS)); if (fragment !== null) { - items.push(fragment); + items.push( + reviewed > offset + ? { ...fragment, reviewedChars: Math.min(reviewed - offset, String(fragment.content).length) } + : fragment, + ); remaining -= typeof fragment.content === "string" ? fragment.content.length : 0; } if (items.length >= max) break; @@ -398,9 +421,10 @@ export async function listDreamingAttention( readonly kind?: string; readonly status?: "pending" | "resolved"; readonly limit?: number; + readonly agentFacing?: boolean; }, ): Promise { - const { agentId: scopeId, kind, status, limit } = params; + const { agentId: scopeId, kind, status, limit, agentFacing } = params; if (kind === "review_due") { if (status === "resolved") return []; const input: DbOwnerDreamingReviewDue = { @@ -455,11 +479,13 @@ export async function listDreamingAttention( kind, status: status ?? "pending", limit: bounded(limit, 20, 100), + agentFacing, }) : getDreamingAttentionAcrossScopes(accessor, { kind, status: status ?? "pending", limit: bounded(limit, 50, 200), + agentFacing, }); } @@ -634,7 +660,7 @@ export function createDreamingCapabilities(params: CreateDreamingCapabilitiesPar capability( "search_evidence", "Search episodic evidence", - "Search immutable episodic memories, artifacts, and transcripts in one agent scope across their full history. A query is split on whitespace into words that match independently as substrings (ASCII case-insensitive; unspaced text such as CJK matches as one phrase); sources matching more words rank first, then newer sources. since and before are optional explicit time bounds. Historical summary records can be requested explicitly with kind=summary, but are not part of the default Dreaming delivery path. Results contain exact bounded excerpts of the rendered evidence with contentOffset/contentLength; use sourceRef for citations, which are validated against the complete canonical source. Each record carries completed: memory, artifact, and summary records are settled captures (true); a transcript is true only after the session-end machinery writes its completion marker, and false while the session is still running — do not file claims from a still-growing transcript, since its states may be contradicted by the session's end. When you look up a specific source and contentTruncated is true, page exact fragments with the same sourceRef and chunkSize: start at offset=0 when contentHasPrevious is true, then use offset=contentOffset+content.length from the fragment just returned until contentHasNext is false. Omit query, since, and before to drain the durable delivery queue: it returns up to limit incomplete source revisions, each resuming at its delivered offset (including fragments already served earlier in this pass), regardless of time watermark. hasMore is true while more of the queue remains. A queued source that is only partly read continues on a later queue page, so do not page it yourself. File what each page establishes before calling again without a query for the next one, and stop when hasMore is false. Partway through a pass the queue closes (deliveryClosed: true): stop reading new sources, file what you have read, and finish so your progress is recorded. Narrow with a query if the list is large; pass an explicit earlier since only when you need older history. Artifacts are deduped by content hash: content-identical files across vault paths collapse to one canonical entry.", + "Search immutable episodic memories, artifacts, and transcripts in one agent scope across their full history. A query is split on whitespace into words that match independently as substrings (ASCII case-insensitive; unspaced text such as CJK matches as one phrase); sources matching more words rank first, then newer sources. since and before are optional explicit time bounds. Historical summary records can be requested explicitly with kind=summary, but are not part of the default Dreaming delivery path. Results contain exact bounded excerpts of the rendered evidence with contentOffset/contentLength; use sourceRef for citations, which are validated against the complete canonical source. Each record carries completed: memory, artifact, and summary records are settled captures (true); a transcript is true only after the session-end machinery writes its completion marker, and false while the session is still running — do not file claims from a still-growing transcript, since its states may be contradicted by the session's end. When you look up a specific source and contentTruncated is true, page exact fragments with the same sourceRef and chunkSize: start at offset=0 when contentHasPrevious is true, then use offset=contentOffset+content.length from the fragment just returned until contentHasNext is false. Omit query, since, and before to drain the durable delivery queue: it returns up to limit source revisions not yet fully reviewed, regardless of time watermark. Each resumes where review_evidence acknowledgements ended, after fragments already served earlier in this pass; a source reviewed in an earlier pass resumes slightly before that point, and reviewedChars counts the leading characters already reviewed, which you need not file again. hasMore is true while more of the queue remains. A queued source that is only partly read continues on a later queue page, so do not page it yourself. File what each page establishes and acknowledge it with review_evidence before calling again without a query for the next one, and stop when hasMore is false. Partway through a pass the queue closes (deliveryClosed: true): stop reading new sources, file what you have read, and finish so your progress is recorded. Narrow with a query if the list is large; pass an explicit earlier since only when you need older history. Artifacts are deduped by content hash: content-identical files across vault paths collapse to one canonical entry.", true, z.object({ agentId: z.string().min(1), @@ -698,6 +724,46 @@ export function createDreamingCapabilities(params: CreateDreamingCapabilitiesPar }); }, ), + capability( + "review_evidence", + "Acknowledge reviewed evidence", + "Record which evidence excerpts this pass has reviewed. A source's evidence cursor moves forward only through text acknowledged here or cited by a successful apply_ontology_ops operation; delivered text that is never acknowledged stays queued and is delivered again to a later pass. For each item, copy sourceRef and contentOffset from a search_evidence result delivered in this pass and scope. Omit through when you reviewed the whole excerpt; otherwise set through to an exact quote from the excerpt, and review ends where its first occurrence ends. A call with any invalid item is rejected as a whole with a code (EXCERPT_NOT_DELIVERED, SCOPE_MISMATCH, QUOTE_NOT_IN_EXCERPT, SOURCE_CHANGED); correct it and call again.", + false, + z.object({ + agentId: z.string().min(1), + items: z + .array( + z.object({ + sourceRef: z.string().min(1), + contentOffset: z.number().int().min(0), + through: z.string().trim().min(1).optional(), + }), + ) + .min(1) + .max(50), + }), + async ({ agentId: scopeId, items }) => { + if (!params.passId) + return { ok: false, code: "PASS_REQUIRED", error: "Evidence review requires a live Dreaming pass" }; + const input: DreamingEvidenceReviewRequest = { agentId: scopeId, passId: params.passId, items }; + return await runDbOwnerDomainOperation(accessor, { + runWithOwner: async (owner) => { + const handle = owner.submit( + { kind: "dreaming_evidence_review", input }, + { + operation: "dreaming.capabilities.review-evidence", + lane: "read", + workloadClass: "foreground", + deadlineMs: 30_000, + estimatedWorkUnits: 200, + }, + ); + return await handle.result; + }, + runInline: ({ read }) => read((db) => reviewDreamingEvidenceInDb(db, input)), + }); + }, + ), capability( "validate_proposal", "Validate proposal", @@ -816,7 +882,7 @@ export function createDreamingCapabilities(params: CreateDreamingCapabilitiesPar }), async ({ agentId: scopeId, kind, status, limit }) => ({ ok: true, - items: await listDreamingAttention(accessor, { agentId: scopeId, kind, status, limit }), + items: await listDreamingAttention(accessor, { agentId: scopeId, kind, status, limit, agentFacing: true }), }), ), capability( diff --git a/platform/daemon/src/pipeline/dreaming-evidence-consumption.ts b/platform/daemon/src/pipeline/dreaming-evidence-consumption.ts index 1a56713bd4..ea219cbfed 100644 --- a/platform/daemon/src/pipeline/dreaming-evidence-consumption.ts +++ b/platform/daemon/src/pipeline/dreaming-evidence-consumption.ts @@ -1,4 +1,4 @@ -import { createHash } from "node:crypto"; +import { createHash, randomUUID } from "node:crypto"; import type { ReadDb, WriteDb } from "../db-accessor"; import { type EpisodicSourceKind, type EpisodicSourceRecord, readEpisodicSource } from "../episodic-sources"; import { renderDreamingEvidence } from "./dreaming-evidence"; @@ -15,8 +15,11 @@ export interface DreamingEvidenceDelivery { readonly end: number; readonly length: number; readonly contentSha256: string; + readonly queue: boolean; } +export const DREAMING_EVIDENCE_STALL_PASSES = 3; + export function evidenceContentSha256(content: string): string { return createHash("sha256").update(content).digest("hex"); } @@ -78,9 +81,15 @@ export function persistedEvidenceDeliveries(db: ReadDb, passId: string): readonl } catch { return []; } - const agentId = text(record(input)?.agentId); + const request = record(input); + const agentId = text(request?.agentId); const data = record(output); if (!agentId || data?.ok !== true || !Array.isArray(data.items)) return []; + const queue = + (typeof request?.query !== "string" || request.query.trim() === "") && + request?.since === undefined && + request?.before === undefined && + request?.sourceRef === undefined; return data.items.flatMap((item) => { const row = record(item); const ref = text(row?.sourceRef); @@ -127,6 +136,7 @@ export function persistedEvidenceDeliveries(db: ReadDb, passId: string): readonl end, length, contentSha256: delivered.sha256, + queue, }, ]; }); @@ -174,10 +184,16 @@ export function extendDeliveredOffset( return offset; } +export interface FiledEvidenceCitation { + readonly key: string; + readonly quote: string; +} + export interface FailedOperationEvidence { readonly sources: ReadonlySet; readonly scopes: ReadonlySet; readonly filedSources: ReadonlySet; + readonly filedCitations: readonly FiledEvidenceCitation[]; } export function failedOperationEvidence( @@ -190,8 +206,9 @@ export function failedOperationEvidence( const scopes = new Set(); const filed = new Set(); const filedQuotes = new Set(); + const filedCitations: FiledEvidenceCitation[] = []; const failedQuotes: Array<{ readonly key: string; readonly quote: string }> = []; - if (!tableExists(db, "dreaming_tool_calls")) return { sources: keys, scopes, filedSources: filed }; + if (!tableExists(db, "dreaming_tool_calls")) return { sources: keys, scopes, filedSources: filed, filedCitations }; const rows = db .prepare( `SELECT input_json AS inputJson, output_json AS outputJson @@ -233,6 +250,7 @@ export function failedOperationEvidence( for (const cited of citations(row.index)) { filed.add(cited.key); filedQuotes.add(`${cited.key}\u0000${cited.quote}`); + if (cited.quote) filedCitations.push(cited); } } } @@ -253,7 +271,7 @@ export function failedOperationEvidence( for (const { key, quote } of failedQuotes) { if (!filedQuotes.has(`${key}\u0000${quote}`)) keys.add(key); } - return { sources: keys, scopes, filedSources: filed }; + return { sources: keys, scopes, filedSources: filed, filedCitations }; } export function verifiedDreamingEvidenceDelivery( @@ -280,6 +298,223 @@ export function verifiedDreamingEvidenceDelivery( } return source; } +function deliveredExcerpt(db: ReadDb, delivery: DreamingEvidenceDelivery): string | null { + const source = verifiedDreamingEvidenceDelivery(db, delivery); + return source === null ? null : renderDreamingEvidence(source).slice(delivery.start, delivery.end); +} + +function sameDelivery(a: DreamingEvidenceDelivery, b: DreamingEvidenceDelivery): boolean { + return ( + a.agentId === b.agentId && + a.kind === b.kind && + a.id === b.id && + a.capturedAt === b.capturedAt && + a.sourceEntryId === b.sourceEntryId && + a.sourceRevision === b.sourceRevision && + a.start === b.start && + a.end === b.end && + a.length === b.length && + a.contentSha256 === b.contentSha256 + ); +} + +export interface DreamingEvidenceReviewRequest { + readonly agentId: string; + readonly passId: string; + readonly items: ReadonlyArray<{ + readonly sourceRef: string; + readonly contentOffset: number; + readonly through?: string; + }>; +} + +export type DreamingEvidenceReviewRejection = + | "EXCERPT_NOT_DELIVERED" + | "SCOPE_MISMATCH" + | "QUOTE_NOT_IN_EXCERPT" + | "SOURCE_CHANGED"; + +const REVIEW_REJECTION_ERRORS: Readonly> = { + EXCERPT_NOT_DELIVERED: + "No excerpt with this sourceRef and contentOffset was delivered by search_evidence in this pass; copy both from a search_evidence result", + SCOPE_MISMATCH: + "This excerpt was delivered in a different agent scope; acknowledge it with the agentId used for search_evidence", + QUOTE_NOT_IN_EXCERPT: "through is not an exact quote from this excerpt; copy it character for character", + SOURCE_CHANGED: "The source changed after this excerpt was delivered; read it again with search_evidence", +}; + +export function reviewDreamingEvidenceInDb( + db: ReadDb, + input: DreamingEvidenceReviewRequest, +): { readonly ok: boolean; readonly [key: string]: unknown } { + const deliveries = persistedEvidenceDeliveries(db, input.passId); + const accepted: Record[] = []; + const rejected: Array & { readonly code: DreamingEvidenceReviewRejection }> = []; + input.items.forEach((item, index) => { + const parsed = sourceRef(item.sourceRef); + const sameExcerpt = + parsed === null + ? [] + : deliveries.filter( + (delivery) => + delivery.kind === parsed.kind && delivery.id === parsed.id && delivery.start === item.contentOffset, + ); + const inScope = sameExcerpt.filter((delivery) => delivery.agentId === input.agentId).reverse(); + const reject = (code: DreamingEvidenceReviewRejection): void => { + rejected.push({ + index, + sourceRef: item.sourceRef, + contentOffset: item.contentOffset, + code, + error: REVIEW_REJECTION_ERRORS[code], + }); + }; + if (sameExcerpt.length === 0) { + reject("EXCERPT_NOT_DELIVERED"); + return; + } + if (inScope.length === 0) { + reject("SCOPE_MISMATCH"); + return; + } + const quote = item.through?.trim(); + let verified = false; + for (const delivery of inScope) { + const excerpt = deliveredExcerpt(db, delivery); + if (excerpt === null) continue; + verified = true; + const at = quote === undefined ? 0 : excerpt.indexOf(quote); + if (quote !== undefined && (quote.length === 0 || at < 0)) continue; + accepted.push({ + sourceRef: item.sourceRef, + contentOffset: delivery.start, + reviewedThrough: quote === undefined ? delivery.end : delivery.start + at + quote.length, + excerptEnd: delivery.end, + contentLength: delivery.length, + capturedAt: delivery.capturedAt, + sourceEntryId: delivery.sourceEntryId, + sourceRevision: delivery.sourceRevision, + contentSha256: delivery.contentSha256, + }); + return; + } + reject(verified ? "QUOTE_NOT_IN_EXCERPT" : "SOURCE_CHANGED"); + }); + if (rejected.length > 0) { + return { + ok: false, + code: rejected[0]?.code, + error: `${rejected.length} of ${input.items.length} acknowledgements were rejected and none were recorded; correct them and call review_evidence again`, + items: rejected, + }; + } + return { ok: true, items: accepted }; +} + +interface ReviewedRange { + readonly delivery: DreamingEvidenceDelivery; + readonly end: number; +} + +function persistedEvidenceReviews( + db: ReadDb, + passId: string, + deliveries: readonly DreamingEvidenceDelivery[], +): readonly ReviewedRange[] { + if (!tableExists(db, "dreaming_tool_calls")) return []; + const rows = db + .prepare( + `SELECT input_json AS inputJson, output_json AS outputJson + FROM dreaming_tool_calls + WHERE pass_id = ? AND tool_name = 'review_evidence' ORDER BY sequence ASC`, + ) + .all(passId) as Array<{ inputJson: string; outputJson: string }>; + const integer = (value: unknown): number | null => + typeof value === "number" && Number.isSafeInteger(value) && value >= 0 ? value : null; + return rows.flatMap(({ inputJson, outputJson }) => { + let input: unknown; + let output: unknown; + try { + input = JSON.parse(inputJson); + output = JSON.parse(outputJson); + } catch { + return []; + } + const agentId = text(record(input)?.agentId); + const data = record(output); + if (!agentId || data?.ok !== true || !Array.isArray(data.items)) return []; + return data.items.flatMap((item): ReviewedRange[] => { + const row = record(item); + const ref = text(row?.sourceRef); + const parsed = ref ? sourceRef(ref) : null; + const capturedAt = text(row?.capturedAt); + const revision = text(row?.sourceRevision); + const contentSha256 = text(row?.contentSha256); + const start = integer(row?.contentOffset); + const end = integer(row?.reviewedThrough); + const excerptEnd = integer(row?.excerptEnd); + const length = integer(row?.contentLength); + if ( + !parsed || + !capturedAt || + !revision || + !contentSha256 || + typeof row?.sourceEntryId !== "string" || + start === null || + end === null || + excerptEnd === null || + length === null || + start > end || + end > excerptEnd || + excerptEnd > length + ) + return []; + const reviewed: DreamingEvidenceDelivery = { + agentId, + kind: parsed.kind, + id: parsed.id, + capturedAt, + sourceEntryId: row.sourceEntryId, + sourceRevision: revision, + start, + end: excerptEnd, + length, + contentSha256, + queue: false, + }; + const delivery = deliveries.find((candidate) => sameDelivery(candidate, reviewed)); + return delivery === undefined ? [] : [{ delivery, end }]; + }); + }); +} + +function citationFloors( + db: ReadDb, + deliveries: readonly DreamingEvidenceDelivery[], + citations: readonly FiledEvidenceCitation[], +): readonly ReviewedRange[] { + const excerpts = new Map(); + return citations.flatMap(({ key, quote }) => + deliveries.flatMap((delivery) => { + if (`${delivery.agentId}\u0000${delivery.kind}:${delivery.id}` !== key) return []; + if (!excerpts.has(delivery)) excerpts.set(delivery, deliveredExcerpt(db, delivery)); + const at = excerpts.get(delivery)?.indexOf(quote) ?? -1; + return at < 0 ? [] : [{ delivery, end: delivery.start + at + quote.length }]; + }), + ); +} + +function revisionKey(delivery: DreamingEvidenceDelivery): string { + return [ + delivery.agentId, + delivery.kind, + delivery.id, + delivery.capturedAt, + delivery.sourceEntryId, + delivery.sourceRevision, + ].join("\u0000"); +} + export function recordDreamingEvidenceConsumptionInTx( db: WriteDb, params: { @@ -287,76 +522,183 @@ export function recordDreamingEvidenceConsumptionInTx( readonly deferredEvidence: ReadonlySet; readonly withheldScopes?: ReadonlySet; readonly filedSources?: ReadonlySet; + readonly filedCitations?: readonly FiledEvidenceCitation[]; }, ): void { if (!tableExists(db, "dreaming_evidence_consumption")) return; - const deliveries = persistedEvidenceDeliveries(db, params.passId) - .filter( - (delivery) => - !params.withheldScopes?.has(delivery.agentId) || - params.filedSources?.has(`${delivery.agentId}\u0000${delivery.kind}:${delivery.id}`) === true, - ) - .filter((delivery) => !params.deferredEvidence.has(`${delivery.agentId}\u0000${delivery.kind}:${delivery.id}`)) - .sort( - (a, b) => - a.agentId.localeCompare(b.agentId) || - a.kind.localeCompare(b.kind) || - a.id.localeCompare(b.id) || - a.capturedAt.localeCompare(b.capturedAt) || - a.start - b.start || - a.end - b.end, - ); + const deliveries = persistedEvidenceDeliveries(db, params.passId); + const revisions = new Map< + string, + { readonly delivery: DreamingEvidenceDelivery; readonly ranges: Array; queued: boolean } + >(); + for (const delivery of deliveries) { + const key = revisionKey(delivery); + const entry = revisions.get(key) ?? { delivery, ranges: [], queued: false }; + entry.queued ||= delivery.queue; + revisions.set(key, entry); + } + const reviewed = [ + ...persistedEvidenceReviews(db, params.passId, deliveries), + ...citationFloors(db, deliveries, params.filedCitations ?? []), + ]; + const progressSuppressed = (delivery: DreamingEvidenceDelivery): boolean => { + const sourceKey = `${delivery.agentId}\u0000${delivery.kind}:${delivery.id}`; + if (params.deferredEvidence.has(sourceKey)) return true; + return params.withheldScopes?.has(delivery.agentId) === true && params.filedSources?.has(sourceKey) !== true; + }; + for (const { delivery, end } of reviewed) { + if (progressSuppressed(delivery)) continue; + if (verifiedDreamingEvidenceDelivery(db, delivery) === null) continue; + revisions.get(revisionKey(delivery))?.ranges.push([delivery.start, end] as const); + } const select = db.prepare( - `SELECT delivered_offset AS deliveredOffset FROM dreaming_evidence_consumption + `SELECT delivered_offset AS deliveredOffset, stalled_passes AS stalledPasses FROM dreaming_evidence_consumption WHERE agent_id = ? AND source_kind = ? AND source_id = ? AND source_captured_at = ? AND source_entry_id = ? AND source_revision = ?`, ); - const upsert = db.prepare( + const advance = db.prepare( `INSERT INTO dreaming_evidence_consumption - (agent_id, source_kind, source_id, source_captured_at, source_entry_id, source_revision, delivered_offset, source_length, pass_id, updated_at) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, datetime('now')) + (agent_id, source_kind, source_id, source_captured_at, source_entry_id, source_revision, delivered_offset, source_length, pass_id, updated_at, cursor_basis, stalled_passes) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, datetime('now'), 'review', 0) ON CONFLICT(agent_id, source_kind, source_id, source_captured_at, source_entry_id, source_revision) DO UPDATE SET delivered_offset = excluded.delivered_offset, source_length = excluded.source_length, pass_id = excluded.pass_id, - updated_at = excluded.updated_at`, + updated_at = excluded.updated_at, + cursor_basis = 'review', + stalled_passes = 0`, ); - for (const delivery of deliveries) { + const stall = db.prepare( + `INSERT INTO dreaming_evidence_consumption + (agent_id, source_kind, source_id, source_captured_at, source_entry_id, source_revision, delivered_offset, source_length, pass_id, updated_at, cursor_basis, stalled_passes) + VALUES (?, ?, ?, ?, ?, ?, 0, ?, ?, datetime('now'), 'review', 1) + ON CONFLICT(agent_id, source_kind, source_id, source_captured_at, source_entry_id, source_revision) DO UPDATE SET + stalled_passes = dreaming_evidence_consumption.stalled_passes + 1`, + ); + const attention = tableExists(db, "dreaming_attention"); + const raiseStall = db.prepare( + `INSERT INTO dreaming_attention (id, agent_id, kind, subject_ref, details_json, priority) + VALUES (?, ?, 'evidence_requeue', ?, ?, 60) + ON CONFLICT(agent_id, kind, subject_ref) DO UPDATE SET + details_json = excluded.details_json, + priority = MAX(dreaming_attention.priority, excluded.priority), + generation = dreaming_attention.generation + 1, + resolved_at = NULL, + resolved_by_pass_id = NULL + WHERE dreaming_attention.resolved_at IS NOT NULL + OR json_extract(dreaming_attention.details_json, '$.reason') = 'evidence-stalled'`, + ); + const ordered = [...revisions.values()].sort( + (a, b) => + a.delivery.agentId.localeCompare(b.delivery.agentId) || + a.delivery.kind.localeCompare(b.delivery.kind) || + a.delivery.id.localeCompare(b.delivery.id) || + a.delivery.capturedAt.localeCompare(b.delivery.capturedAt), + ); + for (const { delivery, ranges, queued } of ordered) { const source = verifiedDreamingEvidenceDelivery(db, delivery); if (source === null) continue; - const identity = sourceIdentity(source); - const revision = sourceRevision(source); - const row = select.get(delivery.agentId, delivery.kind, delivery.id, delivery.capturedAt, identity, revision) as { - deliveredOffset: number; - } | null; - const current = row?.deliveredOffset ?? 0; - if (delivery.start > current) continue; - const next = Math.max(current, delivery.end); - if (next <= current && row != null) continue; - upsert.run( + const identity = [ delivery.agentId, delivery.kind, delivery.id, delivery.capturedAt, - identity, - revision, - Math.min(next, delivery.length), - delivery.length, - params.passId, + sourceIdentity(source), + sourceRevision(source), + ] as const; + const row = select.get(...identity) as { deliveredOffset: number; stalledPasses: number } | null; + const current = Math.max(0, row?.deliveredOffset ?? 0); + const next = Math.min(extendDeliveredOffset(current, ranges), delivery.length); + const ref = `${delivery.kind}:${delivery.id}`; + if (next > current) { + advance.run(...identity, next, delivery.length, params.passId); + if (attention) resolveStalledEvidenceAttentionInTx(db, params.passId, delivery.agentId, ref); + continue; + } + if (!queued || current >= delivery.length || progressSuppressed(delivery)) continue; + stall.run(...identity, delivery.length, params.passId); + const stalledPasses = (row?.stalledPasses ?? 0) + 1; + if (!attention || stalledPasses < DREAMING_EVIDENCE_STALL_PASSES) continue; + raiseStall.run( + randomUUID(), + delivery.agentId, + ref, + JSON.stringify({ + reason: "evidence-stalled", + sourceRef: ref, + sourceRevision: delivery.sourceRevision, + reviewedChars: String(current), + sourceLength: String(delivery.length), + stalledPasses: String(stalledPasses), + }), ); } } -export function deliveredOffsetForSource(db: ReadDb, agentId: string, source: EpisodicSourceRecord): number { - if (!tableExists(db, "dreaming_evidence_consumption")) return 0; +export const STALLED_EVIDENCE_ATTENTION_SQL = + "(kind = 'evidence_requeue' AND json_valid(details_json) AND json_extract(details_json, '$.reason') = 'evidence-stalled')"; + +const IMPORTED_SOURCE_ATTENTION_RAISED_AT_SQL = + "COALESCE(CASE WHEN json_valid(details_json) THEN json_extract(details_json, '$.raisedAt') END, created_at)"; + +export const SEEN_IMPORTED_SOURCE_ATTENTION_SQL = `(kind = 'evidence_requeue' AND subject_ref LIKE 'source:%' AND COALESCE( + julianday(${IMPORTED_SOURCE_ATTENTION_RAISED_AT_SQL}) < ( + SELECT MAX(julianday(dp.started_at)) FROM dreaming_passes dp WHERE dp.agent_id = dreaming_attention.agent_id + ), 0))`; + +export function importedSourceAttentionUpsert( + agentId: string, + sourceId: string, +): { readonly sql: string; readonly params: readonly string[] } { + return { + sql: `INSERT INTO dreaming_attention (id, agent_id, kind, subject_ref, details_json, priority) + VALUES (?, ?, 'evidence_requeue', ?, json_set(?, '$.raisedAt', strftime('%Y-%m-%d %H:%M:%f', 'now')), 50) + ON CONFLICT(agent_id, kind, subject_ref) DO UPDATE SET + details_json = excluded.details_json, + priority = MAX(dreaming_attention.priority, excluded.priority), + generation = dreaming_attention.generation + 1, + resolved_at = NULL, + resolved_by_pass_id = NULL`, + params: [ + randomUUID(), + agentId, + `source:${sourceId}`, + JSON.stringify({ sourceId, reason: "transcript-import-committed" }), + ], + }; +} + +export function resolveStalledEvidenceAttentionInTx(db: WriteDb, passId: string, agentId: string, ref: string): void { + db.prepare( + `UPDATE dreaming_attention SET resolved_at = datetime('now'), resolved_by_pass_id = ? + WHERE agent_id = ? AND subject_ref = ? AND resolved_at IS NULL AND ${STALLED_EVIDENCE_ATTENTION_SQL}`, + ).run(passId, agentId, ref); +} + +export interface DreamingEvidenceCursor { + readonly offset: number; + readonly stalledPasses: number; +} + +export function evidenceCursorForSource( + db: ReadDb, + agentId: string, + source: EpisodicSourceRecord, +): DreamingEvidenceCursor { + if (!tableExists(db, "dreaming_evidence_consumption")) return { offset: 0, stalledPasses: 0 }; const row = db .prepare( - `SELECT delivered_offset AS deliveredOffset FROM dreaming_evidence_consumption + `SELECT delivered_offset AS deliveredOffset, stalled_passes AS stalledPasses FROM dreaming_evidence_consumption WHERE agent_id = ? AND source_kind = ? AND source_id = ? AND source_captured_at = ? AND source_entry_id = ? AND source_revision = ?`, ) .get(agentId, source.kind, source.id, source.capturedAt, sourceIdentity(source), sourceRevision(source)) as { deliveredOffset: number; + stalledPasses: number; } | null; - return Math.max(0, row?.deliveredOffset ?? 0); + return { offset: Math.max(0, row?.deliveredOffset ?? 0), stalledPasses: row?.stalledPasses ?? 0 }; +} + +export function deliveredOffsetForSource(db: ReadDb, agentId: string, source: EpisodicSourceRecord): number { + return evidenceCursorForSource(db, agentId, source).offset; } export function hasDreamingEvidenceContinuation(db: ReadDb, agentId: string, passId: string | null): boolean { if (!passId || !tableExists(db, "dreaming_evidence_consumption")) return false; @@ -415,6 +757,7 @@ export function pendingDreamingEvidenceContinuations( INNER JOIN dreaming_passes pass ON pass.id = dec.pass_id WHERE dec.agent_id = ? AND dec.delivered_offset > 0 AND dec.delivered_offset < dec.source_length + AND dec.stalled_passes < ? ${reviewedPredicate} AND (? IS NULL OR dec.source_kind = ?) AND ( @@ -457,7 +800,7 @@ export function pendingDreamingEvidenceContinuations( ORDER BY pass.rowid ASC, dec.source_kind ASC, dec.source_id ASC, dec.source_captured_at ASC LIMIT ?`, ) - .all(agentId, kind ?? null, kind ?? null, boundedLimit) as Array<{ + .all(agentId, DREAMING_EVIDENCE_STALL_PASSES, kind ?? null, kind ?? null, boundedLimit) as Array<{ kind: EpisodicSourceKind; id: string; capturedAt: string; @@ -476,53 +819,152 @@ export function pendingDreamingEvidenceContinuations( return [source]; }); } -export function countEligibleUnconsumedEvidenceForSource( +function candidateHasUnconsumedEvidence(db: ReadDb, agentId: string, kind: EpisodicSourceKind, id: string): boolean { + const source = readEpisodicSource(db, { agentId, from: `${kind}:${id}` }); + if (source === null) return false; + const reviewed = + tableExists(db, "dreaming_evidence_reviews") && + db + .prepare( + `SELECT 1 FROM dreaming_evidence_reviews WHERE agent_id = ? AND source_kind = ? AND source_id = ? AND source_captured_at = ? AND source_entry_id = ? AND source_revision = ?`, + ) + .get(agentId, source.kind, source.id, source.capturedAt, sourceIdentity(source), sourceRevision(source)) != null; + return !reviewed && deliveredOffsetForSource(db, agentId, source) < renderDreamingEvidence(source).length; +} + +interface SourceCandidateBranch { + readonly kind: "artifact" | "transcript"; + readonly select: string; + readonly args: readonly unknown[]; + readonly id: string; + readonly identity: string; +} + +function sourceCandidateBranches( db: ReadDb, agentId: string, sourceEntryId: string, - _legacyObsidianRoot?: string, -): number { - if (!tableExists(db, "dreaming_evidence_consumption")) return 1; - const legacyRootPrefix = _legacyObsidianRoot?.replace(/\\/g, "/").replace(/\/$/, "") ?? null; - const candidates: Array<{ kind: EpisodicSourceKind; id: string }> = [ - ...( - db - .prepare( - `SELECT source_path AS id FROM memory_artifacts - WHERE agent_id = ? AND COALESCE(is_deleted, 0) = 0 AND length(content) > 0 - AND (source_id = ? OR (? IS NOT NULL AND harness = 'obsidian' AND source_id IS NULL AND source_path >= ? AND source_path < ?))`, - ) - .all( - agentId, - sourceEntryId, - legacyRootPrefix, - legacyRootPrefix ?? "", - `${legacyRootPrefix ?? ""}/\uffff`, - ) as Array<{ id: string }> - ).map((row) => ({ kind: "artifact" as const, id: row.id })), - ...( - db - .prepare( - "SELECT session_key AS id FROM session_transcripts WHERE agent_id = ? AND source_id = ? AND completed_at IS NOT NULL", - ) - .all(agentId, sourceEntryId) as Array<{ id: string }> - ).map((row) => ({ kind: "transcript" as const, id: row.id })), + legacyObsidianRoot: string | undefined, +): readonly SourceCandidateBranch[] { + const legacyRootPrefix = legacyObsidianRoot?.replace(/\\/g, "/").replace(/\/$/, "") ?? null; + const identity = ( + kind: string, + alias: string, + id: string, + capturedAt: string, + entryId: string, + revision: string, + ): string => + `e.agent_id = ${alias}.agent_id AND e.source_kind = '${kind}' AND e.source_id = ${id} + AND e.source_captured_at = ${capturedAt} AND e.source_entry_id = ${entryId} AND e.source_revision = ${revision}`; + const branches: SourceCandidateBranch[] = [ + { + kind: "artifact", + select: `SELECT ma.source_path AS id FROM memory_artifacts ma + WHERE ma.agent_id = ? AND COALESCE(ma.is_deleted, 0) = 0 AND length(ma.content) > 0 + AND (ma.source_id = ? OR (? IS NOT NULL AND ma.harness = 'obsidian' AND ma.source_id IS NULL AND ma.source_path >= ? AND ma.source_path < ?))`, + args: [agentId, sourceEntryId, legacyRootPrefix, legacyRootPrefix ?? "", `${legacyRootPrefix ?? ""}/\uffff`], + id: "ma.source_path", + identity: identity( + "artifact", + "ma", + "ma.source_path", + "ma.captured_at", + "COALESCE(ma.source_id, '')", + "CASE WHEN ma.source_sha256 IS NULL OR ma.source_sha256 = '' THEN ma.captured_at ELSE ma.source_sha256 END", + ), + }, ]; - return candidates.reduce((count, candidate) => { - const source = readEpisodicSource(db, { agentId, from: `${candidate.kind}:${candidate.id}` }); - if (source === null) return count; - const reviewed = - tableExists(db, "dreaming_evidence_reviews") && - db - .prepare( - `SELECT 1 FROM dreaming_evidence_reviews WHERE agent_id = ? AND source_kind = ? AND source_id = ? AND source_captured_at = ? AND source_entry_id = ? AND source_revision = ?`, - ) - .get(agentId, source.kind, source.id, source.capturedAt, sourceIdentity(source), sourceRevision(source)) != - null; - return reviewed || deliveredOffsetForSource(db, agentId, source) >= renderDreamingEvidence(source).length - ? count - : count + 1; - }, 0); + if ( + tableHasColumn(db, "session_transcripts", "source_id") && + tableHasColumn(db, "session_transcripts", "completed_at") + ) { + const updatedAt = tableHasColumn(db, "session_transcripts", "updated_at") ? "st.updated_at" : "NULL"; + const contentHash = tableHasColumn(db, "session_transcripts", "content_hash") ? "st.content_hash" : "NULL"; + const capturedAt = `COALESCE(st.completed_at, ${updatedAt}, st.created_at)`; + branches.push({ + kind: "transcript", + select: `SELECT st.session_key AS id FROM session_transcripts st + WHERE st.agent_id = ? AND st.source_id = ? AND st.completed_at IS NOT NULL`, + args: [agentId, sourceEntryId], + id: "st.session_key", + identity: identity( + "transcript", + "st", + "st.session_key", + capturedAt, + "COALESCE(st.source_id, '')", + `CASE WHEN st.source_id IS NULL OR st.source_id = '' THEN ${capturedAt} ELSE COALESCE(${contentHash}, ${capturedAt}) END`, + ), + }); + } + return branches; +} + +export interface SourceEvidenceDrainProbe { + readonly status: "pending" | "drained" | "undetermined"; + readonly resumeAfter: string | null; + readonly rendered: number; +} + +const SOURCE_DRAIN_PAGE_ROWS = 32; + +export function probeSourceEvidenceDrain( + db: ReadDb, + agentId: string, + sourceEntryId: string, + options: { + readonly legacyObsidianRoot?: string; + readonly maxRenders: number; + readonly resumeAfter?: string | null; + }, +): SourceEvidenceDrainProbe { + let resumeAfter = options.resumeAfter ?? null; + if (!tableExists(db, "dreaming_evidence_consumption")) return { status: "pending", resumeAfter, rendered: 0 }; + const branches = sourceCandidateBranches(db, agentId, sourceEntryId, options.legacyObsidianRoot); + const notReviewed = (branch: SourceCandidateBranch): string => + tableExists(db, "dreaming_evidence_reviews") + ? `AND NOT EXISTS (SELECT 1 FROM dreaming_evidence_reviews e WHERE ${branch.identity})` + : ""; + for (const branch of branches) { + const partial = db + .prepare( + `${branch.select} ${notReviewed(branch)} + AND EXISTS (SELECT 1 FROM dreaming_evidence_consumption e WHERE ${branch.identity} AND e.delivered_offset < e.source_length) + LIMIT 1`, + ) + .get(...branch.args); + if (partial != null) return { status: "pending", resumeAfter, rendered: 0 }; + } + const separator = resumeAfter?.indexOf(":") ?? -1; + const resumeKind = resumeAfter !== null && separator > 0 ? resumeAfter.slice(0, separator) : null; + const resumeId = resumeAfter !== null && separator > 0 ? resumeAfter.slice(separator + 1) : ""; + let rendered = 0; + for (const branch of branches) { + if (resumeKind === "transcript" && branch.kind === "artifact") continue; + let after = resumeKind === branch.kind ? resumeId : ""; + const undelivered = db.prepare( + `${branch.select} AND ${branch.id} > ? ${notReviewed(branch)} + AND NOT EXISTS (SELECT 1 FROM dreaming_evidence_consumption e WHERE ${branch.identity}) + ORDER BY ${branch.id} ASC + LIMIT ?`, + ); + for (;;) { + const pageRows = Math.max(1, Math.min(SOURCE_DRAIN_PAGE_ROWS, options.maxRenders - rendered + 1)); + const rows = undelivered.all(...branch.args, after, pageRows) as Array<{ id: string }>; + for (const row of rows) { + if (rendered >= options.maxRenders) return { status: "undetermined", resumeAfter, rendered }; + rendered += 1; + if (candidateHasUnconsumedEvidence(db, agentId, branch.kind, row.id)) { + return { status: "pending", resumeAfter, rendered }; + } + after = row.id; + resumeAfter = `${branch.kind}:${row.id}`; + } + if (rows.length < pageRows) break; + } + } + return { status: "drained", resumeAfter, rendered }; } export function sourceHasEligibleUnconsumedEvidence( @@ -531,5 +973,68 @@ export function sourceHasEligibleUnconsumedEvidence( sourceEntryId: string, legacyObsidianRoot?: string, ): boolean { - return countEligibleUnconsumedEvidenceForSource(db, agentId, sourceEntryId, legacyObsidianRoot) > 0; + return ( + probeSourceEvidenceDrain(db, agentId, sourceEntryId, { + legacyObsidianRoot, + maxRenders: Number.POSITIVE_INFINITY, + }).status !== "drained" + ); +} + +export const IMPORTED_SOURCE_ATTENTION_ROWS_PER_SCOPE = 20; +export const IMPORTED_SOURCE_ATTENTION_RENDER_BUDGET = 8; + +export function resolveImportedSourceAttentionInTx(db: WriteDb, passId: string, scopes: readonly string[]): number { + if (!tableExists(db, "dreaming_attention")) return 0; + const checkSeq = "(CASE WHEN json_valid(details_json) THEN json_extract(details_json, '$.drainCheckSeq') END)"; + const pendingSource = + "agent_id = ? AND kind = 'evidence_requeue' AND resolved_at IS NULL AND subject_ref LIKE 'source:%'"; + const pending = db.prepare( + `SELECT id, subject_ref AS subjectRef, + CASE WHEN json_valid(details_json) THEN json_extract(details_json, '$.drainResumeAfter') END AS resumeAfter + FROM dreaming_attention + WHERE ${pendingSource} + ORDER BY COALESCE(${checkSeq}, 0) ASC, created_at ASC, id ASC + LIMIT ?`, + ); + const nextCheckSeq = db.prepare( + `SELECT COALESCE(MAX(${checkSeq}), 0) + 1 AS seq FROM dreaming_attention WHERE ${pendingSource}`, + ); + const stamp = db.prepare( + `UPDATE dreaming_attention + SET details_json = json_set(CASE WHEN json_valid(details_json) THEN details_json ELSE '{}' END, + '$.drainCheckSeq', ?, '$.drainResumeAfter', ?) + WHERE id = ?`, + ); + const resolve = db.prepare( + `UPDATE dreaming_attention SET resolved_at = datetime('now'), resolved_by_pass_id = ? + WHERE id = ? AND resolved_at IS NULL`, + ); + let resolved = 0; + let renderBudget = IMPORTED_SOURCE_ATTENTION_RENDER_BUDGET; + for (const agentId of new Set(scopes)) { + const rows = pending.all(agentId, IMPORTED_SOURCE_ATTENTION_ROWS_PER_SCOPE) as Array<{ + id: string; + subjectRef: string; + resumeAfter: unknown; + }>; + const seq = (nextCheckSeq.get(agentId) as { seq: number }).seq; + for (const row of rows) { + const sourceEntryId = row.subjectRef.slice("source:".length); + const probe = sourceEntryId + ? probeSourceEvidenceDrain(db, agentId, sourceEntryId, { + maxRenders: renderBudget, + resumeAfter: typeof row.resumeAfter === "string" ? row.resumeAfter : null, + }) + : null; + renderBudget -= probe?.rendered ?? 0; + if (probe?.status === "drained") { + resolve.run(passId, row.id); + resolved += 1; + continue; + } + stamp.run(seq, probe?.resumeAfter ?? null, row.id); + } + } + return resolved; } diff --git a/platform/daemon/src/pipeline/dreaming-evidence-reviews.ts b/platform/daemon/src/pipeline/dreaming-evidence-reviews.ts index 3075a6303c..cb9f5e7b86 100644 --- a/platform/daemon/src/pipeline/dreaming-evidence-reviews.ts +++ b/platform/daemon/src/pipeline/dreaming-evidence-reviews.ts @@ -7,6 +7,7 @@ import { enqueueDreamingAttentionInTx } from "./dreaming-attention"; import { deliveredOffsetForSource, persistedEvidenceDeliveries, + resolveStalledEvidenceAttentionInTx, verifiedDreamingEvidenceDelivery, } from "./dreaming-evidence-consumption"; import { renderDreamingEvidence } from "./dreaming-evidence"; @@ -130,6 +131,8 @@ export function recordDreamingReviewedExcludedEvidenceInTx( pass_id = excluded.pass_id, reviewed_at = excluded.reviewed_at`, ); + const attention = + db.prepare("SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = 'dreaming_attention'").get() != null; const seen = new Set(); let recorded = 0; for (const entry of params.entries) { @@ -173,6 +176,7 @@ export function recordDreamingReviewedExcludedEvidenceInTx( entry.reason, params.passId, ); + if (attention) resolveStalledEvidenceAttentionInTx(db, params.passId, agentId, `${source.kind}:${source.id}`); recorded += 1; } return recorded; @@ -229,6 +233,6 @@ export async function requestDreamingReviewedEvidenceRequeue( }); return true; }, - { siteToken: "pipeline/dreaming-evidence-reviews.ts:232" }, + { siteToken: "pipeline/dreaming-evidence-reviews.ts:236" }, ); } diff --git a/platform/daemon/src/pipeline/dreaming-evidence.ts b/platform/daemon/src/pipeline/dreaming-evidence.ts index 604ee1eee1..1719cd70ee 100644 --- a/platform/daemon/src/pipeline/dreaming-evidence.ts +++ b/platform/daemon/src/pipeline/dreaming-evidence.ts @@ -15,6 +15,20 @@ export interface DreamingEvidenceFragment { readonly end: number; readonly sourceLength: number; } +function boundaryEndAt(content: string, index: number): number | null { + const character = content[index]; + const previous = content[index - 1]; + if (character === undefined || previous === undefined) return null; + if (!((character === "\n" && previous === "\n") || (/\s/.test(character) && /[.!?]/.test(previous)))) return null; + let boundaryEnd = index + 1; + while (boundaryEnd < content.length) { + const next = content[boundaryEnd]; + if (next === undefined || !/\s/.test(next)) break; + boundaryEnd += 1; + } + return boundaryEnd; +} + export function nextDreamingEvidenceFragment( source: EpisodicSourceRecord, start: number, @@ -26,26 +40,30 @@ export function nextDreamingEvidenceFragment( let end = cappedEnd; if (cappedEnd < content.length) { for (let index = cappedEnd - 1; index > start; index -= 1) { - const character = content[index]; - const previous = content[index - 1]; - if (character === undefined || previous === undefined) continue; - if ((character === "\n" && previous === "\n") || (/\s/.test(character) && /[.!?]/.test(previous))) { - let boundaryEnd = index + 1; - while (boundaryEnd < content.length) { - const next = content[boundaryEnd]; - if (next === undefined || !/\s/.test(next)) break; - boundaryEnd += 1; - } - if (boundaryEnd <= cappedEnd && content.slice(start, boundaryEnd).trim().length > 0) { - end = boundaryEnd; - break; - } + const boundaryEnd = boundaryEndAt(content, index); + if (boundaryEnd !== null && boundaryEnd <= cappedEnd && content.slice(start, boundaryEnd).trim().length > 0) { + end = boundaryEnd; + break; } } } return { source, content: content.slice(start, end), start, end, sourceLength: content.length }; } +export function dreamingEvidenceResumeStart( + source: EpisodicSourceRecord, + reviewedEnd: number, + overlapChars: number, +): number { + const content = renderDreamingEvidence(source); + const start = Math.max(0, Math.min(reviewedEnd, content.length) - Math.max(0, Math.floor(overlapChars))); + for (let index = start - 1; index > Math.max(0, start - overlapChars); index -= 1) { + const boundaryEnd = boundaryEndAt(content, index); + if (boundaryEnd !== null && boundaryEnd <= start) return boundaryEnd; + } + return start; +} + export function completeDreamingEvidenceFragment(source: EpisodicSourceRecord): DreamingEvidenceFragment { const content = renderDreamingEvidence(source); return { source, content, start: 0, end: content.length, sourceLength: content.length }; diff --git a/platform/daemon/src/pipeline/dreaming.test.ts b/platform/daemon/src/pipeline/dreaming.test.ts index cd183a0dd5..1d6d294f83 100644 --- a/platform/daemon/src/pipeline/dreaming.test.ts +++ b/platform/daemon/src/pipeline/dreaming.test.ts @@ -51,7 +51,15 @@ import { getDreamingAttentionSnapshots, resolveDreamingAttentionInTx, } from "./dreaming-attention"; -import { pendingDreamingEvidenceContinuations } from "./dreaming-evidence-consumption"; +import { + DREAMING_EVIDENCE_STALL_PASSES, + IMPORTED_SOURCE_ATTENTION_RENDER_BUDGET, + IMPORTED_SOURCE_ATTENTION_ROWS_PER_SCOPE, + importedSourceAttentionUpsert, + pendingDreamingEvidenceContinuations, + probeSourceEvidenceDrain, + resolveImportedSourceAttentionInTx, +} from "./dreaming-evidence-consumption"; import { renderDreamingEvidence } from "./dreaming-evidence"; import { readEpisodicMemory, searchEpisodicSources, utcTimestampMs } from "../episodic-sources"; import { @@ -262,6 +270,30 @@ async function invokeDreamingTool( return parsed as Record; } +async function reviewDreamingEvidence( + input: DreamingAgentInput, + page: Record, + agentId = AGENT, +): Promise> { + const items = (Array.isArray(page.items) ? page.items : []).map((item) => ({ + sourceRef: String(Reflect.get(item as object, "sourceRef")), + contentOffset: Number(Reflect.get(item as object, "contentOffset")), + })); + const result = await invokeDreamingTool(input, "review_evidence", { agentId, items }); + if (result.ok !== true) throw new Error(`Test evidence review failed: ${JSON.stringify(result)}`); + return result; +} + +async function readAndReviewEvidence( + input: DreamingAgentInput, + args: Record, +): Promise> { + const page = await invokeDreamingTool(input, "search_evidence", args); + if (Array.isArray(page.items) && page.items.length > 0) + await reviewDreamingEvidence(input, page, typeof args.agentId === "string" ? args.agentId : AGENT); + return page; +} + type DreamingTestHeadSupport = { readonly agentId?: string; readonly sourceRef: string; @@ -543,7 +575,7 @@ describe("Dreaming", () => { prompt = input.prompt; const search = input.tools.find((tool) => tool.name === "search_evidence"); if (!search) throw new Error("Missing search_evidence"); - await search.execute("call", { agentId: AGENT }, undefined, undefined, {} as never); + await readAndReviewEvidence(input, { agentId: AGENT }); return { summary: "Done" }; }, }, @@ -581,7 +613,14 @@ describe("Dreaming", () => { (call) => call.toolName === "search_evidence", ); expect(secondDelivery?.output).toMatchObject({ - items: [expect.objectContaining({ sourceRef: "transcript:s1", contentOffset: 16_000, contentLength: 40_000 })], + items: [ + expect.objectContaining({ + sourceRef: "transcript:s1", + contentOffset: 15_800, + contentLength: 40_000, + reviewedChars: 200, + }), + ], }); expect(await getDreamingEpisodicTokenBacklog(accessor, AGENT)).toBeGreaterThan(0); await run(); @@ -696,7 +735,7 @@ describe("Dreaming", () => { const search = input.tools.find((tool) => tool.name === "search_evidence"); const runbook = input.tools.find((tool) => tool.name === "runbook_write"); if (!search || !runbook) throw new Error("Missing Dreaming evidence tools"); - await search.execute("call", { agentId: AGENT }, undefined, undefined, {} as never); + await readAndReviewEvidence(input, { agentId: AGENT }); if (passNumber === 3) { await runbook.execute( "call", @@ -927,8 +966,8 @@ describe("Dreaming", () => { const search = input.tools.find((tool) => tool.name === "search_evidence"); const write = input.tools.find((tool) => tool.name === "runbook_write"); if (!search || !write) throw new Error("Missing Dreaming delivery tools"); - await search.execute("call", { agentId: AGENT }, undefined, undefined, {} as never); - await search.execute("call", { agentId: otherAgent }, undefined, undefined, {} as never); + await readAndReviewEvidence(input, { agentId: AGENT }); + await readAndReviewEvidence(input, { agentId: otherAgent }); await write.execute( "call", { @@ -955,7 +994,9 @@ describe("Dreaming", () => { expect( ( db - .prepare("SELECT COUNT(*) AS count FROM dreaming_evidence_consumption WHERE agent_id = ?") + .prepare( + "SELECT COUNT(*) AS count FROM dreaming_evidence_consumption WHERE agent_id = ? AND delivered_offset > 0", + ) .get(otherAgent) as { count: number; } @@ -963,7 +1004,11 @@ describe("Dreaming", () => { ).toBe(0); expect( ( - db.prepare("SELECT COUNT(*) AS count FROM dreaming_evidence_consumption WHERE agent_id = ?").get(AGENT) as { + db + .prepare( + "SELECT COUNT(*) AS count FROM dreaming_evidence_consumption WHERE agent_id = ? AND delivered_offset > 0", + ) + .get(AGENT) as { count: number; } ).count, @@ -1170,7 +1215,7 @@ describe("Dreaming", () => { `INSERT INTO dreaming_evidence_consumption (agent_id, source_kind, source_id, source_captured_at, source_entry_id, source_revision, delivered_offset, source_length, pass_id, updated_at) - VALUES (?, 'transcript', 'continuation-transcript', ?, '', ?, 2_000, 5_000, ?, ?)`, + VALUES (?, 'transcript', 'continuation-transcript', ?, '', ?, 2000, 5000, ?, ?)`, ).run(AGENT, capturedAt, capturedAt, passId, capturedAt); }); const cfg = defaultCfg({ tokenThreshold: 100_000, backfillOnFirstRun: false }); @@ -1201,7 +1246,7 @@ describe("Dreaming", () => { }); expect(await shouldTriggerDreaming(accessor, cfg, AGENT, now)).toBe(false); accessor.withWriteTx((tx) => { - tx.prepare("UPDATE dreaming_evidence_consumption SET delivered_offset = 2_000 WHERE pass_id = ?").run(passId); + tx.prepare("UPDATE dreaming_evidence_consumption SET delivered_offset = 2000 WHERE pass_id = ?").run(passId); }); await recordDreamingFailure(accessor, AGENT); expect(await shouldTriggerDreaming(accessor, cfg, AGENT, now)).toBe(false); @@ -1218,7 +1263,7 @@ describe("Dreaming", () => { async run(input) { const search = input.tools.find((tool) => tool.name === "search_evidence"); if (!search) throw new Error("Missing search_evidence"); - await search.execute("call", { agentId: AGENT }, undefined, undefined, {} as never); + await readAndReviewEvidence(input, { agentId: AGENT }); await commitDreamingTestHead(input, { sourceRef: "transcript:partial-frontier", text: "The partial frontier transcript was surfaced for review.", @@ -1249,7 +1294,13 @@ describe("Dreaming", () => { (call) => call.toolName === "search_evidence", ); expect(delivery?.output).toMatchObject({ - items: [expect.objectContaining({ sourceRef: "transcript:partial-frontier", contentOffset: 16_000 })], + items: [ + expect.objectContaining({ + sourceRef: "transcript:partial-frontier", + contentOffset: 15_800, + reviewedChars: 200, + }), + ], }); expect( db @@ -1257,7 +1308,7 @@ describe("Dreaming", () => { "SELECT pass_id AS passId, delivered_offset AS offset FROM dreaming_evidence_consumption WHERE source_id = ?", ) .get("partial-frontier"), - ).toEqual({ passId: second.passId, offset: 32_000 }); + ).toEqual({ passId: second.passId, offset: 31_800 }); expect(second.passId).not.toBe(first.passId); expect(await shouldTriggerDreaming(accessor, cfg, AGENT, now + 21_000)).toBe(true); }, 15_000); @@ -1381,7 +1432,7 @@ describe("Dreaming", () => { const first = pendingDreamingEvidenceContinuations(tracedDb, AGENT, 20); expect(continuationQueries).toHaveLength(1); expect(continuationQueries[0]).toContain("LIMIT ?"); - expect(continuationArgs).toEqual([[AGENT, null, null, 20]]); + expect(continuationArgs).toEqual([[AGENT, 3, null, null, 20]]); expect(first.map((source) => source.id)).toEqual( Array.from({ length: 20 }, (_, index) => `prior-partial-${index.toString().padStart(2, "0")}`), ); @@ -1477,7 +1528,7 @@ describe("Dreaming", () => { { async run(input) { for (let page = 0; page < 5; page += 1) { - const result = await invokeDreamingTool(input, "search_evidence", { agentId: AGENT, limit: 50 }); + const result = await readAndReviewEvidence(input, { agentId: AGENT, limit: 50 }); if (result.hasMore !== true) break; } return { summary: "Read every session" }; @@ -2795,8 +2846,8 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what accessor, { async run(input) { - pages.push(await invokeDreamingTool(input, "search_evidence", { agentId: AGENT })); - pages.push(await invokeDreamingTool(input, "search_evidence", { agentId: AGENT })); + pages.push(await readAndReviewEvidence(input, { agentId: AGENT })); + pages.push(await readAndReviewEvidence(input, { agentId: AGENT })); return { summary: "Read two large pages" }; }, }, @@ -2828,7 +2879,7 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what { async run(input) { for (let call = 0; call < 5; call += 1) { - pages.push(await invokeDreamingTool(input, "search_evidence", { agentId: AGENT })); + pages.push(await readAndReviewEvidence(input, { agentId: AGENT })); } return { summary: "Drained the delivery queue" }; }, @@ -2859,6 +2910,792 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what expect(consumed.every((row) => row.offset === row.length)).toBe(true); }); + it("advances the evidence cursor only through reviewed text (#2094)", async () => { + for (const id of ["reviewed-one", "unreviewed-two", "unreviewed-three"]) { + seedTranscript(db, id, `${id} records that the release owner is settled.`); + } + const runPass = (review: (input: DreamingAgentInput, page: Record) => Promise) => + runDreamingAgentPass( + accessor, + { + async run(input) { + const page = await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); + await review(input, page); + await invokeDreamingTool(input, "runbook_write", { + summary: "## No-op\n- Stopped after one source.", + }); + return { summary: "Stopped after one source" }; + }, + }, + defaultCfg(), + "/tmp", + AGENT, + [AGENT], + "incremental", + ); + const first = await runPass(async (input, page) => { + const items = (page.items as Array>).filter( + (item) => item.sourceRef === "transcript:reviewed-one", + ); + expect(page.items).toHaveLength(3); + await reviewDreamingEvidence(input, { items }); + }); + expect(first).toMatchObject({ applied: 0, failed: 0 }); + const cursors = db + .prepare( + `SELECT source_id AS id, delivered_offset AS offset, source_length AS length, cursor_basis AS basis + FROM dreaming_evidence_consumption ORDER BY source_id`, + ) + .all() as Array<{ id: string; offset: number; length: number; basis: string }>; + expect(cursors.map((row) => [row.id, row.offset === row.length, row.offset, row.basis])).toEqual([ + ["reviewed-one", true, cursors[0]?.length, "review"], + ["unreviewed-three", false, 0, "review"], + ["unreviewed-two", false, 0, "review"], + ]); + expect(await getDreamingEpisodicTokenBacklog(accessor, AGENT)).toBeGreaterThan(0); + + let redelivered: string[] = []; + await runPass(async (_input, page) => { + redelivered = (page.items as Array<{ sourceRef: string }>).map((item) => item.sourceRef).sort(); + }); + expect(redelivered).toEqual(["transcript:unreviewed-three", "transcript:unreviewed-two"]); + }); + + it("advances a quote-anchored acknowledgement to the end of the quote and resumes with overlap (#2094)", async () => { + const sentences = Array.from({ length: 120 }, (_, index) => `Sentence ${index} names a settled fact.`).join(" "); + seedTranscript(db, "anchored", sentences); + const quote = "Sentence 60 names a settled fact."; + let reviewedEnd = 0; + await runDreamingAgentPass( + accessor, + { + async run(input) { + const page = await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); + const item = (page.items as Array<{ sourceRef: string; content: string; contentOffset: number }>)[0]; + if (!item) throw new Error("Missing anchored excerpt"); + reviewedEnd = item.contentOffset + item.content.indexOf(quote) + quote.length; + const review = await invokeDreamingTool(input, "review_evidence", { + agentId: AGENT, + items: [{ sourceRef: item.sourceRef, contentOffset: item.contentOffset, through: quote }], + }); + expect(review).toMatchObject({ ok: true, items: [{ reviewedThrough: reviewedEnd }] }); + return { summary: "Reviewed through sentence 60" }; + }, + }, + defaultCfg(), + "/tmp", + AGENT, + [AGENT], + "incremental", + ); + const cursor = () => + db + .prepare( + "SELECT delivered_offset AS offset, stalled_passes AS stalled FROM dreaming_evidence_consumption WHERE source_id = ?", + ) + .get("anchored") as { offset: number; stalled: number }; + expect(cursor()).toEqual({ offset: reviewedEnd, stalled: 0 }); + + let resumed: Record = {}; + await runDreamingAgentPass( + accessor, + { + async run(input) { + const page = await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); + resumed = (page.items as Array>)[0] ?? {}; + return { summary: "Read the resumed excerpt without reviewing it" }; + }, + }, + defaultCfg(), + "/tmp", + AGENT, + [AGENT], + "incremental", + ); + const offset = resumed.contentOffset as number; + expect(offset).toBeLessThanOrEqual(reviewedEnd - 200); + expect(offset).toBeGreaterThan(reviewedEnd - 400); + expect(String(resumed.content).startsWith("Sentence ")).toBe(true); + expect(resumed.reviewedChars).toBe(reviewedEnd - offset); + expect(cursor()).toEqual({ offset: reviewedEnd, stalled: 1 }); + }); + + it("rejects acknowledgements of text this pass and scope did not deliver (#2094)", async () => { + const other = "review-other-scope"; + accessor.withWriteTx((tx) => { + tx.prepare("INSERT OR IGNORE INTO agents (id, name, read_policy) VALUES (?, ?, 'isolated')").run(other, other); + }); + seedTranscript(db, "own-source", "Iris owns the deploy checklist."); + seedTranscript(db, "other-source", "Jove owns the incident rota.", undefined, other); + const rejections: Array> = []; + await runDreamingAgentPass( + accessor, + { + async run(input) { + const own = await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); + await invokeDreamingTool(input, "search_evidence", { agentId: other }); + const item = (own.items as Array<{ sourceRef: string; contentOffset: number }>)[0]; + if (!item) throw new Error("Missing own excerpt"); + const valid = { sourceRef: item.sourceRef, contentOffset: item.contentOffset }; + for (const invalid of [ + { sourceRef: item.sourceRef, contentOffset: item.contentOffset + 5 }, + { sourceRef: "transcript:other-source", contentOffset: 0 }, + { ...valid, through: "Iris owns the release calendar." }, + ]) { + rejections.push( + await invokeDreamingTool(input, "review_evidence", { agentId: AGENT, items: [valid, invalid] }), + ); + } + return { summary: "Every acknowledgement was rejected" }; + }, + }, + defaultCfg(), + "/tmp", + AGENT, + [AGENT, other], + "incremental", + ); + expect(rejections.map((result) => [result.ok, result.code])).toEqual([ + [false, "EXCERPT_NOT_DELIVERED"], + [false, "SCOPE_MISMATCH"], + [false, "QUOTE_NOT_IN_EXCERPT"], + ]); + expect(rejections[0]?.items).toEqual([expect.objectContaining({ index: 1, code: "EXCERPT_NOT_DELIVERED" })]); + expect( + db.prepare("SELECT COUNT(*) AS n FROM dreaming_evidence_consumption WHERE delivered_offset > 0").get(), + ).toEqual({ n: 0 }); + }); + + it("rejects acknowledgements of an excerpt whose source changed after delivery (#2094)", async () => { + seedTranscript(db, "edited-source", "Iris owns the deploy checklist."); + let rejected: Record = {}; + await runDreamingAgentPass( + accessor, + { + async run(input) { + const page = await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); + db.prepare("UPDATE session_transcripts SET content = ? WHERE session_key = ?").run( + "Iris owns the release calendar.", + "edited-source", + ); + rejected = await invokeDreamingTool(input, "review_evidence", { + agentId: AGENT, + items: (page.items as Array<{ sourceRef: string; contentOffset: number }>).map((item) => ({ + sourceRef: item.sourceRef, + contentOffset: item.contentOffset, + })), + }); + return { summary: "The source changed before review" }; + }, + }, + defaultCfg(), + "/tmp", + AGENT, + [AGENT], + "incremental", + ); + expect(rejected).toMatchObject({ ok: false, code: "SOURCE_CHANGED" }); + expect( + db.prepare("SELECT COUNT(*) AS n FROM dreaming_evidence_consumption WHERE delivered_offset > 0").get(), + ).toEqual({ n: 0 }); + }); + + it("floors the evidence cursor at a successfully cited quote (#2094)", async () => { + seedTranscript( + db, + "cited-floor", + "Aster is the durable release project. Later chatter covers the weekend weather.", + ); + const quote = "Aster is the durable release project."; + let floor = 0; + await runDreamingAgentPass( + accessor, + { + async run(input) { + const page = await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); + const item = (page.items as Array<{ content: string; contentOffset: number }>)[0]; + if (!item) throw new Error("Missing cited excerpt"); + floor = item.contentOffset + item.content.indexOf(quote) + quote.length; + const filed = await invokeDreamingTool(input, "apply_ontology_ops", { + agentId: AGENT, + operations: [ + { + operation: "create_entity", + payload: { name: "Aster", type: "project" }, + reason: "The evidence names a durable project.", + evidence: [{ source_ref: "transcript:cited-floor", quote }], + }, + ], + }); + expect(filed.ok).toBe(true); + return { summary: "Filed without acknowledging" }; + }, + }, + defaultCfg(), + "/tmp", + AGENT, + [AGENT], + "incremental", + ); + const row = db + .prepare( + "SELECT delivered_offset AS offset, source_length AS length FROM dreaming_evidence_consumption WHERE source_id = ?", + ) + .get("cited-floor") as { offset: number; length: number }; + expect(row.offset).toBe(floor); + expect(row.offset).toBeLessThan(row.length); + }); + + it("surfaces a stalled source and stops it leading the continuation queue (#2094)", async () => { + const capturedAt = "2026-08-11T00:00:00.000Z"; + accessor.withWriteTx((tx) => { + const insertPass = tx.prepare( + "INSERT INTO dreaming_passes (id, agent_id, mode, status) VALUES (?, ?, 'incremental', 'completed')", + ); + const insertConsumption = tx.prepare( + `INSERT INTO dreaming_evidence_consumption + (agent_id, source_kind, source_id, source_captured_at, source_entry_id, source_revision, + delivered_offset, source_length, pass_id, updated_at) + VALUES (?, 'transcript', ?, ?, '', ?, 2000, 10000, ?, ?)`, + ); + for (const id of ["stuck-first", "waiting-second"]) { + insertPass.run(`${id}-pass`, AGENT); + seedTranscript(db, id, "x".repeat(10_000), capturedAt); + insertConsumption.run(AGENT, id, capturedAt, capturedAt, `${id}-pass`, capturedAt); + } + }); + const readFirst = async (review = false): Promise => { + let ref = ""; + await runDreamingAgentPass( + accessor, + { + async run(input) { + const page = await invokeDreamingTool(input, "search_evidence", { agentId: AGENT, limit: 1 }); + ref = String((page.items as Array<{ sourceRef: string }>)[0]?.sourceRef); + if (review) { + const stuck = await invokeDreamingTool(input, "search_evidence", { agentId: AGENT, limit: 1 }); + await reviewDreamingEvidence(input, stuck); + } + return { summary: "Read one queued excerpt" }; + }, + }, + defaultCfg(), + "/tmp", + AGENT, + [AGENT], + "incremental", + ); + return ref; + }; + const stalledAttention = () => + db + .prepare( + `SELECT details_json AS details, resolved_at AS resolvedAt FROM dreaming_attention + WHERE agent_id = ? AND kind = 'evidence_requeue' AND subject_ref = ?`, + ) + .get(AGENT, "transcript:stuck-first") as { details: string; resolvedAt: string | null } | null; + + for (let pass = 0; pass < 2; pass += 1) expect(await readFirst()).toBe("transcript:stuck-first"); + expect(stalledAttention()).toBeNull(); + expect(await readFirst()).toBe("transcript:stuck-first"); + expect(stalledAttention()).toMatchObject({ resolvedAt: null }); + expect(JSON.parse(stalledAttention()?.details ?? "{}")).toMatchObject({ + reason: "evidence-stalled", + reviewedChars: "2000", + stalledPasses: "3", + }); + expect(await readFirst()).toBe("transcript:waiting-second"); + + expect(await readFirst(true)).toBe("transcript:waiting-second"); + expect(stalledAttention()?.resolvedAt).not.toBeNull(); + expect( + db + .prepare("SELECT stalled_passes AS stalled FROM dreaming_evidence_consumption WHERE source_id = ?") + .get("stuck-first"), + ).toEqual({ stalled: 0 }); + }, 15_000); + + it("does not schedule passes for stall attention and resolves it on reviewed exclusion (#2094)", async () => { + seedTranscript(db, "stalled-noise", "Weekend small talk with no durable fact."); + const runPass = (run: (input: DreamingAgentInput) => Promise) => + runDreamingAgentPass( + accessor, + { + async run(input) { + await run(input); + return { summary: "Read the queue" }; + }, + }, + defaultCfg(), + "/tmp", + AGENT, + [AGENT], + "incremental", + ); + for (let pass = 0; pass < DREAMING_EVIDENCE_STALL_PASSES; pass += 1) { + await runPass(async (input) => { + await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); + }); + } + const stalled = () => + db + .prepare( + "SELECT resolved_at AS resolvedAt FROM dreaming_attention WHERE kind = 'evidence_requeue' AND subject_ref = ?", + ) + .get("transcript:stalled-noise") as { resolvedAt: string | null } | null; + expect(stalled()).toEqual({ resolvedAt: null }); + const cfg = defaultCfg({ tokenThreshold: 100_000, backfillOnFirstRun: false }); + expect( + await evaluateDreamingTrigger(accessor, cfg, AGENT, { + kind: "indeterminate", + tokenLowerBound: 0, + hasBacklog: true, + sourcesScanned: 1, + }), + ).toEqual({ trigger: false }); + + await runPass(async (input) => { + await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); + await invokeDreamingTool(input, "runbook_write", { + summary: "## No-op\n- transcript:stalled-noise is small talk.", + reviewedExcludedEvidence: [ + { agentId: AGENT, sourceRef: "transcript:stalled-noise", reason: "Small talk with no durable fact." }, + ], + }); + }); + expect(stalled()?.resolvedAt).toEqual(expect.any(String)); + }); + + it("does not count deferred or withheld passes toward an evidence stall (#2094)", async () => { + const capturedAt = "2026-08-11T00:00:00.000Z"; + seedTranscript(db, "suppressed-stall", "x".repeat(10_000), capturedAt); + accessor.withWriteTx((tx) => { + tx.prepare( + "INSERT INTO dreaming_passes (id, agent_id, mode, status) VALUES ('suppressed-pass', ?, 'incremental', 'completed')", + ).run(AGENT); + tx.prepare( + `INSERT INTO dreaming_evidence_consumption + (agent_id, source_kind, source_id, source_captured_at, source_entry_id, source_revision, + delivered_offset, source_length, pass_id, updated_at, cursor_basis, stalled_passes) + VALUES (?, 'transcript', 'suppressed-stall', ?, '', ?, 2000, 10000, 'suppressed-pass', ?, 'review', ?)`, + ).run(AGENT, capturedAt, capturedAt, capturedAt, DREAMING_EVIDENCE_STALL_PASSES - 1); + }); + const runPass = (run: (input: DreamingAgentInput) => Promise) => + runDreamingAgentPass( + accessor, + { + async run(input) { + const page = await invokeDreamingTool(input, "search_evidence", { agentId: AGENT, limit: 1 }); + expect((page.items as Array<{ sourceRef: string }>)[0]?.sourceRef).toBe("transcript:suppressed-stall"); + await run(input); + return { summary: "Read the stalled continuation" }; + }, + }, + defaultCfg(), + "/tmp", + AGENT, + [AGENT], + "incremental", + ); + const cursor = () => + db + .prepare( + "SELECT delivered_offset AS offset, stalled_passes AS stalled FROM dreaming_evidence_consumption WHERE source_id = ?", + ) + .get("suppressed-stall"); + const stalledAttention = () => + db + .prepare("SELECT COUNT(*) AS n FROM dreaming_attention WHERE kind = 'evidence_requeue' AND subject_ref = ?") + .get("transcript:suppressed-stall"); + const unchanged = { offset: 2000, stalled: DREAMING_EVIDENCE_STALL_PASSES - 1 }; + + await runPass(async (input) => { + await invokeDreamingTool(input, "runbook_write", { + summary: "## Deferred\n- transcript:suppressed-stall awaits a decision.", + deferredEvidence: [{ agentId: AGENT, sourceRef: "transcript:suppressed-stall" }], + }); + }); + expect(cursor()).toEqual(unchanged); + + await runPass(async (input) => { + const apply = await invokeDreamingTool(input, "apply_ontology_ops", { + agentId: AGENT, + operations: "invalid array", + }); + expect(apply.ok).toBe(false); + }); + expect(cursor()).toEqual(unchanged); + expect(stalledAttention()).toEqual({ n: 0 }); + + await runPass(async () => {}); + expect(cursor()).toEqual({ offset: 2000, stalled: DREAMING_EVIDENCE_STALL_PASSES }); + expect(stalledAttention()).toEqual({ n: 1 }); + }); + + it("serves fresh sources before a stalled continuation (#2094)", async () => { + const capturedAt = "2026-08-11T00:00:00.000Z"; + seedTranscript(db, "stalled-continuation", "x".repeat(10_000), capturedAt); + seedTranscript(db, "fresh-source", "Fresh evidence names a settled owner.", "2026-08-10T00:00:00.000Z"); + accessor.withWriteTx((tx) => { + tx.prepare( + "INSERT INTO dreaming_passes (id, agent_id, mode, status) VALUES ('stalled-pass', ?, 'incremental', 'completed')", + ).run(AGENT); + tx.prepare( + `INSERT INTO dreaming_evidence_consumption + (agent_id, source_kind, source_id, source_captured_at, source_entry_id, source_revision, + delivered_offset, source_length, pass_id, updated_at, cursor_basis, stalled_passes) + VALUES (?, 'transcript', 'stalled-continuation', ?, '', ?, 2000, 10000, 'stalled-pass', ?, 'review', 3)`, + ).run(AGENT, capturedAt, capturedAt, capturedAt); + }); + const pages: Array> = []; + await runDreamingAgentPass( + accessor, + { + async run(input) { + pages.push(await invokeDreamingTool(input, "search_evidence", { agentId: AGENT, limit: 1 })); + pages.push(await invokeDreamingTool(input, "search_evidence", { agentId: AGENT, limit: 1 })); + return { summary: "Read two queue pages" }; + }, + }, + defaultCfg(), + "/tmp", + AGENT, + [AGENT], + "incremental", + ); + const [first, second] = pages.map((page) => (page.items as Array>)[0]); + expect(first?.sourceRef).toBe("transcript:fresh-source"); + expect(pages[0]?.hasMore).toBe(true); + expect(second?.sourceRef).toBe("transcript:stalled-continuation"); + expect(second?.reviewedChars).toBe(2000 - Number(second?.contentOffset)); + }); + + it("resolves imported-source attention only after every member transcript is reviewed (#2094)", async () => { + const importId = "import:batch-2094"; + const insertTranscript = db.prepare( + `INSERT INTO session_transcripts + (session_key, content, agent_id, created_at, updated_at, completed_at, source_id, content_hash) + VALUES (?, ?, ?, ?, ?, ?, ?, ?)`, + ); + const capturedAt = "2026-10-08 00:00:00"; + insertTranscript.run( + "imported-a", + "Kira owns the import pipeline.", + AGENT, + capturedAt, + capturedAt, + capturedAt, + importId, + "hash-a", + ); + insertTranscript.run( + "imported-b", + "Weekend small talk.", + AGENT, + capturedAt, + capturedAt, + capturedAt, + importId, + "hash-b", + ); + accessor.withWriteTx((tx) => { + enqueueDreamingAttentionInTx(tx, { + agentId: AGENT, + kind: "evidence_requeue", + subjectRef: `source:${importId}`, + details: { sourceId: importId, reason: "transcript-import-committed" }, + priority: 50, + }); + }); + const importAttention = () => + db + .prepare( + "SELECT resolved_at AS resolvedAt, resolved_by_pass_id AS passId FROM dreaming_attention WHERE subject_ref = ?", + ) + .get(`source:${importId}`) as { resolvedAt: string | null; passId: string | null }; + + let prompt = ""; + let listed: Record = {}; + await runDreamingAgentPass( + accessor, + { + async run(input) { + prompt = input.prompt; + listed = await invokeDreamingTool(input, "attention_list", { agentId: AGENT, kind: "evidence_requeue" }); + const page = await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); + await reviewDreamingEvidence(input, { + items: (page.items as Array>).filter( + (item) => item.sourceRef === "transcript:imported-a", + ), + }); + return { summary: "Reviewed one imported conversation" }; + }, + }, + defaultCfg(), + "/tmp", + AGENT, + [AGENT], + "incremental", + ); + expect(prompt).not.toContain(`source:${importId}`); + expect(listed).toMatchObject({ ok: true, items: [] }); + expect(importAttention().resolvedAt).toBeNull(); + + const second = await runDreamingAgentPass( + accessor, + { + async run(input) { + await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); + await invokeDreamingTool(input, "runbook_write", { + summary: "## No-op\n- transcript:imported-b is small talk.", + reviewedExcludedEvidence: [ + { agentId: AGENT, sourceRef: "transcript:imported-b", reason: "Small talk with no durable fact." }, + ], + }); + return { summary: "Excluded the other imported conversation" }; + }, + }, + defaultCfg(), + "/tmp", + AGENT, + [AGENT], + "incremental", + ); + expect(importAttention()).toEqual({ resolvedAt: expect.any(String), passId: second.passId }); + }); + + it("schedules a pass for an import nudge only until a pass starts after it (#2094)", async () => { + const importId = "import:trigger-2094"; + const capturedAt = "2026-10-08 00:00:00"; + db.prepare( + `INSERT INTO session_transcripts + (session_key, content, agent_id, created_at, updated_at, completed_at, source_id, content_hash) + VALUES (?, ?, ?, ?, ?, ?, ?, ?)`, + ).run( + "imported-pending", + "Kira owns the import pipeline.", + AGENT, + capturedAt, + capturedAt, + capturedAt, + importId, + "hash", + ); + const raise = async (): Promise => { + await new Promise((resolve) => setTimeout(resolve, 5)); + const nudge = importedSourceAttentionUpsert(AGENT, importId); + db.prepare(nudge.sql).run(...nudge.params); + await new Promise((resolve) => setTimeout(resolve, 5)); + }; + const cfg = defaultCfg({ tokenThreshold: 100_000, backfillOnFirstRun: false }); + const evaluate = () => + evaluateDreamingTrigger(accessor, cfg, AGENT, { + kind: "indeterminate", + tokenLowerBound: 0, + hasBacklog: false, + sourcesScanned: 1, + }); + const pending = () => + db + .prepare("SELECT resolved_at AS resolvedAt FROM dreaming_attention WHERE subject_ref = ?") + .get(`source:${importId}`); + + await raise(); + expect(await evaluate()).toEqual({ trigger: true, reason: "attention" }); + + await runDreamingAgentPass( + accessor, + { + async run(input) { + await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); + return { summary: "Read the import without reviewing it" }; + }, + }, + defaultCfg(), + "/tmp", + AGENT, + [AGENT], + "incremental", + ); + expect(pending()).toEqual({ resolvedAt: null }); + expect(await evaluate()).toEqual({ trigger: false }); + + await raise(); + expect(await evaluate()).toEqual({ trigger: true, reason: "attention" }); + }); + + it("keeps fully reviewed imported transcripts out of the delivery queue (#2094)", async () => { + const insertImported = db.prepare( + `INSERT INTO session_transcripts + (session_key, content, agent_id, created_at, updated_at, completed_at, source_id, content_hash) + VALUES (?, ?, ?, ?, ?, ?, ?, ?)`, + ); + const seedImported = (sessionKey: string, capturedAt: string): void => { + insertImported.run( + sessionKey, + `${sessionKey} owns one durable fact.`, + AGENT, + capturedAt, + capturedAt, + capturedAt, + "import:queue-2094", + `hash-${sessionKey}`, + ); + }; + for (let index = 0; index < 52; index += 1) + seedImported(`reviewed-${String(index).padStart(2, "0")}`, `2026-10-08 01:${String(index).padStart(2, "0")}:00`); + await runDreamingAgentPass( + accessor, + { + async run(input) { + for (let page = 0; page < 10; page += 1) { + const result = await readAndReviewEvidence(input, { agentId: AGENT }); + if (result.hasMore !== true) break; + } + return { summary: "Reviewed every imported conversation" }; + }, + }, + defaultCfg(), + "/tmp", + AGENT, + [AGENT], + "incremental", + ); + expect( + db + .prepare( + "SELECT COUNT(*) AS count FROM dreaming_evidence_consumption WHERE delivered_offset >= source_length AND source_entry_id = ?", + ) + .get("import:queue-2094"), + ).toEqual({ count: 52 }); + + seedImported("older-unreviewed", "2026-10-08 00:00:00"); + let page: Record = {}; + await runDreamingAgentPass( + accessor, + { + async run(input) { + page = await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); + return { summary: "Read the next page" }; + }, + }, + defaultCfg(), + "/tmp", + AGENT, + [AGENT], + "incremental", + ); + expect((page.items as Array>).map((item) => item.sourceRef)).toEqual([ + "transcript:older-unreviewed", + ]); + expect(page.hasMore).toBe(false); + }); + + describe("imported-source attention drain checks (#2094)", () => { + const capturedAt = "2026-10-08 00:00:00"; + const seedImported = (sourceId: string, sessionKey: string, content: string): void => { + db.prepare( + `INSERT INTO session_transcripts + (session_key, content, agent_id, created_at, updated_at, completed_at, source_id, content_hash) + VALUES (?, ?, ?, ?, ?, ?, ?, ?)`, + ).run(sessionKey, content, AGENT, capturedAt, capturedAt, capturedAt, sourceId, `hash-${sessionKey}`); + }; + const recordConsumption = (sourceId: string, sessionKey: string, offset: number, length: number): void => { + db.prepare( + `INSERT INTO dreaming_evidence_consumption + (agent_id, source_kind, source_id, source_captured_at, source_entry_id, source_revision, + delivered_offset, source_length, pass_id, updated_at) + VALUES (?, 'transcript', ?, ?, ?, ?, ?, ?, 'seed-pass', ?)`, + ).run(AGENT, sessionKey, capturedAt, sourceId, `hash-${sessionKey}`, offset, length, capturedAt); + }; + const raise = (sourceId: string, createdAt: string): void => { + db.prepare( + `INSERT INTO dreaming_attention (id, agent_id, kind, subject_ref, details_json, priority, created_at) + VALUES (?, ?, 'evidence_requeue', ?, ?, 50, ?)`, + ).run( + `attention-${sourceId}`, + AGENT, + `source:${sourceId}`, + JSON.stringify({ sourceId, reason: "transcript-import-committed" }), + createdAt, + ); + }; + const resolvedAt = (sourceId: string): string | null => + ( + db + .prepare("SELECT resolved_at AS resolvedAt FROM dreaming_attention WHERE subject_ref = ?") + .get(`source:${sourceId}`) as { resolvedAt: string | null } + ).resolvedAt; + const finalize = (passId: string): number => + accessor.withWriteTx((tx) => resolveImportedSourceAttentionInTx(tx, passId, [AGENT])); + + it("stops at the first unconsumed member instead of rendering the whole source", () => { + for (let index = 0; index < 40; index += 1) { + seedImported("import:wide", `wide-${String(index).padStart(2, "0")}`, "Kira owns the import pipeline."); + } + const read = db as unknown as ReadDb; + expect(probeSourceEvidenceDrain(read, AGENT, "import:wide", { maxRenders: Number.POSITIVE_INFINITY })).toEqual({ + status: "pending", + resumeAfter: null, + rendered: 1, + }); + + recordConsumption("import:wide", "wide-39", 5, 30); + expect(probeSourceEvidenceDrain(read, AGENT, "import:wide", { maxRenders: Number.POSITIVE_INFINITY })).toEqual({ + status: "pending", + resumeAfter: null, + rendered: 0, + }); + }); + + it("resolves drained rows queued behind sources that still have evidence", () => { + const busy = Array.from( + { length: IMPORTED_SOURCE_ATTENTION_ROWS_PER_SCOPE }, + (_, index) => `import:busy-${index}`, + ); + busy.forEach((sourceId, index) => { + seedImported(sourceId, `${sourceId}-a`, "Kira owns the import pipeline."); + recordConsumption(sourceId, `${sourceId}-a`, 5, 30); + raise(sourceId, `2026-10-08 00:00:${String(index).padStart(2, "0")}`); + }); + seedImported("import:drained", "drained-a", "Kira owns the import pipeline."); + recordConsumption("import:drained", "drained-a", 30, 30); + raise("import:drained", "2026-10-08 00:01:00"); + + expect(finalize("pass-1")).toBe(0); + expect(resolvedAt("import:drained")).toBeNull(); + expect(finalize("pass-2")).toBe(1); + expect(resolvedAt("import:drained")).not.toBeNull(); + expect(busy.every((sourceId) => resolvedAt(sourceId) === null)).toBe(true); + }); + + it("bounds renders per finalization and resumes where the last one stopped", () => { + const members = IMPORTED_SOURCE_ATTENTION_RENDER_BUDGET * 3 + 2; + for (let index = 0; index < members; index += 1) { + seedImported("import:blank", `blank-${String(index).padStart(2, "0")}`, " "); + } + raise("import:blank", capturedAt); + const resumeAfter = (): unknown => + JSON.parse( + ( + db + .prepare("SELECT details_json AS details FROM dreaming_attention WHERE subject_ref = ?") + .get("source:import:blank") as { details: string } + ).details, + ).drainResumeAfter; + + expect(finalize("pass-1")).toBe(0); + expect(resumeAfter()).toBe( + `transcript:blank-${String(IMPORTED_SOURCE_ATTENTION_RENDER_BUDGET - 1).padStart(2, "0")}`, + ); + expect(finalize("pass-2")).toBe(0); + expect(finalize("pass-3")).toBe(0); + expect(resolvedAt("import:blank")).toBeNull(); + expect(finalize("pass-4")).toBe(1); + expect(resolvedAt("import:blank")).not.toBeNull(); + }); + }); + it("closes new evidence delivery halfway through the pass timeout", async () => { seedTranscript(db, "late-source", "Ines maintains the on-call rota."); let page: Record = {}; @@ -2895,7 +3732,7 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what accessor, { async run(input) { - await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); + await readAndReviewEvidence(input, { agentId: AGENT }); await invokeDreamingTool(input, "apply_ontology_ops", { agentId: AGENT, operations: [ @@ -2918,7 +3755,9 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what expect(result.failed).toBeGreaterThan(0); const consumed = ( db - .prepare("SELECT source_id AS id FROM dreaming_evidence_consumption WHERE source_kind = 'transcript'") + .prepare( + "SELECT source_id AS id FROM dreaming_evidence_consumption WHERE source_kind = 'transcript' AND delivered_offset > 0", + ) .all() as Array<{ id: string; }> @@ -2932,7 +3771,7 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what accessor, { async run(input) { - await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); + await readAndReviewEvidence(input, { agentId: AGENT }); const apply = await invokeDreamingTool(input, "apply_ontology_ops", { agentId: AGENT, operations: "invalid array", @@ -2949,7 +3788,11 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what ); expect(result).toMatchObject({ applied: 0, failed: 1 }); expect( - db.prepare("SELECT COUNT(*) AS n FROM dreaming_evidence_consumption WHERE source_id = 'schema-rejected'").get(), + db + .prepare( + "SELECT COUNT(*) AS n FROM dreaming_evidence_consumption WHERE source_id = 'schema-rejected' AND delivered_offset > 0", + ) + .get(), ).toEqual({ n: 0 }); }); @@ -2964,8 +3807,8 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what accessor, { async run(input) { - await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); - await invokeDreamingTool(input, "search_evidence", { agentId: other }); + await readAndReviewEvidence(input, { agentId: AGENT }); + await readAndReviewEvidence(input, { agentId: other }); accessor.withWriteTx((tx) => { tx.prepare( `INSERT INTO dreaming_tool_calls @@ -2982,7 +3825,9 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what [AGENT, other], "incremental", ); - expect(db.prepare("SELECT COUNT(*) AS n FROM dreaming_evidence_consumption").get()).toEqual({ n: 0 }); + expect( + db.prepare("SELECT COUNT(*) AS n FROM dreaming_evidence_consumption WHERE delivered_offset > 0").get(), + ).toEqual({ n: 0 }); }); it("keeps progress for sources a pass filed when an uncited write fails in the same agent", async () => { @@ -2992,7 +3837,7 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what accessor, { async run(input) { - await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); + await readAndReviewEvidence(input, { agentId: AGENT }); const filed = await invokeDreamingTool(input, "apply_ontology_ops", { agentId: AGENT, operations: [ @@ -3026,7 +3871,9 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what expect(result.failed).toBeGreaterThan(0); const consumed = ( db - .prepare("SELECT source_id AS id FROM dreaming_evidence_consumption WHERE source_kind = 'transcript'") + .prepare( + "SELECT source_id AS id FROM dreaming_evidence_consumption WHERE source_kind = 'transcript' AND delivered_offset > 0", + ) .all() as Array<{ id: string; }> @@ -3047,7 +3894,7 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what accessor, { async run(input) { - await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); + await readAndReviewEvidence(input, { agentId: AGENT }); const rejected = await invokeDreamingTool(input, "apply_ontology_ops", { agentId: AGENT, operations: [ @@ -3078,7 +3925,9 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what expect(result.failed).toBeGreaterThan(0); const consumed = ( db - .prepare("SELECT source_id AS id FROM dreaming_evidence_consumption WHERE source_kind = 'transcript'") + .prepare( + "SELECT source_id AS id FROM dreaming_evidence_consumption WHERE source_kind = 'transcript' AND delivered_offset > 0", + ) .all() as Array<{ id: string; }> @@ -3135,8 +3984,8 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what accessor, { async run(input) { - await invokeDreamingTool(input, "search_evidence", { agentId: AGENT }); - await invokeDreamingTool(input, "search_evidence", { agentId: other }); + await readAndReviewEvidence(input, { agentId: AGENT }); + await readAndReviewEvidence(input, { agentId: other }); await invokeDreamingTool(input, "apply_ontology_ops", { agentId: other, operations: [{ operation: "not_an_ontology_operation", payload: {} }], @@ -3153,7 +4002,7 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what expect(result.failed).toBeGreaterThan(0); const consumed = db .prepare( - "SELECT agent_id AS agentId, source_id AS id FROM dreaming_evidence_consumption WHERE source_kind = 'transcript'", + "SELECT agent_id AS agentId, source_id AS id FROM dreaming_evidence_consumption WHERE source_kind = 'transcript' AND delivered_offset > 0", ) .all(); expect(consumed).toEqual([{ agentId: AGENT, id: "scope-a-source" }]); @@ -3169,7 +4018,7 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what const search = input.tools.find((tool) => tool.name === "search_evidence"); const apply = input.tools.find((tool) => tool.name === "apply_ontology_ops"); if (!search || !apply) throw new Error("Missing Dreaming delivery tools"); - await search.execute("call", { agentId: AGENT }, undefined, undefined, {} as never); + await readAndReviewEvidence(input, { agentId: AGENT }); await apply.execute( "call", { @@ -3197,7 +4046,9 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what expect( ( db - .prepare("SELECT COUNT(*) AS count FROM dreaming_evidence_consumption WHERE source_id = ?") + .prepare( + "SELECT COUNT(*) AS count FROM dreaming_evidence_consumption WHERE source_id = ? AND delivered_offset > 0", + ) .get("rejected-summary") as { count: number; } @@ -3459,7 +4310,7 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what async run(input) { const search = input.tools.find((tool) => tool.name === "search_evidence"); if (!search) throw new Error("Missing search_evidence"); - await search.execute("call", { agentId: AGENT }, undefined, undefined, {} as never); + await readAndReviewEvidence(input, { agentId: AGENT }); await commitDreamingTestHead(input, { sourceRef: "summary:starved-evidence", text: "New transcript evidence that content passes never reached.", @@ -3492,7 +4343,7 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what seedTranscript(db, "mid-pass-arrival", "Evidence arrived while this content pass was already running."); const search = input.tools.find((tool) => tool.name === "search_evidence"); if (!search) throw new Error("Missing search_evidence"); - await search.execute("call", { agentId: AGENT }, undefined, undefined, {} as never); + await readAndReviewEvidence(input, { agentId: AGENT }); await commitDreamingTestHead(input, { sourceRef: "transcript:mid-pass-arrival", text: "Evidence arrived while this content pass was already running.", @@ -3510,7 +4361,9 @@ It is now Monday, 2026-10-05 18:42 America/Denver (GMT-06:00). Use this for what expect( ( db - .prepare("SELECT COUNT(*) AS count FROM dreaming_evidence_consumption WHERE source_id = ?") + .prepare( + "SELECT COUNT(*) AS count FROM dreaming_evidence_consumption WHERE source_id = ? AND delivered_offset = source_length", + ) .get("mid-pass-arrival") as { count: number; } diff --git a/platform/daemon/src/pipeline/dreaming.ts b/platform/daemon/src/pipeline/dreaming.ts index 6a2d253e20..f81e5898de 100644 --- a/platform/daemon/src/pipeline/dreaming.ts +++ b/platform/daemon/src/pipeline/dreaming.ts @@ -66,6 +66,9 @@ import { evidenceContentSha256, failedOperationEvidence, recordDreamingEvidenceConsumptionInTx, + resolveImportedSourceAttentionInTx, + SEEN_IMPORTED_SOURCE_ATTENTION_SQL, + STALLED_EVIDENCE_ATTENTION_SQL, } from "./dreaming-evidence-consumption"; import { parseDreamingReviewedExcludedEvidence, @@ -712,7 +715,7 @@ export async function getActiveDreamingPasses( } const DREAMING_CODEMODE_PROMPT = - "Lookups (search_entities, get_entity, list_aspect_claims, validate_proposal, attention_list, zoom_history) are available only inside the codemode tool. Batch the lookups a page needs into one script: call them through tools.(args), parse each JSON result, and print only what you need, carrying ids from results instead of retyping them. Reading evidence (search_evidence) and every write (apply_ontology_ops, runbook_write, memory_head_commit) stay direct tool calls; a script cannot call them."; + "Lookups (search_entities, get_entity, list_aspect_claims, validate_proposal, attention_list, zoom_history) are available only inside the codemode tool. Batch the lookups a page needs into one script: call them through tools.(args), parse each JSON result, and print only what you need, carrying ids from results instead of retyping them. Reading and acknowledging evidence (search_evidence, review_evidence) and every write (apply_ontology_ops, runbook_write, memory_head_commit) stay direct tool calls; a script cannot call them."; const MAX_DREAMING_TOOL_TRACE_JSON_CHARS = 128_000; @@ -957,7 +960,7 @@ An install may have several agent scopes (listed in when there is - If you inspect a flagged target and judge it should stay as it is (a deliberate keep — e.g. a live entity with a non-concrete type, or an over-cap aspect you chose not to consolidate), close the record with decline_attention citing its attention id. Declining is an affirmative judgment: only decline records you actually inspected, and never decline records you could not complete this pass — defer those with a named blocker instead. 3. Take the surprisal hints listed in . These are bounded exploration hints, not evidence and not hygiene provenance. For each hint, inspect its memory: subjectRef with search_evidence in the owning scope. Treat the score only as a priority signal: if the source establishes a useful, settled fact, use a normal content operation with an exact quote; if it is valid but not useful or is noise, use decline_attention after inspecting it. Never create a claim or entity from the score alone, and never cite attention: for a content operation. A surprisal hint must not bypass the evidence cursor or audited apply path. 4. Take the review_due records listed in . For expired records, inspect the cited memory with search_evidence using its subjectRef, then supersede the matching active claim with supersede_claim_value. Use the supplied entityId, aspectId, attributeId, and claimKey when present. The replacement must state that the planned event remains unconfirmed; never rewrite it as if the event happened. Cite an exact quote from the original memory. Do not supersede approaching records. When creating or setting a future temporal claim, set payload.reviewAfter to the referenced ISO timestamp. Then take the contested_claim records with reason source_changed: the claim's source was edited after the claim was filed. Read the source's current text with search_evidence using details.sourceRef. If it still states the claim, close the record with decline_attention. If it states a different value, supersede the claim with supersede_claim_value citing an exact quote from the current text, then close the record with decline_attention. If it no longer states the claim, archive it with archive_claim_value and provenance attention:. For reason source_removed the source no longer exists: archive the claim the same way unless search_evidence finds other evidence that still states it, in which case close the record with decline_attention. -5. Only when the hygiene queue is clear: find new evidence since the cutoff. Read unprocessed evidence one page at a time with search_evidence — omit the query, since, and before to take the next page of the durable delivery queue. A page holds one excerpt per source; a source with contentHasNext continues by itself on a later page, so do not page through queued sources with sourceRef and offset. File what each page establishes (or mark a source you have read in full as reviewed or deferred) before asking for the next page, so the work is kept if the pass runs out of time; stop when hasMore is false or the queue reports deliveryClosed. Use a query or sourceRef only to look up specific history or verify a citation. Prefer evidence from completed transcript sessions; historical summary rows are not part of the default delivery path. A transcript with completed: false is mid-stream — defer filing from it with the named blocker "transcript still mid-stream" (re-check completed each pass: a session still active when re-checked is a re-verified blocker, not a repeated one), and note the deferral in the pass log, because its states may be contradicted by the session's end. For each new source: +5. Only when the hygiene queue is clear: find new evidence since the cutoff. Read unprocessed evidence one page at a time with search_evidence — omit the query, since, and before to take the next page of the durable delivery queue. A page holds one excerpt per source; a source with contentHasNext continues by itself on a later page, so do not page through queued sources with sourceRef and offset. File what each page establishes (or mark a source you have read in full as reviewed or deferred), then call review_evidence for that page, before asking for the next page, so the work is kept if the pass runs out of time. Acknowledge each excerpt with the sourceRef and contentOffset from the page, and omit through when you reviewed all of it; set through to an exact quote only when you stopped partway, so the rest is delivered again. Text you never acknowledge or cite stays queued for a later pass. When an excerpt reports reviewedChars, its first reviewedChars characters were reviewed by an earlier pass: do not file them again; stop when hasMore is false or the queue reports deliveryClosed. Use a query or sourceRef only to look up specific history or verify a citation. Prefer evidence from completed transcript sessions; historical summary rows are not part of the default delivery path. A transcript with completed: false is mid-stream — defer filing from it with the named blocker "transcript still mid-stream" (re-check completed each pass: a session still active when re-checked is a re-verified blocker, not a repeated one), and note the deferral in the pass log, because its states may be contradicted by the session's end. For each new source: - search_entities for subjects it establishes. - File claims only for what the source establishes as settled fact: outcomes, decisions, shipped changes, stable behavior. Do not file instructions that were merely suggested, hypotheses or diagnoses, open questions, or intermediate states of an ongoing investigation. When a source shows an attempt and its outcome, file the outcome. - A concrete deliverable the assistant produced for the user's own project, plan, or situation (a budget, schedule, draft, plan, or specific recommendation the user asked for) is a durable fact about that project: file its specifics as a claim on the project, worded as proposed or drafted rather than confirmed. Generic information not tied to the user's own circumstances is not. Work the user asks for on behalf of their own job, business, employer, trip, or event is their own situation. @@ -1077,7 +1080,7 @@ An install may have several agent scopes (listed in when there is 1. Read the pass history below. Establish cutoff: sources viewed, changes applied, deferred items. Zoom (zoom_history) into any line that mentions work you are about to repeat, resume, or re-defer before acting on it. 2. Take the review_due records listed in . For expired records, inspect the cited memory with search_evidence using its subjectRef, then supersede the matching active claim with supersede_claim_value. Use the supplied entityId, aspectId, attributeId, and claimKey when present. The replacement must state that the planned event remains unconfirmed; never rewrite it as if the event happened. Cite an exact quote from the original memory. Do not supersede approaching records. When creating or setting a future temporal claim, set payload.reviewAfter to the referenced ISO timestamp. Then take the contested_claim records with reason source_changed: the claim's source was edited after the claim was filed. Read the source's current text with search_evidence using details.sourceRef. If it still states the claim, close the record with decline_attention. If it states a different value, supersede the claim with supersede_claim_value citing an exact quote from the current text, then close the record with decline_attention. If it no longer states the claim, archive it with archive_claim_value and provenance attention:. For reason source_removed the source no longer exists: archive the claim the same way unless search_evidence finds other evidence that still states it, in which case close the record with decline_attention. 3. Take the surprisal hints listed in . These are bounded exploration hints, not evidence and not hygiene provenance. Inspect each hint's memory: subjectRef with search_evidence in the owning scope. If the source establishes a useful settled fact, use a normal content operation with an exact quote; otherwise decline_attention after inspection. Never create a claim or entity from the score alone, and never cite attention: for a content operation. -4. Find new evidence since the cutoff. Read unprocessed evidence one page at a time with search_evidence — omit the query, since, and before to take the next page of the durable delivery queue. A page holds one excerpt per source; a source with contentHasNext continues by itself on a later page, so do not page through queued sources with sourceRef and offset. File what each page establishes (or mark a source you have read in full as reviewed or deferred) before asking for the next page, so the work is kept if the pass runs out of time; stop when hasMore is false or the queue reports deliveryClosed. Use a query or sourceRef only to look up specific history or verify a citation. Prefer evidence from completed transcript sessions; historical summary rows are not part of the default delivery path. A transcript with completed: false is mid-stream — defer filing from it with the named blocker "transcript still mid-stream" (re-check completed each pass: a session still active when re-checked is a re-verified blocker, not a repeated one), and note the deferral in the pass log, because its states may be contradicted by the session's end. For each new source: +4. Find new evidence since the cutoff. Read unprocessed evidence one page at a time with search_evidence — omit the query, since, and before to take the next page of the durable delivery queue. A page holds one excerpt per source; a source with contentHasNext continues by itself on a later page, so do not page through queued sources with sourceRef and offset. File what each page establishes (or mark a source you have read in full as reviewed or deferred), then call review_evidence for that page, before asking for the next page, so the work is kept if the pass runs out of time. Acknowledge each excerpt with the sourceRef and contentOffset from the page, and omit through when you reviewed all of it; set through to an exact quote only when you stopped partway, so the rest is delivered again. Text you never acknowledge or cite stays queued for a later pass. When an excerpt reports reviewedChars, its first reviewedChars characters were reviewed by an earlier pass: do not file them again; stop when hasMore is false or the queue reports deliveryClosed. Use a query or sourceRef only to look up specific history or verify a citation. Prefer evidence from completed transcript sessions; historical summary rows are not part of the default delivery path. A transcript with completed: false is mid-stream — defer filing from it with the named blocker "transcript still mid-stream" (re-check completed each pass: a session still active when re-checked is a re-verified blocker, not a repeated one), and note the deferral in the pass log, because its states may be contradicted by the session's end. For each new source: - search_entities for subjects it establishes. - File claims only for what the source establishes as settled fact: outcomes, decisions, shipped changes, stable behavior. Do not file instructions that were merely suggested, hypotheses or diagnoses, open questions, or intermediate states of an ongoing investigation. When a source shows an attempt and its outcome, file the outcome. - A concrete deliverable the assistant produced for the user's own project, plan, or situation (a budget, schedule, draft, plan, or specific recommendation the user asked for) is a durable fact about that project: file its specifics as a claim on the project, worded as proposed or drafted rather than confirmed. Generic information not tied to the user's own circumstances is not. Work the user asks for on behalf of their own job, business, employer, trip, or event is their own situation. @@ -1194,6 +1197,7 @@ export async function renderPendingDreamingAttention( kind, status: "pending", limit: ATTENTION_ITEMS_PER_KIND + 1, + agentFacing: true, }); if (items.length === 0) continue; lines.push( @@ -2306,6 +2310,7 @@ export function finalizeDreamingPassInDb(db: WriteDb, input: DbOwnerDreamingPass deferredEvidence, withheldScopes: failedEvidence.scopes, filedSources: failedEvidence.filedSources, + filedCitations: failedEvidence.filedCitations, }); recordDreamingReviewedExcludedEvidenceInTx(db, { passId: input.passId, @@ -2314,6 +2319,7 @@ export function finalizeDreamingPassInDb(db: WriteDb, input: DbOwnerDreamingPass deferredEvidence, }); } + resolveImportedSourceAttentionInTx(db, input.passId, input.scopes); } if ( dreamingModeAdvancesEvidence( @@ -2600,7 +2606,9 @@ export async function evaluateDreamingTrigger( (await ownerQueryOne<{ present: number }>( await getDbOwnerForAccessor(accessor), "dreaming.attention.present", - "SELECT 1 AS present FROM dreaming_attention WHERE agent_id = ? AND resolved_at IS NULL LIMIT 1", + `SELECT 1 AS present FROM dreaming_attention + WHERE agent_id = ? AND resolved_at IS NULL + AND NOT ${STALLED_EVIDENCE_ATTENTION_SQL} AND NOT ${SEEN_IMPORTED_SOURCE_ATTENTION_SQL} LIMIT 1`, [agentId], { deadlineMs: 30_000, estimatedWorkUnits: 1 }, )) !== undefined; diff --git a/scripts/event-loop-contract-baseline.json b/scripts/event-loop-contract-baseline.json index 5ba1faf857..010d61e85e 100644 --- a/scripts/event-loop-contract-baseline.json +++ b/scripts/event-loop-contract-baseline.json @@ -333,196 +333,196 @@ }, { "path": "daemon.ts", - "line": 638, + "line": 639, "api": "existsSync", "source": "if (!existsSync(agentsMdPath)) return;", "category": "hot-path" }, { "path": "daemon.ts", - "line": 655, + "line": 656, "api": "existsSync", "source": "const existingFiles = files.filter((f) => existsSync(join(AGENTS_DIR, f.name)));", "category": "hot-path" }, { "path": "daemon.ts", - "line": 683, + "line": 684, "api": "existsSync", "source": "if (!existsSync(identityPath)) return \"\";", "category": "hot-path" }, { "path": "daemon.ts", - "line": 712, + "line": 713, "api": "existsSync", "source": "if (activeHarnesses.has(\"opencode\") && existsSync(opencodeDir)) {", "category": "hot-path" }, { "path": "daemon.ts", - "line": 739, + "line": 740, "api": "existsSync", "source": "if (activeHarnesses.has(\"opencode\") && existsSync(opencodeDir)) {", "category": "hot-path" }, { "path": "daemon.ts", - "line": 744, + "line": 745, "api": "existsSync", "source": "if (existsSync(geminiDir)) {", "category": "hot-path" }, { "path": "daemon.ts", - "line": 748, + "line": 749, "api": "existsSync", "source": "if (existsSync(settingsPath)) {", "category": "hot-path" }, { "path": "daemon.ts", - "line": 749, + "line": 750, "api": "readFileSync", "source": "const parsed = JSON.parse(readFileSync(settingsPath, \"utf8\"));", "category": "hot-path" }, { "path": "daemon.ts", - "line": 767, + "line": 768, "api": "existsSync", "source": "if (!existsSync(targetPath)) continue;", "category": "hot-path" }, { "path": "daemon.ts", - "line": 1321, + "line": 1322, "api": "existsSync", "source": "if (!existsSync(memoryDir)) {", "category": "hot-path" }, { "path": "daemon.ts", - "line": 1666, + "line": 1667, "api": "existsSync", "source": "if (!existsSync(p)) continue;", "category": "hot-path" }, { "path": "daemon.ts", - "line": 1668, + "line": 1669, "api": "readFileSync", "source": "const yaml = parseSimpleYaml(readFileSync(p, \"utf-8\")) as Record;", "category": "hot-path" }, { "path": "daemon.ts", - "line": 1712, + "line": 1713, "api": "existsSync", "source": "if (!existsSync(path)) continue;", "category": "hot-path" }, { "path": "daemon.ts", - "line": 1714, + "line": 1715, "api": "readFileSync", "source": "const routing = parseRoutingConfig(parseSimpleYaml(readFileSync(path, \"utf-8\")));", "category": "hot-path" }, { "path": "daemon.ts", - "line": 2143, + "line": 2144, "api": "existsSync", "source": "if (existsSync(PID_FILE)) {", "category": "hot-path" }, { "path": "daemon.ts", - "line": 2145, + "line": 2146, "api": "unlinkSync", "source": "unlinkSync(PID_FILE);", "category": "hot-path" }, { "path": "daemon.ts", - "line": 2319, + "line": 2320, "api": "mkdirSync", "source": "mkdirSync(DAEMON_DIR, { recursive: true });", "category": "hot-path" }, { "path": "daemon.ts", - "line": 2320, + "line": 2321, "api": "mkdirSync", "source": "mkdirSync(LOG_DIR, { recursive: true });", "category": "hot-path" }, { "path": "daemon.ts", - "line": 2321, + "line": 2322, "api": "mkdirSync", "source": "mkdirSync(dirname(MEMORY_DB), { recursive: true });", "category": "hot-path" }, { "path": "daemon.ts", - "line": 2466, + "line": 2467, "api": "existsSync", "source": "const workerPath = existsSync(bundled)", "category": "hot-path" }, { "path": "daemon.ts", - "line": 2538, + "line": 2539, "api": "writeFileSync", "source": "writeFileSync(PID_FILE, process.pid.toString());", "category": "hot-path" }, { "path": "daemon.ts", - "line": 2875, + "line": 2860, "api": "existsSync", "source": "isEnabled: async () => existsSync(enabledMarker),", "category": "hot-path" }, { "path": "daemon.ts", - "line": 2877, + "line": 2862, "api": "mkdirSync", "source": "mkdirSync(layout.imports, { recursive: true });", "category": "hot-path" }, { "path": "daemon.ts", - "line": 2879, + "line": 2864, "api": "writeFileSync", "source": "writeFileSync(tmp, \"enabled\\n\", { mode: 0o600 });", "category": "hot-path" }, { "path": "daemon.ts", - "line": 2880, + "line": 2865, "api": "renameSync", "source": "renameSync(tmp, enabledMarker);", "category": "hot-path" }, { "path": "daemon.ts", - "line": 3358, + "line": 3343, "api": "existsSync", "source": "if (existsSync(healthStampPath)) {", "category": "hot-path" }, { "path": "daemon.ts", - "line": 3359, + "line": 3344, "api": "readFileSync", "source": "const prev = JSON.parse(readFileSync(healthStampPath, \"utf-8\"));", "category": "hot-path" }, { "path": "daemon.ts", - "line": 3362, + "line": 3347, "api": "writeFileSync", "source": "writeFileSync(", "category": "hot-path" @@ -917,7 +917,7 @@ }, { "path": "db-owner-worker.ts", - "line": 1002, + "line": 1009, "api": "withReadDbAsync", "source": "const embedding = await getDbAccessor().withReadDbAsync(", "category": "hot-path", @@ -925,7 +925,7 @@ }, { "path": "db-owner-worker.ts", - "line": 1032, + "line": 1039, "api": "withReadDbAsync", "source": "return await getDbAccessor().withReadDbAsync(", "category": "hot-path", @@ -933,21 +933,21 @@ }, { "path": "db-owner-worker.ts", - "line": 1042, + "line": 1049, "api": "writeFileSync", "source": "writeFileSync(startedMarker, \"started\\n\");", "category": "hot-path" }, { "path": "db-owner-worker.ts", - "line": 1043, + "line": 1050, "api": "existsSync", "source": "while (!existsSync(releaseMarker)) await new Promise((resolve) => setTimeout(resolve, 5));", "category": "hot-path" }, { "path": "db-owner-worker.ts", - "line": 1180, + "line": 1188, "api": "writeFileSync", "source": "if (activeFile) writeFileSync(activeFile, `${Date.now()}\\n`);", "category": "hot-path" @@ -2526,35 +2526,35 @@ }, { "path": "pipeline/dreaming-attention.ts", - "line": 138, + "line": 140, "api": "withReadDb", "source": "return accessor.withReadDb(", "category": "hot-path" }, { "path": "pipeline/dreaming-attention.ts", - "line": 150, + "line": 152, "api": "withReadDb", "source": "return accessor.withReadDb((db: import(\"../db-accessor\").ReadDb) => {", "category": "hot-path" }, { "path": "pipeline/dreaming-attention.ts", - "line": 220, + "line": 226, "api": "withReadDb", "source": "return accessor.withReadDb((db: import(\"../db-accessor\").ReadDb) => {", "category": "hot-path" }, { "path": "pipeline/dreaming-attention.ts", - "line": 251, + "line": 257, "api": "withReadDb", "source": "return accessor.withReadDb((db: import(\"../db-accessor\").ReadDb) => {", "category": "hot-path" }, { "path": "pipeline/dreaming-attention.ts", - "line": 285, + "line": 291, "api": "withReadDb", "source": "return accessor.withReadDb(", "category": "hot-path" @@ -2663,14 +2663,14 @@ }, { "path": "pipeline/dreaming.ts", - "line": 919, + "line": 922, "api": "existsSync", "source": "if (!existsSync(path)) return { content: \"\", unreadable: false };", "category": "hot-path" }, { "path": "pipeline/dreaming.ts", - "line": 921, + "line": 924, "api": "readFileSync", "source": "const raw = readFileSync(path, \"utf-8\").trim();", "category": "hot-path" diff --git a/web/docs/src/content/docs/architecture/pipeline-storage.md b/web/docs/src/content/docs/architecture/pipeline-storage.md index 76439bc7c6..3f47e449b5 100644 --- a/web/docs/src/content/docs/architecture/pipeline-storage.md +++ b/web/docs/src/content/docs/architecture/pipeline-storage.md @@ -261,7 +261,7 @@ implementation. Signet uses SQLite in WAL mode. Migrations are numbered sequentially under `platform/core/src/migrations/`, run in order, and recorded in `schema_migrations` with checksum and timing data in -`schema_migrations_audit`. The latest migration is `168-retire-source-paragraph-claims.ts`. +`schema_migrations_audit`. The latest migration is `169-dreaming-evidence-review-cursor.ts`. ### Evidence and semantic state diff --git a/web/docs/src/content/docs/pipeline/extraction-decisions.md b/web/docs/src/content/docs/pipeline/extraction-decisions.md index d37af90c5a..c08592f883 100644 --- a/web/docs/src/content/docs/pipeline/extraction-decisions.md +++ b/web/docs/src/content/docs/pipeline/extraction-decisions.md @@ -80,17 +80,39 @@ Dreaming is an agentic pass, not a fixed per-fact classifier. Its capability registry defines the operations available to the agent, including: - `search_evidence` for immutable episodic memories, artifacts, and transcripts; - without a query it pages through the delivery queue, resuming each source - where the current pass last read it and reporting `hasMore` until the queue - is empty. A page holds one excerpt per source, and a partly read source - continues on a later page by itself. The agent files each page before asking - for the next, so a pass that runs out of time loses at most one page. Once a pass has used half its + without a query it pages through the delivery queue and reports `hasMore` + until the queue is empty. A page holds one excerpt per source, and a partly + reviewed source continues on a later page by itself. Within a pass a source + resumes where the pass last read it; a source reviewed partway in an earlier + pass resumes 200 characters before its reviewed end, snapped back to a + sentence or paragraph start, and the excerpt's `reviewedChars` counts the + leading characters that were already reviewed. The agent files each page and + acknowledges it before asking for the next, so a pass that runs out of time + loses at most one page. Once a pass has used half its timeout, the queue stops handing out new sources (`deliveryClosed`) so the pass can file what it read and record its progress; the rest goes to the next pass. With a query it searches full history, splitting it on whitespace into words that match independently (ASCII case-insensitive) and ranking fuller matches first; unspaced text such as CJK matches as one phrase +- `review_evidence` to acknowledge the excerpts a pass has reviewed. The agent + copies `sourceRef` and `contentOffset` from an excerpt delivered in the same + pass and scope, and either acknowledges the whole excerpt or names an exact + `through` quote where its review stopped. The daemon converts that to a + character offset and rejects the whole call with a specific code when the + excerpt was not delivered in this pass, belongs to another scope, does not + contain the quote, or has changed since delivery. A source's evidence cursor + advances only by text acknowledged this way or cited by a successful + `apply_ontology_ops` operation, extended contiguously from the stored cursor + when the pass finalizes. Delivered text that is never acknowledged stays + queued. A source revision delivered in three successful passes without + progress gets a visible `evidence_requeue` attention record and leaves the + continuation queue, so fresh sources are served before it. A pass that + defers the source, or withholds its progress after a failed write, does not + count toward those three. That record does + not schedule a pass by itself, and it resolves when the source makes progress + or is excluded as reviewed. Cursors written before this rule are marked + with `cursor_basis = 'delivery'`; new writes use `'review'`. - `search_entities` and `get_entity` for scoped graph reads - `list_aspect_claims` for claims with their evidence, and optionally the aspect's contradictions - `attention_list` for queued review and maintenance attention beyond what the pass prompt lists diff --git a/web/docs/src/content/docs/sources.md b/web/docs/src/content/docs/sources.md index 232caf9753..5773050639 100644 --- a/web/docs/src/content/docs/sources.md +++ b/web/docs/src/content/docs/sources.md @@ -88,7 +88,14 @@ canonical writes are verified and replayed without duplicate session, record, or turn IDs. Reconciliation exposes the equation `total = imported + duplicate + rejected + pending`; terminal jobs have zero pending. Import completion adds one Dreaming attention nudge per committed source batch. Dreaming consumption is -separate and uses its normal delivery/review path. +separate and uses its normal delivery/review path. The nudge's `source:` reference +is daemon-owned and is not shown to the Dreaming agent, which reads the imported +conversations through the delivery queue. The nudge schedules a Dreaming pass +only until a pass starts after the batch commits; the next committed batch from +that source re-arms it. A Dreaming pass resolves the nudge once +every conversation in the batch has been reviewed or excluded as reviewed. Each +pass checks a bounded, rotating slice of pending nudges, so resolution can trail +the final review by a pass or two. In a desktop-local session, **Choose from desktop** can return local paths to a loopback daemon. Remote clients must upload file bytes; a remote daemon never treats a path string as permission to read the client’s filesystem.