Complete Phase 1 Tally UI and local agent integration

This commit is contained in:
A R R R Associates
2026-08-18 14:39:19 +05:30
parent b262861608
commit 2ecaa0b088
15 changed files with 837 additions and 7 deletions
+1 -1
View File
@@ -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"
@@ -1,2 +1,2 @@
__version__ = "1.0.0"
__version__ = "1.1.0"
AGENT_NAME = "ERP Local Agent"
@@ -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,
}
@@ -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},
}
@@ -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)
@@ -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 = """<ENVELOPE>
<HEADER>
<VERSION>1</VERSION>
<TALLYREQUEST>Export</TALLYREQUEST>
<TYPE>Collection</TYPE>
<ID>ARRRAccountingLoadedCompanies</ID>
</HEADER>
<BODY>
<DESC>
<STATICVARIABLES><SVEXPORTFORMAT>$$SysName:XML</SVEXPORTFORMAT></STATICVARIABLES>
<TDL><TDLMESSAGE><COLLECTION NAME="ARRRAccountingLoadedCompanies" ISMODIFY="No"><TYPE>Company</TYPE><FETCH>Name,GUID,GSTRegistrationNumber</FETCH></COLLECTION></TDLMESSAGE></TDL>
</DESC>
</BODY>
</ENVELOPE>"""
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"<NAME[^>]*>(.*?)</NAME>", 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
@@ -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
+7 -3
View File
@@ -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: