OSSIA
Open Scenario System for Interactive Application
Loading...
Searching...
No Matches
line_framing.hpp
1#pragma once
2#include <ossia/detail/pod_vector.hpp>
3#include <ossia/network/sockets/writers.hpp>
4
5#include <boost/asio/buffer.hpp>
6#include <boost/asio/error.hpp>
7#include <boost/asio/read.hpp>
8#include <boost/asio/read_until.hpp>
9#include <boost/asio/streambuf.hpp>
10#include <boost/asio/write.hpp>
11#include <boost/endian/conversion.hpp>
12
13#include <algorithm>
14#include <cstring>
15#include <memory>
16#include <string_view>
17
18namespace ossia::net
19{
20
24template <std::size_t N>
25inline std::size_t set_line_delimiter(char (&buffer)[N], std::string_view delim) noexcept
26{
27 static_assert(N >= 1);
28 const std::size_t sz = std::min(delim.size(), N - 1);
29 std::copy_n(delim.data(), sz, buffer);
30 std::fill(buffer + sz, buffer + N, '\0');
31 return sz;
32}
33
34template <typename Socket>
35struct line_framing_decoder
36{
37 using buffer_type = std::vector<char, ossia::pod_allocator_avx2<char>>;
38
39 Socket& socket;
40 char delimiter[8] = {0};
41 int32_t m_next_packet_size{};
42
44
58 std::shared_ptr<buffer_type> m_data{std::make_shared<buffer_type>()};
59 uint8_t m_delimiter_len = 0;
60 lifetime_token m_lifetime;
61
62 explicit line_framing_decoder(Socket& socket)
63 : socket{socket}
64 {
65 m_data->reserve(65535);
66 }
67
68 void set_delimiter(std::string_view delim) noexcept
69 {
70 m_delimiter_len = uint8_t(set_line_delimiter(delimiter, delim));
71 }
72
73 template <typename F>
74 void receive(F f)
75 {
76 if(m_delimiter_len == 0)
77 m_delimiter_len = strnlen(delimiter, 8);
78
79 // Receive until delimiter. m_data starts with the leftover of the previous
80 // read, which async_read_until scans before asking the socket for more.
81 boost::asio::async_read_until(
82 socket, boost::asio::dynamic_buffer(*m_data), (const char*)delimiter,
83 [this, alive = m_lifetime.watch(), buf = m_data,
84 f = std::move(f)](boost::system::error_code ec, std::size_t sz) mutable {
85 // The socket may be gone since this read was armed; see lifetime_token.
86 // `buf` is held only to keep the read buffer alive that long.
87 if(alive.expired())
88 return;
89
90 // Errors are the stream's business: validate_stream is what turns an EOF
91 // into a close notification.
92 read_data(std::move(f), ec, sz);
93 });
94 }
95
97 template <typename F>
98 void read_data(F&& f, boost::system::error_code ec, std::size_t sz)
99 {
100 if(!f.validate_stream(ec))
101 return;
102
103 // Any other failure is persistent on this socket, so stop instead of
104 // spinning on an immediately-failing read.
105 if(ec.failed())
106 return;
107
108 if(sz > 0)
109 dispatch_lines(f, sz);
110
111 this->receive(std::move(f));
112 }
113
115 template <typename F>
116 void dispatch_lines(const F& f, std::size_t sz)
117 {
118 auto& data = *m_data;
119 std::size_t consumed = 0;
120 std::size_t end = sz;
121 for(;;)
122 {
123 // Empty lines carry nothing: consumed, but not dispatched.
124 if(end > consumed + m_delimiter_len)
125 {
126 try
127 {
128 f((const unsigned char*)data.data() + consumed,
129 end - consumed - m_delimiter_len);
130 }
131 catch(...)
132 {
133 }
134 }
135 consumed = end;
136
137 // Next line of the same read, if any: the search resumes where we left.
138 const char* const begin = data.data();
139 const char* const last = begin + data.size();
140 const char* const next
141 = std::search(begin + consumed, last, delimiter, delimiter + m_delimiter_len);
142 if(next == last)
143 break;
144 end = std::size_t(next - begin) + m_delimiter_len;
145 }
146
147 data.erase(data.begin(), data.begin() + consumed);
148 }
149};
150
151template <typename Socket>
152struct line_framing_encoder
153{
154 Socket& socket;
155 char delimiter[8] = {0};
156 uint8_t delimiter_len = 0;
157
158 void set_delimiter(std::string_view delim) noexcept
159 {
160 delimiter_len = uint8_t(set_line_delimiter(delimiter, delim));
161 }
162
163 void write(const char* data, std::size_t sz)
164 {
165 if(delimiter_len == 0)
166 delimiter_len = strnlen(delimiter, 8);
167
168 // Scatter-gather: data + delimiter in single write
169 std::array<boost::asio::const_buffer, 2> bufs = {
170 boost::asio::buffer(data, sz),
171 boost::asio::buffer(delimiter, delimiter_len)};
172 this->do_write(socket, bufs);
173 }
174
175 // Regular socket: scatter-gather (single syscall)
176 template <typename T, std::size_t N>
177 void do_write(T& sock, const std::array<boost::asio::const_buffer, N>& bufs)
178 {
179 boost::asio::write(sock, bufs);
180 }
181
182 // Multi socket: write each buffer to each socket
183 template <typename T, std::size_t N>
184 void do_write(
185 multi_socket_writer<T>& sock,
186 const std::array<boost::asio::const_buffer, N>& bufs)
187 {
188 for(const auto& buf : bufs)
189 sock.write(buf);
190 }
191};
192
193struct line_framing
194{
195 template <typename Socket>
196 using encoder = line_framing_encoder<Socket>;
197 template <typename Socket>
198 using decoder = line_framing_decoder<Socket>;
199};
200
201}