mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-06 21:05:21 +02:00
fix(server): leave unclaimed warm Cloud databases idle (#15314)
## Thinking Path > - Paperclip manages work performed by AI agents. > - Managed deployments prepare empty applications before an owner claims them. > - Those applications start database pollers even though no company can have work. > - Health probes also query SQL, so idle databases cannot remain suspended. > - This pull request adds an explicit standby marker for empty, unclaimed Cloud apps. > - The existing signed, durable claim resumes normal processing without restarting the app. ## Linked Issues or Issue Description **What happened?** An empty, unclaimed warm application runs recurring chat, email, plugin, heartbeat, cleanup, and reconciliation queries. Its health route also opens the database. This prevents idle database compute from suspending. **Expected behavior** An explicitly marked unclaimed application should keep its HTTP process and sandbox provider plugins ready while leaving the database idle. A successful signed claim should resume normal API behavior and background processing. Claimed and self-hosted instances should keep their current behavior. **Steps to reproduce** Start an empty Cloud-managed application and leave it unclaimed. Observe database activity while repeatedly requesting `/api/health`. Before this change, periodic queries continue without company data. Related: #15153 reduces allocation during chat polling. This change suppresses polling only for explicitly marked, empty, unclaimed Cloud apps. ## What Changed - Add `PAPERCLIP_CLOUD_WARM_STANDBY=1`. Check company emptiness once after restoring the persisted Cloud runtime identity. Missing Cloud configuration, existing data, or a persisted claim leaves normal processing active. - Gate recurring database pollers with an in-memory predicate. Keep startup preparation and sandbox provider plugin loading intact. - Serve unclaimed health probes without session or database reads and report `warmStandby: true`. Serve standby pages/assets directly from the UI router, bypassing session, bearer, tenant, and dynamic handlers. Refuse API requests and all WebSocket upgrades before authentication can query SQL or seed company data. - Exit standby after the existing signed identity assertion commits. Normal timers resume at their next tick; a restart restores the claim even with stale provider variables. - Document the marker, readiness semantics, rollout checks, and rollback. ## Verification - `pnpm -r typecheck` passed. A final server typecheck also passed after adding tests. - `pnpm build` passed. - Focused standby, signed claim, restart, health, static/Vite routing, hostname, HMR, and live-events suites: 72 passed after the review fixes. Includes real HTTP upgrade admission before/after claim. - `pnpm test:run` was attempted locally; both superseded runs were stopped after encountering checkout/platform failures. A clean-checkout rerun eliminated ancestor skill-directory lookup failures. The company-skills/runtime-cache families encounter macOS read-only-directory rename failures (`EACCES`); all three company-skills failures reproduce on unmodified base `bf14f803d5`. The initial full run also reported one native runner API test failure; an isolated comparison on both revisions was blocked by local embedded PostgreSQL startup failures. The full [Linux CI run](https://github.com/paperclipai/paperclip/actions/runs/37423263993) passed on final commit `5016c415ea`, including all server and workspace test shards, browser suites, typecheck, build, and release canary. This is not a claim that the full local suite passed. - Isolated full server with local PostgreSQL: after startup and connection expiry, 70 health probes, 70 page requests carrying valid synthetic tenant credentials, and 70 rejected WebSocket upgrades over 70 seconds observed zero app database connections. The signed claim completed in 62 ms and normal polling resumed (475 database transactions over 12 seconds). Restart with stale provider variables restored the durable claim. The latency is local-only, not a provider wake measurement. - Apex review: **5/5** on `5016c415ea`, both earlier threads resolved, no open recommendations. - No live-provider test or production deployment was performed. An actual database suspension/resume canary remains required before enabling the control-plane switch. ## Risks - Standby health reports HTTP readiness rather than current database connectivity. The signed claim still requires a durable database write; claimed health checks retain the SQL probe and 503 failure behavior. - Pollers resume at their usual intervals. A suspended database may add claim latency. Validate the real provider before enabling the marker. - Startup preparation and sandbox plugins remain loaded. New plugins or background loops must respect the same standby contract. - The marker is off by default. Remove it or set it to `0` and restart to roll back. No schema migration or claimed-workspace inactivity policy changes. ## Model Used OpenAI Codex, based on GPT-6. The exact serving snapshot and configured context-window size are not exposed in this session. Assistance included source review, TypeScript changes, command execution, and PostgreSQL tests. ## Checklist - [x] I have included a thinking path that traces from project context to this change - [x] I have specified the model used (with version and capability details) - [x] I have checked ROADMAP.md and confirmed this PR does not duplicate planned core work - [x] I have searched GitHub for duplicate or related PRs and linked them above - [x] I have either (a) linked existing issues with `Fixes: #` / `Closes #` OR (b) described the issue in-PR following the relevant issue template - [x] I have not referenced internal/instance-local Paperclip issues or links (only public GitHub `#NNN` / `github.com/paperclipai/paperclip` URLs) - [x] My branch name describes the change (e.g. `docs/...`, `fix/...`) and contains no internal Paperclip ticket id or instance-derived details - [ ] I have run tests locally and they pass — targeted tests pass; full local suite limitations are documented above, and full Linux CI is green - [x] I have added or updated tests where applicable - [x] I have updated relevant documentation to reflect my changes - [x] I have considered and documented any risks above - [x] All Paperclip CI gates are green - [x] Greptile is 5/5 with no open P2s, recommendations, or follow-ups - [x] I will address all Greptile and reviewer comments before requesting merge --------- Co-authored-by: Paperclip <noreply@paperclip.ing>
This commit is contained in:
1 parent
f2715e02bb
commit
202c2d307e
12 files changed
+420
-29
No files matched your search
@@ -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:
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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<void>((resolve) => server.listen(0, "127.0.0.1", resolve));
|
||||
const address = server.address() as { port: number };
|
||||
const upgrade = (path: string) => new Promise<number>((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<void>((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();
|
||||
}
|
||||
});
|
||||
});
|
||||
+42
-24
@@ -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/<name>.<hash>.<ext>)
|
||||
// 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<typeof setInterval> | 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<typeof setInterval> | null =
|
||||
setInterval(
|
||||
sweepImportTransferSpools,
|
||||
|
||||
+12
-4
@@ -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";
|
||||
/// <reference path="./types/express.d.ts" />
|
||||
@@ -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<typeof createServer>[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.
|
||||
});
|
||||
|
||||
@@ -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() };
|
||||
}
|
||||
@@ -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) {
|
||||
|
||||
@@ -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();
|
||||
});
|
||||
});
|
||||
@@ -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<CloudWarmStandby> {
|
||||
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;
|
||||
};
|
||||
}
|
||||
@@ -74,6 +74,8 @@ export type EmailActor = {
|
||||
type Endpoint = typeof chatEndpoints.$inferSelect;
|
||||
type Tx = Parameters<Parameters<Db["transaction"]>[0]>[0];
|
||||
export interface EmailChannelOptions {
|
||||
/** Suppress periodic database work while an unclaimed Cloud app stands by. */
|
||||
isBackgroundWorkEnabled?: () => boolean;
|
||||
heartbeat: Pick<ReturnType<typeof heartbeatService>, "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 {
|
||||
|
||||
@@ -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<void> {
|
||||
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");
|
||||
|
||||
@@ -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;
|
||||
};
|
||||
|
||||
Reference in new issue
Block a user