feat(ai-workspace): add local relay profiles
This commit is contained in:
@@ -8,12 +8,17 @@ const MAX_EVENTS_PER_AGENT = 5000;
|
||||
const MAX_REQUEST_META_PER_AGENT = 1000;
|
||||
const HEARTBEAT_INTERVAL_MS = 30000;
|
||||
const AGENT_STALE_MS = 75000;
|
||||
const ASSISTANT_RELAY_POLL_TIMEOUT_MS = Number(process.env.AI_WORKSPACE_ASSISTANT_RELAY_POLL_TIMEOUT_MS || 25000);
|
||||
const ASSISTANT_RELAY_ACTION_TIMEOUT_MS = Number(process.env.AI_WORKSPACE_ASSISTANT_RELAY_ACTION_TIMEOUT_MS || 45000);
|
||||
const ASSISTANT_RELAY_MAX_PENDING = Number(process.env.AI_WORKSPACE_ASSISTANT_RELAY_MAX_PENDING || 200);
|
||||
const ASSISTANT_RELAY_MAX_BATCH = Number(process.env.AI_WORKSPACE_ASSISTANT_RELAY_MAX_BATCH || 10);
|
||||
|
||||
const config = readConfig();
|
||||
const app = express();
|
||||
const httpServer = createServer(app);
|
||||
const wss = new WebSocketServer({ noServer: true });
|
||||
const agentsByCode = new Map();
|
||||
const assistantRelaysById = new Map();
|
||||
let eventSeq = 0;
|
||||
|
||||
app.disable("x-powered-by");
|
||||
@@ -25,6 +30,7 @@ app.get("/healthz", (_req, res) => {
|
||||
service: "nodedc-ai-workspace-hub",
|
||||
wsPath: config.wsPath,
|
||||
agentsOnline: Array.from(agentsByCode.values()).filter(isAgentOnline).length,
|
||||
assistantRelays: assistantRelaysById.size,
|
||||
internalApiConfigured: config.internalAccessTokens.length > 0,
|
||||
});
|
||||
});
|
||||
@@ -60,6 +66,12 @@ app.post("/api/ai-workspace/hub/v1/agents/:pairingCode/dispatch", requireInterna
|
||||
res.json({ ok: true, requestId: dispatch.requestId });
|
||||
});
|
||||
|
||||
app.post("/api/ai-workspace/assistant/v1/actions", requireInternalApi, asyncRoute(proxyAssistantActions));
|
||||
|
||||
app.post("/api/ai-workspace/hub/v1/assistant-relays/:relayId/actions", requireInternalApi, asyncRoute(callAssistantRelayAction));
|
||||
app.post("/api/ai-workspace/hub/v1/assistant-relays/:relayId/poll", requireInternalApi, asyncRoute(pollAssistantRelay));
|
||||
app.post("/api/ai-workspace/hub/v1/assistant-relays/:relayId/results/:callId", requireInternalApi, asyncRoute(completeAssistantRelayCall));
|
||||
|
||||
app.use("/api/ai-workspace/hub/v1/ndc-agent-mcp/:pairingCode", asyncRoute(proxyNdcAgentMcp));
|
||||
|
||||
app.use((error, _req, res, _next) => {
|
||||
@@ -256,6 +268,179 @@ async function proxyNdcAgentMcp(req, res) {
|
||||
res.send(text);
|
||||
}
|
||||
|
||||
async function proxyAssistantActions(req, res) {
|
||||
if (!config.assistantInternalUrl || !config.assistantInternalAccessToken) {
|
||||
res.status(503).json({ ok: false, error: "assistant_action_proxy_not_configured" });
|
||||
return;
|
||||
}
|
||||
const targetUrl = `${config.assistantInternalUrl.replace(/\/+$/, "")}/api/ai-workspace/assistant/v1/actions`;
|
||||
const upstream = await fetch(targetUrl, {
|
||||
method: "POST",
|
||||
redirect: "manual",
|
||||
headers: {
|
||||
Accept: "application/json",
|
||||
"Content-Type": "application/json",
|
||||
Authorization: `Bearer ${config.assistantInternalAccessToken}`,
|
||||
...forwardOwnerHeaders(req),
|
||||
},
|
||||
body: JSON.stringify(req.body || {}),
|
||||
});
|
||||
const contentType = upstream.headers.get("content-type") || "application/json; charset=utf-8";
|
||||
const text = await upstream.text();
|
||||
res.status(upstream.status);
|
||||
res.setHeader("content-type", contentType);
|
||||
res.send(text);
|
||||
}
|
||||
|
||||
async function callAssistantRelayAction(req, res) {
|
||||
const relayId = cleanRelayId(req.params.relayId);
|
||||
if (!relayId) {
|
||||
res.status(400).json({ ok: false, error: "assistant_relay_id_required" });
|
||||
return;
|
||||
}
|
||||
|
||||
const relay = getAssistantRelay(relayId);
|
||||
pruneAssistantRelay(relay);
|
||||
if (relay.calls.size >= ASSISTANT_RELAY_MAX_PENDING) {
|
||||
res.status(429).json({ ok: false, error: "assistant_relay_backpressure" });
|
||||
return;
|
||||
}
|
||||
|
||||
const callId = crypto.randomUUID();
|
||||
const timeoutMs = sanitizeTimeoutMs(req.body?.timeoutMs || req.query?.timeoutMs || ASSISTANT_RELAY_ACTION_TIMEOUT_MS, ASSISTANT_RELAY_ACTION_TIMEOUT_MS);
|
||||
const createdAt = new Date().toISOString();
|
||||
const call = {
|
||||
callId,
|
||||
relayId,
|
||||
createdAt,
|
||||
payload: req.body || {},
|
||||
headers: forwardOwnerHeaders(req),
|
||||
};
|
||||
|
||||
const result = await new Promise((resolve, reject) => {
|
||||
const timeout = setTimeout(() => {
|
||||
relay.calls.delete(callId);
|
||||
relay.pending = relay.pending.filter((item) => item.callId !== callId);
|
||||
const error = new Error("assistant_relay_timeout");
|
||||
error.status = 504;
|
||||
reject(error);
|
||||
}, timeoutMs);
|
||||
timeout.unref?.();
|
||||
relay.calls.set(callId, { resolve, reject, timeout, createdAt: Date.now() });
|
||||
relay.pending.push(call);
|
||||
notifyAssistantRelayWaiters(relay);
|
||||
});
|
||||
|
||||
res.status(result.status).json(result.body);
|
||||
}
|
||||
|
||||
async function pollAssistantRelay(req, res) {
|
||||
const relayId = cleanRelayId(req.params.relayId);
|
||||
if (!relayId) {
|
||||
res.status(400).json({ ok: false, error: "assistant_relay_id_required" });
|
||||
return;
|
||||
}
|
||||
const relay = getAssistantRelay(relayId);
|
||||
relay.lastSeenAt = new Date().toISOString();
|
||||
const limit = sanitizeInteger(req.body?.limit, ASSISTANT_RELAY_MAX_BATCH, 1, ASSISTANT_RELAY_MAX_BATCH);
|
||||
const timeoutMs = sanitizeTimeoutMs(req.body?.timeoutMs || req.query?.timeoutMs || ASSISTANT_RELAY_POLL_TIMEOUT_MS, ASSISTANT_RELAY_POLL_TIMEOUT_MS);
|
||||
const calls = await waitForAssistantRelayCalls(relay, limit, timeoutMs);
|
||||
res.json({ ok: true, relayId, calls });
|
||||
}
|
||||
|
||||
async function completeAssistantRelayCall(req, res) {
|
||||
const relayId = cleanRelayId(req.params.relayId);
|
||||
const callId = cleanString(req.params.callId, 120);
|
||||
const relay = assistantRelaysById.get(relayId);
|
||||
const pending = relay?.calls.get(callId);
|
||||
if (!relay || !pending) {
|
||||
res.status(404).json({ ok: false, error: "assistant_relay_call_not_found" });
|
||||
return;
|
||||
}
|
||||
|
||||
relay.calls.delete(callId);
|
||||
clearTimeout(pending.timeout);
|
||||
pending.resolve({
|
||||
status: sanitizeHttpStatus(req.body?.status),
|
||||
body: req.body?.body === undefined ? { ok: true } : req.body.body,
|
||||
});
|
||||
res.json({ ok: true, relayId, callId });
|
||||
}
|
||||
|
||||
function getAssistantRelay(relayId) {
|
||||
const cleanId = cleanRelayId(relayId);
|
||||
let relay = assistantRelaysById.get(cleanId);
|
||||
if (!relay) {
|
||||
relay = {
|
||||
relayId: cleanId,
|
||||
pending: [],
|
||||
waiters: [],
|
||||
calls: new Map(),
|
||||
createdAt: new Date().toISOString(),
|
||||
lastSeenAt: "",
|
||||
};
|
||||
assistantRelaysById.set(cleanId, relay);
|
||||
}
|
||||
return relay;
|
||||
}
|
||||
|
||||
function waitForAssistantRelayCalls(relay, limit, timeoutMs) {
|
||||
const ready = takeAssistantRelayCalls(relay, limit);
|
||||
if (ready.length) return Promise.resolve(ready);
|
||||
return new Promise((resolve) => {
|
||||
const timeout = setTimeout(() => {
|
||||
relay.waiters = relay.waiters.filter((waiter) => waiter.resolve !== resolve);
|
||||
resolve([]);
|
||||
}, timeoutMs);
|
||||
timeout.unref?.();
|
||||
relay.waiters.push({ resolve, timeout, limit });
|
||||
});
|
||||
}
|
||||
|
||||
function notifyAssistantRelayWaiters(relay) {
|
||||
while (relay.waiters.length && relay.pending.length) {
|
||||
const waiter = relay.waiters.shift();
|
||||
clearTimeout(waiter.timeout);
|
||||
waiter.resolve(takeAssistantRelayCalls(relay, waiter.limit));
|
||||
}
|
||||
}
|
||||
|
||||
function takeAssistantRelayCalls(relay, limit) {
|
||||
const calls = [];
|
||||
while (relay.pending.length && calls.length < limit) {
|
||||
const call = relay.pending.shift();
|
||||
if (relay.calls.has(call.callId)) calls.push(call);
|
||||
}
|
||||
return calls;
|
||||
}
|
||||
|
||||
function pruneAssistantRelay(relay) {
|
||||
relay.pending = relay.pending.filter((call) => relay.calls.has(call.callId));
|
||||
}
|
||||
|
||||
function cleanRelayId(value) {
|
||||
return String(value || "").trim().replace(/[^A-Za-z0-9_.:-]/g, "").slice(0, 120);
|
||||
}
|
||||
|
||||
function sanitizeHttpStatus(value) {
|
||||
const status = Number(value || 200);
|
||||
return Number.isInteger(status) && status >= 100 && status <= 599 ? status : 200;
|
||||
}
|
||||
|
||||
function forwardOwnerHeaders(req) {
|
||||
const headers = {};
|
||||
for (const name of [
|
||||
"x-nodedc-user-id",
|
||||
"x-nodedc-user-email",
|
||||
"x-nodedc-user-role",
|
||||
"x-nodedc-user-groups",
|
||||
]) {
|
||||
const value = req.headers[name];
|
||||
if (typeof value === "string" && value.trim()) headers[name] = value.trim();
|
||||
}
|
||||
return headers;
|
||||
}
|
||||
|
||||
function startAgentRequest(pairingCodeRaw, command, payload = {}, timeoutMs = 30000, options = {}) {
|
||||
const pairingCode = cleanPairingCode(pairingCodeRaw);
|
||||
const agent = agentsByCode.get(pairingCode);
|
||||
@@ -432,6 +617,12 @@ function sanitizeTimeoutMs(value, fallback) {
|
||||
return Math.max(1000, Math.min(12 * 60 * 60 * 1000, numeric));
|
||||
}
|
||||
|
||||
function sanitizeInteger(value, fallback, min, max) {
|
||||
const numeric = Number(value || fallback);
|
||||
if (!Number.isFinite(numeric)) return fallback;
|
||||
return Math.min(Math.max(Math.trunc(numeric), min), max);
|
||||
}
|
||||
|
||||
function elapsedLabel(startedAt) {
|
||||
const elapsed = Math.max(0, Date.now() - Number(startedAt || Date.now()));
|
||||
if (elapsed < 1000) return `${elapsed}ms`;
|
||||
@@ -473,6 +664,14 @@ function readConfig() {
|
||||
internalAccessTokens,
|
||||
engineInternalUrl: cleanString(process.env.NODEDC_ENGINE_INTERNAL_URL, 1000).replace(/\/+$/, ""),
|
||||
engineInternalAccessToken: cleanString(process.env.NODEDC_INTERNAL_ACCESS_TOKEN, 1000),
|
||||
assistantInternalUrl: cleanString(
|
||||
process.env.NODEDC_AI_WORKSPACE_ASSISTANT_URL ||
|
||||
process.env.NDC_AI_WORKSPACE_ASSISTANT_INTERNAL_URL ||
|
||||
process.env.AI_WORKSPACE_ASSISTANT_INTERNAL_URL ||
|
||||
"http://ai-workspace-assistant:18082",
|
||||
1000,
|
||||
).replace(/\/+$/, ""),
|
||||
assistantInternalAccessToken: cleanString(process.env.NODEDC_INTERNAL_ACCESS_TOKEN, 1000),
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user