Skip to content

Commit 1287244

Browse files
committed
fix(execution): scope the snapshot FK retry, keep a suspended identity's access read out of the run, tidy tests
1 parent ae920f6 commit 1287244

4 files changed

Lines changed: 33 additions & 16 deletions

File tree

‎apps/sim/lib/core/async-jobs/backends/trigger-dev.test.ts‎

Lines changed: 11 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -274,30 +274,31 @@ describe('TriggerDevJobQueue status mapping', () => {
274274
})
275275
})
276276

277-
it('resolves a caller-chosen job id through its tag without retrieving the id itself', async () => {
278-
mockList.mockReturnValueOnce(createListPage([{ id: 'run_1', tags: ['jobId:schedule_abc'] }]))
279-
mockRetrieve.mockResolvedValueOnce({
280-
id: 'run_1',
281-
payload: {},
282-
status: 'QUEUED',
283-
taskIdentifier: 'schedule-execution',
277+
/** Trigger.dev can only retrieve its own run ids; anything else fails, and slowly. */
278+
function retrieveOnlyRunIds() {
279+
mockRetrieve.mockImplementation(async (id: string) => {
280+
if (!id.startsWith('run_')) throw new Error(`Retrieved a caller-chosen job id: ${id}`)
281+
return { id, payload: {}, status: 'QUEUED', taskIdentifier: 'schedule-execution' }
284282
})
283+
}
284+
285+
it('resolves a caller-chosen job id through its tag', async () => {
286+
retrieveOnlyRunIds()
287+
mockList.mockReturnValueOnce(createListPage([{ id: 'run_1', tags: ['jobId:schedule_abc'] }]))
285288
const queue = new TriggerDevJobQueue()
286289

287290
await expect(queue.getJob('schedule_abc')).resolves.toMatchObject({
288291
id: 'run_1',
289292
status: 'pending',
290293
})
291-
expect(mockRetrieve).toHaveBeenCalledTimes(1)
292-
expect(mockRetrieve).toHaveBeenCalledWith('run_1')
293294
})
294295

295296
it('returns null for a caller-chosen job id with no tagged run', async () => {
297+
retrieveOnlyRunIds()
296298
mockList.mockReturnValueOnce(createListPage([]))
297299
const queue = new TriggerDevJobQueue()
298300

299301
await expect(queue.getJob('schedule_abc')).resolves.toBeNull()
300-
expect(mockRetrieve).not.toHaveBeenCalled()
301302
})
302303

303304
it('falls back to the tag lookup when a run id is not found', async () => {

‎apps/sim/lib/environment/utils.ts‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -538,12 +538,15 @@ export async function getExecutionEnvironment(
538538
return getPersonalAndWorkspaceEnv(personalUserId, workspaceId)
539539
}
540540

541-
const [suspended, actorAccess, personalAccess] = await Promise.all([
541+
const personalAccessRead = checkWorkspaceAccess(workspaceId, personalUserId)
542+
personalAccessRead.catch(() => {})
543+
const [suspended, actorAccess] = await Promise.all([
542544
personalIdentitySuspended,
543545
checkWorkspaceAccess(workspaceId, workspaceUserId),
544-
checkWorkspaceAccess(workspaceId, personalUserId),
545546
])
547+
// A suspended identity's access is never consulted, so its read cannot fail the run.
546548
if (suspended) return resolveSuspendedIdentity()
549+
const personalAccess = await personalAccessRead
547550

548551
/**
549552
* A workspace that no longer exists and one an identity may not read are

‎apps/sim/lib/logs/execution/logger.ts‎

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,12 @@ import {
88
workspace,
99
} from '@sim/db/schema'
1010
import { createLogger } from '@sim/logger'
11-
import { describeError, getErrorMessage, getPostgresErrorCode } from '@sim/utils/errors'
11+
import {
12+
describeError,
13+
getErrorMessage,
14+
getPostgresConstraintName,
15+
getPostgresErrorCode,
16+
} from '@sim/utils/errors'
1217
import { generateId } from '@sim/utils/id'
1318
import { and, eq, inArray, sql } from 'drizzle-orm'
1419
import { checkUsageStatus as checkResolvedUsageStatus } from '@/lib/billing/calculations/usage-monitor'
@@ -95,6 +100,8 @@ const EXECUTION_LOG_IDLE_TIMEOUT_MS = 5_000
95100
// (favor waiting over dropping a charge); only trips on a pathological lock hold.
96101
const USAGE_RECONCILE_LOCK_TIMEOUT_MS = 10_000
97102
const FOREIGN_KEY_VIOLATION = '23503'
103+
/** The log's snapshot foreign key as Postgres names it (identifiers are cut at 63 bytes). */
104+
const STATE_SNAPSHOT_FOREIGN_KEY = 'workflow_execution_logs_state_snapshot_id_workflow_execution_sn'
98105

99106
type ExecutionData = WorkflowExecutionLog['executionData']
100107

@@ -667,7 +674,12 @@ export class ExecutionLogger implements IExecutionLoggerService {
667674
try {
668675
inserted = await insertRunningLog(snapshot.id)
669676
} catch (error) {
670-
if (getPostgresErrorCode(error) !== FOREIGN_KEY_VIOLATION) throw error
677+
if (
678+
getPostgresErrorCode(error) !== FOREIGN_KEY_VIOLATION ||
679+
getPostgresConstraintName(error) !== STATE_SNAPSHOT_FOREIGN_KEY
680+
) {
681+
throw error
682+
}
671683
/**
672684
* A snapshot resolved before the insert can be deleted underneath it by
673685
* orphan cleanup when no log references it yet. Resolve it again from the

‎apps/sim/lib/logs/execution/start-execution.integration.ts‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -102,11 +102,12 @@ beforeAll(async () => {
102102
})
103103

104104
afterAll(async () => {
105-
await db.delete(workspace).where(eq(workspace.id, ids.workspace))
105+
// Snapshots outlive their workflow (workflow_id is set null), so remove them first.
106+
await db.delete(workflowExecutionLogs).where(eq(workflowExecutionLogs.workflowId, ids.workflow))
106107
await db
107108
.delete(workflowExecutionSnapshots)
108109
.where(eq(workflowExecutionSnapshots.workflowId, ids.workflow))
109-
await db.delete(workflow).where(eq(workflow.id, ids.workflow))
110+
await db.delete(workspace).where(eq(workspace.id, ids.workspace))
110111
await db.delete(user).where(eq(user.id, ids.owner))
111112
})
112113

0 commit comments

Comments
 (0)