id: bitcoin-deposit-confirmation-tracker
namespace: company.team
description: |
Track Bitcoin invoice payments from first sighting to final confirmation.
Bitcoin Core calls the webhook whenever the wallet sees a transaction; the
flow re-checks every open invoice, classifies it as paid, overpaid,
underpaid, or expired, catches reorgs, and tells your app and Slack.
concurrency:
limit: 1
behavior: QUEUE
triggers:
- id: wallet_notify
type: io.kestra.plugin.core.trigger.Webhook
description: "Call this from Bitcoin Core's walletnotify so every incoming
transaction triggers a check within seconds, e.g. walletnotify=curl -s -X
POST
https://kestra.example.com/api/v1/main/executions/webhook/company.team/bi\
tcoin-deposit-confirmation-tracker/btc-deposits-change-me?txid=%s"
key: btc-deposits-change-me
- id: confirmation_sweep
type: io.kestra.plugin.core.trigger.Schedule
description: walletnotify fires when a transaction arrives and when it first
confirms, but not for later confirmations or reorgs. This sweep moves
invoices across the confirmation threshold and catches reorgs. Shipped
disabled.
disabled: true
cron: "*/10 * * * *"
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: deposits
description: Bitcoin Core wallet that owns the invoice addresses.
- id: required_confirmations
type: INT
defaults: 3
description: Confirmations a payment needs before the invoice counts as paid.
- id: overpay_tolerance_btc
type: FLOAT
defaults: 0.00001
description: Payments above the amount due by more than this are flagged as
overpaid so you can refund the difference.
- id: reorg_window_hours
type: INT
defaults: 24
description: How long paid invoices stay under watch for chain reorganizations.
- id: max_invoices
type: INT
defaults: 200
description: Maximum number of invoices checked in one run.
- id: callback_url
type: STRING
defaults: ""
description: Optional URL in your application that receives a JSON POST for
every invoice status change. Leave empty to only use Slack.
- id: explorer_address_url
type: STRING
defaults: https://mempool.space/signet/address/
description: Block explorer prefix for address links in notifications.
tasks:
- id: node_check
type: io.kestra.plugin.core.http.Request
description: Fail fast with a clear error when the node is unreachable, the
credentials are wrong, or the wallet is not loaded, before any invoice is
touched.
uri: "{{ inputs.rpc_url }}/wallet/{{ inputs.wallet }}"
method: POST
body: '{"jsonrpc": "1.0", "id": "kestra", "method": "getwalletinfo", "params":
[]}'
options:
auth:
type: BASIC
username: "{{ secret('BITCOIN_RPC_USER') }}"
password: "{{ secret('BITCOIN_RPC_PASSWORD') }}"
retry:
type: exponential
interval: PT2S
maxInterval: PT30S
maxAttempts: 3
- id: ensure_table
type: io.kestra.plugin.jdbc.postgresql.Query
description: Create the invoice table on first run. Your application inserts one
row per invoice with a fresh address from getnewaddress.
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
sql: |
CREATE TABLE IF NOT EXISTS invoices (
id BIGSERIAL PRIMARY KEY,
address TEXT NOT NULL UNIQUE,
amount_due_btc NUMERIC(16, 8) NOT NULL CHECK (amount_due_btc > 0),
status TEXT NOT NULL DEFAULT 'awaiting_payment',
received_btc NUMERIC(16, 8) NOT NULL DEFAULT 0,
confirmed_btc NUMERIC(16, 8) NOT NULL DEFAULT 0,
expires_at TIMESTAMPTZ NOT NULL DEFAULT now() + interval '1 hour',
paid_at TIMESTAMPTZ,
notified_status TEXT,
previous_notified_status TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
)
- id: load_open_invoices
type: io.kestra.plugin.jdbc.postgresql.Query
description: Invoices that can still change, plus recently paid ones that could
be undone by a reorg.
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
fetchType: FETCH
sql: |
SELECT id, address FROM invoices
WHERE status IN ('awaiting_payment', 'seen', 'underpaid')
OR (status IN ('paid', 'overpaid') AND paid_at > now() - interval '{{ inputs.reorg_window_hours }} hours')
OR status IS DISTINCT FROM coalesce(notified_status, 'awaiting_payment')
ORDER BY id
LIMIT {{ inputs.max_invoices }}
- id: nothing_open
type: io.kestra.plugin.core.flow.If
condition: "{{ outputs.load_open_invoices.size == 0 }}"
then:
- id: no_open_invoices
type: io.kestra.plugin.core.execution.Exit
state: SUCCESS
- id: check_invoices
type: io.kestra.plugin.core.flow.Loop
description: For each open invoice, ask the node what the address has received
with and without enough confirmations, then let Postgres decide the new
status. The node recomputes these totals from the current chain, so a
reorg simply shows up as a lower confirmed amount.
values: "{{ outputs.load_open_invoices.rows | toJson }}"
concurrencyLimit: 5
tasks:
- id: received_any
type: io.kestra.plugin.core.http.Request
uri: "{{ inputs.rpc_url }}/wallet/{{ inputs.wallet }}"
method: POST
body: '{"jsonrpc": "1.0", "id": "kestra", "method": "getreceivedbyaddress",
"params": ["{{ fromJson(item.value).address }}", 0]}'
options:
auth:
type: BASIC
username: "{{ secret('BITCOIN_RPC_USER') }}"
password: "{{ secret('BITCOIN_RPC_PASSWORD') }}"
retry:
type: exponential
interval: PT2S
maxInterval: PT30S
maxAttempts: 3
- id: received_confirmed
type: io.kestra.plugin.core.http.Request
uri: "{{ inputs.rpc_url }}/wallet/{{ inputs.wallet }}"
method: POST
body: '{"jsonrpc": "1.0", "id": "kestra", "method": "getreceivedbyaddress",
"params": ["{{ fromJson(item.value).address }}", {{
inputs.required_confirmations }}]}'
options:
auth:
type: BASIC
username: "{{ secret('BITCOIN_RPC_USER') }}"
password: "{{ secret('BITCOIN_RPC_PASSWORD') }}"
retry:
type: exponential
interval: PT2S
maxInterval: PT30S
maxAttempts: 3
- id: update_status
type: io.kestra.plugin.jdbc.postgresql.Query
description: Classify the invoice and claim the notification in one statement.
The UPDATE locks the row, so when two runs check the same invoice at
once only the first sees a change; the second finds the new status
already claimed and stays quiet.
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
fetchType: FETCH_ONE
sql: |
WITH amounts AS (
SELECT {{ outputs.received_any.body | jq('.result') | first }}::numeric AS seen,
{{ outputs.received_confirmed.body | jq('.result') | first }}::numeric AS confirmed
), classified AS (
SELECT i.id, a.seen, a.confirmed,
CASE
WHEN a.confirmed > i.amount_due_btc + {{ inputs.overpay_tolerance_btc }} THEN 'overpaid'
WHEN a.confirmed >= i.amount_due_btc THEN 'paid'
WHEN a.seen >= i.amount_due_btc THEN 'seen'
WHEN a.seen > 0 THEN 'underpaid'
WHEN i.expires_at < now() THEN 'expired'
ELSE 'awaiting_payment'
END AS status
FROM invoices i, amounts a
WHERE i.id = {{ fromJson(item.value).id }}
)
UPDATE invoices i SET
received_btc = c.seen,
confirmed_btc = c.confirmed,
status = c.status,
paid_at = CASE WHEN c.status IN ('paid', 'overpaid') THEN coalesce(i.paid_at, now()) END,
previous_notified_status = coalesce(i.notified_status, 'awaiting_payment'),
notified_status = c.status,
updated_at = now()
FROM classified c
WHERE i.id = c.id
RETURNING i.id, i.address, i.amount_due_btc, i.received_btc, i.confirmed_btc,
greatest(i.confirmed_btc - i.amount_due_btc, 0) AS refund_btc,
greatest(i.amount_due_btc - i.received_btc, 0) AS shortfall_btc,
i.previous_notified_status AS old_status, i.status AS new_status,
i.previous_notified_status <> i.status AS changed,
i.previous_notified_status IN ('paid', 'overpaid') AND i.status NOT IN ('paid', 'overpaid') AS reorged
- id: status_changed
type: io.kestra.plugin.core.flow.If
condition: "{{ outputs.update_status.row.changed }}"
then:
- id: deliver
type: io.kestra.plugin.core.flow.Sequential
description: Tell your application first, then Slack. If a step still fails
after its retries, the local errors handler releases the claim and
reports the failed delivery; allowFailure keeps one unreachable
endpoint from failing the whole run, and the next walletnotify
call or sweep delivers the change again.
allowFailure: true
tasks:
- id: app_callback
type: io.kestra.plugin.core.flow.If
condition: "{{ inputs.callback_url != '' }}"
then:
- id: post_callback
type: io.kestra.plugin.core.http.Request
description: Tell your application about the status change so it can ship goods,
credit a balance, or show the customer an update. The
Idempotency-Key header lets it ignore a change it already
processed.
uri: "{{ inputs.callback_url }}"
method: POST
headers:
Content-Type: application/json
Idempotency-Key: "invoice-{{ outputs.update_status.row.id }}-{{
outputs.update_status.row.old_status }}-{{
outputs.update_status.row.new_status }}"
body: |
{"invoice_id": {{ outputs.update_status.row.id }}, "address": "{{ outputs.update_status.row.address }}", "old_status": "{{ outputs.update_status.row.old_status }}", "status": "{{ outputs.update_status.row.new_status }}", "amount_due_btc": {{ outputs.update_status.row.amount_due_btc }}, "received_btc": {{ outputs.update_status.row.received_btc }}, "confirmed_btc": {{ outputs.update_status.row.confirmed_btc }}}
retry:
type: exponential
interval: PT2S
maxInterval: PT30S
maxAttempts: 3
- id: reorg_check
type: io.kestra.plugin.core.flow.If
description: A paid invoice whose confirmed amount dropped below the amount due
was hit by a reorg.
condition: "{{ outputs.update_status.row.reorged }}"
then:
- id: alert_reorg
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
messageText: "REORG: invoice #{{ outputs.update_status.row.id }} was {{
outputs.update_status.row.old_status }} and is now {{
outputs.update_status.row.new_status }} ({{
outputs.update_status.row.confirmed_btc }} of {{
outputs.update_status.row.amount_due_btc }} BTC
confirmed). Hold any goods or credit released for it. {{
inputs.explorer_address_url }}{{
outputs.update_status.row.address }}"
retry:
type: exponential
interval: PT2S
maxInterval: PT30S
maxAttempts: 3
- id: notify_by_status
type: io.kestra.plugin.core.flow.Switch
value: "{{ outputs.update_status.row.new_status }}"
cases:
paid:
- id: notify_paid
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
messageText: "Invoice #{{ outputs.update_status.row.id }} paid: {{
outputs.update_status.row.confirmed_btc }} BTC with {{
inputs.required_confirmations }}+ confirmations. {{
inputs.explorer_address_url }}{{
outputs.update_status.row.address }}"
retry:
type: exponential
interval: PT2S
maxInterval: PT30S
maxAttempts: 3
overpaid:
- id: notify_overpaid
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
messageText: "Invoice #{{ outputs.update_status.row.id }} overpaid: received {{
outputs.update_status.row.confirmed_btc }} BTC for {{
outputs.update_status.row.amount_due_btc }} due. Refund
{{ outputs.update_status.row.refund_btc }} BTC. {{
inputs.explorer_address_url }}{{
outputs.update_status.row.address }}"
retry:
type: exponential
interval: PT2S
maxInterval: PT30S
maxAttempts: 3
underpaid:
- id: notify_underpaid
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
messageText: "Invoice #{{ outputs.update_status.row.id }} underpaid: {{
outputs.update_status.row.received_btc }} of {{
outputs.update_status.row.amount_due_btc }} BTC
received. Ask the customer for the remaining {{
outputs.update_status.row.shortfall_btc }} BTC. {{
inputs.explorer_address_url }}{{
outputs.update_status.row.address }}"
retry:
type: exponential
interval: PT2S
maxInterval: PT30S
maxAttempts: 3
defaults:
- id: log_transition
type: io.kestra.plugin.core.log.Log
message: "Invoice #{{ outputs.update_status.row.id }}: {{
outputs.update_status.row.old_status }} -> {{
outputs.update_status.row.new_status }} ({{
outputs.update_status.row.received_btc }} BTC seen, {{
outputs.update_status.row.confirmed_btc }} confirmed)."
errors:
- id: release_claim
type: io.kestra.plugin.jdbc.postgresql.Query
description: Delivery failed, so put the last delivered status back. The next
run sees the change again and re-sends it. The WHERE clause
leaves the row alone if a newer change was claimed in the
meantime.
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
sql: |
UPDATE invoices SET notified_status = '{{ outputs.update_status.row.old_status }}'
WHERE id = {{ outputs.update_status.row.id }} AND notified_status = '{{ outputs.update_status.row.new_status }}'
- id: alert_delivery_failed
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
messageText: >-
Could not deliver the change of invoice #{{
outputs.update_status.row.id }} from {{
outputs.update_status.row.old_status }} to {{
outputs.update_status.row.new_status }} after 3 attempts.
Check that {{ inputs.callback_url }} is reachable. The change
stays pending and is sent again on the next walletnotify call
or sweep.
- id: summary
type: io.kestra.plugin.jdbc.postgresql.Query
description: Count invoices by status for the execution logs and outputs.
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
fetchType: FETCH
sql: |
SELECT status, count(*) AS invoices, sum(confirmed_btc) AS confirmed_btc
FROM invoices GROUP BY status ORDER BY status
errors:
- id: alert_failure
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
messageText: >-
Bitcoin deposit check failed in execution {{ execution.id }} at `{{
errorLogs()[0]['taskId'] ?? 'unknown task' }}`: {{
(errorLogs()[0]['message'] ?? 'see the execution logs') | split(' \[\[') |
first }}. Nothing is lost: invoices this run did not reach are checked
again on the next walletnotify call or sweep.
outputs:
- id: invoices_by_status
type: JSON
description: Invoice counts and confirmed BTC per status after this run.
value: "{{ outputs.summary.rows ?? [] }}"