Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 17 additions & 0 deletions platform/core/src/migrations/170-dreaming-evidence-leases.ts
Original file line number Diff line number Diff line change
@@ -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);
`);
}
7 changes: 7 additions & 0 deletions platform/core/src/migrations/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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[] = [
Expand Down Expand Up @@ -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;
Expand Down
2 changes: 1 addition & 1 deletion platform/core/src/migrations/migrations.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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" });
Expand Down
1 change: 1 addition & 0 deletions platform/core/src/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
4 changes: 4 additions & 0 deletions platform/daemon/src/memory-config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -985,6 +986,9 @@ export function loadDreamingConfig(yaml: Record<string, unknown>): 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: {
Expand Down
42 changes: 42 additions & 0 deletions platform/daemon/src/ontology-proposals.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, unknown>) =>
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(), {
Expand Down
2 changes: 1 addition & 1 deletion platform/daemon/src/ontology-proposals.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
})
Expand Down
80 changes: 58 additions & 22 deletions platform/daemon/src/pipeline/dreaming-capabilities.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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?: (
Expand Down Expand Up @@ -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;
Expand All @@ -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)];
Expand Down Expand Up @@ -365,7 +376,7 @@ function drainDreamingEvidenceQueueInDb(db: ReadDb, input: DbOwnerDreamingEviden
}
return items;
};
const continuationRefs = new Set<string>();
const continuationRefs = new Set<string>(leasedElsewhere);
const continuations = page(
pendingDreamingEvidenceContinuations(db, scopeId, 50, input.kind),
limit + 1,
Expand Down Expand Up @@ -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 : {}),
};
}

Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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<DreamingCapabilityOutput>(
{
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<DreamingCapabilityOutput>(
{
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<Record<string, unknown>>;
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(
Expand Down
Loading
Loading