This commit is contained in:
+502
-58
@@ -1,6 +1,7 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import asyncio
|
||||
import json
|
||||
import shutil
|
||||
import subprocess
|
||||
@@ -26,6 +27,7 @@ from traderai.memory import DEFAULT_THREAD_ID, MemoryStore
|
||||
from traderai.plans import ContinualPlanRunner, ContinualPlanStore
|
||||
from traderai.scheduler import WakeScheduler
|
||||
from traderai.scmdb_client import SCMDBClient
|
||||
from traderai.starcitizen_wiki_client import StarCitizenWikiClient
|
||||
from traderai.tools import ToolRegistry
|
||||
from traderai.uex_client import UEXClient
|
||||
from traderai.version import RELEASES_API_URL, RELEASES_URL, __version__
|
||||
@@ -106,34 +108,52 @@ def create_app() -> FastAPI:
|
||||
memory = MemoryStore(settings.traderai_memory_path)
|
||||
plan_store = ContinualPlanStore(memory)
|
||||
scheduler = WakeScheduler(memory)
|
||||
uex = UEXClient(settings.uex_base_url, settings.uex_secret_key, settings.uex_bearer_token)
|
||||
scmdb = SCMDBClient(settings.scmdb_base_url)
|
||||
cornerstone = CornerstoneClient(settings.cornerstone_base_url)
|
||||
tools = ToolRegistry(
|
||||
uex,
|
||||
settings.require_write_approval,
|
||||
memory=memory,
|
||||
scheduler=scheduler,
|
||||
scmdb=scmdb,
|
||||
cornerstone=cornerstone,
|
||||
plan_store=plan_store,
|
||||
)
|
||||
plan_runner = ContinualPlanRunner(plan_store, tools, memory)
|
||||
tools.plan_runner = plan_runner
|
||||
agent = OllamaAgent(
|
||||
settings.openai_base_url if settings.model_provider == "openai" else settings.ollama_base_url,
|
||||
settings.openai_model if settings.model_provider == "openai" else settings.ollama_model,
|
||||
tools,
|
||||
memory=memory,
|
||||
user_name=settings.traderai_user_name,
|
||||
num_ctx=settings.ollama_num_ctx,
|
||||
provider=settings.model_provider,
|
||||
api_key=settings.openai_api_key,
|
||||
)
|
||||
plan_runner.bind_agent(agent)
|
||||
scheduler.bind_agent(agent)
|
||||
scheduler.bind_plan_runner(plan_runner)
|
||||
scheduler.bind_uex_notifications(uex, settings.uex_notification_poll_seconds)
|
||||
runtime: dict[str, Any] = {}
|
||||
|
||||
def configure_runtime(current_settings: Any) -> None:
|
||||
uex = UEXClient(current_settings.uex_base_url, current_settings.uex_secret_key, current_settings.uex_bearer_token)
|
||||
scmdb = SCMDBClient(current_settings.scmdb_base_url)
|
||||
cornerstone = CornerstoneClient(current_settings.cornerstone_base_url)
|
||||
scwiki = StarCitizenWikiClient(current_settings.scwiki_base_url, current_settings.scwiki_api_base_url)
|
||||
tools = ToolRegistry(
|
||||
uex,
|
||||
current_settings.require_write_approval,
|
||||
memory=memory,
|
||||
scheduler=scheduler,
|
||||
scmdb=scmdb,
|
||||
cornerstone=cornerstone,
|
||||
scwiki=scwiki,
|
||||
plan_store=plan_store,
|
||||
)
|
||||
plan_runner = ContinualPlanRunner(plan_store, tools, memory)
|
||||
tools.plan_runner = plan_runner
|
||||
provider_base_url, provider_model, provider_api_key = provider_settings(current_settings)
|
||||
agent = OllamaAgent(
|
||||
provider_base_url,
|
||||
provider_model,
|
||||
tools,
|
||||
memory=memory,
|
||||
user_name=current_settings.traderai_user_name,
|
||||
num_ctx=current_settings.ollama_num_ctx,
|
||||
provider=current_settings.model_provider,
|
||||
api_key=provider_api_key,
|
||||
reasoning_effort=current_settings.model_reasoning_effort,
|
||||
)
|
||||
plan_runner.bind_agent(agent)
|
||||
scheduler.bind_agent(agent)
|
||||
scheduler.bind_plan_runner(plan_runner)
|
||||
scheduler.bind_uex_notifications(uex, current_settings.uex_notification_poll_seconds)
|
||||
runtime.update(
|
||||
{
|
||||
"settings": current_settings,
|
||||
"uex": uex,
|
||||
"tools": tools,
|
||||
"plan_runner": plan_runner,
|
||||
"agent": agent,
|
||||
}
|
||||
)
|
||||
|
||||
configure_runtime(settings)
|
||||
|
||||
app = FastAPI(title="TraderAI")
|
||||
static_dir = resource_path("web")
|
||||
@@ -149,17 +169,20 @@ def create_app() -> FastAPI:
|
||||
scheduler.shutdown()
|
||||
|
||||
async def refresh_user_profile() -> None:
|
||||
if settings.traderai_user_name:
|
||||
memory.set_profile("configured_name", settings.traderai_user_name)
|
||||
agent.user_name = agent.user_name or settings.traderai_user_name
|
||||
current_settings = get_settings()
|
||||
agent = runtime["agent"]
|
||||
uex = runtime["uex"]
|
||||
if current_settings.traderai_user_name:
|
||||
memory.set_profile("configured_name", current_settings.traderai_user_name)
|
||||
agent.user_name = agent.user_name or current_settings.traderai_user_name
|
||||
|
||||
try:
|
||||
response = await uex.get_user(authenticated=True)
|
||||
except Exception as exc:
|
||||
memory.set_profile("uex_user_error", str(exc))
|
||||
if settings.traderai_user_name:
|
||||
if current_settings.traderai_user_name:
|
||||
try:
|
||||
response = await uex.get_user(username=settings.traderai_user_name)
|
||||
response = await uex.get_user(username=current_settings.traderai_user_name)
|
||||
except Exception:
|
||||
return
|
||||
else:
|
||||
@@ -178,9 +201,13 @@ def create_app() -> FastAPI:
|
||||
|
||||
@app.get("/api/health")
|
||||
async def health() -> dict:
|
||||
agent = runtime["agent"]
|
||||
current_settings = get_settings()
|
||||
inference = await agent.health()
|
||||
return {
|
||||
"ollama": await agent.health(),
|
||||
"model_provider": settings.model_provider,
|
||||
"inference": inference,
|
||||
"ollama": inference,
|
||||
"model_provider": current_settings.model_provider,
|
||||
"user": memory.get_profile(),
|
||||
"jobs": scheduler.list_jobs(),
|
||||
"app_data_dir": settings_payload()["app_data_dir"],
|
||||
@@ -193,27 +220,62 @@ def create_app() -> FastAPI:
|
||||
|
||||
@app.post("/api/config")
|
||||
async def update_config(request: ConfigUpdateRequest) -> dict:
|
||||
previous_settings = get_settings()
|
||||
updated = save_settings(request.values)
|
||||
updated["restart_required"] = True
|
||||
updated["message"] = "Configuration saved. Restart TraderAI for all settings to take effect."
|
||||
current_settings = get_settings()
|
||||
configure_runtime(current_settings)
|
||||
await refresh_user_profile()
|
||||
restart_required = (
|
||||
"traderai_memory_path" in request.values
|
||||
and str(request.values.get("traderai_memory_path") or "").strip() != str(previous_settings.traderai_memory_path)
|
||||
)
|
||||
updated["restart_required"] = restart_required
|
||||
updated["message"] = (
|
||||
"Configuration saved. Restart TraderAI to switch memory databases."
|
||||
if restart_required
|
||||
else "Configuration saved and applied."
|
||||
)
|
||||
return updated
|
||||
|
||||
@app.get("/api/ollama/status")
|
||||
async def ollama_status() -> dict:
|
||||
return await inspect_model_provider()
|
||||
|
||||
@app.get("/api/openai/models")
|
||||
async def openai_models() -> dict:
|
||||
status = await inspect_openai()
|
||||
@app.get("/api/provider/models")
|
||||
async def provider_models(provider: str | None = None) -> dict:
|
||||
status = await inspect_provider_models(provider)
|
||||
return {
|
||||
"provider": "openai",
|
||||
"provider": status.get("provider", "openai"),
|
||||
"configured_model": status.get("configured_model"),
|
||||
"models": status.get("models", []),
|
||||
"reasoning_efforts": status.get("reasoning_efforts", reasoning_effort_options()),
|
||||
"configured_reasoning_effort": status.get("configured_reasoning_effort", get_settings().model_reasoning_effort),
|
||||
"message": status.get("message", ""),
|
||||
"detail": status.get("detail", ""),
|
||||
"online": status.get("online", False),
|
||||
}
|
||||
|
||||
@app.post("/api/codex/login")
|
||||
async def launch_codex_login() -> dict:
|
||||
current_settings = get_settings()
|
||||
command = find_codex_cli(current_settings.codex_command)
|
||||
if not command:
|
||||
raise HTTPException(status_code=404, detail="Codex CLI was not found on PATH.")
|
||||
try:
|
||||
login = await start_codex_browser_login(command)
|
||||
except Exception as exc:
|
||||
raise HTTPException(status_code=500, detail=f"Codex App Server login failed: {exception_detail(exc)}") from exc
|
||||
return {
|
||||
"installed": True,
|
||||
"running": False,
|
||||
"online": False,
|
||||
"provider": "codex",
|
||||
"login_id": login.get("loginId"),
|
||||
"auth_url": login.get("authUrl"),
|
||||
"base_url": str(command),
|
||||
"message": "Opened Codex App Server sign-in in your browser. Finish the flow, then TraderAI will detect the new login.",
|
||||
}
|
||||
|
||||
@app.post("/api/ollama/launch")
|
||||
async def launch_ollama() -> dict:
|
||||
command = ollama_launch_command()
|
||||
@@ -319,6 +381,7 @@ def create_app() -> FastAPI:
|
||||
|
||||
@app.post("/api/chat")
|
||||
async def chat(request: ChatRequest) -> dict:
|
||||
agent = runtime["agent"]
|
||||
try:
|
||||
return await agent.chat(
|
||||
request.message,
|
||||
@@ -330,6 +393,8 @@ def create_app() -> FastAPI:
|
||||
|
||||
@app.post("/api/chat/stream")
|
||||
async def chat_stream(request: ChatRequest) -> StreamingResponse:
|
||||
agent = runtime["agent"]
|
||||
|
||||
async def events():
|
||||
async for event in agent.chat_events(
|
||||
request.message,
|
||||
@@ -367,6 +432,7 @@ def create_app() -> FastAPI:
|
||||
|
||||
@app.get("/api/pending-actions")
|
||||
async def pending_actions() -> dict:
|
||||
agent = runtime["agent"]
|
||||
return {"pending_actions": agent._pending_payloads()}
|
||||
|
||||
@app.get("/api/notifications")
|
||||
@@ -393,11 +459,13 @@ def create_app() -> FastAPI:
|
||||
|
||||
@app.get("/api/negotiations/{identifier}/messages")
|
||||
async def negotiation_messages(identifier: str) -> dict:
|
||||
uex = runtime["uex"]
|
||||
params = negotiation_identifier_params(identifier)
|
||||
return await uex.get("marketplace_negotiations_messages", params, authenticated=True)
|
||||
|
||||
@app.post("/api/negotiations/{identifier}/messages")
|
||||
async def send_negotiation_message(identifier: str, request: DirectNegotiationMessageRequest) -> dict:
|
||||
uex = runtime["uex"]
|
||||
params = negotiation_identifier_params(identifier)
|
||||
payload = {**params, "message": request.message, "is_production": 1}
|
||||
return await uex.post("marketplace_negotiations_messages", payload, authenticated=True)
|
||||
@@ -412,6 +480,7 @@ def create_app() -> FastAPI:
|
||||
|
||||
@app.post("/api/plans")
|
||||
async def create_continual_plan(request: ContinualPlanCreateRequest) -> dict:
|
||||
tools = runtime["tools"]
|
||||
result = await tools.create_continual_plan(
|
||||
title=request.title,
|
||||
objective=request.objective,
|
||||
@@ -433,6 +502,7 @@ def create_app() -> FastAPI:
|
||||
|
||||
@app.post("/api/plans/{plan_id}/pause")
|
||||
async def pause_continual_plan(plan_id: str) -> dict:
|
||||
tools = runtime["tools"]
|
||||
result = await tools.pause_continual_plan(plan_id)
|
||||
if result.get("error"):
|
||||
raise HTTPException(status_code=404, detail=result["error"])
|
||||
@@ -440,6 +510,7 @@ def create_app() -> FastAPI:
|
||||
|
||||
@app.post("/api/plans/{plan_id}/resume")
|
||||
async def resume_continual_plan(plan_id: str) -> dict:
|
||||
tools = runtime["tools"]
|
||||
result = await tools.resume_continual_plan(plan_id)
|
||||
if result.get("error"):
|
||||
raise HTTPException(status_code=404, detail=result["error"])
|
||||
@@ -447,13 +518,23 @@ def create_app() -> FastAPI:
|
||||
|
||||
@app.post("/api/plans/{plan_id}/cancel")
|
||||
async def cancel_continual_plan(plan_id: str) -> dict:
|
||||
tools = runtime["tools"]
|
||||
result = await tools.cancel_continual_plan(plan_id)
|
||||
if result.get("error"):
|
||||
raise HTTPException(status_code=404, detail=result["error"])
|
||||
return result
|
||||
|
||||
@app.delete("/api/plans/{plan_id}")
|
||||
async def delete_continual_plan(plan_id: str) -> dict:
|
||||
tools = runtime["tools"]
|
||||
result = await tools.delete_continual_plan(plan_id)
|
||||
if result.get("error"):
|
||||
raise HTTPException(status_code=404, detail=result["error"])
|
||||
return result
|
||||
|
||||
@app.post("/api/plans/{plan_id}/run")
|
||||
async def run_continual_plan(plan_id: str) -> dict:
|
||||
tools = runtime["tools"]
|
||||
result = await tools.run_continual_plan_now(plan_id)
|
||||
if result.get("error"):
|
||||
raise HTTPException(status_code=400, detail=result["error"])
|
||||
@@ -487,10 +568,12 @@ def create_app() -> FastAPI:
|
||||
|
||||
@app.post("/api/approve/{action_id}")
|
||||
async def approve(action_id: str) -> dict:
|
||||
tools = runtime["tools"]
|
||||
return await tools.approve(action_id)
|
||||
|
||||
@app.post("/api/decline/{action_id}")
|
||||
async def decline(action_id: str) -> dict:
|
||||
tools = runtime["tools"]
|
||||
return await tools.decline(action_id)
|
||||
|
||||
return app
|
||||
@@ -509,33 +592,96 @@ async def inspect_model_provider() -> dict[str, Any]:
|
||||
settings = get_settings()
|
||||
if settings.model_provider == "openai":
|
||||
return await inspect_openai()
|
||||
if settings.model_provider == "codex":
|
||||
return await inspect_codex()
|
||||
return await inspect_ollama()
|
||||
|
||||
|
||||
async def inspect_openai() -> dict[str, Any]:
|
||||
settings = get_settings()
|
||||
return await inspect_cloud_provider_config("openai", settings.openai_base_url, settings.openai_api_key, settings.openai_model)
|
||||
|
||||
|
||||
async def inspect_codex() -> dict[str, Any]:
|
||||
settings = get_settings()
|
||||
command = find_codex_cli(settings.codex_command)
|
||||
detail = ""
|
||||
online = False
|
||||
models: list[str] = []
|
||||
effort_map: dict[str, list[str]] = {}
|
||||
if command:
|
||||
try:
|
||||
account, models, effort_map = await inspect_codex_app_server(command)
|
||||
online = bool(account)
|
||||
detail = f"Logged in as {account.get('email')}" if isinstance(account, dict) and account.get("email") else ""
|
||||
except (OSError, RuntimeError, asyncio.TimeoutError) as exc:
|
||||
detail = str(exc)
|
||||
configured_model = settings.codex_model
|
||||
model_available = configured_model in models if models else bool(configured_model)
|
||||
return {
|
||||
"installed": bool(command),
|
||||
"running": online,
|
||||
"online": online,
|
||||
"provider": "codex",
|
||||
"model_available": model_available,
|
||||
"configured_model": configured_model,
|
||||
"configured_reasoning_effort": settings.model_reasoning_effort,
|
||||
"reasoning_efforts": codex_reasoning_efforts(configured_model, effort_map),
|
||||
"base_url": str(command) if command else settings.codex_command,
|
||||
"models": models,
|
||||
"message": codex_status_message(bool(command), online, model_available, configured_model),
|
||||
"detail": detail,
|
||||
}
|
||||
|
||||
|
||||
async def inspect_cloud_provider() -> dict[str, Any]:
|
||||
settings = get_settings()
|
||||
if settings.model_provider == "codex":
|
||||
return await inspect_codex()
|
||||
return await inspect_openai()
|
||||
|
||||
|
||||
async def inspect_provider_models(provider: str | None = None) -> dict[str, Any]:
|
||||
normalized = str(provider or get_settings().model_provider).strip().casefold()
|
||||
if normalized == "codex":
|
||||
return await inspect_codex()
|
||||
if normalized == "ollama":
|
||||
return await inspect_ollama()
|
||||
return await inspect_openai()
|
||||
|
||||
|
||||
async def inspect_cloud_provider_config(
|
||||
provider: str,
|
||||
base_url: str,
|
||||
api_key: str | None,
|
||||
model: str,
|
||||
) -> dict[str, Any]:
|
||||
settings = get_settings()
|
||||
models: list[str] = []
|
||||
online = False
|
||||
detail = ""
|
||||
if not settings.openai_api_key:
|
||||
provider_name = provider_display_name(provider)
|
||||
if not api_key:
|
||||
return {
|
||||
"installed": True,
|
||||
"running": False,
|
||||
"online": False,
|
||||
"provider": "openai",
|
||||
"provider": provider,
|
||||
"model_available": False,
|
||||
"configured_model": settings.openai_model,
|
||||
"base_url": settings.openai_base_url,
|
||||
"configured_model": model,
|
||||
"configured_reasoning_effort": settings.model_reasoning_effort,
|
||||
"reasoning_efforts": reasoning_effort_options(),
|
||||
"base_url": base_url,
|
||||
"models": [],
|
||||
"message": "OpenAI is selected, but no API key is configured.",
|
||||
"message": f"{provider_name} is selected, but no API key is configured.",
|
||||
"detail": "",
|
||||
}
|
||||
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=10) as client:
|
||||
response = await client.get(
|
||||
f"{settings.openai_base_url.rstrip('/')}/models",
|
||||
headers={"Authorization": f"Bearer {settings.openai_api_key}"},
|
||||
f"{base_url.rstrip('/')}/models",
|
||||
headers={"Authorization": f"Bearer {api_key}"},
|
||||
)
|
||||
response.raise_for_status()
|
||||
body = response.json()
|
||||
@@ -544,17 +690,19 @@ async def inspect_openai() -> dict[str, Any]:
|
||||
except (httpx.HTTPError, ValueError) as exc:
|
||||
detail = str(exc)
|
||||
|
||||
model_available = settings.openai_model in models
|
||||
model_available = model in models
|
||||
return {
|
||||
"installed": True,
|
||||
"running": online,
|
||||
"online": online,
|
||||
"provider": "openai",
|
||||
"provider": provider,
|
||||
"model_available": model_available,
|
||||
"configured_model": settings.openai_model,
|
||||
"base_url": settings.openai_base_url,
|
||||
"configured_model": model,
|
||||
"configured_reasoning_effort": settings.model_reasoning_effort,
|
||||
"reasoning_efforts": reasoning_effort_options(),
|
||||
"base_url": base_url,
|
||||
"models": models,
|
||||
"message": openai_status_message(online, bool(settings.openai_api_key), model_available, settings.openai_model),
|
||||
"message": cloud_status_message(provider, online, bool(api_key), model_available, model),
|
||||
"detail": detail,
|
||||
}
|
||||
|
||||
@@ -587,6 +735,8 @@ async def inspect_ollama() -> dict[str, Any]:
|
||||
"provider": "ollama",
|
||||
"model_available": model_available,
|
||||
"configured_model": settings.ollama_model,
|
||||
"configured_reasoning_effort": settings.model_reasoning_effort,
|
||||
"reasoning_efforts": reasoning_effort_options(),
|
||||
"base_url": settings.ollama_base_url,
|
||||
"num_ctx": settings.ollama_num_ctx,
|
||||
"models": models,
|
||||
@@ -599,14 +749,15 @@ async def inspect_ollama() -> dict[str, Any]:
|
||||
}
|
||||
|
||||
|
||||
def openai_status_message(running: bool, configured: bool, model_available: bool, model: str) -> str:
|
||||
def cloud_status_message(provider: str, running: bool, configured: bool, model_available: bool, model: str) -> str:
|
||||
provider_name = provider_display_name(provider)
|
||||
if not configured:
|
||||
return "OpenAI API key is not configured."
|
||||
return f"{provider_name} API key is not configured."
|
||||
if not running:
|
||||
return "OpenAI is not reachable with the configured key."
|
||||
return f"{provider_name} is not reachable with the configured key."
|
||||
if not model_available:
|
||||
return f'OpenAI is reachable, but model "{model}" was not returned by the API.'
|
||||
return "OpenAI is ready."
|
||||
return f'{provider_name} is reachable, but model "{model}" was not returned by the API.'
|
||||
return f"{provider_name} is ready."
|
||||
|
||||
|
||||
def ollama_status_message(installed: bool, running: bool, model_available: bool, model: str) -> str:
|
||||
@@ -619,6 +770,292 @@ def ollama_status_message(installed: bool, running: bool, model_available: bool,
|
||||
return "Ollama is ready."
|
||||
|
||||
|
||||
def codex_status_message(installed: bool, logged_in: bool, model_available: bool, model: str) -> str:
|
||||
if not installed:
|
||||
return "Codex CLI is not installed."
|
||||
if not logged_in:
|
||||
return "Codex CLI is installed, but the Codex App Server is not logged in with ChatGPT."
|
||||
if not model_available:
|
||||
return f'Codex App Server is logged in, but model "{model}" was not returned by the model list.'
|
||||
return "Codex App Server is ready."
|
||||
|
||||
|
||||
def provider_settings(settings: Any) -> tuple[str, str, str | None]:
|
||||
if settings.model_provider == "openai":
|
||||
return settings.openai_base_url, settings.openai_model, settings.openai_api_key
|
||||
if settings.model_provider == "codex":
|
||||
return settings.codex_command, settings.codex_model, None
|
||||
return settings.ollama_base_url, settings.ollama_model, None
|
||||
|
||||
|
||||
def provider_display_name(provider: str) -> str:
|
||||
return {"openai": "OpenAI", "codex": "Codex"}.get(provider, "Ollama")
|
||||
|
||||
|
||||
def find_codex_cli(configured_command: str | None = None) -> Path | None:
|
||||
candidates = [configured_command, shutil.which("codex"), os.path.join(os.environ.get("USERPROFILE", ""), ".codex", ".sandbox-bin", "codex.exe")]
|
||||
for candidate in candidates:
|
||||
if not candidate:
|
||||
continue
|
||||
resolved = shutil.which(candidate) if Path(candidate).name == candidate else candidate
|
||||
if not resolved:
|
||||
continue
|
||||
path = Path(resolved)
|
||||
if path.exists():
|
||||
return path
|
||||
return None
|
||||
|
||||
|
||||
_codex_login_tasks: set[asyncio.Task] = set()
|
||||
|
||||
|
||||
async def start_codex_browser_login(command: Path) -> dict[str, Any]:
|
||||
process = await asyncio.create_subprocess_exec(
|
||||
str(command),
|
||||
"app-server",
|
||||
stdin=asyncio.subprocess.PIPE,
|
||||
stdout=asyncio.subprocess.PIPE,
|
||||
stderr=asyncio.subprocess.PIPE,
|
||||
creationflags=subprocess.CREATE_NO_WINDOW if sys.platform == "win32" else 0,
|
||||
)
|
||||
request_id = 1
|
||||
|
||||
async def write(payload: dict[str, Any]) -> None:
|
||||
if process.stdin is None:
|
||||
raise RuntimeError("Codex App Server stdin is unavailable.")
|
||||
process.stdin.write((json.dumps(payload, ensure_ascii=True) + "\n").encode("utf-8"))
|
||||
await process.stdin.drain()
|
||||
|
||||
async def read(timeout: int = 30) -> dict[str, Any]:
|
||||
if process.stdout is None:
|
||||
raise RuntimeError("Codex App Server stdout is unavailable.")
|
||||
try:
|
||||
line = await asyncio.wait_for(process.stdout.readline(), timeout=timeout)
|
||||
except asyncio.TimeoutError as exc:
|
||||
raise RuntimeError("Codex App Server timed out while starting browser login.") from exc
|
||||
if not line:
|
||||
stderr = ""
|
||||
if process.stderr is not None:
|
||||
try:
|
||||
stderr = (await asyncio.wait_for(process.stderr.read(), timeout=1)).decode("utf-8", errors="replace").strip()
|
||||
except asyncio.TimeoutError:
|
||||
stderr = ""
|
||||
raise RuntimeError(stderr or "Codex App Server exited before login completed.")
|
||||
return json.loads(line.decode("utf-8", errors="replace"))
|
||||
|
||||
async def send(method: str, params: dict[str, Any] | None = None) -> dict[str, Any]:
|
||||
nonlocal request_id
|
||||
current_id = request_id
|
||||
request_id += 1
|
||||
payload: dict[str, Any] = {"jsonrpc": "2.0", "id": current_id, "method": method}
|
||||
if params is not None:
|
||||
payload["params"] = params
|
||||
await write(payload)
|
||||
while True:
|
||||
message = await read()
|
||||
if message.get("id") == current_id:
|
||||
if message.get("error"):
|
||||
error = message["error"]
|
||||
raise RuntimeError(error.get("message") or f"Codex App Server request failed: {error}")
|
||||
return message.get("result") or {}
|
||||
await answer_codex_login_server_request(write, message)
|
||||
|
||||
try:
|
||||
await send(
|
||||
"initialize",
|
||||
{
|
||||
"clientInfo": {"name": "TraderAI", "version": __version__},
|
||||
"capabilities": {"experimentalApi": True},
|
||||
},
|
||||
)
|
||||
await write({"jsonrpc": "2.0", "method": "initialized", "params": {}})
|
||||
login = await send("account/login/start", {"type": "chatgpt"})
|
||||
if login.get("type") != "chatgpt" or not login.get("authUrl"):
|
||||
raise RuntimeError(f"Codex App Server did not return a browser login URL: {login!r}")
|
||||
task = asyncio.create_task(watch_codex_browser_login(process, read, write, login.get("loginId")))
|
||||
_codex_login_tasks.add(task)
|
||||
task.add_done_callback(_codex_login_tasks.discard)
|
||||
return login
|
||||
except Exception:
|
||||
await stop_process(process)
|
||||
raise
|
||||
|
||||
|
||||
async def answer_codex_login_server_request(write: Any, message: dict[str, Any]) -> None:
|
||||
if "id" not in message or "method" not in message:
|
||||
return
|
||||
await write(
|
||||
{
|
||||
"jsonrpc": "2.0",
|
||||
"id": message["id"],
|
||||
"error": {"code": -32601, "message": "TraderAI login does not handle server requests."},
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
async def watch_codex_browser_login(process: asyncio.subprocess.Process, read: Any, write: Any, login_id: str | None) -> None:
|
||||
try:
|
||||
while True:
|
||||
message = await read(timeout=300)
|
||||
if message.get("method") == "account/login/completed":
|
||||
params = message.get("params") or {}
|
||||
if login_id is None or params.get("loginId") == login_id:
|
||||
return
|
||||
await answer_codex_login_server_request(write, message)
|
||||
except Exception:
|
||||
return
|
||||
finally:
|
||||
await stop_process(process)
|
||||
|
||||
|
||||
async def stop_process(process: asyncio.subprocess.Process) -> None:
|
||||
if process.returncode is not None:
|
||||
return
|
||||
process.terminate()
|
||||
try:
|
||||
await asyncio.wait_for(process.wait(), timeout=3)
|
||||
except asyncio.TimeoutError:
|
||||
process.kill()
|
||||
await process.wait()
|
||||
|
||||
|
||||
async def inspect_codex_app_server(command: Path) -> tuple[dict[str, Any] | None, list[str], dict[str, list[str]]]:
|
||||
process = await asyncio.create_subprocess_exec(
|
||||
str(command),
|
||||
"app-server",
|
||||
stdin=asyncio.subprocess.PIPE,
|
||||
stdout=asyncio.subprocess.PIPE,
|
||||
stderr=asyncio.subprocess.PIPE,
|
||||
creationflags=subprocess.CREATE_NO_WINDOW if sys.platform == "win32" else 0,
|
||||
)
|
||||
request_id = 1
|
||||
|
||||
async def write(payload: dict[str, Any]) -> None:
|
||||
if process.stdin is None:
|
||||
raise RuntimeError("Codex App Server stdin is unavailable.")
|
||||
process.stdin.write((json.dumps(payload, ensure_ascii=True) + "\n").encode("utf-8"))
|
||||
await process.stdin.drain()
|
||||
|
||||
async def read(timeout: int = 30) -> dict[str, Any]:
|
||||
if process.stdout is None:
|
||||
raise RuntimeError("Codex App Server stdout is unavailable.")
|
||||
line = await asyncio.wait_for(process.stdout.readline(), timeout=timeout)
|
||||
if not line:
|
||||
stderr = ""
|
||||
if process.stderr is not None:
|
||||
try:
|
||||
stderr = (await asyncio.wait_for(process.stderr.read(), timeout=1)).decode("utf-8", errors="replace").strip()
|
||||
except asyncio.TimeoutError:
|
||||
stderr = ""
|
||||
raise RuntimeError(stderr or "Codex App Server exited without a response.")
|
||||
return json.loads(line.decode("utf-8", errors="replace"))
|
||||
|
||||
async def send(method: str, params: dict[str, Any] | None = None) -> dict[str, Any]:
|
||||
nonlocal request_id
|
||||
current_id = request_id
|
||||
request_id += 1
|
||||
payload: dict[str, Any] = {"jsonrpc": "2.0", "id": current_id, "method": method}
|
||||
if params is not None:
|
||||
payload["params"] = params
|
||||
await write(payload)
|
||||
while True:
|
||||
message = await read()
|
||||
if message.get("id") == current_id:
|
||||
if message.get("error"):
|
||||
error = message["error"]
|
||||
raise RuntimeError(error.get("message") or f"Codex App Server request failed: {error}")
|
||||
return message.get("result") or {}
|
||||
if "id" in message and "method" in message:
|
||||
await write(
|
||||
{
|
||||
"jsonrpc": "2.0",
|
||||
"id": message["id"],
|
||||
"error": {"code": -32601, "message": "TraderAI status checks do not handle server requests."},
|
||||
}
|
||||
)
|
||||
|
||||
try:
|
||||
await send(
|
||||
"initialize",
|
||||
{
|
||||
"clientInfo": {"name": "TraderAI", "version": __version__},
|
||||
"capabilities": {"experimentalApi": True},
|
||||
},
|
||||
)
|
||||
await write({"jsonrpc": "2.0", "method": "initialized", "params": {}})
|
||||
account_result = await send("account/read", {"refreshToken": False})
|
||||
models: list[str] = []
|
||||
effort_map: dict[str, list[str]] = {}
|
||||
cursor: str | None = None
|
||||
for _ in range(20):
|
||||
params: dict[str, Any] = {"limit": 50, "includeHidden": False}
|
||||
if cursor:
|
||||
params["cursor"] = cursor
|
||||
page = await send("model/list", params)
|
||||
for item in page.get("data") or []:
|
||||
model = item.get("id") or item.get("model")
|
||||
if not model:
|
||||
continue
|
||||
models.append(model)
|
||||
efforts = [
|
||||
effort.get("reasoningEffort")
|
||||
for effort in item.get("supportedReasoningEfforts", [])
|
||||
if effort.get("reasoningEffort")
|
||||
]
|
||||
if efforts:
|
||||
effort_map[model] = efforts
|
||||
cursor = page.get("nextCursor")
|
||||
if not cursor:
|
||||
break
|
||||
return account_result.get("account"), sorted(set(models)), effort_map
|
||||
finally:
|
||||
if process.returncode is None:
|
||||
process.terminate()
|
||||
try:
|
||||
await asyncio.wait_for(process.wait(), timeout=3)
|
||||
except asyncio.TimeoutError:
|
||||
process.kill()
|
||||
await process.wait()
|
||||
|
||||
|
||||
def codex_models() -> list[str]:
|
||||
cache_path = Path.home() / ".codex" / "models_cache.json"
|
||||
if not cache_path.exists():
|
||||
return []
|
||||
try:
|
||||
body = json.loads(cache_path.read_text(encoding="utf-8"))
|
||||
except (OSError, ValueError):
|
||||
return []
|
||||
models = []
|
||||
for item in body.get("models", []):
|
||||
slug = item.get("slug")
|
||||
if slug:
|
||||
models.append(slug)
|
||||
return sorted(set(models))
|
||||
|
||||
|
||||
def codex_reasoning_efforts(model: str, effort_map: dict[str, list[str]] | None = None) -> list[str]:
|
||||
if effort_map and effort_map.get(model):
|
||||
return effort_map[model]
|
||||
cache_path = Path.home() / ".codex" / "models_cache.json"
|
||||
if not cache_path.exists():
|
||||
return reasoning_effort_options()
|
||||
try:
|
||||
body = json.loads(cache_path.read_text(encoding="utf-8"))
|
||||
except (OSError, ValueError):
|
||||
return reasoning_effort_options()
|
||||
for item in body.get("models", []):
|
||||
if item.get("slug") != model:
|
||||
continue
|
||||
efforts = [entry.get("effort") for entry in item.get("supported_reasoning_levels", []) if entry.get("effort")]
|
||||
return efforts or reasoning_effort_options()
|
||||
return reasoning_effort_options()
|
||||
|
||||
|
||||
def reasoning_effort_options() -> list[str]:
|
||||
return ["none", "minimal", "low", "medium", "high", "xhigh"]
|
||||
|
||||
|
||||
def find_ollama_executable() -> Path | None:
|
||||
candidates = [
|
||||
shutil.which("ollama"),
|
||||
@@ -671,6 +1108,13 @@ def popen_hidden(command: list[str]) -> subprocess.Popen:
|
||||
return subprocess.Popen(command, **kwargs)
|
||||
|
||||
|
||||
def exception_detail(exc: BaseException) -> str:
|
||||
text = str(exc).strip()
|
||||
if text:
|
||||
return text
|
||||
return f"{type(exc).__name__}: {exc!r}"
|
||||
|
||||
|
||||
async def inspect_update() -> dict[str, Any]:
|
||||
try:
|
||||
latest = await latest_release()
|
||||
|
||||
Reference in New Issue
Block a user