fix(web): keep consultation conclusions across turns and surface cache hits (BUG-555, BUG-556)

Session history was silently clipped to the first 4000 characters of the last 12 messages, so follow-ups could not see timing or audit tables. Keep an append-only tail plus a checkpoint summary, retry overflow in the same request, and expose cache hit rate in admin usage.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
Jesse_Chen
2026-09-06 15:16:18 +08:00
co-authored by Cursor
parent 9ce7a374ed
commit bf8ad0d1ff
29 changed files with 1295 additions and 100 deletions
@@ -65,6 +65,34 @@ function metric(
};
}
type CacheRow = {
days: string;
actual_model_id: string | null;
runs_with_cache: string;
hits: string;
read_tokens: string;
write_tokens: string;
no_cache_tokens: string;
};
function cacheEntry(row: CacheRow) {
const runsWithCache = Number(row.runs_with_cache);
const hits = Number(row.hits);
const readTokens = Number(row.read_tokens);
const writeTokens = Number(row.write_tokens);
const noCacheTokens = Number(row.no_cache_tokens);
const billed = readTokens + writeTokens + noCacheTokens;
return {
actualModelId: row.actual_model_id,
runsWithCache,
hitRate: runsWithCache > 0 ? hits / runsWithCache : null,
readTokens,
writeTokens,
noCacheTokens,
cacheShare: billed > 0 ? readTokens / billed : null,
};
}
export async function GET() {
try {
await requirePermission("billing.orders.read");
@@ -99,6 +127,26 @@ export async function GET() {
group by f.feature_key
order by f.feature_key
`, [FEATURE_KEYS]);
const cacheRows = await queryAdminRows<CacheRow>(`
with windows(days) as (
select 7
union all
select 30
)
select w.days::text as days,
l.actual_model_id,
count(*)::text as runs_with_cache,
count(*) filter (where coalesce((l.metadata->'cache'->>'readTokens')::numeric, 0) > 0)::text as hits,
coalesce(sum(coalesce((l.metadata->'cache'->>'readTokens')::numeric, 0)), 0)::text as read_tokens,
coalesce(sum(coalesce((l.metadata->'cache'->>'writeTokens')::numeric, 0)), 0)::text as write_tokens,
coalesce(sum(coalesce((l.metadata->'cache'->>'noCacheTokens')::numeric, 0)), 0)::text as no_cache_tokens
from windows w
join public.usage_ledger l
on l.created_at >= now() - make_interval(days => w.days)
and l.metadata ? 'cache'
group by w.days, l.actual_model_id
order by w.days, l.actual_model_id
`);
return NextResponse.json({
window: { days: 30, since: new Date(Date.now() - 30 * 24 * 60 * 60 * 1000).toISOString() },
@@ -112,6 +160,10 @@ export async function GET() {
durationMs: metric(row, "avg_duration_ms", "p50_duration_ms", "p95_duration_ms", "max_duration_ms"),
},
})),
cache: {
days7: cacheRows.filter((row) => row.days === "7").map(cacheEntry),
days30: cacheRows.filter((row) => row.days === "30").map(cacheEntry),
},
});
} catch (error) {
return adminErrorResponse(error);
+2 -2
View File
@@ -3,5 +3,5 @@ import { requirePermission } from "@/lib/admin/auth";
import { pageOffset, queryAdminRows } from "@/lib/admin/database";
import { adminErrorResponse, invalidQueryResponse, parseListQuery } from "@/lib/admin/http";
export const runtime="nodejs";
type Row={id:string;user_id:string;email:string|null;request_id:string;feature_key:string;source:string;requested_model_id:string|null;actual_model_id:string|null;model_config_version:number|null;input_tokens:number;output_tokens:number;cost_microusd:string;duration_ms:number|null;created_at:Date;total_count:string};
export async function GET(request:Request){try{await requirePermission("billing.orders.read");const p=parseListQuery(request);if(!p.success)return invalidQueryResponse(p.error.flatten());const q=p.data.q?`%${p.data.q}%`:null;const rows=await queryAdminRows<Row>(`select l.id,l.user_id,u.email,l.request_id,l.feature_key,l.source,l.requested_model_id,l.actual_model_id,l.model_config_version,l.input_tokens,l.output_tokens,l.cost_microusd::text,l.duration_ms,l.created_at,count(*) over()::text total_count from public.usage_ledger l left join identity.users u on u.id=l.user_id where ($1::text is null or u.email ilike $1 or l.request_id ilike $1 or l.actual_model_id ilike $1) and ($2::text is null or l.source=$2 or l.feature_key=$2) order by l.created_at desc limit $3 offset $4`,[q,p.data.status??null,p.data.pageSize,pageOffset(p.data.page,p.data.pageSize)]);return NextResponse.json({data:rows.map(r=>({id:r.id,userId:r.user_id,email:r.email,requestId:r.request_id,featureKey:r.feature_key,source:r.source,requestedModelId:r.requested_model_id,actualModelId:r.actual_model_id,modelConfigVersion:r.model_config_version,inputTokens:r.input_tokens,outputTokens:r.output_tokens,costMicrousd:Number(r.cost_microusd),durationMs:r.duration_ms,createdAt:r.created_at.toISOString()})),total:Number(rows[0]?.total_count??0)});}catch(e){return adminErrorResponse(e)}}
type Row={id:string;user_id:string;email:string|null;request_id:string;feature_key:string;source:string;requested_model_id:string|null;actual_model_id:string|null;model_config_version:number|null;input_tokens:number;output_tokens:number;cost_microusd:string;duration_ms:number|null;created_at:Date;cache_read_tokens:string|null;total_count:string};
export async function GET(request:Request){try{await requirePermission("billing.orders.read");const p=parseListQuery(request);if(!p.success)return invalidQueryResponse(p.error.flatten());const q=p.data.q?`%${p.data.q}%`:null;const rows=await queryAdminRows<Row>(`select l.id,l.user_id,u.email,l.request_id,l.feature_key,l.source,l.requested_model_id,l.actual_model_id,l.model_config_version,l.input_tokens,l.output_tokens,l.cost_microusd::text,l.duration_ms,l.created_at,l.metadata->'cache'->>'readTokens' as cache_read_tokens,count(*) over()::text total_count from public.usage_ledger l left join identity.users u on u.id=l.user_id where ($1::text is null or u.email ilike $1 or l.request_id ilike $1 or l.actual_model_id ilike $1) and ($2::text is null or l.source=$2 or l.feature_key=$2) order by l.created_at desc limit $3 offset $4`,[q,p.data.status??null,p.data.pageSize,pageOffset(p.data.page,p.data.pageSize)]);return NextResponse.json({data:rows.map(r=>({id:r.id,userId:r.user_id,email:r.email,requestId:r.request_id,featureKey:r.feature_key,source:r.source,requestedModelId:r.requested_model_id,actualModelId:r.actual_model_id,modelConfigVersion:r.model_config_version,inputTokens:r.input_tokens,outputTokens:r.output_tokens,costMicrousd:Number(r.cost_microusd),durationMs:r.duration_ms,cacheReadTokens:r.cache_read_tokens==null?null:Number(r.cache_read_tokens),createdAt:r.created_at.toISOString()})),total:Number(rows[0]?.total_count??0)});}catch(e){return adminErrorResponse(e)}}
+168 -57
View File
@@ -29,10 +29,10 @@ import {
shouldLoadGeneralDailyPanchanga,
} from "@/lib/consultation-entrypoint";
import { CreditRpcError } from "@/lib/consultation-billing";
import { cachedSystemMessage, mergePromptCacheUsage, promptCacheUsage } from "@/lib/agent-generation-settings";
import { cachedHistoryMessage, cachedSystemMessage, mergePromptCacheUsage, promptCacheUsage } from "@/lib/agent-generation-settings";
import { FeaturePricingError, resolveFeaturePricing } from "@/lib/feature-pricing";
import { reserveConsultationModel } from "@/lib/consultation-model-selection";
import { resolveSessionLanguageModel } from "@/lib/model-catalog";
import { resolveSessionLanguageModel, loadLanguageModelCatalog } from "@/lib/model-catalog";
import { jsonForSupabaseSetupFailure } from "@/lib/api/service-unavailable";
import { createAdminSupabaseClient } from "@/lib/supabase/admin";
import { createServerSupabaseClient } from "@/lib/supabase/server";
@@ -72,7 +72,17 @@ import {
loadGeneralDailyPanchangaContext,
type GeneralDailyPanchangaContext,
} from "@/lib/general-daily-panchanga";
import { consultationHistoryFromStoredMessages } from "@/lib/consultation-session-history";
import {
consultationHistoryWindow,
consultationUserTurnContent,
isContextOverflowError,
lastConsultationPair,
parseSessionContextSummary,
} from "@/lib/consultation-session-history";
import {
checkpointSessionContextSummary,
generateSessionContextSummaryText,
} from "@/lib/session-context-summary";
import { generateSessionTitle, shouldGenerateSessionTitle } from "@/lib/session-title-agent";
import { z } from "zod";
@@ -279,7 +289,7 @@ export async function POST(request: Request) {
const { data: chatSession, error: chatSessionError } = await supabase
.from("chat_sessions")
.select("id,model_id,model_config_version,session_type,messages,title,theme,chart_profile_role")
.select("id,model_id,model_config_version,session_type,messages,title,theme,chart_profile_role,context_summary")
.eq("id", parsed.data.sessionId)
.eq("user_id", user.id)
.maybeSingle();
@@ -336,11 +346,16 @@ export async function POST(request: Request) {
const visibleQuestion = parsed.data.question;
// Client `history` stays in the request schema for old bundles and is not read.
const storedHistory = consultationHistoryFromStoredMessages(chatSession.messages);
const contextSummary = parseSessionContextSummary(chatSession.context_summary);
const historyWindow = consultationHistoryWindow(chatSession.messages, contextSummary, {
contextWindow: sessionModel.contextWindow,
});
const storedHistory = historyWindow.tail;
const userControlledPrompt = [
parsed.data.question,
historyWindow.summaryText,
...storedHistory.filter((message) => message.role === "user").map((message) => message.text),
].join("\n");
].filter(Boolean).join("\n");
if (blocksPromptExtraction(userControlledPrompt)) {
return NextResponse.json(
{
@@ -551,6 +566,40 @@ export async function POST(request: Request) {
}
const usageStartedAt = Date.now();
async function checkpointConsultationContext() {
try {
const { data: sessionRow, error } = await supabase
.from("chat_sessions")
.select("messages, context_summary")
.eq("id", sessionId)
.eq("user_id", userId)
.maybeSingle();
if (error || !sessionRow) return;
const catalog = await loadLanguageModelCatalog();
const summaryModel = catalog.defaultModelId
? catalog.models.find((model) => model.id === catalog.defaultModelId) ?? selectedModel
: selectedModel;
await checkpointSessionContextSummary({
messages: sessionRow.messages,
summary: sessionRow.context_summary,
generateText: (prompt, signal) => generateSessionContextSummaryText(summaryModel, prompt, signal),
update: async (summary, seenUpdatedAt) => {
let query = supabase.from("chat_sessions")
.update({ context_summary: summary })
.eq("id", sessionId)
.eq("user_id", userId);
query = seenUpdatedAt
? query.eq("context_summary->>updatedAt", seenUpdatedAt)
: query.is("context_summary", null);
const { data, error: writeError } = await query.select("id");
if (writeError) throw writeError;
return Boolean(data?.length);
},
});
} catch (error) {
console.warn("session context summary failed", error);
}
}
async function usagePayload(usage: Promise<{ inputTokens?: number; outputTokens?: number }>) {
const resolved = await usage;
const usageRecord = resolved as Record<string, unknown>;
@@ -615,6 +664,9 @@ export async function POST(request: Request) {
if (!completion.success && completion.error_code !== "request_cancelled") {
throw new CreditRpcError(completion.error_code || "completion_rejected");
}
if (completion.success) {
void checkpointConsultationContext();
}
return "completed";
} catch (error) {
await cancel();
@@ -775,26 +827,36 @@ export async function POST(request: Request) {
}
};
const cacheBoundary = cachedSystemMessage("【上下文缓存边界】后续内容为本轮请求输入。", selectedModel.model);
const baseMessages = [
...(cacheBoundary ? [cacheBoundary] : []),
...history.map((message) => message.role === "user"
const modeInstruction = consultationMode === "general_no_birth_time" || (consultationMode === "declared_birth_window" && generalDailyContext)
? generalNoMinuteInstruction(Boolean(generalDailyContext))
: consultationMode === "declared_birth_window"
? declaredWindowInstruction()
: "先加载 Jyotish Skill;如需新的个人星盘结论,必须调用服务器绑定的排盘工具。事业/财富/婚恋/家庭按 skill Level 2 模板写:原始结构、六步宫位、Yoga 表、时机、综合、文末技法审计表,然后才是现代生活措辞。不要复述内部 JSON 字段。";
const consultationBaseMessages = (overflow: boolean) => {
const tail = overflow ? lastConsultationPair(history) : history;
const mapped = tail.map((message) => message.role === "user"
? { role: "user" as const, content: message.text }
: { role: "assistant" as const, content: message.text }),
{
role: "user" as const,
content: [
currentTimeContext(requestTime),
name ? `用户称呼:${name}` : "",
consultationMode === "general_no_birth_time" || (consultationMode === "declared_birth_window" && generalDailyContext)
? generalNoMinuteInstruction(Boolean(generalDailyContext))
: consultationMode === "declared_birth_window"
? declaredWindowInstruction()
: "先加载 Jyotish Skill;如需新的个人星盘结论,必须调用服务器绑定的排盘工具。事业/财富/婚恋/家庭按 skill Level 2 模板写:原始结构、六步宫位、Yoga 表、时机、综合、文末技法审计表,然后才是现代生活措辞。不要复述内部 JSON 字段。",
generalDailyContextPrompt(generalDailyContext),
resolvedQuestion.modelQuestion,
].filter(Boolean).join("\n"),
},
];
: { role: "assistant" as const, content: message.text });
const cachedTail = mapped.length > 0
? [...mapped.slice(0, -1), cachedHistoryMessage(mapped[mapped.length - 1]!, selectedModel.model)]
: mapped;
return [
...(cacheBoundary ? [cacheBoundary] : []),
...cachedTail,
{
role: "user" as const,
content: consultationUserTurnContent({
currentTime: currentTimeContext(requestTime),
name,
instruction: modeInstruction,
extra: generalDailyContextPrompt(generalDailyContext),
summaryText: historyWindow.summaryText,
question: resolvedQuestion.modelQuestion,
}),
},
];
};
let baseMessages = consultationBaseMessages(false);
const agentAbortSignal = AbortSignal.timeout(AGENT_TIMEOUT_MS);
const streamOptions = {
runId: requestId,
@@ -803,6 +865,27 @@ export async function POST(request: Request) {
hooks,
...consultationGenerationSettings(selectedModel.model),
};
async function streamWithOverflowRetry(agent: {
stream: (
messages: typeof baseMessages,
options: typeof streamOptions,
) => Promise<{ fullStream: AsyncIterable<unknown> | ReadableStream<unknown>; totalUsage: Promise<Usage> }>;
}) {
try {
const result = await agent.stream(baseMessages, streamOptions);
usages.push(result.totalUsage);
return result;
} catch (error) {
if (!isContextOverflowError(error)) throw error;
// Same request, same wait: shrink to summary + last pair. Consultation
// clients cannot parse rectification `attempt.reset`, so the retry stays
// server-side and never opens a second user-visible wait.
baseMessages = consultationBaseMessages(true);
const overflow = await agent.stream(baseMessages, streamOptions);
usages.push(overflow.totalUsage);
return overflow;
}
}
const workflowReceipt: WorkflowReceipt = usesPublicDailyGeneralAgent(consultationMode, generalDailyContext)
? {
route: generalDailyContext ? "general-daily-panchanga" : "general-no-birth-time",
@@ -822,8 +905,7 @@ export async function POST(request: Request) {
if (usesPublicDailyGeneralAgent(consultationMode, generalDailyContext)) {
state.workflowReceipt = workflowReceipt;
const agent = getGeneralJyotishAgent(selectedModel);
const result = await agent.stream(baseMessages, streamOptions);
usages.push(result.totalUsage);
const result = await streamWithOverflowRetry(agent);
// This mode has no calculation to require and no chart method to bind,
// so there is no contract for a retry to repair.
const retryForAnswer = async () => {
@@ -913,8 +995,7 @@ export async function POST(request: Request) {
state,
});
const agent = getWindowJyotishAgent(selectedModel, agentContext);
const result = await agent.stream(baseMessages, streamOptions);
usages.push(result.totalUsage);
const result = await streamWithOverflowRetry(agent);
const retry = async () => {
const retried = await agent.stream([
...baseMessages,
@@ -1013,8 +1094,7 @@ export async function POST(request: Request) {
state,
});
const agent = getJyotishAgent(selectedModel, agentContext);
const result = await agent.stream(baseMessages, streamOptions);
usages.push(result.totalUsage);
const result = await streamWithOverflowRetry(agent);
const retry = async () => {
const retried = await agent.stream([
...baseMessages,
@@ -1138,20 +1218,21 @@ export async function POST(request: Request) {
return await runAgenticConsultation(consultationMode, history, name, generalDailyContext);
}
if (!shouldRunBirthChartWorkflow(consultationMode)) {
const cacheBoundary = cachedSystemMessage("【上下文缓存边界】后续内容为本轮请求输入。", selectedModel.model);
const result = await getGeneralJyotishAgent(selectedModel).stream([
...(cacheBoundary ? [cacheBoundary] : []),
{
role: "user",
content: [
currentTimeContext(requestTime),
name ? `用户称呼:${name}` : "",
generalNoMinuteInstruction(Boolean(generalDailyContext)),
generalDailyContextPrompt(generalDailyContext),
resolvedQuestion.modelQuestion,
].filter(Boolean).join("\n"),
},
]);
const cacheBoundary = cachedSystemMessage("【上下文缓存边界】后续内容为本轮请求输入。", selectedModel.model);
const result = await getGeneralJyotishAgent(selectedModel).stream([
...(cacheBoundary ? [cacheBoundary] : []),
{
role: "user",
content: consultationUserTurnContent({
currentTime: currentTimeContext(requestTime),
name,
instruction: generalNoMinuteInstruction(Boolean(generalDailyContext)),
extra: generalDailyContextPrompt(generalDailyContext),
summaryText: historyWindow.summaryText,
question: resolvedQuestion.modelQuestion,
}),
},
]);
const workflowReceipt: WorkflowReceipt = {
route: generalDailyContext ? "general-daily-panchanga" : "general-no-birth-time",
status: "ready",
@@ -1216,21 +1297,51 @@ export async function POST(request: Request) {
const workflowReceipt = consultationWorkflowReceipt(workflowContext);
const cacheBoundary = cachedSystemMessage("【上下文缓存边界】后续内容为本轮请求输入。", selectedModel.model);
const result = await getLegacyJyotishAgent(selectedModel, workflowContext).stream([
const legacyHistory = history.map((message) => message.role === "user"
? { role: "user" as const, content: message.text }
: { role: "assistant" as const, content: message.text });
const cachedLegacyHistory = legacyHistory.length > 0
? [...legacyHistory.slice(0, -1), cachedHistoryMessage(legacyHistory[legacyHistory.length - 1]!, selectedModel.model)]
: legacyHistory;
const legacyMessages = [
...(cacheBoundary ? [cacheBoundary] : []),
...history.map((message) => message.role === "user"
? { role: "user" as const, content: message.text }
: { role: "assistant" as const, content: message.text }),
...cachedLegacyHistory,
{
role: "user",
content: [
currentTimeContext(requestTime),
name ? `用户称呼:${name}` : "",
"先用 3–6 句口语直接回答下面的问题,不要加标题;然后再按 skill Level 2 骨架写:原始结构、六步宫位、Yoga 表、时机、综合、文末技法审计表,最后才是现代生活。骨架不可省略。星盘事实只使用系统里已经注入的计算结果,不要复述内部字段、JSON 或再跑一遍咨询流程。",
resolvedQuestion.modelQuestion,
].filter(Boolean).join("\n"),
role: "user" as const,
content: consultationUserTurnContent({
currentTime: currentTimeContext(requestTime),
name,
instruction: "先用 3–6 句口语直接回答下面的问题,不要加标题;然后再按 skill Level 2 骨架写:原始结构、六步宫位、Yoga 表、时机、综合、文末技法审计表,最后才是现代生活。骨架不可省略。星盘事实只使用系统里已经注入的计算结果,不要复述内部字段、JSON 或再跑一遍咨询流程。",
summaryText: historyWindow.summaryText,
question: resolvedQuestion.modelQuestion,
}),
},
]);
];
const legacyAgent = getLegacyJyotishAgent(selectedModel, workflowContext);
let result;
try {
result = await legacyAgent.stream(legacyMessages);
} catch (error) {
if (!isContextOverflowError(error)) throw error;
const overflowHistory = lastConsultationPair(legacyHistory);
const overflowCached = overflowHistory.length > 0
? [...overflowHistory.slice(0, -1), cachedHistoryMessage(overflowHistory[overflowHistory.length - 1]!, selectedModel.model)]
: overflowHistory;
result = await legacyAgent.stream([
...(cacheBoundary ? [cacheBoundary] : []),
...overflowCached,
{
role: "user" as const,
content: consultationUserTurnContent({
currentTime: currentTimeContext(requestTime),
name,
instruction: "先用 3–6 句口语直接回答下面的问题,不要加标题;然后再按 skill Level 2 骨架写:原始结构、六步宫位、Yoga 表、时机、综合、文末技法审计表,最后才是现代生活。骨架不可省略。星盘事实只使用系统里已经注入的计算结果,不要复述内部字段、JSON 或再跑一遍咨询流程。",
summaryText: historyWindow.summaryText,
question: resolvedQuestion.modelQuestion,
}),
},
]);
}
const responseWorkflowReceipt = {
route: workflowReceipt.route,
status: workflowReceipt.status,