diff --git a/docs/architecture/memory-quality-gates-and-observability/metrics.md b/docs/architecture/memory-quality-gates-and-observability/metrics.md index bb1836cf34..69e2b4457c 100644 --- a/docs/architecture/memory-quality-gates-and-observability/metrics.md +++ b/docs/architecture/memory-quality-gates-and-observability/metrics.md @@ -18,6 +18,8 @@ SQL, stacks, or other free text. | `agent.retrieval.*.{ftsCandidates,vectorCandidates,selected}` | Count | Authoritative revalidation and fusion completion | Per-Agent counter; same Agent cleanup | Aggregate count | | `agent.retrieval.*.outcomeCounts.*` | Count by closed outcome | Operation `finally` | Per-Agent counter; same Agent cleanup | Closed enum | | `agent.retrieval.*.degradationCounts.*` | Count by closed cause | Operation `finally`; one operation may record multiple causes | Per-Agent counter; same Agent cleanup | Closed enum | +| `agent.queryEmbeddingCircuit.state` | Closed state (`closed`, `open`, `halfOpen`) | Query-embedding circuit transition | Current provider/model only; reset by configuration change and Agent cleanup | Closed enum | +| `agent.queryEmbeddingCircuit.{failures,openCount,skipped}` | Count | Qualifying provider failure, open transition, or skipped query embedding | Per-Agent counter; same configuration and Agent cleanup | Aggregate count | | `agent.extraction.{chunksCompleted,chunksCancelled,chunksFailed,llmCalls,casRetries}` | Count | Chunk settlement or actual second CAS apply | Per-Agent counter; same Agent cleanup | Aggregate count | | `agent.embedding.batchSize` | Row distribution | Embedding drain batch settlement | 256 samples; same Agent cleanup | Aggregate count | | `agent.embedding.drainDurationMs` | Millisecond distribution | Embedding drain batch settlement | 256 samples; same Agent cleanup | Content-free timing | diff --git a/docs/architecture/memory-quality-gates-and-observability/spec.md b/docs/architecture/memory-quality-gates-and-observability/spec.md index 8ac48f71be..6b7dc09cea 100644 --- a/docs/architecture/memory-quality-gates-and-observability/spec.md +++ b/docs/architecture/memory-quality-gates-and-observability/spec.md @@ -102,6 +102,8 @@ metric calculation remain isolated from credentials, providers, and external net - Recorders accept only numbers, booleans, timestamps, and closed enums. - Diagnostics never retain query text, memory content, prompts, vectors, provider responses, API keys, SQL, stacks, exception messages, or other free text. +- Query-embedding circuit diagnostics retain only closed/open/half-open state and aggregate failure, open, and + skip counts; provider/model identity remains internal to the circuit owner. - Collector failures are swallowed and cannot change business results. ### AC-6 — Metric Semantics @@ -117,6 +119,8 @@ metric calculation remain isolated from credentials, providers, and external net - Maintenance reports cheap/heavy outcomes, calls, tokens, and every denied budget step. - Vector warmup distinguishes succeeded, deferred, and failed outcomes. - Provider diagnostics separate admission decisions from deadline, abort, and late-settle race events. +- Query-embedding circuit skips use a closed retrieval degradation cause; cancellation and local control + rejection never increment its provider-health failure count. - Process gauges receive absolute values from resource owners and never aggregate retained Agent state. ### AC-7 — Typed Health Contract and UI diff --git a/docs/architecture/memory-system.md b/docs/architecture/memory-system.md index bf66c68c72..e46b954aab 100644 --- a/docs/architecture/memory-system.md +++ b/docs/architecture/memory-system.md @@ -93,6 +93,12 @@ Memory contribution 必须等待到 soft deadline,成功时限制 token/字符 文本、selection manifest 与成功持久化的 `memory/view_assembled` anchor ID;不能接收或重写 base system prompt。 +Warm recall 的 query embedding 按 Agent 与当前 provider/model identity 使用进程内有界熔断:短窗口内 +连续 deadline/transport failure 会临时跳过 vector path 并直接使用已生成的 FTS candidates;冷却后只 +允许一个 half-open probe,成功自动恢复。取消和本地 capacity rejection 不计 provider health failure; +配置切换、Agent cleanup 与 presenter disposal 清除旧状态。熔断不应用于 embedding batch、warmup、 +dimension discovery 或 text generation,也不改变健康路径的 scoring。 + FTS 在 SQL `LIMIT` 前应用 Agent 和 scope predicate。Vector store 仍按 Agent namespace 查询,使用 有上限的 oversampling,并在 ranking 前通过 SQLite authoritative row 重新校验 owner、scope、 lifecycle、revision 和 embedding identity;不得依赖 vector candidate 本身做授权判断。过滤后不足 diff --git a/docs/issues/memory-query-embedding-circuit-breaker/spec.md b/docs/issues/memory-query-embedding-circuit-breaker/spec.md new file mode 100644 index 0000000000..73dfcb695a --- /dev/null +++ b/docs/issues/memory-query-embedding-circuit-breaker/spec.md @@ -0,0 +1,107 @@ +# Memory Query Embedding Circuit Breaker + +Status: implemented and validated. + +GitHub issue: [#2118](https://github.com/ThinkInAIXYZ/deepchat/issues/2118) + +## Issue + +Warm Memory recall always starts query embedding when the vector store is ready. The provider +gateway and retrieval service bound that call to 800 ms and correctly degrade to FTS, but neither +owner retains provider health across turns. A repeatedly slow embedding path therefore adds the +same pre-stream delay to every Agent turn even though FTS candidates are already available. + +## Impact + +- Every affected turn pays a predictable first-token delay before the existing FTS fallback. +- Requests that are already known to be unhealthy add provider load and repeated warnings. +- Agent and benchmark latency includes avoidable query-embedding deadlines rather than useful work. +- The successful final answer hides the failure mode from users and makes it difficult to diagnose. + +## Root Cause + +- `MemoryProviderGateway` owns per-request deadlines, cancellation, admission, and capacity, but it + intentionally has no cross-request policy and serves all Memory provider purposes. +- `RetrievalService` owns query-embedding de-duplication and FTS degradation, but its retained state + records only in-flight requests. A settled timeout is forgotten before the next turn. +- Runtime diagnostics count individual deadlines and retrieval degradations but expose no current + query-embedding health state or circuit skips. + +## Fix Plan + +Keep the circuit in `RetrievalService`, immediately before query embedding starts. This is the +narrowest owner that can skip only query vectors while preserving the already-computed FTS path. +Do not put the policy in the shared provider gateway, where it could suppress embedding batches, +warmup, dimension discovery, extraction, decisions, or maintenance. + +For each Agent, retain at most the current effective embedding provider/model state: + +- open after two deadline or transport failures in a 30-second window; +- while open, skip new query embedding and continue with FTS immediately; +- after a 30-second cooldown, admit one half-open probe and skip concurrent turns; +- close and reset the failure window after a successful probe; +- reopen after a failed probe without exponential or unbounded cooldown growth. + +Memory cancellation, generic `AbortError`, and local capacity rejection do not count as provider +health failures. Shared in-flight requests settle circuit health only once. Configuration changes, +Agent cleanup, and service disposal remove stale circuit and in-flight state. + +Expose only the closed/open/half-open state and cumulative failure, open, and skip counts through +the existing bounded Agent diagnostics. Add a closed retrieval degradation for circuit skips so an +injection manifest can report that vector recall was intentionally bypassed. Do not retain query +text, memory content, vectors, provider/model identifiers, responses, errors, or other free text in +the circuit or diagnostics. + +## Compatibility And Non-Goals + +- Preserve healthy FTS/vector scoring, candidate limits, adaptive refill, and access accounting. +- Preserve the 800 ms per-request absolute and soft deadlines. +- Do not change provider retry or rate-limit policy. +- Do not persist circuit state or add a database migration. +- Do not share failures across Agents, provider/model identities, or non-query provider purposes. +- Do not add an external telemetry or polling path. + +## Acceptance Criteria + +1. Two qualifying failures open the circuit, and later turns return existing FTS results without + starting or waiting for query embedding. +2. Exactly one probe is admitted after cooldown; success closes the circuit and restores vector + recall, while probe failure reopens it. +3. Agent/provider-model isolation is preserved, and provider changes, Agent cleanup, and disposal + invalidate stale state. +4. Cancellation and local control errors do not increment failure counts or open the circuit. +5. Health diagnostics expose current closed/open/half-open state and content-free failure, open, and + skip counts; skipped retrieval records a closed degradation cause. +6. Deterministic fake-timer tests cover threshold, cooldown, recovery, isolation, cleanup, and + cancellation without changing healthy retrieval results. + +## Task Checklist + +- [x] Validate the issue against the production recall, provider, diagnostics, and lifecycle paths. +- [x] Implement the bounded query-embedding circuit and lifecycle invalidation. +- [x] Extend typed content-free diagnostics and maintained Memory documentation. +- [x] Add focused fake-clock regression coverage. +- [x] Run formatting, i18n, lint, type checking, and relevant Memory/renderer tests. +- [x] Review the staged commit for side effects, compatibility, boundaries, performance, security, + naming, coverage, and maintenance cost before committing. + +## Validation + +Completed on 2026-08-10: + +```bash +pnpm run i18n +pnpm run lint +pnpm run typecheck +pnpm run format:check +pnpm run test:memory +pnpm run test:main +pnpm exec vitest run --config vitest.config.renderer.ts \ + test/renderer/components/MemoryDiagnosticsPanel.test.ts +``` + +- Memory suite: 54 files, 888 tests passed. +- Main suite: 541 files and 6,622 tests passed; 28 files and 371 tests skipped by the existing + suite configuration. +- Renderer diagnostics panel: 1 file, 22 tests passed. +- Formatting, i18n validation, lint, and Node/Web type checking passed. diff --git a/src/main/memory/index.ts b/src/main/memory/index.ts index c51e82bc91..9a50d66425 100644 --- a/src/main/memory/index.ts +++ b/src/main/memory/index.ts @@ -458,6 +458,7 @@ export class MemoryService implements MemoryRuntimePort { agentId, resolveMemoryEmbedding(config) ) + if (embeddingIdentityChanged) this.retrieval.onEmbeddingConfigChanged(agentId) if (embeddingIdentityChanged && observation !== 'changed') { this.runtime.invalidateAgentOperations(agentId) } @@ -742,7 +743,10 @@ export class MemoryService implements MemoryRuntimePort { } getHealth(agentId: string): MemoryHealthDto { - return this.management.getHealth(agentId) + const health = this.management.getHealth(agentId) + health.runtime.agent.queryEmbeddingCircuit.state = + this.retrieval.getQueryEmbeddingCircuitState(agentId) + return health } async deleteMemory(agentId: string, memoryId: string): Promise { diff --git a/src/main/memory/infra/diagnostics/memoryDiagnosticsCollector.ts b/src/main/memory/infra/diagnostics/memoryDiagnosticsCollector.ts index 9247e38d78..07d7c717af 100644 --- a/src/main/memory/infra/diagnostics/memoryDiagnosticsCollector.ts +++ b/src/main/memory/infra/diagnostics/memoryDiagnosticsCollector.ts @@ -74,6 +74,12 @@ type RetrievalDiagnosticsState = { type AgentDiagnosticsState = { lastTouchedAt: number retrieval: Record + queryEmbeddingCircuit: { + state: 'closed' | 'open' | 'halfOpen' + failures: number + openCount: number + skipped: number + } extraction: { chunksCompleted: number chunksCancelled: number @@ -173,6 +179,36 @@ export class MemoryDiagnosticsCollector { }) } + recordQueryEmbeddingCircuitEvent( + agentId: string, + event: 'failure' | 'opened' | 'halfOpen' | 'closed' | 'probeCancelled' | 'skipped' + ): void { + this.safely(() => { + const circuit = this.agent(agentId).queryEmbeddingCircuit + if (event === 'failure') circuit.failures += 1 + else if (event === 'opened') { + circuit.state = 'open' + circuit.openCount += 1 + } else if (event === 'halfOpen') circuit.state = 'halfOpen' + else if (event === 'closed') circuit.state = 'closed' + else if (event === 'probeCancelled') circuit.state = 'open' + else circuit.skipped += 1 + }) + } + + resetQueryEmbeddingCircuit(agentId: string): void { + this.safely(() => { + const state = this.agents.get(agentId) + if (!state) return + state.queryEmbeddingCircuit = { + state: 'closed', + failures: 0, + openCount: 0, + skipped: 0 + } + }) + } + recordExtraction(agentId: string, sample: MemoryExtractionDiagnosticSample): void { this.safely(() => { const counters = this.agent(agentId).extraction @@ -300,6 +336,7 @@ export class MemoryDiagnosticsCollector { ] }) ) as MemoryRuntimeDiagnosticsDto['agent']['retrieval'], + queryEmbeddingCircuit: { ...state.queryEmbeddingCircuit }, extraction: { ...state.extraction }, embedding: { batchSize: distribution(state.embedding.batchSize.snapshot()), @@ -392,6 +429,12 @@ export class MemoryDiagnosticsCollector { retrieval: Object.fromEntries( MEMORY_RETRIEVAL_PURPOSES.map((purpose) => [purpose, retrievalState()]) ) as Record, + queryEmbeddingCircuit: { + state: 'closed', + failures: 0, + openCount: 0, + skipped: 0 + }, extraction: { chunksCompleted: 0, chunksCancelled: 0, diff --git a/src/main/memory/runtimeConstants.ts b/src/main/memory/runtimeConstants.ts index ede89b55fc..84e41d1c1d 100644 --- a/src/main/memory/runtimeConstants.ts +++ b/src/main/memory/runtimeConstants.ts @@ -67,6 +67,9 @@ export const MEMORY_HEALTH_RECENT_FAILURES_LIMIT = 5 export const MEMORY_CREATED_IDS_EVENT_LIMIT = 50 export const RECALL_QUERY_EMBEDDING_TIMEOUT_MS = 800 +export const RECALL_QUERY_EMBEDDING_BREAKER_FAILURE_THRESHOLD = 2 +export const RECALL_QUERY_EMBEDDING_BREAKER_FAILURE_WINDOW_MS = 30 * 1000 +export const RECALL_QUERY_EMBEDDING_BREAKER_COOLDOWN_MS = 30 * 1000 export const RECALL_QUERY_EMBEDDING_STALE_MS = 30 * 1000 export const RECALL_QUERY_EMBEDDING_MAX_CONCURRENT = 2 export const RECALL_VECTOR_QUERY_TIMEOUT_MS = 2 * 1000 diff --git a/src/main/memory/services/retrievalService.ts b/src/main/memory/services/retrievalService.ts index af114a15ed..a598c3c290 100644 --- a/src/main/memory/services/retrievalService.ts +++ b/src/main/memory/services/retrievalService.ts @@ -17,7 +17,10 @@ import { nextMemoryRetrievalCandidateLimit } from '../core/retrievalBudget' import { withSoftDeadline } from '../core/asyncDeadline' -import { isMemoryProviderCancellationError } from '../core/providerCancellation' +import { + MEMORY_PROVIDER_DEADLINE_CODE, + isMemoryProviderCancellationError +} from '../core/providerCancellation' import { evaluateNormalizedMemoryTemporalPolicy, temporalMetadataFromRow } from '../core/temporal' import { createMemoryTopicSuppressionPolicy, @@ -43,6 +46,9 @@ import { import { DECISION_NEIGHBOR_TOP_S, MEMORY_SEARCH_DEFAULT_LIMIT, + RECALL_QUERY_EMBEDDING_BREAKER_COOLDOWN_MS, + RECALL_QUERY_EMBEDDING_BREAKER_FAILURE_THRESHOLD, + RECALL_QUERY_EMBEDDING_BREAKER_FAILURE_WINDOW_MS, RECALL_QUERY_EMBEDDING_MAX_CONCURRENT, RECALL_QUERY_EMBEDDING_STALE_MS, RECALL_QUERY_EMBEDDING_TIMEOUT_MS, @@ -78,10 +84,28 @@ import type { } from '../ports' type QueryEmbeddingInFlight = { + agentId: string startedAt: number promise: Promise + circuit: QueryEmbeddingCircuit + recoveryProbe: boolean + consumerSignals: Set + circuitSettled: boolean } +type QueryEmbeddingCircuit = { + fingerprint: string + failuresInWindow: number + failureWindowStartedAt: number | null + openUntil: number + halfOpenProbe: boolean +} + +type QueryEmbeddingStartResult = + | { status: 'started'; entry: QueryEmbeddingInFlight } + | { status: 'capacity' } + | { status: 'circuitOpen' } + const LEGACY_RETRIEVAL_CANDIDATE_MULTIPLIER = 2 const TEMPORAL_RETRIEVAL_CANDIDATE_MULTIPLIER = 4 const DIRECTIVE_RETRIEVAL_CANDIDATE_MULTIPLIER = 4 @@ -184,6 +208,7 @@ function isStaleExecutionCancellation(error: unknown, isDisposed: boolean): bool export class RetrievalService { private readonly ctx: MemoryRuntimeContext private readonly queryEmbeddingInFlight = new Map>() + private readonly queryEmbeddingCircuits = new Map() constructor( private readonly ports: { @@ -222,6 +247,11 @@ export class RetrievalService { degradations: readonly MemoryRetrievalDegradationCause[] } ): void + recordQueryEmbeddingCircuitEvent?( + agentId: string, + event: 'failure' | 'opened' | 'halfOpen' | 'closed' | 'probeCancelled' | 'skipped' + ): void + resetQueryEmbeddingCircuit?(agentId: string): void } } ) { @@ -560,38 +590,97 @@ export class RetrievalService { private startQueryEmbedding( agentId: string, embedding: MemoryModelRef, - query: string - ): Promise | null { - const key = `${agentId}::${embeddingFingerprint(embedding.providerId, embedding.modelId)}` + query: string, + signal?: AbortSignal + ): QueryEmbeddingStartResult { + const fingerprint = embeddingFingerprint(embedding.providerId, embedding.modelId) + const key = `${agentId}::${fingerprint}` const now = Date.now() + const circuit = this.queryEmbeddingCircuit(agentId, fingerprint) let group = this.queryEmbeddingInFlight.get(key) - if (!group) { - group = new Map() - this.queryEmbeddingInFlight.set(key, group) - } let replacedStale = false - for (const [trackedQuery, entry] of group) { - if (now - entry.startedAt >= RECALL_QUERY_EMBEDDING_STALE_MS) { - group.delete(trackedQuery) - replacedStale = true + if (group) { + for (const [trackedQuery, entry] of group) { + if (entry.circuitSettled) { + group.delete(trackedQuery) + } else if (now - entry.startedAt >= RECALL_QUERY_EMBEDDING_STALE_MS) { + this.settleQueryEmbeddingCircuitCancellation(entry) + group.delete(trackedQuery) + replacedStale = true + } + } + if (group.size === 0) { + this.queryEmbeddingInFlight.delete(key) + group = undefined } } if (replacedStale) logger.warn(`[Memory] stale query embedding replaced for ${agentId}`) - const existing = group.get(query) - if (existing) return existing.promise - if (group.size >= RECALL_QUERY_EMBEDDING_MAX_CONCURRENT) return null + if (circuit.halfOpenProbe || circuit.openUntil > now) { + this.ports.diagnostics?.recordQueryEmbeddingCircuitEvent?.(agentId, 'skipped') + return { status: 'circuitOpen' } + } + + const recoveryProbe = circuit.openUntil > 0 + const existing = group?.get(query) + if (!recoveryProbe && existing) { + existing.consumerSignals.add(signal) + return { status: 'started', entry: existing } + } + if ((group?.size ?? 0) >= RECALL_QUERY_EMBEDDING_MAX_CONCURRENT) { + return { status: 'capacity' } + } - const promise = this.ports.embeddingGateway.getEmbeddings( + if (recoveryProbe) { + circuit.halfOpenProbe = true + this.ports.diagnostics?.recordQueryEmbeddingCircuitEvent?.(agentId, 'halfOpen') + } + if (!group) { + group = new Map() + this.queryEmbeddingInFlight.set(key, group) + } + + let promise: Promise + try { + promise = this.ports.embeddingGateway.getEmbeddings( + agentId, + embedding.providerId, + embedding.modelId, + [query], + 'query-embedding' + ) + } catch (error) { + promise = Promise.reject(error) + } + const entry: QueryEmbeddingInFlight = { agentId, - embedding.providerId, - embedding.modelId, - [query], - 'query-embedding' - ) - const entry: QueryEmbeddingInFlight = { startedAt: now, promise } + startedAt: now, + promise, + circuit, + recoveryProbe, + consumerSignals: new Set([signal]), + circuitSettled: false + } group.set(query, entry) - promise + void promise + .then( + () => { + if (this.queryEmbeddingConsumersCancelled(entry)) { + this.settleQueryEmbeddingCircuitCancellation(entry) + } else { + this.settleQueryEmbeddingCircuitSuccess(entry) + } + }, + (error) => { + if (this.queryEmbeddingConsumersCancelled(entry)) { + this.settleQueryEmbeddingCircuitCancellation(entry) + } else if (this.isQueryEmbeddingCircuitFailure(error)) { + this.settleQueryEmbeddingCircuitFailure(entry) + } else { + this.settleQueryEmbeddingCircuitCancellation(entry) + } + } + ) .finally(() => { const currentGroup = this.queryEmbeddingInFlight.get(key) if (currentGroup?.get(query) === entry) { @@ -600,7 +689,108 @@ export class RetrievalService { } }) .catch(() => undefined) - return promise + return { status: 'started', entry } + } + + private queryEmbeddingCircuit(agentId: string, fingerprint: string): QueryEmbeddingCircuit { + const existing = this.queryEmbeddingCircuits.get(agentId) + if (existing?.fingerprint === fingerprint) return existing + if (existing) { + this.clearQueryEmbeddingInFlight(agentId) + this.ports.diagnostics?.resetQueryEmbeddingCircuit?.(agentId) + } + const circuit: QueryEmbeddingCircuit = { + fingerprint, + failuresInWindow: 0, + failureWindowStartedAt: null, + openUntil: 0, + halfOpenProbe: false + } + this.queryEmbeddingCircuits.set(agentId, circuit) + return circuit + } + + private queryEmbeddingConsumersCancelled(entry: QueryEmbeddingInFlight): boolean { + return ( + entry.consumerSignals.size > 0 && + [...entry.consumerSignals].every((signal) => signal?.aborted === true) + ) + } + + private isQueryEmbeddingCircuitFailure(error: unknown): boolean { + if ((error as { code?: string } | null)?.code === MEMORY_PROVIDER_DEADLINE_CODE) return true + return (error as { name?: string } | null)?.name !== 'AbortError' + } + + private settleQueryEmbeddingCircuitSuccess(entry: QueryEmbeddingInFlight): void { + if (entry.circuitSettled) return + entry.circuitSettled = true + if (this.queryEmbeddingCircuits.get(entry.agentId) !== entry.circuit) return + if (!entry.recoveryProbe && entry.circuit.openUntil > 0) return + entry.circuit.failuresInWindow = 0 + entry.circuit.failureWindowStartedAt = null + entry.circuit.openUntil = 0 + entry.circuit.halfOpenProbe = false + this.ports.diagnostics?.recordQueryEmbeddingCircuitEvent?.(entry.agentId, 'closed') + } + + private settleQueryEmbeddingCircuitFailure(entry: QueryEmbeddingInFlight): void { + if (entry.circuitSettled) return + entry.circuitSettled = true + const circuit = entry.circuit + if (this.queryEmbeddingCircuits.get(entry.agentId) !== circuit) return + this.ports.diagnostics?.recordQueryEmbeddingCircuitEvent?.(entry.agentId, 'failure') + const now = Date.now() + if (entry.recoveryProbe) { + circuit.halfOpenProbe = false + this.openQueryEmbeddingCircuit(entry.agentId, circuit, now) + return + } + if (circuit.openUntil > now) return + if ( + circuit.failureWindowStartedAt === null || + now - circuit.failureWindowStartedAt > RECALL_QUERY_EMBEDDING_BREAKER_FAILURE_WINDOW_MS + ) { + circuit.failureWindowStartedAt = now + circuit.failuresInWindow = 1 + } else { + circuit.failuresInWindow += 1 + } + if (circuit.failuresInWindow >= RECALL_QUERY_EMBEDDING_BREAKER_FAILURE_THRESHOLD) { + this.openQueryEmbeddingCircuit(entry.agentId, circuit, now) + } + } + + private openQueryEmbeddingCircuit( + agentId: string, + circuit: QueryEmbeddingCircuit, + now: number + ): void { + circuit.openUntil = now + RECALL_QUERY_EMBEDDING_BREAKER_COOLDOWN_MS + this.ports.diagnostics?.recordQueryEmbeddingCircuitEvent?.(agentId, 'opened') + logger.warn(`[Memory] query embedding circuit opened for ${agentId}; vector recall paused`) + } + + private settleQueryEmbeddingCircuitCancellation(entry: QueryEmbeddingInFlight): void { + if (entry.circuitSettled) return + entry.circuitSettled = true + if (this.queryEmbeddingCircuits.get(entry.agentId) !== entry.circuit) return + if (!entry.recoveryProbe) return + entry.circuit.halfOpenProbe = false + this.ports.diagnostics?.recordQueryEmbeddingCircuitEvent?.(entry.agentId, 'probeCancelled') + } + + private clearQueryEmbeddingInFlight(agentId: string): void { + for (const key of this.queryEmbeddingInFlight.keys()) { + if (key.startsWith(`${agentId}::`)) this.queryEmbeddingInFlight.delete(key) + } + } + + getQueryEmbeddingCircuitState(agentId: string): 'closed' | 'open' | 'halfOpen' { + const circuit = this.queryEmbeddingCircuits.get(agentId) + if (!circuit) return 'closed' + if (circuit.halfOpenProbe) return 'halfOpen' + return circuit.openUntil > 0 ? 'open' : 'closed' } async searchMemories( @@ -761,9 +951,12 @@ export class RetrievalService { const queryEmbedding = this.startQueryEmbedding( agentId, currentEmbedding, - normalizedQuery + normalizedQuery, + options.signal ) - if (!queryEmbedding) { + if (queryEmbedding.status === 'circuitOpen') { + degradations.add('embeddingCircuitOpen') + } else if (queryEmbedding.status === 'capacity') { logger.warn( `[Memory] query embedding already in flight for ${agentId}; vector recall skipped this turn` ) @@ -771,12 +964,13 @@ export class RetrievalService { const embeddingStartedAt = performance.now() activeStage = 'queryEmbedding' const vectorsResult = await withSoftDeadline( - queryEmbedding, + queryEmbedding.entry.promise, RECALL_QUERY_EMBEDDING_TIMEOUT_MS ) throwIfAborted(options.signal) latencyMs.queryEmbedding = performance.now() - embeddingStartedAt if (vectorsResult.timedOut) { + this.settleQueryEmbeddingCircuitFailure(queryEmbedding.entry) degradations.add('embeddingTimeout') logger.warn( `[Memory] query embedding timed out for ${agentId}; vector recall skipped this turn` @@ -853,6 +1047,7 @@ export class RetrievalService { } const errorName = (error as { name?: string } | null)?.name if ( + activeStage === 'vector' && errorName !== 'AbortError' && !(error instanceof VectorStoreLeaseUnavailableError) ) { @@ -1243,12 +1438,17 @@ export class RetrievalService { } cleanupAgent(agentId: string): void { - for (const key of this.queryEmbeddingInFlight.keys()) { - if (key.startsWith(`${agentId}::`)) this.queryEmbeddingInFlight.delete(key) - } + this.onEmbeddingConfigChanged(agentId) + } + + onEmbeddingConfigChanged(agentId: string): void { + this.clearQueryEmbeddingInFlight(agentId) + this.queryEmbeddingCircuits.delete(agentId) + this.ports.diagnostics?.resetQueryEmbeddingCircuit?.(agentId) } clearAll(): void { this.queryEmbeddingInFlight.clear() + this.queryEmbeddingCircuits.clear() } } diff --git a/src/shared/contracts/routes/memory.routes.ts b/src/shared/contracts/routes/memory.routes.ts index 2f04c83b1f..894210a08e 100644 --- a/src/shared/contracts/routes/memory.routes.ts +++ b/src/shared/contracts/routes/memory.routes.ts @@ -413,6 +413,12 @@ export const MemoryRuntimeDiagnosticsSchema = z.object({ typeof MemoryRetrievalDiagnosticsSchema > ), + queryEmbeddingCircuit: z.object({ + state: z.enum(['closed', 'open', 'halfOpen']), + failures: NonnegativeCountSchema, + openCount: NonnegativeCountSchema, + skipped: NonnegativeCountSchema + }), extraction: z.object({ chunksCompleted: NonnegativeCountSchema, chunksCancelled: NonnegativeCountSchema, @@ -690,6 +696,12 @@ export function createEmptyMemoryRuntimeDiagnostics(): MemoryRuntimeDiagnosticsD } ]) ) as MemoryRuntimeDiagnosticsDto['agent']['retrieval'], + queryEmbeddingCircuit: { + state: 'closed', + failures: 0, + openCount: 0, + skipped: 0 + }, extraction: { chunksCompleted: 0, chunksCancelled: 0, diff --git a/src/shared/types/agent-memory.ts b/src/shared/types/agent-memory.ts index d5ebb343c9..19eb4b100e 100644 --- a/src/shared/types/agent-memory.ts +++ b/src/shared/types/agent-memory.ts @@ -61,6 +61,7 @@ export const MEMORY_RETRIEVAL_DEGRADATION_CAUSES = [ 'vectorCold', 'embeddingTimeout', 'embeddingError', + 'embeddingCircuitOpen', 'storeUnusable', 'storeTimeout', 'storeError', diff --git a/test/main/memory/managementService.test.ts b/test/main/memory/managementService.test.ts index 6a22cd0c71..4790def9d2 100644 --- a/test/main/memory/managementService.test.ts +++ b/test/main/memory/managementService.test.ts @@ -1,8 +1,13 @@ import { describe, expect, it, vi } from 'vitest' import { isSafeAgentId } from '@/memory' +import { MEMORY_PROVIDER_CANCELLATION_CODE } from '@/memory/core/providerCancellation' import { buildMemoryProvenanceKey } from '@/memory/core/scoring' -import { WARM_DIMENSION_FAILURE_COOLDOWN_MS } from '@/memory/runtimeConstants' +import { + RECALL_QUERY_EMBEDDING_BREAKER_COOLDOWN_MS, + RECALL_QUERY_EMBEDDING_BREAKER_FAILURE_WINDOW_MS, + WARM_DIMENSION_FAILURE_COOLDOWN_MS +} from '@/memory/runtimeConstants' import { type IMemoryVectorStore } from '@/memory/types' import { createEmptyMemoryHealth } from '@shared/contracts/routes' import type { DeepChatAgentConfig } from '@shared/types/agent-interface' @@ -1025,20 +1030,34 @@ describe('MemoryService management', () => { } }) - it('starts a fresh query embedding after the prior absolute deadline', async () => { + it('opens after repeated query deadlines and recovers through one half-open probe', async () => { vi.useFakeTimers() try { const repo = createFakeRepository() const store = new FakeVectorStore() - let blockQueryEmbedding = false + let queryMode: 'healthy' | 'timeout' | 'probe' = 'healthy' let queryEmbeddingCalls = 0 - const getEmbeddings = vi.fn((_p: string, _m: string, texts: string[]) => { - if (blockQueryEmbedding && texts[0] !== 'memory warmup') { - queryEmbeddingCalls += 1 - return new Promise(() => undefined) + let resolveProbe: ((vectors: number[][]) => void) | null = null + const getEmbeddings = vi.fn( + (_p: string, _m: string, texts: string[], signal?: AbortSignal) => { + if (queryMode !== 'healthy' && texts[0] !== 'memory warmup') { + queryEmbeddingCalls += 1 + if (queryMode === 'probe') { + return new Promise((resolve) => { + resolveProbe = resolve + }) + } + return new Promise((_resolve, reject) => { + signal?.addEventListener( + 'abort', + () => reject(Object.assign(new Error('cancelled'), { name: 'AbortError' })), + { once: true } + ) + }) + } + return Promise.resolve(texts.map((text) => textToVector(text))) } - return Promise.resolve(texts.map((text) => textToVector(text))) - }) + ) const presenter = new MemoryService({ repository: repo, resolveAgentConfig: () => enabledConfig, @@ -1053,7 +1072,7 @@ describe('MemoryService management', () => { ) await presenter.processPendingEmbeddings('a') - blockQueryEmbedding = true + queryMode = 'timeout' const first = presenter.recall('a', 'redis setup') await vi.advanceTimersByTimeAsync(801) expect((await first).map((item) => item.id)).toEqual([memoryId]) @@ -1064,11 +1083,421 @@ describe('MemoryService management', () => { expect((await second).map((item) => item.id)).toEqual([memoryId]) expect(queryEmbeddingCalls).toBe(2) + + await expect(presenter.recall('a', 'redis setup')).resolves.toEqual([ + expect.objectContaining({ id: memoryId, sources: { fts: true } }) + ]) + expect(queryEmbeddingCalls).toBe(2) + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit).toEqual({ + state: 'open', + failures: 2, + openCount: 1, + skipped: 1 + }) + expect( + presenter.getHealth('a').runtime.agent.retrieval.recall.degradationCounts + .embeddingCircuitOpen + ).toBe(1) + + await vi.advanceTimersByTimeAsync(RECALL_QUERY_EMBEDDING_BREAKER_COOLDOWN_MS) + const failedProbe = presenter.recall('a', 'redis setup') + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit.state).toBe('halfOpen') + await expect(presenter.recall('a', 'redis setup')).resolves.toEqual([ + expect.objectContaining({ id: memoryId, sources: { fts: true } }) + ]) + await vi.advanceTimersByTimeAsync(801) + await expect(failedProbe).resolves.toEqual([ + expect.objectContaining({ id: memoryId, sources: { fts: true } }) + ]) + expect(queryEmbeddingCalls).toBe(3) + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit).toEqual({ + state: 'open', + failures: 3, + openCount: 2, + skipped: 2 + }) + + await vi.advanceTimersByTimeAsync(RECALL_QUERY_EMBEDDING_BREAKER_COOLDOWN_MS - 100) + await expect(presenter.recall('a', 'redis setup')).resolves.toEqual([ + expect.objectContaining({ id: memoryId, sources: { fts: true } }) + ]) + expect(queryEmbeddingCalls).toBe(3) + await vi.advanceTimersByTimeAsync(100) + + queryMode = 'probe' + resolveProbe = null + const probeController = new AbortController() + const cancelledProbe = presenter.buildInjection('a', 'redis setup', { + signal: probeController.signal + }) + await flushMicrotasks(10) + expect(resolveProbe).not.toBeNull() + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit.state).toBe('halfOpen') + probeController.abort() + resolveProbe?.([textToVector('redis setup')]) + await expect(cancelledProbe).rejects.toMatchObject({ name: 'AbortError' }) + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit).toEqual({ + state: 'open', + failures: 3, + openCount: 2, + skipped: 3 + }) + + resolveProbe = null + const recovery = presenter.recall('a', 'redis setup') + await flushMicrotasks(10) + expect(resolveProbe).not.toBeNull() + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit.state).toBe('halfOpen') + expect(queryEmbeddingCalls).toBe(5) + await expect(presenter.recall('a', 'redis setup')).resolves.toEqual([ + expect.objectContaining({ id: memoryId, sources: { fts: true } }) + ]) + + resolveProbe?.([textToVector('redis setup')]) + await expect(recovery).resolves.toEqual([ + expect.objectContaining({ id: memoryId, sources: { vec: true, fts: true } }) + ]) + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit).toEqual({ + state: 'closed', + failures: 3, + openCount: 2, + skipped: 4 + }) + + queryMode = 'healthy' + await expect(presenter.recall('a', 'redis setup')).resolves.toEqual([ + expect.objectContaining({ id: memoryId, sources: { vec: true, fts: true } }) + ]) + } finally { + vi.useRealTimers() + } + }) + + it('stays open when a request admitted before opening succeeds later', async () => { + const repo = createFakeRepository() + const store = new FakeVectorStore() + let queryMode: 'healthy' | 'fail' | 'pending' = 'healthy' + const pendingQueries = new Map< + string, + { resolve: (vectors: number[][]) => void; reject: (error: Error) => void } + >() + const getEmbeddings = vi.fn(async (_p: string, _m: string, texts: string[]) => { + const [text] = texts + if (queryMode === 'fail') throw new Error('transport unavailable') + if (queryMode === 'pending') { + return new Promise((resolve, reject) => { + pendingQueries.set(text, { resolve, reject }) + }) + } + return texts.map((value) => textToVector(value)) + }) + const presenter = new MemoryService({ + repository: repo, + resolveAgentConfig: () => enabledConfig, + getEmbeddings, + getDimensions: embeddingDimensions, + createVectorStore: async () => store, + resetVectorStore: async () => undefined + }) + presenter.writeMemoriesSync( + [ + { kind: 'semantic', content: 'redis alpha' }, + { kind: 'semantic', content: 'redis beta' } + ], + { agentId: 'a' } + ) + await presenter.processPendingEmbeddings('a') + + queryMode = 'fail' + await presenter.recall('a', 'redis seed') + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit.failures).toBe(1) + + queryMode = 'pending' + const failingRecall = presenter.recall('a', 'redis alpha') + const lateSuccessRecall = presenter.recall('a', 'redis beta') + await flushMicrotasks(10) + expect([...pendingQueries.keys()]).toEqual(['redis alpha', 'redis beta']) + + pendingQueries.get('redis alpha')?.reject(new Error('transport unavailable')) + await failingRecall + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit.state).toBe('open') + + pendingQueries.get('redis beta')?.resolve([textToVector('redis beta')]) + await lateSuccessRecall + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit).toEqual({ + state: 'open', + failures: 2, + openCount: 1, + skipped: 0 + }) + + const callsBeforeSkip = getEmbeddings.mock.calls.length + await presenter.recall('a', 'redis alpha') + expect(getEmbeddings).toHaveBeenCalledTimes(callsBeforeSkip) + await presenter.dispose() + }) + + it('expires isolated failures outside the query embedding breaker window', async () => { + vi.useFakeTimers() + try { + vi.setSystemTime(0) + const repo = createFakeRepository() + const store = new FakeVectorStore() + let failQueries = false + const getEmbeddings = vi.fn(async (_p: string, _m: string, texts: string[]) => { + if (failQueries) throw new Error('transport unavailable') + return texts.map((text) => textToVector(text)) + }) + const presenter = new MemoryService({ + repository: repo, + resolveAgentConfig: () => enabledConfig, + getEmbeddings, + getDimensions: embeddingDimensions, + createVectorStore: async () => store, + resetVectorStore: async () => undefined + }) + const [memoryId] = presenter.writeMemoriesSync( + [{ kind: 'semantic', content: 'redis setup' }], + { agentId: 'a' } + ) + await presenter.processPendingEmbeddings('a') + const clearReady = vi.spyOn(memoryRuntimeForTests(presenter).vectorStoreService, 'clearReady') + + failQueries = true + await presenter.recall('a', 'redis setup') + vi.setSystemTime(RECALL_QUERY_EMBEDDING_BREAKER_FAILURE_WINDOW_MS + 1) + await presenter.recall('a', 'redis setup') + + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit).toEqual({ + state: 'closed', + failures: 2, + openCount: 0, + skipped: 0 + }) + expect(clearReady).not.toHaveBeenCalled() + + await presenter.recall('a', 'redis setup') + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit.state).toBe('open') + await expect(presenter.recall('a', 'redis setup')).resolves.toEqual([ + expect.objectContaining({ id: memoryId, sources: { fts: true } }) + ]) + await presenter.dispose() } finally { vi.useRealTimers() } }) + it('isolates circuits by Agent and clears stale state on model changes and cleanup', async () => { + const repo = createFakeRepository() + const modelByAgent = new Map([ + ['a', 'm'], + ['b', 'm'] + ]) + let failQueries = false + const getEmbeddings = vi.fn(async (_p: string, _m: string, texts: string[]) => { + if (failQueries) throw new Error('transport unavailable') + return texts.map((text) => textToVector(text)) + }) + const presenter = new MemoryService({ + repository: repo, + resolveAgentConfig: (agentId) => ({ + memoryEnabled: true, + memoryEmbedding: { providerId: 'p', modelId: modelByAgent.get(agentId) ?? 'm' } + }), + getEmbeddings, + getDimensions: embeddingDimensions, + createVectorStore: async () => new FakeVectorStore(), + resetVectorStore: async () => undefined + }) + const [memoryA] = presenter.writeMemoriesSync([{ kind: 'semantic', content: 'redis alpha' }], { + agentId: 'a' + }) + const [memoryB] = presenter.writeMemoriesSync([{ kind: 'semantic', content: 'redis beta' }], { + agentId: 'b' + }) + await presenter.processPendingEmbeddings('a') + await presenter.processPendingEmbeddings('b') + + failQueries = true + await presenter.recall('a', 'redis alpha') + await presenter.recall('a', 'redis alpha') + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit.state).toBe('open') + + const diagnostics = memoryRuntimeForTests(presenter).diagnosticsCollector + for (let index = 0; index < 64; index += 1) { + diagnostics.recordRecall(`diagnostic-${index}`, { + purpose: 'recall', + latencyMs: { total: 1 }, + ftsCandidates: 0, + vectorCandidates: 0, + selected: 0, + outcome: 'completed', + degradations: [] + }) + } + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit).toEqual({ + state: 'open', + failures: 0, + openCount: 0, + skipped: 0 + }) + await presenter.recall('a', 'redis alpha') + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit).toEqual({ + state: 'open', + failures: 0, + openCount: 0, + skipped: 1 + }) + + failQueries = false + await expect(presenter.recall('b', 'redis beta')).resolves.toEqual([ + expect.objectContaining({ id: memoryB, sources: { vec: true, fts: true } }) + ]) + expect(presenter.getHealth('b').runtime.agent.queryEmbeddingCircuit.state).toBe('closed') + + modelByAgent.set('a', 'm2') + presenter.onAgentMemoryMaintenanceConfigChanged('a') + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit).toEqual({ + state: 'closed', + failures: 0, + openCount: 0, + skipped: 0 + }) + + failQueries = true + await presenter.recall('b', 'redis beta') + await presenter.recall('b', 'redis beta') + expect(presenter.getHealth('b').runtime.agent.queryEmbeddingCircuit.state).toBe('open') + await presenter.cleanupDeletedAgentResources('b') + expect(presenter.getHealth('b').runtime.agent.queryEmbeddingCircuit).toEqual({ + state: 'closed', + failures: 0, + openCount: 0, + skipped: 0 + }) + expect(repo.getById(memoryA)).toBeDefined() + await presenter.dispose() + }) + + it('does not admit an old-model query after immediate configuration invalidation', async () => { + const repo = createFakeRepository() + const store = new FakeVectorStore() + let modelId = 'm' + const getEmbeddings = vi.fn(async (_p: string, _m: string, texts: string[]) => + texts.map((text) => textToVector(text)) + ) + const presenter = new MemoryService({ + repository: repo, + resolveAgentConfig: () => ({ + memoryEnabled: true, + memoryEmbedding: { providerId: 'p', modelId } + }), + getEmbeddings, + getDimensions: embeddingDimensions, + createVectorStore: async () => store, + resetVectorStore: async () => undefined + }) + presenter.writeMemoriesSync([{ kind: 'semantic', content: 'redis setup' }], { agentId: 'a' }) + await presenter.processPendingEmbeddings('a') + getEmbeddings.mockClear() + + const staleRecall = presenter.recall('a', 'redis setup') + modelId = 'm2' + presenter.onAgentMemoryMaintenanceConfigChanged('a') + + await expect(staleRecall).resolves.toEqual([]) + expect(getEmbeddings).not.toHaveBeenCalled() + await presenter.dispose() + }) + + it('does not count query embedding cancellation as a circuit failure', async () => { + const repo = createFakeRepository() + const store = new FakeVectorStore() + let cancelQueries = false + const getEmbeddings = vi.fn(async (_p: string, _m: string, texts: string[]) => { + if (cancelQueries) { + throw Object.assign(new Error('cancelled'), { + name: 'AbortError', + code: MEMORY_PROVIDER_CANCELLATION_CODE + }) + } + return texts.map((text) => textToVector(text)) + }) + const presenter = new MemoryService({ + repository: repo, + resolveAgentConfig: () => enabledConfig, + getEmbeddings, + getDimensions: embeddingDimensions, + createVectorStore: async () => store, + resetVectorStore: async () => undefined + }) + const [memoryId] = presenter.writeMemoriesSync([{ kind: 'semantic', content: 'redis setup' }], { + agentId: 'a' + }) + await presenter.processPendingEmbeddings('a') + + cancelQueries = true + await expect(presenter.recall('a', 'redis setup')).resolves.toEqual([ + expect.objectContaining({ id: memoryId, sources: { fts: true } }) + ]) + await expect(presenter.recall('a', 'redis setup')).resolves.toEqual([ + expect.objectContaining({ id: memoryId, sources: { fts: true } }) + ]) + + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit).toEqual({ + state: 'closed', + failures: 0, + openCount: 0, + skipped: 0 + }) + await presenter.dispose() + }) + + it('settles shared query health from the live consumer when another consumer cancels', async () => { + const repo = createFakeRepository() + const store = new FakeVectorStore() + let resolveQuery: ((vectors: number[][]) => void) | null = null + const getEmbeddings = vi.fn((_p: string, _m: string, texts: string[]) => { + if (texts[0] === 'redis') { + return new Promise((resolve) => { + resolveQuery = resolve + }) + } + return Promise.resolve(texts.map((text) => textToVector(text))) + }) + const presenter = new MemoryService({ + repository: repo, + resolveAgentConfig: () => enabledConfig, + getEmbeddings, + getDimensions: embeddingDimensions, + createVectorStore: async () => store, + resetVectorStore: async () => undefined + }) + presenter.writeMemoriesSync([{ kind: 'semantic', content: 'redis setup' }], { agentId: 'a' }) + await presenter.processPendingEmbeddings('a') + getEmbeddings.mockClear() + + const controller = new AbortController() + const cancelled = presenter.buildInjection('a', 'redis', { signal: controller.signal }) + const live = presenter.buildInjection('a', 'redis') + await flushMicrotasks(10) + expect(getEmbeddings).toHaveBeenCalledTimes(1) + controller.abort() + resolveQuery?.([textToVector('redis')]) + + await expect(cancelled).rejects.toMatchObject({ name: 'AbortError' }) + await expect(live).resolves.toMatchObject({ + payload: { memories: [expect.objectContaining({ content: 'redis setup' })] } + }) + expect(presenter.getHealth('a').runtime.agent.queryEmbeddingCircuit).toEqual({ + state: 'closed', + failures: 0, + openCount: 0, + skipped: 0 + }) + await presenter.dispose() + }) + it('shares identical concurrent query embeddings and returns vector hits to both callers', async () => { const repo = createFakeRepository() const store = new FakeVectorStore() diff --git a/test/main/memory/memoryDiagnosticsCollector.test.ts b/test/main/memory/memoryDiagnosticsCollector.test.ts index eaf0640226..cc257ef76e 100644 --- a/test/main/memory/memoryDiagnosticsCollector.test.ts +++ b/test/main/memory/memoryDiagnosticsCollector.test.ts @@ -152,6 +152,31 @@ describe('MemoryDiagnosticsCollector', () => { }) }) + it('tracks and resets content-free query embedding circuit diagnostics', () => { + const collector = new MemoryDiagnosticsCollector() + collector.recordQueryEmbeddingCircuitEvent('agent', 'failure') + collector.recordQueryEmbeddingCircuitEvent('agent', 'opened') + collector.recordQueryEmbeddingCircuitEvent('agent', 'skipped') + collector.recordQueryEmbeddingCircuitEvent('agent', 'halfOpen') + + expect(collector.snapshot('agent').agent.queryEmbeddingCircuit).toEqual({ + state: 'halfOpen', + failures: 1, + openCount: 1, + skipped: 1 + }) + + collector.recordQueryEmbeddingCircuitEvent('agent', 'closed') + expect(collector.snapshot('agent').agent.queryEmbeddingCircuit.state).toBe('closed') + collector.resetQueryEmbeddingCircuit('agent') + expect(collector.snapshot('agent').agent.queryEmbeddingCircuit).toEqual({ + state: 'closed', + failures: 0, + openCount: 0, + skipped: 0 + }) + }) + it('returns immutable snapshots containing only bounded diagnostic fields', () => { const collector = new MemoryDiagnosticsCollector() collector.recordRecall('agent', recallSample(10)) diff --git a/test/main/memory/serviceTestSupport.ts b/test/main/memory/serviceTestSupport.ts index fdbc5840de..f2e1bf0242 100644 --- a/test/main/memory/serviceTestSupport.ts +++ b/test/main/memory/serviceTestSupport.ts @@ -77,7 +77,7 @@ type MemoryServiceRuntimeTestSeams = { 'repairConflictIntegrity' | 'resolveConflict' | 'runChallengeResolutionPass' > maintenance: Pick - diagnostics: Pick + diagnostics: Pick } export function memoryRuntimeForTests(presenter: MemoryService) { diff --git a/test/main/routes/memoryDto.test.ts b/test/main/routes/memoryDto.test.ts index c6d37b8161..65267e9c15 100644 --- a/test/main/routes/memoryDto.test.ts +++ b/test/main/routes/memoryDto.test.ts @@ -359,6 +359,12 @@ describe('memory.getHealth route contract', () => { expect(Object.keys(runtime.agent.retrieval.recall.degradationCounts)).toEqual( MEMORY_RETRIEVAL_DEGRADATION_CAUSES ) + expect(runtime.agent.queryEmbeddingCircuit).toEqual({ + state: 'closed', + failures: 0, + openCount: 0, + skipped: 0 + }) expect(Object.keys(runtime.agent.maintenance.budgetDeniedByStep)).toEqual( MEMORY_MAINTENANCE_BUDGET_STEPS ) @@ -390,6 +396,12 @@ describe('memory.getHealth route contract', () => { runtime.agent.retrieval.recall.vectorCandidates = 5 runtime.agent.retrieval.recall.selected = 3 runtime.agent.retrieval.recall.degradationCounts.vectorCold = 2 + runtime.agent.queryEmbeddingCircuit = { + state: 'halfOpen', + failures: 3, + openCount: 1, + skipped: 4 + } runtime.agent.extraction = { chunksCompleted: 4, chunksCancelled: 2,