From 93ae6caa543d5f6df08ed3153a1aba9955fd871e Mon Sep 17 00:00:00 2001 From: A R R R Associates Date: Wed, 19 Aug 2026 16:11:03 +0530 Subject: [PATCH] Add Phase 3 Tally accounting master synchronization --- app/core/startup.py | 2 + app/modules/accounting/act_store.py | 44 ++- .../templates/accounting/tally.html | 52 ++- app/modules/accounting/ui.py | 54 ++- app/modules/core/rbac/permissions_registry.py | 2 + app/modules/documents/agent_package.py | 72 +--- .../README_ERP_LOCAL_AGENT.txt | 4 +- .../erp_local_agent/__init__.py | 2 +- .../erp_local_agent/accounting_store.py | 371 ++++++++++-------- .../erp_local_agent/commands.py | 120 ++---- .../erp_local_agent/tally.py | 200 +++++++++- 11 files changed, 586 insertions(+), 337 deletions(-) diff --git a/app/core/startup.py b/app/core/startup.py index f2ad1d2..695226e 100644 --- a/app/core/startup.py +++ b/app/core/startup.py @@ -336,6 +336,7 @@ ROLE_PERMISSION_MAP = { "accounting.tally.view", "accounting.tally.connect", "accounting.tally.map_company", + "accounting.tally.sync_masters", "accounting.act.initialize", "clients.view.own_only", "employees.dashboard.view", @@ -969,3 +970,4 @@ def on_startup(app: FastAPI) -> None: # Phase 7O: start alert notification/escalation automation after schema and seed checks. start_notification_scheduler() + diff --git a/app/modules/accounting/act_store.py b/app/modules/accounting/act_store.py index 09034d6..2f32a18 100644 --- a/app/modules/accounting/act_store.py +++ b/app/modules/accounting/act_store.py @@ -8,7 +8,7 @@ import os import sqlite3 from typing import Iterator, Sequence -ACT_SCHEMA_VERSION = 2 +ACT_SCHEMA_VERSION = 3 class AccountingActStoreError(RuntimeError): @@ -31,8 +31,8 @@ class AccountingActStore: Phase 1 provides accounting storage and read-only Tally discovery. Phase 2 adds durable client/registration -> Tally company mapping keyed by - Tally GUID. Accounting master/transaction sync and write-back remain out of - scope. + Tally GUID. Phase 3 adds read-only Tally accounting master snapshot tables. + Voucher/transaction sync and write-back remain out of scope. """ def __init__(self, root: str | Path) -> None: @@ -136,6 +136,44 @@ class AccountingActStore: details_json TEXT NOT NULL DEFAULT '{}' ); + CREATE TABLE IF NOT EXISTS tally_groups ( + id INTEGER PRIMARY KEY AUTOINCREMENT, tally_guid TEXT NOT NULL, company_name TEXT NOT NULL, master_guid TEXT NOT NULL DEFAULT '', name TEXT NOT NULL, parent TEXT NOT NULL DEFAULT '', synced_at_utc TEXT NOT NULL, payload_json TEXT NOT NULL DEFAULT '{}' + ); + CREATE TABLE IF NOT EXISTS tally_ledgers ( + id INTEGER PRIMARY KEY AUTOINCREMENT, tally_guid TEXT NOT NULL, company_name TEXT NOT NULL, master_guid TEXT NOT NULL DEFAULT '', name TEXT NOT NULL, parent TEXT NOT NULL DEFAULT '', opening_balance REAL NOT NULL DEFAULT 0, closing_balance REAL NOT NULL DEFAULT 0, is_billwise_on TEXT NOT NULL DEFAULT '', gst_applicable TEXT NOT NULL DEFAULT '', synced_at_utc TEXT NOT NULL, payload_json TEXT NOT NULL DEFAULT '{}' + ); + CREATE TABLE IF NOT EXISTS tally_voucher_types ( + id INTEGER PRIMARY KEY AUTOINCREMENT, tally_guid TEXT NOT NULL, company_name TEXT NOT NULL, master_guid TEXT NOT NULL DEFAULT '', name TEXT NOT NULL, parent TEXT NOT NULL DEFAULT '', numbering_method TEXT NOT NULL DEFAULT '', synced_at_utc TEXT NOT NULL, payload_json TEXT NOT NULL DEFAULT '{}' + ); + CREATE TABLE IF NOT EXISTS tally_stock_groups ( + id INTEGER PRIMARY KEY AUTOINCREMENT, tally_guid TEXT NOT NULL, company_name TEXT NOT NULL, master_guid TEXT NOT NULL DEFAULT '', name TEXT NOT NULL, parent TEXT NOT NULL DEFAULT '', synced_at_utc TEXT NOT NULL, payload_json TEXT NOT NULL DEFAULT '{}' + ); + CREATE TABLE IF NOT EXISTS tally_stock_categories ( + id INTEGER PRIMARY KEY AUTOINCREMENT, tally_guid TEXT NOT NULL, company_name TEXT NOT NULL, master_guid TEXT NOT NULL DEFAULT '', name TEXT NOT NULL, parent TEXT NOT NULL DEFAULT '', synced_at_utc TEXT NOT NULL, payload_json TEXT NOT NULL DEFAULT '{}' + ); + CREATE TABLE IF NOT EXISTS tally_stock_items ( + id INTEGER PRIMARY KEY AUTOINCREMENT, tally_guid TEXT NOT NULL, company_name TEXT NOT NULL, master_guid TEXT NOT NULL DEFAULT '', name TEXT NOT NULL, parent TEXT NOT NULL DEFAULT '', category TEXT NOT NULL DEFAULT '', base_units TEXT NOT NULL DEFAULT '', hsn_code TEXT NOT NULL DEFAULT '', synced_at_utc TEXT NOT NULL, payload_json TEXT NOT NULL DEFAULT '{}' + ); + CREATE TABLE IF NOT EXISTS tally_units ( + id INTEGER PRIMARY KEY AUTOINCREMENT, tally_guid TEXT NOT NULL, company_name TEXT NOT NULL, master_guid TEXT NOT NULL DEFAULT '', name TEXT NOT NULL, original_name TEXT NOT NULL DEFAULT '', synced_at_utc TEXT NOT NULL, payload_json TEXT NOT NULL DEFAULT '{}' + ); + CREATE TABLE IF NOT EXISTS tally_cost_categories ( + id INTEGER PRIMARY KEY AUTOINCREMENT, tally_guid TEXT NOT NULL, company_name TEXT NOT NULL, master_guid TEXT NOT NULL DEFAULT '', name TEXT NOT NULL, synced_at_utc TEXT NOT NULL, payload_json TEXT NOT NULL DEFAULT '{}' + ); + CREATE TABLE IF NOT EXISTS tally_cost_centres ( + id INTEGER PRIMARY KEY AUTOINCREMENT, tally_guid TEXT NOT NULL, company_name TEXT NOT NULL, master_guid TEXT NOT NULL DEFAULT '', name TEXT NOT NULL, parent TEXT NOT NULL DEFAULT '', category TEXT NOT NULL DEFAULT '', synced_at_utc TEXT NOT NULL, payload_json TEXT NOT NULL DEFAULT '{}' + ); + + CREATE INDEX IF NOT EXISTS ix_tally_groups_company ON tally_groups(tally_guid, name); + CREATE INDEX IF NOT EXISTS ix_tally_ledgers_company ON tally_ledgers(tally_guid, name); + CREATE INDEX IF NOT EXISTS ix_tally_voucher_types_company ON tally_voucher_types(tally_guid, name); + CREATE INDEX IF NOT EXISTS ix_tally_stock_groups_company ON tally_stock_groups(tally_guid, name); + CREATE INDEX IF NOT EXISTS ix_tally_stock_categories_company ON tally_stock_categories(tally_guid, name); + CREATE INDEX IF NOT EXISTS ix_tally_stock_items_company ON tally_stock_items(tally_guid, name); + CREATE INDEX IF NOT EXISTS ix_tally_units_company ON tally_units(tally_guid, name); + CREATE INDEX IF NOT EXISTS ix_tally_cost_categories_company ON tally_cost_categories(tally_guid, name); + CREATE INDEX IF NOT EXISTS ix_tally_cost_centres_company ON tally_cost_centres(tally_guid, name); + CREATE INDEX IF NOT EXISTS ix_tally_companies_last_seen ON tally_companies(last_seen_at_utc); CREATE INDEX IF NOT EXISTS ix_tally_company_mapping_client_active diff --git a/app/modules/accounting/templates/accounting/tally.html b/app/modules/accounting/templates/accounting/tally.html index 51d58a0..dbe548e 100644 --- a/app/modules/accounting/templates/accounting/tally.html +++ b/app/modules/accounting/templates/accounting/tally.html @@ -1,11 +1,12 @@ {% extends "ui/templates/base/layout.html" %} {% block content %} +{% if synced %}
Tally accounting masters synchronized successfully.
{% endif %}

Tools · Accounting

-

Tally Company Mapping

-

Phase 2 permanently maps an ERP client or registration to a loaded TallyPrime company using the Tally GUID.

+

Tally Accounting Masters

+

Phase 3 permanently maps an ERP client or registration to a loaded TallyPrime company using the Tally GUID.

Refresh Tally Companies
@@ -122,7 +123,7 @@ @@ -149,6 +150,47 @@ {% endif %} + {% if accounting and accounting.exists and tally and tally.connected %} +
+
+
+

Phase 3 · Accounting Master Sync

+

Select the company that is currently open/loaded in TallyPrime. The Local Agent verifies that the selected company is already mapped to this ERP client before reading masters.

+
+ Read-only +
+ +
+ + + + +
+ +
+ {% set last_sync = accounting.latest_master_sync if accounting else none %} +
Latest Sync
{{ last_sync.completed_at_utc if last_sync and last_sync.completed_at_utc else 'Not synced yet' }}
+
Company
{{ last_sync.company_name if last_sync else '-' }}
+
Rows Stored
{{ last_sync.rows_processed if last_sync else 0 }}
+
Status
{{ last_sync.status|title if last_sync else '-' }}
+
+ +
+ Synced into the selected client's .act database: Groups, Ledgers, Voucher Types, Stock Groups, Stock Categories, Stock Items, Units, Cost Categories and Cost Centres. No voucher or transaction data is imported in Phase 3. +
+
+ {% if accounting and accounting.exists %}
@@ -236,7 +278,7 @@ {% endif %}
- Phase 2 only stores company mappings. Ledger masters, vouchers, inventory and Tally write-back remain disabled. + Phase 3 synchronizes accounting masters only. Voucher/transaction sync and Tally write-back remain disabled and are reserved for later phases.
{% endblock %} diff --git a/app/modules/accounting/ui.py b/app/modules/accounting/ui.py index eb1ba69..34e49e6 100644 --- a/app/modules/accounting/ui.py +++ b/app/modules/accounting/ui.py @@ -114,6 +114,7 @@ def tally_tool( initialized: int = 0, mapped: int = 0, unmapped: int = 0, + synced: int = 0, error: str = "", ): db = CommonSessionLocal() @@ -148,7 +149,7 @@ def tally_tool( try: response_data = request_agent_command( node.node_code, - "phase2_status", + "phase3_status", payload, timeout_seconds=20, ) @@ -163,7 +164,7 @@ def tally_tool( request, db, user, - title="Tally Company Mapping", + title="Tally Accounting Masters", clients=clients, selected_client=selected_client, registrations=registrations, @@ -173,6 +174,7 @@ def tally_tool( initialized=bool(initialized), mapped=bool(mapped), unmapped=bool(unmapped), + synced=bool(synced), command_error=command_error, ) finally: @@ -358,3 +360,51 @@ def unmap_tally_company( ) finally: db.close() + + +@router.post("/sync-masters") +def sync_tally_masters( + request: Request, + client_id: int = Form(...), + tally_guid: str = Form(...), + csrf_token: str = Form(...), +): + validate_csrf(request, csrf_token) + db = CommonSessionLocal() + try: + user, response = _require_partner(request, db, "accounting.tally.sync_masters") + if response: + return response + client, _clients, scope = _find_visible_client(db, request, user, client_id) + if not client: + return _denied() + node = get_active_storage_node_for_branch(db, scope.tenant_id, scope.branch_id) + if not node or not _node_online(node): + return RedirectResponse( + url=f"/tools/tally?client_id={client.id}&error={quote('ERP Local Agent is offline for the active branch.')}", + status_code=303, + ) + try: + result = request_agent_command( + node.node_code, + "accounting_sync_masters", + { + "client_id": int(client.id), + "tally_guid": str(tally_guid or "").strip(), + "requested_by_user_id": int(user.id), + }, + timeout_seconds=120, + ) + if not result.get("ok"): + raise RuntimeError(str(result.get("error") or "Tally master synchronization failed.")) + except Exception as exc: + return RedirectResponse( + url=f"/tools/tally?client_id={client.id}&error={quote(str(exc))}", + status_code=303, + ) + return RedirectResponse( + url=f"/tools/tally?client_id={client.id}&refresh=1&synced=1", + status_code=303, + ) + finally: + db.close() diff --git a/app/modules/core/rbac/permissions_registry.py b/app/modules/core/rbac/permissions_registry.py index 5138df4..50140cf 100644 --- a/app/modules/core/rbac/permissions_registry.py +++ b/app/modules/core/rbac/permissions_registry.py @@ -150,6 +150,7 @@ PERMISSIONS = { "accounting.tally.view": "View Tally Accounting Tool", "accounting.tally.connect": "Connect to Tally Through Local Agent", "accounting.tally.map_company": "Map ERP Client or Registration to Tally Company", + "accounting.tally.sync_masters": "Synchronize Tally Accounting Masters", "accounting.act.initialize": "Initialize Client Accounting ACT Storage", "notice_cases.view": "View Notice and Case Management", @@ -179,3 +180,4 @@ def expand_permission_codes(code: str) -> list[str]: codes.append(alias) return codes + diff --git a/app/modules/documents/agent_package.py b/app/modules/documents/agent_package.py index ffae7c6..707d372 100644 --- a/app/modules/documents/agent_package.py +++ b/app/modules/documents/agent_package.py @@ -4,14 +4,9 @@ import io from pathlib import Path import zipfile -ERP_LOCAL_AGENT_VERSION = "1.3.0" +ERP_LOCAL_AGENT_VERSION = "1.4.0" ERP_LOCAL_AGENT_NAME = "ERP Local Agent" RUNTIME_ROOT = Path(__file__).resolve().parent / "local_agent_runtime" - -# ZIP metadata is part of the ZIP byte stream. If writestr() is allowed to use -# the current clock time, two packages built from identical source files have -# different SHA256 hashes. Keep update packages reproducible so the SHA256 -# returned by update-manifest is exactly the SHA256 of update-package. _DETERMINISTIC_ZIP_TIMESTAMP = (2026, 1, 1, 0, 0, 0) @@ -23,24 +18,11 @@ def build_agent_env(*, erp_base_url: str, node_code: str, node_secret: str, stor erp_base_url = "https://" + erp_base_url[7:] storage_root = (storage_root or r"D:\AuditFirmStorage").strip() return ( - f"ERP_BASE_URL={erp_base_url}\n" - f"NODE_CODE={(node_code or '').strip()}\n" - f"NODE_SECRET={(node_secret or '').strip()}\n" - f"STORAGE_ROOT={storage_root}\n" - f"TENANT_ID={'' if tenant_id is None else tenant_id}\n" - f"AUDIT_FIRM_ID={'' if tenant_id is None else tenant_id}\n" - f"BRANCH_ID={'' if branch_id is None else branch_id}\n" - f"SYNC_INTERVAL_SECONDS={int(sync_interval_seconds or 30)}\n" - f"POLL_INTERVAL_SECONDS={int(sync_interval_seconds or 30)}\n" - f"REQUEST_TIMEOUT_SECONDS={int(request_timeout_seconds or 60)}\n" - f"TUNNEL_ENABLED={str(bool(tunnel_enabled)).lower()}\n" - f"TUNNEL_RECONNECT_SECONDS={int(tunnel_reconnect_seconds or 10)}\n" - "AUTO_UPDATE=true\n" - "AUTO_INSTALL_UPDATES=false\n" - "UPDATE_CHECK_INTERVAL_SECONDS=300\n" - "DASHBOARD_ENABLED=true\n" - "DASHBOARD_HOST=127.0.0.1\n" - "DASHBOARD_PORT=8788\n" + f"ERP_BASE_URL={erp_base_url}\n" f"NODE_CODE={(node_code or '').strip()}\n" f"NODE_SECRET={(node_secret or '').strip()}\n" + f"STORAGE_ROOT={storage_root}\n" f"TENANT_ID={'' if tenant_id is None else tenant_id}\n" f"AUDIT_FIRM_ID={'' if tenant_id is None else tenant_id}\n" + f"BRANCH_ID={'' if branch_id is None else branch_id}\n" f"SYNC_INTERVAL_SECONDS={int(sync_interval_seconds or 30)}\n" f"POLL_INTERVAL_SECONDS={int(sync_interval_seconds or 30)}\n" + f"REQUEST_TIMEOUT_SECONDS={int(request_timeout_seconds or 60)}\n" f"TUNNEL_ENABLED={str(bool(tunnel_enabled)).lower()}\n" f"TUNNEL_RECONNECT_SECONDS={int(tunnel_reconnect_seconds or 10)}\n" + "AUTO_UPDATE=true\nAUTO_INSTALL_UPDATES=false\nUPDATE_CHECK_INTERVAL_SECONDS=300\nDASHBOARD_ENABLED=true\nDASHBOARD_HOST=127.0.0.1\nDASHBOARD_PORT=8788\n" ) @@ -48,8 +30,6 @@ def _zip_info(name: str) -> zipfile.ZipInfo: info = zipfile.ZipInfo(filename=name, date_time=_DETERMINISTIC_ZIP_TIMESTAMP) info.compress_type = zipfile.ZIP_DEFLATED info.create_system = 3 - # Regular file with rw-r--r-- permissions. Stable metadata keeps ZIP bytes - # reproducible on every request and across Linux deployments. info.external_attr = (0o100644 & 0xFFFF) << 16 return info @@ -61,50 +41,28 @@ def _write_zip_bytes(dst: zipfile.ZipFile, name: str, payload: bytes) -> None: def _build_zip(*, env_text: str | None, include_env: bool, include_admin_readme: bool) -> bytes: if not RUNTIME_ROOT.exists(): raise RuntimeError(f"ERP Local Agent runtime is missing: {RUNTIME_ROOT}") - buffer = io.BytesIO() - with zipfile.ZipFile( - buffer, - "w", - compression=zipfile.ZIP_DEFLATED, - compresslevel=9, - ) as dst: + with zipfile.ZipFile(buffer, "w", compression=zipfile.ZIP_DEFLATED, compresslevel=9) as dst: for path in sorted(RUNTIME_ROOT.rglob("*"), key=lambda item: item.as_posix()): if not path.is_file() or "__pycache__" in path.parts: continue - name = path.relative_to(RUNTIME_ROOT).as_posix() - _write_zip_bytes(dst, name, path.read_bytes()) - + _write_zip_bytes(dst, path.relative_to(RUNTIME_ROOT).as_posix(), path.read_bytes()) if include_env and env_text is not None: _write_zip_bytes(dst, ".env", env_text.encode("utf-8")) - if include_admin_readme: - readme = ( + text = ( f"ERP Local Agent {ERP_LOCAL_AGENT_VERSION}\n" - "Existing storage, WebSocket tunnel, Tally, accounting .act and local dashboard functionality are preserved.\n" - "Local dashboard: http://127.0.0.1:8788\n" - "The agent checks for updates automatically but installation is always initiated by the local user.\n" - "The existing .env, data, logs, .venv and client storage are preserved during updates.\n" + "Existing storage, WebSocket tunnel, dashboard, Tally mapping and client .act functionality are preserved.\n" + "Phase 3 adds read-only Tally accounting master synchronization.\n" + "Dashboard: http://127.0.0.1:8788\n" ) - _write_zip_bytes(dst, "README_ERP_LOCAL_AGENT.txt", readme.encode("utf-8")) - + _write_zip_bytes(dst, "README_ERP_LOCAL_AGENT.txt", text.encode("utf-8")) return buffer.getvalue() def build_preconfigured_agent_zip(*, env_text: str, include_admin_readme: bool = False) -> bytes: - return _build_zip( - env_text=env_text, - include_env=True, - include_admin_readme=include_admin_readme, - ) + return _build_zip(env_text=env_text, include_env=True, include_admin_readme=include_admin_readme) def build_agent_update_zip() -> bytes: - # Update package deliberately excludes .env, .venv, data, logs and client files. - # The ZIP is deterministic so update-manifest and update-package always agree - # on SHA256 for the same deployed runtime/version. - return _build_zip( - env_text=None, - include_env=False, - include_admin_readme=False, - ) + return _build_zip(env_text=None, include_env=False, include_admin_readme=False) diff --git a/app/modules/documents/local_agent_runtime/README_ERP_LOCAL_AGENT.txt b/app/modules/documents/local_agent_runtime/README_ERP_LOCAL_AGENT.txt index 7f9d5db..0e01ffd 100644 --- a/app/modules/documents/local_agent_runtime/README_ERP_LOCAL_AGENT.txt +++ b/app/modules/documents/local_agent_runtime/README_ERP_LOCAL_AGENT.txt @@ -1,4 +1,4 @@ -ERP Local Agent 1.3.0 +ERP Local Agent 1.4.0 Existing storage, WebSocket tunnel, Tally and client .act functionality are preserved. @@ -18,3 +18,5 @@ Local operational database: Client accounting .act databases remain separate under the configured STORAGE_ROOT. Phase 2: client/registration to Tally company mapping is supported using Tally GUID. + +Phase 3: read-only Tally accounting master sync (Groups, Ledgers, Voucher Types, Stock Groups/Categories/Items, Units and Cost Centres/Categories) into client .act storage. diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/__init__.py b/app/modules/documents/local_agent_runtime/erp_local_agent/__init__.py index 893359a..4a07c3b 100644 --- a/app/modules/documents/local_agent_runtime/erp_local_agent/__init__.py +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/__init__.py @@ -1,2 +1,2 @@ -__version__ = "1.3.0" +__version__ = "1.4.0" AGENT_NAME = "ERP Local Agent" diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/accounting_store.py b/app/modules/documents/local_agent_runtime/erp_local_agent/accounting_store.py index 59e5226..ea6770e 100644 --- a/app/modules/documents/local_agent_runtime/erp_local_agent/accounting_store.py +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/accounting_store.py @@ -7,7 +7,19 @@ import sqlite3 from typing import Sequence -SCHEMA_VERSION = "2" +SCHEMA_VERSION = "3" + +MASTER_TABLES = { + "groups": "tally_groups", + "ledgers": "tally_ledgers", + "voucher_types": "tally_voucher_types", + "stock_groups": "tally_stock_groups", + "stock_categories": "tally_stock_categories", + "stock_items": "tally_stock_items", + "units": "tally_units", + "cost_categories": "tally_cost_categories", + "cost_centres": "tally_cost_centres", +} def _utc_now_iso() -> str: @@ -15,11 +27,10 @@ def _utc_now_iso() -> str: class LocalAccountingStore: - """Client-scoped SQLite .act storage under the existing branch storage root. + """Client-scoped SQLite .act accounting store. - Phase 2 persists ERP client/registration -> Tally company mappings using the - Tally GUID as the durable identifier. Existing Phase 1 databases are upgraded - in place without deleting accounting history. + Phase 3 preserves Phase 1/2 metadata and mappings and adds company-scoped + read-only Tally master snapshots. Existing .act files are upgraded in place. """ def __init__(self, storage_root: Path): @@ -43,7 +54,7 @@ class LocalAccountingStore: def connect(self, client_id: int): path = self.db_path(client_id) path.parent.mkdir(parents=True, exist_ok=True) - db = sqlite3.connect(path, timeout=30) + db = sqlite3.connect(path, timeout=60) db.row_factory = sqlite3.Row db.execute("PRAGMA foreign_keys=ON") db.execute("PRAGMA journal_mode=WAL") @@ -69,81 +80,102 @@ class LocalAccountingStore: if name not in columns: db.execute(f"ALTER TABLE tally_company_mapping ADD COLUMN {name} {ddl}") - def initialize( - self, - client_id: int, - client_name: str = "", - tenant_id: int | None = None, - created_by_user_id: int | None = None, - ) -> Path: + @staticmethod + def _master_table_ddl(table: str) -> str: + return f""" + CREATE TABLE IF NOT EXISTS {table} ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + tally_guid TEXT NOT NULL, + company_name TEXT NOT NULL, + master_guid TEXT NOT NULL DEFAULT '', + name TEXT NOT NULL, + parent TEXT NOT NULL DEFAULT '', + category TEXT NOT NULL DEFAULT '', + reserved_name TEXT NOT NULL DEFAULT '', + base_units TEXT NOT NULL DEFAULT '', + additional_units TEXT NOT NULL DEFAULT '', + original_name TEXT NOT NULL DEFAULT '', + opening_balance REAL NOT NULL DEFAULT 0, + closing_balance REAL NOT NULL DEFAULT 0, + opening_value REAL NOT NULL DEFAULT 0, + opening_rate TEXT NOT NULL DEFAULT '', + numbering_method TEXT NOT NULL DEFAULT '', + tax_type TEXT NOT NULL DEFAULT '', + gst_applicable TEXT NOT NULL DEFAULT '', + gst_registration_type TEXT NOT NULL DEFAULT '', + gst_type_of_supply TEXT NOT NULL DEFAULT '', + hsn_code TEXT NOT NULL DEFAULT '', + is_revenue TEXT NOT NULL DEFAULT '', + is_deemed_positive TEXT NOT NULL DEFAULT '', + is_billwise_on TEXT NOT NULL DEFAULT '', + is_simple_unit TEXT NOT NULL DEFAULT '', + conversion TEXT NOT NULL DEFAULT '', + is_active TEXT NOT NULL DEFAULT '', + synced_at_utc TEXT NOT NULL, + payload_json TEXT NULL + ); + CREATE INDEX IF NOT EXISTS ix_{table}_company ON {table}(tally_guid, name); + CREATE INDEX IF NOT EXISTS ix_{table}_master_guid ON {table}(tally_guid, master_guid); + """ + + def initialize(self, client_id: int, client_name: str = "", tenant_id: int | None = None, created_by_user_id: int | None = None) -> Path: path = self.db_path(client_id) with self.connect(client_id) as db: db.executescript( """ CREATE TABLE IF NOT EXISTS act_meta ( - key TEXT PRIMARY KEY, - value TEXT NOT NULL, - updated_at_utc TEXT NOT NULL + key TEXT PRIMARY KEY, value TEXT NOT NULL, updated_at_utc TEXT NOT NULL ); CREATE TABLE IF NOT EXISTS tally_companies ( id INTEGER PRIMARY KEY AUTOINCREMENT, - tally_guid TEXT NOT NULL DEFAULT '', - company_name TEXT NOT NULL, - gstin TEXT NOT NULL DEFAULT '', - first_seen_at_utc TEXT NOT NULL, - last_seen_at_utc TEXT NOT NULL, - is_currently_loaded INTEGER NOT NULL DEFAULT 0 + tally_guid TEXT NOT NULL DEFAULT '', company_name TEXT NOT NULL, + gstin TEXT NOT NULL DEFAULT '', first_seen_at_utc TEXT NOT NULL, + last_seen_at_utc TEXT NOT NULL, is_currently_loaded INTEGER NOT NULL DEFAULT 0 ); CREATE INDEX IF NOT EXISTS ix_tally_companies_guid ON tally_companies(tally_guid); CREATE INDEX IF NOT EXISTS ix_tally_companies_name ON tally_companies(company_name); CREATE TABLE IF NOT EXISTS tally_company_mapping ( - id INTEGER PRIMARY KEY AUTOINCREMENT, - client_id INTEGER NOT NULL, - registration_id INTEGER NULL, - registration_type_code TEXT NOT NULL DEFAULT '', - registration_number TEXT NOT NULL DEFAULT '', - registration_legal_name TEXT NOT NULL DEFAULT '', - registration_trade_name TEXT NOT NULL DEFAULT '', - business_unit_id INTEGER, - client_branch_id INTEGER, - tally_guid TEXT NOT NULL, - company_name TEXT NOT NULL, - gstin TEXT NOT NULL DEFAULT '', - is_active INTEGER NOT NULL DEFAULT 1, - created_at_utc TEXT NOT NULL, - updated_at_utc TEXT NOT NULL, - mapped_by_user_id INTEGER, - unmapped_at_utc TEXT, - unmapped_by_user_id INTEGER + id INTEGER PRIMARY KEY AUTOINCREMENT, client_id INTEGER NOT NULL, + registration_id INTEGER NULL, registration_type_code TEXT NOT NULL DEFAULT '', + registration_number TEXT NOT NULL DEFAULT '', registration_legal_name TEXT NOT NULL DEFAULT '', + registration_trade_name TEXT NOT NULL DEFAULT '', business_unit_id INTEGER, + client_branch_id INTEGER, tally_guid TEXT NOT NULL, company_name TEXT NOT NULL, + gstin TEXT NOT NULL DEFAULT '', is_active INTEGER NOT NULL DEFAULT 1, + created_at_utc TEXT NOT NULL, updated_at_utc TEXT NOT NULL, + mapped_by_user_id INTEGER, unmapped_at_utc TEXT, unmapped_by_user_id INTEGER ); - CREATE INDEX IF NOT EXISTS ix_tally_company_mapping_client_active - ON tally_company_mapping(client_id, is_active); - CREATE INDEX IF NOT EXISTS ix_tally_company_mapping_guid - ON tally_company_mapping(tally_guid, is_active); + CREATE INDEX IF NOT EXISTS ix_tally_company_mapping_client_active ON tally_company_mapping(client_id, is_active); + CREATE INDEX IF NOT EXISTS ix_tally_company_mapping_guid ON tally_company_mapping(tally_guid, is_active); CREATE TABLE IF NOT EXISTS tally_connection_history ( - id INTEGER PRIMARY KEY AUTOINCREMENT, - checked_at_utc TEXT NOT NULL, - connected INTEGER NOT NULL, - tally_url TEXT NOT NULL DEFAULT '', - company_count INTEGER NOT NULL DEFAULT 0, - error_message TEXT NULL, - payload_json TEXT NULL + id INTEGER PRIMARY KEY AUTOINCREMENT, checked_at_utc TEXT NOT NULL, + connected INTEGER NOT NULL, tally_url TEXT NOT NULL DEFAULT '', + company_count INTEGER NOT NULL DEFAULT 0, error_message TEXT NULL, payload_json TEXT NULL ); CREATE TABLE IF NOT EXISTS tally_sync_runs ( id INTEGER PRIMARY KEY AUTOINCREMENT, - sync_type TEXT NOT NULL, - status TEXT NOT NULL, - started_at_utc TEXT NOT NULL, - completed_at_utc TEXT NULL, - rows_processed INTEGER NOT NULL DEFAULT 0, - error_message TEXT NULL, - details_json TEXT NULL + sync_type TEXT NOT NULL, tally_guid TEXT NOT NULL DEFAULT '', company_name TEXT NOT NULL DEFAULT '', + mapping_id INTEGER, requested_by_user_id INTEGER, + status TEXT NOT NULL, started_at_utc TEXT NOT NULL, completed_at_utc TEXT NULL, + rows_processed INTEGER NOT NULL DEFAULT 0, error_message TEXT NULL, details_json TEXT NULL ); """ ) self._ensure_mapping_columns(db) + # Upgrade older tally_sync_runs created in Phase 1 without dropping history. + sync_columns = {row["name"] for row in db.execute("PRAGMA table_info(tally_sync_runs)").fetchall()} + for name, ddl in { + "tally_guid": "TEXT NOT NULL DEFAULT ''", + "company_name": "TEXT NOT NULL DEFAULT ''", + "mapping_id": "INTEGER", + "requested_by_user_id": "INTEGER", + }.items(): + if name not in sync_columns: + db.execute(f"ALTER TABLE tally_sync_runs ADD COLUMN {name} {ddl}") + for table in MASTER_TABLES.values(): + db.executescript(self._master_table_ddl(table)) + now = _utc_now_iso() meta = { "schema_version": SCHEMA_VERSION, @@ -175,8 +207,7 @@ class LocalAccountingStore: guid = str(company.get("guid") or "").strip() gstin = str(company.get("gstin") or "").strip().upper() existing = db.execute( - "SELECT id FROM tally_companies WHERE tally_guid=? AND company_name=? LIMIT 1", - (guid, name), + "SELECT id FROM tally_companies WHERE tally_guid=? AND company_name=? LIMIT 1", (guid, name) ).fetchone() if existing: db.execute( @@ -192,26 +223,10 @@ class LocalAccountingStore: db.execute( """INSERT INTO tally_connection_history(checked_at_utc, connected, tally_url, company_count, error_message, payload_json) VALUES (?, ?, ?, ?, ?, ?)""", - ( - now, - 1 if status.get("connected") else 0, - str(status.get("url") or ""), - int(status.get("company_count") or 0), - str(status.get("error") or "") or None, - json.dumps(status, ensure_ascii=False, separators=(",", ":")), - ), + (now, 1 if status.get("connected") else 0, str(status.get("url") or ""), int(status.get("company_count") or 0), str(status.get("error") or "") or None, json.dumps(status, ensure_ascii=False, separators=(",", ":"))), ) - def map_company( - self, - client_id: int, - *, - tally_guid: str, - company_name: str, - gstin: str = "", - registration: dict | None = None, - mapped_by_user_id: int | None = None, - ) -> dict: + def map_company(self, client_id: int, *, tally_guid: str, company_name: str, gstin: str = "", registration: dict | None = None, mapped_by_user_id: int | None = None) -> dict: if not self.exists(client_id): raise ValueError("Accounting storage is not initialized for this client.") guid = str(tally_guid or "").strip() @@ -220,136 +235,145 @@ class LocalAccountingStore: raise ValueError("Tally company GUID is required for permanent mapping.") if not name: raise ValueError("Tally company name is required.") - registration = registration or {} registration_id = registration.get("id") registration_id = int(registration_id) if registration_id not in (None, "") else None now = _utc_now_iso() - with self.connect(client_id) as db: if registration_id is None: - db.execute( - """UPDATE tally_company_mapping - SET is_active=0, updated_at_utc=?, unmapped_at_utc=? - WHERE client_id=? AND registration_id IS NULL AND is_active=1""", - (now, now, int(client_id)), - ) + db.execute("UPDATE tally_company_mapping SET is_active=0, updated_at_utc=?, unmapped_at_utc=? WHERE client_id=? AND registration_id IS NULL AND is_active=1", (now, now, int(client_id))) else: - db.execute( - """UPDATE tally_company_mapping - SET is_active=0, updated_at_utc=?, unmapped_at_utc=? - WHERE client_id=? AND registration_id=? AND is_active=1""", - (now, now, int(client_id), registration_id), - ) - - db.execute( - """UPDATE tally_company_mapping - SET is_active=0, updated_at_utc=?, unmapped_at_utc=? - WHERE client_id=? AND tally_guid=? AND is_active=1""", - (now, now, int(client_id), guid), - ) - + db.execute("UPDATE tally_company_mapping SET is_active=0, updated_at_utc=?, unmapped_at_utc=? WHERE client_id=? AND registration_id=? AND is_active=1", (now, now, int(client_id), registration_id)) + db.execute("UPDATE tally_company_mapping SET is_active=0, updated_at_utc=?, unmapped_at_utc=? WHERE client_id=? AND tally_guid=? AND is_active=1", (now, now, int(client_id), guid)) cur = db.execute( - """INSERT INTO tally_company_mapping( - client_id, registration_id, registration_type_code, - registration_number, registration_legal_name, - registration_trade_name, business_unit_id, client_branch_id, - tally_guid, company_name, gstin, is_active, - created_at_utc, updated_at_utc, mapped_by_user_id - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1, ?, ?, ?)""", - ( - int(client_id), - registration_id, - str(registration.get("registration_type_code") or ""), - str(registration.get("registration_number") or ""), - str(registration.get("legal_name") or ""), - str(registration.get("trade_name") or ""), - registration.get("business_unit_id"), - registration.get("client_branch_id"), - guid, - name, - str(gstin or "").strip().upper(), - now, - now, - mapped_by_user_id, - ), + """INSERT INTO tally_company_mapping(client_id, registration_id, registration_type_code, registration_number, + registration_legal_name, registration_trade_name, business_unit_id, client_branch_id, + tally_guid, company_name, gstin, is_active, created_at_utc, updated_at_utc, mapped_by_user_id) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1, ?, ?, ?)""", + (int(client_id), registration_id, str(registration.get("registration_type_code") or ""), str(registration.get("registration_number") or ""), str(registration.get("legal_name") or ""), str(registration.get("trade_name") or ""), registration.get("business_unit_id"), registration.get("client_branch_id"), guid, name, str(gstin or "").strip().upper(), now, now, mapped_by_user_id), ) mapping_id = int(cur.lastrowid) return self.get_mapping(client_id, mapping_id) def get_mapping(self, client_id: int, mapping_id: int) -> dict: with self.connect(client_id) as db: - row = db.execute( - "SELECT * FROM tally_company_mapping WHERE id=? AND client_id=? LIMIT 1", - (int(mapping_id), int(client_id)), - ).fetchone() + row = db.execute("SELECT * FROM tally_company_mapping WHERE id=? AND client_id=? LIMIT 1", (int(mapping_id), int(client_id))).fetchone() if not row: raise ValueError("Tally company mapping was not found.") return dict(row) + def get_active_mapping_by_guid(self, client_id: int, tally_guid: str) -> dict: + with self.connect(client_id) as db: + row = db.execute( + "SELECT * FROM tally_company_mapping WHERE client_id=? AND tally_guid=? AND is_active=1 ORDER BY id DESC LIMIT 1", + (int(client_id), str(tally_guid or "").strip()), + ).fetchone() + if not row: + raise ValueError("The selected open Tally company is not mapped to this ERP client. Complete Phase 2 mapping first.") + return dict(row) + def unmap_company(self, client_id: int, mapping_id: int, unmapped_by_user_id: int | None = None) -> dict: if not self.exists(client_id): raise ValueError("Accounting storage is not initialized for this client.") now = _utc_now_iso() with self.connect(client_id) as db: - row = db.execute( - "SELECT id FROM tally_company_mapping WHERE id=? AND client_id=? AND is_active=1", - (int(mapping_id), int(client_id)), - ).fetchone() + row = db.execute("SELECT id FROM tally_company_mapping WHERE id=? AND client_id=? AND is_active=1", (int(mapping_id), int(client_id))).fetchone() if not row: raise ValueError("Active Tally company mapping was not found.") - db.execute( - """UPDATE tally_company_mapping - SET is_active=0, updated_at_utc=?, unmapped_at_utc=?, unmapped_by_user_id=? - WHERE id=? AND client_id=?""", - (now, now, unmapped_by_user_id, int(mapping_id), int(client_id)), - ) + db.execute("UPDATE tally_company_mapping SET is_active=0, updated_at_utc=?, unmapped_at_utc=?, unmapped_by_user_id=? WHERE id=? AND client_id=?", (now, now, unmapped_by_user_id, int(mapping_id), int(client_id))) return {"mapping_id": int(mapping_id), "unmapped": True} + def replace_master_snapshot(self, client_id: int, *, mapping: dict, masters: dict[str, list[dict]], requested_by_user_id: int | None = None) -> dict: + if not self.exists(client_id): + raise ValueError("Accounting storage is not initialized for this client.") + tally_guid = str(mapping.get("tally_guid") or "").strip() + company_name = str(mapping.get("company_name") or "").strip() + mapping_id = int(mapping.get("id")) + if not tally_guid or not company_name: + raise ValueError("Active Tally mapping is incomplete.") + started = _utc_now_iso() + counts: dict[str, int] = {} + run_id = None + try: + with self.connect(client_id) as db: + cur = db.execute( + """INSERT INTO tally_sync_runs(sync_type, tally_guid, company_name, mapping_id, requested_by_user_id, status, started_at_utc, rows_processed) + VALUES ('masters', ?, ?, ?, ?, 'running', ?, 0)""", + (tally_guid, company_name, mapping_id, requested_by_user_id, started), + ) + run_id = int(cur.lastrowid) + total = 0 + synced_at = _utc_now_iso() + for key, table in MASTER_TABLES.items(): + rows = masters.get(key) or [] + db.execute(f"DELETE FROM {table} WHERE tally_guid=?", (tally_guid,)) + for row in rows: + payload = json.dumps(row, ensure_ascii=False, separators=(",", ":")) + db.execute( + f"""INSERT INTO {table}( + tally_guid, company_name, master_guid, name, parent, category, reserved_name, + base_units, additional_units, original_name, opening_balance, closing_balance, + opening_value, opening_rate, numbering_method, tax_type, gst_applicable, + gst_registration_type, gst_type_of_supply, hsn_code, is_revenue, + is_deemed_positive, is_billwise_on, is_simple_unit, conversion, is_active, + synced_at_utc, payload_json + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""", + (tally_guid, company_name, str(row.get("guid") or ""), str(row.get("name") or ""), str(row.get("parent") or ""), str(row.get("category") or ""), str(row.get("reserved_name") or ""), str(row.get("base_units") or ""), str(row.get("additional_units") or ""), str(row.get("original_name") or ""), float(row.get("opening_balance") or 0), float(row.get("closing_balance") or 0), float(row.get("opening_value") or 0), str(row.get("opening_rate") or ""), str(row.get("numbering_method") or ""), str(row.get("tax_type") or ""), str(row.get("gst_applicable") or ""), str(row.get("gst_registration_type") or ""), str(row.get("gst_type_of_supply") or ""), str(row.get("hsn_code") or ""), str(row.get("is_revenue") or ""), str(row.get("is_deemed_positive") or ""), str(row.get("is_billwise_on") or ""), str(row.get("is_simple_unit") or ""), str(row.get("conversion") or ""), str(row.get("is_active") or ""), synced_at, payload), + ) + counts[key] = len(rows) + total += len(rows) + details = {"counts": counts, "schema_version": SCHEMA_VERSION} + db.execute( + "UPDATE tally_sync_runs SET status='completed', completed_at_utc=?, rows_processed=?, details_json=? WHERE id=?", + (_utc_now_iso(), total, json.dumps(details, ensure_ascii=False, separators=(",", ":")), run_id), + ) + return {"sync_run_id": run_id, "status": "completed", "company_name": company_name, "tally_guid": tally_guid, "counts": counts, "rows_processed": sum(counts.values())} + except Exception as exc: + if run_id is not None: + try: + with self.connect(client_id) as db: + db.execute("UPDATE tally_sync_runs SET status='failed', completed_at_utc=?, error_message=? WHERE id=?", (_utc_now_iso(), str(exc), run_id)) + except Exception: + pass + raise + + def _master_counts(self, db: sqlite3.Connection, tally_guid: str) -> dict[str, int]: + result = {} + for key, table in MASTER_TABLES.items(): + result[key] = int(db.execute(f"SELECT COUNT(*) FROM {table} WHERE tally_guid=?", (tally_guid,)).fetchone()[0]) + return result + def snapshot(self, client_id: int) -> dict: path = self.db_path(client_id) if not path.is_file(): - return { - "exists": False, - "db_path": str(path), - "metadata": {}, - "latest_connection": None, - "companies": [], - "mappings": [], - } - - # initialize() upgrades older Phase 1 schema in-place + return {"exists": False, "db_path": str(path), "metadata": {}, "latest_connection": None, "companies": [], "mappings": [], "latest_master_sync": None} self.initialize(client_id) with self.connect(client_id) as db: meta_rows = db.execute("SELECT key, value FROM act_meta ORDER BY key").fetchall() - companies = db.execute( - """SELECT tally_guid AS guid, company_name AS name, gstin, - first_seen_at_utc, last_seen_at_utc, is_currently_loaded - FROM tally_companies - ORDER BY is_currently_loaded DESC, company_name COLLATE NOCASE""" - ).fetchall() + companies = db.execute("SELECT tally_guid AS guid, company_name AS name, gstin, first_seen_at_utc, last_seen_at_utc, is_currently_loaded FROM tally_companies ORDER BY is_currently_loaded DESC, company_name COLLATE NOCASE").fetchall() mappings = db.execute( - """SELECT id, client_id, registration_id, registration_type_code, - registration_number, registration_legal_name, - registration_trade_name, business_unit_id, client_branch_id, - tally_guid, company_name, gstin, is_active, - created_at_utc AS mapped_at_utc, updated_at_utc, - mapped_by_user_id, unmapped_at_utc, unmapped_by_user_id - FROM tally_company_mapping - WHERE is_active=1 - ORDER BY CASE WHEN registration_id IS NULL THEN 0 ELSE 1 END, - registration_type_code, registration_number, id""" + """SELECT id, client_id, registration_id, registration_type_code, registration_number, + registration_legal_name, registration_trade_name, business_unit_id, client_branch_id, + tally_guid, company_name, gstin, is_active, created_at_utc AS mapped_at_utc, + updated_at_utc, mapped_by_user_id, unmapped_at_utc, unmapped_by_user_id + FROM tally_company_mapping WHERE is_active=1 + ORDER BY CASE WHEN registration_id IS NULL THEN 0 ELSE 1 END, registration_type_code, registration_number, id""" ).fetchall() - latest = db.execute( - "SELECT checked_at_utc, connected, tally_url, company_count, error_message FROM tally_connection_history ORDER BY id DESC LIMIT 1" - ).fetchone() - - loaded_guids = {str(row["guid"] or "") for row in companies if int(row["is_currently_loaded"] or 0)} - mapped = [] - for row in mappings: - item = dict(row) - item["currently_loaded"] = bool(item.get("tally_guid") and item["tally_guid"] in loaded_guids) - mapped.append(item) + latest = db.execute("SELECT checked_at_utc, connected, tally_url, company_count, error_message FROM tally_connection_history ORDER BY id DESC LIMIT 1").fetchone() + latest_sync = db.execute("SELECT id, sync_type, tally_guid, company_name, mapping_id, requested_by_user_id, status, started_at_utc, completed_at_utc, rows_processed, error_message, details_json FROM tally_sync_runs WHERE sync_type='masters' ORDER BY id DESC LIMIT 1").fetchone() + loaded_guids = {str(row["guid"] or "") for row in companies if int(row["is_currently_loaded"] or 0)} + mapped = [] + for row in mappings: + item = dict(row) + item["currently_loaded"] = bool(item.get("tally_guid") and item["tally_guid"] in loaded_guids) + item["master_counts"] = self._master_counts(db, str(item.get("tally_guid") or "")) + mapped.append(item) + latest_sync_dict = dict(latest_sync) if latest_sync else None + if latest_sync_dict and latest_sync_dict.get("details_json"): + try: + latest_sync_dict["details"] = json.loads(latest_sync_dict["details_json"]) + except Exception: + latest_sync_dict["details"] = {} return { "exists": True, "db_path": str(path), @@ -357,4 +381,5 @@ class LocalAccountingStore: "latest_connection": dict(latest) if latest else None, "companies": [dict(row) for row in companies], "mappings": mapped, + "latest_master_sync": latest_sync_dict, } diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/commands.py b/app/modules/documents/local_agent_runtime/erp_local_agent/commands.py index 662fc55..5009860 100644 --- a/app/modules/documents/local_agent_runtime/erp_local_agent/commands.py +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/commands.py @@ -23,7 +23,7 @@ class AgentCommandProcessor: error: str | None = None ok = False try: - if action in {"tally_status", "phase1_status", "phase2_status"}: + if action in {"tally_status", "phase1_status", "phase2_status", "phase3_status"}: result = self._status(payload) elif action == "accounting_initialize": result = self._initialize(payload) @@ -31,28 +31,21 @@ class AgentCommandProcessor: result = self._map_company(payload) elif action == "accounting_unmap_company": result = self._unmap_company(payload) + elif action == "accounting_sync_masters": + result = self._sync_masters(payload) else: raise ValueError(f"Unsupported local-agent command: {action}") ok = True except Exception as exc: error = str(exc) self.logger.exception("Agent command failed action=%s command_id=%s: %s", action, command_id, exc) - return { - "type": "command_result", - "command_id": command_id, - "ok": ok, - "result": result, - "error": error, - "agent_time_utc": datetime.now(timezone.utc).isoformat(), - } + return {"type": "command_result", "command_id": command_id, "ok": ok, "result": result, "error": error, "agent_time_utc": datetime.now(timezone.utc).isoformat()} def _agent_info(self) -> dict[str, Any]: return { - "name": "ERP Local Agent", - "version": __version__, - "tally_capability": True, - "accounting_act_capability": True, - "tally_mapping_capability": True, + "name": "ERP Local Agent", "version": __version__, + "tally_capability": True, "accounting_act_capability": True, + "tally_mapping_capability": True, "tally_master_sync_capability": True, } def _status(self, payload: dict[str, Any]) -> dict[str, Any]: @@ -64,32 +57,17 @@ class AgentCommandProcessor: if self.store.exists(client_id): self.store.record_tally_status(client_id, tally_status) accounting = self.store.snapshot(client_id) - return { - "agent": self._agent_info(), - "tally": tally_status, - "accounting": accounting, - } + return {"agent": self._agent_info(), "tally": tally_status, "accounting": accounting} def _initialize(self, payload: dict[str, Any]) -> dict[str, Any]: client_id = int(payload.get("client_id")) client_name = str(payload.get("client_name") or "").strip() tenant_id = payload.get("tenant_id") created_by_user_id = payload.get("requested_by_user_id") - path = self.store.initialize( - client_id, - client_name, - int(tenant_id) if tenant_id not in (None, "") else None, - int(created_by_user_id) if created_by_user_id not in (None, "") else None, - ) + path = self.store.initialize(client_id, client_name, int(tenant_id) if tenant_id not in (None, "") else None, int(created_by_user_id) if created_by_user_id not in (None, "") else None) tally_status = self.tally.status() self.store.record_tally_status(client_id, tally_status) - return { - "initialized": True, - "db_path": str(path), - "accounting": self.store.snapshot(client_id), - "tally": tally_status, - "agent": self._agent_info(), - } + return {"initialized": True, "db_path": str(path), "accounting": self.store.snapshot(client_id), "tally": tally_status, "agent": self._agent_info()} def _map_company(self, payload: dict[str, Any]) -> dict[str, Any]: client_id = int(payload.get("client_id")) @@ -99,74 +77,54 @@ class AgentCommandProcessor: requested_guid = str(payload.get("tally_guid") or "").strip() registration = payload.get("registration") or None allow_gstin_mismatch = bool(payload.get("allow_gstin_mismatch")) - if not requested_guid: raise ValueError("Select a Tally company before mapping.") - if not self.store.exists(client_id): - self.store.initialize( - client_id, - client_name, - int(tenant_id) if tenant_id not in (None, "") else None, - int(mapped_by_user_id) if mapped_by_user_id not in (None, "") else None, - ) - + self.store.initialize(client_id, client_name, int(tenant_id) if tenant_id not in (None, "") else None, int(mapped_by_user_id) if mapped_by_user_id not in (None, "") else None) tally_status = self.tally.status() if not tally_status.get("connected"): raise ValueError(str(tally_status.get("error") or "TallyPrime is not connected.")) - - companies = tally_status.get("companies") or [] - company = next( - (row for row in companies if str(row.get("guid") or "").strip() == requested_guid), - None, - ) + company = next((row for row in (tally_status.get("companies") or []) if str(row.get("guid") or "").strip() == requested_guid), None) if not company: raise ValueError("The selected Tally company is no longer loaded. Refresh Tally companies and try again.") - company_name = str(company.get("name") or "").strip() company_gstin = str(company.get("gstin") or "").strip().upper() if not company_name: raise ValueError("Tally returned an invalid company name.") - if not requested_guid: - raise ValueError("Tally returned no GUID. Permanent mapping requires a Tally GUID.") - if registration: reg_type = str(registration.get("registration_type_code") or "").strip().upper() reg_number = str(registration.get("registration_number") or "").strip().upper() if reg_type == "GSTIN" and reg_number and company_gstin and reg_number != company_gstin and not allow_gstin_mismatch: - raise ValueError( - f"GSTIN mismatch: ERP registration is {reg_number}, but Tally company reports {company_gstin}. " - "Verify the company or explicitly allow the mismatch." - ) - + raise ValueError(f"GSTIN mismatch: ERP registration is {reg_number}, but Tally company reports {company_gstin}. Verify the company or explicitly allow the mismatch.") self.store.record_tally_status(client_id, tally_status) - mapping = self.store.map_company( - client_id, - tally_guid=requested_guid, - company_name=company_name, - gstin=company_gstin, - registration=registration, - mapped_by_user_id=int(mapped_by_user_id) if mapped_by_user_id not in (None, "") else None, - ) - return { - "mapped": True, - "mapping": mapping, - "accounting": self.store.snapshot(client_id), - "tally": tally_status, - "agent": self._agent_info(), - } + mapping = self.store.map_company(client_id, tally_guid=requested_guid, company_name=company_name, gstin=company_gstin, registration=registration, mapped_by_user_id=int(mapped_by_user_id) if mapped_by_user_id not in (None, "") else None) + return {"mapped": True, "mapping": mapping, "accounting": self.store.snapshot(client_id), "tally": tally_status, "agent": self._agent_info()} def _unmap_company(self, payload: dict[str, Any]) -> dict[str, Any]: client_id = int(payload.get("client_id")) mapping_id = int(payload.get("mapping_id")) unmapped_by_user_id = payload.get("unmapped_by_user_id") - result = self.store.unmap_company( - client_id, - mapping_id, - int(unmapped_by_user_id) if unmapped_by_user_id not in (None, "") else None, - ) - return { - **result, - "accounting": self.store.snapshot(client_id), - "agent": self._agent_info(), - } + result = self.store.unmap_company(client_id, mapping_id, int(unmapped_by_user_id) if unmapped_by_user_id not in (None, "") else None) + return {**result, "accounting": self.store.snapshot(client_id), "agent": self._agent_info()} + + def _sync_masters(self, payload: dict[str, Any]) -> dict[str, Any]: + client_id = int(payload.get("client_id")) + requested_guid = str(payload.get("tally_guid") or "").strip() + requested_by_user_id = payload.get("requested_by_user_id") + if not requested_guid: + raise ValueError("Select a currently open Tally company before synchronizing masters.") + if not self.store.exists(client_id): + raise ValueError("Accounting storage is not initialized for this client.") + mapping = self.store.get_active_mapping_by_guid(client_id, requested_guid) + tally_status = self.tally.status() + if not tally_status.get("connected"): + raise ValueError(str(tally_status.get("error") or "TallyPrime is not connected.")) + company = next((row for row in (tally_status.get("companies") or []) if str(row.get("guid") or "").strip() == requested_guid), None) + if not company: + raise ValueError("The selected mapped Tally company is not currently open in TallyPrime. Open it in Tally and refresh the ERP page.") + company_name = str(company.get("name") or mapping.get("company_name") or "").strip() + masters = self.tally.fetch_accounting_masters(company_name) + self.store.record_tally_status(client_id, tally_status) + sync = self.store.replace_master_snapshot(client_id, mapping={**mapping, "company_name": company_name}, masters=masters, requested_by_user_id=int(requested_by_user_id) if requested_by_user_id not in (None, "") else None) + self.logger.info("Tally master sync completed client_id=%s company=%s rows=%s", client_id, company_name, sync.get("rows_processed")) + return {"synced": True, "sync": sync, "accounting": self.store.snapshot(client_id), "tally": tally_status, "agent": self._agent_info()} diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/tally.py b/app/modules/documents/local_agent_runtime/erp_local_agent/tally.py index 227506d..0a4e9e1 100644 --- a/app/modules/documents/local_agent_runtime/erp_local_agent/tally.py +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/tally.py @@ -6,11 +6,12 @@ import re import urllib.error import urllib.request import xml.etree.ElementTree as ET +from typing import Iterable DEFAULT_TALLY_HOST = "127.0.0.1" DEFAULT_TALLY_PORT = 9000 -DEFAULT_TALLY_TIMEOUT_SECONDS = 8 +DEFAULT_TALLY_TIMEOUT_SECONDS = 20 class TallyConnectionError(RuntimeError): @@ -24,14 +25,35 @@ def _clean_xml_response(xml_text: str) -> str: return re.sub(r"[\x00-\x08\x0B\x0C\x0E-\x1F]", "", xml_text) +def _tag(element: ET.Element) -> str: + return str(element.tag).split("}")[-1].upper() + + def _child_text(element: ET.Element, tag_name: str) -> str: wanted = tag_name.upper() for child in list(element): - if str(child.tag).split("}")[-1].upper() == wanted: + if _tag(child) == wanted: return (child.text or "").strip() return "" +def _first_text(element: ET.Element, names: Iterable[str]) -> str: + wanted = {str(name).upper() for name in names} + for node in element.iter(): + if _tag(node) in wanted and node.text: + return node.text.strip() + return "" + + +def _to_number(value: str) -> float: + text = str(value or "").replace(",", "").strip() + text = re.sub(r"[^0-9.\-]", "", text) + try: + return float(text or 0) + except Exception: + return 0.0 + + @dataclass(frozen=True) class TallyCompany: name: str @@ -43,7 +65,69 @@ class TallyCompany: class TallyLiveConnector: - """Read-only TallyPrime XML/HTTP discovery connector.""" + """Read-only TallyPrime XML/HTTP connector used by the ERP Local Agent. + + Phase 3 adds master discovery for the currently loaded Tally companies. Every + master export is scoped with SVCURRENTCOMPANY and does not create or modify + any Tally data. + """ + + MASTER_SPECS = { + "groups": { + "collection": "ARRRAccountingGroups", + "type": "Group", + "tag": "GROUP", + "fetch": "Name,GUID,Parent,ReservedName,IsRevenue,IsDeemedPositive", + }, + "ledgers": { + "collection": "ARRRAccountingLedgers", + "type": "Ledger", + "tag": "LEDGER", + "fetch": "Name,GUID,Parent,OpeningBalance,ClosingBalance,IsBillWiseOn,TaxType,GSTApplicable,GSTRegistrationType", + }, + "voucher_types": { + "collection": "ARRRAccountingVoucherTypes", + "type": "VoucherType", + "tag": "VOUCHERTYPE", + "fetch": "Name,GUID,Parent,NumberingMethod,IsDeemedPositive,IsActive", + }, + "stock_groups": { + "collection": "ARRRAccountingStockGroups", + "type": "StockGroup", + "tag": "STOCKGROUP", + "fetch": "Name,GUID,Parent,BaseUnits,IsAddable", + }, + "stock_categories": { + "collection": "ARRRAccountingStockCategories", + "type": "StockCategory", + "tag": "STOCKCATEGORY", + "fetch": "Name,GUID,Parent", + }, + "stock_items": { + "collection": "ARRRAccountingStockItems", + "type": "StockItem", + "tag": "STOCKITEM", + "fetch": "Name,GUID,Parent,Category,BaseUnits,AdditionalUnits,OpeningBalance,OpeningValue,OpeningRate,GSTApplicable,GSTTypeOfSupply,HSNCode", + }, + "units": { + "collection": "ARRRAccountingUnits", + "type": "Unit", + "tag": "UNIT", + "fetch": "Name,GUID,OriginalName,IsSimpleUnit,BaseUnits,AdditionalUnits,Conversion", + }, + "cost_categories": { + "collection": "ARRRAccountingCostCategories", + "type": "CostCategory", + "tag": "COSTCATEGORY", + "fetch": "Name,GUID,IsRevenue,IsNonRevenue", + }, + "cost_centres": { + "collection": "ARRRAccountingCostCentres", + "type": "CostCentre", + "tag": "COSTCENTRE", + "fetch": "Name,GUID,Parent,Category", + }, + } def __init__(self, host: str = DEFAULT_TALLY_HOST, port: int = DEFAULT_TALLY_PORT, timeout: int = DEFAULT_TALLY_TIMEOUT_SECONDS): self.host = (host or DEFAULT_TALLY_HOST).strip() @@ -54,6 +138,10 @@ class TallyLiveConnector: def url(self) -> str: return f"http://{self.host}:{self.port}" + @staticmethod + def _xml_escape(value: str) -> str: + return html.escape(str(value or ""), quote=False) + def _post_xml(self, xml_text: str) -> str: request = urllib.request.Request( self.url, @@ -69,20 +157,23 @@ class TallyLiveConnector: f"Could not connect to TallyPrime at {self.url}. Open TallyPrime and ensure its HTTP/XML server is available on port {self.port}. Details: {exc}" ) from exc + def _static_variables(self, company_name: str = "") -> str: + base = "$$SysName:XML" + if not company_name: + return base + return base + f"{self._xml_escape(company_name)}" + def get_loaded_companies(self) -> list[TallyCompany]: xml = """
- 1 - Export - Collection - ARRRAccountingLoadedCompanies + 1ExportCollectionARRRAccountingLoadedCompanies
- - - $$SysName:XML - CompanyName,GUID,GSTRegistrationNumber - - + + $$SysName:XML + + CompanyName,GUID,GSTRegistrationNumber + +
""" return self._parse_companies(self._post_xml(xml)) @@ -109,6 +200,87 @@ class TallyLiveConnector: "error": str(exc), } + def export_master_collection(self, company_name: str, master_key: str) -> list[dict]: + key = str(master_key or "").strip() + spec = self.MASTER_SPECS.get(key) + if not spec: + raise ValueError(f"Unsupported Tally master collection: {key}") + if not str(company_name or "").strip(): + raise ValueError("Tally company name is required for master sync.") + + xml = f""" +
1ExportCollection{spec['collection']}
+ + {self._static_variables(company_name)} + + {spec['type']}{spec['fetch']} + + +
""" + raw = self._post_xml(xml) + return self._parse_master_rows(raw, spec["tag"]) + + def fetch_accounting_masters(self, company_name: str) -> dict[str, list[dict]]: + result: dict[str, list[dict]] = {} + for key in self.MASTER_SPECS: + result[key] = self.export_master_collection(company_name, key) + return result + + @staticmethod + def _parse_master_rows(xml_text: str, expected_tag: str) -> list[dict]: + cleaned = _clean_xml_response(xml_text) + if not cleaned.strip(): + return [] + try: + root = ET.fromstring(cleaned.encode("utf-8")) + except Exception as exc: + raise ValueError(f"Tally returned invalid XML while reading {expected_tag} masters: {exc}") from exc + + rows: list[dict] = [] + seen: set[tuple[str, str]] = set() + for element in root.iter(): + if _tag(element) != expected_tag.upper(): + continue + name = ( + element.attrib.get("NAME") + or element.attrib.get("name") + or _child_text(element, "NAME") + ).strip() + guid = _child_text(element, "GUID") + if not name: + continue + identity = (guid.upper(), name.upper()) + if identity in seen: + continue + seen.add(identity) + rows.append({ + "name": name, + "guid": guid, + "parent": _child_text(element, "PARENT"), + "reserved_name": _child_text(element, "RESERVEDNAME"), + "category": _child_text(element, "CATEGORY"), + "base_units": _child_text(element, "BASEUNITS"), + "additional_units": _child_text(element, "ADDITIONALUNITS"), + "original_name": _child_text(element, "ORIGINALNAME"), + "opening_balance": _to_number(_child_text(element, "OPENINGBALANCE")), + "closing_balance": _to_number(_child_text(element, "CLOSINGBALANCE")), + "opening_value": _to_number(_child_text(element, "OPENINGVALUE")), + "opening_rate": _child_text(element, "OPENINGRATE"), + "numbering_method": _child_text(element, "NUMBERINGMETHOD"), + "tax_type": _child_text(element, "TAXTYPE"), + "gst_applicable": _child_text(element, "GSTAPPLICABLE"), + "gst_registration_type": _child_text(element, "GSTREGISTRATIONTYPE"), + "gst_type_of_supply": _child_text(element, "GSTTYPEOFSUPPLY"), + "hsn_code": _first_text(element, ("HSNCODE", "GSTHSNCODE")), + "is_revenue": _child_text(element, "ISREVENUE"), + "is_deemed_positive": _child_text(element, "ISDEEMEDPOSITIVE"), + "is_billwise_on": _child_text(element, "ISBILLWISEON"), + "is_simple_unit": _child_text(element, "ISSIMPLEUNIT"), + "conversion": _child_text(element, "CONVERSION"), + "is_active": _child_text(element, "ISACTIVE"), + }) + return rows + @staticmethod def _parse_companies(xml_text: str) -> list[TallyCompany]: cleaned = _clean_xml_response(xml_text) @@ -127,7 +299,7 @@ class TallyLiveConnector: companies: list[TallyCompany] = [] seen: set[str] = set() for element in root.iter(): - if str(element.tag).split("}")[-1].upper() != "COMPANY": + if _tag(element) != "COMPANY": continue name = ( element.attrib.get("NAME")