#!/usr/bin/env python3
import argparse
import base64
import json
import os
import time
import urllib.error
import urllib.parse
import urllib.request
def main():
parser = argparse.ArgumentParser()
parser.add_argument("--deadline", type=float, default=120.0)
args = parser.parse_args()
client = Client(args.deadline)
try:
normal_interaction(client)
child_signal(client)
parent_finishes_first(client)
print("Completed input, signal, and parent-shutdown examples.")
finally:
client.cancel_owned()
def normal_interaction(client):
print("\n--- two writes and EOF ---")
parent = client.create_parent(WAIT_FOR_CHILD)
client.wait_ready(parent)
child = client.spawn_child(
parent,
{
"command": READ_TWO_LINES,
"shell": True,
"stdin": {"value": "", "encoding": "ascii", "close": False},
"truncate_output_at": 1048576,
},
)
client.write_stdin(parent, child, "first line\n", close=False)
client.write_stdin(parent, child, "second line\n", close=True)
child_result = client.wait_child(parent, child)
child_state = child_result["state"]
if child_state["exit_code"] != 0 or child_state["signal"] != -1:
raise RuntimeError('Child did not exit successfully')
received = decode_stream(child_result["stdout"])
if received != "received:first line|second line\n":
raise RuntimeError('Child did not return both input lines in order')
if decode_stream(child_result["stderr"]) != "":
raise RuntimeError('Child wrote unexpected stderr output')
print(received, end="")
parent_result = client.wait_parent(parent)
assert_parent_success(parent_result)
print(f"child={child} exit=0; parent={parent_result['status']}")
def child_signal(client):
print("\n--- signal child; parent continues ---")
parent = client.create_parent(WAIT_FOR_RELEASE)
client.wait_ready(parent)
child = client.spawn_child(parent, {"command": "sleep 30", "shell": True})
client.signal_child(parent, child)
child_result = client.wait_child(parent, child)
if child_result["state"]["signal"] != 15:
raise RuntimeError('Child result did not report SIGTERM')
if child_result["state"]["exit_code"] != -1:
raise RuntimeError('Signalled child returned an unexpected exit code')
release_parent(client, parent)
parent_result = client.wait_parent(parent)
assert_parent_success(parent_result)
print(f"child={child} signal=15; parent={parent_result['status']}")
def parent_finishes_first(client):
print("\n--- parent finishes first ---")
parent = client.create_parent(WAIT_FOR_RELEASE)
client.wait_ready(parent)
child = client.spawn_child(parent, {"command": LONG_LIVED_CHILD, "shell": True})
release_parent(client, parent, wait_for_child=True)
parent_result = client.wait_parent(parent)
assert_parent_success(parent_result)
child_result = client.wait_child(parent, child)
child_state = child_result["state"]
if child_state.get("exit_code") == 0 and child_state.get("signal") in (-1, None):
raise RuntimeError("Child unexpectedly survived parent finalization")
print(
f"parent={parent_result['status']}; child exit={child_state.get('exit_code')} "
f"signal={child_state.get('signal')}"
)
def release_parent(client, parent, wait_for_child=False):
command = RELEASE_AFTER_CHILD_STARTS if wait_for_child else RELEASE_PARENT
release = client.spawn_child(parent, {"command": command, "shell": True})
result = client.wait_child(parent, release)
state = result["state"]
if state.get("exit_code") != 0 or state.get("signal") not in (-1, None):
raise RuntimeError(f"Could not release the parent: {state}")
# Each parent waits for a file condition with a bounded remote lifetime.
WAIT_FOR_CHILD = r"""
printf 'parent-ready\n'
i=0
while [ ! -f /tmp/child-complete ] && [ "$i" -lt 20 ]; do
i=$((i + 1))
sleep 1
done
test -f /tmp/child-complete
"""
READ_TWO_LINES = r"""
read first
read second
printf 'received:%s|%s\n' "$first" "$second"
touch /tmp/child-complete
"""
WAIT_FOR_RELEASE = r"""
printf 'parent-ready\n'
i=0
while [ ! -f /tmp/parent-release ] && [ "$i" -lt 20 ]; do
i=$((i + 1))
sleep 1
done
test -f /tmp/parent-release
"""
LONG_LIVED_CHILD = r"""
touch /tmp/child-started
sleep 30
"""
RELEASE_PARENT = "touch /tmp/parent-release"
RELEASE_AFTER_CHILD_STARTS = r"""
i=0
while [ ! -f /tmp/child-started ] && [ "$i" -lt 10 ]; do
i=$((i + 1))
sleep 1
done
test -f /tmp/child-started || exit 1
touch /tmp/parent-release
"""
TERMINAL = {"SUCCESS", "FAILED", "CANCELLED"}
def required(name):
value = os.environ.get(name)
if not value:
raise SystemExit(f"Set {name}; see https://docs.tokenfactory.nebius.com/sandboxes/start/set-up-access")
return value
def retry_delay(headers):
try:
return max(0.1, float(headers.get("Retry-After", "1")))
except ValueError:
return 1.0
def decode_stream(stream):
value = (stream or {}).get("value", "")
encoding = (stream or {}).get("encoding", "ascii")
if encoding == "base64":
return base64.b64decode(value).decode("utf-8", errors="replace")
if encoding == "ascii":
return value
raise RuntimeError(f"Unsupported stream encoding: {encoding}")
class Client:
def __init__(self, deadline_seconds):
self.base = os.environ.get(
"CONTREE_BASE_URL", "https://api.tokenfactory.nebius.com/sandboxes"
).rstrip("/")
self.token = required("NEBIUS_API_KEY")
self.project = required("NEBIUS_PROJECT_ID")
self.image = required("IMAGE_UUID")
self.deadline = time.monotonic() + deadline_seconds
self.owned = set()
def remaining(self):
value = self.deadline - time.monotonic()
if value <= 0:
raise TimeoutError("Recipe deadline exceeded")
return min(15.0, value)
def request(self, method, path, body=None, allowed=()):
headers = {
"Authorization": f"Bearer {self.token}",
"Project": self.project,
"Accept": "application/json",
}
data = None
if body is not None:
data = json.dumps(body).encode()
headers["Content-Type"] = "application/json"
request = urllib.request.Request(
self.base + "/v1" + path, data=data, headers=headers, method=method
)
try:
with urllib.request.urlopen(request, timeout=self.remaining()) as response:
payload = response.read()
decoded = json.loads(payload) if payload else None
return response.status, response.headers, decoded
except urllib.error.HTTPError as error:
payload = error.read()
decoded = json.loads(payload) if payload else None
if error.code in allowed:
return error.code, error.headers, decoded
raise RuntimeError(f"HTTP {error.code} for {method} {path}: {decoded}") from error
def create_parent(self, command):
status, _, result = self.request(
"POST",
"/instances",
{
"image": self.image,
"command": command,
"shell": True,
"disposable": True,
"timeout": 30,
"truncate_output_at": 1048576,
},
)
if status != 201:
raise RuntimeError(f"Expected parent create status 201, got {status}")
operation_id = result["uuid"]
self.owned.add(operation_id)
return operation_id
def operation(self, operation_id, inflight=False):
suffix = "?inflight=1" if inflight else ""
return self.request("GET", f"/operations/{operation_id}{suffix}")
def wait_ready(self, operation_id, marker="parent-ready"):
while True:
_, headers, result = self.operation(operation_id, inflight=True)
status = result["status"]
process = (result.get("metadata") or {}).get("result") or {}
if status in {"ASSIGNED", "EXECUTING"} and marker in decode_stream(
process.get("stdout")
):
return
if status in TERMINAL:
raise RuntimeError(f"Parent ended before readiness: {result}")
time.sleep(min(retry_delay(headers), self.remaining()))
def wait_parent(self, operation_id):
while True:
_, headers, result = self.operation(operation_id)
if result["status"] in TERMINAL:
self.owned.discard(operation_id)
return result
time.sleep(min(retry_delay(headers), self.remaining()))
def spawn_child(self, operation_id, body):
path = f"/operations/{operation_id}/subprocesses"
while True:
status, headers, result = self.request("POST", path, body, allowed=(425, 504))
if status == 201:
return result["spid"]
if status == 425:
time.sleep(min(retry_delay(headers), self.remaining()))
continue
raise RuntimeError(
"Child spawn timed out with uncertain outcome; inspect the event stream before retrying"
)
def write_stdin(self, operation_id, spid, value, close):
status, _, _ = self.request(
"POST",
f"/operations/{operation_id}/subprocesses/{spid}/stdin",
{"value": value, "encoding": "ascii", "close": close},
allowed=(504,),
)
if status == 504:
raise RuntimeError("Stdin delivery is uncertain; do not retry this chunk blindly")
if status != 204:
raise RuntimeError(f"Expected stdin status 204, got {status}")
def signal_child(self, operation_id, spid, signal="SIGTERM"):
query = urllib.parse.urlencode({"signal": signal})
status, _, _ = self.request(
"DELETE", f"/operations/{operation_id}/subprocesses/{spid}?{query}"
)
if status != 204:
raise RuntimeError(f"Expected signal status 204, got {status}")
def wait_child(self, operation_id, spid):
path = f"/operations/{operation_id}/subprocesses/{spid}"
while True:
status, headers, result = self.request("GET", path, allowed=(425,))
if status == 200:
return result
time.sleep(min(retry_delay(headers), self.remaining()))
def cancel_owned(self):
self.deadline = max(self.deadline, time.monotonic() + 15.0)
for operation_id in list(self.owned):
try:
_, _, current = self.operation(operation_id)
if current["status"] not in TERMINAL:
self.request("DELETE", f"/operations/{operation_id}")
self.wait_parent(operation_id)
except Exception as error:
print(f"Cleanup failed for {operation_id}: {error}")
def assert_parent_success(result):
process = (result.get("metadata") or {}).get("result") or {}
state = process.get("state") or {}
if result.get("status") != "SUCCESS" or state.get("exit_code") != 0:
raise RuntimeError(f"Parent did not finish successfully: {result}")
if state.get("signal") not in (-1, None) or state.get("timed_out"):
raise RuntimeError(f"Parent was interrupted: {state}")
if __name__ == "__main__":
main()