/**
* @since 1.0.0
*/
import { SocketServer } from "@effect/platform/SocketServer"
import type * as RpcSerialization from "@effect/rpc/RpcSerialization"
import * as RpcServer from "@effect/rpc/RpcServer"
import * as Effect from "effect/Effect"
import * as Layer from "effect/Layer"
import type { MessageStorage } from "./MessageStorage.js"
import type { RunnerHealth } from "./RunnerHealth.js"
import type * as Runners from "./Runners.js"
import * as RunnerServer from "./RunnerServer.js"
import type * as RunnerStorage from "./RunnerStorage.js"
import type * as Sharding from "./Sharding.js"
import type { ShardingConfig } from "./ShardingConfig.js"
const withLogAddress = (layer: Layer.Layer): Layer.Layer =>
Layer.effectDiscard(Effect.gen(function*() {
const server = yield* SocketServer
const address = server.address._tag === "UnixAddress"
? server.address.path
: `${server.address.hostname}:${server.address.port}`
yield* Effect.annotateLogs(Effect.logInfo(`Listening on: ${address}`), {
package: "@effect/cluster",
service: "Runner"
})
})).pipe(Layer.provideMerge(layer))
/**
* @since 1.0.0
* @category Layers
*/
export const layer: Layer.Layer<
Sharding.Sharding | Runners.Runners,
never,
| Runners.RpcClientProtocol
| ShardingConfig
| RpcSerialization.RpcSerialization
| SocketServer
| MessageStorage
| RunnerStorage.RunnerStorage
| RunnerHealth
> = RunnerServer.layerWithClients.pipe(
withLogAddress,
Layer.provide(RpcServer.layerProtocolSocketServer)
)
/**
* @since 1.0.0
* @category Layers
*/
export const layerClientOnly: Layer.Layer<
Sharding.Sharding | Runners.Runners,
never,
Runners.RpcClientProtocol | ShardingConfig | MessageStorage | RunnerStorage.RunnerStorage
> = RunnerServer.layerClientOnly