cursor reset endpoint - #103
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:
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configurationConfiguration used: Path: .coderabbit.yaml Review profile: ASSERTIVE Plan: Pro Run ID: 📒 Files selected for processing (1)
📝 WalkthroughWalkthroughAdds an admin controller for per-stitch cursor management (DELETE/GET) with Redis per-stream locking, centralizes poll lock key generation, conditionally wires PollSyncRunner behind WINDMILL_ENABLED, changes outbox indexing to a partial index, and updates related scheduler, Windmill client, outbox, tests, and docs. Changes
Sequence Diagram(s)sequenceDiagram
actor Client as HTTP Client
participant Ctrl as CursorResetController
participant Auth as AuthGuard/AuthContext
participant Redis as Redis
participant DB as Database
Client->>Ctrl: DELETE /admin/stitches/:id/cursor/:streamName
Ctrl->>Auth: resolve operator
Ctrl->>Redis: SET key=lock:poll:{stitchId}:{streamName} NX PX <ttl> (token)
alt lock acquired
Redis-->>Ctrl: "OK"
Ctrl->>DB: DELETE FROM syncCursors WHERE stitch_id=id AND stream_name=streamName
DB-->>Ctrl: rowsDeleted
Ctrl->>Redis: EVAL (if value==token then DEL)
Redis-->>Ctrl: release result
Ctrl-->>Client: 204 No Content
else lock not acquired
Redis-->>Ctrl: null
Ctrl-->>Client: 409 Conflict
end
sequenceDiagram
actor Client as HTTP Client
participant Ctrl as CursorResetController
participant DB as Database
Client->>Ctrl: GET /admin/stitches/:id/cursors
Ctrl->>DB: SELECT syncIntervalMinutes, scheduleEnabled FROM integrationStitches WHERE id=...
alt stitch found
DB-->>Ctrl: { syncIntervalMinutes, scheduleEnabled }
Ctrl->>DB: SELECT id, stitchId, streamName, createdAt, updatedAt FROM syncCursors WHERE stitch_id=...
DB-->>Ctrl: [{...}, ...]
Ctrl->>Ctrl: compute ageMs, paused, staleThreshold, stale
Ctrl-->>Client: 200 [{..., ageMs, paused, stale}, ...]
else stitch not found
DB-->>Ctrl: null
Ctrl-->>Client: 404 NotFound
end
Estimated Code Review Effort🎯 4 (Complex) | ⏱️ ~45 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 docstrings
🧪 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: 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 `@apps/api/src/modules/scheduler/cursor-reset.controller.ts`:
- Around line 59-78: The lookup for stale cursors currently only uses
stitch.syncIntervalMinutes (from integrationStitches) so paused stitches
(scheduleEnabled = false) can be marked stale; update the DB select that fetches
the stitch to also pull scheduleEnabled (select({ syncIntervalMinutes:
integrationStitches.syncIntervalMinutes, scheduleEnabled:
integrationStitches.scheduleEnabled })), then change the stale calculation where
staleThresholdMs is used (and any boolean `stale` decision) to short-circuit
when stitch.scheduleEnabled === false (e.g., treat as not stale or set a
separate `paused` flag) so paused stitches are handled distinctly from genuinely
stale ones and downstream logic referencing syncCursors/stitch can see the
paused status.
- Around line 43-50: The current delete(...) call on syncCursors can be undone
by an in-flight PollSyncRunner that later checkpoints/finalizes the old
high-water mark; either reject resets while a poll lock is held or persist a
reset marker the runner checks before any checkpoint. Fix option A: before
deleting, check whether the poll lock for {stitchId, streamName} is held (the
same lock PollSyncRunner uses) and return a 409/validation error if locked. Fix
option B: instead of deleting the row, write a durable reset marker (e.g., an
inserted/updated column or a dedicated reset table row tied to syncCursors with
stitchId and streamName and a reset=true/timestamp) and update PollSyncRunner to
read that marker before every checkpoint/finalize operation and abort/perform
full-refresh if present, then clear the marker after reset is honored; reference
symbols: syncCursors, PollSyncRunner, delete(...), checkpoint/finalize.
In `@docs/architecture/master/tasks.md`:
- Line 433: Update the stale wording in the checklist item that mentions "next
DS run": locate the line containing "`DELETE
/admin/stitches/:id/cursor/:streamName` — deletes the `sync_cursors` row for
`(stitchId, streamName)`; returns `204`; idempotent (no error if row absent);
triggers full refresh on next DS run` and replace "next DS run" with either
"next scheduled run" or "next Windmill-triggered run" to match the Windmill
scheduler terminology used elsewhere.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: b4c0cd23-474a-480f-b92d-c832795327e7
📒 Files selected for processing (4)
apps/api/src/modules/scheduler/cursor-reset.controller.spec.tsapps/api/src/modules/scheduler/cursor-reset.controller.tsapps/api/src/modules/scheduler/scheduler.module.tsdocs/architecture/master/tasks.md
There was a problem hiding this comment.
Actionable comments posted: 8
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In @.claude/agent-memory/code-reviewer/MEMORY.md:
- Line 41: The Markdown line with the TTL formula currently uses bare *
operators which trigger MD037; edit the line containing "Lock key:
`lock:poll:${stitchId}:${streamName}`, TTL = max(syncIntervalMinutes * 2 *
60_000, 5 * 60_000)ms" and wrap the entire TTL expression (the max(...) portion)
in backticks so the asterisks are treated as code rather than emphasis.
In `@apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts`:
- Around line 173-194: The test's redaction assertion is vacuous because the
mocked DB row used in the mockDb.where call for syncCursors does not include a
stateDocument field; add a dummy stateDocument property to that mocked row
fixture and then call controller.listCursors(STITCH_ID) and assert the returned
item does not have stateDocument (i.e., ensure listCursors strips it). Update
the spec where mockDb.where.mockResolvedValueOnce([...]) for the syncCursors
case and strengthen the assertion on the result from listCursors to verify the
stateDocument was removed.
- Around line 72-88: The test and controller currently encode a static Redis
lock value ('admin-reset') and expect an unconditional DEL (mockRedis.del) in
the finally block for deleteCursor; replace that with a per-acquisition token
(e.g., UUID) on SET in deleteCursor (instead of the literal 'admin-reset'), and
perform a compare-and-delete (via a Lua script or redlock library) so the DEL
only removes the lock if the stored token matches; then update the spec
(cursor-reset.controller.spec.ts) to assert that mockRedis.set was called with a
generated token (or that the token argument is passed), and that the final
unlock path uses a compare-and-delete call (assert
mockRedis.eval/mockRedis.evalSha or the redlock unlock method was called) rather
than an unconditional mockRedis.del, and adjust other tests covering lines
100–123 and 126–139 similarly.
- Around line 90-98: The test currently asserts await
expect(controller.deleteCursor(...)).resolves.not.toThrow() which is
inappropriate for a Promise<void>; change the assertion to await
expect(controller.deleteCursor(STITCH_ID, STREAM_NAME,
mockCtx)).resolves.toBeUndefined() so the test verifies the promise resolves
with undefined; update the spec where the "DELETE is idempotent — no error when
row is absent" test calls controller.deleteCursor to use
resolves.toBeUndefined() instead of resolves.not.toThrow().
In `@apps/api/src/modules/scheduler/http-windmill.client.ts`:
- Around line 82-90: When handling the 409 response in the Windmill client (the
block that logs via this.logger.warn about STITCH_RUNNER_PATH), read a bounded
snippet of the response body (e.g., up to a fixed byte/char limit like 1KB)
instead of discarding it, then include that snippet in the logger.warn message
so diagnostics are preserved; ensure you preserve the existing warning text,
handle potential read errors gracefully, and avoid unbounded reads to prevent
memory issues (refer to STITCH_RUNNER_PATH and this.logger.warn in your change).
In `@apps/api/src/modules/scheduler/outbox-worker.service.ts`:
- Around line 130-146: The logs currently interpolate raw lastError (in the
OutboxWorker/OutboxWorkerService failure handling block) which may leak
sensitive connection strings or tokens; update the failure and retry logging in
the method that updates schedulerOutbox (references: lastError, record.id,
record.action, record.stitchId, MAX_OUTBOX_ATTEMPTS, nextRetryAt) to never emit
raw error text — instead pass a sanitized value (e.g. replace sensitive parts or
use a constant like "[redacted]" or a brief sanitizedError) when calling
this.logger.error and this.logger.warn, and persist the full lastError only to
the DB record (not to logs) if needed. Ensure the same sanitized value is used
in both the permanent-failure log and the retry log paths.
In `@apps/api/src/modules/scheduler/scheduler.controller.ts`:
- Around line 41-43: The controller's catch block (in scheduler.controller.ts)
is turning all non-InternalServerErrorException errors into 500s; instead,
change the error handling so permanent domain errors are raised as proper NestJS
HTTP exceptions at the service/runner boundary (update poll-sync-runner.ts to
throw NotFoundException, BadRequestException, etc. for missing
stitch/connection/unregistered piece) and remove the blanket conversion in the
SchedulerController catch (or rethrow non-InternalServerErrorException errors
unchanged) so 4xx errors are returned for permanent failures rather than being
wrapped into InternalServerErrorException.
In `@apps/api/src/modules/scheduler/scheduler.module.ts`:
- Around line 38-72: The current SyncRunner provider eagerly injects
REDIS_CLIENT, TokenManagerService, PieceRegistryService, and
CursorManagerService even when WINDMILL_ENABLED !== 'true'; change the provider
to only inject ConfigService and ModuleRef (or Injector) so heavy deps are
deferred, and inside the useFactory check config.get('WINDMILL_ENABLED') ===
'true' and only then call moduleRef.get(DATABASE_CONNECTION),
moduleRef.get(REDIS_CLIENT), moduleRef.get(TokenManagerService),
moduleRef.get(PieceRegistryService), moduleRef.get(CursorManagerService) and
instantiate new PollSyncRunner(...) ; otherwise return new StubSyncRunner().
Ensure you update the provider's inject array and move resolution of
PollSyncRunner dependencies into the conditional branch in the useFactory
(referencing SyncRunner provider, useFactory, PollSyncRunner, StubSyncRunner,
and ModuleRef/Injector).
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: f062984a-b1ae-4750-a8aa-8a62e17be4d3
📒 Files selected for processing (18)
.claude/agent-memory/code-reviewer/MEMORY.mdapps/api/drizzle/0014_scheduler_outbox_partial_index.sqlapps/api/src/modules/scheduler/cursor-reset.controller.spec.tsapps/api/src/modules/scheduler/cursor-reset.controller.tsapps/api/src/modules/scheduler/http-windmill.client.spec.tsapps/api/src/modules/scheduler/http-windmill.client.tsapps/api/src/modules/scheduler/interval-to-cron.tsapps/api/src/modules/scheduler/lock-keys.tsapps/api/src/modules/scheduler/outbox-worker.service.spec.tsapps/api/src/modules/scheduler/outbox-worker.service.tsapps/api/src/modules/scheduler/poll-sync-runner.tsapps/api/src/modules/scheduler/scheduler.controller.spec.tsapps/api/src/modules/scheduler/scheduler.controller.tsapps/api/src/modules/scheduler/scheduler.module.tsapps/api/src/modules/scheduler/scheduler.service.spec.tsapps/api/src/modules/scheduler/stub-sync-runner.tsdocs/architecture/master/tasks.mdpackages/database/src/schema/stitches.ts
There was a problem hiding this comment.
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)
apps/api/src/modules/scheduler/outbox-worker.service.ts (1)
132-160:⚠️ Potential issue | 🟠 MajorDo not persist the unredacted error text in
scheduler_outbox.This block now explicitly acknowledges that these messages may contain OAuth tokens or connection strings, but both update paths still store the raw
lastErrorin the DB. That turns transient secrets into durable data-at-rest and broadens who can retrieve them. Persist a redacted form instead, and keep any full-fidelity payload only in a secret-safe trace channel if you absolutely need it.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/api/src/modules/scheduler/outbox-worker.service.ts` around lines 132 - 160, The code currently writes the unredacted error string (lastError) into scheduler_outbox; instead persist a redacted/sanitized value instead: use the already-computed safeError (or compute a new redacted variable via sanitizeError(lastError)) for the DB updates in both branches (the update(...) .set({ status: 'failed', lastError, ... }) and the update(...) .set({ status: 'pending', lastError, ... })), i.e. replace lastError with the redacted value when calling this.db.update on schedulerOutbox; if you must keep the full error for secure troubleshooting, send the raw error to a secret-safe trace or error reporting channel (not the DB) while only storing safeError in the scheduler_outbox row.
♻️ Duplicate comments (1)
apps/api/src/modules/scheduler/scheduler.module.ts (1)
26-26:⚠️ Potential issue | 🟠 MajorRedis is still mandatory on the disabled path.
SyncRunneris now lazy, but this module always registersCursorResetController, and that controller's constructor injectsREDIS_CLIENT. SoWINDMILL_ENABLED !== 'true'still does not fully decouple local/dev bootstrap from Redis. Gate the controller behind the same feature flag or move it into a Redis-enabled submodule.Also applies to: 39-64
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@apps/api/src/modules/scheduler/scheduler.module.ts` at line 26, The SchedulerModule currently always registers CursorResetController which injects REDIS_CLIENT, making Redis required even when WINDMILL_ENABLED !== 'true'; update SchedulerModule so CursorResetController is only included when process.env.WINDMILL_ENABLED === 'true' (e.g., compute the controllers array conditionally and only push CursorResetController when the flag is set) or move CursorResetController into a separate Redis-enabled submodule that is imported only when WINDMILL_ENABLED is true; this ensures REDIS_CLIENT is not injected on the disabled path while keeping SyncRunner lazy as-is.
🤖 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/scheduler/cursor-reset.controller.spec.ts`:
- Around line 152-163: This test ("DELETE uses the correct lock key format and
NX flag") is redundant with the earlier DELETE test that already asserts the
lock key and NX flag; remove this test block or repurpose it to cover a distinct
edge case (for example, verify behavior when Redis set returns non-OK, when NX
is absent, when TTL/PX is different, or when mockDb.where throws). Locate the
test by the it(...) description and references to controller.deleteCursor,
mockRedis.set, mockDb.delete, and mockDb.where, then either delete the entire
it(...) block or change its setup/assertions to validate a unique scenario
instead of duplicating the existing assertions.
In `@apps/api/src/modules/scheduler/cursor-reset.controller.ts`:
- Around line 29-30: The ADMIN_RESET_LOCK_TTL_MS (in cursor-reset.controller.ts)
is too short and can expire before the DB delete completes, allowing
PollSyncRunner to reacquire the lock and undo the reset; either increase
ADMIN_RESET_LOCK_TTL_MS to the same multi-minute floor used by the poll runner
or implement lease renewal around the delete operation (acquire lock -> start
renewing until delete finishes -> stop renewing and release). Update the
constant and/or the admin reset path that performs the delete so it references
ADMIN_RESET_LOCK_TTL_MS and ensures the lock is actively renewed for the entire
duration of the DB delete to prevent reacquisition by PollSyncRunner.
- Around line 100-102: The log in cursor-reset.controller.ts uses
this.logger.log and currently emits ctx.user.email (PII); change the logged
actor to a stable internal principal identifier (e.g., ctx.user.id or
ctx.user.principalId) instead of ctx.user.email, update the log key from
"operator" to something like "actorId", and if the codebase requires an audit
trail that includes email keep that in a dedicated secure audit sink rather than
general service logs; ensure you reference and update the this.logger.log call
in the CursorResetController (or the method containing the current log) to use
the non-PII field and add a defensive fallback if the internal id is missing.
In `@docs/architecture/master/tasks.md`:
- Around line 444-459: The checklist in the tasks document skips S4 between S3
and S5; update the checklist in the same section so numbering is contiguous by
either inserting a short S4 entry (e.g., "S4: [note if intentionally omitted]"
or the actual item if it was accidentally left out) or renumber S5→S4 and
subsequent items accordingly; look for the lines listing "S3" and "S5" in the
same checklist block and adjust the labels to restore sequential numbering or
add a clarifying note.
---
Outside diff comments:
In `@apps/api/src/modules/scheduler/outbox-worker.service.ts`:
- Around line 132-160: The code currently writes the unredacted error string
(lastError) into scheduler_outbox; instead persist a redacted/sanitized value
instead: use the already-computed safeError (or compute a new redacted variable
via sanitizeError(lastError)) for the DB updates in both branches (the
update(...) .set({ status: 'failed', lastError, ... }) and the update(...)
.set({ status: 'pending', lastError, ... })), i.e. replace lastError with the
redacted value when calling this.db.update on schedulerOutbox; if you must keep
the full error for secure troubleshooting, send the raw error to a secret-safe
trace or error reporting channel (not the DB) while only storing safeError in
the scheduler_outbox row.
---
Duplicate comments:
In `@apps/api/src/modules/scheduler/scheduler.module.ts`:
- Line 26: The SchedulerModule currently always registers CursorResetController
which injects REDIS_CLIENT, making Redis required even when WINDMILL_ENABLED !==
'true'; update SchedulerModule so CursorResetController is only included when
process.env.WINDMILL_ENABLED === 'true' (e.g., compute the controllers array
conditionally and only push CursorResetController when the flag is set) or move
CursorResetController into a separate Redis-enabled submodule that is imported
only when WINDMILL_ENABLED is true; this ensures REDIS_CLIENT is not injected on
the disabled path while keeping SyncRunner lazy as-is.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: 910aeff7-b2b0-4100-8696-7ac8c008576f
📒 Files selected for processing (11)
.claude/agent-memory/code-reviewer/MEMORY.mdapps/api/src/modules/scheduler/cursor-reset.controller.spec.tsapps/api/src/modules/scheduler/cursor-reset.controller.tsapps/api/src/modules/scheduler/http-windmill.client.tsapps/api/src/modules/scheduler/outbox-worker.service.tsapps/api/src/modules/scheduler/poll-sync-runner.spec.tsapps/api/src/modules/scheduler/poll-sync-runner.tsapps/api/src/modules/scheduler/scheduler.controller.spec.tsapps/api/src/modules/scheduler/scheduler.controller.tsapps/api/src/modules/scheduler/scheduler.module.tsdocs/architecture/master/tasks.md
There was a problem hiding this comment.
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 `@apps/api/src/modules/scheduler/outbox-worker.service.ts`:
- Around line 156-160: The log message in OutboxWorker's retry path uses
inconsistent field naming ("attempt=") vs the other message ("attempts=");
update the logger.warn call in outbox-worker.service.ts (the one referencing
this.logger.warn, record.attempts, MAX_OUTBOX_ATTEMPTS, nextRetryAt, lastError)
to use "attempts=" instead of "attempt=" so both logs consistently emit
attempts=<number>/<MAX_OUTBOX_ATTEMPTS> for easier parsing.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: 7678030b-cac7-4949-b1f2-7b5c8707cdc9
📒 Files selected for processing (5)
apps/api/src/modules/scheduler/cursor-reset.controller.spec.tsapps/api/src/modules/scheduler/cursor-reset.controller.tsapps/api/src/modules/scheduler/outbox-worker.service.tsapps/api/src/modules/scheduler/scheduler.module.tsdocs/architecture/master/tasks.md
Summary by CodeRabbit
New Features
Bug Fixes / Behavior Changes
Tests
Documentation
Database