EdgeAIG's picture
download
raw
4.39 kB
/**
* Stores runner registration and shard-lock state for cluster sharding.
*
* `RunnerStorage` records which runners are registered, whether they are
* healthy, which machine id a runner receives, and which shard locks are held
* by each runner. This module includes the typed storage service, a
* string-encoded backend interface, an adapter from encoded storage to the typed
* service, and an in-memory implementation for tests and local use.
*
* @since 4.0.0
*/
import { isArrayNonEmpty } from "../../Array.js";
import * as Context from "../../Context.js";
import * as Effect from "../../Effect.js";
import * as Layer from "../../Layer.js";
import * as MutableHashMap from "../../MutableHashMap.js";
import * as MachineId from "./MachineId.js";
import { Runner } from "./Runner.js";
import * as ShardId from "./ShardId.js";
/**
* Represents a generic interface to the persistent storage required by the
* cluster.
*
* @category models
* @since 4.0.0
*/
export class RunnerStorage extends /*#__PURE__*/Context.Service()("effect/cluster/RunnerStorage") {}
/**
* Adapts an encoded runner storage implementation into `RunnerStorage`, converting
* runner addresses, runners, machine ids, and shard ids between typed values and
* their string or numeric storage forms.
*
* @category layers
* @since 4.0.0
*/
export const makeEncoded = encoded => RunnerStorage.of({
getRunners: Effect.gen(function* () {
const runners = yield* encoded.getRunners;
const results = [];
for (let i = 0; i < runners.length; i++) {
const [runner, healthy] = runners[i];
// @effect-diagnostics-next-line tryCatchInEffectGen:off
try {
results.push([Runner.decodeSync(runner), healthy]);
} catch {
//
}
}
return results;
}),
register: (runner, healthy) => Effect.map(encoded.register(encodeRunnerAddress(runner.address), Runner.encodeSync(runner), healthy), MachineId.make),
unregister: address => encoded.unregister(encodeRunnerAddress(address)),
setRunnerHealth: (address, healthy) => encoded.setRunnerHealth(encodeRunnerAddress(address), healthy),
acquire: (address, shardIds) => {
const arr = Array.from(shardIds, id => id.toString());
if (!isArrayNonEmpty(arr)) return Effect.succeed([]);
return encoded.acquire(encodeRunnerAddress(address), arr).pipe(Effect.map(shards => shards.map(ShardId.fromString)));
},
refresh: (address, shardIds) => encoded.refresh(encodeRunnerAddress(address), Array.from(shardIds, id => id.toString())).pipe(Effect.map(shards => shards.map(ShardId.fromString))),
release(address, shardId) {
return encoded.release(encodeRunnerAddress(address), shardId.toString());
},
releaseAll(address) {
return encoded.releaseAll(encodeRunnerAddress(address));
}
});
/**
* Creates an in-memory `RunnerStorage` implementation for tests and local use.
*
* **Details**
*
* Registered runners are treated as healthy and shard acquisition is kept only in
* process memory.
*
* @category constructors
* @since 4.0.0
*/
export const makeMemory = /*#__PURE__*/Effect.gen(function* () {
const runners = MutableHashMap.empty();
let acquired = [];
let id = 0;
return RunnerStorage.of({
getRunners: Effect.sync(() => Array.from(MutableHashMap.values(runners), runner => [runner, true])),
register: runner => Effect.sync(() => {
MutableHashMap.set(runners, runner.address, runner);
return MachineId.make(id++);
}),
unregister: address => Effect.sync(() => {
MutableHashMap.remove(runners, address);
}),
setRunnerHealth: () => Effect.void,
acquire: (_address, shardIds) => {
acquired = Array.from(shardIds);
return Effect.succeed(Array.from(shardIds));
},
refresh: () => Effect.sync(() => acquired),
release: () => Effect.void,
releaseAll: () => Effect.void
});
});
/**
* Layer that provides the in-memory `RunnerStorage` implementation.
*
* @category layers
* @since 4.0.0
*/
export const layerMemory = /*#__PURE__*/Layer.effect(RunnerStorage)(makeMemory);
// -------------------------------------------------------------------------------------
// internal
// -------------------------------------------------------------------------------------
const encodeRunnerAddress = runnerAddress => `${runnerAddress.host}:${runnerAddress.port}`;
//# sourceMappingURL=RunnerStorage.js.map

Xet Storage Details

Size:
4.39 kB
·
Xet hash:
bd25bc40e2c1431fd0a918c9af3952a217d36abf00fa164580c43b44966bba23

Xet efficiently stores files, intelligently splitting them into unique chunks and accelerating uploads and downloads. More info.