diff --git a/.github/workflows/apm-integrations.yml b/.github/workflows/apm-integrations.yml index 4ee153f491c..c3ee8bcb327 100644 --- a/.github/workflows/apm-integrations.yml +++ b/.github/workflows/apm-integrations.yml @@ -1169,6 +1169,24 @@ jobs: - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 - uses: ./.github/actions/plugins/test + postgres-js: + runs-on: ubuntu-latest + permissions: + id-token: write + services: + postgres: + image: postgres@sha256:75ebf479151a8fd77bf2fed46ef76ce8d518c23264734c48f2d1de42b4eb40ae # 9.5 + env: + POSTGRES_PASSWORD: postgres + ports: + - 5432:5432 + env: + PLUGINS: postgres + SERVICES: postgres + steps: + - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 + - uses: ./.github/actions/plugins/test + prisma: strategy: fail-fast: false diff --git a/docs/API.md b/docs/API.md index 8811ff367ca..e3497a11e61 100644 --- a/docs/API.md +++ b/docs/API.md @@ -135,6 +135,7 @@ tracer.use('openai', {
+
@@ -221,6 +222,7 @@ tracer.use('openai', { * [pg](./interfaces/export_.plugins.pg.html) * [pino](./interfaces/export_.plugins.pino.html) * [playwright](./interfaces/export_.plugins.playwright.html) +* [postgres](./interfaces/export_.plugins.postgres.html) * [prisma](./interfaces/export_.plugins.prisma.html) * [protobufjs](./interfaces/export_.plugins.protobufjs.html) * [redis](./interfaces/export_.plugins.redis.html) diff --git a/docs/test.ts b/docs/test.ts index 509be07a954..aa5d3fc563e 100644 --- a/docs/test.ts +++ b/docs/test.ts @@ -416,6 +416,10 @@ tracer.use('pg', { appendComment: true }); tracer.use('pg', { truncate: true }); tracer.use('pg', { truncate: 5000 }); tracer.use('pino'); +tracer.use('postgres'); +tracer.use('postgres', { service: 'postgres-service' }); +tracer.use('postgres', { truncate: true }); +tracer.use('postgres', { truncate: 5000 }); tracer.use('prisma'); tracer.use('protobufjs'); tracer.use('redis'); diff --git a/index.d.ts b/index.d.ts index 0c41671f608..41bacbcf725 100644 --- a/index.d.ts +++ b/index.d.ts @@ -298,6 +298,7 @@ interface Plugins { "playwright": tracer.plugins.playwright; "pg": tracer.plugins.pg; "pino": tracer.plugins.pino; + "postgres": tracer.plugins.postgres; "prisma": tracer.plugins.prisma; "protobufjs": tracer.plugins.protobufjs; "redis": tracer.plugins.redis; @@ -3126,6 +3127,17 @@ declare namespace tracer { */ interface pino extends Integration {} + /** + * This plugin automatically instruments the + * [Postgres.js](https://github.com/porsager/postgres) module. + */ + interface postgres extends DatabaseInstrumentation { + /** + * The service name to be used for this plugin. + */ + service?: string; + } + /** * This plugin automatically instruments the * [@prisma/client](https://www.prisma.io/docs/orm/prisma-client) module. diff --git a/index.d.v5.ts b/index.d.v5.ts index b1016d09059..da3dbedb63a 100644 --- a/index.d.v5.ts +++ b/index.d.v5.ts @@ -300,6 +300,7 @@ interface Plugins { "playwright": tracer.plugins.playwright; "pg": tracer.plugins.pg; "pino": tracer.plugins.pino; + "postgres": tracer.plugins.postgres; "prisma": tracer.plugins.prisma; "protobufjs": tracer.plugins.protobufjs; "redis": tracer.plugins.redis; @@ -3296,6 +3297,17 @@ declare namespace tracer { */ interface pino extends Integration {} + /** + * This plugin automatically instruments the + * [Postgres.js](https://github.com/porsager/postgres) module. + */ + interface postgres extends DatabaseInstrumentation { + /** + * The service name to be used for this plugin. + */ + service?: string; + } + /** * This plugin automatically instruments the * [@prisma/client](https://www.prisma.io/docs/orm/prisma-client) module. diff --git a/packages/datadog-instrumentations/src/helpers/hooks.js b/packages/datadog-instrumentations/src/helpers/hooks.js index 5fd596f6337..2058ff6d58f 100644 --- a/packages/datadog-instrumentations/src/helpers/hooks.js +++ b/packages/datadog-instrumentations/src/helpers/hooks.js @@ -141,6 +141,7 @@ module.exports = { 'pino-pretty': () => require('../pino'), playwright: () => require('../playwright'), 'playwright-core': () => require('../playwright'), + postgres: { esmFirst: true, fn: () => require('../postgres') }, 'promise-js': () => require('../promise-js'), promise: () => require('../promise'), protobufjs: () => require('../protobufjs'), diff --git a/packages/datadog-instrumentations/src/helpers/rewriter/index.js b/packages/datadog-instrumentations/src/helpers/rewriter/index.js index 7ce4d41ce53..9f3533086fe 100644 --- a/packages/datadog-instrumentations/src/helpers/rewriter/index.js +++ b/packages/datadog-instrumentations/src/helpers/rewriter/index.js @@ -113,6 +113,10 @@ function createMatcher (moduleType) { configureMercuriusRequest, waitForAsyncEnd, } = require('./transforms') + const { + postgresQueryHandlers, + postgresQueryLifecycle, + } = require('./transforms/postgres') const matcher = create(instrumentations, getDcPolyfillSpecifier(moduleType)) @@ -126,6 +130,8 @@ function createMatcher (moduleType) { matcher.addTransform('configureGraphqlJitExecute', configureGraphqlJitExecute) matcher.addTransform('configureGraphqlJitRuntime', configureGraphqlJitRuntime) matcher.addTransform('configureMercuriusRequest', configureMercuriusRequest) + matcher.addTransform('postgresQueryHandlers', postgresQueryHandlers) + matcher.addTransform('postgresQueryLifecycle', postgresQueryLifecycle) return matcher } diff --git a/packages/datadog-instrumentations/src/helpers/rewriter/instrumentations/index.js b/packages/datadog-instrumentations/src/helpers/rewriter/instrumentations/index.js index 5807827187c..e8f941eb609 100644 --- a/packages/datadog-instrumentations/src/helpers/rewriter/instrumentations/index.js +++ b/packages/datadog-instrumentations/src/helpers/rewriter/instrumentations/index.js @@ -13,6 +13,7 @@ module.exports = [ ...require('./modelcontextprotocol-sdk'), ...require('./openai-agents'), ...require('./playwright'), + ...require('./postgres'), ...require('./webdriverio'), ...require('./aws-durable-execution-sdk-js'), ] diff --git a/packages/datadog-instrumentations/src/helpers/rewriter/instrumentations/postgres.js b/packages/datadog-instrumentations/src/helpers/rewriter/instrumentations/postgres.js new file mode 100644 index 00000000000..7203225c8b9 --- /dev/null +++ b/packages/datadog-instrumentations/src/helpers/rewriter/instrumentations/postgres.js @@ -0,0 +1,25 @@ +'use strict' + +const handlers = ['src/index.js', 'cjs/src/index.js'].map(filePath => ({ + module: { + name: 'postgres', + versionRange: '>=3.0.0', + filePath, + }, + astQuery: 'Program', + transform: 'postgresQueryHandlers', + channelName: 'query', +})) + +const lifecycles = ['src/query.js', 'cjs/src/query.js'].map(filePath => ({ + module: { + name: 'postgres', + versionRange: '>=3.0.0', + filePath, + }, + astQuery: 'Program', + transform: 'postgresQueryLifecycle', + channelName: 'query', +})) + +module.exports = [...handlers, ...lifecycles] diff --git a/packages/datadog-instrumentations/src/helpers/rewriter/targets.json b/packages/datadog-instrumentations/src/helpers/rewriter/targets.json index 04d4ffdc06a..cb21c46550c 100644 --- a/packages/datadog-instrumentations/src/helpers/rewriter/targets.json +++ b/packages/datadog-instrumentations/src/helpers/rewriter/targets.json @@ -66,6 +66,10 @@ "playwright-core/lib/coreBundle.js": "playwright-core", "playwright/lib/index.js": "playwright", "playwright/lib/runner/index.js": "playwright", + "postgres/cjs/src/index.js": "postgres", + "postgres/cjs/src/query.js": "postgres", + "postgres/src/index.js": "postgres", + "postgres/src/query.js": "postgres", "webdriver/build/index.js": "webdriver", "webdriver/build/node.js": "webdriver", "webdriverio/build/index.js": "webdriverio", diff --git a/packages/datadog-instrumentations/src/helpers/rewriter/transforms/postgres.js b/packages/datadog-instrumentations/src/helpers/rewriter/transforms/postgres.js new file mode 100644 index 00000000000..feb56c9230d --- /dev/null +++ b/packages/datadog-instrumentations/src/helpers/rewriter/transforms/postgres.js @@ -0,0 +1,271 @@ +'use strict' + +const assert = require('node:assert') + +const { parse, query } = require('../compiler') + +const POSTGRES_OPTIONS = '__ddTracePostgresOptions' +const POSTGRES_QUERY_REGISTRY = '__ddTracePostgresQueries' +const POSTGRES_QUERY_REGISTRY_SYMBOL = '__ddTracePostgresQueriesSymbol' +const POSTGRES_READY = '__ddTracePostgresQueryReady' +const POSTGRES_STATEMENT = 'this.string ?? (this.tagged && this.strings?.length !== 1 ? undefined : this.strings?.[0])' + +module.exports = { + postgresQueryHandlers, + postgresQueryLifecycle, +} + +/** + * Starts a Postgres.js query at the pool handler that owns its connection options. + * + * @param {object} state + * @param {import('estree').Program} program + */ +function postgresQueryHandlers (state, program) { + const channelVariable = injectPostgresTracingChannel(state, program) + injectPostgresReadyCheck(program, findPostgresQueryIdentifier(program)) + capturePostgresOptions(program) + + const handlers = query(program, 'FunctionDeclaration[id.name="handler"]') + assert(handlers.length >= 2 && handlers.length <= 3, 'postgresQueryHandlers: unexpected handler count') + + for (const handler of handlers) { + const queryParameter = handler.params[0] + assert(queryParameter?.type === 'Identifier', 'postgresQueryHandlers: handler query parameter changed') + wrapPostgresHandler(handler, queryParameter.name, channelVariable) + } +} + +/** + * Publishes the final Postgres.js query settlement without changing its Promise implementation. + * + * @param {object} state + * @param {import('estree').Program} program + */ +function postgresQueryLifecycle (state, program) { + const channelVariable = injectPostgresTracingChannel(state, program) + injectPostgresQueryRegistration(program) + + const resolveAssignments = query( + program, + 'AssignmentExpression[left.object.type="ThisExpression"][left.property.name="resolve"]' + ) + const rejectAssignments = query( + program, + 'AssignmentExpression[left.object.type="ThisExpression"][left.property.name="reject"]' + ) + + assert(resolveAssignments.length === 2, 'postgresQueryLifecycle: unexpected resolve assignment count') + assert(rejectAssignments.length === 2, 'postgresQueryLifecycle: unexpected reject assignment count') + + for (const assignment of resolveAssignments) { + wrapPostgresResolution(assignment, channelVariable) + } + for (const assignment of rejectAssignments) { + wrapPostgresRejection(assignment, channelVariable) + } +} + +/** + * @param {import('estree').FunctionDeclaration|import('estree').ArrowFunctionExpression} handler + * @param {string} queryParameter + * @param {string} channelVariable + */ +function wrapPostgresHandler (handler, queryParameter, channelVariable) { + assert(handler.body.type === 'BlockStatement', 'postgresQueryHandlers: handler body changed') + assert(!handler.async && !handler.generator, 'postgresQueryHandlers: unsupported handler kind') + + const originalBody = handler.body.body + const wrapperBody = parse(` + function wrapper () { + if (${channelVariable}.start.hasSubscribers && ${POSTGRES_READY} && !${queryParameter}.cancelled) { + const __ddTraceContext = { + query: ${queryParameter}, + database: ${POSTGRES_OPTIONS}.database, + user: ${POSTGRES_OPTIONS}.user + }; + if (${POSTGRES_OPTIONS}.host.length === 1 && !${POSTGRES_OPTIONS}.path) { + __ddTraceContext.host = ${POSTGRES_OPTIONS}.host[0]; + __ddTraceContext.port = ${POSTGRES_OPTIONS}.port[0]; + } + return ${channelVariable}.start.runStores(__ddTraceContext, () => {}); + } + } + `).body[0].body.body + + wrapperBody[0].consequent.body.at(-1).argument.arguments[1].body.body = originalBody + handler.body.body = [...wrapperBody, ...originalBody] +} + +/** + * @param {import('estree').Program} program + */ +function capturePostgresOptions (program) { + const postgresFunctions = query(program, 'FunctionDeclaration[id.name="Postgres"]') + assert(postgresFunctions.length === 1, 'postgresQueryHandlers: Postgres function changed') + + const body = postgresFunctions[0].body.body + const optionsIndex = body.findIndex(statement => + statement.type === 'VariableDeclaration' && statement.declarations.some(({ id }) => id.name === 'options') + ) + assert(optionsIndex !== -1, 'postgresQueryHandlers: options declaration changed') + + const capture = parse(`const ${POSTGRES_OPTIONS} = options;`).body[0] + body.splice(optionsIndex + 1, 0, capture) +} + +/** + * @param {import('estree').Program} program + * @param {string} queryIdentifier + */ +function injectPostgresReadyCheck (program, queryIdentifier) { + const statements = parse(` + const ${POSTGRES_QUERY_REGISTRY} = globalThis[Symbol.for('dd-trace:postgres:query')]; + const ${POSTGRES_READY} = ${POSTGRES_QUERY_REGISTRY} instanceof WeakSet && + ${POSTGRES_QUERY_REGISTRY}.has(${queryIdentifier}); + `).body + let importIndex = program.body.length - 1 + while (importIndex >= 0) { + const statement = program.body[importIndex] + if (statement.type === 'ImportDeclaration' || statement.type === 'VariableDeclaration') break + importIndex-- + } + + program.body.splice(importIndex + 1, 0, ...statements) +} + +/** + * @param {import('estree').Program} program + * @returns {string} + */ +function findPostgresQueryIdentifier (program) { + const identifiers = [] + + const imports = query(program, 'ImportDeclaration') + for (const declaration of imports) { + if (declaration.source.value !== './query.js') continue + + for (const specifier of declaration.specifiers) { + if (specifier.type === 'ImportSpecifier' && specifier.imported.name === 'Query') { + identifiers.push(specifier.local.name) + } + } + } + + const declarations = query(program, 'VariableDeclarator[id.type="ObjectPattern"]') + for (const declaration of declarations) { + if (declaration.init?.type !== 'CallExpression' || declaration.init.callee.name !== 'require') continue + if (declaration.init.arguments[0]?.value !== './query.js') continue + + for (const property of declaration.id.properties) { + if (property.type === 'Property' && property.key.name === 'Query' && property.value.type === 'Identifier') { + identifiers.push(property.value.name) + } + } + } + + assert(identifiers.length === 1, 'postgresQueryHandlers: Query import changed') + return identifiers[0] +} + +/** + * @param {import('estree').Program} program + */ +function injectPostgresQueryRegistration (program) { + const queryClassIndex = program.body.findIndex(statement => + statement.type === 'VariableDeclaration' && + statement.declarations.some(({ id }) => id.type === 'Identifier' && id.name === 'Query') || + statement.type === 'ExportNamedDeclaration' && statement.declaration?.type === 'ClassDeclaration' && + statement.declaration.id?.name === 'Query' + ) + assert(queryClassIndex !== -1, 'postgresQueryLifecycle: Query class changed') + + const statements = parse(` + const ${POSTGRES_QUERY_REGISTRY_SYMBOL} = Symbol.for('dd-trace:postgres:query'); + const ${POSTGRES_QUERY_REGISTRY} = globalThis[${POSTGRES_QUERY_REGISTRY_SYMBOL}] instanceof WeakSet + ? globalThis[${POSTGRES_QUERY_REGISTRY_SYMBOL}] + : (globalThis[${POSTGRES_QUERY_REGISTRY_SYMBOL}] = new WeakSet()); + ${POSTGRES_QUERY_REGISTRY}.add(Query); + `).body + program.body.splice(queryClassIndex + 1, 0, ...statements) +} + +/** + * @param {import('estree').AssignmentExpression} assignment + * @param {string} channelVariable + */ +function wrapPostgresResolution (assignment, channelVariable) { + const resolution = assignment.right + assert( + resolution.type === 'ArrowFunctionExpression' && resolution.body.type !== 'BlockStatement', + 'postgresQueryLifecycle: resolve function changed' + ) + + const originalBody = resolution.body + const body = parse(` + function wrapper () { + if (${channelVariable}.asyncEnd.hasSubscribers && (!this.streaming || !this.active)) { + const statement = ${POSTGRES_STATEMENT}; + ${channelVariable}.asyncEnd.publish({ query: this, statement, pid: this.state?.pid }); + } + return undefined; + } + `).body[0].body + resolution.body = body + resolution.expression = false + body.body.at(-1).argument = originalBody +} + +/** + * @param {import('estree').AssignmentExpression} assignment + * @param {string} channelVariable + */ +function wrapPostgresRejection (assignment, channelVariable) { + const rejection = assignment.right + assert( + rejection.type === 'ArrowFunctionExpression' && + rejection.body.type !== 'BlockStatement', + 'postgresQueryLifecycle: reject function changed' + ) + const errorParameter = rejection.params[0] + assert(errorParameter?.type === 'Identifier', 'postgresQueryLifecycle: reject parameter changed') + + const originalBody = rejection.body + const body = parse(` + function wrapper () { + if (${channelVariable}.error.hasSubscribers || ${channelVariable}.asyncEnd.hasSubscribers) { + const statement = ${POSTGRES_STATEMENT}; + const __ddTraceContext = { query: this, error: ${errorParameter.name}, statement, pid: this.state?.pid }; + if (${channelVariable}.error.hasSubscribers) { + ${channelVariable}.error.publish(__ddTraceContext); + } + if (${channelVariable}.asyncEnd.hasSubscribers) { + ${channelVariable}.asyncEnd.publish(__ddTraceContext); + } + } + return undefined; + } + `).body[0].body + rejection.body = body + rejection.expression = false + body.body.at(-1).argument = originalBody +} + +/** + * @param {object} state + * @param {import('estree').Program} program + * @returns {string} + */ +function injectPostgresTracingChannel (state, program) { + state.transforms.tracingChannelDeclaration(state, program) + + const channelName = `orchestrion:${state.module.name}:${state.channelName}` + const declarations = query( + program, + 'VariableDeclarator[init.type="CallExpression"][init.callee.name="tr_ch_apm_tracingChannel"]' + ).filter(({ init }) => init.arguments[0]?.value === channelName) + + assert(declarations.length === 1, 'postgres: tracing channel declaration changed') + assert(declarations[0].id.type === 'Identifier', 'postgres: tracing channel identifier changed') + return declarations[0].id.name +} diff --git a/packages/datadog-instrumentations/src/postgres.js b/packages/datadog-instrumentations/src/postgres.js new file mode 100644 index 00000000000..ebb803a7df5 --- /dev/null +++ b/packages/datadog-instrumentations/src/postgres.js @@ -0,0 +1,7 @@ +'use strict' + +const { addHook, getHooks } = require('./helpers/instrument') + +for (const hook of getHooks('postgres')) { + addHook(hook, exports => exports) +} diff --git a/packages/datadog-instrumentations/test/helpers/rewriter/index.spec.js b/packages/datadog-instrumentations/test/helpers/rewriter/index.spec.js index 2810366e244..bb3e3ddb491 100644 --- a/packages/datadog-instrumentations/test/helpers/rewriter/index.spec.js +++ b/packages/datadog-instrumentations/test/helpers/rewriter/index.spec.js @@ -547,6 +547,46 @@ describe('check-require-cache', () => { callbackName: 'beforeContinue', }, }, + { + module: { + name: 'test', + versionRange: '>=0.1', + filePath: 'postgres-query-alias.js', + }, + astQuery: 'Program', + transform: 'postgresQueryHandlers', + channelName: 'query', + }, + { + module: { + name: 'test', + versionRange: '>=0.1', + filePath: 'postgres-query-async-handler.js', + }, + astQuery: 'Program', + transform: 'postgresQueryHandlers', + channelName: 'query', + }, + { + module: { + name: 'test', + versionRange: '>=0.1', + filePath: 'postgres-query-generator-handler.js', + }, + astQuery: 'Program', + transform: 'postgresQueryHandlers', + channelName: 'query', + }, + { + module: { + name: 'test-esm', + versionRange: '>=0.1', + filePath: 'postgres-query-alias.js', + }, + astQuery: 'Program', + transform: 'postgresQueryHandlers', + channelName: 'query', + }, { module: { name: 'test-esm', @@ -1310,6 +1350,77 @@ describe('check-require-cache', () => { assert.ok(subs.start.calledOnce, 'instrumented start channel should fire once') }) + + it('should ignore unrelated destructured requires and use an aliased Postgres Query binding', () => { + ch = tracingChannel('orchestrion:test:query') + subs = { start: sinon.spy() } + ch.subscribe(subs) + + const Postgres = compileFile('postgres-query-alias') + const queries = Postgres() + + assert.strictEqual(queries.length, 2) + assert.strictEqual(subs.start.callCount, 2) + }) + + it('should use an aliased Postgres Query import binding', async () => { + const fixtureDirectory = resolve(__dirname, 'node_modules', 'test-esm') + const filename = join(fixtureDirectory, 'postgres-query-alias.js') + const source = readFileSync(filename, 'utf8') + const rewritten = rewriter.rewrite(source, filename, 'module', { + moduleName: 'test-esm', + filePath: 'postgres-query-alias.js', + }) + const directory = mkdtempSync(join(tmpdir(), 'dd-rewriter-postgres-esm-')) + const outputFile = join(directory, 'postgres-query-alias.js') + + writeFileSync(join(directory, 'package.json'), '{"type":"module"}') + writeFileSync(join(directory, 'query.js'), readFileSync(join(fixtureDirectory, 'query.js'))) + writeFileSync(outputFile, rewritten) + + ch = tracingChannel('orchestrion:test-esm:query') + subs = { start: sinon.spy() } + ch.subscribe(subs) + + const { default: Postgres } = await import(pathToFileURL(outputFile).href) + const queries = Postgres() + + assert.strictEqual(queries.length, 2) + assert.strictEqual(subs.start.callCount, 2) + }) + + it('should leave async Postgres handlers untouched', async () => { + const filename = resolve(__dirname, 'node_modules', 'test', 'postgres-query-async-handler.js') + const source = readFileSync(filename, 'utf8') + + ch = tracingChannel('orchestrion:test:query') + subs = { start: sinon.spy() } + ch.subscribe(subs) + + const Postgres = compileFile('postgres-query-async-handler') + const queries = await Promise.all(Postgres()) + + assert.strictEqual(content, source) + assert.strictEqual(queries.length, 2) + assert.strictEqual(subs.start.callCount, 0) + }) + + it('should leave generator Postgres handlers untouched', () => { + const filename = resolve(__dirname, 'node_modules', 'test', 'postgres-query-generator-handler.js') + const source = readFileSync(filename, 'utf8') + + ch = tracingChannel('orchestrion:test:query') + subs = { start: sinon.spy() } + ch.subscribe(subs) + + const Postgres = compileFile('postgres-query-generator-handler') + const queries = Postgres() + + assert.strictEqual(content, source) + assert.strictEqual(queries[0].next().value.constructor.name, 'Query') + assert.strictEqual(queries[1].next().value.constructor.name, 'Query') + assert.strictEqual(subs.start.callCount, 0) + }) }) describe('rewriter source-map trailer', () => { diff --git a/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test-esm/postgres-query-alias.js b/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test-esm/postgres-query-alias.js new file mode 100644 index 00000000000..87dc9724ee5 --- /dev/null +++ b/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test-esm/postgres-query-alias.js @@ -0,0 +1,24 @@ +import { Query as PostgresQuery } from './query.js' + +export default function Postgres () { + const options = { + database: 'postgres', + host: ['localhost'], + port: [5432], + user: 'postgres', + } + + function handler (query) { + return query + } + + function dispatch (query) { + function handler (query) { + return query + } + + return handler(query) + } + + return [handler(new PostgresQuery()), dispatch(new PostgresQuery())] +} diff --git a/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test-esm/query.js b/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test-esm/query.js new file mode 100644 index 00000000000..ea1e55b2419 --- /dev/null +++ b/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test-esm/query.js @@ -0,0 +1,5 @@ +export class Query {} + +const registry = globalThis[Symbol.for('dd-trace:postgres:query')] ?? new WeakSet() +registry.add(Query) +globalThis[Symbol.for('dd-trace:postgres:query')] = registry diff --git a/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test/postgres-query-alias.js b/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test/postgres-query-alias.js new file mode 100644 index 00000000000..75fdb266b55 --- /dev/null +++ b/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test/postgres-query-alias.js @@ -0,0 +1,32 @@ +'use strict' + +const { strictEqual } = require('node:assert') + +const { Query: PostgresQuery } = require('./query.js') + +function Postgres () { + strictEqual(typeof PostgresQuery, 'function') + + const options = { + database: 'postgres', + host: ['localhost'], + port: [5432], + user: 'postgres', + } + + function handler (query) { + return query + } + + function dispatch (query) { + function handler (query) { + return query + } + + return handler(query) + } + + return [handler(new PostgresQuery()), dispatch(new PostgresQuery())] +} + +module.exports = Postgres diff --git a/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test/postgres-query-async-handler.js b/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test/postgres-query-async-handler.js new file mode 100644 index 00000000000..80e18ce00b4 --- /dev/null +++ b/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test/postgres-query-async-handler.js @@ -0,0 +1,28 @@ +'use strict' + +const { Query } = require('./query.js') + +function Postgres () { + const options = { + database: 'postgres', + host: ['localhost'], + port: [5432], + user: 'postgres', + } + + async function handler (query) { + return query + } + + async function dispatch (query) { + async function handler (query) { + return query + } + + return handler(query) + } + + return [handler(new Query()), dispatch(new Query())] +} + +module.exports = Postgres diff --git a/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test/postgres-query-generator-handler.js b/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test/postgres-query-generator-handler.js new file mode 100644 index 00000000000..d11ad614c19 --- /dev/null +++ b/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test/postgres-query-generator-handler.js @@ -0,0 +1,28 @@ +'use strict' + +const { Query } = require('./query.js') + +function Postgres () { + const options = { + database: 'postgres', + host: ['localhost'], + port: [5432], + user: 'postgres', + } + + function * handler (query) { + yield query + } + + function dispatch (query) { + function * handler (query) { + yield query + } + + return handler(query) + } + + return [handler(new Query()), dispatch(new Query())] +} + +module.exports = Postgres diff --git a/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test/query.js b/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test/query.js new file mode 100644 index 00000000000..d184519d39f --- /dev/null +++ b/packages/datadog-instrumentations/test/helpers/rewriter/node_modules/test/query.js @@ -0,0 +1,9 @@ +'use strict' + +class Query {} + +const registry = globalThis[Symbol.for('dd-trace:postgres:query')] ?? new WeakSet() +registry.add(Query) +globalThis[Symbol.for('dd-trace:postgres:query')] = registry + +module.exports = { Query } diff --git a/packages/datadog-plugin-postgres/src/index.js b/packages/datadog-plugin-postgres/src/index.js new file mode 100644 index 00000000000..81e005ac997 --- /dev/null +++ b/packages/datadog-plugin-postgres/src/index.js @@ -0,0 +1,87 @@ +'use strict' + +const { CLIENT_PORT_KEY } = require('../../dd-trace/src/constants') +const DatabasePlugin = require('../../dd-trace/src/plugins/database') + +/** + * @typedef {object} PostgresContext + * @property {{ span: import('../../..').Span }} currentStore + * @property {string} database + * @property {unknown} [error] + * @property {string} [host] + * @property {number} [pid] + * @property {number} [port] + * @property {PostgresQuery} query + * @property {string} [statement] + * @property {string} user + * + * @typedef {object} PostgresQuery + */ + +class PostgresPlugin extends DatabasePlugin { + static id = 'postgres' + static prefix = 'tracing:orchestrion:postgres:query' + + /** @type {WeakMap} */ + #contexts = new WeakMap() + + /** + * @param {PostgresContext} ctx + * @returns {object} + */ + bindStart (ctx) { + const { database, host, port, query, user } = ctx + + const span = this.startSpan(this.operationName(), { + service: this.serviceName({ pluginConfig: this.config }), + type: 'sql', + kind: 'client', + meta: { + 'db.type': this.system, + 'db.name': database, + 'db.user': user, + }, + }, ctx) + + if (host !== undefined) { + span.addTags({ + 'out.host': host, + [CLIENT_PORT_KEY]: port, + }) + } + + this.#contexts.set(query, ctx) + return ctx.currentStore + } + + /** + * @param {PostgresContext} ctx + */ + error (ctx) { + const span = this.#contexts.get(ctx.query)?.currentStore.span + if (span !== undefined) { + this.addError(ctx.error, span) + } + } + + /** + * @param {PostgresContext} result + */ + asyncEnd (result) { + const ctx = this.#contexts.get(result.query) + if (ctx === undefined) return + + this.#contexts.delete(result.query) + + const span = ctx.currentStore.span + + if (typeof result.statement === 'string') { + span.setTag('resource.name', this.maybeTruncate(result.statement)) + } + + span.setTag('db.pid', result.pid) + this.finish(ctx) + } +} + +module.exports = PostgresPlugin diff --git a/packages/datadog-plugin-postgres/test/fixtures/select.sql b/packages/datadog-plugin-postgres/test/fixtures/select.sql new file mode 100644 index 00000000000..e2654c8bd5b --- /dev/null +++ b/packages/datadog-plugin-postgres/test/fixtures/select.sql @@ -0,0 +1 @@ +SELECT $1::text AS message diff --git a/packages/datadog-plugin-postgres/test/index.spec.js b/packages/datadog-plugin-postgres/test/index.spec.js new file mode 100644 index 00000000000..79c30ed13ab --- /dev/null +++ b/packages/datadog-plugin-postgres/test/index.spec.js @@ -0,0 +1,602 @@ +'use strict' + +const assert = require('node:assert/strict') +const { once } = require('node:events') +const path = require('node:path') + +const dc = require('dc-polyfill') +const semver = require('semver') + +const { ERROR_MESSAGE, ERROR_STACK, ERROR_TYPE } = require('../../dd-trace/src/constants') +const agent = require('../../dd-trace/test/plugins/agent') +const { withNamingSchema, withPeerService, withVersions } = require('../../dd-trace/test/setup/mocha') +const { expectedSchema, rawExpectedSchema } = require('./naming') + +const postgresStartChannel = dc.channel('tracing:orchestrion:postgres:query:start') + +const POSTGRES_TARGET = { + host: '127.0.0.1', + port: 5432, + user: 'postgres', + password: 'postgres', + database: 'postgres', + max: 1, +} + +/** + * @param {string} resource + */ +function resourcePattern (resource) { + return new RegExp(`^${resource.replace(/[.*+?^${}()|[\]\\]/g, '\\$&')}$`) +} + +/** + * @param {string} resource + * @returns {Promise} + */ +function assertNoQuerySpan (resource) { + return agent.assertNoTraces(() => { + assert.fail(`query was traced: ${resource}`) + }, { + spanResourceMatch: resourcePattern(resource), + timeoutMs: 200, + }) +} + +/** + * @param {string} allowedResource + * @returns {Promise} + */ +function assertNoOtherQuerySpans (allowedResource) { + return agent.assertNoTraces(traces => { + const span = traces.flat().find(span => + span.meta.component === 'postgres' && span.resource !== allowedResource + ) + if (span !== undefined) { + assert.fail(`unexpected query span: ${span.resource}`) + } + }, { timeoutMs: 200 }) +} + +/** + * @template T + * @param {string} resource + * @param {() => Promise} run + * @returns {Promise} + */ +async function assertQuerySpan (resource, run) { + const spanPromise = agent.assertSomeTraces(traces => { + const spans = traces.flat().filter(span => span.meta.component === 'postgres' && span.resource === resource) + assert.strictEqual(spans.length, 1) + }, { spanResourceMatch: resourcePattern(resource) }) + const [result] = await Promise.all([run(), spanPromise]) + return result +} + +/** + * @template T + * @param {Promise} query + * @returns {Promise} + */ +function executeQuery (query) { + return new Promise((resolve, reject) => query.then(resolve, reject)) +} + +describe('Plugin', () => { + describe('postgres', () => { + withVersions('postgres', 'postgres', version => { + let postgres + let resolvedVersion + let sql + let tracer + + before(() => agent.load('postgres')) + + beforeEach(() => { + const postgresModule = require(`../../../versions/postgres@${version}`) + postgres = postgresModule.get() + resolvedVersion = postgresModule.version() + sql = postgres(POSTGRES_TARGET) + tracer = require('../../dd-trace') + tracer.use('postgres', {}) + }) + + afterEach(() => sql.end({ timeout: 0 })) + + after(() => agent.close()) + + withPeerService( + () => tracer, + 'postgres', + () => executeQuery(sql`SELECT 1 AS value`), + 'postgres', + 'db.name', + { resource: 'SELECT 1 AS value' } + ) + + withNamingSchema( + () => sql`SELECT 1 AS value`, + rawExpectedSchema.outbound + ) + + it('instruments tagged queries without changing their result', async () => { + const spanPromise = agent.assertFirstTraceSpan({ + name: expectedSchema.outbound.opName, + service: expectedSchema.outbound.serviceName, + resource: 'SELECT $1::text AS message', + type: 'sql', + meta: { + component: 'postgres', + 'span.kind': 'client', + 'db.type': 'postgres', + 'db.name': 'postgres', + 'db.user': 'postgres', + 'out.host': '127.0.0.1', + }, + metrics: { + 'network.destination.port': 5432, + }, + }) + + const result = await sql`SELECT ${'Hello world!'}::text AS message` + + assert.strictEqual(result[0].message, 'Hello world!') + await spanPromise + }) + + it('does not trace connection bootstrap queries', async () => { + const noTracePromise = assertNoOtherQuerySpans('SELECT 1 AS value') + + await assertQuerySpan('SELECT 1 AS value', () => sql`SELECT 1 AS value`) + await noTracePromise + }) + + it('does not trace an unexecuted lazy query', async () => { + sql`SELECT 2 AS value` + await assertNoQuerySpan('SELECT 2 AS value') + }) + + it('does not trace a query cancelled before dispatch', async () => { + const cancelledSql = postgres(POSTGRES_TARGET) + const query = cancelledSql`SELECT 2 AS value` + const rejection = Promise.prototype.then.call(query) + const noTracePromise = assertNoQuerySpan('SELECT 2 AS value') + let starts = 0 + const onStart = ctx => { + if (ctx.query === query) starts++ + } + + postgresStartChannel.subscribe(onStart) + + try { + query.cancel() + await assert.rejects(rejection, { code: '57014', message: /canceling statement/ }) + + assert.strictEqual(Reflect.apply(Reflect.get(query, 'execute'), query, []), query) + await new Promise(setImmediate) + await cancelledSql.end({ timeout: 0 }) + } finally { + postgresStartChannel.unsubscribe(onStart) + } + + assert.strictEqual(starts, 0) + await noTracePromise + }) + + it('executes queued queries once without tracing while the plugin is disabled', async () => { + await sql`CREATE TEMP TABLE dd_trace_postgres_disabled (value int)` + + tracer.use('postgres', { enabled: false }) + let markStarted + const started = new Promise(resolve => { markStarted = resolve }) + const activeQuery = executeQuery(sql.unsafe('SELECT pg_sleep(0.05)', [], { + onexecute: () => { + markStarted() + return true + }, + })) + await started + + const resource = 'INSERT INTO dd_trace_postgres_disabled VALUES (1)' + const noTracePromise = assertNoQuerySpan(resource) + await Promise.all([activeQuery, sql.unsafe(resource)]) + await noTracePromise + + tracer.use('postgres', {}) + const result = await sql`SELECT count(*)::int AS count FROM dd_trace_postgres_disabled` + assert.strictEqual(result[0].count, 1) + }) + + it('uses the execution context of a lazy query', async () => { + const query = sql`SELECT ${1}::int AS value` + const tracesPromise = agent.assertSomeTraces(traces => { + const spans = traces[0] + const parent = spans.find(span => span.name === 'parent') + const child = spans.find(span => span.resource === 'SELECT $1::int AS value') + + assert.ok(parent) + assert.ok(child) + assert.strictEqual(child.parent_id.toString(), parent.span_id.toString()) + }) + + await tracer.trace('parent', () => executeQuery(query)) + await tracesPromise + }) + + it('preserves Query identity', async () => { + const query = sql`SELECT 1 AS value` + + assert.ok(query instanceof Promise) + assert.strictEqual(Reflect.apply(Reflect.get(query, 'execute'), query, []), query) + + await assertQuerySpan('SELECT 1 AS value', () => query) + }) + + it('uses the final query source after user mutation', async () => { + const query = sql.unsafe('SELECT 1 AS value') + Reflect.get(query, 'strings')[0] = 'SELECT 2 AS value' + Reflect.set(query, 'options', { prepare: false, simple: true }) + + const result = await assertQuerySpan('SELECT 2 AS value', () => query) + + assert.strictEqual(result[0].value, 2) + }) + + it('starts and finishes once when multiple Promise methods observe one query', async () => { + const tracesPromise = agent.assertSomeTraces(traces => { + const spans = traces[0] + const parent = spans.find(span => span.name === 'promise.parent') + const children = spans.filter(span => span.resource === 'SELECT 1 AS value') + + assert.ok(parent) + assert.strictEqual(children.length, 1) + assert.strictEqual(children[0].parent_id.toString(), parent.span_id.toString()) + }) + + await tracer.trace('promise.parent', () => { + const query = sql`SELECT 1 AS value` + return Promise.all([ + query.then(() => {}), + query.catch(() => {}), + query.finally(() => {}), + ]) + }) + await tracesPromise + }) + + it('handles query execution and result modifiers', async () => { + const unsafe = await assertQuerySpan( + 'SELECT $1::text AS message', + () => sql.unsafe('SELECT $1::text AS message', ['unsafe']) + ) + assert.strictEqual(unsafe[0].message, 'unsafe') + + const raw = await assertQuerySpan('SELECT 1 AS value', () => sql`SELECT 1 AS value`.raw()) + assert.ok(Buffer.isBuffer(raw[0][0])) + + const rows = [] + await assertQuerySpan('SELECT 1 AS value', () => sql`SELECT 1 AS value`.forEach(row => rows.push(row))) + assert.strictEqual(rows[0].value, 1) + + const description = await assertQuerySpan( + 'SELECT $1::int AS value', + () => sql`SELECT ${1}::int AS value`.describe() + ) + assert.strictEqual(description.columns[0].name, 'value') + + if (semver.gte(resolvedVersion, '3.4.0')) { + const simple = await assertQuerySpan('SELECT 1 AS value', () => sql`SELECT 1 AS value`.simple()) + assert.strictEqual(simple[0].value, 1) + + const values = await assertQuerySpan('SELECT 1 AS value', () => sql`SELECT 1 AS value`.values()) + assert.deepStrictEqual(values[0], [1]) + } + + const explicit = sql`SELECT 2 AS value`.execute() + const explicitResult = await assertQuerySpan('SELECT 2 AS value', () => explicit) + assert.strictEqual(explicitResult[0].value, 2) + + const notified = await assertQuerySpan( + 'select pg_notify($1, $2)', + () => sql.notify('dd_trace_postgres', 'payload') + ) + assert.ok(Array.isArray(notified)) + }) + + it('finishes callback, multi-batch, and early-return cursors once', async () => { + const callbackRows = [] + await assertQuerySpan( + 'SELECT generate_series(1, 3) AS value', + () => sql`SELECT generate_series(1, 3) AS value`.cursor(1, rows => callbackRows.push(...rows)) + ) + assert.strictEqual(callbackRows.length, 3) + + const iteratorRows = [] + await assertQuerySpan('SELECT generate_series(1, 3) AS value', async () => { + for await (const rows of sql`SELECT generate_series(1, 3) AS value`.cursor(1)) { + iteratorRows.push(...rows) + } + }) + assert.strictEqual(iteratorRows.length, 3) + + const partialRows = [] + await assertQuerySpan('SELECT generate_series(1, 4) AS value', async () => { + const iterator = sql`SELECT generate_series(1, 4) AS value`.cursor(1)[Symbol.asyncIterator]() + const first = await iterator.next() + partialRows.push(...first.value) + await iterator.return() + }) + assert.strictEqual(partialRows.length, 1) + }) + + it('reports errors after an async cursor replaces settlement functions', async () => { + const spanPromise = agent.assertFirstTraceSpan(span => { + assert.strictEqual(span.resource, 'INVALID CURSOR SQL') + assert.strictEqual(span.meta[ERROR_TYPE], 'PostgresError') + assert.match(span.meta[ERROR_MESSAGE], /syntax error/) + }, { spanResourceMatch: /^INVALID CURSOR SQL$/ }) + + await assert.rejects(async () => { + for await (const rows of sql.unsafe('INVALID CURSOR SQL').cursor(1)) { + assert.fail(`unexpected cursor rows: ${rows.length}`) + } + }, { name: 'PostgresError', message: /syntax error/ }) + await spanPromise + }) + + it('preserves COPY readable and writable streams', async () => { + const readableResource = 'COPY (SELECT 1 AS value) TO STDOUT' + const readableSpan = agent.assertFirstTraceSpan({ resource: readableResource }, { + spanResourceMatch: resourcePattern(readableResource), + }) + const readable = await sql.unsafe(readableResource).readable() + let output = '' + for await (const chunk of readable) { + output += chunk + } + assert.strictEqual(output, '1\n') + await readableSpan + + await assertQuerySpan( + 'CREATE TEMP TABLE dd_trace_postgres_copy (value int)', + () => sql.unsafe('CREATE TEMP TABLE dd_trace_postgres_copy (value int)') + ) + const writableResource = 'COPY dd_trace_postgres_copy (value) FROM STDIN' + const writable = await sql.unsafe(writableResource).writable() + await assertNoQuerySpan(writableResource) + + const writableSpan = agent.assertFirstTraceSpan({ resource: writableResource }, { + spanResourceMatch: resourcePattern(writableResource), + }) + const finished = once(writable, 'finish') + writable.end('1\n2\n') + await Promise.all([finished, writableSpan]) + + const result = await assertQuerySpan( + 'SELECT sum(value)::int AS value FROM dd_trace_postgres_copy', + () => sql.unsafe('SELECT sum(value)::int AS value FROM dd_trace_postgres_copy') + ) + assert.strictEqual(result[0].value, 3) + }) + + it('reports COPY errors after the writable stream is acquired', async () => { + await sql.unsafe('CREATE TEMP TABLE dd_trace_postgres_copy_error (value int)') + + const resource = 'COPY dd_trace_postgres_copy_error (value) FROM STDIN' + const writable = await sql.unsafe(resource).writable() + const spanPromise = agent.assertFirstTraceSpan(span => { + assert.strictEqual(span.resource, resource) + assert.strictEqual(span.meta[ERROR_TYPE], 'PostgresError') + assert.match(span.meta[ERROR_MESSAGE], /invalid input syntax for (?:type )?integer/) + }, { spanResourceMatch: resourcePattern(resource) }) + writable.once('error', () => {}) + + writable.end('invalid\n') + + await spanPromise + writable.destroy() + }) + + it('instruments file queries without tracing pre-dispatch file errors', async () => { + const result = await assertQuerySpan( + 'SELECT $1::text AS message\n', + () => sql.file(path.join(__dirname, 'fixtures', 'select.sql'), ['file']) + ) + assert.strictEqual(result[0].message, 'file') + + const missingFile = path.join(__dirname, 'fixtures', 'missing.sql') + const noTracePromise = assertNoOtherQuerySpans('SELECT 1 AS value') + + await assert.rejects(sql.file(missingFile), { code: 'ENOENT' }) + await assertQuerySpan('SELECT 1 AS value', () => sql`SELECT 1 AS value`) + await noTracePromise + }) + + it('reports server and build errors without leaking spans', async () => { + const sqlErrorSpanPromise = agent.assertFirstTraceSpan(span => { + assert.strictEqual(span.resource, 'INVALID SQL') + assert.match(span.meta[ERROR_MESSAGE], /syntax error/) + assert.strictEqual(span.meta[ERROR_TYPE], 'PostgresError') + assert.strictEqual(typeof span.meta[ERROR_STACK], 'string') + }) + + await assert.rejects(sql.unsafe('INVALID SQL'), { name: 'PostgresError', message: /syntax error/ }) + await sqlErrorSpanPromise + + const buildErrorSpanPromise = agent.assertFirstTraceSpan(span => { + assert.strictEqual(span.name, expectedSchema.outbound.opName) + assert.match(span.meta[ERROR_MESSAGE], /Undefined values are not allowed/) + }) + + await assert.rejects(sql`SELECT ${undefined}`, { message: /Undefined values are not allowed/ }) + await buildErrorSpanPromise + }) + + it('preserves one span when Postgres.js retries a prepared statement', async () => { + await assertQuerySpan('SELECT $1::int AS value', () => sql`SELECT ${1}::int AS value`) + await assertQuerySpan('DISCARD ALL', () => sql.unsafe('DISCARD ALL')) + + const result = await assertQuerySpan('SELECT $1::int AS value', () => sql`SELECT ${2}::int AS value`) + assert.strictEqual(result[0].value, 2) + }) + + it('finishes active and queued cancellations with their original errors', async function () { + let startFirst + const firstStarted = new Promise(resolve => { startFirst = resolve }) + const first = sql.unsafe('SELECT pg_sleep(10)', [], { + onexecute: () => { + startFirst() + return true + }, + }) + const firstSpanPromise = agent.assertFirstTraceSpan(span => { + assert.strictEqual(span.resource, 'SELECT pg_sleep(10)') + assert.strictEqual(span.meta[ERROR_TYPE], 'PostgresError') + assert.match(span.meta[ERROR_MESSAGE], /canceling statement/) + }, { spanResourceMatch: /^SELECT pg_sleep\(10\)$/ }) + const firstRejection = assert.rejects(first, { code: '57014', message: /canceling statement/ }) + + await firstStarted + + const queued = sql.unsafe('SELECT 2 AS value') + const queuedSpanPromise = agent.assertFirstTraceSpan(span => { + assert.strictEqual(span.resource, 'SELECT 2 AS value') + assert.strictEqual(span.meta[ERROR_TYPE], 'Error') + assert.match(span.meta[ERROR_MESSAGE], /canceling statement/) + }, { spanResourceMatch: /^SELECT 2 AS value$/ }) + const queuedRejection = assert.rejects(queued, { code: '57014', message: /canceling statement/ }) + + await new Promise(setImmediate) + await queued.cancel() + await Promise.all([queuedRejection, queuedSpanPromise]) + + await first.cancel() + await Promise.all([firstRejection, firstSpanPromise]) + }).timeout(10000) + + it('reports handler errors before a connection executes the query', async () => { + await sql.end({ timeout: 0 }) + + const spanPromise = agent.assertFirstTraceSpan(span => { + assert.strictEqual(span.resource, 'SELECT 1 AS value') + assert.strictEqual(span.meta[ERROR_TYPE], 'Error') + assert.strictEqual(span.meta[ERROR_MESSAGE], 'write CONNECTION_ENDED 127.0.0.1:5432') + }, { spanResourceMatch: /^SELECT 1 AS value$/ }) + + await assert.rejects(sql.unsafe('SELECT 1 AS value'), { code: 'CONNECTION_ENDED' }) + await spanPromise + + const taggedSpanPromise = agent.assertFirstTraceSpan(span => { + assert.strictEqual(span.resource, 'SELECT 2 AS value') + assert.strictEqual(span.meta[ERROR_TYPE], 'Error') + assert.strictEqual(span.meta[ERROR_MESSAGE], 'write CONNECTION_ENDED 127.0.0.1:5432') + }, { spanResourceMatch: /^SELECT 2 AS value$/ }) + + await assert.rejects(sql`SELECT 2 AS value`, { code: 'CONNECTION_ENDED' }) + await taggedSpanPromise + }) + + it('instruments transaction and reserved-connection handlers', async () => { + const transactionSpanPromise = agent.assertFirstTraceSpan( + { resource: 'SELECT 42 AS value' }, + { spanResourceMatch: /^SELECT 42 AS value$/ } + ) + const transactionResult = await sql.begin(transaction => transaction`SELECT 42 AS value`) + + assert.strictEqual(transactionResult[0].value, 42) + await transactionSpanPromise + + if (typeof sql.reserve !== 'function') return + + const reserved = await sql.reserve() + try { + const reservedResult = await assertQuerySpan('SELECT 43 AS value', () => reserved`SELECT 43 AS value`) + assert.strictEqual(reservedResult[0].value, 43) + } finally { + reserved.release() + } + }) + + it('keeps concurrent pipelined queries under their execution parent', async () => { + const tracesPromise = agent.assertSomeTraces(traces => { + const spans = traces[0] + const parent = spans.find(span => span.name === 'pipeline.parent') + const children = spans.filter(span => span.meta.component === 'postgres') + + assert.ok(parent) + assert.strictEqual(children.length, 2) + for (const child of children) { + assert.strictEqual(child.parent_id.toString(), parent.span_id.toString()) + } + }) + + const result = await tracer.trace('pipeline.parent', () => Promise.all([ + sql`SELECT 1 AS value`, + sql`SELECT 2 AS value`, + ])) + + assert.strictEqual(result[0][0].value, 1) + assert.strictEqual(result[1][0].value, 2) + await tracesPromise + }) + + it('omits endpoint tags when Postgres.js can fail over between hosts', async () => { + const multiHost = postgres({ + ...POSTGRES_TARGET, + host: ['127.0.0.1', '127.0.0.2'], + port: [5432, 5432], + }) + const spanPromise = agent.assertFirstTraceSpan(span => { + assert.strictEqual(span.meta['out.host'], undefined) + assert.strictEqual(span.metrics['network.destination.port'], undefined) + }, { spanResourceMatch: /^SELECT 1 AS value$/ }) + + try { + await multiHost`SELECT 1 AS value` + await spanPromise + } finally { + await multiHost.end({ timeout: 0 }) + } + }) + + it('omits TCP endpoint tags for Unix-domain sockets', async () => { + const socketPath = path.join(__dirname, 'fixtures', 'missing.sock') + const unixSocket = postgres({ ...POSTGRES_TARGET, path: socketPath }) + const spanPromise = agent.assertFirstTraceSpan(span => { + assert.strictEqual(span.meta['out.host'], undefined) + assert.strictEqual(span.metrics['network.destination.port'], undefined) + }, { spanResourceMatch: /^SELECT 1 AS value$/ }) + + try { + await Promise.all([ + assert.rejects(unixSocket`SELECT 1 AS value`, { code: 'ENOENT' }), + spanPromise, + ]) + } finally { + await unixSocket.end({ timeout: 0 }) + } + }) + + it('supports a configured service name', async () => { + tracer.use('postgres', { service: 'custom-postgres' }) + const spanPromise = agent.assertFirstTraceSpan({ service: 'custom-postgres' }) + + await sql`SELECT 1 AS value` + await spanPromise + }) + + it('truncates the first resource beyond the configured boundary', async () => { + tracer.use('postgres', { truncate: 12 }) + + await assertQuerySpan('SELECT 12345', () => sql.unsafe('SELECT 12345')) + + const spanPromise = agent.assertFirstTraceSpan( + { resource: 'SELECT 12...' }, + { spanResourceMatch: /^SELECT 12\.\.\.$/ } + ) + await sql.unsafe('SELECT 123456') + await spanPromise + }) + }) + }) +}) diff --git a/packages/datadog-plugin-postgres/test/integration-test/client.spec.js b/packages/datadog-plugin-postgres/test/integration-test/client.spec.js new file mode 100644 index 00000000000..69ce508cdc5 --- /dev/null +++ b/packages/datadog-plugin-postgres/test/integration-test/client.spec.js @@ -0,0 +1,76 @@ +'use strict' + +const assert = require('node:assert/strict') +const { inspect } = require('node:util') + +const { + checkSpansForServiceName, + FakeAgent, + sandboxCwd, + spawnPluginIntegrationTestProcAndExpectExit, + stopProc, + useSandbox, + varySandbox, +} = require('../../../../integration-tests/helpers') +const { withVersions } = require('../../../dd-trace/test/setup/mocha') + +describe('esm', () => { + let agent + let proc + + withVersions('postgres', 'postgres', version => { + useSandbox([`'postgres@${version}'`], false, [ + './packages/datadog-plugin-postgres/test/integration-test/*', + ]) + + const variants = varySandbox('server.mjs', { + bindingName: 'postgres', + defaultExport: true, + namedExports: [], + packageName: 'postgres', + }) + + beforeEach(async () => { + agent = await new FakeAgent().start() + }) + + afterEach(async () => { + await stopProc(proc) + await agent.stop() + }) + + for (const variant of Object.keys(variants)) { + it(`is instrumented loaded with ${variant}`, async () => { + const traceReceived = agent.assertMessageReceived(({ headers, payload }) => { + assert.strictEqual(headers.host, `127.0.0.1:${agent.port}`) + assert.ok(Array.isArray(payload), `Expected array, got ${inspect(payload)}`) + assert.strictEqual(checkSpansForServiceName(payload, 'postgres.query'), true) + }) + + proc = await spawnPluginIntegrationTestProcAndExpectExit(sandboxCwd(), variants[variant], agent.port) + + await traceReceived + }).timeout(20000) + } + + it('does not trace when the Query class loads before tracing', async () => { + const messages = [] + const onMessage = message => messages.push(message) + agent.on('message', onMessage) + + try { + proc = await spawnPluginIntegrationTestProcAndExpectExit( + sandboxCwd(), + 'partial-rewrite.cjs', + agent.port, + { NODE_OPTIONS: '--require=./preload-query.cjs' } + ) + } finally { + agent.removeListener('message', onMessage) + } + + const spans = messages.flatMap(({ payload }) => Array.isArray(payload) ? payload.flat() : []) + assert.strictEqual(spans.some(span => span.meta?.component === 'postgres'), false) + }).timeout(20000) + }) +}) diff --git a/packages/datadog-plugin-postgres/test/integration-test/partial-rewrite.cjs b/packages/datadog-plugin-postgres/test/integration-test/partial-rewrite.cjs new file mode 100644 index 00000000000..d168c91ce0e --- /dev/null +++ b/packages/datadog-plugin-postgres/test/integration-test/partial-rewrite.cjs @@ -0,0 +1,22 @@ +'use strict' + +require('dd-trace').init() // eslint-disable-line n/no-missing-require +const postgres = require('postgres') // eslint-disable-line n/no-missing-require + +const sql = postgres({ + database: 'postgres', + host: 'localhost', + password: 'postgres', + port: 5432, + user: 'postgres', +}) + +async function run () { + const result = await sql`SELECT 1 AS value` + if (result[0].value !== 1) throw new Error('unexpected query result') + await sql.end() +} + +run().catch(error => { + process.nextTick(() => { throw error }) +}) diff --git a/packages/datadog-plugin-postgres/test/integration-test/preload-query.cjs b/packages/datadog-plugin-postgres/test/integration-test/preload-query.cjs new file mode 100644 index 00000000000..1589845e795 --- /dev/null +++ b/packages/datadog-plugin-postgres/test/integration-test/preload-query.cjs @@ -0,0 +1,7 @@ +'use strict' + +const path = require('node:path') + +// eslint-disable-next-line n/no-missing-require +const entrypoint = require.resolve('postgres') +require(path.join(path.dirname(entrypoint), 'query.js')) diff --git a/packages/datadog-plugin-postgres/test/integration-test/server.mjs b/packages/datadog-plugin-postgres/test/integration-test/server.mjs new file mode 100644 index 00000000000..19b12354a06 --- /dev/null +++ b/packages/datadog-plugin-postgres/test/integration-test/server.mjs @@ -0,0 +1,13 @@ +import 'dd-trace/init.js' +import postgres from 'postgres' + +const sql = postgres({ + database: 'postgres', + host: 'localhost', + password: 'postgres', + port: 5432, + user: 'postgres', +}) + +await sql`SELECT 1 AS value` +await sql.end() diff --git a/packages/datadog-plugin-postgres/test/naming.js b/packages/datadog-plugin-postgres/test/naming.js new file mode 100644 index 00000000000..6eeb239ff7d --- /dev/null +++ b/packages/datadog-plugin-postgres/test/naming.js @@ -0,0 +1,21 @@ +'use strict' + +const { resolveNaming } = require('../../dd-trace/test/plugins/helpers') + +const rawExpectedSchema = { + outbound: { + v0: { + opName: 'postgres.query', + serviceName: 'test-postgres', + }, + v1: { + opName: 'postgresql.query', + serviceName: 'test', + }, + }, +} + +module.exports = { + rawExpectedSchema, + expectedSchema: resolveNaming(rawExpectedSchema), +} diff --git a/packages/dd-trace/src/config/generated-config-types.d.ts b/packages/dd-trace/src/config/generated-config-types.d.ts index fb48c8e21a8..f10ca3fa13d 100644 --- a/packages/dd-trace/src/config/generated-config-types.d.ts +++ b/packages/dd-trace/src/config/generated-config-types.d.ts @@ -358,6 +358,7 @@ export interface GeneratedConfig { DD_TRACE_PLAYWRIGHT_CORE_ENABLED: boolean; DD_TRACE_PLAYWRIGHT_ENABLED: boolean; DD_TRACE_PLAYWRIGHT_TEST_ENABLED: boolean; + DD_TRACE_POSTGRES_ENABLED: boolean; DD_TRACE_PRISMA_ENABLED: boolean; DD_TRACE_PROCESS_ENABLED: boolean; DD_TRACE_PROMISE_ENABLED: boolean; @@ -1071,6 +1072,7 @@ export interface GeneratedEnvVarConfig { DD_TRACE_PLAYWRIGHT_CORE_ENABLED: boolean; DD_TRACE_PLAYWRIGHT_ENABLED: boolean; DD_TRACE_PLAYWRIGHT_TEST_ENABLED: boolean; + DD_TRACE_POSTGRES_ENABLED: boolean; DD_TRACE_PRISMA_ENABLED: boolean; DD_TRACE_PROCESS_ENABLED: boolean; DD_TRACE_PROMISE_ENABLED: boolean; diff --git a/packages/dd-trace/src/config/supported-configurations.json b/packages/dd-trace/src/config/supported-configurations.json index dfb55d46419..e8fdefa449b 100644 --- a/packages/dd-trace/src/config/supported-configurations.json +++ b/packages/dd-trace/src/config/supported-configurations.json @@ -3693,6 +3693,13 @@ "default": "true" } ], + "DD_TRACE_POSTGRES_ENABLED": [ + { + "implementation": "A", + "type": "boolean", + "default": "true" + } + ], "DD_TRACE_PRISMA_ENABLED": [ { "implementation": "A", diff --git a/packages/dd-trace/src/plugins/index.js b/packages/dd-trace/src/plugins/index.js index c023683adca..087edbf837d 100644 --- a/packages/dd-trace/src/plugins/index.js +++ b/packages/dd-trace/src/plugins/index.js @@ -120,6 +120,7 @@ const plugins = { get pino () { return require('../../../datadog-plugin-pino/src') }, get 'pino-pretty' () { return require('../../../datadog-plugin-pino/src') }, get playwright () { return require('../../../datadog-plugin-playwright/src') }, + get postgres () { return require('../../../datadog-plugin-postgres/src') }, get protobufjs () { return require('../../../datadog-plugin-protobufjs/src') }, get redis () { return require('../../../datadog-plugin-redis/src') }, get restify () { return require('../../../datadog-plugin-restify/src') }, diff --git a/packages/dd-trace/src/plugins/tracing.js b/packages/dd-trace/src/plugins/tracing.js index dad4c6bdbe2..95d5be3ecac 100644 --- a/packages/dd-trace/src/plugins/tracing.js +++ b/packages/dd-trace/src/plugins/tracing.js @@ -29,6 +29,7 @@ class TracingPlugin extends Plugin { * @param {string} [opts.type] * @param {string} [opts.id] * @param {string} [opts.kind] + * @param {object} [opts.pluginConfig] * @returns {{ name: string, source: string | undefined }} */ serviceName (opts = {}) { diff --git a/packages/dd-trace/src/service-naming/schemas/v0/storage.js b/packages/dd-trace/src/service-naming/schemas/v0/storage.js index 114157c7806..149c65c9774 100644 --- a/packages/dd-trace/src/service-naming/schemas/v0/storage.js +++ b/packages/dd-trace/src/service-naming/schemas/v0/storage.js @@ -156,6 +156,13 @@ const storage = { return optionServiceSource({ tracerService, pluginConfig, connectionName, integration: 'pg' }) }, }, + postgres: { + opName: () => 'postgres.query', + serviceName: withSuffixFunction('postgres'), + serviceSource: ({ tracerService, pluginConfig, connectionName }) => { + return optionServiceSource({ tracerService, pluginConfig, connectionName, integration: 'postgres' }) + }, + }, prisma: { opName: ({ operation }) => `prisma.${operation}`, serviceName: withSuffixFunction('prisma'), diff --git a/packages/dd-trace/src/service-naming/schemas/v1/storage.js b/packages/dd-trace/src/service-naming/schemas/v1/storage.js index 30eccf6065d..f621b151be8 100644 --- a/packages/dd-trace/src/service-naming/schemas/v1/storage.js +++ b/packages/dd-trace/src/service-naming/schemas/v1/storage.js @@ -86,6 +86,11 @@ const storage = { serviceName: withFunction, serviceSource: optionServiceSource, }, + postgres: { + opName: () => 'postgresql.query', + serviceName: configWithFallback, + serviceSource: optionServiceSource, + }, prisma: { opName: ({ operation }) => `prisma.${operation}`, serviceName: configWithFallback, diff --git a/packages/dd-trace/test/plugins/agent.js b/packages/dd-trace/test/plugins/agent.js index d5a49082b99..f2b008cf097 100644 --- a/packages/dd-trace/test/plugins/agent.js +++ b/packages/dd-trace/test/plugins/agent.js @@ -199,9 +199,9 @@ function isMatchingTrace (spans, spanResourceMatch) { * without it having to be `traces[0][0]`. Without a matcher the first span of the * first trace is used. * - * @param {import('../../src/opentracing/span')[][]} traces + * @param {SerializedSpan[][]} traces * @param {RegExp} [spanResourceMatch] - * @returns {import('../../src/opentracing/span')} + * @returns {SerializedSpan} */ function findFirstTraceSpan (traces, spanResourceMatch) { if (spanResourceMatch) { @@ -287,8 +287,8 @@ function dsmStatsExistWithParentHash (agent, expectedParentHash) { /** * Unformats span events. * - * @param {import('../../src/opentracing/span')} span - * @returns {import('../../src/opentracing/span')[]} + * @param {SerializedSpan} span + * @returns {SerializedSpanEvent[]} */ function unformatSpanEvents (span) { if (span.meta?.events) { @@ -339,7 +339,7 @@ function getCurrentIntegrationName () { } /** - * @param {import('../../src/opentracing/span')[][]} traces + * @param {SerializedSpan[][]} traces */ function assertIntegrationName (traces) { // we want to assert that all spans generated by an instrumentation have the right `_dd.integration` tag set @@ -382,9 +382,25 @@ let availableEndpoints = DEFAULT_AVAILABLE_ENDPOINTS * @property {number} [timeoutMs=1000] - The timeout in ms. * @property {boolean} [rejectFirst=false] - If true, reject the first time the callback throws. * @property {RegExp} [spanResourceMatch] - A regex to match against the span resource. - * @typedef {import('../../src/opentracing/span')} Span + * @typedef {object} SerializedSpan + * @property {bigint} trace_id + * @property {bigint} span_id + * @property {bigint} parent_id + * @property {string} name + * @property {string} resource + * @property {string} service + * @property {string} type + * @property {number} error + * @property {Record} meta + * @property {Record} metrics + * @property {number} start + * @property {number} duration + * @typedef {object} SerializedSpanEvent + * @property {string} name + * @property {number} startTime + * @property {Record | undefined} attributes * For a given payload, an array of traces, each trace is an array of spans. - * @typedef {(traces: Span[][]) => void} TracesCallback + * @typedef {(traces: SerializedSpan[][]) => void} TracesCallback * @typedef {(agentlessPayload: {events: Event[]}, request: Request) => void} AgentlessCallback * @typedef {TracesCallback | AgentlessCallback} RunCallbackAgainstTracesCallback */ @@ -772,7 +788,7 @@ module.exports = { * Callback for running test assertions against a span. * * @callback testAssertionSpanCallback - * @param {Span} span - For a given payload, the first span of the first trace. + * @param {SerializedSpan} span - For a given payload, the first span of the first trace. */ /** diff --git a/packages/dd-trace/test/plugins/externals.js b/packages/dd-trace/test/plugins/externals.js index b285d4cc3ee..ef828c0557a 100644 --- a/packages/dd-trace/test/plugins/externals.js +++ b/packages/dd-trace/test/plugins/externals.js @@ -659,6 +659,12 @@ module.exports = { versions: ['8.0.0'], }, ], + postgres: [ + { + name: 'pg', + versions: ['>=8.0.3'], + }, + ], '@prisma/client': [ { name: 'prisma', diff --git a/packages/dd-trace/test/plugins/versions/package.json b/packages/dd-trace/test/plugins/versions/package.json index 19d1feb2ff8..1e64583483e 100644 --- a/packages/dd-trace/test/plugins/versions/package.json +++ b/packages/dd-trace/test/plugins/versions/package.json @@ -214,6 +214,7 @@ "playwright": "1.63.0", "playwright-core": "1.63.0", "pnpm": "11.25.0", + "postgres": "3.4.9", "prisma": "7.10.0", "promise": "8.3.0", "promise-js": "0.0.7", diff --git a/supported_versions_output.json b/supported_versions_output.json index 1fa3e2a5f41..486cecdfd78 100644 --- a/supported_versions_output.json +++ b/supported_versions_output.json @@ -713,6 +713,13 @@ "max_tracer_supported": "1.63.0", "auto-instrumented": "True" }, + { + "dependency": "postgres", + "integration": "postgres", + "minimum_tracer_supported": "3.0.0", + "max_tracer_supported": "3.4.9", + "auto-instrumented": "True" + }, { "dependency": "protobufjs", "integration": "protobufjs", diff --git a/supported_versions_table.csv b/supported_versions_table.csv index 139774e0067..cdfa268ce1d 100644 --- a/supported_versions_table.csv +++ b/supported_versions_table.csv @@ -101,6 +101,7 @@ pg,pg,8.0.3,8.23.0,True pino,pino,2.0.0,10.3.1,True pino-pretty,pino,1.0.0,13.1.3,True playwright,playwright,1.38.0,1.63.0,True +postgres,postgres,3.0.0,3.4.9,True protobufjs,protobufjs,6.8.0,8.8.0,True redis,redis,0.12.0,6.2.1,True restify,restify,3.0.0,12.0.0,True