diff --git a/.github/workflows/release-gate.yml b/.github/workflows/release-gate.yml index d2f2c32..334008a 100644 --- a/.github/workflows/release-gate.yml +++ b/.github/workflows/release-gate.yml @@ -13,8 +13,13 @@ concurrency: cancel-in-progress: true jobs: - release-gate: - runs-on: ubuntu-latest + platform-gate: + name: release-gate (${{ matrix.os }}) + strategy: + fail-fast: false + matrix: + os: [ubuntu-latest, macos-latest, windows-latest] + runs-on: ${{ matrix.os }} timeout-minutes: 10 steps: - name: Check out repository @@ -41,7 +46,7 @@ jobs: run: python -W error::ResourceWarning -m unittest discover -s tests -v - name: Compile Python sources - run: python -m py_compile run_osintai.py src/osintai/*.py tests/*.py + run: python -m compileall -q src tests run_osintai.py - name: Run security analysis run: bandit -q -r src run_osintai.py -ll -ii @@ -51,5 +56,16 @@ jobs: - name: Verify release identity run: | - python run_osintai.py --version | grep -Fx "OSINTai 4.0.0" - python run_osintai.py --help >/dev/null + python -c "import sys; sys.path.insert(0, 'src'); from osintai import __version__; assert f'# OSINTai v{__version__} ' in open('README.md', encoding='utf-8').read()" + python run_osintai.py --help + + release-gate: + name: release-gate + if: ${{ always() }} + needs: platform-gate + runs-on: ubuntu-latest + steps: + - name: Require every platform gate + env: + PLATFORM_GATE_RESULT: ${{ needs.platform-gate.result }} + run: test "$PLATFORM_GATE_RESULT" = "success" diff --git a/.gitignore b/.gitignore index dabc8e7..bd19ff3 100644 --- a/.gitignore +++ b/.gitignore @@ -27,8 +27,7 @@ id_ed25519 *.ovpn # Local OSINT inputs/outputs may contain PII, targets, proxies, or findings -data/runs/ -data/cache/ +/data/ cases/ /OSINTAI_SYNTHESIS/ /TODO @@ -62,6 +61,7 @@ build/ .pytest_cache/ .mypy_cache/ .tox/ +.test-tmp/ # IDEs .vscode/ diff --git a/BENCHMARKS.md b/BENCHMARKS.md new file mode 100644 index 0000000..84b4c9f --- /dev/null +++ b/BENCHMARKS.md @@ -0,0 +1,62 @@ +# OSINTai 4.2.0 offline benchmarks + +Measured on 2026-09-10 using Python 3.12.13, Darwin +arm64. These are single-run observations, not statistical guarantees. +No page fetching or model calls were performed. Source captures and existing reports +were preserved; the benchmark produced new extraction checkpoints and report bundles. + +## Saved crawls + +Default analysis settings: 200,000 characters per page; 16 MB text cache per reader; +100,000 co-occurrence candidate-pair budget; 30-second page and 120-second stage deadlines. +Timing includes source hashing, worker startup, analysis, report validation, and publication. + +| Saved run | Cache | Pages | Elapsed | Cache hits/misses | Parent peak RSS | Maximum child RSS | +| --- | --- | ---: | ---: | ---: | ---: | ---: | +| 20260909_212714 | cold | 152 | 18.133 s | 0/152 | 91.9 MiB | 84.4 MiB | +| 20260909_212714 | warm | 152 | 1.644 s | 152/0 | 86.8 MiB | 84.4 MiB | +| 20260717_144145 | cold | 159 | 18.923 s | 2/157 | 85.5 MiB | 92.6 MiB | + +The 159-page cold run reused two identical page contents within the same invocation. +RSS is reported separately for the parent and the largest individual child; summing the +columns is not a measurement of simultaneous process-tree memory. Warm means extraction +checkpoints were reused, not that every filesystem or OS cache was controlled. + +Both crawls explicitly report partial analytical coverage. The 152-page crawl excluded +92,369 potential pairs on pages exceeding the 60-identifier pairing threshold. The +159-page crawl excluded 161,328 such pairs, omitted 25,125 low-ranked correlation rows +at the output cap, and recorded 46 temporal parsing errors from source date values. +Neither run exhausted its candidate-pair budget. Publication succeeded and these limits +and errors remain visible in each bundle's summary and report. + +## Boilerplate fixture + +The synthetic fixture has 10,000 pages, four shared footer identifiers, and one +unique email/handle pair per page. All four footer identifiers are excluded from pairing. + +- Entity indexing: **0.053 seconds**. +- Correlation: **0.079 seconds**. +- Examined co-occurrence pairs: **10,000**. +- Pairs omitted by budget: **0**. +- Parent peak RSS: **60.8 MiB**. + +## Reproduction + +Run from the repository using its Python environment. Benchmark memory collection uses +`resource`, available on macOS/Linux; the unit-test CI matrix also covers Windows. + +```bash +python tests/benchmark_analysis.py RUN_ID --output .test-tmp/benchmark-cold.json +python tests/benchmark_analysis.py RUN_ID --output .test-tmp/benchmark-warm.json +python tests/benchmark_correlation.py --pages 10000 --pair-budget 100000 +``` + +A run is cold only if no matching extractor-version/content checkpoints exist. No +benchmark deletes prior checkpoints or source material to manufacture a cold run. + +## Verification + +All **119 offline tests** passed locally with `ResourceWarning` treated as an error. +Correctness lint, Bandit at the release gate's severity/confidence thresholds, compilation, +CLI version/help checks, and whitespace checks passed. The release workflow runs the same +gate on Linux, macOS, and Windows. diff --git a/ENHANCEMENTS.md b/ENHANCEMENTS.md new file mode 100644 index 0000000..5924c6d --- /dev/null +++ b/ENHANCEMENTS.md @@ -0,0 +1,96 @@ +# OSINTai enhancements + +## Integrated in 4.2.0 + +All seven enhancement areas from the 4.1.0 follow-up list are implemented: + +1. **Process-isolated deadlines.** Live HTML parsing, indicator extraction, hunt matching, + and Simhash run in disposable spawn workers. Saved-page extraction, entity indexing, + deterministic/model stages, and training export have deadlines. CPU timeouts, crashed + children, cancellation, and large IPC results have offline tests. Workers are terminated + and reaped. Raw HTML is retained before parsing, and failures are recorded. CLI controls: + `--analysis-page-timeout` (30 seconds), `--analysis-stage-timeout` (120 seconds). +2. **Content-addressed extraction checkpoints.** `.extraction_cache/` keys include scanned + text SHA-256, extractor implementation hash, character limit, and truncation state. + Atomic JSON checkpoints carry schema/checksum validation and per-kind coverage counts. + Repeated pages and resumed analyses reuse valid results. Cache writes failing does not + discard successful extraction. Secret values and arbitrary JWT claims are excluded; + secret-shaped values are also removed from cached indicator fields, including values + beyond the findings cap. Reports reconstruct findings from counts and decode validity. +3. **Adversarial regex audit.** Email, ASCII-domain, credential, and JWT extraction consume + maximal tokens or lines before bounded validation. Oversized malformed candidates cannot + restart an unbounded suffix search. Unicode scanning remains linear. External-deadline + tests cover repeated punctuation, failed domain/email suffixes, JWT-like strings, and + long credential lines. Fixed detector limits and their omissions are reported. +4. **Model quality and bounded retry.** Results distinguish `ok`, `empty`, `invalid`, + `missing`, `timed_out`, `error`, and `skipped`. Page attribution comes from saved file + identity rather than model assertions. Manifests record model names, validated page + counts, and optional stage response counts. `--retry-model RUN_ID --retry-limit 20 + --retry-timeout 60` retries unsuccessful/skipped saved pages through local Ollama, + without fetching pages or overwriting original analyses. Later offline runs use the + published `model_retry_latest.json` overlay only when its resolved paths stay inside + the run, its manifest is completed, and its source hashes still match current inputs. +5. **Hunt offsets and URL provenance.** Unicode expansion maps matches back to original + start/end offsets. URL detection uses original text and rejects URLs crossing a snippet + boundary. Indicator provenance identifies `html_attribute`, `page_prose`, and + `html_source`; resolved relative HTML links are included. +6. **Bounded correlation and text caching.** `--correlation-pair-budget` defaults to + 100,000 examined co-occurrence pairs. Omitted pairs, oversized-page exclusions, and + output truncation are explicit. Source membership uses sets, and identity evidence + avoids repeatedly scanning a common handle's entire source list. Each text reader has + a `--text-cache-bytes` LRU budget (16 MB default), accounting for string and entry + overhead. Cache peaks, evictions, missing files, and truncation are reported. +7. **Transactional analysis publication.** A hidden incomplete directory is populated + and validated before an atomic rename publishes an immutable `analysis_*/` bundle. + The bundle includes the human report, JSON/JSONL results, summary, manifest, and optional + training export. `analysis_latest.json` advances only after successful publication. + Attempt status records running/failed/completed state, and manifests retain full source + hashes and options. Source changes during analysis prevent publication. File contents + are flushed; POSIX directory metadata is fsynced. Original captures and prior bundles + remain intact. Completed publication may still report partial analytical coverage. + +The existing 4.1.0 fixes are retained: Unicode scanning, offline recovery, bounded text +reads, failure isolation, clean interruption, HTML text boundaries, complete summaries, +empty artifacts, numeric validation, and explicit profile overrides. + +## Verification and performance + +See `tests/test_enhancements.py` for deterministic offline acceptance checks and +`tests/benchmark_analysis.py` for saved-crawl runtime and memory measurements. CI now +runs the release gate on macOS, Linux, and Windows. Local verification was performed on +macOS with Python 3.12. + +Measured results are recorded in `BENCHMARKS.md`. RSS figures separate parent memory +from the maximum child RSS; they are not a measured concurrent process-tree total. + +Scanning is O(N) for bounded token validators and fixed extraction caps. Entity indexing +is expected O(S), with S source/indicator observations. Correlation generation is bounded +by O(S + B), followed by O(K log K) ordering of K candidates, where B is the pair budget. +Core retained analysis memory is O(S + E + B + C + R + L): E entities, C the text-cache budget, +R retained result rows, and L the per-job input bound. Disposable workers add serialization copies and startup cost. +Hunt searches cost O(TN) for T terms, with at most 500 reported hits; Unicode offset mapping +uses O(N) space only when lowercasing expands characters. Full source hashing streams +files with a 1 MB buffer. Thread timeouts and unbounded caches were rejected because they +cannot provide containment and predictable memory use. + +## Recommended next work + +1. **Public-suffix-aware domain families.** `_registrable()` still uses the last two + labels; names beneath suffixes such as `co.uk` can be incorrectly grouped. Bundle a + versioned Public Suffix List for deterministic offline grouping and ownership caveats. +2. **Cache retention and quotas.** Add explicit age/size-based pruning for extraction + checkpoints and old report/retry bundles, preserving referenced evidence and active + attempts. Current cache memory is bounded, but disk retention is intentionally additive. +3. **Typed streaming artifact ingestion.** Stream large JSONL files and validate record + schemas with line-level corruption counts. Current loaders materialize metadata and + silently skip malformed JSON lines; text-cache limits do not bound total index memory. +4. **Reduce process startup overhead.** Benchmark a supervised, recyclable worker design + that preserves hard per-job termination and isolation. The current spawn-per-job design + is portable and simple, but cold runs pay startup cost for every distinct page. +5. **More precise reproducibility and memory telemetry.** Inject a reference clock for + temporal analysis and generated timestamps; measure concurrent process-tree RSS, and + add Windows benchmark memory collection. Validate filesystem crash durability on the + supported filesystems; Windows does not provide POSIX directory-fsync semantics here. +6. **Raw-only parsing recovery.** Offer an explicit offline command to reparse preserved + HTML from live extraction failures. Current `--analyze-only` consumes saved page text + and reports raw-only failures rather than reconstructing missing text automatically. diff --git a/README.md b/README.md index 0358585..837a60a 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ OSINTai Logo -# OSINTai v4 - Advanced Local-First OSINT Web Crawler +# OSINTai v4.2.0 - Advanced Local-First OSINT Web Crawler [![License: MIT](https://img.shields.io/badge/License-MIT-yellow.svg)](https://opensource.org/licenses/MIT) [![Python 3.10+](https://img.shields.io/badge/python-3.10+-blue.svg)](https://www.python.org/downloads/) @@ -253,7 +253,7 @@ usage: run_osintai.py [-h] [--seed SEED] [--depth DEPTH] [--max MAX] [--no-ollama] [--hunt HUNT] [--hunt-max HUNT_MAX] [--run-id RUN_ID] -OSINTai v4 (async crawling and analysis) +OSINTai 4.2.0 (async crawling and analysis) required arguments: --seed SEED Seed URL (or use seed_urls.txt file) @@ -290,6 +290,80 @@ optional deeper analysis (off by default): ### Analysis Layer +**Recover a completed crawl without fetching pages again:** + +```bash +python run_osintai.py --analyze-only 20260909_212714 +``` + +This mode is fully offline and does not contact Ollama. Each invocation writes to a new +`data/runs/RUN_ID/reanalysis_*/` directory, preserving the source crawl and earlier reports. +It can also use `--evaluate` and `--training-export`; model-assisted modes are rejected. +It uses existing saved page analyses and the latest explicit model-retry overlay. It does +not regenerate missing model responses. + +### Analysis reliability and recovery (4.2) + +Live HTML parsing and saved-page extraction run in disposable worker processes with a +30-second deadline. CPU analysis stages, including indexing and dataset export, have a +120-second deadline. A failed worker is terminated and reaped; the remaining analysis +records partial coverage and continues. Raw HTML is saved before live parsing. + +```bash +python run_osintai.py --analyze-only RUN_ID \ + --analysis-page-timeout 30 --analysis-stage-timeout 120 \ + --analysis-max-chars 200000 --text-cache-bytes 16000000 \ + --correlation-pair-budget 100000 +``` + +Successful page extraction is cached under `RUN_ID/.extraction_cache/`. Keys include the +scanned-text SHA-256, extractor implementation version, character limit, and truncation +state. Repeated or interrupted runs reuse valid results. Checkpoints retain secret counts, +entropy counts, and JWT decode validity; matched secret values and arbitrary JWT claims +are excluded. Original page captures remain unchanged. Cache hits, invalid entries, +extraction failures, text coverage, and per-kind indicator omissions are reported. + +Text readers use an LRU cache with a byte budget that accounts for strings and entry +overhead. Correlation examines at most the configured number of co-occurrence pairs. +Its stage statistics distinguish budget omissions, oversized-page omissions, and output +row truncation. Common footer entities are excluded from pairing and indexed with sets. +Hunt offsets refer to original text even when Unicode lowercasing expands characters; +URLs crossing a snippet boundary are excluded. `url_provenance` distinguishes HTML +links, prose URLs, and URLs found elsewhere in HTML source. + +Model results distinguish `ok`, `empty`, `invalid`, `missing`, `timed_out`, `error`, and +`skipped`, with per-model counts in the manifest. Explicitly retry unsuccessful or skipped +saved-page model analyses using local Ollama: + +```bash +python run_osintai.py --retry-model RUN_ID --retry-limit 20 --retry-timeout 60 +python run_osintai.py --analyze-only RUN_ID --evaluate +``` + +Retries make no page-fetch requests and issue at most one model request per selected page. +They preserve original model outputs and publish a new `model_retry_*/` directory. The +`model_retry_latest.json` pointer selects the overlay used by later offline analyses. +`--model`, `--prompt-profile`, and `--analysis-max-chars` also apply to retries. + +Each analysis first writes a hidden `.analysis_*.incomplete/` staging directory. After +JSON validation and file flushing, it publishes an immutable `analysis_*/` bundle containing +the report, result files, summary, optional training export, and manifest. Consult +`analysis_latest.json` to locate the latest completed bundle; older root-level analysis +files belong to earlier versions and are not updated. Separate `analysis_*.status.json` +files record running, failed, or completed attempts. A hard process kill can leave a +running status and an incomplete directory; these are never selected by the latest pointer. +Completed bundles can report partial analytical coverage: completion means publication +succeeded. Manifests include full source-file hashes and configured limits, and a source +change during analysis prevents publication. POSIX directory metadata is fsynced; Windows +uses atomic replacements and flushed files without POSIX directory-fsync semantics. + +Extraction logs each page before scanning and each analysis stage before starting. +`--analysis-max-chars 200000` controls the per-page text limit (default: 200,000 characters). +The summary records truncation and the number of pages beyond the 2,000-page extended-scan +limit; findings from truncated inputs represent partial coverage. Original saved text is +retained. Ctrl-C exits with status 130 and a recovery hint instead of a traceback. +Unicode domain extraction consumes candidate tokens in linear time to avoid regex hangs. + After the crawl completes, OSINTai runs a deterministic analysis stage over what the crawl collected. It is fast, works fully offline, and is on by default; `--no-analysis` skips it. @@ -368,6 +442,10 @@ Each crawl generates a timestamped directory under `data/runs/` with comprehensi - **`graph_edges.jsonl`** - Graph relationships and connections ### Analysis Results (unless `--no-analysis`) + +These files live in the completed `analysis_*/` bundle selected by +`analysis_latest.json`, or under `reanalysis_*/analysis_*/` for offline recovery. + - **`analysis_report.txt`** - Findings separated by origin, correlations, timeline, hypotheses, leads - **`findings.jsonl`** - Every finding with priority, evidence, sources, confidence, and next step - **`correlations.jsonl`** - Scored candidate entity links with the evidence URLs behind each @@ -705,7 +783,25 @@ pip install black flake8 pytest mypy ## Changelog -### v4.0.0 (2026-08-14) - Current Release +### v4.2.0 (2026-09-10) - Current Release +- Process deadlines for extraction and expensive analysis stages, with child cleanup +- Secret-free content-addressed extraction checkpoints and explicit coverage statistics +- Linear token scans for email, domain, credential, and JWT candidates +- Per-model response quality counts and bounded saved-page retries +- Unicode-correct hunt offsets and URL provenance +- Budgeted correlation, set-backed source tracking, and bounded text caching +- Validated immutable report bundles, source hashes, and durable attempt status +- Offline acceptance tests, saved-crawl benchmarks, and a three-platform CI matrix + +### v4.1.0 (2026-09-10) +- Fix catastrophic backtracking in Unicode domain extraction using a linear token scan +- Add offline `--analyze-only RUN_ID` recovery into a fresh results directory +- Add per-page progress, bounded text reads, coverage statistics, and extraction failure isolation +- Preserve HTML element boundaries to prevent concatenated URL/label artifacts +- Handle interruption cleanly; write complete summary counts and empty result files +- Respect explicit `--flag=value` profile overrides and reject invalid numeric limits + +### v4.0.0 (2026-08-14) - **Evidence-Labelled Analysis**: Findings distinguish observed, derived, model-assisted, and hypothetical statements - **Expanded Deterministic Checks**: Homoglyphs, sensitive infrastructure, secret presence, generated text, temporal gaps, and outliers - **Cross-Source Intelligence**: Entity normalization, evidence-backed candidate correlations, timelines, hypotheses, and pivot leads @@ -755,4 +851,4 @@ Built for the OSINT community with contributions from security researchers, digi --- -*OSINTai v4 - Illuminating the shadows of open source intelligence.* +*OSINTai v4.2.0 - Illuminating the shadows of open source intelligence.* diff --git a/src/osintai/__init__.py b/src/osintai/__init__.py index 674433c..49ff557 100644 --- a/src/osintai/__init__.py +++ b/src/osintai/__init__.py @@ -1,4 +1,4 @@ -__version__ = "4.0.0" +__version__ = "4.2.0" __all__ = [ "cli", diff --git a/src/osintai/checkpoints.py b/src/osintai/checkpoints.py new file mode 100644 index 0000000..59d5328 --- /dev/null +++ b/src/osintai/checkpoints.py @@ -0,0 +1,109 @@ +"""Content-addressed extraction checkpoints containing secret metadata only.""" + +from __future__ import annotations + +import hashlib +import json +from pathlib import Path + +from .entities import extract_extended, API_TOKEN_RE +from .scanners import credentials +from .patterns import _shannon_entropy, decode_jwt_payload +from .storage import read_json, write_json + +# Include implementation hashes so edits cannot silently reuse stale extraction results. +EXTRACTOR_VERSION = hashlib.sha256( + b"".join( + Path(__file__).with_name(name).read_bytes() + for name in ("entities.py", "scanners.py", "patterns.py", "checkpoints.py") + ) +).hexdigest() + + +def extract_checkpoint(text): + extras = extract_extended(text) + coverage = extras.pop("extraction_coverage") + tokens = extras.pop("api_tokens") + jwts = extras.pop("jwts") + pairs = extras.pop("credential_pairs") + # Redaction must scan all matches, including those beyond the findings cap. + secrets = set(API_TOKEN_RE.findall(text)) | set(jwts) | {secret for _, secret in credentials(text)} + # A token can also resemble another indicator. Do not cache that copy either. + extras = { + key: [value for value in values if not any(secret in value for secret in secrets)] + for key, values in extras.items() + } + extras["extraction_coverage"] = coverage + extras["secret_summary"] = { + "token_count": len(tokens), + "high_entropy_count": sum(_shannon_entropy(token) >= 3.5 for token in tokens), + "pair_count": len(pairs), + # Arbitrary JWT claims can themselves contain secrets; retain only validity. + "jwt_decodable": [decode_jwt_payload(token) is not None for token in jwts], + } + return extras + + +def _valid_extras(extras): + list_fields = {"dates", "name_candidates", "addresses", "unicode_domains", "unicode_handles"} + if not isinstance(extras, dict) or set(extras) != list_fields | {"secret_summary", "extraction_coverage"}: + return False + if any( + not isinstance(extras[key], list) or any(not isinstance(v, str) for v in extras[key]) for key in list_fields + ): + return False + summary, coverage = extras["secret_summary"], extras["extraction_coverage"] + if not isinstance(summary, dict) or not isinstance(coverage, dict): + return False + if any( + type(summary.get(key)) is not int or summary[key] < 0 + for key in ("token_count", "high_entropy_count", "pair_count") + ): + return False + if not isinstance(summary.get("jwt_decodable"), list) or any(type(v) is not bool for v in summary["jwt_decodable"]): + return False + return all( + isinstance(counts, dict) + and all(type(counts.get(key)) is int and counts[key] >= 0 for key in ("observed", "retained", "omitted")) + for counts in coverage.values() + ) + + +class ExtractionCache: + def __init__(self, directory): + self.directory = Path(directory) + self.hits = self.misses = self.invalid = 0 + + def key(self, text, limit, truncated): + metadata = { + "text_sha256": hashlib.sha256(text.encode()).hexdigest(), + "extractor_version": EXTRACTOR_VERSION, + "max_text_chars": limit, + "truncated": truncated, + } + key = hashlib.sha256(json.dumps(metadata, sort_keys=True).encode()).hexdigest() + return key, metadata + + def get(self, key, metadata): + path = self.directory / f"{key}.json" + payload = read_json(str(path)) + if isinstance(payload, dict) and payload.get("metadata") == metadata: + extras = payload.get("extras") + if _valid_extras(extras): + checksum = hashlib.sha256(json.dumps(extras, sort_keys=True).encode()).hexdigest() + if checksum == payload.get("checksum"): + self.hits += 1 + return extras + self.invalid += int(path.exists()) + self.misses += 1 + return None + + def put(self, key, metadata, extras): + write_json( + str(self.directory / f"{key}.json"), + { + "metadata": metadata, + "extras": extras, + "checksum": hashlib.sha256(json.dumps(extras, sort_keys=True).encode()).hexdigest(), + }, + ) diff --git a/src/osintai/cli.py b/src/osintai/cli.py index 182e6c1..c747a80 100644 --- a/src/osintai/cli.py +++ b/src/osintai/cli.py @@ -5,6 +5,8 @@ import asyncio import re import time +import tempfile +import math from urllib.parse import urlparse from osintai.storage import safe_mkdir, now_run_id, load_lines, write_json @@ -121,9 +123,53 @@ def _resolve_seed_urls( # Ordered de-duplication prevents redundant work without changing priority. return list(dict.fromkeys(seeds)) +def _analyze_saved_run(args: argparse.Namespace, ap: argparse.ArgumentParser, base_dir: str) -> None: + """Recover offline into a unique directory without touching original outputs.""" + try: + run_id = _validate_run_id(args.analyze_only) + except ValueError as exc: + ap.error(str(exc)) + source = os.path.realpath(os.path.join(base_dir, "data", "runs", run_id)) + runs_root = os.path.realpath(os.path.join(base_dir, "data", "runs")) + if os.path.commonpath([source, runs_root]) != runs_root: + ap.error("saved run must remain inside data/runs") + if not os.path.isfile(os.path.join(source, "urls_crawled.jsonl")): + ap.error(f"No saved crawl found in {source}") + if args.no_analysis or args.deep or args.cross_check or args.experimental_recursive: + ap.error("--analyze-only is offline; incompatible with --no-analysis or model-assisted modes") + destination = tempfile.mkdtemp(prefix="reanalysis_", dir=source) + print(f"OSINTai {__version__}: offline analysis of {run_id}", flush=True) + print(f"RESULTS: {destination}", flush=True) + output = analyze_run( + source, + AnalysisOptions( + use_ollama=False, evaluate=args.evaluate, training_export=args.training_export, + gap_days=args.gap_days, max_leads_per_kind=args.leads_per_kind, + max_text_chars=args.analysis_max_chars, + page_deadline_s=getattr(args, 'analysis_page_timeout', 30.0), + stage_deadline_s=getattr(args, 'analysis_stage_timeout', 120.0), + candidate_pair_budget=getattr(args, 'correlation_pair_budget', 100_000), + text_cache_bytes=getattr(args, 'text_cache_bytes', 16_000_000), + ), + run_id=run_id, output_dir=destination, + log=lambda message: print(message, flush=True), + ) + print(f"DONE. Analysis report: {output.artifacts['analysis_report']}", flush=True) + + + def main(): + try: + _main() + except KeyboardInterrupt: + print("\nInterrupted. Saved crawl files are retained. Recover with " + "--analyze-only RUN_ID.", file=sys.stderr, flush=True) + raise SystemExit(130) + + +def _main(): ap = argparse.ArgumentParser( - description=f"OSINTai v{__version__.split('.')[0]} (async crawling and analysis)" + description=f"OSINTai {__version__} (async crawling and analysis)" ) ap.add_argument("--version", action="version", version=f"OSINTai {__version__}") ap.add_argument( @@ -154,6 +200,17 @@ def main(): ap.add_argument("--hunt-max", type=int, default=50, help="Max lead URLs per page from hunt mode") ap.add_argument("--run-id", default="", help="Optional run id override") + ap.add_argument("--analyze-only", metavar="RUN_ID", + help="Analyze a saved run offline into a new reanalysis directory") + ap.add_argument("--retry-model", metavar="RUN_ID", help="Retry failed/missing model responses over saved pages using local Ollama") + ap.add_argument("--retry-limit", type=int, default=20, help="Maximum saved pages retried (default: 20)") + ap.add_argument("--retry-timeout", type=float, default=60.0, help="Total deadline per model retry in seconds") + ap.add_argument("--analysis-page-timeout", type=float, default=30.0, help="Disposable extraction worker deadline in seconds") + ap.add_argument("--analysis-stage-timeout", type=float, default=120.0, help="Disposable analysis stage deadline in seconds") + ap.add_argument("--correlation-pair-budget", type=int, default=100_000, help="Maximum co-occurrence candidate pairs examined") + ap.add_argument("--text-cache-bytes", type=int, default=16_000_000, help="Maximum page-text LRU cache bytes per process") + ap.add_argument("--analysis-max-chars", type=int, default=200_000, + help="Maximum characters analyzed per saved page (default: 200000)") ap.add_argument( "--profile", @@ -221,7 +278,7 @@ def main(): profile_applied = [] for dest, value in profile_values.items(): flag = "--" + dest.replace("_", "-") - if flag in argv: + if any(arg == flag or arg.startswith(flag + "=") for arg in argv): continue setattr(args, dest, value) profile_applied.append(f"{dest}={value}") @@ -235,6 +292,38 @@ def main(): base_dir = os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "..")) + for flag, value in (("analysis-max-chars", args.analysis_max_chars), + ("retry-limit", args.retry_limit), ("retry-timeout", args.retry_timeout), + ("analysis-page-timeout", args.analysis_page_timeout), + ("analysis-stage-timeout", args.analysis_stage_timeout), + ("correlation-pair-budget", args.correlation_pair_budget), + ("text-cache-bytes", args.text_cache_bytes), + ("concurrency", args.concurrency), ("per-host", args.per_host), + ("max", args.max), ("hunt-max", args.hunt_max), + ("leads-per-kind", args.leads_per_kind)): + if not math.isfinite(value) or value <= 0: + ap.error(f"--{flag} must be finite and greater than zero") + if args.depth < 0: + ap.error("--depth must be nonnegative") + if args.retry_model: + if args.analyze_only or args.no_ollama: + ap.error("--retry-model requires Ollama and cannot be combined with --analyze-only") + from .model_retry import retry_saved + try: + saved_id = _validate_run_id(args.retry_model) + except ValueError as exc: + ap.error(str(exc)) + runs = os.path.realpath(os.path.join(base_dir, "data", "runs")) + source = os.path.realpath(os.path.join(runs, saved_id)) + if os.path.commonpath([runs, source]) != runs or not os.path.isfile(os.path.join(source, "urls_crawled.jsonl")): + ap.error("No saved crawl found inside data/runs") + result = asyncio.run(retry_saved(source, OllamaAPI(), args.model, args.retry_limit, + args.retry_timeout, args.analysis_max_chars, args.prompt_profile)) + print(f"DONE. Saved model retry results: {result}") + return + if args.analyze_only: + return _analyze_saved_run(args, ap, base_dir) + try: seed_urls = _resolve_seed_urls(args.seed, args.seed_file, base_dir) except (OSError, ValueError) as exc: @@ -291,7 +380,8 @@ def main(): model_embed=args.embed_model, hunt_terms=hunt_terms, hunt_max_leads=args.hunt_max, - prompt_profile=args.prompt_profile + prompt_profile=args.prompt_profile, + extraction_timeout_s=args.analysis_page_timeout ) cross_check_models = [m.strip() for m in args.cross_check.split(",") if m.strip()] @@ -312,7 +402,7 @@ def main(): print("") print("=" * 80) - print(f"OSINTai v{__version__.split('.')[0]}") + print(f"OSINTai {__version__}") print(f"RUN DIR: {run_dir}") print(f"SEEDS: {len(seed_urls)} URL(s)") if len(seed_urls) == 1: @@ -363,6 +453,11 @@ def main(): use_ollama=(not args.no_ollama), experimental_recursive=args.experimental_recursive, max_leads_per_kind=args.leads_per_kind, + max_text_chars=args.analysis_max_chars, + page_deadline_s=args.analysis_page_timeout, + stage_deadline_s=args.analysis_stage_timeout, + candidate_pair_budget=args.correlation_pair_budget, + text_cache_bytes=args.text_cache_bytes, ) try: analysis_output = analyze_run( @@ -370,21 +465,9 @@ def main(): options=options, ollama=ollama, run_id=run_id, - log=print, - ) - analysis_report_path = write_analysis_report( - run_dir, - analysis_output, - run_id=run_id, - scope={ - "seeds": len(seed_urls), - "depth": args.depth, - "max urls": args.max, - "same domain only": args.same_domain, - "analysis model": args.model if not args.no_ollama else "none", - "prompt profile": args.prompt_profile, - }, + log=lambda message: print(message, flush=True), ) + analysis_report_path = analysis_output.artifacts["analysis_report"] except Exception as exc: # The crawl is already saved. An analysis failure must not cost the operator it. print(f"[FAIL] analysis stage aborted -> {exc}") @@ -418,6 +501,9 @@ def main(): "experimental_recursive": args.experimental_recursive, "gap_days": args.gap_days, "pages_scored": len(page_scores), + "model_responses": analysis_output.stats.get("model_responses", {}) if analysis_output else {}, + "model_call_counts": crawler.ollama.response_counts if crawler.ollama else {}, + "analysis_bundle": analysis_output.artifacts.get("run_manifest") if analysis_output else None, "analysis_stats": analysis_output.stats if analysis_output else {}, }) diff --git a/src/osintai/correlation.py b/src/osintai/correlation.py index 976c488..ad9c387 100644 --- a/src/osintai/correlation.py +++ b/src/osintai/correlation.py @@ -12,6 +12,7 @@ from __future__ import annotations from collections import Counter, defaultdict +from heapq import nsmallest from dataclasses import dataclass, field from typing import Any, Dict, Iterable, List, Optional, Tuple from urllib.parse import urlparse @@ -82,6 +83,7 @@ def map_domains_to_urls(indicator_rows: Iterable[Dict[str, Any]]) -> Dict[str, A """Group crawled URLs under the domains that reference them, ranked by frequency.""" domain_map: Dict[str, List[str]] = defaultdict(list) frequency: Counter = Counter() + seen = defaultdict(set) for row in indicator_rows or []: source = row.get("url") or "" @@ -92,12 +94,16 @@ def map_domains_to_urls(indicator_rows: Iterable[Dict[str, Any]]) -> Dict[str, A continue if not host: continue - if url not in domain_map[host]: - domain_map[host].append(url) + if url not in seen[host]: + seen[host].add(url) + if len(domain_map[host]) < 200: + domain_map[host].append(url) frequency[host] += 1 source_host = (row.get("domain") or "").lower() - if source_host and source not in domain_map[source_host]: - domain_map[source_host].append(source) + if source_host and source not in seen[source_host]: + seen[source_host].add(source) + if len(domain_map[source_host]) < 200: + domain_map[source_host].append(source) return { "total_domains": len(domain_map), @@ -124,8 +130,8 @@ def _site_wide(index: EntityIndex, page_count: int) -> set: def _co_occurrence_pairs( - index: EntityIndex, page_count: int -) -> Tuple[List[Correlation], int]: + index: EntityIndex, page_count: int, candidate_budget: int = 100_000 +) -> Tuple[List[Correlation], int, dict]: """Identifiers that appeared together on the same page. Site-wide identifiers are excluded from pairing. On a single-site crawl the contact @@ -137,18 +143,27 @@ def _co_occurrence_pairs( linkable = {EMAIL, PHONE, USERNAME, NAME} ubiquitous = _site_wide(index, page_count) + attempted = omitted = oversized = 0 for url, entities in by_source.items(): interesting = [ e for e in entities if e.kind in linkable and e.key not in ubiquitous ] - if len(interesting) < 2 or len(interesting) > MAX_IDENTIFIERS_PER_PAGE: + possible = len(interesting) * (len(interesting) - 1) // 2 + if len(interesting) > MAX_IDENTIFIERS_PER_PAGE: + oversized += possible + continue + remaining = max(0, candidate_budget - attempted) + omitted += max(0, possible - remaining) + if remaining == 0 or possible == 0: continue ordered = sorted(interesting, key=lambda e: (e.kind, e.canonical_value)) for i, left in enumerate(ordered): for right in ordered[i + 1:]: + if attempted >= candidate_budget: + break + attempted += 1 key = (left.key, right.key) - if url not in pair_evidence[key]: - pair_evidence[key].append(url) + pair_evidence[key].append(url) correlations: List[Correlation] = [] for (left_key, right_key), urls in pair_evidence.items(): @@ -165,13 +180,28 @@ def _co_occurrence_pairs( ), score=score, )) - return correlations, len(ubiquitous) + return correlations, len(ubiquitous), { + "candidate_pair_budget": candidate_budget, "candidate_pairs_examined": attempted, + "candidate_pairs_omitted": omitted, "oversized_page_pairs_omitted": oversized, + "partial_coverage": bool(omitted or oversized), + } def _identity_hints(index: EntityIndex) -> List[Correlation]: """Shape-based identity candidates: email local parts, name forms, phone variants.""" correlations: List[Correlation] = [] handles = {e.canonical_value.lstrip("@"): e for e in index.of_kind(USERNAME)} + prefixes = {} + + def evidence(left, right): + small, large = sorted((left, right), key=lambda entity: len(entity.sources)) + shared = nsmallest(25, (url for url in small.sources if url in large._source_set)) + if shared: + return shared, True + for entity in (left, right): + if entity.key not in prefixes: + prefixes[entity.key] = nsmallest(25, entity.sources) + return sorted(set(prefixes[left.key]) | set(prefixes[right.key]))[:25], False for email in index.of_kind(EMAIL): local = email.canonical_value.split("@", 1)[0] @@ -180,12 +210,12 @@ def _identity_hints(index: EntityIndex) -> List[Correlation]: handle = handles.get(local) if handle is None: continue - shared = sorted(set(email.sources) & set(handle.sources)) + urls, shared = evidence(email, handle) correlations.append(Correlation( left_kind=EMAIL, left=email.canonical_value, right_kind=USERNAME, right=handle.value, relation=LOCAL_PART_MATCH, - evidence_urls=(shared or sorted(set(email.sources) | set(handle.sources)))[:25], + evidence_urls=urls, rationale=( f"Email local part {local!r} matches the handle. Common local parts are reused " "widely by unrelated people; treat as a pivot, not an identity." @@ -198,12 +228,12 @@ def _identity_hints(index: EntityIndex) -> List[Correlation]: handle = handles.get(collapsed) if handle is None: continue - shared = sorted(set(name.sources) & set(handle.sources)) + urls, shared = evidence(name, handle) correlations.append(Correlation( left_kind=NAME, left=name.value, right_kind=USERNAME, right=handle.value, relation=NAME_MATCH, - evidence_urls=(shared or sorted(set(name.sources) | set(handle.sources)))[:25], + evidence_urls=urls, rationale="Handle matches the name with separators removed.", score=0.4 if shared else 0.25, )) @@ -240,14 +270,10 @@ def _domain_families(index: EntityIndex) -> List[Correlation]: for base, members in families.items(): if len(members) < 2: continue - sources: List[str] = [] - for member in members: - for url in member.sources: - if url not in sources: - sources.append(url) + sources = list(dict.fromkeys(url for member in members for url in member.sources)) correlations.append(Correlation( left_kind=DOMAIN, left=base, - right_kind=DOMAIN, right=", ".join(sorted(m.canonical_value for m in members)[:8]), + right_kind=DOMAIN, right=", ".join(nsmallest(8, (m.canonical_value for m in members))), relation=SHARES_DOMAIN, evidence_urls=sources[:25], rationale=f"{len(members)} hosts share the registrable domain {base}.", @@ -263,9 +289,12 @@ def _domain_families(index: EntityIndex) -> List[Correlation]: def correlate( - index: EntityIndex, indicator_rows: Iterable[Dict[str, Any]], page_count: int = 0 + index: EntityIndex, indicator_rows: Iterable[Dict[str, Any]], page_count: int = 0, + candidate_budget: int = 100_000 ) -> CheckResult: """Run every correlation method and report the candidates.""" + if candidate_budget < 1: + raise ValueError("candidate_budget must be positive") result = CheckResult(check_name="Cross-Source Correlation") rows = list(indicator_rows or []) @@ -275,13 +304,15 @@ def correlate( correlations: List[Correlation] = [] correlations.extend(_domain_families(index)) correlations.extend(_identity_hints(index)) - co_occurrences, site_wide_count = _co_occurrence_pairs(index, page_count) + co_occurrences, site_wide_count, coverage = _co_occurrence_pairs(index, page_count, candidate_budget) correlations.extend(co_occurrences) correlations.sort(key=lambda c: (-c.score, c.relation, c.left)) truncated = max(0, len(correlations) - MAX_CORRELATION_ROWS) result.rows = [c.to_dict() for c in correlations[:MAX_CORRELATION_ROWS]] result.stats = { + **coverage, + "partial_coverage": coverage["partial_coverage"] or bool(truncated), "correlation_count": len(correlations), "rows_written": len(result.rows), "rows_truncated": truncated, @@ -290,6 +321,9 @@ def correlate( "top_domains": domain_mapping["top_domains"][:10], } + if coverage["partial_coverage"]: + result.notes.append("Candidate pairing coverage is partial; see omitted-pair counts.") + # Only the strongest candidates become findings; the rest stay available in the artifact. promoted = 0 for correlation in correlations: diff --git a/src/osintai/crawler.py b/src/osintai/crawler.py index b57d56e..a73ad06 100644 --- a/src/osintai/crawler.py +++ b/src/osintai/crawler.py @@ -14,6 +14,8 @@ from .extractor import Extractor from .fetcher import AsyncFetcher, FetchRejected from .ollama_api import OllamaAPI +from .model_quality import page_status +from .isolation import async_isolated_call from .dedupe import sha1_text, simhash_64, hamming64 from .analyzer import compute_page_signal from .hunt import hunt_leads @@ -43,6 +45,13 @@ class PageRecord: saved_text_path: str simhash64: int +def _extract_page(extractor, url, html, hunt_terms, hunt_limit): + title, text = extractor.html_to_text(html) + return (title, text, extractor.extract_links(url, html), + extractor.extract_indicators(url, text, html), + hunt_leads(text, hunt_terms, max_leads=hunt_limit), simhash_64(text)) + + class AsyncCrawler: def __init__( self, @@ -61,8 +70,11 @@ def __init__( model_embed: str, hunt_terms: List[str], hunt_max_leads: int, - prompt_profile: str = STANDARD + prompt_profile: str = STANDARD, + extraction_timeout_s: float = 30.0 ): + self.extraction_timeout_s = extraction_timeout_s + self.extraction_sema = asyncio.Semaphore(2) self.prompt_profile = prompt_profile self.seed_urls = seed_urls self.allowed_seed_hosts = {host_of(url) for url in seed_urls if host_of(url)} @@ -190,8 +202,8 @@ def _scoped(self, url: str) -> bool: return False return True - def _dedupe_near(self, text: str) -> Tuple[bool, int]: - sh = simhash_64(text) + def _dedupe_near(self, text: str, sh=None) -> Tuple[bool, int]: + sh = simhash_64(text) if sh is None else sh # if very close to any prior simhash, drop it for prev in self.simhash_seen[-400:]: if hamming64(sh, prev) <= 3: @@ -243,16 +255,33 @@ async def _process_one(self, client: httpx.AsyncClient, url: str, depth: int): return html = resp.text - title, text = self.extractor.html_to_text(html) + rid = sha1(url) + raw_path = os.path.join(self.raw_dir, f"{rid}.html") + # Preserve the response even when a parser defects or exceeds its deadline. + with open(raw_path, "w", encoding="utf-8", errors="ignore") as handle: + handle.write(html) + try: + async with self.extraction_sema: + title, text, out_links, indicators, hunt, sh = await async_isolated_call( + _extract_page, self.extractor, url, html, self.hunt_terms, self.hunt_max_leads, + timeout_s=self.extraction_timeout_s) + except Exception as exc: + append_jsonl(os.path.join(self.run_dir, "extraction_failures.jsonl"), { + "url": url, "saved_raw_path": raw_path, + "status": "timed_out" if isinstance(exc, TimeoutError) else "failed", + "error": str(exc), + }) + self.visited.add(url) + print(f"[FAIL] extraction {url}: {exc}") + return # near-dup check - is_near_dup, sh = self._dedupe_near(text) + is_near_dup, sh = self._dedupe_near(text, sh) if is_near_dup: self.visited.add(url) print(f"[DUP] depth={depth:02d} {url}") return - out_links = self.extractor.extract_links(url, html) for lk in out_links: self._enqueue(lk, depth + 1) @@ -260,8 +289,6 @@ async def _process_one(self, client: httpx.AsyncClient, url: str, depth: int): raw_path = os.path.join(self.raw_dir, f"{rid}.html") text_path = os.path.join(self.text_dir, f"{rid}.txt") - with open(raw_path, "w", encoding="utf-8", errors="ignore") as f: - f.write(html) with open(text_path, "w", encoding="utf-8", errors="ignore") as f: f.write(text) @@ -279,7 +306,6 @@ async def _process_one(self, client: httpx.AsyncClient, url: str, depth: int): ) append_jsonl(self.urls_jsonl, asdict(page)) - indicators = self.extractor.extract_indicators(url, text, html) self.indicators_by_url[url] = indicators append_jsonl(self.indicators_jsonl, indicators) @@ -288,13 +314,21 @@ async def _process_one(self, client: httpx.AsyncClient, url: str, depth: int): prompt = self._analysis_prompt(url, title, text) try: async with self.ollama_sema: - analysis = await self.ollama.async_generate_json(self.model_analyze, prompt, timeout_s=140.0) - except Exception as e: - analysis = {"url": url, "title": title, "error": "ollama_generate_failed", "exception": str(e)} + response = await self.ollama.async_generate_result(self.model_analyze, prompt, timeout_s=140.0) + payload = response["payload"] + model_status = page_status(payload) if response["status"] == "ok" else response["status"] + analysis = dict(payload) if model_status == "ok" else {} + analysis.update(url=url, title=title, _model=self.model_analyze, _model_status=model_status) + except Exception: + analysis = {"url": url, "title": title, "_model": self.model_analyze, "_model_status": "error"} out_json = os.path.join(self.analysis_dir, f"{rid}.analysis.json") - with open(out_json, "w", encoding="utf-8") as f: - json.dump(analysis, f, ensure_ascii=False, indent=2) + write_json(out_json, analysis) + + else: + write_json(os.path.join(self.analysis_dir, f"{rid}.analysis.json"), { + "url": url, "title": title, "_model": self.model_analyze, "_model_status": "skipped", + }) # embeddings (for clustering later) if self.use_ollama and self.ollama and len(text) > 250: @@ -309,7 +343,6 @@ async def _process_one(self, client: httpx.AsyncClient, url: str, depth: int): print(f"[WARN] embedding write failed for {url}: {exc}") # hunt mode on every page (lightweight) - hunt = hunt_leads(text, self.hunt_terms, max_leads=self.hunt_max_leads) if self.hunt_terms else {"hits": [], "lead_urls": []} if hunt.get("hits") or hunt.get("lead_urls"): append_jsonl(self.hunt_jsonl, {"url": url, "depth": depth, **hunt}) for u in hunt.get("lead_urls", [])[: self.hunt_max_leads]: diff --git a/src/osintai/entities.py b/src/osintai/entities.py index 16bc147..4699f89 100644 --- a/src/osintai/entities.py +++ b/src/osintai/entities.py @@ -12,6 +12,8 @@ from __future__ import annotations import re +from heapq import nsmallest +from .scanners import credentials from dataclasses import dataclass, field from typing import Any, Dict, Iterable, List, Optional, Set, Tuple @@ -50,6 +52,7 @@ r"(?i)(?:api[_-]?key|apikey|access[_-]?token|auth[_-]?token|secret|token)" r"[\"'\s:=]{1,5}[\"']?([A-Za-z0-9_\-]{16,64})[\"']?" ) +JWT_TOKEN_RE = re.compile(r"[A-Za-z0-9_.-]+") JWT_RE = re.compile(r"\b(eyJ[A-Za-z0-9_-]{5,}\.[A-Za-z0-9_-]{5,}\.[A-Za-z0-9_-]{5,})\b") # The crawler's own domain and handle patterns are ASCII-only, which means a lookalike @@ -57,17 +60,32 @@ # case homoglyph analysis exists to catch. These patterns accept non-ASCII letters so the # check has something to inspect. Kept here rather than in Extractor so the crawl hot path # and its output schema stay exactly as they were. -UNICODE_DOMAIN_RE = re.compile( - r"(?:[^\W\d_]|[a-zA-Z0-9-])+(?:\.(?:[^\W\d_]|[a-zA-Z0-9-])+)+", re.UNICODE -) +# Consume each token once, then validate it. Requiring a dot in the regex would +# retry long non-domain words at every offset; overlapping letter alternatives +# in the old pattern additionally caused exponential backtracking. +DOMAIN_TOKEN_RE = re.compile(r"(?:[^\W_]|[.\-\u200b\u200c\u200d\ufeff])+", re.UNICODE) + + +def unicode_domains(text: str) -> Iterable[str]: + """Scan in O(n) time, yielding Unicode domain candidates for review.""" + for match in DOMAIN_TOKEN_RE.finditer(text): + value = match.group().strip(".") + if not _has_non_ascii(value) or "." not in value or len(value) > 253: + continue + labels = value.split(".") + if all( + label and len(label) <= 63 and not label.startswith("-") + and not label.endswith("-") for label in labels + ): + yield value + + UNICODE_HANDLE_RE = re.compile(r"(? bool: return any(ord(ch) > 0x7E for ch in value) -CREDENTIAL_PAIR_RE = re.compile( - r"(?m)^[ \t]*([\w.+-]+@[\w.-]+\.[A-Za-z]{2,}|[\w.-]{3,32}):(?!//)(\S{6,64})[ \t]*$" -) + # Left-hand values that make a `word:value` line something other than a credential. Without # these a bare URL on its own line reads as "user https" with a six-character secret. @@ -88,12 +106,12 @@ def detect_type(value: str) -> str: v = (value or "").strip() if not v: return USERNAME - if _EMAIL_SHAPE.match(v): + if len(v) <= 254 and _EMAIL_SHAPE.match(v): return EMAIL digit_count = sum(1 for ch in v if ch.isdigit()) if digit_count >= 7 and _PHONE_SHAPE.match(v): return PHONE - if "@" not in v and _NAME_SHAPE.match(v): + if len(v) <= 128 and "@" not in v and _NAME_SHAPE.match(v): return NAME return USERNAME @@ -167,8 +185,10 @@ class Entity: variants: List[str] = field(default_factory=list) observed_forms: List[str] = field(default_factory=list) count: int = 0 + _source_set: Set[str] = field(default_factory=set, repr=False) def __post_init__(self) -> None: + self._source_set.update(self.sources) if not self.canonical_value: self.canonical_value = canonical(self.value, self.kind) @@ -184,7 +204,8 @@ def observe(self, source: str, raw: str = "") -> None: support a claim about how something was written. """ self.count += 1 - if source and source not in self.sources: + if source and source not in self._source_set: + self._source_set.add(source) self.sources.append(source) raw = (raw or "").strip() if raw and raw not in self.observed_forms and len(self.observed_forms) < 12: @@ -283,39 +304,50 @@ def index_indicators(rows: Iterable[Dict[str, Any]]) -> EntityIndex: return index -def extract_extended(text: str) -> Dict[str, List[str]]: +def extract_extended(text: str) -> Dict[str, Any]: """Indicator classes beyond the crawler's built-in set. Kept separate from Extractor so the crawl hot path and its existing output schema are untouched; the analysis stage calls this over already-saved page text. """ body = text or "" - dates = sorted(set(DATE_RE.findall(body)))[:200] - name_candidates = sorted(set(NAME_CANDIDATE_RE.findall(body)))[:200] - addresses = sorted({m.strip() for m in ADDRESS_RE.findall(body)})[:100] - api_tokens = sorted(set(API_TOKEN_RE.findall(body)))[:100] - jwts = sorted(set(JWT_RE.findall(body)))[:50] - credential_pairs = [ - f"{user}:{secret}" - for user, secret in CREDENTIAL_PAIR_RE.findall(body) - if user.lower() not in _NOT_CREDENTIAL_KEYS - ][:100] - # Only the non-ASCII ones are worth carrying: the ASCII domains and handles are already - # in the crawler's own indicator output. - unicode_domains = sorted({ - m.strip(".") for m in UNICODE_DOMAIN_RE.findall(body) if _has_non_ascii(m) - })[:100] - unicode_handles = sorted({ - "@" + m for m in UNICODE_HANDLE_RE.findall(body) if _has_non_ascii(m) - })[:100] + coverage = {} + + def select(kind, values, limit, unique=True): + if unique: + pool = set(values) + kept = nsmallest(limit, pool) + count = len(pool) + else: + kept, count = [], 0 + for value in values: + count += 1 + if len(kept) < limit: + kept.append(value) + coverage[kind] = {"observed": count, "retained": len(kept), "omitted": count - len(kept)} + return kept + + dates = select("dates", DATE_RE.findall(body), 200) + name_candidates = select("name_candidates", NAME_CANDIDATE_RE.findall(body), 200) + addresses = select("addresses", (m.strip() for m in ADDRESS_RE.findall(body)), 100) + api_tokens = select("api_tokens", API_TOKEN_RE.findall(body), 100) + jwts = select("jwts", (match.group() for match in JWT_TOKEN_RE.finditer(body) + if len(match.group()) <= 16384 and JWT_RE.fullmatch(match.group())), 50) + credential_pairs = select("credential_pairs", ( + f"{user}:{secret}" for user, secret in credentials(body) + if user.lower() not in _NOT_CREDENTIAL_KEYS), 100, unique=False) + domain_values = select("unicode_domains", unicode_domains(body), 100) + unicode_handles = select("unicode_handles", ( + "@" + m for m in UNICODE_HANDLE_RE.findall(body) if _has_non_ascii(m)), 100) return { + "extraction_coverage": coverage, "dates": dates, "name_candidates": name_candidates, "addresses": addresses, "api_tokens": api_tokens, "jwts": jwts, "credential_pairs": credential_pairs, - "unicode_domains": unicode_domains, + "unicode_domains": domain_values, "unicode_handles": unicode_handles, } diff --git a/src/osintai/extractor.py b/src/osintai/extractor.py index b1945e1..cb97f51 100644 --- a/src/osintai/extractor.py +++ b/src/osintai/extractor.py @@ -2,6 +2,7 @@ from bs4 import BeautifulSoup from urllib.parse import urlparse from .normalize import absolutize +from .scanners import emails as scan_emails, domains as scan_domains def _validated_http_url(value: str) -> tuple[str, str] | None: @@ -19,10 +20,8 @@ def _validated_http_url(value: str) -> tuple[str, str] | None: class Extractor: - EMAIL_RE = re.compile(r"\b[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2,}\b") PHONE_RE = re.compile(r"\b\d{3}[-.]?\d{3}[-.]?\d{4}\b") URL_RE = re.compile(r"https?://[^\s<>\"']+") - DOMAIN_RE = re.compile(r"\b(?:[a-zA-Z0-9-]+\.)+[a-zA-Z]{2,}\b") IP_RE = re.compile(r"\b(?:(?:25[0-5]|2[0-4]\d|1?\d?\d)\.){3}(?:25[0-5]|2[0-4]\d|1?\d?\d)\b") BTC_RE = re.compile(r"\b(?:bc1[ac-hj-np-z02-9]{11,71}|[13][a-km-zA-HJ-NP-Z1-9]{25,34})\b") ETH_RE = re.compile(r"\b0x[a-fA-F0-9]{40}\b") @@ -42,7 +41,8 @@ def html_to_text(self, html: str) -> tuple[str, str]: title = soup.title.get_text().strip() # Get text content - text = soup.get_text() + # Preserve element boundaries so a URL cannot absorb the next label. + text = soup.get_text(separator="\n") lines = [line.strip() for line in text.splitlines() if line.strip()] return title, "\n".join(lines) @@ -59,13 +59,18 @@ def extract_links(self, base_url: str, html: str) -> list[str]: def extract_indicators(self, url: str, text: str, html: str) -> dict: """Extract various indicators from content.""" combined = "\n".join([text or "", html or ""]) - emails = sorted(set(self.EMAIL_RE.findall(combined)))[:200] + emails = sorted(set(scan_emails(combined)))[:200] phones = sorted(set(self.PHONE_RE.findall(combined)))[:100] parsed_urls = { parsed for match in self.URL_RE.findall(combined) if (parsed := _validated_http_url(match)) is not None } + prose_urls = {parsed[0] for match in self.URL_RE.findall(text or "") + if (parsed := _validated_http_url(match)) is not None} + attribute_urls = set(self.extract_links(url, html or "")) + parsed_urls.update(parsed for value in attribute_urls + if (parsed := _validated_http_url(value)) is not None) urls = sorted(url for url, _ in parsed_urls)[:100] btc_addresses = sorted(set(self.BTC_RE.findall(combined)))[:100] eth_addresses = sorted({addr.lower() for addr in self.ETH_RE.findall(combined)})[:100] @@ -75,7 +80,7 @@ def extract_indicators(self, url: str, text: str, html: str) -> dict: source = _validated_http_url(url) domain = source[1] if source else "" domains = {domain} if domain else set() - domains.update(self.DOMAIN_RE.findall(combined)) + domains.update(scan_domains(combined)) domains.update(hostname for _, hostname in parsed_urls) # Drop email-only host fragments that appear solely because of address parsing noise. @@ -89,6 +94,9 @@ def extract_indicators(self, url: str, text: str, html: str) -> dict: "emails": emails, "phones": phones, "urls": urls, + "url_provenance": {value: [origin for origin, values in ( + ("html_attribute", attribute_urls), ("page_prose", prose_urls) + ) if value in values] or ["html_source"] for value in urls}, "ip_addresses": ip_addresses, "btc_addresses": btc_addresses, "eth_addresses": eth_addresses, diff --git a/src/osintai/hunt.py b/src/osintai/hunt.py index d059ccd..b32badb 100644 --- a/src/osintai/hunt.py +++ b/src/osintai/hunt.py @@ -1,4 +1,5 @@ import re +from bisect import bisect_left from typing import List, Dict, Any URL_RE = re.compile(r"https?://[^\s\"\'<>]+", re.IGNORECASE) @@ -10,36 +11,50 @@ def hunt_leads(text: str, hunt_terms: List[str], max_leads: int = 50) -> Dict[st hits = [] lead_urls = set() + lowered = text.lower() + offsets = (range(len(text)) if len(lowered) == len(text) else + [i for i, char in enumerate(text) for _ in char.lower()]) + url_matches = list(URL_RE.finditer(text)) + url_starts = [match.start() for match in url_matches] # Search for each hunt term for term in hunt_terms: + if not term or len(hits) >= 500: + continue term_lower = term.lower() pos = 0 - while pos < len(text): - idx = text.lower().find(term_lower, pos) + while pos < len(lowered): + idx = lowered.find(term_lower, pos) if idx == -1: break # Extract snippet around the hit - start = max(0, idx - 100) - end = min(len(text), idx + len(term) + 100) + original_start = offsets[idx] + original_end = offsets[idx + len(term_lower) - 1] + 1 + start = max(0, original_start - 100) + end = min(len(text), original_end + 100) snippet = text[start:end] hits.append({ "term": term, - "position": idx, + "position": original_start, + "end_position": original_end, "snippet": snippet.replace('\n', ' ').strip() }) - # Extract URLs from snippet - for match in URL_RE.findall(snippet): - lead_urls.add(match.strip()) + # Match complete URLs in the original text; a snippet edge is not a URL end. + left = bisect_left(url_starts, start) + right = bisect_left(url_starts, end) + for match in url_matches[left:right]: + if match.end() <= end: + lead_urls.add(match.group().rstrip(".,;:!?)}")) - pos = idx + len(term) + pos = idx + len(term_lower) if len(hits) >= 500: break return { "hits": hits[:500], - "lead_urls": list(lead_urls)[:max_leads] + "lead_urls": sorted(lead_urls)[:max(0, max_leads)], + "url_provenance": {url: "page_prose" for url in sorted(lead_urls)[:max(0, max_leads)]} } diff --git a/src/osintai/isolation.py b/src/osintai/isolation.py new file mode 100644 index 0000000..879aada --- /dev/null +++ b/src/osintai/isolation.py @@ -0,0 +1,92 @@ +"""Disposable spawn workers with wall-clock deadlines and deterministic cleanup.""" + +from __future__ import annotations + +import math +import multiprocessing +import queue +import threading +import time + + +def _worker(connection, function, args, kwargs): + try: + connection.send((True, function(*args, **kwargs))) + except BaseException as exc: + connection.send((False, f"{type(exc).__name__}: {exc}")) + finally: + connection.close() + + +def isolated_call(function, *args, timeout_s=30.0, cancel_event=None, **kwargs): + """Bound one job, including startup/IPC. Arguments must be spawn-picklable. + + A receiver thread drains large results while the parent enforces the deadline; + joining before receiving would deadlock on a full pipe. + """ + if not math.isfinite(timeout_s) or timeout_s <= 0: + raise ValueError("worker deadline must be finite and positive") + started = time.monotonic() + context = multiprocessing.get_context("spawn") + receiver, sender = context.Pipe(duplex=False) + process = context.Process(target=_worker, args=(sender, function, args, kwargs), daemon=True) + messages = queue.Queue(maxsize=1) + + def receive(): + try: + messages.put(receiver.recv()) + except (EOFError, OSError): + messages.put((False, "worker exited without a result")) + + reader = None + try: + process.start() + sender.close() + reader = threading.Thread(target=receive, daemon=True) + reader.start() + while True: + if cancel_event is not None and cancel_event.is_set(): + raise RuntimeError("analysis job cancelled") + remaining = timeout_s - (time.monotonic() - started) + if remaining <= 0: + raise TimeoutError(f"analysis job exceeded {timeout_s:g}s deadline") + try: + ok, result = messages.get(timeout=min(0.1, remaining)) + break + except queue.Empty: + continue + if not ok: + raise RuntimeError(result) + return result + finally: + sender.close() + if process.pid is not None: + if process.is_alive(): + process.terminate() + process.join(timeout=1) + if process.is_alive(): + process.kill() + process.join() + process.close() + if reader is not None: + reader.join(timeout=1) + receiver.close() + + +async def async_isolated_call(function, *args, timeout_s=30.0): + """Keep the event loop responsive and signal worker cleanup on cancellation.""" + import asyncio + + cancel = threading.Event() + task = asyncio.create_task( + asyncio.to_thread(isolated_call, function, *args, timeout_s=timeout_s, cancel_event=cancel) + ) + try: + return await asyncio.shield(task) + except asyncio.CancelledError: + cancel.set() + try: + await task + except RuntimeError: + pass + raise diff --git a/src/osintai/model_quality.py b/src/osintai/model_quality.py new file mode 100644 index 0000000..50c5731 --- /dev/null +++ b/src/osintai/model_quality.py @@ -0,0 +1,41 @@ +"""Model result classification without retaining invalid response bodies.""" + +from collections import Counter, defaultdict + +LIST_FIELDS = ("key_entities", "key_locations", "key_dates", "keywords", "risk_flags", "actionable_leads") +STATUSES = ("ok", "empty", "invalid", "missing", "timed_out", "error", "skipped") + + +def page_status(payload): + if payload is None or payload == {}: + return "empty" + if not isinstance(payload, dict): + return "invalid" + status = payload.get("_model_status") + if status in STATUSES and status != "ok": + return status + if payload.get("error"): + return "error" + if not isinstance(payload.get("summary"), str) or not payload["summary"].strip(): + return "invalid" + if any(not isinstance(payload.get(key), list) for key in LIST_FIELDS): + return "invalid" + if any(not isinstance(item, str) for key in LIST_FIELDS[:-1] for item in payload[key]): + return "invalid" + if any(not isinstance(item, (str, dict)) for item in payload["actionable_leads"]): + return "invalid" + return "ok" + + +def summarize(rows): + by_model = defaultdict(Counter) + counts = Counter() + for row in rows: + status = row["_model_status"] + counts[status] += 1 + by_model[row.get("_model") or "unknown"][status] += 1 + return { + "counts": {status: counts[status] for status in STATUSES}, + "models": {model: dict(count) for model, count in sorted(by_model.items())}, + "pages": len(rows), + } diff --git a/src/osintai/model_retry.py b/src/osintai/model_retry.py new file mode 100644 index 0000000..4598cc9 --- /dev/null +++ b/src/osintai/model_retry.py @@ -0,0 +1,80 @@ +"""Explicit bounded model retries over saved text; original analyses stay intact.""" + +import asyncio +import os +import time +import math +import uuid +from pathlib import Path + +from .model_quality import page_status, summarize +from .pipeline import RunArtifacts +from .prompts import page_prompt +from .publication import retry_source_hashes +from .storage import write_json, sha1 + + +async def retry_saved( + source, ollama, model, limit=20, timeout_s=60.0, max_text_chars=200_000, prompt_profile="standard" +): + if limit < 1 or not math.isfinite(timeout_s) or timeout_s <= 0 or max_text_chars < 1: + raise ValueError("retry limits must be positive") + reader = RunArtifacts(str(source), max_text_chars) + records = reader.model_records() + candidates = [row for row in records if row["_model_status"] != "ok"] + root = Path(source) + name = f"model_retry_{uuid.uuid4().hex}" + staging, final = root / f".{name}.incomplete", root / name + (staging / "analysis").mkdir(parents=True) + status_path = root / f"{name}.status.json" + status = {"status": "running", "model": model, "limit": limit, "timeout_s": timeout_s, "started_at": time.time()} + write_json(str(status_path), status) + outcomes = [] + try: + hashes = retry_source_hashes(str(source)) + # Preserve previous retry results when advancing the pointer after a small batch. + for row in records: + write_json(str(staging / "analysis" / f"{sha1(row['url'])}.analysis.json"), row) + for row in candidates[:limit]: + url = row["url"] + text = reader.text_for(url) + result = {"status": "missing", "payload": None} + if text: + try: + result = await asyncio.wait_for( + ollama.async_generate_result( + model, page_prompt(url, row.get("title", ""), text, prompt_profile), timeout_s=timeout_s + ), + timeout=timeout_s, + ) + except (asyncio.TimeoutError, TimeoutError): + result = {"status": "timed_out", "payload": None} + except Exception: + result = {"status": "error", "payload": None} + status_value = page_status(result["payload"]) if result["status"] == "ok" else result["status"] + payload = dict(result["payload"]) if status_value == "ok" else {} + payload.update(url=url, _model=model, _model_status=status_value) + write_json(str(staging / "analysis" / f"{sha1(url)}.analysis.json"), payload) + outcomes.append(payload) + if retry_source_hashes(str(source)) != hashes: + raise RuntimeError("Source artifacts changed during retry") + status.update( + status="completed", + finished_at=time.time(), + attempted=len(outcomes), + omitted=max(0, len(candidates) - limit), + model_responses=summarize(outcomes), + source_hashes=hashes, + text_coverage=reader.coverage(), + max_text_chars=max_text_chars, + prompt_profile=prompt_profile, + ) + write_json(str(staging / "run_manifest.json"), status) + os.replace(staging, final) + write_json(str(status_path), status) + write_json(str(root / "model_retry_latest.json"), {"directory": name}) + return str(final) + except BaseException as exc: + status.update(status="failed", finished_at=time.time(), error=type(exc).__name__) + write_json(str(status_path), status) + raise diff --git a/src/osintai/ollama_api.py b/src/osintai/ollama_api.py index df61695..b92b597 100644 --- a/src/osintai/ollama_api.py +++ b/src/osintai/ollama_api.py @@ -31,6 +31,7 @@ class OllamaAPI: def __init__(self, base_url: str = "http://localhost:11434"): self.base_url = _validate_local_base_url(base_url) + self.response_counts = {} def generate_json(self, model: str, prompt: str, timeout_s: float = 60.0) -> Optional[Dict[str, Any]]: """Generate JSON response from model.""" @@ -38,6 +39,16 @@ def generate_json(self, model: str, prompt: str, timeout_s: float = 60.0) -> Opt async def async_generate_json(self, model: str, prompt: str, timeout_s: float = 60.0) -> Optional[Dict[str, Any]]: """Generate JSON without blocking the crawler event loop.""" + result = await self.async_generate_result(model, prompt, timeout_s) + return result["payload"] + + async def async_generate_result(self, model: str, prompt: str, timeout_s: float = 60.0): + """Return payload and a distinct outcome for every attempted model call.""" + def outcome(status, payload=None): + counts = self.response_counts.setdefault(model, {}) + counts[status] = counts.get(status, 0) + 1 + return {"status": status, "model": model, "payload": payload} + url = f"{self.base_url}/api/generate" payload = { "model": model, @@ -55,10 +66,23 @@ async def async_generate_json(self, model: str, prompt: str, timeout_s: float = r = await client.post(url, json=payload) r.raise_for_status() data = r.json() - response_text = data.get("response", "") - return self._extract_json(response_text) + if not isinstance(data, dict) or "response" not in data: + return outcome("missing") + response_text = data["response"] + if not isinstance(response_text, str): + return outcome("invalid") + if not response_text.strip(): + return outcome("empty") + parsed = self._extract_json(response_text) + if parsed == {}: + return outcome("empty") + return outcome("ok", parsed) if isinstance(parsed, dict) else outcome("invalid") + except httpx.TimeoutException: + return outcome("timed_out") + except (ValueError, TypeError): + return outcome("invalid") except Exception: - return None + return outcome("error") def embed(self, model: str, input_text: str, timeout_s: float = 30.0) -> Optional[List[float]]: """Generate embeddings for text.""" diff --git a/src/osintai/patterns.py b/src/osintai/patterns.py index 96ea45b..7b5255e 100644 --- a/src/osintai/patterns.py +++ b/src/osintai/patterns.py @@ -285,7 +285,7 @@ def check_sensitive_infrastructure(index) -> CheckResult: return result -def check_secret_exposure(page_extras: Dict[str, Dict[str, List[str]]]) -> CheckResult: +def check_secret_exposure(page_extras: Dict[str, Dict[str, Any]]) -> CheckResult: """Credential-shaped and secret-shaped material on crawled pages. Deliberately reports presence and location only. The matched value is not copied into @@ -298,14 +298,17 @@ def check_secret_exposure(page_extras: Dict[str, Dict[str, List[str]]]) -> Check jwts = extras.get("jwts") or [] pairs = extras.get("credential_pairs") or [] - if tokens: - high_entropy = [t for t in tokens if _shannon_entropy(t) >= 3.5] + summary = extras.get("secret_summary") or {} + token_count = summary.get("token_count", len(tokens)) + high_entropy_count = summary.get("high_entropy_count", sum(_shannon_entropy(t) >= 3.5 for t in tokens)) + pair_count = summary.get("pair_count", len(pairs)) + if token_count: result.findings.append(Finding( check="Potential Secret Exposure", item=url, reason=( - f"{len(tokens)} key/token-shaped value(s) found in page content, " - f"{len(high_entropy)} of them high-entropy." + f"{token_count} key/token-shaped value(s) found in page content, " + f"{high_entropy_count} of them high-entropy." ), next_step=( "Open the page and confirm whether these are live credentials, example " @@ -313,11 +316,11 @@ def check_secret_exposure(page_extras: Dict[str, Dict[str, List[str]]]) -> Check "notify the owner rather than using them." ), origin=DERIVED, - priority=HIGH if high_entropy else MEDIUM, + priority=HIGH if high_entropy_count else MEDIUM, sources=[url], evidence={ - "token_count": len(tokens), - "high_entropy_count": len(high_entropy), + "token_count": token_count, + "high_entropy_count": high_entropy_count, "value_recorded": False, }, method="token_pattern_entropy", @@ -325,8 +328,10 @@ def check_secret_exposure(page_extras: Dict[str, Dict[str, List[str]]]) -> Check 0.6, "pattern match; placeholders and examples also match")], )) - for token in jwts: - claims = decode_jwt_payload(token) + jwt_states = summary.get("jwt_decodable", []) + for token in (jwts or jwt_states): + claims = decode_jwt_payload(token) if isinstance(token, str) else None + decodable = bool(claims) if isinstance(token, str) else token claim_keys = sorted(claims.keys()) if claims else [] result.findings.append(Finding( check="JWT Present", @@ -340,7 +345,7 @@ def check_secret_exposure(page_extras: Dict[str, Dict[str, List[str]]]) -> Check "identify the issuing system. Do not replay the token." ), origin=DERIVED, - priority=HIGH if claims else MEDIUM, + priority=HIGH if decodable else MEDIUM, sources=[url], evidence={ "claim_keys": claim_keys, @@ -351,15 +356,15 @@ def check_secret_exposure(page_extras: Dict[str, Dict[str, List[str]]]) -> Check }, method="jwt_structural_decode", confidence=[deterministic_confidence( - 0.95 if claims else 0.5, - "three-segment structure decoded" if claims else "structure matched, payload unreadable")], + 0.95 if decodable else 0.5, + "three-segment structure decoded" if decodable else "structure matched, payload unreadable")], )) - if pairs: + if pair_count: result.findings.append(Finding( check="Credential-Shaped Content", item=url, - reason=f"{len(pairs)} line(s) matching an identifier:secret layout.", + reason=f"{pair_count} line(s) matching an identifier:secret layout.", next_step=( "Determine whether the page is a leak, a configuration sample, or unrelated " "colon-delimited data. Preserve the page capture before it is removed." @@ -367,7 +372,7 @@ def check_secret_exposure(page_extras: Dict[str, Dict[str, List[str]]]) -> Check origin=DERIVED, priority=HIGH, sources=[url], - evidence={"pair_count": len(pairs), "values_recorded": False}, + evidence={"pair_count": pair_count, "values_recorded": False}, method="credential_pattern", confidence=[deterministic_confidence( 0.5, "layout match only; colon-delimited data is common")], diff --git a/src/osintai/pipeline.py b/src/osintai/pipeline.py index 0bd823e..7145a02 100644 --- a/src/osintai/pipeline.py +++ b/src/osintai/pipeline.py @@ -1,8 +1,8 @@ """The OSINTai analysis lifecycle. -Runs after the crawl, over what the crawl already saved. The crawl itself is untouched: the -per-page fetch, extract, analyze, score path is exactly what it always was, and this stage -reads its output rather than changing it. +Runs after the crawl, over saved artifacts. Disposable workers contain extraction and +expensive stages; content-addressed checkpoints support recovery. Completed reports are +published as immutable bundles without replacing source captures. RUN ARTIFACTS | @@ -26,10 +26,14 @@ from __future__ import annotations import json +import sys +import math +from collections import OrderedDict +from functools import partial import os import time from dataclasses import dataclass, field -from typing import Any, Callable, Dict, List, Optional, Tuple +from typing import Any, Callable, Dict, Iterable, List, Optional, Set from . import correlation as correlation_module from . import evaluation as evaluation_module @@ -41,7 +45,6 @@ from . import training_export as training_export_module from .entities import ( EntityIndex, - extract_extended, index_indicators, ) from .prompts import STANDARD, deep_analysis_prompt, page_prompt @@ -52,7 +55,10 @@ Lead, sort_findings, ) -from .storage import append_jsonl, safe_mkdir, sha1, write_json +from .storage import safe_mkdir, sha1, write_json, read_json +from .model_quality import page_status, summarize +from .isolation import isolated_call +from .checkpoints import ExtractionCache, extract_checkpoint # Cap on page texts read back for extended extraction. A very large crawl should not turn # the analysis stage into a second crawl-length operation. @@ -73,6 +79,11 @@ class AnalysisOptions: use_ollama: bool = True experimental_recursive: int = 0 max_leads_per_kind: int = 10 + max_text_chars: int = 200_000 + page_deadline_s: float = 30.0 + stage_deadline_s: float = 120.0 + candidate_pair_budget: int = 100_000 + text_cache_bytes: int = 16_000_000 @dataclass @@ -96,7 +107,9 @@ def read_jsonl(path: str) -> List[Dict[str, Any]]: if not line: continue try: - rows.append(json.loads(line)) + row = json.loads(line) + if isinstance(row, dict): + rows.append(row) except ValueError: continue return rows @@ -105,11 +118,17 @@ def read_jsonl(path: str) -> List[Dict[str, Any]]: class RunArtifacts: """Reader for one run directory. Loads lazily and caches page text.""" - def __init__(self, run_dir: str): + def __init__(self, run_dir: str, max_text_chars: int = 200_000, text_cache_bytes: int = 16_000_000): self.run_dir = run_dir self.text_dir = os.path.join(run_dir, "pages_text") self.analysis_dir = os.path.join(run_dir, "analysis") - self._text_cache: Dict[str, str] = {} + self._text_cache = OrderedDict() + self.text_cache_bytes = text_cache_bytes + self.cache_bytes = self.cache_peak_bytes = self.cache_evictions = 0 + self.missing_urls = set() + self.unreadable_urls = set() + self.max_text_chars = max_text_chars + self.truncated_urls: Set[str] = set() def indicators(self) -> List[Dict[str, Any]]: return read_jsonl(os.path.join(self.run_dir, "indicators.jsonl")) @@ -125,48 +144,120 @@ def text_for(self, url: str) -> str: if not url: return "" if url in self._text_cache: + self._text_cache.move_to_end(url) return self._text_cache[url] path = os.path.join(self.text_dir, f"{sha1(url)}.txt") text = "" if os.path.exists(path): try: with open(path, "r", encoding="utf-8", errors="ignore") as handle: - text = handle.read() + text = handle.read(self.max_text_chars + 1) + if len(text) > self.max_text_chars: + self.truncated_urls.add(url) + text = text[: self.max_text_chars] except OSError: - text = "" - self._text_cache[url] = text + self.unreadable_urls.add(url) + else: + self.missing_urls.add(url) + size = sys.getsizeof(text) + sys.getsizeof(url) + 128 + if size <= self.text_cache_bytes: + while self._text_cache and self.cache_bytes + size > self.text_cache_bytes: + old_url, old_text = self._text_cache.popitem(last=False) + self.cache_bytes -= sys.getsizeof(old_text) + sys.getsizeof(old_url) + 128 + self.cache_evictions += 1 + self._text_cache[url] = text + self.cache_bytes += size + self.cache_peak_bytes = max(self.cache_peak_bytes, self.cache_bytes) return text - def analyses(self) -> List[Dict[str, Any]]: - if not os.path.isdir(self.analysis_dir): - return [] - rows: List[Dict[str, Any]] = [] - for name in sorted(os.listdir(self.analysis_dir)): - if not name.endswith(".analysis.json"): - continue - try: - with open(os.path.join(self.analysis_dir, name), "r", encoding="utf-8") as handle: - payload = json.load(handle) - except (OSError, ValueError): - continue - if isinstance(payload, dict): - rows.append(payload) + def coverage(self): + return { + "truncated": sorted(self.truncated_urls), + "missing": sorted(self.missing_urls), + "unreadable": sorted(self.unreadable_urls), + "cache_peak_bytes": self.cache_peak_bytes, + "cache_evictions": self.cache_evictions, + } + + def model_records(self): + from pathlib import Path + + from .publication import retry_source_hashes + + manifest = read_json(os.path.join(self.run_dir, "run_manifest.json")) + default_model = manifest.get("analysis_model", "unknown") if isinstance(manifest, dict) else "unknown" + latest = read_json(os.path.join(self.run_dir, "model_retry_latest.json")) + retry_name = latest.get("directory", "") if isinstance(latest, dict) else "" + run_root = Path(self.run_dir).resolve() + retry_dir = None + if retry_name.startswith("model_retry_") and os.path.basename(retry_name) == retry_name: + candidate = (run_root / retry_name).resolve() + if candidate.parent == run_root and candidate.is_dir(): + analysis_dir = (candidate / "analysis").resolve() + retry_manifest_path = (candidate / "run_manifest.json").resolve() + if ( + analysis_dir.is_relative_to(candidate) + and analysis_dir.is_dir() + and retry_manifest_path.parent == candidate + and retry_manifest_path.is_file() + ): + retry_manifest = read_json(str(retry_manifest_path)) + if retry_manifest.get("status") == "completed" and retry_manifest.get( + "source_hashes" + ) == retry_source_hashes(str(run_root)): + retry_dir = analysis_dir + original_analysis_dir = Path(self.analysis_dir).resolve() + original_analysis_allowed = original_analysis_dir.is_relative_to(run_root) + allowed_parents = {original_analysis_dir} + if retry_dir: + allowed_parents.add(retry_dir) + rows = [] + for url in dict.fromkeys(record.get("url") for record in self.page_records() if record.get("url")): + name = f"{sha1(url)}.analysis.json" + path = (original_analysis_dir / name).resolve() if original_analysis_allowed else None + retry_path = (retry_dir / name).resolve() if retry_dir else None + if retry_path and retry_path.parent == retry_dir and retry_path.is_file(): + path = retry_path + if path is None or path.parent not in allowed_parents or not path.is_file(): + payload, status = {}, "missing" + else: + try: + with path.open(encoding="utf-8") as handle: + payload = json.load(handle) + status = page_status(payload) + except (OSError, ValueError): + payload, status = {}, "invalid" + model = payload.get("_model") if isinstance(payload, dict) else None + model = model if isinstance(model, str) and model else default_model + row = dict(payload) if status == "ok" else {} + row.update(url=url, _model_status=status, _model=model if isinstance(model, str) else "unknown") + rows.append(row) return rows + def analyses(self) -> List[Dict[str, Any]]: + return [row for row in self.model_records() if row["_model_status"] == "ok"] + def _run_stage( - name: str, fn: Callable[[], CheckResult], output: AnalysisOutput, log: Callable[[str], None] + name: str, + fn: Callable[[], CheckResult], + output: AnalysisOutput, + log: Callable[[str], None], + timeout_s: float = 120.0, ) -> Optional[CheckResult]: """Execute one stage under isolation. A stage failure is reported, never propagated.""" started = time.time() + log(f"[RUN] {name}") try: - result = fn() + result = isolated_call(fn, timeout_s=timeout_s) except Exception as exc: message = f"{name}: stage failed: {exc}" output.errors.append(message) log(f"[FAIL] analysis stage {name} -> {exc}") failed = CheckResult(check_name=name) failed.errors.append(str(exc)) + failed.stats["status"] = "timed_out" if isinstance(exc, TimeoutError) else "failed" + failed.stats["partial_coverage"] = True failed.notes.append("This stage failed. Remaining stages continued.") output.results.append(failed) return None @@ -180,46 +271,78 @@ def _run_stage( return result -def analyze_run( +def _analyze_run( run_dir: str, options: AnalysisOptions, ollama=None, run_id: str = "", log: Optional[Callable[[str], None]] = None, + output_dir: Optional[str] = None, ) -> AnalysisOutput: """Run the analysis lifecycle over a completed crawl.""" log = log or (lambda message: None) output = AnalysisOutput() - artifacts = RunArtifacts(run_dir) + if options.max_text_chars < 1: + raise ValueError("max_text_chars must be positive") + artifacts = RunArtifacts(run_dir, options.max_text_chars, options.text_cache_bytes) + cache = ExtractionCache(os.path.join(run_dir, ".extraction_cache")) + run_dir = output_dir or run_dir + safe_mkdir(run_dir) indicator_rows = artifacts.indicators() page_scores = artifacts.page_scores() page_records = artifacts.page_records() - analyses = artifacts.analyses() + model_records = artifacts.model_records() + output.stats["model_responses"] = summarize(model_records) + analyses = [row for row in model_records if row["_model_status"] == "ok"] analyses_by_url = {a.get("url", ""): a for a in analyses if isinstance(a, dict)} if not indicator_rows and not page_records: + output.stats["partial_coverage"] = True + output.stats["extraction_failures"] = read_jsonl(os.path.join(artifacts.run_dir, "extraction_failures.jsonl")) output.errors.append("No crawl artifacts found; analysis skipped.") log("[SKIP] analysis: no crawl artifacts to analyze") return output # Entity index over the indicators the crawl already wrote. - index: EntityIndex = index_indicators(indicator_rows) + try: + index: EntityIndex = isolated_call(index_indicators, indicator_rows, timeout_s=options.stage_deadline_s) + except Exception as exc: + output.errors.append(f"entity indexing failed: {exc}") + index = EntityIndex() # Extended indicator classes over saved page text. Kept out of the crawl hot path so # the crawler's own output schema is unchanged. - page_extras: Dict[str, Dict[str, List[str]]] = {} + page_extras: Dict[str, Dict[str, Any]] = {} content_dates: Dict[str, List[str]] = {} scanned = 0 - for record in page_records[:MAX_TEXT_PAGES]: + extraction_failures = read_jsonl(os.path.join(artifacts.run_dir, "extraction_failures.jsonl")) + total = min(len(page_records), MAX_TEXT_PAGES) + for position, record in enumerate(page_records[:MAX_TEXT_PAGES], 1): url = record.get("url") if not url: continue text = artifacts.text_for(url) if not text: continue + log(f"[SCAN] {position}/{total} {url} ({len(text)} chars)") + try: + key, metadata = cache.key(text, options.max_text_chars, url in artifacts.truncated_urls) + extras = cache.get(key, metadata) + if extras is None: + extras = isolated_call(extract_checkpoint, text, timeout_s=options.page_deadline_s) + try: + cache.put(key, metadata, extras) + except OSError as exc: + output.errors.append(f"checkpoint write failed for {url}: {exc}; extracted results retained") + except Exception as exc: + extraction_failures.append( + {"url": url, "status": "timed_out" if isinstance(exc, TimeoutError) else "failed", "error": str(exc)} + ) + output.errors.append(f"extended extraction failed for {url}: {exc}") + log(f"[FAIL] extraction {url}: {exc}") + continue scanned += 1 - extras = extract_extended(text) page_extras[url] = extras if extras["dates"]: content_dates[url] = extras["dates"] @@ -241,114 +364,104 @@ def analyze_run( output.stats["pages_text_scanned"] = scanned log(f"[OK] entity index: {len(index)} entities from {len(indicator_rows)} indicator record(s)") - # Deterministic analysis. Cheap, offline, reproducible, always runs. - _run_stage("homoglyph analysis", lambda: patterns.check_homoglyphs(index), output, log) - _run_stage( + def run_stage(name, fn, output, log): + return _run_stage(name, fn, output, log, timeout_s=options.stage_deadline_s) + + # Independent CPU stages run in disposable processes. + run_stage("homoglyph analysis", partial(patterns.check_homoglyphs, index), output, log) + run_stage( "sensitive infrastructure", - lambda: patterns.check_sensitive_infrastructure(index), output, log, + partial(patterns.check_sensitive_infrastructure, index), + output, + log, ) - _run_stage( + run_stage( "secret exposure", - lambda: patterns.check_secret_exposure(page_extras), output, log, + partial(patterns.check_secret_exposure, page_extras), + output, + log, ) - _run_stage( + run_stage( "generated-content fingerprint", - lambda: patterns.check_generated_text(page_records, artifacts.text_for), output, log, + partial(_text_stage, "generated", artifacts.run_dir, options, page_records), + output, + log, ) - _run_stage( + run_stage( "cross-source correlation", - lambda: correlation_module.correlate(index, indicator_rows, len(page_records)), - output, log, + partial(correlation_module.correlate, index, indicator_rows, len(page_records), options.candidate_pair_budget), + output, + log, ) - def temporal_stage() -> CheckResult: - events, parse_errors = temporal.build_events(page_records, content_dates) - result = temporal.analyze_timeline(events, gap_threshold_days=options.gap_days) - result.errors.extend(parse_errors[:50]) - _write_events(run_dir, events, output) - return result - - _run_stage("temporal analysis", temporal_stage, output, log) - _run_stage( + temporal_result = run_stage( + "temporal analysis", partial(_temporal_stage, page_records, content_dates, options.gap_days), output, log + ) + _write_rows( + os.path.join(run_dir, "timeline.jsonl"), temporal_result.stats.pop("events", []) if temporal_result else [] + ) + output.artifacts["timeline"] = os.path.join(run_dir, "timeline.jsonl") + run_stage( "recurring signals and outliers", - lambda: patterns.check_recurring_and_outliers(index, page_scores), output, log, + partial(patterns.check_recurring_and_outliers, index, page_scores), + output, + log, ) # Optional model-assisted stages. Nothing below runs unless it was asked for. evaluation_result: Optional[CheckResult] = None if options.evaluate: - evaluation_result = _run_stage( + evaluation_result = run_stage( "analysis quality evaluation", - lambda: evaluation_module.evaluate_run( - analyses, artifacts.text_for, model=options.model - ), - output, log, + partial(_text_stage, "evaluation", artifacts.run_dir, options, analyses), + output, + log, ) if options.deep and options.use_ollama and ollama is not None: - _run_stage( + run_stage( "deep analysis", - lambda: _deep_analysis( - ollama, options, index, output, page_scores, analyses, run_id - ), - output, log, + partial(_deep_analysis, ollama, options, index, output, page_scores, analyses, run_id), + output, + log, ) if options.cross_check_models and options.use_ollama and ollama is not None: - _run_stage( + run_stage( "multi-model cross-check", - lambda: _cross_check(ollama, options, analyses), output, log, + partial(_cross_check, ollama, options, analyses), + output, + log, ) # Hypotheses derive from the findings that now exist, so this runs after every stage # that can produce one. - def hypothesis_stage() -> CheckResult: - result = CheckResult(check_name="Hypothesis Generation") - derived = hypotheses_module.from_findings(output.findings, index) - result.hypotheses.extend(derived) - result.notes.append( - f"Derived {len(derived)} hypothesis/hypotheses from deterministic findings. " - "Every entry is labelled HYPOTHESIS and is not an observed fact." - ) - return result - - _run_stage("hypothesis generation", hypothesis_stage, output, log) - - def lead_stage() -> CheckResult: - result = CheckResult(check_name="Lead Generation") - generated = pivots.generate_leads(index, per_kind=options.max_leads_per_kind) - result.leads.extend(generated) - result.notes.append( - f"Generated {len(generated)} pivot lead(s). Each is a candidate lookup path, not a hit; " - "confirmation requires opening it." - ) - return result - - _run_stage("lead generation", lead_stage, output, log) - - _write_artifacts(run_dir, output) + run_stage("hypothesis generation", partial(_hypothesis_stage, output.findings, index), output, log) + run_stage("lead generation", partial(_lead_stage, index, options.max_leads_per_kind), output, log) if options.training_export and evaluation_result is not None: try: - def prompt_loader(url: str) -> str: - text = artifacts.text_for(url) - if not text: - return "" - record = next((r for r in page_records if r.get("url") == url), {}) - return page_prompt(url, record.get("title", ""), text, options.prompt_profile) - - summary = training_export_module.export( - run_dir=run_dir, - evaluation_result=evaluation_result, - analyses_by_url=analyses_by_url, - prompt_loader=prompt_loader, - model=options.model, - run_id=run_id, + summary = isolated_call( + _training_stage, + run_dir, + artifacts.run_dir, + options, + evaluation_result, + analyses_by_url, + page_records, + run_id, + timeout_s=options.stage_deadline_s, ) + artifacts.truncated_urls.update(summary["text_coverage"]["truncated"]) + artifacts.missing_urls.update(summary["text_coverage"]["missing"]) + artifacts.unreadable_urls.update(summary["text_coverage"]["unreadable"]) output.artifacts["training_export"] = summary["dir"] output.stats["training_export"] = summary["counts"] log(f"[OK] training dataset exported: {summary['dir']}") except Exception as exc: + import shutil + + shutil.rmtree(os.path.join(run_dir, "training_export"), ignore_errors=True) output.errors.append(f"training export failed: {exc}") log(f"[FAIL] training export -> {exc}") elif options.training_export: @@ -359,6 +472,42 @@ def prompt_loader(url: str) -> str: output.stats["finding_count"] = len(output.findings) output.stats["hypothesis_count"] = len(output.hypotheses) output.stats["lead_count"] = len(output.leads) + for result in output.results: + coverage = result.stats.get("text_coverage", {}) + artifacts.truncated_urls.update(coverage.get("truncated", [])) + artifacts.missing_urls.update(coverage.get("missing", [])) + artifacts.unreadable_urls.update(coverage.get("unreadable", [])) + output.stats["indicator_values_omitted"] = sum( + count["omitted"] for extras in page_extras.values() for count in extras.get("extraction_coverage", {}).values() + ) + output.stats["extraction_failures"] = extraction_failures + output.stats["extraction_cache"] = {"hits": cache.hits, "misses": cache.misses, "invalid": cache.invalid} + output.stats["text_coverage"] = artifacts.coverage() + output.stats["text_cache_budget_bytes"] = options.text_cache_bytes + output.stats["text_char_limit"] = options.max_text_chars + output.stats["pages_text_truncated"] = len(artifacts.truncated_urls) + output.stats["pages_over_scan_limit"] = max(0, len(page_records) - MAX_TEXT_PAGES) + if artifacts.truncated_urls: + output.errors.append( + f"Text truncated to {options.max_text_chars} characters on " + f"{len(artifacts.truncated_urls)} page(s); analysis coverage is partial." + ) + log(f"[WARN] {output.errors[-1]}") + output.stats["model_coverage_partial"] = any( + output.stats["model_responses"]["counts"][status] + for status in ("empty", "invalid", "missing", "timed_out", "error") + ) + output.stats["partial_coverage"] = bool( + output.errors + or extraction_failures + or artifacts.missing_urls + or artifacts.unreadable_urls + or output.stats["pages_over_scan_limit"] + or output.stats["indicator_values_omitted"] + or ((options.use_ollama or options.evaluate) and output.stats["model_coverage_partial"]) + or any(r.errors or r.stats.get("partial_coverage") for r in output.results) + ) + _write_artifacts(run_dir, output) return output @@ -384,18 +533,17 @@ def _deep_analysis( context = { "run_id": run_id, "pages_analyzed": len(page_scores), - "top_pages": [ - {"url": p.get("url"), "score": p.get("score"), "summary": p.get("summary")} - for p in top_pages - ], + "top_pages": [{"url": p.get("url"), "score": p.get("score"), "summary": p.get("summary")} for p in top_pages], "recurring_entities": [ - {"kind": e.kind, "value": e.value, "sources": len(e.sources)} - for e in index.multi_source(minimum=2)[:40] + {"kind": e.kind, "value": e.value, "sources": len(e.sources)} for e in index.multi_source(minimum=2)[:40] ], "deterministic_findings": [ { - "check": f.check, "item": f.item, "reason": f.reason, - "priority": f.priority, "origin": f.origin, + "check": f.check, + "item": f.item, + "reason": f.reason, + "priority": f.priority, + "origin": f.origin, } for f in sort_findings(output.findings)[:60] ], @@ -408,29 +556,25 @@ def _deep_analysis( for pass_number in range(passes): prompt = deep_analysis_prompt(context) try: - payload = asyncio.run( - ollama.async_generate_json(options.model, prompt, timeout_s=180.0) - ) + payload = asyncio.run(ollama.async_generate_json(options.model, prompt, timeout_s=180.0)) except RuntimeError as exc: result.errors.append(f"deep analysis pass {pass_number + 1} could not run: {exc}") break if not payload: - result.errors.append( - f"deep analysis pass {pass_number + 1} returned no parseable response." - ) + result.errors.append(f"deep analysis pass {pass_number + 1} returned no parseable response.") break - result.rows.append({ - "pass": pass_number + 1, - "model": options.model, - "assessment": payload.get("assessment", ""), - "cross_source_observations": payload.get("cross_source_observations", []), - "recommended_follow_up": payload.get("recommended_follow_up", []), - "gaps": payload.get("gaps", []), - }) - result.hypotheses.extend( - hypotheses_module.from_model(payload, options.model, sources) + result.rows.append( + { + "pass": pass_number + 1, + "model": options.model, + "assessment": payload.get("assessment", ""), + "cross_source_observations": payload.get("cross_source_observations", []), + "recommended_follow_up": payload.get("recommended_follow_up", []), + "gaps": payload.get("gaps", []), + } ) + result.hypotheses.extend(hypotheses_module.from_model(payload, options.model, sources)) if pass_number + 1 >= passes: break @@ -445,7 +589,11 @@ def _deep_analysis( if payload is None and not result.errors: result.errors.append("Deep analysis produced no output.") - result.stats = {"passes_run": len(result.rows), "model": options.model} + result.stats = { + "passes_run": len(result.rows), + "model": options.model, + "model_response_counts": getattr(ollama, "response_counts", {}), + } result.notes.append( f"Deep analysis ran {len(result.rows)} model pass(es) with {options.model or 'the configured model'}. " "All output is model-generated interpretation, labelled as such, and is not observed fact." @@ -470,18 +618,24 @@ def _cross_check(ollama, options: AnalysisOptions, analyses: List[Dict[str, Any] result = CheckResult(check_name="Multi-Model Cross-Check") result.errors.append(f"cross-check could not run: {exc}") return result - return multimodel.summarize_cross_checks(checks) + result = multimodel.summarize_cross_checks(checks) + result.stats["model_response_counts"] = getattr(ollama, "response_counts", {}) + return result def _write_events(run_dir: str, events: List[temporal.Event], output: AnalysisOutput) -> None: path = os.path.join(run_dir, "timeline.jsonl") - if os.path.exists(path): - os.remove(path) - for event in events: - append_jsonl(path, event.to_dict()) + _write_rows(path, (event.to_dict() for event in events)) output.artifacts["timeline"] = path +def _write_rows(path: str, rows: Iterable[Dict[str, Any]]) -> None: + # Always create empty result files too; every advertised artifact must exist. + with open(path, "w", encoding="utf-8") as handle: + for row in rows: + handle.write(json.dumps(row, ensure_ascii=False) + "\n") + + def _write_artifacts(run_dir: str, output: AnalysisOutput) -> None: """Write the analysis artifacts. Existing crawl artifacts are never rewritten.""" for name, rows in ( @@ -490,36 +644,93 @@ def _write_artifacts(run_dir: str, output: AnalysisOutput) -> None: ("leads.jsonl", [l.to_dict() for l in output.leads]), ): path = os.path.join(run_dir, name) - if os.path.exists(path): - os.remove(path) - for row in rows: - append_jsonl(path, row) + _write_rows(path, rows) output.artifacts[name.split(".")[0]] = path - correlations = next( - (r for r in output.results if r.check_name == "Cross-Source Correlation"), None - ) - if correlations is not None: - path = os.path.join(run_dir, "correlations.jsonl") - if os.path.exists(path): - os.remove(path) - for row in correlations.rows: - append_jsonl(path, row) - output.artifacts["correlations"] = path + correlations = next((r for r in output.results if r.check_name == "Cross-Source Correlation"), None) + path = os.path.join(run_dir, "correlations.jsonl") + _write_rows(path, correlations.rows if correlations else []) + output.artifacts["correlations"] = path summary_path = os.path.join(run_dir, "analysis_summary.json") - write_json(summary_path, { - "stages": [ - { - "check_name": r.check_name, - "findings": r.finding_count, - "notes": r.notes, - "errors": r.errors, - "stats": r.stats, - } - for r in output.results - ], - "stats": output.stats, - "errors": output.errors, - }) + write_json( + summary_path, + { + "stages": [ + { + "check_name": r.check_name, + "findings": r.finding_count, + "notes": r.notes, + "errors": r.errors, + "stats": r.stats, + } + for r in output.results + ], + "stats": output.stats, + "errors": output.errors, + }, + ) output.artifacts["analysis_summary"] = summary_path + + +def _text_stage(kind, source, options, rows): + reader = RunArtifacts(source, options.max_text_chars, options.text_cache_bytes) + if kind == "generated": + result = patterns.check_generated_text(rows, reader.text_for) + else: + result = evaluation_module.evaluate_run(rows, reader.text_for, model=options.model) + result.stats["text_coverage"] = reader.coverage() + return result + + +def _temporal_stage(records, dates, gap_days): + events, errors = temporal.build_events(records, dates) + result = temporal.analyze_timeline(events, gap_threshold_days=gap_days) + result.errors.extend(errors[:50]) + result.stats["events"] = [event.to_dict() for event in events] + return result + + +def _hypothesis_stage(findings, index): + result = CheckResult(check_name="Hypothesis Generation") + result.hypotheses = hypotheses_module.from_findings(findings, index) + return result + + +def _lead_stage(index, limit): + result = CheckResult(check_name="Lead Generation") + result.leads = pivots.generate_leads(index, per_kind=limit) + return result + + +def analyze_run(run_dir, options, ollama=None, run_id="", log=None, output_dir=None): + """Publish a validated immutable bundle, then atomically advance its pointer.""" + from .publication import publish_analysis + + for value in ( + options.max_text_chars, + options.candidate_pair_budget, + options.text_cache_bytes, + options.page_deadline_s, + options.stage_deadline_s, + ): + if not math.isfinite(value) or value <= 0: + raise ValueError("analysis limits and deadlines must be finite and positive") + return publish_analysis(run_dir, options, ollama, run_id, log, output_dir) + + +def _training_stage(destination, source, options, evaluation_result, analyses_by_url, records, run_id): + reader = RunArtifacts(source, options.max_text_chars, options.text_cache_bytes) + records_by_url = {row.get("url"): row for row in records} + + def prompt_loader(url): + text = reader.text_for(url) + if not text: + return "" + return page_prompt(url, records_by_url.get(url, {}).get("title", ""), text, options.prompt_profile) + + summary = training_export_module.export( + destination, evaluation_result, analyses_by_url, prompt_loader, model=options.model, run_id=run_id + ) + summary["text_coverage"] = reader.coverage() + return summary diff --git a/src/osintai/publication.py b/src/osintai/publication.py new file mode 100644 index 0000000..21551d1 --- /dev/null +++ b/src/osintai/publication.py @@ -0,0 +1,161 @@ +"""Immutable analysis bundles with atomic publication and durable attempt status.""" + +from __future__ import annotations + +import hashlib +import json +import os +import time +import uuid +from dataclasses import asdict +from pathlib import Path + +from . import __version__ +from .checkpoints import EXTRACTOR_VERSION +from .storage import sync_directory, sync_file, write_json + + +def _hash_paths(root, paths): + """Hash files only after their resolved paths remain inside the saved run.""" + root = Path(root).resolve() + hashes = {} + for path in paths: + resolved = path.resolve() + if not resolved.is_relative_to(root): + raise RuntimeError(f"Source artifact escapes saved run: {path}") + if not resolved.is_file(): + continue + digest = hashlib.sha256() + with resolved.open("rb") as handle: + for chunk in iter(lambda: handle.read(1024 * 1024), b""): + digest.update(chunk) + hashes[str(path.relative_to(root))] = digest.hexdigest() + return hashes + + +def source_hashes(source): + root = Path(source).resolve() + files = [ + root / name + for name in ( + "urls_crawled.jsonl", + "indicators.jsonl", + "page_scores.jsonl", + "model_retry_latest.json", + "run_manifest.json", + "extraction_failures.jsonl", + ) + ] + for name in ("pages_text", "analysis"): + files.extend(sorted((root / name).glob("*"))) + # Retry outputs are inputs whenever the reader overlays saved model responses. + files.extend(sorted(root.glob("model_retry_*/analysis/*.json"))) + return _hash_paths(root, files) + + +def retry_source_hashes(source): + """Hash only immutable inputs used to select and prompt model retries.""" + root = Path(source).resolve() + files = [root / "urls_crawled.jsonl", root / "run_manifest.json"] + files.extend(sorted((root / "pages_text").glob("*"))) + files.extend(sorted((root / "analysis").glob("*"))) + return _hash_paths(root, files) + + +def publish_analysis(source, options, ollama, run_id, log, output_dir): + from .pipeline import _analyze_run, _write_artifacts, _write_rows + from .report import write_analysis_report + + root = Path(output_dir or source).resolve() + root.mkdir(parents=True, exist_ok=True) + attempt = uuid.uuid4().hex + staging = root / f".analysis_{attempt}.incomplete" + final = root / f"analysis_{attempt}" + staging.mkdir() + status_path = root / f"analysis_{attempt}.status.json" + status = { + "status": "running", + "started_at": time.time(), + "source_run": str(Path(source).resolve()), + "staging_directory": staging.name, + "bundle_directory": final.name, + } + write_json(str(status_path), status) + try: + hashes = source_hashes(source) + output = _analyze_run(source, options, ollama, run_id, log, str(staging)) + # Empty/failing runs still have a complete and readable bundle. + if "analysis_summary" not in output.artifacts: + _write_artifacts(str(staging), output) + if "timeline" not in output.artifacts: + _write_rows(str(staging / "timeline.jsonl"), []) + output.artifacts["timeline"] = str(staging / "timeline.jsonl") + output.artifacts["analysis_report"] = write_analysis_report( + str(staging), output, run_id=run_id, scope={"source run": str(source)} + ) + report_path = Path(output.artifacts["analysis_report"]) + report_path.write_text( + report_path.read_text(encoding="utf-8").replace(f"RUN DIR: {staging}", f"RUN DIR: {final}"), + encoding="utf-8", + ) + training_manifest = staging / "training_export" / "manifest.json" + if training_manifest.is_file(): + training = json.loads(training_manifest.read_text(encoding="utf-8")) + training["run_dir"] = str(final) + write_json(str(training_manifest), training) + if source_hashes(source) != hashes: + raise RuntimeError("Source artifacts changed during analysis; retry against a stable saved run") + # Paths in the manifest refer to the final location, never the staging directory. + output.artifacts = { + key: str(final / Path(value).relative_to(staging)) for key, value in output.artifacts.items() + } + manifest = { + "osintai_version": __version__, + "extractor_version": EXTRACTOR_VERSION, + "run_id": run_id, + "status": "completed", + "source_run": str(Path(source).resolve()), + "source_hashes": hashes, + "options": asdict(options), + "analysis_stats": output.stats, + "errors": output.errors, + "artifacts": output.artifacts, + "finished_at": time.time(), + "model_stage_responses": { + r.check_name: r.stats["model_response_counts"] + for r in output.results + if "model_response_counts" in r.stats + }, + } + write_json(str(staging / "run_manifest.json"), manifest) + # Parse every structured output before publication and flush every file to disk. + for path in staging.rglob("*"): + if not path.is_file(): + continue + with path.open("r", encoding="utf-8") as handle: + if path.suffix == ".json": + json.load(handle) + elif path.suffix == ".jsonl": + for line in handle: + if line.strip(): + json.loads(line) + sync_file(path) + for directory in staging.rglob("*"): + if directory.is_dir(): + sync_directory(directory) + sync_directory(staging) + os.replace(staging, final) + sync_directory(root) + status.update( + status="completed", + finished_at=time.time(), + partial_coverage=bool(output.errors or output.stats.get("partial_coverage")), + ) + write_json(str(status_path), status) + write_json(str(root / "analysis_latest.json"), {"directory": final.name, "status": "completed"}) + output.artifacts["run_manifest"] = str(final / "run_manifest.json") + return output + except BaseException as exc: + status.update(status="failed", finished_at=time.time(), error=f"{type(exc).__name__}: {exc}") + write_json(str(status_path), status) + raise diff --git a/src/osintai/report.py b/src/osintai/report.py index 3c77759..3e1c86d 100644 --- a/src/osintai/report.py +++ b/src/osintai/report.py @@ -151,6 +151,16 @@ def write_analysis_report( lines.append(f" findings: {len(findings)}") lines.append(f" hypotheses: {len(output.hypotheses)}") lines.append(f" leads: {len(output.leads)}") + lines.append(f" coverage: {'PARTIAL — review limits and failures' if stats.get('partial_coverage') else 'within configured analysis scope'}") + lines.append(f" truncated page texts: {stats.get('pages_text_truncated', 0)}") + lines.append(f" missing/unreadable texts: {len(stats.get('text_coverage', {}).get('missing', [])) + len(stats.get('text_coverage', {}).get('unreadable', []))}") + lines.append(f" extraction failures: {len(stats.get('extraction_failures', []))}") + lines.append(f" pages beyond extended-scan limit: {stats.get('pages_over_scan_limit', 0)}") + lines.append(f" indicator values omitted by extraction caps: {stats.get('indicator_values_omitted', 0)}") + if stats.get("extraction_cache"): + lines.append(f" extraction checkpoints: {stats['extraction_cache']}") + for model, counts in stats.get("model_responses", {}).get("models", {}).items(): + lines.append(f" model responses ({model}): {counts}") # HOW TO READ lines.extend(_section("HOW TO READ THIS REPORT")) diff --git a/src/osintai/scanners.py b/src/osintai/scanners.py new file mode 100644 index 0000000..92f35be --- /dev/null +++ b/src/osintai/scanners.py @@ -0,0 +1,48 @@ +"""Consume maximal tokens once; validate only protocol-sized candidates.""" + +import re + +EMAIL_TOKEN = re.compile(r"[A-Za-z0-9._%+@-]+") +DOMAIN_TOKEN = re.compile(r"[A-Za-z0-9.-]+") +EMAIL = re.compile(r"[A-Za-z0-9._%+-]{1,64}@[A-Za-z0-9.-]{1,253}\.[A-Za-z]{2,63}") +CREDENTIAL_EMAIL = re.compile(r"[\w.+-]{1,64}@[\w.-]{1,253}\.[A-Za-z]{2,63}") +USER = re.compile(r"[\w.-]{3,32}") + + +def valid_domain(value): + if len(value) > 253: + return False + labels = value.split(".") + return ( + len(labels) >= 2 + and labels[-1].isascii() + and labels[-1].isalpha() + and 2 <= len(labels[-1]) <= 63 + and all(label and len(label) <= 63 and label[0] != "-" and label[-1] != "-" for label in labels) + ) + + +def emails(text): + for match in EMAIL_TOKEN.finditer(text): + value = match.group().strip(".") + if len(value) <= 254 and EMAIL.fullmatch(value) and valid_domain(value.rsplit("@", 1)[1]): + yield value + + +def domains(text): + for match in DOMAIN_TOKEN.finditer(text): + value = match.group().strip(".") + if valid_domain(value): + yield value + + +def credentials(text): + for line in text.splitlines(): + # Check the bound before any regex; malformed megabyte lines remain linear. + value = line.strip(" \t") + if len(value) > 319 or ":" not in value: + continue + user, secret = value.split(":", 1) + if 6 <= len(secret) <= 64 and not secret.startswith("//") and not any(c.isspace() for c in secret): + if USER.fullmatch(user) or (len(user) <= 254 and CREDENTIAL_EMAIL.fullmatch(user)): + yield user, secret diff --git a/src/osintai/storage.py b/src/osintai/storage.py index 13be1fa..923c1ef 100644 --- a/src/osintai/storage.py +++ b/src/osintai/storage.py @@ -57,6 +57,25 @@ def write_json(path: str, data: dict): handle.flush() os.fsync(handle.fileno()) os.replace(temp_path, path) + sync_directory(parent) finally: if temp_path and os.path.exists(temp_path): os.unlink(temp_path) + + +def sync_directory(path): + """Persist rename metadata on POSIX; Windows provides no directory fsync here.""" + if os.name != "posix": + return + descriptor = os.open(path, os.O_RDONLY | getattr(os, "O_DIRECTORY", 0)) + try: + os.fsync(descriptor) + finally: + os.close(descriptor) + + +def sync_file(path): + """Flush a completed file through a descriptor Windows permits fsync to use.""" + with open(path, "rb+") as handle: + handle.flush() + os.fsync(handle.fileno()) diff --git a/tests/benchmark_analysis.py b/tests/benchmark_analysis.py new file mode 100644 index 0000000..38b4ddd --- /dev/null +++ b/tests/benchmark_analysis.py @@ -0,0 +1,51 @@ +"""Offline saved-crawl benchmark. Run from the repository with a saved RUN_ID. + +Reports separate parent/maximum-child RSS (not aggregate concurrent memory). +Use separate invocations for cold and warm extraction-cache measurements. +""" + +import argparse +import json +from pathlib import Path +import resource +import sys +import time + +sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "src")) +from osintai.pipeline import AnalysisOptions, analyze_run, RunArtifacts +from osintai.cli import _validate_run_id + + +def main(): + parser = argparse.ArgumentParser() + parser.add_argument("run_id") + parser.add_argument("--output", required=True) + args = parser.parse_args() + root = Path(__file__).resolve().parents[1] + source = root / "data" / "runs" / _validate_run_id(args.run_id) + destination = (root / args.output).resolve() + if not destination.is_relative_to(root): + parser.error("benchmark output must stay inside the repository") + start = time.perf_counter() + output = analyze_run( + str(source), + AnalysisOptions(use_ollama=False), + run_id=args.run_id, + log=lambda message: print(message, flush=True), + ) + factor = 1 if sys.platform == "darwin" else 1024 + result = { + "run_id": args.run_id, + "pages": len(RunArtifacts(str(source)).page_records()), + "elapsed_seconds": round(time.perf_counter() - start, 3), + "parent_peak_rss_bytes": resource.getrusage(resource.RUSAGE_SELF).ru_maxrss * factor, + "max_child_peak_rss_bytes": resource.getrusage(resource.RUSAGE_CHILDREN).ru_maxrss * factor, + "stats": output.stats, + "errors": output.errors, + } + destination.parent.mkdir(parents=True, exist_ok=True) + destination.write_text(json.dumps(result, indent=2), encoding="utf-8") + + +if __name__ == "__main__": + main() diff --git a/tests/benchmark_correlation.py b/tests/benchmark_correlation.py new file mode 100644 index 0000000..c28ae20 --- /dev/null +++ b/tests/benchmark_correlation.py @@ -0,0 +1,53 @@ +"""Reproducible offline boilerplate benchmark; prints timing and RSS as JSON.""" + +import argparse +import json +from pathlib import Path +import resource +import sys +import time + +sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "src")) +from osintai.correlation import correlate +from osintai.entities import EntityIndex + + +def main(): + parser = argparse.ArgumentParser() + parser.add_argument("--pages", type=int, default=10_000) + parser.add_argument("--pair-budget", type=int, default=100_000) + args = parser.parse_args() + if args.pages < 1 or args.pair_budget < 1: + parser.error("page count and pair budget must be positive") + start = time.perf_counter() + index = EntityIndex() + for i in range(args.pages): + url = f"https://site.test/{i}" + for kind, value in ( + ("email", "footer@site.test"), + ("username", "@footer"), + ("name", "Site Footer"), + ("phone", "5551112222"), + ("email", f"person{i}@site.test"), + ("username", f"@person{i}"), + ): + index.add(kind, value, url) + indexed = time.perf_counter() + result = correlate(index, [], page_count=args.pages, candidate_budget=args.pair_budget) + factor = 1 if sys.platform == "darwin" else 1024 + print( + json.dumps( + { + "pages": args.pages, + "index_seconds": round(indexed - start, 3), + "correlation_seconds": round(time.perf_counter() - indexed, 3), + "parent_peak_rss_bytes": resource.getrusage(resource.RUSAGE_SELF).ru_maxrss * factor, + "stats": result.stats, + }, + indent=2, + ) + ) + + +if __name__ == "__main__": + main() diff --git a/tests/test_analysis_layer.py b/tests/test_analysis_layer.py index 47eccd1..b760a8d 100644 --- a/tests/test_analysis_layer.py +++ b/tests/test_analysis_layer.py @@ -145,7 +145,7 @@ def test_indicator_schema_keys_are_unchanged(self): "eth_count", "social_count", } indicators = Extractor().extract_indicators("https://one.test/", "text", "") - self.assertEqual(set(indicators), expected) + self.assertEqual(set(indicators), expected | {"url_provenance"}) class NoTrainingInOSINTaiTests(unittest.TestCase): @@ -236,7 +236,7 @@ def test_observed_forms_are_separate_from_possible_variants(self): def test_extended_extraction_finds_new_indicator_classes(self): extras = entities.extract_extended( "Meeting on 2024-01-31 with John Martinez at 4421 Troost Ave.\n" - "api_key: AKIAIOSFODNN7EXAMPLEKEY123\n" + "api_key: SYNTHETIC_TEST_TOKEN_0123456789\n" ) self.assertIn("2024-01-31", extras["dates"]) self.assertIn("John Martinez", extras["name_candidates"]) @@ -701,6 +701,10 @@ def _build_run(run_dir: str) -> None: }, handle) +def _deliberate_stage_failure(index): + raise RuntimeError("deliberate stage failure") + + class PipelineTests(unittest.TestCase): def test_full_offline_analysis_produces_every_artifact(self): with tempfile.TemporaryDirectory() as run_dir: @@ -713,7 +717,7 @@ def test_full_offline_analysis_produces_every_artifact(self): self.assertEqual(output.errors, []) for name in ("findings.jsonl", "hypotheses.jsonl", "leads.jsonl", "correlations.jsonl", "timeline.jsonl", "analysis_summary.json"): - self.assertTrue(os.path.exists(os.path.join(run_dir, name)), name) + self.assertTrue(os.path.exists(output.artifacts[name.split(".")[0]]), name) checks = {f.check for f in output.findings} self.assertIn("Unicode / Homoglyph", checks) # Cyrillic domain @@ -746,9 +750,7 @@ def test_a_failing_stage_does_not_abort_the_pipeline(self): from osintai import pipeline as pipeline_module original = pipeline_module.patterns.check_homoglyphs - pipeline_module.patterns.check_homoglyphs = lambda index: (_ for _ in ()).throw( - RuntimeError("deliberate stage failure") - ) + pipeline_module.patterns.check_homoglyphs = _deliberate_stage_failure try: with tempfile.TemporaryDirectory() as run_dir: _build_run(run_dir) diff --git a/tests/test_analysis_recovery.py b/tests/test_analysis_recovery.py new file mode 100644 index 0000000..6fa6b99 --- /dev/null +++ b/tests/test_analysis_recovery.py @@ -0,0 +1,134 @@ +"""Bounded regression checks for the regex hang and saved-run recovery.""" +import contextlib +import argparse +import io +import json +import os +from pathlib import Path +import subprocess +import sys +import tempfile +import unittest +from unittest.mock import patch + +sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "src")) +from osintai import __version__, cli, entities +from osintai.extractor import Extractor +from osintai.hunt import hunt_leads +from osintai.pipeline import AnalysisOptions, analyze_run +from osintai.storage import sha1 + + +def _bad_extraction(text): + raise ValueError("bad page") + + +class RecoveryTests(unittest.TestCase): + def test_adversarial_domain_scan_has_deadline(self): + # An external deadline prevents a reintroduced regex bug hanging CI. + code = """ +from osintai.entities import extract_extended, unicode_domains +assert not list(unicode_domains('a' * 1000000)) +assert not list(unicode_domains(('a.' * 500000) + '!')) +assert 'раypal.com' in extract_extended('a' * 10000 + ' раypal.com')['unicode_domains'] +""" + env = dict(os.environ, PYTHONPATH=str(Path(__file__).resolve().parents[1] / "src")) + subprocess.run([sys.executable, "-c", code], env=env, check=True, timeout=10) + + def test_unicode_domains_and_boundaries(self): + actual = set(entities.unicode_domains( + 'ascii.com раypal.com münich.de example..рф -bad.рф zero\u200bwidth.com' + )) + self.assertEqual(actual, {'раypal.com', 'münich.de', 'zero\u200bwidth.com'}) + + def test_html_elements_do_not_glue_urls(self): + _, text = Extractor().html_to_text( + '

https://example.com/api.md

Postman' + ) + self.assertEqual(text, 'https://example.com/api.md\nProduct\nPostman') + + def build_run(self, directory): + root = Path(directory) + (root / 'pages_text').mkdir() + url = 'https://example.test/' + (root / 'urls_crawled.jsonl').write_text(json.dumps({'url': url}) + '\n') + (root / 'pages_text' / (sha1(url) + '.txt')).write_text('a' * 1000) + return root + + def test_recovery_preserves_source_and_reports_limits(self): + with tempfile.TemporaryDirectory() as directory: + root = self.build_run(directory) + before = {p: p.read_bytes() for p in root.rglob('*') if p.is_file()} + logs = [] + output = analyze_run(str(root), AnalysisOptions(use_ollama=False, max_text_chars=100), + output_dir=str(root / 'recovered'), log=logs.append) + self.assertEqual(output.stats['pages_text_truncated'], 1) + self.assertTrue(any('[SCAN] 1/1' in message for message in logs)) + summary = json.loads(Path(output.artifacts['analysis_summary']).read_text()) + self.assertEqual(summary['stats'], output.stats) + for path in output.artifacts.values(): + self.assertTrue(Path(path).exists(), path) + for path, content in before.items(): + self.assertEqual(path.read_bytes(), content) + + def test_extraction_failure_is_isolated(self): + with tempfile.TemporaryDirectory() as directory: + root = self.build_run(directory) + with patch('osintai.pipeline.extract_checkpoint', _bad_extraction): + output = analyze_run(str(root), AnalysisOptions(use_ollama=False)) + self.assertTrue(any('bad page' in error for error in output.errors)) + self.assertIn('analysis_summary', output.artifacts) + + def test_cli_recovery_never_constructs_network_clients(self): + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + run = root / 'data' / 'runs' / 'saved' + run.mkdir(parents=True) + self.build_run(run) + args = type('Args', (), dict(analyze_only='saved', no_analysis=False, deep=False, + cross_check='', experimental_recursive=0, evaluate=False, training_export=False, + gap_days=90, leads_per_kind=10, analysis_max_chars=200000))() + with patch('osintai.cli.OllamaAPI', side_effect=AssertionError('network')), \ + patch('osintai.cli.AsyncCrawler', side_effect=AssertionError('crawl')), \ + contextlib.redirect_stdout(io.StringIO()): + cli._analyze_saved_run(args, argparse.ArgumentParser(), str(root)) + cli._analyze_saved_run(args, argparse.ArgumentParser(), str(root)) + self.assertEqual(len(list(run.glob('reanalysis_*/analysis_*/analysis_report.txt'))), 2) + + def test_interrupt_is_clean(self): + with patch('osintai.cli._main', side_effect=KeyboardInterrupt), \ + contextlib.redirect_stderr(io.StringIO()) as stderr: + with self.assertRaises(SystemExit) as raised: + cli.main() + self.assertEqual(raised.exception.code, 130) + self.assertIn('--analyze-only', stderr.getvalue()) + + def test_profile_respects_equals_syntax(self): + with patch.object(sys, 'argv', ['run_osintai.py', '--profile=survey', '--max=12', + '--depth=1', '--analyze-only=saved']), \ + patch('osintai.cli._analyze_saved_run') as recover: + cli.main() + args = recover.call_args.args[0] + self.assertEqual(args.max, 12) + self.assertEqual(args.depth, 1) + + def test_invalid_limits_fail_before_recovery(self): + for flag in ('--analysis-max-chars=0', '--concurrency=0', '--depth=-1'): + with self.subTest(flag=flag), \ + patch.object(sys, 'argv', ['run_osintai.py', '--analyze-only=saved', flag]), \ + patch('osintai.cli._analyze_saved_run') as recover, \ + contextlib.redirect_stderr(io.StringIO()): + with self.assertRaises(SystemExit) as raised: + cli.main() + self.assertEqual(raised.exception.code, 2) + recover.assert_not_called() + + def test_hunt_empty_terms_and_global_hit_limit(self): + self.assertEqual(hunt_leads('text', [''])['hits'], []) + result = hunt_leads('api ' * 600, ['api', 'api']) + self.assertEqual(len(result['hits']), 500) + + def test_readme_version_matches_package(self): + readme = (Path(__file__).resolve().parents[1] / 'README.md').read_text() + self.assertIn(f'# OSINTai v{__version__} ', readme) + self.assertIn(f'### v{__version__} (', readme) diff --git a/tests/test_enhancements.py b/tests/test_enhancements.py new file mode 100644 index 0000000..ff2cf5e --- /dev/null +++ b/tests/test_enhancements.py @@ -0,0 +1,380 @@ +"""Offline acceptance tests for containment, resumption, provenance and publication.""" + +import asyncio +import json +import multiprocessing +import os +from pathlib import Path +import subprocess +import sys +import tempfile +import time +import unittest +from unittest.mock import patch + +sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "src")) +from osintai.checkpoints import ExtractionCache, extract_checkpoint +from osintai.correlation import correlate +from osintai.entities import EntityIndex +from osintai.extractor import Extractor +from osintai.hunt import hunt_leads +from osintai.isolation import isolated_call, async_isolated_call +from osintai.model_quality import page_status +from osintai.ollama_api import OllamaAPI +import httpx +from osintai.model_retry import retry_saved +from osintai.pipeline import AnalysisOptions, RunArtifacts, analyze_run +from osintai.publication import retry_source_hashes, source_hashes +from osintai.storage import sha1, sync_file, write_json +from osintai.scanners import credentials + + +def busy_worker(): + while True: + pass + + +def crash_worker(): + os._exit(7) + + +def large_result(): + return "x" * 2_000_000 + + +def make_run(root, texts=("John Smith 2026-01-01",)): + root = Path(root) + (root / "pages_text").mkdir(parents=True, exist_ok=True) + records = [] + for i, text in enumerate(texts): + url = f"https://example.test/{i}" + records.append({"url": url, "title": "Test", "fetched_at": 1_770_000_000}) + (root / "pages_text" / f"{sha1(url)}.txt").write_text(text, encoding="utf-8") + (root / "urls_crawled.jsonl").write_text("".join(json.dumps(row) + "\n" for row in records)) + return records + + +class IsolationTests(unittest.TestCase): + def test_cpu_deadline_reaps_child(self): + before = {p.pid for p in multiprocessing.active_children()} + started = time.monotonic() + with self.assertRaises(TimeoutError): + isolated_call(busy_worker, timeout_s=0.3) + self.assertLess(time.monotonic() - started, 4) + self.assertEqual({p.pid for p in multiprocessing.active_children()}, before) + + def test_child_crash_is_failure(self): + with self.assertRaisesRegex(RuntimeError, "without a result"): + isolated_call(crash_worker, timeout_s=5) + + def test_large_result_does_not_deadlock(self): + self.assertEqual(len(isolated_call(large_result, timeout_s=5)), 2_000_000) + + def test_async_cancellation_cleans_up_worker(self): + async def scenario(): + before = {process.pid for process in multiprocessing.active_children()} + task = asyncio.create_task(async_isolated_call(busy_worker, timeout_s=10)) + await asyncio.sleep(0.2) + task.cancel() + with self.assertRaises(asyncio.CancelledError): + await task + self.assertEqual({process.pid for process in multiprocessing.active_children()}, before) + + asyncio.run(scenario()) + + def test_live_page_extraction_runs_in_worker(self): + from osintai.crawler import _extract_page + + title, text, links, indicators, hunt, fingerprint = isolated_call( + _extract_page, + Extractor(), + "https://example.test", + 'Test

threat a@example.test

contact', + ["threat"], + 10, + timeout_s=5, + ) + self.assertEqual(title, "Test") + self.assertIn("a@example.test", indicators["emails"]) + self.assertIn("https://example.test/contact", links) + self.assertTrue(hunt["hits"]) + self.assertIsInstance(fingerprint, int) + + +class EnhancementTests(unittest.TestCase): + def test_completed_file_can_be_synced_through_writable_descriptor(self): + with tempfile.TemporaryDirectory() as directory: + path = Path(directory) / "validated.json" + path.write_bytes(b"validated") + sync_file(path) + + def test_adversarial_patterns_have_external_deadline(self): + code = """ +from osintai.extractor import Extractor +from osintai.entities import extract_extended, detect_type +for text in ['a.' * 200000 + '!', 'x-' * 200000, 'a@' + 'a.' * 200000 + '!', 'eyJaaaa-' * 50000, ' ' * 400000 + 'a@' + '.' * 400000]: + Extractor().extract_indicators('https://example.test', text, '') + extract_extended(text) + detect_type(text) +assert 'a@example.test' in Extractor().extract_indicators('https://example.test', 'a@example.test', '')['emails'] +""" + env = dict(os.environ, PYTHONPATH=str(Path(__file__).resolve().parents[1] / "src")) + subprocess.run([sys.executable, "-c", code], env=env, check=True, timeout=15) + + def test_unicode_hunt_offsets_and_complete_urls(self): + text = "İ İ " + "threat https://example.test/" + "a" * 200 + result = hunt_leads(text, ["THREAT"]) + self.assertEqual(result["hits"][0]["position"], text.index("threat")) + self.assertEqual(result["lead_urls"], []) + self.assertEqual( + hunt_leads("threat https://example.test/a", ["threat"])["lead_urls"], ["https://example.test/a"] + ) + + def test_attribute_and_prose_provenance(self): + indicators = Extractor().extract_indicators( + "https://example.test", "https://example.test/prose", 'label' + ) + self.assertEqual(indicators["url_provenance"]["https://example.test/link"], ["html_attribute"]) + self.assertEqual(indicators["url_provenance"]["https://example.test/prose"], ["page_prose"]) + + def test_cache_is_content_addressed_and_contains_no_secrets(self): + text = "api_key: SECRETvalue0123456789\nadmin:password123\n" + with tempfile.TemporaryDirectory() as directory: + cache = ExtractionCache(directory) + key, metadata = cache.key(text, 1000, False) + extras = extract_checkpoint(text) + cache.put(key, metadata, extras) + self.assertEqual(cache.get(key, metadata), extras) + self.assertNotEqual(key, cache.key(text + "x", 1000, False)[0]) + self.assertNotEqual(key, cache.key(text, 2000, False)[0]) + self.assertNotEqual(key, cache.key(text, 1000, True)[0]) + data = (Path(directory) / f"{key}.json").read_text() + self.assertNotIn("SECRETvalue0123456789", data) + self.assertNotIn("password123", data) + write_json(str(Path(directory) / f"{key}.json"), {"broken": True}) + self.assertIsNone(cache.get(key, metadata)) + self.assertEqual(cache.invalid, 1) + + def test_unicode_credential_email_is_retained(self): + self.assertEqual(list(credentials("user@пример.com:password123")), [("user@пример.com", "password123")]) + + def test_secret_redaction_extends_beyond_finding_caps(self): + text = "\n".join(f"admin:password{i:03d}" for i in range(100)) + "\nadmin:раypal.com" + extras = extract_checkpoint(text) + self.assertNotIn("раypal.com", json.dumps(extras, ensure_ascii=False)) + self.assertEqual(extras["extraction_coverage"]["credential_pairs"]["omitted"], 1) + + def test_lru_cache_evicts_and_is_byte_bounded(self): + with tempfile.TemporaryDirectory() as directory: + records = make_run(directory, ("a" * 500, "b" * 500, "c" * 500)) + reader = RunArtifacts(directory, text_cache_bytes=1000) + for row in records: + self.assertEqual(len(reader.text_for(row["url"])), 500) + self.assertLessEqual(reader.cache_bytes, 1000) + self.assertGreater(reader.cache_evictions, 0) + self.assertEqual(reader.text_for(records[0]["url"]), "a" * 500) + + def test_pair_budget_accounts_for_omissions(self): + index = EntityIndex() + for i in range(10): + index.add("email", f"user{i}@example.test", "https://example.test") + result = correlate(index, [], page_count=1, candidate_budget=7) + self.assertEqual(result.stats["candidate_pairs_examined"], 7) + self.assertEqual(result.stats["candidate_pairs_omitted"], 38) + self.assertTrue(result.stats["partial_coverage"]) + + def test_resume_reuses_cache_and_preserves_sources(self): + with tempfile.TemporaryDirectory() as directory: + records = make_run(directory) + before = {path: path.read_bytes() for path in Path(directory).rglob("*") if path.is_file()} + first = analyze_run(directory, AnalysisOptions(use_ollama=False)) + second = analyze_run(directory, AnalysisOptions(use_ollama=False)) + self.assertEqual(first.stats["extraction_cache"]["misses"], 1) + self.assertEqual(second.stats["extraction_cache"]["hits"], 1) + self.assertNotEqual(first.artifacts["analysis_summary"], second.artifacts["analysis_summary"]) + for path, data in before.items(): + self.assertEqual(path.read_bytes(), data) + manifest = json.loads(Path(second.artifacts["run_manifest"]).read_text()) + self.assertIn("urls_crawled.jsonl", manifest["source_hashes"]) + self.assertEqual(manifest["status"], "completed") + self.assertEqual(second.stats["model_responses"]["counts"]["missing"], len(records)) + + def test_failed_publication_retains_previous_pointer(self): + with tempfile.TemporaryDirectory() as directory: + make_run(directory) + analyze_run(directory, AnalysisOptions(use_ollama=False)) + pointer = (Path(directory) / "analysis_latest.json").read_bytes() + with patch("osintai.report.write_analysis_report", side_effect=KeyboardInterrupt): + with self.assertRaises(KeyboardInterrupt): + analyze_run(directory, AnalysisOptions(use_ollama=False)) + self.assertEqual((Path(directory) / "analysis_latest.json").read_bytes(), pointer) + states = [json.loads(path.read_text())["status"] for path in Path(directory).glob("*.status.json")] + self.assertIn("failed", states) + self.assertEqual(len(list(Path(directory).glob("*.incomplete"))), 1) + + def test_deadline_failure_is_explicit_partial_coverage(self): + with tempfile.TemporaryDirectory() as directory: + make_run(directory) + output = analyze_run(directory, AnalysisOptions(use_ollama=False, page_deadline_s=0.000001)) + self.assertEqual(output.stats["extraction_failures"][0]["status"], "timed_out") + self.assertTrue(output.stats["partial_coverage"]) + self.assertTrue(Path(output.artifacts["analysis_report"]).is_file()) + + def test_cache_write_failure_keeps_extracted_results(self): + with tempfile.TemporaryDirectory() as directory: + make_run(directory) + with patch("osintai.checkpoints.ExtractionCache.put", side_effect=OSError("read only")): + output = analyze_run(directory, AnalysisOptions(use_ollama=False)) + self.assertEqual(output.stats["pages_text_scanned"], 1) + self.assertTrue(any("extracted results retained" in error for error in output.errors)) + + def test_source_change_prevents_publication(self): + with tempfile.TemporaryDirectory() as directory: + make_run(directory) + with patch("osintai.publication.source_hashes", side_effect=[{"source": "before"}, {"source": "after"}]): + with self.assertRaisesRegex(RuntimeError, "Source artifacts changed"): + analyze_run(directory, AnalysisOptions(use_ollama=False)) + self.assertFalse((Path(directory) / "analysis_latest.json").exists()) + + def test_saved_model_quality_classifies_outcomes(self): + self.assertEqual(page_status(None), "empty") + self.assertEqual(page_status({"summary": 3}), "invalid") + self.assertEqual(page_status({"_model_status": "timed_out"}), "timed_out") + with tempfile.TemporaryDirectory() as directory: + records = make_run(directory, ("a", "b", "c", "d")) + analysis = Path(directory) / "analysis" + analysis.mkdir() + for row, payload in zip(records, (None, {"_model_status": "timed_out"}, [])): + write_json(str(analysis / f"{sha1(row['url'])}.analysis.json"), payload) + rows = RunArtifacts(directory).model_records() + self.assertEqual([r["_model_status"] for r in rows], ["empty", "timed_out", "invalid", "missing"]) + + def test_retry_overlay_rejects_external_symlink(self): + with tempfile.TemporaryDirectory() as directory: + parent = Path(directory) + run = parent / "run" + outside = parent / "outside" + records = make_run(run) + (outside / "analysis").mkdir(parents=True) + write_json( + str(outside / "analysis" / f"{sha1(records[0]['url'])}.analysis.json"), + { + "summary": "outside", + "key_entities": [], + "key_locations": [], + "key_dates": [], + "keywords": [], + "risk_flags": [], + "actionable_leads": [], + }, + ) + write_json( + str(outside / "run_manifest.json"), + {"status": "completed", "source_hashes": retry_source_hashes(run)}, + ) + try: + (run / "model_retry_escape").symlink_to(outside, target_is_directory=True) + except OSError as exc: + self.skipTest(f"directory symlinks unavailable: {exc}") + write_json(str(run / "model_retry_latest.json"), {"directory": "model_retry_escape"}) + self.assertEqual(RunArtifacts(str(run)).model_records()[0]["_model_status"], "missing") + with self.assertRaisesRegex(RuntimeError, "escapes saved run"): + source_hashes(run) + + def test_retry_overlay_requires_current_source_hashes(self): + with tempfile.TemporaryDirectory() as directory: + run = Path(directory) + records = make_run(run) + retry = run / "model_retry_stale" + current_hashes = retry_source_hashes(run) + (retry / "analysis").mkdir(parents=True) + write_json( + str(retry / "analysis" / f"{sha1(records[0]['url'])}.analysis.json"), + { + "summary": "stale", + "key_entities": [], + "key_locations": [], + "key_dates": [], + "keywords": [], + "risk_flags": [], + "actionable_leads": [], + }, + ) + write_json( + str(retry / "run_manifest.json"), + {"status": "completed", "source_hashes": current_hashes}, + ) + write_json(str(run / "model_retry_latest.json"), {"directory": retry.name}) + self.assertEqual(RunArtifacts(str(run)).model_records()[0]["_model_status"], "ok") + (run / "pages_text" / f"{sha1(records[0]['url'])}.txt").write_text("changed", encoding="utf-8") + self.assertEqual(RunArtifacts(str(run)).model_records()[0]["_model_status"], "missing") + + def test_model_transport_distinguishes_failures(self): + async def scenario(): + for body, expected in ( + ({}, "missing"), + ({"response": ""}, "empty"), + ({"response": "{}"}, "empty"), + ({"response": "broken"}, "invalid"), + ({"response": '{"summary":"ok"}'}, "ok"), + ): + transport = httpx.MockTransport(lambda request: httpx.Response(200, json=body)) + client = httpx.AsyncClient(transport=transport) + api = OllamaAPI() + with patch("osintai.ollama_api.httpx.AsyncClient", return_value=client): + result = await api.async_generate_result("fake", "prompt") + self.assertEqual(result["status"], expected) + self.assertEqual(api.response_counts, {"fake": {expected: 1}}) + + def timeout(request): + raise httpx.ReadTimeout("timeout") + + client = httpx.AsyncClient(transport=httpx.MockTransport(timeout)) + api = OllamaAPI() + with patch("osintai.ollama_api.httpx.AsyncClient", return_value=client): + result = await api.async_generate_result("fake", "prompt") + self.assertEqual(result["status"], "timed_out") + + asyncio.run(scenario()) + + def test_model_retry_has_total_deadline(self): + class SlowModel: + async def async_generate_result(self, *args, **kwargs): + await asyncio.sleep(20) + + with tempfile.TemporaryDirectory() as directory: + make_run(directory) + destination = asyncio.run(retry_saved(directory, SlowModel(), "slow", limit=1, timeout_s=0.01)) + manifest = json.loads((Path(destination) / "run_manifest.json").read_text()) + self.assertEqual(manifest["model_responses"]["counts"]["timed_out"], 1) + + def test_stale_extractor_version_invalidates_checkpoint(self): + with tempfile.TemporaryDirectory() as directory: + cache = ExtractionCache(directory) + first, _ = cache.key("some text", 1000, False) + with patch("osintai.checkpoints.EXTRACTOR_VERSION", "new-implementation"): + second, _ = cache.key("some text", 1000, False) + self.assertNotEqual(first, second) + + def test_retry_is_bounded_and_keeps_originals(self): + class FakeModel: + calls = 0 + + async def async_generate_result(self, *args, **kwargs): + self.calls += 1 + return {"status": "empty", "payload": None} + + with tempfile.TemporaryDirectory() as directory: + make_run(directory, ("some text", "more text", "last text")) + model = FakeModel() + destination = asyncio.run(retry_saved(directory, model, "fake", limit=2)) + self.assertEqual(model.calls, 2) + manifest = json.loads((Path(destination) / "run_manifest.json").read_text()) + self.assertEqual(manifest["attempted"], 2) + self.assertEqual(manifest["omitted"], 1) + self.assertEqual(manifest["model_responses"]["counts"]["empty"], 2) + self.assertFalse((Path(directory) / "analysis").exists()) + + +if __name__ == "__main__": + unittest.main()