Files
model-adapt-strategy/main.py

95 lines
3.1 KiB
Python
Raw Permalink Normal View History

"""ModelHub XC Agent Strategy — HTTP service matching platform runtime contract."""
import json
import os
import signal
import traceback
from datetime import datetime
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
HOST = "0.0.0.0"
PORT = 8080
shutdown_requested = False
class Handler(BaseHTTPRequestHandler):
def do_GET(self) -> None:
if self.path == "/health":
self._send_json({"status": "ok"})
elif self.path == "/":
self._send_json({
"name": "p800-vllm-agent",
"status": "running",
"config": self._get_config(),
})
else:
self._send_json({"error": "not found"}, status=404)
def do_POST(self) -> None:
"""Handle task submission from platform."""
if self.path == "/task":
length = int(self.headers.get("Content-Length", 0))
body = self.rfile.read(length) if length else b"{}"
try:
task = json.loads(body)
result = self._handle_task(task)
self._send_json(result)
except Exception:
self._send_json({"error": traceback.format_exc()}, status=500)
else:
self._send_json({"error": "not found"}, status=404)
def _handle_task(self, task: dict) -> dict:
"""Process a model adaptation task."""
model = task.get("model_address", task.get("modelAddress", ""))
config = task.get("config_params", task.get("configParams", {}))
print(f"[{datetime.now()}] Task received: model={model}", flush=True)
return {
"status": "accepted",
"model": model,
"message": "Task received by agent",
}
def _get_config(self) -> dict:
token = os.getenv("EXTERNAL_SERVICE_TOKEN", "")
return {
"external_service_token_present": bool(token),
"model_address": os.getenv("MODEL_ADDRESS", ""),
"config_params": os.getenv("CONFIG_PARAMS", ""),
}
def log_message(self, fmt: str, *args: object) -> None:
print(f"{self.address_string()} - {fmt % args}", flush=True)
def _send_json(self, body: dict, status: int = 200) -> None:
payload = json.dumps(body, ensure_ascii=False).encode()
self.send_response(status)
self.send_header("Content-Type", "application/json")
self.send_header("Content-Length", str(len(payload)))
self.end_headers()
self.wfile.write(payload)
def _handle_signal(signum: int, _frame: object) -> None:
global shutdown_requested
shutdown_requested = True
print(f"Received signal {signum}, shutting down", flush=True)
def main() -> None:
signal.signal(signal.SIGTERM, _handle_signal)
signal.signal(signal.SIGINT, _handle_signal)
server = ThreadingHTTPServer((HOST, PORT), Handler)
server.timeout = 1
print(f"Agent listening on {HOST}:{PORT}", flush=True)
while not shutdown_requested:
server.handle_request()
server.server_close()
print("Agent stopped", flush=True)
if __name__ == "__main__":
main()