178 lines
5.6 KiB
C++
178 lines
5.6 KiB
C++
|
|
/// @file dpf/net/policy.hpp
|
||
|
|
/// @brief Explicit transport, framing, socket, window, and deadline policy.
|
||
|
|
/// @details Library code never reads the environment. Every knob here is a
|
||
|
|
/// field the caller sets; `app::run_config::from_env` is the one place
|
||
|
|
/// that turns `DPF_*` variables into these values.
|
||
|
|
#ifndef LIBDPF_INCLUDE_DPF_NET_POLICY_HPP__
|
||
|
|
#define LIBDPF_INCLUDE_DPF_NET_POLICY_HPP__
|
||
|
|
|
||
|
|
#include <chrono>
|
||
|
|
#include <cstddef>
|
||
|
|
#include <cstdint>
|
||
|
|
#include <stdexcept>
|
||
|
|
#include <string>
|
||
|
|
|
||
|
|
namespace dpf
|
||
|
|
{
|
||
|
|
namespace net
|
||
|
|
{
|
||
|
|
|
||
|
|
/// @brief Which byte transport carries a party-to-party edge.
|
||
|
|
enum class transport : unsigned char
|
||
|
|
{
|
||
|
|
memory_sink, ///< in-process paired `memory_sink` (no stream arrays)
|
||
|
|
memory_stream, ///< in-process synchronous `stream_array`
|
||
|
|
async_memory, ///< in-process event-driven stream array
|
||
|
|
mux, ///< all lanes on one TCP connection
|
||
|
|
parallel, ///< one TCP connection per lane
|
||
|
|
sctp, ///< one SCTP association, lane i = SCTP stream i (Linux)
|
||
|
|
local ///< one unix-domain socket per lane (same host)
|
||
|
|
};
|
||
|
|
|
||
|
|
inline const char * transport_name(transport t) noexcept
|
||
|
|
{
|
||
|
|
switch (t)
|
||
|
|
{
|
||
|
|
case transport::memory_sink:
|
||
|
|
return "memory";
|
||
|
|
case transport::memory_stream:
|
||
|
|
return "stream";
|
||
|
|
case transport::async_memory:
|
||
|
|
return "async";
|
||
|
|
case transport::mux:
|
||
|
|
return "mux";
|
||
|
|
case transport::parallel:
|
||
|
|
return "parallel";
|
||
|
|
case transport::sctp:
|
||
|
|
return "sctp";
|
||
|
|
case transport::local:
|
||
|
|
return "local";
|
||
|
|
}
|
||
|
|
return "async";
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief Parse a transport name. Unknown names throw.
|
||
|
|
inline transport parse_transport(const std::string & s)
|
||
|
|
{
|
||
|
|
if (s == "memory")
|
||
|
|
return transport::memory_sink;
|
||
|
|
if (s == "stream")
|
||
|
|
return transport::memory_stream;
|
||
|
|
if (s == "async" || s == "async_memory")
|
||
|
|
return transport::async_memory;
|
||
|
|
if (s == "mux")
|
||
|
|
return transport::mux;
|
||
|
|
if (s == "parallel")
|
||
|
|
return transport::parallel;
|
||
|
|
if (s == "sctp")
|
||
|
|
return transport::sctp;
|
||
|
|
if (s == "local")
|
||
|
|
return transport::local;
|
||
|
|
throw std::invalid_argument("unknown transport: " + s
|
||
|
|
+ " (memory|stream|async|mux|parallel|sctp|local)");
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief True for transports that cross a kernel socket.
|
||
|
|
inline bool is_socket_transport(transport t) noexcept
|
||
|
|
{
|
||
|
|
return t == transport::mux || t == transport::parallel
|
||
|
|
|| t == transport::sctp || t == transport::local;
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief Whether round payloads on a lane carry a `{round,nbytes}` header.
|
||
|
|
/// @details `automatic` frames only when rounds outnumber lanes. `always`
|
||
|
|
/// frames even with one lane per round (to measure the header, or to
|
||
|
|
/// let partial prefixes and reconnect resume work on every round).
|
||
|
|
/// `never` requires at least one lane per round.
|
||
|
|
enum class framing_mode : unsigned char
|
||
|
|
{
|
||
|
|
automatic,
|
||
|
|
always,
|
||
|
|
never
|
||
|
|
};
|
||
|
|
|
||
|
|
inline const char * framing_name(framing_mode m) noexcept
|
||
|
|
{
|
||
|
|
switch (m)
|
||
|
|
{
|
||
|
|
case framing_mode::automatic:
|
||
|
|
return "auto";
|
||
|
|
case framing_mode::always:
|
||
|
|
return "always";
|
||
|
|
case framing_mode::never:
|
||
|
|
return "never";
|
||
|
|
}
|
||
|
|
return "auto";
|
||
|
|
}
|
||
|
|
|
||
|
|
inline framing_mode parse_framing(const std::string & s)
|
||
|
|
{
|
||
|
|
if (s == "auto" || s == "automatic")
|
||
|
|
return framing_mode::automatic;
|
||
|
|
if (s == "always" || s == "on")
|
||
|
|
return framing_mode::always;
|
||
|
|
if (s == "never" || s == "off")
|
||
|
|
return framing_mode::never;
|
||
|
|
throw std::invalid_argument("unknown framing: " + s + " (auto|always|never)");
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief Options applied to every TCP socket a backend adopts.
|
||
|
|
/// @details Zero means "leave the kernel default".
|
||
|
|
struct socket_options
|
||
|
|
{
|
||
|
|
bool no_delay = true;
|
||
|
|
/// Re-armed after every read and write: Linux clears it after one ACK.
|
||
|
|
bool quickack = true;
|
||
|
|
bool keepalive = true;
|
||
|
|
int keepalive_idle_s = 0;
|
||
|
|
int keepalive_interval_s = 0;
|
||
|
|
int keepalive_count = 0;
|
||
|
|
int send_buffer = 0;
|
||
|
|
int recv_buffer = 0;
|
||
|
|
};
|
||
|
|
|
||
|
|
/// @brief Per-connection wire policy for the async stream backends.
|
||
|
|
struct wire_policy
|
||
|
|
{
|
||
|
|
/// Outstanding-byte high-water mark per lane (per connection on mux and
|
||
|
|
/// SCTP, per socket on parallel). `0` is unlimited.
|
||
|
|
std::size_t window_bytes = std::size_t{1} << 20;
|
||
|
|
/// Largest frame accepted from the peer.
|
||
|
|
std::size_t max_frame = std::size_t{16} << 20;
|
||
|
|
/// Split writes into frames of at most this many payload bytes and
|
||
|
|
/// round-robin them across lanes. `0` sends each write as one frame.
|
||
|
|
std::size_t chunk_bytes = std::size_t{64} << 10;
|
||
|
|
/// Frames and bytes gathered into one write syscall.
|
||
|
|
std::size_t coalesce_frames = 16;
|
||
|
|
std::size_t coalesce_bytes = std::size_t{256} << 10;
|
||
|
|
/// Drop the consumed prefix of an inbox once it passes this many bytes.
|
||
|
|
std::size_t compact_bytes = std::size_t{64} << 10;
|
||
|
|
socket_options socket{};
|
||
|
|
|
||
|
|
void validate() const
|
||
|
|
{
|
||
|
|
if (max_frame == 0)
|
||
|
|
throw std::invalid_argument("wire_policy: max_frame is 0");
|
||
|
|
if (chunk_bytes > max_frame)
|
||
|
|
throw std::invalid_argument(
|
||
|
|
"wire_policy: chunk_bytes exceeds max_frame");
|
||
|
|
if (coalesce_frames == 0)
|
||
|
|
throw std::invalid_argument("wire_policy: coalesce_frames is 0");
|
||
|
|
}
|
||
|
|
};
|
||
|
|
|
||
|
|
/// @brief Bounds for session setup and for draining a write window.
|
||
|
|
struct deadlines
|
||
|
|
{
|
||
|
|
std::chrono::milliseconds join{30000};
|
||
|
|
std::chrono::milliseconds connect{30000};
|
||
|
|
std::chrono::milliseconds accept{30000};
|
||
|
|
std::chrono::milliseconds handshake{30000};
|
||
|
|
std::chrono::milliseconds drain{30000};
|
||
|
|
};
|
||
|
|
|
||
|
|
} // namespace net
|
||
|
|
} // namespace dpf
|
||
|
|
|
||
|
|
#endif // LIBDPF_INCLUDE_DPF_NET_POLICY_HPP__
|