|
|
|
@@ -0,0 +1,345 @@
|
|
|
|
|
from __future__ import annotations
|
|
|
|
|
|
|
|
|
|
import json
|
|
|
|
|
import re
|
|
|
|
|
import sqlite3
|
|
|
|
|
import threading
|
|
|
|
|
import time
|
|
|
|
|
from datetime import datetime, timezone
|
|
|
|
|
from pathlib import Path
|
|
|
|
|
from typing import Any
|
|
|
|
|
|
|
|
|
|
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}"
|
|
|
|
|
|
|
|
|
|
_LOCK = threading.Lock()
|
|
|
|
|
_THREADS: dict[str, threading.Thread] = {}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _now() -> str:
|
|
|
|
|
return datetime.now(timezone.utc).isoformat()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _write_json(path: Path, payload: Any) -> 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[str, Any]:
|
|
|
|
|
try:
|
|
|
|
|
data = json.loads(path.read_text(encoding="utf-8"))
|
|
|
|
|
return data if isinstance(data, dict) else {}
|
|
|
|
|
except Exception:
|
|
|
|
|
return {}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
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 _parse_date(value: str) -> str:
|
|
|
|
|
text = str(value or "").strip()
|
|
|
|
|
for fmt in ("%d-%m-%Y", "%d/%m/%Y", "%Y-%m-%d", "%d-%b-%Y", "%d/%m/%y"):
|
|
|
|
|
try:
|
|
|
|
|
return datetime.strptime(text, fmt).date().isoformat()
|
|
|
|
|
except Exception:
|
|
|
|
|
pass
|
|
|
|
|
return text
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _f(value: Any) -> float:
|
|
|
|
|
try:
|
|
|
|
|
return float(value or 0)
|
|
|
|
|
except Exception:
|
|
|
|
|
return 0.0
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _doc_key(value: str) -> str:
|
|
|
|
|
return "".join(ch for ch in str(value or "").upper() if ch.isalnum())
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _walk(obj: Any):
|
|
|
|
|
if isinstance(obj, dict):
|
|
|
|
|
yield obj
|
|
|
|
|
for v in obj.values():
|
|
|
|
|
yield from _walk(v)
|
|
|
|
|
elif isinstance(obj, list):
|
|
|
|
|
for v in obj:
|
|
|
|
|
yield from _walk(v)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _tax_from_invoice(inv: dict[str, Any]) -> dict[str, float]:
|
|
|
|
|
out = {"taxable": 0.0, "igst": 0.0, "cgst": 0.0, "sgst": 0.0, "cess": 0.0}
|
|
|
|
|
for item in inv.get("itms") or inv.get("items") or []:
|
|
|
|
|
d = item.get("itm_det") or item.get("item_det") or item
|
|
|
|
|
out["taxable"] += _f(d.get("txval"))
|
|
|
|
|
out["igst"] += _f(d.get("iamt"))
|
|
|
|
|
out["cgst"] += _f(d.get("camt"))
|
|
|
|
|
out["sgst"] += _f(d.get("samt"))
|
|
|
|
|
out["cess"] += _f(d.get("csamt"))
|
|
|
|
|
if not (out["taxable"] or out["igst"] or out["cgst"] or out["sgst"] or out["cess"]):
|
|
|
|
|
out["taxable"] = _f(inv.get("txval") or inv.get("taxable_value"))
|
|
|
|
|
out["igst"] = _f(inv.get("iamt") or inv.get("igst"))
|
|
|
|
|
out["cgst"] = _f(inv.get("camt") or inv.get("cgst"))
|
|
|
|
|
out["sgst"] = _f(inv.get("samt") or inv.get("sgst"))
|
|
|
|
|
out["cess"] = _f(inv.get("csamt") or inv.get("cess"))
|
|
|
|
|
return out
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def extract_invoice_rows(data: Any) -> list[dict[str, Any]]:
|
|
|
|
|
rows: list[dict[str, Any]] = []
|
|
|
|
|
seen: set[tuple[str, str, str]] = set()
|
|
|
|
|
for obj in _walk(data):
|
|
|
|
|
invno = str(obj.get("inum") or obj.get("inv_num") or obj.get("doc_no") or obj.get("invoiceNumber") or "").strip()
|
|
|
|
|
if not invno:
|
|
|
|
|
continue
|
|
|
|
|
gstin = str(obj.get("ctin") or obj.get("stin") or obj.get("supplier_gstin") or obj.get("recipientGstin") or "").strip().upper()
|
|
|
|
|
date = _parse_date(str(obj.get("idt") or obj.get("inv_date") or obj.get("doc_date") or obj.get("invoiceDate") or ""))
|
|
|
|
|
key = (gstin, _doc_key(invno), date)
|
|
|
|
|
if key in seen:
|
|
|
|
|
continue
|
|
|
|
|
seen.add(key)
|
|
|
|
|
tax = _tax_from_invoice(obj)
|
|
|
|
|
rows.append({
|
|
|
|
|
"gstin": gstin,
|
|
|
|
|
"invoice_no": invno,
|
|
|
|
|
"invoice_key": _doc_key(invno),
|
|
|
|
|
"invoice_date": date,
|
|
|
|
|
"invoice_value": _f(obj.get("val") or obj.get("invoice_value") or obj.get("invoiceValue")),
|
|
|
|
|
**tax,
|
|
|
|
|
})
|
|
|
|
|
return rows
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def extract_gstr3b_itc(data: Any) -> dict[str, float]:
|
|
|
|
|
totals = {"igst": 0.0, "cgst": 0.0, "sgst": 0.0, "cess": 0.0}
|
|
|
|
|
if not isinstance(data, dict):
|
|
|
|
|
return totals
|
|
|
|
|
itc = data.get("itc_elg") or {}
|
|
|
|
|
for item in itc.get("itc_avl") or []:
|
|
|
|
|
totals["igst"] += _f(item.get("iamt"))
|
|
|
|
|
totals["cgst"] += _f(item.get("camt"))
|
|
|
|
|
totals["sgst"] += _f(item.get("samt"))
|
|
|
|
|
totals["cess"] += _f(item.get("csamt"))
|
|
|
|
|
return totals
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def normalize_downloads(base_dir: Path, period: str) -> dict[str, Any]:
|
|
|
|
|
normalized = base_dir / "normalized"
|
|
|
|
|
normalized.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
result: dict[str, Any] = {}
|
|
|
|
|
for rtype in ("GSTR1", "GSTR2A", "GSTR2B", "GSTR3B"):
|
|
|
|
|
raw_path = base_dir / "raw" / f"{period}_{rtype}.json"
|
|
|
|
|
if not raw_path.exists():
|
|
|
|
|
continue
|
|
|
|
|
try:
|
|
|
|
|
data = json.loads(raw_path.read_text(encoding="utf-8", errors="ignore") or "{}")
|
|
|
|
|
except Exception:
|
|
|
|
|
data = {}
|
|
|
|
|
if rtype in {"GSTR1", "GSTR2A", "GSTR2B"}:
|
|
|
|
|
rows = extract_invoice_rows(data)
|
|
|
|
|
_write_json(normalized / f"{period}_{rtype}_invoices.json", {"period": period, "rows": rows})
|
|
|
|
|
result[rtype] = {"invoice_rows": len(rows)}
|
|
|
|
|
else:
|
|
|
|
|
itc = extract_gstr3b_itc(data)
|
|
|
|
|
_write_json(normalized / f"{period}_{rtype}_itc.json", {"period": period, "itc": itc})
|
|
|
|
|
result[rtype] = {"itc": itc}
|
|
|
|
|
return result
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _book_rows(mirror_db: Path, date_from: str, date_to: str, voucher_type: str) -> list[dict[str, Any]]:
|
|
|
|
|
if not mirror_db.exists():
|
|
|
|
|
return []
|
|
|
|
|
con = sqlite3.connect(str(mirror_db))
|
|
|
|
|
con.row_factory = sqlite3.Row
|
|
|
|
|
try:
|
|
|
|
|
sql = """
|
|
|
|
|
SELECT voucher_date, voucher_type, voucher_number, party_ledger, reference, ABS(COALESCE(voucher_amount,0)) amount
|
|
|
|
|
FROM voucher
|
|
|
|
|
WHERE voucher_date >= ? AND voucher_date <= ? AND UPPER(COALESCE(voucher_type,'')) LIKE ?
|
|
|
|
|
ORDER BY voucher_date, voucher_number
|
|
|
|
|
"""
|
|
|
|
|
rows = con.execute(sql, (date_from, date_to, f"%{voucher_type.upper()}%" )).fetchall()
|
|
|
|
|
return [{"date": r["voucher_date"], "voucher_type": r["voucher_type"], "invoice_no": r["voucher_number"] or r["reference"] or "", "invoice_key": _doc_key(r["voucher_number"] or r["reference"] or ""), "party": r["party_ledger"] or "", "amount": _f(r["amount"])} for r in rows]
|
|
|
|
|
finally:
|
|
|
|
|
con.close()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _match_books_to_portal(books: list[dict[str, Any]], portal: list[dict[str, Any]]) -> dict[str, Any]:
|
|
|
|
|
pmap: dict[str, list[dict[str, Any]]] = {}
|
|
|
|
|
for row in portal:
|
|
|
|
|
if row.get("invoice_key"):
|
|
|
|
|
pmap.setdefault(row["invoice_key"], []).append(row)
|
|
|
|
|
matched=[]; missing_portal=[]; used=set()
|
|
|
|
|
for b in books:
|
|
|
|
|
candidates=pmap.get(b.get("invoice_key") or "", [])
|
|
|
|
|
best=None
|
|
|
|
|
for p in candidates:
|
|
|
|
|
pid=id(p)
|
|
|
|
|
if pid in used: continue
|
|
|
|
|
best=p; break
|
|
|
|
|
if best is None:
|
|
|
|
|
missing_portal.append(b); continue
|
|
|
|
|
used.add(id(best))
|
|
|
|
|
pv=_f(best.get("invoice_value")); bv=_f(b.get("amount")); diff=round(bv-pv,2)
|
|
|
|
|
matched.append({"books":b,"portal":best,"difference":diff,"status":"Matched" if abs(diff)<=1 else "Value mismatch"})
|
|
|
|
|
missing_books=[p for p in portal if id(p) not in used]
|
|
|
|
|
return {"matched":matched,"missing_in_portal":missing_portal,"missing_in_books":missing_books,"counts":{"books":len(books),"portal":len(portal),"matched":len(matched),"missing_in_portal":len(missing_portal),"missing_in_books":len(missing_books),"value_mismatch":sum(1 for r in matched if r["status"]!="Matched")}}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def analyze_gst_storage(base_dir: Path, mirror_db: Path, period: str, date_from: str, date_to: str) -> dict[str, Any]:
|
|
|
|
|
norm=base_dir/"normalized"
|
|
|
|
|
def rows(name):
|
|
|
|
|
p=norm/f"{period}_{name}_invoices.json"
|
|
|
|
|
d=_read_json(p) if p.exists() else {}
|
|
|
|
|
return d.get("rows") or []
|
|
|
|
|
g1=rows("GSTR1"); g2b=rows("GSTR2B")
|
|
|
|
|
sales=_book_rows(mirror_db,date_from,date_to,"SALES")
|
|
|
|
|
purchases=_book_rows(mirror_db,date_from,date_to,"PURCHASE")
|
|
|
|
|
sales_rec=_match_books_to_portal(sales,g1)
|
|
|
|
|
purchase_rec=_match_books_to_portal(purchases,g2b)
|
|
|
|
|
g2b_itc={k:round(sum(_f(r.get(k)) for r in g2b),2) for k in ("igst","cgst","sgst","cess")}
|
|
|
|
|
g3p=norm/f"{period}_GSTR3B_itc.json"; g3=_read_json(g3p).get("itc") if g3p.exists() else {}
|
|
|
|
|
g3={k:round(_f((g3 or {}).get(k)),2) for k in ("igst","cgst","sgst","cess")}
|
|
|
|
|
itc={k:{"gstr2b":g2b_itc[k],"gstr3b":g3[k],"difference":round(g3[k]-g2b_itc[k],2)} for k in g2b_itc}
|
|
|
|
|
result={"period":period,"generated_at_utc":_now(),"sales_reconciliation":sales_rec,"purchase_reconciliation":purchase_rec,"itc_reconciliation":itc}
|
|
|
|
|
reports=base_dir/"reports"; reports.mkdir(parents=True,exist_ok=True); _write_json(reports/f"{period}_GST_reconciliation.json",result)
|
|
|
|
|
return result
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _extract_file_num(text: str) -> str:
|
|
|
|
|
try:
|
|
|
|
|
data=json.loads(text or "{}")
|
|
|
|
|
except Exception:
|
|
|
|
|
return "1"
|
|
|
|
|
vals=[]
|
|
|
|
|
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,""):
|
|
|
|
|
vals.append(str(obj.get(key)))
|
|
|
|
|
return vals[0] if vals else "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","json","payload","content","fileContent","response"):
|
|
|
|
|
value=data.get(key)
|
|
|
|
|
if isinstance(value,(dict,list)):
|
|
|
|
|
return json.dumps(value,ensure_ascii=False)
|
|
|
|
|
if isinstance(value,str) and value.lstrip().startswith(("{","[")):
|
|
|
|
|
return value
|
|
|
|
|
return text or ""
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _browser_download(payload: dict[str, Any], progress_path: Path, storage_root: Path) -> None:
|
|
|
|
|
base_dir = storage_root / Path(str(payload["gst_relative_dir"])) / str(payload["period"])
|
|
|
|
|
raw_dir = base_dir / "raw"; raw_dir.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)
|
|
|
|
|
try:
|
|
|
|
|
progress(status="running",stage="Opening GST Login",message="Opening GST portal. Complete captcha/OTP in the browser.")
|
|
|
|
|
from playwright.sync_api import sync_playwright
|
|
|
|
|
with sync_playwright() as pw:
|
|
|
|
|
browser=None; errors=[]
|
|
|
|
|
profile=storage_root/".erp_gst_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 through Playwright. "+" | ".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. Enter captcha and OTP in the GST browser. The agent will continue 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.")
|
|
|
|
|
downloaded=[]
|
|
|
|
|
period=str(payload["period"])
|
|
|
|
|
fetch_js="async (u)=>{const r=await fetch(u,{credentials:'include'}); return await r.text();}"
|
|
|
|
|
|
|
|
|
|
progress(stage="Downloading GSTR1",message=f"Generating and reading GSTR-1 for {period} from GST portal.")
|
|
|
|
|
generate_text=page.evaluate(fetch_js,GSTR1_URL.format(period=period))
|
|
|
|
|
(raw_dir/f"{period}_GSTR1_GENERATE.json").write_text(generate_text or "",encoding="utf-8",errors="ignore")
|
|
|
|
|
file_num=_extract_file_num(generate_text)
|
|
|
|
|
g1_text=page.evaluate(fetch_js,GSTR1_DOWNLOAD_URL.format(period=period,file_num=file_num))
|
|
|
|
|
g1_text=_extract_payload(g1_text)
|
|
|
|
|
g1_path=raw_dir/f"{period}_GSTR1.json"; g1_path.write_text(g1_text or "",encoding="utf-8",errors="ignore")
|
|
|
|
|
downloaded.append({"return_type":"GSTR1","path":str(g1_path),"bytes":g1_path.stat().st_size})
|
|
|
|
|
|
|
|
|
|
progress(stage="Downloading GSTR2B",message=f"Reading GSTR-2B for {period} from GST portal.")
|
|
|
|
|
g2_text=page.evaluate(fetch_js,GSTR2B_URL.format(period=period))
|
|
|
|
|
g2_path=raw_dir/f"{period}_GSTR2B.json"; g2_path.write_text(g2_text or "",encoding="utf-8",errors="ignore")
|
|
|
|
|
downloaded.append({"return_type":"GSTR2B","path":str(g2_path),"bytes":g2_path.stat().st_size})
|
|
|
|
|
|
|
|
|
|
if payload.get("include_2a"):
|
|
|
|
|
progress(stage="Requesting GSTR2A",message=f"Requesting GSTR-2A offline data for {period}.")
|
|
|
|
|
g2a_text=page.evaluate(fetch_js,GSTR2A_URL.format(period=period))
|
|
|
|
|
g2a_path=raw_dir/f"{period}_GSTR2A.json"; g2a_path.write_text(g2a_text or "",encoding="utf-8",errors="ignore")
|
|
|
|
|
downloaded.append({"return_type":"GSTR2A_GENERATE","path":str(g2a_path),"bytes":g2a_path.stat().st_size})
|
|
|
|
|
|
|
|
|
|
progress(stage="Downloading GSTR3B",message=f"Reading GSTR-3B summary and tax payable for {period}.")
|
|
|
|
|
g3_summary=page.evaluate(fetch_js,GSTR3B_SUMMARY_URL.format(period=period))
|
|
|
|
|
g3_payable=page.evaluate(fetch_js,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 or "",encoding="utf-8",errors="ignore")
|
|
|
|
|
(raw_dir/f"{period}_GSTR3B_TAXPAYBLE.json").write_text(g3_payable or "",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),"bytes":g3_path.stat().st_size})
|
|
|
|
|
|
|
|
|
|
normalized=normalize_downloads(base_dir,period)
|
|
|
|
|
manifest={"gstin":payload.get("gstin"),"period":payload["period"],"financial_year":payload.get("financial_year"),"downloaded_at_utc":_now(),"downloaded":downloaded,"normalized":normalized}
|
|
|
|
|
_write_json(base_dir/"download_manifest.json",manifest)
|
|
|
|
|
browser.close()
|
|
|
|
|
progress(status="completed",stage="Completed",message="GST return data downloaded and stored in the client GST directory.",result=manifest,finished_at_utc=_now())
|
|
|
|
|
except Exception as exc:
|
|
|
|
|
progress(status="failed",stage="Failed",message=str(exc),error=str(exc),finished_at_utc=_now())
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def await_login_probe(page) -> bool:
|
|
|
|
|
try:
|
|
|
|
|
return bool(page.evaluate("async()=>{try{const r=await fetch('https://return.gst.gov.in/returns/auth/api/returns/profile',{credentials:'include'}); return r.status!==401 && r.status!==403;}catch(e){return false;}}"))
|
|
|
|
|
except Exception:
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def start_download(payload: dict[str, Any], storage_root: Path, progress_root: Path) -> dict[str, Any]:
|
|
|
|
|
key=_safe(f"{payload.get('client_id')}_{payload.get('gstin')}_{payload.get('period')}")
|
|
|
|
|
progress_path=progress_root/f"gst_download_{key}.json"
|
|
|
|
|
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":"GST return download queued.","client_id":payload.get("client_id"),"gstin":payload.get("gstin"),"period":payload.get("period"),"started_at_utc":_now(),"updated_at_utc":_now()}
|
|
|
|
|
_write_json(progress_path,initial)
|
|
|
|
|
thread=threading.Thread(target=_browser_download,args=(dict(payload),progress_path,storage_root),daemon=True,name=f"gst-download-{key}")
|
|
|
|
|
with _LOCK: _THREADS[key]=thread
|
|
|
|
|
thread.start()
|
|
|
|
|
return {"started":True,"job":initial}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
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")
|