From 56a4738bdd2007584e6a8049c53a49fbad5a7ef0 Mon Sep 17 00:00:00 2001 From: vesperships Date: Sun, 27 Sep 2026 23:41:45 -0700 Subject: [PATCH] fix(e2e): stabilize ObsessionDB live e2e flakes on main MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Unblock the obsessiondb CI job (run 36370187217 on 7fe44e5) which failed: - plugin-ingest: journal UNKNOWN_TABLE after CREATE (SharedMergeTree lag) - plugin-backfill mv_replay: target missing early buckets after QueryFinish - cli migrate-dictionary: dictGet empty after seed (stale initial load) - cli text-index: suite 15s timeout vs 20–30s sibling ObsessionDB latency Root causes (recurring across prior main runs; not introduced by #219): 1. createClickHouseJournal.ensure() did not wait for table visibility 2. Dictionary HTTP source races initial load vs post-create INSERT 3. mv_replay assert did not poll for SharedMergeTree insert visibility 4. Two text-index tests inherited 15s suite timeout without overrides Fixes: waitForTable after journal CREATE; SYSTEM RELOAD DICTIONARY + waitForRows after seed/replace; waitForRows for mv_replay buckets; 120s/60s timeouts on the two text-index tests. Cannot verify live ObsessionDB here (no secrets). Unit/typecheck/lint OK. Related: #218 batches the escape-check queries (complementary). --- .../src/test/migrate-dictionary.e2e.test.ts | 25 ++++++++++++++----- packages/cli/src/test/text-index.e2e.test.ts | 6 +++-- .../src/mv-replay-plan.e2e.test.ts | 14 ++++++++++- packages/plugin-ingest/src/journal.ts | 5 +++- 4 files changed, 40 insertions(+), 10 deletions(-) diff --git a/packages/cli/src/test/migrate-dictionary.e2e.test.ts b/packages/cli/src/test/migrate-dictionary.e2e.test.ts index 9568ba04..6806e7bc 100644 --- a/packages/cli/src/test/migrate-dictionary.e2e.test.ts +++ b/packages/cli/src/test/migrate-dictionary.e2e.test.ts @@ -14,6 +14,7 @@ import { runCli, runCliWithRetry, waitForDictionary, + waitForRows, waitForTable, } from './e2e-testkit.js' @@ -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') @@ -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') diff --git a/packages/cli/src/test/text-index.e2e.test.ts b/packages/cli/src/test/text-index.e2e.test.ts index f5571f8d..9f05e239 100644 --- a/packages/cli/src/test/text-index.e2e.test.ts +++ b/packages/cli/src/test/text-index.e2e.test.ts @@ -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) @@ -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) diff --git a/packages/plugin-backfill/src/mv-replay-plan.e2e.test.ts b/packages/plugin-backfill/src/mv-replay-plan.e2e.test.ts index acf35a02..0a385dec 100644 --- a/packages/plugin-backfill/src/mv-replay-plan.e2e.test.ts +++ b/packages/plugin-backfill/src/mv-replay-plan.e2e.test.ts @@ -10,6 +10,7 @@ import { createPrefix, createStatelessLiveExecutor, getRequiredEnv, + waitForRows, waitForTable, } from '@chkit/clickhouse/e2e-testkit' @@ -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) }) diff --git a/packages/plugin-ingest/src/journal.ts b/packages/plugin-ingest/src/journal.ts index 5afe24ba..3ec6c5cf 100644 --- a/packages/plugin-ingest/src/journal.ts +++ b/packages/plugin-ingest/src/journal.ts @@ -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' @@ -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) {