Skip to content

Commit f2d3393

Browse files
ovrclaude
andauthored
feat(prestodb-driver): Tag Presto/Trino queries with requestId trace token (#12080)
PrestoDriver (and so TrinoDriver) now sends the Cube request UUID (extractRequestUUID(requestId)) as X-Trino-Trace-Token / X-Presto-Trace-Token on the initial POST /v1/statement for query(), stream(), downloadQueryResults() and the export-bucket unload path, so Trino/Presto queries can be traced back to the Cube request, like BigQuery job labels and ClickHouse query_id. The token is skipped for ids that fail isValidRequestId, since scheduledRefreshContexts ids are only warned about and an invalid header value would make Node fail the query. BaseDriver.downloadQueryResults now forwards requestId to query(), which also fixes BigQuery: its non-streamImport downloads were missing the cube_request_id label. The Trino and Presto docs pages get a "Query attribution" section. Fixes #12072 Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
1 parent cda0f58 commit f2d3393

6 files changed

Lines changed: 182 additions & 21 deletions

File tree

‎docs-mintlify/admin/connect-to-data/data-sources/presto.mdx‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -121,6 +121,13 @@ To enable SSL-encrypted connections between Cube and Presto, set the
121121
configure custom certificates, please check out [Enable SSL Connections to the
122122
Database][ref-recipe-enable-ssl].
123123

124+
## Query attribution
125+
126+
Cube sends the identifier shown for the query in Query History as the
127+
`X-Presto-Trace-Token` header of every statement it runs, so all statements of one Cube
128+
query share it. Presto records it as the session's trace token, returned as
129+
`session.traceToken` by the query info endpoint (`GET /v1/query/<query-id>`).
130+
124131
## Custom headers
125132

126133
The Presto driver supports forwarding custom HTTP headers on every request to

‎docs-mintlify/admin/connect-to-data/data-sources/trino.mdx‎

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -121,6 +121,15 @@ To enable SSL-encrypted connections between Cube and Trino, set the
121121
configure custom certificates, please check out [Enable SSL Connections to the
122122
Database][ref-recipe-enable-ssl].
123123

124+
## Query attribution
125+
126+
Cube sends the identifier shown for the query in Query History as the
127+
`X-Trino-Trace-Token` header of every statement it runs, so all statements of one Cube
128+
query share it. Trino records it as the session's trace token, returned as
129+
`session.traceToken` by the query info endpoint (`GET /v1/query/<query-id>`).
130+
It is also passed to [event listeners][trino-docs-event-listener] as the query context's
131+
`traceToken`.
132+
124133
## Custom headers
125134

126135
The Trino driver supports forwarding custom HTTP headers (e.g., `X-Trino-Source`,
@@ -209,4 +218,5 @@ module.exports = {
209218
[trino-docs-approx-agg-fns]:
210219
https://trino.io/docs/current/functions/aggregate.html#approximate-aggregate-functions
211220
[ref-recipe-enable-ssl]: /recipes/configuration/using-ssl-connections-to-data-source
212-
[ref-schema-ref-types-formats-countdistinctapprox]: /reference/data-modeling/measures#type
221+
[ref-schema-ref-types-formats-countdistinctapprox]: /reference/data-modeling/measures#type
222+
[trino-docs-event-listener]: https://trino.io/docs/current/develop/event-listener.html

‎packages/cubejs-base-driver/src/BaseDriver.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -296,8 +296,8 @@ export abstract class BaseDriver implements DriverInterface {
296296
throw new TypeError('Driver\'s .streamQuery() method is not implemented yet.');
297297
}
298298

299-
public async downloadQueryResults(query: string, values: unknown[], _options: DownloadQueryResultsOptions): Promise<DownloadQueryResultsResult> {
300-
const rows = await this.query<Row>(query, values);
299+
public async downloadQueryResults(query: string, values: unknown[], options: DownloadQueryResultsOptions): Promise<DownloadQueryResultsResult> {
300+
const rows = await this.query<Row>(query, values, { requestId: options?.requestId });
301301
const types = detectTypesFromTabular(rows);
302302

303303
return {

‎packages/cubejs-prestodb-driver/src/PrestoDriver.ts‎

Lines changed: 32 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@
77
import {
88
DownloadQueryResultsOptions, DownloadQueryResultsResult,
99
DriverCapabilities, DriverInterface,
10+
QueryOptions,
1011
StreamOptions,
1112
StreamTableData,
1213
TableStructure,
@@ -17,6 +18,8 @@ import {
1718
getEnv,
1819
assertDataSource,
1920
formatAnsi,
21+
extractRequestUUID,
22+
isValidRequestId,
2023
} from '@cubejs-backend/shared';
2124

2225
import { Transform, TransformCallback } from 'stream';
@@ -82,6 +85,8 @@ export class PrestoDriver extends BaseDriver implements DriverInterface {
8285

8386
protected useSelectTestConnection: boolean;
8487

88+
private readonly traceTokenHeader: string;
89+
8590
/**
8691
* Class constructor.
8792
*/
@@ -127,6 +132,7 @@ export class PrestoDriver extends BaseDriver implements DriverInterface {
127132
...config
128133
};
129134
this.catalog = this.config.catalog;
135+
this.traceTokenHeader = this.config.engine === 'trino' ? 'X-Trino-Trace-Token' : 'X-Presto-Trace-Token';
130136
this.client = new presto.Client({
131137
timeout: this.config.queryTimeout,
132138
engine: 'presto',
@@ -212,15 +218,19 @@ export class PrestoDriver extends BaseDriver implements DriverInterface {
212218
await this.queryPromised('SELECT 1', false);
213219
}
214220

215-
public query(query: string, values: unknown[]): Promise<any[]> {
216-
return <Promise<any[]>> this.queryPromised(this.prepareQueryWithParams(query, values), false);
221+
public query(query: string, values: unknown[], options?: QueryOptions): Promise<any[]> {
222+
return <Promise<any[]>> this.queryPromised(this.prepareQueryWithParams(query, values), false, options?.requestId);
217223
}
218224

219225
protected prepareQueryWithParams(query: string, values: unknown[]) {
220226
return formatAnsi(query, values || []);
221227
}
222228

223-
public queryPromised(query: string, streaming: boolean): Promise<any[] | StreamTableData> {
229+
public queryPromised(query: string, streaming: boolean, requestId?: string): Promise<any[] | StreamTableData> {
230+
const traceToken = requestId && extractRequestUUID(requestId);
231+
const headers = traceToken && isValidRequestId(traceToken)
232+
? { ...this.config.headers, [this.traceTokenHeader]: traceToken }
233+
: this.config.headers;
224234
const toError = (error: any) => new Error(error.error ? `${error.message}\n${error.error}` : error.message);
225235
if (streaming) {
226236
const rowStream = new Transform({
@@ -236,7 +246,7 @@ export class PrestoDriver extends BaseDriver implements DriverInterface {
236246
this.client.execute({
237247
query,
238248
schema: this.config.schema || 'default',
239-
headers: this.config.headers,
249+
headers,
240250
session: this.config.queryTimeout ? `query_max_run_time=${this.config.queryTimeout}s` : undefined,
241251
columns: (error: any, columns: TableStructure) => {
242252
resolve({
@@ -266,7 +276,7 @@ export class PrestoDriver extends BaseDriver implements DriverInterface {
266276
this.client.execute({
267277
query,
268278
schema: this.config.schema || 'default',
269-
headers: this.config.headers,
279+
headers,
270280
data: (error: any, data: any[], columns: TableStructure) => {
271281
const normalData = this.normalizeResultOverColumns(data, columns);
272282
fullData = concat(normalData, fullData);
@@ -347,10 +357,10 @@ export class PrestoDriver extends BaseDriver implements DriverInterface {
347357
return map(arrayToObject, data || []);
348358
}
349359

350-
public stream(query: string, values: unknown[], _options: StreamOptions): Promise<StreamTableData> {
360+
public stream(query: string, values: unknown[], options: StreamOptions): Promise<StreamTableData> {
351361
const queryWithParams = this.prepareQueryWithParams(query, values);
352362

353-
return <Promise<StreamTableData>> this.queryPromised(queryWithParams, true);
363+
return <Promise<StreamTableData>> this.queryPromised(queryWithParams, true, options?.requestId);
354364
}
355365

356366
public capabilities(): DriverCapabilities {
@@ -381,8 +391,8 @@ export class PrestoDriver extends BaseDriver implements DriverInterface {
381391
}
382392

383393
const types = options.query
384-
? await this.unloadWithSql(tableName, options.query.sql, options.query.params)
385-
: await this.unloadWithTable(tableName);
394+
? await this.unloadWithSql(tableName, options.query.sql, options.query.params, options.requestId)
395+
: await this.unloadWithTable(tableName, options.requestId);
386396

387397
const csvFile = await this.getCsvFiles(tableName);
388398

@@ -403,33 +413,36 @@ export class PrestoDriver extends BaseDriver implements DriverInterface {
403413
return types.map((c) => `CAST(${c.name} AS varchar) ${c.name}`).join(', ');
404414
}
405415

406-
private async unloadWithSql(tableFullName: string, sql: string, params: any[]) {
416+
private async unloadWithSql(tableFullName: string, sql: string, params: any[], requestId?: string) {
407417
return this.unloadGeneric({
408418
tableFullName,
409419
typeSql: sql,
410420
typeParams: params,
411421
fromSql: sql,
412-
fromParams: params
422+
fromParams: params,
423+
requestId,
413424
});
414425
}
415426

416-
private async unloadWithTable(tableFullName: string) {
427+
private async unloadWithTable(tableFullName: string, requestId?: string) {
417428
return this.unloadGeneric({
418429
tableFullName,
419430
typeSql: `SELECT * FROM ${tableFullName}`,
420431
typeParams: [],
421432
fromSql: tableFullName,
422-
fromParams: []
433+
fromParams: [],
434+
requestId,
423435
});
424436
}
425437

426-
private async unloadGeneric(params: { tableFullName: string, typeSql: string, typeParams: any[], fromSql: string, fromParams: any[] }) {
438+
private async unloadGeneric(params: { tableFullName: string, typeSql: string, typeParams: any[], fromSql: string, fromParams: any[], requestId?: string }) {
427439
if (!this.config.exportBucket) {
428440
throw new Error('Export bucket is not configured.');
429441
}
430442

431443
const { bucketType, exportBucket } = this.config;
432-
const types = await this.queryColumnTypes(params.typeSql, params.typeParams);
444+
const { requestId } = params;
445+
const types = await this.queryColumnTypes(params.typeSql, params.typeParams, { requestId });
433446

434447
const { schema, tableName } = this.splitTableFullName(params.tableFullName);
435448
const tableWithCatalogAndSchema = `${this.config.catalog}.${schema}.${tableName}`;
@@ -448,16 +461,17 @@ export class PrestoDriver extends BaseDriver implements DriverInterface {
448461
await this.query(
449462
createTableQuery,
450463
params.fromParams,
464+
{ requestId },
451465
);
452466
} finally {
453-
await this.query(`DROP TABLE IF EXISTS ${tableWithCatalogAndSchema}`, []);
467+
await this.query(`DROP TABLE IF EXISTS ${tableWithCatalogAndSchema}`, [], { requestId });
454468
}
455469

456470
return types;
457471
}
458472

459-
public async queryColumnTypes(sql: string, params: unknown[]): Promise<{ name: string; type: string; }[]> {
460-
const response = await this.stream(`${sql} LIMIT 0`, params || [], { highWaterMark: 1 });
473+
public async queryColumnTypes(sql: string, params: unknown[], options?: QueryOptions): Promise<{ name: string; type: string; }[]> {
474+
const response = await this.stream(`${sql} LIMIT 0`, params || [], { highWaterMark: 1, requestId: options?.requestId });
461475
const result = [];
462476

463477
for (const column of response.types || []) {

‎packages/cubejs-prestodb-driver/test/unit/headers.test.ts‎

Lines changed: 114 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -129,4 +129,118 @@ describe('PrestoDriver custom headers', () => {
129129
expect(poll!.headers['X-Custom-Header']).toBe('custom-value');
130130
expect(poll!.headers['Proxy-Authorization']).toBe('Basic dGVzdA==');
131131
});
132+
133+
describe('requestId trace token', () => {
134+
const createDriver = (config: Record<string, unknown> = {}) => new PrestoDriver({
135+
host: 'coordinator.local',
136+
port: '8080',
137+
catalog: 'test',
138+
schema: 'default',
139+
dataSource: 'default',
140+
checkInterval: 1,
141+
...config,
142+
} as any);
143+
144+
it('tags the initial POST of query() with the request UUID', async () => {
145+
const driver = createDriver({ headers: { 'X-Custom-Header': 'custom-value' } });
146+
147+
await driver.query('SELECT 1', [], { requestId: 'abc-123-span-2' });
148+
149+
const post = mockRecorded.find((r) => r.method === 'POST');
150+
const poll = mockRecorded.find((r) => r.method === 'GET');
151+
152+
expect(post!.headers['X-Presto-Trace-Token']).toBe('abc-123');
153+
expect(post!.headers['X-Custom-Header']).toBe('custom-value');
154+
expect(poll!.headers['X-Presto-Trace-Token']).toBeUndefined();
155+
expect(poll!.headers['X-Custom-Header']).toBe('custom-value');
156+
});
157+
158+
it('tags the initial POST of stream()', async () => {
159+
const driver = createDriver();
160+
161+
const { rowStream } = await driver.stream('SELECT 1', [], { highWaterMark: 1, requestId: 'stream-req-span-1' });
162+
163+
for await (const _row of rowStream) {
164+
// drain
165+
}
166+
167+
const post = mockRecorded.find((r) => r.method === 'POST');
168+
expect(post!.headers['X-Presto-Trace-Token']).toBe('stream-req');
169+
});
170+
171+
it('uses the Trino header name for the trino engine', async () => {
172+
const driver = createDriver({ engine: 'trino' });
173+
174+
await driver.query('SELECT 1', [], { requestId: 'trino-req' });
175+
176+
const post = mockRecorded.find((r) => r.method === 'POST');
177+
expect(post!.headers['X-Trino-Trace-Token']).toBe('trino-req');
178+
expect(post!.headers['X-Presto-Trace-Token']).toBeUndefined();
179+
});
180+
181+
it('skips the trace token for an invalid requestId and still runs the query', async () => {
182+
const driver = createDriver();
183+
184+
for (const requestId of ['bad\u0001id', 'refresh-Ж', 'a'.repeat(129)]) {
185+
mockRecorded.length = 0;
186+
const rows = await driver.query('SELECT 1', [], { requestId });
187+
expect(rows).toEqual([{ one: 1 }]);
188+
189+
const post = mockRecorded.find((r) => r.method === 'POST');
190+
expect(post!.headers['X-Presto-Trace-Token']).toBeUndefined();
191+
}
192+
});
193+
194+
it('tags downloadQueryResults() without streamImport', async () => {
195+
const driver = createDriver();
196+
197+
await driver.downloadQueryResults('SELECT 1', [], { highWaterMark: 1, requestId: 'download-req-span-1' });
198+
199+
const post = mockRecorded.find((r) => r.method === 'POST');
200+
expect(post!.headers['X-Presto-Trace-Token']).toBe('download-req');
201+
});
202+
203+
it('tags every statement of an export bucket unload', async () => {
204+
class TestPrestoDriver extends PrestoDriver {
205+
protected override async extractUnloadedFilesFromS3(): Promise<string[]> {
206+
return [];
207+
}
208+
}
209+
210+
const driver = new TestPrestoDriver({
211+
host: 'coordinator.local',
212+
port: '8080',
213+
catalog: 'test',
214+
schema: 'default',
215+
dataSource: 'default',
216+
checkInterval: 1,
217+
bucketType: 's3',
218+
exportBucket: 'bucket',
219+
} as any);
220+
221+
await driver.unload('stb.orders', {
222+
maxFileSize: 64,
223+
query: { sql: 'SELECT 1 AS one', params: [] },
224+
requestId: 'unload-req-span-1',
225+
});
226+
227+
const posts = mockRecorded.filter((r) => r.method === 'POST');
228+
// column type probe, CREATE TABLE ... AS, DROP TABLE
229+
expect(posts).toHaveLength(3);
230+
231+
for (const post of posts) {
232+
expect(post.headers['X-Presto-Trace-Token']).toBe('unload-req');
233+
}
234+
});
235+
236+
it('sends no trace token without a requestId', async () => {
237+
const driver = createDriver();
238+
239+
await driver.query('SELECT 1', []);
240+
241+
const post = mockRecorded.find((r) => r.method === 'POST');
242+
expect(post!.headers['X-Presto-Trace-Token']).toBeUndefined();
243+
expect(post!.headers['X-Trino-Trace-Token']).toBeUndefined();
244+
});
245+
});
132246
});

‎packages/cubejs-trino-driver/test/unit/headers.test.ts‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -85,4 +85,20 @@ describe('TrinoDriver headers', () => {
8585
'X-Trino-Routing-Group': 'etl',
8686
});
8787
});
88+
89+
it('adds X-Trino-Trace-Token from requestId on query()', async () => {
90+
const driver = new TrinoDriver({
91+
host: 'trino.local',
92+
port: '8080',
93+
headers: { 'X-Trino-Source': 'cube' },
94+
});
95+
96+
await driver.query('SELECT 1', [], { requestId: 'req-42-span-3' });
97+
98+
const [executeOpts] = mockExecute.mock.calls[0];
99+
expect(executeOpts.headers).toEqual({
100+
'X-Trino-Source': 'cube',
101+
'X-Trino-Trace-Token': 'req-42',
102+
});
103+
});
88104
});

0 commit comments

Comments
 (0)