Skip to content

Commit 094904f

Browse files
committed
토큰 스트리밍 실제 구현 + 시스템 헬스 엔드포인트
- aiRouter: provider.stream()으로 실제 토큰 단위 SSE 스트리밍 (non-tool path) - asyncio.Queue + 별도 스레드로 동기 Generator → async SSE 브리지 - tool 라운드 마지막 응답도 실제 스트리밍으로 전환 - /api/system/health: 세션 수, 대화 수, 엔진 상태, 메모리 사용량 반환 - Windows/Linux 양쪽 메모리 측정 (resource/psutil 폴백) - testSystemHealth 테스트 추가 (총 251개)
1 parent b00be16 commit 094904f

3 files changed

Lines changed: 88 additions & 12 deletions

File tree

src/codaro/api/aiRouter.py

Lines changed: 34 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -394,20 +394,44 @@ async def apiChatStream(request: Request):
394394
if sessionId:
395395
executor.setActiveSession(sessionId)
396396

397+
async def _yieldTokenStream(messages):
398+
loop = asyncio.get_running_loop()
399+
queue: asyncio.Queue[str | None] = asyncio.Queue()
400+
401+
def _runStream():
402+
try:
403+
for token in llmProvider.stream(messages):
404+
loop.call_soon_threadsafe(queue.put_nowait, token)
405+
finally:
406+
loop.call_soon_threadsafe(queue.put_nowait, None)
407+
408+
thread = threading.Thread(target=_runStream, daemon=True)
409+
thread.start()
410+
411+
accumulated = ""
412+
while True:
413+
token = await queue.get()
414+
if token is None:
415+
break
416+
accumulated += token
417+
yield f"data: {json.dumps({'type': 'token', 'content': accumulated}, ensure_ascii=False)}\n\n"
418+
419+
convManager.addAssistantMessage(conversationId, accumulated)
420+
yield f"data: {json.dumps({'type': 'done', 'answer': accumulated, 'provider': llmProvider.config.provider, 'model': llmProvider.resolvedModel, 'usage': None, 'toolCalls': []}, ensure_ascii=False)}\n\n"
421+
thread.join(timeout=2)
422+
397423
async def _streamGenerate():
398424
nonlocal msgs, conversationId
399425
yield f"data: {json.dumps({'type': 'start', 'conversationId': conversationId}, ensure_ascii=False)}\n\n"
400426

427+
if not (llmProvider.supportsNativeTools and tools):
428+
async for chunk in _yieldTokenStream(msgs):
429+
yield chunk
430+
return
431+
401432
maxToolRounds = 10
402433
for _round in range(maxToolRounds):
403-
if llmProvider.supportsNativeTools and tools:
404-
response = llmProvider.completeWithTools(msgs, tools)
405-
else:
406-
response = llmProvider.complete(msgs)
407-
convManager.addAssistantMessage(conversationId, response.answer)
408-
yield f"data: {json.dumps({'type': 'token', 'content': response.answer}, ensure_ascii=False)}\n\n"
409-
yield f"data: {json.dumps({'type': 'done', 'answer': response.answer, 'provider': response.provider, 'model': response.model, 'usage': response.usage, 'toolCalls': []}, ensure_ascii=False)}\n\n"
410-
return
434+
response = llmProvider.completeWithTools(msgs, tools)
411435

412436
if not response.toolCalls:
413437
convManager.addAssistantMessage(conversationId, response.answer)
@@ -435,10 +459,8 @@ async def _streamGenerate():
435459

436460
yield f"data: {json.dumps({'type': 'tool_results', 'toolCalls': toolResults}, ensure_ascii=False)}\n\n"
437461

438-
finalResponse = llmProvider.complete(msgs)
439-
convManager.addAssistantMessage(conversationId, finalResponse.answer)
440-
yield f"data: {json.dumps({'type': 'token', 'content': finalResponse.answer}, ensure_ascii=False)}\n\n"
441-
yield f"data: {json.dumps({'type': 'done', 'answer': finalResponse.answer, 'provider': finalResponse.provider, 'model': finalResponse.model, 'usage': finalResponse.usage, 'toolCalls': []}, ensure_ascii=False)}\n\n"
462+
async for chunk in _yieldTokenStream(msgs):
463+
yield chunk
442464

443465
return StreamingResponse(_streamGenerate(), media_type="text/event-stream")
444466

src/codaro/api/systemRouter.py

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,7 @@
11
from __future__ import annotations
22

3+
import os
4+
import sys
35
from typing import Any
46

57
from fastapi import APIRouter
@@ -137,4 +139,40 @@ async def apiUninstallPackage(request: PackageRequest) -> dict[str, Any]:
137139
)
138140
return result.model_dump()
139141

142+
@router.get("/api/system/health")
143+
async def apiSystemHealth() -> dict[str, Any]:
144+
processMemoryMb: float | None = None
145+
try:
146+
import resource as _resource
147+
processMemoryMb = round(_resource.getrusage(_resource.RUSAGE_SELF).ru_maxrss / 1024, 1)
148+
except (ImportError, AttributeError):
149+
pass
150+
if processMemoryMb is None:
151+
try:
152+
import psutil
153+
processMemoryMb = round(psutil.Process(os.getpid()).memory_info().rss / (1024 * 1024), 1)
154+
except (ImportError, AttributeError):
155+
pass
156+
157+
from .aiRouter import _getConversationManager
158+
convManager = _getConversationManager()
159+
160+
return {
161+
"status": "ok",
162+
"python": sys.version,
163+
"pid": os.getpid(),
164+
"processMemoryMb": processMemoryMb,
165+
"sessions": {
166+
"active": state.sessionManager.sessionCount,
167+
},
168+
"conversations": {
169+
"active": convManager.conversationCount,
170+
},
171+
"engine": {
172+
"status": workspaceEngine.status,
173+
"executionCount": workspaceEngine.executionCount,
174+
"variableCount": len(workspaceEngine.getVariables()),
175+
},
176+
}
177+
140178
return router

tests/testServerApi.py

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -226,6 +226,22 @@ def testEnvironmentInfo() -> None:
226226
assert "platform" in info
227227

228228

229+
def testSystemHealth() -> None:
230+
client = TestClient(createServerApp())
231+
232+
response = client.get("/api/system/health")
233+
assert response.status_code == 200
234+
data = response.json()
235+
assert data["status"] == "ok"
236+
assert "python" in data
237+
assert "pid" in data
238+
assert "sessions" in data
239+
assert "conversations" in data
240+
assert "engine" in data
241+
assert data["sessions"]["active"] >= 0
242+
assert data["conversations"]["active"] >= 0
243+
244+
229245
def testStructuredErrorEnvelope() -> None:
230246
client = TestClient(createServerApp())
231247

0 commit comments

Comments
 (0)