Files
kumomta/crates/integration-tests/source.lua
T
Wez Furlong 5e1ae20497 Add batching support for log hooks
This really is adding batching support to custom lua delivery
protocol handlers, but the main use case for these today is
to implement log hooks.

The way that it works is that you can specify a `batch_size`
as part of setting up the lua protocol handler.

Then, when it is time to send messages, if the batch_size is
the default of 1, the lua delivery logic will invoke the `send` method
on the connection object returned from the constructor.  This
is the same as the behavior from before this commit.

However, if the batch_size is greater than 1, then the lua delivery
logic will instead attempt to collect up to batch_size messages
that are immediately available from the ready queue, and then pass
those to a new `send_batch` method.

The send_batch method accepts an array of messages; that array will
always have at least one message, and up to batch_size messages,
depending on the throughput and queue size.

If the send_batch method's return value applies equally to all
messages in the batch, so if it indicates that something failed,
that disposition will apply to all messages.

One of the reasons that I'd avoided implementing batching thus far
was that it makes it awkward to resolve persistent/recurring issues
that are due to a single message in that batch.  If the batch is
always retried together then there is a good chance that it will
always fail together.

There's no explicit mitigation for that issue here, but it may
be probablistically mitigated by the jitter that is applied to
messages that transiently fail.  If a batch transiently fails,
each message in that batch will be subject to its own random
jitter which should cause an offending message to be retried
with a different subset of messages next time around.

The integration test included here demonstrates the batching
working with an http log hook implementation.
2024-09-16 18:14:25 -07:00

259 lines
6.6 KiB
Lua

local kumo = require 'kumo'
package.path = '../../assets/?.lua;' .. package.path
log_hooks = require 'policy-extras.log_hooks'
local TEST_DIR = os.getenv 'KUMOD_TEST_DIR'
local SINK_PORT = tonumber(os.getenv 'KUMOD_SMTP_SINK_PORT')
local WEBHOOK_PORT = os.getenv 'KUMOD_WEBHOOK_PORT'
local AMQPHOOK_URL = os.getenv 'KUMOD_AMQPHOOK_URL'
local AMQP_HOST_PORT = os.getenv 'KUMOD_AMQP_HOST_PORT'
local LISTENER_MAP = os.getenv 'KUMOD_LISTENER_DOMAIN_MAP'
kumo.on('init', function()
kumo.configure_accounting_db_path(TEST_DIR .. '/accounting.db')
local relay_hosts = { '0.0.0.0/0' }
local RELAY_HOSTS = os.getenv 'KUMOD_RELAY_HOSTS'
if RELAY_HOSTS then
relay_hosts = kumo.json_parse(RELAY_HOSTS)
end
kumo.start_esmtp_listener {
listen = '127.0.0.1:0',
relay_hosts = relay_hosts,
}
kumo.start_http_listener {
listen = '127.0.0.1:0',
}
kumo.configure_local_logs {
log_dir = TEST_DIR .. '/logs',
max_segment_duration = '1s',
headers = { 'X-*', 'Y-*', 'Subject' },
}
if WEBHOOK_PORT then
kumo.configure_log_hook {
name = 'webhook',
headers = { 'Subject', 'X-*' },
}
end
kumo.define_spool {
name = 'data',
path = TEST_DIR .. '/data-spool',
}
kumo.define_spool {
name = 'meta',
path = TEST_DIR .. '/meta-spool',
}
end)
if AMQPHOOK_URL then
log_hooks:new {
name = 'amqp',
constructor = function(domain, tenant, campaign)
local sender = {}
local client = kumo.amqp.build_client(AMQPHOOK_URL)
function sender:send(msg)
local result = client:publish_with_timeout({
routing_key = 'woot',
payload = msg:get_data(),
}, 20000)
if result.status == 'Ack' or result.status == 'NotRequested' then
return string.format('250 %s', kumo.json_encode(result))
end
-- result.status must be `Nack`; log the full result
kumo.reject(500, kumo.json_encode(result))
end
function sender:close()
client:close()
end
return sender
end,
}
elseif AMQP_HOST_PORT then
log_hooks:new {
name = 'amqp',
constructor = function(domain, tenant, campaign)
local sender = {}
local host, port = table.unpack(kumo.string.split(AMQP_HOST_PORT, ':'))
function sender:send(msg)
kumo.amqp.basic_publish {
routing_key = 'woot',
payload = msg:get_data(),
connection = {
host = host,
port = tonumber(port),
},
}
return '250 ok'
end
return sender
end,
}
end
if WEBHOOK_PORT then
local batch_size = tonumber(os.getenv 'KUMOD_WEBHOOK_BATCH_SIZE')
if batch_size > 1 then
log_hooks:new {
name = 'webhookbatch',
batch_size = batch_size,
constructor = function(domain, tenant, campaign)
local sender = {}
local client = kumo.http.build_client {}
function sender:send_batch(messages)
local payload = {}
for _, msg in ipairs(messages) do
table.insert(payload, msg:get_meta 'log_record')
end
print(string.format('batch size is %d ***************', #payload))
local response = client
:post(
string.format('http://127.0.0.1:%d/log-batch', WEBHOOK_PORT)
)
:header('Content-Type', 'application/json')
:body(kumo.serde.json_encode(payload))
:send()
local disposition = string.format(
'%d %s: %s',
response:status_code(),
response:status_reason(),
response:text()
)
if response:status_is_success() then
return disposition
end
kumo.reject(500, disposition)
end
return sender
end,
}
else
kumo.on('should_enqueue_log_record', function(msg)
local log_record = msg:get_meta 'log_record'
-- avoid an infinite loop caused by logging that we logged
if log_record.queue ~= 'webhook' then
msg:set_meta('queue', 'webhook')
return true
end
return false
end)
kumo.on('make.webhook', function(_domain, _tenant, _campaign)
local sender = {}
local client = kumo.http.build_client {}
function sender:send(message)
local response = client
:post(string.format('http://127.0.0.1:%d/log', WEBHOOK_PORT))
:header('Content-Type', 'application/json')
:body(message:get_data())
:send()
local disposition = string.format(
'%d %s: %s',
response:status_code(),
response:status_reason(),
response:text()
)
if response:status_is_success() then
return disposition
end
kumo.reject(500, disposition)
end
return sender
end)
end
end
kumo.on('get_listener_domain', function(domain, listener, conn_meta)
if LISTENER_MAP then
local map = kumo.json_parse(LISTENER_MAP)
local params = map[domain]
if params then
return kumo.make_listener_domain(params)
end
end
return kumo.make_listener_domain {
relay_to = true,
log_oob = true,
log_arf = true,
}
end)
kumo.on('smtp_server_message_received', function(msg) end)
kumo.on('get_queue_config', function(domain, _tenant, _campaign)
if domain == 'webhook' then
return kumo.make_queue_config {
protocol = {
custom_lua = {
constructor = 'make.webhook',
},
},
}
end
return kumo.make_queue_config {
protocol = {
-- Redirect traffic to the sink
smtp = {
mx_list = { 'localhost' },
},
},
retry_interval = os.getenv 'KUMOD_RETRY_INTERVAL',
strategy = os.getenv 'KUMOD_QUEUE_STRATEGY',
}
end)
kumo.on('get_egress_path_config', function(domain, _source_name, _site_name)
-- Allow sending to a sink
local params = {
enable_tls = os.getenv 'KUMOD_ENABLE_TLS' or 'OpportunisticInsecure',
smtp_port = SINK_PORT,
prohibited_hosts = {},
}
local username = os.getenv 'KUMOD_SMTP_AUTH_USERNAME'
local password = os.getenv 'KUMOD_SMTP_AUTH_PASSWORD'
if username and password then
params.smtp_auth_plain_username = username
params.smtp_auth_plain_password = {
key_data = password,
}
end
if domain == 'webhookbatch.log_hook' then
-- we use a slow throttle here to encourage batches
-- to be used
params.max_message_rate = '3/s'
end
print('get_egress_path_config *******************', domain)
return kumo.make_egress_path(params)
end)
if os.getenv 'KUMOD_WANT_REBIND' then
kumo.on('rebind_message', function(message, data)
message:set_meta('queue', data.queue)
end)
end