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
94 changes: 94 additions & 0 deletions packages/backends/backend-test/src/updateJob.ts
Original file line number Diff line number Diff line change
Expand Up @@ -134,4 +134,98 @@ export default function defineUpdateJobTestSuite() {
await expect(backend.updateJob({ id: -1 })).rejects.toThrow();
});
});

describe("updateJobIfCurrent", () => {
it("should return the current job when a matching update changes no values", async () => {
const insertedJob = await backend.createNewJob({
queue: "default",
class: "TestJob",
args: [],
constructor_args: [],
state: "waiting",
script: "test.js",
attempt: 0,
});

const updatedJob = await backend.updateJobIfCurrent(
{ id: insertedJob.id },
{
state: insertedJob.state,
attempt: insertedJob.attempt,
claimed_by: insertedJob.claimed_by,
claimed_at: insertedJob.claimed_at,
attempted_at: insertedJob.attempted_at,
},
);

expect(updatedJob).toMatchObject(insertedJob);
});

it("should update a job whose execution fingerprint still matches", async () => {
const insertedJob = await backend.createNewJob({
queue: "default",
class: "TestJob",
args: [],
constructor_args: [],
state: "waiting",
script: "test.js",
attempt: 0,
});

const updatedJob = await backend.updateJobIfCurrent(
{ id: insertedJob.id, state: "claimed", claimed_by: "worker-a", claimed_at: new Date() },
{
state: insertedJob.state,
attempt: insertedJob.attempt,
claimed_by: insertedJob.claimed_by,
claimed_at: insertedJob.claimed_at,
attempted_at: insertedJob.attempted_at,
},
);

expect(updatedJob).toMatchObject({ state: "claimed", claimed_by: "worker-a" });
});

it("should not update a newer execution that has the same state", async () => {
const insertedJob = await backend.createNewJob({
queue: "default",
class: "TestJob",
args: [],
constructor_args: [],
state: "waiting",
script: "test.js",
attempt: 0,
});
const oldExecution = await backend.updateJob({
...insertedJob,
state: "running",
attempt: 1,
claimed_by: "worker-a",
claimed_at: new Date(2000, 0, 1),
attempted_at: new Date(2000, 0, 1),
});
const newExecution = await backend.updateJob({
...oldExecution,
state: "running",
attempt: 2,
claimed_by: "worker-b",
claimed_at: new Date(2000, 0, 2),
attempted_at: new Date(2000, 0, 2),
});

const updatedJob = await backend.updateJobIfCurrent(
{ ...oldExecution, state: "waiting" },
{
state: oldExecution.state,
attempt: oldExecution.attempt,
claimed_by: oldExecution.claimed_by,
claimed_at: oldExecution.claimed_at,
attempted_at: oldExecution.attempted_at,
},
);

expect(updatedJob).toBeUndefined();
expect(await backend.getJob(insertedJob.id)).toMatchObject(newExecution);
});
});
}
2 changes: 1 addition & 1 deletion packages/backends/backend/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ This package provides the foundational types, interfaces, and abstract base clas

### Type Definitions

- **Job Types** - `NewJobData`, `UpdateJobData`, `JobCounts` for job operations
- **Job Types** - `NewJobData`, `UpdateJobData`, `JobExecutionFingerprint`, `JobCounts` for job operations
- **Queue Types** - `NewQueueData`, `UpdateQueueData` for queue management
- **Configuration** - `BackendConfig` for backend driver configuration

Expand Down
17 changes: 17 additions & 0 deletions packages/backends/backend/src/backend.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,14 @@ export type NewJobData = Pick<JobData, "queue" | "script" | "class" | "args" | "
*/
export type UpdateJobData = Pick<JobData, "id"> & Partial<Omit<JobData, "id">>;

/**
* Fields that identify the exact execution snapshot a job update was based on.
*
* A backend uses this fingerprint as an atomic update precondition so an old
* worker or stale-job sweep cannot overwrite a newer lifecycle transition.
*/
export type JobExecutionFingerprint = Pick<JobData, "state" | "attempt" | "claimed_by" | "claimed_at" | "attempted_at">;

/**
* Data required to create a new queue.
*/
Expand Down Expand Up @@ -168,6 +176,15 @@ export interface Backend {
*/
updateJob(job: UpdateJobData): Promise<JobData>;

/**
* Updates a job only if its persisted execution fingerprint is unchanged.
*
* @param job The updated job data.
* @param expected The execution fingerprint observed before computing the update.
* @returns The updated job, or undefined if the job no longer matches the expected fingerprint.
*/
updateJobIfCurrent(job: UpdateJobData, expected: JobExecutionFingerprint): Promise<JobData | undefined>;

/**
* Lists jobs with optional filters.
* @param params Optional filter parameters. Where string arrays, they are treated as OR conditions.
Expand Down
15 changes: 15 additions & 0 deletions packages/backends/backend/src/lazy-backend.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ const mockBackend: Backend = vi.hoisted(() => ({
createNewJob: vi.fn().mockResolvedValue({ id: 1 } as JobData),
claimPendingJob: vi.fn().mockResolvedValue([{ id: 1 }] as JobData[]),
updateJob: vi.fn().mockResolvedValue({ id: 1 } as JobData),
updateJobIfCurrent: vi.fn().mockResolvedValue({ id: 1 } as JobData),
listJobs: vi.fn().mockResolvedValue([{ id: 1 }] as JobData[]),
countJobs: vi.fn().mockResolvedValue({ total: 1 } as JobCounts),
countJobsOverTime: vi.fn().mockResolvedValue([{ timestamp: new Date(), total: 1 }]),
Expand Down Expand Up @@ -140,6 +141,20 @@ describe("LazyBackend", () => {
expect(result).toEqual({ id: 1 });
});

it("should proxy updateJobIfCurrent", async () => {
const job = { id: 1 } as UpdateJobData;
const expected = {
state: "running" as const,
attempt: 1,
claimed_by: "worker-a",
claimed_at: new Date(),
attempted_at: new Date(),
};
const result = await lazyBackend.updateJobIfCurrent(job, expected);
expect(mockBackend.updateJobIfCurrent).toHaveBeenCalledWith(job, expected);
expect(result).toEqual({ id: 1 });
});

it("should proxy listJobs", async () => {
const params = { queue: "queue1" };
const result = await lazyBackend.listJobs(params);
Expand Down
15 changes: 14 additions & 1 deletion packages/backends/backend/src/lazy-backend.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,13 @@
import { JobData, JobState, QueueConfig } from "@sidequest/core";
import { Backend, JobCounts, NewJobData, NewQueueData, UpdateJobData, UpdateQueueData } from "./backend";
import {
Backend,
JobCounts,
JobExecutionFingerprint,
NewJobData,
NewQueueData,
UpdateJobData,
UpdateQueueData,
} from "./backend";
import { BackendConfig } from "./config";
import { createBackendFromDriver } from "./factory";

Expand Down Expand Up @@ -121,6 +129,11 @@ export class LazyBackend implements Backend {
return this.backend!.updateJob(job);
}

async updateJobIfCurrent(job: UpdateJobData, expected: JobExecutionFingerprint): Promise<JobData | undefined> {
await this.init();
return this.backend!.updateJobIfCurrent(job, expected);
}

async listJobs(params?: {
queue?: string | string[];
jobClass?: string | string[];
Expand Down
32 changes: 31 additions & 1 deletion packages/backends/backend/src/sql-backend.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,15 @@ import { DuplicatedJobError, JobData, JobState, logger, QueueConfig } from "@sid
import { Knex } from "knex";
import { hostname } from "os";
import { inspect } from "util";
import { Backend, JobCounts, NewJobData, NewQueueData, UpdateJobData, UpdateQueueData } from "./backend";
import {
Backend,
JobCounts,
JobExecutionFingerprint,
NewJobData,
NewQueueData,
UpdateJobData,
UpdateQueueData,
} from "./backend";
import { JOB_FALLBACK, MISC_FALLBACK, QUEUE_FALLBACK } from "./constants";
import { formatDateForBucket, safeParseJobData, whereOrWhereIn } from "./utils";

Expand Down Expand Up @@ -215,6 +223,28 @@ export abstract class SQLBackend implements Backend {
return safeParseJobData(updated);
}

async updateJobIfCurrent(job: UpdateJobData, expected: JobExecutionFingerprint): Promise<JobData | undefined> {
const data = {
...job,
args: job.args ? JSON.stringify(job.args) : job.args,
constructor_args: job.constructor_args ? JSON.stringify(job.constructor_args) : job.constructor_args,
result: job.result ? JSON.stringify(job.result) : job.result,
errors: job.errors ? JSON.stringify(job.errors) : job.errors,
uniqueness_config: job.uniqueness_config ? JSON.stringify(job.uniqueness_config) : job.uniqueness_config,
};

logger("Backend").debug(`Conditionally updating job: ${inspect(data)}`);
const [updated] = (await this.knex("sidequest_jobs")
.where({ id: job.id, ...expected })
.update(data)
.returning("*")) as JobData[];

if (!updated) return undefined;

logger("Backend").debug(`Job updated successfully: ${inspect(updated)}`);
return safeParseJobData(updated);
}

async listJobs(params?: {
queue?: string | string[];
jobClass?: string | string[];
Expand Down
8 changes: 8 additions & 0 deletions packages/backends/mongo/src/mongo-backend.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import {
formatDateForBucket,
JOB_FALLBACK,
JobCounts,
JobExecutionFingerprint,
NewJobData,
NewQueueData,
QUEUE_FALLBACK,
Expand Down Expand Up @@ -199,6 +200,13 @@ export default class MongoBackend implements Backend {
return res as JobData;
}

async updateJobIfCurrent(job: UpdateJobData, expected: JobExecutionFingerprint): Promise<JobData | undefined> {
await this.ensureConnected();
const { id, ...updates } = job;
const res = await this.jobs.findOneAndUpdate({ id, ...expected }, { $set: updates }, { returnDocument: "after" });
return (res as JobData | null) ?? undefined;
}

async listJobs(params?: {
queue?: string | string[];
jobClass?: string | string[];
Expand Down
29 changes: 29 additions & 0 deletions packages/backends/mysql/src/mysql-backend.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import {
JOB_FALLBACK,
JobExecutionFingerprint,
NewJobData,
NewQueueData,
QUEUE_FALLBACK,
Expand Down Expand Up @@ -155,6 +156,34 @@ export default class MysqlBackend extends SQLBackend {
return updatedJob;
}

async updateJobIfCurrent(job: UpdateJobData, expected: JobExecutionFingerprint): Promise<JobData | undefined> {
const data = {
...job,
args: job.args ? JSON.stringify(job.args) : job.args,
constructor_args: job.constructor_args ? JSON.stringify(job.constructor_args) : job.constructor_args,
result: job.result ? JSON.stringify(job.result) : job.result,
errors: job.errors ? JSON.stringify(job.errors) : job.errors,
uniqueness_config: job.uniqueness_config ? JSON.stringify(job.uniqueness_config) : job.uniqueness_config,
};
logger("Backend").debug(`Conditionally updating job: ${inspect(data)}`);

const updatedJob = await this.knex.transaction(async (trx) => {
const updatedCount = await trx("sidequest_jobs")
.where({ id: job.id, ...expected })
.update(data);

if (updatedCount === 0) return undefined;

const updated = await trx<JobData>("sidequest_jobs").where({ id: job.id }).first();
if (!updated) return undefined;

return safeParseJobData(updated);
});

if (updatedJob) logger("Backend").debug(`Job updated successfully: ${inspect(updatedJob)}`);
return updatedJob;
}

async listJobs(params?: {
queue?: string | string[];
jobClass?: string | string[];
Expand Down
30 changes: 23 additions & 7 deletions packages/docs/production/backends.md
Original file line number Diff line number Diff line change
Expand Up @@ -317,8 +317,16 @@ The backend class must be exported as a default export from the module. Sideques
### Implementing the Backend Interface

```typescript
import { Backend, JobData, NewJobData, UpdateJobData, JobCounts } from "@sidequest/backend";
import { JobState, QueueConfig } from "@sidequest/core";
import {
Backend,
JobCounts,
JobExecutionFingerprint,
NewJobData,
NewQueueData,
UpdateJobData,
UpdateQueueData,
} from "@sidequest/backend";
import { JobData, JobState, QueueConfig } from "@sidequest/core";

export class MyCustomBackend implements Backend {
// Required methods to implement
Expand Down Expand Up @@ -372,6 +380,13 @@ export class MyCustomBackend implements Backend {
// Update job data
}

async updateJobIfCurrent(
job: UpdateJobData,
expected: JobExecutionFingerprint,
): Promise<JobData | undefined> {
// Atomically update only when the persisted execution fingerprint matches
}

async listJobs(params?: {
queue?: string | string[];
jobClass?: string | string[];
Expand Down Expand Up @@ -471,11 +486,12 @@ The backend driver is dynamically loaded based on the `driver` string. It will b
When creating a custom backend, ensure:

1. **Atomic job claiming**: Jobs must be claimed atomically to prevent race conditions
2. **Transaction support**: Use transactions for data consistency
3. **Index optimization**: Add appropriate indexes for job and queue queries
4. **Error handling**: Proper error handling and connection management
5. **Migration support**: Implement schema versioning and migrations
6. **JSON serialization**: Handle complex job arguments and results properly
2. **Conditional job updates**: Execution transitions must compare and update atomically so stale workers cannot overwrite newer attempts
3. **Transaction support**: Use transactions for data consistency
4. **Index optimization**: Add appropriate indexes for job and queue queries
5. **Error handling**: Proper error handling and connection management
6. **Migration support**: Implement schema versioning and migrations
7. **JSON serialization**: Handle complex job arguments and results properly

### Testing Your Backend

Expand Down
Loading
Loading