Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 5 additions & 6 deletions backend/agents/create_agent_info.py
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@

from database.a2a_agent_db import PROTOCOL_JSONRPC
from services.memory_config_service import build_memory_context
from services.image_service import get_video_understanding_model, get_vlm_model
from services.model_gateway_service import get_vlm_adapter
from database.agent_db import (
search_agent_info_by_agent_id,
query_sub_agent_relations,
Expand Down Expand Up @@ -254,7 +254,7 @@ def _resolve_safe_input_budget(
except UncertaintyReserveBasisUnknown as exc:
# W2 uncertainty reserve needs context_window_tokens as the 10% basis.
# Falls through here when a model row has max_input_tokens set but
# context_window_tokens is NULL — possible for rows imported before
# context_window_tokens is NULL - possible for rows imported before
# W11 V1 save-time defaults landed, or for rows written directly via
# SQL/legacy import. Degrade to the same "no W2 snapshot" branch the
# caller already handles (falls back to W1 input_budget).
Expand Down Expand Up @@ -289,7 +289,7 @@ def _resolve_input_budget(
Calls ModelCapacityResolver with the catalog + operator overrides. Returns
snapshot.provider_input_limit_tokens and monitoring fields on success.
Falls back to _TOKEN_THRESHOLD_LEGACY_FALLBACK with no snapshot when
capacity is unknown — this is the migration-window behavior before all
capacity is unknown - this is the migration-window behavior before all
model rows are backfilled.
"""
if not isinstance(model_info, dict):
Expand Down Expand Up @@ -1600,15 +1600,14 @@ async def create_tool_config_list(
elif tool_config.class_name == "AnalyzeImageTool":
selected_model_id = param_dict.get("selected_model_id")
tool_config.metadata = {
# get_vlm_model reads the first multimodal slot, now shown as image understanding.
"vlm_model": get_vlm_model(tenant_id=tenant_id, model_id=selected_model_id),
"vlm_model": get_vlm_adapter(tenant_id, selected_model_id, slot="vlm"),
"storage_client": minio_client,
"validate_url_access": lambda urls: validate_urls_access(urls, user_id)
}
elif tool_config.class_name in ["AnalyzeAudioTool", "AnalyzeVideoTool"]:
selected_model_id = param_dict.get("selected_model_id")
tool_config.metadata = {
"vlm_model": get_video_understanding_model(tenant_id=tenant_id, model_id=selected_model_id),
"vlm_model": get_vlm_adapter(tenant_id, selected_model_id, slot="vlm3"),
"storage_client": minio_client,
"validate_url_access": lambda urls: validate_urls_access(urls, user_id)
}
Expand Down
249 changes: 249 additions & 0 deletions backend/services/model_gateway_service.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,249 @@
"""Backend bridge: turn DB model configs into gateway adapters.

Service factory functions keep their signatures but delegate adapter
construction to the gateway via :func:`get_adapter_from_config`.
"""

from __future__ import annotations

import logging
from typing import Any, Dict, Optional

from nexent import MessageObserver
from nexent.core.gateway import (
EmbeddingContext,
LLMContext,
LongContextLLMContext,
ModelContext,
STTContext,
TTSContext,
VLMContext,
get_gateway,
)
from nexent.core.gateway.registry import get_registry
from consts.const import MODEL_CONFIG_MAPPING, TEST_PCM_PATH
from database.model_management_db import get_model_by_model_id, get_model_records
from utils.config_utils import get_model_name_from_config, tenant_config_manager

logger = logging.getLogger("model_gateway_service")

# Normalize vendor aliases to canonical registry factory names.
_FACTORY_NORMALIZE: Dict[str, str] = {
"volc": "volc",
"volcano": "volc",
"volcengine": "volc",
"火山引擎": "volc",
"dashscope": "dashscope",
"ali": "ali",
"alibaba": "ali",
"阿里云": "ali",
"silicon": "siliconflow",
"siliconflow": "siliconflow",
"openai": "openai",
"tokenpony": "tokenpony",
"jina": "jina",
"cohere": "cohere",
"modelengine": "modelengine",
}

# Modality-specific default factory when the raw factory is empty/unknown.
_MODALITY_DEFAULT_FACTORY: Dict[str, str] = {
"llm": "openai",
"llm_long_context": "openai",
"vlm": "openai",
"embedding": "openai",
"rerank": "openai",
"stt": "ali",
"tts": "ali",
"multi_embedding": "jina",
}


def _normalize_factory(raw: Optional[str], modality: str) -> str:
"""Return the canonical registry factory for ``raw`` under ``modality``."""
cleaned = (raw or "").strip().lower()
factory = _FACTORY_NORMALIZE.get(cleaned, cleaned)
# STT/TTS historically route DashScope through the Ali client.
if modality in ("stt", "tts") and factory in ("dashscope", "ali", "alibaba"):
factory = "ali"
if get_registry().has(factory, modality):
return factory
default = _MODALITY_DEFAULT_FACTORY.get(modality, "openai")
if factory:
logger.debug(
"factory %r has no %s adapter; falling back to %r", factory, modality, default
)
return default


def _coalesce(*vals: Any) -> Any:
"""Return the first non-``None`` value, or ``None`` if all are ``None``.

Unlike ``a or b``, this preserves falsy-but-valid values such as
``temperature=0`` or ``top_p=0`` - an explicit ``0`` must reach the
adapter rather than being silently replaced by the cfg/default fallback.
"""
for v in vals:
if v is not None:
return v
return None


def _config_to_context(
cfg: Optional[dict],
modality: str,
slot: str,
tenant_id: Optional[str],
**construct_extras: Any,
) -> ModelContext:
"""Build a modality-specific :class:`ModelContext` from a DB config + per-call extras.

Args:
cfg: Raw model config dict from the DB (may be empty).
modality: The capability family identifier.
slot: The config slot key (e.g. "vlm" / "vlm3").
tenant_id: The tenant identifier.
**construct_extras: Per-call-site tuning (temperature, top_p, stream, ...);
known keys map to subclass fields directly.

Returns:
The modality-specific :class:`ModelContext` instance.
"""
cfg = cfg or {}
factory = _normalize_factory(cfg.get("model_factory"), modality)
needs_observer = modality in ("vlm", "llm", "llm_long_context")
observer = construct_extras.pop("observer", None)
if needs_observer and observer is None:
observer = MessageObserver()

common: Dict[str, Any] = dict(

Check warning on line 119 in backend/services/model_gateway_service.py

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Replace this constructor call with a literal.

See more on https://sonarcloud.io/project/issues?id=ModelEngine-Group_nexent&issues=AZ__B5JFMHc42Ihu-YkP&open=AZ__B5JFMHc42Ihu-YkP&pullRequest=3668
model_name=construct_extras.pop("model_name", None) or get_model_name_from_config(cfg) or "",
base_url=cfg.get("base_url", ""),
api_key=cfg.get("api_key", ""),
modality=modality,
factory=factory,
tenant_id=tenant_id,
slot=slot,
ssl_verify=cfg.get("ssl_verify", True),
observer=observer,
display_name=_coalesce(construct_extras.pop("display_name", None), cfg.get("display_name")),
timeout_seconds=_coalesce(construct_extras.pop("timeout_seconds", None), cfg.get("timeout_seconds")),
)

if modality == "vlm":
caps = construct_extras.pop("capabilities", None) or {}
return VLMContext(
**common,
temperature=_coalesce(construct_extras.pop("temperature", None), cfg.get("temperature")),
top_p=_coalesce(construct_extras.pop("top_p", None), cfg.get("top_p")),
stream=construct_extras.pop("stream", None),
max_output_tokens=_coalesce(construct_extras.pop("max_output_tokens", None), cfg.get("max_output_tokens")),
frequency_penalty=cfg.get("frequency_penalty"),
extra_body=cfg.get("extra_body"),
max_tokens=cfg.get("max_tokens"),
capabilities=caps,
)
raise ValueError(f"Unknown modality: {modality}")


def get_adapter_from_config(
cfg: Optional[dict],
modality: str,
slot: str,
tenant_id: Optional[str] = None,
**construct_extras: Any,
):
"""Resolve and return the adapter for ``cfg`` (cached by the gateway)."""
context = _config_to_context(cfg, modality, slot, tenant_id, **construct_extras)
return get_gateway().get_adapter(context)


def build_adapter_fresh(
cfg: Optional[dict],
modality: str,
slot: str,
tenant_id: Optional[str] = None,
**construct_extras: Any,
):
"""Build a fresh adapter for ``cfg`` WITHOUT the gateway instance cache.

Used by per-call construction sites (e.g. voice streaming sessions) where
vendor config carries per-request params (api_key, ws_url, voice, ...) that
must not collide across tenants under a shared cache key.

Args:
cfg: Raw model config dict.
modality: The capability family identifier.
slot: The config slot key.
tenant_id: The tenant identifier.
**construct_extras: Extra fields forwarded to the context.

Returns:
A newly constructed adapter instance.
"""
context = _config_to_context(cfg, modality, slot, tenant_id, **construct_extras)
cls = get_registry().resolve(context.factory, modality)
return cls(context)


def _fetch_slot_config(tenant_id, model_id, expected_type, slot_key):
"""Fetch a model config by model_id (with type check) or by slot key."""
if model_id:
cfg = get_model_by_model_id(int(model_id), tenant_id)
if not cfg:
raise ValueError(f"Model not found: {model_id}")
if cfg.get("model_type") != expected_type:
raise ValueError(
f"Selected model {model_id} is not a {expected_type} model"
)
return cfg
return tenant_config_manager.get_model_config(
key=MODEL_CONFIG_MAPPING.get(slot_key, slot_key), tenant_id=tenant_id
)


def _fetch_voice_config(tenant_id, model_type):
"""Fetch an STT/TTS config from tenant_config or model_records (voice fallback)."""
try:
cfg = tenant_config_manager.get_model_config(tenant_id, model_type)
if cfg and isinstance(cfg, dict):
return cfg
except Exception:
pass
try:
records = get_model_records({"model_type": model_type}, tenant_id)
if records:
return records[0]
except Exception:
pass
return None


def get_vlm_adapter_from_config(
cfg: Optional[dict],
tenant_id: Optional[str] = None,
slot: str = "vlm",
**construct_extras: Any,
):
return get_adapter_from_config(cfg, "vlm", slot, tenant_id, **construct_extras)


def get_vlm_adapter(tenant_id: str, model_id: Optional[int] = None, slot: str = "vlm"):
"""Resolve the VLM adapter directly (bridge owns config-fetch).

Replaces ``image_service.get_vlm_model`` / ``get_video_understanding_model``.

Args:
tenant_id: The tenant identifier.
model_id: Optional model id; defaults to the slot config when omitted.
slot: "vlm" (image) or "vlm3" (video/audio).

Returns:
A VLM adapter instance, or ``None`` when no config is available.
"""
cfg = _fetch_slot_config(tenant_id, model_id, expected_type=slot, slot_key=slot)
if not cfg:
return None
return get_gateway().get_adapter(_config_to_context(cfg, "vlm", slot, tenant_id))


16 changes: 8 additions & 8 deletions backend/services/model_health_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,11 @@
from typing import Optional

from nexent.core import MessageObserver
from nexent.core.models import OpenAIModel, OpenAIVLModel
from nexent.core.models import OpenAIModel
from nexent.core.models.embedding_model import JinaEmbedding, OpenAICompatibleEmbedding, DashScopeMultimodalEmbedding, SiliconflowMultimodalEmbedding
from nexent.monitor import set_monitoring_context, set_monitoring_operation
from nexent.core.models.rerank_model import OpenAICompatibleRerank
from services.model_gateway_service import build_adapter_fresh

from services.voice_service import get_voice_service
from consts.const import LOCALHOST_IP, LOCALHOST_NAME, DOCKER_INTERNAL_HOST
Expand Down Expand Up @@ -247,13 +248,12 @@ async def _perform_connectivity_check(
observer = MessageObserver()
set_monitoring_operation("connectivity_check",
display_name=display_name)
connectivity = await OpenAIVLModel(
observer,
model_id=model_name,
api_base=model_base_url,
api_key=model_api_key,
ssl_verify=ssl_verify
).check_connectivity()
connectivity = await build_adapter_fresh(
{"base_url": model_base_url, "api_key": model_api_key,
"ssl_verify": ssl_verify},
"vlm", "vlm", None, model_name=model_name,
observer=observer, display_name=display_name,
).health_check()
elif model_type == 'stt':
voice_service = get_voice_service()

Expand Down
7 changes: 3 additions & 4 deletions backend/services/tool_configuration_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@
from services.vectordatabase_service import get_embedding_model_by_index_name, get_rerank_model
from utils.http_client_utils import create_httpx_client
from database.client import minio_client
from services.image_service import get_video_understanding_model, get_vlm_model
from services.model_gateway_service import get_vlm_adapter
from nexent.monitor import set_monitoring_context, set_monitoring_operation
from services.vectordatabase_service import get_vector_db_core
from utils.langchain_utils import discover_langchain_modules
Expand Down Expand Up @@ -920,9 +920,8 @@ def _validate_local_tool(
if not tenant_id or not user_id:
raise ToolExecutionException(
f"Tenant ID and User ID are required for {tool_name} validation")
# get_vlm_model reads the first multimodal slot, now shown as image understanding.
selected_model_id = instantiation_params.get("selected_model_id")
image_to_text_model = get_vlm_model(tenant_id=tenant_id, model_id=selected_model_id)
image_to_text_model = get_vlm_adapter(tenant_id, selected_model_id, slot="vlm")
vlm_display_name = getattr(
image_to_text_model, 'display_name', None)
set_monitoring_context(tenant_id=tenant_id)
Expand All @@ -940,7 +939,7 @@ def _validate_local_tool(
raise ToolExecutionException(
f"Tenant ID and User ID are required for {tool_name} validation")
selected_model_id = instantiation_params.get("selected_model_id")
video_understanding_model = get_video_understanding_model(tenant_id=tenant_id, model_id=selected_model_id)
video_understanding_model = get_vlm_adapter(tenant_id, selected_model_id, slot="vlm3")
model_display_name = getattr(
video_understanding_model, 'display_name', None)
set_monitoring_context(tenant_id=tenant_id)
Expand Down
Loading
Loading