cc-libs/packages/trapos-cloud/apis/libcloud.lua

524 lines
20 KiB
Lua

-- trapos-cloud Lua client.
--
-- Implements the Lua side of the trapos-server WS gateway protocol (see
-- .plans/trapos-server/trapos-server-gateway-spec.md). A boot daemon owns one
-- persistent websocket connection (startSession) and bridges WS frames to in-OS
-- `trapos_cloud_*` events; programs talk to the daemon through those events via
-- request(). Both run as separate coroutines under the same Lua state, so they
-- share the single os event queue.
-- Server-side application close code emitted when a newer registration displaces an
-- older one (see trapos-server registry). Treated as terminal: the loser yields rather
-- than reconnecting, so two daemons sharing a traposId collapse to one (no eviction loop).
local DISPLACED_CLOSE_CODE = 4000;
local function defaultOs()
return os;
end
local function isNonEmptyString(value)
return type(value) == 'string' and value ~= '';
end
-- Empty-payload tolerance: a missing / non-table payload is treated as {}. The
-- TS side emits "[]" for an empty object, which decodes to an empty Lua table.
local function coercePayload(payload)
if type(payload) ~= 'table' then
return {};
end
return payload;
end
-- v4 UUID; `random` is injectable for deterministic tests.
local function makeUuid(random)
random = random or math.random;
return (string.gsub('xxxxxxxx-xxxx-4xxx-yxxx-xxxxxxxxxxxx', '[xy]', function(c)
local value = (c == 'x') and random(0, 15) or random(8, 11);
return string.format('%x', value);
end));
end
local function buildConnectUrl(base)
local trimmed = (string.gsub(base, '/+$', ''));
return trimmed .. '/gateway';
end
local function createCloud()
local api = {};
api.uuid = makeUuid;
function api.traposId(osLike)
osLike = osLike or defaultOs();
return tostring(osLike.getComputerID()) .. ':' .. (osLike.getComputerLabel() or '');
end
-- Envelope builders. computer_* frames carry traposSenderId; gateway_* frames
-- (built only for tests/symmetry) carry traposReceiverId. payload is omitted
-- when empty so the wire stays clean and matches the empty-payload contract.
function api.buildComputerRequest(traposSenderId, messageId, msgType, payload)
local frame = {
event = 'computer_request',
traposSenderId = traposSenderId,
messageId = messageId,
type = msgType,
};
if type(payload) == 'table' and next(payload) ~= nil then
frame.payload = payload;
end
return frame;
end
function api.buildComputerResponse(traposSenderId, messageId, msgType, ok, payload, err)
local frame = {
event = 'computer_response',
traposSenderId = traposSenderId,
messageId = messageId,
type = msgType,
ok = ok,
};
if type(payload) == 'table' and next(payload) ~= nil then
frame.payload = payload;
end
if err ~= nil then
frame.error = err;
end
return frame;
end
function api.buildGatewayRequest(traposReceiverId, messageId, msgType, payload)
local frame = {
event = 'server_request',
traposReceiverId = traposReceiverId,
messageId = messageId,
type = msgType,
};
if type(payload) == 'table' and next(payload) ~= nil then
frame.payload = payload;
end
return frame;
end
function api.buildGatewayResponse(traposReceiverId, messageId, msgType, ok, payload, err)
local frame = {
event = 'server_response',
traposReceiverId = traposReceiverId,
messageId = messageId,
type = msgType,
ok = ok,
};
if type(payload) == 'table' and next(payload) ~= nil then
frame.payload = payload;
end
if err ~= nil then
frame.error = err;
end
return frame;
end
-- The hello handshake frame: a computer_request carrying the registration
-- identity (payload.traposId) and, when set, the shared secret.
function api.buildHello(traposId, secret, messageId)
local payload = { traposId = traposId };
if isNonEmptyString(secret) then
payload.secret = secret;
end
return api.buildComputerRequest(traposId, messageId, 'hello', payload);
end
-- Decode and dispatch one inbound frame. Side effects are injected via ctx so
-- onMessage stays unit-testable in isolation:
-- ctx = { traposId, decode, send=fn(frame), queueEvent=fn, onHello=fn(ok,payload,error) }
function api.onMessage(content, ctx)
ctx = ctx or {};
local decode = ctx.decode or textutils.unserializeJSON;
local ok, frame = pcall(decode, content);
if not ok or type(frame) ~= 'table' then
return;
end
if frame.event == 'server_request' then
-- Reverse direction. `ping` stays built-in so it answers without any handler.
-- Other types are delegated to in-OS daemons (re-emitted via ctx.queueRequest)
-- when a handler has advertised the type (ctx.isServed); an unadvertised type
-- still gets an immediate unknown_type reply.
if frame.type == 'ping' then
if ctx.send then
ctx.send(api.buildComputerResponse(
ctx.traposId, frame.messageId, 'ping', true, { traposId = ctx.traposId }, nil
));
end
return;
end
if ctx.isServed and ctx.isServed(frame.type) and ctx.queueRequest then
ctx.queueRequest(frame.messageId, frame.type, coercePayload(frame.payload));
return;
end
if ctx.send then
ctx.send(api.buildComputerResponse(
ctx.traposId, frame.messageId, frame.type, false, nil,
{ code = 'unknown_type', message = 'unknown type: ' .. tostring(frame.type) }
));
end
return;
end
if frame.event == 'server_response' then
if frame.type == 'hello' then
if ctx.onHello then
ctx.onHello(frame.ok == true, coercePayload(frame.payload), frame.error);
end
elseif ctx.queueEvent then
ctx.queueEvent(
'trapos_cloud_reply', frame.messageId, frame.ok == true,
coercePayload(frame.payload), frame.error
);
end
return;
end
-- Ignore everything else (frames we should never receive, malformed events).
end
-- Boot daemon: own one persistent connection, bridge WS<->os-events. Returns
-- live status getters. Injectables (eventloop/http/encode/decode/os/queueEvent/
-- random + delays) keep it testable.
function api.startSession(opts)
opts = opts or {};
local el = assert(opts.eventloop, 'startSession requires an eventloop');
local baseUrl = assert(opts.url, 'startSession requires a url');
local osLike = opts.os or defaultOs();
local httpLike = opts.http or http;
local encode = opts.encode or textutils.serializeJSON;
local decode = opts.decode or textutils.unserializeJSON;
local queueEvent = opts.queueEvent or function(...) return osLike.queueEvent(...); end;
local reconnectDelay = opts.reconnectDelay or 5;
local helloTimeout = opts.helloTimeout or 10;
local getSecret = opts.getSecret or function() return opts.secret; end;
local secret = getSecret();
-- Lifecycle logger: log(level, msg, fields?). Defaults to a no-op so the daemon
-- stays silent unless the caller wires a sink (e.g. a /logs file appender).
local log = opts.log or function() end;
local url = buildConnectUrl(baseUrl);
local traposId = api.traposId(osLike);
-- state: connecting | awaiting_hello | ready | unauthorized | displaced
-- (unauthorized and displaced are terminal: no reconnect.)
local state = 'connecting';
local lastError = nil;
-- Login verification overrides (set by a trapos_cloud_reconnect carrying a candidate
-- secret + correlation token). secretOverride wins over getSecret() in connect() so the
-- candidate is tried without being persisted; verifyToken is echoed on the verdict events
-- so `cloud login` can match its own attempt. Both survive websocket retries during the
-- verify window; a plain (candidate-less) reconnect clears them.
local secretOverride = nil;
local verifyToken = nil;
local activeWs = nil;
local suppressNextClose = false;
local clearReconnectTimeout = nil;
-- True while looping on retries. Gates the per-attempt chatter (connecting / socket
-- open / connection failed / connection closed) so a long outage logs once on entry
-- and once on recovery instead of flooding the file every reconnectDelay seconds.
local reconnecting = false;
-- Gateway_request types advertised by in-OS daemons (via trapos_cloud_serve).
-- Inbound frames of an advertised type are re-emitted as trapos_cloud_request
-- events; unadvertised (non-ping) types get an immediate unknown_type reply.
local servedTypes = {};
local scheduleReconnect = nil;
local function sendFrame(frame)
if not activeWs then
return false;
end
return pcall(function() activeWs.send(encode(frame)); end);
end
local function closeActive()
if activeWs then
pcall(function() activeWs.close(); end);
end
activeWs = nil;
end
local function connect()
if clearReconnectTimeout then
clearReconnectTimeout();
clearReconnectTimeout = nil;
end
secret = secretOverride or getSecret();
state = 'connecting';
if not reconnecting then
log('info', 'connecting', { url = url });
end
local ok, err = pcall(function() httpLike.websocketAsync(url); end);
if not ok then
if not reconnecting then
log('warn', 'websocket open failed', { url = url, error = tostring(err) });
end
scheduleReconnect();
end
end
function scheduleReconnect()
activeWs = nil;
if state == 'unauthorized' or state == 'displaced' then
return;
end
state = 'connecting';
-- One line on the false->true transition, then silence: every later retry runs
-- through here too, but we keep trying without re-logging.
if not reconnecting then
reconnecting = true;
log('warn', 'reconnect loop started', { delaySeconds = reconnectDelay });
end
clearReconnectTimeout = el.setTimeout(connect, reconnectDelay);
end
local function reconnectNow(reason)
log('info', 'reconnect requested', { reason = reason or 'manual' });
lastError = nil;
reconnecting = false;
if activeWs then
suppressNextClose = true;
closeActive();
end
connect();
end
local function onHello(ok, _payload, err)
if state ~= 'awaiting_hello' then
return;
end
if ok then
state = 'ready';
if reconnecting then
reconnecting = false;
log('info', 'reconnect loop ended; reconnected', { traposId = traposId });
else
log('info', 'connected', { traposId = traposId });
end
-- Echo the candidate's verification token (nil for ordinary reconnects), then drop it
-- so later natural reconnects don't re-emit a stale token. secretOverride is kept: it
-- equals the secret `cloud login` is about to persist, so it stays a safe no-op until a
-- plain reconnect clears it.
queueEvent('trapos_cloud_connected', traposId, verifyToken);
verifyToken = nil;
else
state = 'unauthorized';
lastError = (type(err) == 'table' and err.code) or 'unauthorized';
log('warn', 'hello rejected; yielding', { code = lastError });
queueEvent('trapos_cloud_unauthorized', traposId, lastError, verifyToken);
closeActive();
end
end
el.register('websocket_success', function(eventUrl, ws)
if eventUrl ~= url then return; end
activeWs = ws;
state = 'awaiting_hello';
if not reconnecting then
log('info', 'socket open; sending hello', { traposId = traposId });
end
sendFrame(api.buildHello(traposId, secret, makeUuid(opts.random)));
-- hello-ack watchdog: state-guarded, no stored cancel. If hello already
-- resolved this is a harmless no-op; closing triggers reconnect.
el.setTimeout(function()
if state == 'awaiting_hello' then
closeActive();
end
end, helloTimeout);
end);
el.register('websocket_message', function(eventUrl, content)
if eventUrl ~= url then return; end
api.onMessage(content, {
traposId = traposId,
decode = decode,
send = sendFrame,
queueEvent = queueEvent,
onHello = onHello,
isServed = function(msgType) return servedTypes[msgType] == true; end,
queueRequest = function(messageId, msgType, payload)
queueEvent('trapos_cloud_request', messageId, msgType, payload);
end,
});
end);
el.register('websocket_failure', function(eventUrl)
if eventUrl ~= url then return; end
if not reconnecting then
log('warn', 'connection failed', { url = url });
end
scheduleReconnect();
end);
-- CC:T passes (url, reason, code). A displacement close (code 4000) is only
-- terminal when it hits our *live* (ready) socket: another computer registered
-- our traposId, so we yield rather than reconnect (reconnecting would just
-- displace it back, an endless ping-pong). A displacement that arrives while we
-- are still connecting/awaiting_hello is for a socket our own reconnect has
-- already superseded -- ignoring it keeps the healthy new socket alive instead
-- of bricking the daemon into a permanent 'displaced' state.
el.register('websocket_closed', function(eventUrl, reason, code)
if eventUrl ~= url then return; end
if suppressNextClose then
suppressNextClose = false;
log('info', 'ignored intentional close', { reason = reason, code = code });
return;
end
if code == DISPLACED_CLOSE_CODE or reason == 'displaced' then
if state == 'ready' then
activeWs = nil;
state = 'displaced';
lastError = 'displaced';
log('warn', 'displaced by a newer registration; yielding', { traposId = traposId });
else
log('info', 'ignoring stale displacement during reconnect', { traposId = traposId });
end
return;
end
if not reconnecting then
log('info', 'connection closed', { code = code, reason = reason });
end
scheduleReconnect();
end);
-- Bridge program requests onto the wire. Only emit an immediate reply on
-- failure; on success the real reply arrives later as a server_response.
el.register('trapos_cloud_send', function(messageId, msgType, payload)
if state ~= 'ready' then
queueEvent('trapos_cloud_reply', messageId, false, nil,
{ code = 'not_connected', message = 'cloud daemon not connected' });
return;
end
local sent = sendFrame(api.buildComputerRequest(traposId, messageId, msgType, payload));
if not sent then
queueEvent('trapos_cloud_reply', messageId, false, nil,
{ code = 'not_connected', message = 'cloud send failed' });
end
end);
-- In-OS daemons advertise the server_request types they handle. Recorded so
-- onMessage routes those (and only those) as trapos_cloud_request events.
el.register('trapos_cloud_serve', function(msgType)
if type(msgType) == 'string' and msgType ~= '' then
servedTypes[msgType] = true;
end
end);
-- Bridge a server reply back onto the wire as a computer_response. Dropped when
-- not connected: the gateway-side request then times out, which is the reply.
el.register('trapos_cloud_respond', function(messageId, msgType, ok, payload, err)
if state ~= 'ready' then
return;
end
sendFrame(api.buildComputerResponse(traposId, messageId, msgType, ok == true, payload, err));
end);
el.register('trapos_cloud_reconnect', function(reason, candidateSecret, token)
-- A candidate secret arms login verification; anything else (plain reconnect) clears it,
-- so login-restore / login-saved / manual reconnects fall back to the persisted secret.
if type(candidateSecret) == 'string' and candidateSecret ~= '' then
secretOverride = candidateSecret;
verifyToken = token;
else
secretOverride = nil;
verifyToken = nil;
end
reconnectNow(reason);
end);
connect();
return {
isReady = function() return state == 'ready'; end,
lastError = function() return lastError; end,
};
end
-- Program-side blocking request: queue a send, then wait (unfiltered pull, to
-- catch both the reply and our own timer) for a matching reply or timeout.
function api.request(msgType, payload, opts)
opts = opts or {};
local queueEvent = opts.queueEvent or os.queueEvent;
local pullEvent = opts.pullEvent or os.pullEvent;
local startTimer = opts.startTimer or os.startTimer;
local cancelTimer = opts.cancelTimer or os.cancelTimer;
local timeout = opts.timeout or 10;
local messageId = opts.messageId or makeUuid(opts.random);
queueEvent('trapos_cloud_send', messageId, msgType, payload);
local timerId = startTimer(timeout);
while true do
local event = table.pack(pullEvent());
local name = event[1];
if name == 'timer' and event[2] == timerId then
return false, { code = 'timeout', message = 'request timed out' };
elseif name == 'trapos_cloud_reply' then
if event[2] == messageId then
cancelTimer(timerId);
if event[3] then
return true, coercePayload(event[4]);
end
return false, event[5] or { code = 'internal', message = 'request failed' };
end
-- Not ours: re-queue so a concurrent waiter can claim it.
queueEvent('trapos_cloud_reply', event[2], event[3], event[4], event[5]);
end
-- Ignore any other event and keep waiting.
end
end
-- Server-side reply: queue a computer_response for an inbound server_request that
-- the daemon delivered as a trapos_cloud_request event. The daemon owns the wire.
function api.respond(messageId, msgType, ok, payload, err, opts)
opts = opts or {};
local queueEvent = opts.queueEvent or os.queueEvent;
queueEvent('trapos_cloud_respond', messageId, msgType, ok == true, payload, err);
end
-- Register an in-OS handler for an inbound server_request type. Advertises the
-- type to the daemon (so it routes instead of replying unknown_type) and reacts to
-- matching trapos_cloud_request events. handler(payload, messageId) returns
-- (true, resultTable) on success or (false, { code, message }) on failure; a thrown
-- error is reported as { code = 'internal', message = <error> }.
function api.serve(msgType, handler, opts)
opts = opts or {};
local el = assert(opts.eventloop, 'serve requires an eventloop');
assert(type(msgType) == 'string' and msgType ~= '', 'serve requires a type');
assert(type(handler) == 'function', 'serve requires a handler');
local queueEvent = opts.queueEvent or os.queueEvent;
queueEvent('trapos_cloud_serve', msgType);
el.register('trapos_cloud_request', function(messageId, requestType, payload)
if requestType ~= msgType then
return;
end
local ok, result, errBody = pcall(handler, coercePayload(payload), messageId);
if not ok then
api.respond(messageId, msgType, false, nil,
{ code = 'internal', message = tostring(result) }, { queueEvent = queueEvent });
return;
end
if result == true then
api.respond(messageId, msgType, true, errBody, nil, { queueEvent = queueEvent });
else
if type(errBody) ~= 'table' then
errBody = { code = 'internal', message = tostring(errBody or 'request failed') };
end
api.respond(messageId, msgType, false, nil, errBody, { queueEvent = queueEvent });
end
end);
end
return api;
end
return createCloud;