diff --git a/server/src/routes/chat.ts b/server/src/routes/chat.ts index 04e04dcb..0c503acf 100644 --- a/server/src/routes/chat.ts +++ b/server/src/routes/chat.ts @@ -7,6 +7,8 @@ import { Router, type Router as RouterType } from 'express'; import { z } from 'zod'; import { getChatService } from '../services/chat-service.js'; +import { sendGatewayChat, loadGatewayToken } from '../services/gateway-chat-client.js'; +import { broadcastChatMessage } from '../services/broadcast-service.js'; import type { ChatSendInput } from '@veritas-kanban/shared'; import { asyncHandler } from '../middleware/async-handler.js'; import { NotFoundError, ValidationError } from '../middleware/error-handler.js'; @@ -14,6 +16,9 @@ import { createLogger } from '../lib/logger.js'; const log = createLogger('chat'); +// Load gateway token on startup +loadGatewayToken().catch(() => {}); + const router: RouterType = Router(); const chatService = getChatService(); @@ -80,15 +85,68 @@ router.post( log.info({ sessionId, messageId: userMessage.id, taskId: session.taskId }, 'Chat message sent'); - // Return immediately - agent response will stream via WebSocket + // Return immediately - agent response will arrive async res.status(200).json({ sessionId, messageId: userMessage.id, - message: 'Message sent — agent response will stream via WebSocket', + message: 'Message sent — agent response incoming', }); - // TODO: Trigger agent response (will be handled by WebSocket integration) - // The WebSocket server will listen for 'chat:message' events and stream the response + // Trigger async AI response via Clawdbot Gateway + const gatewaySessionKey = `kanban-chat-${sessionId}`; + + sendGatewayChat(input.message, gatewaySessionKey, { + onDelta: (text) => { + // Broadcast streaming chunk to kanban WebSocket clients + broadcastChatMessage(sessionId, { + type: 'chat:delta', + sessionId, + text, + }); + }, + onFinal: async (response) => { + try { + // Save the assistant response to the session + const assistantMessage = await chatService.addMessage(sessionId, { + role: 'assistant', + content: response.text, + agent: session.agent, + }); + + log.info({ sessionId, messageId: assistantMessage.id }, 'Assistant response saved'); + + // Broadcast final message to kanban WebSocket clients + broadcastChatMessage(sessionId, { + type: 'chat:message', + sessionId, + message: assistantMessage, + }); + } catch (err: any) { + log.error({ err: err.message, sessionId }, 'Failed to save assistant response'); + } + }, + onError: async (error) => { + log.error({ error, sessionId }, 'Gateway chat error'); + + // Save error as system message + try { + await chatService.addMessage(sessionId, { + role: 'system', + content: `Error: ${error}`, + }); + + broadcastChatMessage(sessionId, { + type: 'chat:error', + sessionId, + error, + }); + } catch (err: any) { + log.error({ err: err.message }, 'Failed to save error message'); + } + }, + }).catch((err) => { + log.error({ err: err.message, sessionId }, 'Gateway chat failed'); + }); }) ); diff --git a/server/src/services/broadcast-service.ts b/server/src/services/broadcast-service.ts index f642e7e5..114914b6 100644 --- a/server/src/services/broadcast-service.ts +++ b/server/src/services/broadcast-service.ts @@ -11,7 +11,13 @@ export function initBroadcast(wss: WebSocketServer): void { wssRef = wss; } -export type TaskChangeType = 'created' | 'updated' | 'deleted' | 'archived' | 'restored' | 'reordered'; +export type TaskChangeType = + | 'created' + | 'updated' + | 'deleted' + | 'archived' + | 'restored' + | 'reordered'; export interface TaskChangeEvent { type: 'task:changed'; @@ -42,7 +48,31 @@ export function broadcastTaskChange(changeType: TaskChangeType, taskId?: string) const payload = JSON.stringify(message); wssRef.clients.forEach((client: WebSocket) => { - if (client.readyState === 1) { // WebSocket.OPEN = 1 + if (client.readyState === 1) { + // WebSocket.OPEN = 1 + client.send(payload); + } + }); +} + +export interface ChatBroadcastEvent { + type: 'chat:delta' | 'chat:message' | 'chat:error'; + sessionId: string; + text?: string; + message?: unknown; + error?: string; +} + +/** + * Broadcast a chat message/event to all connected WebSocket clients. + */ +export function broadcastChatMessage(sessionId: string, event: ChatBroadcastEvent): void { + if (!wssRef) return; + + const payload = JSON.stringify(event); + + wssRef.clients.forEach((client: WebSocket) => { + if (client.readyState === 1) { client.send(payload); } }); @@ -63,7 +93,8 @@ export function broadcastTelemetryEvent(event: AnyTelemetryEvent): void { const payload = JSON.stringify(message); wssRef.clients.forEach((client: WebSocket) => { - if (client.readyState === 1) { // WebSocket.OPEN = 1 + if (client.readyState === 1) { + // WebSocket.OPEN = 1 client.send(payload); } }); diff --git a/server/src/services/gateway-chat-client.ts b/server/src/services/gateway-chat-client.ts new file mode 100644 index 00000000..b0b99881 --- /dev/null +++ b/server/src/services/gateway-chat-client.ts @@ -0,0 +1,270 @@ +/** + * Gateway Chat Client + * + * Connects to the Clawdbot Gateway WebSocket to proxy chat messages. + * Handles authentication, message sending, and response collection. + */ + +import WebSocket from 'ws'; +import { randomUUID } from 'crypto'; +import { createLogger } from '../lib/logger.js'; + +const log = createLogger('gateway-chat'); + +const GATEWAY_URL = process.env.CLAWDBOT_GATEWAY || 'http://127.0.0.1:18789'; +const PROTOCOL_VERSION = 3; +const CONNECT_TIMEOUT_MS = 10_000; +const RESPONSE_TIMEOUT_MS = 120_000; // 2 minutes for AI response + +// Cached token — populated lazily +let cachedToken: string | null = null; + +function getToken(): string { + if (cachedToken) return cachedToken; + return process.env.CLAWDBOT_GATEWAY_TOKEN || ''; +} + +interface ChatResponse { + text: string; + usage?: Record; + error?: string; +} + +interface StreamCallbacks { + onDelta?: (text: string) => void; + onFinal?: (response: ChatResponse) => void; + onError?: (error: string) => void; +} + +/** + * Send a message to the Clawdbot Gateway and collect the response. + * Opens a temporary WebSocket connection for each request. + */ +export async function sendGatewayChat( + message: string, + sessionKey: string, + callbacks?: StreamCallbacks +): Promise { + const wsUrl = GATEWAY_URL.replace(/^http/, 'ws'); + + return new Promise((resolve, reject) => { + let connected = false; + let responseText = ''; + let responseUsage: Record | undefined; + let connectTimer: ReturnType; + let responseTimer: ReturnType; + + const ws = new WebSocket(wsUrl); + + const cleanup = () => { + clearTimeout(connectTimer); + clearTimeout(responseTimer); + try { + ws.close(); + } catch { + /* ignore */ + } + }; + + connectTimer = setTimeout(() => { + if (!connected) { + cleanup(); + const err = 'Gateway connection timeout'; + callbacks?.onError?.(err); + reject(new Error(err)); + } + }, CONNECT_TIMEOUT_MS); + + ws.on('error', (err) => { + log.error({ err: err.message }, 'Gateway WebSocket error'); + cleanup(); + const errMsg = `Gateway connection failed: ${err.message}`; + callbacks?.onError?.(errMsg); + reject(new Error(errMsg)); + }); + + ws.on('message', (data) => { + let msg: any; + try { + msg = JSON.parse(data.toString()); + } catch { + return; + } + + // Step 1: Handle challenge → send connect + if (msg.type === 'event' && msg.event === 'connect.challenge') { + ws.send( + JSON.stringify({ + type: 'req', + id: randomUUID(), + method: 'connect', + params: { + minProtocol: PROTOCOL_VERSION, + maxProtocol: PROTOCOL_VERSION, + client: { + id: 'gateway-client', + version: '1.0.0', + platform: 'node', + mode: 'backend', + }, + auth: { token: getToken() }, + }, + }) + ); + return; + } + + // Step 2: Handle connect response → send chat.send + if (msg.type === 'res' && msg.ok && msg.payload?.type === 'hello-ok') { + connected = true; + clearTimeout(connectTimer); + + log.info({ sessionKey }, 'Connected to gateway, sending chat message'); + + // Start response timeout + responseTimer = setTimeout(() => { + cleanup(); + const err = 'Gateway response timeout'; + callbacks?.onError?.(err); + reject(new Error(err)); + }, RESPONSE_TIMEOUT_MS); + + ws.send( + JSON.stringify({ + type: 'req', + id: randomUUID(), + method: 'chat.send', + params: { + sessionKey, + message, + idempotencyKey: randomUUID(), + }, + }) + ); + return; + } + + // Handle chat.send ack + if (msg.type === 'res' && msg.ok && msg.payload?.runId) { + log.debug({ runId: msg.payload.runId }, 'Chat run started'); + return; + } + + // Handle errors + if (msg.type === 'res' && !msg.ok) { + cleanup(); + const errMsg = msg.error?.message || 'Unknown gateway error'; + log.error({ error: msg.error }, 'Gateway error'); + callbacks?.onError?.(errMsg); + reject(new Error(errMsg)); + return; + } + + // Step 3: Handle streaming chat events + if (msg.type === 'event' && msg.event === 'chat') { + const payload = msg.payload || {}; + + if (payload.state === 'delta') { + // Gateway sends full accumulated text in each delta, not incremental chunks + const content = payload.message?.content; + if (Array.isArray(content)) { + let fullText = ''; + for (const block of content) { + if (block.type === 'text' && block.text) { + fullText += block.text; + } + } + // Calculate the new chunk (what was added since last delta) + const newChunk = fullText.slice(responseText.length); + responseText = fullText; + if (newChunk) { + callbacks?.onDelta?.(newChunk); + } + } + } + + if (payload.state === 'final') { + // Extract final text if we didn't get it from deltas + if (!responseText && payload.message?.content) { + const content = payload.message.content; + if (Array.isArray(content)) { + for (const block of content) { + if (block.type === 'text' && block.text) { + responseText += block.text; + } + } + } else if (typeof content === 'string') { + responseText = content; + } + } + + responseUsage = payload.usage; + + const response: ChatResponse = { + text: responseText, + usage: responseUsage, + }; + + log.info({ sessionKey, textLength: responseText.length }, 'Chat response complete'); + cleanup(); + callbacks?.onFinal?.(response); + resolve(response); + return; + } + + if (payload.state === 'error') { + cleanup(); + const errMsg = payload.errorMessage || 'Chat error'; + callbacks?.onError?.(errMsg); + reject(new Error(errMsg)); + return; + } + + if (payload.state === 'aborted') { + cleanup(); + const response: ChatResponse = { + text: responseText || '(response aborted)', + }; + callbacks?.onFinal?.(response); + resolve(response); + return; + } + } + }); + + ws.on('close', () => { + if (!connected) { + reject(new Error('Gateway WebSocket closed before connecting')); + } + }); + }); +} + +/** + * Load the gateway token from config file if not in env + */ +export async function loadGatewayToken(): Promise { + if (cachedToken) return cachedToken; + if (process.env.CLAWDBOT_GATEWAY_TOKEN) { + cachedToken = process.env.CLAWDBOT_GATEWAY_TOKEN; + return cachedToken; + } + + try { + const fs = await import('fs/promises'); + const path = await import('path'); + const configPath = path.join(process.env.HOME || '', '.clawdbot', 'clawdbot.json'); + const raw = await fs.readFile(configPath, 'utf-8'); + const config = JSON.parse(raw); + const token = config?.gateway?.auth?.token; + if (token) { + cachedToken = token; + process.env.CLAWDBOT_GATEWAY_TOKEN = token; + return token; + } + } catch (err: any) { + log.warn({ err: err.message }, 'Failed to load gateway token from config'); + } + + return ''; +} diff --git a/web/src/components/chat/ChatPanel.tsx b/web/src/components/chat/ChatPanel.tsx index f284a00e..c924aec3 100644 --- a/web/src/components/chat/ChatPanel.tsx +++ b/web/src/components/chat/ChatPanel.tsx @@ -3,31 +3,13 @@ import { Sheet, SheetContent, SheetHeader, SheetTitle } from '@/components/ui/sh import { Button } from '@/components/ui/button'; import { Input } from '@/components/ui/input'; import { ScrollArea } from '@/components/ui/scroll-area'; -import { - Select, - SelectContent, - SelectItem, - SelectTrigger, - SelectValue, -} from '@/components/ui/select'; -import { Tooltip, TooltipContent, TooltipProvider, TooltipTrigger } from '@/components/ui/tooltip'; -import { - MessageSquare, - Send, - X, - ChevronDown, - ChevronRight, - Loader2, - Bot, - User, -} from 'lucide-react'; +import { MessageSquare, Send, ChevronDown, ChevronRight, Loader2, Bot, User } from 'lucide-react'; import { useChatSession, useSendChatMessage, useChatStream, useChatSessions, } from '@/hooks/useChat'; -import { useConfig } from '@/hooks/useConfig'; import { useTask } from '@/hooks/useTasks'; import type { ChatMessage } from '@veritas-kanban/shared'; @@ -40,12 +22,7 @@ interface ChatPanelProps { export function ChatPanel({ open, onOpenChange, taskId }: ChatPanelProps) { const [message, setMessage] = useState(''); const [mode, setMode] = useState<'ask' | 'build'>('ask'); - const [selectedAgent, setSelectedAgent] = useState('claude-code'); - const [selectedModel, setSelectedModel] = useState('sonnet'); const [currentSessionId, setCurrentSessionId] = useState(); - - const { data: config } = useConfig(); - const enabledAgents = config?.agents?.filter((a: { enabled: boolean }) => a.enabled) || []; const { data: task } = useTask(taskId || ''); const { data: sessions = [] } = useChatSessions(); const { data: session } = useChatSession(currentSessionId); @@ -87,13 +64,11 @@ export function ChatPanel({ open, onOpenChange, taskId }: ChatPanelProps) { sessionId: currentSessionId, taskId, message: message.trim(), - agent: selectedAgent, - model: selectedModel, mode, }, { - onSuccess: (newSession) => { - setCurrentSessionId(newSession.id); + onSuccess: (response) => { + setCurrentSessionId(response.sessionId); setMessage(''); setShouldAutoScroll(true); }, @@ -108,14 +83,14 @@ export function ChatPanel({ open, onOpenChange, taskId }: ChatPanelProps) { } }; - // Load most recent session on mount if exists + // Load session on mount — task-scoped sessions use a deterministic ID useEffect(() => { - if (!currentSessionId && filteredSessions.length > 0) { + if (taskId && !currentSessionId) { + setCurrentSessionId(`task_${taskId}`); + } else if (!taskId && !currentSessionId && filteredSessions.length > 0) { setCurrentSessionId(filteredSessions[0].id); } - }, [filteredSessions, currentSessionId]); - - const models = ['sonnet', 'opus', 'haiku']; + }, [filteredSessions, currentSessionId, taskId]); return ( @@ -124,47 +99,11 @@ export function ChatPanel({ open, onOpenChange, taskId }: ChatPanelProps) { side="right" > {/* Header */} - -
- - - Agent Chat - -
- - - -
-
+ + + + Agent Chat + {taskId && task && (
@@ -241,20 +180,9 @@ export function ChatPanel({ open, onOpenChange, taskId }: ChatPanelProps) { > Build - - - - ⓘ - - -

- Ask: Read-only queries and questions -
- Build: Make changes, create files, execute commands -

-
-
-
+ + {mode === 'ask' ? '· Read-only queries' : '· Changes, files, commands'} +
diff --git a/web/src/components/task/AgentPanel.tsx b/web/src/components/task/AgentPanel.tsx index 21c70144..41efd58e 100644 --- a/web/src/components/task/AgentPanel.tsx +++ b/web/src/components/task/AgentPanel.tsx @@ -76,10 +76,13 @@ export function AgentPanel({ task }: AgentPanelProps) { const sendMessage = useSendMessage(); const [selectedAgent, setSelectedAgent] = useState(); + const [selectedModel, setSelectedModel] = useState(); const [message, setMessage] = useState(''); const [autoScroll, setAutoScroll] = useState(true); const [viewingAttemptId, setViewingAttemptId] = useState(null); + const models = ['sonnet', 'opus', 'haiku']; + const outputRef = useRef(null); // Fetch log for historical attempt @@ -226,6 +229,21 @@ export function AgentPanel({ task }: AgentPanelProps) { ))} + diff --git a/web/src/components/task/TaskDetailPanel.tsx b/web/src/components/task/TaskDetailPanel.tsx index 1339cfd8..bdd39e8c 100644 --- a/web/src/components/task/TaskDetailPanel.tsx +++ b/web/src/components/task/TaskDetailPanel.tsx @@ -112,83 +112,87 @@ export function TaskDetailPanel({
- -
- - Details - {taskSettings.enableAttachments && ( - - - Attachments - - )} - {isCodeTask && ( - <> - - - Git - - - - Agent - - - - Changes - - - - Review - - - )} - - - Metrics - - + {/* Action buttons above tabs */} +
+ + {!readOnly ? ( - {!readOnly && ( - + ) : ( +
+ )} + {!readOnly && isCodeTask && localTask.git?.repo && agentSettings.enablePreview && ( + + )} +
+ + + + Details + {taskSettings.enableAttachments && ( + + + Attachments + )} - {!readOnly && isCodeTask && localTask.git?.repo && agentSettings.enablePreview && ( - + {isCodeTask && ( + <> + + + Git + + + + Agent + + + + Changes + + + + Review + + )} -
+ + + Metrics + +
{/* Details Tab */} diff --git a/web/src/hooks/useChat.ts b/web/src/hooks/useChat.ts index f84d4964..da3e400e 100644 --- a/web/src/hooks/useChat.ts +++ b/web/src/hooks/useChat.ts @@ -1,6 +1,7 @@ import { useQuery, useMutation, useQueryClient } from '@tanstack/react-query'; import { useState, useEffect } from 'react'; import { chatApi } from '@/lib/api/chat'; +import { chatEventTarget } from '@/hooks/useTaskSync'; import type { ChatMessage, ChatSendInput } from '@veritas-kanban/shared'; /** @@ -34,11 +35,12 @@ export function useSendChatMessage() { return useMutation({ mutationFn: (input: ChatSendInput) => chatApi.sendMessage(input), - onSuccess: (session) => { - // Update the session cache - queryClient.setQueryData(['chat', 'sessions', session.id], session); - // Invalidate sessions list + onSuccess: (response) => { + // Invalidate sessions list and the specific session to refetch queryClient.invalidateQueries({ queryKey: ['chat', 'sessions'] }); + if (response.sessionId) { + queryClient.invalidateQueries({ queryKey: ['chat', 'sessions', response.sessionId] }); + } }, }); } @@ -59,22 +61,53 @@ export function useDeleteChatSession() { /** * Listen for streaming chat messages via WebSocket - * This is a placeholder for when WebSocket streaming is implemented */ export function useChatStream(sessionId: string | undefined) { const [streamingMessage, setStreamingMessage] = useState | null>(null); + const [, setStreamingText] = useState(''); const queryClient = useQueryClient(); + // Listen for chat events from the shared WebSocket (via useTaskSync) useEffect(() => { if (!sessionId) return; - // TODO: Connect to WebSocket and listen for chat:message:chunk events - // For now, this is a placeholder. When backend WebSocket streaming is ready: - // 1. Subscribe to chat:message:chunk events for this sessionId - // 2. Accumulate chunks in streamingMessage state - // 3. On chat:message:complete, clear streaming + invalidate session query + const handler = (e: Event) => { + const msg = (e as CustomEvent).detail; + const msgSessionId = msg.sessionId as string; + if (msgSessionId !== sessionId) return; + + if (msg.type === 'chat:delta') { + const text = msg.text as string; + setStreamingText((prev) => { + const newText = prev + text; + setStreamingMessage({ + id: 'streaming', + role: 'assistant', + content: newText, + timestamp: new Date().toISOString(), + }); + return newText; + }); + } + + if (msg.type === 'chat:message') { + setStreamingMessage(null); + setStreamingText(''); + queryClient.invalidateQueries({ queryKey: ['chat', 'sessions', sessionId] }); + } + + if (msg.type === 'chat:error') { + setStreamingMessage(null); + setStreamingText(''); + queryClient.invalidateQueries({ queryKey: ['chat', 'sessions', sessionId] }); + } + }; + + chatEventTarget.addEventListener('chat', handler); return () => { + chatEventTarget.removeEventListener('chat', handler); setStreamingMessage(null); + setStreamingText(''); }; }, [sessionId, queryClient]); diff --git a/web/src/hooks/useTaskSync.ts b/web/src/hooks/useTaskSync.ts index 16e39102..6871b938 100644 --- a/web/src/hooks/useTaskSync.ts +++ b/web/src/hooks/useTaskSync.ts @@ -2,6 +2,10 @@ import { useCallback } from 'react'; import { useQueryClient } from '@tanstack/react-query'; import { useWebSocket, type WebSocketMessage, type ConnectionState } from './useWebSocket'; +// Global event target for chat WebSocket events +// Chat hooks subscribe to this instead of opening their own WebSocket +export const chatEventTarget = new EventTarget(); + /** * Connects to the Veritas Kanban WebSocket server and listens for * task:changed events. When received, invalidates the React Query @@ -20,6 +24,15 @@ export function useTaskSync(): { const handleMessage = useCallback( (message: WebSocketMessage) => { + // Forward chat events to the chat event target + if ( + message.type === 'chat:delta' || + message.type === 'chat:message' || + message.type === 'chat:error' + ) { + chatEventTarget.dispatchEvent(new CustomEvent('chat', { detail: message })); + } + if (message.type === 'task:changed') { // Invalidate task queries to trigger a refetch queryClient.invalidateQueries({ queryKey: ['tasks'] }); diff --git a/web/src/lib/api/chat.ts b/web/src/lib/api/chat.ts index 98de92d3..4151c497 100644 --- a/web/src/lib/api/chat.ts +++ b/web/src/lib/api/chat.ts @@ -24,17 +24,27 @@ export async function getSession(sessionId: string): Promise { return handleResponse(response); } +/** + * Chat send response from the API + * (Not a full ChatSession — agent response streams via WebSocket) + */ +export interface ChatSendResponse { + sessionId: string; + messageId: string; + message: string; +} + /** * Send a chat message */ -export async function sendMessage(input: ChatSendInput): Promise { +export async function sendMessage(input: ChatSendInput): Promise { const response = await fetch(`${API_BASE}/chat/send`, { method: 'POST', headers: { 'Content-Type': 'application/json' }, credentials: 'include', body: JSON.stringify(input), }); - return handleResponse(response); + return handleResponse(response); } /**