feat(tables): add export, import column creation, infinite row pagination (#4373)

* feat(tables): add export, import column creation, infinite row pagination

- Add `/api/table/[tableId]/export` route streaming CSV/JSON downloads
- Rename `/import-csv` route to `/import` and extend to auto-create new
  columns from unmapped CSV headers via `createColumns` form field
- Switch table view to `useInfiniteQuery` so tables larger than 1000
  rows fully load; reconcile created rows into the paginated cache so
  "New row" past 1000 no longer reverts on invalidate
- Wire scroll-driven prefetch (600px from bottom) and pre-drain pages
  before append to keep new-row position consistent
- Polish import-csv dialog flow and add Export action to header

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>

* fix(table): route boundary validation through Zod contracts

Switch the table query hooks and the import/export routes to the
codebase's contract-based request pattern so the API validation audit
and boundary policy ratchet pass.

- `hooks/queries/tables.ts` now calls every endpoint via
  `requestJson(contract, ...)`; the only remaining raw `fetch` calls are
  the streaming export download and the multipart CSV upload, both
  annotated `boundary-raw-fetch:`.
- The import and export routes parse `params`, the `format` query, and
  the multipart form fields with shared schemas from
  `@/lib/api/contracts/tables`. Two new contract schemas
  (`csvImportCreateColumnsSchema`, `tableExportFormatSchema`) cover the
  fields specific to these routes.
- Bumps the audit baseline by one route to account for the new export
  endpoint.

* fix(table): address PR bot review

- Move the `isAppendingRowRef` reset into the create-row mutation's
  `onSettled` callback. The previous `try/finally` cleared the guard
  immediately after `mutate()` returned, before the request completed,
  so a rapid second click on "New row" could fire a duplicate create.
- Drop the unused `addTableColumns` wrapper from `lib/table/service`.
  The CSV import flow only ever uses the transaction-bound
  `addTableColumnsWithTx`; the standalone wrapper was dead code.

* fix(table): re-throw on infinite-query fetch error in append-row drain

`useInfiniteQuery.fetchNextPage()` resolves (rather than rejects) when a
page request fails — the resolved value carries `status: 'error'` while
`hasNextPage` still reflects the last successful page. The drain loop in
`handleAppendRow` relied on a thrown error to bail, so a failed mid-drain
fetch could spin indefinitely and leave the append guard stuck on.

Re-throw inside `fetchNextPageWrapped` when the result is an error so the
caller's `try/catch` runs as intended.

* fix(table): import TABLE_LIMITS from constants to keep server code out of client bundle

The client hook `use-table-data.ts` was importing `TABLE_LIMITS` as a
value from the `@/lib/table` barrel, which transitively pulls in
`service.ts` and the `postgres` driver. Turbopack then tried to bundle
`fs`, `net`, `tls`, and `perf_hooks` into the client component graph and
the production build failed.

Import `TABLE_LIMITS` directly from `@/lib/table/constants` (a pure
constants module) and keep the type imports against the barrel.

* fix(table): run batch unique check inside import transaction

`checkBatchUniqueConstraintsDb` queried the global `db` connection, so
inside a single import transaction (one tx wrapping all batches) the
constraint lookup couldn't see uncommitted rows from prior batches —
duplicates that crossed `CSV_MAX_BATCH_SIZE` boundaries slipped through.

Accept an optional executor and pass `trx` from `batchInsertRowsWithTx`
so the lookup observes the in-flight transaction state.

---------

Co-authored-by: Claude Opus 4.7 <noreply@anthropic.com>
This commit is contained in:
Waleed
2026-05-01 10:00:24 -07:00
committed by GitHub
co-authored by Claude Opus 4.7
parent 47208e0de6
commit a9c12a2b36
16 changed files with 1142 additions and 267 deletions
@@ -0,0 +1,131 @@
import { createLogger } from '@sim/logger'
import { type NextRequest, NextResponse } from 'next/server'
import { tableExportFormatSchema, tableIdParamsSchema } from '@/lib/api/contracts/tables'
import { getValidationErrorMessage } from '@/lib/api/server'
import { checkSessionOrInternalAuth } from '@/lib/auth/hybrid'
import { generateRequestId } from '@/lib/core/utils/request'
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
import { queryRows } from '@/lib/table/service'
import { accessError, checkAccess } from '@/app/api/table/utils'
const logger = createLogger('TableExport')
const EXPORT_BATCH_SIZE = 1000
type ExportFormat = 'csv' | 'json'
interface RouteParams {
params: Promise<{ tableId: string }>
}
/** GET /api/table/[tableId]/export - Streams the full table contents as CSV or JSON. */
export const GET = withRouteHandler(async (request: NextRequest, { params }: RouteParams) => {
const requestId = generateRequestId()
const { tableId } = tableIdParamsSchema.parse(await params)
const auth = await checkSessionOrInternalAuth(request, { requireWorkflowId: false })
if (!auth.success || !auth.userId) {
return NextResponse.json({ error: 'Authentication required' }, { status: 401 })
}
const { searchParams } = new URL(request.url)
const formatValidation = tableExportFormatSchema.safeParse(
searchParams.get('format') ?? undefined
)
if (!formatValidation.success) {
return NextResponse.json(
{ error: getValidationErrorMessage(formatValidation.error) },
{ status: 400 }
)
}
const format: ExportFormat = formatValidation.data
const access = await checkAccess(tableId, auth.userId, 'read')
if (!access.ok) return accessError(access, requestId, tableId)
const { table } = access
const columns = table.schema.columns
const safeName = sanitizeFilename(table.name)
const filename = `${safeName}.${format}`
const stream = new ReadableStream<Uint8Array>({
async start(controller) {
const encoder = new TextEncoder()
try {
if (format === 'csv') {
controller.enqueue(encoder.encode(`${toCsvRow(columns.map((c) => c.name))}\n`))
} else {
controller.enqueue(encoder.encode('['))
}
let offset = 0
let firstJsonRow = true
while (true) {
const result = await queryRows(
tableId,
table.workspaceId,
{ limit: EXPORT_BATCH_SIZE, offset, includeTotal: false },
requestId
)
for (const row of result.rows) {
if (format === 'csv') {
const values = columns.map((c) => formatCsvValue(row.data[c.name]))
controller.enqueue(encoder.encode(`${toCsvRow(values)}\n`))
} else {
const prefix = firstJsonRow ? '' : ','
firstJsonRow = false
controller.enqueue(encoder.encode(prefix + JSON.stringify({ ...row.data })))
}
}
if (result.rows.length < EXPORT_BATCH_SIZE) break
offset += result.rows.length
}
if (format === 'json') controller.enqueue(encoder.encode(']'))
controller.close()
logger.info(`[${requestId}] Exported table ${tableId}`, {
format,
rowCount: table.rowCount,
})
} catch (err) {
logger.error(`[${requestId}] Export failed for table ${tableId}`, err)
controller.error(err)
}
},
})
return new NextResponse(stream, {
status: 200,
headers: {
'Content-Type': format === 'csv' ? 'text/csv; charset=utf-8' : 'application/json',
'Content-Disposition': `attachment; filename="${filename}"`,
'Cache-Control': 'no-store',
},
})
})
function sanitizeFilename(name: string): string {
const cleaned = name.replace(/[^a-zA-Z0-9_-]+/g, '_').replace(/^_+|_+$/g, '')
return cleaned || 'table'
}
function formatCsvValue(value: unknown): string {
if (value === null || value === undefined) return ''
if (value instanceof Date) return value.toISOString()
if (typeof value === 'object') return JSON.stringify(value)
return String(value)
}
function toCsvRow(values: string[]): string {
return values.map(escapeCsvField).join(',')
}
function escapeCsvField(field: string): string {
if (/[",\n\r]/.test(field)) {
return `"${field.replace(/"/g, '""')}"`
}
return field
}
@@ -6,10 +6,16 @@ import { NextRequest } from 'next/server'
import { beforeEach, describe, expect, it, vi } from 'vitest'
import type { TableDefinition } from '@/lib/table'
const { mockCheckAccess, mockBatchInsertRows, mockReplaceTableRows } = vi.hoisted(() => ({
const {
mockCheckAccess,
mockBatchInsertRowsWithTx,
mockReplaceTableRowsWithTx,
mockAddTableColumnsWithTx,
} = vi.hoisted(() => ({
mockCheckAccess: vi.fn(),
mockBatchInsertRows: vi.fn(),
mockReplaceTableRows: vi.fn(),
mockBatchInsertRowsWithTx: vi.fn(),
mockReplaceTableRowsWithTx: vi.fn(),
mockAddTableColumnsWithTx: vi.fn(),
}))
vi.mock('@sim/utils/id', () => ({
@@ -35,11 +41,12 @@ vi.mock('@/app/api/table/utils', async () => {
* `coerceRowsForTable`, etc.) exported through the barrel.
*/
vi.mock('@/lib/table/service', () => ({
batchInsertRows: mockBatchInsertRows,
replaceTableRows: mockReplaceTableRows,
batchInsertRowsWithTx: mockBatchInsertRowsWithTx,
replaceTableRowsWithTx: mockReplaceTableRowsWithTx,
addTableColumnsWithTx: mockAddTableColumnsWithTx,
}))
import { POST } from '@/app/api/table/[tableId]/import-csv/route'
import { POST } from '@/app/api/table/[tableId]/import/route'
function createCsvFile(contents: string, name = 'data.csv', type = 'text/csv'): File {
return new File([contents], name, { type })
@@ -51,6 +58,7 @@ function createFormData(
workspaceId?: string | null
mode?: string | null
mapping?: unknown
createColumns?: unknown
}
): FormData {
const form = new FormData()
@@ -67,6 +75,14 @@ function createFormData(
typeof options.mapping === 'string' ? options.mapping : JSON.stringify(options.mapping)
)
}
if (options?.createColumns !== undefined) {
form.append(
'createColumns',
typeof options.createColumns === 'string'
? options.createColumns
: JSON.stringify(options.createColumns)
)
}
return form
}
@@ -94,14 +110,14 @@ function buildTable(overrides: Partial<TableDefinition> = {}): TableDefinition {
}
async function callPost(form: FormData, { tableId }: { tableId: string } = { tableId: 'tbl_1' }) {
const req = new NextRequest(`http://localhost:3000/api/table/${tableId}/import-csv`, {
const req = new NextRequest(`http://localhost:3000/api/table/${tableId}/import`, {
method: 'POST',
body: form,
})
return POST(req, { params: Promise.resolve({ tableId }) })
}
describe('POST /api/table/[tableId]/import-csv', () => {
describe('POST /api/table/[tableId]/import', () => {
beforeEach(() => {
vi.clearAllMocks()
hybridAuthMockFns.mockCheckSessionOrInternalAuth.mockResolvedValue({
@@ -110,10 +126,25 @@ describe('POST /api/table/[tableId]/import-csv', () => {
authType: 'session',
})
mockCheckAccess.mockResolvedValue({ ok: true, table: buildTable() })
mockBatchInsertRows.mockImplementation(async (data: { rows: unknown[] }) =>
mockBatchInsertRowsWithTx.mockImplementation(async (_trx, data: { rows: unknown[] }) =>
data.rows.map((_, i) => ({ id: `row_${i}` }))
)
mockReplaceTableRows.mockResolvedValue({ deletedCount: 0, insertedCount: 0 })
mockReplaceTableRowsWithTx.mockResolvedValue({ deletedCount: 0, insertedCount: 0 })
mockAddTableColumnsWithTx.mockImplementation(
async (
_trx,
table: { schema: { columns: { name: string; type: string }[] } },
columns: { name: string; type: string }[]
) => ({
...table,
schema: {
columns: [
...table.schema.columns,
...columns.map((c) => ({ name: c.name, type: c.type as 'string' })),
],
},
})
)
})
it('returns 401 when the user is not authenticated', async () => {
@@ -157,7 +188,7 @@ describe('POST /api/table/[tableId]/import-csv', () => {
const data = await response.json()
expect(data.error).toMatch(/missing required columns/i)
expect(data.details?.missingRequired).toEqual(['name'])
expect(mockBatchInsertRows).not.toHaveBeenCalled()
expect(mockBatchInsertRowsWithTx).not.toHaveBeenCalled()
})
it('appends rows via batchInsertRows', async () => {
@@ -168,13 +199,13 @@ describe('POST /api/table/[tableId]/import-csv', () => {
const data = await response.json()
expect(data.data.mode).toBe('append')
expect(data.data.insertedCount).toBe(2)
expect(mockBatchInsertRows).toHaveBeenCalledTimes(1)
const callArgs = mockBatchInsertRows.mock.calls[0][0] as { rows: unknown[] }
expect(mockBatchInsertRowsWithTx).toHaveBeenCalledTimes(1)
const callArgs = mockBatchInsertRowsWithTx.mock.calls[0][1] as { rows: unknown[] }
expect(callArgs.rows).toEqual([
{ name: 'Alice', age: 30 },
{ name: 'Bob', age: 40 },
])
expect(mockReplaceTableRows).not.toHaveBeenCalled()
expect(mockReplaceTableRowsWithTx).not.toHaveBeenCalled()
})
it('rejects append when it would exceed maxRows', async () => {
@@ -188,11 +219,11 @@ describe('POST /api/table/[tableId]/import-csv', () => {
expect(response.status).toBe(400)
const data = await response.json()
expect(data.error).toMatch(/exceed table row limit/)
expect(mockBatchInsertRows).not.toHaveBeenCalled()
expect(mockBatchInsertRowsWithTx).not.toHaveBeenCalled()
})
it('replaces rows via replaceTableRows', async () => {
mockReplaceTableRows.mockResolvedValueOnce({ deletedCount: 5, insertedCount: 2 })
mockReplaceTableRowsWithTx.mockResolvedValueOnce({ deletedCount: 5, insertedCount: 2 })
const response = await callPost(
createFormData(createCsvFile('name,age\nAlice,30\nBob,40'), { mode: 'replace' })
)
@@ -201,8 +232,8 @@ describe('POST /api/table/[tableId]/import-csv', () => {
expect(data.data.mode).toBe('replace')
expect(data.data.deletedCount).toBe(5)
expect(data.data.insertedCount).toBe(2)
expect(mockReplaceTableRows).toHaveBeenCalledTimes(1)
expect(mockBatchInsertRows).not.toHaveBeenCalled()
expect(mockReplaceTableRowsWithTx).toHaveBeenCalledTimes(1)
expect(mockBatchInsertRowsWithTx).not.toHaveBeenCalled()
})
it('uses an explicit mapping when provided', async () => {
@@ -215,7 +246,7 @@ describe('POST /api/table/[tableId]/import-csv', () => {
expect(response.status).toBe(200)
const data = await response.json()
expect(data.data.mappedColumns).toEqual(['First Name', 'Years'])
const callArgs = mockBatchInsertRows.mock.calls[0][0] as { rows: unknown[] }
const callArgs = mockBatchInsertRowsWithTx.mock.calls[0][1] as { rows: unknown[] }
expect(callArgs.rows).toEqual([
{ name: 'Alice', age: 30 },
{ name: 'Bob', age: 40 },
@@ -247,7 +278,7 @@ describe('POST /api/table/[tableId]/import-csv', () => {
})
it('surfaces unique violations from batchInsertRows as 400', async () => {
mockBatchInsertRows.mockRejectedValueOnce(
mockBatchInsertRowsWithTx.mockRejectedValueOnce(
new Error('Row 1: Column "name" must be unique. Value "Alice" already exists in row row_xxx')
)
const response = await callPost(
@@ -267,7 +298,7 @@ describe('POST /api/table/[tableId]/import-csv', () => {
)
)
expect(response.status).toBe(200)
expect(mockBatchInsertRows).toHaveBeenCalledTimes(1)
expect(mockBatchInsertRowsWithTx).toHaveBeenCalledTimes(1)
})
it('returns 400 for unsupported file extensions', async () => {
@@ -278,4 +309,137 @@ describe('POST /api/table/[tableId]/import-csv', () => {
const data = await response.json()
expect(data.error).toMatch(/CSV and TSV/)
})
describe('createColumns', () => {
it('auto-creates columns for unmapped CSV headers', async () => {
const response = await callPost(
createFormData(createCsvFile('name,age,email\nAlice,30,a@x.io\nBob,40,b@x.io'), {
mode: 'append',
createColumns: ['email'],
})
)
expect(response.status).toBe(200)
expect(mockAddTableColumnsWithTx).toHaveBeenCalledTimes(1)
const [, , columns] = mockAddTableColumnsWithTx.mock.calls[0]
expect(columns).toEqual([{ name: 'email', type: 'string' }])
const callArgs = mockBatchInsertRowsWithTx.mock.calls[0][1] as { rows: unknown[] }
expect(callArgs.rows).toEqual([
{ name: 'Alice', age: 30, email: 'a@x.io' },
{ name: 'Bob', age: 40, email: 'b@x.io' },
])
})
it('infers column type from CSV row values', async () => {
const response = await callPost(
createFormData(createCsvFile('name,score\nAlice,42\nBob,17'), {
mode: 'append',
createColumns: ['score'],
})
)
expect(response.status).toBe(200)
const [, , columns] = mockAddTableColumnsWithTx.mock.calls[0]
expect(columns).toEqual([{ name: 'score', type: 'number' }])
})
it('dedupes when sanitized name collides with an existing column', async () => {
mockCheckAccess.mockResolvedValueOnce({
ok: true,
table: buildTable({
schema: {
columns: [
{ name: 'name', type: 'string', required: true },
{ name: 'age', type: 'number' },
{ name: 'email', type: 'string' },
],
},
}),
})
const response = await callPost(
createFormData(createCsvFile('name,age,Email\nAlice,30,a@x.io'), {
mode: 'append',
createColumns: ['Email'],
})
)
expect(response.status).toBe(200)
const [, , columns] = mockAddTableColumnsWithTx.mock.calls[0]
expect(columns).toEqual([{ name: 'Email_2', type: 'string' }])
})
it('returns 400 when createColumns references a header not in the CSV', async () => {
const response = await callPost(
createFormData(createCsvFile('name,age\nAlice,30'), {
mode: 'append',
createColumns: ['nonexistent'],
})
)
expect(response.status).toBe(400)
const data = await response.json()
expect(data.error).toMatch(/unknown CSV headers/)
expect(mockAddTableColumnsWithTx).not.toHaveBeenCalled()
expect(mockBatchInsertRowsWithTx).not.toHaveBeenCalled()
})
it('returns 400 when createColumns is not an array of strings', async () => {
const response = await callPost(
createFormData(createCsvFile('name,age\nAlice,30'), {
mode: 'append',
createColumns: [1, 2],
})
)
expect(response.status).toBe(400)
const data = await response.json()
expect(data.error).toMatch(/createColumns must be a JSON array/)
expect(mockAddTableColumnsWithTx).not.toHaveBeenCalled()
})
it('returns 400 when createColumns is invalid JSON', async () => {
const response = await callPost(
createFormData(createCsvFile('name,age\nAlice,30'), {
mode: 'append',
createColumns: '{not-json',
})
)
expect(response.status).toBe(400)
const data = await response.json()
expect(data.error).toMatch(/createColumns must be valid JSON/)
})
it('surfaces addTableColumns failures as 400', async () => {
mockAddTableColumnsWithTx.mockRejectedValueOnce(new Error('Column "email" already exists'))
const response = await callPost(
createFormData(createCsvFile('name,age,email\nAlice,30,a@x.io'), {
mode: 'append',
createColumns: ['email'],
})
)
expect(response.status).toBe(400)
const data = await response.json()
expect(data.error).toMatch(/already exists/)
expect(mockBatchInsertRowsWithTx).not.toHaveBeenCalled()
})
it('surfaces row insert failures without success when schema was mutated', async () => {
mockBatchInsertRowsWithTx.mockRejectedValueOnce(new Error('must be unique'))
const response = await callPost(
createFormData(createCsvFile('name,age,email\nAlice,30,a@x.io'), {
mode: 'append',
createColumns: ['email'],
})
)
expect(mockAddTableColumnsWithTx).toHaveBeenCalled()
expect(response.status).toBe(400)
const data = await response.json()
expect(data.success).toBeUndefined()
expect(data.error).toMatch(/must be unique/)
})
it('does not call addTableColumns when createColumns is omitted', async () => {
const response = await callPost(
createFormData(createCsvFile('name,age\nAlice,30'), { mode: 'append' })
)
expect(response.status).toBe(200)
expect(mockAddTableColumnsWithTx).not.toHaveBeenCalled()
})
})
})
@@ -1,26 +1,34 @@
import { db } from '@sim/db'
import { createLogger } from '@sim/logger'
import { toError } from '@sim/utils/errors'
import { generateId } from '@sim/utils/id'
import { type NextRequest, NextResponse } from 'next/server'
import {
csvExtensionSchema,
csvImportCreateColumnsSchema,
csvImportFormSchema,
csvImportMappingSchema,
csvImportModeSchema,
tableIdParamsSchema,
} from '@/lib/api/contracts/tables'
import { getValidationErrorMessage } from '@/lib/api/server'
import { checkSessionOrInternalAuth } from '@/lib/auth/hybrid'
import { generateRequestId } from '@/lib/core/utils/request'
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
import {
batchInsertRows,
addTableColumnsWithTx,
batchInsertRowsWithTx,
buildAutoMapping,
CSV_MAX_BATCH_SIZE,
type CsvHeaderMapping,
CsvImportValidationError,
coerceRowsForTable,
inferColumnType,
parseCsvBuffer,
replaceTableRows,
replaceTableRowsWithTx,
sanitizeName,
type TableDefinition,
type TableSchema,
validateMapping,
} from '@/lib/table'
import { accessError, checkAccess } from '@/app/api/table/utils'
@@ -33,7 +41,7 @@ interface RouteParams {
export const POST = withRouteHandler(async (request: NextRequest, { params }: RouteParams) => {
const requestId = generateRequestId()
const { tableId } = await params
const { tableId } = tableIdParamsSchema.parse(await params)
try {
const authResult = await checkSessionOrInternalAuth(request, { requireWorkflowId: false })
@@ -48,6 +56,7 @@ export const POST = withRouteHandler(async (request: NextRequest, { params }: Ro
})
const rawMode = formData.get('mode') ?? 'append'
const rawMapping = formData.get('mapping')
const rawCreateColumns = formData.get('createColumns')
if (!formValidation.success) {
return NextResponse.json(
@@ -61,7 +70,7 @@ export const POST = withRouteHandler(async (request: NextRequest, { params }: Ro
const modeValidation = csvImportModeSchema.safeParse(rawMode)
if (!modeValidation.success) {
return NextResponse.json(
{ error: `Invalid mode "${rawMode}". Must be "append" or "replace".` },
{ error: `Invalid mode "${String(rawMode)}". Must be "append" or "replace".` },
{ status: 400 }
)
}
@@ -104,18 +113,75 @@ export const POST = withRouteHandler(async (request: NextRequest, { params }: Ro
mapping = mappingValidation.data
}
let createColumns: string[] | undefined
if (rawCreateColumns) {
const createColumnsValidation = csvImportCreateColumnsSchema.safeParse(rawCreateColumns)
if (!createColumnsValidation.success) {
return NextResponse.json(
{ error: getValidationErrorMessage(createColumnsValidation.error) },
{ status: 400 }
)
}
createColumns = createColumnsValidation.data
}
const buffer = Buffer.from(await file.arrayBuffer())
const delimiter = extensionValidation.data === 'tsv' ? '\t' : ','
const { headers, rows } = await parseCsvBuffer(buffer, delimiter)
const effectiveMapping = mapping ?? buildAutoMapping(headers, table.schema)
let effectiveMapping = mapping ?? buildAutoMapping(headers, table.schema)
let prospectiveTable: TableDefinition = table
const additions: { name: string; type: string }[] = []
if (createColumns && createColumns.length > 0) {
const headerSet = new Set(headers)
const unknownHeaders = createColumns.filter((h) => !headerSet.has(h))
if (unknownHeaders.length > 0) {
return NextResponse.json(
{
error: `createColumns references unknown CSV headers: ${unknownHeaders.join(', ')}`,
},
{ status: 400 }
)
}
const usedNames = new Set(table.schema.columns.map((c) => c.name.toLowerCase()))
const updatedMapping: CsvHeaderMapping = { ...effectiveMapping }
const newColumns: TableSchema['columns'] = []
for (const header of createColumns) {
const base = sanitizeName(header)
let columnName = base
let suffix = 2
while (usedNames.has(columnName.toLowerCase())) {
columnName = `${base}_${suffix}`
suffix++
}
usedNames.add(columnName.toLowerCase())
const inferredType = inferColumnType(rows.map((r) => r[header]))
additions.push({ name: columnName, type: inferredType })
newColumns.push({
name: columnName,
type: inferredType as TableSchema['columns'][number]['type'],
required: false,
unique: false,
})
updatedMapping[header] = columnName
}
prospectiveTable = {
...table,
schema: { columns: [...table.schema.columns, ...newColumns] },
}
effectiveMapping = updatedMapping
}
let validation: ReturnType<typeof validateMapping>
try {
validation = validateMapping({
csvHeaders: headers,
mapping: effectiveMapping,
tableSchema: table.schema,
tableSchema: prospectiveTable.schema,
})
} catch (err) {
if (err instanceof CsvImportValidationError) {
@@ -127,47 +193,79 @@ export const POST = withRouteHandler(async (request: NextRequest, { params }: Ro
if (validation.mappedHeaders.length === 0) {
return NextResponse.json(
{
error: `No CSV headers map to columns on the table. CSV headers: ${headers.join(', ')}. Table columns: ${table.schema.columns.map((c) => c.name).join(', ')}`,
error: `No CSV headers map to columns on the table. CSV headers: ${headers.join(', ')}. Table columns: ${prospectiveTable.schema.columns.map((c) => c.name).join(', ')}`,
},
{ status: 400 }
)
}
const coerced = coerceRowsForTable(rows, table.schema, validation.effectiveMap)
const coerced = coerceRowsForTable(rows, prospectiveTable.schema, validation.effectiveMap)
if (mode === 'append') {
if (table.rowCount + coerced.length > table.maxRows) {
const deficit = table.rowCount + coerced.length - table.maxRows
if (prospectiveTable.rowCount + coerced.length > prospectiveTable.maxRows) {
const deficit = prospectiveTable.rowCount + coerced.length - prospectiveTable.maxRows
return NextResponse.json(
{
error: `Append would exceed table row limit (${table.maxRows}). Currently ${table.rowCount} rows, ${coerced.length} new rows, ${deficit} over.`,
error: `Append would exceed table row limit (${prospectiveTable.maxRows}). Currently ${prospectiveTable.rowCount} rows, ${coerced.length} new rows, ${deficit} over.`,
},
{ status: 400 }
)
}
let inserted = 0
try {
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)
const result = await batchInsertRows(
{
tableId: table.id,
rows: batch,
workspaceId,
userId: authResult.userId,
},
table,
batchRequestId
)
inserted += result.length
}
const inserted = await db.transaction(async (trx) => {
let working = table
if (additions.length > 0) {
working = await addTableColumnsWithTx(trx, table, additions, requestId)
}
let total = 0
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)
const result = await batchInsertRowsWithTx(
trx,
{
tableId: working.id,
rows: batch,
workspaceId,
userId: authResult.userId,
},
working,
batchRequestId
)
total += result.length
}
return total
})
logger.info(`[${requestId}] Append CSV imported`, {
tableId: table.id,
fileName: file.name,
mode,
inserted,
createdColumns: additions.length,
mappedColumns: validation.mappedHeaders.length,
skippedHeaders: validation.skippedHeaders.length,
})
return NextResponse.json({
success: true,
data: {
tableId: table.id,
mode,
insertedCount: inserted,
mappedColumns: validation.mappedHeaders,
skippedHeaders: validation.skippedHeaders,
unmappedColumns: validation.unmappedColumns,
sourceFile: file.name,
},
})
} catch (err) {
const message = toError(err).message
logger.warn(`[${requestId}] Append failed mid-import for table ${tableId}`, {
inserted,
logger.warn(`[${requestId}] Append failed for table ${tableId}`, {
total: coerced.length,
createdColumns: additions.length,
error: message,
})
const isClientError =
@@ -176,45 +274,32 @@ export const POST = withRouteHandler(async (request: NextRequest, { params }: Ro
message.includes('Schema validation') ||
message.includes('must be unique') ||
message.includes('Row size exceeds') ||
message.includes('already exists') ||
message.includes('Invalid column name') ||
/^Row \d+:/.test(message)
return NextResponse.json(
{
error: isClientError ? message : 'Failed to import CSV',
data: { insertedCount: inserted },
data: { insertedCount: 0 },
},
{ status: isClientError ? 400 : 500 }
)
}
logger.info(`[${requestId}] Append CSV imported`, {
tableId: table.id,
fileName: file.name,
mode,
inserted,
mappedColumns: validation.mappedHeaders.length,
skippedHeaders: validation.skippedHeaders.length,
})
return NextResponse.json({
success: true,
data: {
tableId: table.id,
mode,
insertedCount: inserted,
mappedColumns: validation.mappedHeaders,
skippedHeaders: validation.skippedHeaders,
unmappedColumns: validation.unmappedColumns,
sourceFile: file.name,
},
})
}
try {
const result = await replaceTableRows(
{ tableId: table.id, rows: coerced, workspaceId, userId: authResult.userId },
table,
requestId
)
const result = await db.transaction(async (trx) => {
let working = table
if (additions.length > 0) {
working = await addTableColumnsWithTx(trx, table, additions, requestId)
}
return replaceTableRowsWithTx(
trx,
{ tableId: working.id, rows: coerced, workspaceId, userId: authResult.userId },
working,
requestId
)
})
logger.info(`[${requestId}] Replace CSV imported`, {
tableId: table.id,
@@ -222,6 +307,7 @@ export const POST = withRouteHandler(async (request: NextRequest, { params }: Ro
mode,
deleted: result.deletedCount,
inserted: result.insertedCount,
createdColumns: additions.length,
mappedColumns: validation.mappedHeaders.length,
})
@@ -245,6 +331,8 @@ export const POST = withRouteHandler(async (request: NextRequest, { params }: Ro
message.includes('Schema validation') ||
message.includes('must be unique') ||
message.includes('Row size exceeds') ||
message.includes('already exists') ||
message.includes('Invalid column name') ||
/^Row \d+:/.test(message)
if (isClientError) {
return NextResponse.json({ error: message }, { status: 400 })
@@ -64,6 +64,7 @@ import { useUsageLimits } from '@/app/workspace/[workspaceId]/w/[workflowId]/com
import { useWorkflowExecution } from '@/app/workspace/[workspaceId]/w/[workflowId]/hooks/use-workflow-execution'
import { useFolders } from '@/hooks/queries/folders'
import { useLogDetail } from '@/hooks/queries/logs'
import { downloadTableExport } from '@/hooks/queries/tables'
import { useWorkflows } from '@/hooks/queries/workflows'
import { useWorkspaceFiles } from '@/hooks/queries/workspace-files'
import { useSettingsNavigation } from '@/hooks/use-settings-navigation'
@@ -222,6 +223,14 @@ export function ResourceActions({ workspaceId, resource }: ResourceActionsProps)
return (
<EmbeddedKnowledgeBaseActions workspaceId={workspaceId} knowledgeBaseId={resource.id} />
)
case 'table':
return (
<EmbeddedTableActions
workspaceId={workspaceId}
tableId={resource.id}
tableName={resource.title}
/>
)
case 'log':
return <EmbeddedLogActions workspaceId={workspaceId} logId={resource.id} />
case 'folder':
@@ -353,6 +362,65 @@ export function EmbeddedKnowledgeBaseActions({
)
}
const tableLogger = createLogger('EmbeddedTableActions')
interface EmbeddedTableActionsProps {
workspaceId: string
tableId: string
tableName: string
}
function EmbeddedTableActions({ workspaceId, tableId, tableName }: EmbeddedTableActionsProps) {
const router = useRouter()
const handleOpenTable = () => {
router.push(`/workspace/${workspaceId}/tables/${tableId}`)
}
const handleExport = async () => {
try {
await downloadTableExport(tableId, tableName)
} catch (err) {
tableLogger.error('Failed to export table:', err)
}
}
return (
<>
<Tooltip.Root>
<Tooltip.Trigger asChild>
<Button
variant='subtle'
onClick={handleOpenTable}
className={RESOURCE_TAB_ICON_BUTTON_CLASS}
aria-label='Open table'
>
<SquareArrowUpRight className={RESOURCE_TAB_ICON_CLASS} />
</Button>
</Tooltip.Trigger>
<Tooltip.Content side='bottom'>
<p>Open table</p>
</Tooltip.Content>
</Tooltip.Root>
<Tooltip.Root>
<Tooltip.Trigger asChild>
<Button
variant='subtle'
onClick={() => void handleExport()}
className={RESOURCE_TAB_ICON_BUTTON_CLASS}
aria-label='Export table as CSV'
>
<Download className={RESOURCE_TAB_ICON_CLASS} />
</Button>
</Tooltip.Trigger>
<Tooltip.Content side='bottom'>
<p>Export CSV</p>
</Tooltip.Content>
</Tooltip.Root>
</>
)
}
const fileLogger = createLogger('EmbeddedFileActions')
interface EmbeddedFileActionsProps {
@@ -1,12 +1,14 @@
'use client'
import React, { useCallback, useEffect, useMemo, useRef, useState } from 'react'
import { createLogger } from '@sim/logger'
import { useParams, useRouter } from 'next/navigation'
import { usePostHog } from 'posthog-js/react'
import {
Button,
Checkbox,
DatePicker,
Download,
DropdownMenu,
DropdownMenuContent,
DropdownMenuItem,
@@ -21,6 +23,7 @@ import {
ModalFooter,
ModalHeader,
Skeleton,
toast,
Upload,
} from '@/components/emcn'
import {
@@ -47,6 +50,7 @@ import { ResourceHeader, ResourceOptionsBar } from '@/app/workspace/[workspaceId
import { useUserPermissionsContext } from '@/app/workspace/[workspaceId]/providers/workspace-permissions-provider'
import { ImportCsvDialog } from '@/app/workspace/[workspaceId]/tables/components/import-csv-dialog'
import {
downloadTableExport,
useAddTableColumn,
useBatchCreateTableRows,
useBatchUpdateTableRows,
@@ -87,6 +91,8 @@ interface NormalizedSelection {
anchorCol: number
}
const logger = createLogger('TableView')
const EMPTY_COLUMNS: never[] = []
const EMPTY_CHECKED_ROWS = new Set<number>()
const COL_WIDTH = 160
@@ -226,12 +232,28 @@ export function Table({
const scrollRef = useRef<HTMLDivElement>(null)
const isDraggingRef = useRef(false)
const { tableData, isLoadingTable, rows, isLoadingRows } = useTableData({
const {
tableData,
isLoadingTable,
rows,
isLoadingRows,
fetchNextPage,
hasNextPage,
isFetchingNextPage,
} = useTableData({
workspaceId,
tableId,
queryOptions,
})
const fetchNextPageRef = useRef(fetchNextPage)
fetchNextPageRef.current = fetchNextPage
const hasNextPageRef = useRef(hasNextPage)
hasNextPageRef.current = hasNextPage
const isFetchingNextPageRef = useRef(isFetchingNextPage)
isFetchingNextPageRef.current = isFetchingNextPage
const isAppendingRowRef = useRef(false)
const userPermissions = useUserPermissionsContext()
const canEditRef = useRef(userPermissions.canEdit)
canEditRef.current = userPermissions.canEdit
@@ -577,7 +599,21 @@ export function Table({
)
}, [contextMenu.row, closeContextMenu])
const handleAppendRow = useCallback(() => {
const handleAppendRow = useCallback(async () => {
if (isAppendingRowRef.current) return
isAppendingRowRef.current = true
try {
while (hasNextPageRef.current) {
const result = await fetchNextPageRef.current()
if (!result.hasNextPage) break
}
} catch (error) {
isAppendingRowRef.current = false
logger.error('Failed to load remaining rows before appending', { error })
toast.error('Failed to load all rows. Try again.', { duration: 5000 })
return
}
createRef.current(
{ data: {} },
{
@@ -591,6 +627,9 @@ export function Table({
})
}
},
onSettled: () => {
isAppendingRowRef.current = false
},
}
)
}, [])
@@ -889,6 +928,30 @@ export function Table({
e.preventDefault()
}
useEffect(() => {
const scrollEl = scrollRef.current
if (!scrollEl) return
const SCROLL_PREFETCH_PX = 600
function maybeFetchNext() {
if (!hasNextPageRef.current || isFetchingNextPageRef.current) return
if (!scrollEl) return
const distanceFromBottom = scrollEl.scrollHeight - scrollEl.scrollTop - scrollEl.clientHeight
if (distanceFromBottom <= SCROLL_PREFETCH_PX) {
fetchNextPageRef.current().catch((error) => {
logger.error('Failed to fetch next page of rows', { error })
})
}
}
maybeFetchNext()
scrollEl.addEventListener('scroll', maybeFetchNext, { passive: true })
return () => {
scrollEl.removeEventListener('scroll', maybeFetchNext)
}
}, [tableData?.id])
useEffect(() => {
if (!tableData?.metadata || metadataSeededRef.current) return
if (!tableData.metadata.columnWidths && !tableData.metadata.columnOrder) return
@@ -1954,6 +2017,16 @@ export function Table({
[handleAddColumn, addColumnMutation.isPending]
)
const handleExportCsv = useCallback(async () => {
if (!tableData) return
try {
await downloadTableExport(tableData.id, tableData.name)
} catch (err) {
logger.error('Failed to export table:', err)
toast.error('Failed to export table')
}
}, [tableData])
const headerActions = useMemo(
() =>
tableData
@@ -1964,9 +2037,15 @@ export function Table({
onClick: () => setIsImportCsvOpen(true),
disabled: userPermissions.canEdit !== true,
},
{
label: 'Export CSV',
icon: Download,
onClick: () => void handleExportCsv(),
disabled: tableData.rowCount === 0,
},
]
: undefined,
[tableData, userPermissions.canEdit]
[tableData, userPermissions.canEdit, handleExportCsv]
)
const activeSortState = useMemo(() => {
@@ -1,5 +1,7 @@
import { useCallback, useMemo } from 'react'
import type { TableDefinition, TableRow } from '@/lib/table'
import { useTable, useTableRows } from '@/hooks/queries/tables'
import { TABLE_LIMITS } from '@/lib/table/constants'
import { useInfiniteTableRows, useTable } from '@/hooks/queries/tables'
import type { QueryOptions } from '../types'
interface UseTableDataParams {
@@ -8,12 +10,24 @@ interface UseTableDataParams {
queryOptions: QueryOptions
}
interface FetchNextPageResult {
hasNextPage: boolean
}
interface UseTableDataReturn {
tableData: TableDefinition | undefined
isLoadingTable: boolean
rows: TableRow[]
isLoadingRows: boolean
refetchRows: () => void
/**
* Fetch the next page of rows. The resolved value's `hasNextPage` reflects the
* post-fetch cache state — read from this rather than the parent's
* `hasNextPage` state, which only updates on the next React render.
*/
fetchNextPage: () => Promise<FetchNextPageResult>
hasNextPage: boolean
isFetchingNextPage: boolean
}
export function useTableData({
@@ -26,19 +40,35 @@ export function useTableData({
const {
data: rowsData,
isLoading: isLoadingRows,
refetch: refetchRows,
} = useTableRows({
refetch,
fetchNextPage,
hasNextPage,
isFetchingNextPage,
} = useInfiniteTableRows({
workspaceId,
tableId,
limit: 1000,
offset: 0,
pageSize: TABLE_LIMITS.MAX_QUERY_LIMIT,
filter: queryOptions.filter,
sort: queryOptions.sort,
includeTotal: false,
enabled: Boolean(workspaceId && tableId),
})
const rows = (rowsData?.rows || []) as TableRow[]
const rows = useMemo<TableRow[]>(
() => rowsData?.pages.flatMap((p) => p.rows) ?? [],
[rowsData?.pages]
)
const refetchRows = useCallback(() => {
void refetch()
}, [refetch])
const fetchNextPageWrapped = useCallback(async () => {
const result = await fetchNextPage()
if (result.status === 'error') {
throw result.error ?? new Error('Failed to fetch next page')
}
return { hasNextPage: Boolean(result.hasNextPage) }
}, [fetchNextPage])
return {
tableData,
@@ -46,5 +76,8 @@ export function useTableData({
rows,
isLoadingRows,
refetchRows,
fetchNextPage: fetchNextPageWrapped,
hasNextPage: Boolean(hasNextPage),
isFetchingNextPage,
}
}
@@ -23,7 +23,7 @@ import {
toast,
} from '@/components/emcn'
import { cn } from '@/lib/core/utils/cn'
import { buildAutoMapping, parseCsvBuffer } from '@/lib/table/csv-import'
import { buildAutoMapping, parseCsvBuffer } from '@/lib/table/import'
import type { TableDefinition } from '@/lib/table/types'
import { type CsvImportMode, useImportCsvIntoTable } from '@/hooks/queries/tables'
@@ -37,6 +37,11 @@ const MAX_EXAMPLES_IN_ERROR = 3
* (`/^[a-z_][a-z0-9_]*$/i`), so no real column can share this value.
*/
const SKIP_VALUE = '__ skip __'
/**
* Sentinel for the "Create new column" option. Same whitespace trick as
* `SKIP_VALUE` to avoid colliding with any valid column name.
*/
const CREATE_VALUE = '__ create __'
/**
* Converts the verbose backend error messages into a short, human-friendly
@@ -102,6 +107,7 @@ export function ImportCsvDialog({
const [submitError, setSubmitError] = useState<string | null>(null)
const [parsing, setParsing] = useState(false)
const [mapping, setMapping] = useState<Record<string, string | null>>({})
const [createHeaders, setCreateHeaders] = useState<Set<string>>(new Set())
const [mode, setMode] = useState<CsvImportMode>('append')
const [isDragging, setIsDragging] = useState(false)
const fileInputRef = useRef<HTMLInputElement>(null)
@@ -112,6 +118,7 @@ export function ImportCsvDialog({
setParseError(null)
setSubmitError(null)
setMapping({})
setCreateHeaders(new Set())
setMode('append')
setIsDragging(false)
setParsing(false)
@@ -130,7 +137,10 @@ export function ImportCsvDialog({
}
const columnOptions: ComboboxOption[] = useMemo(() => {
const options: ComboboxOption[] = [{ label: 'Do not import', value: SKIP_VALUE }]
const options: ComboboxOption[] = [
{ label: 'Do not import', value: SKIP_VALUE },
{ label: '+ Create new column', value: CREATE_VALUE },
]
for (const col of table.schema.columns) {
options.push({
label: col.required ? `${col.name} (required)` : col.name,
@@ -197,22 +207,54 @@ export function ImportCsvDialog({
function handleMappingChange(header: string, value: string) {
setSubmitError(null)
if (value === CREATE_VALUE) {
setCreateHeaders((prev) => {
const next = new Set(prev)
next.add(header)
return next
})
setMapping((prev) => ({ ...prev, [header]: null }))
return
}
setCreateHeaders((prev) => {
if (!prev.has(header)) return prev
const next = new Set(prev)
next.delete(header)
return next
})
setMapping((prev) => ({
...prev,
[header]: value === SKIP_VALUE ? null : value,
}))
}
function handleCreateAllUnmapped() {
if (!parsed) return
setSubmitError(null)
setCreateHeaders((prev) => {
const next = new Set(prev)
for (const header of parsed.headers) {
if (!mapping[header] && !next.has(header)) next.add(header)
}
return next
})
}
function handleModeChange(value: string) {
setSubmitError(null)
setMode(value as CsvImportMode)
}
const { missingRequired, duplicateTargets, mappedCount, skipCount } = useMemo(() => {
const { missingRequired, duplicateTargets, mappedCount, skipCount, createCount } = useMemo(() => {
const mappedTargets = new Map<string, string[]>()
let mapped = 0
let skipped = 0
let creating = 0
for (const header of parsed?.headers ?? []) {
if (createHeaders.has(header)) {
creating++
continue
}
const target = mapping[header]
if (!target) {
skipped++
@@ -235,8 +277,9 @@ export function ImportCsvDialog({
duplicateTargets: dupes,
mappedCount: mapped,
skipCount: skipped,
createCount: creating,
}
}, [mapping, parsed?.headers, table.schema.columns])
}, [mapping, parsed?.headers, table.schema.columns, createHeaders])
const appendCapacityDeficit =
parsed && mode === 'append' && table.rowCount + parsed.totalRows > table.maxRows
@@ -253,7 +296,7 @@ export function ImportCsvDialog({
!importMutation.isPending &&
missingRequired.length === 0 &&
duplicateTargets.length === 0 &&
mappedCount > 0 &&
mappedCount + createCount > 0 &&
appendCapacityDeficit === 0 &&
replaceCapacityDeficit === 0
@@ -267,6 +310,7 @@ export function ImportCsvDialog({
file: parsed.file,
mode,
mapping,
createColumns: createHeaders.size > 0 ? [...createHeaders] : undefined,
})
const data = result.data
if (mode === 'append') {
@@ -365,7 +409,14 @@ export function ImportCsvDialog({
</div>
<div className='flex flex-col gap-2'>
<Label>Column mapping</Label>
<div className='flex items-center justify-between'>
<Label>Column mapping</Label>
{skipCount > 0 && (
<Button variant='ghost' size='sm' onClick={handleCreateAllUnmapped}>
Create columns for {skipCount} unmapped
</Button>
)}
</div>
<div className='overflow-hidden rounded-sm border border-[var(--border)]'>
<div className='max-h-[320px] overflow-auto'>
<Table>
@@ -401,7 +452,11 @@ export function ImportCsvDialog({
<TableCell>
<Combobox
options={columnOptions}
value={mapping[header] ?? SKIP_VALUE}
value={
createHeaders.has(header)
? CREATE_VALUE
: (mapping[header] ?? SKIP_VALUE)
}
onChange={(value) => handleMappingChange(header, value)}
size='sm'
className='w-full'
@@ -415,7 +470,12 @@ export function ImportCsvDialog({
</div>
</div>
<span className='text-[var(--text-tertiary)] text-xs'>
{mappedCount} mapped · {skipCount} skipped
{mappedCount} mapped
{createCount > 0
? ` · ${createCount} new column${createCount === 1 ? '' : 's'}`
: ''}
{' · '}
{skipCount} skipped
</span>
</div>
@@ -1,6 +1,7 @@
'use client'
import {
Download,
DropdownMenu,
DropdownMenuContent,
DropdownMenuItem,
@@ -19,9 +20,11 @@ interface TableContextMenuProps {
onViewSchema?: () => void
onRename?: () => void
onImportCsv?: () => void
onExportCsv?: () => void
disableDelete?: boolean
disableRename?: boolean
disableImport?: boolean
disableExport?: boolean
menuRef?: React.RefObject<HTMLDivElement | null>
}
@@ -34,9 +37,11 @@ export function TableContextMenu({
onViewSchema,
onRename,
onImportCsv,
onExportCsv,
disableDelete = false,
disableRename = false,
disableImport = false,
disableExport = false,
}: TableContextMenuProps) {
return (
<DropdownMenu open={isOpen} onOpenChange={(open) => !open && onClose()} modal={false}>
@@ -78,7 +83,13 @@ export function TableContextMenu({
Import CSV...
</DropdownMenuItem>
)}
{(onViewSchema || onRename || onImportCsv) && (onCopyId || onDelete) && (
{onExportCsv && (
<DropdownMenuItem disabled={disableExport} onSelect={onExportCsv}>
<Download />
Export CSV
</DropdownMenuItem>
)}
{(onViewSchema || onRename || onImportCsv || onExportCsv) && (onCopyId || onDelete) && (
<DropdownMenuSeparator />
)}
{onCopyId && (
@@ -34,6 +34,7 @@ import {
import { TableContextMenu } from '@/app/workspace/[workspaceId]/tables/components/table-context-menu'
import { useContextMenu } from '@/app/workspace/[workspaceId]/w/components/sidebar/hooks'
import {
downloadTableExport,
useCreateTable,
useDeleteTable,
useTablesList,
@@ -530,6 +531,15 @@ export function Tables() {
}}
onDelete={() => setIsDeleteDialogOpen(true)}
onImportCsv={() => setIsImportDialogOpen(true)}
onExportCsv={async () => {
if (!activeTable) return
try {
await downloadTableExport(activeTable.id, activeTable.name)
} catch (err) {
logger.error('Failed to export table:', err)
toast.error('Failed to export table')
}
}}
disableDelete={userPermissions.canEdit !== true}
disableRename={userPermissions.canEdit !== true}
disableImport={userPermissions.canEdit !== true}
+173 -85
View File
@@ -3,7 +3,14 @@
*/
import { createLogger } from '@sim/logger'
import { keepPreviousData, useMutation, useQuery, useQueryClient } from '@tanstack/react-query'
import {
type InfiniteData,
keepPreviousData,
useInfiniteQuery,
useMutation,
useQuery,
useQueryClient,
} from '@tanstack/react-query'
import { toast } from '@/components/emcn'
import { requestJson } from '@/lib/api/client/request'
import type { ContractJsonResponse } from '@/lib/api/contracts'
@@ -58,8 +65,8 @@ export const tableKeys = {
details: () => [...tableKeys.all, 'detail'] as const,
detail: (tableId: string) => [...tableKeys.details(), tableId] as const,
rowsRoot: (tableId: string) => [...tableKeys.detail(tableId), 'rows'] as const,
rows: (tableId: string, paramsKey: string) =>
[...tableKeys.rowsRoot(tableId), paramsKey] as const,
infiniteRows: (tableId: string, paramsKey: string) =>
[...tableKeys.rowsRoot(tableId), 'infinite', paramsKey] as const,
}
type TableRowsParams = Omit<TableRowsQueryInput, 'filter' | 'sort'> &
@@ -88,22 +95,6 @@ type TableRowsDeleteResult = Pick<
'deletedRowIds'
>
function createRowsParamsKey({
limit,
offset,
filter,
sort,
includeTotal,
}: Omit<TableRowsParams, 'workspaceId' | 'tableId'>): string {
return JSON.stringify({
limit,
offset,
filter: filter ?? null,
sort: sort ?? null,
includeTotal: includeTotal ?? true,
})
}
async function fetchTable(
workspaceId: string,
tableId: string,
@@ -140,10 +131,7 @@ async function fetchTableRows({
signal,
})
const { rows, totalCount } = response.data
return {
rows,
totalCount,
}
return { rows, totalCount }
}
function invalidateRowData(queryClient: ReturnType<typeof useQueryClient>, tableId: string) {
@@ -203,37 +191,56 @@ export function useTable(workspaceId: string | undefined, tableId: string | unde
})
}
interface InfiniteTableRowsParams {
workspaceId: string
tableId: string
pageSize: number
filter?: Filter | null
sort?: Sort | null
enabled?: boolean
}
/**
* Fetch rows for a table with pagination/filter/sort.
* Paginated row fetching with `useInfiniteQuery`. Each page requests `pageSize`
* rows at the next offset; `getNextPageParam` returns `undefined` once the last
* page comes back short, signalling end-of-list.
*
* Page 0 includes a server `COUNT(*)`; subsequent pages skip it.
*/
export function useTableRows({
export function useInfiniteTableRows({
workspaceId,
tableId,
limit,
offset,
pageSize,
filter,
sort,
includeTotal,
enabled = true,
}: TableRowsParams & { enabled?: boolean }) {
const paramsKey = createRowsParamsKey({ limit, offset, filter, sort, includeTotal })
}: InfiniteTableRowsParams) {
const paramsKey = JSON.stringify({
pageSize,
filter: filter ?? null,
sort: sort ?? null,
})
return useQuery({
queryKey: tableKeys.rows(tableId, paramsKey),
queryFn: ({ signal }) =>
return useInfiniteQuery({
queryKey: tableKeys.infiniteRows(tableId, paramsKey),
queryFn: ({ pageParam, signal }) =>
fetchTableRows({
workspaceId,
tableId,
limit,
offset,
limit: pageSize,
offset: pageParam,
filter,
sort,
includeTotal,
includeTotal: pageParam === 0,
signal,
}),
initialPageParam: 0,
getNextPageParam: (lastPage, _allPages, lastPageParam) => {
if (lastPage.rows.length < pageSize) return undefined
return lastPageParam + pageSize
},
enabled: Boolean(workspaceId && tableId) && enabled,
staleTime: 30 * 1000, // 30 seconds
placeholderData: keepPreviousData,
staleTime: 30 * 1000,
})
}
@@ -341,22 +348,7 @@ export function useCreateTableRow({ workspaceId, tableId }: RowMutationContext)
const row = response.data.row
if (!row) return
queryClient.setQueriesData<TableRowsResponse>(
{ queryKey: tableKeys.rowsRoot(tableId) },
(old) => {
if (!old) return old
if (old.rows.some((r) => r.id === row.id)) return old
const shifted = old.rows.map((r) =>
r.position >= row.position ? { ...r, position: r.position + 1 } : r
)
const rows: TableRow[] = [...shifted, row].sort((a, b) => a.position - b.position)
return {
...old,
rows,
totalCount: old.totalCount === null ? null : old.totalCount + 1,
}
}
)
reconcileCreatedRow(queryClient, tableId, row)
},
onSettled: () => {
invalidateRowCount(queryClient, workspaceId, tableId)
@@ -364,6 +356,86 @@ export function useCreateTableRow({ workspaceId, tableId }: RowMutationContext)
})
}
/**
* Apply a row-level transformation to all cached infinite row queries for this
* table. Used for cell edits where positions don't change.
*/
function patchCachedRows(
queryClient: ReturnType<typeof useQueryClient>,
tableId: string,
patchRow: (row: TableRow) => TableRow
) {
queryClient.setQueriesData<InfiniteData<TableRowsResponse, number>>(
{ queryKey: tableKeys.rowsRoot(tableId), exact: false },
(old) => {
if (!old) return old
return {
...old,
pages: old.pages.map((page) => ({ ...page, rows: page.rows.map(patchRow) })),
}
}
)
}
/**
* Splice a server-returned new row into the paginated row cache. Bumps the
* `position` of any cached row at or past the new row's position, then inserts
* the row into the overlapping page (or appends to the last page when the
* position lies past everything fetched). `onSettled` invalidation reconciles
* drift after the next refetch.
*/
function reconcileCreatedRow(
queryClient: ReturnType<typeof useQueryClient>,
tableId: string,
row: TableRow
) {
queryClient.setQueriesData<InfiniteData<TableRowsResponse, number>>(
{ queryKey: tableKeys.rowsRoot(tableId), exact: false },
(old) => {
if (!old) return old
if (old.pages.some((p) => p.rows.some((r) => r.id === row.id))) return old
const pages = old.pages.map((page) =>
page.rows.some((r) => r.position >= row.position)
? {
...page,
rows: page.rows.map((r) =>
r.position >= row.position ? { ...r, position: r.position + 1 } : r
),
}
: page
)
let inserted = false
const nextPages = pages.map((page) => {
if (inserted) return page
const last = page.rows[page.rows.length - 1]
const fits = last === undefined || last.position >= row.position
if (!fits) return page
inserted = true
const merged = [...page.rows, row].sort((a, b) => a.position - b.position)
return { ...page, rows: merged }
})
if (!inserted && nextPages.length > 0) {
const lastIdx = nextPages.length - 1
const lastPage = nextPages[lastIdx]
nextPages[lastIdx] = {
...lastPage,
rows: [...lastPage.rows, row].sort((a, b) => a.position - b.position),
}
}
const firstPage = nextPages[0]
if (firstPage && firstPage.totalCount !== null && firstPage.totalCount !== undefined) {
nextPages[0] = { ...firstPage, totalCount: firstPage.totalCount + 1 }
}
return { ...old, pages: nextPages }
}
)
}
type BatchCreateTableRowsParams = Omit<BatchInsertTableRowsBodyInput, 'workspaceId' | 'rows'> & {
rows: Array<Record<string, unknown>>
}
@@ -412,21 +484,12 @@ export function useUpdateTableRow({ workspaceId, tableId }: RowMutationContext)
onMutate: ({ rowId, data }) => {
void queryClient.cancelQueries({ queryKey: tableKeys.rowsRoot(tableId) })
const previousQueries = queryClient.getQueriesData<TableRowsResponse>({
const previousQueries = queryClient.getQueriesData<InfiniteData<TableRowsResponse, number>>({
queryKey: tableKeys.rowsRoot(tableId),
})
queryClient.setQueriesData<TableRowsResponse>(
{ queryKey: tableKeys.rowsRoot(tableId) },
(old) => {
if (!old) return old
return {
...old,
rows: old.rows.map((row) =>
row.id === rowId ? { ...row, data: { ...row.data, ...data } as RowData } : row
),
}
}
patchCachedRows(queryClient, tableId, (row) =>
row.id === rowId ? { ...row, data: { ...row.data, ...data } as RowData } : row
)
return { previousQueries }
@@ -467,26 +530,17 @@ export function useBatchUpdateTableRows({ workspaceId, tableId }: RowMutationCon
onMutate: ({ updates }) => {
void queryClient.cancelQueries({ queryKey: tableKeys.rowsRoot(tableId) })
const previousQueries = queryClient.getQueriesData<TableRowsResponse>({
const previousQueries = queryClient.getQueriesData<InfiniteData<TableRowsResponse, number>>({
queryKey: tableKeys.rowsRoot(tableId),
})
const updateMap = new Map(updates.map((u) => [u.rowId, u.data]))
queryClient.setQueriesData<TableRowsResponse>(
{ queryKey: tableKeys.rowsRoot(tableId) },
(old) => {
if (!old) return old
return {
...old,
rows: old.rows.map((row) => {
const patch = updateMap.get(row.id)
if (!patch) return row
return { ...row, data: { ...row.data, ...patch } as RowData }
}),
}
}
)
patchCachedRows(queryClient, tableId, (row) => {
const patch = updateMap.get(row.id)
if (!patch) return row
return { ...row, data: { ...row.data, ...patch } as RowData }
})
return { previousQueries }
},
@@ -623,7 +677,7 @@ export function useUpdateTableMetadata({ workspaceId, tableId }: RowMutationCont
}
/**
* Delete a column from a table.
* Restore an archived table.
*/
export function useRestoreTable() {
const queryClient = useQueryClient()
@@ -687,9 +741,11 @@ interface ImportCsvIntoTableParams {
file: File
mode: CsvImportMode
mapping?: CsvHeaderMapping
/** CSV headers to auto-create as new columns on the target table. */
createColumns?: string[]
}
interface ImportCsvIntoTableOutcome {
interface ImportCsvIntoTableResponse {
success: boolean
data?: {
tableId: string
@@ -718,7 +774,8 @@ export function useImportCsvIntoTable() {
file,
mode,
mapping,
}: ImportCsvIntoTableParams): Promise<ImportCsvIntoTableOutcome> => {
createColumns,
}: ImportCsvIntoTableParams): Promise<ImportCsvIntoTableResponse> => {
const formData = new FormData()
formData.append('file', file)
formData.append('workspaceId', workspaceId)
@@ -726,9 +783,12 @@ export function useImportCsvIntoTable() {
if (mapping) {
formData.append('mapping', JSON.stringify(mapping))
}
if (createColumns && createColumns.length > 0) {
formData.append('createColumns', JSON.stringify(createColumns))
}
// boundary-raw-fetch: multipart/form-data CSV upload, requestJson only supports JSON bodies
const response = await fetch(`/api/table/${tableId}/import-csv`, {
const response = await fetch(`/api/table/${tableId}/import`, {
method: 'POST',
body: formData,
})
@@ -750,6 +810,34 @@ export function useImportCsvIntoTable() {
})
}
/**
* Downloads the full contents of a table to the user's device by streaming
* `/api/table/[tableId]/export`. Defaults to CSV; pass `'json'` for JSON.
*/
export async function downloadTableExport(
tableId: string,
fileName: string,
format: 'csv' | 'json' = 'csv'
): Promise<void> {
const url = `/api/table/${tableId}/export?format=${format}&t=${Date.now()}`
// boundary-raw-fetch: streaming download to a Blob, requestJson cannot consume non-JSON streams
const response = await fetch(url, { cache: 'no-store' })
if (!response.ok) {
const data = await response.json().catch(() => ({}))
throw new Error(data.error || `Failed to export table: ${response.statusText}`)
}
const blob = await response.blob()
const objectUrl = URL.createObjectURL(blob)
const safeName = fileName.replace(/[^a-zA-Z0-9_-]+/g, '_').replace(/^_+|_+$/g, '') || 'table'
const a = document.createElement('a')
a.href = objectUrl
a.download = `${safeName}.${format}`
document.body.appendChild(a)
a.click()
document.body.removeChild(a)
URL.revokeObjectURL(objectUrl)
}
export function useDeleteColumn({ workspaceId, tableId }: RowMutationContext) {
const queryClient = useQueryClient()
+37 -1
View File
@@ -10,7 +10,7 @@ import type {
TableRow,
} from '@/lib/table'
import { COLUMN_TYPES, NAME_PATTERN, TABLE_LIMITS } from '@/lib/table/constants'
import { CSV_MAX_FILE_SIZE_BYTES } from '@/lib/table/csv-import'
import { CSV_MAX_FILE_SIZE_BYTES } from '@/lib/table/import'
const isRecord = (value: unknown): value is Record<string, unknown> =>
typeof value === 'object' && value !== null && !Array.isArray(value)
@@ -562,6 +562,42 @@ export const csvExtensionSchema = z.enum(['csv', 'tsv'], {
error: 'Only CSV and TSV files are supported',
})
/**
* `createColumns` form field — a JSON-encoded array of CSV header names that
* the import should auto-create as new columns on the target table.
*/
export const csvImportCreateColumnsSchema = z.unknown().transform((value, ctx): string[] => {
if (typeof value !== 'string') {
ctx.addIssue({ code: 'custom', message: 'createColumns must be valid JSON' })
return z.NEVER
}
try {
const parsed: unknown = JSON.parse(value)
if (!Array.isArray(parsed) || parsed.some((h) => typeof h !== 'string')) {
ctx.addIssue({
code: 'custom',
message: 'createColumns must be a JSON array of CSV header names',
})
return z.NEVER
}
return parsed as string[]
} catch {
ctx.addIssue({ code: 'custom', message: 'createColumns must be valid JSON' })
return z.NEVER
}
})
/**
* `format` query param for the table export route. Lower-cases the input
* before validating against the supported formats and defaults to `'csv'`.
*/
export const tableExportFormatSchema = z
.preprocess(
(value) => (typeof value === 'string' ? value.toLowerCase() : value),
z.enum(['csv', 'json'])
)
.default('csv')
/**
* `mapping` form field — a JSON-encoded `CsvHeaderMapping` (CSV header →
* column name, or `null` to skip the header).
@@ -12,10 +12,10 @@ import {
parseCsvBuffer,
sanitizeName,
validateMapping,
} from '@/lib/table/csv-import'
} from '@/lib/table/import'
import type { TableSchema } from '@/lib/table/types'
describe('csv-import', () => {
describe('import', () => {
describe('parseCsvBuffer', () => {
it('parses a CSV string and extracts headers', async () => {
const { headers, rows } = await parseCsvBuffer('a,b\n1,2\n3,4')
@@ -3,7 +3,7 @@
*
* Used by:
* - `POST /api/table/import-csv` (create new table from CSV)
* - `POST /api/table/[tableId]/import-csv` (append/replace into existing table)
* - `POST /api/table/[tableId]/import` (append/replace into existing table)
* - Copilot `user-table` tool (`create_from_file`, `import_file`)
*
* Keeping a single implementation avoids drift between HTTP and agent code paths.
+1 -1
View File
@@ -7,7 +7,7 @@
export * from './billing'
export * from './constants'
export * from './csv-import'
export * from './import'
export * from './llm'
export * from './query-builder'
export * from './service'
+173 -78
View File
@@ -469,6 +469,81 @@ export async function addTableColumn(
}
}
/**
* Adds multiple columns to an existing table inside a caller-provided
* transaction. This is atomic with respect to the surrounding `trx`: either
* all columns are added or none are. Validates each column the same way
* `addTableColumn` does and rejects if any name collides with an existing
* column or another entry in `columns`.
*
* Use this when composing a column addition with other writes (e.g., row
* inserts) that must succeed or roll back together.
*/
export async function addTableColumnsWithTx(
trx: DbTransaction,
table: TableDefinition,
columns: { name: string; type: string; required?: boolean; unique?: boolean }[],
requestId: string
): Promise<TableDefinition> {
if (columns.length === 0) return table
const usedNames = new Set(table.schema.columns.map((c) => c.name.toLowerCase()))
const additions: TableSchema['columns'] = []
for (const column of columns) {
if (!NAME_PATTERN.test(column.name)) {
throw new Error(
`Invalid column name "${column.name}". Must start with a letter or underscore and contain only alphanumeric characters and underscores.`
)
}
if (column.name.length > TABLE_LIMITS.MAX_COLUMN_NAME_LENGTH) {
throw new Error(
`Column name exceeds maximum length (${TABLE_LIMITS.MAX_COLUMN_NAME_LENGTH} characters)`
)
}
if (!COLUMN_TYPES.includes(column.type as (typeof COLUMN_TYPES)[number])) {
throw new Error(
`Invalid column type "${column.type}". Must be one of: ${COLUMN_TYPES.join(', ')}`
)
}
const lower = column.name.toLowerCase()
if (usedNames.has(lower)) {
throw new Error(`Column "${column.name}" already exists`)
}
usedNames.add(lower)
additions.push({
name: column.name,
type: column.type as TableSchema['columns'][number]['type'],
required: column.required ?? false,
unique: column.unique ?? false,
})
}
if (table.schema.columns.length + additions.length > TABLE_LIMITS.MAX_COLUMNS_PER_TABLE) {
throw new Error(
`Adding ${additions.length} column(s) would exceed maximum column limit (${TABLE_LIMITS.MAX_COLUMNS_PER_TABLE})`
)
}
const updatedSchema: TableSchema = { columns: [...table.schema.columns, ...additions] }
const now = new Date()
await trx
.update(userTableDefinitions)
.set({ schema: updatedSchema, updatedAt: now })
.where(eq(userTableDefinitions.id, table.id))
logger.info(
`[${requestId}] Added ${additions.length} column(s) to table ${table.id}: ${additions.map((c) => c.name).join(', ')}`
)
return {
...table,
schema: updatedSchema,
updatedAt: now,
}
}
/**
* Renames a table.
*
@@ -739,7 +814,25 @@ export async function batchInsertRows(
table: TableDefinition,
requestId: string
): Promise<TableRow[]> {
// Validate all rows
return db.transaction((trx) => batchInsertRowsWithTx(trx, data, table, requestId))
}
/**
* Transaction-bound variant of `batchInsertRows`. Validates rows and unique
* constraints, then performs INSERTs inside the provided transaction. Caller
* is responsible for opening the transaction. Use when row inserts must be
* atomic with other writes (e.g., schema mutations) on the same tx.
*
* Capacity enforcement lives in the `increment_user_table_row_count` trigger
* (migration 0198) — fires per row and raises `Maximum row limit (%) reached ...`
* if the cap is hit mid-batch.
*/
export async function batchInsertRowsWithTx(
trx: DbTransaction,
data: BatchInsertData,
table: TableDefinition,
requestId: string
): Promise<TableRow[]> {
for (let i = 0; i < data.rows.length; i++) {
const row = data.rows[i]
@@ -754,10 +847,14 @@ export async function batchInsertRows(
}
}
// Check unique constraints across all rows using optimized database query
const uniqueColumns = getUniqueColumns(table.schema)
if (uniqueColumns.length > 0) {
const uniqueResult = await checkBatchUniqueConstraintsDb(data.tableId, data.rows, table.schema)
const uniqueResult = await checkBatchUniqueConstraintsDb(
data.tableId,
data.rows,
table.schema,
trx
)
if (!uniqueResult.valid) {
const errorMessages = uniqueResult.errors
.map((e) => `Row ${e.row + 1}: ${e.errors.join(', ')}`)
@@ -768,52 +865,41 @@ export async function batchInsertRows(
const now = new Date()
// Capacity enforcement lives in the `increment_user_table_row_count` trigger
// (migration 0198) — fires per row and raises `Maximum row limit (%) reached ...`
// if the cap is hit mid-batch. The outer transaction means a partial batch
// rolls back cleanly.
const insertedRows = await db.transaction(async (trx) => {
await setTableTxTimeouts(trx, { statementMs: 60_000 })
await setTableTxTimeouts(trx, { statementMs: 60_000 })
const buildRow = (rowData: RowData, position: number) => ({
id: `row_${generateId().replace(/-/g, '')}`,
tableId: data.tableId,
workspaceId: data.workspaceId,
data: rowData,
position,
createdAt: now,
updatedAt: now,
...(data.userId ? { createdBy: data.userId } : {}),
})
const buildRow = (rowData: RowData, position: number) => ({
id: `row_${generateId().replace(/-/g, '')}`,
tableId: data.tableId,
workspaceId: data.workspaceId,
data: rowData,
position,
createdAt: now,
updatedAt: now,
...(data.userId ? { createdBy: data.userId } : {}),
})
// Serialize position-aware writes per-table. See `acquireTablePositionLock`
// for why both explicit- and auto-position paths take this lock.
await acquireTablePositionLock(trx, data.tableId)
await acquireTablePositionLock(trx, data.tableId)
if (data.positions && data.positions.length > 0) {
// Position-aware insert: shift existing rows to create gaps, then insert.
// Process positions ascending so each shift preserves gaps created by prior shifts.
// (Descending would cause lower shifts to push higher gaps out of position.)
const sortedPositions = [...data.positions].sort((a, b) => a - b)
let insertedRows
if (data.positions && data.positions.length > 0) {
// Position-aware insert: shift existing rows to create gaps, then insert.
// Process positions ascending so each shift preserves gaps created by prior shifts.
const sortedPositions = [...data.positions].sort((a, b) => a - b)
for (const pos of sortedPositions) {
await trx
.update(userTableRows)
.set({ position: sql`position + 1` })
.where(and(eq(userTableRows.tableId, data.tableId), gte(userTableRows.position, pos)))
}
// Build rows in original input order so RETURNING preserves caller's index correlation
const rowsToInsert = data.rows.map((rowData, i) => buildRow(rowData, data.positions![i]))
return trx.insert(userTableRows).values(rowsToInsert).returning()
for (const pos of sortedPositions) {
await trx
.update(userTableRows)
.set({ position: sql`position + 1` })
.where(and(eq(userTableRows.tableId, data.tableId), gte(userTableRows.position, pos)))
}
const rowsToInsert = data.rows.map((rowData, i) => buildRow(rowData, data.positions![i]))
insertedRows = await trx.insert(userTableRows).values(rowsToInsert).returning()
} else {
const startPos = await nextAutoPosition(trx, data.tableId)
const rowsToInsert = data.rows.map((rowData, i) => buildRow(rowData, startPos + i))
return trx.insert(userTableRows).values(rowsToInsert).returning()
})
insertedRows = await trx.insert(userTableRows).values(rowsToInsert).returning()
}
logger.info(`[${requestId}] Batch inserted ${data.rows.length} rows into table ${data.tableId}`)
@@ -845,6 +931,19 @@ export async function replaceTableRows(
data: ReplaceRowsData,
table: TableDefinition,
requestId: string
): Promise<ReplaceRowsResult> {
return db.transaction((trx) => replaceTableRowsWithTx(trx, data, table, requestId))
}
/**
* Transaction-bound variant of `replaceTableRows`. Caller opens the transaction.
* Use when the replace must be atomic with other writes (e.g., schema mutations).
*/
export async function replaceTableRowsWithTx(
trx: DbTransaction,
data: ReplaceRowsData,
table: TableDefinition,
requestId: string
): Promise<ReplaceRowsResult> {
if (data.tableId !== table.id) {
throw new Error(`Table ID mismatch: ${data.tableId} vs ${table.id}`)
@@ -903,52 +1002,48 @@ export async function replaceTableRows(
perRowMs: 3,
})
const result = await db.transaction(async (trx) => {
await setTableTxTimeouts(trx, { statementMs })
await setTableTxTimeouts(trx, { statementMs })
// Serialize concurrent replaces (and concurrent auto-position inserts) on the
// same table. Without this, two concurrent replaces each see their own MVCC
// snapshot for the DELETE; the second's DELETE would not observe rows the
// first inserted, so both transactions commit and the table ends up with
// the union of both row sets instead of only the last caller's rows.
await acquireTablePositionLock(trx, data.tableId)
// Serialize concurrent replaces (and concurrent auto-position inserts) on the
// same table. Without this, two concurrent replaces each see their own MVCC
// snapshot for the DELETE; the second's DELETE would not observe rows the
// first inserted, so both transactions commit and the table ends up with
// the union of both row sets instead of only the last caller's rows.
await acquireTablePositionLock(trx, data.tableId)
const deletedRows = await trx
.delete(userTableRows)
.where(eq(userTableRows.tableId, data.tableId))
.returning({ id: userTableRows.id })
const deletedRows = await trx
.delete(userTableRows)
.where(eq(userTableRows.tableId, data.tableId))
.returning({ id: userTableRows.id })
let insertedCount = 0
if (data.rows.length > 0) {
const rowsToInsert = data.rows.map((rowData, i) => ({
id: `row_${generateId().replace(/-/g, '')}`,
tableId: data.tableId,
workspaceId: data.workspaceId,
data: rowData,
position: i,
createdAt: now,
updatedAt: now,
...(data.userId ? { createdBy: data.userId } : {}),
}))
let insertedCount = 0
if (data.rows.length > 0) {
const rowsToInsert = data.rows.map((rowData, i) => ({
id: `row_${generateId().replace(/-/g, '')}`,
tableId: data.tableId,
workspaceId: data.workspaceId,
data: rowData,
position: i,
createdAt: now,
updatedAt: now,
...(data.userId ? { createdBy: data.userId } : {}),
}))
const batchSize = TABLE_LIMITS.MAX_BATCH_INSERT_SIZE
for (let i = 0; i < rowsToInsert.length; i += batchSize) {
const chunk = rowsToInsert.slice(i, i + batchSize)
const inserted = await trx.insert(userTableRows).values(chunk).returning({
id: userTableRows.id,
})
insertedCount += inserted.length
}
const batchSize = TABLE_LIMITS.MAX_BATCH_INSERT_SIZE
for (let i = 0; i < rowsToInsert.length; i += batchSize) {
const chunk = rowsToInsert.slice(i, i + batchSize)
const inserted = await trx.insert(userTableRows).values(chunk).returning({
id: userTableRows.id,
})
insertedCount += inserted.length
}
return { deletedCount: deletedRows.length, insertedCount }
})
}
logger.info(
`[${requestId}] Replaced rows in table ${data.tableId}: deleted ${result.deletedCount}, inserted ${result.insertedCount}`
`[${requestId}] Replaced rows in table ${data.tableId}: deleted ${deletedRows.length}, inserted ${insertedCount}`
)
return result
return { deletedCount: deletedRows.length, insertedCount }
}
/**
+14 -2
View File
@@ -379,14 +379,26 @@ export async function checkUniqueConstraintsDb(
return { valid: errors.length === 0, errors }
}
/**
* Minimal executor surface needed by unique-constraint checks. Both `db` and a
* drizzle transaction (`trx`) satisfy this, letting callers run the lookup
* inside an open transaction so it observes uncommitted prior-batch inserts.
*/
type UniqueCheckExecutor = Pick<typeof db, 'select'>
/**
* Checks unique constraints for a batch of rows using targeted database queries.
* Validates both against existing database rows and within the batch itself.
*
* Pass a transaction as `executor` when running inside an open tx so the lookup
* sees rows inserted by earlier batches in the same transaction; otherwise the
* default `db` connection only observes committed state.
*/
export async function checkBatchUniqueConstraintsDb(
tableId: string,
rows: RowData[],
schema: TableSchema
schema: TableSchema,
executor: UniqueCheckExecutor = db
): Promise<{ valid: boolean; errors: Array<{ row: number; errors: string[] }> }> {
const uniqueColumns = getUniqueColumns(schema)
const rowErrors: Array<{ row: number; errors: string[] }> = []
@@ -458,7 +470,7 @@ export async function checkBatchUniqueConstraintsDb(
return sql`(${userTableRows.data}->${sql.raw(`'${columnName}'`)})::jsonb = ${normalizedValue}::jsonb`
})
const conflictingRows = await db
const conflictingRows = await executor
.select({
id: userTableRows.id,
data: userTableRows.data,