Align Tally accounting storage with client storage policy

This commit is contained in:
A R R R Associates
2026-08-20 14:25:57 +05:30
parent f23da90dec
commit 3912d2fd16
6 changed files with 204 additions and 47 deletions
+26 -18
View File
@@ -1,6 +1,7 @@
from __future__ import annotations
from datetime import date, datetime, timezone
from pathlib import Path
from urllib.parse import quote
from fastapi import APIRouter, Form, Request
@@ -14,7 +15,7 @@ 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.documents.services import build_document_scope, get_active_storage_node_for_branch, client_folder_parts
from app.modules.accounting.agent_bridge import request_agent_command
from app.modules.registrations.models import ClientRegistration, RegistrationType
@@ -55,6 +56,18 @@ def _find_visible_client(db, request: Request, user, client_id: int):
return client, clients, scope
def _accounting_storage_payload(client) -> dict:
"""Use the exact client folder naming policy already used by Permanent/Engagement storage."""
letter, client_folder = client_folder_parts(client, int(client.id))
relative_dir = Path("Accounting") / "Clients" / letter / client_folder
return {
"client_id": int(client.id),
"client_name": str(client.client_name or "").strip(),
"client_code": str(getattr(client, "client_code", "") or "").strip(),
"accounting_relative_dir": relative_dir.as_posix(),
}
def _client_registrations(db, client, tenant_id: int):
rows = db.execute(
select(ClientRegistration, RegistrationType)
@@ -143,10 +156,7 @@ def tally_tool(
if should_query_agent and node and online:
payload = {}
if selected_client:
payload = {
"client_id": int(selected_client.id),
"client_name": selected_client.client_name,
}
payload = _accounting_storage_payload(selected_client)
try:
response_data = request_agent_command(
node.node_code,
@@ -216,8 +226,7 @@ def initialize_accounting_storage(
node.node_code,
"accounting_initialize",
{
"client_id": int(client.id),
"client_name": client.client_name,
**_accounting_storage_payload(client),
"tenant_id": int(scope.tenant_id),
"requested_by_user_id": int(user.id),
},
@@ -293,8 +302,7 @@ def map_tally_company(
node.node_code,
"accounting_map_company",
{
"client_id": int(client.id),
"client_name": client.client_name,
**_accounting_storage_payload(client),
"tenant_id": int(scope.tenant_id),
"registration": registration_payload,
"tally_guid": str(tally_guid or "").strip(),
@@ -349,7 +357,7 @@ def unmap_tally_company(
node.node_code,
"accounting_unmap_company",
{
"client_id": int(client.id),
**_accounting_storage_payload(client),
"mapping_id": int(mapping_id),
"unmapped_by_user_id": int(user.id),
},
@@ -398,7 +406,7 @@ def sync_tally_masters(
node.node_code,
"accounting_sync_masters",
{
"client_id": int(client.id),
**_accounting_storage_payload(client),
"tally_guid": str(tally_guid or "").strip(),
"requested_by_user_id": int(user.id),
},
@@ -454,7 +462,7 @@ def sync_tally_transactions(
node.node_code,
"accounting_sync_transactions",
{
"client_id": int(client.id),
**_accounting_storage_payload(client),
"tally_guid": str(tally_guid or "").strip(),
"date_from": start.isoformat(),
"date_to": end.isoformat(),
@@ -496,16 +504,16 @@ def depreciation_it_tool(
live_result = None; preview = None; depreciation_run = None; command_error = error or ""
if selected_client and node and online:
try:
status_response = request_agent_command(node.node_code,"phase6_status",{"client_id":int(selected_client.id)},timeout_seconds=20)
status_response = request_agent_command(node.node_code,"phase6_status",_accounting_storage_payload(selected_client),timeout_seconds=20)
if status_response.get("ok"): live_result=status_response.get("result") or {}
else: command_error=str(status_response.get("error") or "Local Agent status failed.")
chosen_guid=str(tally_guid or "").strip()
if chosen_guid:
preview_response=request_agent_command(node.node_code,"accounting_depreciation_preview",{"client_id":int(selected_client.id),"tally_guid":chosen_guid,"fy_start":start_text,"fy_end":end_text},timeout_seconds=60)
preview_response=request_agent_command(node.node_code,"accounting_depreciation_preview",{**_accounting_storage_payload(selected_client),"tally_guid":chosen_guid,"fy_start":start_text,"fy_end":end_text},timeout_seconds=60)
if preview_response.get("ok"): preview=(preview_response.get("result") or {}).get("preview")
else: command_error=str(preview_response.get("error") or "Depreciation preview failed.")
if run_id:
run_response=request_agent_command(node.node_code,"accounting_get_it_depreciation_run",{"client_id":int(selected_client.id),"run_id":int(run_id)},timeout_seconds=30)
run_response=request_agent_command(node.node_code,"accounting_get_it_depreciation_run",{**_accounting_storage_payload(selected_client),"run_id":int(run_id)},timeout_seconds=30)
if run_response.get("ok"): depreciation_run=(run_response.get("result") or {}).get("depreciation")
except Exception as exc: command_error=str(exc)
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)}
@@ -537,7 +545,7 @@ async def calculate_it_depreciation(request: Request):
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/depreciation?client_id={client.id}&tally_guid={quote(tally_guid)}&fy_start={start.isoformat()}&fy_end={end.isoformat()}&error={quote('ERP Local Agent is offline for the active branch.')}",status_code=303)
try:
result=request_agent_command(node.node_code,"accounting_calculate_it_depreciation",{"client_id":int(client.id),"tally_guid":tally_guid,"fy_start":start.isoformat(),"fy_end":end.isoformat(),"assignments":assignments,"depreciation_expense_ledger":str(form.get("depreciation_expense_ledger") or ""),"depreciation_reserve_ledger":str(form.get("depreciation_reserve_ledger") or ""),"requested_by_user_id":int(user.id)},timeout_seconds=120)
result=request_agent_command(node.node_code,"accounting_calculate_it_depreciation",{**_accounting_storage_payload(client),"tally_guid":tally_guid,"fy_start":start.isoformat(),"fy_end":end.isoformat(),"assignments":assignments,"depreciation_expense_ledger":str(form.get("depreciation_expense_ledger") or ""),"depreciation_reserve_ledger":str(form.get("depreciation_reserve_ledger") or ""),"requested_by_user_id":int(user.id)},timeout_seconds=120)
if not result.get("ok"): raise RuntimeError(str(result.get("error") or "Income-tax depreciation calculation failed."))
dep=(result.get("result") or {}).get("depreciation") or {}; rid=int(dep.get("run_id"))
except Exception as exc:
@@ -569,7 +577,7 @@ async def approve_it_depreciation(request: Request):
try:
result = request_agent_command(
node.node_code, "accounting_approve_it_depreciation",
{"client_id": client_id, "run_id": run_id, "approved_by_user_id": int(user.id), "approval_note": str(form.get("approval_note") or "")},
{**_accounting_storage_payload(client), "run_id": run_id, "approved_by_user_id": int(user.id), "approval_note": str(form.get("approval_note") or "")},
timeout_seconds=30,
)
if not result.get("ok"): raise RuntimeError(str(result.get("error") or "Approval failed."))
@@ -606,7 +614,7 @@ async def post_it_depreciation_to_tally(request: Request):
try:
result = request_agent_command(
node.node_code, "accounting_post_it_depreciation",
{"client_id": client_id, "run_id": run_id, "posted_by_user_id": int(user.id)},
{**_accounting_storage_payload(client), "run_id": run_id, "posted_by_user_id": int(user.id)},
timeout_seconds=120,
)
if not result.get("ok"): raise RuntimeError(str(result.get("error") or "Tally write-back failed."))
+1 -1
View File
@@ -4,7 +4,7 @@ import io
from pathlib import Path
import zipfile
ERP_LOCAL_AGENT_VERSION = "1.8.0"
ERP_LOCAL_AGENT_VERSION = "1.8.1"
ERP_LOCAL_AGENT_NAME = "ERP Local Agent"
RUNTIME_ROOT = Path(__file__).resolve().parent / "local_agent_runtime"
_DETERMINISTIC_ZIP_TIMESTAMP = (2026, 1, 1, 0, 0, 0)
@@ -1,4 +1,4 @@
ERP Local Agent 1.8.0
ERP Local Agent 1.8.1
Existing storage, WebSocket tunnel, Tally and client .act functionality are preserved.
@@ -35,3 +35,10 @@ Depreciation refinement 1.8.0:
- Empty Fixed Asset ledgers are omitted unless they have current-period additions.
- Bulk checkbox + rate assignment is available.
- Nested Tally accounting allocations and case-insensitive ledger matching improve fixed-asset additions detection.
Accounting storage policy 1.8.1:
- .act databases now use the same configured branch Storage Node root and client-folder naming policy as Permanent/Engagement storage.
- ERP builds the client folder through the existing client_folder_parts() helper and sends only the safe relative Accounting/Clients/... directory.
- Canonical layout: <STORAGE_ROOT>\Accounting\Clients\<Letter>\<ExistingClientFolder>\client_XXXXXXXX.act
- Existing databases from <STORAGE_ROOT>\Accounting\client_XXXXXXXX or <ERP_Local_Agent>\data\accounting\client_XXXXXXXX are migrated by SQLite backup on first use.
- Migration sources are retained as recovery copies; no existing .act file is deleted.
@@ -1,2 +1,2 @@
__version__ = "1.8.0"
__version__ = "1.8.1"
AGENT_NAME = "ERP Local Agent"
@@ -39,50 +39,184 @@ class LocalAccountingStore:
"""
def __init__(self, storage_root: Path):
# Accounting databases are application data, not document-storage payloads.
# Keep them under the Local Agent data directory by default. An explicit
# ACCOUNTING_ROOT may override this without changing document STORAGE_ROOT.
configured = str(os.getenv("ACCOUNTING_ROOT", "") or "").strip()
install_root = Path(__file__).resolve().parents[1]
self.root = Path(configured).expanduser().resolve() if configured else (install_root / "data" / "accounting").resolve()
self.legacy_root = Path(storage_root).resolve() / "Accounting"
self.root.mkdir(parents=True, exist_ok=True)
# Accounting follows the same branch Storage Node root used by Permanent
# and Engagement storage. The ERP supplies the exact client-relative
# directory built with its shared client_folder_parts() policy.
self.storage_root = Path(storage_root).resolve()
self.storage_root.mkdir(parents=True, exist_ok=True)
install_root = Path(__file__).resolve().parents[1]
self.agent_data_root = (install_root / "data").resolve()
self.agent_data_root.mkdir(parents=True, exist_ok=True)
self.path_index_file = self.agent_data_root / "accounting_path_index.json"
# Historical locations retained only as migration sources.
self.legacy_flat_root = self.storage_root / "Accounting"
self.legacy_agent_root = self.agent_data_root / "accounting"
self._client_paths: dict[str, str] = self._load_path_index()
def _load_path_index(self) -> dict[str, str]:
if not self.path_index_file.exists():
return {}
try:
raw = json.loads(self.path_index_file.read_text(encoding="utf-8"))
if not isinstance(raw, dict):
return {}
return {str(k): str(v) for k, v in raw.items() if str(v).strip()}
except Exception:
return {}
def _save_path_index(self) -> None:
tmp = self.path_index_file.with_suffix(".json.tmp")
tmp.write_text(json.dumps(self._client_paths, indent=2, sort_keys=True), encoding="utf-8")
tmp.replace(self.path_index_file)
@staticmethod
def _safe_relative_dir(raw: str) -> Path:
text = str(raw or "").strip().replace("\\", "/")
if not text:
raise ValueError("Accounting storage path was not supplied by ERP.")
path = Path(text)
if path.is_absolute():
raise ValueError("Accounting storage path must be relative to the configured Storage Node root.")
safe_parts = []
for part in path.parts:
value = str(part or "").strip()
if value in {"", ".", ".."}:
raise ValueError("Unsafe Accounting storage path received from ERP.")
if "/" in value or "\\" in value:
raise ValueError("Unsafe Accounting storage segment received from ERP.")
safe_parts.append(value)
relative = Path(*safe_parts)
# Accounting paths are intentionally constrained to Accounting/Clients/...
parts_lower = [p.lower() for p in relative.parts]
if len(parts_lower) < 4 or parts_lower[0] != "accounting" or parts_lower[1] != "clients":
raise ValueError("Accounting storage path does not follow the ERP Storage Node client policy.")
return relative
def bind_client_path(self, client_id: int, accounting_relative_dir: str) -> Path:
relative = self._safe_relative_dir(accounting_relative_dir)
target_dir = (self.storage_root / relative).resolve()
if os.path.commonpath([str(self.storage_root), str(target_dir)]) != str(self.storage_root):
raise ValueError("Accounting storage path escapes the configured Storage Node root.")
def _migrate_legacy_client_if_needed(self, client_id: int) -> None:
key = self._client_key(client_id)
new_dir = self.root / key
new_db = new_dir / f"{key}.act"
old_dir = self.legacy_root / key
old_db = old_dir / f"{key}.act"
if new_db.exists() or not old_db.exists():
target_dir.mkdir(parents=True, exist_ok=True)
target_db = target_dir / f"{key}.act"
self._client_paths[str(int(client_id))] = relative.as_posix()
self._save_path_index()
self._migrate_existing_database(client_id, target_db)
return target_db
@staticmethod
def _database_activity_mtime(path: Path) -> float:
stamps = []
for candidate in (path, Path(str(path) + "-wal"), Path(str(path) + "-shm")):
try:
if candidate.exists():
stamps.append(candidate.stat().st_mtime)
except OSError:
pass
return max(stamps) if stamps else 0.0
@staticmethod
def _sqlite_backup(source: Path, destination: Path) -> None:
destination.parent.mkdir(parents=True, exist_ok=True)
tmp = destination.with_suffix(destination.suffix + ".migrating")
if tmp.exists():
tmp.unlink()
src = sqlite3.connect(source, timeout=60)
dst = sqlite3.connect(tmp, timeout=60)
try:
src.execute("PRAGMA busy_timeout=60000")
dst.execute("PRAGMA busy_timeout=60000")
src.backup(dst)
dst.commit()
finally:
dst.close()
src.close()
tmp.replace(destination)
def _migration_candidates(self, client_id: int, target_db: Path) -> list[Path]:
key = self._client_key(client_id)
filename = f"{key}.act"
candidates: list[Path] = []
# Original Phase 1-7 flat Storage Root location.
candidates.append(self.legacy_flat_root / key / filename)
# Temporary v1.8.0 Local-Agent data location.
candidates.append(self.legacy_agent_root / key / filename)
# If a client was renamed, an earlier canonical client-folder path may exist.
accounting_clients = self.storage_root / "Accounting" / "Clients"
if accounting_clients.exists():
for path in accounting_clients.glob(f"*/*/{filename}"):
candidates.append(path)
unique = []
seen = set()
for path in candidates:
try:
resolved = path.resolve()
except Exception:
resolved = path
if resolved == target_db.resolve():
continue
marker = str(resolved).lower()
if marker in seen or not path.exists():
continue
seen.add(marker)
unique.append(path)
return unique
def _migrate_existing_database(self, client_id: int, target_db: Path) -> None:
if target_db.exists():
return
new_dir.mkdir(parents=True, exist_ok=True)
# Copy, do not delete, the legacy database. The old copy remains a recovery
# fallback while all new reads/writes switch to Local Agent data/accounting.
shutil.copy2(old_db, new_db)
for suffix in ("-wal", "-shm"):
src = Path(str(old_db) + suffix)
if src.exists():
shutil.copy2(src, Path(str(new_db) + suffix))
candidates = self._migration_candidates(client_id, target_db)
if not candidates:
return
source = max(candidates, key=self._database_activity_mtime)
self._sqlite_backup(source, target_db)
# Source is deliberately left untouched as a recovery copy.
def _bound_relative_dir(self, client_id: int) -> Path | None:
raw = self._client_paths.get(str(int(client_id)))
if not raw:
return None
try:
return self._safe_relative_dir(raw)
except Exception:
return None
@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)
relative = self._bound_relative_dir(client_id)
if relative is None:
# Do not silently create a new Accounting DB in a non-policy location.
# Existing legacy DBs remain discoverable only after ERP supplies its
# canonical client path on the next status/initialize command.
return self.storage_root / "Accounting" / "Clients" / "_UNBOUND" / self._client_key(client_id)
return self.storage_root / relative
def db_path(self, client_id: int) -> Path:
self._migrate_legacy_client_if_needed(client_id)
key = self._client_key(client_id)
return self.client_dir(client_id) / f"{key}.act"
def exists(self, client_id: int) -> bool:
self._migrate_legacy_client_if_needed(client_id)
return self.db_path(client_id).is_file()
def connect(self, client_id: int):
self._migrate_legacy_client_if_needed(client_id)
if self._bound_relative_dir(client_id) is None:
raise RuntimeError(
"Accounting storage path is not bound to the ERP client folder policy. "
"Refresh Tools → Tally so ERP can supply the canonical Storage Node path."
)
path = self.db_path(client_id)
path.parent.mkdir(parents=True, exist_ok=True)
db = sqlite3.connect(path, timeout=60)
@@ -19,6 +19,7 @@ class AgentCommandProcessor:
command_id = str(command.get("command_id") or "").strip()
action = str(command.get("action") or "").strip()
payload = command.get("payload") or {}
self._bind_accounting_storage(payload)
result: dict[str, Any] | None = None
error: str | None = None
ok = False
@@ -53,6 +54,13 @@ class AgentCommandProcessor:
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 _bind_accounting_storage(self, payload: dict[str, Any]) -> None:
client_id = payload.get("client_id")
relative_dir = str(payload.get("accounting_relative_dir") or "").strip()
if client_id in (None, "") or not relative_dir:
return
self.store.bind_client_path(int(client_id), relative_dir)
def _agent_info(self) -> dict[str, Any]:
return {
"name": "ERP Local Agent", "version": __version__,