diff --git a/packages/@n8n/agents/src/runtime/__tests__/memory-orchestrator-mid-run.test.ts b/packages/@n8n/agents/src/runtime/__tests__/memory-orchestrator-mid-run.test.ts index e933ab431b6..b3140de86e2 100644 --- a/packages/@n8n/agents/src/runtime/__tests__/memory-orchestrator-mid-run.test.ts +++ b/packages/@n8n/agents/src/runtime/__tests__/memory-orchestrator-mid-run.test.ts @@ -219,10 +219,8 @@ describe('MemoryOrchestrator.maybeObserveMidRun', () => { expect(await store.getActiveObservationLog({ observationScopeId: THREAD_ID })).toHaveLength(2); }); - // The per-message transcript adds a `[timestamp] role:` header (~33 chars for - // user, ~38 for assistant), so with a character-count token counter a - // 750-char message lands at ~783 tokens: inside the soft band [700, 1000) - // of a 1000-token hard threshold, and two such messages cross the hard one. + // With the character-count token counter, a 750-char message is inside the + // soft band [700, 1000), and two such messages cross the hard threshold. it('schedules the observer in the background at the soft threshold and activates at a later boundary', async () => { const store = new InMemoryMemory(); @@ -321,6 +319,45 @@ describe('MemoryOrchestrator.maybeObserveMidRun', () => { expect(cursor?.lastObservedMessageId).toBe(list.messages().at(-1)?.id); }); + it('budgets full tool payloads as the model sees them, not the truncated observer rendering', async () => { + const store = new InMemoryMemory(); + const observe = vi.fn( + async () => await Promise.resolve('* CRITICAL (14:30) Large fetch summarized.'), + ); + const { orchestrator } = buildOrchestrator(store, { + observerThresholdTokens: 5_000, + observe, + observationLogTailLimit: 20, + }); + const list = new AgentMessageList(); + list.addInput([userMsg('fetch the report')]); + // One large tool result dominates the window. The budget must count the + // full payload the model receives. + list.addResponse([ + { + role: 'assistant', + content: [ + { + type: 'tool-call', + toolCallId: 'tc1', + toolName: 'fetch_report', + input: { url: 'https://example.com/report' }, + state: 'resolved', + output: { data: 'x'.repeat(20_000) }, + }, + ], + }, + ]); + + await orchestrator.maybeObserveMidRun(list, runOptions()); + + expect(observe).toHaveBeenCalledTimes(1); + expect(await store.getActiveObservationLog({ observationScopeId: THREAD_ID })).toHaveLength(1); + expect(list.forLlm('base').messages).toEqual([ + { role: 'user', content: OBSERVATION_CONTINUATION_REMINDER }, + ]); + }); + it('latches mid-run observation off after repeated non-advancing observer runs', async () => { const store = new InMemoryMemory(); // Runs but never yields a parseable observation, so the cursor never advances. @@ -343,6 +380,132 @@ describe('MemoryOrchestrator.maybeObserveMidRun', () => { }); }); +describe('MemoryOrchestrator.saveToMemory observer gating', () => { + it('does not schedule the observer below the threshold', async () => { + const store = new InMemoryMemory(); + const observe = vi.fn( + async () => await Promise.resolve('* CRITICAL (14:30) Should not appear.'), + ); + const { orchestrator, tracker } = buildOrchestrator(store, { + observerThresholdTokens: 100_000, + observe, + observationLogTailLimit: 20, + }); + const list = new AgentMessageList(); + list.addInput([userMsg('hello')]); + list.addResponse([assistantMsg('hi')]); + + await orchestrator.saveToMemory(list, runOptions()); + await tracker.flush(); + + expect(observe).not.toHaveBeenCalled(); + expect(await store.getActiveObservationLog({ observationScopeId: THREAD_ID })).toEqual([]); + expect(await store.getMessagesForObservationScope(THREAD_ID)).toHaveLength(2); + }); + + it('schedules the observer once the threshold is crossed', async () => { + const store = new InMemoryMemory(); + const observe = vi.fn( + async () => await Promise.resolve('* CRITICAL (14:30) Threshold crossed.'), + ); + const { orchestrator, tracker } = buildOrchestrator(store, { + observerThresholdTokens: 10, + observe, + observationLogTailLimit: 20, + }); + const list = new AgentMessageList(); + list.addInput([userMsg('a message crossing the threshold')]); + + await orchestrator.saveToMemory(list, runOptions()); + await tracker.flush(); + + expect(observe).toHaveBeenCalledTimes(1); + expect(await store.getCursor(THREAD_ID)).not.toBeNull(); + }); + + it('deduplicates an in-flight mid-run task without blocking later post-turn observation', async () => { + const store = new InMemoryMemory(); + const pendingObservation = deferred(); + const observe = vi.fn(async () => await pendingObservation.promise); + const { orchestrator, tracker } = buildOrchestrator(store, { + observerThresholdTokens: 1000, + observe, + observationLogTailLimit: 20, + }); + const list = new AgentMessageList(); + list.addInput([userMsg('x'.repeat(750))]); + + await orchestrator.maybeObserveMidRun(list, runOptions()); + await vi.waitFor(() => expect(observe).toHaveBeenCalledTimes(1)); + + list.addResponse([assistantMsg('y'.repeat(750))]); + await orchestrator.saveToMemory(list, runOptions()); + pendingObservation.resolve('* CRITICAL (14:30) Observed.'); + await tracker.flush(); + + expect(observe).toHaveBeenCalledTimes(1); + + const nextList = new AgentMessageList(); + await orchestrator.loadInto(nextList, runOptions()); + nextList.addInput([userMsg('z'.repeat(250))]); + await orchestrator.saveToMemory(nextList, runOptions()); + await tracker.flush(); + + expect(observe).toHaveBeenCalledTimes(2); + }); + + it('gates normally when a settled mid-run task never advanced the cursor', async () => { + const store = new InMemoryMemory(); + const observe = vi + .fn(async () => await Promise.resolve('* CRITICAL (14:30) Observed post-turn.')) + .mockResolvedValueOnce('not a bullet line'); + const { orchestrator, tracker } = buildOrchestrator(store, { + observerThresholdTokens: 1000, + observe, + observationLogTailLimit: 20, + }); + const list = new AgentMessageList(); + list.addInput([userMsg('x'.repeat(750))]); + + // Soft crossing schedules a background task that settles without + // advancing; no later boundary consumes it before the turn ends. + await orchestrator.maybeObserveMidRun(list, runOptions()); + await tracker.flush(); + expect(await store.getCursor(THREAD_ID)).toBeNull(); + + list.addResponse([assistantMsg('y'.repeat(750))]); + await orchestrator.saveToMemory(list, runOptions()); + await tracker.flush(); + + expect(observe).toHaveBeenCalledTimes(2); + expect(await store.getCursor(THREAD_ID)).not.toBeNull(); + }); + + it('persists the turn when observer budget estimation fails', async () => { + const store = new InMemoryMemory(); + const observe = vi.fn( + async () => await Promise.resolve('* CRITICAL (14:30) Should not appear.'), + ); + const { orchestrator, tracker } = buildOrchestrator( + store, + { + observerThresholdTokens: 10, + observe, + observationLogTailLimit: 20, + }, + async () => await Promise.reject(new Error('token counter exploded')), + ); + const list = new AgentMessageList(); + list.addInput([userMsg('a message crossing the threshold')]); + + await expect(orchestrator.saveToMemory(list, runOptions())).resolves.toBeUndefined(); + await tracker.flush(); + + expect(observe).not.toHaveBeenCalled(); + expect(await store.getMessagesForObservationScope(THREAD_ID)).toHaveLength(1); + }); +}); + describe('MemoryOrchestrator.persistTurnDelta', () => { const observationalMemory: ObservationalMemoryConfig = { observerThresholdTokens: 100_000, diff --git a/packages/@n8n/agents/src/runtime/__tests__/observation-log-observer.test.ts b/packages/@n8n/agents/src/runtime/__tests__/observation-log-observer.test.ts index fddc46329c5..6bb01adaf65 100644 --- a/packages/@n8n/agents/src/runtime/__tests__/observation-log-observer.test.ts +++ b/packages/@n8n/agents/src/runtime/__tests__/observation-log-observer.test.ts @@ -261,7 +261,7 @@ describe('renderObserverTranscript', () => { toolName: 'lookup_workflow', input: { workflow: 'daily-report-prod' }, state: 'resolved', - output: { rows: [{ id: 1 }], blob: 'x'.repeat(80) }, + output: { rows: [{ id: 1 }], notes: 'x'.repeat(80), blob: 'x'.repeat(80) }, }, ], }, @@ -275,6 +275,7 @@ describe('renderObserverTranscript', () => { expect(transcript).toContain('"workflow":"daily-report-prod"'); expect(transcript).toContain('tool_result lookup_workflow'); expect(transcript).toContain('[truncated'); + expect(transcript).toContain('"blob":"[omitted large blob]"'); }); it('redacts credential-looking tool inputs and outputs before serialization', () => { @@ -436,31 +437,6 @@ describe('renderObserverTranscript', () => { }); describe('runObservationLogObserver', () => { - it('waits until the unobserved transcript reaches the token threshold', async () => { - const store = new InMemoryMemory(); - await store.saveThread({ id: 'thread-1', resourceId: 'user-1' }); - await store.saveMessages({ - threadId: 'thread-1', - resourceId: 'user-1', - messages: [message('m1', 'user', 'short turn', new Date(2026, 4, 12, 14, 30))], - }); - - const observe = vi.fn().mockResolvedValue('* CRITICAL (14:30) User said something durable.'); - - const result = await runObservationLogObserver({ - memory: store, - observationScopeId: 'thread-1', - observerThresholdTokens: 999, - observationLogTailLimit: 20, - tokenCounter: () => 1, - observe, - }); - - expect(result).toEqual({ status: 'skipped', reason: 'below-threshold', tokenCount: 1 }); - expect(observe).not.toHaveBeenCalled(); - expect(await store.getCursor('thread-1')).toBeNull(); - }); - it('writes parsed observations and advances the cursor after observing', async () => { const store = new InMemoryMemory(); const parentText = 'User needs the current request remembered.'; @@ -478,7 +454,6 @@ describe('runObservationLogObserver', () => { const result = await runObservationLogObserver({ memory: store, observationScopeId: 'thread-1', - observerThresholdTokens: 1, observationLogTailLimit: 20, tokenCounter, now, @@ -527,7 +502,6 @@ describe('runObservationLogObserver', () => { const result = await runObservationLogObserver({ memory: store, observationScopeId: 'thread-1', - observerThresholdTokens: 1, observationLogTailLimit: 20, tokenCounter: () => 10, now: new Date(2026, 4, 12, 14, 31), @@ -578,7 +552,6 @@ describe('runObservationLogObserver', () => { await runObservationLogObserver({ memory: store, observationScopeId: 'thread-1', - observerThresholdTokens: 1, observationLogTailLimit: 20, tokenCounter: () => 10, now: new Date(2026, 4, 12, 14, 40), @@ -634,7 +607,6 @@ describe('runObservationLogObserver', () => { await runObservationLogObserver({ memory: store, observationScopeId: 'thread-1', - observerThresholdTokens: 1, observationLogTailLimit: 20, tokenCounter: () => 10, now: new Date(2026, 4, 12, 14, 31), diff --git a/packages/@n8n/agents/src/runtime/memory/memory-orchestrator.ts b/packages/@n8n/agents/src/runtime/memory/memory-orchestrator.ts index 2724de67c60..db018b7eb6d 100644 --- a/packages/@n8n/agents/src/runtime/memory/memory-orchestrator.ts +++ b/packages/@n8n/agents/src/runtime/memory/memory-orchestrator.ts @@ -7,7 +7,6 @@ import { import { createFilteredLogger } from '../logger'; import { compareKeyset, saveMessagesToThread } from './memory-store'; import { - renderObserverTranscript, runObservationLogObserver, type ObservationLogObserverMemory, type RunObservationLogObserverResult, @@ -51,7 +50,7 @@ const DEFAULT_MEMORY_TASK_LOCK_TTL_MS = 30_000; const MID_RUN_SOFT_THRESHOLD_RATIO = 0.7; /** * Consecutive observer attempts that failed to advance the cursor (provider - * outage, persistently unparseable output, estimator disagreement) before + * outage, lock contention, persistently unparseable output) before * mid-run observation stops retrying for the rest of the run. Without the * latch, the hard path would fire a blocking observer call at every loop * boundary. Post-turn observation is unaffected. @@ -59,6 +58,29 @@ const MID_RUN_SOFT_THRESHOLD_RATIO = 0.7; const MID_RUN_MAX_NON_ADVANCING_ATTEMPTS = 3; const logger = createFilteredLogger(); +function stringifyForBudget(value: unknown): string { + try { + return JSON.stringify(value) ?? ''; + } catch { + return ''; + } +} + +function serializeMessageForBudget(message: AgentDbMessage): string { + if (!('role' in message) || !Array.isArray(message.content)) return ''; + const parts: string[] = []; + for (const content of message.content) { + if (content.type === 'text') { + parts.push(content.text); + } else if (content.type === 'tool-call') { + parts.push(content.toolName, stringifyForBudget(content.input)); + if (content.state === 'resolved') parts.push(stringifyForBudget(content.output)); + else if (content.state === 'rejected') parts.push(content.error); + } + } + return parts.join('\n'); +} + function hasFunctionProperty( value: object, property: K, @@ -103,10 +125,12 @@ export class MemoryOrchestrator { private episodicMemoryTasksByResource = new Map>(); /** - * Per-message observer-transcript token estimates for the mid-run budget, - * cached by message id so the boundary check never re-encodes messages. + * Per-message model-facing token estimates, cached by message id so + * observation gates never re-encode messages. */ - private midRunTokenCounts = new Map(); + private visibleTokenEstimates = new Map(); + + private visibleTokenEstimateTotal = 0; /** In-flight background mid-run observer task; `result` set on settlement. */ private midRunObserverTask: MidRunObserverTask | undefined; @@ -198,7 +222,7 @@ export class MemoryOrchestrator { list: AgentMessageList, options: (RunOptions & ExecutionOptions) | undefined, ): Promise { - this.resetMidRunState(); + this.resetRunState(); if (this.config.memory && options?.persistence?.threadId) { const telemetry = this.runtimeTelemetry.resolve(options); const memMessages = await this.loadHistoryMessages(options.persistence, telemetry); @@ -388,7 +412,8 @@ export class MemoryOrchestrator { // Memory jobs receive the execution counter so their LLM and embedding // usage contributes to token_count. - const observationTasks = this.scheduleObservationLogJobs( + const observationTasks = await this.scheduleObservationLogJobs( + list, options.persistence, options.executionCounter, telemetry, @@ -402,8 +427,8 @@ export class MemoryOrchestrator { } /** - * Mid-run observation, called at clean loop boundaries. When the - * unobserved-transcript token budget crosses a soft threshold + * Mid-run observation, called at clean loop boundaries. When the visible + * window's estimated model-facing token size crosses a soft threshold * (`MID_RUN_SOFT_THRESHOLD_RATIO` x `observerThresholdTokens`), persist the * turn-so-far (the Observer reads from the store) and start the Observer as * a background task without blocking the loop. A later boundary activates @@ -434,10 +459,12 @@ export class MemoryOrchestrator { } } - /** Per-run mid-run observation state; reset on each run entry (generate via `loadInto`, resume via `applyObservationMask`). */ - private resetMidRunState(): void { + /** Reset per-run observation state on generate and resume entry. */ + private resetRunState(): void { this.midRunNonAdvancingAttempts = 0; this.lastPersistedTurnKeyset = undefined; + this.visibleTokenEstimates.clear(); + this.visibleTokenEstimateTotal = 0; } /** @@ -532,7 +559,6 @@ export class MemoryOrchestrator { options.persistence, options.executionCounter, telemetry, - softThresholdTokens, ); if (!handle) return; const entry: MidRunObserverTask = { handle }; @@ -542,18 +568,30 @@ export class MemoryOrchestrator { this.midRunObserverTask = entry; } - /** Sum cached per-message observer-transcript token estimates for the visible window. */ + /** + * Sum cached per-message estimates of the visible model-facing text and + * full tool payloads. Non-text blocks such as files are excluded. + */ private async estimateVisibleBudget(list: AgentMessageList): Promise { const visible = list.llmVisibleMessages(); for (const message of visible) { - if (!this.midRunTokenCounts.has(message.id)) { - this.midRunTokenCounts.set( - message.id, - await this.tokenCounter(renderObserverTranscript([message])), - ); + if (!this.visibleTokenEstimates.has(message.id)) { + const estimate = await this.tokenCounter(serializeMessageForBudget(message)); + this.visibleTokenEstimates.set(message.id, estimate); + this.visibleTokenEstimateTotal += estimate; + } + } + return this.visibleTokenEstimateTotal; + } + + private pruneVisibleTokenEstimates(list: AgentMessageList): void { + const visibleIds = new Set(list.llmVisibleMessages().map((message) => message.id)); + for (const [messageId, estimate] of this.visibleTokenEstimates) { + if (!visibleIds.has(messageId)) { + this.visibleTokenEstimates.delete(messageId); + this.visibleTokenEstimateTotal -= estimate; } } - return visible.reduce((sum, m) => sum + (this.midRunTokenCounts.get(m.id) ?? 0), 0); } /** Mask the window up to the persisted cursor and refresh the injected log. No LLM call. */ @@ -566,6 +604,7 @@ export class MemoryOrchestrator { const cursor = await memory.getCursor(persistence.threadId); if (!cursor) return; list.maskObservedMessages(cursor); + this.pruneVisibleTokenEstimates(list); await this.setListObservationLogMemory(list, persistence); } @@ -584,7 +623,7 @@ export class MemoryOrchestrator { // mid-run state so a cached runtime cannot carry the suspended run's // watermark or latch into the resumed run — the resolved tool call // mutated in place and must be re-persisted. - this.resetMidRunState(); + this.resetRunState(); const { memory, observationalMemory } = this.config; if (!observationalMemory) return; if (!memory || !persistence || !hasObservationLogObserverMemory(memory)) return; @@ -603,7 +642,6 @@ export class MemoryOrchestrator { persistence: AgentPersistenceOptions, executionCounter?: AgentExecutionCounter, telemetry?: BuiltTelemetry, - effectiveThresholdTokens?: number, ): ScopedMemoryTaskHandle | undefined { const { memory, observationalMemory } = this.config; if (!memory || !observationalMemory || !hasObservationLogObserverMemory(memory)) { @@ -621,10 +659,6 @@ export class MemoryOrchestrator { await runObservationLogObserver({ memory, ...scope, - // The observer re-checks its delta against this threshold; a - // soft-threshold background run must pass the soft value or the - // re-check would skip it as below-threshold. - observerThresholdTokens: effectiveThresholdTokens ?? observerThresholdTokens, observationLogTailLimit: observationalMemory.observationLogTailLimit ?? 0, observe, tokenCounter: this.tokenCounter, @@ -634,11 +668,12 @@ export class MemoryOrchestrator { ); } - private scheduleObservationLogJobs( + private async scheduleObservationLogJobs( + list: AgentMessageList, persistence: AgentPersistenceOptions, executionCounter?: AgentExecutionCounter, telemetry?: BuiltTelemetry, - ): Array> { + ): Promise>> { const { memory, observationalMemory } = this.config; if (!memory || !observationalMemory || !hasObservationLogStore(memory)) return []; @@ -646,8 +681,26 @@ export class MemoryOrchestrator { const runner = this.getMemoryTaskRunner(memory, observationalMemory.lockTtlMs); const tasks: Array> = []; - const observerHandle = this.scheduleObserverTask(persistence, executionCounter, telemetry); - if (observerHandle) tasks.push(observerHandle.done); + // A mid-run task still in flight for this scope already covers the + // messages persisted at its boundary: join it instead of queueing a + // second observer behind it — the post-boundary tail waits for the next + // turn's gate. A task that settled after the last boundary was never + // activated; the run is over, so drop it and let the gauge decide. + const midRunTask = this.midRunObserverTask; + if (midRunTask?.result) this.midRunObserverTask = undefined; + if ( + midRunTask && + !midRunTask.result && + midRunTask.handle.observationScopeId === scope.observationScopeId + ) { + tasks.push(midRunTask.handle.done); + void midRunTask.handle.done.then(() => { + if (this.midRunObserverTask === midRunTask) this.midRunObserverTask = undefined; + }); + } else if (await this.shouldScheduleObserver(list, persistence.threadId)) { + const observerHandle = this.scheduleObserverTask(persistence, executionCounter, telemetry); + if (observerHandle) tasks.push(observerHandle.done); + } const reflect = observationalMemory.reflect; const reflectorThresholdTokens = observationalMemory.reflectorThresholdTokens; @@ -674,6 +727,20 @@ export class MemoryOrchestrator { return tasks; } + private async shouldScheduleObserver(list: AgentMessageList, threadId: string): Promise { + const observerThresholdTokens = this.config.observationalMemory?.observerThresholdTokens; + if (observerThresholdTokens === undefined) return false; + try { + return (await this.estimateVisibleBudget(list)) >= observerThresholdTokens; + } catch (error) { + logger.warn('Post-turn observer gating failed; skipping observation this turn', { + error, + threadId, + }); + return false; + } + } + private scheduleEpisodicMemoryJob( persistence: AgentPersistenceOptions, observationTasks: Array>, diff --git a/packages/@n8n/agents/src/runtime/memory/observation-log-observer.ts b/packages/@n8n/agents/src/runtime/memory/observation-log-observer.ts index 8bcaf729f51..6e77ad01e62 100644 --- a/packages/@n8n/agents/src/runtime/memory/observation-log-observer.ts +++ b/packages/@n8n/agents/src/runtime/memory/observation-log-observer.ts @@ -71,7 +71,6 @@ export interface ObservationLogObserverMemory extends BuiltMemory, BuiltObservat export interface RunObservationLogObserverOpts { memory: ObservationLogObserverMemory; observationScopeId: string; - observerThresholdTokens: number; observationLogTailLimit: number; observe: ObservationLogObserveFn; tokenCounter?: TokenCounter; @@ -83,7 +82,6 @@ export interface RunObservationLogObserverOpts { export type RunObservationLogObserverResult = | { status: 'skipped'; reason: 'no-delta' | 'pending-tool-call' } - | { status: 'skipped'; reason: 'below-threshold'; tokenCount: number } | { status: 'ran'; observationsWritten: number; @@ -204,9 +202,6 @@ export async function runObservationLogObserver( const tokenCounter = opts.tokenCounter ?? estimateObservationTokens; const transcript = renderObserverTranscript(observable); const tokenCount = await tokenCounter(transcript); - if (tokenCount < opts.observerThresholdTokens) { - return { status: 'skipped', reason: 'below-threshold', tokenCount }; - } const observationLogTail = ( await memory.getActiveObservationLog({ @@ -339,7 +334,7 @@ function compactForObserver(value: unknown, options: RenderObserverTranscriptOpt for (const [key, entryValue] of entries.slice(0, maxObjectKeys)) { if (isSensitiveKey(key)) { result[key] = REDACTED_VALUE; - } else if (shouldStripBlob(key, entryValue)) { + } else if (shouldStripBlob(key, entryValue, maxStringChars)) { result[key] = '[omitted large blob]'; } else { result[key] = compactForObserver(entryValue, options); @@ -355,9 +350,9 @@ function isSensitiveKey(key: string): boolean { return SENSITIVE_KEY_PATTERN.test(key); } -function shouldStripBlob(key: string, value: unknown): boolean { +function shouldStripBlob(key: string, value: unknown, maxStringChars: number): boolean { if (typeof value !== 'string') return false; - if (value.length <= DEFAULT_MAX_STRING_CHARS) return false; + if (value.length <= maxStringChars) return false; return /blob|base64|data|file|image/i.test(key); } diff --git a/packages/@n8n/agents/src/types/sdk/memory.ts b/packages/@n8n/agents/src/types/sdk/memory.ts index 42c49c3c655..2cb0c4d4d4b 100644 --- a/packages/@n8n/agents/src/types/sdk/memory.ts +++ b/packages/@n8n/agents/src/types/sdk/memory.ts @@ -289,7 +289,7 @@ export interface ObservationLogMemoryConfig { } export interface ObservationalMemoryConfig { - /** Estimated tokens in unobserved transcript required before the Observer runs. */ + /** Estimated visible-window tokens at which the Observer is scheduled mid-run and post-turn. */ observerThresholdTokens?: number; /** Estimated active observation-log tokens required before the Reflector runs. */ reflectorThresholdTokens?: number;