feat(corpus): add deterministic normalization and chunking
This commit is contained in:
@@ -0,0 +1,39 @@
|
||||
# Evidence Task 3 — deterministic normalization and chunking
|
||||
|
||||
## Outcome
|
||||
|
||||
- Added pure `normalize(acquired, pipeline_version)` and `chunk(document, policy)` transforms.
|
||||
- Normalization enforces UTF-8 (including UTF-8 BOM), a 10 MiB input ceiling, LF line endings,
|
||||
NFC Unicode, safe YAML frontmatter extraction, canonical provenance URIs, and hashes the exact
|
||||
canonical UTF-8 text stored on the document.
|
||||
- Undecodable, unsupported-charset, oversized, and invalid-frontmatter inputs fail explicitly;
|
||||
byte content is never truncated.
|
||||
- Chunking uses a versioned immutable policy, paragraph/word boundaries with deterministic
|
||||
character-count hard splits for long tokens, contiguous ordinals, provenance metadata, exact
|
||||
per-chunk hashes, and IDs derived from document hash + ordinal + policy version.
|
||||
- Empty documents produce no chunks. Non-ASCII, CRLF equivalence, repeatability, policy changes,
|
||||
duplicate-content ordinal collisions, and max-character limits are covered by tests.
|
||||
|
||||
## TDD evidence
|
||||
|
||||
- Initial focused test run failed during collection because both transform modules were absent.
|
||||
- The EOF-frontmatter edge test was separately observed failing before its implementation.
|
||||
- Final focused verification: `12 passed`.
|
||||
|
||||
## Verification
|
||||
|
||||
- `cd harness && .venv/bin/pytest tests/test_corpus_normalize.py tests/test_corpus_chunk.py -q`
|
||||
— **12 passed**.
|
||||
- `cd harness && .venv/bin/pytest -q` — **573 passed, 5 deselected**. The sandboxed attempt could
|
||||
not access Docker; the approved rerun with local Docker access passed.
|
||||
- Targeted Ruff over all four implementation/test files — **clean**.
|
||||
- Full `cd harness && .venv/bin/ruff check .` — reports **34 pre-existing errors** in unrelated
|
||||
legacy tests (unused imports and existing E702 semicolon lines); none are in Task 3 files.
|
||||
|
||||
## Concerns
|
||||
|
||||
- The 10 MiB normalization ceiling is deliberately explicit and independent of adapter download
|
||||
limits. If deployment policy needs a different ceiling, it should become a versioned pipeline
|
||||
configuration before ingestion is wired.
|
||||
- Character limits use Python Unicode code points (`len`), not UTF-8 bytes or tokenizer tokens;
|
||||
this is recorded in the chunk-policy metadata and tested with non-ASCII content.
|
||||
@@ -0,0 +1,71 @@
|
||||
import hashlib
|
||||
|
||||
import pytest
|
||||
|
||||
from tht.corpus.chunk import ChunkPolicy, chunk
|
||||
from tht.corpus.models import CanonicalDocument
|
||||
|
||||
|
||||
def document(content: str) -> CanonicalDocument:
|
||||
normalized = content.replace("\r\n", "\n").replace("\r", "\n")
|
||||
digest = hashlib.sha256(normalized.encode()).hexdigest()
|
||||
return CanonicalDocument(
|
||||
document_id="doc:abc",
|
||||
source_id="source:a",
|
||||
source_uri="https://host/a.md",
|
||||
source_fingerprint="etag:abc",
|
||||
content_hash=f"sha256:{digest}",
|
||||
title="A",
|
||||
content=normalized,
|
||||
media_type="text/markdown",
|
||||
pipeline_version="pipe:v1",
|
||||
metadata={"owner": "docs"},
|
||||
)
|
||||
|
||||
|
||||
def test_chunk_ids_are_stable_for_same_content_and_repeat_runs():
|
||||
policy = ChunkPolicy(version="paragraph:v1", max_chars=8)
|
||||
first = chunk(document("A\n\nB"), policy)
|
||||
second = chunk(document("A\r\n\r\nB"), policy)
|
||||
repeated = chunk(document("A\n\nB"), policy)
|
||||
assert first == second == repeated
|
||||
|
||||
|
||||
def test_policy_version_changes_ids_without_changing_boundaries():
|
||||
doc = document("alpha\n\nbeta")
|
||||
first = chunk(doc, ChunkPolicy(version="paragraph:v1", max_chars=6))
|
||||
second = chunk(doc, ChunkPolicy(version="paragraph:v2", max_chars=6))
|
||||
assert [item.content for item in first] == [item.content for item in second]
|
||||
assert [item.chunk_id for item in first] != [item.chunk_id for item in second]
|
||||
|
||||
|
||||
def test_long_non_ascii_tokens_are_hard_split_by_unicode_characters():
|
||||
chunks = chunk(document("ééééé世界"), ChunkPolicy(version="chars:v1", max_chars=3))
|
||||
assert [item.content for item in chunks] == ["ééé", "éé世", "界"]
|
||||
assert all(len(item.content) <= 3 for item in chunks)
|
||||
|
||||
|
||||
def test_chunks_have_contiguous_ordinals_hashes_and_provenance_metadata():
|
||||
doc = document("alpha beta gamma")
|
||||
chunks = chunk(doc, ChunkPolicy(version="words:v1", max_chars=7))
|
||||
assert [item.ordinal for item in chunks] == list(range(len(chunks)))
|
||||
assert len({item.chunk_id for item in chunks}) == len(chunks)
|
||||
for item in chunks:
|
||||
assert item.source_uri == doc.source_uri
|
||||
assert item.document_id == doc.document_id
|
||||
assert item.pipeline_version == doc.pipeline_version
|
||||
assert item.metadata["chunk_policy"] == {"max_chars": 7, "version": "words:v1"}
|
||||
assert item.metadata["document"] == {"owner": "docs"}
|
||||
assert item.content_hash == "sha256:" + hashlib.sha256(item.content.encode()).hexdigest()
|
||||
|
||||
|
||||
def test_duplicate_chunk_content_cannot_collide_across_ordinals():
|
||||
chunks = chunk(document("same\n\nsame"), ChunkPolicy(version="paragraph:v1", max_chars=8))
|
||||
assert [item.content for item in chunks] == ["same", "same"]
|
||||
assert chunks[0].chunk_id != chunks[1].chunk_id
|
||||
|
||||
|
||||
def test_empty_document_has_no_chunks_and_invalid_policy_is_rejected():
|
||||
assert chunk(document(""), ChunkPolicy(version="v1", max_chars=4)) == []
|
||||
with pytest.raises(ValueError):
|
||||
ChunkPolicy(version="v1", max_chars=0)
|
||||
@@ -0,0 +1,67 @@
|
||||
import hashlib
|
||||
from datetime import UTC, datetime
|
||||
|
||||
import pytest
|
||||
|
||||
from tht.corpus.normalize import MAX_DOCUMENT_BYTES, PermanentNormalizationError, normalize
|
||||
from tht.ports.evidence import AcquiredDocument, SourceObject
|
||||
|
||||
|
||||
def acquired(content: bytes, *, media_type: str = "text/markdown") -> AcquiredDocument:
|
||||
return AcquiredDocument(
|
||||
source=SourceObject(
|
||||
source_id="source:guide",
|
||||
uri="https://host/guide.md?signature=transport#part",
|
||||
fingerprint="etag:abc",
|
||||
modified_at=datetime(2026, 7, 12, 12, 0, tzinfo=UTC),
|
||||
metadata={"owner": "docs"},
|
||||
),
|
||||
content=content,
|
||||
media_type=media_type,
|
||||
metadata={"transport": "http"},
|
||||
)
|
||||
|
||||
|
||||
def test_normalize_utf8_bom_newlines_unicode_and_frontmatter():
|
||||
raw = (
|
||||
"\ufeff---\r\ntitle: Café\r\ntags: [uno, due]\r\n---\r\n"
|
||||
"Cafe\u0301\rBody\r\n"
|
||||
).encode()
|
||||
|
||||
document = normalize(acquired(raw, media_type="text/markdown; charset=UTF-8"), "pipe:v1")
|
||||
|
||||
assert document.content == "Café\nBody\n"
|
||||
assert document.title == "Café"
|
||||
assert document.metadata["frontmatter"] == {"tags": ("uno", "due"), "title": "Café"}
|
||||
assert document.metadata["source"] == {"owner": "docs"}
|
||||
assert document.metadata["acquisition"] == {"transport": "http"}
|
||||
assert document.source_uri == "https://host/guide.md"
|
||||
assert document.modified_at == datetime(2026, 7, 12, 12, 0, tzinfo=UTC)
|
||||
assert document.content_hash == "sha256:" + hashlib.sha256(document.content.encode()).hexdigest()
|
||||
|
||||
|
||||
def test_plain_text_that_only_resembles_frontmatter_is_not_dropped():
|
||||
document = normalize(acquired(b"---\nnot: closed\nbody"), "pipe:v1")
|
||||
assert document.content == "---\nnot: closed\nbody"
|
||||
assert "frontmatter" not in document.metadata
|
||||
|
||||
|
||||
def test_frontmatter_can_end_at_eof_without_inventing_content():
|
||||
document = normalize(acquired(b"---\ntitle: Empty\n---"), "pipe:v1")
|
||||
assert document.title == "Empty"
|
||||
assert document.content == ""
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("content", "media_type", "reason"),
|
||||
[
|
||||
(b"bad: \xff", "text/plain", "undecodable"),
|
||||
(b"hello", "text/plain; charset=iso-8859-1", "unsupported_charset"),
|
||||
(b"x" * (MAX_DOCUMENT_BYTES + 1), "text/plain", "oversized"),
|
||||
],
|
||||
)
|
||||
def test_rejects_invalid_input_as_typed_permanent_error(content, media_type, reason):
|
||||
with pytest.raises(PermanentNormalizationError) as caught:
|
||||
normalize(acquired(content, media_type=media_type), "pipe:v1")
|
||||
assert caught.value.permanent is True
|
||||
assert caught.value.reason == reason
|
||||
@@ -0,0 +1,89 @@
|
||||
"""Versioned deterministic chunking for canonical corpus documents."""
|
||||
|
||||
import hashlib
|
||||
import re
|
||||
from dataclasses import dataclass
|
||||
|
||||
from tht.corpus.models import CanonicalChunk, CanonicalDocument
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class ChunkPolicy:
|
||||
version: str
|
||||
max_chars: int
|
||||
|
||||
def __post_init__(self) -> None:
|
||||
if not self.version:
|
||||
raise ValueError("chunk policy version must not be empty")
|
||||
if self.max_chars <= 0:
|
||||
raise ValueError("max_chars must be greater than zero")
|
||||
|
||||
|
||||
def _hash(text: str) -> str:
|
||||
return hashlib.sha256(text.encode("utf-8")).hexdigest()
|
||||
|
||||
|
||||
def _hard_split(text: str, maximum: int) -> list[str]:
|
||||
return [text[start : start + maximum] for start in range(0, len(text), maximum)]
|
||||
|
||||
|
||||
def _split_block(block: str, maximum: int) -> list[str]:
|
||||
if len(block) <= maximum:
|
||||
return [block]
|
||||
tokens = re.findall(r"\S+", block)
|
||||
chunks: list[str] = []
|
||||
current = ""
|
||||
for token in tokens:
|
||||
if len(token) > maximum:
|
||||
if current:
|
||||
chunks.append(current)
|
||||
current = ""
|
||||
chunks.extend(_hard_split(token, maximum))
|
||||
continue
|
||||
candidate = f"{current} {token}" if current else token
|
||||
if len(candidate) <= maximum:
|
||||
current = candidate
|
||||
else:
|
||||
chunks.append(current)
|
||||
current = token
|
||||
if current:
|
||||
chunks.append(current)
|
||||
return chunks
|
||||
|
||||
|
||||
def _contents(content: str, maximum: int) -> list[str]:
|
||||
if not content:
|
||||
return []
|
||||
result: list[str] = []
|
||||
for block in re.split(r"\n{2,}", content):
|
||||
if block:
|
||||
result.extend(_split_block(block, maximum))
|
||||
return result
|
||||
|
||||
|
||||
def chunk(document: CanonicalDocument, policy: ChunkPolicy) -> list[CanonicalChunk]:
|
||||
"""Split canonical text with stable character-count boundaries and identifiers."""
|
||||
chunks: list[CanonicalChunk] = []
|
||||
for ordinal, content in enumerate(_contents(document.content, policy.max_chars)):
|
||||
identifier = _hash(f"{document.content_hash}:{ordinal}:{policy.version}")
|
||||
chunks.append(
|
||||
CanonicalChunk(
|
||||
chunk_id=f"chunk:{identifier}",
|
||||
document_id=document.document_id,
|
||||
ordinal=ordinal,
|
||||
content=content,
|
||||
content_hash=f"sha256:{_hash(content)}",
|
||||
source_uri=document.source_uri,
|
||||
pipeline_version=document.pipeline_version,
|
||||
metadata={
|
||||
"chunk_policy": {
|
||||
"version": policy.version,
|
||||
"max_chars": policy.max_chars,
|
||||
},
|
||||
"document": document.model_dump(mode="json")["metadata"],
|
||||
"source_fingerprint": document.source_fingerprint,
|
||||
"title": document.title,
|
||||
},
|
||||
)
|
||||
)
|
||||
return chunks
|
||||
@@ -0,0 +1,99 @@
|
||||
"""Pure, deterministic conversion of acquired bytes into canonical text."""
|
||||
|
||||
import hashlib
|
||||
import re
|
||||
import unicodedata
|
||||
from collections.abc import Mapping
|
||||
|
||||
import yaml
|
||||
from pydantic import JsonValue, TypeAdapter, ValidationError
|
||||
|
||||
from tht.corpus.models import CanonicalDocument
|
||||
from tht.ports.evidence import AcquiredDocument, canonical_provenance_uri
|
||||
|
||||
|
||||
MAX_DOCUMENT_BYTES = 10 * 1024 * 1024
|
||||
_CHARSET = re.compile(r"(?:^|;)\s*charset\s*=\s*[\"']?([^;\s\"']+)", re.IGNORECASE)
|
||||
_FRONTMATTER = re.compile(r"\A---\n(.*?)\n---(?:\n|\Z)", re.DOTALL)
|
||||
_JSON_OBJECT = TypeAdapter(dict[str, JsonValue])
|
||||
|
||||
|
||||
class PermanentNormalizationError(ValueError):
|
||||
"""A deterministic input failure which retrying cannot repair."""
|
||||
|
||||
def __init__(self, reason: str) -> None:
|
||||
super().__init__(f"document normalization failed: {reason}")
|
||||
self.reason = reason
|
||||
self.permanent = True
|
||||
|
||||
|
||||
def _sha256(value: str) -> str:
|
||||
return hashlib.sha256(value.encode("utf-8")).hexdigest()
|
||||
|
||||
|
||||
def _decode(acquired: AcquiredDocument) -> str:
|
||||
if len(acquired.content) > MAX_DOCUMENT_BYTES:
|
||||
raise PermanentNormalizationError("oversized")
|
||||
|
||||
media_type = acquired.media_type or "text/plain"
|
||||
charset = _CHARSET.search(media_type)
|
||||
if charset and charset.group(1).lower().replace("_", "-") not in {
|
||||
"utf-8",
|
||||
"utf8",
|
||||
"us-ascii",
|
||||
"ascii",
|
||||
}:
|
||||
raise PermanentNormalizationError("unsupported_charset")
|
||||
try:
|
||||
return acquired.content.decode("utf-8-sig", errors="strict")
|
||||
except UnicodeDecodeError as error:
|
||||
raise PermanentNormalizationError("undecodable") from error
|
||||
|
||||
|
||||
def _frontmatter(text: str) -> tuple[dict[str, JsonValue], str]:
|
||||
match = _FRONTMATTER.match(text)
|
||||
if match is None:
|
||||
return {}, text
|
||||
try:
|
||||
loaded = yaml.safe_load(match.group(1))
|
||||
if loaded is None:
|
||||
loaded = {}
|
||||
if not isinstance(loaded, Mapping):
|
||||
raise TypeError("frontmatter is not a mapping")
|
||||
metadata = _JSON_OBJECT.validate_python(dict(loaded))
|
||||
except (TypeError, UnicodeError, ValidationError, yaml.YAMLError) as error:
|
||||
raise PermanentNormalizationError("invalid_frontmatter") from error
|
||||
return metadata, text[match.end() :]
|
||||
|
||||
|
||||
def normalize(acquired: AcquiredDocument, pipeline_version: str) -> CanonicalDocument:
|
||||
"""Normalize one transport result without I/O or implicit data loss."""
|
||||
if not pipeline_version:
|
||||
raise ValueError("pipeline_version must not be empty")
|
||||
|
||||
decoded = _decode(acquired)
|
||||
canonical = unicodedata.normalize("NFC", decoded.replace("\r\n", "\n").replace("\r", "\n"))
|
||||
frontmatter, content = _frontmatter(canonical)
|
||||
source_uri = canonical_provenance_uri(acquired.source.uri)
|
||||
identity = f"{acquired.source.source_id}\n{source_uri}"
|
||||
media_type = (acquired.media_type or "text/plain").split(";", 1)[0].strip().lower()
|
||||
metadata: dict[str, JsonValue] = {
|
||||
"source": acquired.source.model_dump(mode="json")["metadata"],
|
||||
"acquisition": acquired.model_dump(mode="json")["metadata"],
|
||||
}
|
||||
if frontmatter:
|
||||
metadata["frontmatter"] = frontmatter
|
||||
|
||||
return CanonicalDocument(
|
||||
document_id=f"doc:{_sha256(identity)}",
|
||||
source_id=acquired.source.source_id,
|
||||
source_uri=source_uri,
|
||||
source_fingerprint=acquired.source.fingerprint,
|
||||
content_hash=f"sha256:{_sha256(content)}",
|
||||
title=str(frontmatter.get("title", "")),
|
||||
content=content,
|
||||
media_type=media_type,
|
||||
modified_at=acquired.source.modified_at,
|
||||
pipeline_version=pipeline_version,
|
||||
metadata=metadata,
|
||||
)
|
||||
Reference in New Issue
Block a user