import * as vscode from 'vscode'; import { LiveReasoningFilter } from '../../lib/contextBuilders/liveReasoningFilter'; import { logError, summarizeText } from '../../utils'; import { lmStudioSamplingFromConfig, lmStudioRespondExtrasFromConfig } from '../../lib/contextBuilders/lmStudioSampling'; import type { AgentExecutorOptions, ChatMessage } from '../../agent'; export interface StreamChatOnceDeps { options: AgentExecutorOptions; getWebview: () => vscode.Webview | undefined; isStaleRun: (runId: number) => boolean; createStreamingRequest: (params: { baseUrl: string; modelName: string; reqMessages: ChatMessage[]; temperature: number; /** Dynamic output-token cap computed from the remaining context budget. */ maxTokens?: number; /** Model context window in tokens (used for Ollama's num_ctx). */ contextLength?: number; }) => Promise<{ response: Response; engine: 'lmstudio' | 'ollama'; apiUrl: string }>; } export async function streamChatOnce(deps: StreamChatOnceDeps, params: { runId: number; useLmStudioSdk: boolean; engine: 'lmstudio' | 'ollama'; ollamaUrl: string; modelName: string; messages: ChatMessage[]; temperature: number; maxTokens: number; contextLength: number; contextOverflowPolicy: 'stopAtLimit' | 'truncateMiddle' | 'rollingWindow'; signal: AbortSignal; postLiveDeltas: boolean; }): Promise<{ text: string; stopReason?: string; aborted: boolean }> { let accumulated = ''; let finishStopReason: string | undefined; // 라이브 표시 시 추론 구간(/Harmony thought)은 토큰 단위로 차단 — agent.ts 본 스트림과 동일. const liveFilter = new LiveReasoningFilter(); const post = (token: string) => { if (params.postLiveDeltas && token) { const visible = liveFilter.push(token); if (visible) deps.getWebview()?.postMessage({ type: 'streamChunk', value: visible }); } }; if (params.useLmStudioSdk) { try { const stream = deps.options.lmStudioStreamer!.stream({ modelName: params.modelName, messages: params.messages.map((m) => ({ role: m.role, content: m.content })), temperature: params.temperature, maxTokens: params.maxTokens, contextOverflowPolicy: params.contextOverflowPolicy, ...lmStudioSamplingFromConfig(), ...lmStudioRespondExtrasFromConfig(), signal: params.signal, }); for await (const { token, stopReason } of stream) { if (deps.isStaleRun(params.runId)) { return { text: accumulated, stopReason: finishStopReason, aborted: true }; } if (token) { accumulated += token; post(token); } if (stopReason) finishStopReason = stopReason; } } catch (err: any) { if (err?.name === 'AbortError' || params.signal.aborted) { return { text: accumulated, stopReason: finishStopReason, aborted: true }; } const msg = err?.message ?? String(err); if (/context\s*length|contextlengthreached|exceed|too\s*long/i.test(msg)) { finishStopReason = 'contextLengthReached'; } logError('streamChatOnce SDK path failed.', { engine: params.engine, error: msg }); throw err; } return { text: accumulated, stopReason: finishStopReason, aborted: false }; } const request = await deps.createStreamingRequest({ baseUrl: params.ollamaUrl, modelName: params.modelName, reqMessages: params.messages, temperature: params.temperature, maxTokens: params.maxTokens, contextLength: params.contextLength, }); const reader = request.response.body?.getReader(); if (!reader) throw new Error('Response body is not readable.'); const decoder = new TextDecoder(); let buffer = ''; const consumeJsonLine = (line: string) => { const trimmed = line.trim(); if (!trimmed || trimmed === 'data: [DONE]') return; try { const raw = trimmed.startsWith('data: ') ? trimmed.slice(6) : trimmed; const json = JSON.parse(raw); const token = params.engine === 'lmstudio' ? json.choices?.[0]?.delta?.content || '' : json.message?.content || json.response || ''; if (token) { accumulated += token; post(token); } const fr = params.engine === 'lmstudio' ? json.choices?.[0]?.finish_reason : (json.done_reason ?? (json.done === true ? 'stop' : undefined)); if (fr) finishStopReason = fr; } catch (e: any) { logError('streamChatOnce: failed to parse chunk.', { engine: params.engine, chunk: summarizeText(trimmed, 200), error: e?.message ?? String(e) }); } }; try { while (true) { const { done, value } = await reader.read(); if (done) break; if (deps.isStaleRun(params.runId)) { return { text: accumulated, stopReason: finishStopReason, aborted: true }; } buffer += decoder.decode(value, { stream: true }); const lines = buffer.split('\n'); buffer = lines.pop() || ''; for (const line of lines) consumeJsonLine(line); } if (buffer.trim()) consumeJsonLine(buffer); } catch (err: any) { if (err?.name === 'AbortError') { return { text: accumulated, stopReason: finishStopReason, aborted: true }; } logError('streamChatOnce REST path failed.', { engine: params.engine, error: err?.message ?? String(err) }); throw err; } finally { try { reader.releaseLock(); } catch { /* already released on abort */ } } return { text: accumulated, stopReason: finishStopReason, aborted: false }; }