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 inspectionmove/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",...}\nThe 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.
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
@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.
@type queue() :: :incoming | :active | :deferred | :hold | :corrupt
@type read_error() :: File.posix() | :invalid_header | :checksum_mismatch | :invalid_envelope | :invalid_record
@opaque writer()
A message being written. Use the returned writer after each call.
Functions
@spec abort(writer()) :: :ok
Discards a message that is being written.
@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.
@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.
@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.
@spec list(Path.t(), queue()) :: {:ok, [String.t()]} | {:error, File.posix()}
Lists the queue IDs in queue, oldest first.
@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 totrue. Without it, only the header, envelope, and records are read, which is much faster for large messages.:records- read the records. Defaults totrue.
@spec move(Path.t(), String.t(), queue(), queue()) :: :ok | {:error, File.posix()}
Moves message id from one queue to another.
@spec open(Path.t(), Sovite.Queue.Envelope.t()) :: {:ok, writer()} | {:error, File.posix()}
Starts writing a message for envelope in directory.
Returns the path of message id in queue.
@spec queues() :: [queue()]
The queues, in the order a message normally passes through them.
@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.
@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.
@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.
@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.
@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.
@spec write(writer(), iodata()) :: {:ok, writer()} | {:error, File.posix()}
Appends message data.