mirror of
https://github.com/simstudioai/sim.git
synced 2026-09-24 15:45:35 +08:00
improvement(mothership): bounded delete above the cap runs async, not inline
An explicit delete limit now mirrors update: ≤1000 runs inline, above the cap it escalates to the background worker honoring the limit via maxRows — instead of always staying inline. The worker stops after maxRows (per-page fetch capped to the remaining budget). Bounded background deletes skip pendingDeleteMask: the filter-based mask hides every match, which would over-hide the rows beyond the cap the job never deletes. Unmasked, a bounded delete is eventually consistent like a bounded update (rows disappear as deleted), and doomedCount is omitted from the payload so the count isn't double-subtracted. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Fable 5
parent
d2c93aef04
commit
f1ee3e9068
@@ -732,7 +732,15 @@ describe('userTableServerTool.delete_rows_by_filter', () => {
|
|||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
it('runs an explicit large limit inline without escalating (delete loads only ids)', async () => {
|
it('escalates an explicit limit above the cap to a background delete with maxRows (unmasked)', async () => {
|
||||||
|
mockQueryRows.mockResolvedValueOnce({
|
||||||
|
rows: [],
|
||||||
|
rowCount: 0,
|
||||||
|
totalCount: 20000,
|
||||||
|
limit: 1,
|
||||||
|
offset: 0,
|
||||||
|
})
|
||||||
|
|
||||||
const result = await userTableServerTool.execute(
|
const result = await userTableServerTool.execute(
|
||||||
{
|
{
|
||||||
operation: 'delete_rows_by_filter',
|
operation: 'delete_rows_by_filter',
|
||||||
@@ -740,12 +748,19 @@ describe('userTableServerTool.delete_rows_by_filter', () => {
|
|||||||
},
|
},
|
||||||
{ userId: 'user-1', workspaceId: 'workspace-1' }
|
{ userId: 'user-1', workspaceId: 'workspace-1' }
|
||||||
)
|
)
|
||||||
|
await flushDetached()
|
||||||
|
|
||||||
expect(result.success).toBe(true)
|
expect(result.success).toBe(true)
|
||||||
// An explicit limit never counts/escalates — it deletes inline, bounded by the limit.
|
// target = min(limit 5000, matchCount 20000) = 5000, above the inline cap → background.
|
||||||
expect(mockQueryRows).not.toHaveBeenCalled()
|
expect(result.data?.doomedCount).toBe(5000)
|
||||||
expect(mockDeleteRowsByFilter).toHaveBeenCalledTimes(1)
|
expect(mockDeleteRowsByFilter).not.toHaveBeenCalled()
|
||||||
expect(mockDeleteRowsByFilter.mock.calls[0][1]).toMatchObject({ limit: 5000 })
|
const [, , type, payload] = mockMarkTableJobRunning.mock.calls[0]
|
||||||
|
expect(type).toBe('delete')
|
||||||
|
// Bounded delete carries maxRows and omits doomedCount so the mask is skipped and the count
|
||||||
|
// isn't double-subtracted.
|
||||||
|
expect(payload).toMatchObject({ maxRows: 5000 })
|
||||||
|
expect((payload as { doomedCount?: number }).doomedCount).toBeUndefined()
|
||||||
|
expect(mockRunTableDelete.mock.calls[0][0]).toMatchObject({ maxRows: 5000 })
|
||||||
})
|
})
|
||||||
|
|
||||||
it('deletes inline when the unbounded match count is within the cap', async () => {
|
it('deletes inline when the unbounded match count is within the cap', async () => {
|
||||||
@@ -801,6 +816,8 @@ describe('userTableServerTool.delete_rows_by_filter', () => {
|
|||||||
expect(tableId).toBe('tbl_1')
|
expect(tableId).toBe('tbl_1')
|
||||||
expect(type).toBe('delete')
|
expect(type).toBe('delete')
|
||||||
expect(payload).toMatchObject({ doomedCount: 20000, cutoff: expect.any(String) })
|
expect(payload).toMatchObject({ doomedCount: 20000, cutoff: expect.any(String) })
|
||||||
|
// Unbounded delete masks the whole set — no maxRows cap.
|
||||||
|
expect((payload as { maxRows?: number }).maxRows).toBeUndefined()
|
||||||
expect(mockRunTableDelete).toHaveBeenCalledTimes(1)
|
expect(mockRunTableDelete).toHaveBeenCalledTimes(1)
|
||||||
expect(mockRunTableDelete.mock.calls[0][0]).toMatchObject({
|
expect(mockRunTableDelete.mock.calls[0][0]).toMatchObject({
|
||||||
jobId,
|
jobId,
|
||||||
|
|||||||
@@ -161,8 +161,9 @@ async function dispatchDeleteJob(params: {
|
|||||||
workspaceId: string
|
workspaceId: string
|
||||||
filter: Filter
|
filter: Filter
|
||||||
cutoff: Date
|
cutoff: Date
|
||||||
|
maxRows?: number
|
||||||
}): Promise<void> {
|
}): Promise<void> {
|
||||||
const { jobId, tableId, workspaceId, filter, cutoff } = params
|
const { jobId, tableId, workspaceId, filter, cutoff, maxRows } = params
|
||||||
if (isTriggerDevEnabled) {
|
if (isTriggerDevEnabled) {
|
||||||
try {
|
try {
|
||||||
const [{ tableDeleteTask }, { tasks }] = await Promise.all([
|
const [{ tableDeleteTask }, { tasks }] = await Promise.all([
|
||||||
@@ -171,7 +172,7 @@ async function dispatchDeleteJob(params: {
|
|||||||
])
|
])
|
||||||
await tasks.trigger<typeof tableDeleteTask>(
|
await tasks.trigger<typeof tableDeleteTask>(
|
||||||
'table-delete',
|
'table-delete',
|
||||||
{ jobId, tableId, workspaceId, filter, cutoff: cutoff.toISOString() },
|
{ jobId, tableId, workspaceId, filter, cutoff: cutoff.toISOString(), maxRows },
|
||||||
{ tags: [`tableId:${tableId}`, `jobId:${jobId}`] }
|
{ tags: [`tableId:${tableId}`, `jobId:${jobId}`] }
|
||||||
)
|
)
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
@@ -180,10 +181,12 @@ async function dispatchDeleteJob(params: {
|
|||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
runDetached('table-delete', () =>
|
runDetached('table-delete', () =>
|
||||||
runTableDelete({ jobId, tableId, workspaceId, filter, cutoff }).catch(async (error) => {
|
runTableDelete({ jobId, tableId, workspaceId, filter, cutoff, maxRows }).catch(
|
||||||
|
async (error) => {
|
||||||
await markTableDeleteFailed(tableId, jobId, error)
|
await markTableDeleteFailed(tableId, jobId, error)
|
||||||
throw error
|
throw error
|
||||||
})
|
}
|
||||||
|
)
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -817,27 +820,31 @@ export const userTableServerTool: BaseServerTool<UserTableArgs, UserTableResult>
|
|||||||
const idByName = buildIdByName(table.schema)
|
const idByName = buildIdByName(table.schema)
|
||||||
const idFilter = filterNamesToIds(args.filter, idByName)
|
const idFilter = filterNamesToIds(args.filter, idByName)
|
||||||
|
|
||||||
// An explicit limit runs inline (delete loads only row ids, so even a large bounded
|
// Inline handles up to MAX_BULK_OPERATION_SIZE rows; a larger delete (an explicit limit
|
||||||
// delete is light). Only an unbounded "delete everything matching" measures the blast
|
// above the cap, or unbounded "delete everything matching") hands off to the background
|
||||||
// radius and hands off to the background delete worker (same path as the UI's select-all
|
// delete worker so a broad delete on a huge table doesn't load every matching id into this
|
||||||
// delete) — the read-path mask hides exactly the all-matching set, which a bounded delete
|
// request. A small explicit limit is the fast path.
|
||||||
// would over-hide.
|
const deleteInlineEligible =
|
||||||
if (args.limit === undefined) {
|
args.limit !== undefined && args.limit <= TABLE_LIMITS.MAX_BULK_OPERATION_SIZE
|
||||||
|
if (!deleteInlineEligible) {
|
||||||
const { totalCount } = await queryRows(
|
const { totalCount } = await queryRows(
|
||||||
table,
|
table,
|
||||||
{ filter: idFilter, limit: 1, withExecutions: false },
|
{ filter: idFilter, limit: 1, withExecutions: false },
|
||||||
requestId
|
requestId
|
||||||
)
|
)
|
||||||
const matchCount = totalCount ?? 0
|
const matchCount = totalCount ?? 0
|
||||||
if (matchCount > TABLE_LIMITS.MAX_BULK_OPERATION_SIZE) {
|
const target = args.limit !== undefined ? Math.min(args.limit, matchCount) : matchCount
|
||||||
const doomedCount = Math.min(matchCount, table.rowCount)
|
if (target > TABLE_LIMITS.MAX_BULK_OPERATION_SIZE) {
|
||||||
|
const doomedCount = Math.min(target, table.rowCount)
|
||||||
const cutoff = new Date()
|
const cutoff = new Date()
|
||||||
const jobId = generateId()
|
const jobId = generateId()
|
||||||
const payload: TableDeleteJobPayload = {
|
// Unbounded: mask the whole matching set (instant post-delete view), so `doomedCount`
|
||||||
filter: idFilter,
|
// drives the count adjustment. Bounded (maxRows): no mask — `doomedCount` is omitted so
|
||||||
cutoff: cutoff.toISOString(),
|
// the count isn't double-subtracted; rows disappear progressively as they're deleted.
|
||||||
doomedCount,
|
const bounded = args.limit !== undefined
|
||||||
}
|
const payload: TableDeleteJobPayload = bounded
|
||||||
|
? { filter: idFilter, cutoff: cutoff.toISOString(), maxRows: args.limit }
|
||||||
|
: { filter: idFilter, cutoff: cutoff.toISOString(), doomedCount }
|
||||||
assertNotAborted()
|
assertNotAborted()
|
||||||
const claimed = await markTableJobRunning(table.id, jobId, 'delete', payload)
|
const claimed = await markTableJobRunning(table.id, jobId, 'delete', payload)
|
||||||
if (!claimed) {
|
if (!claimed) {
|
||||||
@@ -849,10 +856,13 @@ export const userTableServerTool: BaseServerTool<UserTableArgs, UserTableResult>
|
|||||||
workspaceId,
|
workspaceId,
|
||||||
filter: idFilter,
|
filter: idFilter,
|
||||||
cutoff,
|
cutoff,
|
||||||
|
maxRows: args.limit,
|
||||||
})
|
})
|
||||||
return {
|
return {
|
||||||
success: true,
|
success: true,
|
||||||
message: `Started background delete of ${doomedCount} matching rows (job ${jobId}). The rows are hidden from reads immediately — query_rows already reflects the post-delete view.`,
|
message: bounded
|
||||||
|
? `Started background delete of up to ${doomedCount} matching rows (job ${jobId}). Rows delete in the background — query_rows to check progress.`
|
||||||
|
: `Started background delete of ${doomedCount} matching rows (job ${jobId}). The rows are hidden from reads immediately — query_rows already reflects the post-delete view.`,
|
||||||
data: { jobId, doomedCount },
|
data: { jobId, doomedCount },
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -82,6 +82,22 @@ describe('runTableDelete', () => {
|
|||||||
)
|
)
|
||||||
})
|
})
|
||||||
|
|
||||||
|
it('stops once maxRows is reached and caps the final page fetch to the remaining budget', async () => {
|
||||||
|
// budget 3 with page size 2: first page fills 2, the second is capped to the remaining 1.
|
||||||
|
mockSelectRowIdPage.mockResolvedValueOnce(['a', 'b']).mockResolvedValueOnce(['c'])
|
||||||
|
|
||||||
|
await runTableDelete(basePayload({ filter: { status: 'old' }, maxRows: 3 }))
|
||||||
|
|
||||||
|
expect(mockSelectRowIdPage).toHaveBeenCalledTimes(2)
|
||||||
|
expect(mockSelectRowIdPage.mock.calls[0][0]).toMatchObject({ limit: 2 })
|
||||||
|
expect(mockSelectRowIdPage.mock.calls[1][0]).toMatchObject({ limit: 1 })
|
||||||
|
expect(mockDeletePageByIds).toHaveBeenCalledTimes(2)
|
||||||
|
expect(mockMarkJobReady).toHaveBeenCalledWith('tbl_1', 'job_1')
|
||||||
|
expect(mockAppendTableEvent).toHaveBeenCalledWith(
|
||||||
|
expect.objectContaining({ status: 'ready', progress: 3 })
|
||||||
|
)
|
||||||
|
})
|
||||||
|
|
||||||
it('skips excluded rows but still advances the keyset cursor past them', async () => {
|
it('skips excluded rows but still advances the keyset cursor past them', async () => {
|
||||||
mockSelectRowIdPage.mockResolvedValueOnce(['keep', 'x']).mockResolvedValueOnce([])
|
mockSelectRowIdPage.mockResolvedValueOnce(['keep', 'x']).mockResolvedValueOnce([])
|
||||||
|
|
||||||
|
|||||||
@@ -36,6 +36,12 @@ export interface TableDeletePayload {
|
|||||||
excludeRowIds?: string[]
|
excludeRowIds?: string[]
|
||||||
/** Only rows created at/before this instant are deleted, so mid-job inserts survive. */
|
/** Only rows created at/before this instant are deleted, so mid-job inserts survive. */
|
||||||
cutoff: Date
|
cutoff: Date
|
||||||
|
/**
|
||||||
|
* Stop after deleting this many rows (an explicit caller-supplied limit). Omitted = every match.
|
||||||
|
* Not combined with `excludeRowIds` (the UI's select-all path uses excludes and no cap; the
|
||||||
|
* copilot tool uses a cap and no excludes), so the per-page fetch can be bounded directly.
|
||||||
|
*/
|
||||||
|
maxRows?: number
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -52,8 +58,9 @@ export interface TableDeletePayload {
|
|||||||
* newer job took the table) returns quietly.
|
* newer job took the table) returns quietly.
|
||||||
*/
|
*/
|
||||||
export async function runTableDelete(payload: TableDeletePayload): Promise<void> {
|
export async function runTableDelete(payload: TableDeletePayload): Promise<void> {
|
||||||
const { jobId, tableId, workspaceId, filter, excludeRowIds, cutoff } = payload
|
const { jobId, tableId, workspaceId, filter, excludeRowIds, cutoff, maxRows } = payload
|
||||||
const requestId = generateId().slice(0, 8)
|
const requestId = generateId().slice(0, 8)
|
||||||
|
const budget = maxRows ?? Number.POSITIVE_INFINITY
|
||||||
|
|
||||||
try {
|
try {
|
||||||
const table = await getTableById(tableId, { includeArchived: true })
|
const table = await getTableById(tableId, { includeArchived: true })
|
||||||
@@ -74,7 +81,7 @@ export async function runTableDelete(payload: TableDeletePayload): Promise<void>
|
|||||||
let lastReported = resumed
|
let lastReported = resumed
|
||||||
let afterId: string | undefined
|
let afterId: string | undefined
|
||||||
|
|
||||||
while (true) {
|
while (processed < budget) {
|
||||||
// Ownership gate before every page: once this run loses the table (cancel/supersede),
|
// Ownership gate before every page: once this run loses the table (cancel/supersede),
|
||||||
// updateJobProgress returns false and we stop before deleting further.
|
// updateJobProgress returns false and we stop before deleting further.
|
||||||
const owns = await updateJobProgress(tableId, processed, jobId)
|
const owns = await updateJobProgress(tableId, processed, jobId)
|
||||||
@@ -86,7 +93,7 @@ export async function runTableDelete(payload: TableDeletePayload): Promise<void>
|
|||||||
cutoff,
|
cutoff,
|
||||||
filterClause,
|
filterClause,
|
||||||
afterId,
|
afterId,
|
||||||
limit: TABLE_LIMITS.DELETE_PAGE_SIZE,
|
limit: Math.min(TABLE_LIMITS.DELETE_PAGE_SIZE, budget - processed),
|
||||||
})
|
})
|
||||||
if (page.length === 0) break
|
if (page.length === 0) break
|
||||||
// Advance the keyset cursor past the whole page — excluded ids are skipped (not deleted),
|
// Advance the keyset cursor past the whole page — excluded ids are skipped (not deleted),
|
||||||
|
|||||||
@@ -213,8 +213,16 @@ export interface TableDeleteJobPayload {
|
|||||||
/** ISO timestamp; rows created after it are spared. */
|
/** ISO timestamp; rows created after it are spared. */
|
||||||
cutoff: string
|
cutoff: string
|
||||||
/** Doomed-row estimate captured at kickoff — display-only: list/detail counts subtract the
|
/** Doomed-row estimate captured at kickoff — display-only: list/detail counts subtract the
|
||||||
* not-yet-deleted remainder (doomedCount - rows_processed) while the job runs. */
|
* not-yet-deleted remainder (doomedCount - rows_processed) while the job runs. Set only for an
|
||||||
|
* unbounded delete (the masked "delete everything matching" path); omitted when `maxRows` is set. */
|
||||||
doomedCount?: number
|
doomedCount?: number
|
||||||
|
/**
|
||||||
|
* Stop after deleting this many rows (an explicit caller-supplied limit above the inline cap).
|
||||||
|
* Omitted = delete every match. When set, reads are NOT masked: the delete is eventually
|
||||||
|
* consistent (rows disappear as they're deleted) like a bounded update, because the filter-based
|
||||||
|
* mask would over-hide the rows beyond the cap that this job never deletes.
|
||||||
|
*/
|
||||||
|
maxRows?: number
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
Reference in New Issue
Block a user