Skip to content

[worker] Reconcile downstream processing (3/7) - #1867

Merged
Asherlc merged 12 commits into
mainfrom
Asherlc/issue-1852-reconciliation
Jul 22, 2026
Merged

Asherlc merged 12 commits into
mainfrom
Asherlc/issue-1852-reconciliation

Conversation

@Asherlc

@Asherlc Asherlc commented Jul 22, 2026 •

Copy link
Copy Markdown
Owner

Part 3 of the #1852 processing-status stack.

Adds CDC, analytics, and cache-refresh reconciliation, health checks, cache warming integration, and operational documentation for downstream processing.

Previous: #1866. Next: #1868.

Changed files: 14.

Summary by CodeRabbit

  • New Features

    • Added processing reconciliation for CDC operations, validating completion markers and metric acknowledgements before advancing analytics work.
    • Added detailed cache-refresh outcomes and per-dataset processing status tracking.
    • CDC health checks now report completed and waiting reconciliation operations.
  • Bug Fixes

    • Improved handling of malformed cache data by evicting invalid entries and treating them as cache misses.
  • Documentation

    • Updated CDC recovery mappings and local development commands.
    • Documented metric-stream acknowledgement handling and processing markers.
  • Tests

    • Added coverage for reconciliation, cache processing, and cache-warming outcomes.

@codereviewbot-ai

Copy link
Copy Markdown

🤖 Review skipped: Repository rate limit exceeded. Free accounts are limited to 2 reviews per 4 hours per repository. Upgrade to a paid plan for unlimited reviews.

@cursor

cursor Bot commented Jul 22, 2026

Copy link
Copy Markdown

Bugbot is not enabled for your account, so this pull request was not reviewed.

Enable Bugbot in the Cursor dashboard to get automatic reviews on future PRs.

@greptile-apps greptile-apps 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.

Your trial has ended. Reactivate Greptile to resume code reviews.

@sourcery-ai sourcery-ai 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.

Sorry @Asherlc, you have reached your weekly rate limit of 500000 diff characters.

Please try again later or upgrade to continue using Sourcery

@qodo-code-review

Copy link
Copy Markdown

Qodo reviews are paused for this user.

Troubleshooting steps vary by plan Learn more →

On a Teams plan?
Reviews resume once this user has a paid seat and their Git account is linked in Qodo.
Link Git account →

Using GitHub Enterprise Server, GitLab Self-Managed, or Bitbucket Data Center?
These require an Enterprise plan - Contact us
Contact us →

@coderabbitai

coderabbitai Bot commented Jul 22, 2026 •

Copy link
Copy Markdown

Review Change Stack

Warning

Review limit reached

@Asherlc, you've reached your PR review limit, so we couldn't start this review.

Next review available in: 46 minutes

Enable usage-based reviews in Billing to review now. Otherwise, wait until the next included review is available.
You're only billed for reviews past your plan's rate limits ($0.25/file).

How can I continue?

After more reviews become available, a review can be triggered using the @coderabbitai review command as a PR comment. Alternatively, push new commits to this PR.

To avoid repeated limits, reduce automatic review volume by pausing incremental auto-reviews earlier, using label-based review opt-in, excluding WIP or generated PR titles, or requesting reviews manually when the PR is ready. If your team needs uninterrupted high-volume reviews, an organization admin can enable usage-based reviews.

How do review limits work?

CodeRabbit enforces per-developer PR review limits for each organization. Most developers receive the normal plan review availability.

For paid Pro and Pro+ PR reviews, CodeRabbit uses adaptive limits for sustained high-volume activity. When a developer's recent PR review activity reaches the 95th percentile or higher among CodeRabbit users, additional reviews become available more gradually as earlier reviews age out of the rolling window.

Please refer docs for additional details.

Review details
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro Plus

Run ID: 5af32aa0-aa35-4aaa-b9ba-5f826c9d07cd

📥 Commits

Reviewing files that changed from the base of the PR and between f31ae4d and b3f4219.

📒 Files selected for processing (4)
  • docs/clickhouse-metric-stream.md
  • scripts/check-clickhouse-cdc.test.ts
  • scripts/check-clickhouse-cdc.ts
  • src/processing/processing-reconciler.test.ts
📝 Walkthrough

Walkthrough

Changes

Processing lifecycle

Layer / File(s) Summary
CDC evidence reconciliation
src/processing/processing-reconciler.ts, src/processing/processing-reconciler.test.ts, src/processing/processing-reconciler.integration.test.ts
Adds relational-marker and metric-acknowledgement reconciliation, CDC and analytics stage events, outbox completion, and unit/integration coverage.
Cache refresh stage recording
src/processing/cache-processing.ts, src/processing/cache-processing.test.ts
Maps cache outcomes to dataset-level succeeded, failed, or skipped events and records pending cache-refresh stages.
Cache warming outcome flow
scripts/warm-query-cache.ts, scripts/warm-query-cache.test.ts
Returns per-query outcomes, records cache runs with generated identifiers, and evaluates warming and processing failures together.
CDC health reconciliation wiring
scripts/check-clickhouse-cdc.ts, scripts/check-clickhouse-cdc.test.ts
Runs processing reconciliation after the ClickHouse health check and reports completed and waiting counts.
Operational documentation and cache test alignment
docs/clickhouse-cdc-health-runbook.md, docs/clickhouse-metric-stream.md, entrypoint.sh, packages/server/src/lib/cache-module.test.ts, packages/server/src/lib/cache-redis.test.ts
Updates marker and acknowledgement documentation, compose commands, analytics build invocation, and Redis cache-key assertions.

Estimated code review effort: 4 (Complex) | ~60 minutes

Sequence Diagram(s)

sequenceDiagram
  participant CDCHealth as check-clickhouse-cdc
  participant Reconciler as reconcilePendingProcessingOperations
  participant Postgres
  participant ClickHouse
  participant EventStore as Processing event store
  CDCHealth->>Reconciler: reconcile pending operations
  Reconciler->>Postgres: read expected processing evidence
  Reconciler->>ClickHouse: read markers and acknowledgements
  Reconciler->>EventStore: persist CDC and analytics stages
  Reconciler->>Postgres: complete reconciled outbox entries
Loading

Possibly related PRs

  • Asherlc/dofek#1315: Overlaps the analytics build command changes in entrypoint.sh.
  • Asherlc/dofek#1475: Covers related Redis malformed-payload eviction and cache-miss behavior.
  • Asherlc/dofek#1862: Introduces processing event ledger infrastructure used by the new reconciliation flows.

Suggested labels: area/server, type/feature

🚥 Pre-merge checks | ✅ 2
✅ Passed checks (2 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title is imperative, under 70 characters, properly prefixed with [worker], and accurately summarizes the downstream processing reconciliation work.

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 commented Jul 22, 2026 •

Copy link
Copy Markdown
Contributor

Storybook previews for 48e1d470 are ready:

This comment updates automatically on each PR push.

@codereviewbot-ai

Copy link
Copy Markdown

🤖 Review skipped: Repository rate limit exceeded. Free accounts are limited to 2 reviews per 4 hours per repository. Upgrade to a paid plan for unlimited reviews.

@greptile-apps greptile-apps 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.

Your trial has ended. Reactivate Greptile to resume code reviews.

@Asherlc Asherlc changed the title Processing status 3/7: downstream reconciliation [worker] Reconcile downstream processing (3/7) Jul 22, 2026
@codereviewbot-ai

Copy link
Copy Markdown

🤖 Review skipped: Repository rate limit exceeded. Free accounts are limited to 2 reviews per 4 hours per repository. Upgrade to a paid plan for unlimited reviews.

@greptile-apps greptile-apps 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.

Your trial has ended. Reactivate Greptile to resume code reviews.

Base automatically changed from Asherlc/issue-1852-metric-ack to main July 22, 2026 22:55
…conciliation-v1

# Conflicts:
#	src/db/clickhouse-cdc.test.ts
#	src/db/clickhouse-cdc.ts
#	src/jobs/process-import-job.test.ts
#	src/jobs/process-sync-job.test.ts
#	src/jobs/process-sync-job.ts
#	src/processing/metric-stream-processing-publisher.test.ts
#	src/processing/metric-stream-processing-publisher.ts
Copilot AI review requested due to automatic review settings July 22, 2026 22:59
@codereviewbot-ai

Copy link
Copy Markdown

🤖 Review skipped: Repository rate limit exceeded. Free accounts are limited to 2 reviews per 4 hours per repository. Upgrade to a paid plan for unlimited reviews.

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

@Asherlc

Asherlc commented Jul 22, 2026

Copy link
Copy Markdown
Owner Author

@greptile review

@greptile-apps

greptile-apps Bot commented Jul 22, 2026 •

Copy link
Copy Markdown

Greptile Summary

This PR adds reconciliation for downstream processing. The main changes are:

  • Validates relational CDC markers and metric-stream acknowledgements.
  • Queues analytics work after required CDC evidence arrives.
  • Records cache-refresh outcomes for pending datasets.
  • Integrates reconciliation into health checks and cache warming.
  • Updates operational documentation and tests.

Confidence Score: 5/5

This looks safe to merge.

  • Dataset identity is included in analytics idempotency keys.
  • Relational markers must match the expected source watermark.
  • Tests cover the corrected reconciliation paths.
  • No blocking issue remains in the updated code.

Important Files Changed

Filename Overview
src/processing/processing-reconciler.ts Adds CDC evidence reconciliation and dataset-scoped analytics queueing.
src/processing/cache-processing.ts Maps cache-warming results to pending dataset processing events.
scripts/warm-query-cache.ts Collects per-query outcomes and records cache-refresh processing status.
scripts/check-clickhouse-cdc.ts Runs processing reconciliation after CDC health validation.
entrypoint.sh Uses the centralized analytics build before warming query caches.

Flowchart

%%{init: {'theme': 'neutral'}}%%
flowchart TD
  A[Pending processing operation] --> B{CDC output path}
  B -->|Relational| C[Validate marker and source watermark]
  B -->|Metric stream| D[Validate batch acknowledgement]
  C --> E{Evidence complete?}
  D --> E
  E -->|No| F[Keep CDC running]
  E -->|Yes| G[Record CDC success]
  G --> H[Queue dataset analytics]
  H --> I[Run analytics build]
  I --> J[Warm query caches]
  J --> K[Record cache-refresh outcome]
Loading

Reviews (4): Last reviewed commit: "fix: clarify reconciliation review behav..." | Re-trigger Greptile

Comment thread src/processing/processing-reconciler.ts Outdated
Comment thread scripts/warm-query-cache.ts
@greptile-apps

greptile-apps Bot commented Jul 22, 2026

Copy link
Copy Markdown

Greptile Summary

This PR adds reconciliation for downstream processing. The main changes are:

  • Matches Postgres processing expectations with ClickHouse CDC evidence.
  • Queues analytics after required CDC outputs complete.
  • Records cache-warming outcomes for pending datasets.
  • Integrates reconciliation into health checks and worker startup.
  • Adds tests and operational documentation.

Confidence Score: 4/5

CDC and cache reconciliation can publish incorrect downstream state, so these paths need fixes before merging.

  • Relational completion does not verify the expected source watermark.
  • Shared cache query families can fail unrelated pending datasets.
  • The surrounding health-check and startup integrations appear consistent with the new workflow.

src/processing/processing-reconciler.ts and src/processing/cache-processing.ts

Important Files Changed

Filename Overview
src/processing/processing-reconciler.ts Adds CDC evidence matching, analytics queueing, and outbox completion; relational matching does not verify the source watermark.
src/processing/cache-processing.ts Maps cache outcomes to pending datasets, but shared query families can create incorrect cross-dataset failures.
scripts/warm-query-cache.ts Adds per-query outcomes and records cache-refresh processing events before reporting failures.
scripts/check-clickhouse-cdc.ts Runs pending processing reconciliation after the ClickHouse CDC health check.
entrypoint.sh Uses the analytics build runner before warming query caches.

Flowchart

%%{init: {'theme': 'neutral'}}%%
flowchart LR
  P[Pending operation] --> R[CDC reconciliation]
  C[(ClickHouse evidence)] --> R
  R -->|Evidence complete| A[Queue analytics]
  A --> B[Run analytics build]
  B --> W[Warm query caches]
  W --> S[Record cache status]
Loading

Fix All in Codex Fix All in Claude Code Fix All in Cursor Fix All in Conductor

Prompt To Fix All With AI
Fix the following 2 code review issues. Work through them one at a time, proposing concise fixes.

---

### Issue 1 of 2
src/processing/processing-reconciler.ts:162
**Source Watermark Is Not Verified**

A ClickHouse marker with the expected operation, dataset, flow, and batch key is accepted even when its `source_watermark` differs from the Postgres expectation. A replayed or inconsistent marker can therefore queue analytics and publish the expected watermark as serving evidence before that exact source position is available.

### Issue 2 of 2
src/processing/cache-processing.ts:45-48
**Shared Families Cross Dataset Boundaries**

Cache outcomes are assigned by query-family name, but `providerDetail` belongs to both the activity and providers contracts. When that query fails for a user with both datasets pending, this code records `cache_refresh_failed` for both operations even if only one dataset owns the affected refresh.

Reviews (2): Last reviewed commit: "Merge remote-tracking branch 'origin/mai..." | Re-trigger Greptile

Comment thread src/processing/processing-reconciler.ts Outdated
Comment thread src/processing/cache-processing.ts
@codereviewbot-ai

Copy link
Copy Markdown

🤖 Review skipped: Repository rate limit exceeded. Free accounts are limited to 2 reviews per 4 hours per repository. Upgrade to a paid plan for unlimited reviews.

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

🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@docs/clickhouse-metric-stream.md`:
- Around line 139-150: The documentation uses incorrect “casual fence”
terminology and leaves the new acknowledgement and mirror-behavior claims
uncited. Update the relevant text in clickhouse-metric-stream.md to use “causal
fence,” and add links to the defining migration, reconciliation code, or other
official primary sources for ingest.metric_stream_processing_acknowledgement and
postgres_fitness behavior.

In `@entrypoint.sh`:
- Line 27: Update the analytics build command in entrypoint.sh to invoke
scripts/run-analytics-build.ts through pnpm tsx instead of the local Node shim,
preserving the existing command chaining behavior.

In `@scripts/check-clickhouse-cdc.test.ts`:
- Around line 205-208: Strengthen the test around
mockedReconcilePendingProcessingOperations by asserting the exact configured
ClickHouse client and database mock instances instead of expect.any(Object).
Retain the mock return values and add an invocation-order assertion proving
reconciliation occurs after assertClickHouseCdcHealth.

In `@scripts/check-clickhouse-cdc.ts`:
- Around line 126-128: Update the reconciliation status log near
reconcilePendingProcessingOperations to include reconciliation.checked and
clearly identify completed and waiting as counts from the current batch,
preventing them from being interpreted as the total backlog.

In `@src/processing/processing-reconciler.test.ts`:
- Around line 255-258: Reorder the mocked outputs in the processing reconciler
test to match the production ORDER BY result, placing the metric_stream entry
before the relational entry. Keep the existing idempotencyKey assertion aligned
with the resulting reconcileProcessingEvidence order and production-generated
key.

In `@src/processing/processing-reconciler.ts`:
- Around line 359-369: Update the pendingOperations claim query in the
processing reconciler to use PostgreSQL row-level locking with FOR UPDATE SKIP
LOCKED, preserving the existing filtering, grouping, ordering, and limit
behavior so concurrent runs claim disjoint pending operations.
🪄 Autofix (Beta)

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro Plus

Run ID: c8255997-b326-4c73-95d5-cd0f53e04eb9

📥 Commits

Reviewing files that changed from the base of the PR and between dc40569 and f31ae4d.

📒 Files selected for processing (14)
  • docs/clickhouse-cdc-health-runbook.md
  • docs/clickhouse-metric-stream.md
  • entrypoint.sh
  • packages/server/src/lib/cache-module.test.ts
  • packages/server/src/lib/cache-redis.test.ts
  • scripts/check-clickhouse-cdc.test.ts
  • scripts/check-clickhouse-cdc.ts
  • scripts/warm-query-cache.test.ts
  • scripts/warm-query-cache.ts
  • src/processing/cache-processing.test.ts
  • src/processing/cache-processing.ts
  • src/processing/processing-reconciler.integration.test.ts
  • src/processing/processing-reconciler.test.ts
  • src/processing/processing-reconciler.ts

Comment thread docs/clickhouse-metric-stream.md Outdated
Comment thread entrypoint.sh
Comment thread scripts/check-clickhouse-cdc.test.ts
Comment thread scripts/check-clickhouse-cdc.ts
Comment thread src/processing/processing-reconciler.test.ts
Comment thread src/processing/processing-reconciler.ts
@codereviewbot-ai

Copy link
Copy Markdown

🤖 Review skipped: Repository rate limit exceeded. Free accounts are limited to 2 reviews per 4 hours per repository. Upgrade to a paid plan for unlimited reviews.

@Asherlc
Asherlc merged commit 8ec133b into main Jul 22, 2026
113 checks passed
@Asherlc
Asherlc deleted the Asherlc/issue-1852-reconciliation branch July 22, 2026 23:48
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.

2 participants