"""Reading the AI-tool list and log files. Log files are read line by line (streaming), so file size is limited by disk, not by memory. Plain and gzip-compressed files are both accepted. Each reader yields records of the form ``(domain, user, timestamp, blocked)``: * ``domain`` normalised host name (lower case, no ``www.``, no trailing dot) * ``user`` user, device or IP address as written in the log, or ``None`` * ``timestamp`` naive ``datetime`` in UTC when the log gives a time zone, otherwise as written in the log; ``None`` when absent * ``blocked`` ``True``/``False`` when the log says whether the request was blocked, otherwise ``None`` """ import codecs import csv import gzip import io import json import os import re import zlib from datetime import datetime, timedelta from urllib.parse import urlsplit DATA_DIR = os.path.join(os.path.dirname(os.path.abspath(__file__)), "data") TOOLS_FILE = os.path.join(DATA_DIR, "ai-tools-free-list.csv") FORMATS = ("auto", "csv", "cloudflare-dns", "text") # How many lines are inspected to detect the format of a file. SNIFF_LINES = 50 class InputError(Exception): """A problem with the input that the user can fix. The message is shown as is.""" # --------------------------------------------------------------------------- # Domain normalisation and matching # --------------------------------------------------------------------------- _HOST_RE = re.compile(r"^[a-z0-9_-]+(?:\.[a-z0-9_-]+)*$") _IPV6_RE = re.compile(r"^[0-9a-f]*:[0-9a-f:.]*$") def normalize_domain(value): """Return the normalised host name in ``value``, or ``None`` if there is none. Accepts bare host names, ``host:port``, ``host/path`` and full URLs. Lower-cases, removes one leading ``www.`` (and ``*.``) and trailing dots. """ if value is None: return None s = str(value).strip().strip("\"'").strip() if not s: return None if "://" in s: try: s = urlsplit(s).hostname or "" except ValueError: return None else: for sep in "/?#": i = s.find(sep) if i >= 0: s = s[:i] if s.startswith("["): # [IPv6]:port s = s[1:].split("]", 1)[0] elif s.count(":") == 1: # host:port s = s.split(":", 1)[0] s = s.strip().lower().rstrip(".") if s.startswith("*."): s = s[2:] if s.startswith("www."): s = s[4:] if not s: return None if not s.isascii(): try: s = s.encode("idna").decode("ascii") except (UnicodeError, ValueError): return None if _HOST_RE.match(s): return s if ":" in s and _IPV6_RE.match(s): return s return None TOOL_FIELDS = ("domain", "name", "category", "subcategory", "ai_type", "risk_level", "data_sovereignty") class Tool(object): """One entry of the AI-tool list.""" __slots__ = TOOL_FIELDS def __init__(self, **kw): for f in TOOL_FIELDS: setattr(self, f, kw.get(f, "")) def as_dict(self): return {f: getattr(self, f) for f in TOOL_FIELDS} def __repr__(self): return "Tool(%r, %r)" % (self.domain, self.name) def load_tools(path=None): """Load the AI-tool list. Returns ``{normalised_domain: Tool}``.""" path = path or TOOLS_FILE tools = {} with open(path, "r", encoding="utf-8-sig", newline="") as f: reader = csv.DictReader(f) missing = [c for c in TOOL_FIELDS if c not in (reader.fieldnames or [])] if missing: raise InputError("AI-tool list %s lacks columns: %s" % (path, ", ".join(missing))) for row in reader: dom = normalize_domain(row.get("domain")) if not dom: continue kw = {f: (row.get(f) or "").strip() for f in TOOL_FIELDS} kw["domain"] = dom for f in ("risk_level", "data_sovereignty"): kw[f] = kw[f].lower() or "not_assessed" tools[dom] = Tool(**kw) return tools class Matcher(object): """Matches a domain to a tool: equal to the tool domain or a subdomain of it.""" def __init__(self, tools): self.tools = tools def match(self, domain): d = domain while d: t = self.tools.get(d) if t is not None: return t i = d.find(".") if i < 0: return None d = d[i + 1:] return None # --------------------------------------------------------------------------- # Timestamps # --------------------------------------------------------------------------- _NUM_RE = re.compile(r"^-?\d+(?:\.\d+)?$") _ISO_RE = re.compile( r"^(\d{4})-(\d{1,2})-(\d{1,2})" r"(?:[Tt ](\d{1,2}):(\d{2})(?::(\d{2}))?(?:[.,](\d+))?)?" r"\s*(Z|UTC|GMT|[+-]\d{2}(?::?\d{2})?)?$", re.IGNORECASE, ) _EPOCH = datetime(1970, 1, 1) def _from_epoch(v): a = abs(v) if a >= 1e17: # nanoseconds (Cloudflare "unixnano") v = v / 1e9 elif a >= 1e14: # microseconds v = v / 1e6 elif a >= 1e11: # milliseconds v = v / 1e3 try: dt = _EPOCH + timedelta(seconds=v) except (OverflowError, ValueError): return None return dt if 1990 <= dt.year <= 2100 else None def parse_timestamp(value): """Parse RFC 3339 / ISO 8601, ``YYYY-MM-DD HH:MM:SS``, Unix time in seconds, milliseconds, microseconds or nanoseconds, and the common log format (``30/Sep/2026:12:00:00 +0000``). Returns a naive datetime (UTC when the input carried a time zone) or ``None``.""" if value is None or isinstance(value, bool): return None if isinstance(value, (int, float)): return _from_epoch(float(value)) s = str(value).strip() if not s: return None if _NUM_RE.match(s): return _from_epoch(float(s)) m = _ISO_RE.match(s) if m: y, mo, d, hh, mi, ss, frac, tz = m.groups() try: dt = datetime(int(y), int(mo), int(d), int(hh or 0), int(mi or 0), int(ss or 0), int((frac or "0")[:6].ljust(6, "0"))) except ValueError: return None if tz and tz.upper() not in ("Z", "UTC", "GMT"): sign = -1 if tz[0] == "-" else 1 digits = tz[1:].replace(":", "") offset = timedelta(hours=int(digits[:2]), minutes=int(digits[2:4] or 0)) dt = dt - sign * offset return dt try: aware = datetime.strptime(s, "%d/%b/%Y:%H:%M:%S %z") return (aware - aware.utcoffset()).replace(tzinfo=None) except ValueError: return None # --------------------------------------------------------------------------- # Blocked / allowed # --------------------------------------------------------------------------- # Cloudflare Gateway resolver decisions, names and numeric values as listed in # https://developers.cloudflare.com/cloudflare-one/insights/logs/dashboard-logs/gateway-logs/ _CF_DECISIONS = { "blockedbycategory": True, "3": True, "allowedonnolocation": False, "4": False, "allowedonnopolicymatch": False, "5": False, "blockedalwayscategory": True, "6": True, "overrideforsafesearch": False, "7": False, "overrideapplied": False, "8": False, "blockedrule": True, "9": True, "allowedrule": False, "10": False, } _BLOCK_WORDS = ("block", "deny", "denied", "drop", "reject", "sinkhole", "refuse") def decision_blocked(value, cloudflare=False): """Interpret an action/decision value. Returns True, False or None (unknown).""" if value is None or isinstance(value, bool): return None s = str(value).strip().lower() if not s: return None if cloudflare and s in _CF_DECISIONS: return _CF_DECISIONS[s] if _NUM_RE.match(s): return None return any(w in s for w in _BLOCK_WORDS) # --------------------------------------------------------------------------- # Column detection for CSV/TSV # --------------------------------------------------------------------------- # Candidate header names, normalised (lower case, letters and digits only), # in order of preference. DOMAIN_COLUMNS = ("domain", "domainname", "queryname", "qname", "query", "fqdn", "hostname", "host", "dsthost", "desthost", "destinationhost", "destinationhostname", "site", "url", "requesturl", "dest", "destination") USER_COLUMNS = ("email", "useremail", "user", "username", "srcuser", "identity", "devicename", "device", "clientname", "client", "srcip", "clientip", "sourceip") TIME_COLUMNS = ("timestamp", "datetime", "time", "date", "ts", "eventtime", "querytime", "logtime") ACTION_COLUMNS = ("action", "decision", "resolverdecision", "policyaction", "verdict") def _norm_header(name): return re.sub(r"[^a-z0-9]", "", str(name).lower()) class Columns(object): """Which columns of a CSV file hold the domain, user, time and action.""" def __init__(self, domain=None, users=(), time=None, action=None, names=None, guessed=False): self.domain = domain self.users = list(users) self.time = time self.action = action self.names = names or [] self.guessed = guessed def _name(self, i): if i is None: return None if i < len(self.names) and self.names[i].strip(): return self.names[i].strip() return "column %d" % (i + 1) def describe(self): d = {"domain": self._name(self.domain)} if self.users: d["user"] = " / ".join(self._name(i) for i in self.users) if self.time is not None: d["time"] = self._name(self.time) if self.action is not None: d["action"] = self._name(self.action) return d def find_columns(header): """Detect columns from a header row. ``domain`` is ``None`` if no known name matches.""" norm = [_norm_header(h) for h in header] def first(cands): for c in cands: if c in norm: return norm.index(c) return None users = [] for c in USER_COLUMNS: if c in norm and norm.index(c) not in users: users.append(norm.index(c)) return Columns(domain=first(DOMAIN_COLUMNS), users=users, time=first(TIME_COLUMNS), action=first(ACTION_COLUMNS), names=list(header)) _TLD_RE = re.compile(r"\.[a-z][a-z0-9-]*[a-z0-9]$") def _guess_columns(rows): """Guess columns of a CSV file without a header row from its content.""" ncols = max((len(r) for r in rows), default=0) n = len(rows) if not n or not ncols: return Columns() def score(col, test): return sum(1 for r in rows if col < len(r) and test(r[col])) def looks_domain(v): d = normalize_domain(v) return bool(d) and bool(_TLD_RE.search(d)) dom_scores = [score(c, looks_domain) for c in range(ncols)] best = max(range(ncols), key=lambda c: dom_scores[c]) if dom_scores[best] < 0.5 * n: return Columns() cols = Columns(domain=best, guessed=True) others = [c for c in range(ncols) if c != best] for c in others: if score(c, lambda v: parse_timestamp(v) is not None) >= 0.8 * n: cols.time = c break for c in others: if c != cols.time and score(c, lambda v: "@" in v and "." in v) >= 0.8 * n: cols.users = [c] break return cols def sniff_delimiter(line): counts = {d: line.count(d) for d in (",", "\t", ";", "|")} best = max(counts, key=lambda d: counts[d]) return best if counts[best] else "," def split_line(line, delim): """Split one physical line into fields. Returns ``None`` for malformed quoting. Records never span lines: a stray quote cannot swallow the rest of the file. """ s = line.rstrip("\r\n") if '"' not in s: return s.split(delim) try: return next(csv.reader([s], delimiter=delim, strict=True)) except (csv.Error, StopIteration): return None # --------------------------------------------------------------------------- # Opening files # --------------------------------------------------------------------------- def open_text(path): """Open a plain or gzip-compressed text file (UTF-8, UTF-8 with BOM or UTF-16).""" try: raw = open(path, "rb") except OSError as e: raise InputError("%s: cannot open file (%s)" % (path, e.strerror or e)) try: if raw.peek(2)[:2] == b"\x1f\x8b": raw.close() raw = gzip.open(path, "rb") head = raw.peek(4)[:4] except (OSError, EOFError, zlib.error) as e: raw.close() raise InputError("%s: cannot read file (%s)" % (path, e)) if head.startswith(codecs.BOM_UTF16_LE) or head.startswith(codecs.BOM_UTF16_BE): enc = "utf-16" else: enc = "utf-8-sig" return io.TextIOWrapper(raw, encoding=enc, errors="replace", newline="") def _nul_free(lines): for line in lines: yield line.replace("\x00", "") if "\x00" in line else line def head_lines(path, n=SNIFF_LINES): """First ``n`` non-blank lines of a file, line endings removed.""" out = [] with open_text(path) as f: try: for line in _nul_free(f): s = line.rstrip("\r\n") if s.strip(): out.append(s) if len(out) >= n: break except (OSError, EOFError, zlib.error) as e: if not out: raise InputError("%s: cannot read file (%s)" % (path, e)) return out def _text_value(line): return line.split("#", 1)[0].strip() def detect_format(path): """Return ``csv``, ``cloudflare-dns`` or ``text`` for a file.""" name = os.path.basename(path) lines = head_lines(path) if not lines: raise InputError("%s: the file is empty" % name) first = lines[0].lstrip() if first.startswith("{"): for ln in lines: try: obj = json.loads(ln) except ValueError: continue if isinstance(obj, dict) and "QueryName" in obj: return "cloudflare-dns" raise InputError( "%s: JSON lines without a QueryName field. JSON input is supported for " "Cloudflare Gateway DNS logs (Logpush dataset gateway_dns); export other " "logs as CSV/TSV or as a plain list of domains." % name) if first.startswith("["): raise InputError( "%s: JSON arrays are not supported. Use JSON lines (one object per line, " "as written by Cloudflare Logpush), CSV/TSV, or a plain list of domains." % name) delim = sniff_delimiter(lines[0]) header = split_line(lines[0], delim) or [] if find_columns(header).domain is not None: return "csv" # No delimiter anywhere in the sample: a plain list, one domain (or URL) per line. body = [_text_value(ln) for ln in lines if not ln.lstrip().startswith("#")] if not any(c in ln for ln in body for c in ",\t;|"): return "text" return "csv" # --------------------------------------------------------------------------- # Readers # --------------------------------------------------------------------------- class FileStats(object): """Counts for one input file, shown in the report so the numbers can be checked.""" def __init__(self, path, fmt=None): self.path = path self.name = os.path.basename(path) self.format = fmt self.columns = {} self.lines = 0 # non-blank data lines (header row not counted) self.records = 0 # lines used in the report self.skipped = 0 # lines that could not be used (lines = records + skipped) self.bad_time = 0 # records whose timestamp could not be read (still counted) self.notes = [] def as_dict(self): return {"file": self.name, "format": self.format, "columns": self.columns, "lines": self.lines, "records": self.records, "skipped": self.skipped, "unreadable_timestamps": self.bad_time, "notes": list(self.notes)} def _read_csv(path, stats): with open_text(path) as f: lines = _nul_free(f) first = None for line in lines: if line.strip(): first = line break if first is None: return delim = sniff_delimiter(first) header = split_line(first, delim) or [] cols = find_columns(header) pending = [] if cols.domain is None: # No header row: guess columns from the first rows, which are data. pending.append(first) for line in lines: if line.strip(): pending.append(line) if len(pending) >= 200: break sample = [r for r in (split_line(ln, delim) for ln in pending) if r] cols = _guess_columns(sample) if cols.domain is None: raise InputError( "%s: could not find a domain column. Add a header row naming it, " "for example 'domain', 'query_name', 'host' or 'url'." % stats.name) stats.notes.append("no header row found; columns guessed from content") stats.columns = cols.describe() if delim == "\t": stats.format = "csv (tab-separated)" elif delim != ",": stats.format = "csv (%r-separated)" % delim def rows(): for ln in pending: yield ln for ln in lines: yield ln dcol, ucols, tcol, acol = cols.domain, cols.users, cols.time, cols.action for line in rows(): if not line.strip(): continue stats.lines += 1 row = split_line(line, delim) if row is None or len(row) <= dcol: stats.skipped += 1 continue dom = normalize_domain(row[dcol]) if dom is None: stats.skipped += 1 continue user = None for i in ucols: if i < len(row): v = row[i].strip() if v: user = v break ts = None if tcol is not None and tcol < len(row): ts = parse_timestamp(row[tcol]) if ts is None: stats.bad_time += 1 blocked = None if acol is not None and acol < len(row): blocked = decision_blocked(row[acol]) stats.records += 1 yield dom, user, ts, blocked def _read_text(path, stats): stats.columns = {"domain": "one domain per line"} with open_text(path) as f: for line in _nul_free(f): s = line.strip() if not s or s.startswith("#"): continue stats.lines += 1 dom = normalize_domain(_text_value(s)) if dom is None: stats.skipped += 1 continue stats.records += 1 yield dom, None, None, None def _cf_user(obj): email = obj.get("Email") if isinstance(email, str): email = email.strip() # Cloudflare writes non_identity@.cloudflareaccess.com when no # user identity was available; fall back to the source IP in that case. if email and not email.lower().startswith("non_identity@"): return email ip = obj.get("SrcIP") if isinstance(ip, str) and ip.strip(): return ip.strip() return None def _read_cloudflare(path, stats): stats.columns = {"domain": "QueryName", "user": "Email / SrcIP", "time": "Datetime", "action": "ResolverDecision"} seen = {"Email": False, "SrcIP": False, "Datetime": False, "ResolverDecision": False} with open_text(path) as f: for line in _nul_free(f): s = line.strip() if not s: continue stats.lines += 1 try: obj = json.loads(s) except ValueError: stats.skipped += 1 continue if not isinstance(obj, dict): stats.skipped += 1 continue qn = obj.get("QueryName") dom = normalize_domain(qn) if isinstance(qn, str) else None if dom is None: stats.skipped += 1 continue for k in seen: if not seen[k] and obj.get(k) not in (None, ""): seen[k] = True ts = None if "Datetime" in obj: ts = parse_timestamp(obj.get("Datetime")) if ts is None: stats.bad_time += 1 stats.records += 1 yield dom, _cf_user(obj), ts, decision_blocked(obj.get("ResolverDecision"), cloudflare=True) users = [k for k in ("Email", "SrcIP") if seen[k]] if users: stats.columns["user"] = " / ".join(users) else: stats.columns.pop("user") for key, field in (("time", "Datetime"), ("action", "ResolverDecision")): if not seen[field]: stats.columns.pop(key) _READERS = {"csv": _read_csv, "text": _read_text, "cloudflare-dns": _read_cloudflare} def read_log(path, stats, fmt="auto"): """Yield records from one log file, updating ``stats`` as it goes.""" if not os.path.exists(path): raise InputError("%s: file not found" % path) if os.path.isdir(path): raise InputError("%s: is a directory, not a log file" % path) if fmt not in FORMATS: raise InputError("unknown format %r (choose from %s)" % (fmt, ", ".join(FORMATS))) if fmt == "auto": fmt = detect_format(path) stats.format = fmt try: for rec in _READERS[fmt](path, stats): yield rec except (OSError, EOFError, zlib.error) as e: if stats.lines == 0: raise InputError("%s: cannot read file (%s)" % (stats.name, e)) stats.notes.append("reading stopped after %d lines: %s (damaged or truncated file?)" % (stats.lines, e))