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
310 lines
11 KiB
TypeScript
310 lines
11 KiB
TypeScript
import { db } from '@sim/db'
|
||
import { createLogger } from '@sim/logger'
|
||
import { chunkArray } from '@sim/utils/helpers'
|
||
import { and, inArray, isNotNull, lt, type SQL, sql } from 'drizzle-orm'
|
||
import type { PgColumn, PgTable } from 'drizzle-orm/pg-core'
|
||
|
||
const logger = createLogger('BatchDelete')
|
||
|
||
/**
|
||
* Structural client surface the delete helpers need. Satisfied by the global
|
||
* `db`, a `dbFor(...)` sub-pool client, and a transaction handle, so callers
|
||
* pick which pool the deletes run on (cleanup jobs pass `dbFor('cleanup')`).
|
||
*/
|
||
export type BatchDeleteClient = Pick<typeof db, 'select' | 'delete'>
|
||
|
||
export const DEFAULT_BATCH_SIZE = 2000
|
||
/** 50 × 2000 = 100K row cap per cleanup run; drains long-tail tenants in days, not weeks. */
|
||
export const DEFAULT_MAX_BATCHES_PER_TABLE = 50
|
||
/**
|
||
* Split workspaceIds into this-sized groups before running SELECT/DELETE. Large
|
||
* IN lists combined with `started_at < X` force Postgres to probe every
|
||
* workspace range in the composite index, which blows the 90s statement timeout
|
||
* at the scale of the full free tier.
|
||
*/
|
||
export const DEFAULT_WORKSPACE_CHUNK_SIZE = 50
|
||
/** Bounds FK cascade trigger queue (per-statement in-memory) and bind-parameter count. */
|
||
export const DEFAULT_DELETE_CHUNK_SIZE = 1000
|
||
|
||
export interface SelectByIdChunksOptions {
|
||
/** Cap on rows returned across all chunks. Defaults to a full per-table cleanup budget. */
|
||
overallLimit?: number
|
||
chunkSize?: number
|
||
}
|
||
|
||
/**
|
||
* Run a SELECT query once per ID chunk and concatenate results up to
|
||
* `overallLimit`. Each chunk's query is passed the remaining row budget so the
|
||
* total never exceeds the cap. Use this when you need the selected row set
|
||
* (e.g. to drive S3 or copilot-backend cleanup alongside the DB delete).
|
||
*
|
||
* Works for any large ID set — workspace IDs, workflow IDs, etc. Avoids
|
||
* sending one massive `IN (...)` list that would blow Postgres's statement
|
||
* timeout.
|
||
*/
|
||
export async function selectRowsByIdChunks<T>(
|
||
ids: string[],
|
||
query: (chunkIds: string[], chunkLimit: number) => Promise<T[]>,
|
||
{
|
||
overallLimit = DEFAULT_BATCH_SIZE * DEFAULT_MAX_BATCHES_PER_TABLE,
|
||
chunkSize = DEFAULT_WORKSPACE_CHUNK_SIZE,
|
||
}: SelectByIdChunksOptions = {}
|
||
): Promise<T[]> {
|
||
if (ids.length === 0) return []
|
||
|
||
const rows: T[] = []
|
||
for (const chunkIds of chunkArray(ids, chunkSize)) {
|
||
if (rows.length >= overallLimit) break
|
||
const remaining = overallLimit - rows.length
|
||
const chunkRows = await query(chunkIds, remaining)
|
||
rows.push(...chunkRows)
|
||
}
|
||
return rows
|
||
}
|
||
|
||
export interface TableCleanupResult {
|
||
table: string
|
||
deleted: number
|
||
failed: number
|
||
}
|
||
|
||
export interface ChunkedBatchDeleteOptions<TRow extends { id: string }> {
|
||
tableDef: PgTable
|
||
workspaceIds: string[]
|
||
tableName: string
|
||
/** SELECT eligible rows for one workspace chunk. The result must include `id`. */
|
||
selectChunk: (chunkIds: string[], limit: number) => Promise<TRow[]>
|
||
/** Runs between SELECT and DELETE; receives the just-selected rows. */
|
||
onBatch?: (rows: TRow[]) => Promise<void>
|
||
/**
|
||
* Re-asserted on the DELETE alongside the id list. A soft-delete sweep whose `onBatch` does
|
||
* real work before the DELETE should pass the same predicate it selected on: a row restored
|
||
* in that window would otherwise be hard-deleted — taking the children a folder restore had
|
||
* just brought back with it. Rows that no longer match are simply not deleted and are
|
||
* counted as failed, so the next run re-evaluates them.
|
||
*
|
||
* Optional because a sweep with no `onBatch` closes a far smaller window; the hand-rolled
|
||
* targets in `cleanup-soft-deletes.ts` re-check inside their own DELETE instead.
|
||
*/
|
||
deleteFilter?: SQL
|
||
batchSize?: number
|
||
/** Max batches per workspace chunk. */
|
||
maxBatches?: number
|
||
/**
|
||
* Hard cap on rows processed (deleted + failed) across all chunks per call.
|
||
* Defaults to `DEFAULT_BATCH_SIZE * DEFAULT_MAX_BATCHES_PER_TABLE`. Cron
|
||
* runs frequently enough to catch up the backlog over multiple invocations.
|
||
*/
|
||
totalRowLimit?: number
|
||
workspaceChunkSize?: number
|
||
/** Client the DELETEs run on. Defaults to the global pool. */
|
||
dbClient?: BatchDeleteClient
|
||
}
|
||
|
||
/**
|
||
* Inner loop primitive for cleanup jobs.
|
||
*
|
||
* For each workspace chunk: SELECT a batch of eligible rows → run optional
|
||
* `onBatch` hook (e.g. to delete S3 files) → DELETE those rows by ID. Repeats
|
||
* until exhausted or `maxBatches` is hit, then moves to the next chunk. Stops
|
||
* the whole call once `totalRowLimit` rows have been processed.
|
||
*
|
||
* Workspace IDs are chunked before the SELECT — see
|
||
* `DEFAULT_WORKSPACE_CHUNK_SIZE` for why.
|
||
*/
|
||
export async function chunkedBatchDelete<TRow extends { id: string }>({
|
||
tableDef,
|
||
workspaceIds,
|
||
tableName,
|
||
selectChunk,
|
||
onBatch,
|
||
deleteFilter,
|
||
batchSize = DEFAULT_BATCH_SIZE,
|
||
maxBatches = DEFAULT_MAX_BATCHES_PER_TABLE,
|
||
totalRowLimit = DEFAULT_BATCH_SIZE * DEFAULT_MAX_BATCHES_PER_TABLE,
|
||
workspaceChunkSize = DEFAULT_WORKSPACE_CHUNK_SIZE,
|
||
dbClient = db,
|
||
}: ChunkedBatchDeleteOptions<TRow>): Promise<TableCleanupResult> {
|
||
const result: TableCleanupResult = { table: tableName, deleted: 0, failed: 0 }
|
||
|
||
if (workspaceIds.length === 0) {
|
||
logger.info(`[${tableName}] Skipped — no workspaces in scope`)
|
||
return result
|
||
}
|
||
|
||
const chunks = chunkArray(workspaceIds, workspaceChunkSize)
|
||
let stoppedEarly = false
|
||
let attempted = 0
|
||
|
||
for (const [chunkIdx, chunkIds] of chunks.entries()) {
|
||
if (attempted >= totalRowLimit) {
|
||
stoppedEarly = true
|
||
break
|
||
}
|
||
|
||
let batchesProcessed = 0
|
||
let hasMore = true
|
||
|
||
while (hasMore && batchesProcessed < maxBatches && attempted < totalRowLimit) {
|
||
let rows: TRow[] = []
|
||
try {
|
||
const remainingLimit = totalRowLimit - attempted
|
||
const effectiveBatchSize = Math.min(batchSize, remainingLimit)
|
||
if (effectiveBatchSize <= 0) {
|
||
hasMore = false
|
||
break
|
||
}
|
||
|
||
rows = await selectChunk(chunkIds, effectiveBatchSize)
|
||
|
||
if (rows.length === 0) {
|
||
hasMore = false
|
||
break
|
||
}
|
||
|
||
attempted += rows.length
|
||
if (onBatch) await onBatch(rows)
|
||
|
||
const ids = rows.map((r) => r.id)
|
||
const deleted = await dbClient
|
||
.delete(tableDef)
|
||
.where(deleteFilter ? and(inArray(sql`id`, ids), deleteFilter) : inArray(sql`id`, ids))
|
||
.returning({ id: sql`id` })
|
||
|
||
result.deleted += deleted.length
|
||
result.failed += rows.length - deleted.length
|
||
hasMore = rows.length === effectiveBatchSize && attempted < totalRowLimit
|
||
batchesProcessed++
|
||
} catch (error) {
|
||
// Count rows we tried to delete; SELECT-stage errors leave rows=[].
|
||
result.failed += rows.length
|
||
logger.error(
|
||
`[${tableName}] Batch failed (chunk ${chunkIdx + 1}/${chunks.length}, ${rows.length} rows):`,
|
||
{ error }
|
||
)
|
||
hasMore = false
|
||
}
|
||
}
|
||
}
|
||
|
||
logger.info(
|
||
`[${tableName}] Complete: ${result.deleted} deleted, ${result.failed} failed across ${chunks.length} chunks${stoppedEarly ? ' (row-limit reached, remaining chunks deferred to next run)' : ''}`
|
||
)
|
||
|
||
return result
|
||
}
|
||
|
||
export interface BatchDeleteOptions {
|
||
tableDef: PgTable
|
||
workspaceIdCol: PgColumn
|
||
timestampCol: PgColumn
|
||
workspaceIds: string[]
|
||
retentionDate: Date
|
||
tableName: string
|
||
/** When true, also requires `timestampCol IS NOT NULL` (soft-delete semantics). */
|
||
requireTimestampNotNull?: boolean
|
||
/**
|
||
* Extra predicate ANDed into the row selection. Needed for tables shared by several
|
||
* resource kinds (e.g. `folder`, which holds workflow/file/knowledge_base/table rows)
|
||
* so a cleanup pass only ever removes the kind it owns.
|
||
*/
|
||
additionalPredicate?: SQL
|
||
/**
|
||
* Runs on each selected batch before its DELETE, for side effects that must observe exactly
|
||
* the rows about to be removed. Forwarded to `chunkedBatchDelete`; see `deleteFilter` there
|
||
* for the restore-race window this opens.
|
||
*/
|
||
onBatch?: (rows: { id: string }[]) => Promise<void>
|
||
batchSize?: number
|
||
maxBatches?: number
|
||
workspaceChunkSize?: number
|
||
/** Client the SELECTs and DELETEs run on. Defaults to the global pool. */
|
||
dbClient?: BatchDeleteClient
|
||
}
|
||
|
||
/**
|
||
* Convenience wrapper around `chunkedBatchDelete` for the common case: delete
|
||
* rows where `workspaceId IN (...) AND timestamp < retentionDate`. Use this
|
||
* when there's no per-row side effect (e.g. no S3 files to clean up alongside).
|
||
*/
|
||
export async function batchDeleteByWorkspaceAndTimestamp({
|
||
tableDef,
|
||
workspaceIdCol,
|
||
timestampCol,
|
||
workspaceIds,
|
||
retentionDate,
|
||
tableName,
|
||
requireTimestampNotNull = false,
|
||
additionalPredicate,
|
||
dbClient = db,
|
||
...rest
|
||
}: BatchDeleteOptions): Promise<TableCleanupResult> {
|
||
/**
|
||
* Re-asserted on the DELETE, not just the SELECT. Every row here is soft-deleted and past
|
||
* retention, so a restore committing between the two statements is exactly the case that must
|
||
* not be hard-deleted — and for `folder` that would take the placement of the children the
|
||
* restore had just brought back with it. Rebuilt rather than reused from `selectChunk` because
|
||
* the id list already scopes the statement; only the eligibility half is re-checked.
|
||
*/
|
||
const eligibility = [lt(timestampCol, retentionDate)]
|
||
if (requireTimestampNotNull) eligibility.push(isNotNull(timestampCol))
|
||
if (additionalPredicate) eligibility.push(additionalPredicate)
|
||
|
||
return chunkedBatchDelete({
|
||
tableDef,
|
||
workspaceIds,
|
||
tableName,
|
||
dbClient,
|
||
deleteFilter: and(...eligibility),
|
||
selectChunk: (chunkIds, limit) => {
|
||
return dbClient
|
||
.select({ id: sql<string>`id` })
|
||
.from(tableDef)
|
||
.where(and(inArray(workspaceIdCol, chunkIds), ...eligibility))
|
||
.limit(limit)
|
||
},
|
||
...rest,
|
||
})
|
||
}
|
||
|
||
/**
|
||
* Delete by explicit ID list, chunked so each statement is its own transaction.
|
||
* Partial progress survives chunk-level failures.
|
||
*/
|
||
export async function deleteRowsById(
|
||
tableDef: PgTable,
|
||
idCol: PgColumn,
|
||
ids: string[],
|
||
tableName: string,
|
||
dbClient: BatchDeleteClient = db,
|
||
chunkSize: number = DEFAULT_DELETE_CHUNK_SIZE
|
||
): Promise<TableCleanupResult> {
|
||
const result: TableCleanupResult = { table: tableName, deleted: 0, failed: 0 }
|
||
if (ids.length === 0) return result
|
||
|
||
const chunks = chunkArray(ids, chunkSize)
|
||
for (const [chunkIdx, chunkIds] of chunks.entries()) {
|
||
try {
|
||
const deleted = await dbClient
|
||
.delete(tableDef)
|
||
.where(inArray(idCol, chunkIds))
|
||
.returning({ id: idCol })
|
||
result.deleted += deleted.length
|
||
} catch (error) {
|
||
// Upper bound: Postgres rolls back the chunk on error, so actual deletes = 0,
|
||
// but we can't tell which IDs in the chunk would have matched. The next cron
|
||
// run picks up whatever's still expired, so this only inflates the metric.
|
||
result.failed += chunkIds.length
|
||
logger.error(
|
||
`[${tableName}] Delete chunk ${chunkIdx + 1}/${chunks.length} failed (up to ${chunkIds.length} rows):`,
|
||
{ error }
|
||
)
|
||
}
|
||
}
|
||
|
||
logger.info(
|
||
`[${tableName}] Deleted ${result.deleted} rows across ${chunks.length} chunk(s)${result.failed > 0 ? `, ${result.failed} failed` : ''}`
|
||
)
|
||
return result
|
||
}
|