Skip to content

Commit 903ee8c

Browse files
committed
fix(chat): keep the chat log's send path free of database reads
1 parent 74df2ab commit 903ee8c

3 files changed

Lines changed: 21 additions & 11 deletions

File tree

‎apps/sim/lib/mothership/chat/chat-log.ts‎

Lines changed: 14 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -10,19 +10,31 @@ const logger = createLogger('ChatLog')
1010

1111
const CHAT_LOG_TIMEOUT_MS = 5_000
1212

13-
/** The turn a Chat send admitted — rebuilt identically when a relay pod recovers it. */
13+
/** The turn a Chat send admitted — rebuilt when a relay pod recovers it. */
1414
export interface ChatTurnLogContext {
1515
chatId: string
1616
messageId: string
1717
requestId: string
1818
userId: string
19+
userEmail?: string
1920
userMessage: string
2021
mode: 'assistant' | 'agent' | 'plan'
2122
startedAt: number
2223
}
2324

2425
export type ChatTurnStatus = 'success' | 'error' | 'aborted'
2526

27+
/** The user's email for a recovered turn, whose send-time session is gone; skipped when the log is off. */
28+
export async function readChatLogEmail(userId: string): Promise<string | undefined> {
29+
if (!env.SIM_LOGGING_WORKFLOW_URL) return undefined
30+
const [row] = await db
31+
.select({ email: user.email })
32+
.from(user)
33+
.where(eq(user.id, userId))
34+
.limit(1)
35+
return row?.email
36+
}
37+
2638
/**
2739
* Operator funnel: posts each finished Chat turn to the Sim workflow at
2840
* `SIM_LOGGING_WORKFLOW_URL` (off when unset). The body keeps the shape the Go
@@ -53,12 +65,6 @@ async function sendChatTurn(
5365
status: ChatTurnStatus,
5466
durationMs: number
5567
): Promise<void> {
56-
const [row] = await db
57-
.select({ email: user.email })
58-
.from(user)
59-
.where(eq(user.id, context.userId))
60-
.limit(1)
61-
6268
const headers: Record<string, string> = { 'Content-Type': 'application/json' }
6369
if (env.SIM_LOGGING_WORKFLOW_API_KEY) headers['X-API-Key'] = env.SIM_LOGGING_WORKFLOW_API_KEY
6470

@@ -69,7 +75,7 @@ async function sendChatTurn(
6975
chatId: context.chatId,
7076
messageId: context.messageId,
7177
userId: context.userId,
72-
userEmail: row?.email,
78+
userEmail: context.userEmail,
7379
userMessage: context.userMessage,
7480
assistantResponse: result.content,
7581
status,

‎apps/sim/lib/mothership/chat/post.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1495,6 +1495,7 @@ export async function handleUnifiedChatPost(req: NextRequest) {
14951495
messageId: userMessageId,
14961496
requestId,
14971497
userId: authenticatedUserId,
1498+
...(authenticatedUserEmail ? { userEmail: authenticatedUserEmail } : {}),
14981499
userMessage: body.message,
14991500
mode: requestMode,
15001501
startedAt: Date.now(),

‎apps/sim/lib/mothership/request/application/recover-stream.ts‎

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,7 @@ import { OrchestrationError } from '@/lib/core/orchestration/types'
1313
import { getLatestRunForStream } from '@/lib/mothership/async-runs/repository'
1414
import { defineAuthorizedChatUseCase } from '@/lib/mothership/chat/application/authorized-chat-use-case'
1515
import { resolveOwnedChatContext } from '@/lib/mothership/chat/application/context'
16-
import type { ChatTurnLogContext } from '@/lib/mothership/chat/chat-log'
16+
import { type ChatTurnLogContext, readChatLogEmail } from '@/lib/mothership/chat/chat-log'
1717
import { buildOnComplete, buildOnError } from '@/lib/mothership/chat/completion'
1818
import { restoreBillingAdmission } from '@/lib/mothership/request/lifecycle/admission'
1919
import { claimRunController } from '@/lib/mothership/request/lifecycle/controller-ownership'
@@ -108,7 +108,8 @@ export const readChatStream = defineAuthorizedChatUseCase({
108108
organizationId,
109109
})
110110
: undefined
111-
const [events, billingAttribution, userPermission] = await Promise.all([
111+
const logsChatTurn = config.data.goRoute === '/api/mothership'
112+
const [events, billingAttribution, userPermission, userEmail] = await Promise.all([
112113
readEvents(run.streamId, '0'),
113114
restoredAdmission
114115
? Promise.resolve(restoredAdmission.attribution)
@@ -118,6 +119,7 @@ export const readChatStream = defineAuthorizedChatUseCase({
118119
workspaceId
119120
? getUserEntityPermissions(userId, 'workspace', workspaceId)
120121
: Promise.resolve(undefined),
122+
logsChatTurn ? readChatLogEmail(userId) : Promise.resolve(undefined),
121123
])
122124
if (workspaceId && !userPermission)
123125
throw new OrchestrationError('forbidden', 'Workspace access revoked')
@@ -134,13 +136,14 @@ export const readChatStream = defineAuthorizedChatUseCase({
134136
requestMode,
135137
notifyWorkspaceStatus: true,
136138
runController: { id: run.id, token: lease.value },
137-
...(config.data.goRoute === '/api/mothership'
139+
...(logsChatTurn
138140
? {
139141
chatLog: {
140142
chatId,
141143
messageId: run.streamId,
142144
requestId,
143145
userId,
146+
...(userEmail ? { userEmail } : {}),
144147
userMessage: intent.message,
145148
mode: requestMode,
146149
startedAt: run.startedAt.getTime(),

0 commit comments

Comments
 (0)