Skip to content
Closed
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
74 changes: 70 additions & 4 deletions backend/agents/create_agent_info.py
Original file line number Diff line number Diff line change
Expand Up @@ -274,6 +274,45 @@ def _operator_overrides_from_model_info(model_info: Optional[dict]) -> dict:
return overrides


def _agent_capacity_overrides(
agent_info: Optional[dict],
model_id: Optional[int],
) -> Dict[str, Any]:
"""Extract per-agent capacity overrides for the selected model.

v2.6.0 model_params_override entries may carry capacity fields
(context_window_tokens / max_input_tokens / max_output_tokens /
default_output_reserve_tokens / tokenizer_family) next to inference
params. When present they win over the model-level capacity columns in
W1/W2 resolution, mirroring how temperature/top_p overrides win.
"""
if not isinstance(agent_info, dict) or model_id is None:
return {}
override_map = agent_info.get("model_params_override")
if not isinstance(override_map, dict):
return {}
entry = override_map.get(str(model_id))
if not isinstance(entry, dict):
return {}
overrides: Dict[str, Any] = {}
for field in _OPERATOR_OVERRIDE_FIELDS:
value = entry.get(field)
if value is not None:
overrides[field] = value
# Per-agent override semantics: a filled value simply replaces the
# model-level value for THIS agent. "最大输出Token数" is the field users
# expect to control the actual per-request max_tokens, so mirror it into
# default_output_reserve_tokens unless the user set the reserve
# explicitly. Without this, the request would keep the model-level
# reserve (4096) and the filled cap alone would change nothing visible.
if (
"max_output_tokens" in overrides
and "default_output_reserve_tokens" not in overrides
):
overrides["default_output_reserve_tokens"] = overrides["max_output_tokens"]
return overrides


def _dominant_capacity_source(field_sources: dict) -> Optional[str]:
values = [value for value in field_sources.values() if value]
if not values:
Expand Down Expand Up @@ -363,6 +402,7 @@ def _resolve_context_budget(

def _resolve_input_budget(
model_info: Optional[dict],
capacity_overrides: Optional[Dict[str, Any]] = None,
) -> tuple[int, Optional[dict], Optional[ModelCapacitySnapshot]]:
"""Resolve the context-manager input budget for a model_record_t row.

Expand All @@ -371,6 +411,9 @@ def _resolve_input_budget(
Falls back to _TOKEN_THRESHOLD_LEGACY_FALLBACK with no snapshot when
capacity is unknown - this is the migration-window behavior before all
model rows are backfilled.

capacity_overrides carries per-agent capacity fields (from
model_params_override) that win over the model-level columns.
"""
if not isinstance(model_info, dict):
return _TOKEN_THRESHOLD_LEGACY_FALLBACK, None, None
Expand All @@ -383,10 +426,13 @@ def _resolve_input_budget(
"model_factory/provider is missing; capacity catalog matching is disabled"
)
try:
operator_overrides = _operator_overrides_from_model_info(model_info)
if capacity_overrides:
operator_overrides.update(capacity_overrides)
snapshot = resolve_capacity(
model_id=model_id,
provider=provider,
operator_overrides=_operator_overrides_from_model_info(model_info),
operator_overrides=operator_overrides,
capability_profiles=CAPABILITY_CATALOG,
)
logger.debug(
Expand Down Expand Up @@ -1423,9 +1469,14 @@ async def create_agent_config(
# W1 step 6: derive input budget via ModelCapacityResolver instead of
# treating model_info["max_tokens"] (a deprecated output cap) as a
# context threshold. Falls back to a safe constant when capacity is
# unknown during the migration window.
# unknown during the migration window. Per-agent capacity overrides
# (model_params_override) win over the model-level columns.
input_budget, capacity_snapshot, resolved_capacity_snapshot = (
_resolve_input_budget(model_info)
_resolve_input_budget(
model_info,
capacity_overrides=_agent_capacity_overrides(
agent_info, model_id_to_use),
)
)
else:
model_name = "main_model"
Expand Down Expand Up @@ -2297,12 +2348,27 @@ async def create_agent_run_info(
mc.temperature = override_entry["temperature"]
if override_entry.get("top_p") is not None:
mc.top_p = override_entry["top_p"]
# v2.6.0: capacity fields are overridable per-agent. The
# authoritative consumer for request shaping is the W1/W2
# resolution (agent-selected model), this keeps each
# ModelConfig consistent with its override entry.
for field in _OPERATOR_OVERRIDE_FIELDS:
if override_entry.get(field) is not None:
setattr(mc, field, override_entry[field])
override_extra = override_entry.get("extra_params")
if override_extra and isinstance(override_extra, dict):
merged = dict(mc.extra_body or {})
for k, v in override_extra.items():
if k == "__custom__" and isinstance(v, dict):
merged.update(v)
for custom_key, custom_value in v.items():
# A null custom value is an explicit
# removal marker: the agent opts out of a
# model-level custom param instead of
# inheriting it.
if custom_value is None:
merged.pop(custom_key, None)
else:
merged[custom_key] = custom_value
else:
merged[k] = v
mc.extra_body = merged if merged else None
Expand Down
1 change: 1 addition & 0 deletions backend/agents/nl2skill_agent.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ def create_nl2skill_agent_config(
tools=[],
max_steps=5,
model_name=model_name,
output_protocol="final_answer_envelope",
provide_run_summary=False,
instructions=system_prompt,
enable_planning=False,
Expand Down
8 changes: 4 additions & 4 deletions backend/apps/agent_evaluation_app.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,9 @@
get_evaluation_stats_impl,
list_agent_evaluation_cases_impl,
list_agent_evaluations_by_agent_impl,
trial_run_evaluator_impl,
)
from services.evaluation_report_service import generate_agent_evaluation_report_impl
from services.runtime_proxy_service import forward_agent_evaluation_trial_run
from utils.auth_utils import get_current_user_id, get_current_user_info


Expand Down Expand Up @@ -552,15 +552,15 @@ async def trial_run_api(
"""
try:
user_id, tenant_id = get_current_user_id(authorization)
result = await trial_run_evaluator_impl(
tenant_id=tenant_id,
user_id=user_id,
result = await forward_agent_evaluation_trial_run(
agent_id=payload.agent_id,
agent_version_no=payload.agent_version_no,
query=payload.query,
judge_model_id=payload.judge_model_id,
evaluator_ids=payload.evaluator_ids,
language=payload.language,
user_id=user_id,
tenant_id=tenant_id,
)
logger.info(
"trial_run_api OK: tenant=%s user=%s agent_id=%s version=%s "
Expand Down
48 changes: 47 additions & 1 deletion backend/apps/agent_evaluation_runtime_app.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
from typing import Annotated

from fastapi import APIRouter, Header, HTTPException
from nexent.core.concurrency import ManagedTaskSpec
from pydantic import BaseModel, Field

from consts.evaluation_status import EvalRunStatus
Expand All @@ -13,7 +14,6 @@
claim_agent_evaluation_run,
get_agent_evaluation,
)
from nexent.core.concurrency import ManagedTaskSpec
from services.thread_lifecycle_service import runtime_thread_manager
from utils.auth_utils import verify_internal_runtime_jwt

Expand All @@ -28,13 +28,59 @@ class EvaluationRunRequest(BaseModel):
agent_evaluation_id: int = Field(gt=0)


class TrialRunRequest(BaseModel):
"""Payload used by Config service for a non-persistent trial evaluation."""

agent_id: int
agent_version_no: int = 1
query: str
judge_model_id: int
evaluator_ids: list[int] | None = None
language: str = "zh"


def _load_evaluation_executor():
"""Load the evaluation service only when a runtime run is dispatched."""
from services.agent_evaluation_service import execute_agent_evaluation_run

return execute_agent_evaluation_run


def _load_trial_executor():
"""Load the trial executor only when Runtime receives a trial request."""
from services.agent_evaluation_service import trial_run_evaluator_impl

return trial_run_evaluator_impl


@router.post("/trial-run", include_in_schema=False)
async def trial_run_evaluation_api(
payload: TrialRunRequest,
authorization: Annotated[str | None, Header()] = None,
):
"""Run one ad-hoc evaluation in the Runtime process."""
try:
user_id, tenant_id = verify_internal_runtime_jwt(authorization)
except Exception as exc:
logger.warning("Rejected unauthenticated trial evaluation: %s", exc)
raise HTTPException(
status_code=HTTPStatus.UNAUTHORIZED,
detail="Invalid internal runtime authorization",
) from exc

trial_run_evaluator_impl = _load_trial_executor()
return await trial_run_evaluator_impl(
tenant_id=tenant_id,
user_id=user_id,
agent_id=payload.agent_id,
agent_version_no=payload.agent_version_no,
query=payload.query,
judge_model_id=payload.judge_model_id,
evaluator_ids=payload.evaluator_ids,
language=payload.language,
)


@router.post("/run", include_in_schema=False, status_code=HTTPStatus.ACCEPTED)
async def dispatch_evaluation_run_api(
payload: EvaluationRunRequest,
Expand Down
4 changes: 2 additions & 2 deletions backend/apps/human_interaction_app.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,12 +50,12 @@ async def capabilities(identity=Depends(identity_dependency)):
async def conversation_snapshot(conversation_id: int, identity=Depends(identity_dependency)):
def read(service, tenant_id, user_id):
run_id = service.repository.latest(tenant_id, user_id, conversation_id)
return service.snapshot(run_id, tenant_id, user_id) if run_id else None
return service.light_snapshot(run_id, tenant_id, user_id) if run_id else None
return await _call(identity, read)

@result.get("/{run_id}")
async def snapshot(run_id: str, identity=Depends(identity_dependency)):
return await _call(identity, lambda service, tenant, user: service.snapshot(run_id, tenant, user))
return await _call(identity, lambda service, tenant, user: service.light_snapshot(run_id, tenant, user))

@result.get("/{run_id}/events")
async def events(run_id: str, after_event: int = Query(0, ge=0), identity=Depends(identity_dependency)):
Expand Down
14 changes: 12 additions & 2 deletions backend/apps/model_managment_app.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
BatchCreateModelsRequest,
CapacitySuggestionFields,
ModelRequest,
ModelProbeRequest,
ModelCapacitySuggestionRequest,
ModelCapacitySuggestionResponse,
ProviderModelRequest,
Expand Down Expand Up @@ -64,6 +65,7 @@
from utils.auth_utils import get_current_user_id
from consts.exceptions import TokenExpiredError
from nexent.core.concurrency import run_blocking
from database.model_management_db import get_model_by_model_id

# Model Catalog loader (with graceful fallback)
try:
Expand Down Expand Up @@ -580,7 +582,7 @@ async def check_model_health(

@router.post("/temporary_healthcheck")
async def check_temporary_model_health(
request: ModelRequest, authorization: Optional[str] = Header(None)
request: ModelProbeRequest, authorization: Optional[str] = Header(None)
):
"""Verify connectivity for the provided model configuration without persisting it.

Expand All @@ -589,7 +591,15 @@ async def check_temporary_model_health(
authorization: Bearer token header used to enforce authentication.
"""
try:
get_current_user_id(authorization)
_, tenant_id = get_current_user_id(authorization)
# Edit-dialog probes arrive without the api_key (the backend never
# returns the persisted key to the client, and the dialog leaves the
# field empty to "keep existing"). Fall back to the stored key so
# verifying does not require retyping it.
if request.probe_model_id is not None and request.api_key in (None, "", "sk-no-api-key"):
stored_model = get_model_by_model_id(request.probe_model_id, tenant_id=tenant_id)
if stored_model and stored_model.get("api_key"):
request.api_key = stored_model["api_key"]
result = await verify_model_config_connectivity(request.model_dump())
if result.get("connectivity") is True:
# suggest_capacity may now issue an LLM self-report HTTP call
Expand Down
2 changes: 1 addition & 1 deletion backend/configs/model_catalog.json
Original file line number Diff line number Diff line change
Expand Up @@ -384,7 +384,7 @@
},
"modelengine": {
"display_name": "ModelEngine",
"base_url": "https://api.modelengine-ai.net/v1/",
"base_url": "https://<your-modelengine-host>/open/router/v1",
"models": {
"qwen3-8b": {
"model_type": "llm",
Expand Down
17 changes: 17 additions & 0 deletions backend/consts/model.py
Original file line number Diff line number Diff line change
Expand Up @@ -580,6 +580,23 @@ class ModelRequest(BaseModel):
accepted_capability_profile_version: Optional[str] = None


class ModelProbeRequest(ModelRequest):
"""Request payload for POST /model/temporary_healthcheck only.

Extends ModelRequest with probe-only fields. Never send this to the
create/update endpoints: they spread model_dump() straight into
INSERT/UPDATE column lists, so any field that is not a real
model_record_t column makes SQLAlchemy raise "Unconsumed column names"
(and a field matching a column name would inject wrong data — the
model_id NULL primary-key bug).
"""
# Edit-dialog connectivity probes omit the stored api_key (the backend
# never returns it to the client). When set, /temporary_healthcheck falls
# back to the persisted key for this model instead of probing with the
# "sk-no-api-key" placeholder.
probe_model_id: Optional[int] = None


class CapacitySuggestionFields(BaseModel):
context_window_tokens: Optional[int] = None
max_input_tokens: Optional[int] = None
Expand Down
5 changes: 5 additions & 0 deletions backend/consts/provider.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,3 +26,8 @@ class ProviderEnum(str, Enum):

# ModelEngine
# Base URL and API key are loaded from environment variables at runtime
# URL path segment identifying ModelEngine northbound endpoints
# (e.g. https://host:port/open/router/v1). Endpoints behind this marker
# serve self-signed certificates, so their records must keep
# ssl_verify=False on both the create and update paths.
MODEL_ENGINE_URL_MARKER = "open/router"
10 changes: 10 additions & 0 deletions backend/database/human_interaction_db.py
Original file line number Diff line number Diff line change
Expand Up @@ -168,6 +168,16 @@ def transaction(self, run_id, tenant_id=None, user_id=None):
if tx:
tx.flush()

@contextmanager
def read_only(self, run_id, tenant_id=None, user_id=None):
"""Snapshot reads must not contend with write paths. No lock, no writes."""
with self.session_factory() as session:
conditions = [HumanRun.run_id == run_id, HumanRun.delete_flag == "N"]
if tenant_id is not None:
conditions.extend([HumanRun.tenant_id == tenant_id, HumanRun.user_id == user_id])
run = session.scalar(select(HumanRun).where(*conditions))
yield RunTransaction(session, run, self.validator, user_id) if run else None

@contextmanager
def creation(self, run):
with self.session_factory() as session, session.no_autoflush:
Expand Down
Loading
Loading