feat(consult): add free model-classified smalltalk fast path
Keep full consultation tool contracts unchanged. Persist short replies and refund the original reservation atomically while recording actual model usage. Verify Linux frontend 3566/3566, database 40/40, Static home and gzip +0.0493%. Co-Authored-By: Claude Code <noreply@anthropic.com>
This commit is contained in:
@@ -1714,6 +1714,7 @@ export default function Home() {
|
||||
streamingThinking={activeStreamingThinking}
|
||||
streamingSections={activeStreamingSections}
|
||||
streamingTimeline={activeStreamingTimeline}
|
||||
streamingResponseKind={streamingReply?.sessionId === activeSession?.id ? streamingReply?.responseKind : undefined}
|
||||
sessionId={activeSession.id}
|
||||
sessionType={activeSession.sessionType}
|
||||
theme={activeSession.theme}
|
||||
|
||||
@@ -38,6 +38,8 @@ import { jsonForSupabaseSetupFailure } from "@/lib/api/service-unavailable";
|
||||
import { createAdminSupabaseClient } from "@/lib/supabase/admin";
|
||||
import { createServerSupabaseClient } from "@/lib/supabase/server";
|
||||
import { streamTextResponse } from "@/lib/stream-text-response";
|
||||
import { classifyConsultationTurn, type SmalltalkUsage } from "@/lib/consultation-smalltalk";
|
||||
import { streamSmalltalkResponse } from "@/lib/stream-smalltalk-response";
|
||||
import { consultationPublicActivityEvent, streamAgentResponse } from "@/lib/stream-agent-response";
|
||||
import type { AgentExecutionReceipt, WorkflowReceipt } from "@/lib/consultation-agent-events";
|
||||
import { consultationComposePrompt, consultationContinuePrompt, natalConsultationThinkingPlan, type PublicThinkingSection } from "@/lib/consultation-thinking-plan";
|
||||
@@ -238,7 +240,7 @@ type Usage = {
|
||||
function mergeUsage(usages: Promise<Usage>[]): Promise<Usage> {
|
||||
return Promise.all(usages).then((items) => items.reduce<Usage>((total, item) => {
|
||||
const usage = item && typeof item === "object" ? item as Record<string, unknown> : {};
|
||||
const cache = promptCacheUsage(usage);
|
||||
const cache = item.cache ?? promptCacheUsage(usage);
|
||||
return {
|
||||
inputTokens: (total.inputTokens ?? 0) + (typeof usage.inputTokens === "number" ? usage.inputTokens : 0),
|
||||
outputTokens: (total.outputTokens ?? 0) + (typeof usage.outputTokens === "number" ? usage.outputTokens : 0),
|
||||
@@ -583,6 +585,74 @@ export async function POST(request: Request) {
|
||||
);
|
||||
}
|
||||
|
||||
const usageStartedAt = Date.now();
|
||||
let classificationUsage: SmalltalkUsage = {};
|
||||
let classificationOutcome: string | undefined;
|
||||
const turn = parsed.data.entrypoint === undefined
|
||||
? await classifyConsultationTurn({
|
||||
model: selectedModel,
|
||||
question: visibleQuestion,
|
||||
history: chatSession.messages,
|
||||
name: parsed.data.name,
|
||||
signal: request.signal,
|
||||
onObservation(observation) {
|
||||
if (!observation.late) {
|
||||
classificationUsage = observation.usage ?? {};
|
||||
classificationOutcome = observation.outcome;
|
||||
}
|
||||
const inputTokens = observation.usage?.inputTokens ?? 0;
|
||||
const outputTokens = observation.usage?.outputTokens ?? 0;
|
||||
logAgentObservability({
|
||||
requestId, sessionId,
|
||||
agentVersion: "consultation-smalltalk-v1",
|
||||
modelVersion: String(selectedModel.configVersion),
|
||||
policyVersion: "consultation-smalltalk-v1",
|
||||
toolCalls: [],
|
||||
contractPhases: [{
|
||||
phase: observation.late ? "classification.late_usage" : `classification.${observation.outcome}`,
|
||||
durationMs: observation.durationMs,
|
||||
status: observation.outcome === "smalltalk" || observation.outcome === "consult" ? "completed" : "failed",
|
||||
}],
|
||||
...(observation.usage ? { inputTokens, outputTokens,
|
||||
costMicrousd: Math.round((inputTokens * (selectedModel.inputCostMicrousdPerMillion ?? 0)
|
||||
+ outputTokens * (selectedModel.outputCostMicrousdPerMillion ?? 0)) / 1_000_000) } : {}),
|
||||
billingSettlementResult: "not_applicable",
|
||||
});
|
||||
},
|
||||
})
|
||||
: { kind: "consult" as const };
|
||||
|
||||
if (turn.kind === "smalltalk") {
|
||||
return streamSmalltalkResponse({
|
||||
requestId,
|
||||
reply: turn.reply,
|
||||
async complete() {
|
||||
const actualUsage = await usagePayload(Promise.resolve({}));
|
||||
const completion = await retryDetachedSettlement(async () => {
|
||||
const { data, error } = await accounting.rpc("complete_consultation_free", {
|
||||
p_user_id: userId,
|
||||
p_request_id: requestId,
|
||||
p_session_id: sessionId,
|
||||
p_response_message: { role: "assistant", text: turn.reply, responseKind: "smalltalk" },
|
||||
p_actual_usage: actualUsage,
|
||||
});
|
||||
const result = consultationCompletionSchema.safeParse(first(data ?? []));
|
||||
if (error || !result.success || !result.data.success) {
|
||||
throw new CreditRpcError(error?.message || (result.success ? result.data.error_code : "invalid_free_completion") || "free_completion_rejected");
|
||||
}
|
||||
return result.data;
|
||||
});
|
||||
logAgentObservability({ requestId, sessionId, agentVersion: "consultation-smalltalk-v1",
|
||||
billingSettlementResult: completion.success ? "completed" : "failed" });
|
||||
},
|
||||
async onError() {
|
||||
await cancel();
|
||||
logAgentObservability({ requestId, sessionId, agentVersion: "consultation-smalltalk-v1",
|
||||
billingSettlementResult: "failed", errorCode: "settlement_failed" });
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
const expectedTitle = typeof chatSession.title === "string" ? chatSession.title : "";
|
||||
const { data: titleRow } = await supabase
|
||||
.from("chat_sessions")
|
||||
@@ -591,7 +661,6 @@ export async function POST(request: Request) {
|
||||
.eq("user_id", userId)
|
||||
.maybeSingle();
|
||||
|
||||
const usageStartedAt = Date.now();
|
||||
async function checkpointConsultationContext() {
|
||||
try {
|
||||
const { data: sessionRow, error } = await supabase
|
||||
@@ -628,9 +697,9 @@ export async function POST(request: Request) {
|
||||
}
|
||||
}
|
||||
async function usagePayload(usage: Promise<{ inputTokens?: number; outputTokens?: number }>) {
|
||||
const resolved = await usage;
|
||||
const resolved = await mergeUsage([usage, Promise.resolve(classificationUsage)]);
|
||||
const usageRecord = resolved as Record<string, unknown>;
|
||||
const cache = promptCacheUsage(usageRecord);
|
||||
const cache = resolved.cache ?? promptCacheUsage(usageRecord);
|
||||
const inputTokens = Math.max(0, Math.trunc(typeof usageRecord.inputTokens === "number" ? usageRecord.inputTokens : 0));
|
||||
const outputTokens = Math.max(0, Math.trunc(typeof usageRecord.outputTokens === "number" ? usageRecord.outputTokens : 0));
|
||||
return {
|
||||
@@ -644,7 +713,15 @@ export async function POST(request: Request) {
|
||||
+ outputTokens * (selectedModel.outputCostMicrousdPerMillion ?? 0)
|
||||
) / 1_000_000),
|
||||
durationMs: Date.now() - usageStartedAt,
|
||||
...(cache ? { metadata: { cache: { ...cache, hit: cache.readTokens > 0 } } } : {}),
|
||||
metadata: {
|
||||
...(cache ? { cache: { ...cache, hit: cache.readTokens > 0 } } : {}),
|
||||
...(classificationOutcome ? { classification: {
|
||||
outcome: classificationOutcome,
|
||||
usageKnown: classificationUsage.inputTokens !== undefined || classificationUsage.outputTokens !== undefined,
|
||||
inputTokens: classificationUsage.inputTokens ?? null,
|
||||
outputTokens: classificationUsage.outputTokens ?? null,
|
||||
} } : {}),
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user