> For the complete documentation index, see [llms.txt](https://docs.arcv.network/llms.txt). Markdown versions of documentation pages are available by appending `.md` to page URLs; this page is available as [Markdown](https://docs.arcv.network/3.-enterprise-and-frontier-ai-labs/3.3-provenance-and-dataloaders.md).

# 3.3 Provenance & PyTorch/Hugging Face Dataloaders

A training loader must establish what its bytes represent before admitting them into a training run. ARCV's intended provenance path links a registry record on Arc to a permanent artifact and then to the shards consumed by the model pipeline. A gateway response, a file extension, or a successful download is not sufficient evidence of that relationship.

The current portal exports raw source data and displays illustrative provenance. The examples here are standalone reference integrations, not a released ARCV SDK or a claim that a public production dataset endpoint exists. They include a deterministic local fixture so the loading path can be exercised without a deployed registry.

## Establish the trust root

1. Obtain the registry address from a verified deployment record, independently of the downloaded dataset.
2. Verify Arc chain ID `5042`, deployed code, and the intended ABI. Pin an accepted block and record its hash according to the deployment's finality policy.
3. Read `datasets(datasetId)` at that block. The current return fields are bounty ID, SHA-256 digest, Arweave transaction ID, contributor, and status.
4. Require an accepted settlement state for a paid-data feed. `Paid` is enum value 2 in the current contract. This establishes an authorized validator's decision, not independent semantic correctness.
5. Retrieve the artifact identified by the registry and verify SHA-256 over its exact bytes before parsing it as a manifest or dataset.
6. If it is a manifest, verify each referenced shard against the manifest's digest, byte length, and row count before yielding records.

The public mapping is named `usedHashes`, not `registeredHashes`. Its boolean only indicates identifier reservation; it does not identify an accepted dataset or prove upload finality. Use the full record and receipt/event history when investigating acceptance.

## Manifest convention used by this example

The following loader defines a small **reference manifest version**, not an existing protocol standard. The manifest binds `version`, `format`, `split`, and a nonempty `shards` list. Each shard contains `url`, `sha256`, `bytes`, and `rows`. The digest commits to exact stored bytes; the loader never normalizes or reserializes downloaded content before hashing it.

In registry mode, the on-chain digest must commit to the manifest itself. The current contract represents one payable sample per record; it does not inherently provide a multi-shard dataset catalog. An application must explicitly define and authorize any manifest-bearing record convention. Do not interpret arbitrary existing annotation bytes as this manifest, or claim one batch record pays every contributor whose work it contains.

For records that commit directly to a single annotation, verify that artifact directly and build a separately authenticated aggregation manifest. That requires lineage from every included annotation; this example's single-root mode does not implement such an aggregation service.

## Retrieval and resource limits

The registry-mode loader retrieves through `https://arweave.net/` by default. An operator may explicitly configure a verified Irys gateway or controlled cache origin through the same gateway argument. The gateway does not change the commitment: the downloaded bytes must still hash to the registry value. A 43-character storage identifier is syntactic evidence only; bundle inclusion and archival finality require separate verification.

This reference accepts only public, unencrypted JSONL with `prompt`, `chosen`, and `rejected` text fields. It rejects redirects and non-HTTPS production URLs, bounds download sizes, requires approved shard hosts, and does not send credentials. Keep outbound-network restrictions in the execution environment as well: a host allowlist is not complete DNS-rebinding protection.

Whole-shard verification requires buffering the shard to temporary disk before records are yielded. Streaming therefore means bounded per-shard staging and incremental iteration, not consumption of unverified prefixes. Temporary storage contains plaintext; run restricted workloads only in an authorized environment with an appropriate encrypted storage and key-management design.

## Complete runnable reference loader

Use Python 3.12. Save the following block as `verified_loader.py`. The standard-library `core` demo needs no third-party packages. Framework and RPC modes require `datasets`, `torch`, and `web3` respectively. Install only the dependencies needed by your chosen mode in an isolated environment; freeze their tested versions for a training release.

```bash
python -m venv .venv
```

Activate that environment using the command appropriate to your shell, then:

```bash
python -m pip install datasets torch web3
python verified_loader.py --demo --framework core
python verified_loader.py --demo --framework hf
python verified_loader.py --demo --framework torch
```

The examples include a local fixture for deterministic execution without a funded wallet, deployed registry, or public artifact. Fixture success verifies the loader and framework integration, not a production deployment. Keep the fixture mode separate from authenticated registry or release-manifest mode.

```python
import argparse
from contextlib import contextmanager
import hashlib
import json
import re
import tempfile
from pathlib import Path
from urllib.parse import urlsplit
from urllib.request import HTTPRedirectHandler, build_opener

MAX_MANIFEST = 1024 * 1024
MAX_SHARD = 100 * 1024 * 1024
MAX_LINE = 1024 * 1024

def unique_object(pairs):
    result = {}
    for key, value in pairs:
        if key in result:
            raise ValueError("Duplicate JSON key")
        result[key] = value
    return result

def reject_constant(value):
    raise ValueError("Non-finite JSON number")

def parse_json(raw):
    if isinstance(raw, bytes):
        raw = raw.decode("utf-8", errors="strict")
    return json.loads(raw, object_pairs_hook=unique_object,
                      parse_constant=reject_constant)

class NoRedirect(HTTPRedirectHandler):
    def redirect_request(self, req, fp, code, msg, headers, newurl):
        raise ValueError("Redirect rejected")

def check_url(url, hosts, demo=False):
    if not isinstance(url, str):
        raise ValueError("URL must be a string")
    parsed = urlsplit(url)
    if demo and parsed.scheme == "file" and not parsed.netloc:
        return
    if (parsed.scheme != "https" or parsed.hostname not in hosts
            or parsed.username or parsed.password or parsed.query
            or parsed.fragment or parsed.port not in (None, 443)):
        raise ValueError("URL outside approved public HTTPS origins")

def verified_file(url, expected, limit, hosts, demo=False):
    if not isinstance(expected, str) or not re.fullmatch(r"[0-9a-f]{64}", expected):
        raise ValueError("Expected lowercase SHA-256")
    check_url(url, hosts, demo)
    output = tempfile.TemporaryFile()
    digest, size = hashlib.sha256(), 0
    try:
        with build_opener(NoRedirect()).open(url, timeout=30) as response:
            while True:
                chunk = response.read(65536)
                if not chunk:
                    break
                size += len(chunk)
                if size > limit:
                    raise ValueError("Artifact exceeds size limit")
                digest.update(chunk)
                output.write(chunk)
        if digest.hexdigest() != expected:
            raise ValueError("SHA-256 mismatch")
        output.seek(0)
        return output, size
    except BaseException:
        output.close()
        raise

def load_manifest(url, digest, hosts, demo=False):
    stream, _ = verified_file(url, digest, MAX_MANIFEST, hosts, demo)
    with stream:
        manifest = parse_json(stream.read())
    if (not isinstance(manifest, dict)
            or set(manifest) != {"version", "format", "split", "shards"}
            or type(manifest["version"]) is not int
            or manifest["version"] != 1 or manifest["format"] != "jsonl"
            or manifest["split"] != "train"):
        raise ValueError("Unsupported manifest")
    shards = manifest["shards"]
    if not isinstance(shards, list) or not 1 <= len(shards) <= 1024:
        raise ValueError("Invalid shard count")
    seen = set()
    for shard in shards:
        if not isinstance(shard, dict) or set(shard) != {"url", "sha256", "bytes", "rows"}:
            raise ValueError("Invalid shard descriptor")
        check_url(shard["url"], hosts, demo)
        digest = shard["sha256"]
        if not isinstance(digest, str) or not re.fullmatch(r"[0-9a-f]{64}", digest):
            raise ValueError("Invalid shard digest")
        if digest in seen:
            raise ValueError("Duplicate shard")
        seen.add(digest)
        if (type(shard["bytes"]) is not int or not 1 <= shard["bytes"] <= MAX_SHARD
                or type(shard["rows"]) is not int or not 1 <= shard["rows"] <= 1000000):
            raise ValueError("Invalid shard bounds")
    return shards

def read_rows(stream):
    while True:
        line = stream.readline(MAX_LINE + 1)
        if not line:
            return
        if len(line) > MAX_LINE or not line.strip():
            raise ValueError("Oversized or blank JSONL record")
        row = parse_json(line)
        if (not isinstance(row, dict) or set(row) != {"prompt", "chosen", "rejected"}
                or any(not isinstance(v, str) or not v.strip() or len(v) > 64000
                       for v in row.values()) or row["chosen"] == row["rejected"]):
            raise ValueError("Invalid preference record")
        yield row

def iter_verified(shards, hosts, demo=False):
    for shard in shards:
        stream, size = verified_file(shard["url"], shard["sha256"], MAX_SHARD, hosts, demo)
        with stream:
            if size != shard["bytes"]:
                raise ValueError("Byte length mismatch")
            count = sum(1 for _ in read_rows(stream))
            if count != shard["rows"]:
                raise ValueError("Row count mismatch")
            stream.seek(0)
            yield from read_rows(stream)

def registry_anchor(rpc, address, dataset_id, block, gateway):
    from web3 import Web3
    if urlsplit(rpc).scheme != "https":
        raise ValueError("HTTPS RPC required")
    w3 = Web3(Web3.HTTPProvider(rpc, request_kwargs={"timeout": 30}))
    if w3.eth.chain_id != 5042:
        raise ValueError("Expected Arc mainnet")
    address = Web3.to_checksum_address(address)
    snapshot = w3.eth.get_block(block)
    if not w3.eth.get_code(address, block_identifier=block):
        raise ValueError("No contract at selected block")
    abi = [{"type": "function", "name": "datasets", "stateMutability": "view",
            "inputs": [{"name": "", "type": "uint256"}],
            "outputs": [{"name": name, "type": kind} for name, kind in [
                ("bountyId", "uint256"), ("dataSha256", "bytes32"),
                ("arweaveTxId", "string"), ("contributor", "address"),
                ("status", "uint8")]]}]
    abi.append({"type": "function", "name": "usedHashes", "stateMutability": "view",
                "inputs": [{"name": "", "type": "bytes32"}],
                "outputs": [{"name": "", "type": "bool"}]})
    registry = w3.eth.contract(address=address, abi=abi)
    result = registry.functions.datasets(dataset_id).call(
        block_identifier=block)
    bounty, digest, txid, contributor, status = result
    if status != 2 or bounty == 0 or not re.fullmatch(r"[A-Za-z0-9_-]{43}", txid):
        raise ValueError("Record is not a paid archival reference")
    if not registry.functions.usedHashes(digest).call(block_identifier=block):
        raise ValueError("Registry hash reservation missing")
    if w3.eth.get_block(block)["hash"] != snapshot["hash"]:
        raise ValueError("Block changed during registry read")
    print("Anchor block:", snapshot["hash"].hex(), "bounty:", bounty,
          "contributor:", contributor)
    return gateway.rstrip("/") + "/" + txid, bytes(digest).hex()

def make_demo(directory):
    row = {"prompt": "What is 2 + 2?", "chosen": "4", "rejected": "5"}
    raw = (json.dumps(row, sort_keys=True) + "\n").encode("utf-8")
    shard = directory / "train.jsonl"
    shard.write_bytes(raw)
    manifest = {"version": 1, "format": "jsonl", "split": "train", "shards": [{
        "url": shard.as_uri(), "sha256": hashlib.sha256(raw).hexdigest(),
        "bytes": len(raw), "rows": 1}]}
    encoded = json.dumps(manifest, sort_keys=True).encode("utf-8")
    path = directory / "manifest.json"
    path.write_bytes(encoded)
    return path.as_uri(), hashlib.sha256(encoded).hexdigest()

def run_framework(shards, hosts, demo, framework, rank=0, world_size=1):
    if framework == "core":
        for row in iter_verified(shards[rank::world_size], hosts, demo):
            print(row)
    elif framework == "hf":
        from datasets import IterableDataset
        dataset = IterableDataset.from_generator(iter_verified, gen_kwargs={
            "shards": shards[rank::world_size], "hosts": tuple(hosts), "demo": demo})
        for row in dataset:
            print(row)
    else:
        from torch.utils.data import DataLoader, IterableDataset, get_worker_info

        class VerifiedPreferences(IterableDataset):
            def __iter__(self):
                worker = get_worker_info()
                workers, worker_id = (worker.num_workers, worker.id) if worker else (1, 0)
                local = shards[rank::world_size][worker_id::workers]
                yield from iter_verified(local, hosts, demo)

        loader = DataLoader(VerifiedPreferences(), batch_size=2, num_workers=0)
        for batch in loader:
            print(batch)

def source_parser():
    parser = argparse.ArgumentParser()
    parser.add_argument("--demo", action="store_true")
    parser.add_argument("--manifest-url")
    parser.add_argument("--sha256")
    parser.add_argument("--registry")
    parser.add_argument("--dataset-id", type=int)
    parser.add_argument("--block", type=int)
    parser.add_argument("--rpc", default="https://rpc.mainnet.arc.io")
    parser.add_argument("--gateway", default="https://gateway.irys.xyz")
    parser.add_argument("--allow-host", action="append",
                        default=["arweave.net", "gateway.irys.xyz"])
    return parser

@contextmanager
def source_from_args(args):
    with tempfile.TemporaryDirectory() as directory:
        if args.demo:
            if any(value is not None for value in (
                    args.manifest_url, args.sha256, args.registry,
                    args.dataset_id, args.block)):
                raise ValueError("Demo must not mix external source arguments")
            url, digest = make_demo(Path(directory).resolve())
        elif args.registry:
            if (args.manifest_url or args.sha256 or args.dataset_id is None
                    or args.dataset_id < 1 or args.block is None or args.block < 0):
                raise ValueError("Registry mode requires dataset ID and block")
            check_url(args.gateway, args.allow_host)
            url, digest = registry_anchor(args.rpc, args.registry,
                                          args.dataset_id, args.block, args.gateway)
        else:
            if (not args.manifest_url or not args.sha256
                    or args.dataset_id is not None or args.block is not None):
                raise ValueError("Supply an authenticated manifest URL and digest")
            url, digest = args.manifest_url, args.sha256
        shards = load_manifest(url, digest, args.allow_host, args.demo)
        yield shards, tuple(args.allow_host), args.demo

def main():
    parser = argparse.ArgumentParser()
    parser.add_argument("--demo", action="store_true")
    parser.add_argument("--framework", choices=["core", "hf", "torch"], default="core")
    parser.add_argument("--manifest-url")
    parser.add_argument("--sha256")
    parser.add_argument("--registry")
    parser.add_argument("--dataset-id", type=int)
    parser.add_argument("--block", type=int)
    parser.add_argument("--rpc", default="https://rpc.mainnet.arc.io")
    parser.add_argument("--gateway", default="https://arweave.net")
    parser.add_argument("--allow-host", action="append", default=["arweave.net"])
    parser.add_argument("--rank", type=int, default=0)
    parser.add_argument("--world-size", type=int, default=1)
    args = parser.parse_args()
    if args.world_size < 1 or not 0 <= args.rank < args.world_size:
        parser.error("Invalid rank/world size")
    with tempfile.TemporaryDirectory() as temp:
        if args.demo:
            if args.registry or args.manifest_url or args.sha256:
                parser.error("Demo cannot use external artifacts")
            url, digest = make_demo(Path(temp).resolve())
        elif args.registry:
            if (args.dataset_id is None or args.dataset_id < 1 or args.block is None
                    or args.block < 0 or args.manifest_url or args.sha256):
                parser.error("Registry mode requires dataset ID and block only")
            check_url(args.gateway, args.allow_host)
            url, digest = registry_anchor(args.rpc, args.registry, args.dataset_id,
                                          args.block, args.gateway)
        else:
            if not args.manifest_url or not args.sha256:
                parser.error("Supply manifest URL and trusted SHA-256, or use --demo")
            url, digest = args.manifest_url, args.sha256
        shards = load_manifest(url, digest, args.allow_host, args.demo)
        run_framework(shards, args.allow_host, args.demo, args.framework,
                      args.rank, args.world_size)

if __name__ == "__main__":
    main()
```

The demo's manifest and data are generated together locally. It demonstrates integrity plumbing and framework integration, **not independent provenance or on-chain verification**. For external manifests, supply `--manifest-url` and `--sha256` from an authenticated release record. For registry mode, supply `--registry`, `--dataset-id`, and `--block` using actual deployment evidence. No invented address is provided.

The RPC helper checks chain identity, nonempty code, paid status, and consistency of the selected block during the read. It does not verify the runtime bytecode against a published release or establish finality independently. Those checks remain mandatory deployment preconditions. A malicious RPC can lie; use trusted infrastructure and independent reconciliation appropriate to the workload.

## Gateway selection and provenance evidence

The Arweave tooling published by Irys documents retrieval through `https://gateway.irys.xyz/` followed by the upload receipt identifier. This is an Arweave gateway convention; it must not be confused with the separate Irys L1 storage network. See the [Irys Arweave SDK reference](https://github.com/Irys-xyz/arweave-js-sdk). The examples allow that gateway and `arweave.net`, while retaining the expected digest across either retrieval path.

An indexing query can discover candidate artifacts by campaign tags, but tags and search results are not the trust root. Resolve the selected artifact against the authenticated registry record or signed release manifest. For bundled uploads, distinguish the data-item identifier from its containing Arweave transaction and establish inclusion and the release's archival-finality policy separately. Successful HTTP retrieval alone does not establish permanent inclusion.

The registry's `usedHashes(bytes32)` result is checked in the RPC helper as an additional consistency condition. A reserved hash can belong to a pending or rejected record. Therefore the helper also requires a full `datasets(datasetId)` record in `Paid` state, preserves its archive identifier, and detects a block-hash change during the read. There is no current `registeredHashes` ABI to invoke.

For a hash committing to an annotation rather than an aggregation manifest, this manifest loader is not the correct interpretation. A production aggregation service must authenticate membership from the paid annotation records to its training shards. A single accepted manifest record does not prove that every member has an independently paid contributor or that its transformations preserved meaning.

## Exact-byte commitments and staged streaming

SHA-256 is calculated over the exact downloaded artifact bytes. Whitespace, line endings, encoding, compression, and key order can change the digest even when two documents appear semantically equivalent. Do not parse, pretty-print, normalize, or decompress an artifact and then compare it with a commitment to its original representation.

This reference consumes uncompressed UTF-8 JSONL. Compressed archives need a defined outer commitment, bounded decompression, and explicit inner-shard commitments. Encrypted archives need authenticated decryption and separate ciphertext/plaintext integrity rules in an authorized environment. Neither format is silently accepted here.

```
Authenticated registry / release record
                |
                v
Download manifest -> verify its SHA-256 -> validate bounded descriptors
                |
                v
Download one shard to private temporary storage
                |
                v
Check hash + byte length + every row + row count
                |
                v
Yield records -> tokenize -> mask labels -> training objective
```

A whole-file hash cannot authenticate its first line before the final byte has arrived. Staging one bounded shard before yielding it is intentional. Earlier verified shards may already have been used if a later shard fails; training jobs should checkpoint and mark the run incomplete rather than pretend a failure rolls back optimizer updates. If all-or-nothing dataset admission is required, preflight every shard before starting training.

## Snippet 1: Hugging Face verified streaming

Save the next complete block as `hf_stream.py` alongside `verified_loader.py`. It uses the shared verification boundary rather than letting a framework fetch unchecked URLs. Local execution is `python hf_stream.py --demo`. For external data, use `--manifest-url` and `--sha256` from an authenticated release, or supply `--registry`, `--dataset-id`, and `--block` from verified deployment evidence. These are required runtime inputs, not invented deployment values.

```python
from datasets import Features, IterableDataset, Value
from verified_loader import iter_verified, source_from_args, source_parser


def main():
    parser = source_parser()
    args = parser.parse_args()
    features = Features({name: Value("string") for name in
                         ("prompt", "chosen", "rejected")})
    with source_from_args(args) as (shards, hosts, demo):
        dataset = IterableDataset.from_generator(
            iter_verified,
            gen_kwargs={"shards": shards, "hosts": hosts, "demo": demo},
            features=features,
        )
        count = 0
        for record in dataset:
            assert set(record) == {"prompt", "chosen", "rejected"}
            count += 1
        if count == 0:
            raise RuntimeError("No verified records")
        print("Verified Hugging Face records:", count)


if __name__ == "__main__":
    main()
```

The generator verifies each shard before exposing records and keeps record memory bounded. Explicit features avoid depending on implicit type inference. For `load_dataset("json", streaming=True)`, first verify artifacts into a protected local cache and supply those verified paths. Remote framework streaming alone does not enforce the registry commitment. [Hugging Face streaming reference](https://huggingface.co/docs/datasets/en/stream).

## Snippet 2: PyTorch tensors and active SFT / DPO training

Save this block as `train_verified.py` alongside the shared loader. It includes a deterministic byte tokenizer and a small causal recurrent model, so it can perform real optimizer steps without downloading model weights. It is an executable integration harness, not a claim that this small model or byte tokenizer is suitable for frontier training.

The SFT path masks prompt and padding targets and learns the chosen completion. The DPO path compares chosen/rejected completion log-probabilities against a frozen reference model using a logistic preference objective. Both consume only verified shards. These are alternatives selectable through `--objective`, not two losses silently combined.

```python
import copy
import torch
from torch import nn
from torch.nn import functional as F
from torch.utils.data import DataLoader, IterableDataset, get_worker_info
from verified_loader import iter_verified, source_from_args, source_parser

PAD, BOS, EOS, VOCAB = 0, 1, 2, 259


class VerifiedPreferences(IterableDataset):
    def __init__(self, shards, hosts, demo, rank=0, world_size=1):
        super().__init__()
        if world_size < 1 or not 0 <= rank < world_size:
            raise ValueError("Invalid rank or world size")
        self.shards = shards
        self.hosts = hosts
        self.demo = demo
        self.rank = rank
        self.world_size = world_size

    def __iter__(self):
        worker = get_worker_info()
        worker_id, workers = (worker.id, worker.num_workers) if worker else (0, 1)
        rank_shards = self.shards[self.rank::self.world_size]
        local_shards = rank_shards[worker_id::workers]
        yield from iter_verified(local_shards, self.hosts, self.demo)


def encode_pair(prompt, completion, max_length=2048):
    prefix = [BOS] + [byte + 3 for byte in (prompt + "\n").encode("utf-8")]
    sequence = prefix + [byte + 3 for byte in completion.encode("utf-8")] + [EOS]
    if len(sequence) > max_length:
        raise ValueError("Record exceeds configured context; no silent truncation")
    inputs, targets = sequence[:-1], sequence[1:]
    targets[:len(prefix) - 1] = [-100] * (len(prefix) - 1)
    return inputs, targets


def collate_preferences(rows):
    output = {}
    for candidate in ("chosen", "rejected"):
        pairs = [encode_pair(row["prompt"], row[candidate]) for row in rows]
        width = max(len(inputs) for inputs, _ in pairs)
        inputs = torch.full((len(rows), width), PAD, dtype=torch.long)
        labels = torch.full((len(rows), width), -100, dtype=torch.long)
        for index, (tokens, targets) in enumerate(pairs):
            inputs[index, :len(tokens)] = torch.tensor(tokens)
            labels[index, :len(targets)] = torch.tensor(targets)
        output[candidate] = (inputs, labels)
    return output


class TinyCausalModel(nn.Module):
    def __init__(self):
        super().__init__()
        self.embedding = nn.Embedding(VOCAB, 32, padding_idx=PAD)
        self.recurrent = nn.GRU(32, 48, batch_first=True)
        self.head = nn.Linear(48, VOCAB)

    def forward(self, tokens):
        hidden, _ = self.recurrent(self.embedding(tokens))
        return self.head(hidden)


def completion_logprob(model, inputs, labels):
    mask = labels.ne(-100)
    safe_labels = labels.masked_fill(~mask, 0)
    log_probs = F.log_softmax(model(inputs), dim=-1)
    selected = log_probs.gather(-1, safe_labels.unsqueeze(-1)).squeeze(-1)
    return (selected * mask).sum(dim=-1)


def main():
    parser = source_parser()
    parser.add_argument("--objective", choices=["sft", "dpo"], default="sft")
    parser.add_argument("--epochs", type=int, default=2)
    parser.add_argument("--workers", type=int, default=0)
    args = parser.parse_args()
    if args.epochs < 1 or args.workers < 0:
        parser.error("epochs must be positive and workers non-negative")
    torch.manual_seed(7)
    device = torch.device("cuda" if torch.cuda.is_available() else "cpu")
    model = TinyCausalModel().to(device)
    reference = copy.deepcopy(model).eval()
    for parameter in reference.parameters():
        parameter.requires_grad_(False)
    optimizer = torch.optim.AdamW(model.parameters(), lr=0.001)
    beta = 0.1
    steps = 0
    with source_from_args(args) as (shards, hosts, demo):
        data = VerifiedPreferences(shards, hosts, demo)
        loader = DataLoader(data, batch_size=2, num_workers=args.workers,
                            collate_fn=collate_preferences)
        for epoch in range(args.epochs):
            epoch_steps = 0
            for batch in loader:
                chosen_x, chosen_y = [t.to(device) for t in batch["chosen"]]
                rejected_x, rejected_y = [t.to(device) for t in batch["rejected"]]
                optimizer.zero_grad(set_to_none=True)
                if args.objective == "sft":
                    logits = model(chosen_x)
                    loss = F.cross_entropy(logits.reshape(-1, VOCAB),
                                           chosen_y.reshape(-1), ignore_index=-100)
                else:
                    policy_margin = (
                        completion_logprob(model, chosen_x, chosen_y)
                        - completion_logprob(model, rejected_x, rejected_y))
                    with torch.no_grad():
                        reference_margin = (
                            completion_logprob(reference, chosen_x, chosen_y)
                            - completion_logprob(reference, rejected_x, rejected_y))
                    loss = -F.logsigmoid(beta * (policy_margin - reference_margin)).mean()
                if not torch.isfinite(loss):
                    raise RuntimeError("Non-finite loss")
                loss.backward()
                nn.utils.clip_grad_norm_(model.parameters(), 1.0, error_if_nonfinite=True)
                optimizer.step()
                steps += 1
                epoch_steps += 1
                print("epoch", epoch, "step", steps, "loss", round(loss.item(), 6))
            if epoch_steps == 0:
                raise RuntimeError("Epoch contained no verified training samples")
    print("Completed verified optimizer steps:", steps)


if __name__ == "__main__":
    main()
```

Run both complete objectives against the local fixture:

```bash
python hf_stream.py --demo
python train_verified.py --demo --objective sft
python train_verified.py --demo --objective dpo
python train_verified.py --demo --objective sft --workers 2
```

The byte vocabulary reserves padding, beginning-of-sequence, and end-of-sequence identifiers separately. Inputs are shifted once against targets. Completion masks include EOS and exclude prompt and padding loss. Long records fail explicitly rather than silently dropping the answer. Replace this demonstrator with the training model's pinned tokenizer, chat template, and causal model when integrating a real run; preserve the verification and masking boundaries.

For standard DPO, preserve a frozen reference, consistent candidate context, and a documented beta. Ties and `neither` labels cannot be converted to strict chosen/rejected pairs without adjudication. A Golden Patch used as the chosen response needs lineage to the accepted correction. This example implements the standard pairwise logistic objective described in the [DPO paper](https://arxiv.org/abs/2305.18290), not every variant or trainer feature.

## Worker, rank, and restart behavior

The dataset class is module-scoped and partitions shards by rank and then worker. It yields CPU text; tensor conversion occurs in collation and device transfer in the main training loop. This avoids repeating every shard in every worker. PyTorch documents iterable-dataset replication and worker-aware partitioning in its [data-loading reference](https://docs.pytorch.org/docs/stable/data.html).

The supplied training harness is single-process training. Its dataset supports rank partitioning for integration, but it does not initialize distributed collectives, coordinate unequal rank lengths, or implement distributed checkpoint recovery. Shard balance matters: disjoint input alone does not guarantee equal optimizer-step counts. Do not enable DDP without an explicit length/join strategy and tested resume semantics.

Each new iterator verifies and rereads the assigned shards. The harness does not persist model checkpoints or a resumable data cursor. A production run must save model, optimizer, scheduler, random-generator state, tokenizer version, manifest digest, rank layout, and shard/record progress. Exactly-once payment identity and exactly-once optimizer updates are separate problems.

## Operational failure policy

Fail closed on hash, schema, size, row-count, or receipt-state mismatches. Do not replace a failed artifact with an unverified alternate URL. A different approved gateway is acceptable only if the same authenticated digest remains mandatory. An RPC can lie about state; use trusted infrastructure, release-bytecode verification, and independent reconciliation appropriate to the value at risk.

Per-read timeouts and artifact limits do not provide a global wall-clock or disk quota. Enforce those at the job runner. The reference uses temporary plaintext storage and supports only public unencrypted data; private archives require verified authorization, encryption at rest, decryption key custody, and secure cleanup beyond this example.

Training-run provenance should record what was verified, the accepted block, artifact identities, schema and transformation versions, and framework versions. Log identifiers and failure categories rather than sensitive record bodies. Passing the local fixture establishes executable integration; it does not prove a deployed registry, public gateway availability, archival permanence, or production training quality.

## Example validation record

The published code blocks were extracted and exercised with Python 3.12, `datasets` 5.1.0, `torch` 2.14.1, and `web3` 8.0.0 in an isolated test environment. Core and Hugging Face fixture iteration passed. SFT and DPO each completed two finite-loss optimizer steps; SFT also passed with two loader workers. Hash, size, and row-count corruption was rejected before yielding a record. Strict JSON parsing rejected duplicate keys, non-finite constants, and invalid encodings. These tests used local synthetic artifacts; live registry reads and gateway retrieval were not end-to-end deployment tests.


---

# Agent Instructions
This documentation is published with GitBook. GitBook is the documentation platform designed so that both humans and AI agents can read, navigate, and reason over technical content effectively. Learn more at gitbook.com.

## Querying This Documentation
If you need additional information that is not directly available in this page, you can query the documentation dynamically by asking a question.

Perform an HTTP GET request on the current page URL with the `ask` query parameter, and the optional `goal` query parameter:

```
GET https://docs.arcv.network/3.-enterprise-and-frontier-ai-labs/3.3-provenance-and-dataloaders.md?ask=<question>&goal=<endgoal>
```

`ask` is the immediate question: it should be specific, self-contained, and written in natural language.
`goal` is optional and describes the broader end goal you are ultimately trying to accomplish on behalf of the user. GitBook uses it to tailor the answer towards what is most useful for that goal.

The response will contain a direct answer to the question and relevant excerpts and sources from the documentation.

Use this mechanism when the answer is not explicitly present in the current page, you need clarification or additional context, or you want to retrieve related documentation sections.
