diff --git a/alembic/versions/20260822_distributed_lean_agent_phase1.py b/alembic/versions/20260822_distributed_lean_agent_phase1.py new file mode 100644 index 0000000..30d896c --- /dev/null +++ b/alembic/versions/20260822_distributed_lean_agent_phase1.py @@ -0,0 +1,48 @@ +"""distributed lean workstation agent phase 1 + +Revision ID: 20260822_lean_agent_p1 +Revises: 20260811_normal_review +Create Date: 2026-08-22 +""" +from alembic import op +import sqlalchemy as sa + +revision = "20260822_lean_agent_p1" +down_revision = "20260811_normal_review" +branch_labels = None +depends_on = None + +def upgrade() -> None: + op.create_table( + "erp_workstation_agents", + sa.Column("id", sa.Integer(), primary_key=True, autoincrement=True), + sa.Column("tenant_id", sa.Integer(), sa.ForeignKey("tenants.id", ondelete="CASCADE"), nullable=False), + sa.Column("branch_id", sa.Integer(), sa.ForeignKey("branches.id", ondelete="SET NULL"), nullable=True), + sa.Column("storage_node_id", sa.Integer(), sa.ForeignKey("branch_storage_nodes.id", ondelete="CASCADE"), nullable=False), + sa.Column("agent_instance_id", sa.String(length=64), nullable=False), + sa.Column("machine_fingerprint", sa.String(length=64), nullable=False), + sa.Column("machine_name", sa.String(length=200), nullable=False), + sa.Column("platform_name", sa.String(length=120), nullable=True), + sa.Column("agent_version", sa.String(length=40), nullable=True), + sa.Column("capabilities_json", sa.Text(), nullable=True), + sa.Column("tally_connected", sa.Boolean(), nullable=False, server_default=sa.false()), + sa.Column("tally_company_count", sa.Integer(), nullable=False, server_default="0"), + sa.Column("tally_companies_json", sa.Text(), nullable=True), + sa.Column("status", sa.String(length=30), nullable=False, server_default="online"), + sa.Column("is_active", sa.Boolean(), nullable=False, server_default=sa.true()), + sa.Column("first_seen_at_utc", sa.DateTime(timezone=True), nullable=False), + sa.Column("last_seen_at_utc", sa.DateTime(timezone=True), nullable=False), + sa.Column("last_seen_ip", sa.String(length=80), nullable=True), + sa.UniqueConstraint("storage_node_id", "agent_instance_id", name="uq_erp_workstation_agent_instance"), + ) + for name, cols in [ + ("ix_erp_workstation_agents_tenant_id", ["tenant_id"]), ("ix_erp_workstation_agents_branch_id", ["branch_id"]), + ("ix_erp_workstation_agents_storage_node_id", ["storage_node_id"]), ("ix_erp_workstation_agents_agent_instance_id", ["agent_instance_id"]), + ("ix_erp_workstation_agents_machine_fingerprint", ["machine_fingerprint"]), ("ix_erp_workstation_agents_agent_version", ["agent_version"]), + ("ix_erp_workstation_agents_tally_connected", ["tally_connected"]), ("ix_erp_workstation_agents_status", ["status"]), + ("ix_erp_workstation_agents_is_active", ["is_active"]), ("ix_erp_workstation_agents_last_seen_at_utc", ["last_seen_at_utc"]), + ]: + op.create_index(name, "erp_workstation_agents", cols) + +def downgrade() -> None: + op.drop_table("erp_workstation_agents") diff --git a/app/modules/documents/agent_package.py b/app/modules/documents/agent_package.py index b551b48..5ce678e 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.9.5" +ERP_LOCAL_AGENT_VERSION = "1.10.0" 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) 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 094c149..c9c5ecc 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.9.5" +__version__ = "1.10.0" AGENT_NAME = "ERP Local Agent" diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/main.py b/app/modules/documents/local_agent_runtime/erp_local_agent/main.py index 2111eb9..b387c56 100644 --- a/app/modules/documents/local_agent_runtime/erp_local_agent/main.py +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/main.py @@ -15,6 +15,7 @@ from .logger import setup_logger from .sync import StorageAgent from .tunnel import StorageAgentTunnel from .updater import AgentUpdater +from .workstation_identity import load_or_create def build_parser() -> argparse.ArgumentParser: @@ -36,6 +37,10 @@ def main() -> int: db.set_meta("agent_version", __version__) db.set_meta("node_code", config.node_code) db.set_meta("erp_base_url", config.erp_base_url) + workstation = load_or_create(root) + db.set_meta("agent_instance_id", workstation["agent_instance_id"]) + db.set_meta("machine_name", workstation["machine_name"]) + db.set_meta("machine_fingerprint", workstation["machine_fingerprint"]) db.record_event("INFO", "agent_started", f"ERP Local Agent {__version__} started") client = ERPClient(config) agent = StorageAgent(config, client, db, logger) 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 d09bc02..08ef417 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 @@ -1,12 +1,17 @@ from __future__ import annotations import asyncio +import time from datetime import datetime, timezone +from pathlib import Path from typing import Any import websockets from .commands import AgentCommandProcessor +from . import __version__ +from .tally import TallyLiveConnector +from .workstation_identity import load_or_create from .sync import StorageAgent @@ -17,6 +22,10 @@ class StorageAgentTunnel: self.logger = agent.logger self.db = agent.db self.command_processor = AgentCommandProcessor(self.config, self.logger) + self.identity = load_or_create(Path.cwd()) + self.tally = TallyLiveConnector() + self._last_tally_status_at = 0.0 + self._cached_tally_status = {"connected": False, "company_count": 0, "companies": []} async def run_forever(self) -> None: self.logger.info("Starting tunnel mode for node=%s", self.config.node_code) @@ -48,7 +57,7 @@ class StorageAgentTunnel: ) as websocket: self.db.record_connection("connected", endpoint) self.db.record_event("INFO", "tunnel_connected", f"WebSocket connected to {endpoint}") - await websocket.send(self._json_text({"type": "ready", "agent_time_utc": self._now()})) + await websocket.send(self._json_text({"type": "ready", "agent_time_utc": self._now(), "workstation": self._workstation_status(force_tally=True)})) self.agent._maybe_heartbeat(force=True) async for raw in websocket: message = self._parse_json(raw) @@ -76,13 +85,43 @@ class StorageAgentTunnel: "jobs_seen": len(jobs), "requests_seen": len(requests_), "commands_seen": len(commands), - "capabilities": ["storage", "tally", "accounting_act", "tally_master_sync", "tally_transaction_sync", "local_dashboard", "manual_updates"], + "capabilities": self._capabilities(), + "workstation": self._workstation_status(), "agent_time_utc": self._now(), })) continue if msg_type == "error": self.logger.error("Tunnel server error: %s", message.get("error")) + def _capabilities(self) -> list[str]: + return [ + "storage", "local_dashboard", "manual_updates", "accounting_act", + "tally.status", "tally.company_identity", "tally.read_groups", + "tally.read_ledgers", "tally.read_stock_items", "tally.read_vouchers", + "tally.read_trial_balance", "tally.master_sync", "tally.transaction_sync", + "tally.write_journal", "accounting.it_depreciation", + ] + + def _workstation_status(self, force_tally: bool = False) -> dict[str, Any]: + now = time.monotonic() + if force_tally or now - self._last_tally_status_at >= 30: + try: + self._cached_tally_status = self.tally.status() + except Exception as exc: + self._cached_tally_status = {"connected": False, "company_count": 0, "companies": [], "error": str(exc)} + self._last_tally_status_at = now + tally = self._cached_tally_status or {} + return { + **self.identity, + "agent_version": __version__, + "capabilities": self._capabilities(), + "tally": { + "connected": bool(tally.get("connected")), + "company_count": int(tally.get("company_count") or 0), + "companies": list(tally.get("companies") or []), + }, + } + def _now(self) -> str: return datetime.now(timezone.utc).isoformat() diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/workstation_identity.py b/app/modules/documents/local_agent_runtime/erp_local_agent/workstation_identity.py new file mode 100644 index 0000000..3e57f02 --- /dev/null +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/workstation_identity.py @@ -0,0 +1,53 @@ +from __future__ import annotations + +import hashlib +import json +import os +import platform +import socket +import subprocess +from pathlib import Path +from uuid import uuid4 + +IDENTITY_FILE = "workstation_identity.json" + +def _windows_machine_guid() -> str: + if os.name != "nt": + return "" + try: + completed = subprocess.run( + ["reg", "query", r"HKLM\SOFTWARE\Microsoft\Cryptography", "/v", "MachineGuid"], + capture_output=True, text=True, timeout=5, check=False, + ) + for line in completed.stdout.splitlines(): + if "MachineGuid" in line: + return line.split()[-1].strip() + except Exception: + pass + return "" + +def _fingerprint(machine_name: str) -> str: + material = "|".join([_windows_machine_guid(), machine_name, platform.system(), platform.machine()]) + return hashlib.sha256(material.encode("utf-8", errors="ignore")).hexdigest() + +def load_or_create(root: Path) -> dict: + data_dir = root / "data" + data_dir.mkdir(parents=True, exist_ok=True) + path = data_dir / IDENTITY_FILE + machine_name = socket.gethostname().strip() or platform.node().strip() or "UNKNOWN-PC" + existing = {} + try: + existing = json.loads(path.read_text(encoding="utf-8")) if path.exists() else {} + except Exception: + existing = {} + instance_id = str(existing.get("agent_instance_id") or "").strip() or uuid4().hex + payload = { + "agent_instance_id": instance_id, + "machine_fingerprint": _fingerprint(machine_name), + "machine_name": machine_name, + "platform_name": f"{platform.system()} {platform.release()} {platform.machine()}".strip(), + } + tmp = path.with_suffix(".tmp") + tmp.write_text(json.dumps(payload, indent=2, sort_keys=True), encoding="utf-8") + tmp.replace(path) + return payload diff --git a/app/modules/documents/models.py b/app/modules/documents/models.py index e25fa4b..5102195 100644 --- a/app/modules/documents/models.py +++ b/app/modules/documents/models.py @@ -183,6 +183,41 @@ class BranchStorageNode(CommonBase): updated_at_utc: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=lambda: datetime.now(timezone.utc), onupdate=lambda: datetime.now(timezone.utc), nullable=False) +class ERPWorkstationAgent(CommonBase): + """A physical workstation running the lean ERP Local Agent. + + Workstations are intentionally separate from BranchStorageNode: a branch keeps + one canonical storage node/root while any number of office PCs may use that + node's authenticated outbound tunnel for Tally/local-tool capabilities. + """ + + __tablename__ = "erp_workstation_agents" + __table_args__ = ( + UniqueConstraint("storage_node_id", "agent_instance_id", name="uq_erp_workstation_agent_instance"), + ) + + id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) + tenant_id: Mapped[int] = mapped_column(ForeignKey("tenants.id", ondelete="CASCADE"), nullable=False, index=True) + branch_id: Mapped[int | None] = mapped_column(ForeignKey("branches.id", ondelete="SET NULL"), nullable=True, index=True) + storage_node_id: Mapped[int] = mapped_column(ForeignKey("branch_storage_nodes.id", ondelete="CASCADE"), nullable=False, index=True) + agent_instance_id: Mapped[str] = mapped_column(String(64), nullable=False, index=True) + machine_fingerprint: Mapped[str] = mapped_column(String(64), nullable=False, index=True) + machine_name: Mapped[str] = mapped_column(String(200), nullable=False) + platform_name: Mapped[str | None] = mapped_column(String(120), nullable=True) + agent_version: Mapped[str | None] = mapped_column(String(40), nullable=True, index=True) + capabilities_json: Mapped[str | None] = mapped_column(Text, nullable=True) + tally_connected: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False, index=True) + tally_company_count: Mapped[int] = mapped_column(Integer, nullable=False, default=0) + tally_companies_json: Mapped[str | None] = mapped_column(Text, nullable=True) + status: Mapped[str] = mapped_column(String(30), nullable=False, default="online", index=True) + is_active: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True, index=True) + first_seen_at_utc: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=lambda: datetime.now(timezone.utc), nullable=False) + last_seen_at_utc: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=lambda: datetime.now(timezone.utc), nullable=False, index=True) + last_seen_ip: Mapped[str | None] = mapped_column(String(80), nullable=True) + + storage_node = relationship("BranchStorageNode") + + class DocumentStorageJob(CommonBase): """Cloud-to-local branch storage transfer job. diff --git a/app/modules/documents/templates/documents/storage_nodes.html b/app/modules/documents/templates/documents/storage_nodes.html index 318372b..a3a6ec9 100644 --- a/app/modules/documents/templates/documents/storage_nodes.html +++ b/app/modules/documents/templates/documents/storage_nodes.html @@ -21,7 +21,7 @@ {% endif %} -
Each PC keeps a persistent agent identity. The branch storage node remains unchanged; multiple Tally workstations can report through the same outbound tunnel credentials.
+| Workstation | Agent | Tally | Open Companies | Last Seen |
|---|---|---|---|---|
{{ ws.machine_name }} {{ ws.platform_name or '-' }} {{ ws.agent_instance_id[:12] }}... |
+ {{ ws.agent_version or '-' }} {{ ws.status }} |
+ {{ 'Connected' if ws.tally_connected else 'Not connected' }} | +{{ ws.tally_company_count or 0 }} | +{{ ws.last_seen_at_utc or '-' }} | +
| No workstation agent has connected since the Phase 1 upgrade. | ||||