In-process pub/sub (L4 graph)

The smallest complete libtracer node: a single-process graph with no transport and no wire bytes. A publisher writes /sensor/temp; three subscribers each receive the value a different way:

  1. a direct in-process callback — the subscribe(src, callback) sugar;

  2. a spec-faithful target vertex — subscribe(src, target) re-dispatches to a handler-backed sink vertex;

  3. a thread blocking in await() — the single-shot readiness primitive.

Delivery to (1) and (2) is a refcount-bump clone of the same rope value — no byte copy. The example finishes by reading the last-known-value back and declaring the STREAM ring depth owner-side (vertex_policy_t::retention, which has no wire surface), then discovering the vertex shape via the :schema control read.

What to notice

  • Handles, not strings, on the hot path — the path is encoded once via the parse-once path_t("…") constructor; register_vertex returns a vertex_handle_t reused for every op.

  • await is join-safe — the publisher join()s the waiter after the write, so every delivery is complete and visible when the checks run.

  • It self-checks — each delivery path is asserted, so the ctest smoke test guards behavior, not just a clean exit.

Source

  1/*
  2 * SPDX-License-Identifier: Apache-2.0
  3 * SPDX-FileCopyrightText: Copyright 2026 avatarsd LLC
  4 */
  5
  6/**
  7 * @file
  8 * @brief In-process publish/subscribe over the L4 graph — the M3 P0 node, end to end.
  9 *
 10 * No transport, no wire bytes leave the process. A publisher writes
 11 * `/sensor/temp`; three subscribers receive the value three different ways:
 12 *   1. a direct in-process callback  — `subscribe(src, callback)` sugar;
 13 *   2. a spec-faithful target vertex — `subscribe(src, target)` → a handler sink;
 14 *   3. a thread blocking in `await()` — the single-shot primitive.
 15 * Delivery to (1) and (2) is a refcount-bump clone of the same rope value — no
 16 * byte copy.
 17 *
 18 * Runs under ctest as `example_in_process_pubsub`: it @ref check "checks" every
 19 * delivered value and returns non-zero on any mismatch, so the smoke test guards
 20 * behavior (each subscriber sees the write) rather than merely exiting cleanly.
 21 */
 22
 23#include <chrono>
 24#include <cstddef>
 25#include <cstdint>
 26#include <cstdio>
 27#include <iterator>
 28#include <thread>
 29
 30#include "libtracer/tracer.hpp"
 31
 32namespace {
 33
 34using namespace std::chrono_literals;
 35using tr::graph::path_t;
 36using tr::graph::role_t;
 37
 38/** @brief A 4-byte little-endian VALUE view over @p v (one heap segment). */
 39tr::view::view_t value_u32(std::uint32_t v) {
 40    tr::view::segment_ptr_t seg = tr::view::heap_alloc(4);
 41    for (int i = 0; i < 4; ++i)
 42        seg->bytes[static_cast<std::size_t>(i)] = static_cast<std::byte>((v >> (8 * i)) & 0xFF);
 43    return tr::view::view_t::over(std::move(seg));
 44}
 45
 46/** @brief Decode a little-endian u32 from @p view's first (≤4) bytes. */
 47std::uint32_t as_u32(const tr::view::view_t& view) {
 48    const auto b = view.bytes();
 49    std::uint32_t v = 0;
 50    for (std::size_t i = 0; i < b.size() && i < 4; ++i)
 51        v |= static_cast<std::uint32_t>(std::to_integer<std::uint8_t>(b[i])) << (8 * i);
 52    return v;
 53}
 54
 55/** @brief Record a failed expectation on @p ok and report it; leaves @p ok false on failure. */
 56void check(bool& ok, bool cond, const char* what) {
 57    if (!cond) {
 58        std::printf("  [FAIL] %s\n", what);
 59        ok = false;
 60    }
 61}
 62
 63}  // namespace
 64
 65int main() {
 66    constexpr std::uint32_t kSent = 23;
 67    tr::graph::graph_t g;
 68
 69    const tr::graph::vertex_handle_t temp =
 70        g.register_vertex(path_t("/sensor/temp"), role_t::STORED_VALUE);
 71
 72    // Values captured from each delivery path, verified at the end.
 73    std::uint32_t cb_got = 0;     // subscriber 1 (callback)
 74    std::uint32_t sink_got = 0;   // subscriber 2 (target vertex handler)
 75    std::uint32_t await_got = 0;  // subscriber 3 (await)
 76
 77    // A "sink" vertex backed by a callback handler — the target of a spec-faithful
 78    // SUBSCRIBER (subscriber 2 below re-dispatches to it).
 79    // Its `on_write` is a `{fn, ctx}` hook: `ctx` is `&sink_got`, which outlives the graph's
 80    // use of it, and the value arrives by reference — the very block /sensor/temp published.
 81    tr::graph::handlers_t sink;
 82    sink.on_write = {[](void* ctx, const tr::graph::value_t& in,
 83                        const tr::graph::write_ctx_t&) -> tr::graph::result_t<void> {
 84                         auto& got = *static_cast<std::uint32_t*>(ctx);
 85                         got = as_u32(in.only());
 86                         std::printf("  [sink vertex /log/temp] received %u\n", got);
 87                         return {};
 88                     },
 89                     &sink_got};
 90    (void)g.register_vertex(path_t("/log/temp"), role_t::HANDLER, sink);
 91
 92    // subscriber_t 1 — direct in-process callback.
 93    auto on_temp = [&cb_got](const tr::graph::value_t& v) {
 94        cb_got = as_u32(v.only());
 95        std::printf("  [callback sub] received %u\n", cb_got);
 96    };
 97    (void)g.subscribe(path_t("/sensor/temp"), on_temp);
 98    // subscriber_t 2 — spec-faithful target-path subscription -> /log/temp.
 99    (void)g.subscribe(path_t("/sensor/temp"), path_t("/log/temp"));
100
101    // subscriber_t 3 — a thread blocking in await().
102    std::thread waiter([&] {
103        auto r = g.await(temp, 2s);
104        if (r) {
105            await_got = as_u32((*r)->only());
106            std::printf("  [await sub] received %u\n", await_got);
107        }
108    });
109    std::this_thread::sleep_for(50ms);  // let the waiter park in await()
110
111    std::printf("publisher: write /sensor/temp = %u\n", kSent);
112    (void)g.write(temp, value_u32(kSent));
113    waiter.join();  // joins-after-write: every delivery is complete + visible here
114
115    // Read back the last-known-value (a clone — keeps the segment alive for us).
116    auto rb = g.read(temp);
117    const std::uint32_t rb_got = rb ? as_u32((*rb)->only()) : 0u;
118    std::printf("read-back /sensor/temp = %u\n", rb_got);
119
120    // Declare the STREAM ring depth OWNER-SIDE (RFC-0022 §3.C — it has no wire surface),
121    // then discover the vertex shape via :schema.
122    (void)g.set_policy(temp, {.retention = tr::graph::retention_t::N, .depth = 8});
123    auto schema = g.read(path_t("/sensor/temp:schema"));
124    std::size_t schema_children = 0;
125    if (schema) {
126        if (auto point = tr::wire::tlv_node_t::over((*schema)->only()))
127            schema_children = static_cast<std::size_t>(std::ranges::distance(point->children()));
128    }
129    std::printf(":schema is a POINT with %zu children\n", schema_children);
130
131    // Guard the semantics, not just a clean exit: every path must see the write.
132    bool ok = true;
133    check(ok, cb_got == kSent, "callback subscriber received the written value");
134    check(ok, sink_got == kSent, "target-vertex subscriber received the written value");
135    check(ok, await_got == kSent, "await subscriber received the written value");
136    check(ok, rb_got == kSent, "read-back returns the last-known-value");
137    check(ok, schema_children == 2, ":schema resolves to a 2-child POINT");
138    return ok ? 0 : 1;
139}

See also: graph module · views · path.