This reduces the boilerplate around declaring metrics
(counters, histograms, gauges) in their various forms,
and more or less standardizes them, making the syntax
more regular regardless of how the metrics are actually
stored.
We move the help out to doc comments, making it easier
to write a multi-line exposition on a given metric (in
the future; we're not doing that yet).
linkme is used to form a registry that can be used to
eagerly collect metadata from the various metrics. This
will be used to drive some automated documentation
extraction for the various metrics in a future commit.
This commit implements a kumomta-specific message transfer
protocol that is intended to be used to migrate messages
from one kumomta node to another.
The transfer is carried out using an HTTP POST request
to the destination node's http listener.
The request includes the full message metadata and body,
in a compressed form.
An xfer request can be made via `kcli xfer` (and thus also via an HTTP API
endpoint). It works similarly to a rebind operation; you specify the
criteria to be used to match scheduled queues, along with the target
node for the xfer, and kumomta will find matching queues, drain out the
messages, make an adjustment to the metadata to capture current
scheduling information, and then place the messages into an xfer queue.
The xfer queue has hard-coded scheduling queue configuration of its own,
with the base retry interval set to 10 seconds, which should be suitably
aggressive for the intended use case.
You may apply shaping to affect the number of concurrent requests in a
similar way to how TSA shaping is configured.
On the receiving side, the incoming xfer sanity check to prohibit
trying to xfer to itself.
The spool id of the Message is not suitable to be reused verbatim on
another node (spool ids include the local mac address and creation
timestamp information, as well as a random component), so the receiving
side will derive an id that should be suitable for use on that node.
The originating node id and spool id will be preserved in metadata to
aid in tracing.
It is possible for an xfer request to target an existing xfer queue, so
that you can correct/update the target in various circumstances. In that
situation the messages will be "simply" moved from the source queue to
the destination queue.
It is possible to cancel an xfer request via `kcli xfer-cancel` (and
thus also via an HTTP API endpoint). You specify the target queue,
which must be an xfer queue, and it will have its messages drained and
the metadata changes that were applied when the xfer was initiated will
be reversed, allowing the messages to then be reinserted into their
originating queue.
refs: https://github.com/KumoCorp/kumomta/issues/311
Under high concurrency, a backlog of store operations can build
up, especially if rocksdb isn't tuned for the workload.
In order to work together with the data_processing_timeout of the smtp
server, we'd like to be able to cancel store operations that are past
the deadline, but the spool operations are not guaranteed to be cancel
safe because both the rocksdb and the local disk implementations can
ultimately call into spawn_blocking to perform the critical portion of
the work, and that is not cancel safe.
This means that we cannot simply apply timeout_at to the store
operation: while it may appear to the caller to have been cancelled at
the appropriate time, the actual work storing the spool may be in
progress and continue to store the message to the spool, which could
lead to us thinking that we are accountable for the message when we next
restart, even though we're presumably about to tell the sender that we
did not take the message.
To deal with this situation more safely, we pass the deadline through
to the underlying operation and have it internally decide whether
it should continue.
In practice, this means that the rocksdb store operation will apply
timeout_at to the semaphore operation it does if
`limit_concurrent_stores` has been configured fr that spool. Otherwise
the deadline is ignored.
The implementation of this is imperfect in a number of ways. A more
robust solution would be to employ something like
<https://docs.rs/cancel-safe-futures/0.1.5/cancel_safe_futures/coop_cancel/index.html>,
but that is a more invasive change and one that is a bit difficult to
wrangle across the spool trait definition in its current form.
If `limit_concurrent_stores` is not in use, this commit doesn't
change the behavior of the spool store operation.
Allow enabling a semaphore to specify concurrency constraints
for each of the spool operations.
This can help to avoid spawning hundreds of tokio blocking threads
that will just sit waiting for rocksdb to complete; the main saving
here is memory, as each blocking thread has a larger amount of
stack sitting idle than the per-task memory overhead that waits
for the semaphore.
Export the memory statistics for each spool database to prometheus for
charting and tracking.
Allow the hosting application to request a cache purge and set up a
monitor task to do that when memory usage is too high. That step will
print the memory that it reclaimed when it kicks in.
The original design of the spooling layer didn't require that
the SpoolId have the creation time encoded within it, which
meant that spool enumeration could consider any and all messages
found in the spool and import them into the queue subsystem.
If we allowed reception of new messages and wrote them to the spool
concurrently with the enumeration process, it would be possible
for the enumerator to observe the newly received messages and
import them into the queue subsystem, even though we had already
placed those newly received messages into the queue subsystem.
The result would be that we might send an additional copy of
each of the messages observed in this way.
To defend against that, the system refused to accept new messages
until spool enumeration was complete.
However, for sites with large spools, the enumeration process could
take some time to complete which could present challenges for
deploying updated configurations without an impact to their
service uptime.
This commit tackles that issue:
* Enumeration now filters out any messages that we created at or
after the start of the enumeration process, so it is not possible
for the duplicate scenario to occur.
* We no longer keep global track of whether spool enumeration is
in progress, but will still log that progress to the diagnostic
log.
* The liveness checks no longer check whether spool enumeration
is in progress.
A potential consequence of this change is that the concurrent writes
to the spool may further reduce the speed at which enumeration
operates, but that's a reasonable trade.
There are certain workloads and traffic patterns that can result
in shutdown taking a long time to complete. It's not generally
clear to the user what is happening there, so it is desirable
to improve that somehow.
During some recent testing I observed that the rust logic had
completed and that the kumod was process was blocked waiting
for an atexit handler that was joining a rocksdb thread.
This commit introduces an explicit shutdown concept to the spool
abstraction and spool manager.
After we have shutdown all in-flight messages and logs, we now
ask the spool manager to shutdown. It will steal away the
global refs to the meta and data spools and, concurrently, ask
them to shutdown, and then drop them.
For rocksdb, the shutdown request consists of asking it to
cancel any background work.
For the plain files spool, shutdown is a NOP.
We print out how long it took to perform the shutdown per spool,
as well as indicate when we start to shutdown the spool, as well
as when we are about to return from main. This should help to
understand when a similar atexit shutdown pause is coming into
play in the future.
We've been troubleshooting an issue where at very high
concurrency we see this panic:
kumomta-16 kumomta thread 'tokio-runtime-worker' panicked at /build-cache/cargo-home/registry/src/index.crates.io-6f17d22bba15001f/sharded-slab-0.1.7/src/tid.rs:163:21:
kumomta-16 kumomta creating a new thread ID (9075) would exceed the maximum number of thread ID bits specified in sharded_slab::cfg::DefaultConfig (8191)
The gist of the issue is that we somehow end up with over 9000 threads
(insert meme here) which is too large for the sharded-slab crate to
use for a thread id. That crate is used by the tracing-subscriber
crate which is vital for our logging/tracing functionality.
What I think is happening is that we have a LOT of concurrent
tasks calling spawn_blocking.
Ordinarily, tokio limits the number of blocking threads that it can
spawn to 512. In kumod we set up a handful of different runtimes
in order to manage the concurrency level for different workloads.
Somewhat ironically this acts as a concurrency multiplier when it comes
to blocking workloads, with each runtime allowing up to 512 threads.
This commit does a couple of things:
* Our auxilliary runtimes now make a point of setting the thread
names for spawned blocking threads appropriately. Previously,
they would use the default which is ambiguous wrt. the main
tokio runtime.
* Add env vars that can be used to override the maximum number of
blocking threads from its default of 512.
* Updates (almost!) all uses of spawn_blocking to use spawn_blocking_on
which takes an explicit runtime handle. We contrive for these to
receive the runtime handle that we set up in main. The idea is
that we don't want our auxilliary runtimes to spawn any additional
blocking threads at all, and we want that work to all run in the
main runtime so that it is easier to reason about the upper bound
of blocking work. There is one use in mod-sqlite and another
in mod-filesystem that I haven't reconciled yet, but those are
not an immediate concern.
This isn't shipped in the packages at this time.
It's a utility for post-mortem or offline analysis
of kumomta spools.
This version assumes rocksdb, but could be tweaked to work
with a plain filesystem spool.
We could add functions in the future to migrate from one format
to another, or perform other sorts of maintenance/admin operations.
This adjusts how we do write and delete operations so that we
first try to immediately submit the operation to the WAL,
but if that would block, then we push the operation to
the tokio blocking thread pool in order not to hard-block
the tokio schedule thread when rocksdb needs to perform writes.
While in there, hook up the force-sync flag to the closest
approximation in the WriteOptions for the batch (of size 1).
Adds `/api-docs/openapi.json` and `/rapidoc` endpoints to both
kumod and tsa-daemon.
The former exposes the subset of the API that is expressable
in the openapi schema as a json file that can be imported into
other tools.
The latter is a single-page web app that consumes the former
to provide an interactive API explorer.
We're using rapidoc for this, because I happen to think it looks
nicest and easiest to use, and we can integrate it into the docs
fairly nicely.
Which leads in nicely to say: I've integrated a read-only version
of rapidoc into the docs, and it even detects and adjusts to the
selected light/dark mode.
The `docs/update-openapi.sh` extracts the openapi.json data from
kumod and tsa-daemon and outputs to the correct place in the docs
directory structure to enable this.
refs: https://github.com/KumoCorp/kumomta/issues/96
The default is 1000 which is really high. This sets it to 10,
and configures rotation of those daily, so that there should be
no more than 10 days worth of LOG files.
Only enable for kumod. This makes it quicker and easier to
iterate on tests for a given crate that might otherwise
indirectly depend on the spool crate for its types.
it pulls in an old version of the time crate that flags many
rust projects as being vulnerable to an old CVE.
We still have other references to time that I'll remove in a follow up.
This commit introduces awareness of memory limits that may
have been established for the process.
The idea is that the hard/soft limits are read from the active
cgroup, falling back to classic ulimit hard/soft limits, falling
back to the system RAM size.
In the absence of an explicitly configured hard limit, the RAM size
is used for the hard limit.
In the absence of an explicitly configured soft limit, 75% of the
hard limit is used for the soft limit.
A background thread monitors the memory usage and (potentially
adjusted) memory limits.
A low memory state is when the usage is within 10% of the soft limit.
In this state, we start to trim back memory usage, shrinking loaded
data when messages are placed into the ready queue.
A more severe memory state is when the usage exceeds the soft limit.
In this state, new message reception is rejected until the usage
recovers.
kumod doesn't do anything directly with the hard limit, but assumes
that the system will terminate the process with no opportunity to
safely clean up if that value is reached.
In order to reduce the memory usage, the allocator has been switched
to jemalloc which has superior management of fragmentation and
a control interface to request release of various memory caches.
We make use of jemalloc specific functions in the tokio thread pools;
when they park (idle), we release the thread local cache, and when
we reach the soft limit we aggressively flush the caches until
memory usage falls below the limit.
There's like some room for tuning of the thresholds around this,
but this is a reasonable start.
Note that there are no explicit config options in kumod to set
the soft and hard limits: those are taken from the process
environment and can be simply set either via ulimit in a calling
shell script, or via systemd service configuration.
We may introduce some knobs for the low memory state threshold
in the future.
I broke spool enumeration for local disk when I adjusted
the path hashing a little while back, and didn't notice
because I was using rocksdb more than regular disk.
This is implemented by allowing the Message:set_due method to
load the message metadata to consult the scheduling constraints,
which should mean that all scheduling updates have the constraints
applied to them.
Looks pretty good compared to Sled.
|kind | flush | throughput |
+-------------+-------+------------+
|RocksDB | false | 102mm/hr |
|RocksDB | true | 96mm/hr | *
|Sled | false | 96mm/hr |
|Sled | true | 34mm/hr |
|LocalDisk | false | 24mm/hr |
|LocalDisk | true | 1mm/hr |
These numbers are from a 5950x (32 core) with an nvme drive,
as reported by:
```
cargo run --release -p traffic-gen -- --target 127.0.0.1:2025 --duration 20 --concurrency 16024
```
Note that the flush implementation with rocksdb just adjusts the setting
of use_fsync when opening the database.
There is an explicit db-wide flush that can be called, but it is very
aggressive and thorougly tanks performance down to 0.25mm/hr.
Note as well that rocksdb has a number of configuration options that may
work better as a write-once spool than the currently selected defaults;
more analysis could be done, but at the time of writing this commit
message, the defaults are the best performing storage option and going
further isn't a priority.
Add a `kind` and `flush` fields when defining a spool. Add a new
[sled](https://docs.rs/sled/latest/sled/index.html) based spool
implementation.
Initial benchmarking, especially at high concurrency, shows
promising numbers:
|kind | flush | throughput |
+-------------+-------+------------+
|Sled | false | 96mm/hr |
|Sled | true | 34mm/hr |
|LocalDisk | false | 24mm/hr |
|LocalDisk | true | 1mm/hr |
These numbers are from a 5950x (32 core) with an nvme drive,
as reported by:
```
cargo run --release -p traffic-gen -- --target 127.0.0.1:2025 --duration 20 --concurrency 16024
```
What's the catch? sled is considered beta by its authors.
https://github.com/spacejam/sled#known-issues-warnings
* Switch uuids to v1 format, so that we can cheaply determine
when a message was created without having to load its metadata
from the spool
* Add some message delivery parameters; retry interval, limit, max age
* Respect those parameters when spooling in and when we encounter
a transient failure.