Add Phase 1 distributed lean workstation agent foundation

This commit is contained in:
A R R R Associates
2026-08-22 11:37:25 +05:30
parent 798a2532c1
commit c8b7dcb805
9 changed files with 269 additions and 7 deletions
@@ -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")
+1 -1
View File
@@ -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)
@@ -1,2 +1,2 @@
__version__ = "1.9.5"
__version__ = "1.10.0"
AGENT_NAME = "ERP Local Agent"
@@ -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)
@@ -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()
@@ -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
+35
View File
@@ -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.
@@ -21,7 +21,7 @@
</div>
{% endif %}
<div class="grid grid-cols-1 md:grid-cols-4 gap-4">
<div class="grid grid-cols-1 md:grid-cols-5 gap-4">
<div class="rounded-2xl border bg-white p-4 shadow-sm">
<div class="text-xs text-slate-500">Storage Nodes</div>
<div class="mt-1 text-2xl font-bold text-slate-900">{{ nodes|length }}</div>
@@ -38,6 +38,10 @@
<div class="text-xs text-slate-500">Recent Download Requests</div>
<div class="mt-1 text-2xl font-bold text-slate-900">{{ recent_download_requests|length if recent_download_requests is defined else 0 }}</div>
</div>
<div class="rounded-2xl border bg-white p-4 shadow-sm">
<div class="text-xs text-slate-500">Workstations Seen</div>
<div class="mt-1 text-2xl font-bold text-slate-900">{{ workstation_agents|length if workstation_agents is defined else 0 }}</div>
</div>
</div>
{% if request.query_params.get('created') == '1' %}
@@ -166,6 +170,31 @@
</table>
</div>
<div class="rounded-2xl border bg-white overflow-hidden shadow-sm">
<div class="p-4 border-b">
<h2 class="font-bold text-slate-900">Registered Workstations</h2>
<p class="text-xs text-slate-500 mt-1">Each PC keeps a persistent agent identity. The branch storage node remains unchanged; multiple Tally workstations can report through the same outbound tunnel credentials.</p>
</div>
<div class="overflow-x-auto">
<table class="w-full text-sm">
<thead class="bg-slate-50 text-slate-600"><tr><th class="p-3 text-left">Workstation</th><th class="p-3 text-left">Agent</th><th class="p-3 text-left">Tally</th><th class="p-3 text-left">Open Companies</th><th class="p-3 text-left">Last Seen</th></tr></thead>
<tbody>
{% for ws in workstation_agents or [] %}
<tr class="border-t align-top">
<td class="p-3"><div class="font-semibold">{{ ws.machine_name }}</div><div class="text-xs text-slate-500">{{ ws.platform_name or '-' }}</div><div class="text-[11px] text-slate-400 font-mono">{{ ws.agent_instance_id[:12] }}...</div></td>
<td class="p-3"><div>{{ ws.agent_version or '-' }}</div><div class="text-xs {{ 'text-emerald-700' if ws.is_active else 'text-slate-500' }}">{{ ws.status }}</div></td>
<td class="p-3"><span class="px-2 py-1 rounded-full text-xs {{ 'bg-emerald-100 text-emerald-700' if ws.tally_connected else 'bg-slate-100 text-slate-600' }}">{{ 'Connected' if ws.tally_connected else 'Not connected' }}</span></td>
<td class="p-3">{{ ws.tally_company_count or 0 }}</td>
<td class="p-3 text-xs">{{ ws.last_seen_at_utc or '-' }}</td>
</tr>
{% else %}
<tr><td colspan="5" class="p-6 text-center text-slate-500">No workstation agent has connected since the Phase 1 upgrade.</td></tr>
{% endfor %}
</tbody>
</table>
</div>
</div>
<div class="grid grid-cols-1 lg:grid-cols-2 gap-6">
<div class="rounded-2xl border bg-white overflow-hidden shadow-sm">
<div class="p-4 border-b"><h2 class="font-bold text-slate-900">Recent Storage Jobs</h2></div>
+55 -2
View File
@@ -1,6 +1,7 @@
from __future__ import annotations
from datetime import date, datetime, timezone
import json
import asyncio
import logging
from pathlib import Path
@@ -81,7 +82,7 @@ from app.modules.services.models import ClientServiceSubscription
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.documents.models import BranchStorageNode, ERPWorkstationAgent
from app.modules.accounting.agent_bridge import list_pending_agent_commands, record_agent_command_result
from app.modules.services.task_documents import (
@@ -1254,6 +1255,53 @@ def _selected_or_forced_branch_id(branch_id: str | None, user, scope) -> int | N
return int(forced_branch_id)
return int(branch_id) if branch_id and str(branch_id).isdigit() else None
def _record_workstation_status(db, node: BranchStorageNode, payload: dict | None, client_ip: str | None = None) -> None:
if not isinstance(payload, dict):
return
instance_id = str(payload.get("agent_instance_id") or "").strip()
fingerprint = str(payload.get("machine_fingerprint") or "").strip()
machine_name = str(payload.get("machine_name") or "").strip()
if not instance_id or not fingerprint or not machine_name:
return
row = db.execute(
select(ERPWorkstationAgent).where(
ERPWorkstationAgent.storage_node_id == node.id,
ERPWorkstationAgent.agent_instance_id == instance_id,
)
).scalar_one_or_none()
if row is None:
row = ERPWorkstationAgent(
tenant_id=node.tenant_id, branch_id=node.branch_id, storage_node_id=node.id,
agent_instance_id=instance_id, machine_fingerprint=fingerprint, machine_name=machine_name,
)
db.add(row)
tally = payload.get("tally") if isinstance(payload.get("tally"), dict) else {}
capabilities = payload.get("capabilities") if isinstance(payload.get("capabilities"), list) else []
row.tenant_id = node.tenant_id
row.branch_id = node.branch_id
row.machine_fingerprint = fingerprint
row.machine_name = machine_name[:200]
row.platform_name = str(payload.get("platform_name") or "")[:120] or None
row.agent_version = str(payload.get("agent_version") or "")[:40] or None
row.capabilities_json = json.dumps(capabilities, ensure_ascii=False, separators=(",", ":"))
row.tally_connected = bool(tally.get("connected"))
row.tally_company_count = int(tally.get("company_count") or 0)
row.tally_companies_json = json.dumps(list(tally.get("companies") or []), ensure_ascii=False, separators=(",", ":"))
row.status = "online"
row.is_active = True
row.last_seen_at_utc = datetime.now(timezone.utc)
row.last_seen_ip = (client_ip or "")[:80] or None
def _visible_workstation_agents(db, tenant_id: int | None, branch_id: int | None) -> list[ERPWorkstationAgent]:
stmt = select(ERPWorkstationAgent).order_by(ERPWorkstationAgent.last_seen_at_utc.desc())
if tenant_id is not None:
stmt = stmt.where(ERPWorkstationAgent.tenant_id == tenant_id)
if branch_id is not None:
stmt = stmt.where(ERPWorkstationAgent.branch_id == branch_id)
return list(db.execute(stmt).scalars().all())
@router.get("/storage-nodes")
def storage_nodes(request: Request):
db = CommonSessionLocal()
@@ -1270,6 +1318,7 @@ def storage_nodes(request: Request):
branches = _visible_branches(db, user, scope)
jobs = list_storage_jobs(db, tenant_id=tenant_filter, branch_id=branch_filter, limit=10)
requests = list_download_requests(db, tenant_id=tenant_filter, branch_id=branch_filter, limit=10)
workstations = _visible_workstation_agents(db, tenant_filter, branch_filter)
return _render(
request,
"modules/documents/templates/documents/storage_nodes.html",
@@ -1281,6 +1330,7 @@ def storage_nodes(request: Request):
branch_map=_branch_name_map(db, branches),
recent_jobs=jobs,
recent_download_requests=requests,
workstation_agents=workstations,
generated_secret=None,
storage_scope_title=_storage_scope_title(user, scope),
erp_local_agent_version=ERP_LOCAL_AGENT_VERSION,
@@ -1535,7 +1585,7 @@ def toggle_storage_node(request: Request, node_id: int, csrf_token: str = Form(.
scope = build_document_scope(request, db, user)
if not _can_manage_branch_storage(scope):
return _redirect_denied()
from app.modules.documents.models import BranchStorageNode
from app.modules.documents.models import BranchStorageNode, ERPWorkstationAgent
node = db.get(BranchStorageNode, node_id)
if not _node_allowed_for_storage_scope(node, user, scope):
return _redirect_denied()
@@ -1839,6 +1889,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 isinstance(message.get("workstation"), dict):
client_ip = websocket.client.host if websocket.client else None
_record_workstation_status(db, node, message.get("workstation"), client_ip)
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"}: