fix(core): restore standalone central backend initialization (#2596)
## Summary - restore the owned PostgreSQL backend bootstrap for layer-less `CentralCore.init()` callers - fix node, mesh, and project CLI commands returning empty state and logging `backendHandle is only available in backend mode` during cleanup - add a hermetic regression test for standalone backend ownership and shutdown PR #2454 accidentally added an unconditional early return immediately before the existing standalone bootstrap. Runtime pool sharing remains unchanged: `attachBackendLayer()` releases the central-only connections before adopting the project store layer. ## Verification - RED: regression failed because `createCentralBackendLayer` had zero calls - GREEN: focused regression passes - `pnpm --filter @fusion/core typecheck` - `pnpm --filter @fusion/core build` - `pnpm test` with `FUSION_PG_TEST_SKIP=1`: 482 engine + 132 Core gate + 71 CLI shape + changed regression passed; isolation clean - changeset format check passed The local PostgreSQL merge-gate harness is unavailable without credentials (`empty password returned by client`), so its 10 tests were explicitly skipped rather than misreported. <!-- This is an auto-generated comment: release notes by coderabbit.ai --> ## Summary by CodeRabbit * **Bug Fixes** * Restored PostgreSQL central registry access for standalone Node/mesh/project CLI commands without marking the host offline on shutdown. * Improved CentralCore lifecycle handling: concurrent `init()` coalesces, and operations are blocked once `close()` is requested/in progress. * Refined embedded PostgreSQL runtime shutdown: owner stop is coordinated with lease release, registrations are rejected while stopping, and shutdown/teardown uses lease lifecycle consistently. Embedded start failures now treat stopping as retryable. * **Tests** * Expanded coverage for CentralCore close/init/attach races and embedded PostgreSQL lease/shutdown coordination scenarios. * **Documentation** * Updated Changeset notes to clarify CLI and shutdown semantics. <!-- end of auto-generated comment: release notes by coderabbit.ai -->
This commit is contained in:
7
.changeset/central-core-cli-backend-init.md
Normal file
7
.changeset/central-core-cli-backend-init.md
Normal file
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@runfusion/fusion": patch
|
||||
---
|
||||
|
||||
summary: Node, mesh, and project CLI commands restore registry access without marking the host offline on exit.
|
||||
category: fix
|
||||
dev: Restore the existing layer-less `CentralCore.init()` backend bootstrap that PostgreSQL dual-path cleanup accidentally left unreachable. Generic `close()` only releases resources; daemon, engine-manager, and dashboard shutdown owners retain their explicitly ordered `markLocalNodeOffline()` writes.
|
||||
@@ -6,6 +6,7 @@ import {
|
||||
} from "../backup.js";
|
||||
import {
|
||||
clearActiveEmbeddedRuntimeUrl,
|
||||
EmbeddedRuntimeStoppingError,
|
||||
getActiveEmbeddedRuntimeUrl,
|
||||
invalidateEmbeddedRuntimeUrl,
|
||||
registerEmbeddedRuntimeUrl,
|
||||
@@ -45,14 +46,76 @@ describe("embedded backup runtime URL registry", () => {
|
||||
);
|
||||
});
|
||||
|
||||
it("keeps an owner URL live when only a joiner releases", () => {
|
||||
it("keeps an owner URL live when only a joiner releases", async () => {
|
||||
const owner = registerEmbeddedRuntimeUrl(embeddedUrl, { ownsProcess: true });
|
||||
const joiner = registerEmbeddedRuntimeUrl(embeddedUrl, { ownsProcess: false });
|
||||
|
||||
releaseEmbeddedRuntimeLease(joiner);
|
||||
await releaseEmbeddedRuntimeLease(joiner);
|
||||
expect(getActiveEmbeddedRuntimeUrl()).toBe(embeddedUrl);
|
||||
|
||||
releaseEmbeddedRuntimeLease(owner);
|
||||
await releaseEmbeddedRuntimeLease(owner);
|
||||
expect(getActiveEmbeddedRuntimeUrl()).toBeUndefined();
|
||||
});
|
||||
|
||||
it("defers owner shutdown until the final joined lease releases", async () => {
|
||||
const stopOwner = vi.fn(async () => undefined);
|
||||
const owner = registerEmbeddedRuntimeUrl(embeddedUrl, { ownsProcess: true });
|
||||
const joiner = registerEmbeddedRuntimeUrl(embeddedUrl, { ownsProcess: false });
|
||||
|
||||
await releaseEmbeddedRuntimeLease(owner, { stopOwner });
|
||||
expect(stopOwner).not.toHaveBeenCalled();
|
||||
expect(getActiveEmbeddedRuntimeUrl()).toBe(embeddedUrl);
|
||||
|
||||
await releaseEmbeddedRuntimeLease(joiner);
|
||||
expect(stopOwner).toHaveBeenCalledOnce();
|
||||
expect(getActiveEmbeddedRuntimeUrl()).toBeUndefined();
|
||||
});
|
||||
|
||||
it("rejects registrations until a deferred owner stop completes", async () => {
|
||||
const owner = registerEmbeddedRuntimeUrl(embeddedUrl, { ownsProcess: true });
|
||||
const joiner = registerEmbeddedRuntimeUrl(embeddedUrl, { ownsProcess: false });
|
||||
let finishStop!: () => void;
|
||||
const stopFinished = new Promise<void>((resolve) => {
|
||||
finishStop = resolve;
|
||||
});
|
||||
const stopOwner = vi.fn(async () => stopFinished);
|
||||
await releaseEmbeddedRuntimeLease(owner, { stopOwner });
|
||||
|
||||
const finalRelease = releaseEmbeddedRuntimeLease(joiner);
|
||||
await vi.waitFor(() => expect(stopOwner).toHaveBeenCalledOnce());
|
||||
expect(getActiveEmbeddedRuntimeUrl()).toBeUndefined();
|
||||
let stoppingError: EmbeddedRuntimeStoppingError | undefined;
|
||||
try {
|
||||
registerEmbeddedRuntimeUrl(embeddedUrl, { ownsProcess: false });
|
||||
} catch (error) {
|
||||
if (error instanceof EmbeddedRuntimeStoppingError) stoppingError = error;
|
||||
}
|
||||
expect(stoppingError).toBeInstanceOf(EmbeddedRuntimeStoppingError);
|
||||
if (!stoppingError) throw new Error("Expected stopping registration to expose completion");
|
||||
const stopCompletion = stoppingError.completion;
|
||||
let stopCompletionSettled = false;
|
||||
void stopCompletion.then(() => {
|
||||
stopCompletionSettled = true;
|
||||
});
|
||||
await Promise.resolve();
|
||||
expect(stopCompletionSettled).toBe(false);
|
||||
|
||||
finishStop();
|
||||
await stopCompletion;
|
||||
await finalRelease;
|
||||
expect(stopCompletionSettled).toBe(true);
|
||||
expect(registerEmbeddedRuntimeUrl(embeddedUrl, { ownsProcess: true })).toBeDefined();
|
||||
});
|
||||
|
||||
it("stops the owner immediately when joined leases already released", async () => {
|
||||
const stopOwner = vi.fn(async () => undefined);
|
||||
const owner = registerEmbeddedRuntimeUrl(embeddedUrl, { ownsProcess: true });
|
||||
const joiner = registerEmbeddedRuntimeUrl(embeddedUrl, { ownsProcess: false });
|
||||
|
||||
await releaseEmbeddedRuntimeLease(joiner);
|
||||
await releaseEmbeddedRuntimeLease(owner, { stopOwner });
|
||||
|
||||
expect(stopOwner).toHaveBeenCalledOnce();
|
||||
expect(getActiveEmbeddedRuntimeUrl()).toBeUndefined();
|
||||
});
|
||||
|
||||
|
||||
171
packages/core/src/__tests__/central-core-layerless-init.test.ts
Normal file
171
packages/core/src/__tests__/central-core-layerless-init.test.ts
Normal file
@@ -0,0 +1,171 @@
|
||||
import { mkdtempSync, rmSync } from "node:fs";
|
||||
import { tmpdir } from "node:os";
|
||||
import { join } from "node:path";
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
|
||||
const mocks = vi.hoisted(() => ({
|
||||
createCentralBackendLayer: vi.fn(),
|
||||
ensureBackendBootstrap: vi.fn(),
|
||||
getLocalNode: vi.fn(),
|
||||
releaseConnections: vi.fn(),
|
||||
shutdown: vi.fn(),
|
||||
}));
|
||||
|
||||
vi.mock("../postgres/startup-factory.js", () => ({
|
||||
createCentralBackendLayer: mocks.createCentralBackendLayer,
|
||||
}));
|
||||
|
||||
vi.mock("../async-central-core.js", async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import("../async-central-core.js")>();
|
||||
return {
|
||||
...actual,
|
||||
ensureBackendBootstrap: mocks.ensureBackendBootstrap,
|
||||
getLocalNode: mocks.getLocalNode,
|
||||
};
|
||||
});
|
||||
|
||||
import { CentralCore } from "../central-core.js";
|
||||
|
||||
const cleanupDirs: string[] = [];
|
||||
const ownedLayer = { db: {} };
|
||||
|
||||
describe("CentralCore layer-less initialization", () => {
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks();
|
||||
mocks.createCentralBackendLayer.mockResolvedValue({
|
||||
asyncLayer: ownedLayer,
|
||||
releaseConnections: mocks.releaseConnections,
|
||||
shutdown: mocks.shutdown,
|
||||
});
|
||||
mocks.getLocalNode.mockResolvedValue(undefined);
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
for (const dir of cleanupDirs.splice(0)) {
|
||||
rmSync(dir, { recursive: true, force: true });
|
||||
}
|
||||
});
|
||||
|
||||
/**
|
||||
* FNXC:CentralCore 2026-07-29-16:10:
|
||||
* Standalone callers must bootstrap the exact layer they own and release that lifecycle on close; concurrent init calls share the same allocation.
|
||||
*/
|
||||
it("boots and owns a PostgreSQL layer for standalone CLI callers", async () => {
|
||||
const globalDir = mkdtempSync(join(tmpdir(), "fusion-central-cli-init-"));
|
||||
cleanupDirs.push(globalDir);
|
||||
const central = new CentralCore(globalDir);
|
||||
|
||||
await central.init();
|
||||
|
||||
expect(mocks.createCentralBackendLayer).toHaveBeenCalledWith({ globalSettingsDir: globalDir });
|
||||
expect(mocks.ensureBackendBootstrap).toHaveBeenCalledWith(ownedLayer);
|
||||
expect(central.asyncLayer).toBe(ownedLayer);
|
||||
expect(central.backendMode).toBe(true);
|
||||
|
||||
await central.close();
|
||||
|
||||
expect(mocks.getLocalNode).not.toHaveBeenCalled();
|
||||
expect(mocks.shutdown).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
it("coalesces concurrent layer-less initialization into one owned backend", async () => {
|
||||
const globalDir = mkdtempSync(join(tmpdir(), "fusion-central-cli-concurrent-init-"));
|
||||
cleanupDirs.push(globalDir);
|
||||
const central = new CentralCore(globalDir);
|
||||
|
||||
await Promise.all([central.init(), central.init()]);
|
||||
|
||||
expect(mocks.createCentralBackendLayer).toHaveBeenCalledOnce();
|
||||
expect(mocks.ensureBackendBootstrap).toHaveBeenCalledOnce();
|
||||
expect(mocks.ensureBackendBootstrap).toHaveBeenCalledWith(ownedLayer);
|
||||
|
||||
await central.close();
|
||||
expect(mocks.shutdown).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
it("waits for in-flight initialization before closing its owned backend", async () => {
|
||||
const globalDir = mkdtempSync(join(tmpdir(), "fusion-central-cli-close-race-"));
|
||||
cleanupDirs.push(globalDir);
|
||||
const central = new CentralCore(globalDir);
|
||||
let finishBackend!: (value: {
|
||||
asyncLayer: typeof ownedLayer;
|
||||
releaseConnections: typeof mocks.releaseConnections;
|
||||
shutdown: typeof mocks.shutdown;
|
||||
}) => void;
|
||||
mocks.createCentralBackendLayer.mockReturnValueOnce(new Promise((resolve) => {
|
||||
finishBackend = resolve;
|
||||
}));
|
||||
|
||||
const initialization = central.init();
|
||||
await vi.waitFor(() => expect(mocks.createCentralBackendLayer).toHaveBeenCalledOnce());
|
||||
const closing = central.close();
|
||||
expect(mocks.shutdown).not.toHaveBeenCalled();
|
||||
finishBackend({
|
||||
asyncLayer: ownedLayer,
|
||||
releaseConnections: mocks.releaseConnections,
|
||||
shutdown: mocks.shutdown,
|
||||
});
|
||||
|
||||
await Promise.all([initialization, closing]);
|
||||
expect(mocks.shutdown).toHaveBeenCalledOnce();
|
||||
expect(central.backendMode).toBe(false);
|
||||
});
|
||||
|
||||
it("serializes close behind an in-flight layer attachment", async () => {
|
||||
const globalDir = mkdtempSync(join(tmpdir(), "fusion-central-cli-attach-close-race-"));
|
||||
cleanupDirs.push(globalDir);
|
||||
const central = new CentralCore(globalDir);
|
||||
await central.init();
|
||||
const sharedLayer = { db: { shared: true } };
|
||||
let finishAttachment!: () => void;
|
||||
mocks.ensureBackendBootstrap.mockImplementationOnce(
|
||||
() => new Promise<void>((resolve) => {
|
||||
finishAttachment = resolve;
|
||||
}),
|
||||
);
|
||||
|
||||
const attachment = central.attachBackendLayer(
|
||||
sharedLayer as unknown as Parameters<CentralCore["attachBackendLayer"]>[0],
|
||||
);
|
||||
await vi.waitFor(() => expect(mocks.ensureBackendBootstrap).toHaveBeenCalledTimes(2));
|
||||
const closing = central.close();
|
||||
expect(mocks.shutdown).not.toHaveBeenCalled();
|
||||
finishAttachment();
|
||||
|
||||
await Promise.all([attachment, closing]);
|
||||
expect(mocks.shutdown).toHaveBeenCalledOnce();
|
||||
expect(central.backendMode).toBe(false);
|
||||
});
|
||||
|
||||
/**
|
||||
* FNXC:CentralPostgresCutover 2026-07-29-17:43:
|
||||
* Closing CentralCore is terminal because cleanup removes listeners and releases owned resources. Later queued initialization or attachment must not revive a partially torn-down instance.
|
||||
*/
|
||||
it("rejects initialization and layer attachment once close is requested", async () => {
|
||||
const globalDir = mkdtempSync(join(tmpdir(), "fusion-central-cli-terminal-close-"));
|
||||
cleanupDirs.push(globalDir);
|
||||
const central = new CentralCore(globalDir);
|
||||
await central.init();
|
||||
let finishShutdown!: () => void;
|
||||
mocks.shutdown.mockImplementationOnce(() => new Promise<void>((resolve) => {
|
||||
finishShutdown = resolve;
|
||||
}));
|
||||
|
||||
const closing = central.close();
|
||||
const sharedLayer = { db: { shared: true } };
|
||||
const initializationAfterClose = central.init();
|
||||
const attachmentAfterClose = central.attachBackendLayer(
|
||||
sharedLayer as unknown as Parameters<CentralCore["attachBackendLayer"]>[0],
|
||||
);
|
||||
|
||||
await expect(initializationAfterClose).rejects.toThrow("CentralCore is closed");
|
||||
await expect(attachmentAfterClose).rejects.toThrow("CentralCore is closed");
|
||||
await vi.waitFor(() => expect(mocks.shutdown).toHaveBeenCalledOnce());
|
||||
|
||||
finishShutdown();
|
||||
await closing;
|
||||
|
||||
expect(mocks.createCentralBackendLayer).toHaveBeenCalledOnce();
|
||||
expect(central.backendMode).toBe(false);
|
||||
});
|
||||
});
|
||||
@@ -215,6 +215,16 @@ export class CentralCore extends EventEmitter<CentralCoreEvents> {
|
||||
private readonly ensureGitRepositoryForProjectPath: typeof ensureGitRepositoryForProjectPath;
|
||||
private ownedBackendShutdown: (() => Promise<void>) | null = null;
|
||||
private ownedBackendReleaseConnections: (() => Promise<void>) | null = null;
|
||||
private initializationPromise: Promise<void> | null = null;
|
||||
private lifecycleOperation: Promise<void> = Promise.resolve();
|
||||
private closeRequested = false;
|
||||
private closed = false;
|
||||
|
||||
private runLifecycleOperation<T>(operation: () => Promise<T>): Promise<T> {
|
||||
const result = this.lifecycleOperation.then(operation, operation);
|
||||
this.lifecycleOperation = result.then(() => undefined, () => undefined);
|
||||
return result;
|
||||
}
|
||||
|
||||
/**
|
||||
* FNXC:CentralCore 2026-06-26-12:30:
|
||||
@@ -242,9 +252,15 @@ export class CentralCore extends EventEmitter<CentralCoreEvents> {
|
||||
* call, backendMode is true and all methods delegate to PostgreSQL.
|
||||
*/
|
||||
async attachBackendLayer(layer: AsyncDataLayer): Promise<void> {
|
||||
this.assertAcceptingOperations();
|
||||
if (!layer) {
|
||||
throw new Error("attachBackendLayer requires a non-null AsyncDataLayer");
|
||||
}
|
||||
return this.runLifecycleOperation(() => this.attachBackendLayerOnce(layer));
|
||||
}
|
||||
|
||||
private async attachBackendLayerOnce(layer: AsyncDataLayer): Promise<void> {
|
||||
this.assertOpen();
|
||||
// Release a central-only pool before adopting the runtime's shared layer.
|
||||
if (this.ownedBackendReleaseConnections) {
|
||||
/*
|
||||
@@ -270,7 +286,7 @@ export class CentralCore extends EventEmitter<CentralCoreEvents> {
|
||||
// post-construction injection point.
|
||||
(this as { asyncLayer: AsyncDataLayer | null }).asyncLayer = layer;
|
||||
this.initialized = false;
|
||||
await this.init();
|
||||
await this.initializeOnce();
|
||||
}
|
||||
|
||||
private readonly onDiscoveryNodeDiscovered = (node: DiscoveredNode): void => {
|
||||
@@ -318,19 +334,33 @@ export class CentralCore extends EventEmitter<CentralCoreEvents> {
|
||||
* Idempotent — safe to call multiple times.
|
||||
*/
|
||||
async init(): Promise<void> {
|
||||
this.assertAcceptingOperations();
|
||||
if (this.initialized) return;
|
||||
|
||||
/*
|
||||
FNXC:SqliteDualPathCleanup 2026-07-26-14:15:
|
||||
CentralCore.init is PostgreSQL-only. When no asyncLayer is attached yet, mark initialized without opening SQLite; attachBackendLayer bootstraps PG later.
|
||||
*/
|
||||
* FNXC:CentralCore 2026-07-29-16:10:
|
||||
* Layer-less initialization allocates an owned PostgreSQL lifecycle. Concurrent callers must share one in-flight attempt so a second backend cannot be orphaned when ownership fields are overwritten. Clear the promise after either outcome so a failed bootstrap remains retryable.
|
||||
*/
|
||||
if (!this.initializationPromise) {
|
||||
this.initializationPromise = this.runLifecycleOperation(() => this.initializeOnce());
|
||||
}
|
||||
const initialization = this.initializationPromise;
|
||||
try {
|
||||
await initialization;
|
||||
} finally {
|
||||
if (this.initializationPromise === initialization) this.initializationPromise = null;
|
||||
}
|
||||
}
|
||||
|
||||
private async initializeOnce(): Promise<void> {
|
||||
this.assertOpen();
|
||||
if (this.initialized) return;
|
||||
|
||||
if (this.asyncLayer) {
|
||||
await asyncCentralCore.ensureBackendBootstrap(this.asyncLayer);
|
||||
this.initialized = true;
|
||||
return;
|
||||
}
|
||||
this.initialized = true;
|
||||
return;
|
||||
|
||||
/*
|
||||
* FNXC:CentralPostgresCutover 2026-07-14-17:14:
|
||||
@@ -364,14 +394,23 @@ export class CentralCore extends EventEmitter<CentralCoreEvents> {
|
||||
* Closes database connections and releases resources.
|
||||
*/
|
||||
async close(): Promise<void> {
|
||||
this.closeRequested = true;
|
||||
return this.runLifecycleOperation(() => this.closeOnce());
|
||||
}
|
||||
|
||||
private async closeOnce(): Promise<void> {
|
||||
/*
|
||||
FNXC:CentralPostgresCutover 2026-07-29-16:26:
|
||||
Initialization, layer replacement, and close share one lifecycle queue. Cleanup must observe and release the backend the preceding operation publishes instead of returning early, leaking it, or letting attachment revive a core after shutdown.
|
||||
|
||||
FNXC:CentralPostgresCutover 2026-07-29-17:43:
|
||||
Close is terminal. Operations queued after cleanup must fail instead of allocating or attaching a backend after listeners and owned resources have been released.
|
||||
*/
|
||||
this.closed = true;
|
||||
if (this.nodeDiscovery) {
|
||||
this.stopDiscovery();
|
||||
}
|
||||
|
||||
await this.markLocalNodeOffline().catch((error) => {
|
||||
severityAuditLog.warn("[central-core] Failed to persist local node offline during close", error);
|
||||
});
|
||||
|
||||
// FNXC:CentralCore 2026-06-26-12:30: In backend mode there is no SQLite
|
||||
// CentralDatabase to close; the shared connection pool is owned by the
|
||||
// TaskStore/startup factory. CentralCore does not close the pool.
|
||||
@@ -389,6 +428,14 @@ export class CentralCore extends EventEmitter<CentralCoreEvents> {
|
||||
this.removeAllListeners();
|
||||
}
|
||||
|
||||
private assertOpen(): void {
|
||||
if (this.closed) throw new Error("CentralCore is closed");
|
||||
}
|
||||
|
||||
private assertAcceptingOperations(): void {
|
||||
if (this.closeRequested) throw new Error("CentralCore is closed");
|
||||
}
|
||||
|
||||
/** Persist the local mesh node's terminal state before its backend closes. */
|
||||
async markLocalNodeOffline(): Promise<void> {
|
||||
if (!this.initialized) return;
|
||||
|
||||
@@ -4,8 +4,8 @@
|
||||
* while backup construction resolves synchronously. This process-local registry
|
||||
* bridges that gap without logging credentials. Leases represent individual
|
||||
* lifecycles within a physical cluster generation: a joiner's release cannot
|
||||
* clear a newer generation, and owner shutdown invalidates every lease because
|
||||
* it is the only lifecycle that actually stops the postmaster.
|
||||
* clear a newer generation, and owner shutdown waits for every live lease
|
||||
* because physical process ownership is not exclusive logical usage.
|
||||
*/
|
||||
|
||||
/** Opaque handle for one embedded-backend lifecycle registration. */
|
||||
@@ -20,6 +20,19 @@ interface Generation {
|
||||
readonly id: number;
|
||||
readonly leases: Set<EmbeddedRuntimeLease>;
|
||||
latestRegistration: number;
|
||||
pendingOwnerStop: (() => Promise<void>) | null;
|
||||
stopCompletion: Promise<void> | null;
|
||||
stopping: boolean;
|
||||
}
|
||||
|
||||
export class EmbeddedRuntimeStoppingError extends Error {
|
||||
constructor(
|
||||
readonly url: string,
|
||||
readonly completion: Promise<void>,
|
||||
) {
|
||||
super("Embedded PostgreSQL runtime is stopping");
|
||||
this.name = "EmbeddedRuntimeStoppingError";
|
||||
}
|
||||
}
|
||||
|
||||
interface LeaseMetadata {
|
||||
@@ -41,12 +54,27 @@ export function registerEmbeddedRuntimeUrl(
|
||||
options: { ownsProcess: boolean },
|
||||
): EmbeddedRuntimeLease {
|
||||
let generation = generationsByUrl.get(url);
|
||||
if (generation?.stopping) {
|
||||
if (!generation.stopCompletion) {
|
||||
throw new Error("Embedded PostgreSQL runtime stop completion is missing");
|
||||
}
|
||||
throw new EmbeddedRuntimeStoppingError(url, generation.stopCompletion);
|
||||
}
|
||||
// FNXC:PostgresBackup 2026-07-16-12:40: An owner started a new postmaster,
|
||||
// so URL reuse must create a new generation rather than retain stale leases.
|
||||
if (!generation || options.ownsProcess) {
|
||||
const id = (nextGenerationByUrl.get(url) ?? 0) + 1;
|
||||
nextGenerationByUrl.set(url, id);
|
||||
generation = { url, epoch: registryEpoch, id, leases: new Set(), latestRegistration: 0 };
|
||||
generation = {
|
||||
url,
|
||||
epoch: registryEpoch,
|
||||
id,
|
||||
leases: new Set(),
|
||||
latestRegistration: 0,
|
||||
pendingOwnerStop: null,
|
||||
stopCompletion: null,
|
||||
stopping: false,
|
||||
};
|
||||
generationsByUrl.set(url, generation);
|
||||
}
|
||||
|
||||
@@ -62,8 +90,16 @@ export function registerEmbeddedRuntimeUrl(
|
||||
return lease;
|
||||
}
|
||||
|
||||
/** Release exactly one lifecycle lease; stale generation handles are inert. */
|
||||
export function releaseEmbeddedRuntimeLease(lease: EmbeddedRuntimeLease): void {
|
||||
/**
|
||||
* Release exactly one lifecycle lease; stale generation handles are inert.
|
||||
*
|
||||
* FNXC:PostgresResourceLifecycle 2026-07-29-16:10:
|
||||
* An embedded-process owner may close before joined consumers. Record its stop callback and run it only after the final lease releases so short-lived central/CLI cleanup cannot terminate PostgreSQL beneath another live store.
|
||||
*/
|
||||
export async function releaseEmbeddedRuntimeLease(
|
||||
lease: EmbeddedRuntimeLease,
|
||||
options: { stopOwner?: () => Promise<void> } = {},
|
||||
): Promise<void> {
|
||||
const metadata = leaseMetadata.get(lease);
|
||||
if (!metadata) return;
|
||||
const generation = generationsByUrl.get(metadata.url);
|
||||
@@ -74,8 +110,33 @@ export function releaseEmbeddedRuntimeLease(lease: EmbeddedRuntimeLease): void {
|
||||
) return;
|
||||
|
||||
generation.leases.delete(lease);
|
||||
if (metadata.ownsProcess && options.stopOwner) {
|
||||
generation.pendingOwnerStop = options.stopOwner;
|
||||
}
|
||||
if (generation.leases.size === 0) {
|
||||
generationsByUrl.delete(metadata.url);
|
||||
const stopOwner = generation.pendingOwnerStop;
|
||||
generation.pendingOwnerStop = null;
|
||||
if (stopOwner) {
|
||||
/*
|
||||
FNXC:PostgresLifecycle 2026-07-29-16:26:
|
||||
Keep the generation visible as stopping until the owner callback completes. A concurrent bootstrap must retry rather than join a postmaster that is already committed to termination.
|
||||
|
||||
FNXC:PostgresLifecycle 2026-07-29-17:43:
|
||||
Publish the actual stop completion to rejected registrants. Startup waits on lifecycle completion instead of exhausting a fixed retry window while an orderly shutdown is still in progress.
|
||||
*/
|
||||
generation.stopping = true;
|
||||
const stopCompletion = Promise.resolve().then(stopOwner);
|
||||
generation.stopCompletion = stopCompletion;
|
||||
try {
|
||||
await stopCompletion;
|
||||
} finally {
|
||||
if (generationsByUrl.get(metadata.url) === generation) {
|
||||
generationsByUrl.delete(metadata.url);
|
||||
}
|
||||
}
|
||||
} else {
|
||||
generationsByUrl.delete(metadata.url);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -103,7 +164,7 @@ export function invalidateEmbeddedRuntimeUrl(url: string, lease?: EmbeddedRuntim
|
||||
export function getActiveEmbeddedRuntimeUrl(): string | undefined {
|
||||
let latest: Generation | undefined;
|
||||
for (const generation of generationsByUrl.values()) {
|
||||
if (generation.leases.size > 0 && (!latest || generation.latestRegistration > latest.latestRegistration)) {
|
||||
if (!generation.stopping && generation.leases.size > 0 && (!latest || generation.latestRegistration > latest.latestRegistration)) {
|
||||
latest = generation;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -58,7 +58,7 @@ import { applySchemaBaseline, MIGRATION_BOOKKEEPING_TABLE } from "./schema-appli
|
||||
import type { PostgresJsDatabase } from "drizzle-orm/postgres-js";
|
||||
import { createAsyncDataLayer, type AsyncDataLayer } from "./data-layer.js";
|
||||
import {
|
||||
invalidateEmbeddedRuntimeUrl,
|
||||
EmbeddedRuntimeStoppingError,
|
||||
registerEmbeddedRuntimeUrl,
|
||||
releaseEmbeddedRuntimeLease,
|
||||
type EmbeddedRuntimeLease,
|
||||
@@ -226,8 +226,8 @@ interface SchemaBackendBootResult {
|
||||
|
||||
/**
|
||||
* FNXC:PostgresBackup 2026-07-16-12:40:
|
||||
* stop() only stops a postmaster owned by this lifecycle. Consequently owner
|
||||
* teardown burns the URL generation for every joiner, while joiner teardown
|
||||
* stop() only stops a postmaster owned by this lifecycle. The runtime registry
|
||||
* defers that owner stop until joined leases release, while joiner teardown
|
||||
* releases only its opaque lease. This keeps synchronous backup resolution
|
||||
* aligned with physical cluster liveness without logging the credential URL.
|
||||
*/
|
||||
@@ -237,16 +237,20 @@ async function stopEmbeddedRuntime(
|
||||
runtimeUrl: string | null,
|
||||
ownsProcess: boolean,
|
||||
): Promise<void> {
|
||||
if (!lease || !runtimeUrl) {
|
||||
await lifecycle?.stop();
|
||||
return;
|
||||
}
|
||||
if (ownsProcess) {
|
||||
await releaseEmbeddedRuntimeLease(lease, {
|
||||
stopOwner: async () => { await lifecycle?.stop(); },
|
||||
});
|
||||
return;
|
||||
}
|
||||
try {
|
||||
await lifecycle?.stop();
|
||||
} finally {
|
||||
if (lease && runtimeUrl) {
|
||||
if (ownsProcess) {
|
||||
invalidateEmbeddedRuntimeUrl(runtimeUrl, lease);
|
||||
} else {
|
||||
releaseEmbeddedRuntimeLease(lease);
|
||||
}
|
||||
}
|
||||
await releaseEmbeddedRuntimeLease(lease);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -382,6 +386,14 @@ async function bootSchemaBackend(
|
||||
try {
|
||||
return await bootSchemaBackendOnce(options, bypassProjectIsolation);
|
||||
} catch (error) {
|
||||
if (error instanceof EmbeddedRuntimeStoppingError) {
|
||||
/*
|
||||
FNXC:PostgresLifecycle 2026-07-29-17:43:
|
||||
A replacement backend waits for the shared registry's actual owner-stop completion. Fixed backoff budgets can expire during a valid slow shutdown and turn orderly handoff into startup failure.
|
||||
*/
|
||||
await error.completion;
|
||||
continue;
|
||||
}
|
||||
if (
|
||||
error instanceof JoinedInstanceUnreachableError &&
|
||||
joinedRetryAttempt < JOINED_INSTANCE_RETRY_DELAYS_MS.length
|
||||
@@ -470,6 +482,7 @@ async function bootSchemaBackendOnce(
|
||||
}
|
||||
} catch (error) {
|
||||
await embeddedLifecycle.stop().catch(() => undefined);
|
||||
if (error instanceof EmbeddedRuntimeStoppingError) throw error;
|
||||
throw new Error(
|
||||
`startup-factory: failed to start embedded PostgreSQL: ${error instanceof Error ? error.message : String(error)}`,
|
||||
);
|
||||
@@ -1371,7 +1384,7 @@ export async function createTaskStoreForBackend(
|
||||
if (shutdownEmbedded) {
|
||||
try {
|
||||
(shutdownEmbedded as unknown as { detachWithoutStop?: () => void }).detachWithoutStop?.();
|
||||
if (embeddedRuntimeLease) releaseEmbeddedRuntimeLease(embeddedRuntimeLease);
|
||||
if (embeddedRuntimeLease) await releaseEmbeddedRuntimeLease(embeddedRuntimeLease);
|
||||
} catch (err) {
|
||||
log.warn(`startup-factory: embedded PostgreSQL detach failed: ${
|
||||
err instanceof Error ? err.message : String(err)
|
||||
|
||||
Reference in New Issue
Block a user