Files
OpenJarvis/src/openjarvis/engine/cloud.py
T
Jon Saad-FalconandClaude Opus 4.6 990d7d8a79 Expand test suite to 1031 tests: new agents, tools, MCP layer, model catalog
Add ReAct and OpenHands agents, WebSearch and CodeInterpreter tools,
full MCP protocol layer (server/client/transport), Gemini cloud engine
support, 12 new model specs (4 local MoE + 8 cloud), trace system,
and comprehensive test coverage across all dimensions (hardware, engine,
memory, agents, tools, MCP, integration).

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-02-21 04:59:11 +00:00

422 lines
13 KiB
Python

"""Cloud inference engine — OpenAI, Anthropic, and Google API backends."""
from __future__ import annotations
import os
from collections.abc import AsyncIterator, Sequence
from typing import Any, Dict, List
from openjarvis.core.registry import EngineRegistry
from openjarvis.core.types import Message
from openjarvis.engine._base import (
EngineConnectionError,
InferenceEngine,
messages_to_dicts,
)
# Pricing per million tokens (input, output)
PRICING: Dict[str, tuple[float, float]] = {
"gpt-4o": (2.50, 10.00),
"gpt-4o-mini": (0.15, 0.60),
"gpt-5": (10.00, 30.00),
"gpt-5-mini": (0.25, 2.00),
"o3-mini": (1.10, 4.40),
"claude-sonnet-4-20250514": (3.00, 15.00),
"claude-opus-4-20250514": (15.00, 75.00),
"claude-haiku-3-5-20241022": (0.80, 4.00),
"claude-opus-4-6": (5.00, 25.00),
"claude-sonnet-4-6": (3.00, 15.00),
"claude-haiku-4-5": (1.00, 5.00),
"gemini-2.5-pro": (1.25, 10.00),
"gemini-2.5-flash": (0.30, 2.50),
"gemini-3-pro": (2.00, 12.00),
"gemini-3-flash": (0.50, 3.00),
}
# Well-known model IDs per provider
_OPENAI_MODELS = ["gpt-4o", "gpt-4o-mini", "gpt-5", "gpt-5-mini", "o3-mini"]
_ANTHROPIC_MODELS = [
"claude-sonnet-4-20250514",
"claude-opus-4-20250514",
"claude-haiku-3-5-20241022",
"claude-opus-4-6",
"claude-sonnet-4-6",
"claude-haiku-4-5",
]
_GOOGLE_MODELS = [
"gemini-2.5-pro",
"gemini-2.5-flash",
"gemini-3-pro",
"gemini-3-flash",
]
def _is_anthropic_model(model: str) -> bool:
return "claude" in model.lower()
def _is_google_model(model: str) -> bool:
return "gemini" in model.lower()
def estimate_cost(model: str, prompt_tokens: int, completion_tokens: int) -> float:
"""Estimate USD cost based on the hardcoded pricing table."""
# Try exact match first, then prefix match
prices = PRICING.get(model)
if prices is None:
for key, val in PRICING.items():
if model.startswith(key):
prices = val
break
if prices is None:
return 0.0
input_cost = (prompt_tokens / 1_000_000) * prices[0]
output_cost = (completion_tokens / 1_000_000) * prices[1]
return input_cost + output_cost
@EngineRegistry.register("cloud")
class CloudEngine(InferenceEngine):
"""Cloud inference via OpenAI, Anthropic, and Google SDKs."""
engine_id = "cloud"
def __init__(self) -> None:
self._openai_client: Any = None
self._anthropic_client: Any = None
self._google_client: Any = None
self._init_clients()
def _init_clients(self) -> None:
if os.environ.get("OPENAI_API_KEY"):
try:
import openai
self._openai_client = openai.OpenAI()
except ImportError:
pass
if os.environ.get("ANTHROPIC_API_KEY"):
try:
import anthropic
self._anthropic_client = anthropic.Anthropic()
except ImportError:
pass
gemini_key = (
os.environ.get("GEMINI_API_KEY")
or os.environ.get("GOOGLE_API_KEY")
)
if gemini_key:
try:
from google import genai
self._google_client = genai.Client(api_key=gemini_key)
except ImportError:
pass
def _generate_openai(
self,
messages: Sequence[Message],
*,
model: str,
temperature: float,
max_tokens: int,
**kwargs: Any,
) -> Dict[str, Any]:
if self._openai_client is None:
raise EngineConnectionError(
"OpenAI client not available — set "
"OPENAI_API_KEY and install "
"openjarvis[inference-cloud]"
)
resp = self._openai_client.chat.completions.create(
model=model,
messages=messages_to_dicts(messages),
temperature=temperature,
max_tokens=max_tokens,
**kwargs,
)
choice = resp.choices[0]
usage = resp.usage
prompt_tokens = usage.prompt_tokens if usage else 0
completion_tokens = usage.completion_tokens if usage else 0
return {
"content": choice.message.content or "",
"usage": {
"prompt_tokens": prompt_tokens,
"completion_tokens": completion_tokens,
"total_tokens": (usage.total_tokens if usage else 0),
},
"model": resp.model,
"finish_reason": choice.finish_reason or "stop",
"cost_usd": estimate_cost(model, prompt_tokens, completion_tokens),
}
def _generate_anthropic(
self,
messages: Sequence[Message],
*,
model: str,
temperature: float,
max_tokens: int,
**kwargs: Any,
) -> Dict[str, Any]:
if self._anthropic_client is None:
raise EngineConnectionError(
"Anthropic client not available — set "
"ANTHROPIC_API_KEY and install "
"openjarvis[inference-cloud]"
)
# Separate system message from conversation messages
system_text = ""
chat_msgs: List[Dict[str, Any]] = []
for m in messages:
if m.role.value == "system":
system_text = m.content
else:
chat_msgs.append({"role": m.role.value, "content": m.content})
create_kwargs: Dict[str, Any] = {
"model": model,
"messages": chat_msgs,
"temperature": temperature,
"max_tokens": max_tokens,
}
if system_text:
create_kwargs["system"] = system_text
resp = self._anthropic_client.messages.create(**create_kwargs)
content = resp.content[0].text if resp.content else ""
prompt_tokens = resp.usage.input_tokens if resp.usage else 0
completion_tokens = resp.usage.output_tokens if resp.usage else 0
return {
"content": content,
"usage": {
"prompt_tokens": prompt_tokens,
"completion_tokens": completion_tokens,
"total_tokens": prompt_tokens + completion_tokens,
},
"model": resp.model,
"finish_reason": resp.stop_reason or "stop",
"cost_usd": estimate_cost(model, prompt_tokens, completion_tokens),
}
def _generate_google(
self,
messages: Sequence[Message],
*,
model: str,
temperature: float,
max_tokens: int,
**kwargs: Any,
) -> Dict[str, Any]:
if self._google_client is None:
raise EngineConnectionError(
"Google client not available — set "
"GEMINI_API_KEY or GOOGLE_API_KEY and install "
"openjarvis[inference-google]"
)
# Build contents from messages
system_text = ""
contents: List[Dict[str, Any]] = []
for m in messages:
if m.role.value == "system":
system_text = m.content
elif m.role.value == "assistant":
contents.append({"role": "model", "parts": [{"text": m.content}]})
else:
contents.append({"role": "user", "parts": [{"text": m.content}]})
from google.genai import types as genai_types
config = genai_types.GenerateContentConfig(
temperature=temperature,
max_output_tokens=max_tokens,
)
if system_text:
config.system_instruction = system_text
resp = self._google_client.models.generate_content(
model=model,
contents=contents,
config=config,
)
content = resp.text or ""
um = resp.usage_metadata
prompt_tokens = (
getattr(um, "prompt_token_count", 0) if um else 0
)
completion_tokens = (
getattr(um, "candidates_token_count", 0) if um else 0
)
return {
"content": content,
"usage": {
"prompt_tokens": prompt_tokens,
"completion_tokens": completion_tokens,
"total_tokens": prompt_tokens + completion_tokens,
},
"model": model,
"finish_reason": "stop",
"cost_usd": estimate_cost(model, prompt_tokens, completion_tokens),
}
def generate(
self,
messages: Sequence[Message],
*,
model: str,
temperature: float = 0.7,
max_tokens: int = 1024,
**kwargs: Any,
) -> Dict[str, Any]:
kw = dict(
model=model,
temperature=temperature,
max_tokens=max_tokens,
**kwargs,
)
if _is_anthropic_model(model):
return self._generate_anthropic(messages, **kw)
if _is_google_model(model):
return self._generate_google(messages, **kw)
return self._generate_openai(messages, **kw)
async def stream(
self,
messages: Sequence[Message],
*,
model: str,
temperature: float = 0.7,
max_tokens: int = 1024,
**kwargs: Any,
) -> AsyncIterator[str]:
kw = dict(
model=model,
temperature=temperature,
max_tokens=max_tokens,
**kwargs,
)
if _is_anthropic_model(model):
async for token in self._stream_anthropic(
messages, **kw
):
yield token
elif _is_google_model(model):
async for token in self._stream_google(
messages, **kw
):
yield token
else:
async for token in self._stream_openai(
messages, **kw
):
yield token
async def _stream_openai(
self,
messages: Sequence[Message],
*,
model: str,
temperature: float,
max_tokens: int,
**kwargs: Any,
) -> AsyncIterator[str]:
if self._openai_client is None:
raise EngineConnectionError("OpenAI client not available")
resp = self._openai_client.chat.completions.create(
model=model,
messages=messages_to_dicts(messages),
temperature=temperature,
max_tokens=max_tokens,
stream=True,
**kwargs,
)
for chunk in resp:
delta = chunk.choices[0].delta if chunk.choices else None
if delta and delta.content:
yield delta.content
async def _stream_anthropic(
self,
messages: Sequence[Message],
*,
model: str,
temperature: float,
max_tokens: int,
**kwargs: Any,
) -> AsyncIterator[str]:
if self._anthropic_client is None:
raise EngineConnectionError("Anthropic client not available")
system_text = ""
chat_msgs: List[Dict[str, Any]] = []
for m in messages:
if m.role.value == "system":
system_text = m.content
else:
chat_msgs.append({"role": m.role.value, "content": m.content})
create_kwargs: Dict[str, Any] = {
"model": model,
"messages": chat_msgs,
"temperature": temperature,
"max_tokens": max_tokens,
}
if system_text:
create_kwargs["system"] = system_text
with self._anthropic_client.messages.stream(**create_kwargs) as stream:
for text in stream.text_stream:
yield text
async def _stream_google(
self,
messages: Sequence[Message],
*,
model: str,
temperature: float,
max_tokens: int,
**kwargs: Any,
) -> AsyncIterator[str]:
if self._google_client is None:
raise EngineConnectionError("Google client not available")
system_text = ""
contents: List[Dict[str, Any]] = []
for m in messages:
if m.role.value == "system":
system_text = m.content
elif m.role.value == "assistant":
contents.append({"role": "model", "parts": [{"text": m.content}]})
else:
contents.append({"role": "user", "parts": [{"text": m.content}]})
from google.genai import types as genai_types
config = genai_types.GenerateContentConfig(
temperature=temperature,
max_output_tokens=max_tokens,
)
if system_text:
config.system_instruction = system_text
for chunk in self._google_client.models.generate_content_stream(
model=model,
contents=contents,
config=config,
):
if chunk.text:
yield chunk.text
def list_models(self) -> List[str]:
models: List[str] = []
if self._openai_client is not None:
models.extend(_OPENAI_MODELS)
if self._anthropic_client is not None:
models.extend(_ANTHROPIC_MODELS)
if self._google_client is not None:
models.extend(_GOOGLE_MODELS)
return models
def health(self) -> bool:
return (
self._openai_client is not None
or self._anthropic_client is not None
or self._google_client is not None
)
__all__ = ["CloudEngine", "PRICING", "estimate_cost"]