diff --git a/platform/core/src/migrations/170-dreaming-evidence-leases.ts b/platform/core/src/migrations/170-dreaming-evidence-leases.ts new file mode 100644 index 0000000000..f234c56b3b --- /dev/null +++ b/platform/core/src/migrations/170-dreaming-evidence-leases.ts @@ -0,0 +1,17 @@ +import type { MigrationDb } from "./contract"; + +export function up(db: MigrationDb): void { + db.exec(` + CREATE TABLE IF NOT EXISTS dreaming_evidence_leases ( + agent_id TEXT NOT NULL, + source_kind TEXT NOT NULL CHECK (source_kind IN ('memory', 'artifact', 'transcript', 'summary')), + source_id TEXT NOT NULL, + pass_id TEXT NOT NULL, + leased_at TEXT NOT NULL, + expires_at TEXT NOT NULL, + PRIMARY KEY (agent_id, source_kind, source_id) + ); + CREATE INDEX IF NOT EXISTS idx_dreaming_evidence_leases_pass + ON dreaming_evidence_leases(pass_id); + `); +} diff --git a/platform/core/src/migrations/index.ts b/platform/core/src/migrations/index.ts index 429f4aaf00..1e8e0b768f 100644 --- a/platform/core/src/migrations/index.ts +++ b/platform/core/src/migrations/index.ts @@ -170,6 +170,7 @@ 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"; +import { up as dreamingEvidenceLeases } from "./170-dreaming-evidence-leases"; export type { Migration, MigrationArtifacts, MigrationDb } from "./contract"; export const MIGRATIONS: readonly Migration[] = [ @@ -1539,6 +1540,12 @@ export const MIGRATIONS: readonly Migration[] = [ ], }, }, + { + version: 170, + name: "dreaming-evidence-leases", + up: dreamingEvidenceLeases, + artifacts: { tables: ["dreaming_evidence_leases"], indexes: ["idx_dreaming_evidence_leases_pass"] }, + }, ]; 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 ab7df601f9..1bfc7b2ca9 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(169); + expect(applied.version).toBe(170); expect( db.query("SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'vector_repair_checkpoints'").get(), ).toEqual({ name: "vector_repair_checkpoints" }); diff --git a/platform/core/src/types.ts b/platform/core/src/types.ts index 35d9ba4a2a..9bb20a06dc 100644 --- a/platform/core/src/types.ts +++ b/platform/core/src/types.ts @@ -458,6 +458,7 @@ export interface DreamingConfig { readonly maxInputTokens: number; readonly maxOutputTokens: number | null; readonly maxConcurrentPasses: number; + readonly maxPassesPerScope?: number; readonly codemode: boolean; readonly backfillOnFirstRun: boolean; readonly surprisal?: DreamingSurprisalConfig; diff --git a/platform/daemon/src/memory-config.ts b/platform/daemon/src/memory-config.ts index 504ab4d667..3fc8e497bf 100644 --- a/platform/daemon/src/memory-config.ts +++ b/platform/daemon/src/memory-config.ts @@ -89,6 +89,7 @@ export const DEFAULT_DREAMING: DreamingConfig = { maxInputTokens: 128_000, maxOutputTokens: null, maxConcurrentPasses: 2, + maxPassesPerScope: 1, codemode: false, backfillOnFirstRun: true, surprisal: DEFAULT_DREAMING_SURPRISAL, @@ -985,6 +986,9 @@ export function loadDreamingConfig(yaml: Record): DreamingConfi maxConcurrentPasses: Math.floor( clampWarn("maxConcurrentPasses", raw.maxConcurrentPasses, 1, 16, dd.maxConcurrentPasses), ), + maxPassesPerScope: Math.floor( + clampWarn("maxPassesPerScope", raw.maxPassesPerScope, 1, 16, dd.maxPassesPerScope ?? 1), + ), codemode: typeof raw.codemode === "boolean" ? raw.codemode : dd.codemode, backfillOnFirstRun: typeof raw.backfillOnFirstRun === "boolean" ? raw.backfillOnFirstRun : dd.backfillOnFirstRun, surprisal: { diff --git a/platform/daemon/src/ontology-proposals.test.ts b/platform/daemon/src/ontology-proposals.test.ts index 10e4ce0ad3..9d24ff8e7c 100644 --- a/platform/daemon/src/ontology-proposals.test.ts +++ b/platform/daemon/src/ontology-proposals.test.ts @@ -3565,6 +3565,48 @@ describe("ontology proposals", () => { expect(attributeTime(morningId)).toMatchObject({ status: "superseded", superseded_by: laterId }); }); + it("orders by capture when only one of two contradicting claims carries an event time", async () => { + getDbAccessor().withWriteTx((db) => { + const source = db.prepare( + `INSERT INTO memories + (id, content, type, agent_id, visibility, memory_kind, created_at, updated_at) + VALUES (?, ?, 'fact', 'default', 'global', 'episodic', ?, ?)`, + ); + source.run("austin-source", "I live in Austin.", "2024-02-01T10:00:00.000Z", "2024-02-01T10:00:00.000Z"); + source.run( + "denver-source", + "I moved to Denver in May.", + "2024-06-01T10:00:00.000Z", + "2024-06-01T10:00:00.000Z", + ); + }); + const residence = { ...slot, aspect: "home", claim_key: "city" }; + const file = async (sourceId: string, payload: Record) => + await applyOntologyOperation(getDbAccessor(), { + agentId: "default", + actor: "dreaming", + operation: "set_claim_value", + payload: { ...residence, ...payload }, + evidence: [{ source_ref: `memory:${sourceId}`, source_kind: "manual", source_id: sourceId }], + sourceKind: "memory", + sourceId, + }); + const denver = await file("denver-source", { + value: "The user moved to Denver in May 2024.", + valid_from: "2024-05-01", + time_precision: "month", + }); + const austin = await file("austin-source", { value: "The user lives in Austin." }); + const denverId = denver.result?.attributeId; + const austinId = austin.result?.attributeId; + if (typeof denverId !== "string" || typeof austinId !== "string") + throw new Error("attribute ids were not returned"); + + expect(austin.result?.supersededByNewerEvidence).toBe(denverId); + expect(attributeTime(denverId)?.status).toBe("active"); + expect(attributeTime(austinId)).toMatchObject({ status: "superseded", superseded_by: denverId }); + }); + it("keeps the newer claim current when an older claim arrives later", async () => { const residence = { ...slot, aspect: "home", claim_key: "city" }; const newer = await applyOntologyOperation(getDbAccessor(), { diff --git a/platform/daemon/src/ontology-proposals.ts b/platform/daemon/src/ontology-proposals.ts index feb84f90b4..7004c355ff 100644 --- a/platform/daemon/src/ontology-proposals.ts +++ b/platform/daemon/src/ontology-proposals.ts @@ -1356,7 +1356,7 @@ function applySetClaimValue( const newerActive = active .filter((row) => { const activeTime = claimEvidenceTime(row); - if (activeTime !== incomingTime) return activeTime !== null && incomingTime !== null && activeTime > incomingTime; + if (activeTime !== null && incomingTime !== null && activeTime !== incomingTime) return activeTime > incomingTime; const activeCapturedAt = claimSourceCapturedAt(db, agentId, row.source_kind, row.source_id); return activeCapturedAt !== null && incomingCapturedAt !== null && activeCapturedAt > incomingCapturedAt; }) diff --git a/platform/daemon/src/pipeline/dreaming-capabilities.ts b/platform/daemon/src/pipeline/dreaming-capabilities.ts index 29b31e6e4e..fec95fdd47 100644 --- a/platform/daemon/src/pipeline/dreaming-capabilities.ts +++ b/platform/daemon/src/pipeline/dreaming-capabilities.ts @@ -36,6 +36,7 @@ import { pendingDreamingEvidenceContinuations, reviewDreamingEvidenceInDb, } from "./dreaming-evidence-consumption"; +import { evidenceLeasedByOtherPasses, leaseDreamingEvidence } from "./dreaming-evidence-leases"; import { DREAMING_ONTOLOGY_OPERATION_SCHEMA } from "./dreaming-operation-contract"; import { type ApplyDreamingOperationsResult, @@ -221,6 +222,7 @@ export interface CreateDreamingCapabilitiesParams { readonly passId?: string; readonly evidenceDeliveryDeadline?: number; readonly evidenceChars?: number; + readonly evidenceLeaseMs?: number; readonly mode?: DreamingCapabilityMode; readonly writeCaps?: GraphWriteCaps; readonly onOperationsApplied?: ( @@ -312,6 +314,11 @@ export function searchDreamingEvidenceInDb(db: ReadDb, input: DbOwnerDreamingEvi } const DELIVERY_QUEUE_SCAN_LIMIT = 51; +const EVIDENCE_LEASE_ATTEMPTS = 4; +const EVIDENCE_HELD_BY_OTHER_PASSES = { + heldByOtherPasses: true, + note: "Nothing unclaimed is left in the queue for this pass: other running Dreaming passes in this scope hold the rest of the queued evidence and will review it. Do not keep paging; file what you have read, write the runbook, and finish.", +} as const; function drainDreamingEvidenceQueueInDb(db: ReadDb, input: DbOwnerDreamingEvidenceSearch): DreamingCapabilityOutput { const scopeId = input.agentId; @@ -328,12 +335,16 @@ function drainDreamingEvidenceQueueInDb(db: ReadDb, input: DbOwnerDreamingEviden }; const stalled = (source: EpisodicSourceRecord): boolean => cursorFor(source).stalledPasses >= DREAMING_EVIDENCE_STALL_PASSES; + const leasedElsewhere = evidenceLeasedByOtherPasses(db, scopeId, input.passId); const scanned = searchEpisodicSources(db, { agentId: scopeId, query: "", kind: input.kind, excludeDelivered: true, - excludeSourceRefs: input.passId ? passFullyServedSourceRefs(db, input.passId, scopeId) : [], + excludeSourceRefs: [ + ...(input.passId ? passFullyServedSourceRefs(db, input.passId, scopeId) : []), + ...leasedElsewhere, + ], limit: DELIVERY_QUEUE_SCAN_LIMIT, }); const fresh = [...scanned.filter((source) => !stalled(source)), ...scanned.filter(stalled)]; @@ -365,7 +376,7 @@ function drainDreamingEvidenceQueueInDb(db: ReadDb, input: DbOwnerDreamingEviden } return items; }; - const continuationRefs = new Set(); + const continuationRefs = new Set(leasedElsewhere); const continuations = page( pendingDreamingEvidenceContinuations(db, scopeId, 50, input.kind), limit + 1, @@ -393,6 +404,7 @@ function drainDreamingEvidenceQueueInDb(db: ReadDb, input: DbOwnerDreamingEviden budgetExhausted || returned.some((item) => item.contentHasNext === true) || (items.length > 0 && fresh.length >= DELIVERY_QUEUE_SCAN_LIMIT), + ...(returned.length === 0 && leasedElsewhere.length > 0 ? EVIDENCE_HELD_BY_OTHER_PASSES : {}), }; } @@ -660,7 +672,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 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.", + "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 for this pass. When other Dreaming passes run in the same scope, each queued source is delivered to only one of them; a page with heldByOtherPasses: true means those passes hold the rest of the queue, so stop paging. 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), @@ -703,25 +715,49 @@ export function createDreamingCapabilities(params: CreateDreamingCapabilitiesPar ...(params.passId === undefined ? {} : { passId: params.passId }), ...(params.evidenceChars === undefined ? {} : { evidenceChars: params.evidenceChars }), }; - return await runDbOwnerDomainOperation(accessor, { - runWithOwner: async (owner) => { - const handle = owner.submit( - { - kind: "dreaming_evidence_search", - input, - }, - { - operation: "dreaming.capabilities.search-evidence", - lane: "read", - workloadClass: "foreground", - deadlineMs: 30_000, - estimatedWorkUnits: 200, - }, - ); - return await handle.result; - }, - runInline: ({ read }) => read((db) => searchDreamingEvidenceInDb(db, input)), - }); + const search = () => + runDbOwnerDomainOperation(accessor, { + runWithOwner: async (owner) => { + const handle = owner.submit( + { + kind: "dreaming_evidence_search", + input, + }, + { + operation: "dreaming.capabilities.search-evidence", + lane: "read", + workloadClass: "foreground", + deadlineMs: 30_000, + estimatedWorkUnits: 200, + }, + ); + return await handle.result; + }, + runInline: ({ read }) => read((db) => searchDreamingEvidenceInDb(db, input)), + }); + const drainsQueue = sourceRef === undefined && !query?.trim() && since === undefined && before === undefined; + const { evidenceLeaseMs, passId } = params; + if (!drainsQueue || evidenceLeaseMs === undefined || passId === undefined) return await search(); + for (let attempt = 1; ; attempt++) { + const result = await search(); + if (result.ok !== true || !Array.isArray(result.items) || result.items.length === 0) return result; + const items = result.items as ReadonlyArray>; + const held = await leaseDreamingEvidence(accessor, { + agentId: scopeId, + passId, + sourceRefs: items.flatMap((item) => (typeof item.sourceRef === "string" ? [item.sourceRef] : [])), + ttlMs: evidenceLeaseMs, + }); + const leased = items.filter((item) => typeof item.sourceRef === "string" && held.has(item.sourceRef)); + if (leased.length === items.length) return result; + if (attempt < EVIDENCE_LEASE_ATTEMPTS) continue; + return { + ...result, + items: leased, + hasMore: true, + note: "Other passes in this scope claimed part of this page first. Call search_evidence again without a query for the next unclaimed page.", + }; + } }, ), capability( diff --git a/platform/daemon/src/pipeline/dreaming-evidence-leases.test.ts b/platform/daemon/src/pipeline/dreaming-evidence-leases.test.ts new file mode 100644 index 0000000000..24f517aa78 --- /dev/null +++ b/platform/daemon/src/pipeline/dreaming-evidence-leases.test.ts @@ -0,0 +1,380 @@ +import { Database } from "bun:sqlite"; +import { afterEach, beforeEach, describe, expect, it } from "bun:test"; +import type { DreamingConfig } from "@signet/core"; +import { runMigrations } from "../../../core/src/migrations"; +import type { DbAccessor } from "../db-accessor"; +import { createDreamingAgentTools } from "./dreaming-agent-tools"; +import { type DreamingAgentExecutor, getDreamingToolCalls, runDreamingAgentPass } from "./dreaming"; +import { leaseDreamingEvidence } from "./dreaming-evidence-leases"; + +const AGENT = "default"; + +function cfg(): DreamingConfig { + return { + enabled: true, + tokenThreshold: 100_000, + maxInterval: 6 * 60 * 60 * 1_000, + maxInputTokens: 32_000, + maxOutputTokens: 16_000, + maxConcurrentPasses: 2, + codemode: false, + timeout: 300_000, + backfillOnFirstRun: false, + }; +} + +function wrapDb(db: Database): DbAccessor { + const tx = (fn: (db: Database) => T): T => { + db.exec("BEGIN IMMEDIATE"); + try { + const result = fn(db); + db.exec("COMMIT"); + return result; + } catch (e) { + db.exec("ROLLBACK"); + throw e; + } + }; + return { + withReadDb: (fn: (db: Database) => T): T => fn(db), + withReadDbAsync: (fn: (db: Database) => Promise): Promise => fn(db), + withWriteTx: tx, + withWriteTxAsync: (fn: (db: Database) => T): Promise => Promise.resolve().then(() => tx(fn)), + } as unknown as DbAccessor; +} + +type Tool = ReturnType[number]; + +function result(res: { content: ReadonlyArray }): Record { + const first = res.content[0] as { text?: string } | undefined; + return JSON.parse(first?.text ?? "{}") as Record; +} + +function refsOf(output: Record): string[] { + return ((output.items as Array<{ sourceRef: string }> | undefined) ?? []).map((item) => item.sourceRef); +} + +function tool(tools: readonly Tool[], name: string): Tool { + const found = tools.find((candidate) => candidate.name === name); + if (!found) throw new Error(`Missing ${name}`); + return found; +} + +async function readAndReview(tools: readonly Tool[], input: Record): Promise { + const output = result(await tool(tools, "search_evidence").execute("read", input, undefined, undefined, {} as never)); + const items = ((output.items as Array<{ sourceRef: string; contentOffset: number }> | undefined) ?? []).map( + (item) => ({ sourceRef: item.sourceRef, contentOffset: item.contentOffset }), + ); + if (items.length === 0) return; + const reviewed = result( + await tool(tools, "review_evidence").execute( + "review", + { agentId: AGENT, items }, + undefined, + undefined, + {} as never, + ), + ); + if (reviewed.ok !== true) throw new Error(`review_evidence failed: ${JSON.stringify(reviewed)}`); +} + +describe("Dreaming evidence leases", () => { + let db: Database; + let accessor: DbAccessor; + + beforeEach(() => { + db = new Database(":memory:"); + runMigrations(db as unknown as Parameters[0]); + accessor = wrapDb(db); + }); + + afterEach(() => { + db.close(); + }); + + function seedTranscripts(count: number): string[] { + const insert = db.prepare( + `INSERT INTO session_transcripts + (session_key, agent_id, content, harness, created_at, updated_at, completed_at) + VALUES (?, ?, ?, 'pi', ?, ?, ?)`, + ); + return Array.from({ length: count }, (_, index) => { + const key = `session-${String(index).padStart(2, "0")}`; + const at = new Date(Date.UTC(2026, 0, 1, index)).toISOString(); + insert.run(key, AGENT, `User: Session ${index} says the user lives in city ${index}.`, at, at, at); + return `transcript:${key}`; + }); + } + + function startPassRow(passId: string): void { + db.prepare( + `INSERT INTO dreaming_passes (id, agent_id, mode, status, started_at, created_at) + VALUES (?, ?, 'incremental', 'running', datetime('now'), datetime('now'))`, + ).run(passId, AGENT); + } + + async function drain(passId: string, limit: number): Promise> { + const tools = createDreamingAgentTools({ + accessor, + agentId: AGENT, + allowedAgentIds: [AGENT], + actor: "dreaming", + passId, + evidenceLeaseMs: 60_000, + }); + return result( + await tool(tools, "search_evidence").execute( + "call", + { agentId: AGENT, limit }, + undefined, + undefined, + {} as never, + ), + ); + } + + it("hands concurrent passes in one scope disjoint evidence", async () => { + const all = seedTranscripts(6); + startPassRow("pass-a"); + startPassRow("pass-b"); + + const first = await drain("pass-a", 3); + const second = await drain("pass-b", 3); + const a = refsOf(first); + const b = refsOf(second); + + expect(a).toHaveLength(3); + expect(b).toHaveLength(3); + expect(a.filter((ref) => b.includes(ref))).toEqual([]); + expect([...a, ...b].sort()).toEqual([...all].sort()); + expect(refsOf(await drain("pass-a", 6)).sort()).toEqual([...a].sort()); + expect(refsOf(await drain("pass-c", 6))).toEqual([]); + }); + + it("gives a contested source to exactly one pass", async () => { + const refs = seedTranscripts(4); + startPassRow("pass-a"); + startPassRow("pass-b"); + + const [a, b] = await Promise.all([ + leaseDreamingEvidence(accessor, { agentId: AGENT, passId: "pass-a", sourceRefs: refs, ttlMs: 60_000 }), + leaseDreamingEvidence(accessor, { agentId: AGENT, passId: "pass-b", sourceRefs: refs, ttlMs: 60_000 }), + ]); + + for (const ref of refs) expect(Number(a.has(ref)) + Number(b.has(ref))).toBe(1); + expect( + await leaseDreamingEvidence(accessor, { agentId: AGENT, passId: "pass-a", sourceRefs: refs, ttlMs: 60_000 }), + ).toEqual(a); + }); + + async function drainWhileRacing(raced: readonly string[], limit: number): Promise> { + const tools = createDreamingAgentTools({ + accessor, + agentId: AGENT, + allowedAgentIds: [AGENT], + actor: "dreaming", + passId: "pass-b", + evidenceLeaseMs: 60_000, + }); + const prepare = db.prepare.bind(db); + let raceRun = false; + (db as unknown as { prepare: Database["prepare"] }).prepare = ((sql: string) => { + if (!raceRun && sql.includes("INSERT INTO dreaming_evidence_leases")) { + raceRun = true; + for (const ref of raced) { + prepare( + `INSERT INTO dreaming_evidence_leases (agent_id, source_kind, source_id, pass_id, leased_at, expires_at) + VALUES (?, 'transcript', ?, 'pass-a', datetime('now'), datetime('now', '+1 hour'))`, + ).run(AGENT, ref.slice("transcript:".length)); + } + } + return prepare(sql); + }) as Database["prepare"]; + try { + return result( + await tool(tools, "search_evidence").execute( + "call", + { agentId: AGENT, limit }, + undefined, + undefined, + {} as never, + ), + ); + } finally { + Reflect.deleteProperty(db, "prepare"); + expect(raceRun).toBe(true); + } + } + + it("fills a page with unclaimed sources when another pass wins a lease race", async () => { + const refs = seedTranscripts(4); + startPassRow("pass-a"); + startPassRow("pass-b"); + const raced = refs.slice(-1); + + const output = await drainWhileRacing(raced, 3); + + expect(refsOf(output)).toHaveLength(3); + expect(refsOf(output).sort()).toEqual(refs.filter((ref) => !raced.includes(ref)).sort()); + expect(output.hasMore).toBe(false); + expect(output.heldByOtherPasses).toBeUndefined(); + expect( + db.prepare("SELECT pass_id AS passId, COUNT(*) AS n FROM dreaming_evidence_leases GROUP BY pass_id").all(), + ).toEqual([ + { passId: "pass-a", n: 1 }, + { passId: "pass-b", n: 3 }, + ]); + }); + + it("ends the queue for a pass when other passes hold everything left", async () => { + const refs = seedTranscripts(2); + startPassRow("pass-a"); + startPassRow("pass-b"); + + const output = await drainWhileRacing(refs, 2); + + expect(refsOf(output)).toEqual([]); + expect(output.hasMore).toBe(false); + expect(output.heldByOtherPasses).toBe(true); + expect(String(output.note)).toContain("other running Dreaming passes"); + expect(await drain("pass-c", 5)).toMatchObject({ items: [], hasMore: false, heldByOtherPasses: true }); + }); + + it("redelivers the evidence of a pass that failed or whose lease expired", async () => { + seedTranscripts(4); + startPassRow("pass-a"); + startPassRow("pass-b"); + startPassRow("pass-c"); + const failed = refsOf(await drain("pass-a", 2)); + const expiring = refsOf(await drain("pass-b", 4)); + expect(failed).toHaveLength(2); + expect(expiring.some((ref) => failed.includes(ref))).toBe(false); + + db.prepare("UPDATE dreaming_passes SET status = 'failed' WHERE id = 'pass-a'").run(); + expect(refsOf(await drain("pass-c", 4)).sort()).toEqual([...failed].sort()); + + db.prepare( + "UPDATE dreaming_evidence_leases SET expires_at = datetime('now', '-1 minute') WHERE pass_id = 'pass-b'", + ).run(); + startPassRow("pass-d"); + expect(refsOf(await drain("pass-d", 4)).sort()).toEqual([...expiring].sort()); + }); + + it("does not settle a source another running pass holds", async () => { + const [held, ...rest] = seedTranscripts(3); + if (held === undefined) throw new Error("missing seed"); + startPassRow("holder"); + await leaseDreamingEvidence(accessor, { agentId: AGENT, passId: "holder", sourceRefs: [held], ttlMs: 60_000 }); + + const executor: DreamingAgentExecutor = { + async run(input) { + const tools = input.tools as readonly Tool[]; + await readAndReview(tools, { agentId: AGENT, sourceRef: held, chunkSize: 4_000 }); + await readAndReview(tools, { agentId: AGENT }); + return { summary: "Read the queue." }; + }, + }; + const pass = await runDreamingAgentPass( + accessor, + executor, + cfg(), + "/tmp", + AGENT, + [AGENT], + "incremental", + undefined, + undefined, + { + sharedScope: { evidenceLeaseMs: 60_000, attentionScopes: [AGENT] }, + }, + ); + + const settled = db + .prepare( + `SELECT source_id AS id FROM dreaming_evidence_consumption + WHERE agent_id = ? AND delivered_offset >= source_length ORDER BY source_id`, + ) + .all(AGENT) as Array<{ id: string }>; + expect(settled.map((row) => `transcript:${row.id}`)).toEqual(rest); + expect(db.prepare("SELECT pass_id AS passId FROM dreaming_evidence_leases").all()).toEqual([{ passId: "holder" }]); + + db.prepare("UPDATE dreaming_passes SET status = 'failed' WHERE id = 'holder'").run(); + startPassRow("next"); + expect(refsOf(await drain("next", 5))).toEqual([held]); + expect(pass.passId).toBeString(); + }); + + it("drains one scope with concurrent passes that each settle only their own evidence", async () => { + const all = seedTranscripts(8); + db.prepare( + `INSERT INTO dreaming_attention (id, agent_id, kind, subject_ref, details_json, priority) + VALUES ('contested', ?, 'contested_claim', 'memory:claim', '{}', 90)`, + ).run(AGENT); + let release: () => void = () => undefined; + const barrier = new Promise((resolve) => { + release = resolve; + }); + let reading = 0; + const prompts = new Map(); + const executor = (name: string): DreamingAgentExecutor => ({ + async run(input) { + prompts.set(name, input.prompt); + await readAndReview(input.tools as readonly Tool[], { agentId: AGENT, limit: 4 }); + reading++; + await barrier; + return { summary: `${name} read its page.` }; + }, + }); + const run = (name: string, attentionScopes: readonly string[]) => + runDreamingAgentPass( + accessor, + executor(name), + cfg(), + "/tmp", + AGENT, + [AGENT], + "incremental", + undefined, + undefined, + { + sharedScope: { evidenceLeaseMs: 60_000, attentionScopes }, + }, + ); + + const first = run("first", [AGENT]); + const second = run("second", []); + const deadline = Date.now() + 5_000; + while (reading < 2 && Date.now() < deadline) await new Promise((resolve) => setTimeout(resolve, 10)); + expect(reading).toBe(2); + release(); + const [a, b] = await Promise.all([first, second]); + + const delivered = async (passId: string) => + (await getDreamingToolCalls(accessor, AGENT, passId)) + .filter((call) => call.toolName === "search_evidence") + .flatMap((call) => refsOf(call.output as Record)); + const aRefs = await delivered(a.passId); + const bRefs = await delivered(b.passId); + expect(aRefs).toHaveLength(4); + expect(bRefs).toHaveLength(4); + expect(aRefs.filter((ref) => bRefs.includes(ref))).toEqual([]); + const settled = db + .prepare( + `SELECT source_id AS id, pass_id AS passId FROM dreaming_evidence_consumption + WHERE agent_id = ? AND delivered_offset >= source_length`, + ) + .all(AGENT) as Array<{ id: string; passId: string }>; + expect(settled.map((row) => `transcript:${row.id}`).sort()).toEqual([...all].sort()); + for (const row of settled) { + expect(row.passId).toBe(aRefs.includes(`transcript:${row.id}`) ? a.passId : b.passId); + } + expect(db.prepare("SELECT COUNT(*) AS n FROM dreaming_evidence_leases").get()).toEqual({ n: 0 }); + const pending = (name: string) => { + const prompt = prompts.get(name) ?? ""; + return prompt.slice(prompt.lastIndexOf(""), prompt.lastIndexOf("")); + }; + expect(pending("first")).toContain('"kind":"contested_claim"'); + expect(pending("second")).toContain("none pending"); + expect(pending("second")).not.toContain("contested_claim"); + }); +}); diff --git a/platform/daemon/src/pipeline/dreaming-evidence-leases.ts b/platform/daemon/src/pipeline/dreaming-evidence-leases.ts new file mode 100644 index 0000000000..eb11e58fa9 --- /dev/null +++ b/platform/daemon/src/pipeline/dreaming-evidence-leases.ts @@ -0,0 +1,71 @@ +import type { DbAccessor, ReadDb, WriteDb } from "../db-accessor"; +import { ownerChanges, ownerRunStatement, ownerTransaction } from "../db-owner-maintenance"; +import { getDbOwnerForAccessor } from "../db-owner-runtime"; +const LIVE_LEASE = `julianday(l.expires_at) > julianday('now') + AND EXISTS (SELECT 1 FROM dreaming_passes p WHERE p.id = l.pass_id AND p.status = 'running')`; + +function leaseTableExists(db: ReadDb): boolean { + return ( + db.prepare("SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = 'dreaming_evidence_leases'").get() != null + ); +} + +export function evidenceLeasedByOtherPasses(db: ReadDb, agentId: string, passId: string | undefined): string[] { + if (!leaseTableExists(db)) return []; + const rows = db + .prepare( + `SELECT l.source_kind AS kind, l.source_id AS id FROM dreaming_evidence_leases l + WHERE l.agent_id = ? AND l.pass_id != ? AND ${LIVE_LEASE}`, + ) + .all(agentId, passId ?? "") as Array<{ kind: string; id: string }>; + return rows.map((row) => `${row.kind}:${row.id}`); +} + +export function evidenceLeasedByOtherPassesInScopes( + db: ReadDb, + passId: string, + scopes: readonly string[], +): ReadonlySet { + return new Set( + scopes.flatMap((scope) => evidenceLeasedByOtherPasses(db, scope, passId).map((ref) => `${scope}\u0000${ref}`)), + ); +} +export function releaseDreamingEvidenceLeasesInTx(db: WriteDb, passId: string): void { + if (!leaseTableExists(db)) return; + db.prepare( + `DELETE FROM dreaming_evidence_leases + WHERE pass_id = ? OR pass_id NOT IN (SELECT id FROM dreaming_passes WHERE status = 'running')`, + ).run(passId); +} +export async function leaseDreamingEvidence( + accessor: DbAccessor, + params: { + readonly agentId: string; + readonly passId: string; + readonly sourceRefs: readonly string[]; + readonly ttlMs: number; + }, +): Promise> { + const refs = [...new Set(params.sourceRefs)].flatMap((ref) => { + const separator = ref.indexOf(":"); + return separator > 0 ? [{ ref, kind: ref.slice(0, separator), id: ref.slice(separator + 1) }] : []; + }); + if (refs.length === 0) return new Set(); + const ttl = `+${Math.max(1, Math.ceil(params.ttlMs / 1000))} seconds`; + const results = await ownerTransaction( + await getDbOwnerForAccessor(accessor), + "dreaming.evidence.lease", + refs.map(({ kind, id }) => + ownerRunStatement( + `INSERT INTO dreaming_evidence_leases AS l (agent_id, source_kind, source_id, pass_id, leased_at, expires_at) + VALUES (?, ?, ?, ?, strftime('%Y-%m-%d %H:%M:%f', 'now'), strftime('%Y-%m-%d %H:%M:%f', 'now', ?)) + ON CONFLICT(agent_id, source_kind, source_id) DO UPDATE SET + pass_id = excluded.pass_id, leased_at = excluded.leased_at, expires_at = excluded.expires_at + WHERE l.pass_id = excluded.pass_id OR NOT (${LIVE_LEASE})`, + [params.agentId, kind, id, params.passId, ttl], + ), + ), + { deadlineMs: 30_000, estimatedWorkUnits: refs.length }, + ); + return new Set(refs.flatMap(({ ref }, index) => (ownerChanges(results[index]) > 0 ? [ref] : []))); +} diff --git a/platform/daemon/src/pipeline/dreaming-worker.test.ts b/platform/daemon/src/pipeline/dreaming-worker.test.ts index db6032d4bf..3bf40a2bdf 100644 --- a/platform/daemon/src/pipeline/dreaming-worker.test.ts +++ b/platform/daemon/src/pipeline/dreaming-worker.test.ts @@ -35,6 +35,7 @@ import { recallThroughDbOwner } from "../db-owner-recall"; import { reportEventLoopLag, resetPressureState } from "../system-pressure"; import { DREAMING_AGENT_PROMPT, + DREAMING_SCHEDULE_BACKLOG_MAX_SOURCES, type DreamingAgentExecutor, type DreamingPassFocus, dreamingFocusOfMode, @@ -1025,6 +1026,241 @@ describe("dreaming worker agent scope", () => { expect(partitionDreamingScopes([{ scope: "busy", tokens: 500 }, ...flagged], 2)).toEqual([["busy"], ["oldest"]]); }); + it("adds passes for a large scope only from slots the other scopes leave free", () => { + expect( + partitionDreamingScopes( + [ + { scope: "large", tokens: 900_000, passes: 3 }, + { scope: "small", tokens: 500, passes: 1 }, + ], + 4, + ), + ).toEqual([["large"], ["small"], ["large"], ["large"]]); + expect( + partitionDreamingScopes( + [ + { scope: "large", tokens: 900_000, passes: 3 }, + { scope: "small", tokens: 500, passes: 1 }, + ], + 2, + ), + ).toEqual([["large"], ["small"]]); + expect( + partitionDreamingScopes( + [ + { scope: "large", tokens: 900_000, passes: 4 }, + { scope: "other", tokens: 800_000, passes: 4 }, + ], + 5, + ), + ).toEqual([["large"], ["other"], ["large"], ["other"], ["large"]]); + expect( + partitionDreamingScopes( + [ + { scope: "held", tokens: 900_000, passes: 0 }, + { scope: "flagged", tokens: 0, attention: true, passes: 0 }, + ], + 3, + ), + ).toEqual([]); + }); + + function seedLargeScope(count: number): void { + const seed = db.prepare( + `INSERT INTO session_transcripts + (session_key, agent_id, content, harness, created_at, updated_at, completed_at) + VALUES (?, 'default', ?, 'pi', datetime('now'), datetime('now'), datetime('now'))`, + ); + for (let index = 0; index < count; index++) { + seed.run(`large-${index}`, `User: in session ${index} I mentioned a durable fact. `.repeat(40)); + } + } + + async function runSharedScopeWorker( + cfg: Partial, + expectedRunning: number, + ): Promise<{ readonly delivered: string[][]; readonly scopes: string[][]; readonly peak: number }> { + const delivered: string[][] = []; + const scopes: string[][] = []; + let running = 0; + let peak = 0; + let releaseBarrier: () => void = () => undefined; + const barrier = new Promise((resolve) => { + releaseBarrier = resolve; + }); + const previousLimit = getLlmConcurrencyLimit(); + configureLlmConcurrency(4); + const worker = startDreamingWorker( + accessor, + defaultCfg({ maxConcurrentPasses: 3, tokenThreshold: 10_000, ...cfg }), + agentsDir, + "default", + { + checkIntervalMs: 60_000, + executorFactory: () => ({ + async run(input) { + running++; + peak = Math.max(peak, running); + const search = input.tools.find((tool) => tool.name === "search_evidence"); + if (!search) throw new Error("Missing search_evidence"); + let refs: string[] = []; + for (let attempt = 0; attempt < 5 && refs.length === 0; attempt++) { + const page = await search.execute( + "drain", + { agentId: "default", limit: 5 }, + undefined, + undefined, + {} as never, + ); + const output = JSON.parse((page.content[0] as { text: string }).text); + refs = (output.items as Array<{ sourceRef: string }>).map((item) => item.sourceRef); + if (output.hasMore !== true) break; + } + delivered.push(refs); + const review = input.tools.find((tool) => tool.name === "review_evidence"); + if (refs.length > 0 && review) { + await review.execute( + "review", + { agentId: "default", items: refs.map((sourceRef) => ({ sourceRef, contentOffset: 0 })) }, + undefined, + undefined, + {} as never, + ); + } + scopes.push([...(worker.activePasses.find((pass) => pass.passId === input.passId)?.scopes ?? [])]); + await Promise.race([barrier, new Promise((resolve) => setTimeout(resolve, 3_000))]); + running--; + return { summary: "Read one page" }; + }, + }), + }, + ); + try { + await worker.triggerAsync("incremental"); + await waitFor(() => delivered.length === expectedRunning, 3_000); + await new Promise((resolve) => setTimeout(resolve, 100)); + if (expectedRunning === 1) { + await expect(worker.triggerAsync("incremental")).rejects.toBeInstanceOf(AlreadyRunningError); + } + releaseBarrier(); + await waitFor(() => !worker.running, 5_000); + return { delivered, scopes, peak }; + } finally { + worker.stop(); + configureLlmConcurrency(previousLimit); + } + } + + it("drains one large scope with several passes that never share a source", async () => { + seedLargeScope(40); + const { delivered, scopes, peak } = await runSharedScopeWorker({ maxPassesPerScope: 3 }, 3); + + expect(peak).toBe(3); + expect(scopes).toEqual([["default"], ["default"], ["default"]]); + expect(delivered.every((refs) => refs.length > 0)).toBe(true); + const all = delivered.flat(); + expect(new Set(all).size).toBe(all.length); + expect(db.prepare("SELECT status, COUNT(*) AS n FROM dreaming_passes GROUP BY status").all()).toEqual([ + { status: "completed", n: 3 }, + ]); + expect( + db + .prepare("SELECT COUNT(*) AS n FROM dreaming_evidence_consumption WHERE delivered_offset >= source_length") + .get(), + ).toEqual({ n: all.length }); + expect(db.prepare("SELECT COUNT(*) AS n FROM dreaming_evidence_leases").get()).toEqual({ n: 0 }); + }); + + it("keeps one pass per scope unless maxPassesPerScope allows more", async () => { + seedLargeScope(40); + const { delivered, peak } = await runSharedScopeWorker({}, 1); + + expect(peak).toBe(1); + expect(delivered).toHaveLength(1); + expect(db.prepare("SELECT COUNT(*) AS n FROM dreaming_evidence_leases").get()).toEqual({ n: 0 }); + }); + + it("hands pending attention to a joining pass once its owner finishes", async () => { + seedLargeScope(40); + db.prepare( + `INSERT INTO dreaming_attention (id, agent_id, kind, subject_ref, details_json, priority) + VALUES ('contested', 'default', 'contested_claim', 'memory:claim', '{}', 90)`, + ).run(); + const prompts: string[] = []; + const finish: Array<() => void> = []; + const previousLimit = getLlmConcurrencyLimit(); + configureLlmConcurrency(4); + const worker = startDreamingWorker( + accessor, + defaultCfg({ maxConcurrentPasses: 2, maxPassesPerScope: 2, tokenThreshold: 10_000 }), + agentsDir, + "default", + { + checkIntervalMs: 60_000, + executorFactory: () => ({ + async run(input) { + const index = prompts.push(input.prompt) - 1; + const done = new Promise((resolve) => { + finish[index] = resolve; + }); + const search = input.tools.find((tool) => tool.name === "search_evidence"); + await search?.execute("drain", { agentId: "default", limit: 1 }, undefined, undefined, {} as never); + await done; + return { summary: "Read one page" }; + }, + }), + }, + ); + const pending = (prompt: string | undefined) => + (prompt ?? "").slice( + (prompt ?? "").lastIndexOf(""), + (prompt ?? "").lastIndexOf(""), + ); + try { + await worker.triggerAsync("incremental"); + await waitFor(() => prompts.length === 2, 3_000); + finish[0]?.(); + await waitFor(() => worker.activePasses.length === 1, 3_000); + + await worker.triggerAsync("incremental"); + await waitFor(() => prompts.length === 3, 3_000); + + expect(pending(prompts[0])).toContain('"kind":"contested_claim"'); + expect(pending(prompts[1])).not.toContain("contested_claim"); + expect(pending(prompts[2])).toContain('"kind":"contested_claim"'); + } finally { + for (const open of finish) open?.(); + await waitFor(() => !worker.running, 5_000); + worker.stop(); + configureLlmConcurrency(previousLimit); + } + }); + + it("adds passes for a scope with more small sources than the backlog probe reads", async () => { + const seed = db.prepare( + `INSERT INTO session_transcripts + (session_key, agent_id, content, harness, created_at, updated_at, completed_at) + VALUES (?, 'default', ?, 'pi', datetime('now'), datetime('now'), datetime('now'))`, + ); + for (let index = 0; index < DREAMING_SCHEDULE_BACKLOG_MAX_SOURCES + 10; index++) { + seed.run(`small-${index}`, `User: small fact number ${index}.`); + } + const { delivered, peak } = await runSharedScopeWorker({ maxPassesPerScope: 3 }, 3); + + expect(peak).toBe(3); + expect(delivered.every((refs) => refs.length > 0)).toBe(true); + const all = delivered.flat(); + expect(new Set(all).size).toBe(all.length); + }); + + it("keeps one pass on a scope whose backlog is below the token threshold", async () => { + seedLargeScope(2); + const { delivered, peak } = await runSharedScopeWorker({ maxPassesPerScope: 3 }, 1); + + expect(peak).toBe(1); + expect(delivered).toHaveLength(1); + }); + it("runs disjoint agent groups concurrently after the first pass reaches a tool", async () => { const seed = db.prepare( `INSERT INTO session_transcripts diff --git a/platform/daemon/src/pipeline/dreaming-worker.ts b/platform/daemon/src/pipeline/dreaming-worker.ts index 37e5ec2978..e797402d78 100644 --- a/platform/daemon/src/pipeline/dreaming-worker.ts +++ b/platform/daemon/src/pipeline/dreaming-worker.ts @@ -72,6 +72,8 @@ interface RunningDreamingPass { readonly mode: DreamingMode; readonly scopes: readonly string[]; readonly exclusive: boolean; + readonly shared: boolean; + readonly attentionScopes: readonly string[]; readonly slots: number; readonly settled: Promise; } @@ -83,14 +85,16 @@ interface StartedDreamingPass { type DreamingPassResult = { passId: string; applied: number; skipped: number; failed: number; summary: string }; export function partitionDreamingScopes( - backlogs: ReadonlyArray<{ + work: ReadonlyArray<{ readonly scope: string; readonly tokens: number; readonly attention?: boolean; readonly oldestAttentionAt?: string | null; + readonly passes?: number; }>, slots: number, ): string[][] { + const backlogs = work.filter((item) => (item.passes ?? 1) > 0); const withBacklog = backlogs .filter((item) => item.tokens > 0) .sort((a, b) => b.tokens - a.tokens || a.scope.localeCompare(b.scope)); @@ -111,6 +115,15 @@ export function partitionDreamingScopes( target.tokens += item.tokens; } for (const scope of attentionOnly.slice(0, free - groups.length)) groups.push({ scopes: [scope], tokens: 0 }); + const extra = new Map(withBacklog.map((item) => [item.scope, Math.max(0, Math.floor(item.passes ?? 1) - 1)])); + while (groups.length < free && [...extra.values()].some((count) => count > 0)) { + for (const item of withBacklog) { + const count = extra.get(item.scope) ?? 0; + if (count <= 0 || groups.length >= free) continue; + extra.set(item.scope, count - 1); + groups.push({ scopes: [item.scope], tokens: item.tokens }); + } + } return groups.map((group) => [...group.scopes].sort()).filter((scopes) => scopes.length > 0); } export interface DreamingSchedulerStatus { @@ -343,6 +356,9 @@ export function startDreamingWorker( const runningPasses = new Set(); const configuredConcurrentPasses = Math.max(1, Math.floor(cfg.maxConcurrentPasses ?? 1)); const maxPasses = (): number => Math.max(1, Math.min(configuredConcurrentPasses, getLlmConcurrencyLimit())); + const passesPerScope = Math.max(1, Math.floor(cfg.maxPassesPerScope ?? 1)); + const sharesScopes = passesPerScope > 1; + const evidenceLeaseMs = 2 * cfg.timeout; let scheduler: DreamingSchedulerStatus = { status: "idle", reason: null, checkedAt: null }; let nextScheduledFocus: DreamingPassFocus | null = null; const getAgentScopes = createAgentScopeSnapshot(AGENT_SCOPE_SNAPSHOT_REFRESH_MS, () => @@ -446,6 +462,19 @@ export function startDreamingWorker( const usedSlots = (): number => [...runningPasses].reduce((sum, pass) => sum + pass.slots, 0); const leasedScopes = (): ReadonlySet => new Set([...runningPasses].flatMap((pass) => pass.scopes)); + const scopePassCounts = (): ReadonlyMap => { + const counts = new Map(); + for (const scope of [...runningPasses].flatMap((pass) => pass.scopes)) + counts.set(scope, (counts.get(scope) ?? 0) + 1); + return counts; + }; + const closedScopes = (): ReadonlySet => { + if (!sharesScopes) return leasedScopes(); + const counts = scopePassCounts(); + const closed = new Set([...runningPasses].filter((pass) => !pass.shared).flatMap((pass) => pass.scopes)); + for (const [scope, count] of counts) if (count >= passesPerScope) closed.add(scope); + return closed; + }; const exclusiveRunning = (): boolean => [...runningPasses].some((pass) => pass.exclusive); const listScopes = async (): Promise => { knownScopes = await getDreamingWorkerAgentIds(accessor, defaultAgentId, options.ownerMaintenance); @@ -454,8 +483,8 @@ export function startDreamingWorker( const knownScopesLeased = (): boolean => { if (runningPasses.size === 0) return false; if (exclusiveRunning() || usedSlots() >= maxPasses()) return true; - const leased = leasedScopes(); - return knownScopes.length > 0 && knownScopes.every((scope) => leased.has(scope)); + const closed = closedScopes(); + return knownScopes.length > 0 && knownScopes.every((scope) => closed.has(scope)); }; async function admit(fn: () => Promise): Promise { @@ -500,12 +529,20 @@ export function startDreamingWorker( resolveToolCall = resolve; }); let release: () => void = () => undefined; + const shared = sharesScopes && !exclusive && live?.userRequest === undefined; + const attentionHeld = new Set([...runningPasses].flatMap((pass) => pass.attentionScopes)); + const attentionScopes = shared ? scopes.filter((scope) => !attentionHeld.has(scope)) : scopes; + const passLive: DreamingPassLiveOptions | undefined = shared + ? { ...live, sharedScope: { evidenceLeaseMs, attentionScopes } } + : live; const entry: RunningDreamingPass = { passId: null, agentId: runAgentId, mode, scopes, exclusive, + shared, + attentionScopes, slots: 1, settled: new Promise((resolve) => { release = resolve; @@ -530,7 +567,7 @@ export function startDreamingWorker( mode, id, caps, - live, + passLive, options.ownerMaintenance, ).catch((error: unknown) => { recordDreamingFailureOrLog(runAgentId); @@ -564,8 +601,10 @@ export function startDreamingWorker( readonly tokens: number; readonly attention: boolean; readonly oldestAttentionAt: string | null; + readonly passes: number; }> > { + const running = scopePassCounts(); const owner = await getDbOwnerForAccessor(accessor); const attention = new Map( ( @@ -585,7 +624,16 @@ export function startDreamingWorker( const probe = await probeDreamingEpisodicBacklog(accessor, scope, cfg.tokenThreshold, options.ownerMaintenance); const tokens = probe.hasBacklog === false ? 0 : Math.max(1, probe.kind === "exact" ? probe.tokens : probe.tokenLowerBound); - return { scope, tokens, attention: attention.has(scope), oldestAttentionAt: attention.get(scope) ?? null }; + const held = running.get(scope) ?? 0; + const deep = tokens >= cfg.tokenThreshold || (probe.kind === "indeterminate" && probe.hasBacklog === true); + const wanted = sharesScopes && deep ? passesPerScope : 1; + return { + scope, + tokens, + attention: attention.has(scope), + oldestAttentionAt: attention.get(scope) ?? null, + passes: Math.max(0, wanted - held), + }; }), ); } @@ -594,8 +642,13 @@ export function startDreamingWorker( const slots = maxPasses() - usedSlots(); if (scopes.length === 0 || slots <= 0) throw new AlreadyRunningError(); const measured = - scopes.length === 1 ? [[...scopes]] : partitionDreamingScopes(await measureScopeWork(scopes), slots); - const groups = measured.length > 0 ? measured : [[...scopes]]; + scopes.length === 1 && !sharesScopes + ? [[...scopes]] + : partitionDreamingScopes(await measureScopeWork(scopes), slots); + const held = leasedScopes(); + const idle = scopes.filter((scope) => !held.has(scope)); + if (measured.length === 0 && idle.length === 0) throw new AlreadyRunningError(); + const groups = measured.length > 0 ? measured : [idle]; const [firstGroup, ...rest] = groups; const first = startPass(runAgentId, "incremental", firstGroup ?? [...scopes], false); if (rest.length === 0) return first; @@ -606,6 +659,8 @@ export function startDreamingWorker( mode: "incremental", scopes: rest.flat(), exclusive: false, + shared: sharesScopes, + attentionScopes: [], slots: rest.length, settled: new Promise((resolve) => { releaseReservation = resolve; @@ -671,8 +726,8 @@ export function startDreamingWorker( return; } scheduler = { status: "idle", reason: null, checkedAt }; - const leased = leasedScopes(); - const scopes = (await getAgentScopes()).filter((scope) => !leased.has(scope)); + const closed = closedScopes(); + const scopes = (await getAgentScopes()).filter((scope) => !closed.has(scope)); if (scopes.length === 0) return; const autoRequeued = await autoRequeueRepairedDreamingEvidence(accessor, evidenceRetry); if (autoRequeued > 0) { @@ -802,10 +857,10 @@ export function startDreamingWorker( if (runningPasses.size > 0) throw new AlreadyRunningError(); return startPass(runAgentId, mode, scopes, true); } - const leased = leasedScopes(); + const closed = closedScopes(); return await startIncrementalPasses( runAgentId, - scopes.filter((scope) => !leased.has(scope)), + scopes.filter((scope) => !closed.has(scope)), ); }); return await started.passId; diff --git a/platform/daemon/src/pipeline/dreaming.ts b/platform/daemon/src/pipeline/dreaming.ts index f81e5898de..d521f47433 100644 --- a/platform/daemon/src/pipeline/dreaming.ts +++ b/platform/daemon/src/pipeline/dreaming.ts @@ -70,6 +70,7 @@ import { SEEN_IMPORTED_SOURCE_ATTENTION_SQL, STALLED_EVIDENCE_ATTENTION_SQL, } from "./dreaming-evidence-consumption"; +import { evidenceLeasedByOtherPassesInScopes, releaseDreamingEvidenceLeasesInTx } from "./dreaming-evidence-leases"; import { parseDreamingReviewedExcludedEvidence, recordDreamingReviewedExcludedEvidenceInTx, @@ -489,13 +490,24 @@ function resetDreamingTokens( `UPDATE dreaming_state SET consecutive_failures = 0, last_failure_at = NULL, - last_pass_at = ?, + last_pass_at = CASE + WHEN last_pass_at IS NOT NULL AND (? IS NULL OR julianday(?) < julianday(last_pass_at)) THEN last_pass_at + ELSE ? + END, evidence_cursor = ?, last_pass_id = ?, last_pass_mode = ?, updated_at = datetime('now') WHERE agent_id = ?`, - ).run(lastPassAt, evidenceCursor === null ? null : JSON.stringify(evidenceCursor), passId, mode, agentId); + ).run( + lastPassAt, + lastPassAt, + lastPassAt, + evidenceCursor === null ? null : JSON.stringify(evidenceCursor), + passId, + mode, + agentId, + ); } else { db.prepare( `INSERT INTO dreaming_state @@ -1262,6 +1274,7 @@ export function selectDreamingPassMode( export interface DreamingPassLiveOptions { readonly hub?: DreamingLiveEventHub; + readonly sharedScope?: { readonly evidenceLeaseMs: number; readonly attentionScopes: readonly string[] }; readonly userRequest?: { readonly sourceRef: string; readonly content: string }; readonly memoryHeadReader?: (agentId: string, passId?: string) => Promise>; } @@ -1715,9 +1728,10 @@ ${JSON.stringify(liveOptions.userRequest)} { deadlineMs: 30_000, estimatedWorkUnits: 1 }, ); const cutoff = cutoffRow?.now ?? new Date().toISOString(); + const attentionScopes = liveOptions?.sharedScope?.attentionScopes ?? scopes; const [hasPendingHygieneAttention, hasPendingContentAttention, hasPendingAttention] = await Promise.all([ Promise.all( - scopes.map( + attentionScopes.map( async (scope) => (await ownerQueryOne<{ present: number }>( await getDbOwnerForAccessor(accessor), @@ -1729,7 +1743,7 @@ ${JSON.stringify(liveOptions.userRequest)} ), ).then((values) => values.some(Boolean)), Promise.all( - scopes.map( + attentionScopes.map( async (scope) => (await ownerQueryOne<{ present: number }>( await getDbOwnerForAccessor(accessor), @@ -1741,7 +1755,7 @@ ${JSON.stringify(liveOptions.userRequest)} ), ).then((values) => values.some(Boolean)), Promise.all( - scopes.map( + attentionScopes.map( async (scope) => (await ownerQueryOne<{ present: number }>( await getDbOwnerForAccessor(accessor), @@ -1862,6 +1876,7 @@ ${JSON.stringify(liveOptions.userRequest)} }), evidenceDeliveryDeadline: Date.now() + Math.floor(cfg.timeout / 2), evidenceChars: dreamingEvidencePageChars(cfg.maxInputTokens), + ...(liveOptions?.sharedScope ? { evidenceLeaseMs: liveOptions.sharedScope.evidenceLeaseMs } : {}), accessor, agentId, memoryHeadCommitter, @@ -1952,7 +1967,11 @@ ${JSON.stringify(liveOptions.userRequest)} const passPrompt = dreamingPassPrompt( prompt, await renderDreamingHistoryForPass(accessor, agentId, historyScopes), - await renderPendingDreamingAttention(accessor, historyScopes, mode), + await renderPendingDreamingAttention( + accessor, + historyScopes.filter((scope) => attentionScopes.includes(scope)), + mode, + ), ).concat("\n\n", dreamingPassClock(new Date(), detectLocalTimeZone())); logger.info("dreaming", "Starting agentic dreaming pass", { mode, @@ -2301,7 +2320,9 @@ export function finalizeDreamingPassInDb(db: WriteDb, input: DbOwnerDreamingPass } const runbookDeferred = parsedRunbook === null ? null : deferredEvidenceKeys(parsedRunbook, input.agentId); const failedEvidence = failedOperationEvidence(db, input.passId, input.agentId, input.scopes); - const deferredEvidence = runbookDeferred === null ? null : new Set([...runbookDeferred, ...failedEvidence.sources]); + const heldElsewhere = evidenceLeasedByOtherPassesInScopes(db, input.passId, input.scopes); + const deferredEvidence = + runbookDeferred === null ? null : new Set([...runbookDeferred, ...failedEvidence.sources, ...heldElsewhere]); const reviewedExcludedEvidence = parsedRunbook === null ? null : parseDreamingReviewedExcludedEvidence(parsedRunbook); if (deferredEvidence !== null && reviewedExcludedEvidence !== null) { @@ -2321,6 +2342,7 @@ export function finalizeDreamingPassInDb(db: WriteDb, input: DbOwnerDreamingPass } resolveImportedSourceAttentionInTx(db, input.passId, input.scopes); } + releaseDreamingEvidenceLeasesInTx(db, input.passId); if ( dreamingModeAdvancesEvidence( input.mode as DreamingMode, diff --git a/platform/daemon/src/routes/pipeline-routes.ts b/platform/daemon/src/routes/pipeline-routes.ts index 82ec2e62ea..d37298c0df 100644 --- a/platform/daemon/src/routes/pipeline-routes.ts +++ b/platform/daemon/src/routes/pipeline-routes.ts @@ -726,6 +726,7 @@ export function registerPipelineRoutes(app: Hono): void { maxInputTokens: cfg.dreaming.maxInputTokens, maxOutputTokens: cfg.dreaming.maxOutputTokens, maxConcurrentPasses: cfg.dreaming.maxConcurrentPasses, + maxPassesPerScope: cfg.dreaming.maxPassesPerScope ?? 1, codemode: cfg.dreaming.codemode, timeout: cfg.dreaming.timeout, surprisal: cfg.dreaming.surprisal, @@ -892,7 +893,7 @@ export function registerPipelineRoutes(app: Hono): void { async (maintenance) => (await ownerQueryOne<{ present: number }>( maintenance.owner, - "routes/pipeline-routes.ts:895", + "routes/pipeline-routes.ts:896", "SELECT 1 AS present FROM dreaming_evidence_exclusions WHERE agent_id = ? AND source_kind = 'summary' AND source_id = ? AND resolved_at IS NULL", [agentId, sourceId], )) != null, diff --git a/scripts/event-loop-contract-baseline.json b/scripts/event-loop-contract-baseline.json index 010d61e85e..c3fcea2e4a 100644 --- a/scripts/event-loop-contract-baseline.json +++ b/scripts/event-loop-contract-baseline.json @@ -1655,14 +1655,14 @@ }, { "path": "memory-config.ts", - "line": 359, + "line": 360, "api": "existsSync", "source": ".find((path) => existsSync(path));", "category": "hot-path" }, { "path": "memory-config.ts", - "line": 361, + "line": 362, "api": "readFileSync", "source": "const yaml = path ? parseRuntimeYaml(readFileSync(path, \"utf8\")) : {};", "category": "hot-path" @@ -2631,7 +2631,7 @@ }, { "path": "pipeline/dreaming-worker.ts", - "line": 185, + "line": 198, "api": "withReadDbAsync", "source": "return await accessor.withReadDbAsync((db) => getQueueHealth(db).status !== \"healthy\", {", "category": "hot-path", @@ -2639,7 +2639,7 @@ }, { "path": "pipeline/dreaming-worker.ts", - "line": 235, + "line": 248, "api": "withReadDbAsync", "source": ": await accessor.withReadDbAsync(", "category": "hot-path", @@ -2647,7 +2647,7 @@ }, { "path": "pipeline/dreaming-worker.ts", - "line": 293, + "line": 306, "api": "withReadDbAsync", "source": ": accessor.withReadDbAsync((db) => hasDreamingAttentionKindInDb(db, scope, [\"hygiene\"]), {", "category": "hot-path", @@ -2655,7 +2655,7 @@ }, { "path": "pipeline/dreaming-worker.ts", - "line": 313, + "line": 326, "api": "withReadDbAsync", "source": ": accessor.withReadDbAsync(", "category": "hot-path", @@ -2663,14 +2663,14 @@ }, { "path": "pipeline/dreaming.ts", - "line": 922, + "line": 934, "api": "existsSync", "source": "if (!existsSync(path)) return { content: \"\", unreadable: false };", "category": "hot-path" }, { "path": "pipeline/dreaming.ts", - "line": 924, + "line": 936, "api": "readFileSync", "source": "const raw = readFileSync(path, \"utf-8\").trim();", "category": "hot-path" diff --git a/web/docs/src/content/docs/api/knowledge-ontology.md b/web/docs/src/content/docs/api/knowledge-ontology.md index ee9c51e94d..89c524850a 100644 --- a/web/docs/src/content/docs/api/knowledge-ontology.md +++ b/web/docs/src/content/docs/api/knowledge-ontology.md @@ -520,6 +520,7 @@ an exact count when at most 50 sources are waiting, `null` when more are. "maxInputTokens": 128000, "maxOutputTokens": null, "maxConcurrentPasses": 2, + "maxPassesPerScope": 1, "codemode": false, "timeout": 300000 }, @@ -770,7 +771,7 @@ When a claim arrives for a slot whose active claim has a later evidence time (`validFrom`, else `occurredAt`), the incoming claim is recorded as already superseded by the active one and the result names it in `supersededByNewerEvidence`. When the evidence times are equal (two updates on -the same day) or neither claim has one, the claim whose cited source was +the same day) or either claim lacks one, the claim whose cited source was captured later stays current, so the order Dreaming happens to file sources in does not decide the current value. A claim whose source has no capture time falls back to the newest write. @@ -923,8 +924,23 @@ groups (default 2), balanced by evidence backlog, and runs one pass per group. It never starts more passes than `worker.maxLlmConcurrency` allows, because a pass waiting for a shared LLM permit would spend its own timeout waiting. The daemon allows that many Pi agent workers plus three for retained dashboard chats. -Each pass may read and write only its own group's agents, and an agent is in at -most one running pass. The first pass starts immediately; the other groups start +Each pass may read and write only its own group's agents, and by default an agent +is in at most one running pass. `memory.dreaming.maxPassesPerScope` (default 1, +up to 16) lets one agent take several incremental passes at once: when slots are +left after every agent with work has a pass, an agent whose backlog reaches +`tokenThreshold`, or holds more pending sources than the 50 the backlog probe +reads, gets up to that many passes, still within `maxConcurrentPasses`. These +passes lease the evidence they draw from the delivery queue, so no source is +handed to two running passes. A pass that loses a source to another pass reads +the next unclaimed sources instead; when the other passes hold everything left, +its queue page is empty with `hasMore: false` and `heldByOtherPasses: true`. A lease ends when +its pass finishes, fails, or is cancelled, or after twice the pass timeout, and +a source whose pass did not finish it is delivered again. One running pass on +an agent works its pending attention and the passes that join it only read +evidence; the next pass to start after it finishes takes the attention over. A claim +written from older evidence after a newer contradicting claim lands as +superseded, so the current value follows the evidence rather than which pass +wrote last. The first pass starts immediately; the other groups start only after it completes a tool call, so an unavailable provider is not called once per group. The response's `passId` is the first pass. `worker.activePasses` in `GET /api/dream/status` lists every running pass; the trigger is complete when it diff --git a/web/docs/src/content/docs/architecture/pipeline-storage.md b/web/docs/src/content/docs/architecture/pipeline-storage.md index 3f47e449b5..4d0a713bfd 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 `169-dreaming-evidence-review-cursor.ts`. +`schema_migrations_audit`. The latest migration is `170-dreaming-evidence-leases.ts`. ### Evidence and semantic state