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
105 lines
3.3 KiB
TypeScript
105 lines
3.3 KiB
TypeScript
import { createLogger } from '@sim/logger'
|
|
import { getRedisClient } from '@/lib/core/config/redis'
|
|
import { createPubSubChannel, type PubSubChannel } from '@/lib/events/pubsub'
|
|
|
|
const logger = createLogger('ExecutionCancellation')
|
|
|
|
const EXECUTION_CANCEL_PREFIX = 'execution:cancel:'
|
|
export const EXECUTION_CANCEL_MIN_RETENTION_MS = 14 * 24 * 60 * 60 * 1000
|
|
const EXECUTION_CANCEL_CHANNEL = 'execution:cancel'
|
|
|
|
export interface MarkExecutionCancelledOptions {
|
|
executionDeadlineAt?: Date | null
|
|
}
|
|
|
|
export interface ExecutionCancelEvent {
|
|
executionId: string
|
|
}
|
|
|
|
export type ExecutionCancellationRecordResult =
|
|
| { durablyRecorded: true; reason: 'recorded' }
|
|
| {
|
|
durablyRecorded: false
|
|
reason: 'redis_unavailable' | 'redis_write_failed'
|
|
}
|
|
|
|
type CancellationGlobal = typeof globalThis & {
|
|
_executionCancelChannel?: PubSubChannel<ExecutionCancelEvent>
|
|
}
|
|
|
|
const _g = globalThis as CancellationGlobal
|
|
|
|
export function getCancellationChannel(): PubSubChannel<ExecutionCancelEvent> {
|
|
if (!_g._executionCancelChannel) {
|
|
_g._executionCancelChannel = createPubSubChannel<ExecutionCancelEvent>({
|
|
channel: EXECUTION_CANCEL_CHANNEL,
|
|
label: 'execution-cancel',
|
|
})
|
|
}
|
|
return _g._executionCancelChannel
|
|
}
|
|
|
|
export function isRedisCancellationEnabled(): boolean {
|
|
return getRedisClient() !== null
|
|
}
|
|
|
|
/** Writes the durable key first, then publishes — so a late subscriber still sees the flag on backstop check. */
|
|
export async function markExecutionCancelled(
|
|
executionId: string,
|
|
options: MarkExecutionCancelledOptions = {}
|
|
): Promise<ExecutionCancellationRecordResult> {
|
|
const redis = getRedisClient()
|
|
if (!redis) {
|
|
getCancellationChannel().publish({ executionId })
|
|
return { durablyRecorded: false, reason: 'redis_unavailable' }
|
|
}
|
|
|
|
try {
|
|
const minimumExpiryAt = Date.now() + EXECUTION_CANCEL_MIN_RETENTION_MS
|
|
const deadlineExpiryAt = options.executionDeadlineAt?.getTime()
|
|
const expiryAt =
|
|
deadlineExpiryAt !== undefined && Number.isFinite(deadlineExpiryAt)
|
|
? Math.max(minimumExpiryAt, deadlineExpiryAt)
|
|
: minimumExpiryAt
|
|
await redis.set(`${EXECUTION_CANCEL_PREFIX}${executionId}`, '1', 'PXAT', expiryAt)
|
|
logger.info('Marked execution as cancelled', {
|
|
executionId,
|
|
expiresAt: new Date(expiryAt).toISOString(),
|
|
})
|
|
getCancellationChannel().publish({ executionId })
|
|
return { durablyRecorded: true, reason: 'recorded' }
|
|
} catch (error) {
|
|
logger.error('Failed to mark execution as cancelled', { executionId, error })
|
|
getCancellationChannel().publish({ executionId })
|
|
return { durablyRecorded: false, reason: 'redis_write_failed' }
|
|
}
|
|
}
|
|
|
|
export async function isExecutionCancelled(executionId: string): Promise<boolean> {
|
|
const redis = getRedisClient()
|
|
if (!redis) {
|
|
return false
|
|
}
|
|
|
|
try {
|
|
const result = await redis.exists(`${EXECUTION_CANCEL_PREFIX}${executionId}`)
|
|
return result === 1
|
|
} catch (error) {
|
|
logger.error('Failed to check execution cancellation', { executionId, error })
|
|
return false
|
|
}
|
|
}
|
|
|
|
export async function clearExecutionCancellation(executionId: string): Promise<void> {
|
|
const redis = getRedisClient()
|
|
if (!redis) {
|
|
return
|
|
}
|
|
|
|
try {
|
|
await redis.del(`${EXECUTION_CANCEL_PREFIX}${executionId}`)
|
|
} catch (error) {
|
|
logger.error('Failed to clear execution cancellation', { executionId, error })
|
|
}
|
|
}
|