Skip to content
Open
10 changes: 5 additions & 5 deletions docs/event-loop-contract-audit.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
@@ -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)",
);
}
}
12 changes: 12 additions & 0 deletions platform/core/src/migrations/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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[] = [
Expand Down Expand Up @@ -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;
Expand Down
26 changes: 25 additions & 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(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" });
Expand Down Expand Up @@ -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);
Expand Down
27 changes: 6 additions & 21 deletions platform/daemon/src/daemon.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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",
});
},
});
}
Expand Down
11 changes: 11 additions & 0 deletions platform/daemon/src/db-owner-protocol.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 }
Expand Down Expand Up @@ -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;
Expand Down
1 change: 1 addition & 0 deletions platform/daemon/src/db-owner-runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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":
Expand Down
8 changes: 8 additions & 0 deletions platform/daemon/src/db-owner-worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -560,6 +560,13 @@ export function runDbOwnerWorker(): void {
return readDreamingEvidenceSourceInDb(db as never, request.input);
}

async function executeDreamingEvidenceReview(
request: Extract<DbOwnerJob["request"], { readonly kind: "dreaming_evidence_review" }>,
): Promise<unknown> {
const { reviewDreamingEvidenceInDb } = await import("./pipeline/dreaming-evidence-consumption");
return reviewDreamingEvidenceInDb(db as never, request.input);
}

async function executeDreamingPassFinalize(
request: Extract<DbOwnerJob["request"], { readonly kind: "dreaming_pass_finalize" }>,
context: JobExecutionContext,
Expand Down Expand Up @@ -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);
Expand Down
15 changes: 13 additions & 2 deletions platform/daemon/src/episodic-sources.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 =
Expand Down Expand Up @@ -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],
Expand Down
20 changes: 13 additions & 7 deletions platform/daemon/src/pipeline/dreaming-attention.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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",
);
}

Expand All @@ -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,
Expand All @@ -172,11 +174,13 @@ export async function getDreamingAttentionScoped(
readonly kind?: string;
readonly status?: "pending" | "resolved";
readonly limit?: number;
readonly agentFacing?: boolean;
},
): Promise<readonly DreamingAttention[]> {
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<{
Expand All @@ -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],
Expand All @@ -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);
Expand All @@ -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 ?`,
)
Expand All @@ -240,7 +246,7 @@ export function getDreamingAttentionAcrossScopes(
...attention,
details: parseDetails(detailsJson),
}));
}, "pipeline/dreaming-attention.ts:220");
}, "pipeline/dreaming-attention.ts:226");
}

export function getDreamingAttentionById(
Expand Down Expand Up @@ -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(
Expand All @@ -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",
);
}

Expand Down
Loading
Loading