diff --git a/App/backend/local-api-contracts/src/memory-runtime.ts b/App/backend/local-api-contracts/src/memory-runtime.ts index e7fd1b334..612c29b45 100644 --- a/App/backend/local-api-contracts/src/memory-runtime.ts +++ b/App/backend/local-api-contracts/src/memory-runtime.ts @@ -59,6 +59,7 @@ export const JobTypeSchema = z.enum([ "skill_trial_resolve", "decision_repair", "work_memory_extract", + "work_memory_idle_flush", "feedback_experience" ]); export type JobType = z.infer; diff --git a/Memory/src/contracts/memory-runtime.ts b/Memory/src/contracts/memory-runtime.ts index 3367f0d77..25cea8112 100644 --- a/Memory/src/contracts/memory-runtime.ts +++ b/Memory/src/contracts/memory-runtime.ts @@ -61,6 +61,7 @@ export const JobTypeSchema = z.enum([ "skill_trial_resolve", "decision_repair", "work_memory_extract", + "work_memory_idle_flush", "feedback_experience" ]); export type JobType = z.infer; diff --git a/Memory/src/logging/logger.ts b/Memory/src/logging/logger.ts index 09b779920..3a4ef0728 100644 --- a/Memory/src/logging/logger.ts +++ b/Memory/src/logging/logger.ts @@ -318,7 +318,8 @@ function jobTypeTag(value: string): string { episode_idle_close: "episode.close", skill_trial_resolve: "skill.trial_resolve", decision_repair: "decision.repair", - l2_association: "l2.association" + l2_association: "l2.association", + work_memory_idle_flush: "work.memory.idle_flush" }; return tags[value] ?? value.replace(/_/g, "."); } diff --git a/Memory/src/service/memory-service.ts b/Memory/src/service/memory-service.ts index 184d74437..2312838e6 100644 --- a/Memory/src/service/memory-service.ts +++ b/Memory/src/service/memory-service.ts @@ -383,7 +383,10 @@ export class MemoryService { embedUserMemory: (job) => this.embeddingJobs.embedUserMemory(job) }, workMemory: { - extract: (job) => this.workMemory.extract(job) + extract: (job) => this.workMemory.extract(job), + flushIdle: (job) => { + this.workMemory.flushIdle(job); + } }, episodeTitle: { generate: (job) => this.episodeTitle.generate(job) @@ -661,6 +664,8 @@ export class MemoryService { shouldDeferBudgetedEvolutionLlm: () => this.shouldDeferBudgetedEvolutionLlm(), firstLine, memoryLayersForIntent, + armWorkMemoryIdleFlush: this.armWorkMemoryIdleFlush.bind(this), + extractUnextractedWorkMemory: this.extractUnextractedWorkMemory.bind(this), namespaceIdFromContext, namespaceIdFromMemory, namespaceIdFromSession, @@ -1183,6 +1188,16 @@ export class MemoryService { return this.sessionTurns.closeSession(sessionId, this.withTimeZone(request)); } + /** Arm the Work Memory idle flush for a Session inside the caller's transaction. */ + private armWorkMemoryIdleFlush(sessionId: string, at: string): void { + this.workMemory.armIdleFlush(sessionId, at); + } + + /** Extract unextracted Work Memory for a Session inside the caller's transaction. */ + private extractUnextractedWorkMemory(sessionId: string, throughTraceSeq: number, at: string): void { + this.workMemory.extractUnextracted(sessionId, throughTraceSeq, at); + } + l3WorldModelTraceHead( sessionId: string, request: L3WorldModelRequestEnvelope @@ -1208,8 +1223,8 @@ export class MemoryService { trigger: request.trigger, throughL1MemoryId: request.throughL1MemoryId }, (frozen) => { - if (request.trigger === "token_compaction" && frozen.batchIds.length > 0) { - this.workMemory.scheduleBatchesInTransaction(frozen.batchIds, nowIso()); + if (request.trigger === "token_compaction" && frozen.throughTraceSeq) { + this.workMemory.extractUnextracted(sessionId, frozen.throughTraceSeq, nowIso()); } }); if (!result.throughTraceSeq) { diff --git a/Memory/src/service/session/session-turn-service.ts b/Memory/src/service/session/session-turn-service.ts index eb2e46138..6c74f8d2a 100644 --- a/Memory/src/service/session/session-turn-service.ts +++ b/Memory/src/service/session/session-turn-service.ts @@ -815,6 +815,13 @@ export class SessionTurnService { }); this.deps.finalizeClosedEpisode(episode, at, "session_closed"); } + if (session.meta.l3_world_model_protocol_version === 2) { + this.deps.extractUnextractedWorkMemory( + session.id, + this.deps.repos.l3WorldModels.maxInputTraceSeq(session.id), + at + ); + } this.deps.repos.runtime.appendChange({ memoryId: session.id, namespaceId: this.deps.namespaceIdFromSession(closedWithMeta), @@ -895,6 +902,11 @@ export class SessionTurnService { trigger: "session_close", at }); + this.deps.extractUnextractedWorkMemory( + sessionId, + this.deps.repos.l3WorldModels.maxInputTraceSeq(sessionId), + at + ); } const changeSeq = this.deps.repos.runtime.appendChange({ memoryId: sessionId, @@ -2255,6 +2267,13 @@ export class SessionTurnService { })); } } + if ( + rawTurnFirstCompleted && + !completedEndTopicDecision && + session.meta.l3_world_model_protocol_version === 2 + ) { + this.deps.armWorkMemoryIdleFlush(session.id, at); + } const uniqueClosedEpisodeIds = uniq(closedEpisodeIds); const responseChangeSeq = this.deps.repos.runtime.latestChangeSeq(session.userId, this.deps.namespaceIdFromSession(session)); const body: CompleteTurnResponse = { diff --git a/Memory/src/service/work-memory/work-memory-pipeline.ts b/Memory/src/service/work-memory/work-memory-pipeline.ts index bf9167cc2..da59a64aa 100644 --- a/Memory/src/service/work-memory/work-memory-pipeline.ts +++ b/Memory/src/service/work-memory/work-memory-pipeline.ts @@ -4,9 +4,10 @@ import { type JsonValue } from "../../contracts/index.js"; import type { Embedder, LlmClient } from "../../model/types.js"; +import { splitL3TracesByRawTurn } from "../../storage/repositories.js"; import type { EvolutionJobRecord, - L3WorldModelEvidenceBatchRecord, + L3WorldModelInputTraceRecord, Repositories } from "../../storage/repositories.js"; import type { MemoryFilter, MemoryRow } from "../../types.js"; @@ -18,6 +19,22 @@ export interface WorkMemoryQaPair { assistant: string; } +/** Result of one incremental Work Memory extraction pass. */ +export interface WorkMemoryExtractionResult { + windows: number; + enqueued: number; + cursorAdvancedTo: number; +} + +/** Raw turns per extraction window; a single turn is never split across windows. */ +const WORK_MEMORY_WINDOW_MAX_RAW_TURNS = 20; + +/** Q&A pairs per extraction window; an oversized turn overflows into the next window. */ +const WORK_MEMORY_WINDOW_MAX_QA_PAIRS = 20; + +/** Idle delay before an untouched Session flushes its unextracted Work Memory. */ +export const WORK_MEMORY_IDLE_TIMEOUT_MS = 2 * 60 * 60 * 1000; + export interface WorkMemoryCandidate { requirement: string; reason: string; @@ -75,6 +92,147 @@ export class WorkMemoryPipeline { return this.deps.repos.transaction(() => this.scheduleBatchesInTransaction(batchIds, at)); } + /** + * Enqueue `work_memory_extract` jobs for every L3 input trace of a Session + * that has not been scheduled yet, up to `throughTraceSeq`. + * + * Callers must already hold a transaction: the cursor advance and the job + * enqueue have to commit together, otherwise a failed enqueue would mark the + * window extracted and no later trigger would ever pick it up again. + * + * @param sessionId Session whose unextracted traces should be scheduled. + * @param throughTraceSeq Inclusive trace sequence ceiling for this trigger. + * @param at Timestamp to stamp cursor and job rows with. + * @returns Window, enqueue, and cursor counters for the caller. + */ + extractUnextracted( + sessionId: string, + throughTraceSeq: number, + at = this.deps.nowIso() + ): WorkMemoryExtractionResult { + const session = this.deps.repos.runtime.getSession(sessionId); + if (!session || session.meta.l3_world_model_protocol_version !== 2) { + return { windows: 0, enqueued: 0, cursorAdvancedTo: 0 }; + } + // The first Work Memory pass of a Session starts from a clean slate: the L3 + // cursor has already advanced past this compaction window by the time the + // boundary callback runs, so inheriting it would swallow the whole delta. + const cursor = this.deps.repos.runtime.ensureWorkMemoryCursor(sessionId, 0, at); + const endTraceSeq = Math.min(throughTraceSeq, this.deps.repos.l3WorldModels.maxInputTraceSeq(sessionId)); + if (endTraceSeq <= cursor.lastExtractedSeq) { + return { windows: 0, enqueued: 0, cursorAdvancedTo: cursor.lastExtractedSeq }; + } + const traces = this.deps.repos.l3WorldModels.listInputTracesInRange( + sessionId, + cursor.lastExtractedSeq, + endTraceSeq + ); + if (traces.length === 0) { + return { windows: 0, enqueued: 0, cursorAdvancedTo: cursor.lastExtractedSeq }; + } + const scope = { userId: session.userId, sessionId }; + let enqueued = 0; + let windows = 0; + let lastExtractedSeq = cursor.lastExtractedSeq; + for (const chunk of splitL3TracesByRawTurn(traces, WORK_MEMORY_WINDOW_MAX_RAW_TURNS)) { + const window = this.takeWindow(chunk, scope); + if (window.qa.length > 0) { + const job = this.scheduleQaInTransaction({ + qa: window.qa, + userId: session.userId, + sessionId, + projectId: session.projectId ?? null + }, at); + if (job) enqueued += 1; + } + this.deps.repos.runtime.setWorkMemoryCursor(sessionId, window.throughTraceSeq, at); + lastExtractedSeq = window.throughTraceSeq; + windows += 1; + if (window.truncated) break; + } + return { windows, enqueued, cursorAdvancedTo: lastExtractedSeq }; + } + + /** + * Build the Q&A window for one chunk of traces, capped at + * `WORK_MEMORY_WINDOW_MAX_QA_PAIRS` pairs without splitting a Raw turn. + * + * The cursor stops on the last turn that fits, so an oversized turn leaves + * its remainder to the next trigger instead of losing it. + */ + private takeWindow( + chunk: readonly L3WorldModelInputTraceRecord[], + scope: { userId: string; sessionId: string } + ): { qa: WorkMemoryQaPair[]; throughTraceSeq: number; truncated: boolean } { + const rawTurnIds = [...new Set(chunk.map((trace) => trace.rawTurnId))]; + const lastTraceByRawTurnId = new Map(); + for (const trace of chunk) { + lastTraceByRawTurnId.set(trace.rawTurnId, trace.traceSeq); + } + const qa: WorkMemoryQaPair[] = []; + let throughTraceSeq = chunk[0]!.traceSeq; + let truncated = false; + for (const rawTurnId of rawTurnIds) { + const pair = this.qaForRawTurn(rawTurnId, scope); + if (pair && qa.length >= WORK_MEMORY_WINDOW_MAX_QA_PAIRS) { + truncated = true; + break; + } + if (pair) qa.push(pair); + throughTraceSeq = lastTraceByRawTurnId.get(rawTurnId)!; + } + return { qa, throughTraceSeq, truncated }; + } + + /** + * Arm (or re-arm) the idle flush for a Session so its unextracted Work Memory + * is scheduled once the Session has been quiet for `WORK_MEMORY_IDLE_TIMEOUT_MS`. + * + * Callers must already hold a transaction. + */ + armIdleFlush( + sessionId: string, + at = this.deps.nowIso() + ): EvolutionJobRecord | undefined { + const session = this.deps.repos.runtime.getSession(sessionId); + if (!session || session.meta.l3_world_model_protocol_version !== 2) return undefined; + return this.deps.repos.runtime.armWorkMemoryIdleFlush({ + sessionId, + userId: session.userId, + lastActivityAt: at, + runAfter: new Date(Date.parse(at) + WORK_MEMORY_IDLE_TIMEOUT_MS).toISOString(), + at + }); + } + + /** + * Handle a due `work_memory_idle_flush`: extract the delta when the Session + * really has been quiet for the full timeout, and re-arm otherwise. + */ + flushIdle(job: EvolutionJobRecord): WorkMemoryExtractionResult { + if (job.jobType !== "work_memory_idle_flush") { + throw new Error(`invalid work memory idle flush job type: ${job.id}`); + } + const sessionId = job.sessionId; + if (!sessionId) throw new Error(`work memory idle flush job has no session: ${job.id}`); + const empty: WorkMemoryExtractionResult = { windows: 0, enqueued: 0, cursorAdvancedTo: 0 }; + const session = this.deps.repos.runtime.getSession(sessionId); + if (!session || session.status !== "open" || session.meta.l3_world_model_protocol_version !== 2) { + return empty; + } + const at = this.deps.nowIso(); + const lastActivityAt = this.deps.repos.l3WorldModels.latestInputTraceCreatedAt(sessionId) + ?? (typeof job.payload.lastActivityAt === "string" ? job.payload.lastActivityAt : at); + if (Date.parse(at) < Date.parse(lastActivityAt) + WORK_MEMORY_IDLE_TIMEOUT_MS) { + this.deps.repos.transaction(() => { + this.armIdleFlush(sessionId, lastActivityAt); + }); + return empty; + } + const throughTraceSeq = this.deps.repos.l3WorldModels.maxInputTraceSeq(sessionId); + return this.deps.repos.transaction(() => this.extractUnextracted(sessionId, throughTraceSeq, at)); + } + async extract(job: EvolutionJobRecord): Promise { if (job.jobType !== "work_memory_extract") { throw new Error(`invalid work memory job type: ${job.id}`); @@ -138,13 +296,59 @@ export class WorkMemoryPipeline { if (!session || session.userId !== batch.userId || normalizeProjectId(session.projectId) !== normalizeProjectId(batch.projectId)) { throw new Error(`work memory batch session scope mismatch: ${batchId}`); } - const qa = this.qaForBatch(batch); + return this.scheduleQaInTransaction({ + qa: batch.rawTurnIds + .map((rawTurnId) => this.qaForRawTurn(rawTurnId, { + userId: batch.userId, + sessionId: batch.sessionId + })) + .filter((pair): pair is WorkMemoryQaPair => Boolean(pair)), + userId: batch.userId, + sessionId: batch.sessionId, + projectId: batch.projectId ?? null + }, at); + } + + /** + * Build the Q&A pair for one Raw turn, dropping tool traffic. Returns + * undefined for turns that were deleted, redacted, or carry no text. + */ + private qaForRawTurn( + rawTurnId: string, + scope: { userId: string; sessionId: string } + ): WorkMemoryQaPair | undefined { + const rawTurn = this.deps.repos.runtime.getRawTurn(rawTurnId); + if (!rawTurn || rawTurn.deletedAt || rawTurn.redactedAt) return undefined; + if (rawTurn.userId !== scope.userId || rawTurn.sessionId !== scope.sessionId) { + throw new Error(`work memory RawTurn scope mismatch: ${rawTurnId}`); + } + const user = normalizeQaText(rawTurn.userText ?? ""); + const assistant = normalizeQaText(rawTurn.assistantText ?? ""); + if (!user && !assistant) return undefined; + return { user, assistant }; + } + + /** Enqueue one `work_memory_extract` job for a Q&A window. */ + private scheduleQaInTransaction( + input: { + qa: readonly WorkMemoryQaPair[]; + userId: string; + sessionId: string; + projectId: string | null; + }, + at: string + ): EvolutionJobRecord | undefined { + const session = this.deps.repos.runtime.getSession(input.sessionId); + if (!session || session.userId !== input.userId) { + throw new Error(`work memory session scope mismatch: ${input.sessionId}`); + } + const qa = input.qa; const trajectoryHash = trajectoryHashForQa(qa); - const projectId = normalizeProjectId(batch.projectId); + const projectId = normalizeProjectId(input.projectId); const source = session.source.trim() || "unknown"; const dedupeKey = stableHash([ "work_memory_extract", - batch.userId, + input.userId, source, projectId, trajectoryHash @@ -153,15 +357,15 @@ export class WorkMemoryPipeline { if (existing?.status === "queued" || existing?.status === "leased" || existing?.status === "succeeded" || existing?.status === "dead_letter") { return existing; } - const scopeKey = stableHash(["work_memory", batch.userId, projectId]); + const scopeKey = stableHash(["work_memory", input.userId, projectId]); const scopeSeq = existing?.scopeSeq ?? this.deps.repos.runtime.nextWorkMemoryScopeSeq(scopeKey); const job = this.deps.repos.runtime.enqueueJobInTransaction({ id: existing?.id ?? newId("job"), jobType: "work_memory_extract", status: "queued", dedupeKey, - userId: batch.userId, - sessionId: batch.sessionId, + userId: input.userId, + sessionId: input.sessionId, scopeKey, scopeSeq, payload: { @@ -191,22 +395,6 @@ export class WorkMemoryPipeline { return job; } - private qaForBatch(batch: L3WorldModelEvidenceBatchRecord): WorkMemoryQaPair[] { - const result: WorkMemoryQaPair[] = []; - for (const rawTurnId of batch.rawTurnIds) { - const rawTurn = this.deps.repos.runtime.getRawTurn(rawTurnId); - if (!rawTurn || rawTurn.deletedAt || rawTurn.redactedAt) continue; - if (rawTurn.userId !== batch.userId || rawTurn.sessionId !== batch.sessionId) { - throw new Error(`work memory RawTurn scope mismatch: ${rawTurnId}`); - } - const user = normalizeQaText(rawTurn.userText ?? ""); - const assistant = normalizeQaText(rawTurn.assistantText ?? ""); - if (!user && !assistant) continue; - result.push({ user, assistant }); - } - return result; - } - private async retrieveCandidates( candidates: WorkMemoryCandidate[], filter: MemoryFilter diff --git a/Memory/src/service/worker/job-handlers.ts b/Memory/src/service/worker/job-handlers.ts index 0230bc456..8eb2d8bc2 100644 --- a/Memory/src/service/worker/job-handlers.ts +++ b/Memory/src/service/worker/job-handlers.ts @@ -85,6 +85,7 @@ export interface WorkerJobProcessors { }; workMemory: { extract(job: EvolutionJobRecord): MaybePromise; + flushIdle(job: EvolutionJobRecord): MaybePromise; }; episodeTitle: { generate(job: EvolutionJobRecord): MaybePromise; @@ -304,6 +305,9 @@ export async function processJob( case "work_memory_extract": await deps.processors.workMemory.extract(job); return; + case "work_memory_idle_flush": + await deps.processors.workMemory.flushIdle(job); + return; case "episode_title": await deps.processors.episodeTitle.generate(job); return; @@ -699,6 +703,10 @@ export function evolutionJobDedupeKey(input: Pick= 0), updated_at TIMESTAMPTZ NOT NULL )`, + `CREATE TABLE IF NOT EXISTS work_memory_session_cursors ( + session_id TEXT PRIMARY KEY REFERENCES sessions(id) ON DELETE CASCADE, + last_extracted_seq BIGINT NOT NULL DEFAULT 0 CHECK (last_extracted_seq >= 0), + updated_at TIMESTAMPTZ NOT NULL + )`, `CREATE TABLE IF NOT EXISTS l3_world_model_input_traces ( session_id TEXT NOT NULL REFERENCES sessions(id) ON DELETE CASCADE, trace_seq BIGINT NOT NULL CHECK (trace_seq >= 1), diff --git a/Memory/src/storage/repositories.ts b/Memory/src/storage/repositories.ts index eb64caef4..8b6f08eba 100644 --- a/Memory/src/storage/repositories.ts +++ b/Memory/src/storage/repositories.ts @@ -63,6 +63,7 @@ const BUNDLE_TABLES = [ "user_memories", "sessions", "l3_world_model_session_cursors", + "work_memory_session_cursors", "episodes", "raw_turns", "l3_world_model_input_traces", @@ -350,6 +351,12 @@ export interface L3WorldModelInputTraceRecord { createdAt: string; } +export interface WorkMemorySessionCursorRecord { + sessionId: string; + lastExtractedSeq: number; + updatedAt: string; +} + export interface L3WorldModelEvidenceBatchRecord { id: string; scopeKey: string; @@ -3252,6 +3259,97 @@ export class RuntimeRepository { return Number(row.next_seq); } + /** + * Read the Work Memory extraction cursor for a Session. + * + * A missing row is seeded from the L3 cursor so Sessions that predate this + * table do not re-extract windows the L3 boundary already covered. + */ + getWorkMemoryCursor(sessionId: string, at = nowIso()): WorkMemorySessionCursorRecord { + const seed = this.db.prepare( + `SELECT COALESCE( + (SELECT last_scheduled_seq FROM l3_world_model_session_cursors WHERE session_id = ?), + 0 + ) AS last_seq` + ).get(sessionId) as { last_seq: number }; + return this.ensureWorkMemoryCursor(sessionId, Number(seed.last_seq), at); + } + + /** Create the Work Memory cursor row for a Session when it is still missing. */ + ensureWorkMemoryCursor( + sessionId: string, + lastExtractedSeq: number, + at = nowIso() + ): WorkMemorySessionCursorRecord { + this.db.prepare( + `INSERT INTO work_memory_session_cursors (session_id, last_extracted_seq, updated_at) + VALUES (?, ?, ?) + ON CONFLICT(session_id) DO NOTHING` + ).run(sessionId, lastExtractedSeq, at); + const row = this.db.prepare( + `SELECT session_id, last_extracted_seq, updated_at + FROM work_memory_session_cursors WHERE session_id = ?` + ).get(sessionId) as { session_id: string; last_extracted_seq: number; updated_at: string }; + return { + sessionId: row.session_id, + lastExtractedSeq: Number(row.last_extracted_seq), + updatedAt: row.updated_at + }; + } + + /** Advance the Work Memory extraction cursor. */ + setWorkMemoryCursor(sessionId: string, lastExtractedSeq: number, at = nowIso()): void { + this.db.prepare( + `UPDATE work_memory_session_cursors + SET last_extracted_seq = ?, updated_at = ? + WHERE session_id = ?` + ).run(lastExtractedSeq, at, sessionId); + } + + /** + * Arm the Work Memory idle flush for a Session, pushing `runAfter` forward. + * + * The generic enqueue path merges `runAfter` by keeping the earlier value, + * which is the opposite of re-arming. This upsert therefore reuses terminal + * rows as well and clears the retry bookkeeping, so the auto worker keeps + * scheduling the job. + */ + armWorkMemoryIdleFlush(input: { + sessionId: string; + userId: string; + lastActivityAt: string; + runAfter: string; + at?: string; + }): EvolutionJobRecord { + const at = input.at ?? nowIso(); + const dedupeKey = `work_memory_idle_flush:${input.sessionId}`; + const payload = toJson({ lastActivityAt: input.lastActivityAt, runAfter: input.runAfter }); + const existing = this.getJobByDedupeKey(dedupeKey); + if (existing) { + this.db.prepare( + `UPDATE evolution_jobs + SET status = 'queued', + payload_json = ?, + attempts = 0, + leased_until = NULL, + last_error = NULL, + updated_at = ? + WHERE id = ?` + ).run(payload, at, existing.id); + } else { + this.db.prepare( + `INSERT INTO evolution_jobs ( + id, job_type, status, dedupe_key, user_id, session_id, episode_id, + target_memory_id, scope_key, scope_seq, payload_json, attempts, + max_attempts, leased_until, last_error, created_at, updated_at + ) VALUES (?, 'work_memory_idle_flush', 'queued', ?, ?, ?, NULL, NULL, NULL, NULL, ?, 0, 3, NULL, NULL, ?, ?)` + ).run(newId("job"), dedupeKey, input.userId, input.sessionId, payload, at, at); + } + const job = this.getJobByDedupeKey(dedupeKey); + if (!job) throw new Error(`failed to arm work memory idle flush: ${input.sessionId}`); + return job; + } + listJobs(status?: JobStatus, limit = 50, userId?: string): EvolutionJobRecord[] { void userId; const clauses: string[] = []; @@ -4985,6 +5083,39 @@ export class L3WorldModelRepository { return row ? l3WorldModelInputTraceFromSql(row) : undefined; } + /** Highest trace sequence registered for a Session, or 0 when it has none. */ + maxInputTraceSeq(sessionId: string): number { + const row = this.db.prepare( + `SELECT COALESCE(MAX(trace_seq), 0) AS trace_seq + FROM l3_world_model_input_traces WHERE session_id = ?` + ).get(sessionId) as { trace_seq: number }; + return Number(row.trace_seq); + } + + /** Input traces for a Session in an inclusive trace sequence range, ascending. */ + listInputTracesInRange( + sessionId: string, + afterTraceSeq: number, + throughTraceSeq: number + ): L3WorldModelInputTraceRecord[] { + return (this.db.prepare( + `SELECT * FROM l3_world_model_input_traces + WHERE session_id = ? AND trace_seq > ? AND trace_seq <= ? + ORDER BY trace_seq ASC` + ).all(sessionId, afterTraceSeq, throughTraceSeq) as SqlL3WorldModelInputTraceRow[]) + .map(l3WorldModelInputTraceFromSql); + } + + /** Most recent input trace timestamp, used as the real last activity of a Session. */ + latestInputTraceCreatedAt(sessionId: string): string | undefined { + const row = this.db.prepare( + `SELECT created_at FROM l3_world_model_input_traces + WHERE session_id = ? + ORDER BY trace_seq DESC LIMIT 1` + ).get(sessionId) as { created_at: string } | undefined; + return row?.created_at; + } + freezeBatches(input: { sessionId: string; trigger: L3WorldModelBatchTrigger; @@ -6039,7 +6170,7 @@ function l3WorldModelSourceMemoryIds(memory?: MemoryRow): string[] { return Array.isArray(value) ? value.filter((item): item is string => typeof item === "string" && Boolean(item)) : []; } -function splitL3TracesByRawTurn( +export function splitL3TracesByRawTurn( traces: L3WorldModelInputTraceRecord[], maxRawTurns: number ): L3WorldModelInputTraceRecord[][] { @@ -7646,6 +7777,7 @@ function bundleIdentity( source_turn_captures: ["user_id", "source", "profile_id", "namespace_key", "conversation_id", "turn_id"], l3_world_model_scopes: ["scope_key"], l3_world_model_session_cursors: ["session_id"], + work_memory_session_cursors: ["session_id"], l3_world_model_input_traces: ["session_id", "trace_seq"], l3_world_model_evidence_batches: ["id"], l3_world_model_batch_targets: ["batch_id", "target_field"], diff --git a/Memory/src/storage/schema.ts b/Memory/src/storage/schema.ts index bea6d90f4..9755c4b92 100644 --- a/Memory/src/storage/schema.ts +++ b/Memory/src/storage/schema.ts @@ -164,6 +164,12 @@ const statements = [ updated_at TEXT NOT NULL )`, + `CREATE TABLE IF NOT EXISTS work_memory_session_cursors ( + session_id TEXT PRIMARY KEY REFERENCES sessions(id) ON DELETE CASCADE, + last_extracted_seq INTEGER NOT NULL DEFAULT 0 CHECK (last_extracted_seq >= 0), + updated_at TEXT NOT NULL + )`, + `CREATE TABLE IF NOT EXISTS episodes ( id TEXT PRIMARY KEY, session_id TEXT NOT NULL REFERENCES sessions(id) ON DELETE CASCADE, diff --git a/Memory/src/types.ts b/Memory/src/types.ts index d270b9c72..b8ad7cb9b 100644 --- a/Memory/src/types.ts +++ b/Memory/src/types.ts @@ -87,6 +87,7 @@ export type JobType = | "skill_batch_evolve" | "skill_trial_resolve" | "work_memory_extract" + | "work_memory_idle_flush" | "feedback_experience"; export interface RuntimeNamespace { diff --git a/Memory/tests/repository/polardb-schema.test.ts b/Memory/tests/repository/polardb-schema.test.ts index aa29aeb00..5dcff859e 100644 --- a/Memory/tests/repository/polardb-schema.test.ts +++ b/Memory/tests/repository/polardb-schema.test.ts @@ -48,6 +48,7 @@ describe("repository PolarDB schema contract", () => { expect(sql).toContain("idx_embedding_retry_due"); expect(sql).toContain("CREATE TABLE IF NOT EXISTS l3_world_model_scopes"); expect(sql).toContain("CREATE TABLE IF NOT EXISTS l3_world_model_session_cursors"); + expect(sql).toContain("CREATE TABLE IF NOT EXISTS work_memory_session_cursors"); expect(sql).toContain("CREATE TABLE IF NOT EXISTS l3_world_model_input_traces"); expect(sql).toContain("CREATE TABLE IF NOT EXISTS l3_world_model_evidence_batches"); expect(sql).toContain("CREATE TABLE IF NOT EXISTS l3_world_model_batch_targets"); diff --git a/Memory/tests/repository/sqlite-schema.test.ts b/Memory/tests/repository/sqlite-schema.test.ts index 4f2bc3018..2e11d9077 100644 --- a/Memory/tests/repository/sqlite-schema.test.ts +++ b/Memory/tests/repository/sqlite-schema.test.ts @@ -112,6 +112,7 @@ describe("repository sqlite schema contract", () => { "l3_world_model_scopes", "sessions", "l3_world_model_session_cursors", + "work_memory_session_cursors", "episodes", "raw_turns", "l3_world_model_input_traces", @@ -763,6 +764,9 @@ describe("repository sqlite schema contract", () => { ).get()).toEqual({ status: "open" }); expect((migrated.db.prepare(`PRAGMA table_info(l3_world_model_scopes)`).all() as Array<{ name: string }>) .map((column) => column.name)).toContain("workspace_uri"); + expect(migrated.db.prepare( + `SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'work_memory_session_cursors'` + ).get()).toEqual({ name: "work_memory_session_cursors" }); const projectEnvironmentColumns = migrated.db.prepare( `PRAGMA table_info(l3_world_model_project_environment_state)` ).all() as Array<{ name: string }>; diff --git a/Memory/tests/service/session/session-lifecycle.test.ts b/Memory/tests/service/session/session-lifecycle.test.ts index 47b126633..df5d8ac18 100644 --- a/Memory/tests/service/session/session-lifecycle.test.ts +++ b/Memory/tests/service/session/session-lifecycle.test.ts @@ -458,6 +458,17 @@ describe("MemoryService / session / lifecycle", () => { query: "add a project rule", answer: "the project rule was added" }); + expect(complete.jobs.map((job) => job.jobType)).toEqual([ + "trace_summary", + "episode_idle_close", + "episode_title" + ]); + expect(db.db.prepare( + `SELECT job_type, status, session_id + FROM evolution_jobs WHERE job_type = 'work_memory_idle_flush'` + ).all()).toEqual([ + { job_type: "work_memory_idle_flush", status: "queued", session_id: opened.sessionId } + ]); expect(db.db.prepare( `SELECT l1_memory_id, raw_turn_id, trace_seq FROM l3_world_model_input_traces WHERE session_id = ?` diff --git a/Memory/tests/service/work-memory/work-memory-pipeline.test.ts b/Memory/tests/service/work-memory/work-memory-pipeline.test.ts index 4f26a1c77..30cb9d001 100644 --- a/Memory/tests/service/work-memory/work-memory-pipeline.test.ts +++ b/Memory/tests/service/work-memory/work-memory-pipeline.test.ts @@ -2,12 +2,17 @@ import { afterEach, describe, expect, it, vi } from "vitest"; import type { LlmClient } from "../../../src/model/types.js"; import { Repositories, type EvolutionJobRecord } from "../../../src/storage/repositories.js"; import type { MemoryRow } from "../../../src/types.js"; +import type { MemoryService } from "../../../src/service/memory-service.js"; import { canonicalWorkMemoryText, WorkMemoryPipeline } from "../../../src/service/work-memory/work-memory-pipeline.js"; import { stableHash } from "../../../src/utils/id.js"; -import { createCapturingEmbedder, createMemoryServiceFixture } from "../../fixtures/memory-service-fixture.js"; +import { + createCapturingEmbedder, + createMemoryServiceFixture, + runWorkerRounds +} from "../../fixtures/memory-service-fixture.js"; import { upsertMemoryVectorForTest } from "../../fixtures/evolution-fixture.js"; const { cleanup, createTestService } = createMemoryServiceFixture(); @@ -259,8 +264,277 @@ describe("Work Memory pipeline", () => { expect(recall.injectedContext.markdown).toContain("Requirement: 固定 SFT 数据清洗流程"); expect(recall.injectedContext.markdown).toContain("Requirement rationale: 保证训练结果可复现"); }); + + it("从没压缩过的会话在 session close 时抽完全部 trace", () => { + const { db, service } = createTestService(); + const opened = openWorkMemorySession(service, "work-memory-close-session", "work-memory-close-user"); + service.completeTurn("work-memory-close-turn-1", { + sessionId: opened.sessionId, + query: "第一条要求。", + answer: "已记录。", + status: "succeeded" + }); + service.completeTurn("work-memory-close-turn-2", { + sessionId: opened.sessionId, + query: "第二条要求。", + answer: "也记录了。", + status: "succeeded" + }); + const repos = new Repositories(db.db); + expect(workMemoryJobs(repos)).toHaveLength(0); + + service.closeSession(opened.sessionId, { namespace: opened.namespace }); + + const jobs = workMemoryJobs(repos); + expect(jobs).toHaveLength(1); + expect(jobs[0]?.payload.qa).toEqual([ + { user: "第一条要求。", assistant: "已记录。" }, + { user: "第二条要求。", assistant: "也记录了。" } + ]); + expect(repos.runtime.getWorkMemoryCursor(opened.sessionId).lastExtractedSeq).toBe(2); + }); + + it("压缩之后再 close 只抽尾巴,已抽过的窗口不再入队", () => { + const { db, service } = createTestService(); + const opened = openWorkMemorySession(service, "work-memory-tail-session", "work-memory-tail-user"); + const first = service.completeTurn("work-memory-tail-turn-1", { + sessionId: opened.sessionId, + query: "压缩前的要求。", + answer: "已记录。", + status: "succeeded" + }); + service.l3WorldModelBoundary(opened.sessionId, { + requestId: "work-memory-tail-request-1", + adapterId: "codex-memory", + source: "codex", + namespace: opened.namespace, + trigger: "token_compaction", + throughL1MemoryId: first.l1MemoryId + }); + const repos = new Repositories(db.db); + expect(workMemoryJobs(repos)).toHaveLength(1); + + service.completeTurn("work-memory-tail-turn-2", { + sessionId: opened.sessionId, + query: "压缩后的要求。", + answer: "也记录了。", + status: "succeeded" + }); + service.closeSession(opened.sessionId, { namespace: opened.namespace }); + + const jobs = workMemoryJobs(repos); + expect(jobs).toHaveLength(2); + expect(jobs[1]?.payload.qa).toEqual([{ user: "压缩后的要求。", assistant: "也记录了。" }]); + expect(repos.runtime.getWorkMemoryCursor(opened.sessionId).lastExtractedSeq).toBe(2); + }); + + it("压缩之后再 idle,不再产生新的 extract", async () => { + const { db, service } = createTestService(); + const opened = openWorkMemorySession(service, "work-memory-idle-session", "work-memory-idle-user"); + const completed = service.completeTurn("work-memory-idle-turn", { + sessionId: opened.sessionId, + query: "完整要求。", + answer: "已记录。", + status: "succeeded" + }); + service.l3WorldModelBoundary(opened.sessionId, { + requestId: "work-memory-idle-request", + adapterId: "codex-memory", + source: "codex", + namespace: opened.namespace, + trigger: "token_compaction", + throughL1MemoryId: completed.l1MemoryId + }); + const repos = new Repositories(db.db); + makeIdleFlushDue(db, opened.sessionId); + backdateInputTraces(db, opened.sessionId); + + await runWorkerRounds(service, 3, 20); + + expect(workMemoryJobs(repos)).toHaveLength(1); + }); + + it("idle 到期时抽出未抽取的增量,游标推进到当前最大 trace", async () => { + const { db, service } = createTestService(); + const opened = openWorkMemorySession(service, "work-memory-idle-flush-session", "work-memory-idle-flush-user"); + service.completeTurn("work-memory-idle-flush-turn", { + sessionId: opened.sessionId, + query: "idle 才抽的要求。", + answer: "已记录。", + status: "succeeded" + }); + const repos = new Repositories(db.db); + expect(workMemoryJobs(repos)).toHaveLength(0); + + makeIdleFlushDue(db, opened.sessionId); + backdateInputTraces(db, opened.sessionId); + await runWorkerRounds(service, 3, 20); + + const jobs = workMemoryJobs(repos); + expect(jobs).toHaveLength(1); + expect(jobs[0]?.payload.qa).toEqual([{ user: "idle 才抽的要求。", assistant: "已记录。" }]); + expect(repos.runtime.getWorkMemoryCursor(opened.sessionId).lastExtractedSeq).toBe(1); + }); + + it("idle 之后 close 不再重复抽取,也不关 episode", async () => { + const { db, service } = createTestService(); + const opened = openWorkMemorySession(service, "work-memory-idle-close-session", "work-memory-idle-close-user"); + service.completeTurn("work-memory-idle-close-turn", { + sessionId: opened.sessionId, + query: "先 idle 再关闭。", + answer: "已记录。", + status: "succeeded" + }); + const repos = new Repositories(db.db); + makeIdleFlushDue(db, opened.sessionId); + backdateInputTraces(db, opened.sessionId); + await runWorkerRounds(service, 3, 20); + expect(workMemoryJobs(repos)).toHaveLength(1); + + service.closeSession(opened.sessionId, { namespace: opened.namespace }); + + expect(workMemoryJobs(repos)).toHaveLength(1); + expect(repos.runtime.listJobs(undefined, 100).some((job) => job.jobType === "l3_world_model_update")).toBe(true); + }); + + it("同一 session 两次 turn 武装同一个 dedupeKey,runAfter 以第二次为准", () => { + const { db, service } = createTestService(); + const opened = openWorkMemorySession(service, "work-memory-arm-session", "work-memory-arm-user"); + service.completeTurn("work-memory-arm-turn-1", { + sessionId: opened.sessionId, + query: "第一次说话。", + answer: "好。", + status: "succeeded" + }); + const repos = new Repositories(db.db); + const first = idleFlushJobs(repos); + expect(first).toHaveLength(1); + const firstRunAfter = String(first[0]?.payload.runAfter); + + service.completeTurn("work-memory-arm-turn-2", { + sessionId: opened.sessionId, + query: "第二次说话。", + answer: "好。", + status: "succeeded" + }); + + const second = idleFlushJobs(repos); + expect(second).toHaveLength(1); + expect(second[0]?.id).toBe(first[0]?.id); + expect(Date.parse(String(second[0]?.payload.runAfter))).toBeGreaterThan(Date.parse(firstRunAfter)); + }); + + it("session 已关闭时 idle flush 不抽取", async () => { + const { db, service } = createTestService(); + const opened = openWorkMemorySession(service, "work-memory-closed-idle-session", "work-memory-closed-idle-user"); + service.completeTurn("work-memory-closed-idle-turn", { + sessionId: opened.sessionId, + query: "先关闭再 idle。", + answer: "已记录。", + status: "succeeded" + }); + const repos = new Repositories(db.db); + service.closeSession(opened.sessionId, { namespace: opened.namespace }); + const before = workMemoryJobs(repos).length; + + makeIdleFlushDue(db, opened.sessionId); + backdateInputTraces(db, opened.sessionId); + await runWorkerRounds(service, 3, 20); + + expect(workMemoryJobs(repos)).toHaveLength(before); + }); + + it("失败后重新武装的 idle flush 仍然可被调度", () => { + const { db, service } = createTestService(); + const opened = openWorkMemorySession(service, "work-memory-retry-session", "work-memory-retry-user"); service.completeTurn("work-memory-retry-turn", { + sessionId: opened.sessionId, + query: "重试路径。", + answer: "好。", + status: "succeeded" + }); + const repos = new Repositories(db.db); + const job = idleFlushJobs(repos)[0]!; + db.db.prepare( + `UPDATE evolution_jobs SET status = 'failed', payload_json = '{}', last_error = 'boom' WHERE id = ?` + ).run(job.id); + + service.completeTurn("work-memory-retry-turn-2", { + sessionId: opened.sessionId, + query: "再来一次。", + answer: "好。", + status: "succeeded" + }); + + const revived = idleFlushJobs(repos)[0]!; + expect(revived.status).toBe("queued"); + expect(revived.payload.lastActivityAt).toEqual(expect.any(String)); + expect(revived.payload.runAfter).toEqual(expect.any(String)); + }); }); +function workMemoryNamespace( + opened: ReturnType, + sessionKey: string +): { + source: string; + profileId: string; + userId: string; + sessionKey: string; + projectId?: string; +} { + const base = { source: "codex", profileId: "default", userId: opened.userId, sessionKey }; + return opened.projectId ? { ...base, projectId: opened.projectId } : base; +} + +function openWorkMemorySession( + service: MemoryService, + sessionKey: string, + userId: string +): ReturnType & { sessionKey: string; namespace: ReturnType } { + const opened = service.openSession({ + l3WorldModelProtocolVersion: 2, + l3WorldModelTransition: "resume_only", + workspaceUri: `file:///tmp/${sessionKey}`, + workspaceHostId: stableHash(sessionKey).slice(0, 64), + namespace: { source: "codex", profileId: "default", sessionKey, userId } + }); + return { ...opened, sessionKey, namespace: workMemoryNamespace(opened, sessionKey) }; +} + +function workMemoryJobs(repos: Repositories): EvolutionJobRecord[] { + return repos.runtime.listJobs(undefined, 100).filter((job) => job.jobType === "work_memory_extract"); +} + +function idleFlushJobs(repos: Repositories): EvolutionJobRecord[] { + return repos.runtime.listJobs(undefined, 100).filter((job) => job.jobType === "work_memory_idle_flush"); +} + +/** Pull the armed idle flush forward so the next worker pass leases it. */ +function makeIdleFlushDue(db: { db: import("better-sqlite3").Database }, sessionId: string): void { + const due = new Date(Date.now() - 1000).toISOString(); + idleFlushJobs(new Repositories(db.db)) + .filter((job) => job.sessionId === sessionId) + .forEach((job) => { + db.db.prepare( + `UPDATE evolution_jobs SET payload_json = ? WHERE id = ?` + ).run(JSON.stringify({ ...job.payload, runAfter: due }), job.id); + }); +} + +/** + * Backdate a Session's input traces so the idle handler sees a real quiet gap. + * The handler reads the newest trace timestamp, not the clock, as last activity. + */ +function backdateInputTraces( + db: { db: import("better-sqlite3").Database }, + sessionId: string, + idleMs = 3 * 60 * 60 * 1000 +): void { + db.db.prepare( + `UPDATE l3_world_model_input_traces SET created_at = ? WHERE session_id = ?` + ).run(new Date(Date.now() - idleMs).toISOString(), sessionId); +} + function findWorkMemoryJob(repos: Repositories): EvolutionJobRecord { const job = repos.runtime.listJobs(undefined, 100).find((item) => item.jobType === "work_memory_extract"); if (!job) throw new Error("work memory job not found");