Skip to content

feat(pipes): resolve table dependencies for cache invalidation - #343

Open
taitelee wants to merge 32 commits into
mainfrom
pipe-cache-invalidation
Open

feat(pipes): resolve table dependencies for cache invalidation#343
taitelee wants to merge 32 commits into
mainfrom
pipe-cache-invalidation

Conversation

@taitelee

@taitelee taitelee commented Jun 11, 2026

Copy link
Copy Markdown
Member

Summary

Pipes were cached TTL-only because they couldn't report which tables they read, so an ingest never invalidated a stale pipe result (#178). This teaches pipes dependency-aware invalidation on the same namespace-versioned path structured queries already use — by asking ClickHouse, never by parsing SQL.

On a pipe's first execution with a given parameter binding, PipesHandler.resolveDeps runs EXPLAIN QUERY TREE over the bound SQL and reads the table set off ClickHouse's own analysis (discovery.ResolveTables). Resolution is per bound query because the table set can depend on parameter values: a parameter-gated UNION arm ClickHouse folds to constant-false is pruned, so ?source=web depends only on web_events. Resolutions are memoized per bound SQL (singleflight-coalesced; a generation counter keeps a schema refresh that lands mid-resolution from being silently overwritten by the in-flight result) and re-resolve on every schema refresh. Each resolved table folds into the cache key as its own versioned namespace, encoded exactly as the ingest worker encodes them, so a write to any of them evicts the result; the registry's refresh-time base-table→view cascade extends the same eviction to views and materialized views, and refreshes are serialized end-to-end so an older refresh can't install its cascade over a newer one.

Degraded paths fail toward expiry or over-invalidation, never staleness. A query ClickHouse can't analyze (a write/DDL pipe, a missing table, an unreachable server) falls back to a single database-wide version that every write bumps — any write evicts, O(1) per request — and a transient failure is retried on the next execution rather than memoized. A result that can't be reliably version-invalidated — an unfoldable view, an external read (a table function, a cross-database table, a non-local dictionary source), or a dead-branch-pruned dependency set — is TTL-capped (~10 s) so it self-expires. The version-folded cache key is snapshotted before the query executes and reused for its Set, so a write landing mid-query orphans the entry instead of re-homing pre-write data under the post-bump key (#382, fixed for both cached read paths).

Related Issues

Closes #178
Closes #382

@coderabbitai

coderabbitai Bot commented Jun 11, 2026

Copy link
Copy Markdown

Review 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
  • ✅ Review completed - (🔄 Check again to review again)
📝 Walkthrough

Walkthrough

The PR implements dependency-aware pipe cache invalidation. VersionManager is refactored to a two-level tableVersions/namespaceVersions model with BumpTable/BumpNamespace APIs. A new pipe_deps.go module resolves a pipe's referenced base tables via EXPLAIN QUERY TREE with recursive view expansion and backtick-aware identifier parsing. NamedQuery stores resolved tables and DummyBind generates safe runnable SQL for analysis. PipesHandler populates resolved tables at Put time and uses them to key cache entries in Execute, enabling ingest-triggered invalidation to evict stale pipe results.

Changes

Pipe dependency-based cache invalidation

Layer / File(s) Summary
VersionManager two-level versioning refactor
internal/cache/version_manager.go
Replaces single versions map with tableVersions and namespaceVersions; adds Namespace struct, NamespaceKey, QueryKey, BumpTable, BumpNamespace; removes GetCacheKey, GetVersion, IncrementVersion.
NamedQuery.ResolvedTables and DummyBind helper
internal/pipes/pipes.go, internal/pipes/dummybind_test.go
NamedQuery gains a server-owned ResolvedTables []string field; DummyBind generates runnable SQL for EXPLAIN analysis by substituting type-appropriate dummy values for all parameters; tests validate typed, boolean, bare inline, inline-default, and no-param cases.
Pipe SQL dependency extraction via EXPLAIN QUERY TREE
internal/api/pipe_deps.go
New pipe_deps.go implements best-effort base-table resolution through EXPLAIN QUERY TREE with recursive view expansion, cycle prevention, identifier parsing (backtick/escape handling), and schema-registry filtering; helpers include resolvePipeDeps, collectReadTables, explainQueryTreeTables, parseQueryTreeTables, filterKnownTables, splitQualified, unquoteIdent, and pipeDeps.
Pipe deps unit tests
internal/api/pipe_deps_test.go
Unit tests cover all parsing/normalization helpers and the pipeDeps namespace converter; fakeCache test double records cache dependency flow; handler tests assert deps pass through on cache HIT/MISS and server-side ownership of resolved_tables.
PipesHandler dependency-based cache wiring
internal/api/pipes.go
PipesHandler adds Registry and Database fields for best-effort resolution; Put calls resolvePipeDeps to populate ResolvedTables (server-owned); Execute computes deps from resolved tables and passes them to Cache.Get/Cache.Set instead of always nil, enabling ingest invalidation to reach pipe cache entries.
Main wiring and integration tests
cmd/wavehouse/main.go, tests/integration/pipe_deps_test.go, clients/ts/src/types.ts, docs/src/content/docs/pipes.mdx, tests/integration/setup_test.go, CHANGELOG.md
main.go constructs pipesHandler separately to assign Registry/Database; integration tests exercise view→base-table resolution, materialized view source tracking, direct table references, UNION multi-table resolution, and MISS→HIT→invalidation→MISS cache cycles against a real ClickHouse instance; client type adds resolved_tables field; docs describe cache invalidation behavior and TTL fallback; setup improves table-name sanitization; CHANGELOG documents two-level versioning and dependency resolution.

Sequence Diagram(s)

sequenceDiagram
    participant Client
    participant PipesHandler
    participant resolvePipeDeps
    participant ClickHouse
    participant SchemaRegistry
    participant Cache

    rect rgba(100, 149, 237, 0.5)
        Note over Client,Cache: PUT /v1/pipes/{name}
        Client->>PipesHandler: PUT pipe SQL
        PipesHandler->>resolvePipeDeps: pipe definition
        resolvePipeDeps->>ClickHouse: DummyBind SQL → EXPLAIN QUERY TREE
        ClickHouse-->>resolvePipeDeps: table_name identifiers
        resolvePipeDeps->>ClickHouse: system.tables as_select (view expansion)
        ClickHouse-->>resolvePipeDeps: view SQL (recursive)
        resolvePipeDeps->>SchemaRegistry: filterKnownTables
        SchemaRegistry-->>resolvePipeDeps: resolved base table names
        resolvePipeDeps-->>PipesHandler: ResolvedTables []string
        PipesHandler->>Cache: store pipe with ResolvedTables
    end

    rect rgba(60, 179, 113, 0.5)
        Note over Client,Cache: GET /v1/pipes/{name}/execute
        Client->>PipesHandler: Execute request
        PipesHandler->>PipesHandler: pipeDeps(q.ResolvedTables) → []Namespace
        PipesHandler->>Cache: Get(sha, deps)
        Cache-->>PipesHandler: HIT or MISS
        alt MISS
            PipesHandler->>ClickHouse: run pipe SQL
            ClickHouse-->>PipesHandler: results
            PipesHandler->>Cache: Set(sha, deps, results)
        end
        PipesHandler-->>Client: results + X-Cache header
    end
Loading

Estimated code review effort

🎯 4 (Complex) | ⏱️ ~60 minutes

Possibly related PRs

  • Wave-RF/WaveHouse#314: Introduces the cache.Namespace struct and foundational VersionManager.QueryKey/BumpTable/BumpNamespace APIs that this PR's pipe invalidation logic is built on.
  • Wave-RF/WaveHouse#177: Modifies PipesHandler cache wiring in internal/api/pipes.go, directly overlapping with this PR's changes to Execute and Put cache integration.
  • Wave-RF/WaveHouse#172: Changes PipesHandler authorization in Execute, colliding with the same handler code paths modified here for dependency-aware cache keying.

Suggested labels

area/query, area/ingest

Suggested reviewers

  • EricAndrechek
🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 62.04% which is insufficient. The required threshold is 80.00%. Write docstrings for the functions missing them to satisfy the coverage threshold.
✅ Passed checks (4 passed)
Check name Status Explanation
Linked Issues check ✅ Passed The implementation satisfies issues [#178] and [#382] through dependency resolution, invalidation, read-your-writes protection, and cache-key snapshots.
Out of Scope Changes check ✅ Passed The supporting refactors, tests, documentation, changelog, and encoding updates are related to the stated cache invalidation objectives.
Title check ✅ Passed The title clearly and concisely describes the main change: resolving table dependencies for pipe cache invalidation.
Description check ✅ Passed The description directly explains dependency-aware pipe cache invalidation, implementation details, fallback behavior, and test coverage.
✨ Finishing Touches
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch pipe-cache-invalidation
✨ Simplify code
  • Create PR with simplified code
  • Commit simplified code in branch pipe-cache-invalidation

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 go Pull requests that update go code area/api HTTP handlers, routing, middleware area/ingest Ingest pipeline (Bento, batching, DLQ) area/cache Local / shared / tiered caching area/pipes Named query pipes area/docs Documentation, site/, README area/infra CI, build, deploy, Docker, release labels Jun 11, 2026

@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: 4


ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro Plus

Run ID: d7ff65f4-c449-4c3d-a7b3-4ef586eb7bbe

📥 Commits

Reviewing files that changed from the base of the PR and between 9661c82 and 27ae1f1.

📒 Files selected for processing (18)
  • CHANGELOG.md
  • cmd/wavehouse/main.go
  • internal/api/pipe_deps.go
  • internal/api/pipe_deps_test.go
  • internal/api/pipes.go
  • internal/api/structured_query.go
  • internal/cache/cache.go
  • internal/cache/cache_test.go
  • internal/cache/local.go
  • internal/cache/local_test.go
  • internal/cache/version_manager.go
  • internal/cache/version_manager_test.go
  • internal/ingest/worker.go
  • internal/ingest/worker_test.go
  • internal/pipes/dummybind_test.go
  • internal/pipes/pipes.go
  • internal/testutil/mocks.go
  • tests/integration/pipe_deps_test.go
💤 Files with no reviewable changes (1)
  • internal/cache/cache_test.go
📜 Review details
🧰 Additional context used
📓 Path-based instructions (7)
**/*.go

📄 CodeRabbit inference engine (AGENTS.md)

**/*.go: Use Go 1.26 with strict formatting enforced by gofumpt
Use structured logging with log/slog (JSON handler)
Use Chi v5 for HTTP routing
Return errors, don't panic. Wrap with fmt.Errorf("context: %w", err)
Use package naming: lowercase, single word (or abbreviated). internal/ enforces module privacy
No global state: Dependencies are passed explicitly (constructor injection)
Comment the why, not the what. Add a comment only when the reason isn't obvious from the code; a line that matches the surrounding pattern needs none. Keep comments to 1–2 lines
DRY — one source of truth. Before adding logic, look for an existing helper, type, or constant to reuse; before duplicating a rule, factor it into one place every caller reads
Leave it neater than you found it — within reason. Fix small, safe things in passing: a stale comment, an obvious typo, a misnamed local, dead code on your path

Files:

  • internal/api/structured_query.go
  • internal/testutil/mocks.go
  • internal/ingest/worker.go
  • internal/pipes/pipes.go
  • internal/pipes/dummybind_test.go
  • cmd/wavehouse/main.go
  • tests/integration/pipe_deps_test.go
  • internal/ingest/worker_test.go
  • internal/api/pipe_deps.go
  • internal/cache/local.go
  • internal/api/pipe_deps_test.go
  • internal/cache/local_test.go
  • internal/cache/cache.go
  • internal/api/pipes.go
  • internal/cache/version_manager.go
  • internal/cache/version_manager_test.go
internal/api/**/*.go

📄 CodeRabbit inference engine (AGENTS.md)

Chi HTTP router, JWT/JWKS middleware (from auth/), ingest/query/structured-query/SSE/schema/DLQ/policy/pipes handlers, Hub

Files:

  • internal/api/structured_query.go
  • internal/api/pipe_deps.go
  • internal/api/pipe_deps_test.go
  • internal/api/pipes.go
internal/ingest/**/*.go

📄 CodeRabbit inference engine (AGENTS.md)

Ingest worker pipeline (worker.go): JetStream input → per-table batch INSERT with DLQ output. The pipeline is insert-only. Wire format EventMessage carries {table_name, scope, received_timestamp, data} and nothing else

Files:

  • internal/ingest/worker.go
  • internal/ingest/worker_test.go
internal/pipes/**/*.go

📄 CodeRabbit inference engine (AGENTS.md)

Named query pipes: NamedQuery type + NATS KV store (WAVEHOUSE_PIPES) + .sql file bootstrap. Pre-defined SQL templates with param binding + caching; per-pipe allowed_roles is the only execute-path gate via policy.RoleAllowed

Files:

  • internal/pipes/pipes.go
  • internal/pipes/dummybind_test.go
**/*_test.go

📄 CodeRabbit inference engine (AGENTS.md)

**/*_test.go: Use table-driven tests with tests := []struct{ name string; ... } and t.Run(tt.name, ...)
Use shared mocks from internal/testutil/ (MockPublisher, MockCache, MockDeduplicator, MockSubscriber) instead of creating ad-hoc mocks
Use testutil.MakeJWT(t, claims) and testutil.MakeExpiredJWT(t, claims) for auth tests
Use testutil.NewTestSchemaRegistry(tables) or discovery.NewSchemaRegistryFromMap(tables) for schema-aware tests
Use policy.NewMemoryStore(p) for in-memory policy testing without NATS
Use pipes.NewMemoryStore(queries...) for in-memory pipes testing without NATS
Use testutil.AssertJSONResponse(t, rec, status, expected) and testutil.AssertJSONContains(t, rec, status, substring) for response assertions

Files:

  • internal/pipes/dummybind_test.go
  • tests/integration/pipe_deps_test.go
  • internal/ingest/worker_test.go
  • internal/api/pipe_deps_test.go
  • internal/cache/local_test.go
  • internal/cache/version_manager_test.go
tests/integration/**/*_test.go

📄 CodeRabbit inference engine (AGENTS.md)

Go integration tests (//go:build integration; ClickHouse testcontainer); run via make test-integration

Files:

  • tests/integration/pipe_deps_test.go
internal/cache/**/*.go

📄 CodeRabbit inference engine (AGENTS.md)

Implement Cache interface → LocalCache (Ristretto) + SharedCache (TBD) + TieredCache (singleflight)

Files:

  • internal/cache/local.go
  • internal/cache/local_test.go
  • internal/cache/cache.go
  • internal/cache/version_manager.go
  • internal/cache/version_manager_test.go
🧠 Learnings (5)
📚 Learning: 2026-06-10T15:01:09.027Z
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse PR: 312
File: docs/src/content/docs/development.md:0-0
Timestamp: 2026-06-10T15:01:09.027Z
Learning: In this repo’s Markdown review (all .md files), do not flag capitalization/style issues for literal paths starting with ".github/" (or any substring that is a path beginning with ".github/"). Treat ".github" as the correct lowercase dotfile directory name, even when it appears inside prose or code spans; automated checks such as LanguageTool’s "(GITHUB)" rule commonly produce false positives for this literal filesystem path.

Applied to files:

  • CHANGELOG.md
📚 Learning: 2026-05-25T11:24:21.130Z
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse PR: 180
File: internal/cache/local.go:0-0
Timestamp: 2026-05-25T11:24:21.130Z
Learning: In WaveHouse’s cache packages (e.g., internal/cache/local.go), it’s acceptable to define package-level `var` constants that hold immutable OpenTelemetry metric attribute sets / `metric.MeasurementOption` values (for example: `cacheL1Attrs = metric.WithAttributes(attribute.String("tier","L1"))`). Treat these as stateless, pre-allocated option values (analogous to `regexp.MustCompile(...)`), not mutable global state. When applying the AGENTS.md “no global state / constructor injection” guideline, apply it to application dependencies (e.g., Cache, Publisher, Deduplicator) rather than to these immutable OTel attribute/measurement option variables—do not flag them as constructor-injection violations.

Applied to files:

  • internal/cache/local.go
  • internal/cache/local_test.go
  • internal/cache/cache.go
  • internal/cache/version_manager.go
  • internal/cache/version_manager_test.go
📚 Learning: 2026-05-20T01:02:00.784Z
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse PR: 164
File: internal/api/router_test.go:289-350
Timestamp: 2026-05-20T01:02:00.784Z
Learning: In WaveHouse’s internal API tests (files matching internal/api/**/*_test.go), follow the existing separation-of-concerns convention for testing the RequireRole middleware: inject `ContextKeyRole` directly into the request `context.Context` instead of using `testutil.MakeJWT`/JWT-driven flows. Do not refactor role-gate tests to use JWT tokens—JWT parsing and token handling are covered separately in `middleware_test.go` (the dedicated JWT parsing tests), and mixing those concerns would expand the failure surface and reduce isolation.

Applied to files:

  • internal/api/pipe_deps_test.go
📚 Learning: 2026-05-23T01:23:59.268Z
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse PR: 174
File: internal/api/ingest_test.go:111-111
Timestamp: 2026-05-23T01:23:59.268Z
Learning: In WaveHouse Go tests in internal/api/**/*_test.go, use internal/testutil.AssertJSONErrorResponse(t, w) for HTTP error-path JSON assertions. Do not use (or reintroduce) package-local assertJSONErrorResponse helpers. AssertJSONErrorResponse verifies the response Content-Type is application/json, includes the X-Content-Type-Options: nosniff header, and that the JSON body contains an "error" field.

Applied to files:

  • internal/api/pipe_deps_test.go
📚 Learning: 2026-05-20T20:30:15.808Z
Learnt from: taitelee
Repo: Wave-RF/WaveHouse PR: 172
File: internal/api/pipes_test.go:106-118
Timestamp: 2026-05-20T20:30:15.808Z
Learning: For WaveHouse pipes authorization allowlist checks, fix the empty-role fail-open behavior by (1) removing any outer guard that prevents allowlist evaluation when the incoming `role` is `""` (e.g., don’t short-circuit with `if role != "" { ... }`), and (2) during allowlist scanning, ensure only non-empty allowlist entries can match—e.g., require `ar != "" && ar == role` (so a malformed allowlist like `["" ]` cannot grant access to an empty incoming role via `"" == ""`).

Applied to files:

  • internal/api/pipes.go
🔇 Additional comments (30)
internal/cache/cache.go (1)

8-32: LGTM!

internal/cache/version_manager.go (3)

16-32: LGTM!


34-51: LGTM!


76-94: LGTM!

internal/cache/version_manager_test.go (1)

9-72: LGTM!

internal/cache/local.go (1)

31-75: LGTM!

internal/cache/local_test.go (1)

12-161: LGTM!

internal/api/structured_query.go (2)

129-143: LGTM!


172-174: LGTM!

internal/testutil/mocks.go (1)

117-136: LGTM!

internal/ingest/worker_test.go (4)

224-232: LGTM!


449-543: LGTM!


704-731: LGTM!


733-764: LGTM!

internal/ingest/worker.go (1)

461-509: LGTM!

internal/pipes/pipes.go (2)

29-37: LGTM!


204-243: LGTM!

internal/api/pipe_deps.go (7)

28-35: LGTM!


43-56: LGTM!


64-83: LGTM!


90-103: LGTM!


134-153: LGTM!


196-205: LGTM!


179-187: ⚠️ Potential issue | 🟠 Major | ⚡ Quick win

Incorrect escape sequence handling for identifiers with both backslashes and backticks.

Lines 183-184 unescape ClickHouse backtick-quoted identifiers in the wrong order, causing incorrect results when an identifier contains both a literal backslash and a literal backtick. ClickHouse uses \\ for a backslash and ``` for a backtick inside backtick-quoted identifiers.

Example: The ClickHouse identifier a\b(a, backslash, backtick, b) is written in SQL as ``a\`b`` (where\` and \`` → `` ``). After stripping the outer backticks, the string is a\\\b`. The current code:

  1. Line 183 replaces \`` with `` ``: a\\\b` → `a\b`
  2. Line 184 tries to replace \\\\ with \\: no match (only two backslashes remain)
  3. Result: a\\b (wrong; expected a\b`)

Sequential ReplaceAll is unsafe here because the replacements can interact. The correct approach is to process escape sequences in a single pass:

  • Iterate byte-by-byte; when \ is seen, check the next byte:
    • If \, output one \ and advance two bytes
    • If `, output one ` and advance two bytes
    • Otherwise, output \ (or handle as an error)
🔧 Proposed fix for correct escape handling
 func unquoteIdent(s string) string {
 	s = strings.TrimSpace(s)
 	if len(s) >= 2 && s[0] == '`' && s[len(s)-1] == '`' {
-		s = s[1 : len(s)-1]
-		s = strings.ReplaceAll(s, "\\`", "`")
-		s = strings.ReplaceAll(s, "\\\\", "\\")
+		s = s[1 : len(s)-1] // strip outer backticks
+		// Unescape in a single pass to avoid interaction between \\ and \`
+		var out strings.Builder
+		for i := 0; i < len(s); i++ {
+			if s[i] == '\\' && i+1 < len(s) {
+				next := s[i+1]
+				if next == '\\' || next == '`' {
+					out.WriteByte(next)
+					i++ // skip the next byte (already consumed)
+					continue
+				}
+			}
+			out.WriteByte(s[i])
+		}
+		s = out.String()
 	}
 	return s
 }
			> Likely an incorrect or invalid review comment.
internal/api/pipes.go (3)

27-33: LGTM!


76-80: LGTM!


155-193: LGTM!

cmd/wavehouse/main.go (1)

367-372: LGTM!

internal/pipes/dummybind_test.go (1)

1-63: LGTM!

internal/api/pipe_deps_test.go (1)

1-192: LGTM!

Comment thread CHANGELOG.md Outdated
Comment thread internal/api/pipe_deps.go Outdated
Comment thread internal/cache/version_manager.go
Comment thread tests/integration/pipe_deps_test.go Outdated
@github-project-automation github-project-automation Bot moved this from Backlog to In review in WaveHouse Task Board Jun 11, 2026
coderabbitai[bot]
coderabbitai Bot previously approved these changes Jun 14, 2026
coderabbitai[bot]
coderabbitai Bot previously approved these changes Jun 14, 2026
coderabbitai[bot]
coderabbitai Bot previously approved these changes Jun 29, 2026
@EricAndrechek

Copy link
Copy Markdown
Member

Docs note (posting as a normal comment — api.md is unchanged on this branch, so there's no diff line to attach a review comment to).

docs/src/content/docs/api.md:467 — the GET/POST /v1/pipes/{name} "Execute Named Pipe" section documents the L1/singleflight caching and the X-Cache header, but not that a pipe's cached result is now invalidated by writes to the tables it reads — the headline feature of this PR. It's covered thoroughly in pipes.mdx, but the endpoint reference is the first place a reader looks for "when does my cached pipe go stale," and it's silent there. Since api.md isn't touched by this branch, it's really a small follow-up docs task: add one line cross-referencing the pipes guide, e.g. "Cached results are invalidated by writes to the tables the pipe reads — see Named Pipes for the resolution/invalidation semantics."

@EricAndrechek EricAndrechek left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Overall much better than the last pass, but I think it still has some room for improvement. A few nit-picky things are included, some docs tweaks, etc, but also some efficiency concerns and the main bit being on caching and pipe dependencies. Specifically, I think we need to consider a method for pipe/table dependencies like we have for the namespace dependencies and versioning using a global table version for easy table-wide invalidation for pipe dependencies/tables and their schemas, potentially on the whole database or something if that makes sense.

Comment thread cmd/wavehouse/main.go
Comment thread cmd/wavehouse/main.go Outdated
Comment thread internal/discovery/discovery.go
Comment thread docs/src/content/docs/architecture.md Outdated
Comment thread internal/api/pipes.go Outdated
Comment thread docs/src/content/docs/architecture.md
Comment thread internal/api/pipes.go Outdated
Comment thread internal/api/pipes.go Outdated
Comment thread internal/cache/version_manager.go Outdated
Comment thread internal/discovery/explain.go
taitelee and others added 3 commits July 10, 2026 17:28
…base-version fallback; TTL-cap external/pruned reads; serialize schema refreshes

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
# Conflicts:
#	AGENTS.md
#	CHANGELOG.md
#	internal/api/stream.go
coderabbitai[bot]
coderabbitai Bot previously approved these changes Jul 16, 2026
@github-actions github-actions Bot added the area/sdk TypeScript SDK (clients/ts/) label Jul 16, 2026
@taitelee
taitelee requested a review from EricAndrechek July 16, 2026 20:03
@taitelee taitelee moved this from In review to Backlog in WaveHouse Task Board Jul 16, 2026
@taitelee taitelee moved this from Backlog to Ready in WaveHouse Task Board Jul 16, 2026

@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: 10

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (3)
internal/api/ingest.go (1)

58-61: 📐 Maintainability & Code Quality | 🟠 Major | ⚡ Quick win

Inject the deduplication counter instead of creating global state.

dedupeMissingIDCounter is package-global and binds to the global OpenTelemetry meter during package initialization. Add the counter to the metrics/observer wiring and pass it into NewIngestHandler as an *IngestHandler dependency. This avoids uncontrolled global OpenTelemetry instrumentation and keeps the handler dependencies explicit.

Source: Coding guidelines

AGENTS.md (1)

29-31: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value

Extend the chsql/ gloss to cover the NATS encoder.

This PR moves SafeEncodeNATS/SafeDecodeNATS into internal/chsql/nats.go, and docs/src/content/docs/architecture.md now documents them. The chsql/ bullet in this package list and the internal/chsql/ line in the File Structure block still describe only QuoteIdent and BindUnsafe. Add the NATS-safe name encoder (and QuoteString) so the two documents agree.

As per coding guidelines: "Every code change must update its corresponding documentation and CHANGELOG.md in the same PR."

Source: Coding guidelines

CHANGELOG.md (1)

77-77: 📐 Maintainability & Code Quality | 🟡 Minor | ⚡ Quick win

Update the trailing sentence: it contradicts the new entry in the same release.

The last sentence of this bullet states that pipes pass no dependencies and rely on TTL-only invalidation until #178 lands. The Added entry at line 14 in the same Unreleased section closes #178 and documents per-table versioned namespaces for pipes. A reader of one release section gets two opposite statements.

♻️ Proposed change
-Pipes currently pass no dependencies (TTL-only invalidation) until they can report the tables they read ([`#178`](https://github.com/Wave-RF/WaveHouse/issues/178)).
+Pipes gained their table dependencies in the same release — see the dependency-aware pipe invalidation entry under `Added` ([`#178`](https://github.com/Wave-RF/WaveHouse/issues/178)).

ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro Plus

Run ID: 7c338981-aac0-4953-b5c4-d56094681546

📥 Commits

Reviewing files that changed from the base of the PR and between 394420d and 4d856a5.

📒 Files selected for processing (36)
  • AGENTS.md
  • CHANGELOG.md
  • cmd/wavehouse/main.go
  • docs/src/content/docs/api.md
  • docs/src/content/docs/architecture.md
  • docs/src/content/docs/pipes.mdx
  • internal/api/boot_chain_test.go
  • internal/api/dlq.go
  • internal/api/errors_test.go
  • internal/api/ingest.go
  • internal/api/pipe_deps_test.go
  • internal/api/pipes.go
  • internal/api/pipes_test.go
  • internal/api/stream.go
  • internal/api/structured_query.go
  • internal/cache/cache.go
  • internal/cache/local.go
  • internal/cache/local_test.go
  • internal/cache/version_manager.go
  • internal/chsql/chsql.go
  • internal/chsql/nats.go
  • internal/chsql/nats_test.go
  • internal/chsql/quote_test.go
  • internal/discovery/deps_test.go
  • internal/discovery/discovery.go
  • internal/discovery/discovery_test.go
  • internal/discovery/explain.go
  • internal/discovery/explain_test.go
  • internal/discovery/fuzz_deps_test.go
  • internal/ingest/worker.go
  • internal/ingest/worker_test.go
  • internal/pipes/pipes.go
  • tests/e2e/sdk/cache.test.ts
  • tests/integration/dlq_test.go
  • tests/integration/pipe_cache_test.go
  • tests/integration/setup_test.go
📜 Review details
🧰 Additional context used
📓 Path-based instructions (10)
**/*

📄 CodeRabbit inference engine (AGENTS.md)

**/*: Run make ci locally before every push, using the documented background execution method and a Docker daemon.
Every code change must update its corresponding documentation and CHANGELOG.md in the same PR.
Address every review finding with a substantive reply, a fix or tracking issue, and resolved review threads before merge.
Create agent PRs as drafts with Conventional-Commits titles of at most 72 characters; validate titles with scripts/lint-pr-title.sh.
Never force-push or rebase PR branches; merge origin/main instead.
Do not hand-write review markers or bypass hooks with --no-verify; use the prescribed tooling.

Files:

  • internal/chsql/quote_test.go
  • docs/src/content/docs/api.md
  • internal/api/boot_chain_test.go
  • AGENTS.md
  • internal/api/dlq.go
  • internal/discovery/explain_test.go
  • internal/api/ingest.go
  • internal/api/errors_test.go
  • docs/src/content/docs/architecture.md
  • tests/integration/setup_test.go
  • tests/e2e/sdk/cache.test.ts
  • internal/pipes/pipes.go
  • CHANGELOG.md
  • docs/src/content/docs/pipes.mdx
  • internal/api/stream.go
  • internal/api/pipes_test.go
  • tests/integration/dlq_test.go
  • internal/chsql/nats_test.go
  • internal/discovery/discovery_test.go
  • internal/api/pipe_deps_test.go
  • internal/chsql/nats.go
  • internal/ingest/worker.go
  • tests/integration/pipe_cache_test.go
  • internal/discovery/deps_test.go
  • internal/discovery/explain.go
  • internal/cache/local_test.go
  • cmd/wavehouse/main.go
  • internal/cache/version_manager.go
  • internal/discovery/fuzz_deps_test.go
  • internal/api/structured_query.go
  • internal/cache/local.go
  • internal/ingest/worker_test.go
  • internal/cache/cache.go
  • internal/api/pipes.go
  • internal/chsql/chsql.go
  • internal/discovery/discovery.go
**/*.go

📄 CodeRabbit inference engine (AGENTS.md)

**/*.go: Use Go 1.26 with strict gofumpt formatting.
Return errors instead of panicking, and wrap errors with fmt.Errorf("context: %w", err).
Use explicit dependency passing and constructor injection; do not introduce global state.
Use structured logging with log/slog and JSON handlers.

Files:

  • internal/chsql/quote_test.go
  • internal/api/boot_chain_test.go
  • internal/api/dlq.go
  • internal/discovery/explain_test.go
  • internal/api/ingest.go
  • internal/api/errors_test.go
  • tests/integration/setup_test.go
  • internal/pipes/pipes.go
  • internal/api/stream.go
  • internal/api/pipes_test.go
  • tests/integration/dlq_test.go
  • internal/chsql/nats_test.go
  • internal/discovery/discovery_test.go
  • internal/api/pipe_deps_test.go
  • internal/chsql/nats.go
  • internal/ingest/worker.go
  • tests/integration/pipe_cache_test.go
  • internal/discovery/deps_test.go
  • internal/discovery/explain.go
  • internal/cache/local_test.go
  • cmd/wavehouse/main.go
  • internal/cache/version_manager.go
  • internal/discovery/fuzz_deps_test.go
  • internal/api/structured_query.go
  • internal/cache/local.go
  • internal/ingest/worker_test.go
  • internal/cache/cache.go
  • internal/api/pipes.go
  • internal/chsql/chsql.go
  • internal/discovery/discovery.go
internal/**/*.go

📄 CodeRabbit inference engine (AGENTS.md)

internal/**/*.go: Keep internal packages lowercase and single-word or abbreviated; use internal/ to enforce module privacy.
Core interchangeable behaviors should be represented by interfaces, including Cache, Deduplicator, Publisher, and Subscriber.

Files:

  • internal/chsql/quote_test.go
  • internal/api/boot_chain_test.go
  • internal/api/dlq.go
  • internal/discovery/explain_test.go
  • internal/api/ingest.go
  • internal/api/errors_test.go
  • internal/pipes/pipes.go
  • internal/api/stream.go
  • internal/api/pipes_test.go
  • internal/chsql/nats_test.go
  • internal/discovery/discovery_test.go
  • internal/api/pipe_deps_test.go
  • internal/chsql/nats.go
  • internal/ingest/worker.go
  • internal/discovery/deps_test.go
  • internal/discovery/explain.go
  • internal/cache/local_test.go
  • internal/cache/version_manager.go
  • internal/discovery/fuzz_deps_test.go
  • internal/api/structured_query.go
  • internal/cache/local.go
  • internal/ingest/worker_test.go
  • internal/cache/cache.go
  • internal/api/pipes.go
  • internal/chsql/chsql.go
  • internal/discovery/discovery.go
**/*_test.go

📄 CodeRabbit inference engine (AGENTS.md)

**/*_test.go: Use table-driven tests with t.Run(tt.name, ...) for multiple scenarios, and add tests for every new function.
Reuse shared test helpers and mocks from internal/testutil/, including JWT, schema, policy, pipe, response, and mock helpers.

Files:

  • internal/chsql/quote_test.go
  • internal/api/boot_chain_test.go
  • internal/discovery/explain_test.go
  • internal/api/errors_test.go
  • tests/integration/setup_test.go
  • internal/api/pipes_test.go
  • tests/integration/dlq_test.go
  • internal/chsql/nats_test.go
  • internal/discovery/discovery_test.go
  • internal/api/pipe_deps_test.go
  • tests/integration/pipe_cache_test.go
  • internal/discovery/deps_test.go
  • internal/cache/local_test.go
  • internal/discovery/fuzz_deps_test.go
  • internal/ingest/worker_test.go
**/*.md

📄 CodeRabbit inference engine (AGENTS.md)

Keep CLAUDE.md as a short pointer to AGENTS.md and do not duplicate the full agent instructions.

Files:

  • docs/src/content/docs/api.md
  • AGENTS.md
  • docs/src/content/docs/architecture.md
  • CHANGELOG.md
internal/api/**/*.go

📄 CodeRabbit inference engine (AGENTS.md)

internal/api/**/*.go: Use Chi v5 for HTTP routing, and ensure all /v1/* routes run the always-on JWT middleware.
Maintain bearer-token-only CORS: never emit Access-Control-Allow-Credentials, and do not reintroduce cookie or session authentication.
Enforce column-level access control on every read path, including structured queries and live streams, through the shared policy decision function.

Files:

  • internal/api/boot_chain_test.go
  • internal/api/dlq.go
  • internal/api/ingest.go
  • internal/api/errors_test.go
  • internal/api/stream.go
  • internal/api/pipes_test.go
  • internal/api/pipe_deps_test.go
  • internal/api/structured_query.go
  • internal/api/pipes.go
tests/e2e/sdk/*.test.ts

📄 CodeRabbit inference engine (AGENTS.md)

tests/e2e/sdk/*.test.ts: E2E tests must use the TypeScript SDK and suite-specific tables obtained from suiteTables; never use bare shared table names.
Keep E2E tests sequential with maxWorkers: 1 because they share mutable global policy state; policy-mutating tests must snapshot and restore the full policy.

Files:

  • tests/e2e/sdk/cache.test.ts
internal/pipes/**/*.go

📄 CodeRabbit inference engine (AGENTS.md)

Named pipes must enforce exact allowed_roles membership, with no wildcard authorization and admin-only behavior when roles are absent; pipe SQL runs as-is and unresolved dependencies must fall back safely to database-wide invalidation or a TTL cap.

Files:

  • internal/pipes/pipes.go
docs/src/content/docs/**/*.mdx

📄 CodeRabbit inference engine (AGENTS.md)

Author Mermaid diagrams vertically by default, avoid large side-by-side diagrams, and keep labels short and readable.

Files:

  • docs/src/content/docs/pipes.mdx
internal/ingest/**/*.go

📄 CodeRabbit inference engine (AGENTS.md)

Ingest payloads must be validated against discovered ClickHouse schemas, and failed batch inserts must publish to the enabled DLQ without silent data loss.

Files:

  • internal/ingest/worker.go
  • internal/ingest/worker_test.go
🧠 Learnings (8)
📚 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/chsql/quote_test.go
  • internal/api/boot_chain_test.go
  • internal/discovery/explain_test.go
  • internal/api/errors_test.go
  • tests/integration/setup_test.go
  • internal/api/pipes_test.go
  • tests/integration/dlq_test.go
  • internal/chsql/nats_test.go
  • internal/discovery/discovery_test.go
  • internal/api/pipe_deps_test.go
  • tests/integration/pipe_cache_test.go
  • internal/discovery/deps_test.go
  • internal/cache/local_test.go
  • internal/discovery/fuzz_deps_test.go
  • internal/ingest/worker_test.go
📚 Learning: 2026-07-07T12:38:12.052Z
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse PR: 378
File: internal/auth/auth.go:119-132
Timestamp: 2026-07-07T12:38:12.052Z
Learning: In this repo, do not add or recommend logging/tracing client IP addresses using naive or untrusted sources (e.g., `r.RemoteAddr` or directly trusting/deriving `X-Forwarded-For`) anywhere in the Go codebase. `middleware.RealIP` was removed due to IP-spoofing risks, and proper trusted-proxy-aware client-IP handling is intentionally deferred to issue `#333`. During code review, if proposed changes would record client IPs (including in audit paths such as `internal/auth/auth.go`), reject/redirect until `#333` lands with correct trusted-proxy configuration and safeguards.

Applied to files:

  • internal/chsql/quote_test.go
  • internal/api/boot_chain_test.go
  • internal/api/dlq.go
  • internal/discovery/explain_test.go
  • internal/api/ingest.go
  • internal/api/errors_test.go
  • tests/integration/setup_test.go
  • internal/pipes/pipes.go
  • internal/api/stream.go
  • internal/api/pipes_test.go
  • tests/integration/dlq_test.go
  • internal/chsql/nats_test.go
  • internal/discovery/discovery_test.go
  • internal/api/pipe_deps_test.go
  • internal/chsql/nats.go
  • internal/ingest/worker.go
  • tests/integration/pipe_cache_test.go
  • internal/discovery/deps_test.go
  • internal/discovery/explain.go
  • internal/cache/local_test.go
  • cmd/wavehouse/main.go
  • internal/cache/version_manager.go
  • internal/discovery/fuzz_deps_test.go
  • internal/api/structured_query.go
  • internal/cache/local.go
  • internal/ingest/worker_test.go
  • internal/cache/cache.go
  • internal/api/pipes.go
  • internal/chsql/chsql.go
  • internal/discovery/discovery.go
📚 Learning: 2026-06-10T15:01:09.027Z
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse PR: 312
File: docs/src/content/docs/development.md:0-0
Timestamp: 2026-06-10T15:01:09.027Z
Learning: In this repo’s Markdown review (all .md files), do not flag capitalization/style issues for literal paths starting with ".github/" (or any substring that is a path beginning with ".github/"). Treat ".github" as the correct lowercase dotfile directory name, even when it appears inside prose or code spans; automated checks such as LanguageTool’s "(GITHUB)" rule commonly produce false positives for this literal filesystem path.

Applied to files:

  • docs/src/content/docs/api.md
  • AGENTS.md
  • docs/src/content/docs/architecture.md
  • CHANGELOG.md
📚 Learning: 2026-05-20T01:02:00.784Z
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse PR: 164
File: internal/api/router_test.go:289-350
Timestamp: 2026-05-20T01:02:00.784Z
Learning: In WaveHouse’s internal API tests (files matching internal/api/**/*_test.go), follow the existing separation-of-concerns convention for testing the RequireRole middleware: inject `ContextKeyRole` directly into the request `context.Context` instead of using `testutil.MakeJWT`/JWT-driven flows. Do not refactor role-gate tests to use JWT tokens—JWT parsing and token handling are covered separately in `middleware_test.go` (the dedicated JWT parsing tests), and mixing those concerns would expand the failure surface and reduce isolation.

Applied to files:

  • internal/api/boot_chain_test.go
  • internal/api/errors_test.go
  • internal/api/pipes_test.go
  • internal/api/pipe_deps_test.go
📚 Learning: 2026-05-23T01:23:59.268Z
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse PR: 174
File: internal/api/ingest_test.go:111-111
Timestamp: 2026-05-23T01:23:59.268Z
Learning: In WaveHouse Go tests in internal/api/**/*_test.go, use internal/testutil.AssertJSONErrorResponse(t, w) for HTTP error-path JSON assertions. Do not use (or reintroduce) package-local assertJSONErrorResponse helpers. AssertJSONErrorResponse verifies the response Content-Type is application/json, includes the X-Content-Type-Options: nosniff header, and that the JSON body contains an "error" field.

Applied to files:

  • internal/api/boot_chain_test.go
  • internal/api/errors_test.go
  • internal/api/pipes_test.go
  • internal/api/pipe_deps_test.go
📚 Learning: 2026-07-09T16:20:50.620Z
Learnt from: taitelee
Repo: Wave-RF/WaveHouse PR: 402
File: internal/api/ingest_test.go:1434-1448
Timestamp: 2026-07-09T16:20:50.620Z
Learning: In WaveHouse’s Go schema/timestamp handling (internal/discovery), resolve any timestamp column zone/precision requirements exactly once during SchemaRegistry.Refresh() and cache the resolved specs on the column. Do not resolve zones per-record on the hot path. Ensure timezone resolution uses embedded tzdata (e.g., via cmd/wavehouse’s time.LoadLocation + embedded tzdata) so minimal containers don’t silently fall back to UTC for non-UTC ClickHouse servers, which would corrupt stored instants. Treat timezone-resolution failures during refresh as non-fatal: keep boot alive, RetryRefresh should continue retrying, and readiness/health (/livez) should surface a degraded state. If a column’s timezone cannot be resolved after refresh, emit a warning and reject only that column’s timestamp values per-record rather than failing the entire refresh.

Applied to files:

  • internal/discovery/explain_test.go
  • internal/discovery/discovery_test.go
  • internal/discovery/deps_test.go
  • internal/discovery/explain.go
  • internal/discovery/fuzz_deps_test.go
  • internal/discovery/discovery.go
📚 Learning: 2026-05-20T20:30:15.808Z
Learnt from: taitelee
Repo: Wave-RF/WaveHouse PR: 172
File: internal/api/pipes_test.go:106-118
Timestamp: 2026-05-20T20:30:15.808Z
Learning: For WaveHouse pipes authorization allowlist checks, fix the empty-role fail-open behavior by (1) removing any outer guard that prevents allowlist evaluation when the incoming `role` is `""` (e.g., don’t short-circuit with `if role != "" { ... }`), and (2) during allowlist scanning, ensure only non-empty allowlist entries can match—e.g., require `ar != "" && ar == role` (so a malformed allowlist like `["" ]` cannot grant access to an empty incoming role via `"" == ""`).

Applied to files:

  • internal/api/pipes_test.go
  • internal/api/pipes.go
📚 Learning: 2026-05-25T11:24:21.130Z
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse PR: 180
File: internal/cache/local.go:0-0
Timestamp: 2026-05-25T11:24:21.130Z
Learning: In WaveHouse’s cache packages (e.g., internal/cache/local.go), it’s acceptable to define package-level `var` constants that hold immutable OpenTelemetry metric attribute sets / `metric.MeasurementOption` values (for example: `cacheL1Attrs = metric.WithAttributes(attribute.String("tier","L1"))`). Treat these as stateless, pre-allocated option values (analogous to `regexp.MustCompile(...)`), not mutable global state. When applying the AGENTS.md “no global state / constructor injection” guideline, apply it to application dependencies (e.g., Cache, Publisher, Deduplicator) rather than to these immutable OTel attribute/measurement option variables—do not flag them as constructor-injection violations.

Applied to files:

  • internal/cache/local_test.go
  • internal/cache/version_manager.go
  • internal/cache/local.go
  • internal/cache/cache.go
🪛 ast-grep (0.45.0)
tests/integration/pipe_cache_test.go

[warning] 73-74: Detected a SQL statement built with 'fmt.Sprintf' and passed directly to 'db.Exec'/'db.ExecContext'. Interpolating values into a query string lets an attacker inject arbitrary SQL. Use parameterized queries instead: pass the SQL with placeholders ('?' or '') as the query argument and supply the values as separate arguments, e.g. 'db.Exec("UPDATE t SET x = ? WHERE id = ?", x, id)'.
Context: e.chConn.Exec(ctx,
fmt.Sprintf("INSERT INTO %s SELECT number, number%%3 FROM numbers(60)", base))
Note: [CWE-89] Improper Neutralization of Special Elements used in an SQL Command ('SQL Injection').

(sql-injection-exec-sprintf-go)


[warning] 77-78: Detected a SQL statement built with 'fmt.Sprintf' and passed directly to 'db.Exec'/'db.ExecContext'. Interpolating values into a query string lets an attacker inject arbitrary SQL. Use parameterized queries instead: pass the SQL with placeholders ('?' or '') as the query argument and supply the values as separate arguments, e.g. 'db.Exec("UPDATE t SET x = ? WHERE id = ?", x, id)'.
Context: e.chConn.Exec(ctx,
fmt.Sprintf("CREATE VIEW %s AS SELECT user_id, org_id FROM %s", view, base))
Note: [CWE-89] Improper Neutralization of Special Elements used in an SQL Command ('SQL Injection').

(sql-injection-exec-sprintf-go)


[warning] 120-121: Detected a SQL statement built with 'fmt.Sprintf' and passed directly to 'db.Exec'/'db.ExecContext'. Interpolating values into a query string lets an attacker inject arbitrary SQL. Use parameterized queries instead: pass the SQL with placeholders ('?' or '') as the query argument and supply the values as separate arguments, e.g. 'db.Exec("UPDATE t SET x = ? WHERE id = ?", x, id)'.
Context: e.chConn.Exec(ctx,
fmt.Sprintf("INSERT INTO %s SELECT number FROM numbers(60)", tbl))
Note: [CWE-89] Improper Neutralization of Special Elements used in an SQL Command ('SQL Injection').

(sql-injection-exec-sprintf-go)


[warning] 183-184: Detected a SQL statement built with 'fmt.Sprintf' and passed directly to 'db.Exec'/'db.ExecContext'. Interpolating values into a query string lets an attacker inject arbitrary SQL. Use parameterized queries instead: pass the SQL with placeholders ('?' or '') as the query argument and supply the values as separate arguments, e.g. 'db.Exec("UPDATE t SET x = ? WHERE id = ?", x, id)'.
Context: e.chConn.Exec(ctx,
fmt.Sprintf("INSERT INTO %s SELECT number, number FROM numbers(100)", src))
Note: [CWE-89] Improper Neutralization of Special Elements used in an SQL Command ('SQL Injection').

(sql-injection-exec-sprintf-go)


[warning] 187-189: Detected a SQL statement built with 'fmt.Sprintf' and passed directly to 'db.Exec'/'db.ExecContext'. Interpolating values into a query string lets an attacker inject arbitrary SQL. Use parameterized queries instead: pass the SQL with placeholders ('?' or '') as the query argument and supply the values as separate arguments, e.g. 'db.Exec("UPDATE t SET x = ? WHERE id = ?", x, id)'.
Context: e.chConn.Exec(ctx, fmt.Sprintf(
"CREATE MATERIALIZED VIEW %s ENGINE=AggregatingMergeTree ORDER BY id AS SELECT id, sumState(v) AS s FROM %s GROUP BY id",
mv, src))
Note: [CWE-89] Improper Neutralization of Special Elements used in an SQL Command ('SQL Injection').

(sql-injection-exec-sprintf-go)


[warning] 229-231: Detected a SQL statement built with 'fmt.Sprintf' and passed directly to 'db.Exec'/'db.ExecContext'. Interpolating values into a query string lets an attacker inject arbitrary SQL. Use parameterized queries instead: pass the SQL with placeholders ('?' or '') as the query argument and supply the values as separate arguments, e.g. 'db.Exec("UPDATE t SET x = ? WHERE id = ?", x, id)'.
Context: e.chConn.Exec(ctx, fmt.Sprintf(
"CREATE MATERIALIZED VIEW %s TO %s AS SELECT id FROM %s",
mv, chQuoteIdent(target), src))
Note: [CWE-89] Improper Neutralization of Special Elements used in an SQL Command ('SQL Injection').

(sql-injection-exec-sprintf-go)


[error] 79-79: SQL query is built by concatenating a string literal with a variable and passed to a database/sql call (Query, Exec, QueryRow, Prepare, or their Context variants). String concatenation lets attacker-controlled input alter the query structure, enabling SQL injection. Use parameterized queries with placeholders ('?' or '') and pass the values as separate arguments instead of concatenating them into the query string.
Context: e.chConn.Exec(context.Background(), "DROP VIEW IF EXISTS "+view+"")
Note: [CWE-89] Improper Neutralization of Special Elements used in an SQL Command ('SQL Injection').

(sql-injection-query-string-concat-go)


[error] 190-190: SQL query is built by concatenating a string literal with a variable and passed to a database/sql call (Query, Exec, QueryRow, Prepare, or their Context variants). String concatenation lets attacker-controlled input alter the query structure, enabling SQL injection. Use parameterized queries with placeholders ('?' or '') and pass the values as separate arguments instead of concatenating them into the query string.
Context: e.chConn.Exec(context.Background(), "DROP VIEW IF EXISTS "+mv+"")
Note: [CWE-89] Improper Neutralization of Special Elements used in an SQL Command ('SQL Injection').

(sql-injection-query-string-concat-go)


[error] 232-232: SQL query is built by concatenating a string literal with a variable and passed to a database/sql call (Query, Exec, QueryRow, Prepare, or their Context variants). String concatenation lets attacker-controlled input alter the query structure, enabling SQL injection. Use parameterized queries with placeholders ('?' or '') and pass the values as separate arguments instead of concatenating them into the query string.
Context: e.chConn.Exec(context.Background(), "DROP VIEW IF EXISTS "+mv+"")
Note: [CWE-89] Improper Neutralization of Special Elements used in an SQL Command ('SQL Injection').

(sql-injection-query-string-concat-go)

internal/discovery/explain.go

[error] 70-70: SQL query is built by concatenating a string literal with a variable and passed to a database/sql call (Query, Exec, QueryRow, Prepare, or their Context variants). String concatenation lets attacker-controlled input alter the query structure, enabling SQL injection. Use parameterized queries with placeholders ('?' or '') and pass the values as separate arguments instead of concatenating them into the query string.
Context: conn.Query(ctx, "EXPLAIN QUERY TREE "+sql)
Note: [CWE-89] Improper Neutralization of Special Elements used in an SQL Command ('SQL Injection').

(sql-injection-query-string-concat-go)

internal/discovery/fuzz_deps_test.go

[error] 60-60: SQL query is built by concatenating a string literal with a variable and passed to a database/sql call (Query, Exec, QueryRow, Prepare, or their Context variants). String concatenation lets attacker-controlled input alter the query structure, enabling SQL injection. Use parameterized queries with placeholders ('?' or '') and pass the values as separate arguments instead of concatenating them into the query string.
Context: boot.Exec(ctx, "CREATE DATABASE IF NOT EXISTS "+fdb)
Note: [CWE-89] Improper Neutralization of Special Elements used in an SQL Command ('SQL Injection').

(sql-injection-query-string-concat-go)


[error] 61-61: SQL query is built by concatenating a string literal with a variable and passed to a database/sql call (Query, Exec, QueryRow, Prepare, or their Context variants). String concatenation lets attacker-controlled input alter the query structure, enabling SQL injection. Use parameterized queries with placeholders ('?' or '') and pass the values as separate arguments instead of concatenating them into the query string.
Context: boot.Exec(ctx, "CREATE DATABASE IF NOT EXISTS "+fother)
Note: [CWE-89] Improper Neutralization of Special Elements used in an SQL Command ('SQL Injection').

(sql-injection-query-string-concat-go)


[error] 105-105: SQL query is built by concatenating a string literal with a variable and passed to a database/sql call (Query, Exec, QueryRow, Prepare, or their Context variants). String concatenation lets attacker-controlled input alter the query structure, enabling SQL injection. Use parameterized queries with placeholders ('?' or '') and pass the values as separate arguments instead of concatenating them into the query string.
Context: conn.Exec(ctx, "DROP DATABASE IF EXISTS "+fdb)
Note: [CWE-89] Improper Neutralization of Special Elements used in an SQL Command ('SQL Injection').

(sql-injection-query-string-concat-go)


[error] 106-106: SQL query is built by concatenating a string literal with a variable and passed to a database/sql call (Query, Exec, QueryRow, Prepare, or their Context variants). String concatenation lets attacker-controlled input alter the query structure, enabling SQL injection. Use parameterized queries with placeholders ('?' or '') and pass the values as separate arguments instead of concatenating them into the query string.
Context: conn.Exec(ctx, "DROP DATABASE IF EXISTS "+fother)
Note: [CWE-89] Improper Neutralization of Special Elements used in an SQL Command ('SQL Injection').

(sql-injection-query-string-concat-go)


[error] 249-249: SQL query is built by concatenating a string literal with a variable and passed to a database/sql call (Query, Exec, QueryRow, Prepare, or their Context variants). String concatenation lets attacker-controlled input alter the query structure, enabling SQL injection. Use parameterized queries with placeholders ('?' or '') and pass the values as separate arguments instead of concatenating them into the query string.
Context: conn.Query(ctx, "EXPLAIN QUERY TREE "+sql)
Note: [CWE-89] Improper Neutralization of Special Elements used in an SQL Command ('SQL Injection').

(sql-injection-query-string-concat-go)

🪛 LanguageTool
docs/src/content/docs/api.md

[style] ~480-~480: Redundant conjunctions can lead to confusion; consider removing a conjunction here.
Context: ...meters can be supplied via query string and/or JSON body. Results are cached in the sh...

(AND_OR)

docs/src/content/docs/architecture.md

[style] ~73-~73: Since ownership is already implied, this phrasing may be redundant.
Context: ...e JWT/JWKS authentication middleware is its own package, [auth/](#auth--authenticatio...

(PRP_OWN)


[style] ~157-~157: Since ownership is already implied, this phrasing may be redundant.
Context: ...ared by query/ and policy/, kept in their own package to break an import cycle. `Quot...

(PRP_OWN)

docs/src/content/docs/pipes.mdx

[typographical] ~178-~178: The word ‘WHERE’ starts a question. Add a question mark (“?”) at the end of the sentence.
Context: ...its data, so this is precise, not stale). The source=mobile result is resolved ...

(WRB_QUESTION_MARK)


[style] ~178-~178: Try elevating your writing by using a synonym here.
Context: ...solved and cached separately, depending only on mobile_events. The pruning is deli...

(ONLY_SOLELY)


[style] ~184-~184: Since ownership is already implied, this phrasing may be redundant.
Context: ...e detects the read (a table function is its own node in the query tree) and **caps the ...

(PRP_OWN)

🔇 Additional comments (33)
internal/api/stream.go (1)

24-25: 🩺 Stability & Availability

No change needed. Production wiring assigns streamHandler.Metrics = sseMetrics after construction, and production metrics instrumentation is nil-safe.

internal/api/pipes.go (6)

26-63: LGTM!


74-88: LGTM!

Also applies to: 98-114


251-266: LGTM!


307-343: LGTM!

Also applies to: 355-403


405-448: LGTM!

Also applies to: 450-461


247-249: 🩺 Stability & Availability

No change needed. pipes.BindParams inlines supplied {{name}} and {{name:default}} placeholders as literals before dependency resolution runs, so resolveDeps passes a fully bound query to EXPLAIN QUERY TREE.

internal/api/structured_query.go (1)

15-15: LGTM!

Also applies to: 154-168, 177-177, 221-227

internal/api/pipe_deps_test.go (1)

5-119: LGTM!

Also applies to: 121-135, 137-186, 188-228, 230-257, 259-311, 313-346

internal/api/pipes_test.go (1)

45-45: LGTM!

Also applies to: 62-62, 77-77, 91-91, 107-107, 131-131, 158-158, 179-179, 200-200, 224-224, 249-249, 278-278, 297-297, 317-317, 334-334, 353-353, 370-370, 392-392, 415-415, 441-441, 462-462, 483-483

internal/api/errors_test.go (1)

17-17: LGTM!

Also applies to: 140-140

internal/cache/cache.go (1)

10-72: LGTM!

internal/cache/version_manager.go (1)

32-84: LGTM!

internal/cache/local_test.go (1)

137-239: LGTM!

internal/discovery/discovery.go (2)

246-367: LGTM!

Also applies to: 571-664


461-464: 🎯 Functional Correctness

viewDef is declared once. No change is needed.

internal/discovery/explain.go (1)

158-309: LGTM!

Also applies to: 340-435

internal/discovery/deps_test.go (1)

61-134: LGTM!

Also applies to: 199-290, 297-374, 504-602

internal/discovery/discovery_test.go (1)

158-160: LGTM!

Also applies to: 178-180

internal/discovery/explain_test.go (1)

217-301: LGTM!

Also applies to: 303-343

internal/discovery/fuzz_deps_test.go (1)

1-357: LGTM!

internal/cache/local.go (1)

81-101: 🩺 Stability & Availability

No change needed. Production invalidation passes encoded table inputs where cached dependency keys are NATS-encoded.

cmd/wavehouse/main.go (1)

186-204: LGTM!

Also applies to: 269-273, 282-282, 297-301, 344-368, 421-422

internal/api/boot_chain_test.go (1)

115-118: LGTM!

AGENTS.md (1)

56-62: LGTM!

CHANGELOG.md (1)

14-14: LGTM!

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

480-480: LGTM!

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

55-55: LGTM!

Also applies to: 117-117, 148-148, 157-158

docs/src/content/docs/pipes.mdx (1)

122-122: LGTM!

Also applies to: 174-192

tests/e2e/sdk/cache.test.ts (2)

2-10: LGTM!


74-95: 🩺 Stability & Availability

Clarify whether suiteTables should create materialized-view dependencies.

tests/e2e/sdk/cache.test.ts creates mv_src_${stamp}, mv_tgt_${stamp}, and mv_${stamp} with raw DDL even though T = suiteTables("cache") exists. The suite-table path only creates predefined clicks/events/users tables, so either document these raw object names as acceptable for this edge case or extend suiteTables to include cache-specific dependencies if all ClickHouse objects should be owned centrally.

tests/integration/pipe_cache_test.go (1)

26-57: LGTM!

Also applies to: 66-105, 176-209, 222-256, 271-298, 306-322

tests/integration/setup_test.go (1)

81-89: LGTM!

Also applies to: 320-333

Comment thread cmd/wavehouse/main.go
Comment on lines +381 to +391
registry.SetOnRefresh(func(snap discovery.DependencySnapshot) {
resultCache.SetDependents(snap.Cascade)
if len(snap.ChangedViews) > 0 {
nss := make([]cache.Namespace, len(snap.ChangedViews))
for i, v := range snap.ChangedViews {
nss[i] = cache.Namespace{Table: v}
}
_, _ = resultCache.Invalidate(ctx, nss)
}
pipesHandler.ClearResolvedDeps()
})

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win

Log the failure of the refresh-time invalidation.

resultCache.Invalidate returns an error that this callback discards. If invalidation fails after a view definition changes, pipes keep serving results keyed on the pre-change version and no signal reaches the operator. Log the error at WARN with the affected namespaces.

♻️ Proposed change
 		if len(snap.ChangedViews) > 0 {
 			nss := make([]cache.Namespace, len(snap.ChangedViews))
 			for i, v := range snap.ChangedViews {
 				nss[i] = cache.Namespace{Table: v}
 			}
-			_, _ = resultCache.Invalidate(ctx, nss)
+			if _, err := resultCache.Invalidate(ctx, nss); err != nil {
+				logger.Warn("invalidate changed view namespaces", "error", err, "views", len(nss))
+			}
 		}

As per coding guidelines: "Use structured logging with log/slog and JSON handlers."

Source: Coding guidelines

### `discovery/` — Schema Discovery & Validation

- **discovery.go** — `SchemaRegistry` queries `system.columns` to discover ClickHouse table schemas. Supports periodic auto-refresh, on-demand refresh, and `RetryRefresh` (boot-time exponential backoff loop used by `cmd/wavehouse` so a transiently unreachable ClickHouse doesn't crash-loop the binary). Thread-safe via `sync.RWMutex`.
- **discovery.go** — `SchemaRegistry` queries `system.columns` to discover ClickHouse table schemas. Supports periodic auto-refresh, on-demand refresh, and `RetryRefresh` (boot-time exponential backoff loop used by `cmd/wavehouse` so a transiently unreachable ClickHouse doesn't crash-loop the binary). The registry also owns the cache-invalidation dependency graph: on each content-changed refresh it resolves every view and materialized-view definition (via `ResolveTables`, below) into a base-table→dependents cascade — carrying an MV's trigger edge through to its `TO` target, parsed off the rendered `create_table_query` (`parseMVTarget`) and propagated along `system.tables.dependencies_table` — exposed as `Dependents()` and pushed into the query cache via `SetOnRefresh`; `IsKnown` reports whether a name resolved cleanly into that cascade (the pipe handler's test for a dependency version invalidation can be trusted to watch — see the `pipes/` section for the full invalidation model). Thread-safe via `sync.RWMutex`.

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

📐 Maintainability & Code Quality | 🟡 Minor | ⚡ Quick win

Fix the garbled clause describing IsKnown.

The parenthetical reads "the pipe handler's test for a dependency version invalidation can be trusted to watch". The sentence has no coherent subject-verb structure, so the meaning of IsKnown is lost at the point the reader needs it.

✏️ Proposed wording
-`IsKnown` reports whether a name resolved cleanly into that cascade (the pipe handler's test for a dependency version invalidation can be trusted to watch — see the `pipes/` section for the full invalidation model)
+`IsKnown` reports whether a name resolved cleanly into that cascade — the pipe handler's test for whether a dependency's version can be trusted to invalidate on write (see the `pipes/` section for the full invalidation model)
📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
- **discovery.go** — `SchemaRegistry` queries `system.columns` to discover ClickHouse table schemas. Supports periodic auto-refresh, on-demand refresh, and `RetryRefresh` (boot-time exponential backoff loop used by `cmd/wavehouse` so a transiently unreachable ClickHouse doesn't crash-loop the binary). The registry also owns the cache-invalidation dependency graph: on each content-changed refresh it resolves every view and materialized-view definition (via `ResolveTables`, below) into a base-table→dependents cascade — carrying an MV's trigger edge through to its `TO` target, parsed off the rendered `create_table_query` (`parseMVTarget`) and propagated along `system.tables.dependencies_table` — exposed as `Dependents()` and pushed into the query cache via `SetOnRefresh`; `IsKnown` reports whether a name resolved cleanly into that cascade (the pipe handler's test for a dependency version invalidation can be trusted to watch — see the `pipes/` section for the full invalidation model). Thread-safe via `sync.RWMutex`.
- **discovery.go** — `SchemaRegistry` queries `system.columns` to discover ClickHouse table schemas. Supports periodic auto-refresh, on-demand refresh, and `RetryRefresh` (boot-time exponential backoff loop used by `cmd/wavehouse` so a transiently unreachable ClickHouse doesn't crash-loop the binary). The registry also owns the cache-invalidation dependency graph: on each content-changed refresh it resolves every view and materialized-view definition (via `ResolveTables`, below) into a base-table→dependents cascade — carrying an MV's trigger edge through to its `TO` target, parsed off the rendered `create_table_query` (`parseMVTarget`) and propagated along `system.tables.dependencies_table` — exposed as `Dependents()` and pushed into the query cache via `SetOnRefresh`; `IsKnown` reports whether a name resolved cleanly into that cascade — the pipe handler's test for whether a dependency's version can be trusted to invalidate on write (see the `pipes/` section for the full invalidation model). Thread-safe via `sync.RWMutex`.

Comment on lines +376 to +385
<-started
// The first resolution is now blocked in flight; give the remaining waiters
// time to reach the singleflight and join it. A straggler that misses even
// this window would still be served by the memo after the flight completes,
// so the assertion below stays stable.
time.Sleep(100 * time.Millisecond)
close(release)
wg.Wait()

assert.Equal(t, int32(1), calls.Load(), "concurrent cold misses must coalesce onto one EXPLAIN")

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🩺 Stability & Availability | 🟡 Minor | ⚡ Quick win

The sleep-based barrier can produce a false failure, and the comment overstates the guarantee.

The comment says a straggler is still served by the memo. That is not true for this assertion. resolveDeps reads the memo BEFORE it joins resolveSF. A goroutine that reads the empty memo and is then descheduled past close(release) reaches resolveSF.Do after the first flight finished, so it starts a second EXPLAIN. calls becomes 2 and assert.Equal(t, int32(1), ...) fails. A loaded CI machine makes this reachable.

Gate the release on all waiters actually reaching resolveDeps, and keep the sleep only as a short settle window.

🩹 Proposed fix
 	const waiters = 8
 	var wg sync.WaitGroup
+	var entered sync.WaitGroup
+	entered.Add(waiters)
 	results := make([][]cache.Namespace, waiters)
 	for i := range waiters {
 		wg.Add(1)
 		go func() {
 			defer wg.Done()
+			entered.Done()
 			results[i], _ = h.resolveDeps(ctx, "SELECT * FROM events")
 		}()
 	}
 	<-started
-	// The first resolution is now blocked in flight; give the remaining waiters
-	// time to reach the singleflight and join it. A straggler that misses even
-	// this window would still be served by the memo after the flight completes,
-	// so the assertion below stays stable.
+	// Every waiter goroutine is now scheduled and the first resolution is blocked
+	// in flight; give the remaining waiters time to reach the singleflight and
+	// join it before the flight is allowed to complete.
+	entered.Wait()
 	time.Sleep(100 * time.Millisecond)
 	close(release)

Comment thread internal/api/pipes.go
Comment on lines 284 to +294
if h.Cache != nil {
_ = h.Cache.Set(r.Context(), cacheKey, nil, data, ttl)
ttl := cache.QueryTimeToTTL(queryDuration)
if depsUnresolved {
// The result can't be reliably version-invalidated — an unfoldable
// view, an external read (table function/cross-database), or a
// dead-branch-pruned set (belt-and-suspenders): cap the TTL so it
// self-expires rather than serving stale on a write the folded
// versions never see.
ttl = min(ttl, cache.UnresolvedDepsTTLCap)
}
_ = h.Cache.Set(r.Context(), entryKey, data, ttl)

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🩺 Stability & Availability | 🟡 Minor | ⚡ Quick win

Cache writes are bound to the leader request's cancellable context in both handlers. Both Set calls run inside a singleflight closure owned by whichever request won the race. If that leader disconnects after the query completes, its context is already canceled, so the write is dropped while the waiters still receive the data. The next reader then re-executes the query, which weakens the caching this PR is built around.

  • internal/api/pipes.go#L284-L294: wrap the Set context with context.WithoutCancel(r.Context()); keep the depsUnresolved TTL cap as written.
  • internal/api/structured_query.go#L229-L230: wrap the Set context with context.WithoutCancel(r.Context()); keep the IsKnown TTL cap as written.
📍 Affects 2 files
  • internal/api/pipes.go#L284-L294 (this comment)
  • internal/api/structured_query.go#L229-L230

Comment thread internal/api/pipes.go
Comment on lines +344 to +354
// Coalesce concurrent cold misses for the same bound query: one EXPLAIN,
// every waiter shares the outcome. The generation is captured INSIDE the
// flight, before the EXPLAIN runs, so it dates the schema state the
// resolution was computed against for every waiter alike.
v, _, _ := h.resolveSF.Do(boundSQL, func() (any, error) {
h.pipeDepsMu.RLock()
gen := h.pipeDepsGen
h.pipeDepsMu.RUnlock()
rd, memoize := h.resolvePipe(ctx, boundSQL)
return resolveOutcome{rd: rd, memoize: memoize, gen: gen}, nil
})

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🩺 Stability & Availability | 🟡 Minor | ⚡ Quick win

The shared resolution inherits only the leader's request context.

resolveSF.Do runs resolvePipe with ctx from whichever request won the race. resolvePipe derives its 5s deadline from that context. If the leader client disconnects mid-EXPLAIN, the resolution is canceled for every coalesced waiter. All of them then serve on the database-version fallback, and the transient classification means nothing is memoized, so the next request pays another EXPLAIN.

The resolution is already bounded to 5s inside resolvePipe, so detaching cancellation does not create an unbounded operation. Pass a cancellation-detached context into the flight.

🩹 Proposed fix
 		v, _, _ := h.resolveSF.Do(boundSQL, func() (any, error) {
 			h.pipeDepsMu.RLock()
 			gen := h.pipeDepsGen
 			h.pipeDepsMu.RUnlock()
-			rd, memoize := h.resolvePipe(ctx, boundSQL)
+			// The flight is shared, so it must not die with the leader's request.
+			// resolvePipe still bounds itself to 5s.
+			rd, memoize := h.resolvePipe(context.WithoutCancel(ctx), boundSQL)
 			return resolveOutcome{rd: rd, memoize: memoize, gen: gen}, nil
 		})
📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
// Coalesce concurrent cold misses for the same bound query: one EXPLAIN,
// every waiter shares the outcome. The generation is captured INSIDE the
// flight, before the EXPLAIN runs, so it dates the schema state the
// resolution was computed against for every waiter alike.
v, _, _ := h.resolveSF.Do(boundSQL, func() (any, error) {
h.pipeDepsMu.RLock()
gen := h.pipeDepsGen
h.pipeDepsMu.RUnlock()
rd, memoize := h.resolvePipe(ctx, boundSQL)
return resolveOutcome{rd: rd, memoize: memoize, gen: gen}, nil
})
// Coalesce concurrent cold misses for the same bound query: one EXPLAIN,
// every waiter shares the outcome. The generation is captured INSIDE the
// flight, before the EXPLAIN runs, so it dates the schema state the
// resolution was computed against for every waiter alike.
v, _, _ := h.resolveSF.Do(boundSQL, func() (any, error) {
h.pipeDepsMu.RLock()
gen := h.pipeDepsGen
h.pipeDepsMu.RUnlock()
// The flight is shared, so it must not die with the leader's request.
// resolvePipe still bounds itself to 5s.
rd, memoize := h.resolvePipe(context.WithoutCancel(ctx), boundSQL)
return resolveOutcome{rd: rd, memoize: memoize, gen: gen}, nil
})

Comment on lines +6 to +16
tests := []struct{ in, want string }{
{"acme", "'acme'"},
{"a'b", "'a''b'"},
{"a\\b", "'a\\\\b'"},
{"org_id = '1'", "'org_id = ''1'''"},
}
for _, tt := range tests {
if got := QuoteString(tt.in); got != tt.want {
t.Errorf("QuoteString(%q) = %q, want %q", tt.in, got, tt.want)
}
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

📐 Maintainability & Code Quality | 🟠 Major | ⚡ Quick win

Run each quoting case as a named subtest.

This table has multiple scenarios but does not call t.Run(tt.name, ...). Add a name field and run each assertion in its own subtest.

Proposed change
-	tests := []struct{ in, want string }{
-		{"acme", "'acme'"},
-		{"a'b", "'a''b'"},
-		{"a\\b", "'a\\\\b'"},
-		{"org_id = '1'", "'org_id = ''1'''"},
+	tests := []struct{ name, in, want string }{
+		{"plain", "acme", "'acme'"},
+		{"quote", "a'b", "'a''b'"},
+		{"backslash", "a\\b", "'a\\\\b'"},
+		{"sql_like", "org_id = '1'", "'org_id = ''1'''"},
 	}
 	for _, tt := range tests {
-		if got := QuoteString(tt.in); got != tt.want {
-			t.Errorf("QuoteString(%q) = %q, want %q", tt.in, got, tt.want)
-		}
+		t.Run(tt.name, func(t *testing.T) {
+			if got := QuoteString(tt.in); got != tt.want {
+				t.Errorf("QuoteString(%q) = %q, want %q", tt.in, got, tt.want)
+			}
+		})
 	}

As per coding guidelines, “Use table-driven tests with t.Run(tt.name, ...) for multiple scenarios.”

📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
tests := []struct{ in, want string }{
{"acme", "'acme'"},
{"a'b", "'a''b'"},
{"a\\b", "'a\\\\b'"},
{"org_id = '1'", "'org_id = ''1'''"},
}
for _, tt := range tests {
if got := QuoteString(tt.in); got != tt.want {
t.Errorf("QuoteString(%q) = %q, want %q", tt.in, got, tt.want)
}
}
tests := []struct{ name, in, want string }{
{"plain", "acme", "'acme'"},
{"quote", "a'b", "'a''b'"},
{"backslash", "a\\b", "'a\\\\b'"},
{"sql_like", "org_id = '1'", "'org_id = ''1'''"},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := QuoteString(tt.in); got != tt.want {
t.Errorf("QuoteString(%q) = %q, want %q", tt.in, got, tt.want)
}
})
}

Source: Coding guidelines

Comment on lines +457 to +459
}
return infos, rows.Err()
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win

Unwrapped rows.Err() returns in the new discovery paths. Three new row-iteration paths return the driver error without context, while every sibling error path in the same functions wraps with fmt.Errorf. All three propagate up through Refresh or ResolveTables, so an operator sees a bare driver message with no indication of which read failed.

  • internal/discovery/discovery.go#L457-L459: wrap the discoverColumns return as fmt.Errorf("iterate system.columns rows: %w", err).
  • internal/discovery/discovery.go#L520-L524: wrap rerr in discoverViewMeta as fmt.Errorf("iterate system.tables rows: %w", rerr).
  • internal/discovery/explain.go#L85-L87: wrap the ResolveTables return as fmt.Errorf("iterate explain rows: %w", err).

As per coding guidelines: "Return errors instead of panicking, and wrap errors with fmt.Errorf("context: %w", err)."

📍 Affects 2 files
  • internal/discovery/discovery.go#L457-L459 (this comment)
  • internal/discovery/discovery.go#L520-L524
  • internal/discovery/explain.go#L85-L87

Source: Coding guidelines

Comment on lines +140 to +172
{
name: "dictGet: dictionary captured; the projection's stray '' before it is ignored",
lines: []string{
" CONSTANT id: 2, constant_value: '', constant_value_type: String",
" FUNCTION id: 3, function_name: dictGet, function_type: ordinary, result_type: String",
" CONSTANT id: 5, constant_value: 'default.mydict', constant_value_type: String",
" CONSTANT id: 6, constant_value: 'val', constant_value_type: String",
},
db: "default",
wantDicts: []string{"mydict"},
},
{
name: "dictGet cross-db dropped and marks the read external",
lines: []string{
" FUNCTION id: 3, function_name: dictGet, function_type: ordinary, result_type: String",
" CONSTANT id: 5, constant_value: 'otherdb.mydict', constant_value_type: String",
},
db: "default",
wantExternal: true,
},
{
name: "joinGet and dictGet together",
lines: []string{
" TABLE id: 3, table_name: default.users",
" FUNCTION id: 4, function_name: joinGet, function_type: ordinary, result_type: String",
" CONSTANT id: 6, constant_value: 'default.jt', constant_value_type: String",
" FUNCTION id: 8, function_name: dictGetOrDefault, function_type: ordinary, result_type: String",
" CONSTANT id: 9, constant_value: 'default.mydict', constant_value_type: String",
},
db: "default",
wantTables: []string{"jt", "users"},
wantDicts: []string{"mydict"},
},

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win

Add unit cases for dictHas and dictIsIn.

parseExplainTables matches three dict prefixes: dictGet, dictHas, and dictIsIn. The table only covers dictGet and dictGetOrDefault. dictHas and dictIsIn are covered only by the fuzzdeps suite, which is build-tagged and needs a live ClickHouse. A regression that drops either match would pass make ci. The earlier review thread settled on keeping these two names outside the dictGet prefix, so pin that decision here.

💚 Proposed additional cases
 		{
 			name: "dictGet cross-db dropped and marks the read external",

Add before that entry:

{
	name: "dictHas captures the dictionary",
	lines: []string{
		"          FUNCTION id: 3, function_name: dictHas, function_type: ordinary, result_type: UInt8",
		"                CONSTANT id: 5, constant_value: 'default.mydict', constant_value_type: String",
	},
	db:        "default",
	wantDicts: []string{"mydict"},
},
{
	name: "dictIsIn captures the dictionary",
	lines: []string{
		"          FUNCTION id: 3, function_name: dictIsIn, function_type: ordinary, result_type: UInt8",
		"                CONSTANT id: 5, constant_value: 'default.mydict', constant_value_type: String",
	},
	db:        "default",
	wantDicts: []string{"mydict"},
},
📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
{
name: "dictGet: dictionary captured; the projection's stray '' before it is ignored",
lines: []string{
" CONSTANT id: 2, constant_value: '', constant_value_type: String",
" FUNCTION id: 3, function_name: dictGet, function_type: ordinary, result_type: String",
" CONSTANT id: 5, constant_value: 'default.mydict', constant_value_type: String",
" CONSTANT id: 6, constant_value: 'val', constant_value_type: String",
},
db: "default",
wantDicts: []string{"mydict"},
},
{
name: "dictGet cross-db dropped and marks the read external",
lines: []string{
" FUNCTION id: 3, function_name: dictGet, function_type: ordinary, result_type: String",
" CONSTANT id: 5, constant_value: 'otherdb.mydict', constant_value_type: String",
},
db: "default",
wantExternal: true,
},
{
name: "joinGet and dictGet together",
lines: []string{
" TABLE id: 3, table_name: default.users",
" FUNCTION id: 4, function_name: joinGet, function_type: ordinary, result_type: String",
" CONSTANT id: 6, constant_value: 'default.jt', constant_value_type: String",
" FUNCTION id: 8, function_name: dictGetOrDefault, function_type: ordinary, result_type: String",
" CONSTANT id: 9, constant_value: 'default.mydict', constant_value_type: String",
},
db: "default",
wantTables: []string{"jt", "users"},
wantDicts: []string{"mydict"},
},
{
name: "dictGet: dictionary captured; the projection's stray '' before it is ignored",
lines: []string{
" CONSTANT id: 2, constant_value: '', constant_value_type: String",
" FUNCTION id: 3, function_name: dictGet, function_type: ordinary, result_type: String",
" CONSTANT id: 5, constant_value: 'default.mydict', constant_value_type: String",
" CONSTANT id: 6, constant_value: 'val', constant_value_type: String",
},
db: "default",
wantDicts: []string{"mydict"},
},
{
name: "dictHas captures the dictionary",
lines: []string{
" FUNCTION id: 3, function_name: dictHas, function_type: ordinary, result_type: UInt8",
" CONSTANT id: 5, constant_value: 'default.mydict', constant_value_type: String",
},
db: "default",
wantDicts: []string{"mydict"},
},
{
name: "dictIsIn captures the dictionary",
lines: []string{
" FUNCTION id: 3, function_name: dictIsIn, function_type: ordinary, result_type: UInt8",
" CONSTANT id: 5, constant_value: 'default.mydict', constant_value_type: String",
},
db: "default",
wantDicts: []string{"mydict"},
},
{
name: "dictGet cross-db dropped and marks the read external",
lines: []string{
" FUNCTION id: 3, function_name: dictGet, function_type: ordinary, result_type: String",
" CONSTANT id: 5, constant_value: 'otherdb.mydict', constant_value_type: String",
},
db: "default",
wantExternal: true,
},
{
name: "joinGet and dictGet together",
lines: []string{
" TABLE id: 3, table_name: default.users",
" FUNCTION id: 4, function_name: joinGet, function_type: ordinary, result_type: String",
" CONSTANT id: 6, constant_value: 'default.jt', constant_value_type: String",
" FUNCTION id: 8, function_name: dictGetOrDefault, function_type: ordinary, result_type: String",
" CONSTANT id: 9, constant_value: 'default.mydict', constant_value_type: String",
},
db: "default",
wantTables: []string{"jt", "users"},
wantDicts: []string{"mydict"},
},

Comment on lines +147 to +151
await admin.pipes.delete(pipeName);
await chQuery(`DROP VIEW IF EXISTS default.\`${mv}\``);
await chQuery(`DROP TABLE IF EXISTS default.\`${tgt}\``);
await chQuery(`DROP TABLE IF EXISTS default.\`${src}\``);
}, 30_000);

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🩺 Stability & Availability | 🟠 Major | ⚡ Quick win

Move the teardown so it also runs when an assertion fails.

The four cleanup calls sit at the end of the test body. If any expect fails or either waitForCondition times out, the test throws before line 147. The pipe stays registered in the pipes KV store, and the two tables plus the materialized view stay in ClickHouse. This suite runs with maxWorkers: 1 against shared global state, so leaked objects accumulate across runs and a leaked pipe remains visible to later pipes.list assertions.

Wrap the body in try/finally, or register the teardown right after each object is created.

♻️ Proposed change
-    await admin.pipes.delete(pipeName);
-    await chQuery(`DROP VIEW IF EXISTS default.\`${mv}\``);
-    await chQuery(`DROP TABLE IF EXISTS default.\`${tgt}\``);
-    await chQuery(`DROP TABLE IF EXISTS default.\`${src}\``);
   }, 30_000);

and wrap the steps that follow object creation:

try {
  // refresh, pipe set, exec/assert, ingest, waits, final assertions
} finally {
  await admin.pipes.delete(pipeName);
  await chQuery(`DROP VIEW IF EXISTS default.\`${mv}\``);
  await chQuery(`DROP TABLE IF EXISTS default.\`${tgt}\``);
  await chQuery(`DROP TABLE IF EXISTS default.\`${src}\``);
}

As per path instructions: "Keep E2E tests sequential with maxWorkers: 1 because they share mutable global policy state."

Source: Path instructions

Comment on lines +143 to +148
boundSQL := strings.ReplaceAll(sql, "{{source:web}}", "'web'")
res, err := discovery.ResolveTables(ctx, e.chConn, testCHDatabase, boundSQL)
require.NoError(t, err)
require.Contains(t, res.Tables, web, "the live arm's table stays tracked")
require.NotContains(t, res.Tables, mobile, "the constant-false arm's table must be pruned")
require.True(t, res.Pruned, "ResolveTables must report the prune (it drives the TTL cap)")

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
# Description: Locate the pipe parameter binder and its call site in the handler.
set -euo pipefail

fd -t f 'pipes.go' internal | while IFS= read -r f; do
  echo "=== outline $f ==="
  ast-grep outline "$f" --items all
done

echo "=== BindParams definition ==="
ast-grep run --pattern 'func $$$BindParams($$$) { $$$ }' --lang go internal || true
rg -nP --type=go -C 5 '\bfunc\s+(\([^)]*\)\s+)?BindParams\s*\('

echo "=== BindParams call sites ==="
rg -nP --type=go -C 4 '\bBindParams\s*\('

Repository: Wave-RF/WaveHouse

Length of output: 3856


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "=== tests/integration/pipe_cache_test.go focus ==="
sed -n '120,170p' tests/integration/pipe_cache_test.go

echo "=== internal/pipes/pipes.go BindParams and formatParamValue ==="
sed -n '140,245p' internal/pipes/pipes.go

echo "=== internal/api/pipes.go Execute resolution path ==="
sed -n '180,320p' internal/api/pipes.go

echo "=== tests references for pin/coupling SQL handling ==="
rg -n --type=go -C 3 'pin|prune|BindParams|ResolveTables|ResolveTablesFn|source:web' tests internal

Repository: Wave-RF/WaveHouse

Length of output: 50374


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "=== internal/pipes/pipes.go imports and QuoteString ==="
sed -n '1,20p' internal/pipes/pipes.go
rg -n --type=go -C 4 'func QuoteString|QuoteString\(' internal

echo "=== pipes helper used by tests maybe ==="
sed -n '1,170p' tests/integration/pipe_cache_test.go

Repository: Wave-RF/WaveHouse

Length of output: 11041


Bind the pin SQL with pipes.BindParams.

streams.ReplaceAll(sql, "{{source:web}}", "'web'") duplicates the binder’s current output. Use pipes.BindParams(&pipes.NamedQuery{SQL: sql}, nil) so the pruning pin stays tied to the same BindParams rendering that /v1/pipes/{name}/execute uses.

@github-project-automation github-project-automation Bot moved this from Ready to In review in WaveHouse Task Board Aug 7, 2026
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/cache Local / shared / tiered caching area/docs Documentation, site/, README area/infra CI, build, deploy, Docker, release area/ingest Ingest pipeline (Bento, batching, DLQ) area/pipes Named query pipes area/query Structured query AST, SQL builder area/sdk TypeScript SDK (clients/ts/) documentation Improvements or additions to documentation go Pull requests that update go code

Projects

Status: In review

2 participants