diff --git a/alembic/versions/20260918_accounting_mirror_versioning.py b/alembic/versions/20260918_accounting_mirror_versioning.py new file mode 100644 index 0000000..1a53325 --- /dev/null +++ b/alembic/versions/20260918_accounting_mirror_versioning.py @@ -0,0 +1,72 @@ +"""Version Accounting Mirrors per Client + Financial Year. + +Revision ID: 20260918_accounting_mirror_ver +Revises: 20260918_accounting_mirror_reg +""" +from alembic import op +import sqlalchemy as sa + +revision = "20260918_accounting_mirror_ver" +down_revision = "20260918_accounting_mirror_reg" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + bind = op.get_bind() + inspector = sa.inspect(bind) + columns = {c["name"] for c in inspector.get_columns("accounting_mirror_registry")} + + if "version_no" not in columns: + op.add_column("accounting_mirror_registry", sa.Column("version_no", sa.Integer(), nullable=False, server_default="1")) + if "is_current" not in columns: + op.add_column("accounting_mirror_registry", sa.Column("is_current", sa.Boolean(), nullable=False, server_default=sa.true())) + if "supersedes_mirror_id" not in columns: + op.add_column("accounting_mirror_registry", sa.Column("supersedes_mirror_id", sa.Integer(), nullable=True)) + op.create_foreign_key( + "fk_accounting_mirror_registry_supersedes", + "accounting_mirror_registry", + "accounting_mirror_registry", + ["supersedes_mirror_id"], + ["id"], + ondelete="SET NULL", + ) + + uniques = {u.get("name") for u in inspector.get_unique_constraints("accounting_mirror_registry")} + if "uq_accounting_mirror_registry_client_fy" in uniques: + op.drop_constraint("uq_accounting_mirror_registry_client_fy", "accounting_mirror_registry", type_="unique") + if "uq_accounting_mirror_registry_client_fy_version" not in uniques: + op.create_unique_constraint( + "uq_accounting_mirror_registry_client_fy_version", + "accounting_mirror_registry", + ["tenant_id", "client_id", "financial_year", "version_no"], + ) + + indexes = {i.get("name") for i in inspector.get_indexes("accounting_mirror_registry")} + if "ix_accounting_mirror_registry_version_no" not in indexes: + op.create_index("ix_accounting_mirror_registry_version_no", "accounting_mirror_registry", ["version_no"], unique=False) + if "ix_accounting_mirror_registry_is_current" not in indexes: + op.create_index("ix_accounting_mirror_registry_is_current", "accounting_mirror_registry", ["is_current"], unique=False) + if "ix_accounting_mirror_registry_supersedes_mirror_id" not in indexes: + op.create_index("ix_accounting_mirror_registry_supersedes_mirror_id", "accounting_mirror_registry", ["supersedes_mirror_id"], unique=False) + + # Existing registry rows become version 1/current without touching their + # file path, company metadata, dates or status. + op.execute(sa.text("UPDATE accounting_mirror_registry SET version_no = 1 WHERE version_no IS NULL")) + op.execute(sa.text("UPDATE accounting_mirror_registry SET is_current = TRUE WHERE is_current IS NULL")) + + +def downgrade() -> None: + op.drop_index("ix_accounting_mirror_registry_supersedes_mirror_id", table_name="accounting_mirror_registry") + op.drop_index("ix_accounting_mirror_registry_is_current", table_name="accounting_mirror_registry") + op.drop_index("ix_accounting_mirror_registry_version_no", table_name="accounting_mirror_registry") + op.drop_constraint("uq_accounting_mirror_registry_client_fy_version", "accounting_mirror_registry", type_="unique") + op.create_unique_constraint( + "uq_accounting_mirror_registry_client_fy", + "accounting_mirror_registry", + ["tenant_id", "client_id", "financial_year"], + ) + op.drop_constraint("fk_accounting_mirror_registry_supersedes", "accounting_mirror_registry", type_="foreignkey") + op.drop_column("accounting_mirror_registry", "supersedes_mirror_id") + op.drop_column("accounting_mirror_registry", "is_current") + op.drop_column("accounting_mirror_registry", "version_no") diff --git a/app/modules/accounting/accounting_mirror_models.py b/app/modules/accounting/accounting_mirror_models.py index 14cd01c..0b83c15 100644 --- a/app/modules/accounting/accounting_mirror_models.py +++ b/app/modules/accounting/accounting_mirror_models.py @@ -21,7 +21,8 @@ class AccountingMirrorRegistry(CommonBase): "tenant_id", "client_id", "financial_year", - name="uq_accounting_mirror_registry_client_fy", + "version_no", + name="uq_accounting_mirror_registry_client_fy_version", ), ) @@ -40,9 +41,14 @@ class AccountingMirrorRegistry(CommonBase): voucher_to_date: Mapped[date | None] = mapped_column(Date, nullable=True) file_size_bytes: Mapped[int] = mapped_column(BigInteger, nullable=False, default=0) + version_no: Mapped[int] = mapped_column(Integer, nullable=False, default=1, index=True) + is_current: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True, index=True) status: Mapped[str] = mapped_column(String(30), nullable=False, default="active", index=True) is_active: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True, index=True) replacement_count: Mapped[int] = mapped_column(Integer, nullable=False, default=0) + supersedes_mirror_id: Mapped[int | None] = mapped_column( + ForeignKey("accounting_mirror_registry.id", ondelete="SET NULL"), nullable=True, index=True + ) created_by_user_id: Mapped[int | None] = mapped_column(ForeignKey("users.id", ondelete="SET NULL"), nullable=True) last_synced_by_user_id: Mapped[int | None] = mapped_column(ForeignKey("users.id", ondelete="SET NULL"), nullable=True) diff --git a/app/modules/accounting/accounting_mirror_registry.py b/app/modules/accounting/accounting_mirror_registry.py deleted file mode 100644 index b389a3b..0000000 --- a/app/modules/accounting/accounting_mirror_registry.py +++ /dev/null @@ -1,181 +0,0 @@ -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 diff --git a/app/modules/accounting/accounting_mirror_service.py b/app/modules/accounting/accounting_mirror_service.py index 685bbd8..fd6b7ea 100644 --- a/app/modules/accounting/accounting_mirror_service.py +++ b/app/modules/accounting/accounting_mirror_service.py @@ -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 diff --git a/app/modules/accounting/creditors_aging_ui.py b/app/modules/accounting/creditors_aging_ui.py index 115d60d..ecbfa31 100644 --- a/app/modules/accounting/creditors_aging_ui.py +++ b/app/modules/accounting/creditors_aging_ui.py @@ -13,7 +13,7 @@ from app.core.db.common import CommonSessionLocal from app.core.security.csrf import get_or_create_csrf_token from app.core.templating import templates from app.modules.accounting.accounting_mirror_service import AccountingMirrorError, sundry_creditors_aging -from app.modules.accounting.accounting_mirror_registry import ( +from app.modules.accounting.accounting_mirror_service import ( get_registered_mirror, list_registered_mirrors, sync_discovered_mirrors, diff --git a/app/modules/accounting/templates/accounting/tally.html b/app/modules/accounting/templates/accounting/tally.html index bd8e11d..4593b23 100644 --- a/app/modules/accounting/templates/accounting/tally.html +++ b/app/modules/accounting/templates/accounting/tally.html @@ -292,7 +292,7 @@