| 1 | #!/usr/bin/env python3 |
| 2 | # intake: append-only json line log. POST any json to /<name> with the shared |
| 3 | # key and it lands in <name>.jsonl; GET reads it back. no schema, no setup — |
| 4 | # a new name creates a new file on first write. |
| 5 | import hmac |
| 6 | import json |
| 7 | import os |
| 8 | import re |
| 9 | import threading |
| 10 | import time |
| 11 | |
| 12 | from flask import Flask, Response, jsonify, request |
| 13 | |
| 14 | os.umask(0o007) |
| 15 | |
| 16 | DATA_DIR = os.environ.get("DATA_DIR", "/data") |
| 17 | KEY = os.environ["INTAKE_KEY"] |
| 18 | # filenames are user-supplied path segments; anything outside this can escape |
| 19 | NAME_RE = re.compile(r"[a-z0-9][a-z0-9._-]{0,63}\Z") |
| 20 | |
| 21 | app = Flask(__name__) |
| 22 | app.config["MAX_CONTENT_LENGTH"] = 4 * 1024 * 1024 |
| 23 | write_lock = threading.Lock() |
| 24 | |
| 25 | |
| 26 | def authorized(): |
| 27 | given = request.headers.get("X-Intake-Key") or "" |
| 28 | if not given: |
| 29 | auth = request.headers.get("Authorization", "") |
| 30 | if auth.startswith("Bearer "): |
| 31 | given = auth[7:] |
| 32 | return hmac.compare_digest(given, KEY) |
| 33 | |
| 34 | |
| 35 | def path_for(name): |
| 36 | if not NAME_RE.match(name): |
| 37 | return None |
| 38 | return os.path.join(DATA_DIR, name + ".jsonl") |
| 39 | |
| 40 | |
| 41 | @app.before_request |
| 42 | def check_key(): |
| 43 | if request.path == "/health": |
| 44 | return |
| 45 | if not authorized(): |
| 46 | return jsonify(error="bad or missing key"), 401 |
| 47 | |
| 48 | |
| 49 | @app.get("/health") |
| 50 | def health(): |
| 51 | return "", 204 |
| 52 | |
| 53 | |
| 54 | @app.get("/") |
| 55 | def index(): |
| 56 | dbs = [] |
| 57 | for f in sorted(os.listdir(DATA_DIR)): |
| 58 | if f.endswith(".jsonl"): |
| 59 | st = os.stat(os.path.join(DATA_DIR, f)) |
| 60 | dbs.append({"name": f[:-6], "bytes": st.st_size, "modified": int(st.st_mtime)}) |
| 61 | return jsonify(databases=dbs) |
| 62 | |
| 63 | |
| 64 | @app.post("/<name>") |
| 65 | def append(name): |
| 66 | path = path_for(name) |
| 67 | if path is None: |
| 68 | return jsonify(error="name must match [a-z0-9][a-z0-9._-]{0,63}"), 400 |
| 69 | body = request.get_data() |
| 70 | try: |
| 71 | row = json.loads(body) |
| 72 | except ValueError: |
| 73 | # shortcuts and curl one-liners often send plain text; keep it rather |
| 74 | # than rejecting, so a mis-typed shortcut still logs something usable |
| 75 | row = {"text": body.decode("utf-8", "replace")} |
| 76 | rows = row if isinstance(row, list) else [row] |
| 77 | now = time.time() |
| 78 | lines = [] |
| 79 | for r in rows: |
| 80 | if not isinstance(r, dict): |
| 81 | r = {"value": r} |
| 82 | lines.append(json.dumps({"_at": now, **r}, ensure_ascii=False) + "\n") |
| 83 | with write_lock, open(path, "a", encoding="utf-8") as fh: |
| 84 | fh.write("".join(lines)) |
| 85 | return jsonify(ok=True, db=name, appended=len(lines)) |
| 86 | |
| 87 | |
| 88 | @app.get("/<name>") |
| 89 | def read(name): |
| 90 | path = path_for(name) |
| 91 | if path is None or not os.path.exists(path): |
| 92 | return jsonify(error="no such database"), 404 |
| 93 | with open(path, "rb") as fh: |
| 94 | return Response(fh.read(), mimetype="application/x-ndjson") |
| 95 | |
| 96 | |
| 97 | if __name__ == "__main__": |
| 98 | os.makedirs(DATA_DIR, exist_ok=True) |
| 99 | app.run(host="0.0.0.0", port=8000) |