Skip to content

cursor manager engine - #101

Merged
pramodnarayana merged 4 commits into
developmentfrom
feat/cursor-manager-engine
Mar 24, 2026
Merged

pramodnarayana merged 4 commits into
developmentfrom
feat/cursor-manager-engine

Conversation

@pramodnarayana

@pramodnarayana pramodnarayana commented Mar 24, 2026 •

Copy link
Copy Markdown
Owner

Summary by CodeRabbit

  • New Features

    • Added a cursor management service and types to improve windowing, high-water marks, and checkpointing.
    • Introduced a new engine package with build and test configuration.
  • Database

    • Added a persistent cursor table with indexes, foreign key, and update trigger to store per-stream bookmarks.
  • Tests

    • Added comprehensive tests for window calculation and high-water mark tracking.
  • Documentation

    • Updated architecture and orchestration docs to reflect Windmill-based changes.
  • Chores

    • Updated a development dependency.

@pramodnarayana
pramodnarayana marked this pull request as draft March 24, 2026 08:51
@pramodnarayana
pramodnarayana marked this pull request as ready for review March 24, 2026 08:51
@coderabbitai

coderabbitai Bot commented Mar 24, 2026 •

Copy link
Copy Markdown
Contributor
📝 Walkthrough

Walkthrough

Adds a new @nexiom/engine package (exports, build/test config, types), implements CursorManagerService with tests and exported constants/types, adds a DB migration for sync_cursors (indexes + trigger), updates architecture docs for Windmill orchestration, and bumps jsdom in apps/web/package.json.

Changes

Cohort / File(s) Summary
Engine package manifest & build/test config
packages/engine/package.json, packages/engine/tsconfig.json, packages/engine/vitest.config.ts
Create new @nexiom/engine package: ESM exports, build/test scripts, TypeScript config, and Vitest setup.
Engine public surface
packages/engine/src/index.ts
Re-export cursor-manager types and expose CursorManagerService and DEFAULT_CURSOR_CHECKPOINT_INTERVAL.
Cursor manager types
packages/engine/src/state/cursor-manager.types.ts
Add re-exports from connectors/framework and define StreamBookmark, SyncStateDocument, ExecuteStitchPayload/Result, and StreamResult types.
Cursor manager implementation
packages/engine/src/state/cursor-manager.service.ts
Add CursorManagerService (OnModuleInit), DEFAULT_CURSOR_CHECKPOINT_INTERVAL, calculateWindow() and trackHighWaterMark() with env-driven config and validation/warning/error paths.
Cursor manager tests
packages/engine/src/state/cursor-manager.service.spec.ts
Add comprehensive Vitest suite covering init, calculateWindow() across replication methods/types, and trackHighWaterMark() edge cases.
Database migration & journal
packages/database/drizzle/0004_sync_cursors.sql, packages/database/drizzle/meta/_journal.json
Add sync_cursors table (UUID PK, stitch_id FK → integration_stitch), stream_name, JSONB state_document, unique/indexes, updated_at trigger, and journal entry.
Docs & small dep bump
docs/architecture/master/tasks.md, apps/web/package.json
Replace DolphinScheduler references with Windmill orchestration in docs; bump devDependency jsdom from ^25.0.1 → ^27.4.0.

Estimated Code Review Effort

🎯 4 (Complex) | ⏱️ ~45 minutes

Possibly Related PRs

Poem

🐰 I hop through checkpoints with a curious twirl,

Bookmarks snug in pockets where cursors unfurl.
Timestamps, numbers, opaque tokens in line,
Windmill whispers schedules — syncs hum in time.
I nibble the bytes and send the markers to shine.

🚥 Pre-merge checks | ✅ 2 | ❌ 1

❌ Failed checks (1 inconclusive)

Check name Status Explanation Resolution
Title check ❓ Inconclusive The title 'cursor manager engine' is vague and doesn't clearly convey the scope and purpose of the changes, which involve implementing cursor management infrastructure including database schema, service layer, types, and configuration. Consider a more descriptive title such as 'Add cursor manager service and database schema for state tracking' or 'Implement CursorManagerService with database persistence' to better communicate the primary changes.
✅ Passed checks (2 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Docstring Coverage ✅ Passed No functions found in the changed files to evaluate docstring coverage. Skipping docstring coverage check.

✏️ 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/cursor-manager-engine

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

Caution

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

⚠️ Outside diff range comments (1)
docs/architecture/master/tasks.md (1)

543-543: ⚠️ Potential issue | 🟡 Minor

Inconsistent orchestrator reference in summary table.

The summary table still references "DolphinScheduler orchestration" for Phase 3.5, but the task descriptions (T029, T030, T049, T050) have been updated to reference Windmill. Update the summary for consistency.

📝 Proposed fix
-| 3.5 — Stateful Sync | T046–T050 + T029 + T030 | DolphinScheduler orchestration + Singer-style cursor engine |
+| 3.5 — Stateful Sync | T046–T050 + T029 + T030 | Windmill orchestration + Singer-style cursor engine |
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@docs/architecture/master/tasks.md` at line 543, Update the Phase 3.5 summary
table row that currently reads "DolphinScheduler orchestration + Singer-style
cursor engine" so the orchestrator reference matches the task descriptions
(T029, T030, T049, T050) by replacing "DolphinScheduler orchestration" with
"Windmill orchestration" while preserving the rest of the cell (e.g., keep
"Singer-style cursor engine") so the table entry for "3.5 — Stateful Sync |
T046–T050 + T029 + T030 | ..." is consistent with the task docs.
🤖 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/master/tasks.md`:
- Around line 417-428: The documentation for T047 incorrectly states there are
24 unit tests; update the single-line summary under "T047 · package:
`packages/engine/` — `CursorManagerService`" to state 28 tests instead of 24 to
match the actual `cursor-manager.service.spec.ts` test count (grep -c "it('").
Ensure the sentence referencing "Unit tests: 24 tests" is changed to "Unit
tests: 28 tests" and leave all other descriptions (functions like
calculateWindow, trackHighWaterMark, onModuleInit,
DEFAULT_CURSOR_CHECKPOINT_INTERVAL) unchanged.

In `@packages/database/drizzle/0004_sync_cursors.sql`:
- Around line 15-21: The trigger function set_updated_at() is generic and
created with CREATE OR REPLACE which can conflict with other migrations; rename
it to a table-scoped name like sync_cursors_set_updated_at_fn and update the
trigger to call that function (or alternatively add an existence/body check
before creating), so replace references to set_updated_at() with
sync_cursors_set_updated_at_fn and ensure the trigger definition for the
sync_cursors table uses the new function name to avoid accidental overwrites.
- Line 11: The migration creates only the composite unique index
sync_cursors_stitch_stream_unique_idx on table sync_cursors but misses the
single-column index sync_cursors_stitch_idx referenced in the TypeScript schema
(stitches.ts); update the migration to also create a btree index named
sync_cursors_stitch_idx on sync_cursors(stitch_id) so queries that filter by
stitch_id (e.g., "get all cursors for a stitch") use the index and match the
schema definition.

---

Outside diff comments:
In `@docs/architecture/master/tasks.md`:
- Line 543: Update the Phase 3.5 summary table row that currently reads
"DolphinScheduler orchestration + Singer-style cursor engine" so the
orchestrator reference matches the task descriptions (T029, T030, T049, T050) by
replacing "DolphinScheduler orchestration" with "Windmill orchestration" while
preserving the rest of the cell (e.g., keep "Singer-style cursor engine") so the
table entry for "3.5 — Stateful Sync | T046–T050 + T029 + T030 | ..." is
consistent with the task docs.

ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: ef0d8699-e61a-489c-beb9-a2b164395848

📥 Commits

Reviewing files that changed from the base of the PR and between 01c00d8 and fa1f90c.

⛔ Files ignored due to path filters (1)
  • pnpm-lock.yaml is excluded by !**/pnpm-lock.yaml
📒 Files selected for processing (11)
  • apps/web/package.json
  • docs/architecture/master/tasks.md
  • packages/database/drizzle/0004_sync_cursors.sql
  • packages/database/drizzle/meta/_journal.json
  • packages/engine/package.json
  • packages/engine/src/index.ts
  • packages/engine/src/state/cursor-manager.service.spec.ts
  • packages/engine/src/state/cursor-manager.service.ts
  • packages/engine/src/state/cursor-manager.types.ts
  • packages/engine/tsconfig.json
  • packages/engine/vitest.config.ts

Comment on lines 417 to 428
### T047 · package: `packages/engine/` — `CursorManagerService`

- [ ] `cursor-manager.types.ts` — `StreamBookmark` (with `replication_key_type`, `offset`), `SyncStateDocument` (with `versions`, `currently_syncing`), `ExecuteStitchPayload`, `ExecuteStitchResult`, `StreamResult` interfaces; re-exports `ReplicationKeyType`, `StreamDescriptor` from `@nexiom/connectors/framework`
- [ ] `cursor-manager.service.ts` — `CursorManagerService` with three phases:
- `calculateWindow(bookmark: StreamBookmark | undefined, catalog: StreamDescriptor): PollWindow` — type-aware (timestamp/numeric/opaque); applies safety buffer for timestamps; `'0'` lower bound for numeric on first run; epoch for timestamp on first run
- `trackHighWaterMark(records: PollRecord[], currentMax: string, replicationKeyType: ReplicationKeyType): string` — numeric comparison for `numeric` keys; lexicographic for `timestamp`; last-write-wins for `opaque`
- Checkpoint is the caller's responsibility (write `state_document` to `public.sync_cursors` after successful batch commit)
- [ ] `onModuleInit()` logs configured safety buffer value
- [ ] `CURSOR_CHECKPOINT_INTERVAL` env var (default 10) — exported constant consumed by `SchedulerWorker`
- [ ] Unit tests: full-refresh (no bookmark), incremental with safety buffer, high-water mark across multiple pages, configurable buffer via `ConfigService`
- Files: `packages/engine/src/state/cursor-manager.types.ts`, `packages/engine/src/state/cursor-manager.service.ts`, `packages/engine/src/state/cursor-manager.spec.ts`, `packages/engine/src/index.ts`, `packages/engine/package.json`
- [x] `cursor-manager.types.ts` — `StreamBookmark`, `SyncStateDocument`, `ExecuteStitchPayload`, `ExecuteStitchResult`, `StreamResult`; re-exports `ReplicationKeyType`, `StreamDescriptor`, `PollWindow`, `PollRecord`, `PollPage` from `@nexiom/connectors/framework`
- [x] `cursor-manager.service.ts` — `CursorManagerService`:
- `calculateWindow(bookmark, catalog)` — timestamp/numeric/opaque; 5-min default safety buffer; epoch / '0' / '' on first run; FULL_TABLE returns opaque empty window
- `trackHighWaterMark(records, currentMax, type)` — numeric max, lexicographic ISO-8601 max, opaque last-write-wins; returns `currentMax` on empty input
- Checkpoint is the caller's responsibility
- [x] `onModuleInit()` logs safety buffer and checkpoint interval
- [x] `DEFAULT_CURSOR_CHECKPOINT_INTERVAL = 10` exported constant; configurable via `CURSOR_CHECKPOINT_INTERVAL` env var
- [x] Unit tests: 24 tests covering all key types, first-run, incremental, safety buffer, multi-page accumulation, configurable buffer
- Files: `packages/engine/src/state/cursor-manager.types.ts`, `packages/engine/src/state/cursor-manager.service.ts`, `packages/engine/src/state/cursor-manager.service.spec.ts`, `packages/engine/src/index.ts`, `packages/engine/package.json`
- Depends: T046

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

⚠️ Potential issue | 🟡 Minor

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
# Count the number of test cases in the cursor-manager spec file
grep -c "it('" packages/engine/src/state/cursor-manager.service.spec.ts

Repository: pramodnarayana/nexiom

Length of output: 67


Update test count in T047 documentation from 24 to 28.

The spec file contains 28 test cases (verified via grep -c "it('"), not 24 as currently documented.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@docs/architecture/master/tasks.md` around lines 417 - 428, The documentation
for T047 incorrectly states there are 24 unit tests; update the single-line
summary under "T047 · package: `packages/engine/` — `CursorManagerService`" to
state 28 tests instead of 24 to match the actual
`cursor-manager.service.spec.ts` test count (grep -c "it('"). Ensure the
sentence referencing "Unit tests: 24 tests" is changed to "Unit tests: 28 tests"
and leave all other descriptions (functions like calculateWindow,
trackHighWaterMark, onModuleInit, DEFAULT_CURSOR_CHECKPOINT_INTERVAL) unchanged.

Comment thread packages/database/drizzle/0004_sync_cursors.sql
Comment thread packages/database/drizzle/0004_sync_cursors.sql 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: 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/master/tasks.md`:
- Around line 283-293: Remove or retarget all stale "pending T047" notes
referring to the stubbed sync path: update descriptions for executeStitch, the
Abstract SyncRunner and StubSyncRunner to no longer say "stub until T047" or
"pending T047"; either drop the T047 reference entirely or replace it with the
current remaining work (e.g., "real implementation in progress" or a new ticket
number). Specifically edit the entries mentioning SyncRunner, StubSyncRunner and
executeStitch so they reflect the completed T047 state (or new target) and
ensure the duplicate T047 status language that contradicts the completed section
is removed or reconciled.
- Line 274: The task T028 references the wrong task ID (T046); update the
parenthetical reference to point to T047 so readers are directed to the
`packages/engine` re-export work — locate the T028 entry that mentions
`PollRecord`, `PollPage`, `PollWindow` and replace the `(see T046)` text with
`(see T047)`.
- Around line 381-396: The phase heading still reads "DolphinScheduler + Cursor
Manager" while the tasks below (T049, T050) and content reference Windmill;
update the enclosing section title to reference Windmill (e.g., "Windmill +
Cursor Manager" or "Windmill Integration + Cursor Manager") so it matches
T049/T050 and surrounding text, and scan the section for any remaining
"DolphinScheduler" mentions to rename them to "Windmill" for consistency.

ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: 7fcbe0b9-3516-4767-ac19-1c7d736b925b

📥 Commits

Reviewing files that changed from the base of the PR and between fa1f90c and 068319b.

📒 Files selected for processing (2)
  • docs/architecture/master/tasks.md
  • packages/database/drizzle/0004_sync_cursors.sql

Comment thread docs/architecture/master/tasks.md Outdated
Comment thread docs/architecture/master/tasks.md Outdated
Comment thread docs/architecture/master/tasks.md

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Actionable comments posted: 2

♻️ Duplicate comments (1)
docs/architecture/master/tasks.md (1)

302-314: ⚠️ Potential issue | 🟡 Minor

Remove stale "until T047" reference and resolve completion status inconsistency.

Two issues:

  1. Stale reference: Line 305 states the stub is "until T047," but T047 is marked complete at lines 417-428 in this same update. Either remove the "until T047" phrase or update it to reflect the current state.

  2. Inconsistent completion: Line 302 marks the endpoint complete ([x]), but lines 306-313 show incomplete items ([ ]) that are core parts of the same endpoint's functionality (lock acquisition, cursor reading, pagination, etc.). These should all share the same completion status.

📝 Suggested edit

If T030 is incomplete, change line 302:

-- [x] `POST /internal/scheduler/execute-stitch` — receives `{ stitchId }` from Windmill; returns `ExecuteStitchResult`
+- [ ] `POST /internal/scheduler/execute-stitch` — receives `{ stitchId }` from Windmill; returns `ExecuteStitchResult`

And update line 305:

-- [x] Delegates to `SyncRunner.run(stitchId)` (stub returning `{ status: 'started' }` until T047)
+- [ ] Delegates to `SyncRunner.run(stitchId)` (currently stubbed; full poll-run sequence below)

As per past review comments, the "pending T047" language should be removed or retargeted.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@docs/architecture/master/tasks.md` around lines 302 - 314, Remove the stale
"until T047" phrase from the bullet that mentions SyncRunner.run and make the
checklist consistent by changing the top-level endpoint item (`POST
/internal/scheduler/execute-stitch`) from checked to unchecked (or mark all
subitems checked if the work is actually done); ensure the subitems (lock
acquisition, piece.describeStreams, reading SyncStateDocument,
CursorManagerService.calculateWindow, pagination/checkpointing, final
checkpoint, unit tests) reflect their true completion state, and update any T047
reference text to match the current task status; check related symbols/files
(`executeStitch`, `SyncRunner.run`, `CursorManagerService.calculateWindow`, and
the listed files under Files) to verify which items remain pending and then
update the checklist 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 `@docs/architecture/master/tasks.md`:
- Around line 381-406: The checklist shows a dependency mismatch: T050
(WindmillClient) is marked complete while its dependency T049 (docker-compose
Windmill setup) remains incomplete; update docs/tasks.md to keep statuses
consistent by either marking all T049 checklist items as completed (change the
T049 boxes to [x]) if the infra was actually provisioned, or mark T050 as
incomplete (change the T050 boxes from [x] to [ ]) until T049 is done; edit the
entries for T049 and T050 in the tasks.md file so the "Depends: T049"
relationship and checklist states reflect the true current progress.
- Around line 283-295: Update the checklist and description to be consistent:
either remove the "placeholder" wording and note that SyncRunner (and T030) are
implemented if executeStitch and T030 are truly complete, or mark executeStitch
and T030 as incomplete (`[ ]`) if the production SyncRunner is still pending;
adjust the lines referencing executeStitch, SyncRunner, and T030 so the
completion checkboxes and the "placeholder/real implementation" text match.

---

Duplicate comments:
In `@docs/architecture/master/tasks.md`:
- Around line 302-314: Remove the stale "until T047" phrase from the bullet that
mentions SyncRunner.run and make the checklist consistent by changing the
top-level endpoint item (`POST /internal/scheduler/execute-stitch`) from checked
to unchecked (or mark all subitems checked if the work is actually done); ensure
the subitems (lock acquisition, piece.describeStreams, reading
SyncStateDocument, CursorManagerService.calculateWindow,
pagination/checkpointing, final checkpoint, unit tests) reflect their true
completion state, and update any T047 reference text to match the current task
status; check related symbols/files (`executeStitch`, `SyncRunner.run`,
`CursorManagerService.calculateWindow`, and the listed files under Files) to
verify which items remain pending and then update the checklist accordingly.

ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: f782b7ae-f5eb-4757-9b05-cba514fa7a40

📥 Commits

Reviewing files that changed from the base of the PR and between 068319b and 59fe82c.

📒 Files selected for processing (1)
  • docs/architecture/master/tasks.md

Comment thread docs/architecture/master/tasks.md
Comment thread docs/architecture/master/tasks.md

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

♻️ Duplicate comments (1)
docs/architecture/master/tasks.md (1)

385-406: ⚠️ Potential issue | 🟡 Minor

Resolve dependency-status mismatch between T049 and T050.

T050 is fully checked while T049 still has an unchecked item at Line 389, and T050 depends on T049. Keep completion states consistent (or split manual bootstrap into a separate task).

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@docs/architecture/master/tasks.md` around lines 385 - 406, The docs show a
dependency mismatch: T050 is marked complete while its dependency T049 still has
an unchecked item; update the task state so they are consistent by either
marking the remaining T049 checkbox (the manual bootstrap item at T049) as done
if completed, or remove the T049 dependency from T050 and/or split the manual
bootstrap step into its own task so T050 can be checked independently; ensure
you update the checklist entries for T049 and the dependency list for T050
accordingly (references: task identifiers "T049" and "T050" and the unchecked
bootstrap item in T049).
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.

Duplicate comments:
In `@docs/architecture/master/tasks.md`:
- Around line 385-406: The docs show a dependency mismatch: T050 is marked
complete while its dependency T049 still has an unchecked item; update the task
state so they are consistent by either marking the remaining T049 checkbox (the
manual bootstrap item at T049) as done if completed, or remove the T049
dependency from T050 and/or split the manual bootstrap step into its own task so
T050 can be checked independently; ensure you update the checklist entries for
T049 and the dependency list for T050 accordingly (references: task identifiers
"T049" and "T050" and the unchecked bootstrap item in T049).

ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Pro

Run ID: 23f31f3f-5077-4188-95a6-e0d5c240bea4

📥 Commits

Reviewing files that changed from the base of the PR and between 59fe82c and 28e6170.

📒 Files selected for processing (1)
  • docs/architecture/master/tasks.md

@pramodnarayana
pramodnarayana merged commit 22db6d0 into development Mar 24, 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