Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 7 additions & 2 deletions .claude/agent-memory/code-reviewer/MEMORY.md
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, unknown>` 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

Expand Down
17 changes: 17 additions & 0 deletions apps/api/drizzle/0014_scheduler_outbox_partial_index.sql
Original file line number Diff line number Diff line change
@@ -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';
314 changes: 314 additions & 0 deletions apps/api/src/modules/scheduler/cursor-reset.controller.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,314 @@
import { describe, it, expect, vi, beforeEach } from 'vitest';
import { Test } from '@nestjs/testing';
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';

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() {
const chain = {
select: vi.fn().mockReturnThis(),
from: vi.fn().mockReturnThis(),
where: vi.fn().mockReturnThis(),
limit: vi.fn().mockReturnThis(),
delete: vi.fn().mockReturnThis(),
};
return chain;
}

function createMockRedis() {
return {
// SET NX: returns 'OK' (lock acquired) or null (already held)
set: vi.fn().mockResolvedValue('OK'),
// Lua compare-and-delete used to release the lock
eval: vi.fn().mockResolvedValue(1),
};
}

type MockDb = ReturnType<typeof createMockDb>;
type MockRedis = ReturnType<typeof createMockRedis>;

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 },
{ provide: REDIS_CLIENT, useValue: mockRedis },
],
})
.overrideGuard(AuthGuard)
.useValue({ canActivate: () => true })
.overrideGuard(SystemAdminGuard)
.useValue({ canActivate: () => true })
.compile();

controller = module.get(CursorResetController);
vi.clearAllMocks();
});

// ── DELETE ──────────────────────────────────────────────────────────────

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);

const result = await controller.deleteCursor(
STITCH_ID,
STREAM_NAME,
mockCtx,
);

expect(result).toBeUndefined();

// 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 — 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.toBeUndefined();
});

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'));

await expect(
controller.deleteCursor(STITCH_ID, STREAM_NAME, mockCtx),
).rejects.toThrow('DB error');

// Compare-and-delete must run in finally
expect(mockRedis.eval).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.eval).not.toHaveBeenCalled();
});

it('DELETE accepts streamName at the maximum valid length (200 characters)', async () => {
mockRedis.set.mockResolvedValue('OK');
mockDb.delete.mockReturnValue(mockDb);
mockDb.where.mockResolvedValue(undefined);

const maxName = 'a'.repeat(200);
await expect(
controller.deleteCursor(STITCH_ID, maxName, mockCtx),
).resolves.toBeUndefined();
});

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);

mockDb.from.mockImplementation((table: unknown) => {
if (table === integrationStitches) {
mockDb.where.mockReturnValueOnce(mockDb);
mockDb.limit.mockResolvedValueOnce([
{ 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,
},
]);
}
return mockDb;
});

const result = await controller.listCursors(STITCH_ID);

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 be stripped — it may contain opaque vendor cursor tokens.
expect(result[0]).not.toHaveProperty('stateDocument');
Comment thread
coderabbitai[bot] marked this conversation as resolved.
});

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, scheduleEnabled: true },
]);
} else if (table === syncCursors) {
mockDb.where.mockResolvedValueOnce([
{
id: 'c1',
stitchId: STITCH_ID,
streamName: STREAM_NAME,
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].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, scheduleEnabled: true },
]);
} 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,
);
});
});
Loading
Loading