From de6b3a26cb3d82ee2f662417a687a0f03310866d Mon Sep 17 00:00:00 2001 From: Pramod Date: Wed, 25 Mar 2026 11:05:17 +0530 Subject: [PATCH 1/5] cursor reset endpoint --- .../scheduler/cursor-reset.controller.spec.ts | 192 ++++++++++++++++++ .../scheduler/cursor-reset.controller.ts | 89 ++++++++ .../src/modules/scheduler/scheduler.module.ts | 3 +- docs/architecture/master/tasks.md | 11 +- 4 files changed, 290 insertions(+), 5 deletions(-) create mode 100644 apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts create mode 100644 apps/api/src/modules/scheduler/cursor-reset.controller.ts diff --git a/apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts b/apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts new file mode 100644 index 00000000..1720ffa0 --- /dev/null +++ b/apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts @@ -0,0 +1,192 @@ +import { describe, it, expect, vi, beforeEach } from 'vitest'; +import { Test } from '@nestjs/testing'; +import { BadRequestException, NotFoundException } from '@nestjs/common'; +import { AuthGuard } from '@nexiom/auth'; +import { DATABASE_CONNECTION } from '@nexiom/database'; +import { integrationStitches, syncCursors } from '@nexiom/database'; +import { SystemAdminGuard } from '../identity/auth/system-admin.guard.js'; +import { CursorResetController } from './cursor-reset.controller.js'; + +const STITCH_ID = '550e8400-e29b-41d4-a716-446655440000'; +const STREAM_NAME = 'Account'; + +const mockCtx = { + user: { id: 'u1', email: 'admin@example.com' }, +} as unknown as import('@nexiom/auth').RequestAuthContext; + +function createMockDb() { + let currentTable: unknown = null; + + const chain = { + select: vi.fn().mockReturnThis(), + from: vi.fn().mockImplementation((table: unknown) => { + currentTable = table; + return chain; + }), + where: vi.fn().mockReturnThis(), + limit: vi.fn().mockReturnThis(), + delete: vi.fn().mockReturnThis(), + // Expose currentTable for assertions + getCurrentTable: () => currentTable, + }; + return chain; +} + +type MockDb = ReturnType; + +describe('CursorResetController', () => { + let controller: CursorResetController; + let mockDb: MockDb; + + beforeEach(async () => { + mockDb = createMockDb(); + + const module = await Test.createTestingModule({ + controllers: [CursorResetController], + providers: [{ provide: DATABASE_CONNECTION, useValue: mockDb }], + }) + .overrideGuard(AuthGuard) + .useValue({ canActivate: () => true }) + .overrideGuard(SystemAdminGuard) + .useValue({ canActivate: () => true }) + .compile(); + + controller = module.get(CursorResetController); + vi.clearAllMocks(); + }); + + // ── DELETE ────────────────────────────────────────────────────────────── + + it('DELETE returns void when row exists', async () => { + mockDb.delete.mockReturnValue(mockDb); + mockDb.where.mockResolvedValue(undefined); + + const result = await controller.deleteCursor( + STITCH_ID, + STREAM_NAME, + mockCtx, + ); + + expect(result).toBeUndefined(); + expect(mockDb.delete).toHaveBeenCalledOnce(); + expect(mockDb.where).toHaveBeenCalledOnce(); + }); + + it('DELETE is idempotent — no-op when row is absent', async () => { + mockDb.delete.mockReturnValue(mockDb); + mockDb.where.mockResolvedValue(undefined); + + await expect( + controller.deleteCursor(STITCH_ID, STREAM_NAME, mockCtx), + ).resolves.not.toThrow(); + }); + + it('DELETE rejects streamName with newline characters', async () => { + await expect( + controller.deleteCursor(STITCH_ID, 'Account\nevil-log-line', mockCtx), + ).rejects.toThrow(BadRequestException); + }); + + it('DELETE rejects streamName with null bytes', async () => { + await expect( + controller.deleteCursor(STITCH_ID, 'Account\x00', mockCtx), + ).rejects.toThrow(BadRequestException); + }); + + it('DELETE rejects streamName longer than 200 characters', async () => { + const longName = 'a'.repeat(201); + await expect( + controller.deleteCursor(STITCH_ID, longName, mockCtx), + ).rejects.toThrow(BadRequestException); + }); + + // ── GET ───────────────────────────────────────────────────────────────── + + it('GET returns cursors with stale: true when age exceeds 2x interval', async () => { + const now = Date.now(); + const sixtyOneMinutesAgo = new Date(now - 61 * 60_000); + + // Table-based dispatch: stitch lookup (integrationStitches) → cursor listing (syncCursors) + mockDb.from.mockImplementation((table: unknown) => { + if (table === integrationStitches) { + mockDb.where.mockReturnValueOnce(mockDb); + mockDb.limit.mockResolvedValueOnce([{ syncIntervalMinutes: 30 }]); + } else if (table === syncCursors) { + mockDb.where.mockResolvedValueOnce([ + { + id: 'c1', + stitchId: STITCH_ID, + streamName: STREAM_NAME, + stateDocument: {}, + createdAt: sixtyOneMinutesAgo, + updatedAt: sixtyOneMinutesAgo, + }, + ]); + } + return mockDb; + }); + + const result = await controller.listCursors(STITCH_ID); + + expect(result).toHaveLength(1); + expect(result[0].stale).toBe(true); + expect(result[0].ageMs).toBeGreaterThan(2 * 30 * 60_000); + }); + + it('GET returns stale: false when age is within 2x interval', async () => { + const now = Date.now(); + const fiveMinutesAgo = new Date(now - 5 * 60_000); + + mockDb.from.mockImplementation((table: unknown) => { + if (table === integrationStitches) { + mockDb.where.mockReturnValueOnce(mockDb); + mockDb.limit.mockResolvedValueOnce([{ syncIntervalMinutes: 30 }]); + } else if (table === syncCursors) { + mockDb.where.mockResolvedValueOnce([ + { + id: 'c1', + stitchId: STITCH_ID, + streamName: STREAM_NAME, + stateDocument: {}, + createdAt: fiveMinutesAgo, + updatedAt: fiveMinutesAgo, + }, + ]); + } + return mockDb; + }); + + const result = await controller.listCursors(STITCH_ID); + + expect(result).toHaveLength(1); + expect(result[0].stale).toBe(false); + expect(result[0].ageMs).toBeLessThanOrEqual(2 * 30 * 60_000); + }); + + it('GET returns empty array when stitch has no cursors', async () => { + mockDb.from.mockImplementation((table: unknown) => { + if (table === integrationStitches) { + mockDb.where.mockReturnValueOnce(mockDb); + mockDb.limit.mockResolvedValueOnce([{ syncIntervalMinutes: 30 }]); + } else if (table === syncCursors) { + mockDb.where.mockResolvedValueOnce([]); + } + return mockDb; + }); + + const result = await controller.listCursors(STITCH_ID); + + expect(result).toEqual([]); + }); + + it('GET throws NotFoundException when stitch does not exist', async () => { + mockDb.select.mockReturnValue(mockDb); + mockDb.from.mockReturnValue(mockDb); + mockDb.where.mockReturnValue(mockDb); + mockDb.limit.mockResolvedValue([]); + + await expect(controller.listCursors(STITCH_ID)).rejects.toThrow( + NotFoundException, + ); + }); +}); diff --git a/apps/api/src/modules/scheduler/cursor-reset.controller.ts b/apps/api/src/modules/scheduler/cursor-reset.controller.ts new file mode 100644 index 00000000..3574b5c5 --- /dev/null +++ b/apps/api/src/modules/scheduler/cursor-reset.controller.ts @@ -0,0 +1,89 @@ +import { + BadRequestException, + Controller, + Delete, + Get, + Inject, + Logger, + Param, + ParseUUIDPipe, + HttpCode, + HttpStatus, + NotFoundException, + UseGuards, +} from '@nestjs/common'; +import { AuthContext, AuthGuard, type RequestAuthContext } from '@nexiom/auth'; +import { SystemAdminGuard } from '../identity/auth/system-admin.guard.js'; +import { + DATABASE_CONNECTION, + syncCursors, + integrationStitches, +} from '@nexiom/database'; +import type { DrizzleDb } from '@nexiom/database'; +import { eq, and } from 'drizzle-orm'; + +@UseGuards(AuthGuard, SystemAdminGuard) +@Controller('admin/stitches') +export class CursorResetController { + private readonly logger = new Logger(CursorResetController.name); + + constructor(@Inject(DATABASE_CONNECTION) private readonly db: DrizzleDb) {} + + @Delete(':id/cursor/:streamName') + @HttpCode(HttpStatus.NO_CONTENT) + async deleteCursor( + @Param('id', ParseUUIDPipe) id: string, + @Param('streamName') streamName: string, + @AuthContext() ctx: RequestAuthContext, + ): Promise { + if (!/^[\w.-]{1,200}$/.test(streamName)) { + throw new BadRequestException('streamName contains invalid characters'); + } + + await this.db + .delete(syncCursors) + .where( + and( + eq(syncCursors.stitchId, id), + eq(syncCursors.streamName, streamName), + ), + ); + + this.logger.log( + `Cursor reset: stitchId=${id}, streamName=${JSON.stringify(streamName)}, operator=${ctx.user.email}`, + ); + } + + @Get(':id/cursors') + async listCursors(@Param('id', ParseUUIDPipe) id: string) { + const [stitch] = await this.db + .select({ syncIntervalMinutes: integrationStitches.syncIntervalMinutes }) + .from(integrationStitches) + .where(eq(integrationStitches.id, id)) + .limit(1); + + if (!stitch) { + throw new NotFoundException(`Stitch not found: ${id}`); + } + + const rows = await this.db + .select() + .from(syncCursors) + .where(eq(syncCursors.stitchId, id)); + + const now = Date.now(); + const staleThresholdMs = + stitch.syncIntervalMinutes > 0 + ? 2 * stitch.syncIntervalMinutes * 60_000 + : Number.POSITIVE_INFINITY; + + return rows.map((row) => { + const ageMs = now - row.updatedAt.getTime(); + return { + ...row, + ageMs, + stale: ageMs > staleThresholdMs, + }; + }); + } +} diff --git a/apps/api/src/modules/scheduler/scheduler.module.ts b/apps/api/src/modules/scheduler/scheduler.module.ts index 7dc54c4b..2c5fb110 100644 --- a/apps/api/src/modules/scheduler/scheduler.module.ts +++ b/apps/api/src/modules/scheduler/scheduler.module.ts @@ -11,12 +11,13 @@ import { SyncRunner } from './sync-runner.js'; import { PollSyncRunner } from './poll-sync-runner.js'; import { SchedulerService } from './scheduler.service.js'; import { SchedulerController } from './scheduler.controller.js'; +import { CursorResetController } from './cursor-reset.controller.js'; import { InternalSchedulerGuard } from './internal-scheduler.guard.js'; import { OutboxWorkerService } from './outbox-worker.service.js'; @Module({ imports: [DbModule, ConnectionsModule, PiecesModule], - controllers: [SchedulerController], + controllers: [SchedulerController, CursorResetController], providers: [ { provide: WindmillClient, diff --git a/docs/architecture/master/tasks.md b/docs/architecture/master/tasks.md index 8c16aeb4..ebd3cddc 100644 --- a/docs/architecture/master/tasks.md +++ b/docs/architecture/master/tasks.md @@ -430,10 +430,13 @@ Each task is one commit (or one small PR). Checkboxes track completion. ### T048 · api: `CursorResetEndpoint` — admin full-refresh trigger -- [ ] `DELETE /admin/stitches/:id/cursor/:streamName` — deletes the `sync_cursors` row for `(stitchId, streamName)`; returns `204`; triggers full refresh on next DS run -- [ ] Superadmin guard only; logs the reset with operator identity for audit -- [ ] `GET /admin/stitches/:id/cursors` — lists all `sync_cursors` rows for the stitch with `updated_at` age and stale flag (`age > 2 × syncIntervalMinutes`) -- [ ] Unit tests for both endpoints +- [x] `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 +- [x] Superadmin guard only (`AuthGuard` + `SystemAdminGuard`); logs the reset with operator email for audit +- [x] `GET /admin/stitches/:id/cursors` — lists all `sync_cursors` rows for the stitch enriched with `ageMs` and `stale` flag (`ageMs > 2 × syncIntervalMinutes × 60_000`); throws `NotFoundException` if stitch missing +- [x] `streamName` validated against `/^[\w.-]{1,200}$/`; invalid values rejected with `BadRequestException`; value JSON-encoded in audit log to prevent log injection +- [x] `staleThresholdMs` guarded for non-positive `syncIntervalMinutes` (returns `Infinity`) +- [x] `row.updatedAt.getTime()` used directly (Drizzle returns JS `Date` for `timestamptz`) +- [x] Unit tests for both endpoints (9 tests: DELETE row exists, DELETE idempotent, DELETE invalid streamName x3, GET stale true, GET stale false, GET empty cursors, GET not found); test mock rewritten to table-based dispatch - Files: `apps/api/src/modules/scheduler/cursor-reset.controller.ts`, `apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts` - Depends: T046, T029 From 4796d4b83f8d49be222a6517c5ba9effaea66413 Mon Sep 17 00:00:00 2001 From: Pramod Date: Wed, 25 Mar 2026 11:48:23 +0530 Subject: [PATCH 2/5] cursor reset endpoint code review --- .claude/agent-memory/code-reviewer/MEMORY.md | 9 +- .../0014_scheduler_outbox_partial_index.sql | 17 +++ .../scheduler/cursor-reset.controller.spec.ts | 135 +++++++++++++++--- .../scheduler/cursor-reset.controller.ts | 70 +++++++-- .../scheduler/http-windmill.client.spec.ts | 26 +++- .../modules/scheduler/http-windmill.client.ts | 39 ++--- .../src/modules/scheduler/interval-to-cron.ts | 6 + apps/api/src/modules/scheduler/lock-keys.ts | 12 ++ .../scheduler/outbox-worker.service.spec.ts | 24 +++- .../scheduler/outbox-worker.service.ts | 21 +-- .../src/modules/scheduler/poll-sync-runner.ts | 18 ++- .../scheduler/scheduler.controller.spec.ts | 18 ++- .../modules/scheduler/scheduler.controller.ts | 18 ++- .../src/modules/scheduler/scheduler.module.ts | 41 +++++- .../scheduler/scheduler.service.spec.ts | 4 +- .../src/modules/scheduler/stub-sync-runner.ts | 6 +- docs/architecture/master/tasks.md | 25 +++- packages/database/src/schema/stitches.ts | 8 +- 18 files changed, 395 insertions(+), 102 deletions(-) create mode 100644 apps/api/drizzle/0014_scheduler_outbox_partial_index.sql create mode 100644 apps/api/src/modules/scheduler/lock-keys.ts diff --git a/.claude/agent-memory/code-reviewer/MEMORY.md b/.claude/agent-memory/code-reviewer/MEMORY.md index a820508b..2291d7e2 100644 --- a/.claude/agent-memory/code-reviewer/MEMORY.md +++ b/.claude/agent-memory/code-reviewer/MEMORY.md @@ -38,9 +38,14 @@ - `PollSyncRunner` implements Singer-style poll with Redis NX locks, crash-resume via `currently_syncing` + `bookmark.offset` - State persisted in `public.sync_cursors` (stitchId + streamName composite key, JSONB stateDocument) - `CursorManagerService` in `packages/engine` -- stateless, computes windows + tracks HWM -- Lock key: `lock:poll:${stitchId}:${streamName}`, TTL = syncIntervalMinutes (potential issue: long syncs exceed TTL) +- Lock key: `lock:poll:${stitchId}:${streamName}`, TTL = max(syncIntervalMinutes * 2 * 60_000, 5 * 60_000)ms +- Lock uses Lua-atomic renew/release with owner token verification - `OAuthCredentialBlob` double-cast to `Record` is a recurring pattern in piece calls -- Test mock pattern: `makeDb()` with sequential callCount dispatching -- fragile, order-dependent +- Test mock pattern: `makeDb()` with table-aware `.from()` dispatching (improved from callCount) +- OutboxWorkerService: off-by-one risk in retry count vs back-off comment (OW-1 flagged 2026-03-25) +- CursorResetController: TOCTOU race between lock check and cursor delete (CR-1 flagged 2026-03-25) +- HttpWindmillClient.ensureStitchScript: does not update existing scripts on content change (HW-1 flagged 2026-03-25) +- PollSyncRunner always instantiated even when WINDMILL_ENABLED=false (SM-1 flagged 2026-03-25) ## Schema Notes diff --git a/apps/api/drizzle/0014_scheduler_outbox_partial_index.sql b/apps/api/drizzle/0014_scheduler_outbox_partial_index.sql new file mode 100644 index 00000000..b0eb5091 --- /dev/null +++ b/apps/api/drizzle/0014_scheduler_outbox_partial_index.sql @@ -0,0 +1,17 @@ +-- Migration: Upgrade scheduler_outbox poll index to a partial index (S1). +-- +-- The composite index on (status, next_retry_at) scans all rows including +-- the large succeeded/failed population. A partial index covering only +-- pending rows is much smaller and keeps the poll query fast at scale. +-- +-- The poll query shape: WHERE status = 'pending' AND next_retry_at <= NOW() +-- ORDER BY next_retry_at ASC LIMIT N FOR UPDATE SKIP LOCKED + +-- Drop old composite index if it exists (created by schema push in earlier envs). +DROP INDEX IF EXISTS "scheduler_outbox_poll_idx"; +--> statement-breakpoint + +-- Partial index: only pending rows, ordered by earliest retry time. +CREATE INDEX IF NOT EXISTS "scheduler_outbox_poll_idx" + ON "scheduler_outbox" ("next_retry_at" ASC) + WHERE status = 'pending'; diff --git a/apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts b/apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts index 1720ffa0..8ed0afb8 100644 --- a/apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts +++ b/apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts @@ -1,9 +1,14 @@ import { describe, it, expect, vi, beforeEach } from 'vitest'; import { Test } from '@nestjs/testing'; -import { BadRequestException, NotFoundException } from '@nestjs/common'; +import { + BadRequestException, + ConflictException, + NotFoundException, +} from '@nestjs/common'; import { AuthGuard } from '@nexiom/auth'; import { DATABASE_CONNECTION } from '@nexiom/database'; import { integrationStitches, syncCursors } from '@nexiom/database'; +import { REDIS_CLIENT } from '@nexiom/cache'; import { SystemAdminGuard } from '../identity/auth/system-admin.guard.js'; import { CursorResetController } from './cursor-reset.controller.js'; @@ -15,35 +20,42 @@ const mockCtx = { } as unknown as import('@nexiom/auth').RequestAuthContext; function createMockDb() { - let currentTable: unknown = null; - const chain = { select: vi.fn().mockReturnThis(), - from: vi.fn().mockImplementation((table: unknown) => { - currentTable = table; - return chain; - }), + from: vi.fn().mockReturnThis(), where: vi.fn().mockReturnThis(), limit: vi.fn().mockReturnThis(), delete: vi.fn().mockReturnThis(), - // Expose currentTable for assertions - getCurrentTable: () => currentTable, }; return chain; } +function createMockRedis() { + return { + // SET NX: returns 'OK' (lock acquired) or null (already held) + set: vi.fn().mockResolvedValue('OK'), + del: vi.fn().mockResolvedValue(1), + }; +} + type MockDb = ReturnType; +type MockRedis = ReturnType; describe('CursorResetController', () => { let controller: CursorResetController; let mockDb: MockDb; + let mockRedis: MockRedis; beforeEach(async () => { mockDb = createMockDb(); + mockRedis = createMockRedis(); const module = await Test.createTestingModule({ controllers: [CursorResetController], - providers: [{ provide: DATABASE_CONNECTION, useValue: mockDb }], + providers: [ + { provide: DATABASE_CONNECTION, useValue: mockDb }, + { provide: REDIS_CLIENT, useValue: mockRedis }, + ], }) .overrideGuard(AuthGuard) .useValue({ canActivate: () => true }) @@ -57,7 +69,8 @@ describe('CursorResetController', () => { // ── DELETE ────────────────────────────────────────────────────────────── - it('DELETE returns void when row exists', async () => { + it('DELETE acquires lock, deletes cursor, releases lock', async () => { + mockRedis.set.mockResolvedValue('OK'); mockDb.delete.mockReturnValue(mockDb); mockDb.where.mockResolvedValue(undefined); @@ -68,11 +81,14 @@ describe('CursorResetController', () => { ); expect(result).toBeUndefined(); + // Lock acquired then released + expect(mockRedis.set).toHaveBeenCalledOnce(); + expect(mockRedis.del).toHaveBeenCalledOnce(); expect(mockDb.delete).toHaveBeenCalledOnce(); - expect(mockDb.where).toHaveBeenCalledOnce(); }); - it('DELETE is idempotent — no-op when row is absent', async () => { + it('DELETE is idempotent — no error when row is absent', async () => { + mockRedis.set.mockResolvedValue('OK'); mockDb.delete.mockReturnValue(mockDb); mockDb.where.mockResolvedValue(undefined); @@ -81,6 +97,48 @@ describe('CursorResetController', () => { ).resolves.not.toThrow(); }); + it('DELETE releases lock even when DB delete throws', async () => { + mockRedis.set.mockResolvedValue('OK'); + mockDb.delete.mockReturnValue(mockDb); + mockDb.where.mockRejectedValue(new Error('DB error')); + + await expect( + controller.deleteCursor(STITCH_ID, STREAM_NAME, mockCtx), + ).rejects.toThrow('DB error'); + + // Lock must be released in finally + expect(mockRedis.del).toHaveBeenCalledOnce(); + }); + + it('DELETE throws ConflictException when poll lock is already held (NX fails)', async () => { + mockRedis.set.mockResolvedValue(null); // NX fails — lock is held + + await expect( + controller.deleteCursor(STITCH_ID, STREAM_NAME, mockCtx), + ).rejects.toThrow(ConflictException); + + // Must not attempt DB delete + expect(mockDb.delete).not.toHaveBeenCalled(); + // Must not attempt lock release (never acquired it) + expect(mockRedis.del).not.toHaveBeenCalled(); + }); + + it('DELETE uses the correct lock key format', async () => { + mockRedis.set.mockResolvedValue('OK'); + mockDb.delete.mockReturnValue(mockDb); + mockDb.where.mockResolvedValue(undefined); + + await controller.deleteCursor(STITCH_ID, STREAM_NAME, mockCtx); + + expect(mockRedis.set).toHaveBeenCalledWith( + `lock:poll:${STITCH_ID}:${STREAM_NAME}`, + 'admin-reset', + 'PX', + expect.any(Number), + 'NX', + ); + }); + it('DELETE rejects streamName with newline characters', async () => { await expect( controller.deleteCursor(STITCH_ID, 'Account\nevil-log-line', mockCtx), @@ -106,18 +164,18 @@ describe('CursorResetController', () => { const now = Date.now(); const sixtyOneMinutesAgo = new Date(now - 61 * 60_000); - // Table-based dispatch: stitch lookup (integrationStitches) → cursor listing (syncCursors) mockDb.from.mockImplementation((table: unknown) => { if (table === integrationStitches) { mockDb.where.mockReturnValueOnce(mockDb); - mockDb.limit.mockResolvedValueOnce([{ syncIntervalMinutes: 30 }]); + mockDb.limit.mockResolvedValueOnce([ + { syncIntervalMinutes: 30, scheduleEnabled: true }, + ]); } else if (table === syncCursors) { mockDb.where.mockResolvedValueOnce([ { id: 'c1', stitchId: STITCH_ID, streamName: STREAM_NAME, - stateDocument: {}, createdAt: sixtyOneMinutesAgo, updatedAt: sixtyOneMinutesAgo, }, @@ -130,7 +188,10 @@ describe('CursorResetController', () => { expect(result).toHaveLength(1); expect(result[0].stale).toBe(true); + expect(result[0].paused).toBe(false); expect(result[0].ageMs).toBeGreaterThan(2 * 30 * 60_000); + // stateDocument must not be present in the response + expect(result[0]).not.toHaveProperty('stateDocument'); }); it('GET returns stale: false when age is within 2x interval', async () => { @@ -140,14 +201,15 @@ describe('CursorResetController', () => { mockDb.from.mockImplementation((table: unknown) => { if (table === integrationStitches) { mockDb.where.mockReturnValueOnce(mockDb); - mockDb.limit.mockResolvedValueOnce([{ syncIntervalMinutes: 30 }]); + mockDb.limit.mockResolvedValueOnce([ + { syncIntervalMinutes: 30, scheduleEnabled: true }, + ]); } else if (table === syncCursors) { mockDb.where.mockResolvedValueOnce([ { id: 'c1', stitchId: STITCH_ID, streamName: STREAM_NAME, - stateDocument: {}, createdAt: fiveMinutesAgo, updatedAt: fiveMinutesAgo, }, @@ -160,14 +222,49 @@ describe('CursorResetController', () => { expect(result).toHaveLength(1); expect(result[0].stale).toBe(false); + expect(result[0].paused).toBe(false); expect(result[0].ageMs).toBeLessThanOrEqual(2 * 30 * 60_000); }); + it('GET returns stale: false and paused: true when stitch is paused', async () => { + const now = Date.now(); + const twoHoursAgo = new Date(now - 120 * 60_000); + + mockDb.from.mockImplementation((table: unknown) => { + if (table === integrationStitches) { + mockDb.where.mockReturnValueOnce(mockDb); + mockDb.limit.mockResolvedValueOnce([ + { syncIntervalMinutes: 30, scheduleEnabled: false }, + ]); + } else if (table === syncCursors) { + mockDb.where.mockResolvedValueOnce([ + { + id: 'c1', + stitchId: STITCH_ID, + streamName: STREAM_NAME, + createdAt: twoHoursAgo, + updatedAt: twoHoursAgo, + }, + ]); + } + return mockDb; + }); + + const result = await controller.listCursors(STITCH_ID); + + expect(result).toHaveLength(1); + // Age (120 min) exceeds 2×30 min threshold but stitch is paused — not stale. + expect(result[0].stale).toBe(false); + expect(result[0].paused).toBe(true); + }); + it('GET returns empty array when stitch has no cursors', async () => { mockDb.from.mockImplementation((table: unknown) => { if (table === integrationStitches) { mockDb.where.mockReturnValueOnce(mockDb); - mockDb.limit.mockResolvedValueOnce([{ syncIntervalMinutes: 30 }]); + mockDb.limit.mockResolvedValueOnce([ + { syncIntervalMinutes: 30, scheduleEnabled: true }, + ]); } else if (table === syncCursors) { mockDb.where.mockResolvedValueOnce([]); } diff --git a/apps/api/src/modules/scheduler/cursor-reset.controller.ts b/apps/api/src/modules/scheduler/cursor-reset.controller.ts index 3574b5c5..438bd09b 100644 --- a/apps/api/src/modules/scheduler/cursor-reset.controller.ts +++ b/apps/api/src/modules/scheduler/cursor-reset.controller.ts @@ -1,5 +1,6 @@ import { BadRequestException, + ConflictException, Controller, Delete, Get, @@ -20,14 +21,22 @@ import { integrationStitches, } from '@nexiom/database'; import type { DrizzleDb } from '@nexiom/database'; +import { REDIS_CLIENT, type Redis } from '@nexiom/cache'; import { eq, and } from 'drizzle-orm'; +import { pollLockKey } from './lock-keys.js'; + +/** Short-lived TTL (ms) for the admin reset lock — long enough to cover the DB delete. */ +const ADMIN_RESET_LOCK_TTL_MS = 5_000; @UseGuards(AuthGuard, SystemAdminGuard) @Controller('admin/stitches') export class CursorResetController { private readonly logger = new Logger(CursorResetController.name); - constructor(@Inject(DATABASE_CONNECTION) private readonly db: DrizzleDb) {} + constructor( + @Inject(DATABASE_CONNECTION) private readonly db: DrizzleDb, + @Inject(REDIS_CLIENT) private readonly redis: Redis, + ) {} @Delete(':id/cursor/:streamName') @HttpCode(HttpStatus.NO_CONTENT) @@ -40,14 +49,39 @@ export class CursorResetController { throw new BadRequestException('streamName contains invalid characters'); } - await this.db - .delete(syncCursors) - .where( - and( - eq(syncCursors.stitchId, id), - eq(syncCursors.streamName, streamName), - ), + // Atomically acquire the per-stream poll lock (NX) for the duration of the + // delete. SET NX is a single atomic operation, so there is no TOCTOU window + // between checking and holding the lock. If the lock is already held by a + // PollSyncRunner the NX fails and we return 409 — the admin retries once + // the run completes. Holding the lock during the delete prevents a new poll + // from starting and immediately re-creating the cursor row we just removed. + const lockKey = pollLockKey(id, streamName); + const acquired = await this.redis.set( + lockKey, + 'admin-reset', + 'PX', + ADMIN_RESET_LOCK_TTL_MS, + 'NX', + ); + if (!acquired) { + throw new ConflictException( + `Stream "${streamName}" is currently being polled — retry after the run completes`, ); + } + + try { + await this.db + .delete(syncCursors) + .where( + and( + eq(syncCursors.stitchId, id), + eq(syncCursors.streamName, streamName), + ), + ); + } finally { + // Release immediately — we only needed the lock to prevent concurrent starts. + await this.redis.del(lockKey); + } this.logger.log( `Cursor reset: stitchId=${id}, streamName=${JSON.stringify(streamName)}, operator=${ctx.user.email}`, @@ -57,7 +91,10 @@ export class CursorResetController { @Get(':id/cursors') async listCursors(@Param('id', ParseUUIDPipe) id: string) { const [stitch] = await this.db - .select({ syncIntervalMinutes: integrationStitches.syncIntervalMinutes }) + .select({ + syncIntervalMinutes: integrationStitches.syncIntervalMinutes, + scheduleEnabled: integrationStitches.scheduleEnabled, + }) .from(integrationStitches) .where(eq(integrationStitches.id, id)) .limit(1); @@ -66,8 +103,16 @@ export class CursorResetController { throw new NotFoundException(`Stitch not found: ${id}`); } + // Select only safe, non-sensitive columns. stateDocument is intentionally + // excluded — it may contain opaque vendor cursor tokens. const rows = await this.db - .select() + .select({ + id: syncCursors.id, + stitchId: syncCursors.stitchId, + streamName: syncCursors.streamName, + createdAt: syncCursors.createdAt, + updatedAt: syncCursors.updatedAt, + }) .from(syncCursors) .where(eq(syncCursors.stitchId, id)); @@ -79,10 +124,13 @@ export class CursorResetController { return rows.map((row) => { const ageMs = now - row.updatedAt.getTime(); + // Paused stitches are never stale — cursors are not expected to advance. + const paused = !stitch.scheduleEnabled; return { ...row, ageMs, - stale: ageMs > staleThresholdMs, + paused, + stale: paused ? false : ageMs > staleThresholdMs, }; }); } diff --git a/apps/api/src/modules/scheduler/http-windmill.client.spec.ts b/apps/api/src/modules/scheduler/http-windmill.client.spec.ts index abbca39e..4078105e 100644 --- a/apps/api/src/modules/scheduler/http-windmill.client.spec.ts +++ b/apps/api/src/modules/scheduler/http-windmill.client.spec.ts @@ -77,7 +77,9 @@ describe('HttpWindmillClient', () => { expect(getCallInit(fetchSpy, 0).method).toBe('POST'); }); - it('treats 409 Conflict as success (idempotent — script already exists)', async () => { + it('treats 409 as non-fatal — same content hash already deployed', async () => { + // Windmill returns 409 when the exact same content hash exists at the path. + // A new STITCH_RUNNER_CONTENT hash would produce 200 (new version created). fetchSpy.mockResolvedValue(mockResponse(409, 'conflict')); await expect(client.ensureStitchScript()).resolves.toBeUndefined(); }); @@ -252,5 +254,27 @@ describe('HttpWindmillClient', () => { const init = getCallInit(fetchSpy, 0); expect(init.signal).toBeDefined(); }); + + it('omits Content-Type on GET requests (no body)', async () => { + fetchSpy.mockResolvedValue(mockResponse(200, '{}')); + await client.scheduleExists(STITCH_ID); + + const headers = getCallInit(fetchSpy, 0).headers as Record< + string, + string + >; + expect(headers['Content-Type']).toBeUndefined(); + }); + + it('includes Content-Type: application/json on POST requests with a body', async () => { + fetchSpy.mockResolvedValue(mockResponse(200, '')); + await client.setScheduleEnabled(STITCH_ID, true); + + const headers = getCallInit(fetchSpy, 0).headers as Record< + string, + string + >; + expect(headers['Content-Type']).toBe('application/json'); + }); }); }); diff --git a/apps/api/src/modules/scheduler/http-windmill.client.ts b/apps/api/src/modules/scheduler/http-windmill.client.ts index 35061f5f..804ace8f 100644 --- a/apps/api/src/modules/scheduler/http-windmill.client.ts +++ b/apps/api/src/modules/scheduler/http-windmill.client.ts @@ -79,8 +79,14 @@ export class HttpWindmillClient extends WindmillClient { } if (res.status === 409) { await res.body?.cancel(); - this.logger.debug( - `Stitch-runner script already exists at ${STITCH_RUNNER_PATH} (409)`, + // Windmill returns 409 when an identical content hash already exists at this + // path (truly idempotent). If STITCH_RUNNER_CONTENT changed since the last + // deploy, the hash differs and Windmill creates a new version (200). + // A persistent 409 after a content change indicates the Windmill workspace + // needs a manual redeploy (delete the script at the path and redeploy). + this.logger.warn( + `Stitch-runner script at ${STITCH_RUNNER_PATH} returned 409 — ` + + 'script content matches an existing version or a manual redeploy is needed.', ); return; } @@ -221,7 +227,9 @@ export class HttpWindmillClient extends WindmillClient { method, headers: { Authorization: `Bearer ${this.token}`, - 'Content-Type': 'application/json', + // Only set Content-Type when there is a body — some API gateways reject + // GET requests that carry a Content-Type header. + ...(body !== undefined && { 'Content-Type': 'application/json' }), }, signal: AbortSignal.timeout(REQUEST_TIMEOUT_MS), ...(body !== undefined && { body: JSON.stringify(body) }), @@ -246,29 +254,4 @@ export class HttpWindmillClient extends WindmillClient { } await res.body?.cancel(); } - - /** - * Sends a request and parses the response body as JSON. - * Throws a descriptive Error on non-OK status or an empty body. - */ - private async requestJson( - method: string, - path: string, - body?: unknown, - ): Promise { - const res = await this.rawRequest(method, path, body); - if (!res.ok) { - const text = await res.text(); - throw new Error( - `Windmill API ${method} ${path} failed [${res.status}]: ${text}`, - ); - } - const text = await res.text(); - if (!text) { - throw new Error( - `Windmill API ${method} ${path} returned an empty body where JSON was expected`, - ); - } - return JSON.parse(text) as T; - } } diff --git a/apps/api/src/modules/scheduler/interval-to-cron.ts b/apps/api/src/modules/scheduler/interval-to-cron.ts index 12a0e0a6..d3dd77a3 100644 --- a/apps/api/src/modules/scheduler/interval-to-cron.ts +++ b/apps/api/src/modules/scheduler/interval-to-cron.ts @@ -4,6 +4,12 @@ import type { SyncIntervalMinutes } from '@nexiom/database'; * Maps a syncIntervalMinutes value to a 6-field Quartz/Windmill cron expression. * Field order: seconds minutes hours day-of-month month day-of-week * + * Windmill's scheduler is built on the Rust `cron` crate which supports both + * 5-field POSIX cron and 6-field Quartz-style cron (with a leading seconds + * field). The 6-field format is used here to allow second-level precision; + * all expressions pin the seconds field to 0 so runs trigger at the top of + * each interval boundary, matching POSIX semantics. + * * 30m -> every-30-minutes, 60m -> hourly, 120m -> every-2h, ..., 1440m -> daily midnight. */ export function intervalToCron(minutes: SyncIntervalMinutes): string { diff --git a/apps/api/src/modules/scheduler/lock-keys.ts b/apps/api/src/modules/scheduler/lock-keys.ts new file mode 100644 index 00000000..4bc885f8 --- /dev/null +++ b/apps/api/src/modules/scheduler/lock-keys.ts @@ -0,0 +1,12 @@ +/** + * Shared Redis lock key helpers for the scheduler module. + * + * These functions must be imported by every component that reads or writes + * a scheduler lock (PollSyncRunner, CursorResetController). A single + * definition here prevents the key format from diverging across files. + */ + +/** Redis key that guards per-stream poll execution. */ +export function pollLockKey(stitchId: string, streamName: string): string { + return `lock:poll:${stitchId}:${streamName}`; +} diff --git a/apps/api/src/modules/scheduler/outbox-worker.service.spec.ts b/apps/api/src/modules/scheduler/outbox-worker.service.spec.ts index 7a58b17f..b91e3749 100644 --- a/apps/api/src/modules/scheduler/outbox-worker.service.spec.ts +++ b/apps/api/src/modules/scheduler/outbox-worker.service.spec.ts @@ -202,8 +202,9 @@ describe('OutboxWorkerService', () => { expect(setCall.lastError).toContain('Windmill down'); }); - it('marks record as failed after MAX_OUTBOX_ATTEMPTS attempts', async () => { - // attempts=5 means we are AT the limit — next failure should permanently fail + it('retries with 32s delay when attempts=5 (5th attempt, not yet exhausted)', async () => { + // With MAX_OUTBOX_ATTEMPTS=6 (1 initial + 5 retries), attempt 5 still retries. + // Back-off delay = 2^5 * 1000 = 32 000 ms. const record = makeRecord('created', 5); mocks.returningClaim.mockResolvedValue([record]); mocks.findFirstStitch.mockResolvedValue(STITCH); @@ -211,6 +212,25 @@ describe('OutboxWorkerService', () => { await service.processOutbox(); + const setCall = mocks.markChain.set.mock.calls[0][0] as { + status: string; + nextRetryAt: Date; + }; + expect(setCall.status).toBe('pending'); + // 32s delay: nextRetryAt should be approximately 32s in the future. + const delayMs = setCall.nextRetryAt.getTime() - Date.now(); + expect(delayMs).toBeGreaterThan(30_000); + expect(delayMs).toBeLessThan(34_000); + }); + + it('permanently fails when attempts=6 (all 6 attempts exhausted)', async () => { + const record = makeRecord('created', 6); + mocks.returningClaim.mockResolvedValue([record]); + mocks.findFirstStitch.mockResolvedValue(STITCH); + scheduler.onStitchCreated.mockRejectedValue(new Error('still down')); + + await service.processOutbox(); + const setCall = mocks.markChain.set.mock.calls[0][0] as { status: string; lastError: string; diff --git a/apps/api/src/modules/scheduler/outbox-worker.service.ts b/apps/api/src/modules/scheduler/outbox-worker.service.ts index 5e97c9e8..b8577398 100644 --- a/apps/api/src/modules/scheduler/outbox-worker.service.ts +++ b/apps/api/src/modules/scheduler/outbox-worker.service.ts @@ -9,7 +9,9 @@ import { } from '@nexiom/database'; import { SchedulerService } from './scheduler.service.js'; -const MAX_OUTBOX_ATTEMPTS = 5; +// 1 initial attempt + 5 retries = 6 total attempts. +// Back-off delays between attempts: 2s, 4s, 8s, 16s, 32s. +const MAX_OUTBOX_ATTEMPTS = 6; const BATCH_SIZE = 20; /** @@ -60,7 +62,7 @@ export class OutboxWorkerService { if (claimed.length === 0) return; this.logger.debug( - `Claimed ${claimed.length} outbox record(s) for processing`, + `Claimed ${claimed.length} outbox record(s) for processing (ids=${claimed.map((r) => r.id).join(',')})`, ); await Promise.allSettled( @@ -99,7 +101,7 @@ export class OutboxWorkerService { await this.markSucceeded(record.id); this.logger.debug( - `Outbox record ${record.id}: action=${record.action} stitch=${record.stitchId} succeeded`, + `Outbox record succeeded: id=${record.id} action=${record.action} stitchId=${record.stitchId} attempts=${record.attempts}`, ); } catch (err) { await this.handleFailure(record, err); @@ -126,11 +128,12 @@ export class OutboxWorkerService { .set({ status: 'failed', lastError, processedAt: new Date() }) .where(eq(schedulerOutbox.id, record.id)); this.logger.error( - `Outbox record ${record.id} exhausted ${MAX_OUTBOX_ATTEMPTS} attempts ` + - `(action=${record.action} stitch=${record.stitchId}): ${lastError}`, + `Outbox record permanently failed: id=${record.id} action=${record.action} ` + + `stitchId=${record.stitchId} attempts=${record.attempts}/${MAX_OUTBOX_ATTEMPTS} ` + + `error="${lastError}"`, ); } else { - // Exponential back-off: 2s, 4s, 8s, 16s, 32s. + // Exponential back-off between attempts: 2s, 4s, 8s, 16s, 32s. const delayMs = Math.pow(2, record.attempts) * 1_000; const nextRetryAt = new Date(Date.now() + delayMs); await this.db @@ -138,9 +141,9 @@ export class OutboxWorkerService { .set({ status: 'pending', lastError, nextRetryAt }) .where(eq(schedulerOutbox.id, record.id)); this.logger.warn( - `Outbox record ${record.id} failed (attempt ${record.attempts}/${MAX_OUTBOX_ATTEMPTS}) ` + - `(action=${record.action} stitch=${record.stitchId}) — ` + - `retry at ${nextRetryAt.toISOString()}: ${lastError}`, + `Outbox record will retry: id=${record.id} action=${record.action} ` + + `stitchId=${record.stitchId} attempt=${record.attempts}/${MAX_OUTBOX_ATTEMPTS} ` + + `nextRetryAt=${nextRetryAt.toISOString()} error="${lastError}"`, ); } } diff --git a/apps/api/src/modules/scheduler/poll-sync-runner.ts b/apps/api/src/modules/scheduler/poll-sync-runner.ts index 5ae38d9f..c4e2855a 100644 --- a/apps/api/src/modules/scheduler/poll-sync-runner.ts +++ b/apps/api/src/modules/scheduler/poll-sync-runner.ts @@ -22,14 +22,20 @@ import { } from '@nexiom/engine'; import { PieceRegistryService } from '../trigger/piece-registry.service.js'; import { SyncRunner, type SyncResult } from './sync-runner.js'; +import { pollLockKey } from './lock-keys.js'; // --------------------------------------------------------------------------- // Internal helpers // --------------------------------------------------------------------------- -/** Redis key for a per-stream poll lock. */ -function lockKey(stitchId: string, streamName: string): string { - return `lock:poll:${stitchId}:${streamName}`; +/** + * Casts OAuthCredentialBlob to the generic Record the Piece interface expects. + * Centralised here so any future Piece interface change only needs one update. + */ +function toCredentialsRecord( + credentials: OAuthCredentialBlob, +): Record { + return credentials as unknown as Record; } /** Starting high-water mark for a stream (before any pages are processed). */ @@ -176,7 +182,7 @@ export class PollSyncRunner extends SyncRunner { piece: Piece, credentials: OAuthCredentialBlob, ): Promise { - const key = lockKey(stitchId, descriptor.streamName); + const key = pollLockKey(stitchId, descriptor.streamName); // Use 2× the sync interval as the lock TTL so that a slow poll run that // approaches the full interval does not lose the lock mid-pagination. // Floor at 5 minutes to protect very short intervals. @@ -271,7 +277,7 @@ export class PollSyncRunner extends SyncRunner { } const page = await piece.poll!( - credentials as unknown as Record, + toCredentialsRecord(credentials), streamName, window, nextCursor, @@ -413,7 +419,7 @@ export class PollSyncRunner extends SyncRunner { ): Promise { if (typeof piece.describeStreams === 'function') { const streams = await piece.describeStreams( - credentials as unknown as Record, + toCredentialsRecord(credentials), ); const match = streams.find((s) => s.streamName === sourceObject); if (match) return match; diff --git a/apps/api/src/modules/scheduler/scheduler.controller.spec.ts b/apps/api/src/modules/scheduler/scheduler.controller.spec.ts index 55f12b48..bf0fb7c9 100644 --- a/apps/api/src/modules/scheduler/scheduler.controller.spec.ts +++ b/apps/api/src/modules/scheduler/scheduler.controller.spec.ts @@ -66,13 +66,21 @@ describe('SchedulerController', () => { ).rejects.toBeInstanceOf(InternalServerErrorException); }); - it('propagates thrown errors from the service without wrapping', async () => { + it('wraps unexpected service throws in InternalServerErrorException to prevent raw error leakage', async () => { + // Raw errors (e.g. DB connection strings, vendor tokens) must never reach + // the Windmill worker response body. service.executeStitch.mockRejectedValue( - new Error(`Stitch not found: ${STITCH_ID}`), + new Error(`DB connection string: postgres://secret@host/db`), ); - await expect( - controller.executeStitch({ stitchId: STITCH_ID }), - ).rejects.toThrow(`Stitch not found: ${STITCH_ID}`); + const err = await controller + .executeStitch({ stitchId: STITCH_ID }) + .catch((e: unknown) => e); + + expect(err).toBeInstanceOf(InternalServerErrorException); + // The raw message must not be forwarded. + expect((err as InternalServerErrorException).message).not.toContain( + 'postgres://', + ); }); }); diff --git a/apps/api/src/modules/scheduler/scheduler.controller.ts b/apps/api/src/modules/scheduler/scheduler.controller.ts index 024b7db8..d3b2d5de 100644 --- a/apps/api/src/modules/scheduler/scheduler.controller.ts +++ b/apps/api/src/modules/scheduler/scheduler.controller.ts @@ -29,12 +29,18 @@ export class SchedulerController { @Post('execute-stitch') @HttpCode(HttpStatus.OK) async executeStitch(@Body() body: ExecuteStitchBody) { - const result = await this.schedulerService.executeStitch(body.stitchId); - if (result.status === 'failed') { - throw new InternalServerErrorException( - `Stitch execution failed: ${body.stitchId}`, - ); + // Wrap the entire call so that raw internal errors (DB connection strings, + // vendor tokens, stack traces) are never surfaced in the HTTP response body + // seen by Windmill workers. + try { + const result = await this.schedulerService.executeStitch(body.stitchId); + if (result.status === 'failed') { + throw new InternalServerErrorException('Stitch execution failed'); + } + return result; + } catch (err) { + if (err instanceof InternalServerErrorException) throw err; + throw new InternalServerErrorException('Stitch execution failed'); } - return result; } } diff --git a/apps/api/src/modules/scheduler/scheduler.module.ts b/apps/api/src/modules/scheduler/scheduler.module.ts index 2c5fb110..50383f4a 100644 --- a/apps/api/src/modules/scheduler/scheduler.module.ts +++ b/apps/api/src/modules/scheduler/scheduler.module.ts @@ -1,14 +1,19 @@ import { Module } from '@nestjs/common'; import { ConfigService } from '@nestjs/config'; import { CursorManagerService } from '@nexiom/engine'; +import { TokenManagerService } from '@nexiom/connectors'; +import { REDIS_CLIENT, type Redis } from '@nexiom/cache'; +import { DATABASE_CONNECTION, type DrizzleDb } from '@nexiom/database'; import { DbModule } from '../../db/db.module.js'; import { ConnectionsModule } from '../connections/connections.module.js'; import { PiecesModule } from '../pieces/pieces.module.js'; +import { PieceRegistryService } from '../trigger/piece-registry.service.js'; import { WindmillClient } from './windmill.client.js'; import { HttpWindmillClient } from './http-windmill.client.js'; import { StubWindmillClient } from './stub-windmill.client.js'; import { SyncRunner } from './sync-runner.js'; import { PollSyncRunner } from './poll-sync-runner.js'; +import { StubSyncRunner } from './stub-sync-runner.js'; import { SchedulerService } from './scheduler.service.js'; import { SchedulerController } from './scheduler.controller.js'; import { CursorResetController } from './cursor-reset.controller.js'; @@ -30,7 +35,41 @@ import { OutboxWorkerService } from './outbox-worker.service.js'; }, }, CursorManagerService, - { provide: SyncRunner, useClass: PollSyncRunner }, + { + // Only instantiate the full PollSyncRunner (with its Redis + token deps) + // when Windmill is enabled. In local dev (WINDMILL_ENABLED=false) the + // stub is returned so a missing Redis or credential provider does not + // crash the process on startup. + provide: SyncRunner, + inject: [ + ConfigService, + DATABASE_CONNECTION, + REDIS_CLIENT, + TokenManagerService, + PieceRegistryService, + CursorManagerService, + ], + useFactory: ( + config: ConfigService, + db: DrizzleDb, + redis: Redis, + tokenManager: TokenManagerService, + pieceRegistry: PieceRegistryService, + cursorManager: CursorManagerService, + ): SyncRunner => { + if (config.get('WINDMILL_ENABLED') === 'true') { + return new PollSyncRunner( + db, + redis, + config, + tokenManager, + pieceRegistry, + cursorManager, + ); + } + return new StubSyncRunner(); + }, + }, SchedulerService, InternalSchedulerGuard, OutboxWorkerService, diff --git a/apps/api/src/modules/scheduler/scheduler.service.spec.ts b/apps/api/src/modules/scheduler/scheduler.service.spec.ts index 1867304a..75d78591 100644 --- a/apps/api/src/modules/scheduler/scheduler.service.spec.ts +++ b/apps/api/src/modules/scheduler/scheduler.service.spec.ts @@ -43,7 +43,7 @@ describe('SchedulerService', () => { const mockSyncRunner = { run: vi .fn() - .mockResolvedValue({ stitchId: STITCH_ID, status: 'started' }), + .mockResolvedValue({ stitchId: STITCH_ID, status: 'succeeded' }), }; const module = await Test.createTestingModule({ @@ -220,7 +220,7 @@ describe('SchedulerService', () => { describe('executeStitch', () => { it('delegates to SyncRunner and returns the result', async () => { const result = await service.executeStitch(STITCH_ID); - expect(result).toEqual({ stitchId: STITCH_ID, status: 'started' }); + expect(result).toEqual({ stitchId: STITCH_ID, status: 'succeeded' }); }); }); }); diff --git a/apps/api/src/modules/scheduler/stub-sync-runner.ts b/apps/api/src/modules/scheduler/stub-sync-runner.ts index 30d2853f..bf5aa69a 100644 --- a/apps/api/src/modules/scheduler/stub-sync-runner.ts +++ b/apps/api/src/modules/scheduler/stub-sync-runner.ts @@ -2,8 +2,8 @@ import { Injectable, Logger } from '@nestjs/common'; import { SyncRunner, type SyncResult } from './sync-runner.js'; /** - * StubSyncRunner — no-op placeholder used in unit tests. - * Returns 'started' immediately without executing the sync pipeline. + * StubSyncRunner — no-op placeholder used when WINDMILL_ENABLED=false. + * Returns 'succeeded' immediately without executing the sync pipeline. */ @Injectable() export class StubSyncRunner extends SyncRunner { @@ -11,6 +11,6 @@ export class StubSyncRunner extends SyncRunner { run(stitchId: string): Promise { this.logger.debug(`StubSyncRunner: run stitchId=${stitchId} (no-op)`); - return Promise.resolve({ stitchId, status: 'started' }); + return Promise.resolve({ stitchId, status: 'succeeded' }); } } diff --git a/docs/architecture/master/tasks.md b/docs/architecture/master/tasks.md index ebd3cddc..7a3599bd 100644 --- a/docs/architecture/master/tasks.md +++ b/docs/architecture/master/tasks.md @@ -430,16 +430,31 @@ Each task is one commit (or one small PR). Checkboxes track completion. ### T048 · api: `CursorResetEndpoint` — admin full-refresh trigger -- [x] `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 +- [x] `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 scheduled run - [x] Superadmin guard only (`AuthGuard` + `SystemAdminGuard`); logs the reset with operator email for audit -- [x] `GET /admin/stitches/:id/cursors` — lists all `sync_cursors` rows for the stitch enriched with `ageMs` and `stale` flag (`ageMs > 2 × syncIntervalMinutes × 60_000`); throws `NotFoundException` if stitch missing +- [x] `GET /admin/stitches/:id/cursors` — lists all `sync_cursors` rows enriched with `ageMs`, `stale`, and `paused` flags; `scheduleEnabled=false` stitches never marked stale; column allowlist excludes `stateDocument` (may contain vendor tokens); throws `NotFoundException` if stitch missing - [x] `streamName` validated against `/^[\w.-]{1,200}$/`; invalid values rejected with `BadRequestException`; value JSON-encoded in audit log to prevent log injection +- [x] DELETE atomically acquires the per-stream Redis lock (SET NX) before deleting to close TOCTOU race with `PollSyncRunner`; returns `409 ConflictException` if lock is held; lock released in `finally` +- [x] Lock key format centralised in `lock-keys.ts` (shared by `PollSyncRunner` and `CursorResetController`) - [x] `staleThresholdMs` guarded for non-positive `syncIntervalMinutes` (returns `Infinity`) -- [x] `row.updatedAt.getTime()` used directly (Drizzle returns JS `Date` for `timestamptz`) -- [x] Unit tests for both endpoints (9 tests: DELETE row exists, DELETE idempotent, DELETE invalid streamName x3, GET stale true, GET stale false, GET empty cursors, GET not found); test mock rewritten to table-based dispatch -- Files: `apps/api/src/modules/scheduler/cursor-reset.controller.ts`, `apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts` +- [x] Unit tests: 13 tests (atomic NX lock, lock release on DB error, ConflictException, streamName validation ×3, stale/paused/fresh/empty cursor scenarios, NotFoundException) +- Files: `cursor-reset.controller.ts`, `cursor-reset.controller.spec.ts`, `lock-keys.ts` - Depends: T046, T029 +### Enterprise-grade quality pass (all scheduler module files) + +- [x] **C1** `outbox-worker.service.ts`: `MAX_OUTBOX_ATTEMPTS` changed from 5 → 6 (1 initial + 5 retries) to match documented back-off schedule "2s, 4s, 8s, 16s, 32s"; tests updated for new boundary (attempts=6 permanently fails, attempts=5 retries with 32s delay) +- [x] **C3** `http-windmill.client.ts`: 409 on `ensureStitchScript` now logs WARN with explanation; removes silent swallow; dead `requestJson` method removed (W7) +- [x] **W1** `scheduler.module.ts`: `SyncRunner` binding converted from `useClass` to conditional `useFactory`; returns `StubSyncRunner` when `WINDMILL_ENABLED=false` to avoid heavyweight dependency instantiation on dev startup +- [x] **W2** `scheduler.controller.ts`: `executeStitch` wrapped in try/catch; raw service errors (DB strings, tokens) never reach Windmill worker HTTP response +- [x] **W4** `poll-sync-runner.ts`: `credentials as unknown as Record` double-cast centralised into `toCredentialsRecord()` helper +- [x] **W5** `lock-keys.ts`: `pollLockKey()` extracted to shared module; `PollSyncRunner` and `CursorResetController` both import from it +- [x] **W6** `http-windmill.client.ts`: `Content-Type: application/json` only set when request has a body (GET requests omit it) +- [x] **S1** `packages/database/src/schema/stitches.ts` + `drizzle/0014_scheduler_outbox_partial_index.sql`: composite index on `(status, next_retry_at)` upgraded to partial index on `(next_retry_at)` WHERE `status='pending'` +- [x] **S2** `stub-sync-runner.ts`: changed terminal status from `'started'` → `'succeeded'` to match real runner +- [x] **S3** `interval-to-cron.ts`: comment added explaining 6-field Quartz cron support in Windmill's Rust `cron` crate +- [x] **S5** `outbox-worker.service.ts`: log messages include structured context fields (`id=`, `action=`, `stitchId=`, `attempts=`, `nextRetryAt=`, `error=`) + --- ## Phase 4 — Route Intelligence Dashboard diff --git a/packages/database/src/schema/stitches.ts b/packages/database/src/schema/stitches.ts index 26d06f62..ba5a68dd 100644 --- a/packages/database/src/schema/stitches.ts +++ b/packages/database/src/schema/stitches.ts @@ -262,8 +262,12 @@ export const schedulerOutbox = pgTable('scheduler_outbox', { foreignColumns: [integrationStitches.id], name: 'scheduler_outbox_stitch_fk', }).onDelete('cascade'), - // Primary poll query: WHERE status = 'pending' AND next_retry_at <= NOW() - index('scheduler_outbox_poll_idx').on(table.status, table.nextRetryAt), + // Partial index covering only pending rows — excludes the large succeeded/failed + // population so the poll query (WHERE status='pending' AND next_retry_at<=NOW()) + // stays fast as the table grows. + index('scheduler_outbox_poll_idx') + .on(table.nextRetryAt) + .where(sql`status = 'pending'`), index('scheduler_outbox_stitch_idx').on(table.stitchId), ]); From 047002bf95939e037e2775309f02ecce931fdd47 Mon Sep 17 00:00:00 2001 From: Pramod Date: Wed, 25 Mar 2026 12:23:05 +0530 Subject: [PATCH 3/5] cursor reset endpoint code review --- .claude/agent-memory/code-reviewer/MEMORY.md | 2 +- .../scheduler/cursor-reset.controller.spec.ts | 67 +++++++++++++------ .../scheduler/cursor-reset.controller.ts | 28 ++++++-- .../modules/scheduler/http-windmill.client.ts | 12 +++- .../scheduler/outbox-worker.service.ts | 18 ++++- .../scheduler/poll-sync-runner.spec.ts | 13 ++-- .../src/modules/scheduler/poll-sync-runner.ts | 17 +++-- .../scheduler/scheduler.controller.spec.ts | 15 ++++- .../modules/scheduler/scheduler.controller.ts | 6 +- .../src/modules/scheduler/scheduler.module.ts | 38 +++++------ docs/architecture/master/tasks.md | 8 ++- 11 files changed, 158 insertions(+), 66 deletions(-) diff --git a/.claude/agent-memory/code-reviewer/MEMORY.md b/.claude/agent-memory/code-reviewer/MEMORY.md index 2291d7e2..c1669be0 100644 --- a/.claude/agent-memory/code-reviewer/MEMORY.md +++ b/.claude/agent-memory/code-reviewer/MEMORY.md @@ -38,7 +38,7 @@ - `PollSyncRunner` implements Singer-style poll with Redis NX locks, crash-resume via `currently_syncing` + `bookmark.offset` - State persisted in `public.sync_cursors` (stitchId + streamName composite key, JSONB stateDocument) - `CursorManagerService` in `packages/engine` -- stateless, computes windows + tracks HWM -- Lock key: `lock:poll:${stitchId}:${streamName}`, TTL = max(syncIntervalMinutes * 2 * 60_000, 5 * 60_000)ms +- Lock key: `lock:poll:${stitchId}:${streamName}`, TTL = `max(syncIntervalMinutes * 2 * 60_000, 5 * 60_000)`ms - Lock uses Lua-atomic renew/release with owner token verification - `OAuthCredentialBlob` double-cast to `Record` is a recurring pattern in piece calls - Test mock pattern: `makeDb()` with table-aware `.from()` dispatching (improved from callCount) diff --git a/apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts b/apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts index 8ed0afb8..6bdc459a 100644 --- a/apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts +++ b/apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts @@ -34,7 +34,8 @@ function createMockRedis() { return { // SET NX: returns 'OK' (lock acquired) or null (already held) set: vi.fn().mockResolvedValue('OK'), - del: vi.fn().mockResolvedValue(1), + // Lua compare-and-delete used to release the lock + eval: vi.fn().mockResolvedValue(1), }; } @@ -69,7 +70,7 @@ describe('CursorResetController', () => { // ── DELETE ────────────────────────────────────────────────────────────── - it('DELETE acquires lock, deletes cursor, releases lock', async () => { + it('DELETE acquires lock with a UUID token, deletes cursor, releases lock via compare-and-delete', async () => { mockRedis.set.mockResolvedValue('OK'); mockDb.delete.mockReturnValue(mockDb); mockDb.where.mockResolvedValue(undefined); @@ -81,23 +82,48 @@ describe('CursorResetController', () => { ); expect(result).toBeUndefined(); - // Lock acquired then released - expect(mockRedis.set).toHaveBeenCalledOnce(); - expect(mockRedis.del).toHaveBeenCalledOnce(); + + // SET NX called with a UUID token (not a static string) + const setArgs = mockRedis.set.mock.calls[0] as [ + string, + string, + string, + number, + string, + ]; + expect(setArgs[0]).toBe(`lock:poll:${STITCH_ID}:${STREAM_NAME}`); + expect(setArgs[1]).toMatch( + /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i, + ); + expect(setArgs[2]).toBe('PX'); + expect(typeof setArgs[3]).toBe('number'); + expect(setArgs[4]).toBe('NX'); + + // Compare-and-delete called with the same token + expect(mockRedis.eval).toHaveBeenCalledOnce(); + const evalArgs = mockRedis.eval.mock.calls[0] as [ + string, + number, + string, + string, + ]; + expect(evalArgs[2]).toBe(`lock:poll:${STITCH_ID}:${STREAM_NAME}`); + expect(evalArgs[3]).toBe(setArgs[1]); // same token passed to both SET and eval + expect(mockDb.delete).toHaveBeenCalledOnce(); }); - it('DELETE is idempotent — no error when row is absent', async () => { + it('DELETE is idempotent — resolves with undefined when row is absent', async () => { mockRedis.set.mockResolvedValue('OK'); mockDb.delete.mockReturnValue(mockDb); mockDb.where.mockResolvedValue(undefined); await expect( controller.deleteCursor(STITCH_ID, STREAM_NAME, mockCtx), - ).resolves.not.toThrow(); + ).resolves.toBeUndefined(); }); - it('DELETE releases lock even when DB delete throws', async () => { + it('DELETE releases lock via compare-and-delete even when DB delete throws', async () => { mockRedis.set.mockResolvedValue('OK'); mockDb.delete.mockReturnValue(mockDb); mockDb.where.mockRejectedValue(new Error('DB error')); @@ -106,8 +132,8 @@ describe('CursorResetController', () => { controller.deleteCursor(STITCH_ID, STREAM_NAME, mockCtx), ).rejects.toThrow('DB error'); - // Lock must be released in finally - expect(mockRedis.del).toHaveBeenCalledOnce(); + // Compare-and-delete must run in finally + expect(mockRedis.eval).toHaveBeenCalledOnce(); }); it('DELETE throws ConflictException when poll lock is already held (NX fails)', async () => { @@ -120,23 +146,20 @@ describe('CursorResetController', () => { // Must not attempt DB delete expect(mockDb.delete).not.toHaveBeenCalled(); // Must not attempt lock release (never acquired it) - expect(mockRedis.del).not.toHaveBeenCalled(); + expect(mockRedis.eval).not.toHaveBeenCalled(); }); - it('DELETE uses the correct lock key format', async () => { + it('DELETE uses the correct lock key format and NX flag', async () => { mockRedis.set.mockResolvedValue('OK'); mockDb.delete.mockReturnValue(mockDb); mockDb.where.mockResolvedValue(undefined); await controller.deleteCursor(STITCH_ID, STREAM_NAME, mockCtx); - expect(mockRedis.set).toHaveBeenCalledWith( - `lock:poll:${STITCH_ID}:${STREAM_NAME}`, - 'admin-reset', - 'PX', - expect.any(Number), - 'NX', - ); + const setArgs = mockRedis.set.mock.calls[0] as unknown[]; + expect(setArgs[0]).toBe(`lock:poll:${STITCH_ID}:${STREAM_NAME}`); + expect(setArgs[2]).toBe('PX'); + expect(setArgs[4]).toBe('NX'); }); it('DELETE rejects streamName with newline characters', async () => { @@ -171,11 +194,15 @@ describe('CursorResetController', () => { { syncIntervalMinutes: 30, scheduleEnabled: true }, ]); } else if (table === syncCursors) { + // Include stateDocument in the raw DB row to prove the controller strips it. mockDb.where.mockResolvedValueOnce([ { id: 'c1', stitchId: STITCH_ID, streamName: STREAM_NAME, + stateDocument: { + bookmarks: { Account: { replication_key_value: 'secret-token' } }, + }, createdAt: sixtyOneMinutesAgo, updatedAt: sixtyOneMinutesAgo, }, @@ -190,7 +217,7 @@ describe('CursorResetController', () => { expect(result[0].stale).toBe(true); expect(result[0].paused).toBe(false); expect(result[0].ageMs).toBeGreaterThan(2 * 30 * 60_000); - // stateDocument must not be present in the response + // stateDocument must be stripped — it may contain opaque vendor cursor tokens. expect(result[0]).not.toHaveProperty('stateDocument'); }); diff --git a/apps/api/src/modules/scheduler/cursor-reset.controller.ts b/apps/api/src/modules/scheduler/cursor-reset.controller.ts index 438bd09b..53ff435a 100644 --- a/apps/api/src/modules/scheduler/cursor-reset.controller.ts +++ b/apps/api/src/modules/scheduler/cursor-reset.controller.ts @@ -23,6 +23,7 @@ import { import type { DrizzleDb } from '@nexiom/database'; import { REDIS_CLIENT, type Redis } from '@nexiom/cache'; import { eq, and } from 'drizzle-orm'; +import { randomUUID } from 'node:crypto'; import { pollLockKey } from './lock-keys.js'; /** Short-lived TTL (ms) for the admin reset lock — long enough to cover the DB delete. */ @@ -56,9 +57,13 @@ export class CursorResetController { // the run completes. Holding the lock during the delete prevents a new poll // from starting and immediately re-creating the cursor row we just removed. const lockKey = pollLockKey(id, streamName); + // Use a unique token per acquisition so the compare-and-delete release + // cannot accidentally free a lock held by a concurrent PollSyncRunner + // that acquired it between our NX set and our eventual release. + const lockToken = randomUUID(); const acquired = await this.redis.set( lockKey, - 'admin-reset', + lockToken, 'PX', ADMIN_RESET_LOCK_TTL_MS, 'NX', @@ -79,8 +84,17 @@ export class CursorResetController { ), ); } finally { - // Release immediately — we only needed the lock to prevent concurrent starts. - await this.redis.del(lockKey); + // Compare-and-delete: only remove the lock if we still own it. + // Prevents releasing a lock that expired and was re-acquired by a + // PollSyncRunner during a slow DB delete. + const lua = [ + "if redis.call('get', KEYS[1]) == ARGV[1] then", + " return redis.call('del', KEYS[1])", + 'else', + ' return 0', + 'end', + ].join('\n'); + await this.redis.eval(lua, 1, lockKey, lockToken); } this.logger.log( @@ -126,8 +140,14 @@ export class CursorResetController { const ageMs = now - row.updatedAt.getTime(); // Paused stitches are never stale — cursors are not expected to advance. const paused = !stitch.scheduleEnabled; + // Explicitly enumerate fields rather than spreading — stateDocument is + // intentionally absent and must never appear in the HTTP response. return { - ...row, + id: row.id, + stitchId: row.stitchId, + streamName: row.streamName, + createdAt: row.createdAt, + updatedAt: row.updatedAt, ageMs, paused, stale: paused ? false : ageMs > staleThresholdMs, diff --git a/apps/api/src/modules/scheduler/http-windmill.client.ts b/apps/api/src/modules/scheduler/http-windmill.client.ts index 804ace8f..8468bc8d 100644 --- a/apps/api/src/modules/scheduler/http-windmill.client.ts +++ b/apps/api/src/modules/scheduler/http-windmill.client.ts @@ -78,15 +78,23 @@ export class HttpWindmillClient extends WindmillClient { return; } if (res.status === 409) { - await res.body?.cancel(); // Windmill returns 409 when an identical content hash already exists at this // path (truly idempotent). If STITCH_RUNNER_CONTENT changed since the last // deploy, the hash differs and Windmill creates a new version (200). // A persistent 409 after a content change indicates the Windmill workspace // needs a manual redeploy (delete the script at the path and redeploy). + // Read a bounded snippet (≤1 KB) for diagnostics — discard the rest. + let snippet = ''; + try { + const full = await res.text(); + snippet = full.length > 1024 ? `${full.slice(0, 1024)}…` : full; + } catch { + // Ignore body-read errors — the 409 itself is sufficient signal. + } this.logger.warn( `Stitch-runner script at ${STITCH_RUNNER_PATH} returned 409 — ` + - 'script content matches an existing version or a manual redeploy is needed.', + `script content matches an existing version or a manual redeploy is needed` + + (snippet ? `. Response: ${snippet}` : '.'), ); return; } diff --git a/apps/api/src/modules/scheduler/outbox-worker.service.ts b/apps/api/src/modules/scheduler/outbox-worker.service.ts index b8577398..4788eadc 100644 --- a/apps/api/src/modules/scheduler/outbox-worker.service.ts +++ b/apps/api/src/modules/scheduler/outbox-worker.service.ts @@ -27,6 +27,16 @@ const BATCH_SIZE = 20; * MAX_OUTBOX_ATTEMPTS, after which the record is marked 'failed' for * human/alerting review. */ +/** + * Returns a log-safe version of an error message. + * Truncates to 200 characters and strips URL credentials + * (e.g. https://user:token@host) that may appear in vendor API errors. + */ +function sanitizeError(message: string): string { + const stripped = message.replaceAll(/\/\/[^@\s]*@/g, '//[REDACTED]@'); + return stripped.length > 200 ? `${stripped.slice(0, 200)}…` : stripped; +} + @Injectable() export class OutboxWorkerService { private readonly logger = new Logger(OutboxWorkerService.name); @@ -119,7 +129,11 @@ export class OutboxWorkerService { record: typeof schedulerOutbox.$inferSelect, err: unknown, ): Promise { + // Full error text is persisted to the DB column for human/alerting review. + // Logs only emit a sanitized snippet — raw vendor error messages may contain + // OAuth tokens, connection strings, or other sensitive material. const lastError = err instanceof Error ? err.message : String(err); + const safeError = sanitizeError(lastError); if (record.attempts >= MAX_OUTBOX_ATTEMPTS) { // Permanently failed — mark for alerting/human review. @@ -130,7 +144,7 @@ export class OutboxWorkerService { this.logger.error( `Outbox record permanently failed: id=${record.id} action=${record.action} ` + `stitchId=${record.stitchId} attempts=${record.attempts}/${MAX_OUTBOX_ATTEMPTS} ` + - `error="${lastError}"`, + `error="${safeError}"`, ); } else { // Exponential back-off between attempts: 2s, 4s, 8s, 16s, 32s. @@ -143,7 +157,7 @@ export class OutboxWorkerService { this.logger.warn( `Outbox record will retry: id=${record.id} action=${record.action} ` + `stitchId=${record.stitchId} attempt=${record.attempts}/${MAX_OUTBOX_ATTEMPTS} ` + - `nextRetryAt=${nextRetryAt.toISOString()} error="${lastError}"`, + `nextRetryAt=${nextRetryAt.toISOString()} error="${safeError}"`, ); } } diff --git a/apps/api/src/modules/scheduler/poll-sync-runner.spec.ts b/apps/api/src/modules/scheduler/poll-sync-runner.spec.ts index 0f588dc1..0795a600 100644 --- a/apps/api/src/modules/scheduler/poll-sync-runner.spec.ts +++ b/apps/api/src/modules/scheduler/poll-sync-runner.spec.ts @@ -1,5 +1,6 @@ import { describe, it, expect, vi, beforeEach } from 'vitest'; import { Test } from '@nestjs/testing'; +import { NotFoundException, BadRequestException } from '@nestjs/common'; import { ConfigService } from '@nestjs/config'; import { DATABASE_CONNECTION, @@ -421,35 +422,39 @@ describe('PollSyncRunner', () => { // ── Error paths ──────────────────────────────────────────────────────────── - it('throws when the stitch does not exist', async () => { + it('throws NotFoundException when the stitch does not exist', async () => { const noStitch = makeDb({ stitch: null }); await build({ db: noStitch }); + await expect(runner.run(STITCH_ID)).rejects.toThrow(NotFoundException); await expect(runner.run(STITCH_ID)).rejects.toThrow( `Stitch not found: ${STITCH_ID}`, ); }); - it('throws when the connection does not exist', async () => { + it('throws NotFoundException when the connection does not exist', async () => { const noConn = makeDb({ connection: null }); await build({ db: noConn }); + await expect(runner.run(STITCH_ID)).rejects.toThrow(NotFoundException); await expect(runner.run(STITCH_ID)).rejects.toThrow( `Connection not found: ${CONN_ID}`, ); }); - it('throws when the piece is not registered', async () => { + it('throws BadRequestException when the piece is not registered', async () => { pieceRegistry.getPiece.mockReturnValue(undefined); + await expect(runner.run(STITCH_ID)).rejects.toThrow(BadRequestException); await expect(runner.run(STITCH_ID)).rejects.toThrow( 'Piece not registered: "salesforce"', ); }); - it('throws when the piece does not support poll()', async () => { + it('throws BadRequestException when the piece does not support poll()', async () => { pieceRegistry.getPiece.mockReturnValue({ name: 'salesforce' }); + await expect(runner.run(STITCH_ID)).rejects.toThrow(BadRequestException); await expect(runner.run(STITCH_ID)).rejects.toThrow( 'does not support polling', ); diff --git a/apps/api/src/modules/scheduler/poll-sync-runner.ts b/apps/api/src/modules/scheduler/poll-sync-runner.ts index c4e2855a..baccf819 100644 --- a/apps/api/src/modules/scheduler/poll-sync-runner.ts +++ b/apps/api/src/modules/scheduler/poll-sync-runner.ts @@ -1,4 +1,10 @@ -import { Injectable, Logger, Inject } from '@nestjs/common'; +import { + Injectable, + Logger, + Inject, + NotFoundException, + BadRequestException, +} from '@nestjs/common'; import { ConfigService } from '@nestjs/config'; import { randomUUID } from 'node:crypto'; import { z } from 'zod'; @@ -331,7 +337,7 @@ export class PollSyncRunner extends SyncRunner { .from(integrationStitches) .where(eq(integrationStitches.id, stitchId)) .limit(1); - if (!stitch) throw new Error(`Stitch not found: ${stitchId}`); + if (!stitch) throw new NotFoundException(`Stitch not found: ${stitchId}`); return stitch; } @@ -344,7 +350,8 @@ export class PollSyncRunner extends SyncRunner { .from(appConnections) .where(eq(appConnections.id, connectionId)) .limit(1); - if (!conn) throw new Error(`Connection not found: ${connectionId}`); + if (!conn) + throw new NotFoundException(`Connection not found: ${connectionId}`); return conn; } @@ -402,10 +409,10 @@ export class PollSyncRunner extends SyncRunner { private resolvePiece(appName: string): Piece { const piece = this.pieceRegistry.getPiece(appName); if (!piece) { - throw new Error(`Piece not registered: "${appName}"`); + throw new BadRequestException(`Piece not registered: "${appName}"`); } if (typeof piece.poll !== 'function') { - throw new Error( + throw new BadRequestException( `Piece "${appName}" does not support polling (no poll() method)`, ); } diff --git a/apps/api/src/modules/scheduler/scheduler.controller.spec.ts b/apps/api/src/modules/scheduler/scheduler.controller.spec.ts index bf0fb7c9..dd1dcfe0 100644 --- a/apps/api/src/modules/scheduler/scheduler.controller.spec.ts +++ b/apps/api/src/modules/scheduler/scheduler.controller.spec.ts @@ -1,6 +1,9 @@ import { describe, it, expect, vi, beforeEach } from 'vitest'; import { Test } from '@nestjs/testing'; -import { InternalServerErrorException } from '@nestjs/common'; +import { + InternalServerErrorException, + NotFoundException, +} from '@nestjs/common'; import { SchedulerController } from './scheduler.controller.js'; import { SchedulerService } from './scheduler.service.js'; import { InternalSchedulerGuard } from './internal-scheduler.guard.js'; @@ -66,6 +69,16 @@ describe('SchedulerController', () => { ).rejects.toBeInstanceOf(InternalServerErrorException); }); + it('propagates NotFoundException from service as 404 (not wrapped as 500)', async () => { + service.executeStitch.mockRejectedValue( + new NotFoundException('Stitch not found: ' + STITCH_ID), + ); + + await expect( + controller.executeStitch({ stitchId: STITCH_ID }), + ).rejects.toBeInstanceOf(NotFoundException); + }); + it('wraps unexpected service throws in InternalServerErrorException to prevent raw error leakage', async () => { // Raw errors (e.g. DB connection strings, vendor tokens) must never reach // the Windmill worker response body. diff --git a/apps/api/src/modules/scheduler/scheduler.controller.ts b/apps/api/src/modules/scheduler/scheduler.controller.ts index d3b2d5de..0580da04 100644 --- a/apps/api/src/modules/scheduler/scheduler.controller.ts +++ b/apps/api/src/modules/scheduler/scheduler.controller.ts @@ -5,6 +5,7 @@ import { HttpCode, HttpStatus, UseGuards, + HttpException, InternalServerErrorException, } from '@nestjs/common'; import { InternalSchedulerGuard } from './internal-scheduler.guard.js'; @@ -39,7 +40,10 @@ export class SchedulerController { } return result; } catch (err) { - if (err instanceof InternalServerErrorException) throw err; + // Re-throw NestJS HTTP exceptions (NotFoundException, BadRequestException, + // etc.) unchanged so domain errors surface as the correct 4xx status. + // Only wrap truly unexpected errors as 500 to avoid leaking internals. + if (err instanceof HttpException) throw err; throw new InternalServerErrorException('Stitch execution failed'); } } diff --git a/apps/api/src/modules/scheduler/scheduler.module.ts b/apps/api/src/modules/scheduler/scheduler.module.ts index 50383f4a..c0cab5cf 100644 --- a/apps/api/src/modules/scheduler/scheduler.module.ts +++ b/apps/api/src/modules/scheduler/scheduler.module.ts @@ -1,9 +1,10 @@ import { Module } from '@nestjs/common'; +import { ModuleRef } from '@nestjs/core'; import { ConfigService } from '@nestjs/config'; import { CursorManagerService } from '@nexiom/engine'; import { TokenManagerService } from '@nexiom/connectors'; -import { REDIS_CLIENT, type Redis } from '@nexiom/cache'; -import { DATABASE_CONNECTION, type DrizzleDb } from '@nexiom/database'; +import { REDIS_CLIENT } from '@nexiom/cache'; +import { DATABASE_CONNECTION } from '@nexiom/database'; import { DbModule } from '../../db/db.module.js'; import { ConnectionsModule } from '../connections/connections.module.js'; import { PiecesModule } from '../pieces/pieces.module.js'; @@ -40,31 +41,22 @@ import { OutboxWorkerService } from './outbox-worker.service.js'; // when Windmill is enabled. In local dev (WINDMILL_ENABLED=false) the // stub is returned so a missing Redis or credential provider does not // crash the process on startup. + // + // Heavy dependencies (Redis, DB, TokenManager, etc.) are resolved lazily + // via ModuleRef so NestJS does not eagerly instantiate them when + // WINDMILL_ENABLED=false — prevents startup failures in environments where + // those providers are absent. provide: SyncRunner, - inject: [ - ConfigService, - DATABASE_CONNECTION, - REDIS_CLIENT, - TokenManagerService, - PieceRegistryService, - CursorManagerService, - ], - useFactory: ( - config: ConfigService, - db: DrizzleDb, - redis: Redis, - tokenManager: TokenManagerService, - pieceRegistry: PieceRegistryService, - cursorManager: CursorManagerService, - ): SyncRunner => { + inject: [ConfigService, ModuleRef], + useFactory: (config: ConfigService, moduleRef: ModuleRef): SyncRunner => { if (config.get('WINDMILL_ENABLED') === 'true') { return new PollSyncRunner( - db, - redis, + moduleRef.get(DATABASE_CONNECTION, { strict: false }), + moduleRef.get(REDIS_CLIENT, { strict: false }), config, - tokenManager, - pieceRegistry, - cursorManager, + moduleRef.get(TokenManagerService, { strict: false }), + moduleRef.get(PieceRegistryService, { strict: false }), + moduleRef.get(CursorManagerService, { strict: false }), ); } return new StubSyncRunner(); diff --git a/docs/architecture/master/tasks.md b/docs/architecture/master/tasks.md index 7a3599bd..f341e5ef 100644 --- a/docs/architecture/master/tasks.md +++ b/docs/architecture/master/tasks.md @@ -444,9 +444,11 @@ Each task is one commit (or one small PR). Checkboxes track completion. ### Enterprise-grade quality pass (all scheduler module files) - [x] **C1** `outbox-worker.service.ts`: `MAX_OUTBOX_ATTEMPTS` changed from 5 → 6 (1 initial + 5 retries) to match documented back-off schedule "2s, 4s, 8s, 16s, 32s"; tests updated for new boundary (attempts=6 permanently fails, attempts=5 retries with 32s delay) -- [x] **C3** `http-windmill.client.ts`: 409 on `ensureStitchScript` now logs WARN with explanation; removes silent swallow; dead `requestJson` method removed (W7) -- [x] **W1** `scheduler.module.ts`: `SyncRunner` binding converted from `useClass` to conditional `useFactory`; returns `StubSyncRunner` when `WINDMILL_ENABLED=false` to avoid heavyweight dependency instantiation on dev startup -- [x] **W2** `scheduler.controller.ts`: `executeStitch` wrapped in try/catch; raw service errors (DB strings, tokens) never reach Windmill worker HTTP response +- [x] **C2** `outbox-worker.service.ts`: `sanitizeError()` helper — logs emit truncated (≤200 char) URL-credential-stripped error message; full raw text persisted only to DB `last_error` column for human/alerting review +- [x] **C3** `http-windmill.client.ts`: 409 on `ensureStitchScript` now reads bounded body snippet (≤1 KB) and includes it in WARN log for diagnostics; dead `requestJson` method removed (W7) +- [x] **W1** `scheduler.module.ts`: `SyncRunner` factory changed to `inject: [ConfigService, ModuleRef]`; heavy deps (DB, Redis, TokenManager, PieceRegistry, CursorManager) resolved lazily via `moduleRef.get(..., { strict: false })` only when `WINDMILL_ENABLED=true` — prevents startup failures when those providers are absent +- [x] **W2** `scheduler.controller.ts`: catch block changed from `instanceof InternalServerErrorException` to `instanceof HttpException` so `NotFoundException`, `BadRequestException`, etc. propagate as the correct 4xx status instead of being wrapped as 500; test added for `NotFoundException` propagation +- [x] **W3** `poll-sync-runner.ts`: `loadStitch()` throws `NotFoundException`, `loadConnection()` throws `NotFoundException`, `resolvePiece()` throws `BadRequestException` — these now propagate through `SchedulerController` as the correct HTTP status codes; spec updated to assert exception types - [x] **W4** `poll-sync-runner.ts`: `credentials as unknown as Record` double-cast centralised into `toCredentialsRecord()` helper - [x] **W5** `lock-keys.ts`: `pollLockKey()` extracted to shared module; `PollSyncRunner` and `CursorResetController` both import from it - [x] **W6** `http-windmill.client.ts`: `Content-Type: application/json` only set when request has a body (GET requests omit it) From 2e491ad2120b79ebef50bce7181e9ea356331cd5 Mon Sep 17 00:00:00 2001 From: Pramod Date: Wed, 25 Mar 2026 12:43:44 +0530 Subject: [PATCH 4/5] cursor reset endpoint code review --- .../scheduler/cursor-reset.controller.spec.ts | 12 +++++------- .../modules/scheduler/cursor-reset.controller.ts | 11 ++++++++--- .../src/modules/scheduler/outbox-worker.service.ts | 13 ++++++------- apps/api/src/modules/scheduler/scheduler.module.ts | 14 ++++++++++++-- docs/architecture/master/tasks.md | 2 +- 5 files changed, 32 insertions(+), 20 deletions(-) diff --git a/apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts b/apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts index 6bdc459a..04f8eb17 100644 --- a/apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts +++ b/apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts @@ -149,17 +149,15 @@ describe('CursorResetController', () => { expect(mockRedis.eval).not.toHaveBeenCalled(); }); - it('DELETE uses the correct lock key format and NX flag', async () => { + it('DELETE accepts streamName at the maximum valid length (200 characters)', async () => { mockRedis.set.mockResolvedValue('OK'); mockDb.delete.mockReturnValue(mockDb); mockDb.where.mockResolvedValue(undefined); - await controller.deleteCursor(STITCH_ID, STREAM_NAME, mockCtx); - - const setArgs = mockRedis.set.mock.calls[0] as unknown[]; - expect(setArgs[0]).toBe(`lock:poll:${STITCH_ID}:${STREAM_NAME}`); - expect(setArgs[2]).toBe('PX'); - expect(setArgs[4]).toBe('NX'); + const maxName = 'a'.repeat(200); + await expect( + controller.deleteCursor(STITCH_ID, maxName, mockCtx), + ).resolves.toBeUndefined(); }); it('DELETE rejects streamName with newline characters', async () => { diff --git a/apps/api/src/modules/scheduler/cursor-reset.controller.ts b/apps/api/src/modules/scheduler/cursor-reset.controller.ts index 53ff435a..b79bf040 100644 --- a/apps/api/src/modules/scheduler/cursor-reset.controller.ts +++ b/apps/api/src/modules/scheduler/cursor-reset.controller.ts @@ -26,8 +26,11 @@ import { eq, and } from 'drizzle-orm'; import { randomUUID } from 'node:crypto'; import { pollLockKey } from './lock-keys.js'; -/** Short-lived TTL (ms) for the admin reset lock — long enough to cover the DB delete. */ -const ADMIN_RESET_LOCK_TTL_MS = 5_000; +// TTL for the admin cursor-reset lock. Must be long enough to cover the DB +// delete even under elevated database latency. Matches the minimum floor used +// by PollSyncRunner (5 minutes) so that a concurrent poll that starts just +// after we acquire the lock cannot complete and reacquire before we release. +const ADMIN_RESET_LOCK_TTL_MS = 5 * 60_000; @UseGuards(AuthGuard, SystemAdminGuard) @Controller('admin/stitches') @@ -97,8 +100,10 @@ export class CursorResetController { await this.redis.eval(lua, 1, lockKey, lockToken); } + // Use the stable internal principal ID rather than email (PII) in service + // logs. Audit trails that require email belong in a dedicated secure sink. this.logger.log( - `Cursor reset: stitchId=${id}, streamName=${JSON.stringify(streamName)}, operator=${ctx.user.email}`, + `Cursor reset: stitchId=${id}, streamName=${JSON.stringify(streamName)}, actorId=${ctx.user.id ?? '[unknown]'}`, ); } diff --git a/apps/api/src/modules/scheduler/outbox-worker.service.ts b/apps/api/src/modules/scheduler/outbox-worker.service.ts index 4788eadc..61c15792 100644 --- a/apps/api/src/modules/scheduler/outbox-worker.service.ts +++ b/apps/api/src/modules/scheduler/outbox-worker.service.ts @@ -129,11 +129,10 @@ export class OutboxWorkerService { record: typeof schedulerOutbox.$inferSelect, err: unknown, ): Promise { - // Full error text is persisted to the DB column for human/alerting review. - // Logs only emit a sanitized snippet — raw vendor error messages may contain - // OAuth tokens, connection strings, or other sensitive material. - const lastError = err instanceof Error ? err.message : String(err); - const safeError = sanitizeError(lastError); + // Sanitize before persisting or logging — raw vendor error messages may + // contain OAuth tokens, connection strings, or other sensitive material. + const rawError = err instanceof Error ? err.message : String(err); + const lastError = sanitizeError(rawError); if (record.attempts >= MAX_OUTBOX_ATTEMPTS) { // Permanently failed — mark for alerting/human review. @@ -144,7 +143,7 @@ export class OutboxWorkerService { this.logger.error( `Outbox record permanently failed: id=${record.id} action=${record.action} ` + `stitchId=${record.stitchId} attempts=${record.attempts}/${MAX_OUTBOX_ATTEMPTS} ` + - `error="${safeError}"`, + `error="${lastError}"`, ); } else { // Exponential back-off between attempts: 2s, 4s, 8s, 16s, 32s. @@ -157,7 +156,7 @@ export class OutboxWorkerService { this.logger.warn( `Outbox record will retry: id=${record.id} action=${record.action} ` + `stitchId=${record.stitchId} attempt=${record.attempts}/${MAX_OUTBOX_ATTEMPTS} ` + - `nextRetryAt=${nextRetryAt.toISOString()} error="${safeError}"`, + `nextRetryAt=${nextRetryAt.toISOString()} error="${lastError}"`, ); } } diff --git a/apps/api/src/modules/scheduler/scheduler.module.ts b/apps/api/src/modules/scheduler/scheduler.module.ts index c0cab5cf..3ad3126f 100644 --- a/apps/api/src/modules/scheduler/scheduler.module.ts +++ b/apps/api/src/modules/scheduler/scheduler.module.ts @@ -1,4 +1,4 @@ -import { Module } from '@nestjs/common'; +import { Module, type Type } from '@nestjs/common'; import { ModuleRef } from '@nestjs/core'; import { ConfigService } from '@nestjs/config'; import { CursorManagerService } from '@nexiom/engine'; @@ -21,9 +21,19 @@ import { CursorResetController } from './cursor-reset.controller.js'; import { InternalSchedulerGuard } from './internal-scheduler.guard.js'; import { OutboxWorkerService } from './outbox-worker.service.js'; +// Evaluated once at module load time — env vars are set before app bootstrap. +const WINDMILL_ENABLED = process.env['WINDMILL_ENABLED'] === 'true'; + +// CursorResetController injects REDIS_CLIENT; only register it when Redis is +// expected to be present (i.e. WINDMILL_ENABLED=true). + +const schedulerControllers: Type[] = WINDMILL_ENABLED + ? [SchedulerController, CursorResetController] + : [SchedulerController]; + @Module({ imports: [DbModule, ConnectionsModule, PiecesModule], - controllers: [SchedulerController, CursorResetController], + controllers: schedulerControllers, providers: [ { provide: WindmillClient, diff --git a/docs/architecture/master/tasks.md b/docs/architecture/master/tasks.md index f341e5ef..923b7290 100644 --- a/docs/architecture/master/tasks.md +++ b/docs/architecture/master/tasks.md @@ -455,7 +455,7 @@ Each task is one commit (or one small PR). Checkboxes track completion. - [x] **S1** `packages/database/src/schema/stitches.ts` + `drizzle/0014_scheduler_outbox_partial_index.sql`: composite index on `(status, next_retry_at)` upgraded to partial index on `(next_retry_at)` WHERE `status='pending'` - [x] **S2** `stub-sync-runner.ts`: changed terminal status from `'started'` → `'succeeded'` to match real runner - [x] **S3** `interval-to-cron.ts`: comment added explaining 6-field Quartz cron support in Windmill's Rust `cron` crate -- [x] **S5** `outbox-worker.service.ts`: log messages include structured context fields (`id=`, `action=`, `stitchId=`, `attempts=`, `nextRetryAt=`, `error=`) +- [x] **S4** `outbox-worker.service.ts`: log messages include structured context fields (`id=`, `action=`, `stitchId=`, `attempts=`, `nextRetryAt=`, `error=`) --- From c54ee3964485bbffb0d39218fb58bc1653af9c4c Mon Sep 17 00:00:00 2001 From: Pramod Date: Wed, 25 Mar 2026 14:03:18 +0530 Subject: [PATCH 5/5] cursor reset endpoint code review --- apps/api/src/modules/scheduler/outbox-worker.service.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/apps/api/src/modules/scheduler/outbox-worker.service.ts b/apps/api/src/modules/scheduler/outbox-worker.service.ts index 61c15792..9eec71af 100644 --- a/apps/api/src/modules/scheduler/outbox-worker.service.ts +++ b/apps/api/src/modules/scheduler/outbox-worker.service.ts @@ -155,7 +155,7 @@ export class OutboxWorkerService { .where(eq(schedulerOutbox.id, record.id)); this.logger.warn( `Outbox record will retry: id=${record.id} action=${record.action} ` + - `stitchId=${record.stitchId} attempt=${record.attempts}/${MAX_OUTBOX_ATTEMPTS} ` + + `stitchId=${record.stitchId} attempts=${record.attempts}/${MAX_OUTBOX_ATTEMPTS} ` + `nextRetryAt=${nextRetryAt.toISOString()} error="${lastError}"`, ); }