|
| 1 | +export function createProcessMessage({ queries, activeExecutions, rateLimitState, execMachine, broadcastSync, runClaudeWithStreaming, cleanupExecution, checkpointManager, discoveredAgents, ownedSessionIds, STARTUP_CWD, buildSystemPrompt, parseRateLimitResetTime, eagerTTS, touchACP, createChunkBatcher, debugLog, logError, scheduleRetry, drainMessageQueue, createEventHandler }) { |
| 2 | + async function processMessageWithStreaming(conversationId, messageId, sessionId, content, agentId, model, subAgent) { |
| 3 | + const startTime = Date.now(); |
| 4 | + touchACP(agentId); |
| 5 | + const conv = queries.getConversation(conversationId); |
| 6 | + if (!conv) { |
| 7 | + console.error(`[stream] Conversation ${conversationId} not found, aborting`); |
| 8 | + queries.updateSession(sessionId, { status: 'error', error: 'Conversation not found' }); |
| 9 | + queries.setIsStreaming(conversationId, false); |
| 10 | + return; |
| 11 | + } |
| 12 | + if (activeExecutions.has(conversationId)) { |
| 13 | + const existing = activeExecutions.get(conversationId); |
| 14 | + if (existing.sessionId !== sessionId) { |
| 15 | + debugLog(`[stream] Conversation ${conversationId} already has active execution (different session), aborting duplicate`); |
| 16 | + return; |
| 17 | + } |
| 18 | + } |
| 19 | + if (rateLimitState.has(conversationId)) { |
| 20 | + const rlState = rateLimitState.get(conversationId); |
| 21 | + if (rlState.retryAt > Date.now()) { |
| 22 | + debugLog(`[stream] Conversation ${conversationId} is in rate limit cooldown, aborting`); |
| 23 | + return; |
| 24 | + } |
| 25 | + } |
| 26 | + activeExecutions.set(conversationId, { pid: null, startTime, sessionId, lastActivity: startTime }); |
| 27 | + execMachine.send(conversationId, { type: 'START', sessionId }); |
| 28 | + queries.setIsStreaming(conversationId, true); |
| 29 | + queries.updateSession(sessionId, { status: 'active' }); |
| 30 | + const batcher = createChunkBatcher(queries, debugLog); |
| 31 | + const cwd = conv?.workingDirectory || STARTUP_CWD; |
| 32 | + const allBlocksRef = { val: [] }; |
| 33 | + const currentSequenceRef = { val: queries.getMaxSequence(sessionId) ?? -1 }; |
| 34 | + const batcherRef = { batcher, eventCount: 0, resumeSessionId: conv?.claudeSessionId || null }; |
| 35 | + const onEvent = createEventHandler({ queries, activeExecutions, broadcastSync, rateLimitState, batcherRef, sessionId, conversationId, messageId, content, agentId, model, subAgent, ownedSessionIds, allBlocksRef, currentSequenceRef, scheduleRetry, eagerTTS, debugLog, parseRateLimitResetTime }); |
| 36 | + try { |
| 37 | + debugLog(`[stream] Starting: conversationId=${conversationId}, sessionId=${sessionId}`); |
| 38 | + let resolvedAgentId = agentId || 'claude-code'; |
| 39 | + const wrapperAgent = discoveredAgents.find(a => a.id === resolvedAgentId && a.protocol === 'cli-wrapper' && a.acpId); |
| 40 | + if (wrapperAgent) resolvedAgentId = wrapperAgent.acpId; |
| 41 | + const resolvedModel = model || conv?.model || null; |
| 42 | + const resolvedSubAgent = subAgent || conv?.subAgent || null; |
| 43 | + const config = { |
| 44 | + verbose: true, outputFormat: 'stream-json', timeout: 1800000, print: true, |
| 45 | + resumeSessionId: batcherRef.resumeSessionId, |
| 46 | + systemPrompt: buildSystemPrompt(agentId, resolvedModel, resolvedSubAgent), |
| 47 | + model: resolvedModel || undefined, subAgent: resolvedSubAgent || undefined, onEvent, |
| 48 | + onPid: (pid) => { const e = activeExecutions.get(conversationId); if (e) e.pid = pid; execMachine.send(conversationId, { type: 'SET_PID', pid }); }, |
| 49 | + onProcess: (proc) => { const e = activeExecutions.get(conversationId); if (e) e.proc = proc; execMachine.send(conversationId, { type: 'SET_PROC', proc }); } |
| 50 | + }; |
| 51 | + const { outputs, sessionId: claudeSessionId } = await runClaudeWithStreaming(content, cwd, resolvedAgentId, config); |
| 52 | + if (rateLimitState.get(conversationId)?.isStreamDetected) { |
| 53 | + debugLog(`[rate-limit] Rate limit already handled in stream for conv ${conversationId}, skipping success handler`); |
| 54 | + return; |
| 55 | + } |
| 56 | + activeExecutions.delete(conversationId); |
| 57 | + execMachine.send(conversationId, { type: 'COMPLETE' }); |
| 58 | + batcher.drain(); |
| 59 | + if (claudeSessionId) ownedSessionIds.delete(claudeSessionId); |
| 60 | + debugLog(`[stream] Claude returned ${outputs.length} outputs, sessionId=${claudeSessionId}`); |
| 61 | + queries.updateSession(sessionId, { status: 'complete', response: JSON.stringify({ outputs, eventCount: batcherRef.eventCount }), completed_at: Date.now() }); |
| 62 | + broadcastSync({ type: 'streaming_complete', sessionId, conversationId, agentId, eventCount: batcherRef.eventCount, seq: currentSequenceRef.val, timestamp: Date.now() }); |
| 63 | + debugLog(`[stream] Completed: ${outputs.length} outputs, ${batcherRef.eventCount} events`); |
| 64 | + } catch (error) { |
| 65 | + const elapsed = Date.now() - startTime; |
| 66 | + debugLog(`[stream] Error after ${elapsed}ms: ${error.message}`); |
| 67 | + const conv2 = queries.getConversation(conversationId); |
| 68 | + if (conv2?.claudeSessionId) ownedSessionIds.delete(conv2.claudeSessionId); |
| 69 | + if (rateLimitState.get(conversationId)?.isStreamDetected) { |
| 70 | + debugLog(`[rate-limit] Rate limit already handled in stream for conv ${conversationId}, skipping catch handler`); |
| 71 | + return; |
| 72 | + } |
| 73 | + const isAuthError = error.authError || error.nonRetryable || /401|unauthorized|invalid.*auth|invalid.*token|auth.*failed|permission denied|access denied/i.test(error.message); |
| 74 | + const isRateLimit = error.rateLimited || /rate.?limit|429|too many requests|overloaded|throttl/i.test(error.message); |
| 75 | + queries.updateSession(sessionId, { status: 'error', error: error.message, completed_at: Date.now() }); |
| 76 | + if (isAuthError) { |
| 77 | + debugLog(`[auth-error] Auth error for conv ${conversationId}: ${error.message}`); |
| 78 | + broadcastSync({ type: 'streaming_error', sessionId, conversationId, error: `Authentication failed: ${error.message}. Please check your API credentials.`, recoverable: false, isAuthError: true, timestamp: Date.now() }); |
| 79 | + const errMsg = queries.createMessage(conversationId, 'assistant', `Error: Authentication failed. ${error.message}. Please update your credentials and try again.`); |
| 80 | + broadcastSync({ type: 'message_created', conversationId, message: errMsg, timestamp: Date.now() }); |
| 81 | + queries.setIsStreaming(conversationId, false); |
| 82 | + batcher.drain(); |
| 83 | + activeExecutions.delete(conversationId); |
| 84 | + return; |
| 85 | + } |
| 86 | + if (isRateLimit) { |
| 87 | + const existingState = rateLimitState.get(conversationId) || {}; |
| 88 | + const retryCount = (existingState.retryCount || 0) + 1; |
| 89 | + const maxRateLimitRetries = 3; |
| 90 | + if (retryCount > maxRateLimitRetries) { |
| 91 | + broadcastSync({ type: 'streaming_error', sessionId, conversationId, error: `Rate limit exceeded after ${retryCount} attempts. Please try again later.`, recoverable: false, timestamp: Date.now() }); |
| 92 | + const errMsg = queries.createMessage(conversationId, 'assistant', `Error: Rate limit exceeded after ${retryCount} attempts. Please try again later.`); |
| 93 | + broadcastSync({ type: 'message_created', conversationId, message: errMsg, timestamp: Date.now() }); |
| 94 | + queries.setIsStreaming(conversationId, false); |
| 95 | + return; |
| 96 | + } |
| 97 | + const cooldownMs = (error.retryAfterSec || 60) * 1000; |
| 98 | + const retryAt = Date.now() + cooldownMs; |
| 99 | + rateLimitState.set(conversationId, { retryAt, cooldownMs, retryCount }); |
| 100 | + broadcastSync({ type: 'rate_limit_hit', sessionId, conversationId, retryAfterMs: cooldownMs, retryAt, retryCount, timestamp: Date.now() }); |
| 101 | + batcher.drain(); |
| 102 | + debugLog(`[rate-limit] Scheduling retry for conv ${conversationId} in ${cooldownMs}ms (attempt ${retryCount + 1})`); |
| 103 | + setTimeout(() => { |
| 104 | + debugLog(`[rate-limit] Timeout fired for conv ${conversationId}, calling scheduleRetry`); |
| 105 | + rateLimitState.delete(conversationId); |
| 106 | + broadcastSync({ type: 'rate_limit_clear', conversationId, timestamp: Date.now() }); |
| 107 | + scheduleRetry(conversationId, messageId, content, agentId, model, subAgent); |
| 108 | + }, cooldownMs); |
| 109 | + return; |
| 110 | + } |
| 111 | + const isSessionConflict = error.exitCode === null && batcherRef.eventCount === 0; |
| 112 | + broadcastSync({ type: 'streaming_error', sessionId, conversationId, error: error.message, isPrematureEnd: error.isPrematureEnd || false, exitCode: error.exitCode, stderrText: error.stderrText, recoverable: elapsed < 60000, isSessionConflict, timestamp: Date.now() }); |
| 113 | + if (!isSessionConflict) { |
| 114 | + const errMsg = queries.createMessage(conversationId, 'assistant', `Error: ${error.message}`); |
| 115 | + broadcastSync({ type: 'message_created', conversationId, message: errMsg, timestamp: Date.now() }); |
| 116 | + } |
| 117 | + } finally { |
| 118 | + batcher.drain(); |
| 119 | + if (!rateLimitState.has(conversationId)) { |
| 120 | + cleanupExecution(conversationId); |
| 121 | + drainMessageQueue(conversationId); |
| 122 | + } |
| 123 | + } |
| 124 | + } |
| 125 | + return { processMessageWithStreaming }; |
| 126 | +} |
0 commit comments