id: bitcoin-fee-aware-payout-batching
namespace: company.team
description: |
Batch queued Bitcoin payouts into a single transaction, but only when mempool
fees are below your cap. Invalid addresses are rejected, a human approves the
batch, Bitcoin Core broadcasts it, and every payout row records its txid.
concurrency:
limit: 1
behavior: QUEUE
triggers:
- id: queue_threshold
type: io.kestra.plugin.jdbc.postgresql.Trigger
description: Start a batch as soon as enough payouts are waiting. The query only
returns a row when the pending queue reaches the threshold. Shipped
disabled; enable it after your first manual run.
disabled: true
interval: PT5M
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
fetchType: FETCH_ONE
sql: |
SELECT count(*) AS pending FROM payouts WHERE status = 'pending' HAVING count(*) >= 10
- id: every_30_minutes
type: io.kestra.plugin.core.trigger.Schedule
description: Sweep the queue regularly so small batches do not wait forever for
the threshold. Shipped disabled.
disabled: true
cron: "*/30 * * * *"
inputs:
- id: rpc_url
type: STRING
defaults: http://127.0.0.1:38332
description: Bitcoin Core JSON-RPC endpoint. The default is the signet port; use
18443 for regtest and 8332 for mainnet.
- id: wallet
type: STRING
defaults: payouts
description: Name of the Bitcoin Core wallet that funds the payouts.
- id: allow_mainnet
type: BOOL
defaults: false
description: The flow asks the node which chain it is on and refuses to send on
mainnet unless this is true.
- id: fee_api_url
type: URI
defaults: https://mempool.space/signet/api/v1/fees/recommended
description: Fee estimate endpoint in the mempool.space format. Use
https://mempool.space/api/v1/fees/recommended for mainnet.
- id: fee_target
type: SELECT
values:
- fastestFee
- halfHourFee
- hourFee
- economyFee
defaults: hourFee
description: Which mempool.space estimate to pay. Payouts are rarely urgent, so
the default trades speed for cost.
- id: max_fee_rate
type: FLOAT
defaults: 5
description: Highest fee rate in sat/vB you are willing to pay. Above it, the
run defers and the queue waits for the next trigger.
- id: max_outputs
type: INT
defaults: 50
description: Maximum number of payouts in one transaction.
- id: max_batch_btc
type: FLOAT
defaults: 0.1
description: Safety cap on the total value of one batch, in BTC. A larger batch
fails before approval and is returned to the queue.
- id: require_approval
type: BOOL
defaults: true
description: Pause for a human to approve each batch before it is broadcast.
- id: kestra_url
type: STRING
defaults: http://localhost:8080
description: Base URL of your Kestra UI, used in the approval message so the
approver can open the paused execution.
- id: explorer_tx_url
type: STRING
defaults: https://mempool.space/signet/tx/
description: Block explorer prefix for transaction links in notifications.
tasks:
- id: ensure_table
type: io.kestra.plugin.jdbc.postgresql.Query
description: Create the payout queue on first run. Your application inserts rows
with status 'pending'.
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
sql: |
CREATE TABLE IF NOT EXISTS payouts (
id BIGSERIAL PRIMARY KEY,
address TEXT NOT NULL,
amount_btc NUMERIC(16, 8) NOT NULL CHECK (amount_btc > 0),
status TEXT NOT NULL DEFAULT 'pending',
batch_id TEXT,
txid TEXT,
note TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
)
- id: node_info
type: io.kestra.plugin.core.http.Request
description: Ask the node which chain it is on, so the mainnet guard does not
depend on how the inputs were filled in.
uri: "{{ inputs.rpc_url }}"
method: POST
body: '{"jsonrpc": "1.0", "id": "kestra", "method": "getblockchaininfo",
"params": []}'
options:
auth:
type: BASIC
username: "{{ secret('BITCOIN_RPC_USER') }}"
password: "{{ secret('BITCOIN_RPC_PASSWORD') }}"
- id: mainnet_guard
type: io.kestra.plugin.core.flow.If
condition: "{{ (outputs.node_info.body | jq('.result.chain') | first) == 'main'
and not inputs.allow_mainnet }}"
then:
- id: refuse_mainnet
type: io.kestra.plugin.core.execution.Fail
errorMessage: The node is on mainnet and allow_mainnet is false. Nothing was sent.
- id: load_queue
type: io.kestra.plugin.jdbc.postgresql.Query
description: Read the oldest pending payouts, up to max_outputs.
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
fetchType: FETCH
sql: |
SELECT id, address FROM payouts
WHERE status = 'pending'
ORDER BY created_at, id
LIMIT {{ inputs.max_outputs }}
- id: queue_empty
type: io.kestra.plugin.core.flow.If
condition: "{{ outputs.load_queue.size == 0 }}"
then:
- id: nothing_to_pay
type: io.kestra.plugin.core.execution.Exit
state: SUCCESS
- id: validate_addresses
type: io.kestra.plugin.core.flow.Loop
description: Check every address with the node and reject invalid ones, so one
bad address cannot block the whole queue.
values: "{{ outputs.load_queue.rows | toJson }}"
concurrencyLimit: 10
tasks:
- id: validate
type: io.kestra.plugin.core.http.Request
uri: "{{ inputs.rpc_url }}"
method: POST
body: '{"jsonrpc": "1.0", "id": "kestra", "method": "validateaddress", "params":
["{{ fromJson(item.value).address }}"]}'
options:
auth:
type: BASIC
username: "{{ secret('BITCOIN_RPC_USER') }}"
password: "{{ secret('BITCOIN_RPC_PASSWORD') }}"
- id: invalid_address
type: io.kestra.plugin.core.flow.If
condition: "{{ not (outputs.validate.body | jq('.result.isvalid') | first) }}"
then:
- id: reject_payout
type: io.kestra.plugin.jdbc.postgresql.Query
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
sql: |
UPDATE payouts SET status = 'rejected', note = 'invalid address for this network', updated_at = now()
WHERE id = {{ fromJson(item.value).id }} AND status = 'pending'
- id: fees
type: io.kestra.plugin.core.http.Request
description: Current fee estimates from mempool.space.
uri: "{{ inputs.fee_api_url }}"
- id: fee_gate
type: io.kestra.plugin.core.flow.If
description: Defer the batch when fees are above the cap. The queue is untouched
and the next trigger tries again.
condition: "{{ (outputs.fees.body | jq('.' ~ inputs.fee_target) | first) >
inputs.max_fee_rate }}"
then:
- id: log_deferred
type: io.kestra.plugin.core.log.Log
message: "Deferred: {{ inputs.fee_target }} is {{ outputs.fees.body | jq('.' ~
inputs.fee_target) | first }} sat/vB, above the {{ inputs.max_fee_rate
}} sat/vB cap."
- id: deferred
type: io.kestra.plugin.core.execution.Exit
state: SUCCESS
- id: lock_batch
type: io.kestra.plugin.jdbc.postgresql.Query
description: Claim the oldest valid pending payouts for this execution, so no
other run can pay them.
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
sql: |
UPDATE payouts p SET status = 'locked', batch_id = '{{ execution.id }}', updated_at = now()
FROM (
SELECT id FROM payouts WHERE status = 'pending'
ORDER BY created_at, id LIMIT {{ inputs.max_outputs }}
FOR UPDATE SKIP LOCKED
) batch
WHERE p.id = batch.id
- id: batch_summary
type: io.kestra.plugin.jdbc.postgresql.Query
description: Total the batch and build the sendmany amounts, merging payouts
that go to the same address.
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
fetchType: FETCH_ONE
sql: |
SELECT count(*) AS outputs,
coalesce(sum(amount), 0) AS total_btc,
coalesce(sum(payouts), 0) AS payouts,
coalesce(json_object_agg(address, amount), '{}')::text AS amounts
FROM (
SELECT address, sum(amount_btc) AS amount, count(*) AS payouts
FROM payouts WHERE batch_id = '{{ execution.id }}' AND status = 'locked'
GROUP BY address
) per_address
- id: batch_empty
type: io.kestra.plugin.core.flow.If
description: Every queued address was rejected; nothing left to send.
condition: "{{ outputs.batch_summary.row.outputs == 0 }}"
then:
- id: nothing_valid
type: io.kestra.plugin.core.execution.Exit
state: SUCCESS
- id: value_cap
type: io.kestra.plugin.core.flow.If
description: Refuse batches above the safety cap. The errors block returns the
payouts to the queue.
condition: "{{ outputs.batch_summary.row.total_btc > inputs.max_batch_btc }}"
then:
- id: over_cap
type: io.kestra.plugin.core.execution.Fail
errorMessage: "Batch total {{ outputs.batch_summary.row.total_btc }} BTC is
above the {{ inputs.max_batch_btc }} BTC cap. Payouts were returned to
the queue."
- id: approval
type: io.kestra.plugin.core.flow.If
condition: "{{ inputs.require_approval }}"
then:
- id: request_approval
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "Bitcoin payout batch waiting for approval: {{ outputs.batch_summary.row.payouts }} payouts to {{ outputs.batch_summary.row.outputs }} addresses, {{ outputs.batch_summary.row.total_btc }} BTC at {{ outputs.fees.body | jq('.' ~ inputs.fee_target) | first }} sat/vB. Approve or reject: {{ inputs.kestra_url }}/ui/main/executions/{{ flow.namespace }}/{{ flow.id }}/{{ execution.id }}"
}
- id: wait_for_approval
type: io.kestra.plugin.core.flow.Pause
pauseDuration: PT24H
behavior: FAIL
onResume:
- id: approved
type: BOOL
defaults: false
description: Approve this batch for broadcast
- id: approver_note
type: STRING
required: false
description: Optional note stored on the payouts
- id: rejected_by_approver
type: io.kestra.plugin.core.flow.If
condition: "{{ inputs.require_approval and not
(outputs.wait_for_approval.onResume.approved ?? false) }}"
then:
- id: release_rejected
type: io.kestra.plugin.jdbc.postgresql.Query
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
sql: |
UPDATE payouts SET status = 'pending', batch_id = NULL, updated_at = now()
WHERE batch_id = '{{ execution.id }}' AND status = 'locked'
- id: batch_rejected
type: io.kestra.plugin.core.execution.Exit
state: WARNING
- id: mark_broadcasting
type: io.kestra.plugin.jdbc.postgresql.Query
description: From here on the payouts are never returned to the queue
automatically, so a failure after broadcast cannot cause a double payment.
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
sql: |
UPDATE payouts SET status = 'broadcasting', updated_at = now()
WHERE batch_id = '{{ execution.id }}' AND status = 'locked'
- id: send_batch
type: io.kestra.plugin.core.http.Request
description: One sendmany call pays every address at the approved fee rate.
Transactions are replaceable (RBF), so a stuck batch can be fee-bumped.
uri: "{{ inputs.rpc_url }}/wallet/{{ inputs.wallet }}"
method: POST
body: |
{"jsonrpc": "1.0", "id": "kestra", "method": "sendmany", "params": {
"dummy": "",
"amounts": {{ outputs.batch_summary.row.amounts }},
"comment": "kestra payout batch {{ execution.id }}",
"replaceable": true,
"fee_rate": {{ outputs.fees.body | jq('.' ~ inputs.fee_target) | first }}
}}
options:
auth:
type: BASIC
username: "{{ secret('BITCOIN_RPC_USER') }}"
password: "{{ secret('BITCOIN_RPC_PASSWORD') }}"
- id: mark_sent
type: io.kestra.plugin.jdbc.postgresql.Query
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
sql: |
UPDATE payouts
SET status = 'sent', txid = '{{ outputs.send_batch.body | jq('.result') | first }}',
note = '{{ (outputs.wait_for_approval.onResume.approver_note ?? '') | replace({"'": "''"}) }}', updated_at = now()
WHERE batch_id = '{{ execution.id }}' AND status = 'broadcasting'
- id: notify_sent
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "Bitcoin payout batch sent: {{ outputs.batch_summary.row.payouts }} payouts, {{ outputs.batch_summary.row.total_btc }} BTC at {{ outputs.fees.body | jq('.' ~ inputs.fee_target) | first }} sat/vB. {{ inputs.explorer_tx_url }}{{ outputs.send_batch.body | jq('.result') | first }}"
}
errors:
- id: release_locked
type: io.kestra.plugin.jdbc.postgresql.Query
description: Return payouts that were claimed but never broadcast to the queue.
Rows already marked broadcasting are left alone for manual review.
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
sql: |
UPDATE payouts SET status = 'pending', batch_id = NULL, updated_at = now()
WHERE batch_id = '{{ execution.id }}' AND status = 'locked'
- id: alert_failure
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "Bitcoin payout batch failed in execution {{ execution.id }}. Payouts that were not broadcast are back in the queue. Any rows still marked 'broadcasting' for this batch need a manual check in the wallet before retrying."
}
outputs:
- id: txid
type: STRING
description: Transaction id of the broadcast batch, empty when nothing was sent.
value: "{{ (outputs.send_batch.body ?? '{}') | jq('.result') | first ?? '' }}"
- id: batch
type: JSON
description: Number of payouts and addresses, total BTC, and the amounts sent.
value: "{{ outputs.batch_summary.row ?? {} }}"