[packages/astgrep] Initial

arekm arekm at pld-linux.org
Wed Aug 19 15:28:53 CEST 2026


commit a320b81605b6c91d825f1bc2c876b8350b318e3d
Author: Arkadiusz Miśkiewicz <arekm at maven.pl>
Date:   Wed Aug 19 15:28:05 2026 +0200

    Initial

 astgrep.py   | 481 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
 astgrep.spec |  49 ++++++
 2 files changed, 530 insertions(+)
---
diff --git a/astgrep.spec b/astgrep.spec
new file mode 100644
index 0000000..ff3ac48
--- /dev/null
+++ b/astgrep.spec
@@ -0,0 +1,49 @@
+Summary:	astgrep - exigrep for Asterisk logs
+Summary(pl.UTF-8):	astgrep - odpowiednik exigrep dla logów Asteriska
+Name:		astgrep
+Version:	0.1
+Release:	1
+License:	AGPL v3+
+Group:		Applications/System
+Source0:	%{name}.py
+BuildRequires:	rpm-pythonprov
+Requires:	python3-modules >= 1:3.9
+Suggests:	python3-backports.zstd
+BuildArch:	noarch
+BuildRoot:	%{tmpdir}/%{name}-%{version}-root-%(id -u -n)
+
+%description
+astgrep matches a pattern against Asterisk log entries and prints
+every entry belonging to the same call, not just the line that
+matched. Calls are identified by the call id Asterisk stamps on every
+line it logs for a call ([C-xxxxxxxx]). Unlike a plain grep it keeps
+multi-line entries such as SIP packet dumps whole, and it expands a
+match to the complete call. Command line options follow exigrep(8).
+Compressed logs (.gz/.bz2/.xz/.zst) are read directly.
+
+%description -l pl.UTF-8
+astgrep dopasowuje wzorzec do wpisów w logach Asteriska i wypisuje
+wszystkie wpisy należące do tego samego połączenia, a nie tylko
+pasującą linię. Połączenia są identyfikowane po identyfikatorze
+([C-xxxxxxxx]), który Asterisk umieszcza w każdej linii logu
+dotyczącej połączenia. W odróżnieniu od zwykłego grepa zachowuje w
+całości wpisy wielolinijkowe, na przykład zrzuty pakietów SIP, oraz
+rozszerza dopasowanie do całego połączenia. Opcje wiersza poleceń są
+zgodne z exigrep(8). Skompresowane logi (.gz/.bz2/.xz/.zst) są
+odczytywane bezpośrednio.
+
+%prep
+%setup -q -c -T
+
+%install
+rm -rf $RPM_BUILD_ROOT
+install -d $RPM_BUILD_ROOT%{_bindir}
+
+install -p %{SOURCE0} $RPM_BUILD_ROOT%{_bindir}/%{name}
+
+%clean
+rm -rf $RPM_BUILD_ROOT
+
+%files
+%defattr(644,root,root,755)
+%attr(755,root,root) %{_bindir}/%{name}
diff --git a/astgrep.py b/astgrep.py
new file mode 100755
index 0000000..ffd7be2
--- /dev/null
+++ b/astgrep.py
@@ -0,0 +1,481 @@
+#!/usr/bin/python3
+# AGPL v3+ PLD Linux Team
+"""astgrep - exigrep for Asterisk logs.
+
+Matches a pattern against log *entries* and prints every entry belonging to the
+same call, not just the line that matched.  Calls are identified by the call
+identifier Asterisk stamps on every line it emits for a call ([C-xxxxxxxx],
+logger.conf use_callids, on by default since Asterisk 11).
+
+Two things a plain grep cannot do:
+  * follow the pattern back to the whole call (grep gives one line out of 35);
+  * keep multi-line entries whole - SIP packet dumps span many lines and only
+    the header line carries the [C-...] tag, so grep silently drops the body.
+
+Options follow exigrep(8), a call standing in for a message: -l literal,
+-I case-sensitive, -M related (calls sharing a channel, i.e. transfers),
+-t<n> only calls longer than n seconds, -v invert.  Three deliberate departures:
+
+  * matching is per entry, not per line, because Asterisk entries are
+    multi-line while Exim's are not;
+  * -v selects calls with *no* matching entry.  exigrep inverts the per-line
+    test instead, which selects almost every message - not useful here;
+  * -M is on by default (--no-follow to disable).  A transfer or a callback is
+    part of the same conversation, which is not true of Exim's related messages.
+
+The log is read twice rather than held in memory the way exigrep does it:
+Asterisk 'full' logs with verbose/debug on run to hundreds of MB, and only the
+matched calls are worth buffering.  Non-seekable input is spooled to a tmpfile.
+"""
+
+import argparse
+import bz2
+import contextlib
+import gzip
+import importlib
+import lzma
+import os
+import re
+import shutil
+import signal
+import subprocess
+import sys
+import tempfile
+from datetime import datetime
+
+HEADER_RE = re.compile(r"^\[(?P<ts>[^\]]+)\]")
+GROUP_RE = re.compile(r"^# (?:C-[0-9a-fA-F]+ |entries with no call id)")
+CALLID_RE = re.compile(r"\[(C-[0-9a-fA-F]+)\]")
+CHANNEL_RE = re.compile(
+    r"\b(?:SIP|PJSIP|IAX2|Local|DAHDI|Motif|PJ)/[^\s\"',()\[\]]*?-[0-9a-f]{8}(?:;[12])?"
+)
+
+# Markers emitted by stock Asterisk core modules, used only by --summary.
+EXECUTING_RE = re.compile(r"Executing \[([^\]]+)\] \w+\(\"([^\"]*)\"")
+CALLED_RE = re.compile(r"app_dial\.c: Called (\S+)")
+ANSWERED_RE = re.compile(r"(\S+) answered (\S+)")
+STATUS_RE = re.compile(r"status is '([A-Z]+)'")
+SIPRESP_RE = re.compile(r"Got SIP response (\d{3}) \"([^\"]*)\"")
+HANGUP_RE = re.compile(r"Spawn extension \([^)]*\) exited non-zero")
+
+# logger.conf dateformat is configurable.  The yearless syslog-style one gets a
+# placeholder year (a leap year, so Feb 29 parses); only differences are used.
+TS_FORMATS = (
+    ("", "%Y-%m-%d %H:%M:%S.%f"),
+    ("", "%Y-%m-%d %H:%M:%S"),
+    ("1972 ", "%Y %b %d %H:%M:%S"),
+)
+
+
+class Entry:
+    """One log entry: its header line plus any continuation lines."""
+
+    __slots__ = ("callid", "lines", "ts")
+
+    def __init__(self, callid, ts, lines):
+        self.callid = callid
+        self.ts = ts
+        self.lines = lines
+
+    @property
+    def text(self):
+        return "\n".join(self.lines)
+
+
+def parse_ts(ts):
+    for prefix, fmt in TS_FORMATS:
+        try:
+            return datetime.strptime(prefix + ts, fmt)  # noqa: DTZ007 - only differences are used
+        except ValueError:
+            continue
+    return None
+
+
+def open_log(path):
+    ext = os.path.splitext(path)[1].lower()
+    if ext in (".zst", ".zstd"):
+        return open_zstd(path)
+    opener = {".gz": gzip.open, ".bz2": bz2.open, ".xz": lzma.open}.get(ext, open)
+    return opener(path, "rt", errors="replace")
+
+
+# All of these expose open(file, mode, ..., errors=...), so any one will do.
+# PyPI 'zstd' (PLD python3-zstd) is deliberately absent: it has no file or
+# stream API, only whole-buffer decompress(), which would mean slurping the
+# entire log into memory - twice, since the log is read in two passes.
+ZSTD_MODULES = ("compression.zstd", "backports.zstd", "zstandard", "pyzstd")
+
+
+def open_zstd(path):
+    for name in ZSTD_MODULES:
+        try:
+            module = importlib.import_module(name)
+        except ImportError:
+            continue
+        return module.open(path, "rt", errors="replace")
+    return zstd_command(path)
+
+
+ at contextlib.contextmanager
+def zstd_command(path):
+    """Decompress through the zstd binary, when no bindings are installed."""
+    try:
+        proc = subprocess.Popen(
+            ["zstd", "-dcq", "--", path],
+            stdout=subprocess.PIPE, text=True, errors="replace",
+        )
+    except FileNotFoundError:
+        sys.exit(f"astgrep: {path}: needs the zstd command, Python 3.14+, or one of "
+                 + ", ".join(ZSTD_MODULES[1:]))
+    try:
+        yield proc.stdout
+    finally:
+        proc.stdout.close()
+        if proc.wait() not in (0, -signal.SIGPIPE):
+            print(f"astgrep: zstd failed on {path} (exit {proc.returncode})",
+                  file=sys.stderr)
+
+
+def iter_entries(fh):
+    ts = callid = None
+    lines = []
+    for line in fh:
+        line = line.rstrip("\n")
+        if GROUP_RE.match(line):
+            # our own group header, fed back in through a pipe - not log content
+            while lines and not lines[-1]:
+                lines.pop()
+            continue
+        header = HEADER_RE.match(line)
+        if header:
+            if lines:
+                yield Entry(callid, ts, lines)
+            ts = header.group("ts")
+            found = CALLID_RE.search(line)
+            callid = found.group(1) if found else None
+            lines = [line]
+        elif lines:
+            lines.append(line)
+        # else: continuation before the first header (truncated log) - drop
+    if lines:
+        yield Entry(callid, ts, lines)
+
+
+def first_stamp(path):
+    """Timestamp of a log's first dated line, for putting rotated files in order."""
+    try:
+        with open_log(path) as fh:
+            for line in fh:
+                header = HEADER_RE.match(line)
+                if header:
+                    stamp = parse_ts(header.group("ts"))
+                    if stamp:
+                        return stamp.timestamp()
+    except OSError:
+        pass
+    return float("inf")
+
+
+def order_files(paths):
+    """Oldest first, so 'messages*' glob order (newest first) reads correctly."""
+    return sorted(paths, key=first_stamp) if len(paths) > 1 else paths
+
+
+def iter_all(paths):
+    for path in paths:
+        with open_log(path) as fh:
+            yield from iter_entries(fh)
+
+
+def make_matcher(pattern, literal, case_sensitive):
+    flags = 0 if case_sensitive else re.IGNORECASE
+    if literal:
+        pattern = re.escape(pattern)
+    try:
+        rx = re.compile(pattern, flags)
+    except re.error as exc:
+        sys.exit(f"astgrep: bad regex {pattern!r}: {exc} (-l matches it literally)")
+    return lambda text: rx.search(text) is not None
+
+
+def collect(paths, matches, args):
+    """Pass 1: which call ids to print, plus channel -> call ids for -M."""
+    hits = set()
+    seen = set()
+    orphans = 0
+    chan_ids = {}
+    span = {}
+    for entry in iter_all(paths):
+        text = entry.text
+        hit = matches(text)
+        if entry.callid is None:
+            orphans += hit != args.invert
+            continue
+        seen.add(entry.callid)
+        if hit:
+            hits.add(entry.callid)
+        if args.min_duration:
+            widen(span, entry.callid, entry.ts)
+        if args.follow:
+            for chan in CHANNEL_RE.findall(text):
+                chan_ids.setdefault(chan, set()).add(entry.callid)
+
+    wanted = seen - hits if args.invert else hits
+    # Expanding under -v could only re-add calls that *do* match, which is
+    # exactly what -v excludes; every non-matching call is already selected.
+    if args.follow and not args.invert:
+        wanted = expand(wanted, chan_ids)
+    if args.min_duration:
+        wanted, undated = filter_by_duration(wanted, span, args.min_duration)
+        if undated:
+            print(f"astgrep: {undated} calls kept, timestamps not understood for -t",
+                  file=sys.stderr)
+    return wanted, orphans
+
+
+def filter_by_duration(wanted, span, minimum):
+    keep = set()
+    undated = 0
+    for callid in wanted:
+        secs = seconds(*span[callid])
+        if secs is None:
+            undated += 1
+            keep.add(callid)
+        elif secs >= minimum:
+            keep.add(callid)
+    return keep, undated
+
+
+def widen(span, callid, ts):
+    """Track earliest/latest timestamp of a call; files may arrive out of order."""
+    stamp = parse_ts(ts)
+    if stamp is None:
+        span.setdefault(callid, (None, None))
+        return
+    low, high = span.get(callid, (None, None))
+    span[callid] = (stamp if low is None or stamp < low else low,
+                    stamp if high is None or stamp > high else high)
+
+
+def seconds(low, high):
+    """Length of a span, or None if its timestamps did not parse."""
+    return None if low is None or high is None else (high - low).total_seconds()
+
+
+def in_time_order(entries):
+    """Sort a call's entries chronologically, or keep file order if unparsable."""
+    stamps = [parse_ts(e.ts) for e in entries]
+    if any(stamp is None for stamp in stamps):
+        return entries
+    return [e for _, e in sorted(zip(stamps, entries), key=lambda pair: pair[0])]
+
+
+def expand(wanted, chan_ids):
+    """Pull in call ids sharing a channel with a wanted call (transfers, pickups)."""
+    id_ids = {}
+    for ids in chan_ids.values():
+        if len(ids) > 1:
+            for callid in ids:
+                id_ids.setdefault(callid, set()).update(ids)
+    result = set(wanted)
+    queue = list(wanted)
+    while queue:
+        for linked in id_ids.get(queue.pop(), ()):
+            if linked not in result:
+                result.add(linked)
+                queue.append(linked)
+    return result
+
+
+def summarize(callid, entries):
+    entry_point = caller = answered = status = sip_resp = None
+    dialed = []
+    hung_up = False
+    for entry in entries:
+        text = entry.text
+        if entry_point is None:
+            found = EXECUTING_RE.search(text)
+            if found:
+                entry_point, caller = found.group(1), found.group(2)
+        found = CALLED_RE.search(text)
+        if found and found.group(1) not in dialed:
+            dialed.append(found.group(1))
+        found = ANSWERED_RE.search(text)
+        if found:
+            answered = found.group(1)
+        found = STATUS_RE.search(text)
+        if found:
+            status = found.group(1)
+        found = SIPRESP_RE.search(text)
+        if found and int(found.group(1)) >= 400:
+            sip_resp = f"{found.group(1)} {found.group(2)}"
+        hung_up = hung_up or HANGUP_RE.search(text) is not None
+
+    if answered:
+        result = f"ANSWERED by {answered}"
+    elif status:
+        result = f"{status} ({sip_resp})" if sip_resp else status
+    elif sip_resp:
+        result = sip_resp
+    elif hung_up:
+        result = "NO ANSWER (hung up)"
+    else:
+        result = "-"
+
+    stamps = [parse_ts(e.ts) for e in entries]
+    known = [stamp for stamp in stamps if stamp is not None]
+    secs = seconds(min(known), max(known)) if known else None
+    span = f"{secs:.1f}s" if secs is not None else "?"
+    return (
+        f"{callid}  {entries[0].ts}  {span:>7}  {len(entries):<4}  "
+        f"{entry_point or '-':<30} {caller or '-':<34} -> "
+        f"{', '.join(dialed) or '-':<26} {result}"
+    )
+
+
+def start_key(entries):
+    """Sort key: call start, with unparsable timestamps sorted last."""
+    stamp = parse_ts(entries[0].ts)
+    return stamp.timestamp() if stamp else float("inf")
+
+
+def main():
+    ap = argparse.ArgumentParser(
+        description="Print every Asterisk log entry of each call matching PATTERN.",
+        epilog="Compressed logs (.gz/.bz2/.xz/.zst) are read directly; "
+        "'-' or no file reads stdin.",
+    )
+    ap.add_argument("pattern")
+    ap.add_argument("files", nargs="*", default=["-"])
+    ap.add_argument(
+        "-l", "--literal", action="store_true", help="PATTERN is literal text, not a regex"
+    )
+    ap.add_argument(
+        "-I", "--case-sensitive", action="store_true", help="match case-sensitively"
+    )
+    ap.add_argument(
+        "-v",
+        "--invert",
+        action="store_true",
+        help="select calls with no matching entry",
+    )
+    ap.add_argument(
+        "-M",
+        "--follow",
+        action=argparse.BooleanOptionalAction,
+        default=True,
+        help="also select calls sharing a channel with a selected one - transfers, "
+             "callbacks (default: on; --no-follow for one call id only)",
+    )
+    ap.add_argument(
+        "-t",
+        "--min-duration",
+        type=float,
+        metavar="N",
+        help="only calls lasting N seconds or longer",
+    )
+    ap.add_argument("-s", "--summary", action="store_true", help="one line per selected call")
+    ap.add_argument("--ids", action="store_true", help="print selected call ids only")
+    ap.add_argument(
+        "--interleave",
+        action="store_true",
+        help="one stream in log order, not grouped (no per-call buffering)",
+    )
+    ap.add_argument(
+        "--header",
+        action=argparse.BooleanOptionalAction,
+        default=None,
+        help="group header lines (default: on to a terminal, off into a pipe so "
+             "output can be fed to astgrep again)",
+    )
+    ap.add_argument(
+        "--orphans",
+        action="store_true",
+        help="also print matching entries that carry no call id",
+    )
+    args = ap.parse_args()
+    if args.header is None:
+        args.header = sys.stdout.isatty()
+
+    paths = args.files or ["-"]
+    stdin_copy = None
+    if "-" in paths:
+        # stdin is read twice, so spool it
+        with tempfile.NamedTemporaryFile("w+", prefix="astgrep.", delete=False) as spool:
+            shutil.copyfileobj(sys.stdin, spool)
+            stdin_copy = spool.name
+        paths = [stdin_copy if p == "-" else p for p in paths]
+    paths = order_files(paths)
+
+    try:
+        matches = make_matcher(args.pattern, args.literal, args.case_sensitive)
+        wanted, orphans = collect(paths, matches, args)
+
+        if orphans and not args.orphans:
+            print(
+                f"astgrep: {orphans} matching entries have no call id "
+                "(--orphans to show them)",
+                file=sys.stderr,
+            )
+
+        if args.ids:
+            for callid in sorted(wanted):
+                print(callid)
+            return 0 if wanted else 1
+
+        if not wanted and not (orphans and args.orphans):
+            return 1
+
+        if args.interleave and not args.summary:
+            for entry in iter_all(paths):
+                if entry.callid in wanted or (args.orphans and entry.callid is None
+                                              and matches(entry.text) != args.invert):
+                    print(entry.text)
+            return 0
+
+        calls = {}
+        loose = []
+        for entry in iter_all(paths):
+            if entry.callid in wanted:
+                calls.setdefault(entry.callid, []).append(entry)
+            elif (args.orphans and entry.callid is None
+                  and matches(entry.text) != args.invert):
+                loose.append(entry)
+
+        for callid, entries in calls.items():
+            calls[callid] = in_time_order(entries)
+        ordered = sorted(calls.items(), key=lambda kv: start_key(kv[1]))
+
+        first = True
+        for callid, entries in ordered:
+            if args.summary:
+                print(summarize(callid, entries))
+                continue
+            if args.header:
+                if not first:
+                    print()
+                print(f"# {callid}  {entries[0].ts} .. {entries[-1].ts} "
+                      f"({len(entries)} entries)")
+            first = False
+            for entry in entries:
+                print(entry.text)
+
+        if loose:
+            if args.header and not args.summary:
+                print("\n# entries with no call id")
+            for entry in loose:
+                print(entry.text)
+        return 0
+    finally:
+        if stdin_copy:
+            os.unlink(stdin_copy)
+
+
+if __name__ == "__main__":
+    try:
+        sys.exit(main())
+    except BrokenPipeError:
+        os.dup2(os.open(os.devnull, os.O_WRONLY), sys.stdout.fileno())
+        sys.exit(1)
+    except KeyboardInterrupt:
+        sys.exit(130)
================================================================

---- gitweb:

http://git.pld-linux.org/gitweb.cgi/packages/astgrep.git/commitdiff/a320b81605b6c91d825f1bc2c876b8350b318e3d



More information about the pld-cvs-commit mailing list