Skip to content

Group pipeline

Modules > pipeline

More...

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 .

dp_pull_t * dp_pull_create (
    const char * endpoint
) 

Parameters:

  • endpoint NATS 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.

void dp_pull_destroy (
    dp_pull_t * ctx
) 

Parameters:

  • ctx Pull context (may be NULL).

function dp_pull_recv

Receive one frame from a Pull socket (zero-copy).

int dp_pull_recv (
    dp_pull_t * ctx,
    dp_msg_t ** msg,
    dp_header_t * header
) 

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:

  • ctx Subscriber context.
  • msg Set to a zero-copy message handle.
  • header Set 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.

void dp_pull_set_timeout (
    dp_pull_t * ctx,
    int timeout_ms
) 

Parameters:

  • ctx Pull context.
  • timeout_ms Timeout in milliseconds (-1 = infinite, 0 = non-blocking).

function dp_push_create

Create a Push producer and connect to endpoint .

dp_push_t * dp_push_create (
    const char * endpoint,
    dp_sample_type_t sample_type
) 

Parameters:

  • endpoint NATS endpoint, e.g. "nats://127.0.0.1:4222/work".
  • sample_type Sample 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.

void dp_push_destroy (
    dp_push_t * ctx
) 

Parameters:

  • ctx Push 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:

  • ctx Publisher context.
  • samples Interleaved int32_t I/Q pairs; length 2×num_samples.
  • num_samples Number of complex samples.
  • sample_rate Sample rate in Hz.
  • center_freq Centre 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:

  • ctx Publisher context.
  • samples Interleaved int32_t I/Q pairs; length 2×num_samples.
  • num_samples Number of complex samples.
  • sample_rate Sample rate in Hz.
  • center_freq Centre 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:

  • ctx Publisher context.
  • samples Interleaved int32_t I/Q pairs; length 2×num_samples.
  • num_samples Number of complex samples.
  • sample_rate Sample rate in Hz.
  • center_freq Centre 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:

  • ctx Publisher context.
  • samples Interleaved int32_t I/Q pairs; length 2×num_samples.
  • num_samples Number of complex samples.
  • sample_rate Sample rate in Hz.
  • center_freq Centre 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:

  • ctx Publisher context.
  • samples Interleaved int32_t I/Q pairs; length 2×num_samples.
  • num_samples Number of complex samples.
  • sample_rate Sample rate in Hz.
  • center_freq Centre frequency in Hz.

Returns:

DP_OK (0) on success, negative error code on failure.