e2e-sync-pipeline testing initial commit - #128
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:
📝 WalkthroughWalkthroughAdds tenant-scoped inbound transactional outbox and processor, registry-backed replica extractors/normalizers with a Revenova app package, a Trace Explorer (backend + UI), propagated connectionId through trigger/ingestion flows, tests, schema updates, and module/tsconfig/workspace wiring. Changes
Sequence Diagram(s)sequenceDiagram
actor Client
participant WebhooksController as "WebhooksController"
participant TriggerExecutor as "TriggerExecutorService"
participant StorageResolver as "StorageResolverService"
participant TenantDB as "Tenant DB (tenant schema)"
participant InboundOutbox as "InboundOutboxService"
participant Queue as "QueueService"
Client->>WebhooksController: POST /v1/webhooks/{connectionId}
WebhooksController->>TriggerExecutor: runWebhook(payload, connectionId)
TriggerExecutor->>StorageResolver: resolveSchemaName(connectionId)
StorageResolver-->>TriggerExecutor: schemaName
TriggerExecutor->>TenantDB: INSERT inbound_gateway (SET LOCAL search_path)
TenantDB-->>TriggerExecutor: inserted (traceId)
TriggerExecutor->>TenantDB: INSERT inbound_outbox (same tx)
TenantDB-->>TriggerExecutor: outbox row
TriggerExecutor-->>WebhooksController: 200 OK
par Background: Outbox drain
InboundOutbox->>TenantDB: SELECT ... FOR UPDATE SKIP LOCKED
TenantDB-->>InboundOutbox: claimed rows
InboundOutbox->>Queue: send(InboundQueue, {traceId, connectionId})
alt send succeeds
Queue-->>InboundOutbox: ack
InboundOutbox->>TenantDB: UPDATE status=SUCCESS
else send fails
Queue-->>InboundOutbox: error
InboundOutbox->>TenantDB: UPDATE status=RETRY/FAIL, nextRetryAt
end
end
sequenceDiagram
actor ReplicaWorker as "ReplicaWorker"
participant ReplicaService as "ReplicaService"
participant AppConns as "app_connections"
participant Registry as "Framework Registry"
participant TenantDB as "Tenant DB"
ReplicaWorker->>ReplicaService: handle(inboundRow)
ReplicaService->>AppConns: SELECT appName, metadata FROM app_connections WHERE id=...
AppConns-->>ReplicaService: {appName, metadata:{appProfile}}
ReplicaService->>Registry: getReplicaExtractor(appName, appProfile)
Registry-->>ReplicaService: extractorFn | undefined
alt extractorFn exists
ReplicaService->>ReplicaService: envelope = extractorFn(payload)
alt envelope != null
ReplicaService->>TenantDB: upsert replica_entity (entityType, data)
else envelope == null
ReplicaService-->>ReplicaWorker: throw extraction failed
end
else no extractor
ReplicaService->>TenantDB: upsert replica_entity (from inbound.objectType,payload)
end
Estimated code review effort🎯 4 (Complex) | ⏱️ ~75 minutes Possibly related PRs
Poem
🚥 Pre-merge checks | ✅ 2 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (2 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: 26
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (5)
apps/api/src/modules/trigger/poller.service.ts (2)
63-74:⚠️ Potential issue | 🟠 MajorMake the keyset cursor unique.
The cursor uses
(created_at, workspace_id), but those columns are not unique. If more than one active connection shares that tuple and the page boundary lands inside that group, the next query’s>predicate skips the remaining rows. Sinceconnection_idis now selected, includeac.idas the final tie-breaker in the cursor, predicate, and ordering.🐛 Proposed fix
- let cursor: { created_at: string; workspace_id: string } | undefined; + let cursor: + | { created_at: string; workspace_id: string; connection_id: string } + | undefined; @@ - cursor = { created_at: last.created_at, workspace_id: last.workspace_id }; + cursor = { + created_at: last.created_at, + workspace_id: last.workspace_id, + connection_id: last.connection_id, + }; @@ async fetchActivePollingConnectionsBatch(after?: { created_at: string; workspace_id: string; + connection_id: string; }): Promise<ActiveConnection[]> { @@ - where += ` AND (ac.created_at, ac.workspace_id) > ($2, $3)`; - params.push(after.created_at, after.workspace_id); + where += ` AND (ac.created_at, ac.workspace_id, ac.id) > ($2, $3, $4)`; + params.push(after.created_at, after.workspace_id, after.connection_id); @@ - ORDER BY ac.created_at ASC, ac.workspace_id ASC + ORDER BY ac.created_at ASC, ac.workspace_id ASC, ac.id ASCAlso applies to: 155-181
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/api/src/modules/trigger/poller.service.ts` around lines 63 - 74, The keyset cursor currently uses { created_at, workspace_id } which is not unique; update the cursor to include the connection id (ac.id) as a final tie-breaker and propagate that through fetchActivePollingConnectionsBatch and the loop (the cursor variable and where/predicate logic), change the ordering to include ac.id as the last ordering key, and adjust the next-page predicate to compare the three-field tuple (created_at, workspace_id, id) so rows with the same created_at and workspace_id aren’t skipped.
123-132:⚠️ Potential issue | 🟠 MajorScope the polling lock by
connectionId.The lock key at
apps/api/src/modules/trigger/trigger-executor.service.ts:89uses onlyworkspaceIdandtriggerName. AlthoughconnectionIdis now passed torunPoll()and available inparams, it is not included in the lock key. This allows multiple active connections for the same workspace/trigger to contend for the same lock and be skipped during concurrent poll cycles.🔒 Required fix
-const lockKey = `lock:poll:${params.workspaceId}:${params.triggerName}`; +const lockKey = `lock:poll:${params.workspaceId}:${params.connectionId}:${params.triggerName}`;🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/api/src/modules/trigger/poller.service.ts` around lines 123 - 132, The polling lock key in TriggerExecutorService currently uses only workspaceId and triggerName, causing concurrent connections to collide; update the lock key construction inside TriggerExecutorService (the method that builds/acquires the lock for runPoll) to include the connectionId from the runPoll params (use params.connectionId or connectionId) so the lock is scoped per connection. Locate the lock acquisition code referenced around the lock-key creation (the function invoked by runPoll) and append connectionId to the key string/tuple so each connection gets a distinct lock while leaving workspaceId and triggerName intact.apps/worker/src/modules/pipeline/normalization.service.ts (1)
130-151:⚠️ Potential issue | 🟠 MajorFall back when a custom normalizer returns no result.
If
customNormalizerexists but returnsnull/undefinedfor an unsupported entity, the currentelse ifskipspiece.normalizeand stores the record asRAW. Preserve the existing piece fallback unless the custom normalizer actually produced a normalized record.🐛 Proposed fix
- if (customNormalizer) { - const normalized = customNormalizer({ + const customNormalized = customNormalizer?.({ entityType: replica.entityType, data: replica.data as Record<string, unknown>, - }); - if (normalized) { - canonicalType = normalized.canonicalType; - canonicalData = normalized.data; - } - } else if (piece.normalize) { - const normalized = await piece.normalize( - replica.entityType, - replica.data as Record<string, unknown>, - ); + }); + + const normalized = + customNormalized ?? + (piece.normalize + ? await piece.normalize( + replica.entityType, + replica.data as Record<string, unknown>, + ) + : undefined); + + if (normalized) { + canonicalType = normalized.canonicalType; + canonicalData = normalized.data; - if (normalized) { - canonicalType = normalized.canonicalType; - canonicalData = normalized.data; - } }🤖 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 130 - 151, The current branch uses `else if` so when `getNormalizer(connectionAppName, appProfile)` returns a function that yields no result, `piece.normalize` is skipped and the record is stored as RAW; change the flow in `normalization.service.ts` to call `customNormalizer` first (via `getNormalizer`), capture its return in a `normalized` variable, and only if that `normalized` is null/undefined and `piece.normalize` exists, call `await piece.normalize(replica.entityType, replica.data)`; then set `canonicalType` and `canonicalData` from whichever `normalized` produced a value (referring to `getNormalizer`, `customNormalizer`, `piece.normalize`, `canonicalType`, `canonicalData`, and `replica` to locate the logic).apps/worker/src/modules/pipeline/replica.service.ts (1)
57-102: 🧹 Nitpick | 🔵 TrivialMove
appConnectionslookup out of the transaction.The
this.db.select(...)at Lines 83–90 executes inside thedb.transaction(async (tx) => …)callback but usesthis.dbrather thantx. This pulls a second connection from the pool while the transaction's connection is still held, increasing contention and risking pool exhaustion under load. Since the lookup isconnectionId-scoped and independent of the tenant schema'sSET LOCAL search_path, hoist it above the transaction:♻️ Proposed refactor
+ // Fetch application metadata (independent of tenant schema, run outside tx) + const connRows = await this.db + .select({ + appName: appConnections.appName, + metadata: appConnections.metadata, + }) + .from(appConnections) + .where(eq(appConnections.id, connectionId)) + .limit(1); + + const appName = connRows[0]?.appName; + if (!appName) { + throw new Error( + `Connection ${connectionId} not found in appConnections!`, + ); + } + const metadata = connRows[0]?.metadata as + | Record<string, unknown> + | undefined; + const appProfile = (metadata?.appProfile as string) || "default"; + await this.db.transaction(async (tx) => { ... - // Fetch application metadata - const connRows = await this.db - .select({ ... }) - ... - if (!appName) { throw ... }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/worker/src/modules/pipeline/replica.service.ts` around lines 57 - 102, The appConnections lookup (the this.db.select(...) that selects appName and metadata from appConnections using connectionId) is executed inside the db.transaction callback but uses this.db, which pulls a second pool connection; move that entire lookup out of the transaction so it runs before awaiting this.db.transaction(...), then pass the derived values (appName, metadata, appProfile, and the existing existence check for connectionId) into the transaction block; alternatively, if you prefer to keep it inside the transaction, change this.db.select(...) to tx.select(...) so it uses the transaction connection—update references to appConnections, connectionId, appName, metadata, and appProfile accordingly and remove the redundant this.db.select from inside the transaction.apps/api/src/modules/trigger/trigger-executor.service.ts (1)
415-433:⚠️ Potential issue | 🔴 Critical
pushToDlqdropsconnectionId— DLQ retries will fail with undefined schema.
DlqJob(dlq-processor.service.ts Line 29) requiresconnectionId, and the DLQ processor propagates it intoTriggerRunParams.connectionId(Line 243 of that file), which is then consumed byexecuteAndIngestviastorageResolver.resolveSchemaName(params.connectionId)(Line 271 here). However, thispushToDlqserializer omitsconnectionIdentirely, so every first-time DLQ job will deserialize withconnectionId === undefined. On the first retryresolveSchemaName(undefined)will throw, and the DLQ layer will only increment the attempt counter until the job is moved todlq:triggers:failed— without ever actually retrying the work.🐛 Proposed fix
private async pushToDlq( params: TriggerRunParams | WebhookRunParams, err: unknown, ): Promise<void> { const job = JSON.stringify({ appName: params.appName, triggerName: params.triggerName, workspaceId: params.workspaceId, + connectionId: params.connectionId, objectType: params.objectType, propsValue: params.propsValue, auth: params.auth, // required for credential reconstruction on retry // Preserve webhook payload for retry (may contain remaining unprocessed records) payload: 'payload' in params ? params.payload : undefined, failedAt: new Date().toISOString(), error: err instanceof Error ? err.message : String(err), attempt: 1, }); await this.redis.lpush('dlq:triggers', job); }A test in
trigger-executor.service.spec.tsasserting that the serialized DLQ job includesconnectionIdwould prevent regressions.🤖 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 415 - 433, pushToDlq currently omits connectionId from the serialized DLQ job which causes DlqJob retries to have connectionId === undefined and fail when storageResolver.resolveSchemaName(params.connectionId) is called; fix by including params.connectionId in the JSON payload created in pushToDlq so the DLQ processor can reconstruct TriggerRunParams.connectionId on retry, and add a unit test in trigger-executor.service.spec.ts asserting the serialized DLQ job contains connectionId to prevent regressions.
♻️ Duplicate comments (1)
apps/api/src/modules/pipeline/replica-outbox.service.spec.ts (1)
67-103: 🧹 Nitpick | 🔵 TrivialSame weak assertions as
inbound-outbox.service.spec.ts.This file mirrors the inbound-outbox spec almost verbatim. The same assertion gaps apply: retry/permanent-failure tests only check
db.updatewas called (true on the happy path too), and the graceful-rejection test ends withexpect(true).toBe(true). Please apply the same tightening here — captureset(...)arguments and assertstatus: 'RETRY'/status: 'FAIL', and for the rejection test assert the non-rejecting workspace still delivered viaqueueService.sendwithQueueName.ReplicaQueue.Given the near-identical shape of the two specs, consider extracting a shared test helper (e.g.,
runOutboxServiceContract) parameterized by the service class and expected queue name, to avoid further drift between the two suites.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/api/src/modules/pipeline/replica-outbox.service.spec.ts` around lines 67 - 103, Update the three tests in replica-outbox.service.spec.ts to assert the actual status payloads instead of only that db.update was called: when simulating a retry in processOutbox/processOutboxRow capture the arguments passed to the mock chain (the object returned by db.transaction's update/set calls) and assert set(...) was called with { status: 'RETRY' } (and similarly assert { status: 'FAIL' } for the max-attempts case where returning shows attempts: 6); for the "rejecting drainWorkspace" test, instead of expect(true) assert that the non-failing workspace still invoked queueService.send with QueueName.ReplicaQueue and the expected message shape; finally, factor these common assertions into a shared helper (e.g., runOutboxServiceContract) parameterized by the service under test and the expected queue name to avoid duplication between inbound and replica specs.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In @.agent/docs/architecture/nexiom-architecture.md:
- Around line 20-21: The TOC entry "Custom Mapping — Configuration vs GitOps
Sharding" does not match the actual heading "Custom Mapping — The Security &
Scale Boundary"; update one to match the other so the anchor resolves: either
rename the TOC line to "Custom Mapping — The Security & Scale Boundary" (and
adjust its anchor to `#custom-mapping--the-security--scale-boundary`) or rename
the heading to "Custom Mapping — Configuration vs GitOps Sharding" so the
existing TOC anchor (`#11-custom-mapping--configuration-vs-gitops-sharding`) is
correct; ensure the visible text and generated markdown anchor for the TOC entry
and the heading are identical.
In `@apps/api/src/modules/pipeline/inbound-outbox.service.spec.ts`:
- Around line 67-103: The tests for processOutbox/processOutboxRow are too
weak—change assertions to validate actual state transitions and behavior: in the
retry test (where queueService.send is mocked to reject) assert that the
transaction's set/update was called with an object containing status: 'RETRY'
(inspect the mocked tx.set or outer db.update().set calls), in the
permanent-failure test assert that tx.set/db.update().set was called with
status: 'FAIL' and that queueService.send was not invoked for the row with
attempts: 6, and in the graceful-rejection test replace the tautology with a
real check that the second workspace progressed (e.g., assert queueService.send
was called with the expected traceId or spy on the injected logger and assert an
error was logged) so failures in
processOutbox/drainWorkspaceOutbox/processOutboxRow are detected.
In `@apps/api/src/modules/pipeline/inbound-outbox.service.ts`:
- Around line 65-84: The current claim logic in
tx.update(inboundOutbox).set(...) increments attempts at claim time which counts
restarts as delivery attempts; either add a separate claim counter or change
semantics: (A) Add a new column claim_attempts (or claim_count), increment
claim_attempts in the claim SQL block (the tx.update(...) that sets status =
'PROCESSING'), keep attempts untouched there, and move the attempts increment
into processOutboxRow when a real delivery fails; update all references to
attempts/nextRetryAt and any queries filtering by attempts; or (B) if you prefer
the existing behavior, rename attempts to claim_attempts across the codebase and
increase MAX_ATTEMPTS accordingly (update constant MAX_ATTEMPTS and any
metrics/alerts that read attempts); ensure the RETURNING and WHERE clauses, the
processOutboxRow function, and any monitoring/alerting logic are updated to use
the new names/semantics.
- Around line 26-54: processOutbox (the `@Cron` handler) can re-enter and run
overlapping executions; add a re-entrancy guard so a new tick returns early
while one is running (either a simple in-memory boolean like isProcessingOutbox
checked at the top of processOutbox and reset in a finally block, or register
the job with SchedulerRegistry and skip if already active), ensure the guard
covers the entire body including the Promise.allSettled call and calls to
drainWorkspaceOutbox, and additionally modify drainWorkspaceOutbox to bound
concurrent outbound sends (replace unconstrained Promise.allSettled over
queueService.send with a concurrency limiter such as p-limit or a semaphore to
cap parallel queueService.send invocations).
In `@apps/api/src/modules/trace/data-explorer.controller.spec.ts`:
- Around line 40-44: The test currently asserts synchronous throws for the async
controller method listInbound; change the assertion to await
expect(controller.listInbound(mockNoOrgCtx, 's1', 1,
10)).rejects.toThrow(BadRequestException) so the rejected Promise is asserted
correctly, and apply the same pattern to the sibling list endpoints (e.g.,
listOutbound, listSomethingElse) that validate org presence to ensure all
organization-guard tests use await expect(...).rejects.toThrow().
In `@apps/api/src/modules/trace/data-explorer.controller.ts`:
- Around line 19-21: Remove the redundant `@Inject` decorator in the controller
constructor: rely on NestJS's automatic type-based injection by keeping the
parameter signature private readonly explorer: DataExplorerService in the
constructor of the class that currently uses `@Inject`(DataExplorerService), i.e.
update the constructor that defines explorer to remove the `@Inject` import and
decorator usage and leave only the typed parameter to clean up the code.
- Around line 30-113: Add UUID validation for workspaceId and consolidate the
five near-identical handlers: update each handler signature (listInbound,
listReplica, listNormalized, listEntityMap, listOutbound) to use
`@Query`('workspaceId', new ParseUUIDPipe({ optional: true })) workspaceId?:
string so bad UUIDs return 400 instead of silently causing a NotFoundException
in resolveStitch; then replace the five individual `@Get` routes with a single
generic endpoint (e.g. `@Get`(':tab') listByTab) that whitelists allowed tab names
and dispatches to the corresponding explorer methods (explorer.listInbound,
explorer.listReplica, explorer.listNormalized, explorer.listEntityMap,
explorer.listOutbound) or use a small dispatch map to call the correct method,
preserving requireOrg(ctx), stitchId, page, limit, workspaceId parameters.
In `@apps/api/src/modules/trace/data-explorer.service.spec.ts`:
- Around line 97-105: The test captures originalSelect but never uses it,
leaving tx.select mocked without a proper fallthrough which breaks composability
(e.g., listInbound expects chained select behavior); either remove the unused
originalSelect captures or change the mock to use them as the fallback:
bind/capture originalSelect = tx.select (or tx.select.bind(tx)) then in the
vi.fn().mockImplementation for tx.select return the special count branch when
args.count is present and otherwise call/return originalSelect(...) (preserving
mockReturnThis/chaining semantics); also remove the file-level no-unused-vars
pragma if you delete the captures so linter warnings are no longer suppressed.
In `@apps/api/src/modules/trace/data-explorer.service.ts`:
- Around line 154-190: listNormalized currently queries the tenant schema's
normalizedEntity without scoping to the stitch/connection, leaking
cross-connection data; check the normalizedEntity model for a
connectionId/stitchId/sourceConnectionId column and, if present, add a WHERE
clause in listNormalized (the transaction block that sets search_path and builds
the tx.select().from(normalizedEntity) query) to filter by the appropriate
column (e.g., connectionId = stitch.srcConnectionId or stitchId = stitchId)
similar to listReplica and listOutbound; if no such column exists, add a
comment/docstring to listNormalized explaining that L3 is tenant-global and this
endpoint intentionally returns rows from all source connections.
- Around line 81-106: This transaction block needlessly sets the local
search_path and repeats assertValidSchemaName even though buildTenantSchema() /
pgSchema() produces schema-qualified table references; remove the
tx.execute(sql`SET LOCAL search_path ...`) call(s) and eliminate duplicate
assertValidSchemaName calls in the functions using db.transaction (e.g., the
block containing inboundGateway queries) so queries rely on Drizzle's
schema-qualified identifiers, and run/adjust any callers or tests that assumed
the search_path side-effect.
In `@apps/api/src/modules/trigger/dlq-processor.service.ts`:
- Line 29: DlqJob.connectionId is now required but
TriggerExecutorService.pushToDlq is not serializing connectionId for first-time
failures, causing deserialized DLQ jobs to have connectionId === undefined and
later fail when runPoll calls executeAndIngest which uses
StorageResolverService.resolveSchemaName; update
TriggerExecutorService.pushToDlq to include the current connectionId in the
serialized DLQ payload (the same value used when creating DlqJob) so that new
DLQ entries always contain connectionId and subsequent runPoll retries receive a
valid connectionId.
In `@apps/api/src/modules/trigger/trigger-executor.service.ts`:
- Around line 380-408: Move the schema validation so
assertValidSchemaName(schemaName) runs before any schema-derived handles are
created (i.e., call assertValidSchemaName before buildTenantSchema or at the
very start of the transaction) to prevent constructing
inboundGateway/inboundOutbox with an invalid identifier; also make the
inboundOutbox insert defensive by adding an onConflictDoNothing with the same
conflict target used in replica.service (target: [traceId, connectionId]) so the
outbox insert cannot raise a conflict error and roll back the gateway insert.
In `@apps/api/src/scripts/test-e2e-ingestion.ts`:
- Around line 83-109: Replace the local ws_<id> derivation with the canonical
resolver and validate the schema before interpolating into SQL: import and call
StorageResolverService (or the exported function that returns the tenant schema)
from `@nexiom/database` to produce the schemaName instead of manually building
'ws_' + conn.workspace_id..., then call assertValidSchemaName(schemaName) to
fail-fast on invalid values and only use the validated schemaName in the
pool.query calls that currently interpolate ${schemaName}.
- Around line 1-11: The script hardcodes a local DB and localhost webhook and
unconditionally mutates connection metadata; update the run() startup to refuse
execution against production unless an explicit confirmation flag (e.g. --yes)
is provided or NODE_ENV !== 'production', accept host/connection id and
DATABASE_URL via CLI args or env vars instead of falling back to a localhost
literal (replace usage of process.env.DATABASE_URL ||
'postgres://postgres:postgres@localhost:5432/nexiom_local'), and before
executing the UPDATE app_connection ... metadata or calling
http://localhost:3000/v1/webhooks/${conn.id} print a clear, interactive warning
(or require the --yes flag) showing target DB URL and conn.id and bail if not
confirmed; follow these changes inside the run() function and around the Pool
creation and the code that performs the UPDATE to ensure no accidental
production mutation.
- Around line 77-81: Replace the fixed 5s sleep (the console.log and await new
Promise((r) => setTimeout(r, 5000))) with a polling loop that queries the
pipeline state (e.g., check for presence/processed status in replica_entity
and/or normalized_entity) at a short interval (e.g., 250–500ms) and stops as
soon as the expected entity appears or a global timeout is reached (e.g., 30s);
use the same DB client or service calls already used in this script to fetch
replica_entity/normalized_entity, log progress each poll if helpful, and throw
an error when the timeout elapses so the test fails instead of being flaky.
In `@apps/web/src/modules/trace/api/data-explorer.api.ts`:
- Around line 66-119: The buildQuery function contains dead URLSearchParams and
callers duplicate request-building logic; remove the unused q creation from
buildQuery and replace the five near-identical functions (listInbound,
listReplica, listNormalized, listEntityMap, listOutbound) with a single generic
helper (e.g., listExplorer<T>(stitchId: string, params: Params = {}, segment:
string)): inside that helper build a URLSearchParams using strict checks
(params.page !== undefined, params.limit !== undefined, params.workspaceId) and
call apiClient.get<ExplorerPage<T>> with
`${encodeURIComponent(stitchId)}/explorer/${segment}?${q}` (or keep buildQuery
to return only the base path and let the helper append the segment and query);
update existing callers to call listExplorer with the appropriate segment and
generic type.
In `@apps/web/src/modules/trace/pages/TraceExplorerPage.tsx`:
- Around line 66-102: The DataTable currently computes columns from
Object.keys(rows[0]) which drops keys that only appear in later rows; update
DataTable to derive the full column union (or accept an explicit columns prop) —
e.g., compute a Set by iterating all rows and collecting Object.keys for each
row in DataTable, then convert to an array (preserving first-seen order) and use
that array instead of Object.keys(rows[0]) when rendering headers and cells
(references: function DataTable, variable columns, prop rows, and the
headers/cell rendering logic).
- Around line 41-60: JsonCell currently calls JSON.stringify multiple times and
will throw on circular references; fix it by computing a safe, memoized
serialization once using useMemo inside JsonCell (dependent on value) and
replace the raw JSON.stringify calls with the memoized results; implement a
safeStringify helper (e.g., use a replacer that tracks seen objects with a Set
or use a small try/catch that falls back to a placeholder like '[Circular]') to
produce both a pretty-printed str and a compact preview, and ensure any errors
are caught so the component renders a non-throwing fallback (e.g., '—' or
'[Unserializable]') instead of crashing the row.
- Around line 205-213: The current useEffect loader swallows errors when calling
listStitches in load(), leaving StitchSelector empty and giving no feedback;
update the load function to catch and handle errors by logging the error (e.g.,
console.error or app logger) and setting an error state (e.g., stitchesError)
via a new state hook, setStitches([]) or keep prior state, and expose that error
to the UI so StitchSelector can render an inline error message and a retry
action (call load again) while still calling setLoading(false) in finally;
modify references to load, listStitches, setStitches, setLoading, and the
StitchSelector render to read stitchesError and show a retry button.
In `@apps/web/src/shared/components/layout/WorkspaceExplorer.tsx`:
- Around line 116-140: The Trace link's active check currently uses
location.pathname.startsWith(`${wsHref}/trace`) which will also match routes
like `${wsHref}/traceability`; change it to use the same exact-or-descendant
pattern as Stitches: define a traceHref (`const traceHref = `${wsHref}/trace``)
and an isTraceActive boolean (`location.pathname === traceHref ||
location.pathname.startsWith(`${traceHref}/`)`) then use isTraceActive for the
aria-current and className on the Trace <Link> inside WorkspaceExplorer so only
the exact trace route or its subpaths mark it active.
In `@apps/worker/src/modules/pipeline/replica.service.ts`:
- Around line 104-121: The ternary fallbacks for resolvedEntityType/resolvedData
are dead because extracted is always populated (either by the inline fallback
when extractor is undefined or by a non-null extractor result, with extractor &&
!extracted already throwing); simplify by assigning resolvedEntityType =
extracted.entityType and resolvedData = extracted.data directly. Update the code
around getReplicaExtractor, extractor, extracted, resolvedEntityType, and
resolvedData to remove the redundant ternaries and any unnecessary null checks,
keeping the existing extractor && !extracted throw logic intact.
In `@packages/pieces/application/revenova/src/index.ts`:
- Around line 5-8: The initializer initializeRevenovaApplicationRegistry
currently registers handlers (via registerReplicaExtractor and
registerNormalizer) but is never invoked, so getReplicaExtractor/getNormalizer
lookups for 'salesforce:revenova' fail; fix by invoking
initializeRevenovaApplicationRegistry at application startup—either call it from
the central startup/bootstrap sequence (e.g., inside the main initialization
function that configures registries) or invoke it as a module-side effect (call
it at the end of the same module) so the registrations are executed before any
getReplicaExtractor/getNormalizer calls.
In `@packages/pieces/application/revenova/src/upsertRevenovaObject.ts`:
- Around line 8-29: In upsertRevenovaObject.ts, the current rename only checks
lowercase r_obj.name so PascalCase Salesforce fields (e.g., Name) never get
renamed; update the logic after the sanitization loop that handles the name
rename to check both r_obj.name and r_obj.Name, assign r_obj.name_ = r_obj.name
?? r_obj.Name, then delete both r_obj.name and r_obj.Name; while here, tighten
typing for rawObj/value after unwrapping notification.sobject (replace payload
as any/value as any with a narrower interface or guarded casts) to avoid loose
any usage.
In `@packages/pieces/application/revenova/src/upsertTMSObject.ts`:
- Line 38: Replace the 'unknown' string fallback for sourceId with a null (or
throw) so distinct vendor records without Id/id don't collapse; in the mapping
that sets sourceId (the line using sourceId: String(replica.data.Id ||
replica.data.id || 'unknown')), return null when neither replica.data.Id nor
replica.data.id exists (or throw an explicit error) and ensure you don't coerce
null into the string "null" (i.e., remove the String(...) wrapper when
preserving null); update any affected callers of upsertTMSObject/normalizer to
accept null per the normalizer contract.
- Around line 6-22: Change the metadataDictionary entry type so the "type"
property is typed as the CanonicalType union instead of string and remove the
runtime cast (the `as CanonicalType` usage) where NormalizedRecord is built;
update the declaration of metadataDictionary to Record<string, { type:
CanonicalType; fields: Record<string,string> }>, then update any code that
previously relied on the cast (e.g., the place that constructs a
NormalizedRecord or uses canonicalType) to rely on the now-typed value so typos
like 'TMS_CARIER' fail at compile time.
In `@packages/pieces/platform/framework/src/normalizer.ts`:
- Around line 6-36: The registries (extractorRegistry, normalizerRegistry)
silently overwrite duplicate registrations and the string key built by
buildKey(appName, appProfile) is ambiguous when inputs contain ":"; update
registerReplicaExtractor and registerNormalizer to detect existing entries for
the given appName/appProfile and emit a dev warning (e.g., console.warn or
processLogger.warn) when an overwrite would occur, and change the underlying
storage from a single Map<string,Fn> to a nested Map<string, Map<string,Fn>> (or
use a non-printable delimiter like '\u0001' in buildKey) so keys are
unambiguous; update getReplicaExtractor and getNormalizer to read from the new
nested structure (or the new delimiter scheme) so lookups remain correct.
---
Outside diff comments:
In `@apps/api/src/modules/trigger/poller.service.ts`:
- Around line 63-74: The keyset cursor currently uses { created_at, workspace_id
} which is not unique; update the cursor to include the connection id (ac.id) as
a final tie-breaker and propagate that through
fetchActivePollingConnectionsBatch and the loop (the cursor variable and
where/predicate logic), change the ordering to include ac.id as the last
ordering key, and adjust the next-page predicate to compare the three-field
tuple (created_at, workspace_id, id) so rows with the same created_at and
workspace_id aren’t skipped.
- Around line 123-132: The polling lock key in TriggerExecutorService currently
uses only workspaceId and triggerName, causing concurrent connections to
collide; update the lock key construction inside TriggerExecutorService (the
method that builds/acquires the lock for runPoll) to include the connectionId
from the runPoll params (use params.connectionId or connectionId) so the lock is
scoped per connection. Locate the lock acquisition code referenced around the
lock-key creation (the function invoked by runPoll) and append connectionId to
the key string/tuple so each connection gets a distinct lock while leaving
workspaceId and triggerName intact.
In `@apps/api/src/modules/trigger/trigger-executor.service.ts`:
- Around line 415-433: pushToDlq currently omits connectionId from the
serialized DLQ job which causes DlqJob retries to have connectionId ===
undefined and fail when storageResolver.resolveSchemaName(params.connectionId)
is called; fix by including params.connectionId in the JSON payload created in
pushToDlq so the DLQ processor can reconstruct TriggerRunParams.connectionId on
retry, and add a unit test in trigger-executor.service.spec.ts asserting the
serialized DLQ job contains connectionId to prevent regressions.
In `@apps/worker/src/modules/pipeline/normalization.service.ts`:
- Around line 130-151: The current branch uses `else if` so when
`getNormalizer(connectionAppName, appProfile)` returns a function that yields no
result, `piece.normalize` is skipped and the record is stored as RAW; change the
flow in `normalization.service.ts` to call `customNormalizer` first (via
`getNormalizer`), capture its return in a `normalized` variable, and only if
that `normalized` is null/undefined and `piece.normalize` exists, call `await
piece.normalize(replica.entityType, replica.data)`; then set `canonicalType` and
`canonicalData` from whichever `normalized` produced a value (referring to
`getNormalizer`, `customNormalizer`, `piece.normalize`, `canonicalType`,
`canonicalData`, and `replica` to locate the logic).
In `@apps/worker/src/modules/pipeline/replica.service.ts`:
- Around line 57-102: The appConnections lookup (the this.db.select(...) that
selects appName and metadata from appConnections using connectionId) is executed
inside the db.transaction callback but uses this.db, which pulls a second pool
connection; move that entire lookup out of the transaction so it runs before
awaiting this.db.transaction(...), then pass the derived values (appName,
metadata, appProfile, and the existing existence check for connectionId) into
the transaction block; alternatively, if you prefer to keep it inside the
transaction, change this.db.select(...) to tx.select(...) so it uses the
transaction connection—update references to appConnections, connectionId,
appName, metadata, and appProfile accordingly and remove the redundant
this.db.select from inside the transaction.
---
Duplicate comments:
In `@apps/api/src/modules/pipeline/replica-outbox.service.spec.ts`:
- Around line 67-103: Update the three tests in replica-outbox.service.spec.ts
to assert the actual status payloads instead of only that db.update was called:
when simulating a retry in processOutbox/processOutboxRow capture the arguments
passed to the mock chain (the object returned by db.transaction's update/set
calls) and assert set(...) was called with { status: 'RETRY' } (and similarly
assert { status: 'FAIL' } for the max-attempts case where returning shows
attempts: 6); for the "rejecting drainWorkspace" test, instead of expect(true)
assert that the non-failing workspace still invoked queueService.send with
QueueName.ReplicaQueue and the expected message shape; finally, factor these
common assertions into a shared helper (e.g., runOutboxServiceContract)
parameterized by the service under test and the expected queue name to avoid
duplication between inbound and replica specs.
🪄 Autofix (Beta)
✅ Autofix completed
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: de6a99b0-2034-4b9a-bbf6-ad8a6a37b636
⛔ Files ignored due to path filters (1)
pnpm-lock.yamlis excluded by!**/pnpm-lock.yaml
📒 Files selected for processing (124)
.agent/docs/architecture/nexiom-architecture.mdapps/api/src/modules/pipeline/inbound-outbox.service.spec.tsapps/api/src/modules/pipeline/inbound-outbox.service.tsapps/api/src/modules/pipeline/pipeline.module.tsapps/api/src/modules/pipeline/replica-outbox.service.spec.tsapps/api/src/modules/trace/data-explorer.controller.spec.tsapps/api/src/modules/trace/data-explorer.controller.tsapps/api/src/modules/trace/data-explorer.service.spec.tsapps/api/src/modules/trace/data-explorer.service.tsapps/api/src/modules/trace/trace.module.tsapps/api/src/modules/trigger/dlq-processor.service.tsapps/api/src/modules/trigger/poller.service.tsapps/api/src/modules/trigger/trigger-executor.service.spec.tsapps/api/src/modules/trigger/trigger-executor.service.tsapps/api/src/modules/trigger/trigger.module.tsapps/api/src/modules/trigger/webhooks.controller.tsapps/api/src/scripts/test-e2e-ingestion.tsapps/web/src/app/routes/TenantRoutes.tsxapps/web/src/modules/trace/api/data-explorer.api.tsapps/web/src/modules/trace/pages/TraceExplorerPage.tsxapps/web/src/shared/components/layout/WorkspaceExplorer.tsxapps/worker/src/modules/pipeline/normalization.service.tsapps/worker/src/modules/pipeline/replica.service.spec.tsapps/worker/src/modules/pipeline/replica.service.tspackages/database/src/schema/pipeline.tspackages/pieces/application/revenova/package.jsonpackages/pieces/application/revenova/src/index.tspackages/pieces/application/revenova/src/upsertRevenovaObject.tspackages/pieces/application/revenova/src/upsertTMSObject.tspackages/pieces/application/revenova/tsconfig.jsonpackages/pieces/application/revenova/tsconfig.lib.jsonpackages/pieces/platform/framework/package.jsonpackages/pieces/platform/framework/src/action.tspackages/pieces/platform/framework/src/auth.spec.tspackages/pieces/platform/framework/src/auth.tspackages/pieces/platform/framework/src/canonical/index.tspackages/pieces/platform/framework/src/discovery/igt-logger.tspackages/pieces/platform/framework/src/discovery/index.tspackages/pieces/platform/framework/src/discovery/interfaces.tspackages/pieces/platform/framework/src/discovery/optimization-registry.tspackages/pieces/platform/framework/src/discovery/smart-cursor-selector.tspackages/pieces/platform/framework/src/discovery/universal-trigger.spec.tspackages/pieces/platform/framework/src/discovery/universal-trigger.tspackages/pieces/platform/framework/src/http-client.tspackages/pieces/platform/framework/src/index.tspackages/pieces/platform/framework/src/normalizer.tspackages/pieces/platform/framework/src/piece.tspackages/pieces/platform/framework/src/property.tspackages/pieces/platform/framework/src/retryable-exception.tspackages/pieces/platform/framework/src/trigger.tspackages/pieces/platform/framework/tsconfig.jsonpackages/pieces/platform/framework/tsconfig.spec.jsonpackages/pieces/platform/framework/vitest.config.tspackages/pieces/platform/quickbooks/.eslintrc.jsonpackages/pieces/platform/quickbooks/README.mdpackages/pieces/platform/quickbooks/openapi.jsonpackages/pieces/platform/quickbooks/package.jsonpackages/pieces/platform/quickbooks/project.jsonpackages/pieces/platform/quickbooks/src/i18n/de.jsonpackages/pieces/platform/quickbooks/src/i18n/es.jsonpackages/pieces/platform/quickbooks/src/i18n/fr.jsonpackages/pieces/platform/quickbooks/src/i18n/ja.jsonpackages/pieces/platform/quickbooks/src/i18n/nl.jsonpackages/pieces/platform/quickbooks/src/i18n/pt.jsonpackages/pieces/platform/quickbooks/src/i18n/ru.jsonpackages/pieces/platform/quickbooks/src/i18n/translation.jsonpackages/pieces/platform/quickbooks/src/i18n/vi.jsonpackages/pieces/platform/quickbooks/src/i18n/zh.jsonpackages/pieces/platform/quickbooks/src/index.tspackages/pieces/platform/quickbooks/src/lib/auth.tspackages/pieces/platform/quickbooks/src/lib/common.tspackages/pieces/platform/quickbooks/src/lib/types.tspackages/pieces/platform/quickbooks/src/triggers/quickbooks-polling.helper.tspackages/pieces/platform/quickbooks/src/triggers/quickbooks-query.adapter.spec.tspackages/pieces/platform/quickbooks/src/triggers/quickbooks-query.adapter.tspackages/pieces/platform/quickbooks/src/triggers/universal-trigger.tspackages/pieces/platform/quickbooks/tsconfig.jsonpackages/pieces/platform/quickbooks/tsconfig.lib.jsonpackages/pieces/platform/registry/package.jsonpackages/pieces/platform/registry/src/index.tspackages/pieces/platform/registry/src/metadata/metadata-discovery.service.spec.tspackages/pieces/platform/registry/src/metadata/metadata-discovery.service.tspackages/pieces/platform/registry/src/metadata/metadata.module.tspackages/pieces/platform/registry/src/pieces/piece-loader.service.tspackages/pieces/platform/registry/src/pieces/piece-registry.service.spec.tspackages/pieces/platform/registry/src/pieces/piece-registry.service.tspackages/pieces/platform/registry/src/pieces/pieces.module.tspackages/pieces/platform/registry/tsconfig.jsonpackages/pieces/platform/salesforce/.eslintrc.jsonpackages/pieces/platform/salesforce/README.mdpackages/pieces/platform/salesforce/openapi.jsonpackages/pieces/platform/salesforce/package.jsonpackages/pieces/platform/salesforce/project.jsonpackages/pieces/platform/salesforce/src/i18n/ca.jsonpackages/pieces/platform/salesforce/src/i18n/de.jsonpackages/pieces/platform/salesforce/src/i18n/es.jsonpackages/pieces/platform/salesforce/src/i18n/fr.jsonpackages/pieces/platform/salesforce/src/i18n/hi.jsonpackages/pieces/platform/salesforce/src/i18n/id.jsonpackages/pieces/platform/salesforce/src/i18n/ja.jsonpackages/pieces/platform/salesforce/src/i18n/nl.jsonpackages/pieces/platform/salesforce/src/i18n/pt.jsonpackages/pieces/platform/salesforce/src/i18n/ru.jsonpackages/pieces/platform/salesforce/src/i18n/translation.jsonpackages/pieces/platform/salesforce/src/i18n/vi.jsonpackages/pieces/platform/salesforce/src/i18n/zh.jsonpackages/pieces/platform/salesforce/src/index.tspackages/pieces/platform/salesforce/src/lib/auth.tspackages/pieces/platform/salesforce/src/lib/common/index.tspackages/pieces/platform/salesforce/src/lib/discovery/index.tspackages/pieces/platform/salesforce/src/lib/discovery/salesforce-bulk.adapter.spec.tspackages/pieces/platform/salesforce/src/lib/discovery/salesforce-bulk.adapter.tspackages/pieces/platform/salesforce/src/lib/discovery/salesforce-discovery.adapter.spec.tspackages/pieces/platform/salesforce/src/lib/discovery/salesforce-discovery.adapter.tspackages/pieces/platform/salesforce/src/lib/discovery/salesforce-query.adapter.spec.tspackages/pieces/platform/salesforce/src/lib/discovery/salesforce-query.adapter.tspackages/pieces/platform/salesforce/src/lib/salesforce-types.tspackages/pieces/platform/salesforce/src/lib/sf-fetch.tspackages/pieces/platform/salesforce/src/lib/trigger/index.tspackages/pieces/platform/salesforce/src/lib/trigger/salesforce-polling.helper.tspackages/pieces/platform/salesforce/src/lib/trigger/universal-trigger.tspackages/pieces/platform/salesforce/tsconfig.jsonpackages/pieces/platform/salesforce/tsconfig.lib.jsonpnpm-workspace.yaml
| return tx | ||
| .update(inboundOutbox) | ||
| .set({ | ||
| status: 'PROCESSING', | ||
| attempts: sql`${inboundOutbox.attempts} + 1`, | ||
| nextRetryAt: sql`NOW() + INTERVAL '5 minutes'`, | ||
| }) | ||
| .where( | ||
| sql`${inboundOutbox.id} IN ( | ||
| SELECT id FROM ${sql.identifier(schemaName)}.inbound_outbox | ||
| WHERE status = 'PENDING' | ||
| OR (status = 'RETRY' AND next_retry_at <= NOW()) | ||
| OR (status = 'PROCESSING' AND next_retry_at <= NOW()) | ||
| ORDER BY next_retry_at ASC | ||
| LIMIT ${BATCH_SIZE} | ||
| FOR UPDATE SKIP LOCKED | ||
| )`, | ||
| ) | ||
| .returning(); | ||
| }); |
There was a problem hiding this comment.
attempts is incremented on every claim, not every delivery — a mid-flight crash burns an attempt.
Line 69 bumps attempts at claim time. If the process crashes (or the SUCCESS/RETRY update on lines 117–143 fails) after the claim commits, the row sits in PROCESSING with nextRetryAt = NOW()+5min, then gets re-claimed and attempts increments again — even though zero delivery attempts happened in that interval. With MAX_ATTEMPTS=6 you still have headroom, but any downstream on-call alert tied to "attempts" will be misleading and a run of unlucky restarts could mark a row FAIL without ever making 6 real delivery attempts.
Two options:
- Track claims and deliveries separately (
claim_countvsattempts), incrementingattemptsonly insideprocessOutboxRowon failure. - Leave as-is but rename to
claim_attemptsand raiseMAX_ATTEMPTSto reflect the new semantics.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@apps/api/src/modules/pipeline/inbound-outbox.service.ts` around lines 65 -
84, The current claim logic in tx.update(inboundOutbox).set(...) increments
attempts at claim time which counts restarts as delivery attempts; either add a
separate claim counter or change semantics: (A) Add a new column claim_attempts
(or claim_count), increment claim_attempts in the claim SQL block (the
tx.update(...) that sets status = 'PROCESSING'), keep attempts untouched there,
and move the attempts increment into processOutboxRow when a real delivery
fails; update all references to attempts/nextRetryAt and any queries filtering
by attempts; or (B) if you prefer the existing behavior, rename attempts to
claim_attempts across the codebase and increase MAX_ATTEMPTS accordingly (update
constant MAX_ATTEMPTS and any metrics/alerts that read attempts); ensure the
RETURNING and WHERE clauses, the processOutboxRow function, and any
monitoring/alerting logic are updated to use the new names/semantics.
|
Note Autofix is a beta feature. Expect some limitations and changes as we gather feedback and continue to improve it. Fixes Applied SuccessfullyFixed 17 file(s) based on 26 unresolved review comments. Files modified:
Commit: The changes have been pushed to the Time taken: |
Fixed 17 file(s) based on 26 unresolved review comments. Co-authored-by: CodeRabbit <noreply@coderabbit.ai>
There was a problem hiding this comment.
Actionable comments posted: 15
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
apps/api/src/modules/trigger/trigger-executor.service.ts (1)
89-90:⚠️ Potential issue | 🔴 CriticalScope locks, deduplication, and cursors by
connectionIdto prevent cross-connection collisions.
connectionIdis stored in the gateway table, but poll/webhook locks,sourceEventIddeduplication, and Redis cursor keys omit it. Multiple connections in the same workspace with the same trigger can collide, causing one delivery to be skipped by the lock, another byonConflictDoNothingon duplicateextReqId, and both to share the samelast_cursorposition.Proposed scoping fix
- const lockKey = `lock:poll:${params.workspaceId}:${params.triggerName}`; + const lockKey = `lock:poll:${params.workspaceId}:${params.connectionId}:${params.triggerName}`; @@ - const lockKey = `lock:webhook:${params.workspaceId}:${params.triggerName}:${bodyHash}`; + const lockKey = `lock:webhook:${params.workspaceId}:${params.connectionId}:${params.triggerName}:${bodyHash}`; @@ const sourceEventId = this.buildSourceEventId( params.workspaceId, + params.connectionId, params.triggerName, record, ); @@ private buildSourceEventId( workspaceId: string, + connectionId: string, triggerName: string, record: unknown, ): string { @@ return createHash('sha256') - .update(`${workspaceId}:${triggerName}:${bounded}`) + .update(`${workspaceId}:${connectionId}:${triggerName}:${bounded}`) .digest('hex'); }Also verify
RedisBackedTriggerStorecursor keys includeconnectionId; the current patterncursor:${workspaceId}:${appName}:${objectType}:${triggerName}can be shared across connections.🤖 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 89 - 90, The lock and dedupe scope must include connectionId to avoid cross-connection collisions: update the lock key construction used before acquireLock (currently building lockKey with params.workspaceId and params.triggerName) to append params.connectionId (or fetch connectionId from the gateway row if not in params); ensure any DB deduplication that relies on extReqId (the upsert/onConflictDoNothing path) namespaces extReqId with connectionId so duplicates are per-connection; and update RedisBackedTriggerStore cursor keys (currently cursor:${workspaceId}:${appName}:${objectType}:${triggerName}) to include connectionId so last_cursor is connection-scoped. Touch the code paths around trigger-executor.service.ts where lockKey is created and acquireLock is called, the DB insert/upsert that uses extReqId/onConflictDoNothing, and RedisBackedTriggerStore methods that build cursor keys to include connectionId.
♻️ Duplicate comments (2)
packages/pieces/application/revenova/src/upsertTMSObject.ts (1)
36-42:⚠️ Potential issue | 🟠 MajorReject blank vendor IDs before returning a normalized record.
?? nullstill treats""as a valid ID, so malformed records with blank IDs can share the samesourceId. Returnnullfor blank IDs so the platform falls back instead of canonicalizing under a colliding key.🛡️ Proposed fix
- const sourceId = replica.data.Id ?? replica.data.id ?? null; + const sourceId = replica.data.Id ?? replica.data.id; + if ( + sourceId === undefined || + sourceId === null || + (typeof sourceId === 'string' && sourceId.trim() === '') + ) { + return null; + } return { canonicalType: meta.type, - sourceId: sourceId !== null ? String(sourceId) : undefined, + sourceId: String(sourceId), data: canonicalFields, };🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@packages/pieces/application/revenova/src/upsertTMSObject.ts` around lines 36 - 42, The current sourceId assignment (using replica.data.Id ?? replica.data.id ?? null) treats empty strings as valid IDs; change logic in upsertTMSObject (the sourceId variable) to treat blank or whitespace-only values as null by checking the chosen value (replica.data.Id or replica.data.id), trimming it (or testing length) and returning null when it’s empty/whitespace-only so that the returned sourceId becomes undefined fallback rather than a colliding empty string key.apps/api/src/scripts/test-e2e-ingestion.ts (1)
94-99:⚠️ Potential issue | 🟡 MinorResolve the schema-derivation TODO before relying on raw schema interpolation.
This still duplicates tenant schema derivation and interpolates
schemaNameinto SQL without canonical validation; use the resolver/validator before querying tenant tables.I can help generate the resolver-backed version if you want to track it as a follow-up.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/api/src/scripts/test-e2e-ingestion.ts` around lines 94 - 99, Replace the manual schema derivation that builds schemaName from conn.workspace_id with the canonical resolver/validator: import and call the StorageResolverService (or at minimum assertValidSchemaName) from `@nexiom/database` to derive/validate the tenant schema, then use that validated value when constructing queries instead of the raw 'ws_'+conn.workspace_id.replace(...) string; ensure the validated schemaName is used for interpolation and that you handle/throw errors from the resolver so no unvalidated schema ever reaches SQL execution.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In @.agent/docs/architecture/nexiom-architecture.md:
- Line 28: Fix markdownlint issues in nexiom-architecture.md by (1) adding
language tags to all fenced code blocks (use appropriate languages like text,
typescript, sql, json, or bash), (2) wrapping all occurrences of literal
identifiers such as billing_* and BillAddr.* in backticks so emphasis/parsing is
correct, and (3) ensuring the document ends with a single newline; apply these
changes to every code fence and every instance of billing_* / BillAddr.*
referenced throughout the file (including the other reported locations).
- Around line 239-260: Update the docs to match the current API: the normalizer
registry is keyed by (appName, appProfile) and the worker reads
metadata?.appProfile, so change the text that says "(appProfile, objectType)" /
"metadata.app_profile" / "getNormalizer(appProfile, objectType)" to reflect the
actual signatures; update the example registerNormalizer calls to the real API
(include appName and appProfile ordering) and show getNormalizer being called
with appName and appProfile (and objectType) to retrieve the normalizer (e.g.,
reference registerNormalizer, getNormalizer, and metadata?.appProfile and the
normalize* function names to locate the code to change).
- Around line 72-79: Update the tenant schema inventory to include the newly
added tenant-scoped outbox: add "inbound_outbox" to the list under the
ws_sf_abc123 schema so it reads (for example) "inbound_gateway, inbound_outbox,
replica_entity, normalized_entity, outbound_gateway, sync_log / sync_cursor,
replica_outbox / normalized_outbox / delivery_outbox", ensuring the document's
inventory matches the implemented L1→L2 outbox transition.
In `@apps/api/src/modules/pipeline/inbound-outbox.service.ts`:
- Around line 153-170: The code writes a lastError field into inboundOutbox rows
but the inbound_outbox schema lacks that column; add a nullable text column
"lastError" to the inbound_outbox table in the schema (so the TypeScript model
inboundOutbox includes lastError) and create a DB migration to add the column,
or alternatively remove the lastError assignments in the update calls inside
inbound-outbox.service.ts (the update block that sets { status: 'FAIL',
lastError } and the block that sets { status: 'RETRY', lastError, nextRetryAt
}); ensure inboundOutbox's TypeScript definition, the DB schema, and migrations
stay consistent with whichever approach you choose.
- Line 3: Update the update queries that change a row from PROCESSING to
RETRY/FAIL so they include the claim-specific predicate to avoid stomping newer
claims: add the claimed attempt identifier and the expected current status to
the WHERE clause (e.g., AND claim_attempt_id = :claimedAttemptId AND status =
'PROCESSING' or match the claimed status variable) for the UPDATE statements
that set status = 'FAIL' and status = 'RETRY' (and any other status transitions
performed in inbound-outbox.service.ts), so only the worker that holds the exact
claim can change that row.
- Around line 108-122: The code discards rejected results from processWithLimit
causing drainWorkspaceOutbox to treat failures as success; update
processWithLimit (and/or drainWorkspaceOutbox) to collect Promise.allSettled
results for each chunk, filter for entries with status === "rejected", and
surface them: log per-row failures with identifying info and error (from
processOutboxRow), and if any rejections remain decide/implement a non-silent
outcome (e.g., rethrow an aggregated error or return failure metadata) so the
tenant drain does not silently succeed; reference processWithLimit,
processOutboxRow, claimed, and drainWorkspaceOutbox when making the change.
In `@apps/api/src/modules/trace/data-explorer.controller.spec.ts`:
- Around line 40-114: Tests call non-existent controller methods
(listInbound/listReplica/listNormalized/listEntityMap/listOutbound); update each
to call the single exposed method listByTab(...) and pass the correct tab string
so the controller routes to the service. Replace controller.listInbound(mockCtx,
...) etc. with controller.listByTab(mockCtx, 'inbound', 's1', 1, 10, 'ws1') and
similarly use 'replica', 'normalized', 'entity-map', 'outbound' for the other
specs, and keep assertions that the service
(service.listInbound/listReplica/...) was called with the org id and same args.
In `@apps/api/src/modules/trace/data-explorer.service.ts`:
- Around line 57-61: The safePagination helper should defensively normalize its
inputs: coerce page and limit to numeric values, guard against NaN and Infinity
(using Number() and Number.isFinite), convert to integers (e.g., Math.trunc) and
then apply the existing bounds (Math.max(1, ...) and Math.min(..., MAX_LIMIT))
before computing offset; update the safePagination function to perform these
checks on the incoming page and limit values so calls that bypass controller
validation remain safe.
In `@apps/api/src/scripts/test-e2e-ingestion.ts`:
- Around line 58-65: The e2e uses a static payload variable and a poll that only
checks “latest row exists,” so reruns can pass using old rows; add a per-run
unique marker to the payload (e.g., payload.TestRunId or payload.__testRunId set
to a UUID/timestamp) and capture a runStart timestamp before sending; then
update the poll/query logic that asserts ingestion (the code that checks for the
latest L2/L3 rows) to filter for rows containing that TestRunId or rows created
after runStart so the assertion proves this run produced the webhook-derived
rows.
- Around line 78-84: The e2e script is logging webhook failures and letting
run() resolve which makes CI treat failures as success; update the failure
branches (the response.ok check inside test-e2e-ingestion.ts and any
timeout/fatal error handlers referenced around run()) to terminate with a
non-zero exit code instead of just logging—e.g., after logging the error in the
response.ok false branch and in the timeout/fatal catch paths, call
process.exit(1) (or rethrow so the caller exits non-zero) so the process fails
CI when webhook ingestion or timeouts occur.
- Around line 14-15: The code logs the full DATABASE_URL (variable databaseUrl)
which may expose credentials; update the logging to sanitize userinfo before
printing: parse process.env.DATABASE_URL (or databaseUrl) with the URL
constructor or a parser, replace or remove url.username and url.password (e.g.,
set to 'REDACTED' or omit userinfo), then log the reconstructed sanitizedUrl
instead of databaseUrl (keep the console.log call but reference the sanitized
value) so credentials are never emitted to logs.
In `@apps/web/src/modules/trace/pages/TraceExplorerPage.tsx`:
- Around line 41-53: safeStringify is using a single Set named seen for both
JSON.stringify calls, so the first pretty stringify marks objects as visited and
the subsequent compact stringify treats them as circular; fix by using a fresh
circular-reference Set for each stringify pass (e.g., create a factory function
for the replacer or reinitialize seen before creating compact) so that replacer
used for pretty and replacer used for compact each close over their own Set;
update references to replacer, seen, pretty and compact inside safeStringify
accordingly.
- Around line 235-265: The Retry button currently only sets loading and clears
error but does not re-run the async load logic inside the useEffect; extract the
async load logic into a stable function (e.g., move load out of the inline
useEffect or wrap it in useCallback like load = useCallback(async () => { ... },
[workspaceId])) that calls listStitches and updates
setStitches/setError/setLoading, then call that same load function from the
Retry button onClick (invoke load() after setLoading(true)/setError(null) or
have load set loading itself) so clicking Retry actually reloads stitches for
the current workspaceId.
In `@apps/worker/src/modules/pipeline/replica.service.ts`:
- Around line 93-104: Replace the unsafe cast of metadata?.appProfile with a
runtime check so only a real string is used as the registry key: e.g. compute
appProfile by checking typeof metadata?.appProfile === "string" &&
metadata.appProfile.trim() !== "" ? metadata.appProfile : undefined, then pass
that validated value into getReplicaExtractor(appName, appProfile) (do the same
change in normalization.service.ts). This prevents non-string values from being
treated as valid keys and avoids silently falling back to incorrect defaults.
In `@packages/pieces/application/revenova/src/upsertRevenovaObject.ts`:
- Around line 9-20: The upsertRevenovaObject extractor currently type-casts
nested Salesforce envelopes without validating the shape, allowing malformed
payloads to fall through and produce a DEFAULT replica; update
upsertRevenovaObject to validate that envelope.notification and
envelope.notification.sobject are objects (and that required fields exist)
before assigning rawObj, and if the envelope is present but malformed return
null (so ReplicaService fails the bad webhook) — specifically add runtime checks
around envelope.notification and envelope.notification.sobject used in the
upsertRevenovaObject function (instead of just casting) and ensure rawObj is
only set from envelope.notification.sobject when it passes those checks.
---
Outside diff comments:
In `@apps/api/src/modules/trigger/trigger-executor.service.ts`:
- Around line 89-90: The lock and dedupe scope must include connectionId to
avoid cross-connection collisions: update the lock key construction used before
acquireLock (currently building lockKey with params.workspaceId and
params.triggerName) to append params.connectionId (or fetch connectionId from
the gateway row if not in params); ensure any DB deduplication that relies on
extReqId (the upsert/onConflictDoNothing path) namespaces extReqId with
connectionId so duplicates are per-connection; and update
RedisBackedTriggerStore cursor keys (currently
cursor:${workspaceId}:${appName}:${objectType}:${triggerName}) to include
connectionId so last_cursor is connection-scoped. Touch the code paths around
trigger-executor.service.ts where lockKey is created and acquireLock is called,
the DB insert/upsert that uses extReqId/onConflictDoNothing, and
RedisBackedTriggerStore methods that build cursor keys to include connectionId.
---
Duplicate comments:
In `@apps/api/src/scripts/test-e2e-ingestion.ts`:
- Around line 94-99: Replace the manual schema derivation that builds schemaName
from conn.workspace_id with the canonical resolver/validator: import and call
the StorageResolverService (or at minimum assertValidSchemaName) from
`@nexiom/database` to derive/validate the tenant schema, then use that validated
value when constructing queries instead of the raw
'ws_'+conn.workspace_id.replace(...) string; ensure the validated schemaName is
used for interpolation and that you handle/throw errors from the resolver so no
unvalidated schema ever reaches SQL execution.
In `@packages/pieces/application/revenova/src/upsertTMSObject.ts`:
- Around line 36-42: The current sourceId assignment (using replica.data.Id ??
replica.data.id ?? null) treats empty strings as valid IDs; change logic in
upsertTMSObject (the sourceId variable) to treat blank or whitespace-only values
as null by checking the chosen value (replica.data.Id or replica.data.id),
trimming it (or testing length) and returning null when it’s
empty/whitespace-only so that the returned sourceId becomes undefined fallback
rather than a colliding empty string key.
🪄 Autofix (Beta)
✅ Autofix completed
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: 6cc07da5-f19c-448c-8236-22945b96734a
📒 Files selected for processing (17)
.agent/docs/architecture/nexiom-architecture.mdapps/api/src/modules/pipeline/inbound-outbox.service.spec.tsapps/api/src/modules/pipeline/inbound-outbox.service.tsapps/api/src/modules/trace/data-explorer.controller.spec.tsapps/api/src/modules/trace/data-explorer.controller.tsapps/api/src/modules/trace/data-explorer.service.spec.tsapps/api/src/modules/trace/data-explorer.service.tsapps/api/src/modules/trigger/trigger-executor.service.tsapps/api/src/scripts/test-e2e-ingestion.tsapps/web/src/modules/trace/api/data-explorer.api.tsapps/web/src/modules/trace/pages/TraceExplorerPage.tsxapps/web/src/shared/components/layout/WorkspaceExplorer.tsxapps/worker/src/modules/pipeline/replica.service.tspackages/pieces/application/revenova/src/index.tspackages/pieces/application/revenova/src/upsertRevenovaObject.tspackages/pieces/application/revenova/src/upsertTMSObject.tspackages/pieces/platform/framework/src/normalizer.ts
|
Note Autofix is a beta feature. Expect some limitations and changes as we gather feedback and continue to improve it. Fixes Applied SuccessfullyFixed 10 file(s) based on 15 unresolved review comments. Files modified:
Commit: The changes have been pushed to the Time taken: |
Fixed 10 file(s) based on 15 unresolved review comments. Co-authored-by: CodeRabbit <noreply@coderabbit.ai>
There was a problem hiding this comment.
Actionable comments posted: 7
♻️ Duplicate comments (2)
.agent/docs/architecture/nexiom-architecture.md (2)
240-260:⚠️ Potential issue | 🟡 MinorAlign this section with the current normalizer registry API.
The implementation in
normalizer.tsexposesregisterNormalizer(appName, appProfile, fn)andgetNormalizer(appName, appProfile), but this doc still shows an object-type overload. Update the examples so contributors do not copy calls that TypeScript will reject.🧪 Read-only verification of the documented API mismatch
#!/bin/bash # Description: Compare the documented normalizer calls with the exported framework function signatures. # Expected: Docs should not contain 4-argument registerNormalizer or 3-argument getNormalizer examples # unless the framework exports those overloads. rg -n -C2 'export function (registerNormalizer|getNormalizer)\s*\(' packages/pieces/platform/framework/src/normalizer.ts rg -n -C2 'registerNormalizer\(|getNormalizer\(' .agent/docs/architecture/nexiom-architecture.md🛠️ Proposed fix
-1. **One function per object type** — no if-chains. Each object has its own normalizer file. +1. **One app-profile normalizer entry point** — application packages can fan out internally by `entityType`. 2. **Registry keyed by (appName, appProfile)** — NOT (platform, objectType). Two companies can use Salesforce with different data models. 3. **`appProfile` comes from `app_connection.metadata?.appProfile`** — set when the customer registers their connection (e.g. `"revenova"`). @@ -registerNormalizer('salesforce', 'revenova', 'Account', normalizeAccount); -registerNormalizer('salesforce', 'revenova', 'TransportationProfile__c', normalizeTransportationProfile); -registerNormalizer('salesforce', 'revenova', 'VendorInvoice__c', normalizeVendorInvoice); -registerNormalizer('salesforce', 'revenova', 'Invoice__c', normalizeCustomerInvoice); +registerNormalizer('salesforce', 'revenova', upsertTMSObject); @@ -const normalizer = getNormalizer(appName, appProfile, objectType); -const result = normalizer?.(rawData) ?? null; +const normalizer = getNormalizer(appName, appProfile); +const result = normalizer?.({ entityType, data: rawData }) ?? null;🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In @.agent/docs/architecture/nexiom-architecture.md around lines 240 - 260, The docs show a 4-arg registerNormalizer and 3-arg getNormalizer but the real API is registerNormalizer(appName, appProfile, fn) and getNormalizer(appName, appProfile); update the examples in this document to call registerNormalizer with (appName, appProfile, fn) where the fn returns/contains the object-type dispatch (or registers per-object inside its closure) and to call getNormalizer(appName, appProfile) (two args) and then lookup by objectType from the returned normalizer map or dispatch function; update the code block examples (the registration snippet and the worker call) to use the actual registerNormalizer and getNormalizer signatures and remove the extra objectType argument so they compile against the normalizer.ts API.
571-571:⚠️ Potential issue | 🟡 MinorEnd the document with a single trailing newline.
markdownlintreports MD047 here; add exactly one final newline to keep the doc lint-clean.🧹 Proposed fix
-*Outcome:* Mathematical security traps the untrusted code while maintaining microsecond native latency. +*Outcome:* Mathematical security traps the untrusted code while maintaining microsecond native latency. +🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In @.agent/docs/architecture/nexiom-architecture.md at line 571, The document currently ends without a trailing newline after the line "Mathematical security traps the untrusted code while maintaining microsecond native latency."; add exactly one final newline character at the end of the file so the document ends with a single trailing newline to satisfy markdownlint MD047.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In @.agent/docs/architecture/nexiom-architecture.md:
- Around line 112-147: The fenced TypeScript examples (TMS_CARRIER,
TMS_CUSTOMER, TMS_TRANSPORTATION_PROFILE) use backticked property keys like
`billing_street` which are invalid in TS; edit the fenced code blocks to remove
the backticks so properties read billing_street: string, billing_city: string,
etc., for all occurrences across those three type examples and ensure comments
and notes remain unchanged.
In `@apps/api/src/modules/pipeline/inbound-outbox.service.ts`:
- Around line 142-146: The terminal throw in drainWorkspaceOutbox after
Promise.allSettled causes a noisy generic error to bubble up even when per-row
processing already logged detailed failures; update drainWorkspaceOutbox to stop
throwing that generic Error and instead either (A) remove the throw so the
per-row logs (from the loop over rejections) are the signal, or (B) if you must
throw, wrap the thrown Error with the actual rejection details (include the
rejections array / per-row identifiers and messages) so processOutbox receives
actionable context; locate drainWorkspaceOutbox and processOutboxRow to
implement the chosen approach.
- Line 3: The import list in inbound-outbox.service.ts includes an unused symbol
`eq` (imported from 'drizzle-orm') which breaks CI; remove `eq` from the import
statement (leaving `notInArray` and `sql`) or use `eq` where intended, by
editing the import on the top of the file to only import the actually used
helpers (`notInArray`, `sql`) so the `@typescript-eslint/no-unused-vars` error
is resolved.
- Around line 87-200: The code references inboundOutbox.claimAttemptId in claim
and stale-claim guards (inside the claim logic and processOutboxRow) but the
inbound_outbox Drizzle schema lacks that column; either add a claimAttemptId
column and DB migration to packages/database/src/schema/pipeline (and update any
types) so the SQL conditions and update(...).set(...) calls using claimAttemptId
compile and the DB has the column, or remove all claimAttemptId usage (in the
claim update, the sql IN (...) selector, and the WHERE clauses inside
processOutboxRow) and replace the stale-claim guards with checks based only on
(id, status, attempts) so the WHEREs become sql`${inboundOutbox.id} = ${row.id}
AND ${inboundOutbox.status} = 'PROCESSING'` (and equivalent for retry/fail
paths); update types for the row param in processOutboxRow accordingly.
In `@apps/api/src/modules/trace/data-explorer.service.ts`:
- Around line 90-111: The four methods listInbound, listReplica, listNormalized,
and listOutbound currently wrap two concurrent read-only queries in
this.db.transaction using Promise.all(tx.select(...), tx.select(...)); remove
the transaction wrapper and run both selects directly on this.db (e.g.,
Promise.all([this.db.select()... , this.db.select()...])) so the queries remain
concurrent but without transactional overhead, ensure you stop using the tx
variable, keep the same ordering/limit/offset and count query, and preserve the
returned shape (data, total, page, limit) and types (satisfies ExplorerPage) for
each method.
In `@apps/api/src/scripts/test-e2e-ingestion.ts`:
- Around line 128-149: The poll queries currently filter by created_at and a
top-level JSON key that doesn't exist after extraction, causing timeouts on
reruns; update the SELECTs in the polling loop (the two pool.query calls using
schemaName on replica_entity and normalized_entity) to use updated_at >= $1
instead of created_at >= $1, and change the replica_entity JSON test to check
the nested path preserved by upsertRevenovaObject (or better: embed __testRunId
at the top level of the webhook payload) so the condition (data->>'__testRunId'
= $2) actually matches; alternatively, ensure tests truncate or reset the
replica/normalized rows before running to avoid relying on created_at.
In `@apps/worker/src/modules/pipeline/replica.service.ts`:
- Around line 82-108: The appConnections lookup currently uses this.db inside
the transaction callback (the tx passed into the transaction started around the
tenant transaction) which makes the metadata read occur on a different
connection/transaction; either move the select that queries appConnections (the
block that populates connRows, appName, metadata and appProfile using
connectionId) to run before calling this.db.transaction(...) so it clearly
executes outside the tenant transaction and allows failing fast on missing
connection, or change the query to use
tx.select(...).from(appConnections).where(...) inside the transaction to make it
share the same isolation; also remove the duplicated lookup inside the
transaction callback so only one consistent read exists.
---
Duplicate comments:
In @.agent/docs/architecture/nexiom-architecture.md:
- Around line 240-260: The docs show a 4-arg registerNormalizer and 3-arg
getNormalizer but the real API is registerNormalizer(appName, appProfile, fn)
and getNormalizer(appName, appProfile); update the examples in this document to
call registerNormalizer with (appName, appProfile, fn) where the fn
returns/contains the object-type dispatch (or registers per-object inside its
closure) and to call getNormalizer(appName, appProfile) (two args) and then
lookup by objectType from the returned normalizer map or dispatch function;
update the code block examples (the registration snippet and the worker call) to
use the actual registerNormalizer and getNormalizer signatures and remove the
extra objectType argument so they compile against the normalizer.ts API.
- Line 571: The document currently ends without a trailing newline after the
line "Mathematical security traps the untrusted code while maintaining
microsecond native latency."; add exactly one final newline character at the end
of the file so the document ends with a single trailing newline to satisfy
markdownlint MD047.
🪄 Autofix (Beta)
✅ Autofix completed
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: 3b3f16e8-bc0f-44d0-8e3c-7299dc3ea33c
📒 Files selected for processing (10)
.agent/docs/architecture/nexiom-architecture.mdapps/api/src/modules/pipeline/inbound-outbox.service.tsapps/api/src/modules/trace/data-explorer.controller.spec.tsapps/api/src/modules/trace/data-explorer.service.tsapps/api/src/scripts/test-e2e-ingestion.tsapps/web/src/modules/trace/pages/TraceExplorerPage.tsxapps/worker/src/modules/pipeline/normalization.service.tsapps/worker/src/modules/pipeline/replica.service.tspackages/pieces/application/revenova/src/upsertRevenovaObject.tspackages/pieces/platform/framework/src/normalizer.ts
|
Note Autofix is a beta feature. Expect some limitations and changes as we gather feedback and continue to improve it. Fixes Applied SuccessfullyFixed 5 file(s) based on 5 unresolved review comments. Files modified:
Commit: The changes have been pushed to the Time taken: |
Fixed 5 file(s) based on 5 unresolved review comments. Co-authored-by: CodeRabbit <noreply@coderabbit.ai>
There was a problem hiding this comment.
Actionable comments posted: 6
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
apps/api/src/modules/trigger/trigger-executor.service.ts (1)
280-292:⚠️ Potential issue | 🔴 CriticalScope the dedup key by
connectionIdto avoid cross-connection drops.Line 400 deduplicates solely on
extReqId, butbuildSourceEventId()does not includeconnectionId. Two source connections in the same workspace with the sametriggerNameand equivalent record payload can produce the sameextReqId; the second insert becomesDO NOTHING, and no outbox row is created.🛡️ Proposed fix
const sourceEventId = this.buildSourceEventId( + params.connectionId, params.workspaceId, params.triggerName, record, ); @@ private buildSourceEventId( + connectionId: string, workspaceId: string, triggerName: string, record: unknown, ): string { @@ return createHash('sha256') - .update(`${workspaceId}:${triggerName}:${bounded}`) + .update(`${connectionId}:${workspaceId}:${triggerName}:${bounded}`) .digest('hex'); }Alternatively, make the database uniqueness and Drizzle conflict target composite on
(connectionId, extReqId)ifextReqIdmust remain source-scoped.Also applies to: 390-401, 490-532
🤖 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 280 - 292, The deduplication uses extReqId built by buildSourceEventId but that ID isn't scoped to connectionId, so inserts via insertGatewayRow can be dropped across different connections; include connectionId in the dedup key by changing buildSourceEventId (or the value passed to insertGatewayRow) to incorporate params.connectionId (or alternatively make the DB uniqueness/conflict target composite on (connectionId, extReqId)) so the INSERT ... ON CONFLICT/DO NOTHING is scoped per connection; update all usages (calls around insertGatewayRow and any DB unique index/conflict targets) so deduplication uses the composite (connectionId, extReqId).
♻️ Duplicate comments (2)
.agent/docs/architecture/nexiom-architecture.md (1)
151-156:⚠️ Potential issue | 🟡 MinorBackticks in TypeScript fenced block are still invalid syntax.
The previous fix cleaned up the TMS types block at Lines 112-147, but the ACCT block here retains backtick-quoted property keys inside a ```typescript fence — backticks are literal characters in TS and these are not valid identifiers.
🛠️ Proposed fix
-ACCT_VENDOR = { name, `billing_*`, phone, email, tax_id } -ACCT_CUSTOMER = { name, `billing_*`, phone, email, credit_limit } +ACCT_VENDOR = { name, billing_*, phone, email, tax_id } +ACCT_CUSTOMER = { name, billing_*, phone, email, credit_limit }Or convert the fence language to
textsince this is a schema sketch rather than real TS.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In @.agent/docs/architecture/nexiom-architecture.md around lines 151 - 156, The TypeScript fenced block contains invalid backtick-quoted property keys for ACCT_VENDOR, ACCT_CUSTOMER, ACCT_INVOICE and ACCT_BILL; remove the backticks around property names like `billing_*` (e.g., change to billing_* or billingFields) so the snippet is valid TypeScript-like syntax, or alternatively change the code fence language from ```typescript to ```text to indicate this is a schema sketch rather than real TS; update the ACCT_* lines accordingly to reflect either valid identifiers or plain-text schema.apps/api/src/scripts/test-e2e-ingestion.ts (1)
139-160:⚠️ Potential issue | 🟠 MajorL3 poll will still time out on reruns.
normalized_entityinserts useonConflictDoNothingonreplicaId, so on a 2nd run with the same sourceId: '0015Y00002bcdefGHI'no new row is created andupdated_atremains frozen at the first run's timestamp. The poll query at Line 145 filtersupdated_at >= runStart, so it will never match on reruns and the loop will time out.Additionally, for L2, the
__testRunIdmarker at Line 91 is nested insideAccount. Depending on howupsertRevenovaObjectunwraps the payload,data->>'__testRunId'may or may not resolve (if the extractor stores the innerAccountobject asdata, it would resolve at the top level — please verify). The OR fallback ondata->>'Id'masks this becauseupdated_atis refreshed on upsert for L2, but L3 has no such safety net.Consider one of:
- Making the source
Idunique per run (e.g., derived fromtestRunId) so both L2 and L3 insert fresh rows.- Cleaning up prior replica/normalized rows before the test.
- Promoting
__testRunIdto the top level of the webhook payload so it survives extraction intodataunambiguously.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/api/src/scripts/test-e2e-ingestion.ts` around lines 139 - 160, The L3 polling query on normalized_entity (used in the while loop around runStart/testRunId) will never match on reruns because normalized rows use onConflictDoNothing on replicaId so updated_at doesn't change; also data->>'__testRunId' may be nested under Account and not visible for replica_entity checks in upsertRevenovaObject. Fix by one of: (a) make the test source Id unique per run (derive Id from testRunId) so normalized_entity and replica_entity insert fresh rows; or (b) delete prior rows for the test Id from schemaName.replica_entity and schemaName.normalized_entity before running the poll; or (c) promote __testRunId to the top-level webhook payload so data->>'__testRunId' consistently exists after extraction by upsertRevenovaObject—choose one approach and update the test setup and/or payload construction and the while-loop queries accordingly.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In @.agent/docs/architecture/nexiom-architecture.md:
- Line 571: The file .agent/docs/architecture/nexiom-architecture.md is missing
a final newline after the last line ("*Outcome:* Mathematical security traps the
untrusted code while maintaining microsecond native latency."); update the file
so it ends with a single trailing newline character (ensure the last line is
terminated) to satisfy markdownlint MD047.
In `@apps/api/src/modules/pipeline/inbound-outbox.service.ts`:
- Around line 165-201: The updates to inboundOutbox (in the try success path and
the catch FAIL/RETRY paths) do not set the lastError column; capture the
computed errorMessage and persist it on failure/retry and clear it on success by
adding lastError: errorMessage for the FAIL and RETRY .set(...) calls and
lastError: null for the SUCCESS .set(...); update the .set calls on the
inboundOutbox updates that set status to 'SUCCESS', 'FAIL', and 'RETRY' (and
keep nextRetryAt logic intact) so operators can see error context for traceId
rows.
- Around line 91-99: The claim selector currently reclaims PROCESSING rows
purely by next_retry_at and can be overwritten by a late worker; modify the SQL
fragment that selects rows to claim (the sql`... IN ( SELECT id FROM
${sql.identifier(schemaName)}.inbound_outbox WHERE ... )` block) to include the
existing attempts value in the WHERE clause (e.g., AND attempts =
${row.attempts}) so a newer claim with a different attempts cannot be clobbered,
and ensure the claim/update that sets status/next_retry_at also increments/sets
attempts atomically; also update the finalizer/UPDATE logic that writes
status='PROCESSING' to include an attempts check (only apply the finalizer if
attempts still match) to make the late writer a no-op.
In `@apps/api/src/modules/trace/data-explorer.service.spec.ts`:
- Around line 86-106: The tests currently mock db.transaction
(db.transaction.mockImplementationOnce(...)) but the service now calls
this.db.select() directly, so those transaction callbacks never execute and
return empty results; replace each db.transaction.mockImplementationOnce block
with a direct select mock helper (use the suggested mockPageSelect helper) that
stubs this.db.select() to return rows and count pairs: e.g., use
mockPageSelect([{ id: 'norm_1' }], 3, false) for listNormalized (its count query
has no .where), and for other tests call mockPageSelect with the appropriate
rows/count and the third argument true/false depending on whether the count uses
a .where() path; update instances referenced in the tests (db.transaction,
listNormalized, mockPageSelect) accordingly.
In `@apps/worker/src/modules/pipeline/normalization.service.ts`:
- Around line 104-109: appProfile is validated with trim() but the untrimmed
value is returned, causing registry lookups like
getNormalizer(connectionAppName, appProfile) to miss keys; change the assignment
so you capture and use the trimmed string (e.g., compute const trimmedAppProfile
= metadata.appProfile.trim() and set appProfile = trimmedAppProfile when
non-empty) and ensure all subsequent calls (notably
getNormalizer(connectionAppName, appProfile)) use that trimmed value.
In `@apps/worker/src/modules/pipeline/replica.service.ts`:
- Around line 72-77: The current validation sets appProfile to
metadata.appProfile when non-empty but preserves surrounding whitespace; change
the assignment so the variable holds the normalized (trimmed) string instead of
the original metadata value and then pass that normalized appProfile into
getReplicaExtractor(appName, appProfile). In other words, when validating
metadata.appProfile, compute and store the trimmed result (e.g., trimmed
non-empty string) into the appProfile variable and ensure all subsequent calls
(notably getReplicaExtractor) use that normalized appProfile so registry lookups
use the exact key.
---
Outside diff comments:
In `@apps/api/src/modules/trigger/trigger-executor.service.ts`:
- Around line 280-292: The deduplication uses extReqId built by
buildSourceEventId but that ID isn't scoped to connectionId, so inserts via
insertGatewayRow can be dropped across different connections; include
connectionId in the dedup key by changing buildSourceEventId (or the value
passed to insertGatewayRow) to incorporate params.connectionId (or alternatively
make the DB uniqueness/conflict target composite on (connectionId, extReqId)) so
the INSERT ... ON CONFLICT/DO NOTHING is scoped per connection; update all
usages (calls around insertGatewayRow and any DB unique index/conflict targets)
so deduplication uses the composite (connectionId, extReqId).
---
Duplicate comments:
In @.agent/docs/architecture/nexiom-architecture.md:
- Around line 151-156: The TypeScript fenced block contains invalid
backtick-quoted property keys for ACCT_VENDOR, ACCT_CUSTOMER, ACCT_INVOICE and
ACCT_BILL; remove the backticks around property names like `billing_*` (e.g.,
change to billing_* or billingFields) so the snippet is valid TypeScript-like
syntax, or alternatively change the code fence language from ```typescript to
```text to indicate this is a schema sketch rather than real TS; update the
ACCT_* lines accordingly to reflect either valid identifiers or plain-text
schema.
In `@apps/api/src/scripts/test-e2e-ingestion.ts`:
- Around line 139-160: The L3 polling query on normalized_entity (used in the
while loop around runStart/testRunId) will never match on reruns because
normalized rows use onConflictDoNothing on replicaId so updated_at doesn't
change; also data->>'__testRunId' may be nested under Account and not visible
for replica_entity checks in upsertRevenovaObject. Fix by one of: (a) make the
test source Id unique per run (derive Id from testRunId) so normalized_entity
and replica_entity insert fresh rows; or (b) delete prior rows for the test Id
from schemaName.replica_entity and schemaName.normalized_entity before running
the poll; or (c) promote __testRunId to the top-level webhook payload so
data->>'__testRunId' consistently exists after extraction by
upsertRevenovaObject—choose one approach and update the test setup and/or
payload construction and the while-loop queries accordingly.
🪄 Autofix (Beta)
✅ Autofix completed
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: f3b58a65-76a7-498e-ade0-186566c5283f
📒 Files selected for processing (11)
.agent/docs/architecture/nexiom-architecture.mdapps/api/src/modules/pipeline/inbound-outbox.service.spec.tsapps/api/src/modules/pipeline/inbound-outbox.service.tsapps/api/src/modules/trace/data-explorer.controller.spec.tsapps/api/src/modules/trace/data-explorer.controller.tsapps/api/src/modules/trace/data-explorer.service.spec.tsapps/api/src/modules/trace/data-explorer.service.tsapps/api/src/modules/trigger/trigger-executor.service.tsapps/api/src/scripts/test-e2e-ingestion.tsapps/worker/src/modules/pipeline/normalization.service.tsapps/worker/src/modules/pipeline/replica.service.ts
|
|
||
| **Model B: V8 Isolates / WebAssembly (The Figma Pattern)** | ||
| If network latency is unacceptable, the physical tenant directories are compiled into standalone, dependency-free JavaScript bundles (or Wasm binaries) and uploaded to an S3 bucket. The platform worker securely evaluates the fetched code using a constrained `isolated-vm` memory trap inside the native worker process. | ||
| *Outcome:* Mathematical security traps the untrusted code while maintaining microsecond native latency. No newline at end of file |
There was a problem hiding this comment.
Add trailing newline.
markdownlint MD047 is flagging the missing final newline.
🧰 Tools
🪛 markdownlint-cli2 (0.22.0)
[warning] 571-571: Files should end with a single newline character
(MD047, single-trailing-newline)
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In @.agent/docs/architecture/nexiom-architecture.md at line 571, The file
.agent/docs/architecture/nexiom-architecture.md is missing a final newline after
the last line ("*Outcome:* Mathematical security traps the untrusted code while
maintaining microsecond native latency."); update the file so it ends with a
single trailing newline character (ensure the last line is terminated) to
satisfy markdownlint MD047.
|
Note Autofix is a beta feature. Expect some limitations and changes as we gather feedback and continue to improve it. Fixes Applied SuccessfullyFixed 5 file(s) based on 6 unresolved review comments. Files modified:
Commit: The changes have been pushed to the Time taken: |
Fixed 5 file(s) based on 6 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 (1)
apps/worker/src/modules/pipeline/normalization.service.ts (1)
137-158:⚠️ Potential issue | 🟡 MinorConfirm intended semantics when
customNormalizerreturns a falsy result.Current behavior: if
customNormalizeris registered and returnsnull/undefined, normalization silently keepscanonicalType = "RAW"/canonicalData = replica.dataand does not fall through topiece.normalize. This is asymmetric withreplica.service.ts(Lines 119–123), where an extractor returning falsy throws an explicit "envelope mismatch" error.Two failure modes worth considering:
- A bug in
customNormalizerthat returnsnullfor some entity types silently stores RAW rather than surfacing the problem, so downstream consumers ofnormalizedEntitycan't distinguish "intentional RAW passthrough" from "normalizer silently failed".- If the contract is "normalizer returning falsy means the entity type isn't handled by this app/profile," falling back to
piece.normalize(by restructuring to try the piece whencustomNormalizerreturns falsy) would preserve the prior behavior rather than masking it.Please confirm intent. If RAW-on-falsy is intentional, a brief comment documenting the contract would help; otherwise consider either throwing (matching replica) or falling through to
piece.normalize:♻️ Option A: fall through to piece.normalize on falsy custom result
- if (customNormalizer) { - const normalized = customNormalizer({ - entityType: replica.entityType, - data: replica.data as Record<string, unknown>, - }); - if (normalized) { - canonicalType = normalized.canonicalType; - canonicalData = normalized.data; - } - } else if (piece.normalize) { + let normalized = customNormalizer + ? customNormalizer({ + entityType: replica.entityType, + data: replica.data as Record<string, unknown>, + }) + : null; + if (!normalized && piece.normalize) { const normalized = await piece.normalize( replica.entityType, replica.data as Record<string, unknown>, ); - if (normalized) { - canonicalType = normalized.canonicalType; - canonicalData = normalized.data; - } + } + if (normalized) { + canonicalType = normalized.canonicalType; + canonicalData = normalized.data; }🤖 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 137 - 158, The current code short-circuits when customNormalizer exists but returns a falsy value, leaving canonicalType as RAW; decide and implement one of two fixes: either (A) fall through to piece.normalize when customNormalizer returns falsy by changing the if/else structure so that after calling customNormalizer you only assign canonicalType/canonicalData when normalized is truthy and otherwise continue to the piece.normalize branch (look for customNormalizer and piece.normalize in this function), or (B) make the behavior explicit by throwing an error on a falsy customNormalizer result (e.g., throw new Error("normalizer mismatch") or similar) to match the extractor behavior (see replica.service.ts extractor semantics) and add a brief comment documenting the chosen contract; update the code around canonicalType/canonicalData assignments accordingly.
♻️ Duplicate comments (1)
apps/api/src/modules/pipeline/inbound-outbox.service.ts (1)
180-198:⚠️ Potential issue | 🟠 MajorTruncate
lastErrorbefore writing tovarchar(500).
lastErroris defined asvarchar(500)inpackages/database/src/schema/pipeline.ts:223-237, but Lines 185 and 198 write the full error message. A long broker/DB error can make the RETRY/FAIL update fail, leaving the row stuck inPROCESSINGuntil it is reclaimed.🛡️ Proposed fix
} catch (err) { const errorMessage = err instanceof Error ? err.message : String(err); + const lastError = errorMessage.slice(0, 500); if (row.attempts >= MAX_ATTEMPTS) { await this.db .update(inboundOutbox) - .set({ status: 'FAIL', lastError: errorMessage }) + .set({ status: 'FAIL', lastError }) @@ await this.db .update(inboundOutbox) - .set({ status: 'RETRY', nextRetryAt, lastError: errorMessage }) + .set({ status: 'RETRY', nextRetryAt, lastError })🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/api/src/modules/pipeline/inbound-outbox.service.ts` around lines 180 - 198, The updates that set lastError in the inbound-outbox flow (see the errorMessage variable and the two this.db.update(...).set({ ..., lastError: errorMessage }) calls in the branch handling MAX_ATTEMPTS and the retry branch) can fail because lastError column is varchar(500); truncate errorMessage to 500 chars before the DB write. Modify code around errorMessage (used by inboundOutbox.update in the FAIL and RETRY paths) to compute a safeLastError = errorMessage.length > 500 ? errorMessage.slice(0, 497) + '...' : errorMessage and use safeLastError in the .set({ lastError: safeLastError }) values so long errors won’t cause the UPDATE to fail.
🤖 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/pipeline/inbound-outbox.service.ts`:
- Around line 91-103: The current predicate compares scalar inboundOutbox.id to
a subquery that returns two columns (id, attempts), which is invalid; change it
to a row-value comparison so both id and attempts are compared atomically:
replace the scalar IN clause with something like (inboundOutbox.id,
inboundOutbox.attempts) IN (SELECT id, attempts FROM
${sql.identifier(schemaName)}.inbound_outbox WHERE status = 'PENDING' OR (status
= 'RETRY' AND next_retry_at <= NOW()) OR (status = 'PROCESSING' AND
next_retry_at <= NOW()) ORDER BY next_retry_at ASC LIMIT ${BATCH_SIZE} FOR
UPDATE SKIP LOCKED), and remove the separate attempts subquery that selects
attempts FROM ... subq; keep references to inboundOutbox, schemaName, and
BATCH_SIZE so the change is applied to the same predicate.
In `@apps/api/src/modules/trace/data-explorer.service.spec.ts`:
- Line 1: Remove the broad file-level eslint-disable comment at the top of
data-explorer.service.spec.ts and replace it with targeted inline disables only
where necessary: run ESLint to find the exact offending lines and add `//
eslint-disable-next-line` with the specific rule(s) (e.g.
`@typescript-eslint/no-unsafe-assignment`,
`@typescript-eslint/no-unsafe-member-access`, `@typescript-eslint/no-unsafe-return`,
`@typescript-eslint/no-unsafe-call`, `@typescript-eslint/no-unused-vars`,
`@typescript-eslint/require-await`) immediately above the specific statements
(assignments, member accesses, returns, calls, unused test variables, or async
stubs) that trigger each rule; ensure no blanket disable remains and keep
comments narrowly scoped to the exact expressions or helper functions in this
spec file.
- Around line 163-180: The test for listEntityMap duplicates the custom
db.select mock which can drift from other tests; replace the ad-hoc db.select
implementation with the shared mockPageSelect helper to set up the paginated
count and results (instead of overriding db.select), then use the existing
db.offset.mockResolvedValueOnce([{ id: 'gem_1' }]) as the page payload;
reference the mockPageSelect helper to configure the count result and ensure the
test uses that helper rather than redefining db.select.
---
Outside diff comments:
In `@apps/worker/src/modules/pipeline/normalization.service.ts`:
- Around line 137-158: The current code short-circuits when customNormalizer
exists but returns a falsy value, leaving canonicalType as RAW; decide and
implement one of two fixes: either (A) fall through to piece.normalize when
customNormalizer returns falsy by changing the if/else structure so that after
calling customNormalizer you only assign canonicalType/canonicalData when
normalized is truthy and otherwise continue to the piece.normalize branch (look
for customNormalizer and piece.normalize in this function), or (B) make the
behavior explicit by throwing an error on a falsy customNormalizer result (e.g.,
throw new Error("normalizer mismatch") or similar) to match the extractor
behavior (see replica.service.ts extractor semantics) and add a brief comment
documenting the chosen contract; update the code around
canonicalType/canonicalData assignments accordingly.
---
Duplicate comments:
In `@apps/api/src/modules/pipeline/inbound-outbox.service.ts`:
- Around line 180-198: The updates that set lastError in the inbound-outbox flow
(see the errorMessage variable and the two this.db.update(...).set({ ...,
lastError: errorMessage }) calls in the branch handling MAX_ATTEMPTS and the
retry branch) can fail because lastError column is varchar(500); truncate
errorMessage to 500 chars before the DB write. Modify code around errorMessage
(used by inboundOutbox.update in the FAIL and RETRY paths) to compute a
safeLastError = errorMessage.length > 500 ? errorMessage.slice(0, 497) + '...' :
errorMessage and use safeLastError in the .set({ lastError: safeLastError })
values so long errors won’t cause the UPDATE to fail.
🪄 Autofix (Beta)
✅ Autofix completed
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: 0aa3ca30-1dae-4a97-976a-c35cb7d0e5bb
📒 Files selected for processing (4)
apps/api/src/modules/pipeline/inbound-outbox.service.tsapps/api/src/modules/trace/data-explorer.service.spec.tsapps/worker/src/modules/pipeline/normalization.service.tsapps/worker/src/modules/pipeline/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 2 file(s) based on 1 unresolved review comment. Files modified:
Commit: The changes have been pushed to the Time taken: |
Fixed 2 file(s) based on 1 unresolved review comment. Co-authored-by: CodeRabbit <noreply@coderabbit.ai>
Summary by CodeRabbit
New Features
Documentation
Tests