Skip to content
36 changes: 34 additions & 2 deletions apps/server/src/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,7 @@ import * as GitManager from "./git/GitManager.ts";
import * as EnvironmentTheme from "./environmentTheme.ts";
import * as Keybindings from "./keybindings.ts";
import * as ServerRuntimeStartup from "./serverRuntimeStartup.ts";
import * as ServerSingleton from "./serverSingleton.ts";
import { OrchestrationReactorLive } from "./orchestration/Layers/OrchestrationReactor.ts";
import { RuntimeReceiptBusLive } from "./orchestration/Layers/RuntimeReceiptBus.ts";
import { ProviderRuntimeIngestionLive } from "./orchestration/Layers/ProviderRuntimeIngestion.ts";
Expand Down Expand Up @@ -221,6 +222,22 @@ const RelayClientLive = Layer.unwrap(
}),
);

/**
* Claims the data directory before anything binds a port or opens the database.
*
* Provided into both `HttpServerLive` and `RuntimeDependenciesLive` (same layer
* identity, so one acquire) so the lock precedes the HTTP bind *and* SQLite
* open. Providing it only into the HTTP layer left persistence free to
* initialize in parallel and write `state.sqlite` before a second process
* refused.
*/
const ServerSingletonLive = Layer.effectDiscard(
Effect.gen(function* () {
const config = yield* ServerConfig.ServerConfig;
yield* ServerSingleton.acquireServerSingleton(config.stateDir);
}),
);

const HttpServerLive = Layer.unwrap(
Effect.gen(function* () {
const config = yield* ServerConfig.ServerConfig;
Expand Down Expand Up @@ -550,6 +567,7 @@ const RuntimeDependenciesLive = RuntimeCoreDependenciesLive.pipe(
Layer.provideMerge(RemoteOpenTargets.layer),
Layer.provideMerge(ServerLifecycleEvents.layer),
Layer.provide(NetService.layer),
Layer.provide(ServerSingletonLive),
);

const commandReadinessLayer = HttpRouter.middleware(
Expand Down Expand Up @@ -605,9 +623,22 @@ const makeServerLayer = Layer.unwrap(

const httpListeningLayer = Layer.effectDiscard(
Effect.gen(function* () {
yield* HttpServer.HttpServer;
const server = yield* HttpServer.HttpServer;
const startup = yield* ServerRuntimeStartup.ServerRuntimeStartup;
yield* startup.markHttpListening;
const address = server.address;
if (typeof address === "string" || !("port" in address)) {
return;
}
// Stamp the port as soon as the socket is bound. Waiting until
// runtimeStateLayer's awaitActivation left a window where the server
// was listening but a second launch's refusal could not name the port.
yield* ServerSingleton.serverLockPath(config.stateDir).pipe(
Effect.flatMap((lockPath) =>
ServerSingleton.recordServerLockPort(lockPath, address.port),
),
Effect.ignore,
);
}),
);
const runtimeStateLayer = Layer.effectDiscard(
Expand All @@ -622,6 +653,7 @@ const makeServerLayer = Layer.unwrap(
}

const launcher = yield* ServiceLauncherClient.ServiceLauncherClient;

const state = yield* makePersistedServerRuntimeState({
config,
port: address.port,
Expand Down Expand Up @@ -794,7 +826,7 @@ const makeServerLayer = Layer.unwrap(
Layer.provideMerge(runtimeServicesLive),
Layer.provide(activationLayer),
Layer.provideMerge(serverRelayBrokerTracingLayer),
Layer.provideMerge(HttpServerLive),
Layer.provideMerge(HttpServerLive.pipe(Layer.provide(ServerSingletonLive))),
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Layer.provide(ApplicationObservabilityLive),
Layer.provideMerge(FetchHttpClient.layer),
// PR reads, Git operations, and WebSocket discovery share one process limiter.
Expand Down
Loading
Loading