fix(preprocess): verify canonical generations and cleanup

This commit is contained in:
2026-07-12 06:18:18 +02:00
parent 4a29086fe4
commit 7d41c4cefc
6 changed files with 86 additions and 18 deletions
@@ -53,3 +53,12 @@ and absent validators fail closed. Acquisition accepts only the exact stored `So
compares the response ETag with the stored discovery validator. The real smoke snapshots generation compares the response ETag with the stored discovery validator. The real smoke snapshots generation
directory counts after every run and has an injected-failure cleanup mode; cleanup fails if Compose directory counts after every run and has an injected-failure cleanup mode; cleanup fails if Compose
down fails or any owned container, volume, or network remains. down fails or any owned container, volume, or network remains.
The canonical smoke correction counts only root-level `corpus/gen-<32 hex>` directories. It exposed
that the durable job path still published an empty unchanged generation; the pipeline now returns
the existing ACTIVE generation without staging a directory when compatibility and all source
fingerprints are unchanged. The smoke therefore proves directory deltas `+1`, `+0`, `+1`.
Failure injection runs a real exit-97 command after resources exist and reaches the EXIT trap.
Cleanup aggregates Compose-down, residual container/volume/network, and temp-directory failures
while preserving the original failure status. S3 prefixes are validated before any client request
for leading slash, UTF-8 byte length, controls, and DEL.
+14
View File
@@ -431,6 +431,20 @@ def test_unchanged_documents_skip_acquire_normalize_chunk_and_embed(tmp_path):
assert second_embedder.calls == [] assert second_embedder.calls == []
def test_unchanged_job_reuses_active_generation_without_new_directory(tmp_path):
source = Source([(item("one", "a"), "hello")])
candidate = pipeline(tmp_path, source)
args = dict(workspace_id="demo", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64)
first = candidate.run_as_job(**args)
count = len(candidate.store.list_generations())
second = candidate.run_as_job(**args)
assert second.generation == first.generation
assert second.published is False
assert len(candidate.store.list_generations()) == count
def test_removed_documents_are_marked_and_absent_from_new_manifest(tmp_path): def test_removed_documents_are_marked_and_absent_from_new_manifest(tmp_path):
one, two = item("one", "a"), item("two", "b") one, two = item("one", "a"), item("two", "b")
pipeline(tmp_path, Source([(one, "one"), (two, "two")])).run() pipeline(tmp_path, Source([(one, "one"), (two, "two")])).run()
+9
View File
@@ -104,6 +104,15 @@ def test_s3_rejects_leading_slash_prefix_empty_and_control_keys():
list(S3EvidenceSource(bucket="evidence", prefix="clinical/", client=client).discover()) list(S3EvidenceSource(bucket="evidence", prefix="clinical/", client=client).discover())
@pytest.mark.parametrize("prefix", ["/bad", "x" * 1025, "bad\x00prefix", "bad\x7fprefix"])
def test_s3_rejects_invalid_prefix_before_client_request(prefix):
from tht.adapters.evidence.s3 import S3EvidenceSource
client = Client()
with pytest.raises(ValueError, match="prefix"):
S3EvidenceSource(bucket="evidence", prefix=prefix, client=client)
assert client.list_calls == 0
def test_s3_hard_page_limit_never_requests_page_max_plus_one(): def test_s3_hard_page_limit_never_requests_page_max_plus_one():
from tht.adapters.evidence.s3 import S3EvidenceSource from tht.adapters.evidence.s3 import S3EvidenceSource
client = Client() client = Client()
+3 -2
View File
@@ -29,8 +29,9 @@ class S3EvidenceSource:
if (not bucket_valid or bucket_is_ip if (not bucket_valid or bucket_is_ip
or any(value < 1 for value in (max_bytes, max_objects, max_pages, page_size))): 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") raise ValueError("S3 evidence limits and bucket must be non-empty and positive")
if prefix.startswith("/"): if (prefix.startswith("/") or len(prefix.encode()) > 1024
raise ValueError("S3 prefix must not start with a slash") or any(ord(char) < 32 or ord(char) == 127 for char in prefix)):
raise ValueError("S3 prefix is invalid")
if endpoint_url: if endpoint_url:
parsed = urlsplit(endpoint_url) parsed = urlsplit(endpoint_url)
if parsed.username or parsed.password: if parsed.username or parsed.password:
+10
View File
@@ -220,6 +220,16 @@ class CorpusPipeline:
"dimensions": self.embedding_dimensions, "dimensions": self.embedding_dimensions,
"chunk_policy": asdict(self.chunk_policy), "chunk_policy": asdict(self.chunk_policy),
}) })
previous = self.store.active_manifest()
if (not dry_run and resume_run_id is None and previous is not None
and previous.metadata.get("compatibility_fingerprint") == compatibility
and previous.metadata.get("fingerprints") == {
item.source_id: item.fingerprint for _, item in discovered
}):
return PipelineResult(
"succeeded", previous.manifest_id, False, (),
tuple(sorted(item.source_id for _, item in discovered)), (), previous,
)
spec = JobSpec( spec = JobSpec(
workspace_id=workspace_id, workspace_id=workspace_id,
job_type="evidence", job_type="evidence",
+41 -16
View File
@@ -2,20 +2,50 @@
set -eu set -eu
cd "$(dirname "$0")/.." cd "$(dirname "$0")/.."
tmp=$(mktemp -d "${TMPDIR:-/tmp}/thoth-preprocess.XXXXXX") if [ "${1:-}" = "--cleanup-failure" ] && [ -z "${PREPROCESS_SMOKE_CHILD:-}" ]; then
project="thoth-preprocess-$$" child_project="thoth-preprocess-failure-$$"
child_tmp=$(mktemp -d "${TMPDIR:-/tmp}/thoth-preprocess-failure.XXXXXX")
set +e
PREPROCESS_SMOKE_CHILD=1 PREPROCESS_SMOKE_INJECT_FAILURE=1 \
PREPROCESS_SMOKE_PROJECT="$child_project" PREPROCESS_SMOKE_TMP="$child_tmp" "$0"
child_status=$?
set -e
test "$child_status" -eq 97
test ! -e "$child_tmp"
test -z "$(docker ps -aq --filter "label=com.docker.compose.project=$child_project")"
test -z "$(docker volume ls -q --filter "label=com.docker.compose.project=$child_project")"
test -z "$(docker network ls -q --filter "label=com.docker.compose.project=$child_project")"
echo "injected preprocessing failure preserved status and cleaned every owned resource."
exit 0
fi
tmp=${PREPROCESS_SMOKE_TMP:-$(mktemp -d "${TMPDIR:-/tmp}/thoth-preprocess.XXXXXX")}
project=${PREPROCESS_SMOKE_PROJECT:-thoth-preprocess-$$}
compose="" compose=""
cleanup() { cleanup() {
trap - EXIT HUP INT TERM original_status=$?
cleanup_failed=0
set +e
if [ -n "$compose" ]; then if [ -n "$compose" ]; then
$compose down --volumes >/dev/null $compose down --volumes >/dev/null
test "$?" -eq 0 || cleanup_failed=1
fi fi
test -z "$(docker ps -aq --filter "label=com.docker.compose.project=$project")" test -z "$(docker ps -aq --filter "label=com.docker.compose.project=$project")" || cleanup_failed=1
test -z "$(docker volume ls -q --filter "label=com.docker.compose.project=$project")" test -z "$(docker volume ls -q --filter "label=com.docker.compose.project=$project")" || cleanup_failed=1
test -z "$(docker network ls -q --filter "label=com.docker.compose.project=$project")" test -z "$(docker network ls -q --filter "label=com.docker.compose.project=$project")" || cleanup_failed=1
rm -rf "$tmp" rm -rf "$tmp"
test ! -e "$tmp" || cleanup_failed=1
if [ "$original_status" -ne 0 ]; then
exit "$original_status"
fi
if [ "$cleanup_failed" -ne 0 ]; then
exit 1
fi
} }
trap cleanup EXIT HUP INT TERM trap cleanup EXIT
trap 'exit 129' HUP
trap 'exit 130' INT
trap 'exit 143' TERM
mkdir -p "$tmp/source/evidence" mkdir -p "$tmp/source/evidence"
printf '%s\n' '# Evidence' 'generation one' >"$tmp/source/evidence/a.md" printf '%s\n' '# Evidence' 'generation one' >"$tmp/source/evidence/a.md"
for name in bootstrap migrator reader writer; do for name in bootstrap migrator reader writer; do
@@ -61,15 +91,12 @@ $compose build preprocess-evidence
$compose up -d vector-db mock-embeddings $compose up -d vector-db mock-embeddings
$compose run --rm vector-reconcile >/dev/null $compose run --rm vector-reconcile >/dev/null
$compose run --rm --no-deps vector-migrate >/dev/null $compose run --rm --no-deps vector-migrate >/dev/null
if [ "${1:-}" = "--cleanup-failure" ]; then if [ "${PREPROCESS_SMOKE_INJECT_FAILURE:-0}" = "1" ]; then
cleanup sh -c 'exit 97'
test ! -e "$tmp"
echo "injected preprocessing failure cleanup passed."
exit 0
fi fi
generation_count() { generation_count() {
$compose run --rm --no-deps --entrypoint sh preprocess-evidence -c \ $compose run --rm --no-deps --entrypoint /opt/venv/bin/python preprocess-evidence -c \
'find /data/workspaces/preprocess-evidence/corpus/generations -mindepth 1 -maxdepth 1 -type d 2>/dev/null | wc -l' 'import pathlib,re; root=pathlib.Path("/data/workspaces/preprocess-evidence/corpus"); print(sum(1 for p in root.iterdir() if p.is_dir() and re.fullmatch(r"gen-[0-9a-f]{32}",p.name)) if root.exists() else 0)'
} }
before=$(generation_count) before=$(generation_count)
first=$($compose run --rm preprocess-evidence) first=$($compose run --rm preprocess-evidence)
@@ -103,5 +130,3 @@ import json, sys
assert json.loads(sys.argv[1])["generation"] == sys.argv[2].strip() assert json.loads(sys.argv[1])["generation"] == sys.argv[2].strip()
PY PY
echo "real Compose preprocessing unchanged rerun, mutation, DWH job, and ACTIVE publish passed." echo "real Compose preprocessing unchanged rerun, mutation, DWH job, and ACTIVE publish passed."
cleanup
test ! -e "$tmp"