Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/run-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -66,5 +66,5 @@ jobs:
run: make run-pgindent-diff

- name: Run Test
timeout-minutes: 1
timeout-minutes: 5
run: make run-test
4 changes: 2 additions & 2 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand Down
6 changes: 5 additions & 1 deletion doc/pg_tracing.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
9 changes: 9 additions & 0 deletions expected/parallel.out
Original file line number Diff line number Diff line change
Expand Up @@ -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?
Expand Down
9 changes: 9 additions & 0 deletions regress/16/expected/parallel.out
Original file line number Diff line number Diff line change
Expand Up @@ -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?
Expand Down
9 changes: 9 additions & 0 deletions regress/17/expected/parallel.out
Original file line number Diff line number Diff line change
Expand Up @@ -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?
Expand Down
9 changes: 9 additions & 0 deletions regress/18/expected/parallel.out
Original file line number Diff line number Diff line change
Expand Up @@ -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?
Expand Down
10 changes: 10 additions & 0 deletions sql/parallel.sql
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down
55 changes: 42 additions & 13 deletions src/pg_tracing.c
Original file line number Diff line number Diff line change
Expand Up @@ -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 */

Expand Down Expand Up @@ -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.",
Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -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) +
Expand All @@ -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));
Expand Down Expand Up @@ -1729,14 +1746,19 @@ 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();
span_end_time = GetCurrentTimestamp();
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);

Expand Down Expand Up @@ -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;
}

Expand Down
7 changes: 4 additions & 3 deletions src/pg_tracing.h
Original file line number Diff line number Diff line change
Expand Up @@ -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 */

Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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);
Expand Down
14 changes: 11 additions & 3 deletions src/pg_tracing_active_spans.c
Original file line number Diff line number Diff line change
Expand Up @@ -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 */
Expand Down Expand Up @@ -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;
}

Expand All @@ -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)
{
Expand Down
2 changes: 1 addition & 1 deletion src/pg_tracing_json.c
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
14 changes: 8 additions & 6 deletions src/pg_tracing_operation_hash.c
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -82,18 +82,20 @@ 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);

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;

Expand Down
10 changes: 9 additions & 1 deletion src/pg_tracing_otel.c
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,7 @@ static CURLcode
send_json_trace(OtelContext * octx, const char *json_span)
{
CURLcode res;
long http_status;

if (octx->curl == NULL)
{
Expand All @@ -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;
}

Expand Down
Loading
Loading