diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 8578d31..3110dd6 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -24,6 +24,9 @@ jobs: run: | uv venv uv pip install -e ".[dev]" + # ia:// lives in fsspec's ia module, unreleased; installed on top of the released fsspec + # (a git dependency cannot resolve against s3fs's exact fsspec pin). Drop once fsspec ships it. + uv pip install "fsspec[http] @ git+https://github.com/lfoppiano/filesystem_spec@ia-filesystem" - name: Smoke test run: | diff --git a/.github/workflows/long-tests.yml b/.github/workflows/long-tests.yml new file mode 100644 index 0000000..2996c4e --- /dev/null +++ b/.github/workflows/long-tests.yml @@ -0,0 +1,55 @@ +# Manual only: streams every README example archive end to end (0.4-10 GB each), which is far too +# slow and bandwidth-heavy for the per-push CI. Start it from the Actions tab ("Run workflow") or +# with `gh workflow run long-tests.yml`. One job per output format so the two halves run in parallel. +name: Long tests + +on: + workflow_dispatch: + inputs: + keyword: + description: "pytest -k expression to narrow the run (e.g. 'not megawarc', 'cc-us-federal')" + required: false + default: "" + python-version: + description: "Python version" + required: false + default: "3.12" + +jobs: + long: + runs-on: ubuntu-latest + timeout-minutes: 360 # the hosted-runner maximum + strategy: + fail-fast: false + matrix: + format: [flat, sidecar] + steps: + - uses: actions/checkout@v4 + + - uses: actions/setup-python@v5 + with: + python-version: ${{ inputs.python-version }} + + - uses: astral-sh/setup-uv@v4 + + - name: Install dependencies + run: | + uv venv + uv pip install -e ".[dev]" + # ia:// lives in fsspec's ia module, unreleased; installed on top of the released fsspec + # (a git dependency cannot resolve against s3fs's exact fsspec pin). Drop once fsspec ships it. + uv pip install "fsspec[http] @ git+https://github.com/lfoppiano/filesystem_spec@ia-filesystem" + + - name: Disk before + run: df -h / + + - name: Long tests (${{ matrix.format }}) + env: + KEYWORD: ${{ inputs.keyword }} + run: | + source .venv/bin/activate + if [ -n "$KEYWORD" ]; then + pytest -m long -k "${{ matrix.format }} and ($KEYWORD)" -v --durations=0 + else + pytest -m long -k "${{ matrix.format }}" -v --durations=0 + fi diff --git a/README.md b/README.md index 6ab7c0d..87b537c 100644 --- a/README.md +++ b/README.md @@ -8,8 +8,6 @@ collection, perhaps a repackage of Common Crawl containing just webpages labeled Swahili. `warc2zip` then converts these WARCs into zip files containing the web capture payloads, with metadata stored as csv (spreadsheet) files. -FIXME: should we always write WARC (all caps) other than filenames? - Each response record's payload is stored as a individual file with a proper extension, derived from its Content-Type. Metadata (both WARC and http) from the request, response, and metadata records are written to CSV (spreadsheet) files. @@ -26,7 +24,15 @@ Metadata (both WARC and http) from the request, response, and metadata records a pip install . ``` -By default, pip will install remote access tools, namely `fsspec` configured to talk to https and s3 remote files. +By default, pip will install remote access tools, namely `fsspec` configured to talk to https, s3, and Internet Archive (`ia://`) remote files. + +> [!NOTE] +> `ia://` relies on `fsspec.implementations.ia`, which is not in a released `fsspec` yet. Until it is, install the +> branch behind the upstream pull request on top: +> +> ```bash +> pip install "fsspec[http] @ git+https://github.com/lfoppiano/filesystem_spec@ia-filesystem" +> ``` ## Usage @@ -34,13 +40,21 @@ By default, pip will install remote access tools, namely `fsspec` configured to warc2zip ``` -Input can be a local path or a remote URI (S3, HTTP, etc.): +Input can be a local path or a remote URI (S3, HTTP, an Internet Archive item, etc.): ```bash warc2zip s3://commoncrawl/crawl-data/.../CC-MAIN-....warc.gz warc2zip https://data.commoncrawl.org/crawl-data/.../CC-MAIN-....warc.gz +warc2zip ia://EOT24PRE-20240926175758-crawl808/EOT24PRE-20240926175758-00032.warc.gz ``` +`ia:///` names a file in an [archive.org](https://archive.org) item; it is read from +`https://archive.org/download//`, so the two spellings are interchangeable. Public items +need no account. For a restricted item, log in once with the `internetarchive` package — `pip install +internetarchive && ia configure` writes `~/.config/internetarchive/ia.ini`, which `warc2zip` reads the way `ia` +does (`$IA_CONFIG_FILE` first) — or set `IA_ACCESS_KEY_ID` and `IA_SECRET_ACCESS_KEY`. A refused item says which +of those it tried. + **Note**: Please use s3 inside of AWS and https outside. **Note**: `s3://commoncrawl` does **not** allow anonymous access — requests must be signed with credentials from any AWS account. To fetch Common Crawl data without an AWS account, use the HTTPS endpoint (`https://data.commoncrawl.org/...`) instead. @@ -49,12 +63,15 @@ warc2zip https://data.commoncrawl.org/crawl-data/.../CC-MAIN-....warc.gz | Flag | Description | Default | |---------------------------|----------------------------------------------------------------------------------------|-----------------------------------------| -| `input_file` | Path or URI to a `.warc.gz` file (positional, required) | | -| `--output` | Path to the output zip file | Replace `.warc.gz` with `.zip` | +| `input_file` | Path or URI to a `.warc.gz` file: local, `s3://`, `http(s)://` or `ia:///` (positional, required) | | +| `--output` | Path to the output zip file (the output `.warc.gz` with `--fetch`) | `{basename}_{hex}.zip` in the current directory, same hex and `_partial` rule as the root directory inside (`{basename}_{hex}.warc.gz` with `--fetch`) | | `--dry-run` | Print summary without creating output. The scan always stops after at most 10 capture records, so it never streams the whole file; a lower `--limit` is respected | | | `--limit ` | Limit to N capture records, with their full set of associated request/metadata records | No limit, all records are processed | | `--format {flat,sidecar}` | Output format (see [Output Formats](#output-formats) below) | `flat` | | `--metadata-only` | Write every CSV, manifest and sidecar but no payload files | Off, payloads are written | +| `--fetch` | Treat `input_file` as a `manifest.csv` and download every row's byte range into one `.warc.gz` (see [Building and downloading a subset](#building-and-downloading-a-subset)). Not combinable with `--limit` or `--format` | Off | +| `--rate ` | `--fetch` only: requests per second per host, `0` for unlimited | [cdx_toolkit](https://github.com/commoncrawl/cdx_toolkit)'s per-host pacing for http(s), `2` for s3 and ia | +| `--retries ` | `--fetch` only: connection failures tolerated per request (http(s), throttling is retried without limit) or retries per request (s3, ia) | `100` / `8` | ### Small Examples @@ -146,7 +163,7 @@ Warnings never change the exit status; warc2zip exits 1 only when a CSV row coul ## Output Formats -All files are placed under a unique root directory inside the zip to prevent collisions when extracting multiple archives into the same folder. The directory name is derived from the WARC-Filename header (in the `warcinfo` record), the current timestamp, and a random suffix: `{crawl_name}_{YYYYMMDDTHHMMSS}_{hex}`. +All files are placed under a unique root directory inside the zip to prevent collisions when extracting multiple archives into the same folder. The directory name is derived from the WARC-Filename header (in the `warcinfo` record), the current timestamp, and a random suffix: `{crawl_name}_{YYYYMMDDTHHMMSS}_{hex}`. Without `--output`, the zip carries the same hex (`archive_{hex}.zip`), so same-named inputs such as Common Crawl's `warc/`, `crawldiagnostics/` and `robotstxt/` files never overwrite each other's zip. A `--limit` run appends `_partial` to both names (`archive_{hex}_partial.zip`). This testing release of the software supports 2 output formats: flat and sidecar. Flat puts the metadata into a small number of large files, and sidecar instead @@ -185,7 +202,7 @@ FOO.zip - **Denormalized CSVs** (`*_headers.csv`, `metadata.csv`, `warcinfo.csv`): multiple rows per file — columns: `filename, header_name, header_value`. Header names are normalized to lowercase with `-` replaced by `_`. - **Multiline CSVs** (`*_multi.csv`): one row per file — columns: `filename, headers` (headers as a multiline string) - **`manifest.csv`**: the wide mirror of `manifest.jsonl` — one row per response, one column per key. Both are written from the same entries, so they cannot drift. -- **`warcinfo.*`**: crawl-level provenance (`isPartOf`, `publisher`, `software`, `hostname`, `conformsTo`, …). The `filename` column carries the synthetic key `warcinfo`, since the record belongs to no payload file; a concatenated WARC with several warcinfo records numbers the extras `warcinfo.1`, `warcinfo.2`, …. +- **`warcinfo.*`**: crawl-level provenance (`isPartOf`, `publisher`, `software`, `hostname`, `conformsTo`, …). The `filename` column carries the synthetic key `warcinfo`, since the record belongs to no payload file; a concatenated WARC with several `warcinfo` records numbers the extras `warcinfo.1`, `warcinfo.2`, …. ### Sidecar format (`--format sidecar`) @@ -333,13 +350,13 @@ A `revisit` row names no file in the zip (its `filename` is a synthetic join key The last four columns are what make a row re-fetchable on its own, and they are repeated on every row on purpose — no join against another file, no knowledge of the zip they came from: -| Column | Meaning | -|---|---| -| `warc_filename` | What the WARC **calls itself** (the warcinfo record's `WARC-Filename`), falling back to the input's basename. A label — for Common Crawl it is a bare basename, not a fetchable path. | -| `source_uri` | Where warc2zip **read the file from** — the `input_file` argument, verbatim. **This is the fetch target.** | -| `warc_record_offset` / `warc_record_length` | Byte range of the record. The same triple a CDX index carries. | +| Column | Meaning | +|---------------------------------------------|---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------| +| `warc_filename` | What the WARC **calls itself** (the `warcinfo` record's `WARC-Filename`), falling back to the input's basename. A label — for Common Crawl it is a bare basename, not a fetchable path. | +| `source_uri` | Where warc2zip **read the file from** — the `input_file` argument, verbatim. **This is the fetch target.** | +| `warc_record_offset` / `warc_record_length` | Byte range of the record. The same triple a CDX index carries. | -The offsets are positions in the file named by `source_uri` — that is the only file they are guaranteed to address. `warc_filename` may name a *different* file: a derived WARC (an extract, or the output of `tools/warc_limit.py`) copies the original warcinfo record, so it keeps advertising the original WARC's name while its byte offsets refer to the derived file. Use `source_uri` to fetch, `warc_filename` to say where the records originated. See [Building and downloading a subset](#building-and-downloading-a-subset). +The offsets are positions in the file named by `source_uri` — that is the only file they are guaranteed to address. `warc_filename` may name a *different* file: a derived WARC (an extract, or the output of `tools/warc_limit.py`) copies the original `warcinfo` record, so it keeps advertising the original WARC's name while its byte offsets refer to the derived file. Use `source_uri` to fetch, `warc_filename` to say where the records originated. See [Building and downloading a subset](#building-and-downloading-a-subset). ### warcinfo.csv @@ -360,7 +377,7 @@ Crawl-level provenance, with the record body flattened the same way as `metadata `warcinfo.warc` and `warcinfo.warc-fields` hold the same record as raw wire bytes. -`source_uri` is the input this zip was converted from. It is also a column on every `manifest.csv` row; the copy here is the crawl-level one, written even when the WARC carries no warcinfo record at all, so the provenance is never lost. +`source_uri` is the input this zip was converted from. It is also a column on every `manifest.csv` row; the copy here is the crawl-level one, written even when the WARC carries no `warcinfo` record at all, so the provenance is never lost. Both `source_uri` and the `warc_record_*` fields are pseudo-headers — computed while streaming, not read off the wire. A real warcinfo header named `Source-URI` would normalize to the same name; none exists in practice, but the collision is worth knowing about, same as for `status_code`. @@ -372,50 +389,42 @@ The point of `warc_filename` + `warc_record_offset` + `warc_record_length` is th The intended workflow is: convert once with `--metadata-only` (a few tens of KB instead of gigabytes), filter the CSV however you like, then pull only the captures you kept. +### Get metadata only — no payloads ```bash -# 1. metadata only — no payloads warc2zip 'https://data.commoncrawl.org/crawl-data/.../CC-MAIN-....warc.gz' --metadata-only --output meta.zip unzip -p meta.zip '*/manifest.csv' > manifest.csv ``` -FIXME: replace 2. with grep and csv sorts of instructions +### Filter it with whatever you already use +List the column names: +```bash +head -1 manifest.csv | tr ',' '\n' | tr -d '"\r' | nl ``` -# 2. filter it with whatever you already use — here, everything that came back 200 -python - <<'EOF' -import csv -with open("manifest.csv") as fh: - rows = [r for r in csv.DictReader(fh) if r["http_status_code"] == "200"] -with open("subset.csv", "w", newline="") as fh: - w = csv.DictWriter(fh, fieldnames=rows[0].keys(), quoting=csv.QUOTE_ALL) - w.writeheader() - w.writerows(rows) -EOF -``` - -FIXME: replace 3. by expending warc2zip to do this download, with -appropriate retrying and rate limits. That might mean installing -cdx_toolkit and using it for this step. +Filter by column name and value, e.g. get all the `200` responses from the column `http_status_code`: +```bash +awk -F'","' -v col="http_status_code" -v val="200" ' +NR==1 { for (i=1; i<=NF; i++) { h=$i; gsub(/["\r]/,"",h); if (h==col) c=i } print; next } +{ v=$c; gsub(/["\r]/,"",v) } v==val +' manifest.csv > subset.csv +```` +or with [miller](https://miller.readthedocs.io/en/latest/): ``` -# 3. fetch each surviving row — every row already knows where it came from -python - <<'EOF' -import csv, urllib.request -with open("subset.csv") as fh, open("subset.warc.gz", "wb") as out: - for row in csv.DictReader(fh): - start = int(row["warc_record_offset"]) - end = start + int(row["warc_record_length"]) - 1 - req = urllib.request.Request(row["source_uri"], headers={"Range": f"bytes={start}-{end}"}) - out.write(urllib.request.urlopen(req).read()) -EOF +mlr --csv filter '$http_status_code == 200' manifest.csv > subset.csv ``` -Concatenated gzip members are themselves a valid `.warc.gz`, so appending the fetched ranges into one file produces a WARC you can feed straight back into `warc2zip` — or into any other WARC tool. +### Fetch what survived +```bash +warc2zip subset.csv --fetch --output subset.warc.gz +``` + +Every row is fetched by byte range from its own `source_uri` (https, s3, ia or a local file). Nearby rows share one request, http(s) requests go through [cdx_toolkit](https://github.com/commoncrawl/cdx_toolkit)'s Common Crawl-aware retry and pacing (`--rate`, `--retries`), and each source's own `warcinfo` record leads the output. Each record is stamped with `WARC-Source-URI` and `WARC-Source-Range`, cdx_toolkit's convention, so a re-converted subset says on every row where the record sat in the original. `subset.warc.gz` goes straight back into `warc2zip` — or into any other WARC tool. Three caveats: - Fetch with `source_uri`, not `warc_filename`. For Common Crawl the latter is a bare basename like `CC-MAIN-20260618163205-20260618193205-00999.warc.gz`; the full path is `crawl-data/{crawl}/segments/{segment}/warc/{basename}`, and **the segment is not recorded anywhere in the WARC** — `warcinfo.csv` gives you the crawl (`_body.ispartof`) but you would need the crawl's `warc.paths.gz` to resolve the rest. -- Offsets address the file named by `source_uri`, nothing else. A derived WARC keeps the original's warcinfo record, so its `warc_filename` names a file its offsets do not index. +- Offsets address the file named by `source_uri`, nothing else. A fetched subset keeps the original's `warcinfo` record, so its `warc_filename` names the original while its offsets index the subset. - Offsets stay valid under `--limit`: limiting only stops the read early, it never rewrites them. ## WARC examples for testing @@ -426,11 +435,11 @@ below read directly from that bucket. These examples are `--format flat` ... you Details: -- Example WARC from CC-MAIN-2026-25 with 500 records (response, request and metadata, 13 MBytes) [CC-MAIN-2026-30-500_records.warc.gz](https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/CC-MAIN-2026-30-500_records.warc.gz?download=true) +- Example WARC from CC-MAIN-2026-25 with 500 records (response, request and metadata, 13 MBytes) [500_RECORDS-REPACKAGE-CC-MAIN-2026-30.warc.gz](https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/500_RECORDS-REPACKAGE-CC-MAIN-2026-30.warc.gz?download=true) - make the zip ``` - warc2zip 'https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/CC-MAIN-2026-30-500_records.warc.gz?download=true' --format flat + warc2zip 'https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/500_RECORDS-REPACKAGE-CC-MAIN-2026-30.warc.gz?download=true' --format flat ``` - here is the warcinfo @@ -446,10 +455,10 @@ Details: conformsTo: https://iipc.github.io/warc-specifications/specifications/warc-format/warc-1.1/ ``` -- Homepages extracted from CC-MAIN-2026-21 (response records only, 1 GByte) [homepages_CC-MAIN-2026-21.warc.gz](https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/homepages_CC-MAIN-2026-21.warc.gz?download=true) +- Homepages extracted from CC-MAIN-2026-21 (response records only, 1 GByte) [HOMEPAGES-REPACKAGE-CC-MAIN-2026-21.warc.gz](https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/HOMEPAGES-REPACKAGE-CC-MAIN-2026-21.warc.gz?download=true) - make the zip, note the limit ``` -warc2zip 'https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/homepages_CC-MAIN-2026-21.warc.gz?download=true' --format flat --limit 1000 +warc2zip 'https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/HOMEPAGES-REPACKAGE-CC-MAIN-2026-21.warc.gz?download=true' --format flat --limit 1000 ``` - here is the warcinfo ``` @@ -460,10 +469,10 @@ warc2zip 'https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/h creator: Common Crawl Foundation operator: Malte Ostendorff ``` -- URLs of federal institutions (response records only, 1/2 GByte), as part of the [End Of Term Archive](https://eotarchive.org/) project: [is_us_federal_CC-MAIN-2025-13.warc.gz](https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/is_us_federal_CC-MAIN-2025-13.warc.gz?download=true) +- URLs of federal institutions (response records only, 1/2 GByte), as part of the [End Of Term Archive](https://eotarchive.org/) project: [IS_US_FEDERAL-REPACKAGE-CC-MAIN-2025-13.warc.gz](https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/IS_US_FEDERAL-REPACKAGE-CC-MAIN-2025-13.warc.gz?download=true) - make the zip, note the limit ``` -warc2zip 'https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/is_us_federal_CC-MAIN-2025-13.warc.gz?download=true' --format flat --limit 1000 +warc2zip 'https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/IS_US_FEDERAL-REPACKAGE-CC-MAIN-2025-13.warc.gz?download=true' --format flat --limit 1000 ``` - here is the warcinfo ``` @@ -477,30 +486,37 @@ warc2zip 'https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/i ## Many more WARC examples for testing -### Common Crawl style repackaged warcs (intended for testing) +Every archive in this section and the previous one is a test case in `tests/test_readme_warcs.py`: +`pytest` converts the first 20 captures of each (this is what CI runs), `pytest -m long` converts +the whole files (also available as the manual "Long tests" workflow under Actions). -- https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/CC-MAIN-2026-30-500_records.warc.gz?download=true (13 MBytes) -- https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/homepages_CC-MAIN-2026-21.warc.gz?download=true (1 GByte) -- https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/is_us_federal_CC-MAIN-2025-13.warc.gz?download=true (1/2 GByte) +### Common Crawl style repackaged WARCs (intended for testing) -FIXME REPACKAGE should be in the name. Are these also on s3://commoncrawl/ ? projects/warc2zip-examples ? +#### HuggingFace: +- https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/500_RECORDS-REPACKAGE-CC-MAIN-2026-30.warc.gz?download=true +- https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/HOMEPAGES-REPACKAGE-CC-MAIN-2026-21.warc.gz?download=true +- https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/IS_US_FEDERAL-REPACKAGE-CC-MAIN-2025-13.warc.gz?download=true -### Normal Common Crawl CC-MAIN warcs +#### AWS: +**Prefixes**: https://data.commoncrawl.org/ or `s3://commoncrawl/`: +- /projects/warc2zip-examples/500_RECORDS-REPACKAGE-CC-MAIN-2026-30.warc.gz (13 MBytes) +- /projects/warc2zip-examples/HOMEPAGES-REPACKAGE-CC-MAIN-2026-21.warc.gz (1 GByte) +- /projects/warc2zip-examples/IS_US_FEDERAL-REPACKAGE-CC-MAIN-2025-13.warc.gz (1/2 GByte) -- prefixes: https://data.commoncrawl.org/ or s3://commoncrawl/ +### Normal Common Crawl CC-MAIN WARCs + +**Prefixes**: https://data.commoncrawl.org/ or `s3://commoncrawl/`: - crawl-data/CC-MAIN-2026-34/segments/1786091384908.68/warc/CC-MAIN-20260807101845-20260807131845-00000.warc.gz - crawl-data/CC-MAIN-2026-34/segments/1786091384908.68/crawldiagnostics/CC-MAIN-20260807101845-20260807131845-00000.warc.gz - crawl-data/CC-MAIN-2026-34/segments/1786091384908.68/robotstxt/CC-MAIN-20260807101845-20260807131845-00000.warc.gz -FIXME: all 3 of these download to the same zip name - ### End Of Term Archive (https://eotarchive.org/data/) -- prefixes: https://eotarchive.s3.amazonaws.com/ or s3://eotarchive/ +**Prefixes**: https://eotarchive.s3.amazonaws.com/ or `s3://eotarchive/`: #### Heretrix/IA style warcs from EOT 2024 -- crawl-data/EOT-2024/segments/IA-000/EOT24PRE-20240926172119-crawl804_EOT24PRE-20240926172119-00000.warc.gz +- crawl-data/EOT-2024/segments/IA-000/warc/EOT24PRE-20240926172119-crawl804_EOT24PRE-20240926172119-00000.warc.gz #### Nutch/CCF style warcs from EOT 2024 @@ -520,6 +536,14 @@ FIXME: all 3 of these download to the same zip name - crawl-data/EOT-2004/segments/NARA-000/warc/NARA-PEOT-2004-20041014205819-00000-crawling009-c_NARA-PEOT-2004-20041014205819-00000-crawling009.archive.org.arc.gz +### Internet Archive items (`ia://`) + +A public archive.org item holding an EOT 2024 Heritrix WARC (1.6 GBytes). The two spellings read the same bytes; +the `ia://` one also works for restricted items once you are logged in (see [Usage](#usage)): + +- https://archive.org/download/EOT24PRE-20240926175758-crawl808/EOT24PRE-20240926175758-00032.warc.gz +- ia://EOT24PRE-20240926175758-crawl808/EOT24PRE-20240926175758-00032.warc.gz + ## Old CCF ARCs - prefix: s3://commoncrawl/ diff --git a/pyproject.toml b/pyproject.toml index 06d4200..76548bf 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -8,6 +8,11 @@ version = "0.1.0" description = "Convert gzipped WARC files into zip-of-zips archives" requires-python = ">=3.10" dependencies = [ + "cdx_toolkit>=0.9.39", + # ia:// needs fsspec.implementations.ia, not in a fsspec release yet. It cannot be a git dependency + # here: s3fs (via warcio[s3]) pins fsspec to its own exact release, which a branch build never + # matches. Until it ships, the branch behind the upstream PR is installed *on top* (README, CI). + "fsspec[http]", "tqdm>=4.67.3", "warcio[s3]==1.8.1", ] @@ -21,3 +26,12 @@ warc2zip = "warc2zip:cli" [tool.ruff] target-version = "py310" line-length = 120 + +[tool.pytest.ini_options] +# `long` streams every README example WARC end to end (0.4–10 GB each) and is opt-in: +# `pytest -m long`. The default run keeps the `--limit` tier, which CI runs. +addopts = "-m 'not long'" +markers = [ + "network: downloads real archives listed in the README (WARC2ZIP_OFFLINE=1 skips)", + "long: full-file conversions of the README archives; hours, never in CI", +] diff --git a/tests/conftest.py b/tests/conftest.py new file mode 100644 index 0000000..38d5bfa --- /dev/null +++ b/tests/conftest.py @@ -0,0 +1,198 @@ +"""Shared fixtures: the synthetic three-capture WARC used by the end-to-end and --fetch tests.""" + +import io +from pathlib import Path + +import pytest +from warcio.statusandheaders import StatusAndHeaders +from warcio.warcwriter import WARCWriter + +# The shape CC writes into a metadata record's warc-fields body. +CLD2 = '{"reliable":true,"languages":[{"code":"zh","text-covered":0.87,"name":"Chinese"}]}' + +CAPTURES = [ + ( + "https://example.com/", + "text/html", + b"hello", + [("Content-Type", "text/html; charset=UTF-8"), ("Connection", "close\x00")], + ), + ( + "https://cloudflare-ish.example.org/index.html", + "text/html", + b"report-to", + [ + ("Content-Type", "text/html"), + ("Report-To", '{"group":"cf-nel","max_age":604800}'), + ("Server-Timing", 'cfCacheStatus;desc="DYNAMIC"'), + ("Cache-Control", "no-store, must-revalidate, no-cache"), + ], + ), + ( + "https://plain.example.net/data.json", + "application/json", + b'{"ok": true}', + [("Content-Type", "application/json"), ("X-Fold", "a\r\n b")], + ), +] + + +@pytest.fixture +def warc_path(tmp_path): + """A three-capture WARC: warcinfo + response/request/metadata per capture.""" + path = tmp_path / "test.warc.gz" + with open(path, "wb") as fh: + writer = WARCWriter(fh, gzip=True) + writer.write_record(writer.create_warcinfo_record("test.warc.gz", {"software": "warc2zip-tests"})) + for i, (uri, _mime, payload, headers) in enumerate(CAPTURES): + request_id = f"" + + http_headers = StatusAndHeaders("200 OK", headers, protocol="HTTP/1.1") + response = writer.create_warc_record( + uri, + "response", + payload=io.BytesIO(payload), + length=len(payload), + http_headers=http_headers, + warc_headers_dict={"WARC-Concurrent-To": request_id}, + ) + writer.write_record(response) + response_id = response.rec_headers.get_header("WARC-Record-ID") + + request_headers = StatusAndHeaders( + "GET / HTTP/1.1", [("Host", "example.com"), ("User-Agent", 'cc-bot/1.0 "test"')], is_http_request=True + ) + writer.write_record( + writer.create_warc_record( + uri, + "request", + http_headers=request_headers, + warc_headers_dict={"WARC-Record-ID": request_id, "WARC-Concurrent-To": response_id}, + ) + ) + + body = ( + b"fetchTimeMs: 42\r\n" + b"charset-detected: utf-8\x00\r\n" + + f"languages-cld2: {CLD2}\r\n".encode() + + b"http-header-user-agent: cc-bot/1.0 (X11; Linux)\r\n" + b" continued-on-the-next-line\r\n" + ) + writer.write_record( + writer.create_warc_record( + uri, + "metadata", + payload=io.BytesIO(body), + length=len(body), + warc_headers_dict={ + "WARC-Concurrent-To": response_id, + "Content-Type": "application/warc-fields", + }, + ) + ) + return path + + +# --- a stand-in for archive.org ---------------------------------------------------------------- +# +# ia:// reads go to https://archive.org/download//, which 302s to a data node on +# another origin that honours Range. This server reproduces exactly that: /download/... on +# 127.0.0.1 redirects to /items/... on `localhost` (same server, different origin, so aiohttp +# applies its cross-origin rule and drops any Cookie/Authorization *headers*), and the item route +# serves byte ranges — optionally only to requests carrying a login cookie. + +import shutil +import threading +from http.server import BaseHTTPRequestHandler, HTTPServer +from types import SimpleNamespace +from urllib.parse import urlsplit + + +class _IAHandler(BaseHTTPRequestHandler): + state = None # set per server: directory, require_cookie, hits + + def log_message(self, *args): # keep pytest output clean + pass + + def _serve(self, send_body): + path = urlsplit(self.path).path + self.state.hits.append(SimpleNamespace(path=path, headers=dict(self.headers))) + if path.startswith("/download/"): + self.send_response(302) + self.send_header("Location", f"http://localhost:{self.server.server_port}/items/{path[len('/download/'):]}") + self.send_header("Content-Length", "0") + self.end_headers() + return + if not path.startswith("/items/"): + self.send_error(404) + return + required = self.state.require_cookie + if required and f"{required[0]}={required[1]}" not in self.headers.get("Cookie", ""): + self.send_error(403) + return + target = self.state.directory / path[len("/items/") :] + if not target.is_file(): + self.send_error(404) + return + data = target.read_bytes() + start, end = 0, len(data) - 1 + status = 200 + range_header = self.headers.get("Range") + if range_header and range_header.startswith("bytes="): + first, _, last = range_header[len("bytes=") :].partition("-") + start = int(first) + end = min(int(last), len(data) - 1) if last else len(data) - 1 + if start >= len(data): + self.send_response(416) + self.send_header("Content-Range", f"bytes */{len(data)}") + self.send_header("Content-Length", "0") + self.end_headers() + return + status = 206 + body = data[start : end + 1] + self.send_response(status) + self.send_header("Content-Type", "application/octet-stream") + self.send_header("Accept-Ranges", "bytes") + self.send_header("Content-Length", str(len(body))) + if status == 206: + self.send_header("Content-Range", f"bytes {start}-{end}/{len(data)}") + self.end_headers() + if send_body: + self.wfile.write(body) + + def do_GET(self): + self._serve(send_body=True) + + def do_HEAD(self): + self._serve(send_body=False) + + +@pytest.fixture +def ia_server(tmp_path, monkeypatch): + """A local archive.org: `download_url` and `cookie_domain` of InternetArchiveFileSystem are + pointed at it for the test. `.add(item, path)` publishes a file; `.require_cookie` gates the + data-node route; `.hits` records every request.""" + from fsspec.implementations.ia import InternetArchiveFileSystem + + directory = tmp_path / "items" + directory.mkdir() + state = SimpleNamespace(directory=directory, require_cookie=None, hits=[]) + handler = type("Handler", (_IAHandler,), {"state": state}) + server = HTTPServer(("127.0.0.1", 0), handler) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + + def add(item, path): + (directory / item).mkdir(exist_ok=True) + shutil.copy(path, directory / item / Path(path).name) + return f"ia://{item}/{Path(path).name}" + + state.add = add + state.download_url = f"http://127.0.0.1:{server.server_port}/download/" + monkeypatch.setattr(InternetArchiveFileSystem, "download_url", state.download_url) + monkeypatch.setattr(InternetArchiveFileSystem, "cookie_domain", "localhost") + try: + yield state + finally: + server.shutdown() + server.server_close() diff --git a/tests/test_end_to_end.py b/tests/test_end_to_end.py index c776624..383e45e 100644 --- a/tests/test_end_to_end.py +++ b/tests/test_end_to_end.py @@ -9,13 +9,15 @@ import json import re import zipfile +from pathlib import Path import pytest +from conftest import CAPTURES from warcio.archiveiterator import ArchiveIterator from warcio.statusandheaders import StatusAndHeaders from warcio.warcwriter import WARCWriter -from warc2zip import MANIFEST_COLUMNS, format_record_type_summary, main +from warc2zip import MANIFEST_COLUMNS, default_output_path, format_record_type_summary, main # Payload files are named {counter}{ext}. Sidecars share the payload's full name and add a second # suffix (1000000.html.request.json), so a plain endswith(".json") would confuse the two. @@ -37,92 +39,6 @@ "warcinfo_multi.csv", } -# The shape CC writes into a metadata record's warc-fields body. -CLD2 = '{"reliable":true,"languages":[{"code":"zh","text-covered":0.87,"name":"Chinese"}]}' - -CAPTURES = [ - ( - "https://example.com/", - "text/html", - b"hello", - [("Content-Type", "text/html; charset=UTF-8"), ("Connection", "close\x00")], - ), - ( - "https://cloudflare-ish.example.org/index.html", - "text/html", - b"report-to", - [ - ("Content-Type", "text/html"), - ("Report-To", '{"group":"cf-nel","max_age":604800}'), - ("Server-Timing", 'cfCacheStatus;desc="DYNAMIC"'), - ("Cache-Control", "no-store, must-revalidate, no-cache"), - ], - ), - ( - "https://plain.example.net/data.json", - "application/json", - b'{"ok": true}', - [("Content-Type", "application/json"), ("X-Fold", "a\r\n b")], - ), -] - - -@pytest.fixture -def warc_path(tmp_path): - """A three-capture WARC: warcinfo + response/request/metadata per capture.""" - path = tmp_path / "test.warc.gz" - with open(path, "wb") as fh: - writer = WARCWriter(fh, gzip=True) - writer.write_record(writer.create_warcinfo_record("test.warc.gz", {"software": "warc2zip-tests"})) - for i, (uri, _mime, payload, headers) in enumerate(CAPTURES): - request_id = f"" - - http_headers = StatusAndHeaders("200 OK", headers, protocol="HTTP/1.1") - response = writer.create_warc_record( - uri, - "response", - payload=io.BytesIO(payload), - length=len(payload), - http_headers=http_headers, - warc_headers_dict={"WARC-Concurrent-To": request_id}, - ) - writer.write_record(response) - response_id = response.rec_headers.get_header("WARC-Record-ID") - - request_headers = StatusAndHeaders( - "GET / HTTP/1.1", [("Host", "example.com"), ("User-Agent", 'cc-bot/1.0 "test"')], is_http_request=True - ) - writer.write_record( - writer.create_warc_record( - uri, - "request", - http_headers=request_headers, - warc_headers_dict={"WARC-Record-ID": request_id, "WARC-Concurrent-To": response_id}, - ) - ) - - body = ( - b"fetchTimeMs: 42\r\n" - b"charset-detected: utf-8\x00\r\n" - + f"languages-cld2: {CLD2}\r\n".encode() - + b"http-header-user-agent: cc-bot/1.0 (X11; Linux)\r\n" - b" continued-on-the-next-line\r\n" - ) - writer.write_record( - writer.create_warc_record( - uri, - "metadata", - payload=io.BytesIO(body), - length=len(body), - warc_headers_dict={ - "WARC-Concurrent-To": response_id, - "Content-Type": "application/warc-fields", - }, - ) - ) - return path - - def csv_members(zf): return [n for n in zf.namelist() if n.endswith(".csv")] @@ -169,6 +85,82 @@ def test_conversion_produces_parseable_csvs(warc_path, tmp_path, output_format): assert {entry["warc_target_uri"] for entry in manifest} == {uri for uri, _, _, _ in CAPTURES} +RUN_ID_RE = r"\d{8}T\d{6}_[0-9a-f]{4}" + + +@pytest.mark.parametrize("limit", [None, 2]) +def test_default_output_name_follows_the_root_dirs_rule(warc_path, tmp_path, monkeypatch, limit): + """No --output: the zip is {basename}_{hex}[_partial].zip, with _partial iff --limit was given. + + That is the root directory's rule, so the zip name tells you whether it holds a sample, and + the hex is the *same* one the root directory carries, so a zip on disk can be matched to the + directory it extracts to. + """ + monkeypatch.chdir(tmp_path) + + main(str(warc_path), output_path=None, limit=limit) + + suffix = "_partial" if limit else "" + zips = [p.name for p in tmp_path.glob("*.zip")] + assert len(zips) == 1 + zip_match = re.fullmatch(rf"test_([0-9a-f]{{4}}){suffix}\.zip", zips[0]) + assert zip_match, zips[0] + + with zipfile.ZipFile(tmp_path / zips[0]) as zf: + root_dirs = {n.split("/", 1)[0] for n in zf.namelist()} + assert len(root_dirs) == 1 + dir_match = re.fullmatch(rf"test_\d{{8}}T\d{{6}}_([0-9a-f]{{4}}){suffix}", next(iter(root_dirs))) + assert dir_match, root_dirs + assert dir_match.group(1) == zip_match.group(1) + + +def test_default_output_names_do_not_collide_for_same_basename(warc_path, tmp_path, monkeypatch): + """CC's warc/, crawldiagnostics/ and robotstxt/ files share a basename: two runs, two zips.""" + monkeypatch.chdir(tmp_path) + + main(str(warc_path), output_path=None) + main(str(warc_path), output_path=None) + + assert len(list(tmp_path.glob("test_*.zip"))) == 2 + + +@pytest.mark.parametrize( + ("input_file", "expected"), + [ + # Every input shape the README shows + ("archive.warc.gz", "archive.zip"), + ("/data/crawls/archive.warc.gz", "archive.zip"), + ("s3://commoncrawl/crawl-data/CC-MAIN-2026-34/segments/x/warc/CC-MAIN-0000.warc.gz", "CC-MAIN-0000.zip"), + ("https://data.commoncrawl.org/crawl-data/CC-MAIN-2026-34/warc/CC-MAIN-0000.warc.gz", "CC-MAIN-0000.zip"), + ( + ( + "https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/" + "CC-MAIN-2026-30-500_records.warc.gz?download=true" + ), + "CC-MAIN-2026-30-500_records.zip", + ), + ("https://eotarchive.s3.amazonaws.com/crawl-data/EOT-2004/NARA-PEOT-2004.arc.gz", "NARA-PEOT-2004.zip"), + ("plain.warc", "plain.zip"), + ("-", "stdin.zip"), + ], +) +def test_default_output_path_handles_every_readme_input_shape(input_file, expected): + label = expected[: -len(".zip")] + assert default_output_path(input_file, run_id="abcd").name == f"{label}_abcd.zip" + assert default_output_path(input_file, partial=True, run_id="abcd").name == f"{label}_abcd_partial.zip" + # Without a run id one is minted, and it is 4 hex chars like the root directory's + assert re.fullmatch(rf"{re.escape(label)}_[0-9a-f]{{4}}\.zip", default_output_path(input_file).name) + # Always the current directory, never the input's + assert default_output_path(input_file).parent == Path(".") + + +def test_explicit_output_path_is_used_verbatim(warc_path, tmp_path): + out = tmp_path / "chosen-name.zip" + main(str(warc_path), str(out), limit=1) + assert out.exists() + assert list(tmp_path.glob("*.zip")) == [out] + + def test_nul_from_the_wire_is_escaped_in_the_csv(warc_path, tmp_path): out = tmp_path / "flat.zip" main(str(warc_path), str(out), output_format="flat") diff --git a/tests/test_fetch.py b/tests/test_fetch.py new file mode 100644 index 0000000..daf42ed --- /dev/null +++ b/tests/test_fetch.py @@ -0,0 +1,581 @@ +"""--fetch: rebuilding a subset .warc.gz from manifest rows. + +Everything runs offline: the sources are local paths, which fsspec serves through the same +cat_file(start, end) range read that https:// and s3:// use. +""" + +import csv +import io +import shutil +import sys +import zipfile +from pathlib import Path + +import pytest +from conftest import CAPTURES +from fsspec.implementations.local import LocalFileSystem +from warcio.archiveiterator import ArchiveIterator +from warcio.statusandheaders import StatusAndHeaders +from warcio.warcwriter import WARCWriter + +import warc2zip +from warc2zip import ( + FETCH_MAX_BACKOFF, + FetchLengthMismatch, + FetchRow, + RateLimiter, + annotate_record_slice, + cli, + coalesce_ranges, + default_fetch_output_path, + fetch_main, + fetch_with_retry, + http_range, + leading_warcinfo, + main, + set_host_interval, + validate_record_slice, +) + + +def zip_csv(zip_path, suffix): + with zipfile.ZipFile(zip_path) as zf: + name = next(n for n in zf.namelist() if n.endswith(suffix)) + return list(csv.reader(io.StringIO(zf.read(name).decode("utf-8")))) + + +def manifest_rows(zip_path): + header, *rows = zip_csv(zip_path, "manifest.csv") + return [dict(zip(header, row)) for row in rows] + + +def warcinfo_range(zip_path): + values = {name: value for key, name, value in zip_csv(zip_path, "warcinfo.csv")[1:] if key == "warcinfo"} + return int(values["warc_record_offset"]), int(values["warc_record_length"]) + + +def request_range(zip_path, filename): + """(offset, length) of the request record joined to a payload file.""" + values = {name: value for key, name, value in zip_csv(zip_path, "request_warc_headers.csv")[1:] if key == filename} + return int(values["warc_record_offset"]), int(values["warc_record_length"]) + + +def slice_of(raw, row): + offset, length = int(row["warc_record_offset"]), int(row["warc_record_length"]) + return raw[offset : offset + length] + + +def fetch_row(row, **overrides): + values = { + "source_uri": row["source_uri"], + "offset": int(row["warc_record_offset"]), + "length": int(row["warc_record_length"]), + "record_id": row["warc_record_id"], + "target_uri": row["warc_target_uri"], + "line": 2, + } + values.update(overrides) + return FetchRow(**values) + + +def write_subset(path, rows, fieldnames=None): + with open(path, "w", newline="", encoding="utf-8") as fh: + writer = csv.DictWriter(fh, fieldnames=fieldnames or list(rows[0].keys()), quoting=csv.QUOTE_ALL) + writer.writeheader() + writer.writerows(rows) + return path + + +def record_types_and_ids(path): + with open(path, "rb") as fh: + return [(r.rec_type, r.rec_headers.get_header("WARC-Record-ID")) for r in ArchiveIterator(fh)] + + +def parsed_records(data): + """(rec_type, WARC headers, payload) per record; the payload is read before advancing.""" + return [ + (r.rec_type, dict(r.rec_headers.headers), r.content_stream().read()) + for r in ArchiveIterator(io.BytesIO(data)) + ] + + +def without_source_headers(headers): + return {k: v for k, v in headers.items() if not k.startswith("WARC-Source-")} + + +def raw_blocks(data): + """Each record's content block as written, HTTP headers and body, without warcio's HTTP parse.""" + return [r.content_stream().read() for r in ArchiveIterator(io.BytesIO(data), no_record_parse=True)] + + +def assert_is_stamped_copy(fetched, original, source_uri, row): + """A fetched record is the original plus the two provenance headers, same payload.""" + rec_type, headers, payload = fetched + o, n = int(row["warc_record_offset"]), int(row["warc_record_length"]) + assert rec_type == "response" + assert without_source_headers(headers) == original[1] + assert headers["WARC-Source-URI"] == source_uri + assert headers["WARC-Source-Range"] == f"bytes={o}-{o + n - 1}" + assert payload == original[2] + + +def assert_blocks_untouched(fetched_member, original_member): + """The stamped record's content block is the wire bytes: obs-folds, NULs and all.""" + assert raw_blocks(fetched_member) == raw_blocks(original_member) + + +def write_warc_without_warcinfo(path): + with open(path, "wb") as fh: + writer = WARCWriter(fh, gzip=True) + payload = b"bare" + writer.write_record( + writer.create_warc_record( + "https://bare.example.com/", + "response", + payload=io.BytesIO(payload), + length=len(payload), + http_headers=StatusAndHeaders("200 OK", [("Content-Type", "text/html")], protocol="HTTP/1.1"), + ) + ) + return path + + +@pytest.fixture +def metadata_zip(warc_path, tmp_path): + out = tmp_path / "meta.zip" + assert main(str(warc_path), str(out), metadata_only=True) == 0 + return out + + +# --- end to end --------------------------------------------------------------------------- + + +def test_fetch_rebuilds_the_subset(warc_path, metadata_zip, tmp_path): + """The source's warcinfo member verbatim, then each kept row's record stamped with its origin.""" + rows = manifest_rows(metadata_zip) + kept = [rows[0], rows[2]] + subset = write_subset(tmp_path / "subset.csv", kept) + out = tmp_path / "subset.warc.gz" + + assert fetch_main(str(subset), str(out)) == 0 + + raw = warc_path.read_bytes() + offset, length = warcinfo_range(metadata_zip) + fetched = out.read_bytes() + assert fetched.startswith(raw[offset : offset + length]) + records = parsed_records(fetched) + assert [t for t, _, _ in records] == ["warcinfo", "response", "response"] + for record, row in zip(records[1:], kept): + assert_is_stamped_copy(record, parsed_records(slice_of(raw, row))[0], str(warc_path), row) + + +def test_fetched_subset_reconverts_and_still_names_the_original_warc(warc_path, metadata_zip, tmp_path): + """The README caveat: a derived WARC keeps the original's warcinfo, while its offsets index itself.""" + rows = manifest_rows(metadata_zip) + subset = write_subset(tmp_path / "subset.csv", [rows[1]]) + out = tmp_path / "subset.warc.gz" + assert fetch_main(str(subset), str(out)) == 0 + + check = tmp_path / "check.zip" + assert main(str(out), str(check)) == 0 + check_rows = manifest_rows(check) + + assert [row["warc_record_id"] for row in check_rows] == [rows[1]["warc_record_id"]] + assert check_rows[0]["warc_filename"] == "test.warc.gz" + assert check_rows[0]["source_uri"] == str(out) + record = next(iter(ArchiveIterator(io.BytesIO(slice_of(out.read_bytes(), check_rows[0]))))) + assert record.rec_headers.get_header("WARC-Record-ID") == rows[1]["warc_record_id"] + # ...and the CSV now says where the record sat in the original. + o, n = int(rows[1]["warc_record_offset"]), int(rows[1]["warc_record_length"]) + provenance = {name: value for _, name, value in zip_csv(check, "response_warc_headers.csv")[1:]} + assert provenance["warc_source_uri"] == str(warc_path) + assert provenance["warc_source_range"] == f"bytes={o}-{o + n - 1}" + + +def test_rows_are_grouped_by_source_and_ordered_by_offset(warc_path, metadata_zip, tmp_path): + """Two sources: each gets its own warcinfo, in first-appearance order, rows sorted within.""" + other = tmp_path / "other.warc.gz" + shutil.copy(warc_path, other) + rows = manifest_rows(metadata_zip) + from_other = [dict(row, source_uri=str(other)) for row in rows] + # CSV order deliberately scrambled: other first, then this file's rows descending. + subset = write_subset(tmp_path / "subset.csv", [from_other[1], rows[2], rows[0], from_other[0]]) + out = tmp_path / "subset.warc.gz" + + assert fetch_main(str(subset), str(out)) == 0 + + warcinfo_id = record_types_and_ids(warc_path)[0][1] + assert record_types_and_ids(out) == [ + ("warcinfo", warcinfo_id), + ("response", rows[0]["warc_record_id"]), + ("response", rows[1]["warc_record_id"]), + ("warcinfo", warcinfo_id), + ("response", rows[0]["warc_record_id"]), + ("response", rows[2]["warc_record_id"]), + ] + + +def test_source_without_warcinfo_yields_responses_only(tmp_path, capsys): + bare = write_warc_without_warcinfo(tmp_path / "bare.warc.gz") + meta = tmp_path / "meta.zip" + assert main(str(bare), str(meta), metadata_only=True) == 0 + subset = write_subset(tmp_path / "subset.csv", manifest_rows(meta)) + out = tmp_path / "subset.warc.gz" + + assert fetch_main(str(subset), str(out)) == 0 + + assert [t for t, _ in record_types_and_ids(out)] == ["response"] + assert "no leading warcinfo record" in capsys.readouterr().err + + +def test_blank_offset_row_is_skipped_with_a_warning(metadata_zip, tmp_path, capsys): + rows = manifest_rows(metadata_zip) + rows[1]["warc_record_offset"] = "" + subset = write_subset(tmp_path / "subset.csv", rows) + out = tmp_path / "subset.warc.gz" + + assert fetch_main(str(subset), str(out)) == 1 + + err = capsys.readouterr().err + assert "line 3" in err and "1 row(s) could not be fetched" in err + assert [t for t, _ in record_types_and_ids(out)] == ["warcinfo", "response", "response"] + + +def test_duplicate_rows_are_fetched_once(metadata_zip, tmp_path, capsys): + rows = manifest_rows(metadata_zip) + subset = write_subset(tmp_path / "subset.csv", [rows[0], rows[0]]) + out = tmp_path / "subset.warc.gz" + + assert fetch_main(str(subset), str(out)) == 0 + + assert [t for t, _ in record_types_and_ids(out)] == ["warcinfo", "response"] + assert "1 duplicate row(s) dropped" in capsys.readouterr().err + + +def test_a_server_that_ignores_range_is_caught_not_retried(metadata_zip, tmp_path, monkeypatch, capsys): + """fsspec never checks for a 206, so the length check is the only guard against a whole-file answer.""" + monkeypatch.setattr(LocalFileSystem, "cat_file", lambda self, path, start=None, end=None, **kw: Path(path).read_bytes()) + subset = write_subset(tmp_path / "subset.csv", manifest_rows(metadata_zip)[:1]) + out = tmp_path / "subset.warc.gz" + sleeps = [] + + assert fetch_main(str(subset), str(out), sleep=sleeps.append) == 1 + + assert sleeps == [] + assert "got" in capsys.readouterr().err + assert [t for t, _ in record_types_and_ids(out)] == ["warcinfo"] + + +def test_csv_without_manifest_columns_is_a_usage_error(tmp_path, capsys): + subset = write_subset(tmp_path / "subset.csv", [{"url": "x", "warc_record_offset": "0"}]) + with pytest.raises(SystemExit) as exc: + fetch_main(str(subset), str(tmp_path / "out.warc.gz")) + assert exc.value.code == 2 + assert "source_uri" in capsys.readouterr().err + + +# --- pure pieces -------------------------------------------------------------------------- + + +def row_at(offset, length): + return FetchRow("src", offset, length, "", "", 0) + + +def test_coalesce_ranges_merges_near_rows_and_splits_on_gap_or_span(): + a, b, c = row_at(0, 100), row_at(150, 100), row_at(10_000, 50) + assert coalesce_ranges([a, b, c], max_gap=100, max_span=10**6) == [(0, 250, [a, b]), (10_000, 10_050, [c])] + assert coalesce_ranges([a, b], max_gap=10, max_span=10**6) == [(0, 100, [a]), (150, 250, [b])] + assert coalesce_ranges([a, b], max_gap=100, max_span=200) == [(0, 100, [a]), (150, 250, [b])] + assert coalesce_ranges([a]) == [(0, 100, [a])] + assert coalesce_ranges([]) == [] + + +class FakeHTTPError(Exception): + """Shaped like aiohttp.ClientResponseError: a .status and .headers.""" + + def __init__(self, status, headers=None): + super().__init__(f"HTTP {status}") + self.status = status + self.headers = headers or {} + + +def failing(exc): + def fetch(): + raise exc + + return fetch + + +def test_fetch_with_retry_honours_retry_after_then_succeeds(capsys): + outcomes = [FakeHTTPError(503, {"Retry-After": "3"}), FakeHTTPError(429, {"Retry-After": "3"}), b"ok"] + + def fetch(): + outcome = outcomes.pop(0) + if isinstance(outcome, Exception): + raise outcome + return outcome + + sleeps = [] + assert fetch_with_retry(fetch, retries=3, label="x", sleep=sleeps.append) == b"ok" + assert sleeps == [3.0, 3.0] + assert capsys.readouterr().err.count("retrying in 3.0 s") == 2 + + +def test_fetch_with_retry_backs_off_exponentially_then_gives_up(): + sleeps = [] + with pytest.raises(FakeHTTPError): + fetch_with_retry(failing(FakeHTTPError(503)), retries=3, sleep=sleeps.append, rng=lambda: 0.5) + assert sleeps == [2.0, 4.0, 8.0] # 2**attempt * (0.5 + rng) + + +def test_fetch_with_retry_caps_the_backoff(): + sleeps = [] + with pytest.raises(FakeHTTPError): + fetch_with_retry(failing(FakeHTTPError(503)), retries=8, sleep=sleeps.append, rng=lambda: 0.5) + assert max(sleeps) == FETCH_MAX_BACKOFF + + +@pytest.mark.parametrize( + "exc", + [FileNotFoundError("404"), PermissionError("403"), FakeHTTPError(416), FakeHTTPError(400), FetchLengthMismatch("short")], +) +def test_deterministic_errors_are_not_retried(exc): + calls = [] + + def fetch(): + calls.append(1) + raise exc + + with pytest.raises(type(exc)): + fetch_with_retry(fetch, retries=5, sleep=lambda s: pytest.fail("slept")) + assert len(calls) == 1 + + +@pytest.mark.parametrize("exc", [ConnectionResetError("reset"), TimeoutError("timeout"), OSError("eio")]) +def test_connection_errors_are_retried(exc): + sleeps = [] + with pytest.raises(type(exc)): + fetch_with_retry(failing(exc), retries=1, sleep=sleeps.append) + assert len(sleeps) == 1 + + +def test_rate_limiter_spaces_requests_per_key(): + clock = [100.0] + sleeps = [] + + def sleep(seconds): + sleeps.append(seconds) + clock[0] += seconds + + limiter = RateLimiter(2.0, clock=lambda: clock[0], sleep=sleep) + limiter.wait("a") + limiter.wait("b") + assert sleeps == [] + limiter.wait("a") + assert sleeps == [0.5] + clock[0] += 10 + limiter.wait("a") + assert sleeps == [0.5] + + +def test_rate_limiter_zero_disables(): + limiter = RateLimiter(0, clock=lambda: 0.0, sleep=lambda s: pytest.fail("slept")) + limiter.wait("a") + limiter.wait("a") + + +def test_leading_warcinfo_cuts_the_first_member_only_when_it_is_complete(warc_path, metadata_zip, tmp_path): + raw = warc_path.read_bytes() + offset, length = warcinfo_range(metadata_zip) + assert offset == 0 + assert leading_warcinfo(raw[: 64 * 1024]) == raw[:length] + assert leading_warcinfo(raw[: length + 10]) == raw[:length] # probe ends inside the second member + assert leading_warcinfo(raw[: length - 5]) is None # probe ends inside the warcinfo itself + assert leading_warcinfo(b"not a warc") is None + assert leading_warcinfo(b"") is None + bare = write_warc_without_warcinfo(tmp_path / "bare.warc.gz") + assert leading_warcinfo(bare.read_bytes()) is None + + +def test_validate_record_slice_rejects_impostors(warc_path, metadata_zip): + rows = manifest_rows(metadata_zip) + raw = warc_path.read_bytes() + row = fetch_row(rows[0]) + good = slice_of(raw, rows[0]) + + assert validate_record_slice(good, row) is None + assert validate_record_slice(b"not a warc at all", row) # warcio sniffs this as an ARC response + assert "WARC-Record-ID" in validate_record_slice(slice_of(raw, rows[1]), row) + assert validate_record_slice(good + b"trailing garbage bytes", row) + assert validate_record_slice(good[:-1], row) + offset, length = request_range(metadata_zip, rows[0]["filename"]) + assert "request" in validate_record_slice(raw[offset : offset + length], row) + + # ARC rows carry no record id, so the target URI is what identifies the record. + assert validate_record_slice(good, fetch_row(rows[0], record_id="")) is None + assert "WARC-Target-URI" in validate_record_slice(good, fetch_row(rows[0], record_id="", target_uri="https://x/")) + + +def test_default_fetch_output_path(): + assert default_fetch_output_path("subset.csv", run_id="abcd") == Path("subset_abcd.warc.gz") + assert default_fetch_output_path("/tmp/dir/Manifest.CSV", run_id="abcd") == Path("Manifest_abcd.warc.gz") + assert default_fetch_output_path("-", run_id="abcd") == Path("stdin_abcd.warc.gz") + assert default_fetch_output_path("subset.csv").parent == Path(".") + + +# --- cli ---------------------------------------------------------------------------------- + + +@pytest.mark.parametrize( + "argv", + [ + ["subset.csv", "--fetch", "--limit", "1"], + ["subset.csv", "--fetch", "--format", "sidecar"], + ["x.warc.gz", "--rate", "1"], + ["x.warc.gz", "--retries", "2"], + ["subset.csv", "--fetch", "--dry-run"], + ], +) +def test_cli_refuses_flags_that_would_otherwise_be_ignored(argv, monkeypatch): + monkeypatch.setattr(sys, "argv", ["warc2zip", *argv]) + with pytest.raises(SystemExit) as exc: + cli() + assert exc.value.code == 2 + + +def test_cli_fetch_and_default_format_still_flat(warc_path, metadata_zip, tmp_path, monkeypatch): + subset = write_subset(tmp_path / "subset.csv", manifest_rows(metadata_zip)[:1]) + out = tmp_path / "subset.warc.gz" + monkeypatch.setattr(sys, "argv", ["warc2zip", str(subset), "--fetch", "--output", str(out), "--rate", "0"]) + assert cli() == 0 + assert [t for t, _ in record_types_and_ids(out)] == ["warcinfo", "response"] + + check = tmp_path / "check.zip" + monkeypatch.setattr(sys, "argv", ["warc2zip", str(out), "--output", str(check)]) + assert cli() == 0 + with zipfile.ZipFile(check) as zf: + names = {n.rsplit("/", 1)[-1] for n in zf.namelist()} + assert "manifest.csv" in names and not any(".request." in n for n in names) # flat, not sidecar + assert len(CAPTURES) == 3 # the fixture shape the row indexes above rely on + + +# --- http(s) transport through cdx_toolkit ------------------------------------------------ + + +class FakeResponse: + def __init__(self, content=b"", status_code=200, headers=None): + self.content = content + self.status_code = status_code + self.headers = headers or {} + + +def fake_myrequests_get(warc_bytes, log, redirect_from=None, ignore_range=False): + """A `myrequests_get` stand-in serving Range requests for `warc_bytes`, optionally after a 302.""" + + def get(url, params=None, headers=None, **kwargs): + log.append((url, headers["Range"], kwargs)) + if url == redirect_from: + return FakeResponse(b"moved", 302, {"Location": "/cdn/x.warc.gz"}) + if ignore_range: + return FakeResponse(warc_bytes, 200) + start, end = (int(v) for v in headers["Range"][len("bytes=") :].split("-")) + return FakeResponse(warc_bytes[start : end + 1], 206) + + return get + + +def https_rows(rows, uri): + return [dict(row, source_uri=uri) for row in rows] + + +def test_https_sources_go_through_cdx_toolkit(warc_path, metadata_zip, tmp_path, monkeypatch): + log = [] + monkeypatch.setattr(warc2zip, "myrequests_get", fake_myrequests_get(warc_path.read_bytes(), log)) + uri = "https://data.example.org/crawl-data/x.warc.gz" + rows = manifest_rows(metadata_zip) + subset = write_subset(tmp_path / "subset.csv", https_rows([rows[0], rows[2]], uri)) + out = tmp_path / "subset.warc.gz" + + assert fetch_main(str(subset), str(out), retries=3) == 0 + + raw = warc_path.read_bytes() + offset, length = warcinfo_range(metadata_zip) + fetched = out.read_bytes() + assert fetched.startswith(raw[offset : offset + length]) + records = parsed_records(fetched) + assert [t for t, _, _ in records] == ["warcinfo", "response", "response"] + for record, row in zip(records[1:], [rows[0], rows[2]]): + assert_is_stamped_copy(record, parsed_records(slice_of(raw, row))[0], uri, row) + # One probe plus one coalesced span, both Range requests to cdx_toolkit's transport. + assert [(u, r) for u, r, _ in log] == [ + (uri, "bytes=0-65535"), + (uri, f"bytes={rows[0]['warc_record_offset']}-{int(rows[2]['warc_record_offset']) + int(rows[2]['warc_record_length']) - 1}"), + ] + assert all(kw == {"raise_error_after_n_errors": 3} for _, _, kw in log) + + +def test_http_range_follows_redirects_by_hand(warc_path): + """myrequests_get passes allow_redirects=False, so a 302 would otherwise come back as the slice.""" + log = [] + raw = warc_path.read_bytes() + get = fake_myrequests_get(raw, log, redirect_from="https://hub.example/resolve/x.warc.gz?download=true") + + data = http_range("https://hub.example/resolve/x.warc.gz?download=true", 10, 20, get=get) + + assert data == raw[10:20] + assert [u for u, _, _ in log] == ["https://hub.example/resolve/x.warc.gz?download=true", "https://hub.example/cdn/x.warc.gz"] + assert {r for _, r, _ in log} == {"bytes=10-19"} + + +def test_http_range_gives_up_on_a_redirect_loop(warc_path): + log = [] + get = fake_myrequests_get(warc_path.read_bytes(), log, redirect_from="https://hub.example/cdn/x.warc.gz") + with pytest.raises(RuntimeError, match="redirects"): + http_range("https://hub.example/cdn/x.warc.gz", 0, 10, redirects=2, get=get) + assert len(log) == 3 + + +def test_https_server_ignoring_range_is_caught(warc_path, metadata_zip, tmp_path, monkeypatch, capsys): + log = [] + monkeypatch.setattr(warc2zip, "myrequests_get", fake_myrequests_get(warc_path.read_bytes(), log, ignore_range=True)) + rows = https_rows(manifest_rows(metadata_zip)[:1], "https://data.example.org/x.warc.gz") + subset = write_subset(tmp_path / "subset.csv", rows) + out = tmp_path / "subset.warc.gz" + + assert fetch_main(str(subset), str(out)) == 1 + + assert "got" in capsys.readouterr().err + assert [t for t, _ in record_types_and_ids(out)] == ["warcinfo"] # the probe tolerates a long answer + + +def test_set_host_interval_maps_rate_onto_cdx_toolkit(monkeypatch): + from cdx_toolkit.myrequests import retry_info + + host = "data.example.org" + monkeypatch.delitem(retry_info, host, raising=False) + set_host_interval(host, None) + assert host not in retry_info # None leaves cdx_toolkit's own defaults alone + set_host_interval(host, 4) + assert retry_info[host]["minimum_interval"] == 0.25 + set_host_interval(host, 0) + assert retry_info[host]["minimum_interval"] == 0.0 + assert retry_info["data.commoncrawl.org"]["minimum_interval"] == 0.55 + + +def test_stamping_keeps_every_original_header_and_the_payload(warc_path, metadata_zip): + rows = manifest_rows(metadata_zip) + raw = warc_path.read_bytes() + row = rows[1] + original = parsed_records(slice_of(raw, row))[0] + + stamped = annotate_record_slice(slice_of(raw, row), str(warc_path), int(row["warc_record_offset"]), int(row["warc_record_length"])) + + records = parsed_records(stamped) + assert len(records) == 1 + assert_is_stamped_copy(records[0], original, str(warc_path), row) + assert_blocks_untouched(stamped, slice_of(raw, row)) + # The fixture's obs-folded header is the case a parsed re-serialisation would unfold. + folded = annotate_record_slice(slice_of(raw, rows[2]), str(warc_path), 0, 1) + assert b"X-Fold: a\r\n b" in raw_blocks(folded)[0] + assert records[0][1]["WARC-Target-URI"] == CAPTURES[1][0] + assert records[0][2] == CAPTURES[1][2] diff --git a/tests/test_ia.py b/tests/test_ia.py new file mode 100644 index 0000000..04c1e2f --- /dev/null +++ b/tests/test_ia.py @@ -0,0 +1,137 @@ +"""ia://item/file inputs through fsspec's InternetArchiveFileSystem. + +The class itself (URI mapping, ia.ini parsing, cookie-jar scoping) is tested in fsspec; what is +tested here is warc2zip's use of it: that the installed fsspec provides the protocol, a conversion +and a --fetch through it against the `ia_server` fixture (a stand-in archive.org whose /download/ +URL 302s to a data node on another origin, serving byte ranges, optionally only to a logged-in +cookie), and the CLI's explanation of a refused item. The one network test lives in +test_readme_warcs.py. +""" + +import pytest +from conftest import CAPTURES +from fsspec.implementations.ia import InternetArchiveFileSystem, load_ia_credentials +from test_fetch import assert_is_stamped_copy, manifest_rows, parsed_records, slice_of, warcinfo_range, write_subset + +import warc2zip +from warc2zip import ( + cli, + default_output_path, + fetch_main, + format_ia_permission_error, + input_basename, + main, +) + +INI = """[s3] +access = AKIA-TEST +secret = s3cr3t +[cookies] +logged-in-user = someone%40example.org; expires=Sat, 28-Aug-2027 19:39:52 GMT; Max-Age=31536000; path=/; domain=.archive.org +logged-in-sig = 1756000000-abcdef0123456789; expires=Sat, 28-Aug-2027 19:39:52 GMT; Max-Age=31536000; path=/; domain=.archive.org +[general] +screenname = Some One +""" + + +@pytest.fixture(autouse=True) +def isolated_credentials(tmp_path, monkeypatch): + """Never read the developer's real ia.ini, and never reuse a cached filesystem (fsspec caches + instances by constructor arguments, so a no-argument instance would carry the previous test's + credentials).""" + monkeypatch.setenv("HOME", str(tmp_path / "home")) + for name in ("IA_CONFIG_FILE", "XDG_CONFIG_HOME", "IA_ACCESS_KEY_ID", "IA_SECRET_ACCESS_KEY"): + monkeypatch.delenv(name, raising=False) + InternetArchiveFileSystem.clear_instance_cache() + yield + InternetArchiveFileSystem.clear_instance_cache() + + +def write_ini(path, text=INI): + path.write_text(text, encoding="utf-8") + return path + + +# --- the contract with fsspec --------------------------------------------------------------- + + +def test_ia_is_registered_with_fsspec(): + """The installed fsspec must provide the protocol: nothing in warc2zip registers it any more.""" + import fsspec + + fs, path = fsspec.core.url_to_fs("ia://item/file.warc.gz") + assert isinstance(fs, InternetArchiveFileSystem) + assert path == "https://archive.org/download/item/file.warc.gz" + + +def test_ia_input_names_the_zip_like_any_other_uri(): + assert input_basename("ia://item/CRAWL-00032.warc.gz") == "CRAWL-00032.warc.gz" + assert str(default_output_path("ia://item/CRAWL-00032.warc.gz", run_id="a1b2")) == "CRAWL-00032_a1b2.zip" + + +def test_ia_uri_converts_and_refetches(ia_server, warc_path, tmp_path): + """main() streams the file by ranges through the redirect; the manifest carries the ia:// + URI verbatim, so --fetch can go back to it.""" + uri = ia_server.add("testitem", warc_path) + out = tmp_path / "out.zip" + assert main(uri, str(out)) == 0 + + rows = manifest_rows(out) + assert [row["warc_target_uri"] for row in rows] == [capture[0] for capture in CAPTURES] + assert {row["source_uri"] for row in rows} == {uri} + item_reads = [hit for hit in ia_server.hits if hit.path.startswith("/items/") and "Range" in hit.headers] + assert item_reads, "data-node reads should be range requests" + assert all(hit.path.startswith(("/download/", "/items/")) for hit in ia_server.hits) + + subset = write_subset(tmp_path / "subset.csv", [rows[0], rows[2]]) + fetched_path = tmp_path / "subset.warc.gz" + assert fetch_main(str(subset), str(fetched_path)) == 0 + raw = warc_path.read_bytes() + offset, length = warcinfo_range(out) + fetched = fetched_path.read_bytes() + assert fetched.startswith(raw[offset : offset + length]) + records = parsed_records(fetched) + assert [t for t, _, _ in records] == ["warcinfo", "response", "response"] + for record, row in zip(records[1:], [rows[0], rows[2]]): + assert_is_stamped_copy(record, parsed_records(slice_of(raw, row))[0], uri, row) + + +def test_logged_in_user_converts_a_restricted_item(ia_server, warc_path, tmp_path, monkeypatch): + """The issue's acceptance: anonymous, a restricted item is refused; once ia.ini holds the + account's cookies it converts, with nothing passed on the command line. (How the cookies + reach the data node behind the cross-origin redirect is fsspec's business and tested there.)""" + uri = ia_server.add("restricted", warc_path) + ia_server.require_cookie = ("logged-in-sig", "1756000000-abcdef0123456789") + with pytest.raises(PermissionError): + main(uri, str(tmp_path / "anonymous.zip")) + + monkeypatch.setenv("IA_CONFIG_FILE", str(write_ini(tmp_path / "ia.ini"))) + InternetArchiveFileSystem.clear_instance_cache() + out = tmp_path / "out.zip" + assert main(uri, str(out)) == 0 + assert len(manifest_rows(out)) == len(CAPTURES) + + +def test_cli_explains_a_refused_ia_item(ia_server, warc_path, tmp_path, monkeypatch, capsys): + uri = ia_server.add("restricted", warc_path) + ia_server.require_cookie = ("logged-in-sig", "nope") + monkeypatch.setattr("sys.argv", ["warc2zip", uri, "--output", str(tmp_path / "out.zip")]) + assert cli() == 1 + err = capsys.readouterr().err + assert "refused access to item 'restricted'" in err + assert "anonymous" in err and "ia configure" in err + + +def test_permission_message_names_the_credentials_it_used(tmp_path): + anonymous = format_ia_permission_error("ia://item/x.warc.gz", load_ia_credentials(str(tmp_path / "none.ini"))) + assert "anonymous" in anonymous and "no ia.ini" in anonymous + logged_in = format_ia_permission_error("ia://item/x.warc.gz", load_ia_credentials(str(write_ini(tmp_path / "ia.ini")))) + assert "were refused" in logged_in and str(tmp_path / "ia.ini") in logged_in + assert "https://archive.org/download/item/x.warc.gz" in logged_in + + +def test_permission_error_elsewhere_is_not_blamed_on_archive_org(tmp_path, monkeypatch): + monkeypatch.setattr(warc2zip, "main", lambda *a, **k: (_ for _ in ()).throw(PermissionError("disk"))) + monkeypatch.setattr("sys.argv", ["warc2zip", "local.warc.gz"]) + with pytest.raises(PermissionError, match="disk"): + cli() diff --git a/tests/test_readme_warcs.py b/tests/test_readme_warcs.py new file mode 100644 index 0000000..c3a964f --- /dev/null +++ b/tests/test_readme_warcs.py @@ -0,0 +1,233 @@ +"""Crash tests over the real archives listed in the README's "WARC examples for testing". + +Every archive there was written by a different tool (Nutch, cdx_toolkit repackaging, Heritrix, +Browsertrix, ArchiveTeam's megawarc, 2004/2008-era ARC writers), and each one has broken warc2zip +at some point in a way the synthetic fixtures could not: a producer that writes no metadata +records, one that writes only ``resource`` records, one with four warcinfo records, one with +``dns:`` captures and no HTTP layer. The check is deliberately minimal — exit 0, a valid zip, and +the manifest agreeing with the zip members — because the files are the point, not the assertions. + +Two tiers, chosen with pytest markers (registered in ``pyproject.toml``): + +- **short** — the default and what CI runs. ``--limit`` keeps every run to a few MB over HTTP + (~3 s each), so all fourteen archives in both formats finish in about a minute. +- **long** — ``pytest -m long``. Streams each archive end to end: 0.4–1 GB each, and the + ArchiveTeam megawarc is 10 GB, so this is an hour-plus run for a developer's machine, never CI. + ``-k flat`` or ``-k sidecar`` halves it. Every zip is deleted once checked, so the disk high-water + mark is one output at a time rather than ~20 GB under pytest's tmp dir. + +Both tiers need the network. ``WARC2ZIP_OFFLINE=1`` skips them. +""" + +import csv +import io +import os +import posixpath +import subprocess +import sys +import zipfile +from dataclasses import dataclass + +import pytest + +# Captures per run in the short tier. Small enough that fsspec's first HTTP block (5 MiB) usually +# covers it, large enough that the drain-until-next-capture logic and the request/metadata joins +# see more than one group. +SHORT_LIMIT = 20 + +# Wall-clock ceiling for one short run: a hung connection must fail CI, not stall it. +SHORT_TIMEOUT = 300 + +HF = "https://huggingface.co/buckets/commoncrawl/warc2zip-examples/resolve/" +CC = "https://data.commoncrawl.org/" +EOT = "https://eotarchive.s3.amazonaws.com/crawl-data/" +# An archive.org item, through the ia:// scheme (see InternetArchiveFileSystem). +IA = "ia://" + + +@dataclass(frozen=True) +class Example: + name: str # pytest id + url: str + producer: str # what wrote the archive, for the failure message + size: str # from the README, for humans + # Total captures in the file when it has fewer than SHORT_LIMIT: the Browsertrix pages WARC + # has 3, and the screenshots/text WARCs have 0 because they carry only `resource` records. + # None means "at least SHORT_LIMIT", which is every real crawl WARC. + captures: int | None = None + + +EXAMPLES = [ + Example( + "cc-500-records", + HF + "500_RECORDS-REPACKAGE-CC-MAIN-2026-30.warc.gz?download=true", + "Common Crawl Nutch slice, response+request+metadata", + "13 MB", + ), + Example( + "cc-homepages", + HF + "HOMEPAGES-REPACKAGE-CC-MAIN-2026-21.warc.gz?download=true", + "cdx_toolkit repackage, response records only", + "1 GB", + ), + Example( + "cc-us-federal", + HF + "IS_US_FEDERAL-REPACKAGE-CC-MAIN-2025-13.warc.gz?download=true", + "cdx_toolkit repackage, response records only", + "427 MB", + ), + Example( + "cc-main-warc", + CC + "crawl-data/CC-MAIN-2026-34/segments/1786091384908.68/warc/CC-MAIN-20260807101845-20260807131845-00000.warc.gz", + "Common Crawl Nutch, main WARC", + "~1 GB", + ), + Example( + "cc-main-crawldiagnostics", + CC + "crawl-data/CC-MAIN-2026-34/segments/1786091384908.68/crawldiagnostics/CC-MAIN-20260807101845-20260807131845-00000.warc.gz", + "Common Crawl Nutch, crawldiagnostics (revisit records)", + "~1 GB", + ), + Example( + "cc-main-robotstxt", + CC + "crawl-data/CC-MAIN-2026-34/segments/1786091384908.68/robotstxt/CC-MAIN-20260807101845-20260807131845-00000.warc.gz", + "Common Crawl Nutch, robotstxt (no metadata records)", + "~1 GB", + ), + Example( + "eot2024-heritrix", + EOT + "EOT-2024/segments/IA-000/warc/EOT24PRE-20240926172119-crawl804_EOT24PRE-20240926172119-00000.warc.gz", + "Heritrix / Internet Archive", + "1 GB", + ), + Example( + "ia-eot24pre-heritrix", + IA + "EOT24PRE-20240926175758-crawl808/EOT24PRE-20240926175758-00032.warc.gz", + "Heritrix, read from an Internet Archive item over ia://", + "1.6 GB", + ), + Example( + "eot2024-nutch-repackage", + EOT + "EOT-2024/segments/CC-000/warc/EOT-2024-REPACKAGE-CC-MAIN-2024-42-GOV-000000-001.warc.gz", + "Common Crawl repackage for EOT", + "~1 GB", + ), + Example( + "eot2024-browsertrix-pages", + EOT + "EOT-2024/segments/WR-000/warc/EOT24WR-0015_20250114215650265-8c53efcc-e2d-0_eot-http-energy-gov-eere-office-energy-efficiency-renewable-energy-manual-20250114215335-8c53efcc-e2d-20250114215647018-0.warc.gz", + "Browsertrix, pages", + "small", + captures=3, + ), + Example( + "eot2024-browsertrix-screenshots", + EOT + "EOT-2024/segments/WR-000/warc/EOT24WR-0015_20250114215650265-8c53efcc-e2d-0_eot-http-energy-gov-eere-office-energy-efficiency-renewable-energy-manual-20250114215335-8c53efcc-e2d-screenshots-20250114215649547.warc.gz", + "Browsertrix, screenshots (resource records only)", + "small", + captures=0, + ), + Example( + "eot2024-browsertrix-text", + EOT + "EOT-2024/segments/WR-000/warc/EOT24WR-0015_20250114215650265-8c53efcc-e2d-0_eot-http-energy-gov-eere-office-energy-efficiency-renewable-energy-manual-20250114215335-8c53efcc-e2d-text-20250114215649747.warc.gz", + "Browsertrix, extracted text (resource records only)", + "small", + captures=0, + ), + Example( + "eot2024-archiveteam-megawarc", + EOT + "EOT-2024/segments/AT-000/warc/archiveteam_usgovernment_20250131232111_96ad506d_usgovernment_20250131232111_96ad506d.1738361595.megawarc.warc.gz", + "ArchiveTeam megawarc (several concatenated WARCs, 4 warcinfo records)", + "10 GB", + ), + Example( + "eot2004-heritrix-arc", + EOT + "EOT-2004/segments/NARA-000/warc/NARA-PEOT-2004-20041014205819-00000-crawling009-c_NARA-PEOT-2004-20041014205819-00000-crawling009.archive.org.arc.gz", + "Heritrix ARC/1.1 with dns: records", + "~100 MB", + ), + Example( + "ccf2008-arc", + CC + "crawl-001/2008/06/19/0/1213886083018_0.arc.gz", + "Common Crawl 2008 ARC", + "~100 MB", + ), +] + +FORMATS = ["flat", "sidecar"] + +pytestmark = [ + pytest.mark.network, + pytest.mark.skipif(bool(os.environ.get("WARC2ZIP_OFFLINE")), reason="WARC2ZIP_OFFLINE is set"), +] + + +def run_warc2zip(example, output, output_format, limit=None, timeout=None): + """Run the CLI as a subprocess (the exit code is what a crash test is about) and return the + completed process. Stdout/stderr are captured so a failure shows the tool's own summary.""" + cmd = [sys.executable, "-m", "warc2zip", example.url, "--output", str(output), "--format", output_format] + if limit is not None: + cmd += ["--limit", str(limit)] + return subprocess.run(cmd, capture_output=True, text=True, timeout=timeout, check=False) + + +def check_output(example, proc, output, captures, exact=True): + """Exit 0, a zip that passes testzip(), one root directory, and a manifest whose row count + matches `captures` (exactly, or as a lower bound when `exact` is False) and whose `filename` + column agrees with the members: every response row names a payload file, no revisit row does. + The manifest is the CSV-only user's index into the zip, so that agreement is the one property + worth checking on every producer's output.""" + context = f"{example.name} ({example.producer}, {example.size})\n--- stdout\n{proc.stdout}\n--- stderr\n{proc.stderr}" + assert proc.returncode == 0, f"exit {proc.returncode}: {context}" + assert output.exists(), f"no output zip: {context}" + assert output.stat().st_size > 0, f"empty zip: {context}" + + with zipfile.ZipFile(output) as zf: + assert zf.testzip() is None, f"corrupt member in zip: {context}" + names = zf.namelist() + roots = {name.split("/", 1)[0] for name in names} + assert len(roots) == 1, f"expected one root directory, got {sorted(roots)}: {context}" + (root,) = roots + manifest = list(csv.DictReader(io.TextIOWrapper(zf.open(f"{root}/manifest.csv"), encoding="utf-8"))) + + if exact: + assert len(manifest) == captures, f"{len(manifest)} manifest rows, expected {captures}: {context}" + else: + assert len(manifest) >= captures, f"{len(manifest)} manifest rows, expected at least {captures}: {context}" + + # `filename` is a bare name in both formats; sidecar mode nests it under a domain directory, + # and a sidecar member's basename ("1000000.html.request.json") never equals a payload's. + basenames = {posixpath.basename(name) for name in names} + for row in manifest: + if row["warc_type"] == "response": + assert row["filename"] in basenames, f"manifest names {row['filename']}, not in the zip: {context}" + else: + assert row["filename"] not in basenames, f"revisit row {row['filename']} has a payload: {context}" + + +@pytest.mark.parametrize("output_format", FORMATS) +@pytest.mark.parametrize("example", EXAMPLES, ids=lambda e: e.name) +def test_short(example, output_format, tmp_path): + """The CI tier: the first SHORT_LIMIT captures of every README archive, both formats.""" + output = tmp_path / "out.zip" + proc = run_warc2zip(example, output, output_format, limit=SHORT_LIMIT, timeout=SHORT_TIMEOUT) + captures = SHORT_LIMIT if example.captures is None else min(example.captures, SHORT_LIMIT) + check_output(example, proc, output, captures) + + +@pytest.mark.long +@pytest.mark.parametrize("output_format", FORMATS) +@pytest.mark.parametrize("example", EXAMPLES, ids=lambda e: e.name) +def test_long(example, output_format, tmp_path): + """The full-file tier, ``pytest -m long``. See the module docstring for the cost.""" + output = tmp_path / "out.zip" + try: + proc = run_warc2zip(example, output, output_format) + if example.captures is None: + # A real crawl WARC holds thousands of captures; pinning the number buys nothing. + check_output(example, proc, output, SHORT_LIMIT, exact=False) + else: + check_output(example, proc, output, example.captures) + finally: + # A full-file zip is up to 10 GB; do not leave it under pytest's tmp dir. + if output.exists(): + output.unlink() diff --git a/tools/warc_limit.py b/tools/warc_limit.py new file mode 100644 index 0000000..00f02b2 --- /dev/null +++ b/tools/warc_limit.py @@ -0,0 +1,85 @@ +import argparse + +from warcio.archiveiterator import ArchiveIterator +from warcio.utils import fsspec_open +from warcio.warcwriter import WARCWriter + + +def limit_warc(input_file, output_path, limit): + """Copy the first `limit` captures from a WARC into a new gzipped WARC. + + A "capture" is a `response` record plus its associated `request`/`metadata` + records (matching warc2zip.py's --limit semantics). The `warcinfo` record is + preserved but not counted toward the limit. + + Uses the drain-until-next-response pattern (see warc2zip.py:240-257): the Nth + response sets `limit_reached`, and we keep writing the trailing + request/metadata records that belong to capture N until the *next* response + arrives, at which point we stop. A `limit` of 0 (or less) writes only the + leading `warcinfo` record; a `limit` larger than the file simply copies + everything. + + Returns the number of response records (captures) written. + """ + response_count = 0 + limit_reached = False + + with fsspec_open(input_file, "rb") as stream, open(output_path, "wb") as out: + writer = WARCWriter(out, gzip=True) + for record in ArchiveIterator(stream): + if limit_reached and record.rec_type == "response": + break + # limit=0: stop before the very first capture, but keep leading warcinfo. + if limit <= 0 and record.rec_type == "response": + break + writer.write_record(record) # write BEFORE advancing the iterator + if record.rec_type == "response": + response_count += 1 + if response_count >= limit: + limit_reached = True + + return response_count + + +def default_output_path(input_file, limit): + """Derive `foo.warc.gz` -> `foo.limit{N}.warc.gz`.""" + for suffix in (".warc.gz", ".warc", ".gz"): + if input_file.endswith(suffix): + return f"{input_file[: -len(suffix)]}.limit{limit}.warc.gz" + return f"{input_file}.limit{limit}.warc.gz" + + +def main(input_file, output_path, limit): + if output_path is None: + output_path = default_output_path(input_file, limit) + + response_count = limit_warc(input_file, output_path, limit) + + print(f"Wrote {response_count} captures to {output_path}") + + +if __name__ == "__main__": + parser = argparse.ArgumentParser( + description="Create a new WARC keeping only the first N captures of an existing WARC" + ) + parser.add_argument( + "input_file", + type=str, + help="Local path or remote URI (s3://, http://, ...) to a .warc.gz file", + ) + parser.add_argument( + "-n", + "--limit", + type=int, + required=True, + help="Number of captures (response records) to keep", + ) + parser.add_argument( + "--output", + type=str, + default=None, + help="Output WARC path (default: .limit.warc.gz)", + ) + args = parser.parse_args() + + main(args.input_file, args.output, args.limit) diff --git a/uv.lock b/uv.lock index f12c19e..f5cdd63 100644 --- a/uv.lock +++ b/uv.lock @@ -208,6 +208,19 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/fb/51/08f32aea872253173f513ba68122f4300966290677c8e59887b4ffd5d957/botocore-1.42.70-py3-none-any.whl", hash = "sha256:54ed9d25f05f810efd22b0dfda0bb9178df3ad8952b2e4359e05156c9321bd3c", size = 14671393, upload-time = "2026-03-17T19:43:06.777Z" }, ] +[[package]] +name = "cdx-toolkit" +version = "0.9.39" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "requests" }, + { name = "warcio" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/c2/92/ab18e28629867f3af01f64974ef75075dd5775ad9cbd1dca42faba1504af/cdx_toolkit-0.9.39.tar.gz", hash = "sha256:7fd85cd5b00ec3ac57bfa21e7f70ba633e1df59760eb54dae3a85e63596af4bc", size = 92815, upload-time = "2026-06-01T11:49:02.752Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/93/22/5cda1ba55a88188b5b19ae60961244759b8b0cc5fe97fb29516afef32adf/cdx_toolkit-0.9.39-py3-none-any.whl", hash = "sha256:577ae9a9578f33ce4bae6c887074f0f70331bf4d6a66a26d1e20836b3004e431", size = 32685, upload-time = "2026-06-01T11:49:01.507Z" }, +] + [[package]] name = "certifi" version = "2026.2.25" @@ -961,6 +974,8 @@ name = "warc2zip" version = "0.1.0" source = { editable = "." } dependencies = [ + { name = "cdx-toolkit" }, + { name = "fsspec" }, { name = "tqdm" }, { name = "warcio", extra = ["s3"] }, ] @@ -973,6 +988,8 @@ dev = [ [package.metadata] requires-dist = [ + { name = "cdx-toolkit", specifier = ">=0.9.39" }, + { name = "fsspec" }, { name = "pytest", marker = "extra == 'dev'", specifier = ">=8.0" }, { name = "ruff", marker = "extra == 'dev'" }, { name = "tqdm", specifier = ">=4.67.3" }, diff --git a/warc2zip.py b/warc2zip.py index d6b1481..e5c6f87 100644 --- a/warc2zip.py +++ b/warc2zip.py @@ -1,23 +1,66 @@ import argparse +import asyncio import csv +import functools import io import json import mimetypes import posixpath +import random import re import secrets import sys +import time import zipfile +import zlib from collections import Counter -from datetime import datetime, timezone from dataclasses import dataclass, field +from datetime import datetime, timezone from pathlib import Path -from urllib.parse import urlparse +from urllib.parse import urljoin, urlparse +from cdx_toolkit.myrequests import get_retries, myrequests_get, retry_info +from fsspec.core import url_to_fs +from fsspec.implementations.ia import ENV_ACCESS_KEY, ENV_SECRET_KEY, InternetArchiveFileSystem, load_ia_credentials from tqdm import tqdm from warcio.archiveiterator import ArchiveIterator from warcio.recordloader import ARC2WARCHeadersParser from warcio.utils import fsspec_open +from warcio.warcwriter import WARCWriter + +# --- Internet Archive input: ia:/// ------------------------------------------ +# +# `InternetArchiveFileSystem` lives in fsspec (fsspec/implementations/ia.py, protocol "ia"): an +# HTTPFileSystem reading https://archive.org/download//, with credentials from +# the ia.ini that `ia configure` writes. fsspec registers the protocol itself, so fsspec_open(), +# url_to_fs() and therefore main(), --fetch and input_basename() take ia:// with no code here. Until +# a fsspec release ships it, the branch behind the upstream PR is installed on top of the released +# fsspec (see pyproject.toml). What stays in warc2zip is the CLI's explanation of a refused item. + + +def format_ia_permission_error(input_file, credentials=None): + """The stderr message for a 401/403 from archive.org on an ia:// input.""" + credentials = credentials or load_ia_credentials() + item = urlparse(input_file).netloc + if credentials.anonymous: + where = credentials.config_file or "no ia.ini found" + how = f"no archive.org credentials ({where}), so the request was anonymous" + else: + where = credentials.config_file or "the environment" + how = f"the credentials from {where} were refused (expired cookies look the same)" + return "\n".join( + [ + f"Error: archive.org refused access to item '{item}': {how}.", + "", + "A restricted item needs the account that can see it to be logged in:", + "", + " pip install internetarchive && ia configure # writes ~/.config/internetarchive/ia.ini", + "", + f"or set {ENV_ACCESS_KEY} and {ENV_SECRET_KEY}. Public items need neither; check the item name.", + f"Underlying URL: {InternetArchiveFileSystem._strip_protocol(input_file)}", + ] + ) + MIME_EXTENSION_OVERRIDES = { "text/html": ".html", @@ -517,18 +560,60 @@ def get_file_size(input_file): return None -def build_root_dir_name(crawl_name, partial=False): +def new_run_id(): + """4-char hex that makes one run's outputs unique: shared by the zip name and the root dir.""" + return secrets.token_hex(2) + + +def build_root_dir_name(crawl_name, partial=False, run_id=None): """Build a unique root directory name from a crawl name. - Format: {crawl_name}_{YYYYMMDDTHHMMSS}_{4-char hex suffix} - The suffix doesn't affect sort order since it comes after the timestamp. + Format: {crawl_name}_{YYYYMMDDTHHMMSS}_{4-char hex suffix}[_partial] + The hex suffix doesn't affect sort order since it comes after the timestamp. `partial` is + set iff --limit was given. `run_id` is the hex; main() passes the same one it gave + default_output_path(), so the zip on disk and the directory inside it carry the same tag. """ timestamp = datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%S") - suffix = secrets.token_hex(2) if not partial else secrets.token_hex(2) + "_partial" + suffix = run_id or new_run_id() + if partial: + suffix += "_partial" return f"{crawl_name}_{timestamp}_{suffix}" +def input_basename(input_file): + """Basename of the input, whatever shape it comes in. + + The README's inputs are local paths, `s3://bucket/key`, `https://host/path`, the same with a + query string (`https://huggingface.co/.../X.warc.gz?download=true`), and `-` for stdin. A + plain posixpath.basename() keeps the query string, which would put `?download=true` in the + zip name and the root directory. URIs are parsed and only the path's last segment kept; + stdin has no name and is labelled "stdin". + """ + if input_file == "-": + return "stdin" + parsed = urlparse(input_file) + path = parsed.path if parsed.scheme and parsed.netloc else input_file + return posixpath.basename(path.rstrip("/")) + + +def default_output_path(input_file, partial=False, run_id=None): + """Default zip path when --output is not given: {input_basename}_{hex}[_partial].zip. + + The .warc.gz / .warc / .arc.gz / .arc suffix is stripped and .zip appended. The hex is the + run id (see new_run_id) — without it Common Crawl's `warc/`, `crawldiagnostics/` and + `robotstxt/` files, which share a basename and differ only by directory, overwrite each + other's zip. `_partial` follows the same rule as the root directory inside the zip (set iff + --limit was given). Written to the current directory. + """ + basename = input_basename(input_file) + label = extract_crawl_name(basename) if basename else "unknown" + suffix = run_id or new_run_id() + if partial: + suffix += "_partial" + return Path(f"{label}_{suffix}.zip") + + def extract_crawl_name(warc_filename): """Extract a clean crawl name from a WARC-Filename header value. @@ -692,8 +777,9 @@ def write_sidecar_files(zip_file, root_dir, group): zip_file.writestr(f"{base}.metadata.warc-fields", "\n\n".join(body_parts)) -def main(input_file, output_path, dry_run=False, limit=None, output_format="flat", metadata_only=False): +def main(input_file, output_path=None, dry_run=False, limit=None, output_format="flat", metadata_only=False): file_size = get_file_size(input_file) + partial = limit is not None if dry_run: capped = limit is None or limit > DRY_RUN_MAX @@ -751,13 +837,20 @@ def main(input_file, output_path, dry_run=False, limit=None, output_format="flat warcinfos = [] # list[(warc_header_pairs, body_text, offset, length)] - crawl-level provenance warc_filename = "" # WARC-Filename from the warcinfo record; every manifest row repeats it - # Fallback crawl name from input filename - input_basename = posixpath.basename(input_file.rstrip("/")) - fallback_crawl_name = extract_crawl_name(input_basename) if input_basename else "unknown" + # Fallback crawl name from input filename (same label default_output_path uses) + basename = input_basename(input_file) + fallback_crawl_name = extract_crawl_name(basename) if basename else "unknown" record_types = Counter() # every record read, by WARC-Type — extracted or not limit_reached = False + # One run id for the zip name and the root directory, so the two can be matched on disk. + # An explicit --output is used verbatim; the default shares the root directory's hex and + # _partial rule. + run_id = new_run_id() + if output_path is None: + output_path = str(default_output_path(input_file, partial, run_id)) + # response -> Record id <-> metadata - with zipfile.ZipFile(output_path, "w", zipfile.ZIP_DEFLATED) as outer_zip: # Pass 1: Read WARC, write payloads immediately, buffer only headers @@ -791,13 +884,13 @@ def main(input_file, output_path, dry_run=False, limit=None, output_format="flat if root_dir is None: warc_filename = record.rec_headers.get_header("WARC-Filename") or "" crawl_name = extract_crawl_name(warc_filename) if warc_filename else fallback_crawl_name - root_dir = build_root_dir_name(crawl_name, limit is not None) + root_dir = build_root_dir_name(crawl_name, partial, run_id) pbar.update(stream.tell() - pbar.n) continue # Resolve root_dir before first payload write if no warcinfo appeared if root_dir is None: - root_dir = build_root_dir_name(fallback_crawl_name, limit is not None) + root_dir = build_root_dir_name(fallback_crawl_name, partial, run_id) if rec_type in ("response", "revisit"): record_id = record.rec_headers.get_header("WARC-Record-ID") @@ -940,7 +1033,7 @@ def main(input_file, output_path, dry_run=False, limit=None, output_format="flat # A WARC without a warcinfo record, or without WARC-Filename on it, still needs a name # in the manifest: fall back to what the input was called. - warc_filename = warc_filename or input_basename + warc_filename = warc_filename or basename # Pass 2: Build metadata from buffered headers (payloads already in zip) sorted_groups = sorted( @@ -1030,10 +1123,519 @@ def main(input_file, output_path, dry_run=False, limit=None, output_format="flat return skipped +# --------------------------------------------------------------------------- +# --fetch: re-download the rows of a (filtered) manifest.csv as one .warc.gz +# --------------------------------------------------------------------------- + +FETCH_COLUMNS = ("source_uri", "warc_record_offset", "warc_record_length") +FETCH_DEFAULT_RATE = 2.0 # requests per second, per host +FETCH_DEFAULT_RETRIES = 8 # retries after the first attempt +FETCH_MAX_BACKOFF = 60.0 # seconds; the cap on one exponential-backoff wait +FETCH_RETRYABLE_STATUSES = frozenset({408, 425, 429, 500, 502, 503, 504}) +FETCH_MAX_GAP = 64 * 1024 # bytes between two rows that still get fetched in one request +FETCH_MAX_SPAN = 64 * 1024 * 1024 # bytes held in memory for one request +FETCH_HEAD_BYTES = 64 * 1024 # probe size for the source's leading warcinfo record +FETCH_MAX_REDIRECTS = 5 # cdx_toolkit's myrequests_get does not follow redirects; http_range does +FETCH_REDIRECT_STATUSES = frozenset({301, 302, 303, 307, 308}) + + +@dataclass +class FetchRow: + source_uri: str + offset: int + length: int + record_id: str # blank for ARC rows (the id was minted, so the manifest leaves it out) + target_uri: str + line: int # CSV line number, for warnings + + +class FetchLengthMismatch(Exception): + """The server returned a different number of bytes than the range asked for. + + fsspec sends the Range header but never checks for a 206, so a server that ignores Range + hands back the whole file as the "slice". Deterministic, so never retried. + """ + + +def read_fetch_rows(csv_path): + """Parse a manifest.csv (filtered or not) into FetchRows. + + Missing a required column is a usage error (exit 2). A row with a blank or non-numeric + offset/length (main() writes "" when the offset was unknown), or with no fetchable + `source_uri` (stdin input), is warned about and counted, never fatal. + Returns (rows, skipped). + """ + if csv_path == "-": + return _read_fetch_rows(io.TextIOWrapper(sys.stdin.buffer, encoding="utf-8", newline=""), "stdin") + with open(csv_path, newline="", encoding="utf-8") as fh: + return _read_fetch_rows(fh, csv_path) + + +def _read_fetch_rows(fh, label): + reader = csv.DictReader(fh) + missing = [c for c in FETCH_COLUMNS if c not in (reader.fieldnames or [])] + if missing: + print(f"error: {label}: not a manifest.csv, missing column(s): {', '.join(missing)}", file=sys.stderr) + raise SystemExit(2) + + rows, skipped = [], 0 + for row in reader: + source_uri = (row.get("source_uri") or "").strip() + try: + offset = int(row.get("warc_record_offset") or "") + length = int(row.get("warc_record_length") or "") + except ValueError: + offset, length = -1, 0 + + if not source_uri or source_uri == "-": + problem = "no fetchable source_uri" + elif offset < 0 or length <= 0: + problem = "no usable warc_record_offset/warc_record_length" + else: + problem = None + if problem: + print(f"warning: {label} line {reader.line_num}: {problem}, skipped", file=sys.stderr) + skipped += 1 + continue + + rows.append( + FetchRow( + source_uri=source_uri, + offset=offset, + length=length, + record_id=(row.get("warc_record_id") or "").strip(), + target_uri=(row.get("warc_target_uri") or "").strip(), + line=reader.line_num, + ) + ) + return rows, skipped + + +def group_fetch_rows(rows): + """Rows keyed by source_uri (first-appearance order), each list sorted by offset. + + Exact (offset, length) duplicates are dropped with one warning per source: fetching the + same member twice would put the same record in the output twice. Output order is + therefore source-then-offset, not CSV order. + """ + groups = {} + for row in rows: + groups.setdefault(row.source_uri, []).append(row) + + for source_uri, source_rows in groups.items(): + source_rows.sort(key=lambda r: (r.offset, r.length)) + deduped, seen = [], set() + for row in source_rows: + key = (row.offset, row.length) + if key not in seen: + seen.add(key) + deduped.append(row) + if len(deduped) != len(source_rows): + print(f"warning: {source_uri}: {len(source_rows) - len(deduped)} duplicate row(s) dropped", file=sys.stderr) + groups[source_uri] = deduped + return groups + + +def coalesce_ranges(rows, max_gap=FETCH_MAX_GAP, max_span=FETCH_MAX_SPAN): + """Merge offset-sorted rows into (start, end, rows) spans fetched with one request each. + + A response record is never adjacent to the next one in a Common Crawl WARC (its request + and metadata records sit in between), so merging only touching ranges would merge nothing. + Instead rows closer than `max_gap` share a request and each row is sliced back out of the + span locally: the output stays response-only and byte-exact, and a status-filtered subset + collapses into a fraction of the requests, which is what data.commoncrawl.org throttles on. + `max_gap` bounds the bytes wasted per merge, `max_span` the bytes held in memory. + """ + spans = [] + for row in rows: + end = row.offset + row.length + if spans: + start, span_end, span_rows = spans[-1] + if row.offset - span_end <= max_gap and end - start <= max_span: + spans[-1] = (start, max(span_end, end), span_rows + [row]) + continue + spans.append((row.offset, end, [row])) + return spans + + +class RateLimiter: + """Minimum interval between requests, tracked per key (a host). rate <= 0 disables it.""" + + def __init__(self, rate, clock=time.monotonic, sleep=time.sleep): + self.interval = 1.0 / rate if rate and rate > 0 else 0.0 + self.clock = clock + self.sleep = sleep + self.next_allowed = {} + + def wait(self, key): + if not self.interval: + return + now = self.clock() + delay = self.next_allowed.get(key, now) - now + if delay > 0: + self.sleep(delay) + now += delay + self.next_allowed[key] = now + self.interval + + +def error_status(exc): + """HTTP status carried by an exception (aiohttp's ClientResponseError), else None.""" + status = getattr(exc, "status", None) + return status if isinstance(status, int) else None + + +def is_retryable(exc): + """Transient failures only: throttling, server errors, timeouts, dropped connections. + + A 404/403 (fsspec raises FileNotFoundError / PermissionError), any other 4xx (416 means + the offsets do not fit this file) and a length mismatch are deterministic and re-raised + at once. aiohttp's non-OSError exceptions (payload/disconnect) are recognised by module, + so aiohttp is not imported here. + """ + if isinstance(exc, FetchLengthMismatch): + return False + status = error_status(exc) + if status is not None: + return status in FETCH_RETRYABLE_STATUSES + if isinstance(exc, (FileNotFoundError, PermissionError)): + return False + if isinstance(exc, (OSError, TimeoutError, asyncio.TimeoutError)): + return True + return type(exc).__module__.split(".")[0] == "aiohttp" + + +def retry_after_seconds(exc): + """The Retry-After delay the server asked for, in seconds; only the delta-seconds form.""" + headers = getattr(exc, "headers", None) + if not headers: + return None + try: + return max(0.0, float(headers.get("Retry-After"))) + except (TypeError, ValueError): + return None + + +def fetch_with_retry(fetch, *, retries=FETCH_DEFAULT_RETRIES, label="", sleep=time.sleep, rng=random.random): + """Call the zero-argument `fetch` until it returns, retrying transient failures. + + Waits Retry-After when the server sends one, otherwise an exponential backoff (2^attempt + seconds, capped at FETCH_MAX_BACKOFF) with jitter in [0.5, 1.5). One stderr line per retry. + Re-raises on a non-retryable error or once `retries` retries are used up. + """ + attempts = max(1, retries + 1) + for attempt in range(1, attempts + 1): + try: + return fetch() + except Exception as exc: + if attempt == attempts or not is_retryable(exc): + raise + delay = retry_after_seconds(exc) + if delay is None: + delay = min(FETCH_MAX_BACKOFF, 2.0**attempt) * (0.5 + rng()) + print( + f"warning: {label}: attempt {attempt}/{attempts} failed ({exc}), retrying in {delay:.1f} s", + file=sys.stderr, + ) + sleep(delay) + + +def leading_warcinfo(head): + """The first gzip member of `head`, verbatim, if it is a warcinfo record; else None. + + `head` is the first FETCH_HEAD_BYTES of the source. Going through open_archive_iterator() + means an ARC's filedesc:// record qualifies too. A warcinfo record that does not fit in + the probe, or a source whose first record is something else, yields None: nothing is + ever synthesised in its place. + """ + try: + iterator = open_archive_iterator(io.BytesIO(head)) + record = next(iter(iterator), None) + if record is None: + return None + rec_type = record.rec_type + record.content_stream().read() + length = iterator.get_record_length() + except Exception: # noqa: BLE001 - any parse failure means "no usable warcinfo" + return None + if rec_type != "warcinfo" or not length or length > len(head): + return None + member = head[:length] + # warcio reads a truncated gzip member without complaint; only a member whose trailer is + # inside the probe is a complete record. + if member.startswith(b"\x1f\x8b") and not _complete_gzip_member(member): + return None + return member + + +def _complete_gzip_member(data): + decompressor = zlib.decompressobj(16 + zlib.MAX_WBITS) + try: + decompressor.decompress(data) + except zlib.error: + return False + return decompressor.eof + + +def validate_record_slice(data, row): + """Check that `data` is exactly the response record the manifest row describes. + + Returns None when it is, else a short reason. "Parses as a response" is not enough: + warcio's sniffer reads arbitrary bytes as an ARC record of type response, so the + WARC-Record-ID (or, for ARC rows whose id is blank by design, the WARC-Target-URI) is + compared against the row, the parsed record must span the whole slice, and the gzip + member must be complete (warcio reads a truncated one without complaint). + """ + try: + iterator = open_archive_iterator(io.BytesIO(data)) + records = [] + for record in iterator: + record.content_stream().read() + records.append( + ( + record.rec_type, + record.rec_headers.get_header("WARC-Record-ID") or "", + record.rec_headers.get_header("WARC-Target-URI") or "", + iterator.get_record_length(), + ) + ) + except Exception as exc: # noqa: BLE001 - the reason is reported, not swallowed + return f"not a WARC record: {exc}" + + if len(records) != 1: + return f"expected one record, found {len(records)}" + rec_type, record_id, target_uri, length = records[0] + if rec_type != "response": + return f"expected a response record, found {rec_type}" + if length != len(data): + return f"record spans {length} of {len(data)} bytes" + if data.startswith(b"\x1f\x8b") and not _complete_gzip_member(data): + return "truncated gzip member" + if row.record_id: + if record_id != row.record_id: + return f"WARC-Record-ID {record_id} is not {row.record_id}" + elif target_uri != row.target_uri: + return f"WARC-Target-URI {target_uri} is not {row.target_uri}" + return None + + +def default_fetch_output_path(csv_path, run_id=None): + """Default .warc.gz path for --fetch: {csv basename minus .csv}_{hex}.warc.gz, in cwd. + + Same hex as default_output_path() and for the same reason: filtering one manifest twice + into `subset.csv` must not overwrite the first subset. `-` is labelled stdin. + """ + basename = input_basename(csv_path) + label = basename[:-4] if basename.lower().endswith(".csv") else basename + return Path(f"{label or 'unknown'}_{run_id or new_run_id()}.warc.gz") + + +def http_range(url, start, end, retries=None, redirects=FETCH_MAX_REDIRECTS, get=None): + """Bytes [start, end) over http(s), through cdx_toolkit's `myrequests_get`. + + That function carries Common Crawl's own retry policy (429/5xx backed off without limit, + connection failures up to `retries`) and per-host pacing, so neither is reimplemented here. + It does not follow redirects (`allow_redirects=False`), so 3xx answers are followed by hand, + re-entering `myrequests_get` each time so the new host is paced too. A 200 with the whole + file, from a server that ignores Range, is left to the caller's length check. + """ + get = get or myrequests_get # resolved at call time so tests can substitute the transport + headers = {"Range": f"bytes={start}-{end - 1}"} + for _ in range(redirects + 1): + resp = get(url, headers=headers, raise_error_after_n_errors=retries) + location = resp.headers.get("Location") + if resp.status_code in FETCH_REDIRECT_STATUSES and location: + url = urljoin(url, location) + continue + return resp.content + raise RuntimeError(f"more than {redirects} redirects") + + +def set_host_interval(host, rate): + """Map --rate onto cdx_toolkit's per-host schedule; rate=None keeps cdx_toolkit's own defaults. + + cdx_toolkit paces requests from a module-level table keyed by hostname (0.55 s between + requests to data.commoncrawl.org, 3 s to a host it does not know). `get_retries()` creates + the host's entry from the default one, and the interval is then overwritten in place — the + table is the only place the pacing can be set. + """ + if rate is None: + return + get_retries(host) + retry_info[host]["minimum_interval"] = 1.0 / rate if rate > 0 else 0.0 + + +def check_length(data, start, end): + if len(data) != end - start: + raise FetchLengthMismatch(f"asked for {end - start} bytes at offset {start}, got {len(data)}") + return data + + +def fetch_range(fs, path, start, end, exact=True): + """Bytes [start, end) of a file through fsspec (s3 and local paths). + + With `exact`, a short or long answer raises FetchLengthMismatch (see there). The warcinfo + probe passes exact=False because it deliberately asks past the end of small files. + """ + data = fs.cat_file(path, start=start, end=end) + return check_length(data, start, end) if exact else data + + +def annotate_record_slice(data, source_uri, offset, length): + """Re-serialise one record with WARC-Source-URI / WARC-Source-Range added. + + The names and values are cdx_toolkit's convention for extracts, so a subset made this way + matches one from `cdxt warc`, and a re-converted subset carries every record's original + coordinates in response_warc_headers.csv. The WARC header block is rewritten by warcio and + the member recompressed, so it is no longer the source's bytes, but the content block (HTTP + headers and body) is copied through untouched (a verbatim switch was tried and dropped: nothing + in the workflow needs the wire bytes, and the curl loop in the README gives them anyway). + Cost measured on CC records: ~1.4 ms CPU and ~43 bytes of output per record. + """ + # no_record_parse=True leaves the HTTP layer unparsed, so warcio writes the content block + # through byte for byte (a parsed one is re-serialised: obs-folds unfolded, Content-Length + # recomputed, WARC-Block-Digest silently wrong). Only the WARC header block is rewritten. + record = next(iter(ArchiveIterator(io.BytesIO(data), no_record_parse=True, arc2warc=True))) + record.rec_headers.replace_header("WARC-Source-URI", source_uri) + record.rec_headers.replace_header("WARC-Source-Range", f"bytes={offset}-{offset + length - 1}") + buffer = io.BytesIO() + WARCWriter(buffer, gzip=True).write_record(record) + return buffer.getvalue() + + +def _is_local_filesystem(fs): + protocol = fs.protocol if isinstance(fs.protocol, (tuple, list)) else (fs.protocol,) + return "file" in protocol + + +class _FetchSource: + """One source WARC: its transport, retry policy, pacing key and transfer counters. + + http(s) goes through cdx_toolkit (`http_range`), which retries and paces on its own; s3 and + local paths go through fsspec with `fetch_with_retry` and `RateLimiter` around them. + """ + + def __init__(self, source_uri, limiter, rate=None, retries=None): + self.source_uri = source_uri + parsed = urlparse(source_uri) + self.via_http = parsed.scheme.lower() in ("http", "https") + self.retries = retries + self.requests = 0 + self.transferred = 0 + if self.via_http: + self.fs = self.path = self.limiter = None + set_host_interval(parsed.hostname, rate) + self.host = parsed.hostname + else: + self.fs, self.path = url_to_fs(source_uri) + # Local files bypass the limiter; for s3 the key is the bucket. + self.host = None if _is_local_filesystem(self.fs) else parsed.netloc + self.limiter = limiter + + def get(self, start, end, exact=True): + self.requests += 1 + if self.via_http: + data = http_range(self.source_uri, start, end, retries=self.retries) + if exact: + check_length(data, start, end) + else: + if self.host: + self.limiter.wait(self.host) + data = fetch_range(self.fs, self.path, start, end, exact) + self.transferred += len(data) + return data + + def fetch(self, start, end, exact=True, label="", sleep=time.sleep): + if self.via_http: + return self.get(start, end, exact) + retries = FETCH_DEFAULT_RETRIES if self.retries is None else self.retries + return fetch_with_retry( + functools.partial(self.get, start, end, exact), retries=retries, label=label, sleep=sleep + ) + + +def fetch_main(csv_path, output_path=None, rate=None, retries=None, sleep=time.sleep): + """--fetch: download every manifest row's byte range and concatenate the raw members. + + Each source's own warcinfo record is copied first (see leading_warcinfo), then its rows + in offset order. A span that fails after retries skips all its rows with a warning and + the run goes on, like the CSV writers: the file is a valid .warc.gz at every member + boundary, and the return value (skipped rows) makes cli() exit 1. `rate` and `retries` + are None by default so each transport keeps its own defaults (see _FetchSource). Every + record is stamped with WARC-Source-URI / WARC-Source-Range (see annotate_record_slice); the + warcinfo record is copied as-is. + """ + rows, skipped = read_fetch_rows(csv_path) + groups = group_fetch_rows(rows) + if output_path is None: + output_path = str(default_fetch_output_path(csv_path, new_run_id())) + + limiter = RateLimiter(FETCH_DEFAULT_RATE if rate is None else rate, sleep=sleep) + total_rows = sum(len(source_rows) for source_rows in groups.values()) + written = warcinfo_count = 0 + sources = [] + + with open(output_path, "wb") as out, tqdm(total=total_rows, unit="rec", desc="Fetching") as pbar: + for source_uri, source_rows in groups.items(): + source = _FetchSource(source_uri, limiter, rate=rate, retries=retries) + sources.append(source) + + try: + head = source.fetch(0, FETCH_HEAD_BYTES, exact=False, label=f"{source_uri} (warcinfo probe)", sleep=sleep) + except Exception as exc: # noqa: BLE001 - reported below; the rows are still attempted + print(f"warning: {source_uri}: could not read the leading warcinfo record ({exc})", file=sys.stderr) + head = b"" + warcinfo = leading_warcinfo(head) + if warcinfo: + out.write(warcinfo) + warcinfo_count += 1 + else: + print(f"note: {source_uri}: no leading warcinfo record, subset starts at its first response", + file=sys.stderr) + + for start, end, span_rows in coalesce_ranges(source_rows): + label = f"{source_uri} bytes {start}-{end - 1}" + try: + data = source.fetch(start, end, label=label, sleep=sleep) + except Exception as exc: # noqa: BLE001 - skip-and-warn, the run continues + print(f"warning: {label}: could not be fetched ({exc}), {len(span_rows)} row(s) skipped", + file=sys.stderr) + skipped += len(span_rows) + pbar.update(len(span_rows)) + continue + for row in span_rows: + chunk = data[row.offset - start : row.offset - start + row.length] + problem = validate_record_slice(chunk, row) + if problem: + print(f"warning: {csv_path} line {row.line}: fetched bytes are not the expected record " + f"({problem}), skipped", file=sys.stderr) + skipped += 1 + else: + out.write(annotate_record_slice(chunk, source_uri, row.offset, row.length)) + written += 1 + pbar.update(1) + + request_count = sum(s.requests for s in sources) + transferred = sum(s.transferred for s in sources) + print( + f"Created {output_path}: {written} records from {len(groups)} source(s), " + f"{request_count} requests, {tqdm.format_sizeof(transferred, suffix='B')} transferred " + f"({warcinfo_count} warcinfo record{'' if warcinfo_count == 1 else 's'})" + ) + if skipped: + print(f"warning: {skipped} row(s) could not be fetched (see warnings above)", file=sys.stderr) + return skipped + + def cli(): parser = argparse.ArgumentParser(description="Convert a gzipped WARC file into a zip-of-zips archive.") - parser.add_argument("input_file", help="Path to a .warc.gz file, or '-' for stdin") - parser.add_argument("--output", default=None, help="Output zip path (default: replace .warc.gz with .zip)") + parser.add_argument("input_file", help="Path to a .warc.gz file, or '-' for stdin (a manifest.csv with --fetch)") + parser.add_argument( + "--output", + default=None, + help="Output zip path (default: {basename}_{hex}.zip in the current directory, with the same hex as the " + "root directory inside and _partial appended when --limit is set). With --fetch: the output .warc.gz " + "(default: {basename}_{hex}.warc.gz)", + ) mode = parser.add_mutually_exclusive_group() mode.add_argument("--dry-run", action="store_true", @@ -1044,6 +1646,12 @@ def cli(): help="Write every CSV, manifest and sidecar but no payload files. The WARC is still " "streamed in full, so this saves output size, not transfer (use --limit for that).", ) + mode.add_argument( + "--fetch", + action="store_true", + help="Treat input_file as a manifest.csv (filtered or not) and download every row's byte range from its " + "source_uri into one .warc.gz, with retries and per-host rate limiting.", + ) parser.add_argument( "--limit", help="Limit to N capture records, with their full set of associated request/metadata records", @@ -1053,31 +1661,53 @@ def cli(): parser.add_argument( "--format", choices=["flat", "sidecar"], - default="flat", - help="Output format: 'flat' (counter-named files + global CSVs) or 'sidecar' (domain dirs + per-file metadata)", + default=None, + help="Output format: 'flat' (counter-named files + global CSVs) or 'sidecar' (domain dirs + per-file " + "metadata). Default: flat", + ) + parser.add_argument( + "--rate", + type=float, + default=None, + help="--fetch only: requests per second per host, 0 for unlimited. Default: cdx_toolkit's per-host " + f"pacing for http(s) (about 1.8/s for data.commoncrawl.org, 0.33/s elsewhere), {FETCH_DEFAULT_RATE:g}/s for s3", + ) + parser.add_argument( + "--retries", + type=int, + default=None, + help="--fetch only: for http(s), connection failures tolerated per request (default 100; throttling " + f"and server errors are retried without limit, as cdx_toolkit does); for s3, retries per request " + f"(default {FETCH_DEFAULT_RETRIES})", ) args = parser.parse_args() - input_file = args.input_file - if args.output: - output_path = Path(args.output) - else: - # Extract basename from local path or remote URI - name = posixpath.basename(input_file.rstrip("/")) - if name.endswith(".warc.gz"): - name = name[: -len(".warc.gz")] + ".zip" - else: - name = name + ".zip" - output_path = Path(name) - - skipped = main( - input_file, - str(output_path), - dry_run=args.dry_run, - limit=args.limit, - output_format=args.format, - metadata_only=args.metadata_only, - ) + # Flags that would be silently ignored are refused instead: --format defaults to None so an + # explicit value is detectable here, and becomes "flat" only when it reaches main(). + if args.fetch: + if args.limit is not None or args.format is not None: + parser.error("--limit and --format do not apply to --fetch") + return 1 if fetch_main(args.input_file, args.output, rate=args.rate, retries=args.retries) else 0 + if args.rate is not None or args.retries is not None: + parser.error("--rate and --retries only apply to --fetch") + + # Default output naming lives in main() (see default_output_path). + try: + skipped = main( + args.input_file, + args.output, + dry_run=args.dry_run, + limit=args.limit, + output_format=args.format or "flat", + metadata_only=args.metadata_only, + ) + except PermissionError as exc: + # archive.org's 401/403 (see InternetArchiveFileSystem); other inputs keep the traceback. + if urlparse(args.input_file).scheme != "ia": + raise + print(format_ia_permission_error(args.input_file), file=sys.stderr) + print(f"Underlying error: {exc!r}", file=sys.stderr) + return 1 return 1 if skipped else 0