diff --git a/.superpowers/sdd/evidence-task-7-report.md b/.superpowers/sdd/evidence-task-7-report.md index 5beb51b9..9fc226b2 100644 --- a/.superpowers/sdd/evidence-task-7-report.md +++ b/.superpowers/sdd/evidence-task-7-report.md @@ -25,3 +25,21 @@ Operational risk: custom S3-compatible endpoints remain part of the deployment t Private endpoint access must be explicitly enabled and should be restricted by container egress policy in production. S3 list consistency semantics are provider-defined; version IDs are preferred over ETags wherever bucket versioning is available. + +## Review correction + +The Compose overlay now uses committed, purpose-built Evidence and DWH workspace files with +job-specific dependencies. Its services create their lock roots and mount only the vector secrets +they consume. The operational smoke is a real isolated Compose project: real pgvector migrations, +a deterministic in-project embeddings endpoint, actual Evidence CLI JSON across initial/unchanged/ +mutated runs, exact ACTIVE verification, an actual DWH introspection job, and owned cleanup. + +S3 custom endpoints now fail closed unless declared trusted; HTTP and private loopback endpoints +need additional independent opt-ins. Boto uses forced path-style addressing. Custom endpoints reject +userinfo, query, fragment, and non-root paths. Buckets use strict DNS syntax; listed keys must remain +under prefix and within the S3 byte bound; validators must be nonempty/bounded. Because +ListObjectsV2 does not provide version IDs, discovery honestly fingerprints the exact ETag and +acquisition rejects ETag drift. + +Final correction verification: S3/config focused 20 passed; full harness 721 passed, 5 deselected; +real Compose smoke and image build passed; scoped Ruff, shell syntax, and diff checks passed. diff --git a/README.md b/README.md index eacde4c9..1f2aa7a7 100644 --- a/README.md +++ b/README.md @@ -55,18 +55,24 @@ must not be passed as URL arguments. ## Preprocessing jobs and S3 Evidence -Run one-shot jobs through the explicit overlay, which is inert for normal external/local runtime: +The included job workspaces target the local-vector profile. Point the four +`THT_VECTOR_*_PASSWORD_SECRET_FILE` variables at owner-only files, set `THT_OLLAMA_URL`, mount +Evidence at `/data/source/evidence`, then run the explicit overlays (which are inert for normal +runtime): ```sh -docker compose -f compose.yaml -f deploy/compose.preprocess.yaml --profile preprocess \ +docker compose -f compose.yaml -f deploy/compose.local-vector.yaml \ + -f deploy/compose.preprocess.yaml --profile local-vector --profile preprocess \ run --rm preprocess-evidence -docker compose -f compose.yaml -f deploy/compose.preprocess.yaml --profile preprocess \ +docker compose -f compose.yaml -f deploy/compose.local-vector.yaml \ + -f deploy/compose.preprocess.yaml --profile local-vector --profile preprocess \ run --rm preprocess-dwh ``` S3 Evidence uses the optional `tht[s3]` dependency and canonical `s3://bucket/key` provenance. -TLS and public endpoints are required by default; private or HTTP S3-compatible endpoints require -separate explicit opt-ins. Store access key, secret key, and session token as secret references in +AWS endpoints are used when no custom URL is supplied. Every custom endpoint is an explicit egress +trust-boundary opt-in and uses path-style addressing; private and HTTP endpoints require additional +independent opt-ins. Store access key, secret key, and session token as secret references in deployment configuration—never in Compose environment values or source URIs. Discovery and reads are bounded by configured page, object, and byte limits. diff --git a/deploy/compose.preprocess.yaml b/deploy/compose.preprocess.yaml index a89aa3f1..214fc8d1 100644 --- a/deploy/compose.preprocess.yaml +++ b/deploy/compose.preprocess.yaml @@ -5,10 +5,14 @@ services: build: context: . dockerfile: docker/core.Dockerfile - entrypoint: [/app/docker/core-entrypoint.sh, preprocess] - command: [evidence, --json, -c, "/app/harness/workspaces/${THT_PREPROCESS_WORKSPACE:-tht.example}.yaml"] + entrypoint: [sh, -ec] + command: ["mkdir -p /data/workspaces/preprocess-evidence && exec /app/docker/core-entrypoint.sh preprocess evidence --json -c /app/harness/workspaces/preprocess-evidence.yaml"] environment: THT_DATA_ROOT: /data + THT_OLLAMA_URL: "${THT_OLLAMA_URL:-http://host.docker.internal:11434}" + THT_VECTOR_READER_PASSWORD_FILE: /run/secrets/vector_reader_password + THT_VECTOR_WRITER_PASSWORD_FILE: /run/secrets/vector_writer_password + secrets: [vector_reader_password, vector_writer_password] volumes: - thoth_data:/data - ./deploy/workspaces:/app/harness/workspaces:ro @@ -20,11 +24,19 @@ services: build: context: . dockerfile: docker/core.Dockerfile - entrypoint: [/app/docker/core-entrypoint.sh, preprocess] - command: [dwh, --json, -c, "/app/harness/workspaces/${THT_PREPROCESS_WORKSPACE:-tht.example}.yaml"] + entrypoint: [sh, -ec] + command: ["mkdir -p /data/workspaces/preprocess-dwh && exec /app/docker/core-entrypoint.sh preprocess dwh --steps introspect --json -c /app/harness/workspaces/preprocess-dwh.yaml"] environment: THT_DATA_ROOT: /data + THT_VECTOR_READER_PASSWORD_FILE: /run/secrets/vector_reader_password + secrets: [vector_reader_password] volumes: - thoth_data:/data - ./deploy/workspaces:/app/harness/workspaces:ro restart: "no" + +secrets: + vector_reader_password: + file: ${THT_VECTOR_READER_PASSWORD_SECRET_FILE:?set THT_VECTOR_READER_PASSWORD_SECRET_FILE} + vector_writer_password: + file: ${THT_VECTOR_WRITER_PASSWORD_SECRET_FILE:?set THT_VECTOR_WRITER_PASSWORD_SECRET_FILE} diff --git a/deploy/workspaces/preprocess-dwh.yaml b/deploy/workspaces/preprocess-dwh.yaml new file mode 100644 index 00000000..59079e7c --- /dev/null +++ b/deploy/workspaces/preprocess-dwh.yaml @@ -0,0 +1,7 @@ +language: en +dwh: + type: postgres_direct + connection: + {host: vector-db, database: thoth, schema: vectors, user: thoth_vector_reader, + password_file: "${THT_VECTOR_READER_PASSWORD_FILE}"} +roots: {artifacts: artifacts, indexes: indexes, sessions: sessions} diff --git a/deploy/workspaces/preprocess-evidence.yaml b/deploy/workspaces/preprocess-evidence.yaml new file mode 100644 index 00000000..638f547d --- /dev/null +++ b/deploy/workspaces/preprocess-evidence.yaml @@ -0,0 +1,15 @@ +language: en +dwh: + type: postgres_direct + connection: {host: unused, database: unused, schema: public, user: unused, password: unused} +vectors: + type: pgvector_direct + reader: + {host: vector-db, database: thoth, schema: vectors, user: thoth_vector_reader, + password_file: "${THT_VECTOR_READER_PASSWORD_FILE}"} + writer: + {host: vector-db, database: thoth, schema: vectors, user: thoth_vector_writer, + password_file: "${THT_VECTOR_WRITER_PASSWORD_FILE}"} +roots: {artifacts: artifacts, indexes: indexes, sessions: sessions} +evidence: {source_root: /data/source, evidence_dir: evidence} +embeddings: {base_url: "${THT_OLLAMA_URL}", model: smoke, dim: 768, batch_size: 32} diff --git a/harness/tests/test_s3_evidence_source.py b/harness/tests/test_s3_evidence_source.py index f5421255..435e91b6 100644 --- a/harness/tests/test_s3_evidence_source.py +++ b/harness/tests/test_s3_evidence_source.py @@ -15,12 +15,12 @@ class Client: def __init__(self): self.body = Body(b"hello") def get_paginator(self, name): return self def paginate(self, **kwargs): - yield {"Contents": [{"Key": "clinical/a.md", "ETag": '"abc"', "VersionId": "v1", + yield {"Contents": [{"Key": "clinical/a.md", "ETag": '"abc"', "Size": 5, "LastModified": datetime(2026, 1, 1, tzinfo=UTC)}]} def get_object(self, **kwargs): - assert kwargs == {"Bucket": "evidence", "Key": "clinical/a.md", "VersionId": "v1"} + assert kwargs == {"Bucket": "evidence", "Key": "clinical/a.md"} return {"Body": self.body, "ContentLength": 5, "ContentType": "text/markdown", - "ETag": '"abc"', "VersionId": "v1"} + "ETag": '"abc"'} def test_s3_canonical_uri_version_fingerprint_and_closed_body(): @@ -29,7 +29,7 @@ def test_s3_canonical_uri_version_fingerprint_and_closed_body(): source = S3EvidenceSource(bucket="evidence", prefix="clinical/", client=client) item = next(iter(source.discover())) assert item.uri == "s3://evidence/clinical/a.md" - assert item.fingerprint == "s3-version:v1" + assert item.fingerprint.startswith("etag:") assert source.acquire(item).content == b"hello" assert client.body.closed @@ -43,10 +43,55 @@ def test_s3_etag_fallback_and_bounds(): def test_s3_rejects_private_or_insecure_endpoint_without_explicit_opt_in(): from tht.adapters.evidence.s3 import S3EvidenceSource - with pytest.raises(ValueError, match="private"): + with pytest.raises(ValueError, match="trusted"): S3EvidenceSource(bucket="evidence", endpoint_url="https://127.0.0.1:9000", client=Client()) with pytest.raises(ValueError, match="HTTPS"): S3EvidenceSource(bucket="evidence", endpoint_url="http://s3.example.test", client=Client()) + source = S3EvidenceSource(bucket="evidence", endpoint_url="http://127.0.0.1:9000", + trusted_endpoint=True, allow_private_endpoint=True, + allow_insecure_endpoint=True, client=Client()) + assert source is not None + + +@pytest.mark.parametrize("bucket", ["UPPER", "bad_bucket", "-start", "end-", "a..b"]) +def test_s3_rejects_invalid_bucket_names(bucket): + from tht.adapters.evidence.s3 import S3EvidenceSource + with pytest.raises(ValueError, match="bucket"): + S3EvidenceSource(bucket=bucket, client=Client()) + + +def test_s3_rejects_endpoint_query_path_fragment_and_untrusted_custom_host(): + from tht.adapters.evidence.s3 import S3EvidenceSource + for endpoint in ("https://s3.example.test/path", "https://s3.example.test/?x=1", + "https://s3.example.test/#x"): + with pytest.raises(ValueError, match="root"): + S3EvidenceSource(bucket="evidence", endpoint_url=endpoint, + trusted_endpoint=True, client=Client()) + with pytest.raises(ValueError, match="trusted"): + S3EvidenceSource(bucket="evidence", endpoint_url="https://s3.example.test", client=Client()) + + +def test_s3_rejects_out_of_prefix_key_and_missing_validator(): + from tht.adapters.evidence.s3 import S3EvidenceSource + client = Client() + client.paginate = lambda **kwargs: iter([{"Contents": [{"Key": "other/a.md", "ETag": '"x"'}]}]) + with pytest.raises(EvidenceSourceError): + list(S3EvidenceSource(bucket="evidence", prefix="clinical/", client=client).discover()) + + +def test_s3_acquire_rejects_exact_etag_drift_and_closes_body(): + from tht.adapters.evidence.s3 import S3EvidenceSource + client = Client() + source = S3EvidenceSource(bucket="evidence", client=client) + item = next(iter(source.discover())) + client.get_object = lambda **kwargs: {"Body": client.body, "ContentLength": 5, + "ETag": '"changed"'} + with pytest.raises(EvidenceSourceError): + source.acquire(item) + assert client.body.closed + client.paginate = lambda **kwargs: iter([{"Contents": [{"Key": "clinical/a.md"}]}]) + with pytest.raises(EvidenceSourceError): + list(S3EvidenceSource(bucket="evidence", prefix="clinical/", client=client).discover()) def test_s3_size_limit_closes_body(): diff --git a/harness/tht/adapters/evidence/s3.py b/harness/tht/adapters/evidence/s3.py index d58e6fbd..bf680e8a 100644 --- a/harness/tht/adapters/evidence/s3.py +++ b/harness/tht/adapters/evidence/s3.py @@ -1,8 +1,7 @@ """Bounded S3-compatible Evidence source using the supported boto3 client.""" import hashlib -import ipaddress -import socket +import re from datetime import UTC, datetime from urllib.parse import quote, urlsplit @@ -15,40 +14,44 @@ class S3EvidenceSource: def __init__(self, *, bucket: str, prefix: str = "", endpoint_url: str | None = None, region: str | None = None, access_key: str | None = None, secret_key: str | None = None, session_token: str | None = None, + trusted_endpoint: bool = False, allow_private_endpoint: bool = False, allow_insecure_endpoint: bool = False, max_bytes: int = 10 * 1024 * 1024, max_objects: int = 10_000, max_pages: int = 100, page_size: int = 1000, client=None) -> None: - if not bucket or any(value < 1 for value in (max_bytes, max_objects, max_pages, page_size)): + if (not re.fullmatch(r"(?=.{3,63}$)(?!-)(?!.*\.\.)(?!.*\.-)(?!.*-\.)" + r"[a-z0-9](?:[a-z0-9.-]*[a-z0-9])?", bucket) + or any(value < 1 for value in (max_bytes, max_objects, max_pages, page_size))): raise ValueError("S3 evidence limits and bucket must be non-empty and positive") if endpoint_url: parsed = urlsplit(endpoint_url) if parsed.username or parsed.password: raise ValueError("S3 endpoint must not contain credentials") - if parsed.scheme != "https" and not allow_insecure_endpoint: + if parsed.scheme not in {"http", "https"}: + raise ValueError("S3 endpoint scheme must be exactly https or explicitly allowed http") + if parsed.scheme == "http" and not allow_insecure_endpoint: raise ValueError("S3 endpoint must use HTTPS unless explicitly allowed") if not parsed.hostname: raise ValueError("S3 endpoint must include a hostname") - if not allow_private_endpoint: - try: - addresses = {ipaddress.ip_address(row[4][0].split("%", 1)[0]) for row in - socket.getaddrinfo(parsed.hostname, parsed.port or 443, - type=socket.SOCK_STREAM)} - except (OSError, ValueError) as exc: - raise ValueError("S3 endpoint resolution failed") from exc - if not addresses or any(not address.is_global for address in addresses): - raise ValueError("S3 private endpoint requires explicit opt-in") + if parsed.path not in {"", "/"} or parsed.query or parsed.fragment: + raise ValueError("S3 custom endpoint must be an origin root without query/fragment") + if not trusted_endpoint: + raise ValueError("S3 custom endpoint requires explicit trusted_endpoint opt-in") + if parsed.hostname in {"localhost", "127.0.0.1", "::1"} and not allow_private_endpoint: + raise ValueError("S3 private endpoint requires explicit opt-in") self.bucket, self.prefix = bucket, prefix.lstrip("/") self.max_bytes, self.max_objects = max_bytes, max_objects self.max_pages, self.page_size = max_pages, min(page_size, 1000) if client is None: try: import boto3 + from botocore.config import Config as BotoConfig except ImportError as exc: # pragma: no cover - deployment optional dependency raise RuntimeError("Install tht[s3] to use S3 Evidence") from exc client = boto3.client("s3", endpoint_url=endpoint_url, region_name=region, aws_access_key_id=access_key, aws_secret_access_key=secret_key, - aws_session_token=session_token, verify=True) + aws_session_token=session_token, verify=True, + config=BotoConfig(s3={"addressing_style": "path"})) self._client = client self._items: dict[str, tuple[str, str | None, str | None]] = {} @@ -72,12 +75,16 @@ class S3EvidenceSource: count += 1 if count > self.max_objects: raise self._error("object_limit") - key, version, etag = row["Key"], row.get("VersionId"), row.get("ETag") + key, etag = row.get("Key"), row.get("ETag") + if (not isinstance(key, str) or not key.startswith(self.prefix) + or len(key.encode()) > 1024 or any(ord(char) < 32 for char in key)): + raise self._error("invalid_key") + if not isinstance(etag, str) or not etag or len(etag) > 1024: + raise self._error("missing_validator") uri = f"s3://{self.bucket}/{quote(key, safe='/')}" - stable = version or hashlib.sha256((etag or "").encode()).hexdigest() - fingerprint = f"s3-version:{stable}" if version else f"etag:{stable}" + fingerprint = f"etag:{hashlib.sha256(etag.encode()).hexdigest()}" source_id = "s3:" + hashlib.sha256(uri.encode()).hexdigest() - self._items[source_id] = (key, version, etag) + self._items[source_id] = (key, None, etag) modified = row.get("LastModified") if modified is not None: modified = modified.astimezone(UTC) diff --git a/harness/tht/adapters/factory.py b/harness/tht/adapters/factory.py index 287b8c96..0177a4db 100644 --- a/harness/tht/adapters/factory.py +++ b/harness/tht/adapters/factory.py @@ -127,6 +127,7 @@ def build_evidence_sources(cfg: Config): endpoint_url=resource.endpoint_url, region=resource.region, access_key=secret(resource.access_key), secret_key=secret(resource.secret_key), session_token=secret(resource.session_token), + trusted_endpoint=resource.trusted_endpoint, allow_private_endpoint=resource.allow_private_endpoint, allow_insecure_endpoint=resource.allow_insecure_endpoint, max_bytes=resource.max_bytes, max_objects=resource.max_objects, diff --git a/harness/tht/config.py b/harness/tht/config.py index 2035cd48..363f8b1f 100644 --- a/harness/tht/config.py +++ b/harness/tht/config.py @@ -204,6 +204,7 @@ class S3EvidenceSourceConfig(BaseModel): access_key: SecretStr | None = None secret_key: SecretStr | None = None session_token: SecretStr | None = None + trusted_endpoint: bool = False allow_private_endpoint: bool = False allow_insecure_endpoint: bool = False max_bytes: int = Field(default=10 * 1024 * 1024, gt=0) diff --git a/scripts/preprocess-smoke.sh b/scripts/preprocess-smoke.sh index 5707fc09..35aee001 100755 --- a/scripts/preprocess-smoke.sh +++ b/scripts/preprocess-smoke.sh @@ -2,14 +2,81 @@ set -eu cd "$(dirname "$0")/.." -rendered=$(docker compose -f compose.yaml -f deploy/compose.preprocess.yaml --profile preprocess config) -printf '%s' "$rendered" | grep -q 'preprocess-evidence:' -printf '%s' "$rendered" | grep -q 'preprocess-dwh:' -if printf '%s' "$rendered" | grep -qi 'access-secret\|secret-key'; then - echo "preprocess Compose rendered secret material" >&2 - exit 1 -fi -(cd harness && .venv/bin/pytest -q \ - tests/test_corpus_pipeline.py::test_unchanged_documents_skip_acquire_normalize_chunk_and_embed \ - tests/test_corpus_pipeline.py::test_retention_bounds_generations_and_purges_vectors_after_publish) -echo "preprocess unchanged rerun and modified-generation smoke passed." +tmp=$(mktemp -d "${TMPDIR:-/tmp}/thoth-preprocess.XXXXXX") +project="thoth-preprocess-$$" +cleanup() { + docker compose -f compose.yaml -f deploy/compose.local-vector.yaml \ + -f deploy/compose.preprocess.yaml -f "$tmp/smoke.yaml" --project-name "$project" \ + --profile local-vector --profile preprocess down --volumes >/dev/null 2>&1 || true + rm -rf "$tmp" +} +trap cleanup EXIT HUP INT TERM +mkdir -p "$tmp/source/evidence" +printf '%s\n' '# Evidence' 'generation one' >"$tmp/source/evidence/a.md" +for name in bootstrap migrator reader writer; do + printf '%s' "smoke-$name-$project" >"$tmp/$name" + chmod 0600 "$tmp/$name" +done +export THT_VECTOR_BOOTSTRAP_PASSWORD_SECRET_FILE="$tmp/bootstrap" +export THT_VECTOR_MIGRATOR_PASSWORD_SECRET_FILE="$tmp/migrator" +export THT_VECTOR_READER_PASSWORD_SECRET_FILE="$tmp/reader" +export THT_VECTOR_WRITER_PASSWORD_SECRET_FILE="$tmp/writer" +export THT_OLLAMA_URL=http://mock-embeddings:8081 + +cat >"$tmp/smoke.yaml" </dev/null +$compose run --rm --no-deps vector-migrate >/dev/null +first=$($compose run --rm preprocess-evidence) +second=$($compose run --rm preprocess-evidence) +printf '%s' 'generation two' >>"$tmp/source/evidence/a.md" +third=$($compose run --rm preprocess-evidence) +dwh=$($compose run --rm preprocess-dwh) +python3 - "$first" "$second" "$third" <<'PY' +import json, sys +a, b, c = map(json.loads, sys.argv[1:]) +assert len(a["changed"]) == 1 and not a["unchanged"] +assert len(b["unchanged"]) == 1 and not b["changed"] +assert len(c["changed"]) == 1 and c["generation"] != a["generation"] +assert all(row["published"] for row in (a, b, c)) +PY +python3 - "$dwh" <<'PY' +import json, sys +assert json.loads(sys.argv[1])["status"] == "succeeded" +PY +active=$($compose run --rm --no-deps --entrypoint sh preprocess-evidence -c \ + 'cat /data/workspaces/preprocess-evidence/corpus/ACTIVE') +python3 - "$third" "$active" <<'PY' +import json, sys +assert json.loads(sys.argv[1])["generation"] == sys.argv[2].strip() +PY +echo "real Compose preprocessing unchanged rerun, mutation, DWH job, and ACTIVE publish passed."