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";
|
} from "../backup.js";
|
||||||
import {
|
import {
|
||||||
clearActiveEmbeddedRuntimeUrl,
|
clearActiveEmbeddedRuntimeUrl,
|
||||||
|
EmbeddedRuntimeStoppingError,
|
||||||
getActiveEmbeddedRuntimeUrl,
|
getActiveEmbeddedRuntimeUrl,
|
||||||
invalidateEmbeddedRuntimeUrl,
|
invalidateEmbeddedRuntimeUrl,
|
||||||
registerEmbeddedRuntimeUrl,
|
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 owner = registerEmbeddedRuntimeUrl(embeddedUrl, { ownsProcess: true });
|
||||||
const joiner = registerEmbeddedRuntimeUrl(embeddedUrl, { ownsProcess: false });
|
const joiner = registerEmbeddedRuntimeUrl(embeddedUrl, { ownsProcess: false });
|
||||||
|
|
||||||
releaseEmbeddedRuntimeLease(joiner);
|
await releaseEmbeddedRuntimeLease(joiner);
|
||||||
expect(getActiveEmbeddedRuntimeUrl()).toBe(embeddedUrl);
|
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();
|
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 readonly ensureGitRepositoryForProjectPath: typeof ensureGitRepositoryForProjectPath;
|
||||||
private ownedBackendShutdown: (() => Promise<void>) | null = null;
|
private ownedBackendShutdown: (() => Promise<void>) | null = null;
|
||||||
private ownedBackendReleaseConnections: (() => 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:
|
* 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.
|
* call, backendMode is true and all methods delegate to PostgreSQL.
|
||||||
*/
|
*/
|
||||||
async attachBackendLayer(layer: AsyncDataLayer): Promise<void> {
|
async attachBackendLayer(layer: AsyncDataLayer): Promise<void> {
|
||||||
|
this.assertAcceptingOperations();
|
||||||
if (!layer) {
|
if (!layer) {
|
||||||
throw new Error("attachBackendLayer requires a non-null AsyncDataLayer");
|
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.
|
// Release a central-only pool before adopting the runtime's shared layer.
|
||||||
if (this.ownedBackendReleaseConnections) {
|
if (this.ownedBackendReleaseConnections) {
|
||||||
/*
|
/*
|
||||||
@@ -270,7 +286,7 @@ export class CentralCore extends EventEmitter<CentralCoreEvents> {
|
|||||||
// post-construction injection point.
|
// post-construction injection point.
|
||||||
(this as { asyncLayer: AsyncDataLayer | null }).asyncLayer = layer;
|
(this as { asyncLayer: AsyncDataLayer | null }).asyncLayer = layer;
|
||||||
this.initialized = false;
|
this.initialized = false;
|
||||||
await this.init();
|
await this.initializeOnce();
|
||||||
}
|
}
|
||||||
|
|
||||||
private readonly onDiscoveryNodeDiscovered = (node: DiscoveredNode): void => {
|
private readonly onDiscoveryNodeDiscovered = (node: DiscoveredNode): void => {
|
||||||
@@ -318,19 +334,33 @@ export class CentralCore extends EventEmitter<CentralCoreEvents> {
|
|||||||
* Idempotent — safe to call multiple times.
|
* Idempotent — safe to call multiple times.
|
||||||
*/
|
*/
|
||||||
async init(): Promise<void> {
|
async init(): Promise<void> {
|
||||||
|
this.assertAcceptingOperations();
|
||||||
if (this.initialized) return;
|
if (this.initialized) return;
|
||||||
|
|
||||||
/*
|
/*
|
||||||
FNXC:SqliteDualPathCleanup 2026-07-26-14:15:
|
* FNXC:CentralCore 2026-07-29-16:10:
|
||||||
CentralCore.init is PostgreSQL-only. When no asyncLayer is attached yet, mark initialized without opening SQLite; attachBackendLayer bootstraps PG later.
|
* 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) {
|
if (this.asyncLayer) {
|
||||||
await asyncCentralCore.ensureBackendBootstrap(this.asyncLayer);
|
await asyncCentralCore.ensureBackendBootstrap(this.asyncLayer);
|
||||||
this.initialized = true;
|
this.initialized = true;
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
this.initialized = true;
|
|
||||||
return;
|
|
||||||
|
|
||||||
/*
|
/*
|
||||||
* FNXC:CentralPostgresCutover 2026-07-14-17:14:
|
* FNXC:CentralPostgresCutover 2026-07-14-17:14:
|
||||||
@@ -364,14 +394,23 @@ export class CentralCore extends EventEmitter<CentralCoreEvents> {
|
|||||||
* Closes database connections and releases resources.
|
* Closes database connections and releases resources.
|
||||||
*/
|
*/
|
||||||
async close(): Promise<void> {
|
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) {
|
if (this.nodeDiscovery) {
|
||||||
this.stopDiscovery();
|
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
|
// 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
|
// CentralDatabase to close; the shared connection pool is owned by the
|
||||||
// TaskStore/startup factory. CentralCore does not close the pool.
|
// TaskStore/startup factory. CentralCore does not close the pool.
|
||||||
@@ -389,6 +428,14 @@ export class CentralCore extends EventEmitter<CentralCoreEvents> {
|
|||||||
this.removeAllListeners();
|
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. */
|
/** Persist the local mesh node's terminal state before its backend closes. */
|
||||||
async markLocalNodeOffline(): Promise<void> {
|
async markLocalNodeOffline(): Promise<void> {
|
||||||
if (!this.initialized) return;
|
if (!this.initialized) return;
|
||||||
|
|||||||
@@ -4,8 +4,8 @@
|
|||||||
* while backup construction resolves synchronously. This process-local registry
|
* while backup construction resolves synchronously. This process-local registry
|
||||||
* bridges that gap without logging credentials. Leases represent individual
|
* bridges that gap without logging credentials. Leases represent individual
|
||||||
* lifecycles within a physical cluster generation: a joiner's release cannot
|
* lifecycles within a physical cluster generation: a joiner's release cannot
|
||||||
* clear a newer generation, and owner shutdown invalidates every lease because
|
* clear a newer generation, and owner shutdown waits for every live lease
|
||||||
* it is the only lifecycle that actually stops the postmaster.
|
* because physical process ownership is not exclusive logical usage.
|
||||||
*/
|
*/
|
||||||
|
|
||||||
/** Opaque handle for one embedded-backend lifecycle registration. */
|
/** Opaque handle for one embedded-backend lifecycle registration. */
|
||||||
@@ -20,6 +20,19 @@ interface Generation {
|
|||||||
readonly id: number;
|
readonly id: number;
|
||||||
readonly leases: Set<EmbeddedRuntimeLease>;
|
readonly leases: Set<EmbeddedRuntimeLease>;
|
||||||
latestRegistration: number;
|
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 {
|
interface LeaseMetadata {
|
||||||
@@ -41,12 +54,27 @@ export function registerEmbeddedRuntimeUrl(
|
|||||||
options: { ownsProcess: boolean },
|
options: { ownsProcess: boolean },
|
||||||
): EmbeddedRuntimeLease {
|
): EmbeddedRuntimeLease {
|
||||||
let generation = generationsByUrl.get(url);
|
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,
|
// 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.
|
// so URL reuse must create a new generation rather than retain stale leases.
|
||||||
if (!generation || options.ownsProcess) {
|
if (!generation || options.ownsProcess) {
|
||||||
const id = (nextGenerationByUrl.get(url) ?? 0) + 1;
|
const id = (nextGenerationByUrl.get(url) ?? 0) + 1;
|
||||||
nextGenerationByUrl.set(url, id);
|
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);
|
generationsByUrl.set(url, generation);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -62,8 +90,16 @@ export function registerEmbeddedRuntimeUrl(
|
|||||||
return lease;
|
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);
|
const metadata = leaseMetadata.get(lease);
|
||||||
if (!metadata) return;
|
if (!metadata) return;
|
||||||
const generation = generationsByUrl.get(metadata.url);
|
const generation = generationsByUrl.get(metadata.url);
|
||||||
@@ -74,8 +110,33 @@ export function releaseEmbeddedRuntimeLease(lease: EmbeddedRuntimeLease): void {
|
|||||||
) return;
|
) return;
|
||||||
|
|
||||||
generation.leases.delete(lease);
|
generation.leases.delete(lease);
|
||||||
|
if (metadata.ownsProcess && options.stopOwner) {
|
||||||
|
generation.pendingOwnerStop = options.stopOwner;
|
||||||
|
}
|
||||||
if (generation.leases.size === 0) {
|
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 {
|
export function getActiveEmbeddedRuntimeUrl(): string | undefined {
|
||||||
let latest: Generation | undefined;
|
let latest: Generation | undefined;
|
||||||
for (const generation of generationsByUrl.values()) {
|
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;
|
latest = generation;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -58,7 +58,7 @@ import { applySchemaBaseline, MIGRATION_BOOKKEEPING_TABLE } from "./schema-appli
|
|||||||
import type { PostgresJsDatabase } from "drizzle-orm/postgres-js";
|
import type { PostgresJsDatabase } from "drizzle-orm/postgres-js";
|
||||||
import { createAsyncDataLayer, type AsyncDataLayer } from "./data-layer.js";
|
import { createAsyncDataLayer, type AsyncDataLayer } from "./data-layer.js";
|
||||||
import {
|
import {
|
||||||
invalidateEmbeddedRuntimeUrl,
|
EmbeddedRuntimeStoppingError,
|
||||||
registerEmbeddedRuntimeUrl,
|
registerEmbeddedRuntimeUrl,
|
||||||
releaseEmbeddedRuntimeLease,
|
releaseEmbeddedRuntimeLease,
|
||||||
type EmbeddedRuntimeLease,
|
type EmbeddedRuntimeLease,
|
||||||
@@ -226,8 +226,8 @@ interface SchemaBackendBootResult {
|
|||||||
|
|
||||||
/**
|
/**
|
||||||
* FNXC:PostgresBackup 2026-07-16-12:40:
|
* FNXC:PostgresBackup 2026-07-16-12:40:
|
||||||
* stop() only stops a postmaster owned by this lifecycle. Consequently owner
|
* stop() only stops a postmaster owned by this lifecycle. The runtime registry
|
||||||
* teardown burns the URL generation for every joiner, while joiner teardown
|
* defers that owner stop until joined leases release, while joiner teardown
|
||||||
* releases only its opaque lease. This keeps synchronous backup resolution
|
* releases only its opaque lease. This keeps synchronous backup resolution
|
||||||
* aligned with physical cluster liveness without logging the credential URL.
|
* aligned with physical cluster liveness without logging the credential URL.
|
||||||
*/
|
*/
|
||||||
@@ -237,16 +237,20 @@ async function stopEmbeddedRuntime(
|
|||||||
runtimeUrl: string | null,
|
runtimeUrl: string | null,
|
||||||
ownsProcess: boolean,
|
ownsProcess: boolean,
|
||||||
): Promise<void> {
|
): Promise<void> {
|
||||||
|
if (!lease || !runtimeUrl) {
|
||||||
|
await lifecycle?.stop();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (ownsProcess) {
|
||||||
|
await releaseEmbeddedRuntimeLease(lease, {
|
||||||
|
stopOwner: async () => { await lifecycle?.stop(); },
|
||||||
|
});
|
||||||
|
return;
|
||||||
|
}
|
||||||
try {
|
try {
|
||||||
await lifecycle?.stop();
|
await lifecycle?.stop();
|
||||||
} finally {
|
} finally {
|
||||||
if (lease && runtimeUrl) {
|
await releaseEmbeddedRuntimeLease(lease);
|
||||||
if (ownsProcess) {
|
|
||||||
invalidateEmbeddedRuntimeUrl(runtimeUrl, lease);
|
|
||||||
} else {
|
|
||||||
releaseEmbeddedRuntimeLease(lease);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -382,6 +386,14 @@ async function bootSchemaBackend(
|
|||||||
try {
|
try {
|
||||||
return await bootSchemaBackendOnce(options, bypassProjectIsolation);
|
return await bootSchemaBackendOnce(options, bypassProjectIsolation);
|
||||||
} catch (error) {
|
} 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 (
|
if (
|
||||||
error instanceof JoinedInstanceUnreachableError &&
|
error instanceof JoinedInstanceUnreachableError &&
|
||||||
joinedRetryAttempt < JOINED_INSTANCE_RETRY_DELAYS_MS.length
|
joinedRetryAttempt < JOINED_INSTANCE_RETRY_DELAYS_MS.length
|
||||||
@@ -470,6 +482,7 @@ async function bootSchemaBackendOnce(
|
|||||||
}
|
}
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
await embeddedLifecycle.stop().catch(() => undefined);
|
await embeddedLifecycle.stop().catch(() => undefined);
|
||||||
|
if (error instanceof EmbeddedRuntimeStoppingError) throw error;
|
||||||
throw new Error(
|
throw new Error(
|
||||||
`startup-factory: failed to start embedded PostgreSQL: ${error instanceof Error ? error.message : String(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) {
|
if (shutdownEmbedded) {
|
||||||
try {
|
try {
|
||||||
(shutdownEmbedded as unknown as { detachWithoutStop?: () => void }).detachWithoutStop?.();
|
(shutdownEmbedded as unknown as { detachWithoutStop?: () => void }).detachWithoutStop?.();
|
||||||
if (embeddedRuntimeLease) releaseEmbeddedRuntimeLease(embeddedRuntimeLease);
|
if (embeddedRuntimeLease) await releaseEmbeddedRuntimeLease(embeddedRuntimeLease);
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
log.warn(`startup-factory: embedded PostgreSQL detach failed: ${
|
log.warn(`startup-factory: embedded PostgreSQL detach failed: ${
|
||||||
err instanceof Error ? err.message : String(err)
|
err instanceof Error ? err.message : String(err)
|
||||||
|
|||||||
Reference in New Issue
Block a user