-
Notifications
You must be signed in to change notification settings - Fork 180
Expand file tree
/
Copy pathindex.js
More file actions
210 lines (183 loc) · 6.31 KB
/
Copy pathindex.js
File metadata and controls
210 lines (183 loc) · 6.31 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
import express from "express";
import cors from "cors";
import config, { validateConfig } from "./config.js";
import logger from "./lib/logger.js";
import {
requestLogger,
requestContextMiddleware,
} from "./middleware/requestContext.js";
import { checkReadiness } from "./lib/readiness.js";
import {
getSubmitQueueDepth,
drainSubmitQueue,
getPendingTransactionCount,
getPendingTransactions,
dumpPendingTransactions,
resumePendingTransactions,
} from "./lib/contract.js";
import requestIdMiddleware from "./lib/requestId.js";
import registryRouter from "./routes/registry.js";
import servicesRouter from "./routes/services.js";
import demoRouter from "./routes/demo.js";
import agentsRouter from "./routes/agents.js";
if (process.argv.includes("--print-config")) {
console.log(
JSON.stringify(
{
nodeEnv: config.nodeEnv,
port: config.port,
logLevel: config.logLevel,
stellar: config.stellar,
contract: config.contract,
x402: {
facilitatorUrl: config.x402.facilitatorUrl,
searchPrice: config.x402.searchPrice,
weatherPrice: config.x402.weatherPrice,
payTo: config.x402.payTo,
},
corsOrigin: config.corsOrigin,
jsonBodyLimit: config.jsonBodyLimit,
trustProxy: config.trustProxy,
rateLimit: config.rateLimit,
demoRun: config.demoRun,
},
null,
2,
),
);
process.exit(0);
}
validateConfig(logger);
logger.info({ corsOrigin: config.corsOrigin }, "Resolved CORS origin allowlist");
logger.info(
{ demoRun: config.demoRun },
"Effective activity-poll configuration (override with DEMO_RUN_POLL_* env vars)",
);
const app = express();
// Trust the configured number of proxy hops so req.ip reflects the real client
// (via X-Forwarded-For) behind a reverse proxy — required for correct IP-based
// rate limiting. Defaults to false (no proxy) to avoid X-Forwarded-For spoofing.
app.set("trust proxy", config.trustProxy);
app.use(requestLogger);
app.use(requestContextMiddleware);
app.use(cors({ origin: config.corsOrigin, credentials: true }));
app.use(requestIdMiddleware);
app.use(express.json({ limit: config.jsonBodyLimit }));
// Liveness (#841): dependency-free by design. Answers only "is this process
// running?" so an orchestrator never restarts a healthy instance because an
// upstream is slow. Dependency state lives on /readyz.
app.get("/healthz", (req, res) => {
res.status(200).json({
status: "ok",
uptimeSeconds: Math.floor(process.uptime()),
queueDepth: getSubmitQueueDepth(),
pendingTransactions: getPendingTransactionCount(),
timestamp: new Date().toISOString(),
});
});
// Readiness (#841): can this instance serve traffic right now? Checks RPC and
// Redis under a short timeout and returns 503 when a required dependency is
// unreachable, so traffic stops without the process being restarted.
app.get("/readyz", async (req, res) => {
try {
const result = await checkReadiness();
res.status(result.ready ? 200 : 503).json(result);
} catch (err) {
logger.error({ err }, "GET /readyz failed");
res.status(503).json({
ready: false,
status: "not_ready",
error: "Readiness check failed",
timestamp: new Date().toISOString(),
});
}
});
app.use("/api", registryRouter);
app.use("/api", agentsRouter);
// Demo routes are backed by server-custodied keys and should not be reachable
// in production deployments. Gate them behind ENABLE_DEMO_ROUTES, defaulting to
// enabled only when NODE_ENV is not "production".
const enableDemoRoutes =
process.env.ENABLE_DEMO_ROUTES === 'true' ||
(process.env.ENABLE_DEMO_ROUTES === undefined && config.nodeEnv !== 'production');
if (enableDemoRoutes) {
logger.info({ nodeEnv: config.nodeEnv }, 'Demo routes enabled');
app.use("/api", demoRouter);
app.use("/demo", servicesRouter);
} else {
logger.info({ nodeEnv: config.nodeEnv }, 'Demo routes disabled (set ENABLE_DEMO_ROUTES=true to enable)');
}
app.use((err, req, res, _next) => {
if (err.type === "entity.too.large") {
req.log.warn({ expected: config.jsonBodyLimit }, "Request body too large");
return res.status(413).json({
error: `Request body too large. Maximum size is ${config.jsonBodyLimit}.`,
code: "PAYLOAD_TOO_LARGE",
});
}
res.status(500).json({
error: "Internal server error",
code: "INTERNAL_ERROR",
requestId: _req.requestId,
});
});
let server;
let shuttingDown = false;
async function shutdown() {
if (shuttingDown) return;
shuttingDown = true;
logger.info("Shutting down gracefully...");
// Force-exit after the configured timeout so the process never hangs forever
const forceExitTimer = setTimeout(() => {
const pending = getPendingTransactions();
if (pending.length > 0) {
logger.warn(
{ count: pending.length, timeout: config.shutdownTimeoutMs },
"Shutdown timeout reached — dumping pending transactions and force-exiting",
);
dumpPendingTransactions();
}
process.exit(1);
}, config.shutdownTimeoutMs);
forceExitTimer.unref();
// If server hasn't been created yet (signal during startup resume) skip close
if (!server) {
logger.warn("Server was not yet listening — draining queue directly");
await doDrainAndDump();
clearTimeout(forceExitTimer);
process.exit(0);
return;
}
// Stop accepting new connections
server.close(async (closeErr) => {
if (closeErr) {
logger.error({ err: closeErr }, "Error closing HTTP server");
} else {
logger.info("HTTP server closed — no longer accepting new connections");
}
await doDrainAndDump();
clearTimeout(forceExitTimer);
process.exit(0);
});
}
async function doDrainAndDump() {
try {
await drainSubmitQueue();
logger.info("Submit queue drained successfully");
} catch (err) {
logger.error({ err }, "Error draining submit queue");
}
const pending = getPendingTransactions();
if (pending.length > 0) {
logger.warn(
{ count: pending.length, hashes: pending.map((t) => t.hash) },
"Pending transactions remain after queue drain — dumped to pending-transactions.json for manual verification",
);
dumpPendingTransactions();
} else {
logger.info("No pending transactions — clean shutdown");
}
}
process.on("SIGTERM", shutdown);
process.on("SIGINT", shutdown);
start();