Skip to content

Commit 57c2f1f

Browse files
committed
cron
1 parent de5f849 commit 57c2f1f

8 files changed

Lines changed: 3585 additions & 3450 deletions

File tree

‎openapi/v1.json‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2188,4 +2188,4 @@
21882188
"description": "Dry-run schedule expressions"
21892189
}
21902190
]
2191-
}
2191+
}

‎sdk/ts/cron/ARCHITECTURE.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -106,7 +106,7 @@ Only `__cron_*` topic targets are touched. Schedules targeting queues, user-defi
106106

107107
### Management client
108108

109-
`createCronClient()` / the lazy `cron` singleton provide direct schedule management without a handler or HTTP listener. All methods follow the same `cmd()` → `wrap()` → `camelize` pipeline as `@beyond.dev/queue` and `@beyond.dev/events`. `schedules.sync()` is a client-side operation: upsert all provided specs in parallel, then list all schedules and delete those not in the desired set. There is no server-side sync endpoint.
109+
`createCronClient()` / the lazy `cron` singleton provide direct schedule management without a handler or HTTP listener. All methods follow the same `cmd()` → `wrap()` → `camelize` pipeline as `@beyond.dev/queue` and `@beyond.dev/events`.
110110

111111
## State Machine
112112

Lines changed: 77 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,77 @@
1+
import { createServer } from "node:http";
2+
import { describe, expect, it } from "vitest";
3+
import { createCronClient } from "../src/client.js";
4+
import { findFreePort } from "./harness.js";
5+
6+
type MockResponse = { status: number; body: string };
7+
8+
async function withMockServer(
9+
responses: MockResponse[],
10+
fn: (url: string, requestCount: () => number) => Promise<void>,
11+
) {
12+
const port = await findFreePort();
13+
let count = 0;
14+
const server = createServer((_, res) => {
15+
const r = responses[Math.min(count++, responses.length - 1)]!;
16+
res.writeHead(r.status, { "content-type": "application/json" });
17+
res.end(r.body);
18+
});
19+
await new Promise<void>((r) => server.listen(port, "127.0.0.1", r));
20+
try {
21+
await fn(`http://127.0.0.1:${port}`, () => count);
22+
} finally {
23+
await new Promise<void>((r) => server.close(() => r()));
24+
}
25+
}
26+
27+
describe("buildFetch — retry behavior", () => {
28+
it("retries on 5xx and returns the successful response", async () => {
29+
await withMockServer(
30+
[
31+
{
32+
status: 503,
33+
body: "{\"error\":{\"code\":\"unavailable\",\"message\":\"down\"}}",
34+
},
35+
{ status: 200, body: "{\"name\":\"t\",\"status\":\"active\"}" },
36+
],
37+
async (url, requestCount) => {
38+
const client = createCronClient({ url, retries: 2 });
39+
const { data, error } = await client.schedules.get("t");
40+
expect(requestCount()).toBe(2);
41+
expect(error).toBeUndefined();
42+
expect(data?.name).toBe("t");
43+
},
44+
);
45+
});
46+
47+
it("returns error after exhausting retries", async () => {
48+
await withMockServer(
49+
[{
50+
status: 503,
51+
body: "{\"error\":{\"code\":\"unavailable\",\"message\":\"down\"}}",
52+
}],
53+
async (url, requestCount) => {
54+
const client = createCronClient({ url, retries: 2 });
55+
const { error } = await client.schedules.get("t");
56+
expect(requestCount()).toBe(3); // 1 initial + 2 retries
57+
expect(error?.status).toBe(503);
58+
expect(error?.code).toBe("unavailable");
59+
},
60+
);
61+
});
62+
63+
it("does not retry on 4xx errors", async () => {
64+
await withMockServer(
65+
[{
66+
status: 404,
67+
body: "{\"error\":{\"code\":\"not_found\",\"message\":\"missing\"}}",
68+
}],
69+
async (url, requestCount) => {
70+
const client = createCronClient({ url, retries: 2 });
71+
const { error } = await client.schedules.get("t");
72+
expect(requestCount()).toBe(1);
73+
expect(error?.status).toBe(404);
74+
},
75+
);
76+
});
77+
});

‎sdk/ts/cron/__tests__/schedules.test.ts‎

Lines changed: 6 additions & 39 deletions
Original file line numberDiff line numberDiff line change
@@ -179,55 +179,20 @@ describe("schedules — preview", () => {
179179
expect(error).toBeUndefined();
180180
expect(typeof data?.humanReadable).toBe("string");
181181
});
182-
});
183-
184-
describe("schedules — sync", () => {
185-
const client = cronClient();
186-
187-
it("sync upserts all specs and removes others", async () => {
188-
const keep = uniqueName("keep");
189-
const remove = uniqueName("remove");
190182

191-
// Seed a schedule that should be removed
192-
await client.schedules.upsert({
193-
name: remove,
194-
every: "1h",
195-
target: { queue: "test-q", message: {} },
183+
it("preview returns humanReadable for fireAt one-shot", async () => {
184+
const { data, error } = await client.schedules.preview({
185+
fireAt: "2099-01-01T00:00:00Z",
196186
});
197-
198-
const { data, error } = await client.schedules.sync([
199-
{
200-
name: keep,
201-
every: "1h",
202-
target: { queue: "test-q", message: {} },
203-
},
204-
]);
205187
expect(error).toBeUndefined();
206-
expect(data?.upserted).toBe(1);
207-
expect(data?.removed).toBeGreaterThanOrEqual(1);
208-
209-
const { data: all } = await client.schedules.list();
210-
expect(all?.some((s) => s.name === keep)).toBe(true);
211-
expect(all?.some((s) => s.name === remove)).toBe(false);
212-
213-
// Cleanup
214-
await client.schedules.delete(keep);
188+
expect(typeof data?.humanReadable).toBe("string");
215189
});
216190
});
217191

218192
describe("schedules — observability hooks", () => {
219193
it("fires onRequest and onResponse for each call", async () => {
220194
const requests: string[] = [];
221195
const responses: string[] = [];
222-
const c = cronClient();
223-
// Wrap with hooks via a new client instance
224-
const hooked = {
225-
...c,
226-
schedules: {
227-
...c.schedules,
228-
},
229-
};
230-
void hooked; // just verifying createCronClient accepts hooks
231196
const name = uniqueName();
232197
const hClient = (await import("../src/client.js")).createCronClient({
233198
url: process.env["QUEUE_TEST_URL"]!,
@@ -241,6 +206,8 @@ describe("schedules — observability hooks", () => {
241206
});
242207
await hClient.schedules.delete(name);
243208
expect(requests).toContain("schedules.upsert");
209+
expect(requests).toContain("schedules.delete");
244210
expect(responses).toContain("schedules.upsert");
211+
expect(responses).toContain("schedules.delete");
245212
});
246213
});

‎sdk/ts/cron/__tests__/start.test.ts‎

Lines changed: 146 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,53 @@
11
import { afterEach, beforeEach, describe, expect, it } from "vitest";
2-
import { start } from "../src/client.js";
2+
import { schedule, start } from "../src/client.js";
33
import { cronClient, findFreePort, uniqueName } from "./harness.js";
44

5+
describe("start() — input validation", () => {
6+
const serverUrl = process.env["QUEUE_TEST_URL"]!;
7+
8+
it("throws for an invalid job name", async () => {
9+
const job = schedule({
10+
name: "Bad Name",
11+
every: "1h",
12+
run: async () => {},
13+
});
14+
await expect(start([job], { url: serverUrl })).rejects.toThrow(
15+
/Invalid schedule name/,
16+
);
17+
});
18+
19+
it("throws for duplicate job names", async () => {
20+
const job = schedule({ name: "my-job", every: "1h", run: async () => {} });
21+
await expect(start([job, job], { url: serverUrl })).rejects.toThrow(
22+
/Duplicate schedule name/,
23+
);
24+
});
25+
26+
it("throws when BEYOND_QUEUE_URL is missing and no url option provided", async () => {
27+
const job = schedule({ name: "my-job", every: "1h", run: async () => {} });
28+
const saved = process.env["BEYOND_QUEUE_URL"];
29+
delete process.env["BEYOND_QUEUE_URL"];
30+
try {
31+
await expect(start([job])).rejects.toThrow(/BEYOND_QUEUE_URL/);
32+
} finally {
33+
if (saved !== undefined) process.env["BEYOND_QUEUE_URL"] = saved;
34+
}
35+
});
36+
37+
it("throws when BEYOND_INTERNAL_URL is missing", async () => {
38+
const job = schedule({ name: "my-job", every: "1h", run: async () => {} });
39+
const saved = process.env["BEYOND_INTERNAL_URL"];
40+
delete process.env["BEYOND_INTERNAL_URL"];
41+
try {
42+
await expect(start([job], { url: serverUrl })).rejects.toThrow(
43+
/BEYOND_INTERNAL_URL/,
44+
);
45+
} finally {
46+
if (saved !== undefined) process.env["BEYOND_INTERNAL_URL"] = saved;
47+
}
48+
});
49+
});
50+
551
describe("start() — worker lifecycle", () => {
652
const client = cronClient();
753
const serverUrl = process.env["QUEUE_TEST_URL"]!;
@@ -250,6 +296,105 @@ describe("start() — worker lifecycle", () => {
250296
}
251297
});
252298

299+
it("passes CronContext with correct name, ISO-8601 scheduledFor, and boolean outOfBand", async () => {
300+
const name = uniqueName("ctx");
301+
cleanupNames.push(name);
302+
const port = await findFreePort();
303+
const ac = new AbortController();
304+
let capturedCtx:
305+
| { name: string; scheduledFor: string; outOfBand: boolean }
306+
| undefined;
307+
308+
const done = start(
309+
[{
310+
name,
311+
spec: { name, every: "1h" },
312+
handler: async (ctx) => {
313+
capturedCtx = ctx;
314+
},
315+
}],
316+
{ url: serverUrl, port, signal: ac.signal },
317+
);
318+
319+
await new Promise<void>((r) => setTimeout(r, 200));
320+
321+
try {
322+
const res = await fetch(`http://127.0.0.1:${port}/__cron/${name}`, {
323+
method: "POST",
324+
body: "{}",
325+
headers: { "content-type": "application/json" },
326+
});
327+
expect(res.status).toBe(200);
328+
expect(capturedCtx?.name).toBe(name);
329+
expect(capturedCtx?.scheduledFor).toMatch(
330+
/^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}/,
331+
);
332+
expect(typeof capturedCtx?.outOfBand).toBe("boolean");
333+
} finally {
334+
ac.abort();
335+
await done;
336+
}
337+
});
338+
339+
it("returns 405 for non-POST requests", async () => {
340+
const name = uniqueName("405");
341+
cleanupNames.push(name);
342+
const port = await findFreePort();
343+
const ac = new AbortController();
344+
345+
const done = start(
346+
[{ name, spec: { name, every: "1h" }, handler: async () => {} }],
347+
{ url: serverUrl, port, signal: ac.signal },
348+
);
349+
350+
await new Promise<void>((r) => setTimeout(r, 200));
351+
352+
try {
353+
const res = await fetch(`http://127.0.0.1:${port}/__cron/${name}`, {
354+
method: "GET",
355+
});
356+
expect(res.status).toBe(405);
357+
} finally {
358+
ac.abort();
359+
await done;
360+
}
361+
});
362+
363+
it("resolves cleanly when aborted during an in-flight handler", async () => {
364+
const name = uniqueName("inflight");
365+
cleanupNames.push(name);
366+
const port = await findFreePort();
367+
const ac = new AbortController();
368+
let handlerStarted = false;
369+
370+
const done = start(
371+
[{
372+
name,
373+
spec: { name, every: "1h" },
374+
handler: async () => {
375+
handlerStarted = true;
376+
await new Promise<void>((r) => setTimeout(r, 500));
377+
},
378+
}],
379+
{ url: serverUrl, port, signal: ac.signal },
380+
);
381+
382+
await new Promise<void>((r) => setTimeout(r, 200));
383+
384+
// Fire the handler without awaiting the response
385+
fetch(`http://127.0.0.1:${port}/__cron/${name}`, {
386+
method: "POST",
387+
body: "{}",
388+
}).catch(() => {});
389+
390+
// Let the handler begin, then abort
391+
await new Promise<void>((r) => setTimeout(r, 50));
392+
expect(handlerStarted).toBe(true);
393+
394+
ac.abort();
395+
await expect(done).resolves.toBeUndefined();
396+
});
397+
253398
it("shuts down cleanly on signal abort", async () => {
254399
const name = uniqueName("shut");
255400
cleanupNames.push(name);

‎sdk/ts/cron/src/client.ts‎

Lines changed: 0 additions & 50 deletions
Original file line numberDiff line numberDiff line change
@@ -136,14 +136,6 @@ export interface CronClient {
136136
/** Create or replace a schedule (PUT semantics — idempotent). */
137137
upsert(spec: ScheduleSpec): CronResult<Schedule>;
138138

139-
/**
140-
* Upsert all provided specs and delete any schedules not in the list.
141-
* Always declarative — the provided list is the desired state.
142-
*/
143-
sync(
144-
specs: ScheduleSpec[],
145-
): CronResult<{ upserted: number; removed: number }>;
146-
147139
/** Dry-run a schedule expression. No schedule is created. */
148140
preview(input: PreviewSpec): CronResult<Preview>;
149141

@@ -552,48 +544,6 @@ export function createCronClient(opts?: CronClientOptions): CronClient {
552544
}),
553545
)),
554546

555-
sync: cmd("schedules.sync", async (specs) => {
556-
const desired = new Map(specs.map((s) => [s.name, s]));
557-
558-
// Upsert all
559-
await Promise.all(
560-
specs.map((spec) =>
561-
client.PUT("/v1/schedules/{name}", {
562-
params: { path: { name: spec.name } },
563-
body: snakenize(spec) as components["schemas"]["ScheduleSpec"],
564-
})
565-
),
566-
);
567-
568-
// List all and delete those not in desired set
569-
const { data: all, error, response } = await client.GET(
570-
"/v1/schedules",
571-
{},
572-
);
573-
if (error) {
574-
return {
575-
data: undefined,
576-
error: toCronError(error, response),
577-
response,
578-
};
579-
}
580-
581-
const toDelete = (all ?? []).filter((s) => !desired.has(s.name));
582-
await Promise.all(
583-
toDelete.map((s) =>
584-
client.DELETE("/v1/schedules/{name}", {
585-
params: { path: { name: s.name } },
586-
})
587-
),
588-
);
589-
590-
return {
591-
data: { upserted: specs.length, removed: toDelete.length },
592-
error: undefined,
593-
response,
594-
};
595-
}),
596-
597547
preview: cmd(
598548
"schedules.preview",
599549
(input) =>

0 commit comments

Comments
 (0)