feat(runner): unify plugin execution across agents and event processors

This commit is contained in:
RockChinQ
2026-09-10 18:04:38 +08:00
parent 8903a40c41
commit f24a7c9bb2
223 changed files with 4091 additions and 3068 deletions
@@ -9,7 +9,7 @@ import json
import quart
from .....agent.runner.errors import (
AgentRunnerError,
RunnerError,
RunnerExecutionError,
RunnerNotAuthorizedError,
RunnerNotFoundError,
@@ -39,7 +39,7 @@ def debug_stream_response(service, context, agent_uuid: str, payload: dict) -> q
code, message = 'runner_protocol_error', 'The Agent runner returned an invalid response'
elif isinstance(exc, ValueError):
code, message = 'invalid_request', str(exc)
elif isinstance(exc, AgentRunnerError):
elif isinstance(exc, RunnerError):
code, message = 'runner_error', 'The Agent runner could not complete this test'
else:
code, message = 'runner_error', 'The Agent debug execution failed'
@@ -3,7 +3,7 @@ from __future__ import annotations
import quart
from .....agent.runner.errors import (
AgentRunnerError,
RunnerError,
RunnerExecutionError,
RunnerNotAuthorizedError,
RunnerNotFoundError,
@@ -144,7 +144,7 @@ class AgentsRouterGroup(group.RouterGroup):
'runner_protocol_error',
'The Agent runner returned an invalid response',
)
except AgentRunnerError:
except RunnerError:
return self.http_status(
502,
'runner_error',
+25 -20
View File
@@ -8,13 +8,13 @@ import uuid
import typing
import sqlalchemy
from langbot_plugin.api.entities.builtin.agent_runner.delivery import DeliveryContext
from langbot_plugin.api.entities.builtin.agent_runner.event import (
from langbot_plugin.api.entities.builtin.runner.delivery import DeliveryContext
from langbot_plugin.api.entities.builtin.runner.event import (
ActorContext,
RawEventRef,
SubjectContext,
)
from langbot_plugin.api.entities.builtin.agent_runner.input import AgentInput
from langbot_plugin.api.entities.builtin.runner.input import AgentInput
from ....core import app
from ....agent.runner.config_resolver import RunnerConfigResolver
@@ -69,13 +69,13 @@ class AgentService:
except Exception as exc:
self.ap.logger.warning(f'Failed to load Agent Host tool catalog: {exc}')
event_processors = []
registry = getattr(self.ap, 'agent_runner_registry', None)
registry = getattr(self.ap, 'runner_registry', None)
if registry is not None:
event_processors = [
item.model_dump(mode='json')
for item in await registry.list_runners(
context,
component_kind='EventProcessor',
usage='event',
use_cache=False,
)
]
@@ -168,7 +168,7 @@ class AgentService:
config = agent.get('config')
if not isinstance(config, dict):
raise ValueError('Agent configuration is invalid')
_, runner_id, runner_config = RunnerConfigResolver.resolve_agent_runner_config(config)
_, runner_id, runner_config = RunnerConfigResolver.resolve_agent_config(config)
if not runner_id:
raise ValueError('Agent has no configured runner')
@@ -405,9 +405,8 @@ class AgentService:
config, runner_id, patterns = await self._prepare_event_processor(context, agent_data)
else:
config = agent_data['config'] if 'config' in agent_data else await self._get_default_agent_config(context)
config, runner_id, _ = RunnerConfigResolver.resolve_agent_runner_config(config)
if (runner_id or '').startswith('event_processor:'):
raise ValueError('EventProcessor components require an event processor instance')
config, runner_id, _ = RunnerConfigResolver.resolve_agent_config(config)
await self._validate_runner_for_agent(context, runner_id)
patterns = agent_data.get('supported_event_patterns', AGENT_DEFAULT_EVENT_PATTERNS)
new_uuid = str(uuid.uuid4())
values = {
@@ -444,13 +443,13 @@ class AgentService:
config, runner_id, patterns = await self._prepare_event_processor(context, agent_data, existing_agent)
update_data.update(config=config, component_ref=runner_id, supported_event_patterns=patterns)
if 'config' in update_data:
config, runner_id, _ = RunnerConfigResolver.resolve_agent_runner_config(update_data['config'])
config, runner_id, _ = RunnerConfigResolver.resolve_agent_config(update_data['config'])
update_data['config'] = config
else:
_, runner_id, _ = RunnerConfigResolver.resolve_agent_runner_config(existing_agent.config)
_, runner_id, _ = RunnerConfigResolver.resolve_agent_config(existing_agent.config)
update_data['component_ref'] = runner_id
if existing_agent.kind == AGENT_KIND_AGENT and (runner_id or '').startswith('event_processor:'):
raise ValueError('EventProcessor components require an event processor instance')
if existing_agent.kind == AGENT_KIND_AGENT:
await self._validate_runner_for_agent(context, runner_id)
result = await self.ap.persistence_mgr.execute_async(
scope_statement(
sqlalchemy.update(persistence_agent.Agent)
@@ -482,6 +481,12 @@ class AgentService:
raise ValueError(f'Agent {agent_uuid} not found')
await self.ap.pipeline_service.delete_pipeline(context, agent_uuid)
async def _validate_runner_for_agent(self, context, runner_id):
if runner_id:
descriptor = await self.ap.runner_registry.get(context, runner_id)
if 'agent' not in descriptor.usages:
raise ValueError('The selected Runner does not support agent usage')
async def _prepare_event_processor(self, context, data, existing=None):
"""Resolve an installed component and keep its capability declaration authoritative."""
config = copy.deepcopy(data.get('config', existing.config if existing is not None else {}))
@@ -491,17 +496,17 @@ class AgentService:
if component_ref is None and not config and not data.get('parameters'):
# An unconfigured instance cannot subscribe to or execute any events.
return {}, None, []
if not isinstance(component_ref, str) or not component_ref.startswith('event_processor:'):
raise ValueError('Select an installed EventProcessor component')
if not isinstance(component_ref, str) or not component_ref.startswith('plugin:'):
raise ValueError('Select an installed Runner component')
try:
descriptor = await self.ap.agent_runner_registry.get(context, component_ref)
descriptor = await self.ap.runner_registry.get(context, component_ref)
except Exception as exc:
from ....agent.runner.errors import RunnerNotFoundError
if isinstance(exc, RunnerNotFoundError):
raise ValueError('EventProcessor component is unavailable') from exc
raise ValueError('Runner component is unavailable') from exc
raise
if descriptor.component_kind != 'EventProcessor' or not descriptor.supported_event_patterns:
if 'event' not in descriptor.usages or not descriptor.supported_event_patterns:
raise ValueError('The component does not declare supported events')
config['runner'] = {'id': component_ref}
parameters = data.get('parameters')
@@ -583,9 +588,9 @@ class AgentService:
async def _get_default_agent_config(self, context: TenantContext) -> dict[str, typing.Any]:
runners = []
if getattr(self.ap, 'agent_runner_registry', None) is not None:
if getattr(self.ap, 'runner_registry', None) is not None:
try:
runners = await self.ap.agent_runner_registry.list_runners(context, bound_plugins=None)
runners = await self.ap.runner_registry.list_runners(context, bound_plugins=None)
except Exception as e:
if getattr(self.ap, 'logger', None):
self.ap.logger.warning(f'Failed to load plugin agent runners for default agent config: {e}')
+16 -49
View File
@@ -42,30 +42,22 @@ class PipelineService:
def _get_default_values_from_schema(
config_schema: list[dict[str, typing.Any]],
) -> dict[str, typing.Any]:
return {
item['name']: item['default']
for item in config_schema
if item.get('name') and 'default' in item
}
return {item['name']: item['default'] for item in config_schema if item.get('name') and 'default' in item}
async def get_default_pipeline_config(self, context: TenantContext) -> dict[str, typing.Any]:
from ....utils import paths as path_utils
template_path = path_utils.get_resource_path(
'templates/default-pipeline-config.json'
)
template_path = path_utils.get_resource_path('templates/default-pipeline-config.json')
with open(template_path, 'r', encoding='utf-8') as f:
config = json.load(f)
registry = getattr(self.ap, 'agent_runner_registry', None)
registry = getattr(self.ap, 'runner_registry', None)
if registry is None:
return config
try:
runners = await registry.list_runners(context, bound_plugins=None)
except Exception as exc:
self.ap.logger.warning(
f'Failed to load AgentRunner defaults for pipeline config: {exc}'
)
self.ap.logger.warning(f'Failed to load Runner defaults for pipeline config: {exc}')
return config
if not runners:
return config
@@ -75,11 +67,7 @@ class PipelineService:
runner_config = ai_config.setdefault('runner', {})
runner_config['id'] = selected.id
runner_config.setdefault('expire-time', 0)
ai_config['runner_config'] = {
selected.id: self._get_default_values_from_schema(
selected.config_schema
)
}
ai_config['runner_config'] = {selected.id: self._get_default_values_from_schema(selected.config_schema)}
return config
async def get_pipeline_metadata(self, context: TenantContext) -> list[dict]:
@@ -88,11 +76,7 @@ class PipelineService:
ai_metadata = copy.deepcopy(self.ap.pipeline_config_meta_ai)
runner_stage = next(
(
stage
for stage in ai_metadata.get('stages', [])
if stage.get('name') == 'runner'
),
(stage for stage in ai_metadata.get('stages', []) if stage.get('name') == 'runner'),
None,
)
if runner_stage:
@@ -100,25 +84,18 @@ class PipelineService:
if config_item.get('name') != 'id':
continue
try:
runner_options, runner_stages = (
await self.ap.agent_runner_registry.get_runner_metadata_for_pipeline(context)
runner_options, runner_stages = await self.ap.runner_registry.get_runner_metadata_for_pipeline(
context
)
config_item['options'] = runner_options
if runner_options and 'default' not in config_item:
config_item['default'] = runner_options[0]['name']
existing = {
stage.get('name')
for stage in ai_metadata.get('stages', [])
}
existing = {stage.get('name') for stage in ai_metadata.get('stages', [])}
ai_metadata.setdefault('stages', []).extend(
stage
for stage in runner_stages
if stage.get('name') not in existing
stage for stage in runner_stages if stage.get('name') not in existing
)
except Exception as exc:
self.ap.logger.warning(
f'Failed to load AgentRunner pipeline metadata: {exc}'
)
self.ap.logger.warning(f'Failed to load Runner pipeline metadata: {exc}')
return [
self.ap.pipeline_config_meta_trigger,
self.ap.pipeline_config_meta_safety,
@@ -187,9 +164,7 @@ class PipelineService:
async def create_pipeline(self, context: TenantContext, pipeline_data: dict, default: bool = False) -> str:
workspace_uuid = require_workspace_uuid(context)
if 'extensions_preferences' in pipeline_data:
self._validate_extension_preferences(
pipeline_data['extensions_preferences']
)
self._validate_extension_preferences(pipeline_data['extensions_preferences'])
if 'config' in pipeline_data:
RunnerConfigResolver.validate_pipeline_config(pipeline_data['config'])
# Check limitation
@@ -250,9 +225,7 @@ class PipelineService:
)
RunnerConfigResolver.validate_pipeline_config(pipeline_data['config'])
if 'extensions_preferences' in pipeline_data:
self._validate_extension_preferences(
pipeline_data['extensions_preferences']
)
self._validate_extension_preferences(pipeline_data['extensions_preferences'])
result = await self.ap.persistence_mgr.execute_async(
scope_statement(
@@ -333,9 +306,7 @@ class PipelineService:
'stages': original_pipeline.stages.copy() if original_pipeline.stages else default_stage_order.copy(),
'config': original_pipeline.config.copy() if original_pipeline.config else {},
'is_default': False,
'extensions_preferences': normalize_extension_preferences(
original_pipeline.extensions_preferences
),
'extensions_preferences': normalize_extension_preferences(original_pipeline.extensions_preferences),
}
# Insert the new pipeline
@@ -377,9 +348,7 @@ class PipelineService:
if bound_mcp_resources is not None:
extension_updates['mcp_resources'] = bound_mcp_resources
if mcp_resource_agent_read_enabled is not None:
extension_updates['mcp_resource_agent_read_enabled'] = (
mcp_resource_agent_read_enabled
)
extension_updates['mcp_resource_agent_read_enabled'] = mcp_resource_agent_read_enabled
self._validate_extension_preferences(
extension_updates,
context='Pipeline extension',
@@ -406,9 +375,7 @@ class PipelineService:
raise WorkspaceNotFoundError(f'Pipeline {pipeline_uuid} not found')
# Update extensions_preferences
extensions_preferences = normalize_extension_preferences(
pipeline.extensions_preferences
)
extensions_preferences = normalize_extension_preferences(pipeline.extensions_preferences)
extensions_preferences['enable_all_plugins'] = enable_all_plugins
extensions_preferences['enable_all_mcp_servers'] = enable_all_mcp_servers
extensions_preferences['enable_all_skills'] = enable_all_skills