import json, os, sys, ssl, urllib.parse, urllib.request, urllib.error def _cfg(): return json.loads(os.environ.get("INTEGRATION_SECRETS", "{}")) def _inputs(): return json.loads(os.environ.get("INTEGRATION_INPUTS", "{}")) def _ctx(cfg): if cfg.get("insecure"): c = ssl.create_default_context() c.check_hostname = False c.verify_mode = ssl.CERT_NONE return c return None def request(method, path, cfg, body=None, params=None): url = str(cfg.get("url", "")).rstrip("/") + path if params: clean = {k: v for k, v in params.items() if v not in (None, "")} if clean: url += "?" + urllib.parse.urlencode(clean) data = json.dumps(body).encode("utf-8") if body is not None else None headers = {"Authorization": "ApiKey " + str(cfg.get("api_key", "")), "Accept": "application/json"} if data is not None: headers["Content-Type"] = "application/json" req = urllib.request.Request(url, data=data, headers=headers, method=method) with urllib.request.urlopen(req, timeout=90, context=_ctx(cfg)) as r: raw = r.read() return json.loads(raw) if raw else {} def _run(fn): try: print(json.dumps(fn(_cfg(), _inputs()))) except urllib.error.HTTPError as e: print(json.dumps({"error": "HTTP " + str(e.code), "detail": e.read().decode("utf-8", "replace")})) sys.exit(1) except Exception as e: print(json.dumps({"error": str(e)})) sys.exit(1) q = lambda v: urllib.parse.quote(str(v), safe="") def _parse_json(s, field): try: return json.loads(s) except Exception: raise Exception(field + " must be a valid JSON object") def main(cfg, inputs): index = inputs.get("index") if not index: raise Exception("index is required") document_json = inputs.get("document_json") if not document_json: raise Exception("document_json is required") doc = _parse_json(document_json, "document_json") doc_id = inputs.get("doc_id") if doc_id: return request("PUT", "/" + q(index) + "/_doc/" + q(doc_id), cfg, body=doc) return request("POST", "/" + q(index) + "/_doc", cfg, body=doc) _run(main)