Add Phase 3 Tally accounting master synchronization

This commit is contained in:
A R R R Associates
2026-08-19 16:11:03 +05:30
parent 2f061d29b4
commit 93ae6caa54
11 changed files with 586 additions and 337 deletions
@@ -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.
@@ -1,2 +1,2 @@
__version__ = "1.3.0"
__version__ = "1.4.0"
AGENT_NAME = "ERP Local Agent"
@@ -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,
}
@@ -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()}
@@ -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 = "<SVEXPORTFORMAT>$$SysName:XML</SVEXPORTFORMAT>"
if not company_name:
return base
return base + f"<SVCURRENTCOMPANY>{self._xml_escape(company_name)}</SVCURRENTCOMPANY>"
def get_loaded_companies(self) -> list[TallyCompany]:
xml = """<ENVELOPE>
<HEADER>
<VERSION>1</VERSION>
<TALLYREQUEST>Export</TALLYREQUEST>
<TYPE>Collection</TYPE>
<ID>ARRRAccountingLoadedCompanies</ID>
<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>
<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))
@@ -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"""<ENVELOPE>
<HEADER><VERSION>1</VERSION><TALLYREQUEST>Export</TALLYREQUEST><TYPE>Collection</TYPE><ID>{spec['collection']}</ID></HEADER>
<BODY><DESC>
<STATICVARIABLES>{self._static_variables(company_name)}</STATICVARIABLES>
<TDL><TDLMESSAGE>
<COLLECTION NAME="{spec['collection']}" ISMODIFY="No"><TYPE>{spec['type']}</TYPE><FETCH>{spec['fetch']}</FETCH></COLLECTION>
</TDLMESSAGE></TDL>
</DESC></BODY>
</ENVELOPE>"""
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")