Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

#27 Implement ETA (estimated time of arrival) on client-side and worker-side #54

Open
wants to merge 5 commits into
base: master
Choose a base branch
from
Open
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
22 changes: 22 additions & 0 deletions examples/eta/client.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
"use strict";
const celery = require("../../dist");

const client = celery.createClient("redis://", "redis://");
// client.conf.TASK_PROTOCOL = 1;

const task = client.createTask("tasks.add");

Promise.all([
task
.applyAsync([1, 2], {}, {
countdown: 10,
})
.get()
.then(console.log),
task
.applyAsync([1, 2], {}, {
eta: new Date(Date.now() + 3 * 1000),
})
.get()
.then(console.log),
]).then(() => client.disconnect());
1,052 changes: 413 additions & 639 deletions package-lock.json

Large diffs are not rendered by default.

34 changes: 17 additions & 17 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -34,29 +34,29 @@
},
"homepage": "https://github.com/actumn/celery.node#readme",
"devDependencies": {
"@types/amqplib": "^0.5.13",
"@types/chai": "^4.2.9",
"@types/ioredis": "^4.14.6",
"@types/mocha": "^7.0.1",
"@types/sinon": "^9.0.9",
"@types/uuid": "^3.4.7",
"@typescript-eslint/eslint-plugin": "^2.23.0",
"@typescript-eslint/parser": "^2.23.0",
"chai": "^4.2.0",
"eslint": "^7.16.0",
"eslint-config-prettier": "^7.1.0",
"@types/amqplib": "^0.5.17",
"@types/chai": "^4.2.15",
"@types/ioredis": "^4.22.0",
"@types/mocha": "^7.0.2",
"@types/sinon": "^9.0.11",
"@types/uuid": "^3.4.9",
"@typescript-eslint/eslint-plugin": "^2.34.0",
"@typescript-eslint/parser": "^2.34.0",
"chai": "^4.3.4",
"eslint": "^7.22.0",
"eslint-config-prettier": "^7.2.0",
"eslint-plugin-import": "^2.22.1",
"eslint-plugin-mocha": "^8.0.0",
"eslint-plugin-mocha": "^8.1.0",
"express": "^4.17.1",
"mocha": "^8.2.1",
"mocha": "^8.3.2",
"prettier": "^1.19.1",
"sinon": "^7.2.2",
"ts-node": "^8.10.2",
"typescript": "^3.7.2"
"typescript": "^3.9.9"
},
"dependencies": {
"amqplib": "^0.5.5",
"ioredis": "^4.14.0",
"uuid": "^3.3.2"
"amqplib": "^0.7.1",
"ioredis": "^4.24.2",
"uuid": "^3.4.0"
}
}
51 changes: 39 additions & 12 deletions src/app/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -44,17 +44,23 @@ export default class Client extends Base {
taskId: string,
taskName: string,
args?: Array<any>,
kwargs?: object
kwargs?: object,
countdown?: number,
eta?: Date,
): TaskMessage {
if (countdown) {
eta = new Date(Date.now() + countdown * 1000);
}

const message: TaskMessage = {
headers: {
lang: "js",
task: taskName,
id: taskId
id: taskId,
eta: eta?.toISOString(),
/*
'shadow': shadow,
'eta': eta,
'expires': expires,
'shadow': shadow,
'group': group_id,
'retries': retries,
'timelimit': [time_limit, soft_time_limit],
Expand All @@ -66,8 +72,8 @@ export default class Client extends Base {
*/
},
properties: {
correlationId: taskId,
replyTo: ""
correlation_id: taskId,
reply_to: "",
},
body: [args, kwargs, {}],
sentEvent: null
Expand All @@ -90,19 +96,38 @@ export default class Client extends Base {
taskId: string,
taskName: string,
args?: Array<any>,
kwargs?: object
kwargs?: object,
countdown?: number,
eta?: Date,
): TaskMessage {
// const expires = null;
if (countdown) {
eta = new Date(Date.now() + countdown * 1000);
}

const message: TaskMessage = {
headers: {},
properties: {
correlationId: taskId,
replyTo: ""
correlation_id: taskId,
reply_to: ""
},
body: {
task: taskName,
id: taskId,
args: args,
kwargs: kwargs
kwargs: kwargs,
eta: eta.toISOString(),
/*
'expires': expires,
'group': group_id,
'retries': retries,
'timelimit': [time_limit, soft_time_limit],
'root_id': root_id,
'parent_id': parent_id,
'argsrepr': argsrepr,
'kwargsrepr': kwargsrepr,
'origin': origin or anon_nodename()
*/
},
sentEvent: null
};
Expand Down Expand Up @@ -136,10 +161,12 @@ export default class Client extends Base {
taskName: string,
args?: Array<any>,
kwargs?: object,
taskId?: string
taskId?: string,
countdown?: number,
eta?: Date,
): AsyncResult {
taskId = taskId || v4();
const message = this.createTaskMessage(taskId, taskName, args, kwargs);
const message = this.createTaskMessage(taskId, taskName, args, kwargs, countdown, eta);
this.sendTaskMessage(taskName, message);

const result = new AsyncResult(taskId, this.backend);
Expand Down
36 changes: 34 additions & 2 deletions src/app/task.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,37 @@
import Client from "./client";
import { AsyncResult } from "./result";

export interface TaskOptions {

/**
* @member {string}
*/
taskId: string;

/**
* @member {number}
* Number of seconds into the future that the task should execute.
* Defaults to immediate execution.
*/
countdown: number;

/**
* @member {Date}
* Absolute time and date of when the task should be executed.
* May not be specified if `countdown` is also supplied.
*/
eta: Date;

// /**
// * @member {Number | Date}
// * Datetime or seconds in the future for the task should expire.
// * The task won't be executed after the expiration time.
// */
// expires: number | Date;

// TODO:: retry, retryPolicy
}

export default class Task {
client: Client;
name: string;
Expand Down Expand Up @@ -28,7 +59,7 @@ export default class Task {
return this.applyAsync([...args]);
}

public applyAsync(args: Array<any>, kwargs?: object): AsyncResult {
public applyAsync(args: Array<any>, kwargs?: object, taskOptions?: TaskOptions): AsyncResult {
if (args && !Array.isArray(args)) {
throw new Error("args is not array");
}
Expand All @@ -37,6 +68,7 @@ export default class Task {
throw new Error("kwargs is not object");
}

return this.client.sendTask(this.name, args || [], kwargs || {});

return this.client.sendTask(this.name, args || [], kwargs || {}, taskOptions?.taskId, taskOptions?.countdown, taskOptions?.eta);
}
}
35 changes: 23 additions & 12 deletions src/app/worker.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import Base from "./base";
import { Message } from "../kombu/message";
import { parseISO8601 } from "../util/time";

export default class Worker extends Base {
handlers: object = {};
Expand Down Expand Up @@ -136,25 +137,35 @@ export default class Worker extends Base {
throw new Error(`Missing process handler for task ${taskName}`);
}

const eta = headers.eta ? parseISO8601(headers.eta) - Date.now() : 0;

console.info(
`celery.node Received task: ${taskName}[${taskId}], args: ${args}, kwargs: ${JSON.stringify(
kwargs
)}`
);

const timeStart = process.hrtime();
const taskPromise = handler(...args, kwargs).then(result => {
const diff = process.hrtime(timeStart);
console.info(
`celery.node Task ${taskName}[${taskId}] succeeded in ${diff[0] +
diff[1] / 1e9}s: ${result}`
);
this.backend.storeResult(taskId, result, "SUCCESS");
this.activeTasks.delete(taskPromise);
}).catch(err => {
console.info(`celery.node Task ${taskName}[${taskId}] failed: [${err}]`);
this.backend.storeResult(taskId, err, "FAILURE");
this.activeTasks.delete(taskPromise);
const taskPromise = new Promise((resolve, reject) => {
setTimeout(() => {
handler(...args, kwargs)
.then(result => {
const diff = process.hrtime(timeStart);
console.info(
`celery.node Task ${taskName}[${taskId}] succeeded in ${diff[0] +
diff[1] / 1e9}s: ${result}`
);
this.backend.storeResult(taskId, result, "SUCCESS");
this.activeTasks.delete(taskPromise);
resolve();
})
.catch(err => {
console.info(`celery.node Task ${taskName}[${taskId}] failed: [${err}]`);
this.backend.storeResult(taskId, err, "FAILURE");
this.activeTasks.delete(taskPromise);
reject();
});
}, eta ?? 0 );
});

// record the executing task
Expand Down
40 changes: 40 additions & 0 deletions src/util/time.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
// https://gist.github.com/dimkalinux/1478408
const ISO8601_REGEX = /^(\d{4})-(\d\d)-(\d\d)([T ](\d\d):(\d\d):(\d\d)(\.\d+)?(Z|([+-])(\d\d)(:(\d\d))?)?)?$/;

export function parseISO8601(dateString: string): number {
// 0 = whole string
// 1 = year
// 2 = month
// 3 = day
// 4 = whole time part
// 5 = hour
// 6 = minute
// 7 = second
// 8 = fractional (with dot)
// 9 = whole timezone (possibly Z)
// 10 = offset sign (+ or -)
// 11 = offset hours
// 12 = offset minutes (with colon)
// 13 = offset minutes

const r = ISO8601_REGEX.exec(dateString);
if (!r)
return Date.parse(dateString);

const year = Number(r[1]);
const month = Number(r[2]) - 1;
const day = Number(r[3]);
if (!r[4])
return Date.UTC(year, month, day);

const hour = Number(r[5]);
const minute = Number(r[6]);
const second = Number(r[7]);
const ms = r[8]? Number((r[8] + "000").substr(1, 3)): 0;
if (!r[9])
return Date.UTC(year, month, day, hour, minute, second, ms);

const oh = r[11]? Number(r[10]) + Number(r[11]): 0;
const om = r[13]? Number(r[10]) + Number(r[13]): 0;
return Date.UTC(year, month, day, hour - oh, minute - om, second, ms);
}
11 changes: 11 additions & 0 deletions test/util/test_time.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
import { assert } from "chai";
import { parseISO8601 } from "../../src/util/time";

describe("time util test", () => {
describe("parseISO8601", () => {
it("asserts below", () => {
assert.equal(1618622133596, parseISO8601("2021-04-17T01:15:33.596988"));
assert.equal(1618622133596, parseISO8601("2021-04-17T01:15:33.596Z"));
});
});
});