Add customer mass import service with API route, test coverage, and e-conomic integration
This commit is contained in:
+199
-154
@@ -1,5 +1,6 @@
|
||||
import http from "node:http";
|
||||
import { randomUUID } from "node:crypto";
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { WebSocketServer } from "ws";
|
||||
|
||||
function parseJsonBody(req) {
|
||||
@@ -515,179 +516,183 @@ export function createBrokerServer(options = {}) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (ws.gatewayId) {
|
||||
if (message.type === "COMMAND_RESULT") {
|
||||
const pending = pendingCommands.get(message.commandId);
|
||||
if (!pending) {
|
||||
return;
|
||||
}
|
||||
clearTimeout(pending.timeout);
|
||||
pendingCommands.delete(message.commandId);
|
||||
pending.resolve({
|
||||
ok: Boolean(message.ok),
|
||||
payload: message.payload,
|
||||
error: message.error,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "TELEMETRY") {
|
||||
const payload = message.payload || {};
|
||||
const ingested = await ingestTelemetry(String(ws.gatewayId), payload);
|
||||
broadcastGatewayEvent(String(ws.gatewayId), {
|
||||
type: "gateway.telemetry",
|
||||
gatewayId: String(ws.gatewayId),
|
||||
telemetry: payload,
|
||||
gateway: ingested?.gateway || ingested || null,
|
||||
});
|
||||
broadcastGatewayEvent(String(ws.gatewayId), {
|
||||
type: "stats.updated",
|
||||
gatewayId: String(ws.gatewayId),
|
||||
statistics: ingested?.statistics || ingested || null,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "TASK_EVENT") {
|
||||
const operationId = Number(message.operationId ?? message.payload?.operation_id ?? 0);
|
||||
if (!Number.isFinite(operationId) || operationId <= 0) {
|
||||
try {
|
||||
if (ws.gatewayId) {
|
||||
if (message.type === "COMMAND_RESULT") {
|
||||
const pending = pendingCommands.get(message.commandId);
|
||||
if (!pending) {
|
||||
return;
|
||||
}
|
||||
clearTimeout(pending.timeout);
|
||||
pendingCommands.delete(message.commandId);
|
||||
pending.resolve({
|
||||
ok: Boolean(message.ok),
|
||||
payload: message.payload,
|
||||
error: message.error,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
const operation = await ingestTaskEvent(String(ws.gatewayId), operationId, message.payload || {});
|
||||
broadcastGatewayEvent(String(ws.gatewayId), {
|
||||
type: "task.updated",
|
||||
gatewayId: String(ws.gatewayId),
|
||||
operationId,
|
||||
operation,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "TASK_RESULT") {
|
||||
const operationId = Number(message.operationId ?? message.payload?.operation_id ?? 0);
|
||||
if (!Number.isFinite(operationId) || operationId <= 0) {
|
||||
if (message.type === "TELEMETRY") {
|
||||
const payload = message.payload || {};
|
||||
const ingested = await ingestTelemetry(String(ws.gatewayId), payload);
|
||||
broadcastGatewayEvent(String(ws.gatewayId), {
|
||||
type: "gateway.telemetry",
|
||||
gatewayId: String(ws.gatewayId),
|
||||
telemetry: payload,
|
||||
gateway: ingested?.gateway || ingested || null,
|
||||
});
|
||||
broadcastGatewayEvent(String(ws.gatewayId), {
|
||||
type: "stats.updated",
|
||||
gatewayId: String(ws.gatewayId),
|
||||
statistics: ingested?.statistics || ingested || null,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
const operation = await ingestTaskResult(String(ws.gatewayId), operationId, message.payload || {});
|
||||
broadcastGatewayEvent(String(ws.gatewayId), {
|
||||
type: "task.updated",
|
||||
gatewayId: String(ws.gatewayId),
|
||||
operationId,
|
||||
operation,
|
||||
});
|
||||
await syncGatewayBacklog(String(ws.gatewayId), ws).catch(() => {});
|
||||
if (message.type === "TASK_EVENT") {
|
||||
const operationId = Number(message.operationId ?? message.payload?.operation_id ?? 0);
|
||||
if (!Number.isFinite(operationId) || operationId <= 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
const operation = await ingestTaskEvent(String(ws.gatewayId), operationId, message.payload || {});
|
||||
broadcastGatewayEvent(String(ws.gatewayId), {
|
||||
type: "task.updated",
|
||||
gatewayId: String(ws.gatewayId),
|
||||
operationId,
|
||||
operation,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "TASK_RESULT") {
|
||||
const operationId = Number(message.operationId ?? message.payload?.operation_id ?? 0);
|
||||
if (!Number.isFinite(operationId) || operationId <= 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
const operation = await ingestTaskResult(String(ws.gatewayId), operationId, message.payload || {});
|
||||
broadcastGatewayEvent(String(ws.gatewayId), {
|
||||
type: "task.updated",
|
||||
gatewayId: String(ws.gatewayId),
|
||||
operationId,
|
||||
operation,
|
||||
});
|
||||
await syncGatewayBacklog(String(ws.gatewayId), ws).catch(() => {});
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "LOG_FRAME") {
|
||||
const logEntry = await ingestLogEntry(String(ws.gatewayId), message.payload || {});
|
||||
broadcastGatewayEvent(String(ws.gatewayId), {
|
||||
type: "log.append",
|
||||
gatewayId: String(ws.gatewayId),
|
||||
entry: logEntry,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
if (["SHELL_OUTPUT", "SHELL_OPENED", "SHELL_EXIT"].includes(message.type)) {
|
||||
const sessionRecord = browserShellSessions.get(String(message.sessionId));
|
||||
if (!sessionRecord) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "SHELL_OUTPUT") {
|
||||
sessionRecord.transcript += String(message.data || "");
|
||||
sendJson(sessionRecord.ws, { type: "output", data: String(message.data || "") });
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "SHELL_OPENED") {
|
||||
await markShellSessionOpened(sessionRecord.ws.sessionToken, ws.connectionId || null).catch(() => {});
|
||||
sendJson(sessionRecord.ws, { type: "opened" });
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "SHELL_EXIT") {
|
||||
sendJson(sessionRecord.ws, { type: "closed", code: message.code ?? 0 });
|
||||
await closeBrowserShellSession(sessionRecord, "agent_exit");
|
||||
browserShellSessions.delete(String(message.sessionId));
|
||||
if (sessionRecord.ws.readyState < 2) {
|
||||
sessionRecord.ws.close();
|
||||
}
|
||||
}
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "LOG_FRAME") {
|
||||
const logEntry = await ingestLogEntry(String(ws.gatewayId), message.payload || {});
|
||||
broadcastGatewayEvent(String(ws.gatewayId), {
|
||||
type: "log.append",
|
||||
gatewayId: String(ws.gatewayId),
|
||||
entry: logEntry,
|
||||
});
|
||||
if (ws.sessionInfo) {
|
||||
const sessionId = String(ws.sessionInfo.id);
|
||||
const agent = agents.get(String(ws.sessionInfo.gateway_id));
|
||||
if (!agent || agent.readyState !== 1) {
|
||||
return;
|
||||
}
|
||||
if (message.type === "input") {
|
||||
sendJson(agent, {
|
||||
type: "SHELL_INPUT",
|
||||
payload: {
|
||||
sessionId,
|
||||
data: String(message.data || ""),
|
||||
},
|
||||
});
|
||||
return;
|
||||
}
|
||||
if (message.type === "resize") {
|
||||
sendJson(agent, {
|
||||
type: "RESIZE_ROOT_SHELL",
|
||||
payload: {
|
||||
sessionId,
|
||||
cols: Number(message.cols || 0),
|
||||
rows: Number(message.rows || 0),
|
||||
},
|
||||
});
|
||||
return;
|
||||
}
|
||||
if (message.type === "close") {
|
||||
sendJson(agent, {
|
||||
type: "CLOSE_ROOT_SHELL",
|
||||
payload: { sessionId },
|
||||
});
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
if (["SHELL_OUTPUT", "SHELL_OPENED", "SHELL_EXIT"].includes(message.type)) {
|
||||
const sessionRecord = browserShellSessions.get(String(message.sessionId));
|
||||
if (ws.streamSessionInfo) {
|
||||
const sessionRecord = browserStreamSessions.get(String(ws.streamSessionInfo.id));
|
||||
if (!sessionRecord) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "SHELL_OUTPUT") {
|
||||
sessionRecord.transcript += String(message.data || "");
|
||||
sendJson(sessionRecord.ws, { type: "output", data: String(message.data || "") });
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "SHELL_OPENED") {
|
||||
await markShellSessionOpened(sessionRecord.ws.sessionToken, ws.connectionId || null).catch(() => {});
|
||||
sendJson(sessionRecord.ws, { type: "opened" });
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "SHELL_EXIT") {
|
||||
sendJson(sessionRecord.ws, { type: "closed", code: message.code ?? 0 });
|
||||
await closeBrowserShellSession(sessionRecord, "agent_exit");
|
||||
browserShellSessions.delete(String(message.sessionId));
|
||||
if (sessionRecord.ws.readyState < 2) {
|
||||
sessionRecord.ws.close();
|
||||
if (message.type === "SUBSCRIBE") {
|
||||
for (const scope of parseScopes(message.scopes || message.subscriptions || [])) {
|
||||
sessionRecord.subscriptions.add(scope);
|
||||
}
|
||||
sendJson(sessionRecord.ws, {
|
||||
type: "subscribed",
|
||||
subscriptions: Array.from(sessionRecord.subscriptions.values()),
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "UNSUBSCRIBE") {
|
||||
for (const scope of parseScopes(message.scopes || message.subscriptions || [])) {
|
||||
sessionRecord.subscriptions.delete(scope);
|
||||
}
|
||||
sendJson(sessionRecord.ws, {
|
||||
type: "unsubscribed",
|
||||
subscriptions: Array.from(sessionRecord.subscriptions.values()),
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "PING") {
|
||||
sendJson(sessionRecord.ws, { type: "PONG" });
|
||||
}
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
if (ws.sessionInfo) {
|
||||
const sessionId = String(ws.sessionInfo.id);
|
||||
const agent = agents.get(String(ws.sessionInfo.gateway_id));
|
||||
if (!agent || agent.readyState !== 1) {
|
||||
return;
|
||||
}
|
||||
if (message.type === "input") {
|
||||
sendJson(agent, {
|
||||
type: "SHELL_INPUT",
|
||||
payload: {
|
||||
sessionId,
|
||||
data: String(message.data || ""),
|
||||
},
|
||||
});
|
||||
return;
|
||||
}
|
||||
if (message.type === "resize") {
|
||||
sendJson(agent, {
|
||||
type: "RESIZE_ROOT_SHELL",
|
||||
payload: {
|
||||
sessionId,
|
||||
cols: Number(message.cols || 0),
|
||||
rows: Number(message.rows || 0),
|
||||
},
|
||||
});
|
||||
return;
|
||||
}
|
||||
if (message.type === "close") {
|
||||
sendJson(agent, {
|
||||
type: "CLOSE_ROOT_SHELL",
|
||||
payload: { sessionId },
|
||||
});
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
if (ws.streamSessionInfo) {
|
||||
const sessionRecord = browserStreamSessions.get(String(ws.streamSessionInfo.id));
|
||||
if (!sessionRecord) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "SUBSCRIBE") {
|
||||
for (const scope of parseScopes(message.scopes || message.subscriptions || [])) {
|
||||
sessionRecord.subscriptions.add(scope);
|
||||
}
|
||||
sendJson(sessionRecord.ws, {
|
||||
type: "subscribed",
|
||||
subscriptions: Array.from(sessionRecord.subscriptions.values()),
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "UNSUBSCRIBE") {
|
||||
for (const scope of parseScopes(message.scopes || message.subscriptions || [])) {
|
||||
sessionRecord.subscriptions.delete(scope);
|
||||
}
|
||||
sendJson(sessionRecord.ws, {
|
||||
type: "unsubscribed",
|
||||
subscriptions: Array.from(sessionRecord.subscriptions.values()),
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "PING") {
|
||||
sendJson(sessionRecord.ws, { type: "PONG" });
|
||||
}
|
||||
} catch {
|
||||
// Ignore stale gateway/session delivery errors without killing the broker process.
|
||||
}
|
||||
});
|
||||
|
||||
@@ -788,3 +793,43 @@ export function createBrokerServer(options = {}) {
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
async function runBrokerFromCli() {
|
||||
const broker = createBrokerServer();
|
||||
let shuttingDown = false;
|
||||
const shutdown = async (signal) => {
|
||||
if (shuttingDown) {
|
||||
return;
|
||||
}
|
||||
shuttingDown = true;
|
||||
try {
|
||||
await broker.close();
|
||||
process.exit(0);
|
||||
} catch (error) {
|
||||
console.error(`Failed to shut down broker after ${signal}:`, error);
|
||||
process.exit(1);
|
||||
}
|
||||
};
|
||||
|
||||
process.on("SIGINT", () => {
|
||||
void shutdown("SIGINT");
|
||||
});
|
||||
process.on("SIGTERM", () => {
|
||||
void shutdown("SIGTERM");
|
||||
});
|
||||
|
||||
const requestedPort = Number(process.env.PORT || 4300);
|
||||
const address = await broker.listen(Number.isFinite(requestedPort) ? requestedPort : 4300);
|
||||
const normalizedPort =
|
||||
typeof address === "object" && address !== null && "port" in address
|
||||
? address.port
|
||||
: requestedPort;
|
||||
console.log(`TruckWash edge broker listening on ${normalizedPort}`);
|
||||
}
|
||||
|
||||
if (process.argv[1] && fileURLToPath(import.meta.url) === process.argv[1]) {
|
||||
runBrokerFromCli().catch((error) => {
|
||||
console.error("TruckWash edge broker failed to start:", error);
|
||||
process.exit(1);
|
||||
});
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user