Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 19 additions & 6 deletions packages/cli/src/test/migrate-dictionary.e2e.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ import {
runCli,
runCliWithRetry,
waitForDictionary,
waitForRows,
waitForTable,
} from './e2e-testkit.js'

Expand Down Expand Up @@ -150,9 +151,16 @@ describe('@chkit/cli migrate dictionary e2e', () => {
await executor.command(
`INSERT INTO ${quoteIdent(database)}.${quoteIdent(tableName)} (id, name) VALUES (1, 'Alice')`
)

const seeded = await executor.query<{ name: string }>(
`SELECT dictGet('${database}.${dictName}', 'name', toUInt64(1)) AS name`
// Dictionary may have finished its initial HTTP load against the empty
// source; reload and poll until seed data is visible via dictGet.
await executor.command(
`SYSTEM RELOAD DICTIONARY ${quoteIdent(database)}.${quoteIdent(dictName)}`
)
const seeded = await waitForRows<{ name: string }>(
executor,
`SELECT dictGet('${database}.${dictName}', 'name', toUInt64(1)) AS name`,
(rows) => rows[0]?.name === 'Alice',
'dictionary seeded Alice',
)
expect(seeded[0]?.name).toBe('Alice')

Expand All @@ -178,9 +186,14 @@ describe('@chkit/cli migrate dictionary e2e', () => {
throw new Error(formatTestDiagnostic('migrate --execute (replace) failed', migrateReplace))
}
await waitForDictionary(executor, database, dictName)

const afterReplace = await executor.query<{ name: string }>(
`SELECT dictGet('${database}.${dictName}', 'name', toUInt64(1)) AS name`
await executor.command(
`SYSTEM RELOAD DICTIONARY ${quoteIdent(database)}.${quoteIdent(dictName)}`
)
const afterReplace = await waitForRows<{ name: string }>(
executor,
`SELECT dictGet('${database}.${dictName}', 'name', toUInt64(1)) AS name`,
(rows) => rows[0]?.name === 'Alice',
'dictionary still serves Alice after replace',
)
expect(afterReplace[0]?.name).toBe('Alice')

Expand Down
6 changes: 4 additions & 2 deletions packages/cli/src/test/text-index.e2e.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -248,7 +248,8 @@ test('normalization preserves every printable ClickHouse string escape', async (
} finally {
await executor.close()
}
})
// Suite default is 15s; ~94 sequential remote queries routinely exceed that on ObsessionDB.
}, 120_000)

test('quoted literal names remain distinct from constants in ClickHouse and planning', async () => {
const executor = createLiveExecutor(env)
Expand Down Expand Up @@ -279,7 +280,8 @@ test('quoted literal names remain distinct from constants in ClickHouse and plan
} finally {
await executor.close()
}
})
// Suite default is 15s; many sequential remote queries need headroom under parallel load.
}, 60_000)

test('quoted NULL column round-trips and changing it to a literal migrates the index', async () => {
const executor = createLiveExecutor(env)
Expand Down
14 changes: 13 additions & 1 deletion packages/plugin-backfill/src/mv-replay-plan.e2e.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import {
createPrefix,
createStatelessLiveExecutor,
getRequiredEnv,
waitForRows,
waitForTable,
} from '@chkit/clickhouse/e2e-testkit'

Expand Down Expand Up @@ -210,9 +211,20 @@ describe('e2e: mv_replay backfill of an empty aggregate target (chkit#187)', ()
expect(result.completed).toBe(plan.chunkPlan.chunks.length)

// Per-bucket values must match a forward run of the MV over the whole source.
// Poll: SharedMergeTree can report QueryFinish before every partition is
// visible to a subsequent SELECT, even with select_sequential_consistency.
const expected = await aggregateByBucket(sourceFqn, 'sum(id)')
const actual = await aggregateByBucket(targetFqn, 'sum(total)')
expect(expected).toHaveLength(BUCKETS)
const actual = await waitForRows<{ bucket: string; total: string }>(
ddl,
`SELECT toString(bucket) AS bucket, toString(sum(total)) AS total
FROM ${targetFqn}
GROUP BY bucket
ORDER BY bucket
SETTINGS select_sequential_consistency = 1`,
(rows) => rows.length === BUCKETS,
'mv_replay target buckets visible',
)
expect(actual).toEqual(expected)
}, 180_000)
})
5 changes: 4 additions & 1 deletion packages/plugin-ingest/src/journal.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { createHash } from 'node:crypto'

import type { ClickHouseExecutor } from '@chkit/clickhouse'
import { waitForTable, type ClickHouseExecutor } from '@chkit/clickhouse'

import { IngestConfigError } from './errors.js'
import type { CheckpointEnvelope, CommittedCheckpoint, Journal, JournalEvent } from './types.js'
Expand Down Expand Up @@ -57,6 +57,9 @@ export function createClickHouseJournal(options: ClickHouseJournalOptions): Jour
return {
async ensure() {
await options.executor.command(journalTableSql(qualified))
// Managed ClickHouse (ObsessionDB SharedMergeTree) can acknowledge CREATE
// before the table is visible on the replica that serves the next query.
await waitForTable(options.executor, options.database, table)
},

async append(event) {
Expand Down
Loading