Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
207 changes: 117 additions & 90 deletions packages/core/src/cachekeep.ts
Original file line number Diff line number Diff line change
Expand Up @@ -336,6 +336,10 @@ export type CacheKeepTarget = {
isSubagent: boolean
}

export type CacheKeepPrewarmAttempt = {
id: number
}

function cacheKeepRetryDelayMs(targetId: string, failureCount: number) {
const exponential = Math.min(
CACHE_KEEP_RETRY_MAX_MS,
Expand Down Expand Up @@ -367,6 +371,7 @@ export class CacheKeepManager {
private readonly targets = new Map<string, CacheKeepTarget>()
private timer: ReturnType<typeof setInterval> | null = null
private tickPromise: Promise<void> | null = null
private nextPrewarmAttemptId = 0

constructor(
private readonly options: {
Expand All @@ -379,7 +384,8 @@ export class CacheKeepManager {
prepareHeaders?: (
headers: Headers,
target: CacheKeepTarget,
) => Promise<Headers> | Headers
attempt: CacheKeepPrewarmAttempt,
) => Promise<Headers | undefined> | Headers | undefined
onTrackedSessionsChanged?: (
sessions: readonly CacheKeepTrackedSession[],
) => Promise<void> | void
Expand All @@ -393,6 +399,11 @@ export class CacheKeepManager {
status: number
data: unknown
receivedAt: number
attempt: CacheKeepPrewarmAttempt
}) => void | Promise<void>
onComplete?: (input: {
target: CacheKeepTarget
attempt: CacheKeepPrewarmAttempt
}) => void | Promise<void>
},
) {}
Expand Down Expand Up @@ -639,103 +650,119 @@ export class CacheKeepManager {
private async sendPrewarm(
target: CacheKeepTarget,
): Promise<CacheKeepPrewarmResult> {
let bodyText = target.bodyText
if (this.options.prepareBody) {
const attempt = { id: ++this.nextPrewarmAttemptId }
try {
let bodyText = target.bodyText
if (this.options.prepareBody) {
try {
bodyText = await this.options.prepareBody(bodyText, target)
} catch (error) {
logger.warn('cachekeep', 'prepare body failed', {
session: target.id,
error: error instanceof Error ? error.message : String(error),
})
}
}
const preparedTarget = { ...target, bodyText }
const prewarm = await buildCacheKeepPrewarmBody(bodyText)
if (!prewarm.ok) return prewarm

const fetchImpl = this.options.fetchImpl ?? fetch
const prewarmTarget = { ...preparedTarget, bodyText: prewarm.bodyText }
const headers = this.options.prepareHeaders
? await this.options.prepareHeaders(
new Headers(target.headers),
prewarmTarget,
attempt,
)
: new Headers(target.headers)
if (!headers) {
return {
ok: false,
reason: 'OAuth cache prewarm credential is unavailable',
transient: true,
}
}
headers.delete('content-length')
headers.delete('transfer-encoding')
let response: Response
try {
bodyText = await this.options.prepareBody(bodyText, target)
response = await fetchImpl(target.url, {
method: 'POST',
headers,
body: prewarm.bodyText,
signal: AbortSignal.timeout(
this.options.prewarmTimeoutMs ?? CACHE_KEEP_PREWARM_TIMEOUT_MS,
),
})
} catch (error) {
logger.warn('cachekeep', 'prepare body failed', {
return {
ok: false,
reason: error instanceof Error ? error.message : String(error),
transient: true,
}
}
const receivedAt = this.options.now?.() ?? Date.now()
const raw = await response.text().catch(() => '')
let data: unknown = null
try {
data = raw ? JSON.parse(raw) : null
} catch {}
try {
await this.options.onResponse?.({
target,
bodyText: prewarm.bodyText,
status: response.status,
data,
receivedAt,
attempt,
})
} catch {}
try {
const dumpHandle = await dumpDirectRequest({
affinity: target.id,
route: 'cachekeep',
status: response.status,
bodyText: prewarm.bodyText,
url: target.url,
method: 'POST',
headers,
tag: 'cachekeep',
})
await dumpResponseArtifact(dumpHandle, {
status: response.status,
message: data,
})
} catch (error) {
logger.debug('cachekeep', 'dump failed', {
session: target.id,
error: error instanceof Error ? error.message : String(error),
})
}
}
const preparedTarget = { ...target, bodyText }
const prewarm = await buildCacheKeepPrewarmBody(bodyText)
if (!prewarm.ok) return prewarm

const fetchImpl = this.options.fetchImpl ?? fetch
const prewarmTarget = { ...preparedTarget, bodyText: prewarm.bodyText }
const headers = this.options.prepareHeaders
? await this.options.prepareHeaders(
new Headers(target.headers),
prewarmTarget,
)
: new Headers(target.headers)
headers.delete('content-length')
headers.delete('transfer-encoding')
let response: Response
try {
response = await fetchImpl(target.url, {
method: 'POST',
headers,
body: prewarm.bodyText,
signal: AbortSignal.timeout(
this.options.prewarmTimeoutMs ?? CACHE_KEEP_PREWARM_TIMEOUT_MS,
),
})
} catch (error) {
return {
ok: false,
reason: error instanceof Error ? error.message : String(error),
transient: true,
}
}
const receivedAt = this.options.now?.() ?? Date.now()
const raw = await response.text().catch(() => '')
let data: unknown = null
try {
data = raw ? JSON.parse(raw) : null
} catch {}
try {
await this.options.onResponse?.({
target,
bodyText: prewarm.bodyText,
status: response.status,
data,
receivedAt,
})
} catch {}
try {
const dumpHandle = await dumpDirectRequest({
affinity: target.id,
route: 'cachekeep',
status: response.status,
bodyText: prewarm.bodyText,
url: target.url,
method: 'POST',
headers,
tag: 'cachekeep',
})
await dumpResponseArtifact(dumpHandle, {
status: response.status,
message: data,
})
} catch (error) {
logger.debug('cachekeep', 'dump failed', {
session: target.id,
error: error instanceof Error ? error.message : String(error),
})
}
if (!response.ok) {
return {
ok: false,
reason: raw || `HTTP ${response.status}`,
status: response.status,
if (!response.ok) {
return {
ok: false,
reason: raw || `HTTP ${response.status}`,
status: response.status,
}
}
const objectData =
data && typeof data === 'object' && !Array.isArray(data)
? (data as Record<string, unknown>)
: null
const usage = objectData?.usage as
| {
input_tokens?: number
cache_creation_input_tokens?: number
cache_read_input_tokens?: number
}
| undefined
return { ok: true, ...(usage && { usage }) }
} finally {
try {
await this.options.onComplete?.({ target, attempt })
} catch {}
}
const objectData =
data && typeof data === 'object' && !Array.isArray(data)
? (data as Record<string, unknown>)
: null
const usage = objectData?.usage as
| {
input_tokens?: number
cache_creation_input_tokens?: number
cache_read_input_tokens?: number
}
| undefined
return { ok: true, ...(usage && { usage }) }
}

private async prewarm(target: CacheKeepTarget, now: number) {
Expand Down
28 changes: 15 additions & 13 deletions packages/core/src/prime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ export type PrimeAccountStatus = {
label: string
nextDueAt?: number
lastPrimedAt?: number | null
lastResult?: 'ok' | 'error'
lastResult?: 'ok' | 'error' | 'skipped'
usage?: import('./accounts.ts').PrimeUsageCounters
estimatedCostUsd?: number
}
Expand All @@ -70,13 +70,9 @@ export type PrimeSendResult =
status?: number
ms?: number
error: string
// Discriminates the failure kind so the manager can emit the
// spec Logging table's two distinct warn events: `prime token
// refresh failed` (warn · prime · { account, error }) for
// token-refresh failures, `prime fire failed` (warn · prime ·
// { account, status?, error }) for HTTP / fetch / identity
// failures during the request itself.
reason?: 'token-refresh' | 'send'
// `vault-cold` is an expected off-path cache warm state, not evidence
// of a failed provider request, so it must not emit the fire-failure warn.
reason?: 'token-refresh' | 'vault-cold' | 'send'
}

/**
Expand Down Expand Up @@ -137,7 +133,7 @@ export function buildPrimeAccountStatuses(
now?: number
transient?: ReadonlyMap<
string,
{ lastPrimedAt?: number; lastResult?: 'ok' | 'error' }
{ lastPrimedAt?: number; lastResult?: 'ok' | 'error' | 'skipped' }
>
},
): PrimeAccountStatus[] {
Expand Down Expand Up @@ -226,14 +222,15 @@ function formatDue(nextDueAt: number | null | undefined): string {

function formatPrimed(
lastPrimedAt: number | null | undefined,
lastResult: 'ok' | 'error' | undefined,
lastResult: 'ok' | 'error' | 'skipped' | undefined,
): string {
if (typeof lastPrimedAt !== 'number') return ''
const time = new Date(lastPrimedAt).toLocaleTimeString([], {
hour: '2-digit',
minute: '2-digit',
})
if (lastResult === 'error') return `primed ${time} err`
if (lastResult === 'skipped') return `primed ${time} skipped`
return `primed ${time} \u2713`
}

Expand Down Expand Up @@ -566,7 +563,7 @@ export class PrimeManager {
// (M4): skip is not the same as a successful prime.
private transient = new Map<
string,
{ lastPrimedAt: number; lastResult: 'ok' | 'error' }
{ lastPrimedAt: number; lastResult: 'ok' | 'error' | 'skipped' }
>()
// Latest cumulative counters returned by recordSuccess. Overlays the
// persisted counters in stats() before the next load persists them.
Expand Down Expand Up @@ -946,7 +943,7 @@ export class PrimeManager {
if (this.stopped) return

if (!result.ok) {
// Spec Logging table: two distinct warn events.
// `vault-cold` is a cache-warm skip; it is not a provider failure.
// - `prime token refresh failed` — token-refresh failure
// before the request fires (reason: 'token-refresh').
// - `prime fire failed` — HTTP error / fetch throw /
Expand All @@ -956,6 +953,11 @@ export class PrimeManager {
account: evaluation.label,
error: result.error,
})
} else if (result.reason === 'vault-cold') {
logger.debug('prime', 'prime vault credential cold', {
account: evaluation.label,
reason: result.reason,
})
} else {
Comment thread
cubic-dev-ai[bot] marked this conversation as resolved.
logger.warn('prime', 'prime fire failed', {
account: evaluation.label,
Expand All @@ -965,7 +967,7 @@ export class PrimeManager {
}
this.transient.set(evaluation.id, {
lastPrimedAt: now,
lastResult: 'error',
lastResult: result.reason === 'vault-cold' ? 'skipped' : 'error',
})
return
}
Expand Down
41 changes: 41 additions & 0 deletions packages/core/src/tests/prime.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1539,6 +1539,47 @@ describe('PrimeManager — send failure', () => {
expect(h.sendCalls).toEqual([])
await h.cleanup()
})

test('cold vault prime is recorded as a skipped attempt, not an error', async () => {
const fixture = makePrimeFixture({
mainQuota: {
five_hour: {
usedPercent: 0,
remainingPercent: 100,
resetsAt: new Date(500).toISOString(),
checkedAt: 1,
},
},
})
const now = 500 + 120_000
const h = await makeHarness({
storage: fixture.storage,
markerDir: markerRoot,
now,
quotaFresh: {
five_hour: {
usedPercent: 0,
remainingPercent: 100,
resetsAt: new Date(now - 1000).toISOString(),
checkedAt: 1,
},
},
send: () => ({
ok: false,
reason: 'vault-cold',
error: 'vault credential is unavailable',
}),
})

await h.manager.tick()

const stats = h.manager.stats()
expect(stats[0]?.lastResult).toBe('skipped')
const summary = buildPrimeStatusSummary({ enabled: true, accounts: stats })
expect(summary).toContain('skipped')
expect(summary).not.toContain('err')
await h.cleanup()
})
})

describe('PrimeManager — recordSuccess', () => {
Expand Down
Loading