diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index c50244d64973..521482f69fff 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -56,6 +56,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"; @@ -186,6 +187,21 @@ const RelayClientLive = Layer.unwrap( }), ); +/** + * Claims the data directory before anything binds a port or opens the database. + * + * Provided once into `HttpServerLive`, which is then provided into the runtime + * services. That dependency chain makes the ordering structural: claim the state + * directory, bind HTTP, then open persistence/provider state. A second server + * cannot reach either protected phase while another holds the directory. + */ +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; @@ -515,6 +531,16 @@ export const makeServerLayer = Layer.unwrap( return; } + // Stamp the port onto the lock we already hold. It is only ever read + // by a *later* server's refusal message, which turns "something else + // is running" into an address the user can open. + yield* ServerSingleton.serverLockPath(config.stateDir).pipe( + Effect.flatMap((lockPath) => + ServerSingleton.recordServerLockPort(lockPath, address.port), + ), + Effect.ignore, + ); + const state = yield* makePersistedServerRuntimeState({ config, port: address.port, @@ -672,6 +698,10 @@ export const makeServerLayer = Layer.unwrap( { concurrency: "unbounded" }, ).pipe(Effect.asVoid), }).pipe(Layer.provideMerge(RuntimeDependenciesLive), Layer.provide(launcherLayer)); + const ownedHttpServerLive = HttpServerLive.pipe(Layer.provide(ServerSingletonLive)); + const ownedRuntimeServicesLive = runtimeServicesLive.pipe( + Layer.provideMerge(ownedHttpServerLive), + ); const routesLayer = HttpRouter.serve(makeRoutesLayer.pipe(Layer.provide(launcherLayer)), { disableLogger: !config.logWebSocketEvents, @@ -685,10 +715,9 @@ export const makeServerLayer = Layer.unwrap( ); return serverApplicationLayer.pipe( - Layer.provideMerge(runtimeServicesLive), + Layer.provideMerge(ownedRuntimeServicesLive), Layer.provide(activationLayer), Layer.provideMerge(serverRelayBrokerTracingLayer), - Layer.provideMerge(HttpServerLive), Layer.provide(ApplicationObservabilityLive), Layer.provideMerge(FetchHttpClient.layer), Layer.provideMerge(VcsProcess.layer), diff --git a/apps/server/src/serverSingleton.test.ts b/apps/server/src/serverSingleton.test.ts new file mode 100644 index 000000000000..287fdea12409 --- /dev/null +++ b/apps/server/src/serverSingleton.test.ts @@ -0,0 +1,325 @@ +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { assert, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +import * as FileSystem from "effect/FileSystem"; +import * as Layer from "effect/Layer"; +import * as Path from "effect/Path"; +import * as Schema from "effect/Schema"; + +import * as ServerConfig from "./config.ts"; +import { makeServerLayer } from "./server.ts"; +import { PersistedServerRuntimeState } from "./serverRuntimeState.ts"; +import { + SERVER_LOCK_FILENAME, + SERVER_LOCK_DATABASE_FILENAME, + SERVER_RUNTIME_STATE_FILENAME, + acquireServerSingleton, + isServerLockBusyError, + processIsAlive, + processIsAliveWith, + recordServerLockPort, + serverLockDatabasePath, + serverLockPath, +} from "./serverSingleton.ts"; + +const layer = it.layer(NodeServices.layer); + +const makeStateDir = Effect.fn("test.makeStateDir")(function* () { + const fs = yield* FileSystem.FileSystem; + return yield* fs.makeTempDirectory({ prefix: "t3-singleton-" }); +}); + +/** A stale lock file, written without the module's own encoder on purpose. */ +const staleHolder = (pid: number) => + `{"version":1,"ownerId":"stale-owner","pid":${pid},"startedAt":"2026-01-01T00:00:00.000Z"}`; + +const PersistedServerRuntimeStateFromJson = Schema.fromJsonString(PersistedServerRuntimeState); +const encodeRuntimeState = Schema.encodeSync(PersistedServerRuntimeStateFromJson); + +/** What a pre-lock server persists: its live pid and port, and no lock file. */ +const legacyRuntimeState = (pid: number, port: number) => + `${encodeRuntimeState({ + version: 1, + pid, + port, + origin: `http://127.0.0.1:${port}`, + startedAt: "2026-08-27T12:00:00.000Z", + })}\n`; + +layer("serverSingleton", (it) => { + it.effect("the full server layer refuses before opening persistence", () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const baseDir = yield* fs.makeTempDirectory({ prefix: "t3-singleton-server-layer-" }); + const stateDir = path.join(baseDir, "userdata"); + yield* fs.makeDirectory(stateDir, { recursive: true }); + yield* Effect.scoped( + Effect.gen(function* () { + yield* acquireServerSingleton(stateDir); + const failure = yield* Layer.build( + makeServerLayer.pipe(Layer.provide(ServerConfig.layerTest(process.cwd(), baseDir))), + ).pipe(Effect.scoped, Effect.flip); + + assert.strictEqual(failure._tag, "ServerAlreadyRunningError"); + assert.isFalse(yield* fs.exists(path.join(stateDir, "state.sqlite"))); + }), + ); + }), + ); + + it.effect("claims a free directory and releases it on scope exit", () => + Effect.gen(function* () { + const stateDir = yield* makeStateDir(); + const fs = yield* FileSystem.FileSystem; + const lockPath = yield* serverLockPath(stateDir); + + yield* Effect.scoped( + Effect.gen(function* () { + yield* acquireServerSingleton(stateDir); + assert.isTrue(yield* fs.exists(lockPath)); + }), + ); + // Released on scope exit, so a restart is not blocked by its predecessor. + assert.isFalse(yield* fs.exists(lockPath)); + }), + ); + + it.effect("refuses a second server while the first holds the directory", () => + Effect.gen(function* () { + const stateDir = yield* makeStateDir(); + const failure = yield* Effect.scoped( + Effect.gen(function* () { + yield* acquireServerSingleton(stateDir); + // The incident: a second server started against a held directory, found + // its port taken, silently bound another, and corrupted shared state. + return yield* acquireServerSingleton(stateDir).pipe(Effect.flip); + }), + ); + assert.strictEqual(failure._tag, "ServerAlreadyRunningError"); + if (failure._tag === "ServerAlreadyRunningError") { + assert.strictEqual(failure.holderPid, process.pid); + assert.include(failure.message, stateDir); + assert.include(failure.message, "overwrite each other"); + } + }), + ); + + it.effect("overwrites stale display metadata after its owner is gone", () => + Effect.gen(function* () { + const stateDir = yield* makeStateDir(); + const fs = yield* FileSystem.FileSystem; + const lockPath = yield* serverLockPath(stateDir); + // pid 2^22 is above every /proc/sys/kernel/pid_max default, so it cannot + // be live. A crashed server must not lock its own directory forever. + yield* fs.writeFileString(lockPath, staleHolder(4194304)); + + yield* Effect.scoped( + Effect.gen(function* () { + const held = yield* acquireServerSingleton(stateDir); + assert.strictEqual(held, lockPath); + const metadata = yield* fs.readFileString(lockPath); + assert.notInclude(metadata, "stale-owner"); + }), + ); + }), + ); + + it.effect("overwrites half-written display metadata after a crash", () => + Effect.gen(function* () { + const stateDir = yield* makeStateDir(); + const fs = yield* FileSystem.FileSystem; + const lockPath = yield* serverLockPath(stateDir); + yield* fs.writeFileString(lockPath, '{"version":1,"pid":'); + + yield* Effect.scoped( + Effect.gen(function* () { + yield* acquireServerSingleton(stateDir); + }), + ); + }), + ); + + it.effect("does not remove display metadata replaced by another owner", () => + Effect.gen(function* () { + const stateDir = yield* makeStateDir(); + const fs = yield* FileSystem.FileSystem; + const lockPath = yield* serverLockPath(stateDir); + yield* Effect.scoped( + Effect.gen(function* () { + yield* acquireServerSingleton(stateDir); + yield* fs.writeFileString(lockPath, staleHolder(process.pid)); + }), + ); + + assert.isTrue(yield* fs.exists(lockPath)); + assert.include(yield* fs.readFileString(lockPath), "stale-owner"); + }), + ); + + it.effect("records the bound port so the next server can name it", () => + Effect.gen(function* () { + const stateDir = yield* makeStateDir(); + yield* Effect.scoped( + Effect.gen(function* () { + const lockPath = yield* acquireServerSingleton(stateDir); + yield* recordServerLockPort(lockPath, 3775); + const failure = yield* acquireServerSingleton(stateDir).pipe(Effect.flip); + assert.strictEqual(failure._tag, "ServerAlreadyRunningError"); + if (failure._tag === "ServerAlreadyRunningError") { + assert.strictEqual(failure.holderPort, 3775); + assert.include(failure.message, "listening on port 3775"); + } + }), + ); + }), + ); + + it.effect("keeps separate directories independent", () => + Effect.gen(function* () { + const first = yield* makeStateDir(); + const second = yield* makeStateDir(); + yield* Effect.scoped( + Effect.gen(function* () { + yield* acquireServerSingleton(first); + // A dev server and the real one use different state dirs and must both run. + yield* acquireServerSingleton(second); + }), + ); + }), + ); + + it.effect("keeps ownership and display metadata inside the state directory", () => + Effect.gen(function* () { + const stateDir = yield* makeStateDir(); + const path = yield* Path.Path; + const lockPath = yield* serverLockPath(stateDir); + const databasePath = yield* serverLockDatabasePath(stateDir); + assert.strictEqual(lockPath, path.join(stateDir, SERVER_LOCK_FILENAME)); + assert.strictEqual(databasePath, path.join(stateDir, SERVER_LOCK_DATABASE_FILENAME)); + }), + ); + + it("treats the current process as alive and an impossible pid as dead", () => { + assert.isTrue(processIsAlive(process.pid)); + assert.isFalse(processIsAlive(4194304)); + assert.isFalse(processIsAlive(0)); + assert.isFalse(processIsAlive(-1)); + assert.isTrue( + processIsAliveWith(123, () => { + throw Object.assign(new Error("not permitted"), { code: "EPERM" }); + }), + ); + assert.isFalse( + processIsAliveWith(123, () => { + throw Object.assign(new Error("not found"), { code: "ESRCH" }); + }), + ); + }); + + it("recognizes Node and Bun SQLite busy errors", () => { + assert.isTrue( + isServerLockBusyError( + Object.assign(new Error("database is locked"), { code: "SQLITE_BUSY" }), + ), + ); + assert.isTrue(isServerLockBusyError(Object.assign(new Error("busy"), { errno: 5 }))); + assert.isFalse(isServerLockBusyError(new Error("disk I/O error"))); + }); + + it.effect("refuses next to a live pre-lock server that wrote no lock", () => + Effect.gen(function* () { + const stateDir = yield* makeStateDir(); + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + // 0.0.34 and earlier persist their live pid here but never claim a lock: + // the desktop auto-update transition. The upgrade must still refuse. + yield* fs.writeFileString( + path.join(stateDir, SERVER_RUNTIME_STATE_FILENAME), + legacyRuntimeState(process.pid, 3775), + ); + + const failure = yield* Effect.scoped(acquireServerSingleton(stateDir)).pipe(Effect.flip); + assert.strictEqual(failure._tag, "ServerAlreadyRunningError"); + if (failure._tag === "ServerAlreadyRunningError") { + assert.strictEqual(failure.holderPid, process.pid); + assert.strictEqual(failure.holderPort, 3775); + assert.include(failure.message, SERVER_RUNTIME_STATE_FILENAME); + assert.include(failure.message, "compatibility guard for this older"); + assert.include(failure.message, "Do not remove it while the recorded process is live"); + assert.notInclude(failure.message, "contains display metadata only"); + } + // And it must not have claimed the directory it refused. + assert.isFalse(yield* fs.exists(yield* serverLockPath(stateDir))); + }), + ); + + it.effect("claims the directory when the pre-lock server's pid is gone", () => + Effect.gen(function* () { + const stateDir = yield* makeStateDir(); + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + // Leftover state from a crashed legacy server must not wedge the upgrade. + yield* fs.writeFileString( + path.join(stateDir, SERVER_RUNTIME_STATE_FILENAME), + legacyRuntimeState(4194304, 3775), + ); + + yield* Effect.scoped( + Effect.gen(function* () { + const held = yield* acquireServerSingleton(stateDir); + assert.strictEqual(held, yield* serverLockPath(stateDir)); + }), + ); + }), + ); + + it.effect("claims the directory when legacy runtime state is corrupt", () => + Effect.gen(function* () { + const stateDir = yield* makeStateDir(); + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + yield* fs.writeFileString( + path.join(stateDir, SERVER_RUNTIME_STATE_FILENAME), + '{"version":1,"pid":', + ); + + yield* Effect.scoped( + Effect.gen(function* () { + const held = yield* acquireServerSingleton(stateDir); + assert.strictEqual(held, yield* serverLockPath(stateDir)); + }), + ); + }), + ); + + it.effect("updates port metadata atomically", () => + Effect.gen(function* () { + const stateDir = yield* makeStateDir(); + const fs = yield* FileSystem.FileSystem; + const lockPath = yield* serverLockPath(stateDir); + yield* Effect.scoped( + Effect.gen(function* () { + yield* acquireServerSingleton(stateDir); + yield* recordServerLockPort(lockPath, 3775); + + const failure = yield* acquireServerSingleton(stateDir).pipe(Effect.flip); + assert.strictEqual(failure._tag, "ServerAlreadyRunningError"); + if (failure._tag === "ServerAlreadyRunningError") { + assert.strictEqual(failure.holderPort, 3775); + } + }), + ); + + // No temp staging directory outlives the metadata update. + const leftovers = yield* fs + .readDirectory(stateDir) + .pipe( + Effect.map((entries) => + entries.filter((entry) => entry.startsWith(`.${SERVER_LOCK_FILENAME}.`)), + ), + ); + assert.deepStrictEqual(leftovers, []); + }), + ); +}); diff --git a/apps/server/src/serverSingleton.ts b/apps/server/src/serverSingleton.ts new file mode 100644 index 000000000000..7057037e0103 --- /dev/null +++ b/apps/server/src/serverSingleton.ts @@ -0,0 +1,339 @@ +/** + * One server per data directory. + * + * The server already relies on SQLite for cross-platform filesystem locking. + * A dedicated lock database holds one write transaction for the process + * lifetime: a crash releases the OS lock automatically, while another process + * receives SQLITE_BUSY before it can open T3's real persistence or bind HTTP. + * + * `server.lock` is display-only metadata for the refusal message. It never + * decides ownership; deleting it cannot release the SQLite lock. + * + * A pre-lock server has no lock database, so the first upgraded process also + * checks `server-runtime.json` and refuses while that recorded pid is live. + * That compatibility check starts only after the old binary publishes its + * runtime state; an unmodified old binary cannot participate in a lock that + * did not exist when it shipped. + */ +import * as Data from "effect/Data"; +import * as Crypto from "effect/Crypto"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as FileSystem from "effect/FileSystem"; +import * as Option from "effect/Option"; +import * as Path from "effect/Path"; +import * as Schema from "effect/Schema"; + +import { writeFileStringAtomically } from "./atomicWrite.ts"; +import { readPersistedServerRuntimeState } from "./serverRuntimeState.ts"; + +export const SERVER_LOCK_FILENAME = "server.lock"; +export const SERVER_LOCK_DATABASE_FILENAME = "server-lock.sqlite"; +export const SERVER_RUNTIME_STATE_FILENAME = "server-runtime.json"; + +export const ServerLockHolder = Schema.Struct({ + version: Schema.Literal(1), + ownerId: Schema.String, + pid: Schema.Int, + startedAt: Schema.String, + /** Absent until the server binds; ownership is already held in SQLite. */ + port: Schema.optional(Schema.Int), +}); +export type ServerLockHolder = typeof ServerLockHolder.Type; + +const ServerLockHolderFromJson = Schema.fromJsonString(ServerLockHolder); +const decodeHolder = Schema.decodeUnknownOption(ServerLockHolderFromJson); +const encodeHolder = Schema.encodeSync(ServerLockHolderFromJson); + +export class ServerAlreadyRunningError extends Schema.TaggedErrorClass()( + "ServerAlreadyRunningError", + { + stateDir: Schema.String, + lockPath: Schema.String, + holderPid: Schema.optional(Schema.Int), + holderPort: Schema.optional(Schema.Int), + holderStartedAt: Schema.optional(Schema.String), + legacyRuntimeStatePath: Schema.optional(Schema.String), + }, +) { + override get message(): string { + const holder = + this.holderPid === undefined + ? "another live T3 Code server" + : this.holderPort === undefined + ? `pid ${this.holderPid}` + : `pid ${this.holderPid}, listening on port ${this.holderPort}`; + const common = [ + "Another T3 Code server is already using this data directory.", + "", + ` data directory: ${this.stateDir}`, + ` held by: ${holder}`, + ...(this.holderStartedAt === undefined ? [] : [` since: ${this.holderStartedAt}`]), + "", + "Two servers sharing one data directory overwrite each other's state.sqlite", + "and settings.json. Connect to the running server, stop it cleanly, or use", + "a different --base-dir.", + ]; + if (this.legacyRuntimeStatePath !== undefined) { + return [ + ...common, + "", + `${this.legacyRuntimeStatePath} is the compatibility guard for this older`, + "server. Do not remove it while the recorded process is live.", + ].join("\n"); + } + return [ + ...common, + "", + `Ownership is released automatically when the process exits. ${this.lockPath}`, + "contains display metadata only; removing it does not release a live owner.", + ].join("\n"); + } +} + +export class ServerLockUnavailableError extends Schema.TaggedErrorClass()( + "ServerLockUnavailableError", + { + lockPath: Schema.String, + operation: Schema.String, + cause: Schema.Defect(), + }, +) { + override get message(): string { + return `Could not ${this.operation} the server ownership database at ${this.lockPath}.`; + } +} + +export const processIsAliveWith = (pid: number, sendSignal: (pid: number) => void): boolean => { + if (!Number.isInteger(pid) || pid <= 0) return false; + try { + sendSignal(pid); + return true; + } catch (cause) { + return (cause as NodeJS.ErrnoException).code === "EPERM"; + } +}; + +export const processIsAlive = (pid: number): boolean => + processIsAliveWith(pid, (targetPid) => process.kill(targetPid, 0)); + +interface LockDatabase { + readonly exec: (sql: string) => void; + readonly close: () => void; +} + +interface HeldServerLock { + readonly database: LockDatabase; + readonly lockPath: string; + readonly ownerId: string; +} + +class ServerLockAttemptError extends Data.TaggedError("ServerLockAttemptError")<{ + readonly cause: unknown; +}> {} + +const readHolder = Effect.fn("serverSingleton.readHolder")(function* (lockPath: string) { + const fs = yield* FileSystem.FileSystem; + const raw = yield* fs + .readFileString(lockPath) + .pipe( + Effect.catch((error) => + error.reason._tag === "NotFound" ? Effect.succeed("") : Effect.fail(error), + ), + ); + return Option.getOrUndefined(decodeHolder(raw)); +}); + +const closeDatabase = (database: LockDatabase) => + Effect.sync(() => { + try { + database.exec("ROLLBACK"); + } catch { + // A failed BEGIN has no transaction to roll back. + } + database.close(); + }).pipe(Effect.ignore); + +const openLockDatabase = Effect.fn("serverSingleton.openDatabase")(function* (lockPath: string) { + if (process.versions.bun !== undefined) { + const { Database } = yield* Effect.promise(() => import("bun:sqlite")); + return yield* Effect.try({ + try: () => { + const database = new Database(lockPath, { create: true }); + return { + exec: (sql: string) => database.exec(sql), + close: () => database.close(), + } satisfies LockDatabase; + }, + catch: (cause) => new ServerLockUnavailableError({ lockPath, operation: "open", cause }), + }); + } + + const { DatabaseSync } = yield* Effect.promise(() => import("node:sqlite")); + return yield* Effect.try({ + try: () => { + const database = new DatabaseSync(lockPath); + return { + exec: (sql: string) => database.exec(sql), + close: () => database.close(), + } satisfies LockDatabase; + }, + catch: (cause) => new ServerLockUnavailableError({ lockPath, operation: "open", cause }), + }); +}); + +export const isServerLockBusyError = (cause: unknown): boolean => { + const code = + cause instanceof Error && "code" in cause ? String((cause as NodeJS.ErrnoException).code) : ""; + const errno = + cause instanceof Error && "errno" in cause + ? Number((cause as NodeJS.ErrnoException).errno) + : undefined; + const message = cause instanceof Error ? cause.message : String(cause); + return ( + code === "SQLITE_BUSY" || + code === "SQLITE_BUSY_SNAPSHOT" || + errno === 5 || + /database is (?:locked|busy)/i.test(message) + ); +}; + +/** A pre-lock server is live against this directory. */ +export class LiveLegacyServerRuntime extends Data.TaggedError("LiveLegacyServerRuntime")<{ + readonly state: { + readonly pid: number; + readonly port: number; + readonly startedAt: string; + }; +}> {} + +const beginExclusiveWrite = Effect.fn("serverSingleton.beginExclusiveWrite")(function* ( + database: LockDatabase, + lockPath: string, +) { + const result = yield* Effect.try({ + try: () => { + database.exec("PRAGMA busy_timeout = 0"); + database.exec("BEGIN IMMEDIATE"); + }, + catch: (cause) => new ServerLockAttemptError({ cause }), + }).pipe( + Effect.as({ ok: true as const }), + Effect.catch((error) => Effect.succeed({ ok: false as const, cause: error.cause })), + ); + if (result.ok) return true; + if (isServerLockBusyError(result.cause)) return false; + return yield* new ServerLockUnavailableError({ + lockPath, + operation: "lock", + cause: result.cause, + }); +}); + +export const serverLockPath = Effect.fn("serverSingleton.lockPath")(function* (stateDir: string) { + const path = yield* Path.Path; + return path.join(stateDir, SERVER_LOCK_FILENAME); +}); + +export const serverLockDatabasePath = Effect.fn("serverSingleton.databasePath")(function* ( + stateDir: string, +) { + const path = yield* Path.Path; + return path.join(stateDir, SERVER_LOCK_DATABASE_FILENAME); +}); + +export const legacyServerRuntimeStatePath = Effect.fn("serverSingleton.legacyStatePath")(function* ( + stateDir: string, +) { + const path = yield* Path.Path; + return path.join(stateDir, SERVER_RUNTIME_STATE_FILENAME); +}); + +const acquireLock = Effect.fn("serverSingleton.acquireLock")(function* (stateDir: string) { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const crypto = yield* Crypto.Crypto; + const lockPath = yield* serverLockPath(stateDir); + const databasePath = yield* serverLockDatabasePath(stateDir); + const legacyRuntimeStatePath = yield* legacyServerRuntimeStatePath(stateDir); + + const legacyState = yield* readPersistedServerRuntimeState(legacyRuntimeStatePath); + if (Option.isSome(legacyState) && processIsAlive(legacyState.value.pid)) { + return yield* new LiveLegacyServerRuntime({ state: legacyState.value }); + } + + yield* fs.makeDirectory(path.dirname(databasePath), { recursive: true }); + const database = yield* openLockDatabase(databasePath); + return yield* Effect.gen(function* () { + if (!(yield* beginExclusiveWrite(database, databasePath))) { + const holder = yield* readHolder(lockPath); + return yield* new ServerAlreadyRunningError({ + stateDir, + lockPath, + ...(holder === undefined + ? {} + : { + holderPid: holder.pid, + holderStartedAt: holder.startedAt, + ...(holder.port === undefined ? {} : { holderPort: holder.port }), + }), + }); + } + + const ownerId = yield* crypto.randomUUIDv4; + const startedAt = DateTime.formatIso(yield* DateTime.now); + yield* writeFileStringAtomically({ + filePath: lockPath, + contents: encodeHolder({ version: 1, ownerId, pid: process.pid, startedAt }), + }); + return { database, lockPath, ownerId } satisfies HeldServerLock; + }).pipe(Effect.onError(() => closeDatabase(database))); +}); + +const releaseLock = Effect.fn("serverSingleton.releaseLock")(function* (held: HeldServerLock) { + const fs = yield* FileSystem.FileSystem; + const holder = yield* readHolder(held.lockPath).pipe(Effect.orElseSucceed(() => undefined)); + if (holder?.ownerId === held.ownerId) { + // Remove metadata while the SQLite transaction still excludes successors. + yield* fs.remove(held.lockPath, { force: true }).pipe(Effect.ignore); + } + yield* closeDatabase(held.database); +}); + +/** Records the bound port in display metadata while SQLite owns the lock. */ +export const recordServerLockPort = Effect.fn("serverSingleton.recordPort")(function* ( + lockPath: string, + port: number, +) { + const holder = yield* readHolder(lockPath); + if (holder === undefined || holder.pid !== process.pid) return; + yield* writeFileStringAtomically({ + filePath: lockPath, + contents: encodeHolder({ ...holder, port }), + }).pipe(Effect.ignore); +}); + +/** Holds the state directory until the caller's scope closes. */ +export const acquireServerSingleton = Effect.fn("serverSingleton.acquire")(function* ( + stateDir: string, +) { + const legacyRuntimeStatePath = yield* legacyServerRuntimeStatePath(stateDir); + return yield* Effect.acquireRelease( + acquireLock(stateDir).pipe( + Effect.catchTags({ + LiveLegacyServerRuntime: (legacy) => + Effect.fail( + new ServerAlreadyRunningError({ + stateDir, + lockPath: legacyRuntimeStatePath, + holderPid: legacy.state.pid, + holderPort: legacy.state.port, + holderStartedAt: legacy.state.startedAt, + legacyRuntimeStatePath, + }), + ), + }), + ), + releaseLock, + ).pipe(Effect.map((held) => held.lockPath)); +});