Skip to content
Merged
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
3 changes: 3 additions & 0 deletions packages/cli/src/runtime/plugin-runtime/executor-debug.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void> {
debug('clickhouse', `command: ${truncate(sql)}`)
return trace('command', null, () => executor.command(sql))
Expand Down
128 changes: 128 additions & 0 deletions packages/clickhouse/src/index.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import {
createStatelessClickHouseClient,
formatConnectionError,
inferSchemaKindFromEngine,
observableSystemTable,
parseCommentFromCreateDictionaryQuery,
parseDictionaryAttributesFromCreateDictionaryQuery,
parseDictionaryPrimaryKeyFromCreateDictionaryQuery,
Expand Down Expand Up @@ -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<typeof createStatelessClickHouseClient>

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')
})
})
33 changes: 31 additions & 2 deletions packages/clickhouse/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,8 @@ export interface ClickHouseJsonQueryResult<
query_id?: string
}

export type ObservableSystemTable = 'processes' | 'query_log'

export interface ClickHouseExecutor {
command(sql: string): Promise<void>
query<T>(sql: string, settings?: ClickHouseSettings): Promise<T[]>
Expand Down Expand Up @@ -92,9 +94,25 @@ export interface ClickHouseExecutor {
options?: { afterTime?: string },
): Promise<QueryStatus>

/** 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<void>
}

/**
* Resolve the system table expression used for async query observation.
* Prefers `executor.systemTableSource` when present; otherwise local `system.<table>`.
*/
export function observableSystemTable(
executor: Pick<ClickHouseExecutor, 'systemTableSource'>,
table: ObservableSystemTable,
): string {
return executor.systemTableSource?.(table) ?? `system.${table}`
}

export interface SchemaObjectRef {
kind: 'table' | 'view' | 'materialized_view' | 'dictionary'
database: string
Expand Down Expand Up @@ -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<QueryStatus> {
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',
})
Expand All @@ -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')
Expand Down
67 changes: 67 additions & 0 deletions packages/plugin-backfill/src/async-backfill.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 <T>(sql: string): Promise<T[]> => {
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 <T>(sql: string): Promise<T[]> => {
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 <T>(sql: string): Promise<T[]> => {
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/)
})
})
9 changes: 6 additions & 3 deletions packages/plugin-backfill/src/async-backfill.ts
Original file line number Diff line number Diff line change
@@ -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 {
Expand Down Expand Up @@ -125,7 +125,7 @@
* to discover queries that were submitted but whose status was never
* persisted locally (e.g. client crash between submit and state write).
*/
export async function syncProgress(

Check warning on line 128 in packages/plugin-backfill/src/async-backfill.ts

View workflow job for this annotation

GitHub Actions / verify

High cognitive complexity (moderate)

Function 'syncProgress' is hard to understand (cognitive: 22, threshold: 15). • Severity: moderate • Cyclomatic: 12 • Cognitive: 22 • Lines: 93 High cognitive complexity means deeply nested or interleaved logic. Consider flattening control flow or extracting helper functions.
executor: ClickHouseExecutor,
planId: string,
chunks: Array<{ id: string }>,
Expand All @@ -146,8 +146,11 @@
// 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))

Expand All @@ -160,7 +163,7 @@
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')
Expand Down
Loading