add API and Openrouter
This commit is contained in:
@@ -0,0 +1 @@
|
||||
"""HTTP API helpers (FastAPI) for running Graph of Thoughts with OpenRouter."""
|
||||
@@ -0,0 +1,4 @@
|
||||
from graph_of_thoughts.api.app import run
|
||||
|
||||
if __name__ == "__main__":
|
||||
run()
|
||||
@@ -0,0 +1,192 @@
|
||||
# Copyright (c) 2023 ETH Zurich.
|
||||
# All rights reserved.
|
||||
#
|
||||
# Use of this source code is governed by a BSD-style license that can be
|
||||
# found in the LICENSE file.
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
import time
|
||||
import uuid
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from graph_of_thoughts.api.got_openai_pipeline import (
|
||||
ChatCompletionParser,
|
||||
ChatCompletionPrompter,
|
||||
build_default_chat_graph,
|
||||
extract_assistant_text,
|
||||
format_chat_messages,
|
||||
)
|
||||
from graph_of_thoughts.controller import Controller
|
||||
from graph_of_thoughts.language_models.openrouter import (
|
||||
OpenRouter,
|
||||
OpenRouterBadRequestError,
|
||||
OpenRouterRateLimitError,
|
||||
)
|
||||
|
||||
try:
|
||||
from fastapi import FastAPI, HTTPException
|
||||
from fastapi.responses import JSONResponse
|
||||
from pydantic import BaseModel, Field
|
||||
except ImportError as e:
|
||||
raise ImportError(
|
||||
"FastAPI and Pydantic are required for the HTTP API. "
|
||||
'Install with: pip install "graph_of_thoughts[api]"'
|
||||
) from e
|
||||
|
||||
|
||||
class ChatMessage(BaseModel):
|
||||
role: str
|
||||
content: str
|
||||
|
||||
|
||||
class ChatCompletionRequest(BaseModel):
|
||||
model: Optional[str] = None
|
||||
messages: List[ChatMessage]
|
||||
temperature: Optional[float] = None
|
||||
max_tokens: Optional[int] = None
|
||||
stream: Optional[bool] = False
|
||||
n: Optional[int] = Field(default=1, ge=1, le=1)
|
||||
|
||||
|
||||
def _get_config_path() -> str:
|
||||
return os.environ.get(
|
||||
"OPENROUTER_CONFIG",
|
||||
os.path.join(
|
||||
os.path.dirname(os.path.dirname(os.path.abspath(__file__))),
|
||||
"language_models",
|
||||
"config.openrouter.yaml",
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def _run_controller(lm: OpenRouter, user_text: str) -> str:
|
||||
graph = build_default_chat_graph(num_candidates=3)
|
||||
ctrl = Controller(
|
||||
lm,
|
||||
graph,
|
||||
ChatCompletionPrompter(),
|
||||
ChatCompletionParser(),
|
||||
{"input": user_text},
|
||||
)
|
||||
ctrl.run()
|
||||
return extract_assistant_text(ctrl.get_final_thoughts())
|
||||
|
||||
|
||||
app = FastAPI(
|
||||
title="Graph of Thoughts (OpenRouter)",
|
||||
version="0.1.0",
|
||||
description="OpenAI-compatible chat completions backed by Graph of Operations + OpenRouter.",
|
||||
)
|
||||
|
||||
|
||||
@app.on_event("startup")
|
||||
def _startup() -> None:
|
||||
logging.basicConfig(level=os.environ.get("LOG_LEVEL", "INFO"))
|
||||
|
||||
|
||||
@app.get("/v1/models")
|
||||
def list_models() -> Dict[str, Any]:
|
||||
path = _get_config_path()
|
||||
if not os.path.isfile(path):
|
||||
return {"object": "list", "data": []}
|
||||
from graph_of_thoughts.language_models.openrouter import load_openrouter_config
|
||||
|
||||
cfg = load_openrouter_config(path)
|
||||
models = cfg.get("models") or []
|
||||
if isinstance(models, str):
|
||||
models = [models]
|
||||
data = [
|
||||
{
|
||||
"id": m,
|
||||
"object": "model",
|
||||
"created": int(time.time()),
|
||||
"owned_by": "openrouter",
|
||||
}
|
||||
for m in models
|
||||
]
|
||||
return {"object": "list", "data": data}
|
||||
|
||||
|
||||
@app.post("/v1/chat/completions")
|
||||
def chat_completions(body: ChatCompletionRequest) -> JSONResponse:
|
||||
if body.stream:
|
||||
raise HTTPException(
|
||||
status_code=400,
|
||||
detail="stream=true is not supported; use stream=false.",
|
||||
)
|
||||
if body.n != 1:
|
||||
raise HTTPException(status_code=400, detail="Only n=1 is supported.")
|
||||
|
||||
path = _get_config_path()
|
||||
if not os.path.isfile(path):
|
||||
raise HTTPException(
|
||||
status_code=500,
|
||||
detail=f"OpenRouter config not found at {path}. Set OPENROUTER_CONFIG.",
|
||||
)
|
||||
|
||||
lm = OpenRouter(config_path=path)
|
||||
user_text = format_chat_messages(
|
||||
[{"role": m.role, "content": m.content} for m in body.messages]
|
||||
)
|
||||
try:
|
||||
lm.set_request_overrides(
|
||||
model=body.model,
|
||||
temperature=body.temperature,
|
||||
max_tokens=body.max_tokens,
|
||||
)
|
||||
try:
|
||||
answer = _run_controller(lm, user_text)
|
||||
finally:
|
||||
lm.clear_request_overrides()
|
||||
except OpenRouterRateLimitError as e:
|
||||
raise HTTPException(status_code=429, detail=str(e)) from e
|
||||
except OpenRouterBadRequestError as e:
|
||||
raise HTTPException(status_code=400, detail=str(e)) from e
|
||||
|
||||
model_id = (
|
||||
body.model
|
||||
or lm.generation_model_id
|
||||
or lm.last_model_id
|
||||
or (lm.models[0] if lm.models else "openrouter")
|
||||
)
|
||||
resp_id = f"chatcmpl-{uuid.uuid4().hex}"
|
||||
now = int(time.time())
|
||||
payload = {
|
||||
"id": resp_id,
|
||||
"object": "chat.completion",
|
||||
"created": now,
|
||||
"model": model_id,
|
||||
"choices": [
|
||||
{
|
||||
"index": 0,
|
||||
"message": {"role": "assistant", "content": answer},
|
||||
"finish_reason": "stop",
|
||||
}
|
||||
],
|
||||
"usage": {
|
||||
"prompt_tokens": lm.prompt_tokens,
|
||||
"completion_tokens": lm.completion_tokens,
|
||||
"total_tokens": lm.prompt_tokens + lm.completion_tokens,
|
||||
},
|
||||
}
|
||||
return JSONResponse(content=payload)
|
||||
|
||||
|
||||
def run() -> None:
|
||||
import uvicorn
|
||||
|
||||
host = os.environ.get("HOST", "0.0.0.0")
|
||||
port = int(os.environ.get("PORT", "8000"))
|
||||
uvicorn.run(
|
||||
"graph_of_thoughts.api.app:app",
|
||||
host=host,
|
||||
port=port,
|
||||
reload=os.environ.get("RELOAD", "").lower() in ("1", "true", "yes"),
|
||||
)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
run()
|
||||
@@ -0,0 +1,123 @@
|
||||
# Copyright (c) 2023 ETH Zurich.
|
||||
# All rights reserved.
|
||||
#
|
||||
# Use of this source code is governed by a BSD-style license that can be
|
||||
# found in the LICENSE file.
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import re
|
||||
from typing import Any, Dict, List, Union
|
||||
|
||||
from graph_of_thoughts.operations import GraphOfOperations, operations
|
||||
from graph_of_thoughts.operations.thought import Thought
|
||||
from graph_of_thoughts.parser import Parser
|
||||
from graph_of_thoughts.prompter import Prompter
|
||||
|
||||
|
||||
def format_chat_messages(messages: List[Dict[str, str]]) -> str:
|
||||
parts: List[str] = []
|
||||
for m in messages:
|
||||
role = m.get("role", "user")
|
||||
content = m.get("content", "")
|
||||
if not isinstance(content, str):
|
||||
content = str(content)
|
||||
parts.append(f"{role.upper()}:\n{content}")
|
||||
return "\n\n".join(parts)
|
||||
|
||||
|
||||
class ChatCompletionPrompter(Prompter):
|
||||
"""Prompter for a small generate → score → keep-best Graph of Operations."""
|
||||
|
||||
def generate_prompt(self, num_branches: int, **kwargs: Any) -> str:
|
||||
problem = kwargs.get("input", "")
|
||||
return (
|
||||
"You are a careful assistant. Read the conversation below and produce "
|
||||
"one candidate answer for the USER's latest needs.\n\n"
|
||||
f"{problem}\n\n"
|
||||
"Reply with your answer only, no preamble."
|
||||
)
|
||||
|
||||
def score_prompt(self, state_dicts: List[Dict], **kwargs: Any) -> str:
|
||||
lines = [
|
||||
"You evaluate candidate answers for the same problem. "
|
||||
"Score each candidate from 0 (worst) to 10 (best) on correctness, "
|
||||
"completeness, and relevance.",
|
||||
"",
|
||||
"Return ONLY a JSON array of numbers, one score per candidate in order, e.g. [7, 5, 9].",
|
||||
"",
|
||||
]
|
||||
for i, st in enumerate(state_dicts):
|
||||
cand = st.get("candidate", "")
|
||||
lines.append(f"Candidate {i}:\n{cand}\n")
|
||||
return "\n".join(lines)
|
||||
|
||||
def aggregation_prompt(self, state_dicts: List[Dict], **kwargs: Any) -> str:
|
||||
raise RuntimeError("aggregation_prompt is not used by the chat completion pipeline")
|
||||
|
||||
def improve_prompt(self, **kwargs: Any) -> str:
|
||||
raise RuntimeError("improve_prompt is not used by the chat completion pipeline")
|
||||
|
||||
def validation_prompt(self, **kwargs: Any) -> str:
|
||||
raise RuntimeError("validation_prompt is not used by the chat completion pipeline")
|
||||
|
||||
|
||||
class ChatCompletionParser(Parser):
|
||||
def parse_generate_answer(self, state: Dict, texts: List[str]) -> List[Dict]:
|
||||
out: List[Dict] = []
|
||||
for i, t in enumerate(texts):
|
||||
out.append({"candidate": (t or "").strip(), "branch_index": i})
|
||||
return out
|
||||
|
||||
def parse_score_answer(self, states: List[Dict], texts: List[str]) -> List[float]:
|
||||
raw = texts[0] if texts else ""
|
||||
scores = self._scores_from_text(raw, len(states))
|
||||
if len(scores) < len(states):
|
||||
scores.extend([0.0] * (len(states) - len(scores)))
|
||||
return scores[: len(states)]
|
||||
|
||||
def _scores_from_text(self, raw: str, n: int) -> List[float]:
|
||||
raw = raw.strip()
|
||||
try:
|
||||
data = json.loads(raw)
|
||||
if isinstance(data, list):
|
||||
return [float(x) for x in data]
|
||||
except (json.JSONDecodeError, ValueError, TypeError):
|
||||
pass
|
||||
nums = re.findall(r"-?\d+(?:\.\d+)?", raw)
|
||||
return [float(x) for x in nums[:n]]
|
||||
|
||||
def parse_aggregation_answer(
|
||||
self, states: List[Dict], texts: List[str]
|
||||
) -> Union[Dict, List[Dict]]:
|
||||
raise RuntimeError("parse_aggregation_answer is not used")
|
||||
|
||||
def parse_improve_answer(self, state: Dict, texts: List[str]) -> Dict:
|
||||
raise RuntimeError("parse_improve_answer is not used")
|
||||
|
||||
def parse_validation_answer(self, state: Dict, texts: List[str]) -> bool:
|
||||
raise RuntimeError("parse_validation_answer is not used")
|
||||
|
||||
|
||||
def build_default_chat_graph(num_candidates: int = 3) -> GraphOfOperations:
|
||||
g = GraphOfOperations()
|
||||
g.append_operation(
|
||||
operations.Generate(
|
||||
num_branches_prompt=1, num_branches_response=num_candidates
|
||||
)
|
||||
)
|
||||
g.append_operation(operations.Score(combined_scoring=True))
|
||||
g.append_operation(operations.KeepBestN(1))
|
||||
return g
|
||||
|
||||
|
||||
def extract_assistant_text(final_thoughts_list: List[List[Thought]]) -> str:
|
||||
"""``get_final_thoughts`` returns one list per leaf operation; we take the first leaf's first thought."""
|
||||
if not final_thoughts_list:
|
||||
return ""
|
||||
thoughts = final_thoughts_list[0]
|
||||
if not thoughts:
|
||||
return ""
|
||||
state = thoughts[0].state or {}
|
||||
return str(state.get("candidate", ""))
|
||||
Reference in New Issue
Block a user