id: ai-invoice-math-discrepancy-gate
namespace: company.finance
description: |
Automated invoice auditing and mathematical reconciliation gate that ingests vendor invoice data,
calculates line item totals, verifies tax calculations via in-memory DuckDB, and halts closing
on penny discrepancies before ERP synchronization.
triggers:
- id: scheduled_sweep
type: io.kestra.plugin.core.trigger.Schedule
description: Daily reconciliation sweep. Shipped disabled by default.
cron: "0 18 * * 1-5"
disabled: true
- id: webhook_invoice_ingest
type: io.kestra.plugin.core.trigger.Webhook
description: Authenticated webhook for accounts payable and ERP ingestion pipelines.
key: "{{ secret('WEBHOOK_KEY') }}"
inputs:
- id: invoice_payload
type: STRING
defaults: |
{
"invoice_id": "INV-2026-9042",
"vendor_name": "Apex Cloud Infrastructure Inc",
"tax_rate": 0.0825,
"stated_subtotal": 12500.00,
"stated_tax": 1031.25,
"stated_total": 13531.25,
"line_items": [
{"item_id": "ITEM-101", "description": "Bare Metal Server Nodes", "quantity": 2, "unit_price": 4500.00},
{"item_id": "ITEM-102", "description": "High Throughput NVMe Volume 20TB", "quantity": 1, "unit_price": 1500.00},
{"item_id": "ITEM-103", "description": "Premium 24/7 SLA Engineering Support", "quantity": 1, "unit_price": 2000.00}
]
}
description: Raw JSON payload representing vendor invoice details and line items.
tasks:
- id: prepare_invoice_tables
type: io.kestra.plugin.scripts.python.Script
description: Ingests raw invoice JSON and writes normalized CSV tables for
DuckDB relational processing.
taskRunner:
type: io.kestra.plugin.core.runner.Process
script: |
import json
import csv
raw_payload = """{{ inputs.invoice_payload }}"""
data = json.loads(raw_payload)
# 1. Write Header CSV
with open("invoice_header.csv", "w", newline="") as f:
writer = csv.writer(f)
writer.writerow(["invoice_id", "vendor_name", "tax_rate", "stated_subtotal", "stated_tax", "stated_total"])
writer.writerow([
data["invoice_id"],
data["vendor_name"],
data["tax_rate"],
data["stated_subtotal"],
data["stated_tax"],
data["stated_total"]
])
# 2. Write Line Items CSV
with open("invoice_line_items.csv", "w", newline="") as f:
writer = csv.writer(f)
writer.writerow(["item_id", "description", "quantity", "unit_price"])
for item in data.get("line_items", []):
writer.writerow([
item["item_id"],
item["description"],
item["quantity"],
item["unit_price"]
])
print("Successfully generated invoice_header.csv and invoice_line_items.csv")
outputFiles:
- invoice_header.csv
- invoice_line_items.csv
- id: audit_invoice_math
type: io.kestra.plugin.jdbc.duckdb.Query
description: Executes in-memory DuckDB CTE query to calculate line item sums,
compute expected tax, and identify rounding or calculation variances.
communityExtensions: []
inputFiles:
invoice_header.csv: "{{ outputs.prepare_invoice_tables.outputFiles['invoice_header.csv'] }}"
invoice_line_items.csv: "{{
outputs.prepare_invoice_tables.outputFiles['invoice_line_items.csv'] }}"
fetchType: FETCH_ONE
sql: |
WITH header AS (
SELECT * FROM read_csv_auto('invoice_header.csv', header = true)
),
lines AS (
SELECT * FROM read_csv_auto('invoice_line_items.csv', header = true)
),
calc_subtotal AS (
SELECT
ROUND(SUM(quantity * unit_price), 2) AS calculated_subtotal,
COUNT(*) AS item_count
FROM lines
),
reconciliation AS (
SELECT
h.invoice_id,
h.vendor_name,
h.stated_subtotal,
c.calculated_subtotal,
ROUND(ABS(h.stated_subtotal - c.calculated_subtotal), 2) AS subtotal_variance,
h.stated_tax,
ROUND(c.calculated_subtotal * h.tax_rate, 2) AS calculated_tax,
ROUND(ABS(h.stated_tax - ROUND(c.calculated_subtotal * h.tax_rate, 2)), 2) AS tax_variance,
h.stated_total,
ROUND(c.calculated_subtotal + ROUND(c.calculated_subtotal * h.tax_rate, 2), 2) AS calculated_total,
ROUND(ABS(h.stated_total - (c.calculated_subtotal + ROUND(c.calculated_subtotal * h.tax_rate, 2))), 2) AS total_variance,
c.item_count
FROM header h
CROSS JOIN calc_subtotal c
)
SELECT
*,
CASE
WHEN subtotal_variance > 0.01 OR tax_variance > 0.01 OR total_variance > 0.01 THEN true
ELSE false
END AS has_discrepancy
FROM reconciliation;
- id: discrepancy_gate
type: io.kestra.plugin.core.flow.If
description: Halts closing and routes to audit escalation if mathematical
variance exceeds $0.01 tolerance.
condition: "{{ outputs.audit_invoice_math.row.has_discrepancy == false }}"
then:
- id: approve_for_erp_posting
type: io.kestra.plugin.core.log.Log
description: Logs validation confirmation and passes invoice to ERP accounting sync.
message: "Invoice {{ outputs.audit_invoice_math.row.invoice_id }} approved.
Stated Total: ${{ outputs.audit_invoice_math.row.stated_total }},
Calculated Total: ${{ outputs.audit_invoice_math.row.calculated_total
}}. Variance within tolerance."
- id: notify_accounts_payable
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Dispatches confirmation to Accounts Payable channel.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"channel": "#accounts-payable",
"text": "✅ *Invoice Math Verified*: `{{ outputs.audit_invoice_math.row.invoice_id }}` from *{{ outputs.audit_invoice_math.row.vendor_name }}* approved for ERP posting.\n*Subtotal*: ${{ outputs.audit_invoice_math.row.calculated_subtotal }} | *Tax*: ${{ outputs.audit_invoice_math.row.calculated_tax }} | *Total*: ${{ outputs.audit_invoice_math.row.calculated_total }}"
}
else:
- id: generate_discrepancy_report
type: io.kestra.plugin.scripts.python.Script
description: Generates markdown audit artifact detailing line item and tax
variances.
taskRunner:
type: io.kestra.plugin.core.runner.Process
script: |
invoice_id = "{{ outputs.audit_invoice_math.row.invoice_id }}"
vendor = "{{ outputs.audit_invoice_math.row.vendor_name }}"
subtotal_var = "{{ outputs.audit_invoice_math.row.subtotal_variance }}"
tax_var = "{{ outputs.audit_invoice_math.row.tax_variance }}"
total_var = "{{ outputs.audit_invoice_math.row.total_variance }}"
stated_subtotal = "{{ outputs.audit_invoice_math.row.stated_subtotal }}"
calc_subtotal = "{{ outputs.audit_invoice_math.row.calculated_subtotal }}"
stated_tax = "{{ outputs.audit_invoice_math.row.stated_tax }}"
calc_tax = "{{ outputs.audit_invoice_math.row.calculated_tax }}"
stated_total = "{{ outputs.audit_invoice_math.row.stated_total }}"
calc_total = "{{ outputs.audit_invoice_math.row.calculated_total }}"
report = f"""# INVOICE MATHEMATICAL DISCREPANCY AUDIT
## Executive Summary
- **Invoice ID**: {invoice_id}
- **Vendor**: {vendor}
- **Status**: REJECTED AT GATE (Discrepancy > $0.01)
## Mathematical Breakdown
| Category | Stated ($) | Calculated ($) | Variance ($) |
|---|---|---|---|
| Subtotal | {stated_subtotal} | {calc_subtotal} | {subtotal_var} |
| Tax | {stated_tax} | {calc_tax} | {tax_var} |
| Grand Total | {stated_total} | {calc_total} | {total_var} |
## Required Next Action
Notify vendor accounts receivable of computation variance and withhold automatic payment release.
"""
with open("invoice-discrepancy-report.md", "w") as f:
f.write(report)
print("Saved invoice-discrepancy-report.md")
outputFiles:
- invoice-discrepancy-report.md
- id: escalate_discrepancy_slack
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Alerts Finance Discrepancies channel with exact variance breakdown.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"channel": "#finance-discrepancies",
"text": "🚨 *Invoice Math Discrepancy Gate Triggered*: `{{ outputs.audit_invoice_math.row.invoice_id }}` rejected.\n*Vendor*: {{ outputs.audit_invoice_math.row.vendor_name }}\n*Stated Total*: ${{ outputs.audit_invoice_math.row.stated_total }} vs *Calculated*: ${{ outputs.audit_invoice_math.row.calculated_total }}\n*Total Variance*: ${{ outputs.audit_invoice_math.row.total_variance }}\nAutomatic ERP sync has been halted."
}
errors:
- id: alert_on_failure
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Alerts on-call finance engineering channel if pipeline execution
encounters an unhandled error.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"channel": "#finance-ops",
"text": "⚠️ *Pipeline Error*: Flow `{{ flow.id }}` failed on execution `{{ execution.id }}`."
}