I noticed this recently during some testing; for large queue sizes
and large connection limits, we could end up opening more connections
than we currently have queued messages to deliver.
The issue was that the `.min()` constraint was placed on the wrong
term of the calculation, clamping prior to scaling, instead of
after scaling.
When `log_arf` or `log_oob` are set to true with `relay_to=false`, we
now return a 550 error response for messages that are not ARF or OOB
reports. Previously, we would return a 250 response and silently drop
the message in this case, which gave the false impression that it was
accepted for relaying.
Expand integration test to explicitly assert that the right things
are allowed/denied/relayed/parsed.
When this was added, it was to enable TSA daemon to express
and convey that a given path should be suspended.
The implementation was objectively unpleasant, turning the
fast path of egress source assignment from a simple single
iteration of the weighted round robin logic to N (where N
is the number of sources in a pool) iterations to try
and figure out if a suspension is active.
Now that we have TSA subscriptions that update the admin
suspension information in realtime, we can now simply
rely on that source of information.
This commit removes the logic associated with the deprecated
suspended flag and tidies up the code, making that fast
path a little more fast, and making the code a bit easier
to reason about.
This sequence of events:
1. Messages are flowing
2. The remote site starts to return a 421 for mail at connection
time, prior to MAIL FROM, on new connections
3. The local site employs a TSA rule that suspends the corresponding
ready queue
Could result in the contents of that ready_queue getting stuck.
Here's a representation of the logs:
```
Apr 24 19:03:44 maintain SITENAME: computed ideal connection count as 3
Apr 24 19:03:44 Error in Dispatcher::run for SITENAME: connect to ResolvedAddress { name: "smtp-in.orange.fr.", addr: 80.12.26.32 } port 25 and read initial banner: Command rejected Response { code: 421, enhanced_code: None, content: "XXXXX smtp.orange.fr XXXXX Service refuse. Veuillez essayer plus tard. Service refused, please try later. OFR_999 [999]", command: None } (consecutive_connection_failures=1)
Apr 24 19:03:45 maintain SITENAME: there are now 3 connections, suspended(via config)=false, suspended(admin)=false, queue_size=691
Apr 24 19:03:45 maintain SITENAME: computed ideal connection count as 3
Apr 24 19:03:45 maintain SITENAME: there are now 3 connections, suspended(via config)=false, suspended(admin)=true, queue_size=746
... the above repeats every minute for a while ...
Apr 24 19:03:45 maintain SITENAME: computed ideal connection count as 0
Apr 24 19:13:45 maintain SITENAME: there are now 3 connections, suspended(via config)=false, suspended(admin)=true, queue_size=746
Apr 24 19:13:45 maintain SITENAME: computed ideal connection count as 0
Apr 24 19:13:45 reaping site SITENAME
Apr 24 19:23:44 maintain SITENAME: there are now 0 connections, suspended(via config)=false, suspended(admin)=false, queue_size=1
```
There are two problems:
1. The main mechanism that acts to sweep the read queue into the
scheduled queue is only triggered once we get past the initial
connection phase, causing messages to be retained in the
ready queue instead of moving back into the scheduled queue.
2. The is-reapable logic doesn't consider the number of messages
that may be in the ready queue. It was written before we had
a suspension concept and only considers that an ideal connection
count of 0 implies that there are no messages.
The result is, in combination with 1, is that we can reap the ready
queue while it contains messages and forget about them until
the server is restarted.
This commit resolves both of these issues:
1. The maintain routine will now act to sweep messages into the
corresponding scheduled queue if an admin suspension is active.
In addition, this same sweep is considered when we reap the
ready queue, just in case some other logic bug manifests in
the future with the same consequences.
2. The reap logic now also requires that the ready queue be empty.
This script is triggered as part of `make test` which is not
really a great place for it.
Its purpose is to extract the auto-generated openapi spec
from the kumod and tsa binaries and update the snapshot
that is present in the docs.
It needs kumod and tsa-daemon to have been built in debug mode
to run successfully.
It piggy-backs on `make test` on the assumption that it will
cause the person who is making changes to it to include those
spec updates in their commit/PR.
Since `make test` invokes it, the various builders may try
and fail to execute jq in the `test` step. This is mostly
harmless, but looks noisy in the logs.
The ideal situation for this would be:
* at PR time and push time: add a check that runs this script and that
fails if the specs are updated for the docs and are not part of the PR
itself. (eg: status is dirty after running it).
This way it will be visible from the CI state that something is awry.
If many calls are made to the memoized function with the same parameters
at the time that the cache is empty/expired, then each of those
concurrent calls will proceed to compute and populate the cache. This
is known as a thundering herd, and can be painful if the amount of work
performed by the cache population function is high, or alternatively, if
the level of concurrency is high and some system resource is required to
satisfy the call, this can put the system under higher pressure.
This commit adds a simple mechanism to mitigate this: each combination
of cache and cache parameters is paired with an optional semaphore that
is used to constrain concurrency to a single call. This could be a
mutex, but using a semaphore allows future explicit control over the
concurrency level.
Adds a couple of options that provide more control over the default
timeouts and logging.
These are exposed to the shaping helper as `publish_timeout`,
`publish_pool_idle_timeout` and `publish_connection_verbose`.
I think we've gone back and forth on this a couple of times.
The motivation for this commit is that we've been troubleshooting
an issue that seems to be triggered by bouncing all the mail.
The symptom is that the system becomes unresponsive after
triggering the bounce.
Here in the context of `bounce_all`, the queue lock is held
so that we can capture the messages, but then we serially
remove them from the spool, so that we can report the total
number back to the originating HTTP request.
For a large queue size that presents a big point of contention
for other tasks or requests that may need to operate on the
queue.
This commit moves that spool removal into another task so that
that task can run asynchronously and independently from the
bounce HTTP request.
The consequence of this is that the numbers reported by the
bounce request will likely be lower than the final count.
The `bounce-list` kcli command can be used to check up on
those numbers.
We've been troubleshooting a lockup on a system with a low core count.
What we found in a stack trace was that one of the threads was blocking
on the CACHE mutex in the sts logic. That mutex is intended to be
short-duration in scope, managing the direct lookup or insertion
into the cache, and no more.
However, due to the the way that rust scopes the lifetime of the
MutexGuard that is acquired from the CACHE, we were holding it
across the async DNS operation that is used to validate that
a cached policy is still current.
This can result in a deadlock on a system with a sufficiently
low core count/high enough concurrent volume of traffic to sites
with MTA-STS enabled.
This commit resolves this by introducing a lookup function that
explicitly clones the cached policy without returning a MutexGuard.
We skip logging the 421 we generate while shutting down because
it feels a bit redundant; you'll see the server shutting down
in the journal anyway.
refs: https://github.com/KumoCorp/kumomta/issues/88
This commit connects the new websocket based suspension feed
up to shaping.lua. This allows ready-q suspensions to be
enacted in realtime, as well as sets things up to support
scheduled queue suspensions in a later commit.
refs: https://github.com/KumoCorp/kumomta/issues/113
This is very similar to the HTTP suspension API, with the
difference that the suspend method returns just the uuid rather
than the entire suspension object.
refs: https://github.com/KumoCorp/kumomta/issues/113
This function will spawn a new thread that runs a tokio
localset, which in turn will trigger the specified event
and run it.
On it's own it doesn't do a lot, but it provides a way
to perform background tasks in lua.
```lua
kumo.on('init', function()
kumo.spawn_task {
event_name = 'my.task',
args = { 'hello', 'there' },
}
end)
kumo.on('my.task', function(args)
-- Prints: `I am the task. ["hello","there"]`
print('I am the task.', kumo.json_encode(args))
end)
```
The intent of the idle behavior is to linger for up to the configured
idle_timeout value, waiting for new messages to arrive in the ready
queue.
The actual behavior was to initiate a wait, but when woken up,
if there were no messages in the ready queue, the connection
would close out, even if there was still time that could be
waited out before the idle period was up.
This commit adds in a loop that will keep the connection open
until the idle timeout is reached, which in turn will improve
throughput, especially if the traffic is a little bursty.
Previously we would only remove it from the first line. This commit
removes it from the second and subseqent lines as well.
The bulk of this commit was a little bit of refactoring to favilitate
testing this change.
refs: #157
This commit refactors listener_domains.lua to facilitate unit testing
and adds a couple of basic test cases.
The functional change here is that we were missing an explicit
fallback step in the case where no listener or domain matches
the provided values directly; we need to explicitly add a check
against the `*` listener AND `*` domain for the final step.
Previously, we would only look at the `*` listener for the final
step.
closes: #128
These allow performing arbitrary rate limiting operations
at both reception time and when messages are moved from the
scheduled queue and into the ready queue.
refs: https://github.com/KumoCorp/kumomta/issues/149