Skip to content

chore: remove obsolete files — engine/, packages/platform, worker pipeline/db modules, vitest.workspace - #155

Merged
pramodnarayana merged 7 commits into
developmentfrom
refactor/monolith-boundaries
Jun 11, 2026
Merged

pramodnarayana merged 7 commits into
developmentfrom
refactor/monolith-boundaries

Conversation

@pramodnarayana

@pramodnarayana pramodnarayana commented Jun 11, 2026 •

Copy link
Copy Markdown
Owner

Summary by CodeRabbit

  • Removed Features

    • CLI-based database lifecycle tools (reset/seed/migrate/provision/debug) and assorted local DB seeding/fixtures removed.
  • New Features

    • New observability and AI packages added.
    • Pipeline package exposes expanded delivery, retry, and global-entity mapping capabilities.
  • Refactor

    • Database responsibilities moved behind repository/port abstractions and shared packages; worker/pipeline modules reorganized to use centralized packages.

@coderabbitai

coderabbitai Bot commented Jun 11, 2026 •

Copy link
Copy Markdown
Contributor

Review Change Stack

Note

Reviews paused

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

Use the following commands to manage reviews:

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

Use the checkboxes below for quick actions:

  • ▶️ Resume reviews
  • 🔍 Trigger review
📝 Walkthrough

Walkthrough

Removes app-local DB CLI, providers, and seeding; rewires API/worker imports to shared packages; refactors pipeline (delivery, fanout, normalization, dependency sweeper) to repository/transaction ports; adds observability, DB migrations/snapshots, test helpers, and pipeline import-fix scripts.

Changes

Platform extraction and pipeline migration

Layer / File(s) Summary
Monorepo extraction & pipeline refactor
apps/api/src/db/..., apps/worker/src/..., packages/pipeline/src/..., packages/database/..., packages/observability/..., packages/ai/*, packages/*/vitest.config.ts, packages/pipeline/fix-*.cjs
Removes local DB CLI/provider/seeder modules and schema barrels; adds shared package manifests, observability logger, DB migrations/snapshots, TestDatabaseManager, pipeline import-rewrite scripts; refactors delivery, fanout, normalization, dependency-sweeper flows to repository/transaction ports and updates/adds tests.
Developer-facing checklist
engine/*, apps/*/bootstrap/*, apps/*/consumers/*, apps/*/pollers/*
Worker bootstrap and consumer/poller wiring moved to use shared DatabaseModule/ObservabilityModule and @soopa/pipeline exports; consumer behaviors unchanged aside from dependency import sources.

Sequence Diagram(s)

sequenceDiagram
  participant ClaimDeliveryUseCase
  participant DeliveryRetryService
  participant DeliveryService
  participant OutboundGatewayRepositoryPort
  participant TransactionManagerPort
  participant SyncLogRepositoryPort
  participant PipelineStateRepositoryPort

  ClaimDeliveryUseCase->>DeliveryRetryService: retrySourceFinalization(params...)
  DeliveryRetryService->>OutboundGatewayRepositoryPort: fetchOutboundGatewayResult(outboundGatewayId)
  DeliveryRetryService->>DeliveryService: writeL6Result(retryParams...)
  DeliveryService->>OutboundGatewayRepositoryPort: markResult(result...)
  DeliveryService->>TransactionManagerPort: runInTenantTransaction(tenantId, schema, fn)
  TransactionManagerPort->>SyncLogRepositoryPort: writeSyncLog(syncLogRow)
  TransactionManagerPort->>PipelineStateRepositoryPort: releaseSyncLockByTraceId(traceId)
Loading

Estimated code review effort

🎯 5 (Critical) | ⏱️ ~90+ minutes

Possibly related PRs

Poem

🐰 I hopped through code both near and far,
Removed old CLI leaves, lit a new shared star.
Logs now hum and pipelines sing,
Migrations march — a tidy spring.
A small rabbit cheers: refactor, done!

✨ Finishing Touches
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch refactor/monolith-boundaries

@pramodnarayana
pramodnarayana force-pushed the refactor/monolith-boundaries branch from 45428c5 to 50473df Compare June 11, 2026 06:10
@pramodnarayana

Copy link
Copy Markdown
Owner Author

@coderabbitai review

@coderabbitai

coderabbitai Bot commented Jun 11, 2026 •

Copy link
Copy Markdown
Contributor
✅ Action performed

Review finished.

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

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

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

Inline comments:
In `@apps/worker/src/bootstrap/dbmanager/dbmanager.module.ts`:
- Around line 109-122: The log may leak full DSNs because dbName falls back to
connectionString; change the logging to never include raw credentials by
deriving a sanitized identifier instead of the raw connectionString. In the
block using dbName and connectionString (the URL parsing and the console.error
call), ensure you extract only non-sensitive parts (e.g., hostname and
pathname/dbname via new URL(connectionString) or, on parse failure, replace
user:pass@ in the DSN with "[REDACTED_CREDENTIALS]" or use a fixed placeholder
like "<redacted-connection>") and log that sanitized value (keep err as-is); do
not log the original connectionString anywhere.

In `@apps/worker/src/consumers/copilot.worker.ts`:
- Around line 116-124: The non-OK branch currently calls this.redis.publish(...)
twice and then throws, which duplicates the terminal events because the outer
catch also publishes; remove the two redis.publish(...) calls in the
webResponse.ok false branch (keep only throw new Error(`Non-OK response from
orchestrator: ${errorText}`)) so that the outer catch (the handler around the
call that references data.jobId and this.redis.publish) is solely responsible
for emitting the `error:` and `[DONE]` messages to the
`job:stream:${data.jobId}` channel.
- Around line 114-209: When webResponse.body is falsy the code currently does
nothing, leaving clients hanging; add an explicit handling branch after the
webResponse check that publishes an error message and a termination marker to
the SSE channel (use this.redis.publish with `job:stream:${data.jobId}` for both
an error payload and `[DONE]\n`), persist a failure/system message via
this.chatPersistence.appendMessage (include tenantId, conversationId,
role:"system" and a short error content), and then throw an Error to stop
processing; locate the logic around the webResponse.body check (references:
webResponse, this.redis.publish, this.chatPersistence.appendMessage, data.jobId)
and implement this early-return error flow.

In `@apps/worker/src/consumers/gitops-sync.worker.ts`:
- Around line 16-19: SHARD_BASE_PATH's fallback still points at the removed
"engine/sync/application" which causes the watcher to miss shards; update the
default fallback used in the SHARD_BASE_PATH constant (which reads
process.env.SHARD_APPLICATION_PATH) to the new location of the sync application
(for example replace "../../engine/sync/application" with the correct
"../../sync/application" or the repo's actual sync/application path), keeping
the environment override behavior intact so SHARD_APPLICATION_PATH still takes
precedence.

In `@apps/worker/src/pollers/inbound-outbox.poller.ts`:
- Line 1: Change the import so OutboxTable is explicitly a type-only import:
update the import that currently brings in BaseOutboxPoller and OutboxTable to
import OutboxTable with the type modifier (so only BaseOutboxPoller is a value
import and OutboxTable is imported as a type). This affects the import line that
references BaseOutboxPoller and OutboxTable and the place where OutboxTable is
used as a type assertion; ensure only the type is imported for OutboxTable
(e.g., import { BaseOutboxPoller, type OutboxTable } ...) so the intent is
unambiguous to TypeScript.

In `@apps/worker/src/pollers/normalized-outbox.poller.ts`:
- Line 1: Change the import of OutboxTable to a type-only import so it’s not
treated as a runtime value: update the import statement that currently brings in
BaseOutboxPoller and OutboxTable to import OutboxTable using the `type` keyword,
and ensure the type assertion `normalizedOutbox as unknown as OutboxTable`
continues to compile without bringing OutboxTable into the emitted JS; locate
the import in this file (referencing BaseOutboxPoller and OutboxTable) and make
the OutboxTable import type-only to match its usage.

In `@packages/database/drizzle/global/0018_gorgeous_alice.sql`:
- Line 1: Before adding the unique constraint on workspace_pieces
(ux_workspace_pieces_workspaceId_pieceId over columns workspace_id and
piece_id), add a deterministic pre-migration step that either: 1) detects
duplicates with a GROUP BY workspace_id, piece_id and fails fast with an
explicit error listing (or count) so the deploy can be fixed manually, or 2)
performs a deterministic dedupe/backfill that deletes or consolidates duplicate
rows (keeping the row with the lowest id or latest updated_at) and records the
removed ids in an audit table; implement this logic in the same migration prior
to the ALTER TABLE so that the ADD CONSTRAINT will not fail due to existing
duplicates. Ensure you reference workspace_pieces, workspace_id, piece_id and
ux_workspace_pieces_workspaceId_pieceId when adding the precheck/dedupe.

In `@packages/database/drizzle/tenant/0012_military_black_queen.sql`:
- Around line 1-2: The migration should remove any existing representation of
gem_unique_mapping_idx before adding it as a UNIQUE constraint on table
global_entity_map; change the script to first DROP CONSTRAINT IF EXISTS
"gem_unique_mapping_idx" (on global_entity_map) and also DROP INDEX IF EXISTS
"gem_unique_mapping_idx" if needed, then run ALTER TABLE "global_entity_map" ADD
CONSTRAINT "gem_unique_mapping_idx"
UNIQUE("stitch_id","source_data_source_id","source_entity_id","dest_data_source_id","dest_entity_type")
so the migration succeeds whether the object exists as an index or as a
constraint.

In `@packages/observability/package.json`:
- Around line 12-23: The package manifest's "dependencies" and "devDependencies"
blocks were changed causing lockfile drift; regenerate the workspace lockfile
and commit it so CI's frozen-lockfile check passes. From the repository root run
the pnpm workspace install command (e.g., pnpm install or pnpm -w install) to
update pnpm-lock.yaml to match the new packages/observability/package.json
specifiers, verify the lockfile changes, and commit the updated pnpm-lock.yaml
alongside the package.json change.
- Around line 13-15: observability package pins older NestJS deps causing mixed
major versions; update packages in packages/observability/package.json to match
workspace majors by bumping "`@nestjs/common`" to the workspace version (e.g.
^11.1.11) and "`@nestjs/config`" to the matching major (e.g. ^4.0.2), and ensure
peerDependencies (if present) for Nest packages align with those versions;
verify dev/build tooling (e.g. nestjs-pino) remains compatible with the newer
Nest major and run install/typecheck to catch any remaining incompatibilities.

In `@packages/pipeline/fix-imports.cjs`:
- Around line 14-30: The replacement regexes in
packages/pipeline/fix-imports.cjs currently only match double-quoted imports
(e.g. patterns like /from "\.\.\/utils\.js"/g) and miss single-quoted imports;
update each affected pattern (e.g. the ones referencing ports, adapters,
interfaces, outbox.utils.js, ../ports/, ../adapters/, ../interfaces/,
../outbox.utils.js, /from "\.\.\/utils\.js"/, /from "\.\.\/index\.js"/, /from
"\.\.\/storage-resolver\//, /from "\.\.\/sharding\//) to be quote-agnostic by
matching either single or double quotes (use a regex like /from
['"]\.\.\/...['"]/ style) so all import variants are rewritten consistently
while leaving the gem-hydration.service.js rule unchanged.

In `@packages/pipeline/src/delivery/delivery.service.integration.spec.ts`:
- Around line 46-48: Remove the duplicate call to sqlManager.applyPlan for
SchemaPlan.OUTBOUND_ACTIVE: locate the consecutive calls to
sqlManager.applyPlan(currentSchemaName, SchemaPlan.OUTBOUND_ACTIVE, { appName:
"testApp", appProfile: "online" }) and delete the redundant second invocation so
the plan is applied only once during test setup; ensure only the NAMESPACE_ONLY
and a single OUTBOUND_ACTIVE call remain.
- Around line 239-296: The test is passing an obsolete tenantDb argument to
writeL6Result, which no longer accepts a DrizzleDb parameter; remove the final
testDbManager.db! argument from the writeL6Result call in this spec so the
argument list matches the updated writeL6Result signature (refer to
writeL6Result and the test that asserts gemService.writeGemMapping to locate the
call), and run the spec to ensure no other tests still pass tenantDb.

In `@packages/pipeline/src/delivery/use-cases/claim-delivery.use-case.spec.ts`:
- Around line 136-169: Add explicit assertions that
retryService.isSourceFinalized was called with the correct argument ordering to
guard against regression: after calling useCase.execute in the two
finalized-path tests (the ones setting outboundGatewayRepository entries to
status 'SUCCESS' and 'FAIL'), add
expect(retryService.isSourceFinalized).toHaveBeenCalledWith(tenantId,
srcSchemaName, traceId, routeId) using the variables from the test (or values
from defaultInput) to assert tenant/schema/trace/route ordering; ensure this
same assertion is added for the other similar tests around lines 100-134 as
well.

In `@packages/pipeline/src/delivery/use-cases/claim-delivery.use-case.ts`:
- Around line 179-185: The calls to retryService.isSourceFinalized are passing
arguments in the wrong order (currently calling with srcSchemaName, traceId,
routeId, tenantId) which breaks the finalization check; update both calls that
set sourceFinalized (the one guarded by currentStatus === "SUCCESS" and the
later call around the retry path) to pass arguments in the same order as the
retryService.isSourceFinalized signature—reorder srcSchemaName, traceId,
routeId, tenantId to match that signature (e.g., if the method expects traceId,
routeId, srcSchemaName, tenantId, swap the arguments accordingly) so the
parameters traceId, routeId, srcSchemaName, tenantId align with the method
definition. Ensure you change both usages of retryService.isSourceFinalized,
leaving variable names (currentStatus, sourceFinalized, srcSchemaName, traceId,
routeId, tenantId) intact.

In `@packages/pipeline/src/fanout/fanout-batch-processor.ts`:
- Around line 176-204: The SyncToken extraction in the destState handling is
overly nested and should be moved into a small helper to simplify the block:
implement a function (e.g., extractSyncTokenFromState(state): string |
undefined) that first returns state.SyncToken if present, otherwise finds the
first key not equal to "time" whose value is a non-null object and returns that
object's SyncToken; replace the inline logic that sets debugSyncToken in the
block that assigns destState and call this helper, then use the returned value
in the this.logger.log call (symbols to change: debugSyncToken, destState, and
the logging call inside the if (state) branch). Ensure behavior remains
identical and that this helper is only used for debug logging.
- Around line 282-289: The current code calls
outboxRepo.markOutboundGatewayFailed(...) and then throws sendErr, which ends up
swallowed by the outer catch and leaves outbound_gateway rows permanently
FAILED; fix by either propagating the error to the upstream queue (so the
message is retried) or immediately re-scheduling the gateway for delivery: after
catching the QueueName.DeliveryQueue publish failure in fanout-batch-processor,
replace the terminal FAILED-only path with a recovery action such as calling
outboxRepo.upsertPendingOutboundGateway(...) to set status to PENDING/RETRY or
re-enqueue the delivery via this.mq.publish(...) with the same payload, and only
mark FAILED when you've exhausted retries; adjust logic around
markOutboundGatewayFailed, sendErr handling, and the outer catch so failures are
retriable.
- Around line 340-358: The deferred-retry flow is broken because
claimDeferredTrace() only keys by traceId (not routeId) and marks rows PENDING
while upsertPendingOutboundGateway() only treats prior statuses
('DEFERRED_DEPENDENCY','FAILED',NULL) as publishable; fix by (1) changing
DependencySweeperService.claimDeferredTrace(...) to include routeId in its
WHERE/claim predicate so claims are per-route, and (2) update
upsertPendingOutboundGateway(...) to accept rows whose prior status is 'PENDING'
(in addition to 'DEFERRED_DEPENDENCY','FAILED',NULL) so re-queued work that was
claimed can be published to QueueName.DeliveryQueue; refer to
fanout-batch-processor.ts symbols claimDeferredTrace,
upsertPendingOutboundGateway, QueueName.ActiveFetchQueue and DeliveryQueue when
making these changes.

In `@packages/pipeline/src/fanout/fanout-router.service.ts`:
- Around line 157-159: Remove the temporary compatibility plumbing: delete the
tenantDb = await this.dbManager.getTenantDb(tenantId) and const { syncLog } =
buildTenantSchema(schemaName) lines in fanout-router.service.ts and remove the
corresponding unused parameters from batchProcessor.processSingleStitch's
signature and implementation in the fanout-batch-processor (the function/method
named processSingleStitch). Update any other call sites of processSingleStitch
to stop passing tenantDb/syncLog and remove internal references to those
parameters in the BatchProcessor/FanoutBatchProcessor code so the new
implementation no longer accepts or expects those unused args.
- Around line 87-117: The supersession check is not race-safe:
runInTenantTransaction (used in FanoutRouterService) doesn't set isolation and
code calls evaluateSuperseded then reads
getNormalizedData/getReplicaSourceVendorId without any DB locking or lock
acquisition, and the pipeline never inserts into active_sync_locks (only
releases). Fix by acquiring a DB-level lock for the (dataSourceId, entityId)
before evaluateSuperseded and holding it for the transaction—either (preferred)
insert/select-for-update on active_sync_locks (INSERT ... ON CONFLICT DO NOTHING
then SELECT FOR UPDATE) inside txManager.runInTenantTransaction or use
pg_advisory_xact_lock keyed by dataSourceId+entityId; ensure the lock
acquisition happens in FanoutRouterService before calling
routingDecisionEngine.evaluateSuperseded and
stateRepo.getNormalizedData/getReplicaSourceVendorId, and ensure
releaseSyncLock* remains consistent (or remove if using advisory locks); also
adjust runInTenantTransaction to allow explicit transaction isolation or ensure
lock acquisition occurs within the same transaction context so the supersession
check + subsequent reads and eventual upsert
(drizzle-outbound-gateway.repository) are atomic.

In `@packages/pipeline/src/fanout/target-builder.service.ts`:
- Around line 118-122: The current check in target-builder.service.ts only tests
Object.keys(hydrated).length === 0 which misses payloads like {a: undefined, b:
null}; update the validation around the hydrated object (the variable named
hydrated and the error that references normalizedEntityType) to ensure at least
one property has a non-null/undefined value (e.g. use
Object.values(hydrated).some(v => v != null)) before accepting the payload, and
throw the same error including normalizedEntityType if no valid values remain;
optionally consider stripping null/undefined properties from hydrated before
further processing to avoid sending invalid fields downstream.
- Line 28: The constructor parameter on TargetBuilderService uses
`@Inject`(PipelineHookBrokerService) on the hookBroker parameter; confirm that
PipelineHookBrokerService is indeed registered/exported in
ApplicationLoaderModule and then remove the redundant `@Inject` decorator to use
the typed injection idiom (i.e., keep the constructor param named hookBroker:
PipelineHookBrokerService and remove the `@Inject`(...) decorator) unless you have
a custom provider token requiring explicit injection.

In `@packages/pipeline/src/index.ts`:
- Line 16: The export statements are concatenated into a single invalid line;
split them into two separate export statements so each module is exported on its
own line—separate "export * from './sharding/pipeline-hook-broker.service.js';"
and "export * from \"./pipeline-core.module.js\";" by adding a newline (or a
semicolon + newline) between them so both exports are valid.
- Around line 21-23: Remove the duplicate export of './utils.js' by keeping a
single export statement for utils (remove one of the two lines that read export
* from './utils.js';), so only one export * from './utils.js' remains in the
module index (ensure no other references to duplicate export remain).

In `@packages/pipeline/src/normalization/dependency-sweeper.service.ts`:
- Around line 107-141: The current claim-then-send flow (claimDeferredTrace ->
queueService.send) can leave traces permanently stuck if send fails; change the
logic so the DB claim is reverted on publish failure and only mark
processedTraceIds after a successful send: after calling
this.sweeperRepo.claimDeferredTrace(tenant.tenantId, schemaName, traceId) and
before adding to processedTraceIds, wrap this.queueService.send in a try/catch
that on failure calls a new or existing sweeperRepo.unclaimDeferredTrace (or
releaseClaim/rollbackClaim) for the same tenant/schema/traceId, logs the
failure, and rethrows or continues without adding to processedTraceIds; ensure
processedTraceIds.add(traceId) happens only after send succeeds.

In `@packages/pipeline/src/normalization/normalization.service.ts`:
- Around line 189-196: The .catch that re-throws in the call to
this.hookBroker.normalize (which assigns normalizedFromShard) is misleading and
redundant; remove the .catch((err:any)=>{ throw err; }) and the accompanying
comment about "fall through to piece.normalize" so errors correctly propagate
from hookBroker.normalize, or alternatively if you intended to swallow errors
return null in the catch — update the code in normalization.service.ts around
the this.hookBroker.normalize call and adjust the normalizedFromShard handling
and comment accordingly (referenced symbols: this.hookBroker.normalize,
normalizedFromShard, piece.normalize).
- Around line 262-278: The queueService.send calls that re-queue parent traces
(inside the loop over parentTraceIds) are currently executed within the
transaction started in thisNormalization flow, causing a dual-write risk if the
transaction later rolls back (e.g., insertNormalizedOutboxPending). Move the
reverse-lookup queue sends out of the transaction: collect parentTraceIds (and
associated metadata like replica.entityId and canonicalType) during the
transactional section, commit the transaction, then iterate and call
queueService.send and logger.debug afterwards (mirroring the L3→L4 handoff
pattern). Ensure references to queueService.send, parentTraceIds,
replica.entityId, canonicalType, and insertNormalizedOutboxPending are used to
locate and refactor the code.
🪄 Autofix (Beta)

❌ Autofix failed (check again to retry)

Fix all unresolved CodeRabbit comments on this PR:

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

ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: 1ffc8f77-3818-423a-b5dd-aa25581c289d

📥 Commits

Reviewing files that changed from the base of the PR and between 50473df and f1698ac.

📒 Files selected for processing (120)
  • apps/worker/src/bootstrap/dbmanager/dbmanager.module.ts
  • apps/worker/src/bootstrap/observability/observability.module.ts
  • apps/worker/src/consumers/active-fetch.worker.spec.ts
  • apps/worker/src/consumers/active-fetch.worker.ts
  • apps/worker/src/consumers/app-installer.processor.spec.ts
  • apps/worker/src/consumers/app-installer.processor.ts
  • apps/worker/src/consumers/copilot.worker.ts
  • apps/worker/src/consumers/gitops-sync.worker.spec.ts
  • apps/worker/src/consumers/gitops-sync.worker.ts
  • apps/worker/src/cron/app-updater.cron.spec.ts
  • apps/worker/src/cron/app-updater.cron.ts
  • apps/worker/src/pollers/base-outbox.poller.spec.ts
  • apps/worker/src/pollers/base-outbox.poller.ts
  • apps/worker/src/pollers/inbound-outbox.poller.spec.ts
  • apps/worker/src/pollers/inbound-outbox.poller.ts
  • apps/worker/src/pollers/normalized-outbox.poller.spec.ts
  • apps/worker/src/pollers/normalized-outbox.poller.ts
  • apps/worker/src/pollers/registry-outbox.poller.spec.ts
  • apps/worker/src/pollers/registry-outbox.poller.ts
  • apps/worker/src/pollers/replica-outbox.poller.spec.ts
  • apps/worker/src/pollers/replica-outbox.poller.ts
  • packages/ai/package.json
  • packages/ai/src/ai-engine.module.spec.ts
  • packages/ai/src/ai-engine.module.ts
  • packages/ai/src/categories/mapping.service.ts
  • packages/ai/src/contracts/chat-request.types.ts
  • packages/ai/src/contracts/prompts.ts
  • packages/ai/src/index.ts
  • packages/ai/src/planner/intent-classifier.service.ts
  • packages/ai/src/runtime/orchestrator.service.ts
  • packages/ai/src/services/chat-persistence.service.ts
  • packages/ai/src/services/transformer-simulation.service.ts
  • packages/ai/src/tools/action-tool.factory.ts
  • packages/ai/src/tools/hydrator-tool.factory.ts
  • packages/ai/src/tools/tool-zod.wrapper.ts
  • packages/ai/src/transformers/token-optimizer.util.ts
  • packages/ai/src/transformers/transformer.engine.ts
  • packages/ai/tsconfig.json
  • packages/ai/vitest.config.ts
  • packages/auth/vitest.config.ts
  • packages/database/drizzle/global/0018_gorgeous_alice.sql
  • packages/database/drizzle/global/meta/0018_snapshot.json
  • packages/database/drizzle/tenant/0012_military_black_queen.sql
  • packages/database/drizzle/tenant/meta/0012_snapshot.json
  • packages/dbmanager/vitest.config.ts
  • packages/domain/core/vitest.config.ts
  • packages/domain/tms/vitest.config.ts
  • packages/observability/package.json
  • packages/observability/src/index.ts
  • packages/observability/src/logger.module.ts
  • packages/observability/tsconfig.json
  • packages/pipeline/fix-imports.cjs
  • packages/pipeline/fix-module.cjs
  • packages/pipeline/fix-module2.cjs
  • packages/pipeline/revert-imports.cjs
  • packages/pipeline/src/delivery/delivery-retry.service.integration.spec.ts
  • packages/pipeline/src/delivery/delivery-retry.service.spec.ts
  • packages/pipeline/src/delivery/delivery-retry.service.ts
  • packages/pipeline/src/delivery/delivery.service.integration.spec.ts
  • packages/pipeline/src/delivery/delivery.service.spec.ts
  • packages/pipeline/src/delivery/delivery.service.ts
  • packages/pipeline/src/delivery/gem-hydration.service.integration.spec.ts
  • packages/pipeline/src/delivery/gem-hydration.service.ts
  • packages/pipeline/src/delivery/piece-outbound.dispatcher.spec.ts
  • packages/pipeline/src/delivery/piece-outbound.dispatcher.ts
  • packages/pipeline/src/delivery/use-cases/claim-delivery.use-case.integration.spec.ts
  • packages/pipeline/src/delivery/use-cases/claim-delivery.use-case.spec.ts
  • packages/pipeline/src/delivery/use-cases/claim-delivery.use-case.ts
  • packages/pipeline/src/evaluator.spec.ts
  • packages/pipeline/src/evaluator.ts
  • packages/pipeline/src/fanout/fanout-batch-processor.spec.ts
  • packages/pipeline/src/fanout/fanout-batch-processor.ts
  • packages/pipeline/src/fanout/fanout-router.service.spec.ts
  • packages/pipeline/src/fanout/fanout-router.service.ts
  • packages/pipeline/src/fanout/routing-decision.engine.spec.ts
  • packages/pipeline/src/fanout/routing-decision.engine.ts
  • packages/pipeline/src/fanout/target-builder.service.integration.spec.ts
  • packages/pipeline/src/fanout/target-builder.service.spec.ts
  • packages/pipeline/src/fanout/target-builder.service.ts
  • packages/pipeline/src/hydrator.spec.ts
  • packages/pipeline/src/hydrator.ts
  • packages/pipeline/src/index.ts
  • packages/pipeline/src/normalization/dependency-sweeper.service.spec.ts
  • packages/pipeline/src/normalization/dependency-sweeper.service.ts
  • packages/pipeline/src/normalization/normalization.service.spec.ts
  • packages/pipeline/src/normalization/normalization.service.ts
  • packages/pipeline/src/pipeline-core.module.ts
  • packages/pipeline/src/replication/registry-replication.service.integration.spec.ts
  • packages/pipeline/src/replication/registry-replication.service.spec.ts
  • packages/pipeline/src/replication/registry-replication.service.ts
  • packages/pipeline/src/replication/registry-token-refresh.service.integration.spec.ts
  • packages/pipeline/src/replication/registry-token-refresh.service.ts
  • packages/pipeline/src/replication/replica.service.integration.spec.ts
  • packages/pipeline/src/replication/replica.service.spec.ts
  • packages/pipeline/src/replication/replica.service.ts
  • packages/pipeline/src/sharding/application-loader.module.ts
  • packages/pipeline/src/sharding/application-loader.service.spec.ts
  • packages/pipeline/src/sharding/application-loader.service.ts
  • packages/pipeline/src/sharding/application-shard-event.handler.spec.ts
  • packages/pipeline/src/sharding/application-shard-event.handler.ts
  • packages/pipeline/src/sharding/application-shard.types.ts
  • packages/pipeline/src/sharding/pipeline-hook-broker.service.spec.ts
  • packages/pipeline/src/sharding/pipeline-hook-broker.service.ts
  • packages/pipeline/src/shared/adapters/drizzle-connection.repository.ts
  • packages/pipeline/src/shared/adapters/drizzle-dependency-sweeper.repository.ts
  • packages/pipeline/src/shared/adapters/drizzle-fanout-repositories.integration.spec.ts
  • packages/pipeline/src/shared/adapters/drizzle-field-mapping.repository.ts
  • packages/pipeline/src/shared/adapters/drizzle-global-entity-map.repository.ts
  • packages/pipeline/src/shared/adapters/drizzle-normalization.repository.ts
  • packages/pipeline/src/shared/adapters/drizzle-outbound-gateway.repository.ts
  • packages/pipeline/src/shared/adapters/drizzle-pipeline-state.repository.ts
  • packages/pipeline/src/shared/adapters/drizzle-routing.repository.integration.spec.ts
  • packages/pipeline/src/shared/adapters/drizzle-routing.repository.ts
  • packages/pipeline/src/shared/adapters/drizzle-stitch.repository.ts
  • packages/pipeline/src/shared/adapters/drizzle-sync-log.repository.ts
  • packages/pipeline/src/shared/adapters/drizzle-transaction-manager.adapter.ts
  • packages/pipeline/src/shared/adapters/outbound-gateway.adapter.integration.spec.ts
  • packages/pipeline/src/shared/adapters/outbound-gateway.adapter.ts
  • packages/pipeline/src/shared/adapters/registry-replication.adapter.integration.spec.ts
  • packages/pipeline/src/shared/adapters/registry-replication.adapter.ts

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

Caution

Inline review comments failed to post. This is likely due to GitHub's internal server error or limits when posting large numbers of comments. If you are seeing this consistently it is likely a permissions issue. Please check "Moderation" -> "Code review limits" under your organization settings.

🛑 Comments failed to post (27)
apps/worker/src/bootstrap/dbmanager/dbmanager.module.ts (1)

109-122: ⚠️ Potential issue | 🟠 Major | ⚡ Quick win

Avoid credential leakage in pool error logs.

On Line 109, dbName defaults to the full connectionString; if URL parsing fails, Line 120 can log credentials. Redact or avoid raw DSNs in logs.

Suggested fix
-              let dbName = connectionString;
+              let dbName = "unknown-db";
               try {
                 const url = new URL(connectionString);
                 dbName =
                   url.pathname.replace(/^\/+/, "") ||
                   url.searchParams.get("dbname") ||
-                  connectionString;
+                  "unknown-db";
               } catch {
                 // Fall back to original connectionString if URL parsing fails
               }
📝 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.

              let dbName = "unknown-db";
              try {
                const url = new URL(connectionString);
                dbName =
                  url.pathname.replace(/^\/+/, "") ||
                  url.searchParams.get("dbname") ||
                  "unknown-db";
              } catch {
                // Fall back to original connectionString if URL parsing fails
              }
              console.error(
                `Unexpected error on idle tenant DB client [${dbName}]`,
                err,
              );
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@apps/worker/src/bootstrap/dbmanager/dbmanager.module.ts` around lines 109 -
122, The log may leak full DSNs because dbName falls back to connectionString;
change the logging to never include raw credentials by deriving a sanitized
identifier instead of the raw connectionString. In the block using dbName and
connectionString (the URL parsing and the console.error call), ensure you
extract only non-sensitive parts (e.g., hostname and pathname/dbname via new
URL(connectionString) or, on parse failure, replace user:pass@ in the DSN with
"[REDACTED_CREDENTIALS]" or use a fixed placeholder like
"<redacted-connection>") and log that sanitized value (keep err as-is); do not
log the original connectionString anywhere.
apps/worker/src/consumers/copilot.worker.ts (2)

114-209: ⚠️ Potential issue | 🟠 Major | ⚡ Quick win

Handle missing response bodies explicitly to prevent silent job completion.

At Line 114, when webResponse.body is falsy, the method exits without throwing, without publishing a terminal event, and without persistence. That can acknowledge a job while clients wait indefinitely.

Proposed fix
-      if (webResponse.body) {
+      if (!webResponse.body) {
+        throw new Error("Orchestrator response has no body stream");
+      }
+      {
         // Check for non-OK response
         if (!webResponse.ok) {
📝 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.

      if (!webResponse.body) {
        throw new Error("Orchestrator response has no body stream");
      }
      {
        // Check for non-OK response
        if (!webResponse.ok) {
          const errorText = await webResponse.text();
          await this.redis.publish(
            `job:stream:${data.jobId}`,
            `error: ${errorText}\n`,
          );
          await this.redis.publish(`job:stream:${data.jobId}`, `[DONE]\n`);
          throw new Error(`Non-OK response from orchestrator: ${errorText}`);
        }

        // Read stream to exhaustion
        const reader = webResponse.body.getReader();
        const decoder = new TextDecoder("utf-8", { fatal: false });
        let finalResponseBuilder = "";

        while (true) {
          const { done, value } = await reader.read();
          if (done) break;

          // Process streams here for SSE broadcast (Step 4)
          // `value` is a Uint8Array containing Vercel Stream Parts.
          if (value) {
            const decodedChunk = decoder.decode(value, { stream: true });
            finalResponseBuilder += decodedChunk;
            // Publish standard Vercel Stream Parts out to SSE bridge
            await this.redis.publish(`job:stream:${data.jobId}`, decodedChunk);
          }
        }

        // Flush any remaining bytes from the decoder
        const finalChunk = decoder.decode();
        if (finalChunk) {
          finalResponseBuilder += finalChunk;
        }

        // Publish termination marker for SSE Client
        await this.redis.publish(`job:stream:${data.jobId}`, `[DONE]\n`);

        this.logger.debug(
          `Stream fully consumed. Payload length: ${finalResponseBuilder.length}`,
        );

        // Reconstruct human-readable response and step metadata for persistent storage
        let humanResponse = "";
        const stepMetadata: Record<string, unknown>[] = [];
        const lines = finalResponseBuilder.split("\n");
        for (const line of lines) {
          if (line.trim().startsWith("0:")) {
            try {
              humanResponse += JSON.parse(line.trim().substring(2));
            } catch {
              // Ignore partial parsing errors
            }
          } else if (line.trim().startsWith("8:")) {
            try {
              const stepData = JSON.parse(line.trim().substring(2)) as Record<
                string,
                unknown
              >;
              stepMetadata.push(stepData);
            } catch {
              // Ignore partial parsing errors
            }
          }
        }

        if (humanResponse.trim()) {
          await this.chatPersistence.appendMessage({
            tenantId: data.tenantId,
            conversationId: data.conversationId,
            role: "assistant",
            content: humanResponse.trim(),
            status: "completed",
          });
        }

        // Persist step metadata alongside the final response if any steps were captured
        if (stepMetadata.length > 0) {
          this.logger.debug(
            `Captured ${stepMetadata.length} step metadata entries for conversation ${data.conversationId}`,
          );
          // Store step metadata in the same persistence layer
          // Note: This could be stored as a system message or in a dedicated step metadata table
          await this.chatPersistence.appendMessage({
            tenantId: data.tenantId,
            conversationId: data.conversationId,
            role: "system",
            content: JSON.stringify({ steps: stepMetadata }),
            status: "completed",
          });
        }

        this.logger.log(`Successfully completed AI Job ${data.jobId}`);
      }
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@apps/worker/src/consumers/copilot.worker.ts` around lines 114 - 209, When
webResponse.body is falsy the code currently does nothing, leaving clients
hanging; add an explicit handling branch after the webResponse check that
publishes an error message and a termination marker to the SSE channel (use
this.redis.publish with `job:stream:${data.jobId}` for both an error payload and
`[DONE]\n`), persist a failure/system message via
this.chatPersistence.appendMessage (include tenantId, conversationId,
role:"system" and a short error content), and then throw an Error to stop
processing; locate the logic around the webResponse.body check (references:
webResponse, this.redis.publish, this.chatPersistence.appendMessage, data.jobId)
and implement this early-return error flow.

116-124: ⚠️ Potential issue | 🟡 Minor | ⚡ Quick win

Avoid double-publishing error/[DONE] events on non-OK responses.

At Lines 116–124, the non-OK branch publishes to Redis and then throws; the catch block at Lines 210–220 publishes again. This sends duplicate terminal events for one failure.

Proposed fix
         if (!webResponse.ok) {
           const errorText = await webResponse.text();
-          await this.redis.publish(
-            `job:stream:${data.jobId}`,
-            `error: ${errorText}\n`,
-          );
-          await this.redis.publish(`job:stream:${data.jobId}`, `[DONE]\n`);
           throw new Error(`Non-OK response from orchestrator: ${errorText}`);
         }

Also applies to: 210-220

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

In `@apps/worker/src/consumers/copilot.worker.ts` around lines 116 - 124, The
non-OK branch currently calls this.redis.publish(...) twice and then throws,
which duplicates the terminal events because the outer catch also publishes;
remove the two redis.publish(...) calls in the webResponse.ok false branch (keep
only throw new Error(`Non-OK response from orchestrator: ${errorText}`)) so that
the outer catch (the handler around the call that references data.jobId and
this.redis.publish) is solely responsible for emitting the `error:` and `[DONE]`
messages to the `job:stream:${data.jobId}` channel.
apps/worker/src/consumers/gitops-sync.worker.ts (1)

16-19: ⚠️ Potential issue | 🟠 Major | ⚡ Quick win

Update the fallback shard path; it still targets removed engine/ modules.

At Lines 16–19, fallback path is ../../engine/sync/application, but this PR removes engine/. If SHARD_APPLICATION_PATH is unset, sync/watcher points to an obsolete directory and misses real shard updates.

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

In `@apps/worker/src/consumers/gitops-sync.worker.ts` around lines 16 - 19,
SHARD_BASE_PATH's fallback still points at the removed "engine/sync/application"
which causes the watcher to miss shards; update the default fallback used in the
SHARD_BASE_PATH constant (which reads process.env.SHARD_APPLICATION_PATH) to the
new location of the sync application (for example replace
"../../engine/sync/application" with the correct "../../sync/application" or the
repo's actual sync/application path), keeping the environment override behavior
intact so SHARD_APPLICATION_PATH still takes precedence.
apps/worker/src/pollers/inbound-outbox.poller.ts (1)

1-1: 🧹 Nitpick | 🔵 Trivial

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "== Inspect OutboxTable export kind =="
rg -n -C2 'export\s+(type|interface|class|const|let|var)\s+OutboxTable' apps/worker/src/pollers/base-outbox.poller.ts

echo
echo "== Inspect current import site =="
rg -n -C2 'import\s*\{[^}]*OutboxTable[^}]*\}\s*from\s*"\./base-outbox\.poller\.js"' apps/worker/src/pollers/inbound-outbox.poller.ts

echo
echo "== Inspect TS module import-preservation settings =="
rg -n --iglob 'tsconfig*.json' '"verbatimModuleSyntax"|"preserveValueImports"|"importsNotUsedAsValues"'

Repository: pramodnarayana/nexiom

Length of output: 602


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "== Find tsconfig files in worker =="
fd -t f -e json 'tsconfig*.json' apps/worker | sort

echo
echo "== Relevant compiler options in worker tsconfig(s) =="
for f in $(fd -t f -e json 'tsconfig*.json' apps/worker | sort); do
  echo "--- $f ---"
  node -e "const fs=require('fs'); const p=process.argv[1]; const j=JSON.parse(fs.readFileSync(p,'utf8')); console.log(JSON.stringify({verbatimModuleSyntax:j.compilerOptions?.verbatimModuleSyntax, preserveValueImports:j.compilerOptions?.preserveValueImports, importsNotUsedAsValues:j.compilerOptions?.importsNotUsedAsValues, isolatedModules:j.compilerOptions?.isolatedModules, verbatimModuleSyntax:j.compilerOptions?.verbatimModuleSyntax, module:j.compilerOptions?.module, moduleResolution:j.compilerOptions?.moduleResolution, target:j.compilerOptions?.target}, null, 2));" "$f"
done

echo
echo "== Inspect if inbound-outbox.poller uses OutboxTable in a type-only way =="
rg -n "OutboxTable" apps/worker/src/pollers/inbound-outbox.poller.ts

Repository: pramodnarayana/nexiom

Length of output: 536


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "== Search for OutboxTable imports elsewhere =="
rg -n --iglob 'apps/worker/src/**/*.ts' 'OutboxTable' apps/worker/src | head -n 50

echo
echo "== Search for 'type OutboxTable' usage =="
rg -n --iglob 'apps/worker/src/**/*.ts' 'type\s+OutboxTable' apps/worker/src || true

echo
echo "== Search for existing pattern 'type' specifier imports in worker pollers =="
rg -n --iglob 'apps/worker/src/pollers/**/*.ts' 'import\s*\{[^}]*,\s*type\s' apps/worker/src/pollers || true

echo
echo "== Inspect base-outbox.poller.ts around OutboxTable export =="
sed -n '1,80p' apps/worker/src/pollers/base-outbox.poller.ts

Repository: pramodnarayana/nexiom

Length of output: 4045


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "== Show apps/worker/tsconfig.json =="
cat apps/worker/tsconfig.json

echo
echo "== Find root tsconfig files =="
ls -1 tsconfig*.json || true
ls -1 **/tsconfig*.json 2>/dev/null | head -n 50 || true

echo
echo "== Grep for verbatimModuleSyntax / preserveValueImports / importsNotUsedAsValues across repo tsconfigs =="
rg -n --glob '*/tsconfig*.json' '"verbatimModuleSyntax"|"preserveValueImports"|"importsNotUsedAsValues"' tsconfig*.json apps/worker apps -S || true

echo
echo "== Locate tsconfig base/extends chain from apps/worker/tsconfig.json =="
python3 - <<'PY'
import json, pathlib
p=pathlib.Path("apps/worker/tsconfig.json")
j=json.loads(p.read_text())
print("extends:", j.get("extends"))
PY

Repository: pramodnarayana/nexiom

Length of output: 1136


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "== Show tsconfig.base.json =="
cat tsconfig.base.json

echo
echo "== Search for tsconfig.*worker overrides in repo =="
rg -n --glob 'tsconfig*.json' '"importsNotUsedAsValues"|"preserveValueImports"|"verbatimModuleSyntax"' .

Repository: pramodnarayana/nexiom

Length of output: 637


Make OutboxTable import explicitly type-only (low current risk).

OutboxTable is exported as a TS interface in apps/worker/src/pollers/base-outbox.poller.ts, and inbound-outbox.poller.ts uses it only in a type assertion at line 135. With apps/worker/tsconfig.json not enabling verbatimModuleSyntax/preserveValueImports, TS should elide the value import under normal settings—still, prefer import { BaseOutboxPoller, type OutboxTable } ... at line 1 to make this intention unambiguous.

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

In `@apps/worker/src/pollers/inbound-outbox.poller.ts` at line 1, Change the
import so OutboxTable is explicitly a type-only import: update the import that
currently brings in BaseOutboxPoller and OutboxTable to import OutboxTable with
the type modifier (so only BaseOutboxPoller is a value import and OutboxTable is
imported as a type). This affects the import line that references
BaseOutboxPoller and OutboxTable and the place where OutboxTable is used as a
type assertion; ensure only the type is imported for OutboxTable (e.g., import {
BaseOutboxPoller, type OutboxTable } ...) so the intent is unambiguous to
TypeScript.
apps/worker/src/pollers/normalized-outbox.poller.ts (1)

1-1: 🧹 Nitpick | 🔵 Trivial

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "1) Check how OutboxTable is exported:"
fd -i 'base-outbox.poller.ts' --exec sed -n '1,240p' {}

echo "2) Check TS emit/import flags affecting type-only imports:"
fd -i 'tsconfig*.json' --exec rg -n '"verbatimModuleSyntax"|"preserveValueImports"|"importsNotUsedAsValues"' {}

echo "3) Find pollers using non-type OutboxTable imports:"
rg -n 'import\s+\{\s*BaseOutboxPoller,\s*OutboxTable\s*\}\s+from\s+"\.\/base-outbox\.poller\.js"' -C2

Repository: pramodnarayana/nexiom

Length of output: 4788


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# 1) Inspect the target poller file import + OutboxTable usage
echo "1) normalized-outbox.poller.ts import/usages:"
norm_file="$(fd -i 'normalized-outbox\.poller\.ts$' . | head -n 1 || true)"
if [[ -z "${norm_file}" ]]; then
  echo "NOT_FOUND: normalized-outbox.poller.ts"
else
  echo "FOUND: ${norm_file}"
  sed -n '1,200p' "${norm_file}"
  echo
  echo "---- OutboxTable type/value usages in ${norm_file} ----"
  rg -n '\bOutboxTable\b' "${norm_file}" || true
fi

# 2) Check TS config flags that affect whether non-type imports are preserved
echo
echo "2) tsconfig compiler options affecting type-only imports:"
mapfile -t tsconfigs < <(fd -i 'tsconfig*.json$' . || true)
if [[ ${`#tsconfigs`[@]} -eq 0 ]]; then
  echo "NOT_FOUND: no tsconfig*.json files"
else
  printf "Found %s tsconfig file(s):\n" "${`#tsconfigs`[@]}"
  printf "%s\n" "${tsconfigs[@]}"
  echo
  # print only relevant lines (and some surrounding context)
  for f in "${tsconfigs[@]}"; do
    echo "---- ${f} ----"
    rg -n '"verbatimModuleSyntax"|"preserveValueImports"|"importsNotUsedAsValues"' "${f}" || echo "no relevant flags"
  done
fi

# 3) Search for non-type import pattern across pollers
echo
echo "3) Non-type import pattern of OutboxTable from ./base-outbox.poller.js:"
rg -n 'import\s*\{\s*[^}]*\bBaseOutboxPoller\b[^}]*,\s*[^}]*\bOutboxTable\b[^}]*\}\s*from\s*"\.\/base-outbox\.poller\.js"' -S . || true

# 4) Confirm runtime exports for OutboxTable (in TS source)
echo
echo "4) base-outbox.poller.ts OutboxTable declaration kind:"
base_file="$(fd -i 'base-outbox\.poller\.ts$' . | head -n 1 || true)"
if [[ -n "${base_file}" ]]; then
  rg -n 'export (interface|type|class|const|function)\s+OutboxTable\b' "${base_file}" || true
fi

Repository: pramodnarayana/nexiom

Length of output: 5850


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "2) tsconfig compiler options affecting type-only imports:"
# Collect tsconfig files safely (avoid process substitution)
tsconfigs="$(fd -i 'tsconfig*.json$' . || true)"
if [[ -z "${tsconfigs}" ]]; then
  echo "NOT_FOUND: no tsconfig*.json files"
else
  echo "${tsconfigs}" | sed 's/^/ - /'
  echo
  while IFS= read -r f; do
    [[ -z "${f}" ]] && continue
    echo "---- ${f} ----"
    rg -n '"verbatimModuleSyntax"|"preserveValueImports"|"importsNotUsedAsValues"' "${f}" || echo "no relevant flags"
    echo
  done <<< "${tsconfigs}"
fi

echo
echo "3) Non-type import pattern of OutboxTable from ./base-outbox.poller.js (any file):"
rg -n 'import\s*\{\s*BaseOutboxPoller\s*,\s*OutboxTable\s*\}\s*from\s*"\.\/base-outbox\.poller\.js"' -S . || true

echo
echo "4) base-outbox.poller.ts OutboxTable declaration kind:"
base_file="$(fd -i 'base-outbox\.poller\.ts$' . | head -n 1 || true)"
if [[ -n "${base_file}" ]]; then
  rg -n 'export (interface|type|class|const|function)\s+OutboxTable\b' "${base_file}" || true
else
  echo "NOT_FOUND: base-outbox.poller.ts"
fi

Repository: pramodnarayana/nexiom

Length of output: 3535


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "1) Search all JSON for strict module/type-import flags:"
rg -n --hidden -S '"verbatimModuleSyntax"|"preserveValueImports"|"importsNotUsedAsValues"' . --glob '*.json' || true

echo
echo "2) Show worker tsconfig compilerOptions (module/moduleResolution + any type/import flags):"
worker_tsconfig="$(fd -i 'apps/worker/tsconfig.json$|tsconfig.json$' . | rg -n 'apps/worker/tsconfig\.json$|^./tsconfig\.json$' -o || true)"
# safer: explicitly read apps/worker/tsconfig.json if it exists
if [[ -f ./apps/worker/tsconfig.json ]]; then
  sed -n '1,220p' ./apps/worker/tsconfig.json | cat -n | rg -n 'compilerOptions|module|moduleResolution|verbatimModuleSyntax|preserveValueImports|importsNotUsedAsValues' || true
  echo "---- full compilerOptions excerpt (around match) ----"
  rg -n 'compilerOptions' ./apps/worker/tsconfig.json | head -n 5 || true
  rg -n '("module"|"moduleResolution"|"verbatimModuleSyntax"|"preserveValueImports"|"importsNotUsedAsValues")' ./apps/worker/tsconfig.json || true
else
  echo "NOT_FOUND: ./apps/worker/tsconfig.json"
fi

echo
echo "3) Find non-type OutboxTable imports from base-outbox.poller.js:"
rg -n --hidden -S 'import\s*\{\s*BaseOutboxPoller\s*,\s*OutboxTable\s*\}\s*from\s*"\.\/base-outbox\.poller\.js"' apps packages --glob '*.ts' || true
rg -n --hidden -S 'import\s*\{\s*BaseOutboxPoller\s*,\s*OutboxTable\s*\}\s*from\s*"\.\/base-outbox\.poller\.js"' . --glob '*.ts' || true

echo
echo "4) Inspect exact import/use of OutboxTable in normalized-outbox.poller.ts:"
sed -n '1,40p' ./apps/worker/src/pollers/normalized-outbox.poller.ts | cat -n
rg -n '\bOutboxTable\b' ./apps/worker/src/pollers/normalized-outbox.poller.ts

Repository: pramodnarayana/nexiom

Length of output: 2989


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "1) Search for relevant compiler options keys across repo (any json):"
rg -n --hidden -S '"verbatimModuleSyntax"' . --glob '*.json' || true
rg -n --hidden -S '"preserveValueImports"' . --glob '*.json' || true
rg -n --hidden -S '"importsNotUsedAsValues"' . --glob '*.json' || true

echo
echo "2) Find all imports of OutboxTable from base-outbox.poller.js (non-type import pattern):"
rg -n --hidden -S 'from\s*"\.\/base-outbox\.poller\.js"\s*;\s*$' apps packages --glob '*.ts' || true
rg -n --hidden -S 'import\s*\{\s*[^}]*\bOutboxTable\b[^}]*\}\s*from\s*"\.\/base-outbox\.poller\.js"' . --glob '*.ts' || true

echo
echo "3) Show OutboxTable occurrences with imports (small context):"
rg -n --hidden -S '\bOutboxTable\b' apps/worker/src/pollers/*.ts --glob '*.ts' || true

Repository: pramodnarayana/nexiom

Length of output: 1860


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "apps/worker/tsconfig.json (relevant compilerOptions excerpt):"
sed -n '1,120p' ./apps/worker/tsconfig.json | cat -n | rg -n 'compilerOptions|module|moduleResolution|verbatimModuleSyntax|preserveValueImports|importsNotUsedAsValues|importsNotUsedAsValues|typeRoots|types|verbatimModuleSyntax' || true
echo
echo "Full compilerOptions block (for worker tsconfig):"
rg -n '"compilerOptions"\s*:' ./apps/worker/tsconfig.json -n
# print only around compilerOptions to keep output small
perl -0777 -ne 'if(/"compilerOptions"\s*:\s*\{(.*?)\}\s*(,|\n\s*\})/s){print $&}' ./apps/worker/tsconfig.json | head -n 120

echo
echo "Look for tsconfig.json (root) that might be extended by apps/worker/tsconfig.json:"
ls -1 . | rg -n 'tsconfig.*\.json$' || true
for f in $(ls -1 . | rg 'tsconfig.*\.json$' || true); do
  echo "---- $f ----"
  rg -n '"verbatimModuleSyntax"|"preserveValueImports"|"importsNotUsedAsValues"|"extends"' "$f" || true
done

Repository: pramodnarayana/nexiom

Length of output: 1180


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "tsconfig.base.json (relevant flags):"
sed -n '1,220p' ./tsconfig.base.json | cat -n | rg -n '"verbatimModuleSyntax"|"preserveValueImports"|"importsNotUsedAsValues"|"module"|"moduleResolution"|"target"|"moduleDetection"|"\bjsx\b"|"\bisolatedModules\b"|"\bemitDecoratorMetadata\b"' || true

echo
echo "---- full tsconfig.base.json ----"
sed -n '1,260p' ./tsconfig.base.json | cat -n

Repository: pramodnarayana/nexiom

Length of output: 716


Use a type-only import for OutboxTable

apps/worker/src/pollers/normalized-outbox.poller.ts imports OutboxTable as a value even though it’s only referenced in a type assertion (normalizedOutbox as unknown as OutboxTable). The worker config doesn’t enable verbatimModuleSyntax/preserveValueImports/importsNotUsedAsValues, so runtime mismatch risk is low, but using type keeps the intent clear and avoids future config regressions (same pattern exists in other outbox pollers).

Proposed fix
-import { BaseOutboxPoller, OutboxTable } from "./base-outbox.poller.js";
+import { BaseOutboxPoller, type OutboxTable } from "./base-outbox.poller.js";
📝 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.

import { BaseOutboxPoller, type OutboxTable } from "./base-outbox.poller.js";
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@apps/worker/src/pollers/normalized-outbox.poller.ts` at line 1, Change the
import of OutboxTable to a type-only import so it’s not treated as a runtime
value: update the import statement that currently brings in BaseOutboxPoller and
OutboxTable to import OutboxTable using the `type` keyword, and ensure the type
assertion `normalizedOutbox as unknown as OutboxTable` continues to compile
without bringing OutboxTable into the emitted JS; locate the import in this file
(referencing BaseOutboxPoller and OutboxTable) and make the OutboxTable import
type-only to match its usage.
packages/database/drizzle/global/0018_gorgeous_alice.sql (1)

1-1: ⚠️ Potential issue | 🟠 Major | 🏗️ Heavy lift

Protect this migration from existing duplicate rows before adding the unique constraint.

Line 1 applies the unique constraint immediately; if duplicates already exist, deployment will fail at migration time. Add a deterministic dedupe/backfill (or explicit precheck/fail-fast query) before ADD CONSTRAINT.

Suggested migration shape
+-- Example: keep one row per (workspace_id, piece_id) and remove older duplicates
+WITH ranked AS (
+  SELECT ctid, workspace_id, piece_id,
+         ROW_NUMBER() OVER (PARTITION BY workspace_id, piece_id ORDER BY ctid) AS rn
+  FROM workspace_pieces
+)
+DELETE FROM workspace_pieces wp
+USING ranked r
+WHERE wp.ctid = r.ctid
+  AND r.rn > 1;
+
 ALTER TABLE "workspace_pieces" ADD CONSTRAINT "ux_workspace_pieces_workspaceId_pieceId" UNIQUE("workspace_id","piece_id");
📝 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.

-- Example: keep one row per (workspace_id, piece_id) and remove older duplicates
WITH ranked AS (
  SELECT ctid, workspace_id, piece_id,
         ROW_NUMBER() OVER (PARTITION BY workspace_id, piece_id ORDER BY ctid) AS rn
  FROM workspace_pieces
)
DELETE FROM workspace_pieces wp
USING ranked r
WHERE wp.ctid = r.ctid
  AND r.rn > 1;

ALTER TABLE "workspace_pieces" ADD CONSTRAINT "ux_workspace_pieces_workspaceId_pieceId" UNIQUE("workspace_id","piece_id");
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@packages/database/drizzle/global/0018_gorgeous_alice.sql` at line 1, Before
adding the unique constraint on workspace_pieces
(ux_workspace_pieces_workspaceId_pieceId over columns workspace_id and
piece_id), add a deterministic pre-migration step that either: 1) detects
duplicates with a GROUP BY workspace_id, piece_id and fails fast with an
explicit error listing (or count) so the deploy can be fixed manually, or 2)
performs a deterministic dedupe/backfill that deletes or consolidates duplicate
rows (keeping the row with the lowest id or latest updated_at) and records the
removed ids in an audit table; implement this logic in the same migration prior
to the ALTER TABLE so that the ADD CONSTRAINT will not fail due to existing
duplicates. Ensure you reference workspace_pieces, workspace_id, piece_id and
ux_workspace_pieces_workspaceId_pieceId when adding the precheck/dedupe.
packages/database/drizzle/tenant/0012_military_black_queen.sql (1)

1-2: ⚠️ Potential issue | 🔴 Critical

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
# Verify whether prior tenant migrations created gem_unique_mapping_idx as an index or as a constraint.
# Expected: exactly one historical definition style across branches/environments.
rg -n -C3 'gem_unique_mapping_idx|global_entity_map' packages/database/drizzle/tenant

Repository: pramodnarayana/nexiom

Length of output: 50377


Fix mixed-schema handling for gem_unique_mapping_idx migration (avoid failing when it already exists as a UNIQUE constraint).

Meta snapshots show gem_unique_mapping_idx is represented as an index in earlier states (e.g., packages/database/drizzle/tenant/meta/0010_snapshot.json) but as a UNIQUE constraint in packages/database/drizzle/tenant/meta/0012_snapshot.json. With your current order (DROP INDEX IF EXISTS then ALTER TABLE ... ADD CONSTRAINT), Line 1 can be a no-op when the constraint already exists, causing Line 2 to fail due to the existing constraint. Drop the constraint (IF EXISTS) before adding it (your suggested change directionally covers this).

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

In `@packages/database/drizzle/tenant/0012_military_black_queen.sql` around lines
1 - 2, The migration should remove any existing representation of
gem_unique_mapping_idx before adding it as a UNIQUE constraint on table
global_entity_map; change the script to first DROP CONSTRAINT IF EXISTS
"gem_unique_mapping_idx" (on global_entity_map) and also DROP INDEX IF EXISTS
"gem_unique_mapping_idx" if needed, then run ALTER TABLE "global_entity_map" ADD
CONSTRAINT "gem_unique_mapping_idx"
UNIQUE("stitch_id","source_data_source_id","source_entity_id","dest_data_source_id","dest_entity_type")
so the migration succeeds whether the object exists as an index or as a
constraint.
packages/observability/package.json (2)

12-23: ⚠️ Potential issue | 🟠 Major | ⚡ Quick win

Fix lockfile drift for this new workspace manifest.

CI already fails with frozen-lockfile because pnpm-lock.yaml specifiers don’t match this package manifest. Regenerate and commit the lockfile in this PR.

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

In `@packages/observability/package.json` around lines 12 - 23, The package
manifest's "dependencies" and "devDependencies" blocks were changed causing
lockfile drift; regenerate the workspace lockfile and commit it so CI's
frozen-lockfile check passes. From the repository root run the pnpm workspace
install command (e.g., pnpm install or pnpm -w install) to update pnpm-lock.yaml
to match the new packages/observability/package.json specifiers, verify the
lockfile changes, and commit the updated pnpm-lock.yaml alongside the
package.json change.

Source: Pipeline failures


13-15: ⚠️ Potential issue | 🟠 Major

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
python - <<'PY'
import json, pathlib
for p in sorted(pathlib.Path("packages").rglob("package.json")):
    d = json.loads(p.read_text())
    deps = d.get("dependencies", {})
    peers = d.get("peerDependencies", {})
    if any(k in deps or k in peers for k in ("`@nestjs/common`","`@nestjs/config`","nestjs-pino")):
        print(f"{p}")
        print("  name:", d.get("name"))
        print("  dep `@nestjs/common`:", deps.get("`@nestjs/common`"))
        print("  peer `@nestjs/common`:", peers.get("`@nestjs/common`"))
        print("  dep `@nestjs/config`:", deps.get("`@nestjs/config`"))
        print("  dep nestjs-pino:", deps.get("nestjs-pino"))
PY

Repository: pramodnarayana/nexiom

Length of output: 1867


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "=== packages/observability/package.json ==="
if [ -f packages/observability/package.json ]; then
  cat -n packages/observability/package.json
else
  echo "missing packages/observability/package.json"
fi

echo
echo "=== Root package.json (workspaces context) ==="
if [ -f package.json ]; then
  cat -n package.json
else
  echo "missing root package.json"
fi

echo
echo "=== Collect `@nestjs/common/`@nestjs/config/nestjs-pino declarations across workspace ==="
python3 - <<'PY'
import json, pathlib
keys = ("`@nestjs/common`","`@nestjs/config`","nestjs-pino")
for p in sorted(pathlib.Path("packages").rglob("package.json")):
    d = json.loads(p.read_text())
    name = d.get("name", str(p))
    deps = d.get("dependencies", {}) or {}
    peers = d.get("peerDependencies", {}) or {}
    hits = []
    for k in keys:
        if k in deps: hits.append((k,"dependencies",deps.get(k)))
        if k in peers: hits.append((k,"peerDependencies",peers.get(k)))
    if hits:
        print(f"{p} :: {name}")
        for k, where, v in hits:
            print(f"  {k} ({where}): {v}")
PY

Repository: pramodnarayana/nexiom

Length of output: 5476


Align NestJS major versions across workspace packages (observability is pinned to older NestJS).

packages/observability/package.json pins @nestjs/common to ^10.4.1 (deps) while other workspace packages use ^11.1.11 (mostly as peer deps); it also pins @nestjs/config to ^3.2.3 while others use ^4.0.2—this can lead to multiple NestJS major versions being installed and runtime/type incompatibilities across packages.

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

In `@packages/observability/package.json` around lines 13 - 15, observability
package pins older NestJS deps causing mixed major versions; update packages in
packages/observability/package.json to match workspace majors by bumping
"`@nestjs/common`" to the workspace version (e.g. ^11.1.11) and "`@nestjs/config`"
to the matching major (e.g. ^4.0.2), and ensure peerDependencies (if present)
for Nest packages align with those versions; verify dev/build tooling (e.g.
nestjs-pino) remains compatible with the newer Nest major and run
install/typecheck to catch any remaining incompatibilities.
packages/pipeline/fix-imports.cjs (1)

14-30: ⚠️ Potential issue | 🟠 Major

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
set -euo pipefail

rg -nP "from\\s+'(\\./(ports|adapters|interfaces)/|\\./outbox\\.utils\\.js|\\.\\./(ports/|utils\\.js|index\\.js|storage-resolver/|sharding/))" packages/pipeline/src || true

Repository: pramodnarayana/nexiom

Length of output: 663


Regex rewrites are quote-style fragile and can skip valid imports.

packages/pipeline/fix-imports.cjs replacement patterns only match from "...", leaving single-quoted imports untouched (e.g. packages/pipeline/src/replication/replica.service.spec.ts and packages/pipeline/src/delivery/delivery.service.spec.ts import from '../storage-resolver/... and '../sharding/...), resulting in a partially migrated state.

Proposed hardening
-  { from: /from "\.\/ports\//g, to: 'from "../shared/ports/' },
+  { from: /from\s+(["'])\.\/ports\//g, to: 'from $1../shared/ports/' },

-  { from: /from "\.\/adapters\//g, to: 'from "../shared/adapters/' },
+  { from: /from\s+(["'])\.\/adapters\//g, to: 'from $1../shared/adapters/' },

-  { from: /from "\.\/interfaces\//g, to: 'from "../shared/interfaces/' },
+  { from: /from\s+(["'])\.\/interfaces\//g, to: 'from $1../shared/interfaces/' },

Apply the same quote-agnostic approach to the other from "\.\.\/..." replacement patterns as well (e.g. ../storage-resolver/, ../sharding/, etc.).

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

In `@packages/pipeline/fix-imports.cjs` around lines 14 - 30, The replacement
regexes in packages/pipeline/fix-imports.cjs currently only match double-quoted
imports (e.g. patterns like /from "\.\.\/utils\.js"/g) and miss single-quoted
imports; update each affected pattern (e.g. the ones referencing ports,
adapters, interfaces, outbox.utils.js, ../ports/, ../adapters/, ../interfaces/,
../outbox.utils.js, /from "\.\.\/utils\.js"/, /from "\.\.\/index\.js"/, /from
"\.\.\/storage-resolver\//, /from "\.\.\/sharding\//) to be quote-agnostic by
matching either single or double quotes (use a regex like /from
['"]\.\.\/...['"]/ style) so all import variants are rewritten consistently
while leaving the gem-hydration.service.js rule unchanged.
packages/pipeline/src/delivery/delivery.service.integration.spec.ts (2)

46-48: ⚠️ Potential issue | 🟡 Minor | ⚡ Quick win

Remove duplicate applyPlan call.

Line 48 is an exact duplicate of line 47. Applying the same plan (SchemaPlan.OUTBOUND_ACTIVE) twice in succession is likely a copy-paste error and unnecessary. This wastes test setup time and may cause unintended side effects if the plan isn't idempotent.

🔧 Proposed fix
     await sqlManager.applyPlan(currentSchemaName, SchemaPlan.NAMESPACE_ONLY, { appName: "testApp", appProfile: "online" });
     await sqlManager.applyPlan(currentSchemaName, SchemaPlan.OUTBOUND_ACTIVE, { appName: "testApp", appProfile: "online" });
-    await sqlManager.applyPlan(currentSchemaName, SchemaPlan.OUTBOUND_ACTIVE, { appName: "testApp", appProfile: "online" });
📝 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.

    await sqlManager.applyPlan(currentSchemaName, SchemaPlan.NAMESPACE_ONLY, { appName: "testApp", appProfile: "online" });
    await sqlManager.applyPlan(currentSchemaName, SchemaPlan.OUTBOUND_ACTIVE, { appName: "testApp", appProfile: "online" });
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@packages/pipeline/src/delivery/delivery.service.integration.spec.ts` around
lines 46 - 48, Remove the duplicate call to sqlManager.applyPlan for
SchemaPlan.OUTBOUND_ACTIVE: locate the consecutive calls to
sqlManager.applyPlan(currentSchemaName, SchemaPlan.OUTBOUND_ACTIVE, { appName:
"testApp", appProfile: "online" }) and delete the redundant second invocation so
the plan is applied only once during test setup; ensure only the NAMESPACE_ONLY
and a single OUTBOUND_ACTIVE call remain.

239-296: ⚠️ Potential issue | 🔴 Critical | ⚡ Quick win

Remove obsolete tenantDb parameter from writeL6Result call.

According to the AI summary, the writeL6Result method signature was updated to remove the tenantDb: DrizzleDb parameter. However, line 293 still passes testDbManager.db! as the final argument. This parameter count mismatch will cause the test to fail or pass incorrect values to other parameters.

🐛 Proposed fix
       const res = await (service as any).writeL6Result(
         currentSchemaName,
         currentSchemaName,
         uuidv4(),
         1,
         "conn",
         uuidv4(),
         uuidv4(),
         null,
         null,
         200,
         "SUCCESS",
         Date.now(),
         "destVendorId",
         "RAW",
         "app",
         currentWorkspaceId,
         "srcVendorId",
         "tgt",
         "targetApp",
         "targetTenant",
         "targetObject",
         currentWorkspaceId,
-        testDbManager.db!,
       );
📝 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.

    it("should write GEM mapping if all conditions are met", async () => {
      vi.spyOn(gemService, "writeGemMapping").mockResolvedValue(undefined);

      const res = await (service as any).writeL6Result(
        currentSchemaName,
        currentSchemaName,
        uuidv4(),
        1,
        "conn",
        uuidv4(),
        uuidv4(),
        null,
        null,
        200,
        "SUCCESS",
        Date.now(),
        "destVendorId",
        "RAW",
        "app",
        currentWorkspaceId,
        "srcVendorId",
        "tgt",
        "targetApp",
        "targetTenant",
        "targetObject",
        currentWorkspaceId,
      );
      expect(res).toBe(true);
      expect(gemService.writeGemMapping).toHaveBeenCalled();
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@packages/pipeline/src/delivery/delivery.service.integration.spec.ts` around
lines 239 - 296, The test is passing an obsolete tenantDb argument to
writeL6Result, which no longer accepts a DrizzleDb parameter; remove the final
testDbManager.db! argument from the writeL6Result call in this spec so the
argument list matches the updated writeL6Result signature (refer to
writeL6Result and the test that asserts gemService.writeGemMapping to locate the
call), and run the spec to ensure no other tests still pass tenantDb.
packages/pipeline/src/delivery/use-cases/claim-delivery.use-case.spec.ts (1)

136-169: 🛠️ Refactor suggestion | 🟠 Major | ⚡ Quick win

Add explicit isSourceFinalized argument-order assertions.

These tests assert outcomes but not the call contract, so a tenant/schema/trace/route ordering regression can slip through. Add toHaveBeenCalledWith(tenantId, srcSchemaName, traceId, routeId) assertions in the SUCCESS/FAIL finalized paths.

Suggested test assertion shape
   expect(result).toEqual({ status: 'TERMINATED' });
+  expect(retryService.isSourceFinalized).toHaveBeenCalledWith(
+    'ten-1',
+    'ws_src-1',
+    'tr-1',
+    'rt-1',
+  );
   expect(retryService.retrySourceFinalization).not.toHaveBeenCalled();

Also applies to: 100-134

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

In `@packages/pipeline/src/delivery/use-cases/claim-delivery.use-case.spec.ts`
around lines 136 - 169, Add explicit assertions that
retryService.isSourceFinalized was called with the correct argument ordering to
guard against regression: after calling useCase.execute in the two
finalized-path tests (the ones setting outboundGatewayRepository entries to
status 'SUCCESS' and 'FAIL'), add
expect(retryService.isSourceFinalized).toHaveBeenCalledWith(tenantId,
srcSchemaName, traceId, routeId) using the variables from the test (or values
from defaultInput) to assert tenant/schema/trace/route ordering; ensure this
same assertion is added for the other similar tests around lines 100-134 as
well.
packages/pipeline/src/delivery/use-cases/claim-delivery.use-case.ts (1)

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

Fix isSourceFinalized parameter order (currently miswired).

Line 180 and Line 242 pass arguments in the wrong order. This breaks the finalization check contract and can incorrectly force retry/failure paths.

Proposed fix
-      sourceFinalized = await this.retryService.isSourceFinalized(
-        srcSchemaName,
-        traceId,
-        routeId,
-        tenantId,
-      );
+      sourceFinalized = await this.retryService.isSourceFinalized(
+        tenantId,
+        srcSchemaName,
+        traceId,
+        routeId,
+      );
...
-      sourceFinalized = await this.retryService.isSourceFinalized(
-        srcSchemaName,
-        traceId,
-        routeId,
-        tenantId,
-      );
+      sourceFinalized = await this.retryService.isSourceFinalized(
+        tenantId,
+        srcSchemaName,
+        traceId,
+        routeId,
+      );

Also applies to: 242-247

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

In `@packages/pipeline/src/delivery/use-cases/claim-delivery.use-case.ts` around
lines 179 - 185, The calls to retryService.isSourceFinalized are passing
arguments in the wrong order (currently calling with srcSchemaName, traceId,
routeId, tenantId) which breaks the finalization check; update both calls that
set sourceFinalized (the one guarded by currentStatus === "SUCCESS" and the
later call around the retry path) to pass arguments in the same order as the
retryService.isSourceFinalized signature—reorder srcSchemaName, traceId,
routeId, tenantId to match that signature (e.g., if the method expects traceId,
routeId, srcSchemaName, tenantId, swap the arguments accordingly) so the
parameters traceId, routeId, srcSchemaName, tenantId align with the method
definition. Ensure you change both usages of retryService.isSourceFinalized,
leaving variable names (currentStatus, sourceFinalized, srcSchemaName, traceId,
routeId, tenantId) intact.
packages/pipeline/src/fanout/fanout-batch-processor.ts (3)

176-204: 🧹 Nitpick | 🔵 Trivial | ⚡ Quick win

Simplify SyncToken extraction logic.

The nested logic for extracting debugSyncToken from state is complex and fragile. It first checks state.SyncToken, then falls back to searching for the first object-typed key that isn't "time", then accesses that nested object's SyncToken. This logic appears purely for debug logging and doesn't affect control flow.

Consider extracting this to a helper function to improve readability and maintainability.

♻️ Proposed refactor
+  private extractDebugSyncToken(state: Record<string, unknown>): string {
+    if ((state as { SyncToken?: string })?.SyncToken) {
+      return (state as { SyncToken?: string }).SyncToken!;
+    }
+    const entityKey = Object.keys(state).find(
+      (k) => k !== "time" && typeof state[k] === "object" && state[k] !== null
+    );
+    if (entityKey) {
+      return (state[entityKey] as { SyncToken?: string })?.SyncToken ?? "not_found";
+    }
+    return "not_found";
+  }

   if (state) {
     destState = state;
-    let debugSyncToken = (state as { SyncToken?: string })?.SyncToken;
-    if (!debugSyncToken) {
-      const entityKey = Object.keys(state).find(
-        (k) =>
-          k !== "time" &&
-          typeof state[k] === "object" &&
-          state[k] !== null,
-      );
-      if (entityKey)
-        debugSyncToken = (state[entityKey] as { SyncToken?: string })
-          ?.SyncToken;
-    }
+    const debugSyncToken = this.extractDebugSyncToken(state);
     this.logger.log(
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@packages/pipeline/src/fanout/fanout-batch-processor.ts` around lines 176 -
204, The SyncToken extraction in the destState handling is overly nested and
should be moved into a small helper to simplify the block: implement a function
(e.g., extractSyncTokenFromState(state): string | undefined) that first returns
state.SyncToken if present, otherwise finds the first key not equal to "time"
whose value is a non-null object and returns that object's SyncToken; replace
the inline logic that sets debugSyncToken in the block that assigns destState
and call this helper, then use the returned value in the this.logger.log call
(symbols to change: debugSyncToken, destState, and the logging call inside the
if (state) branch). Ensure behavior remains identical and that this helper is
only used for debug logging.

282-289: ⚠️ Potential issue | 🟠 Major

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
# Description: Check how failed outbound gateways are retried

# Search for retry logic or sweeper services that process failed gateways
rg -nP --type=ts -C5 'markOutboundGatewayFailed|FAILED.*gateway|gateway.*retry' packages/pipeline/

# Check if there's a background job that processes failed gateways
ast-grep --pattern 'class $NAME {
  $$$
  processFailed$_($$$) {
    $$$
  }
  $$$
}'

Repository: pramodnarayana/nexiom

Length of output: 8086


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# Show the relevant control-flow around the suggested throw/outer-catch region
sed -n '230,460p' packages/pipeline/src/fanout/fanout-batch-processor.ts | nl -ba | sed -n '1,240p'

# Inspect the repository method that marks gateways as failed
sed -n '1,240p' packages/pipeline/src/shared/adapters/drizzle-outbound-gateway.repository.ts | nl -ba | sed -n '1,240p'

# Search for code paths that select/claim failed/failed-gateway rows for retry
rg -n --type=ts "outbound_gateway|markOutboundGatewayFailed|FAILED|DEFERRED_DEPENDENCY|PENDING|attempts|delivery-retry|retry.*gateway|gateway.*retry|process.*outbound gateway|sweep.*outbound" packages/pipeline/src | head -n 200

# Search for any job/scheduler/worker entrypoints that might run retries
rg -n --type=ts "(schedule|cron|setInterval|Bull|queue|scheduler|worker|processor)" packages/pipeline/src | head -n 200

Repository: pramodnarayana/nexiom

Length of output: 108


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# 1) Fanout batch processor: inspect around the failure mark + surrounding try/catch
cat -n packages/pipeline/src/fanout/fanout-batch-processor.ts | sed -n '240,460p'

# 2) Outbound gateway repo: inspect markOutboundGatewayFailed implementation and any attempt/status handling
cat -n packages/pipeline/src/shared/adapters/drizzle-outbound-gateway.repository.ts | sed -n '60,240p'

# 3) Find any code that processes/claims outbound gateways with FAILED/attempts or that runs delivery retries
rg -n --type=ts "markOutboundGatewayFailed|DEFERRED_DEPENDENCY|FAILED|status.*FAILED|attempts|delivery-retry|retry.*gateway|gateway.*retry|process.*outbound gateway|sweep.*outbound|outbound_gateway" packages/pipeline/src | head -n 200

# 4) Look specifically for any background job/worker/scheduler related to delivery retry
rg -n --type=ts "(cron|schedule|setInterval|Bull|queue|scheduler|worker|processor)" packages/pipeline/src | head -n 200

Repository: pramodnarayana/nexiom

Length of output: 50377


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# Fanout: confirm outer catch behavior is only for logging/continue (not requeueing)
sed -n '240,430p' packages/pipeline/src/fanout/fanout-batch-processor.ts | cat

# Delivery: find how it claims work (what statuses are eligible)
sed -n '1,260p' packages/pipeline/src/delivery/use-cases/claim-delivery.use-case.ts | cat

# Delivery: inspect main consumer loop
sed -n '1,260p' packages/pipeline/src/delivery/delivery.service.ts | cat

# Retry: inspect DeliveryRetryService implementation for what it retries and how it transitions status
sed -n '1,240p' packages/pipeline/src/delivery/delivery-retry.service.ts | cat

# Outbound gateway adapter: inspect claim/transition logic relevant to retries (PENDING/RETRY/FAILED)
sed -n '1,220p' packages/pipeline/src/shared/adapters/drizzle-outbound-gateway.repository.ts | cat

Repository: pramodnarayana/nexiom

Length of output: 33301


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# Find where outbound_gateway rows are claimed/queued for (re)processing
rg -n --type=ts "claimForProcessing|QueueName\.DeliveryQueue|outbound_gateway.*status|status.*='RETRY'|status.*='FAILED'|outbound_gateway.*status.*'RETRY'|outbound_gateway.*status.*'FAILED'" packages/pipeline/src | head -n 200

# Find any code that transitions FAILED -> PENDING/RETRY or requeues FAILED outbound gateways
rg -n --type=ts "FAILED.*(RETRY|PENDING)|RETRY.*FAILED|markOutboundGatewayFailed|status\s*=\s*'RETRY'|status\s*=\s*'PENDING'|outbound_gateway\s+SET\s+status" packages/pipeline/src | head -n 200

# Locate the DeliveryQueue worker entrypoint to see what it does with gateway.status
rg -n --type=ts "consume\(QueueName\.DeliveryQueue|processMessage\(rawMsg|claimForProcessing" packages/pipeline/src/delivery | head -n 200

# Inspect any repository method(s) that likely move FAILED rows back into schedulable states
rg -n --type=ts "outbound.*gateway.*(RETRY|PENDING)" packages/pipeline/src/shared/adapters | head -n 200

Repository: pramodnarayana/nexiom

Length of output: 5777


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# Inspect how FanoutRouter handles errors from FanoutBatchProcessor
sed -n '1,240p' packages/pipeline/src/fanout/fanout-router.service.ts | cat -n

# Inspect if FanoutBatchProcessor errors are caught/handled at router level (e.g., try/catch)
rg -n --type=ts "consume\\(QueueName\\.NormalizedQueue|FanoutBatchProcessor|try \\{|catch \\(" packages/pipeline/src/fanout/fanout-router.service.ts packages/pipeline/src/fanout/fanout-batch-processor.ts

# Inspect any sweeping/cron-like job that looks at outbound_gateway.status = 'FAILED'
rg -n --type=ts "Cron|schedule|setInterval|sweep|worker|processor" packages/pipeline/src/normalization packages/pipeline/src/fanout packages/pipeline/src/delivery packages/pipeline/src/replication | head -n 200

# Inspect any logic that transitions status = 'FAILED' to 'PENDING'/'RETRY'
rg -n --type=ts "status\\s*=\\s*'FAILED'|\\bFAILED\\b" packages/pipeline/src | head -n 200

Repository: pramodnarayana/nexiom

Length of output: 18663


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# Enqueue points for DeliveryQueue (should reveal any recovery worker)
rg -n --type=ts "send\\(QueueName\\.DeliveryQueue" packages/pipeline/src

# Queue consumers for DeliveryQueue
rg -n --type=ts "consume\\(QueueName\\.DeliveryQueue" packages/pipeline/src

# Any periodic sweeper/cron/worker touching outbound_gateway
rg -n --type=ts "(Cron\\(|`@Cron`|EVERY_|schedule|setInterval|worker|processor)" packages/pipeline/src | head -n 200

# Any code that transitions outbound_gateway.status FROM FAILED to something schedulable (PENDING/RETRY/PROCESSING)
rg -n --type=ts "status\\s*=\\s*'PENDING'|status\\s*=\\s*'RETRY'|status\\s*=\\s*'PROCESSING'|status\\s*=\\s*'FAILED'" packages/pipeline/src/shared packages/pipeline/src/fanout packages/pipeline/src/delivery packages/pipeline/src/normalization packages/pipeline/src | head -n 200

# Find any usage of outbound_gateway claims/processing that isn't driven by QueueName.DeliveryQueue messages
rg -n --type=ts "claimForProcessing\\(|fetchOutboundGatewayResult\\(|markOutboundGatewayFailed\\(" packages/pipeline/src | head -n 200

Repository: pramodnarayana/nexiom

Length of output: 6015


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# Show where the delivery layer branches on outbound_gateway status
sed -n '260,420p' packages/pipeline/src/delivery/use-cases/claim-delivery.use-case.ts | cat -n

# Find all occurrences of 'FAILED' in the pipeline code
rg -n "'FAILED'" packages/pipeline/src | head -n 200

# Find all occurrences of status checks for 'FAIL' vs 'FAILED' vs 'RETRY'
rg -n "status\\s*===\\s*\"(FAIL|FAILED|RETRY)\"|status\\s+IN \\('PENDING', 'RETRY'\\)|status = 'FAILED'|status = 'FAIL'" packages/pipeline/src/shared packages/pipeline/src/delivery packages/pipeline/src/fanout | head -n 200

Repository: pramodnarayana/nexiom

Length of output: 3923


Potential delivery attempt loss when MQ publish fails

In packages/pipeline/src/fanout/fanout-batch-processor.ts (around lines 282-289), a QueueName.DeliveryQueue publish error marks outbound_gateway.status = 'FAILED' and rethrows—but the outer catch (around lines 312-406) only logs FAIL and does not rethrow/retry.

Delivery claim only processes outbound_gateway rows in ('PENDING','RETRY') (claimForProcessing), so FAILED rows won’t be picked up by the delivery worker. The only transition back to a schedulable state is upsertPendingOutboundGateway converting FAILED → PENDING, which requires another L4 fanout run for the same route (no dedicated recovery/sweeper for FAILED outbound gateways was found).

await this.outboxRepo.markOutboundGatewayFailed(
  srcTenantId,
  destSchemaName,
  traceId,
  stitch.id
);

throw sendErr;

Add a recovery path for outbound_gateway.status = 'FAILED' (e.g., enqueue/requeue or transition to PENDING/RETRY), or propagate the error so the upstream queue message is retried.

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

In `@packages/pipeline/src/fanout/fanout-batch-processor.ts` around lines 282 -
289, The current code calls outboxRepo.markOutboundGatewayFailed(...) and then
throws sendErr, which ends up swallowed by the outer catch and leaves
outbound_gateway rows permanently FAILED; fix by either propagating the error to
the upstream queue (so the message is retried) or immediately re-scheduling the
gateway for delivery: after catching the QueueName.DeliveryQueue publish failure
in fanout-batch-processor, replace the terminal FAILED-only path with a recovery
action such as calling outboxRepo.upsertPendingOutboundGateway(...) to set
status to PENDING/RETRY or re-enqueue the delivery via this.mq.publish(...) with
the same payload, and only mark FAILED when you've exhausted retries; adjust
logic around markOutboundGatewayFailed, sendErr handling, and the outer catch so
failures are retriable.

340-358: ⚠️ Potential issue | 🔴 Critical

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
# Description: Verify DependencySweeperService processes deferred gateways

# Search for DependencySweeperService implementation
rg -nP --type=ts -C10 'class DependencySweeperService' packages/pipeline/

# Check for scheduled/cron jobs that process deferred gateways
rg -nP --type=ts 'deferred.*gateway|upsertDeferredOutboundGateway' packages/pipeline/ apps/worker/

Repository: pramodnarayana/nexiom

Length of output: 3017


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# Show relevant sections around the call site in fanout-batch-processor.ts
echo "===== fanout-batch-processor.ts (around upsertDeferredOutboundGateway + ActiveFetchQueue send) ====="
sed -n '280,420p' packages/pipeline/src/fanout/fanout-batch-processor.ts | nl -ba | sed -n '1,180p'

# Show core logic of DependencySweeperService
echo "===== dependency-sweeper.service.ts (sweepDeferredDependencies) ====="
sed -n '1,220p' packages/pipeline/src/normalization/dependency-sweeper.service.ts | nl -ba

# Show repository method upsertDeferredOutboundGateway
echo "===== drizzle-outbound-gateway.repository.ts (upsertDeferredOutboundGateway) ====="
sed -n '1,220p' packages/pipeline/src/shared/adapters/drizzle-outbound-gateway.repository.ts | nl -ba | sed -n '1,220p'

Repository: pramodnarayana/nexiom

Length of output: 227


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "===== fanout-batch-processor.ts (around upsertDeferredOutboundGateway + ActiveFetchQueue send) ====="
sed -n '300,420p' packages/pipeline/src/fanout/fanout-batch-processor.ts | cat -n

echo "===== dependency-sweeper.service.ts (sweepDeferredDependencies) ====="
sed -n '1,260p' packages/pipeline/src/normalization/dependency-sweeper.service.ts | cat -n

echo "===== drizzle-outbound-gateway.repository.ts (upsertDeferredOutboundGateway) ====="
sed -n '1,260p' packages/pipeline/src/shared/adapters/drizzle-outbound-gateway.repository.ts | cat -n

Repository: pramodnarayana/nexiom

Length of output: 23617


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "===== Find DependencySweeperRepositoryPort implementation (getDeferredTraces, claimDeferredTrace) ====="
rg -n --type=ts "getDeferredTraces\\(|claimDeferredTrace\\(" packages/pipeline/src

echo "===== Find all usages of QueueName.ActiveFetchQueue ====="
rg -n --type=ts "QueueName\\.ActiveFetchQueue" packages/pipeline/src

Repository: pramodnarayana/nexiom

Length of output: 1114


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "===== drizzle-dependency-sweeper.repository.ts (getDeferredTraces, claimDeferredTrace) ====="
sed -n '1,220p' packages/pipeline/src/shared/adapters/drizzle-dependency-sweeper.repository.ts | cat -n

echo "===== Find usages of QueueName.NormalizedQueue (pipeline + apps/worker) ====="
rg -n --type=ts "QueueName\\.NormalizedQueue|NormalizedQueue" packages/pipeline apps/worker

Repository: pramodnarayana/nexiom

Length of output: 6001


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "===== fanout-router.service.ts (NormalizedQueue consumer) ====="
sed -n '1,120p' packages/pipeline/src/fanout/fanout-router.service.ts | cat -n

echo "===== normalized-outbox.poller.ts (worker side) ====="
sed -n '1,240p' apps/worker/src/pollers/normalized-outbox.poller.ts | cat -n

Repository: pramodnarayana/nexiom

Length of output: 12169


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "===== ActiveFetchQueue usages across repo ====="
rg -n --type=ts "QueueName\\.ActiveFetchQueue" .

echo "===== outbound_gateway status transitions / usages ====="
rg -n --type=ts "'DEFERRED_DEPENDENCY'|'PENDING'|'PROCESSING'|'SUCCESS'|'FAILED'|DEFERRED_DEPENDENCY|markOutboundGatewayFailed|claimForProcessing" packages/pipeline apps/worker

echo "===== upsertPendingOutboundGateway and upsertDeferredOutboundGateway call sites ====="
rg -n --type=ts "upsertPendingOutboundGateway\\(|upsertDeferredOutboundGateway\\(" packages/pipeline apps/worker

Repository: pramodnarayana/nexiom

Length of output: 10763


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "===== fanout-batch-processor.ts (broader context around outbound gateway + DependenciesMissingError) ====="
sed -n '200,420p' packages/pipeline/src/fanout/fanout-batch-processor.ts | cat -n

echo "===== active-fetch.worker.ts (consumer behavior) ====="
sed -n '1,220p' apps/worker/src/consumers/active-fetch.worker.ts | cat -n

Repository: pramodnarayana/nexiom

Length of output: 12215


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "===== Search for any other writes to outbound_gateway status='DEFERRED_DEPENDENCY' ====="
rg -n --type=ts "'DEFERRED_DEPENDENCY'" packages/pipeline apps/worker

echo "===== Inspect fanout-batch-processor.ts around where it decides route already processed / outbound_gateway state ====="
# Use targeted range likely earlier than 200; adjust after quick peek
sed -n '1,260p' packages/pipeline/src/fanout/fanout-batch-processor.ts | cat -n

Repository: pramodnarayana/nexiom

Length of output: 12484


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "===== QueueName.DeliveryQueue consumers (apps/worker) ====="
rg -n --type=ts "QueueName\\.DeliveryQueue" apps/worker

echo "===== Find outbound gateway processing / claimForProcessing usages ====="
rg -n --type=ts "claimForProcessing\\(|markResult\\(|markOutboundGatewayFailed\\(" packages/pipeline apps/worker

Repository: pramodnarayana/nexiom

Length of output: 125


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "===== Search repo for 'DeliveryQueue' (any usage) ====="
rg -n --type=ts "DeliveryQueue" .

echo "===== Search apps/worker/src for queue.consume( ... ) patterns ====="
rg -n --type=ts "queueService\\.consume\\(" apps/worker/src

echo "===== List files under apps/worker/src/consumers ====="
fd . apps/worker/src/consumers -t f | head -n 50

Repository: pramodnarayana/nexiom

Length of output: 4161


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "===== Search for outbound_gateway table definition / primary key ====="
rg -n --type=ts "outbound_gateway" packages -S

echo "===== Inspect delivery.service.ts (DeliveryQueue consumer) ====="
sed -n '1,260p' packages/pipeline/src/delivery/delivery.service.ts | cat -n

Repository: pramodnarayana/nexiom

Length of output: 16396


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "===== Search for outbound_gateway usage in apps/worker ====="
rg -n --type=ts "outbound_gateway" apps/worker/src

echo "===== Inspect DB schema for outbound_gateway constraints (sql-database-manager.ts) ====="
sed -n '520,720p' packages/dbmanager/src/impl/sql-database-manager.ts | cat -n

echo "===== Inspect tenant pipeline schema for outbound_gateway (packages/database/src/schema/tenant/pipeline.ts) ====="
sed -n '150,380p' packages/database/src/schema/tenant/pipeline.ts | cat -n

echo "===== Inspect ClaimDeliveryUseCase to see which outbound_gateway statuses are claimed ====="
sed -n '1,240p' packages/pipeline/src/delivery/use-cases/claim-delivery.use-case.ts | cat -n

Repository: pramodnarayana/nexiom

Length of output: 126


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "===== Search for outbound_gateway usage across repo (non-schema) ====="
rg -n --type=ts "outbound_gateway" packages/pipeline apps/worker apps/api packages/queue packages/dbmanager packages/database -S

echo "===== Search for 'processOutbox' or 'outbox' patterns tied to outbound_gateway ====="
rg -n --type=ts "(outbox|poller).*outbound|outbound.*(outbox|poller)|drain.*outbox" packages/pipeline apps/worker apps/api -S

Repository: pramodnarayana/nexiom

Length of output: 12158


Deferred dependency retry after ActiveFetchQueue publish failure is broken (gateways can get stuck).

  • fanout-batch-processor.ts logs QueueName.ActiveFetchQueue send failures but swallows them and then writes SKIPPED.
  • The deferred record is written as outbound_gateway.status = 'DEFERRED_DEPENDENCY', but DependencySweeperService claims it by updating outbound_gateway to status = 'PENDING' (via claimDeferredTrace, keyed only by traceId).
  • FanoutBatchProcessor only publishes to QueueName.DeliveryQueue when upsertPendingOutboundGateway() can update rows whose prior status is in ('DEFERRED_DEPENDENCY','FAILED') (or NULL). If the row is already PENDING, it will not publish, and subsequent sweeps won’t pick it up because they only select DEFERRED_DEPENDENCY. So the deferred retry path likely never progresses.

Fix the state transition contract between claimDeferredTrace() and upsertPendingOutboundGateway() (and include routeId in the deferred claim predicate) so re-queued NormalizedQueue work can actually dispatch DeliveryQueue and continue retrying.

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

In `@packages/pipeline/src/fanout/fanout-batch-processor.ts` around lines 340 -
358, The deferred-retry flow is broken because claimDeferredTrace() only keys by
traceId (not routeId) and marks rows PENDING while
upsertPendingOutboundGateway() only treats prior statuses
('DEFERRED_DEPENDENCY','FAILED',NULL) as publishable; fix by (1) changing
DependencySweeperService.claimDeferredTrace(...) to include routeId in its
WHERE/claim predicate so claims are per-route, and (2) update
upsertPendingOutboundGateway(...) to accept rows whose prior status is 'PENDING'
(in addition to 'DEFERRED_DEPENDENCY','FAILED',NULL) so re-queued work that was
claimed can be published to QueueName.DeliveryQueue; refer to
fanout-batch-processor.ts symbols claimDeferredTrace,
upsertPendingOutboundGateway, QueueName.ActiveFetchQueue and DeliveryQueue when
making these changes.
packages/pipeline/src/fanout/fanout-router.service.ts (2)

87-117: ⚠️ Potential issue | 🟠 Major

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
# Description: Check transaction isolation level and concurrent processing safeguards

# Search for transaction manager implementation and isolation settings
rg -nP --type=ts -C5 'runInTenantTransaction|isolation.*level|SERIALIZABLE|REPEATABLE' packages/dbmanager/

# Check if there's entity-level locking or queuing to prevent concurrent processing
rg -nP --type=ts 'entity.*lock|processing.*lock|concurrent.*entity' packages/pipeline/ packages/queue/

Repository: pramodnarayana/nexiom

Length of output: 47


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# 1) Locate runInTenantTransaction implementation + any isolation configuration
rg -n "runInTenantTransaction" -S packages | head -n 50

# Find transaction manager / db manager files likely containing isolation level configuration
rg -n "class .*Tx|txManager|TransactionManager|isolationLevel|SERIALIZABLE|REPEATABLE|READ COMMITTED|read committed" -S packages | head -n 200

# 2) Inspect the likely implementation file(s)
# (Resolve file paths via the first search; then print relevant portions with line numbers.)
files="$(rg -l "runInTenantTransaction" -S packages || true)"
echo "$files" | head -n 20
for f in $(echo "$files" | head -n 20); do
  echo "---- $f (context around runInTenantTransaction) ----"
  rg -n -C 5 "runInTenantTransaction" "$f" || true
done

# 3) Search for concurrency safeguards around entity/replica processing
rg -n "lock|locked|advisory|pg_advisory|FOR UPDATE|SKIP LOCKED|SKIPLOCKED|queue|enqueue|dedup|idempotent|inflight|concurrent|singleflight|semaphore" -S packages/pipeline packages/queue packages/dbmanager packages | head -n 200

Repository: pramodnarayana/nexiom

Length of output: 35140


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# 1) Show the transaction manager adapter implementation to see isolation config.
echo "---- drizzle-transaction-manager.adapter.ts ----"
sed -n '1,220p' packages/pipeline/src/shared/adapters/drizzle-transaction-manager.adapter.ts

# Also check the transaction-manager port + tx context typing (for hints about isolation/locking).
echo "---- transaction-manager.port.js/ts ----"
ls -1 packages/pipeline/src/shared/ports
sed -n '1,200p' packages/pipeline/src/shared/ports/transaction-manager.port.js 2>/dev/null || true
sed -n '1,200p' packages/pipeline/src/shared/ports/transaction-manager.port.ts 2>/dev/null || true

# 2) Inspect fanout-router.service.ts around the supersession check + any locking/dedup it uses.
echo "---- fanout-router.service.ts ----"
sed -n '1,220p' packages/pipeline/src/fanout/fanout-router.service.ts

# 3) Find what "releaseSyncLock" is (and any related locking functions) in pipeline code.
echo "---- find releaseSyncLock usages/defs ----"
rg -n "releaseSyncLock" packages/pipeline/src | head -n 100
rg -n "sync lock|SyncLock|advisory|pg_advisory|FOR UPDATE|SKIP LOCKED|lock contention|isLockContention|inflight|dedup|idempotent" packages/pipeline/src | head -n 200

Repository: pramodnarayana/nexiom

Length of output: 1647


🏁 Script executed:

#!/bin/bash
set -uo pipefail

# 1) Locate transaction-manager port file (path may differ from earlier assumption)
echo "---- locate transaction-manager.port ----"
fd "transaction-manager.port." packages/pipeline/src -a -t f

# 2) Locate fanout-router dependencies (any sync lock / repository calls)
echo "---- fanout-router.service.ts (first 220 lines) ----"
sed -n '1,240p' packages/pipeline/src/fanout/fanout-router.service.ts

echo "---- fanout-router.service.ts (around releaseSyncLock call/usage) ----"
rg -n "releaseSyncLock" packages/pipeline/src/fanout/fanout-router.service.ts -n
rg -n "evaluateSuperseded|superseded|normalizedObj|getReplicaSourceVendorId" packages/pipeline/src/fanout/fanout-router.service.ts -n -C 2

# 3) Find releaseSyncLock definition/implementation
echo "---- releaseSyncLock defs/usages across pipeline ----"
rg -n "releaseSyncLock" packages/pipeline/src | head -n 200

Repository: pramodnarayana/nexiom

Length of output: 10716


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "---- drizzle-pipeline-state.repository.ts (lock methods) ----"
sed -n '1,220p' packages/pipeline/src/shared/adapters/drizzle-pipeline-state.repository.ts

echo "---- search for sync lock acquisition/claiming ----"
rg -n "acquireSyncLock|getSyncLock|syncLock|releaseSyncLock|pg_advisory|advisory|FOR UPDATE|SKIP LOCKED|lockContention|lock_contention" \
  packages/pipeline/src/shared/adapters packages/pipeline/src | head -n 200

echo "---- routing-decision.engine.ts (evaluateSuperseded) ----"
rg -n "evaluateSuperseded" packages/pipeline/src/fanout packages/pipeline/src | head -n 50
sed -n '1,260p' packages/pipeline/src/fanout/routing-decision.engine.ts

echo "---- pipeline-state.repository port (relevant methods) ----"
sed -n '1,200p' packages/pipeline/src/shared/ports/pipeline-state.repository.port.js 2>/dev/null || true
sed -n '1,200p' packages/pipeline/src/shared/ports/pipeline-state.repository.port.ts 2>/dev/null || true

Repository: pramodnarayana/nexiom

Length of output: 9612


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# 1) Locate dbmanager package + transaction config
echo "---- packages list (to find dbmanager) ----"
ls -1 packages

echo "---- locate getTenantDb and tenantDb.transaction configuration ----"
rg -n "getTenantDb\\(" packages | head -n 50
rg -n "tenantDb\\.transaction\\(" packages | head -n 50
rg -n "isolationLevel|SERIALIZABLE|REPEATABLE READ|READ COMMITTED|READ_UNCOMMITTED" packages | head -n 200

# 2) Inspect routing repository (superseding logic) for query semantics
echo "---- locate routing repository methods used by evaluateSuperseded ----"
rg -n "hasNormalizedRecord|hasSupersedingNormalizedRecord|getReplicaIdByTraceId" packages/pipeline/src -S
# Try to open the port+adapter implementation
rg -n "class .*Routing.*Repository|implements .*RoutingRepositoryPort" packages/pipeline/src -S | head -n 50
rg -n "hasNormalizedRecord\\(" packages/pipeline/src -S
rg -n "hasSupersedingNormalizedRecord\\(" packages/pipeline/src -S

# 3) Inspect activeSyncLocks usage (lock acquisition/claiming)
echo "---- activeSyncLocks usage across pipeline ----"
rg -n "activeSyncLocks" packages/pipeline/src | head -n 200

# Look for any "lock" or "lockedByTraceId" fields beyond release.
rg -n "lockedByTraceId|INSERT.*activeSyncLocks|upsert.*activeSyncLocks|SKIP LOCKED|FOR UPDATE" packages/pipeline/src | head -n 200

Repository: pramodnarayana/nexiom

Length of output: 7777


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# 1) Inspect dbmanager tenant database manager for transaction/isolation configuration
echo "---- packages/dbmanager/src/impl/tenant-database-manager.ts (relevant parts) ----"
sed -n '1,260p' packages/dbmanager/src/impl/tenant-database-manager.ts
echo "---- search within that file for transaction/isolation options ----"
rg -n "transaction\\(|isolation|SERIALIZABLE|REPEATABLE|READ COMMITTED|READ UNCOMMITTED|setTransaction|SET TRANSACTION" packages/dbmanager/src/impl/tenant-database-manager.ts

# 2) Inspect routing repository + adapter used by RoutingDecisionEngine
echo "---- locate RoutingRepositoryPort implementation ----"
rg -n "RoutingRepositoryPort|implements .*RoutingRepositoryPort" packages/pipeline/src | head -n 50

# Print routing decision repo adapter (if found)
routing_files="$(rg -l "hasNormalizedRecord|hasSupersedingNormalizedRecord|getReplicaIdByTraceId" packages/pipeline/src || true)"
echo "---- routing files containing superseding queries ----"
echo "$routing_files" | head -n 50

for f in $(echo "$routing_files" | head -n 20); do
  echo "==== $f ===="
  sed -n '1,240p' "$f"
done

# 3) Inspect activeSyncLocks schema usage/queries beyond release
echo "---- activeSyncLocks usage/queries ----"
rg -n "activeSyncLocks" packages/pipeline/src | head -n 200
rg -n "lockedByTraceId|dataSourceId.*entityId|lock.*activeSyncLocks|insert.*activeSyncLocks|upsert.*activeSyncLocks" packages/pipeline/src | head -n 200

# 4) Inspect where FanoutRouterService calls batchProcessor for lock-related flow.
echo "---- fanout-batch-processor.ts (maybe acquires locks) ----"
ls -1 packages/pipeline/src/fanout
sed -n '1,260p' packages/pipeline/src/fanout/fanout-batch-processor.ts

Repository: pramodnarayana/nexiom

Length of output: 12682


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "---- activeSyncLocks insert/upsert/claim (all pipeline adapters) ----"
rg -n "activeSyncLocks" packages/pipeline/src/shared/adapters packages/pipeline/src | head -n 200

echo "---- activeSyncLocks lockedByTraceId usage ----"
rg -n "lockedByTraceId" packages/pipeline/src | head -n 200

echo "---- any acquire/claim sync lock helpers (keywords) ----"
rg -n "acquire.*SyncLock|claim.*SyncLock|sync lock|activeSyncLocks.*(insert|upsert|update)|SKIP LOCKED|FOR UPDATE|pg_advisory" packages/pipeline/src | head -n 200

echo "---- routing repository port implementation ----"
rg -n "class .*Routing.*(Repository|Adapter)|implements .*RoutingRepositoryPort|hasNormalizedRecord|hasSupersedingNormalizedRecord|getReplicaIdByTraceId" packages/pipeline/src -S | head -n 200

echo "---- print matching routing repository implementation files (small) ----"
routing_files="$(rg -l "hasNormalizedRecord|hasSupersedingNormalizedRecord|getReplicaIdByTraceId" packages/pipeline/src -S || true)"
echo "$routing_files" | head -n 20
for f in $(echo "$routing_files" | head -n 20); do
  echo "==== $f ===="
  sed -n '1,260p' "$f"
done

echo "---- fanout-batch-processor (for any locking/dedup around stitches) ----"
sed -n '1,320p' packages/pipeline/src/fanout/fanout-batch-processor.ts

echo "---- search for sync lock usage in fanout/batchprocessor folders ----"
rg -n "SyncLock|releaseSyncLock|activeSyncLocks|lockedByTraceId" packages/pipeline/src/fanout | head -n 200

Repository: pramodnarayana/nexiom

Length of output: 2212


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# 1) Routing repo methods implementation (search across all packages)
echo "---- routingRepo method implementations across repo ----"
rg -n "hasNormalizedRecord\\(|hasSupersedingNormalizedRecord\\(|getReplicaIdByTraceId\\(" packages | head -n 200

echo "---- class/adapter implementing RoutingRepositoryPort ----"
rg -n "RoutingRepositoryPort" packages/pipeline/src packages | head -n 200

# 2) activeSyncLocks creation/claiming (insert/upsert/update)
echo "---- activeSyncLocks insert/upsert/update across repo ----"
rg -n "activeSyncLocks" packages | head -n 200

echo "---- lockedByTraceId writes (insert/update/upsert) across repo ----"
rg -n "lockedByTraceId" packages | head -n 200

echo "---- any \"sync lock\" acquisition/claiming keywords across repo ----"
rg -n "sync lock|SyncLock|activeSyncLocks|acquire.*Lock|claim.*Lock|lockedByTraceId|SKIP LOCKED|FOR UPDATE" packages | head -n 250

# 3) If schema defines activeSyncLocks columns, search in database schema builder code
echo "---- database schema references for activeSyncLocks ----"
rg -n "activeSyncLocks|ActiveSyncLocks|active_sync_locks" packages/database packages | head -n 200

Repository: pramodnarayana/nexiom

Length of output: 10731


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "---- drizzle-routing.repository.ts ----"
sed -n '1,220p' packages/pipeline/src/shared/adapters/drizzle-routing.repository.ts

echo "---- tenant pipeline schema: active_sync_locks definition ----"
# Print around where activeSyncLocks is defined
rg -n "activeSyncLocks =|active_sync_locks" packages/database/src/schema/tenant/pipeline.ts
sed -n '80,170p' packages/database/src/schema/tenant/pipeline.ts
sed -n '320,380p' packages/database/src/schema/tenant/pipeline.ts

echo "---- any other active_sync_locks usage (DDL/DML) ----"
rg -n "active_sync_locks|activeSyncLocks" packages | head -n 200
rg -n "INSERT.*active_sync_locks|upsert.*active_sync_locks|UPDATE.*active_sync_locks" packages | head -n 200
rg -n "locked_by_trace_id" packages | head -n 200

echo "---- look for explicit transaction isolation configuration ----"
rg -n "SET TRANSACTION ISOLATION LEVEL|isolationLevel:|transaction\\(\\{[^}]*isolation|SERIALIZABLE|REPEATABLE READ|READ COMMITTED|READ_UNCOMMITTED" packages | head -n 200

echo "---- fanout-batch-processor.ts (lock acquisition may live here) ----"
sed -n '1,260p' packages/pipeline/src/fanout/fanout-batch-processor.ts

Repository: pramodnarayana/nexiom

Length of output: 10420


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "---- inspect fanout-batch-processor.ts ----"
sed -n '1,320p' packages/pipeline/src/fanout/fanout-batch-processor.ts

echo "---- search for writes into active_sync_locks by column names ----"
rg -n "locked_by_trace_id|expires_at|active_sync_locks|sync_lock" packages/pipeline/src | head -n 200

echo "---- search for INSERT/UPDATE targeting activeSyncLocks table alias/object ----"
rg -n "activeSyncLocks|active_sync_locks" packages/pipeline/src -S

echo "---- search for locks with dataSourceId+entityId pattern ----"
rg -n "dataSourceId.*entityId|entityId.*dataSourceId" packages/pipeline/src/fanout packages/pipeline/src/normalization packages/pipeline/src/delivery | head -n 200

Repository: pramodnarayana/nexiom

Length of output: 11073


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# 1) Inspect outbox/outbound gateway repository used by fanout batch processor
echo "---- locate outbox repository methods ----"
rg -n "upsertPendingOutboundGateway|markOutboundGatewayFailed" packages/pipeline/src -S

echo "---- outbound gateway repository adapter (if any) ----"
rg -n "class .*Outbound|implements .*OutboundGatewayRepositoryPort|OutboundGatewayRepositoryPort" packages/pipeline/src -S
out_files="$(rg -l "upsertPendingOutboundGateway" packages/pipeline/src -S || true)"
echo "$out_files" | head -n 20
for f in $(echo "$out_files" | head -n 20); do
  echo "==== $f ===="
  sed -n '1,260p' "$f"
done

# 2) Inspect tenant pipeline schema for outbound_outbox uniqueness/locking semantics
echo "---- tenant pipeline schema: outbound_outbox/outbound_gateway constraints ----"
rg -n "outbound_outbox|outbound_gateway|UNIQUE|uniqueIndex" packages/database/src/schema/tenant/pipeline.ts
# Print around outbound_outbox definition
sed -n '260,420p' packages/database/src/schema/tenant/pipeline.ts

# 3) Inspect QueueService.consume behavior: concurrency and FIFO/message grouping
echo "---- queue.service.ts (consume implementation) ----"
sed -n '1,280p' packages/queue/src/queue.service.ts

echo "---- search queue.service for FIFO / MessageGroupId / ordering hints ----"
rg -n "FIFO|MessageGroupId|dedup|group|MessageDeduplicationId|maxConcurrency|concurrency|parallel" packages/queue/src/queue.service.ts packages/queue/src/queue.module.ts packages/queue/src/queue.options.ts

Repository: pramodnarayana/nexiom

Length of output: 50377


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# 1) Show QueueService.poll() / message processing concurrency controls (small extraction around inFlight)
echo "---- queue.service.ts: search inFlight usage ----"
rg -n "inFlight" packages/queue/src/queue.service.ts

echo "---- queue.service.ts: poll/processMessage relevant block ----"
rg -n "private async poll\\(|processMessage\\(|MaxNumberOfMessages|ReceiveMessageCommand|Promise\\.all|inFlight\\.add|while \\(|await Promise" packages/queue/src/queue.service.ts
# Print from poll() start to just after handler invocation
start=$(rg -n "async poll\\(" packages/queue/src/queue.service.ts | head -n 1 | cut -d: -f1)
if [ -n "${start:-}" ]; then
  end=$((start+260))
  sed -n "${start},${end}p" packages/queue/src/queue.service.ts
fi

# 2) Check queue constants/options for FIFO/grouping
echo "---- queue constants/options ----"
ls -1 packages/queue/src
sed -n '1,220p' packages/queue/src/constants.ts 2>/dev/null || true
sed -n '1,260p' packages/queue/src/queue.module.ts
sed -n '1,260p' packages/queue/src/queue.options.ts
sed -n '1,260p' packages/queue/src/queue.options.spec.ts

# 3) Search pipeline delivery stage for supersession/out-of-order protection
echo "---- pipeline: search superseded/supersession in delivery/replication ----"
rg -n "supersed|supersession|evaluateSuperseded" packages/pipeline/src | head -n 200

echo "---- pipeline: search for replica_entity overwrites / version logic ----"
rg -n "INSERT INTO replica_entity|ON CONFLICT.*replica_entity|version\\s*=|updated_at\\s*=.*replica_entity" packages/pipeline/src | head -n 200

# 4) Inspect delivery claim-delivery/use-case for any ordering/supersession checks
rg -n "claim-delivery|ClaimDelivery|supersed|evaluate|locked|out-of-order|traceId" packages/pipeline/src/delivery -S | head -n 200
sed -n '1,240p' packages/pipeline/src/delivery/use-cases/claim-delivery.use-case.ts

# 5) Confirm there is no active_sync_locks insert anywhere besides schema/release
echo "---- repo-wide search for locked_by_trace_id inserts/updates (targeted) ----"
rg -n "locked_by_trace_id" packages | head -n 200
echo "---- repo-wide search for active_sync_locks insert/upsert (targeted) ----"
rg -n "active_sync_locks" packages | head -n 200

Repository: pramodnarayana/nexiom

Length of output: 31605


Ensure fanout supersession check is race-safe across concurrent workers (isolation/locking).

  • runInTenantTransaction (packages/pipeline/src/shared/adapters/drizzle-transaction-manager.adapter.ts) does not set a transaction isolation level (only SET LOCAL search_path), so evaluateSuperseded + subsequent normalized/replica reads run under the DB/driver default.
  • FanoutRouterService (packages/pipeline/src/fanout/fanout-router.service.ts) calls QueueService.consume with no options; QueueService defaults maxConcurrent to 1 per consumer, but this doesn’t prevent the same entity from being processed concurrently by other process instances/pods.
  • The “sync lock” table (active_sync_locks) is defined with UNIQUE (dataSourceId, entityId), but the pipeline code only ever releases locks (releaseSyncLock*)—there’s no corresponding insert/claim of active_sync_locks anywhere in the pipeline.
  • Delivery ultimately upserts into replica_entity and overwrites trace_id/data on conflict (packages/pipeline/src/shared/adapters/drizzle-outbound-gateway.repository.ts), so an older trace that passes the supersession check can overwrite a newer trace.

Fix by adding DB-level protection (explicit isolation +/or row/advisory locking) and/or implementing the missing active_sync_locks acquisition/hold around supersession evaluation and subsequent processing for the same (dataSourceId, entityId).

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

In `@packages/pipeline/src/fanout/fanout-router.service.ts` around lines 87 - 117,
The supersession check is not race-safe: runInTenantTransaction (used in
FanoutRouterService) doesn't set isolation and code calls evaluateSuperseded
then reads getNormalizedData/getReplicaSourceVendorId without any DB locking or
lock acquisition, and the pipeline never inserts into active_sync_locks (only
releases). Fix by acquiring a DB-level lock for the (dataSourceId, entityId)
before evaluateSuperseded and holding it for the transaction—either (preferred)
insert/select-for-update on active_sync_locks (INSERT ... ON CONFLICT DO NOTHING
then SELECT FOR UPDATE) inside txManager.runInTenantTransaction or use
pg_advisory_xact_lock keyed by dataSourceId+entityId; ensure the lock
acquisition happens in FanoutRouterService before calling
routingDecisionEngine.evaluateSuperseded and
stateRepo.getNormalizedData/getReplicaSourceVendorId, and ensure
releaseSyncLock* remains consistent (or remove if using advisory locks); also
adjust runInTenantTransaction to allow explicit transaction isolation or ensure
lock acquisition occurs within the same transaction context so the supersession
check + subsequent reads and eventual upsert
(drizzle-outbound-gateway.repository) are atomic.

157-159: 🛠️ Refactor suggestion | 🟠 Major | ⚡ Quick win

Remove temporary compatibility code.

Lines 158-159 build tenantDb and syncLog objects that are passed to batchProcessor.processSingleStitch but marked as "kept for signature compatibility" in the batch processor (lines 58-59 of fanout-batch-processor.ts). These parameters are not used in the new implementation.

Since this is a refactoring PR to remove obsolete code, consider removing these unused parameters from both the caller and callee to complete the migration.

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

In `@packages/pipeline/src/fanout/fanout-router.service.ts` around lines 157 -
159, Remove the temporary compatibility plumbing: delete the tenantDb = await
this.dbManager.getTenantDb(tenantId) and const { syncLog } =
buildTenantSchema(schemaName) lines in fanout-router.service.ts and remove the
corresponding unused parameters from batchProcessor.processSingleStitch's
signature and implementation in the fanout-batch-processor (the function/method
named processSingleStitch). Update any other call sites of processSingleStitch
to stop passing tenantDb/syncLog and remove internal references to those
parameters in the BatchProcessor/FanoutBatchProcessor code so the new
implementation no longer accepts or expects those unused args.
packages/pipeline/src/fanout/target-builder.service.ts (2)

28-28: 🧹 Nitpick | 🔵 Trivial

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
# Description: Check PipelineHookBrokerService provider configuration

# Search for how PipelineHookBrokerService is provided in modules
rg -nP --type=ts -C5 'PipelineHookBrokerService.*provide|providers.*PipelineHookBrokerService' packages/pipeline/

# Check if PipelineHookBrokerService is exported/configured as a token
ast-grep --pattern 'class PipelineHookBrokerService {
  $$$
}'

Repository: pramodnarayana/nexiom

Length of output: 13421


Confirm @Inject(PipelineHookBrokerService) token is valid (and consider removing for clarity).

PipelineHookBrokerService is registered and exported as a concrete provider in ApplicationLoaderModule (providers and exports both include it), so injecting it via the class token with @Inject(PipelineHookBrokerService) is supported. The decorator is likely redundant when the constructor parameter is already typed as PipelineHookBrokerService, so it can be omitted for a more idiomatic NestJS style.

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

In `@packages/pipeline/src/fanout/target-builder.service.ts` at line 28, The
constructor parameter on TargetBuilderService uses
`@Inject`(PipelineHookBrokerService) on the hookBroker parameter; confirm that
PipelineHookBrokerService is indeed registered/exported in
ApplicationLoaderModule and then remove the redundant `@Inject` decorator to use
the typed injection idiom (i.e., keep the constructor param named hookBroker:
PipelineHookBrokerService and remove the `@Inject`(...) decorator) unless you have
a custom provider token requiring explicit injection.

118-122: ⚠️ Potential issue | 🟡 Minor | ⚡ Quick win

Empty object check may not catch all invalid payloads.

Line 118 checks Object.keys(hydrated).length === 0 to validate the payload. However, this won't catch payloads where all values are undefined or null, such as { field1: undefined, field2: null }, which may still be invalid for downstream systems.

Consider adding additional validation or document that downstream systems must handle null/undefined values.

🛡️ Enhanced validation
  if (Object.keys(hydrated).length === 0) {
    throw new Error(
      `Mapping rules failed to produce a valid payload for ${normalizedEntityType}. Check your field mapping configuration.`,
    );
  }
+
+ // Optional: Check if all values are null/undefined
+ const hasNonNullValue = Object.values(hydrated).some(v => v != null);
+ if (!hasNonNullValue) {
+   this.logger.warn(
+     { event: "target_builder.all_null_values", normalizedEntityType },
+     "Hydrated payload contains only null/undefined values"
+   );
+ }

  return hydrated;
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@packages/pipeline/src/fanout/target-builder.service.ts` around lines 118 -
122, The current check in target-builder.service.ts only tests
Object.keys(hydrated).length === 0 which misses payloads like {a: undefined, b:
null}; update the validation around the hydrated object (the variable named
hydrated and the error that references normalizedEntityType) to ensure at least
one property has a non-null/undefined value (e.g. use
Object.values(hydrated).some(v => v != null)) before accepting the payload, and
throw the same error including normalizedEntityType if no valid values remain;
optionally consider stripping null/undefined properties from hydrated before
further processing to avoid sending invalid fields downstream.
packages/pipeline/src/index.ts (2)

16-16: ⚠️ Potential issue | 🔴 Critical | ⚡ Quick win

Fix syntax error: missing line break between exports.

Line 16 concatenates two export statements without proper separation, which creates a syntax error. The exports should be on separate lines.

🐛 Proposed fix
-export * from './sharding/pipeline-hook-broker.service.js';export * from "./pipeline-core.module.js";
+export * from './sharding/pipeline-hook-broker.service.js';
+export * from "./pipeline-core.module.js";
📝 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.

export * from './sharding/pipeline-hook-broker.service.js';
export * from "./pipeline-core.module.js";
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@packages/pipeline/src/index.ts` at line 16, The export statements are
concatenated into a single invalid line; split them into two separate export
statements so each module is exported on its own line—separate "export * from
'./sharding/pipeline-hook-broker.service.js';" and "export * from
\"./pipeline-core.module.js\";" by adding a newline (or a semicolon + newline)
between them so both exports are valid.

21-23: ⚠️ Potential issue | 🟡 Minor | ⚡ Quick win

Remove duplicate export of utils.js.

Lines 21 and 23 both export from './utils.js'. While JavaScript/TypeScript handles duplicate exports without error, this is redundant and may confuse readers.

♻️ Proposed fix
 export * from './normalization/normalization.service.js';
 export * from './utils.js';

-export * from "./utils.js";
-
 export * from './shared/outbox.utils.js';
📝 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.

export * from './utils.js';
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@packages/pipeline/src/index.ts` around lines 21 - 23, Remove the duplicate
export of './utils.js' by keeping a single export statement for utils (remove
one of the two lines that read export * from './utils.js';), so only one export
* from './utils.js' remains in the module index (ensure no other references to
duplicate export remain).
packages/pipeline/src/normalization/dependency-sweeper.service.ts (1)

107-141: ⚠️ Potential issue | 🟠 Major | 🏗️ Heavy lift

Claim-then-send lacks atomicity — traces can get stuck if queue publish fails.

If claimDeferredTrace succeeds (line 112) but queueService.send throws (line 116), the trace is marked as claimed but never enters the queue. Subsequent sweeps skip it because it's already claimed, leaving the trace stuck permanently until manual intervention.

Consider one of these mitigations:

  1. Claim with TTL: Add an expiry/timeout to the claim so stuck traces become eligible again after a period.
  2. Unclaim on send failure: Wrap the send in a try/catch that reverts the claim on failure.
  3. Outbox pattern: Insert an outbox row inside the claim transaction, then publish and mark complete outside.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@packages/pipeline/src/normalization/dependency-sweeper.service.ts` around
lines 107 - 141, The current claim-then-send flow (claimDeferredTrace ->
queueService.send) can leave traces permanently stuck if send fails; change the
logic so the DB claim is reverted on publish failure and only mark
processedTraceIds after a successful send: after calling
this.sweeperRepo.claimDeferredTrace(tenant.tenantId, schemaName, traceId) and
before adding to processedTraceIds, wrap this.queueService.send in a try/catch
that on failure calls a new or existing sweeperRepo.unclaimDeferredTrace (or
releaseClaim/rollbackClaim) for the same tenant/schema/traceId, logs the
failure, and rethrows or continues without adding to processedTraceIds; ensure
processedTraceIds.add(traceId) happens only after send succeeds.
packages/pipeline/src/normalization/normalization.service.ts (2)

189-196: ⚠️ Potential issue | 🟡 Minor | ⚡ Quick win

Misleading .catch — re-throwing doesn't "fall through" as comment suggests.

The .catch((err: any) => { throw err; }) is a no-op: it catches the error and immediately re-throws it, so errors still propagate up. The comment "shard may not exist yet — fall through to piece.normalize" only applies when hookBroker.normalize returns null, not when it throws.

If the intent is to suppress errors and fall through to piece.normalize:

-.catch((err: any) => {
-  throw err;
-}); // shard may not exist yet — fall through to piece.normalize
+.catch(() => null); // shard may not exist yet — return null to fall through to piece.normalize

If the intent is to propagate errors (current behavior), remove the misleading catch and comment:

-  .catch((err: any) => {
-    throw err;
-  }); // shard may not exist yet — fall through to piece.normalize
📝 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.

      const normalizedFromShard = await this.hookBroker
        .normalize(connectionAppName, appProfile, {
          entityType: replica.entityType,
          data: replica.data as Record<string, unknown>,
        })
        .catch(() => null); // shard may not exist yet — return null to fall through to piece.normalize
      const normalizedFromShard = await this.hookBroker
        .normalize(connectionAppName, appProfile, {
          entityType: replica.entityType,
          data: replica.data as Record<string, unknown>,
        });
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@packages/pipeline/src/normalization/normalization.service.ts` around lines
189 - 196, The .catch that re-throws in the call to this.hookBroker.normalize
(which assigns normalizedFromShard) is misleading and redundant; remove the
.catch((err:any)=>{ throw err; }) and the accompanying comment about "fall
through to piece.normalize" so errors correctly propagate from
hookBroker.normalize, or alternatively if you intended to swallow errors return
null in the catch — update the code in normalization.service.ts around the
this.hookBroker.normalize call and adjust the normalizedFromShard handling and
comment accordingly (referenced symbols: this.hookBroker.normalize,
normalizedFromShard, piece.normalize).

262-278: ⚠️ Potential issue | 🔴 Critical | ⚡ Quick win

Queue sends inside transaction violate the stated design principle — dual-write hazard.

Lines 312-323 explicitly warn: "IMPORTANT: publish must happen AFTER the transaction — never inside. Publishing inside the transaction creates a dual-write problem."

However, queueService.send at line 265 is called inside the transaction started at line 212. If the transaction rolls back after these sends (e.g., insertNormalizedOutboxPending at line 295 fails), the queue messages for parent traces are orphaned.

Move the reverse-lookup queue sends outside the transaction, similar to the L3→L4 handoff pattern at lines 324-350:

Suggested approach
+      let parentTraceIds: string[] = [];
       await this.transactionManager.runInTenantTransaction(tenantId, schemaName, async (tx) => {
         // ... upsert logic ...
         if (normalizedEntityId) {
-          let parentTraceIds: string[] = [];
           try {
             await this.transactionManager.runNestedTransaction(tx, async (sp) => {
               await this.hookBroker.writeNormalized(...);
               parentTraceIds = await this.hookBroker.reverseLookup(...);
             });
-
-            for (const pTraceId of parentTraceIds) {
-              await this.queueService.send(QueueName.NormalizedQueue, {...});
-            }
           } catch (hookErr) {
             this.logger.warn(...);
           }
           // ... outbox insert ...
         }
       });
+
+      // Send reverse-lookup requeues AFTER transaction commits
+      for (const pTraceId of parentTraceIds) {
+        await this.queueService.send(QueueName.NormalizedQueue, {
+          traceId: pTraceId,
+          dataSourceId,
+        });
+      }
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@packages/pipeline/src/normalization/normalization.service.ts` around lines
262 - 278, The queueService.send calls that re-queue parent traces (inside the
loop over parentTraceIds) are currently executed within the transaction started
in thisNormalization flow, causing a dual-write risk if the transaction later
rolls back (e.g., insertNormalizedOutboxPending). Move the reverse-lookup queue
sends out of the transaction: collect parentTraceIds (and associated metadata
like replica.entityId and canonicalType) during the transactional section,
commit the transaction, then iterate and call queueService.send and logger.debug
afterwards (mirroring the L3→L4 handoff pattern). Ensure references to
queueService.send, parentTraceIds, replica.entityId, canonicalType, and
insertNormalizedOutboxPending are used to locate and refactor the code.

@coderabbitai

coderabbitai Bot commented Jun 11, 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.

Autofix skipped. No unresolved CodeRabbit review comments with fix instructions found.

@coderabbitai

coderabbitai Bot commented Jun 11, 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.

Autofix skipped. No unresolved CodeRabbit review comments with fix instructions found.

@coderabbitai

coderabbitai Bot commented Jun 11, 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.

Autofix skipped. No unresolved CodeRabbit review comments with fix instructions found.

@pramodnarayana

Copy link
Copy Markdown
Owner Author

@coderabbitai post the review findings as inline comments

@coderabbitai

coderabbitai Bot commented Jun 11, 2026 •

Copy link
Copy Markdown
Contributor

@pramodnarayana I’ll run the review and post any actionable findings as inline comments where applicable.

✅ Action performed

Full review finished.

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

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

Inline comments:
In `@apps/api/package.json`:
- Line 52: The lockfile is out of sync with the observability package (the
dependency entry like "`@soopa/pipeline`" in apps/api's package.json and the
manifest at packages/observability/package.json); run pnpm install at repo root
to regenerate pnpm-lock.yaml (ensuring pnpm vX used by CI), verify the install
succeeds with --frozen-lockfile, and commit the updated pnpm-lock.yaml so CI no
longer fails (repeat for any other affected package.json mentioned by the
review).
- Around line 37-82: Several files still import from `@soopa/engine` (e.g.,
StorageResolverService, StorageResolverModule, TriggerExecutorService,
WebhooksController, DataExplorerService, TraceService) which will break
compilation; open each listed file and replace imports from '`@soopa/engine`' with
the appropriate '`@soopa/pipeline`' exports (or remove unused engine imports),
update any renamed symbols to the pipeline equivalents (e.g.,
StorageResolverService/StorageResolverModule → pipeline-provided resolver if
renamed), adjust import paths and types accordingly in trigger-executor.service,
trigger.module, webhooks.controller/module, trace.* and their spec files, then
run the app/tests to confirm compilation and update any broken usages or mocks
in the spec files to match the new pipeline API.

In `@apps/api/src/modules/observability/observability.module.ts`:
- Line 9: Remove the misleading inline comment in the ObservabilityModule
exports line; update the exports to only export MetricsService and either delete
the uncertain comment or replace it with a factual note stating that
LoggerModule is global and PinoLogger is available app-wide (so
ObservabilityModule does not need to export LoggerModule). Edit the exports
array in ObservabilityModule (the line exporting MetricsService) and remove or
correct the comment referencing LoggerModule and PinoLogger so it no longer
implies uncertainty (also ensure any references to MappingsModule or
mappings.service.ts remain unchanged).
🪄 Autofix (Beta)

✅ Autofix completed


ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: b897da13-b9da-4c7f-9b71-20a49bb1bca2

📥 Commits

Reviewing files that changed from the base of the PR and between f1698ac and ca622a2.

📒 Files selected for processing (66)
  • apps/api/package.json
  • apps/api/src/app/app.module.ts
  • apps/api/src/modules/ai/ai.module.ts
  • apps/api/src/modules/ai/controllers/ai.controller.spec.ts
  • apps/api/src/modules/ai/controllers/ai.controller.ts
  • apps/api/src/modules/ai/controllers/transformer-simulation.controller.ts
  • apps/api/src/modules/connections/connection-lifecycle.service.spec.ts
  • apps/api/src/modules/connections/connection-lifecycle.service.ts
  • apps/api/src/modules/connections/connections.module.ts
  • apps/api/src/modules/connections/connectors.service.spec.ts
  • apps/api/src/modules/connections/connectors.service.ts
  • apps/api/src/modules/dbmanager/dbmanager.module.ts
  • apps/api/src/modules/exceptions/exception.service.spec.ts
  • apps/api/src/modules/exceptions/exception.service.ts
  • apps/api/src/modules/exceptions/exceptions.module.ts
  • apps/api/src/modules/identity/auth/auth.module.ts
  • apps/api/src/modules/observability/observability.module.ts
  • apps/api/src/modules/pipeline/pipeline.module.ts
  • apps/api/src/modules/scheduler/connection-sync-runner.spec.ts
  • apps/api/src/modules/scheduler/connection-sync-runner.ts
  • apps/api/src/modules/scheduler/scheduler.module.ts
  • apps/api/src/modules/scheduler/sync-runner.ts
  • apps/api/src/modules/stitches/stitches.module.ts
  • packages/pipeline/src/shared/adapters/replica-state.adapter.integration.spec.ts
  • packages/pipeline/src/shared/adapters/replica-state.adapter.ts
  • packages/pipeline/src/shared/fakes/fake-connection.repository.ts
  • packages/pipeline/src/shared/fakes/fake-dependency-sweeper.repository.ts
  • packages/pipeline/src/shared/fakes/fake-normalization.repository.ts
  • packages/pipeline/src/shared/fakes/fake-outbound-gateway.repository.ts
  • packages/pipeline/src/shared/fakes/fake-pipeline-state.repository.ts
  • packages/pipeline/src/shared/fakes/fake-sync-log.repository.ts
  • packages/pipeline/src/shared/fakes/fake-transaction-manager.ts
  • packages/pipeline/src/shared/interfaces/outbound-dispatcher.interface.ts
  • packages/pipeline/src/shared/outbox.utils.spec.ts
  • packages/pipeline/src/shared/outbox.utils.ts
  • packages/pipeline/src/shared/ports/connection.repository.port.ts
  • packages/pipeline/src/shared/ports/dependency-sweeper.repository.port.ts
  • packages/pipeline/src/shared/ports/field-mapping.repository.port.ts
  • packages/pipeline/src/shared/ports/global-entity-map.repository.port.ts
  • packages/pipeline/src/shared/ports/normalization.repository.port.ts
  • packages/pipeline/src/shared/ports/outbound-gateway.repository.port.ts
  • packages/pipeline/src/shared/ports/pipeline-state.repository.port.ts
  • packages/pipeline/src/shared/ports/routing.repository.port.ts
  • packages/pipeline/src/shared/ports/stitch.repository.port.ts
  • packages/pipeline/src/shared/ports/sync-log.repository.port.ts
  • packages/pipeline/src/shared/ports/transaction-manager.port.ts
  • packages/pipeline/src/state/cursor-manager.service.spec.ts
  • packages/pipeline/src/state/cursor-manager.service.ts
  • packages/pipeline/src/state/cursor-manager.types.ts
  • packages/pipeline/src/storage-resolver/storage-resolver.module.ts
  • packages/pipeline/src/storage-resolver/storage-resolver.service.spec.ts
  • packages/pipeline/src/storage-resolver/storage-resolver.service.ts
  • packages/pipeline/src/test-utils/mock-db-manager.ts
  • packages/pipeline/src/utils.ts
  • packages/pipeline/src/vitest-env.d.ts
  • packages/pipeline/tsconfig.build.json
  • packages/pipeline/tsconfig.json
  • packages/pipeline/vitest.config.ts
  • packages/utils/package.json
  • packages/utils/src/index.ts
  • packages/utils/src/path-utils.ts
  • packages/utils/tsconfig.json
  • plugins/quickbooks/core/vitest.config.ts
  • plugins/salesforce/core/vitest.config.ts
  • vitest.integration.config.ts
  • vitest.shared.ts

Comment thread apps/api/package.json
Comment thread apps/api/package.json
Comment thread apps/api/src/modules/observability/observability.module.ts
@pramodnarayana

Copy link
Copy Markdown
Owner Author

@coderabbitai please re-review this PR. The previous review appears incomplete. Review all 60 changed files and provide findings for every file that requires attention. If review limits or context limits prevent a full review, please state which files were reviewed and which files were skipped.

@coderabbitai

coderabbitai Bot commented Jun 11, 2026 •

Copy link
Copy Markdown
Contributor

@pramodnarayana Understood — I'll kick off a complete re-review of all changed files in this PR. I'll post inline comments on every file that requires attention and will note any files that are skipped due to context or size constraints.

✅ Action performed

Full review finished.

@coderabbitai

coderabbitai Bot commented Jun 11, 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 12 file(s) based on 1 unresolved review comment.

Files modified:

  • apps/api/src/modules/trace/data-explorer.service.spec.ts
  • apps/api/src/modules/trace/data-explorer.service.ts
  • apps/api/src/modules/trace/trace.module.ts
  • apps/api/src/modules/trace/trace.service.spec.ts
  • apps/api/src/modules/trace/trace.service.ts
  • apps/api/src/modules/trigger/trigger-executor.service.spec.ts
  • apps/api/src/modules/trigger/trigger-executor.service.ts
  • apps/api/src/modules/trigger/trigger.module.ts
  • apps/api/src/modules/webhooks/webhooks.controller.spec.ts
  • apps/api/src/modules/webhooks/webhooks.controller.ts
  • apps/api/src/modules/webhooks/webhooks.module.ts
  • apps/worker/package.json

Commit: e4abc3818d9fff9e1f0e40dda2a757130bd038dc

The changes have been pushed to the refactor/monolith-boundaries branch.

Time taken: 7m 11s

coderabbitai Bot and others added 2 commits June 11, 2026 09:32
Fixed 12 file(s) based on 1 unresolved review comment.

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: 8

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

Inline comments:
In `@apps/api/src/scripts/test-e2e-ingestion.ts`:
- Around line 57-63: The query that retrieves a test connection (in
test-e2e-ingestion.ts where pool.query is used to populate res and the similar
query later) currently selects any ACTIVE connection and can pick incompatible
apps; change the WHERE clause to restrict to a known-compatible app/profile
(e.g., add "AND app_name = '<expected_app>'" or "AND connector_profile =
'<expected_profile>'") and add a deterministic ORDER BY (e.g., ORDER BY id ASC)
before LIMIT 1 so the chosen connection is both compatible with the fixed
payload and deterministic across runs; update both occurrences (the query that
fills res and the later similar query) to use the same filters and ordering.

In `@apps/worker/package.json`:
- Around line 21-26: Update the worker package.json db scripts that currently
point to the removed entrypoint "src/db/db-cli.ts": change each "db:*" script
(db:drop, db:seed, db:fresh, db:reset, db:provision:local, db:check-role) to
invoke the existing CLI entrypoint used by the API (e.g. "tsx src/db/cli/main.ts
<command>") or remove any unused db:* scripts; ensure you preserve the same
command arguments (drop, seed, fresh, reset, provision:local, check-role) so the
scripts map one-to-one to the corresponding handlers in the new CLI.

In `@apps/worker/src/consumers/copilot.worker.ts`:
- Around line 204-218: The else branch currently publishes both the error and
the trailing “[DONE]” to the job stream via this.redis.publish and then throws,
which causes the catch handler to publish them again; remove the two
this.redis.publish calls that send `error: ${errorMsg}\n` and `[DONE]\n` from
that else branch so the code only appends the failure to chatPersistence
(this.chatPersistence.appendMessage) and throws the Error, allowing the existing
catch-path publishers to emit the stream error and terminal marker once.

In `@apps/worker/src/consumers/gitops-sync.worker.ts`:
- Around line 16-18: The default SHARD_BASE_PATH constant currently falls back
to "../../sync/application" which diverges from ApplicationLoaderService's
fallback; update the SHARD_BASE_PATH declaration (private readonly
SHARD_BASE_PATH) to use the same default path used by ApplicationLoaderService
(i.e. "../../engine/sync/application") or otherwise reference the same shared
fallback logic so both consumers use an identical default when
SHARD_APPLICATION_PATH is unset.

In `@packages/database/src/test-db.ts`:
- Around line 72-75: The createSchema method in TestDatabaseManager uses a raw
object with sql string and params when calling this.db.execute; instead use
Drizzle's sql tagged template and identifier helpers to preserve typing and safe
quoting (e.g. replace the object call in createSchema with
this.db.execute(sql`CREATE SCHEMA IF NOT EXISTS ${sql.identifier(schemaName)}`)
so the call uses the Drizzle sql object and sql.identifier(schemaName) for the
schema name.

In `@packages/dbmanager/src/impl/sql-database-manager.ts`:
- Line 562: The domain port type OutboundGatewayRecord in
packages/domain/core/src/ports/outbound-gateway.port.ts is missing the
'DEFERRED_DEPENDENCY' status even though the DB constraint and repository
(drizzle-outbound-gateway.repository.ts) and DependencySweeperService already
use it; update the OutboundGatewayRecord.status union to include
'DEFERRED_DEPENDENCY' so consumers modeling DB rows accept that value, and run
TypeScript build/tests to catch any places (e.g., delivery claiming logic) that
must explicitly handle or intentionally ignore this status.

In `@packages/pipeline/src/fanout/fanout-batch-processor.ts`:
- Around line 280-288: The catch block that handles MQ publish errors should not
call upsertPendingOutboundGateway(...) again (which keeps the outbox in pending
and causes immediate reprocessing); instead, call the outbox repository's
explicit failed-state transition method on the same identifiers (e.g., a method
like
markOutboundGatewayFailed/transitionOutboundToFailed/upsertFailedOutboundGateway)
using srcTenantId, destSchemaName, traceId, stitch.id, stitch.destDataSourceId,
dataSourceId and any error details, then rethrow the error. Update the catch
path in fanout-batch-processor.ts to remove the upsertPendingOutboundGateway
call and invoke the repository's failure transition so the outbox record moves
to failed state before throwing.

In `@packages/pipeline/src/fanout/fanout-router.service.ts`:
- Around line 89-107: The code acquires a sync lock (via insert/select on
active_sync_locks after calling stateRepo.getReplicaSourceVendorId) but some
early exits (returns or throws before the stitch-loop finally) skip calling
releaseSyncLock; wrap the lock acquisition and all subsequent logic in a
try/finally (or use a disposer returned by a new acquireSyncLock helper) and
ensure releaseSyncLock is invoked in the finally block so every path after the
INSERT/SELECT calls releaseSyncLock; update the method in
fanout-router.service.ts that contains getReplicaSourceVendorId, the
INSERT/SELECT block, and any early return points (including the
superseded-return and exception paths noted around lines ~89-205) to use this
pattern.
🪄 Autofix (Beta)

✅ Autofix completed


ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: b0c11a07-ff8f-44a5-ad82-3cf3fc0acd95

📥 Commits

Reviewing files that changed from the base of the PR and between e4abc38 and 61d3b1d.

⛔ Files ignored due to path filters (2)
  • packages/dbmanager/src/impl/__snapshots__/sql-database-manager.spec.ts.snap is excluded by !**/*.snap
  • pnpm-lock.yaml is excluded by !**/pnpm-lock.yaml
📒 Files selected for processing (56)
  • apps/api/src/modules/trace/trace.module.ts
  • apps/api/src/modules/trigger/trigger.module.ts
  • apps/api/src/modules/webhooks/webhooks.module.ts
  • apps/api/src/modules/workspaces/workspaces.module.ts
  • apps/api/src/scripts/admin-bootstrap.ts
  • apps/api/src/scripts/test-e2e-ingestion.ts
  • apps/api/test/invitations.e2e-spec.ts
  • apps/api/test/system-admin.tenants.e2e-spec.ts
  • apps/api/vitest.config.mts
  • apps/web/vitest.config.ts
  • apps/worker/package.json
  • apps/worker/src/app.module.ts
  • apps/worker/src/bootstrap/dbmanager/dbmanager.module.ts
  • apps/worker/src/consumers/copilot.worker.ts
  • apps/worker/src/consumers/gitops-sync.worker.ts
  • apps/worker/src/pollers/inbound-outbox.poller.ts
  • apps/worker/src/pollers/normalized-outbox.poller.ts
  • apps/worker/vitest.config.mts
  • package.json
  • packages/database/drizzle/global/0018_gorgeous_alice.sql
  • packages/database/drizzle/global/meta/_journal.json
  • packages/database/drizzle/tenant/0012_military_black_queen.sql
  • packages/database/drizzle/tenant/meta/_journal.json
  • packages/database/package.json
  • packages/database/src/client.ts
  • packages/database/src/schema/tenant/gem.ts
  • packages/database/src/test-db.ts
  • packages/dbmanager/package.json
  • packages/dbmanager/src/impl/sql-database-manager.spec.ts
  • packages/dbmanager/src/impl/sql-database-manager.ts
  • packages/domain/core/package.json
  • packages/domain/tms/package.json
  • packages/observability/package.json
  • packages/pipeline/fix-imports.cjs
  • packages/pipeline/package.json
  • packages/pipeline/src/delivery/delivery.service.integration.spec.ts
  • packages/pipeline/src/delivery/use-cases/claim-delivery.use-case.spec.ts
  • packages/pipeline/src/delivery/use-cases/claim-delivery.use-case.ts
  • packages/pipeline/src/fanout/fanout-batch-processor.spec.ts
  • packages/pipeline/src/fanout/fanout-batch-processor.ts
  • packages/pipeline/src/fanout/fanout-router.service.spec.ts
  • packages/pipeline/src/fanout/fanout-router.service.ts
  • packages/pipeline/src/fanout/target-builder.service.ts
  • packages/pipeline/src/index.ts
  • packages/pipeline/src/normalization/dependency-sweeper.service.spec.ts
  • packages/pipeline/src/normalization/dependency-sweeper.service.ts
  • packages/pipeline/src/normalization/normalization.service.ts
  • packages/pipeline/src/replication/registry-token-refresh.service.spec.ts
  • packages/pipeline/src/shared/adapters/drizzle-dependency-sweeper.repository.ts
  • packages/pipeline/src/shared/adapters/drizzle-outbound-gateway.repository.ts
  • plugins/quickbooks/core/package.json
  • plugins/salesforce/core/package.json
  • plugins/shared/mapping/src/mapping-engine.ts
  • pnpm-workspace.yaml
  • test/integration/setup/global-setup.ts
  • test/integration/setup/test-env-setup.ts
💤 Files with no reviewable changes (1)
  • packages/pipeline/src/delivery/delivery.service.integration.spec.ts

Comment thread apps/api/src/scripts/test-e2e-ingestion.ts
Comment thread apps/worker/package.json Outdated
Comment thread apps/worker/src/consumers/copilot.worker.ts
Comment on lines +16 to +18
private readonly SHARD_BASE_PATH =
process.env.SHARD_APPLICATION_PATH ||
path.resolve(process.cwd(), "../../engine/sync/application");
path.resolve(process.cwd(), "../../sync/application");

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 | ⚡ Quick win

Default shard path now diverges from ApplicationLoaderService fallback.

Line 18 defaults to ../../sync/application, while packages/pipeline/src/sharding/application-loader.service.ts still defaults to ../../engine/sync/application. With SHARD_APPLICATION_PATH unset, sync/invalidation and loading will operate on different directories.

Proposed fix
   private readonly SHARD_BASE_PATH =
     process.env.SHARD_APPLICATION_PATH ||
-    path.resolve(process.cwd(), "../../sync/application");
+    path.resolve(process.cwd(), "../../engine/sync/application");
📝 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
private readonly SHARD_BASE_PATH =
process.env.SHARD_APPLICATION_PATH ||
path.resolve(process.cwd(), "../../engine/sync/application");
path.resolve(process.cwd(), "../../sync/application");
private readonly SHARD_BASE_PATH =
process.env.SHARD_APPLICATION_PATH ||
path.resolve(process.cwd(), "../../engine/sync/application");
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@apps/worker/src/consumers/gitops-sync.worker.ts` around lines 16 - 18, The
default SHARD_BASE_PATH constant currently falls back to
"../../sync/application" which diverges from ApplicationLoaderService's
fallback; update the SHARD_BASE_PATH declaration (private readonly
SHARD_BASE_PATH) to use the same default path used by ApplicationLoaderService
(i.e. "../../engine/sync/application") or otherwise reference the same shared
fallback logic so both consumers use an identical default when
SHARD_APPLICATION_PATH is unset.

Comment thread packages/database/src/test-db.ts
Comment thread packages/dbmanager/src/impl/sql-database-manager.ts
Comment thread packages/pipeline/src/fanout/fanout-batch-processor.ts Outdated
Comment thread packages/pipeline/src/fanout/fanout-router.service.ts
@coderabbitai

coderabbitai Bot commented Jun 11, 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 8 file(s) based on 8 unresolved review comments.

Files modified:

  • apps/api/src/scripts/test-e2e-ingestion.ts
  • apps/worker/package.json
  • apps/worker/src/consumers/copilot.worker.ts
  • apps/worker/src/consumers/gitops-sync.worker.ts
  • packages/database/src/test-db.ts
  • packages/domain/core/src/ports/outbound-gateway.port.ts
  • packages/pipeline/src/fanout/fanout-batch-processor.ts
  • packages/pipeline/src/fanout/fanout-router.service.ts

Commit: b3dd0619acab2e6cd6f19be26b5623fe967b217f

The changes have been pushed to the refactor/monolith-boundaries branch.

Time taken: 8m 46s

coderabbitai Bot and others added 2 commits June 11, 2026 10:36
Fixed 8 file(s) based on 8 unresolved review comments.

Co-authored-by: CodeRabbit <noreply@coderabbit.ai>
@pramodnarayana
pramodnarayana merged commit ebbdeab into development Jun 11, 2026
2 checks passed
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