Skip to content

feat(capture): each output reads its own topic and producer - #106535

Draft
pl wants to merge 6 commits into
masterfrom
pl/ingestion/capture-output-config
Draft

pl wants to merge 6 commits into
masterfrom
pl/ingestion/capture-output-config

Conversation

@pl

@pl pl commented Sep 25, 2026 •

Copy link
Copy Markdown
Contributor

Problem

Moving one capture output to another Kafka cluster needs a code change. Every v0 output publishes through the one INGESTION producer, and the topics come from shared KAFKA_*_TOPIC variables. This is Step 10 of the capture outputs plan, and the emergency fallback (Step 17) builds on it.

Changes

  • Each v0 Destination is an output with a topic and a named producer: CAPTURE_OUTPUT_<OUTPUT>_TOPIC and CAPTURE_OUTPUT_<OUTPUT>_PRODUCER. This is the same shape as Node ingestion's INGESTION_OUTPUT_<OUTPUT>_{TOPIC,PRODUCER}.
  • Outputs: ANALYTICS_MAIN, ANALYTICS_OVERFLOW, ANALYTICS_HISTORICAL, SESSION_REPLAY_MAIN, SESSION_REPLAY_OVERFLOW, HEATMAPS, CLIENT_WARNINGS, ERROR_TRACKING, DLQ, AI_MAIN, AI_OVERFLOW. Custom redirects use CAPTURE_OUTPUT_CUSTOM_PRODUCER.
  • Explicit migration: capture no longer reads KAFKA_TOPIC, KAFKA_*_TOPIC, or CAPTURE_ANALYTICS_AI_EVENTS_*_TOPIC. charts#16337 sets the new names first and must be merged before this deploys.
  • Producers default to INGESTION. Unknown producer names fail boot.
  • Topic defaults are the local dev and hobby topics. Hobby (--pull always) and local dev pull capture:master with compose files that can predate these variables. With only the legacy names set, the new image uses the defaults, so they must route like the old compose:
    • SESSION_REPLAY_MAIN defaults to session_recording_snapshot_item_events, not events_plugin_ingestion. Otherwise recordings would land in the analytics events topic.
    • ERROR_TRACKING defaults to ingestion-errortracking-main, not error_tracking_events. That's the topic the Node consumer reads.
    • Every other default already matched. Production sets every reachable output in charts, so it's unaffected.
  • Session replay's main topic is its own output. It used to share KAFKA_TOPIC with analytics. Recordings pods mount only /s, so each deployment's KAFKA_TOPIC maps to exactly one of the two outputs.
  • TopicTable becomes OutputTable, generic over the producer: names in config, handles in the sink. The sink maps names to handles once, at construction.
  • PreparedPayload carries the Destination, as v1's PreparedEvent does. At enqueue, the sink publishes through that target's own producer, with no second lookup.
  • The Kafka sink holds only its output table and no longer knows which producers exist. Shutdown flushes the ProducerRegistry directly, and flush comes off PublishEvents, Output, OutputRegistry and Sink. Only the Kafka leaf ever flushed through that path.
  • Topics are Arc<str>, so enqueue allocates once per record, down from two allocations across prep and enqueue before.
  • docker-compose.base.yml keeps the legacy names beside the new ones. A stale local capture:master image would otherwise send replay to the analytics default topic. The cleanup PR removes the legacy names.
  • bin/start-rust-service builds from source and switches to the new names outright.
  • The plan doc marks Steps 9 and 10 as shipped. Step 17's fallback becomes a secondary target per output, as in Node's dual-write outputs.

How did you test this code?

  • cargo test -p capture --lib: all pass except 5 event_restrictions::repository tests that need a local Redis.
  • cargo clippy -p capture --tests -- -D warnings and cargo fmt are clean.
  • The integration suites under rust/capture/tests/ compile but were not run locally (they need Kafka and Redis). CI runs them. Their assertions are unchanged; only config field names changed.
  • New tests catch:
    • a legacy topic variable that is still read (legacy_topic_env_var_is_not_read)
    • an undeclared producer name that boots (producer_name_parses_only_declared_slots, output_producers_default_to_ingestion_and_parse_their_env_var)
    • replay main sharing the analytics topic again (session_replay_main_does_not_share_analytics_main)
    • a blank replay main topic that passes the completeness check
    • an output wired to another output's producer (each_output_resolves_to_its_own_producer)
    • a default that no longer matches the local stack, which breaks hobby with an old compose file (output_topic_defaults_match_the_local_stack)
  • Hobby and local dev compatibility was checked by reading the latest release's docker-compose.base.yml against the new defaults. No hobby stack was run.
  • Not checked: a helm render of the charts PR against this image, and a dev deploy.

Release status

  • No feature flag controls this change
  • This change is behind a feature flag and is not available to users
  • This change makes a previously flagged feature available to everyone

Automatic notifications

  • Publish to changelog?

Docs update

None. The outputs plan doc is updated in this PR.

🤖 Agent context

Autonomy: Human-driven (agent-assisted)

Agent: Claude Code, Claude Opus 5.5

  • Design decisions were made with the assignee: follow Node ingestion's per-output pattern, not one topic block per cluster, and migrate charts first.
  • CodeRabbit CLI: not run. The CLI was signed out and the assignee chose to skip it.
  • No duplicate: no other open PR implements outputs plan Step 10.
  • Public artifact: no session-only material is in the diff.

pl added 3 commits September 25, 2026 12:09
An output names its producer in an environment variable, so a slot name has
to parse. Only declared slots parse; anything else fails boot and lists the
valid names.
Every v0 destination is now an output with a topic and a named producer, read
from CAPTURE_OUTPUT_<OUTPUT>_TOPIC and CAPTURE_OUTPUT_<OUTPUT>_PRODUCER, the
same shape as Node ingestion's INGESTION_OUTPUT_<OUTPUT>_{TOPIC,PRODUCER}.
Moving one output to another cluster is now configuration.

Explicit migration: capture stops reading KAFKA_TOPIC, KAFKA_*_TOPIC and
CAPTURE_ANALYTICS_AI_EVENTS_*_TOPIC. charts#16337 sets the new names first.
Topic defaults are unchanged; producers default to INGESTION.

KAFKA_TOPIC splits in two. Session replay's main topic is its own output
instead of sharing the analytics main topic; recordings pods only reach replay
outputs, so each deployment maps its KAFKA_TOPIC to exactly one of them.

TopicTable becomes OutputTable. PreparedPayload carries the Destination, as
v1's PreparedEvent does, and the Kafka sink resolves it to a topic and producer
at enqueue. Topics are Arc<str>, so the serial enqueue allocates once per
record for the ProduceRecord topic. Custom redirects publish through
CAPTURE_OUTPUT_CUSTOM_PRODUCER.

docker-compose keeps the legacy names beside the new ones until the new image
is the published one; a stale capture:master would otherwise send replay to
the analytics default topic.
Step 17's fallback follows the per-output shape: a secondary target on each
output, as in Node ingestion's dual-write outputs.
@pl pl self-assigned this Sep 25, 2026
@trunk-io

trunk-io Bot commented Sep 25, 2026

Copy link
Copy Markdown

Merging to master in this repository is managed by Trunk.

  • To merge this pull request, check the box to the left or comment /trunk merge below.

After your PR is submitted to the merge queue, this comment will be automatically updated with its status. If the PR fails, failure details will also be posted here

@github-actions

github-actions Bot commented Sep 25, 2026 •

Copy link
Copy Markdown
Contributor

🤖 CI report

🚨 Trunk lane — universal lane

This PR is assigned to the universal lane. It cannot merge in parallel with other PRs, so it can take longer to merge. Ask dev-ex if you think this is wrong.

🚨 Comment density — 11% of added code lines are comments (70 of 662)

This section warns when comments are more than 3% of the code lines a PR adds, and alerts above 6%. Before agent-assisted PRs, the typical share was about 2%. Only full-line comments count. Docstrings, generated files, snapshots, migrations, and workflow files are left out.

Comments that restate the code, record how the change came about, or narrate the next line add noise for the next reader. Keep the comments that explain a reason the code cannot show, and remove the rest. See .agents/skills/writing-code-comments/SKILL.md for the house rules.

Files with the most added comment lines:

File Comment lines Added lines
rust/capture/src/sinks/registry.rs 22 200
rust/capture/src/sinks/kafka.rs 16 57
rust/capture/src/config.rs 9 157
rust/capture/src/setup.rs 6 74
rust/capture/src/v1/sinks/kafka/config.rs 4 4
rust/capture/src/outputs.rs 2 2
rust/capture/src/sinks/sink.rs 2 4
rust/capture/src/events/analytics.rs 1 11

This check does not block merging. It updates on every push and clears when the share drops.

⚠️ Bundle size — 🔺 +25.2 KiB (+0.0%)

Uncompressed size of every built .js bundle, compared against the base branch.

Total: 68.82 MiB · 🔺 +25.2 KiB (+0.0%)

File Size Δ vs base
render-query/src/render-query/render-query.js 20.23 MiB 🔺 +9.1 KiB (+0.0%)
posthog-app/_parent/products/engineering_analytics/frontend/scenes/EngineeringAnalyticsAuthorScene.js 23.4 KiB 🔺 +3.6 KiB (+18.4%)
posthog-app/_parent/products/engineering_analytics/frontend/scenes/PullRequestDetailScene.js 31.2 KiB 🔺 +2.5 KiB (+8.8%)
posthog-app/_parent/products/engineering_analytics/frontend/scenes/EngineeringAnalyticsScene.js 78.1 KiB 🔺 +2.3 KiB (+3.1%)
posthog-app/_parent/products/workflows/frontend/Broadcasts/BroadcastScene.js 42.4 KiB 🔺 +1.8 KiB (+4.4%)
posthog-app/_parent/products/engineering_analytics/frontend/scenes/EngineeringAnalyticsTeamScene.js 16.6 KiB 🔺 +1.5 KiB (+10.3%)
posthog-app/_parent/products/workflows/frontend/Broadcasts/BroadcastsScene.js 14.8 KiB 🔺 +1.2 KiB (+8.5%)
posthog-app/_parent/products/ai_observability/frontend/evaluations/AIObservabilityEvaluation.js 80.9 KiB 🔺 +1.1 KiB (+1.4%)

Posted automatically by build-bundle-size-report · uncompressed bytes from dist-report

✅ Eager graph — within budget

How much code each root ships on the eager path — downloaded and parsed before the surface is interactive. Measured from the esbuild output chunks (post-tree-shake, static imports only); lazy import() / React.lazy chunks are not counted.

Root Eager (shipped) Δ vs base Budget
entry (logged-out pages, app bootstrap)
src/index.tsx
1.57 MiB · 22 files 🔺 +916 B (+0.1%) █████████░ 85.1% of 1.84 MiB
logged-out boot: index + App + bootApp (preloaded by every page, including /login)
src/index.tsx + src/scenes/App.tsx + src/scenes/bootApp.ts
3.57 MiB · 628 files 🔺 +1.2 KiB (+0.0%) █████████░ 88.6% of 4.03 MiB
authenticated shell (every logged-in page)
src/scenes/AuthenticatedShell.tsx
7.27 MiB · 2,301 files 🔺 +9.1 KiB (+0.1%) █████████░ 87.2% of 8.34 MiB

🟢 node_modules/monaco-editor/ stays out of src/index.tsx
🟢 src/lib/components/ActivityLog/describers stays out of src/index.tsx
🟢 [object Object] stays out of src/index.tsx
🟢 [object Object] stays out of src/index.tsx
🟢 node_modules/monaco-editor/ stays out of src/index.tsx + src/scenes/App.tsx + src/scenes/bootApp.ts
🟢 src/layout/navigation-3000/navigationLogic.tsx stays out of src/index.tsx + src/scenes/App.tsx + src/scenes/bootApp.ts
🟢 src/scenes/dashboard/dashboardLogic.tsx stays out of src/index.tsx + src/scenes/App.tsx + src/scenes/bootApp.ts
🟢 src/lib/lemon-ui/LemonMarkdown/ stays out of src/index.tsx + src/scenes/App.tsx + src/scenes/bootApp.ts
🟢 src/lib/components/RichContentEditor/ stays out of src/index.tsx + src/scenes/App.tsx + src/scenes/bootApp.ts
🟢 src/lib/components/CodeSnippet/ stays out of src/index.tsx + src/scenes/App.tsx + src/scenes/bootApp.ts
🟢 src/taxonomy/core-filter-definitions-by-group.json stays out of src/index.tsx + src/scenes/App.tsx + src/scenes/bootApp.ts
🟢 node_modules/monaco-editor/ stays out of src/scenes/AuthenticatedShell.tsx
🟢 src/lib/components/ActivityLog/describers stays out of src/scenes/AuthenticatedShell.tsx
🟢 [object Object] stays out of src/scenes/AuthenticatedShell.tsx
🟢 src/scenes/session-recordings/player/sessionRecordingPlayerLogic.ts stays out of src/scenes/AuthenticatedShell.tsx
🟢 [object Object] stays out of src/scenes/AuthenticatedShell.tsx
🟢 [object Object] stays out of src/scenes/AuthenticatedShell.tsx
🟢 [object Object] stays out of src/scenes/AuthenticatedShell.tsx

Largest files eagerly shipped from src/index.tsx
Size File
126.8 KiB ../node_modules/.pnpm/react-dom@18.3.1_react@18.3.1/node_modules/react-dom/cjs/react-dom.production.min.js
24.6 KiB ../node_modules/.pnpm/buffer@6.0.3/node_modules/buffer/index.js
6.3 KiB ../node_modules/.pnpm/react@18.3.1/node_modules/react/cjs/react.production.min.js
4.5 KiB ../node_modules/.pnpm/@jspm+core@2.1.0/node_modules/@jspm/core/nodelibs/browser/process.js
3.9 KiB ../node_modules/.pnpm/scheduler@0.23.2/node_modules/scheduler/cjs/scheduler.production.min.js
1.4 KiB ../node_modules/.pnpm/base64-js@1.5.1/node_modules/base64-js/index.js
1.3 KiB src/index.tsx
1.3 KiB src/RootErrorBoundary.tsx
912 B ../node_modules/.pnpm/ieee754@1.2.1/node_modules/ieee754/index.js
854 B src/scenes/ChunkLoadErrorBoundary.tsx
Largest files eagerly shipped from src/index.tsx + src/scenes/App.tsx + src/scenes/bootApp.ts
Size File
301.8 KiB ../node_modules/.pnpm/posthog-js@1.434.13_@types+react@18.3.27_react@18.3.1/node_modules/posthog-js/dist/module.mjs
267.6 KiB ../node_modules/.pnpm/@posthog+icons@0.38.0_react-dom@18.3.1_react@18.3.1__react@18.3.1/node_modules/@posthog/icons/dist/posthog-icons.es.js
126.8 KiB ../node_modules/.pnpm/react-dom@18.3.1_react@18.3.1/node_modules/react-dom/cjs/react-dom.production.min.js
100.4 KiB src/lib/api.ts
84.9 KiB src/products.tsx
69.1 KiB src/lib/lemon-ui/icons/icons.tsx
63.9 KiB src/lib/utils/eventUsageLogic.ts
38.7 KiB ../node_modules/.pnpm/@dnd-kit+core@6.0.8_react-dom@18.3.1_react@18.3.1__react@18.3.1/node_modules/@dnd-kit/core/dist/core.esm.js
33.9 KiB ../node_modules/.pnpm/kea@4.0.0-pre.6_patch_hash=139b8d1f1304f9d9da452a9a1244c94ea679dbcb85687d8999563146879fb6f5_react@18.3.1/node_modules/kea/lib/index.cjs.js
28.3 KiB src/scenes/scenes.ts
Largest files eagerly shipped from src/scenes/AuthenticatedShell.tsx
Size File
301.8 KiB ../node_modules/.pnpm/posthog-js@1.434.13_@types+react@18.3.27_react@18.3.1/node_modules/posthog-js/dist/module.mjs
271.7 KiB src/taxonomy/core-filter-definitions-by-group.json
267.6 KiB ../node_modules/.pnpm/@posthog+icons@0.38.0_react-dom@18.3.1_react@18.3.1__react@18.3.1/node_modules/@posthog/icons/dist/posthog-icons.es.js
153.7 KiB ../node_modules/.pnpm/re2js@0.4.1/node_modules/re2js/build/index.esm.js
126.8 KiB ../node_modules/.pnpm/react-dom@18.3.1_react@18.3.1/node_modules/react-dom/cjs/react-dom.production.min.js
100.4 KiB src/lib/api.ts
98.5 KiB ../packages/quill/packages/quill/dist/index.js
93.3 KiB ../node_modules/.pnpm/prosemirror-view@1.40.1/node_modules/prosemirror-view/dist/index.js
90.6 KiB ../node_modules/.pnpm/@tiptap+core@3.20.6_@tiptap+pm@3.20.6/node_modules/@tiptap/core/dist/index.js
84.9 KiB src/products.tsx

Posted automatically by check-eager-graph · sizes are eager output bytes (shipped, post-tree-shake) from the esbuild metafile · part of #32479

✅ Toolbar bundle — eager 2.37 MiB within budget

What the toolbar ships to customer pages, measured from the esbuild output (minified, post-tree-shake). The eager set is the entry plus everything statically imported from it — fetched before any feature runs; deferred chunks load lazily. The eager guardrail is 5.72 MiB. Each output file must also stay below 10 MB, where CloudFront stops compressing it. The module boundary is enforced separately by check-toolbar-graph.

Metric Size Δ vs base Budget
Eager (shipped)
entry + static imports
2.37 MiB · 19 files 🔺 +556 B (+0.0%) ████░░░░░░ 41.4% of 5.72 MiB
Deferred (lazy) 2.10 MiB · 44 files no change n/a — loads on demand
Loader dist/toolbar.js 1.2 KiB no change █░░░░░░░░░ 6.0% of 19.5 KiB
Largest eagerly-shipped chunks
Size File
790.2 KiB dist/toolbar/toolbar-app-ZNXS4H6W.css
650.4 KiB dist/toolbar/chunk-chunk-DPTOFHTZ.js
483.6 KiB dist/toolbar/chunk-chunk-DPLA2GAY.js
138.3 KiB dist/toolbar/chunk-chunk-2EHPTPPA.js
131.8 KiB dist/toolbar/chunk-chunk-FDH2IBXT.js
75.2 KiB dist/toolbar/toolbar-app-FOEGQ3V3.js
69.0 KiB dist/toolbar/chunk-chunk-TSAL54PB.js
35.6 KiB dist/toolbar/chunk-chunk-U5CQIMDB.js
21.0 KiB dist/toolbar/chunk-chunk-OYNHVBHM.js
6.8 KiB dist/toolbar/chunk-chunk-DV7IWQNF.js

Posted automatically by check-toolbar-size · sizes are toolbar output bytes (shipped, post-tree-shake) from the esbuild metafile

✅ Dist folder size — 🟢 -871.16 MiB (-48.0%)

Total size of the built frontend/dist folder (all assets), compared against the base branch.

Total: 943.47 MiB · 🟢 -871.16 MiB (-48.0%)

✅ Hobby preview — passed

Hobby deployment smoke test passed successfully.


Run 36151367037

@greptile-apps

greptile-apps Bot commented Sep 25, 2026 •

Copy link
Copy Markdown
Contributor

Retrigger

[High risk] Refactors how the capture service routes events to Kafka topics.

The PR appears safe to merge, subject to the stated charts-first deployment order.

Reviews (1) · Last reviewed commit: "docs(capture): record steps 9 and 10 as ..."

@coderabbitai

coderabbitai Bot commented Sep 25, 2026 •

Copy link
Copy Markdown
Contributor

Review in Change Stack →

Navigate logical layers of code changes, visualize relationships, and explore their blast radius.

No actionable comments were generated in the recent review. 🎉

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Repository: PostHog/posthog/.coderabbit.yaml

Review profile: QUIET

Plan: Enterprise

Run ID: 70b2b2ef-082e-4adb-be9e-189f1e223d7f

📥 Commits

Reviewing files that changed from the base of the PR and between 8b451b8 and 0fbaf3d.

📒 Files selected for processing (8)
  • rust/capture/OUTPUTS_REFACTOR_PLAN.md
  • rust/capture/src/outputs.rs
  • rust/capture/src/producers.rs
  • rust/capture/src/server.rs
  • rust/capture/src/setup.rs
  • rust/capture/src/sinks/kafka.rs
  • rust/capture/src/sinks/registry.rs
  • rust/capture/src/sinks/sink.rs
💤 Files with no reviewable changes (2)
  • rust/capture/src/sinks/sink.rs
  • rust/capture/src/outputs.rs

Included review availability: Your plan provides up to 12 included reviews per hour; 11 remain after this review.


📝 Walkthrough

Walkthrough

Capture now configures topics and producer names per output. OutputTable resolves each destination to its topic and producer. The Kafka sink retains destinations in prepared payloads, selects a producer during enqueue, and flushes all registered producers. Startup validation, service environment variables, and test configuration now use the output configuration. The changes also add exact-name parsing for producer names.

Priority: ➖ Normal

Merge Risk: ⚪ Minimal · up to 0fbaf

The shutdown flush handles the currently configured producer. No actionable merge-blocking issue remains after normal checks.

Security Architecture Review

Security architecture risk: 🟡 Moderate · up to 0fbaf

Per-output routing is a meaningful design change. The main risk is a deployment that still supplies customized legacy topic settings: the new capture process ignores them and may publish to default topics instead. The production configuration change and deployment order could not be verified here.

Retained concerns

  • Medium · security · inferred: A deployment retaining only a customized legacy topic variable will silently use the new output default instead. Depending on the deployment’s topics and consumers, that can put capture data on an unintended stream; whether production is exposed depends on an unverified configuration rollout.
Security review details

Security Blast Radius

  • inferred — A wrong topic setting can affect whichever capture outputs use that setting and their downstream consumers. The evidence does not establish cross-tenant exposure or a second usable producer cluster in this version.

Security Findings and Attack Paths

  • inferred — No introduced sender-controlled producer-selection path was established. The material exposure is conditional configuration drift during cutover, not a verified external attack path.

Trust Boundaries and Controls

  • observed — Producer identity is parsed against a fixed allowlist, and startup constructs configured handles before exposing the sink. Output selection is resolved from a pipeline destination rather than a producer name in the event.

Resilience and Maintainability Implications

  • inferred — Default-compatible local settings limit the demonstrated migration risk, but they do not contain an incorrect custom production topic setting: a nonempty default passes the startup completeness check.

Hardening Proposals

  • proposed — Before deployment, compare each deployed legacy topic value with its new output value, retain settings needed for rollback, and verify the chart’s ordering across old and new images.
🚥 Pre-merge checks | ✅ 1
✅ Passed checks (1 passed)
Check name Status Explanation
Description check ✅ Passed The description is complete and stand-alone. It covers the problem, user-visible and mechanical changes, migration requirements, testing results and limitations, release status, documentation, and age…
✨ Finishing Touches
📝 Generate docstrings
  • Commit to this branch
  • Create a new PR
🧪 Generate unit tests (beta)
  • Commit to this branch
  • Create a new PR

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

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Actionable comments posted: 1


ℹ️ Review info
⚙️ Run configuration

Configuration used: Repository: PostHog/posthog/.coderabbit.yaml

Review profile: QUIET

Plan: Enterprise

Run ID: c3582124-8ff0-49f6-9f53-248c54a61228

📥 Commits

Reviewing files that changed from the base of the PR and between 1738ae8 and be89082.

📒 Files selected for processing (26)
  • bin/start-rust-service
  • docker-compose.base.yml
  • rust/capture/OUTPUTS_REFACTOR_PLAN.md
  • rust/capture/src/config.rs
  • rust/capture/src/events/analytics.rs
  • rust/capture/src/events/overflow_stamping.rs
  • rust/capture/src/events/recordings.rs
  • rust/capture/src/outputs.rs
  • rust/capture/src/overflow_parity.rs
  • rust/capture/src/producers.rs
  • rust/capture/src/router.rs
  • rust/capture/src/setup.rs
  • rust/capture/src/sinks/kafka.rs
  • rust/capture/src/sinks/registry.rs
  • rust/capture/src/sinks/sink.rs
  • rust/capture/src/v0_request.rs
  • rust/capture/src/v1/analytics/process.rs
  • rust/capture/src/v1/quota_limiter_shim.rs
  • rust/capture/src/v1/sinks/kafka/config.rs
  • rust/capture/src/v1/sinks/types.rs
  • rust/capture/src/v1/test_utils.rs
  • rust/capture/tests/common/utils.rs
  • rust/capture/tests/events.rs
  • rust/capture/tests/integration_person_processing_matrix.rs
  • rust/capture/tests/recordings.rs
  • rust/capture/tests/routing_e2e.rs

Included review availability: Your plan provides up to 12 included reviews per hour; 10 remain after this review.

Comment thread rust/capture/src/config.rs
@trunk-io

trunk-io Bot commented Sep 25, 2026 •

Copy link
Copy Markdown

Static Badge   Static Badge   Static Badge

View Full Report ↗︎ ⋅ Docs

pl added 3 commits September 25, 2026 14:54
…topics

Hobby and local dev pull capture:master with compose files that can predate
the CAPTURE_OUTPUT_* variables. With only the legacy names set, the new image
falls back to its defaults, and two of them differed from what those compose
files set: replay would land in events_plugin_ingestion and exceptions in
error_tracking_events, which nothing locally consumes.

The defaults are now session_recording_snapshot_item_events and
ingestion-errortracking-main, so an old compose file routes as before.
Production is unaffected: every deployment that reaches these outputs sets
them in charts.
OutputTable is generic over how a target names its producer: a ProducerName
in config, the producer handle in the Kafka sink. The sink maps names to
handles once at construction, so enqueue publishes through the target it
resolved, with no second lookup. The Producers wrapper is gone; the sink keeps
the distinct producers only to flush them.
Producers are created once by the ProducerRegistry and shared by every output
that names them, so flushing them at shutdown is the registry's job. The server
now flushes the registry setup hands it, instead of walking the output tree
down to the Kafka leaf.

The Kafka sink then holds only its output table, each target carrying its
producer handle, and flush comes off PublishEvents, Output, OutputRegistry and
the Sink trait. S3, print and noop already flushed nothing, and failover only
ever reached the Kafka primary.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant