182 lines
6.2 KiB
Python
182 lines
6.2 KiB
Python
from __future__ import annotations
|
|
|
|
from datetime import date, datetime, timezone
|
|
from pathlib import Path
|
|
import re
|
|
from typing import Any
|
|
|
|
from sqlalchemy import select
|
|
|
|
from app.modules.accounting.accounting_mirror_models import AccountingMirrorRegistry
|
|
|
|
|
|
def _utcnow() -> datetime:
|
|
return datetime.now(timezone.utc)
|
|
|
|
|
|
def _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 _get_any_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(),
|
|
)
|
|
).scalar_one_or_none()
|
|
|
|
|
|
def get_registered_mirror(db, tenant_id: int, client_id: int, financial_year: str) -> AccountingMirrorRegistry | None:
|
|
row = _get_any_mirror(db, tenant_id, client_id, financial_year)
|
|
return row if row and row.is_active and row.status == "active" else None
|
|
|
|
|
|
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_active.is_(True),
|
|
)
|
|
).scalars().all()
|
|
return sorted(rows, key=lambda row: _fy_key(row.financial_year), reverse=True)
|
|
|
|
|
|
def _parse_iso_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_iso_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 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:
|
|
mirror = dict(mirror or {})
|
|
job = dict(job or {})
|
|
company = dict(mirror.get("company") or {})
|
|
voucher_period = dict(mirror.get("voucher_period") or {})
|
|
|
|
row = _get_any_mirror(db, tenant_id, client_id, financial_year)
|
|
created = row is None
|
|
if row is None:
|
|
row = AccountingMirrorRegistry(
|
|
tenant_id=int(tenant_id),
|
|
client_id=int(client_id),
|
|
financial_year=str(financial_year).strip(),
|
|
accounting_relative_dir=str(accounting_relative_dir or "").strip(),
|
|
mirror_file_name=f"client_{int(client_id):08d}.act",
|
|
created_by_user_id=requested_by_user_id,
|
|
)
|
|
db.add(row)
|
|
|
|
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_iso_date(voucher_period.get("from_date")) or row.voucher_from_date
|
|
row.voucher_to_date = _parse_iso_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.last_synced_by_user_id = requested_by_user_id
|
|
row.last_synced_at_utc = (
|
|
_parse_iso_datetime(job.get("finished_at_utc"))
|
|
or _parse_iso_datetime(job.get("updated_at_utc"))
|
|
or _utcnow()
|
|
)
|
|
if replacement and not created:
|
|
row.replacement_count = int(row.replacement_count or 0) + 1
|
|
row.updated_at_utc = _utcnow()
|
|
db.flush()
|
|
return row
|
|
|
|
|
|
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)
|
|
|
|
existing = db.execute(
|
|
select(AccountingMirrorRegistry).where(
|
|
AccountingMirrorRegistry.tenant_id == int(tenant_id),
|
|
AccountingMirrorRegistry.client_id == int(client_id),
|
|
AccountingMirrorRegistry.is_active.is_(True),
|
|
)
|
|
).scalars().all()
|
|
for row in existing:
|
|
if row.financial_year not in discovered_fys:
|
|
row.status = "missing"
|
|
row.is_active = False
|
|
row.updated_at_utc = _utcnow()
|
|
db.commit()
|
|
return output
|