Skip to content

feat: emit parent_query_id to link nested SPI queries - #95

Open
JoshDreamland wants to merge 10 commits into
mainfrom
parent-query-id-surgical
Open

feat: emit parent_query_id to link nested SPI queries#95
JoshDreamland wants to merge 10 commits into
mainfrom
parent-query-id-surgical

Conversation

@JoshDreamland

@JoshDreamland JoshDreamland commented May 13, 2026

Copy link
Copy Markdown
Contributor

Assigns a parent_query_id to events to avoid polluting metrics for nested queries (when plpgsql functions issue SPI statements that emit events within events).

Changes

  • PschEvent::top_levelPschEvent::parent_query_id (same slot in event/ring_entry; static asserts still hold).
  • Three single statics (rusage_start, query_start_ts, current_query_is_top_level) consolidated into one fixed-size query_stack[16] of {queryid, rusage_start, query_start_ts} indexed by nesting_level. The existing PG_TRY/nesting_level pattern in ExecutorRun/Finish/ProcessUtility is unchanged.
  • 2-line RecordUInt64+Append for the new column in stats_exporter.cc.
  • parent_query_id UInt64 column in docker/init/00-schema.sql; migrations/001_add_parent_query_id.sql for existing deployments.

Bug fix

The prior single static rusage_start was overwritten by nested getrusage() calls, so the outer query's CPU delta only measured the post-SPI portion. Per-frame baselines fix that alongside parent linkage.

Test plan

  • Local build (cmake + ninja)
  • Regression: mise run test:regress
  • Isolation: mise run test:isolation
  • CI green across PG 16/17/18 × amd64/arm64 + Clang job

Each event now carries the queryid of its calling query (0 for top-level).
Lets aggregations filter to top-level work with WHERE parent_query_id = 0,
avoiding the double-counting that happens when plpgsql functions issue
SPI statements that themselves emit events.

The Event payload's prior `bool top_level` is replaced in-place by
`uint64 parent_query_id`; static asserts in ring_entry.h continue to verify
layout equivalence with the on-wire ring slot.

Single statics for rusage_start, query_start_ts, and current_query_is_top_level
are consolidated into a fixed-size query_stack[16] indexed by nesting_level.
The previous single rusage_start was clobbered on nested SPI getrusage()
calls, so nested CPU deltas were measuring just the inner portion and the
outer's CPU was wrong; per-frame baselines fix that alongside parent linkage.
The PG_TRY/nesting_level pattern in ExecutorRun/Finish/ProcessUtility is
untouched.

Schema gains a parent_query_id UInt64 column; migrations/001_add_parent_query_id.sql
covers existing deployments. This supersedes #61 (which mixed the feature
with broader refactor churn — vector→array→array+counter iterations,
helper extraction, depth-cap retunes.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
EOF
)
Copilot AI review requested due to automatic review settings May 13, 2026 15:56

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

This PR adds parent_query_id to emitted events so nested SPI queries can be linked to their caller and excluded from “top-level only” aggregations, while also fixing nested CPU timing by switching from single baselines to a per-nesting frame stack.

Changes:

  • Replaced top_level with parent_query_id in the in-memory event/ring-entry payload.
  • Introduced a fixed-size query_stack to track per-nesting queryid, CPU baseline, and start timestamp.
  • Exported and persisted parent_query_id via ClickHouse schema changes plus a migration.

Reviewed changes

Copilot reviewed 6 out of 6 changed files in this pull request and generated 5 comments.

Show a summary per file
File Description
src/queue/ring_entry.h Replaces top_level with parent_query_id in the ring buffer entry while preserving the block-copy prefix layout checks.
src/queue/event.h Replaces top_level with parent_query_id in the public event struct.
src/hooks/hooks.c Implements per-nesting query frames and populates parent_query_id during executor/utility/log event emission.
src/export/stats_exporter.cc Adds a parent_query_id column to exported event rows.
docker/init/00-schema.sql Adds parent_query_id column to the canonical ClickHouse schema.
migrations/001_add_parent_query_id.sql Adds a migration for adding parent_query_id to existing ClickHouse tables.

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment thread docker/init/00-schema.sql
Comment thread migrations/001_add_parent_query_id.sql Outdated
Comment on lines +11 to +12
ALTER TABLE pg_stat_ch.events_raw
ADD COLUMN IF NOT EXISTS parent_query_id UInt64 DEFAULT 0;
Comment thread src/hooks/hooks.c Outdated
Comment thread src/export/stats_exporter.cc
Comment thread src/hooks/hooks.c
Comment on lines 888 to 892
static void CaptureLogEvent(ErrorData* edata) {
PschEvent event;
InitBaseEvent(&event, GetCurrentTimestamp(), (nesting_level == 0), PSCH_CMD_UNKNOWN);
InitBaseEvent(&event, GetCurrentTimestamp(), GetParentQueryId(), PSCH_CMD_UNKNOWN);

UnpackSqlState(edata->sqlerrcode, event.err_sqlstate);
JoshDreamland and others added 2 commits May 13, 2026 12:47
clang-format wants the column-aligned trailing comments above
parent_query_id to share its (wider) alignment column.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
Copilot review on #95:

1. parent_query_id was UInt64 while query_id is Int64; comparisons/joins
   between the two would require explicit casts in ClickHouse.  Switch the
   ClickHouse column, migration, and exporter to Int64.  Keep the in-memory
   struct field as uint64 (matches PG's queryId type) and cast at append
   time the same way query_id already does.

2. At nesting_level >= PSCH_MAX_NESTING_DEPTH, the PschExecutorEnd frame
   pointer is NULL, leaving start_ts at 0 (the PG epoch).  Any path that
   computed GetCurrentTimestamp() - start_ts would yield ~25 years of us.
   Fall back to GetCurrentTimestamp() so deltas are ~0 instead.  In practice
   query_desc->totaltime supplies a real duration via instrumentation, so
   the fallback subtraction is only used when totaltime wasn't allocated.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
Copilot AI review requested due to automatic review settings May 13, 2026 16:54

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot encountered an error and was unable to review this pull request. You can try again by re-requesting a review.

JoshDreamland and others added 5 commits May 13, 2026 15:58
028_parent_query_id.pl exercises the three semantic invariants by
inspecting actual exported rows in ClickHouse:

  1. Top-level queries report parent_query_id = 0.
  2. Nested SPI queries report parent_query_id = outer caller's query_id
     (verified by self-join on the events_raw table).
  3. A log/error event captured during nested execution reports
     parent_query_id = outer caller (not the running query itself) and a
     non-zero query_id (the running statement).  This catches the
     CaptureLogEvent off-by-one: nesting_level is bumped inside
     ExecutorRun, so slot nesting_level - 1 holds the running query and
     its caller lives at nesting_level - 2.

Filters key off distinctive table/function names rather than constants
or comments — those survive query normalization, where literals would
be replaced with $N placeholders.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
The CaptureLogEvent off-by-one came from nesting_level moving in places
where the frame stack didn't — Run/Finish bumped/decremented it around
the body, while the frame slot was written in Start and only read in
End.  Inside the body (where CaptureLogEvent fires) nesting_level was
one higher than it was at End, so the same "slot[nesting_level - 1]"
expression meant different things in each place.

Bind the two together: nesting_level moves exactly where a frame is
pushed or popped, and parent_query_id is captured into the frame at
push time.  Now slot[nesting_level - 1] is the currently-active query
at every emit site — End, ProcessUtility, CaptureLogEvent — and parent
is read directly from frame->parent_query_id without offset arithmetic.

Mechanics:
  - PushQueryFrame in PschExecutorStart bumps nesting_level after
    writing the slot.  PopQueryFrame in PschExecutorEnd decrements.
  - PschExecutorRun and PschExecutorFinish no longer touch nesting_level
    on the success path; each just wraps its chain in PG_TRY/PG_CATCH
    so a longjmp out of the body pops the frame on the unwind and
    keeps depth balanced.  Each PG_CATCH on the unwind path decrements
    once per level — three deep nesting with an error at the bottom
    cascades through three PG_CATCH blocks for three pops, exactly
    matching three pushes.  No subxact callback needed.
  - PschProcessUtility brackets its body with PG_TRY/PG_FINALLY in the
    same function, so push at the top and PopQueryFrame in PG_FINALLY
    balance regardless of how the body exits.  ExecuteUtilityWithNesting
    is gone.
  - CaptureLogEvent now reads TopQueryFrame() and uses both its queryid
    (the running query, attached as event.queryid for attribution) and
    its parent_query_id (the running query's caller, attached as
    event.parent_query_id with strict semantics — previously this slot
    was misread as the running query's own id).

PeekQueryStack and its offset parameter are gone.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
The totaltime allocation runs after standard_ExecutorStart returns,
between our PG_END_TRY and the close of PschExecutorStart.  InstrAlloc
goes through MemoryContextAllocZero, which can ereport(ERROR) on OOM —
that longjmp would unwind past our PG_END_TRY without firing the
PG_CATCH, leaving the pushed frame unpopped.

Move the allocation inside the existing PG_TRY block so the same
PG_CATCH cleans up the frame.  No behavior change on the happy path;
new robustness against any future code added between the push and the
function's return.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
The push had been doing more conditional work than necessary to figure
out whether a parent existed.  Convert nesting_level so that:

  - -1 is the resting state (no active query),
  - n in [0, PSCH_MAX_NESTING_DEPTH-1] is the slot index of the
    currently-active frame on top of the stack,
  - >= PSCH_MAX_NESTING_DEPTH is the overflow region.

Push now increments first and bails on overflow, so the cap check
collapses to a single comparison.  The parent_query_id lookup still
needs one conditional ("do we have a parent at all"), but the previous
double-condition (positive AND within cap) is gone.  TopQueryFrame /
PopQueryFrame benefit too: < 0 doubles as a nullity check.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
Two cleanups around the runaway-nesting case:

  - PushQueryFrame now emits a WARNING when it returns NULL because the
    cap was exceeded.  Previously the overflow was silent — telemetry
    just went missing for the frame.  WARNING is loud enough to notice
    if it happens in practice but doesn't fail the query.

  - PschProcessUtility's start_ts fallback was still 0 (PG epoch) when
    frame was NULL.  Match PschExecutorEnd by falling back to
    GetCurrentTimestamp() so the emitted event carries an approximately
    correct ts_start rather than 1970.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
@JoshDreamland
JoshDreamland force-pushed the parent-query-id-surgical branch from 3232e8e to c632cdc Compare May 14, 2026 14:42
Copilot AI review requested due to automatic review settings May 14, 2026 14:42

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot encountered an error and was unable to review this pull request. You can try again by re-requesting a review.

Earlier commits on this branch added parent_query_id to PschEvent and to
the ClickHouse-native exporter, plus the ClickHouse schema migration, but
arrow_batch.cc keeps its own hard-coded column list and was missed.  The
result: the production export path (OTel + Arrow IPC) silently dropped
parent_query_id on every event, so downstream consumers saw 0 across the
board.

Caught via the OTel/Arrow quickstart-validate harness (see PR #96).
Local script flips from 4/8 passing to all-green once this lands.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
@JoshDreamland
JoshDreamland requested a review from serprex May 14, 2026 21:06
Copilot AI review requested due to automatic review settings May 15, 2026 15:06
@JoshDreamland
JoshDreamland force-pushed the parent-query-id-surgical branch from 31eadfe to 18a6eae Compare May 15, 2026 15:06

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot encountered an error and was unable to review this pull request. You can try again by re-requesting a review.

@JoshDreamland
JoshDreamland force-pushed the parent-query-id-surgical branch from 18a6eae to 85266f0 Compare May 15, 2026 15:22
t/028 already exercises parent_query_id linkage through the
ClickHouse-native exporter.  Add t/029 to exercise the same
invariants through the Arrow/OTel export path (the production
pathway), using pg_stat_ch.debug_arrow_dump_dir to capture each
Arrow IPC batch to disk before the gRPC send (which we point at a
non-existent collector so it fails harmlessly).  Same trick as
t/026_arrow_dump.

t/029 guards specifically against the earlier surgical-PR regression
where arrow_batch.cc was missing parent_query_id entirely (the
ClickHouse-native exporter had it, so t/028 alone would have let
the gap through).

Test 3 in both files uses RAISE WARNING inside a nested SPI call to
exercise CaptureLogEvent's queryid/parent_query_id assignment.  We
cannot use a caught ERROR here: errfinish PG_RE_THROWs at elog.c:539
for ERROR-level events without calling EmitErrorReport, so
emit_log_hook only fires later from PostgresMain's top-level catch
(after all frames have been popped via PG_CATCH unwinding) — or
never at all, for errors caught by a plpgsql EXCEPTION block.
WARNING-level events go through EmitErrorReport directly, so our
hook fires while the inner SPI's frame is still on the stack, which
is the only scenario where CaptureLogEvent's slot choice is
observable.

The assertions distinguish "inner SPI queryid" from "outer caller
queryid" so a regression that attributes the log event to the wrong
slot (the off-by-one we're guarding against) would be caught.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
Copilot AI review requested due to automatic review settings May 15, 2026 18:12
@JoshDreamland
JoshDreamland force-pushed the parent-query-id-surgical branch from 85266f0 to f52d76c Compare May 15, 2026 18:12

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot encountered an error and was unable to review this pull request. You can try again by re-requesting a review.

@JoshDreamland

Copy link
Copy Markdown
Contributor Author

I missed stuff, because Arrow is a different code path with different implications. I think that's addressed but I'm feeling less sure of myself.

JoshDreamland added a commit that referenced this pull request Jun 8, 2026
In-place patch of docker/init/00-schema.sql, the CH-native exporter, and
the TAP tests so the docker quickstart schema aligns with what prod
actually writes to (datagres_otel.query_logs_arrow in clickgres-platform).
This is the pre-cutover unification: pg_stat_ch's CH-native path was
previously isolated from prod, and the two schemas had drifted apart on
both column naming and types.

Column renames (prod-side naming wins; closer to OTel semantic
conventions and minimizes downstream churn):
  ts_start    -> ts
  db          -> db_name
  username    -> db_user
  cmd_type    -> db_operation
  query       -> query_text

Type fix:
  err_sqlstate FixedString(5) -> LowCardinality(String)
    FixedString does not round-trip through Arrow IPC cleanly, and ~270
    SQLSTATE codes are dictionary-friendly. The CH-native exporter is
    updated to write the column via TagString (clickhouse-cpp's
    ColumnString -> CH LowCardinality(String) is fine on the wire).

Envelope columns added (with DEFAULT '' so the CH-native exporter, which
does not yet emit these, continues to insert successfully):
  instance_ubid, server_ubid, server_role, region, cell,
  service_version, host_id, pod_name

Engine/partitioning aligned with prod:
  ORDER BY ts -> ORDER BY (instance_ubid, ts)   (tenant locality)
  TTL added: toDate(ts) + INTERVAL 180 DAY
  SETTINGS index_granularity = 8192, ttl_only_drop_parts = 1

Materialized views (events_recent_1h, query_stats_5m, db_app_user_1m,
errors_recent) updated to reference the new column names and to include
instance_ubid in their ORDER BY / GROUP BY / SELECT projections so they
remain consistent with the events_raw partitioning strategy.

Test fixtures updated to query the new column names:
  t/010_clickhouse_export.pl, t/012_timing_accuracy.pl,
  t/021_cmd_type_counts.pl, t/027_query_normalization.pl,
  t/031_normalize_cache.pl

parent_query_id is intentionally NOT included here — it's the subject of
PR #95 (parent-query-id-surgical) and lands as its own follow-up
migration after this PR.

Validated end-to-end: docker/init/00-schema.sql applies cleanly on
clickhouse/clickhouse-server:26.1 (the version pinned in
docker/docker-compose.test.yml); INSERTs that omit the envelope columns
fill them via DEFAULT ''; all 4 MVs build. CI will run the TAP suite.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
JoshDreamland added a commit that referenced this pull request Jun 8, 2026
In-place patch of docker/init/00-schema.sql, the CH-native exporter, and
the TAP tests so the docker quickstart schema aligns with what prod
actually writes to (datagres_otel.query_logs_arrow in clickgres-platform).
This is the pre-cutover unification: pg_stat_ch's CH-native path was
previously isolated from prod, and the two schemas had drifted apart on
both column naming and types.

Column renames (prod-side naming wins; closer to OTel semantic
conventions and minimizes downstream churn):
  ts_start    -> ts
  db          -> db_name
  username    -> db_user
  cmd_type    -> db_operation
  query       -> query_text

Type fix:
  err_sqlstate FixedString(5) -> LowCardinality(String)
    FixedString does not round-trip through Arrow IPC cleanly, and ~270
    SQLSTATE codes are dictionary-friendly. The CH-native exporter is
    updated to write the column via TagString (clickhouse-cpp's
    ColumnString -> CH LowCardinality(String) is fine on the wire).

Envelope columns added (with DEFAULT '' so the CH-native exporter, which
does not yet emit these, continues to insert successfully):
  instance_ubid, server_ubid, server_role, region, cell,
  service_version, host_id, pod_name

Engine/partitioning aligned with prod:
  ORDER BY ts -> ORDER BY (instance_ubid, ts)   (tenant locality)
  TTL added: toDate(ts) + INTERVAL 180 DAY
  SETTINGS index_granularity = 8192, ttl_only_drop_parts = 1

Materialized views (events_recent_1h, query_stats_5m, db_app_user_1m,
errors_recent) updated to reference the new column names and to include
instance_ubid in their ORDER BY / GROUP BY / SELECT projections so they
remain consistent with the events_raw partitioning strategy.

Test fixtures updated to query the new column names:
  t/010_clickhouse_export.pl, t/012_timing_accuracy.pl,
  t/021_cmd_type_counts.pl, t/027_query_normalization.pl,
  t/031_normalize_cache.pl

parent_query_id is intentionally NOT included here — it's the subject of
PR #95 (parent-query-id-surgical) and lands as its own follow-up
migration after this PR.

Validated end-to-end: docker/init/00-schema.sql applies cleanly on
clickhouse/clickhouse-server:26.1 (the version pinned in
docker/docker-compose.test.yml); INSERTs that omit the envelope columns
fill them via DEFAULT ''; all 4 MVs build. CI will run the TAP suite.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
JoshDreamland added a commit that referenced this pull request Jun 18, 2026
…_raw

Closes the test gap that's existed since the Arrow path went live: no
existing test proves that pg_stat_ch's Arrow IPC output can actually be
ingested by ClickHouse against the unified events_raw schema. t/026
asserts on the IPC schema shape via pyarrow but never pushes the bytes
into CH; t/010 etc. exercise the CH-native Block path, not Arrow.

The new test wires the full producer-to-CH chain locally, bypassing
the OTel collector + receiver service entirely:

  1. Spin up a node with use_unified_arrow_exporter=on +
     debug_arrow_dump_dir set, an OTel endpoint that doesn't resolve so
     gRPC send fails — MaybeDumpArrowBatch fires BEFORE send so IPC
     files land on disk regardless.
  2. Run a deliberately-shaped workload (SELECT, CREATE, INSERT,
     SELECT count, DROP — five distinct statements).
  3. Force pg_stat_ch_flush(), wait for IPC files in $dump_dir.
  4. TRUNCATE pg_stat_ch.events_raw, then for each IPC file:
       curl -X POST --data-binary @$f \
         'http://localhost:18123/?query=INSERT INTO pg_stat_ch.events_raw FORMAT ArrowStream'
     A type mismatch on the wire (e.g. if the producer regressed to
     writing query_id as String) would surface here as a 4xx with a
     clear error rather than silently corrupting data.
  5. SELECT count() FROM events_raw, assert >= 5 rows.
  6. Pull system.columns and assert each id/counter column has the
     declared type from PR #99's schema (no silent string-typed regressions).
  7. Pinpoint the marker SELECT row and assert db_name/db_operation/
     query_text values match what we sent.
  8. Assert envelope columns (instance_ubid, server_role, region, cell,
     read_replica_type) carry the values from pg_stat_ch.extra_attributes.
  9. Assert parent_query_id is 0 across all rows (synthesized by the
     exporter until PR #95 lands).

Skips cleanly when Docker / the test CH container / the events_raw
schema aren't available — same patterns as t/010, t/013, t/021.

The "no OTel collector required" property makes this test purely a
producer⇄CH wire-format check. The clickgres-platform Go receiver is
not exercised here, since for verifying that the bytes match the
schema, a curl invocation is the simplest possible expression of "POST
this Arrow IPC body to CH" — the receiver's only added value over
curl in prod is OTel-collector pipeline integration, which we don't
care about for wire-format correctness.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
JoshDreamland added a commit that referenced this pull request Jun 18, 2026
…to-end

Introduce a new exporter that builds Arrow IPC RecordBatches through the
typed StatsExporter column-factory interface (StatLC/StatHC/StatTimestamp)
instead of the open-coded ArrowBatchBuilder used by arrow_batch.cc.
Composition over inheritance: the new exporter holds an OTelExporter for
gRPC transport (SendArrowBatch) but doesn't extend it, so the per-row
LogRecord state machine in OTelExporter — which is unused on this path
post-PR-#72 — stays out of scope.

Wire shape targets events_raw (the unified schema authored in PR #99),
not the legacy query_logs_arrow:
  * query_id, parent_query_id: Int64 (no sprintf decimal-string encoding)
  * pid: Int32
  * err_elevel: UInt8
  * buffer counters (shared/local/temp_blks_*, *_blk_*_time_us,
    wal_*, cpu_*_time_us): Int64
  * parallel_workers_planned/launched: Int16
  * jit_*: Int32
  * LC strings (db_*, err_sqlstate, app, server_role, region, cell,
    service_version, read_replica_type) -> DictionaryUtf8
  * HC strings (query_text, err_message, client_addr, instance_ubid,
    server_ubid, host_id, pod_name) -> plain utf8
  * ts: arrow::timestamp(MICRO, "UTC") matching DateTime64(6, 'UTC')

Column<T> wrappers are nested private types inside OTelArrowExporter
(not at namespace scope) so they can inherit from the protected
Column<T> base — same convention OTelExporter and ClickHouseExporter use
for their own column types.

Columns the caller doesn't explicitly populate are synthesized in
BeginRow by the exporter itself, so stats_exporter.cc's column-emission
loop stays unchanged:
  * parent_query_id (hardcoded 0 until PR #95 lands and PschEvent carries
    the field — events_raw requires the column on every insert, no DEFAULT)
  * 8 envelope columns from pg_stat_ch.extra_attributes (instance_ubid,
    server_ubid, server_role, region, cell, host_id, pod_name) plus
    read_replica_type (default 'none' if extra_attributes didn't supply)
  * service_version pinned to the compile-time PG_STAT_CH_VERSION macro

This commit only adds the exporter file (no dispatcher wiring yet) —
the next commit adds the GUC and routes batches through it when on.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
JoshDreamland added a commit that referenced this pull request Jun 18, 2026
…_raw

Closes the test gap that's existed since the Arrow path went live: no
existing test proves that pg_stat_ch's Arrow IPC output can actually be
ingested by ClickHouse against the unified events_raw schema. t/026
asserts on the IPC schema shape via pyarrow but never pushes the bytes
into CH; t/010 etc. exercise the CH-native Block path, not Arrow.

The new test wires the full producer-to-CH chain locally, bypassing
the OTel collector + receiver service entirely:

  1. Spin up a node with use_unified_arrow_exporter=on +
     debug_arrow_dump_dir set, an OTel endpoint that doesn't resolve so
     gRPC send fails — MaybeDumpArrowBatch fires BEFORE send so IPC
     files land on disk regardless.
  2. Run a deliberately-shaped workload (SELECT, CREATE, INSERT,
     SELECT count, DROP — five distinct statements).
  3. Force pg_stat_ch_flush(), wait for IPC files in $dump_dir.
  4. TRUNCATE pg_stat_ch.events_raw, then for each IPC file:
       curl -X POST --data-binary @$f \
         'http://localhost:18123/?query=INSERT INTO pg_stat_ch.events_raw FORMAT ArrowStream'
     A type mismatch on the wire (e.g. if the producer regressed to
     writing query_id as String) would surface here as a 4xx with a
     clear error rather than silently corrupting data.
  5. SELECT count() FROM events_raw, assert >= 5 rows.
  6. Pull system.columns and assert each id/counter column has the
     declared type from PR #99's schema (no silent string-typed regressions).
  7. Pinpoint the marker SELECT row and assert db_name/db_operation/
     query_text values match what we sent.
  8. Assert envelope columns (instance_ubid, server_role, region, cell,
     read_replica_type) carry the values from pg_stat_ch.extra_attributes.
  9. Assert parent_query_id is 0 across all rows (synthesized by the
     exporter until PR #95 lands).

Skips cleanly when Docker / the test CH container / the events_raw
schema aren't available — same patterns as t/010, t/013, t/021.

The "no OTel collector required" property makes this test purely a
producer⇄CH wire-format check. The clickgres-platform Go receiver is
not exercised here, since for verifying that the bytes match the
schema, a curl invocation is the simplest possible expression of "POST
this Arrow IPC body to CH" — the receiver's only added value over
curl in prod is OTel-collector pipeline integration, which we don't
care about for wire-format correctness.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
JoshDreamland added a commit that referenced this pull request Jun 18, 2026
…_raw

Closes the test gap that's existed since the Arrow path went live: no
existing test proves that pg_stat_ch's Arrow IPC output can actually be
ingested by ClickHouse against the unified events_raw schema. t/026
asserts on the IPC schema shape via pyarrow but never pushes the bytes
into CH; t/010 etc. exercise the CH-native Block path, not Arrow.

The new test wires the full producer-to-CH chain locally, bypassing
the OTel collector + receiver service entirely:

  1. Spin up a node with use_unified_arrow_exporter=on +
     debug_arrow_dump_dir set, an OTel endpoint that doesn't resolve so
     gRPC send fails — MaybeDumpArrowBatch fires BEFORE send so IPC
     files land on disk regardless.
  2. Run a deliberately-shaped workload (SELECT, CREATE, INSERT,
     SELECT count, DROP — five distinct statements).
  3. Force pg_stat_ch_flush(), wait for IPC files in $dump_dir.
  4. TRUNCATE pg_stat_ch.events_raw, then for each IPC file:
       curl -X POST --data-binary @$f \
         'http://localhost:18123/?query=INSERT INTO pg_stat_ch.events_raw FORMAT ArrowStream'
     A type mismatch on the wire (e.g. if the producer regressed to
     writing query_id as String) would surface here as a 4xx with a
     clear error rather than silently corrupting data.
  5. SELECT count() FROM events_raw, assert >= 5 rows.
  6. Pull system.columns and assert each id/counter column has the
     declared type from PR #99's schema (no silent string-typed regressions).
  7. Pinpoint the marker SELECT row and assert db_name/db_operation/
     query_text values match what we sent.
  8. Assert envelope columns (instance_ubid, server_role, region, cell,
     read_replica_type) carry the values from pg_stat_ch.extra_attributes.
  9. Assert parent_query_id is 0 across all rows (synthesized by the
     exporter until PR #95 lands).

Skips cleanly when Docker / the test CH container / the events_raw
schema aren't available — same patterns as t/010, t/013, t/021.

The "no OTel collector required" property makes this test purely a
producer⇄CH wire-format check. The clickgres-platform Go receiver is
not exercised here, since for verifying that the bytes match the
schema, a curl invocation is the simplest possible expression of "POST
this Arrow IPC body to CH" — the receiver's only added value over
curl in prod is OTel-collector pipeline integration, which we don't
care about for wire-format correctness.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
JoshDreamland added a commit that referenced this pull request Jun 18, 2026
…to-end

Introduce a new exporter that builds Arrow IPC RecordBatches through the
typed StatsExporter column-factory interface (StatLC/StatHC/StatTimestamp)
instead of the open-coded ArrowBatchBuilder used by arrow_batch.cc.
Composition over inheritance: the new exporter holds an OTelExporter for
gRPC transport (SendArrowBatch) but doesn't extend it, so the per-row
LogRecord state machine in OTelExporter — which is unused on this path
post-PR-#72 — stays out of scope.

Wire shape targets events_raw (the unified schema authored in PR #99),
not the legacy query_logs_arrow:
  * query_id, parent_query_id: Int64 (no sprintf decimal-string encoding)
  * pid: Int32
  * err_elevel: UInt8
  * buffer counters (shared/local/temp_blks_*, *_blk_*_time_us,
    wal_*, cpu_*_time_us): Int64
  * parallel_workers_planned/launched: Int16
  * jit_*: Int32
  * LC strings (db_*, err_sqlstate, app, server_role, region, cell,
    service_version, read_replica_type) -> DictionaryUtf8
  * HC strings (query_text, err_message, client_addr, instance_ubid,
    server_ubid, host_id, pod_name) -> plain utf8
  * ts: arrow::timestamp(MICRO, "UTC") matching DateTime64(6, 'UTC')

Column<T> wrappers are nested private types inside OTelArrowExporter
(not at namespace scope) so they can inherit from the protected
Column<T> base — same convention OTelExporter and ClickHouseExporter use
for their own column types.

Columns the caller doesn't explicitly populate are synthesized in
BeginRow by the exporter itself, so stats_exporter.cc's column-emission
loop stays unchanged:
  * parent_query_id (hardcoded 0 until PR #95 lands and PschEvent carries
    the field — events_raw requires the column on every insert, no DEFAULT)
  * 8 envelope columns from pg_stat_ch.extra_attributes (instance_ubid,
    server_ubid, server_role, region, cell, host_id, pod_name) plus
    read_replica_type (default 'none' if extra_attributes didn't supply)
  * service_version pinned to the compile-time PG_STAT_CH_VERSION macro

This commit only adds the exporter file (no dispatcher wiring yet) —
the next commit adds the GUC and routes batches through it when on.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
JoshDreamland added a commit that referenced this pull request Jun 18, 2026
…_raw

Closes the test gap that's existed since the Arrow path went live: no
existing test proves that pg_stat_ch's Arrow IPC output can actually be
ingested by ClickHouse against the unified events_raw schema. t/026
asserts on the IPC schema shape via pyarrow but never pushes the bytes
into CH; t/010 etc. exercise the CH-native Block path, not Arrow.

The new test wires the full producer-to-CH chain locally, bypassing
the OTel collector + receiver service entirely:

  1. Spin up a node with use_unified_arrow_exporter=on +
     debug_arrow_dump_dir set, an OTel endpoint that doesn't resolve so
     gRPC send fails — MaybeDumpArrowBatch fires BEFORE send so IPC
     files land on disk regardless.
  2. Run a deliberately-shaped workload (SELECT, CREATE, INSERT,
     SELECT count, DROP — five distinct statements).
  3. Force pg_stat_ch_flush(), wait for IPC files in $dump_dir.
  4. TRUNCATE pg_stat_ch.events_raw, then for each IPC file:
       curl -X POST --data-binary @$f \
         'http://localhost:18123/?query=INSERT INTO pg_stat_ch.events_raw FORMAT ArrowStream'
     A type mismatch on the wire (e.g. if the producer regressed to
     writing query_id as String) would surface here as a 4xx with a
     clear error rather than silently corrupting data.
  5. SELECT count() FROM events_raw, assert >= 5 rows.
  6. Pull system.columns and assert each id/counter column has the
     declared type from PR #99's schema (no silent string-typed regressions).
  7. Pinpoint the marker SELECT row and assert db_name/db_operation/
     query_text values match what we sent.
  8. Assert envelope columns (instance_ubid, server_role, region, cell,
     read_replica_type) carry the values from pg_stat_ch.extra_attributes.
  9. Assert parent_query_id is 0 across all rows (synthesized by the
     exporter until PR #95 lands).

Skips cleanly when Docker / the test CH container / the events_raw
schema aren't available — same patterns as t/010, t/013, t/021.

The "no OTel collector required" property makes this test purely a
producer⇄CH wire-format check. The clickgres-platform Go receiver is
not exercised here, since for verifying that the bytes match the
schema, a curl invocation is the simplest possible expression of "POST
this Arrow IPC body to CH" — the receiver's only added value over
curl in prod is OTel-collector pipeline integration, which we don't
care about for wire-format correctness.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
JoshDreamland added a commit that referenced this pull request Jul 20, 2026
…#116)

* feat(exporter): add OTelArrowExporter implementing StatsExporter end-to-end

Introduce a new exporter that builds Arrow IPC RecordBatches through the
typed StatsExporter column-factory interface (StatLC/StatHC/StatTimestamp)
instead of the open-coded ArrowBatchBuilder used by arrow_batch.cc.
Composition over inheritance: the new exporter holds an OTelExporter for
gRPC transport (SendArrowBatch) but doesn't extend it, so the per-row
LogRecord state machine in OTelExporter — which is unused on this path
post-PR-#72 — stays out of scope.

Wire shape targets events_raw (the unified schema authored in PR #99),
not the legacy query_logs_arrow:
  * query_id, parent_query_id: Int64 (no sprintf decimal-string encoding)
  * pid: Int32
  * err_elevel: UInt8
  * buffer counters (shared/local/temp_blks_*, *_blk_*_time_us,
    wal_*, cpu_*_time_us): Int64
  * parallel_workers_planned/launched: Int16
  * jit_*: Int32
  * LC strings (db_*, err_sqlstate, app, server_role, region, cell,
    service_version, read_replica_type) -> DictionaryUtf8
  * HC strings (query_text, err_message, client_addr, instance_ubid,
    server_ubid, host_id, pod_name) -> plain utf8
  * ts: arrow::timestamp(MICRO, "UTC") matching DateTime64(6, 'UTC')

Column<T> wrappers are nested private types inside OTelArrowExporter
(not at namespace scope) so they can inherit from the protected
Column<T> base — same convention OTelExporter and ClickHouseExporter use
for their own column types.

Columns the caller doesn't explicitly populate are synthesized in
BeginRow by the exporter itself, so stats_exporter.cc's column-emission
loop stays unchanged:
  * parent_query_id (hardcoded 0 until PR #95 lands and PschEvent carries
    the field — events_raw requires the column on every insert, no DEFAULT)
  * 8 envelope columns from pg_stat_ch.extra_attributes (instance_ubid,
    server_ubid, server_role, region, cell, host_id, pod_name) plus
    read_replica_type (default 'none' if extra_attributes didn't supply)
  * service_version pinned to the compile-time PG_STAT_CH_VERSION macro

This commit only adds the exporter file (no dispatcher wiring yet) —
the next commit adds the GUC and routes batches through it when on.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

* feat(exporter): add use_unified_arrow_exporter GUC + dispatcher wiring

New bool GUC pg_stat_ch.use_unified_arrow_exporter (default off, PGC_SIGHUP)
opts a producer instance into the OTelArrowExporter added in the previous
commit. When on (together with use_otel and otel_arrow_passthrough), the
bgworker constructs an OTelArrowExporter at init time and PschExportBatch
goes through ExportEventStats (the typed column interface) instead of
calling ExportEventsAsArrow (the legacy bypass that uses arrow_batch.cc
directly).

When off — the default — behavior is preserved bit-for-bit:
  * arrow_batch.cc / ExportEventsAsArrow remain reachable on the
    arrow_passthrough path
  * The OTelExporter column-emission path (off-arrow OTel logs) remains
    reachable when arrow_passthrough is off
  * The CH-native ClickHouseExporter remains reachable when use_otel is off

The sprintf-decimal-string ID encoding in arrow_batch.cc lives on for
the legacy path; the new exporter neither calls into it nor perpetuates
it. After the GUC has been on in prod long enough to retire the legacy
query_logs_arrow table, the arrow_batch.cc + ExportEventsAsArrow* + the
otel_arrow_passthrough GUC itself can be deleted in a single follow-on.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

* test: end-to-end round-trip of the unified Arrow exporter into events_raw

Closes the test gap that's existed since the Arrow path went live: no
existing test proves that pg_stat_ch's Arrow IPC output can actually be
ingested by ClickHouse against the unified events_raw schema. t/026
asserts on the IPC schema shape via pyarrow but never pushes the bytes
into CH; t/010 etc. exercise the CH-native Block path, not Arrow.

The new test wires the full producer-to-CH chain locally, bypassing
the OTel collector + receiver service entirely:

  1. Spin up a node with use_unified_arrow_exporter=on +
     debug_arrow_dump_dir set, an OTel endpoint that doesn't resolve so
     gRPC send fails — MaybeDumpArrowBatch fires BEFORE send so IPC
     files land on disk regardless.
  2. Run a deliberately-shaped workload (SELECT, CREATE, INSERT,
     SELECT count, DROP — five distinct statements).
  3. Force pg_stat_ch_flush(), wait for IPC files in $dump_dir.
  4. TRUNCATE pg_stat_ch.events_raw, then for each IPC file:
       curl -X POST --data-binary @$f \
         'http://localhost:18123/?query=INSERT INTO pg_stat_ch.events_raw FORMAT ArrowStream'
     A type mismatch on the wire (e.g. if the producer regressed to
     writing query_id as String) would surface here as a 4xx with a
     clear error rather than silently corrupting data.
  5. SELECT count() FROM events_raw, assert >= 5 rows.
  6. Pull system.columns and assert each id/counter column has the
     declared type from PR #99's schema (no silent string-typed regressions).
  7. Pinpoint the marker SELECT row and assert db_name/db_operation/
     query_text values match what we sent.
  8. Assert envelope columns (instance_ubid, server_role, region, cell,
     read_replica_type) carry the values from pg_stat_ch.extra_attributes.
  9. Assert parent_query_id is 0 across all rows (synthesized by the
     exporter until PR #95 lands).

Skips cleanly when Docker / the test CH container / the events_raw
schema aren't available — same patterns as t/010, t/013, t/021.

The "no OTel collector required" property makes this test purely a
producer⇄CH wire-format check. The clickgres-platform Go receiver is
not exercised here, since for verifying that the bytes match the
schema, a curl invocation is the simplest possible expression of "POST
this Arrow IPC body to CH" — the receiver's only added value over
curl in prod is OTel-collector pipeline integration, which we don't
care about for wire-format correctness.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

* fix(exporter): address Cursor + Copilot review nits

Three real findings worth addressing, plus a TODO for one that's worth
flagging but bigger than this PR:

  src/export/otel_arrow_exporter.cc:
    * NumExported() now returns row_count_ instead of delegating to the
      inner OTelExporter. The inner exporter's exported_count_ accumulates
      across batches because we never call inner_->BeginBatch() (its
      per-record LogRecord state machine is unused on this path), and
      PschExportBatch fetch_add's NumExported() into shared stats once per
      successful commit — so the cumulative inner count would cause
      quadratic over-reporting. (Cursor Bugbot, medium.)
    * ExtraAttrs class comment said "last write wins on duplicate keys"
      but Get() does a linear scan from the front and returns the first
      match. Comment now reflects what the code actually does. Duplicate
      keys in pg_stat_ch.extra_attributes are a user error in practice;
      no behavior change. (Copilot.)
    * Added a TODO(memory-budget) at CommitBatch documenting the missing
      mid-batch flush against psch_otel_max_block_bytes. The legacy
      ExportEventsAsArrowInternal flushes when the builder exceeds 3 MiB;
      this exporter ships the whole batch in one IPC. Acceptable while
      the GUC defaults off, but the budget check has to land before the
      default flips. Plumbing it through requires either invalidating
      caller-held column shared_ptrs at the flush boundary or threading
      a per-row size hook through the StatsExporter interface. (Cursor
      Bugbot, high — disagreeing on severity for the shadow-rollout
      window but acknowledging the underlying gap.)
    * Comment note on MaybeDumpArrowBatch acknowledging the intentional
      duplication with stats_exporter.cc — both copies die when the
      legacy path retires, so extracting a shared helper now would be a
      header just to delete it later. (Copilot, push-back.)

  t/036_unified_arrow_e2e.pl:
    * curl now uses --fail-with-body so HTTP 4xx/5xx from CH surfaces as
      a non-zero exit at the assertion site. The downstream SELECT count()
      assertion would catch a real ingestion failure too, but the sharper
      error at the source is a strict improvement. (Copilot.)

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

* fix(guc): make use_unified_arrow_exporter PGC_POSTMASTER

The bgworker constructs the exporter implementation once at init time
based on this GUC, but PschExportBatch was re-reading it every cycle to
choose between the legacy ExportEventsAsArrow bypass and the new
unified ExportEventStats path. Under PGC_SIGHUP, the two reads can
disagree:

  - Init=on, runtime=off: dispatcher takes the legacy bypass and calls
    SendArrowBatch on an OTelArrowExporter, which doesn't override it,
    so every batch silently drops.
  - Init=off, runtime=on: dispatcher runs the unified path through a
    plain OTelExporter, which ships OTLP log records instead of Arrow
    IPC. events_raw never sees the data.

Flag is a producer-shape feature toggle, not an operational tunable;
flipping it should be a deliberate restart.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

* feat(exporter): mid-batch flush at row boundary against psch_otel_max_block_bytes

OTelArrowExporter accumulated rows into a single Arrow IPC up to
psch_batch_max (default 200000) with no byte ceiling, so a queue
backlog could produce a payload exceeding gRPC's 4 MiB wire cap or
otelcol's HTTP body cap. Closes the gap relative to
ExportEventsAsArrowInternal which flushes mid-batch when the builder's
estimated bytes cross psch_otel_max_block_bytes.

Column wrappers each take a pointer to bytes_estimate_ and bump it per
Append (sizeof for fixed-width, value bytes + 4 for var-length offsets).
BeginRow samples bytes_estimate_ at the row boundary; if it has
crossed max_block_bytes_ and at least one row is already accumulated,
Flush() runs, ships the chunk, and resets per-flush state. Arrow
ArrayBuilder::Finish leaves builders empty and reusable, so subsequent
chunks share the same slot vector + column shared_ptrs without any
re-registration.

NumExported now reports exported_in_batch_ (cumulative across all
chunks in the active batch) instead of the per-chunk row_count_,
otherwise mid-batch flushes would erase the dispatcher's view of work
done. A sticky batch_failed_ flag poisons CommitBatch after any
mid-batch Flush failure so partial batches never count as a success.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

* feat(exporter): emit arrow_events_raw block_format on the unified path

The central OTel collector's routingconnector matches on the
pg_stat_ch.block_format OTLP attribute to fan batches between the
legacy receiver (writing query_logs_arrow) and a new receiver
configured for events_raw. Producer-side migration is therefore a
single-value swap rather than an endpoint reconfiguration — keeps the
extension dumb about deployment topology and the routing decision in
the central collector where it belongs.

The change widens StatsExporter::SendArrowBatch with an explicit
string_view block_format argument (no default — Google Style forbids
default args on virtuals because they dispatch on static type, not
dynamic). Threaded through PopulateResource and the per-record
attribute setter in OTelExporter::SendArrowBatch. Two callers:

- ExportEventsAsArrowInternal (legacy ArrowBatchBuilder path) passes
  "arrow_ipc" — preserves the existing wire-format marker so the
  legacy receiver continues to be the routing target.
- OTelArrowExporter::Flush (unified path) passes "arrow_events_raw"
  — the distinct value the routingconnector will key on to dispatch
  to the events_raw receiver. Coordinated with collector values in
  clickgres-platform/services/datagres-otelcol (separate PR).

No behavior change for the legacy path. No-op for any deployment
whose collector doesn't have the routingconnector configured (the
attribute is just additional metadata the receiver can ignore).

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

* fix(exporter): drop dangling Crunch() overrides on Arrow column wrappers

Serp's interface refactor moved Crunch() and Clear() out of the
StatsExporter::Column<T> base into a CH-private ChColumn interface (only
ClickHouseExporter's columns implement it, via multiple inheritance with
Column<T>). OTelArrowExporter's nested column wrappers still declared
'void Crunch() final {}' as no-op overrides — with the virtual gone from
the base, 'final' now applies to a nothing, which is a compile error
(only virtual member functions can be marked 'final'). CI turned red
across the build matrix on the merge commit for this reason.

Fix: delete the five no-op Crunch() overrides from ArrowNumColumn,
ArrowDictStrColumn, ArrowDictSvColumn, ArrowUtf8SvColumn, and
ArrowTimestampColumn. Builder finish already lives centrally in
OTelArrowExporter::Flush() (invoked mid-batch on the block-budget
crossing and once at CommitBatch for the residual chunk); the column
wrappers only own Append, so having them empty-implement a nonexistent
virtual was never load-bearing.

The rest of the refactor is sound: ClickHouseExporter's FixedCol /
StringCol correctly multi-inherit from Column<T> + ChColumn; columns_
stores shared_ptr<ChColumn> upcasts safely; both bases have virtual
destructors so lifetime is correct regardless of which pointer holds
the last reference.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants