1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
|
// Inspired from https://github.com/restatedev/examples/blob/main/typescript/patterns-use-cases/src/priorityqueue/queue.ts
import * as restate from "@restatedev/restate-sdk";
import { Context, object, ObjectContext } from "@restatedev/restate-sdk";
interface QueueItem {
awakeable: string;
priority: number;
}
interface QueueState {
items: QueueItem[];
inFlight: number;
}
export const semaphore = object({
name: "Semaphore",
handlers: {
acquire: async (
ctx: ObjectContext<QueueState>,
req: { awakeableId: string; priority: number; capacity: number },
): Promise<void> => {
const state = await getState(ctx);
state.items.push({
awakeable: req.awakeableId,
priority: req.priority,
});
tick(ctx, state, req.capacity);
setState(ctx, state);
},
release: async (
ctx: ObjectContext<QueueState>,
capacity: number,
): Promise<void> => {
const state = await getState(ctx);
state.inFlight--;
tick(ctx, state, capacity);
setState(ctx, state);
},
},
options: {
ingressPrivate: true,
journalRetention: 0,
},
});
// Lower numbers represent higher priority, mirroring Liteque’s semantics.
function selectAndPopItem(items: QueueItem[]): QueueItem {
let selected = { priority: Number.MAX_SAFE_INTEGER, index: 0 };
for (const [i, item] of items.entries()) {
if (item.priority < selected.priority) {
selected.priority = item.priority;
selected.index = i;
}
}
const [item] = items.splice(selected.index, 1);
return item;
}
function tick(
ctx: ObjectContext<QueueState>,
state: QueueState,
capacity: number,
) {
while (state.inFlight < capacity && state.items.length > 0) {
const item = selectAndPopItem(state.items);
state.inFlight++;
ctx.resolveAwakeable(item.awakeable);
}
}
async function getState(ctx: ObjectContext<QueueState>): Promise<QueueState> {
return {
items: (await ctx.get("items")) ?? [],
inFlight: (await ctx.get("inFlight")) ?? 0,
};
}
function setState(ctx: ObjectContext<QueueState>, state: QueueState) {
ctx.set("items", state.items);
ctx.set("inFlight", state.inFlight);
}
export class RestateSemaphore {
constructor(
private readonly ctx: Context,
private readonly id: string,
private readonly capacity: number,
) {}
async acquire(priority: number) {
const awk = this.ctx.awakeable();
await this.ctx
.objectClient<typeof semaphore>({ name: "Semaphore" }, this.id)
.acquire({
awakeableId: awk.id,
priority,
capacity: this.capacity,
});
try {
await awk.promise;
} catch (e) {
if (e instanceof restate.CancelledError) {
await this.release();
throw e;
}
}
}
async release() {
await this.ctx
.objectClient<typeof semaphore>({ name: "Semaphore" }, this.id)
.release(this.capacity);
}
}
|