From 2ecaa0b08851f58a94d5057dc4963533cde58f2c Mon Sep 17 00:00:00 2001 From: A R R R Associates Date: Tue, 18 Aug 2026 14:39:19 +0530 Subject: [PATCH] Complete Phase 1 Tally UI and local agent integration --- app/core/startup.py | 3 + app/modules/accounting/agent_bridge.py | 132 +++++++++++++ .../templates/accounting/tally.html | 105 +++++++++++ app/modules/accounting/ui.py | 174 ++++++++++++++++++ app/modules/core/rbac/permissions_registry.py | 4 + app/modules/documents/agent_package.py | 2 +- .../erp_local_agent/__init__.py | 2 +- .../erp_local_agent/accounting_store.py | 170 +++++++++++++++++ .../erp_local_agent/commands.py | 79 ++++++++ .../erp_local_agent/sync.py | 2 +- .../erp_local_agent/tally.py | 143 ++++++++++++++ .../erp_local_agent/tunnel.py | 8 + app/modules/documents/ui.py | 10 +- app/ui/app.py | 2 + .../components/partner_navigation_v2.html | 8 +- 15 files changed, 837 insertions(+), 7 deletions(-) create mode 100644 app/modules/accounting/agent_bridge.py create mode 100644 app/modules/accounting/templates/accounting/tally.html create mode 100644 app/modules/accounting/ui.py create mode 100644 app/modules/documents/local_agent_runtime/erp_local_agent/accounting_store.py create mode 100644 app/modules/documents/local_agent_runtime/erp_local_agent/commands.py create mode 100644 app/modules/documents/local_agent_runtime/erp_local_agent/tally.py diff --git a/app/core/startup.py b/app/core/startup.py index 0757b29..36cc94c 100644 --- a/app/core/startup.py +++ b/app/core/startup.py @@ -333,6 +333,9 @@ ROLE_PERMISSION_MAP = { "documents.upload", "documents.download", "documents.delete", + "accounting.tally.view", + "accounting.tally.connect", + "accounting.act.initialize", "clients.view.own_only", "employees.dashboard.view", "employees.view", diff --git a/app/modules/accounting/agent_bridge.py b/app/modules/accounting/agent_bridge.py new file mode 100644 index 0000000..9fd0617 --- /dev/null +++ b/app/modules/accounting/agent_bridge.py @@ -0,0 +1,132 @@ +from __future__ import annotations + +import json +import os +from pathlib import Path +import tempfile +import time +from typing import Any +from uuid import uuid4 + + +COMMAND_TIMEOUT_SECONDS = 25 +COMMAND_MAX_AGE_SECONDS = 120 + + +def _bus_root() -> Path: + configured = (os.getenv("ERP_AGENT_COMMAND_BUS_ROOT") or "").strip() + if configured: + root = Path(configured).expanduser().resolve() + else: + root = Path(tempfile.gettempdir()) / "arrr_erp_agent_command_bus" + root.mkdir(parents=True, exist_ok=True) + return root + + +def _safe_node_code(node_code: str) -> str: + return "".join(ch for ch in str(node_code or "") if ch.isalnum() or ch in "-_.")[:100] + + +def _node_dir(node_code: str) -> Path: + safe = _safe_node_code(node_code) + if not safe: + raise ValueError("Storage node code is required.") + path = _bus_root() / safe + path.mkdir(parents=True, exist_ok=True) + return path + + +def _atomic_json_write(path: Path, payload: dict[str, Any]) -> None: + temp = path.with_suffix(path.suffix + ".tmp") + temp.write_text(json.dumps(payload, ensure_ascii=False, separators=(",", ":")), encoding="utf-8") + temp.replace(path) + + +def _read_json(path: Path) -> dict[str, Any] | None: + try: + data = json.loads(path.read_text(encoding="utf-8")) + return data if isinstance(data, dict) else None + except Exception: + return None + + +def _cleanup(node_code: str) -> None: + now = time.time() + for path in _node_dir(node_code).glob("*.json"): + try: + if now - path.stat().st_mtime > COMMAND_MAX_AGE_SECONDS: + path.unlink(missing_ok=True) + except Exception: + pass + + +def enqueue_agent_command(node_code: str, action: str, payload: dict[str, Any] | None = None) -> str: + _cleanup(node_code) + command_id = uuid4().hex + body = { + "command_id": command_id, + "action": str(action or "").strip(), + "payload": payload or {}, + "created_at_epoch": time.time(), + } + _atomic_json_write(_node_dir(node_code) / f"command_{command_id}.json", body) + return command_id + + +def list_pending_agent_commands(node_code: str, limit: int = 10) -> list[dict[str, Any]]: + _cleanup(node_code) + base = _node_dir(node_code) + items: list[dict[str, Any]] = [] + for path in sorted(base.glob("command_*.json"), key=lambda p: p.stat().st_mtime): + command_id = path.stem.removeprefix("command_") + if (base / f"response_{command_id}.json").exists(): + continue + payload = _read_json(path) + if payload: + items.append(payload) + if len(items) >= max(1, int(limit)): + break + return items + + +def record_agent_command_result(node_code: str, result: dict[str, Any]) -> None: + command_id = str(result.get("command_id") or "").strip() + if not command_id: + return + response = { + "command_id": command_id, + "ok": bool(result.get("ok")), + "result": result.get("result"), + "error": result.get("error"), + "agent_time_utc": result.get("agent_time_utc"), + "received_at_epoch": time.time(), + } + _atomic_json_write(_node_dir(node_code) / f"response_{command_id}.json", response) + + +def wait_for_agent_result(node_code: str, command_id: str, timeout_seconds: int = COMMAND_TIMEOUT_SECONDS) -> dict[str, Any]: + base = _node_dir(node_code) + response_path = base / f"response_{command_id}.json" + command_path = base / f"command_{command_id}.json" + deadline = time.monotonic() + max(1, int(timeout_seconds)) + while time.monotonic() < deadline: + result = _read_json(response_path) + if result: + try: + response_path.unlink(missing_ok=True) + command_path.unlink(missing_ok=True) + except Exception: + pass + return result + time.sleep(0.2) + raise TimeoutError("ERP Local Agent did not respond before the command timeout.") + + +def request_agent_command( + node_code: str, + action: str, + payload: dict[str, Any] | None = None, + timeout_seconds: int = COMMAND_TIMEOUT_SECONDS, +) -> dict[str, Any]: + command_id = enqueue_agent_command(node_code, action, payload) + return wait_for_agent_result(node_code, command_id, timeout_seconds=timeout_seconds) diff --git a/app/modules/accounting/templates/accounting/tally.html b/app/modules/accounting/templates/accounting/tally.html new file mode 100644 index 0000000..ee72ae2 --- /dev/null +++ b/app/modules/accounting/templates/accounting/tally.html @@ -0,0 +1,105 @@ +{% extends "ui/templates/base/layout.html" %} +{% block content %} +
+
+
+

Tools · Accounting

+

Tally Connection

+

Read-only Phase 1 connection through the existing ERP Local Agent. Tally port 9000 is never exposed to the internet.

+
+ Test Tally Connection +
+ + {% if initialized %} +
Client accounting storage was initialized successfully.
+ {% endif %} + {% if command_error %} +
{{ command_error }}
+ {% endif %} + +
+
+
ERP Local Agent
+
{% if agent_online %}Connected{% else %}Offline{% endif %}
+
{{ storage_node.node_name if storage_node else 'No active branch agent' }}
+
+
+
Tally Module
+ {% set agent = live_result.agent if live_result else None %} +
{{ 'Available' if agent and agent.tally_capability else ('Check connection' if agent_online else 'Unavailable') }}
+
Agent {{ agent.version if agent and agent.version else '-' }}
+
+
+
TallyPrime
+ {% set tally = live_result.tally if live_result else None %} +
{{ 'Connected' if tally and tally.connected else ('Not connected' if tally else 'Not checked') }}
+
{{ tally.url if tally and tally.url else '127.0.0.1:9000' }}
+
+
+
Loaded Companies
+
{{ tally.company_count if tally else '-' }}
+
Read-only discovery
+
+
+ +
+
+
+ + +
+ {% if selected_client %} +
+ + + +
+ {% endif %} +
+

Only clients assigned to the logged-in Partner in the active branch are shown.

+
+ + {% if tally %} +
+
+

Loaded Tally Companies

+

Company name, GUID and GSTIN returned directly by the local TallyPrime instance.

+
+ {% if tally.companies %} +
+ + + {% for company in tally.companies %}{% endfor %} +
CompanyGUIDGSTIN
{{ company.name }}{{ company.guid or '-' }}{{ company.gstin or '-' }}
+
+ {% else %} +
TallyPrime responded, but no loaded company was returned.
+ {% endif %} +
+ {% endif %} + + {% if selected_client %} + {% set accounting = live_result.accounting if live_result else None %} +
+
+

Accounting Storage

{{ selected_client.client_name }}

+ {{ 'Initialized' if accounting and accounting.exists else 'Not initialized / not checked' }} +
+ {% if accounting and accounting.exists %} +
+
.act Database
{{ accounting.db_path }}
+
Schema Version
{{ accounting.metadata.schema_version or '1' }}
+
+ {% endif %} +
+ {% endif %} + +
Phase 1 is read-only. No ledger, voucher, inventory or write-back operation is performed.
+
+{% endblock %} diff --git a/app/modules/accounting/ui.py b/app/modules/accounting/ui.py new file mode 100644 index 0000000..80ef990 --- /dev/null +++ b/app/modules/accounting/ui.py @@ -0,0 +1,174 @@ +from __future__ import annotations + +from datetime import datetime, timezone +from urllib.parse import quote + +from fastapi import APIRouter, Form, Request +from fastapi.responses import RedirectResponse +from sqlalchemy import select + +from app.core.db.common import CommonSessionLocal +from app.core.security.csrf import get_or_create_csrf_token, validate_csrf +from app.core.security.session_auth import get_current_user +from app.core.templating import templates +from app.modules.clients.models import Client +from app.modules.core.rbac.deps import get_user_permissions, get_user_roles +from app.modules.core.rbac.permission_guard import require_permission +from app.modules.documents.services import build_document_scope, get_active_storage_node_for_branch +from app.modules.accounting.agent_bridge import request_agent_command + + +router = APIRouter(prefix="/tools/tally", tags=["accounting-tally-ui"]) + + +def _denied(): + from app.core.http_responses import ui_access_denied + return ui_access_denied() + + +def _require_partner(request: Request, db, permission: str): + user = get_current_user(request, db) + if not user: + return None, RedirectResponse(url="/login", status_code=303) + roles = set(get_user_roles(db, user.id)) + if "Partner" not in roles: + return None, _denied() + try: + require_permission(db, user, permission) + except Exception: + return None, _denied() + return user, None + + +def _visible_clients(db, request: Request, user): + scope = build_document_scope(request, db, user) + stmt = select(Client).where(Client.tenant_id == scope.tenant_id, Client.partner_id == user.id) + if scope.branch_id is not None: + stmt = stmt.where(Client.branch_id == scope.branch_id) + return db.execute(stmt.order_by(Client.client_name.asc(), Client.id.asc())).scalars().all(), scope + + +def _find_visible_client(db, request: Request, user, client_id: int): + clients, scope = _visible_clients(db, request, user) + client = next((row for row in clients if int(row.id) == int(client_id)), None) + return client, clients, scope + + +def _node_online(node) -> bool: + if not node or not node.last_seen_at_utc: + return False + seen = node.last_seen_at_utc + if seen.tzinfo is None: + seen = seen.replace(tzinfo=timezone.utc) + return (datetime.now(timezone.utc) - seen).total_seconds() <= 180 + + +def _render(request: Request, db, user, **context): + base = { + "request": request, + "current_user": user, + "current_user_roles": get_user_roles(db, user.id), + "current_user_permissions": get_user_permissions(db, user.id), + "csrf_token": get_or_create_csrf_token(request), + } + base.update(context) + return templates.TemplateResponse( + "modules/accounting/templates/accounting/tally.html", + base, + ) + + +@router.get("") +def tally_tool(request: Request, client_id: int | None = None, refresh: int = 0, initialized: int = 0, error: str = ""): + db = CommonSessionLocal() + try: + user, response = _require_partner(request, db, "accounting.tally.view") + if response: + return response + + clients, scope = _visible_clients(db, request, user) + selected_client = next((row for row in clients if client_id and int(row.id) == int(client_id)), None) + node = get_active_storage_node_for_branch(db, scope.tenant_id, scope.branch_id) + online = _node_online(node) + + live_result = None + command_error = error or "" + if refresh: + try: + require_permission(db, user, "accounting.tally.connect") + except Exception: + return _denied() + if refresh and node and online: + payload = {} + if selected_client: + payload = { + "client_id": int(selected_client.id), + "client_name": selected_client.client_name, + } + try: + response_data = request_agent_command(node.node_code, "phase1_status", payload, timeout_seconds=20) + if response_data.get("ok"): + live_result = response_data.get("result") or {} + else: + command_error = str(response_data.get("error") or "Local agent command failed.") + except Exception as exc: + command_error = str(exc) + + return _render( + request, + db, + user, + title="Tally Connection", + clients=clients, + selected_client=selected_client, + storage_node=node, + agent_online=online, + live_result=live_result, + initialized=bool(initialized), + command_error=command_error, + ) + finally: + db.close() + + +@router.post("/initialize") +def initialize_accounting_storage( + request: Request, + client_id: int = Form(...), + csrf_token: str = Form(...), +): + validate_csrf(request, csrf_token) + db = CommonSessionLocal() + try: + user, response = _require_partner(request, db, "accounting.act.initialize") + 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_initialize", + {"client_id": int(client.id), "client_name": client.client_name}, + timeout_seconds=20, + ) + if not result.get("ok"): + raise RuntimeError(str(result.get("error") or "Accounting storage initialization 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&initialized=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 8d81e68..b052473 100644 --- a/app/modules/core/rbac/permissions_registry.py +++ b/app/modules/core/rbac/permissions_registry.py @@ -147,6 +147,10 @@ PERMISSIONS = { "documents.delete": "Archive Engagement Documents", "documents.audit.view": "View Document Access Logs", + "accounting.tally.view": "View Tally Accounting Tool", + "accounting.tally.connect": "Connect to Tally Through Local Agent", + "accounting.act.initialize": "Initialize Client Accounting ACT Storage", + "notice_cases.view": "View Notice and Case Management", "notice_cases.create": "Create Notices and Cases", "notice_cases.edit": "Edit Notices and Cases", diff --git a/app/modules/documents/agent_package.py b/app/modules/documents/agent_package.py index 543bb43..84bc9f5 100644 --- a/app/modules/documents/agent_package.py +++ b/app/modules/documents/agent_package.py @@ -4,7 +4,7 @@ import io from pathlib import Path import zipfile -ERP_LOCAL_AGENT_VERSION = "1.0.0" +ERP_LOCAL_AGENT_VERSION = "1.1.0" ERP_LOCAL_AGENT_NAME = "ERP Local Agent" RUNTIME_ROOT = Path(__file__).resolve().parent / "local_agent_runtime" 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 212ad5f..c94f386 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.0.0" +__version__ = "1.1.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 new file mode 100644 index 0000000..0814a75 --- /dev/null +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/accounting_store.py @@ -0,0 +1,170 @@ +from __future__ import annotations + +from datetime import datetime, timezone +import json +from pathlib import Path +import sqlite3 +from typing import Sequence + + +SCHEMA_VERSION = "1" + + +def _utc_now_iso() -> str: + return datetime.now(timezone.utc).isoformat() + + +class LocalAccountingStore: + """Client-scoped SQLite .act storage under the existing branch storage root.""" + + def __init__(self, storage_root: Path): + self.root = Path(storage_root).resolve() / "Accounting" + self.root.mkdir(parents=True, exist_ok=True) + + @staticmethod + def _client_key(client_id: int) -> str: + return f"client_{int(client_id):08d}" + + def client_dir(self, client_id: int) -> Path: + return self.root / self._client_key(client_id) + + def db_path(self, client_id: int) -> Path: + key = self._client_key(client_id) + return self.client_dir(client_id) / f"{key}.act" + + def exists(self, client_id: int) -> bool: + return self.db_path(client_id).is_file() + + def connect(self, client_id: int): + path = self.db_path(client_id) + path.parent.mkdir(parents=True, exist_ok=True) + db = sqlite3.connect(path) + db.row_factory = sqlite3.Row + return db + + def initialize(self, client_id: int, client_name: str = "") -> Path: + path = self.db_path(client_id) + with self.connect(client_id) as db: + db.executescript( + """ + PRAGMA journal_mode=WAL; + PRAGMA foreign_keys=ON; + 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 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, + 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 + ); + 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 + ); + 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 + ); + """ + ) + now = _utc_now_iso() + meta = { + "schema_version": SCHEMA_VERSION, + "client_id": str(int(client_id)), + "client_name": str(client_name or "").strip(), + "storage_kind": "client_accounting_act", + } + for key, value in meta.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: + if not self.exists(client_id): + return + 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: + path = self.db_path(client_id) + if not path.is_file(): + return {"exists": False, "db_path": str(path), "metadata": {}, "latest_connection": None} + with self.connect(client_id) as db: + meta_rows = db.execute("SELECT key, value FROM act_meta ORDER BY key").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 { + "exists": True, + "db_path": str(path), + "metadata": {row["key"]: row["value"] for row in meta_rows}, + "latest_connection": dict(latest) if latest else None, + } 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 new file mode 100644 index 0000000..a1eb803 --- /dev/null +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/commands.py @@ -0,0 +1,79 @@ +from __future__ import annotations + +from datetime import datetime, timezone +from typing import Any + +from . import __version__ +from .accounting_store import LocalAccountingStore +from .tally import TallyLiveConnector + + +class AgentCommandProcessor: + def __init__(self, config, logger): + self.config = config + self.logger = logger + self.store = LocalAccountingStore(config.storage_root) + self.tally = TallyLiveConnector() + + def process(self, command: dict[str, Any]) -> dict[str, Any]: + command_id = str(command.get("command_id") or "").strip() + action = str(command.get("action") or "").strip() + payload = command.get("payload") or {} + result: dict[str, Any] | None = None + error: str | None = None + ok = False + try: + if action == "tally_status": + result = self._status(payload) + elif action == "phase1_status": + result = self._status(payload) + elif action == "accounting_initialize": + result = self._initialize(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(), + } + + def _status(self, payload: dict[str, Any]) -> dict[str, Any]: + tally_status = self.tally.status() + client_id = payload.get("client_id") + accounting = None + if client_id not in (None, ""): + client_id = int(client_id) + if self.store.exists(client_id): + self.store.record_tally_status(client_id, tally_status) + accounting = self.store.snapshot(client_id) + return { + "agent": { + "name": "ERP Local Agent", + "version": __version__, + "tally_capability": True, + "accounting_act_capability": True, + }, + "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() + path = self.store.initialize(client_id, client_name) + 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": {"name": "ERP Local Agent", "version": __version__, "tally_capability": True, "accounting_act_capability": True}, + } diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/sync.py b/app/modules/documents/local_agent_runtime/erp_local_agent/sync.py index a0428ed..07f0a84 100644 --- a/app/modules/documents/local_agent_runtime/erp_local_agent/sync.py +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/sync.py @@ -49,7 +49,7 @@ class StorageAgent: "free_bytes": free, "agent_version": __version__, "agent_name": "ERP Local Agent", - "capabilities": ["storage"], + "capabilities": ["storage", "tally", "accounting_act"], "agent_time_utc": datetime.now(timezone.utc).isoformat(), } self.client.heartbeat(payload) 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 new file mode 100644 index 0000000..227506d --- /dev/null +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/tally.py @@ -0,0 +1,143 @@ +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): + pass + + +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) + return re.sub(r"[\x00-\x08\x0B\x0C\x0E-\x1F]", "", 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 TallyPrime XML/HTTP discovery connector.""" + + 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() + 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}" + + 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}. Open TallyPrime and ensure its HTTP/XML server is available on port {self.port}. Details: {exc}" + ) from exc + + def get_loaded_companies(self) -> list[TallyCompany]: + xml = """ +
+ 1 + Export + Collection + ARRRAccountingLoadedCompanies +
+ + + $$SysName:XML + CompanyName,GUID,GSTRegistrationNumber + + +
""" + return self._parse_companies(self._post_xml(xml)) + + 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: + 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 diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/tunnel.py b/app/modules/documents/local_agent_runtime/erp_local_agent/tunnel.py index 790b5e3..c3bc2ba 100644 --- a/app/modules/documents/local_agent_runtime/erp_local_agent/tunnel.py +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/tunnel.py @@ -7,6 +7,7 @@ from typing import Any import websockets from .sync import StorageAgent +from .commands import AgentCommandProcessor class StorageAgentTunnel: @@ -14,6 +15,7 @@ class StorageAgentTunnel: self.agent = agent self.config = agent.config self.logger = agent.logger + self.command_processor = AgentCommandProcessor(self.config, self.logger) async def run_forever(self) -> None: self.logger.info("Starting tunnel mode for node=%s", self.config.node_code) @@ -47,12 +49,18 @@ class StorageAgentTunnel: if msg_type == "sync": jobs = message.get("jobs") or [] requests_ = message.get("download_requests") or message.get("requests") or [] + commands = message.get("commands") or [] await asyncio.to_thread(self.agent.process_storage_jobs_from_payload, jobs) await asyncio.to_thread(self.agent.process_download_requests_from_payload, requests_) + for command in commands: + result = await asyncio.to_thread(self.command_processor.process, command) + await websocket.send(self._json_text(result)) await websocket.send(self._json_text({ "type": "agent_status", "jobs_seen": len(jobs), "requests_seen": len(requests_), + "commands_seen": len(commands), + "capabilities": ["storage", "tally", "accounting_act"], "agent_time_utc": self._now(), })) continue diff --git a/app/modules/documents/ui.py b/app/modules/documents/ui.py index 6855c74..d386f4c 100644 --- a/app/modules/documents/ui.py +++ b/app/modules/documents/ui.py @@ -1,4 +1,4 @@ -from __future__ import annotations +from __future__ import annotations from datetime import date, datetime, timezone import asyncio @@ -82,6 +82,7 @@ from app.modules.core.tenancy.models import Branch, Tenant from app.modules.core.tenancy.year_control import is_row_financial_year_locked from app.modules.documents.agent_package import ERP_LOCAL_AGENT_VERSION, build_agent_env, build_preconfigured_agent_zip, build_agent_update_zip from app.modules.documents.models import BranchStorageNode +from app.modules.accounting.agent_bridge import list_pending_agent_commands, record_agent_command_result from app.modules.services.task_documents import ( get_task_with_subscription, @@ -1693,7 +1694,8 @@ def _storage_agent_sync_payload(db, node): jobs += _permanent_storage_jobs_payload(list_pending_permanent_storage_jobs(db, node)) requests_ = _normal_download_requests_payload(list_pending_download_requests(db, node)) requests_ += _permanent_download_requests_payload(list_pending_permanent_download_requests(db, node)) - return {"ok": True, "jobs": jobs, "download_requests": requests_, "requests": requests_} + commands = list_pending_agent_commands(node.node_code, limit=10) + return {"ok": True, "jobs": jobs, "download_requests": requests_, "requests": requests_, "commands": commands} @router.get("/storage-agent/jobs/pending") @@ -1837,7 +1839,9 @@ async def storage_agent_tunnel(websocket: WebSocket): await websocket.send_json({"ok": False, "type": "error", "error": "node_deactivated_or_invalid"}) await websocket.close(code=1008) return - if message and message.get("type") in {"heartbeat", "agent_status"}: + if message and message.get("type") == "command_result": + record_agent_command_result(node.node_code, message) + if message and message.get("type") in {"heartbeat", "agent_status", "command_result"}: try: node.storage_mode = "tunnel" except Exception: diff --git a/app/ui/app.py b/app/ui/app.py index 569bba9..0f6bb7d 100644 --- a/app/ui/app.py +++ b/app/ui/app.py @@ -36,6 +36,7 @@ from app.modules.firm_admin_dashboard.ui import router as firm_admin_dashboard_r from app.modules.aqmm_dashboard.ui import router as aqmm_dashboard_router from app.modules.peer_review_export.ui import router as peer_review_export_router from app.modules.bank_statement_analyzer.ui import router as bank_statement_analyzer_router +from app.modules.accounting.ui import router as accounting_ui_router from app.modules.registrations.ui import router as registrations_ui_router from app.modules.credential_vault.ui import router as credential_vault_ui_router from app.modules.client_identity.ui import router as client_identity_ui_router @@ -58,6 +59,7 @@ def mount_ui(app: FastAPI) -> None: app.include_router(aqmm_dashboard_router) app.include_router(peer_review_export_router) app.include_router(bank_statement_analyzer_router) + app.include_router(accounting_ui_router) app.include_router(work_tracker_ui_router) app.include_router(billing_ui_router) app.include_router(platform_billing_ui_router) diff --git a/app/ui/templates/components/partner_navigation_v2.html b/app/ui/templates/components/partner_navigation_v2.html index 9d09165..29fe80c 100644 --- a/app/ui/templates/components/partner_navigation_v2.html +++ b/app/ui/templates/components/partner_navigation_v2.html @@ -1,4 +1,4 @@ -{% if navigation_render_location_v2|default('page') == 'global-host' or not global_role_navigation_host_v2|default(false) %} +{% if navigation_render_location_v2|default('page') == 'global-host' or not global_role_navigation_host_v2|default(false) %} {# Partner operations navigation. Presentation only: all links use existing production routes and route-level permissions remain authoritative. @@ -212,6 +212,11 @@ 'visible': true, 'active': (_partner_path.startswith('/partner/dashboard') and _partner_tab == 'wizards') or _partner_path.startswith('/tools/bank-statement-analyzer') or _partner_path.startswith('/documents/storage-nodes') or _partner_path.startswith('/documents/storage-jobs'), 'children': [ + { + 'label': 'Tally', + 'url': '/tools/tally', + 'active': _partner_path.startswith('/tools/tally') + }, { 'label': 'Bank Analyzer', 'url': '/tools/bank-statement-analyzer', @@ -297,3 +302,4 @@ })(); {% endif %} +