mirror of
https://github.com/simstudioai/sim.git
synced 2026-09-24 15:45:35 +08:00
improvement(execution): stop rewriting execution snapshots on reuse + skip redundant actor lookup (#5242)
* improvement(execution): stop rewriting execution snapshots on reuse + skip redundant actor lookup - SnapshotService.createSnapshotWithDeduplication: switch the per-execution dedup write from onConflictDoUpdate(set state_data) to onConflictDoNothing + a conditional select. A (workflowId, stateHash) row is byte-identical by hash, so rewriting the full state jsonb every run only churned a dead tuple + TOAST/WAL under Postgres MVCC. The reuse path (the common case) now performs no write. - preprocessExecution: add an optional resolvedActorUserId so a caller that already resolved the billing actor upstream can skip the redundant workspace billed-account lookup. The ban/usage/rate/archived gates still run against the actor — only the resolution is reused, never a gate. The webhook background job passes the route-resolved payload.userId. * fix(webhooks): scope actor reuse to inline execution only Addresses review: a queued/Trigger.dev webhook can outlive a workspace billed-account change, so reusing the route-resolved actor there could gate against a stale account. Set resolvedActorUserId only on the in-process inline payload (sub-second after resolution); queued and persisted payloads omit it, so the background pass re-resolves the current billed account. Gates unchanged. * docs(webhooks): convert inline comments on actor-reuse to TSDoc * fix(logs): keep snapshot dedup a single atomic upsert (no select race) Addresses review: the DO NOTHING + follow-up select could fail if cleanup deletes the conflicting (orphaned, aged) snapshot between the no-op insert and the select. Revert to one atomic upsert but SET only state_hash, so RETURNING always yields the row (no race) while the unchanged TOASTed state_data jsonb is still not rewritten under MVCC — keeping the per-execution write tiny.
This commit is contained in:
@@ -241,6 +241,14 @@ export type WebhookExecutionPayload = {
|
||||
webhookReceivedAt?: number
|
||||
/** Epoch ms of the originating provider interaction (e.g. Slack x-slack-request-timestamp). */
|
||||
triggerTimestampMs?: number
|
||||
/**
|
||||
* Billing actor resolved by the webhook route, set ONLY for in-process inline
|
||||
* execution that runs microseconds after resolution. The background pass reuses
|
||||
* it to skip the redundant billed-account lookup. Deliberately absent on queued
|
||||
* (Trigger.dev) and persisted payloads — a deferred run could outlive a
|
||||
* billed-account change, so it re-resolves the current actor instead.
|
||||
*/
|
||||
resolvedActorUserId?: string
|
||||
}
|
||||
|
||||
export async function executeWebhookJob(payload: WebhookExecutionPayload) {
|
||||
@@ -367,6 +375,13 @@ async function executeWebhookJobInternal(
|
||||
skipUsageLimits: true,
|
||||
workspaceId: payload.workspaceId,
|
||||
loggingSession,
|
||||
/**
|
||||
* Reuse the route-resolved actor only for inline execution (set on the
|
||||
* in-process payload). When absent — queued/Trigger.dev runs — preprocessing
|
||||
* re-resolves the current billed account. Either way the ban and
|
||||
* archived-workflow gates run fresh against the resolved actor.
|
||||
*/
|
||||
resolvedActorUserId: payload.resolvedActorUserId,
|
||||
})
|
||||
|
||||
if (!preprocessResult.success) {
|
||||
|
||||
@@ -311,3 +311,63 @@ describe('preprocessExecution ban gate', () => {
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
describe('preprocessExecution resolvedActorUserId reuse', () => {
|
||||
const baseOptions = {
|
||||
workflowId: 'workflow-1',
|
||||
userId: 'owner-1',
|
||||
triggerType: 'webhook' as const,
|
||||
executionId: 'execution-1',
|
||||
requestId: 'request-1',
|
||||
checkDeployment: false,
|
||||
checkRateLimit: false,
|
||||
skipConcurrencyReservation: true,
|
||||
workspaceId: 'workspace-1',
|
||||
workflowRecord: { id: 'workflow-1', workspaceId: 'workspace-1', isDeployed: true } as any,
|
||||
}
|
||||
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks()
|
||||
mockGetWorkspaceBilledAccountUserId.mockResolvedValue('billed-account-1')
|
||||
mockGetActivelyBannedUserIds.mockResolvedValue([])
|
||||
vi.mocked(getHighestPrioritySubscription).mockResolvedValue({ plan: 'free' } as any)
|
||||
vi.mocked(checkServerSideUsageLimits).mockResolvedValue({
|
||||
isExceeded: false,
|
||||
currentUsage: 1,
|
||||
limit: 10,
|
||||
} as any)
|
||||
})
|
||||
|
||||
it('skips the workspace billed-account lookup when an actor is pre-resolved', async () => {
|
||||
const result = await preprocessExecution({
|
||||
...baseOptions,
|
||||
resolvedActorUserId: 'pre-resolved-actor',
|
||||
})
|
||||
|
||||
expect(result.success).toBe(true)
|
||||
expect(result.actorUserId).toBe('pre-resolved-actor')
|
||||
expect(mockGetWorkspaceBilledAccountUserId).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('still runs the ban gate against the pre-resolved actor', async () => {
|
||||
mockGetActivelyBannedUserIds.mockResolvedValue(['pre-resolved-actor'])
|
||||
|
||||
const result = await preprocessExecution({
|
||||
...baseOptions,
|
||||
resolvedActorUserId: 'pre-resolved-actor',
|
||||
})
|
||||
|
||||
expect(result).toMatchObject({ success: false, error: { statusCode: 403 } })
|
||||
expect(mockGetActivelyBannedUserIds).toHaveBeenCalledWith(
|
||||
expect.arrayContaining(['pre-resolved-actor'])
|
||||
)
|
||||
})
|
||||
|
||||
it('falls back to the billed-account lookup when no actor is pre-resolved', async () => {
|
||||
const result = await preprocessExecution(baseOptions)
|
||||
|
||||
expect(result.success).toBe(true)
|
||||
expect(result.actorUserId).toBe('billed-account-1')
|
||||
expect(mockGetWorkspaceBilledAccountUserId).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -60,6 +60,14 @@ export interface PreprocessExecutionOptions {
|
||||
useDraftState?: boolean
|
||||
/** Pre-fetched workflow row for caller context; preprocessing still re-checks active state. */
|
||||
workflowRecord?: WorkflowRecord
|
||||
/**
|
||||
* Billing actor already resolved by an upstream gate earlier in the same
|
||||
* request (e.g. the webhook route's preprocessing pass, whose result is carried
|
||||
* as the job's userId). When provided, the redundant workspace billed-account
|
||||
* lookup is skipped. The ban, deployment, usage, and rate-limit gates still run
|
||||
* against this actor — only the resolution is reused, never a gate.
|
||||
*/
|
||||
resolvedActorUserId?: string
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -111,6 +119,7 @@ export async function preprocessExecution(
|
||||
isResumeContext: _isResumeContext = false,
|
||||
useAuthenticatedUserAsActor = false,
|
||||
workflowRecord: prefetchedWorkflowRecord,
|
||||
resolvedActorUserId,
|
||||
} = options
|
||||
|
||||
// When `logPreprocessingErrors` is false the caller surfaces failures itself
|
||||
@@ -260,6 +269,16 @@ export async function preprocessExecution(
|
||||
logger.info(`[${requestId}] Using authenticated user as actor: ${actorUserId}`)
|
||||
}
|
||||
|
||||
/**
|
||||
* Reuse an actor already resolved upstream this request (e.g. the webhook
|
||||
* route's preprocessing) to skip the redundant workspace billed-account
|
||||
* lookup. Gates below still run against this actor.
|
||||
*/
|
||||
if (!actorUserId && resolvedActorUserId) {
|
||||
actorUserId = resolvedActorUserId
|
||||
logger.info(`[${requestId}] Using pre-resolved billing actor: ${actorUserId}`)
|
||||
}
|
||||
|
||||
if (!actorUserId && workspaceId) {
|
||||
actorUserId = await getWorkspaceBilledAccountUserId(workspaceId)
|
||||
if (actorUserId) {
|
||||
|
||||
@@ -373,11 +373,32 @@ describe('SnapshotService', () => {
|
||||
})
|
||||
|
||||
describe('createSnapshotWithDeduplication', () => {
|
||||
it('should use upsert to insert a new snapshot', async () => {
|
||||
type SnapshotRow = {
|
||||
id: string
|
||||
workflowId: string
|
||||
stateHash: string
|
||||
stateData: WorkflowState
|
||||
createdAt: Date
|
||||
}
|
||||
|
||||
/** Mock the insert → values → onConflictDoUpdate → returning chain. */
|
||||
function mockUpsertReturning(rows: SnapshotRow[]) {
|
||||
let capturedConflictConfig: Record<string, unknown> | undefined
|
||||
const onConflictDoUpdate = vi.fn().mockImplementation((config: Record<string, unknown>) => {
|
||||
capturedConflictConfig = config
|
||||
return { returning: vi.fn().mockResolvedValue(rows) }
|
||||
})
|
||||
const values = vi.fn().mockReturnValue({ onConflictDoUpdate })
|
||||
databaseMock.db.insert = vi.fn().mockReturnValue({ values })
|
||||
databaseMock.db.select = vi.fn()
|
||||
return { values, onConflictDoUpdate, getConflictConfig: () => capturedConflictConfig }
|
||||
}
|
||||
|
||||
it('inserts a new snapshot in a single atomic upsert', async () => {
|
||||
const service = new SnapshotService()
|
||||
const workflowId = 'wf-123'
|
||||
|
||||
const mockReturning = vi.fn().mockResolvedValue([
|
||||
const { values } = mockUpsertReturning([
|
||||
{
|
||||
id: 'generated-uuid-1',
|
||||
workflowId,
|
||||
@@ -386,35 +407,23 @@ describe('SnapshotService', () => {
|
||||
createdAt: new Date('2026-02-19T00:00:00Z'),
|
||||
},
|
||||
])
|
||||
const mockOnConflictDoUpdate = vi.fn().mockReturnValue({ returning: mockReturning })
|
||||
const mockValues = vi.fn().mockReturnValue({ onConflictDoUpdate: mockOnConflictDoUpdate })
|
||||
const mockInsert = vi.fn().mockReturnValue({ values: mockValues })
|
||||
databaseMock.db.insert = mockInsert
|
||||
|
||||
const result = await service.createSnapshotWithDeduplication(workflowId, mockState)
|
||||
|
||||
expect(mockInsert).toHaveBeenCalled()
|
||||
expect(mockValues).toHaveBeenCalledWith(
|
||||
expect.objectContaining({
|
||||
id: 'generated-uuid-1',
|
||||
workflowId,
|
||||
stateData: mockState,
|
||||
})
|
||||
)
|
||||
expect(mockOnConflictDoUpdate).toHaveBeenCalledWith(
|
||||
expect.objectContaining({
|
||||
set: expect.any(Object),
|
||||
})
|
||||
expect(values).toHaveBeenCalledWith(
|
||||
expect.objectContaining({ id: 'generated-uuid-1', workflowId, stateData: mockState })
|
||||
)
|
||||
expect(result.snapshot.id).toBe('generated-uuid-1')
|
||||
expect(result.isNew).toBe(true)
|
||||
// Single atomic statement — never a follow-up select (which would race with cleanup).
|
||||
expect(databaseMock.db.select).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('should detect reused snapshot when returned id differs from generated id', async () => {
|
||||
it('reuses the existing snapshot atomically when the returned id differs', async () => {
|
||||
const service = new SnapshotService()
|
||||
const workflowId = 'wf-123'
|
||||
|
||||
const mockReturning = vi.fn().mockResolvedValue([
|
||||
mockUpsertReturning([
|
||||
{
|
||||
id: 'existing-snapshot-id',
|
||||
workflowId,
|
||||
@@ -423,22 +432,19 @@ describe('SnapshotService', () => {
|
||||
createdAt: new Date('2026-02-19T00:00:00Z'),
|
||||
},
|
||||
])
|
||||
const mockOnConflictDoUpdate = vi.fn().mockReturnValue({ returning: mockReturning })
|
||||
const mockValues = vi.fn().mockReturnValue({ onConflictDoUpdate: mockOnConflictDoUpdate })
|
||||
const mockInsert = vi.fn().mockReturnValue({ values: mockValues })
|
||||
databaseMock.db.insert = mockInsert
|
||||
|
||||
const result = await service.createSnapshotWithDeduplication(workflowId, mockState)
|
||||
|
||||
expect(result.snapshot.id).toBe('existing-snapshot-id')
|
||||
expect(result.isNew).toBe(false)
|
||||
expect(databaseMock.db.select).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('should not throw on concurrent inserts with the same hash', async () => {
|
||||
it('SET targets only state_hash on conflict, never the large state_data', async () => {
|
||||
const service = new SnapshotService()
|
||||
const workflowId = 'wf-123'
|
||||
|
||||
const mockReturningNew = vi.fn().mockResolvedValue([
|
||||
const { onConflictDoUpdate, getConflictConfig } = mockUpsertReturning([
|
||||
{
|
||||
id: 'generated-uuid-1',
|
||||
workflowId,
|
||||
@@ -447,23 +453,38 @@ describe('SnapshotService', () => {
|
||||
createdAt: new Date('2026-02-19T00:00:00Z'),
|
||||
},
|
||||
])
|
||||
const mockReturningExisting = vi.fn().mockResolvedValue([
|
||||
{
|
||||
id: 'existing-snapshot-id',
|
||||
workflowId,
|
||||
stateHash: 'abc123',
|
||||
stateData: mockState,
|
||||
createdAt: new Date('2026-02-19T00:00:00Z'),
|
||||
},
|
||||
])
|
||||
|
||||
let callCount = 0
|
||||
await service.createSnapshotWithDeduplication(workflowId, mockState)
|
||||
|
||||
expect(onConflictDoUpdate).toHaveBeenCalledTimes(1)
|
||||
const config = getConflictConfig()
|
||||
expect(config?.target).toBeDefined()
|
||||
// The crux of this change: the SET touches state_hash only, so the unchanged
|
||||
// TOASTed state_data jsonb is never rewritten.
|
||||
expect(config?.set).toHaveProperty('stateHash')
|
||||
expect(config?.set).not.toHaveProperty('stateData')
|
||||
})
|
||||
|
||||
it('does not throw on concurrent inserts with the same hash', async () => {
|
||||
const service = new SnapshotService()
|
||||
const workflowId = 'wf-123'
|
||||
|
||||
const newRow: SnapshotRow = {
|
||||
id: 'generated-uuid-1',
|
||||
workflowId,
|
||||
stateHash: 'abc123',
|
||||
stateData: mockState,
|
||||
createdAt: new Date('2026-02-19T00:00:00Z'),
|
||||
}
|
||||
const existingRow: SnapshotRow = { ...newRow, id: 'existing-snapshot-id' }
|
||||
|
||||
let upsertCall = 0
|
||||
databaseMock.db.insert = vi.fn().mockImplementation(() => ({
|
||||
values: vi.fn().mockImplementation(() => ({
|
||||
onConflictDoUpdate: vi.fn().mockImplementation(() => ({
|
||||
returning: callCount++ === 0 ? mockReturningNew : mockReturningExisting,
|
||||
})),
|
||||
})),
|
||||
values: vi.fn().mockReturnValue({
|
||||
onConflictDoUpdate: vi.fn().mockReturnValue({
|
||||
returning: vi.fn().mockResolvedValue(upsertCall++ === 0 ? [newRow] : [existingRow]),
|
||||
}),
|
||||
}),
|
||||
}))
|
||||
|
||||
const [result1, result2] = await Promise.all([
|
||||
@@ -476,62 +497,6 @@ describe('SnapshotService', () => {
|
||||
expect(result2.snapshot.id).toBe('existing-snapshot-id')
|
||||
expect(result2.isNew).toBe(false)
|
||||
})
|
||||
|
||||
it('should pass state_data in the ON CONFLICT SET clause', async () => {
|
||||
const service = new SnapshotService()
|
||||
const workflowId = 'wf-123'
|
||||
|
||||
let capturedConflictConfig: Record<string, unknown> | undefined
|
||||
const mockReturning = vi.fn().mockResolvedValue([
|
||||
{
|
||||
id: 'generated-uuid-1',
|
||||
workflowId,
|
||||
stateHash: 'abc123',
|
||||
stateData: mockState,
|
||||
createdAt: new Date('2026-02-19T00:00:00Z'),
|
||||
},
|
||||
])
|
||||
|
||||
databaseMock.db.insert = vi.fn().mockReturnValue({
|
||||
values: vi.fn().mockReturnValue({
|
||||
onConflictDoUpdate: vi.fn().mockImplementation((config: Record<string, unknown>) => {
|
||||
capturedConflictConfig = config
|
||||
return { returning: mockReturning }
|
||||
}),
|
||||
}),
|
||||
})
|
||||
|
||||
await service.createSnapshotWithDeduplication(workflowId, mockState)
|
||||
|
||||
expect(capturedConflictConfig).toBeDefined()
|
||||
expect(capturedConflictConfig!.target).toBeDefined()
|
||||
expect(capturedConflictConfig!.set).toBeDefined()
|
||||
expect(capturedConflictConfig!.set).toHaveProperty('stateData')
|
||||
})
|
||||
|
||||
it('should always call insert, never a separate select for deduplication', async () => {
|
||||
const service = new SnapshotService()
|
||||
const workflowId = 'wf-123'
|
||||
|
||||
const mockReturning = vi.fn().mockResolvedValue([
|
||||
{
|
||||
id: 'generated-uuid-1',
|
||||
workflowId,
|
||||
stateHash: 'abc123',
|
||||
stateData: mockState,
|
||||
createdAt: new Date('2026-02-19T00:00:00Z'),
|
||||
},
|
||||
])
|
||||
const mockOnConflictDoUpdate = vi.fn().mockReturnValue({ returning: mockReturning })
|
||||
const mockValues = vi.fn().mockReturnValue({ onConflictDoUpdate: mockOnConflictDoUpdate })
|
||||
databaseMock.db.insert = vi.fn().mockReturnValue({ values: mockValues })
|
||||
databaseMock.db.select = vi.fn()
|
||||
|
||||
await service.createSnapshotWithDeduplication(workflowId, mockState)
|
||||
|
||||
expect(databaseMock.db.insert).toHaveBeenCalledTimes(1)
|
||||
expect(databaseMock.db.select).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
|
||||
describe('cleanupOrphanedSnapshots', () => {
|
||||
|
||||
@@ -37,13 +37,27 @@ export class SnapshotService implements ISnapshotService {
|
||||
stateData: state,
|
||||
}
|
||||
|
||||
/**
|
||||
* Insert the snapshot, or — when an identical (workflowId, stateHash) row
|
||||
* already exists — return it without rewriting the large stateData jsonb.
|
||||
*
|
||||
* The hash is a sha256 of the normalized state, so an existing row's stateData
|
||||
* is byte-identical; there is nothing to update. The previous implementation
|
||||
* SET state_data on conflict, which rewrote the full (tens-of-KB) jsonb every
|
||||
* run. We keep a single atomic upsert — so RETURNING always yields the row and
|
||||
* there is no race with snapshot cleanup (unlike DO NOTHING + a follow-up
|
||||
* select) — but SET only the small state_hash column to itself. Under Postgres
|
||||
* MVCC the unchanged, TOASTed stateData is not rewritten: its existing
|
||||
* out-of-line storage is reused, so the per-execution write drops from the
|
||||
* full blob to a tiny heap tuple.
|
||||
*/
|
||||
const [upsertedSnapshot] = await db
|
||||
.insert(workflowExecutionSnapshots)
|
||||
.values(snapshotData)
|
||||
.onConflictDoUpdate({
|
||||
target: [workflowExecutionSnapshots.workflowId, workflowExecutionSnapshots.stateHash],
|
||||
set: {
|
||||
stateData: sql`excluded.state_data`,
|
||||
stateHash: sql`excluded.state_hash`,
|
||||
},
|
||||
})
|
||||
.returning()
|
||||
|
||||
@@ -667,10 +667,18 @@ export async function queueWebhookExecution(
|
||||
`[${options.requestId}] Queued ${foundWebhook.provider} webhook execution ${jobId} via inline backend`
|
||||
)
|
||||
|
||||
/**
|
||||
* Inline runs in-process microseconds after the route resolved the actor,
|
||||
* so reuse it to skip a redundant billed-account lookup. Set only on the
|
||||
* in-process payload — the enqueued/persisted copy omits it so any deferred
|
||||
* re-run re-resolves the current billed account.
|
||||
*/
|
||||
const inlinePayload = { ...payload, resolvedActorUserId: actorUserId }
|
||||
|
||||
void (async () => {
|
||||
try {
|
||||
await jobQueue.startJob(jobId)
|
||||
const output = await executeWebhookJob(payload)
|
||||
const output = await executeWebhookJob(inlinePayload)
|
||||
await jobQueue.completeJob(jobId, output)
|
||||
} catch (error) {
|
||||
const errorMessage = toError(error).message
|
||||
|
||||
Reference in New Issue
Block a user