import { randomUUID } from 'node:crypto'; import { WORKS_SQUARE_CONFIG } from '../api/works-config'; import { proxyAwareFetch } from '../utils/proxy-fetch'; import { getValidWorksSquareAccessToken } from './works-square-session'; const COURSE_ID_PATTERN = /^[A-Za-z0-9_-]{1,64}$/; const MODULE_ID_PATTERN = /^(?!\.{1,2}$)[A-Za-z0-9._-]{1,128}$/; const SHA256_PATTERN = /^[0-9a-f]{64}$/; const REMOTE_ID_PATTERN = /^[A-Za-z0-9_.:-]{1,128}$/; const MAX_MESSAGE_LENGTH = 4_000; const MAX_HISTORY_ITEMS = 8; const MAX_HISTORY_CONTENT_LENGTH = 4_000; const MAX_SCENE_TEXT_LENGTH = 256; const MAX_ANSWER_LENGTH = 64_000; const MAX_SSE_BUFFER_LENGTH = 256_000; export type LearningAgentRequest = { courseId: string; contentHash: string; message: string; history?: Array<{ role: 'user' | 'assistant'; content: string }>; anchor?: { sceneId?: string; sceneOrder?: number; sceneTitle?: string; actionIndex?: number; moduleId?: string | null; moduleContentHash?: string; }; }; type Dependencies = { fetchImpl?: typeof fetch; getAccessToken?: typeof getValidWorksSquareAccessToken; apiBaseUrl?: string; }; class LearningAgentHttpError extends Error { constructor(readonly status: number, message: string) { super(message); this.name = 'LearningAgentHttpError'; } } function record(value: unknown): Record { return value && typeof value === 'object' && !Array.isArray(value) ? value as Record : {}; } function isLearningAccountChanged(error: unknown): boolean { return record(error).code === 'LEARNING_ACCOUNT_CHANGED'; } export function createLearningAgentClient(dependencies: Dependencies = {}) { const fetchImpl = dependencies.fetchImpl ?? proxyAwareFetch; const getAccessToken = dependencies.getAccessToken ?? getValidWorksSquareAccessToken; const apiBaseUrl = (dependencies.apiBaseUrl ?? WORKS_SQUARE_CONFIG.apiBaseUrl).replace(/\/+$/, ''); const sessions = new Map(); const sessionFlights = new Map>(); const turnTails = new Map>(); async function authorizedJson( path: string, init: RequestInit, assertCurrentAccount: () => void, ): Promise> { assertCurrentAccount(); try { const token = await getAccessToken({ fetchImpl }); assertCurrentAccount(); if (!token) throw new Error('请先登录'); const request = (accessToken: string) => { assertCurrentAccount(); return fetchImpl(`${apiBaseUrl}${path}`, { ...init, headers: { Accept: 'application/json', 'Content-Type': 'application/json', Authorization: `Bearer ${accessToken}` }, }); }; let response = await request(token); assertCurrentAccount(); if (response.status === 401) { await response.body?.cancel().catch(() => undefined); assertCurrentAccount(); const refreshed = await getAccessToken({ fetchImpl, forceRefresh: true }); assertCurrentAccount(); if (refreshed) { response = await request(refreshed); assertCurrentAccount(); } } const payload = await response.json().catch(() => null); if (!response.ok) throw new LearningAgentHttpError(response.status, '助教服务暂时不可用'); return record(payload); } catch (error) { if (isLearningAccountChanged(error)) throw error; if (error instanceof LearningAgentHttpError || (error instanceof Error && error.message === '请先登录')) throw error; throw new Error('助教服务暂时不可用', { cause: error }); } } async function createSession( bindingKey: string, courseId: string, assertCurrentAccount: () => void, ): Promise { assertCurrentAccount(); const existing = sessions.get(bindingKey); if (existing) return existing; const inFlight = sessionFlights.get(bindingKey); if (inFlight) return inFlight; const pending = authorizedJson('/api/agents/sessions', { method: 'POST', body: JSON.stringify({ client_session_id: `learning-${courseId}-${randomUUID()}`, runtime: 'learning', runtime_version: 'v1', binding: { kind: 'learning-course', key: bindingKey }, }), }, assertCurrentAccount).then((session) => { assertCurrentAccount(); const sessionId = typeof session.session_id === 'string' && REMOTE_ID_PATTERN.test(session.session_id) ? session.session_id : ''; if (!sessionId) throw new Error('助教服务暂时不可用'); sessions.set(bindingKey, sessionId); return sessionId; }).finally(() => sessionFlights.delete(bindingKey)); sessionFlights.set(bindingKey, pending); return pending; } function safeTurnInput(input: LearningAgentRequest) { if (!input || typeof input.message !== 'string') throw new Error('助教请求无效'); const message = input.message.trim(); const history = Array.isArray(input.history) ? input.history.slice(-MAX_HISTORY_ITEMS).map((item) => { if ((item?.role !== 'user' && item?.role !== 'assistant') || typeof item.content !== 'string' || !item.content.trim() || item.content.length > MAX_HISTORY_CONTENT_LENGTH) throw new Error('助教请求无效'); return { role: item.role, content: item.content }; }) : []; const anchor = input.anchor ?? {}; if ((anchor.sceneId !== undefined && (typeof anchor.sceneId !== 'string' || !anchor.sceneId || anchor.sceneId.length > MAX_SCENE_TEXT_LENGTH)) || (anchor.sceneTitle !== undefined && (typeof anchor.sceneTitle !== 'string' || anchor.sceneTitle.length > MAX_SCENE_TEXT_LENGTH)) || (anchor.sceneOrder !== undefined && (!Number.isSafeInteger(anchor.sceneOrder) || anchor.sceneOrder < 0 || anchor.sceneOrder > 100_000)) || (anchor.actionIndex !== undefined && (!Number.isSafeInteger(anchor.actionIndex) || anchor.actionIndex < 0 || anchor.actionIndex > 100_000)) || (anchor.moduleId !== undefined && anchor.moduleId !== null && !MODULE_ID_PATTERN.test(anchor.moduleId)) || (anchor.moduleContentHash !== undefined && !SHA256_PATTERN.test(anchor.moduleContentHash))) { throw new Error('助教请求无效'); } return { message, history, anchor: { ...(anchor.sceneId === undefined ? {} : { sceneId: anchor.sceneId }), ...(anchor.sceneOrder === undefined ? {} : { sceneOrder: anchor.sceneOrder }), ...(anchor.sceneTitle === undefined ? {} : { sceneTitle: anchor.sceneTitle }), ...(anchor.actionIndex === undefined ? {} : { actionIndex: anchor.actionIndex }), ...(anchor.moduleId === undefined ? {} : { moduleId: anchor.moduleId }), ...(anchor.moduleContentHash === undefined ? {} : { moduleContentHash: anchor.moduleContentHash }), }, }; } async function runTurn( sessionId: string, input: ReturnType, assertCurrentAccount: () => void, ): Promise<{ text: string }> { assertCurrentAccount(); const command = await authorizedJson(`/api/agents/sessions/${encodeURIComponent(sessionId)}/commands`, { method: 'POST', body: JSON.stringify({ client_command_id: `turn-${randomUUID()}`, name: 'turn.submit', input: { message: input.message, // Compatibility context for an older runtime. The durable Works // Agent session remains authoritative across turns. history: input.history, anchor: input.anchor, }, }), }, assertCurrentAccount); const runId = typeof command.run_id === 'string' && REMOTE_ID_PATTERN.test(command.run_id) ? command.run_id : ''; const ticket = await authorizedJson(`/api/agents/sessions/${encodeURIComponent(sessionId)}/stream-tickets`, { method: 'POST', body: JSON.stringify({ transport: 'sse' }), }, assertCurrentAccount); const streamUrl = typeof ticket.stream_url === 'string' && ticket.stream_url.length <= 2_048 ? ticket.stream_url : ''; const expectedStreamPath = `/api/agents/sessions/${encodeURIComponent(sessionId)}/events`; let streamTarget: URL | null = null; try { const candidate = new URL(streamUrl, `${apiBaseUrl}/`); const base = new URL(apiBaseUrl); if (candidate.origin === base.origin && candidate.pathname === expectedStreamPath) streamTarget = candidate; } catch { streamTarget = null; } if (!runId || !streamTarget) throw new Error('助教服务暂时不可用'); const controller = new AbortController(); const timeout = setTimeout(() => controller.abort(), 120_000); let response: Response; try { assertCurrentAccount(); response = await fetchImpl(streamTarget, { headers: { Accept: 'text/event-stream' }, signal: controller.signal }); assertCurrentAccount(); } catch (error) { clearTimeout(timeout); if (isLearningAccountChanged(error)) throw error; throw new Error('助教服务暂时不可用', { cause: error }); } if (!response.ok || !response.body) { clearTimeout(timeout); throw new LearningAgentHttpError(response.status, '助教事件通道不可用'); } const reader = response.body.getReader(); const decoder = new TextDecoder(); let buffer = ''; let text = ''; try { while (true) { const chunk = await reader.read(); assertCurrentAccount(); if (chunk.done) break; buffer += decoder.decode(chunk.value, { stream: true }); if (buffer.length > MAX_SSE_BUFFER_LENGTH) throw new Error('助教服务暂时不可用'); const frames = buffer.split('\n\n'); buffer = frames.pop() || ''; for (const frame of frames) { const data = frame.split('\n').find((line) => line.startsWith('data:'))?.slice(5).trim(); if (!data) continue; let envelope: Record; try { envelope = record(JSON.parse(data)); } catch { throw new Error('助教服务暂时不可用'); } if (envelope.run_id !== runId) continue; const payload = record(envelope.payload); if (envelope.type === 'learning.assistant.delta' && typeof payload.delta === 'string') { if (text.length + payload.delta.length > MAX_ANSWER_LENGTH) throw new Error('助教服务暂时不可用'); text += payload.delta; } if (envelope.type === 'learning.assistant.failed') throw new Error('助教回答失败'); if (envelope.type === 'learning.assistant.completed') { controller.abort(); const answer = text.trim(); if (!answer) throw new Error('助教回答失败'); return { text: answer }; } } } } finally { clearTimeout(timeout); await reader.cancel().catch(() => undefined); } throw new Error('助教回答意外中断'); } async function ask( input: LearningAgentRequest, assertCurrentAccount: () => void = () => undefined, ): Promise<{ text: string }> { assertCurrentAccount(); const safeInput = safeTurnInput(input); if (!COURSE_ID_PATTERN.test(input.courseId) || !SHA256_PATTERN.test(input.contentHash) || !safeInput.message || safeInput.message.length > MAX_MESSAGE_LENGTH) throw new Error('助教请求无效'); const bindingKey = `${input.courseId}:${input.contentHash}`; const previous = turnTails.get(bindingKey) ?? Promise.resolve(); const turn = previous.catch(() => undefined).then(async () => { assertCurrentAccount(); let sessionId = await createSession(bindingKey, input.courseId, assertCurrentAccount); try { return await runTurn(sessionId, safeInput, assertCurrentAccount); } catch (error) { if (!(error instanceof LearningAgentHttpError) || error.status !== 404) throw error; assertCurrentAccount(); sessions.delete(bindingKey); sessionId = await createSession(bindingKey, input.courseId, assertCurrentAccount); assertCurrentAccount(); return runTurn(sessionId, safeInput, assertCurrentAccount); } }); const tail = turn.then(() => undefined, () => undefined); turnTails.set(bindingKey, tail); try { return await turn; } finally { if (turnTails.get(bindingKey) === tail) turnTails.delete(bindingKey); } } async function reset( courseId?: string, contentHash?: string, assertCurrentAccount: () => void = () => undefined, ): Promise { assertCurrentAccount(); const exactKey = courseId && contentHash ? `${courseId}:${contentHash}` : null; const entries = [...sessions.entries()].filter(([key]) => exactKey ? key === exactKey : courseId ? key.startsWith(`${courseId}:`) : true); for (const [key, sessionId] of entries) { assertCurrentAccount(); sessions.delete(key); void authorizedJson( `/api/agents/sessions/${encodeURIComponent(sessionId)}`, { method: 'DELETE' }, assertCurrentAccount, ).catch(() => undefined); } } return { ask, reset }; }