feat(teacher): stream Yuxi reasoning and answer previews

This commit is contained in:
brother7 committed 2026-09-29 18:18:56 +08:00
1 parent f443f3bc66
commit 8ccc4650a2
12 files changed
+364 -42

No files matched your search

+57 -35
View File
@@ -39,6 +39,13 @@ export interface TeacherCloudTransport {
): Promise<void>;
}
export interface TeacherCloudStreamUpdate {
/** Current assistant message snapshot, replaced across messages and Runs. */
response?: string;
/** New reasoning text, accumulated across the question's tool continuations. */
reasoning?: string;
}
export function teacherCloudTransport(account: TeacherAccount): TeacherCloudTransport {
let session: TeacherSession | undefined;
const fetchCloud = async (path: string, body?: unknown, signal?: AbortSignal) => {
@@ -181,6 +188,7 @@ export function prepareCloudTeacher(
saveRequestId: (id: string) => Promise<void>,
transport: TeacherCloudTransport = teacherCloudTransport(account),
onToolActivity: (activity: TeacherToolActivity) => void = () => undefined,
onStream: (update: TeacherCloudStreamUpdate) => void = () => undefined,
) {
const tools = createTeacherReadTools(access);
return {
@@ -281,6 +289,44 @@ export function prepareCloudTeacher(
readBytes = 0;
let runText = '',
messageId = '';
let reasoningMessage = '';
const readEvents = async (threadId: string) => {
await transport.events(
'/runs/' + encodeURIComponent(runId) + '/events?after_seq=' + encodeURIComponent(cursor),
bounded,
(_event, envelope, id) => {
if (id) cursor = id;
// Parent Runs include child events; only the main thread is visible.
if (envelope.thread_id !== threadId) return;
const payload = envelope.payload ? object(envelope.payload) : {};
for (const item of Array.isArray(payload.items) ? payload.items : payload.chunk ? [payload.chunk] : []) {
const chunk = object(item);
const activity = cloudToolActivity(chunk, requestId);
if (activity) reportActivity(activity);
const event = chunk.stream_event ? object(chunk.stream_event) : {};
if (event.type !== 'message_delta') continue;
if (typeof event.message_id === 'string' && event.message_id !== messageId) {
if (!structuredReply && messageId && runText) onText('\n\n');
messageId = event.message_id;
runText = '';
if (structuredReply) onStream({ response: '' });
}
const reasoning = typeof event.reasoning_content === 'string' ? event.reasoning_content
: typeof event.additional_reasoning_content === 'string' ? event.additional_reasoning_content : '';
if (reasoning) {
const identity = runId + ':' + messageId;
onStream({ reasoning: (reasoningMessage && reasoningMessage !== identity ? '\n\n' : '') + reasoning });
reasoningMessage = identity;
}
if (typeof event.content === 'string') {
runText += event.content;
if (structuredReply) onStream({ response: runText });
else onText(event.content);
}
}
}
);
};
while (Date.now() < deadline) {
bounded.throwIfAborted();
access.assertCurrent();
@@ -289,11 +335,18 @@ export function prepareCloudTeacher(
undefined,
bounded
);
// A fast run/continuation may already be terminal before our first GET.
// Drain its retained events before moving on so reasoning is not lost.
if (view.status !== 'running' && view.status !== 'pending' && typeof view.thread_id === 'string') {
try { await readEvents(view.thread_id); }
catch { bounded.throwIfAborted(); access.assertCurrent(); }
}
if (view.continued_run_id) {
runId = identifier(view.continued_run_id);
cursor = '0-0';
runText = '';
messageId = '';
if (structuredReply) onStream({ response: '' });
continue;
}
if (view.status === 'interrupted') {
@@ -346,6 +399,7 @@ export function prepareCloudTeacher(
cursor = '0-0';
runText = '';
messageId = '';
if (structuredReply) onStream({ response: '' });
onProgress('智能体正在继续思考…');
continue;
}
@@ -353,8 +407,8 @@ export function prepareCloudTeacher(
if (typeof view.output !== 'string')
throw new TeacherError(502, 'teacher_protocol_invalid', '智能体返回的正文格式不受支持,请重试。');
const output = view.output;
// Structured UI replies must contain only the final answer. A cloud
// run can stream a preamble or draft before reading and continuing.
// Commit only the final answer; onStream carries replaceable previews
// from messages that may precede a tool read or another draft.
if (structuredReply) onText(output);
else if (output.startsWith(runText)) onText(output.slice(runText.length));
else if (output) onText('\n\n' + output);
@@ -374,39 +428,7 @@ export function prepareCloudTeacher(
const threadId = identifier(view.thread_id);
onProgress('智能体正在思考…');
try {
await transport.events(
'/runs/' +
encodeURIComponent(runId) +
'/events?after_seq=' +
encodeURIComponent(cursor),
bounded,
(_event, envelope, id) => {
if (id) cursor = id;
// Yuxi 的父 Run 也包含子线程事件,智能体正文只接收云端主线程文本。
if (envelope.thread_id !== threadId) return;
const payload = envelope.payload ? object(envelope.payload) : {};
for (const item of Array.isArray(payload.items)
? payload.items
: payload.chunk
? [payload.chunk]
: []) {
const chunk = object(item);
// A checkpoint continuation changes Run, not the logical tool call.
const activity = cloudToolActivity(chunk, requestId);
if (activity) reportActivity(activity);
const event = chunk.stream_event ? object(chunk.stream_event) : {};
if (event.type === 'message_delta' && typeof event.content === 'string') {
if (typeof event.message_id === 'string' && event.message_id !== messageId) {
if (!structuredReply && messageId && runText) onText('\n\n');
messageId = event.message_id;
runText = '';
}
runText += event.content;
if (!structuredReply) onText(event.content);
}
}
}
);
await readEvents(threadId);
} catch (error) {
bounded.throwIfAborted();
access.assertCurrent();