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
99 lines
2.3 KiB
TypeScript
99 lines
2.3 KiB
TypeScript
import { createLogger } from '@sim/logger'
|
|
import { generateId } from '@sim/utils/id'
|
|
import { client } from '@/lib/auth/auth-client'
|
|
import { useOperationQueueStore } from '@/stores/operation-queue/store'
|
|
import type { WorkflowState } from '@/stores/workflows/workflow/types'
|
|
import { normalizeWorkflowState } from '@/stores/workflows/workflow/validation'
|
|
|
|
const logger = createLogger('WorkflowSocketOperations')
|
|
|
|
async function resolveUserId(): Promise<string> {
|
|
try {
|
|
const sessionResult = await client.getSession()
|
|
const userId = sessionResult.data?.user?.id
|
|
if (userId) {
|
|
return userId
|
|
}
|
|
} catch (error) {
|
|
logger.warn('Failed to resolve session user id for workflow operation', { error })
|
|
}
|
|
|
|
return 'unknown'
|
|
}
|
|
|
|
interface EnqueueWorkflowOperationArgs {
|
|
operation: string
|
|
target: string
|
|
payload: any
|
|
workflowId: string
|
|
operationId?: string
|
|
}
|
|
|
|
/**
|
|
* Queues a workflow socket operation so it flows through the standard operation queue,
|
|
* ensuring consistent retries, confirmations, and telemetry.
|
|
*/
|
|
async function enqueueWorkflowOperation({
|
|
operation,
|
|
target,
|
|
payload,
|
|
workflowId,
|
|
operationId,
|
|
}: EnqueueWorkflowOperationArgs): Promise<string> {
|
|
const userId = await resolveUserId()
|
|
const opId = operationId ?? generateId()
|
|
|
|
useOperationQueueStore.getState().addToQueue({
|
|
id: opId,
|
|
operation: {
|
|
operation,
|
|
target,
|
|
payload,
|
|
},
|
|
workflowId,
|
|
userId,
|
|
})
|
|
|
|
logger.debug('Queued workflow operation', {
|
|
workflowId,
|
|
operation,
|
|
target,
|
|
operationId: opId,
|
|
})
|
|
|
|
return opId
|
|
}
|
|
|
|
interface EnqueueReplaceStateArgs {
|
|
workflowId: string
|
|
state: WorkflowState
|
|
operationId?: string
|
|
}
|
|
|
|
/**
|
|
* Convenience wrapper for broadcasting a full workflow state replacement via the queue.
|
|
*/
|
|
export async function enqueueReplaceWorkflowState({
|
|
workflowId,
|
|
state,
|
|
operationId,
|
|
}: EnqueueReplaceStateArgs): Promise<string> {
|
|
const { state: validatedState, warnings } = normalizeWorkflowState(state)
|
|
|
|
if (warnings.length > 0) {
|
|
logger.warn('Normalized state before enqueuing replace-state', {
|
|
workflowId,
|
|
warningCount: warnings.length,
|
|
warnings,
|
|
})
|
|
}
|
|
|
|
return enqueueWorkflowOperation({
|
|
workflowId,
|
|
operation: 'replace-state',
|
|
target: 'workflow',
|
|
payload: { state: validatedState },
|
|
operationId,
|
|
})
|
|
}
|