| 1 | #!/usr/bin/env python3 |
| 2 | import datetime |
| 3 | import json |
| 4 | import os |
| 5 | from pathlib import Path |
| 6 | import re |
| 7 | import subprocess |
| 8 | import time |
| 9 | from urllib.request import Request, urlopen |
| 10 | |
| 11 | |
| 12 | STATE = Path("/var/lib/studio/log-shipper.json") |
| 13 | LOGS = Path("/var/lib/nomad/alloc") |
| 14 | CRI = re.compile(r"^(\S+) (stdout|stderr) [FP] ?(.*)$") |
| 15 | |
| 16 | |
| 17 | def nomad(path): |
| 18 | token = Path("/var/lib/studio/nomad.token").read_text().strip() |
| 19 | with urlopen(Request("http://127.0.0.1:4646" + path, headers={"X-Nomad-Token": token}), timeout=10) as response: |
| 20 | return json.load(response) |
| 21 | |
| 22 | |
| 23 | def send(base, rows, fields): |
| 24 | if not rows: |
| 25 | return |
| 26 | data = b"\n".join(json.dumps(row, ensure_ascii=False).encode() for row in rows) + b"\n" |
| 27 | request = Request(base + "/insert/jsonline", data=data, headers={ |
| 28 | "Content-Type": "application/stream+json", "VL-Stream-Fields": fields, |
| 29 | }) |
| 30 | with urlopen(request, timeout=30) as response: |
| 31 | response.read() |
| 32 | |
| 33 | |
| 34 | def once(state): |
| 35 | services = nomad("/v1/service/victoria-logs") |
| 36 | if not services: |
| 37 | return |
| 38 | target = services[0] |
| 39 | base = f"http://{target['Address']}:{target['Port']}" |
| 40 | jobs = {item["ID"]: item["JobID"] for item in nomad("/v1/allocations")} |
| 41 | offsets = state.setdefault("files", {}) |
| 42 | rows = [] |
| 43 | pending = {} |
| 44 | for file in sorted(LOGS.glob("*/alloc/logs/*")): |
| 45 | match = re.fullmatch(r"(.+)\.(stdout|stderr)\.\d+", file.name) |
| 46 | job = jobs.get(file.parent.parent.parent.name) |
| 47 | if not match or not job: |
| 48 | continue |
| 49 | info = file.stat() |
| 50 | old = offsets.get(str(file), {}) |
| 51 | offset = old.get("offset", 0) if old.get("inode") == info.st_ino and old.get("offset", 0) <= info.st_size else 0 |
| 52 | with file.open("rb") as stream: |
| 53 | stream.seek(offset) |
| 54 | chunk = stream.read(262144) |
| 55 | if not chunk: |
| 56 | continue |
| 57 | if not chunk.endswith(b"\n"): |
| 58 | chunk = chunk[:chunk.rfind(b"\n") + 1] |
| 59 | if not chunk: |
| 60 | continue |
| 61 | pending[str(file)] = {"inode": info.st_ino, "offset": offset + len(chunk)} |
| 62 | for raw in chunk.decode("utf8", "replace").splitlines(): |
| 63 | parsed = CRI.match(raw) |
| 64 | if not parsed: |
| 65 | continue |
| 66 | stamp, stream, message = parsed.groups() |
| 67 | level = "error" if re.search(r"\b(ERROR|ERR|FATAL|CRITICAL|PANIC)\b|level=(error|fatal)", message, re.I) else "warn" if re.search(r"\b(WARN|WARNING|WRN)\b|level=warn", message, re.I) else "info" |
| 68 | rows.append({"_time": stamp, "_msg": message, "source": "nomad", "job": job, |
| 69 | "task": match.group(1), "stream": stream, "level": level}) |
| 70 | send(base, rows, "source,job,task,stream") |
| 71 | offsets.update(pending) |
| 72 | |
| 73 | command = ["journalctl", "--no-pager", "--output=json"] |
| 74 | command += ["--after-cursor=" + state["cursor"]] if state.get("cursor") else ["--since=-1 hour"] |
| 75 | system = [] |
| 76 | cursor = None |
| 77 | process = subprocess.Popen(command, stdout=subprocess.PIPE, text=True) |
| 78 | try: |
| 79 | for index, line in enumerate(process.stdout): |
| 80 | if index == 5000: |
| 81 | process.terminate() |
| 82 | break |
| 83 | entry = json.loads(line) |
| 84 | cursor = entry.get("__CURSOR", cursor) |
| 85 | message = entry.get("MESSAGE") |
| 86 | if not isinstance(message, str): |
| 87 | continue |
| 88 | stamp = int(entry["__REALTIME_TIMESTAMP"]) / 1_000_000 |
| 89 | system.append({"_time": datetime.datetime.fromtimestamp(stamp, datetime.timezone.utc).isoformat(), |
| 90 | "_msg": message, "source": "journald", "unit": entry.get("_SYSTEMD_UNIT", "system"), |
| 91 | "identifier": entry.get("SYSLOG_IDENTIFIER", ""), "priority": entry.get("PRIORITY", "")}) |
| 92 | finally: |
| 93 | process.stdout.close() |
| 94 | process.wait() |
| 95 | if process.returncode not in (0, -15): |
| 96 | raise RuntimeError(f"journalctl exited {process.returncode}") |
| 97 | send(base, system, "source,unit") |
| 98 | if cursor: |
| 99 | state["cursor"] = cursor |
| 100 | pending_state = STATE.with_suffix(".pending") |
| 101 | pending_state.write_text(json.dumps(state)) |
| 102 | os.replace(pending_state, STATE) |
| 103 | |
| 104 | |
| 105 | def main(): |
| 106 | STATE.parent.mkdir(parents=True, exist_ok=True) |
| 107 | while True: |
| 108 | try: |
| 109 | state = json.loads(STATE.read_text()) if STATE.exists() else {} |
| 110 | once(state) |
| 111 | except Exception as error: |
| 112 | print(f"log shipping paused: {error}", flush=True) |
| 113 | time.sleep(5) |
| 114 | |
| 115 | |
| 116 | if __name__ == "__main__": |
| 117 | main() |