Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion MODULE.bazel
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
module(
name = "subspace",
version = "3.0.3",
version = "3.1.0",
)

bazel_dep(name = "bazel_skylib", version = "1.9.0")
Expand Down
2 changes: 1 addition & 1 deletion MODULE.bazel.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

28 changes: 28 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ It has the following features:
1. Optional split payload buffers for external allocators and memory pools.
1. Explicit multi-slot publisher buffer leases with exact-slot reclamation.
1. Server-enforced publisher and subscriber limits with lease-aware channel capacity.
1. Server-generated channel telemetry for participant, drop, and resize changes.
1. Automatic UDP discovery and TCP bridging of channels between servers, plus optional TCP unicast discovery for bridging across NAT/VMs/emulators.
1. Shadow process for crash recovery -- the server can restart and resume without losing shared memory state.
1. Shared and weak pointers for message references.
Expand All @@ -41,6 +42,7 @@ It has the following features:

See the file docs/subspace.pdf for full documentation. Additional documentation:
- [Checksums and User Metadata](docs/checksums-and-metadata.md)
- [Channel Telemetry](docs/channel-telemetry.md)
- [Split Buffers](docs/split-buffers.md)
- [Publisher Buffer Leases](docs/publisher-buffer-leases.md)
- [C Client API](docs/c-client.md)
Expand Down Expand Up @@ -733,6 +735,28 @@ auto msg_ptr = msg_ptr_or.value();
// Message is automatically released when msg_ptr goes out of scope
```

**Method 4: Read server-generated channel telemetry**
```cpp
auto telemetry_sub = client->CreateSubscriber(
"my_channel",
subspace::SubscriberOptions().SetTelemetry(true)).value();

auto telemetry_or = telemetry_sub.ReadTelemetryMessage();
if (!telemetry_or.ok()) {
// Handle a read or protobuf decoding error
return;
}
std::shared_ptr<subspace::Telemetry> telemetry = *telemetry_or;
if (telemetry == nullptr) {
// No telemetry message is currently available
return;
}
```

The monitored channel must already exist. The server sends an initial
participant snapshot and then batches participant, drop, and resize changes at
one-second intervals. See [Channel Telemetry](docs/channel-telemetry.md).

### Waiting for Messages

```cpp
Expand Down Expand Up @@ -771,6 +795,8 @@ public:
absl::StatusOr<shared_ptr<T>> ReadMessage(
ReadMode mode = ReadMode::kReadNext,
ClearTrigger clear_trigger = ClearTrigger::kClearTrigger);
absl::StatusOr<std::shared_ptr<Telemetry>>
ReadTelemetryMessage(ReadMode mode = ReadMode::kReadNext);

// Find message by timestamp
absl::StatusOr<Message> FindMessage(uint64_t timestamp);
Expand Down Expand Up @@ -1153,6 +1179,7 @@ auto sub = client->CreateSubscriber("channel",
| Field/Method | Type | Default | Description |
|--------------|------|---------|-------------|
| `reliable` / `SetReliable()` | `bool` | `false` | If true, reliable delivery (see Reliable Channels section). |
| `telemetry` / `SetTelemetry()` | `bool` | `false` | Subscribe to server-generated telemetry for the named existing channel instead of its payloads. |
| `type` / `SetType()` | `std::string` | `""` | User-defined message type identifier. Must match publisher type. |
| `max_active_messages` / `SetMaxActiveMessages()` | `int` | `1` | Maximum number of active messages (shared_ptrs) that can be held simultaneously. |
| `max_active_messages` / `SetMaxSharedPtrs()` | `int` | `0` | Alias: sets max_active_messages to n+1. |
Expand All @@ -1168,6 +1195,7 @@ auto sub = client->CreateSubscriber("channel",

**Getter Methods:**
- `bool IsReliable() const`
- `bool Telemetry() const`
- `const std::string& Type() const`
- `int MaxActiveMessages() const`
- `int MaxSharedPtrs() const`
Expand Down
112 changes: 112 additions & 0 deletions c_client/client_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
#include "toolbelt/hexdump.h"
#include "toolbelt/pipe.h"
#include <algorithm>
#include <chrono>
#include <gtest/gtest.h>
#include <inttypes.h>
#include <memory>
Expand Down Expand Up @@ -372,6 +373,60 @@ TEST_F(ClientTest, PublishSingleMessageAndRead) {
ASSERT_TRUE(subspace_remove_client(&sub_client));
}

TEST_F(ClientTest, ReadTelemetryMessage) {
SubspaceClient publisher_client =
subspace_create_client_with_socket_and_name(Socket().c_str(),
"c-telemetry-publisher");
SubspaceClient watcher_client =
subspace_create_client_with_socket_and_name(Socket().c_str(),
"c-telemetry-watcher");
ASSERT_NE(nullptr, publisher_client.client);
ASSERT_NE(nullptr, watcher_client.client);

SubspacePublisher publisher = subspace_create_publisher(
publisher_client, "c_telemetry",
CPublisherOptionsDefault(/*slot_size=*/128, /*num_slots=*/4));
ASSERT_NE(nullptr, publisher.publisher) << subspace_get_last_error();
SubspaceSubscriberOptions options = CSubscriberOptionsDefault();
options.telemetry = true;
SubspaceSubscriber subscriber =
subspace_create_subscriber(watcher_client, "c_telemetry", options);
ASSERT_NE(nullptr, subscriber.subscriber) << subspace_get_last_error();

SubspaceTelemetry telemetry = {};
auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds(5);
while (telemetry.telemetry == nullptr &&
std::chrono::steady_clock::now() < deadline) {
telemetry = subspace_read_telemetry_message(subscriber);
ASSERT_FALSE(subspace_has_error()) << subspace_get_last_error();
if (telemetry.telemetry == nullptr) {
std::this_thread::sleep_for(std::chrono::milliseconds(50));
}
}
ASSERT_NE(nullptr, telemetry.telemetry);
bool found_publisher = false;
for (size_t i = 0; i < telemetry.num_publishers; ++i) {
const SubspaceTelemetryParticipant &entry = telemetry.publishers[i];
found_publisher |=
std::string(entry.name.data, entry.name.length) ==
"c-telemetry-publisher" &&
entry.change == kSubspaceTelemetryNoChange;
}
EXPECT_TRUE(found_publisher);
EXPECT_TRUE(subspace_free_telemetry(&telemetry));
EXPECT_EQ(nullptr, telemetry.telemetry);

telemetry = subspace_read_telemetry_message_with_mode(
subscriber, kSubspaceReadNewest);
EXPECT_EQ(nullptr, telemetry.telemetry);
EXPECT_FALSE(subspace_has_error());

EXPECT_TRUE(subspace_remove_subscriber(&subscriber));
EXPECT_TRUE(subspace_remove_publisher(&publisher));
EXPECT_TRUE(subspace_remove_client(&publisher_client));
EXPECT_TRUE(subspace_remove_client(&watcher_client));
}

TEST_F(ClientTest, ReadMessageNoClearTriggerLeavesEventFdReadable) {
auto pub_client = subspace_create_client_with_socket(Socket().c_str());
ASSERT_NE(nullptr, pub_client.client);
Expand Down Expand Up @@ -1810,6 +1865,63 @@ void DroppedMessageCallback(SubspaceSubscriber /*subscriber*/,
num_dropped_messages += num_dropped;
}

TEST_F(ClientTest, SubscriberOptionsTelemetry) {
SubspaceSubscriberOptions options = subspace_subscriber_options_default();
ASSERT_FALSE(options.telemetry);

options.telemetry = true;
ASSERT_TRUE(options.telemetry);
}

TEST_F(ClientTest, TelemetrySubscriberSmoke) {
auto pub_client = subspace_create_client_with_socket(Socket().c_str());
ASSERT_NE(nullptr, pub_client.client);
ASSERT_FALSE(subspace_has_error());
auto watcher_client = subspace_create_client_with_socket(Socket().c_str());
ASSERT_NE(nullptr, watcher_client.client);
ASSERT_FALSE(subspace_has_error());

SubspacePublisher pub = subspace_create_publisher(
pub_client, "c_telemetry_smoke", CPublisherOptionsDefault(128, 4));
ASSERT_NE(nullptr, pub.publisher);
ASSERT_FALSE(subspace_has_error());

SubspaceSubscriberOptions telemetry_opts = CSubscriberOptionsDefault();
telemetry_opts.telemetry = true;
SubspaceSubscriber telemetry = subspace_create_subscriber(
watcher_client, "c_telemetry_smoke", telemetry_opts);
ASSERT_NE(nullptr, telemetry.subscriber);
ASSERT_FALSE(subspace_has_error());

SubspaceTypeInfo type_info = subspace_get_subscriber_type(telemetry);
ASSERT_FALSE(subspace_has_error());
ASSERT_EQ(strlen("subspace.Telemetry"), type_info.type_length);
ASSERT_EQ(
0, memcmp(type_info.type, "subspace.Telemetry", type_info.type_length));

SubspaceMessage msg = {};
bool got_message = false;
const auto deadline =
std::chrono::steady_clock::now() + std::chrono::seconds(5);
while (std::chrono::steady_clock::now() < deadline) {
msg = subspace_read_message(telemetry);
ASSERT_FALSE(subspace_has_error());
if (msg.length > 0) {
got_message = true;
break;
}
std::this_thread::sleep_for(std::chrono::milliseconds(50));
}
ASSERT_TRUE(got_message) << "Timed out waiting for telemetry payload";
ASSERT_GT(msg.length, 0);
subspace_free_message(&msg);

ASSERT_TRUE(subspace_remove_subscriber(&telemetry));
ASSERT_TRUE(subspace_remove_publisher(&pub));
ASSERT_TRUE(subspace_remove_client(&pub_client));
ASSERT_TRUE(subspace_remove_client(&watcher_client));
}

TEST_F(ClientTest, DroppedMessage) {
auto pub_client = subspace_create_client_with_socket(Socket().c_str());
ASSERT_NE(nullptr, pub_client.client);
Expand Down
86 changes: 86 additions & 0 deletions c_client/subspace.cc
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,14 @@ struct HandleCache {
std::vector<SubspaceMessage> messages;
};

struct TelemetryStorage {
std::shared_ptr<subspace::Telemetry> message;
std::vector<SubspaceTelemetryParticipant> publishers;
std::vector<SubspaceTelemetryParticipant> subscribers;
std::vector<SubspaceTelemetryDrop> drops;
std::vector<SubspaceTelemetryResize> resizes;
};

std::unordered_map<void *, ClientCache> client_caches;
std::unordered_map<void *, HandleCache> publisher_caches;
std::unordered_map<void *, HandleCache> subscriber_caches;
Expand Down Expand Up @@ -61,6 +69,48 @@ SubspaceString ToCString(const std::string &s) {
return {.data = s.data(), .length = s.size()};
}

SubspaceTelemetry TakeCTelemetry(
std::shared_ptr<subspace::Telemetry> message) {
if (message == nullptr) {
return {};
}

auto *storage = new TelemetryStorage;
storage->message = std::move(message);
storage->publishers.reserve(storage->message->publishers_size());
for (const auto &publisher : storage->message->publishers()) {
storage->publishers.push_back(
{.name = ToCString(publisher.name()),
.change =
static_cast<SubspaceTelemetryChange>(publisher.change())});
}
storage->subscribers.reserve(storage->message->subscribers_size());
for (const auto &subscriber : storage->message->subscribers()) {
storage->subscribers.push_back(
{.name = ToCString(subscriber.name()),
.change =
static_cast<SubspaceTelemetryChange>(subscriber.change())});
}
storage->drops.reserve(storage->message->drops_size());
for (const auto &drop : storage->message->drops()) {
storage->drops.push_back({.num_drops = drop.num_drops()});
}
storage->resizes.reserve(storage->message->resizes_size());
for (const auto &resize : storage->message->resizes()) {
storage->resizes.push_back({.new_size = resize.new_size()});
}

return {.telemetry = storage,
.publishers = storage->publishers.data(),
.num_publishers = storage->publishers.size(),
.subscribers = storage->subscribers.data(),
.num_subscribers = storage->subscribers.size(),
.drops = storage->drops.data(),
.num_drops = storage->drops.size(),
.resizes = storage->resizes.data(),
.num_resizes = storage->resizes.size()};
}

SubspaceChannelCounters ToCCounters(const subspace::ChannelCounters &counters) {
return {
.num_pub_updates = counters.num_pub_updates,
Expand Down Expand Up @@ -530,6 +580,7 @@ subspace_create_subscriber(SubspaceClient client, const char *channel_name,
.SetSubscriberQueueSize(options.subscriber_queue_size)
.SetBridge(options.bridge)
.SetForTunnel(options.for_tunnel)
.SetTelemetry(options.telemetry)
.SetType(StringFromPointer(options.type.type, options.type.type_length))
.SetMaxActiveMessages(options.max_active_messages)
.SetMaxSubscribers(options.max_subscribers)
Expand Down Expand Up @@ -646,6 +697,41 @@ SubspaceMessage subspace_read_message(SubspaceSubscriber subscriber) {
return subspace_read_message_with_mode(subscriber, kSubspaceReadNext);
}

SubspaceTelemetry subspace_read_telemetry_message_with_mode(
SubspaceSubscriber subscriber, SubspaceReadMode mode) {
subspace_clear_error();
if (subscriber.subscriber == nullptr) {
subspace_set_error("Invalid subscriber");
return {};
}

auto sub_ptr = SubscriberPtr(subscriber);
absl::StatusOr<std::shared_ptr<subspace::Telemetry>> telemetry =
(*sub_ptr)->ReadTelemetryMessage(ToCppReadMode(mode));
if (!telemetry.ok()) {
subspace_set_error(telemetry.status().ToString().c_str());
return {};
}
return TakeCTelemetry(std::move(*telemetry));
}

SubspaceTelemetry
subspace_read_telemetry_message(SubspaceSubscriber subscriber) {
return subspace_read_telemetry_message_with_mode(subscriber,
kSubspaceReadNext);
}

bool subspace_free_telemetry(SubspaceTelemetry *telemetry) {
subspace_clear_error();
if (telemetry == nullptr || telemetry->telemetry == nullptr) {
subspace_set_error("Invalid telemetry parameter");
return false;
}
delete reinterpret_cast<TelemetryStorage *>(telemetry->telemetry);
*telemetry = {};
return true;
}

SubspaceMessage subspace_find_message(SubspaceSubscriber subscriber,
uint64_t timestamp) {
subspace_clear_error();
Expand Down
41 changes: 41 additions & 0 deletions c_client/subspace.h
Original file line number Diff line number Diff line change
Expand Up @@ -182,6 +182,41 @@ typedef struct {
bool checksum_error;
} SubspaceMessage;

typedef enum {
kSubspaceTelemetryNoChange = 0,
kSubspaceTelemetryAdded = 1,
kSubspaceTelemetryRemoved = 2,
} SubspaceTelemetryChange;

typedef struct {
SubspaceString name;
SubspaceTelemetryChange change;
} SubspaceTelemetryParticipant;

typedef struct {
int32_t num_drops;
} SubspaceTelemetryDrop;

typedef struct {
int64_t new_size;
} SubspaceTelemetryResize;

// A decoded server-generated telemetry message. The telemetry pointer owns the
// arrays and strings exposed by this struct. Release a non-empty result with
// subspace_free_telemetry. A null telemetry pointer means no message is
// currently available.
typedef struct {
void *telemetry;
const SubspaceTelemetryParticipant *publishers;
size_t num_publishers;
const SubspaceTelemetryParticipant *subscribers;
size_t num_subscribers;
const SubspaceTelemetryDrop *drops;
size_t num_drops;
const SubspaceTelemetryResize *resizes;
size_t num_resizes;
} SubspaceTelemetry;

typedef struct {
void *slot;
} SubspaceMessageSlot;
Expand Down Expand Up @@ -269,6 +304,7 @@ typedef struct {
int32_t subscriber_queue_size;
bool bridge; // This subscriber is for the bridge.
bool for_tunnel; // Mark subscriptions for external tunnels.
bool telemetry; // Subscribe to server-generated telemetry.
SubspaceTypeInfo type; // Type of the message. This is an opaque string.
int max_active_messages; // Max number of message that can be active at once.
int32_t max_subscribers; // 0 means no explicit subscriber limit.
Expand Down Expand Up @@ -413,6 +449,11 @@ SubspaceMessage subspace_read_message_with_mode(SubspaceSubscriber subscriber,
SubspaceMessage subspace_read_message_with_mode_and_trigger(
SubspaceSubscriber subscriber, SubspaceReadMode mode,
SubspaceClearTrigger clear_trigger);
SubspaceTelemetry
subspace_read_telemetry_message(SubspaceSubscriber subscriber);
SubspaceTelemetry subspace_read_telemetry_message_with_mode(
SubspaceSubscriber subscriber, SubspaceReadMode mode);
bool subspace_free_telemetry(SubspaceTelemetry *telemetry);
SubspaceMessage subspace_find_message(SubspaceSubscriber subscriber,
uint64_t timestamp);
bool subspace_get_all_messages(SubspaceSubscriber subscriber,
Expand Down
Loading
Loading