diff --git a/docs/architecture/04-server-api-endpoints.md b/docs/architecture/04-server-api-endpoints.md index 2618a7317..c8b17880c 100644 --- a/docs/architecture/04-server-api-endpoints.md +++ b/docs/architecture/04-server-api-endpoints.md @@ -85,7 +85,7 @@ Every `/api/files/*` operation resolves its requested path and refuses the opera | GET | `/api/tasks/{task_id}` | `gateway.tasks.api_task_get` | | GET | `/api/tasks/{task_id}/events` | `gateway.tasks.api_task_events` (legacy integer rank) | | POST | `/api/tasks/{task_id}/events` | `gateway.tasks.api_task_events` (read-only v2 cursor) | -| GET | `/api/tasks/{task_id}/artifacts/{name}` | `gateway.tasks.api_task_artifact` (the task's own stores via `task_archive`: a bare name is a top-level file, `?relpath=` a nested one, `?archive=` a directory ZIP; the detail's `artifact_archives` says what each ZIP holds. A row that records a digest is served only when its bytes still match it — a changed mutable file is 409 `artifact_identity_changed` naming the recorded digest, a failed capture 404 `artifact_unverified` — and the response says `x-ouroboros-artifact-identity: verified` or `unmeasured`; a ZIP member follows the same rule) | +| GET | `/api/tasks/{task_id}/artifacts/{name}` | `gateway.tasks.api_task_artifact` (the task's own stores via `task_archive`: a bare name is a top-level file, `?relpath=` a nested one, `?archive=` a directory ZIP; the detail's `artifact_archives` says what each ZIP holds. A row that records a digest is served only when its bytes still match it — a changed mutable file is 409 `artifact_identity_changed` naming the recorded digest, a failed capture 404 `artifact_unverified` — and the response says `x-ouroboros-artifact-identity: verified` or `unmeasured`; a ZIP member follows the same rule. Windows ordinary file/chat-media/ZIP downloads return HTTP 503 pending #1297; bound `?source=` review downloads remain available) | | POST | `/api/tasks/{task_id}/cancel` | `gateway.tasks.api_task_cancel` | | POST | `/api/tasks/{task_id}/hurry` | `gateway.tasks.api_task_hurry` | | POST | `/api/tasks/{task_id}/resume` | `gateway.tasks.api_task_resume` | diff --git a/ouroboros/skill_loader.py b/ouroboros/skill_loader.py index 9f98fda90..46ba8e190 100644 --- a/ouroboros/skill_loader.py +++ b/ouroboros/skill_loader.py @@ -11,6 +11,7 @@ from __future__ import annotations import hashlib import logging import pathlib +import stat from dataclasses import asdict, dataclass, field from typing import Any, Callable, Dict, Iterable, List, Optional @@ -356,10 +357,11 @@ def _iter_payload_files( resolved = (skill_dir / rel).resolve() try: resolved.relative_to(resolved_root) - except ValueError: + # is_file() hides ELOOP; an unreadable declared entry cannot be omitted from its hash. + if stat.S_ISREG(resolved.stat().st_mode): + _add(resolved) + except (ValueError, FileNotFoundError, NotADirectoryError): return - if resolved.is_file(): - _add(resolved) # Broad walk: everything runtime-reachable, minus metadata/cache names. # Every candidate is resolved back under skill_dir so symlinks cannot leak diff --git a/ouroboros/tools/plan_evidence.py b/ouroboros/tools/plan_evidence.py index 225f56b6e..55422714f 100644 --- a/ouroboros/tools/plan_evidence.py +++ b/ouroboros/tools/plan_evidence.py @@ -14,6 +14,7 @@ from __future__ import annotations import codecs import ast +import errno from hashlib import sha256 import pathlib import json @@ -214,13 +215,15 @@ def _read_evidence( def _path_kind(path: pathlib.Path) -> str: - """``missing`` · ``directory`` · ``file`` (regular) · ``unreadable`` (stat failure or a + """``missing`` · ``directory`` · ``file`` · ``symlink_loop`` · ``unreadable`` (stat failure or a non-regular node: fifo/device/socket would block or never end a read).""" try: mode = path.stat().st_mode except FileNotFoundError: return "missing" - except (OSError, ValueError, RuntimeError): + except OSError as exc: + return "symlink_loop" if exc.errno == errno.ELOOP else "unreadable" + except (ValueError, RuntimeError): return "unreadable" if stat.S_ISDIR(mode): return "directory" diff --git a/tests/conftest.py b/tests/conftest.py index 34648369d..429b76471 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -586,13 +586,16 @@ def _rebind_runtime_roots_between_tests(): @pytest.fixture(autouse=True) def _reset_custody_memo_between_tests(): - """The custody row memo is process-local and keyed by events-log path; a test - that rewrites its log in place (``write_text``) or reuses a path must never - inherit another test's consumed prefix (``delegate_custody_memo``).""" + """Isolate both custody caches: the row memo is keyed by events-log path, + while active custody is keyed only by run ID. Tests reuse both identities; + neither a consumed prefix nor a previous run's first-wins binding may leak.""" + from ouroboros import delegate_custody from ouroboros.delegate_custody_memo import reset_custody_memo + delegate_custody._CUSTODY.clear() reset_custody_memo() yield + delegate_custody._CUSTODY.clear() reset_custody_memo() diff --git a/tests/test_author_ui_kit_browser.py b/tests/test_author_ui_kit_browser.py index 642cb10f3..b99e10ad9 100644 --- a/tests/test_author_ui_kit_browser.py +++ b/tests/test_author_ui_kit_browser.py @@ -196,7 +196,7 @@ def test_author_kit_authenticated_mount_and_lifetime(author_kit_server, tmp_path assert route.evaluate("window.kitCspViolations") == [] assert module.evaluate("document.documentElement.dataset.theme") == "dark" page.evaluate("() => window.ouroTheme.set('light')") - module.wait_for_function("document.documentElement.dataset.theme === 'light'") + module.locator('html[data-theme="light"]').wait_for(state="attached") assert page.evaluate( "node => document.querySelector('[data-widget-key=\\\"module-old\\\"] iframe') === node", module_node, @@ -205,7 +205,7 @@ def test_author_kit_authenticated_mount_and_lifetime(author_kit_server, tmp_path "getComputedStyle(document.querySelector('.ui-control')).backgroundColor" ) == "rgb(245, 246, 248)" page.evaluate("() => window.ouroTheme.set('dark')") - module.wait_for_function("document.documentElement.dataset.theme === 'dark'") + module.locator('html[data-theme="dark"]').wait_for(state="attached") for frame in (module, route): frame.get_by_role("button", name="Preview", exact=True).wait_for() assert frame.get_by_label("Title", exact=True).input_value() == "My notes" @@ -322,7 +322,7 @@ def test_author_kit_authenticated_mount_and_lifetime(author_kit_server, tmp_path }""") page.wait_for_function("window.handlers.size === 1") page.evaluate("() => handlers.forEach(handler => handler({type:'ext:export_widget:tick', data:{value:'delivered'}}))") - module.wait_for_function("document.getElementById('root').dataset.event === 'delivered'") + module.locator('#root[data-event="delivered"]').wait_for(state="attached") with page.expect_download() as download: module.get_by_role("button", name="Export example", exact=True).click() assert Path(download.value.path()).read_text(encoding="utf-8") == "author kit export" @@ -335,7 +335,7 @@ def test_author_kit_authenticated_mount_and_lifetime(author_kit_server, tmp_path evidence = Path(os.environ.get("OUROBOROS_UI_EVIDENCE_OUT", str(tmp_path / "evidence"))) evidence.mkdir(parents=True, exist_ok=True) page.evaluate("() => window.ouroTheme.set('light')") - module.wait_for_function("document.documentElement.dataset.theme === 'light'") + module.locator('html[data-theme="light"]').wait_for(state="attached") module.locator('body').screenshot( path=str(evidence / f"author-kit-{browser_name}-module-light.png") ) diff --git a/tests/test_child_drive_settlement.py b/tests/test_child_drive_settlement.py index 9d537e740..30d7c40eb 100644 --- a/tests/test_child_drive_settlement.py +++ b/tests/test_child_drive_settlement.py @@ -117,7 +117,7 @@ def test_a_recorded_capture_whose_source_is_gone_stays_held_by_its_verified_cano @pytest.mark.parametrize("material", ["child_result", "registration", "file", "directory", "mailbox", "receipts"]) -def test_unreadable_material_retains_the_drive(tmp_path, material): +def test_unreadable_material_retains_the_drive(tmp_path, monkeypatch, material): data, drive, record = _cancelled_with_capture(tmp_path) store = artifacts.task_artifact_dir_path(drive, TASK) if material == "child_result": @@ -129,12 +129,26 @@ def test_unreadable_material_retains_the_drive(tmp_path, material): elif material == "file": extra = store / "notes.txt" extra.write_text("mutable note", encoding="utf-8") - extra.chmod(0) + original_open = Path.open + + def unreadable_open(path, *args, **kwargs): + if path == extra: + raise PermissionError("file read denied") + return original_open(path, *args, **kwargs) + + monkeypatch.setattr(Path, "open", unreadable_open) expected = "child_artifact_store_unreadable" elif material == "directory": (store / "tree").mkdir() (store / "tree" / "leaf.txt").write_text("leaf", encoding="utf-8") - (store / "tree").chmod(0) + original_scandir = os.scandir + + def unreadable_scandir(path): + if Path(path) == store / "tree": + raise PermissionError("directory read denied") + return original_scandir(path) + + monkeypatch.setattr(os, "scandir", unreadable_scandir) expected = "child_artifact_store_unreadable" elif material == "mailbox": mailbox = owner_mailbox._mailbox_path(drive, TASK) @@ -144,16 +158,22 @@ def test_unreadable_material_retains_the_drive(tmp_path, material): else: verification_receipts_path(drive, TASK, create=True).write_text("{torn\n", encoding="utf-8") expected = "verification_receipts_uncustodied" - try: - if os.name != "nt" and material in {"file", "directory"} and os.geteuid() == 0: - pytest.skip("root reads unreadable files") - outcome = _settle(data, drive) - finally: - for path in (store / "notes.txt", store / "tree"): - if path.exists(): - path.chmod(0o755) + # chmod(0) does not deny Windows reads and is bypassed by POSIX root. Inject + # the actual read failure so retention and recovery are tested on every host. + outcome = _settle(data, drive) assert outcome["status"] == "retained" and outcome["reason"] == expected, outcome assert drive.is_dir() and not _canonical(data, record["name"]).exists() + if material in {"file", "directory"}: + if material == "file": + monkeypatch.setattr(Path, "open", original_open) + else: + monkeypatch.setattr(os, "scandir", original_scandir) + assert _settle(data, drive)["status"] == "removed" + assert not drive.exists() + extra_name = "notes.txt" if material == "file" else "tree/leaf.txt" + assert _canonical(data, extra_name).read_text(encoding="utf-8") == ( + "mutable note" if material == "file" else "leaf" + ) def test_a_write_that_does_not_land_or_read_back_deletes_and_seals_nothing(tmp_path, monkeypatch): diff --git a/tests/test_cybergym_dispatch.py b/tests/test_cybergym_dispatch.py index b54c49a32..19e0fede9 100644 --- a/tests/test_cybergym_dispatch.py +++ b/tests/test_cybergym_dispatch.py @@ -443,11 +443,18 @@ def _transport_row(task_id: str) -> dict: } -def test_breaker_pauses_probes_and_resumes_instead_of_abandoning(tmp_path): +def test_breaker_pauses_probes_and_resumes_instead_of_abandoning(tmp_path, monkeypatch): """full1507: three transport rows from a ~100 s stall, not a dead isolate.""" + from concurrent.futures import ALL_COMPLETED, wait + from devtools.benchmarks.cybergym import cybergym_dispatch from devtools.benchmarks.cybergym.cybergym_dispatch import run_dispatched + # This scenario delivers a whole failed wave before any later success can + # reset the consecutive-failure streak. Worker scheduling is not its oracle. + monkeypatch.setattr(cybergym_dispatch, "wait", lambda futures, **kwargs: wait( + futures, **{**kwargs, "return_when": ALL_COMPLETED})) + clock = _PausingClock() probes: list[float] = [] # Gateway is "stalled" for the first two probes and answers on the third. @@ -457,17 +464,16 @@ def test_breaker_pauses_probes_and_resumes_instead_of_abandoning(tmp_path): probes.append(clock.now) return next(probe_answers) - outcomes = iter(["transport", "transport", "transport", "ok", "ok", "ok"]) + tasks = [_Task(f"arvo:{index}") for index in range(1, 7)] + failed_task_ids = {task.task_id for task in tasks[:3]} def run_one(task): - if next(outcomes) == "ok": + if task.task_id not in failed_task_ids: return {"task_id": task.task_id, "status": "completed"} return _transport_row(task.task_id) events: list[dict] = [] - tasks = [_Task(f"arvo:{index}") for index in range(1, 7)] for workers in (1, 3): - outcomes = iter(["transport", "transport", "transport", "ok", "ok", "ok"]) probe_answers = iter([False, False, True]) probes.clear() events.clear() @@ -486,6 +492,7 @@ def test_breaker_pauses_probes_and_resumes_instead_of_abandoning(tmp_path): clock=clock.monotonic, ) + assert [row["task_id"] for row in rows] == [task.task_id for task in tasks] assert [row["status"] for row in rows] == ["infra_failed"] * 3 + ["completed"] * 3 # Backoff schedule: first probe after 30 s, then +60 s, then +120 s. assert probes == [30.0, 90.0, 210.0] diff --git a/tests/test_evolution_commit_receipt.py b/tests/test_evolution_commit_receipt.py index 0ea50aa3e..b412eb3b5 100644 --- a/tests/test_evolution_commit_receipt.py +++ b/tests/test_evolution_commit_receipt.py @@ -226,6 +226,7 @@ def test_rescue_link_uses_shared_campaign_cas_and_preserves_commit_receipt( def test_commit_receipt_uses_campaign_sidecar_before_rescue(tmp_path, monkeypatch): + from ouroboros import platform_layer from ouroboros.platform_layer import ( acquire_exclusive_file_lock, release_exclusive_file_lock, @@ -239,8 +240,20 @@ def test_commit_receipt_uses_campaign_sidecar_before_rescue(tmp_path, monkeypatc lock_fd = acquire_exclusive_file_lock(lock_path, timeout_sec=1.0) assert lock_fd is not None done = threading.Event() + acquiring = threading.Event() + released = threading.Event() result = {} + def _acquire_after_release(path, **kwargs): + if path == lock_path: + acquiring.set() + released.wait() + return acquire_exclusive_file_lock(path, **kwargs) + + # Keep real contention, but start the acquisition timeout after the fixture + # releases its lock, as in the rescue interleaving test above. + monkeypatch.setattr(platform_layer, "acquire_exclusive_file_lock", _acquire_after_release) + def _record() -> None: result.update(evolution_lifecycle.record_evolution_commit( campaign["id"], tx["transaction_id"], tx["task_id"], "4" * 40, @@ -250,11 +263,15 @@ def test_commit_receipt_uses_campaign_sidecar_before_rescue(tmp_path, monkeypatc thread = threading.Thread(target=_record, daemon=True) thread.start() try: + assert acquire_exclusive_file_lock(lock_path, timeout_sec=0.001) is None + assert acquiring.wait(2.0) is True assert done.wait(0.1) is False finally: release_exclusive_file_lock(lock_path, lock_fd) + released.set() + thread.join(timeout=2.0) assert done.wait(2.0) is True - thread.join(timeout=1.0) + assert not thread.is_alive() assert result["ok"] is True assert evolution_lifecycle._read_evolution_campaign()["active_transaction"][ "commit_receipt" diff --git a/tests/test_headless_task_artifacts.py b/tests/test_headless_task_artifacts.py index 64906ab5f..965f90cf4 100644 --- a/tests/test_headless_task_artifacts.py +++ b/tests/test_headless_task_artifacts.py @@ -17,6 +17,7 @@ from starlette.applications import Starlette from starlette.routing import Route from starlette.testclient import TestClient +from ouroboros.gateway import task_archive from ouroboros.gateway.tasks import ( api_task_artifact, ) @@ -39,6 +40,16 @@ from tests._headless_cli_shared import ( # noqa: F401 (autouse fixture applies ) +def _assert_file_response(response, content: bytes) -> None: + # Owner-approved #1297: platforms without confined opens return a typed 503. + if task_archive.CONFINED: + assert response.status_code == 200 + assert response.content == content + else: + assert response.status_code == 503 + assert response.json()["reason_code"] == "artifact_unavailable" + + def test_copy_child_result_cannot_overwrite_finalized_accounting(tmp_path): """F2: once the root's terminal checkpoint has finalized accounting (task_cost_finalized rides the same write as post_task_synthesis), a late @@ -386,7 +397,7 @@ def test_task_artifact_endpoint_serves_only_declared_artifacts(tmp_path): app.state.drive_root = data client = TestClient(app) - assert client.get("/api/tasks/task-artifact/artifacts/workspace.patch").text.startswith("diff --git") + _assert_file_response(client.get("/api/tasks/task-artifact/artifacts/workspace.patch"), b"diff --git a/a b/a\n") assert client.get("/api/tasks/task-artifact/artifacts/missing.patch").status_code == 404 assert client.get("/api/tasks/task-artifact/artifacts/bad%5Cname").status_code == 400 @@ -415,8 +426,7 @@ def test_task_artifact_endpoint_serves_manifest_artifact_after_status_repair(tmp response = TestClient(app).get("/api/tasks/orphaned/artifacts/report.html") - assert response.status_code == 200 - assert response.text == "

ok

" + _assert_file_response(response, b"

ok

") def test_task_artifact_endpoint_serves_child_drive_artifact_read_only_after_status_repair(tmp_path): @@ -457,9 +467,8 @@ def test_task_artifact_endpoint_serves_child_drive_artifact_read_only_after_stat response = TestClient(app).get("/api/tasks/childart/artifacts/report.html") - # Two-root read: the task's own child store serves it; nothing is copied or created. - assert response.status_code == 200 - assert response.text == "

child

" + # The own child store is read when supported; neither outcome copies or creates it. + _assert_file_response(response, b"

child

") assert not task_artifacts_dir(data, "childart", create=False).exists() @@ -664,8 +673,7 @@ def test_task_artifact_endpoint_serves_exact_chat_media_without_task_result(tmp_ client = TestClient(app) response = client.get(f"/api/tasks/ephemeral1/artifacts/{stored['name']}") - assert response.status_code == 200 - assert response.content == b"photo-bytes" + _assert_file_response(response, b"photo-bytes") assert collect_task_artifact_records(data, "ephemeral1") == [] assert client.get("/api/tasks/ephemeral1/artifacts/chat-media-bad.png").status_code == 404 diff --git a/tests/test_large_task_artifacts.py b/tests/test_large_task_artifacts.py index 3faac8bf9..57503c868 100644 --- a/tests/test_large_task_artifacts.py +++ b/tests/test_large_task_artifacts.py @@ -629,7 +629,7 @@ def test_file_verification_does_not_run_on_the_asgi_loop(tmp_path, monkeypatch, import uvicorn from starlette.applications import Starlette from starlette.routing import Route - from ouroboros.gateway import tasks + from ouroboros.gateway import task_archive, tasks from ouroboros.task_results import write_task_result api_task_artifact = tasks.api_task_artifact @@ -667,14 +667,24 @@ def test_file_verification_does_not_run_on_the_asgi_loop(tmp_path, monkeypatch, async def request(): async with httpx.AsyncClient(timeout=10) as client: reply = await client.request(method, f'http://127.0.0.1:{sock.getsockname()[1]}/api/tasks/download/artifacts/{record["name"]}') - assert reply.status_code == 200, reply.text - assert reply.content == (source.read_bytes() if method == 'GET' else b'') + if task_archive.CONFINED: + assert reply.status_code == 200, reply.text + assert reply.content == (source.read_bytes() if method == 'GET' else b'') + else: + assert reply.status_code == 503, reply.text + if method == 'GET': + assert reply.json()['reason_code'] == 'artifact_unavailable' + else: + assert reply.content == b'' asyncio.run(request()) assert bool(materialization_calls) is (not registered) assert all(identity != thread.ident for identity in materialization_calls) - if registered: - assert len(calls) == 1, "registered downloads verify once per request" - assert calls and all(identity != thread.ident for identity in calls), 'whole-file hashing ran synchronously on the ASGI event loop' + if task_archive.CONFINED: + if registered: + assert len(calls) == 1, "registered downloads verify once per request" + assert calls and all(identity != thread.ident for identity in calls), 'whole-file hashing ran synchronously on the ASGI event loop' + else: + assert calls == [], "unsupported platforms refuse before reading file bytes (#1297)" finally: server.should_exit = True thread.join(10) diff --git a/tests/test_net_transport_extra_ca.py b/tests/test_net_transport_extra_ca.py index 6f4451933..cdb401ffb 100644 --- a/tests/test_net_transport_extra_ca.py +++ b/tests/test_net_transport_extra_ca.py @@ -45,11 +45,18 @@ def _throwaway_ca(tmp_path: pathlib.Path): now = _dt.datetime.now(_dt.timezone.utc) ca_key = rsa.generate_private_key(public_exponent=65537, key_size=2048) ca_name = x509.Name([x509.NameAttribute(NameOID.COMMON_NAME, "Ouroboros throwaway test CA")]) + # Python 3.13 enables strict X.509 verification: CA SKI/key usage and leaf AKI are required. ca = ( x509.CertificateBuilder().subject_name(ca_name).issuer_name(ca_name) .public_key(ca_key.public_key()).serial_number(x509.random_serial_number()) .not_valid_before(now - _dt.timedelta(days=1)).not_valid_after(now + _dt.timedelta(days=2)) .add_extension(x509.BasicConstraints(ca=True, path_length=None), critical=True) + .add_extension(x509.SubjectKeyIdentifier.from_public_key(ca_key.public_key()), critical=False) + .add_extension(x509.KeyUsage( + digital_signature=False, content_commitment=False, key_encipherment=False, + data_encipherment=False, key_agreement=False, key_cert_sign=True, crl_sign=True, + encipher_only=False, decipher_only=False, + ), critical=True) .sign(ca_key, hashes.SHA256()) ) leaf_key = rsa.generate_private_key(public_exponent=65537, key_size=2048) @@ -59,6 +66,7 @@ def _throwaway_ca(tmp_path: pathlib.Path): .issuer_name(ca_name).public_key(leaf_key.public_key()).serial_number(x509.random_serial_number()) .not_valid_before(now - _dt.timedelta(days=1)).not_valid_after(now + _dt.timedelta(days=2)) .add_extension(x509.SubjectAlternativeName([x509.IPAddress(ipaddress.ip_address("127.0.0.1"))]), critical=False) + .add_extension(x509.AuthorityKeyIdentifier.from_issuer_public_key(ca_key.public_key()), critical=False) .sign(ca_key, hashes.SHA256()) ) tmp_path.mkdir(parents=True, exist_ok=True) diff --git a/tests/test_review_source_handles.py b/tests/test_review_source_handles.py index 966d31a7e..14f37869a 100644 --- a/tests/test_review_source_handles.py +++ b/tests/test_review_source_handles.py @@ -9,6 +9,7 @@ from starlette.routing import Route from starlette.testclient import TestClient from ouroboros import artifacts, review_projection +from ouroboros.gateway import task_archive from ouroboros.gateway.tasks import api_task_artifact from ouroboros.headless import copy_child_task_result, prepare_task_drive, remove_subagent_task_drive from ouroboros.task_results import write_task_result @@ -82,7 +83,12 @@ def test_source_download_is_bound_and_distinct_from_same_named_user_file(tmp_pat app.state.drive_root = tmp_path with TestClient(app) as client: url = f"/api/tasks/applied/artifacts/{name}" - assert client.get(url).content == b"user result" + artifact = client.get(url) + if task_archive.CONFINED: + assert artifact.status_code == 200 and artifact.content == b"user result" + else: + assert artifact.status_code == 503 + assert artifact.json()["reason_code"] == "artifact_unavailable" source = client.get(url, params={"source": ref["path"]}) assert source.status_code == 200 assert hashlib.sha256(source.content).hexdigest() == ref["sha256"] diff --git a/tests/test_test_environment.py b/tests/test_test_environment.py index f8d84054b..6a83b46df 100644 --- a/tests/test_test_environment.py +++ b/tests/test_test_environment.py @@ -294,7 +294,7 @@ _SESSION_PROBE = """ import json, os, pathlib, tempfile def pytest_sessionstart(session): - worker = os.environ.get("PYTEST_XDIST_WORKER", "controller") + worker = getattr(session.config, "workerinput", {}).get("workerid", "controller") record = {"tmpdir": os.environ.get("TMPDIR", ""), "gettempdir": tempfile.gettempdir(), "session_root": str(pathlib.Path(os.environ["OUROBOROS_DATA_DIR"]).parent)} pathlib.Path(os.environ["SESSION_PROBE_OUT"], worker + ".json").write_text(json.dumps(record))