jsonata transformation inbound gateway data sync - #131
Conversation
|
Note Reviews pausedIt 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 Use the following commands to manage reviews:
Use the checkboxes below for quick actions:
📝 WalkthroughWalkthroughReworks pipeline schemas and outbox flow (payload→request/response), adds inbound/replica outbox tables and migrations, removes API-side ReplicaService, adds DB provisioning/migration CLI, introduces app-specific synchronous webhook response registry (used by Revenova SOAP ACK), updates worker extractors/normalizers, and adjusts webhook ingestion and UI webhook UX. Changes
Sequence Diagram(s)sequenceDiagram
autonumber
participant Client as Client
participant Controller as WebhooksController
participant Registry as AppResponseRegistry
participant DB as Database
participant HTTP as HTTPResponse
Client->>Controller: POST /webhooks/{connectionId} (body + headers)
Controller->>Controller: Normalize body → request (capture raw)
Controller->>Registry: executeAppWebhookResponses(body, headers)
Registry->>Registry: Invoke handlers in insertion order
alt handler returns non-null
Registry-->>Controller: {status, contentType, body}
Controller->>DB: Insert inbound_gateway (request + response) + inbound_outbox
DB-->>Controller: Inserted
Controller->>HTTP: Send custom status/content-type/body
HTTP-->>Client: Custom response
else no handler matched
Controller->>DB: Insert inbound_gateway (request) + inbound_outbox
DB-->>Controller: Inserted
Controller->>HTTP: Send 202 Accepted
HTTP-->>Client: 202 Accepted
end
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes Possibly related PRs
Poem
🚥 Pre-merge checks | ✅ 3 | ❌ 2❌ Failed checks (1 warning, 1 inconclusive)
✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches🧪 Generate unit tests (beta)
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. Comment |
There was a problem hiding this comment.
Actionable comments posted: 11
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (4)
packages/pieces/platform/salesforce/src/index.ts (1)
206-229:⚠️ Potential issue | 🔴 CriticalDo not drop Salesforce webhook authentication without a replacement.
With no
webhookconfig on the Salesforce piece,WebhookSignatureGuardreturnstrueat line 44 and skips signature verification. The piece does not register an app-level webhook response handler, andexecuteAppWebhookResponsesonly shapes the ACK response body—it performs no authentication. Forged inbound events can reach L1 unless piece-level webhook verification is restored.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@packages/pieces/platform/salesforce/src/index.ts` around lines 206 - 229, The Salesforce piece currently created by createPiece (symbol: salesforce) has no webhook config so WebhookSignatureGuard short-circuits verification; add a webhook configuration to the piece that restores signature verification (or supply a webhook.verify handler) so inbound webhooks are validated before handlers run, and ensure the piece either registers an app-level webhook response handler or provides the same authentication logic used by WebhookSignatureGuard (reference: WebhookSignatureGuard, executeAppWebhookResponses, and salesforceUniversalTrigger) to validate signatures and reject forged events.apps/api/src/modules/webhooks/webhooks.controller.spec.ts (1)
86-293: 🧹 Nitpick | 🔵 TrivialConsider adding minimal req/res mocks instead of
{} as any.Current scenarios happen to avoid the
req.rawBodybranch (bodies are objects) and the custom-response branch (noNotificationkey triggersexecuteAppWebhookResponses), so{} as anyslips through. This leaves two important controller paths untested:
- Non-JSON payloads where
req.rawBodyis used to populatenormalizedPayload.- App-defined synchronous responses (Revenova SOAP ACK) where
res.status().set().send()is invoked andinboundGateway.responseis persisted.Adding a typed
reqmock withrawBody: Buffer.from(...)and aresmock withstatus/set/sendspies would cover these paths and prevent future regressions.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/api/src/modules/webhooks/webhooks.controller.spec.ts` around lines 86 - 293, Add focused req/res mocks to the webhooks.controller.spec.ts tests to cover the req.rawBody and synchronous app-response branches: create a typed req mock (with rawBody: Buffer.from(...) and headers) and a res mock implementing status(), set(), send() spies that return res for chaining, then pass these mocks to controller.ingest in one or two tests so the code path that uses req.rawBody for normalizedPayload and the branch that calls executeAppWebhookResponses / persists inboundGateway.response is exercised; update the relevant tests that currently call controller.ingest with "{} as any" to use the new mocks and assert res.status/set/send calls and that inboundGateway.response persistence behavior occurs.apps/api/src/modules/trigger/trigger-executor.service.ts (1)
430-478: 🧹 Nitpick | 🔵 TrivialOptional: rename local
row.payload→row.requestfor naming consistency.The DB column is now
request, and theWebhookRunParams.payload/TriggerContext.payloadfield retains its meaning at the caller layer — but withininsertGatewayRow, the row object'spayloadkey is now only used to map torequest. Renaming the local field reduces cognitive friction when reading this function.♻️ Proposed rename
private async insertGatewayRow( schemaName: string, row: { connectionId: string; objectType: string | undefined; - payload: unknown; + request: unknown; extReqId: string; }, ): Promise<boolean> { @@ - request: row.payload, + request: row.request,Update the call site in
executeAndIngestaccordingly (payload: record→request: record).🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/api/src/modules/trigger/trigger-executor.service.ts` around lines 430 - 478, Rename the local row field used in insertGatewayRow from payload to request so the object key matches the DB column and reduces confusion: update the parameter shape in insertGatewayRow (the row object) to use request: unknown instead of payload, replace usages of row.payload with row.request inside insertGatewayRow (including the .values mapping to request), and update the caller executeAndIngest (and any other callers) to pass request: record instead of payload: record so all call sites and the function signature stay consistent.apps/api/src/modules/webhooks/webhooks.controller.ts (1)
176-287:⚠️ Potential issue | 🟠 MajorTwo concerns around
executeAppWebhookResponses.1)
bodyargument may not reflect the raw request.
For content types NestJS doesn't natively parse (e.g.text/xmlfrom Salesforce outbound messaging),@Body()typically arrives as{}or similar, while the raw payload is only inreq.rawBody. Handlers such as Revenova's SOAP ACK inspecttypeof body === 'string'and will therefore never match. The method already constructsnormalizedPayloadwith correct fallback logic; pass that to the handler instead of rawbody.Suggested fix (apply to both call sites at lines 176 and 272)
- const customResponse = executeAppWebhookResponses(body, headers); + const customResponse = executeAppWebhookResponses(normalizedPayload, headers);2) Inconsistent response persistence between primary and idempotency paths.
The happy path persists{ status, contentType, body }ontoinbound_gateway.response(lines 187–202). The idempotency-collision path (lines 272–287) sends the same custom response to the caller but never updates the pre-existing row. Support/audit tooling will see a NULLresponsefor any retried delivery.Suggested fix
if (customResponse) { this.logger.debug( { event: 'l1.app_response', contentType: customResponse.contentType, status: customResponse.status, }, 'Returning app-defined synchronous response (idempotency path)', ); + if (existingTraceIdOutside) { + await this.db.transaction(async (tx) => { + assertValidSchemaName(schemaName); + await tx.execute( + sql`SET LOCAL search_path TO ${sql.raw('"' + schemaName + '"')}`, + ); + await tx + .update(inboundGateway) + .set({ + response: { + status: customResponse.status, + contentType: customResponse.contentType, + body: customResponse.body, + }, + }) + .where(sql`${inboundGateway.traceId} = ${existingTraceIdOutside}`); + }); + } res .status(customResponse.status) .set('Content-Type', customResponse.contentType) .send(customResponse.body); return; }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/api/src/modules/webhooks/webhooks.controller.ts` around lines 176 - 287, Call executeAppWebhookResponses with the normalized payload instead of the raw body at both call sites (replace passing body with the previously-constructed normalizedPayload variable) so handlers that expect the raw string (e.g., SOAP ACK) can match; and in the idempotency path, when a customResponse is returned, persist the response onto the existing inbound_gateway row (use the same transaction pattern as the primary path: assertValidSchemaName(schemaName), SET LOCAL search_path, then tx.update(inboundGateway).set({ response: { status: customResponse.status, contentType: customResponse.contentType, body: customResponse.body } }).where(inboundGateway.traceId = existingTraceIdOutside)) before sending the HTTP response. Ensure you reference executeAppWebhookResponses and inboundGateway when making these edits.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@apps/api/src/main.ts`:
- Line 3: Replace the namespace import with the idiomatic default import: change
the import from "import * as express from 'express'" to "import express from
'express'", and verify usages like express.text() and any direct calls express()
still work (update call sites if any rely on namespace behavior); ensure
TypeScript compiler options esModuleInterop/allowSyntheticDefaultImports are
enabled so the default import resolves correctly.
- Around line 42-43: The express.text() middleware mounted at '/webhooks'
currently only sets req.body and not req.rawBody, causing WebhookSignatureGuard
to throw a 403; update the app.use('/webhooks', express.text(...)) call to
include a verify callback in the options that assigns the raw buffer to
req.rawBody (e.g., verify: (req, _res, buf) => { req.rawBody = buf; }) so
signature verification can read the raw payload; keep the existing type and
limit options.
In `@apps/api/src/modules/trigger/trigger.module.ts`:
- Line 23: Remove the redundant empty controllers array from the TriggerModule
configuration: delete the controllers: [] entry in the TriggerModule (where
WebhooksController was previously registered and is now moved to WebhooksModule)
so the module metadata is cleaner and avoids an unnecessary empty property.
In `@apps/web/src/modules/connections/components/ActiveConnectionCard.tsx`:
- Around line 234-241: The onClick handler in ActiveConnectionCard currently
calls navigator.clipboard.writeText(webhookUrl) without awaiting or catching
errors and will always show a success toast; change the handler (the inline
onClick in ActiveConnectionCard.tsx) to first check for navigator.clipboard,
then perform an awaited writeText in a try/catch (or use .then/.catch) so
failures (permission denied, unavailable API) are caught, and show a failure
toast with a descriptive message in the catch branch while preserving
e.stopPropagation(); ensure the success toast is only shown after writeText
resolves.
- Around line 153-158: The webhookUrl derivation is fragile and inconsistent
with callbackUrl; stop reverse‑engineering VITE_API_URL with a regex and instead
use a canonical source: either add a backend-provided webhook base in the
ActiveConnectionResponse (e.g., response.webhookBaseUrl) and build the full URL
as `${webhookBaseUrl}/webhooks/${connection.id}`, or introduce and read an
explicit VITE_WEBHOOK_BASE_URL env var and construct webhookUrl from that
(falling back only to window.location.origin if neither exists). Update the code
paths that use webhookUrl (reference: webhookUrl, callbackUrl, connection.id,
ActiveConnectionResponse) to consume the new property/env var so external users
always get a correct, stable webhook URL.
In `@apps/web/src/modules/trace/pages/PipelineTracePage.tsx`:
- Around line 127-130: The conditional rendering currently uses
`!!details.inboundGateway?.response`, which treats empty objects like {} as
truthy and renders an empty JsonViewer; change the guard to an explicit
null/undefined check so only null-ish values are treated as absent. Update the
conditional around `JsonViewer` in PipelineTracePage (the expression using
`details.inboundGateway?.response`) to check `!= null` (or `!== undefined && !==
null`) instead of coercing to boolean, keeping the component `JsonViewer`'s
internal null/undefined behavior intact.
In `@package.json`:
- Line 7: The dev:light script is missing the workspace package
`@nexiom/application-revenova` which the API imports; either add
`@nexiom/application-revenova` to the parallel watch filters in the "dev:light"
script (so pnpm runs it alongside api, web, ai-engine, etc.) or add a pre-start
build step that runs a workspace build for `@nexiom/application-revenova` (e.g., a
pnpm/pnpm -w build or equivalent) before the existing parallel dev command;
update the "dev:light" npm script accordingly to include the chosen approach and
ensure the API can resolve dist/index.js from `@nexiom/application-revenova`.
In `@packages/dbmanager/src/impl/sql-database-manager.ts`:
- Around line 58-163: The GATEWAY_ACTIVE test's call-count and assertions are
outdated after provisionGatewayTables now emits 12 queries; update the test
referencing GATEWAY_ACTIVE to expect toHaveBeenCalledTimes(12) and add
assertions that the emitted SQL includes the inbound_outbox creation and its
associated DDL (look for inbound_outbox, idx_inbound_outbox_claim index and the
UNIQUE constraint DO-block); this change corresponds to the
provisionGatewayTables implementation which creates inbound_gateway and
inbound_outbox plus indexes and DO blocks, so assert those specific SQL snippets
appear in the recorded db.$client.query calls.
In `@packages/pieces/platform/framework/src/app-response.ts`:
- Around line 16-22: executeAppWebhookResponses currently iterates
appResponseRegistry and calls each handler directly, so a thrown error from any
handler will bubble up; modify executeAppWebhookResponses to wrap each call to
the handlers in a try/catch (around the invocation of fn(body, headers)) that
logs the error (include the error object and ideally an identifier for the
handler) and continues to the next handler, returning the first successful
non-null result or null if none succeed; additionally add a module-level helper
(e.g., resetAppResponseRegistry or unregisterAppResponse) to manage
appResponseRegistry for tests to avoid cross-test state leakage.
In `@TECHNICAL_DEBT.md`:
- Around line 200-216: Add a migration strategy and rollout plan to the
Capability URL proposal: update the docs and
apps/api/src/modules/webhooks/webhooks.controller.ts and
apps/api/src/guards/tenant-rate-limit.guard.ts to support dual-mode routing
(legacy POST /webhooks/:connectionId and new POST
/webhooks/:orgSlug/:connectionSlug/:secretToken) during a transition window,
implement transparent automatic redirects or deprecation warnings from the
legacy route to the new route, add a customer communication plan entry and an
estimated deprecation timeline (e.g., 90 days) in the recommended solution, and
ensure the migration steps include backfilling secret tokens for existing
app_connection records and a plan to revoke/rotate tokens safely.
- Around line 44-46: Update the documentation to reference the actual class
name: replace the incorrect `NormalizedOutboxService` with
`NormalizedOutboxWorker` in the cleanup task description; also ensure the
three-step checklist still mentions updating `trigger-executor.service.ts` (to
include `normalized_outbox` in the ALTER PUBLICATION script) and adding the
`else if (__table === 'normalized_outbox')` branch in `cdc-relay.controller.ts`
so the doc consistently names the real worker class (`NormalizedOutboxWorker`)
that should be deleted.
---
Outside diff comments:
In `@apps/api/src/modules/trigger/trigger-executor.service.ts`:
- Around line 430-478: Rename the local row field used in insertGatewayRow from
payload to request so the object key matches the DB column and reduces
confusion: update the parameter shape in insertGatewayRow (the row object) to
use request: unknown instead of payload, replace usages of row.payload with
row.request inside insertGatewayRow (including the .values mapping to request),
and update the caller executeAndIngest (and any other callers) to pass request:
record instead of payload: record so all call sites and the function signature
stay consistent.
In `@apps/api/src/modules/webhooks/webhooks.controller.spec.ts`:
- Around line 86-293: Add focused req/res mocks to the
webhooks.controller.spec.ts tests to cover the req.rawBody and synchronous
app-response branches: create a typed req mock (with rawBody: Buffer.from(...)
and headers) and a res mock implementing status(), set(), send() spies that
return res for chaining, then pass these mocks to controller.ingest in one or
two tests so the code path that uses req.rawBody for normalizedPayload and the
branch that calls executeAppWebhookResponses / persists inboundGateway.response
is exercised; update the relevant tests that currently call controller.ingest
with "{} as any" to use the new mocks and assert res.status/set/send calls and
that inboundGateway.response persistence behavior occurs.
In `@apps/api/src/modules/webhooks/webhooks.controller.ts`:
- Around line 176-287: Call executeAppWebhookResponses with the normalized
payload instead of the raw body at both call sites (replace passing body with
the previously-constructed normalizedPayload variable) so handlers that expect
the raw string (e.g., SOAP ACK) can match; and in the idempotency path, when a
customResponse is returned, persist the response onto the existing
inbound_gateway row (use the same transaction pattern as the primary path:
assertValidSchemaName(schemaName), SET LOCAL search_path, then
tx.update(inboundGateway).set({ response: { status: customResponse.status,
contentType: customResponse.contentType, body: customResponse.body }
}).where(inboundGateway.traceId = existingTraceIdOutside)) before sending the
HTTP response. Ensure you reference executeAppWebhookResponses and
inboundGateway when making these edits.
In `@packages/pieces/platform/salesforce/src/index.ts`:
- Around line 206-229: The Salesforce piece currently created by createPiece
(symbol: salesforce) has no webhook config so WebhookSignatureGuard
short-circuits verification; add a webhook configuration to the piece that
restores signature verification (or supply a webhook.verify handler) so inbound
webhooks are validated before handlers run, and ensure the piece either
registers an app-level webhook response handler or provides the same
authentication logic used by WebhookSignatureGuard (reference:
WebhookSignatureGuard, executeAppWebhookResponses, and
salesforceUniversalTrigger) to validate signatures and reject forged events.
🪄 Autofix (Beta)
✅ Autofix completed
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: beeb2d04-d1f6-4d7b-a785-b2004c540079
⛔ Files ignored due to path filters (1)
pnpm-lock.yamlis excluded by!**/pnpm-lock.yaml
📒 Files selected for processing (26)
TECHNICAL_DEBT.mdapps/api/package.jsonapps/api/src/db/database-manager.tsapps/api/src/db/db-cli.tsapps/api/src/main.tsapps/api/src/modules/pipeline/pipeline.module.tsapps/api/src/modules/pipeline/replica.service.spec.tsapps/api/src/modules/pipeline/replica.service.tsapps/api/src/modules/trace/trace.service.tsapps/api/src/modules/trigger/trigger-executor.service.tsapps/api/src/modules/trigger/trigger.module.tsapps/api/src/modules/webhooks/webhooks.controller.spec.tsapps/api/src/modules/webhooks/webhooks.controller.tsapps/web/src/modules/connections/components/ActiveConnectionCard.tsxapps/web/src/modules/trace/api/trace.api.tsapps/web/src/modules/trace/pages/PipelineTracePage.tsxapps/worker/src/modules/pipeline/replica.service.spec.tsapps/worker/src/modules/pipeline/replica.service.tspackage.jsonpackages/database/src/schema/pipeline.tspackages/dbmanager/src/impl/sql-database-manager.tspackages/pieces/application/revenova/package.jsonpackages/pieces/application/revenova/src/index.tspackages/pieces/platform/framework/src/app-response.tspackages/pieces/platform/framework/src/index.tspackages/pieces/platform/salesforce/src/index.ts
💤 Files with no reviewable changes (2)
- apps/api/src/modules/pipeline/replica.service.spec.ts
- apps/api/src/modules/pipeline/replica.service.ts
| private async provisionGatewayTables(schemaName: string): Promise<void> { | ||
| // We execute these independently so that they are idempotent. | ||
| // If they error, the whole task fails, preventing partial corruption. | ||
| // ── LAYER 1 — INBOUND GATEWAY ──────────────────────────────────────── | ||
| // Must match pipeline.ts buildTenantSchema > inboundGateway exactly. | ||
| // Uses TEXT + CHECK instead of public.pipeline_status_enum so that | ||
| // tenant schemas have no cross-schema type dependency. | ||
| await this.db.$client.query(` | ||
| CREATE TABLE IF NOT EXISTS "${schemaName}".inbound_gateway ( | ||
| id UUID PRIMARY KEY DEFAULT gen_random_uuid(), | ||
| trace_id UUID NOT NULL UNIQUE, | ||
| connection_id UUID NOT NULL, | ||
| object_type VARCHAR(100), | ||
| request JSONB NOT NULL, | ||
| response JSONB, | ||
| headers JSONB, | ||
| ext_req_id VARCHAR(255), | ||
| status TEXT NOT NULL DEFAULT 'RECEIVED' | ||
| CHECK (status IN ('RECEIVED','PROCESSING','REPLICATED', | ||
| 'NORMALIZED','SKIPPED','PENDING', | ||
| 'SUCCESS','FAIL','RETRY','DISMISSED')), | ||
| created_at TIMESTAMPTZ NOT NULL DEFAULT NOW() | ||
| ); | ||
| `); | ||
|
|
||
| // Idempotently rename payload → request for pre-existing schemas. | ||
| await this.db.$client.query(` | ||
| DO $$ BEGIN | ||
| IF EXISTS ( | ||
| SELECT 1 FROM information_schema.columns | ||
| WHERE table_schema = '${schemaName}' | ||
| AND table_name = 'inbound_gateway' | ||
| AND column_name = 'payload' | ||
| ) THEN | ||
| ALTER TABLE "${schemaName}".inbound_gateway RENAME COLUMN payload TO request; | ||
| END IF; | ||
| END $$; | ||
| `); | ||
|
|
||
| // Idempotently add response column for pre-existing schemas. | ||
| await this.db.$client.query(` | ||
| CREATE TABLE IF NOT EXISTS "${schemaName}".inbound_gateway ( | ||
| id UUID PRIMARY KEY DEFAULT gen_random_uuid(), | ||
| source_event_id TEXT NOT NULL, | ||
| trigger_name TEXT NOT NULL, | ||
| app_name TEXT NOT NULL, | ||
| object_type TEXT, | ||
| payload JSONB NOT NULL, | ||
| status TEXT NOT NULL DEFAULT 'pending', | ||
| trace_id UUID NOT NULL DEFAULT gen_random_uuid(), | ||
| created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), | ||
|
|
||
| CONSTRAINT uq_inbound_source_event UNIQUE (source_event_id) | ||
| ); | ||
| `); | ||
| ALTER TABLE "${schemaName}".inbound_gateway | ||
| ADD COLUMN IF NOT EXISTS response JSONB; | ||
| `); | ||
|
|
||
| // Idempotency: unique (connection_id, ext_req_id) prevents duplicate | ||
| // vendor events from being ingested twice. | ||
| await this.db.$client.query(` | ||
| CREATE INDEX IF NOT EXISTS idx_inbound_gateway_status | ||
| ON "${schemaName}".inbound_gateway (status); | ||
| `); | ||
| CREATE UNIQUE INDEX IF NOT EXISTS idx_l1_ext_id | ||
| ON "${schemaName}".inbound_gateway (connection_id, ext_req_id) | ||
| WHERE ext_req_id IS NOT NULL; | ||
| `); | ||
|
|
||
| await this.db.$client.query(` | ||
| CREATE INDEX IF NOT EXISTS idx_inbound_gateway_object_type | ||
| ON "${schemaName}".inbound_gateway (object_type) | ||
| WHERE object_type IS NOT NULL; | ||
| `); | ||
| CREATE INDEX IF NOT EXISTS idx_l1_object_type | ||
| ON "${schemaName}".inbound_gateway (object_type) | ||
| WHERE object_type IS NOT NULL; | ||
| `); | ||
|
|
||
| await this.db.$client.query(` | ||
| CREATE INDEX IF NOT EXISTS idx_inbound_gateway_created_at | ||
| ON "${schemaName}".inbound_gateway (created_at DESC); | ||
| `); | ||
| CREATE INDEX IF NOT EXISTS idx_l1_status | ||
| ON "${schemaName}".inbound_gateway (status); | ||
| `); | ||
|
|
||
| // Drop old payload GIN index if present, create new request GIN index. | ||
| await this.db.$client.query(` | ||
| DROP INDEX IF EXISTS "${schemaName}".idx_l1_payload_gin; | ||
| `); | ||
|
|
||
| await this.db.$client.query(` | ||
| CREATE INDEX IF NOT EXISTS idx_l1_request_gin | ||
| ON "${schemaName}".inbound_gateway USING gin (request); | ||
| `); | ||
|
|
||
| // ── INBOUND OUTBOX — L1 → L2 transactional outbox ─────────────────── | ||
| // Must match pipeline.ts buildTenantSchema > inboundOutbox exactly. | ||
| // Written atomically with inbound_gateway in the same transaction so | ||
| // a process crash between DB commit and SQS publish cannot lose events. | ||
| await this.db.$client.query(` | ||
| CREATE TABLE IF NOT EXISTS "${schemaName}".inbound_outbox ( | ||
| id UUID PRIMARY KEY DEFAULT gen_random_uuid(), | ||
| trace_id UUID NOT NULL, | ||
| connection_id UUID NOT NULL, | ||
| schema_name VARCHAR(128) NOT NULL DEFAULT current_schema(), | ||
| status TEXT NOT NULL DEFAULT 'PENDING' | ||
| CHECK (status IN ('PENDING','PROCESSING','SUCCESS','FAIL','RETRY')), | ||
| attempts INTEGER NOT NULL DEFAULT 0, | ||
| last_error VARCHAR(500), | ||
| next_retry_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), | ||
| created_at TIMESTAMPTZ NOT NULL DEFAULT NOW() | ||
| ); | ||
| `); | ||
|
|
||
| await this.db.$client.query(` | ||
| CREATE INDEX IF NOT EXISTS idx_inbound_outbox_claim | ||
| ON "${schemaName}".inbound_outbox (status, next_retry_at ASC) | ||
| WHERE status IN ('PENDING', 'PROCESSING', 'RETRY'); | ||
| `); | ||
|
|
||
| await this.db.$client.query(` | ||
| DO $$ BEGIN | ||
| ALTER TABLE "${schemaName}".inbound_outbox | ||
| ADD CONSTRAINT idx_inbound_outbox_trace UNIQUE (trace_id, connection_id); | ||
| EXCEPTION WHEN duplicate_table THEN NULL; | ||
| WHEN duplicate_object THEN NULL; | ||
| END $$; | ||
| `); | ||
| } |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
# Inspect the spec's GATEWAY_ACTIVE expectations
fd 'sql-database-manager\.spec' --type f --exec cat {}Repository: pramodnarayana/nexiom
Length of output: 4273
Update the spec assertion for GATEWAY_ACTIVE to reflect the expanded DDL and verify outbox table provisioning.
The GATEWAY_ACTIVE test currently asserts toHaveBeenCalledTimes(5), expecting 1 CREATE SCHEMA + 4 gateway statements. The implementation now issues 12 total queries: 1 CREATE SCHEMA + 11 statements from provisionGatewayTables (inbound_gateway table + rename payload→request DO-block + ADD COLUMN response + 3 indexes + DROP old GIN + new GIN index + inbound_outbox table + claim index + UNIQUE constraint DO-block). Update the assertion to toHaveBeenCalledTimes(12) and add assertions to verify that inbound_outbox and its associated DDL (claim index, UNIQUE constraint) are present in the query calls.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@packages/dbmanager/src/impl/sql-database-manager.ts` around lines 58 - 163,
The GATEWAY_ACTIVE test's call-count and assertions are outdated after
provisionGatewayTables now emits 12 queries; update the test referencing
GATEWAY_ACTIVE to expect toHaveBeenCalledTimes(12) and add assertions that the
emitted SQL includes the inbound_outbox creation and its associated DDL (look
for inbound_outbox, idx_inbound_outbox_claim index and the UNIQUE constraint
DO-block); this change corresponds to the provisionGatewayTables implementation
which creates inbound_gateway and inbound_outbox plus indexes and DO blocks, so
assert those specific SQL snippets appear in the recorded db.$client.query
calls.
|
Note Autofix is a beta feature. Expect some limitations and changes as we gather feedback and continue to improve it. Fixes Applied SuccessfullyFixed 8 file(s) based on 11 unresolved review comments. Files modified:
Commit: The changes have been pushed to the Time taken: |
Fixed 8 file(s) based on 11 unresolved review comments. Co-authored-by: CodeRabbit <noreply@coderabbit.ai>
There was a problem hiding this comment.
Actionable comments posted: 3
♻️ Duplicate comments (2)
TECHNICAL_DEBT.md (1)
40-40:⚠️ Potential issue | 🟡 MinorFix class name inconsistency.
Line 40 references
NormalizedOutboxService, but the actual class isNormalizedOutboxWorker(as correctly stated on line 46). This inconsistency could confuse developers implementing the cleanup task.📝 Proposed fix
-- As a result, the L3->L4 handoff still relies on a legacy cron-polling service (`NormalizedOutboxService`), which wastes database CPU and prevents true real-time elasticity. +- As a result, the L3->L4 handoff still relies on a legacy cron-polling worker (`NormalizedOutboxWorker`), which wastes database CPU and prevents true real-time elasticity.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@TECHNICAL_DEBT.md` at line 40, The doc references the wrong class name: replace the incorrect `NormalizedOutboxService` with the actual class `NormalizedOutboxWorker` in the sentence on line 40 so the README consistently references `NormalizedOutboxWorker` (matching the usage on line 46) to avoid developer confusion when locating the cleanup task implementation.package.json (1)
7-7:⚠️ Potential issue | 🟠 MajorPre-build Revenova before starting the parallel dev stack.
Adding
@nexiom/application-revenovato the parallel filters still leaves a clean-checkout race: the API can start before Revenova’sdevtask emits thedistentrypoint it imports. Add a small prebuild before the parallel command.Proposed fix
- "dev:light": "QUEUE_ENABLED=false dotenv -e .env -- pnpm --parallel --filter=api --filter=web --filter=@nexiom/ai-engine --filter=@nexiom/piece-quickbooks --filter=@nexiom/piece-salesforce --filter=@nexiom/piece-registry --filter=@nexiom/application-revenova dev", + "dev:light": "pnpm --filter=@nexiom/application-revenova build && QUEUE_ENABLED=false dotenv -e .env -- pnpm --parallel --filter=api --filter=web --filter=@nexiom/ai-engine --filter=@nexiom/piece-quickbooks --filter=@nexiom/piece-salesforce --filter=@nexiom/piece-registry --filter=@nexiom/application-revenova dev",Verify whether the package resolves through built output and whether
dev:lightguarantees that output first:#!/bin/bash # Description: Confirm Revenova package entrypoints/scripts and the root dev:light startup order. # Expected: if Revenova exports dist/*, dev:light should build it before starting api in parallel. jq -r '.scripts["dev:light"]' package.json jq -r '.main, .exports["."].import, .scripts.build, .scripts.dev' packages/pieces/application/revenova/package.json jq -r '.dependencies["@nexiom/application-revenova"]' apps/api/package.json🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@package.json` at line 7, The dev:light script can race because `@nexiom/application-revenova`’s built output may not exist when the parallel dev tasks start; update the "dev:light" script so it first runs a prebuild step for `@nexiom/application-revenova` (e.g., run its build or dev that produces dist) and only then launches the existing parallel pnpm --parallel --filter=... dev command; ensure you reference the "dev:light" npm script and the package name "@nexiom/application-revenova" when making the change so the API can resolve the built entrypoint before parallel startup.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@apps/api/src/main.ts`:
- Around line 42-49: The fallback text parser registered in the express
middleware for '/webhooks' (the express.text(...) call) is currently using type:
'*/*' which causes application/json payloads to be parsed as plain text; change
the type to a function that returns true only for non-JSON content types (i.e.,
inspect req.headers['content-type'] and return false for 'application/json' and
any '+json' vendor types like 'application/*+json') so JSON requests are left
for NestJS's JSON parser and only non-JSON/plain XML or other raw payloads are
parsed and assigned to (req as any).rawBody.
In `@apps/web/src/modules/connections/components/ActiveConnectionCard.tsx`:
- Around line 153-160: The webhookUrl useMemo can produce a double-slash when
VITE_WEBHOOK_BASE_URL or the derived VITE_API_URL ends with a trailing slash;
update the useMemo that computes webhookUrl so it normalizes webhookBase by
stripping any trailing slashes (e.g. replace trailing / characters) before
appending `/webhooks/${connection.id}`, ensuring you still fall back to
window.location.origin when needed; reference the webhookUrl useMemo, the
webhookBase variable, import.meta.env.VITE_WEBHOOK_BASE_URL,
import.meta.env.VITE_API_URL, and connection.id when making the change.
In `@TECHNICAL_DEBT.md`:
- Around line 224-227: Update the migration notes to explicitly state storage
and validation security: clarify that the new webhook_secret column on
app_connection is stored in plaintext (since these are capability URLs and need
direct comparison) and add a rate-limiting requirement for failed token
validations (e.g., max 10 failed attempts per orgSlug/connectionSlug per minute
with exponential backoff or temporary lockout) plus mandatory logging of failed
attempts (IP and timestamp); reference the webhook_secret column name and the
orgSlug/connectionSlug keying approach so readers know where to apply the
protections.
---
Duplicate comments:
In `@package.json`:
- Line 7: The dev:light script can race because `@nexiom/application-revenova`’s
built output may not exist when the parallel dev tasks start; update the
"dev:light" script so it first runs a prebuild step for
`@nexiom/application-revenova` (e.g., run its build or dev that produces dist) and
only then launches the existing parallel pnpm --parallel --filter=... dev
command; ensure you reference the "dev:light" npm script and the package name
"@nexiom/application-revenova" when making the change so the API can resolve the
built entrypoint before parallel startup.
In `@TECHNICAL_DEBT.md`:
- Line 40: The doc references the wrong class name: replace the incorrect
`NormalizedOutboxService` with the actual class `NormalizedOutboxWorker` in the
sentence on line 40 so the README consistently references
`NormalizedOutboxWorker` (matching the usage on line 46) to avoid developer
confusion when locating the cleanup task implementation.
🪄 Autofix (Beta)
✅ Autofix completed
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: 5e727c13-668c-4161-90fa-de2af24591eb
📒 Files selected for processing (8)
TECHNICAL_DEBT.mdapps/api/src/main.tsapps/api/src/modules/trigger/trigger.module.tsapps/web/src/modules/connections/components/ActiveConnectionCard.tsxapps/web/src/modules/trace/pages/PipelineTracePage.tsxpackage.jsonpackages/dbmanager/src/impl/sql-database-manager.spec.tspackages/pieces/platform/framework/src/app-response.ts
|
Note Autofix is a beta feature. Expect some limitations and changes as we gather feedback and continue to improve it. Fixes Applied SuccessfullyFixed 3 file(s) based on 3 unresolved review comments. Files modified:
Commit: The changes have been pushed to the Time taken: |
Fixed 3 file(s) based on 3 unresolved review comments. Co-authored-by: CodeRabbit <noreply@coderabbit.ai>
There was a problem hiding this comment.
Actionable comments posted: 13
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (5)
packages/database/src/schema/pipeline.ts (2)
104-113:⚠️ Potential issue | 🟡 MinorStale docstring on
replica_entity.The block comment still references
(entityType, sourceId)uniqueness andsrcReqTraceId, both of which were removed. Please update to reflect the new(connectionId, entityType, entityId)uniqueness and the fact that onlytraceIdnow links back to L1.Proposed edit
/** * LAYER 2 — UNIVERSAL REPLICA * * Parsed, structured state store. Batches from L1 are split into - * individual entity rows here. Unique on (entityType, sourceId) - * for deduplication — re-processing the same record is an upsert. + * individual entity rows here. Unique on (connectionId, entityType, entityId) + * for deduplication — re-processing the same record is an upsert. * - * `srcReqTraceId` links back to the L1 transmission that produced it. + * `traceId` links back to the L1 transmission that produced it. * `version` increments on every update for optimistic concurrency. */🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@packages/database/src/schema/pipeline.ts` around lines 104 - 113, Update the stale block comment for replica_entity to reflect the new schema: replace references to uniqueness on (entityType, sourceId) with uniqueness on (connectionId, entityType, entityId), remove mention of srcReqTraceId and state that traceId now links back to L1, and keep the note about version incrementing for optimistic concurrency; locate the docstring above the replica_entity definition and edit the text to mention "connectionId, entityType, entityId" uniqueness and "traceId" as the L1 link while preserving the overall purpose description.
84-102:⚠️ Potential issue | 🟡 MinorMigration strategy is correctly implemented, but update one stale comment.
The migration properly uses
ALTER TABLE ... RENAME COLUMN payload TO request(not DROP+ADD) and correctly handles index renaming (DROP idx_l1_payload_gin→CREATE idx_l1_request_gin). All code paths have been updated:replica.service.tsandfanout.service.tsboth reference therequestcolumn, and no active code references the oldpayloadcolumn name. The idempotency guards (IF EXISTS) are in place for pre-existing schemas.One minor fix: the comment at
apps/worker/src/modules/pipeline/replica.service.ts:112says "Read inbound_gateway payload" but the code correctly readsinbound.request. Update this to reflect the renamed column.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@packages/database/src/schema/pipeline.ts` around lines 84 - 102, Update the stale inline comment that still says "Read inbound_gateway payload" to reflect the renamed column by changing it to "Read inbound_gateway request" in replica.service.ts where the code accesses inbound.request (the comment adjacent to the code that reads inbound.request should be the only change).apps/worker/src/app.module.ts (1)
9-22:⚠️ Potential issue | 🔴 CriticalAiWorkerModule removal breaks the AI copilot queue consumer.
The
AiCopilotQueueproducer (apps/api/src/modules/ai/controllers/ai.controller.ts) still sends messages, but the only consumer (CopilotWorker, which requiresAiWorkerModule) is no longer instantiated. Messages will accumulate without being processed. Either re-enableAiWorkerModulein the worker imports, stop producing to the queue from the API, or remove the orphaned AI module files from the worker.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/worker/src/app.module.ts` around lines 9 - 22, The worker no longer imports AiWorkerModule so the CopilotWorker consumer that listens to AiCopilotQueue never runs; re-enable AiWorkerModule by adding AiWorkerModule back into the imports array in app.module.ts (so CopilotWorker is instantiated and consumes AiCopilotQueue), or if you intend to stop processing these jobs instead remove or disable the API producer (AiCopilotQueue usage in apps/api/src/modules/ai/controllers/ai.controller.ts) and delete the orphaned AiWorkerModule/CopilotWorker files; update whichever path you choose consistently (either restore AiWorkerModule import or remove queue production and module/class files) so producers and consumers match.apps/worker/src/modules/pipeline/fanout.service.ts (1)
126-137:⚠️ Potential issue | 🟡 MinorUpdate the remaining
sourceIdwording around this GEM lookup.The implementation now uses
replicaEntity.entityId, but the surrounding comments still describesourceId. Keeping those comments in sync avoids confusion in the L4/L5 GEM threading path.Suggested cleanup
- // ── Read normalized data + sourceId (for GEM) in one transaction ────── - // sourceId comes from replica_entity and is needed by DeliveryService (L5/L6) + // ── Read normalized data + entityId (for GEM) in one transaction ────── + // entityId comes from replica_entity and is needed by DeliveryService (L5/L6)🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/worker/src/modules/pipeline/fanout.service.ts` around lines 126 - 137, Update the outdated comment wording that still references "sourceId" around the GEM lookup using replicaEntity.entityId: change comments near the replicaEntity/select block (the lines that fetch replicaRows by traceId and assign srcVendorId) to refer to entityId or srcVendorId/vendorEntityId as appropriate so they match the implementation that uses replicaEntity.entityId and sets srcVendorId; verify any nearby variable names like replicaRows, traceId, and srcVendorId are described consistently (e.g., "Fetch entityId from replicaEntity for GEM threading" instead of "sourceId").apps/api/src/modules/webhooks/webhooks.controller.ts (1)
187-298: 🧹 Nitpick | 🔵 TrivialExtract the app-response send logic into a helper.
The happy-path (187-219) and idempotency-path (283-298) blocks duplicate the same
executeAppWebhookResponses→ log →res.status().set().send()sequence, differing only in log message and whether the DB is updated. A small private helper would keep the two call sites in sync and make the persist vs. no-persist decision explicit at a single site.♻️ Sketch
private sendAppResponse( res: Response, customResponse: { body: string; contentType: string; status: number }, context: 'normal' | 'idempotency', ) { this.logger.debug( { event: 'l1.app_response', contentType: customResponse.contentType, status: customResponse.status, path: context, }, `Returning app-defined synchronous response (${context} path)`, ); res .status(customResponse.status) .set('Content-Type', customResponse.contentType) .send(customResponse.body); }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/api/src/modules/webhooks/webhooks.controller.ts` around lines 187 - 298, Extract the duplicated response-send logic into a private helper (e.g., sendAppResponse) and call it from both the happy-path and idempotency-path; the helper should accept the Express Response, the customResponse object returned by executeAppWebhookResponses, and a context flag ('normal' | 'idempotency'), perform the unified this.logger.debug call (including contentType, status and path/context) and then call res.status(...).set('Content-Type', ...).send(...); replace the inline logging + res.status(...).set(...).send(...) at the two call sites (where executeAppWebhookResponses is used) with calls to sendAppResponse, keeping the DB persist logic only in the original happy-path before invoking the helper.
♻️ Duplicate comments (1)
TECHNICAL_DEBT.md (1)
40-40:⚠️ Potential issue | 🟡 MinorUse the correct worker class name in this cleanup task.
Line 40 references
NormalizedOutboxService, but the actual class in code isNormalizedOutboxWorker(apps/worker/src/modules/pipeline/normalized-outbox.worker.ts). Please align the name to avoid implementation confusion.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@TECHNICAL_DEBT.md` at line 40, Update the cleanup task text to use the actual worker class name: replace the incorrect reference to NormalizedOutboxService with NormalizedOutboxWorker so documentation matches the implementation (the class is defined as NormalizedOutboxWorker in the codebase).
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@apps/api/src/db/database-manager.ts`:
- Around line 821-869: Extract the repeated dynamic-import + client + drizzle +
SqlDatabaseManager setup into a private helper like withSchemaMgr<T>(cb: (mgr:
SqlDatabaseManager) => Promise<T>) that (1) dynamically imports
'@nexiom/dbmanager' and 'drizzle-orm/node-postgres', imports ./schema.js, calls
this.getPgClient(), constructs const db = drizzle(client, { schema: dbSchema }),
builds new SqlDatabaseManager(db as unknown as
import('@nexiom/database').DrizzleDb), and ensures client.end() in a finally
block; then refactor provisionGateway, provisionOutbound, migrateAllSchemas and
provisionLocal to call withSchemaMgr and invoke schemaMgr.applyPlan(schemaName,
SchemaPlan.<...>) or other per-method logic, removing the duplicated 5-line
boilerplate and keeping existing logging and SchemaPlan usages.
- Around line 876-919: migrateAllSchemas currently swallows per-schema failures
and calls SqlDatabaseManager.applyPlan(schema_name, SchemaPlan.REPLICA_ACTIVE)
which also re-applies gateway DDL; fix by switching to the new method
schemaMgr.migrateReplicaTables(schema_name) (to avoid gateway DDL side-effects)
and accumulate a failure counter or list during the for loop, then after the
loop log a summary and throw an error (or return non-zero) if any schema failed
so callers can detect partial failures; update the success message accordingly
and reference migrateAllSchemas, SqlDatabaseManager.applyPlan,
SchemaPlan.REPLICA_ACTIVE, and SqlDatabaseManager.migrateReplicaTables to find
the changes.
In `@apps/api/src/modules/webhooks/webhooks.controller.ts`:
- Around line 119-134: The current normalization uses Object.keys(body).length
to decide whether to treat parsed JSON as present, causing empty {} or [] to
fall back to req.rawBody; update the condition in the webhooks.controller.ts
normalization logic so parsed bodies are recognized if body is not
null/undefined and typeof body === 'object' (use Array.isArray(body) to allow
arrays and do not cast arrays to Record), i.e.: treat any object or array
(including empty ones) as the parsed payload (assign normalizedPayload = body),
only fall back to using req.rawBody when body is strictly undefined/null (and
the earlier string branch didn't match); keep contentType handling for the raw
branch and avoid using Object.keys(body).length and unsafe casts to Record for
arrays.
- Around line 283-298: The idempotency/duplicate path calls
executeAppWebhookResponses and returns the vendor response but does not persist
it, causing inboundGateway.response to remain null and breaking audit parity;
modify the idempotency branch (after executeAppWebhookResponses and before
sending the response in the block handling existingTraceIdOutside) to check if
existingTraceIdOutside.inboundGateway.response is null and, if so,
update/persist inboundGateway.response with the same object (status,
contentType, body) used for the outgoing response so support can see what was
sent; ensure you use the same persistence method/DB update used on the happy
path so transactional semantics match, or add a clear comment explaining why
persistence is intentionally omitted if you choose not to persist.
In `@apps/worker/src/modules/pipeline/normalization.service.ts`:
- Around line 187-219: The upsert currently updates traceId on conflict (in the
tx.insert(normalizedEntity)...onConflictDoUpdate block) which breaks downstream
FanOutService lookups that still query by the original traceId; change the
behavior so traceId is not overwritten on conflict — keep traceId immutable
after first insert (or switch to insert-once semantics) and instead expose a
stable handoff key (e.g., normalizedEntity.id/replicaId) to the outbox; update
the onConflictDoUpdate set to exclude traceId (only update canonicalType and
data) and ensure normalizedOutbox/FanOutService use normalizedEntityId or
replicaId as the lookup key.
- Around line 193-199: The JSON round-trip around canonicalData (used to create
safeData) can throw on non-serializable outputs (BigInt, cyclic refs); wrap the
JSON.stringify/parse conversion in a try/catch inside the normalization logic
that produces safeData, and on error throw a clear normalization-specific error
(e.g., NormalizationError or a new Error with a descriptive message) so the
caller handling L3 can record the FAIL state instead of bubbling a raw
serialization exception; update the block that computes safeData from
canonicalData to perform the defensive conversion and rethrow a normalized error
that includes context (which normalizer/function produced canonicalData).
In `@apps/worker/src/modules/pipeline/outbox.utils.spec.ts`:
- Around line 5-15: Add a new assertion to the processInChunks spec that
verifies chunking/concurrency rather than just output: create tasks that expose
asynchronous start and settle events (e.g., tasks that return a Promise whose
resolver you control or that push to a "started" array when invoked and only
resolve after a signal), call processInChunks([...], concurrency = 2, taskFn)
and assert that only 2 tasks enter the "started" array before you resolve them,
then resolve the first chunk and assert the next chunk starts; reference the
existing test helpers and the processInChunks function to locate where to add
this controlled-start/controlled-resolve test so it proves bounded parallel
execution rather than just final values.
In `@apps/worker/src/modules/pipeline/replica.service.ts`:
- Around line 137-160: The fallback object returned when extractor is falsy (in
the extractor ? ... : {...} expression that sets extracted) is dead code because
later the function unconditionally throws if resolvedEntityId is missing; update
replica.service.ts so this path is consistent: either detect missing extractor
earlier and throw a clear error (e.g., when extractor is falsy and appProfile is
"default", throw a descriptive message referencing traceId/appProfile) or
implement a real fallback entityId derivation (derive a stable ID from
inbound.extReqId or from inbound.objectType + a deterministic hash of
inbound.request) and ensure extracted.entityId is set before using
resolvedEntityId; reference the variables/functions extractor, extracted,
inbound, resolvedEntityId, traceId, and appProfile to locate and modify the
code.
In `@packages/dbmanager/src/impl/sql-database-manager.ts`:
- Around line 64-67: migrateReplicaTables is dead code because migrateAllSchemas
currently routes through schemaMgr.applyPlan and
provisionGatewayTables→provisionReplicaTables; either wire migrateAllSchemas to
call migrateReplicaTables (replace or delegate the existing
provisionReplicaTables call) so the new helper is used, or remove
migrateReplicaTables to avoid an unused API surface. Update the caller path
(migrateAllSchemas in apps/api/src/db/database-manager.ts) to invoke
dbmanager.impl.sqlDatabaseManager.migrateReplicaTables(schemaName) when handling
SchemaPlan.REPLICA_ACTIVE, or delete the migrateReplicaTables method and any
related tests/exports if you prefer removal.
- Around line 213-242: Duplicate DDL for the replica_outbox table and its
associated index/constraint exists in provisionReplicaTables and
provisionOutboundTables; keep the canonical definition in provisionReplicaTables
(which includes schema_name) and remove the duplicate CREATE TABLE, CREATE INDEX
(idx_replica_outbox_claim) and DO $$... ADD CONSTRAINT (unique on trace_id,
connection_id) block from provisionOutboundTables so there is a single source of
truth for replica_outbox; if provisionOutboundTables referenced
idx_replica_outbox_claim or the constraint anywhere else, update those
references to rely on the provisionReplicaTables definition instead.
In `@packages/pieces/application/revenova/src/upsertRevenovaObject.ts`:
- Around line 45-57: The XML extraction loop stores raw entity text (value)
without decoding XML entities, so fields like "<sf:Name>Acme & Co</sf:Name>"
become "Acme & Co"; fix by decoding standard XML entities and numeric
character references before assigning into data: after computing value, run an
entity/charref decoder (handle &, <, >, ", ' and numeric
refs like &#xHH; and &#DD;) and use the decoded string when setting data[key];
keep existing namespace stripping logic (rawKey, key) and the
SOAP_STRUCTURAL_TAGS filter unchanged. Ensure the decoder is used inside the
same loop that uses pattern.exec so data stores decoded values.
- Around line 26-64: parseSalesforceSoapXml currently flattens the entire XML so
repeated <Notification> blocks overwrite earlier records and entityId is taken
from the last match; update it to detect and reject multi-notification payloads
(e.g., search for multiple <Notification> occurrences and return null or throw)
and change extraction to iterate per sObject subtree: locate each <sObject> (or
<Notification>/sObject) block, parse leaf elements scoped to that subtree (so
sf:Id is extracted from the same sObject), and only accept single-sObject
payloads—use the existing pattern/SOAP_STRUCTURAL_TAGS logic inside the
per-sObject loop and derive entityId from that scoped data instead of the global
data map.
In `@TECHNICAL_DEBT.md`:
- Line 235: Update the section header text "Rate Limiting Updates" in
TECHNICAL_DEBT.md to the hyphenated form "Rate-Limiting Updates"; locate the
header exactly as written and replace it so the compound modifier is hyphenated
for consistent style across section headings.
---
Outside diff comments:
In `@apps/api/src/modules/webhooks/webhooks.controller.ts`:
- Around line 187-298: Extract the duplicated response-send logic into a private
helper (e.g., sendAppResponse) and call it from both the happy-path and
idempotency-path; the helper should accept the Express Response, the
customResponse object returned by executeAppWebhookResponses, and a context flag
('normal' | 'idempotency'), perform the unified this.logger.debug call
(including contentType, status and path/context) and then call
res.status(...).set('Content-Type', ...).send(...); replace the inline logging +
res.status(...).set(...).send(...) at the two call sites (where
executeAppWebhookResponses is used) with calls to sendAppResponse, keeping the
DB persist logic only in the original happy-path before invoking the helper.
In `@apps/worker/src/app.module.ts`:
- Around line 9-22: The worker no longer imports AiWorkerModule so the
CopilotWorker consumer that listens to AiCopilotQueue never runs; re-enable
AiWorkerModule by adding AiWorkerModule back into the imports array in
app.module.ts (so CopilotWorker is instantiated and consumes AiCopilotQueue), or
if you intend to stop processing these jobs instead remove or disable the API
producer (AiCopilotQueue usage in
apps/api/src/modules/ai/controllers/ai.controller.ts) and delete the orphaned
AiWorkerModule/CopilotWorker files; update whichever path you choose
consistently (either restore AiWorkerModule import or remove queue production
and module/class files) so producers and consumers match.
In `@apps/worker/src/modules/pipeline/fanout.service.ts`:
- Around line 126-137: Update the outdated comment wording that still references
"sourceId" around the GEM lookup using replicaEntity.entityId: change comments
near the replicaEntity/select block (the lines that fetch replicaRows by traceId
and assign srcVendorId) to refer to entityId or srcVendorId/vendorEntityId as
appropriate so they match the implementation that uses replicaEntity.entityId
and sets srcVendorId; verify any nearby variable names like replicaRows,
traceId, and srcVendorId are described consistently (e.g., "Fetch entityId from
replicaEntity for GEM threading" instead of "sourceId").
In `@packages/database/src/schema/pipeline.ts`:
- Around line 104-113: Update the stale block comment for replica_entity to
reflect the new schema: replace references to uniqueness on (entityType,
sourceId) with uniqueness on (connectionId, entityType, entityId), remove
mention of srcReqTraceId and state that traceId now links back to L1, and keep
the note about version incrementing for optimistic concurrency; locate the
docstring above the replica_entity definition and edit the text to mention
"connectionId, entityType, entityId" uniqueness and "traceId" as the L1 link
while preserving the overall purpose description.
- Around line 84-102: Update the stale inline comment that still says "Read
inbound_gateway payload" to reflect the renamed column by changing it to "Read
inbound_gateway request" in replica.service.ts where the code accesses
inbound.request (the comment adjacent to the code that reads inbound.request
should be the only change).
---
Duplicate comments:
In `@TECHNICAL_DEBT.md`:
- Line 40: Update the cleanup task text to use the actual worker class name:
replace the incorrect reference to NormalizedOutboxService with
NormalizedOutboxWorker so documentation matches the implementation (the class is
defined as NormalizedOutboxWorker in the codebase).
🪄 Autofix (Beta)
✅ Autofix completed
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: caf3eeb1-cb1e-468f-9dcd-2176dbb4d02f
⛔ Files ignored due to path filters (1)
pnpm-lock.yamlis excluded by!**/pnpm-lock.yaml
📒 Files selected for processing (30)
TECHNICAL_DEBT.mdapps/api/package.jsonapps/api/src/app/app.module.tsapps/api/src/db/database-manager.tsapps/api/src/db/db-cli.tsapps/api/src/main.tsapps/api/src/modules/trace/trace.service.tsapps/api/src/modules/trigger/trigger.module.tsapps/api/src/modules/webhooks/webhooks.controller.spec.tsapps/api/src/modules/webhooks/webhooks.controller.tsapps/web/src/modules/connections/components/ActiveConnectionCard.tsxapps/worker/package.jsonapps/worker/src/app.module.tsapps/worker/src/main.tsapps/worker/src/modules/pipeline/fanout.service.tsapps/worker/src/modules/pipeline/normalization.service.spec.tsapps/worker/src/modules/pipeline/normalization.service.tsapps/worker/src/modules/pipeline/outbox.utils.spec.tsapps/worker/src/modules/pipeline/replica.service.spec.tsapps/worker/src/modules/pipeline/replica.service.tsdocker-compose.ymlinfra/debezium/application.propertiesinfra/debezium/application.properties.prodpackage.jsonpackages/database/src/schema/pipeline.tspackages/dbmanager/src/impl/sql-database-manager.tspackages/pieces/application/revenova/src/upsertRevenovaObject.tspackages/pieces/application/revenova/src/upsertTMSObject.tspackages/pieces/platform/framework/src/jsonata-normalizer.tspackages/pieces/platform/framework/src/normalizer.ts
💤 Files with no reviewable changes (1)
- apps/api/src/modules/trigger/trigger.module.ts
|
Note Autofix is a beta feature. Expect some limitations and changes as we gather feedback and continue to improve it. Fixes Applied SuccessfullyFixed 8 file(s) based on 13 unresolved review comments. Files modified:
Commit: The changes have been pushed to the Time taken: |
Fixed 8 file(s) based on 13 unresolved review comments. Co-authored-by: CodeRabbit <noreply@coderabbit.ai>
There was a problem hiding this comment.
Actionable comments posted: 8
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
apps/api/src/modules/webhooks/webhooks.controller.ts (2)
282-325: 🧹 Nitpick | 🔵 TrivialReuse one duplicate-path app response for persistence and sending.
executeAppWebhookResponsesis called once for persistence and again for the HTTP response. If a handler adds timestamps, IDs, or has side effects,inbound_gateway.responsemay not match what was actually sent.Proposed refactor
- // If the existing record has no response and we can generate one, persist it - const customResponse = executeAppWebhookResponses(body, headers); - if (customResponse && existingRecord && existingRecord.response === null) { + // If the existing record has no response and we can generate one, persist it. + const duplicateCustomResponse = executeAppWebhookResponses(body, headers); + if (duplicateCustomResponse && existingRecord && existingRecord.response === null) { await this.db.transaction(async (tx) => { @@ response: { - status: customResponse.status, - contentType: customResponse.contentType, - body: customResponse.body, + status: duplicateCustomResponse.status, + contentType: duplicateCustomResponse.contentType, + body: duplicateCustomResponse.body, }, @@ - const customResponse = executeAppWebhookResponses(body, headers); + const customResponse = + typeof duplicateCustomResponse !== 'undefined' + ? duplicateCustomResponse + : executeAppWebhookResponses(body, headers);🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/api/src/modules/webhooks/webhooks.controller.ts` around lines 282 - 325, Call executeAppWebhookResponses once and reuse its result for both persistence and sending so the stored inbound_gateway.response exactly matches what is returned; compute const customResponse = executeAppWebhookResponses(body, headers) once before the try/catch (or at the top of the duplicate-path handling block), then use that same customResponse inside the transaction update (referencing existingRecord, existingTraceIdOutside, inboundGateway) and later when writing the HTTP response via res (status, Content-Type, send). Ensure no second call to executeAppWebhookResponses remains.
154-182:⚠️ Potential issue | 🟠 MajorDo not reject the webhook after the transactional outbox commit.
The DB transaction already writes
inbound_outbox; if the direct queue send fails, throwing here returns an error to the vendor even though the event is durably stored for CDC/outbox relay. Make this send best-effort, matching the L2 pattern.Proposed fix
await this.queueService .send(QueueName.InboundQueue, { traceId, connectionId }) .catch((err: unknown) => { - this.logger.error( + this.logger.warn( { event: 'l1.enqueue_failed', traceId, err: err instanceof Error ? err.message : String(err), }, - 'Failed to enqueue L1 event — rejecting webhook', + 'Failed to best-effort enqueue L1 event — relying on inbound_outbox', ); - throw err; });🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/api/src/modules/webhooks/webhooks.controller.ts` around lines 154 - 182, The webhook handler currently sends to queueService.send(QueueName.InboundQueue, { traceId: inboundGatewayId, connectionId }) and on failure logs and re-throws, which causes the webhook to be rejected even though inbound_outbox was committed; change the send to be best-effort: call queueService.send and in the .catch handler log the failure (use this.logger.error with event 'l1.enqueue_failed', traceId: inboundGatewayId and the error message) but do NOT re-throw—simply swallow the error so the HTTP 202 is still returned after the DB transaction that inserted inboundOutbox completes.
♻️ Duplicate comments (3)
TECHNICAL_DEBT.md (1)
37-40:⚠️ Potential issue | 🟡 MinorUse one canonical worker name in this section.
Line 40 says
NormalizedOutboxService, while Line 46 saysNormalizedOutboxWorker. This inconsistency can cause incorrect cleanup work; align both references to the actual class name.Also applies to: 46-46
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@TECHNICAL_DEBT.md` around lines 37 - 40, The comment points out an inconsistent reference: the section uses both NormalizedOutboxService and NormalizedOutboxWorker; pick the actual canonical class name (NormalizedOutboxWorker) and update all occurrences in this section so references to the legacy cron-polling component, any cleanup tasks, and documentation uniformly use NormalizedOutboxWorker; while updating, double-check nearby mentions of normalized_outbox, CdcRelayController, and Debezium publication to ensure wording remains accurate (e.g., that normalized_outbox is not registered in Debezium and CdcRelayController does not listen) and adjust any sentences that referenced the old name so they still convey the same meaning.apps/worker/src/modules/pipeline/normalization.service.ts (1)
219-243:⚠️ Potential issue | 🔴 CriticalUse a stable L3→L4 handoff key before upserting by
replicaId.Line 223 keeps
normalizedEntity.traceIdimmutable, but lines 234-243 still enqueue the currenttraceId. SinceFanOutServicestill looks up normalized rows bytraceId, a re-ingest of the same replica under a new trace publishes a message that cannot find its normalized row. ThreadnormalizedEntity.id/replicaIdthroughnormalizedOutboxand the NormalizedQueue, or keep per-trace insert semantics until L4 is changed.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/worker/src/modules/pipeline/normalization.service.ts` around lines 219 - 243, The outbox currently enqueues using the transient traceId which breaks FanOutService lookups after a re-ingest; change the normalized outbox insert to include and use the stable normalizedEntity.id (or normalizedEntity.replicaId) as the handoff key instead of the per-ingest traceId by threading normalizedEntity.id/replicaId into the normalizedOutbox insert (and into the NormalizedQueue message) so downstream lookups use the stable replica/id key; alternatively, if you cannot change downstream consumers yet, keep the per-trace insert semantics but ensure the outbox also persists the stable normalizedEntity.id/replicaId alongside traceId so FanOutService can resolve the normalized row.apps/api/src/modules/webhooks/webhooks.controller.ts (1)
119-130:⚠️ Potential issue | 🟡 MinorPreserve parsed JSON arrays instead of wrapping them.
requestis meant to be the inbound payload. Wrapping arrays as{ items: body }changes the vendor payload shape and can break downstream extractors/auditing that expect the original JSON array.Proposed fix
- let normalizedPayload: Record<string, unknown>; + let normalizedPayload: Record<string, unknown> | unknown[]; @@ } else if (body != null && typeof body === 'object') { - // Accept both objects and arrays as parsed payloads - if (Array.isArray(body)) { - normalizedPayload = { items: body }; - } else { - normalizedPayload = body as Record<string, unknown>; - } + // Preserve parsed JSON exactly, including [] and {}. + normalizedPayload = body as Record<string, unknown> | unknown[];🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/api/src/modules/webhooks/webhooks.controller.ts` around lines 119 - 130, The code currently wraps parsed JSON arrays into an object ({ items: body }) which alters the original vendor payload shape; update the logic in the webhooks controller where normalizedPayload is assigned (the body handling block that checks Array.isArray(body)) to preserve arrays as-is by assigning normalizedPayload = body (or body as unknown) instead of wrapping, while retaining the existing branches for string/raw handling and object handling so downstream consumers that expect the original JSON array receive the unchanged array.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@apps/worker/src/modules/pipeline/outbox.utils.spec.ts`:
- Around line 17-33: Consolidate the three duplicate tests that assert
processInChunks rejects for invalid concurrency by replacing them with a single
parameterized test using Jest's it.each: supply an array of invalid concurrency
values (e.g., [0, -1, 1.5]) and for each value call processInChunks([1], value,
(n) => Promise.resolve(n)) and expect rejection with "concurrency must be a
positive integer"; update the spec around the existing tests in
outbox.utils.spec.ts to use this it.each pattern so adding more invalid values
(NaN, Infinity, -0.5) becomes trivial.
- Around line 35-70: The test is using real timers (setTimeout) which makes
concurrency assertions flaky; replace those sleeps with deterministic microtask
flushes instead: after creating resultPromise and after calling resolvers, use a
microtask flush (e.g., await Promise.resolve() or a shared helper like
flushPromises / queueMicrotask) to let processInChunks advance, so update the
three await new Promise(resolve => setTimeout(resolve, 10)) points to await a
microtask flush; keep references to started, resolvers and the call to
processInChunks unchanged.
In `@packages/dbmanager/src/impl/sql-database-manager.ts`:
- Around line 176-190: The CREATE TABLE IF NOT EXISTS in provisionReplicaTables
won’t add the new entity_id column or the UNIQUE (connection_id, entity_type,
entity_id) constraint to existing tenant schemas, so migrateReplicaTables must
be run and perform an ALTER for existing tables before any L2 writes in
replica.service.ts; update provisionReplicaTables to call or invoke
migrateReplicaTables (or inline its logic) after creating the table, and
implement migrateReplicaTables to: 1) check for existence of entity_id column on
"${schemaName}".replica_entity, 2) ALTER TABLE to ADD COLUMN entity_id
VARCHAR(255) DEFAULT ''/NULL as appropriate, 3) populate entity_id for existing
rows if determinable, and 4) add the UNIQUE constraint (or create a unique
index) safely (using IF NOT EXISTS / DROP CONSTRAINT IF EXISTS patterns) so old
schemas are migrated before replication inserts occur in replica.service.ts.
- Around line 92-139: The migration must add any new inbound_gateway columns
before creating indexes or relying on them: alter the inbound_gateway table
(function/logic around schemaName handling in sql-database-manager.ts) to
idempotently ADD COLUMN IF NOT EXISTS headers JSONB and ADD COLUMN IF NOT EXISTS
ext_req_id TEXT (and any other new controller-written columns still missing)
prior to creating the idx_l1_ext_id unique index and any ext_req_id-dependent
logic; keep the existing payload→request rename, response JSONB add, and
GIN/index creation but move or add these ALTER TABLE ADD COLUMN IF NOT EXISTS
statements before creating indexes so provisioning/ingest won't fail for tenants
missing headers/ext_req_id.
In `@packages/pieces/application/revenova/src/upsertRevenovaObject.ts`:
- Around line 33-38: The current regex used to compute notificationMatches
(const notificationMatches = xml.match(/<[^:>]*:?Notification[^>]*>/gi)) also
matches the enclosing <notifications> wrapper; update the regex used in the
xml.match call so it only matches element names that are exactly "Notification"
(not "notifications") by requiring the name to end (e.g., use a lookahead that
asserts a boundary like whitespace, '>', or '/'): keep the global and
case-insensitive flags, then keep the existing logic that returns null if
notificationMatches && notificationMatches.length > 1.
- Around line 81-89: In decodeXmlEntities, the replacement order causes
double-decoding of sequences like "&lt;" and uses fromCharCode which
mishandles non-BMP code points; update the function (decodeXmlEntities) to
perform entity replacements for numeric hex (&#x...;), numeric decimal (&#...;),
and named entities (<, >, ", ') first, and apply the ampersand
replacement (&) last, and replace String.fromCharCode(...) calls with
String.fromCodePoint(...) for both hex and decimal numeric entity handlers so
high Unicode code points are decoded correctly.
In `@TECHNICAL_DEBT.md`:
- Line 222: Replace the hardcoded `YYYY-MM-DD` placeholder in the
`X-Deprecation-Warning` header on legacy webhook routes with a concrete,
computed or config-driven sunset date: read the date from a central config/env
value (e.g., config.SUNSET_DATE or a helper like getSunsetDate()), format it as
YYYY-MM-DD, and inject that value into the header for responses on the legacy
route `/webhooks/:orgSlug/:connectionSlug/:secretToken` (and any other places
emitting `X-Deprecation-Warning`) so the header becomes `X-Deprecation-Warning:
"Legacy endpoint; migrate to /webhooks/:orgSlug/:connectionSlug/:secretToken by
<sunset-date>"`.
- Line 228: The doc currently advises storing webhook_secret in plaintext;
instead update docs and code to store a derived secure hash (e.g., HMAC-SHA256
or salted SHA-256) in the webhook_secret column and change any
resolution/validation code (e.g., resolveWebhookSecret, validateWebhookToken, or
equivalent lookup by orgSlug/connectionSlug) to compute the same HMAC/hash from
the presented token and perform a constant-time comparison; add a DB migration
to replace plaintext column semantics (or add webhook_secret_hash), update
creation/rotation flows to store only the derived hash, and update tests and
examples accordingly.
---
Outside diff comments:
In `@apps/api/src/modules/webhooks/webhooks.controller.ts`:
- Around line 282-325: Call executeAppWebhookResponses once and reuse its result
for both persistence and sending so the stored inbound_gateway.response exactly
matches what is returned; compute const customResponse =
executeAppWebhookResponses(body, headers) once before the try/catch (or at the
top of the duplicate-path handling block), then use that same customResponse
inside the transaction update (referencing existingRecord,
existingTraceIdOutside, inboundGateway) and later when writing the HTTP response
via res (status, Content-Type, send). Ensure no second call to
executeAppWebhookResponses remains.
- Around line 154-182: The webhook handler currently sends to
queueService.send(QueueName.InboundQueue, { traceId: inboundGatewayId,
connectionId }) and on failure logs and re-throws, which causes the webhook to
be rejected even though inbound_outbox was committed; change the send to be
best-effort: call queueService.send and in the .catch handler log the failure
(use this.logger.error with event 'l1.enqueue_failed', traceId: inboundGatewayId
and the error message) but do NOT re-throw—simply swallow the error so the HTTP
202 is still returned after the DB transaction that inserted inboundOutbox
completes.
---
Duplicate comments:
In `@apps/api/src/modules/webhooks/webhooks.controller.ts`:
- Around line 119-130: The code currently wraps parsed JSON arrays into an
object ({ items: body }) which alters the original vendor payload shape; update
the logic in the webhooks controller where normalizedPayload is assigned (the
body handling block that checks Array.isArray(body)) to preserve arrays as-is by
assigning normalizedPayload = body (or body as unknown) instead of wrapping,
while retaining the existing branches for string/raw handling and object
handling so downstream consumers that expect the original JSON array receive the
unchanged array.
In `@apps/worker/src/modules/pipeline/normalization.service.ts`:
- Around line 219-243: The outbox currently enqueues using the transient traceId
which breaks FanOutService lookups after a re-ingest; change the normalized
outbox insert to include and use the stable normalizedEntity.id (or
normalizedEntity.replicaId) as the handoff key instead of the per-ingest traceId
by threading normalizedEntity.id/replicaId into the normalizedOutbox insert (and
into the NormalizedQueue message) so downstream lookups use the stable
replica/id key; alternatively, if you cannot change downstream consumers yet,
keep the per-trace insert semantics but ensure the outbox also persists the
stable normalizedEntity.id/replicaId alongside traceId so FanOutService can
resolve the normalized row.
In `@TECHNICAL_DEBT.md`:
- Around line 37-40: The comment points out an inconsistent reference: the
section uses both NormalizedOutboxService and NormalizedOutboxWorker; pick the
actual canonical class name (NormalizedOutboxWorker) and update all occurrences
in this section so references to the legacy cron-polling component, any cleanup
tasks, and documentation uniformly use NormalizedOutboxWorker; while updating,
double-check nearby mentions of normalized_outbox, CdcRelayController, and
Debezium publication to ensure wording remains accurate (e.g., that
normalized_outbox is not registered in Debezium and CdcRelayController does not
listen) and adjust any sentences that referenced the old name so they still
convey the same meaning.
🪄 Autofix (Beta)
✅ Autofix completed
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: bed07a8f-87fe-4c0d-9b32-47828b34c565
📒 Files selected for processing (8)
TECHNICAL_DEBT.mdapps/api/src/db/database-manager.tsapps/api/src/modules/webhooks/webhooks.controller.tsapps/worker/src/modules/pipeline/normalization.service.tsapps/worker/src/modules/pipeline/outbox.utils.spec.tsapps/worker/src/modules/pipeline/replica.service.tspackages/dbmanager/src/impl/sql-database-manager.tspackages/pieces/application/revenova/src/upsertRevenovaObject.ts
| // Idempotently rename payload → request for pre-existing schemas. | ||
| await this.db.$client.query(` | ||
| DO $$ BEGIN | ||
| IF EXISTS ( | ||
| SELECT 1 FROM information_schema.columns | ||
| WHERE table_schema = '${schemaName}' | ||
| AND table_name = 'inbound_gateway' | ||
| AND column_name = 'payload' | ||
| ) THEN | ||
| ALTER TABLE "${schemaName}".inbound_gateway RENAME COLUMN payload TO request; | ||
| END IF; | ||
| END $$; | ||
| `); | ||
|
|
||
| // Idempotently add response column for pre-existing schemas. | ||
| await this.db.$client.query(` | ||
| CREATE INDEX IF NOT EXISTS idx_inbound_gateway_status | ||
| ON "${schemaName}".inbound_gateway (status); | ||
| `); | ||
| ALTER TABLE "${schemaName}".inbound_gateway | ||
| ADD COLUMN IF NOT EXISTS response JSONB; | ||
| `); | ||
|
|
||
| // Idempotency: unique (connection_id, ext_req_id) prevents duplicate | ||
| // vendor events from being ingested twice. | ||
| await this.db.$client.query(` | ||
| CREATE INDEX IF NOT EXISTS idx_inbound_gateway_object_type | ||
| ON "${schemaName}".inbound_gateway (object_type) | ||
| WHERE object_type IS NOT NULL; | ||
| `); | ||
| CREATE UNIQUE INDEX IF NOT EXISTS idx_l1_ext_id | ||
| ON "${schemaName}".inbound_gateway (connection_id, ext_req_id) | ||
| WHERE ext_req_id IS NOT NULL; | ||
| `); | ||
|
|
||
| await this.db.$client.query(` | ||
| CREATE INDEX IF NOT EXISTS idx_inbound_gateway_created_at | ||
| ON "${schemaName}".inbound_gateway (created_at DESC); | ||
| `); | ||
| CREATE INDEX IF NOT EXISTS idx_l1_object_type | ||
| ON "${schemaName}".inbound_gateway (object_type) | ||
| WHERE object_type IS NOT NULL; | ||
| `); | ||
|
|
||
| await this.db.$client.query(` | ||
| CREATE INDEX IF NOT EXISTS idx_l1_status | ||
| ON "${schemaName}".inbound_gateway (status); | ||
| `); | ||
|
|
||
| // Drop old payload GIN index if present, create new request GIN index. | ||
| await this.db.$client.query(` | ||
| DROP INDEX IF EXISTS "${schemaName}".idx_l1_payload_gin; | ||
| `); | ||
|
|
||
| await this.db.$client.query(` | ||
| CREATE INDEX IF NOT EXISTS idx_l1_request_gin | ||
| ON "${schemaName}".inbound_gateway USING gin (request); | ||
| `); |
There was a problem hiding this comment.
Add missing inbound_gateway ALTERs for existing tenant schemas.
For existing schemas, CREATE TABLE IF NOT EXISTS at Line 75 is skipped. The migration only renames payload and adds response, but the controller now writes headers/ext_req_id, and Line 115 indexes ext_req_id. Tenants missing those columns will fail during provisioning or ingest.
Proposed migration hardening
// Idempotently add response column for pre-existing schemas.
await this.db.$client.query(`
ALTER TABLE "${schemaName}".inbound_gateway
+ ADD COLUMN IF NOT EXISTS request JSONB NOT NULL DEFAULT '{}'::jsonb,
ADD COLUMN IF NOT EXISTS response JSONB;
+
+ ALTER TABLE "${schemaName}".inbound_gateway
+ ALTER COLUMN request DROP DEFAULT;
+ `);
+
+ await this.db.$client.query(`
+ ALTER TABLE "${schemaName}".inbound_gateway
+ ADD COLUMN IF NOT EXISTS headers JSONB,
+ ADD COLUMN IF NOT EXISTS ext_req_id VARCHAR(255),
+ ADD COLUMN IF NOT EXISTS object_type VARCHAR(100);
`);🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@packages/dbmanager/src/impl/sql-database-manager.ts` around lines 92 - 139,
The migration must add any new inbound_gateway columns before creating indexes
or relying on them: alter the inbound_gateway table (function/logic around
schemaName handling in sql-database-manager.ts) to idempotently ADD COLUMN IF
NOT EXISTS headers JSONB and ADD COLUMN IF NOT EXISTS ext_req_id TEXT (and any
other new controller-written columns still missing) prior to creating the
idx_l1_ext_id unique index and any ext_req_id-dependent logic; keep the existing
payload→request rename, response JSONB add, and GIN/index creation but move or
add these ALTER TABLE ADD COLUMN IF NOT EXISTS statements before creating
indexes so provisioning/ingest won't fail for tenants missing
headers/ext_req_id.
| private async provisionReplicaTables(schemaName: string): Promise<void> { | ||
| await this.db.$client.query(` | ||
| CREATE TABLE IF NOT EXISTS "${schemaName}".replica_entity ( | ||
| id UUID PRIMARY KEY DEFAULT gen_random_uuid(), | ||
| connection_id UUID NOT NULL, | ||
| trace_id UUID NOT NULL, | ||
| src_req_trace_id UUID NOT NULL, | ||
| source_id VARCHAR(255) NOT NULL, | ||
| entity_type VARCHAR(100) NOT NULL, | ||
| data JSONB NOT NULL, | ||
| version INTEGER NOT NULL DEFAULT 1, | ||
| updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), | ||
| CONSTRAINT uq_l2_entity UNIQUE (connection_id, entity_type, source_id) | ||
| id UUID PRIMARY KEY DEFAULT gen_random_uuid(), | ||
| connection_id UUID NOT NULL, | ||
| trace_id UUID NOT NULL, | ||
| entity_id VARCHAR(255) NOT NULL, | ||
| entity_type VARCHAR(100) NOT NULL, | ||
| data JSONB NOT NULL, | ||
| version INTEGER NOT NULL DEFAULT 1, | ||
| created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), | ||
| updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), | ||
| CONSTRAINT uq_l2_entity UNIQUE (connection_id, entity_type, entity_id) | ||
| ); | ||
| `); |
There was a problem hiding this comment.
Migrate existing replica_entity tables before the worker writes entity_id.
CREATE TABLE IF NOT EXISTS will not add entity_id or the new (connection_id, entity_type, entity_id) uniqueness to existing tenant schemas. After migrateReplicaTables, L2 inserts at replica.service.ts will still fail on old schemas with no entity_id column.
Proposed migration block
private async provisionReplicaTables(schemaName: string): Promise<void> {
await this.db.$client.query(`
CREATE TABLE IF NOT EXISTS "${schemaName}".replica_entity (
@@
);
`);
+
+ await this.db.$client.query(`
+ ALTER TABLE "${schemaName}".replica_entity
+ ADD COLUMN IF NOT EXISTS entity_id VARCHAR(255);
+
+ DO $$ BEGIN
+ IF EXISTS (
+ SELECT 1 FROM information_schema.columns
+ WHERE table_schema = '${schemaName}'
+ AND table_name = 'replica_entity'
+ AND column_name = 'source_id'
+ ) THEN
+ UPDATE "${schemaName}".replica_entity
+ SET entity_id = source_id
+ WHERE entity_id IS NULL
+ AND source_id IS NOT NULL;
+ END IF;
+
+ IF EXISTS (
+ SELECT 1
+ FROM "${schemaName}".replica_entity
+ WHERE entity_id IS NULL
+ LIMIT 1
+ ) THEN
+ RAISE EXCEPTION 'Cannot migrate %.replica_entity: entity_id has NULL rows after backfill', '${schemaName}';
+ END IF;
+
+ ALTER TABLE "${schemaName}".replica_entity
+ ALTER COLUMN entity_id SET NOT NULL;
+
+ ALTER TABLE "${schemaName}".replica_entity
+ DROP CONSTRAINT IF EXISTS uq_l2_entity;
+
+ ALTER TABLE "${schemaName}".replica_entity
+ ADD CONSTRAINT uq_l2_entity UNIQUE (connection_id, entity_type, entity_id);
+ EXCEPTION WHEN duplicate_object THEN NULL;
+ END $$;
+ `);🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@packages/dbmanager/src/impl/sql-database-manager.ts` around lines 176 - 190,
The CREATE TABLE IF NOT EXISTS in provisionReplicaTables won’t add the new
entity_id column or the UNIQUE (connection_id, entity_type, entity_id)
constraint to existing tenant schemas, so migrateReplicaTables must be run and
perform an ALTER for existing tables before any L2 writes in replica.service.ts;
update provisionReplicaTables to call or invoke migrateReplicaTables (or inline
its logic) after creating the table, and implement migrateReplicaTables to: 1)
check for existence of entity_id column on "${schemaName}".replica_entity, 2)
ALTER TABLE to ADD COLUMN entity_id VARCHAR(255) DEFAULT ''/NULL as appropriate,
3) populate entity_id for existing rows if determinable, and 4) add the UNIQUE
constraint (or create a unique index) safely (using IF NOT EXISTS / DROP
CONSTRAINT IF EXISTS patterns) so old schemas are migrated before replication
inserts occur in replica.service.ts.
|
Note Autofix is a beta feature. Expect some limitations and changes as we gather feedback and continue to improve it. Fixes Applied SuccessfullyFixed 4 file(s) based on 8 unresolved review comments. Files modified:
Commit: The changes have been pushed to the Time taken: |
Fixed 4 file(s) based on 8 unresolved review comments. Co-authored-by: CodeRabbit <noreply@coderabbit.ai>
There was a problem hiding this comment.
Actionable comments posted: 3
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
apps/api/src/modules/webhooks/webhooks.controller.ts (2)
238-324:⚠️ Potential issue | 🟠 MajorReuse the stored idempotency response instead of recomputing it.
existingRecord.responseis selected but never used, andexecuteAppWebhookResponsescan run twice. Send the stored response when present; otherwise generate once, persist it if needed, and send that same object.🛠️ Proposed direction
+ let idempotentResponse: + | { body: string; contentType: string; status: number } + | null = null; ... - const customResponse = executeAppWebhookResponses(body, headers); - if (customResponse && existingRecord && (existingRecord as { traceId: string; response: unknown | null }).response === null) { + idempotentResponse = + asStoredWebhookResponse(existingRecord?.response) ?? + executeAppWebhookResponses(body, headers); + + if (idempotentResponse && existingRecord?.response == null) { await this.db.transaction(async (tx) => { ... response: { - status: customResponse.status, - contentType: customResponse.contentType, - body: customResponse.body, + status: idempotentResponse.status, + contentType: idempotentResponse.contentType, + body: idempotentResponse.body, }, ... - const customResponse = executeAppWebhookResponses(body, headers); + const customResponse = + idempotentResponse ?? executeAppWebhookResponses(body, headers);Add a small shape guard for
asStoredWebhookResponsebefore trusting the JSON value from the database.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/api/src/modules/webhooks/webhooks.controller.ts` around lines 238 - 324, The code currently ignores the stored idempotent response (existingRecord.response) and calls executeAppWebhookResponses twice; change the flow to first check existingRecord.response and, if present and valid, use that stored response to send to the client (validate its shape with a small guard like asStoredWebhookResponse before trusting JSON), otherwise call executeAppWebhookResponses exactly once, persist that generated response into inboundGateway when existingRecord exists with response === null (in the same transaction path where you update the row), and then send that same response object to the client; refer to existingRecord, executeAppWebhookResponses, inboundGateway, and the db.transaction/update block to locate where to read, validate, persist, and reuse the response.
154-182:⚠️ Potential issue | 🟠 MajorDon’t abort the webhook after the durable outbox write succeeds.
The transaction already commits a
PENDINGinboundOutboxrow, but Line 181 still throws on queue send failure. That contradicts the recovery comment and can make vendors retry an event that is already durably stored, creating duplicates when noextReqIdis present.🛠️ Proposed fix
await this.queueService .send(QueueName.InboundQueue, { traceId, connectionId }) .catch((err: unknown) => { - this.logger.error( + this.logger.warn( { event: 'l1.enqueue_failed', traceId, err: err instanceof Error ? err.message : String(err), }, - 'Failed to enqueue L1 event — rejecting webhook', + 'Failed to enqueue L1 event — inbound outbox will retry', ); - throw err; });🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/api/src/modules/webhooks/webhooks.controller.ts` around lines 154 - 182, The code currently rethrows inside the queueService.send(...).catch(...) block which will abort the webhook even though the inboundOutbox row was committed; remove the rethrow so failures are logged but do not propagate. Specifically, after inserting the inboundOutbox row (using inboundGatewayId/connectionId) and calling this.queueService.send(QueueName.InboundQueue, { traceId, connectionId }), update the catch handler on queueService.send to log the error via this.logger.error (keeping event 'l1.enqueue_failed' and traceId) but do not throw err — let the request continue so the webhook returns 202 and the pending outbox record can be picked up by the worker.
♻️ Duplicate comments (1)
TECHNICAL_DEBT.md (1)
40-40:⚠️ Potential issue | 🟡 MinorUse a single correct class name in the L3→L4 description.
Line 40 still mentions
NormalizedOutboxService, while Line 46 correctly saysNormalizedOutboxWorker. Please align Line 40 toNormalizedOutboxWorkerto avoid implementation confusion.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@TECHNICAL_DEBT.md` at line 40, Update the L3→L4 description to use the correct class name: replace the mistaken `NormalizedOutboxService` reference with `NormalizedOutboxWorker` so the doc consistently refers to `NormalizedOutboxWorker` (ensure any other occurrences in this section also match).
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@apps/api/src/modules/webhooks/webhooks.controller.ts`:
- Around line 121-132: The idempotency/ACK path still passes the original body
to synchronous app-response handlers even when normalization used req.rawBody;
update the code so the same normalized/appResponseBody value is used everywhere.
Concretely, after computing normalizedPayload (and creating appResponseBody from
it), replace any usage of the raw body variable (body) in the idempotency/ack
code paths (the places that call the app response/ACK handlers) with
appResponseBody (or the normalizedPayload fallback that includes req.rawBody),
so XML/text requests that fell back to req.rawBody are forwarded to the handlers
consistently.
In `@packages/pieces/application/revenova/src/upsertRevenovaObject.ts`:
- Around line 160-168: The code currently casts r_obj['id'] or r_obj['Id'] to
string (entityId) without verifying type or emptiness; update the logic in
upsertRevenovaObject so that you check typeof r_obj['id'] === 'string' || typeof
r_obj['Id'] === 'string' and that the resulting entityId.trim().length > 0
before using it (and only delete r_obj['id']/r_obj['Id'] after successful
validation), returning null if the value is missing, not a string, or empty to
preserve ReplicaExtractorFn invariants.
- Around line 81-89: The decodeXmlEntities function currently calls
String.fromCodePoint on parsed numeric references without validation; update the
hex and decimal replacement callbacks in decodeXmlEntities to parse the code
point, check that it is within 0..0x10FFFF, and only call String.fromCodePoint
for valid values—otherwise return a safe fallback (e.g., the Unicode replacement
character '\uFFFD' or the original entity string). Modify the callbacks for
/&#x([0-9A-Fa-f]+);/ and /&#(\d+);/ accordingly so invalid code points do not
throw RangeError.
---
Outside diff comments:
In `@apps/api/src/modules/webhooks/webhooks.controller.ts`:
- Around line 238-324: The code currently ignores the stored idempotent response
(existingRecord.response) and calls executeAppWebhookResponses twice; change the
flow to first check existingRecord.response and, if present and valid, use that
stored response to send to the client (validate its shape with a small guard
like asStoredWebhookResponse before trusting JSON), otherwise call
executeAppWebhookResponses exactly once, persist that generated response into
inboundGateway when existingRecord exists with response === null (in the same
transaction path where you update the row), and then send that same response
object to the client; refer to existingRecord, executeAppWebhookResponses,
inboundGateway, and the db.transaction/update block to locate where to read,
validate, persist, and reuse the response.
- Around line 154-182: The code currently rethrows inside the
queueService.send(...).catch(...) block which will abort the webhook even though
the inboundOutbox row was committed; remove the rethrow so failures are logged
but do not propagate. Specifically, after inserting the inboundOutbox row (using
inboundGatewayId/connectionId) and calling
this.queueService.send(QueueName.InboundQueue, { traceId, connectionId }),
update the catch handler on queueService.send to log the error via
this.logger.error (keeping event 'l1.enqueue_failed' and traceId) but do not
throw err — let the request continue so the webhook returns 202 and the pending
outbox record can be picked up by the worker.
---
Duplicate comments:
In `@TECHNICAL_DEBT.md`:
- Line 40: Update the L3→L4 description to use the correct class name: replace
the mistaken `NormalizedOutboxService` reference with `NormalizedOutboxWorker`
so the doc consistently refers to `NormalizedOutboxWorker` (ensure any other
occurrences in this section also match).
🪄 Autofix (Beta)
✅ Autofix completed
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: b51e6e48-7dd1-4530-8345-e07d21fac89b
📒 Files selected for processing (4)
TECHNICAL_DEBT.mdapps/api/src/modules/webhooks/webhooks.controller.tsapps/worker/src/modules/pipeline/outbox.utils.spec.tspackages/pieces/application/revenova/src/upsertRevenovaObject.ts
|
Note Autofix is a beta feature. Expect some limitations and changes as we gather feedback and continue to improve it. Fixes Applied SuccessfullyFixed 2 file(s) based on 3 unresolved review comments. Files modified:
Commit: The changes have been pushed to the Time taken: |
Fixed 2 file(s) based on 3 unresolved review comments. Co-authored-by: CodeRabbit <noreply@coderabbit.ai>
Summary by CodeRabbit
New Features
Chores