Add versioned Client FY Accounting Mirrors
This commit is contained in:
@@ -1,8 +1,14 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
from datetime import date, datetime, timezone
|
||||
from pathlib import Path
|
||||
import re
|
||||
|
||||
from sqlalchemy import select, func
|
||||
|
||||
from app.modules.accounting.agent_bridge import request_agent_command
|
||||
from app.modules.accounting.accounting_mirror_models import AccountingMirrorRegistry
|
||||
|
||||
|
||||
DEFAULT_QUERY_TIMEOUT_SECONDS = 45
|
||||
@@ -728,3 +734,299 @@ def sundry_creditors_aging(
|
||||
"current_mirror": current.get("mirror") or {},
|
||||
"follow_up_mirror": later.get("mirror") or {},
|
||||
}
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Accounting Mirror metadata/version helpers
|
||||
# Kept inside the existing Accounting module/service so every downstream tool
|
||||
# resolves Client + FY mirrors through one source of truth.
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def _mirror_utcnow() -> datetime:
|
||||
return datetime.now(timezone.utc)
|
||||
|
||||
|
||||
def _mirror_fy_key(value: str) -> int:
|
||||
match = re.fullmatch(r"(\d{4})-(\d{2})", str(value or "").strip())
|
||||
return int(match.group(1)) if match else -1
|
||||
|
||||
|
||||
def _parse_mirror_date(value: Any) -> date | None:
|
||||
text = str(value or "").strip()
|
||||
if not text:
|
||||
return None
|
||||
try:
|
||||
return date.fromisoformat(text[:10])
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
|
||||
def _parse_mirror_datetime(value: Any) -> datetime | None:
|
||||
text = str(value or "").strip()
|
||||
if not text:
|
||||
return None
|
||||
try:
|
||||
parsed = datetime.fromisoformat(text.replace("Z", "+00:00"))
|
||||
return parsed if parsed.tzinfo else parsed.replace(tzinfo=timezone.utc)
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
|
||||
def get_current_accounting_mirror(db, tenant_id: int, client_id: int, financial_year: str) -> AccountingMirrorRegistry | None:
|
||||
return db.execute(
|
||||
select(AccountingMirrorRegistry).where(
|
||||
AccountingMirrorRegistry.tenant_id == int(tenant_id),
|
||||
AccountingMirrorRegistry.client_id == int(client_id),
|
||||
AccountingMirrorRegistry.financial_year == str(financial_year or "").strip(),
|
||||
AccountingMirrorRegistry.is_current.is_(True),
|
||||
AccountingMirrorRegistry.is_active.is_(True),
|
||||
AccountingMirrorRegistry.status == "active",
|
||||
)
|
||||
).scalar_one_or_none()
|
||||
|
||||
|
||||
def get_registered_mirror(db, tenant_id: int, client_id: int, financial_year: str) -> AccountingMirrorRegistry | None:
|
||||
"""Backward-compatible name used by existing Accounting tools."""
|
||||
return get_current_accounting_mirror(db, tenant_id, client_id, financial_year)
|
||||
|
||||
|
||||
def list_accounting_mirror_versions(db, tenant_id: int, client_id: int, financial_year: str) -> list[AccountingMirrorRegistry]:
|
||||
rows = db.execute(
|
||||
select(AccountingMirrorRegistry).where(
|
||||
AccountingMirrorRegistry.tenant_id == int(tenant_id),
|
||||
AccountingMirrorRegistry.client_id == int(client_id),
|
||||
AccountingMirrorRegistry.financial_year == str(financial_year or "").strip(),
|
||||
).order_by(AccountingMirrorRegistry.version_no.desc(), AccountingMirrorRegistry.id.desc())
|
||||
).scalars().all()
|
||||
return list(rows)
|
||||
|
||||
|
||||
def list_registered_mirrors(db, tenant_id: int, client_id: int) -> list[AccountingMirrorRegistry]:
|
||||
rows = db.execute(
|
||||
select(AccountingMirrorRegistry).where(
|
||||
AccountingMirrorRegistry.tenant_id == int(tenant_id),
|
||||
AccountingMirrorRegistry.client_id == int(client_id),
|
||||
AccountingMirrorRegistry.is_current.is_(True),
|
||||
AccountingMirrorRegistry.is_active.is_(True),
|
||||
AccountingMirrorRegistry.status == "active",
|
||||
)
|
||||
).scalars().all()
|
||||
return sorted(rows, key=lambda row: _mirror_fy_key(row.financial_year), reverse=True)
|
||||
|
||||
|
||||
def next_accounting_mirror_version(db, tenant_id: int, client_id: int, financial_year: str) -> int:
|
||||
value = db.execute(
|
||||
select(func.max(AccountingMirrorRegistry.version_no)).where(
|
||||
AccountingMirrorRegistry.tenant_id == int(tenant_id),
|
||||
AccountingMirrorRegistry.client_id == int(client_id),
|
||||
AccountingMirrorRegistry.financial_year == str(financial_year or "").strip(),
|
||||
)
|
||||
).scalar_one_or_none()
|
||||
return int(value or 0) + 1
|
||||
|
||||
|
||||
def register_accounting_mirror_version(
|
||||
db,
|
||||
*,
|
||||
tenant_id: int,
|
||||
client_id: int,
|
||||
financial_year: str,
|
||||
accounting_relative_dir: str,
|
||||
storage_node_id: int | None,
|
||||
mirror: dict[str, Any] | None = None,
|
||||
job: dict[str, Any] | None = None,
|
||||
requested_by_user_id: int | None = None,
|
||||
version_no: int | None = None,
|
||||
archived_previous_file_name: str = "",
|
||||
) -> AccountingMirrorRegistry:
|
||||
mirror = dict(mirror or {})
|
||||
job = dict(job or {})
|
||||
company = dict(mirror.get("company") or {})
|
||||
voucher_period = dict(mirror.get("voucher_period") or {})
|
||||
fy = str(financial_year or "").strip()
|
||||
|
||||
current = get_current_accounting_mirror(db, tenant_id, client_id, fy)
|
||||
target_version = int(version_no or next_accounting_mirror_version(db, tenant_id, client_id, fy))
|
||||
|
||||
# Idempotent status polling: if this version was already registered, update it
|
||||
# instead of creating another row.
|
||||
row = db.execute(
|
||||
select(AccountingMirrorRegistry).where(
|
||||
AccountingMirrorRegistry.tenant_id == int(tenant_id),
|
||||
AccountingMirrorRegistry.client_id == int(client_id),
|
||||
AccountingMirrorRegistry.financial_year == fy,
|
||||
AccountingMirrorRegistry.version_no == target_version,
|
||||
)
|
||||
).scalar_one_or_none()
|
||||
|
||||
if row is None:
|
||||
row = AccountingMirrorRegistry(
|
||||
tenant_id=int(tenant_id),
|
||||
client_id=int(client_id),
|
||||
financial_year=fy,
|
||||
version_no=target_version,
|
||||
is_current=True,
|
||||
is_active=True,
|
||||
status="active",
|
||||
accounting_relative_dir=str(accounting_relative_dir or "").strip(),
|
||||
mirror_file_name=Path(str(job.get("accounting_db_path") or f"client_{int(client_id):08d}.act")).name,
|
||||
created_by_user_id=requested_by_user_id,
|
||||
supersedes_mirror_id=(int(current.id) if current else None),
|
||||
)
|
||||
db.add(row)
|
||||
db.flush()
|
||||
|
||||
if current and int(current.id) != int(row.id):
|
||||
current.is_current = False
|
||||
current.status = "superseded"
|
||||
current.updated_at_utc = _mirror_utcnow()
|
||||
if archived_previous_file_name:
|
||||
current.mirror_file_name = Path(str(archived_previous_file_name)).name
|
||||
current.replacement_count = int(current.replacement_count or 0) + 1
|
||||
|
||||
# Ensure no other row remains current for the same Client/FY.
|
||||
others = db.execute(
|
||||
select(AccountingMirrorRegistry).where(
|
||||
AccountingMirrorRegistry.tenant_id == int(tenant_id),
|
||||
AccountingMirrorRegistry.client_id == int(client_id),
|
||||
AccountingMirrorRegistry.financial_year == fy,
|
||||
AccountingMirrorRegistry.id != int(row.id),
|
||||
AccountingMirrorRegistry.is_current.is_(True),
|
||||
)
|
||||
).scalars().all()
|
||||
for other in others:
|
||||
other.is_current = False
|
||||
if other.status == "active":
|
||||
other.status = "superseded"
|
||||
other.updated_at_utc = _mirror_utcnow()
|
||||
|
||||
row.storage_node_id = int(storage_node_id) if storage_node_id else None
|
||||
row.accounting_relative_dir = str(accounting_relative_dir or row.accounting_relative_dir or "").strip()
|
||||
row.mirror_file_name = Path(str(job.get("accounting_db_path") or row.mirror_file_name or f"client_{int(client_id):08d}.act")).name
|
||||
row.company_name = str(job.get("company_name") or company.get("company_name") or row.company_name or "").strip()
|
||||
row.company_guid = str(job.get("tally_guid") or company.get("company_guid") or row.company_guid or "").strip()
|
||||
row.voucher_from_date = _parse_mirror_date(voucher_period.get("from_date")) or row.voucher_from_date
|
||||
row.voucher_to_date = _parse_mirror_date(voucher_period.get("to_date")) or row.voucher_to_date
|
||||
try:
|
||||
row.file_size_bytes = int(mirror.get("size_bytes") or row.file_size_bytes or 0)
|
||||
except Exception:
|
||||
pass
|
||||
row.status = "active"
|
||||
row.is_active = True
|
||||
row.is_current = True
|
||||
row.last_synced_by_user_id = requested_by_user_id
|
||||
row.last_synced_at_utc = _parse_mirror_datetime(job.get("finished_at_utc")) or _parse_mirror_datetime(job.get("updated_at_utc")) or _mirror_utcnow()
|
||||
row.updated_at_utc = _mirror_utcnow()
|
||||
db.flush()
|
||||
return row
|
||||
|
||||
|
||||
def upsert_registered_mirror(
|
||||
db,
|
||||
*,
|
||||
tenant_id: int,
|
||||
client_id: int,
|
||||
financial_year: str,
|
||||
accounting_relative_dir: str,
|
||||
storage_node_id: int | None,
|
||||
mirror: dict[str, Any] | None = None,
|
||||
job: dict[str, Any] | None = None,
|
||||
requested_by_user_id: int | None = None,
|
||||
replacement: bool = False,
|
||||
) -> AccountingMirrorRegistry:
|
||||
"""Compatibility/backfill helper.
|
||||
|
||||
Discovery of a pre-versioning canonical mirror creates/updates v1 only when no
|
||||
current row exists. A successful explicit replacement is registered through
|
||||
register_accounting_mirror_version() with the job's requested version.
|
||||
"""
|
||||
fy = str(financial_year or "").strip()
|
||||
current = get_current_accounting_mirror(db, tenant_id, client_id, fy)
|
||||
if replacement:
|
||||
return register_accounting_mirror_version(
|
||||
db,
|
||||
tenant_id=tenant_id,
|
||||
client_id=client_id,
|
||||
financial_year=fy,
|
||||
accounting_relative_dir=accounting_relative_dir,
|
||||
storage_node_id=storage_node_id,
|
||||
mirror=mirror,
|
||||
job=job,
|
||||
requested_by_user_id=requested_by_user_id,
|
||||
version_no=int((job or {}).get("mirror_version_no") or 0) or None,
|
||||
archived_previous_file_name=str((job or {}).get("archived_previous_file_name") or ""),
|
||||
)
|
||||
if current:
|
||||
# Refresh metadata for the current physical canonical mirror without
|
||||
# manufacturing another version during discovery/status polling.
|
||||
current.storage_node_id = int(storage_node_id) if storage_node_id else None
|
||||
current.accounting_relative_dir = str(accounting_relative_dir or current.accounting_relative_dir or "").strip()
|
||||
current.company_name = str(((mirror or {}).get("company") or {}).get("company_name") or current.company_name or "").strip()
|
||||
current.company_guid = str(((mirror or {}).get("company") or {}).get("company_guid") or current.company_guid or "").strip()
|
||||
current.file_size_bytes = int((mirror or {}).get("size_bytes") or current.file_size_bytes or 0)
|
||||
current.updated_at_utc = _mirror_utcnow()
|
||||
db.flush()
|
||||
return current
|
||||
return register_accounting_mirror_version(
|
||||
db,
|
||||
tenant_id=tenant_id,
|
||||
client_id=client_id,
|
||||
financial_year=fy,
|
||||
accounting_relative_dir=accounting_relative_dir,
|
||||
storage_node_id=storage_node_id,
|
||||
mirror=mirror,
|
||||
job=job,
|
||||
requested_by_user_id=requested_by_user_id,
|
||||
version_no=1,
|
||||
)
|
||||
|
||||
|
||||
def sync_discovered_mirrors(
|
||||
db,
|
||||
*,
|
||||
tenant_id: int,
|
||||
client_id: int,
|
||||
storage_node_id: int | None,
|
||||
discovered: list[dict[str, Any]],
|
||||
requested_by_user_id: int | None = None,
|
||||
) -> list[AccountingMirrorRegistry]:
|
||||
output: list[AccountingMirrorRegistry] = []
|
||||
discovered_fys: set[str] = set()
|
||||
for item in discovered or []:
|
||||
fy = str(item.get("financial_year") or "").strip()
|
||||
relative_dir = str(item.get("accounting_relative_dir") or "").strip()
|
||||
mirror = dict(item.get("mirror") or {})
|
||||
if not re.fullmatch(r"\d{4}-\d{2}", fy) or not relative_dir or not mirror.get("ready"):
|
||||
continue
|
||||
discovered_fys.add(fy)
|
||||
row = upsert_registered_mirror(
|
||||
db,
|
||||
tenant_id=tenant_id,
|
||||
client_id=client_id,
|
||||
financial_year=fy,
|
||||
accounting_relative_dir=relative_dir,
|
||||
storage_node_id=storage_node_id,
|
||||
mirror=mirror,
|
||||
job=item,
|
||||
requested_by_user_id=requested_by_user_id,
|
||||
replacement=False,
|
||||
)
|
||||
output.append(row)
|
||||
|
||||
# Only current rows are availability indicators. Historical/superseded
|
||||
# versions are intentionally left untouched.
|
||||
current_rows = db.execute(
|
||||
select(AccountingMirrorRegistry).where(
|
||||
AccountingMirrorRegistry.tenant_id == int(tenant_id),
|
||||
AccountingMirrorRegistry.client_id == int(client_id),
|
||||
AccountingMirrorRegistry.is_current.is_(True),
|
||||
AccountingMirrorRegistry.is_active.is_(True),
|
||||
)
|
||||
).scalars().all()
|
||||
for row in current_rows:
|
||||
if row.financial_year not in discovered_fys:
|
||||
row.status = "missing"
|
||||
row.is_active = False
|
||||
row.updated_at_utc = _mirror_utcnow()
|
||||
db.commit()
|
||||
return output
|
||||
|
||||
Reference in New Issue
Block a user