id: intercompany-reconciliation-elimination
namespace: company.team
description: |
Month-end intercompany close for a group of entities: pair every intercompany balance with
its counterparty, translate both sides to the group currency, explain each difference
(timing, FX, missing booking, wrong counterparty), propose the matching entry, generate the
consolidation eliminations, and lock the period only when the group nets to zero.
inputs:
- id: period
type: STRING
displayName: Period (YYYY-MM)
defaults: "2026-09"
validator: ^\d{4}-(0[1-9]|1[0-2])$
- id: scenario
type: SELECT
displayName: Demo scenario
description: >
CLEAN: entities agree up to timing and FX. ISSUES: a management fee booked
by the sender only, a charge posted against the wrong counterparty and a
loan whose interest rates disagree, which must block the consolidation.
values:
- CLEAN
- ISSUES
defaults: CLEAN
- id: group_currency
type: STRING
displayName: Group currency
defaults: EUR
- id: fx_tolerance
type: FLOAT
displayName: FX tolerance (group currency)
description: A difference after translation that is within this amount, between
parties that agree in transaction currency, is an FX difference and goes
to the translation reserve.
defaults: 2500.0
- id: approved_by
type: STRING
displayName: Approved by
description: Needed to consolidate with goods in transit or unbooked charges
carried as proposed entries.
defaults: ""
variables:
lock_key: intercompany_close
pack_dir: "intercompany/{{ inputs.period }}"
concurrency:
limit: 1
triggers:
- id: fourth_business_day
type: io.kestra.plugin.core.trigger.Schedule
description: The 4th of each month, after the entity closes. Disabled until the
entity ledgers are wired in.
cron: "0 8 4 * *"
disabled: true
tasks:
- id: period_control
type: io.kestra.plugin.core.output.OutputValues
values:
last_closed: "{{ (kv(vars.lock_key, errorOnMissing=false) ?? {'period': ''}) |
jq('.period') | first }}"
previous_period: "{{ (inputs.period ~ '-01') | date('yyyy-MM-dd') | dateAdd(-1,
'MONTHS') | date('yyyy-MM') }}"
period_end: "{{ (inputs.period ~ '-01') | date('yyyy-MM-dd') | dateAdd(1,
'MONTHS') | dateAdd(-1, 'DAYS') | date('yyyy-MM-dd') }}"
- id: period_gate
type: io.kestra.plugin.core.flow.If
condition: >-
{{ outputs.period_control.values.last_closed != ''
and (outputs.period_control.values.last_closed >= inputs.period
or outputs.period_control.values.last_closed != outputs.period_control.values.previous_period) }}
then:
- id: period_refused
type: io.kestra.plugin.core.execution.Fail
errorMessage: >-
{{ outputs.period_control.values.last_closed >= inputs.period
? 'Period ' ~ inputs.period ~ ' is already closed (last closed: ' ~ outputs.period_control.values.last_closed ~ ').'
: 'Period ' ~ outputs.period_control.values.previous_period ~ ' must be closed before ' ~ inputs.period ~ ' (last closed: ' ~ outputs.period_control.values.last_closed ~ ').' }}
- id: reconcile
type: io.kestra.plugin.jdbc.duckdb.Queries
description: >
One DuckDB session: entity ledgers (demo generator or your exports),
intercompany pairs, translation at closing and average rates, difference
analysis, proposed entries, eliminations and the consolidation check.
outputFiles:
- pairs
- differences
- proposed
- eliminations
- matrix
- blockers
fetchType: FETCH_ONE
sql: |
SET VARIABLE p_end = DATE '{{ outputs.period_control.values.period_end }}';
SET VARIABLE p_start = DATE '{{ inputs.period }}-01';
SET VARIABLE p_tag = '{{ inputs.period | replace({"-": ""}) }}';
-- =======================================================================
-- The group. Replace with your entity master and the group's rates.
-- Rates are group currency per unit. drift is how far the entity's own
-- month-end revaluation rate sits from the group rate (local bank rates,
-- another fixing time), which is where real translation differences come from.
-- =======================================================================
CREATE TABLE entities (entity VARCHAR, name VARCHAR, currency VARCHAR, drift DOUBLE);
INSERT INTO entities VALUES
('E100', 'Holding SA (FR)', 'EUR', 0), ('E200', 'Sales GmbH (DE)', 'EUR', 0),
('E300', 'Ops Ltd (UK)', 'GBP', 0.0008), ('E400', 'Inc (US)', 'USD', -0.0006),
('E500', 'Pvt Ltd (IN)', 'INR', 0.0011);
CREATE TABLE rates (currency VARCHAR, closing_rate DOUBLE, average_rate DOUBLE);
INSERT INTO rates VALUES ('EUR', 1.0, 1.0), ('GBP', 1.1701, 1.1665), ('USD', 0.8807, 0.8714), ('INR', 0.009189, 0.009162);
-- =======================================================================
-- Intercompany lines as each entity booked them, in transaction currency.
-- The sender books + (receivable, revenue), the receiver - (payable,
-- expense), both with the sender's document number. Replace with the
-- intercompany-flagged lines of each ledger.
-- =======================================================================
CREATE TABLE ic_lines (entity VARCHAR, counterparty VARCHAR, kind VARCHAR, doc VARCHAR,
booked_on DATE, txn_currency VARCHAR, txn_amount DECIMAL(16, 2));
INSERT INTO ic_lines
WITH flows (sender, receiver, kind, txn_currency, base, ship_day) AS (VALUES
('E100', 'E200', 'MGMT_FEE', 'EUR', 42000.00, 31), ('E100', 'E300', 'MGMT_FEE', 'EUR', 38000.00, 31),
('E100', 'E400', 'MGMT_FEE', 'EUR', 51000.00, 31), ('E100', 'E500', 'MGMT_FEE', 'EUR', 12000.00, 31),
('E300', 'E200', 'GOODS', 'EUR', 118000.00, 31), ('E400', 'E300', 'GOODS', 'USD', 96000.00, 24),
('E200', 'E400', 'GOODS', 'EUR', 63000.00, 20), ('E500', 'E400', 'SERVICES', 'USD', 74000.00, 31),
('E500', 'E100', 'SERVICES', 'EUR', 21500.00, 31), ('E400', 'E100', 'ROYALTY', 'USD', 15400.00, 31)
), docs AS (
SELECT f.*, least(m + INTERVAL (ship_day - 1) DAY, m + INTERVAL 1 MONTH - INTERVAL 1 DAY)::DATE AS doc_date,
kind || '-' || sender || '-' || receiver || '-' || strftime(m, '%Y%m') AS doc,
round(base * (0.9 + (hash(sender || receiver || m) % 21)::INT / 100.0), 2) AS amount
FROM flows f, generate_series(DATE '2026-07-01', DATE '2026-11-01', INTERVAL 1 MONTH) g(m)
)
SELECT sender, receiver, kind, doc, doc_date, txn_currency, amount FROM docs
UNION ALL
-- Goods are booked by the receiver on arrival, 3 days after shipping.
SELECT receiver, sender, kind, doc, CASE WHEN kind = 'GOODS' THEN doc_date + 3 ELSE doc_date END, txn_currency, -amount FROM docs
{% if inputs.scenario == 'ISSUES' %}
-- The Indian entity never booked this month's management fee.
WHERE NOT (doc = 'MGMT_FEE-E100-E500-' || getvariable('p_tag') AND receiver = 'E500')
{% endif %};
-- A EUR 2,000,000 loan from the holding to the US entity, interest at 4.5% a month in arrears.
INSERT INTO ic_lines VALUES
('E100', 'E400', 'LOAN', 'LOAN-E100-E400', DATE '2026-01-15', 'EUR', 2000000.00),
('E400', 'E100', 'LOAN', 'LOAN-E100-E400', DATE '2026-01-15', 'EUR', -2000000.00);
INSERT INTO ic_lines
SELECT CASE WHEN side = 1 THEN 'E100' ELSE 'E400' END, CASE WHEN side = 1 THEN 'E400' ELSE 'E100' END, 'INTEREST',
'INTEREST-E100-E400-' || strftime(m, '%Y%m'), (m + INTERVAL 1 MONTH - INTERVAL 1 DAY)::DATE, 'EUR',
CASE WHEN side = 1 THEN round(2000000 * 0.045 / 12, 2)
ELSE -round(2000000 * {% if inputs.scenario == 'ISSUES' %}CASE WHEN strftime(m, '%Y%m') = getvariable('p_tag') THEN 0.040 ELSE 0.045 END{% else %}0.045{% endif %} / 12, 2) END
FROM generate_series(DATE '2026-07-01', DATE '2026-11-01', INTERVAL 1 MONTH) g(m), range(1, 3) t(side);
{% if inputs.scenario == 'ISSUES' %}
-- The German entity booked a UK goods invoice against the US entity.
UPDATE ic_lines SET counterparty = 'E400' WHERE entity = 'E200' AND doc = 'GOODS-E300-E200-' || getvariable('p_tag');
{% endif %}
-- Earlier months are settled. Open at period end: the loan and this month's documents.
CREATE TABLE open_lines AS
SELECT * FROM ic_lines
WHERE booked_on <= getvariable('p_end') AND (kind = 'LOAN' OR right(doc, 6) = getvariable('p_tag'));
-- =======================================================================
-- Document by document: who booked what, when, against whom.
-- =======================================================================
CREATE TABLE docs AS
SELECT doc, any_value(kind) AS kind, any_value(txn_currency) AS txn_currency,
any_value(entity) FILTER (WHERE txn_amount > 0) AS sender,
any_value(counterparty) FILTER (WHERE txn_amount > 0) AS sender_says,
any_value(entity) FILTER (WHERE txn_amount < 0) AS receiver,
any_value(counterparty) FILTER (WHERE txn_amount < 0) AS receiver_says,
max(booked_on) FILTER (WHERE txn_amount > 0) AS sent_on,
max(booked_on) FILTER (WHERE txn_amount < 0) AS received_on,
sum(txn_amount) FILTER (WHERE txn_amount > 0) AS sent_amount,
-sum(txn_amount) FILTER (WHERE txn_amount < 0) AS received_amount
FROM ic_lines WHERE kind = 'LOAN' OR right(doc, 6) = getvariable('p_tag')
GROUP BY doc;
CREATE TABLE differences AS
SELECT doc, kind, sender, coalesce(receiver, sender_says) AS receiver, txn_currency, sent_amount, received_amount, sent_on, received_on,
CASE
WHEN receiver IS NULL THEN 'NOT_BOOKED_BY_RECEIVER'
WHEN receiver_says <> sender OR sender_says <> receiver THEN 'WRONG_COUNTERPARTY'
WHEN sent_amount <> received_amount THEN 'AMOUNT_MISMATCH'
WHEN received_on > getvariable('p_end') THEN 'IN_TRANSIT'
END AS cause,
CASE
WHEN receiver IS NULL THEN sender_says || ' has not booked ' || sent_amount || ' ' || txn_currency || '. Book it, or ask ' || sender || ' to reverse it'
WHEN receiver_says <> sender OR sender_says <> receiver THEN receiver || ' booked it against ' || receiver_says || ' instead of ' || sender || '. Change the counterparty'
WHEN sent_amount <> received_amount THEN sender || ' booked ' || sent_amount || ', ' || receiver || ' booked ' || received_amount || ' ' || txn_currency || '. Agree the amount'
ELSE 'Shipped ' || sent_on || ', received ' || received_on || '. Book goods in transit in ' || receiver
END AS action
FROM docs
WHERE receiver IS NULL OR receiver_says <> sender OR sender_says <> receiver OR sent_amount <> received_amount OR received_on > getvariable('p_end');
CREATE TABLE blockers AS
SELECT cause AS code, doc AS ref, action AS detail FROM differences WHERE cause <> 'IN_TRANSIT';
-- Timing items are booked in the receiving entity at period end (periodic inventory:
-- purchases in transit), so both sides agree.
CREATE TABLE proposed AS
SELECT receiver AS entity, sender AS counterparty, 'PURCHASES_IN_TRANSIT' AS kind, doc, getvariable('p_end') AS booked_on, txn_currency, -sent_amount AS txn_amount
FROM differences WHERE cause = 'IN_TRANSIT';
-- =======================================================================
-- Pairs after the proposed entries, in transaction currency.
-- =======================================================================
CREATE TABLE pairs AS
SELECT least(entity, counterparty) AS entity_a, greatest(entity, counterparty) AS entity_b, txn_currency,
sum(txn_amount) FILTER (WHERE entity = least(entity, counterparty)) AS balance_a,
sum(txn_amount) FILTER (WHERE entity = greatest(entity, counterparty)) AS balance_b,
sum(txn_amount) AS difference
FROM (SELECT entity, counterparty, txn_currency, txn_amount FROM open_lines
UNION ALL SELECT entity, counterparty, txn_currency, txn_amount FROM proposed)
GROUP BY ALL;
-- A pair still apart that no document explains is a blocker of its own.
INSERT INTO blockers
SELECT 'UNEXPLAINED_PAIR_DIFFERENCE', p.entity_a || '/' || p.entity_b || ' ' || p.txn_currency, 'Still ' || p.difference || ' apart'
FROM pairs p
WHERE p.difference <> 0 AND NOT EXISTS (
SELECT 1 FROM differences d WHERE d.cause <> 'IN_TRANSIT'
AND (p.entity_a IN (d.sender, d.receiver) OR p.entity_b IN (d.sender, d.receiver)));
-- =======================================================================
-- Eliminations in group currency. Each side translates its balance as its
-- own revaluation left it (transaction amount at the group closing rate,
-- moved by the entity's drift when it is a foreign currency balance). The
-- residual is the translation difference and goes to the CTA. Income and
-- expense are eliminated at the average rate.
-- =======================================================================
CREATE TABLE eliminations AS
SELECT 'BALANCE_SHEET' AS kind, l.entity, l.counterparty,
CASE WHEN l.kind = 'LOAN' THEN 'Intercompany loan' WHEN l.txn_amount > 0 THEN 'Intercompany receivable' ELSE 'Intercompany payable' END AS account,
l.txn_currency, sum(l.txn_amount) AS txn_amount,
round(-sum(l.txn_amount) * r.closing_rate * (1 + CASE WHEN l.txn_currency <> e.currency THEN e.drift ELSE 0 END), 2)::DECIMAL(16, 2) AS amount_group
FROM (SELECT * FROM open_lines UNION ALL SELECT * FROM proposed) l
JOIN entities e USING (entity) JOIN rates r ON r.currency = l.txn_currency
GROUP BY l.entity, l.counterparty, 4, l.txn_currency, r.closing_rate, e.drift, e.currency
UNION ALL
SELECT 'PROFIT_AND_LOSS', l.entity, l.counterparty,
CASE WHEN l.txn_amount > 0 THEN 'Intercompany revenue' ELSE 'Intercompany expense' END,
l.txn_currency, sum(l.txn_amount), round(-sum(l.txn_amount) * r.average_rate, 2)::DECIMAL(16, 2)
-- By document month: last month's goods in transit are received this month, but the
-- in-transit entry booked last month reverses on day one, so the receipt nets out.
FROM (SELECT * FROM ic_lines WHERE kind <> 'LOAN' AND right(doc, 6) = getvariable('p_tag') AND booked_on <= getvariable('p_end')
UNION ALL SELECT * FROM proposed) l
JOIN rates r ON r.currency = l.txn_currency
GROUP BY l.entity, l.counterparty, 4, l.txn_currency, r.average_rate;
CREATE TABLE totals AS
SELECT coalesce(sum(amount_group) FILTER (WHERE kind = 'BALANCE_SHEET'), 0) AS translation_difference,
coalesce(sum(amount_group) FILTER (WHERE kind = 'PROFIT_AND_LOSS'), 0) AS pl_difference,
coalesce(sum(amount_group) FILTER (WHERE kind = 'BALANCE_SHEET' AND amount_group > 0), 0) AS bs_eliminated,
coalesce(sum(amount_group) FILTER (WHERE kind = 'PROFIT_AND_LOSS' AND amount_group > 0), 0) AS pl_eliminated
FROM eliminations;
INSERT INTO eliminations
SELECT 'TRANSLATION', 'GROUP', '', 'Translation reserve (CTA)', '{{ inputs.group_currency }}', NULL, -translation_difference
FROM totals WHERE translation_difference <> 0;
-- Who owes whom, in group currency, for the review pack.
CREATE TABLE matrix AS
PIVOT (SELECT entity, counterparty, round(-amount_group, 0) AS v FROM eliminations WHERE kind = 'BALANCE_SHEET')
ON counterparty USING sum(v) GROUP BY entity ORDER BY entity;
COPY (SELECT * FROM pairs ORDER BY entity_a, entity_b, txn_currency) TO '{{ outputFiles.pairs }}' (HEADER, DELIMITER ',');
COPY (SELECT * FROM differences ORDER BY cause, doc) TO '{{ outputFiles.differences }}' (HEADER, DELIMITER ',');
COPY (SELECT entity, counterparty, booked_on AS posting_date, txn_currency, doc,
'5150 Purchases in transit (intercompany)' AS debit_account, '2310 Intercompany payable' AS credit_account, -txn_amount AS amount
FROM proposed ORDER BY doc) TO '{{ outputFiles.proposed }}' (HEADER, DELIMITER ',');
COPY (SELECT * FROM eliminations ORDER BY kind, entity, counterparty, account) TO '{{ outputFiles.eliminations }}' (HEADER, DELIMITER ',');
COPY matrix TO '{{ outputFiles.matrix }}' (HEADER, DELIMITER ',');
COPY blockers TO '{{ outputFiles.blockers }}' (HEADER, DELIMITER ',');
SELECT
(SELECT count(*) FROM entities) AS entities,
(SELECT count(*) FROM docs) AS documents,
(SELECT count(*) FROM pairs) AS pairs,
(SELECT count(*) FROM pairs WHERE difference = 0) AS pairs_agreed,
(SELECT count(*) FROM differences WHERE cause = 'IN_TRANSIT') AS in_transit,
(SELECT coalesce(string_agg(doc || ' ' || sent_amount || ' ' || txn_currency, ', '), '') FROM differences WHERE cause = 'IN_TRANSIT') AS in_transit_detail,
t.translation_difference, t.pl_difference, t.bs_eliminated, t.pl_eliminated,
(SELECT sum(amount_group) FROM eliminations) AS group_net,
abs(t.translation_difference) <= {{ inputs.fx_tolerance }} AS fx_within_tolerance,
(SELECT count(*) FROM blockers) AS blockers,
(SELECT coalesce(string_agg(code || ' ' || ref || ': ' || detail, '; ' ORDER BY code, ref), '') FROM blockers) AS blocker_detail
FROM totals t;
- id: result
type: io.kestra.plugin.core.output.OutputValues
values:
r: "{{ (outputs.reconcile.outputs | last).row | toJson }}"
- id: log_close
type: io.kestra.plugin.core.log.Log
message: |
Intercompany close {{ inputs.period }} ({{ inputs.scenario }}), group currency {{ inputs.group_currency }}.
{{ outputs.result.values.r | jq('"\(.entities) entities, \(.documents) intercompany documents, \(.pairs) counterparty pairs, \(.pairs_agreed) agreed after timing items."') | first }}
{{ outputs.result.values.r | jq('"Goods in transit: \(.in_transit) \(.in_transit_detail)"') | first }}
{{ outputs.result.values.r | jq('"Eliminated: balance sheet \(.bs_eliminated), P&L \(.pl_eliminated). Translation difference to CTA \(.translation_difference) (within tolerance: \(.fx_within_tolerance)). P&L difference \(.pl_difference). After the CTA the group nets to \(.group_net)."') | first }}
{{ outputs.result.values.r | jq('"Blockers: \(.blockers). \(.blocker_detail)"') | first }}
- id: evidence
type: io.kestra.plugin.core.namespace.UploadFiles
namespace: "{{ flow.namespace }}"
filesMap:
"{{ render(vars.pack_dir) }}/pairs.csv": "{{ outputs.reconcile.outputFiles.pairs }}"
"{{ render(vars.pack_dir) }}/differences.csv": "{{ outputs.reconcile.outputFiles.differences }}"
"{{ render(vars.pack_dir) }}/proposed-entries.csv": "{{ outputs.reconcile.outputFiles.proposed }}"
"{{ render(vars.pack_dir) }}/eliminations.csv": "{{ outputs.reconcile.outputFiles.eliminations }}"
"{{ render(vars.pack_dir) }}/counterparty-matrix.csv": "{{ outputs.reconcile.outputFiles.matrix }}"
"{{ render(vars.pack_dir) }}/blockers.csv": "{{ outputs.reconcile.outputFiles.blockers }}"
- id: decide
type: io.kestra.plugin.core.flow.Switch
value: >-
{%- set r = outputs.result.values.r -%} {%- if (r | jq('.blockers') |
first) > 0 -%}BLOCKED {%- elseif not (r | jq('.fx_within_tolerance') |
first) or (r | jq('.pl_difference') | first) != 0 -%}OUT_OF_TOLERANCE {%-
elseif (r | jq('.in_transit') | first) > 0 and inputs.approved_by == ''
-%}NEEDS_APPROVAL {%- else -%}READY{%- endif -%}
cases:
READY:
- id: lock_period
type: io.kestra.plugin.core.kv.Set
key: "{{ vars.lock_key }}"
kvType: JSON
value: |
{"period": "{{ inputs.period }}", "closed_at": "{{ now() }}", "execution_id": "{{ execution.id }}",
"bs_eliminated": {{ outputs.result.values.r | jq('.bs_eliminated') | first }},
"pl_eliminated": {{ outputs.result.values.r | jq('.pl_eliminated') | first }},
"translation_difference": {{ outputs.result.values.r | jq('.translation_difference') | first }},
"approved_by": "{{ inputs.approved_by }}"}
- id: log_closed
type: io.kestra.plugin.core.log.Log
message: "Intercompany for {{ inputs.period }} closed. Post proposed-entries.csv
in the entities and eliminations.csv in consolidation, from {{
render(vars.pack_dir) }}."
NEEDS_APPROVAL:
- id: approval_required
type: io.kestra.plugin.core.execution.Fail
errorMessage: >-
Every pair agrees once goods in transit are booked: {{
outputs.result.values.r | jq('.in_transit_detail') | first }}.
Review {{ render(vars.pack_dir) }}/proposed-entries.csv and run
again with approved_by.
OUT_OF_TOLERANCE:
- id: out_of_tolerance
type: io.kestra.plugin.core.execution.Fail
errorMessage: >-
The eliminations do not net: translation difference {{
outputs.result.values.r | jq('.translation_difference') | first }}
(tolerance {{ inputs.fx_tolerance }}), P&L difference {{
outputs.result.values.r | jq('.pl_difference') | first }}.
defaults:
- id: blocked
type: io.kestra.plugin.core.execution.Fail
errorMessage: "Intercompany for {{ inputs.period }} not closed: {{
outputs.result.values.r | jq('.blocker_detail') | first }}"