Skip to content

Commit ac35656

Browse files
committed
fix(execution): publish a pause only after its run log is finalized
1 parent 1b20406 commit ac35656

2 files changed

Lines changed: 193 additions & 0 deletions

File tree

Lines changed: 187 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,187 @@
1+
/**
2+
* Pause publication against real PostgreSQL: a paused run becomes resumable only after
3+
* its log has been finalized out of `running`, so an immediate resume finds a claimable log.
4+
*/
5+
import { db } from '@sim/db'
6+
import {
7+
pausedExecutions,
8+
resumeQueue,
9+
user,
10+
workflow,
11+
workflowExecutionLogs,
12+
workflowExecutionSnapshots,
13+
workspace,
14+
} from '@sim/db/schema'
15+
import { createDeferred } from '@sim/testing'
16+
import { sleep } from '@sim/utils/helpers'
17+
import { generateId } from '@sim/utils/id'
18+
import { eq } from 'drizzle-orm'
19+
import { afterAll, beforeAll, describe, expect, it } from 'vitest'
20+
import {
21+
type BillingAttributionSnapshot,
22+
resolveBillingAttribution,
23+
} from '@/lib/billing/core/billing-attribution'
24+
import { LoggingSession } from '@/lib/logs/execution/logging-session'
25+
import type { WorkflowState } from '@/lib/logs/types'
26+
import { PauseResumeManager } from '@/lib/workflows/executor/human-in-the-loop-manager'
27+
import { handlePostExecutionPauseState } from '@/lib/workflows/executor/pause-persistence'
28+
import type { ExecutionResult } from '@/executor/types'
29+
30+
const ids = {
31+
owner: `pause-publish-owner-${generateId()}`,
32+
workspace: generateId(),
33+
workflow: generateId(),
34+
execution: generateId(),
35+
}
36+
37+
const CONTEXT_ID = 'approval'
38+
39+
/** Long enough for an unguarded publish to commit; a guarded one never resolves while held. */
40+
const UNGUARDED_PUBLISH_WINDOW_MS = 250
41+
42+
const workflowState: WorkflowState = {
43+
blocks: {
44+
start: {
45+
id: 'start',
46+
type: 'starter',
47+
name: 'Start',
48+
position: { x: 0, y: 0 },
49+
subBlocks: {},
50+
outputs: {},
51+
enabled: true,
52+
},
53+
},
54+
edges: [],
55+
loops: {},
56+
parallels: {},
57+
}
58+
59+
function pausedResult(billingAttribution: BillingAttributionSnapshot): ExecutionResult {
60+
return {
61+
success: true,
62+
output: {},
63+
status: 'paused',
64+
pausePoints: [
65+
{
66+
contextId: CONTEXT_ID,
67+
blockId: CONTEXT_ID,
68+
response: {},
69+
registeredAt: new Date().toISOString(),
70+
resumeStatus: 'paused',
71+
snapshotReady: true,
72+
pauseKind: 'human',
73+
},
74+
],
75+
snapshotSeed: {
76+
snapshot: JSON.stringify({
77+
metadata: {
78+
workflowId: ids.workflow,
79+
workspaceId: ids.workspace,
80+
executionId: ids.execution,
81+
userId: ids.owner,
82+
billingAttribution,
83+
},
84+
}),
85+
triggerIds: [],
86+
},
87+
}
88+
}
89+
90+
async function logStatus() {
91+
const [row] = await db
92+
.select({ status: workflowExecutionLogs.status })
93+
.from(workflowExecutionLogs)
94+
.where(eq(workflowExecutionLogs.executionId, ids.execution))
95+
return row?.status
96+
}
97+
98+
function resume() {
99+
return PauseResumeManager.enqueueOrStartResume({
100+
executionId: ids.execution,
101+
workflowId: ids.workflow,
102+
contextId: CONTEXT_ID,
103+
resumeInput: {},
104+
userId: ids.owner,
105+
})
106+
}
107+
108+
beforeAll(async () => {
109+
const now = new Date()
110+
await db.insert(user).values({
111+
id: ids.owner,
112+
name: 'Pause Publish',
113+
email: `${ids.owner}@pause-publish.test`,
114+
emailVerified: true,
115+
createdAt: now,
116+
updatedAt: now,
117+
})
118+
await db.insert(workspace).values({
119+
id: ids.workspace,
120+
name: 'Pause Publish',
121+
ownerId: ids.owner,
122+
billedAccountUserId: ids.owner,
123+
})
124+
await db.insert(workflow).values({
125+
id: ids.workflow,
126+
userId: ids.owner,
127+
workspaceId: ids.workspace,
128+
name: 'Pause Publish',
129+
lastSynced: now,
130+
createdAt: now,
131+
updatedAt: now,
132+
})
133+
})
134+
135+
afterAll(async () => {
136+
await db.delete(resumeQueue).where(eq(resumeQueue.parentExecutionId, ids.execution))
137+
await db.delete(pausedExecutions).where(eq(pausedExecutions.workflowId, ids.workflow))
138+
await db.delete(workflowExecutionLogs).where(eq(workflowExecutionLogs.workflowId, ids.workflow))
139+
await db
140+
.delete(workflowExecutionSnapshots)
141+
.where(eq(workflowExecutionSnapshots.workflowId, ids.workflow))
142+
await db.delete(workspace).where(eq(workspace.id, ids.workspace))
143+
await db.delete(user).where(eq(user.id, ids.owner))
144+
})
145+
146+
describe('handlePostExecutionPauseState', () => {
147+
it('publishes a pause only after the run log is finalized, so an immediate resume finds a claimable log', async () => {
148+
const billingAttribution = await resolveBillingAttribution({
149+
actorUserId: ids.owner,
150+
workspaceId: ids.workspace,
151+
})
152+
const loggingSession = new LoggingSession(ids.workflow, ids.execution, 'api', 'pause-publish')
153+
await loggingSession.safeStart({
154+
userId: ids.owner,
155+
workspaceId: ids.workspace,
156+
billingAttribution,
157+
workflowState,
158+
})
159+
expect(await logStatus()).toBe('running')
160+
161+
/** Holds the core's background log finalizer open, as a slow trace projection would. */
162+
const finalizer = createDeferred<void>()
163+
loggingSession.setPostExecutionPromise(
164+
finalizer.promise.then(() => loggingSession.safeCompleteWithPause({ traceSpans: [] }))
165+
)
166+
167+
const publish = handlePostExecutionPauseState({
168+
result: pausedResult(billingAttribution),
169+
workflowId: ids.workflow,
170+
executionId: ids.execution,
171+
loggingSession,
172+
})
173+
await Promise.race([publish, sleep(UNGUARDED_PUBLISH_WINDOW_MS)])
174+
175+
expect(await logStatus()).toBe('running')
176+
await expect(resume()).rejects.toMatchObject({
177+
name: 'ResumeAdmissionError',
178+
statusCode: 404,
179+
})
180+
181+
finalizer.resolve()
182+
await publish
183+
184+
expect(await logStatus()).toBe('pending')
185+
await expect(resume()).resolves.toMatchObject({ status: 'starting' })
186+
})
187+
})

‎apps/sim/lib/workflows/executor/pause-persistence.ts‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,11 @@ interface HandlePostExecutionPauseStateArgs {
2424
* - If execution is paused with a valid snapshot: persists to `paused_executions` table
2525
* - If execution is paused without a snapshot: marks execution as failed
2626
* - If execution is not paused: processes any queued resume entries
27+
*
28+
* A pause is published only after the core's post-execution logging settles. The
29+
* core finalizes the run log in the background, and a resume claims that log only
30+
* once it has left `running`, so publishing first lets an immediate resume be
31+
* rejected as no longer resumable.
2732
*/
2833
export async function handlePostExecutionPauseState({
2934
result,
@@ -33,6 +38,7 @@ export async function handlePostExecutionPauseState({
3338
loggingSession,
3439
}: HandlePostExecutionPauseStateArgs): Promise<void> {
3540
if (result.status === 'paused') {
41+
await loggingSession.waitForPostExecution()
3642
if (!result.snapshotSeed) {
3743
logger.error('Missing snapshot seed for paused execution', { executionId })
3844
await loggingSession.markAsFailed('Missing snapshot seed for paused execution')

0 commit comments

Comments
 (0)