mirror of
https://github.com/stablyai/orca.git
synced 2026-09-28 16:02:45 +00:00
fix(relay): keep failed rehome polls out of the durable failure budget
The regional rehome worker polls claimRegionalRehome about once a second. Any error thrown before an attempt was claimed - in practice a director pool timeout on the pre-claim control read, 52-74 a day against a pool of 3 - was charged to relay_region_rehome_worker_state.consecutive_failures, which durably disables the control at three. That counter only ever resets on a drain receipt, so while the control is disabled it never resets: production sits at 1068 and still climbing. Enabling the control leaves the stale counter in place, so the next pool timeout latches it straight back off. That is what ended the 2026-08-28 enable after ten minutes. - A poll that never claimed an attempt drained nothing, so it no longer feeds the dispatch-failure budget and logs .._poll_failed instead of .._dispatch_failed. recordRegionalRehomeWorkerFailure had no other caller and is removed. - Enabling the control clears consecutive_failures and paused_until, so a budget spent under a previous enable cannot kill a fresh one. The dispatch interval in next_dispatch_at is deliberately left alone. - The budget's auto-disable now emits orca_relay_regional_rehome_failure_budget_disabled, matching the existing .._safety_disabled precedent. It wrote no event before, which is why this went unnoticed for two weeks. No change to region selection, the candidate query, or host eligibility.
This commit is contained in:
@@ -362,6 +362,9 @@ export const REGIONAL_REHOME_QUARANTINE_FAILURES = 3
|
||||
export const REGIONAL_REHOME_QUARANTINE_MS = 15 * 60_000
|
||||
const REGIONAL_REHOME_QUARANTINE_EXCLUSION_LIMIT = 50
|
||||
const REGIONAL_REHOME_QUARANTINE_MEMORY_LIMIT = 1_000
|
||||
// Consecutive drain-dispatch failures that latch the durable control off. Any
|
||||
// drain receipt and any enable reset it, so it reads "dispatch is broken now".
|
||||
const REGIONAL_REHOME_FAILURE_BUDGET = 3
|
||||
const REGIONAL_REHOME_OBSERVATION_MS = 24 * 60 * 60_000
|
||||
const ASSIGNMENT_LOCK_RETRY_MAX_DELAY_MS = 50
|
||||
type AssignmentInventoryScope = 'none' | 'general' | 'all'
|
||||
@@ -4992,6 +4995,20 @@ export class RelayAssignmentStore {
|
||||
now
|
||||
]
|
||||
)
|
||||
if (input.enabled) {
|
||||
// A budget spent under a previous enable is not evidence about this one.
|
||||
// Without this an old counter latches the fresh enable straight back off
|
||||
// on its first transient failure.
|
||||
await transaction.query(
|
||||
`INSERT INTO relay_region_rehome_worker_state
|
||||
(worker_id, next_dispatch_at, paused_until, consecutive_failures, updated_at)
|
||||
VALUES ('global', 0, 0, 0, ?)
|
||||
ON CONFLICT (worker_id) DO UPDATE
|
||||
SET paused_until = 0, consecutive_failures = 0,
|
||||
updated_at = excluded.updated_at`,
|
||||
[now]
|
||||
)
|
||||
}
|
||||
const updated = (
|
||||
await transaction.query(
|
||||
`SELECT * FROM relay_region_rehome_control WHERE control_id = 'global'`
|
||||
@@ -5940,7 +5957,7 @@ export class RelayAssignmentStore {
|
||||
|
||||
async recordRegionalRehomeDispatchFailure(attemptId: string): Promise<void> {
|
||||
const now = this.now()
|
||||
await this.database.transaction(async (transaction) => {
|
||||
const disableLog = await this.database.transaction(async (transaction) => {
|
||||
const worker = (
|
||||
await transaction.queryLocked(
|
||||
`SELECT * FROM relay_region_rehome_worker_state WHERE worker_id = 'global'`
|
||||
@@ -5952,49 +5969,44 @@ export class RelayAssignmentStore {
|
||||
[attemptId]
|
||||
)
|
||||
)[0]
|
||||
if (!worker || !attempt) return
|
||||
await this.incrementRegionalRehomeWorkerFailure(transaction, worker, now)
|
||||
})
|
||||
}
|
||||
|
||||
async recordRegionalRehomeWorkerFailure(): Promise<void> {
|
||||
const now = this.now()
|
||||
await this.database.transaction(async (transaction) => {
|
||||
await transaction.query(
|
||||
`INSERT INTO relay_region_rehome_worker_state
|
||||
(worker_id, next_dispatch_at, paused_until, consecutive_failures, updated_at)
|
||||
VALUES ('global', 0, 0, 0, ?)
|
||||
ON CONFLICT (worker_id) DO NOTHING`,
|
||||
[now]
|
||||
)
|
||||
const worker = (
|
||||
await transaction.queryLocked(
|
||||
`SELECT * FROM relay_region_rehome_worker_state WHERE worker_id = 'global'`
|
||||
)
|
||||
)[0]!
|
||||
await this.incrementRegionalRehomeWorkerFailure(transaction, worker, now)
|
||||
if (!worker || !attempt) return null
|
||||
return await this.incrementRegionalRehomeWorkerFailure(transaction, worker, now)
|
||||
})
|
||||
// Logged after the commit so a rollback cannot fabricate the record.
|
||||
if (disableLog) console.warn(JSON.stringify(disableLog))
|
||||
}
|
||||
|
||||
// Returns the durable disable this failure caused, for the caller to log once
|
||||
// its transaction commits; null when the budget survives or was already spent.
|
||||
private async incrementRegionalRehomeWorkerFailure(
|
||||
transaction: RelayDatabase,
|
||||
worker: SqlRow,
|
||||
now: number
|
||||
): Promise<void> {
|
||||
): Promise<Record<string, string | number> | null> {
|
||||
const failures = integer(worker, 'consecutive_failures') + 1
|
||||
const spent = failures >= REGIONAL_REHOME_FAILURE_BUDGET
|
||||
await transaction.query(
|
||||
`UPDATE relay_region_rehome_worker_state
|
||||
SET consecutive_failures = ?, paused_until = ?, updated_at = ?
|
||||
WHERE worker_id = 'global'`,
|
||||
[failures, failures >= 3 ? now + 5 * 60_000 : 0, now]
|
||||
[failures, spent ? now + 5 * 60_000 : 0, now]
|
||||
)
|
||||
if (failures >= 3) {
|
||||
await transaction.query(
|
||||
`UPDATE relay_region_rehome_control
|
||||
SET generation = generation + 1, enabled = 0, updated_at = ?
|
||||
WHERE control_id = 'global' AND enabled = 1`,
|
||||
[now]
|
||||
)
|
||||
if (!spent) return null
|
||||
const disabled = await transaction.query(
|
||||
`UPDATE relay_region_rehome_control
|
||||
SET generation = generation + 1, enabled = 0, updated_at = ?
|
||||
WHERE control_id = 'global' AND enabled = 1
|
||||
RETURNING generation`,
|
||||
[now]
|
||||
)
|
||||
// The disable is otherwise invisible: inspection only shows enabled=false and
|
||||
// nothing records that the failure budget, not an operator, turned it off.
|
||||
if (disabled.length === 0) return null
|
||||
return {
|
||||
event: 'orca_relay_regional_rehome_failure_budget_disabled',
|
||||
controlGeneration: integer(disabled[0]!, 'generation'),
|
||||
consecutiveFailures: failures,
|
||||
now
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1709,6 +1709,77 @@ describe('regional rehome assignment state', () => {
|
||||
expect(await context.store.claimRegionalRehome()).toBeNull()
|
||||
await context.database.close()
|
||||
})
|
||||
|
||||
it('clears a stale failure budget when the control is enabled again', async () => {
|
||||
const context = await setup()
|
||||
await activatePreferredSource(context, {
|
||||
userId: 'user-1',
|
||||
relayHostId: 'abcdefghijklmnop'
|
||||
})
|
||||
const attempt = await context.store.claimRegionalRehome()
|
||||
for (let index = 0; index < 3; index++) {
|
||||
await context.store.recordRegionalRehomeDispatchFailure(attempt!.attemptId)
|
||||
}
|
||||
expect(await workerState(context)).toMatchObject({ consecutiveFailures: 3 })
|
||||
const latched = await context.store.inspectRegionalRehomeControl()
|
||||
expect(latched).toMatchObject({ generation: 2, enabled: false })
|
||||
|
||||
await context.store.applyRegionalRehomeControl({
|
||||
expectedGeneration: latched.generation,
|
||||
enabled: true,
|
||||
notBefore: context.now(),
|
||||
ratePerMinute: 10,
|
||||
preferenceMaxAgeMs: 24 * 60 * 60_000,
|
||||
hostCooldownMs: 7 * 24 * 60 * 60_000,
|
||||
drainGraceMs: 60 * 60_000
|
||||
})
|
||||
|
||||
// A budget spent under the previous enable is not evidence about this one.
|
||||
expect(await workerState(context)).toMatchObject({
|
||||
consecutiveFailures: 0,
|
||||
pausedUntil: 0
|
||||
})
|
||||
// One transient failure must not latch the fresh enable straight back off.
|
||||
await context.store.recordRegionalRehomeDispatchFailure(attempt!.attemptId)
|
||||
expect(await context.store.inspectRegionalRehomeControl()).toMatchObject({
|
||||
generation: 3,
|
||||
enabled: true
|
||||
})
|
||||
await context.database.close()
|
||||
})
|
||||
|
||||
it('reports the durable disable when the failure budget latches the control off', async () => {
|
||||
const context = await setup()
|
||||
await activatePreferredSource(context, {
|
||||
userId: 'user-1',
|
||||
relayHostId: 'abcdefghijklmnop'
|
||||
})
|
||||
const attempt = await context.store.claimRegionalRehome()
|
||||
const warnings = collectEventWarnings(
|
||||
'orca_relay_regional_rehome_failure_budget_disabled'
|
||||
)
|
||||
try {
|
||||
for (let index = 0; index < 5; index++) {
|
||||
await context.store.recordRegionalRehomeDispatchFailure(attempt!.attemptId)
|
||||
}
|
||||
} finally {
|
||||
warnings.restore()
|
||||
}
|
||||
|
||||
// Only the transition is reported; later failures find the control already off.
|
||||
expect(warnings.entries).toEqual([
|
||||
expect.objectContaining({
|
||||
event: 'orca_relay_regional_rehome_failure_budget_disabled',
|
||||
controlGeneration: 2,
|
||||
consecutiveFailures: 3
|
||||
})
|
||||
])
|
||||
expect(await context.store.inspectRegionalRehomeControl()).toMatchObject({
|
||||
generation: 2,
|
||||
enabled: false
|
||||
})
|
||||
await context.database.close()
|
||||
})
|
||||
})
|
||||
|
||||
class TransactionCountingDatabase implements RelayDatabase {
|
||||
@@ -2066,3 +2137,18 @@ class CellInventoryLockProbe {
|
||||
return decorate(database)
|
||||
}
|
||||
}
|
||||
|
||||
async function workerState(
|
||||
context: Context
|
||||
): Promise<{ consecutiveFailures: number; pausedUntil: number }> {
|
||||
const row = (
|
||||
await context.database.query(
|
||||
`SELECT consecutive_failures, paused_until
|
||||
FROM relay_region_rehome_worker_state WHERE worker_id = 'global'`
|
||||
)
|
||||
)[0]!
|
||||
return {
|
||||
consecutiveFailures: Number(row.consecutive_failures),
|
||||
pausedUntil: Number(row.paused_until)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -30,11 +30,9 @@ describe('regional rehome worker', () => {
|
||||
}
|
||||
const claimRegionalRehome = vi.fn().mockResolvedValueOnce(null).mockResolvedValue(attempt)
|
||||
const recordRegionalRehomeDrainReceipt = vi.fn().mockResolvedValue(true)
|
||||
const recordRegionalRehomeWorkerFailure = vi.fn().mockResolvedValue(undefined)
|
||||
const assignments = {
|
||||
claimRegionalRehome,
|
||||
recordRegionalRehomeDrainReceipt,
|
||||
recordRegionalRehomeWorkerFailure
|
||||
recordRegionalRehomeDrainReceipt
|
||||
} as unknown as RelayAssignmentStore
|
||||
const requests: Array<{ url: string; init?: RequestInit }> = []
|
||||
const warn = vi.spyOn(console, 'warn').mockImplementation(() => {})
|
||||
@@ -99,8 +97,7 @@ describe('regional rehome worker', () => {
|
||||
const claimRegionalRehome = vi.fn().mockResolvedValueOnce(null).mockResolvedValue(attempt)
|
||||
const assignments = {
|
||||
claimRegionalRehome,
|
||||
recordRegionalRehomeDispatchFailure: vi.fn().mockResolvedValue(undefined),
|
||||
recordRegionalRehomeWorkerFailure: vi.fn().mockResolvedValue(undefined)
|
||||
recordRegionalRehomeDispatchFailure: vi.fn().mockResolvedValue(undefined)
|
||||
} as unknown as RelayAssignmentStore
|
||||
const worker = startRegionalRehomeWorker(config(), assignments, {
|
||||
now: () => now,
|
||||
@@ -118,7 +115,36 @@ describe('regional rehome worker', () => {
|
||||
expect(assignments.recordRegionalRehomeDispatchFailure).toHaveBeenCalledWith(
|
||||
'11111111-1111-4111-8111-111111111111'
|
||||
)
|
||||
expect(assignments.recordRegionalRehomeWorkerFailure).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('keeps a failed poll out of the durable dispatch-failure budget', async () => {
|
||||
let now = 0
|
||||
const claimRegionalRehome = vi
|
||||
.fn()
|
||||
.mockResolvedValueOnce(null)
|
||||
.mockRejectedValue(new Error('Connection terminated due to connection timeout'))
|
||||
const recordRegionalRehomeDispatchFailure = vi.fn().mockResolvedValue(undefined)
|
||||
const assignments = {
|
||||
claimRegionalRehome,
|
||||
recordRegionalRehomeDispatchFailure
|
||||
} as unknown as RelayAssignmentStore
|
||||
const warn = vi.spyOn(console, 'warn').mockImplementation(() => {})
|
||||
const worker = startRegionalRehomeWorker(config(), assignments, {
|
||||
now: () => now,
|
||||
safetySnapshot: () => safety(now),
|
||||
intervalMs: 60_000
|
||||
})!
|
||||
await settleWorker()
|
||||
now = 1_000
|
||||
await expect(worker.run()).resolves.toBeUndefined()
|
||||
worker.stop()
|
||||
|
||||
// The poll never claimed an attempt, so nothing was drained and nothing may
|
||||
// be charged to the budget that latches the durable control off.
|
||||
expect(recordRegionalRehomeDispatchFailure).not.toHaveBeenCalled()
|
||||
expect(warn.mock.calls.map((call) => JSON.parse(String(call[0])).event)).toEqual([
|
||||
'orca_relay_regional_rehome_poll_failed'
|
||||
])
|
||||
})
|
||||
|
||||
it('passes unsafe process telemetry to the durable claim gate', async () => {
|
||||
|
||||
@@ -97,13 +97,19 @@ export function startRegionalRehomeWorker(
|
||||
})
|
||||
)
|
||||
} catch (error) {
|
||||
await (attemptId
|
||||
? assignments.recordRegionalRehomeDispatchFailure(attemptId)
|
||||
: assignments.recordRegionalRehomeWorkerFailure()
|
||||
).catch(() => undefined)
|
||||
// Only a claimed attempt was drained. A poll that failed before the claim
|
||||
// - a pool timeout on the once-a-second control read - dispatched nothing,
|
||||
// so it must not spend the budget that latches the durable control off.
|
||||
if (attemptId) {
|
||||
await assignments
|
||||
.recordRegionalRehomeDispatchFailure(attemptId)
|
||||
.catch(() => undefined)
|
||||
}
|
||||
console.warn(
|
||||
JSON.stringify({
|
||||
event: 'orca_relay_regional_rehome_dispatch_failed',
|
||||
event: attemptId
|
||||
? 'orca_relay_regional_rehome_dispatch_failed'
|
||||
: 'orca_relay_regional_rehome_poll_failed',
|
||||
reason: error instanceof Error ? error.message : 'unknown'
|
||||
})
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user