Sovite.Core.QueueManager (sovite v0.2.0)

Copy Markdown View Source

Schedules queued messages for delivery, like Postfix's qmgr.

Lifecycle

On startup, messages left in active/ by a crash or kill -9 go back to incoming/ (Sovite.Queue.Spool.recover/1), and the retry times of deferred/ messages are read from their records.

A message is picked up from incoming/ (right after it is queued, and by a periodic scan) or from deferred/ (when its retry time comes) and moved to active/. Its pending recipients are routed (Sovite.Core.Router) and grouped into jobs: one per destination, with at most delivery.max_recipients recipients each. Jobs run in Sovite.Core.Delivery workers.

Each job's results are appended to the queue file and fsynced before anything else happens, so a crash never loses a delivery result: at worst, the recipients of the jobs in flight are delivered again, which SMTP allows. When all jobs of a message are done:

  1. Failed recipients are reported to the sender (Sovite.Core.Bounce).
  2. If no recipient is pending, the message is deleted.
  3. If the message is older than queue.max_lifetime, its pending recipients fail and are reported, and it is deleted.
  4. Otherwise it is deferred: a delay warning is sent if due (queue.delay_warning), and the next attempt is scheduled with exponential backoff (Sovite.Queue.Backoff).

Limits

  • delivery.max_deliveries workers run at once.
  • delivery.destination_concurrency of them per destination.
  • delivery.destination_rate_delay between two deliveries to the same destination.

A worker that finishes a job may get the next job for the same destination and send it over the same connection.

Summary

Functions

Returns a specification to start this module under a supervisor.

Retries every deferred message now, like postqueue -f.

Tells the queue manager that message queue_id was queued in incoming/.

Queue manager options from the running configuration (repo is Sovite's database, for database tables). start_link/1 takes these, plus

Counts of messages and deliveries, for status output and tests.

Functions

child_spec(init_arg)

Returns a specification to start this module under a supervisor.

See Supervisor.

flush(server)

@spec flush(GenServer.server()) :: :ok

Retries every deferred message now, like postqueue -f.

notify(server, queue_id)

@spec notify(GenServer.server(), String.t()) :: :ok

Tells the queue manager that message queue_id was queued in incoming/.

opts(config, repo \\ nil)

@spec opts(Sovite.Core.Config.t(), Sovite.Core.Repo.t() | nil) :: keyword()

Queue manager options from the running configuration (repo is Sovite's database, for database tables). start_link/1 takes these, plus:

  • :name - registered name, or nil for none. Defaults to this module.
  • :resolver - DNS resolver. Defaults to Sovite.DNS.default_resolver/0.
  • :port - SMTP port for MX and address-literal deliveries. Defaults to 25.
  • :client - extra Sovite.SMTP.Client.connect/3 options.
  • :scan_interval - milliseconds between scans of incoming/. Defaults to one minute.
  • :max_active - messages in memory at once. Defaults to 10000.

start_link(opts)

@spec start_link(keyword()) :: GenServer.on_start()

stats(server)

@spec stats(GenServer.server()) :: %{
  active: non_neg_integer(),
  deferred: non_neg_integer(),
  deliveries: non_neg_integer()
}

Counts of messages and deliveries, for status output and tests.