Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions apps/sim/lib/core/config/env.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
119 changes: 119 additions & 0 deletions apps/sim/lib/mothership/chat/chat-log.ts
Original file line number Diff line number Diff line change
@@ -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<string | undefined> {
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(
Comment thread
waleedlatif1 marked this conversation as resolved.
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<void> {
const headers: Record<string, string> = { '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, {
Comment thread
waleedlatif1 marked this conversation as resolved.
Comment thread
waleedlatif1 marked this conversation as resolved.
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,
})
}
}
48 changes: 32 additions & 16 deletions apps/sim/lib/mothership/chat/completion.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand All @@ -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
Expand All @@ -46,6 +49,7 @@ export function buildOnComplete(params: {
userId,
runController,
otelRoot,
chatLog,
} = params
const notifyChatStatus = params.notifyChatStatus ?? params.notifyWorkspaceStatus ?? false

Expand Down Expand Up @@ -78,6 +82,7 @@ export function buildOnComplete(params: {
finalization.updated ||
finalization.outcome === CopilotChatFinalizeOutcome.AssistantAlreadyPersisted

if (chatLog && shouldPublishCompletion) logChatTurn(chatLog, result, 'aborted')
Comment thread
waleedlatif1 marked this conversation as resolved.
if (notifyChatStatus && shouldPublishCompletion) {
publishChatStatusChanged(
{ workspaceId, organizationId, userId },
Expand All @@ -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,
Expand All @@ -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 },
Expand Down Expand Up @@ -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) => {
Expand All @@ -151,26 +169,24 @@ 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,
assistantMessage,
streamMarkerPolicy: 'active-or-cleared',
})

if (chatLog && finalization.updated) logChatTurn(chatLog, failed, 'error')
if (notifyChatStatus) {
publishChatStatusChanged(
{ workspaceId, organizationId, userId },
Expand Down
24 changes: 20 additions & 4 deletions apps/sim/lib/mothership/chat/post.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
}),
Comment thread
waleedlatif1 marked this conversation as resolved.
onError: buildOnError({
runController,
Expand All @@ -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,
}),
},
})
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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',
Expand Down
28 changes: 21 additions & 7 deletions apps/sim/lib/mothership/request/application/recover-stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -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)
Expand All @@ -117,25 +119,37 @@ export const readChatStream = defineAuthorizedChatUseCase({
workspaceId
? getUserEntityPermissions(userId, 'workspace', workspaceId)
: Promise.resolve(undefined),
logsChatTurn ? readChatLogEmail(userId) : Promise.resolve(undefined),
Comment thread
waleedlatif1 marked this conversation as resolved.
])
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,
requestId,
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,
Expand Down
Loading