diff --git a/doc/DATABASE.md b/doc/DATABASE.md index f194265751..2c87b1a3c7 100644 --- a/doc/DATABASE.md +++ b/doc/DATABASE.md @@ -247,6 +247,38 @@ public identity, and the existing unique singleton-key index makes concurrent or later attempts to replace it fail closed. The server loads the row before constructing URL-dependent runtime services on every boot. +### Unclaimed Cloud warm standby + +`PAPERCLIP_CLOUD_WARM_STANDBY=1` is an opt-in control-plane marker for an +unclaimed warm application. It requires Cloud configuration, a stack identity, +and runtime identity verification keys. After restoring the durable identity, +startup checks once for company data. Existing companies or a persisted claim +keep the application fully active. A failed database check fails startup. + +An empty unclaimed application keeps its HTTP server and sandbox plugins ready, +but skips recurring database work: chat/email delivery, plugin jobs, browser +cleanup, feedback export, import cleanup, execution reconciliation, heartbeat +schedules, and automatic backups. Startup migrations, plugin installation, and +other one-time preparation still run. Normal anonymous health probes return +`warmStandby: true` without SQL or session lookups. This is application liveness, +not a current database connectivity check. Other API requests and WebSocket upgrades return 503 until +the signed claim succeeds. Page and asset requests serve only the static UI +router, bypassing session, bearer-key, tenant, and other dynamic handlers. + +The existing signed claim on `GET /api/health` writes the identity durably before +normal requests and polling resume. No polling discovers claims and no process +restart is required. Timers resume on their next normal tick; request-driven +work can proceed immediately. Claim failure leaves standby intact. A restart +restores the claim even if provider environment alignment has not completed. +Deleting a claimed workspace's last company never puts it back into standby. + +Deploy support before enabling the marker. Validate idle database transactions +and health probes, claim/bootstrap latency, recurring work after claim, and a +restart with stale provider variables on an isolated warm application first. +Rollback by setting the marker to `0` and restarting. No schema changes are +required. This mechanism does not sleep claimed workspaces or replace a durable +scheduler for their background work. + ## Resource membership tables Paperclip stores current-user sidebar membership state in: diff --git a/server/src/__tests__/cloud-runtime-identity.test.ts b/server/src/__tests__/cloud-runtime-identity.test.ts index ae2cd9f9da..08749398ea 100644 --- a/server/src/__tests__/cloud-runtime-identity.test.ts +++ b/server/src/__tests__/cloud-runtime-identity.test.ts @@ -1,7 +1,9 @@ +import { createCloudWarmStandby } from "../services/cloud-warm-standby.js"; +import { cloudWarmStandbyMiddleware } from "../middleware/cloud-warm-standby.js"; import { generateKeyPairSync, sign } from "node:crypto"; import express from "express"; import request from "supertest"; -import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it } from "vitest"; +import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; import { createDb, instanceSettings } from "@paperclipai/db"; import { applyCloudRuntimeIdentityAssertion, @@ -106,6 +108,7 @@ describeEmbeddedPostgres("Cloud runtime identity", () => { }); afterEach(() => { + vi.restoreAllMocks(); for (const key of [ "PAPERCLIP_CLOUD_TENANT_SERVER_TOKEN", "PAPERCLIP_CLOUD_STACK_ID", @@ -183,6 +186,39 @@ describeEmbeddedPostgres("Cloud runtime identity", () => { }); }); + it("keeps probes idle until a verified durable claim and stays active across restart", async () => { + const env = { ...process.env, PAPERCLIP_CLOUD_WARM_STANDBY: "1" }; + const isStandby = await createCloudWarmStandby(db, env); + const execute = vi.spyOn(db, "execute"); + const health = healthRoutes(db, { + deploymentMode: "authenticated", deploymentExposure: "public", + authReady: true, companyDeletionEnabled: false, isWarmStandby: isStandby, + }); + const app = express(); + app.use(cloudRuntimeIdentityMiddleware(db)); + app.use(cloudWarmStandbyMiddleware(isStandby, health)); + app.use("/api/health", health); + expect((await request(app).get("/api/health")).body.warmStandby).toBe(true); + expect(execute).not.toHaveBeenCalled(); + const requestTime = Math.floor(Date.now() / 1000); + const invalid = await request(app).get("/api/health").set("x-paperclip-cloud-runtime-identity", assertion({ + claims: { sub: "another-stack", iat: requestTime, exp: requestTime + 300 }, + })); + expect(invalid.status).toBe(401); + expect(isStandby()).toBe(true); + const signed = assertion({ claims: { iat: requestTime, exp: requestTime + 300 } }); + const claimed = await request(app).get("/api/health").set("x-paperclip-cloud-runtime-identity", signed); + expect(claimed.status).toBe(200); + expect(claimed.body.warmStandby).toBeUndefined(); + expect(isStandby()).toBe(false); + expect(execute).toHaveBeenCalledOnce(); + expect((await request(app).get("/api/health").set("x-paperclip-cloud-runtime-identity", signed)).status).toBe(200); + resetCloudRuntimeIdentityForTests(); + process.env.PAPERCLIP_PUBLIC_URL = POOL_ORIGIN; + await initializeCloudRuntimeIdentity(db); + expect((await createCloudWarmStandby(db, env))()).toBe(false); + }); + it("restores the canonical identity before consumers read stale startup variables", async () => { await applyCloudRuntimeIdentityAssertion({ db, compactJws: assertion(), now: NOW }); resetCloudRuntimeIdentityForTests(); diff --git a/server/src/__tests__/cloud-warm-standby.test.ts b/server/src/__tests__/cloud-warm-standby.test.ts new file mode 100644 index 0000000000..c4c76476ae --- /dev/null +++ b/server/src/__tests__/cloud-warm-standby.test.ts @@ -0,0 +1,127 @@ +import { createServer, request as httpRequest } from "node:http"; +import express from "express"; +import request from "supertest"; +import { afterEach, describe, expect, it, vi } from "vitest"; +import type { Db } from "@paperclipai/db"; +import { cloudWarmStandbyMiddleware, cloudWarmStandbyServerOptions } from "../middleware/cloud-warm-standby.js"; +import { healthRoutes } from "../routes/health.js"; +import { emailChannelService } from "../services/email-channels.js"; +import { createPluginJobScheduler } from "../services/plugin-job-scheduler.js"; + +const healthOptions = { + deploymentMode: "authenticated" as const, + deploymentExposure: "public" as const, + authReady: true, + companyDeletionEnabled: false, + runtimeEnv: { PAPERCLIP_CLOUD_TENANT_SERVER_TOKEN: "test-token" }, +}; +afterEach(() => vi.useRealTimers()); + +describe("unclaimed Cloud background work", () => { + it("keeps probes and tenant requests out of SQL and auth until claim", async () => { + let standby = true; + const execute = vi.fn().mockResolvedValue([]); + const db = { execute } as unknown as Db; + const health = healthRoutes(db, { ...healthOptions, isWarmStandby: () => standby }); + const app = express(); + const auth = vi.fn((_req, _res, next) => next()); + const ui = express.Router(); + ui.get("/", (_req, res) => res.type("html").send("ready")); + ui.get("/assets/app.js", (_req, res) => res.type("js").send("ready")); + app.use(cloudWarmStandbyMiddleware(() => standby, health, ui)); + app.use(auth); + app.use("/api/health", health); + app.get("/api/companies", (_req, res) => res.json([])); + for (let i = 0; i < 3; i++) { + const probe = await request(app).get("/api/health").set("authorization", "Bearer ignored-in-standby"); + expect(probe.status).toBe(200); + expect(probe.body.warmStandby).toBe(true); + expect(probe.body.serverInfo).toBeUndefined(); + } + expect((await request(app).head("/api/health")).status).toBe(200); + expect((await request(app).get("/api/companies")).status).toBe(503); + expect((await request(app).post("/api/health/dev-server/restart")).status).toBe(503); + for (const path of ["/", "/assets/app.js"]) { + const page = await request(app).get(path) + .set("authorization", "Bearer test-bearer") + .set("x-paperclip-cloud-tenant-token", "test-tenant") + .set("cookie", "paperclip.session_token=test-session"); + expect(page.status).toBe(200); + } + expect((await request(app).post("/mcp/gateways/gw_test")).status).toBe(503); + expect((await request(app).get("/unknown-dynamic-handler")).status).toBe(503); + expect(auth).not.toHaveBeenCalled(); + expect(execute).not.toHaveBeenCalled(); + standby = false; + expect((await request(app).get("/api/companies")).status).toBe(200); + const claimed = await request(app).get("/api/health"); + expect(claimed.status).toBe(200); + expect(claimed.body.warmStandby).toBeUndefined(); + expect(execute).toHaveBeenCalledOnce(); + execute.mockRejectedValueOnce(new Error("offline")); + expect((await request(app).get("/api/health")).status).toBe(503); + }); + + it("rejects every WebSocket upgrade before auth and admits upgrades after claim", async () => { + let standby = true; + const app = express(); + const health = healthRoutes({} as Db, { ...healthOptions, isWarmStandby: () => standby }); + const authenticate = vi.fn(); + app.use(cloudWarmStandbyMiddleware(() => standby, health)); + const server = createServer(cloudWarmStandbyServerOptions(() => standby), app); + server.on("upgrade", (_req, socket) => { + authenticate(); + socket.end("HTTP/1.1 101 Switching Protocols\r\nConnection: Upgrade\r\nUpgrade: websocket\r\n\r\n"); + }); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + const address = server.address() as { port: number }; + const upgrade = (path: string) => new Promise((resolve, reject) => { + const req = httpRequest({ + host: "127.0.0.1", port: address.port, path, + headers: { connection: "Upgrade", upgrade: "websocket", authorization: "Bearer test-bearer" }, + }); + req.on("response", (res) => { res.resume(); resolve(res.statusCode!); }); + req.on("upgrade", (res, socket) => { socket.destroy(); resolve(res.statusCode!); }); + req.on("error", reject); + req.end(); + }); + try { + for (const path of ["/api/companies/company/events/ws?token=anything", "/api/health", "/prp", "/"]) { + expect(await upgrade(path)).toBe(503); + } + expect(authenticate).not.toHaveBeenCalled(); + standby = false; + expect(await upgrade("/api/companies/company/events/ws?token=anything")).toBe(101); + expect(authenticate).toHaveBeenCalledOnce(); + } finally { + await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve())); + } + }); + + it("email and plugin timers leave SQL idle, then resume without restarting", async () => { + vi.useFakeTimers(); + let enabled = false; + const select = vi.fn(() => { throw new Error("SQL probe"); }); + const db = { select } as unknown as Db; + const email = emailChannelService(db, { heartbeat: { wakeup: vi.fn() }, isBackgroundWorkEnabled: () => enabled }); + const scheduler = createPluginJobScheduler({ + db, + jobStore: {} as never, + workerManager: {} as never, + tickIntervalMs: 1000, + isBackgroundWorkEnabled: () => enabled, + }); + try { + email.start(); + scheduler.start(); + await vi.advanceTimersByTimeAsync(10 * 60_000); + expect(select).not.toHaveBeenCalled(); + enabled = true; + await vi.advanceTimersByTimeAsync(1000); + expect(select.mock.calls.length).toBeGreaterThanOrEqual(2); + } finally { + scheduler.stop(); + await email.shutdown(); + } + }); +}); diff --git a/server/src/app.ts b/server/src/app.ts index c99839e4fa..32a4534ce6 100644 --- a/server/src/app.ts +++ b/server/src/app.ts @@ -1,3 +1,5 @@ +import { cloudWarmStandbyMiddleware } from "./middleware/cloud-warm-standby.js"; +import type { CloudWarmStandby } from "./services/cloud-warm-standby.js"; import { browserUseRoutes } from "./routes/browser-use.js"; import { browserUseService } from "./services/browser-use.js"; import { slackToolRoutes } from "./routes/slack-tools.js"; @@ -461,6 +463,7 @@ export function createManagedBundledPluginWorkerRecovery(input: { export async function createApp( db: Db, opts: { + cloudWarmStandby?: CloudWarmStandby; uiMode: UiMode; serverPort: number; storageService: StorageService; @@ -505,7 +508,17 @@ export async function createApp( }, ) { const app = express(); + const staticUi = express.Router(); app.locals.paperclipDb = db; + const isWarmStandby = opts.cloudWarmStandby ?? (() => false); + const health = healthRoutes(db, { + deploymentMode: opts.deploymentMode, + deploymentExposure: opts.deploymentExposure, + authReady: opts.authReady, + companyDeletionEnabled: opts.companyDeletionEnabled, + databaseBackupHealth: opts.databaseBackupHealth, + isWarmStandby, + }); const captureRawBody = ( req: express.Request, _res: express.Response, @@ -558,6 +571,9 @@ export async function createApp( }), ); app.use(cloudRuntimeIdentityMiddleware(db)); + // A signed claim above commits identity before any normal request can seed + // company data. Unclaimed probes bypass session resolution as well as SQL. + app.use(cloudWarmStandbyMiddleware(isWarmStandby, health, staticUi)); // Connection-intent tools carry their own short-lived, run-bound bearer and // must be reachable by remote adapters that intentionally do not receive an // agent API key. Every request revalidates the active heartbeat row. @@ -595,7 +611,12 @@ export async function createApp( // Provider-authenticated ingress is intentionally outside the board // mutation guard. The Chat SDK adapter verifies the provider signature // before Paperclip persists or acts on any event. - const emailChannels = emailChannelService(db, { heartbeat: connectionIntentHeartbeat, storage: opts.storageService, publicBaseUrl: opts.chatWebhookPublicBaseUrl ?? opts.authPublicBaseUrl }); + const emailChannels = emailChannelService(db, { + isBackgroundWorkEnabled: () => !isWarmStandby(), + heartbeat: connectionIntentHeartbeat, + storage: opts.storageService, + publicBaseUrl: opts.chatWebhookPublicBaseUrl ?? opts.authPublicBaseUrl, + }); app.use(emailWebhookRoutes(emailChannels)); app.use(chatWebhookRoutes(chatChannels)); // The instance validates single-use registration state and its trusted @@ -642,16 +663,7 @@ export async function createApp( const agentAvatars = agentAvatarRoutes(); api.use(agentAvatars.router); api.use(boardMutationGuard()); - api.use( - "/health", - healthRoutes(db, { - deploymentMode: opts.deploymentMode, - deploymentExposure: opts.deploymentExposure, - authReady: opts.authReady, - companyDeletionEnabled: opts.companyDeletionEnabled, - databaseBackupHealth: opts.databaseBackupHealth, - }), - ); + api.use("/health", health); api.use(openApiRoutes()); api.use("/cloud", cloudRoutes()); api.use("/companies", companyRoutes(db, opts.storageService)); @@ -810,6 +822,7 @@ export async function createApp( const jobStore = pluginJobStore(db); const lifecycle = pluginLifecycleManager(db, { workerManager }); const scheduler = createPluginJobScheduler({ + isBackgroundWorkEnabled: () => !isWarmStandby(), db, jobStore, workerManager, @@ -982,7 +995,7 @@ export async function createApp( if (uiDist) { // Hashed asset files (Vite emits them under /assets/..) // never change once built, so they can be cached aggressively. - app.use( + staticUi.use( "/assets", express.static(path.join(uiDist, "assets"), { maxAge: "1y", @@ -990,14 +1003,14 @@ export async function createApp( }), ); // Serve root/index through the same runtime HTML transform as SPA routes. - app.get(["/", "/index.html"], (_req, res) => { + staticUi.get(["/", "/index.html"], (_req, res) => { res.type("html").set("Cache-Control", "no-cache").send(readBrandedStaticIndexHtml(uiDist)); }); // Non-hashed static files (favicon.ico, manifest, robots.txt, etc.): // short cache so operators who swap them out see the new version // reasonably fast, with must-revalidate overrides for index.html and // sw.js (see staticUiCacheControl for why those two). - app.use( + staticUi.use( express.static(uiDist, { maxAge: "1h", setHeaders(res, filePath) { @@ -1014,7 +1027,7 @@ export async function createApp( // with a MIME-type error, and cache that broken response. Return 404 // instead. The index.html response itself is no-cache so a subsequent // deploy's updated asset hashes are picked up on next load. - app.get(/.*/, (req, res) => { + staticUi.get(/.*/, (req, res) => { if (req.path.startsWith("/assets/")) { res.status(404).end(); return; @@ -1106,9 +1119,9 @@ export async function createApp( const renderViteHtml = viteHtmlRenderer; if (fs.existsSync(publicUiRoot)) { - app.use(express.static(publicUiRoot, { index: false })); + staticUi.use(express.static(publicUiRoot, { index: false })); } - app.get(/.*/, async (req, res, next) => { + staticUi.get(/.*/, async (req, res, next) => { if (!shouldServeViteDevHtml(req)) { next(); return; @@ -1120,9 +1133,10 @@ export async function createApp( next(err); } }); - app.use(vite.middlewares); + staticUi.use(vite.middlewares); } + app.use(staticUi); app.use(errorHandler); jobCoordinator.start(); @@ -1137,7 +1151,7 @@ export async function createApp( } }; const flushPendingFeedbackExports = async () => { - if (feedbackExportShuttingDown) return; + if (feedbackExportShuttingDown || isWarmStandby()) return; try { await opts.feedbackExportService?.flushPendingFeedbackTraces(); } catch (err) { @@ -1186,24 +1200,25 @@ export async function createApp( }); const unsubscribeChatPublicationSignals = subscribeAllCompanyLiveEvents( (event) => { - if (isChatPublicationCommitSignal(event)) + if (!isWarmStandby() && isChatPublicationCommitSignal(event)) chatReconciliation.notifyPublications(); }, ); let chatPublicationTimer: ReturnType | null = setInterval( () => { - chatReconciliation.reconcile(); + if (!isWarmStandby()) chatReconciliation.reconcile(); }, CHAT_PUBLICATION_FLUSH_INTERVAL_MS, ); chatPublicationTimer.unref?.(); - chatReconciliation.reconcile(); + if (!isWarmStandby()) chatReconciliation.reconcile(); // Abandoned chunked-import spool sweep: hourly (plus once at startup), // deleting spool dirs whose transfer saw no activity for 24h and cancelling // their still-open ledger runs. Same setInterval + unref + shutdown-clear // shape as the feedback export flush above. const importTransferSpoolRoot = resolveDefaultImportTransferSpoolRoot(); const sweepImportTransferSpools = () => { + if (isWarmStandby()) return; sweepAbandonedImportTransferSpools(db, importTransferSpoolRoot) .then((result) => { if (result.swept > 0) { @@ -1217,9 +1232,12 @@ export async function createApp( ); }); }; - const browserUseTimer = setInterval(() => { void browserUse.sweep().catch(() => logger.warn("Browser Use reconciliation failed; retrying.")); }, 3000); + const browserUseTimer = setInterval(() => { + if (isWarmStandby()) return; + void browserUse.sweep().catch(() => logger.warn("Browser Use reconciliation failed; retrying.")); + }, 3000); browserUseTimer.unref?.(); - void browserUse.sweep().catch(() => logger.warn("Browser Use startup reconciliation failed; retrying.")); + if (!isWarmStandby()) void browserUse.sweep().catch(() => logger.warn("Browser Use startup reconciliation failed; retrying.")); let importTransferSweepTimer: ReturnType | null = setInterval( sweepImportTransferSpools, diff --git a/server/src/index.ts b/server/src/index.ts index 92c57e4b95..91a19f0cbd 100644 --- a/server/src/index.ts +++ b/server/src/index.ts @@ -1,3 +1,5 @@ +import { cloudWarmStandbyServerOptions } from "./middleware/cloud-warm-standby.js"; +import { createCloudWarmStandby } from "./services/cloud-warm-standby.js"; import { subscribeAllCompanyLiveEvents } from "./services/live-events.js"; import { chatCompletionDeliveryService } from "./services/chat-completion-delivery.js"; /// @@ -17,7 +19,7 @@ import { reconcileAbandonedExecutionControl } from "./services/execution-control import { EXECUTION_RECONCILIATION_INTERVAL_MS } from "./services/execution-control-deadline.js"; import { connectionIntentDeliveryService } from "./services/connection-intent-delivery.js"; import { existsSync, readFileSync, rmSync } from "node:fs"; -import { createServer } from "node:http"; +import { createServer, type RequestListener } from "node:http"; import { resolve } from "node:path"; import { createInterface } from "node:readline/promises"; import { stdin, stdout } from "node:process"; @@ -652,6 +654,7 @@ async function startServerWithDatabaseTeardown( // Auth, routes, or child-runtime configuration capture any public URL. const restoredCloudRuntimeIdentity = await initializeCloudRuntimeIdentity(db as any); if (restoredCloudRuntimeIdentity) config = loadConfig(); + const isWarmStandby = await createCloudWarmStandby(db as any); if (config.deploymentMode === "local_trusted" && !isLoopbackHost(config.host)) { throw new Error( @@ -884,6 +887,7 @@ async function startServerWithDatabaseTeardown( // self-hosted: createApp falls back to its built-in kubernetes-only default. const managedPluginAutoInstall = managedConfig?.plugins.autoInstall ?? null; const app = await createApp(db as any, { + cloudWarmStandby: isWarmStandby, uiMode, serverPort: listenPort, storageService, @@ -922,7 +926,8 @@ async function startServerWithDatabaseTeardown( decisionServiceOptions, managedPluginAutoInstall, }); - const server = createServer(app as unknown as Parameters[0]); + // Upgrade admission runs before every WebSocket listener, outside Express. + const server = createServer(cloudWarmStandbyServerOptions(isWarmStandby), app as unknown as RequestListener); // Increase keep-alive timeouts to safely outlive default idle timeouts // of common reverse proxies and load balancers (like AWS ALB, Nginx, or Traefik). @@ -1165,7 +1170,7 @@ async function startServerWithDatabaseTeardown( ["local_ai_login_cleanup", () => localAiLoginService(db).reapExpired()], ] as const; const sweepExecutionControl = () => { - if (heartbeatSchedulerStopped) return; + if (heartbeatSchedulerStopped || isWarmStandby()) return; // Independent durable queues must not block one another. Each queue remains // single-flight; a later sweep observes committed transitions from its peers. for (const [queue, work] of executionControlSweeps) { @@ -1180,7 +1185,9 @@ async function startServerWithDatabaseTeardown( executionControlInterval.unref?.(); sweepExecutionControl(); const startHeartbeatSchedulerInterval = (callback: () => void) => { - heartbeatSchedulerInterval = setInterval(callback, config.heartbeatSchedulerIntervalMs); + heartbeatSchedulerInterval = setInterval(() => { + if (!isWarmStandby()) callback(); + }, config.heartbeatSchedulerIntervalMs); heartbeatSchedulerInterval?.unref?.(); }; const externalObjects = externalObjectService(db as any, { @@ -1853,6 +1860,7 @@ async function startServerWithDatabaseTeardown( "Automatic database backups enabled", ); setInterval(() => { + if (isWarmStandby()) return; void runServerDatabaseBackup("scheduled").catch(() => { // runServerDatabaseBackup already logs the failure with context. }); diff --git a/server/src/middleware/cloud-warm-standby.ts b/server/src/middleware/cloud-warm-standby.ts new file mode 100644 index 0000000000..73a25bc46c --- /dev/null +++ b/server/src/middleware/cloud-warm-standby.ts @@ -0,0 +1,41 @@ +import type { ServerOptions } from "node:http"; +import { Router, type RequestHandler } from "express"; +import type { CloudWarmStandby } from "../services/cloud-warm-standby.js"; + +/** Mount after signed identity handoff and before session or tenant resolution. */ +export function cloudWarmStandbyMiddleware( + isStandby: CloudWarmStandby, + health: RequestHandler, + staticUi: RequestHandler = Router(), +) { + const router = Router(); + router.use("/api/health", (req, res, next) => { + if (isStandby() && !req.headers.upgrade && (req.method === "GET" || req.method === "HEAD") && req.path === "/") { + req.actor = { type: "none", source: "none" }; + health(req, res, next); + return; + } + next(); + }); + router.use((req, res, next) => { + if (!isStandby()) return next(); + req.actor = { type: "none", source: "none" }; + const unavailable = () => res.status(503).json({ error: "workspace_unclaimed" }); + if (/^\/api(?:\/|$)/i.test(req.path) || req.headers.upgrade || (req.method !== "GET" && req.method !== "HEAD")) { + unavailable(); + return; + } + // Serve only the UI router, never fall through to session, tenant, bearer, + // MCP, or other dynamic handlers, even when the request carries credentials. + staticUi(req, res, (error?: unknown) => { + if (error) next(error); + else unavailable(); + }); + }); + return router; +} + +/** Node 24.9+ routes rejected upgrades through the guarded HTTP request path. */ +export function cloudWarmStandbyServerOptions(isStandby: CloudWarmStandby): ServerOptions { + return { shouldUpgradeCallback: () => !isStandby() }; +} diff --git a/server/src/routes/health.ts b/server/src/routes/health.ts index bf7dc4a4a0..b641a98ac8 100644 --- a/server/src/routes/health.ts +++ b/server/src/routes/health.ts @@ -127,6 +127,7 @@ export function healthRoutes( serverInfo?: ServerInfoSnapshot; databaseBackupHealth?: InspectDatabaseBackupHealthOptions; runtimeEnv?: CloudInstanceEnv; + isWarmStandby?: () => boolean; } = { deploymentMode: "local_trusted", deploymentExposure: "private", @@ -293,6 +294,23 @@ export function healthRoutes( return; } + // Startup has already validated the empty database. During warm standby, + // readiness means the HTTP app can accept a signed claim; probing SQL would + // keep the idle database awake. The claim itself still performs durable SQL, + // and claimed instances retain the normal live database check below. + if (opts.isWarmStandby?.()) { + res.json({ + status: healthStatus, + deploymentMode: opts.deploymentMode, + deploymentExposure: opts.deploymentExposure, + commit, + warmStandby: true, + ...(cloud ? { cloud } : {}), + ...(hiddenSettings.length ? { hiddenSettings } : {}), + }); + return; + } + try { await db.execute(sql`SELECT 1`); } catch (error) { diff --git a/server/src/services/cloud-warm-standby.test.ts b/server/src/services/cloud-warm-standby.test.ts new file mode 100644 index 0000000000..cd51cc419f --- /dev/null +++ b/server/src/services/cloud-warm-standby.test.ts @@ -0,0 +1,66 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; +import type { Db } from "@paperclipai/db"; +import * as identity from "./cloud-runtime-identity.js"; +import { createCloudWarmStandby } from "./cloud-warm-standby.js"; + +const env = { + PAPERCLIP_CLOUD_WARM_STANDBY: "1", + PAPERCLIP_CLOUD_TENANT_SERVER_TOKEN: "test-token", + PAPERCLIP_CLOUD_STACK_ID: "stack-test", + PAPERCLIP_CLOUD_RUNTIME_IDENTITY_JWKS: "test-keys", +}; +function database(rows: unknown[] = []) { + const limit = vi.fn().mockResolvedValue(rows); + const select = vi.fn(() => ({ from: () => ({ limit }) })); + return { db: { select } as unknown as Db, select, limit }; +} + +afterEach(() => vi.restoreAllMocks()); + +describe("Cloud warm standby", () => { + it.each([ + {}, + { ...env, PAPERCLIP_CLOUD_WARM_STANDBY: "0" }, + { ...env, PAPERCLIP_CLOUD_TENANT_SERVER_TOKEN: undefined }, + { ...env, PAPERCLIP_CLOUD_STACK_ID: undefined }, + { ...env, PAPERCLIP_CLOUD_RUNTIME_IDENTITY_JWKS: undefined }, + ])("leaves ordinary or incomplete configurations active without querying", async (config) => { + const { db, select } = database(); + expect((await createCloudWarmStandby(db, config))()).toBe(false); + expect(select).not.toHaveBeenCalled(); + }); + + it("protects existing company data even when the marker is set", async () => { + const { db } = database([{ id: "company-test" }]); + expect((await createCloudWarmStandby(db, env))()).toBe(false); + }); + + it("never mistakes a failed empty-database check for standby", async () => { + const { db, limit } = database(); + limit.mockRejectedValue(new Error("database unavailable")); + await expect(createCloudWarmStandby(db, env)).rejects.toThrow("database unavailable"); + }); + + it("checks emptiness once and exits standby monotonically on the matching committed claim", async () => { + const getter = vi.spyOn(identity, "getCloudRuntimeIdentity").mockReturnValue(null); + const { db, select } = database(); + const isStandby = await createCloudWarmStandby(db, env); + for (let i = 0; i < 100; i++) expect(isStandby()).toBe(true); + expect(select).toHaveBeenCalledOnce(); + const claimed = { stackId: "stack-other" } as identity.CloudRuntimeIdentitySnapshot; + getter.mockReturnValue(claimed); + expect(isStandby()).toBe(true); + getter.mockReturnValue({ ...claimed, stackId: env.PAPERCLIP_CLOUD_STACK_ID }); + expect(isStandby()).toBe(false); + getter.mockReturnValue(null); + expect(isStandby()).toBe(false); + expect(select).toHaveBeenCalledOnce(); + }); + + it("starts active after a persisted identity is restored, even with no companies", async () => { + vi.spyOn(identity, "getCloudRuntimeIdentity").mockReturnValue({ stackId: env.PAPERCLIP_CLOUD_STACK_ID } as identity.CloudRuntimeIdentitySnapshot); + const { db, select } = database(); + expect((await createCloudWarmStandby(db, env))()).toBe(false); + expect(select).not.toHaveBeenCalled(); + }); +}); diff --git a/server/src/services/cloud-warm-standby.ts b/server/src/services/cloud-warm-standby.ts new file mode 100644 index 0000000000..8770e1ec55 --- /dev/null +++ b/server/src/services/cloud-warm-standby.ts @@ -0,0 +1,37 @@ +import { companies, type Db } from "@paperclipai/db"; +import { isCloudManagedInstance } from "./cloud-instance.js"; +import { getCloudRuntimeIdentity } from "./cloud-runtime-identity.js"; + +/** In-memory predicate. Calling it never opens a database connection. */ +export type CloudWarmStandby = () => boolean; + +/** + * Call after initializeCloudRuntimeIdentity and before starting any pollers. + * Only explicitly marked, empty, unclaimed Cloud instances may stand by. + * A signed claim is committed before the identity getter changes, and restores + * on boot even if the provider still carries the old warm-pool environment. + */ +export async function createCloudWarmStandby( + db: Db, + env: NodeJS.ProcessEnv = process.env, +): Promise { + const stackId = env.PAPERCLIP_CLOUD_STACK_ID?.trim(); + if ( + env.PAPERCLIP_CLOUD_WARM_STANDBY !== "1" + || !isCloudManagedInstance(env) + || !stackId + || !env.PAPERCLIP_CLOUD_RUNTIME_IDENTITY_JWKS?.trim() + || getCloudRuntimeIdentity() + ) return () => false; + + // Protect existing tenants if an operator accidentally sets the marker. + // Fail startup on a read error; never infer emptiness from an unavailable DB. + const existing = await db.select({ id: companies.id }).from(companies).limit(1); + if (existing.length > 0) return () => false; + + let standby = true; + return () => { + if (standby && getCloudRuntimeIdentity()?.stackId === stackId) standby = false; + return standby; + }; +} diff --git a/server/src/services/email-channels.ts b/server/src/services/email-channels.ts index d7d8b08d4c..de3e9eb300 100644 --- a/server/src/services/email-channels.ts +++ b/server/src/services/email-channels.ts @@ -74,6 +74,8 @@ export type EmailActor = { type Endpoint = typeof chatEndpoints.$inferSelect; type Tx = Parameters[0]>[0]; export interface EmailChannelOptions { + /** Suppress periodic database work while an unclaimed Cloud app stands by. */ + isBackgroundWorkEnabled?: () => boolean; heartbeat: Pick, "wakeup">; storage?: StorageService; publicBaseUrl?: string; @@ -2041,6 +2043,7 @@ export function emailChannelService(db: Db, options: EmailChannelOptions) { await Promise.all(items.slice(i, i + 4).map(work)); } async function tick() { + if (options.isBackgroundWorkEnabled?.() === false) return; if (activeTick) return activeTick; activeTick = runTick(); try { diff --git a/server/src/services/plugin-job-scheduler.ts b/server/src/services/plugin-job-scheduler.ts index 09a6b8782b..a7044f2e8d 100644 --- a/server/src/services/plugin-job-scheduler.ts +++ b/server/src/services/plugin-job-scheduler.ts @@ -63,6 +63,8 @@ const DEFAULT_MAX_CONCURRENT_JOBS = 10; * Options for creating a PluginJobScheduler. */ export interface PluginJobSchedulerOptions { + /** Suppress periodic database work while an unclaimed Cloud app stands by. */ + isBackgroundWorkEnabled?: () => boolean; /** Drizzle database instance. */ db: Db; /** Persistence layer for jobs and runs. */ @@ -244,6 +246,7 @@ export function createPluginJobScheduler( * A single scheduler tick. Queries for due jobs and dispatches them. */ async function tick(): Promise { + if (options.isBackgroundWorkEnabled?.() === false) return; // Prevent overlapping ticks (in case a tick takes longer than the interval) if (tickInProgress) { log.debug("skipping tick — previous tick still in progress"); diff --git a/ui/src/api/health.ts b/ui/src/api/health.ts index 2b6a2a08aa..2612f019e0 100644 --- a/ui/src/api/health.ts +++ b/ui/src/api/health.ts @@ -36,6 +36,8 @@ export type HealthStatus = { authReady?: boolean; bootstrapStatus?: "ready" | "bootstrap_pending"; bootstrapInviteActive?: boolean; + /** The unclaimed Cloud app is ready; this response did not probe SQL. */ + warmStandby?: boolean; features?: { companyDeletionEnabled?: boolean; };