Group pipeline¶
Public Functions¶
| Type | Name |
|---|---|
| dp_pull_t * | dp_pull_create (const char * endpoint) Create a Pull consumer and connect to endpoint . |
| void | dp_pull_destroy (dp_pull_t * ctx) Destroy a Pull context and release all resources. |
| int | dp_pull_recv (dp_pull_t * ctx, dp_msg_t ** msg, dp_header_t * header) Receive one frame from a Pull socket (zero-copy). |
| void | dp_pull_set_timeout (dp_pull_t * ctx, int timeout_ms) Set receive timeout for a Pull socket. |
| dp_push_t * | dp_push_create (const char * endpoint, dp_sample_type_t sample_type) Create a Push producer and connect to endpoint . |
| void | dp_push_destroy (dp_push_t * ctx) Destroy a Push context and release all resources. |
| int | dp_push_send_cf32 (dp_push_t * ctx, const float _Complex * samples, size_t num_samples, double sample_rate, double center_freq) Send CF32 samples via a Push socket. |
| int | dp_push_send_cf64 (dp_push_t * ctx, const double _Complex * samples, size_t num_samples, double sample_rate, double center_freq) Send CF64 samples via a Push socket. |
| int | dp_push_send_ci16 (dp_push_t * ctx, const int16_t * samples, size_t num_samples, double sample_rate, double center_freq) Send CI16 samples via a Push socket. |
| int | dp_push_send_ci32 (dp_push_t * ctx, const int32_t * samples, size_t num_samples, double sample_rate, double center_freq) Send CI32 samples via a Push socket. |
| int | dp_push_send_ci8 (dp_push_t * ctx, const int8_t * samples, size_t num_samples, double sample_rate, double center_freq) Send CI8 samples via a Push socket. |
Detailed Description¶
Push sockets distribute work across all connected Pull workers in a round-robin fashion. Unlike PUB/SUB, each frame is delivered to exactly one Pull consumer.
Public Functions Documentation¶
function dp_pull_create¶
Create a Pull consumer and connect to endpoint .
Parameters:
endpointNATS endpoint, e.g."nats://127.0.0.1:4222/work".
Returns:
Non-NULL context on success, NULL on failure.
function dp_pull_destroy¶
Destroy a Pull context and release all resources.
Parameters:
ctxPull context (may be NULL).
function dp_pull_recv¶
Receive one frame from a Pull socket (zero-copy).
On success, *msg is set to a message handle whose data buffer is valid until dp_msg_free() is called. Use dp_msg_data() to access the sample pointer.
Parameters:
ctxSubscriber context.msgSet to a zero-copy message handle.headerSet to the frame metadata.
Returns:
DP_OK on success, DP_ERR_TIMEOUT on timeout, DP_ERR_EOF when the sender has finished (no message is produced, so there is nothing to free), negative on error.
function dp_pull_set_timeout¶
Set receive timeout for a Pull socket.
Parameters:
ctxPull context.timeout_msTimeout in milliseconds (-1 = infinite, 0 = non-blocking).
function dp_push_create¶
Create a Push producer and connect to endpoint .
Parameters:
endpointNATS endpoint, e.g."nats://127.0.0.1:4222/work".sample_typeSample format that will be sent.
Returns:
Non-NULL context on success, NULL on failure.
Note:
On the NATS (nats://) work-queue tier the per-frame payload must fit one message: header + data <= server max_payload (default 1 MiB). Unlike PUB/SUB, PUSH does not chunk (the work-queue load-balances frames across workers, which cannot reassemble a split frame), so an oversized dp_push_send_* returns DP_ERR_TOO_LARGE. Raise the broker max_payload for larger durable frames, or use PUB/SUB (which chunks).
function dp_push_destroy¶
Destroy a Push context and release all resources.
Parameters:
ctxPush context (may be NULL).
function dp_push_send_cf32¶
Send CF32 samples via a Push socket.
int dp_push_send_cf32 (
dp_push_t * ctx,
const float _Complex * samples,
size_t num_samples,
double sample_rate,
double center_freq
)
Parameters:
ctxPublisher context.samplesInterleaved int32_t I/Q pairs; length 2×num_samples.num_samplesNumber of complex samples.sample_rateSample rate in Hz.center_freqCentre frequency in Hz.
Returns:
DP_OK (0) on success, negative error code on failure.
function dp_push_send_cf64¶
Send CF64 samples via a Push socket.
int dp_push_send_cf64 (
dp_push_t * ctx,
const double _Complex * samples,
size_t num_samples,
double sample_rate,
double center_freq
)
Parameters:
ctxPublisher context.samplesInterleaved int32_t I/Q pairs; length 2×num_samples.num_samplesNumber of complex samples.sample_rateSample rate in Hz.center_freqCentre frequency in Hz.
Returns:
DP_OK (0) on success, negative error code on failure.
function dp_push_send_ci16¶
Send CI16 samples via a Push socket.
int dp_push_send_ci16 (
dp_push_t * ctx,
const int16_t * samples,
size_t num_samples,
double sample_rate,
double center_freq
)
Parameters:
ctxPublisher context.samplesInterleaved int32_t I/Q pairs; length 2×num_samples.num_samplesNumber of complex samples.sample_rateSample rate in Hz.center_freqCentre frequency in Hz.
Returns:
DP_OK (0) on success, negative error code on failure.
function dp_push_send_ci32¶
Send CI32 samples via a Push socket.
int dp_push_send_ci32 (
dp_push_t * ctx,
const int32_t * samples,
size_t num_samples,
double sample_rate,
double center_freq
)
Parameters:
ctxPublisher context.samplesInterleaved int32_t I/Q pairs; length 2×num_samples.num_samplesNumber of complex samples.sample_rateSample rate in Hz.center_freqCentre frequency in Hz.
Returns:
DP_OK (0) on success, negative error code on failure.
function dp_push_send_ci8¶
Send CI8 samples via a Push socket.
int dp_push_send_ci8 (
dp_push_t * ctx,
const int8_t * samples,
size_t num_samples,
double sample_rate,
double center_freq
)
Parameters:
ctxPublisher context.samplesInterleaved int32_t I/Q pairs; length 2×num_samples.num_samplesNumber of complex samples.sample_rateSample rate in Hz.center_freqCentre frequency in Hz.
Returns:
DP_OK (0) on success, negative error code on failure.