fix(core): Gate observation on the model-facing window budget (no-changelog) (#37438)

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: Matsu <matias.huhta@n8n.io>
This commit is contained in:
Robin Braumann
2026-09-02 11:13:37 +00:00
committed by GitHub
co-authored by Cursor Matsu
parent d5652d8829
commit 6912024af8
5 changed files with 269 additions and 72 deletions
@@ -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<string>();
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,
@@ -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),
@@ -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<K extends PropertyKey>(
value: object,
property: K,
@@ -103,10 +125,12 @@ export class MemoryOrchestrator {
private episodicMemoryTasksByResource = new Map<string, Promise<unknown>>();
/**
* 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<string, number>();
private visibleTokenEstimates = new Map<string, number>();
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<void> {
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<number> {
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<RunObservationLogObserverResult> | 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<unknown>> {
): Promise<Array<Promise<unknown>>> {
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<Promise<unknown>> = [];
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<boolean> {
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<Promise<unknown>>,
@@ -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);
}
+1 -1
View File
@@ -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;