From ac8d25de235e161f822900f8271e9fe573510a7c Mon Sep 17 00:00:00 2001 From: vesperships Date: Sun, 27 Sep 2026 19:12:15 -0700 Subject: [PATCH] fix: use clickhouse.cluster for async migration polling (#207) Hardcoded clusterAllReplicas('cluster', ...) broke async mode when clickhouse.cluster was set to a real name. Add optional systemTableSource on ClickHouseExecutor, shared observableSystemTable helper, and wire queryStatus + backfill syncProgress through it. --- .../runtime/plugin-runtime/executor-debug.ts | 3 + packages/clickhouse/src/index.test.ts | 128 ++++++++++++++++++ packages/clickhouse/src/index.ts | 33 ++++- .../src/async-backfill.test.ts | 67 +++++++++ .../plugin-backfill/src/async-backfill.ts | 9 +- 5 files changed, 235 insertions(+), 5 deletions(-) diff --git a/packages/cli/src/runtime/plugin-runtime/executor-debug.ts b/packages/cli/src/runtime/plugin-runtime/executor-debug.ts index fd32d2c2..f743e361 100644 --- a/packages/cli/src/runtime/plugin-runtime/executor-debug.ts +++ b/packages/cli/src/runtime/plugin-runtime/executor-debug.ts @@ -33,7 +33,10 @@ export function wrapExecutorWithDebug( if (!isDebugEnabled()) return executor const queryJson = executor.queryJson?.bind(executor) + const systemTableSource = executor.systemTableSource?.bind(executor) + return { + ...(systemTableSource ? { systemTableSource } : {}), command(sql: string): Promise { debug('clickhouse', `command: ${truncate(sql)}`) return trace('command', null, () => executor.command(sql)) diff --git a/packages/clickhouse/src/index.test.ts b/packages/clickhouse/src/index.test.ts index c6c45a39..24b11ef8 100644 --- a/packages/clickhouse/src/index.test.ts +++ b/packages/clickhouse/src/index.test.ts @@ -10,6 +10,7 @@ import { createStatelessClickHouseClient, formatConnectionError, inferSchemaKindFromEngine, + observableSystemTable, parseCommentFromCreateDictionaryQuery, parseDictionaryAttributesFromCreateDictionaryQuery, parseDictionaryPrimaryKeyFromCreateDictionaryQuery, @@ -705,3 +706,130 @@ RANGE(MIN start_date MAX end_date)` ) }) }) + + +describe('systemTableSource / observableSystemTable', () => { + test('named cluster uses clusterAllReplicas with that name', () => { + const executor = createExecutorWithClient( + { + url: 'http://localhost:8123', + database: 'default', + cluster: 'moebel_cluster', + }, + createMockClient('status', []), + ) + + expect(executor.systemTableSource?.('processes')).toBe( + "clusterAllReplicas('moebel_cluster', system.processes)", + ) + expect(executor.systemTableSource?.('query_log')).toBe( + "clusterAllReplicas('moebel_cluster', system.query_log)", + ) + expect(observableSystemTable(executor, 'processes')).toBe( + "clusterAllReplicas('moebel_cluster', system.processes)", + ) + }) + + test('no cluster uses local system tables (does not emit literal cluster)', () => { + const executor = createExecutorWithClient( + { + url: 'http://localhost:8123', + database: 'default', + }, + createMockClient('status', []), + ) + + expect(executor.systemTableSource?.('processes')).toBe('system.processes') + expect(executor.systemTableSource?.('query_log')).toBe('system.query_log') + expect(observableSystemTable(executor, 'processes')).toBe('system.processes') + expect(executor.systemTableSource?.('processes')).not.toContain("'cluster'") + expect(observableSystemTable(executor, 'query_log')).not.toContain( + "clusterAllReplicas('cluster'", + ) + }) + + test('macro {cluster} is preserved in clusterAllReplicas', () => { + const executor = createExecutorWithClient( + { + url: 'http://localhost:8123', + database: 'default', + cluster: '{cluster}', + }, + createMockClient('status', []), + ) + + expect(executor.systemTableSource?.('processes')).toBe( + "clusterAllReplicas('{cluster}', system.processes)", + ) + }) + + test("literal name 'cluster' still works when configured as such", () => { + const executor = createExecutorWithClient( + { + url: 'http://localhost:8123', + database: 'default', + cluster: 'cluster', + }, + createMockClient('status', []), + ) + + expect(executor.systemTableSource?.('processes')).toBe( + "clusterAllReplicas('cluster', system.processes)", + ) + }) + + test('observableSystemTable falls back to local tables when capability is absent', () => { + expect(observableSystemTable({}, 'processes')).toBe('system.processes') + expect(observableSystemTable({}, 'query_log')).toBe('system.query_log') + }) + + test('queryStatus SQL uses configured cluster (and local tables when unset)', async () => { + const queries: string[] = [] + const capturingClient = { + async command() { + return { query_id: 'cmd', response_headers: {} } + }, + async query(params: { query: string }) { + queries.push(params.query) + return { + query_id: 'q', + response_headers: {}, + async json() { + return [] + }, + } + }, + async insert() { + return { query_id: 'i', response_headers: {} } + }, + async close() {}, + } as unknown as ReturnType + + const clustered = createExecutorWithClient( + { + url: 'http://localhost:8123', + database: 'default', + cluster: 'moebel_cluster', + }, + capturingClient, + ) + await clustered.queryStatus('qid-1') + expect(queries[0]).toContain("clusterAllReplicas('moebel_cluster', system.processes)") + expect(queries[1]).toContain("clusterAllReplicas('moebel_cluster', system.query_log)") + expect(queries.join('\n')).not.toContain("clusterAllReplicas('cluster'") + + queries.length = 0 + const local = createExecutorWithClient( + { + url: 'http://localhost:8123', + database: 'default', + }, + capturingClient, + ) + await local.queryStatus('qid-2') + expect(queries[0]).toContain('FROM system.processes') + expect(queries[0]).not.toContain('clusterAllReplicas') + expect(queries[1]).toContain('FROM system.query_log') + expect(queries[1]).not.toContain('clusterAllReplicas') + }) +}) diff --git a/packages/clickhouse/src/index.ts b/packages/clickhouse/src/index.ts index 29c14dc5..58eb6640 100644 --- a/packages/clickhouse/src/index.ts +++ b/packages/clickhouse/src/index.ts @@ -62,6 +62,8 @@ export interface ClickHouseJsonQueryResult< query_id?: string } +export type ObservableSystemTable = 'processes' | 'query_log' + export interface ClickHouseExecutor { command(sql: string): Promise query(sql: string, settings?: ClickHouseSettings): Promise @@ -92,9 +94,25 @@ export interface ClickHouseExecutor { options?: { afterTime?: string }, ): Promise + /** Optional capability: FROM-clause source for cluster-aware system table polling. + * Native executors return `clusterAllReplicas(...)` when `clickhouse.cluster` is set. + * Remote / ObsessionDB executors may omit this and keep local `system.*` tables. */ + systemTableSource?(table: ObservableSystemTable): string + close(): Promise } +/** + * Resolve the system table expression used for async query observation. + * Prefers `executor.systemTableSource` when present; otherwise local `system.`. + */ +export function observableSystemTable( + executor: Pick, + table: ObservableSystemTable, +): string { + return executor.systemTableSource?.(table) ?? `system.${table}` +} + export interface SchemaObjectRef { kind: 'table' | 'view' | 'materialized_view' | 'dictionary' database: string @@ -767,13 +785,23 @@ export function createExecutorWithClient( } return id }, + systemTableSource(table: ObservableSystemTable): string { + // `config.cluster` is validated at resolveConfig (assertValidClusterName), so it is + // safe to interpolate into the single-quoted clusterAllReplicas argument — same + // contract as onClusterClause. Never guess a default name like 'cluster'. + if (config.cluster) { + return `clusterAllReplicas('${config.cluster}', system.${table})` + } + return `system.${table}` + }, async queryStatus( queryId: string, options?: { afterTime?: string }, ): Promise { try { + const processesFrom = observableSystemTable(this, 'processes') const running = await client.query({ - query: `SELECT read_rows, read_bytes, written_rows, written_bytes, elapsed FROM clusterAllReplicas('cluster', system.processes) WHERE user = currentUser() AND query_id = {qid:String} SETTINGS skip_unavailable_shards = 1`, + query: `SELECT read_rows, read_bytes, written_rows, written_bytes, elapsed FROM ${processesFrom} WHERE user = currentUser() AND query_id = {qid:String} SETTINGS skip_unavailable_shards = 1`, query_params: { qid: queryId }, format: 'JSONEachRow', }) @@ -797,9 +825,10 @@ export function createExecutorWithClient( } const afterTime = options?.afterTime ?? '1970-01-01T00:00:00Z' + const queryLogFrom = observableSystemTable(this, 'query_log') const log = await client.query({ query: `SELECT type, written_rows, written_bytes, query_duration_ms, exception -FROM clusterAllReplicas('cluster', system.query_log) +FROM ${queryLogFrom} WHERE user = currentUser() AND query_id = {qid:String} AND type IN ('QueryFinish', 'ExceptionWhileProcessing') diff --git a/packages/plugin-backfill/src/async-backfill.test.ts b/packages/plugin-backfill/src/async-backfill.test.ts index 7ed13ebe..73591774 100644 --- a/packages/plugin-backfill/src/async-backfill.test.ts +++ b/packages/plugin-backfill/src/async-backfill.test.ts @@ -349,4 +349,71 @@ describe('syncProgress', () => { expect(synced.c1.status).toBe('failed') expect(synced.c1.error).toBe('Memory limit exceeded') }) + + test('uses systemTableSource when executor provides clustered tables', async () => { + const queries: string[] = [] + const executor = createMockExecutor(new Map()) + executor.systemTableSource = (table) => + `clusterAllReplicas('moebel_cluster', system.${table})` + executor.query = async (sql: string): Promise => { + queries.push(sql) + return [] as T[] + } + + await syncProgress( + executor, + PLAN_ID, + [{ id: 'c1' }], + { c1: { status: 'pending' } }, + ) + + expect(queries[0]).toContain("clusterAllReplicas('moebel_cluster', system.processes)") + expect(queries[1]).toContain("clusterAllReplicas('moebel_cluster', system.query_log)") + expect(queries.join('\n')).not.toContain("clusterAllReplicas('cluster'") + }) + + test('falls back to local system tables when capability is absent', async () => { + const queries: string[] = [] + const executor = createMockExecutor(new Map()) + executor.query = async (sql: string): Promise => { + queries.push(sql) + return [] as T[] + } + + await syncProgress( + executor, + PLAN_ID, + [{ id: 'c1' }], + { c1: { status: 'pending' } }, + ) + + expect(queries[0]).toContain('FROM system.processes') + expect(queries[0]).not.toContain('clusterAllReplicas') + expect(queries[1]).toContain('FROM system.query_log') + expect(queries[1]).not.toContain('clusterAllReplicas') + }) + + test('observable fallback still works when systemTableSource is undefined', async () => { + const queries: string[] = [] + const executor = createMockExecutor(new Map()) + // Explicitly no capability (default mock omits it) + expect(executor.systemTableSource).toBeUndefined() + executor.query = async (sql: string): Promise => { + queries.push(sql) + if (sql.includes('system.processes')) { + return [{ query_id: `backfill-${PLAN_ID}-c1` }] as T[] + } + return [] as T[] + } + + const synced = await syncProgress( + executor, + PLAN_ID, + [{ id: 'c1' }], + { c1: { status: 'pending' } }, + ) + + expect(synced.c1.status).toBe('running') + expect(queries[0]).toMatch(/FROM system\.processes/) + }) }) diff --git a/packages/plugin-backfill/src/async-backfill.ts b/packages/plugin-backfill/src/async-backfill.ts index 6c78543b..81d55277 100644 --- a/packages/plugin-backfill/src/async-backfill.ts +++ b/packages/plugin-backfill/src/async-backfill.ts @@ -1,6 +1,6 @@ import { setTimeout as sleep } from 'node:timers/promises' -import type { ClickHouseExecutor, QueryStatus } from '@chkit/clickhouse' +import { observableSystemTable, type ClickHouseExecutor, type QueryStatus } from '@chkit/clickhouse' import pMap from 'p-map' export interface BackfillOptions { @@ -146,8 +146,11 @@ export async function syncProgress( // Escape single-quotes in the prefix for safe SQL embedding const safePrefix = prefix.replace(/'/g, "''").replace(/%/g, '\\%').replace(/_/g, '\\_') + const processesFrom = observableSystemTable(executor, 'processes') + const queryLogFrom = observableSystemTable(executor, 'query_log') + const runningRows = await executor.query<{ query_id: string }>( - `SELECT query_id FROM clusterAllReplicas('cluster', system.processes) WHERE user = currentUser() AND query_id LIKE '${safePrefix}%' SETTINGS skip_unavailable_shards = 1` + `SELECT query_id FROM ${processesFrom} WHERE user = currentUser() AND query_id LIKE '${safePrefix}%' SETTINGS skip_unavailable_shards = 1` ) const runningSet = new Set(runningRows.map((r) => r.query_id)) @@ -160,7 +163,7 @@ export async function syncProgress( exception: string }>( `SELECT query_id, type, written_rows, written_bytes, query_duration_ms, exception -FROM clusterAllReplicas('cluster', system.query_log) +FROM ${queryLogFrom} WHERE user = currentUser() AND query_id LIKE '${safePrefix}%' AND type IN ('QueryFinish', 'ExceptionWhileProcessing')