Skip to content

Commit 7c023cf

Browse files
committed
fix(knowledge): record a terminal status for every Slack Assistant run
The Slack Assistant admits its own run row and only ever marked it as an error, so every completed or stopped Slack turn stayed active. It now records the terminal status once, after the turn ends, through the shared run-status update: complete on success, cancelled when its user stopped it (in Slack or in Sim), and error otherwise. Every other run-creating path already settles its run: interactive turns through their controller's finalize, and headless turns that admit their own run in the lifecycle's own finally.
1 parent 34768f4 commit 7c023cf

2 files changed

Lines changed: 255 additions & 1 deletion

File tree

Lines changed: 240 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,240 @@
1+
/**
2+
* The Slack Assistant's run record against real PostgreSQL: the run row it admits and
3+
* the terminal status it records are production code. Slack delivery, the worker
4+
* lifecycle, identity, and chat locking are stood in, since only the record's outcome
5+
* is under test.
6+
*/
7+
import { billingAttributionMock } from '@sim/testing/mocks/billing-attribution.mock'
8+
import {
9+
mothershipChatPayloadMock,
10+
mothershipChatPayloadMockFns,
11+
} from '@sim/testing/mocks/mothership-chat-payload.mock'
12+
import { mothershipEnvironmentContextMock } from '@sim/testing/mocks/mothership-environment-context.mock'
13+
import { organizationAuthorizationMock } from '@sim/testing/mocks/organization-authorization.mock'
14+
import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'
15+
16+
const hoisted = vi.hoisted(() => ({
17+
chat: { id: '', model: 'default' },
18+
turnId: '',
19+
userId: '',
20+
organizationId: '',
21+
lifecycle: vi.fn(),
22+
/** The turn's own controller, which a Stop aborts through its registered stream. */
23+
controller: new AbortController(),
24+
stopped: vi.fn(async () => false),
25+
}))
26+
vi.mock('@/lib/knowledge/application/slack-search/authorization', () => ({
27+
authorizeSlackSearchInstallation: async () => ({
28+
installation: { id: 'i1', organizationId: hoisted.organizationId, teamId: 'T1' },
29+
secret: { botToken: 'token' },
30+
}),
31+
}))
32+
vi.mock('@/lib/internal/slack/search-client', () => ({
33+
getSlackSearchSender: async () => ({ email: 'member@example.com' }),
34+
}))
35+
vi.mock('@/lib/knowledge/application/slack-search/identity', () => ({
36+
resolveSlackSearchMember: async () => hoisted.userId,
37+
SlackSearchIdentityError: class extends Error {},
38+
}))
39+
vi.mock('@/lib/knowledge/application/slack-search/chat', () => ({
40+
resolveSlackSearchChat: async () => hoisted.chat,
41+
persistSlackSearchQuestion: async () => undefined,
42+
slackSearchChatOperation: { id: 'organization.chats.slack' },
43+
}))
44+
vi.mock('@/lib/core/application/organization-authorization', () => organizationAuthorizationMock)
45+
vi.mock('@/lib/knowledge/application/operations', () => ({
46+
knowledgeOperations: { search: { organizationOperation: { id: 'knowledge.search' } } },
47+
}))
48+
vi.mock('@/lib/knowledge/application/slack-search/repository', () => ({
49+
recordSlackSearchOutcome: async () => undefined,
50+
}))
51+
vi.mock('@/lib/knowledge/application/slack-search/turns', () => ({
52+
requireSlackSearchTurnLease: async () => undefined,
53+
wasSlackSearchTurnStopped: hoisted.stopped,
54+
}))
55+
vi.mock('@/lib/knowledge/application/slack-search/onboarding', () => ({
56+
sendSlackSearchOnboarding: vi.fn(),
57+
}))
58+
vi.mock('@/lib/knowledge/application/slack-search/title', () => ({
59+
generateSlackSearchChatTitle: async () => undefined,
60+
}))
61+
vi.mock('@/lib/billing/core/billing-attribution', () => billingAttributionMock)
62+
vi.mock('@/lib/mothership/application/load-search-integrations', () => ({
63+
loadCopilotSearchIntegrations: async () => '{"connections":[],"available":[]}',
64+
}))
65+
vi.mock('@/lib/mothership/chat/payload', () => mothershipChatPayloadMock)
66+
vi.mock('@/lib/mothership/chat/terminal-state', () => ({
67+
finalizeAssistantTurn: async () => ({ appendedAssistant: true }),
68+
}))
69+
vi.mock('@/lib/mothership/environment-context', () => mothershipEnvironmentContextMock)
70+
vi.mock('@/lib/mothership/request/lifecycle/headless', () => ({
71+
runHeadlessCopilotLifecycle: hoisted.lifecycle,
72+
}))
73+
vi.mock('@/lib/mothership/request/session/abort', () => ({
74+
acquirePendingChatStream: async () => true,
75+
cleanupAbortMarker: async () => undefined,
76+
getChatStreamLockOwners: async () => ({
77+
status: 'verified',
78+
ownersByChatId: new Map([[hoisted.chat.id, hoisted.turnId]]),
79+
}),
80+
registerActiveStream: vi.fn(),
81+
releasePendingChatStream: async () => undefined,
82+
startAbortPoller: () => 0,
83+
unregisterActiveStream: vi.fn(),
84+
}))
85+
vi.mock('@/executor/utils/resolved-secret-content-projection', () => ({
86+
projectResolvedSecretDiagnosticContent: (value: unknown) => ({ safe: true, value }),
87+
}))
88+
vi.mock('@/lib/slack-search/connections', () => ({ deliverSlackSearchConnections: vi.fn() }))
89+
vi.mock('@/lib/slack-search/assistant-stream', () => ({
90+
SlackSearchAssistantStream: class {
91+
start = async () => undefined
92+
finish = async () => undefined
93+
finishWithError = async () => undefined
94+
terminateAfterFailure = async () => undefined
95+
onEvent = async () => undefined
96+
assertHealthy = () => undefined
97+
},
98+
}))
99+
100+
import { db } from '@sim/db'
101+
import { copilotChats, copilotRuns, organization, user } from '@sim/db/schema'
102+
import { generateId } from '@sim/utils/id'
103+
import { eq } from 'drizzle-orm'
104+
import { runSlackSearchAssistant } from '@/lib/knowledge/application/slack-search/assistant'
105+
import { AbortReason } from '@/lib/mothership/request/session/abort-reason'
106+
107+
const principal = {
108+
kind: 'slack_installation',
109+
credentialId: 'c1',
110+
credentialVersion: 'v1',
111+
appId: 'A1',
112+
teamId: 'T1',
113+
eventId: 'Ev1',
114+
receivedAt: new Date(),
115+
} as const
116+
117+
function job() {
118+
return {
119+
installationId: 'i1',
120+
revision: 'r1',
121+
credentialId: 'c1',
122+
credentialVersion: 'v1',
123+
receivedAt: Date.now(),
124+
message: {
125+
appId: 'A1',
126+
teamId: 'T1',
127+
eventId: 'Ev1',
128+
channelId: 'D1',
129+
userId: 'U1',
130+
messageTs: '1800000000.000001',
131+
query: 'release notes',
132+
queryTooLong: false,
133+
},
134+
}
135+
}
136+
137+
/** Runs one Slack turn in a fresh private chat and returns the run it recorded. */
138+
async function slackTurn() {
139+
const chatId = generateId()
140+
const turnId = generateId()
141+
hoisted.chat = { id: chatId, model: 'default' }
142+
hoisted.turnId = turnId
143+
await db.insert(copilotChats).values({
144+
id: chatId,
145+
userId: hoisted.userId,
146+
organizationId: hoisted.organizationId,
147+
type: 'mothership',
148+
})
149+
const controller = new AbortController()
150+
hoisted.controller = controller
151+
const outcome = await runSlackSearchAssistant(principal, {
152+
job: job(),
153+
turnId,
154+
leaseId: generateId(),
155+
controller,
156+
}).then(
157+
() => undefined,
158+
(error: unknown) => error
159+
)
160+
const [run] = await db.select().from(copilotRuns).where(eq(copilotRuns.chatId, chatId))
161+
return { run, outcome }
162+
}
163+
164+
describe('Slack Assistant run record', () => {
165+
beforeAll(async () => {
166+
hoisted.userId = generateId()
167+
hoisted.organizationId = generateId()
168+
const now = new Date()
169+
await db.insert(user).values({
170+
id: hoisted.userId,
171+
name: 'Slack Assistant fixture',
172+
email: `${hoisted.userId}@slack-assistant.test`,
173+
emailVerified: true,
174+
createdAt: now,
175+
updatedAt: now,
176+
})
177+
await db.insert(organization).values({
178+
id: hoisted.organizationId,
179+
name: 'Slack Assistant fixture',
180+
slug: `slack-assistant-${hoisted.organizationId}`,
181+
})
182+
mothershipChatPayloadMockFns.mockBuildCopilotRequestPayload.mockResolvedValue({
183+
mode: 'assistant',
184+
})
185+
})
186+
187+
afterAll(async () => {
188+
await db.delete(copilotChats).where(eq(copilotChats.userId, hoisted.userId))
189+
await db.delete(organization).where(eq(organization.id, hoisted.organizationId))
190+
await db.delete(user).where(eq(user.id, hoisted.userId))
191+
})
192+
193+
beforeEach(() => {
194+
hoisted.stopped.mockResolvedValue(false)
195+
})
196+
197+
it('records a completed turn as complete', async () => {
198+
hoisted.lifecycle.mockResolvedValueOnce({
199+
success: true,
200+
content: 'Answer',
201+
contentBlocks: [],
202+
toolCalls: [],
203+
})
204+
205+
const { run, outcome } = await slackTurn()
206+
207+
expect(outcome).toBeUndefined()
208+
expect(run.status).toBe('complete')
209+
expect(run.completedAt).not.toBeNull()
210+
})
211+
212+
it('records a turn its user stopped as cancelled', async () => {
213+
hoisted.lifecycle.mockImplementationOnce(async () => {
214+
/** A Slack Stop marks the turn stopped, then aborts its registered stream. */
215+
hoisted.stopped.mockResolvedValue(true)
216+
hoisted.controller.abort(AbortReason.UserStop)
217+
return { success: false, cancelled: true, content: '', contentBlocks: [], toolCalls: [] }
218+
})
219+
220+
const { run } = await slackTurn()
221+
222+
expect(run.status).toBe('cancelled')
223+
expect(run.completedAt).not.toBeNull()
224+
})
225+
226+
it('records a failed turn as an error', async () => {
227+
hoisted.lifecycle.mockResolvedValueOnce({
228+
success: false,
229+
error: 'worker failed',
230+
content: '',
231+
contentBlocks: [],
232+
toolCalls: [],
233+
})
234+
235+
const { run, outcome } = await slackTurn()
236+
237+
expect(outcome).toBeInstanceOf(Error)
238+
expect(run.status).toBe('error')
239+
})
240+
})

‎apps/sim/lib/knowledge/application/slack-search/assistant.ts‎

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,7 @@ import {
4949
startAbortPoller,
5050
unregisterActiveStream,
5151
} from '@/lib/mothership/request/session/abort'
52+
import { isExplicitStopReason } from '@/lib/mothership/request/session/abort-reason'
5253
import type { OrchestratorResult } from '@/lib/mothership/request/types'
5354
import { organizationRoutes } from '@/lib/navigation/paths'
5455
import { SlackSearchAssistantStream } from '@/lib/slack-search/assistant-stream'
@@ -303,7 +304,6 @@ export async function runSlackSearchAssistant(
303304
]
304305
: []),
305306
recordSlackSearchOutcome(installation, 'assistant_or_delivery_failed'),
306-
...(runId ? [updateRunStatus(runId, 'error')] : []),
307307
])
308308
const errors = outcomes.flatMap((outcome) =>
309309
outcome.status === 'rejected' ? [outcome.reason] : []
@@ -317,6 +317,20 @@ export async function runSlackSearchAssistant(
317317
await titleTask
318318
clearInterval(accessPoller)
319319
clearInterval(abortPoller)
320+
try {
321+
if (runId) {
322+
/** This turn admitted its own run, so it records the terminal status no other path will. */
323+
const cancelled =
324+
failed &&
325+
(isExplicitStopReason(controller.signal.reason) ||
326+
(await wasSlackSearchTurnStopped(turnId, leaseId)))
327+
await updateRunStatus(runId, failed ? (cancelled ? 'cancelled' : 'error') : 'complete')
328+
}
329+
} catch (error) {
330+
failure = failure
331+
? new AggregateError([failure, error], 'Slack turn run status could not be recorded')
332+
: toError(error)
333+
}
320334
try {
321335
if (questionPersisted) {
322336
const stopped = failed && (await wasSlackSearchTurnStopped(turnId, leaseId))

0 commit comments

Comments
 (0)