Files
WeHub Mirror 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
WeHub snapshot of cb28d14c6f2c081de7a0d8729a8c816c9adef67a
2026-08-10 11:17:50 +08:00

267 lines
10 KiB
TypeScript

import { createLogger } from '@sim/logger'
import { getErrorMessage, toError } from '@sim/utils/errors'
import { generateId } from '@sim/utils/id'
import { truncate } from '@sim/utils/string'
import type { Filter, RowData, TableDefinition } from '@/lib/table'
import { TABLE_LIMITS, USER_TABLE_ROWS_SQL_NAME } from '@/lib/table/constants'
import { appendTableEvent } from '@/lib/table/events'
import {
getJobProgress,
markJobCanceled,
markJobFailed,
markJobReady,
updateJobProgress,
} from '@/lib/table/jobs/service'
import {
assertRowUpdate,
type MutationProof,
patchColumnIds,
TableLockedError,
} from '@/lib/table/mutation-locks'
import type { DbTransaction } from '@/lib/table/planner'
import { selectRowDataPage, updatePageByIds } from '@/lib/table/rows/ordering'
import { createExactEmptyTableRowSecretProvenance } from '@/lib/table/rows/secret-provenance'
import { getTableById } from '@/lib/table/service'
import { buildFilterClause } from '@/lib/table/sql'
import { coerceRowToSchema, coerceRowValues, validateRowSize } from '@/lib/table/validation'
const logger = createLogger('TableUpdateRunner')
/** Emit a progress event / heartbeat at most every this many rows. */
const PROGRESS_INTERVAL_ROWS = 5000
/**
* Thrown when this worker discovers it no longer owns the table's job (canceled, or the
* stale-job janitor marked it failed and a newer job took over). The worker stops updating.
*/
class JobSupersededError extends Error {}
export interface TableUpdatePayload {
jobId: string
tableId: string
workspaceId: string
/** Rows matching this filter get the patch. */
filter: Filter
/** Column-id-keyed partial patch merged into every matched row. */
data: RowData
/** Only rows created at/before this instant are patched, so mid-job inserts are spared. */
cutoff: Date
/** Stop after updating this many rows (an explicit caller-supplied limit). Omitted = every match. */
maxRows?: number
}
/**
* Background worker for large filtered row updates (trigger.dev task, or detached on the web
* container when trigger.dev is disabled — see the update dispatch in the user_table tool).
* Applies the same `data` patch (JSONB merge) to every row matching `filter` with
* `created_at <= cutoff`, in keyset-paginated pages. Each page validates the merged result per
* row, then commits in batches — **best-effort, not atomic**: committed pages persist even if a
* later page fails validation (unlike the inline `updateRowsByFilter`, which pre-validates all
* rows in one transaction). Reads are not masked: updated rows still exist, so mid-job reads are
* eventually consistent. Ownership-gated per page so a cancel/supersede stops within one page.
*
* Unlike the inline path, the worker does NOT fire per-row table triggers or auto-recompute
* workflow/enrichment columns — that would be a runaway cascade across thousands of rows. Run
* the affected columns explicitly afterward if downstream recompute is needed.
*
* Unexpected errors are rethrown for the caller's retry machinery; the caller marks the job
* failed via `markTableUpdateFailed`. A superseded run returns quietly.
*/
export async function runTableUpdate(payload: TableUpdatePayload): Promise<void> {
const { jobId, tableId, workspaceId, filter, data, cutoff, maxRows } = payload
const requestId = generateId().slice(0, 8)
const budget = maxRows ?? Number.POSITIVE_INFINITY
try {
const table = await getTableById(tableId, { includeArchived: true })
if (!table) throw new Error(`Update target table ${tableId} not found`)
// Gate the run on the update lock, then re-gate it before every page (see
// the loop below), so enabling the lock stops a job that is already
// running. Runs through `assertRowUpdate` rather than reading
// `updateLocked` directly so the enqueue site and the worker apply
// identical rules. This is a user-driven bulk patch, so it deliberately
// does not pass `computedWrite` — the workflow-output carve-out belongs to
// the cell-write path alone.
const cancelForLock = async (processedSoFar: number): Promise<void> => {
logger.info(`[${requestId}] Update job stopped — table is update-locked`, {
tableId,
jobId,
processedSoFar,
})
await markJobCanceled(tableId, jobId)
void appendTableEvent({ kind: 'job', type: 'update', tableId, jobId, status: 'canceled' })
}
const stopIfLocked = async (
fresh: TableDefinition,
processedSoFar: number
): Promise<MutationProof<'update'> | null> => {
try {
return assertRowUpdate(fresh, patchColumnIds(data))
} catch (err) {
if (!(err instanceof TableLockedError)) throw err
await cancelForLock(processedSoFar)
return null
}
}
if ((await stopIfLocked(table, 0)) === null) return
// Runs inside each batch's transaction, under the same advisory lock the
// lock toggle holds, so no page can be written after a lock commits.
const revalidate = async (trx: DbTransaction) => {
const fresh = await getTableById(tableId, { tx: trx, includeArchived: true })
if (fresh) assertRowUpdate(fresh, patchColumnIds(data))
return fresh ?? undefined
}
const filterClause = buildFilterClause(filter, USER_TABLE_ROWS_SQL_NAME, table.schema.columns)
if (!filterClause) throw new Error('Filter is required for bulk update')
// Coerce the patch once to the schema's types — the merged validation below and the persisted
// JSONB merge both use this normalized copy.
coerceRowValues(data, table.schema)
const patchJson = JSON.stringify(data)
// Resume the persisted count: a retried attempt's earlier pages are already committed, so
// starting at zero would overwrite cumulative progress. Doubles as the initial ownership gate.
const resumed = await getJobProgress(tableId, jobId)
if (resumed === null) throw new JobSupersededError()
let processed = resumed
let lastReported = resumed
let afterId: string | undefined
while (processed < budget) {
const owns = await updateJobProgress(tableId, processed, jobId)
if (!owns) throw new JobSupersededError()
// Cheap early-out before selecting a page we may not be allowed to
// write. The authoritative gate is `revalidate` below, which re-asserts
// inside each batch transaction. Pages already applied stay applied, as
// with an explicit cancel.
const current = await getTableById(tableId, { includeArchived: true })
if (!current) throw new JobSupersededError()
const pageProof = await stopIfLocked(current, processed)
if (pageProof === null) return
const page = await selectRowDataPage({
tableId,
workspaceId,
cutoff,
filterClause,
afterId,
limit: Math.min(TABLE_LIMITS.DELETE_PAGE_SIZE, budget - processed),
// Skip rows already carrying the patch so a retried run resumes without re-walking /
// double-counting the rows an earlier attempt updated (updated rows still exist and may
// still match the filter, unlike deletes).
excludeIfPatched: patchJson,
})
if (page.length === 0) break
afterId = page[page.length - 1].id
// Validate each merged result before writing the page — a row that would overflow the size
// cap or violate the schema fails the job (earlier pages stay applied; best-effort).
for (const row of page) {
const merged = { ...row.data, ...data }
const sizeValidation = validateRowSize(merged)
if (!sizeValidation.valid) {
throw new Error(`Row ${row.id}: ${sizeValidation.errors.join(', ')}`)
}
const schemaValidation = coerceRowToSchema(merged, table.schema)
if (!schemaValidation.valid) {
throw new Error(`Row ${row.id}: ${schemaValidation.errors.join(', ')}`)
}
}
try {
processed += await updatePageByIds(
tableId,
workspaceId,
page.map((r) => r.id),
patchJson,
createExactEmptyTableRowSecretProvenance(data),
pageProof,
revalidate
)
} catch (err) {
if (!(err instanceof TableLockedError)) throw err
// A lock landed between batches. Batches already committed stay
// applied; `processed` undercounts them, which only affects the final
// progress number on an already-canceled job.
await cancelForLock(processed)
return
}
if (
processed - lastReported >= PROGRESS_INTERVAL_ROWS ||
(lastReported === 0 && processed > 0)
) {
lastReported = processed
void appendTableEvent({
kind: 'job',
type: 'update',
tableId,
jobId,
status: 'running',
progress: processed,
})
}
}
await updateJobProgress(tableId, processed, jobId)
const becameReady = await markJobReady(tableId, jobId)
if (becameReady) {
void appendTableEvent({
kind: 'job',
type: 'update',
tableId,
jobId,
status: 'ready',
progress: processed,
})
logger.info(`[${requestId}] Update complete`, { tableId, rows: processed })
} else {
logger.info(
`[${requestId}] Update finished but no longer owns the run (canceled/superseded)`,
{
tableId,
jobId,
}
)
}
} catch (err) {
if (err instanceof JobSupersededError) {
logger.info(`[${requestId}] Update superseded by a newer run; stopping`, { tableId, jobId })
return
}
const cause = toError(err).cause
const error = cause ? toError(cause) : toError(err)
logger.error(`[${requestId}] Update failed for table ${tableId}:`, error)
throw error
}
}
/**
* Marks the update job failed and emits the failed SSE event. Called once the caller gives up on
* the run (trigger.dev `onFailure` after retries, or the detached fallback). Scoped to jobId — a
* no-op if a newer job has taken over.
*/
export async function markTableUpdateFailed(
tableId: string,
jobId: string,
error: unknown
): Promise<void> {
const message = truncate(getErrorMessage(toError(error).cause ?? error, 'Update failed'), 500)
await markJobFailed(tableId, jobId, message).catch(() => {})
void appendTableEvent({
kind: 'job',
type: 'update',
tableId,
jobId,
status: 'failed',
error: message,
})
}