Add Phase 1 accounting Tally connector and ACT storage foundation
This commit is contained in:
@@ -0,0 +1,248 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from contextlib import contextmanager
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
import json
|
||||
import os
|
||||
import sqlite3
|
||||
from typing import Iterator, Sequence
|
||||
|
||||
ACT_SCHEMA_VERSION = 1
|
||||
|
||||
|
||||
class AccountingActStoreError(RuntimeError):
|
||||
pass
|
||||
|
||||
|
||||
def _utc_now_iso() -> str:
|
||||
return datetime.now(timezone.utc).replace(microsecond=0).isoformat()
|
||||
|
||||
|
||||
def _safe_client_id(client_id: int) -> int:
|
||||
value = int(client_id)
|
||||
if value <= 0:
|
||||
raise AccountingActStoreError("client_id must be a positive integer")
|
||||
return value
|
||||
|
||||
|
||||
class AccountingActStore:
|
||||
"""Client-scoped SQLite storage for ERP accounting data.
|
||||
|
||||
The `.act` extension is intentional; the underlying file format is SQLite.
|
||||
Phase 1 stores metadata, discovered Tally companies, connection history and
|
||||
sync-run control records only. Transaction/master tables arrive in later
|
||||
phases.
|
||||
"""
|
||||
|
||||
def __init__(self, root: str | Path) -> None:
|
||||
self.root = Path(root).expanduser().resolve()
|
||||
|
||||
def client_dir(self, client_id: int) -> Path:
|
||||
cid = _safe_client_id(client_id)
|
||||
return self.root / "Accounting" / f"client_{cid:08d}"
|
||||
|
||||
def db_path(self, client_id: int) -> Path:
|
||||
cid = _safe_client_id(client_id)
|
||||
return self.client_dir(cid) / f"client_{cid:08d}.act"
|
||||
|
||||
@contextmanager
|
||||
def connect(self, client_id: int) -> Iterator[sqlite3.Connection]:
|
||||
path = self.db_path(client_id)
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
connection = sqlite3.connect(path, timeout=30)
|
||||
connection.row_factory = sqlite3.Row
|
||||
try:
|
||||
connection.execute("PRAGMA foreign_keys = ON")
|
||||
connection.execute("PRAGMA journal_mode = WAL")
|
||||
connection.execute("PRAGMA synchronous = NORMAL")
|
||||
yield connection
|
||||
connection.commit()
|
||||
except Exception:
|
||||
connection.rollback()
|
||||
raise
|
||||
finally:
|
||||
connection.close()
|
||||
|
||||
def initialize(
|
||||
self,
|
||||
client_id: int,
|
||||
*,
|
||||
tenant_id: int | None = None,
|
||||
client_name: str = "",
|
||||
created_by_user_id: int | None = None,
|
||||
) -> Path:
|
||||
path = self.db_path(client_id)
|
||||
with self.connect(client_id) as db:
|
||||
db.executescript(
|
||||
"""
|
||||
CREATE TABLE IF NOT EXISTS act_meta (
|
||||
key TEXT PRIMARY KEY,
|
||||
value TEXT NOT NULL,
|
||||
updated_at_utc TEXT NOT NULL
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS tally_companies (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
tally_guid TEXT NOT NULL DEFAULT '',
|
||||
company_name TEXT NOT NULL,
|
||||
gstin TEXT NOT NULL DEFAULT '',
|
||||
first_seen_at_utc TEXT NOT NULL,
|
||||
last_seen_at_utc TEXT NOT NULL,
|
||||
is_currently_loaded INTEGER NOT NULL DEFAULT 1,
|
||||
UNIQUE(tally_guid, company_name)
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS tally_company_mapping (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
client_id INTEGER NOT NULL,
|
||||
registration_id INTEGER,
|
||||
tally_guid TEXT NOT NULL DEFAULT '',
|
||||
company_name TEXT NOT NULL DEFAULT '',
|
||||
gstin TEXT NOT NULL DEFAULT '',
|
||||
is_active INTEGER NOT NULL DEFAULT 1,
|
||||
mapped_at_utc TEXT NOT NULL,
|
||||
mapped_by_user_id INTEGER
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS tally_connection_history (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
checked_at_utc TEXT NOT NULL,
|
||||
connected INTEGER NOT NULL,
|
||||
tally_url TEXT NOT NULL DEFAULT '',
|
||||
company_count INTEGER NOT NULL DEFAULT 0,
|
||||
error_message TEXT,
|
||||
payload_json TEXT NOT NULL DEFAULT '{}'
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS tally_sync_runs (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
sync_type TEXT NOT NULL,
|
||||
status TEXT NOT NULL,
|
||||
started_at_utc TEXT NOT NULL,
|
||||
completed_at_utc TEXT,
|
||||
requested_by_user_id INTEGER,
|
||||
records_received INTEGER NOT NULL DEFAULT 0,
|
||||
error_message TEXT,
|
||||
details_json TEXT NOT NULL DEFAULT '{}'
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS ix_tally_companies_last_seen
|
||||
ON tally_companies(last_seen_at_utc);
|
||||
CREATE INDEX IF NOT EXISTS ix_tally_connection_history_checked
|
||||
ON tally_connection_history(checked_at_utc);
|
||||
CREATE INDEX IF NOT EXISTS ix_tally_sync_runs_started
|
||||
ON tally_sync_runs(started_at_utc);
|
||||
"""
|
||||
)
|
||||
now = _utc_now_iso()
|
||||
metadata = {
|
||||
"schema_version": str(ACT_SCHEMA_VERSION),
|
||||
"client_id": str(_safe_client_id(client_id)),
|
||||
"tenant_id": "" if tenant_id is None else str(int(tenant_id)),
|
||||
"client_name": str(client_name or ""),
|
||||
"created_by_user_id": "" if created_by_user_id is None else str(int(created_by_user_id)),
|
||||
"storage_kind": "ERP_ACCOUNTING_ACT_SQLITE",
|
||||
}
|
||||
for key, value in metadata.items():
|
||||
db.execute(
|
||||
"""
|
||||
INSERT INTO act_meta(key, value, updated_at_utc)
|
||||
VALUES (?, ?, ?)
|
||||
ON CONFLICT(key) DO UPDATE SET
|
||||
value=excluded.value,
|
||||
updated_at_utc=excluded.updated_at_utc
|
||||
""",
|
||||
(key, value, now),
|
||||
)
|
||||
return path
|
||||
|
||||
def record_tally_status(self, client_id: int, status: dict) -> None:
|
||||
self.initialize(client_id)
|
||||
now = _utc_now_iso()
|
||||
companies: Sequence[dict] = status.get("companies") or []
|
||||
with self.connect(client_id) as db:
|
||||
db.execute("UPDATE tally_companies SET is_currently_loaded = 0")
|
||||
for company in companies:
|
||||
name = str(company.get("name") or "").strip()
|
||||
if not name:
|
||||
continue
|
||||
guid = str(company.get("guid") or "").strip()
|
||||
gstin = str(company.get("gstin") or "").strip().upper()
|
||||
existing = db.execute(
|
||||
"""
|
||||
SELECT id FROM tally_companies
|
||||
WHERE tally_guid = ? AND company_name = ?
|
||||
LIMIT 1
|
||||
""",
|
||||
(guid, name),
|
||||
).fetchone()
|
||||
if existing:
|
||||
db.execute(
|
||||
"""
|
||||
UPDATE tally_companies
|
||||
SET gstin=?, last_seen_at_utc=?, is_currently_loaded=1
|
||||
WHERE id=?
|
||||
""",
|
||||
(gstin, now, int(existing["id"])),
|
||||
)
|
||||
else:
|
||||
db.execute(
|
||||
"""
|
||||
INSERT INTO tally_companies(
|
||||
tally_guid, company_name, gstin,
|
||||
first_seen_at_utc, last_seen_at_utc, is_currently_loaded
|
||||
) VALUES (?, ?, ?, ?, ?, 1)
|
||||
""",
|
||||
(guid, name, gstin, now, now),
|
||||
)
|
||||
|
||||
db.execute(
|
||||
"""
|
||||
INSERT INTO tally_connection_history(
|
||||
checked_at_utc, connected, tally_url, company_count,
|
||||
error_message, payload_json
|
||||
) VALUES (?, ?, ?, ?, ?, ?)
|
||||
""",
|
||||
(
|
||||
now,
|
||||
1 if status.get("connected") else 0,
|
||||
str(status.get("url") or ""),
|
||||
int(status.get("company_count") or 0),
|
||||
str(status.get("error") or "") or None,
|
||||
json.dumps(status, ensure_ascii=False, separators=(",", ":")),
|
||||
),
|
||||
)
|
||||
|
||||
def snapshot(self, client_id: int) -> dict:
|
||||
self.initialize(client_id)
|
||||
with self.connect(client_id) as db:
|
||||
meta_rows = db.execute("SELECT key, value FROM act_meta ORDER BY key").fetchall()
|
||||
companies = db.execute(
|
||||
"""
|
||||
SELECT tally_guid AS guid, company_name AS name, gstin,
|
||||
first_seen_at_utc, last_seen_at_utc, is_currently_loaded
|
||||
FROM tally_companies
|
||||
ORDER BY is_currently_loaded DESC, company_name COLLATE NOCASE
|
||||
"""
|
||||
).fetchall()
|
||||
latest = db.execute(
|
||||
"""
|
||||
SELECT checked_at_utc, connected, tally_url, company_count, error_message
|
||||
FROM tally_connection_history
|
||||
ORDER BY id DESC LIMIT 1
|
||||
"""
|
||||
).fetchone()
|
||||
return {
|
||||
"db_path": str(self.db_path(client_id)),
|
||||
"metadata": {row["key"]: row["value"] for row in meta_rows},
|
||||
"companies": [dict(row) for row in companies],
|
||||
"latest_connection": dict(latest) if latest else None,
|
||||
}
|
||||
|
||||
|
||||
def default_accounting_root() -> Path:
|
||||
env = os.getenv("AUDIT_ACCOUNTING_STORAGE_ROOT", "").strip()
|
||||
if env:
|
||||
return Path(env).expanduser().resolve()
|
||||
return Path.cwd() / "data" / "accounting"
|
||||
Reference in New Issue
Block a user