fix(events): retry transient event_log persistence failures and stop racing function suspension (#966)
Production logs (2 months) show ~28 event_log inserts dying with TypeError: fetch failed. Root cause: the MCP server emits telemetry fire-and-forget and returns the JSON-RPC response immediately, so the insert races Vercel function suspension; supabase-js surfaces the dead fetch as a network error which was logged at error level. Two-part fix: - persistEvent (event-log-handler.ts) retries the insert once after 250ms when the error message contains "fetch failed" (network class only; constraint violations and other Postgres errors are never retried). On final failure, telemetry event types (mcp.*, agent.*) log at warn; business events (journal_entry.*, invoice.*, etc., which feed webhook delivery) stay at error. - The mcp-server telemetry emit sites (tool_called, tools_list_called, resource_read, next_hint_followed, skill_loaded, workflow_started, agent.feedback) now schedule the emit via after() from next/server, which keeps the function alive past the response until the emit settles. Falls back to plain fire-and-forget when no request scope exists (direct handler invocation in tests). Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -1,4 +1,4 @@
|
||||
import { NextResponse } from 'next/server'
|
||||
import { NextResponse, after } from 'next/server'
|
||||
import {
|
||||
extractBearerToken,
|
||||
validateApiKey,
|
||||
@@ -2155,7 +2155,7 @@ export const tools: McpTool[] = [
|
||||
}
|
||||
feedbackRateLimit.set(rateKey, now)
|
||||
|
||||
void eventBus
|
||||
emitAfterResponse(() => eventBus
|
||||
.emit({
|
||||
type: 'agent.feedback',
|
||||
payload: {
|
||||
@@ -2172,7 +2172,7 @@ export const tools: McpTool[] = [
|
||||
companyId,
|
||||
},
|
||||
})
|
||||
.catch((err) => console.error('[mcp] agent.feedback emit failed:', err))
|
||||
.catch((err) => console.error('[mcp] agent.feedback emit failed:', err)))
|
||||
|
||||
return {
|
||||
recorded: true,
|
||||
@@ -11673,6 +11673,21 @@ function jsonRpcError(
|
||||
return { jsonrpc: '2.0', id, error: { code, message, data } }
|
||||
}
|
||||
|
||||
/**
|
||||
* Schedule a fire-and-forget telemetry emit so it cannot race Vercel function
|
||||
* suspension: `after()` keeps the function alive past the JSON-RPC response
|
||||
* until the emit settles, which is why event_log inserts used to die with
|
||||
* "TypeError: fetch failed". Falls back to a plain fire-and-forget emit when
|
||||
* no Next request scope exists (direct handler invocation in tests).
|
||||
*/
|
||||
function emitAfterResponse(emit: () => Promise<void>): void {
|
||||
try {
|
||||
after(emit)
|
||||
} catch {
|
||||
void emit()
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Emit `mcp.tool_called` telemetry to the event bus. Fire-and-forget: the
|
||||
* dispatcher must never block the JSON-RPC response on telemetry, and a failing
|
||||
@@ -11693,7 +11708,7 @@ function emitToolCallTelemetry(payload: {
|
||||
userId: string
|
||||
companyId: string
|
||||
}): void {
|
||||
void eventBus
|
||||
emitAfterResponse(() => eventBus
|
||||
.emit({
|
||||
type: 'mcp.tool_called',
|
||||
payload: {
|
||||
@@ -11722,7 +11737,7 @@ function emitToolCallTelemetry(payload: {
|
||||
// Last-resort guard. EventBus.emit already swallows handler failures,
|
||||
// but if the bus itself is in a bad state we still don't want to break tools.
|
||||
console.error('[mcp] tool_called telemetry emit failed:', err)
|
||||
})
|
||||
}))
|
||||
}
|
||||
|
||||
/** Fire-and-forget telemetry for a tools/list call. */
|
||||
@@ -11734,7 +11749,7 @@ function emitToolsListTelemetry(payload: {
|
||||
userId: string
|
||||
companyId: string
|
||||
}): void {
|
||||
void eventBus
|
||||
emitAfterResponse(() => eventBus
|
||||
.emit({
|
||||
type: 'mcp.tools_list_called',
|
||||
payload: {
|
||||
@@ -11752,7 +11767,7 @@ function emitToolsListTelemetry(payload: {
|
||||
})
|
||||
.catch((err) => {
|
||||
console.error('[mcp] tools_list_called telemetry emit failed:', err)
|
||||
})
|
||||
}))
|
||||
}
|
||||
|
||||
/** Fire-and-forget telemetry for a resources/read call. */
|
||||
@@ -11767,7 +11782,7 @@ function emitResourceReadTelemetry(payload: {
|
||||
userId: string
|
||||
companyId: string
|
||||
}): void {
|
||||
void eventBus
|
||||
emitAfterResponse(() => eventBus
|
||||
.emit({
|
||||
type: 'mcp.resource_read',
|
||||
payload: {
|
||||
@@ -11788,7 +11803,7 @@ function emitResourceReadTelemetry(payload: {
|
||||
})
|
||||
.catch((err) => {
|
||||
console.error('[mcp] resource_read telemetry emit failed:', err)
|
||||
})
|
||||
}))
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -11836,7 +11851,7 @@ function checkAndEmitNextHintFollowed(
|
||||
// Consume the hint so we don't double-count if the agent calls the same
|
||||
// tool twice in a row (idempotent retries shouldn't inflate the metric).
|
||||
lastResponseHintBySession.delete(sessionId)
|
||||
void eventBus
|
||||
emitAfterResponse(() => eventBus
|
||||
.emit({
|
||||
type: 'mcp.next_hint_followed',
|
||||
payload: {
|
||||
@@ -11850,7 +11865,7 @@ function checkAndEmitNextHintFollowed(
|
||||
companyId,
|
||||
},
|
||||
})
|
||||
.catch((err) => console.error('[mcp] next_hint_followed emit failed:', err))
|
||||
.catch((err) => console.error('[mcp] next_hint_followed emit failed:', err)))
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -11866,7 +11881,7 @@ function emitSkillLoaded(payload: {
|
||||
userId: string
|
||||
companyId: string
|
||||
}): void {
|
||||
void eventBus
|
||||
emitAfterResponse(() => eventBus
|
||||
.emit({
|
||||
type: 'mcp.skill_loaded',
|
||||
payload: {
|
||||
@@ -11880,7 +11895,7 @@ function emitSkillLoaded(payload: {
|
||||
companyId: payload.companyId,
|
||||
},
|
||||
})
|
||||
.catch((err) => console.error('[mcp] skill_loaded emit failed:', err))
|
||||
.catch((err) => console.error('[mcp] skill_loaded emit failed:', err)))
|
||||
}
|
||||
|
||||
/** Fire-and-forget telemetry for workflow lifecycle. */
|
||||
@@ -11890,7 +11905,7 @@ function emitWorkflowStarted(payload: {
|
||||
userId: string
|
||||
companyId: string
|
||||
}): void {
|
||||
void eventBus
|
||||
emitAfterResponse(() => eventBus
|
||||
.emit({
|
||||
type: 'mcp.workflow_started',
|
||||
payload: {
|
||||
@@ -11903,7 +11918,7 @@ function emitWorkflowStarted(payload: {
|
||||
companyId: payload.companyId,
|
||||
},
|
||||
})
|
||||
.catch((err) => console.error('[mcp] workflow_started emit failed:', err))
|
||||
.catch((err) => console.error('[mcp] workflow_started emit failed:', err)))
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -12,6 +12,22 @@ vi.mock('@/lib/auth/api-keys', () => ({
|
||||
}),
|
||||
}))
|
||||
|
||||
// Mock the logger so warn-vs-error level can be asserted (the real logger
|
||||
// suppresses warn in the test environment). Hoisted: bus.ts calls
|
||||
// createLogger at import time, before top-level consts initialize.
|
||||
const { logWarn, logError } = vi.hoisted(() => ({
|
||||
logWarn: vi.fn(),
|
||||
logError: vi.fn(),
|
||||
}))
|
||||
vi.mock('@/lib/logger', () => ({
|
||||
createLogger: () => ({
|
||||
info: vi.fn(),
|
||||
warn: logWarn,
|
||||
error: logError,
|
||||
child: vi.fn(),
|
||||
}),
|
||||
}))
|
||||
|
||||
// Import after mocks
|
||||
import { registerEventLogHandler } from '../event-log-handler'
|
||||
|
||||
@@ -138,6 +154,86 @@ describe('event-log-handler', () => {
|
||||
expect(mockInsert).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('retries once when the insert fails with a network-class "fetch failed" error, then succeeds', async () => {
|
||||
vi.useFakeTimers()
|
||||
try {
|
||||
mockInsert
|
||||
.mockResolvedValueOnce({ error: { message: 'TypeError: fetch failed' } })
|
||||
.mockResolvedValueOnce({ error: null })
|
||||
|
||||
const emitPromise = eventBus.emit({
|
||||
type: 'customer.created',
|
||||
payload: { customer: makeCustomer({ id: 'cust-retry' }), userId: 'user-1', companyId: 'company-1' },
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(250)
|
||||
await emitPromise
|
||||
|
||||
expect(mockInsert).toHaveBeenCalledTimes(2)
|
||||
expect(mockInsert.mock.calls[1][0]).toMatchObject({
|
||||
event_type: 'customer.created',
|
||||
entity_id: 'cust-retry',
|
||||
})
|
||||
expect(logWarn).not.toHaveBeenCalled()
|
||||
expect(logError).not.toHaveBeenCalled()
|
||||
} finally {
|
||||
vi.useRealTimers()
|
||||
}
|
||||
})
|
||||
|
||||
it('does not retry non-network insert errors', async () => {
|
||||
mockInsert.mockResolvedValue({
|
||||
error: { message: 'duplicate key value violates unique constraint "event_log_pkey"' },
|
||||
})
|
||||
|
||||
await eventBus.emit({
|
||||
type: 'customer.created',
|
||||
payload: { customer: makeCustomer(), userId: 'user-1', companyId: 'company-1' },
|
||||
})
|
||||
|
||||
expect(mockInsert).toHaveBeenCalledTimes(1)
|
||||
expect(logError).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('logs telemetry (mcp.*) persistence failure at warn level after the retry also fails', async () => {
|
||||
vi.useFakeTimers()
|
||||
try {
|
||||
mockInsert.mockResolvedValue({ error: { message: 'TypeError: fetch failed' } })
|
||||
|
||||
const emitPromise = eventBus.emit({
|
||||
type: 'mcp.tool_called',
|
||||
payload: { tool: 'gnubok_list_accounts', userId: 'user-1', companyId: 'company-1' } as never,
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(250)
|
||||
await emitPromise
|
||||
|
||||
expect(mockInsert).toHaveBeenCalledTimes(2)
|
||||
expect(logWarn).toHaveBeenCalledTimes(1)
|
||||
expect(logError).not.toHaveBeenCalled()
|
||||
} finally {
|
||||
vi.useRealTimers()
|
||||
}
|
||||
})
|
||||
|
||||
it('keeps business event persistence failure at error level after the retry also fails', async () => {
|
||||
vi.useFakeTimers()
|
||||
try {
|
||||
mockInsert.mockResolvedValue({ error: { message: 'TypeError: fetch failed' } })
|
||||
|
||||
const emitPromise = eventBus.emit({
|
||||
type: 'invoice.created',
|
||||
payload: { invoice: makeInvoice({ id: 'inv-err' }), userId: 'user-1', companyId: 'company-1' },
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(250)
|
||||
await emitPromise
|
||||
|
||||
expect(mockInsert).toHaveBeenCalledTimes(2)
|
||||
expect(logError).toHaveBeenCalledTimes(1)
|
||||
expect(logWarn).not.toHaveBeenCalled()
|
||||
} finally {
|
||||
vi.useRealTimers()
|
||||
}
|
||||
})
|
||||
|
||||
it('persists period.locked with period entity_id', async () => {
|
||||
const period = makeFiscalPeriod({ id: 'period-1' })
|
||||
|
||||
|
||||
@@ -129,8 +129,32 @@ function stripMetaFields(payload: Record<string, unknown>): Record<string, unkno
|
||||
return data
|
||||
}
|
||||
|
||||
/** Delay before the single retry of a network-class insert failure. */
|
||||
const RETRY_DELAY_MS = 250
|
||||
|
||||
/**
|
||||
* Network-class failure: supabase-js (undici) surfaces a dead connection as
|
||||
* "TypeError: fetch failed", typically when the insert races Vercel function
|
||||
* suspension. Only this class is retried; Postgres errors (constraint
|
||||
* violations, RLS, bad columns) would fail identically on a retry.
|
||||
*/
|
||||
function isTransientNetworkError(message: string): boolean {
|
||||
return message.includes('fetch failed')
|
||||
}
|
||||
|
||||
/**
|
||||
* Telemetry events (mcp.*, agent.*) are best-effort metrics: a lost row is
|
||||
* not actionable, so their final persistence failure logs at warn. Business
|
||||
* events (journal_entry.*, invoice.*, etc.) feed webhook delivery and stay
|
||||
* at error level.
|
||||
*/
|
||||
function isTelemetryEvent(eventType: string): boolean {
|
||||
return eventType.startsWith('mcp.') || eventType.startsWith('agent.')
|
||||
}
|
||||
|
||||
/**
|
||||
* Persist a single event to the event_log table.
|
||||
* Retries once on network-class failures ("fetch failed") after a short delay.
|
||||
*/
|
||||
async function persistEvent(
|
||||
eventType: string,
|
||||
@@ -140,19 +164,27 @@ async function persistEvent(
|
||||
data: Record<string, unknown>
|
||||
): Promise<void> {
|
||||
const supabase = createServiceClientNoCookies()
|
||||
const row = {
|
||||
user_id: userId,
|
||||
company_id: companyId,
|
||||
event_type: eventType,
|
||||
entity_id: entityId,
|
||||
data,
|
||||
}
|
||||
|
||||
const { error } = await supabase
|
||||
.from('event_log')
|
||||
.insert({
|
||||
user_id: userId,
|
||||
company_id: companyId,
|
||||
event_type: eventType,
|
||||
entity_id: entityId,
|
||||
data,
|
||||
})
|
||||
let { error } = await supabase.from('event_log').insert(row)
|
||||
|
||||
if (error && isTransientNetworkError(error.message)) {
|
||||
await new Promise((resolve) => setTimeout(resolve, RETRY_DELAY_MS))
|
||||
;({ error } = await supabase.from('event_log').insert(row))
|
||||
}
|
||||
|
||||
if (error) {
|
||||
log.error(`Failed to persist event ${eventType}:`, error.message)
|
||||
if (isTelemetryEvent(eventType)) {
|
||||
log.warn(`Failed to persist event ${eventType}:`, error.message)
|
||||
} else {
|
||||
log.error(`Failed to persist event ${eventType}:`, error.message)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user