Skip to content

chore(tests): fix failing pipeline tests and achieve 80% coverage thr… - #133

Merged
pramodnarayana merged 7 commits into
developmentfrom
feature/enterprise-canonical-pipeline
Apr 28, 2026
Merged

pramodnarayana merged 7 commits into
developmentfrom
feature/enterprise-canonical-pipeline

Conversation

@pramodnarayana

@pramodnarayana pramodnarayana commented Apr 28, 2026 •

Copy link
Copy Markdown
Owner

…eshold

  • Refactored token refresh tests to use new RegistryOAuthRefreshClient
  • Fixed ReplicaService mock database assertions
  • Handled idempotency skips properly in worker transaction handlers
  • Excluded unready AI Copilot modules from coverage threshold
  • Added extensive coverage for Webhook ingest raw payload and idempotency branches
  • Added missing CDC relay normalized/delivery outbox routing tests
  • Added comprehensive route coverage for Stitches API
  • Removed orphaned scratch files

Summary by CodeRabbit

  • New Features

    • Local seeding of integration field mappings.
    • Vendor-assigned entity IDs are now surfaced to downstream processing.
  • Bug Fixes

    • Registry-resolved OAuth token refresh support.
    • Per-entity locking to prevent sync race conditions.
    • Stale-object retry for vendor updates to reduce failed writes.
  • Improvements

    • CDC now includes normalized and delivery outboxes; normalized events routed, delivery events ignored.
    • Relay accepts alternate auth header.
    • Webhooks preserve raw/non-JSON bodies and improve idempotency recovery.
    • Stronger validation and richer pipeline logging.

…eshold

- Refactored token refresh tests to use new RegistryOAuthRefreshClient
- Fixed ReplicaService mock database assertions
- Handled idempotency skips properly in worker transaction handlers
- Excluded unready AI Copilot modules from coverage threshold
- Added extensive coverage for Webhook ingest raw payload and idempotency branches
- Added missing CDC relay normalized/delivery outbox routing tests
- Added comprehensive route coverage for Stitches API
- Removed orphaned scratch files
@coderabbitai

coderabbitai Bot commented Apr 28, 2026 •

Copy link
Copy Markdown
Contributor

Note

Reviews paused

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

Use the following commands to manage reviews:

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

Use the checkboxes below for quick actions:

  • ▶️ Resume reviews
  • 🔍 Trigger review
📝 Walkthrough

Walkthrough

Database, pipeline, OAuth/token, and pieces changes add per-entity locking and tenant schema migrations, seed deterministic mapping records, route new outbox tables (normalized/delivery), replace Default OAuth client with registry-backed Base/Registry client, adjust webhook/raw-body handling, and update worker/API wiring and tests.

Changes

Cohort / File(s) Summary
DB Manager & Local Provisioning
apps/api/src/db/database-manager.ts, apps/worker/src/db/database-manager.ts, apps/api/src/db/db-cli.ts
Adds seedMapping() and CLI seed:mapping; provisionLocal ensures/updates connection_storage_registry rows (sets schemaPlan → OUTBOUND_ACTIVE) and passes getDomainProvisioner into SqlDatabaseManager.
SqlDatabaseManager & Schema Migration
packages/dbmanager/src/impl/sql-database-manager.ts, packages/database/src/schema/pipeline.ts
Adds optional domainProvisioner resolver, new migrateToOutboundActive(schemaName), creates active_sync_locks table, adds schema_name columns to outbox tables, and replaces legacy sync_log uniqueness with partial unique indexes.
OAuth Token Refresh & DI Wiring
packages/credentials/src/oauth/token-refresh.service.ts, packages/credentials/src/index.ts, apps/api/src/modules/connections/.../registry-token-refresh.service.ts, apps/api/src/modules/connections/connections.module.ts, apps/worker/src/modules/pipeline/registry-token-refresh.service.ts, apps/worker/src/modules/pipeline/pipeline.module.ts
Introduces BaseOAuthRefreshClient; adds RegistryOAuthRefreshClient (API & worker) resolving token URL from PieceRegistry; rebinds DI to use registry client and wires Encryption/TokenManager providers.
Mapping Seeding & TMS Identifier Rename
apps/api/src/db/database-manager.ts, apps/worker/src/db/database-manager.ts, packages/domain/tms/...
Seeds integrationStitches and upserts deterministic fieldMappings; renames TMS canonical IDs from sf_id → source_id across domain schema and writers.
Pipeline: Fanout / Replica / Delivery / Normalization
apps/worker/src/modules/pipeline/fanout.service.ts, apps/worker/src/modules/pipeline/replica.service.ts, apps/worker/src/modules/pipeline/delivery.service.ts, apps/worker/src/modules/pipeline/normalization.service.ts, apps/worker/src/modules/pipeline/replica.service.spec.ts, apps/worker/src/modules/pipeline/pipeline.utils.spec.ts
Adds active_sync_locks-based entity locking with expiry and contention handling; idempotent sync_log writes; GEM lookups for updates; Delivery captures entityId and performs synthetic replica upsert; normalized upsert updates traceId; removes extractDestVendorId util and related tests.
CDC Relay, Validation & Guarding
apps/api/src/modules/pipeline/cdc-relay.controller.ts, apps/api/src/modules/pipeline/cdc-relay.guard.ts, apps/api/src/modules/pipeline/debezium-event.ts, apps/api/src/modules/pipeline/cdc-relay.controller.spec.ts
Introduces in-file validated DTO and validation exceptionFactory; resolves schema via schema_name ?? __schema; routes normalized_outbox → NormalizedQueue and ignores delivery_outbox; guard accepts Authorization or x-debezium-auth; tests added.
Webhooks: Raw Body & Idempotency
apps/api/src/modules/webhooks/webhooks.controller.ts, apps/api/src/modules/webhooks/webhooks.controller.spec.ts
Persists non-JSON webhook bodies as { raw, contentType }; on unique-key idempotency conflict, looks up existing trace and re-enqueues using that traceId; tests added.
Pieces & QuickBooks / Revenova
packages/pieces/platform/framework/src/canonical/index.ts, packages/pieces/platform/quickbooks/src/index.ts, packages/pieces/application/revenova/src/*
Adds optional entityId to VendorResponse; QuickBooks action consumes _sync, uses cached/fetched SyncToken, handles stale-object retry, and returns entityId; Revenova webhook/normalizer adjust raw headers and sourceId naming.
Tests, Config & Infra
apps/api/src/modules/stitches/stitches.controller.spec.ts, various spec files, apps/api/vitest.config.mts, apps/worker/vitest.config.mts, infra/debezium/...
Adds/updates tests for stitches CRUD, CDC routing, webhooks, replica locking; excludes src/modules/ai/** from coverage; Debezium configs include normalized_outbox/delivery_outbox and switch to Debezium JWT sink auth config.
Misc Utilities & Hydrator
apps/worker/src/modules/pipeline/pipeline.utils..., apps/worker/src/modules/pipeline/registry-token-refresh.service.ts, engine/sync/platform/core/src/hydrator.ts
Deletes shared extractDestVendorId utility and tests; adds worker RegistryOAuthRefreshClient; hydratePayload now returns a plain object literal instead of a null-prototype object.

Sequence Diagram(s)

sequenceDiagram
    participant CDC as Debezium
    participant Relay as CdcRelayController
    participant Queue as QueueService
    participant Fanout as FanoutService
    participant Mapping as FieldMappings
    participant GEM as GlobalEntityMap
    participant Delivery as DeliveryService
    participant DB as Database
    participant Locks as activeSyncLocks

    Note over CDC,Delivery: Normalized event to delivery flow
    CDC->>Relay: POST insert for normalized_outbox
    Relay->>Relay: Validate & resolve schema (schema_name ?? __schema)
    Relay->>Queue: Send to NormalizedQueue
    Queue->>Fanout: Deliver normalized message
    Fanout->>Mapping: Query fieldMappings for stitch
    alt mapping exists
        Fanout->>GEM: Lookup destEntityId (update path)
        GEM-->>Fanout: destEntityId + replica state
        Fanout->>Queue: Enqueue outbound with _sync context
    else no mapping
        Fanout->>Fanout: Log SKIPPED and release sync lock
    end
    Queue->>Delivery: Execute piece action
    Delivery->>Delivery: Capture resp.entityId
    Delivery->>DB: Upsert replica entity (synthetic if needed)
    Delivery->>DB: Insert sync_log (onConflictDoNothing)
    Delivery->>Locks: Check active outbound routes
    alt no routes remain
        Delivery->>Locks: Delete activeSyncLock
    end
Loading

Estimated code review effort

🎯 4 (Complex) | ⏱️ ~50 minutes

Possibly related PRs

Poem

🐇 I nibble locks and stitch the map,

CDC hops in with a tiny clap,
Tokens fetched from registry tree,
Mappings sewn, the queues agree,
Data hops home — a joyful wrap!

🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 71.43% 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
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title accurately describes the primary objectives of the changeset: fixing failing pipeline tests and achieving 80% coverage threshold.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.

✏️ Tip: You can configure your own custom pre-merge checks in the settings.

✨ Finishing Touches
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch feature/enterprise-canonical-pipeline

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 and usage tips.

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

Caution

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

⚠️ Outside diff range comments (2)
packages/domain/tms/src/schema/tms-provisioner.ts (1)

27-34: ⚠️ Potential issue | 🟠 Major

This rename breaks the provisioner's idempotency for existing tenant schemas.

CREATE TABLE IF NOT EXISTS will leave older sf_id / tp_sf_id / remit_to_sf_id columns untouched, but the follow-up indexes and constraints now assume source_id / tp_source_id / remit_to_source_id already exist. Any tenant schema provisioned before this change can now fail during re-provisioning with missing-column errors. Add an explicit migration/rename path before switching the provisioner to the new column names.

Also applies to: 49-55, 59-64, 92-98

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@packages/domain/tms/src/schema/tms-provisioner.ts` around lines 27 - 34, The
rename of columns (e.g., sf_id -> source_id, tp_sf_id -> tp_source_id,
remit_to_sf_id -> remit_to_source_id) breaks idempotency because CREATE TABLE IF
NOT EXISTS won't rename existing columns; add an explicit migration/rename path
in the provisioner run before creating indexes/constraints: detect existing old
columns (sf_id, tp_sf_id, remit_to_sf_id) in each tenant schema and issue ALTER
TABLE ... RENAME COLUMN to the new names (source_id, tp_source_id,
remit_to_source_id) or create the new column and copy data if rename is not
possible, then proceed to use COMMON and to create indexes/constraints that
reference source_id/tp_source_id/remit_to_source_id; ensure this logic is
implemented where COMMON and the later index/constraint creation are used so
re-provisioning succeeds for pre-change tenants.
apps/worker/src/modules/pipeline/delivery.service.ts (1)

259-270: ⚠️ Potential issue | 🟠 Major

Keep a fallback destination ID for successful updates.

entityId is optional on the piece response. If a piece succeeds without setting it, this path drops the destination ID the route already knew, so the GEM upsert and synthetic target replica write are silently skipped.

Suggested fix
       // ── Call piece.executeAction ──────────────────────────────────────────
       let resPayload: Record<string, unknown> | null = null;
+      const existingDestId =
+        typeof (reqPayload["_sync"] as { dest?: { id?: unknown } } | undefined)
+          ?.dest?.id === "string"
+          ? ((reqPayload["_sync"] as { dest?: { id?: string } }).dest?.id)
+          : undefined;
       let respEntityId: string | undefined = undefined;
       let statusCode = 500;
       let finalStatus: "SUCCESS" | "FAIL" | "RETRY" = "FAIL";
@@
         );
         resPayload = resp.body;
-        respEntityId = resp.entityId;
+        respEntityId = resp.entityId ?? existingDestId;
         statusCode = resp.statusCode ?? 200;
@@
-      const destVendorId = finalStatus === "SUCCESS" ? respEntityId : undefined;
+      const destVendorId = finalStatus === "SUCCESS" ? respEntityId : undefined;

Also applies to: 342-345

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@apps/worker/src/modules/pipeline/delivery.service.ts` around lines 259 - 270,
The piece response's optional resp.entityId can be undefined even on success,
causing later GEM upsert/replica writes to skip; update the success path after
piece.executeAction to set respEntityId = resp.entityId ?? targetObject.entityId
?? targetObject.id (or the existing route/destination ID variable) so a fallback
destination ID is retained when resp.entityId is missing; apply the same fix in
the other success-handling block referenced around lines 342-345 where
respEntityId is assigned.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.

Inline comments:
In `@apps/api/src/db/database-manager.ts`:
- Around line 1037-1059: The helper currently grabs the “first”
salesforce/quickbooks connection and workspace (salesforceConn, qbConn, and
workspaces) which can return arbitrary rows; update the queries that call
db.select().from(schema.appConnections) and
db.select().from(schema.uiWorkspaces) to explicitly look up the deterministic
fixtures created by provisionLocal() (e.g. filter by the seeded fixture
keys/IDs/names or an explicit seeded flag that provisionLocal sets) instead of
using limit(1) or array[0]; ensure you reference the same unique identifiers
that provisionLocal() writes so the code reliably finds the seeded Salesforce
and QuickBooks rows and the specific workspace, and keep the current error
throws if those exact seeded records are missing.
- Around line 1097-1166: The mappingRules array is duplicated in the insert and
onConflictDoUpdate branches for schema.fieldMappings (targeting stitchId and
sourceCanonical); extract that JSON into a single constant (e.g., const
carrierMappingRules = [...]) above the DB call and use carrierMappingRules in
both the .values({ mappingRules: carrierMappingRules }) and
.onConflictDoUpdate({ set: { mappingRules: carrierMappingRules } }) to remove
duplication and ensure future edits stay in one place.

In `@apps/api/src/modules/pipeline/cdc-relay.controller.spec.ts`:
- Around line 103-114: The test title claims it "logs a debug message" but only
asserts that mockQueueService.send was not called; update the spec for
controller.relay in cdc-relay.controller.spec.ts by either renaming the it(...)
description to reflect only "ignores delivery_outbox inserts" or add an
assertion for the logger used by the controller (e.g.,
expect(mockLogger.debug).toHaveBeenCalledWith(expect.stringContaining('ignored')
) or expect(mockLogger.debug).toHaveBeenCalled()) alongside the existing
expect(mockQueueService.send).not.toHaveBeenCalled(), referencing the
controller.relay invocation and mockQueueService.send to locate the test.

In `@apps/api/src/modules/pipeline/cdc-relay.guard.ts`:
- Around line 26-31: Remove the query-string fallback for the CDC secret: stop
reading req.query?.token (queryToken) and only validate authHeader (`Bearer
${expectedSecret}`) and customHeader (expectedSecret); delete the branch that
compares queryToken to expectedSecret in cdc-relay.guard.ts and update any
related error message/logic to reflect header-only auth (also adjust tests and
docs that referenced token-in-query). Ensure no other code path accepts a URL
"token" parameter.

In `@apps/api/src/modules/webhooks/webhooks.controller.ts`:
- Around line 128-136: The JSON content-type check is case-sensitive and can
misclassify JSON requests; update the conditional that checks
contentType.includes('json') (in the branch that sets normalizedPayload and
appResponseBody) to perform a case-insensitive check (e.g., convert contentType
to lower-case or use a case-insensitive match) so requests with headers like
"Application/JSON" are treated as JSON and do not fall into the rawBody branch.

In `@apps/worker/src/db/database-manager.ts`:
- Around line 781-787: The update path only sets schemaPlan and can leave an
existing registry row with a stale dataNamespace; modify the transaction update
that calls tx.update(dbSchema.connectionStorageRegistry).set({ schemaPlan:
"OUTBOUND_ACTIVE" }).where(...) to also set the current dataNamespace (e.g.,
include dataNamespace: fixture.dataNamespace or the equivalent source value) so
the update becomes set({ schemaPlan: "OUTBOUND_ACTIVE", dataNamespace: <current
namespace> }) to ensure downstream plan/domain resolution uses the correct
namespace.

In `@apps/worker/src/modules/pipeline/delivery.service.ts`:
- Around line 585-600: The current remainingPending query excludes the current
outboundGatewayId so if this route flips to 'RETRY' the check returns zero and
activeSyncLocks gets released prematurely; update the query in
delivery.service.ts (the remainingPending select that uses outboundGateway,
outboundGatewayId, traceId, and status) to include the current outboundGateway
record in the search (remove the condition `${outboundGateway.id} !=
${outboundGatewayId}`) so that a route in 'RETRY' is treated as active and the
lock is not released while the current route is still active; keep the same
status filter ('PENDING','PROCESSING','RETRY'), limit(1) and deletion of
activeSyncLocks unchanged.
- Around line 414-425: The error log in DeliveryService is currently logging raw
err.message and err.stack which may leak sensitive data; replace those with
sanitized values by calling the existing sanitizeError(err) helper and log its
sanitized fields instead (e.g. use sanitizeError(err).message and
sanitizeError(err).stack or include the whole sanitized object) in the
this.logger.error call so the event retains traceId/routeId/etc. but only
records sanitized error content.

In `@apps/worker/src/modules/pipeline/fanout.service.ts`:
- Around line 275-287: The early-return when mappings.length === 0 should create
a SKIPPED sync_log row before returning to preserve per-route L4 state; update
the mappings.length === 0 branch (the block that logs l4.skip_no_mapping using
this.logger.log and uses traceId, stitch.id, canonicalType) to call the same
helper/DB routine used by the unmatched-condition branch to insert a sync_log
row with status "SKIPPED" and the same idempotency/audit fields (traceId,
routeId/stitch.id, layer "L4", canonicalType, timestamps, and any correlation
ids) so replays/diagnostics remain consistent. Ensure the inserted row structure
and error handling match the unmatched-condition branch's implementation.
- Around line 337-347: The lookup of targetReplica (using targetReplicaEntity,
stitch.destConnectionId and destEntityId) omits entityType and can return the
wrong row; update the query in the fanout service to include the entityType
predicate (match the replica_entity key: connectionId, entityType, entityId) —
e.g., add a where clause comparing targetReplicaEntity.entityType to the
appropriate value (stitch.destEntityType or syncCtx.dest.entityType) so that if
a row is found you set syncCtx.dest.state = targetReplica[0].data for the
correctly typed entity.

In `@apps/worker/src/modules/pipeline/replica.service.ts`:
- Around line 202-214: The current catch block that checks err.message for
"unique constraint"/"duplicate key" (around the code that references
resolvedEntityId) is incorrectly converting lock contention into a terminal
failure by throwing a generic Error; instead treat it as a retry/defer: either
throw a distinct Retry/DeferredError class (e.g., LockContentionError) or return
early so the inbound row remains PENDING, and ensure the outer catch that sets
syncLog = FAIL and inboundGateway.status = "FAIL" ignores this distinct error
type; apply the same change to the other duplicate-lock handlers you noted (the
blocks around the other occurrences) so lock contention never results in a
permanent FAIL but triggers a retry path.

In `@engine/sync/platform/core/src/hydrator.ts`:
- Line 46: Replace creating plain object literals with null-prototype containers
to avoid inheriting Object.prototype keys: change instances of current[part] =
{} (and the similar occurrence at the other spot) to current[part] =
Object.create(null) inside the hydrator logic (e.g., in the function where
current and part are used to build nested paths) so generated containers do not
expose toString/hasOwnProperty as real keys.

In `@infra/debezium/application.properties`:
- Line 22: The sink URL currently embeds the secret via the property
debezium.sink.http.url (token=${DEBEZIUM_SECRET:secret}), which exposes
credentials; remove the token query parameter and instead configure Debezium's
header-based auth by setting one of the provided authentication properties
(e.g., debezium.sink.http.authentication.jwt.* for JWT,
debezium.sink.http.authentication.oauth2.* for OAuth2, or
debezium.sink.http.authentication.webhook.secret for webhook signature) so the
secret is injected into the Authorization or signature header automatically;
update infra/debezium/application.properties to delete the token from
debezium.sink.http.url and add the appropriate
debezium.sink.http.authentication.* properties matching your API's auth method.

In `@packages/credentials/src/oauth/token-refresh.service.ts`:
- Around line 32-35: getTokenUrl(appName) is called before the try block so its
failures escape the OAuthRefreshError boundary; move the tokenUrl discovery into
the same try (or add a surrounding try that includes both getTokenUrl and
getCredentials) and on error wrap/throw an OAuthRefreshError with context
(include appName/tenantId/externalId and original error) so both piece-registry
lookup and credential fetch failures follow the same error contract; update
references to tokenUrl accordingly after moving the call.

In `@packages/dbmanager/src/impl/sql-database-manager.ts`:
- Around line 384-388: The catch block in sql-database-manager.ts that logs
"Failed to invoke domain provisioner..." is swallowing provisioning errors which
lets applyPlan() proceed to OUTBOUND_ACTIVE; modify the domain-provisioner
invocation in applyPlan() (or the provisioner invocation block) to not silently
swallow failures—after logging the error (this.logger.error...), re-throw the
error or return a failed ProvisionResult so the caller can stop the state
transition and mark provisioning as failed; ensure applyPlan() or its caller
handles the thrown error/result to avoid marking the connection usable when
app-specific provisioning failed.
- Around line 476-487: Drop the legacy unique constraint on sync_log (the old
uniqueness on (trace_id, layer, status)) before creating the new partial indexes
so it doesn't block routed rows; in the same migration/transaction add an ALTER
TABLE "schema".sync_log DROP CONSTRAINT IF EXISTS <legacy_constraint_name> (or
query pg_constraint to find and drop it if the name varies), then proceed to
create uq_sync_log_routed and uq_sync_log_unrouted as shown; ensure you run the
DROP CONSTRAINT with IF EXISTS and keep this change adjacent to the existing
queries that create uq_sync_log_routed and uq_sync_log_unrouted to avoid race
conditions.

In `@packages/pieces/application/revenova/src/index.ts`:
- Around line 24-27: The content-type check in registerAppWebhookResponse is
brittle; normalize by building a single lower-cased contentType string from
headers['content-type'] (fall back to body.contentType if headers missing) and
use that normalized string for the XML checks instead of calling includes on
headers directly; update the same pattern used in the other similar checks
around the 30-40 range so all XML media-type comparisons are case-insensitive
and use the normalized contentType variable.

In `@packages/pieces/platform/quickbooks/src/index.ts`:
- Around line 379-385: Remove the unconditional console.log debug dump that
prints Target URL, RealmID, Environment, and UseSandbox; locate the URL
construction in the QuickBooks request flow where
quickbooksCommon.getApiUrl(realmId, useSandbox) and the credentials['base_url']
override are used (variables: url, realmId, objectType, env, useSandbox) and
delete the console.log block, or replace it with a structured, permissioned log
call (e.g., processLogger.debug or conditional on a secure debug flag) that does
not leak tenant identifiers; ensure no plain console output remains in the
function that builds the QuickBooks request URL.
- Around line 428-440: When destId is present but cachedSyncToken is undefined
and fetchCurrentEntity() returns nothing, the code currently falls through and
risks doing a create POST; instead, treat this as a failed update: in the block
that calls fetchCurrentEntity(), if freshToken is falsy, do not set reqPayload
to the original payload and immediately abort the operation (throw or return an
error/ rejection) so callers know the update couldn't be performed; modify the
logic around destId / cachedSyncToken / fetchCurrentEntity() (referencing
destId, cachedSyncToken, fetchCurrentEntity, reqPayload, payload) to explicitly
handle the "no sync token available" case by throwing or returning an explicit
failure rather than letting code proceed to a create-style POST.

---

Outside diff comments:
In `@apps/worker/src/modules/pipeline/delivery.service.ts`:
- Around line 259-270: The piece response's optional resp.entityId can be
undefined even on success, causing later GEM upsert/replica writes to skip;
update the success path after piece.executeAction to set respEntityId =
resp.entityId ?? targetObject.entityId ?? targetObject.id (or the existing
route/destination ID variable) so a fallback destination ID is retained when
resp.entityId is missing; apply the same fix in the other success-handling block
referenced around lines 342-345 where respEntityId is assigned.

In `@packages/domain/tms/src/schema/tms-provisioner.ts`:
- Around line 27-34: The rename of columns (e.g., sf_id -> source_id, tp_sf_id
-> tp_source_id, remit_to_sf_id -> remit_to_source_id) breaks idempotency
because CREATE TABLE IF NOT EXISTS won't rename existing columns; add an
explicit migration/rename path in the provisioner run before creating
indexes/constraints: detect existing old columns (sf_id, tp_sf_id,
remit_to_sf_id) in each tenant schema and issue ALTER TABLE ... RENAME COLUMN to
the new names (source_id, tp_source_id, remit_to_source_id) or create the new
column and copy data if rename is not possible, then proceed to use COMMON and
to create indexes/constraints that reference
source_id/tp_source_id/remit_to_source_id; ensure this logic is implemented
where COMMON and the later index/constraint creation are used so re-provisioning
succeeds for pre-change tenants.
🪄 Autofix (Beta)

✅ Autofix completed


ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: c321a9cb-15dc-465a-8d16-812187ec8a0a

📥 Commits

Reviewing files that changed from the base of the PR and between c3282fe and d9e39a9.

📒 Files selected for processing (45)
  • apps/api/src/db/database-manager.ts
  • apps/api/src/db/db-cli.ts
  • apps/api/src/modules/connections/connections.module.ts
  • apps/api/src/modules/connections/connections/registry-token-refresh.service.spec.ts
  • apps/api/src/modules/connections/connections/registry-token-refresh.service.ts
  • apps/api/src/modules/connections/connections/token-refresh.service.ts
  • apps/api/src/modules/dbmanager/dbmanager.module.ts
  • apps/api/src/modules/pipeline/cdc-relay.controller.spec.ts
  • apps/api/src/modules/pipeline/cdc-relay.controller.ts
  • apps/api/src/modules/pipeline/cdc-relay.guard.ts
  • apps/api/src/modules/pipeline/debezium-event.ts
  • apps/api/src/modules/pipeline/dto/debezium-event.dto.ts
  • apps/api/src/modules/stitches/stitches.controller.spec.ts
  • apps/api/src/modules/trigger/trigger-executor.service.ts
  • apps/api/src/modules/webhooks/webhooks.controller.spec.ts
  • apps/api/src/modules/webhooks/webhooks.controller.ts
  • apps/api/vitest.config.mts
  • apps/worker/src/app.module.ts
  • apps/worker/src/db/database-manager.ts
  • apps/worker/src/modules/pipeline/delivery.service.ts
  • apps/worker/src/modules/pipeline/fanout.service.ts
  • apps/worker/src/modules/pipeline/normalization.service.ts
  • apps/worker/src/modules/pipeline/pipeline.module.ts
  • apps/worker/src/modules/pipeline/pipeline.utils.spec.ts
  • apps/worker/src/modules/pipeline/registry-token-refresh.service.ts
  • apps/worker/src/modules/pipeline/replica.service.spec.ts
  • apps/worker/src/modules/pipeline/replica.service.ts
  • apps/worker/src/shared/pipeline.utils.ts
  • apps/worker/vitest.config.mts
  • engine/sync/platform/core/src/hydrator.ts
  • infra/debezium/application.properties
  • infra/debezium/application.properties.prod
  • packages/credentials/src/index.ts
  • packages/credentials/src/oauth/token-refresh.service.ts
  • packages/database/src/schema/pipeline.ts
  • packages/dbmanager/src/impl/sql-database-manager.ts
  • packages/domain/tms/src/schema/tms-provisioner.ts
  • packages/domain/tms/src/schema/tms-schema.ts
  • packages/domain/tms/src/tms-normalized-writer.ts
  • packages/domain/tms/src/tms-target-builder.ts
  • packages/pieces/application/revenova/src/index.ts
  • packages/pieces/application/revenova/src/normalizeRevenovaToTms.ts
  • packages/pieces/application/revenova/src/upsertRevenovaObject.ts
  • packages/pieces/platform/framework/src/canonical/index.ts
  • packages/pieces/platform/quickbooks/src/index.ts
💤 Files with no reviewable changes (4)
  • apps/worker/src/shared/pipeline.utils.ts
  • apps/worker/src/modules/pipeline/pipeline.utils.spec.ts
  • apps/api/src/modules/connections/connections/token-refresh.service.ts
  • apps/api/src/modules/pipeline/dto/debezium-event.dto.ts

Comment thread apps/api/src/db/database-manager.ts Outdated
Comment thread apps/api/src/modules/pipeline/cdc-relay.controller.spec.ts Outdated
Comment thread apps/api/src/modules/pipeline/cdc-relay.guard.ts Outdated
Comment thread apps/api/src/modules/webhooks/webhooks.controller.ts
Comment thread apps/worker/src/db/database-manager.ts
Comment thread packages/dbmanager/src/impl/sql-database-manager.ts
Comment thread packages/dbmanager/src/impl/sql-database-manager.ts
Comment thread packages/pieces/application/revenova/src/index.ts
Comment thread packages/pieces/platform/quickbooks/src/index.ts Outdated
Comment thread packages/pieces/platform/quickbooks/src/index.ts
@coderabbitai

coderabbitai Bot commented Apr 28, 2026 •

Copy link
Copy Markdown
Contributor

Note

Autofix is a beta feature. Expect some limitations and changes as we gather feedback and continue to improve it.

Fixes Applied Successfully

Fixed 14 file(s) based on 19 unresolved review comments.

Files modified:

  • apps/api/src/db/database-manager.ts
  • apps/api/src/modules/pipeline/cdc-relay.controller.spec.ts
  • apps/api/src/modules/pipeline/cdc-relay.guard.ts
  • apps/api/src/modules/webhooks/webhooks.controller.ts
  • apps/worker/src/db/database-manager.ts
  • apps/worker/src/modules/pipeline/delivery.service.ts
  • apps/worker/src/modules/pipeline/fanout.service.ts
  • apps/worker/src/modules/pipeline/replica.service.ts
  • engine/sync/platform/core/src/hydrator.ts
  • infra/debezium/application.properties
  • packages/credentials/src/oauth/token-refresh.service.ts
  • packages/dbmanager/src/impl/sql-database-manager.ts
  • packages/pieces/application/revenova/src/index.ts
  • packages/pieces/platform/quickbooks/src/index.ts

Commit: df1093af3111ed93b915e11569cc8051e09f7081

The changes have been pushed to the feature/enterprise-canonical-pipeline branch.

Time taken: 10m 12s

Fixed 14 file(s) based on 19 unresolved review comments.

Co-authored-by: CodeRabbit <noreply@coderabbit.ai>

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

♻️ Duplicate comments (6)
engine/sync/platform/core/src/hydrator.ts (1)

73-73: ⚠️ Potential issue | 🟠 Major

Restore null-prototype root payload to avoid inherited-key collisions (Line 73).

Using {} here reintroduces Object.prototype keys into hydration, so destinations like toString.id/hasOwnProperty.value can throw instead of being mapped. Use a null-prototype object for the root payload, consistent with the nested container strategy.

🔧 Proposed fix
-  const payload: any = {};
+  const payload: any = Object.create(null);
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@engine/sync/platform/core/src/hydrator.ts` at line 73, The root payload is
being created with a normal object ({}) which allows Object.prototype keys to
collide; change the root payload creation in hydrator.ts (the payload variable)
to use a null-prototype object (e.g., via Object.create(null)) so it matches the
existing null-prototype nested container strategy and avoids inherited-key
collisions like toString.id or hasOwnProperty.value during hydration.
apps/api/src/modules/pipeline/cdc-relay.guard.ts (1)

24-29: 🧹 Nitpick | 🔵 Trivial

Add guard tests for the new x-debezium-auth branch.

Current snippet in apps/api/src/modules/pipeline/cdc-relay.guard.spec.ts (Line 33-51) covers only Authorization: Bearer .... Please add positive and negative tests for x-debezium-auth to lock this auth path.

🧪 Suggested test additions
+it('allows access when x-debezium-auth header matches secret', () => {
+  mockRequest.headers['x-debezium-auth'] = 'test-secret';
+  expect(guard.canActivate(mockExecutionContext as ExecutionContext)).toBe(true);
+});
+
+it('throws UnauthorizedException when x-debezium-auth header is invalid', () => {
+  mockRequest.headers['x-debezium-auth'] = 'wrong-secret';
+  expect(() =>
+    guard.canActivate(mockExecutionContext as ExecutionContext),
+  ).toThrow(UnauthorizedException);
+});
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@apps/api/src/modules/pipeline/cdc-relay.guard.ts` around lines 24 - 29, Add
two tests to the CdcRelayGuard spec to cover the x-debezium-auth branch: one
positive test that sets req.headers['x-debezium-auth'] to the expectedSecret and
asserts CdcRelayGuard.canActivate returns true, and one negative test that sets
req.headers['x-debezium-auth'] to an invalid value (or omits it) and asserts
canActivate returns false (keeping the existing expectedSecret test
setup/fixtures and using the same guard instantiation and mock execution context
as the existing Authorization tests); ensure you reference and reuse the
expectedSecret variable and the canActivate method so the new tests exercise the
alternate header path.
packages/credentials/src/oauth/token-refresh.service.ts (1)

30-35: ⚠️ Potential issue | 🟠 Major

Keep input validation inside the OAuthRefreshError boundary.

validateInputs() is now the only failure path in refresh() that still escapes as a raw TypeError. Any caller that only handles OAuthRefreshError can still miss this branch, and it also skips the structured error log.

Suggested fix
-    this.validateInputs(tenantId, appName, externalId, refreshToken);
-
     try {
+      this.validateInputs(tenantId, appName, externalId, refreshToken);
       const tokenUrl = await this.getTokenUrl(appName);
       const { clientId, clientSecret, vendorParams } = await this.getCredentials(tenantId, appName, externalId);
       const resolvedTokenUrl = resolveOAuth2Url(tokenUrl, vendorParams);
@@
     } catch (error) {
       if (error instanceof OAuthRefreshError) throw error;
+      if (error instanceof TypeError && error.message.startsWith('Invalid refresh input:')) {
+        throw new OAuthRefreshError(error.message, 400);
+      }
       this.logger.error(`[TokenRefresh] Unexpected error for ${appName} on tenant ${tenantId}:`, error);
       throw new OAuthRefreshError(`Unexpected error during token refresh for ${appName}: ${error instanceof Error ? error.message : String(error)}`, 500);
     }

Also applies to: 51-55, 58-63

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@packages/credentials/src/oauth/token-refresh.service.ts` around lines 30 -
35, validateInputs(...) in refresh() can throw a raw TypeError and escape the
existing OAuthRefreshError boundary; wrap the call to validateInputs inside the
same try/catch that converts all failures into OAuthRefreshError so callers only
see the structured error. Concretely, move or include validateInputs(tenantId,
appName, externalId, refreshToken) into the try block that also calls
getTokenUrl, getCredentials and resolveOAuth2Url (or catch its errors and
rethrow an OAuthRefreshError with the original error attached), so all
exceptions from validateInputs, getTokenUrl, getCredentials, and
resolveOAuth2Url are handled uniformly by the OAuthRefreshError path.
apps/worker/src/modules/pipeline/replica.service.ts (1)

129-135: ⚠️ Potential issue | 🟠 Major

Exit processMessage() on idempotent skip, not just the transaction callback.

This return only leaves the inner transaction. The method still falls through to the best-effort ReplicaQueue enqueue at Lines 286-297, so traces that L2 intentionally skipped are replayed downstream anyway. Bubble a sentinel out of the transaction and guard the enqueue/log path on it.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@apps/worker/src/modules/pipeline/replica.service.ts` around lines 129 - 135,
The idempotent-skip return inside the transaction for processMessage only exits
the transaction callback but not the outer method, causing skipped traces to
still be enqueued to ReplicaQueue; modify processMessage so the transaction
callback returns a sentinel (e.g., SKIPPED_IDEMPOTENT) or sets a local flag when
the inbound.status is not "RECEIVED"/"PENDING", propagate that sentinel/flag out
of the transaction, and then guard the best-effort ReplicaQueue enqueue/log path
(the code around ReplicaQueue enqueue at the previously shown block) to skip
enqueueing and logging when the sentinel/flag indicates an idempotent skip.
Ensure you use the existing function/method names (processMessage, the
transaction callback, and the ReplicaQueue enqueue call) to locate and implement
the guard.
apps/api/src/db/database-manager.ts (1)

1055-1059: ⚠️ Potential issue | 🟠 Major

Stop attaching the seeded stitch to an arbitrary workspace.

This helper still grabs ui_workspaces.limit(1), so once a local DB has multiple workspaces the seeded stitch can land under the wrong org/workspace pair. Since seedMapping() is meant for deterministic local fixtures, resolve the exact seeded workspace (or require it explicitly) instead of taking the first row.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@apps/api/src/db/database-manager.ts` around lines 1055 - 1059, The helper
currently selects an arbitrary workspace via
db.select().from(schema.uiWorkspaces).limit(1) which can attach the seeded
stitch to the wrong org; update seedMapping() to resolve the exact seeded
workspace instead: either require a workspace identifier parameter (workspaceId
or workspaceSlug) and use it to query uiWorkspaces, or query uiWorkspaces by a
deterministic seeded field (e.g., name or is_seeded flag) instead of limit(1);
replace the existing db.select() call with a targeted lookup and throw a clear
error if the requested workspace is not found.
apps/worker/src/modules/pipeline/fanout.service.ts (1)

344-351: ⚠️ Potential issue | 🟠 Major

Use stitch.targetObject in this predicate.

integration_stitch exposes targetObject, not destEntityType, so this filter resolves to the wrong value and the target replica lookup still misses. That means _sync.dest.state never hydrates even when the destination cache already has the entity, which pushes downstream update flows back onto the slower/failing live-fetch path.

Proposed fix
             const targetReplica = await this.db
               .select()
               .from(targetReplicaEntity)
               .where(
-                sql`${targetReplicaEntity.connectionId} = ${stitch.destConnectionId} AND ${targetReplicaEntity.entityType} = ${stitch.destEntityType} AND ${targetReplicaEntity.entityId} = ${destEntityId}`,
+                sql`${targetReplicaEntity.connectionId} = ${stitch.destConnectionId} AND ${targetReplicaEntity.entityType} = ${stitch.targetObject} AND ${targetReplicaEntity.entityId} = ${destEntityId}`,
               )
               .limit(1);
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@apps/worker/src/modules/pipeline/fanout.service.ts` around lines 344 - 351,
The WHERE predicate for finding the target replica uses the wrong stitch field:
replace uses of stitch.destEntityType with stitch.targetObject so the lookup
compares targetReplicaEntity.entityType to stitch.targetObject (keep the other
predicates using stitch.destConnectionId and destEntityId as-is); update the SQL
in the targetReplica lookup (the query built off targetReplicaEntity in the
fanout service) so the filter matches the integration_stitch's exposed
targetObject, which will allow _sync.dest.state to hydrate correctly when the
destination cache already contains the entity.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.

Inline comments:
In `@apps/api/src/db/database-manager.ts`:
- Around line 851-860: The existing-row update only sets schemaPlan to
'OUTBOUND_ACTIVE' but fails to update dataNamespace, leaving the registry
pointing at an old namespace; update the db.update call against
dbSchema.connectionStorageRegistry so it also sets dataNamespace: schemaName
(the same field updated in the worker-side fix) when where(...) matches
resolved.id, ensuring both schemaPlan and dataNamespace are persisted for that
connection.

In `@apps/worker/src/modules/pipeline/replica.service.ts`:
- Around line 181-217: The lock in activeSyncLocks is acquired too early (during
the transaction that inserts with tx.insert) causing entities to remain locked
on all L4 "no outbound work" paths because locks are only released in
DeliveryService.writeL6Result; modify the flow so either (A) move the lock
acquisition to after route resolution/after confirming at least one outbound
route exists (i.e. perform route/stitch checks before calling tx.insert on
activeSyncLocks using resolvedEntityId/connectionId/traceId), or (B) add
explicit release logic that deletes the activeSyncLocks row on every L4
skip/fail branch (ensure the delete uses connectionId and resolvedEntityId and
runs in the same transactional context or a safe compensating transaction), and
keep the existing TTL as a safety fallback. Ensure DeliveryService.writeL6Result
remains a path that also removes the lock to avoid double-locks.

In `@infra/debezium/application.properties`:
- Line 15: Remove the unnecessary delivery_outbox table from the Debezium
include list by editing the debezium.source.table.include.list property: delete
the ",.*\\.delivery_outbox" segment so the list becomes only inbound_outbox,
replica_outbox, and normalized_outbox; this aligns with the CDC relay logic that
intentionally ignores delivery_outbox (see cdc-relay.controller.ts handling) and
avoids extra CDC load and log noise.

In `@packages/credentials/src/oauth/token-refresh.service.ts`:
- Around line 73-83: The decrypted JSON (valueBlob) must be validated as a
non-null plain object before accessing properties or using object spreads: after
JSON.parse(decryptedValue) in token-refresh.service.ts, check that valueBlob is
an object (typeof valueBlob === 'object' && valueBlob !== null &&
!Array.isArray(valueBlob)); if not, throw a clear "invalid credential payload"
error; ensure clientId/clientSecret checks only run when valueBlob is an object
and when building vendorParams only spread valueBlob.vendorParams if it is
itself an object and coerce environment safely (use 'environment' in valueBlob
only after the object check).

In `@packages/dbmanager/src/impl/sql-database-manager.ts`:
- Around line 202-213: migrateAllSchemas() currently only replays
migrateReplicaTables(), so existing tenants miss newly provisioned objects
(active_sync_locks table, outbox schema_name columns, sync_log partial indexes)
that ReplicaService and DeliveryService expect; add a migration entrypoint that,
when upgrading existing tenant schemas, reapplies all provisioner layers up to
the OUTBOUND_ACTIVE state (i.e., invoke the same sequence used for fresh
OUTBOUND_ACTIVE creation rather than only migrateReplicaTables()), ensuring
migrations create "active_sync_locks", add the outbox schema_name columns, and
create the sync_log partial indexes for each tenant schema; update the migration
dispatcher to call this new entrypoint from migrateAllSchemas() and reference
the existing provisioner functions that create those objects so no schema is
left missing after upgrade.

---

Duplicate comments:
In `@apps/api/src/db/database-manager.ts`:
- Around line 1055-1059: The helper currently selects an arbitrary workspace via
db.select().from(schema.uiWorkspaces).limit(1) which can attach the seeded
stitch to the wrong org; update seedMapping() to resolve the exact seeded
workspace instead: either require a workspace identifier parameter (workspaceId
or workspaceSlug) and use it to query uiWorkspaces, or query uiWorkspaces by a
deterministic seeded field (e.g., name or is_seeded flag) instead of limit(1);
replace the existing db.select() call with a targeted lookup and throw a clear
error if the requested workspace is not found.

In `@apps/api/src/modules/pipeline/cdc-relay.guard.ts`:
- Around line 24-29: Add two tests to the CdcRelayGuard spec to cover the
x-debezium-auth branch: one positive test that sets
req.headers['x-debezium-auth'] to the expectedSecret and asserts
CdcRelayGuard.canActivate returns true, and one negative test that sets
req.headers['x-debezium-auth'] to an invalid value (or omits it) and asserts
canActivate returns false (keeping the existing expectedSecret test
setup/fixtures and using the same guard instantiation and mock execution context
as the existing Authorization tests); ensure you reference and reuse the
expectedSecret variable and the canActivate method so the new tests exercise the
alternate header path.

In `@apps/worker/src/modules/pipeline/fanout.service.ts`:
- Around line 344-351: The WHERE predicate for finding the target replica uses
the wrong stitch field: replace uses of stitch.destEntityType with
stitch.targetObject so the lookup compares targetReplicaEntity.entityType to
stitch.targetObject (keep the other predicates using stitch.destConnectionId and
destEntityId as-is); update the SQL in the targetReplica lookup (the query built
off targetReplicaEntity in the fanout service) so the filter matches the
integration_stitch's exposed targetObject, which will allow _sync.dest.state to
hydrate correctly when the destination cache already contains the entity.

In `@apps/worker/src/modules/pipeline/replica.service.ts`:
- Around line 129-135: The idempotent-skip return inside the transaction for
processMessage only exits the transaction callback but not the outer method,
causing skipped traces to still be enqueued to ReplicaQueue; modify
processMessage so the transaction callback returns a sentinel (e.g.,
SKIPPED_IDEMPOTENT) or sets a local flag when the inbound.status is not
"RECEIVED"/"PENDING", propagate that sentinel/flag out of the transaction, and
then guard the best-effort ReplicaQueue enqueue/log path (the code around
ReplicaQueue enqueue at the previously shown block) to skip enqueueing and
logging when the sentinel/flag indicates an idempotent skip. Ensure you use the
existing function/method names (processMessage, the transaction callback, and
the ReplicaQueue enqueue call) to locate and implement the guard.

In `@engine/sync/platform/core/src/hydrator.ts`:
- Line 73: The root payload is being created with a normal object ({}) which
allows Object.prototype keys to collide; change the root payload creation in
hydrator.ts (the payload variable) to use a null-prototype object (e.g., via
Object.create(null)) so it matches the existing null-prototype nested container
strategy and avoids inherited-key collisions like toString.id or
hasOwnProperty.value during hydration.

In `@packages/credentials/src/oauth/token-refresh.service.ts`:
- Around line 30-35: validateInputs(...) in refresh() can throw a raw TypeError
and escape the existing OAuthRefreshError boundary; wrap the call to
validateInputs inside the same try/catch that converts all failures into
OAuthRefreshError so callers only see the structured error. Concretely, move or
include validateInputs(tenantId, appName, externalId, refreshToken) into the try
block that also calls getTokenUrl, getCredentials and resolveOAuth2Url (or catch
its errors and rethrow an OAuthRefreshError with the original error attached),
so all exceptions from validateInputs, getTokenUrl, getCredentials, and
resolveOAuth2Url are handled uniformly by the OAuthRefreshError path.
🪄 Autofix (Beta)

✅ Autofix completed


ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: 2185376d-11ad-4885-b7de-36dc88a21600

📥 Commits

Reviewing files that changed from the base of the PR and between d9e39a9 and df1093a.

📒 Files selected for processing (14)
  • apps/api/src/db/database-manager.ts
  • apps/api/src/modules/pipeline/cdc-relay.controller.spec.ts
  • apps/api/src/modules/pipeline/cdc-relay.guard.ts
  • apps/api/src/modules/webhooks/webhooks.controller.ts
  • apps/worker/src/db/database-manager.ts
  • apps/worker/src/modules/pipeline/delivery.service.ts
  • apps/worker/src/modules/pipeline/fanout.service.ts
  • apps/worker/src/modules/pipeline/replica.service.ts
  • engine/sync/platform/core/src/hydrator.ts
  • infra/debezium/application.properties
  • packages/credentials/src/oauth/token-refresh.service.ts
  • packages/dbmanager/src/impl/sql-database-manager.ts
  • packages/pieces/application/revenova/src/index.ts
  • packages/pieces/platform/quickbooks/src/index.ts

Comment thread apps/api/src/db/database-manager.ts
Comment on lines +181 to +217
// ── ACQUIRE ENTITY LOCK ───────────────────────────────────────────────
// Prevents an UPDATE event from starting while a CREATE event is still
// in-flight (L2 -> L6), ensuring the UPDATE has access to the GEM mapping.
try {
// Self-healing: clear any stale locks that have expired (e.g. from permanently crashed workers)
await tx
.delete(activeSyncLocks)
.where(
and(
eq(activeSyncLocks.connectionId, connectionId),
eq(activeSyncLocks.entityId, resolvedEntityId),
sql`${activeSyncLocks.expiresAt} < NOW()`,
),
);

await tx.insert(activeSyncLocks).values({
connectionId,
entityId: resolvedEntityId,
lockedByTraceId: traceId,
expiresAt: sql`NOW() + INTERVAL '10 minutes'`,
});
} catch (err: unknown) {
if (
err instanceof Error &&
(err.message.includes("unique constraint") ||
err.message.includes("duplicate key"))
) {
// Lock contention — treat as a deferral (retry later), NOT a terminal FAIL
const lockContentionError = new Error(
`Entity ${resolvedEntityId} is currently locked by an in-flight sync. ` +
`Delaying processing to maintain FIFO order.`,
);
(lockContentionError as any).isLockContention = true;
throw lockContentionError;
}
throw err;
}

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.

⚠️ Potential issue | 🟠 Major

These locks survive every “no outbound work” path.

The lock is acquired before route resolution, but the only release path in the provided changes is in DeliveryService.writeL6Result(). If L4 finds no active stitches, unmatched conditions, or no mapping, the trace never reaches L6 and this entity stays blocked until the 10-minute TTL expires. Either acquire the lock only after at least one outbound route exists, or add explicit release paths for all L4 skip/fail branches.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@apps/worker/src/modules/pipeline/replica.service.ts` around lines 181 - 217,
The lock in activeSyncLocks is acquired too early (during the transaction that
inserts with tx.insert) causing entities to remain locked on all L4 "no outbound
work" paths because locks are only released in DeliveryService.writeL6Result;
modify the flow so either (A) move the lock acquisition to after route
resolution/after confirming at least one outbound route exists (i.e. perform
route/stitch checks before calling tx.insert on activeSyncLocks using
resolvedEntityId/connectionId/traceId), or (B) add explicit release logic that
deletes the activeSyncLocks row on every L4 skip/fail branch (ensure the delete
uses connectionId and resolvedEntityId and runs in the same transactional
context or a safe compensating transaction), and keep the existing TTL as a
safety fallback. Ensure DeliveryService.writeL6Result remains a path that also
removes the lock to avoid double-locks.

Comment thread infra/debezium/application.properties Outdated
Comment thread packages/credentials/src/oauth/token-refresh.service.ts
Comment on lines +202 to +213
// ── ACTIVE SYNC LOCKS — L1 → L6 concurrency control ─────────────────
await this.db.$client.query(`
CREATE TABLE IF NOT EXISTS "${schemaName}".active_sync_locks (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
connection_id UUID NOT NULL,
entity_id VARCHAR(255) NOT NULL,
locked_by_trace_id UUID NOT NULL,
expires_at TIMESTAMPTZ NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
CONSTRAINT uq_sync_lock UNIQUE (connection_id, entity_id)
);
`);

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.

⚠️ Potential issue | 🔴 Critical

Provision these new tenant-schema changes from the migration path as well.

migrateAllSchemas() still replays only migrateReplicaTables(), so existing tenants upgraded through the schema migration command will not get active_sync_locks, the new outbox schema_name columns, or the sync_log partial indexes. ReplicaService and DeliveryService already depend on those objects, so this rollout can break pre-existing schemas even though fresh OUTBOUND_ACTIVE schemas work. Please add a migration entrypoint that reapplies all layer provisioners up to OUTBOUND_ACTIVE for existing tenants.

Also applies to: 337-342, 476-507, 541-605

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@packages/dbmanager/src/impl/sql-database-manager.ts` around lines 202 - 213,
migrateAllSchemas() currently only replays migrateReplicaTables(), so existing
tenants miss newly provisioned objects (active_sync_locks table, outbox
schema_name columns, sync_log partial indexes) that ReplicaService and
DeliveryService expect; add a migration entrypoint that, when upgrading existing
tenant schemas, reapplies all provisioner layers up to the OUTBOUND_ACTIVE state
(i.e., invoke the same sequence used for fresh OUTBOUND_ACTIVE creation rather
than only migrateReplicaTables()), ensuring migrations create
"active_sync_locks", add the outbox schema_name columns, and create the sync_log
partial indexes for each tenant schema; update the migration dispatcher to call
this new entrypoint from migrateAllSchemas() and reference the existing
provisioner functions that create those objects so no schema is left missing
after upgrade.

@coderabbitai

coderabbitai Bot commented Apr 28, 2026 •

Copy link
Copy Markdown
Contributor

Note

Autofix is a beta feature. Expect some limitations and changes as we gather feedback and continue to improve it.

Fixes Applied Successfully

Fixed 5 file(s) based on 5 unresolved review comments.

Files modified:

  • apps/api/src/db/database-manager.ts
  • apps/worker/src/modules/pipeline/fanout.service.ts
  • infra/debezium/application.properties
  • packages/credentials/src/oauth/token-refresh.service.ts
  • packages/dbmanager/src/impl/sql-database-manager.ts

Commit: 77e1147f620c8f6e7085477d7f885f913e919878

The changes have been pushed to the feature/enterprise-canonical-pipeline branch.

Time taken: 7m 40s

Fixed 5 file(s) based on 5 unresolved review comments.

Co-authored-by: CodeRabbit <noreply@coderabbit.ai>
@pramodnarayana
pramodnarayana marked this pull request as draft April 28, 2026 16:56
@pramodnarayana
pramodnarayana marked this pull request as ready for review April 28, 2026 16:56

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

Caution

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

⚠️ Outside diff range comments (2)
packages/dbmanager/src/impl/sql-database-manager.ts (1)

341-357: ⚠️ Potential issue | 🟠 Major

Backfill schema_name with the tenant schema, not current_schema().

These statements run on the default connection context, not inside SET LOCAL search_path TO "${schemaName}". On existing normalized_outbox / delivery_outbox rows, DEFAULT current_schema() will therefore materialize the session schema (typically public) instead of the tenant schema being migrated, so legacy outbox rows can be routed to the wrong namespace. Use the validated schemaName literal for the default and the migration backfill before enforcing NOT NULL.

🛠️ Suggested fix
-            schema_name   VARCHAR(128) NOT NULL DEFAULT current_schema(),
+            schema_name   VARCHAR(128) NOT NULL DEFAULT '${schemaName}',
-        DO $$ BEGIN
-            ALTER TABLE "${schemaName}".normalized_outbox ADD COLUMN IF NOT EXISTS schema_name VARCHAR(128) NOT NULL DEFAULT current_schema();
-        EXCEPTION WHEN duplicate_column THEN NULL;
-        END $$;
+        ALTER TABLE "${schemaName}".normalized_outbox
+            ADD COLUMN IF NOT EXISTS schema_name VARCHAR(128);
+        UPDATE "${schemaName}".normalized_outbox
+           SET schema_name = '${schemaName}'
+         WHERE schema_name IS NULL;
+        ALTER TABLE "${schemaName}".normalized_outbox
+            ALTER COLUMN schema_name SET DEFAULT '${schemaName}',
+            ALTER COLUMN schema_name SET NOT NULL;

Apply the same pattern to delivery_outbox.

Also applies to: 553-583, 614-620

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@packages/dbmanager/src/impl/sql-database-manager.ts` around lines 341 - 357,
The migration uses DEFAULT current_schema() and will capture the session schema
(e.g., public) instead of the tenant being migrated; update the ALTER/CREATE
statements in sql-database-manager.ts (the normalized_outbox and delivery_outbox
DDL/ALTER blocks) to use the validated schemaName literal as the DEFAULT (e.g.,
DEFAULT 'schemaName') and perform an explicit UPDATE backfill setting
schema_name = 'schemaName' for existing rows before adding or enforcing NOT
NULL/DEFAULT constraints; apply the same pattern to delivery_outbox and the
other occurrences referenced (around the other ALTER/CREATE blocks).
apps/worker/src/modules/pipeline/fanout.service.ts (1)

266-285: ⚠️ Potential issue | 🔴 Critical

Don't drop the entity lock from a single skipped route.

processInChunks() runs stitches concurrently. If one route hits a SKIPPED path while another route for the same trace is still creating outbound work, this deletes the L2 lock early and allows a new CDC event for the same entity to start before the current trace has drained. Release only after you know the trace produced no outbound rows at all, or leave lock ownership with delivery.service.ts after it verifies there are no remaining routes.

Also applies to: 297-326

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@apps/worker/src/modules/pipeline/fanout.service.ts` around lines 266 - 285,
The code currently calls releaseSyncLock from the per-route SKIPPED branch
(inside processInChunks / the stitch handling in fanout.service.ts), which can
drop the L2 lock while other concurrent routes for the same trace may still be
producing outbound work; remove the releaseSyncLock call from this per-route
SKIPPED path (and the similar calls in the 297-326 region) so a single skipped
route does not release the entity lock; instead leave lock ownership to
delivery.service.ts (or a higher-level caller) to release only after it verifies
the entire trace produced no outbound rows (i.e., after all routes/stitches are
processed), and keep writeSyncLog and return as-is in the per-route branch.
♻️ Duplicate comments (2)
packages/credentials/src/oauth/token-refresh.service.ts (1)

80-88: ⚠️ Potential issue | 🟡 Minor

Normalize vendorParams values to strings before returning.

vendorParams is declared as Record<string, string>, but raw spread can carry numbers/booleans/objects from stored JSON. That can produce invalid OAuth URL substitutions and violates the method contract.

🔧 Suggested fix
-      const vendorParamsSpread = typeof valueBlob.vendorParams === 'object' && valueBlob.vendorParams !== null && !Array.isArray(valueBlob.vendorParams) ? valueBlob.vendorParams : {};
+      const vendorParamsSpread =
+        typeof valueBlob.vendorParams === 'object' &&
+        valueBlob.vendorParams !== null &&
+        !Array.isArray(valueBlob.vendorParams)
+          ? Object.fromEntries(
+              Object.entries(valueBlob.vendorParams).map(([key, value]) => [key, String(value)]),
+            )
+          : {};
       const environmentEntry = 'environment' in valueBlob && valueBlob.environment ? { environment: String(valueBlob.environment) } : {};
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@packages/credentials/src/oauth/token-refresh.service.ts` around lines 80 -
88, The returned vendorParams must be normalized to strings: when building the
return object (the variables vendorParamsSpread and environmentEntry in
token-refresh.service.ts), coerce every value in vendorParamsSpread and
environmentEntry to String(value) so the merged vendorParams object is
Record<string,string>; iterate the keys (or map entries) of vendorParamsSpread
and environmentEntry and replace non-string values with their stringified
equivalents before spreading into the returned vendorParams to ensure URL
substitutions and the method contract remain valid.
apps/api/src/db/database-manager.ts (1)

1058-1060: ⚠️ Potential issue | 🟠 Major

This still binds the seed to an arbitrary workspace.

limit(1) on uiWorkspaces will attach the local stitch to whichever workspace happens to come back first once a dev/test database has multiple rows. This helper is meant to use the deterministic local fixtures, so it should resolve the seeded workspace explicitly instead of taking the first row.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@apps/api/src/db/database-manager.ts` around lines 1058 - 1060, The current
lookup uses db.select().from(schema.uiWorkspaces).limit(1) which can bind the
seed to an arbitrary workspace; replace this non-deterministic query with an
explicit lookup for the seeded/local workspace (e.g., query by the known fixture
identifier/name). Update the code that defines workspaces (the
db.select...from(schema.uiWorkspaces).limit(1) call) to use a where clause
referencing the deterministic seed value (e.g., where({ id: LOCAL_WORKSPACE_ID
}) or where('name', '=', LOCAL_WORKSPACE_NAME)), or fetch the seed id from the
local fixtures helper and use that in the query so the local stitch always
attaches to the intended seeded workspace.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.

Inline comments:
In `@apps/api/src/db/database-manager.ts`:
- Around line 843-863: The registry is being marked OUTBOUND_ACTIVE before
provisioning completes; change the order so you first insert or update
dataNamespace (seed/update dbSchema.connectionStorageRegistry.dataNamespace for
resolved.id/schemaName as you already do) but do NOT set schemaPlan there, then
call schemaMgr.applyPlan(schemaName, SchemaPlan.OUTBOUND_ACTIVE), and only after
that succeeds perform an update on dbSchema.connectionStorageRegistry for
connectionId == resolved.id to set schemaPlan = SchemaPlan.OUTBOUND_ACTIVE (and
update dataNamespace if needed); adjust the branches around existReg,
db.insert/db.update and the final db.update to flip schemaPlan after applyPlan
returns.
- Around line 1128-1144: The seed upsert for schema.fieldMappings is missing the
non-null destCanonical column; update the insert values and the
onConflictDoUpdate.set to include destCanonical: 'Vendor' so the initial insert
and any conflict-update both persist destCanonical correctly for the QuickBooks
Vendor mapping (referencing schema.fieldMappings, stitchId, sourceCanonical,
mappingRules, destCanonical).

In `@apps/worker/src/modules/pipeline/fanout.service.ts`:
- Around line 563-578: The releaseSyncLock function must only delete the lock
owned by the current trace to avoid removing newer locks; update
releaseSyncLock(schemaName, connectionId, entityId, activeSyncLocks) to accept a
traceId parameter and modify the delete predicate in the
tx.delete(activeSyncLocks).where(...) to include a condition that
activeSyncLocks.locked_by_trace_id = traceId (in addition to connectionId and
entityId), and update any callers to pass the current traceId through to
releaseSyncLock.

In `@infra/debezium/application.properties`:
- Line 23: The property debezium.sink.http.headers=Authorization: Bearer
${DEBEZIUM_SECRET:secret} is unsupported so Debezium will not send the
Authorization header; replace this with Debezium's supported authentication
configuration by removing debezium.sink.http.headers and adding the appropriate
debezium.sink.http.authentication.* settings (for example set
debezium.sink.http.authentication.type to jwt and provide
debezium.sink.http.authentication.jwt.username, jwt.password and jwt.url, or use
type=standard-webhooks with debezium.sink.http.authentication.webhook.secret) so
the sink will send valid auth to your relay guard.

In `@packages/credentials/src/oauth/token-refresh.service.ts`:
- Line 91: The thrown error in token-refresh.service.ts when credentials can't
be retrieved only includes tenantId and appName; update the error message in the
throw inside the credential lookup (the same place that currently references
tenantId and appName) to also include externalId so the message reads something
like "Failed to retrieve credentials for tenantId=... appName=...
externalId=...: <error>". Ensure you use the existing error formatting (error
instanceof Error ? error.message : String(error)) and reference the existing
variables tenantId, appName, and externalId in the message.

---

Outside diff comments:
In `@apps/worker/src/modules/pipeline/fanout.service.ts`:
- Around line 266-285: The code currently calls releaseSyncLock from the
per-route SKIPPED branch (inside processInChunks / the stitch handling in
fanout.service.ts), which can drop the L2 lock while other concurrent routes for
the same trace may still be producing outbound work; remove the releaseSyncLock
call from this per-route SKIPPED path (and the similar calls in the 297-326
region) so a single skipped route does not release the entity lock; instead
leave lock ownership to delivery.service.ts (or a higher-level caller) to
release only after it verifies the entire trace produced no outbound rows (i.e.,
after all routes/stitches are processed), and keep writeSyncLog and return as-is
in the per-route branch.

In `@packages/dbmanager/src/impl/sql-database-manager.ts`:
- Around line 341-357: The migration uses DEFAULT current_schema() and will
capture the session schema (e.g., public) instead of the tenant being migrated;
update the ALTER/CREATE statements in sql-database-manager.ts (the
normalized_outbox and delivery_outbox DDL/ALTER blocks) to use the validated
schemaName literal as the DEFAULT (e.g., DEFAULT 'schemaName') and perform an
explicit UPDATE backfill setting schema_name = 'schemaName' for existing rows
before adding or enforcing NOT NULL/DEFAULT constraints; apply the same pattern
to delivery_outbox and the other occurrences referenced (around the other
ALTER/CREATE blocks).

---

Duplicate comments:
In `@apps/api/src/db/database-manager.ts`:
- Around line 1058-1060: The current lookup uses
db.select().from(schema.uiWorkspaces).limit(1) which can bind the seed to an
arbitrary workspace; replace this non-deterministic query with an explicit
lookup for the seeded/local workspace (e.g., query by the known fixture
identifier/name). Update the code that defines workspaces (the
db.select...from(schema.uiWorkspaces).limit(1) call) to use a where clause
referencing the deterministic seed value (e.g., where({ id: LOCAL_WORKSPACE_ID
}) or where('name', '=', LOCAL_WORKSPACE_NAME)), or fetch the seed id from the
local fixtures helper and use that in the query so the local stitch always
attaches to the intended seeded workspace.

In `@packages/credentials/src/oauth/token-refresh.service.ts`:
- Around line 80-88: The returned vendorParams must be normalized to strings:
when building the return object (the variables vendorParamsSpread and
environmentEntry in token-refresh.service.ts), coerce every value in
vendorParamsSpread and environmentEntry to String(value) so the merged
vendorParams object is Record<string,string>; iterate the keys (or map entries)
of vendorParamsSpread and environmentEntry and replace non-string values with
their stringified equivalents before spreading into the returned vendorParams to
ensure URL substitutions and the method contract remain valid.
🪄 Autofix (Beta)

✅ Autofix completed


ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: 9d01b8d9-7b75-4632-81b5-b204e74b157b

📥 Commits

Reviewing files that changed from the base of the PR and between df1093a and 77e1147.

📒 Files selected for processing (5)
  • apps/api/src/db/database-manager.ts
  • apps/worker/src/modules/pipeline/fanout.service.ts
  • infra/debezium/application.properties
  • packages/credentials/src/oauth/token-refresh.service.ts
  • packages/dbmanager/src/impl/sql-database-manager.ts

Comment thread apps/api/src/db/database-manager.ts
Comment thread apps/api/src/db/database-manager.ts
Comment thread apps/worker/src/modules/pipeline/fanout.service.ts
Comment thread infra/debezium/application.properties Outdated
Comment thread packages/credentials/src/oauth/token-refresh.service.ts Outdated
@coderabbitai

coderabbitai Bot commented Apr 28, 2026 •

Copy link
Copy Markdown
Contributor

Note

Autofix is a beta feature. Expect some limitations and changes as we gather feedback and continue to improve it.

Fixes Applied Successfully

Fixed 4 file(s) based on 5 unresolved review comments.

Files modified:

  • apps/api/src/db/database-manager.ts
  • apps/worker/src/modules/pipeline/fanout.service.ts
  • infra/debezium/application.properties
  • packages/credentials/src/oauth/token-refresh.service.ts

Commit: 7d9ccd2b3a71df1bc556706f52b13d315b841223

The changes have been pushed to the feature/enterprise-canonical-pipeline branch.

Time taken: 4m 57s

Fixed 4 file(s) based on 5 unresolved review comments.

Co-authored-by: CodeRabbit <noreply@coderabbit.ai>

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

Caution

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

⚠️ Outside diff range comments (1)
apps/worker/src/modules/pipeline/fanout.service.ts (1)

192-210: ⚠️ Potential issue | 🔴 Critical

Hold the L2 lock until all stitches are evaluated.

Line 277 and Line 319 release the entity lock inside per-stitch skip branches. With concurrent processInChunks, one skipped route can unlock while another route for the same entity is still in-flight, allowing overlapping traces/retries.

🛠️ Proposed fix (release once, after aggregate outcome)
-  const stitchResults = await processInChunks(stitches, 5, (stitch) =>
+  const stitchResults = await processInChunks(stitches, 5, (stitch) =>
     this.processSingleStitch(
       schemaName,
       traceId,
       connectionId,
@@
       syncLog,
       activeSyncLocks,
     ),
   );
+
+  const anyOutboundEnqueued = stitchResults
+    .filter(
+      (r): r is PromiseFulfilledResult<boolean> => r.status === "fulfilled",
+    )
+    .some((r) => r.value);
+
+  if (!anyOutboundEnqueued && srcVendorId) {
+    await this.releaseSyncLock(
+      schemaName,
+      connectionId,
+      srcVendorId,
+      activeSyncLocks,
+      traceId,
+    );
+  }
-  ): Promise<void> {
+  ): Promise<boolean> {
@@
-      if (!matched) {
+      if (!matched) {
         await this.writeSyncLog(...);
-        if (srcVendorId) {
-          await this.releaseSyncLock(...);
-        }
-        return;
+        return false;
       }
@@
-      if (mappings.length === 0) {
+      if (mappings.length === 0) {
         await this.writeSyncLog(...);
-        if (srcVendorId) {
-          await this.releaseSyncLock(...);
-        }
-        return;
+        return false;
       }
@@
       await this.db.transaction(async (tx) => {
         ...
         await tx.insert(deliveryOutbox)...
       });
+      return true;
@@
-    } catch (err) {
+    } catch (err) {
       ...
       await this.writeSyncLog(...);
+      return false;
     }
   }

Also applies to: 277-286, 319-328

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@apps/worker/src/modules/pipeline/fanout.service.ts` around lines 192 - 210,
The L2 entity lock is being released inside per-stitch skip branches which,
under concurrent processInChunks calls, allows premature unlocks; remove any
unlock/release calls from inside processSingleStitch's skip branches and instead
propagate a skip/outcome status back to the caller (the processInChunks
aggregation) and perform a single activeSyncLocks release once all stitches have
been evaluated (after stitchResults are collected). Locate uses of
processSingleStitch, processInChunks and activeSyncLocks and ensure lock release
happens once, after the aggregate outcome, not inside per-stitch early-return
paths.
♻️ Duplicate comments (2)
apps/api/src/db/database-manager.ts (1)

1065-1069: ⚠️ Potential issue | 🟠 Major

Avoid non-deterministic workspace selection in seedMapping().

Line 1066 still uses select().from(schema.uiWorkspaces).limit(1), which can attach the stitch to an arbitrary workspace/org when multiple rows exist. Resolve the exact seeded workspace (same deterministic identifier created by seeding) and fail if that specific row is missing.

Suggested direction
-      // Check if a workspace exists
-      const workspaces = await db.select().from(schema.uiWorkspaces).limit(1);
-      if (workspaces.length === 0) {
-        throw new Error('No workspace found. Run pnpm db:seed first.');
-      }
+      // Resolve the deterministic seeded workspace written by db:seed
+      const workspace = await db.query.uiWorkspaces.findFirst({
+        // Replace with the exact seeded selector used by your seed routine
+        // (e.g. seeded workspace ID/slug/name + orgId guard).
+        // where: and(eq(schema.uiWorkspaces.id, SEEDED_WORKSPACE_ID), eq(schema.uiWorkspaces.orgId, salesforceConn[0].tenantId)),
+      });
+      if (!workspace) {
+        throw new Error('Seeded local workspace not found. Run pnpm db:seed first.');
+      }
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@apps/api/src/db/database-manager.ts` around lines 1065 - 1069, seedMapping()
currently queries an arbitrary workspace via
db.select().from(schema.uiWorkspaces).limit(1), which is non-deterministic;
update it to look up the exact seeded workspace using the deterministic
identifier used by your seed data (for example the seeded workspace's id, slug,
or name) instead of limit(1), e.g. query schema.uiWorkspaces for that specific
identifier via db.select().from(schema.uiWorkspaces).where(...).limit(1) and if
the row is not found throw an explicit error indicating the expected seeded
workspace is missing so the operation fails loudly.
packages/credentials/src/oauth/token-refresh.service.ts (1)

80-88: ⚠️ Potential issue | 🟠 Major

Normalize vendorParams values to strings before returning.

Raw spread from decrypted JSON allows non-string values into a Record<string, string> contract. This can break URL templating/resolution when values are numbers/booleans/objects.

🔧 Proposed fix
-      const vendorParamsSpread = typeof valueBlob.vendorParams === 'object' && valueBlob.vendorParams !== null && !Array.isArray(valueBlob.vendorParams) ? valueBlob.vendorParams : {};
+      const vendorParamsSpread =
+        typeof valueBlob.vendorParams === 'object' &&
+        valueBlob.vendorParams !== null &&
+        !Array.isArray(valueBlob.vendorParams)
+          ? Object.fromEntries(
+              Object.entries(valueBlob.vendorParams).map(([key, value]) => [key, String(value)]),
+            )
+          : {};
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@packages/credentials/src/oauth/token-refresh.service.ts` around lines 80 -
88, vendorParamsSpread may contain non-string values but the return expects
Record<string,string>; before returning build a normalizedVendorParams by
iterating Object.entries(vendorParamsSpread) and mapping each [k,v] to [k,
String(v)] (or JSON.stringify(v) for objects if you prefer), then spread
normalizedVendorParams and environmentEntry into the returned vendorParams;
update the return site that currently spreads vendorParamsSpread to use this
normalized object so all vendorParams values are strings.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.

Inline comments:
In `@infra/debezium/application.properties`:
- Around line 23-26: The JWT auth URL is currently pointing at the relay
endpoint (debezium.sink.http.authentication.jwt.url -> /internal/cdc/relay),
which causes Debezium to append auth paths and fail token acquisition; change
debezium.sink.http.authentication.jwt.url to an auth-base URL (introduce/use an
env var like DEBEZIUM_AUTH_URL) that points to your authentication service
(e.g., /internal/auth or /internal/cdc/auth) while keeping the relay endpoint
value (DEBEZIUM_SINK_URL / debezium.sink.http.url) unchanged so token requests
go to the correct auth server instead of the relay. Ensure the properties
referenced are debezium.sink.http.authentication.jwt.url (new auth base) and the
existing debezium.sink.http.authentication.jwt.password /
debezium.sink.http.authentication.jwt.username remain the same.

In `@packages/credentials/src/oauth/token-refresh.service.ts`:
- Line 50: The return value from response.json() must be validated as an object
before casting to Record<string, unknown>; replace the unchecked cast around
response.json() with a runtime check that the parsed value is !== null, typeof
=== 'object', and !Array.isArray(value) (reference the response.json() call and
the current return statement). If the check passes return the value as
Record<string, unknown>, otherwise throw a descriptive error (or handle the
unexpected shape) so downstream code does not receive null, scalar, or array
values.

---

Outside diff comments:
In `@apps/worker/src/modules/pipeline/fanout.service.ts`:
- Around line 192-210: The L2 entity lock is being released inside per-stitch
skip branches which, under concurrent processInChunks calls, allows premature
unlocks; remove any unlock/release calls from inside processSingleStitch's skip
branches and instead propagate a skip/outcome status back to the caller (the
processInChunks aggregation) and perform a single activeSyncLocks release once
all stitches have been evaluated (after stitchResults are collected). Locate
uses of processSingleStitch, processInChunks and activeSyncLocks and ensure lock
release happens once, after the aggregate outcome, not inside per-stitch
early-return paths.

---

Duplicate comments:
In `@apps/api/src/db/database-manager.ts`:
- Around line 1065-1069: seedMapping() currently queries an arbitrary workspace
via db.select().from(schema.uiWorkspaces).limit(1), which is non-deterministic;
update it to look up the exact seeded workspace using the deterministic
identifier used by your seed data (for example the seeded workspace's id, slug,
or name) instead of limit(1), e.g. query schema.uiWorkspaces for that specific
identifier via db.select().from(schema.uiWorkspaces).where(...).limit(1) and if
the row is not found throw an explicit error indicating the expected seeded
workspace is missing so the operation fails loudly.

In `@packages/credentials/src/oauth/token-refresh.service.ts`:
- Around line 80-88: vendorParamsSpread may contain non-string values but the
return expects Record<string,string>; before returning build a
normalizedVendorParams by iterating Object.entries(vendorParamsSpread) and
mapping each [k,v] to [k, String(v)] (or JSON.stringify(v) for objects if you
prefer), then spread normalizedVendorParams and environmentEntry into the
returned vendorParams; update the return site that currently spreads
vendorParamsSpread to use this normalized object so all vendorParams values are
strings.
🪄 Autofix (Beta)

✅ Autofix completed


ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: dc68389a-7534-499d-ac96-002abd9dfdba

📥 Commits

Reviewing files that changed from the base of the PR and between 77e1147 and 7d9ccd2.

📒 Files selected for processing (4)
  • apps/api/src/db/database-manager.ts
  • apps/worker/src/modules/pipeline/fanout.service.ts
  • infra/debezium/application.properties
  • packages/credentials/src/oauth/token-refresh.service.ts

Comment thread infra/debezium/application.properties Outdated
Comment thread packages/credentials/src/oauth/token-refresh.service.ts Outdated
@coderabbitai

coderabbitai Bot commented Apr 28, 2026 •

Copy link
Copy Markdown
Contributor

Note

Autofix is a beta feature. Expect some limitations and changes as we gather feedback and continue to improve it.

Fixes Applied Successfully

Fixed 2 file(s) based on 2 unresolved review comments.

Files modified:

  • infra/debezium/application.properties
  • packages/credentials/src/oauth/token-refresh.service.ts

Commit: a6be5514e10ef21dd48206260c4181fd5c75eec7

The changes have been pushed to the feature/enterprise-canonical-pipeline branch.

Time taken: 3m 58s

Fixed 2 file(s) based on 2 unresolved review comments.

Co-authored-by: CodeRabbit <noreply@coderabbit.ai>
@pramodnarayana
pramodnarayana marked this pull request as draft April 28, 2026 17:33
@pramodnarayana
pramodnarayana marked this pull request as ready for review April 28, 2026 17:34

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

Caution

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

⚠️ Outside diff range comments (1)
apps/worker/src/modules/pipeline/fanout.service.ts (1)

508-529: ⚠️ Potential issue | 🟠 Major

Missing lock release in stitch error catch block.

When processSingleStitch fails (catch block at line 508), a FAIL sync_log is written but the entity lock is not released. If all stitches fail for a trace, the lock remains held until TTL expiry (10 minutes), blocking subsequent events for this entity.

Consider releasing the lock in the error path as well, or letting the caller handle lock cleanup after all stitches complete.

🔧 Suggested fix
     } catch (err) {
       this.logger.error(
         {
           event: "l4.stitch_error",
           stitchId: stitch.id,
           traceId,
           connectionId,
           layer: "L4",
           err: sanitizeError(err),
         },
         "L4 stitch fan-out failed — recording failure and continuing to next route",
       );
       await this.writeSyncLog(
         schemaName,
         traceId,
         stitch.id,
         "L4",
         "FAIL",
         Date.now() - start,
         syncLog,
       );
+      // Release lock — no successful outbound work will occur for this route
+      if (srcVendorId) {
+        await this.releaseSyncLock(
+          schemaName,
+          connectionId,
+          srcVendorId,
+          activeSyncLocks,
+          traceId,
+        );
+      }
     }
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@apps/worker/src/modules/pipeline/fanout.service.ts` around lines 508 - 529,
The catch block for processSingleStitch logs the failure and writes a FAIL sync
log but never releases the entity lock, leaving the lock held until TTL expiry;
add the same lock-release logic used in the successful path (the method you use
when acquiring/releasing entity locks — e.g., this.releaseEntityLock or
this.unlockEntity) and call it (await) in the catch block (before or after
writeSyncLog) for the same entity identifier used when acquiring (referencing
stitch.id or the entity id variable), handling any errors from the release so
failures here don't mask the original error; ensure the release call mirrors the
success-path release semantics and uses the same schemaName/trace/entity id
parameters.
♻️ Duplicate comments (2)
packages/credentials/src/oauth/token-refresh.service.ts (1)

84-92: ⚠️ Potential issue | 🟡 Minor

Coerce vendorParams values to strings for type safety.

Per the OAuthCredentialBlob interface in token-manager.service.ts, vendorParams should be Record<string, string>. However, the stored blob may contain non-string values. Currently, the spread passes values through unchanged, potentially violating the return type contract.

🔧 Suggested fix
-      const vendorParamsSpread = typeof valueBlob.vendorParams === 'object' && valueBlob.vendorParams !== null && !Array.isArray(valueBlob.vendorParams) ? valueBlob.vendorParams : {};
+      const rawVendorParams = typeof valueBlob.vendorParams === 'object' && valueBlob.vendorParams !== null && !Array.isArray(valueBlob.vendorParams) ? valueBlob.vendorParams : {};
+      const vendorParamsSpread = Object.fromEntries(
+        Object.entries(rawVendorParams).map(([k, v]) => [k, String(v)])
+      );
       const environmentEntry = 'environment' in valueBlob && valueBlob.environment ? { environment: String(valueBlob.environment) } : {};
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@packages/credentials/src/oauth/token-refresh.service.ts` around lines 84 -
92, The returned vendorParams currently spreads vendorParamsSpread directly
which can leave non-string values and violate the OAuthCredentialBlob type; in
the return where vendorParams is built (around vendorParamsSpread and
environmentEntry) replace the spread of vendorParamsSpread with a string-coerced
map (e.g. transform Object.entries(vendorParamsSpread) to
Object.fromEntries(entries.map(([k,v]) => [k, String(v)]))) so every value is a
string, keep environmentEntry as String(...) and ensure the final vendorParams
is typed/constructed as Record<string,string>.
apps/worker/src/modules/pipeline/replica.service.ts (1)

129-135: ⚠️ Potential issue | 🟠 Major

Idempotent skip still enqueues to L3 (ReplicaQueue).

When status is not RECEIVED/PENDING, the code returns early from the transaction callback at line 134, but processMessage() continues execution and calls queueService.send(QueueName.ReplicaQueue, ...) at line 286. This means already-processed traces get re-enqueued downstream.

Consider returning a sentinel from the transaction or setting a flag to skip the best-effort enqueue.

🔧 Suggested fix
+      let shouldSkipL3 = false;
       await this.db.transaction(async (tx) => {
         // ... existing code ...

         if (inbound.status !== "RECEIVED" && inbound.status !== "PENDING") {
           this.logger.debug(
             { traceId, status: inbound.status },
             "L2 already processed this trace (idempotent redelivery). Skipping.",
           );
+          shouldSkipL3 = true;
           return;
         }
         // ... rest of transaction ...
       });

+      if (shouldSkipL3) {
+        this.logger.debug({ traceId }, "Skipping L3 enqueue for already-processed trace");
+        return;
+      }
+
       // Best-effort enqueue to L3 ...
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@apps/worker/src/modules/pipeline/replica.service.ts` around lines 129 - 135,
The idempotent check inside the transaction in processMessage() skips work but
doesn’t prevent the later enqueue to ReplicaQueue; modify processMessage() so
the transaction returns a sentinel or sets a local flag (e.g., shouldEnqueue =
false) when inbound.status !== "RECEIVED" && inbound.status !== "PENDING" inside
the transaction callback, then after the transaction completes only call
queueService.send(QueueName.ReplicaQueue, ...) if the sentinel/shouldEnqueue is
true; update the transaction callback and the post-transaction call site to
honor that flag to avoid re-enqueueing already-processed traces.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.

Inline comments:
In `@apps/api/src/db/database-manager.ts`:
- Around line 1071-1075: The workspace check uses an unfiltered query
(db.select().from(schema.uiWorkspaces).limit(1)) which may pick any workspace;
either document this behavior in the surrounding JSDoc to state "any existing
workspace will be used" or make the lookup deterministic by creating a seeded
workspace with a known ID in seed() and changing the check to filter by that ID
(e.g., query schema.uiWorkspaces where id = <seededWorkspaceId>); update
references to seed() and the workspace lookup to use the deterministic fixture
ID so the check is explicit and repeatable.

In `@apps/worker/src/modules/pipeline/delivery.service.ts`:
- Around line 414-427: The current log in DeliveryService still exposes raw
stack via "err.stack"; update the logger.error call to stop logging the raw
stack by using the sanitized error's stack instead (e.g., replace the "stack:
err instanceof Error ? err.stack : undefined" entry with "stack: sanitized.stack
?? undefined") or remove the stack field entirely; ensure sanitizeError (and its
callers) returns a safe stack string if you keep it, and reference the
sanitizeError function and the logger.error invocation in delivery.service.ts
when making the change.

In `@apps/worker/src/modules/pipeline/pipeline.module.ts`:
- Around line 71-73: The PipelineModule currently injects DeliveryService in its
constructor but never uses it; either remove the unused constructor parameter
(delete "private readonly deliveryService: DeliveryService" from the
PipelineModule constructor) or, if injection is intentional to force eager
instantiation, keep the parameter and add an explicit comment above the
constructor (e.g., "// intentionally injected to force eager instantiation of
DeliveryService") so the purpose is clear; locate the change in the
PipelineModule class constructor to apply one of these two fixes.

In `@infra/debezium/application.properties`:
- Around line 23-26: Debezium's config expects JWT endpoints but only a
Bearer-token relay exists; either implement the JWT endpoints or switch the
config. Add a controller with POST routes matching
/internal/cdc/auth/authenticate and /internal/cdc/auth/refreshToken that accept
the expected payloads and return JWTs (and refresh tokens) using the same secret
logic as the existing `@Post`('relay') guard, or alternatively change the
properties in application.properties to use a bearer/static token auth type to
match the current `@Post`('relay') validation; update the authentication logic to
reuse the existing validation (token generation/verification) functions so the
new endpoints produce tokens compatible with the guard.

In `@packages/credentials/src/oauth/token-refresh.service.ts`:
- Around line 62-67: The validateInputs function currently throws raw TypeError
which bypasses the OAuthRefreshError contract; update the caller flow so
validation failures are converted to OAuthRefreshError by either moving the
validateInputs(…) call into the existing try block in the public refresh routine
or wrapping the validateInputs call in a try/catch that catches TypeError and
rethrows a new OAuthRefreshError (include original error/message as
cause/context). Ensure references to validateInputs and OAuthRefreshError are
used so callers always receive OAuthRefreshError on invalid inputs.

---

Outside diff comments:
In `@apps/worker/src/modules/pipeline/fanout.service.ts`:
- Around line 508-529: The catch block for processSingleStitch logs the failure
and writes a FAIL sync log but never releases the entity lock, leaving the lock
held until TTL expiry; add the same lock-release logic used in the successful
path (the method you use when acquiring/releasing entity locks — e.g.,
this.releaseEntityLock or this.unlockEntity) and call it (await) in the catch
block (before or after writeSyncLog) for the same entity identifier used when
acquiring (referencing stitch.id or the entity id variable), handling any errors
from the release so failures here don't mask the original error; ensure the
release call mirrors the success-path release semantics and uses the same
schemaName/trace/entity id parameters.

---

Duplicate comments:
In `@apps/worker/src/modules/pipeline/replica.service.ts`:
- Around line 129-135: The idempotent check inside the transaction in
processMessage() skips work but doesn’t prevent the later enqueue to
ReplicaQueue; modify processMessage() so the transaction returns a sentinel or
sets a local flag (e.g., shouldEnqueue = false) when inbound.status !==
"RECEIVED" && inbound.status !== "PENDING" inside the transaction callback, then
after the transaction completes only call
queueService.send(QueueName.ReplicaQueue, ...) if the sentinel/shouldEnqueue is
true; update the transaction callback and the post-transaction call site to
honor that flag to avoid re-enqueueing already-processed traces.

In `@packages/credentials/src/oauth/token-refresh.service.ts`:
- Around line 84-92: The returned vendorParams currently spreads
vendorParamsSpread directly which can leave non-string values and violate the
OAuthCredentialBlob type; in the return where vendorParams is built (around
vendorParamsSpread and environmentEntry) replace the spread of
vendorParamsSpread with a string-coerced map (e.g. transform
Object.entries(vendorParamsSpread) to Object.fromEntries(entries.map(([k,v]) =>
[k, String(v)]))) so every value is a string, keep environmentEntry as
String(...) and ensure the final vendorParams is typed/constructed as
Record<string,string>.
🪄 Autofix (Beta)

✅ Autofix completed


ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: 1634a932-15db-4c56-ba6b-d0456211a262

📥 Commits

Reviewing files that changed from the base of the PR and between 7d9ccd2 and 2bf80f1.

📒 Files selected for processing (7)
  • apps/api/src/db/database-manager.ts
  • apps/worker/src/modules/pipeline/delivery.service.ts
  • apps/worker/src/modules/pipeline/fanout.service.ts
  • apps/worker/src/modules/pipeline/pipeline.module.ts
  • apps/worker/src/modules/pipeline/replica.service.ts
  • infra/debezium/application.properties
  • packages/credentials/src/oauth/token-refresh.service.ts

Comment on lines +1071 to +1075
// Check if a workspace exists
const workspaces = await db.select().from(schema.uiWorkspaces).limit(1);
if (workspaces.length === 0) {
throw new Error('No workspace found. Run pnpm db:seed first.');
}

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.

🧹 Nitpick | 🔵 Trivial

Consider documenting the workspace requirement.

The workspace lookup uses limit(1) without filtering by a deterministic ID, which could bind to any existing workspace. While the connection lookups now correctly use deterministic fixture IDs, the workspace requirement is less explicit. If a specific seeded workspace is expected, consider either:

  1. Documenting in the JSDoc that any workspace will be used, or
  2. Creating a deterministic workspace in seed() and filtering by its ID here
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@apps/api/src/db/database-manager.ts` around lines 1071 - 1075, The workspace
check uses an unfiltered query (db.select().from(schema.uiWorkspaces).limit(1))
which may pick any workspace; either document this behavior in the surrounding
JSDoc to state "any existing workspace will be used" or make the lookup
deterministic by creating a seeded workspace with a known ID in seed() and
changing the check to filter by that ID (e.g., query schema.uiWorkspaces where
id = <seededWorkspaceId>); update references to seed() and the workspace lookup
to use the deterministic fixture ID so the check is explicit and repeatable.

Comment on lines +414 to +427
const sanitized = sanitizeError(err);
this.logger.error(
{
event: "l5.error",
traceId,
routeId,
connectionId,
outboundGatewayId,
layer: "L5",
err: sanitized,
stack: err instanceof Error ? err.stack : undefined,
},
"DeliveryService encountered an unexpected error",
);

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.

⚠️ Potential issue | 🟡 Minor

Raw stack trace may still leak sensitive information.

While err is now sanitized at line 414, line 424 still logs the raw err.stack. Stack traces can contain file paths, credentials in closure variables, or vendor URLs. Consider omitting the stack or sanitizing it as well.

🔧 Suggested fix
       const sanitized = sanitizeError(err);
       this.logger.error(
         {
           event: "l5.error",
           traceId,
           routeId,
           connectionId,
           outboundGatewayId,
           layer: "L5",
           err: sanitized,
-          stack: err instanceof Error ? err.stack : undefined,
         },
         "DeliveryService encountered an unexpected error",
       );
📝 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
const sanitized = sanitizeError(err);
this.logger.error(
{
event: "l5.error",
traceId,
routeId,
connectionId,
outboundGatewayId,
layer: "L5",
err: sanitized,
stack: err instanceof Error ? err.stack : undefined,
},
"DeliveryService encountered an unexpected error",
);
const sanitized = sanitizeError(err);
this.logger.error(
{
event: "l5.error",
traceId,
routeId,
connectionId,
outboundGatewayId,
layer: "L5",
err: sanitized,
},
"DeliveryService encountered an unexpected error",
);
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@apps/worker/src/modules/pipeline/delivery.service.ts` around lines 414 - 427,
The current log in DeliveryService still exposes raw stack via "err.stack";
update the logger.error call to stop logging the raw stack by using the
sanitized error's stack instead (e.g., replace the "stack: err instanceof Error
? err.stack : undefined" entry with "stack: sanitized.stack ?? undefined") or
remove the stack field entirely; ensure sanitizeError (and its callers) returns
a safe stack string if you keep it, and reference the sanitizeError function and
the logger.error invocation in delivery.service.ts when making the change.

Comment thread apps/worker/src/modules/pipeline/pipeline.module.ts
Comment thread infra/debezium/application.properties Outdated
Comment thread packages/credentials/src/oauth/token-refresh.service.ts
@coderabbitai

coderabbitai Bot commented Apr 28, 2026 •

Copy link
Copy Markdown
Contributor

Note

Autofix is a beta feature. Expect some limitations and changes as we gather feedback and continue to improve it.

Fixes Applied Successfully

Fixed 4 file(s) based on 5 unresolved review comments.

Files modified:

  • apps/api/src/db/database-manager.ts
  • apps/worker/src/modules/pipeline/pipeline.module.ts
  • infra/debezium/application.properties
  • packages/credentials/src/oauth/token-refresh.service.ts

Commit: c93e10b009716520a628bb8371956638d8bb2372

The changes have been pushed to the feature/enterprise-canonical-pipeline branch.

Time taken: 5m 23s

Fixed 4 file(s) based on 5 unresolved review comments.

Co-authored-by: CodeRabbit <noreply@coderabbit.ai>
@pramodnarayana
pramodnarayana merged commit 7aa4346 into development Apr 28, 2026
2 checks passed
@coderabbitai coderabbitai Bot mentioned this pull request Jun 2, 2026
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