Terminate workers after timeout, handle errors
This commit is contained in:
parent
31bddb98c6
commit
49ec724049
@ -38,11 +38,22 @@ interface WorkerState {
|
|||||||
/**
|
/**
|
||||||
* The actual worker thread.
|
* The actual worker thread.
|
||||||
*/
|
*/
|
||||||
w: Worker;
|
w: Worker|null;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Are we currently running a task on this worker?
|
* Are we currently running a task on this worker?
|
||||||
*/
|
*/
|
||||||
busy: boolean;
|
busy: boolean;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Timer to terminate the worker if it's not busy enough.
|
||||||
|
*/
|
||||||
|
terminationTimerHandle: number|null;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Index of this worker in the list of workers.
|
||||||
|
*/
|
||||||
|
workerIndex: number;
|
||||||
}
|
}
|
||||||
|
|
||||||
interface WorkItem {
|
interface WorkItem {
|
||||||
@ -74,53 +85,99 @@ export class CryptoApi {
|
|||||||
private numWaiting: number = 0;
|
private numWaiting: number = 0;
|
||||||
|
|
||||||
|
|
||||||
constructor() {
|
/**
|
||||||
let handler = (msg: MessageEvent) => {
|
* Start a worker (if not started) and set as busy.
|
||||||
let id = msg.data.id;
|
*/
|
||||||
if (typeof id !== "number") {
|
wake(ws: WorkerState) {
|
||||||
console.error("rpc id must be number");
|
if (ws.busy) {
|
||||||
return;
|
throw Error("assertion failed");
|
||||||
}
|
}
|
||||||
if (!this.rpcRegistry[id]) {
|
ws.busy = true;
|
||||||
console.error(`RPC with id ${id} has no registry entry`);
|
this.numBusy++;
|
||||||
return;
|
if (!ws.w) {
|
||||||
}
|
let w = new Worker("/lib/wallet/cryptoWorker.js");
|
||||||
let {resolve, workerIndex} = this.rpcRegistry[id];
|
w.onmessage = (m: MessageEvent) => this.handleWorkerMessage(ws, m);
|
||||||
delete this.rpcRegistry[id];
|
w.onerror = (e: ErrorEvent) => this.handleWorkerError(ws, e);
|
||||||
let ws = this.workers[workerIndex];
|
ws.w = w;
|
||||||
if (!ws.busy) {
|
}
|
||||||
throw Error("assertion failed");
|
if (ws.terminationTimerHandle != null) {
|
||||||
|
clearTimeout(ws.terminationTimerHandle);
|
||||||
|
}
|
||||||
|
let destroy = () => {
|
||||||
|
if (ws.w && !ws.busy) {
|
||||||
|
ws.w!.terminate();
|
||||||
|
ws.w = null;
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
ws.terminationTimerHandle = setTimeout(destroy, 20 * 1000);
|
||||||
|
}
|
||||||
|
|
||||||
|
handleWorkerError(ws: WorkerState, e: ErrorEvent) {
|
||||||
|
console.error("error in worker", e);
|
||||||
|
try {
|
||||||
|
ws.w!.terminate();
|
||||||
|
} catch (e) {
|
||||||
|
console.error(e);
|
||||||
|
}
|
||||||
|
if (ws.busy) {
|
||||||
ws.busy = false;
|
ws.busy = false;
|
||||||
this.numBusy--;
|
this.numBusy--;
|
||||||
resolve(msg.data.result);
|
}
|
||||||
|
this.findWork(ws);
|
||||||
|
}
|
||||||
|
|
||||||
// try to find more work for this worker
|
findWork(ws: WorkerState) {
|
||||||
for (let i = 0; i < NUM_PRIO; i++) {
|
// try to find more work for this worker
|
||||||
let q = this.workQueues[NUM_PRIO - i - 1];
|
for (let i = 0; i < NUM_PRIO; i++) {
|
||||||
if (q.length != 0) {
|
let q = this.workQueues[NUM_PRIO - i - 1];
|
||||||
let work: WorkItem = q.shift()!;
|
if (q.length != 0) {
|
||||||
let msg: any = {
|
let work: WorkItem = q.shift()!;
|
||||||
operation: work.operation,
|
let msg: any = {
|
||||||
args: work.args,
|
operation: work.operation,
|
||||||
id: this.registerRpcId(work.resolve, work.reject, workerIndex),
|
args: work.args,
|
||||||
};
|
id: this.registerRpcId(work.resolve, work.reject, ws.workerIndex),
|
||||||
ws.w.postMessage(msg);
|
};
|
||||||
ws.busy = true;
|
|
||||||
this.numBusy++;
|
this.wake(ws);
|
||||||
return;
|
ws.w!.postMessage(msg);
|
||||||
}
|
return;
|
||||||
}
|
}
|
||||||
};
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
handleWorkerMessage(ws: WorkerState, msg: MessageEvent) {
|
||||||
|
let id = msg.data.id;
|
||||||
|
if (typeof id !== "number") {
|
||||||
|
console.error("rpc id must be number");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (!this.rpcRegistry[id]) {
|
||||||
|
console.error(`RPC with id ${id} has no registry entry`);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
let {resolve, workerIndex} = this.rpcRegistry[id];
|
||||||
|
delete this.rpcRegistry[id];
|
||||||
|
if (workerIndex != ws.workerIndex) {
|
||||||
|
throw Error("assertion failed");
|
||||||
|
}
|
||||||
|
if (!ws.busy) {
|
||||||
|
throw Error("assertion failed");
|
||||||
|
}
|
||||||
|
ws.busy = false;
|
||||||
|
this.numBusy--;
|
||||||
|
resolve(msg.data.result);
|
||||||
|
this.findWork(ws);
|
||||||
|
}
|
||||||
|
|
||||||
|
constructor() {
|
||||||
this.workers = new Array<WorkerState>((navigator as any)["hardwareConcurrency"] || 2);
|
this.workers = new Array<WorkerState>((navigator as any)["hardwareConcurrency"] || 2);
|
||||||
|
|
||||||
for (let i = 0; i < this.workers.length; i++) {
|
for (let i = 0; i < this.workers.length; i++) {
|
||||||
let w = new Worker("/lib/wallet/cryptoWorker.js");
|
|
||||||
w.onmessage = handler;
|
|
||||||
this.workers[i] = {
|
this.workers[i] = {
|
||||||
w,
|
w: null,
|
||||||
busy: false,
|
busy: false,
|
||||||
|
terminationTimerHandle: null,
|
||||||
|
workerIndex: i,
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
this.workQueues = [];
|
this.workQueues = [];
|
||||||
@ -156,15 +213,14 @@ export class CryptoApi {
|
|||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|
||||||
ws.busy = true;
|
this.wake(ws);
|
||||||
this.numBusy++;
|
|
||||||
|
|
||||||
return new Promise<T>((resolve, reject) => {
|
return new Promise<T>((resolve, reject) => {
|
||||||
let msg: any = {
|
let msg: any = {
|
||||||
operation, args,
|
operation, args,
|
||||||
id: this.registerRpcId(resolve, reject, i),
|
id: this.registerRpcId(resolve, reject, i),
|
||||||
};
|
};
|
||||||
ws.w.postMessage(msg);
|
ws.w!.postMessage(msg);
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
Loading…
Reference in New Issue
Block a user