MOQX.Transport behaviour (moqx v0.10.0)

Copy Markdown View Source

QUIC transport boundary for MOQT-family implementations.

Protocol code should use this façade instead of backend modules directly. New code uses caller-owned %MOQX.Transport.Context{} values and opaque wrapper handles.

Summary

Functions

Aborts local receive side of stream.

Aborts local send side of stream.

Accepts server connection from listener.

Accepts peer stream and records exact stream metadata.

Returns normalized capabilities for a negotiated transport connection.

Closes connection with application error code.

Closes listener where backend supports it.

Connects client connection through context backend.

Transfers ownership of every known backend handle in context to pid.

Gracefully finishes local send side of stream.

Completes backend handshake.

Starts listener through context backend.

Returns local address for listener or connection where backend supports it.

Creates caller-owned transport context for backend module.

Normalizes one already-received backend message through a caller-owned context.

Normalizes a message for a known connection, adopting a peer-created stream when an active backend reports :new_stream before accept_stream/4 runs.

Opens local stream and records exact stream metadata.

Receives one backend message and normalizes it through context backend.

Receives bytes from stream in passive mode.

Schedules an unreliable datagram on connection.

Schedules bytes on stream from the context owner.

Configures active stream delivery.

Returns exact metadata for a stream wrapper.

Types

connection()

@type connection() :: term()

event()

@type event() ::
  {:listener_event, listener() | connection(), atom(), term()}
  | {:connection_event, connection(), atom(), term()}
  | {:stream_event, stream(), atom(), term()}
  | {:stream_data, stream(), binary(), map()}
  | {:datagram, connection(), binary(), term()}

listener()

@type listener() :: term()

stream()

@type stream() :: term()

Callbacks

abort_receiving(stream, non_neg_integer)

@callback abort_receiving(stream(), non_neg_integer()) :: :ok | {:error, term()}

abort_sending(stream, non_neg_integer)

@callback abort_sending(stream(), non_neg_integer()) :: :ok | {:error, term()}

accept(listener, opts, timeout)

@callback accept(listener(), opts :: keyword() | map(), timeout()) ::
  {:ok, connection()} | {:error, term()}

accept_stream(connection, opts, timeout)

@callback accept_stream(connection(), opts :: keyword() | map(), timeout()) ::
  {:ok, stream()} | {:error, term()}

capabilities(connection)

@callback capabilities(connection()) :: MOQX.Transport.Capabilities.t() | {:error, term()}

close_connection(connection, non_neg_integer)

@callback close_connection(connection(), non_neg_integer()) :: :ok | {:error, term()}

close_listener(listener, timeout)

(optional)
@callback close_listener(listener(), timeout()) :: :ok | {:error, term()}

connect(arg1, port_number, arg3, timeout)

@callback connect(
  String.t() | :inet.ip_address(),
  :inet.port_number(),
  keyword() | map(),
  timeout()
) ::
  {:ok, connection()} | {:error, term()}

controlling_process(arg1, pid)

@callback controlling_process(listener() | connection() | stream(), pid()) ::
  :ok | {:error, term()}

finish_sending(stream)

@callback finish_sending(stream()) :: :ok | {:error, term()}

handshake(connection, timeout)

@callback handshake(connection(), timeout()) :: {:ok, connection()} | {:error, term()}

listen(port, opts)

@callback listen(port :: non_neg_integer() | String.t(), opts :: keyword() | map()) ::
  {:ok, listener()} | {:error, term()}

local_address(arg1)

(optional)
@callback local_address(listener() | connection()) ::
  {:ok, {:inet.ip_address(), :inet.port_number()}} | {:error, term()}

normalize_message(term)

@callback normalize_message(term()) :: event() | :unknown

open_stream(connection, opts)

@callback open_stream(connection(), opts :: keyword() | map()) ::
  {:ok, stream()} | {:error, term()}

recv_stream(stream, byte_count)

@callback recv_stream(stream(), byte_count :: non_neg_integer()) ::
  {:ok, binary()} | {:error, term()}

send_datagram(connection, binary)

@callback send_datagram(connection(), binary()) :: :ok | {:error, term()}

send_datagram(connection, binary, opts)

(optional)
@callback send_datagram(connection(), binary(), opts :: keyword() | map()) ::
  :ok | {:error, term()}

send_stream(stream, iodata, opts)

@callback send_stream(stream(), iodata(), opts :: keyword() | map()) ::
  :ok | {:error, term()}

set_active(stream, arg2)

@callback set_active(stream(), boolean() | :once | non_neg_integer()) ::
  :ok | {:error, term()}

stream_info(stream, arg2, arg3)

(optional)
@callback stream_info(stream(), :client | :server, :local | :peer) ::
  {:ok, MOQX.Transport.Conn.Stream.Info.t()} | {:error, term()}

Functions

abort_receiving(ctx, stream, error_code)

Aborts local receive side of stream.

Intent: caller no longer wants to receive bytes on this stream. QUIC mapping: STOP_SENDING with application error code. Peer observation: peer receives :peer_aborted_receiving. Completion: returns after backend accepts request; lifecycle events arrive later.

abort_sending(ctx, stream, error_code)

Aborts local send side of stream.

Intent: caller cannot or will not finish sending bytes. QUIC mapping: RESET_STREAM with application error code. Peer observation: peer receives :peer_aborted_sending. Completion: returns after backend accepts request; lifecycle events arrive later.

accept(ctx, listener, opts \\ [], timeout \\ :infinity)

Accepts server connection from listener.

accept_stream(ctx, connection, opts \\ [], timeout \\ :infinity)

Accepts peer stream and records exact stream metadata.

capabilities(ctx, connection)

Returns normalized capabilities for a negotiated transport connection.

close_connection(ctx, connection, error_code)

Closes connection with application error code.

Intent: caller closes whole transport connection. QUIC mapping: CONNECTION_CLOSE with application error code. Peer observation: peer receives connection :closed event where backend exposes it. Completion: returns after backend accepts request; lifecycle events arrive later.

close_listener(ctx, listener, timeout \\ 0)

Closes listener where backend supports it.

connect(ctx, host, port, opts \\ [], timeout \\ 5000)

Connects client connection through context backend.

controlling_process(ctx, pid)

Transfers ownership of every known backend handle in context to pid.

finish_sending(ctx, stream)

Gracefully finishes local send side of stream.

Intent: caller has sent all bytes successfully. QUIC mapping: FIN. Peer observation: peer receives :peer_finished_sending. Completion: returns after backend accepts request; lifecycle events arrive later.

handshake(ctx, connection, timeout \\ 5000)

Completes backend handshake.

listen(ctx, port, opts \\ [])

Starts listener through context backend.

local_address(ctx, listener)

Returns local address for listener or connection where backend supports it.

new(backend, opts \\ [])

@spec new(module(), keyword() | map()) ::
  {:ok, MOQX.Transport.Context.t()} | {:error, term()}

Creates caller-owned transport context for backend module.

normalize_event(ctx, message)

@spec normalize_event(MOQX.Transport.Context.t(), term()) ::
  {:ok, event(), MOQX.Transport.Context.t()}
  | {:error, term(), MOQX.Transport.Context.t()}
  | {:unknown, term(), MOQX.Transport.Context.t()}

Normalizes one already-received backend message through a caller-owned context.

This is the non-blocking counterpart to receive_event/2 for process runtimes whose receive loop already removed the message from the mailbox.

normalize_event(ctx, connection, message)

@spec normalize_event(MOQX.Transport.Context.t(), MOQX.Transport.Conn.t(), term()) ::
  {:ok, event(), MOQX.Transport.Context.t()}
  | {:error, term(), MOQX.Transport.Context.t()}
  | {:unknown, term(), MOQX.Transport.Context.t()}

Normalizes a message for a known connection, adopting a peer-created stream when an active backend reports :new_stream before accept_stream/4 runs.

open_stream(ctx, connection, opts \\ [])

Opens local stream and records exact stream metadata.

receive_event(ctx, timeout \\ :infinity)

Receives one backend message and normalizes it through context backend.

recv_stream(ctx, stream, byte_count)

Receives bytes from stream in passive mode.

send_datagram(ctx, connection, data)

Schedules an unreliable datagram on connection.

Completion, loss, or cancellation is reported asynchronously by backend connection events where available.

send_stream(ctx, stream, data, opts \\ [])

Schedules bytes on stream from the context owner.

Returns a send token once the backend accepts the send request. Context-owned receive loops can observe later completion or cancellation events, but token correlation is stream-local. Use MOQX.Transport.Conn.Stream.Sender when the caller needs completion feedback as backend credit for accepted sends.

Pass finish: true to attach FIN to this accepted payload.

set_active(ctx, stream, active)

Configures active stream delivery.

stream_info(ctx, stream)

Returns exact metadata for a stream wrapper.