mirror of
https://github.com/simstudioai/sim.git
synced 2026-09-24 15:45:35 +08:00
fix(table): preserve workflow groups on CSV column-add and dispatch after tx commit (#4503)
Two bugs in the CSV-import path: - addTableColumnsWithTx rebuilt the schema with only `columns`, dropping `workflowGroups` (and any other top-level schema fields). Importing CSV into a table that has workflow groups erased the group config. Spread `table.schema` first so siblings survive. - batchInsertRowsWithTx fired fireTableTrigger and scheduleRunsForRows from inside the caller's transaction. Both read through the global db connection, so they could run before the inserts committed and see no rows. Extracted the dispatch into dispatchAfterBatchInsert; non-tx wrapper fires it after `db.transaction(...)` resolves, and the CSV import route does the same after its tx.
This commit is contained in:
@@ -11,11 +11,13 @@ const {
|
||||
mockBatchInsertRowsWithTx,
|
||||
mockReplaceTableRowsWithTx,
|
||||
mockAddTableColumnsWithTx,
|
||||
mockDispatchAfterBatchInsert,
|
||||
} = vi.hoisted(() => ({
|
||||
mockCheckAccess: vi.fn(),
|
||||
mockBatchInsertRowsWithTx: vi.fn(),
|
||||
mockReplaceTableRowsWithTx: vi.fn(),
|
||||
mockAddTableColumnsWithTx: vi.fn(),
|
||||
mockDispatchAfterBatchInsert: vi.fn(),
|
||||
}))
|
||||
|
||||
vi.mock('@sim/utils/id', () => ({
|
||||
@@ -44,6 +46,7 @@ vi.mock('@/lib/table/service', () => ({
|
||||
batchInsertRowsWithTx: mockBatchInsertRowsWithTx,
|
||||
replaceTableRowsWithTx: mockReplaceTableRowsWithTx,
|
||||
addTableColumnsWithTx: mockAddTableColumnsWithTx,
|
||||
dispatchAfterBatchInsert: mockDispatchAfterBatchInsert,
|
||||
}))
|
||||
|
||||
import { POST } from '@/app/api/table/[tableId]/import/route'
|
||||
|
||||
@@ -29,11 +29,13 @@ import {
|
||||
type CsvHeaderMapping,
|
||||
CsvImportValidationError,
|
||||
coerceRowsForTable,
|
||||
dispatchAfterBatchInsert,
|
||||
inferColumnType,
|
||||
parseCsvBuffer,
|
||||
replaceTableRowsWithTx,
|
||||
sanitizeName,
|
||||
type TableDefinition,
|
||||
type TableRow,
|
||||
type TableSchema,
|
||||
validateMapping,
|
||||
} from '@/lib/table'
|
||||
@@ -228,13 +230,13 @@ export const POST = withRouteHandler(async (request: NextRequest, { params }: Ro
|
||||
}
|
||||
|
||||
try {
|
||||
const inserted = await db.transaction(async (trx) => {
|
||||
const txResult = await db.transaction(async (trx) => {
|
||||
let working = table
|
||||
if (additions.length > 0) {
|
||||
working = await addTableColumnsWithTx(trx, table, additions, requestId)
|
||||
}
|
||||
|
||||
let total = 0
|
||||
const allInserted: TableRow[] = []
|
||||
for (let i = 0; i < coerced.length; i += CSV_MAX_BATCH_SIZE) {
|
||||
const batch = coerced.slice(i, i + CSV_MAX_BATCH_SIZE)
|
||||
const batchRequestId = generateId().slice(0, 8)
|
||||
@@ -249,10 +251,15 @@ export const POST = withRouteHandler(async (request: NextRequest, { params }: Ro
|
||||
working,
|
||||
batchRequestId
|
||||
)
|
||||
total += result.length
|
||||
allInserted.push(...result)
|
||||
}
|
||||
return total
|
||||
return { inserted: allInserted, working }
|
||||
})
|
||||
const { inserted: insertedRows, working: finalTable } = txResult
|
||||
const inserted = insertedRows.length
|
||||
// Fire trigger + scheduler AFTER the tx commits — both read through the
|
||||
// global db connection and would otherwise see no rows.
|
||||
dispatchAfterBatchInsert(finalTable, insertedRows, requestId)
|
||||
|
||||
logger.info(`[${requestId}] Append CSV imported`, {
|
||||
tableId: table.id,
|
||||
|
||||
@@ -643,7 +643,12 @@ export async function addTableColumnsWithTx(
|
||||
)
|
||||
}
|
||||
|
||||
const updatedSchema: TableSchema = { columns: [...table.schema.columns, ...additions] }
|
||||
// Spread `table.schema` first so workflow groups (and any future top-level
|
||||
// schema fields) survive a CSV import that only adds plain columns.
|
||||
const updatedSchema: TableSchema = {
|
||||
...table.schema,
|
||||
columns: [...table.schema.columns, ...additions],
|
||||
}
|
||||
const now = new Date()
|
||||
|
||||
await trx
|
||||
@@ -1067,7 +1072,9 @@ export async function batchInsertRows(
|
||||
table: TableDefinition,
|
||||
requestId: string
|
||||
): Promise<TableRow[]> {
|
||||
return db.transaction((trx) => batchInsertRowsWithTx(trx, data, table, requestId))
|
||||
const result = await db.transaction((trx) => batchInsertRowsWithTx(trx, data, table, requestId))
|
||||
dispatchAfterBatchInsert(table, result, requestId)
|
||||
return result
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1165,7 +1172,21 @@ export async function batchInsertRowsWithTx(
|
||||
updatedAt: r.updatedAt,
|
||||
}))
|
||||
|
||||
void fireTableTrigger(data.tableId, table.name, 'insert', result, null, table.schema, requestId)
|
||||
return result
|
||||
}
|
||||
|
||||
/**
|
||||
* Side-effect dispatch for an insert batch. Caller fires this AFTER the
|
||||
* surrounding transaction commits — `fireTableTrigger` and `runWorkflowColumn`
|
||||
* both read through the global db connection, so firing inside the tx can see
|
||||
* no rows and no-op.
|
||||
*/
|
||||
export function dispatchAfterBatchInsert(
|
||||
table: TableDefinition,
|
||||
result: TableRow[],
|
||||
requestId: string
|
||||
): void {
|
||||
void fireTableTrigger(table.id, table.name, 'insert', result, null, table.schema, requestId)
|
||||
// Scope to the newly-inserted row ids so the dispatcher doesn't walk every
|
||||
// row in the table. After the sidecar migration, all existing rows have
|
||||
// zero entries → `mode:'new'`'s `NOT EXISTS` filter would otherwise include
|
||||
@@ -1178,8 +1199,6 @@ export async function batchInsertRowsWithTx(
|
||||
isManualRun: false,
|
||||
requestId,
|
||||
}).catch((err) => logger.error(`[${requestId}] auto-dispatch (batchInsertRows) failed:`, err))
|
||||
|
||||
return result
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
Reference in New Issue
Block a user