# Evidence Sources and Preprocessing Implementation Plan > **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. **Goal:** Materialize heterogeneous Evidence sources into a versioned canonical corpus and index it through resumable, idempotent jobs. **Architecture:** Source adapters only discover and acquire. Pure normalization/chunking stages create immutable version artifacts; vector indexing writes a staging generation; publish atomically switches the active manifest. DWH preprocessing uses the same job envelope but a separate pipeline. **Tech Stack:** Python 3.11+, Pydantic 2, requests, optional `fsspec`/S3 client, existing embeddings/vector ports, Typer, pytest. ## Global Constraints - Plans 1-3 are complete. - Runtime search reads only the active canonical corpus and vector generation. - Filesystem and HTTP sources are MVP; S3-compatible follows on the same port. - Source credentials never enter corpus metadata or logs. - Failed runs never replace the last valid published generation. - Document and DWH preprocessing are separate jobs with shared lock/report infrastructure. --- ### Task 1: Define Evidence source port and canonical records **Files:** - Create: `harness/tht/ports/evidence.py` - Create: `harness/tht/corpus/models.py` - Test: `harness/tests/test_evidence_port_contract.py` - Test: `harness/tests/test_corpus_models.py` **Interfaces:** - Produces: `EvidenceSource`, `SourceObject`, `AcquiredDocument`, `CanonicalDocument`, `CanonicalChunk`, `CorpusManifest`. - [ ] **Step 1: Write serialization and secret-exclusion tests** ```python def test_manifest_contains_provenance_without_credentials(): manifest = CorpusManifest(documents=[document(source_uri="https://host/a.md")]) payload = manifest.model_dump_json() assert "https://host/a.md" in payload assert "api_key" not in payload ``` - [ ] **Step 2: Verify failure** Run: `cd harness && .venv/bin/pytest tests/test_evidence_port_contract.py tests/test_corpus_models.py -q` Expected: FAIL. - [ ] **Step 3: Define immutable records and source protocol** ```python @runtime_checkable class EvidenceSource(Protocol): def discover(self) -> Iterable[SourceObject]: ... def acquire(self, item: SourceObject) -> AcquiredDocument: ... class SourceObject(BaseModel, frozen=True): source_id: str uri: str fingerprint: str modified_at: datetime | None = None metadata: dict[str, JsonValue] = {} ``` - [ ] **Step 4: Run model and protocol tests** Run: `cd harness && .venv/bin/pytest tests/test_evidence_port_contract.py tests/test_corpus_models.py -q` Expected: PASS. - [ ] **Step 5: Commit** ```bash git add harness/tht/ports/evidence.py harness/tht/corpus/models.py harness/tests/test_evidence_port_contract.py harness/tests/test_corpus_models.py git commit -m "feat(evidence): define source and corpus contracts" ``` ### Task 2: Implement filesystem and HTTP source adapters **Files:** - Create: `harness/tht/adapters/evidence/__init__.py` - Create: `harness/tht/adapters/evidence/filesystem.py` - Create: `harness/tht/adapters/evidence/http.py` - Modify: `harness/tht/config.py` - Modify: `harness/tht/adapters/factory.py` - Test: `harness/tests/test_filesystem_evidence_source.py` - Test: `harness/tests/test_http_evidence_source.py` **Interfaces:** - Produces: `FilesystemEvidenceSource`, `HttpManifestEvidenceSource`, `build_evidence_sources(cfg)`. - [ ] **Step 1: Test deterministic discovery and conditional HTTP acquisition** ```python def test_filesystem_discovery_is_stable(source): assert [x.uri for x in source.discover()] == sorted(x.uri for x in source.discover()) def test_http_uses_etag_for_fingerprint(http_source): assert next(http_source.discover()).fingerprint == 'etag:"abc"' ``` - [ ] **Step 2: Verify failure** Run: `cd harness && .venv/bin/pytest tests/test_filesystem_evidence_source.py tests/test_http_evidence_source.py -q` Expected: FAIL. - [ ] **Step 3: Implement adapters with bounded reads** ```python class FilesystemEvidenceSource: def discover(self): for path in sorted(self.root.rglob("*.md")): yield SourceObject(source_id=stable_id(path), uri=path.as_uri(), fingerprint=sha256_file(path)) ``` HTTP uses a declared manifest of URLs, connect/read timeouts, maximum bytes, ETag/Last-Modified where available, and content hashing as fallback. - [ ] **Step 4: Run adapter tests** Run: `cd harness && .venv/bin/pytest tests/test_filesystem_evidence_source.py tests/test_http_evidence_source.py tests/test_config_resources.py -q` Expected: PASS. - [ ] **Step 5: Commit** ```bash git add harness/tht/adapters/evidence harness/tht/config.py harness/tht/adapters/factory.py harness/tests/test_filesystem_evidence_source.py harness/tests/test_http_evidence_source.py git commit -m "feat(evidence): add filesystem and HTTP sources" ``` ### Task 3: Build deterministic normalization and chunking **Files:** - Create: `harness/tht/corpus/normalize.py` - Create: `harness/tht/corpus/chunk.py` - Test: `harness/tests/test_corpus_normalize.py` - Test: `harness/tests/test_corpus_chunk.py` **Interfaces:** - Produces: `normalize(acquired, pipeline_version) -> CanonicalDocument` and `chunk(document, policy) -> list[CanonicalChunk]`. - [ ] **Step 1: Pin UTF-8, frontmatter, line-ending, and stable chunk-id behavior** ```python def test_chunk_ids_are_stable_for_same_content(): first = chunk(document("A\n\nB"), policy(max_chars=8)) second = chunk(document("A\r\n\r\nB"), policy(max_chars=8)) assert [x.chunk_id for x in first] == [x.chunk_id for x in second] ``` - [ ] **Step 2: Verify failure** Run: `cd harness && .venv/bin/pytest tests/test_corpus_normalize.py tests/test_corpus_chunk.py -q` Expected: FAIL. - [ ] **Step 3: Implement pure deterministic transforms** ```python chunk_id = sha256(f"{document.content_hash}:{ordinal}:{policy.version}".encode()).hexdigest() ``` Reject undecodable or oversized content with a typed permanent error; never silently truncate source documents. - [ ] **Step 4: Run normalization/chunk tests** Run: `cd harness && .venv/bin/pytest tests/test_corpus_normalize.py tests/test_corpus_chunk.py -q` Expected: PASS. - [ ] **Step 5: Commit** ```bash git add harness/tht/corpus/normalize.py harness/tht/corpus/chunk.py harness/tests/test_corpus_normalize.py harness/tests/test_corpus_chunk.py git commit -m "feat(corpus): add deterministic normalization and chunking" ``` ### Task 4: Add shared job envelope, locking, and reports **Files:** - Create: `harness/tht/jobs/models.py` - Create: `harness/tht/jobs/runner.py` - Create: `harness/tht/jobs/locking.py` - Test: `harness/tests/test_job_runner.py` - Test: `harness/tests/test_job_locking.py` **Interfaces:** - Produces: `JobSpec`, `JobRun`, `JobReport`, `WorkspaceJobLock`, `run_job(spec, stages)`. - [ ] **Step 1: Test lock exclusion, resume, and JSON report schema** ```python def test_failed_stage_is_resumable(tmp_path): first = run_job(spec, [ok_stage, failing_stage]) second = run_job(spec.with_resume(first.run_id), [ok_stage, recovered_stage]) assert second.resumed_from == first.run_id assert second.status == "succeeded" ``` - [ ] **Step 2: Verify failure** Run: `cd harness && .venv/bin/pytest tests/test_job_runner.py tests/test_job_locking.py -q` Expected: FAIL. - [ ] **Step 3: Implement atomic report writes and workspace-scoped locks** ```python tmp = report_path.with_suffix(".tmp") tmp.write_text(report.model_dump_json(indent=2)) tmp.replace(report_path) ``` Persist checkpoints after each stage; ensure a crashed process leaves the active corpus untouched. - [ ] **Step 4: Run job infrastructure tests** Run: `cd harness && .venv/bin/pytest tests/test_job_runner.py tests/test_job_locking.py -q` Expected: PASS. - [ ] **Step 5: Commit** ```bash git add harness/tht/jobs harness/tests/test_job_runner.py harness/tests/test_job_locking.py git commit -m "feat(jobs): add resumable preprocessing envelope" ``` ### Task 5: Implement incremental document pipeline and atomic publish **Files:** - Create: `harness/tht/corpus/pipeline.py` - Create: `harness/tht/corpus/store.py` - Create: `harness/tht/cli/preprocess_cmd.py` - Modify: `harness/tht/cli/__init__.py` - Create: `harness/tht/search/evidence.py` - Test: `harness/tests/test_corpus_pipeline.py` - Test: `harness/tests/test_corpus_publish.py` - Test: `harness/tests/test_preprocess_cli.py` **Interfaces:** - Produces: `tht preprocess evidence [--dry-run] [--resume RUN_ID] [--json]`; active pointer `corpus//ACTIVE`. - Consumes: Evidence sources, canonical transforms, embedder, and `VectorStore`. - [ ] **Step 1: Test incremental skip and failed-run isolation** ```python def test_failed_generation_does_not_replace_active(corpus_store, pipeline): old = corpus_store.publish(valid_generation()) with pytest.raises(StageError): pipeline.run(source_with_failure()) assert corpus_store.active_generation() == old ``` - [ ] **Step 2: Verify failure** Run: `cd harness && .venv/bin/pytest tests/test_corpus_pipeline.py tests/test_corpus_publish.py tests/test_preprocess_cli.py -q` Expected: FAIL. - [ ] **Step 3: Implement staged generations and compare fingerprints** ```python changed = { item.source_id for item in discovered if previous.fingerprints.get(item.source_id) != item.fingerprint } ``` Index changed chunks, mark removed documents, validate counts and embedding dimensions, then atomically replace `ACTIVE`. Make runtime Evidence lookup resolve files through the active manifest instead of `rglob` on the source directory. - [ ] **Step 4: Run pipeline and existing search/session tests** Run: `cd harness && .venv/bin/pytest tests/test_corpus_pipeline.py tests/test_corpus_publish.py tests/test_preprocess_cli.py tests/test_search_pack.py tests/test_session_documents.py -q` Expected: PASS. - [ ] **Step 5: Commit** ```bash git add harness/tht/corpus harness/tht/cli/preprocess_cmd.py harness/tht/cli/__init__.py harness/tht/search/evidence.py harness/tests/test_corpus_pipeline.py harness/tests/test_corpus_publish.py harness/tests/test_preprocess_cli.py git commit -m "feat(preprocess): publish incremental Evidence corpus" ``` ### Task 6: Move DWH introspection and LSH into separate jobs **Files:** - Create: `harness/tht/jobs/dwh_pipeline.py` - Modify: `harness/tht/cli/schema_cmd.py` - Modify: `harness/tht/cli/lsh_cmd.py` - Modify: `harness/tht/cli/preprocess_cmd.py` - Test: `harness/tests/test_dwh_preprocess_job.py` - Test: `harness/tests/test_lsh_job_resume.py` **Interfaces:** - Produces: `tht preprocess dwh --steps introspect,lsh [--resume RUN_ID] [--json]`. - [ ] **Step 1: Test independent document/DWH locks and LSH resume** ```python def test_dwh_and_evidence_jobs_have_distinct_lock_names(): assert job_lock_name("demo", "dwh") != job_lock_name("demo", "evidence") ``` - [ ] **Step 2: Verify failure** Run: `cd harness && .venv/bin/pytest tests/test_dwh_preprocess_job.py tests/test_lsh_job_resume.py -q` Expected: FAIL. - [ ] **Step 3: Wrap existing commands as stages without changing their core algorithms** ```python stages = { "introspect": lambda ctx: refresh_catalog(ctx.dwh, ctx.paths.artifacts), "lsh": lambda ctx: build_lsh(ctx.dwh, ctx.paths.indexes, ctx.checkpoint), } ``` - [ ] **Step 4: Run DWH/LSH regression and full harness suite** Run: `cd harness && .venv/bin/pytest tests/test_dwh_preprocess_job.py tests/test_lsh_job_resume.py tests/test_schema_introspect_guard.py tests/l0/test_db_sampling.py -q` Run: `cd harness && .venv/bin/pytest -q` Expected: PASS. - [ ] **Step 5: Commit** ```bash git add harness/tht/jobs/dwh_pipeline.py harness/tht/cli/schema_cmd.py harness/tht/cli/lsh_cmd.py harness/tht/cli/preprocess_cmd.py harness/tests/test_dwh_preprocess_job.py harness/tests/test_lsh_job_resume.py git commit -m "feat(preprocess): add resumable DWH jobs" ``` ### Task 7: Add Compose job profiles, S3 follow-up adapter, and operational gates **Files:** - Modify: `compose.yaml` - Create: `harness/tht/adapters/evidence/s3.py` - Modify: `harness/pyproject.toml` - Test: `harness/tests/test_s3_evidence_source.py` - Create: `scripts/preprocess-smoke.sh` - Modify: `README.md` **Interfaces:** - Produces Compose profile `preprocess`; optional `s3` source type using endpoint URL, bucket, prefix, and secret references. - [ ] **Step 1: Add S3 contract tests against a fake endpoint and a Compose job smoke** ```python def test_s3_uri_and_fingerprint(source): item = next(source.discover()) assert item.uri == "s3://evidence/clinical/a.md" assert item.fingerprint.startswith("etag:") ``` - [ ] **Step 2: Verify failure** Run: `cd harness && .venv/bin/pytest tests/test_s3_evidence_source.py -q` Expected: FAIL. - [ ] **Step 3: Implement S3 on the established port and Compose one-shot services** ```yaml preprocess-evidence: image: thothii-core:${THOTHII_TAG:-latest} profiles: ["preprocess"] command: ["preprocess", "evidence", "--json"] volumes: ["thoth_data:/data"] ``` - [ ] **Step 4: Run all operational gates** Run: `./scripts/preprocess-smoke.sh` Expected: second run reports all documents unchanged; a modified file creates and publishes one new generation. Run: `cd harness && .venv/bin/pytest -q` Run: `cd harness && .venv/bin/ruff check tht tests/test_*evidence* tests/test_corpus* tests/test_*job*` Expected: PASS. - [ ] **Step 5: Commit** ```bash git add compose.yaml harness/tht/adapters/evidence/s3.py harness/pyproject.toml harness/tests/test_s3_evidence_source.py scripts/preprocess-smoke.sh README.md git commit -m "feat(preprocess): add deployment jobs and S3 source" ```