5 #include "sol3/core/msg/exchange_config.h"
10 #include <boost/asio.hpp>
11 #include <boost/asio/local/datagram_protocol.hpp>
20 #include <unordered_set>
50 char const* app_name, cpp::fs::path
const& sol3_root_path);
83 uint32_t peer_port = 0)
override;
84 void configure(msg::ExchangeConfigT
const& config)
override;
88 std::shared_ptr<IBufferConst>
get(uint32_t port, uint32_t idx)
const override;
89 std::shared_ptr<IBufferConst>
get(
90 Endpoint const& ep, uint32_t port, uint32_t idx)
const override;
91 void dispose(std::shared_ptr<IBufferConst>&& buffer)
override;
92 void listPorts(std::vector<EndpointPort>& ports)
const override;
94 std::map<uint32_t, std::vector<std::pair<Endpoint, uint32_t>>>& mappings)
97 std::map<
PortIdx, std::shared_ptr<IBufferConst>>& buffers)
const override;
101 struct ExchangePacket {
104 } __attribute__((packed));
106 static size_t constexpr kMaxRecvBuffer = 8192;
107 static size_t constexpr kMaxPayload =
sizeof(uint64_t);
108 static size_t constexpr kMaxMissedPollsBeforeEvict = 4;
109 static size_t constexpr kMaxBuffers =
110 (kMaxRecvBuffer -
sizeof(ExchangePacket)) /
sizeof(BufferInfo);
113 using PeerSessionId = uint64_t;
115 using Socket = boost::asio::local::datagram_protocol::socket;
116 using ShmemBufferMap = std::map<PortIdx, std::shared_ptr<IBufferConst>>;
122 ShmemBufferMap buffers_to_share;
125 std::vector<Endpoint> peer_endpoints;
128 std::map<Endpoint, PeerSessionId> peer_session_ids;
131 std::map<Endpoint, std::unordered_set<PortIdx>> sent_buffers_per_peer;
136 struct ReceivedState {
140 std::map<Endpoint, PeerSessionId> peer_session_ids;
144 std::set<Endpoint> peer_endpoints;
149 std::map<Endpoint, ShmemBufferMap> buffers_per_peer;
153 std::map<Endpoint, uint32_t> peer_missed_polls;
158 std::map<uint32_t, std::vector<std::pair<Endpoint, uint32_t>>>
159 local_port_to_peer_ports;
174 std::map<PortIdx, ResolvedCacheEntry> resolved_buffers;
177 uint64_t resolve_generation = 0;
180 void bumpResolvedGeneration();
185 std::shared_ptr<IBufferConst> resolveFromPortMappings(
186 uint32_t mapping_local_port,
187 uint32_t requested_local_port,
193 std::shared_ptr<IBufferConst> resolveEntry(
PortIdx local_port_idx)
const;
202 std::shared_ptr<IBufferConst>
get(
203 Endpoint const& ep, uint32_t port, uint32_t idx)
const;
206 template <
typename TFunction>
207 auto accessReceivedState(TFunction&& cb)
const {
208 return received_state_.
access(std::forward<TFunction>(cb));
215 void startPeerPoll();
216 void handleReceive();
217 void sendSomeBuffers(
Endpoint const& sender_ep);
218 void requestPeerBuffers(
Endpoint const& ep);
219 bool sendPacketOnStrand(
222 uint8_t
const* payload =
nullptr,
223 size_t payload_size = 0);
224 void sendExchangePacket(
227 void const* payload =
nullptr,
228 size_t payload_size = 0);
230 void handleBufferRequest(
232 std::array<char, ShmemExchange::kMaxRecvBuffer>
const& msgbuf);
233 void handleIncomingBuffers(
235 std::array<char, ShmemExchange::kMaxRecvBuffer>
const& msgbuf,
236 const struct msghdr& msgh);
237 void handleBufferRemoved(
239 std::array<char, ShmemExchange::kMaxRecvBuffer>
const& msgbuf);
240 void handlePeerDisconnect(
242 std::array<char, ShmemExchange::kMaxRecvBuffer>
const& msgbuf);
244 void sendDisconnect();
246 PeerSessionId peer_session_id_;
247 boost::asio::io_context ctx_;
248 boost::asio::io_context::strand strand_;
249 boost::asio::steady_timer peer_poll_timer_;
254 std::atomic<bool> stop_requested_{
false};
256 std::thread run_thread_;
Interface for registering, discovering, and disposing shared buffers.
Definition: shmem_buffer.h:112
boost::asio::local::datagram_protocol::endpoint Endpoint
Definition: shmem_buffer.h:114
Mutable view of shared memory buffer.
Definition: shmem_buffer.h:90
Definition: shmem_exchange.h:54
void start() override
Start running in a background thread.
void unregisterBuffer(IBufferMutable const &buffer) override
Unregister a previously shared mutable buffer.
Endpoint const & localEndpoint() const override
Local endpoint used by this exchange.
Definition: shmem_exchange.h:87
ShmemExchange(msg::ExchangeConfigT exchange_config)
Create an exchange from an exchange configuration.
void listPorts(std::vector< EndpointPort > &ports) const override
List all ports this exchange has received from peers or shared locally.
void mapLocalToPeer(uint32_t local_port, Endpoint peer_endpoint, uint32_t peer_port=0) override
void listResolvedBuffers(std::map< PortIdx, std::shared_ptr< IBufferConst >> &buffers) const override
IBufferExchange::Endpoint Endpoint
Definition: shmem_exchange.h:56
ShmemExchange(ShmemExchange const &)=delete
bool stopped() const override
True if the exchange has been stopped.
void registerBuffer(IBufferMutable const &buffer) override
~ShmemExchange() override
Stops and releases socket/IO resources.
ShmemExchange & operator=(ShmemExchange &&)=delete
void stop() override
Request the event loop to stop.
ShmemExchange & operator=(ShmemExchange const &)=delete
ShmemExchange(ShmemExchange &&)=delete
std::shared_ptr< IBufferConst > get(uint32_t port, uint32_t idx) const override
ShmemExchange(Endpoint const &local_endpoint)
void addPeer(Endpoint peer_endpoint) override
void dispose(std::shared_ptr< IBufferConst > &&buffer) override
Release resources associated with a previously retrieved buffer.
std::shared_ptr< IBufferConst > get(Endpoint const &ep, uint32_t port, uint32_t idx) const override
void listMappings(std::map< uint32_t, std::vector< std::pair< Endpoint, uint32_t >>> &mappings) const override
Endpoint resolveEndpoint(Endpoint const &endpoint) const override
void configure(msg::ExchangeConfigT const &config) override
auto access(TFunction &&cb) const
Shared/read-only access; executes cb(TData const&) under a shared_lock.
Definition: thread_safe_value.h:34
enum ArtifactFormat string
Camera/host holding the original file.
Definition: any_message_input.h:15
msg::ExchangeConfigT loadExchangeConfigFromSol3Root(char const *app_name, cpp::fs::path const &sol3_root_path)
msg::ExchangeConfigT loadExchangeConfigFromPath(std::string const &config_path)
msg::ExchangeConfigT loadExchangeConfigFromEnv()
boost::asio::local::datagram_protocol::endpoint Endpoint
Definition: shmem_buffer.h:28
ExchangeCode
Definition: shmem_exchange.h:25
@ buffers_request
Definition: shmem_exchange.h:26
@ buffers_response
Definition: shmem_exchange.h:27
@ peer_disconnect
Definition: shmem_exchange.h:29
@ buffer_removed
Definition: shmem_exchange.h:28
ExchangeCode what
Definition: shmem_exchange.h:0
uint64_t peer_session_id
Definition: shmem_exchange.h:1
Definition: shmem_buffer.h:41
Definition: shmem_buffer.h:62
Definition: shmem_exchange.h:166
PortIdx port_idx
Definition: shmem_exchange.h:167
uint64_t generation
Definition: shmem_exchange.h:170
std::shared_ptr< IBufferConst > buffer
Definition: shmem_exchange.h:168
bool cache_needs_update
Definition: shmem_exchange.h:169
Definition: shmem_exchange.h:161
std::shared_ptr< IBufferConst > buffer
Definition: shmem_exchange.h:162
uint64_t generation
Definition: shmem_exchange.h:163