Skip to content

feat(stream)!: positional rows on NATS and SSE, columns sent apart - #554

Merged
EricAndrechek merged 115 commits into
mainfrom
stack/5-positional-wire
Sep 8, 2026
Merged

EricAndrechek merged 115 commits into
mainfrom
stack/5-positional-wire

Conversation

@EricAndrechek

@EricAndrechek EricAndrechek commented Sep 2, 2026 •

Copy link
Copy Markdown
Member

Summary

NATS envelope v2 (BREAKING). EventMessage's data object is replaced by
format ("JSONCompactEachRow"), columns (the table's declaration order) and
row (one compact line — a positional JSON array). The INSERT the worker emits carries the
column names once for a whole group instead of every row repeating every key
(each NATS envelope still carries its own columns, since a message must stand
alone), and a reader can tell a schema change mid-stream from a reordering.

In-flight messages published by an older version are not readable by the new
worker — an envelope whose format is absent or unknown, or whose columns
and row cannot be paired, carries no way to say which value belongs to which
column. DRAIN THE INGEST QUEUE BEFORE DEPLOYING. Such an envelope is parked on
the DLQ, or acked-and-dropped with an ERROR log where the DLQ is off, and never
inserted. Either disposition increments
wavehouse_ingest_poison_total{table,reason,disposition}, counted once the park
or the ack actually lands: a DLQ publish that fails, or an ack that fails, both
deliberately leave the message unacked and redelivering rather than losing the
row, so those are the two cases where it is not yet accounted for and not yet
gone.

The worker groups a batch by column list, so a schema change mid-stream splits
the INSERT rather than corrupting it, and writes
INSERT INTO {table} (cols) FORMAT JSONCompactEachRow. Row cells are copied as
their original bytes rather than re-encoded at each hop, so a 64-bit id past
2^53 keeps every digit end to end.

Computed columns (BREAKING). Naming columns explicitly in the INSERT makes a
computed column fatal — discovery reads every row of system.columns with no
default_kind filter, so a MATERIALIZED or ALIAS column landed in the
envelope and then the statement. Verified on the pinned ClickHouse 26.6.3.62:
MATERIALIZED is code 44, ALIAS is code 16. A table carrying either could
ingest under the previous column-less FORMAT JSONEachRow and could not ingest
at all after this change — every row to the DLQ, or redelivered forever where
the DLQ is off.

That classification (DefaultKind, IsInsertable, InsertableColumns,
InsertableColumnNames) originally sat in the NEXT PR of the stack, which
would have left main broken for any such table between the two merges. It is
here instead, with the guard that goes with it: a record supplying a value for
a computed column is a 400 naming the column, not a silent drop behind a
200. GET /v1/ops/schema still reports the whole table, now including
default_kind — a computed column stays queryable, it just cannot be written.
EPHEMERAL remains insertable; it is insert-only by construction.

No fixture in the suite declared a computed column, which is why every gate was
green while this was broken. tests/integration now creates one and drives
HTTP ingest → NATS → the worker's INSERT end to end.

SSE (BREAKING for raw consumers). A data frame's data object is replaced by
row, and each connection is sent an event: schema frame naming its
projected columns before its first row and again on drift. The TS SDK consumes
it and still yields row objects, so .stream(), .liveQuery() and
StreamEvent.data are unchanged; a raw EventSource consumer must now zip
rows itself.

Upgrade the SDK and the server together — they share this wire protocol, and a
skewed pair delivers no usable rows and raises no error.

A row is never queued without its announcement: if the queue is full when the
announcement is offered, the row is dropped too and the signature is not
recorded, so the next event announces again. A slow consumer loses a row rather
than receiving one it would zip against a stale column list. Both drops are
counted under their own frame kinds — the withheld row explicitly, since it is
never offered to Send.

Replay and the live fan-out track drift in separate state and do not reconcile
it; the residual same-length case is #543.

Stacked PR

This is part 5 of 7 in a stack that replaces #540. Each PR is based on the one above it, so review this PR's own diff against its base — GitHub shows only this layer's changes.

# Branch Base
0 stack/0-classify-paths main
1 stack/1-discovery stack/0-classify-paths
2 stack/2-policy stack/1-discovery
3 stack/3-content-type stack/2-policy
4 stack/4-seams stack/3-content-type
5 stack/5-positional-wire stack/4-seams → this PR
6 stack/6-computed-columns stack/5-positional-wire

Merge in order, top to bottom. Rebasing or squashing out of order will make the later PRs' diffs unreadable.

Test plan

  • make ci green on this branch's exact tree (verify, unit, integration against live ClickHouse, e2e, all coverage gates)
  • The branch descends from its base and carries only this layer's change (plus any follow-up commits answering review)
  • go.mod / go.sum untouched; no new dependencies

Review

Both pre-push reviewers ran against this branch directly, over several rounds:

  • pre-push-reviewer — ship_it. It independently verified the merge from refactor(ingest): buffer the body, add type-layer seams and encoder #553 lost nothing, re-ran the mutation proving the computed-column guard is load-bearing, and walked every .Columns consumer to confirm nothing else assumes a full-column envelope.
  • docs-reviewer — four rounds. It found the silent SSE replay hole across a v2 upgrade, an undocumented change to stored data on Nullable(T) DEFAULT columns, a missing drain runbook, and a schema-frame contract stated as a guarantee the wire does not provide. All applied. Its final sweep for "claims true before this branch and false after it" came back clean across every docs-prose file.

The markers are recorded via scripts/skip-pre-push-review.sh with the verdicts in tmp/review-skips-*.log, because .claude/hooks/review-marker.sh keys a marker to the session worktree's HEAD rather than the reviewed one, so a review run against another worktree cannot write its own marker. The reviews ran; only the bookkeeping needed the manual step.

🤖 Generated with Claude Code

https://claude.ai/code/session_018Epn88jTEw4ZkXrvTKzZXQ

EricAndrechek and others added 6 commits September 1, 2026 21:03
Both decisions in classify-paths.sh were `if printf … | grep -qE …; then A;
else B; fi`. grep exits 0 on match, 1 on no match, and 2 on error — can't
fork/exec, read error, bad pattern — and the else branch collapsed 1 and 2 into
the same answer. `set -euo pipefail` does not help: `set -e` is suppressed for a
command used as an `if` condition.

Observed twice while gating this stack, in runs whose static checks run at
-j 14. A different single case failed each time — `mixed-docs-go` answering
docs=false, then `dep-bump-go` answering code=false — while every other case
passed. That is the signature of a transient grep failure under load, not a
pattern bug; the script and its test are unchanged from main and both pass
standalone.

The test caught it because it asserts expected values. The production path has
no such check: CI's `changes` job gates the docs pipeline on this answer, so a
docs=false produced by an errored grep skips the docs build and still reports
success.

Both greps now go through a `matches` helper that aborts with a diagnostic on
any exit above 1. The test suite stubs grep onto PATH to prove the abort fires
— that case fails against the previous script, which answers confidently
instead.

Closes #545.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018Epn88jTEw4ZkXrvTKzZXQ
Additive metadata the upcoming native type layer needs, captured on the
same refresh as the columns so they can never describe different servers:

- Column gains DefaultExpression and Position, both scanned from the
  widened system.columns select. TableSchema.Columns was already ordered
  by position, so declaration order needed no new structure.
- TableSchema gains DDL from system.tables.create_table_query. It is
  json:"-": the schema endpoint marshals TableSchema straight to the
  client and an external-engine table (S3, MySQL, Kafka) carries its
  credentials in that statement. A table listed in system.tables with no
  system.columns rows is skipped, never published column-less.
- SchemaRegistry gains ServerVersion(), from a SELECT version() probe
  next to the existing SELECT timezone().

Both new queries fail the refresh on error, matching timezone() and
system.columns: callers keep the prior cache and retry.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018Epn88jTEw4ZkXrvTKzZXQ
`tables.<table>.select.<role>` becomes `tables.<table>.<role>.select`. Field
names and semantics are unchanged; only the nesting moves. `select` and
`insert` are now distinct types (`SelectPermissions`/`InsertPermissions`), so a
field on the wrong side is a validation error instead of being accepted and
ignored.

There is no automatic conversion — convert the file by hand and run
`wavehouse validate` before restarting. A pre-v2 document is reported as one
clear finding pointing at the migration note rather than a confusing
strict-decode error.

Three fail-closed fixes the new shape made visible or possible:

- `ResolvedPermissions` now marks the side `Evaluate` did not resolve. An
  unresolved side is zero, and a zero side reads as an empty allow list plus an
  empty deny list — which every accessor would answer as "unrestricted". Each
  one now denies instead, `HasRowFilter` included: it is the reachable gate in
  front of `RowVisible`, so without it the guard behind it was dead code and
  both call sites took the whole-bucket fast path.
- `evaluateInsert` resolved a check using an operator it does not honor
  (`_neq`/`_gt`/`_lt`, or the ambiguous `_eq`+`_in`) to no clause at all,
  authorizing the insert with the rule silently gone — where `evaluateSelect`
  denies outright in the mirror situation.
- A `filter` or `check` entry naming no operator (`"tenant_id": {}`) survived a
  strict decode, passed `Validate`, and matched no case in either resolver.
  Reproduced before fixing: `Allowed=true`, `HasRowFilter()=false`, and
  `RowVisible` true for another tenant's row — a declared row-level restriction
  applying on neither the query path nor the live stream.

Note on scope: an earlier revision of this work also made
`POST /v1/ops/policy/validate` agree with file adoption. #541 deleted the
policy HTTP surface entirely, so that fix is gone with it — the legacy-layout
detector it relied on is still reached through file adoption's own pipeline.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018Epn88jTEw4ZkXrvTKzZXQ
`POST /v1/ingest` used to sniff the body and treat the header as a hint: the
first non-whitespace byte chose between a single object and an array, and an
NDJSON body whose first line happened to start with `[` was re-framed as a JSON
array — silently, as a whole-request reinterpretation rather than a per-record
error.

The header is now required and authoritative. A request declaring nothing, or
something ingest does not read, is `415` before the body is parsed, with the
supported types named in the body. The declared type chooses the format
*family* — `application/json` versus the four NDJSON spellings — and within the
JSON family the first non-whitespace byte still picks array versus single
object. The bytes never choose the family, so an NDJSON body is read as NDJSON
whatever its first byte and a bad line fails as a per-record error.

Parameters are ignored (`application/json; charset=utf-8` is `application/json`),
matching `mime.ParseMediaType`.

The TS SDK already sent `application/json` on every ingest call, so
`.insert()` is unaffected. A hand-rolled client that relied on sniffing must
now declare the type.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018Epn88jTEw4ZkXrvTKzZXQ
…encoder

Observable behavior is unchanged. This is the seam work the native type layer
lands against, plus the memory trade that comes with it.

The handler now reads the whole (already `MaxBytesReader`-capped) body into a
pooled `*bytes.Buffer` and runs the record readers over those bytes rather than
the live connection, so the `413` surfaces at that read instead of
mid-iteration and the `415` is decided from the header before a byte is read.

The memory profile is NOT unchanged, and that is the deliberate part. Streaming
meant peak resident bytes on the order of one record — NDJSON scanned line by
line, the array path let `json.Decoder` compact after each element. Peak is now
O(body) per in-flight request, and `bytes.Buffer` doubles, so peak allocation
can exceed the cap before `MaxBytesReader` errors. `maxPooledBufferBytes` (1
MiB) caps what a request hands back to the pool, not its peak, and nothing
bounds total in-flight bytes — the ceiling is concurrency × the 16 MiB
data-plane cap, which has no operator knob, so the outer limit is the proxy's.
Kept because the type layer needs the body addressable rather than consumed; a
server-side bound is tracked in #544, and the reverse-proxy guide now says to
size the container for concurrency rather than for one request.

Three decision points became interfaces with default implementations that
delegate to today's code unchanged: `RecordValidator` (schema validation +
timestamp canonicalization — the two calls stay where they are, with the
check-clause block between them, since merging them would move checks onto
canonicalized values), `InsertChecker` (the `_eq` and `_in` comparisons), and
`stream.RowEvaluator` (row visibility, reached by both the live fan-out and
replay through the one shared admission step).

`ingest.EncodeCompactRow` lands here unused — the positional encoder the wire
change ahead of it will call.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018Epn88jTEw4ZkXrvTKzZXQ
NATS envelope v2 (BREAKING). `EventMessage`'s `data` object is replaced by
`format` ("JSONCompactEachRow"), `columns` (the table's declaration order) and
`row` (one compact line — a positional JSON array). A batch carries the column
names once instead of repeating every key on every row, and a reader can tell a
schema change mid-stream from a reordering.

In-flight messages published by an older version are not readable by the new
worker — an envelope whose `format` is absent or unknown, or whose `columns`
and `row` cannot be paired, carries no way to say which value belongs to which
column. DRAIN THE INGEST QUEUE BEFORE DEPLOYING. Such an envelope is parked on
the DLQ, or acked-and-dropped with an ERROR log and a
`wavehouse_ingest_poison_dropped_total{table,reason}` increment where the DLQ
is off — never left unacked to redeliver forever, and never inserted.

The worker groups a batch by column list, so a schema change mid-stream splits
the INSERT rather than corrupting it, and writes
`INSERT INTO {table} (cols) FORMAT JSONCompactEachRow`. Row cells are copied as
their original bytes rather than re-encoded at each hop, so a 64-bit id past
2^53 keeps every digit end to end.

SSE (BREAKING for raw consumers). A data frame's `data` object is replaced by
`row`, and each connection is sent an `event: schema` frame naming its
projected columns before its first row and again on drift. The TS SDK consumes
it and still yields row objects, so `.stream()`, `.liveQuery()` and
`StreamEvent.data` are unchanged; a raw `EventSource` consumer must now zip
rows itself.

Upgrade the SDK and the server together — they share this wire protocol, and a
skewed pair delivers no usable rows and raises no error.

A row is never queued without its announcement: if the queue is full when the
announcement is offered, the row is dropped too and the signature is not
recorded, so the next event announces again. A slow consumer loses a row rather
than receiving one it would zip against a stale column list. Both drops are
counted under their own frame kinds — the withheld row explicitly, since it is
never offered to Send.

Replay and the live fan-out track drift in separate state and do not reconcile
it; the residual same-length case is #543.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018Epn88jTEw4ZkXrvTKzZXQ
@coderabbitai

coderabbitai Bot commented Sep 2, 2026 •

Copy link
Copy Markdown

Review Change StackReview Change Stack

Note

Reviews paused

It looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the reviews.auto_review.auto_pause_after_reviewed_commits setting.

Use the following commands to manage reviews:

  • @coderabbitai resume to resume automatic reviews.
  • @coderabbitai review to trigger a single review.

Use the checkboxes below for quick actions:

  • ▶️ Resume reviews
  • 🔍 Trigger review

No actionable comments were generated in the recent review. 🎉

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Advanced

Run ID: f5f166a0-3839-4f4c-82ce-650ab2087e17

📥 Commits

Reviewing files that changed from the base of the PR and between 4c27c54 and 2b89ac8.

📒 Files selected for processing (4)
  • internal/ingest/worker.go
  • internal/ingest/worker_test.go
  • internal/stream/hub.go
  • internal/stream/hub_test.go

Included review availability: Your plan provides up to 1 included review per hour; 0 remain after this review.

📜 Recent review details
🧰 Additional context used
🧠 Learnings (2)
📓 Common learnings
Learnt from: EricAndrechek
URL: https://github.com/Wave-RF/WaveHouse/pull/554

Timestamp: 2026-09-05T00:43:57.044Z
Learning: In WaveHouse positional ingest, `internal/api/ingest.go` derives `EventMessage.Columns` from `discovery.TableSchema.InsertableColumnNames()`, which reflects unique `system.columns` names, and the embedded NATS server uses `DontListen: true`; therefore duplicate column names are unreachable from the current in-process producer. `internal/ingest/worker.go`, `internal/stream/hub.go`, and `clients/ts/src/stream/sse.ts` must still reject duplicate column names defensively because replayed DLQ messages or future external publishers can supply them, and `pairRow` otherwise overwrites an earlier cell in its name-keyed map before row-level policy evaluation.
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse

Timestamp: 2026-09-08T13:17:06.219Z
Learning: For WaveHouse PR stacks, document a protocol caveat in the PR that introduces the protocol behavior, not only in a later stacked PR. For the positional SSE protocol introduced by PR `#554`, documentation must state that gap-fill across a schema change can deliver live rows without a new schema frame; raw consumers must use the latest received schema, drop rows with mismatched arity, and reconnect because arity checks cannot detect same-length schema changes such as `RENAME COLUMN`.
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse PR: 554
File: docs/src/content/docs/api.md:613-613
Timestamp: 2026-09-05T00:26:39.936Z
Learning: In WaveHouse SSE streaming, `SubscribeSchemaFrame` is best-effort and is used for schema de-duplication and quiet-table connections. The live-delivery guarantee comes from `internal/stream/hub.go` `deliver` checking `Subscriber.needsSchema` before sending a row. `ReplayProjector` tracks schema state separately with closure-local `lastSig`; replay and live schema state do not reconcile across a schema change (`#543`). Raw SSE consumers must check row arity and reconnect to resynchronize a same-length schema change.
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse

Timestamp: 2026-09-08T12:06:54.852Z
Learning: In WaveHouse documentation that describes raw SSE positional rows, state that the stream sends an `event: schema` frame before rows and re-sends it when the column list changes. State that each `row` array matches the most recent schema. This prevents raw consumers from pairing later rows with a stale column list and silently mislabeling values.
Learnt from: CR
Repo: Wave-RF/WaveHouse

Timestamp: 2026-09-08T19:54:33.659Z
Learning: Address and resolve every review finding
Learnt from: CR
Repo: Wave-RF/WaveHouse

Timestamp: 2026-09-08T19:54:33.659Z
Learning: Run `make lint` and `make test` before considering work complete.
📚 Learning: 2026-06-26T12:23:22.696Z
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse PR: 346
File: internal/stream/subscriber_test.go:9-28
Timestamp: 2026-06-26T12:23:22.696Z
Learning: In this Go repository, prefer table-driven tests (e.g., `[]struct{...}` with `t.Run(...)`) only for tests that cover multiple scenarios/inputs and can be cleanly enumerated. Do not artificially rewrite a clear single-scenario sequential behavioral-flow test into a table-driven form just to fit the pattern; if there’s only one meaningful scenario, keep the test as a straightforward linear flow (as in `TestSubscriber_SendDeliversThenDropsWhenFull`).

Applied to files:

  • internal/stream/hub_test.go
🔇 Additional comments (4)
internal/ingest/worker.go (1)

720-727: LGTM!

Also applies to: 746-755

internal/ingest/worker_test.go (1)

1543-1555: LGTM!

internal/stream/hub.go (1)

358-362: LGTM!

Also applies to: 370-378

internal/stream/hub_test.go (1)

171-171: LGTM!

Also applies to: 185-191


📝 Summary

Summary by CodeRabbit

  • New Features

    • Streaming responses announce schemas separately and deliver positional rows.
    • Computed columns are excluded from inserts and calculated by the database.
    • Omitted nullable default columns are stored as NULL.
  • Bug Fixes

    • Unreadable ingest messages are routed to the DLQ or logged and counted when disabled.
    • Invalid policy checks return per-record 403 responses.
    • Stream rows preserve columns named __proto__.
  • Breaking Changes

    • Ingest and SSE wire formats require compatible server and SDK versions.
  • Documentation

    • Added migration, compatibility, streaming, and error-handling guidance.

Walkthrough

This change introduces positional v2 ingest envelopes and schema-aware SSE frames. It excludes computed columns from inserts, adds policy guards and poison-message handling, updates stream projection and replay, and updates the TypeScript SDK, tests, and documentation.

Changes

Positional ingest and schema handling

Layer / File(s) Summary
Insertable schema and v2 envelope
internal/discovery/*, internal/ingest/*, internal/api/ingest.go
Schemas preserve default_kind. Computed columns are excluded from inserts. Events use format, columns, and positional row fields.
Worker validation and batching
internal/ingest/worker.go, internal/ingest/worker_test.go
The worker validates row shape, groups messages by column list, inserts with JSONCompactEachRow, and routes unreadable messages to the DLQ or drop metric.
API and integration validation
internal/api/ingest_test.go, tests/integration/*
Tests cover positional publishing, computed columns, mixed column lists, policy check guards, and DLQ behavior.

Schema-aware streaming

Layer / File(s) Summary
Live and replay frame delivery
internal/stream/hub.go, internal/stream/subscriber.go, internal/api/stream.go
The hub sends projected schema frames before rows and re-announces schemas when column signatures change. Replay returns schema and data frame slices.
Projection and stream tests
internal/stream/*_test.go
Tests cover positional projection, row visibility, schema drift, replay, backpressure, fail-closed behavior, and computed-column exclusion.

TypeScript SDK

Layer / File(s) Summary
SSE decoding and row projection
clients/ts/src/stream/sse.ts, clients/ts/src/query-builder.ts, clients/ts/src/types.ts
The SDK validates schema frames, zips positional rows into null-prototype objects, bounds warnings, and preserves __proto__ columns.
SDK tests and documentation
clients/ts/src/stream/sse.test.ts, clients/ts/src/query-builder.test.ts, clients/ts/README.md
Tests cover schema changes, invalid rows, reconnects, hostile column names, and SDK/server version pairing.

Protocol documentation and build configuration

Layer / File(s) Summary
Protocol and upgrade documentation
docs/src/content/docs/*, CHANGELOG.md, AGENTS.md
Documentation describes the v2 ingest envelope, schema-frame SSE contract, computed-column behavior, DLQ handling, SDK pairing, and upgrade constraints.
Published SDK documentation build
Makefile, docs/package.json, .github/*, pnpm-workspace.yaml, scripts/classify-paths.sh
The docs site consumes the published SDK, docs targets no longer require build-ts, and release-age and Dependabot rules are updated.

Priority: ➖ Normal

Estimated code review effort: 5 (Critical) | ~120 minutes

Merge Risk: ⚪ Minimal · up to 2b89a

This change moves ingest and streaming to positional rows with schema frames, while preserving SDK row-object behavior. Invalid envelopes are rejected safely, and the current implementation is ready to merge.

Suggested reviewers: taitelee

🚥 Pre-merge checks | ✅ 5
✅ Passed checks (5 passed)
Check name Status Explanation
Title check ✅ Passed The title clearly identifies the main breaking change: positional rows for NATS and SSE with separately transmitted columns.
Description check ✅ Passed The description is directly related to the changeset and explains the wire protocol, computed-column handling, SSE behavior, deployment impact, tests, and review scope.
Docstring Coverage ✅ Passed Docstring coverage is 84.75% which is sufficient. The required threshold is 80.00%. Docstring coverage is scoped to functions touched by this diff. Analyzed 118 functions across 30 files.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
✨ Finishing Touches
📝 Generate docstrings
  • Create stacked PR
  • Commit on current branch
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch stack/5-positional-wire
✨ Simplify code
  • Create PR with simplified code
  • Commit simplified code in branch stack/5-positional-wire

Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@github-actions github-actions Bot added documentation Improvements or additions to documentation go Pull requests that update go code area/observability Metrics, logs, traces, health, profiling area/api HTTP handlers, routing, middleware area/ingest Ingest pipeline (Bento, batching, DLQ) area/policy Access control policies (Hasura-style) area/sdk TypeScript SDK (clients/ts/) area/docs Documentation, site/, README labels Sep 2, 2026
@EricAndrechek

Copy link
Copy Markdown
Member Author

@coderabbitai review

Auto-review is off for PRs whose base isn't the default branch, and this one sits in a stack (base stack/4-seams). Requesting a review explicitly.

Note for context: this PR's diff against its base is one layer of a seven-PR stack replacing #540. Reviewing it in isolation is the intent — the layers below it are already in its history.

@coderabbitai

coderabbitai Bot commented Sep 2, 2026 •

Copy link
Copy Markdown

@EricAndrechek I will review #554 against its stack/4-seams base. I will limit the review to this stack layer.

⚠️ Action not completed

Review rate limited.

Note: CodeRabbit is an incremental review system and does not re-review already reviewed commits. This command is applicable only when automatic reviews are paused.

EricAndrechek and others added 10 commits September 2, 2026 07:13
`database` is a live getter so a ClickHouse reconfigure is honored on the next
refresh. Reading it once per query rather than once per refresh let a
reconfigure land between the system.columns and system.tables scans, attaching
DDL from the new database to same-named schemas discovered from the old one —
silently, and it would survive into GET /v1/ops/schema. Refresh now snapshots
the name before the first query and threads it to both.

Both refresh-failure tests asserted emptiness after a single failed refresh,
which is satisfied by the regression they exist to prevent — a failure that
wipes an already-published registry. Both now seed a successful refresh first.
TestRefresh_DatabaseSnapshottedForWholeRefresh flips the getter between the two
scans and fails without the snapshot.

CORRECTS TWO CLAIMS IN THIS BRANCH'S EARLIER COMMIT (e10eb04), whose message
is already pushed and cannot be amended under the no-force-push rule:

1. "an external-engine table carries its credentials in that statement" is
   FALSE on the ClickHouse this repo ships. Verified on 26.7.3.19: MySQL,
   PostgreSQL, S3, Kafka SASL, MongoDB URIs and S3 presigned signatures all
   render the secret as `[HIDDEN]`. What create_table_query does leak,
   unconditionally, is topology — endpoint, bucket or host, database, username,
   S3 access key id. `json:"-"` is still right, for that stronger reason.
   The password is exposed only on a pre-masking server or one with the
   server-level display_secrets_in_show_and_select enabled.

2. "captured on the same refresh so they can never describe different servers"
   overstated the guarantee. chconn.Manager resolves the connection per call,
   so a reload changing clickhouse.addr mid-refresh can pair a version from one
   server with schemas from another. It is a publication guarantee, not a
   same-server one. The read side straddles too: ServerVersion() and Get() take
   separate RLocks.

The rationale is corrected at every site it was replicated — discovery.go,
AGENTS.md, architecture.md, api.md, the CHANGELOG — and the test fixture, which
had been asserting a literal secret no real server emits, now uses masked DDL
and asserts on the bucket and access key id instead.

Adds tests/integration/discovery_metadata_test.go: position is 1-based and
contiguous against live ClickHouse, DefaultExpression and HasDefault behave as
documented for a MATERIALIZED column, DDL is captured, ServerVersion is set.

Raised by CodeRabbit on #550 and by the pre-push reviewer.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018Epn88jTEw4ZkXrvTKzZXQ
The unsupported-media-type test checked the status and two of the five accepted
types. It now asserts the JSON error envelope and its security headers via
testutil.AssertJSONErrorResponse — which this branch was missing, though the
same assertion exists further up the stack, so the split dropped it onto the
wrong layer — and iterates supportedContentTypes.

On what that loop pins: it cannot catch an alias being dropped from
supportedContentTypes, because the message is built from the same slice and the
loop would simply check one fewer. Verified. What it does catch is the message
diverging from the list, so a client is never told to use a type the server
accepts but never names; verified by truncating the Join and watching it fail.

Also rewords the documented 415 cause in three places: "or one ingest doesn't
read" was not a condition a reader could act on. Now "or an unsupported one".

Raised by CodeRabbit on #552.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018Epn88jTEw4ZkXrvTKzZXQ
The pre-v2 layout detector treated `tables.<table>.select` as an operation map
whenever its value's keys were not permission fields, so a valid v2 grant for a
role literally named `select` was rejected before the strict decode. Raised by
CodeRabbit on #551.

The first fix bailed out whenever every key was an operation name. That was
wrong, and wrong in the dangerous direction — it made this document adopt:

    {"tables":{"clicks":{"select":{"select":{"allow_columns":["a"]},
                                   "insert":{"allow_columns":["b"]}}}}}

Pre-v2 that grants two read-only roles. Read as v2 it grants ONE role read and
write: role `select` gained an INSERT the file never contained, and role
`insert` silently lost its SELECT. No error, no warning. Caught by the pre-push
reviewer; reproduced before this fix.

The discriminator is narrower. When the inner key set is exactly the outer
operation (`select.select`), the v1 and v2 readings describe identical access,
so accepting the v2 decode is free. When it names the OTHER operation
(`select.insert`) the readings diverge, and that is precisely the escalating
case — so it is refused, with its own message asking for a rename rather than
the migration message, which would tell an operator to convert a file that may
already be v2.

TestReportLegacyPolicyLayout_RoleNamedAfterAnOperation pins all three outcomes
and asserts the divergent case does not get the migration message. The previous
version of that test asserted the escalating shape as safe.

Also corrects `settings-directory.mdx`, which said a leftover `"select": {}`
block always fails the undeclared-role check first — that holds only when
`roles.json` does not declare a role named `select` — and qualifies the
"reported as one clear finding" claims in the CHANGELOG and access-control.mdx,
which are not true of the ambiguous shape.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018Epn88jTEw4ZkXrvTKzZXQ
Four items from the chtypes-side review of #551:

- `looksLikeRoleMap`'s doc comment still described the withdrawn first
  attempt, where that function itself special-cased an operation-named
  key and the ambiguous case "lost" to the v2 reading. Neither is true:
  the body only rejects permission-field keys, and the ambiguous case is
  refused by `operationNamedRoles`, not resolved. Rewritten to match.
- Re-add the two operator-less rows dropped in the role-first rewrite of
  `validate_test.go`. `policy_test.go` still pinned the rule, but the
  settings layer is where an operator actually hits it.
- Drop `ResolvedInsert.CheckPredicates`. It had no reader and a rule that
  deliberately diverges from the one ingest enforces (fail-closed vs
  auto-inject ""), which is a trap to leave lying around unread. It can
  come back in the PR that reads it, reviewed against a real consumer.
- Note at `ingest.go`'s bare `CheckClauses` read why the `unresolved`
  marker cannot defend that line, at the call site rather than only on
  the type.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018Epn88jTEw4ZkXrvTKzZXQ
Pre-push review findings on the previous commit:

- The `CheckClauses` call-site comment I just added claimed "Evaluate only
  returns Allowed with both sides resolved". That is false, and false in
  the direction that matters: `evaluateSelect` returns `Allowed: true`
  with `Insert` marked `unresolved`, which is the whole point of the
  marker. A maintainer trusting that sentence would conclude a cross-side
  bare read is safe anywhere. The real reason this line is safe is
  narrower — the handler resolved this grant for `insert` — so say that.
- `ResolvedPredicate`/`ResolvePredicates` were exported by this PR only
  because `ResolvedInsert.CheckPredicates` needed a public type. That
  field is gone and nothing outside `internal/policy` names either
  identifier, so they go back to main's unexported spelling, along with
  the two prose mentions that moved with them.
- `architecture.md`'s rowfilter bullet lost its third list item to a
  spliced-in clause: `ColumnSpec` read as a coordinate of the
  `HasRowFilter` gate explanation with no verb attaching it. The gate
  explanation moves to the end of the bullet.
- The CHANGELOG's "without it" attached the `RowVisible`-true consequence
  to the insert-side bind-unsafe deny; that consequence belongs to the
  select-side operator-less deny. Check clauses never reach `RowVisible`.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018Epn88jTEw4ZkXrvTKzZXQ
Three CodeRabbit threads on #551:

- `access-control.mdx` claimed an old-layout document is always reported
  as one clear error. An *empty* operation block (`"select": {}`) carries
  no role names to give it away, so it reads as a grant for a role named
  `select` and either trips the undeclared-role error or — if `select` is
  declared — warns and adopts. State the exception.
- `settings-directory.mdx` described only the migration rejection. The
  ambiguity refusal (a role named after the *other* operation) is a
  second, distinct one with its own message, and the operator who hits it
  was being told their file uses the pre-v2 layout. Both spellings are
  now stated, including that the same-operation form is accepted.
- Pin the mixed-key classification in `validate_test.go`: a pre-v2 block
  carrying one operation-named role beside a real one. The real name
  settles the reading, so it is pre-v2 and gets the migration message,
  not the ambiguity refusal.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018Epn88jTEw4ZkXrvTKzZXQ
EricAndrechek and others added 2 commits September 8, 2026 14:33
The anchor I wrote named the danger block by its title, but it is a
:::danger[...] admonition, which renders no anchor -- so the link failed
starlight-links-validator and took make ci down with it. Repointed at
Transport Behavior, the heading that actually contains that block.

The lesson is the one already in the notes: make verify does not check
docs links, so a link added after the last build-docs run is unverified
until the next one. I added this link after the build-docs I had already
run, and did not re-run it before committing.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Both reviewers, final round. No MUSTs from either; these are the SHOULDs
and MAYs.

The docs SDK pin was defeatable by our own automation. dependabot.yml
points the npm ecosystem at the workspace root with patterns ["*"] and
ignores only typescript and eventsource-parser, so the moment 0.2.0
publishes it would have proposed ^0.1.1 -> ^0.2.0 on docs/package.json
in the weekly grouped deps: PR. That bump would then land as routine
dependency hygiene, reviewed by someone not thinking about which wire
the demo backend speaks -- reintroducing #568 by the exact path the pin
exists to close. @wavehouse/sdk is now in the ignore list for minor and
major (for a 0.x package the minor is what crosses the wire boundary;
patches inside 0.1.x already satisfy the caret, so nothing routine is
lost), and the release runbook says so rather than implying the bump is
the only way the range can move.

Three more copies of claims this stack has been sweeping. ci.yml's
comment still described the docs job's chain as build-docs -> check-docs
-> build-ts, which is the coupling b379d4f removed -- the conclusion
(no Go toolchain) was right, the stated reason would have taught the
next editor the opposite of what classify-paths.sh was just reworded to
say. poisonCounter's own doc comment, the canonical statement of the
metric's semantics, stated the ack-gating rule for "parked" only, though
b379d4f made it govern "dropped" too. And api.md described the schema
frame as delivered once per connection, which the drift re-announcement
contradicts -- the same file says the right thing 200 lines earlier and
tells raw consumers to keep the most recent list, so this summary was
the one sentence that would have led an implementer to cache the first
announcement forever.

Also finished a rewrap applied to types.go but only half-applied in
hub.go, and moved the docs-pin step below the changelog step in the
release runbook so the section reads in the order you do it.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
@github-actions github-actions Bot added dependencies Pull requests that update a dependency file github_actions Pull requests that update GitHub Actions code area/infra CI, build, deploy, Docker, release labels Sep 8, 2026
@EricAndrechek

Copy link
Copy Markdown
Member Author

@coderabbitai rate limit

@EricAndrechek

Copy link
Copy Markdown
Member Author

@coderabbitai review

Six commits since your last pass (141a54ff → 4c27c547), answering taitelee's two review threads and three rounds of local reviewers. Both his threads are now replied to and resolved. Summary of what changed and why, so the review can be targeted:

  • 750a4aee — deleted wavehouse_sse_rows_unpairable_total outright. Its two call sites had inconsistent denominators (once per envelope in Broadcast, once per connection in ReplayProjector), and the premise it was added on was wrong: ingest.parseMsg refuses the identical conditions on the identical bytes first and parks them on the DLQ. Also added a Format check in newEventView, so the hub refuses an envelope whose declared format it doesn't know — it's a separate NATS consumer that acks independently of the worker, so without it a future second format would stream to SSE clients while the worker DLQ'd it. usable=false not decoded=false, because decoded=false routes into the legacy passthrough.
  • 3a374195 — removing that counter left a hole: poison_dropped_total fired only where the DLQ is off, so with it on (the recommended setting) an unreadable envelope had no metric at all. Renamed to wavehouse_ingest_poison_total with disposition="parked"|"dropped", counted only once the park actually succeeds.
  • 8277bf6e — docs/package.json moves @wavehouse/sdk from workspace:* to ^0.1.1. The docs site streams against a separately-deployed backend, so binding it to the workspace would ship this PR's unreleased wire to the live landing page on merge (docs site streams the workspace SDK against a pre-v2 backend on merge #568). Also exempted @wavehouse/* from minimumReleaseAge, since depending on our own package from the registry newly subjected it to the 7-day cooldown.
  • b379d4f1 — the drop path had the mirror of the parked bug: counted before a DoubleAck whose error was discarded, so a failed ack would re-count on every redelivery while reporting the row as gone. Now counted inside the goroutine after a successful ack.
  • bbf6f68f — link-anchor fix (a :::danger[…] title renders no anchor).
  • 4c27c547 — ignored @wavehouse/sdk in .github/dependabot.yml; the weekly grouped deps: PR would otherwise have proposed ^0.1.1 → ^0.2.0 and silently undone the pin above.

Worth your attention specifically: the ack/park ordering in rejectPoison, whether the Format check's fail-closed lever is the right one, and whether the dependabot ignore's update-types correctly covers a 0.x minor as the wire-crossing bump.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Actionable comments posted: 2


ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Advanced

Run ID: 68c1e4a1-29ac-46c9-8ee8-d23392027249

📥 Commits

Reviewing files that changed from the base of the PR and between a5d3a56 and 4c27c54.

⛔ Files ignored due to path filters (1)
  • pnpm-lock.yaml is excluded by !**/pnpm-lock.yaml
📒 Files selected for processing (20)
  • .github/dependabot.yml
  • .github/workflows/ci.yml
  • CHANGELOG.md
  • Makefile
  • docs/package.json
  • docs/src/content/docs/api.md
  • docs/src/content/docs/architecture.md
  • docs/src/content/docs/deployment.md
  • docs/src/content/docs/development.md
  • docs/src/content/docs/getting-started.md
  • docs/src/content/docs/sdk/streaming.md
  • docs/src/content/docs/settings-directory.mdx
  • internal/ingest/types.go
  • internal/ingest/worker.go
  • internal/ingest/worker_test.go
  • internal/stream/hub.go
  • internal/stream/hub_test.go
  • internal/stream/metrics.go
  • pnpm-workspace.yaml
  • scripts/classify-paths.sh

Included review availability: Your plan provides up to 1 included review per hour; 0 remain after this review.

📜 Review details
🧰 Additional context used
📓 Path-based instructions (2)
WH001 applies to every tracked Markdown file, with no carve-out

📄 CodeRabbit inference engine (AGENTS.md)

Files:

  • docs/src/content/docs/development.md
  • docs/src/content/docs/architecture.md
  • docs/src/content/docs/getting-started.md
  • docs/src/content/docs/settings-directory.mdx
  • docs/src/content/docs/sdk/streaming.md
  • docs/src/content/docs/api.md
  • docs/src/content/docs/deployment.md
  • CHANGELOG.md
In MDX, leave a blank line between a JSX tag and a code fence.

📄 CodeRabbit inference engine (AGENTS.md)

Files:

  • docs/src/content/docs/settings-directory.mdx
🧠 Learnings (3)
📓 Common learnings
Learnt from: EricAndrechek
URL: https://github.com/Wave-RF/WaveHouse/pull/554

Timestamp: 2026-09-05T00:43:57.044Z
Learning: In WaveHouse positional ingest, `internal/api/ingest.go` derives `EventMessage.Columns` from `discovery.TableSchema.InsertableColumnNames()`, which reflects unique `system.columns` names, and the embedded NATS server uses `DontListen: true`; therefore duplicate column names are unreachable from the current in-process producer. `internal/ingest/worker.go`, `internal/stream/hub.go`, and `clients/ts/src/stream/sse.ts` must still reject duplicate column names defensively because replayed DLQ messages or future external publishers can supply them, and `pairRow` otherwise overwrites an earlier cell in its name-keyed map before row-level policy evaluation.
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse

Timestamp: 2026-09-05T00:26:39.915Z
Learning: In `docs/src/content/docs/deployment.md`, a pre-v2 ingest envelope is permanently acknowledged and dropped when the table DLQ is disabled. It can be manually recovered only when the DLQ is enabled and the envelope is parked on `dlq.{table}`. During the `stream.gap_window_minutes` retention window, SSE gap-fill silently withholds pre-v2 envelopes because they cannot be paired into positional rows.
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse

Timestamp: 2026-09-08T12:06:56.369Z
Learning: In Wave-RF/WaveHouse, insert policy-check validation reports every invalid checked column in sorted order. This prevents nondeterministic Go map-iteration error messages when multiple policy `check` columns are invalid.
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse

Timestamp: 2026-09-08T13:17:06.219Z
Learning: For WaveHouse PR stacks, document a protocol caveat in the PR that introduces the protocol behavior, not only in a later stacked PR. For the positional SSE protocol introduced by PR `#554`, documentation must state that gap-fill across a schema change can deliver live rows without a new schema frame; raw consumers must use the latest received schema, drop rows with mismatched arity, and reconnect because arity checks cannot detect same-length schema changes such as `RENAME COLUMN`.
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse PR: 554
File: docs/src/content/docs/api.md:613-613
Timestamp: 2026-09-05T00:26:39.936Z
Learning: In WaveHouse SSE streaming, `SubscribeSchemaFrame` is best-effort and is used for schema de-duplication and quiet-table connections. The live-delivery guarantee comes from `internal/stream/hub.go` `deliver` checking `Subscriber.needsSchema` before sending a row. `ReplayProjector` tracks schema state separately with closure-local `lastSig`; replay and live schema state do not reconcile across a schema change (`#543`). Raw SSE consumers must check row arity and reconnect to resynchronize a same-length schema change.
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse

Timestamp: 2026-09-08T12:06:54.852Z
Learning: In WaveHouse documentation that describes raw SSE positional rows, state that the stream sends an `event: schema` frame before rows and re-sends it when the column list changes. State that each `row` array matches the most recent schema. This prevents raw consumers from pairing later rows with a stale column list and silently mislabeling values.
Learnt from: CR
Repo: Wave-RF/WaveHouse

Timestamp: 2026-09-08T18:51:30.141Z
Learning: Every new function should have corresponding test cases.
Learnt from: CR
Repo: Wave-RF/WaveHouse

Timestamp: 2026-09-08T18:51:30.141Z
Learning: Validate locally before pushing.
Learnt from: CR
Repo: Wave-RF/WaveHouse

Timestamp: 2026-09-08T18:51:30.141Z
Learning: Validate locally before every push
Learnt from: CR
Repo: Wave-RF/WaveHouse

Timestamp: 2026-09-08T18:51:30.141Z
Learning: Run `make lint` and `make test` before considering work complete.
📚 Learning: 2026-08-13T12:17:52.620Z
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse PR: 470
File: docs/src/content/docs/reverse-proxy.mdx:137-144
Timestamp: 2026-08-13T12:17:52.620Z
Learning: For Wave-RF/WaveHouse documentation, verify claims about implementation control flow against the authoritative implementation source (for example, internal/auth/auth.go) rather than relying solely on docs/** content. Documentation may lag behind or paraphrase behavior, so control-flow claims should be confirmed in source code.

Applied to files:

  • docs/src/content/docs/settings-directory.mdx
📚 Learning: 2026-06-26T12:23:22.696Z
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse PR: 346
File: internal/stream/subscriber_test.go:9-28
Timestamp: 2026-06-26T12:23:22.696Z
Learning: In this Go repository, prefer table-driven tests (e.g., `[]struct{...}` with `t.Run(...)`) only for tests that cover multiple scenarios/inputs and can be cleanly enumerated. Do not artificially rewrite a clear single-scenario sequential behavioral-flow test into a table-driven form just to fit the pattern; if there’s only one meaningful scenario, keep the test as a straightforward linear flow (as in `TestSubscriber_SendDeliversThenDropsWhenFull`).

Applied to files:

  • internal/stream/hub_test.go
🪛 LanguageTool
docs/src/content/docs/development.md

[uncategorized] ~604-~604: The official name of this software platform is spelled with a capital “H”.
Context: ...use/sdkis in the npmignorelist in.github/dependabot.yml` for exactly this reason...

(GITHUB)

CHANGELOG.md

[typographical] ~23-~23: Consider using an em dash in dialogues and enumerations.
Context: - **The docs site now consumes the *publi...

(DASH_RULE)


[style] ~23-~23: Since ownership is already implied, this phrasing may be redundant.
Context: ...gainst a separately-deployed backend on its own release cadence, but took its SDK from ...

(PRP_OWN)


[style] ~23-~23: ‘for want of’ might be wordy. Consider a shorter alternative.
Context: ...ing "live" while every frame is dropped for want of a schema announcement, with no error ...

(EN_WORDINESS_PREMIUM_FOR_WANT_OF)


[style] ~23-~23: Since ownership is already implied, this phrasing may be redundant.
Context: ...ately keeps workspace:*. Depending on our own package from the registry also made `mi...

(PRP_OWN)


[style] ~58-~58: Consider an alternative for the overused word “exactly”.
Context: ...JSON, an unknown row format (which is exactly what a pre-v2 message looks like), or c...

(EXACTLY_PRECISELY)

🔇 Additional comments (18)
.github/dependabot.yml (1)

95-96: LGTM!

.github/workflows/ci.yml (1)

144-144: LGTM!

CHANGELOG.md (1)

23-23: LGTM!

Also applies to: 58-58

Makefile (1)

286-286: LGTM!

Also applies to: 637-641

docs/src/content/docs/architecture.md (1)

210-211: LGTM!

docs/src/content/docs/deployment.md (1)

338-338: LGTM!

Also applies to: 349-349, 352-352

internal/ingest/worker.go (1)

86-88: LGTM!

Also applies to: 689-693, 698-710

docs/package.json (1)

24-24: 🗄️ Data Integrity & Integration

No lockfile change is required. @wavehouse/sdk@0.1.1 is published, the root lockfile records its registry integrity, and pnpm-workspace.yaml excludes @wavehouse/* from minimumReleaseAge.

docs/src/content/docs/api.md (1)

821-821: 🗄️ Data Integrity & Integration

No documentation change is needed. internal/stream.newEventView rejects unsupported format values before pairing columns and row in both live delivery and replay.

internal/ingest/worker_test.go (1)

1473-1476: 📐 Maintainability & Code Quality

No change is required for poisonCounter.

In go.opentelemetry.io/otel v1.46.0, the first otel.SetMeterProvider call delegates instruments created before registration to the new provider. This test's first call therefore connects poisonCounter to its ManualReader. Later calls are ignored, but the package has no other meter-provider installation. The claimed zero-data-point failure does not apply.

docs/src/content/docs/development.md (1)

32-32: LGTM!

Also applies to: 185-185, 187-189, 604-604

docs/src/content/docs/getting-started.md (1)

96-98: LGTM!

docs/src/content/docs/sdk/streaming.md (1)

100-100: LGTM!

Also applies to: 119-119, 121-121, 125-125, 127-127

docs/src/content/docs/settings-directory.mdx (1)

186-186: LGTM!

internal/ingest/types.go (1)

20-22: LGTM!

pnpm-workspace.yaml (1)

22-27: LGTM!

Also applies to: 30-30

scripts/classify-paths.sh (1)

12-15: LGTM!

internal/stream/metrics.go (1)

23-29: LGTM!

Also applies to: 49-49

Comment thread internal/ingest/worker.go
Comment thread internal/stream/hub.go
@coderabbitai

coderabbitai Bot commented Sep 8, 2026

Copy link
Copy Markdown

Your plan includes PR reviews subject to rate limits. More reviews will be available in 45 minutes.

@coderabbitai

coderabbitai Bot commented Sep 8, 2026

Copy link
Copy Markdown

Tip

For best results, initiate chat on the files or code changes.

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
set -euo pipefail

repo="Wave-RF/WaveHouse"
pr="554"

echo "--- PR state ---"
gh pr view "$pr" --repo "$repo" \
  --json url,baseRefName,baseRefOid,headRefName,headRefOid,mergeable,reviewDecision,statusCheckRollup \
  --jq '{url,baseRefName,baseRefOid,headRefName,headRefOid,mergeable,reviewDecision,checks:[.statusCheckRollup[] | {name,conclusion,status}]}'

base="$(gh pr view "$pr" --repo "$repo" --json baseRefOid --jq .baseRefOid)"
head="$(gh pr view "$pr" --repo "$repo" --json headRefOid --jq .headRefOid)"
printf 'base=%s\nhead=%s\n' "$base" "$head"

echo "--- stated commit availability and range ---"
for rev in 141a54ff 750a4aee 3a374195 8277bf6e b379d4f1 bbf6f68f 4c27c547; do
  git show -s --format='commit=%H%nparents=%P%ntree=%T%nsubject=%s' "$rev"
done
git merge-base --is-ancestor 141a54ff "$head" && echo "PASS: start commit is in head history"
git merge-base --is-ancestor 4c27c547 "$head" && echo "PASS: end commit is in head history"

echo "--- effective diff hygiene and changed paths ---"
git diff --check "$base" "$head"
git diff --name-status "$base" "$head"
git diff --stat "$base" "$head"

echo "--- six stated commit patches ---"
git show --format=fuller --find-renames 750a4aee 3a374195 8277bf6e b379d4f1 bbf6f68f 4c27c547 -- \
  internal/ingest/worker.go internal/ingest/worker_test.go \
  internal/stream/hub.go internal/stream/hub_test.go internal/stream/metrics.go \
  docs/package.json pnpm-workspace.yaml .github/dependabot.yml \
  docs/src/content/docs docs/src/content/docs/api.md CHANGELOG.md

echo "--- outlines ---"
for file in internal/ingest/worker.go internal/ingest/worker_test.go internal/stream/hub.go internal/stream/hub_test.go internal/ingest/types.go internal/stream/metrics.go; do
  echo "### $file"
  ast-grep outline "$file" --items all || true
done

echo "--- poison acknowledgement, DLQ, metrics, and tests ---"
rg -n -C12 'rejectPoison|parkOnDLQ|sendToDLQ|DoubleAck|poison.*(parked|dropped)|Poison' \
  internal/ingest --glob '*.go'

echo "--- stream envelope decode and fallback paths ---"
rg -n -C14 'newEventView|FormatJSONCompactEachRow|decoded|usable|legacy|ReplayProjector|Broadcast|pairRow' \
  internal/stream internal/ingest/types.go --glob '*.go'

echo "--- exact implementations ---"
sed -n '160,400p' internal/ingest/worker.go
sed -n '400,720p' internal/ingest/worker.go
sed -n '120,310p' internal/stream/hub.go
sed -n '310,610p' internal/stream/hub.go
cat -n internal/stream/metrics.go
cat -n internal/ingest/types.go

echo "--- docs dependency and Dependabot configuration ---"
cat -n docs/package.json
cat -n pnpm-workspace.yaml
cat -n .github/dependabot.yml
rg -n -C5 '"`@wavehouse/sdk`"|`@wavehouse/sdk`|minimumReleaseAgeExclude|update-types|version-update:semver' \
  docs clients .github pnpm-workspace.yaml --glob 'package.json' --glob '*.yml' --glob '*.yaml' --glob '*.md' --glob '*.mdx'

echo "--- all direct poison metric references ---"
rg -n -C5 'wavehouse_ingest_poison|Poison.*(Parked|Dropped)|Poison' \
  . --glob '*.go' --glob '*.md' --glob '*.mdx' --glob '*.yaml' --glob '*.yml'

Length of output: 181332


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "--- complete DLQ helper and immediate callers ---"
sed -n '670,770p' internal/ingest/worker.go
rg -n -C12 'parkOnDLQ\s*\(|DLQ.*(Ack|ack)|Ack.*DLQ|DoubleAck.*(Fail|fail)|Poison.*(Count|count|Disposition|disposition)' \
  internal/ingest/worker_test.go internal/testutil --glob '*.go'

echo "--- poison metric test in full ---"
sed -n '1450,1560p' internal/ingest/worker_test.go

echo "--- mock JetStream message acknowledgement behavior ---"
rg -n -C15 'type MockJetStreamMsg|func \(.*MockJetStreamMsg.*DoubleAck|DoubleAckErr|AckErr' \
  internal/testutil internal/ingest/worker_test.go --glob '*.go'

echo "--- tests for unknown-format stream envelopes and legacy fallback ---"
rg -n -C12 'UnknownFormat|unknown format|FormatJSONCompactEachRow|newEventView|legacy.*passthrough|passthrough' \
  internal/stream/hub_test.go internal/stream --glob '*_test.go'

echo "--- patch context for the six requested commits ---"
git show --format= --unified=20 750a4aee 3a374195 b379d4f1 -- \
  internal/ingest/worker.go internal/ingest/worker_test.go internal/stream/hub.go internal/stream/hub_test.go

Length of output: 124073


🤖 web_search executed:

GitHub Dependabot options ignore update-types version-update:semver-minor documentation. Does Dependabot classify a 0.x minor version increase, such as 0.1.1 to 0.2.0, as version-update:semver-minor?

Length of output: 10201


@EricAndrechek Full review found one issue.

  • P2 — Do not count a parked envelope before its original message is acknowledged. In internal/ingest/worker.go, parkOnDLQ publishes the DLQ copy, ignores an error from natsMsg.DoubleAck(ctx), and returns true. rejectPoison then records wavehouse_ingest_poison_total{disposition="parked"}. If that acknowledgement fails, the original message remains in the ingest stream and is redelivered. Each retry can publish another DLQ copy and increment parked. This is the same counting-before-ack error that disposition="dropped" now avoids. Make parkOnDLQ return false when DoubleAck fails, log the error, and add a test for successful DLQ publish plus failed acknowledgement. That test should assert that no parked increment occurs.

The newEventView format check is correct. It leaves decoded=true and usable=false, so planForRole rejects the envelope instead of entering legacy passthrough. The live and replay paths both use this behavior.

The Dependabot rule is correct. Dependabot treats 0.1.1 to 0.2.0 as version-update:semver-minor. The semver-minor and semver-major ignore entries block wire-crossing range updates while allowing patch updates. GitHub’s Dependabot options reference confirms the update-type classification.

I used the existing ingest shutdown invariant during this review. The ackWg ownership remains race-safe.


🧠 Learnings used
Learnt from: EricAndrechek
URL: https://github.com/Wave-RF/WaveHouse/pull/554

Timestamp: 2026-09-08T13:09:42.927Z
Learning: In WaveHouse `internal/stream/hub.go`, `Hub.Broadcast` and `Hub.ReplayProjector` must each call `Metrics.RowUnpairable` when `newEventView` decodes an `ingest.EventMessage` but `pairRow` cannot pair its columns and row. `internal/stream/hub_test.go` verifies this contract with an OpenTelemetry `ManualReader`: one duplicate-column envelope through each path produces `wavehouse_sse_rows_unpairable_total` value 2 with `table=clicks` and `reason=duplicate_column`. `TestPairRow_Reasons` pins the reason-label strings separately.

Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse

Timestamp: 2026-09-05T00:26:46.671Z
Learning: In `internal/ingest/worker.go`, when a row still fails during `flushGroup` row isolation and the table DLQ is disabled, the worker does not `DoubleAck` the message. NATS redelivers it while successful rows can still reach ClickHouse. `wavehouse_ingest_poison_dropped_total` is incremented by `rejectPoison` only for unreadable envelopes dropped with the DLQ disabled; it does not report twice-failed ClickHouse inserts. For the latter case, the `ERROR` log `isolated bad row, DLQ disabled for table` is the only signal.

Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse

Timestamp: 2026-09-05T00:26:42.746Z
Learning: In `internal/ingest/worker.go`, `IngestWorker.rejectPoison` backgrounds both DLQ parking and disabled-DLQ `DoubleAck` operations with `IngestWorker.ackWg.Go`. `IngestWorker.parseMsg` runs only on `IngestWorker.dispatchLoop`, and that same dispatch-loop goroutine calls shutdown and `IngestWorker.ackWg.Wait()`. Therefore, `ackWg.Add` can originate from either a table loop through `handleSuccess` or the dispatch loop through `rejectPoison`, without racing the shutdown wait.

You are interacting with an AI system.

EricAndrechek and others added 2 commits September 8, 2026 15:37
Both CodeRabbit findings on 4c27c54. Both real; one fixed at a wider
point than proposed.

parkOnDLQ returned true when the publish succeeded but the DoubleAck
failed, contradicting the contract its own doc comment states. That
failed ack leaves the envelope in the ingest stream, so NATS redelivers
it, parseMsg refuses it again, and parkOnDLQ publishes a SECOND copy to
dlq.<table> and counts another "parked". An operator draining the DLQ
then re-ingests the same envelope more than once, and the parked total
the upgrade runbook points them at overstates by one per redelivery. It
is the same hazard b379d4f closed on the drop path, on the branch I did
not revisit. The duplicate DLQ copy is the residual cost: PublishMsg is
not idempotent, so only the count and the log can be corrected here.

pairRow accepted a zero-column envelope. CodeRabbit reported it as a
JSON null row -- null unmarshals to a nil slice, length zero, which
matches a zero-length column list -- and proposed rejecting a nil cells
slice. That fixes the spelling rather than the fault: `[]` unmarshals to
an EMPTY slice, not a nil one, and pairs with zero columns just as
happily. I probed all four shapes before choosing, and `cells == nil`
leaves {cols: [], row: []} accepted. Rejecting len(cols) == 0 closes
both, and matches what parseMsg already refuses on the write side, so
the two consumers agree about an envelope neither can read. The test
substituting the narrower guard fails on the empty-array case, which is
what makes the difference load-bearing rather than stylistic.

Neither shape is reachable from our own producer: a zero-insertable-
column table cannot be created in ClickHouse 26.6.3.62, and the worker
refuses an empty column list before publish. pairRow is defence against
an envelope we did not write, which is the whole reason it re-checks
what the worker already checked.

Both new cases fail under mutation: restoring the fire-and-forget ack,
and swapping the column guard for the nil-slice one.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Comment-only. The reviewer's findings on 5798365: that commit changed
behaviour under three comments and left all three describing the old one.

pairRow's doc comment enumerates why pairing fails, and the enumeration
reads as exhaustive. The empty column list this commit added is not
covered by any of the three it lists -- a zero-length row against a
zero-length list is a length that AGREES, which is precisely how the
shape slipped through in the first place. Named it.

parkOnDLQ's doc comment described false as the failed-publish case, so a
reader would take false to mean nothing was published. It now has two
meanings, and the second is the opposite: published successfully, ack
failed, a DLQ copy exists but the original is still in the stream and
the next redelivery will park another. Both stated, with the shared
conclusion -- do not count this as a parking -- said once.

And the test's control comment counted four negative rows when there are
now seven.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
@EricAndrechek

Copy link
Copy Markdown
Member Author

@coderabbitai review

Both of your findings on 4c27c547 are fixed and their threads answered and resolved. Two commits, now pushed (head 2b89ac84):

  • 5798365c — parkOnDLQ returns false on a failed DoubleAck after a successful publish, so a parking that will repeat on redelivery is no longer counted. Test case added with DoubleAckErr set, as you suggested.
  • 5798365c — for the null row: I fixed len(cols) == 0 rather than cells == nil. Probing the pre-fix function showed {cols: [], row: []} was accepted too — [] unmarshals to an empty, non-nil slice — so the nil guard would have closed one spelling of a two-spelling fault. Details and the probe output are in the thread. TestPairRow_Verdict gained three rows, and substituting cells == nil makes the empty-array case fail.
  • 2b89ac84 — comment-only: pairRow's and parkOnDLQ's doc comments both described the pre-fix behaviour (an enumeration that read as exhaustive, and false implying nothing was published).

Please re-review these two commits. If it comes back clean I'll follow with run pre-merge checks to clear the standing CHANGES_REQUESTED.

@coderabbitai

coderabbitai Bot commented Sep 8, 2026

Copy link
Copy Markdown

Tip

For best results, initiate chat on the files or code changes.

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
set -euo pipefail

repo="Wave-RF/WaveHouse"
pr="554"

echo "--- PR state ---"
gh pr view "$pr" --repo "$repo" \
  --json url,baseRefName,baseRefOid,headRefName,headRefOid,mergeable,reviewDecision,statusCheckRollup \
  --jq '{url,baseRefName,baseRefOid,headRefName,headRefOid,mergeable,reviewDecision,checks:[.statusCheckRollup[]|{name,conclusion,status}]}'

base="$(gh pr view "$pr" --repo "$repo" --json baseRefOid --jq .baseRefOid)"
head="$(gh pr view "$pr" --repo "$repo" --json headRefOid --jq .headRefOid)"
printf 'base=%s\nhead=%s\n' "$base" "$head"

echo "--- stated commits and ancestry ---"
for rev in 4c27c547 5798365c 2b89ac84; do
  git show -s --format='commit=%H%nparents=%P%ntree=%T%nsubject=%s' "$rev"
  git merge-base --is-ancestor "$rev" "$head" && echo "PASS: $rev is an ancestor of HEAD"
done

echo "--- effective diff and hygiene ---"
git diff --check "$base" "$head"
git diff --name-status "$base" "$head"
git diff --stat "$base" "$head"

echo "--- exact fix commits ---"
git show --format=fuller --find-renames 5798365c 2b89ac84 -- \
  internal/ingest/worker.go internal/ingest/worker_test.go \
  internal/stream/hub.go internal/stream/hub_test.go

echo "--- outlines ---"
for file in internal/ingest/worker.go internal/ingest/worker_test.go internal/stream/hub.go internal/stream/hub_test.go internal/ingest/types.go internal/stream/metrics.go; do
  echo "### $file"
  ast-grep outline "$file" --items all || true
done

echo "--- poison handling, acknowledgement ordering, and metric calls ---"
rg -n -C14 'func \(w \*IngestWorker\) (rejectPoison|parkOnDLQ)|rejectPoison\s*\(|parkOnDLQ\s*\(|sendToDLQ\s*\(|DoubleAck\s*\(|Poison.*(Parked|Dropped)|poison.*(parked|dropped)' \
  internal/ingest --glob '*.go'

echo "--- zero-column validation in all envelope consumers ---"
rg -n -C12 'func pairRow|pairRow\s*\(|len\(.*cols.*\)\s*==\s*0|zero.?column|empty.?column|nil.*row|FormatJSONCompactEachRow' \
  internal/ingest internal/stream --glob '*.go'

echo "--- focused implementations ---"
sed -n '180,390p' internal/ingest/worker.go
sed -n '650,780p' internal/ingest/worker.go
sed -n '135,260p' internal/stream/hub.go
sed -n '320,460p' internal/stream/hub.go

echo "--- focused tests and mock acknowledgement behavior ---"
rg -n -C18 'Test.*(DLQ|Poison|DoubleAck|Parked|PairRow|Zero|Empty|Nil)|DoubleAckErr|DoubleAck\s*\(' \
  internal/ingest/worker_test.go internal/stream/hub_test.go internal/testutil --glob '*.go'

echo "--- effective changed production hunks ---"
git diff --unified=20 "$base" "$head" -- \
  internal/ingest/worker.go internal/stream/hub.go internal/ingest/types.go

echo "--- comment association scan for changed Go files ---"
python3 - <<'PY'
import re, subprocess
base = subprocess.check_output(
    ["gh", "pr", "view", "554", "--repo", "Wave-RF/WaveHouse", "--json", "baseRefOid", "--jq", ".baseRefOid"],
    text=True,
).strip()
head = subprocess.check_output(
    ["gh", "pr", "view", "554", "--repo", "Wave-RF/WaveHouse", "--json", "headRefOid", "--jq", ".headRefOid"],
    text=True,
).strip()
paths = subprocess.check_output(
    ["git", "diff", "--name-only", base, head, "--", "*.go"], text=True
).splitlines()
for path in paths:
    lines = open(path, encoding="utf-8").read().splitlines()
    for i, line in enumerate(lines):
        match = re.match(r'func\s+(?:\([^)]*\)\s*)?([A-Za-z_]\w*)\s*\(', line)
        if not match:
            continue
        name = match.group(1)
        j, comments = i - 1, []
        while j >= 0 and (not lines[j].strip() or lines[j].startswith("//")):
            if lines[j].startswith("//"):
                comments.append(lines[j][2:].strip())
            j -= 1
        if comments:
            word = re.match(r'([A-Za-z_]\w*)', comments[-1])
            if word and word.group(1) != name and not comments[-1].startswith(("TODO", "NOTE", "BUG")):
                print(f"{path}:{i+1}: comment begins {word.group(1)!r}; function is {name!r}")
PY

Length of output: 215699


@EricAndrechek The re-review is complete.

The DLQ acknowledgement fix is correct. parkOnDLQ returns false after a failed DoubleAck. rejectPoison does not increment wavehouse_ingest_poison_total{disposition="parked"} in that case. TestRejectPoison_CountedByDisposition covers this path.

The zero-column validation is correct. pairRow rejects len(cols) == 0. This rejects both row:null and row:[] for an empty column list. The test cases would not pass with a cells == nil check.

The updated comments describe both parkOnDLQ failure cases and the empty-column cause.

git diff --check is clean. GitHub reports the PR as mergeable. All required checks succeeded.

I found no additional actionable issues.

[approve]


✏️ Learnings added
Learnt from: EricAndrechek
URL: https://github.com/Wave-RF/WaveHouse/pull/554

Timestamp: 2026-09-08T20:03:34.463Z
Learning: In WaveHouse `internal/stream/hub.go`, `pairRow` must reject `len(cols) == 0`, not only a nil decoded row slice. A zero-column positional envelope can encode its row as either `null`, which unmarshals to a nil slice, or `[]`, which unmarshals to an empty non-nil slice. Both forms must fail closed so stream handling matches `internal/ingest/worker.go` `parseMsg`.

Note: Learnings are effective only in the context of similar code segments. To apply general coding guidelines, please configure review instructions. You can manage existing learnings in the app.

You are interacting with an AI system.

@github-project-automation github-project-automation Bot moved this from In review to In progress in WaveHouse Task Board Sep 8, 2026
@EricAndrechek
EricAndrechek added this pull request to the merge queue Sep 8, 2026
Merged via the queue into main with commit 4f05b18 Sep 8, 2026
19 checks passed
@github-project-automation github-project-automation Bot moved this from In progress to Done in WaveHouse Task Board Sep 8, 2026
@EricAndrechek
EricAndrechek deleted the stack/5-positional-wire branch September 8, 2026 23:55
EricAndrechek added a commit that referenced this pull request Sep 9, 2026
#554 landed as squash 4f05b18. This branch had merged #554 at 141a54f,
BEFORE the eight review commits that PR accumulated, so main carried work
the branch did not have.

That makes `-s ours` wrong here, and the interlock is what caught it: its
second condition is that the squash's content already be present in the
branch, and it was not. `-s ours` discards main's side, so it would have
silently reverted the whole #554 review round the moment #555 merged --
the deleted SSE unpairable metric, the newEventView format check, the
poison counter rename and both of its ack-ordering fixes, the docs SDK
pin, the dependabot ignore, and the zero-column pairRow guard. A normal
merge instead, resolved per hunk.

Resolution rule: main is authoritative for anything #554 owns, and the
branch is authoritative for what its own commits authored (b09dbbe,
d686fa5, db3d8e4, 8f017f4, 4a33fd1). Establishing which was which
mattered -- diffing the branch tip against 141a54f conflates #555's work
with #554 refinements the branch never received, so authorship came from
`git log -S` per contested string rather than from the diff direction.

Two things that would not have surfaced on their own. metrics.go
auto-merged CLEANLY into a stale state, keeping the unpairable counter
and RowUnpairable that #554 deleted -- orphaned dead code, no conflict
marker, compiles and passes tests; the branch's own tip already matched
main there, so the staleness came from its earlier merge of 141a54f.
And taking main's CHANGELOG wholesale dropped the computed-columns
bullet this PR exists to add, because a whole-file checkout takes the
file, not the hunk; re-applied from b09dbbe.

Verified after resolving: every #554 review artifact present, both
removals actually gone, #555's computed-column exclusion intact, and the
only files differing from main among #554's 50 are the six carrying
#555's own authored content.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
EricAndrechek added a commit that referenced this pull request Sep 9, 2026
Both from the reviewer, both introduced by my merge resolution.

AGENTS.md still named wavehouse_ingest_poison_dropped_total. It was the
only surviving instance in the repo, and I missed it because my
stale-claim sweep grepped internal/, docs/src and CHANGELOG.md and never
AGENTS.md -- the one file whose whole job is to be the invariant index a
future reader trusts. Two more errors in the same sentence, fixed in the
same pass: it attached the counter to the dropped case alone, though
rejectPoison counts both dispositions and the disposition label is the
entire point of the rename; and it listed two of the three unreadable
cases, omitting malformed JSON.

The computed-columns CHANGELOG bullet I re-applied after the whole-file
checkout dropped it turned out to duplicate the entry main already
carries -- #554 had absorbed this branch's work, so main's version is
the fuller and corrected one. The two disagreed on the ClickHouse
version: mine said 26.7.3, main's 26.6.3, and 26.6.3.62 is what the
compose files, CI and the integration setup actually pin. Deleted mine,
but folded in the one claim only it carried and which appears nowhere in
main's changelog: that a record SUPPLYING a value for a computed column
is now rejected 400 rather than accepted and silently dropped. Marked
that entry BREAKING, which it was not.

The other 26.7.3 in this file is pre-existing on main, in the
schema-discovery entry, and is left alone.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area/api HTTP handlers, routing, middleware area/docs Documentation, site/, README area/infra CI, build, deploy, Docker, release area/ingest Ingest pipeline (Bento, batching, DLQ) area/observability Metrics, logs, traces, health, profiling area/policy Access control policies (Hasura-style) area/query Structured query AST, SQL builder area/sdk TypeScript SDK (clients/ts/) dependencies Pull requests that update a dependency file documentation Improvements or additions to documentation github_actions Pull requests that update GitHub Actions code go Pull requests that update go code

Projects

Archived in project

Development

Successfully merging this pull request may close these issues.

2 participants