Sovite.Queue.Spool (sovite v0.2.0)

Copy Markdown View Source

Durable storage for queued messages.

A message is streamed into tmp/, then committed: the file is fsynced, renamed into incoming/, and incoming/ itself is fsynced. Only after commit/1 returns may the server reply 250. A crash at any earlier point leaves at most a file in tmp/, which init/1 deletes.

Each message is one file, named by its queue ID, in one of these queues:

<directory>/tmp/        being received
<directory>/incoming/   accepted, not yet picked up for delivery
<directory>/active/     being delivered
<directory>/deferred/   waiting for the next attempt
<directory>/hold/       held by the administrator, not delivered
<directory>/corrupt/    failed verification, kept for inspection

move/4 moves a file between queues with an atomic rename, so after a crash every message is in exactly one queue. recover/1 puts messages that were being delivered back in incoming/.

File format

Each queue file is a header line, the envelope as one line of JSON, the message, and any number of delivery records:

SOVITE-QUEUE 2 <envelope bytes:10> <message bytes:20> <sha256:64>\n
{"queue_id":"...","sender":"...","recipients":[...],...}\n
<message, CRLF line endings, as received>
R <sha256:64> {"type":"recipient",...}\n
R <sha256:64> {"type":"retry",...}\n

The header has a fixed length. It is written last, so a file whose header does not parse was never committed. The SHA-256 in the header covers the envelope and message. Files are created with mode 0600.

Records (Sovite.Queue.Record) are appended with append/3 and fsynced; each line carries the SHA-256 of its JSON. Only the last line can be incomplete, after a crash during an append: it is ignored and overwritten by the next append. Version 1 files have no records.

Summary

Types

A queue file read by load/2. The message is the message_size bytes at message_offset; end_offset is where the next record goes.

A message being written. Use the returned writer after each call.

Functions

Discards a message that is being written.

Appends delivery records to the queue file at path and fsyncs it. end_offset comes from load/2 or the previous append/3; anything after it (an incomplete record from a crash) is overwritten. Returns the new end offset.

Makes the message durable and moves it to incoming/. Returns the final path and the message size in bytes.

Creates the spool directories with mode 0700 and deletes files left in tmp/ by an earlier crash.

Lists the queue IDs in queue, oldest first.

Reads the queue file at path: envelope, message position, and delivery records.

Moves message id from one queue to another.

Starts writing a message for envelope in directory.

Returns the path of message id in queue.

The queues, in the order a message normally passes through them.

Reads and verifies the queue file at path. Returns the envelope and the byte offset at which the message starts. The checksum is verified without loading the message into memory. Records are not read; use load/2 for those.

Returns the header fields of the message in the queue file at path, up to the empty line that ends them (not included), and at most limit bytes, cut at the end of a line.

Moves every message in active/ back to incoming/. Run at startup: those messages were being delivered when the previous run stopped. Returns the number of messages moved.

Deletes message id from queue, for example after delivery, and emits [:sovite, :queue, :message, :removed] with reason.

Streams the message of the queue file at path, as returned by load/2, in chunks of up to 64 KiB.

Appends message data.

Types

loaded()

@type loaded() :: %{
  envelope: Sovite.Queue.Envelope.t(),
  message_offset: non_neg_integer(),
  message_size: non_neg_integer(),
  records: [Sovite.Queue.Record.t()],
  end_offset: non_neg_integer()
}

A queue file read by load/2. The message is the message_size bytes at message_offset; end_offset is where the next record goes.

queue()

@type queue() :: :incoming | :active | :deferred | :hold | :corrupt

read_error()

@type read_error() ::
  File.posix()
  | :invalid_header
  | :checksum_mismatch
  | :invalid_envelope
  | :invalid_record

writer()

@opaque writer()

A message being written. Use the returned writer after each call.

Functions

abort(writer)

@spec abort(writer()) :: :ok

Discards a message that is being written.

append(path, end_offset, records)

@spec append(Path.t(), non_neg_integer(), [Sovite.Queue.Record.t()]) ::
  {:ok, non_neg_integer()} | {:error, File.posix()}

Appends delivery records to the queue file at path and fsyncs it. end_offset comes from load/2 or the previous append/3; anything after it (an incomplete record from a crash) is overwritten. Returns the new end offset.

commit(writer)

@spec commit(writer()) :: {:ok, Path.t(), non_neg_integer()} | {:error, File.posix()}

Makes the message durable and moves it to incoming/. Returns the final path and the message size in bytes.

On error the temporary file is deleted.

init(directory)

@spec init(Path.t()) :: :ok | {:error, File.posix()}

Creates the spool directories with mode 0700 and deletes files left in tmp/ by an earlier crash.

list(directory, queue)

@spec list(Path.t(), queue()) :: {:ok, [String.t()]} | {:error, File.posix()}

Lists the queue IDs in queue, oldest first.

load(path, opts \\ [])

@spec load(Path.t(), keyword()) :: {:ok, loaded()} | {:error, read_error()}

Reads the queue file at path: envelope, message position, and delivery records.

Options

  • :verify - check the message checksum. Defaults to true. Without it, only the header, envelope, and records are read, which is much faster for large messages.
  • :records - read the records. Defaults to true.

move(directory, id, from, to)

@spec move(Path.t(), String.t(), queue(), queue()) :: :ok | {:error, File.posix()}

Moves message id from one queue to another.

open(directory, envelope)

@spec open(Path.t(), Sovite.Queue.Envelope.t()) ::
  {:ok, writer()} | {:error, File.posix()}

Starts writing a message for envelope in directory.

path(directory, queue, id)

@spec path(Path.t(), queue(), String.t()) :: Path.t()

Returns the path of message id in queue.

queues()

@spec queues() :: [queue()]

The queues, in the order a message normally passes through them.

read(path)

@spec read(Path.t()) ::
  {:ok, Sovite.Queue.Envelope.t(), non_neg_integer()} | {:error, read_error()}

Reads and verifies the queue file at path. Returns the envelope and the byte offset at which the message starts. The checksum is verified without loading the message into memory. Records are not read; use load/2 for those.

read_headers(path, offset, size, limit \\ 65536)

@spec read_headers(Path.t(), non_neg_integer(), non_neg_integer(), pos_integer()) ::
  {:ok, binary()} | {:error, File.posix()}

Returns the header fields of the message in the queue file at path, up to the empty line that ends them (not included), and at most limit bytes, cut at the end of a line.

recover(directory)

@spec recover(Path.t()) :: {:ok, non_neg_integer()} | {:error, File.posix()}

Moves every message in active/ back to incoming/. Run at startup: those messages were being delivered when the previous run stopped. Returns the number of messages moved.

remove(directory, queue, id, reason)

@spec remove(Path.t(), queue(), String.t(), atom()) :: :ok | {:error, File.posix()}

Deletes message id from queue, for example after delivery, and emits [:sovite, :queue, :message, :removed] with reason.

stream_message(path, offset, size)

@spec stream_message(Path.t(), non_neg_integer(), non_neg_integer()) ::
  Enumerable.t(binary())

Streams the message of the queue file at path, as returned by load/2, in chunks of up to 64 KiB.

write(writer, data)

@spec write(writer(), iodata()) :: {:ok, writer()} | {:error, File.posix()}

Appends message data.