EctoPGMQ.Message (ecto_pgmq v2.0.0)

Copy Markdown View Source

Read-only schema for PGMQ messages.

Summary

Types

A PGMQ message group.

PGMQ message headers.

A PGMQ message ID.

A PGMQ message specification.

A PGMQ message payload.

A PGMQ Message payload type.

t()

A PGMQ message.

Message API

Constructs a message specification for the given payload, group, and headers.

Returns the group/0 for the given message or nil if the given message has no group.

Query API

Returns a query for the messages in the archive for the given queue.

Returns a query for the messages in the given queue.

Types

group()

@type group() :: String.t()

A PGMQ message group.

For more information about FIFO message groups, see FIFO Message Groups.

headers()

@type headers() :: %{optional(String.Chars.t()) => term()}

PGMQ message headers.

id()

@type id() :: pos_integer()

A PGMQ message ID.

message()

@opaque message()

A PGMQ message specification.

payload()

@type payload() :: EctoPGMQ.PGMQ.payload() | term()

A PGMQ message payload.

For more information about valid PGMQ message payloads, see Custom Payload Types.

payload_type()

@type payload_type() :: module() | {module(), keyword()} | :map

A PGMQ Message payload type.

This can take any of the following forms:

For more information about custom PGMQ payloads, see Custom Payload Types guide.

t()

@type t() :: %EctoPGMQ.Message{
  archived_at: DateTime.t() | nil,
  enqueued_at: DateTime.t(),
  headers: headers() | nil,
  id: id(),
  last_read_at: DateTime.t() | nil,
  payload: payload() | nil,
  reads: non_neg_integer(),
  visible_at: DateTime.t()
}

A PGMQ message.

Functions

build(payload, group)

@spec build(payload() | nil, group() | headers() | nil) :: message()

build(payload, group, headers)

@spec build(payload() | nil, group() | nil, headers() | nil) :: message()

Message API

build(payload)

@spec build(payload() | nil) :: message()

Constructs a message specification for the given payload, group, and headers.

Multiple Groups

If the group is not nil, it will override any group that may already be specified in the headers.

Examples

iex> Message.build(%{"id" => 1})
{:message, %{"id" => 1}, nil}

iex> Message.build(%{"id" => 1}, %{"header" => "foo"})
{:message, %{"id" => 1}, %{"header" => "foo"}}

iex> Message.build(%{"id" => 1}, "A")
{:message, %{"id" => 1}, %{"x-pgmq-group" => "A"}}

iex> Message.build(%{"id" => 1}, "A", %{"x-pgmq-group" => "B"})
{:message, %{"id" => 1}, %{"x-pgmq-group" => "A"}}

group(message)

@spec group(t() | message() | headers() | nil) :: group() | nil

Returns the group/0 for the given message or nil if the given message has no group.

Examples

iex> messages = [Message.build(%{"id" => 1}, "A")]
iex> %{"my_queue" => [id]} = EctoPGMQ.send_messages(Repo, "my_queue", messages)
iex> message = Repo.get(Message.queue_query("my_queue"), id)
iex> Message.group(message)
"A"

iex> message_spec = Message.build(%{"id" => 1}, "A")
iex> Message.group(message_spec)
"A"

iex> headers = %{"header" => "foo"}
iex> Message.group(headers)
nil

Query API

archive_query(queue, opts \\ [])

@spec archive_query(EctoPGMQ.Queue.name(), [{:payload_type, payload_type()}]) ::
  Ecto.Query.t()

Returns a query for the messages in the archive for the given queue.

Options

An archive message query can be built with the following options:

  • :payload_type - An optional payload_type/0 for the message payloads. Defaults to :map.

Examples

iex> messages = [Message.build(%{"id" => 1})]
iex> %{"my_queue" => ids} = EctoPGMQ.send_messages(Repo, "my_queue", messages)
iex> EctoPGMQ.archive_messages(Repo, "my_queue", ids)
iex> [%Message{}] = Repo.all(Message.archive_query("my_queue"))

queue_query(queue, opts \\ [])

@spec queue_query(EctoPGMQ.Queue.name(),
  archived_at?: boolean(),
  payload_type: payload_type()
) ::
  Ecto.Query.t()

Returns a query for the messages in the given queue.

Options

A queue message query can be built with the following options:

  • :archived_at? - An optional boolean/0 denoting whether or not to select a NULL :archived_at column. This can be used to make the query structure match that of archive_query/1. Defaults to false.

  • :payload_type - An optional payload_type/0 for the message payloads. Defaults to :map.

Examples

iex> messages = [Message.build(%{"id" => 1})]
iex> EctoPGMQ.send_messages(Repo, "my_queue", messages)
iex> [%Message{}] = Repo.all(Message.queue_query("my_queue"))