#!/usr/bin/env python3
"""ferry-log-shipper — tail the ferry litellm proxy log into VictoriaLogs.

The push half of the llm-ferry observability stack (observ/CONTRACT.md):
the exporter turns the proxy log into `ferry_` METRICS for VictoriaMetrics,
this daemon ships the same log's LINES into VictoriaLogs (:9428) so Grafana
can show per-model, per-client, per-status logs next to those metrics.

  proxy log ──(this daemon, push)──▶ VictoriaLogs :9428 ──▶ Grafana :3001

It reads only the local access log — no model calls, no token spend — and it
is standard library only, so it runs under any python3 on macOS and Linux
(no venv, no pip, no secrets, no env).

What it emits — one JSON object per log line, POSTed as JSON-lines to
`{vlogs}/insert/jsonline?_stream_fields=source,level`:

    {"_time": "2026-08-22T16:02:03.123456Z",     # ingest time, UTC RFC3339
     "_msg":  "INFO: 1.2.3.4:5001 - \\"POST /v1/chat/completions ...\\" 200 OK",
     "source": "proxy",                          # stream field
     "level":  "info" | "warn" | "error",        # stream field
     "status": "200",                            # HTTP code, "" if not an access line
     "model":  "openrouter/openai/gpt-5.6-luna", # backend model, "" if absent
     "requested_model": "heavy",                 # model group asked for, "" if absent
     "attribution": "request" | "mention",       # did this line SERVE, or merely name?
     "client_ip": "1.2.3.4"}                     # "" if not an access line

`model` answers "which model does this line name", which is not the same
question as "which model served a request" — litellm names models at startup
too. `attribution` separates them, so a per-lane panel can filter
`attribution:request` instead of trusting every mention.

`_time` is INGEST time, not the line's own clock: uvicorn access lines carry
no timestamp at all and litellm's `15:51:27 - LiteLLM:WARNING` prefix carries
no date, so there is nothing dateable to trust. Successive records are nudged
1µs apart so Grafana's ordering stays stable inside a batch.

`_msg` is the line with ANSI colour escapes (litellm colours its output) and
C0 control bytes removed — otherwise every litellm warning renders as
`[92m...[0m` garbage in the Grafana logs panel. Nothing else is altered.

Discovery, and the access-line regex, mirror the sibling `ferry-dash` /
`ferry-metrics-exporter` (`${TMPDIR:-/tmp}/ferry-logs/cloud-proxy-<port>.log`,
plus the macOS `/var/folders/*/*/T` TMPDIR) so all three read the same file.

Robustness (this is a nohup daemon — it must never die):
  * truncation  — file shrinks below our offset  -> reread from byte 0
  * rotation    — (st_dev, st_ino) changes        -> reread from byte 0
  * in-place rewrite — `ferry up` reopens the log with `>`, so the inode is
                  unchanged and the fresh content can already be LONGER than
                  our offset: the first 256 bytes are fingerprinted, and a
                  changed head -> reread from byte 0
  * deletion    — file vanishes                   -> rediscover each poll,
                  and start the replacement from byte 0 (unless --log pinned)
  * partial line— a line still being written      -> buffered, shipped once
                  its newline lands
  * VictoriaLogs down -> the batch is kept and retried with exponential
                  backoff (1s→30s); the pending buffer is capped and drops
                  OLDEST lines first, so memory can't grow without bound
  * a bad line / a bad response never propagates: the poll loop catches
    everything, reports it on stderr, and keeps going.

Usage:
  ferry-log-shipper                                  # tail from END, ship to :9428
  ferry-log-shipper --vlogs http://127.0.0.1:9428
  ferry-log-shipper --log /tmp/ferry-logs/cloud-proxy-8090.log
  ferry-log-shipper --from-start                     # also ship existing history
  ferry-log-shipper --once --dry-run                 # parse one poll, print, exit
"""
import argparse
import datetime
import json
import os
import re
import sys
import time
import urllib.error
import urllib.request

DEFAULT_VLOGS = "http://127.0.0.1:9428"
DEFAULT_PORT = "8090"                # the litellm proxy port ferry runs (CONTRACT.md)

POLL_INTERVAL = 1.0                  # seconds between log reads
FLUSH_INTERVAL = 1.0                 # seconds a partial batch may wait
BATCH_LINES = 100                    # ...or this many lines, whichever comes first
MAX_PENDING = 20000                  # cap on un-shipped records (drops oldest)
BACKOFF_START = 1.0
BACKOFF_MAX = 30.0
MAX_READ = 4 * 1024 * 1024           # bytes per poll (catch up over several polls)
MAX_PARTIAL = 256 * 1024             # a "line" this long is shipped without its newline
SIG_LEN = 256                        # head bytes fingerprinted to spot an in-place rewrite
HTTP_TIMEOUT = 10
CONTENT_TYPE = "application/stream+json"


# ── line scrubbing ──────────────────────────────────────────────────────────
# litellm colours its output; uvicorn does not. Strip CSI/OSC escapes and the
# leftover C0 controls so the Grafana logs panel shows readable text.
ANSI = re.compile(r"\x1b\[[0-9;:?]*[ -/]*[@-~]|\x1b\][^\x07\x1b]*(?:\x07|\x1b\\)|\x1b[@-Z\\-_]")
CTRL = re.compile(r"[\x00-\x08\x0b\x0c\x0e-\x1f\x7f]")


def scrub(line):
    """The line as a human reads it: no colour escapes, no control bytes."""
    return CTRL.sub("", ANSI.sub("", line)).rstrip("\r\n")


# ── parsing ─────────────────────────────────────────────────────────────────
# uvicorn access line, verbatim from observ/CONTRACT.md / ferry-dash Activity:
#   `INFO: 1.2.3.4:5678 - "POST /v1/chat/completions HTTP/1.1" 200 OK`
ACCESS = re.compile(r'(\d+\.\d+\.\d+\.\d+):\d+ - "(\w+) (\S+) [^"]*" (\d+)')

# error/warn keyword hints (case-sensitive on purpose: "ERROR"/"Exception" are
# log levels and class names, while a lowercase "error" is ordinary prose and a
# field name in every JSON body litellm echoes back).
ERROR_HINT = re.compile(r"ERROR|CRITICAL|Exception|Traceback|RESOURCE_EXHAUSTED"
                        r"|\b(?:429|403|500)\b")
# \b keeps the numeric codes off ephemeral ports (":63500" has no boundary
# before "500") while still matching `... HTTP/1.1" 500` and `429 Too Many`.
WARN_HINT = re.compile(r"WARN|[Ff]allback|cool")

_IDCHARS = r"[\w./:@-]+"
# `litellm_model_name=x` / `model=x` / `model: x`. The lookbehind is what keeps
# `register_model: model=openrouter/openai/gpt-5.6-luna` from matching on the
# "model:" inside "register_model:" (which would capture the literal word "model").
MODEL_PATTERNS = [
    re.compile(r"litellm_model_name['\"]?\s*[=:]\s*['\"]?(" + _IDCHARS + r")"),
    re.compile(r"(?<![\w.-])model['\"]?\s*[=:]\s*['\"]?(" + _IDCHARS + r")"),
    re.compile(r"(?<![\w.-])model\s+(" + _IDCHARS + r")"),
]
REQUESTED_PATTERNS = [
    re.compile(r"requested_model['\"]?\s*[=:]\s*['\"]?(" + _IDCHARS + r")"),
    re.compile(r"model_group['\"]?\s*[=:]\s*['\"]?(" + _IDCHARS + r")"),
    re.compile(r"model group[:=]?\s*['\"]?(" + _IDCHARS + r")"),
]
# Known ferry backend ids (CONTRACT.md route topology), optionally provider-
# prefixed. CURRENT (2026-09-04): chatgpt/responses/gpt-5.6-sol (heavy driver),
# openrouter/google/gemini-3.8-flash (flash / super-flash), openrouter/openai/
# gpt-5.6-luna (the flash-luna / super-flash-luna hops), and the two local GPU
# ids mlx-community/Qwen3.8-27B-nvfp4 / mlx-community/NVIDIA-Nemotron-3-Nano-
# 30B-A3B-NVFP4. The glm-/deepseek/k3/kimi alternatives below are RETIRED
# vendors (Z.ai GLM, DeepSeek, Kimi K3) — kept only so a log line from an
# older ferry build (or the test fixtures that model one) still scans; no
# current deployment can emit them.
KNOWN_MODEL = re.compile(
    r"(?<![\w.-])((?:[\w-]+/)*"
    # `k3` is a WHOLE id on its own (the retired 1M-context Kimi deployment) as
    # well as a prefix (`k3-256k`), so the suffix is optional and a trailing
    # boundary keeps the bare form from biting a chunk out of some longer word.
    r"(?:gemini-[\w.-]+|gpt-[\w.-]+|[Qq]wen[\w.-]*|NVIDIA-Nemotron[\w.-]*"
    r"|glm-[\w.-]+|deepseek[\w.-]*"
    r"|k3(?:-[\w.-]+)?(?![\w-])"
    r"|kimi[\w.-]*|claude-[\w.-]+))")
# Ferry's model GROUPS (what a client asks for) as opposed to a backend id.
# CURRENT lane names (2026-09-04): `heavy` (the driver), `flash` and
# `super-flash` (each with its own `*-luna` fallback hop), and the two GPU
# lanes `local-orch` / `local-sub`. `orch` / `orchestrator` (+ `orch-*` hops)
# are the pre-rename driver names, kept so a log line from an older config, or
# a session still pinned mid-run, still resolves to a group.
KNOWN_GROUP = re.compile(
    r"(?<![\w.-])((?:local-)?orch(?:estrator)?(?:-[\w.-]+)?"
    r"|local-sub|heavy|super-flash(?:-luna)?|flash(?:-luna)?)")

# A line that NAMES a model and a line that reports a model SERVING A REQUEST
# are different facts, and `model` alone cannot tell them apart. litellm names
# models at startup (`register_model:`, the cost-map warnings, the catalogue
# dump) and those captures are deliberate — see the tests — but the Grafana
# `$model` filter reads the field as routing truth. Observed 2026-08-29: the
# newest record carrying a populated `model` was a `register_model:` cost-map
# warning, and its neighbour was a bare indented model name from a catalogue
# dump. Neither describes a request. `attribution` is the field that says which
# kind of line this was, so a per-lane panel can ask for the request ones.
#
# Evidence a line concerns a real call: an access line's HTTP status (handled by
# the caller), a versioned API path, a litellm call id, or one of the routing
# events that only happen mid-request.
REQUEST_SCOPED = re.compile(
    r"/v\d+/(?:chat/completions|completions|messages|responses|embeddings)"
    r"|litellm_call_id|x-litellm-call-id"
    r"|[Ff]allback|cool(?:ed|ing)?[ _]?down|RESOURCE_EXHAUSTED|[Rr]ate ?limit"
    r"|selected .* for this call")

# Startup and catalogue lines: they name models without serving anything.
NON_REQUEST = re.compile(
    r"register_model|cost map|cache cost fields|model_cost"
    r"|Proxy initialized|Loaded config|Set models:")


def attribution_for(text, status):
    """"request" | "mention" — whether this line describes a served call.

    Deliberately NOT a filter on `model`: the loose id scan is load-bearing and
    tested, so the model fields keep their existing meaning. This only labels
    which kind of line produced them.
    """
    if status:
        return "request"
    if NON_REQUEST.search(text):
        return "mention"
    return "request" if REQUEST_SCOPED.search(text) else "mention"


# Captures that are grammar, not identity ("model not in built-in cost map").
_STOPWORDS = {"model", "name", "group", "none", "null", "not", "in", "the",
              "id", "info", "is", "for", "to", "and", "a", "an", "with",
              "true", "false", "list", "params", "unknown"}
_HEX = re.compile(r"^[0-9a-f]{32,}$")


def _clean_id(value):
    """Normalise a captured model id, or "" if the capture is not an id.

    litellm's `register_model: model=<64 hex>` lines carry the ANONYMISED
    cost-map key for a keyed deployment, not a routable model id — shipping
    those as `model` would bury the real ids in the Grafana model filter, so
    they are treated as absent.
    """
    if not value:
        return ""
    v = value.strip().strip("'\"").rstrip(".,;:)]}")
    if len(v) < 2 or v.lower() in _STOPWORDS or _HEX.match(v.lower()):
        return ""
    return v


def _first_match(patterns, text):
    for pat in patterns:
        m = pat.search(text)
        if m:
            got = _clean_id(m.group(1))
            if got:
                return got
    return ""


def find_model(text):
    """The BACKEND model the line is about ("" when the line names none)."""
    got = _first_match(MODEL_PATTERNS, text)
    if got:
        return got
    m = KNOWN_MODEL.search(text)
    return _clean_id(m.group(1)) if m else ""


def find_requested_model(text):
    """The model GROUP the client asked for ("" when the line names none)."""
    got = _first_match(REQUESTED_PATTERNS, text)
    if got:
        return got
    m = KNOWN_GROUP.search(text)
    return _clean_id(m.group(1)) if m else ""


def level_for(text, status):
    """info | warn | error. An access line's HTTP status OVERRIDES keywords."""
    if status:
        try:
            code = int(status)
        except ValueError:
            code = 0
        if code >= 500:
            return "error"
        if code >= 400:
            return "warn"
        return "info"
    if ERROR_HINT.search(text):
        return "error"
    if WARN_HINT.search(text):
        return "warn"
    return "info"


def parse_line(line):
    """One raw log line -> the VictoriaLogs record (minus `_time`).

    Best effort and total: every field defaults to "" and nothing here raises.
    """
    text = scrub(line)
    client_ip = status = ""
    m = ACCESS.search(text)
    if m:
        client_ip, _method, _path, status = m.groups()
    return {
        "_msg": text,
        "source": "proxy",
        "level": level_for(text, status),
        "status": status,
        "model": find_model(text),
        "requested_model": find_requested_model(text),
        "attribution": attribution_for(text, status),
        "client_ip": client_ip,
    }


# ── log discovery (the ferry-dash find_log ladder, glob-free) ───────────────
def _var_folders_candidates(port):
    """macOS: ferry's $TMPDIR is /var/folders/<xx>/<hash>/T — walk two levels."""
    out = []
    base = "/var/folders"
    try:
        level1 = os.listdir(base)
    except OSError:
        return out
    for a in level1:
        try:
            level2 = os.listdir(os.path.join(base, a))
        except OSError:
            continue
        for b in level2:
            p = os.path.join(base, a, b, "T", "ferry-logs",
                             "cloud-proxy-%s.log" % port)
            if os.path.exists(p):
                out.append(p)
    return out


def _dir_candidates(directory, port):
    """Every cloud-proxy-*.log in `directory`, exact-port matches flagged first."""
    exact, loose = [], []
    try:
        names = os.listdir(directory)
    except OSError:
        return exact, loose
    for n in names:
        if not (n.startswith("cloud-proxy-") and n.endswith(".log")):
            continue
        p = os.path.join(directory, n)
        if not os.path.isfile(p):
            continue
        (exact if n == "cloud-proxy-%s.log" % port else loose).append(p)
    return exact, loose


def find_log(port=DEFAULT_PORT):
    """The ferry proxy log, or None. Exact-port match wins; else newest match."""
    tmp = os.environ.get("TMPDIR") or "/tmp"
    exact = [os.path.join(tmp, "ferry-logs", "cloud-proxy-%s.log" % port),
             os.path.join("/tmp", "ferry-logs", "cloud-proxy-%s.log" % port)]
    exact = [p for p in exact if os.path.exists(p)]
    exact += _var_folders_candidates(port)
    loose = []
    for d in (os.path.join(tmp, "ferry-logs"), "/tmp/ferry-logs"):
        e, l = _dir_candidates(d, port)
        exact += e
        loose += l
    for group in (exact, loose):
        group = [p for p in dict.fromkeys(group) if os.path.exists(p)]
        if group:
            try:
                return max(group, key=os.path.getmtime)
            except OSError:
                return group[0]
    return None


# ── incremental tail ────────────────────────────────────────────────────────
class Tailer:
    """`tail -f` over a path: survives truncation, rotation, and deletion.

    Offsets are BYTE offsets from a binary read, so a multi-byte or invalid
    UTF-8 sequence (the proxy log has both) can never desync the position.
    A trailing fragment with no newline yet is held back until it completes.

    Three distinct restart shapes are handled, because a size check alone is
    not enough: `ferry up` reopens the log with `>`, which TRUNCATES IN PLACE
    (same inode) and may immediately write MORE bytes than we had already
    read — so `size < offset` never fires and we would resume mid-line of the
    new content. The first SIG_LEN bytes are therefore fingerprinted; if that
    head changes, the file was rewritten and we start over at byte 0.
    """

    def __init__(self, path, from_start=False, port=DEFAULT_PORT, pinned=False,
                 report=None):
        self.path = path
        self.port = port
        self.pinned = pinned            # --log given: never rediscover elsewhere
        self.from_start = from_start
        self.report = report or (lambda msg: None)
        self.fid = None                 # (st_dev, st_ino) of the file we're on
        self.offset = 0
        self.buf = ""
        self.sig = b""                  # first SIG_LEN bytes of the file we're on
        self.first_attach = True

    def _attach(self, st):
        """Position ourselves on a file we have not read yet."""
        if self.first_attach and not self.from_start:
            self.offset = st.st_size    # tail from END: a restart re-ships nothing
        else:
            self.offset = 0             # --from-start, or a brand-new file after rotation
        self.buf = ""
        self.fid = (st.st_dev, st.st_ino)
        self.first_attach = False
        self._capture_sig(st.st_size)
        self.report("tailing %s from byte %d" % (self.path, self.offset))

    def _rewind(self, why):
        self.report("%s — rereading from byte 0" % why)
        self.offset = 0
        self.buf = ""

    def _capture_sig(self, size):
        """Fingerprint the head of the current file (never raises)."""
        n = min(SIG_LEN, max(size, 0))
        if n <= 0:
            self.sig = b""
            return
        try:
            with open(self.path, "rb") as f:
                self.sig = f.read(n)
        except OSError:
            self.sig = b""

    def _same_file_content(self, size):
        """False once the head bytes differ, i.e. the file was rewritten."""
        n = len(self.sig)
        if n == 0:                       # nothing to compare against yet
            return True
        if size < n:                     # truncated below the fingerprint
            return False
        try:
            with open(self.path, "rb") as f:
                head = f.read(n)
        except OSError:
            return True                  # unreadable right now; decide next poll
        return head == self.sig

    def poll(self):
        """Every complete line written since the last poll (never raises)."""
        if not self.path:
            self.path = None if self.pinned else find_log(self.port)
            if not self.path:
                return []
        try:
            st = os.stat(self.path)
        except OSError:
            # The file went away (ferry stopped, tmpdir reaped). Forget where we
            # were and look again next poll; a replacement is read from byte 0.
            if self.fid is not None:
                self.report("log %s disappeared — waiting for it to come back"
                            % self.path)
            self.fid = None
            self.buf = ""
            self.offset = 0
            self.sig = b""
            self.first_attach = False   # a replacement file is new content: read it all
            if not self.pinned:
                self.path = None
            return []

        fid = (st.st_dev, st.st_ino)
        if self.fid is None:
            self._attach(st)
        elif fid != self.fid:                       # rotated: different file, same name
            self.fid = fid
            self._rewind("log rotated (inode changed)")
            self._capture_sig(st.st_size)
        elif st.st_size < self.offset:              # truncated: ferry restarted the proxy
            self._rewind("log truncated (%d < %d)" % (st.st_size, self.offset))
            self._capture_sig(st.st_size)
        elif not self._same_file_content(st.st_size):
            # same inode, same-or-larger size, but a different head: `ferry up`
            # reopened the log with `>` and wrote fresh content over the old.
            self._rewind("log rewritten in place (head changed)")
            self._capture_sig(st.st_size)
        elif not self.sig and st.st_size > 0:
            self._capture_sig(st.st_size)           # file was empty at attach

        if st.st_size <= self.offset:
            return []
        try:
            with open(self.path, "rb") as f:
                f.seek(self.offset)
                chunk = f.read(MAX_READ)             # catch up over several polls
                self.offset = f.tell()
        except OSError as e:
            self.report("read failed: %s" % e)
            return []
        if not chunk:
            return []

        data = self.buf + chunk.decode("utf-8", "replace")
        parts = data.split("\n")
        self.buf = parts.pop()                       # incomplete tail line
        if len(self.buf) > MAX_PARTIAL:              # a "line" that will never end
            parts.append(self.buf)
            self.buf = ""
        return parts


# ── VictoriaLogs sink ───────────────────────────────────────────────────────
class Shipper:
    """Batches records and POSTs them as JSON-lines, retrying with backoff."""

    def __init__(self, vlogs, batch_size=BATCH_LINES, flush_interval=FLUSH_INTERVAL,
                 dry_run=False, report=None):
        base = vlogs.rstrip("/")
        self.url = base + "/insert/jsonline?_stream_fields=source,level"
        self.batch_size = batch_size
        self.flush_interval = flush_interval
        self.dry_run = dry_run
        self.report = report or (lambda msg: None)
        self.pending = []
        self.shipped = 0
        self.dropped = 0
        self.failures = 0
        self.last_flush = time.monotonic()
        self.retry_at = 0.0
        self.backoff = BACKOFF_START
        self.healthy = None                          # None until the first attempt
        self._last_time = 0.0

    # -- record timestamps -------------------------------------------------
    def stamp(self):
        """RFC3339 UTC, strictly increasing so a batch keeps its order."""
        t = time.time()
        if t <= self._last_time:
            t = self._last_time + 1e-6
        self._last_time = t
        dt = datetime.datetime.fromtimestamp(t, datetime.timezone.utc)
        return dt.strftime("%Y-%m-%dT%H:%M:%S.%f") + "Z"

    # -- buffering ---------------------------------------------------------
    def add(self, record):
        record["_time"] = self.stamp()
        self.pending.append(record)
        if len(self.pending) > MAX_PENDING:          # bound memory: oldest first
            over = len(self.pending) - MAX_PENDING
            del self.pending[:over]
            self.dropped += over

    def due(self, now=None):
        if not self.pending:
            return False
        now = time.monotonic() if now is None else now
        if now < self.retry_at:                      # still backing off
            return False
        return (len(self.pending) >= self.batch_size
                or (now - self.last_flush) >= self.flush_interval)

    @staticmethod
    def build_body(records):
        """Newline-delimited JSON — one object per line, one POST per batch."""
        return "\n".join(json.dumps(r, ensure_ascii=False, sort_keys=True)
                         for r in records).encode("utf-8")

    def flush(self, force=False):
        """Ship one batch. Returns True if it landed (or there was nothing)."""
        now = time.monotonic()
        if not self.pending:
            self.last_flush = now
            return True
        if not force and now < self.retry_at:
            return False
        batch = self.pending[:self.batch_size]
        if self.dry_run:
            sys.stdout.write(self.build_body(batch).decode("utf-8") + "\n")
            sys.stdout.flush()
            ok, err = True, None
        else:
            ok, err = self.post(batch)
        if ok:
            del self.pending[:len(batch)]
            self.shipped += len(batch)
            self.backoff = BACKOFF_START
            self.retry_at = 0.0
            if self.healthy is not True:
                if not self.dry_run:
                    self.report("VictoriaLogs accepting inserts (%s)" % self.url)
                self.healthy = True
        else:
            self.failures += 1
            self.healthy = False
            self.report("insert failed (%s) — retrying in %.0fs, %d line(s) queued"
                        % (err, self.backoff, len(self.pending)))
            self.retry_at = now + self.backoff
            self.backoff = min(self.backoff * 2, BACKOFF_MAX)
        self.last_flush = now
        return ok

    def post(self, records):
        """POST one batch. Never raises — returns (ok, error_string)."""
        req = urllib.request.Request(self.url, data=self.build_body(records),
                                     method="POST")
        req.add_header("Content-Type", CONTENT_TYPE)
        try:
            with urllib.request.urlopen(req, timeout=HTTP_TIMEOUT) as r:
                code = r.status
                r.read()
            return (200 <= code < 300), (None if 200 <= code < 300 else "HTTP %d" % code)
        except urllib.error.HTTPError as e:
            body = ""
            try:
                body = e.read().decode("utf-8", "replace")[:200]
            except Exception:
                pass
            return False, ("HTTP %s %s" % (e.code, body)).strip()
        except Exception as e:                       # connection refused, DNS, timeout…
            return False, str(e)

    def drain(self, deadline_seconds=5.0):
        """Best-effort final flush on shutdown."""
        end = time.monotonic() + deadline_seconds
        while self.pending and time.monotonic() < end:
            if not self.flush(force=True):
                break


# ── daemon ──────────────────────────────────────────────────────────────────
def run(args):
    def report(msg):
        sys.stderr.write("ferry-log-shipper: %s\n" % msg)
        sys.stderr.flush()

    port = args.port
    logpath = args.log or find_log(port)
    tailer = Tailer(logpath, from_start=args.from_start, port=port,
                    pinned=bool(args.log), report=report)
    shipper = Shipper(args.vlogs, batch_size=args.batch_size,
                      flush_interval=args.flush_interval, dry_run=args.dry_run,
                      report=report)

    print("ferry-log-shipper -> %s" % shipper.url)
    print("  log       : %s" % (logpath or "not found yet (will keep looking)"))
    print("  mode      : %s%s" % ("from start" if args.from_start else "tail from end",
                                  "  [dry-run: printing, not posting]" if args.dry_run else ""))
    print("  batching  : up to %d lines or %.1fs per POST, poll every %.1fs"
          % (args.batch_size, args.flush_interval, args.poll_interval))
    print("Ctrl-C to stop." if not args.once else "single pass (--once).")
    sys.stdout.flush()

    heartbeat = time.monotonic()
    try:
        while True:
            try:
                for line in tailer.poll():
                    if not line.strip():             # blank/banner padding
                        continue
                    shipper.add(parse_line(line))
                    if len(shipper.pending) >= shipper.batch_size:
                        shipper.flush()
                if shipper.due():
                    shipper.flush()
            except Exception as e:                   # a daemon never dies on one bad poll
                report("poll error: %s" % e)
            if args.once:
                shipper.drain()
                break
            now = time.monotonic()
            if args.verbose and (now - heartbeat) >= 60:
                heartbeat = now
                report("shipped=%d pending=%d dropped=%d failures=%d"
                       % (shipper.shipped, len(shipper.pending),
                          shipper.dropped, shipper.failures))
            time.sleep(args.poll_interval)
    except KeyboardInterrupt:
        print("\nferry-log-shipper stopping — flushing %d pending line(s)."
              % len(shipper.pending))
        shipper.drain()
    print("ferry-log-shipper: shipped=%d pending=%d dropped=%d failures=%d"
          % (shipper.shipped, len(shipper.pending), shipper.dropped, shipper.failures))
    return 0


def parse_args(argv=None):
    ap = argparse.ArgumentParser(
        description="Tail the ferry litellm proxy log into VictoriaLogs.")
    ap.add_argument("--vlogs", default=DEFAULT_VLOGS,
                    help="VictoriaLogs base URL (default %s)" % DEFAULT_VLOGS)
    ap.add_argument("--log", default=None,
                    help="proxy log path (auto-discovered if omitted)")
    ap.add_argument("--port", default=DEFAULT_PORT,
                    help="litellm proxy port used for log discovery (default %s)"
                         % DEFAULT_PORT)
    ap.add_argument("--from-start", action="store_true",
                    help="ship the whole existing log first (default: tail from the end, "
                         "so a restart never re-ships history)")
    ap.add_argument("--batch-size", type=int, default=BATCH_LINES,
                    help="max lines per POST (default %d)" % BATCH_LINES)
    ap.add_argument("--flush-interval", type=float, default=FLUSH_INTERVAL,
                    help="max seconds a partial batch waits (default %.1f)" % FLUSH_INTERVAL)
    ap.add_argument("--poll-interval", type=float, default=POLL_INTERVAL,
                    help="seconds between log reads (default %.1f)" % POLL_INTERVAL)
    ap.add_argument("--dry-run", action="store_true",
                    help="print the JSON-lines batches instead of POSTing them")
    ap.add_argument("--once", action="store_true",
                    help="do a single poll+flush and exit (smoke test)")
    ap.add_argument("--verbose", action="store_true",
                    help="log a shipped/pending/dropped heartbeat every 60s")
    return ap.parse_args(argv)


def main(argv=None):
    # Nothing above parse_args touches the filesystem, the network, or env, so
    # `--help` works on a bare shell with no ferry, no config, and no secrets.
    return run(parse_args(argv))


if __name__ == "__main__":
    sys.exit(main())
