6bf8bebf51
CI / Test and Build (push) Failing after 1s
CI / Migrate Dev DB (push) Has been skipped
CI / Migrate DB (push) Has been skipped
CodeQL / Analyze actions (push) Has been cancelled
CodeQL / Analyze javascript-typescript (push) Has been cancelled
CI / Detect Version (push) Has been cancelled
CI / Detect Desktop Changes (push) Has been cancelled
CI / Build AMD64 (blacksmith-2vcpu-ubuntu-2404, ./docker/cron.Dockerfile, ubuntu-latest, ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Build AMD64 (blacksmith-2vcpu-ubuntu-2404, ./docker/db.Dockerfile, ECR_MIGRATIONS, ubuntu-latest, ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Build AMD64 (blacksmith-4vcpu-ubuntu-2404, ./docker/pii.Dockerfile, ECR_PII, ubuntu-latest, ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Build AMD64 (blacksmith-4vcpu-ubuntu-2404, ./docker/realtime.Dockerfile, ECR_REALTIME, ubuntu-latest, ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build AMD64 (blacksmith-8vcpu-ubuntu-2404, ./docker/app.Dockerfile, ECR_APP, linux-x64-8-core, ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/cron.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/db.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/pii.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/realtime.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-8vcpu-ubuntu-2404-arm, ./docker/app.Dockerfile, linux-arm64-8-core, ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Check Docs Changes (push) Has been cancelled
Publish CLI Package / publish-npm (push) Has been cancelled
Publish Python SDK / publish-pypi (push) Has been cancelled
CI / Deploy Trigger.dev (Dev) (push) Has been cancelled
Helm Chart / Lint, test, and validate chart (push) Has been cancelled
Helm Chart / Chart version bumped (push) Has been cancelled
Publish TypeScript SDK / publish-npm (push) Has been cancelled
CI / Build Dev ECR (blacksmith-8vcpu-ubuntu-2404, ./docker/app.Dockerfile, ECR_APP, linux-x64-8-core) (push) Has been cancelled
CI / Promote Images (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-2vcpu-ubuntu-2404, ./docker/db.Dockerfile, ECR_MIGRATIONS, ubuntu-latest) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-4vcpu-ubuntu-2404, ./docker/pii.Dockerfile, ECR_PII, ubuntu-latest) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-4vcpu-ubuntu-2404, ./docker/realtime.Dockerfile, ECR_REALTIME, ubuntu-latest) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Process Docs (push) Has been cancelled
CI / Create GitHub Release (push) Has been cancelled
CI / Check Desktop Signing Secrets (push) Has been cancelled
CI / Desktop Release (push) Has been cancelled
CI / Create Desktop Prerelease (push) Has been cancelled
CI / Desktop Prerelease Build (push) Has been cancelled
CI / Publish Desktop Prerelease (push) Has been cancelled
CI / Prune Desktop Prereleases (push) Has been cancelled
Helm Chart / Install on kind and run helm test (push) Has been cancelled
94 lines
2.5 KiB
TypeScript
94 lines
2.5 KiB
TypeScript
import { db } from '@sim/db'
|
|
import { idempotencyKey, workflowExecutionLogs } from '@sim/db/schema'
|
|
import { generateId } from '@sim/utils/id'
|
|
import { and, eq, sql } from 'drizzle-orm'
|
|
import type { DbOrTx } from '@/lib/db/types'
|
|
|
|
const EXECUTION_ID_CLAIM_PREFIX = 'workflow-execution-id'
|
|
|
|
interface ExecutionIdClaimRecord {
|
|
status: 'claimed'
|
|
claimToken: string
|
|
executionId: string
|
|
}
|
|
|
|
export interface ExecutionIdClaim {
|
|
key: string
|
|
token: string
|
|
}
|
|
|
|
async function hasDurableExecutionOwnerWithExecutor(
|
|
executor: DbOrTx,
|
|
executionId: string
|
|
): Promise<boolean> {
|
|
const existingExecution = await executor
|
|
.select({ id: workflowExecutionLogs.id })
|
|
.from(workflowExecutionLogs)
|
|
.where(eq(workflowExecutionLogs.executionId, executionId))
|
|
.limit(1)
|
|
|
|
return existingExecution.length > 0
|
|
}
|
|
|
|
/**
|
|
* Atomically reserves a globally unique workflow execution ID.
|
|
*
|
|
* The PostgreSQL primary key is the serialization point. Claims remain as
|
|
* durable tombstones after a run starts so deleting execution logs cannot make
|
|
* a previously consumed ID reusable.
|
|
*/
|
|
export async function claimExecutionId(executionId: string): Promise<ExecutionIdClaim | null> {
|
|
const key = `${EXECUTION_ID_CLAIM_PREFIX}:${executionId}`
|
|
const token = generateId()
|
|
|
|
return db.transaction(async (tx) => {
|
|
const inserted = await tx
|
|
.insert(idempotencyKey)
|
|
.values({
|
|
key,
|
|
result: {
|
|
status: 'claimed',
|
|
claimToken: token,
|
|
executionId,
|
|
} satisfies ExecutionIdClaimRecord,
|
|
createdAt: new Date(),
|
|
})
|
|
.onConflictDoNothing()
|
|
.returning({ key: idempotencyKey.key })
|
|
|
|
if (inserted.length === 0) {
|
|
return null
|
|
}
|
|
|
|
if (await hasDurableExecutionOwnerWithExecutor(tx, executionId)) {
|
|
return null
|
|
}
|
|
|
|
return { key, token }
|
|
})
|
|
}
|
|
|
|
/**
|
|
* Checks whether a durable execution log has taken ownership of an ID.
|
|
*/
|
|
export async function hasDurableExecutionOwner(executionId: string): Promise<boolean> {
|
|
return hasDurableExecutionOwnerWithExecutor(db, executionId)
|
|
}
|
|
|
|
/**
|
|
* Releases only the transient claim owned by this request.
|
|
*
|
|
* Token matching prevents a stale request from deleting another owner's claim.
|
|
*/
|
|
export async function releaseExecutionIdClaim(claim: ExecutionIdClaim): Promise<void> {
|
|
await db
|
|
.delete(idempotencyKey)
|
|
.where(
|
|
and(
|
|
eq(idempotencyKey.key, claim.key),
|
|
sql`${idempotencyKey.result} ->> 'status' = 'claimed'`,
|
|
sql`${idempotencyKey.result} ->> 'claimToken' = ${claim.token}`
|
|
)
|
|
)
|
|
}
|