84#include <cx/net/net_shared.h>
89typedef struct NetFlow_WeakRef NetFlow_WeakRef;
91typedef struct NetFlow_WeakRef NetFlow_WeakRef;
94#define _sti_NetFlow _sti_object
95#define SType_NetFlow NetFlow*
96#define STStorageType_NetFlow NetFlow*
97#define STypeArg_NetFlow(type, val) stgeneric(object, (ObjInst*)objInstCheckClass(NetFlow, val))
98#define STypeArgPtr_NetFlow(type, val) (stgeneric*)objInstCheckClassPtr(NetFlow, val)
99#define STypeCheckedArg_NetFlow(type, val) stType(type), stArg(type, val)
100#define STypeCheckedPtrArg_NetFlow(type, val) stType(type), stArgPtr(type, val)
104typedef struct NetFlow_ClassIf {
110 bool (*close)(_In_
void* self);
111 bool (*send)(_In_
void* self, _In_
const uint8* data,
size_t len, flags_t flags);
113extern NetFlow_ClassIf NetFlow_ClassIf_tmpl;
123 atomic(uintptr) _ref;
124 atomic(ptr) _weakref;
201#define NetFlow(inst) objInstCheckClass(NetFlow, inst)
202#define NetFlowNone ((NetFlow*)NULL)
204typedef struct NetFlow_WeakRef {
207 void* _is_NetFlow_WeakRef;
208 void* _is_ObjInst_WeakRef;
210 atomic(uintptr) _ref;
213#define NetFlow_WeakRef(inst) objWeakRefCheckClass(NetFlow, inst)
217#define netflowCreate(socket, peer) NetFlow_create(NetSocket(socket), peer)
219void NetFlow_setHandlers(_In_
NetFlow* self, _In_opt_
const NetHandlers* handlers, _In_opt_
void* ctx);
230#define netflowSetHandlers(self, handlers, ctx) NetFlow_setHandlers(NetFlow(self), handlers, ctx)
241#define netflowSetHandlersObj(self, handlers, ctx) NetFlow_setHandlersObj(NetFlow(self), handlers, ObjInst(ctx))
260#define netflowAddTimer(self, delay, flags) NetFlow_addTimer(NetFlow(self), delay, flags)
271#define netflowCancelTimer(self, id) NetFlow_cancelTimer(NetFlow(self), id)
281#define netflowRearmTimer(self, id, delay) NetFlow_rearmTimer(NetFlow(self), id, delay)
296#define netflow_push(self, msg) NetFlow__push(NetFlow(self), msg)
303#define netflow_pop(self) NetFlow__pop(NetFlow(self))
312#define netflow_unpop(self, msg) NetFlow__unpop(NetFlow(self), msg)
319#define netflow_close(self, reason) NetFlow__close(NetFlow(self), reason)
326#define netflow_queue(self) NetFlow__queue(NetFlow(self))
328void NetFlow__snapshotFlows(_In_
NetSocket* sock, _Out_ sa_NetFlow* out);
339#define netflow_snapshotFlows(sock, out) NetFlow__snapshotFlows(NetSocket(sock), out)
352#define netflow_addFilter(self, sock, filter) NetFlow__addFilter(NetFlow(self), NetSocket(sock), NetFilter(filter))
359#define netflow_buildFilters(self, sock) NetFlow__buildFilters(NetFlow(self), NetSocket(sock))
361void NetFlow__clearFilters(_In_
NetFlow* self);
365#define netflow_clearFilters(self) NetFlow__clearFilters(NetFlow(self))
367void NetFlow__filterShutdown(_In_
NetFlow* self);
372#define netflow_filterShutdown(self) NetFlow__filterShutdown(NetFlow(self))
382#define netflow_filterNotify(self, q, sock, onWorker) NetFlow__filterNotify(NetFlow(self), NetQueue(q), NetSocket(sock), onWorker)
384bool NetFlow__filterStreamSend(_In_
NetFlow* self, _In_opt_
NetQueue* q, _Inout_
NetSocket* sock, _In_reads_bytes_opt_(len)
const uint8* data,
size_t len);
397#define netflow_filterStreamSend(self, q, sock, data, len) NetFlow__filterStreamSend(NetFlow(self), NetQueue(q), NetSocket(sock), data, len)
407#define netflow_filterStreamRecv(self, q, sock) NetFlow__filterStreamRecv(NetFlow(self), NetQueue(q), NetSocket(sock))
409bool NetFlow__filterDatagramEncode(_In_
NetFlow* self, _In_opt_
NetQueue* q, _In_reads_bytes_opt_(len)
const uint8* data,
size_t len, _Inout_
NetMsgQueue* out, _Out_opt_
bool* fatalp);
418#define netflow_filterDatagramEncode(self, q, data, len, out, fatalp) NetFlow__filterDatagramEncode(NetFlow(self), NetQueue(q), data, len, out, fatalp)
420bool NetFlow__filterDatagramSend(_In_
NetFlow* self, _In_opt_
NetQueue* q, _Inout_
NetSocket* sock, _In_reads_bytes_opt_(len)
const uint8* data,
size_t len);
427#define netflow_filterDatagramSend(self, q, sock, data, len) NetFlow__filterDatagramSend(NetFlow(self), NetQueue(q), NetSocket(sock), data, len)
436#define netflow_filterDatagramRecv(self, q, sock, msg, out) NetFlow__filterDatagramRecv(NetFlow(self), NetQueue(q), NetSocket(sock), msg, out)
445#define netflow_primeFilters(self, q, sock) NetFlow__primeFilters(NetFlow(self), NetQueue(q), NetSocket(sock))
461#define netflowClose(self) (self)->_->close(NetFlow(self))
475#define netflowSend(self, data, len, flags) (self)->_->send(NetFlow(self), data, len, flags)
#define saDeclarePtr(name)
uint64 NetTimerId
Handle to an armed timer, unique for the lifetime of its queue.
#define _objfactory_guaranteed
CX Struct System - Introspectable, serializable POD structures in C.
CX Object System - Object-oriented programming in C.
Socket-level factory for per-flow filters.
A single ordering domain: one connection, or one datagram peer.
sa_stvar user
Application state for this peer, allocated lazily.
atomic(uint32) lastActive
Last time a packet was ingested for this flow, for approximate LRU.
NetMsgQueue encInMsgs
Datagram: staging queue the encode chain consumes from, the encIn of the message side.
uint64 key
QUIC: the stream id this flow is keyed on. Zero on every other kind of flow.
sa_uint64 timers
Ids of the timers currently armed on this flow, guarded by the queue's timerLock.
NetPool * pool
Pool every message on this flow was drawn from, held strongly.
Weak(ObjInst) *handlerWeak
Context passed to per-flow handlers, set by setHandlersObj() – NULL when handlerCtx is in use instead...
NetMessage * ready
Consumer-private FIFO head, refilled from inbox.
sa_NetFlowFilter filters
This flow's filter chain, ordered application -> wire, or empty for none.
Mutex filterLock
Serializes filter chain access for this flow.
atomic(uint32) claimed
A worker is currently draining this flow.
BufRing * encIn
Stream: staging ring the encode chain consumes from, allocated with the chain.
atomic(uint32) queued
Present on the queue's runqueue.
atomic(ptr) inbox
Pending messages for this flow, in an intrusive lock-free stack.
uint8 closeReason
NetCloseReason once the flow is dying.
NetHandlers * handlers
Per-flow handler overrides, optional.
RWLock handlerLock
Guards handlers/handlerCtx/handlerWeak against a concurrent setHandlers()/setHandlersObj()
Weak(NetSocket) *socket
Socket that owns this flow.
Set of event handlers, registered per flow, per socket, or queue-wide.
A received packet, as delivered to a handler.
Simple intrusive FIFO of NetMessage, linked through NetMessage::next.
Shared, capped pool of network message buffers and headers.
NetQueue manages one or more sockets and a thread pool of workers.