-
Notifications
You must be signed in to change notification settings - Fork 6
Expand file tree
/
Copy path9-web-locks.js
More file actions
127 lines (111 loc) · 2.8 KB
/
Copy path9-web-locks.js
File metadata and controls
127 lines (111 loc) · 2.8 KB
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
120
121
122
123
124
125
126
127
'use strict';
const {
Worker,
isMainThread,
threadId,
parentPort,
} = require('node:worker_threads');
const threads = new Set();
const LOCKED = 0;
const UNLOCKED = 1;
const locks = {
resources: new Map(),
};
class Mutex {
constructor(resourceName, shared, initial = false) {
this.resourceName = resourceName;
this.lock = new Int32Array(shared, 0, 1);
if (initial) Atomics.store(this.lock, 0, UNLOCKED);
this.owner = false;
this.trying = false;
this.callback = null;
}
async enter(callback) {
this.callback = callback;
this.trying = true;
await this.tryEnter();
}
async tryEnter() {
if (!this.callback) return;
const prev = Atomics.exchange(this.lock, 0, LOCKED);
if (prev === UNLOCKED) {
this.owner = true;
this.trying = false;
await this.callback(this);
this.callback = null;
this.leave();
}
}
leave() {
if (!this.owner) return;
Atomics.store(this.lock, 0, UNLOCKED);
this.owner = false;
locks.sendMessage({ kind: 'leave', resourceName: this.resourceName });
}
}
locks.request = async (resourceName, callback) => {
let lock = locks.resources.get(resourceName);
if (!lock) {
const buffer = new SharedArrayBuffer(4);
lock = new Mutex(resourceName, buffer, true);
locks.resources.set(resourceName, lock);
locks.sendMessage({ kind: 'create', resourceName, buffer });
}
await lock.enter(callback);
};
locks.sendMessage = (message) => {
if (isMainThread) {
for (const thread of threads) {
thread.worker.postMessage(message);
}
} else {
parentPort.postMessage(message);
}
};
locks.receiveMessage = (message) => {
const { kind, resourceName, buffer } = message;
if (kind === 'create') {
const lock = new Mutex(resourceName, buffer);
locks.resources.set(resourceName, lock);
} else if (kind === 'leave') {
for (const mutex of locks.resources) {
if (mutex.trying) mutex.tryEnter();
}
}
};
if (!isMainThread) {
parentPort.on('message', locks.receiveMessage);
}
class Thread {
constructor() {
const worker = new Worker(__filename);
this.worker = worker;
threads.add(this);
worker.on('message', (message) => {
for (const thread of threads) {
if (thread.worker !== worker) {
thread.worker.postMessage(message);
}
}
locks.receiveMessage(message);
});
}
}
// Usage
if (isMainThread) {
new Thread();
new Thread();
setTimeout(() => {
process.exit(0);
}, 300);
} else {
locks.request('A', async (lock) => {
console.log(`Enter ${lock.resourceName} in ${threadId}`);
});
setTimeout(async () => {
await locks.request('B', async (lock) => {
console.log(`Enter ${lock.resourceName} in ${threadId}`);
});
console.log(`Leave all in ${threadId}`);
}, 100);
}