832 lines
27 KiB
Python
832 lines
27 KiB
Python
from __future__ import annotations
|
|
|
|
import json
|
|
import re
|
|
from datetime import date, datetime, timezone
|
|
from difflib import SequenceMatcher
|
|
|
|
from sqlalchemy import delete, func, select
|
|
|
|
from app.modules.accounting.bank_models import AccountingBankTransaction
|
|
from app.modules.accounting.bank_reconciliation_models import (
|
|
AccountingBankLedgerMapping,
|
|
BankReconciliationItem,
|
|
BankReconciliationRun,
|
|
)
|
|
from app.modules.accounting.chart_models import AccountingChartLedger
|
|
from app.modules.accounting.chart_service import effective_role
|
|
from app.modules.bank_statement_analyzer.models import BankStatementAnalysisJob
|
|
from app.modules.documents.agent_jobs import enqueue_agent_job
|
|
from app.modules.documents.models import ERPAgentJob, ERPWorkstationAgent
|
|
|
|
ACTION = "accounting_bank_reconciliation_extract"
|
|
|
|
|
|
def _utcnow():
|
|
return datetime.now(timezone.utc)
|
|
|
|
|
|
def _s(value):
|
|
return str(value or "").strip()
|
|
|
|
|
|
def _loads(value, default=None):
|
|
try:
|
|
return json.loads(value or "")
|
|
except Exception:
|
|
return {} if default is None else default
|
|
|
|
|
|
def _norm_ref(value):
|
|
return re.sub(r"[^A-Z0-9]", "", _s(value).upper())
|
|
|
|
|
|
def _norm_text(value):
|
|
return re.sub(r"[^A-Z0-9]+", " ", _s(value).upper()).strip()
|
|
|
|
|
|
def _date_obj(value):
|
|
try:
|
|
return date.fromisoformat(_s(value)[:10])
|
|
except Exception:
|
|
return None
|
|
|
|
|
|
def visible_bank_ledgers(db, *, tenant_id: int, client_id: int, tally_guid: str):
|
|
rows = list(
|
|
db.execute(
|
|
select(AccountingChartLedger)
|
|
.where(
|
|
AccountingChartLedger.tenant_id == int(tenant_id),
|
|
AccountingChartLedger.client_id == int(client_id),
|
|
AccountingChartLedger.tally_guid == _s(tally_guid),
|
|
)
|
|
.order_by(AccountingChartLedger.name)
|
|
).scalars().all()
|
|
)
|
|
# Prefer the semantic role, but also trust Tally's immediate parent group.
|
|
# This keeps reconciliation usable even before a previously mis-resolved
|
|
# Phase 16 root hierarchy is manually refreshed.
|
|
return [
|
|
row
|
|
for row in rows
|
|
if effective_role(row) == "BANK"
|
|
or _norm_text(row.parent_group_name) == "BANK ACCOUNTS"
|
|
or _norm_text(row.root_group_name) == "BANK ACCOUNTS"
|
|
]
|
|
|
|
|
|
def visible_workstations(db, *, tenant_id: int, branch_id: int | None = None):
|
|
stmt = select(ERPWorkstationAgent).where(
|
|
ERPWorkstationAgent.tenant_id == int(tenant_id),
|
|
ERPWorkstationAgent.is_active.is_(True),
|
|
ERPWorkstationAgent.tally_connected.is_(True),
|
|
)
|
|
if branch_id is not None:
|
|
stmt = stmt.where(ERPWorkstationAgent.branch_id == int(branch_id))
|
|
return list(
|
|
db.execute(
|
|
stmt.order_by(
|
|
ERPWorkstationAgent.last_seen_at_utc.desc(),
|
|
ERPWorkstationAgent.id.desc(),
|
|
)
|
|
).scalars().all()
|
|
)
|
|
|
|
|
|
def completed_client_jobs(db, *, tenant_id: int, client_id: int, limit: int = 50):
|
|
return list(
|
|
db.execute(
|
|
select(BankStatementAnalysisJob)
|
|
.where(
|
|
BankStatementAnalysisJob.tenant_id == int(tenant_id),
|
|
BankStatementAnalysisJob.client_id == int(client_id),
|
|
BankStatementAnalysisJob.status == "completed",
|
|
BankStatementAnalysisJob.ownership_status == "confirmed",
|
|
)
|
|
.order_by(BankStatementAnalysisJob.completed_at_utc.desc())
|
|
.limit(limit)
|
|
).scalars().all()
|
|
)
|
|
|
|
|
|
def source_accounts(db, *, tenant_id: int, client_id: int, job_id: str):
|
|
rows = list(
|
|
db.execute(
|
|
select(
|
|
AccountingBankTransaction.account_number,
|
|
AccountingBankTransaction.bank_name,
|
|
func.count(AccountingBankTransaction.id),
|
|
)
|
|
.where(
|
|
AccountingBankTransaction.tenant_id == int(tenant_id),
|
|
AccountingBankTransaction.client_id == int(client_id),
|
|
AccountingBankTransaction.source_job_id == job_id,
|
|
)
|
|
.group_by(
|
|
AccountingBankTransaction.account_number,
|
|
AccountingBankTransaction.bank_name,
|
|
)
|
|
.order_by(
|
|
AccountingBankTransaction.bank_name,
|
|
AccountingBankTransaction.account_number,
|
|
)
|
|
).all()
|
|
)
|
|
return [
|
|
{
|
|
"account_number": _s(account),
|
|
"bank_name": _s(bank),
|
|
"transaction_count": int(count or 0),
|
|
}
|
|
for account, bank, count in rows
|
|
if _s(account)
|
|
]
|
|
|
|
|
|
def list_bank_mappings(db, *, tenant_id: int, client_id: int, tally_guid: str = ""):
|
|
stmt = select(AccountingBankLedgerMapping).where(
|
|
AccountingBankLedgerMapping.tenant_id == int(tenant_id),
|
|
AccountingBankLedgerMapping.client_id == int(client_id),
|
|
)
|
|
if _s(tally_guid):
|
|
stmt = stmt.where(AccountingBankLedgerMapping.tally_guid == _s(tally_guid))
|
|
return list(
|
|
db.execute(
|
|
stmt.order_by(
|
|
AccountingBankLedgerMapping.bank_name,
|
|
AccountingBankLedgerMapping.account_number,
|
|
)
|
|
).scalars().all()
|
|
)
|
|
|
|
|
|
def save_bank_mapping(
|
|
db,
|
|
*,
|
|
tenant_id: int,
|
|
client_id: int,
|
|
account_number: str,
|
|
bank_name: str,
|
|
tally_guid: str,
|
|
company_name: str,
|
|
bank_ledger_name: str,
|
|
user_id: int,
|
|
):
|
|
account_number = _s(account_number)
|
|
if not account_number:
|
|
raise ValueError("Bank account number is required.")
|
|
|
|
allowed = {
|
|
row.name
|
|
for row in visible_bank_ledgers(
|
|
db,
|
|
tenant_id=tenant_id,
|
|
client_id=client_id,
|
|
tally_guid=tally_guid,
|
|
)
|
|
}
|
|
if bank_ledger_name not in allowed:
|
|
raise ValueError("Select a valid Tally Bank ledger from Chart of Accounts.")
|
|
|
|
mapping = db.execute(
|
|
select(AccountingBankLedgerMapping).where(
|
|
AccountingBankLedgerMapping.tenant_id == int(tenant_id),
|
|
AccountingBankLedgerMapping.client_id == int(client_id),
|
|
AccountingBankLedgerMapping.account_number == account_number,
|
|
AccountingBankLedgerMapping.tally_guid == _s(tally_guid),
|
|
)
|
|
).scalar_one_or_none()
|
|
|
|
if mapping is None:
|
|
mapping = AccountingBankLedgerMapping(
|
|
tenant_id=int(tenant_id),
|
|
client_id=int(client_id),
|
|
account_number=account_number,
|
|
tally_guid=_s(tally_guid),
|
|
created_by_user_id=int(user_id),
|
|
)
|
|
|
|
mapping.bank_name = _s(bank_name)
|
|
mapping.company_name = _s(company_name)
|
|
mapping.bank_ledger_name = _s(bank_ledger_name)
|
|
mapping.updated_by_user_id = int(user_id)
|
|
mapping.updated_at_utc = _utcnow()
|
|
db.add(mapping)
|
|
db.commit()
|
|
db.refresh(mapping)
|
|
return mapping
|
|
|
|
|
|
def _job_period(db, *, tenant_id: int, client_id: int, job_id: str, account_number: str = ""):
|
|
stmt = select(AccountingBankTransaction.transaction_date).where(
|
|
AccountingBankTransaction.tenant_id == int(tenant_id),
|
|
AccountingBankTransaction.client_id == int(client_id),
|
|
AccountingBankTransaction.source_job_id == job_id,
|
|
)
|
|
if _s(account_number):
|
|
stmt = stmt.where(AccountingBankTransaction.account_number == _s(account_number))
|
|
values = list(
|
|
db.execute(
|
|
stmt.order_by(AccountingBankTransaction.transaction_date)
|
|
).scalars().all()
|
|
)
|
|
values = [v for v in values if _date_obj(v)]
|
|
if not values:
|
|
raise ValueError(
|
|
"No imported bank transactions were found for this Bank Analyzer job/account. "
|
|
"Use a Bank Reconciliation purpose job or import the completed job into Accounting first."
|
|
)
|
|
return values[0][:10], values[-1][:10]
|
|
|
|
|
|
def queue_reconciliation(
|
|
db,
|
|
*,
|
|
tenant_id: int,
|
|
client_id: int,
|
|
source_job_id: str,
|
|
tally_guid: str,
|
|
company_name: str,
|
|
bank_ledger_name: str,
|
|
workstation_id: int,
|
|
user_id: int,
|
|
account_number: str = "",
|
|
date_tolerance_days: int = 15,
|
|
):
|
|
source_job = db.get(BankStatementAnalysisJob, source_job_id)
|
|
if (
|
|
not source_job
|
|
or source_job.status != "completed"
|
|
or int(source_job.client_id or 0) != int(client_id)
|
|
or int(source_job.tenant_id or 0) != int(tenant_id)
|
|
or source_job.ownership_status != "confirmed"
|
|
):
|
|
raise ValueError("The selected completed Bank Analyzer job is not valid for this client.")
|
|
|
|
if source_job.accounting_import_status != "completed":
|
|
raise ValueError(
|
|
"The selected Bank Analyzer job has not been imported successfully into the Accounting bank queue."
|
|
)
|
|
|
|
account_number = _s(account_number)
|
|
date_tolerance_days = max(0, min(90, int(date_tolerance_days or 15)))
|
|
|
|
date_from, date_to = _job_period(
|
|
db,
|
|
tenant_id=tenant_id,
|
|
client_id=client_id,
|
|
job_id=source_job_id,
|
|
account_number=account_number,
|
|
)
|
|
|
|
allowed = {
|
|
row.name
|
|
for row in visible_bank_ledgers(
|
|
db,
|
|
tenant_id=tenant_id,
|
|
client_id=client_id,
|
|
tally_guid=tally_guid,
|
|
)
|
|
}
|
|
if bank_ledger_name not in allowed:
|
|
raise ValueError("Select a Bank ledger from the client's synchronized Chart of Accounts.")
|
|
|
|
ws = db.get(ERPWorkstationAgent, int(workstation_id))
|
|
if (
|
|
not ws
|
|
or not ws.is_active
|
|
or not ws.tally_connected
|
|
or int(ws.tenant_id) != int(tenant_id)
|
|
):
|
|
raise ValueError("Selected workstation is unavailable or Tally is not connected.")
|
|
|
|
run = BankReconciliationRun(
|
|
tenant_id=int(tenant_id),
|
|
client_id=int(client_id),
|
|
source_job_id=source_job_id,
|
|
tally_guid=_s(tally_guid),
|
|
company_name=_s(company_name),
|
|
bank_ledger_name=bank_ledger_name,
|
|
account_number=account_number,
|
|
date_from=date_from,
|
|
date_to=date_to,
|
|
date_tolerance_days=date_tolerance_days,
|
|
workstation_agent_id=ws.id,
|
|
status="queued",
|
|
created_by_user_id=int(user_id),
|
|
)
|
|
db.add(run)
|
|
db.commit()
|
|
db.refresh(run)
|
|
|
|
job = enqueue_agent_job(
|
|
db,
|
|
workstation_agent_id=ws.id,
|
|
action=ACTION,
|
|
payload={
|
|
"tenant_id": int(tenant_id),
|
|
"client_id": int(client_id),
|
|
"tally_guid": _s(tally_guid),
|
|
"company_name": _s(company_name),
|
|
"bank_ledger_name": bank_ledger_name,
|
|
"account_number": account_number,
|
|
"date_from": date_from,
|
|
"date_to": date_to,
|
|
"reconciliation_run_id": run.id,
|
|
},
|
|
idempotency_key=f"bank-recon:{tenant_id}:{client_id}:{source_job_id}:{account_number}:{tally_guid}:{bank_ledger_name}:{run.id}",
|
|
priority=8,
|
|
max_attempts=2,
|
|
created_by_user_id=user_id,
|
|
)
|
|
run.agent_job_id = job.id
|
|
db.add(run)
|
|
db.commit()
|
|
return run
|
|
|
|
|
|
def queue_all_mapped_reconciliations(
|
|
db,
|
|
*,
|
|
tenant_id: int,
|
|
client_id: int,
|
|
source_job_id: str,
|
|
tally_guid: str,
|
|
company_name: str,
|
|
workstation_id: int,
|
|
user_id: int,
|
|
date_tolerance_days: int = 15,
|
|
):
|
|
accounts = source_accounts(
|
|
db,
|
|
tenant_id=tenant_id,
|
|
client_id=client_id,
|
|
job_id=source_job_id,
|
|
)
|
|
if not accounts:
|
|
raise ValueError("No bank accounts were imported from the selected Bank Analyzer job.")
|
|
|
|
mappings = {
|
|
row.account_number: row
|
|
for row in list_bank_mappings(
|
|
db,
|
|
tenant_id=tenant_id,
|
|
client_id=client_id,
|
|
tally_guid=tally_guid,
|
|
)
|
|
}
|
|
|
|
missing = [
|
|
row
|
|
for row in accounts
|
|
if row["account_number"] not in mappings
|
|
]
|
|
if missing:
|
|
labels = ", ".join(
|
|
f"{row['bank_name']} {row['account_number']}"
|
|
for row in missing
|
|
)
|
|
raise ValueError(
|
|
"Save Tally bank-ledger mapping for every bank account before "
|
|
f"running all accounts together. Missing: {labels}"
|
|
)
|
|
|
|
runs = []
|
|
for account in accounts:
|
|
mapping = mappings[account["account_number"]]
|
|
runs.append(
|
|
queue_reconciliation(
|
|
db,
|
|
tenant_id=tenant_id,
|
|
client_id=client_id,
|
|
source_job_id=source_job_id,
|
|
tally_guid=tally_guid,
|
|
company_name=company_name,
|
|
bank_ledger_name=mapping.bank_ledger_name,
|
|
workstation_id=workstation_id,
|
|
user_id=user_id,
|
|
account_number=account["account_number"],
|
|
date_tolerance_days=date_tolerance_days,
|
|
)
|
|
)
|
|
return runs
|
|
|
|
|
|
|
|
def _tally_side(voucher: dict, bank_ledger_name: str):
|
|
bank_key = bank_ledger_name.casefold()
|
|
entries = list(voucher.get("ledger_entries") or [])
|
|
bank_entries = [
|
|
row for row in entries
|
|
if _s(row.get("ledger_name")).casefold() == bank_key
|
|
]
|
|
if not bank_entries:
|
|
return None
|
|
|
|
# Tally exports in this project use negative amount for a Debit ledger and
|
|
# positive amount for Credit. Bank statement direction is opposite from the
|
|
# bank ledger accounting side: bank Debit statement = bank ledger Credit.
|
|
amount = round(abs(float(bank_entries[0].get("amount") or 0)), 2)
|
|
signed = float(bank_entries[0].get("amount") or 0)
|
|
direction = "DEBIT" if signed > 0 else "CREDIT"
|
|
|
|
return {
|
|
"guid": _s(voucher.get("guid")),
|
|
"voucher_number": _s(voucher.get("voucher_number")),
|
|
"voucher_type": _s(voucher.get("voucher_type_name")),
|
|
"date": _s(voucher.get("date"))[:10],
|
|
"reference": _s(voucher.get("reference")),
|
|
"narration": _s(voucher.get("narration")),
|
|
"amount": amount,
|
|
"direction": direction,
|
|
}
|
|
|
|
|
|
def _pair_score(bank: AccountingBankTransaction, tally: dict, tolerance_days: int = 15):
|
|
if round(float(bank.amount or 0), 2) != round(float(tally["amount"] or 0), 2):
|
|
return 0, ""
|
|
if _s(bank.direction).upper() != _s(tally["direction"]).upper():
|
|
return 0, ""
|
|
|
|
bd = _date_obj(bank.transaction_date)
|
|
td = _date_obj(tally["date"])
|
|
if not bd or not td:
|
|
return 0, ""
|
|
|
|
gap = abs((bd - td).days)
|
|
tolerance_days = max(0, min(90, int(tolerance_days or 15)))
|
|
if gap > tolerance_days:
|
|
return 0, ""
|
|
|
|
bank_ref = _norm_ref(bank.transfer_reference or bank.reference_no)
|
|
tally_ref = _norm_ref(tally.get("reference"))
|
|
ref_exact = bool(bank_ref and tally_ref and bank_ref == tally_ref)
|
|
|
|
bank_text = _norm_text(f"{bank.auto_party} {bank.narration}")
|
|
tally_text = _norm_text(tally.get("narration"))
|
|
similarity = int(round(SequenceMatcher(None, bank_text, tally_text).ratio() * 100)) if bank_text and tally_text else 0
|
|
|
|
score = 60
|
|
reasons = ["same amount", "same bank direction"]
|
|
|
|
if gap == 0:
|
|
score += 22
|
|
reasons.append("same date")
|
|
elif gap <= 2:
|
|
score += 12
|
|
reasons.append(f"{gap}-day timing difference")
|
|
elif gap <= 4:
|
|
score += 7
|
|
reasons.append(f"{gap}-day timing difference")
|
|
elif gap <= 7:
|
|
score += 2
|
|
reasons.append(f"{gap}-day timing difference")
|
|
else:
|
|
# Long timing differences remain match candidates only because amount
|
|
# and bank direction are exact. They are surfaced separately rather
|
|
# than silently treated as a normal probable match.
|
|
reasons.append(f"{gap}-day timing difference")
|
|
|
|
if ref_exact:
|
|
score += 18
|
|
reasons.append("same reference")
|
|
elif similarity >= 85:
|
|
score += 10
|
|
reasons.append("high narration similarity")
|
|
elif similarity >= 65:
|
|
score += 5
|
|
reasons.append("narration similarity")
|
|
|
|
return min(100, score), ", ".join(reasons)
|
|
|
|
|
|
def build_reconciliation(db, run: BankReconciliationRun, vouchers: list[dict]):
|
|
db.execute(delete(BankReconciliationItem).where(BankReconciliationItem.run_id == run.id))
|
|
db.commit()
|
|
|
|
bank_rows = list(
|
|
db.execute(
|
|
select(AccountingBankTransaction)
|
|
.where(
|
|
AccountingBankTransaction.tenant_id == run.tenant_id,
|
|
AccountingBankTransaction.client_id == run.client_id,
|
|
AccountingBankTransaction.source_job_id == run.source_job_id,
|
|
*(
|
|
[AccountingBankTransaction.account_number == run.account_number]
|
|
if _s(run.account_number)
|
|
else []
|
|
),
|
|
)
|
|
.order_by(
|
|
AccountingBankTransaction.transaction_date,
|
|
AccountingBankTransaction.id,
|
|
)
|
|
).scalars().all()
|
|
)
|
|
|
|
tally_rows = []
|
|
for voucher in vouchers:
|
|
if _s(voucher.get("is_cancelled")).lower() in {"yes", "true", "1"}:
|
|
continue
|
|
row = _tally_side(voucher, run.bank_ledger_name)
|
|
if row and row["amount"] > 0:
|
|
tally_rows.append(row)
|
|
|
|
candidates = {}
|
|
for bank in bank_rows:
|
|
scored = []
|
|
for index, tally in enumerate(tally_rows):
|
|
score, reason = _pair_score(bank, tally, run.date_tolerance_days)
|
|
if score:
|
|
scored.append((score, index, reason))
|
|
scored.sort(key=lambda item: (-item[0], item[1]))
|
|
candidates[bank.id] = scored
|
|
|
|
used_tally = set()
|
|
exact = probable = timing = bank_only = duplicates = 0
|
|
|
|
for bank in bank_rows:
|
|
options = [
|
|
item for item in candidates.get(bank.id, [])
|
|
if item[1] not in used_tally
|
|
]
|
|
|
|
if not options:
|
|
status = "bank_only"
|
|
confidence = 0
|
|
reason = "No Tally bank-ledger voucher matched amount, direction and permitted date window."
|
|
tally = None
|
|
bank_only += 1
|
|
else:
|
|
top_score, tally_index, reason = options[0]
|
|
tied = [item for item in options if item[0] == top_score]
|
|
if len(tied) > 1 and top_score < 95:
|
|
status = "duplicate_candidate"
|
|
confidence = top_score
|
|
tally = tally_rows[tally_index]
|
|
duplicates += 1
|
|
else:
|
|
used_tally.add(tally_index)
|
|
tally = tally_rows[tally_index]
|
|
confidence = top_score
|
|
bank_date = _date_obj(bank.transaction_date)
|
|
tally_date = _date_obj(tally["date"])
|
|
gap = abs((bank_date - tally_date).days) if bank_date and tally_date else 999
|
|
if top_score >= 90:
|
|
status = "matched"
|
|
exact += 1
|
|
elif gap > 7:
|
|
status = "timing_difference"
|
|
timing += 1
|
|
else:
|
|
status = "probable_match"
|
|
probable += 1
|
|
|
|
item = BankReconciliationItem(
|
|
run_id=run.id,
|
|
bank_transaction_id=bank.id,
|
|
match_status=status,
|
|
confidence=int(confidence),
|
|
match_reason=reason,
|
|
bank_date=_s(bank.transaction_date)[:10],
|
|
bank_direction=_s(bank.direction).upper(),
|
|
bank_amount=float(bank.amount or 0),
|
|
bank_reference=_s(bank.transfer_reference or bank.reference_no),
|
|
bank_narration=_s(bank.narration),
|
|
)
|
|
if tally:
|
|
item.tally_guid = tally["guid"]
|
|
item.tally_voucher_number = tally["voucher_number"]
|
|
item.tally_voucher_type = tally["voucher_type"]
|
|
item.tally_date = tally["date"]
|
|
item.tally_reference = tally["reference"]
|
|
item.tally_narration = tally["narration"]
|
|
item.tally_amount = tally["amount"]
|
|
item.tally_direction = tally["direction"]
|
|
|
|
bank.reconciliation_status = status
|
|
bank.last_reconciliation_run_id = run.id
|
|
db.add(bank)
|
|
db.add(item)
|
|
|
|
books_only = 0
|
|
for index, tally in enumerate(tally_rows):
|
|
if index in used_tally:
|
|
continue
|
|
books_only += 1
|
|
db.add(
|
|
BankReconciliationItem(
|
|
run_id=run.id,
|
|
bank_transaction_id=None,
|
|
match_status="books_only",
|
|
confidence=0,
|
|
match_reason="Tally bank-ledger voucher has no matching transaction in the uploaded bank statement set.",
|
|
tally_guid=tally["guid"],
|
|
tally_voucher_number=tally["voucher_number"],
|
|
tally_voucher_type=tally["voucher_type"],
|
|
tally_date=tally["date"],
|
|
tally_reference=tally["reference"],
|
|
tally_narration=tally["narration"],
|
|
tally_amount=tally["amount"],
|
|
tally_direction=tally["direction"],
|
|
)
|
|
)
|
|
|
|
bank_only_rows = [
|
|
row for row in db.execute(
|
|
select(BankReconciliationItem).where(
|
|
BankReconciliationItem.run_id == run.id,
|
|
BankReconciliationItem.match_status == "bank_only",
|
|
)
|
|
).scalars().all()
|
|
]
|
|
books_only_rows = [
|
|
row for row in db.execute(
|
|
select(BankReconciliationItem).where(
|
|
BankReconciliationItem.run_id == run.id,
|
|
BankReconciliationItem.match_status == "books_only",
|
|
)
|
|
).scalars().all()
|
|
]
|
|
|
|
def _direction_totals(rows, prefix):
|
|
result = {
|
|
f"{prefix}_debit_amount": 0.0,
|
|
f"{prefix}_credit_amount": 0.0,
|
|
}
|
|
for row in rows:
|
|
direction = (
|
|
row.bank_direction if prefix == "bank_only"
|
|
else row.tally_direction
|
|
)
|
|
amount = (
|
|
row.bank_amount if prefix == "bank_only"
|
|
else row.tally_amount
|
|
)
|
|
key = (
|
|
f"{prefix}_debit_amount"
|
|
if _s(direction).upper() == "DEBIT"
|
|
else f"{prefix}_credit_amount"
|
|
)
|
|
result[key] = round(result[key] + float(amount or 0), 2)
|
|
return result
|
|
|
|
summary = {
|
|
"bank_transactions": len(bank_rows),
|
|
"tally_bank_vouchers": len(tally_rows),
|
|
"matched": exact,
|
|
"probable_match": probable,
|
|
"timing_difference": timing,
|
|
"bank_only": bank_only,
|
|
"books_only": books_only,
|
|
"duplicate_candidate": duplicates,
|
|
"account_number": run.account_number,
|
|
"bank_ledger_name": run.bank_ledger_name,
|
|
"date_tolerance_days": int(run.date_tolerance_days or 15),
|
|
}
|
|
summary.update(_direction_totals(bank_only_rows, "bank_only"))
|
|
summary.update(_direction_totals(books_only_rows, "books_only"))
|
|
summary["bank_only_net"] = round(
|
|
summary["bank_only_credit_amount"] - summary["bank_only_debit_amount"],
|
|
2,
|
|
)
|
|
summary["books_only_net"] = round(
|
|
summary["books_only_credit_amount"] - summary["books_only_debit_amount"],
|
|
2,
|
|
)
|
|
summary["unreconciled_net_difference"] = round(
|
|
summary["bank_only_net"] - summary["books_only_net"],
|
|
2,
|
|
)
|
|
|
|
run.summary_json = json.dumps(summary, ensure_ascii=False)
|
|
run.status = "completed"
|
|
run.completed_at_utc = _utcnow()
|
|
run.last_error = ""
|
|
db.add(run)
|
|
db.commit()
|
|
return run
|
|
|
|
|
|
def sync_run(db, run: BankReconciliationRun):
|
|
if not run.agent_job_id or run.status == "completed":
|
|
return run
|
|
job = db.get(ERPAgentJob, int(run.agent_job_id))
|
|
if not job:
|
|
return run
|
|
|
|
if job.status in {"queued", "claimed"}:
|
|
next_status = "extracting" if job.status == "claimed" else "queued"
|
|
if run.status != next_status:
|
|
run.status = next_status
|
|
db.add(run)
|
|
db.commit()
|
|
return run
|
|
|
|
if job.status == "succeeded":
|
|
result = _loads(job.result_json, {})
|
|
vouchers = list(result.get("vouchers") or [])
|
|
return build_reconciliation(db, run, vouchers)
|
|
|
|
if job.status in {"failed", "cancelled"}:
|
|
run.status = "failed"
|
|
run.last_error = _s(job.last_error) or f"Local Agent reconciliation extraction {job.status}."
|
|
run.completed_at_utc = _utcnow()
|
|
db.add(run)
|
|
db.commit()
|
|
|
|
return run
|
|
|
|
|
|
def resolve_reconciliation_item(
|
|
db,
|
|
*,
|
|
run_id: int,
|
|
item_id: int,
|
|
action: str,
|
|
note: str,
|
|
user_id: int,
|
|
):
|
|
item = db.execute(
|
|
select(BankReconciliationItem).where(
|
|
BankReconciliationItem.id == int(item_id),
|
|
BankReconciliationItem.run_id == int(run_id),
|
|
)
|
|
).scalar_one_or_none()
|
|
if not item:
|
|
raise ValueError("Reconciliation item was not found.")
|
|
|
|
action = _s(action)
|
|
allowed = {
|
|
"confirm_match",
|
|
"confirm_timing",
|
|
"reject_match_bank_only",
|
|
"confirm_bank_only",
|
|
"confirm_books_only",
|
|
"needs_follow_up",
|
|
"reopen",
|
|
}
|
|
if action not in allowed:
|
|
raise ValueError("Unsupported reconciliation resolution.")
|
|
|
|
item.resolution_status = action
|
|
item.resolution_note = _s(note)
|
|
item.resolved_by_user_id = int(user_id)
|
|
item.resolved_at_utc = _utcnow()
|
|
|
|
if item.bank_transaction_id:
|
|
tx = db.get(AccountingBankTransaction, int(item.bank_transaction_id))
|
|
if tx:
|
|
if action == "confirm_match":
|
|
tx.reconciliation_status = "matched"
|
|
elif action == "confirm_timing":
|
|
tx.reconciliation_status = "timing_difference_confirmed"
|
|
elif action in {"reject_match_bank_only", "confirm_bank_only"}:
|
|
tx.reconciliation_status = "bank_only"
|
|
elif action == "needs_follow_up":
|
|
tx.reconciliation_status = "needs_review"
|
|
elif action == "reopen":
|
|
tx.reconciliation_status = item.match_status
|
|
tx.last_reconciliation_run_id = int(run_id)
|
|
db.add(tx)
|
|
|
|
db.add(item)
|
|
db.commit()
|
|
return item
|
|
|
|
|
|
|
|
def list_runs(db, *, tenant_id: int, client_id: int, limit: int = 30):
|
|
rows = list(
|
|
db.execute(
|
|
select(BankReconciliationRun)
|
|
.where(
|
|
BankReconciliationRun.tenant_id == int(tenant_id),
|
|
BankReconciliationRun.client_id == int(client_id),
|
|
)
|
|
.order_by(BankReconciliationRun.id.desc())
|
|
.limit(limit)
|
|
).scalars().all()
|
|
)
|
|
for row in rows:
|
|
sync_run(db, row)
|
|
return rows
|
|
|
|
|
|
def run_items(db, *, run_id: int, status: str = ""):
|
|
stmt = select(BankReconciliationItem).where(
|
|
BankReconciliationItem.run_id == int(run_id)
|
|
)
|
|
if _s(status):
|
|
stmt = stmt.where(BankReconciliationItem.match_status == _s(status))
|
|
return list(
|
|
db.execute(
|
|
stmt.order_by(
|
|
BankReconciliationItem.bank_date,
|
|
BankReconciliationItem.tally_date,
|
|
BankReconciliationItem.id,
|
|
)
|
|
).scalars().all()
|
|
)
|