Add operator workstation GST browser and full FY return download

This commit is contained in:
A R R R Associates
2026-09-10 22:18:32 +05:30
parent 90d7769bee
commit 0ce2954674
7 changed files with 497 additions and 115 deletions
@@ -1,2 +1,2 @@
__version__ = "1.26.14"
__version__ = "1.26.15"
AGENT_NAME = "ERP Local Agent"
@@ -6,6 +6,10 @@ import time
import threading
import json
import uuid
import shutil
import zipfile
import tempfile
import requests
from pathlib import Path
from typing import Any
@@ -189,6 +193,8 @@ class AgentCommandProcessor:
result = self._gst_return_download_status(payload)
elif action == "gst_reconciliation_analyze":
result = self._gst_reconciliation_analyze(payload)
elif action == "gst_return_package_store":
result = self._gst_return_package_store(payload)
elif action == "accounting_analysis_save":
result = self._accounting_analysis_save(payload)
elif action == "accounting_analysis_history":
@@ -600,6 +606,52 @@ class AgentCommandProcessor:
result = gst_download_status(payload, Path(__file__).resolve().parents[1] / "data")
return {"job": result, "agent": self._agent_info()}
def _gst_return_package_store(self, payload: dict[str, Any]) -> dict[str, Any]:
client_id = int(payload.get("client_id") or 0)
gstin = str(payload.get("gstin") or "").strip().upper()
gst_relative_dir = str(payload.get("gst_relative_dir") or "").strip()
package_url = str(payload.get("package_url") or "").strip()
if client_id <= 0 or len(gstin) != 15 or not gst_relative_dir or not package_url:
raise ValueError("Client, GSTIN, GST storage path and package URL are required.")
root = Path(self.config.storage_root).resolve()
target = (root / Path(gst_relative_dir)).resolve()
if root != target and root not in target.parents:
raise ValueError("GST storage path is outside the configured storage root.")
target.mkdir(parents=True, exist_ok=True)
temp_dir = Path(tempfile.mkdtemp(prefix="gst_store_"))
package_path = temp_dir / "gst_returns.zip"
try:
with requests.get(package_url, stream=True, timeout=180) as response:
response.raise_for_status()
total = 0
with package_path.open("wb") as handle:
for chunk in response.iter_content(chunk_size=1024 * 1024):
if not chunk:
continue
total += len(chunk)
if total > 250 * 1024 * 1024:
raise ValueError("GST return package exceeds the 250 MB safety limit.")
handle.write(chunk)
stored=[]
with zipfile.ZipFile(package_path, "r") as archive:
for info in archive.infolist():
if info.is_dir():
continue
rel = Path(info.filename.replace("\\", "/"))
if rel.is_absolute() or ".." in rel.parts:
raise ValueError(f"Unsafe GST package path: {info.filename}")
destination = (target / rel).resolve()
if target != destination and target not in destination.parents:
raise ValueError(f"Unsafe GST package destination: {info.filename}")
destination.parent.mkdir(parents=True, exist_ok=True)
with archive.open(info, "r") as src, destination.open("wb") as dst:
shutil.copyfileobj(src, dst)
stored.append(str(destination))
return {"stored_files": len(stored), "gst_storage_path": str(target), "periods": payload.get("periods") or [], "return_types": payload.get("return_types") or [], "agent": self._agent_info()}
finally:
shutil.rmtree(temp_dir, ignore_errors=True)
def _gst_reconciliation_analyze(self, payload: dict[str, Any]) -> dict[str, Any]:
client_id = int(payload.get("client_id") or 0)
period = re.sub(r"\D", "", str(payload.get("period") or ""))
@@ -9,9 +9,11 @@ from http import HTTPStatus
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path
from typing import Any
from urllib.parse import parse_qs, urlsplit
from . import AGENT_NAME, __version__
from .tally import TallyLiveConnector
from .gst_portal_runtime import start_operator_download, operator_download_status
class AgentDashboard:
@@ -84,6 +86,16 @@ class AgentDashboard:
self.send_response(status)
self.send_header("Content-Type", "application/json; charset=utf-8")
self.send_header("Cache-Control", "no-store")
origin = self.headers.get("Origin") or ""
erp_origin = ""
try:
parts = urlsplit(str(dashboard.config.erp_base_url or ""))
erp_origin = f"{parts.scheme}://{parts.netloc}" if parts.scheme and parts.netloc else ""
except Exception:
erp_origin = ""
if origin and erp_origin and origin.rstrip("/") == erp_origin.rstrip("/"):
self.send_header("Access-Control-Allow-Origin", origin)
self.send_header("Vary", "Origin")
self.send_header("Content-Length", str(len(data)))
self.end_headers()
self.wfile.write(data)
@@ -97,6 +109,22 @@ class AgentDashboard:
self.end_headers()
self.wfile.write(data)
def do_OPTIONS(self):
origin = self.headers.get("Origin") or ""
erp_origin = ""
try:
parts = urlsplit(str(dashboard.config.erp_base_url or ""))
erp_origin = f"{parts.scheme}://{parts.netloc}" if parts.scheme and parts.netloc else ""
except Exception:
erp_origin = ""
self.send_response(204)
if origin and erp_origin and origin.rstrip("/") == erp_origin.rstrip("/"):
self.send_header("Access-Control-Allow-Origin", origin)
self.send_header("Vary", "Origin")
self.send_header("Access-Control-Allow-Methods", "GET, POST, OPTIONS")
self.send_header("Access-Control-Allow-Headers", "Content-Type")
self.end_headers()
def do_GET(self):
if self.path == "/" or self.path.startswith("/?"):
return self._html(dashboard._page())
@@ -106,6 +134,11 @@ class AgentDashboard:
return self._json(dashboard.history())
if self.path.startswith("/api/mirror-progress"):
return self._json(dashboard.mirror_progress())
if self.path.startswith("/api/gst/operator/status"):
query=parse_qs(urlsplit(self.path).query)
job_id=str((query.get("job_id") or [""])[0])
data_root=Path(__file__).resolve().parents[1] / "data"
return self._json({"ok":True,"job":operator_download_status(job_id,data_root)})
if self.path.split("?", 1)[0] == "/tally/ERP_Accounting_Mirror.tdl":
info = dashboard._ensure_managed_tdl()
payload = Path(info["path"]).read_bytes()
@@ -121,6 +154,14 @@ class AgentDashboard:
def do_POST(self):
try:
if self.path == "/api/gst/operator/start":
length=int(self.headers.get("Content-Length") or 0)
raw=self.rfile.read(length) if length else b"{}"
body=json.loads(raw.decode("utf-8") or "{}")
token=str(body.get("token") or "")
data_root=Path(__file__).resolve().parents[1] / "data"
result=start_operator_download(token,str(dashboard.config.erp_base_url or ""),data_root)
return self._json({"ok":True,**result})
if self.path == "/api/update/check":
return self._json({"ok": True, "update": dashboard.updater.check_for_update(force=True)})
if self.path == "/api/update/download":
@@ -5,6 +5,9 @@ import re
import sqlite3
import threading
import time
import shutil
import zipfile
import requests
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
@@ -343,3 +346,146 @@ def start_download(payload: dict[str, Any], storage_root: Path, progress_root: P
def download_status(payload: dict[str, Any], progress_root: Path) -> dict[str, Any]:
key=_safe(f"{payload.get('client_id')}_{payload.get('gstin')}_{payload.get('period')}")
return _read_json(progress_root/f"gst_download_{key}.json")
# ---------------------------------------------------------------------------
# Interactive GST browser mode (operator workstation)
# ---------------------------------------------------------------------------
_OPERATOR_THREADS: dict[str, threading.Thread] = {}
def _operator_job_path(progress_root: Path, job_id: str) -> Path:
return progress_root / f"gst_operator_{_safe(job_id)}.json"
def _fetch_text(page, url: str) -> str:
return page.evaluate("async (u)=>{const r=await fetch(u,{credentials:'include'}); return await r.text();}", url) or ""
def _download_selected_period(page, work_root: Path, period: str, return_types: list[str], progress) -> dict[str, Any]:
base_dir = work_root / period
raw_dir = base_dir / "raw"
raw_dir.mkdir(parents=True, exist_ok=True)
downloaded=[]
selected={str(x or '').upper() for x in return_types}
if "GSTR1" in selected:
progress(stage=f"Downloading GSTR-1 {period}", message=f"Reading GSTR-1 for {period}.")
generate_text=_fetch_text(page,GSTR1_URL.format(period=period))
(raw_dir/f"{period}_GSTR1_GENERATE.json").write_text(generate_text,encoding="utf-8",errors="ignore")
file_num=_extract_file_num(generate_text)
g1_text=_extract_payload(_fetch_text(page,GSTR1_DOWNLOAD_URL.format(period=period,file_num=file_num)))
g1_path=raw_dir/f"{period}_GSTR1.json"; g1_path.write_text(g1_text,encoding="utf-8",errors="ignore")
downloaded.append({"return_type":"GSTR1","path":str(g1_path.relative_to(work_root)),"bytes":g1_path.stat().st_size})
if "GSTR2B" in selected:
progress(stage=f"Downloading GSTR-2B {period}", message=f"Reading GSTR-2B for {period}.")
g2_text=_fetch_text(page,GSTR2B_URL.format(period=period))
g2_path=raw_dir/f"{period}_GSTR2B.json"; g2_path.write_text(g2_text,encoding="utf-8",errors="ignore")
downloaded.append({"return_type":"GSTR2B","path":str(g2_path.relative_to(work_root)),"bytes":g2_path.stat().st_size})
if "GSTR2A" in selected:
progress(stage=f"Requesting GSTR-2A {period}", message=f"Requesting GSTR-2A for {period}.")
g2a_text=_fetch_text(page,GSTR2A_URL.format(period=period))
g2a_path=raw_dir/f"{period}_GSTR2A.json"; g2a_path.write_text(g2a_text,encoding="utf-8",errors="ignore")
downloaded.append({"return_type":"GSTR2A_GENERATE","path":str(g2a_path.relative_to(work_root)),"bytes":g2a_path.stat().st_size})
if "GSTR3B" in selected:
progress(stage=f"Downloading GSTR-3B {period}", message=f"Reading GSTR-3B for {period}.")
g3_summary=_fetch_text(page,GSTR3B_SUMMARY_URL.format(period=period))
g3_payable=_fetch_text(page,GSTR3B_URL.format(period=period))
try: combined=json.loads(g3_summary or "{}")
except Exception: combined={"summary_raw":g3_summary or ""}
if not isinstance(combined,dict): combined={"summary":combined}
try: combined["taxpayble"]=json.loads(g3_payable or "{}")
except Exception: combined["taxpayble"]={"raw":g3_payable or ""}
(raw_dir/f"{period}_GSTR3B_SUMMARY.json").write_text(g3_summary,encoding="utf-8",errors="ignore")
(raw_dir/f"{period}_GSTR3B_TAXPAYBLE.json").write_text(g3_payable,encoding="utf-8",errors="ignore")
g3_path=raw_dir/f"{period}_GSTR3B.json"; g3_path.write_text(json.dumps(combined,ensure_ascii=False,indent=2),encoding="utf-8")
downloaded.append({"return_type":"GSTR3B","path":str(g3_path.relative_to(work_root)),"bytes":g3_path.stat().st_size})
normalized=normalize_downloads(base_dir,period)
_write_json(base_dir/"download_manifest.json",{"period":period,"downloaded_at_utc":_now(),"downloaded":downloaded,"normalized":normalized})
return {"period":period,"downloaded":downloaded,"normalized":normalized}
def _operator_browser_worker(payload: dict[str, Any], progress_path: Path, progress_root: Path) -> None:
job_id=str(payload.get("jti") or payload.get("job_id") or "job")
work_root=progress_root / "gst_operator_work" / _safe(job_id)
try:
if work_root.exists(): shutil.rmtree(work_root,ignore_errors=True)
work_root.mkdir(parents=True,exist_ok=True)
def progress(**kw):
cur=_read_json(progress_path); cur.update(kw); cur["updated_at_utc"]=_now(); _write_json(progress_path,cur)
progress(status="running",stage="Opening GST Login",message="Opening GST Portal on this computer. Complete CAPTCHA/OTP in the visible browser.")
from playwright.sync_api import sync_playwright
with sync_playwright() as pw:
browser=None; errors=[]
profile=progress_root/"gst_operator_browser_profile"; profile.mkdir(parents=True,exist_ok=True)
for channel in ("chrome","msedge"):
try:
browser=pw.chromium.launch_persistent_context(user_data_dir=str(profile/channel),channel=channel,headless=False,no_viewport=True,accept_downloads=True,args=["--start-maximized","--no-first-run","--disable-blink-features=AutomationControlled"])
break
except Exception as exc: errors.append(f"{channel}: {exc}")
if browser is None: raise RuntimeError("Could not open Chrome/Edge on this workstation. "+" | ".join(errors))
page=browser.pages[0] if browser.pages else browser.new_page()
page.goto(GST_LOGIN_URL,wait_until="domcontentloaded",timeout=90000)
page.wait_for_selector("#username",timeout=30000)
page.fill("#username",str(payload.get("username") or "")); page.fill("#user_pass",str(payload.get("password") or ""))
progress(stage="Waiting for GST Login",message="Username/password filled. Complete CAPTCHA/OTP in this browser; download will continue automatically after login.")
deadline=time.time()+int(payload.get("login_timeout_seconds") or 900)
while time.time()<deadline:
try:
if "/dashboard" in page.url.lower() or await_login_probe(page): break
except Exception: pass
time.sleep(2)
else: raise RuntimeError("GST login was not completed within 15 minutes.")
results=[]; periods=list(payload.get("periods") or []); return_types=list(payload.get("return_types") or [])
for index,period in enumerate(periods,1):
progress(stage=f"Downloading period {index}/{len(periods)}",message=f"Downloading selected returns for {period}.",period=period,period_index=index,period_total=len(periods))
results.append(_download_selected_period(page,work_root,str(period),return_types,progress))
browser.close()
package_path=progress_root/f"gst_operator_{_safe(job_id)}.zip"
if package_path.exists(): package_path.unlink()
with zipfile.ZipFile(package_path,"w",zipfile.ZIP_DEFLATED,compresslevel=6) as z:
for f in sorted(work_root.rglob("*")):
if f.is_file(): z.write(f,f.relative_to(work_root).as_posix())
progress(stage="Transferring to Client Storage",message="Uploading downloaded GST return package to ERP for transfer to the configured Local Storage Agent.")
with package_path.open("rb") as handle:
response=requests.post(str(payload.get("upload_url") or ""),data={"token":str(payload.get("job_token") or "")},files={"package":(package_path.name,handle,"application/zip")},timeout=180)
try: body=response.json()
except Exception: body={"ok":False,"error":response.text[:1000]}
if not response.ok or not body.get("ok"): raise RuntimeError(body.get("error") or f"ERP package transfer failed with HTTP {response.status_code}.")
progress(status="completed",stage="Completed",message="GST returns downloaded on this computer and stored in the configured client local storage.",result={"periods":periods,"return_types":return_types,"storage":body.get("stored") or {}},finished_at_utc=_now())
except Exception as exc:
cur=_read_json(progress_path); cur.update({"status":"failed","stage":"Failed","message":str(exc),"error":str(exc),"finished_at_utc":_now(),"updated_at_utc":_now()}); _write_json(progress_path,cur)
finally:
try: shutil.rmtree(work_root,ignore_errors=True)
except Exception: pass
try:
pkg=progress_root/f"gst_operator_{_safe(job_id)}.zip"
pkg.unlink(missing_ok=True)
except Exception: pass
def start_operator_download(job_token: str, erp_base_url: str, progress_root: Path) -> dict[str, Any]:
if not job_token: raise ValueError("GST operator job token is required.")
redeem_url=str(erp_base_url or "").rstrip("/")+"/tools/accounting/gst-reconciliation/operator/redeem"
response=requests.post(redeem_url,json={"token":job_token},timeout=30)
try: body=response.json()
except Exception: body={"ok":False,"error":response.text[:1000]}
if not response.ok or not body.get("ok"): raise RuntimeError(body.get("error") or "ERP could not authorize the GST operator job.")
payload=dict(body.get("payload") or {}); payload["job_token"]=job_token
job_id=str(payload.get("jti") or "")
if not job_id: raise RuntimeError("ERP did not return a GST operator job id.")
progress_path=_operator_job_path(progress_root,job_id)
current=_read_json(progress_path)
if current.get("status") in {"queued","running"}: return {"started":False,"already_running":True,"job":current}
initial={"status":"queued","stage":"Queued","message":"Interactive GST browser job queued on this computer.","job_id":job_id,"periods":payload.get("periods") or [],"return_types":payload.get("return_types") or [],"started_at_utc":_now(),"updated_at_utc":_now()}
_write_json(progress_path,initial)
thread=threading.Thread(target=_operator_browser_worker,args=(payload,progress_path,progress_root),daemon=True,name=f"gst-operator-{job_id}")
_OPERATOR_THREADS[job_id]=thread; thread.start()
return {"started":True,"job":initial}
def operator_download_status(job_id: str, progress_root: Path) -> dict[str, Any]:
return _read_json(_operator_job_path(progress_root,job_id))