/* * antpeer.cpp — Peer library for Antheos bus communication. * * Pure C++17. Single compilation unit. * Uses C++ APIs: membus::Bus, sockbus::Bus, antheos::Context. * * Copyright (c) 2025-2026 Are Bjørby * SPDX-License-Identifier: MIT */ #ifndef _GNU_SOURCE #define _GNU_SOURCE #endif #include "antpeer.hpp" #include "Ed25519Utils.hpp" #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include namespace antpeer { /* ── Constants ─────────────────────────────────────────────────────── */ namespace { constexpr size_t READ_BUF = 65536; constexpr size_t RESP_BUF = 65536; constexpr int SLEEP_US = 1000; constexpr int ESTABLISH_TIMEOUT_MS = 50; /* Set to true inside worker threads so reply() routes to response queue. */ thread_local bool g_is_worker_thread = false; /* BLOB tail data for the current request (per worker thread). */ thread_local std::vector g_worker_tail; } // anonymous namespace namespace { /* ── BID generation ────────────────────────────────────────────────── */ std::string generate_bid(size_t len = antheos::id::BID_MIN_LEN) { auto bid = antheos::id::bid_generate(len); if (!bid) throw std::runtime_error("BID generation failed"); return *bid; } /* ═══════════════════════════════════════════════════════════════════════ * Bus implementations * ═══════════════════════════════════════════════════════════════════════ */ /* ── Stale SHM detection (membus only) ─────────────────────────────── */ bool membus_is_mapped(const std::string& bus_name) { std::string needle = "/dev/shm/" + bus_name + "_shm"; struct stat st; if (stat(needle.c_str(), &st) != 0) return false; struct timespec now; clock_gettime(CLOCK_REALTIME, &now); double age_s = static_cast(now.tv_sec - st.st_ctim.tv_sec) + static_cast(now.tv_nsec - st.st_ctim.tv_nsec) / 1e9; if (age_s < 5.0) return true; DIR* proc = opendir("/proc"); if (!proc) return true; struct dirent* entry; while ((entry = readdir(proc)) != nullptr) { if (entry->d_name[0] < '0' || entry->d_name[0] > '9') continue; char maps_path[280]; std::snprintf(maps_path, sizeof(maps_path), "/proc/%s/maps", entry->d_name); FILE* f = fopen(maps_path, "r"); if (!f) continue; char line[512]; while (fgets(line, sizeof(line), f)) { if (strstr(line, needle.c_str())) { fclose(f); closedir(proc); return true; } } fclose(f); } closedir(proc); return false; } /* ── MemBus ────────────────────────────────────────────────────────── */ class MemBusBus : public Bus { std::optional bus_; std::string bus_name_; size_t buf_size_; bool ensure_and_open() { std::string shm_path = "/dev/shm/" + bus_name_ + "_shm"; struct stat st; if (stat(shm_path.c_str(), &st) == 0 && !membus_is_mapped(bus_name_)) try { membus::destroy(bus_name_); } catch (...) {} mode_t old = umask(0007); bool create_failed = false; try { membus::create(bus_name_, buf_size_); } catch (const std::system_error& e) { if (e.code().value() != EEXIST) { create_failed = true; std::fprintf(stderr, "[bus] membus::create('%s') failed: %s\n", bus_name_.c_str(), e.what()); } } umask(old); if (create_failed) return false; try { bus_.emplace(bus_name_); } catch (const std::system_error& e) { std::fprintf(stderr, "[bus] membus::Bus('%s') failed: %s\n", bus_name_.c_str(), e.what()); return false; } return true; } public: MemBusBus(std::string_view name, size_t size) : bus_name_(name) , buf_size_(size == 0 ? MEMBUS_DEFAULT_SIZE : size) { if (bus_name_.empty()) { errno = EINVAL; throw std::system_error(EINVAL, std::system_category()); } if (!ensure_and_open()) throw std::system_error(errno ? errno : EIO, std::system_category()); } ~MemBusBus() override = default; ssize_t write(const void* data, size_t len) override { if (!bus_) return -1; try { return static_cast( bus_->write(static_cast(data), len)); } catch (...) { return -1; } } ssize_t read(void* buffer, size_t len) override { if (!bus_) return -1; try { return static_cast( bus_->read(static_cast(buffer), len)); } catch (...) { return -1; } } bool reopen() override { bus_.reset(); return ensure_and_open(); } bool is_open() const override { return bus_.has_value(); } std::string_view name() const override { return bus_name_; } }; /* ── SockBus ───────────────────────────────────────────────────────── */ class SockBusBus : public Bus { std::optional bus_; std::string bus_name_; bool create_and_open() { try { sockbus::create(bus_name_); } catch (const std::system_error& e) { if (e.code().value() != EEXIST && e.code().value() != EADDRINUSE && e.code().value() != EADDRNOTAVAIL) { std::fprintf(stderr, "[bus] sockbus::create('%s') failed: %s\n", bus_name_.c_str(), e.what()); return false; } } try { bus_.emplace(bus_name_); } catch (const std::system_error& e) { std::fprintf(stderr, "[bus] sockbus::Bus('%s') failed: %s\n", bus_name_.c_str(), e.what()); return false; } return true; } public: explicit SockBusBus(std::string_view name) : bus_name_(name) { if (bus_name_.empty()) { errno = EINVAL; throw std::system_error(EINVAL, std::system_category()); } if (!create_and_open()) throw std::system_error(errno ? errno : EIO, std::system_category()); } ~SockBusBus() override = default; ssize_t write(const void* data, size_t len) override { if (!bus_) return -1; try { return static_cast( bus_->write(static_cast(data), len)); } catch (...) { return -1; } } ssize_t read(void* buffer, size_t len) override { if (!bus_) return -1; try { return static_cast( bus_->read(static_cast(buffer), len)); } catch (...) { return -1; } } bool reopen() override { bus_.reset(); return create_and_open(); } bool is_open() const override { return bus_.has_value(); } std::string_view name() const override { return bus_name_; } }; class BleBusBus : public Bus { std::optional bus_; std::string bus_name_; bool create_and_open() { try { blebus::create(bus_name_); } catch (const std::system_error& e) { if (e.code().value() != EEXIST && e.code().value() != EADDRINUSE && e.code().value() != EADDRNOTAVAIL) { std::fprintf(stderr, "[bus] blebus::create('%s') failed: %s\n", bus_name_.c_str(), e.what()); return false; } } try { bus_.emplace(bus_name_); } catch (const std::system_error& e) { std::fprintf(stderr, "[bus] blebus::Bus('%s') failed: %s\n", bus_name_.c_str(), e.what()); return false; } return true; } public: explicit BleBusBus(std::string_view name) : bus_name_(name) { if (bus_name_.empty()) { errno = EINVAL; throw std::system_error(EINVAL, std::system_category()); } if (!create_and_open()) throw std::system_error(errno ? errno : EIO, std::system_category()); } ~BleBusBus() override = default; ssize_t write(const void* data, size_t len) override { if (!bus_) return -1; try { return static_cast( bus_->write(static_cast(data), len)); } catch (...) { return -1; } } ssize_t read(void* buffer, size_t len) override { if (!bus_) return -1; try { return static_cast( bus_->read(static_cast(buffer), len)); } catch (...) { return -1; } } bool reopen() override { bus_.reset(); return create_and_open(); } bool is_open() const override { return bus_.has_value(); } std::string_view name() const override { return bus_name_; } }; } // anonymous namespace /* ── Bus factories ─────────────────────────────────────────────────── */ std::unique_ptr membus_open(std::string_view name, size_t size) { return std::make_unique(name, size); } std::unique_ptr sockbus_open(std::string_view name) { return std::make_unique(name); } std::unique_ptr blebus_open(std::string_view name) { return std::make_unique(name); } /* ── Named-object removal ──────────────────────────────────────────── */ /* * Thin forwards, and thin on purpose. The `*_open` calls above wrap real * behaviour — stale-segment reclaim, create-if-missing, EEXIST tolerance — * because opening is where the transports differ. Removing is where they * agree: each takes a name and unlinks it, and each treats a name that is * not there as success. Adding logic here would be inventing a difference * the layer underneath does not have. * * One asymmetry is worth knowing about rather than hiding. MemBusBus's * `ensure_and_open` already calls `membus::destroy` when the shm file * exists and no process has it mapped — a stale-segment reclaim, not a * remove-half, and the only one of the three transports that has such a * path. So a caller who never calls `membus_destroy` may still find a * membus bus quietly recreated under them, where a sockbus or blebus * broker in the same state simply stays stale. */ void membus_destroy(std::string_view name) { membus::destroy(name); } void sockbus_destroy(std::string_view name) { sockbus::destroy(name); } void blebus_destroy(std::string_view name) { blebus::destroy(name); } /* ═══════════════════════════════════════════════════════════════════════ * FrameExtractor * ═══════════════════════════════════════════════════════════════════════ */ namespace { size_t decode_blob_size(antheos::wire::Radix radix, const uint8_t* body, size_t len) { if (len == 0) return 0; std::string s(reinterpret_cast(body), len); try { switch (radix) { case antheos::wire::Radix::Binary: return std::stoull(s, nullptr, 2); case antheos::wire::Radix::Octal: return std::stoull(s, nullptr, 8); case antheos::wire::Radix::Decimal: return std::stoull(s, nullptr, 10); case antheos::wire::Radix::Hex: return std::stoull(s, nullptr, 16); case antheos::wire::Radix::Base32: { size_t val = 0; for (char c : s) { const char* p = std::strchr( antheos::id::BASE32_ALPHABET, c); if (!p) return 0; val = val * 32 + static_cast(p - antheos::id::BASE32_ALPHABET); } return val; } default: return 0; } } catch (...) { return 0; } } } // anonymous namespace struct FrameExtractor::Impl { antheos::Parser parser; std::vector frame_buf; FrameCb frame_cb; size_t blob_tail_total = 0; size_t tail_remaining = 0; bool in_frame = false; bool in_tail = false; bool head_done = false; size_t total_frames_ = 0; size_t parse_errors_ = 0; Impl() { parser.on_word( [this](antheos::wire::WordType type, antheos::wire::Radix radix, antheos::wire::Unit /*unit*/, const uint8_t* body, size_t len) { if (type == antheos::wire::WordType::Blob) { blob_tail_total += decode_blob_size(radix, body, len); } }); parser.on_message([this](size_t /*word_count*/) { head_done = true; }); } void emit_frame() { if (frame_cb && !frame_buf.empty()) frame_cb(frame_buf.data(), frame_buf.size()); frame_buf.clear(); in_frame = false; in_tail = false; head_done = false; blob_tail_total = 0; tail_remaining = 0; total_frames_++; } void feed(const uint8_t* data, size_t len) { for (size_t i = 0; i < len; i++) { uint8_t byte = data[i]; /* Tail consumption — bypass parser entirely */ if (in_tail) { frame_buf.push_back(byte); tail_remaining--; if (tail_remaining == 0) emit_frame(); continue; } auto prev = parser.state(); head_done = false; auto next = parser.feed(byte); /* SOM detected — start new frame */ if (prev == antheos::ParseState::WaitSom && next == antheos::ParseState::WaitSow) { frame_buf.clear(); blob_tail_total = 0; in_frame = true; } if (in_frame) frame_buf.push_back(byte); /* Head complete (parser saw EOM, no parser-level tail) */ if (head_done) { if (blob_tail_total > 0) { tail_remaining = blob_tail_total; in_tail = true; } else { emit_frame(); } continue; } /* Parse error — discard partial frame */ if (next == antheos::ParseState::Error) { /* Parser auto-recovers on next SOM, but if this byte IS a SOM (error recovery), start a new frame */ frame_buf.clear(); in_frame = false; blob_tail_total = 0; parse_errors_++; } } } void reset() { parser.reset(); frame_buf.clear(); blob_tail_total = 0; tail_remaining = 0; in_frame = false; in_tail = false; head_done = false; } }; FrameExtractor::FrameExtractor() : impl_(std::make_unique()) {} FrameExtractor::~FrameExtractor() = default; FrameExtractor::FrameExtractor(FrameExtractor&&) noexcept = default; FrameExtractor& FrameExtractor::operator=(FrameExtractor&&) noexcept = default; void FrameExtractor::on_frame(FrameCb cb) { impl_->frame_cb = std::move(cb); } void FrameExtractor::feed(const uint8_t* data, size_t len) { impl_->feed(data, len); } void FrameExtractor::reset() { impl_->reset(); } size_t FrameExtractor::total_frames() const { return impl_->total_frames_; } size_t FrameExtractor::parse_errors() const { return impl_->parse_errors_; } /* ═══════════════════════════════════════════════════════════════════════ * Server * ═══════════════════════════════════════════════════════════════════════ */ struct Server::Impl { Server* owner = nullptr; Bus* bus = nullptr; std::unique_ptr ctx_; Server::RequestFn request_fn; Server::SessionFn session_fn; Server::FinishFn finish_fn; Server::NotifyFn notify_fn; Server::AcceptFn accept_fn; Server::VerifyFn verify_fn; Server::AuthFn auth_fn; /* Auth (Z-verb, Level 2) */ bool auth_required = false; std::vector trusted_keys; std::unordered_map pending_auth; /* BID → nonce_hex */ std::unordered_set authenticated_peers; /* Relay routes for multi-hop auth */ std::unordered_map relay_routes; /* target BID → Route */ std::string service_name; std::string log_prefix; std::string self_bid; /* Session tracking */ struct SessionEntry { std::string sid; int slot; }; std::vector inbound; std::vector known_sids; std::vector accepted_peers; /* Peer identity (Verify protocol) */ std::unordered_map peer_identities; std::deque pending_verify; std::unordered_map sid_owner; /* SID→BID */ struct SessionRecord { std::string sid; uint32_t last_mid; time_t last_activity; }; std::vector registry; int session_expiry_s = 60; std::vector additional_offers; std::string last_trace_id_; bool bid_established = false; bool conflict_received = false; std::atomic running{false}; std::mutex session_mutex; /* ── Threaded dispatch (opt-in via set_dispatch_threads) ─────── */ size_t dispatch_threads_ = 0; struct WorkItem { std::string sid; std::string command; std::string trace_id; std::vector tail; }; struct RespItem { std::string sid; std::string body; std::vector blob; }; std::mutex work_mu_; std::condition_variable work_cv_; std::deque work_queue_; std::mutex resp_mu_; std::deque resp_queue_; std::vector workers_; void start_workers() { for (size_t i = 0; i < dispatch_threads_; i++) { workers_.emplace_back([this] { worker_loop(); }); } } void stop_workers() { /* Signal all workers to wake and check running flag */ work_cv_.notify_all(); for (auto& t : workers_) { if (t.joinable()) t.join(); } workers_.clear(); } void worker_loop() { g_is_worker_thread = true; while (running.load()) { WorkItem item; { std::unique_lock lock(work_mu_); work_cv_.wait_for(lock, std::chrono::milliseconds(50), [this] { return !work_queue_.empty() || !running.load(); }); if (!running.load() && work_queue_.empty()) break; if (work_queue_.empty()) continue; item = std::move(work_queue_.front()); work_queue_.pop_front(); } /* Thread-local state for the handler to retrieve via accessors */ g_worker_tail = std::move(item.tail); if (request_fn) { request_fn(*owner, item.sid, item.command); } g_worker_tail.clear(); } g_is_worker_thread = false; } /* Drain response queue — called from poll loop thread only. */ void drain_responses() { std::deque batch; { std::lock_guard lock(resp_mu_); if (resp_queue_.empty()) return; batch.swap(resp_queue_); } for (auto& r : batch) { internal_reply_and_finish(r.sid, r.body, r.blob.data(), r.blob.size()); } } /* Reply + finish on the poll loop thread (touches ctx_). */ void internal_reply_and_finish(const std::string& sid, const std::string& body, const uint8_t* blob = nullptr, size_t blob_len = 0) { std::lock_guard lock(session_mutex); int idx = find_inbound(sid); if (idx < 0) return; int slot = inbound[idx].slot; auto state = ctx_->session_state(slot); if (state != antheos::SessionState::Active) { remove_inbound(idx); remove_known_sid(sid); return; } auto slot_sid = ctx_->session_sid(slot); if (slot_sid.empty() || sid != slot_sid) { remove_inbound(idx); remove_known_sid(sid); return; } /* Reply (text-only or with BLOB tail) */ std::optional frame; if (blob && blob_len > 0) frame = ctx_->session_call_blob(slot, {}, body, blob, blob_len); else frame = ctx_->session_call(slot, {}, body); if (frame) bus_send(*frame); /* Finish */ remove_inbound(idx); remove_known_sid(sid); auto close_frame = ctx_->session_close(slot); if (close_frame) bus_send(*close_frame); } /* ── Helpers ───────────────────────────────────────────────────── */ bool bus_send(const antheos::Frame& frame) { if (!bus || frame.empty()) return false; return bus->write(frame.data(), frame.size()) > 0; } int find_inbound(std::string_view sid) const { for (size_t i = 0; i < inbound.size(); i++) if (inbound[i].sid == sid) return static_cast(i); return -1; } void remove_inbound(int idx) { if (idx < 0 || static_cast(idx) >= inbound.size()) return; inbound[idx] = std::move(inbound.back()); inbound.pop_back(); } bool has_known_sid(std::string_view sid) const { for (auto& s : known_sids) if (s == sid) return true; return false; } void add_known_sid(std::string_view sid) { if (known_sids.size() >= static_cast(MAX_SESSIONS * 2)) return; known_sids.emplace_back(sid); } void remove_known_sid(std::string_view sid) { for (size_t i = 0; i < known_sids.size(); i++) { if (known_sids[i] == sid) { known_sids[i] = std::move(known_sids.back()); known_sids.pop_back(); return; } } } int find_registry(std::string_view sid) const { for (size_t i = 0; i < registry.size(); i++) if (registry[i].sid == sid) return static_cast(i); return -1; } void remove_registry(int idx) { if (idx < 0 || static_cast(idx) >= registry.size()) return; registry[idx] = std::move(registry.back()); registry.pop_back(); } void sweep_expired() { time_t now = time(nullptr); for (int i = static_cast(registry.size()) - 1; i >= 0; i--) { if (now - registry[i].last_activity > session_expiry_s) remove_registry(i); } } void upsert_registry(std::string_view sid, uint32_t mid) { int idx = find_registry(sid); if (idx >= 0) { registry[idx].last_mid = mid; registry[idx].last_activity = time(nullptr); } else if (registry.size() < static_cast(MAX_SESSIONS * 2)) { registry.push_back({std::string(sid), mid, time(nullptr)}); } if (registry.size() > static_cast(MAX_SESSIONS)) sweep_expired(); } /* Record SID→BID from O: ownership header. session_mutex must be held. */ void register_sid_owner(std::string_view sid, std::string_view bid) { if (sid.empty() || bid.empty()) return; sid_owner[std::string(sid)] = std::string(bid); } void setup_callbacks() { ctx_->on_message( [this](std::string_view verb, std::string_view sid, std::string_view target_bid, uint32_t mid, std::string_view body) { handle_message(verb, sid, target_bid, mid, body); }); ctx_->on_event( [this](std::string_view verb, std::string_view id, std::string_view id2, std::string_view detail) { handle_event(verb, id, id2, detail); }); ctx_->on_offer([](std::string_view, std::string_view) {}); ctx_->on_relay( [this](char ref, std::string_view id, std::string_view id2, std::string_view body, uint32_t /*index*/, std::string_view /*path*/) { handle_relay(ref, id, id2, body); }); } void handle_relay(char ref, std::string_view id, std::string_view id2, std::string_view body) { if (ref != 'Z') return; /* Relayed Z-response: treat as direct Z-response */ handle_event("Z", id, id2, body); } bool establish_with_conflict(std::string_view oid, std::string_view did, std::string_view iid) { size_t bid_len = antheos::id::BID_MIN_LEN; bid_established = false; while (bid_len <= antheos::id::BID_MAX_LEN) { try { ctx_ = std::make_unique( oid, did, iid, generate_bid(bid_len)); } catch (const std::exception& e) { std::fprintf(stderr, "%s antheos::Context init failed: %s\n", log_prefix.c_str(), e.what()); return false; } setup_callbacks(); self_bid = std::string(ctx_->bid()); conflict_received = false; auto frame = ctx_->establish(); if (!frame || !bus_send(*frame)) return false; auto deadline = std::chrono::steady_clock::now() + std::chrono::milliseconds(ESTABLISH_TIMEOUT_MS); uint8_t buf[READ_BUF]; while (std::chrono::steady_clock::now() < deadline) { ssize_t n = bus->read(buf, sizeof(buf)); if (n > 0) ctx_->feed(buf, static_cast(n)); else usleep(SLEEP_US); if (conflict_received) break; } if (!conflict_received) { bid_established = true; return true; } std::fprintf(stderr, "%s BID conflict at length %zu, retrying\n", log_prefix.c_str(), bid_len); bid_len++; } std::fprintf(stderr, "%s BID_OVERFLOW — all lengths exhausted\n", log_prefix.c_str()); auto frame = antheos::bus::exception("BID_OVERFLOW"); if (frame) bus_send(*frame); ctx_.reset(); self_bid.clear(); return false; } bool offer_service() { auto frame = ctx_->offer(service_name); if (!frame) return false; return bus_send(*frame); } /* ── Antheos callback handlers ─────────────────────────────────── */ void handle_message(std::string_view verb, std::string_view sid, std::string_view target_bid, uint32_t mid, std::string_view body) { if (verb.empty()) return; /* N-verb broadcast notifications */ if (verb[0] == 'N' && !body.empty()) { if (notify_fn) notify_fn(body); return; } if (verb[0] != 'K' || body.empty() || sid.empty()) return; if (mid % 2 == 0) return; if (!target_bid.empty() && !self_bid.empty() && self_bid != target_bid) return; { std::lock_guard lock(session_mutex); if (has_known_sid(sid)) { /* continuation */ } else if (accepted_peers.empty()) { return; } else { add_known_sid(sid); } upsert_registry(sid, mid); } /* Extract trace ID header if present: "T:\n" */ std::string_view actual_body = body; if (body.size() > 2 && body[0] == 'T' && body[1] == ':') { auto nl = body.find('\n'); if (nl != std::string_view::npos) { last_trace_id_ = std::string(body.substr(2, nl - 2)); actual_body = body.substr(nl + 1); } } else { last_trace_id_.clear(); } /* Extract ownership header if present: "O:\n" */ if (actual_body.size() > 2 && actual_body[0] == 'O' && actual_body[1] == ':') { auto nl = actual_body.find('\n'); if (nl != std::string_view::npos) { auto owner_bid = actual_body.substr(2, nl - 2); { std::lock_guard lock(session_mutex); register_sid_owner(sid, owner_bid); } actual_body = actual_body.substr(nl + 1); } } if (session_fn) { int slot = ctx_->session_accept(sid, mid); if (slot < 0) { std::lock_guard lock(session_mutex); remove_known_sid(sid); return; } { std::lock_guard lock(session_mutex); if (inbound.size() < static_cast(MAX_SESSIONS)) inbound.push_back({std::string(sid), slot}); } session_fn(*owner, sid, mid, actual_body); } else if (request_fn) { int slot = ctx_->session_accept(sid, mid); if (slot < 0) { std::lock_guard lock(session_mutex); remove_known_sid(sid); return; } { std::lock_guard lock(session_mutex); if (inbound.size() < static_cast(MAX_SESSIONS)) inbound.push_back({std::string(sid), slot}); } if (dispatch_threads_ > 0) { /* Threaded: push to work queue, poll loop drains responses */ auto [tptr, tlen] = ctx_->last_tail(); std::lock_guard lock(work_mu_); work_queue_.push_back({std::string(sid), std::string(actual_body), last_trace_id_, std::vector(tptr, tptr + tlen)}); work_cv_.notify_one(); } else { /* Inline: existing single-threaded path */ request_fn(*owner, sid, actual_body); owner->finish(sid); } } } void handle_event(std::string_view verb, std::string_view id, std::string_view id2, std::string_view detail) { if (verb.empty()) return; switch (verb[0]) { case 'A': if (!id.empty() && !id2.empty() && !self_bid.empty() && self_bid == id) { { std::lock_guard lock(session_mutex); if (accepted_peers.size() < static_cast(MAX_PEERS)) accepted_peers.emplace_back(id2); } std::fprintf(stderr, "%s Accept from peer %.*s\n", log_prefix.c_str(), static_cast(id2.size()), id2.data()); if (accept_fn) accept_fn(id2); /* Auto-send Verify to learn peer identity */ auto vframe = ctx_->verify(id2); if (vframe) { bus_send(*vframe); std::lock_guard lock(session_mutex); pending_verify.emplace_back(id2); } } break; case 'V': if (!id.empty() && !id2.empty()) { /* Verify response: id=OID, id2=DID, detail=IID */ std::string peer_bid; { std::lock_guard lock(session_mutex); if (!pending_verify.empty()) { peer_bid = std::move(pending_verify.front()); pending_verify.pop_front(); } } if (!peer_bid.empty()) { PeerIdentity pi; pi.bid = peer_bid; pi.oid = std::string(id); pi.did = std::string(id2); pi.iid = std::string(detail); std::fprintf(stderr, "%s Verified peer %s: %s:%s:%s\n", log_prefix.c_str(), peer_bid.c_str(), pi.oid.c_str(), pi.did.c_str(), pi.iid.c_str()); { std::lock_guard lock(session_mutex); peer_identities[peer_bid] = pi; } if (verify_fn) verify_fn(pi); /* Auth: send Z challenge after identity is known */ if (auth_required && !peer_bid.empty()) { std::string challenge_nonce = crypto::nonce(32); std::optional zframe; auto rit = relay_routes.find(peer_bid); if (rit != relay_routes.end()) { zframe = ctx_->relay_auth_challenge( peer_bid, challenge_nonce, rit->second.my_index, rit->second.path); } else { zframe = ctx_->auth_challenge(peer_bid, challenge_nonce); } if (zframe) { bus_send(*zframe); std::lock_guard lock2(session_mutex); pending_auth[peer_bid] = std::move(challenge_nonce); } } } } else if (!id.empty() && id2.empty()) { /* Verify request: id=target BID — respond if it's us */ if (!self_bid.empty() && self_bid == id) { auto vr = ctx_->verify_response(); if (vr) bus_send(*vr); } } break; case 'Z': /* Auth response: id=target(our BID), id2=key_id, detail=sig_hex */ if (!id.empty() && !id2.empty() && !detail.empty() && self_bid == id && auth_required) { std::string key_id_str(id2); std::string sig_str(detail); /* Find which pending peer this key_id belongs to */ std::string auth_bid; std::string auth_nonce; PeerIdentity auth_pi; { std::lock_guard lock(session_mutex); auto* tk = crypto::find_trusted_key(trusted_keys, {}, key_id_str); if (tk) { for (auto& [bid, nonce_val] : pending_auth) { auto pit = peer_identities.find(bid); if (pit != peer_identities.end() && pit->second.oid == tk->oid) { auth_bid = bid; auth_nonce = nonce_val; auth_pi = pit->second; break; } } } } if (!auth_bid.empty()) { auto* tk = crypto::find_trusted_key(trusted_keys, auth_pi.oid, key_id_str); bool ok = tk && crypto::verify(auth_nonce, sig_str, tk->public_key); { std::lock_guard lock(session_mutex); pending_auth.erase(auth_bid); if (ok) authenticated_peers.insert(auth_bid); } std::fprintf(stderr, "%s Auth %s peer %s (%s)\n", log_prefix.c_str(), ok ? "OK" : "FAILED", auth_bid.c_str(), auth_pi.oid.c_str()); if (auth_fn) auth_fn(auth_pi, ok); if (!ok) { auto xf = ctx_->exception("AUTH_FAILED"); if (xf) bus_send(*xf); } } } break; case 'F': if (!id.empty()) { int slot = -1; { std::lock_guard lock(session_mutex); int idx = find_inbound(id); if (idx >= 0) { slot = inbound[idx].slot; remove_inbound(idx); } remove_known_sid(id); int ri = find_registry(id); if (ri >= 0) remove_registry(ri); } if (slot >= 0) ctx_->session_close(slot); if (finish_fn) finish_fn(id); } break; case 'L': if (!id.empty() && id2.empty()) { std::lock_guard lock(session_mutex); int ri = find_registry(id); if (ri >= 0) { time_t age = time(nullptr) - registry[ri].last_activity; if (age < session_expiry_s) { auto frame = antheos::session::locate_response( id, self_bid); if (frame) bus_send(*frame); } else { remove_registry(ri); } } } break; case 'U': if (!id.empty()) { std::lock_guard lock(session_mutex); int ri = find_registry(id); if (ri >= 0) { time_t age = time(nullptr) - registry[ri].last_activity; if (age < session_expiry_s) { registry[ri].last_activity = time(nullptr); auto frame = antheos::session::status_response( id, {}, registry[ri].last_mid, "ACTIVE"); if (frame) bus_send(*frame); std::fprintf(stderr, "%s session resumed: %.*s\n", log_prefix.c_str(), static_cast(id.size()), id.data()); } else { remove_registry(ri); auto frame = antheos::bus::exception( "SESSION_EXPIRED"); if (frame) bus_send(*frame); } } else { auto frame = antheos::bus::exception( "SESSION_NOT_FOUND"); if (frame) bus_send(*frame); } } break; case 'Q': if (!detail.empty()) { if (service_name == detail) { offer_service(); break; } for (auto& ao : additional_offers) { if (ao == detail) { auto frame = ctx_->offer(ao); if (frame) bus_send(*frame); break; } } } break; case 'P': if (!id.empty()) { auto our_bid = ctx_->bid(); if (our_bid == id) { auto frame = ctx_->acknowledge(id); if (frame) bus_send(*frame); } } break; case 'E': if (!id.empty() && !self_bid.empty() && self_bid == id && bid_established) { auto frame = antheos::bus::conflict(self_bid); if (frame) bus_send(*frame); } break; case 'C': if (!id.empty() && !self_bid.empty() && self_bid == id) conflict_received = true; break; case 'X': std::fprintf(stderr, "%s Exception: %.*s %.*s\n", log_prefix.c_str(), static_cast(id.size()), id.data(), static_cast(detail.size()), detail.data()); if (notify_fn) { std::string reason = "!X"; if (!detail.empty()) { reason += ' '; reason.append(detail.data(), detail.size()); } notify_fn(reason); } break; default: break; } } }; Server::Server() : impl_(std::make_unique()) { impl_->owner = this; } Server::~Server() { if (impl_) { impl_->running.store(false); impl_->work_cv_.notify_all(); impl_->stop_workers(); impl_->ctx_.reset(); } } int Server::init(Bus& bus, std::string_view oid, std::string_view did, std::string_view iid, std::string_view service_name) { if (oid.empty() || did.empty() || iid.empty() || service_name.empty() || !bus.is_open()) return -1; impl_->bus = &bus; impl_->service_name = service_name; impl_->log_prefix = "[" + std::string(service_name) + "]"; if (!impl_->establish_with_conflict(oid, did, iid)) { impl_->bus = nullptr; return -1; } std::fprintf(stderr, "%s ready on bus '%.*s', service '%.*s'\n", impl_->log_prefix.c_str(), static_cast(bus.name().size()), bus.name().data(), static_cast(service_name.size()), service_name.data()); return 0; } void Server::on_request(RequestFn fn) { impl_->request_fn = std::move(fn); } void Server::on_session(SessionFn fn) { impl_->session_fn = std::move(fn); } void Server::on_finish(FinishFn fn) { impl_->finish_fn = std::move(fn); } void Server::on_notify(NotifyFn fn) { impl_->notify_fn = std::move(fn); } void Server::set_session_expiry(int seconds) { impl_->session_expiry_s = seconds; } void Server::set_dispatch_threads(size_t n) { impl_->dispatch_threads_ = n; } void Server::run() { impl_->running.store(true); if (impl_->dispatch_threads_ > 0) impl_->start_workers(); uint8_t buf[READ_BUF]; while (impl_->running.load()) { ssize_t n = impl_->bus->read(buf, sizeof(buf)); if (n > 0) impl_->ctx_->feed(buf, static_cast(n)); /* Drain worker responses (no-op if no threaded dispatch). */ if (impl_->dispatch_threads_ > 0) impl_->drain_responses(); if (n <= 0) usleep(SLEEP_US); } if (impl_->dispatch_threads_ > 0) impl_->stop_workers(); } void Server::stop() { impl_->running.store(false); impl_->work_cv_.notify_all(); } int Server::reply(std::string_view sid, std::string_view body) { if (!impl_->ctx_) return -1; /* Worker thread: queue for poll loop to send (avoids ctx_ races). */ if (g_is_worker_thread) { std::lock_guard lock(impl_->resp_mu_); impl_->resp_queue_.push_back({std::string(sid), std::string(body), {}}); return 0; } /* Poll thread (inline mode): send directly. */ std::lock_guard lock(impl_->session_mutex); int idx = impl_->find_inbound(sid); if (idx < 0) return -1; int slot = impl_->inbound[idx].slot; auto state = impl_->ctx_->session_state(slot); if (state != antheos::SessionState::Active) { impl_->remove_inbound(idx); return -1; } auto slot_sid = impl_->ctx_->session_sid(slot); if (slot_sid.empty() || sid != slot_sid) { impl_->remove_inbound(idx); return -1; } auto frame = impl_->ctx_->session_call(slot, {}, body); if (!frame) return -1; return impl_->bus_send(*frame) ? 0 : -1; } int Server::reply_with_blob(std::string_view sid, std::string_view body, const uint8_t* blob, size_t blob_len) { if (!impl_->ctx_) return -1; /* Worker thread: queue for poll loop to send. */ if (g_is_worker_thread) { std::lock_guard lock(impl_->resp_mu_); impl_->resp_queue_.push_back({std::string(sid), std::string(body), std::vector(blob, blob + blob_len)}); return 0; } /* Poll thread (inline mode): send directly. */ std::lock_guard lock(impl_->session_mutex); int idx = impl_->find_inbound(sid); if (idx < 0) return -1; int slot = impl_->inbound[idx].slot; auto state = impl_->ctx_->session_state(slot); if (state != antheos::SessionState::Active) { impl_->remove_inbound(idx); return -1; } auto slot_sid = impl_->ctx_->session_sid(slot); if (slot_sid.empty() || sid != slot_sid) { impl_->remove_inbound(idx); return -1; } auto frame = impl_->ctx_->session_call_blob(slot, {}, body, blob, blob_len); if (!frame) return -1; return impl_->bus_send(*frame) ? 0 : -1; } int Server::notify(std::string_view sid, std::string_view event) { if (!impl_->ctx_) return -1; std::lock_guard lock(impl_->session_mutex); int idx = impl_->find_inbound(sid); if (idx < 0) return -1; int slot = impl_->inbound[idx].slot; auto state = impl_->ctx_->session_state(slot); if (state != antheos::SessionState::Active) { impl_->remove_inbound(idx); return -1; } auto frame = impl_->ctx_->session_notify(slot, {}, event); if (!frame) return -1; return impl_->bus_send(*frame) ? 0 : -1; } int Server::broadcast(std::string_view event) { if (!impl_->ctx_ || !impl_->bus) return -1; std::lock_guard lock(impl_->session_mutex); int slot = impl_->ctx_->session_open(); if (slot < 0) { std::fprintf(stderr, "%s broadcast: no free session slot\n", impl_->log_prefix.c_str()); return -1; } auto frame = impl_->ctx_->session_notify(slot, {}, event); bool ok = false; if (frame) ok = impl_->bus_send(*frame); auto close_frame = impl_->ctx_->session_close(slot); if (close_frame) impl_->bus_send(*close_frame); return ok ? 0 : -1; } int Server::finish(std::string_view sid) { if (!impl_->ctx_) return -1; std::lock_guard lock(impl_->session_mutex); int idx = impl_->find_inbound(sid); if (idx < 0) return -1; int slot = impl_->inbound[idx].slot; auto state = impl_->ctx_->session_state(slot); if (state != antheos::SessionState::Active) { impl_->remove_inbound(idx); impl_->remove_known_sid(sid); return -1; } impl_->remove_inbound(idx); impl_->remove_known_sid(sid); auto frame = impl_->ctx_->session_close(slot); if (!frame) return -1; return impl_->bus_send(*frame) ? 0 : -1; } int Server::offer_additional(std::string_view description) { if (!impl_->ctx_ || !impl_->bus || description.empty()) return -1; if (impl_->additional_offers.size() >= static_cast(MAX_OFFERS)) return -1; impl_->additional_offers.emplace_back(description); return 0; } const char* Server::bid() const { if (impl_->self_bid.empty()) return nullptr; return impl_->self_bid.c_str(); } std::string_view Server::last_trace_id() const { return impl_->last_trace_id_; } std::pair Server::last_blob_tail() const { if (!impl_->ctx_) return {nullptr, 0}; /* Inline dispatch: tail is valid in ctx_ during the callback. * Threaded dispatch: tail is in the thread-local worker copy. */ if (g_is_worker_thread) return {g_worker_tail.data(), g_worker_tail.size()}; return impl_->ctx_->last_tail(); } void Server::on_accept(AcceptFn fn) { impl_->accept_fn = std::move(fn); } void Server::on_verify(VerifyFn fn) { impl_->verify_fn = std::move(fn); } void Server::on_auth(AuthFn fn) { impl_->auth_fn = std::move(fn); } int Server::require_auth(const char* trusted_keys_path) { if (!trusted_keys_path) return -1; try { impl_->trusted_keys = crypto::load_trusted_keys(trusted_keys_path); impl_->auth_required = true; return 0; } catch (...) { return -1; } } void Server::add_route(std::string_view target_bid, std::string_view path, uint32_t my_index) { Route r; r.path = std::string(path); r.my_index = my_index; impl_->relay_routes[std::string(target_bid)] = std::move(r); } int Server::send_verify(std::string_view bid) { if (!impl_->ctx_ || !impl_->bus || bid.empty()) return -1; auto frame = impl_->ctx_->verify(bid); if (!frame) return -1; if (!impl_->bus_send(*frame)) return -1; std::lock_guard lock(impl_->session_mutex); impl_->pending_verify.emplace_back(bid); return 0; } const PeerIdentity* Server::peer_identity(std::string_view bid) const { std::lock_guard lock(impl_->session_mutex); auto it = impl_->peer_identities.find(std::string(bid)); if (it == impl_->peer_identities.end()) return nullptr; return &it->second; } std::string_view Server::peer_bid_for_session(std::string_view sid) const { std::lock_guard lock(impl_->session_mutex); auto it = impl_->sid_owner.find(std::string(sid)); if (it == impl_->sid_owner.end()) return {}; return it->second; } bool Server::is_authenticated(std::string_view bid) const { if (!impl_->auth_required) return true; // no auth = all peers trusted std::lock_guard lock(impl_->session_mutex); return impl_->authenticated_peers.count(std::string(bid)) > 0; } /* ═══════════════════════════════════════════════════════════════════════ * Client * ═══════════════════════════════════════════════════════════════════════ */ struct Client::Impl { Client* owner = nullptr; Bus* bus = nullptr; std::unique_ptr ctx_; std::string oid; std::string did; std::string iid; std::string bus_name; std::string service_name; std::string log_prefix{"[AntClient]"}; std::string peer_bid; std::string trace_id_; /* Session */ int session_slot = -1; std::string session_sid; uint32_t next_mid = 1; SessionState session_state = SessionState::Disconnected; bool ownership_sent = false; /* Response wait */ std::mutex resp_mutex; std::condition_variable resp_cv; bool resp_ready = false; char resp_buf[RESP_BUF]; size_t resp_len = 0; std::string active_sid; /* Recovery */ std::string locate_bid; std::atomic ping_ack{false}; bool resume_ack = false; /* Establish/Conflict */ bool bid_established = false; bool conflict_received = false; /* Callbacks */ Client::NotifyFn notify_fn; Client::ShutdownFn shutdown_fn; bool required = false; /* Auth (Z-verb, Level 2) */ std::string auth_private_key; /* hex */ std::string auth_key_id; /* hex */ /* Relay route for multi-hop auth */ Route relay_route; bool has_relay_route = false; /* ── Helpers ───────────────────────────────────────────────────── */ bool bus_send(const antheos::Frame& frame) { if (!bus || frame.empty()) return false; return bus->write(frame.data(), frame.size()) > 0; } bool pump_bus() { uint8_t buf[READ_BUF]; ssize_t n = bus->read(buf, sizeof(buf)); if (n > 0) { ctx_->feed(buf, static_cast(n)); return true; } return false; } using clock = std::chrono::steady_clock; using time_point = clock::time_point; static time_point deadline_from_ms(int timeout_ms) { return clock::now() + std::chrono::milliseconds(timeout_ms); } static bool deadline_reached(const time_point& deadline) { return clock::now() >= deadline; } /* ── Session management ────────────────────────────────────────── */ bool open_session() { if (session_slot >= 0) return true; session_slot = ctx_->session_open(); if (session_slot < 0) { std::fprintf(stderr, "%s session_open failed\n", log_prefix.c_str()); return false; } session_sid = std::string(ctx_->session_sid(session_slot)); next_mid = 1; session_state = SessionState::Connected; ownership_sent = false; std::fprintf(stderr, "%s session opened (SID: %s)\n", log_prefix.c_str(), session_sid.c_str()); return true; } void close_session() { if (session_slot >= 0 && ctx_) { std::string_view tbid = peer_bid.empty() ? std::string_view{} : std::string_view{peer_bid}; auto frame = ctx_->session_close(session_slot, tbid); if (frame && bus) bus_send(*frame); } session_slot = -1; session_sid.clear(); next_mid = 1; } void suspend_session() { session_state = SessionState::Suspended; std::fprintf(stderr, "%s session SUSPENDED\n", log_prefix.c_str()); } bool resume_session(int timeout_ms) { if (session_sid.empty() || !ctx_ || !bus) return false; auto frame = antheos::session::locate(session_sid); if (!frame) return false; locate_bid.clear(); bus_send(*frame); int locate_timeout = std::max(timeout_ms / 2, 500); auto dl = deadline_from_ms(locate_timeout); while (!deadline_reached(dl)) { pump_bus(); if (!locate_bid.empty()) break; usleep(SLEEP_US); } if (locate_bid.empty()) { std::fprintf(stderr, "%s Locate failed (no response)\n", log_prefix.c_str()); return false; } auto resume_frame = antheos::session::resume(session_sid); if (!resume_frame) return false; resume_ack = false; bus_send(*resume_frame); int resume_timeout = std::max(timeout_ms / 2, 500); dl = deadline_from_ms(resume_timeout); while (!deadline_reached(dl)) { pump_bus(); if (resume_ack) break; usleep(SLEEP_US); } if (!resume_ack) { std::fprintf(stderr, "%s Resume failed\n", log_prefix.c_str()); return false; } peer_bid = locate_bid; session_state = SessionState::Connected; std::fprintf(stderr, "%s session resumed (BID: %s)\n", log_prefix.c_str(), peer_bid.c_str()); return true; } bool recover_connection(int timeout_ms); void do_reconnect(); /* ── Antheos callback handlers ─────────────────────────────────── */ void handle_offer(std::string_view bid, std::string_view description) { if (!description.empty() && !service_name.empty() && service_name == description) { peer_bid = std::string(bid); } } void handle_message(std::string_view verb, std::string_view sid, std::string_view, uint32_t mid, std::string_view body) { /* K-verb response (even MID) */ if (!verb.empty() && verb[0] == 'K' && mid % 2 == 0 && !body.empty() && !sid.empty()) { std::lock_guard lock(resp_mutex); if (!active_sid.empty() && active_sid != sid) return; size_t blen = body.size(); if (blen >= RESP_BUF) blen = RESP_BUF - 1; std::memcpy(resp_buf, body.data(), blen); resp_buf[blen] = '\0'; resp_len = blen; resp_ready = true; resp_cv.notify_one(); return; } /* N-verb notification */ if (!verb.empty() && verb[0] == 'N' && !body.empty()) { if (notify_fn) { auto sp = body.find(' '); if (sp != std::string_view::npos) { notify_fn(body.substr(0, sp), body.substr(sp + 1)); } else { notify_fn(body, ""); } } else { std::fprintf(stderr, "%s N-verb dropped (no handler): %.*s\n", log_prefix.c_str(), static_cast(std::min(body.size(), size_t{60})), body.data()); } } } void handle_event(std::string_view verb, std::string_view id, std::string_view id2, std::string_view detail) { if (verb.empty()) return; switch (verb[0]) { case 'W': ping_ack.store(true, std::memory_order_relaxed); break; case 'L': if (!id.empty() && !id2.empty() && !session_sid.empty() && session_sid == id) locate_bid = std::string(id2); break; case 'T': if (!id.empty() && !session_sid.empty() && session_sid == id) { std::string_view state_str = !detail.empty() ? detail : (!id2.empty() ? id2 : std::string_view{}); if (state_str.find("ACTIVE") != std::string_view::npos) resume_ack = true; } break; case 'X': std::fprintf(stderr, "[AntClient] Exception: %.*s %.*s\n", static_cast(id.size()), id.data(), static_cast(detail.size()), detail.data()); break; case 'E': if (!id.empty() && ctx_ && bid_established) { auto our_bid = ctx_->bid(); if (!our_bid.empty() && our_bid == id) { auto frame = antheos::bus::conflict(our_bid); if (frame) bus_send(*frame); } } break; case 'V': /* Verify request: respond with our identity */ if (!id.empty() && id2.empty() && ctx_) { auto our_bid = ctx_->bid(); if (!our_bid.empty() && our_bid == id) { auto vr = ctx_->verify_response(); if (vr) bus_send(*vr); } } break; case 'Z': /* Auth challenge: id=our BID, id2=empty, detail=nonce */ if (!id.empty() && id2.empty() && !detail.empty() && ctx_) { auto our_bid = ctx_->bid(); if (!our_bid.empty() && our_bid == id && !auth_private_key.empty()) { std::string sig = crypto::sign(detail, auth_private_key); auto zr = ctx_->auth_response( peer_bid, auth_key_id, sig); if (zr) bus_send(*zr); } } break; case 'C': if (!id.empty() && ctx_) { auto our_bid = ctx_->bid(); if (!our_bid.empty() && our_bid == id) conflict_received = true; } break; default: break; } } void setup_callbacks() { ctx_->on_offer([this](std::string_view bid, std::string_view desc) { handle_offer(bid, desc); }); ctx_->on_message( [this](std::string_view verb, std::string_view sid, std::string_view target, uint32_t mid, std::string_view body) { handle_message(verb, sid, target, mid, body); }); ctx_->on_event( [this](std::string_view verb, std::string_view id, std::string_view id2, std::string_view detail) { handle_event(verb, id, id2, detail); }); ctx_->on_relay( [this](char ref, std::string_view id, std::string_view id2, std::string_view body, uint32_t /*index*/, std::string_view /*path*/) { handle_relay(ref, id, id2, body); }); } void handle_relay(char ref, std::string_view id, std::string_view id2, std::string_view body) { if (ref != 'Z') return; /* Relayed Z-challenge: id=our BID, id2=empty, body=nonce */ if (!id.empty() && id2.empty() && !body.empty() && ctx_) { auto our_bid = ctx_->bid(); if (!our_bid.empty() && our_bid == id && !auth_private_key.empty() && has_relay_route) { std::string sig = crypto::sign(body, auth_private_key); auto zr = ctx_->relay_auth_response( peer_bid, auth_key_id, sig, relay_route.my_index, relay_route.path); if (zr) bus_send(*zr); } } } bool establish_with_conflict() { size_t bid_len = antheos::id::BID_MIN_LEN; bid_established = false; while (bid_len <= antheos::id::BID_MAX_LEN) { try { ctx_ = std::make_unique( oid, did, iid, generate_bid(bid_len)); } catch (const std::exception& e) { std::fprintf(stderr, "%s antheos::Context init failed: %s\n", log_prefix.c_str(), e.what()); return false; } setup_callbacks(); conflict_received = false; auto frame = ctx_->establish(); if (!frame || !bus_send(*frame)) return false; auto deadline = std::chrono::steady_clock::now() + std::chrono::milliseconds(ESTABLISH_TIMEOUT_MS); while (std::chrono::steady_clock::now() < deadline) { if (!pump_bus()) usleep(SLEEP_US); if (conflict_received) break; } if (!conflict_received) { bid_established = true; return true; } std::fprintf(stderr, "%s BID conflict at length %zu, retrying\n", log_prefix.c_str(), bid_len); bid_len++; } std::fprintf(stderr, "%s BID_OVERFLOW — all lengths exhausted\n", log_prefix.c_str()); auto frame = antheos::bus::exception("BID_OVERFLOW"); if (frame) bus_send(*frame); ctx_.reset(); return false; } }; /* ── Client::Impl deferred methods ─────────────────────────────────── */ void Client::Impl::do_reconnect() { if (bus_name.empty() || !bus) { close_session(); peer_bid.clear(); session_state = SessionState::Disconnected; return; } close_session(); ctx_.reset(); bus->reopen(); if (!bus->is_open()) { std::fprintf(stderr, "%s reconnect: bus reopen failed\n", log_prefix.c_str()); bus = nullptr; peer_bid.clear(); session_state = SessionState::Disconnected; return; } if (!establish_with_conflict()) { bus = nullptr; peer_bid.clear(); session_state = SessionState::Disconnected; return; } peer_bid.clear(); session_state = SessionState::Disconnected; std::fprintf(stderr, "%s reconnect — bus reopened\n", log_prefix.c_str()); } bool Client::Impl::recover_connection(int timeout_ms) { std::fprintf(stderr, "%s recovering connection...\n", log_prefix.c_str()); int resume_budget = std::max(timeout_ms / 3, 1000); if (resume_session(resume_budget)) return true; close_session(); peer_bid.clear(); if (!bus || !ctx_) { if (bus_name.empty()) { session_state = SessionState::Disconnected; return false; } do_reconnect(); if (!bus || !ctx_) { session_state = SessionState::Disconnected; return false; } } if (owner->discover(service_name, timeout_ms) != 0) { session_state = SessionState::Disconnected; return false; } std::fprintf(stderr, "%s recovered via rediscovery\n", log_prefix.c_str()); return true; } /* ── Client public API ─────────────────────────────────────────────── */ Client::Client() : impl_(std::make_unique()) { impl_->owner = this; } Client::~Client() { if (impl_) { close(); } } int Client::init(Bus& bus, std::string_view oid, std::string_view did, std::string_view iid) { if (oid.empty() || did.empty() || iid.empty() || !bus.is_open()) return -1; impl_->bus = &bus; impl_->bus_name = std::string(bus.name()); impl_->oid = oid; impl_->did = did; impl_->iid = std::string(iid); impl_->log_prefix = "[AntClient]"; if (!impl_->establish_with_conflict()) { impl_->bus = nullptr; return -1; } std::fprintf(stderr, "%s initialized on bus '%s'\n", impl_->log_prefix.c_str(), impl_->bus_name.c_str()); return 0; } int Client::discover(std::string_view service, int timeout_ms) { if (!impl_->ctx_ || !impl_->bus || service.empty()) return -1; impl_->service_name = service; impl_->log_prefix = "[AntClient:" + std::string(service) + "]"; impl_->peer_bid.clear(); impl_->close_session(); auto frame = impl_->ctx_->query(service); if (!frame) { std::fprintf(stderr, "%s query failed\n", impl_->log_prefix.c_str()); return -1; } impl_->bus_send(*frame); auto deadline = Impl::deadline_from_ms(timeout_ms); while (impl_->peer_bid.empty()) { if (Impl::deadline_reached(deadline)) { std::fprintf(stderr, "%s discovery timeout\n", impl_->log_prefix.c_str()); auto xframe = antheos::bus::exception("SERVICE_UNKNOWN"); if (xframe) impl_->bus_send(*xframe); return -1; } impl_->pump_bus(); if (impl_->peer_bid.empty()) usleep(SLEEP_US); } auto our_bid = impl_->ctx_->bid(); auto accept_frame = impl_->ctx_->accept(impl_->peer_bid, our_bid); if (accept_frame) impl_->bus_send(*accept_frame); impl_->session_state = SessionState::Connected; std::fprintf(stderr, "%s discovered (BID: %s)\n", impl_->log_prefix.c_str(), impl_->peer_bid.c_str()); return 0; } int Client::call(std::string_view command, int timeout_ms, char* reply, size_t reply_max) { if (command.data() == nullptr) return -1; /* Build wire headers: T: (trace), O: (ownership) — matches server extraction order */ std::string header_cmd; std::string_view wire_cmd = command; bool need_trace = !impl_->trace_id_.empty(); bool need_ownership = !impl_->ownership_sent && impl_->ctx_; if (need_trace || need_ownership) { header_cmd.reserve( (need_trace ? 2 + impl_->trace_id_.size() + 1 : 0) + (need_ownership ? 2 + impl_->ctx_->bid().size() + 1 : 0) + command.size()); if (need_trace) { header_cmd += "T:"; header_cmd += impl_->trace_id_; header_cmd += '\n'; } if (need_ownership) { std::string_view bid = impl_->ctx_->bid(); header_cmd += "O:"; header_cmd.append(bid.data(), bid.size()); header_cmd += '\n'; impl_->ownership_sent = true; } header_cmd.append(command.data(), command.size()); wire_cmd = header_cmd; } for (int attempt = 0; attempt < 2; ++attempt) { if (impl_->session_state == SessionState::Suspended) { if (!impl_->recover_connection(timeout_ms)) break; } if (!impl_->ctx_ || !impl_->bus || impl_->peer_bid.empty()) break; if (impl_->session_slot < 0) { if (!impl_->open_session()) break; } auto frame = antheos::session::call( impl_->session_sid, impl_->peer_bid, impl_->next_mid, wire_cmd); if (!frame) { std::fprintf(stderr, "%s session_call encode failed\n", impl_->log_prefix.c_str()); break; } impl_->bus_send(*frame); { std::lock_guard lock(impl_->resp_mutex); impl_->active_sid = impl_->session_sid; impl_->resp_ready = false; impl_->resp_len = 0; } auto deadline = Impl::deadline_from_ms(timeout_ms); bool got_response = false; while (!Impl::deadline_reached(deadline)) { { std::lock_guard lock(impl_->resp_mutex); if (impl_->resp_ready) { got_response = true; break; } } impl_->pump_bus(); usleep(SLEEP_US); } if (!got_response) { std::lock_guard lock(impl_->resp_mutex); if (impl_->resp_ready) got_response = true; } if (got_response) { impl_->next_mid += 2; std::lock_guard lock(impl_->resp_mutex); impl_->active_sid.clear(); int written = 0; if (reply && reply_max > 0) { size_t copy = impl_->resp_len; if (copy >= reply_max) copy = reply_max - 1; std::memcpy(reply, impl_->resp_buf, copy); reply[copy] = '\0'; written = static_cast(copy); } return written; } /* Timeout — peer unreachable */ std::fprintf(stderr, "%s call timeout (%dms), SUSPENDED: %.*s\n", impl_->log_prefix.c_str(), timeout_ms, static_cast(std::min(command.size(), size_t{40})), command.data()); { std::lock_guard lock(impl_->resp_mutex); impl_->active_sid.clear(); } impl_->suspend_session(); } if (impl_->required) { std::fprintf(stderr, "%s SERVICE_UNKNOWN: required peer '%s' on bus '%s' " "not found after recovery — shutting down\n", impl_->log_prefix.c_str(), impl_->service_name.c_str(), impl_->bus_name.c_str()); auto xframe = antheos::bus::exception("SERVICE_UNKNOWN"); if (xframe) impl_->bus_send(*xframe); if (impl_->shutdown_fn) impl_->shutdown_fn(); return -1; } std::fprintf(stderr, "%s recovery failed after retry (non-required peer)\n", impl_->log_prefix.c_str()); return 0; } std::string Client::call(std::string_view command, int timeout_ms) { char buf[READ_BUF]; int n = call(command, timeout_ms, buf, sizeof(buf)); if (n <= 0) return {}; return std::string(buf, static_cast(n)); } int Client::fire(std::string_view command) { if (command.data() == nullptr || !impl_->ctx_ || !impl_->bus || impl_->peer_bid.empty()) return -1; if (impl_->session_slot < 0) { if (!impl_->open_session()) return -1; } std::string traced_cmd; std::string_view wire_cmd = command; if (!impl_->trace_id_.empty()) { traced_cmd.reserve(2 + impl_->trace_id_.size() + 1 + command.size()); traced_cmd += "T:"; traced_cmd += impl_->trace_id_; traced_cmd += '\n'; traced_cmd.append(command.data(), command.size()); wire_cmd = traced_cmd; } auto frame = antheos::session::call( impl_->session_sid, impl_->peer_bid, impl_->next_mid, wire_cmd); if (!frame) return -1; bool ok = impl_->bus_send(*frame); impl_->next_mid += 2; return ok ? 0 : -1; } bool Client::connected() const { return impl_->session_state == SessionState::Connected && !impl_->peer_bid.empty(); } SessionState Client::state() const { return impl_->session_state; } int Client::ensure_connected(int timeout_ms) { if (impl_->session_state == SessionState::Connected && !impl_->peer_bid.empty()) return 0; if (impl_->session_state == SessionState::Suspended) { if (impl_->recover_connection(timeout_ms)) return 0; if (impl_->required) { std::fprintf(stderr, "%s SERVICE_UNKNOWN: required peer '%s' on bus '%s' " "not found after recovery — shutting down\n", impl_->log_prefix.c_str(), impl_->service_name.c_str(), impl_->bus_name.c_str()); auto xframe = antheos::bus::exception("SERVICE_UNKNOWN"); if (xframe) impl_->bus_send(*xframe); if (impl_->shutdown_fn) impl_->shutdown_fn(); } return -1; } if (impl_->service_name.empty()) { std::fprintf(stderr, "%s ensureConnected: no service name set\n", impl_->log_prefix.c_str()); return -1; } if (!impl_->bus) { if (impl_->bus_name.empty()) { std::fprintf(stderr, "%s ensureConnected: no bus name set\n", impl_->log_prefix.c_str()); return -1; } impl_->do_reconnect(); if (!impl_->bus) return -1; } std::fprintf(stderr, "%s ensureConnected: rediscovering '%s'...\n", impl_->log_prefix.c_str(), impl_->service_name.c_str()); if (discover(impl_->service_name, timeout_ms) == 0) return 0; if (impl_->required) { std::fprintf(stderr, "%s SERVICE_UNKNOWN: required peer '%s' on bus '%s' " "not found after discovery — shutting down\n", impl_->log_prefix.c_str(), impl_->service_name.c_str(), impl_->bus_name.c_str()); auto xframe = antheos::bus::exception("SERVICE_UNKNOWN"); if (xframe) impl_->bus_send(*xframe); if (impl_->shutdown_fn) impl_->shutdown_fn(); } return -1; } void Client::reconnect() { impl_->do_reconnect(); } void Client::close() { impl_->close_session(); impl_->ctx_.reset(); impl_->bus = nullptr; impl_->peer_bid.clear(); impl_->bus_name.clear(); impl_->session_state = SessionState::Disconnected; } void Client::on_notify(NotifyFn fn) { impl_->notify_fn = std::move(fn); } void Client::on_shutdown(ShutdownFn fn) { impl_->shutdown_fn = std::move(fn); } void Client::set_required(bool required) { impl_->required = required; } int Client::set_auth_key(const char* private_key_path) { if (!private_key_path) return -1; try { auto kp = crypto::load_key(private_key_path); impl_->auth_private_key = std::move(kp.private_key); impl_->auth_key_id = std::move(kp.key_id); return 0; } catch (...) { return -1; } } void Client::set_route(std::string_view path, uint32_t my_index) { impl_->relay_route.path = std::string(path); impl_->relay_route.my_index = my_index; impl_->has_relay_route = true; } const char* Client::peer_bid() const { if (impl_->peer_bid.empty()) return nullptr; return impl_->peer_bid.c_str(); } void Client::set_trace_id(std::string_view id) { impl_->trace_id_ = std::string(id); } void Client::clear_trace_id() { impl_->trace_id_.clear(); } std::string_view Client::trace_id() const { return impl_->trace_id_; } } // namespace antpeer