Something went wrong. Try again.
The agentic engineering control plane for the posthuman future
Something went wrong. Try again.
44 kB · 1089 lines
C++
at main
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090// DaemonCore.cpp — see DaemonCore.h for the wiring summary.#include "DaemonCore.h"
#include "ByteLog.h"#include "XenoPaths.h"#include "XenoWire.h"
#include <QCoreApplication>#include <QDir>#include <QFileInfo>#include <QJsonArray>#include <QJsonDocument>#include <QJsonObject>#include <QJsonParseError>#include <QRandomGenerator>#include <QTemporaryFile>#include <QTimer>#include <QUrlQuery>#include <QWebSocket>#include <QWebSocketProtocol>#include <QWebSocketServer>
#include <cerrno>#include <chrono>#include <cstdio>#include <cstring>#include <fcntl.h>#include <sys/file.h>#include <unistd.h>
#include <ranges>#include <utility>
namespace {
constexpr int kPtyChunkBytes = 16 * 1024;
QString jsonValueStr(const QJsonObject& o, const char* key){ return o.value(QLatin1String(key)).toString();}
qint64 jsonValueInt(const QJsonObject& o, const char* key){ return static_cast<qint64>(o.value(QLatin1String(key)).toDouble());}
QStringList jsonStringList(const QJsonObject& o, const char* key){ QStringList out; const QJsonArray arr = o.value(QLatin1String(key)).toArray(); out.reserve(arr.size()); for (const QJsonValueConstRef& v : arr) { out.append(v.toString()); } return out;}
} // namespace
DaemonCore::DaemonCore(QObject* parent) : QObject(parent), m_panes(new PaneTable(this)){ // Reclaim exited, non-restart panes after the grace period. A 5 s sweep is // far tighter than the 60 s grace and cheap (panes are few). auto* reaper = new QTimer(this); reaper->setInterval(5'000); connect(reaper, &QTimer::timeout, this, [this]() { m_panes->reclaimExpired(); }); reaper->start();}
DaemonCore::~DaemonCore(){ if (m_lockFd >= 0) { ::close(m_lockFd); }}
QString DaemonCore::makeToken(){ // 32 CSPRNG bytes as 64 lowercase hex chars. QByteArray raw(32, Qt::Uninitialized); auto* words = reinterpret_cast<quint32*>(raw.data()); // NOLINT(cppcoreguidelines-pro-type-reinterpret-cast) QRandomGenerator* gen = QRandomGenerator::system(); for (int i : std::views::iota(0, 8)) { words[i] = gen->generate(); } return QString::fromLatin1(raw.toHex());}
QString DaemonCore::bootStamp(){ return QDateTime::currentDateTimeUtc().toString(Qt::ISODate);}
bool DaemonCore::claimLock(QString* errorOut){ const QString lockPath = xeno::paths::daemonLockPath(); const QDir dir = QFileInfo(lockPath).dir(); if (!dir.exists() && !QDir().mkpath(dir.absolutePath())) { if (errorOut != nullptr) { *errorOut = QStringLiteral("cannot create config dir: %1") .arg(dir.absolutePath()); } return false; } m_lockFd = ::open(lockPath.toUtf8().constData(), O_RDWR | O_CREAT | O_CLOEXEC, 0600); if (m_lockFd < 0) { if (errorOut != nullptr) { *errorOut = QStringLiteral("cannot open lock file %1").arg(lockPath); } return false; } // LOCK_NB: a second daemon must fail immediately, not queue. if (::flock(m_lockFd, LOCK_EX | LOCK_NB) != 0) { if (errorOut != nullptr) { *errorOut = QStringLiteral( "another xenomorphicd holds the lock (%1); exiting") .arg(lockPath); } ::close(m_lockFd); m_lockFd = -1; return false; } return true;}
bool DaemonCore::writeDiscoveryFile(QString* errorOut){ const QString path = xeno::paths::daemonJsonPath(); const QFileInfo fi(path); if (!QDir().mkpath(fi.dir().absolutePath())) { if (errorOut != nullptr) { *errorOut = QStringLiteral("cannot create dir for %1").arg(path); } return false; }
QJsonObject obj; obj.insert(QStringLiteral("version"), 1); obj.insert(QStringLiteral("port"), m_port); obj.insert(QStringLiteral("token"), m_token); obj.insert(QStringLiteral("pid"), QCoreApplication::applicationPid()); obj.insert(QStringLiteral("boot"), m_boot); const QByteArray body = QJsonDocument(obj).toJson(QJsonDocument::Compact);
// Atomic write: temp file in the SAME directory, then rename(). // Disable auto-remove so WE control the rename (and the temp does not // vanish before it lands at the final path). QTemporaryFile tmp(fi.dir().filePath(QStringLiteral("daemon.json.tmp.XXXXXX"))); tmp.setAutoRemove(false); if (!tmp.open()) { if (errorOut != nullptr) { *errorOut = QStringLiteral("cannot open temp file beside %1").arg(path); } return false; } if (tmp.write(body) != body.size()) { if (errorOut != nullptr) { *errorOut = QStringLiteral("short write to discovery temp file"); } return false; } tmp.flush(); ::fsync(tmp.handle()); // 0600 before the rename so the token is never world-readable. if (!tmp.setPermissions( QFile::Permissions(QFile::ReadOwner | QFile::WriteOwner))) { if (errorOut != nullptr) { *errorOut = QStringLiteral("cannot chmod discovery temp file"); } return false; } const QString tmpPath = tmp.fileName(); tmp.close(); // POSIX rename(2), NOT QFile::rename(). The whole point of the // write-temp-then-rename idiom is that the rename atomically REPLACES any // existing daemon.json — which is the normal case after a crash or a daemon // restart, where a stale file from the previous boot is still there. // QFile::rename() deliberately refuses to overwrite ("Destination file // exists"), so it fails exactly when a stale file is present. ::rename() is // atomic within one directory (the temp is beside the target by design) and // overwrites, giving readers either the old file or the new one, never a // partial one. const QByteArray tmpUtf8 = QFile::encodeName(tmpPath); const QByteArray dstUtf8 = QFile::encodeName(path); if (::rename(tmpUtf8.constData(), dstUtf8.constData()) != 0) { if (errorOut != nullptr) { *errorOut = QStringLiteral("rename(%1 -> %2) failed: %3") .arg(tmpPath, path, QString::fromLocal8Bit(strerror(errno))); } QFile::remove(tmpPath); return false; } return true;}
bool DaemonCore::start(QString* errorOut){ if (m_server != nullptr) { return true; // idempotent } if (!claimLock(errorOut)) { return false; }
m_token = makeToken(); m_boot = bootStamp();
m_server = new QWebSocketServer(QStringLiteral("xenomorphicd"), QWebSocketServer::SslMode::NonSecureMode, this); // NOTE: incoming/outgoing size limits live on QWebSocket (per socket), // not QWebSocketServer; applied in onNewConnection().
if (!m_server->listen(QHostAddress::LocalHost, 0)) { // ephemeral port if (errorOut != nullptr) { *errorOut = QStringLiteral("listen(127.0.0.1, 0) failed: %1") .arg(m_server->errorString()); } return false; } m_port = m_server->serverPort(); connect(m_server, &QWebSocketServer::newConnection, this, &DaemonCore::onNewConnection);
if (!writeDiscoveryFile(errorOut)) { return false; } qInfo("xenomorphicd: pid=%lld listening on 127.0.0.1:%u boot=%s", static_cast<long long>(QCoreApplication::applicationPid()), m_port, qPrintable(m_boot)); return true;}
void DaemonCore::shutdown(const QString& reason){ if (m_shuttingDown) { return; } m_shuttingDown = true;
// Broadcast daemonDown FIRST so clients stop auto-respawning, then close // listeners. Clients get CloseCodeGoingAway (1001) on their sockets. // // Both loops iterate a SNAPSHOT of socket pointers, never the live // containers. close() can re-enter our disconnected() handlers, which call // m_ctlSockets.remove() / m_dataConns.removeOne() plus `delete conn` — that // would invalidate the iterator and free the DataConn under us. The snapshot // is safe because every socket is destroyed via deleteLater(), so the // QWebSocket* stays valid until this synchronous call stack unwinds even if // its DataConn is already gone. const QList<QWebSocket*> ctlSockets = m_ctlSockets.values(); for (QWebSocket* ctl : ctlSockets) { QJsonObject notif; notif.insert(QStringLiteral("op"), QStringLiteral("daemonDown")); notif.insert(QStringLiteral("reason"), reason); ctl->sendTextMessage( QString::fromUtf8(QJsonDocument(notif).toJson(QJsonDocument::Compact))); ctl->close(QWebSocketProtocol::CloseCodeGoingAway, reason); } QVector<QWebSocket*> paneSockets; paneSockets.reserve(m_dataConns.size()); for (const DataConn* conn : std::as_const(m_dataConns)) { if (conn->socket != nullptr) { paneSockets.append(conn->socket); } } for (QWebSocket* sock : std::as_const(paneSockets)) { sock->close(QWebSocketProtocol::CloseCodeGoingAway, reason); } if (m_server != nullptr) { m_server->close(); }}
static void closeWith(QWebSocket* socket, quint16 code, const QString& reason){ socket->close(static_cast<QWebSocketProtocol::CloseCode>(code), reason); // The socket is destroyed on disconnected() below.}
// -- connection routing and auth -------------------------------------------
void DaemonCore::onNewConnection(){ while (m_server->hasPendingConnections()) { QWebSocket* socket = m_server->nextPendingConnection(); // Bound what a client can send and tune outgoing framing for the 16 KB // pty chunk. These live on the socket, not the server. socket->setMaxAllowedIncomingFrameSize( static_cast<quint64>(kMaxIncomingFrame)); socket->setMaxAllowedIncomingMessageSize( static_cast<quint64>(kMaxIncomingFrame)); socket->setOutgoingFrameSize(static_cast<quint64>(kPtyChunkBytes) * 4);
// Auth on BOTH paths, from the handshake header only. const QByteArray token = socket->request().rawHeader("X-Xeno-Token"); if (token != m_token.toLatin1()) { closeWith(socket, static_cast<quint16>( QWebSocketProtocol::CloseCodePolicyViolated), QStringLiteral("unauthorized")); connect(socket, &QWebSocket::disconnected, socket, &QWebSocket::deleteLater); continue; }
const QString path = socket->requestUrl().path(); if (path == QStringLiteral("/ctl")) { m_ctlSockets.insert(socket); connect(socket, &QWebSocket::textMessageReceived, this, [this, socket](const QString& msg) { onCtlText(socket, msg); }); connect(socket, &QWebSocket::disconnected, this, [this, socket]() { m_ctlSockets.remove(socket); socket->deleteLater(); }); } else if (path.startsWith(QStringLiteral("/pane/"))) { const QString paneId = path.mid(QString(QStringLiteral("/pane/")).size()); QUrlQuery query(socket->requestUrl().query()); quint64 fromSeq = 0; if (query.hasQueryItem(QStringLiteral("fromSeq"))) { fromSeq = query.queryItemValue(QStringLiteral("fromSeq")).toULongLong(); } acceptData(socket, paneId, fromSeq); } else { closeWith(socket, static_cast<quint16>(QWebSocketProtocol::CloseCodePolicyViolated), QStringLiteral("bad-path")); connect(socket, &QWebSocket::disconnected, socket, &QWebSocket::deleteLater); } }}
// -- control plane ----------------------------------------------------------
QJsonObject DaemonCore::messageObject(const QString& message){ QJsonParseError err{}; const QJsonDocument doc = QJsonDocument::fromJson(message.toUtf8(), &err); if (!doc.isObject()) { return {}; } return doc.object();}
static void sendOkImpl(QWebSocket* socket, quint64 id, const QJsonObject& result){ QJsonObject reply; reply.insert(QStringLiteral("op"), QStringLiteral("ok")); reply.insert(QStringLiteral("id"), static_cast<qint64>(id)); reply.insert(QStringLiteral("result"), result); socket->sendTextMessage( QString::fromUtf8(QJsonDocument(reply).toJson(QJsonDocument::Compact)));}
static void sendErrImpl(QWebSocket* socket, quint64 id, const QString& error){ QJsonObject reply; reply.insert(QStringLiteral("op"), QStringLiteral("err")); reply.insert(QStringLiteral("id"), static_cast<qint64>(id)); reply.insert(QStringLiteral("error"), error); socket->sendTextMessage( QString::fromUtf8(QJsonDocument(reply).toJson(QJsonDocument::Compact)));}
void DaemonCore::onCtlText(QWebSocket* socket, const QString& message){ const QJsonObject msg = messageObject(message); if (msg.isEmpty()) { sendErrImpl(socket, 0, QStringLiteral("bad-json")); return; } const QString op = jsonValueStr(msg, "op"); const auto id = static_cast<quint64>(jsonValueInt(msg, "id")); if (op == QStringLiteral("hello")) { handleHello(socket, msg); } else if (op == QStringLiteral("spawn")) { handleSpawn(socket, msg); } else if (op == QStringLiteral("kill")) { handleKill(socket, msg); } else if (op == QStringLiteral("resize")) { handleResize(msg); sendOkImpl(socket, id, {}); } else if (op == QStringLiteral("list")) { handleList(socket, msg); } else if (op == QStringLiteral("shutdown")) { handleShutdown(socket, msg); } else { sendErrImpl(socket, id, QStringLiteral("unknown-op")); }}
void DaemonCore::handleHello(QWebSocket* socket, const QJsonObject& msg){ const auto id = static_cast<quint64>(jsonValueInt(msg, "id")); QJsonObject result; result.insert(QStringLiteral("version"), 1); result.insert(QStringLiteral("boot"), m_boot); result.insert(QStringLiteral("pid"), QCoreApplication::applicationPid()); sendOkImpl(socket, id, result);}
void DaemonCore::handleSpawn(QWebSocket* socket, const QJsonObject& msg){ const auto id = static_cast<quint64>(jsonValueInt(msg, "id")); const QJsonObject spec = msg.value(QStringLiteral("spec")).toObject(); const QString paneId = jsonValueStr(spec, "paneId"); if (paneId.isEmpty()) { sendErrImpl(socket, id, QStringLiteral("missing-paneId")); return; } // Reject only a LIVE duplicate, per the contract ("duplicate live id"). // An EXITED-but-retained pane (inside the 60 s grace window, or held for // restart) still occupies its paneId in the table, but its shell is gone, so // a respawn for that id must succeed rather than fail with "pane-exists". // Reaching this with a retained exited pane is the normal relaunch path: // the client's paneAlive() reads `alive = !exited`, sees false, and sends // spawn instead of reattaching. Reclaim the dead shell first so adopt() does // not leak the old Pane* and the fresh pty gets a clean ByteLog. if (m_panes->find(paneId, /*liveOnly=*/true) != nullptr) { sendErrImpl(socket, id, QStringLiteral("pane-exists")); return; } if (m_panes->contains(paneId)) { // Exited and retained: drop any lingering data connection and the pane // itself (closes the dead pty, frees its log) before respawning fresh // under the same id. Without dropPaneConnections a stale DataConn could // feed the new pane's bytes under an old streamId; not reachable with a // single client (the relaunched app's previous connections are gone), but // the invariant should hold regardless of client count. dropPaneConnections(paneId); m_panes->killPane(paneId); } if (m_panes->full()) { sendErrImpl(socket, id, QStringLiteral("too-many-panes")); return; }
PaneTable::SpawnSpec rec; rec.paneId = paneId; rec.cwd = jsonValueStr(spec, "cwd"); rec.command = jsonValueStr(spec, "command"); rec.args = jsonStringList(spec, "args"); // env passes through VERBATIM: Pty implements the override semantics // (an empty value REMOVES the inherited variable). Never pre-process it. rec.env = jsonStringList(spec, "env"); rec.initialInput = jsonValueStr(spec, "initialInput").toUtf8(); rec.initialCommands = jsonStringList(spec, "initialCommands"); rec.title = jsonValueStr(spec, "title"); rec.restart = spec.value(QStringLiteral("restart")).toBool(); rec.restartArgv = jsonStringList(spec, "restartArgv");
auto pty = std::make_unique<Pty>(); QString err; // command empty = default $SHELL in login form, per Pty::start. A non-empty // command is argv[0]; Pty::start does the login-shell dash-prefix itself. const QString shell = rec.command.isEmpty() ? QString() : rec.command; if (!pty->start(shell, rec.args, rec.cwd, rec.env, &err)) { sendErrImpl(socket, id, QStringLiteral("spawn-failed: %1").arg(err)); return; }
PaneTable::Pane* pane = m_panes->adopt(rec, std::move(pty));
// Daemon-side initial input: initialInput, then each initialCommands entry // newline-terminated, empty/whitespace-only lines skipped — the same order // TerminalQuickItem::ensureTerminal() used to write client-side. if (!rec.initialInput.isEmpty()) { pane->pty->write(rec.initialInput); } for (const QString& cmd : std::as_const(rec.initialCommands)) { if (cmd.trimmed().isEmpty()) { continue; } pane->pty->write(cmd.toUtf8() + '\n'); }
// Capture paneId ONLY — never the Pty pointer. This is a QUEUED // cross-thread signal, so the lambda body runs later, on the event loop, // as a QMetaCallEvent whose receiver is `this` (DaemonCore). Destroying // the Pty in killPane does NOT remove events already posted to DaemonCore, // so a captured raw Pty* would be dangling by the time the event ran: // ackChunk() then locks a freed mutex and std::mutex::lock() throws // system_error(EINVAL), which nothing catches, so the daemon SIGABRTs. // Reproduced: flood a pane, then /ctl kill it -> exit 134 with the faulting // stack in Pty::ackChunk under sendPostedEvents. // appendChunk re-resolves the pane from the table instead, exactly as // appendBytes already does, so a killed pane is simply not found. connect(pane->pty.get(), &Pty::readyRead, this, [this, paneId](const QByteArray& data) { appendChunk(paneId, data); }); connect(pane->pty.get(), &Pty::processExited, this, [this, paneId](int code) { onPaneExited(paneId, code); });
QJsonObject result; result.insert(QStringLiteral("paneId"), paneId); result.insert(QStringLiteral("seq"), static_cast<qint64>(pane->log.tailSeq())); result.insert(QStringLiteral("pid"), pane->pty->pid()); result.insert(QStringLiteral("shellPath"), pane->pty->shellPath()); sendOkImpl(socket, id, result);}
void DaemonCore::handleKill(QWebSocket* socket, const QJsonObject& msg){ const auto id = static_cast<quint64>(jsonValueInt(msg, "id")); const QString paneId = jsonValueStr(msg, "paneId"); PaneTable::Pane* pane = m_panes->find(paneId); if (pane == nullptr) { sendErrImpl(socket, id, QStringLiteral("no-such-pane")); return; } // Drop data connections for this pane first so no further frames cross, // then reclaim the pane itself. dropPaneConnections(paneId); m_panes->killPane(paneId); sendOkImpl(socket, id, {});}
void DaemonCore::handleResize(const QJsonObject& msg){ const QString paneId = jsonValueStr(msg, "paneId"); const int cols = static_cast<int>(jsonValueInt(msg, "cols")); const int rows = static_cast<int>(jsonValueInt(msg, "rows")); PaneTable::Pane* pane = m_panes->find(paneId, /*liveOnly=*/true); if (pane == nullptr || cols <= 0 || rows <= 0) { return; } // A resize is a REQUEST: apply it to the pty, then append the RESIZE // record to the log under the serializer lock so it lands at exactly the // right position relative to concurrent BYTES. pane->pty->resize(cols, rows); appendRecord(paneId, static_cast<quint8>(xeno::wire::Tag::Resize), xeno::wire::resizePayload(cols, rows));}
void DaemonCore::handleList(QWebSocket* socket, const QJsonObject& msg){ const auto id = static_cast<quint64>(jsonValueInt(msg, "id")); QJsonArray panes; for (const PaneTable::Pane* pane : std::as_const(*m_panes)) { QJsonObject o; o.insert(QStringLiteral("paneId"), pane->paneId); o.insert(QStringLiteral("pid"), pane->pty != nullptr ? pane->pty->pid() : 0); o.insert(QStringLiteral("foregroundPid"), pane->pty != nullptr ? pane->pty->foregroundPid() : 0); // cwd is the SPAWN-TIME value only (set in handleSpawn from the spec), NOT // a live lookup. It used to come from Pty::foregroundCwd(), which spawned // `lsof` under waitForFinished() ON THE EVENT LOOP for every live pane // (~117 ms per pane measured; a 15-pane session made /ctl list cost ~1.8 s // against the 8 s handshake budget, which contributed to relaunch failures). // // IMPORTANT, and a real limitation: the daemon does NOT parse OSC-7 and // nothing else updates pane->cwd, so a shell that `cd`s elsewhere still // reports its spawn directory here. Do not describe this field as tracking // the foreground process. Live cwd is tracked CLIENT-side instead: the // ghostty VT parses OSC-7 and raises workingDirectoryChanged // (TerminalQuickItem.cpp) -> PanePool -> AppModel::notePaneCwd, which is // what session save uses. So `list`'s cwd is a coarse spawn-time hint, and // the authoritative cwd lives in AppModel's PaneRecord. // // That is acceptable because no code path consumes this field's liveness: // PtyClient::foregroundCwd() (served from this cache) has no callers // anywhere in the tree, and the quit-time restart scan matches on argv via // ProcessScanner, not on cwd. If a consumer ever needs live cwd from the // daemon, the daemon must learn to parse OSC-7 out of the byte stream — // do not reintroduce the per-pane `lsof` on the event loop to get it. o.insert(QStringLiteral("cwd"), pane->cwd); o.insert(QStringLiteral("title"), pane->title); o.insert(QStringLiteral("shellPath"), pane->pty != nullptr ? pane->pty->shellPath() : QString()); o.insert(QStringLiteral("alive"), !pane->exited); o.insert(QStringLiteral("headSeq"), static_cast<qint64>(pane->log.headSeq())); o.insert(QStringLiteral("tailSeq"), static_cast<qint64>(pane->log.tailSeq())); panes.append(o); } QJsonObject result; result.insert(QStringLiteral("panes"), panes); // Open /ctl connections (including the requester). Diagnostics surface for // the settings dialog: a stuck client shows up as a nonzero count. result.insert(QStringLiteral("ctlClients"), m_ctlSockets.size()); sendOkImpl(socket, id, result);}
void DaemonCore::handleShutdown(QWebSocket* socket, const QJsonObject& msg){ const auto id = static_cast<quint64>(jsonValueInt(msg, "id")); sendOkImpl(socket, id, {}); // Queue the teardown one event-loop turn so the ok reply (and the // daemonDown notification shutdown() broadcasts) are written and flushed // before exit(0) stops the loop — a synchronous shutdown here would drop // both frames mid-flush. shutdown() is idempotent (m_shuttingDown), so a // second shutdown op or an in-flight SIGTERM teardown is harmless. QMetaObject::invokeMethod( this, [this]() { shutdown(QStringLiteral("shutdown-requested")); QCoreApplication::exit(0); }, Qt::QueuedConnection);}
// -- serialization + data plane ----------------------------------------------
// THE serializer. Single-thread confinement, NOT a lock — there is no mutex// here and none is needed, because every caller already runs on the daemon// event-loop thread://// - pty bytes: Pty::readyRead is emitted from Pty's reader thread// (Pty.cpp), connected with AutoConnection to a DaemonCore receiver,// so Qt QUEUES it and appendChunk runs on the event loop. The queue is// BOUNDED by Pty's high/low water-mark producer throttle: the reader// thread pauses reading (filling the kernel pty buffer, blocking the// flooding program) until the consumer acks each chunk. The small mutex// guarding that handoff counter lives entirely inside Pty — the// event-loop side never blocks on it — and guards no pane/log state, so// the no-lock argument below is untouched.// - exits: Pty::processExited, also reader-thread (Pty.cpp:270/327), also// queued -> onPaneExited on the event loop.// - resizes: the /ctl socket handler (handleResize) is on the event loop.// - reclamation: the QTimer reaper is on the event loop.//// Per-pane seq is assigned inside ByteLog::append, so the log's record order IS// the total serialization order and clients only ever replay that order. This// is what makes cross-connection ordering safe with no cross-socket sequencing.//// DO NOT call this from another thread (e.g. directly from a pty reader thread)// and DO NOT add a lock "to be safe" without first making fanOut's socket writes// thread-safe too: nothing here is guarded, and the confinement is the whole// correctness argument. If a future path needs to append off-thread, marshal it// to the event loop (QMetaObject::invokeMethod Qt::QueuedConnection), do not// reach in.void DaemonCore::appendRecord(const QString& paneId, quint8 tag, const QByteArray& payload){ PaneTable::Pane* pane = m_panes->find(paneId); if (pane == nullptr) { return; } pane->log.append(tag, payload); // Fan-out happens on the event loop, never inline from the reader thread. fanOut(paneId);}
void DaemonCore::appendBytes(const QString& paneId, const QByteArray& data){ // The stamp is taken HERE — on the thread that received readyRead, which // for Pty's queued cross-thread signal is the daemon event loop, right // before the record is appended (mirrors Pty::stampPerfChunk()). const quint64 nowUs = static_cast<quint64>( std::chrono::duration_cast<std::chrono::microseconds>( std::chrono::steady_clock::now().time_since_epoch()) .count()); appendRecord(paneId, static_cast<quint8>(xeno::wire::Tag::Bytes), xeno::wire::bytesPayload(nowUs, data));}
void DaemonCore::appendChunk(const QString& paneId, const QByteArray& data){ // Flow-control consumer side: consume the chunk, then acknowledge it to the // pane's Pty so the reader thread resumes after a high-water pause. The ack // happens HERE — after actual consumption (log append + fanOut), not at // queue time — so it measures real drain, which is the whole point of the // producer throttle. // // The pane is re-resolved from the table rather than passed in, because this // runs from a QUEUED cross-thread signal: events posted before a /ctl kill // still arrive after killPane destroyed the Pty. Looking up means a killed // pane is simply not found and we neither append nor ack — there is no Pty // left to ack, and its reader thread is gone with it. // // One lookup suffices for the whole call. Nothing in appendBytes -> fanOut // can destroy the Pane: fanOut's 1009 closeWith only deletes a DataConn, and // pane destruction happens in killPane, called from handleKill or the QTimer // reaper — both a separate event-loop turn, never inline from here. PaneTable::Pane* pane = m_panes->find(paneId); if (pane == nullptr) { return; } appendBytes(paneId, data); if (pane->pty != nullptr) { pane->pty->ackChunk(); }}
void DaemonCore::onPaneExited(const QString& paneId, int code){ // The EXIT record participates in the same total order as bytes/resizes. PaneTable::Pane* pane = m_panes->find(paneId); const quint64 exitSeq = pane != nullptr ? pane->log.tailSeq() : 0; appendRecord(paneId, static_cast<quint8>(xeno::wire::Tag::Exit), xeno::wire::exitPayload(code)); m_panes->markExited(paneId, code);
// Notify every control connection. QJsonObject notif; notif.insert(QStringLiteral("op"), QStringLiteral("exit")); notif.insert(QStringLiteral("paneId"), paneId); notif.insert(QStringLiteral("code"), code); notif.insert(QStringLiteral("seq"), static_cast<qint64>(exitSeq)); const QString text = QString::fromUtf8(QJsonDocument(notif).toJson(QJsonDocument::Compact)); for (QWebSocket* ctl : std::as_const(m_ctlSockets)) { ctl->sendTextMessage(text); }}
void DaemonCore::acceptData(QWebSocket* socket, const QString& paneId, quint64 fromSeqRaw){ PaneTable::Pane* pane = m_panes->find(paneId); if (pane == nullptr) { closeWith(socket, static_cast<quint16>(QWebSocketProtocol::CloseCodePolicyViolated), QStringLiteral("no-such-pane")); connect(socket, &QWebSocket::disconnected, socket, &QWebSocket::deleteLater); return; }
auto* conn = new DataConn; // owned by m_dataConns, freed on disconnect conn->socket = socket; conn->paneId = paneId; conn->streamId = m_nextStreamId; ++m_nextStreamId;
const quint64 floor = pane->log.floor(); if (fromSeqRaw < floor) { conn->fromSeq = floor; // clamp; report the real floor } else { conn->fromSeq = fromSeqRaw; } conn->sentUpTo = conn->fromSeq;
// Exactly ONE text frame: the attach receipt. Everything after is binary. QJsonObject attached; attached.insert(QStringLiteral("op"), QStringLiteral("attached")); attached.insert(QStringLiteral("paneId"), paneId); attached.insert(QStringLiteral("streamId"), static_cast<qint64>(conn->streamId)); attached.insert(QStringLiteral("fromSeq"), static_cast<qint64>(conn->fromSeq)); attached.insert(QStringLiteral("floor"), static_cast<qint64>(floor)); attached.insert(QStringLiteral("tailSeq"), static_cast<qint64>(pane->log.tailSeq())); socket->sendTextMessage( QString::fromUtf8(QJsonDocument(attached).toJson(QJsonDocument::Compact)));
// Replay [fromSeq, tailSeq) as binary frames — CHUNKED. The old code // replayed the whole retained log in one synchronous loop: with a flooded // pane (measured: 340k retained records) that allocated the entire set as a // QVector copy, starved the event loop for minutes, and let RSS balloon // past 20 GB until the daemon was OOM-killed. Instead: send ONE bounded // first batch now, then pumpReplay() re-arms on the event loop (driven by // the socket's bytesWritten drain signal, with a queued safety-net re-arm) // until everything below the attach-time tailSeq is queued. fanOut skips a // replaying conn so the two write paths never interleave out of order. const quint64 attachTailSeq = pane->log.tailSeq(); conn->replaying = true; conn->replayTailSeq = attachTailSeq;
connect(socket, &QWebSocket::bytesWritten, this, [this, socket](qint64) { // Drain-driven resume: bytes were flushed to the OS, so there is // backlog head-room. This is the PRIMARY resume trigger and it // runs the batch INLINE (one bounded batch per turn, see // pumpReplay) — it does not post another event, so it cannot // grow any queue. pumpReplay re-validates the conn. DataConn* c = connFor(socket); if (c != nullptr && c->replaying) { pumpReplay(c); } }); connect(socket, &QWebSocket::binaryMessageReceived, this, [this, socket](const QByteArray& frame) { DataConn* c = connFor(socket); if (c != nullptr) { onDataBinary(c, frame); } }); connect(socket, &QWebSocket::disconnected, this, [this, socket]() { DataConn* c = connFor(socket); if (c != nullptr) { onDataDisconnect(c); } }); m_dataConns.append(conn); pane->attachedCount++; // First batch of the chunked replay. pumpReplay sends at most // kReplayBatchRecords and re-arms itself via the event loop; the // bytesWritten connection above resumes it as the socket drains. pumpReplay(conn);}
DaemonCore::DataConn* DaemonCore::connFor(QWebSocket* socket){ for (DataConn* conn : m_dataConns) { if (conn->socket == socket) { return conn; } } return nullptr;}
void DaemonCore::pumpReplay(DataConn* conn){ // Re-validate membership before touching conn: closeWith() during the last // batch (or a disconnected() re-entry) deletes the DataConn IMMEDIATELY // (see fanOut's comment and onDataDisconnect), and this call can be a // queued lambda or a bytesWritten callback from a socket that outlived it. if (!m_dataConns.contains(conn) || !conn->replaying || conn->socket == nullptr) { return; } QWebSocket* sock = conn->socket;
// Disarm the queued re-arm FIRST: exactly one queued pump event may exist // per connection at any time. A first version let every yield post another // queued call while bytesWritten ALSO re-entered inline; with a fast // local socket the two triggers multiplied and the posted-event queue // grew without bound (measured: RSS 6+ GB, event-loop starvation). With // rearmPending=false the queue depth for this conn is bounded at one. conn->rearmPending = false;
// Pause, never close, while the socket is congested (same cap as fanOut's // live path). bytesWritten resumes us; the re-armed queued call is the // safety net against a missed drain signal. NO 1009 close here: the client // does not yet re-attach after a backlog close, so closing mid-replay // would blank the pane instead of pausing it. if (sock->bytesToWrite() > kSocketBacklogCap) { return; }
PaneTable::Pane* pane = m_panes->find(conn->paneId); if (pane == nullptr) { // Pane reclaimed mid-replay: nothing left to replay; hand off to fanOut. conn->replaying = false; return; }
if (conn->sentUpTo >= conn->replayTailSeq) { // Everything below the attach-time tailSeq is queued (or was trimmed // past — the floor clamp gives the same semantics as the old synchronous // read). Release the conn to fanOut for live records; strictly- // increasing seq is preserved because sentUpTo only ever advanced over // consecutively queued records. conn->replaying = false; return; }
// ONE bounded batch per event-loop turn. This call runs either from // acceptData, from a queued re-arm, or from the bytesWritten drain signal — // whichever comes first — but never in a loop of its own. Between calls the // event loop runs and /ctl, pty appends, and new connections get serviced. // Record bound: kReplayBatchRecords (256, see the constant's comment — // roughly one 16 KB pty-chunk worth of payload per turn). Byte bound: stop // before the batch could exceed kSocketBacklogCap even if every remaining // record were max-sized (kMaxIncomingFrame), so bytesToWrite() can never // blow past the cap on this path either. qint64 batchBytes = 0; const QVector<ByteLog::Record> batch = pane->log.readFrom(conn->sentUpTo, kReplayBatchRecords); bool hitBoundary = false; for (const ByteLog::Record& r : batch) { // Records at/beyond the attach-time tailSeq belong to fanOut's live path. if (r.seq >= conn->replayTailSeq) { hitBoundary = true; break; } xeno::wire::Header h; h.streamId = conn->streamId; h.seq = r.seq; h.tag = r.tag; QByteArray out; xeno::wire::encode(out, h, r.payload); batchBytes += out.size(); sock->sendBinaryMessage(out); conn->sentUpTo = r.seq + 1; if (batchBytes + kMaxIncomingFrame > kSocketBacklogCap || sock->bytesToWrite() > kSocketBacklogCap) { break; // Byte bound: yield; drain resumes us. } }
if (!hitBoundary && conn->sentUpTo < conn->replayTailSeq) { // More to replay. Rearm via the event loop: the PRIMARY trigger is the // socket's bytesWritten drain (connected in acceptData), which fires per // write flush; this queued call is the safety net so a batch that fits // entirely in the OS socket buffer (no observable bytesWritten) still // progresses. Rearm only if none is already queued (guarded above). conn->rearmPending = true; QMetaObject::invokeMethod( this, [this, conn]() { pumpReplay(conn); }, Qt::QueuedConnection); } else if (hitBoundary || conn->sentUpTo >= conn->replayTailSeq) { // Replayed everything below the attach-time tailSeq (the boundary hit is // just the earliest form of that). Release the conn to fanOut. conn->replaying = false; }}
void DaemonCore::onDataBinary(DataConn* conn, const QByteArray& frame){ // Upstream: STDIN frames only. seq is 0 and ignored; streamId must match. QVector<xeno::wire::Frame> frames; QByteArray buf = frame; if (!xeno::wire::decode(buf, frames)) { closeWith(conn->socket, static_cast<quint16>(QWebSocketProtocol::CloseCodeBadOperation), QStringLiteral("bad-frame")); return; } PaneTable::Pane* pane = m_panes->find(conn->paneId, /*liveOnly=*/true); for (const xeno::wire::Frame& f : std::as_const(frames)) { if (f.header.tag != static_cast<quint8>(xeno::wire::UpTag::Stdin) || f.header.streamId != conn->streamId) { closeWith(conn->socket, static_cast<quint16>(QWebSocketProtocol::CloseCodeBadOperation), QStringLiteral("bad-tag")); return; } if (pane != nullptr && pane->pty != nullptr) { pane->pty->write(f.payload); // straight to the master fd, no interpretation } }}
void DaemonCore::onDataDisconnect(DataConn* conn){ PaneTable::Pane* pane = m_panes->find(conn->paneId); if (pane != nullptr && pane->attachedCount > 0) { pane->attachedCount--; } detachConn(conn);}
void DaemonCore::detachConn(DataConn* conn){ QWebSocket* sock = conn->socket; conn->socket = nullptr; if (sock != nullptr) { // The socket's deleteLater is deferred, so its bytesWritten/disconnected // lambdas could still fire after `delete conn` below. Sever them now: // this object is the context of every pump/fan connection made in // acceptData, so disconnecting it disarms all queued pump re-arms and // drain callbacks for this socket (the deleteLater connection targets // the socket itself, which still needs to run). sock->disconnect(this); sock->deleteLater(); } m_dataConns.removeOne(conn); delete conn;}
void DaemonCore::dropPaneConnections(const QString& paneId){ // Snapshot the sockets rather than iterating the live list: close() can // re-enter onDataDisconnect, which removes from m_dataConns and deletes the // DataConn immediately. The QWebSocket* itself stays valid because its // teardown is deleteLater(), so the snapshot loop never dereferences freed // memory. QVector<QWebSocket*> doomed; for (const DataConn* conn : std::as_const(m_dataConns)) { if (conn->paneId == paneId && conn->socket != nullptr) { doomed.append(conn->socket); } } for (QWebSocket* sock : std::as_const(doomed)) { sock->close(QWebSocketProtocol::CloseCodeNormal, QStringLiteral("pane-killed")); }}
void DaemonCore::fanOut(const QString& paneId){ PaneTable::Pane* pane = m_panes->find(paneId); if (pane == nullptr) { return; } // Snapshot the connection pointers, then re-validate membership before each // dereference. The backlog branch calls closeWith(), whose disconnected() // handler runs onDataDisconnect -> m_dataConns.removeOne(conn) + `delete // conn` (the delete is IMMEDIATE, unlike the socket's deleteLater). Iterating // the live container by reference would then walk freed memory and a // stale iterator. removeOne precedes delete, so contains() is a valid // liveness test: a conn evicted mid-loop is simply skipped. const QVector<DataConn*> conns = m_dataConns; for (DataConn* conn : conns) { if (!m_dataConns.contains(conn)) { continue; // Evicted by an earlier closeWith() in this same loop. } if (conn->paneId != paneId || conn->socket == nullptr) { continue; } // Mutually exclusive with the cold-replay pump: a connection still // replaying its retained log must not ALSO get fanOut's readFrom( // sentUpTo) — that would duplicate/reorder records against the pump's // batches. The pump advances sentUpTo and clears `replaying` once it has // queued everything below the attach-time tailSeq; from then on fanOut // owns the live path. if (conn->replaying) { continue; } QWebSocket* sock = conn->socket; // Backpressure: the log is the source of truth; a lagging socket is a // hint. Over the cap, close 1009 so the client re-attaches from its // last seq. Lossless, because the log still holds the records. // // MEASURED WEDGE (14-byte line-buffered flood): plain close() here is NOT // enough. close() must first FLUSH the queued 4 MB before it can send the // close frame, and a peer that never reads never drains that buffer, so // disconnected() never fires, the DataConn stays live, and EVERY later // append re-enters this branch and re-closes the same socket. Counted on // a 6 s flood: 553,000 closeWith(1009) calls on ONE connection; each // close() allocates Qt timer state, so the heap grew ~700 MB while the // process sat in ~QTimer teardown. The fix detaches the conn first (full // removal from m_dataConns, so no later append can even reach this // branch for that socket), then abort() drops the write buffer and lets // the close frame go out immediately. if (sock->bytesToWrite() > kSocketBacklogCap) { // MEASURED WEDGE (14-byte line-buffered flood): plain close() here is // NOT enough. close() must first FLUSH the queued 4 MB before it can // send the close frame, and a peer that never reads never drains that // buffer, so disconnected() never fires, the DataConn stays live, and // EVERY later append re-enters this branch and re-closes the same // socket. Counted on a 6 s flood: 553,000 closeWith(1009) calls on ONE // connection; each close() allocates Qt timer state, so the heap grew // ~700 MB while the process sat in ~QTimer teardown. // // Order is load-bearing: detachConn FIRST. abort() can emit // disconnected() synchronously, and that handler calls connFor(socket) // -> onDataDisconnect -> detachConn -> delete conn. Detaching before // the abort means the signal path finds the conn already removed // (connFor returns nullptr) and there is no use-after-free; the // attachedCount decrement happens exactly once, here. PaneTable::Pane* p = m_panes->find(conn->paneId); if (p != nullptr && p->attachedCount > 0) { p->attachedCount--; } detachConn(conn); // abort() drops the write buffer so the close frame is not queued // behind 4 MB of undeliverable data; the socket object itself is // deleteLater'd (by detachConn) and dies on the next event-loop turn. sock->abort(); closeWith(sock, static_cast<quint16>(QWebSocketProtocol::CloseCodeTooMuchData), QStringLiteral("backlog")); continue; } // Bound what THIS call enqueues: the cap check above only gates the next // call, not the loop below. Use the bounded read (kReplayBatchRecords per // turn) and stop early on the byte bound, so one fanOut call can never // queue an unbounded backlog behind a single passing cap check. Anything // left pending is delivered by the next append's fanOut (the flood case: // the next append arrives within milliseconds) or on the next drain. for (;;) { const QVector<ByteLog::Record> pending = pane->log.readFrom(conn->sentUpTo, kReplayBatchRecords); if (pending.isEmpty()) { break; } qint64 batchBytes = 0; bool byteBound = false; for (const ByteLog::Record& r : pending) { xeno::wire::Header h; h.streamId = conn->streamId; h.seq = r.seq; h.tag = r.tag; QByteArray out; xeno::wire::encode(out, h, r.payload); batchBytes += out.size(); sock->sendBinaryMessage(out); conn->sentUpTo = r.seq + 1; if (batchBytes + kMaxIncomingFrame > kSocketBacklogCap || sock->bytesToWrite() > kSocketBacklogCap) { byteBound = true; break; } } if (byteBound || sock->bytesToWrite() > kSocketBacklogCap || std::cmp_less(pending.size(), kReplayBatchRecords)) { break; } } }}