diff --git a/alembic/versions/20260822_distributed_lean_agent_phase2.py b/alembic/versions/20260822_distributed_lean_agent_phase2.py new file mode 100644 index 0000000..a5be7fd --- /dev/null +++ b/alembic/versions/20260822_distributed_lean_agent_phase2.py @@ -0,0 +1,92 @@ +"""distributed lean workstation agent phase 2 durable jobs + +Revision ID: 20260822_lean_agent_p2 +Revises: 20260822_lean_agent_p1 +Create Date: 2026-08-22 +""" +from alembic import op +import sqlalchemy as sa + +revision = "20260822_lean_agent_p2" +down_revision = "20260822_lean_agent_p1" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.create_table( + "erp_agent_jobs", + sa.Column("id", sa.Integer(), primary_key=True, autoincrement=True), + sa.Column("job_uuid", sa.String(length=64), nullable=False), + 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("workstation_agent_id", sa.Integer(), sa.ForeignKey("erp_workstation_agents.id", ondelete="CASCADE"), nullable=False), + sa.Column("action", sa.String(length=120), nullable=False), + sa.Column("payload_json", sa.Text(), nullable=False, server_default="{}"), + sa.Column("idempotency_key", sa.String(length=200), nullable=True), + sa.Column("status", sa.String(length=30), nullable=False, server_default="queued"), + sa.Column("priority", sa.Integer(), nullable=False, server_default="5"), + sa.Column("max_attempts", sa.Integer(), nullable=False, server_default="3"), + sa.Column("attempts", sa.Integer(), nullable=False, server_default="0"), + sa.Column("last_error", sa.Text(), nullable=True), + sa.Column("result_json", sa.Text(), nullable=True), + sa.Column("created_by_user_id", sa.Integer(), sa.ForeignKey("users.id", ondelete="SET NULL"), nullable=True), + sa.Column("created_at_utc", sa.DateTime(timezone=True), nullable=False), + sa.Column("updated_at_utc", sa.DateTime(timezone=True), nullable=False), + sa.Column("claimed_at_utc", sa.DateTime(timezone=True), nullable=True), + sa.Column("lease_expires_at_utc", sa.DateTime(timezone=True), nullable=True), + sa.Column("agent_completed_at_utc", sa.DateTime(timezone=True), nullable=True), + sa.Column("completed_at_utc", sa.DateTime(timezone=True), nullable=True), + sa.Column("failed_at_utc", sa.DateTime(timezone=True), nullable=True), + sa.Column("cancelled_at_utc", sa.DateTime(timezone=True), nullable=True), + sa.UniqueConstraint("job_uuid", name="uq_erp_agent_jobs_uuid"), + sa.UniqueConstraint("workstation_agent_id", "idempotency_key", name="uq_erp_agent_jobs_workstation_idempotency"), + ) + for name, cols in [ + ("ix_erp_agent_jobs_job_uuid", ["job_uuid"]), + ("ix_erp_agent_jobs_tenant_id", ["tenant_id"]), + ("ix_erp_agent_jobs_branch_id", ["branch_id"]), + ("ix_erp_agent_jobs_storage_node_id", ["storage_node_id"]), + ("ix_erp_agent_jobs_workstation_agent_id", ["workstation_agent_id"]), + ("ix_erp_agent_jobs_action", ["action"]), + ("ix_erp_agent_jobs_idempotency_key", ["idempotency_key"]), + ("ix_erp_agent_jobs_status", ["status"]), + ("ix_erp_agent_jobs_priority", ["priority"]), + ("ix_erp_agent_jobs_created_by_user_id", ["created_by_user_id"]), + ("ix_erp_agent_jobs_created_at_utc", ["created_at_utc"]), + ("ix_erp_agent_jobs_claimed_at_utc", ["claimed_at_utc"]), + ("ix_erp_agent_jobs_lease_expires_at_utc", ["lease_expires_at_utc"]), + ("ix_erp_agent_jobs_completed_at_utc", ["completed_at_utc"]), + ("ix_erp_agent_jobs_failed_at_utc", ["failed_at_utc"]), + ("ix_erp_agent_jobs_cancelled_at_utc", ["cancelled_at_utc"]), + ]: + op.create_index(name, "erp_agent_jobs", cols) + + op.create_table( + "erp_agent_job_events", + sa.Column("id", sa.Integer(), primary_key=True, autoincrement=True), + sa.Column("job_id", sa.Integer(), sa.ForeignKey("erp_agent_jobs.id", ondelete="CASCADE"), nullable=False), + 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("workstation_agent_id", sa.Integer(), sa.ForeignKey("erp_workstation_agents.id", ondelete="CASCADE"), nullable=False), + sa.Column("event_type", sa.String(length=40), nullable=False), + sa.Column("detail", sa.Text(), nullable=True), + sa.Column("actor_user_id", sa.Integer(), sa.ForeignKey("users.id", ondelete="SET NULL"), nullable=True), + sa.Column("occurred_at_utc", sa.DateTime(timezone=True), nullable=False), + ) + for name, cols in [ + ("ix_erp_agent_job_events_job_id", ["job_id"]), + ("ix_erp_agent_job_events_tenant_id", ["tenant_id"]), + ("ix_erp_agent_job_events_branch_id", ["branch_id"]), + ("ix_erp_agent_job_events_workstation_agent_id", ["workstation_agent_id"]), + ("ix_erp_agent_job_events_event_type", ["event_type"]), + ("ix_erp_agent_job_events_actor_user_id", ["actor_user_id"]), + ("ix_erp_agent_job_events_occurred_at_utc", ["occurred_at_utc"]), + ]: + op.create_index(name, "erp_agent_job_events", cols) + + +def downgrade() -> None: + op.drop_table("erp_agent_job_events") + op.drop_table("erp_agent_jobs") diff --git a/app/modules/documents/agent_jobs.py b/app/modules/documents/agent_jobs.py new file mode 100644 index 0000000..341761c --- /dev/null +++ b/app/modules/documents/agent_jobs.py @@ -0,0 +1,257 @@ +from __future__ import annotations + +import json +from datetime import datetime, timedelta, timezone +from typing import Any +from uuid import uuid4 + +from sqlalchemy import or_, select +from sqlalchemy.exc import IntegrityError + +from app.modules.documents.models import ERPAgentJob, ERPAgentJobEvent, ERPWorkstationAgent + + +LEASE_SECONDS = 45 +DEFAULT_MAX_ATTEMPTS = 3 +TERMINAL_STATUSES = {"succeeded", "failed", "cancelled"} + + +def _utcnow() -> datetime: + return datetime.now(timezone.utc) + + +def _json_dump(value: Any) -> str: + return json.dumps(value, ensure_ascii=False, separators=(",", ":"), default=str) + + +def _json_load(value: str | None, default: Any) -> Any: + if not value: + return default + try: + parsed = json.loads(value) + return parsed + except Exception: + return default + + +def _event(db, job: ERPAgentJob, event_type: str, *, detail: str | None = None, actor_user_id: int | None = None) -> None: + db.add( + ERPAgentJobEvent( + job_id=job.id, + tenant_id=job.tenant_id, + branch_id=job.branch_id, + workstation_agent_id=job.workstation_agent_id, + event_type=str(event_type or "event")[:40], + detail=(detail or None), + actor_user_id=actor_user_id, + occurred_at_utc=_utcnow(), + ) + ) + + +def enqueue_agent_job( + db, + *, + workstation_agent_id: int, + action: str, + payload: dict[str, Any] | None = None, + idempotency_key: str | None = None, + priority: int = 5, + max_attempts: int = DEFAULT_MAX_ATTEMPTS, + created_by_user_id: int | None = None, +) -> ERPAgentJob: + workstation = db.get(ERPWorkstationAgent, int(workstation_agent_id)) + if not workstation or not workstation.is_active: + raise ValueError("The selected ERP workstation agent is unavailable or inactive.") + action_value = str(action or "").strip() + if not action_value: + raise ValueError("Agent job action is required.") + idem = str(idempotency_key or "").strip() or None + if idem: + existing = db.execute( + select(ERPAgentJob).where( + ERPAgentJob.workstation_agent_id == workstation.id, + ERPAgentJob.idempotency_key == idem, + ) + ).scalar_one_or_none() + if existing: + return existing + job = ERPAgentJob( + job_uuid=uuid4().hex, + tenant_id=workstation.tenant_id, + branch_id=workstation.branch_id, + storage_node_id=workstation.storage_node_id, + workstation_agent_id=workstation.id, + action=action_value[:120], + payload_json=_json_dump(payload or {}), + idempotency_key=idem[:200] if idem else None, + status="queued", + priority=max(0, min(int(priority), 100)), + max_attempts=max(1, min(int(max_attempts), 20)), + attempts=0, + created_by_user_id=created_by_user_id, + created_at_utc=_utcnow(), + updated_at_utc=_utcnow(), + ) + db.add(job) + try: + db.flush() + except IntegrityError: + db.rollback() + if idem: + existing = db.execute( + select(ERPAgentJob).where( + ERPAgentJob.workstation_agent_id == workstation.id, + ERPAgentJob.idempotency_key == idem, + ) + ).scalar_one_or_none() + if existing: + return existing + raise + _event(db, job, "queued", actor_user_id=created_by_user_id) + return job + + +def claim_jobs_for_workstation(db, *, storage_node_id: int, agent_instance_id: str, limit: int = 10) -> list[dict[str, Any]]: + instance_id = str(agent_instance_id or "").strip() + if not instance_id: + return [] + workstation = db.execute( + select(ERPWorkstationAgent).where( + ERPWorkstationAgent.storage_node_id == int(storage_node_id), + ERPWorkstationAgent.agent_instance_id == instance_id, + ERPWorkstationAgent.is_active.is_(True), + ) + ).scalar_one_or_none() + if not workstation: + return [] + + now = _utcnow() + # Expired claims become retryable. Jobs that have exhausted their attempts fail. + expired = list( + db.execute( + select(ERPAgentJob).where( + ERPAgentJob.workstation_agent_id == workstation.id, + ERPAgentJob.status == "claimed", + ERPAgentJob.lease_expires_at_utc.is_not(None), + ERPAgentJob.lease_expires_at_utc < now, + ) + ).scalars().all() + ) + for job in expired: + if int(job.attempts or 0) >= int(job.max_attempts or DEFAULT_MAX_ATTEMPTS): + job.status = "failed" + job.failed_at_utc = now + job.last_error = job.last_error or "Agent job lease expired after maximum retry attempts." + _event(db, job, "failed", detail=job.last_error) + else: + job.status = "queued" + job.claimed_at_utc = None + job.lease_expires_at_utc = None + _event(db, job, "lease_expired", detail="Job returned to queue for retry.") + + stmt = ( + select(ERPAgentJob) + .where( + ERPAgentJob.workstation_agent_id == workstation.id, + ERPAgentJob.status == "queued", + ERPAgentJob.attempts < ERPAgentJob.max_attempts, + ) + .order_by(ERPAgentJob.priority.desc(), ERPAgentJob.created_at_utc.asc(), ERPAgentJob.id.asc()) + .limit(max(1, min(int(limit), 25))) + ) + rows = list(db.execute(stmt).scalars().all()) + payloads: list[dict[str, Any]] = [] + for job in rows: + job.status = "claimed" + job.attempts = int(job.attempts or 0) + 1 + job.claimed_at_utc = now + job.lease_expires_at_utc = now + timedelta(seconds=LEASE_SECONDS) + job.updated_at_utc = now + _event(db, job, "claimed", detail=f"attempt={job.attempts}") + payloads.append( + { + "job_uuid": job.job_uuid, + "action": job.action, + "payload": _json_load(job.payload_json, {}), + "idempotency_key": job.idempotency_key or job.job_uuid, + "attempt": int(job.attempts or 0), + "max_attempts": int(job.max_attempts or DEFAULT_MAX_ATTEMPTS), + "lease_seconds": LEASE_SECONDS, + } + ) + return payloads + + +def complete_agent_job( + db, + *, + storage_node_id: int, + agent_instance_id: str, + message: dict[str, Any], +) -> ERPAgentJob | None: + instance_id = str(agent_instance_id or "").strip() + job_uuid = str(message.get("job_uuid") or "").strip() + if not instance_id or not job_uuid: + return None + workstation = db.execute( + select(ERPWorkstationAgent).where( + ERPWorkstationAgent.storage_node_id == int(storage_node_id), + ERPWorkstationAgent.agent_instance_id == instance_id, + ERPWorkstationAgent.is_active.is_(True), + ) + ).scalar_one_or_none() + if not workstation: + return None + job = db.execute( + select(ERPAgentJob).where( + ERPAgentJob.job_uuid == job_uuid, + ERPAgentJob.workstation_agent_id == workstation.id, + ) + ).scalar_one_or_none() + if not job: + return None + if job.status in TERMINAL_STATUSES: + return job + + now = _utcnow() + ok = bool(message.get("ok")) + job.result_json = _json_dump(message.get("result") or {}) if ok else None + job.last_error = None if ok else str(message.get("error") or "Agent job failed.")[:4000] + job.agent_completed_at_utc = now + job.updated_at_utc = now + job.lease_expires_at_utc = None + if ok: + job.status = "succeeded" + job.completed_at_utc = now + _event(db, job, "succeeded") + elif int(job.attempts or 0) < int(job.max_attempts or DEFAULT_MAX_ATTEMPTS) and bool(message.get("retryable", False)): + job.status = "queued" + job.claimed_at_utc = None + _event(db, job, "retry_queued", detail=job.last_error) + else: + job.status = "failed" + job.failed_at_utc = now + _event(db, job, "failed", detail=job.last_error) + return job + + +def cancel_agent_job(db, job: ERPAgentJob, *, actor_user_id: int | None = None, reason: str | None = None) -> ERPAgentJob: + if job.status in TERMINAL_STATUSES: + return job + job.status = "cancelled" + job.cancelled_at_utc = _utcnow() + job.updated_at_utc = job.cancelled_at_utc + job.lease_expires_at_utc = None + _event(db, job, "cancelled", detail=reason, actor_user_id=actor_user_id) + return job + + +def list_recent_agent_jobs(db, *, tenant_id: int | None = None, branch_id: int | None = None, limit: int = 25) -> list[ERPAgentJob]: + stmt = select(ERPAgentJob).order_by(ERPAgentJob.created_at_utc.desc(), ERPAgentJob.id.desc()) + if tenant_id is not None: + stmt = stmt.where(ERPAgentJob.tenant_id == int(tenant_id)) + if branch_id is not None: + stmt = stmt.where(ERPAgentJob.branch_id == int(branch_id)) + stmt = stmt.limit(max(1, min(int(limit), 200))) + return list(db.execute(stmt).scalars().all()) diff --git a/app/modules/documents/agent_package.py b/app/modules/documents/agent_package.py index 5ce678e..ec32202 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.10.0" +ERP_LOCAL_AGENT_VERSION = "1.11.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 c9c5ecc..34d1752 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.10.0" +__version__ = "1.11.0" AGENT_NAME = "ERP Local Agent" diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/db.py b/app/modules/documents/local_agent_runtime/erp_local_agent/db.py index 8280964..46b8e3a 100644 --- a/app/modules/documents/local_agent_runtime/erp_local_agent/db.py +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/db.py @@ -71,6 +71,16 @@ class LocalDB: status TEXT NOT NULL DEFAULT 'received', error_message TEXT NULL ); + CREATE TABLE IF NOT EXISTS processed_agent_jobs ( + job_uuid TEXT PRIMARY KEY, + idempotency_key TEXT NOT NULL UNIQUE, + action TEXT NOT NULL, + execution_status TEXT NOT NULL DEFAULT 'processing', + ok INTEGER NULL, + result_json TEXT NULL, + error_message TEXT NULL, + processed_at_utc TEXT NOT NULL + ); CREATE TABLE IF NOT EXISTS update_history ( id INTEGER PRIMARY KEY AUTOINCREMENT, event_type TEXT NOT NULL, @@ -192,6 +202,42 @@ class LocalDB: ) conn.commit() + def get_processed_agent_job(self, job_uuid: str, idempotency_key: str) -> dict[str, Any] | None: + with self.connect() as conn: + row = conn.execute( + "SELECT * FROM processed_agent_jobs WHERE job_uuid=? OR idempotency_key=? ORDER BY processed_at_utc DESC LIMIT 1", + (str(job_uuid or ""), str(idempotency_key or "")), + ).fetchone() + if not row: + return None + result = dict(row) + try: + result["result"] = json.loads(result.get("result_json") or "{}") + except Exception: + result["result"] = {} + return result + + def begin_agent_job(self, job_uuid: str, idempotency_key: str, action: str) -> None: + with self.connect() as conn: + conn.execute( + """INSERT OR IGNORE INTO processed_agent_jobs + (job_uuid, idempotency_key, action, execution_status, ok, result_json, error_message, processed_at_utc) + VALUES (?, ?, ?, 'processing', NULL, NULL, NULL, ?)""", + (str(job_uuid), str(idempotency_key), str(action), _utc_now()), + ) + conn.commit() + + def record_processed_agent_job(self, job_uuid: str, idempotency_key: str, action: str, ok: bool, result: Any = None, error_message: str | None = None) -> None: + result_json = None if result is None else json.dumps(result, ensure_ascii=False, default=str) + with self.connect() as conn: + conn.execute( + """INSERT OR REPLACE INTO processed_agent_jobs + (job_uuid, idempotency_key, action, execution_status, ok, result_json, error_message, processed_at_utc) + VALUES (?, ?, ?, 'completed', ?, ?, ?, ?)""", + (str(job_uuid), str(idempotency_key), str(action), 1 if ok else 0, result_json, error_message, _utc_now()), + ) + conn.commit() + def record_update(self, event_type: str, from_version: str, to_version: str, status: str = "", sha256: str = "", package_path: str = "", error_message: str | None = None, details: Any = None) -> None: details_json = None if details is None else json.dumps(details, ensure_ascii=False, default=str) with self.connect() as conn: 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 08ef417..c35a82b 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 @@ -71,6 +71,7 @@ class StorageAgentTunnel: jobs = message.get("jobs") or [] requests_ = message.get("download_requests") or message.get("requests") or [] commands = message.get("commands") or [] + agent_jobs = message.get("agent_jobs") or [] await asyncio.to_thread(self.agent.process_storage_jobs_from_payload, jobs) await asyncio.to_thread(self.agent.process_download_requests_from_payload, requests_) for command in commands: @@ -80,11 +81,14 @@ class StorageAgentTunnel: result = await asyncio.to_thread(self.command_processor.process, command) self.db.command_completed(history_id, bool(result.get("ok")), result.get("error")) await websocket.send(self._json_text(result)) + for job in agent_jobs: + await self._process_agent_job(websocket, job) await websocket.send(self._json_text({ "type": "agent_status", "jobs_seen": len(jobs), "requests_seen": len(requests_), "commands_seen": len(commands), + "agent_jobs_seen": len(agent_jobs), "capabilities": self._capabilities(), "workstation": self._workstation_status(), "agent_time_utc": self._now(), @@ -93,9 +97,100 @@ class StorageAgentTunnel: if msg_type == "error": self.logger.error("Tunnel server error: %s", message.get("error")) + async def _process_agent_job(self, websocket, job: dict[str, Any]) -> None: + job_uuid = str(job.get("job_uuid") or "").strip() + idempotency_key = str(job.get("idempotency_key") or job_uuid).strip() + action = str(job.get("action") or "").strip() + if not job_uuid or not idempotency_key or not action: + return + + cached = self.db.get_processed_agent_job(job_uuid, idempotency_key) + if cached: + if str(cached.get("execution_status") or "").lower() == "processing": + await websocket.send(self._json_text({ + "type": "agent_job_result", + "job_uuid": job_uuid, + "idempotency_key": idempotency_key, + "ok": False, + "result": {}, + "error": "Previous execution was interrupted before acknowledgement. Verify Tally before retrying this job.", + "retryable": False, + "indeterminate": True, + "replayed": True, + "agent_time_utc": self._now(), + "workstation": self._workstation_status(), + })) + return + await websocket.send(self._json_text({ + "type": "agent_job_result", + "job_uuid": job_uuid, + "idempotency_key": idempotency_key, + "ok": bool(cached.get("ok")), + "result": cached.get("result") or {}, + "error": cached.get("error_message"), + "retryable": False, + "replayed": True, + "agent_time_utc": self._now(), + "workstation": self._workstation_status(), + })) + return + + # Persist 'processing' before touching Tally. If Windows/Tally/agent dies in + # the narrow post-before-ack window, the next connection refuses to blindly + # execute the same accounting job again. + self.db.begin_agent_job(job_uuid, idempotency_key, action) + history_id = self.db.command_received(job_uuid, action) + command = { + "command_id": job_uuid, + "action": action, + "payload": job.get("payload") or {}, + } + try: + result = await asyncio.to_thread(self.command_processor.process, command) + ok = bool(result.get("ok")) + error = result.get("error") + result_payload = result.get("result") or {} + self.db.command_completed(history_id, ok, error) + # Both success and deterministic application errors are cached so a + # reconnect cannot execute the same accounting write twice. + self.db.record_processed_agent_job( + job_uuid, idempotency_key, action, ok, result_payload, error + ) + await websocket.send(self._json_text({ + "type": "agent_job_result", + "job_uuid": job_uuid, + "idempotency_key": idempotency_key, + "ok": ok, + "result": result_payload, + "error": error, + "retryable": False, + "replayed": False, + "agent_time_utc": self._now(), + "workstation": self._workstation_status(), + })) + except Exception as exc: + error = str(exc) + self.db.command_completed(history_id, False, error) + # Execution has already crossed the durable local 'processing' marker. + # We cannot prove that Tally did not accept a write before this exception, + # therefore this is deliberately non-retryable and requires verification. + await websocket.send(self._json_text({ + "type": "agent_job_result", + "job_uuid": job_uuid, + "idempotency_key": idempotency_key, + "ok": False, + "result": {}, + "error": f"Agent execution became indeterminate: {error}", + "retryable": False, + "indeterminate": True, + "replayed": False, + "agent_time_utc": self._now(), + "workstation": self._workstation_status(), + })) + def _capabilities(self) -> list[str]: return [ - "storage", "local_dashboard", "manual_updates", "accounting_act", + "storage", "local_dashboard", "manual_updates", "accounting_act", "agent.jobs.v2", "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", diff --git a/app/modules/documents/models.py b/app/modules/documents/models.py index 5102195..8c98e60 100644 --- a/app/modules/documents/models.py +++ b/app/modules/documents/models.py @@ -218,6 +218,62 @@ class ERPWorkstationAgent(CommonBase): storage_node = relationship("BranchStorageNode") +class ERPAgentJob(CommonBase): + """Durable server-side job routed to one specific ERP workstation agent.""" + + __tablename__ = "erp_agent_jobs" + __table_args__ = ( + UniqueConstraint("job_uuid", name="uq_erp_agent_jobs_uuid"), + UniqueConstraint("workstation_agent_id", "idempotency_key", name="uq_erp_agent_jobs_workstation_idempotency"), + ) + + id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) + job_uuid: Mapped[str] = mapped_column(String(64), nullable=False, index=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) + workstation_agent_id: Mapped[int] = mapped_column(ForeignKey("erp_workstation_agents.id", ondelete="CASCADE"), nullable=False, index=True) + action: Mapped[str] = mapped_column(String(120), nullable=False, index=True) + payload_json: Mapped[str] = mapped_column(Text, nullable=False, default="{}") + idempotency_key: Mapped[str | None] = mapped_column(String(200), nullable=True, index=True) + status: Mapped[str] = mapped_column(String(30), nullable=False, default="queued", index=True) + priority: Mapped[int] = mapped_column(Integer, nullable=False, default=5, index=True) + max_attempts: Mapped[int] = mapped_column(Integer, nullable=False, default=3) + attempts: Mapped[int] = mapped_column(Integer, nullable=False, default=0) + last_error: Mapped[str | None] = mapped_column(Text, nullable=True) + result_json: Mapped[str | None] = mapped_column(Text, nullable=True) + created_by_user_id: Mapped[int | None] = mapped_column(ForeignKey("users.id", ondelete="SET NULL"), nullable=True, index=True) + created_at_utc: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=lambda: datetime.now(timezone.utc), nullable=False, index=True) + updated_at_utc: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=lambda: datetime.now(timezone.utc), nullable=False) + claimed_at_utc: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True, index=True) + lease_expires_at_utc: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True, index=True) + agent_completed_at_utc: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + completed_at_utc: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True, index=True) + failed_at_utc: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True, index=True) + cancelled_at_utc: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True, index=True) + + workstation = relationship("ERPWorkstationAgent") + storage_node = relationship("BranchStorageNode") + + +class ERPAgentJobEvent(CommonBase): + """Immutable audit event for an ERP workstation-agent job.""" + + __tablename__ = "erp_agent_job_events" + + id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) + job_id: Mapped[int] = mapped_column(ForeignKey("erp_agent_jobs.id", ondelete="CASCADE"), nullable=False, index=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) + workstation_agent_id: Mapped[int] = mapped_column(ForeignKey("erp_workstation_agents.id", ondelete="CASCADE"), nullable=False, index=True) + event_type: Mapped[str] = mapped_column(String(40), nullable=False, index=True) + detail: Mapped[str | None] = mapped_column(Text, nullable=True) + actor_user_id: Mapped[int | None] = mapped_column(ForeignKey("users.id", ondelete="SET NULL"), nullable=True, index=True) + occurred_at_utc: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=lambda: datetime.now(timezone.utc), nullable=False, index=True) + + job = relationship("ERPAgentJob") + + 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 a3a6ec9..1f51a32 100644 --- a/app/modules/documents/templates/documents/storage_nodes.html +++ b/app/modules/documents/templates/documents/storage_nodes.html @@ -177,7 +177,7 @@
- + {% for ws in workstation_agents or [] %} @@ -186,9 +186,44 @@ + {% else %} - + + {% endfor %} + +
WorkstationAgentTallyOpen CompaniesLast Seen
WorkstationAgentTallyOpen CompaniesLast SeenPhase 2
{{ 'Connected' if ws.tally_connected else 'Not connected' }} {{ ws.tally_company_count or 0 }} {{ ws.last_seen_at_utc or '-' }} + {% if can_manage_branch_storage %} +
+ + +
+ {% else %}-{% endif %} +
No workstation agent has connected since the Phase 1 upgrade.
No workstation agent has connected since the Phase 1 upgrade.
+
+ + + +
+
+

Recent Workstation Agent Jobs

+

Phase 2 jobs are durable, routed to one workstation, leased for safe retry, and protected by an idempotency key.

+
+
+ + + + {% for job in recent_agent_jobs or [] %} + + + + + + + + + {% else %} + {% endfor %}
JobWorkstationActionStatusAttemptsCreated
{{ job.job_uuid[:10] }}...{{ job.workstation.machine_name if job.workstation else ('#' ~ job.workstation_agent_id) }}{{ job.action }}{{ job.status }}{% if job.last_error %}
{{ job.last_error }}
{% endif %}
{{ job.attempts }}/{{ job.max_attempts }}{{ job.created_at_utc or '-' }}
No Phase 2 workstation jobs yet. Use Test route on a connected workstation to verify routing.
diff --git a/app/modules/documents/ui.py b/app/modules/documents/ui.py index a86fe52..68b4777 100644 --- a/app/modules/documents/ui.py +++ b/app/modules/documents/ui.py @@ -1,6 +1,7 @@ from __future__ import annotations from datetime import date, datetime, timezone +from uuid import uuid4 import json import asyncio import logging @@ -83,6 +84,7 @@ 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, ERPWorkstationAgent +from app.modules.documents.agent_jobs import claim_jobs_for_workstation, complete_agent_job, list_recent_agent_jobs, enqueue_agent_job from app.modules.accounting.agent_bridge import list_pending_agent_commands, record_agent_command_result from app.modules.services.task_documents import ( @@ -1319,6 +1321,7 @@ def storage_nodes(request: Request): 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) + agent_jobs = list_recent_agent_jobs(db, tenant_id=tenant_filter, branch_id=branch_filter, limit=20) return _render( request, "modules/documents/templates/documents/storage_nodes.html", @@ -1331,6 +1334,7 @@ def storage_nodes(request: Request): recent_jobs=jobs, recent_download_requests=requests, workstation_agents=workstations, + recent_agent_jobs=agent_jobs, generated_secret=None, storage_scope_title=_storage_scope_title(user, scope), erp_local_agent_version=ERP_LOCAL_AGENT_VERSION, @@ -1341,6 +1345,43 @@ def storage_nodes(request: Request): db.close() +@router.post("/storage-workstations/{workstation_id}/test-route") +def test_workstation_job_route(request: Request, workstation_id: int, csrf_token: str = Form(...)): + """Queue a harmless workstation-specific Tally status job for Phase 2 routing verification.""" + validate_csrf(request, csrf_token) + db = CommonSessionLocal() + try: + user, response = _require_user(request, db, "documents.upload") + if response: + return response + scope = build_document_scope(request, db, user) + if not _can_manage_branch_storage(scope): + return _redirect_denied() + workstation = db.get(ERPWorkstationAgent, int(workstation_id)) + if not workstation: + return _redirect_denied() + tenant_filter = _storage_tenant_filter(user, scope) + branch_filter = _storage_branch_filter(user, scope) + if tenant_filter is not None and int(workstation.tenant_id) != int(tenant_filter): + return _redirect_denied() + if branch_filter is not None and int(workstation.branch_id or 0) != int(branch_filter): + return _redirect_denied() + enqueue_agent_job( + db, + workstation_agent_id=workstation.id, + action="tally_status", + payload={}, + idempotency_key=f"route-test:{workstation.id}:{uuid4().hex}", + created_by_user_id=int(user.id), + priority=10, + max_attempts=3, + ) + db.commit() + return RedirectResponse(url="/documents/storage-nodes", status_code=303) + finally: + db.close() + + @router.get("/branch-storage-dashboard") def partner_branch_storage_dashboard(request: Request): """Partner/Branch dashboard card for local storage setup and status. @@ -1737,15 +1778,28 @@ def _permanent_download_requests_payload(items): ] -def _storage_agent_sync_payload(db, node): +def _storage_agent_sync_payload(db, node, workstation_instance_id: str | None = None): # Retry cleanup in tunnel mode as well as normal heartbeat mode. retry_verified_vps_cleanup_for_node(db, node, limit=25) jobs = _normal_storage_jobs_payload(list_pending_storage_jobs(db, node)) jobs += _permanent_storage_jobs_payload(list_pending_permanent_storage_jobs(db, node)) requests_ = _normal_download_requests_payload(list_pending_download_requests(db, node)) requests_ += _permanent_download_requests_payload(list_pending_permanent_download_requests(db, node)) + # Legacy node-level commands are retained unchanged for backward compatibility. commands = list_pending_agent_commands(node.node_code, limit=10) - return {"ok": True, "jobs": jobs, "download_requests": requests_, "requests": requests_, "commands": commands} + agent_jobs = [] + if workstation_instance_id: + agent_jobs = claim_jobs_for_workstation( + db, storage_node_id=node.id, agent_instance_id=workstation_instance_id, limit=10 + ) + return { + "ok": True, + "jobs": jobs, + "download_requests": requests_, + "requests": requests_, + "commands": commands, + "agent_jobs": agent_jobs, + } @router.get("/storage-agent/jobs/pending") @@ -1875,6 +1929,7 @@ async def storage_agent_tunnel(websocket: WebSocket): db.close() last_push = 0.0 + connection_workstation_instance_id = "" while True: try: try: @@ -1891,10 +1946,21 @@ async def storage_agent_tunnel(websocket: WebSocket): 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) + workstation_payload = message.get("workstation") + _record_workstation_status(db, node, workstation_payload, client_ip) + advertised_instance_id = str(workstation_payload.get("agent_instance_id") or "").strip() + if advertised_instance_id: + connection_workstation_instance_id = advertised_instance_id 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"}: + if message and message.get("type") == "agent_job_result": + complete_agent_job( + db, + storage_node_id=node.id, + agent_instance_id=connection_workstation_instance_id, + message=message, + ) + if message and message.get("type") in {"heartbeat", "agent_status", "command_result", "agent_job_result"}: try: node.storage_mode = "tunnel" except Exception: @@ -1902,7 +1968,7 @@ async def storage_agent_tunnel(websocket: WebSocket): now = datetime.now(timezone.utc).timestamp() force = bool(message and message.get("type") in {"ready", "sync_now", "agent_status"}) if force or now - last_push >= 5: - payload = _storage_agent_sync_payload(db, node) + payload = _storage_agent_sync_payload(db, node, connection_workstation_instance_id) payload.update({"type": "sync", "node_code": node.node_code}) await websocket.send_json(payload) last_push = now