Skip to content

Scheduler foundation: Singer-compliant schema, poll framework, and ar… - #98

Merged
pramodnarayana merged 4 commits into
developmentfrom
feat/scheduler-foundation
Mar 23, 2026
Merged

pramodnarayana merged 4 commits into
developmentfrom
feat/scheduler-foundation

Conversation

@pramodnarayana

@pramodnarayana pramodnarayana commented Mar 23, 2026 •

Copy link
Copy Markdown
Owner

…chitecture spec

  • Add Singer-style sync_cursors schema keyed on (stitch_id, stream_name) instead of connection_id — prevents cursor collision between stitches sharing a source connection; state_document default includes bookmarks, versions, currently_syncing
  • Add dsProjectCode (nullable bigint) to organization table with CHECK constraint capping at JS MAX_SAFE_INTEGER to fail loudly on truncation
  • Add poll framework types to Piece interface: PollWindow, PollRecord, PollPage, StreamDescriptor (discriminated union enforcing replicationKey required on INCREMENTAL at compile time), ReplicationKeyType, poll(), describeStreams()
  • Add full DolphinScheduler + CursorManager architecture spec aligned with singer-python state.py conventions (currently_syncing, offset, versions at top level); covers DS tenant isolation, lifecycle, security, and edge cases
  • Update T029/T030/T046/T047/T048/T049 task descriptions to reflect corrected stitch_id key, DS_INTERNAL_SECRET auth, and Singer-compliant interfaces
  • Fix identity test fixture to include dsProjectCode: null

Summary by CodeRabbit

  • New Features

    • Stateful sync with cursor-based, pageable polling and explicit replication-key metadata for reliable incremental syncs
    • DolphinScheduler-backed scheduler for stitch lifecycle, retries, and orchestrated execution
  • Database

    • Added persistent sync-cursor storage and a nullable scheduler project identifier for organizations
  • Documentation

    • New scheduling architecture doc and updated task/phase plans
  • Tests

    • Updated test fixtures to include the new scheduler project field

…chitecture spec

- Add Singer-style sync_cursors schema keyed on (stitch_id, stream_name) instead
  of connection_id — prevents cursor collision between stitches sharing a source
  connection; state_document default includes bookmarks, versions, currently_syncing
- Add dsProjectCode (nullable bigint) to organization table with CHECK constraint
  capping at JS MAX_SAFE_INTEGER to fail loudly on truncation
- Add poll framework types to Piece interface: PollWindow, PollRecord, PollPage,
  StreamDescriptor (discriminated union enforcing replicationKey required on
  INCREMENTAL at compile time), ReplicationKeyType, poll(), describeStreams()
- Add full DolphinScheduler + CursorManager architecture spec aligned with
  singer-python state.py conventions (currently_syncing, offset, versions at
  top level); covers DS tenant isolation, lifecycle, security, and edge cases
- Update T029/T030/T046/T047/T048/T049 task descriptions to reflect corrected
  stitch_id key, DS_INTERNAL_SECRET auth, and Singer-compliant interfaces
- Fix identity test fixture to include dsProjectCode: null

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
@coderabbitai

coderabbitai Bot commented Mar 23, 2026 •

Copy link
Copy Markdown
Contributor

Warning

Rate limit exceeded

@pramodnarayana has exceeded the limit for the number of commits that can be reviewed per hour. Please wait 3 minutes and 32 seconds before requesting another review.

⌛ How to resolve this issue?

After the wait time has elapsed, a review can be triggered using the @coderabbitai review command as a PR comment. Alternatively, push new commits to this PR.

We recommend that you space out your commits to avoid hitting the rate limit.

🚦 How do rate limits work?

CodeRabbit enforces hourly rate limits for each developer per organization.

Our paid plans have higher rate limits than the trial, open-source and free plans. In all cases, we re-allow further reviews after a brief timeout.

Please see our FAQ for further information.

ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: 91c95375-f649-406e-9607-4b90285ff74c

📥 Commits

Reviewing files that changed from the base of the PR and between 7322977 and 2249eea.

📒 Files selected for processing (1)
  • docs/architecture/scheduling/scheduling.md
📝 Walkthrough

Walkthrough

Adds stateful cursor-based polling: new sync_cursors DB table and re-export, connector polling contract/types, DolphinScheduler-backed scheduler/execute-stitch flow and cursor manager design, organization dsProjectCode column, and related docs/tests updates.

Changes

Cohort / File(s) Summary
DB schema additions
packages/database/src/schema/stitches.ts, packages/database/src/schema/identity.ts
Added syncCursors (sync_cursors) table with JSONB stateDocument, FK to integrationStitches, unique (stitchId, streamName) and index. Added dsProjectCode bigint column to organization with safe-integer CHECK and partial unique index.
API re-exports
apps/api/src/db/schema.ts
Re-exported syncCursors from @nexiom/database to include the new table in centralized ORM exports.
Connector framework
packages/connectors/src/framework/piece.ts
Introduced polling types (ReplicationKeyType, PollWindow, PollRecord, PollPage, StreamDescriptor) and optional describeStreams / poll signatures on Piece and CreatePieceParams; createPiece conditionally attaches them.
Architecture docs
docs/architecture/master/tasks.md, docs/architecture/scheduling/scheduling.md
Documented DolphinScheduler-backed SchedulerService, execute-stitch HTTP callback flow, CursorManager semantics, Singer-style bookmark model, checkpointing, and updated task list/phases.
Tests / Fixtures
packages/identity/src/adapters/drizzle-tenant.adapter.spec.ts
Added dsProjectCode: null to mock schema.Organization fixtures to align with updated schema.

Sequence Diagram

sequenceDiagram
    participant DS as DolphinScheduler
    participant API as POST /internal/scheduler/execute-stitch
    participant CM as CursorManagerService
    participant Piece as Connector Piece
    participant DB as Database (sync_cursors)
    participant Redis as Redis Lock

    DS->>API: Trigger stitch (stitch_id)
    API->>API: Validate DS bearer secret
    API->>Redis: Attempt lock for stitch_id
    alt lock acquired
        API->>DB: Load stateDocument for stitch_id
        DB-->>API: stateDocument (bookmarks, versions)
        API->>CM: calculateWindow(stream, stateDocument)
        CM-->>API: PollWindow { lowerBound, upperBound, replicationKeyType }
        loop paginate
            API->>Piece: poll(credentials, streamName, window, nextPageCursor?)
            Piece-->>API: PollPage { records[], nextPageCursor? }
            API->>DB: Write intermediate checkpoint to sync_cursors
            API->>Downstream: write records
        end
        API->>DB: Final checkpoint commit (stateDocument)
        API->>DS: Return 200 { status: "SUCCESS" }
    else lock contention
        API->>DS: Return 200 { status: "SKIPPED" }
    end
Loading

Estimated code review effort

🎯 4 (Complex) | ⏱️ ~45 minutes

Possibly related PRs

Poem

🐇 I nibble bookmarks in tidy rows,

windows of time where replication grows,
DolphinScheduler hums, the cursor hops light,
checkpoints keep safe through day and night,
a rabbit cheers: stateful sync takes flight!

🚥 Pre-merge checks | ✅ 3
✅ Passed checks (3 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title clearly and concisely summarizes the main changes: adding Singer-compliant schema, poll framework, and architecture specification for the scheduler foundation.
Docstring Coverage ✅ Passed Docstring coverage is 100.00% which is sufficient. The required threshold is 80.00%.

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

✨ Finishing Touches
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch feat/scheduler-foundation

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

❤️ Share

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

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Actionable comments posted: 4

Caution

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

⚠️ Outside diff range comments (1)
packages/database/src/schema/identity.ts (1)

258-276: ⚠️ Potential issue | 🟠 Major

Add a uniqueness guard for dsProjectCode to preserve tenant isolation.

dsProjectCode is now a persisted org→DS mapping, but there is no unique constraint/index on this column. Two organizations can accidentally share one DS project code, which is a cross-tenant isolation risk.

🔧 Proposed fix
   (table) => [
     check("organization_id_not_sentinel", sql`${table.id} <> '__NULL__'`),
+    uniqueIndex("organization_ds_project_code_unique_idx")
+      .on(table.dsProjectCode)
+      .where(sql`${table.dsProjectCode} IS NOT NULL`),
     check(
       "ds_project_code_safe_integer",
       sql`${table.dsProjectCode} IS NULL OR ${table.dsProjectCode} <= 9007199254740991`,
     ),
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@packages/database/src/schema/identity.ts` around lines 258 - 276, Add a
uniqueness constraint/index for dsProjectCode to prevent multiple organizations
from sharing the same DolphinScheduler project code: modify the table definition
that declares dsProjectCode (identifier: dsProjectCode on the table variable) to
add a unique index (similar to organization_slug_unique_idx) that enforces
uniqueness only for non-null dsProjectCode values (a partial unique index WHERE
"dsProjectCode" IS NOT NULL) so NULLs remain allowed but any assigned project
code is unique across organizations; ensure the new index has a clear name
(e.g., organization_ds_project_code_unique_idx) and coexists with the existing
checks (organization_id_not_sentinel and ds_project_code_safe_integer).
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.

Inline comments:
In `@docs/architecture/scheduling/scheduling.md`:
- Around line 354-386: The calculateWindow function currently reads
catalog.replicationKeyType when bookmark?.replication_key_type is undefined,
which can be invalid for non-INCREMENTAL streams; update calculateWindow to
assert at runtime that the stream is INCREMENTAL (or that a replication key type
exists) before using catalog.replicationKeyType — e.g., check the stream's sync
mode or throw/log and bail if not INCREMENTAL — and add a short comment in
calculateWindow explaining this precondition; reference calculateWindow,
StreamDescriptor.replicationKeyType, and StreamBookmark to locate and fix the
code.
- Around line 172-173: The spec's table entry for the state_document default is
out of sync with the DB schema; update the table row for `state_document` to use
the actual default used in the schema (`{ bookmarks: {}, versions: {},
currently_syncing: null }`) so the spec matches the code in `stitches.ts`;
locate the `state_document` table row in scheduling.md and replace the existing
`{"bookmarks":{}}` default with the schema default object to keep documentation
consistent with the `stitches.ts` definition.

In `@packages/connectors/src/framework/piece.ts`:
- Around line 43-51: Add a JSDoc note to PollRecord.replicationKeyValue
explaining the coercion contract: implementers may return string or number, but
the system coerces numeric values to string for storage/comparison (e.g.,
PollWindow.lowerBound/upperBound are strings) and
CursorManagerService.trackHighWaterMark may call String(value) / Number(value)
when comparing or storing high-water marks; this informs piece authors they can
return either type and how it will be treated internally.

In `@packages/database/src/schema/identity.ts`:
- Around line 268-271: The constraint named "ds_project_code_safe_integer" only
enforces an upper bound on table.dsProjectCode; update the SQL expression used
in the check() call to enforce the full JS safe integer range by requiring the
column to be NULL or between -9007199254740991 and 9007199254740991 (or
equivalently add a >= -9007199254740991 check alongside the existing <= check)
so values below the negative safe limit are rejected; locate the check(...)
invocation for "ds_project_code_safe_integer" and modify its sql`${...}`
expression accordingly.

---

Outside diff comments:
In `@packages/database/src/schema/identity.ts`:
- Around line 258-276: Add a uniqueness constraint/index for dsProjectCode to
prevent multiple organizations from sharing the same DolphinScheduler project
code: modify the table definition that declares dsProjectCode (identifier:
dsProjectCode on the table variable) to add a unique index (similar to
organization_slug_unique_idx) that enforces uniqueness only for non-null
dsProjectCode values (a partial unique index WHERE "dsProjectCode" IS NOT NULL)
so NULLs remain allowed but any assigned project code is unique across
organizations; ensure the new index has a clear name (e.g.,
organization_ds_project_code_unique_idx) and coexists with the existing checks
(organization_id_not_sentinel and ds_project_code_safe_integer).

ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: 8206bd63-b189-40b8-8ac0-3507d63642e0

📥 Commits

Reviewing files that changed from the base of the PR and between 07ebf1f and 6497b06.

📒 Files selected for processing (7)
  • apps/api/src/db/schema.ts
  • docs/architecture/master/tasks.md
  • docs/architecture/scheduling/scheduling.md
  • packages/connectors/src/framework/piece.ts
  • packages/database/src/schema/identity.ts
  • packages/database/src/schema/stitches.ts
  • packages/identity/src/adapters/drizzle-tenant.adapter.spec.ts

Comment thread docs/architecture/scheduling/scheduling.md Outdated
Comment thread docs/architecture/scheduling/scheduling.md
Comment thread packages/connectors/src/framework/piece.ts
Comment thread packages/database/src/schema/identity.ts

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Actionable comments posted: 3

🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.

Inline comments:
In `@docs/architecture/scheduling/scheduling.md`:
- Around line 167-176: The current docs show the sync_cursors table and its
state_document column holding multi-stream Singer state, which contradicts the
one-row-per-(stitch_id, stream_name) invariant; update the documentation for
sync_cursors (mentioning sync_cursors, state_document, stream_name, stitch_id)
to state explicitly that each row's state_document must contain only the
bookmark/version entry for that row's stream_name (not multiple streams),
replace the multi-stream JSON example with a single-stream JSON containing only
the corresponding stream entry, and make the same clarification/change in the
later section (the part covering lines 186–217) so the invariant is consistent
throughout.

In `@packages/connectors/src/framework/piece.ts`:
- Around line 67-72: The poll() contract is ambiguous because it receives only
credentials, window, and a string cursor which cannot identify which stream
(from describeStreams()) it targets nor preserve connector-specific
bookmark.offset shape; update the API so poll() accepts a structured cursor
including streamName and the connector bookmark (e.g., an object with streamName
and bookmark fields) and update PollPage to include streamName and a typed
nextPageCursor that can round-trip the bookmark; specifically modify the poll()
signature and any usages, adjust the PollPage interface (and related types
around lines ~135-156 and ~188-191) to carry streamName and a connector bookmark
payload instead of a plain string, and ensure scheduler state keys (stitchId,
streamName) align with the new cursor shape so the first page can
deterministically target a stream and subsequent pages can restore
bookmark.offset.
- Around line 29-40: Change PollWindow from a single interface into a
discriminated union keyed by replicationKeyType: define separate types (e.g.,
PollWindowTimestamp, PollWindowNumeric, PollWindowOpaque) each with
replicationKeyType set to the appropriate ReplicationKeyType literal and only
the bounds that make sense (Timestamp: lowerBound and upperBound as ISO strings;
Numeric: lowerBound and optional/typed numeric bound if applicable; Opaque:
lowerBound/token only, no mandatory upperBound). Replace the existing PollWindow
export with the union of those types and update any usages (including the poll()
signature and any code referencing upperBound/lowerBound) to handle the union
via the replicationKeyType discriminator.

ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: 88e5d8f0-b979-4e9d-a22c-d4a2e7c0def1

📥 Commits

Reviewing files that changed from the base of the PR and between 6497b06 and e454078.

📒 Files selected for processing (3)
  • docs/architecture/scheduling/scheduling.md
  • packages/connectors/src/framework/piece.ts
  • packages/database/src/schema/identity.ts

Comment thread docs/architecture/scheduling/scheduling.md
Comment thread packages/connectors/src/framework/piece.ts Outdated
Comment thread packages/connectors/src/framework/piece.ts Outdated

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Actionable comments posted: 1

🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.

Inline comments:
In `@docs/architecture/scheduling/scheduling.md`:
- Around line 563-592: The compose snippet uses simple depends_on entries
(ds-master, ds-worker, ds-api, ds-alert) which only wait for container start,
not readiness; update the snippet to reference health-based dependencies for the
backing services (postgres, zookeeper) and add corresponding healthcheck
definitions for those services so Docker Compose can use condition:
service_healthy; specifically, modify the depends_on blocks for ds-master
(postgres, zookeeper) and for ds-worker/ds-api to reference postgres (and
ds-master where applicable) with condition: service_healthy, and add healthcheck
configurations for postgres and zookeeper services to probe readiness (e.g.,
pg_isready/HTTP/TCP checks), ensuring service names ds-master, ds-worker,
ds-api, ds-alert, postgres, and zookeeper are updated accordingly.

ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: 17ed8d8a-473d-4d1c-be0d-86bfa5dda4ed

📥 Commits

Reviewing files that changed from the base of the PR and between e454078 and 7322977.

📒 Files selected for processing (2)
  • docs/architecture/scheduling/scheduling.md
  • packages/connectors/src/framework/piece.ts

Comment thread docs/architecture/scheduling/scheduling.md
@pramodnarayana
pramodnarayana merged commit f3c5745 into development Mar 23, 2026
2 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant