diff --git a/app/modules/accounting/__init__.py b/app/modules/accounting/__init__.py new file mode 100644 index 0000000..5d1b92a --- /dev/null +++ b/app/modules/accounting/__init__.py @@ -0,0 +1,18 @@ +"""Accounting Phase 1 foundation. + +This package is intentionally additive. Phase 1 provides client-scoped .act +SQLite storage plus read-only local TallyPrime discovery. ERP UI/tunnel +registration is wired only after matching the current production router/agent +source so existing features are not disturbed. +""" + +from .act_store import AccountingActStore, AccountingActStoreError +from .tally_connector import TallyCompany, TallyConnectionError, TallyLiveConnector + +__all__ = [ + "AccountingActStore", + "AccountingActStoreError", + "TallyCompany", + "TallyConnectionError", + "TallyLiveConnector", +] diff --git a/app/modules/accounting/act_store.py b/app/modules/accounting/act_store.py new file mode 100644 index 0000000..ed892df --- /dev/null +++ b/app/modules/accounting/act_store.py @@ -0,0 +1,248 @@ +from __future__ import annotations + +from contextlib import contextmanager +from datetime import datetime, timezone +from pathlib import Path +import json +import os +import sqlite3 +from typing import Iterator, Sequence + +ACT_SCHEMA_VERSION = 1 + + +class AccountingActStoreError(RuntimeError): + pass + + +def _utc_now_iso() -> str: + return datetime.now(timezone.utc).replace(microsecond=0).isoformat() + + +def _safe_client_id(client_id: int) -> int: + value = int(client_id) + if value <= 0: + raise AccountingActStoreError("client_id must be a positive integer") + return value + + +class AccountingActStore: + """Client-scoped SQLite storage for ERP accounting data. + + The `.act` extension is intentional; the underlying file format is SQLite. + Phase 1 stores metadata, discovered Tally companies, connection history and + sync-run control records only. Transaction/master tables arrive in later + phases. + """ + + def __init__(self, root: str | Path) -> None: + self.root = Path(root).expanduser().resolve() + + def client_dir(self, client_id: int) -> Path: + cid = _safe_client_id(client_id) + return self.root / "Accounting" / f"client_{cid:08d}" + + def db_path(self, client_id: int) -> Path: + cid = _safe_client_id(client_id) + return self.client_dir(cid) / f"client_{cid:08d}.act" + + @contextmanager + def connect(self, client_id: int) -> Iterator[sqlite3.Connection]: + path = self.db_path(client_id) + path.parent.mkdir(parents=True, exist_ok=True) + connection = sqlite3.connect(path, timeout=30) + connection.row_factory = sqlite3.Row + try: + connection.execute("PRAGMA foreign_keys = ON") + connection.execute("PRAGMA journal_mode = WAL") + connection.execute("PRAGMA synchronous = NORMAL") + yield connection + connection.commit() + except Exception: + connection.rollback() + raise + finally: + connection.close() + + def initialize( + self, + client_id: int, + *, + tenant_id: int | None = None, + client_name: str = "", + 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 + ); + + 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 1, + UNIQUE(tally_guid, company_name) + ); + + CREATE TABLE IF NOT EXISTS tally_company_mapping ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + client_id INTEGER NOT NULL, + registration_id INTEGER, + tally_guid TEXT NOT NULL DEFAULT '', + company_name TEXT NOT NULL DEFAULT '', + gstin TEXT NOT NULL DEFAULT '', + is_active INTEGER NOT NULL DEFAULT 1, + mapped_at_utc TEXT NOT NULL, + mapped_by_user_id INTEGER + ); + + 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, + payload_json TEXT NOT NULL DEFAULT '{}' + ); + + 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, + requested_by_user_id INTEGER, + records_received INTEGER NOT NULL DEFAULT 0, + error_message TEXT, + details_json TEXT NOT NULL DEFAULT '{}' + ); + + CREATE INDEX IF NOT EXISTS ix_tally_companies_last_seen + ON tally_companies(last_seen_at_utc); + CREATE INDEX IF NOT EXISTS ix_tally_connection_history_checked + ON tally_connection_history(checked_at_utc); + CREATE INDEX IF NOT EXISTS ix_tally_sync_runs_started + ON tally_sync_runs(started_at_utc); + """ + ) + now = _utc_now_iso() + metadata = { + "schema_version": str(ACT_SCHEMA_VERSION), + "client_id": str(_safe_client_id(client_id)), + "tenant_id": "" if tenant_id is None else str(int(tenant_id)), + "client_name": str(client_name or ""), + "created_by_user_id": "" if created_by_user_id is None else str(int(created_by_user_id)), + "storage_kind": "ERP_ACCOUNTING_ACT_SQLITE", + } + for key, value in metadata.items(): + db.execute( + """ + INSERT INTO act_meta(key, value, updated_at_utc) + VALUES (?, ?, ?) + ON CONFLICT(key) DO UPDATE SET + value=excluded.value, + updated_at_utc=excluded.updated_at_utc + """, + (key, value, now), + ) + return path + + def record_tally_status(self, client_id: int, status: dict) -> None: + self.initialize(client_id) + now = _utc_now_iso() + companies: Sequence[dict] = status.get("companies") or [] + with self.connect(client_id) as db: + db.execute("UPDATE tally_companies SET is_currently_loaded = 0") + for company in companies: + name = str(company.get("name") or "").strip() + if not name: + continue + 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), + ).fetchone() + if existing: + db.execute( + """ + UPDATE tally_companies + SET gstin=?, last_seen_at_utc=?, is_currently_loaded=1 + WHERE id=? + """, + (gstin, now, int(existing["id"])), + ) + else: + db.execute( + """ + INSERT INTO tally_companies( + tally_guid, company_name, gstin, + first_seen_at_utc, last_seen_at_utc, is_currently_loaded + ) VALUES (?, ?, ?, ?, ?, 1) + """, + (guid, name, gstin, now, now), + ) + + 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=(",", ":")), + ), + ) + + def snapshot(self, client_id: int) -> dict: + 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() + 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() + return { + "db_path": str(self.db_path(client_id)), + "metadata": {row["key"]: row["value"] for row in meta_rows}, + "companies": [dict(row) for row in companies], + "latest_connection": dict(latest) if latest else None, + } + + +def default_accounting_root() -> Path: + env = os.getenv("AUDIT_ACCOUNTING_STORAGE_ROOT", "").strip() + if env: + return Path(env).expanduser().resolve() + return Path.cwd() / "data" / "accounting" diff --git a/app/modules/accounting/phase1_probe.py b/app/modules/accounting/phase1_probe.py new file mode 100644 index 0000000..219b788 --- /dev/null +++ b/app/modules/accounting/phase1_probe.py @@ -0,0 +1,41 @@ +from __future__ import annotations + +import argparse +import json + +from .act_store import AccountingActStore, default_accounting_root +from .tally_connector import TallyLiveConnector + + +def main() -> int: + parser = argparse.ArgumentParser(description="Phase 1 .act + local TallyPrime smoke test") + parser.add_argument("--client-id", type=int, required=True) + parser.add_argument("--client-name", default="") + parser.add_argument("--tenant-id", type=int) + parser.add_argument("--storage-root", default="") + parser.add_argument("--host", default="127.0.0.1") + parser.add_argument("--port", type=int, default=9000) + args = parser.parse_args() + + root = args.storage_root or str(default_accounting_root()) + store = AccountingActStore(root) + path = store.initialize( + args.client_id, + tenant_id=args.tenant_id, + client_name=args.client_name, + ) + + connector = TallyLiveConnector(args.host, args.port) + status = connector.status() + store.record_tally_status(args.client_id, status) + + print(json.dumps({ + "act_db": str(path), + "tally": status, + "snapshot": store.snapshot(args.client_id), + }, indent=2, ensure_ascii=False)) + return 0 if status.get("connected") else 2 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/app/modules/accounting/tally_connector.py b/app/modules/accounting/tally_connector.py new file mode 100644 index 0000000..49722bd --- /dev/null +++ b/app/modules/accounting/tally_connector.py @@ -0,0 +1,178 @@ +from __future__ import annotations + +from dataclasses import asdict, dataclass +import html +import re +import urllib.error +import urllib.request +import xml.etree.ElementTree as ET + +DEFAULT_TALLY_HOST = "127.0.0.1" +DEFAULT_TALLY_PORT = 9000 +DEFAULT_TALLY_TIMEOUT_SECONDS = 8 + + +class TallyConnectionError(RuntimeError): + """Raised when the local TallyPrime XML/HTTP endpoint cannot be reached.""" + + +def _clean_xml_response(xml_text: str) -> str: + if not xml_text: + return "" + xml_text = re.sub(r"&#(0?[0-8]|1[0-9]|2[0-9]|3[01]);", "", xml_text) + xml_text = re.sub(r"[\x00-\x08\x0B\x0C\x0E-\x1F]", "", xml_text) + return xml_text + + +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: + return (child.text or "").strip() + return "" + + +@dataclass(frozen=True) +class TallyCompany: + name: str + guid: str = "" + gstin: str = "" + + def as_dict(self) -> dict[str, str]: + return asdict(self) + + +class TallyLiveConnector: + """Read-only connector for a running local TallyPrime instance. + + Uses Tally XML-over-HTTP. Phase 1 performs only discovery/status reads and + never creates, updates or deletes data in Tally. + """ + + def __init__( + self, + host: str = DEFAULT_TALLY_HOST, + port: int = DEFAULT_TALLY_PORT, + timeout: int = DEFAULT_TALLY_TIMEOUT_SECONDS, + ) -> None: + self.host = (host or DEFAULT_TALLY_HOST).strip() + self.port = int(port or DEFAULT_TALLY_PORT) + self.timeout = max(1, int(timeout or DEFAULT_TALLY_TIMEOUT_SECONDS)) + + @property + 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, + data=xml_text.encode("utf-8"), + headers={"Content-Type": "application/xml; charset=utf-8"}, + method="POST", + ) + try: + with urllib.request.urlopen(request, timeout=self.timeout) as response: + return response.read().decode("utf-8", errors="replace") + except (urllib.error.URLError, TimeoutError, OSError) as exc: + raise TallyConnectionError( + f"Could not connect to TallyPrime at {self.url}. " + f"Open TallyPrime and enable its HTTP/XML server on port {self.port}. " + f"Details: {exc}" + ) from exc + + def get_loaded_companies(self) -> list[TallyCompany]: + xml = """ +
+ 1 + Export + Collection + ARRRAccountingLoadedCompanies +
+ + + + $$SysName:XML + + + + + Company + Name,GUID,GSTRegistrationNumber + + + + + +
""" + raw = self._post_xml(xml) + return self._parse_companies(raw) + + def status(self) -> dict: + try: + companies = self.get_loaded_companies() + return { + "connected": True, + "url": self.url, + "host": self.host, + "port": self.port, + "company_count": len(companies), + "companies": [company.as_dict() for company in companies], + "error": None, + } + except TallyConnectionError as exc: + return { + "connected": False, + "url": self.url, + "host": self.host, + "port": self.port, + "company_count": 0, + "companies": [], + "error": str(exc), + } + + @staticmethod + def _parse_companies(xml_text: str) -> list[TallyCompany]: + cleaned = _clean_xml_response(xml_text) + if not cleaned.strip(): + return [] + + try: + root = ET.fromstring(cleaned.encode("utf-8")) + except Exception: + # Keep the same conservative fallback behaviour as the working + # Office Automation connector: discover names even if Tally emits + # malformed XML around otherwise useful response data. + names: list[str] = [] + for match in re.finditer(r"]*>(.*?)", cleaned, flags=re.I | re.S): + value = re.sub(r"<.*?>", "", match.group(1)).strip() + if value and value not in names: + names.append(value) + return [TallyCompany(name=name) for name in names] + + companies: list[TallyCompany] = [] + seen: set[str] = set() + for element in root.iter(): + if str(element.tag).split("}")[-1].upper() != "COMPANY": + continue + name = ( + element.attrib.get("NAME") + or element.attrib.get("name") + or _child_text(element, "NAME") + or _child_text(element, "BASICCOMPANYFORMALNAME") + ).strip() + key = name.upper() + if not name or key in seen: + continue + seen.add(key) + companies.append( + TallyCompany( + name=name, + guid=_child_text(element, "GUID"), + gstin=_child_text(element, "GSTREGISTRATIONNUMBER"), + ) + ) + return companies