Files

1209 lines
46 KiB
Python

from __future__ import annotations
import hashlib
import math
import uuid
from datetime import UTC, datetime, timedelta
from typing import Literal
from unittest.mock import Mock
import pytest
from fastapi.testclient import TestClient
from sqlalchemy import create_engine, select, update
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from sqlalchemy.pool import StaticPool
from modelforge_api.api.routes.serving import get_serving_service
from modelforge_api.domain.runtime import AgentRuntimeProbeComplete
from modelforge_api.domain.serving import (
AgentResidentState,
AgentServingJobComplete,
AgentServingStateReport,
CapabilityExperimentCreate,
CapabilityPromotionCreate,
CoResidencyEvidenceCreate,
EmbeddingInvokeRequest,
PlacementPlanRequest,
ProductionApprovalCreate,
ProjectFitEvidenceCreate,
RerankDocument,
RerankingInvokeRequest,
ServiceClientCreate,
SupplyChainReview,
)
from modelforge_api.main import app
from modelforge_api.persistence.models import (
AcceleratorTelemetryLatest,
ArtifactLocation,
ArtifactSet,
Base,
ComputeNode,
EmbeddingSpace,
GatewayRequest,
ServiceCredential,
ServingGpuLease,
ServingJob,
)
from modelforge_api.services.manifest_registry import ManifestRegistry
from modelforge_api.services.project_registry import sync_project_registry
from modelforge_api.services.serving import ServingError, ServingService
from modelforge_api.services.transient_payloads import MemoryPayloadStore
from modelforge_api.settings import Settings, get_settings
from tests.test_runtime import prepare_probe, setup_runtime
@pytest.fixture
def session() -> Session:
engine = create_engine(
"sqlite+pysqlite:///:memory:",
connect_args={"check_same_thread": False},
poolclass=StaticPool,
)
Base.metadata.create_all(engine)
with Session(engine) as value:
yield value
def review() -> SupplyChainReview:
return SupplyChainReview(
exact_revision_reviewed=True,
artifact_hashes_reviewed=True,
safetensors_only=True,
remote_code_required=False,
pickle_present=False,
scanner_evidence_reviewed=True,
provenance_complete=True,
runtime_offline_reviewed=True,
dependency_provenance_reviewed=True,
license_reviewed=True,
license_identifier="Apache-2.0",
)
def test_non_authentication_integrity_conflicts_retain_the_state_conflict_contract(
session: Session,
monkeypatch: pytest.MonkeyPatch,
) -> None:
service = ServingService(
session,
Settings(_env_file=None),
ManifestRegistry(),
MemoryPayloadStore(),
)
failed_commit = Mock(
side_effect=IntegrityError(
"INSERT INTO capability_deployments ...",
{},
RuntimeError("genuine state conflict"),
)
)
rollback = Mock(wraps=session.rollback)
monkeypatch.setattr(session, "commit", failed_commit)
monkeypatch.setattr(session, "rollback", rollback)
with pytest.raises(ServingError) as raised:
service._commit()
assert raised.value.status_code == 409
assert raised.value.code == "STATE_CONFLICT"
assert failed_commit.call_count == 1
assert rollback.call_count == 1
def setup_m5(session: Session) -> tuple[ServingService, object, ComputeNode]:
runtime, artifact_set, node, _artifacts = setup_runtime(session)
node.production_eligible = True
node.agent_capabilities = [
*node.agent_capabilities,
"deployment.load.v1",
"deployment.invoke.v1",
"deployment.health.v1",
"deployment.drain.v1",
"deployment.unload.v1",
]
session.commit()
_profile, _assessment, _approval, probe = prepare_probe(runtime, artifact_set, node)
lease = runtime.claim_next(node)
assert lease is not None
runtime.complete(
probe.id,
node,
AgentRuntimeProbeComplete(
lease_token=lease.lease_token,
load_result={"status": "passed", "load_time_ms": 3690.5},
health_result={
"process": "healthy",
"runtime": "healthy",
"model": "healthy",
"capability": "not_routed_in_m4",
},
inference_result={
"status": "passed",
"shape": [1, 1024],
"dimension": 1024,
"finite": True,
"latency_ms": 1203.2,
},
unload_result={"status": "passed", "reclaimed": True},
measured_resources={
"kind": "measured",
"samples": {
"before_load": {"used_vram_bytes": 1_130_168_320},
"loaded": {"used_vram_bytes": 2_610_036_736},
"probe_peak": {"used_vram_bytes": 2_668_756_992},
"after_unload": {"used_vram_bytes": 1_130_168_320},
},
"baseline_vram_bytes": 1_130_168_320,
"loaded_vram_bytes": 2_610_036_736,
"peak_vram_bytes": 2_668_756_992,
"delta_peak_vram_bytes": 1_538_588_672,
"batch_size": 1,
"max_sequence_length": 128,
},
runtime_facts={
"offline_local_only": True,
"trust_remote_code": False,
"network_attempts": 0,
"gpu_uuid": "GPU-test",
},
environment_fingerprint=hashlib.sha256(b"m5-runtime").hexdigest(),
),
)
accelerator = runtime.repo.accelerators(node.id)[0]
session.add(
AcceleratorTelemetryLatest(
accelerator_id=accelerator.id,
used_vram_bytes=1_130_168_320,
free_vram_bytes=16_041_312_256,
gpu_utilization_percent=0,
memory_utilization_percent=0,
temperature_c=50,
power_draw_w=10.0,
power_limit_w=320.0,
graphics_clock_mhz=0,
memory_clock_mhz=0,
fan_speed_percent=0,
performance_state="P8",
availability={},
observed_at=datetime.now(UTC),
)
)
session.commit()
service = ServingService(
session,
Settings(
database_url="sqlite+pysqlite:///:memory:",
runtime_artifact_root="/models",
scheduler_safety_reserve_bytes=1_073_741_824,
gateway_request_timeout_seconds=5,
gateway_queue_timeout_seconds=2,
),
ManifestRegistry(),
MemoryPayloadStore(),
)
return service, runtime.candidates()[0], node
def approve_and_promote(service: ServingService, candidate_id: uuid.UUID):
approval = service.approve_production(
candidate_id,
ProductionApprovalCreate(
approved_by="security-operator",
reason="Exact artifact, runtime and capability production review completed",
supply_chain_review=review(),
deployment_config={"normalization": True, "dimension": 1024},
),
)
deployment = service.promote(
candidate_id,
CapabilityPromotionCreate(
production_approval_id=approval.id,
keep_warm_seconds=900,
),
)
return approval, deployment
def test_operational_embedding_contract_is_versioned_and_requires_reindex() -> None:
contract = next(
item for item in ManifestRegistry().capabilities() if item.capability == "rag.embedding"
)
assert contract.version == 1
assert contract.vector is not None
assert contract.vector.dimensionality == 1024
assert contract.vector.normalized is True
assert contract.upgrade_class.value == "requires_reindex"
assert "input" in contract.input_schema["properties"]
def test_production_approval_requires_completed_probe_and_current_node(session: Session) -> None:
runtime, artifact_set, node, _artifacts = setup_runtime(session)
node.production_eligible = True
session.commit()
_profile, _assessment, _lab, _probe = prepare_probe(runtime, artifact_set, node)
service = ServingService(session, Settings(), ManifestRegistry(), MemoryPayloadStore())
candidate_id = uuid.uuid4()
with pytest.raises(ServingError, match="LAB_READY"):
service.approve_production(
candidate_id,
ProductionApprovalCreate(
approved_by="operator",
reason="Production review cannot bypass runtime evidence",
supply_chain_review=review(),
),
)
def test_exact_production_approval_and_promotion_are_idempotent(session: Session) -> None:
service, candidate, _node = setup_m5(session)
first_approval, first = approve_and_promote(service, candidate.id)
second_approval = service.approve_production(
candidate.id,
ProductionApprovalCreate(
approved_by="security-operator",
reason="Exact artifact, runtime and capability production review completed",
supply_chain_review=review(),
deployment_config={"normalization": True, "dimension": 1024},
),
)
second = service.promote(
candidate.id,
CapabilityPromotionCreate(
production_approval_id=first_approval.id,
keep_warm_seconds=900,
),
)
assert first_approval.id == second_approval.id
assert first.id == second.id
assert first.status == "stable" and first.production is True
assert first.resource_envelope.required_vram_bytes == 1_538_588_672
assert first.embedding_space.migration_class == "requires_reindex"
def test_resource_envelope_prefers_process_owned_cuda_allocation(session: Session) -> None:
service, candidate, _node = setup_m5(session)
probe = service.repo.probe(candidate.runtime_probe_id)
assert probe is not None
measured = dict(probe.measured_resources)
measured["samples"] = {
**measured["samples"],
"before_load": {
**measured["samples"]["before_load"],
"process_used_vram_bytes": 0,
},
"loaded": {
**measured["samples"]["loaded"],
"process_used_vram_bytes": 1_216_348_160,
},
"probe_peak": {
**measured["samples"]["probe_peak"],
"process_used_vram_bytes": 1_216_348_160,
},
}
probe.measured_resources = measured
session.commit()
_approval, deployment = approve_and_promote(service, candidate.id)
assert deployment.resource_envelope.required_vram_bytes == 1_216_348_160
assert deployment.resource_envelope.resident_vram_bytes == 1_216_348_160
def test_same_dimension_does_not_imply_same_embedding_space(session: Session) -> None:
service, candidate, _node = setup_m5(session)
_approval, deployment = approve_and_promote(service, candidate.id)
original = session.get(EmbeddingSpace, deployment.embedding_space.id)
assert original is not None
other = EmbeddingSpace(
capability_contract_id=original.capability_contract_id,
artifact_set_id=original.artifact_set_id,
runtime_profile_id=original.runtime_profile_id,
identity_digest="b" * 64,
dimension=1024,
normalized=True,
migration_class="requires_reindex",
identity_facts={"artifact_revision": "different"},
created_at=datetime.now(UTC),
immutable_at=datetime.now(UTC),
)
session.add(other)
session.commit()
assert other.dimension == original.dimension
assert other.identity_digest != original.identity_digest
def test_service_credentials_are_scoped_hash_only_and_revocable(session: Session) -> None:
service, _candidate, _node = setup_m5(session)
service.ensure_contract()
created = service.create_client(
ServiceClientCreate(
name="m5-acceptance",
allowed_capabilities=["rag.embedding@1"],
requests_per_minute=2,
)
)
stored = session.scalar(select(ServiceCredential))
assert stored is not None
assert stored.secret_hash == hashlib.sha256(created.credential.encode()).hexdigest()
assert created.credential not in str(service.clients())
assert service.authenticate(f"Bearer {created.credential}", "rag.embedding@1").id == created.id
with pytest.raises(ServingError, match="scope"):
service.authenticate(f"Bearer {created.credential}", "assistant.general@1")
service.revoke_client_credential(created.id)
with pytest.raises(ServingError, match="revoked"):
service.authenticate(f"Bearer {created.credential}", "rag.embedding@1")
def test_project_client_binding_rotation_usage_and_fit_are_operational(session: Session) -> None:
service, _candidate, _node = setup_m5(session)
service.ensure_contract("vision.embedding", 1)
sync_project_registry(session, ManifestRegistry())
created = service.create_client(
ServiceClientCreate(
name="examplevision-vision-shadow",
allowed_capabilities=["vision.embedding@1"],
requests_per_minute=120,
workload_priority="background",
project_key="examplevision",
integration_environment="shadow",
purpose="card visual candidate retrieval",
)
)
assert created.project_key == "examplevision"
assert created.project_binding_id is not None
assert (
service.authenticate(f"Bearer {created.credential}", "vision.embedding@1").id == created.id
)
rotated = service.rotate_client_credential(created.id)
with pytest.raises(ServingError, match="revoked"):
service.authenticate(f"Bearer {created.credential}", "vision.embedding@1")
assert (
service.authenticate(f"Bearer {rotated.credential}", "vision.embedding@1").id == created.id
)
fit = service.record_project_fit(
ProjectFitEvidenceCreate(
project_key="examplevision",
capability="vision.embedding@1",
environment="evaluation",
recommendation="REQUIRES_MORE_EVIDENCE",
evidence_class="CATALOG_REFERENCE",
engineering_integration="PASS",
production_validation="DEFERRED_EXTERNAL_VALIDATION",
case_count=0,
metric_values={},
critical_errors=0,
blockers=["real_photo_dataset_unavailable"],
evidence={"source": "owner_media_inventory"},
)
)
assert fit.project_key == "examplevision"
assert fit.evidence_digest
integration = service.project_integrations()[0]
assert integration.project_key == "examplevision"
assert integration.client_name == "examplevision-vision-shadow"
assert integration.usage.request_volume == 0
assert integration.project_fit is not None
assert integration.project_fit.recommendation == "REQUIRES_MORE_EVIDENCE"
assert integration.project_fit.evidence_class == "CATALOG_REFERENCE"
assert integration.project_fit.engineering_integration == "PASS"
assert integration.project_fit.production_validation == "DEFERRED_EXTERNAL_VALIDATION"
assert integration.project_fit.production_action == "NONE"
session.add_all(
[
GatewayRequest(
request_id=uuid.uuid4(),
service_client_id=created.id,
capability_key="vision.embedding",
capability_version=1,
status="completed",
priority="background",
input_sha256="a" * 64,
input_count=1,
total_latency_ms=20.0,
),
GatewayRequest(
request_id=uuid.uuid4(),
service_client_id=created.id,
capability_key="vision.embedding",
capability_version=1,
status="failed",
priority="background",
input_sha256="b" * 64,
input_count=1,
total_latency_ms=30.0,
failure_code="CAPACITY_CONSTRAINED",
),
]
)
session.commit()
usage = service.project_integrations()[0].usage
assert usage.request_volume == 2
assert usage.successful_requests == 1
assert usage.error_count == 1
assert usage.latency_p50_ms == 20.0
def test_project_fit_cannot_claim_eligibility_with_blockers(session: Session) -> None:
service, _candidate, _node = setup_m5(session)
service.ensure_contract("vision.embedding", 1)
sync_project_registry(session, ManifestRegistry())
with pytest.raises(ServingError, match="production eligibility"):
service.record_project_fit(
ProjectFitEvidenceCreate(
project_key="examplevision",
capability="vision.embedding@1",
environment="evaluation",
recommendation="PROMOTION_ELIGIBLE",
evidence_class="OWNER_PHOTO",
engineering_integration="PASS",
production_validation="SATISFIED",
case_count=1,
metric_values={"recall_at_1": 1.0},
critical_errors=0,
blockers=["dataset_not_representative"],
)
)
def test_project_fit_satisfies_production_gate_only_with_owner_photo_evidence(
session: Session,
) -> None:
service, _candidate, _node = setup_m5(session)
service.ensure_contract("vision.embedding", 1)
sync_project_registry(session, ManifestRegistry())
evidence = service.record_project_fit(
ProjectFitEvidenceCreate(
project_key="examplevision",
capability="vision.embedding@1",
environment="evaluation",
recommendation="PROMOTION_ELIGIBLE",
evidence_class="OWNER_PHOTO",
engineering_integration="PASS",
production_validation="SATISFIED",
case_count=30,
metric_values={"recall_at_1": 1.0},
critical_errors=0,
)
)
assert evidence.production_validation == "SATISFIED"
assert evidence.production_action == "NONE"
@pytest.mark.parametrize("evidence_class", ["CATALOG_REFERENCE", "PUBLIC_PHYSICAL_CAPTURE"])
def test_project_fit_cannot_satisfy_production_gate_without_owner_photos(
session: Session,
evidence_class: Literal["CATALOG_REFERENCE", "PUBLIC_PHYSICAL_CAPTURE"],
) -> None:
service, _candidate, _node = setup_m5(session)
service.ensure_contract("vision.embedding", 1)
sync_project_registry(session, ManifestRegistry())
with pytest.raises(ServingError, match="owner-photo"):
service.record_project_fit(
ProjectFitEvidenceCreate(
project_key="examplevision",
capability="vision.embedding@1",
environment="evaluation",
recommendation="PROMOTION_ELIGIBLE",
evidence_class=evidence_class,
engineering_integration="PASS",
production_validation="SATISFIED",
case_count=30,
metric_values={"recall_at_1": 1.0},
critical_errors=0,
)
)
def test_scheduler_accounts_external_usage_reserve_and_workstation_rejection(
session: Session,
) -> None:
service, candidate, _node = setup_m5(session)
_approval, deployment = approve_and_promote(service, candidate.id)
workstation = ComputeNode(
key="workstation",
hostname="DESKTOP",
display_name="DESKTOP",
enabled=True,
liveness_state="online",
production_eligible=False,
)
session.add(workstation)
session.commit()
budget = service.scheduler_overview()[0]
assert len(service.scheduler_overview()) == 1
assert budget.external_vram_bytes == 1_130_168_320
assert budget.safety_reserve_bytes == 1_073_741_824
stored = service.repo.deployment(deployment.id)
assert stored is not None
evidence = service._placement_evidence(stored)
rejected = next(item for item in evidence["candidate_nodes"] if item["node"] == "DESKTOP")
assert "NODE_NOT_PRODUCTION_ELIGIBLE" in rejected["reasons"]
class InlineWorkerServingService(ServingService):
def _wait_job(self, job_id: uuid.UUID, deadline: float):
del deadline
job = self.repo.job(job_id)
assert job is not None
node = self.repo.node(job.compute_node_id)
assert node is not None
lease = self.claim_next(node)
assert lease is not None and lease.job_id == job_id
if lease.operation == "load":
result = {}
metrics = {
"load_time_ms": 200.0,
"baseline_vram_bytes": 1_130_168_320,
"loaded_vram_bytes": 2_610_036_736,
}
elif lease.operation == "unload":
result = {"reclaimed": True}
metrics = {
"after_unload_vram_bytes": 1_130_168_320,
"reclaimed_vram_bytes": 1_479_868_416,
}
elif lease.operation == "drain":
result = {"drained": True}
metrics = {}
elif lease.capability == "rag.embedding":
assert lease.input is not None
vector = [1.0] + [0.0] * 1023
result = {
"vectors": [vector for _ in lease.input],
"count": len(lease.input),
"normalized": True,
"prompt_tokens": len(lease.input) * 4,
}
metrics = {"inference_time_ms": 25.0}
else:
assert lease.rerank_query is not None
assert lease.rerank_documents is not None
assert lease.top_n is not None
documents = list(reversed(lease.rerank_documents))[: lease.top_n]
result = {
"results": [
{"id": document.id, "score": 1.0 - index / 10, "rank": index + 1}
for index, document in enumerate(documents)
],
"count": len(lease.rerank_documents),
"top_n": lease.top_n,
"finite": True,
"prompt_tokens": 18,
}
metrics = {
"inference_time_ms": 30.0,
"worker_total_time_ms": 35.0,
}
self.complete_job(
job_id,
node,
AgentServingJobComplete(
lease_token=lease.lease_token,
worker_instance_id="inline-worker",
state="warm",
result=result,
metrics=metrics,
health={"process": "healthy", "runtime": "healthy", "model": "healthy"},
runtime_facts={"offline_local_only": True},
),
)
completed = self.repo.job(job_id)
assert completed is not None
payload = (
self.payload_store.get(completed.result_reference) if completed.result_reference else {}
)
return completed, payload or {}
def test_real_gateway_path_uses_scheduler_cold_load_then_warm_reuse(session: Session) -> None:
base, candidate, _node = setup_m5(session)
_approval, deployment = approve_and_promote(base, candidate.id)
service = InlineWorkerServingService(
session,
base.settings,
base.manifests,
base.payload_store,
)
created = service.create_client(
ServiceClientCreate(
name="gateway-client",
allowed_capabilities=["rag.embedding@1"],
requests_per_minute=10,
)
)
client = service.authenticate(f"Bearer {created.credential}", "rag.embedding@1")
cold = service.invoke(EmbeddingInvokeRequest(input="ModelForge cold"), client)
client = service.authenticate(f"Bearer {created.credential}", "rag.embedding@1")
warm = service.invoke(EmbeddingInvokeRequest(input=["ModelForge warm"]), client)
app.dependency_overrides[get_settings] = lambda: service.settings
app.dependency_overrides[get_serving_service] = lambda: service
try:
http_response = TestClient(app).post(
"/api/v1/capabilities/rag.embedding@1/invoke",
headers={"Authorization": f"Bearer {created.credential}"},
json={"input": "ModelForge authenticated HTTP control"},
)
finally:
app.dependency_overrides.clear()
assert cold.execution.cold is True and warm.execution.cold is False
assert cold.dimension == warm.dimension == 1024
assert cold.embedding_space_id == warm.embedding_space_id == deployment.embedding_space.id
assert cold.execution.load_count == warm.execution.load_count == 1
assert http_response.status_code == 200, http_response.text
assert http_response.json()["dimension"] == 1024
assert len(cold.data[0]) == 1024 and all(math.isfinite(value) for value in cold.data[0])
stored_request = session.scalar(
select(GatewayRequest).where(GatewayRequest.request_id == cold.request_id)
)
assert stored_request is not None and stored_request.input_sha256
latency = stored_request.decision_evidence["latency_ms"]
assert latency["total"] >= latency["inference"]
assert latency["dispatch"] >= 0
assert cold.execution.timings.scheduling_ms == cold.execution.timings.queue_ms
history = service.request_history()
assert history[0].latency_breakdown["worker_total"] >= 0
def test_lab_experiment_has_opaque_route_own_space_and_no_production_cutover(
session: Session,
) -> None:
base, candidate, _node = setup_m5(session)
_approval, stable = approve_and_promote(base, candidate.id)
probe = base.repo.probe(candidate.runtime_probe_id)
assert probe is not None
experiment = base.create_experiment(
candidate.id,
CapabilityExperimentCreate(
route_key="m7-retrieval-aware",
execution_approval_id=probe.execution_approval_id,
purpose="Compare a retrieval-aware input profile without changing stable routing.",
),
)
assert experiment.deployment.production is False
assert experiment.deployment.production_approval_id is None
assert experiment.deployment.execution_approval_id == probe.execution_approval_id
assert experiment.deployment.embedding_space.id != stable.embedding_space.id
stable_row = base.repo.deployment(stable.id)
assert stable_row is not None and stable_row.status == "stable" and stable_row.production
service = InlineWorkerServingService(session, base.settings, base.manifests, base.payload_store)
created = service.create_client(
ServiceClientCreate(
name="experiment-client",
allowed_capabilities=["rag.embedding@1"],
workload_priority="benchmark",
)
)
client = service.authenticate(f"Bearer {created.credential}", "rag.embedding@1")
deployment = service.repo.deployment(experiment.capability_deployment_id)
assert deployment is not None
result = service.invoke(
EmbeddingInvokeRequest(input="typed query", input_type="query"),
client,
deployment=deployment,
experiment_route=experiment.route_key,
)
assert result.embedding_space_id == experiment.deployment.embedding_space.id
resolved_stable = base.repo.stable_deployment(stable_row.capability_contract_id)
assert resolved_stable is not None and resolved_stable.id == stable.id
deactivated = base.deactivate_experiment(experiment.id)
assert deactivated.status == "inactive"
assert base.deactivate_experiment(experiment.id).status == "inactive"
def test_candidate_unload_remains_ready_on_demand_and_live_pair_can_be_evidenced(
session: Session,
) -> None:
base, candidate, _node = setup_m5(session)
_approval, stable = approve_and_promote(base, candidate.id)
probe = base.repo.probe(candidate.runtime_probe_id)
assert probe is not None
experiment = base.create_experiment(
candidate.id,
CapabilityExperimentCreate(
route_key="m10-live-pair",
execution_approval_id=probe.execution_approval_id,
purpose="Prove that exact co-resident deployments remain schedulable after unload.",
),
)
service = InlineWorkerServingService(session, base.settings, base.manifests, base.payload_store)
stable_row = service.repo.deployment(stable.id)
candidate_row = service.repo.deployment(experiment.capability_deployment_id)
stable_residency = service.repo.residency(stable.id)
candidate_residency = service.repo.residency(experiment.capability_deployment_id)
telemetry = service.repo.telemetry(stable.accelerator_id)
assert stable_row and candidate_row and stable_residency and candidate_residency and telemetry
for residency in (stable_residency, candidate_residency):
residency.state = "warm"
residency.worker_instance_id = "exact-m10-worker"
residency.generation = 7
residency.measured_resident_vram_bytes = 1_479_868_416
telemetry.used_vram_bytes = 5_000_000_000
session.commit()
evidence = service.record_co_residency_evidence(
CoResidencyEvidenceCreate(
left_deployment_id=stable.id,
right_deployment_id=experiment.capability_deployment_id,
)
)
assert evidence.status == "PROVEN_SAFE"
assert evidence.measured_combined_bytes == 2_959_736_832
assert evidence.evidence["content_persisted"] is False
service.request_unload(experiment.capability_deployment_id)
queued = list(
session.scalars(
select(ServingJob)
.where(
ServingJob.capability_deployment_id == experiment.capability_deployment_id,
ServingJob.status == "queued",
)
.order_by(ServingJob.created_at)
)
)
assert [item.operation for item in queued] == ["drain", "unload"]
for job in queued:
service._wait_job(job.id, 0)
session.refresh(candidate_row)
assert candidate_row.health_status == "ready_on_demand"
candidate_row.health_status = "unavailable"
session.commit()
service.reconcile()
session.refresh(candidate_row)
assert candidate_row.health_status == "ready_on_demand"
plan = service.dry_run_placement(
candidate_row.id,
PlacementPlanRequest(priority="lab"),
)
assert plan.verdict != "REJECT_HEALTH"
def test_worker_generation_accepts_64_bit_identity(session: Session) -> None:
service, candidate, node = setup_m5(session)
_approval, deployment = approve_and_promote(service, candidate.id)
stored = service.repo.deployment(deployment.id)
residency = service.repo.residency(deployment.id)
profile = service.repo.profile(deployment.runtime_profile_id)
assert stored and residency and profile and profile.runtime_environment_id
environment = service.repo.environment(profile.runtime_environment_id)
assert environment is not None
generation = 1_787_748_198_899_462_925
service.report_state(
node,
AgentServingStateReport(
worker_instance_id="m10-64-bit-worker",
deployment_id=None,
generation=generation,
state="warm",
load_count=1,
health={"process": "healthy", "runtime": "healthy", "model": "healthy"},
observed_at=datetime.now(UTC),
residencies=[
AgentResidentState(
deployment_id=deployment.id,
artifact_set_id=deployment.artifact_set_id,
runtime_profile_id=deployment.runtime_profile_id,
runtime_environment_fingerprint=environment.fingerprint,
revision_sha="a" * 40,
state="warm",
resident_vram_bytes=1_479_868_416,
load_count=1,
health={
"process": "healthy",
"runtime": "healthy",
"model": "healthy",
},
)
],
),
)
session.refresh(residency)
assert residency.generation == generation
def test_reranking_gateway_is_scoped_bounded_and_content_free(session: Session) -> None:
base, candidate, _node = setup_m5(session)
probe = base.repo.probe(candidate.runtime_probe_id)
assert probe is not None
probe.inference_result = {
"status": "passed",
"output_type": "ranked_scores",
"count": 1,
"finite": True,
"latency_ms": 30.0,
}
probe.runtime_facts = {
**probe.runtime_facts,
"adapter": "qwen3_reranker",
"offline_local_only": True,
"trust_remote_code": False,
"network_attempts": 0,
}
session.commit()
experiment = base.create_experiment(
candidate.id,
CapabilityExperimentCreate(
route_key="m8-qwen-reranker",
execution_approval_id=probe.execution_approval_id,
capability="rag.reranking",
purpose="Evaluate fixed top-forty reranking without production cutover.",
),
)
assert experiment.capability == "rag.reranking"
assert experiment.deployment.embedding_space is None
service = InlineWorkerServingService(session, base.settings, base.manifests, base.payload_store)
created = service.create_client(
ServiceClientCreate(
name="examplerag-reranking",
allowed_capabilities=["rag.reranking@1"],
workload_priority="benchmark",
)
)
with pytest.raises(ServingError, match="scope"):
service.authenticate(f"Bearer {created.credential}", "rag.embedding@1")
client = service.authenticate(f"Bearer {created.credential}", "rag.reranking@1")
response = service.invoke_reranking(
RerankingInvokeRequest(
query="private operator query",
documents=[
RerankDocument(id="chunk-a", text="private first document"),
RerankDocument(id="chunk-b", text="private second document"),
],
top_n=2,
),
client,
)
assert [item.id for item in response.results] == ["chunk-b", "chunk-a"]
assert all(math.isfinite(item.score) for item in response.results)
stored = session.scalar(
select(GatewayRequest).where(GatewayRequest.request_id == response.request_id)
)
assert stored is not None
serialized = str(stored.decision_evidence)
assert "private operator query" not in serialized
assert "private first document" not in serialized
assert stored.decision_evidence["content_logged"] is False
assert base.payload_store.values == {}
def test_reranking_request_rejects_duplicate_ids_and_oversized_pool() -> None:
with pytest.raises(ValueError, match="unique"):
RerankingInvokeRequest(
query="query",
documents=[
RerankDocument(id="same", text="one"),
RerankDocument(id="same", text="two"),
],
top_n=1,
)
with pytest.raises(ValueError):
RerankingInvokeRequest(
query="query",
documents=[RerankDocument(id=str(index), text="text") for index in range(41)],
top_n=10,
)
def test_expired_gpu_lease_is_reaped_without_capacity_leak(session: Session) -> None:
service, candidate, _node = setup_m5(session)
_approval, deployment = approve_and_promote(service, candidate.id)
lease = ServingGpuLease(
accelerator_id=deployment.accelerator_id,
capability_deployment_id=deployment.id,
request_id=uuid.uuid4(),
reserved_vram_bytes=100,
priority="production",
state="active",
owner="crashed-client",
created_at=datetime.now(UTC) - timedelta(minutes=2),
acquired_at=datetime.now(UTC) - timedelta(minutes=2),
expires_at=datetime.now(UTC) - timedelta(minutes=1),
)
session.add(lease)
session.commit()
service.reconcile()
assert lease.state == "expired"
assert service.scheduler_overview()[0].leased_vram_bytes == 0
def test_worker_restart_reconciles_failed_residency_to_cold(session: Session) -> None:
service, candidate, node = setup_m5(session)
_approval, deployment = approve_and_promote(service, candidate.id)
residency = service.repo.residency(deployment.id)
assert residency is not None
residency.state = "failed"
residency.worker_instance_id = "crashed-worker"
residency.measured_resident_vram_bytes = 1_479_868_416
residency.resident_since = datetime.now(UTC)
residency.failure_message = "previous worker failed"
session.commit()
service.report_state(
node,
AgentServingStateReport(
worker_instance_id="replacement-worker",
deployment_id=None,
state="cold",
load_count=0,
health={"process": "healthy", "runtime": "unavailable", "model": "unavailable"},
observed_at=datetime.now(UTC),
generation=99,
),
)
assert residency.state == "cold"
assert residency.measured_resident_vram_bytes == 0
assert residency.failure_code == "WORKER_RESTART_RECONCILED"
assert residency.generation == 99
assert residency.transition_reason == "WORKER_ABSENCE_RECONCILED"
assert residency.failure_message is None and residency.resident_since is None
def test_worker_residency_is_adopted_only_with_exact_immutable_identity(session: Session) -> None:
service, candidate, node = setup_m5(session)
_approval, deployment = approve_and_promote(service, candidate.id)
row = service.repo.deployment(deployment.id)
residency = service.repo.residency(deployment.id)
assert row is not None and residency is not None
profile = service.repo.profile(row.runtime_profile_id)
assert profile is not None and profile.runtime_environment_id is not None
environment = service.repo.environment(profile.runtime_environment_id)
assert environment is not None
report = AgentServingStateReport(
worker_instance_id="surviving-worker",
deployment_id=None,
state="warm",
load_count=1,
health={"process": "healthy"},
generation=42,
observed_at=datetime.now(UTC),
residencies=[
AgentResidentState(
deployment_id=row.id,
artifact_set_id=row.artifact_set_id,
runtime_profile_id=row.runtime_profile_id,
runtime_environment_fingerprint=environment.fingerprint,
revision_sha="a" * 40,
state="warm",
load_count=1,
resident_vram_bytes=1_479_868_416,
health={"runtime": "healthy", "model": "healthy"},
)
],
)
service.report_state(node, report)
assert residency.state == "warm"
assert residency.transition_reason == "WORKER_RESIDENCY_ADOPTED"
assert residency.generation == 42
wrong = report.model_copy(
update={
"residencies": [
report.residencies[0].model_copy(update={"artifact_set_id": uuid.uuid4()})
]
}
)
with pytest.raises(ServingError) as caught:
service.report_state(node, wrong)
assert caught.value.code == "RESIDENCY_CONFLICT"
def test_gateway_rate_limit_is_enforced_with_typed_failure(session: Session) -> None:
base, candidate, _node = setup_m5(session)
approve_and_promote(base, candidate.id)
service = InlineWorkerServingService(session, base.settings, base.manifests, base.payload_store)
created = service.create_client(
ServiceClientCreate(
name="rate-limited-client",
allowed_capabilities=["rag.embedding@1"],
requests_per_minute=1,
)
)
client = service.authenticate(f"Bearer {created.credential}", "rag.embedding@1")
service.invoke(EmbeddingInvokeRequest(input="first request"), client)
with pytest.raises(ServingError) as caught:
service.authenticate(f"Bearer {created.credential}", "rag.embedding@1")
assert caught.value.status_code == 429
assert caught.value.code == "RATE_LIMITED"
def test_gateway_rejects_full_bounded_queue(session: Session) -> None:
service, candidate, _node = setup_m5(session)
_approval, deployment = approve_and_promote(service, candidate.id)
stored_deployment = service.repo.deployment(deployment.id)
assert stored_deployment is not None
stored_deployment.max_queue_depth = 1
session.commit()
client = service.create_client(
ServiceClientCreate(
name="queue-client",
allowed_capabilities=["rag.embedding@1"],
)
)
payload_reference = service.payload_store.put({"input": ["queued"]}, 60)
service._enqueue_job(
stored_deployment,
"invoke",
payload_reference=payload_reference,
key_suffix="queue-capacity",
priority="background",
)
authenticated = service.authenticate(f"Bearer {client.credential}", "rag.embedding@1")
with pytest.raises(ServingError) as caught:
service.invoke(EmbeddingInvokeRequest(input="rejected"), authenticated)
assert caught.value.status_code == 429
assert caught.value.code == "QUEUE_FULL"
def test_scheduler_blocks_stale_resource_envelope(session: Session) -> None:
service, candidate, _node = setup_m5(session)
_approval, deployment = approve_and_promote(service, candidate.id)
envelope = service.repo.envelope(deployment.id)
assert envelope is not None
envelope.stale = True
envelope.stale_reason = "runtime image changed"
session.commit()
created = service.create_client(
ServiceClientCreate(
name="stale-envelope-client",
allowed_capabilities=["rag.embedding@1"],
)
)
client = service.authenticate(f"Bearer {created.credential}", "rag.embedding@1")
with pytest.raises(ServingError) as caught:
service.invoke(EmbeddingInvokeRequest(input="must not route"), client)
assert caught.value.status_code == 503
assert caught.value.code == "SCHEDULER_STATE_STALE"
stored = session.scalar(select(GatewayRequest).order_by(GatewayRequest.created_at.desc()))
assert stored is not None and stored.failure_code == "SCHEDULER_STATE_STALE"
def test_scheduler_rejects_request_when_external_usage_consumes_budget(
session: Session,
) -> None:
service, candidate, _node = setup_m5(session)
_approval, deployment = approve_and_promote(service, candidate.id)
telemetry = service.repo.telemetry(deployment.accelerator_id)
accelerator = service.repo.accelerator(deployment.accelerator_id)
assert telemetry is not None and accelerator is not None
assert accelerator.total_vram_bytes is not None
telemetry.used_vram_bytes = accelerator.total_vram_bytes
telemetry.free_vram_bytes = 0
session.commit()
created = service.create_client(
ServiceClientCreate(
name="capacity-client",
allowed_capabilities=["rag.embedding@1"],
)
)
client = service.authenticate(f"Bearer {created.credential}", "rag.embedding@1")
with pytest.raises(ServingError) as caught:
service.invoke(EmbeddingInvokeRequest(input="no overcommit"), client)
assert caught.value.status_code == 503
assert caught.value.code == "EXTERNAL_GPU_PRESSURE"
def test_serving_jobs_are_claimed_in_priority_order(session: Session) -> None:
service, candidate, node = setup_m5(session)
_approval, deployment = approve_and_promote(service, candidate.id)
background = service._enqueue_job(
deployment,
"health",
key_suffix="background-first",
priority="background",
)
production = service._enqueue_job(
deployment,
"health",
key_suffix="production-second",
priority="production",
)
lease = service.claim_next(node)
assert lease is not None and lease.job_id == production.id
assert session.get(ServingJob, background.id).status == "queued"
def test_serving_manifest_preserves_duplicate_content_paths(session: Session) -> None:
service, candidate, _node = setup_m5(session)
_approval, deployment = approve_and_promote(service, candidate.id)
artifact_set = service.repo.artifact_set(deployment.artifact_set_id)
assert artifact_set is not None
artifacts = service.repo.set_artifacts(artifact_set.id)
tokenizer = next(item for item in artifacts if item.filename == "tokenizer.json")
existing_location = service.repo.artifact_locations([tokenizer.id])[0]
duplicate_path = "tokenizer_config.json"
selected_paths = [*artifact_set.selected_paths, duplicate_path]
session.execute(
update(ArtifactSet)
.where(ArtifactSet.id == artifact_set.id)
.values(selected_paths=selected_paths, file_count=len(selected_paths))
)
session.add(
ArtifactLocation(
artifact_id=tokenizer.id,
storage_root_id=existing_location.storage_root_id,
relative_path=existing_location.relative_path.rsplit("/", 1)[0] + f"/{duplicate_path}",
status="verified",
size_bytes=tokenizer.size_bytes,
observed_sha256=tokenizer.sha256,
last_checked_at=datetime.now(UTC),
)
)
session.commit()
manifest = service._job_manifest(deployment)
assert [item["path"] for item in manifest["files"]] == selected_paths
assert len(manifest["files"]) == 6
def test_serving_job_long_poll_claims_job_without_worker_sleep(
session: Session, monkeypatch: pytest.MonkeyPatch
) -> None:
service, candidate, node = setup_m5(session)
_approval, deployment = approve_and_promote(service, candidate.id)
queued = service._enqueue_job(
deployment,
"health",
key_suffix="long-poll",
priority="production",
)
original = service.repo.claimable_job
calls = 0
def delayed_claim(node_id: uuid.UUID, now: datetime):
nonlocal calls
calls += 1
return None if calls == 1 else original(node_id, now)
monkeypatch.setattr(service.repo, "claimable_job", delayed_claim)
lease = service.claim_next(node, wait_seconds=0.1)
assert lease is not None and lease.job_id == queued.id
assert calls >= 2
def test_background_service_identity_cannot_escalate_scheduler_priority(
session: Session,
) -> None:
service, _candidate, _node = setup_m5(session)
created = service.create_client(
ServiceClientCreate(
name="examplerag-shadow-reindex",
allowed_capabilities=["rag.embedding@1"],
workload_priority="background",
)
)
assert created.workload_priority == "background"
stored = service.repo.client(created.id)
assert stored is not None and stored.workload_priority == "background"