From de0664adae85395b160cee168999ead4ec9fbcbb Mon Sep 17 00:00:00 2001 From: Devon Stewart Date: Fri, 12 Jun 2026 00:37:16 -0700 Subject: [PATCH 1/2] Fix parallel hash span creation --- expected/parallel.out | 9 +++++++++ regress/16/expected/parallel.out | 9 +++++++++ regress/17/expected/parallel.out | 9 +++++++++ regress/18/expected/parallel.out | 9 +++++++++ sql/parallel.sql | 10 ++++++++++ src/pg_tracing_planstate.c | 8 +++++--- 6 files changed, 51 insertions(+), 3 deletions(-) diff --git a/expected/parallel.out b/expected/parallel.out index 2e0b053..eeb49b8 100644 --- a/expected/parallel.out +++ b/expected/parallel.out @@ -70,6 +70,15 @@ set parallel_setup_cost=0; set parallel_tuple_cost=0; set min_parallel_table_scan_size=0; set max_parallel_workers_per_gather=2; +-- Test a sampled parallel hash join. Parallel Hash nodes may have +-- instrumentation even when their child was executed by another worker. +set enable_nestloop=false; +set enable_mergejoin=false; +/*dddbs='postgres.db',traceparent='00-00000000000000000000000000000001-0000000000000005-01'*/ select count(*) from pg_tracing_test r join pg_tracing_test s using (a) \gset +\echo :count +10000 +CALL clean_spans(); +-- Test leaderless parallel query set parallel_leader_participation=false; /*dddbs='postgres.db',traceparent='00-00000000000000000000000000000001-0000000000000001-01'*/ select 1 from pg_class limit 1; ?column? diff --git a/regress/16/expected/parallel.out b/regress/16/expected/parallel.out index 2e0b053..eeb49b8 100644 --- a/regress/16/expected/parallel.out +++ b/regress/16/expected/parallel.out @@ -70,6 +70,15 @@ set parallel_setup_cost=0; set parallel_tuple_cost=0; set min_parallel_table_scan_size=0; set max_parallel_workers_per_gather=2; +-- Test a sampled parallel hash join. Parallel Hash nodes may have +-- instrumentation even when their child was executed by another worker. +set enable_nestloop=false; +set enable_mergejoin=false; +/*dddbs='postgres.db',traceparent='00-00000000000000000000000000000001-0000000000000005-01'*/ select count(*) from pg_tracing_test r join pg_tracing_test s using (a) \gset +\echo :count +10000 +CALL clean_spans(); +-- Test leaderless parallel query set parallel_leader_participation=false; /*dddbs='postgres.db',traceparent='00-00000000000000000000000000000001-0000000000000001-01'*/ select 1 from pg_class limit 1; ?column? diff --git a/regress/17/expected/parallel.out b/regress/17/expected/parallel.out index 2e0b053..eeb49b8 100644 --- a/regress/17/expected/parallel.out +++ b/regress/17/expected/parallel.out @@ -70,6 +70,15 @@ set parallel_setup_cost=0; set parallel_tuple_cost=0; set min_parallel_table_scan_size=0; set max_parallel_workers_per_gather=2; +-- Test a sampled parallel hash join. Parallel Hash nodes may have +-- instrumentation even when their child was executed by another worker. +set enable_nestloop=false; +set enable_mergejoin=false; +/*dddbs='postgres.db',traceparent='00-00000000000000000000000000000001-0000000000000005-01'*/ select count(*) from pg_tracing_test r join pg_tracing_test s using (a) \gset +\echo :count +10000 +CALL clean_spans(); +-- Test leaderless parallel query set parallel_leader_participation=false; /*dddbs='postgres.db',traceparent='00-00000000000000000000000000000001-0000000000000001-01'*/ select 1 from pg_class limit 1; ?column? diff --git a/regress/18/expected/parallel.out b/regress/18/expected/parallel.out index 2e0b053..eeb49b8 100644 --- a/regress/18/expected/parallel.out +++ b/regress/18/expected/parallel.out @@ -70,6 +70,15 @@ set parallel_setup_cost=0; set parallel_tuple_cost=0; set min_parallel_table_scan_size=0; set max_parallel_workers_per_gather=2; +-- Test a sampled parallel hash join. Parallel Hash nodes may have +-- instrumentation even when their child was executed by another worker. +set enable_nestloop=false; +set enable_mergejoin=false; +/*dddbs='postgres.db',traceparent='00-00000000000000000000000000000001-0000000000000005-01'*/ select count(*) from pg_tracing_test r join pg_tracing_test s using (a) \gset +\echo :count +10000 +CALL clean_spans(); +-- Test leaderless parallel query set parallel_leader_participation=false; /*dddbs='postgres.db',traceparent='00-00000000000000000000000000000001-0000000000000001-01'*/ select 1 from pg_class limit 1; ?column? diff --git a/sql/parallel.sql b/sql/parallel.sql index 9a139f1..7d52b42 100644 --- a/sql/parallel.sql +++ b/sql/parallel.sql @@ -41,6 +41,16 @@ set parallel_setup_cost=0; set parallel_tuple_cost=0; set min_parallel_table_scan_size=0; set max_parallel_workers_per_gather=2; + +-- Test a sampled parallel hash join. Parallel Hash nodes may have +-- instrumentation even when their child was executed by another worker. +set enable_nestloop=false; +set enable_mergejoin=false; +/*dddbs='postgres.db',traceparent='00-00000000000000000000000000000001-0000000000000005-01'*/ select count(*) from pg_tracing_test r join pg_tracing_test s using (a) \gset +\echo :count +CALL clean_spans(); + +-- Test leaderless parallel query set parallel_leader_participation=false; /*dddbs='postgres.db',traceparent='00-00000000000000000000000000000001-0000000000000001-01'*/ select 1 from pg_class limit 1; diff --git a/src/pg_tracing_planstate.c b/src/pg_tracing_planstate.c index 184b457..4a78e18 100644 --- a/src/pg_tracing_planstate.c +++ b/src/pg_tracing_planstate.c @@ -395,10 +395,12 @@ create_spans_from_planstate(PlanState *planstate, planstateTraceContext * planst span_id = generate_rnd_uint64(); break; case T_Hash: - /* For hash node, use the child's start */ + /* The child may be executed by another parallel worker */ traced_planstate = get_traced_planstate(outerPlanState(planstate)); - Assert(traced_planstate != NULL); - span_start = traced_planstate->node_start; + if (traced_planstate != NULL) + span_start = traced_planstate->node_start; + else + span_start = parent_start; /* * We still need to generate a dedicated span_id since From fe511e8455c690290b27c948acce88f2660afc82 Mon Sep 17 00:00:00 2001 From: Devon Stewart Date: Fri, 18 Sep 2026 16:06:35 -0700 Subject: [PATCH 2/2] Fix tracing lifecycle and exporter reliability --- .github/workflows/run-tests.yml | 2 +- Makefile | 4 +- doc/pg_tracing.md | 6 +- src/pg_tracing.c | 55 +++++++++++++---- src/pg_tracing.h | 7 ++- src/pg_tracing_active_spans.c | 14 ++++- src/pg_tracing_json.c | 2 +- src/pg_tracing_operation_hash.c | 14 +++-- src/pg_tracing_otel.c | 10 ++- src/pg_tracing_planstate.c | 37 ++++++++--- src/pg_tracing_span.c | 4 +- src/pg_tracing_sql_functions.c | 2 +- t/004_reliability.pl | 106 ++++++++++++++++++++++++++++++++ t/005_export_retry.pl | 79 ++++++++++++++++++++++++ 14 files changed, 298 insertions(+), 44 deletions(-) create mode 100644 t/004_reliability.pl create mode 100644 t/005_export_retry.pl diff --git a/.github/workflows/run-tests.yml b/.github/workflows/run-tests.yml index 3c4f7a4..c427f3c 100644 --- a/.github/workflows/run-tests.yml +++ b/.github/workflows/run-tests.yml @@ -66,5 +66,5 @@ jobs: run: make run-pgindent-diff - name: Run Test - timeout-minutes: 1 + timeout-minutes: 5 run: make run-test diff --git a/Makefile b/Makefile index 51a1386..4d13ebb 100644 --- a/Makefile +++ b/Makefile @@ -27,6 +27,8 @@ ifeq ($(shell test $(PG_VERSION) -le 14; echo $$?),0) OBJS += src/pg_prng.o endif +TAP_TESTS = 1 + ifdef PG_CONFIG_EXISTS PGXS := $(shell $(PG_CONFIG) --pgxs) include $(PGXS) @@ -61,8 +63,6 @@ REGRESSCHECKS += sample planstate planstate_bitmap planstate_hash \ planstate_union parallel subxact full_buffer \ guc nested wal cleanup -TAP_TESTS = 1 - regresscheck_noinstall: $(pg_regress_check) $(REGRESSCHECKS_OPTS) $(REGRESSCHECKS) || \ (cat regression.diffs && exit 1) diff --git a/doc/pg_tracing.md b/doc/pg_tracing.md index 79eddb3..980abe5 100644 --- a/doc/pg_tracing.md +++ b/doc/pg_tracing.md @@ -203,4 +203,8 @@ Service name to set in traces sent to otel collector. This parameter can only be ### pg_tracing.otel_connect_timeout_ms (integer) -Maximum time in milliseconds to connect to the otel collector. This includes DNS resolution and protocol handshake. This parameter can only be set at server start. +Maximum time in milliseconds to connect to the otel collector. This includes DNS resolution and protocol handshake. The default is 1000. Changes take effect on configuration reload. + +### pg_tracing.otel_timeout_ms (integer) + +Maximum time in milliseconds for a complete export request, including connection establishment and waiting for the collector response. The default is 10000. Changes take effect on configuration reload. Only HTTP 2xx responses count as successful exports; failed or timed-out requests retain their payload for retry. diff --git a/src/pg_tracing.c b/src/pg_tracing.c index 8935c4a..e976e7a 100644 --- a/src/pg_tracing.c +++ b/src/pg_tracing.c @@ -130,6 +130,7 @@ int pg_tracing_otel_naptime; /* Delay between upload of spans to * otel collector */ int pg_tracing_otel_connect_timeout_ms; /* Connect timeout to the otel * collector */ +int pg_tracing_otel_timeout_ms; /* Total timeout for an otel request */ static char *guc_tracecontext_str = NULL; /* Trace context string propagated * through GUC variable */ @@ -475,6 +476,19 @@ _PG_init(void) &otel_config_int_assign_hook, NULL); + DefineCustomIntVariable("pg_tracing.otel_timeout_ms", + "Maximum time in milliseconds for a complete otel export request.", + NULL, + &pg_tracing_otel_timeout_ms, + 10000, + 100, + 600000, + PGC_SIGHUP, + 0, + NULL, + &otel_config_int_assign_hook, + NULL); + DefineCustomStringVariable("pg_tracing.otel_endpoint", "Otel endpoint to send spans.", "If unset, no background worker to export to otel is created.", @@ -1343,12 +1357,15 @@ end_tracing(void) if (span->type >= SPAN_TOP_SELECT && span->type <= SPAN_TOP_UNKNOWN) span->operation_name_offset = lookup_operation_name(span, operation_name_buffer->data + span->operation_name_offset); else - span->operation_name_offset += plan_name_pos; + span->operation_name_offset = plan_name_pos == (Size) -1 ? -1 : + span->operation_name_offset + plan_name_pos; } if (span->parameter_offset != -1) - span->parameter_offset += parameter_pos; + span->parameter_offset = parameter_pos == (Size) -1 ? -1 : + span->parameter_offset + parameter_pos; if (span->deparse_info_offset != -1) - span->deparse_info_offset += parameter_pos; + span->deparse_info_offset = parameter_pos == (Size) -1 ? -1 : + span->deparse_info_offset + parameter_pos; add_span_to_shared_buffer_locked(span); } @@ -1429,7 +1446,7 @@ initialize_trace_level(void) oldcxt = MemoryContextSwitchTo(pg_tracing_mem_ctx); /* initial allocation */ - allocated_nested_level = 1; + allocated_nested_level = nested_level + 1; per_level_infos = palloc0(allocated_nested_level * sizeof(PerLevelInfos)); current_trace_spans = palloc0(sizeof(pgTracingSpans) + @@ -1446,7 +1463,7 @@ initialize_trace_level(void) /* New nested level, allocate more memory */ int old_allocated_nested_level = allocated_nested_level; - allocated_nested_level++; + allocated_nested_level = nested_level + 1; /* repalloc uses the pointer's memory context, no need to switch */ per_level_infos = repalloc0(per_level_infos, old_allocated_nested_level * sizeof(PerLevelInfos), allocated_nested_level * sizeof(PerLevelInfos)); @@ -1729,7 +1746,8 @@ pg_tracing_planner_hook(Query *query, const char *query_string, int cursorOption { nested_level--; span_end_time = GetCurrentTimestamp(); - handle_pg_error(traceparent, NULL, span_end_time); + if (current_trace_spans != NULL) + handle_pg_error(traceparent, NULL, span_end_time); PG_RE_THROW(); } PG_END_TRY(); @@ -1737,6 +1755,10 @@ pg_tracing_planner_hook(Query *query, const char *query_string, int cursorOption end_nested_level(span_end_time); nested_level--; + /* A function evaluated by the planner may have stopped tracing. */ + if (current_trace_spans == NULL) + return result; + /* End planner span */ pop_and_store_active_span(span_end_time); @@ -2155,15 +2177,22 @@ pg_tracing_ProcessUtility(PlannedStmt *pstmt, const char *queryString, { /* No sampling, just go through the standard process utility */ nested_level++; - if (prev_ProcessUtility) - prev_ProcessUtility(pstmt, queryString, readOnlyTree, - context, params, queryEnv, - dest, qc); - else - standard_ProcessUtility(pstmt, queryString, readOnlyTree, + PG_TRY(); + { + if (prev_ProcessUtility) + prev_ProcessUtility(pstmt, queryString, readOnlyTree, context, params, queryEnv, dest, qc); - nested_level--; + else + standard_ProcessUtility(pstmt, queryString, readOnlyTree, + context, params, queryEnv, + dest, qc); + } + PG_FINALLY(); + { + nested_level--; + } + PG_END_TRY(); return; } diff --git a/src/pg_tracing.h b/src/pg_tracing.h index 87e6453..4d95c27 100644 --- a/src/pg_tracing.h +++ b/src/pg_tracing.h @@ -194,9 +194,9 @@ typedef struct Span SpanType type; /* Type of the span. Used to generate the * span's name */ - int8 nested_level; /* Nested level of this span this span. + int nested_level; /* Nested level of this span this span. * Internal usage only */ - int8 parent_planstate_index; /* Index to the parent planstate of + int parent_planstate_index; /* Index to the parent planstate of * this span. Internal usage only */ uint8 subxact_count; /* Active count of backend's subtransaction */ @@ -397,7 +397,7 @@ extern void begin_span(TraceId trace_id, Span * span, SpanType type, TimestampTz start_span); extern void end_span(Span * span, const TimestampTz *end_time); extern void reset_span(Span * span); -extern const char *get_operation_name(const Span * span); +extern const char *get_operation_name(const Span * span, const char *spans_str); extern bool traceid_zero(TraceId trace_id); extern bool traceid_equal(TraceId trace_id_1, TraceId trace_id_2); extern const char *span_type_to_str(SpanType span_type); @@ -429,6 +429,7 @@ extern char *pg_tracing_otel_service_name; extern char *pg_tracing_otel_endpoint; extern int pg_tracing_otel_naptime; extern int pg_tracing_otel_connect_timeout_ms; +extern int pg_tracing_otel_timeout_ms; extern void store_span(const Span * span); diff --git a/src/pg_tracing_active_spans.c b/src/pg_tracing_active_spans.c index f947452..0089652 100644 --- a/src/pg_tracing_active_spans.c +++ b/src/pg_tracing_active_spans.c @@ -132,7 +132,7 @@ begin_active_span(const SpanContext * span_context, Span * span, int query_len; const char *normalised_query; uint64 parent_id; - int8 parent_planstate_index = -1; + int parent_planstate_index = -1; if (nested_level == 0) /* Root active span, use the parent id from the trace context */ @@ -223,11 +223,11 @@ Span * push_child_active_span(MemoryContext context, const SpanContext * span_context, SpanType span_type) { - Span *parent_span = peek_active_span(); + uint64 parent_id = peek_active_span()->span_id; Span *span = allocate_new_active_span(context); begin_span(span_context->traceparent->trace_id, span, span_type, NULL, - parent_span->span_id, span_context->query_id, span_context->start_time); + parent_id, span_context->query_id, span_context->start_time); return span; } @@ -248,6 +248,14 @@ push_active_span(MemoryContext context, const SpanContext * span_context, SpanTy { Span *span = peek_active_span_for_current_level(); Span *parent_span = peek_active_span(); + Span parent_copy; + + /* Allocating a child can move the active span array. */ + if (parent_span != NULL) + { + parent_copy = *parent_span; + parent_span = &parent_copy; + } if (span == NULL) { diff --git a/src/pg_tracing_json.c b/src/pg_tracing_json.c index 006071e..9498a4b 100644 --- a/src/pg_tracing_json.c +++ b/src/pg_tracing_json.c @@ -302,7 +302,7 @@ append_span(const JsonContext * json_ctx, const Span * span) const char *operation_name; StringInfo str = json_ctx->str; - operation_name = get_operation_name(span); + operation_name = get_operation_name(span, json_ctx->spans_str); pg_snprintf(trace_id, 33, UINT64_HEX_PADDED_FORMAT UINT64_HEX_PADDED_FORMAT, span->trace_id.traceid_left, diff --git a/src/pg_tracing_operation_hash.c b/src/pg_tracing_operation_hash.c index 4c4a173..cc7447a 100644 --- a/src/pg_tracing_operation_hash.c +++ b/src/pg_tracing_operation_hash.c @@ -70,7 +70,7 @@ lookup_operation_name(const Span * span, const char *txt) { bool found; Size offset; - operationKey key; + operationKey key = {0}; operationEntry *entry; key.query_id = span->query_id; @@ -82,9 +82,7 @@ lookup_operation_name(const Span * span, const char *txt) * correctly propagate queryId. In those cases, we can't use the hash * and need to fallback to write text without hash tracking */ - offset = pg_tracing_shared_state->extent; - append_str_to_shared_str(txt, strlen(txt) + 1); - return offset; + return append_str_to_shared_str(txt, strlen(txt) + 1); } entry = (operationEntry *) hash_search(operation_name_hash, &key, HASH_ENTER, &found); @@ -92,8 +90,12 @@ lookup_operation_name(const Span * span, const char *txt) if (found) return entry->query_offset; - offset = pg_tracing_shared_state->extent; - append_str_to_shared_str(txt, strlen(txt) + 1); + offset = append_str_to_shared_str(txt, strlen(txt) + 1); + if (offset == (Size) -1) + { + hash_search(operation_name_hash, &key, HASH_REMOVE, NULL); + return offset; + } /* Update hash's entry */ entry->query_offset = offset; diff --git a/src/pg_tracing_otel.c b/src/pg_tracing_otel.c index 43bc9e8..3813f81 100644 --- a/src/pg_tracing_otel.c +++ b/src/pg_tracing_otel.c @@ -90,6 +90,7 @@ static CURLcode send_json_trace(OtelContext * octx, const char *json_span) { CURLcode res; + long http_status; if (octx->curl == NULL) { @@ -109,13 +110,20 @@ send_json_trace(OtelContext * octx, const char *json_span) if (octx->config_changed) { curl_easy_setopt(octx->curl, CURLOPT_URL, pg_tracing_otel_endpoint); - curl_easy_setopt(octx->curl, CURLOPT_CONNECTTIMEOUT_MS, pg_tracing_otel_connect_timeout_ms); + curl_easy_setopt(octx->curl, CURLOPT_CONNECTTIMEOUT_MS, (long) pg_tracing_otel_connect_timeout_ms); + curl_easy_setopt(octx->curl, CURLOPT_TIMEOUT_MS, (long) pg_tracing_otel_timeout_ms); octx->config_changed = false; } curl_easy_setopt(octx->curl, CURLOPT_POSTFIELDS, json_span); curl_easy_setopt(octx->curl, CURLOPT_POSTFIELDSIZE, (long) strlen(json_span)); res = curl_easy_perform(octx->curl); + if (res == CURLE_OK) + { + res = curl_easy_getinfo(octx->curl, CURLINFO_RESPONSE_CODE, &http_status); + if (res == CURLE_OK && (http_status < 200 || http_status >= 300)) + res = CURLE_HTTP_RETURNED_ERROR; + } return res; } diff --git a/src/pg_tracing_planstate.c b/src/pg_tracing_planstate.c index 4a78e18..3d3c49b 100644 --- a/src/pg_tracing_planstate.c +++ b/src/pg_tracing_planstate.c @@ -328,7 +328,7 @@ is_subplan_executed(PlanState *planstate) return splanstate->need_to_scan_locally; } - case T_GatherMerge: + case T_GatherMergeState: { GatherMergeState *splanstate = (GatherMergeState *) planstate; @@ -415,9 +415,16 @@ create_spans_from_planstate(PlanState *planstate, planstateTraceContext * planst * span_id */ traced_planstate = get_traced_planstate(planstate); - Assert(traced_planstate != NULL); - span_start = traced_planstate->node_start; - span_id = traced_planstate->span_id; + if (traced_planstate != NULL) + { + span_start = traced_planstate->node_start; + span_id = traced_planstate->span_id; + } + else + { + span_start = parent_start; + span_id = generate_rnd_uint64(); + } break; } Assert(span_start > 0); @@ -460,7 +467,10 @@ create_spans_from_planstate(PlanState *planstate, planstateTraceContext * planst TracedPlanstate *initplan_traced_planstate; TimestampTz initplan_span_end; uint64 init_plan_span_id; + TimestampTz initplan_span_start; + if (splan == NULL || splan->instrument == NULL) + continue; InstrEndLoop(splan->instrument); if (splan->instrument->total == 0) continue; @@ -471,15 +481,17 @@ create_spans_from_planstate(PlanState *planstate, planstateTraceContext * planst * to get the span's end. */ initplan_traced_planstate = get_traced_planstate(splan); + initplan_span_start = initplan_traced_planstate != NULL ? + initplan_traced_planstate->node_start : span_start; init_plan_span_id = generate_rnd_uint64(); - initplan_span_end = get_span_end_from_planstate(splan, initplan_traced_planstate->node_start, root_end); + initplan_span_end = get_span_end_from_planstate(splan, initplan_span_start, root_end); initplan_span = create_span_node(splan, planstateTraceContext, &init_plan_span_id, span_id, query_id, - SPAN_NODE_INIT_PLAN, sstate->subplan->plan_name, initplan_traced_planstate->node_start, initplan_span_end); + SPAN_NODE_INIT_PLAN, sstate->subplan->plan_name, initplan_span_start, initplan_span_end); store_span(&initplan_span); /* Use the initplan span as a parent */ - create_spans_from_planstate(splan, planstateTraceContext, initplan_span.span_id, query_id, initplan_traced_planstate->node_start, root_end, latest_end); + create_spans_from_planstate(splan, planstateTraceContext, initplan_span.span_id, query_id, initplan_span_start, root_end, latest_end); } /* Handle sub plans */ @@ -492,7 +504,10 @@ create_spans_from_planstate(PlanState *planstate, planstateTraceContext * planst TracedPlanstate *subplan_traced_planstate; TimestampTz subplan_span_end; uint64 subplan_span_id; + TimestampTz subplan_span_start; + if (splan == NULL || splan->instrument == NULL) + continue; InstrEndLoop(splan->instrument); if (splan->instrument->total == 0) continue; @@ -502,8 +517,10 @@ create_spans_from_planstate(PlanState *planstate, planstateTraceContext * planst * still use the tracedplan to get the end. */ subplan_traced_planstate = get_traced_planstate(splan); + subplan_span_start = subplan_traced_planstate != NULL ? + subplan_traced_planstate->node_start : span_start; subplan_span_id = generate_rnd_uint64(); - subplan_span_end = get_span_end_from_planstate(subplan_traced_planstate->planstate, subplan_traced_planstate->node_start, root_end); + subplan_span_end = get_span_end_from_planstate(splan, subplan_span_start, root_end); /* * Push subplan as an ancestor so that we can resolve referent of @@ -514,9 +531,9 @@ create_spans_from_planstate(PlanState *planstate, planstateTraceContext * planst subplan_span = create_span_node(splan, planstateTraceContext, &subplan_span_id, span_id, query_id, - SPAN_NODE_SUBPLAN, sstate->subplan->plan_name, subplan_traced_planstate->node_start, subplan_span_end); + SPAN_NODE_SUBPLAN, sstate->subplan->plan_name, subplan_span_start, subplan_span_end); store_span(&subplan_span); - child_end = create_spans_from_planstate(splan, planstateTraceContext, subplan_span.span_id, query_id, subplan_traced_planstate->node_start, + child_end = create_spans_from_planstate(splan, planstateTraceContext, subplan_span.span_id, query_id, subplan_span_start, root_end, latest_end); if (planstateTraceContext->deparse_ctx) diff --git a/src/pg_tracing_span.c b/src/pg_tracing_span.c index 40f5190..bae7f03 100644 --- a/src/pg_tracing_span.c +++ b/src/pg_tracing_span.c @@ -282,7 +282,7 @@ is_span_top(SpanType span_type) * For node span, the name may be pulled from the stat file. */ const char * -get_operation_name(const Span * span) +get_operation_name(const Span * span, const char *spans_str) { const char *span_type_str; const char *operation_str = NULL; @@ -297,7 +297,7 @@ get_operation_name(const Span * span) span_type_str = span_type_to_str(span->type); /* TODO: Check for maximum offset */ if (span->operation_name_offset != -1) - operation_str = shared_str + span->operation_name_offset; + operation_str = spans_str + span->operation_name_offset; else return span_type_str; diff --git a/src/pg_tracing_sql_functions.c b/src/pg_tracing_sql_functions.c index be2afea..af23b3b 100644 --- a/src/pg_tracing_sql_functions.c +++ b/src/pg_tracing_sql_functions.c @@ -190,7 +190,7 @@ add_result_span(ReturnSetInfo *rsinfo, Span * span) char span_id[17]; span_type = span_type_to_str(span->type); - operation_name = get_operation_name(span); + operation_name = get_operation_name(span, shared_str); sql_error_code = unpack_sql_state(span->sql_error_code); pg_snprintf(trace_id, 33, UINT64_HEX_PADDED_FORMAT UINT64_HEX_PADDED_FORMAT, diff --git a/t/004_reliability.pl b/t/004_reliability.pl new file mode 100644 index 0000000..99b27c0 --- /dev/null +++ b/t/004_reliability.pl @@ -0,0 +1,106 @@ +use strict; +use warnings; + +use PostgreSQL::Test::Cluster; +use PostgreSQL::Test::Utils; +use Test::More; + +my $node = PostgreSQL::Test::Cluster->new('reliability'); +$node->init; +$node->append_conf('postgresql.conf', qq{ +shared_preload_libraries = 'pg_tracing' +pg_tracing.max_span = 10000 +statement_timeout = '10s' +}); +$node->start; +$node->safe_psql('postgres', 'CREATE EXTENSION pg_tracing'); + +# An ordinary utility error must not change sampling of the next statement in +# the same connection. +my ($stdout, $stderr); +my $status = $node->psql('postgres', q{ +SET pg_tracing.sample_rate = 0; +DROP TABLE pg_tracing_missing_table; +/*traceparent='00-00000000000000000000000000000001-0000000000000001-01'*/ SELECT 42; +SELECT count(*) > 0 FROM pg_tracing_consume_spans; +}, stdout => \$stdout, stderr => \$stderr, on_error_stop => 0); +is($status, 0, 'connection remains usable after a utility error'); +like($stderr, qr/does not exist/, 'the original SQL error is preserved'); +like($stdout, qr/42\s+t\s*$/, 'subsequent sampled query still produces spans'); + +$node->safe_psql('postgres', q{ +CREATE FUNCTION tracing_nested(n integer) RETURNS integer LANGUAGE plpgsql AS $$ +BEGIN + IF n = 0 THEN RETURN 1; END IF; + RETURN tracing_nested(n - 1) + 1; +END +$$; +}); +is($node->safe_psql('postgres', q{ +SET pg_tracing.sample_rate = 1; +SELECT tracing_nested(20); +}), '21', 'nested queries preserve their result across span storage growth'); +is($node->safe_psql('postgres', q{ +WITH spans AS MATERIALIZED (SELECT * FROM pg_tracing_consume_spans) +SELECT count(*) = 0 FROM spans child +WHERE child.span_type = 'ExecutorRun' +AND NOT EXISTS (SELECT FROM spans parent WHERE parent.span_id = child.parent_id); +}), 't', 'nested executor spans keep their parent links'); + +# JSON output remains readable after nested execution. +$node->safe_psql('postgres', q{ +/*traceparent='00-00000000000000000000000000000001-0000000000000001-01'*/ SELECT 7; +}); +is($node->safe_psql('postgres', q{ +SELECT count(*) > 0 FROM json_array_elements( + pg_tracing_json_spans()::json -> 'resourceSpans' -> 0 -> 'scopeSpans'); +}), 't', 'JSON output remains readable'); + +$node->safe_psql('postgres', q{ +CREATE TABLE tracing_data AS SELECT g AS id FROM generate_series(1, 20000) g; +ANALYZE tracing_data; +}); +my $parallel_settings = q{ +SET max_parallel_workers_per_gather = 2; +SET min_parallel_table_scan_size = 0; +SET parallel_setup_cost = 0; +SET parallel_tuple_cost = 0; +SET parallel_leader_participation = off; +SET enable_nestloop = off; +SET enable_mergejoin = off; +}; +my $sorted_join = 'SELECT a.id FROM tracing_data a JOIN tracing_data b USING (id) ORDER BY a.id'; +like($node->safe_psql('postgres', $parallel_settings . 'EXPLAIN ' . $sorted_join), + qr/Gather Merge/, 'parallel test uses the sorted parallel plan'); +is($node->safe_psql('postgres', $parallel_settings . q{ +SET pg_tracing.sample_rate = 1; +} . "SELECT count(*) FROM ($sorted_join) s;"), + '20000', 'sorted parallel join returns all rows with tracing enabled'); +is($node->safe_psql('postgres', q{ +SELECT count(*) > 0 FROM pg_tracing_consume_spans WHERE span_type = 'GatherMerge'; +}), 't', 'sorted parallel plan produces a GatherMerge span'); +$node->stop; + +# Keep span capacity available while exercising the independent text capacity. +my $small = PostgreSQL::Test::Cluster->new('small_strings'); +$small->init; +$small->append_conf('postgresql.conf', qq{ +shared_preload_libraries = 'pg_tracing' +pg_tracing.shared_str_size = 4 +pg_tracing.max_span = 10000 +statement_timeout = '10s' +}); +$small->start; +$small->safe_psql('postgres', 'CREATE EXTENSION pg_tracing'); +$small->safe_psql('postgres', q{ +/*traceparent='00-00000000000000000000000000000001-0000000000000001-01'*/ SELECT 'ordinary parameter longer than the text buffer'; +}); +is($small->safe_psql('postgres', q{ +SELECT span_operation = 'Select query' AND parameters IS NULL +FROM pg_tracing_peek_spans WHERE span_type = 'Select query'; +}), 't', 'unavailable text is omitted while the statement span is retained'); +is($small->safe_psql('postgres', q{ +SELECT json_typeof(pg_tracing_json_spans()::json); +}), 'object', 'JSON remains valid when the text buffer is full'); +$small->stop; +done_testing(); diff --git a/t/005_export_retry.pl b/t/005_export_retry.pl new file mode 100644 index 0000000..c2a04f5 --- /dev/null +++ b/t/005_export_retry.pl @@ -0,0 +1,79 @@ +use strict; +use warnings; + +use PostgreSQL::Test::Cluster; +use PostgreSQL::Test::Utils; +use Test::More; + +{ + package ReliabilityCollector; + use HTTP::Server::Simple::CGI; + use base qw(HTTP::Server::Simple::CGI); + + sub handle_request { + my ($self, $cgi) = @_; + if ($cgi->path_info eq '/reject') { + print "HTTP/1.0 503 Service Unavailable\r\nContent-Length: 0\r\n\r\n"; + } else { + sleep 2 if $cgi->path_info eq '/slow'; + print "HTTP/1.0 200 OK\r\nContent-Length: 0\r\n\r\n"; + } + } +} + +my $webserver_pid = ReliabilityCollector->new(4319)->background(); +END { kill 'TERM', $webserver_pid if defined $webserver_pid; } + +my $node = PostgreSQL::Test::Cluster->new('export_retry'); +$node->init; +$node->append_conf('postgresql.conf', qq{ +shared_preload_libraries = 'pg_tracing' +pg_tracing.otel_naptime = 1000 +pg_tracing.otel_endpoint = 'http://127.0.0.1:4319/reject' +pg_tracing.otel_timeout_ms = 200 +}); +$node->start; +$node->safe_psql('postgres', 'CREATE EXTENSION pg_tracing'); +$node->safe_psql('postgres', q{ +/*traceparent='00-00000000000000000000000000000001-0000000000000001-01'*/ SELECT 1; +}); +ok($node->poll_query_until('postgres', + 'SELECT otel_failures >= 1 FROM pg_tracing_info'), + 'collector HTTP failure is recorded'); +is($node->safe_psql('postgres', 'SELECT otel_sent_spans FROM pg_tracing_info'), + '0', 'rejected payload is not counted as sent'); + +$node->append_conf('postgresql.conf', + "pg_tracing.otel_endpoint = 'http://127.0.0.1:4319/ok'"); +$node->reload; +ok($node->poll_query_until('postgres', + 'SELECT otel_sent_spans = 4 FROM pg_tracing_info'), + 'the original payload is retried after collector recovery'); + +$node->safe_psql('postgres', 'SELECT pg_tracing_reset()'); +$node->append_conf('postgresql.conf', + "pg_tracing.otel_endpoint = 'http://127.0.0.1:4319/slow'"); +$node->reload; +ok($node->poll_query_until('postgres', + q{SELECT current_setting('pg_tracing.otel_endpoint') LIKE '%/slow'}), + 'collector configuration reload is visible'); +$node->safe_psql('postgres', q{ +/*traceparent='00-00000000000000000000000000000001-0000000000000001-01'*/ SELECT 2; +}); +ok($node->poll_query_until('postgres', + 'SELECT otel_failures >= 1 FROM pg_tracing_info'), + 'a connected collector with a slow response respects the request timeout'); +is($node->safe_psql('postgres', 'SELECT otel_sent_spans FROM pg_tracing_info'), + '0', 'timed-out payload is retained for retry'); + +# Allow the normal response after reloading both endpoint and total timeout. +$node->append_conf('postgresql.conf', qq{ +pg_tracing.otel_endpoint = 'http://127.0.0.1:4319/ok' +pg_tracing.otel_timeout_ms = 10000 +}); +$node->reload; +ok($node->poll_query_until('postgres', + 'SELECT otel_sent_spans = 4 FROM pg_tracing_info'), + 'export recovers after a request timeout'); +$node->stop; +done_testing();