from __future__ import annotations import json import os import time from concurrent.futures import ThreadPoolExecutor from pathlib import Path from types import SimpleNamespace from unittest.mock import MagicMock, patch from uuid import uuid4 import pytest from langchain_core.language_models.fake_chat_models import FakeListChatModel from langchain_core.messages import AIMessage, AIMessageChunk from langchain_core.outputs import ChatGeneration, ChatGenerationChunk, LLMResult from pydantic import ValidationError from EvoScientist.usage.callback import ( UsageCaptureCallback, UsageModelIdentity, _scope, attach_usage_callback, ) from EvoScientist.usage.identity import ( _load_or_create, normalize_workspace_path_v1, prepare_usage_environment, workspace_identity, workspace_identity_from_normalized, ) from EvoScientist.usage.schema import UsageEventV1 from EvoScientist.usage.spool import UsageSpool IDENTITY = UsageModelIdentity( provider_profile_id="profile-a", provider_revision="revision-a", provider_adapter="openai", model_alias="chat-main", upstream_model_id="upstream-a", ) FIXTURES = Path( os.getenv( "EVOSCIENTIST_USAGE_FIXTURES", str( Path(__file__).parents[2] / "EvoScientist-WebUI" / "docs" / "schemas" / "fixtures" ), ) ) def _fixture_event() -> UsageEventV1: fixture = FIXTURES / "accepted" / "confirmed.json" return UsageEventV1.model_validate_json(fixture.read_text(encoding="utf-8")) def _usage_environment(monkeypatch: pytest.MonkeyPatch, tmp_path: Path) -> None: values = { "EVOSCIENTIST_USAGE_TRACKING": "true", "EVOSCIENTIST_USAGE_SINK_URL": "http://127.0.0.1:1/api/usage/events", "EVOSCIENTIST_USAGE_SINK_TOKEN": "test-token", "EVOSCIENTIST_USAGE_DEPLOYMENT_ID": "11111111-1111-4111-8111-111111111111", "EVOSCIENTIST_WORKSPACE_ID": "ws1_fixture", "EVOSCIENTIST_USAGE_SPOOL_DIR": str(tmp_path / "spool"), "EVOSCIENTIST_WORKSPACE_DIR": str(tmp_path), } for key, value in values.items(): monkeypatch.setenv(key, value) def test_prepare_usage_environment_rejects_invalid_webui_port(tmp_path: Path) -> None: with pytest.raises(ValueError, match="webui_port"): prepare_usage_environment(tmp_path, webui_port=0) @pytest.mark.parametrize("name", ["confirmed", "unknown"]) def test_shared_accepted_fixtures_match_python_contract(name: str) -> None: fixture = FIXTURES / "accepted" / f"{name}.json" event = UsageEventV1.model_validate_json(fixture.read_text(encoding="utf-8")) assert event.usage_status == name assert event.model_dump(mode="json")["completed_at"].endswith("Z") def test_shared_rejected_fixtures_match_python_contract() -> None: rejected = json.loads( (FIXTURES / "rejected" / "unsafe-token-sum.json").read_text(encoding="utf-8") ) base = json.loads( (FIXTURES / "rejected" / rejected["base_fixture"]).read_text(encoding="utf-8") ) base.update(rejected["patch"]) with pytest.raises(ValidationError, match="input plus output tokens"): UsageEventV1.model_validate(base) def test_unknown_usage_rejects_numeric_tokens() -> None: data = _fixture_event().model_dump() data.update( usage_status="unknown", input_tokens=0, output_tokens=None, provider_total_tokens=None, ) with pytest.raises(ValidationError, match="unknown usage"): UsageEventV1.model_validate(data) def test_python_contract_rejects_missing_and_coerced_fields() -> None: data = _fixture_event().model_dump(mode="json") data.pop("turn_id") with pytest.raises(ValidationError, match="turn_id"): UsageEventV1.model_validate(data) data = _fixture_event().model_dump(mode="json") data["input_tokens"] = "1000" with pytest.raises(ValidationError, match="input_tokens"): UsageEventV1.model_validate(data) data = _fixture_event().model_dump(mode="json") data.update(started_at=None, observed_at=0, completed_at=0) with pytest.raises(ValidationError, match="observed_at"): UsageEventV1.model_validate(data) def test_model_copy_preflight_attaches_callback_without_mutating_source( monkeypatch: pytest.MonkeyPatch, tmp_path: Path ) -> None: _usage_environment(monkeypatch, tmp_path) source = FakeListChatModel(responses=["ok"]) with patch("EvoScientist.usage.callback.get_usage_spool") as get_spool: tracked = attach_usage_callback(source, IDENTITY) get_spool.assert_called_once_with() assert tracked is not source assert not source.callbacks assert any( isinstance(item, UsageCaptureCallback) for item in tracked.callbacks or [] ) selector = tracked.model_copy( update={ "metadata": {**(tracked.metadata or {}), "usage_scope": "tool_selector"} } ) assert selector.metadata["usage_scope"] == "tool_selector" def test_async_run_inherits_usage_context_only_when_tracking_enabled( monkeypatch: pytest.MonkeyPatch, ) -> None: from EvoScientist.llm import patches monkeypatch.setenv("EVOSCIENTIST_USAGE_TRACKING", "true") active = { "metadata": { "turn_id": "turn-a", "thread_id": "thread-a", "source_agent": "EvoScientist", }, "configurable": {"thread_id": "thread-a"}, } with patch("langgraph.config.get_config", return_value=active): merged = patches._merge_runs_config_kwargs({"metadata": {"name": "writer"}}) assert merged["metadata"] == { "turn_id": "turn-a", "source_agent": "EvoScientist", "source_session_id": "thread-a", "name": "writer", "usage_scope": "async_subagent", } def test_callback_emits_one_confirmed_terminal_event( monkeypatch: pytest.MonkeyPatch, tmp_path: Path ) -> None: _usage_environment(monkeypatch, tmp_path) spool = MagicMock() callback = UsageCaptureCallback(IDENTITY) run_id = uuid4() callback.on_chat_model_start( {}, [], run_id=run_id, metadata={ "thread_id": "thread-a", "turn_id": "turn-a", "usage_scope": "tool_selector", }, ) message = AIMessage( content="ok", usage_metadata={"input_tokens": 12, "output_tokens": 3, "total_tokens": 15}, response_metadata={"request_id": "provider-request"}, ) result = LLMResult(generations=[[ChatGeneration(message=message)]]) with patch("EvoScientist.usage.callback.get_usage_spool", return_value=spool): callback.on_llm_end(result, run_id=run_id) callback.on_llm_end(result, run_id=run_id) spool.enqueue.assert_called_once() event = spool.enqueue.call_args.args[0] assert event.model_call_id == str(run_id) assert event.scope == "tool_selector" assert event.turn_id == "turn-a" assert event.provider_request_id == "provider-request" assert event.input_tokens == 12 assert event.output_tokens == 3 def test_callback_error_emits_unknown( monkeypatch: pytest.MonkeyPatch, tmp_path: Path ) -> None: _usage_environment(monkeypatch, tmp_path) spool = MagicMock() callback = UsageCaptureCallback(IDENTITY) run_id = uuid4() callback.on_chat_model_start( {}, [], run_id=run_id, metadata={"run_kind": "scheduled_task"} ) with patch("EvoScientist.usage.callback.get_usage_spool", return_value=spool): callback.on_llm_error(RuntimeError("provider failed"), run_id=run_id) event = spool.enqueue.call_args.args[0] assert event.usage_status == "unknown" assert event.input_tokens is None assert event.scope == "scheduler" def test_stream_usage_survives_terminal_error( monkeypatch: pytest.MonkeyPatch, tmp_path: Path ) -> None: _usage_environment(monkeypatch, tmp_path) spool = MagicMock() callback = UsageCaptureCallback(IDENTITY) run_id = uuid4() callback.on_chat_model_start({}, [], run_id=run_id, metadata={}) chunk = ChatGenerationChunk( message=AIMessageChunk( content="", usage_metadata={"input_tokens": 8, "output_tokens": 2, "total_tokens": 10}, ) ) callback.on_llm_new_token("", run_id=run_id, chunk=chunk) with patch("EvoScientist.usage.callback.get_usage_spool", return_value=spool): callback.on_llm_error(RuntimeError("stream interrupted"), run_id=run_id) event = spool.enqueue.call_args.args[0] assert event.usage_status == "confirmed" assert event.input_tokens == 8 assert event.output_tokens == 2 def test_callback_uses_deepagents_name_for_sync_subagent_scope( monkeypatch: pytest.MonkeyPatch, tmp_path: Path ) -> None: _usage_environment(monkeypatch, tmp_path) spool = MagicMock() callback = UsageCaptureCallback(IDENTITY) run_id = uuid4() callback.on_chat_model_start( {}, [], run_id=run_id, metadata={"lc_agent_name": "writing-agent"} ) with patch("EvoScientist.usage.callback.get_usage_spool", return_value=spool): callback.on_llm_error(RuntimeError("provider failed"), run_id=run_id) assert spool.enqueue.call_args.args[0].scope == "sync_subagent" @pytest.mark.parametrize( ("metadata", "expected"), [ ({"usage_scope": "tool_selector"}, "tool_selector"), ({"usage_scope": "diagnostic"}, "diagnostic"), ({"usage_scope": "skill_eval"}, "skill_eval"), ({"lc_source": "summarization"}, "summarizer"), ({"run_kind": "scheduled_task"}, "scheduler"), ({"run_kind": "evomemory_autoskills"}, "autoskills"), ({"run_kind": "evomemory_turn_worker"}, "memory"), ({"lc_agent_name": "EvoScientist"}, "main"), ({"lc_agent_name": "writing-agent"}, "sync_subagent"), ({"source_session_id": "origin-thread"}, "async_subagent"), ({"thread_id": "thread-a"}, "main"), ({}, "unattributed"), ], ) def test_scope_mapping_contract(metadata: dict, expected: str) -> None: assert _scope(metadata) == expected def test_provider_compatibility_replays_all_terminal_cases( monkeypatch: pytest.MonkeyPatch, tmp_path: Path ) -> None: _usage_environment(monkeypatch, tmp_path) compatibility = json.loads( (FIXTURES / "providers" / "compatibility.json").read_text(encoding="utf-8") ) for provider in compatibility["providers"]: for case in compatibility["cases"]: identity = UsageModelIdentity( provider_profile_id=provider["name"], provider_revision="fixture", provider_adapter=provider["adapter"], model_alias="fixture", upstream_model_id="fixture", ) callback = UsageCaptureCallback(identity) spool = MagicMock() run_id = uuid4() callback.on_chat_model_start({}, [], run_id=run_id, metadata={}) request_field = provider["request_id_field"].split(".")[-1] if case.get("stream_usage"): chunk = ChatGenerationChunk( message=AIMessageChunk( content="", usage_metadata=case["stream_usage"], response_metadata={request_field: "request-fixture"}, ) ) callback.on_llm_new_token("", run_id=run_id, chunk=chunk) with patch( "EvoScientist.usage.callback.get_usage_spool", return_value=spool ): if case["terminal"] == "error": callback.on_llm_error(RuntimeError("fixture error"), run_id=run_id) else: message = AIMessage( content="ok", usage_metadata=case.get("usage"), response_metadata={request_field: "request-fixture"}, ) callback.on_llm_end( LLMResult(generations=[[ChatGeneration(message=message)]]), run_id=run_id, ) event = spool.enqueue.call_args.args[0] assert event.usage_status == case["expected_status"] if case["expected_status"] == "confirmed": assert event.input_tokens is not None assert event.output_tokens is not None def test_workspace_identity_is_stable_and_deployment_scoped(tmp_path: Path) -> None: first = workspace_identity("11111111-1111-4111-8111-111111111111", tmp_path) second = workspace_identity("11111111-1111-4111-8111-111111111111", tmp_path) other = workspace_identity("22222222-2222-4222-8222-222222222222", tmp_path) assert first == second assert first.startswith("ws1_") assert first != other def test_identity_file_creation_is_atomic_across_launchers(tmp_path: Path) -> None: identity_path = tmp_path / "deployment-id" with ThreadPoolExecutor(max_workers=16) as pool: values = list( pool.map( lambda index: _load_or_create( identity_path, lambda: f"launcher-{index}" ), range(64), ) ) assert len(set(values)) == 1 assert identity_path.read_text(encoding="utf-8").strip() == values[0] if os.name != "nt": assert identity_path.stat().st_mode & 0o777 == 0o600 def test_callback_remains_fail_open_when_spool_write_fails( monkeypatch: pytest.MonkeyPatch, tmp_path: Path ) -> None: _usage_environment(monkeypatch, tmp_path) callback = UsageCaptureCallback(IDENTITY) run_id = uuid4() callback.on_chat_model_start({}, [], run_id=run_id, metadata={}) message = AIMessage( content="ok", usage_metadata={"input_tokens": 1, "output_tokens": 1, "total_tokens": 2}, ) result = LLMResult(generations=[[ChatGeneration(message=message)]]) spool = MagicMock() spool.enqueue.side_effect = OSError("disk full") with patch("EvoScientist.usage.callback.get_usage_spool", return_value=spool): callback.on_llm_end(result, run_id=run_id) def test_workspace_identity_normalization_contract(tmp_path: Path) -> None: deployment = "11111111-1111-4111-8111-111111111111" assert ( normalize_workspace_path_v1("/tmp/research/", windows=False) == "/tmp/research" ) assert normalize_workspace_path_v1("/", windows=False) == "/" assert normalize_workspace_path_v1("/tmp/Cafe\u0301", windows=False) == "/tmp/Café" assert normalize_workspace_path_v1("C:\\Research\\", windows=True) == "c:/research" assert normalize_workspace_path_v1("C:\\", windows=True) == "c:/" assert ( normalize_workspace_path_v1("\\\\Server\\Share\\Research\\", windows=True) == "//server/share/research" ) normalized = normalize_workspace_path_v1("C:\\Research\\", windows=True) assert workspace_identity_from_normalized(deployment, normalized).startswith("ws1_") real = tmp_path / "real" real.mkdir() link = tmp_path / "link" link.symlink_to(real, target_is_directory=True) assert workspace_identity(deployment, real) == workspace_identity(deployment, link) posix = json.loads( (FIXTURES / "identity" / "workspace-posix.json").read_text(encoding="utf-8") ) assert ( workspace_identity_from_normalized(deployment, posix["normalized_path"]) == posix["expected_workspace_id"] ) windows = json.loads( (FIXTURES / "identity" / "workspace-windows.json").read_text(encoding="utf-8") ) for case in windows["cases"]: normalized_case = normalize_workspace_path_v1(case["input"], windows=True) assert normalized_case == case["normalized_path"] assert ( workspace_identity_from_normalized(deployment, normalized_case) == case["expected_workspace_id"] ) unicode_case = json.loads( (FIXTURES / "identity" / "workspace-unicode.json").read_text(encoding="utf-8") ) normalized_unicode = normalize_workspace_path_v1( unicode_case["input"], windows=False ) assert normalized_unicode == unicode_case["normalized_path"] assert ( workspace_identity_from_normalized(deployment, normalized_unicode) == unicode_case["expected_workspace_id"] ) root_case = json.loads( (FIXTURES / "identity" / "workspace-root-and-symlink.json").read_text( encoding="utf-8" ) )["posix_root"] assert ( workspace_identity_from_normalized(deployment, root_case["normalized_path"]) == root_case["expected_workspace_id"] ) def test_spool_enqueue_is_durable_without_waiting_for_http( monkeypatch: pytest.MonkeyPatch, tmp_path: Path ) -> None: _usage_environment(monkeypatch, tmp_path) monkeypatch.setattr(UsageSpool, "_run", lambda self: self._stop.wait()) spool = UsageSpool() try: spool.enqueue(_fixture_event()) pending = list(spool.pending.glob("*.json")) assert len(pending) == 1 assert json.loads(pending[0].read_text(encoding="utf-8"))["schema_version"] == 1 finally: spool.close() def _event_with_id(event: UsageEventV1, model_call_id: str) -> UsageEventV1: return event.model_copy( update={ "model_call_id": model_call_id, "event_id": (f"{event.deployment_id}:{model_call_id}:callback_final:1"), } ) def test_spool_multi_worker_same_event_is_idempotent( monkeypatch: pytest.MonkeyPatch, tmp_path: Path ) -> None: _usage_environment(monkeypatch, tmp_path) monkeypatch.setattr(UsageSpool, "_run", lambda self: self._stop.wait()) first = UsageSpool() second = UsageSpool() try: event = _fixture_event() with ThreadPoolExecutor(max_workers=2) as pool: list(pool.map(lambda spool: spool.enqueue(event), (first, second))) assert len(list(first.pending.glob("*.json"))) == 1 assert list(first.quarantine.glob("*.json")) == [] finally: first.close() second.close() def test_spool_recovers_stale_inflight_and_rejects_at_soft_limit( monkeypatch: pytest.MonkeyPatch, tmp_path: Path ) -> None: _usage_environment(monkeypatch, tmp_path) monkeypatch.setenv("EVOSCIENTIST_USAGE_INFLIGHT_LEASE_SECONDS", "1") monkeypatch.setenv("EVOSCIENTIST_USAGE_SPOOL_MAX_FILES", "1") monkeypatch.setattr(UsageSpool, "_run", lambda self: self._stop.wait()) spool = UsageSpool() try: event = _fixture_event() spool.enqueue(event) pending = next(spool.pending.glob("*.json")) inflight = spool.inflight / pending.name pending.replace(inflight) old = time.time() - 5 os.utime(inflight, (old, old)) spool._recover_stale_inflight() assert (spool.pending / inflight.name).exists() second = _event_with_id(event, "77777777-7777-4777-8777-777777777777") spool.enqueue(second) assert len(list(spool.pending.glob("*.json"))) == 1 assert spool.first_loss_at is not None assert spool.degraded_reason == "spool_soft_limit_reached" finally: spool.close() @pytest.mark.parametrize( ("status_code", "response_status", "expected_location", "progressed"), [ (200, "accepted", "deleted", True), (200, "duplicate", "deleted", True), (409, "conflict", "quarantine", True), (422, "rejected", "quarantine", True), (401, "unauthorized", "pending", False), (429, "rate_limited", "pending", False), (500, "error", "pending", False), ], ) def test_sender_response_state_machine( monkeypatch: pytest.MonkeyPatch, tmp_path: Path, status_code: int, response_status: str, expected_location: str, progressed: bool, ) -> None: _usage_environment(monkeypatch, tmp_path) monkeypatch.setattr(UsageSpool, "_run", lambda self: self._stop.wait()) spool = UsageSpool() response = SimpleNamespace( status_code=status_code, json=lambda: {"status": response_status}, ) client = SimpleNamespace(post=lambda *args, **kwargs: response) try: spool.enqueue(_fixture_event()) assert spool._send_one(client) is progressed if expected_location == "deleted": assert list(spool.pending.glob("*.json")) == [] assert list(spool.inflight.glob("*.json")) == [] else: directory = getattr(spool, expected_location) assert len(list(directory.glob("*.json"))) == 1 finally: spool.close()