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
12 changes: 11 additions & 1 deletion .cursor/skills/subspace-clients/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,10 @@ for (;;) {
auto bytes = static_cast<const char *>(msg_or->buffer);
// Process bytes[0..msg_or->length).
}

// Leave the subscriber trigger fd unread:
// sub.ReadMessage(subspace::ReadMode::kReadNext,
// subspace::ClearTrigger::kNoClearTrigger);
```

## Python Client
Expand All @@ -108,6 +112,9 @@ while True:
if len(data) == 0:
break
# Process data as bytes.

# Leave the subscriber trigger fd unread:
# subscriber.read_message(clear_trigger=subspace.ClearTrigger.NO_CLEAR_TRIGGER)
```

For message metadata and ordinals, use `read_message_object()`:
Expand Down Expand Up @@ -142,7 +149,7 @@ channel subscriber limit.
The Rust crate name is `subspace_client`.

```rust
use subspace_client::{Client, PublisherOptions, ReadMode, SubscriberOptions};
use subspace_client::{ClearTrigger, Client, PublisherOptions, ReadMode, SubscriberOptions};

let client = Client::new("/tmp/subspace", "rust-client").unwrap();

Expand Down Expand Up @@ -173,6 +180,9 @@ loop {
let data = unsafe { msg.as_slice() };
// Process data.
}

// Leave the subscriber trigger fd unread:
// subscriber.read_message_with_trigger(ReadMode::ReadNext, ClearTrigger::NoClearTrigger);
```

For multiple in-flight unpublished slots:
Expand Down
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,11 @@

## Unreleased

### ReadMessage trigger control
- Added `ClearTrigger` (`kClearTrigger` / `kNoClearTrigger`) so `ReadMessage`
can leave the subscriber trigger fd unread. Exposed in C++, C, Python, and
Rust.

### Publisher Buffer Leases
- Added explicit C++, C, Python, and Rust APIs to acquire multiple unpublished
publisher slots, publish or release individual leases, and reject stale lease
Expand Down
28 changes: 26 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -710,6 +710,16 @@ auto msg_or = sub.ReadMessage(subspace::ReadMode::kReadNewest);
// This skips to the most recent message, discarding older ones
```

**Leave the subscriber trigger fd unread**
```cpp
auto msg_or = sub.ReadMessage(subspace::ReadMode::kReadNext,
subspace::ClearTrigger::kNoClearTrigger);
```
By default `ReadMessage` consumes the subscriber trigger fd (eventfd or pipe)
so a later `poll`/`Wait` blocks until a new message is published. Pass
`ClearTrigger::kNoClearTrigger` when the caller is managing that fd from an
external event loop.

**Method 3: Typed read (returns shared_ptr)**
```cpp
auto msg_ptr_or = sub.ReadMessage<MyMessageType>();
Expand Down Expand Up @@ -754,9 +764,13 @@ if (fd_or.ok()) {
class Subscriber {
public:
// Read messages
absl::StatusOr<Message> ReadMessage(ReadMode mode = ReadMode::kReadNext);
absl::StatusOr<Message> ReadMessage(
ReadMode mode = ReadMode::kReadNext,
ClearTrigger clear_trigger = ClearTrigger::kClearTrigger);
template <typename T>
absl::StatusOr<shared_ptr<T>> ReadMessage(ReadMode mode = ReadMode::kReadNext);
absl::StatusOr<shared_ptr<T>> ReadMessage(
ReadMode mode = ReadMode::kReadNext,
ClearTrigger clear_trigger = ClearTrigger::kClearTrigger);

// Find message by timestamp
absl::StatusOr<Message> FindMessage(uint64_t timestamp);
Expand Down Expand Up @@ -1465,6 +1479,12 @@ if (newest.length > 0) {
// Process message
subspace_free_message(&newest);
}

// Leave the subscriber trigger fd unread (for example when poll/epoll
// already consumed it, or another waiter should still observe it).
SubspaceMessage kept = subspace_read_message_with_mode_and_trigger(
sub, kSubspaceReadNext, kSubspaceNoClearTrigger);
subspace_free_message(&kept);
```

**Important:** You must call `subspace_free_message()` when done with a message. The `max_active_messages` option determines how many messages you can hold simultaneously. If you don't free messages, the subscriber will run out of slots and be unable to read more messages.
Expand Down Expand Up @@ -1682,6 +1702,7 @@ This is a quick reference for the most common calls. See
- `SubspaceSubscriber subspace_create_subscriber(SubspaceClient client, const char *channel_name, SubspaceSubscriberOptions options)`
- `SubspaceMessage subspace_read_message(SubspaceSubscriber subscriber)`
- `SubspaceMessage subspace_read_message_with_mode(SubspaceSubscriber subscriber, SubspaceReadMode mode)`
- `SubspaceMessage subspace_read_message_with_mode_and_trigger(SubspaceSubscriber subscriber, SubspaceReadMode mode, SubspaceClearTrigger clear_trigger)`
- `SubspaceMessage subspace_find_message(SubspaceSubscriber subscriber, uint64_t timestamp)`
- `bool subspace_get_all_messages(SubspaceSubscriber subscriber, SubspaceReadMode mode, SubspaceMessage **messages, size_t *count)`
- `bool subspace_free_message(SubspaceMessage *message)`
Expand Down Expand Up @@ -1743,6 +1764,9 @@ let subscriber = client.create_subscriber("sensor_data", &sub_opts)?;
// Read a message.
let msg = subscriber.read_message(ReadMode::ReadNext)?;
assert_eq!(msg.length, 5);

// Leave the subscriber trigger fd unread:
// subscriber.read_message_with_trigger(ReadMode::ReadNext, ClearTrigger::NoClearTrigger)?;
```

### Features
Expand Down
42 changes: 42 additions & 0 deletions c_client/client_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
#include <memory>
#include <signal.h>
#include <string.h>
#include <sys/poll.h>
#include <sys/resource.h>
#include <thread>
#include <unistd.h>
Expand Down Expand Up @@ -371,6 +372,47 @@ TEST_F(ClientTest, PublishSingleMessageAndRead) {
ASSERT_TRUE(subspace_remove_client(&sub_client));
}

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

SubspacePublisher pub = subspace_create_publisher(
pub_client, "c_no_clear", CPublisherOptionsDefault(256, 10));
ASSERT_NE(nullptr, pub.publisher);
SubspaceSubscriber sub = subspace_create_subscriber(
sub_client, "c_no_clear", CSubscriberOptionsDefault());
ASSERT_NE(nullptr, sub.subscriber);

struct pollfd pfd = subspace_get_subscriber_poll_fd(sub);
ASSERT_GT(pfd.fd, 0);

SubspaceMessageBuffer buffer = subspace_get_message_buffer(pub, 6);
ASSERT_NE(nullptr, buffer.buffer);
memcpy(buffer.buffer, "foobar", 6);
const SubspaceMessage pub_status = subspace_publish_message(pub, 6);
ASSERT_NE(0, pub_status.length);

ASSERT_EQ(1, ::poll(&pfd, 1, 1000));

SubspaceMessage msg = subspace_read_message_with_mode_and_trigger(
sub, kSubspaceReadNext, kSubspaceNoClearTrigger);
ASSERT_FALSE(subspace_has_error());
ASSERT_EQ(6, msg.length);
subspace_free_message(&msg);

pfd.revents = 0;
ASSERT_EQ(1, ::poll(&pfd, 1, 0));

ASSERT_TRUE(subspace_remove_subscriber(&sub));
ASSERT_TRUE(subspace_remove_publisher(&pub));
ASSERT_TRUE(subspace_remove_client(&pub_client));
ASSERT_TRUE(subspace_remove_client(&sub_client));
}

TEST_F(ClientTest, MetadataAddressesAndDescriptorHelpers) {
auto pub_client = subspace_create_client_with_socket(Socket().c_str());
ASSERT_NE(nullptr, pub_client.client);
Expand Down
21 changes: 17 additions & 4 deletions c_client/subspace.cc
Original file line number Diff line number Diff line change
Expand Up @@ -146,6 +146,12 @@ subspace::ReadMode ToCppReadMode(SubspaceReadMode mode) {
: subspace::ReadMode::kReadNext;
}

subspace::ClearTrigger ToCppClearTrigger(SubspaceClearTrigger clear_trigger) {
return clear_trigger == kSubspaceNoClearTrigger
? subspace::ClearTrigger::kNoClearTrigger
: subspace::ClearTrigger::kClearTrigger;
}

subspace::ChecksumCallback
ToCppChecksumCallback(SubspaceChecksumCallback callback, void *user_data) {
return [callback,
Expand Down Expand Up @@ -602,8 +608,9 @@ SubspacePublisher subspace_create_publisher(SubspaceClient client,
return publisher;
}

SubspaceMessage subspace_read_message_with_mode(SubspaceSubscriber subscriber,
SubspaceReadMode mode) {
SubspaceMessage subspace_read_message_with_mode_and_trigger(
SubspaceSubscriber subscriber, SubspaceReadMode mode,
SubspaceClearTrigger clear_trigger) {
SubspaceMessage message = EmptyMessage();
if (subscriber.subscriber == nullptr) {
return message;
Expand All @@ -613,8 +620,8 @@ SubspaceMessage subspace_read_message_with_mode(SubspaceSubscriber subscriber,
// subspace::Subscriber.
auto sub_ptr = reinterpret_cast<std::shared_ptr<subspace::Subscriber> *>(
subscriber.subscriber);
absl::StatusOr<subspace::Message> status_or_msg =
(*sub_ptr)->ReadMessage(ToCppReadMode(mode));
absl::StatusOr<subspace::Message> status_or_msg = (*sub_ptr)->ReadMessage(
ToCppReadMode(mode), ToCppClearTrigger(clear_trigger));
if (!status_or_msg.ok()) {
subspace_set_error(status_or_msg.status().ToString().c_str());
return message;
Expand All @@ -629,6 +636,12 @@ SubspaceMessage subspace_read_message_with_mode(SubspaceSubscriber subscriber,
return TakeCMessage(std::move(*status_or_msg));
}

SubspaceMessage subspace_read_message_with_mode(SubspaceSubscriber subscriber,
SubspaceReadMode mode) {
return subspace_read_message_with_mode_and_trigger(subscriber, mode,
kSubspaceClearTrigger);
}

SubspaceMessage subspace_read_message(SubspaceSubscriber subscriber) {
return subspace_read_message_with_mode(subscriber, kSubspaceReadNext);
}
Expand Down
15 changes: 15 additions & 0 deletions c_client/subspace.h
Original file line number Diff line number Diff line change
Expand Up @@ -294,6 +294,14 @@ typedef enum {
kSubspaceReadNewest = 1, // Read the newest message.
} SubspaceReadMode;

// Controls whether a read consumes the subscriber trigger fd (eventfd or
// pipe). kSubspaceClearTrigger is the default and matches existing
// subspace_read_message / subspace_read_message_with_mode behavior.
typedef enum {
kSubspaceClearTrigger = 0, // Read (clear) the subscriber trigger fd.
kSubspaceNoClearTrigger = 1, // Leave the subscriber trigger fd unread.
} SubspaceClearTrigger;

// This is a message buffer that is used to publish a message. The 'buffer'
// member is a pointer to the message buffer that can be used to publish a
// message which has size 'buffer_size' bytes. The buffer is read/write and can
Expand Down Expand Up @@ -398,6 +406,13 @@ bool subspace_remove_client(SubspaceClient *client);
SubspaceMessage subspace_read_message(SubspaceSubscriber subscriber);
SubspaceMessage subspace_read_message_with_mode(SubspaceSubscriber subscriber,
SubspaceReadMode mode);
// Same as subspace_read_message_with_mode, with control over whether the
// subscriber trigger fd is consumed. Pass kSubspaceNoClearTrigger to leave
// the fd unread, for example when the caller is managing it from an
// external event loop.
SubspaceMessage subspace_read_message_with_mode_and_trigger(
SubspaceSubscriber subscriber, SubspaceReadMode mode,
SubspaceClearTrigger clear_trigger);
SubspaceMessage subspace_find_message(SubspaceSubscriber subscriber,
uint64_t timestamp);
bool subspace_get_all_messages(SubspaceSubscriber subscriber,
Expand Down
11 changes: 8 additions & 3 deletions client/client.cc
Original file line number Diff line number Diff line change
Expand Up @@ -1453,9 +1453,12 @@ ClientImpl::ReadMessageInternal(SubscriberImpl *subscriber, ReadMode mode,
}

absl::StatusOr<Message> ClientImpl::ReadMessage(SubscriberImpl *subscriber,
ReadMode mode) {
ReadMode mode,
ClearTrigger clear_trigger) {

ClientLockGuard guard(this);
const bool should_clear_trigger =
clear_trigger == ClearTrigger::kClearTrigger;
// If the channel is a placeholder (no publishers present), look
// in the SCB to see if a new publisher has been created and if so,
// talk to the server to get the information to reload the shared
Expand All @@ -1464,7 +1467,9 @@ absl::StatusOr<Message> ClientImpl::ReadMessage(SubscriberImpl *subscriber,
if (subscriber->IsPlaceholder()) {
absl::Status status = ReloadSubscriber(subscriber);
if (!status.ok() || subscriber->IsPlaceholder()) {
subscriber->ClearPollFd();
if (should_clear_trigger) {
subscriber->ClearPollFd();
}
return Message();
}
subscriber->TriggerReliablePublishers();
Expand All @@ -1479,7 +1484,7 @@ absl::StatusOr<Message> ClientImpl::ReadMessage(SubscriberImpl *subscriber,

return ReadMessageInternal(subscriber, mode,
subscriber->options_.pass_activation,
/*clear_trigger=*/true);
should_clear_trigger);
}

absl::StatusOr<Message>
Expand Down
40 changes: 30 additions & 10 deletions client/client.h
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,16 @@ enum class ReadMode {
kReadNewest,
};

// Controls whether ReadMessage consumes the subscriber trigger fd (eventfd
// or pipe). The default, kClearTrigger, reads the fd so a later poll/Wait
// blocks until a new message is published. kNoClearTrigger leaves the fd
// unread, which is useful when the caller is managing it from an external
// event loop.
enum class ClearTrigger {
kClearTrigger,
kNoClearTrigger,
};

struct ChannelInfo {
std::string channel_name;
int num_publishers;
Expand Down Expand Up @@ -602,15 +612,19 @@ class ClientImpl : public std::enable_shared_from_this<ClientImpl> {
// memory which is read-only. If the read is triggered by the PollFd,
// you must read all the avaiable messages from the subscriber as the
// PollFd is only triggered when a new message is published.
absl::StatusOr<Message> ReadMessage(details::SubscriberImpl *subscriber,
ReadMode mode = ReadMode::kReadNext);
// Pass ClearTrigger::kNoClearTrigger to leave the subscriber trigger fd
// unread.
absl::StatusOr<Message> ReadMessage(
details::SubscriberImpl *subscriber, ReadMode mode = ReadMode::kReadNext,
ClearTrigger clear_trigger = ClearTrigger::kClearTrigger);

// As ReadMessage above but returns a shared_ptr to the typed message.
// NOTE: this is subspace::shared_ptr, not std::shared_ptr.
template <typename T, typename Aliaser = DefaultAliaser>
absl::StatusOr<shared_ptr<T, Aliaser>>
ReadMessage(details::SubscriberImpl *subscriber,
ReadMode mode = ReadMode::kReadNext);
ReadMode mode = ReadMode::kReadNext,
ClearTrigger clear_trigger = ClearTrigger::kClearTrigger);

// Find a message given a timestamp.
absl::StatusOr<Message> FindMessage(details::SubscriberImpl *subscriber,
Expand Down Expand Up @@ -829,8 +843,9 @@ class ClientImpl : public std::enable_shared_from_this<ClientImpl> {
// you need to as it may prevent a publisher getting a slot.
template <typename T, typename Aliaser>
inline absl::StatusOr<::subspace::shared_ptr<T, Aliaser>>
ClientImpl::ReadMessage(details::SubscriberImpl *subscriber, ReadMode mode) {
absl::StatusOr<Message> msg = ReadMessage(subscriber, mode);
ClientImpl::ReadMessage(details::SubscriberImpl *subscriber, ReadMode mode,
ClearTrigger clear_trigger) {
absl::StatusOr<Message> msg = ReadMessage(subscriber, mode, clear_trigger);
if (!msg.ok()) {
return msg.status();
}
Expand Down Expand Up @@ -1442,15 +1457,20 @@ class Subscriber {
// memory which is read-only. If the read is triggered by the PollFd,
// you must read all the avaiable messages from the subscriber as the
// PollFd is only triggered when a new message is published.
absl::StatusOr<Message> ReadMessage(ReadMode mode = ReadMode::kReadNext) {
return client_->ReadMessage(impl_.get(), mode);
// Pass ClearTrigger::kNoClearTrigger to leave the subscriber trigger fd
// unread.
absl::StatusOr<Message>
ReadMessage(ReadMode mode = ReadMode::kReadNext,
ClearTrigger clear_trigger = ClearTrigger::kClearTrigger) {
return client_->ReadMessage(impl_.get(), mode, clear_trigger);
}

// As ReadMessage above but returns a shared_ptr to the typed message.
// NOTE: this is subspace::shared_ptr, not std::shared_ptr.
template <typename T, typename Aliaser = DefaultAliaser>
absl::StatusOr<shared_ptr<T, Aliaser>>
ReadMessage(ReadMode mode = ReadMode::kReadNext);
ReadMessage(ReadMode mode = ReadMode::kReadNext,
ClearTrigger clear_trigger = ClearTrigger::kClearTrigger);

bool AddActiveMessage(int32_t slot_id) {
return impl_->AddActiveMessage(impl_->GetSlot(slot_id));
Expand Down Expand Up @@ -1728,8 +1748,8 @@ class Subscriber {

template <typename T, typename Aliaser>
inline absl::StatusOr<::subspace::shared_ptr<T, Aliaser>>
Subscriber::ReadMessage(ReadMode mode) {
return client_->ReadMessage<T, Aliaser>(impl_.get(), mode);
Subscriber::ReadMessage(ReadMode mode, ClearTrigger clear_trigger) {
return client_->ReadMessage<T, Aliaser>(impl_.get(), mode, clear_trigger);
}

template <typename T, typename Aliaser>
Expand Down
Loading
Loading