Add handling for self-serve machine signals in edge agent

This commit is contained in:
Jeppe Bundgaard
2026-06-30 16:19:43 +02:00
parent eb21405a3d
commit b0ea771e6a
6 changed files with 99 additions and 1 deletions
+11
View File
@@ -296,6 +296,12 @@ export function createBrokerServer(options = {}) {
? async (_gatewayId, payload = {}) => payload
: async (gatewayId, payload = {}) =>
managerRequest(`/edge-agent/internal/gateways/${gatewayId}/logs`, payload));
const ingestMachineSignal =
options.ingestMachineSignal ||
(authMode === "stub"
? async (_gatewayId, payload = {}) => payload
: async (gatewayId, payload = {}) =>
managerRequest(`/edge-agent/internal/gateways/${gatewayId}/selfserve/machine-signal`, payload));
const broadcastGatewayEvent = (gatewayId, message) => {
const sessionIds = gatewayStreamSessions.get(String(gatewayId));
@@ -853,6 +859,11 @@ export function createBrokerServer(options = {}) {
return;
}
if (message.type === "MACHINE_SIGNAL") {
await ingestMachineSignal(String(ws.gatewayId), message.payload || {});
return;
}
if (["SHELL_OUTPUT", "SHELL_OPENED", "SHELL_EXIT"].includes(message.type)) {
const sessionRecord = browserShellSessions.get(String(message.sessionId));
if (!sessionRecord) {
+45
View File
@@ -766,6 +766,51 @@ test("broker fans out telemetry, task, log, and presence updates to browser gate
await broker.close();
});
test("broker ingests self-serve machine signals from connected agents", async () => {
const machineSignals = [];
const broker = createBrokerServer({
authMode: "stub",
validateAgent: async () => ({ id: "701", gateway_id: "701", label: "CPH Edge 01" }),
ingestMachineSignal: async (gatewayId, payload) => {
machineSignals.push({ gatewayId, payload });
return { recorded: true, lane_id: payload.lane_id };
},
});
const address = await broker.listen(0);
const port = address.port;
const agent = new WebSocket(`ws://127.0.0.1:${port}/ws/agent?gatewayId=701&token=agent-token`);
await new Promise((resolve) => agent.once("open", resolve));
agent.send(
JSON.stringify({
type: "MACHINE_SIGNAL",
payload: {
lane_id: 3,
relay_id: "machine-relay",
component: "input",
channel: 0,
event: "input.toggle_on",
state: true,
},
})
);
await waitFor(() => machineSignals.length === 1, { description: "machine signal ingestion" });
assert.equal(machineSignals[0].gatewayId, "701");
assert.deepEqual(machineSignals[0].payload, {
lane_id: 3,
relay_id: "machine-relay",
component: "input",
channel: 0,
event: "input.toggle_on",
state: true,
});
agent.terminate();
await broker.close();
});
test("broker survives telemetry ingestion failures for stale gateways", async () => {
const broker = createBrokerServer({
authMode: "stub",
File diff suppressed because one or more lines are too long
@@ -109,6 +109,7 @@ class edgeGatewaysRoute
$this->post('/edge-agent/internal/gateways/{id}/operations/{operationId}/events', fn() => $this->handleBrokerOperationEvent());
$this->post('/edge-agent/internal/gateways/{id}/operations/{operationId}/complete', fn() => $this->handleBrokerOperationComplete());
$this->post('/edge-agent/internal/gateways/{id}/logs', fn() => $this->handleBrokerGatewayLogEntry());
$this->post('/edge-agent/internal/gateways/{id}/selfserve/machine-signal', fn() => $this->handleBrokerSelfserveMachineSignal());
$this->post('/edge-agent/internal/browser-streams/validate', fn() => $this->handleBrokerBrowserStreamValidate());
$this->post('/edge-agent/internal/shell-sessions/validate', fn() => $this->handleBrokerShellSessionValidate());
$this->post('/edge-agent/internal/shell-sessions/opened', fn() => $this->handleBrokerShellSessionOpened());
@@ -706,6 +707,21 @@ class edgeGatewaysRoute
));
}
private function handleBrokerSelfserveMachineSignal(): void
{
global /** @var response $response */ $response;
$this->requireBrokerSecret();
$gatewayId = (int)$this->fromRoute('id');
$payload = self::getParametersAsArray();
try {
$result = (new selfserve_machine_signal())->recordBrokerEdgeGatewaySignal($gatewayId, $payload);
$response->success($result, !empty($result['recorded']) ? 201 : 202);
} catch (\Throwable $exception) {
$response->error($exception->getMessage(), 400);
}
}
private function handleBrokerBrowserStreamValidate(): void
{
global /** @var response $response */ $response;
@@ -148,6 +148,25 @@ class selfserve_machine_signal
$gateway = (new edge_gateway_manager())->authenticateGateway($gatewayId, $agentToken);
$departmentId = (int)$gateway->department_id->value();
return $this->recordEdgeGatewaySignalForDepartment($gatewayId, $departmentId, $payload);
}
/**
* @param array<string,mixed> $payload
* @return array<string,mixed>
*/
public function recordBrokerEdgeGatewaySignal(int $gatewayId, array $payload): array
{
$gateway = (new edge_gateway_manager())->getGateway($gatewayId);
return $this->recordEdgeGatewaySignalForDepartment($gatewayId, (int)$gateway['department_id'], $payload);
}
/**
* @param array<string,mixed> $payload
* @return array<string,mixed>
*/
private function recordEdgeGatewaySignalForDepartment(int $gatewayId, int $departmentId, array $payload): array
{
return $this->recordCloudShellySignal(
$departmentId,
isset($payload['lane_id']) ? (int)$payload['lane_id'] : null,
@@ -2950,6 +2950,13 @@ final class TruckwashEdgeAgent
]);
}
if (preg_match('#/edge-agent/gateways/\d+/selfserve/machine-signal$#', $endpoint)) {
return $this->sendBrokerMessage([
'type' => 'MACHINE_SIGNAL',
'payload' => $this->stripAgentAuthentication($payload),
]);
}
return null;
}