| /** | |
| * 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.