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:
a direct in-process callback — the
subscribe(src, callback)sugar;a spec-faithful target vertex —
subscribe(src, target)re-dispatches to a handler-backed sink vertex;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_vertexreturns avertex_handle_treused for every op.awaitis join-safe — the publisherjoin()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.