1209 lines
46 KiB
Python
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"
|