|
|
|
@@ -0,0 +1,370 @@
|
|
|
|
|
from __future__ import annotations
|
|
|
|
|
|
|
|
|
|
import json
|
|
|
|
|
import re
|
|
|
|
|
import shutil
|
|
|
|
|
import threading
|
|
|
|
|
import time
|
|
|
|
|
import zipfile
|
|
|
|
|
from datetime import datetime, timezone
|
|
|
|
|
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
|
|
|
|
from pathlib import Path
|
|
|
|
|
from urllib.parse import parse_qs, urlsplit
|
|
|
|
|
|
|
|
|
|
import requests
|
|
|
|
|
|
|
|
|
|
VERSION = "1.0.0"
|
|
|
|
|
PORT = 8791
|
|
|
|
|
ROOT = Path(__file__).resolve().parent
|
|
|
|
|
DATA = ROOT / "data"
|
|
|
|
|
CONFIG_PATH = ROOT / "config.json"
|
|
|
|
|
GST_LOGIN_URL = "https://services.gst.gov.in/services/login"
|
|
|
|
|
GSTR1_URL = "https://return.gst.gov.in/returns/auth/api/offline/download/generate?flag=0&rtn_prd={period}&rtn_typ=GSTR1"
|
|
|
|
|
GSTR2A_URL = "https://return.gst.gov.in/returns/auth/api/offline/download/generate?flag=0&rtn_prd={period}&rtn_typ=GSTR2A"
|
|
|
|
|
GSTR1_DOWNLOAD_URL = "https://return.gst.gov.in/returns/auth/api/offline/download/url?rtn_prd={period}&rtn_typ=GSTR1&file_num={file_num}"
|
|
|
|
|
GSTR2B_URL = "https://gstr2b.gst.gov.in/gstr2b/auth/api/gstr2b/getjson?rtnprd={period}"
|
|
|
|
|
GSTR3B_SUMMARY_URL = "https://return.gst.gov.in/returns/auth/api/gstr3b/summary?rtn_prd={period}"
|
|
|
|
|
GSTR3B_URL = "https://return.gst.gov.in/returns/auth/api/gstr3b/taxpayble?rtn_prd={period}"
|
|
|
|
|
THREADS: dict[str, threading.Thread] = {}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def now() -> str:
|
|
|
|
|
return datetime.now(timezone.utc).isoformat()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def read_config() -> dict:
|
|
|
|
|
try:
|
|
|
|
|
return json.loads(CONFIG_PATH.read_text(encoding="utf-8"))
|
|
|
|
|
except Exception:
|
|
|
|
|
return {"erp_base_url": "https://office.arrrassociates.com", "port": PORT}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def safe(value: str, fallback: str = "item") -> str:
|
|
|
|
|
value = re.sub(r"[^A-Za-z0-9._ -]+", "_", str(value or "").strip()).strip(" ._")
|
|
|
|
|
return value or fallback
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def write_json(path: Path, payload) -> None:
|
|
|
|
|
path.parent.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
tmp = path.with_suffix(path.suffix + ".tmp")
|
|
|
|
|
tmp.write_text(json.dumps(payload, ensure_ascii=False, indent=2, default=str), encoding="utf-8")
|
|
|
|
|
tmp.replace(path)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def read_json(path: Path) -> dict:
|
|
|
|
|
try:
|
|
|
|
|
value = json.loads(path.read_text(encoding="utf-8"))
|
|
|
|
|
return value if isinstance(value, dict) else {}
|
|
|
|
|
except Exception:
|
|
|
|
|
return {}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def job_path(job_id: str) -> Path:
|
|
|
|
|
return DATA / f"gst_job_{safe(job_id)}.json"
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def walk(obj):
|
|
|
|
|
if isinstance(obj, dict):
|
|
|
|
|
yield obj
|
|
|
|
|
for value in obj.values():
|
|
|
|
|
yield from walk(value)
|
|
|
|
|
elif isinstance(obj, list):
|
|
|
|
|
for value in obj:
|
|
|
|
|
yield from walk(value)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def extract_file_num(text: str) -> str:
|
|
|
|
|
try:
|
|
|
|
|
data = json.loads(text or "{}")
|
|
|
|
|
except Exception:
|
|
|
|
|
return "1"
|
|
|
|
|
for obj in walk(data):
|
|
|
|
|
if isinstance(obj, dict):
|
|
|
|
|
for key in ("file_num", "fileNum", "filenum", "file_number", "fileNumber"):
|
|
|
|
|
if obj.get(key) not in (None, ""):
|
|
|
|
|
return str(obj.get(key))
|
|
|
|
|
return "1"
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def extract_payload(text: str) -> str:
|
|
|
|
|
try:
|
|
|
|
|
data = json.loads(text or "{}")
|
|
|
|
|
except Exception:
|
|
|
|
|
return text or ""
|
|
|
|
|
if isinstance(data, dict):
|
|
|
|
|
for key in ("data", "payload", "json", "response"):
|
|
|
|
|
value = data.get(key)
|
|
|
|
|
if isinstance(value, (dict, list)):
|
|
|
|
|
return json.dumps(value, ensure_ascii=False)
|
|
|
|
|
if isinstance(value, str) and value.strip().startswith(("{", "[")):
|
|
|
|
|
return value
|
|
|
|
|
return text or ""
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def api_text(context, url: str) -> str:
|
|
|
|
|
response = context.request.get(url, timeout=90000)
|
|
|
|
|
if not response.ok:
|
|
|
|
|
raise RuntimeError(f"GST portal request failed ({response.status}) for {url.split('?')[0]}")
|
|
|
|
|
return response.text()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def download_period(context, work_root: Path, period: str, return_types: list[str], progress) -> list[dict]:
|
|
|
|
|
raw_dir = work_root / period / "raw"
|
|
|
|
|
raw_dir.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
selected = {str(x or "").upper() for x in return_types}
|
|
|
|
|
downloaded: list[dict] = []
|
|
|
|
|
|
|
|
|
|
if "GSTR1" in selected:
|
|
|
|
|
progress(stage=f"Downloading GSTR-1 {period}", message=f"Reading GSTR-1 for {period}.")
|
|
|
|
|
generated = api_text(context, GSTR1_URL.format(period=period))
|
|
|
|
|
(raw_dir / f"{period}_GSTR1_GENERATE.json").write_text(generated, encoding="utf-8")
|
|
|
|
|
file_num = extract_file_num(generated)
|
|
|
|
|
content = extract_payload(api_text(context, GSTR1_DOWNLOAD_URL.format(period=period, file_num=file_num)))
|
|
|
|
|
path = raw_dir / f"{period}_GSTR1.json"
|
|
|
|
|
path.write_text(content, encoding="utf-8")
|
|
|
|
|
downloaded.append({"return_type": "GSTR1", "path": str(path.relative_to(work_root)), "bytes": path.stat().st_size})
|
|
|
|
|
|
|
|
|
|
if "GSTR2B" in selected:
|
|
|
|
|
progress(stage=f"Downloading GSTR-2B {period}", message=f"Reading GSTR-2B for {period}.")
|
|
|
|
|
content = api_text(context, GSTR2B_URL.format(period=period))
|
|
|
|
|
path = raw_dir / f"{period}_GSTR2B.json"
|
|
|
|
|
path.write_text(content, encoding="utf-8")
|
|
|
|
|
downloaded.append({"return_type": "GSTR2B", "path": str(path.relative_to(work_root)), "bytes": path.stat().st_size})
|
|
|
|
|
|
|
|
|
|
if "GSTR3B" in selected:
|
|
|
|
|
progress(stage=f"Downloading GSTR-3B {period}", message=f"Reading GSTR-3B for {period}.")
|
|
|
|
|
summary = api_text(context, GSTR3B_SUMMARY_URL.format(period=period))
|
|
|
|
|
payable = api_text(context, GSTR3B_URL.format(period=period))
|
|
|
|
|
(raw_dir / f"{period}_GSTR3B_SUMMARY.json").write_text(summary, encoding="utf-8")
|
|
|
|
|
(raw_dir / f"{period}_GSTR3B_TAXPAYBLE.json").write_text(payable, encoding="utf-8")
|
|
|
|
|
try:
|
|
|
|
|
combined = json.loads(summary or "{}")
|
|
|
|
|
except Exception:
|
|
|
|
|
combined = {"summary_raw": summary}
|
|
|
|
|
if not isinstance(combined, dict):
|
|
|
|
|
combined = {"summary": combined}
|
|
|
|
|
try:
|
|
|
|
|
combined["taxpayble"] = json.loads(payable or "{}")
|
|
|
|
|
except Exception:
|
|
|
|
|
combined["taxpayble"] = {"raw": payable}
|
|
|
|
|
path = raw_dir / f"{period}_GSTR3B.json"
|
|
|
|
|
path.write_text(json.dumps(combined, ensure_ascii=False, indent=2), encoding="utf-8")
|
|
|
|
|
downloaded.append({"return_type": "GSTR3B", "path": str(path.relative_to(work_root)), "bytes": path.stat().st_size})
|
|
|
|
|
|
|
|
|
|
if "GSTR2A" in selected:
|
|
|
|
|
progress(stage=f"Requesting GSTR-2A {period}", message=f"Requesting GSTR-2A for {period}.")
|
|
|
|
|
content = api_text(context, GSTR2A_URL.format(period=period))
|
|
|
|
|
path = raw_dir / f"{period}_GSTR2A.json"
|
|
|
|
|
path.write_text(content, encoding="utf-8")
|
|
|
|
|
downloaded.append({"return_type": "GSTR2A", "path": str(path.relative_to(work_root)), "bytes": path.stat().st_size})
|
|
|
|
|
|
|
|
|
|
write_json(work_root / period / "download_manifest.json", {
|
|
|
|
|
"period": period, "downloaded_at_utc": now(), "source": "gst_lightweight_operator_agent", "downloaded": downloaded
|
|
|
|
|
})
|
|
|
|
|
return downloaded
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def login_complete(page) -> bool:
|
|
|
|
|
url = (page.url or "").lower()
|
|
|
|
|
if "/dashboard" in url or "/returns" in url:
|
|
|
|
|
return True
|
|
|
|
|
try:
|
|
|
|
|
return bool(page.locator("text=Search Taxpayer").count() and not page.locator("#username").count())
|
|
|
|
|
except Exception:
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def browser_worker(payload: dict, token: str, path: Path) -> None:
|
|
|
|
|
job_id = str(payload.get("jti") or "")
|
|
|
|
|
work_root = DATA / "work" / safe(job_id)
|
|
|
|
|
package_path = DATA / f"gst_{safe(job_id)}.zip"
|
|
|
|
|
try:
|
|
|
|
|
shutil.rmtree(work_root, ignore_errors=True)
|
|
|
|
|
work_root.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
|
|
|
|
|
def progress(**updates):
|
|
|
|
|
current = read_json(path)
|
|
|
|
|
current.update(updates)
|
|
|
|
|
current["updated_at_utc"] = now()
|
|
|
|
|
write_json(path, current)
|
|
|
|
|
|
|
|
|
|
progress(status="running", percent=5, stage="Opening GST Login", message="Opening GST Portal in visible Chrome/Edge on this workstation.")
|
|
|
|
|
from playwright.sync_api import sync_playwright
|
|
|
|
|
with sync_playwright() as pw:
|
|
|
|
|
context = None
|
|
|
|
|
errors = []
|
|
|
|
|
profile_root = DATA / "browser_profile"
|
|
|
|
|
profile_root.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
for channel in ("chrome", "msedge"):
|
|
|
|
|
try:
|
|
|
|
|
context = pw.chromium.launch_persistent_context(
|
|
|
|
|
user_data_dir=str(profile_root / 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 context is None:
|
|
|
|
|
raise RuntimeError("Could not open installed Chrome or Edge. " + " | ".join(errors))
|
|
|
|
|
page = context.pages[0] if context.pages else context.new_page()
|
|
|
|
|
page.goto(GST_LOGIN_URL, wait_until="domcontentloaded", timeout=90000)
|
|
|
|
|
page.locator("#username").wait_for(state="visible", timeout=30000)
|
|
|
|
|
page.fill("#username", str(payload.get("username") or ""))
|
|
|
|
|
# GST portal has used user_pass for the password field; fall back to generic password locator.
|
|
|
|
|
password_filled = False
|
|
|
|
|
for selector in ("#user_pass", "input[type=password]"):
|
|
|
|
|
try:
|
|
|
|
|
locator = page.locator(selector).first
|
|
|
|
|
if locator.count():
|
|
|
|
|
locator.fill(str(payload.get("password") or ""))
|
|
|
|
|
password_filled = True
|
|
|
|
|
break
|
|
|
|
|
except Exception:
|
|
|
|
|
pass
|
|
|
|
|
if not password_filled:
|
|
|
|
|
raise RuntimeError("GST password field could not be located. The portal login page may have changed.")
|
|
|
|
|
progress(percent=10, stage="Waiting for CAPTCHA / OTP", message="GST username/password filled. Complete CAPTCHA/OTP in the visible GST browser.")
|
|
|
|
|
deadline = time.time() + int(payload.get("login_timeout_seconds") or 900)
|
|
|
|
|
while time.time() < deadline:
|
|
|
|
|
if login_complete(page):
|
|
|
|
|
break
|
|
|
|
|
time.sleep(2)
|
|
|
|
|
else:
|
|
|
|
|
raise RuntimeError("GST login was not completed within 15 minutes.")
|
|
|
|
|
|
|
|
|
|
periods = [str(p) for p in payload.get("periods") or []]
|
|
|
|
|
return_types = [str(r) for r in payload.get("return_types") or []]
|
|
|
|
|
all_downloaded = []
|
|
|
|
|
for index, period in enumerate(periods, 1):
|
|
|
|
|
base_pct = 10 + int((index - 1) * 75 / max(1, len(periods)))
|
|
|
|
|
progress(percent=base_pct, period=period, period_index=index, period_total=len(periods), stage=f"Period {index}/{len(periods)}", message=f"Downloading selected GST returns for {period}.")
|
|
|
|
|
all_downloaded.extend(download_period(context, work_root, period, return_types, progress))
|
|
|
|
|
context.close()
|
|
|
|
|
|
|
|
|
|
if package_path.exists():
|
|
|
|
|
package_path.unlink()
|
|
|
|
|
with zipfile.ZipFile(package_path, "w", zipfile.ZIP_DEFLATED, compresslevel=6) as archive:
|
|
|
|
|
for item in sorted(work_root.rglob("*")):
|
|
|
|
|
if item.is_file():
|
|
|
|
|
archive.write(item, item.relative_to(work_root).as_posix())
|
|
|
|
|
|
|
|
|
|
progress(percent=90, stage="Transferring to Client Storage", message="Uploading GST return package to ERP for transfer to the configured Storage Agent.")
|
|
|
|
|
with package_path.open("rb") as handle:
|
|
|
|
|
response = requests.post(
|
|
|
|
|
str(payload.get("upload_url") or ""), data={"token": token},
|
|
|
|
|
files={"package": (package_path.name, handle, "application/zip")}, timeout=240,
|
|
|
|
|
)
|
|
|
|
|
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 storage transfer failed with HTTP {response.status_code}.")
|
|
|
|
|
progress(status="completed", percent=100, stage="Completed", message="GST returns downloaded and stored in the configured client local storage.", result={"stored": body.get("stored") or {}}, finished_at_utc=now())
|
|
|
|
|
except Exception as exc:
|
|
|
|
|
progress = read_json(path)
|
|
|
|
|
progress.update({"status": "failed", "percent": 100, "stage": "Failed", "message": str(exc), "error": str(exc), "finished_at_utc": now(), "updated_at_utc": now()})
|
|
|
|
|
write_json(path, progress)
|
|
|
|
|
finally:
|
|
|
|
|
shutil.rmtree(work_root, ignore_errors=True)
|
|
|
|
|
try:
|
|
|
|
|
package_path.unlink(missing_ok=True)
|
|
|
|
|
except Exception:
|
|
|
|
|
pass
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def start_job(token: str) -> dict:
|
|
|
|
|
cfg = read_config()
|
|
|
|
|
erp = str(cfg.get("erp_base_url") or "").rstrip("/")
|
|
|
|
|
if not token:
|
|
|
|
|
raise ValueError("GST operator token is required.")
|
|
|
|
|
response = requests.post(erp + "/tools/accounting/gst-reconciliation/operator/redeem", json={"token": 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 this GST download.")
|
|
|
|
|
payload = dict(body.get("payload") or {})
|
|
|
|
|
job_id = str(payload.get("jti") or "")
|
|
|
|
|
if not job_id:
|
|
|
|
|
raise RuntimeError("ERP did not return a GST job id.")
|
|
|
|
|
path = job_path(job_id)
|
|
|
|
|
current = read_json(path)
|
|
|
|
|
if current.get("status") in {"queued", "running"}:
|
|
|
|
|
return current
|
|
|
|
|
initial = {
|
|
|
|
|
"status": "queued", "percent": 1, "stage": "Queued", "message": "GST browser job queued.",
|
|
|
|
|
"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(path, initial)
|
|
|
|
|
thread = threading.Thread(target=browser_worker, args=(payload, token, path), daemon=True, name=f"gst-{job_id}")
|
|
|
|
|
THREADS[job_id] = thread
|
|
|
|
|
thread.start()
|
|
|
|
|
return initial
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class Handler(BaseHTTPRequestHandler):
|
|
|
|
|
def log_message(self, fmt, *args):
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
def cors(self):
|
|
|
|
|
cfg = read_config()
|
|
|
|
|
origin = self.headers.get("Origin") or ""
|
|
|
|
|
allowed = str(cfg.get("erp_base_url") or "").rstrip("/")
|
|
|
|
|
if origin.rstrip("/") == allowed:
|
|
|
|
|
self.send_header("Access-Control-Allow-Origin", origin)
|
|
|
|
|
self.send_header("Vary", "Origin")
|
|
|
|
|
if (self.headers.get("Access-Control-Request-Private-Network") or "").lower() == "true":
|
|
|
|
|
self.send_header("Access-Control-Allow-Private-Network", "true")
|
|
|
|
|
|
|
|
|
|
def reply(self, payload, status=200):
|
|
|
|
|
data = json.dumps(payload, ensure_ascii=False, default=str).encode("utf-8")
|
|
|
|
|
self.send_response(status)
|
|
|
|
|
self.send_header("Content-Type", "application/json; charset=utf-8")
|
|
|
|
|
self.send_header("Cache-Control", "no-store")
|
|
|
|
|
self.cors()
|
|
|
|
|
self.send_header("Content-Length", str(len(data)))
|
|
|
|
|
self.end_headers()
|
|
|
|
|
self.wfile.write(data)
|
|
|
|
|
|
|
|
|
|
def do_OPTIONS(self):
|
|
|
|
|
self.send_response(204)
|
|
|
|
|
self.cors()
|
|
|
|
|
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):
|
|
|
|
|
path = urlsplit(self.path)
|
|
|
|
|
if path.path == "/api/status":
|
|
|
|
|
return self.reply({"ok": True, "name": "ARRR GST Operator Agent", "version": VERSION, "port": PORT})
|
|
|
|
|
if path.path == "/api/gst/status":
|
|
|
|
|
job_id = str((parse_qs(path.query).get("job_id") or [""])[0])
|
|
|
|
|
return self.reply({"ok": True, "job": read_json(job_path(job_id))})
|
|
|
|
|
return self.reply({"ok": False, "error": "Not found"}, 404)
|
|
|
|
|
|
|
|
|
|
def do_POST(self):
|
|
|
|
|
try:
|
|
|
|
|
if urlsplit(self.path).path != "/api/gst/start":
|
|
|
|
|
return self.reply({"ok": False, "error": "Not found"}, 404)
|
|
|
|
|
length = int(self.headers.get("Content-Length") or 0)
|
|
|
|
|
body = json.loads((self.rfile.read(length) if length else b"{}").decode("utf-8"))
|
|
|
|
|
job = start_job(str(body.get("token") or ""))
|
|
|
|
|
return self.reply({"ok": True, "job": job})
|
|
|
|
|
except Exception as exc:
|
|
|
|
|
return self.reply({"ok": False, "error": str(exc)}, 400)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def main():
|
|
|
|
|
DATA.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
cfg = read_config()
|
|
|
|
|
port = int(cfg.get("port") or PORT)
|
|
|
|
|
server = ThreadingHTTPServer(("127.0.0.1", port), Handler)
|
|
|
|
|
server.serve_forever()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
if __name__ == "__main__":
|
|
|
|
|
main()
|