Skip to content
Merged
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
57 changes: 57 additions & 0 deletions packages/cloud/api/__tests__/apps-chat-stream-refund.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@
*/
import { describe, expect, mock, test } from "bun:test";
import {
reconcileNonStreamingSettleError,
reconcileStreamProcessingError,
type StreamRefundCredits,
} from "../v1/apps/[id]/chat/stream-refund";
Expand Down Expand Up @@ -106,3 +107,59 @@ describe("reconcileStreamProcessingError (#10837)", () => {
expect(credits.calls.every((c) => c.actualBaseCost === 0)).toBe(true);
});
});

const nonStreamBase = {
appId: "app-1",
userId: "user-1",
reservedBaseCost: 0.05,
model: "openai/gpt-oss-120b",
provider: "openai",
billingSource: "openai",
errorMessage: "provider body was not valid JSON",
};

describe("reconcileNonStreamingSettleError (#11169 part 1)", () => {
test("throw BEFORE the settle reconcile was invoked → full refund (actualBaseCost 0)", async () => {
const credits = makeCredits();
const result = await reconcileNonStreamingSettleError(
{ ...nonStreamBase, settleStarted: false },
credits,
);
expect(result.refunded).toBe(true);
expect(credits.reconcileCredits).toHaveBeenCalledTimes(1);
expect(credits.calls[0]).toEqual({
estimatedBaseCost: 0.05,
actualBaseCost: 0,
});
});

test("throw at/after the settle reconcile (incl. from INSIDE it — movement may have committed) → NO refund (no double-credit)", async () => {
const credits = makeCredits();
const result = await reconcileNonStreamingSettleError(
{ ...nonStreamBase, settleStarted: true },
credits,
);
expect(result.refunded).toBe(false);
expect(credits.reconcileCredits).not.toHaveBeenCalled();
});

test("the refund is tagged non-streaming (streaming:false) so it's distinguishable in the ledger", async () => {
const metaCalls: Array<Record<string, unknown> | undefined> = [];
const credits = {
reconcileCredits: mock(
async (args: { metadata?: Record<string, unknown> }) => {
metaCalls.push(args.metadata);
return null;
},
),
} as unknown as StreamRefundCredits;
await reconcileNonStreamingSettleError(
{ ...nonStreamBase, settleStarted: false },
credits,
);
expect(metaCalls[0]).toMatchObject({
streaming: false,
refundReason: "non_streaming_settle_error",
});
});
});
173 changes: 113 additions & 60 deletions packages/cloud/api/v1/apps/[id]/chat/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,10 @@ import { logger } from "@/lib/utils/logger";
import { getRouteTimeoutMs } from "@/lib/utils/request-timeout";
import type { AppEnv } from "@/types/cloud-worker-env";
import { reservationOutputTokens } from "./chat-reservation";
import { reconcileStreamProcessingError } from "./stream-refund";
import {
reconcileNonStreamingSettleError,
reconcileStreamProcessingError,
} from "./stream-refund";

const ROUTE_MAX_DURATION = 800;

Expand Down Expand Up @@ -659,73 +662,123 @@ async function handlePOST(
);
}

// Non-streaming response
const responseData = (await providerResponse.json()) as {
usage?: { prompt_tokens?: number; completion_tokens?: number };
choices?: Array<{ message?: { content?: string } }>;
};

// Calculate actual cost - use fallback estimation if provider doesn't return usage
let actualInputTokens = responseData.usage?.prompt_tokens || 0;
let actualOutputTokens = responseData.usage?.completion_tokens || 0;

// Fallback: estimate tokens if usage not provided (matching streaming behavior)
if (actualInputTokens === 0 && actualOutputTokens === 0) {
const outputContent = responseData.choices?.[0]?.message?.content || "";
actualInputTokens = estimatedInputTokens; // Use pre-calculated estimate
actualOutputTokens = estimateTokens(outputContent);

logger.warn("[App Chat] No usage data in response, using estimates", {
appId,
// Non-streaming response. The body-read + cost-calc + reconcile below run
// AFTER the upfront debit, and the outer catch returns 500 WITHOUT refunding
// — so a malformed body / calculateCost throw here stranded the reserved
// hold (#11169 part 1; the streaming branch was already guarded).
// Guard the settle path: any throw BEFORE the reconcile is invoked refunds
// the hold, then rethrows to the outer error handler. Once the reconcile
// has been invoked we never refund — reconcileCredits is not transactional,
// so a mid-flight throw may have already committed the org-balance movement
// and a blind refund would double-credit (same reason the streaming branch
// flips streamCompleted before ITS settle). A hold stranded by that rare
// window is recovered by the stranded-reservation sweep (#11169 part 3).
let nonStreamingSettleStarted = false;
try {
const responseData = (await providerResponse.json()) as {
usage?: { prompt_tokens?: number; completion_tokens?: number };
choices?: Array<{ message?: { content?: string } }>;
};

// Calculate actual cost - use fallback estimation if provider doesn't return usage
let actualInputTokens = responseData.usage?.prompt_tokens || 0;
let actualOutputTokens = responseData.usage?.completion_tokens || 0;

// Fallback: estimate tokens if usage not provided (matching streaming behavior)
if (actualInputTokens === 0 && actualOutputTokens === 0) {
const outputContent = responseData.choices?.[0]?.message?.content || "";
actualInputTokens = estimatedInputTokens; // Use pre-calculated estimate
actualOutputTokens = estimateTokens(outputContent);

logger.warn("[App Chat] No usage data in response, using estimates", {
appId,
actualInputTokens,
actualOutputTokens,
});
}

const { totalCost: actualBaseCost } = await calculateCost(
normalizedModel,
provider,
actualInputTokens,
actualOutputTokens,
});
}
billingSource,
);

const { totalCost: actualBaseCost } = await calculateCost(
normalizedModel,
provider,
actualInputTokens,
actualOutputTokens,
billingSource,
);
// Reconcile the difference between reserved and actual costs
// Pass app to avoid N+1 query (app already fetched above)
// Flag BEFORE the call: reconcileCredits commits its org-balance movement
// before its (non-co-transactional) earnings/counter writes, so a throw
// from inside it must NOT trigger the refund below (double-credit).
nonStreamingSettleStarted = true;
const reconciliation = await appCreditsService.reconcileCredits({
appId,
userId: user.id,
estimatedBaseCost: reservedBaseCost,
actualBaseCost,
description: `Chat reconciliation: ${model}`,
metadata: {
model,
provider,
billingSource,
inputTokens: actualInputTokens,
outputTokens: actualOutputTokens,
streaming: false,
},
app,
});

// Reconcile the difference between reserved and actual costs
// Pass app to avoid N+1 query (app already fetched above)
const reconciliation = await appCreditsService.reconcileCredits({
appId,
userId: user.id,
estimatedBaseCost: reservedBaseCost,
actualBaseCost,
description: `Chat reconciliation: ${model}`,
metadata: {
const duration = Date.now() - startTime;
logger.info("[App Chat] Request completed", {
appId,
userId: user.id,
model,
provider,
billingSource,
duration,
inputTokens: actualInputTokens,
outputTokens: actualOutputTokens,
streaming: false,
},
app,
});

const duration = Date.now() - startTime;
logger.info("[App Chat] Request completed", {
appId,
userId: user.id,
model,
duration,
inputTokens: actualInputTokens,
outputTokens: actualOutputTokens,
reservedBaseCost,
actualBaseCost,
reconciliation: {
action: reconciliation.action,
amount: reconciliation.adjustedAmount,
},
});
reservedBaseCost,
actualBaseCost,
reconciliation: {
action: reconciliation.action,
amount: reconciliation.adjustedAmount,
},
});

return withCors(Response.json(responseData));
return withCors(Response.json(responseData));
} catch (nonStreamingError) {
// Refund the upfront hold if the settle reconcile was never invoked
// (#11169 part 1). Best-effort: a refund failure is logged but never
// masks the original error surfaced to the client.
await reconcileNonStreamingSettleError(
{
settleStarted: nonStreamingSettleStarted,
appId,
userId: user.id,
reservedBaseCost,
model,
provider,
billingSource,
errorMessage:
nonStreamingError instanceof Error
? nonStreamingError.message
: String(nonStreamingError),
},
appCreditsService,
).catch((refundError) => {
logger.error(
"[App Chat] refund after non-streaming settle failure ALSO failed — hold stranded",
{
appId,
userId: user.id,
error:
refundError instanceof Error
? refundError.message
: String(refundError),
},
);
});
throw nonStreamingError;
}
} catch (error) {
logger.error("[App Chat] Error:", error);

Expand Down
73 changes: 73 additions & 0 deletions packages/cloud/api/v1/apps/[id]/chat/stream-refund.ts
Original file line number Diff line number Diff line change
Expand Up @@ -68,3 +68,76 @@ export async function reconcileStreamProcessingError(
});
return { refunded: true };
}

/**
* Money-critical (#11169 part 1): the NON-streaming app-chat path debits the
* upfront hold, then reads the provider body + runs `calculateCost` +
* `reconcileCredits`. If the body-read or cost-calc throws AFTER the debit, the
* route's outer catch returns 500 WITHOUT refunding — stranding the reserved
* hold. Refund it: the caller received no billable answer and no settle was
* ever attempted.
*
* `settleStarted` is true once `reconcileCredits` has been INVOKED — not once
* it returned. `reconcileCredits` is not transactional: it commits its
* org-balance movement (refund or extra charge) before its earnings/counter
* writes, and that movement carries no idempotency key. A throw from inside it
* may therefore have already moved money, so refunding blindly would
* double-credit the org (mint credits during a DB blip, systemically across
* concurrent requests). Mirror of the streaming branch, which flips
* `streamCompleted` before ITS settle for the same reason. A hold stranded by
* that rare window is recovered by the stranded-reservation sweep
* (#11169 part 3).
*/
export async function reconcileNonStreamingSettleError(
params: {
settleStarted: boolean;
appId: string;
userId: string;
reservedBaseCost: number;
model: string;
provider: string;
billingSource: string;
errorMessage: string;
},
credits: StreamRefundCredits,
): Promise<{ refunded: boolean }> {
const {
settleStarted,
appId,
userId,
reservedBaseCost,
model,
provider,
billingSource,
errorMessage,
} = params;

if (settleStarted) {
logger.error(
"[App Chat] Non-streaming throw at/after the settle reconcile; NOT refunding (movement may have committed — sweep recovers a stranded hold)",
{ appId, userId, reservedBaseCost, error: errorMessage },
);
return { refunded: false };
}

logger.error(
"[App Chat] Non-streaming settle never started after debit; refunding reserved hold (#11169)",
{ appId, userId, reservedBaseCost, error: errorMessage },
);
await credits.reconcileCredits({
appId,
userId,
estimatedBaseCost: reservedBaseCost,
actualBaseCost: 0, // Full refund — nothing was billed.
description: `Chat refund (non-streaming settle failed): ${model}`,
metadata: {
error: true,
streaming: false,
model,
provider,
billingSource,
refundReason: "non_streaming_settle_error",
},
});
return { refunded: true };
}
Loading