refactor: simpler streaming preview and small AI SDK cleanups

- useChat throttles streamed message updates (experimental_throttle,
  150 ms), replacing the two hand-written 150 ms timers of the
  display_diagram and edit_diagram previews (94 lines less). The preview
  now only runs while the input streams; once it is complete the tool
  handler takes over, so a queued preview can no longer redraw an edit
  the handler rejected and rolled back. Measured on a streamed 60-cell
  diagram: 41 redraws at least 97 ms apart, before 37 with gaps down to
  48 ms
- The diagram check endpoint uses streamText with Output.object instead
  of the deprecated streamObject, and returns its fixed result as a plain
  text response; new route test
- Import createGateway/gateway from ai and drop the direct
  @ai-sdk/gateway dependency
- The per-request message structure logs only print with
  DEBUG_LLM_PAYLOAD=true
- Remove an empty onFinish callback
This commit is contained in:
dayuan.jiang
2026-10-04 13:16:38 +09:00
parent ac62a58c9f
commit 22a1d3f03b
10 changed files with 206 additions and 222 deletions
+66 -57
View File
@@ -95,6 +95,8 @@ function createCachedStreamResponse(xml: string): Response {
const modelStreamResponses = new WeakSet<Response>()
// Inner handler function
const DEBUG_LLM_PAYLOAD = process.env.DEBUG_LLM_PAYLOAD === "true"
async function handleChatRequest(req: Request): Promise<Response> {
// Check for access code
const accessDenied = checkAccessCode(req)
@@ -335,35 +337,37 @@ ${userInputText}
// Convert UIMessages to ModelMessages and add system message
const modelMessages = await convertToModelMessages(messages)
// DEBUG: Log incoming messages structure
console.log("[route.ts] Incoming messages count:", messages.length)
messages.forEach((msg: any, idx: number) => {
console.log(
`[route.ts] Message ${idx} role:`,
msg.role,
"parts count:",
msg.parts?.length,
)
if (msg.parts) {
msg.parts.forEach((part: any, partIdx: number) => {
if (
part.type === "tool-invocation" ||
part.type === "tool-result"
) {
console.log(`[route.ts] Part ${partIdx}:`, {
type: part.type,
toolName: part.toolName,
hasInput: !!part.input,
inputType: typeof part.input,
inputKeys:
part.input && typeof part.input === "object"
? Object.keys(part.input)
: null,
})
}
})
}
})
// DEBUG_LLM_PAYLOAD=true logs the incoming message structure
if (DEBUG_LLM_PAYLOAD) {
console.log("[route.ts] Incoming messages count:", messages.length)
messages.forEach((msg: any, idx: number) => {
console.log(
`[route.ts] Message ${idx} role:`,
msg.role,
"parts count:",
msg.parts?.length,
)
if (msg.parts) {
msg.parts.forEach((part: any, partIdx: number) => {
if (
part.type === "tool-invocation" ||
part.type === "tool-result"
) {
console.log(`[route.ts] Part ${partIdx}:`, {
type: part.type,
toolName: part.toolName,
hasInput: !!part.input,
inputType: typeof part.input,
inputKeys:
part.input && typeof part.input === "object"
? Object.keys(part.input)
: null,
})
}
})
}
})
}
// Replace historical tool call XML with placeholders to reduce tokens
// Disabled by default - some models (e.g. minimax) copy placeholders instead of generating XML
@@ -385,34 +389,39 @@ ${userInputText}
// JSON object, and every provider rejects a tool result whose call is gone.
enhancedMessages = dropInvalidToolCalls(enhancedMessages)
// DEBUG: Log modelMessages structure (what's being sent to AI)
console.log("[route.ts] Model messages count:", enhancedMessages.length)
enhancedMessages.forEach((msg: any, idx: number) => {
console.log(
`[route.ts] ModelMsg ${idx} role:`,
msg.role,
"content count:",
msg.content?.length,
)
if (msg.content) {
msg.content.forEach((part: any, partIdx: number) => {
if (part.type === "tool-call" || part.type === "tool-result") {
console.log(`[route.ts] Content ${partIdx}:`, {
type: part.type,
toolName: part.toolName,
hasInput: !!part.input,
inputType: typeof part.input,
inputValue:
part.input === undefined
? "undefined"
: part.input === null
? "null"
: "object",
})
}
})
}
})
// DEBUG_LLM_PAYLOAD=true logs what is sent to the model
if (DEBUG_LLM_PAYLOAD) {
console.log("[route.ts] Model messages count:", enhancedMessages.length)
enhancedMessages.forEach((msg: any, idx: number) => {
console.log(
`[route.ts] ModelMsg ${idx} role:`,
msg.role,
"content count:",
msg.content?.length,
)
if (msg.content) {
msg.content.forEach((part: any, partIdx: number) => {
if (
part.type === "tool-call" ||
part.type === "tool-result"
) {
console.log(`[route.ts] Content ${partIdx}:`, {
type: part.type,
toolName: part.toolName,
hasInput: !!part.input,
inputType: typeof part.input,
inputValue:
part.input === undefined
? "undefined"
: part.input === null
? "null"
: "object",
})
}
})
}
})
}
// Update the last message with user input only (XML moved to separate cached system message)
if (enhancedMessages.length >= 1) {
+8 -19
View File
@@ -3,7 +3,7 @@
* Accepts a PNG image and streams validation results using useObject-compatible format.
*/
import { streamObject } from "ai"
import { Output, streamText } from "ai"
import { checkAccessCode } from "@/lib/access-code"
import { getValidationModel } from "@/lib/ai-providers"
import { VALIDATION_SYSTEM_PROMPT } from "@/lib/validation-prompts"
@@ -29,20 +29,9 @@ const DEFAULT_VALID_RESULT: ValidationResult = {
suggestions: [],
}
/**
* Create a streaming response for useObject compatibility.
* useObject expects text stream format, not plain JSON.
*/
/** A fixed result in the text format useObject reads */
function createStreamingResponse(result: ValidationResult): Response {
const encoder = new TextEncoder()
const stream = new ReadableStream({
start(controller) {
// Stream the JSON as text (useObject parses this)
controller.enqueue(encoder.encode(JSON.stringify(result)))
controller.close()
},
})
return new Response(stream, {
return new Response(JSON.stringify(result), {
headers: { "Content-Type": "text/plain; charset=utf-8" },
})
}
@@ -108,9 +97,9 @@ export async function POST(req: Request): Promise<Response> {
) || 10000
// Stream the VLM response for useObject consumption
const result = streamObject({
const result = streamText({
model,
schema: ValidationResultSchema,
output: Output.object({ schema: ValidationResultSchema }),
system: VALIDATION_SYSTEM_PROMPT,
messages: [
{
@@ -129,10 +118,10 @@ export async function POST(req: Request): Promise<Response> {
],
maxOutputTokens: 1024,
abortSignal: AbortSignal.timeout(timeout),
onFinish: ({ object }) => {
if (sessionId && object) {
onFinish: ({ output }) => {
if (sessionId && output) {
console.log(
`[validate-diagram] Session ${sessionId}: valid=${object.valid}, issues=${object.issues?.length ?? 0}`,
`[validate-diagram] Session ${sessionId}: valid=${output.valid}, issues=${output.issues?.length ?? 0}`,
)
}
},
+1 -2
View File
@@ -1,13 +1,12 @@
import { createAmazonBedrock } from "@ai-sdk/amazon-bedrock"
import { createAnthropic } from "@ai-sdk/anthropic"
import { createDeepSeek, deepseek } from "@ai-sdk/deepseek"
import { createGateway } from "@ai-sdk/gateway"
import { createGoogleGenerativeAI } from "@ai-sdk/google"
import { createVertex } from "@ai-sdk/google-vertex"
import { createOpenAI } from "@ai-sdk/openai"
import { createAihubmix } from "@aihubmix/ai-sdk-provider"
import { createOpenRouter } from "@openrouter/ai-sdk-provider"
import { generateText } from "ai"
import { createGateway, generateText } from "ai"
import { NextResponse } from "next/server"
import { createOllama } from "ollama-ai-provider-v2"
import { checkAccessCode } from "@/lib/access-code"
+42 -137
View File
@@ -210,20 +210,6 @@ export function ChatMessageDisplay({
scrollTopRef.current?.scrollIntoView({ behavior: "instant" })
}
}, [messages.length, processedToolCalls])
// Debounce streaming diagram updates - store pending XML and timeout
const pendingXmlRef = useRef<string | null>(null)
const debounceTimeoutRef = useRef<ReturnType<typeof setTimeout> | null>(
null,
)
const STREAMING_DEBOUNCE_MS = 150 // Only update diagram every 150ms during streaming
// Refs for edit_diagram streaming
const pendingEditRef = useRef<{
operations: DiagramOperation[]
toolCallId: string
} | null>(null)
const editDebounceTimeoutRef = useRef<ReturnType<typeof setTimeout> | null>(
null,
)
const [expandedTools, setExpandedTools] = useState<Record<string, boolean>>(
{},
)
@@ -454,42 +440,17 @@ export function ChatMessageDisplay({
return // Skip redundant processing
}
// Messages update at most every 150 ms while
// streaming (useChat throttle in chat-panel)
if (state === "input-streaming") {
// Debounce streaming updates - queue the XML and process after delay
pendingXmlRef.current = xml
if (!debounceTimeoutRef.current) {
// No pending timeout - set one up
debounceTimeoutRef.current = setTimeout(
() => {
const pendingXml =
pendingXmlRef.current
debounceTimeoutRef.current = null
pendingXmlRef.current = null
if (pendingXml) {
handleDisplayChart(pendingXml)
lastProcessedXmlRef.current.set(
toolCallId,
pendingXml,
)
}
},
STREAMING_DEBOUNCE_MS,
)
}
handleDisplayChart(xml)
lastProcessedXmlRef.current.set(toolCallId, xml)
} else if (
!processedToolCalls.current.has(toolCallId)
) {
// Input complete: the tool handler loads the
// validated diagram, so drop a queued preview
// that would draw the raw cells over it
if (debounceTimeoutRef.current) {
clearTimeout(debounceTimeoutRef.current)
debounceTimeoutRef.current = null
pendingXmlRef.current = null
}
// validated diagram
processedToolCalls.current.add(toolCallId)
// Clean up the ref entry - tool is complete, no longer needed
lastProcessedXmlRef.current.delete(toolCallId)
}
}
@@ -500,19 +461,10 @@ export function ChatMessageDisplay({
part.type === "tool-edit_diagram" &&
input?.operations
) {
// Failed or stopped: drop the queued preview. If the original
// XML is still stored, the tool handler never ran (user pressed
// Failed or stopped: if the original XML is still
// stored, the tool handler never ran (user pressed
// stop), so undo the streamed preview here.
if (state === "output-error") {
if (
pendingEditRef.current?.toolCallId ===
toolCallId &&
editDebounceTimeoutRef.current
) {
clearTimeout(editDebounceTimeoutRef.current)
editDebounceTimeoutRef.current = null
pendingEditRef.current = null
}
const originalXml =
editDiagramOriginalXmlRef.current.get(
toolCallId,
@@ -526,10 +478,23 @@ export function ChatMessageDisplay({
return
}
if (state !== "input-streaming") {
// Input complete: the tool handler applies the
// checked edit (it reads the original XML too)
if (
!processedToolCalls.current.has(toolCallId)
) {
lastProcessedXmlRef.current.delete(
toolCallId + "-opCount",
)
processedToolCalls.current.add(toolCallId)
}
return
}
const completeOps = getCompleteOperations(
input.operations as DiagramOperation[],
)
if (completeOps.length === 0) return
// Capture original XML when streaming starts (store in shared ref)
@@ -549,7 +514,6 @@ export function ChatMessageDisplay({
chartXML,
)
}
const originalXml =
editDiagramOriginalXmlRef.current.get(
toolCallId,
@@ -557,95 +521,36 @@ export function ChatMessageDisplay({
if (!originalXml) return
// Skip if no change from last processed state
const lastCount = lastProcessedXmlRef.current.get(
toolCallId + "-opCount",
)
if (lastCount === String(completeOps.length)) return
const countKey = `${toolCallId}-opCount`
const opCount = String(completeOps.length)
if (
state === "input-streaming" ||
state === "input-available"
lastProcessedXmlRef.current.get(countKey) ===
opCount
) {
// Queue the operations for debounced processing
pendingEditRef.current = {
operations: completeOps,
toolCallId,
}
if (!editDebounceTimeoutRef.current) {
editDebounceTimeoutRef.current = setTimeout(
() => {
const pending =
pendingEditRef.current
editDebounceTimeoutRef.current =
null
pendingEditRef.current = null
if (pending) {
const origXml =
editDiagramOriginalXmlRef.current.get(
pending.toolCallId,
)
if (!origXml) return
try {
const {
result: editedXml,
} = applyDiagramOperations(
origXml,
pending.operations,
)
// Load the full document so other pages stay intact
onDisplayChart(
editedXml,
true,
)
lastProcessedXmlRef.current.set(
pending.toolCallId +
"-opCount",
String(
pending.operations
.length,
),
)
} catch (e) {
console.warn(
`[edit_diagram streaming] Operation failed:`,
e instanceof Error
? e.message
: e,
)
}
}
},
STREAMING_DEBOUNCE_MS,
)
}
} else if (
state === "output-available" &&
!processedToolCalls.current.has(toolCallId)
) {
// Final state - cleanup streaming refs (tool handler does final application)
if (editDebounceTimeoutRef.current) {
clearTimeout(editDebounceTimeoutRef.current)
editDebounceTimeoutRef.current = null
}
lastProcessedXmlRef.current.delete(
toolCallId + "-opCount",
return
}
try {
const { result } = applyDiagramOperations(
originalXml,
completeOps,
)
// Load the full document so other pages stay intact
onDisplayChart(result, true)
lastProcessedXmlRef.current.set(
countKey,
opCount,
)
} catch (e) {
console.warn(
"[edit_diagram streaming] Operation failed:",
e instanceof Error ? e.message : e,
)
processedToolCalls.current.add(toolCallId)
// Note: Don't delete editDiagramOriginalXmlRef here - tool handler needs it
}
}
}
})
}
})
// NOTE: Don't cleanup debounce timeouts here!
// The cleanup runs on every re-render (when messages changes),
// which would cancel the timeout before it fires.
// Let the timeouts complete naturally - they're harmless if component unmounts.
}, [messages, handleDisplayChart, chartXML])
return (
+3 -2
View File
@@ -472,7 +472,9 @@ export default function ChatPanel({
setShowSettingsDialog(true)
}
},
onFinish: () => {},
// Re-render streamed messages at most every 150 ms. The streaming
// diagram preview draws on each update, so this also limits redraws
experimental_throttle: 150,
sendAutomaticallyWhen: ({ messages }) => {
const isInContinuationMode = partialXmlRef.current.length > 0
@@ -883,7 +885,6 @@ export default function ChatPanel({
await sendWithCurrentDiagram(parts)
// Token count is tracked in onFinish with actual server usage
setInput("")
sessionStorage.removeItem(SESSION_STORAGE_INPUT_KEY)
setFiles([])
+6 -1
View File
@@ -123,9 +123,14 @@ AI_MODEL=global.anthropic.claude-sonnet-4-5-20250929-v1:0
# Temperature (Optional)
# Controls randomness in AI responses. Lower = more deterministic.
# Leave unset for models that don't support temperature (e.g., GPT-5.1 reasoning models)
# Leave unset for models that don't support temperature (e.g., GPT-5.1 reasoning models).
# Claude 4.7 and later reject it; the request is then retried without it.
# TEMPERATURE=0
# Debug Logging (Optional)
# Log the structure of the messages each chat request sends to the model
# DEBUG_LLM_PAYLOAD=true
# Access Control (Optional)
# ACCESS_CODE_LIST=your-secret-code,another-code
+6 -2
View File
@@ -2,14 +2,18 @@ import { createAmazonBedrock } from "@ai-sdk/amazon-bedrock"
import { createAnthropic } from "@ai-sdk/anthropic"
import { azure, createAzure } from "@ai-sdk/azure"
import { createDeepSeek, deepseek } from "@ai-sdk/deepseek"
import { createGateway, gateway } from "@ai-sdk/gateway"
import { createGoogleGenerativeAI, google } from "@ai-sdk/google"
import { createVertex } from "@ai-sdk/google-vertex"
import { createOpenAI, openai } from "@ai-sdk/openai"
import { aihubmix, createAihubmix } from "@aihubmix/ai-sdk-provider"
import { fromNodeProviderChain } from "@aws-sdk/credential-providers"
import { createOpenRouter } from "@openrouter/ai-sdk-provider"
import { defaultSettingsMiddleware, wrapLanguageModel } from "ai"
import {
createGateway,
defaultSettingsMiddleware,
gateway,
wrapLanguageModel,
} from "ai"
import { createOllama, ollama } from "ollama-ai-provider-v2"
import {
adminProvidersToConfig,
-1
View File
@@ -13,7 +13,6 @@
"@ai-sdk/anthropic": "^3.0.127",
"@ai-sdk/azure": "^3.0.133",
"@ai-sdk/deepseek": "^2.0.71",
"@ai-sdk/gateway": "^3.0.209",
"@ai-sdk/google": "^3.0.130",
"@ai-sdk/google-vertex": "^4.0.210",
"@ai-sdk/openai": "^3.0.124",
-1
View File
@@ -35,7 +35,6 @@
"@ai-sdk/anthropic": "^3.0.127",
"@ai-sdk/azure": "^3.0.133",
"@ai-sdk/deepseek": "^2.0.71",
"@ai-sdk/gateway": "^3.0.209",
"@ai-sdk/google": "^3.0.130",
"@ai-sdk/google-vertex": "^4.0.210",
"@ai-sdk/openai": "^3.0.124",
+74
View File
@@ -0,0 +1,74 @@
// @vitest-environment node
import { simulateReadableStream } from "ai"
import { MockLanguageModelV3 } from "ai/test"
import { afterEach, describe, expect, it, vi } from "vitest"
import { POST as validateDiagram } from "@/app/api/validate-diagram/route"
const RESULT = {
valid: false,
issues: [
{
type: "overlap",
severity: "critical",
description: "Box A covers box B",
},
],
suggestions: ["Move box B to the right"],
}
// A vision model that answers with the JSON in a few text chunks
vi.mock("@/lib/ai-providers", () => ({
getValidationModel: () =>
new MockLanguageModelV3({
doStream: (async () => {
const json = JSON.stringify(RESULT)
return {
stream: simulateReadableStream({
chunks: [
{ type: "text-start", id: "t" },
...[json.slice(0, 20), json.slice(20)].map(
(delta) => ({ type: "text-delta", id: "t", delta }),
),
{ type: "text-end", id: "t" },
{
type: "finish",
finishReason: { unified: "stop", raw: "stop" },
usage: {
inputTokens: { total: 1 },
outputTokens: { total: 1 },
},
},
],
}),
}
}) as any,
}),
}))
const post = () =>
new Request("http://localhost/api/validate-diagram", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ imageData: "data:image/png;base64,AAAA" }),
})
afterEach(() => {
delete process.env.ENABLE_VLM_VALIDATION
})
describe("POST /api/validate-diagram", () => {
it("streams the model's result as JSON text for useObject", async () => {
const res = await validateDiagram(post())
expect(JSON.parse(await res.text())).toEqual(RESULT)
})
it("answers valid when the check is turned off", async () => {
process.env.ENABLE_VLM_VALIDATION = "false"
const res = await validateDiagram(post())
expect(JSON.parse(await res.text())).toEqual({
valid: true,
issues: [],
suggestions: [],
})
})
})