From dd1723426e1f6f48f49604a90e58f86ce0261866 Mon Sep 17 00:00:00 2001 From: Derek Date: Sat, 19 Sep 2026 13:40:43 +1000 Subject: [PATCH] fix(docs): give the README a Context section and write the architecture doc The single most surprising fact about this repo was nowhere in it: dfe-engine imports logreducer at runtime and declares NO dependency on it, deliberately, in three lazy import sites inside the gated sampler modes. That is now the whole of "Where this sits", with the sites named. Also recorded, because all four are true and none was written down: nothing is deselected by default, so tests/integration/ is collected and then SKIPS without an endpoint or Docker; the version in the tree is a release behind the tag and PyPI; min_coverage gates at 75 while its own comment says 80; and vulture warns rather than fails. Verified against origin/main this pass: the three import sites read out of dfe-engine at sampling/service.py:200 and :236 and sampling/kafka_reader.py:134, with the non-declaration recorded at its pyproject.toml:114. Tag v3.5.0 against pyproject version 3.4.0, and PyPI serving 3.5.0 with 3.4.0 published 2026-07-03. Every Don't/Do/Why row traces to a comment in core.py, to CONTRIBUTING.md, or to open issue #5. Dropped "right now" from the CLI bullet while the file was open -- it was a DOC-TIME-1 finding and the next sentence already says it. Not verified: no test run. uv sync was not run and pytest was not run, so the commands are quoted from CONTRIBUTING.md, .hyperi-ci.yaml and pyproject.toml rather than from a run. --- README.md | 74 +++++++++++++++++++++- docs/architecture.md | 147 +++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 220 insertions(+), 1 deletion(-) create mode 100644 docs/architecture.md diff --git a/README.md b/README.md index 6d0574b..ea5a95c 100644 --- a/README.md +++ b/README.md @@ -8,7 +8,7 @@ Reduce gigabytes of logs to a small, representative sample - keeping the pattern **LogReducer is two tools in one package:** -- **A CLI you can use right now.** `logreducer app.log` reduces a file - or a SQL, ClickHouse, or Kafka source - straight from the shell. No code to write. +- **A CLI.** `logreducer app.log` reduces a file - or a SQL, ClickHouse, or Kafka source - straight from the shell. No code to write. - **A library with an IO-agnostic core.** The engine has zero IO dependencies and reduces any re-iterable stream of `str` lines (a `Source`). Embed it in your own pipeline; the engine never manages the connection. ## Features @@ -341,3 +341,75 @@ See [CONTRIBUTING.md](CONTRIBUTING.md) for the toolchain, test layout (including [Apache-2.0](LICENSE). Third-party attributions are recorded in [NOTICE](NOTICE). Copyright 2026 HYPERI PTY LIMITED. + +## Context + +Depth: [docs/architecture.md](docs/architecture.md). + +### What this is + +A reduction engine over an abstraction, plus a CLI on top of it. The core carries +zero IO dependencies and reduces any re-iterable stream of `str` lines, so it +never owns a connection or a loading path. It is not a log shipper, not a parser, +and not an aggregator front end -- it picks a representative subset of the lines +it is given and drops the rest. + +### Where things live + +| Path | Holds | +|---|---| +| `src/logreducer/core.py` | `LogReducer` -- the four modes, the multi-pass orchestration, stats, output | +| `src/logreducer/sources.py`, `sinks.py` | the two structural protocols that ARE the integration surface, plus `FileSource` and `FileSink` | +| `src/logreducer/sql.py`, `clickhouse.py`, `kafka.py` | the optional adapters, one extra each | +| `src/logreducer/patterns.py`, `anomaly.py`, `temporal.py` | Drain3 mining and fuzzy dedup, Isolation Forest, time-window grouping | +| `src/logreducer/memory.py` | `MemoryMonitor`, `BoundedDeduplicator`, the file read strategies | +| `src/logreducer/sampling.py`, `target.py` | dialect-aware SQL sampling, and the `reduce_to_target` batch loop | +| `src/logreducer/config.py` | `BigDialConfig`, the level presets, the `from_env` cascade | +| `src/logreducer/cli.py` | the `logreducer` entry point, dispatching on the `--dsn` scheme | +| `tests/unit/`, `tests/integration/` | server-free, and `integration`-marked against real services | +| `tests/testdata/` | gzipped, PII-cleansed slices of real public log datasets, with a manifest and a rebuild tool | + +### Commands that prove a change + +```bash +uv sync --all-extras +hyperi-ci check # the canonical local gate -- there is no Makefile +``` + +Individually: `uv run ruff format`, `uv run ruff check --fix`, +`uv run mypy src/logreducer`, `uv run ty check src/logreducer`, +`uv run pytest -q`, `uv build`. Three ways green lies: + +| It looks like | It is actually | +|---|---| +| `pytest -q` covered the adapters | Nothing is deselected by default -- `addopts` is only `-ra -q --strict-markers --strict-config`, so `tests/integration/` is collected and then SKIPS when there is neither a configured endpoint nor Docker. CI sets no env vars, so CI always takes the docker path and a local run may not have | +| the tree tells you the version | It does not. `pyproject.toml` and `VERSION` both say 3.4.0 while the newest tag is v3.5.0 and PyPI serves 3.5.0. semantic-release resolves the version at release time | +| coverage is at the stated floor | `min_coverage: 75` gates, and its own comment says actual is 80. `vulture` is set to `warn`, not fail, because a library's public API reads as unused to dead-code analysis | + +### What tends to bite + +| Don't | Do | Why | +|---|---|---| +| Hand `reduce()` a bare generator | Hand it something whose `__iter__` returns a FRESH iterator each call | `reduce()` counts lines in one pass then re-reads to process, and hybrid reads again. A one-shot iterator is drained by the count and leaves nothing behind, so it now raises `TypeError` rather than returning an empty result (`core.py`) | +| Assume a reused `LogReducer` starts clean | Nothing -- `_reset_components()` rebuilds them per run | The dedup seen-set, both Drain3 miners and the fuzzy LSH accumulate per line. Carried into the next run they drop every line as already-seen and return empty (`core.py`) | +| Share one deduplicator across hybrid's two passes | Rebuild it before the anomaly pass | The pattern pass consumes the seen-set, so the anomaly pass sees everything as seen and no anomalies survive. Hybrid silently degrades to pattern-only (`core.py`) | +| Assume an unknown config keyword is ignored | Read the `ValueError` and fix the name | A silently dropped override means a user's tuning quietly does nothing (`core.py`) | +| Read `LogPattern.template` as the final Drain3 template | Treat it as a first-seen snapshot until #5 lands | `extract_patterns` captures `cluster.get_template()` when a cluster FIRST appears, before later lines widen slots to `<*>`, so a varying cluster stores its raw first line. The `<*>`-count priority boost then runs on that stale snapshot (#5, open) | +| Hand-edit the version | Land a conventional commit and let semantic-release do it | The old regime wrote `VERSION` by hand and took five commits to get tagging working. CONTRIBUTING makes it a standing rule | + +### Where this sits + +One declared edge in `dfe-infra/suite.yaml`, and it runs the surprising way. + +| Repo | Direction | Kind | Mechanism | +|---|---|---|---| +| dfe-engine | outbound | `python-dep-undeclared` | dfe-engine imports this library at runtime and declares NO dependency on it, deliberately. `from logreducer import LogReducer` at `sampling/service.py:200`, `from logreducer.clickhouse import ClickHouseSource` at `:236`, and `from logreducer.kafka import KafkaSource` at `sampling/kafka_reader.py:134`. All three are lazy, inside the gated `smart` and `anomaly` sampler modes, and an `ImportError` is turned into an error telling the caller to use `mode=recent` or `mode=random` instead. There is no version range to test -- install a new version alongside dfe-engine and run the tests that exercise those three sites | + +`dfe-stack suite --consumer logreducer` returns an empty edge list: no repo in +that graph is something this library depends on. The runtime dependencies +(drain3, psutil, loguru, numpy, scikit-learn, typer) are packages, declared in +`pyproject.toml`, which is where they belong. + +The reason dfe-engine gives for not declaring the dependency -- pin it "once it +hits public PyPI" -- has been met since 3.4.0 was published on 2026-07-03, so +expect that edge to become an ordinary pinned dependency. diff --git a/docs/architecture.md b/docs/architecture.md new file mode 100644 index 0000000..0a75a9a --- /dev/null +++ b/docs/architecture.md @@ -0,0 +1,147 @@ +# logreducer architecture + +Why the engine is shaped the way it is. The README covers what it does and how to +drive it -- this page is the reasoning and the invariants a caller cannot infer +from the API. + +## The problem + +You have gigabytes of logs and a question. Reading all of it is not an option, +and a random sample is worse than it sounds: random throws away the rare lines, +which are usually the interesting ones, and keeps thousands of copies of the line +that repeats. What a reader actually wants is one of each SHAPE, plus the +outliers. + +Two constraints make that harder than picking a sample. + +- The input does not fit in memory, and neither does the set of UNIQUE lines. A + 10 GB log with 40 million distinct lines will not fit either way, so the engine + must hold neither. +- Some reduction modes need more than one look at the data, and the data might be + a database query or a Kafka topic rather than a file. + +## The shape: a reduction engine behind two protocols + +The core does no IO at all. Input is a `Source` -- anything whose `__iter__` +yields `str` and can be iterated more than once. Output is a `Sink` -- +`write(lines) -> int`. Those two structural protocols are the whole integration +surface. + +They are protocols rather than base classes on purpose. A `list[str]` is already +a `Source`, and an application's own object becomes one without importing +anything from here. + +The point of that boundary is ownership. The engine never manages a connection, +so there is nothing for it to leak, nothing for it to time out, and no +credential inside it. The adapters that DO own connections -- SQLAlchemy, +clickhouse-connect, confluent-kafka -- sit behind optional extras, and the core +never imports them. + +```mermaid +flowchart LR + A["any re-iterable
of str lines"] -->|Source| C["core
zero IO deps"] + C -->|"list[str]"| R["caller"] + C -->|"Sink.write()"| S["anywhere"] + AD["sql / clickhouse / kafka
adapters (extras)"] -.->|"also just Sources"| A +``` + +## Re-iterability is the load that protocol carries + +`reduce()` counts lines in one pass to compute the reduction ratio, then re-reads +the source to process it. Hybrid mode reads it a third time. So `__iter__` has to +return a FRESH iterator every call. + +A bare generator satisfies the type and breaks the contract, and the failure is +silent in the worst way -- the counting pass drains it and the processing pass +finds nothing, so the result is an empty list rather than an error. That is why +`reduce()` tests `iter(source) is source` and raises `TypeError` up front. + +The same requirement shapes every adapter. `SQLSource` re-runs its query per +pass. `KafkaSource` re-reads from the earliest offset and never commits. +`FileSource` re-opens the file and streams it again. A seeded sample is required +for the same reason: an unseeded one would hand pass two a different set of rows +than pass one, and the reducer would be comparing two populations. + +## What streams, and what cannot + +Pattern and temporal modes stream end to end. Exact dedup, optional fuzzy dedup +and Drain3 mining run as one generator pipeline, so the unique-line set is never +collected into a list. Peak memory is the bounded dedup cache plus the Drain3 +template store, and neither grows with the number of unique lines the source has. + +Anomaly mode cannot stream. Isolation Forest over a TF-IDF matrix is batch ML and +needs its rows at once. That is the one place a hard memory ceiling needs an +explicit sample, which is exactly what `anomaly_max_rows` is -- a reservoir cap +with a fixed seed, trading anomaly recall for a bounded matrix, and off unless +asked for. + +So picking a mode is also picking a memory shape. Worth knowing before reaching +for the tuning knobs. + +## Collecting a target instead of reducing everything + +`reduce_to_target` is a different loop for a different question: about N +representative lines, rather than all of the input reduced. It pulls fresh random +batches, reduces each, and accumulates distinct representatives until one of five +stop conditions fires -- target, exhausted, max_fetches, plateau, or memory. +Batch size is resized from the observed average row bytes, so peak memory stays +around one batch plus the accumulator. + +The stop reason is reported rather than hidden, because "plateau" and "target" +mean very different things about the answer you are holding. + +## Determinism + +Every sample in the engine is seeded. The anomaly reservoir and the +normal-context sample both use a fixed seed, and SQL sampling takes an explicit +`sample_seed`. + +That is not a security property, and there is no security requirement in picking +log lines. It is a correctness property the multi-pass design depends on: pass +two has to see what pass one saw. + +It is also why SQLite raises `SamplingNotSupported` for a seeded sample instead +of quietly giving an unseeded one. SQLite has no seedable RNG, so it cannot keep +the promise, and pretending otherwise would produce a result that looks +reproducible and is not. + +## Embedding in a host application + +Three seams, none of which needs the library to know anything about the host. + +| Seam | How | +|---|---| +| Config | Build a `BigDialConfig` from the host's own cascade and inject it, or call `BigDialConfig.from_env(prefix, fallback)` where the prefixed name beats the bare one. Keyword arguments still win on top | +| Logging | Off by default. `own_sinks=False` registers nothing, so records flow through the host's handlers formatted by the host's standard | +| IO | The `Source` and `Sink` protocols above | + +## Invariants + +1. **`__iter__` returns a fresh iterator.** Checked at the entry point, not + documented and hoped for. +2. **The core imports no IO library.** An adapter's dependency is an extra, and + the core never imports the adapter modules. +3. **Analysis state is per run.** `_reset_components()` rebuilds the dedup + seen-set, both Drain3 miners and the fuzzy LSH, so one `LogReducer` is + reusable across calls and hybrid's two passes cannot poison each other. + Without it every line reads as already-seen and the result is empty. +4. **An unknown config keyword raises.** A silently dropped override is tuning + that quietly does nothing, which is worse than a failure. +5. **The row-to-line convention is shared.** The database adapters take the FIRST + column, skip NULLs and skip blank lines, matching `FileSource`, so the same + data reduces identically whichever source carried it. +6. **Every sample is seeded.** See above. +7. **Config enums serialise as their `.value`**, so metadata survives + `json.dumps` and no `OutputFormat.LINE`-style repr leaks into output. +8. **A memory ceiling above 70% of available RAM is clamped with a warning** + rather than accepted. + +### Where the numbers are approximate + +`input_lines` and the reduction ratio come from the counting pass over the +source. A size-sampled `FileSource` yields its SAMPLED line count, so on a very +large file both figures describe the sample rather than the file. Say which when +quoting them. + +`estimate_processing` is a size-based estimate and nothing more -- it reads the +file size and the chosen read strategy, and does not look at the content.