From a1c00a23b0ffe36c4290117dca4c6197f52e9871 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Sat, 26 Sep 2026 09:49:48 -0700 Subject: [PATCH 1/5] fix(chat): restore the operator chat log after the TypeScript worker cutover --- apps/sim/lib/core/config/env.ts | 2 + apps/sim/lib/mothership/chat/chat-log.ts | 88 ++++++++++++++++++++++ apps/sim/lib/mothership/chat/completion.ts | 42 +++++++---- apps/sim/lib/mothership/chat/post.ts | 26 ++++++- 4 files changed, 140 insertions(+), 18 deletions(-) create mode 100644 apps/sim/lib/mothership/chat/chat-log.ts 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..693247bae7e --- /dev/null +++ b/apps/sim/lib/mothership/chat/chat-log.ts @@ -0,0 +1,88 @@ +import { createLogger } from '@sim/logger' +import { getErrorMessage } from '@sim/utils/errors' +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 finished turn's identity, captured when the user's message is admitted. */ +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' + +/** + * 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 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: Date.now() - context.startedAt, + usage: { + inputTokens: result.usage?.prompt ?? 0, + outputTokens: result.usage?.completion ?? 0, + }, + } + + void fetch(url, { + method: 'POST', + headers, + body: JSON.stringify({ input }), + signal: AbortSignal.timeout(CHAT_LOG_TIMEOUT_MS), + }) + .then(async (response) => { + 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, + }) + } + }) + .catch((error) => { + logger.warn('Chat log workflow request failed', { + chatId: context.chatId, + requestId: context.requestId, + error: getErrorMessage(error), + }) + }) +} diff --git a/apps/sim/lib/mothership/chat/completion.ts b/apps/sim/lib/mothership/chat/completion.ts index a7a0c051bfd..e90c2d30794 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 Chat turns that 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 }, @@ -107,6 +112,7 @@ export function buildOnComplete(params: { ...(result.success ? {} : { streamMarkerPolicy: 'active-or-cleared' as const }), }) + if (chatLog) logChatTurn(chatLog, result, result.success ? 'success' : 'error') if (notifyChatStatus) { publishChatStatusChanged( { workspaceId, organizationId, userId }, @@ -138,9 +144,19 @@ export function buildOnError(params: { organizationId?: string userId?: string requestMode?: 'assistant' | 'agent' | 'plan' + /** Present for Chat turns that 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,18 +167,15 @@ 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 - ) + const failed: OrchestratorResult = { + content: '', + contentBlocks: [], + toolCalls: [], + ...result, + success: false, + error: result?.error || getErrorMessage(error), + } + const assistantMessage = buildPersistedAssistantMessage(failed, requestId, params.requestMode) await finalizeAssistantTurn({ runController, chatId, @@ -171,6 +184,7 @@ export function buildOnError(params: { streamMarkerPolicy: 'active-or-cleared', }) + if (chatLog) 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..4b96a07e23d 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, @@ -985,6 +986,7 @@ export async function handleUnifiedChatPost(req: NextRequest) { let requestId = '' const executionId = generateId() const runId = generateId() + const startedAt = Date.now() try { const session = await getSession() @@ -1485,6 +1487,22 @@ 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' + /** Workspace and organization Chat feed the operator chat log; the workflow panel does not. */ + const chatLog: ChatTurnLogContext | undefined = + branch.kind !== 'workflow' && actualChatId + ? { + chatId: actualChatId, + messageId: userMessageId, + requestId, + userId: authenticatedUserId, + ...(authenticatedUserEmail ? { userEmail: authenticatedUserEmail } : {}), + userMessage: body.message, + mode: requestMode, + startedAt, + } + : undefined const stream = createSSEStream({ requestPayload, admittedRun, @@ -1526,9 +1544,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 +1557,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, }), }, }) From 74df2ab2b0e5811f95b882429276e3116049d9ff Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Sat, 26 Sep 2026 10:01:29 -0700 Subject: [PATCH 2/5] fix(chat): log only turns this call finalized, include recovered turns, and resolve the email at send time --- apps/sim/lib/mothership/chat/chat-log.ts | 57 ++++++++++++------- apps/sim/lib/mothership/chat/completion.ts | 14 +++-- apps/sim/lib/mothership/chat/post.ts | 5 +- .../application/recover-stream.test.ts | 1 + .../request/application/recover-stream.ts | 23 ++++++-- 5 files changed, 63 insertions(+), 37 deletions(-) diff --git a/apps/sim/lib/mothership/chat/chat-log.ts b/apps/sim/lib/mothership/chat/chat-log.ts index 693247bae7e..2da1b609519 100644 --- a/apps/sim/lib/mothership/chat/chat-log.ts +++ b/apps/sim/lib/mothership/chat/chat-log.ts @@ -1,5 +1,8 @@ +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' @@ -7,13 +10,12 @@ const logger = createLogger('ChatLog') const CHAT_LOG_TIMEOUT_MS = 5_000 -/** The finished turn's identity, captured when the user's message is admitted. */ +/** The turn a Chat send admitted — rebuilt identically 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 @@ -34,6 +36,28 @@ export function logChatTurn( ): 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 [row] = await db + .select({ email: user.email }) + .from(user) + .where(eq(user.id, context.userId)) + .limit(1) const headers: Record = { 'Content-Type': 'application/json' } if (env.SIM_LOGGING_WORKFLOW_API_KEY) headers['X-API-Key'] = env.SIM_LOGGING_WORKFLOW_API_KEY @@ -45,7 +69,7 @@ export function logChatTurn( chatId: context.chatId, messageId: context.messageId, userId: context.userId, - userEmail: context.userEmail, + userEmail: row?.email, userMessage: context.userMessage, assistantResponse: result.content, status, @@ -55,34 +79,25 @@ export function logChatTurn( mode: context.mode, source: 'workspace-chat', startedAt: new Date(context.startedAt).toISOString(), - durationMs: Date.now() - context.startedAt, + durationMs, usage: { inputTokens: result.usage?.prompt ?? 0, outputTokens: result.usage?.completion ?? 0, }, } - void fetch(url, { + const response = await fetch(url, { method: 'POST', headers, body: JSON.stringify({ input }), signal: AbortSignal.timeout(CHAT_LOG_TIMEOUT_MS), }) - .then(async (response) => { - 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, - }) - } - }) - .catch((error) => { - logger.warn('Chat log workflow request failed', { - chatId: context.chatId, - requestId: context.requestId, - error: getErrorMessage(error), - }) + 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 e90c2d30794..16ea7df062d 100644 --- a/apps/sim/lib/mothership/chat/completion.ts +++ b/apps/sim/lib/mothership/chat/completion.ts @@ -24,7 +24,7 @@ export function buildOnComplete(params: { organizationId?: string userId?: string requestMode?: 'assistant' | 'agent' | 'plan' - /** Present for Chat turns that feed the operator chat log. */ + /** 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 @@ -102,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, @@ -112,7 +112,9 @@ export function buildOnComplete(params: { ...(result.success ? {} : { streamMarkerPolicy: 'active-or-cleared' as const }), }) - if (chatLog) logChatTurn(chatLog, result, result.success ? 'success' : 'error') + if (chatLog && finalization.updated) { + logChatTurn(chatLog, result, result.success ? 'success' : 'error') + } if (notifyChatStatus) { publishChatStatusChanged( { workspaceId, organizationId, userId }, @@ -144,7 +146,7 @@ export function buildOnError(params: { organizationId?: string userId?: string requestMode?: 'assistant' | 'agent' | 'plan' - /** Present for Chat turns that feed the operator chat log. */ + /** Present for workspace and organization Chat turns, which feed the operator chat log. */ chatLog?: ChatTurnLogContext }) { const { @@ -176,7 +178,7 @@ export function buildOnError(params: { error: result?.error || getErrorMessage(error), } const assistantMessage = buildPersistedAssistantMessage(failed, requestId, params.requestMode) - await finalizeAssistantTurn({ + const finalization = await finalizeAssistantTurn({ runController, chatId, userMessageId, @@ -184,7 +186,7 @@ export function buildOnError(params: { streamMarkerPolicy: 'active-or-cleared', }) - if (chatLog) logChatTurn(chatLog, failed, 'error') + 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 4b96a07e23d..60940bed99e 100644 --- a/apps/sim/lib/mothership/chat/post.ts +++ b/apps/sim/lib/mothership/chat/post.ts @@ -986,7 +986,6 @@ export async function handleUnifiedChatPost(req: NextRequest) { let requestId = '' const executionId = generateId() const runId = generateId() - const startedAt = Date.now() try { const session = await getSession() @@ -1489,7 +1488,6 @@ export async function handleUnifiedChatPost(req: NextRequest) { } const requestMode = body.mode === 'plan' ? 'plan' : body.mode === 'assistant' ? 'assistant' : 'agent' - /** Workspace and organization Chat feed the operator chat log; the workflow panel does not. */ const chatLog: ChatTurnLogContext | undefined = branch.kind !== 'workflow' && actualChatId ? { @@ -1497,10 +1495,9 @@ export async function handleUnifiedChatPost(req: NextRequest) { messageId: userMessageId, requestId, userId: authenticatedUserId, - ...(authenticatedUserEmail ? { userEmail: authenticatedUserEmail } : {}), userMessage: body.message, mode: requestMode, - startedAt, + startedAt: Date.now(), } : undefined const stream = createSSEStream({ 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..1e99e286b8a 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 } 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' @@ -121,6 +122,8 @@ export const readChatStream = defineAuthorizedChatUseCase({ 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 +131,22 @@ 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 }, + ...(config.data.goRoute === '/api/mothership' + ? { + chatLog: { + chatId, + messageId: run.streamId, + requestId, + userId, + userMessage: intent.message, + mode: requestMode, + startedAt: run.startedAt.getTime(), + }, + } + : {}), } const stream = createSSEStream({ userId, From 903ee8c34122acfe8e13f36b31284e50a878c36a Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Sat, 26 Sep 2026 10:12:49 -0700 Subject: [PATCH 3/5] fix(chat): keep the chat log's send path free of database reads --- apps/sim/lib/mothership/chat/chat-log.ts | 22 ++++++++++++------- apps/sim/lib/mothership/chat/post.ts | 1 + .../request/application/recover-stream.ts | 9 +++++--- 3 files changed, 21 insertions(+), 11 deletions(-) diff --git a/apps/sim/lib/mothership/chat/chat-log.ts b/apps/sim/lib/mothership/chat/chat-log.ts index 2da1b609519..c6f6d34fd52 100644 --- a/apps/sim/lib/mothership/chat/chat-log.ts +++ b/apps/sim/lib/mothership/chat/chat-log.ts @@ -10,12 +10,13 @@ const logger = createLogger('ChatLog') const CHAT_LOG_TIMEOUT_MS = 5_000 -/** The turn a Chat send admitted — rebuilt identically when a relay pod recovers it. */ +/** 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 @@ -23,6 +24,17 @@ export interface ChatTurnLogContext { 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. */ +export async function readChatLogEmail(userId: string): Promise { + if (!env.SIM_LOGGING_WORKFLOW_URL) return undefined + const [row] = await db + .select({ email: user.email }) + .from(user) + .where(eq(user.id, userId)) + .limit(1) + return row?.email +} + /** * 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 @@ -53,12 +65,6 @@ async function sendChatTurn( status: ChatTurnStatus, durationMs: number ): Promise { - const [row] = await db - .select({ email: user.email }) - .from(user) - .where(eq(user.id, context.userId)) - .limit(1) - const headers: Record = { 'Content-Type': 'application/json' } if (env.SIM_LOGGING_WORKFLOW_API_KEY) headers['X-API-Key'] = env.SIM_LOGGING_WORKFLOW_API_KEY @@ -69,7 +75,7 @@ async function sendChatTurn( chatId: context.chatId, messageId: context.messageId, userId: context.userId, - userEmail: row?.email, + userEmail: context.userEmail, userMessage: context.userMessage, assistantResponse: result.content, status, diff --git a/apps/sim/lib/mothership/chat/post.ts b/apps/sim/lib/mothership/chat/post.ts index 60940bed99e..2d0d3c73046 100644 --- a/apps/sim/lib/mothership/chat/post.ts +++ b/apps/sim/lib/mothership/chat/post.ts @@ -1495,6 +1495,7 @@ export async function handleUnifiedChatPost(req: NextRequest) { messageId: userMessageId, requestId, userId: authenticatedUserId, + ...(authenticatedUserEmail ? { userEmail: authenticatedUserEmail } : {}), userMessage: body.message, mode: requestMode, startedAt: Date.now(), diff --git a/apps/sim/lib/mothership/request/application/recover-stream.ts b/apps/sim/lib/mothership/request/application/recover-stream.ts index 1e99e286b8a..d049331d755 100644 --- a/apps/sim/lib/mothership/request/application/recover-stream.ts +++ b/apps/sim/lib/mothership/request/application/recover-stream.ts @@ -13,7 +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 } from '@/lib/mothership/chat/chat-log' +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' @@ -108,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) @@ -118,6 +119,7 @@ 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') @@ -134,13 +136,14 @@ export const readChatStream = defineAuthorizedChatUseCase({ requestMode, notifyWorkspaceStatus: true, runController: { id: run.id, token: lease.value }, - ...(config.data.goRoute === '/api/mothership' + ...(logsChatTurn ? { chatLog: { chatId, messageId: run.streamId, requestId, userId, + ...(userEmail ? { userEmail } : {}), userMessage: intent.message, mode: requestMode, startedAt: run.startedAt.getTime(), From 29f5c64898ce67c1d6d16d510f3c3feac4b69bc9 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Sat, 26 Sep 2026 10:23:08 -0700 Subject: [PATCH 4/5] fix(chat): make the recovered turn's email lookup best-effort --- apps/sim/lib/mothership/chat/chat-log.ts | 22 +++++++++++++++------- 1 file changed, 15 insertions(+), 7 deletions(-) diff --git a/apps/sim/lib/mothership/chat/chat-log.ts b/apps/sim/lib/mothership/chat/chat-log.ts index c6f6d34fd52..365fb8c6c49 100644 --- a/apps/sim/lib/mothership/chat/chat-log.ts +++ b/apps/sim/lib/mothership/chat/chat-log.ts @@ -24,15 +24,23 @@ export interface ChatTurnLogContext { 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. */ +/** + * 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 - const [row] = await db - .select({ email: user.email }) - .from(user) - .where(eq(user.id, userId)) - .limit(1) - return row?.email + 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 + } } /** From 1eed2167b5fb89ff96146471e415d4db3bd31396 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Sat, 26 Sep 2026 10:32:08 -0700 Subject: [PATCH 5/5] fix(chat): refuse redirects on the authenticated chat log POST --- apps/sim/lib/mothership/chat/chat-log.ts | 2 ++ 1 file changed, 2 insertions(+) diff --git a/apps/sim/lib/mothership/chat/chat-log.ts b/apps/sim/lib/mothership/chat/chat-log.ts index 365fb8c6c49..508d3341abc 100644 --- a/apps/sim/lib/mothership/chat/chat-log.ts +++ b/apps/sim/lib/mothership/chat/chat-log.ts @@ -100,8 +100,10 @@ async function sendChatTurn( }, } + // 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),