Skip to content

Commit 3ff5a3f

Browse files
committed
fix(knowledge): record the Slack run's status after its turn is saved
- The Slack Assistant now writes its run's one terminal status after its outcome and response are persisted, from the final outcome, so a failed save ends the run as an error instead of complete. - A Stop lookup that fails no longer skips that write: the turn is treated as not stopped, logged, and settled as an error. - The orphaned-run suite deletes the sweep cursor before each test and in teardown, so no later suite starts from its leftover position.
1 parent 7c023cf commit 3ff5a3f

3 files changed

Lines changed: 84 additions & 18 deletions

File tree

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

Lines changed: 49 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,8 @@ const hoisted = vi.hoisted(() => ({
2222
/** The turn's own controller, which a Stop aborts through its registered stream. */
2323
controller: new AbortController(),
2424
stopped: vi.fn(async () => false),
25+
outcome: vi.fn(async () => undefined),
26+
finalize: vi.fn(async () => ({ appendedAssistant: true })),
2527
}))
2628
vi.mock('@/lib/knowledge/application/slack-search/authorization', () => ({
2729
authorizeSlackSearchInstallation: async () => ({
@@ -46,7 +48,7 @@ vi.mock('@/lib/knowledge/application/operations', () => ({
4648
knowledgeOperations: { search: { organizationOperation: { id: 'knowledge.search' } } },
4749
}))
4850
vi.mock('@/lib/knowledge/application/slack-search/repository', () => ({
49-
recordSlackSearchOutcome: async () => undefined,
51+
recordSlackSearchOutcome: hoisted.outcome,
5052
}))
5153
vi.mock('@/lib/knowledge/application/slack-search/turns', () => ({
5254
requireSlackSearchTurnLease: async () => undefined,
@@ -64,7 +66,7 @@ vi.mock('@/lib/mothership/application/load-search-integrations', () => ({
6466
}))
6567
vi.mock('@/lib/mothership/chat/payload', () => mothershipChatPayloadMock)
6668
vi.mock('@/lib/mothership/chat/terminal-state', () => ({
67-
finalizeAssistantTurn: async () => ({ appendedAssistant: true }),
69+
finalizeAssistantTurn: hoisted.finalize,
6870
}))
6971
vi.mock('@/lib/mothership/environment-context', () => mothershipEnvironmentContextMock)
7072
vi.mock('@/lib/mothership/request/lifecycle/headless', () => ({
@@ -192,8 +194,17 @@ describe('Slack Assistant run record', () => {
192194

193195
beforeEach(() => {
194196
hoisted.stopped.mockResolvedValue(false)
197+
hoisted.outcome.mockResolvedValue(undefined)
198+
hoisted.finalize.mockResolvedValue({ appendedAssistant: true })
195199
})
196200

201+
const answered = {
202+
success: true,
203+
content: 'Answer',
204+
contentBlocks: [],
205+
toolCalls: [],
206+
}
207+
197208
it('records a completed turn as complete', async () => {
198209
hoisted.lifecycle.mockResolvedValueOnce({
199210
success: true,
@@ -237,4 +248,40 @@ describe('Slack Assistant run record', () => {
237248
expect(outcome).toBeInstanceOf(Error)
238249
expect(run.status).toBe('error')
239250
})
251+
252+
it('records a failed turn as an error even when its Stop cannot be looked up', async () => {
253+
hoisted.stopped.mockRejectedValue(new Error('database unavailable'))
254+
hoisted.lifecycle.mockResolvedValueOnce({
255+
success: false,
256+
error: 'worker failed',
257+
content: '',
258+
contentBlocks: [],
259+
toolCalls: [],
260+
})
261+
262+
const { run, outcome } = await slackTurn()
263+
264+
expect(outcome).toBeInstanceOf(Error)
265+
expect(run.status).toBe('error')
266+
})
267+
268+
it('records an answered turn as an error when its response is not saved', async () => {
269+
hoisted.lifecycle.mockResolvedValueOnce(answered)
270+
hoisted.finalize.mockResolvedValue({ appendedAssistant: false })
271+
272+
const { run, outcome } = await slackTurn()
273+
274+
expect(outcome).toBeInstanceOf(Error)
275+
expect(run.status).toBe('error')
276+
})
277+
278+
it('records an answered turn as an error when its outcome is not saved', async () => {
279+
hoisted.lifecycle.mockResolvedValueOnce(answered)
280+
hoisted.outcome.mockRejectedValueOnce(new Error('outcome write failed'))
281+
282+
const { run, outcome } = await slackTurn()
283+
284+
expect(outcome).toBeInstanceOf(Error)
285+
expect(run.status).toBe('error')
286+
})
240287
})

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

Lines changed: 26 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@ import type {
33
SlackInstallationPrincipal,
44
} from '@sim/auth/principal'
55
import { createLogger } from '@sim/logger'
6-
import { toError } from '@sim/utils/errors'
6+
import { getErrorMessage, toError } from '@sim/utils/errors'
77
import { generateId } from '@sim/utils/id'
88
import { isRecordLike } from '@sim/utils/object'
99
import { resolveOrganizationBillingAttribution } from '@/lib/billing/core/billing-attribution'
@@ -317,20 +317,6 @@ 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-
}
334320
try {
335321
if (questionPersisted) {
336322
const stopped = failed && (await wasSlackSearchTurnStopped(turnId, leaseId))
@@ -392,6 +378,31 @@ export async function runSlackSearchAssistant(
392378
? new AggregateError([failure, error], 'Slack turn and history persistence failed')
393379
: toError(error)
394380
} finally {
381+
if (runId) {
382+
/**
383+
* This turn admitted its own run, so it records the one terminal status no other
384+
* path will, after its outcome and response were saved: any failure, including
385+
* a failed save, ends it as an error unless its user stopped it.
386+
*/
387+
let cancelled = failed && isExplicitStopReason(controller.signal.reason)
388+
if (failed && !cancelled) {
389+
try {
390+
cancelled = await wasSlackSearchTurnStopped(turnId, leaseId)
391+
} catch (error) {
392+
logger.warn('Slack turn Stop could not be read; recording its run as an error', {
393+
turnId,
394+
error: getErrorMessage(error),
395+
})
396+
}
397+
}
398+
try {
399+
await updateRunStatus(runId, !failure ? 'complete' : cancelled ? 'cancelled' : 'error')
400+
} catch (error) {
401+
failure = failure
402+
? new AggregateError([failure, error], 'Slack turn run status could not be recorded')
403+
: toError(error)
404+
}
405+
}
395406
unregisterActiveStream(messageId)
396407
await releasePendingChatStream(chat.id, messageId)
397408
await cleanupAbortMarker(messageId)

‎apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@
44
* use case are production code. A local HTTP server stands in for the worker's abort
55
* endpoint.
66
*/
7-
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
7+
import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'
88

99
const { redisUrl, inheritedEnv, worker } = await vi.hoisted(async () => {
1010
const { readTestRedisUrl } = await import('@sim/db/testing/test-infrastructure')
@@ -68,8 +68,16 @@ function redis() {
6868
return client
6969
}
7070

71+
/** The sweep's resume point lives in shared Redis; each test and the next suite start fresh. */
72+
async function resetSweepCursor() {
73+
await getRedisClient()?.del('copilot:orphaned-runs:sweep-cursor')
74+
}
75+
76+
beforeEach(resetSweepCursor)
77+
7178
afterAll(async () => {
7279
chatPubSub?.dispose()
80+
await resetSweepCursor()
7381
await closeRedisConnection()
7482
await new Promise<void>((resolve) => worker.server.close(() => resolve()))
7583
for (const [key, value] of Object.entries(inheritedEnv)) {

0 commit comments

Comments
 (0)