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
65 changes: 62 additions & 3 deletions handwritten/spanner/observability-test/spanner.ts
Original file line number Diff line number Diff line change
Expand Up @@ -416,6 +416,44 @@
expectedEventNames,
);
});

it('runPartitionedUpdate with string query', async () => {
await database.runPartitionedUpdate(updateSql);

const expectedSpanNames = [
'CloudSpanner.Snapshot.begin',
'CloudSpanner.Snapshot.runStream',
'CloudSpanner.Snapshot.run',
'CloudSpanner.Dml.runUpdate',
'CloudSpanner.PartitionedDml.runUpdate',
'CloudSpanner.Database.runPartitionedUpdate',
];
const expectedEventNames = [
'Begin Transaction',
'Transaction Creation Done',
'Starting stream',
'Acquiring session',
'Cache hit: has usable session',
'Acquired session',
];
verifySpansAndEvents(
traceExporter,
expectedSpanNames,
expectedEventNames,
);

const finishedSpans = traceExporter.getFinishedSpans();
for (const span of finishedSpans) {
const hasNumericAttributeKeys = Object.keys(span.attributes).some(
key => /^\d+$/.test(key),
);
assert.strictEqual(
hasNumericAttributeKeys,
false,
`Span ${span.name} should not contain numeric attributes from string spreading`,
);
}
});
});
});
});
Expand Down Expand Up @@ -843,6 +881,28 @@
expectedEventNames,
`Unexpected events:\n\tGot: ${actualEventNames}\n\tWant: ${expectedEventNames}`,
);

for (const span of spansFromInjected) {
if (
span.name === 'CloudSpanner.Database.run' ||
span.name === 'CloudSpanner.Database.runStream' ||
span.name === 'CloudSpanner.Snapshot.runStream'
) {
assert.strictEqual(
span.attributes['db.statement'],
'SELECT 1',
`Span ${span.name} should have db.statement set to 'SELECT 1'`,
);
const hasNumericAttributeKeys = Object.keys(span.attributes).some(
key => /^\d+$/.test(key),
);
assert.strictEqual(
hasNumericAttributeKeys,
false,
`Span ${span.name} should not contain numeric attributes from string spreading`,
);
}
}
} catch (err) {
assert.ifError(err);
} finally {
Expand Down Expand Up @@ -1683,10 +1743,10 @@
database.run(selectSql, (err, rows) => {
assert.ifError(err);
assert.strictEqual(rows!.length, 3);
database

Check warning on line 1746 in handwritten/spanner/observability-test/spanner.ts

View workflow job for this annotation

GitHub Actions / lint

Avoid using promises inside of callbacks

Check warning on line 1746 in handwritten/spanner/observability-test/spanner.ts

View workflow job for this annotation

GitHub Actions / lint

Avoid using promises inside of callbacks
.close()
.then(() => done())

Check warning on line 1748 in handwritten/spanner/observability-test/spanner.ts

View workflow job for this annotation

GitHub Actions / lint

Avoid calling back inside of a promise
.catch(err => done(err));

Check warning on line 1749 in handwritten/spanner/observability-test/spanner.ts

View workflow job for this annotation

GitHub Actions / lint

Avoid calling back inside of a promise
});
});

Expand All @@ -1703,10 +1763,10 @@
database.run(selectSql, err => {
assert.ok(err, 'Missing expected error');
assert.strictEqual(err!.message, '2 UNKNOWN: Non-retryable error');
database

Check warning on line 1766 in handwritten/spanner/observability-test/spanner.ts

View workflow job for this annotation

GitHub Actions / lint

Avoid using promises inside of callbacks

Check warning on line 1766 in handwritten/spanner/observability-test/spanner.ts

View workflow job for this annotation

GitHub Actions / lint

Avoid using promises inside of callbacks
.close()
.then(() => done())

Check warning on line 1768 in handwritten/spanner/observability-test/spanner.ts

View workflow job for this annotation

GitHub Actions / lint

Avoid calling back inside of a promise
.catch(err => done(err));

Check warning on line 1769 in handwritten/spanner/observability-test/spanner.ts

View workflow job for this annotation

GitHub Actions / lint

Avoid calling back inside of a promise
});
});

Expand Down Expand Up @@ -1734,7 +1794,7 @@
})
.on('error', err => {
assert.strictEqual(err.message, '2 UNKNOWN: Non-retryable error');
database

Check warning on line 1797 in handwritten/spanner/observability-test/spanner.ts

View workflow job for this annotation

GitHub Actions / lint

Avoid using promises inside of callbacks

Check warning on line 1797 in handwritten/spanner/observability-test/spanner.ts

View workflow job for this annotation

GitHub Actions / lint

Avoid using promises inside of callbacks
.close()
.then(() => {
traceExporter.forceFlush();
Expand Down Expand Up @@ -1778,6 +1838,7 @@
);

done();
return null;
})
.catch(err => done(err));
});
Expand Down Expand Up @@ -1938,9 +1999,7 @@
assert.strictEqual(attempts, 1);
tx!
.commit()
.then(() => {
database.close().catch(assert.ifError);
})
.then(() => database.close())
.catch(assert.ifError);
});
});
Expand Down
56 changes: 42 additions & 14 deletions handwritten/spanner/src/batch-transaction.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,12 @@ import {
ResourceCallback,
addLeaderAwareRoutingHeader,
} from '../src/common';
import {startTrace, setSpanError, traceConfig} from './instrument';
import {
startTrace,
setSpanError,
traceConfig,
getQueryTraceConfig,
} from './instrument';
import {injectRequestIDIntoHeaders} from './request_id_header';
import {isString} from './helper';

Expand Down Expand Up @@ -169,15 +174,24 @@ class BatchTransaction extends Snapshot {
const request: ExecuteSqlRequest =
typeof query === 'string' ? {sql: query} : query;

const reqOpts = Object.assign({}, request, Snapshot.encodeParams(request));
const {
gaxOptions: _omittedGaxOptions,
types: _omittedTypes,
...cleanRequest
} = request as ExecuteSqlRequest & {types?: unknown};
void _omittedGaxOptions;
void _omittedTypes;

delete (reqOpts as any).gaxOptions;
delete (reqOpts as any).types;
const reqOpts = Object.assign(
{},
cleanRequest,
Snapshot.encodeParams(request),
);

const traceConfig: traceConfig = {
sql: request.sql,
opts: this._observabilityOptions,
dbName: this.getDBName(),
...getQueryTraceConfig(query),
};
return startTrace(
'BatchTransaction.createQueryPartitions',
Expand Down Expand Up @@ -233,17 +247,21 @@ class BatchTransaction extends Snapshot {
'BatchTransaction.createPartitions_',
traceConfig,
span => {
const query = Object.assign({}, config.reqOpts, {
const baseRequest = Object.assign({}, config.reqOpts, {
session: this.session.formattedName_,
transaction: {id: this.id},
});
config.reqOpts = Object.assign({}, query);
config.reqOpts = baseRequest;
const headers = {
[CLOUD_RESOURCE_HEADER]: (this.session.parent as Database)
.formattedName_,
};
config.headers = injectRequestIDIntoHeaders(headers, this.session);
delete query.partitionOptions;
const {
partitionOptions: _omittedPartitionOptions,
...baseRequestWithoutPartitionOptions
} = baseRequest;
void _omittedPartitionOptions;
Comment thread
olavloite marked this conversation as resolved.
this.session.request(config, (err, resp) => {
if (err) {
setSpanError(span, err);
Expand All @@ -253,7 +271,11 @@ class BatchTransaction extends Snapshot {
}

const partitions = resp.partitions.map(partition => {
return Object.assign({}, query, partition);
return Object.assign(
{},
baseRequestWithoutPartitionOptions,
partition,
);
});

if (resp.transaction) {
Expand Down Expand Up @@ -323,14 +345,20 @@ class BatchTransaction extends Snapshot {
'BatchTransaction.createReadPartitions',
traceConfig,
span => {
const reqOpts = Object.assign({}, options, {
const {
gaxOptions: _omittedGaxOptions,
keys: _omittedKeys,
ranges: _omittedRanges,
...cleanOptions
} = options;
void _omittedGaxOptions;
void _omittedKeys;
void _omittedRanges;

const reqOpts = Object.assign({}, cleanOptions, {
keySet: Snapshot.encodeKeySet(options),
});

delete reqOpts.gaxOptions;
delete reqOpts.keys;
delete reqOpts.ranges;

const headers: {[k: string]: string} = {};
if (this._getSpanner().routeToLeaderEnabled) {
addLeaderAwareRoutingHeader(headers);
Expand Down
14 changes: 6 additions & 8 deletions handwritten/spanner/src/database.ts
Original file line number Diff line number Diff line change
Expand Up @@ -118,6 +118,7 @@ import {
setSpanError,
setSpanErrorAndException,
traceConfig,
getQueryTraceConfig,
} from './instrument';
import {
AtomicCounter,
Expand Down Expand Up @@ -2944,8 +2945,8 @@ class Database extends common.GrpcServiceObject {
startTrace(
'Database.run',
{
...(query as ExecuteSqlRequest),
...this._traceConfig,
...getQueryTraceConfig(query),
},
span => {
this.runStream(query, options)
Expand Down Expand Up @@ -2985,9 +2986,9 @@ class Database extends common.GrpcServiceObject {
options: TimestampBounds,
callback: RunCallback,
): void {
const traceConfig = {
...(query as ExecuteSqlRequest),
const traceConfig: traceConfig = {
...this._traceConfig,
...getQueryTraceConfig(query),
};

startTrace('Database.run', traceConfig, runSpan => {
Expand Down Expand Up @@ -3142,10 +3143,8 @@ class Database extends common.GrpcServiceObject {
return startTrace(
'Database.runPartitionedUpdate',
{
...(query as RunPartitionedUpdateOptions),
...this._traceConfig,
requestTag: (query as RunPartitionedUpdateOptions)?.requestOptions
?.requestTag,
...getQueryTraceConfig(query),
},
span => {
this.sessionFactory_.getSessionForPartitionedOps((err, session) => {
Expand Down Expand Up @@ -3335,9 +3334,8 @@ class Database extends common.GrpcServiceObject {
return startTrace(
'Database.runStream',
{
...(query as ExecuteSqlRequest),
...this._traceConfig,
requestTag: (query as ExecuteSqlRequest)?.requestOptions?.requestTag,
...getQueryTraceConfig(query),
},
span => {
this.sessionFactory_.getSession((err, session) => {
Expand Down
31 changes: 30 additions & 1 deletion handwritten/spanner/src/instrument.ts
Original file line number Diff line number Diff line change
Expand Up @@ -90,8 +90,37 @@ interface traceConfig {
opts?: ObservabilityOptions;
}

interface QueryWithRequestOptions {
sql?: string | SQLStatement;
requestOptions?: {
requestTag?: string | null;
transactionTag?: string | null;
} | null;
}

function getQueryTraceConfig(query?: string | QueryWithRequestOptions | null): {
sql?: string | SQLStatement;
requestTag?: string | null;
} {
Comment thread
olavloite marked this conversation as resolved.
if (typeof query === 'string') {
return {sql: query};
}
if (query && typeof query === 'object') {
return {
sql: query.sql,
requestTag: query.requestOptions?.requestTag,
};
}
return {};
}

const SPAN_NAMESPACE_PREFIX = 'CloudSpanner'; // TODO: discuss & standardize this prefix.
export {SPAN_NAMESPACE_PREFIX, traceConfig};
export {
SPAN_NAMESPACE_PREFIX,
traceConfig,
QueryWithRequestOptions,
getQueryTraceConfig,
};

const {
AsyncHooksContextManager,
Expand Down
Loading
Loading