Merge codex/pr2-final-response into staging integration

This commit is contained in:
Jesse_Chen
2026-08-15 12:27:00 +08:00
26 changed files with 1643 additions and 73 deletions
+44 -6
View File
@@ -10,6 +10,12 @@ import {
import { blocksPromptExtraction } from "@/lib/consult-safety";
import { consultationDomainSchema } from "@/lib/consultation-domain-registry";
import { parseAgentReply } from "@/lib/agent-reply";
import { createConsultationReplyMetadata } from "@/lib/consultation-reply-metadata";
import {
ConsultationPlanValidationError,
createConsultationPlan,
type ConsultationPlan,
} from "@/lib/consultation-plan";
import {
logAgentObservability,
settlementTelemetryOutcome,
@@ -30,6 +36,7 @@ import { streamAgentResponse } from "@/lib/stream-agent-response";
import type { AgentExecutionReceipt, WorkflowReceipt } from "@/lib/consultation-agent-events";
import {
createConsultationAgentContext,
consultationStepBudgetReceipt,
createConsultationRuntimeHooks,
createConsultationRuntimeState,
} from "@/mastra/consultation-tools";
@@ -43,6 +50,7 @@ import {
import {
ConsultationProfileTruthError,
prepareConsultationRoute,
type PreparedConsultationRoute,
} from "@/lib/consultation-route-service";
import { z } from "zod";
@@ -142,9 +150,10 @@ function mergeUsage(usages: Promise<Usage>[]): Promise<Usage> {
}
function shouldUseAgenticRuntime(user: { id: string; app_metadata?: Record<string, unknown> }) {
const mode = process.env.CONSULTATION_AGENTIC_RUNTIME?.trim().toLowerCase() ?? "legacy";
const mode = process.env.CONSULTATION_AGENTIC_RUNTIME?.trim().toLowerCase() ?? "enabled";
if (mode === "legacy") return false;
if (mode === "enabled") return true;
if (mode !== "canary") return false;
if (mode !== "canary") return true;
const ids = new Set((process.env.CONSULTATION_AGENTIC_CANARY_USER_IDS ?? "")
.split(",").map((value) => value.trim()).filter(Boolean));
const roles = Array.isArray(user.app_metadata?.roles) ? user.app_metadata.roles : [];
@@ -249,6 +258,7 @@ export async function POST(request: Request) {
const requestId = parsed.data.requestId;
const sessionId = parsed.data.sessionId;
const consultationTheme = parsed.data.theme;
const visibleQuestion = parsed.data.question;
const userControlledPrompt = [
parsed.data.question,
@@ -269,9 +279,9 @@ export async function POST(request: Request) {
type SelectedModel = NonNullable<Awaited<ReturnType<typeof resolveSessionLanguageModel>>>;
type ReservationResult = { success: boolean; credits: number | null; error_code: string | null };
type ModelSelection = Awaited<ReturnType<typeof reserveConsultationModel<SelectedModel, ReservationResult>>>;
let prepared: Awaited<ReturnType<typeof prepareConsultationRoute<ModelSelection>>>;
let prepared: PreparedConsultationRoute<ModelSelection, ConsultationPlan>;
try {
prepared = await prepareConsultationRoute({
prepared = await prepareConsultationRoute<ModelSelection, ConsultationPlan>({
userId,
mode: parsed.data.consultationMode,
async loadProfile(profileUserId) {
@@ -283,6 +293,12 @@ export async function POST(request: Request) {
if (error || !data) throw new ConsultationProfileTruthError("profile_unavailable");
return data;
},
beforeReserve: ({ consultationMode }) => createConsultationPlan({
userIntent: resolvedQuestion.modelQuestion,
theme: consultationTheme,
consultationMode,
modelCreditCost: sessionModel.creditCost,
}),
reserve: () => reserveConsultationModel(
chatSession.model_id,
(modelId) => sessionModel?.id === modelId ? sessionModel : null,
@@ -307,6 +323,17 @@ export async function POST(request: Request) {
),
});
} catch (error) {
if (error instanceof ConsultationPlanValidationError) {
return NextResponse.json(
{
error: "当前咨询计划不可用",
message: error.code === "plan_cost_ceiling_exceeded"
? "当前会话模型超出本次咨询的点数上限,请选择标准模型后重试,本次不会扣点。"
: "当前咨询模式与回答精度边界不一致,请刷新后重试,本次不会扣点。",
},
{ status: 409 },
);
}
if (error instanceof ConsultationProfileTruthError) {
const modeChanged = error.code === "mode_changed";
return NextResponse.json(
@@ -414,7 +441,11 @@ export async function POST(request: Request) {
agentExecutionReceipt?: AgentExecutionReceipt,
): Promise<AgentSettlementResult> {
try {
const reply = parseAgentReply(rawTransformedText, consultationTheme);
const reply = parseAgentReply(
rawTransformedText,
consultationTheme,
createConsultationReplyMetadata({ theme: consultationTheme, question: visibleQuestion }),
);
if (!reply.text) throw new Error("empty_agent_reply");
const responseMessage = {
role: "assistant" as const,
@@ -596,6 +627,7 @@ export async function POST(request: Request) {
runtime: "mastra-agentic",
skill: { name: "jyotish-vedic-astrology", loaded: state.jyotishSkillLoaded },
steps: state.steps,
stepBudget: consultationStepBudgetReceipt(state),
workflow: workflowReceipt,
techniqueTruth: "not-applicable",
});
@@ -634,6 +666,8 @@ export async function POST(request: Request) {
sessionId,
requestId,
consultationMode,
plan: prepared.preReserveResult,
theme: consultationTheme,
serverChart: prepared.serverChart,
abortSignal: agentAbortSignal,
state,
@@ -657,6 +691,7 @@ export async function POST(request: Request) {
runtime: "mastra-agentic",
skill: { name: "jyotish-vedic-astrology", loaded: state.jyotishSkillLoaded },
steps: state.steps,
stepBudget: consultationStepBudgetReceipt(state),
workflow: state.workflowReceipt ?? workflowReceipt,
techniqueTruth: state.techniqueTruth ?? "unknown",
});
@@ -761,7 +796,10 @@ export async function POST(request: Request) {
theme: parsed.data.theme,
});
const workflowContext = applyBirthTimeModeToWorkflowContext(
await runConsultationWorkflow(toolInput, { foreground: true }),
await runConsultationWorkflow(toolInput, {
foreground: true,
plan: prepared.preReserveResult,
}),
consultationMode,
);
const workflowReceipt = consultationWorkflowReceipt(workflowContext);