import asyncio
import base64
import os
import time
from dataclasses import dataclass
import httpx
BASE = os.environ.get(
"CONTREE_BASE_URL", "https://api.tokenfactory.nebius.com/sandboxes"
).rstrip("/")
HEADERS = {
"Authorization": f"Bearer {os.environ['NEBIUS_API_KEY']}",
"Project": os.environ["NEBIUS_PROJECT_ID"],
}
BASE_IMAGE = os.environ["IMAGE_UUID"]
FIXTURES = {
"pass": "/fixture/evaluate.sh pass",
"test-fail": "/fixture/evaluate.sh fail",
"slow-cancel": "/fixture/evaluate.sh slow",
}
TERMINAL = {"SUCCESS", "FAILED", "CANCELLED"}
@dataclass
class Record:
item_id: str
input_image_id: str
operation_id: str | None = None
result_image_id: str | None = None
outcome: str | None = None
evidence: str = ""
async def main():
deadline = time.monotonic() + 90
limits = httpx.Limits(max_connections=3)
async with httpx.AsyncClient(timeout=20, limits=limits) as http:
preparation = await http.post(
f"{BASE}/v1/instances",
headers=HEADERS,
json={
"image": BASE_IMAGE,
"command": PREPARE_FIXTURE,
"shell": True,
"disposable": False,
"timeout": 30,
"truncate_output_at": 1048576,
},
)
preparation.raise_for_status()
preparation_id = preparation.json()["uuid"]
print(f"accepted preparation={preparation_id} input_image={BASE_IMAGE}")
prepared = await wait_operation(http, preparation_id, deadline)
if prepared["status"] != "SUCCESS":
raise RuntimeError('Preparation operation did not complete successfully')
if prepared["metadata"]["result"]["state"]["exit_code"] != 0:
raise RuntimeError('Preparation command failed')
image = prepared["result_image_uuid"]
if not image:
raise RuntimeError('Preparation returned no saved fixture image')
records = {item: Record(item, image) for item in FIXTURES}
owned_operation_ids = set()
try:
for item, command in FIXTURES.items():
response = await http.post(
f"{BASE}/v1/instances",
headers=HEADERS,
json={
"image": image,
"command": command,
"shell": True,
"disposable": True,
"timeout": 120,
"truncate_output_at": 1048576,
},
)
if response.status_code == 429:
records[item].outcome = "rate-or-quota-rejection"
continue
if response.status_code == 503:
records[item].outcome = "service-unavailable-rejection"
continue
response.raise_for_status()
operation_id = response.json()["uuid"]
records[item].operation_id = operation_id
owned_operation_ids.add(operation_id)
print(f"accepted item={item} operation={operation_id} image={image}")
observed_items = [
item for item in ("pass", "test-fail")
if records[item].operation_id is not None
]
completed = await asyncio.gather(
*(
wait_operation(http, records[item].operation_id, deadline)
for item in observed_items
)
)
for item, data in zip(observed_items, completed, strict=True):
records[item].outcome = classify(data)
records[item].result_image_id = data.get("result_image_uuid")
result = (data.get("metadata") or {}).get("result")
if result:
records[item].evidence = (
decoded(result.get("stdout") or {})
+ decoded(result.get("stderr") or {})
)[:80].replace("\n", " ")
slow = records["slow-cancel"]
if slow.operation_id in owned_operation_ids:
response = await http.delete(
f"{BASE}/v1/operations/{slow.operation_id}",
headers=HEADERS,
)
response.raise_for_status()
slow_data = await wait_operation(http, slow.operation_id, deadline)
slow.outcome = classify(slow_data)
slow.result_image_id = slow_data.get("result_image_uuid")
finally:
for record in records.values():
if record.operation_id in owned_operation_ids and record.outcome is None:
try:
response = await http.delete(
f"{BASE}/v1/operations/{record.operation_id}",
headers=HEADERS,
)
response.raise_for_status()
data = await wait_operation(
http,
record.operation_id,
time.monotonic() + 30,
)
record.outcome = classify(data)
record.result_image_id = data.get("result_image_uuid")
except (httpx.HTTPError, TimeoutError):
record.outcome = "cleanup-uncertain"
print("\nBatch summary:")
for record in records.values():
print(f" {record.item_id}: {record.outcome or 'unresolved'}")
print(f" operation_id: {record.operation_id}")
print(f" input_image_id: {record.input_image_id}")
print(f" result_image_id: {record.result_image_id}")
if record.evidence:
print(f" evidence: {record.evidence.strip()}")
if records["pass"].outcome != "passed":
raise RuntimeError('Passing case did not complete as passed')
if records["test-fail"].outcome != "test-failure":
raise RuntimeError('Failing case did not report a test failure')
if records["slow-cancel"].outcome != "cancelled":
raise RuntimeError('Slow case was not cancelled')
if any(record.result_image_id is not None for record in records.values()):
raise RuntimeError('Disposable evaluation returned an unexpected saved image')
if "case=pass" not in records["pass"].evidence:
raise RuntimeError('Passing case output marker was missing')
if "case=test-fail" not in records["test-fail"].evidence:
raise RuntimeError('Failing case output marker was missing')
async def wait_operation(http, operation_id, deadline):
while time.monotonic() < deadline:
response = await http.get(
f"{BASE}/v1/operations/{operation_id}",
headers=HEADERS,
)
response.raise_for_status()
data = response.json()
if data["status"] in TERMINAL:
return data
await asyncio.sleep(1)
raise TimeoutError(operation_id)
def classify(data):
if data["status"] == "CANCELLED":
return "cancelled"
if data["status"] == "FAILED":
return "platform-failure"
state = data["metadata"]["result"]["state"]
if state.get("timed_out"):
return "execution-timeout"
if state.get("signal", -1) > 0:
return "process-signal"
if state["exit_code"] != 0:
return "test-failure"
return "passed"
def decoded(stream):
value = stream.get("value", "")
if stream.get("encoding") == "base64":
return base64.b64decode(value).decode("utf-8", errors="replace")
return value
# This shell program runs inside the sandbox to prepare the retained image.
PREPARE_FIXTURE = r"""
set -e
mkdir -p /fixture
cat > /fixture/evaluate.sh <<'SH'
#!/bin/sh
case "$1" in
pass)
test -f /etc/os-release && printf 'case=pass\n'
;;
fail)
printf 'case=test-fail\n' >&2
test -f /fixture/missing
;;
slow)
printf 'case=slow\n'
while :; do sleep 1; done
;;
*) exit 64 ;;
esac
SH
chmod +x /fixture/evaluate.sh
"""
asyncio.run(main())