diff --git a/assets/policy-extras/policy_utils.lua b/assets/policy-extras/policy_utils.lua index 72e2661e..a0c12b2f 100644 --- a/assets/policy-extras/policy_utils.lua +++ b/assets/policy-extras/policy_utils.lua @@ -1,6 +1,24 @@ local mod = {} local kumo = require 'kumo' +function mod.assert_eq(a, b, message) + if a ~= b then + if message then + message = string.format('%s. ', message) + else + message = '' + end + error( + string.format( + '%sexpected %s == %s', + message, + kumo.json_encode(a), + kumo.json_encode(b) + ) + ) + end +end + -- Helper function that merges the values from `src` into `dest` function mod.merge_into(src, dest) if src then @@ -10,6 +28,23 @@ function mod.merge_into(src, dest) end end +-- Helper function that merges the values from `src` into `dest` +function mod.recursive_merge_into(src, dest) + if src then + for k, v in pairs(src) do + if type(v) == 'table' then + if not dest[k] then + dest[k] = v + else + mod.recursive_merge_into(v, dest[k]) + end + else + dest[k] = v + end + end + end +end + -- Returns true if the table tbl contains value v function mod.table_contains(tbl, v) if v == nil then @@ -33,6 +68,14 @@ function mod.table_keys(tbl) return result end +-- Returns true if table is empty, false otherwise +function mod.table_is_empty(tbl) + for _k, _v in pairs(tbl) do + return false + end + return true +end + function mod.load_json_or_toml_file(filename) if filename:match '.toml$' then return kumo.toml_load(filename) diff --git a/assets/policy-extras/queue.lua b/assets/policy-extras/queue.lua new file mode 100644 index 00000000..b858b119 --- /dev/null +++ b/assets/policy-extras/queue.lua @@ -0,0 +1,326 @@ +local mod = {} +local kumo = require 'kumo' +local utils = require 'policy-extras.policy_utils' + +--[[ +Usage: + +Create a `/opt/kumomta/etc/queue.toml` with contents like: + +```toml +# Allow optional scheduled sends based on this header +# https://docs.kumomta.com/reference/message/import_scheduling_header +scheduling_header = "X-Schedule" +# Set the tenant from this header +tenant_header = "X-Tenant" +# Set the campaign from this header +campaign_header = "X-Campaign" +# The tenant to use if no tenant_header is present +default_tenant = "default-tenant" + +[tenant.'default-tenant'] +egress_pool = 'pool-0' + +[tenant.'mytenant'] +# Which pool should be used for this tenant +egress_pool = 'pool-1' +# Override age based on tenant; this overrides +# settings at the domain level +max_age = '10 hours' +# Only the authorized identities are allowed to use this tenant +# via the tenant_header +require_authz = ["scott"] + +# The default set of parameters +[queue.default] +max_age = '24 hours' + +# Base settings for a given destination domain. +# These are overridden by more specific settings +# in a tenant or more specific queue +[queue.'gmail.com'] +max_age = '22 hours' +retry_interval = '17 mins' + +[queue.'gmail.com'.'mytenant'] +# options here for domain=gmail.com AND tenant=mytenant for any unmatched campaign + +[queue.'gmail.com'.'mytenant'.'welcome-campaign'] +# options here for domain=gmail.com, tenant=mytenant, and campaign='welcome-campaign' +``` + +--- +local queue_module = require 'policy-extras.queue' +local queue_helper = queue_module:setup({'/opt/kumomta/etc/queue.toml'}) + +kumo.on('smtp_server_message_received', function(msg) + queue_helper:apply(msg) + + -- do your dkim signing here +end) + +--- +]] + +local function parse_one(data) + local options = {} + local tenants = {} + local queues = {} + + for k, v in pairs(data) do + if k == 'tenant' then + for tenant_name, tenant_options in pairs(v) do + local tenant = tenants[tenant_name] or {} + utils.merge_into(tenant_options, tenant) + tenants[tenant_name] = tenant + end + elseif k == 'queue' then + for domain_name, queue_options in pairs(v) do + local domain_options = queues[domain_name] or {} + utils.recursive_merge_into(queue_options, domain_options) + queues[domain_name] = domain_options + end + else + options[k] = v + end + end + + local result = { + options = options, + tenants = tenants, + queues = queues, + } + + return result +end + +local function merge_data(loaded_files) + local result = {} + for _, data in ipairs(loaded_files) do + utils.recursive_merge_into(parse_one(data), result) + end + -- print(kumo.json_encode_pretty(result)) + result.queues = kumo.domain_map.new(result.queues) + return result +end + +local function load_queue_config(file_names) + local data = {} + for _, file_name in ipairs(file_names) do + table.insert(data, utils.load_json_or_toml_file(file_name)) + end + + return merge_data(data) +end + +local function resolve_config(data, domain, tenant, campaign) + -- print('resolve_config', domain, tenant, campaign) + + local params = {} + + local default_config = data.queues.default + if default_config then + utils.merge_into(default_config, params) + end + + local domain_config = data.queues[domain] + if domain_config then + for k, v in pairs(domain_config) do + if type(v) ~= 'table' then + params[k] = v + end + end + end + + local tenant_definition = data.tenants[tenant] + if tenant_definition then + for k, v in pairs(tenant_definition) do + if type(v) ~= 'table' then + params[k] = v + end + end + end + + if domain_config then + local tenant_config = domain_config[tenant] + if tenant_config then + for k, v in pairs(tenant_config) do + if type(v) ~= 'table' then + params[k] = v + end + end + + local campaign = tenant_config[campaign] + + if campaign then + for k, v in pairs(campaign) do + if type(v) ~= 'table' then + params[k] = v + end + end + end + end + end + + if utils.table_is_empty(params) then + return nil + end + + -- print(kumo.json_encode_pretty(params)) + return params +end + +function mod:setup(file_names) + local cached_load_data = kumo.memoize(load_queue_config, { + name = 'queue_helper_data', + ttl = '1 minute', + capacity = 10, + }) + + local test = cached_load_data(file_names) + + kumo.on( + 'get_queue_config', + function(domain, tenant, campaign, _routing_domain) + local data = cached_load_data(file_names) + local params = resolve_config(data, domain, tenant, campaign) + if params then + return kumo.make_queue_config(params) + end + end + ) + + local helper = { + file_names = file_names, + } + + function helper:apply(msg) + local data = cached_load_data(file_names) + if data.options.scheduling_header then + msg:import_scheduling_header(data.options.scheduling_header, true) + end + if data.options.campaign_header then + local campaign = + msg:get_first_named_header_value(data.options.campaign_header) + if campaign then + msg:set_meta('campaign', campaign) + end + end + local tenant = nil + local tenant_source = nil + if data.options.tenant_header then + tenant = msg:get_first_named_header_value(data.options.tenant_header) + if tenant then + tenant_source = + string.format("'%s' header", data.options.tenant_header) + end + end + if not tenant and data.options.default_tenant then + tenant = data.options.default_tenant + tenant_source = 'default_tenant option' + end + if tenant then + local tenant_config = data.tenants[tenant] + + if not tenant_config then + kumo.reject( + 500, + string.format( + "tenant '%s' specified by %s is not defined in your queue config", + tenant, + tenant_source + ) + ) + end + + if tenant_config.require_authz then + local my_auth = msg:get_meta 'authz_id' + + if not my_auth then + kumo.reject( + 500, + string.format("tenant '%s' requires SMTP AUTH", tenant) + ) + end + + if not utils.table_contains(tenant_config.require_authz, my_auth) then + kumo.reject( + 500, + string.format( + "authz_id '%s' is not authorized to send as tenant '%s'", + my_auth, + tenant + ) + ) + end + end + + msg:set_meta('tenant', tenant) + end + end + + return helper +end + +--[[ +Run some basic unit tests for the data parsing/merging; use it like this: + +``` +local queue_module = require 'policy-extras.queue' +queue_module:test() +``` +]] +function mod:test() + local base_data = [=[ +default_tenant = 'mytenant' + +[queue.default] +max_age = '24 hours' + +[tenant.'mytenant'] +egress_pool = 'tpool' + +[queue.'gmail.com'.'mytenant'] +egress_pool = "foo" +]=] + + local user_data = [=[ +[queue.'gmail.com'] +max_age = '6 hours' + +[queue.'gmail.com'.'mytenant'.'campaign'] +egress_pool = "bar" +]=] + + local data = merge_data { + kumo.toml_parse(base_data), + kumo.toml_parse(user_data), + } + + utils.assert_eq( + resolve_config(data, 'foo.com', nil, nil, nil).max_age, + '24 hours' + ) + utils.assert_eq( + resolve_config(data, 'gmail.com', nil, nil, nil).max_age, + '6 hours' + ) + utils.assert_eq( + resolve_config(data, 'gmail.com', nil, nil, nil).egress_pool, + nil + ) + utils.assert_eq( + resolve_config(data, 'example.com', 'mytenant', nil, nil).egress_pool, + 'tpool' + ) + utils.assert_eq( + resolve_config(data, 'gmail.com', 'mytenant', nil, nil).egress_pool, + 'foo' + ) + utils.assert_eq( + resolve_config(data, 'gmail.com', 'mytenant', 'campaign', nil).egress_pool, + 'bar' + ) +end + +return mod