Skip to content

initial sync testing; decouple pipeline and provision - #169

Merged
pramodnarayana merged 5 commits into
developmentfrom
initial-sync-testing
Jun 20, 2026
Merged

pramodnarayana merged 5 commits into
developmentfrom
initial-sync-testing

Conversation

@pramodnarayana

@pramodnarayana pramodnarayana commented Jun 20, 2026 •

Copy link
Copy Markdown
Owner

Summary by CodeRabbit

  • New Features

    • Strengthened OAuth credential storage and reconnection by encrypting persisted values and reusing previously saved client credentials.
    • Added schema provisioning support that registers CDC publication tables during provisioning.
    • Improved ingestion throughput with batch-based gateway/outbox writes and continued outbox processing until empty.
  • Improvements

    • Refreshed Salesforce polling to use cursor-based pagination with stricter response validation and better field discovery.
    • Enhanced plugin metadata extraction and workspace bootstrapping.
  • Configuration

    • Updated CDC relay/debezium configuration (sink URL, header key) and added support for REGISTRY_PLUGIN, DEBEZIUM_SECRET, and CDC_RELAY_DISABLE_AUTH.

@coderabbitai

coderabbitai Bot commented Jun 20, 2026 •

Copy link
Copy Markdown
Contributor

Review Change Stack

Warning

Review limit reached

@pramodnarayana, we couldn't start this review because you've reached your PR review rate limit.

More reviews will be available in 40 minutes and 3 seconds. Learn how PR review limits work.

Your organization has run out of usage credits. Purchase more credits in the billing tab to continue.

⌛ How to resolve this issue?

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

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

🚦 How do rate limits work?

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

For paid Pro and Pro+ PR reviews, CodeRabbit uses adaptive limits for sustained high-volume activity. When a developer's recent PR review activity reaches the 95th percentile or higher among CodeRabbit users, the refill rate gradually slows as usage increases. The highest same-day bursts are limited more strictly.

Please see our Fair Usage Limits Policy for further information.

ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: 3fb2c195-6bb8-4e79-93cf-65725c671140

📥 Commits

Reviewing files that changed from the base of the PR and between 6faae43 and 1c26b4a.

📒 Files selected for processing (5)
  • apps/api/src/modules/scheduler/connection-sync-runner.ts
  • packages/pipeline/src/fanout/fanout-router.service.spec.ts
  • packages/pipeline/src/fanout/fanout-router.service.ts
  • packages/pipeline/src/plugin-hooks/pipeline-hook-broker.service.spec.ts
  • packages/pipeline/src/plugin-hooks/pipeline-hook-broker.service.ts
📝 Walkthrough

Walkthrough

This PR extracts a new @soopa/provision package from @soopa/pipeline housing replication, schema provisioning, and OAuth refresh services. PipelineHookBrokerService is rewritten to use direct @soopa/piece-framework registry lookups instead of EventEmitter2 shard events. OAuth credential encryption is moved into StoreOAuthConnectionUseCase/GetExistingOAuthCredentialsUseCase. The SDK plugin registry gains WorkspacePluginRegistrar, extractPieceMetadata, and a bootstrapper abstraction. Salesforce polling is rewritten to use nextRecordsUrl cursor pagination.

Changes

@soopa/provision extraction and pipeline infrastructure cleanup

Layer / File(s) Summary
@soopa/provision package scaffolding
packages/provision/package.json, packages/provision/tsconfig.json, packages/provision/tsconfig.build.json, packages/provision/vitest.config.ts, packages/provision/src/index.ts
Creates the @soopa/provision package manifest, TypeScript configs, Vitest config, and the public barrel index re-exporting provisioning/registry services, the replication port, and ProvisionModule.
RegistryReplicationPort.registerCdcTables contract and adapter
packages/provision/src/shared/ports/registry-replication.port.ts, packages/provision/src/shared/adapters/outbound/registry-replication.adapter.ts, packages/provision/src/registry/registry-replication.service.ts, packages/provision/src/provisioning/schema-provision.worker.ts, packages/provision/src/registry/registry-replication.service.spec.ts, packages/pipeline/src/shared/domain.ts
Extends RegistryReplicationPort with registerCdcTables; implements it in RegistryReplicationAdapter using PostgreSQL DO $$ DDL blocks with duplicate_object suppression; updates import paths across provision services and removes the port re-export from pipeline's domain barrel.
ProvisionModule wiring, ProvisionSchemaUseCase CDC step, token refresh cleanup
packages/provision/src/provision.module.ts, packages/provision/src/provisioning/use-cases/provision-schema.use-case.ts, packages/provision/src/provisioning/use-cases/provision-schema.use-case.spec.ts, packages/provision/src/registry/registry-token-refresh.service.ts, packages/provision/src/registry/registry-token-refresh.service.integration.spec.ts, packages/provision/src/test-utils/mock-db-manager.ts
Creates ProvisionModule binding RegistryReplicationPort to RegistryReplicationAdapter; adds registerCdcTables call in ProvisionSchemaUseCase after applyPlan; removes dummy constructor injection from RegistryOAuthRefreshClient; improves integration test with a safe test subclass; adds MockDatabaseManager utility.
PipelineHookBrokerService: EventEmitter2 → direct registry lookup
packages/pipeline/src/plugin-hooks/pipeline-hook-broker.service.ts, packages/pipeline/src/plugin-hooks/pipeline-hook-broker.service.spec.ts, packages/pipeline/src/fanout/fanout-batch-processor.ts, packages/pipeline/src/fanout/fanout-batch-processor.spec.ts
Rewrites PipelineHookBrokerService to call hook handlers directly from the piece-framework registry; removes EventEmitter2 injection from FanoutBatchProcessor; adds complete test coverage for the new registry-based dispatch; adds extractSyncTokenFromState helper.
PipelineCoreModule cleanup and fanout router hardening
packages/pipeline/src/pipeline-core.module.ts, packages/pipeline/src/index.ts, packages/pipeline/src/normalization/normalization.service.ts, packages/pipeline/src/fanout/target-builder.service.ts, packages/pipeline/src/fanout/target-builder.service.*.spec.ts, packages/pipeline/src/fanout/fanout-router.service.ts, packages/pipeline/src/shared/adapters/outbound/drizzle-stitch.adapter.ts, packages/pipeline/src/delivery/delivery.service.ts, packages/pipeline/src/replication/replica.service.spec.ts
Removes ApplicationLoaderModule, RegistryReplicationService, SchemaProvisionWorker, RegistryOAuthRefreshClient from PipelineCoreModule; updates barrel index to plugin-hooks exports; adds schema name validation and sync-lock locked_by_trace_id/expires_at in FanoutRouterService; switches DrizzleStitchAdapter from DB_MANAGER to DATABASE_CONNECTION with explicit tenant filters.
DatabaseProvisioner and SchemaProvisioner adapter/port removal
apps/api/src/modules/trigger/core/use-cases/enable-trigger.use-case.ts, apps/api/src/modules/trigger/core/use-cases/enable-trigger.use-case.spec.ts, apps/api/src/modules/trigger/trigger.module.ts, apps/api/src/modules/stitches/stitches.module.ts, apps/api/src/modules/pipeline/pipeline.module.ts
Deletes PipelineDatabaseProvisionerAdapter, DatabaseProvisionerPort, FakeDatabaseProvisioner, DbManagerSchemaProvisionerAdapter, and SchemaProvisionerPort; simplifies EnableTriggerUseCase to only validate schema name and call onEnable; removes PipelineCdcListener from PipelineModule.
Worker and API app wiring, Debezium config
apps/api/package.json, apps/api/src/app/app.module.ts, apps/api/src/config/env.validation.ts, apps/worker/package.json, apps/worker/src/app.module.ts, apps/worker/src/config/env.validation.ts, apps/worker/src/adapters/outbound/nest-cache-invalidator.adapter.ts, infra/debezium/application.properties, package.json, packages/database/drizzle/tenant/0013_bright_molten_man.sql
Adds @soopa/provision to worker/API packages; imports ProvisionModule in worker replacing ApplicationLoaderModule; switches RegistryOAuthRefreshClient import to @soopa/provision; makes NestCacheInvalidatorAdapter a no-op; fixes Debezium sink URL and Authorization header key; updates env schemas for REGISTRY_PLUGIN/DEBEZIUM_SECRET; fixes migration DROP INDEX.
Scheduler batched gateway ingestion and outbox drain loop
apps/api/src/modules/scheduler/connection-sync-runner.ts, apps/api/src/modules/scheduler/connection-sync-runner.spec.ts, apps/worker/src/core/use-cases/outbox/process-outbox.use-case.ts, apps/worker/src/core/use-cases/outbox/process-outbox.use-case.spec.ts, apps/worker/src/core/fakes/outbox-repository.fake.ts, apps/worker/vitest.config.mts
Replaces per-record insertGatewayRow with insertGatewayBatch doing chunked transactional inserts with deterministic extReqId; changes ProcessOutboxUseCase.execute to a while(true) drain loop; fixes FakeOutboxRepository claimNextBatch to allow only PENDING/undefined rows; adds onPermanentFailure tests; configures vitest forks pool.

OAuth connection credential encryption rework

Layer / File(s) Summary
GetExistingOAuthCredentialsUseCase and StoreOAuthConnectionUseCase encryption
apps/api/src/modules/connections/core/use-cases/get-existing-oauth-credentials.use-case.ts, apps/api/src/modules/connections/core/use-cases/store-oauth-connection.use-case.ts, apps/api/src/modules/connections/core/use-cases/store-oauth-connection.use-case.spec.ts
Adds GetExistingOAuthCredentialsUseCase to decrypt and parse stored credentials returning clientId/clientSecret; injects IEncryptionService into StoreOAuthConnectionUseCase to encrypt the connection value before persistence; updates tests with mock crypto.
OAuthController reconnect flow and ConnectionsModule wiring
apps/api/src/modules/connections/connections/oauth.controller.ts, apps/api/src/modules/connections/connections.module.ts
Removes crypto injection from OAuthController; uses GetExistingOAuthCredentialsUseCase for reconnect credential lookup; stores raw JSON string instead of encrypting in controller; adds clientId/clientSecret trimming; tightens metadata typing; registers both use-cases in ConnectionsModule.

SDK registry plugin system: WorkspacePluginRegistrar, extractPieceMetadata, and bootstrapper abstraction

Layer / File(s) Summary
extractPieceMetadata utility and UpsertPieceDto extension
sdk/registry/src/pieces/piece-metadata.util.ts, sdk/registry/src/pieces/piece-metadata.util.spec.ts, sdk/registry/src/pieces/piece-repository.port.ts, sdk/registry/src/pieces/drizzle-piece.repository.ts, sdk/registry/src/index.ts
Adds ExtractedPieceMetadata interface and extractPieceMetadata function with four discovery strategies and runtime type validation; extends UpsertPieceDto with five optional fields; updates DrizzlePieceRepository to persist all new fields; exports the utility from the SDK index.
WorkspacePluginRegistrar rename and LocalFilePieceResolver update
sdk/registry/src/pieces/workspace-plugin-registrar.ts, sdk/registry/src/pieces/workspace-plugin-registrar.spec.ts, sdk/registry/src/pieces/local-file.piece-resolver.ts, sdk/registry/src/pieces/local-file.piece-resolver.spec.ts
Renames LocalDevPluginSyncService to WorkspacePluginRegistrar; uses extractPieceMetadata in the scan loop; adds initPromise-based idempotency guard; updates LocalFilePieceResolver to await initialize() before path lookup; updates specs.
IPluginBootstrapper abstraction and two implementations
sdk/registry/src/pieces/plugin-bootstrapper.port.ts, sdk/registry/src/pieces/plugin-metadata.bootstrapper.ts, sdk/registry/src/pieces/plugin-registry.bootstrapper.ts, sdk/registry/src/metadata/metadata-discovery.service.ts
Defines PLUGIN_BOOTSTRAPPER token and IPluginBootstrapper interface; adds PluginMetadataBootstrapper and PluginRegistryBootstrapper; hardens MetadataDiscoveryService with ServiceUnavailableException wrapping on describeObjects/describeFields errors.
PiecesModule bootstrapper wiring and install-piece use-case
sdk/registry/src/pieces/pieces.module.ts, apps/worker/src/core/use-cases/app-installer/install-piece.use-case.ts, apps/worker/src/core/use-cases/app-installer/app-installer.spec.ts
Rewires PiecesModule to call bootstrapper.bootstrap() via PLUGIN_BOOTSTRAPPER before loading pieces; introduces BOOTSTRAPPER_PROVIDER selecting implementation based on env; updates install-piece use-case to use extractPieceMetadata; updates tests.

Salesforce plugin: pagination overhaul, fetch adapter refactor, Revenova updates

Layer / File(s) Summary
Salesforce poll pagination and describeObjects hardening
plugins/salesforce/core/src/use-cases/salesforce.use-cases.ts, plugins/salesforce/core/src/use-cases/salesforce.use-cases.spec.ts, plugins/salesforce/core/src/index.spec.ts
Rewrites poll to use describeFields for first-page field selection and nextRecordsUrl for subsequent pages; removes lastId cursor logic; adds describeObjects array validation; updates unit tests with describeStub helper and new paging coverage.
NativeFetchAdapter getSignal() centralization
plugins/salesforce/core/src/adapters/native-fetch.adapter.ts
Extracts a private getSignal() helper consolidating AbortSignal/AbortController/timeout setup; updates get(), post(), and patch() to use signal: this.getSignal(signal) removing per-request controller construction and combineSignals().
Revenova piece registration keys, async normalizer, and attributes handling
plugins/salesforce/revenova-piece/src/index.ts, plugins/salesforce/revenova-piece/src/normalizeRevenovaToTms.ts, plugins/salesforce/revenova-piece/src/upsertRevenovaObject.ts, plugins/salesforce/*/package.json
Changes Revenova hook registration keys to ('salesforce-revenova', 'standard', ...); converts normalizeRevenovaToTms to async with awaited JSONata evaluation; adds attributes key handling in upsertRevenovaObject for doctype derivation; bumps package versions to 0.6.8.

Estimated code review effort

🎯 4 (Complex) | ⏱️ ~75 minutes

Possibly related PRs

  • pramodnarayana/nexiom#133: Both PRs touch RegistryOAuthRefreshClient wiring—this PR moves it to @soopa/provision while #133 introduced/used it in the connections token-refresh setup.
  • pramodnarayana/nexiom#158: Both PRs refactor OAuth/connection crypto wiring in apps/api/src/modules/connections using IEncryptionService via the ENCRYPTION_SERVICE NestJS token.
  • pramodnarayana/nexiom#166: Both PRs modify the same OAuth connection exchange flow in OAuthController.exchangeCode, with changes to credential retrieval, encryption handling, and storage behavior.

Poem

🐇 A rabbit once juggled a pipeline too wide,
So she carved out a provision to run by its side.
The shards and the emitters? Now shelved on a shelf!
The hooks speak to registry, direct, by themself.
With cursors for Salesforce and crypto tucked neat,
The warren runs leaner—oh, isn't it sweet! 🌿

🚥 Pre-merge checks | ✅ 5
✅ Passed checks (5 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title clearly identifies the main objective of decoupling pipeline and provision components with a secondary focus on initial sync testing.
Docstring Coverage ✅ Passed Docstring coverage is 100.00% which is sufficient. The required threshold is 80.00%.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.

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

✨ Finishing Touches
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch initial-sync-testing

Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out.

❤️ Share

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

@coderabbitai coderabbitai Bot left a comment •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Actionable comments posted: 17

Caution

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

⚠️ Outside diff range comments (2)
packages/pipeline/src/fanout/fanout-batch-processor.ts (1)

40-54: 🧹 Nitpick | 🔵 Trivial | ⚡ Quick win

Inject PipelineHookBrokerService instead of constructing it manually.

new PipelineHookBrokerService() bypasses Nest DI and will break this class as soon as the broker needs injected dependencies. Prefer constructor injection here.

♻️ Proposed refactor
 export class FanoutBatchProcessor {
   private readonly logger = new Logger(FanoutBatchProcessor.name);
-  private readonly broker: PipelineHookBrokerService;
 
   constructor(
     private readonly queueService: QueueService,
     private readonly storageResolver: StorageResolverService,
     private readonly targetBuilder: TargetBuilderService,
+    private readonly broker: PipelineHookBrokerService,
     `@Inject`(CONNECTION_REPOSITORY_PORT) private readonly connRepo: ConnectionRepositoryPort,
     `@Inject`(PIPELINE_STATE_REPOSITORY_PORT) private readonly stateRepo: PipelineStateRepositoryPort,
     `@Inject`(GLOBAL_ENTITY_MAP_REPOSITORY_PORT) private readonly gemRepo: GlobalEntityMapRepositoryPort,
     `@Inject`(FIELD_MAPPING_REPOSITORY_PORT) private readonly fieldMappingRepo: FieldMappingRepositoryPort,
     `@Inject`(SYNC_LOG_REPOSITORY_PORT) private readonly syncLogRepo: SyncLogRepositoryPort,
     `@Inject`(OUTBOUND_GATEWAY_REPOSITORY_PORT) private readonly outboxRepo: OutboundGatewayRepositoryPort,
-  ) {
-    this.broker = new PipelineHookBrokerService();
-  }
+  ) {}
🤖 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 40 - 54,
In the FanoutBatchProcessor constructor, replace the manual instantiation of
PipelineHookBrokerService in the constructor body with dependency injection. Add
PipelineHookBrokerService as a constructor parameter (with an `@Inject` decorator
if it has a provider token, similar to how the other repository services are
injected), and remove the line that creates a new instance with this.broker =
new PipelineHookBrokerService(). Instead, assign the injected parameter directly
to the private broker field, following the same pattern used for queueService,
storageResolver, and other dependencies.
sdk/registry/src/pieces/local-file.piece-resolver.spec.ts (1)

24-31: ⚠️ Potential issue | 🟡 Minor | ⚡ Quick win

Assert the new initialize() call in resolver tests.

The tests were updated with an initialize mock, but they never assert it is invoked. A regression removing Line 24 would currently still pass.

Suggested assertion
   it('should resolve a piece from local file system using dynamic import', async () => {
@@
     const result = await resolver.resolve('`@soopa/existing`');
+    expect(mockLocalSync.initialize).toHaveBeenCalledTimes(1);
     expect(result.module).toBe('test-module');
🤖 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 `@sdk/registry/src/pieces/local-file.piece-resolver.spec.ts` around lines 24 -
31, The test case for resolving a piece from the local file system does not
assert that the initialize() method is invoked on the mock. Add an assertion
after the resolver.resolve() call to verify that mockLocalSync.initialize() was
called. This will prevent regressions where the initialize() call could be
accidentally removed from the resolver without the test failing.
🤖 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/modules/connections/core/use-cases/store-oauth-connection.use-case.ts`:
- Around line 7-13: The `@Inject`(ENCRYPTION_SERVICE) decorator on the crypto
parameter in the StoreOAuthConnectionUseCase constructor is ineffective and
should be removed because this class is instantiated via a factory in
connections.module.ts rather than through Nest's DI container, and the crypto
dependency is passed explicitly by the factory. Remove the `@Inject` decorator
from the crypto parameter in the constructor and update the imports to remove
the Inject decorator if it is no longer used elsewhere in the file.

In `@apps/api/src/modules/scheduler/connection-sync-runner.spec.ts`:
- Around line 50-66: The mocked returning() function generates trace IDs that
restart from trace-0 for each invocation, causing duplicate IDs across multiple
chunked insert calls. To fix this, introduce a counter variable outside the
returning mock that tracks the total number of traces generated across all
calls, then update the mockImplementation to use this counter and increment it
by lastInsertedCount after each returning() call, ensuring each simulated insert
produces unique and sequential trace IDs.

In `@apps/api/src/modules/scheduler/connection-sync-runner.ts`:
- Around line 611-615: The recordId computation in the connection-sync-runner
contains a non-deterministic fallback using Date.now() and Math.random(), which
breaks idempotency guarantees for the (data_source_id, ext_req_id) contract.
Replace the random fallback in the recordId assignment with a deterministic
value generated from the record payload itself (such as a hash of the payload
content) so that identical records always generate the same recordId across
retries and pagination, ensuring proper deduplication. This same fix applies to
both occurrences mentioned in the review (line 614 and line 627).

In `@apps/worker/src/adapters/outbound/nest-cache-invalidator.adapter.ts`:
- Around line 6-10: The invalidate method in NestCacheInvalidatorAdapter is a
complete no-op that does nothing in any environment. Replace the
Promise.resolve() placeholder with actual cache invalidation logic that works in
production. The method should use the _shardName parameter to invalidate the
specific shard's cache and trigger a reload or re-initialization of the affected
modules so that SyncGitopsShardUseCase can properly reload shard updates when
new commits are pulled. This must function in production environments, not just
during development, since the file watcher mechanism is restricted to
development mode only and the queue handler path needs a working invalidation
mechanism.

In `@infra/debezium/application.properties`:
- Line 25: The property debezium.sink.http.header.Authorization in the
application.properties file contains a fallback value of "secret" when the
DEBEZIUM_SECRET environment variable is unset, creating a security
vulnerability. Remove the colon and fallback value (the ":secret" portion) from
the property value so that it reads
debezium.sink.http.header.Authorization=Bearer ${DEBEZIUM_SECRET} without any
default, ensuring the configuration fails safely if the required environment
variable is not properly set.
- Around line 23-25: The configuration property
debezium.sink.http.header.Authorization uses the incorrect singular form
"header" when it should use the plural form "headers" to match Debezium's HTTP
sink configuration specification. Update the property key from
debezium.sink.http.header.Authorization to
debezium.sink.http.headers.Authorization to ensure Debezium properly recognizes
and sends the custom Authorization header with the shared secret value in CDC
relay calls.

In `@packages/pipeline/src/fanout/fanout-router.service.ts`:
- Around line 98-99: The schemaName parameter is being directly interpolated
into sql.raw() in the INSERT INTO statement without validation, creating a SQL
injection vulnerability. Add schema validation before the transaction in the
fanout router service using the same protection pattern applied in
connection-sync-runner.ts (reference Line 591), ensuring schemaName is validated
to be a safe identifier before it is used in sql.raw() to construct the table
reference in the INSERT statement.

In `@packages/pipeline/src/plugin-hooks/pipeline-hook-broker.service.ts`:
- Around line 35-39: The normalize method in PipelineHookBroker is throwing an
error when no normalizer is found for a given appName, which blocks the fallback
normalization flow in NormalizationService. Instead of throwing when
getNormalizer returns null, return the replica unchanged to allow
NormalizationService to proceed with its piece.normalize fallback path. Remove
the throw statement and ensure the replica is returned when no registered
normalizer exists.
- Around line 135-136: The reverseLookup hook implementation in the
pipeline-hook-broker.service.ts file is silently disabled by returning an empty
array when the hook is not implemented in the registry. This breaks the parent
trace requeue mechanism that NormalizationService depends on. Instead of
returning an empty array when the reverseLookup hook is not found, implement the
hook properly or change the logic to ensure that the dependency replay mechanism
is not silently disabled. The fix should ensure that NormalizationService
receives the correct data needed to requeue parent traces rather than an empty
result that breaks this critical functionality.

In
`@packages/provision/src/shared/adapters/outbound/registry-replication.adapter.ts`:
- Around line 216-227: The ALTER PUBLICATION ADD TABLE statement in the
registry-replication.adapter.ts file attempts to add all four tables
(inbound_outbox, replica_outbox, normalized_outbox, outbound_outbox) in a single
atomic operation with one exception handler. If any single table is already
published, the entire statement fails and none of the tables are added. To fix
this, wrap each individual table addition in its own separate BEGIN-EXCEPTION
block so that each table (inbound_outbox, replica_outbox, normalized_outbox,
outbound_outbox) can be independently attempted and independently catch the
duplicate_object exception. This allows tables that weren't already registered
to still be successfully added even if others already exist, preserving
idempotency across retries.
- Around line 213-224: The registerCdcTables method has two critical issues:
First, validate schemaName before concatenating it into sql.raw() calls (lines
221-224) to prevent SQL injection vulnerabilities—implement a whitelist check or
regex validation for valid PostgreSQL identifier patterns. Second, refactor the
exception handling to avoid idempotency issues where a duplicate table error on
any single ADD TABLE statement prevents the remaining tables from being
added—execute each ADD TABLE statement individually with its own exception
handler, or query the PostgreSQL information_schema to check which tables
already exist in the publication before attempting to add them. Additionally,
upgrade Drizzle ORM to version ≥0.45.2 in package.json to address the
CVE-2026-39356 vulnerability.

In `@plugins/salesforce/core/src/adapters/native-fetch.adapter.ts`:
- Around line 24-28: The setTimeout timer created in the AbortController
fallback path is never cleared if the request completes before the 10-second
timeout, causing timer accumulation in high-throughput scenarios. Either
refactor the code to use AbortSignal.timeout() API (available in Node 17.3+ and
modern browsers) when available instead of the setTimeout fallback, or modify
the implementation to return both the AbortSignal and the timer ID so that
callers can clear the timer after the fetch resolves. Choose the approach based
on your minimum Node.js version requirements.

In `@plugins/salesforce/core/src/use-cases/salesforce.use-cases.ts`:
- Around line 171-172: The debug console.log statements in the
SalesforceUseCases.poll method write directly to stdout and will appear in
production logs, making them unsuitable for production environments. Replace
these console.log calls with calls to a configured logger instance (if one is
available in the class or injected) that supports environment-specific logging
levels, or remove them entirely if they are only needed during development.
Ensure any retained logging uses the logger's appropriate method (such as debug
or info) rather than console.log.

In `@sdk/registry/src/pieces/piece-metadata.util.spec.ts`:
- Around line 32-50: The test should include assertions that verify the
extraction logic properly handles invalid auth metadata shapes. Add additional
expect statements after the existing ones to check that when auth.type is
invalid or auth.props is not an object, the result.authSchema is undefined or
properly rejected. This ensures the extractPieceMetadata function validates the
auth object structure according to its runtime-shape guards and prevents
regressions in auth extraction logic.

In `@sdk/registry/src/pieces/piece-metadata.util.ts`:
- Around line 94-95: The authType and authSchema assignments on lines 94-95
perform type casts (as string and as Record<string, unknown>) without validating
the actual shapes of the values. Add runtime validation checks to confirm that
piece.auth.type is actually a string type and piece.auth.props is actually an
object/Record type before assigning them to authType and authSchema
respectively. This ensures only valid shapes are accepted and prevents invalid
metadata from passing through the DTO contract.

In `@sdk/registry/src/pieces/sandbox.ts`:
- Around line 18-25: Remove the network-capable Web APIs from the sandbox
globals that are being exposed in the sandbox.ts file. Specifically, delete
fetch, AbortController, AbortSignal, Headers, Request, Response, FormData, and
Blob from the list of globals being added to the sandbox context. These APIs
provide direct outbound HTTP capability that bypasses the whitelist-based
isolation model, so they should not be exposed to evaluated plugin modules in
the sandbox environment.

In `@sdk/registry/src/pieces/workspace-plugin-registrar.ts`:
- Around line 24-30: The initialize() method in workspace-plugin-registrar.ts
stores the promise from _doInitialize() in this.initPromise but never clears it
after completion. If the first initialization attempt fails or exits early
(e.g., missing plugins directory), subsequent calls will return the same stale
failed promise instead of retrying. After awaiting this.initPromise on line 28,
reset this.initPromise to null or undefined to ensure that future calls can
retry initialization instead of returning a cached failed promise.

---

Outside diff comments:
In `@packages/pipeline/src/fanout/fanout-batch-processor.ts`:
- Around line 40-54: In the FanoutBatchProcessor constructor, replace the manual
instantiation of PipelineHookBrokerService in the constructor body with
dependency injection. Add PipelineHookBrokerService as a constructor parameter
(with an `@Inject` decorator if it has a provider token, similar to how the other
repository services are injected), and remove the line that creates a new
instance with this.broker = new PipelineHookBrokerService(). Instead, assign the
injected parameter directly to the private broker field, following the same
pattern used for queueService, storageResolver, and other dependencies.

In `@sdk/registry/src/pieces/local-file.piece-resolver.spec.ts`:
- Around line 24-31: The test case for resolving a piece from the local file
system does not assert that the initialize() method is invoked on the mock. Add
an assertion after the resolver.resolve() call to verify that
mockLocalSync.initialize() was called. This will prevent regressions where the
initialize() call could be accidentally removed from the resolver without the
test failing.
🪄 Autofix (Beta)

✅ Autofix completed


ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: cd7a9893-acc0-4085-b3e7-889504af8dc8

📥 Commits

Reviewing files that changed from the base of the PR and between db91d1c and 70da71a.

⛔ Files ignored due to path filters (1)
  • pnpm-lock.yaml is excluded by !**/pnpm-lock.yaml
📒 Files selected for processing (99)
  • apps/api/package.json
  • apps/api/src/app/app.module.ts
  • apps/api/src/config/env.validation.ts
  • apps/api/src/modules/connections/connections.module.ts
  • apps/api/src/modules/connections/connections/oauth.controller.ts
  • apps/api/src/modules/connections/core/use-cases/get-existing-oauth-credentials.use-case.ts
  • apps/api/src/modules/connections/core/use-cases/store-oauth-connection.use-case.spec.ts
  • apps/api/src/modules/connections/core/use-cases/store-oauth-connection.use-case.ts
  • apps/api/src/modules/pipeline/pipeline-cdc.listener.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/stitches/adapters/outbound/db-manager-schema-provisioner.adapter.ts
  • apps/api/src/modules/stitches/core/ports/outbound/schema-provisioner.port.ts
  • apps/api/src/modules/stitches/stitches.module.ts
  • apps/api/src/modules/trigger/adapters/outbound/pipeline-database-provisioner.adapter.ts
  • apps/api/src/modules/trigger/core/fakes/database-provisioner.fake.ts
  • apps/api/src/modules/trigger/core/ports/outbound/database-provisioner.port.ts
  • apps/api/src/modules/trigger/core/use-cases/enable-trigger.use-case.spec.ts
  • apps/api/src/modules/trigger/core/use-cases/enable-trigger.use-case.ts
  • apps/api/src/modules/trigger/trigger.module.ts
  • apps/worker/package.json
  • apps/worker/src/adapters/outbound/nest-cache-invalidator.adapter.ts
  • apps/worker/src/app.module.ts
  • apps/worker/src/config/env.validation.ts
  • apps/worker/src/core/fakes/outbox-repository.fake.ts
  • apps/worker/src/core/use-cases/app-installer/app-installer.spec.ts
  • apps/worker/src/core/use-cases/app-installer/install-piece.use-case.ts
  • apps/worker/src/core/use-cases/outbox/process-outbox.use-case.spec.ts
  • apps/worker/src/core/use-cases/outbox/process-outbox.use-case.ts
  • apps/worker/vitest.config.mts
  • infra/debezium/application.properties
  • package.json
  • packages/database/drizzle/tenant/0013_bright_molten_man.sql
  • packages/pipeline/src/delivery/delivery.service.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.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/index.ts
  • packages/pipeline/src/normalization/normalization.service.ts
  • packages/pipeline/src/pipeline-core.module.ts
  • packages/pipeline/src/plugin-hooks/application-shard.types.ts
  • packages/pipeline/src/plugin-hooks/pipeline-hook-broker.service.spec.ts
  • packages/pipeline/src/plugin-hooks/pipeline-hook-broker.service.ts
  • packages/pipeline/src/replication/replica.service.spec.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/pipeline-hook-broker.service.spec.ts
  • packages/pipeline/src/shared/adapters/outbound/drizzle-stitch.adapter.ts
  • packages/pipeline/src/shared/domain.ts
  • packages/provision/package.json
  • packages/provision/src/index.ts
  • packages/provision/src/provision.module.ts
  • packages/provision/src/provisioning/schema-provision.worker.ts
  • packages/provision/src/provisioning/use-cases/provision-schema.use-case.spec.ts
  • packages/provision/src/provisioning/use-cases/provision-schema.use-case.ts
  • packages/provision/src/registry/registry-replication.service.integration.spec.ts
  • packages/provision/src/registry/registry-replication.service.spec.ts
  • packages/provision/src/registry/registry-replication.service.ts
  • packages/provision/src/registry/registry-token-refresh.service.integration.spec.ts
  • packages/provision/src/registry/registry-token-refresh.service.spec.ts
  • packages/provision/src/registry/registry-token-refresh.service.ts
  • packages/provision/src/shared/adapters/outbound/registry-replication.adapter.integration.spec.ts
  • packages/provision/src/shared/adapters/outbound/registry-replication.adapter.ts
  • packages/provision/src/shared/ports/registry-replication.port.ts
  • packages/provision/src/test-utils/mock-db-manager.ts
  • packages/provision/tsconfig.build.json
  • packages/provision/tsconfig.json
  • packages/provision/vitest.config.ts
  • plugins/salesforce/core/package.json
  • plugins/salesforce/core/src/adapters/native-fetch.adapter.ts
  • plugins/salesforce/core/src/index.spec.ts
  • plugins/salesforce/core/src/use-cases/salesforce.use-cases.spec.ts
  • plugins/salesforce/core/src/use-cases/salesforce.use-cases.ts
  • plugins/salesforce/revenova-piece/package.json
  • plugins/salesforce/revenova-piece/src/index.ts
  • plugins/salesforce/revenova-piece/src/normalizeRevenovaToTms.ts
  • plugins/salesforce/revenova-piece/src/upsertRevenovaObject.ts
  • sdk/registry/src/index.ts
  • sdk/registry/src/metadata/metadata-discovery.service.ts
  • sdk/registry/src/pieces/drizzle-piece.repository.ts
  • sdk/registry/src/pieces/local-file.piece-resolver.spec.ts
  • sdk/registry/src/pieces/local-file.piece-resolver.ts
  • sdk/registry/src/pieces/piece-metadata.util.spec.ts
  • sdk/registry/src/pieces/piece-metadata.util.ts
  • sdk/registry/src/pieces/piece-repository.port.ts
  • sdk/registry/src/pieces/pieces.module.ts
  • sdk/registry/src/pieces/plugin-bootstrapper.port.ts
  • sdk/registry/src/pieces/plugin-metadata.bootstrapper.ts
  • sdk/registry/src/pieces/plugin-registry.bootstrapper.ts
  • sdk/registry/src/pieces/sandbox.ts
  • sdk/registry/src/pieces/workspace-plugin-registrar.spec.ts
  • sdk/registry/src/pieces/workspace-plugin-registrar.ts
💤 Files with no reviewable changes (16)
  • apps/api/src/modules/trigger/core/fakes/database-provisioner.fake.ts
  • packages/pipeline/src/sharding/pipeline-hook-broker.service.spec.ts
  • apps/api/src/modules/trigger/adapters/outbound/pipeline-database-provisioner.adapter.ts
  • apps/api/src/modules/stitches/core/ports/outbound/schema-provisioner.port.ts
  • packages/pipeline/src/sharding/application-shard-event.handler.spec.ts
  • packages/pipeline/src/sharding/application-loader.module.ts
  • apps/api/src/modules/trigger/trigger.module.ts
  • packages/pipeline/src/sharding/application-shard-event.handler.ts
  • apps/api/src/modules/trigger/core/ports/outbound/database-provisioner.port.ts
  • packages/pipeline/src/sharding/application-loader.service.ts
  • packages/provision/src/registry/registry-token-refresh.service.ts
  • apps/api/src/modules/stitches/stitches.module.ts
  • packages/pipeline/src/sharding/application-loader.service.spec.ts
  • packages/pipeline/src/fanout/fanout-batch-processor.spec.ts
  • apps/api/src/modules/pipeline/pipeline-cdc.listener.ts
  • apps/api/src/modules/stitches/adapters/outbound/db-manager-schema-provisioner.adapter.ts

Comment thread apps/api/src/modules/scheduler/connection-sync-runner.spec.ts Outdated
Comment thread apps/api/src/modules/scheduler/connection-sync-runner.ts Outdated
Comment on lines +6 to 10
async invalidate(_shardName: string): Promise<void> {
// Legacy application loader cache invalidation is no longer required.
// Plugin hot-loading handles this automatically in dev.
await Promise.resolve();
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

⚠️ Potential issue | 🔴 Critical

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
set -euo pipefail

# 1) Locate all cache invalidator implementations.
rg -n "implements\\s+CacheInvalidatorPort|class\\s+.*CacheInvalidator" apps/worker/src -C3

# 2) Trace all call sites of invalidate(...).
rg -n "\\.invalidate\\(" apps/worker/src -C3

# 3) Inspect GitOps sync path/logs around invalidation intent.
rg -n "invalidating module cache|pulled new commits|GitOps|sync-gitops" apps/worker/src -C3

Repository: pramodnarayana/nexiom

Length of output: 8438


🏁 Script executed:

# Get full context of the cache invalidator adapter
cat -n apps/worker/src/adapters/outbound/nest-cache-invalidator.adapter.ts

Repository: pramodnarayana/nexiom

Length of output: 591


🏁 Script executed:

# Check for git history to see what the previous implementation looked like
git log --oneline -p apps/worker/src/adapters/outbound/nest-cache-invalidator.adapter.ts | head -100

Repository: pramodnarayana/nexiom

Length of output: 887


🏁 Script executed:

# Search for plugin hot-loading or module reload mechanisms
rg -n "hot.?load|reload|dynamic.?require|require.cache|Module|import.meta.hot" apps/worker/src -i

Repository: pramodnarayana/nexiom

Length of output: 7761


🏁 Script executed:

# Check for environment-specific configurations or conditions in sync worker/use-case
rg -n "NODE_ENV|process.env|isDev|isProd|development|production" apps/worker/src/consumers/gitops-sync.worker.ts apps/worker/src/core/use-cases/gitops-sync/sync-gitops-shard.use-case.ts -C3

Repository: pramodnarayana/nexiom

Length of output: 1166


🏁 Script executed:

# Check what's in setupLocalFileWatcher and the actual hot-reload mechanism
cat -n apps/worker/src/consumers/gitops-sync.worker.ts | head -80

Repository: pramodnarayana/nexiom

Length of output: 3405


🏁 Script executed:

# Search for plugin reloading in production context
rg -n "system:plugins:reloaded|plugins:updated|invalidate|reload" apps/worker/src/core/use-cases/app-installer -A2 -B2

Repository: pramodnarayana/nexiom

Length of output: 3782


🏁 Script executed:

# Check the use-case that calls invalidate to see if there's prod-specific handling
cat -n apps/worker/src/core/use-cases/gitops-sync/sync-gitops-shard.use-case.ts

Repository: pramodnarayana/nexiom

Length of output: 2208


🏁 Script executed:

# Check if there's any other mechanism that reloads modules after git pull in production
rg -n "git pull|new commits|shard.*reload|plugin.*reload" apps/worker/src/core/use-cases/gitops-sync -C5

Repository: pramodnarayana/nexiom

Length of output: 2360


🏁 Script executed:

# Search for any event emission or reload trigger in the gitops sync path for production
rg -n "publishSystemEvent|pubSub|EventEmitter" apps/worker/src/consumers/gitops-sync.worker.ts apps/worker/src/core/use-cases/gitops-sync/

Repository: pramodnarayana/nexiom

Length of output: 47


🏁 Script executed:

# Check if there's any documentation or comments about production module reloading
rg -n "production|prod|reload.*prod|prod.*reload" apps/worker/src/adapters/outbound/nest-cache-invalidator.adapter.ts apps/worker/src/consumers/gitops-sync.worker.ts apps/worker/src/core/use-cases/gitops-sync/

Repository: pramodnarayana/nexiom

Length of output: 47


Production GitOps sync will silently fail to reload shard updates after new commits are pulled.

The invalidate() method is now a complete no-op in all environments. In production, when SyncGitopsShardUseCase executes (triggered by queue messages), it successfully pulls new commits but then calls invalidate(shardName) which does nothing, leaving the updated code unloaded. The comment mentions "Plugin hot-loading handles this automatically in dev," but the file watcher setup is restricted to NODE_ENV === "development" (gitops-sync.worker.ts:24), and there is no reload mechanism in the production queue handler path. Workers will require manual restart to pick up shard changes.

🤖 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/adapters/outbound/nest-cache-invalidator.adapter.ts` around
lines 6 - 10, The invalidate method in NestCacheInvalidatorAdapter is a complete
no-op that does nothing in any environment. Replace the Promise.resolve()
placeholder with actual cache invalidation logic that works in production. The
method should use the _shardName parameter to invalidate the specific shard's
cache and trigger a reload or re-initialization of the affected modules so that
SyncGitopsShardUseCase can properly reload shard updates when new commits are
pulled. This must function in production environments, not just during
development, since the file watcher mechanism is restricted to development mode
only and the queue handler path needs a working invalidation mechanism.

Comment thread infra/debezium/application.properties Outdated
Comment thread plugins/salesforce/core/src/use-cases/salesforce.use-cases.ts Outdated
Comment thread sdk/registry/src/pieces/piece-metadata.util.spec.ts
Comment thread sdk/registry/src/pieces/piece-metadata.util.ts Outdated
Comment thread sdk/registry/src/pieces/sandbox.ts Outdated
Comment thread sdk/registry/src/pieces/workspace-plugin-registrar.ts
@coderabbitai

coderabbitai Bot commented Jun 20, 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 17 unresolved review comments.

Files modified:

  • apps/api/src/modules/connections/core/use-cases/store-oauth-connection.use-case.ts
  • apps/api/src/modules/scheduler/connection-sync-runner.spec.ts
  • apps/api/src/modules/scheduler/connection-sync-runner.ts
  • infra/debezium/application.properties
  • packages/pipeline/src/fanout/fanout-router.service.ts
  • packages/pipeline/src/plugin-hooks/pipeline-hook-broker.service.ts
  • packages/provision/src/shared/adapters/outbound/registry-replication.adapter.ts
  • plugins/salesforce/core/src/use-cases/salesforce.use-cases.ts
  • sdk/registry/src/pieces/piece-metadata.util.spec.ts
  • sdk/registry/src/pieces/piece-metadata.util.ts
  • sdk/registry/src/pieces/sandbox.ts
  • sdk/registry/src/pieces/workspace-plugin-registrar.ts

Commit: baeb6ecf9523913c6d66abcc35ad43103d4f6104

The changes have been pushed to the initial-sync-testing branch.

Time taken: 8m 53s

Fixed 12 file(s) based on 17 unresolved review comments.

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

@coderabbitai coderabbitai Bot left a comment •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Actionable comments posted: 2

Caution

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

⚠️ Outside diff range comments (5)
apps/api/src/modules/scheduler/connection-sync-runner.ts (2)

322-324: ⚠️ Potential issue | 🟠 Major | ⚡ Quick win

Do not log raw pagination cursors.

nextPageCursor is connector-controlled and can contain opaque tokens, URLs, or customer identifiers. Log only whether another page exists, or sanitize the cursor before emitting it.

Proposed log change
-      this.logger.log(
-        `[ConnectionSyncRunner.runPollLoop] Polled page: ${page.records.length} records. nextPageCursor: ${JSON.stringify(page.nextPageCursor)}`,
-      );
+      this.logger.log(
+        {
+          event: 'connection_sync.polled_page',
+          records: page.records.length,
+          hasNextPage: Boolean(page.nextPageCursor),
+        },
+        '[ConnectionSyncRunner.runPollLoop] Polled page',
+      );
🤖 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/api/src/modules/scheduler/connection-sync-runner.ts` around lines 322 -
324, The logger.log call in the runPollLoop method of ConnectionSyncRunner is
logging the raw page.nextPageCursor, which may contain sensitive information
like opaque tokens, URLs, or customer identifiers. Replace the logging of
JSON.stringify(page.nextPageCursor) with a safer approach that logs only whether
a next page exists (for example, checking if page.nextPageCursor is truthy and
logging hasNextPage: true/false) or by sanitizing the cursor value before
logging to avoid exposing connector-controlled sensitive data.

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

Namespace extReqId by stream/object type.

The conflict target is (data_source_id, ext_req_id), but extReqId does not include objectType. Two streams from the same connection with the same record id and payload hash can suppress each other, dropping one gateway row and its outbox event.

Proposed deterministic key
-        const extReqId = `${recordId}-${payloadHash}`;
+        const extReqId = createHash('sha256')
+          .update(`${objectType}\0${recordId}\0${payloadHash}`)
+          .digest('hex');
🤖 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/api/src/modules/scheduler/connection-sync-runner.ts` at line 638, The
extReqId variable construction is missing the objectType namespace, which can
cause different streams with the same recordId and payloadHash but different
object types to conflict on the (data_source_id, ext_req_id) unique constraint.
Update the extReqId assignment to include objectType in the template string so
that it follows the pattern of including recordId, payloadHash, and objectType
separated by delimiters. This ensures each unique object type combination gets
its own unique extReqId and prevents rows from different streams from
incorrectly suppressing each other.
packages/pipeline/src/plugin-hooks/pipeline-hook-broker.service.ts (2)

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

Do not silently pass through unimplemented update preparation.

prepareUpdate now ignores appName, appProfile, destId, and destState, so destination updates can be sent without the connector-specific shaping that this hook exists to perform. Wire a registry-backed prepare hook or fail closed for connectors that require transformation.

🤖 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/plugin-hooks/pipeline-hook-broker.service.ts` around
lines 102 - 103, The prepareUpdate hook in the pipeline-hook-broker.service.ts
is silently passing through when not implemented in the registry, allowing
destination updates to be sent without required connector-specific
transformation. Instead of just logging and returning the payload unchanged,
either wire a registry-backed prepare hook that performs the necessary
transformation, or implement fail-closed behavior that prevents the update from
proceeding for connectors that require this hook's shaping. Ensure that appName,
appProfile, destId, and destState are properly handled through the actual hook
implementation rather than being ignored.

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

Restore dependency fetch behavior instead of skipping it.

activeFetch receives missingDependencies but now only logs and returns, so dependency backfills are silently skipped and downstream normalization/fanout can proceed with unresolved parent data. Reintroduce a registry-backed active-fetch hook or surface an explicit failure when dependencies cannot be fetched.

🤖 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/plugin-hooks/pipeline-hook-broker.service.ts` at line
125, The activeFetch hook in the pipeline-hook-broker.service.ts file currently
only logs a message when missingDependencies are detected and then returns
without performing any actual dependency fetching, causing silent skipping of
critical backfill operations. Replace the logging-only approach in the
hook.activeFetch event handler with either a proper implementation that uses a
registry-backed mechanism to actively fetch the missingDependencies, or
explicitly throw or return an error to signal failure when the hook cannot
resolve dependencies, ensuring downstream operations do not proceed with
unresolved parent data.
packages/pipeline/src/fanout/fanout-router.service.ts (1)

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

Verify lock ownership before processing the entity.

ON CONFLICT DO NOTHING can leave an existing active_sync_locks row owned by another trace, but the following SELECT ... FOR UPDATE does not check locked_by_trace_id or expires_at. After the short transaction commits, another trace can proceed as if it acquired the lock, causing concurrent fanout for the same entity.

Safer acquisition pattern
-        await tx.execute(sql`
+        const lockRows = await tx.execute(sql`
           INSERT INTO ${sql.raw('"' + schemaName + '"')}.active_sync_locks (data_source_id, entity_id, locked_by_trace_id, expires_at)
           VALUES (${dataSourceId}, ${entityId}, ${traceId}, NOW() + INTERVAL '5 minutes')
-          ON CONFLICT (data_source_id, entity_id) DO NOTHING
+          ON CONFLICT (data_source_id, entity_id) DO UPDATE
+            SET locked_by_trace_id = EXCLUDED.locked_by_trace_id,
+                expires_at = EXCLUDED.expires_at
+            WHERE active_sync_locks.expires_at < NOW()
+          RETURNING locked_by_trace_id
         `);
 
-        await tx.execute(sql`
-          SELECT 1 FROM ${sql.raw('"' + schemaName + '"')}.active_sync_locks
-          WHERE data_source_id = ${dataSourceId} AND entity_id = ${entityId}
-          FOR UPDATE
-        `);
+        if (!lockRows.length) {
+          return { kind: "superseded" as const };
+        }
🤖 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 98 - 108,
In the fanout-router.service.ts file, the SELECT ... FOR UPDATE query in the
lock acquisition block does not verify lock ownership. After the INSERT with ON
CONFLICT DO NOTHING, modify the SELECT statement's WHERE clause to include an
additional condition checking that locked_by_trace_id equals the current traceId
(in addition to the existing data_source_id and entity_id checks). This ensures
the current trace actually owns the lock before proceeding, preventing
concurrent processing when another trace owns an existing lock.
🤖 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
`@packages/provision/src/shared/adapters/outbound/registry-replication.adapter.ts`:
- Line 219: The tables array in registerCdcTables includes outbound_outbox, but
the Debezium configuration in infra/debezium/application.properties only filters
for inbound_outbox, replica_outbox, and normalized_outbox. Either remove
outbound_outbox from the tables array to match Debezium's current
table.include.list configuration, or update the Debezium configuration to
include outbound_outbox if it should be CDC-relayed. Ensure the tables
registered for CDC in the registry-replication.adapter.ts file are aligned with
what Debezium is configured to capture.

In `@sdk/registry/src/pieces/piece-metadata.util.ts`:
- Around line 96-99: The authSchema assignment in the piece metadata utility
does not validate that auth.type is valid before returning the props. Currently,
the code at line 96-99 checks if piece.auth exists and has props as an object,
but it should also verify that piece.auth.type is valid and non-empty. Modify
the condition to additionally check that piece.auth.type exists and is truthy,
ensuring that authSchema is only assigned when both a valid auth.type and valid
props are present.

---

Outside diff comments:
In `@apps/api/src/modules/scheduler/connection-sync-runner.ts`:
- Around line 322-324: The logger.log call in the runPollLoop method of
ConnectionSyncRunner is logging the raw page.nextPageCursor, which may contain
sensitive information like opaque tokens, URLs, or customer identifiers. Replace
the logging of JSON.stringify(page.nextPageCursor) with a safer approach that
logs only whether a next page exists (for example, checking if
page.nextPageCursor is truthy and logging hasNextPage: true/false) or by
sanitizing the cursor value before logging to avoid exposing
connector-controlled sensitive data.
- Line 638: The extReqId variable construction is missing the objectType
namespace, which can cause different streams with the same recordId and
payloadHash but different object types to conflict on the (data_source_id,
ext_req_id) unique constraint. Update the extReqId assignment to include
objectType in the template string so that it follows the pattern of including
recordId, payloadHash, and objectType separated by delimiters. This ensures each
unique object type combination gets its own unique extReqId and prevents rows
from different streams from incorrectly suppressing each other.

In `@packages/pipeline/src/fanout/fanout-router.service.ts`:
- Around line 98-108: In the fanout-router.service.ts file, the SELECT ... FOR
UPDATE query in the lock acquisition block does not verify lock ownership. After
the INSERT with ON CONFLICT DO NOTHING, modify the SELECT statement's WHERE
clause to include an additional condition checking that locked_by_trace_id
equals the current traceId (in addition to the existing data_source_id and
entity_id checks). This ensures the current trace actually owns the lock before
proceeding, preventing concurrent processing when another trace owns an existing
lock.

In `@packages/pipeline/src/plugin-hooks/pipeline-hook-broker.service.ts`:
- Around line 102-103: The prepareUpdate hook in the
pipeline-hook-broker.service.ts is silently passing through when not implemented
in the registry, allowing destination updates to be sent without required
connector-specific transformation. Instead of just logging and returning the
payload unchanged, either wire a registry-backed prepare hook that performs the
necessary transformation, or implement fail-closed behavior that prevents the
update from proceeding for connectors that require this hook's shaping. Ensure
that appName, appProfile, destId, and destState are properly handled through the
actual hook implementation rather than being ignored.
- Line 125: The activeFetch hook in the pipeline-hook-broker.service.ts file
currently only logs a message when missingDependencies are detected and then
returns without performing any actual dependency fetching, causing silent
skipping of critical backfill operations. Replace the logging-only approach in
the hook.activeFetch event handler with either a proper implementation that uses
a registry-backed mechanism to actively fetch the missingDependencies, or
explicitly throw or return an error to signal failure when the hook cannot
resolve dependencies, ensuring downstream operations do not proceed with
unresolved parent data.
🪄 Autofix (Beta)

✅ Autofix completed


ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: 6185e5df-b3f6-415e-8644-aa199abd357d

📥 Commits

Reviewing files that changed from the base of the PR and between 70da71a and baeb6ec.

📒 Files selected for processing (11)
  • apps/api/src/modules/connections/core/use-cases/store-oauth-connection.use-case.ts
  • apps/api/src/modules/scheduler/connection-sync-runner.spec.ts
  • apps/api/src/modules/scheduler/connection-sync-runner.ts
  • infra/debezium/application.properties
  • packages/pipeline/src/fanout/fanout-router.service.ts
  • packages/pipeline/src/plugin-hooks/pipeline-hook-broker.service.ts
  • packages/provision/src/shared/adapters/outbound/registry-replication.adapter.ts
  • plugins/salesforce/core/src/use-cases/salesforce.use-cases.ts
  • sdk/registry/src/pieces/piece-metadata.util.spec.ts
  • sdk/registry/src/pieces/piece-metadata.util.ts
  • sdk/registry/src/pieces/workspace-plugin-registrar.ts
💤 Files with no reviewable changes (1)
  • plugins/salesforce/core/src/use-cases/salesforce.use-cases.ts

Comment thread packages/provision/src/shared/adapters/outbound/registry-replication.adapter.ts Outdated
Comment thread sdk/registry/src/pieces/piece-metadata.util.ts Outdated
@coderabbitai

coderabbitai Bot commented Jun 20, 2026 •

Copy link
Copy Markdown
Contributor

Note

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

Fixes Applied Successfully

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

Files modified:

  • packages/provision/src/shared/adapters/outbound/registry-replication.adapter.ts
  • sdk/registry/src/pieces/piece-metadata.util.ts

Commit: 104c09c91ff28bad58bd1875ab65f069402e08df

The changes have been pushed to the initial-sync-testing branch.

Time taken: 2m 14s

coderabbitai Bot and others added 2 commits June 20, 2026 17:07
Fixed 2 file(s) based on 2 unresolved review comments.

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

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Caution

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

⚠️ Outside diff range comments (3)
apps/api/src/modules/scheduler/connection-sync-runner.ts (1)

611-634: 🧹 Nitpick | 🔵 Trivial | ⚡ Quick win

Payload hash is computed twice when recordId falls back to hash.

When recordId is empty and falls back to the payload hash (lines 615-622), the same hash is computed again for payloadHash (lines 628-634). Extract the hash computation once before determining recordId.

♻️ Proposed refactor to compute hash once
         const payload = record.data as Record<string, unknown>;

+        // Compute payload hash once for both deduplication and extReqId
+        const payloadHash = createHash('sha256')
+          .update(
+            typeof payload === 'object' && payload !== null
+              ? stringify(payload)
+              : String(payload),
+          )
+          .digest('hex');
+
         // Extract the primary identifier of the record (e.g. Salesforce Id)
         let recordId = cursorValue || this.extractRecordCursor(payload);

         // If no deterministic ID exists, generate one from payload hash
         if (!recordId) {
-          recordId = createHash('sha256')
-            .update(
-              typeof payload === 'object' && payload !== null
-                ? stringify(payload)
-                : String(payload),
-            )
-            .digest('hex')
-            .substring(0, 16);
+          recordId = payloadHash.substring(0, 16);
         }

-        // Hash the payload. This ensures that:
-        // 1. Identical polls (overlapping pages) have the exact same extReqId and are dropped as duplicates.
-        // 2. Updated records have the same recordId but a different hash, generating a new L1 trace.
-        const payloadHash = createHash('sha256')
-          .update(
-            typeof payload === 'object' && payload !== null
-              ? stringify(payload)
-              : String(payload),
-          )
-          .digest('hex');
-
         const extReqId = `${objectType}-${recordId}-${payloadHash}`;
🤖 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/api/src/modules/scheduler/connection-sync-runner.ts` around lines 611 -
634, The payload hash is being computed twice in this block - once as a fallback
for recordId and again for payloadHash. Extract the hash computation (using
createHash and stringify) before the recordId assignment, store it in a
variable, and then reuse that variable for both the recordId substring and
payloadHash assignments to eliminate the duplicate computation.
packages/pipeline/src/fanout/fanout-router.service.ts (1)

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

Lock ownership is not verified after ON CONFLICT DO NOTHING.

If another trace already holds (data_source_id, entity_id), the insert no-ops and the SELECT ... locked_by_trace_id = ${traceId} can return zero rows. The code still proceeds and later calls releaseSyncLock(...), which can release a lock it never acquired.

Also applies to: 200-203

🤖 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 98 - 109,
The lock acquisition logic in the fanout-router.service.ts file does not verify
lock ownership after the INSERT with ON CONFLICT DO NOTHING clause. When another
trace already holds the lock for the same (data_source_id, entity_id) pair, the
INSERT is a no-op and the subsequent SELECT query for locked_by_trace_id
verification returns zero rows. Check the result of the SELECT query to confirm
that the current trace owns the lock (i.e., the query returns at least one row)
before proceeding. If the SELECT returns no rows indicating another trace holds
the lock, the code should handle this conflict appropriately (e.g., throw an
error or take corrective action) instead of proceeding to call releaseSyncLock
for a lock that was never acquired.
packages/pipeline/src/plugin-hooks/pipeline-hook-broker.service.ts (1)

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

prepareUpdate now hard-fails the normal fanout path.

FanoutBatchProcessor.processSingleStitch(...) always invokes broker.prepareUpdate(...); this implementation always throws, so routes fail before outbox publish even when payload mapping succeeds.

Suggested fix
  async prepareUpdate(
    appName: string,
    appProfile: string,
    payload: Record<string, any>,
    destId?: string,
    destState?: Record<string, any>,
  ): Promise<Record<string, any>> {
-    throw new Error(`prepareUpdate hook is not implemented for ${appName}:${appProfile}. Cannot proceed with update without connector-specific transformation.`);
+    this.logger.debug(
+      { event: "hook.prepareUpdate", appName, appProfile },
+      "No prepareUpdate hook registered. Passing payload through.",
+    );
+    return payload;
  }
🤖 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/plugin-hooks/pipeline-hook-broker.service.ts` around
lines 95 - 103, The prepareUpdate method always throws an error, which causes
FanoutBatchProcessor.processSingleStitch(...) to fail before outbox publish even
when payload mapping succeeds. Instead of always throwing, modify the
prepareUpdate method to return the payload unchanged as a default behavior when
no connector-specific implementation is available, allowing the normal fanout
path to proceed. Only throw an error if the transformation is actually required
for that specific appName and appProfile combination.
🤖 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.

Outside diff comments:
In `@apps/api/src/modules/scheduler/connection-sync-runner.ts`:
- Around line 611-634: The payload hash is being computed twice in this block -
once as a fallback for recordId and again for payloadHash. Extract the hash
computation (using createHash and stringify) before the recordId assignment,
store it in a variable, and then reuse that variable for both the recordId
substring and payloadHash assignments to eliminate the duplicate computation.

In `@packages/pipeline/src/fanout/fanout-router.service.ts`:
- Around line 98-109: The lock acquisition logic in the fanout-router.service.ts
file does not verify lock ownership after the INSERT with ON CONFLICT DO NOTHING
clause. When another trace already holds the lock for the same (data_source_id,
entity_id) pair, the INSERT is a no-op and the subsequent SELECT query for
locked_by_trace_id verification returns zero rows. Check the result of the
SELECT query to confirm that the current trace owns the lock (i.e., the query
returns at least one row) before proceeding. If the SELECT returns no rows
indicating another trace holds the lock, the code should handle this conflict
appropriately (e.g., throw an error or take corrective action) instead of
proceeding to call releaseSyncLock for a lock that was never acquired.

In `@packages/pipeline/src/plugin-hooks/pipeline-hook-broker.service.ts`:
- Around line 95-103: The prepareUpdate method always throws an error, which
causes FanoutBatchProcessor.processSingleStitch(...) to fail before outbox
publish even when payload mapping succeeds. Instead of always throwing, modify
the prepareUpdate method to return the payload unchanged as a default behavior
when no connector-specific implementation is available, allowing the normal
fanout path to proceed. Only throw an error if the transformation is actually
required for that specific appName and appProfile combination.

ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: 5b55d095-33a1-4b46-a961-a5bc51074ec1

📥 Commits

Reviewing files that changed from the base of the PR and between baeb6ec and 6faae43.

📒 Files selected for processing (10)
  • apps/api/src/modules/connections/core/use-cases/store-oauth-connection.use-case.ts
  • apps/api/src/modules/scheduler/connection-sync-runner.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.ts
  • packages/pipeline/src/plugin-hooks/pipeline-hook-broker.service.spec.ts
  • packages/pipeline/src/plugin-hooks/pipeline-hook-broker.service.ts
  • packages/provision/src/shared/adapters/outbound/registry-replication.adapter.ts
  • sdk/registry/src/pieces/local-file.piece-resolver.spec.ts
  • sdk/registry/src/pieces/piece-metadata.util.ts

@pramodnarayana
pramodnarayana merged commit 1657877 into development Jun 20, 2026
2 checks passed
@coderabbitai coderabbitai Bot mentioned this pull request Jun 21, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant