diff --git a/apps/sim/lib/core/config/env.ts b/apps/sim/lib/core/config/env.ts index d213f8fe180..1151274e206 100644 --- a/apps/sim/lib/core/config/env.ts +++ b/apps/sim/lib/core/config/env.ts @@ -312,6 +312,8 @@ export const env = createEnv({ // Monitoring & Analytics TELEMETRY_ENDPOINT: z.string().url().optional(), // Custom telemetry/analytics endpoint + SIM_LOGGING_WORKFLOW_URL: z.string().url().optional(), // Sim workflow execute URL that receives each finished Chat turn (operator chat log; off when unset) + SIM_LOGGING_WORKFLOW_API_KEY: z.string().min(1).optional(), // X-API-Key for SIM_LOGGING_WORKFLOW_URL COST_MULTIPLIER: z.number().optional(), // Multiplier for cost calculations LOG_LEVEL: z.enum(['DEBUG', 'INFO', 'WARN', 'ERROR']).optional(), // Minimum log level to display (defaults to ERROR in production, DEBUG in development) GRAFANA_OTLP_ENDPOINT: z.string().url().optional(), // Grafana Cloud OTLP HTTP gateway base URL (e.g., https://otlp-gateway-prod-us-east-0.grafana.net/otlp). Trigger.dev exporters append /v1/traces, /v1/logs, /v1/metrics. diff --git a/apps/sim/lib/mothership/chat/chat-log.ts b/apps/sim/lib/mothership/chat/chat-log.ts new file mode 100644 index 00000000000..508d3341abc --- /dev/null +++ b/apps/sim/lib/mothership/chat/chat-log.ts @@ -0,0 +1,119 @@ +import { db } from '@sim/db' +import { user } from '@sim/db/schema' +import { createLogger } from '@sim/logger' +import { getErrorMessage } from '@sim/utils/errors' +import { eq } from 'drizzle-orm' +import { env } from '@/lib/core/config/env' +import type { OrchestratorResult } from '@/lib/mothership/request/types' + +const logger = createLogger('ChatLog') + +const CHAT_LOG_TIMEOUT_MS = 5_000 + +/** The turn a Chat send admitted — rebuilt when a relay pod recovers it. */ +export interface ChatTurnLogContext { + chatId: string + messageId: string + requestId: string + userId: string + userEmail?: string + userMessage: string + mode: 'assistant' | 'agent' | 'plan' + startedAt: number +} + +export type ChatTurnStatus = 'success' | 'error' | 'aborted' + +/** + * The user's email for a recovered turn, whose send-time session is gone. Skipped when the + * log is off; best-effort, so it never fails the recovery it enriches. + */ +export async function readChatLogEmail(userId: string): Promise { + if (!env.SIM_LOGGING_WORKFLOW_URL) return undefined + try { + const [row] = await db + .select({ email: user.email }) + .from(user) + .where(eq(user.id, userId)) + .limit(1) + return row?.email + } catch (error) { + logger.warn('Chat log email lookup failed', { userId, error: getErrorMessage(error) }) + return undefined + } +} + +/** + * Operator funnel: posts each finished Chat turn to the Sim workflow at + * `SIM_LOGGING_WORKFLOW_URL` (off when unset). The body keeps the shape the Go + * copilot sent — `{ input: { event: 'copilot_request_completed', ... } }` — so the + * existing workflow reads it unchanged. Fire-and-forget: never awaited, never throws. + */ +export function logChatTurn( + context: ChatTurnLogContext, + result: OrchestratorResult, + status: ChatTurnStatus +): void { + const url = env.SIM_LOGGING_WORKFLOW_URL + if (!url) return + const durationMs = Date.now() - context.startedAt + void sendChatTurn(url, context, result, status, durationMs).catch((error) => { + logger.warn('Chat log workflow request failed', { + chatId: context.chatId, + requestId: context.requestId, + error: getErrorMessage(error), + }) + }) +} + +async function sendChatTurn( + url: string, + context: ChatTurnLogContext, + result: OrchestratorResult, + status: ChatTurnStatus, + durationMs: number +): Promise { + const headers: Record = { 'Content-Type': 'application/json' } + if (env.SIM_LOGGING_WORKFLOW_API_KEY) headers['X-API-Key'] = env.SIM_LOGGING_WORKFLOW_API_KEY + + const input = { + event: 'copilot_request_completed', + idempotencyKey: context.messageId, + requestId: context.requestId, + chatId: context.chatId, + messageId: context.messageId, + userId: context.userId, + userEmail: context.userEmail, + userMessage: context.userMessage, + assistantResponse: result.content, + status, + errored: status === 'error', + aborted: status === 'aborted', + errorMessage: result.error ?? result.errors?.join('\n'), + mode: context.mode, + source: 'workspace-chat', + startedAt: new Date(context.startedAt).toISOString(), + durationMs, + usage: { + inputTokens: result.usage?.prompt ?? 0, + outputTokens: result.usage?.completion ?? 0, + }, + } + + // A redirect would forward X-API-Key to another host and turn the POST into a GET. + const response = await fetch(url, { + method: 'POST', + redirect: 'error', + headers, + body: JSON.stringify({ input }), + signal: AbortSignal.timeout(CHAT_LOG_TIMEOUT_MS), + }) + await response.body?.cancel() + if (!response.ok) { + logger.warn('Chat log workflow returned a non-2xx status', { + status: response.status, + chatId: context.chatId, + requestId: context.requestId, + }) + } +} diff --git a/apps/sim/lib/mothership/chat/completion.ts b/apps/sim/lib/mothership/chat/completion.ts index a7a0c051bfd..16ea7df062d 100644 --- a/apps/sim/lib/mothership/chat/completion.ts +++ b/apps/sim/lib/mothership/chat/completion.ts @@ -1,5 +1,6 @@ import { createLogger } from '@sim/logger' import { getErrorMessage } from '@sim/utils/errors' +import { type ChatTurnLogContext, logChatTurn } from '@/lib/mothership/chat/chat-log' import { buildPersistedAssistantMessage, withStoppedContentBlock, @@ -23,6 +24,8 @@ export function buildOnComplete(params: { organizationId?: string userId?: string requestMode?: 'assistant' | 'agent' | 'plan' + /** Present for workspace and organization Chat turns, which feed the operator chat log. */ + chatLog?: ChatTurnLogContext /** * Root agent span for this request. When present, the final * assistant message + invoked tool calls are recorded as @@ -46,6 +49,7 @@ export function buildOnComplete(params: { userId, runController, otelRoot, + chatLog, } = params const notifyChatStatus = params.notifyChatStatus ?? params.notifyWorkspaceStatus ?? false @@ -78,6 +82,7 @@ export function buildOnComplete(params: { finalization.updated || finalization.outcome === CopilotChatFinalizeOutcome.AssistantAlreadyPersisted + if (chatLog && shouldPublishCompletion) logChatTurn(chatLog, result, 'aborted') if (notifyChatStatus && shouldPublishCompletion) { publishChatStatusChanged( { workspaceId, organizationId, userId }, @@ -97,7 +102,7 @@ export function buildOnComplete(params: { const assistantMessage = buildPersistedAssistantMessage(result, requestId, params.requestMode) const hasPartial = !!assistantMessage.content?.trim() || (assistantMessage.contentBlocks?.length ?? 0) > 0 - await finalizeAssistantTurn({ + const finalization = await finalizeAssistantTurn({ runController, chatId, userMessageId, @@ -107,6 +112,9 @@ export function buildOnComplete(params: { ...(result.success ? {} : { streamMarkerPolicy: 'active-or-cleared' as const }), }) + if (chatLog && finalization.updated) { + logChatTurn(chatLog, result, result.success ? 'success' : 'error') + } if (notifyChatStatus) { publishChatStatusChanged( { workspaceId, organizationId, userId }, @@ -138,9 +146,19 @@ export function buildOnError(params: { organizationId?: string userId?: string requestMode?: 'assistant' | 'agent' | 'plan' + /** Present for workspace and organization Chat turns, which feed the operator chat log. */ + chatLog?: ChatTurnLogContext }) { - const { chatId, userMessageId, requestId, workspaceId, organizationId, userId, runController } = - params + const { + chatId, + userMessageId, + requestId, + workspaceId, + organizationId, + userId, + runController, + chatLog, + } = params const notifyChatStatus = params.notifyChatStatus ?? params.notifyWorkspaceStatus ?? false return async (error: Error, result?: OrchestratorResult) => { @@ -151,19 +169,16 @@ export function buildOnError(params: { // cancelled / non-success completion path, so the partial assistant turn // (text + tool calls + subagent work) survives the refetch instead of the // chat collapsing to an empty assistant row. - const assistantMessage = buildPersistedAssistantMessage( - { - content: '', - contentBlocks: [], - toolCalls: [], - ...result, - success: false, - error: result?.error || getErrorMessage(error), - }, - requestId, - params.requestMode - ) - await finalizeAssistantTurn({ + const failed: OrchestratorResult = { + content: '', + contentBlocks: [], + toolCalls: [], + ...result, + success: false, + error: result?.error || getErrorMessage(error), + } + const assistantMessage = buildPersistedAssistantMessage(failed, requestId, params.requestMode) + const finalization = await finalizeAssistantTurn({ runController, chatId, userMessageId, @@ -171,6 +186,7 @@ export function buildOnError(params: { streamMarkerPolicy: 'active-or-cleared', }) + if (chatLog && finalization.updated) logChatTurn(chatLog, failed, 'error') if (notifyChatStatus) { publishChatStatusChanged( { workspaceId, organizationId, userId }, diff --git a/apps/sim/lib/mothership/chat/post.ts b/apps/sim/lib/mothership/chat/post.ts index 9822e89d962..2d0d3c73046 100644 --- a/apps/sim/lib/mothership/chat/post.ts +++ b/apps/sim/lib/mothership/chat/post.ts @@ -36,6 +36,7 @@ import { type AssistantImageContent, prepareOrganizationChatAttachments, } from '@/lib/mothership/chat/assistant-images' +import type { ChatTurnLogContext } from '@/lib/mothership/chat/chat-log' import { buildOnComplete, buildOnError } from '@/lib/mothership/chat/completion' import { DESKTOP_TERMINAL_HINT_ID_MAX_LENGTH, @@ -1485,6 +1486,21 @@ export async function handleUnifiedChatPost(req: NextRequest) { // Admission committed. A failure to attach this HTTP sink must leave the turn recoverable. sendClaim = undefined } + const requestMode = + body.mode === 'plan' ? 'plan' : body.mode === 'assistant' ? 'assistant' : 'agent' + const chatLog: ChatTurnLogContext | undefined = + branch.kind !== 'workflow' && actualChatId + ? { + chatId: actualChatId, + messageId: userMessageId, + requestId, + userId: authenticatedUserId, + ...(authenticatedUserEmail ? { userEmail: authenticatedUserEmail } : {}), + userMessage: body.message, + mode: requestMode, + startedAt: Date.now(), + } + : undefined const stream = createSSEStream({ requestPayload, admittedRun, @@ -1526,9 +1542,9 @@ export async function handleUnifiedChatPost(req: NextRequest) { notifyChatStatus: branch.notifyChatStatus, organizationId: branch.kind === 'organization' ? branch.organizationId : undefined, userId: authenticatedUserId, - requestMode: - body.mode === 'plan' ? 'plan' : body.mode === 'assistant' ? 'assistant' : 'agent', + requestMode, otelRoot, + chatLog, }), onError: buildOnError({ runController, @@ -1539,8 +1555,8 @@ export async function handleUnifiedChatPost(req: NextRequest) { notifyChatStatus: branch.notifyChatStatus, organizationId: branch.kind === 'organization' ? branch.organizationId : undefined, userId: authenticatedUserId, - requestMode: - body.mode === 'plan' ? 'plan' : body.mode === 'assistant' ? 'assistant' : 'agent', + requestMode, + chatLog, }), }, }) diff --git a/apps/sim/lib/mothership/request/application/recover-stream.test.ts b/apps/sim/lib/mothership/request/application/recover-stream.test.ts index 2a1eedf766b..afc38fd9e51 100644 --- a/apps/sim/lib/mothership/request/application/recover-stream.test.ts +++ b/apps/sim/lib/mothership/request/application/recover-stream.test.ts @@ -73,6 +73,7 @@ const run = { workspaceId: '33333333-3333-4333-8333-333333333333', status: 'active', workflowId: null, + startedAt: new Date('2026-01-01T00:00:00Z'), requestContext: { requestId: 'request', controllerToken: 'old-controller', diff --git a/apps/sim/lib/mothership/request/application/recover-stream.ts b/apps/sim/lib/mothership/request/application/recover-stream.ts index 441750260e8..d049331d755 100644 --- a/apps/sim/lib/mothership/request/application/recover-stream.ts +++ b/apps/sim/lib/mothership/request/application/recover-stream.ts @@ -13,6 +13,7 @@ import { OrchestrationError } from '@/lib/core/orchestration/types' import { getLatestRunForStream } from '@/lib/mothership/async-runs/repository' import { defineAuthorizedChatUseCase } from '@/lib/mothership/chat/application/authorized-chat-use-case' import { resolveOwnedChatContext } from '@/lib/mothership/chat/application/context' +import { type ChatTurnLogContext, readChatLogEmail } from '@/lib/mothership/chat/chat-log' import { buildOnComplete, buildOnError } from '@/lib/mothership/chat/completion' import { restoreBillingAdmission } from '@/lib/mothership/request/lifecycle/admission' import { claimRunController } from '@/lib/mothership/request/lifecycle/controller-ownership' @@ -107,7 +108,8 @@ export const readChatStream = defineAuthorizedChatUseCase({ organizationId, }) : undefined - const [events, billingAttribution, userPermission] = await Promise.all([ + const logsChatTurn = config.data.goRoute === '/api/mothership' + const [events, billingAttribution, userPermission, userEmail] = await Promise.all([ readEvents(run.streamId, '0'), restoredAdmission ? Promise.resolve(restoredAdmission.attribution) @@ -117,10 +119,13 @@ export const readChatStream = defineAuthorizedChatUseCase({ workspaceId ? getUserEntityPermissions(userId, 'workspace', workspaceId) : Promise.resolve(undefined), + logsChatTurn ? readChatLogEmail(userId) : Promise.resolve(undefined), ]) if (workspaceId && !userPermission) throw new OrchestrationError('forbidden', 'Workspace access revoked') const requestId = typeof saved?.requestId === 'string' ? saved.requestId : generateId() + const requestMode: ChatTurnLogContext['mode'] = + intent.mode === 'plan' ? 'plan' : intent.mode === 'assistant' ? 'assistant' : 'agent' const completion = { chatId, userMessageId: run.streamId, @@ -128,14 +133,23 @@ export const readChatStream = defineAuthorizedChatUseCase({ workspaceId, organizationId, userId, - requestMode: - intent.mode === 'plan' - ? ('plan' as const) - : intent.mode === 'assistant' - ? ('assistant' as const) - : ('agent' as const), + requestMode, notifyWorkspaceStatus: true, runController: { id: run.id, token: lease.value }, + ...(logsChatTurn + ? { + chatLog: { + chatId, + messageId: run.streamId, + requestId, + userId, + ...(userEmail ? { userEmail } : {}), + userMessage: intent.message, + mode: requestMode, + startedAt: run.startedAt.getTime(), + }, + } + : {}), } const stream = createSSEStream({ userId,