Files
olp/lib/cache/store.mjs
T
taodengandClaude Opus 4.7 bdfea6884b feat+docs+test: D39 — D16 follow-ups (issue #3, 4 parts)
D16 reviewer (commit `bafa6d1`) left 4 non-blocking suggestions
batched into issue #3 as a tracker. D39 closes all 4.

**Part 1 — CacheStore.delete(keyId, cacheKey) API** (lib/cache/store.mjs)

D16 originally evicted truncated entries via
`cacheStore.set(keyId, hopCacheKey, result, ttlMs=0)` — a TTL=0
tombstone purged lazily on next access. D39 introduces an explicit
delete primitive that removes the entry immediately.

- Synchronous: `delete(keyId, cacheKey) → boolean`. Returns true if
  the entry was present and removed, false if absent. Sync (not
  async) for the simplest in-memory Map contract — mirrors clear().
  Other CacheStore methods are async to leave room for a Phase 2
  file-backed adapter; delete being sync was a deliberate choice.
- Namespace cleanup: when the inner Map becomes empty after delete,
  the outer Map's per-keyId entry is also removed (memory hygiene;
  mirrors the _activeSpawns cleanup pattern from D38).
- Behavior: peek/get/getOrCompute see no trace after delete; the
  subsequent getOrCompute triggers a fresh compute.

**Part 2 — `cache_evicted_truncated` observability log** (server.mjs)

After the D16 eviction call in collectAllChunks, emit:
```js
logEvent('info', 'cache_evicted_truncated', {
  provider, model, cache_eviction_hit,
});
```

Dashboard sees salvage frequency per (provider, model). The
cache_eviction_hit boolean distinguishes "we evicted an entry" (true)
from "we tried to evict but it was already gone" (false — race with
concurrent eviction or TTL purge), preserving observability accuracy
under concurrency.

**Part 3 — Sticky-cache regression test** (test-features.mjs)

Defense-in-depth around the eviction path. Two consecutive identical
buffered requests; the first triggers SPAWN_FAILED after partial
chunks → Case B salvage returns `{ chunks..., finish_reason: 'length' }`
to the client and evicts via delete(). The second identical request
must trigger a fresh spawn (NOT serve the salvaged response from a
stale cache entry).

Asserts on BOTH invariants for defense-in-depth:
- Mock provider spawn count == 2 across the 2 identical requests
- Second request's X-OLP-Cache header is 'miss'

If eviction silently breaks in a future regression, both assertions
catch it independently.

**Part 4 — SPAWN_TIMEOUT salvage parity: DOCUMENT ASYMMETRY**
(docs/adr/0004-fallback-engine.md)

Maintainer decision: SPAWN_TIMEOUT is NOT salvaged. Document the
asymmetry rather than implementing parity. ADR 0004 Amendment 1 is
extended with a new section "Why SPAWN_TIMEOUT is excluded from
salvage" with 4-point rationale:

1. SPAWN_FAILED indicates the provider crashed mid-stream — there's
   nothing more coming; partial > nothing. Next-hop spawn has no
   advantage (same input may crash same way).
2. SPAWN_TIMEOUT indicates the provider was slow (deadline exceeded
   per `hints.maxSpawnTimeMs`). Fallback advancement to a DIFFERENT
   provider is more likely to give a complete response than salvaging
   a partial from a slow provider.
3. The "user paid for partial content" framing from D16 captures only
   SPAWN_FAILED. For SPAWN_TIMEOUT the user actually paid for "result
   within time T" — partial-at-time-T is not what was paid for;
   "full result soon after T" via fallback is closer.
4. Code-level inspection confirms the asymmetry: collectAllChunks
   catch matches ONLY `code === 'SPAWN_FAILED'` (server.mjs:563).
   SPAWN_TIMEOUT propagates via re-throw and hits evaluateHardTriggers.

v1.x re-evaluation trigger: if real usage shows users want partial-
on-timeout for very long deadlines, add a v1.x design ADR.

Stale comment fix: `lib/providers/anthropic.mjs:369` previously said
"SPAWN_TIMEOUT salvage parity is tracked in issue #3". D39 closes
that issue, so the comment is updated to point at ADR 0004 Amendment 1.

**Tests** (test-features.mjs): 447 → 452 (+5):
- 3 unit tests on CacheStore.delete (Suite 9): present-key true, absent-key
  false, namespace cleanup at empty
- 1 D16 integration test: cache_evicted_truncated log fires with
  correct fields during salvage
- 1 sticky-cache regression: spawn count 2 across 2 identical requests,
  X-OLP-Cache miss on second

Pre-commit fold-ins (per evidence-first checkpoint #4):

- **Reviewer Suggestion #1**: cacheStore.delete() return value was
  discarded at the call site → log inflated salvage metric under
  concurrent-eviction race. Folded: captured `evicted` boolean and
  added to log payload as `cache_eviction_hit`.
- **Reviewer Suggestion #2**: anthropic.mjs:369 stale comment
  pointing at now-closed issue #3. Folded: rewrote to point at
  ADR 0004 Amendment 1 § "Why SPAWN_TIMEOUT is excluded from
  salvage".
- **Reviewer Suggestion #3**: ADR 0004 attribution ambiguity —
  parenthetical "(per Amendment 3 — SPAWN_TIMEOUT is one of the 4
  live hard-trigger codes alongside SPAWN_FAILED, CLI_NOT_FOUND, and
  CONCURRENCY_LIMIT from Amendment 4)" could mis-parse as Amendment 3
  covering all four. Folded: split to
  "(per Amendment 3: SPAWN_FAILED, CLI_NOT_FOUND, SPAWN_TIMEOUT;
  per Amendment 4: CONCURRENCY_LIMIT)".

**CHANGELOG**: D39 sub-entry appended under the existing D38 entry
in Unreleased section. No package.json bump (phase_rolling_mode).

Authority:
- ADR 0005 § Cache layer — CacheStore API extension (Part 1)
- ADR 0004 Amendment 1 update — SPAWN_TIMEOUT asymmetry rationale (Part 4)
- GitHub issue #3 — closed by this commit
- D16 commit bafa6d1 non-blocking suggestions — batched here
- CC 开发铁律 v1.6 § 10.x — fresh-context opus reviewer independent
- CLAUDE.md release_kit_overlay phase_rolling_mode — under Unreleased

Reviewer (fresh-context opus, Iron Rule 10): APPROVE_WITH_MINOR.
Verified: delete() sync + call-site no-await correct; namespace
cleanup guarded on (had && ns.size === 0); SPAWN_TIMEOUT NOT in
salvage catch; ADR section 4-point rationale internally consistent
with code state; hygiene clean; 452/452 tests pass independently.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-05-24 21:13:57 +10:00

399 lines
14 KiB
JavaScript

/**
* lib/cache/store.mjs — In-process cache store with per-key isolation + singleflight
*
* Authority: ADR 0005 § D1 (per-key isolation) + D4 (singleflight)
*
* This module implements the OLP response cache. At D5 the backing store is
* an in-memory Map (not file-backed). File backing lands in a later Phase.
* The architecture mirrors OCP v3.13.0 (D1+D4) generalized to OLP's
* (provider, model) keyed cache key space.
*
* D1 — Per-key isolation:
* Outer Map: keyId → inner Map
* Inner Map: cacheKey (SHA-256 hex) → CacheEntry
*
* keyId is the OLP API key identifier (Phase 2 multi-key). At D5, the server
* passes '__anonymous__' as keyId for all requests.
*
* D4 — Singleflight:
* Concurrent requests with the same (keyId, cacheKey) share one compute
* invocation. The first caller triggers computeFn(); subsequent callers
* await the same Promise. After completion, all callers receive the same
* value, and the result is cached for future callers.
*
* D2/D3 are placeholders at D5:
* D2 (cache_control bypass) — checked in server.mjs before entering getOrCompute.
* D3 (chunked stream replay with timing fidelity) — at D5 the cache stores
* full collected chunks and replays them sequentially without timing.
* Full D3 (timing-accurate replay) lands in a later Phase.
*
* Per OCP learning (MEMORY.md learnings/concurrency_dedup_test_signals.md):
* Singleflight is verified by observing that N concurrent requests for the
* same key return within ~0ms of each other (not by counting spawn log lines).
* The content-hash for the cache key must NOT include randomUUID or any
* per-request random component — only the semantic content that affects output.
*/
// ── Types ─────────────────────────────────────────────────────────────────
/**
* @typedef {Object} CacheEntry
* @property {*} value - the cached response (IR chunk array for streaming, response object for non-streaming)
* @property {number} createdAt - epoch ms (Date.now()) when the entry was created
* @property {number} ttlMs - time-to-live in milliseconds
*/
/**
* @typedef {Object} CacheStats
* @property {number} hits
* @property {number} misses
* @property {number} size
* @property {number} inflightCount
*/
// ── CacheStore ────────────────────────────────────────────────────────────
export class CacheStore {
/**
* @param {object} [opts]
* @param {number} [opts.maxEntriesPerKey=1000] — evict oldest entries when exceeded
* @param {number} [opts.maxAgeMs=86400000] — default TTL: 24 hours
* @param {number} [opts.maxEntryBytes=10485760] — max serialized size per entry (default 10 MB).
* Entries whose JSON.stringify serialization exceeds this limit are not persisted.
* ADR 0005 § "Cache write conditions" item 4 (D23). Cache is for hot-path repeat
* requests, not bulk archive. Configurable so tests can exercise the cap cheaply.
* @param {() => number} [opts._nowFn=Date.now] — injectable for TTL testing
* @param {(msg: string, meta?: object) => void} [opts._warnFn] — injectable for warn testing
*/
constructor(opts = {}) {
this._maxEntriesPerKey = opts.maxEntriesPerKey ?? 1000;
this._maxAgeMs = opts.maxAgeMs ?? 24 * 60 * 60 * 1000; // 24 hours
// ADR 0005 § "Cache write conditions" item 4 (D23): default 10 MB = 10 * 1024 * 1024 bytes.
this._maxEntryBytes = opts.maxEntryBytes ?? 10 * 1024 * 1024;
this._nowFn = opts._nowFn ?? (() => Date.now());
// Injectable warn function for testing (defaults to console.warn).
this._warnFn = opts._warnFn ?? ((msg, meta) => console.warn(msg, meta ?? ''));
// D1 per-key isolation: keyId → Map<cacheKey, CacheEntry>
/** @type {Map<string, Map<string, CacheEntry>>} */
this._store = new Map();
// D4 singleflight: `${keyId}:${cacheKey}` → Promise<value>
/** @type {Map<string, Promise<*>>} */
this._inflight = new Map();
// Stats per keyId
/** @type {Map<string, { hits: number, misses: number }>} */
this._stats = new Map();
}
// ── Internal helpers ──────────────────────────────────────────────────
/**
* Returns or creates the inner Map for a given keyId.
* @param {string} keyId
* @returns {Map<string, CacheEntry>}
*/
_getNamespace(keyId) {
if (!this._store.has(keyId)) {
this._store.set(keyId, new Map());
}
return this._store.get(keyId);
}
/**
* Returns or creates the stats object for a given keyId.
* @param {string} keyId
* @returns {{ hits: number, misses: number }}
*/
_getStats(keyId) {
if (!this._stats.has(keyId)) {
this._stats.set(keyId, { hits: 0, misses: 0 });
}
return this._stats.get(keyId);
}
/**
* Returns true if entry has not yet expired.
* @param {CacheEntry} entry
* @returns {boolean}
*/
_isAlive(entry) {
const age = this._nowFn() - entry.createdAt;
return age < entry.ttlMs;
}
/**
* Evicts oldest entries in a namespace when it exceeds maxEntriesPerKey.
* @param {Map<string, CacheEntry>} ns
*/
_evictIfNeeded(ns) {
if (ns.size <= this._maxEntriesPerKey) return;
// Collect entries sorted by createdAt ascending (oldest first)
const entries = [...ns.entries()].sort((a, b) => a[1].createdAt - b[1].createdAt);
const toEvict = ns.size - this._maxEntriesPerKey;
for (let i = 0; i < toEvict; i++) {
ns.delete(entries[i][0]);
}
}
// ── Public API ────────────────────────────────────────────────────────
/**
* Retrieves a cached entry. Returns null on miss or TTL expiry.
* Updates stats.
*
* @param {string} keyId - OLP API key namespace (use '__anonymous__' at D5)
* @param {string} cacheKey - SHA-256 hex from computeCacheKey()
* @returns {Promise<CacheEntry|null>}
*/
async get(keyId, cacheKey) {
const ns = this._getNamespace(keyId);
const entry = ns.get(cacheKey);
if (!entry) {
this._getStats(keyId).misses++;
return null;
}
if (!this._isAlive(entry)) {
ns.delete(cacheKey);
this._getStats(keyId).misses++;
return null;
}
this._getStats(keyId).hits++;
return entry;
}
/**
* Stores a value in the cache.
*
* ADR 0005 § "Cache write conditions" item 4 (D23): if the serialized size of `value`
* exceeds `maxEntryBytes`, the entry is NOT persisted. A `cache_skip_oversize` warn
* event is logged. The caller still receives the value (this method returns undefined
* either way); only persistent caching is skipped.
*
* @param {string} keyId
* @param {string} cacheKey
* @param {*} value
* @param {number} [ttlMs] - overrides default maxAgeMs
* @returns {Promise<void>}
*/
async set(keyId, cacheKey, value, ttlMs) {
// ADR 0005 § "Cache write conditions" item 4 (D23): size cap check.
const byteLength = Buffer.byteLength(JSON.stringify(value));
if (byteLength > this._maxEntryBytes) {
this._warnFn('cache_skip_oversize', {
byteLength,
maxEntryBytes: this._maxEntryBytes,
keyId,
cacheKey,
});
return; // Do not persist; caller still receives the value from computeFn.
}
const ns = this._getNamespace(keyId);
const entry = {
value,
createdAt: this._nowFn(),
ttlMs: ttlMs ?? this._maxAgeMs,
};
ns.set(cacheKey, entry);
this._evictIfNeeded(ns);
}
/**
* Returns true if a valid (non-expired) entry exists for (keyId, cacheKey).
*
* Stats side-effect: this method calls get(), which increments hit/miss
* counters. Callers that need to peek WITHOUT touching stats should use
* peek() instead. Server-side cache-status reporting MUST use peek() to
* avoid double-counting when getOrCompute() is called immediately after.
*
* @param {string} keyId
* @param {string} cacheKey
* @returns {Promise<boolean>}
*/
async has(keyId, cacheKey) {
const entry = await this.get(keyId, cacheKey);
return entry !== null;
}
/**
* Stats-neutral cache peek. Returns true if a valid entry exists without
* incrementing hit/miss counters. Use this from server.mjs for pre-check
* cache-status reporting; the subsequent getOrCompute() call is the one
* that increments the canonical stats counter at the semantic boundary.
*
* Fixes the codex-flagged double-count bug where has() + getOrCompute()
* both incremented the same counter on the same logical request.
*
* @param {string} keyId
* @param {string} cacheKey
* @returns {Promise<boolean>}
*/
async peek(keyId, cacheKey) {
const ns = this._getNamespace(keyId);
const entry = ns.get(cacheKey);
if (!entry) return false;
if (!this._isAlive(entry)) {
ns.delete(cacheKey);
return false;
}
return true;
}
/**
* D4 singleflight: get from cache or compute exactly once for concurrent callers.
*
* - Cache hit: return cached value immediately (no computeFn invocation).
* - Inflight hit: await the existing inflight Promise (singleflight — computeFn
* runs exactly once for N concurrent callers with the same key).
* - Miss: register inflight promise, invoke computeFn(), cache the result,
* release the inflight promise, return the value.
* - On computeFn error: release the inflight promise (so subsequent calls retry),
* rethrow the error to all awaiting callers.
*
* Per OCP MEMORY.md learning (concurrency_dedup_test_signals.md):
* Singleflight is verified by timing: N concurrent calls return within ~ms
* of each other, and computeFn is called exactly once.
*
* @param {string} keyId
* @param {string} cacheKey
* @param {() => Promise<*>} computeFn - async function producing the value
* @param {number} [ttlMs]
* @returns {Promise<*>}
*/
async getOrCompute(keyId, cacheKey, computeFn, ttlMs) {
// 1. Cache hit — return immediately, no singleflight overhead
const ns = this._getNamespace(keyId);
const existing = ns.get(cacheKey);
if (existing && this._isAlive(existing)) {
this._getStats(keyId).hits++;
return existing.value;
}
// 2. Inflight hit — singleflight: await the in-progress compute
const inflightKey = `${keyId}:${cacheKey}`;
const inFlight = this._inflight.get(inflightKey);
if (inFlight) {
// Another caller is already computing this key. Await their result.
// The first caller will cache the result; we return the same value.
return inFlight;
}
// 3. Miss — we are the first caller; register inflight promise
const computePromise = (async () => {
try {
const value = await computeFn();
// Cache the result so future calls are cache hits
await this.set(keyId, cacheKey, value, ttlMs);
this._getStats(keyId).misses++;
return value;
} finally {
// Always release the inflight slot, even on error
this._inflight.delete(inflightKey);
}
})();
this._inflight.set(inflightKey, computePromise);
return computePromise;
}
/**
* Returns stats for a specific keyId, or aggregate stats across all keyIds.
*
* @param {string} [keyId] - if omitted, returns aggregate across all keys
* @returns {CacheStats}
*/
stats(keyId) {
if (keyId !== undefined) {
const s = this._stats.get(keyId) ?? { hits: 0, misses: 0 };
const ns = this._store.get(keyId);
const size = ns ? ns.size : 0;
return {
hits: s.hits,
misses: s.misses,
size,
inflightCount: this._inflight.size,
};
}
// Aggregate
let totalHits = 0;
let totalMisses = 0;
let totalSize = 0;
for (const s of this._stats.values()) {
totalHits += s.hits;
totalMisses += s.misses;
}
for (const ns of this._store.values()) {
totalSize += ns.size;
}
return {
hits: totalHits,
misses: totalMisses,
size: totalSize,
inflightCount: this._inflight.size,
};
}
/**
* Removes a specific (keyId, cacheKey) entry immediately.
*
* ADR 0005 § "Cache write conditions" item 1 (D39, issue #3 Part 1):
* D16 truncation-eviction previously used `set(..., ttlMs=0)` to leave a
* tombstone that the next `get`/`peek` would lazily purge. That pattern
* left dead entries in the namespace Map until next access, accruing
* memory if no follow-up read ever fires. `delete(keyId, cacheKey)` makes
* the eviction explicit and immediate.
*
* Memory hygiene: if the per-keyId namespace becomes empty after delete,
* the namespace Map entry itself is removed (mirrors the pattern in D38
* `_activeSpawns` so empty namespaces don't accumulate in `_store`).
*
* Stats: this method does NOT touch hit/miss counters — it is an eviction
* primitive, not a read. Aggregate `size` reported by `stats()` reflects
* the removal on the next call.
*
* @param {string} keyId
* @param {string} cacheKey
* @returns {boolean} true if the entry was present and removed; false if absent.
*/
delete(keyId, cacheKey) {
const ns = this._store.get(keyId);
if (!ns) return false;
const had = ns.delete(cacheKey);
// Memory hygiene: drop empty namespace Map entries.
if (had && ns.size === 0) {
this._store.delete(keyId);
}
return had;
}
/**
* Clears cache entries (and stats) for a specific keyId, or ALL entries.
*
* @param {string} [keyId] - if omitted, clears all namespaces
*/
clear(keyId) {
if (keyId !== undefined) {
this._store.delete(keyId);
this._stats.delete(keyId);
// Also remove any inflight entries for this keyId
for (const k of this._inflight.keys()) {
if (k.startsWith(`${keyId}:`)) {
this._inflight.delete(k);
}
}
} else {
this._store.clear();
this._stats.clear();
this._inflight.clear();
}
}
}