feat(server/openai/v1): structured request/delivered logging with requestId
The `/v1/openai/{chat,audio}` handlers used to be silent past
`hono/logger`'s `<-- POST` / `--> 502` lines — no userId, no model,
no token counts, no flux billed. Operators looking at a real gateway
incident had to cross-reference traces, request-log rows, and billing
ledger entries by timestamp alone. For a gateway whose value is
auditable per-request metering, that's not enough.
Hoists `requestId = nanoid()` to handler entry so the same
correlation id flows through:
- the inbound log line (model / stream / messageCount or inputChars
for TTS)
- the per-stream / non-streaming delivered log (status / durationMs /
promptTokens / completionTokens / fluxConsumed)
- the upstream-error degraded path (warn level)
- the partial-debit and debit-failure paths (already used requestId)
- `billingService.consumeFluxForLLM` / `ttsMeter.accumulate` for
DB-level idempotency (replaces the previous local nanoid() calls)
handleListVoices gets a `debug`-level line — the route is high-
frequency from UI voice pickers and we don't want it in the regular
audit feed, but it's useful when chasing voice-picker drift bugs.
No new schemas, no metric emission changes; this is purely logger
output. Pairs with the cause-propagation change so errors carry the
upstream snippet AND the request can be traced end-to-end by id.
This commit is contained in:
@@ -119,6 +119,12 @@ export function createV1CompletionsRoutes(
|
|||||||
// the gate, concurrent requests are rejected before the upstream call.
|
// the gate, concurrent requests are rejected before the upstream call.
|
||||||
async function handleCompletion(c: Context<HonoEnv>) {
|
async function handleCompletion(c: Context<HonoEnv>) {
|
||||||
const user = c.get('user')!
|
const user = c.get('user')!
|
||||||
|
// Generated up-front so incoming, completion, partial-debit, debit-failure,
|
||||||
|
// and request-log entries all carry the same correlation id. Re-used as
|
||||||
|
// the billing requestId (both streaming and non-streaming branches) for
|
||||||
|
// DB-level idempotency.
|
||||||
|
const requestId = nanoid()
|
||||||
|
|
||||||
// Read billing rates before pre-flight so the gate can compare against
|
// Read billing rates before pre-flight so the gate can compare against
|
||||||
// the realistic per-request cost (fallback rate), not just `> 0`.
|
// the realistic per-request cost (fallback rate), not just `> 0`.
|
||||||
const fallbackRate = await configKV.getOrThrow('FLUX_PER_REQUEST')
|
const fallbackRate = await configKV.getOrThrow('FLUX_PER_REQUEST')
|
||||||
@@ -136,6 +142,15 @@ export function createV1CompletionsRoutes(
|
|||||||
requestModel = await configKV.getOrThrow('DEFAULT_CHAT_MODEL')
|
requestModel = await configKV.getOrThrow('DEFAULT_CHAT_MODEL')
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const stream = !!body.stream
|
||||||
|
logger.withFields({
|
||||||
|
requestId,
|
||||||
|
userId: user.id,
|
||||||
|
model: requestModel,
|
||||||
|
stream,
|
||||||
|
messageCount: Array.isArray(body.messages) ? body.messages.length : undefined,
|
||||||
|
}).log('chat completion request')
|
||||||
|
|
||||||
// Server-connection attrs come from the router (which knows the actual
|
// Server-connection attrs come from the router (which knows the actual
|
||||||
// upstream baseURL it dispatched to) — it enriches the active span with
|
// upstream baseURL it dispatched to) — it enriches the active span with
|
||||||
// its own `airi.gen_ai.gateway.*` attrs on success.
|
// its own `airi.gen_ai.gateway.*` attrs on success.
|
||||||
@@ -143,7 +158,7 @@ export function createV1CompletionsRoutes(
|
|||||||
attributes: {
|
attributes: {
|
||||||
[GEN_AI_ATTR_OPERATION_NAME]: 'chat',
|
[GEN_AI_ATTR_OPERATION_NAME]: 'chat',
|
||||||
[GEN_AI_ATTR_REQUEST_MODEL]: requestModel,
|
[GEN_AI_ATTR_REQUEST_MODEL]: requestModel,
|
||||||
[AIRI_ATTR_GEN_AI_STREAM]: !!body.stream,
|
[AIRI_ATTR_GEN_AI_STREAM]: stream,
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
|
|
||||||
@@ -190,6 +205,10 @@ export function createV1CompletionsRoutes(
|
|||||||
stream: !!body.stream,
|
stream: !!body.stream,
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
|
|
||||||
|
logger.withFields({ requestId, userId: user.id, model: requestModel, status: response.status, durationMs })
|
||||||
|
.warn('chat completion delivered with upstream error status')
|
||||||
|
|
||||||
return new Response(response.body, {
|
return new Response(response.body, {
|
||||||
status: response.status,
|
status: response.status,
|
||||||
headers: buildSafeResponseHeaders(response),
|
headers: buildSafeResponseHeaders(response),
|
||||||
@@ -200,7 +219,7 @@ export function createV1CompletionsRoutes(
|
|||||||
// `fallbackRate` / `fluxPer1kTokens` were hoisted to the top of this
|
// `fallbackRate` / `fluxPer1kTokens` were hoisted to the top of this
|
||||||
// function so the pre-flight gate can use them too.
|
// function so the pre-flight gate can use them too.
|
||||||
|
|
||||||
if (body.stream) {
|
if (stream) {
|
||||||
// Streaming: return response immediately, bill after stream ends
|
// Streaming: return response immediately, bill after stream ends
|
||||||
const { readable, writable } = new TransformStream()
|
const { readable, writable } = new TransformStream()
|
||||||
const reader = response.body!.getReader()
|
const reader = response.body!.getReader()
|
||||||
@@ -302,7 +321,6 @@ export function createV1CompletionsRoutes(
|
|||||||
// race) or real DB errors. Partial debits are signalled via the
|
// race) or real DB errors. Partial debits are signalled via the
|
||||||
// returned `charged < requested` and accounted to the same
|
// returned `charged < requested` and accounted to the same
|
||||||
// `fluxUnbilled` counter (different `reason` label).
|
// `fluxUnbilled` counter (different `reason` label).
|
||||||
const requestId = nanoid()
|
|
||||||
let actualCharged = 0
|
let actualCharged = 0
|
||||||
try {
|
try {
|
||||||
const result = await billingService.consumeFluxForLLM({
|
const result = await billingService.consumeFluxForLLM({
|
||||||
@@ -367,6 +385,18 @@ export function createV1CompletionsRoutes(
|
|||||||
stream_interrupted: streamInterrupted,
|
stream_interrupted: streamInterrupted,
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
|
|
||||||
|
logger.withFields({
|
||||||
|
requestId,
|
||||||
|
userId: user.id,
|
||||||
|
model: requestModel,
|
||||||
|
status: response.status,
|
||||||
|
durationMs,
|
||||||
|
promptTokens: usage.promptTokens,
|
||||||
|
completionTokens: usage.completionTokens,
|
||||||
|
fluxConsumed: actualCharged,
|
||||||
|
stream: true,
|
||||||
|
}).log('chat completion delivered')
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
})()
|
})()
|
||||||
@@ -394,7 +424,6 @@ export function createV1CompletionsRoutes(
|
|||||||
// The upstream call has already happened (cost incurred), so partial
|
// The upstream call has already happened (cost incurred), so partial
|
||||||
// debit + `fluxUnbilled` is the only sane recovery — same shape as the
|
// debit + `fluxUnbilled` is the only sane recovery — same shape as the
|
||||||
// streaming path. `balance <= 0` still throws and bubbles up as 402.
|
// streaming path. `balance <= 0` still throws and bubbles up as 402.
|
||||||
const requestId = nanoid()
|
|
||||||
const result = await billingService.consumeFluxForLLM({
|
const result = await billingService.consumeFluxForLLM({
|
||||||
userId: user.id,
|
userId: user.id,
|
||||||
amount: fluxConsumed,
|
amount: fluxConsumed,
|
||||||
@@ -443,11 +472,24 @@ export function createV1CompletionsRoutes(
|
|||||||
},
|
},
|
||||||
})
|
})
|
||||||
|
|
||||||
|
logger.withFields({
|
||||||
|
requestId,
|
||||||
|
userId: user.id,
|
||||||
|
model: requestModel,
|
||||||
|
status: response.status,
|
||||||
|
durationMs,
|
||||||
|
promptTokens: usage.promptTokens,
|
||||||
|
completionTokens: usage.completionTokens,
|
||||||
|
fluxConsumed: result.charged,
|
||||||
|
stream: false,
|
||||||
|
}).log('chat completion delivered')
|
||||||
|
|
||||||
return c.json(responseBody)
|
return c.json(responseBody)
|
||||||
}
|
}
|
||||||
|
|
||||||
async function handleTTS(c: Context<HonoEnv>) {
|
async function handleTTS(c: Context<HonoEnv>) {
|
||||||
const user = c.get('user')!
|
const user = c.get('user')!
|
||||||
|
const requestId = nanoid()
|
||||||
const flux = await fluxService.getFlux(user.id)
|
const flux = await fluxService.getFlux(user.id)
|
||||||
if (flux.flux <= 0) {
|
if (flux.flux <= 0) {
|
||||||
throw createPaymentRequiredError('Insufficient flux')
|
throw createPaymentRequiredError('Insufficient flux')
|
||||||
@@ -463,6 +505,14 @@ export function createV1CompletionsRoutes(
|
|||||||
requestModel = await configKV.getOrThrow('DEFAULT_TTS_MODEL')
|
requestModel = await configKV.getOrThrow('DEFAULT_TTS_MODEL')
|
||||||
}
|
}
|
||||||
|
|
||||||
|
logger.withFields({
|
||||||
|
requestId,
|
||||||
|
userId: user.id,
|
||||||
|
model: requestModel,
|
||||||
|
inputChars: inputText.length,
|
||||||
|
voice: typeof body.voice === 'string' ? body.voice : undefined,
|
||||||
|
}).log('tts speech request')
|
||||||
|
|
||||||
// Pre-flight: refuse before hitting upstream if this segment would push the
|
// Pre-flight: refuse before hitting upstream if this segment would push the
|
||||||
// user past their balance. Cheap-path requests below the Flux threshold
|
// user past their balance. Cheap-path requests below the Flux threshold
|
||||||
// still pass when the user has at least 1 Flux.
|
// still pass when the user has at least 1 Flux.
|
||||||
@@ -506,6 +556,8 @@ export function createV1CompletionsRoutes(
|
|||||||
span.setStatus({ code: SpanStatusCode.ERROR, message: `Gateway ${response.status}` })
|
span.setStatus({ code: SpanStatusCode.ERROR, message: `Gateway ${response.status}` })
|
||||||
span.end()
|
span.end()
|
||||||
recordMetrics({ model: requestModel, status: response.status, type: 'tts', durationMs, fluxConsumed: 0 })
|
recordMetrics({ model: requestModel, status: response.status, type: 'tts', durationMs, fluxConsumed: 0 })
|
||||||
|
logger.withFields({ requestId, userId: user.id, model: requestModel, status: response.status, durationMs })
|
||||||
|
.warn('tts speech delivered with upstream error status')
|
||||||
return new Response(response.body, {
|
return new Response(response.body, {
|
||||||
status: response.status,
|
status: response.status,
|
||||||
headers: buildSafeResponseHeaders(response),
|
headers: buildSafeResponseHeaders(response),
|
||||||
@@ -527,7 +579,7 @@ export function createV1CompletionsRoutes(
|
|||||||
userId: user.id,
|
userId: user.id,
|
||||||
units: inputText.length,
|
units: inputText.length,
|
||||||
currentBalance: flux.flux,
|
currentBalance: flux.flux,
|
||||||
requestId: nanoid(),
|
requestId,
|
||||||
metadata: { model: requestModel },
|
metadata: { model: requestModel },
|
||||||
})
|
})
|
||||||
fluxConsumed = result.fluxDebited
|
fluxConsumed = result.fluxDebited
|
||||||
@@ -546,6 +598,16 @@ export function createV1CompletionsRoutes(
|
|||||||
fluxConsumed,
|
fluxConsumed,
|
||||||
})
|
})
|
||||||
|
|
||||||
|
logger.withFields({
|
||||||
|
requestId,
|
||||||
|
userId: user.id,
|
||||||
|
model: requestModel,
|
||||||
|
status: response.status,
|
||||||
|
durationMs,
|
||||||
|
inputChars: inputText.length,
|
||||||
|
fluxConsumed,
|
||||||
|
}).log('tts speech delivered')
|
||||||
|
|
||||||
return new Response(response.body, {
|
return new Response(response.body, {
|
||||||
status: response.status,
|
status: response.status,
|
||||||
headers: buildSafeResponseHeaders(response),
|
headers: buildSafeResponseHeaders(response),
|
||||||
@@ -566,6 +628,10 @@ export function createV1CompletionsRoutes(
|
|||||||
|
|
||||||
const voices = await llmRouter.listTtsVoices(model)
|
const voices = await llmRouter.listTtsVoices(model)
|
||||||
const recommended = (await configKV.getOptional('DEFAULT_TTS_VOICES')) ?? {}
|
const recommended = (await configKV.getOptional('DEFAULT_TTS_VOICES')) ?? {}
|
||||||
|
// Debug level: high-frequency catalog poll from UI selectors, no
|
||||||
|
// billing / user-facing side effect — useful only when debugging
|
||||||
|
// voice-picker drift, never as a permanent audit trail line.
|
||||||
|
logger.withFields({ model, voiceCount: voices.length }).debug('list tts voices')
|
||||||
return Response.json({ voices, recommended })
|
return Response.json({ voices, recommended })
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user