server: cluster fork tracking

This commit is contained in:
Koushik Dutta
2024-11-15 22:00:51 -08:00
parent df249c554c
commit 882709ea51
3 changed files with 11 additions and 3 deletions

View File

@@ -1,12 +1,12 @@
{
"name": "@scrypted/server",
"version": "0.123.22",
"version": "0.123.23",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"name": "@scrypted/server",
"version": "0.123.22",
"version": "0.123.23",
"hasInstallScript": true,
"license": "ISC",
"dependencies": {

View File

@@ -68,6 +68,7 @@ export interface ClusterWorkerProperties {
export interface ClusterWorker extends ClusterWorkerProperties {
peer: RpcPeer;
forks: Set<ClusterForkOptions>;
}
export class PeerLiveness {
@@ -287,6 +288,7 @@ export function createClusterServer(runtime: ScryptedRuntime, certificate: Retur
const worker: ClusterWorker = {
...properties,
peer,
forks: new Set(),
};
runtime.clusterWorkers.add(worker);
peer.killed.then(() => {

View File

@@ -18,7 +18,12 @@ export class ClusterFork {
throw new Error(`no worker found for cluster labels ${JSON.stringify(options.labels)}`);
const fork: ClusterForkParam = await worker.peer.getParam('fork');
return fork(peerLiveness, options.runtime, packageJson, zipHash, getZip);
const forkResult = await fork(peerLiveness, options.runtime, packageJson, zipHash, getZip);
worker.forks.add(options);
forkResult.waitKilled().catch(() => {}).finally(() => {
worker.forks.delete(options);
});
return forkResult;
}
async getClusterWorkers() {
@@ -26,6 +31,7 @@ export class ClusterFork {
for (const worker of this.runtime.clusterWorkers) {
ret[worker.peer.peerName] = {
labels: worker.labels,
forks: [...worker.forks],
};
}
return ret;