use dashmap::DashMap;
use std::collections::HashMap;
use std::sync::Arc;
use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader};
use tokio::sync::{mpsc, oneshot};
use zerolaunch_plugin_api::config::Configurable;
use zerolaunch_plugin_api::{
ActionExecutor, CachedCandidateData, DataSource, KeywordInjector, KeywordOptimizer, Plugin,
ScoreBooster, SearchEngine,
};
use zerolaunch_plugin_protocol::codec::{
encode_frame, summarize_message, MAX_FRAME_SIZE, MAX_HEADER_SIZE,
};
use zerolaunch_plugin_protocol::jsonrpc::{Message, Request, Response};
use zerolaunch_plugin_protocol::messages::*;
use zerolaunch_plugin_protocol::methods::plugin as plugin_methods;
use zerolaunch_plugin_protocol::{codes, JsonRpcError, PROTOCOL_VERSION};
use crate::host_proxy::HostProxy;
use crate::logging;
use std::sync::OnceLock;
static HOST_PROXY: OnceLock<Arc<HostProxy>> = OnceLock::new();
pub fn host() -> Arc<HostProxy> {
HOST_PROXY
.get()
.expect("host() 必须在 run() 之后调用")
.clone()
}
struct IncomingRequest {
id: u64,
method: String,
params: serde_json::Value,
}
pub struct PluginApp {
pub(crate) plugin: Arc<dyn Plugin>,
pub(crate) data_sources: Vec<Arc<dyn DataSource>>,
pub(crate) executors: Vec<Arc<dyn ActionExecutor>>,
pub(crate) search_engines: Vec<Arc<dyn SearchEngine>>,
pub(crate) score_boosters: Vec<Arc<dyn ScoreBooster>>,
pub(crate) keyword_optimizers: Vec<Arc<dyn KeywordOptimizer>>,
pub(crate) keyword_injectors: Vec<Arc<dyn KeywordInjector>>,
by_id: HashMap<String, ComponentEntry>,
}
enum ComponentEntry {
Plugin(Arc<dyn Plugin>),
DataSource(Arc<dyn DataSource>),
Executor(Arc<dyn ActionExecutor>),
SearchEngine(Arc<dyn SearchEngine>),
ScoreBooster(Arc<dyn ScoreBooster>),
KeywordOptimizer(Arc<dyn KeywordOptimizer>),
KeywordInjector(Arc<dyn KeywordInjector>),
}
impl PluginApp {
pub fn new(plugin: impl Plugin + 'static) -> Self {
let plugin = Arc::new(plugin);
let mut by_id = HashMap::new();
by_id.insert(
plugin.component_id().to_string(),
ComponentEntry::Plugin(plugin.clone()),
);
Self {
plugin,
data_sources: Vec::new(),
executors: Vec::new(),
search_engines: Vec::new(),
score_boosters: Vec::new(),
keyword_optimizers: Vec::new(),
keyword_injectors: Vec::new(),
by_id,
}
}
pub fn with_data_source(mut self, ds: impl DataSource + 'static) -> Self {
let ds = Arc::new(ds);
let component_id = ds.component_id().to_string();
assert!(
self.by_id
.insert(component_id, ComponentEntry::DataSource(ds.clone()))
.is_none(),
"DataSource 组件 id 重复:{}",
ds.component_id()
);
self.data_sources.push(ds);
self
}
pub fn with_executor(mut self, ex: impl ActionExecutor + 'static) -> Self {
let ex = Arc::new(ex);
let component_id = ex.component_id().to_string();
assert!(
self.by_id
.insert(component_id, ComponentEntry::Executor(ex.clone()))
.is_none(),
"ActionExecutor 组件 id 重复:{}",
ex.component_id()
);
self.executors.push(ex);
self
}
pub fn with_search_engine(mut self, engine: impl SearchEngine + 'static) -> Self {
let engine = Arc::new(engine);
let component_id = engine.component_id().to_string();
assert!(
self.by_id
.insert(component_id, ComponentEntry::SearchEngine(engine.clone()))
.is_none(),
"SearchEngine 组件 id 重复:{}",
engine.component_id()
);
self.search_engines.push(engine);
self
}
pub fn with_score_booster(mut self, booster: impl ScoreBooster + 'static) -> Self {
let booster = Arc::new(booster);
let component_id = booster.component_id().to_string();
assert!(
self.by_id
.insert(component_id, ComponentEntry::ScoreBooster(booster.clone()))
.is_none(),
"ScoreBooster 组件 id 重复:{}",
booster.component_id()
);
self.score_boosters.push(booster);
self
}
pub fn with_keyword_optimizer(mut self, optimizer: impl KeywordOptimizer + 'static) -> Self {
let optimizer = Arc::new(optimizer);
let component_id = optimizer.component_id().to_string();
assert!(
self.by_id
.insert(
component_id,
ComponentEntry::KeywordOptimizer(optimizer.clone()),
)
.is_none(),
"KeywordOptimizer 组件 id 重复:{}",
optimizer.component_id()
);
self.keyword_optimizers.push(optimizer);
self
}
pub fn with_keyword_injector(mut self, injector: impl KeywordInjector + 'static) -> Self {
let injector = Arc::new(injector);
let component_id = injector.component_id().to_string();
assert!(
self.by_id
.insert(
component_id,
ComponentEntry::KeywordInjector(injector.clone()),
)
.is_none(),
"KeywordInjector 组件 id 重复:{}",
injector.component_id()
);
self.keyword_injectors.push(injector);
self
}
pub fn run(self) {
let rt = tokio::runtime::Builder::new_multi_thread()
.worker_threads(2)
.enable_all()
.build()
.expect("failed to build tokio runtime");
rt.block_on(async move {
run_async(self).await;
});
}
}
pub fn run(plugin: impl Plugin + 'static) {
PluginApp::new(plugin).run()
}
async fn run_async(mut app: PluginApp) {
let mut log_rx = logging::init_logging();
let stdin = tokio::io::stdin();
let stdout = tokio::io::stdout();
let (request_tx, mut request_rx) = mpsc::channel::<IncomingRequest>(64);
let (outbound_tx, mut outbound_rx) = mpsc::channel::<Vec<u8>>(64);
let pending: Arc<DashMap<u64, oneshot::Sender<Result<serde_json::Value, JsonRpcError>>>> =
Arc::new(DashMap::new());
let host_proxy = Arc::new(HostProxy::new(pending.clone(), outbound_tx.clone()));
let _ = HOST_PROXY.set(host_proxy);
let hp_for_logs = HOST_PROXY.get().expect("HOST_PROXY 已设置").clone();
let mut plugin_context: Option<zerolaunch_plugin_api::PluginContext> = None;
tokio::spawn(async move {
while let Some(entry) = log_rx.recv().await {
hp_for_logs.log_no_wait(&entry.level, &entry.message);
}
});
let pending_r = pending.clone();
let request_tx_clone = request_tx.clone();
let mut read_handle = tokio::spawn(async move {
let reader = BufReader::new(stdin);
let mut stdin = reader;
while let Ok(body) = read_frame(&mut stdin).await {
tracing::debug!(
"收到宿主原始帧 ({} bytes): {:?}",
body.len(),
String::from_utf8_lossy(&body)
);
let msg: Message = match serde_json::from_slice(&body) {
Ok(m) => m,
Err(e) => {
tracing::debug!(
"解析宿主消息失败: {}; raw frame: {:?}",
e,
String::from_utf8_lossy(&body)
);
continue;
}
};
match msg {
Message::Response(resp) => {
if let Some((_, tx)) = pending_r.remove(&resp.id) {
let result = if let Some(err) = resp.error {
Err(err)
} else {
Ok(resp.result.unwrap_or(serde_json::Value::Null))
};
let _ = tx.send(result);
} else {
tracing::warn!(
"收到未知的响应 id={},可能已超时或被取消: {:?}",
resp.id,
resp
);
}
}
Message::Request(req) => {
let ret = request_tx_clone
.send(IncomingRequest {
id: req.id,
method: req.method,
params: req.params,
})
.await;
if ret.is_err() {
tracing::warn!(
"无法将请求发送到 dispatch task,可能 dispatch task 已退出: {:?}",
ret
);
}
}
Message::Notification(_) => {
tracing::trace!("忽略通知");
}
}
}
});
let outbound_dispatch = outbound_tx.clone();
let mut dispatch_handle = tokio::spawn(async move {
while let Some(incoming) = request_rx.recv().await {
let req = Request::new(incoming.id, &incoming.method, incoming.params);
tracing::debug!(
"dispatch 收到请求: {:?}",
summarize_message(&Message::Request(req.clone()))
);
let result = handle_request(&mut app, &req, &mut plugin_context).await;
if let Ok(payload) = serde_json::to_vec(&result) {
if outbound_dispatch.send(payload).await.is_err() {
tracing::warn!(
"dispatch 发送响应失败(outbound channel 关闭),dispatch task 退出"
);
break;
}
} else {
tracing::error!("dispatch 序列化响应失败,请求被丢弃");
}
}
tracing::warn!("dispatch task 退出(request channel 关闭或发送失败)");
});
let write_handle = tokio::spawn(async move {
let mut writer = stdout;
while let Some(payload) = outbound_rx.recv().await {
match serde_json::from_slice::<Message>(&payload) {
Ok(msg) => tracing::debug!("发送给宿主的消息: {:?}", summarize_message(&msg)),
Err(e) => tracing::debug!(
"发送给宿主的原始负载 ({} bytes, 解析失败: {}): {:?}",
payload.len(),
e,
String::from_utf8_lossy(&payload)
),
}
let frame = encode_frame(&payload);
if writer.write_all(&frame).await.is_err() {
break;
}
if writer.flush().await.is_err() {
break;
}
}
});
tokio::select! {
_ = &mut read_handle => {
drop(request_tx);
write_handle.abort();
let _ = (&mut dispatch_handle).await;
}
r = &mut dispatch_handle => {
if r.is_err() {
tracing::error!("dispatch task 因 panic 退出,插件进程终止(宿主将按 auto_restart 重启)");
}
}
}
}
async fn read_frame<R: tokio::io::AsyncBufRead + Unpin>(reader: &mut R) -> Result<Vec<u8>, String> {
let mut content_length: Option<usize> = None;
let mut total_header_len = 0usize;
loop {
let mut line = String::new();
let n = reader
.read_line(&mut line)
.await
.map_err(|e| format!("read error: {}", e))?;
if n == 0 {
return Err("transport closed".into());
}
total_header_len += n;
if total_header_len > MAX_HEADER_SIZE {
return Err("header too long".into());
}
let trimmed = line.trim();
if trimmed.is_empty() {
break;
}
if let Some(value) = trimmed.strip_prefix("Content-Length:") {
content_length = Some(
value
.trim()
.parse::<usize>()
.map_err(|e| format!("bad Content-Length: {}", e))?,
);
}
}
let len = content_length.ok_or("missing Content-Length")?;
if len > MAX_FRAME_SIZE {
return Err(format!("Content-Length too large: {}", len));
}
let mut body = vec![0u8; len];
reader
.read_exact(&mut body)
.await
.map_err(|e| format!("read body: {}", e))?;
Ok(body)
}
async fn handle_request(
app: &mut PluginApp,
req: &Request,
plugin_ctx: &mut Option<zerolaunch_plugin_api::PluginContext>,
) -> Message {
let id = req.id;
let result = dispatch(app, &req.method, &req.params, plugin_ctx).await;
match result {
Ok(value) => Message::Response(Response::ok(id, value)),
Err(err) => Message::Response(Response::err(id, err)),
}
}
async fn dispatch(
app: &mut PluginApp,
method: &str,
params: &serde_json::Value,
plugin_ctx: &mut Option<zerolaunch_plugin_api::PluginContext>,
) -> Result<serde_json::Value, JsonRpcError> {
match method {
plugin_methods::INITIALIZE => {
let p: InitializeParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
crate::set_plugin_id(&p.plugin_id);
*plugin_ctx = Some(zerolaunch_plugin_api::PluginContext {
trace_id: "init".into(),
query_id: None,
plugin_id: Some(p.plugin_id),
query_revision_gate: None,
query_channel: zerolaunch_plugin_api::QueryChannel::Ui,
locale: p.locale,
});
let result = InitializeResult {
protocol_version: PROTOCOL_VERSION.to_string(),
};
Ok(serde_json::to_value(result).unwrap_or_default())
}
plugin_methods::GET_COMPONENTS => {
let mut components = vec![ComponentDescriptor {
component_id: app.plugin.component_id().to_string(),
component_name: app.plugin.component_name().to_string(),
component_description: app.plugin.component_description().to_string(),
component_type: app.plugin.component_type(),
kind: ComponentKind::Plugin,
priority: app.plugin.priority(),
}];
for ds in &app.data_sources {
components.push(ComponentDescriptor {
component_id: ds.component_id().to_string(),
component_name: ds.component_name().to_string(),
component_description: ds.component_description().to_string(),
component_type: ds.component_type(),
kind: ComponentKind::DataSource,
priority: ds.priority(),
});
}
for ex in &app.executors {
components.push(ComponentDescriptor {
component_id: ex.component_id().to_string(),
component_name: ex.component_name().to_string(),
component_description: ex.component_description().to_string(),
component_type: ex.component_type(),
kind: ComponentKind::ActionExecutor {
target_types: ex.supported_target_types(),
},
priority: ex.priority(),
});
}
for engine in &app.search_engines {
components.push(ComponentDescriptor {
component_id: engine.component_id().to_string(),
component_name: engine.component_name().to_string(),
component_description: engine.component_description().to_string(),
component_type: engine.component_type(),
kind: ComponentKind::SearchEngine,
priority: engine.priority(),
});
}
for booster in &app.score_boosters {
components.push(ComponentDescriptor {
component_id: booster.component_id().to_string(),
component_name: booster.component_name().to_string(),
component_description: booster.component_description().to_string(),
component_type: booster.component_type(),
kind: ComponentKind::ScoreBooster,
priority: booster.priority(),
});
}
for optimizer in &app.keyword_optimizers {
components.push(ComponentDescriptor {
component_id: optimizer.component_id().to_string(),
component_name: optimizer.component_name().to_string(),
component_description: optimizer.component_description().to_string(),
component_type: optimizer.component_type(),
kind: ComponentKind::KeywordOptimizer,
priority: optimizer.priority(),
});
}
for injector in &app.keyword_injectors {
components.push(ComponentDescriptor {
component_id: injector.component_id().to_string(),
component_name: injector.component_name().to_string(),
component_description: injector.component_description().to_string(),
component_type: injector.component_type(),
kind: ComponentKind::KeywordInjector,
priority: injector.priority(),
});
}
Ok(serde_json::to_value(components).unwrap_or_default())
}
plugin_methods::GET_SETTINGS_SCHEMA => {
let p: GetSettingsSchemaParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let conf = find_configurable(app, &p.component_id)?;
Ok(serde_json::to_value(conf.setting_schema()).unwrap_or(serde_json::Value::Null))
}
plugin_methods::GET_SETTINGS => {
let p: GetSettingsParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let conf = find_configurable(app, &p.component_id)?;
Ok(conf.get_settings())
}
plugin_methods::GET_DEFAULT_ENABLED => {
let p: GetDefaultEnabledParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let conf = find_configurable(app, &p.component_id)?;
Ok(serde_json::to_value(conf.default_enabled()).unwrap_or_default())
}
plugin_methods::APPLY_SETTINGS => {
let p: ApplySettingsParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let conf = find_configurable(app, &p.component_id)?;
conf.apply_settings(p.settings)
.await
.map_err(|e| JsonRpcError::new(codes::PLUGIN_ERROR, e.to_string()))?;
Ok(serde_json::Value::Null)
}
plugin_methods::VALIDATE_SETTINGS => {
let p: ValidateSettingsParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let conf = find_configurable(app, &p.component_id)?;
let result = match conf.validate_settings(&p.settings).await {
Ok(()) => ValidateSettingsResult { error: None },
Err(e) => ValidateSettingsResult {
error: Some(e.to_string()),
},
};
Ok(serde_json::to_value(result).unwrap_or_default())
}
plugin_methods::CONFIG_ACTIONS => {
let p: ConfigActionsParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let conf = find_configurable(app, &p.component_id)?;
Ok(serde_json::to_value(conf.config_actions()).unwrap_or_default())
}
plugin_methods::EXECUTE_CONFIG_ACTION => {
let p: ExecuteConfigActionParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let conf = find_configurable(app, &p.component_id)?;
conf.execute_config_action(&p.action, &p.params)
.await
.map_err(|e| JsonRpcError::new(codes::PLUGIN_ERROR, e.to_string()))
}
plugin_methods::QUERY => {
let p: QueryParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let response = app
.plugin
.query(&p.ctx, &p.query)
.await
.map_err(|e| JsonRpcError::new(codes::PLUGIN_ERROR, e.to_string()))?;
Ok(serde_json::to_value(response).unwrap_or_default())
}
plugin_methods::MATCH_QUERY => {
let p: MatchQueryParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
Ok(serde_json::Value::Bool(
app.plugin
.match_query(&p.raw_query, &p.trigger_keywords)
.await,
))
}
plugin_methods::EXECUTE_ACTION => {
let p: ExecuteActionParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
app.plugin
.execute_action(&p.ctx, &p.action_id, p.payload)
.await
.map_err(|e| JsonRpcError::new(codes::PLUGIN_ERROR, e.to_string()))?;
Ok(serde_json::Value::Null)
}
plugin_methods::INIT => {
let p: InitParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
app.plugin
.init(&p.ctx, None)
.await
.map_err(|e| JsonRpcError::new(codes::PLUGIN_ERROR, e.to_string()))?;
Ok(serde_json::Value::Null)
}
plugin_methods::INTERACTION_POLICY => {
let _p: InteractionPolicyParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
Ok(serde_json::to_value(app.plugin.interaction_policy()).unwrap_or_default())
}
plugin_methods::FETCH_CANDIDATES => {
let p: FetchCandidatesParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let ds = find_data_source(app, &p.component_id)?;
let cache = ds.fetch_candidates().await;
Ok(serde_json::to_value(FetchCandidatesResult {
candidates: cache.get_candidates().clone(),
})
.unwrap_or_default())
}
plugin_methods::SUPPORTED_TARGET_TYPES => {
let p: SupportedTargetTypesParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let ex = find_executor(app, &p.component_id)?;
Ok(serde_json::to_value(ex.supported_target_types()).unwrap_or_default())
}
plugin_methods::SUPPORTED_ACTIONS => {
let p: SupportedActionsParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let ex = find_executor(app, &p.component_id)?;
Ok(serde_json::to_value(ex.supported_actions()).unwrap_or_default())
}
plugin_methods::EXECUTOR_EXECUTE => {
let p: ExecutorExecuteParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let ex = find_executor(app, &p.component_id)?;
let result = match ex.execute(&p.execution_ctx, &p.action_id).await {
Ok(()) => ExecutorExecuteResult { error: None },
Err(e) => ExecutorExecuteResult {
error: Some(e.to_string()),
},
};
Ok(serde_json::to_value(result).unwrap_or_default())
}
plugin_methods::CALCULATE_SCORES => {
let p: CalculateScoresParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let engine = find_search_engine(app, &p.component_id)?;
let cache = CachedCandidateData::from_data(p.candidates);
let scored = engine.calculate_scores(&cache, &p.query).await;
Ok(serde_json::to_value(scored).unwrap_or_default())
}
plugin_methods::BOOSTER_BOOST => {
let p: BoosterBoostParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let booster = find_score_booster(app, &p.component_id)?;
let cache = CachedCandidateData::from_data(p.candidates);
let mut scored = p.scored;
booster.boost(&mut scored, &cache, &p.query).await;
Ok(serde_json::to_value(scored).unwrap_or_default())
}
plugin_methods::BOOSTER_RECORD => {
let p: BoosterRecordParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let booster = find_score_booster(app, &p.component_id)?;
let cache = CachedCandidateData::from_data(p.candidates);
booster.record(p.candidate_id, &cache, &p.query).await;
Ok(serde_json::Value::Null)
}
plugin_methods::KEYWORD_OPTIMIZER_INFO => {
let p: KeywordOptimizerInfoParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let optimizer = find_keyword_optimizer(app, &p.component_id)?;
Ok(serde_json::to_value(KeywordOptimizerInfo {
input_source: optimizer.input_source(),
priority: optimizer.get_priority(),
})
.unwrap_or_default())
}
plugin_methods::KEYWORD_OPTIMIZE => {
let p: KeywordOptimizeParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let optimizer = find_keyword_optimizer(app, &p.component_id)?;
Ok(serde_json::to_value(optimizer.optimize(&p.keyword).await).unwrap_or_default())
}
plugin_methods::KEYWORD_INJECT => {
let p: KeywordInjectParams = serde_json::from_value(params.clone())
.map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
let injector = find_keyword_injector(app, &p.component_id)?;
Ok(
serde_json::to_value(injector.inject_keywords(&p.candidate).await)
.unwrap_or_default(),
)
}
_ => Err(JsonRpcError::new(
codes::METHOD_NOT_FOUND,
format!("method not found: {}", method),
)),
}
}
fn find_configurable<'a>(
app: &'a PluginApp,
component_id: &str,
) -> Result<&'a dyn Configurable, JsonRpcError> {
let entry = app.by_id.get(component_id).ok_or_else(|| {
JsonRpcError::new(
codes::METHOD_NOT_FOUND,
format!("component not found: {component_id}"),
)
})?;
Ok(match entry {
ComponentEntry::Plugin(p) => p.as_ref() as &dyn Configurable,
ComponentEntry::DataSource(ds) => ds.as_ref() as &dyn Configurable,
ComponentEntry::Executor(ex) => ex.as_ref() as &dyn Configurable,
ComponentEntry::SearchEngine(engine) => engine.as_ref() as &dyn Configurable,
ComponentEntry::ScoreBooster(booster) => booster.as_ref() as &dyn Configurable,
ComponentEntry::KeywordOptimizer(optimizer) => optimizer.as_ref() as &dyn Configurable,
ComponentEntry::KeywordInjector(injector) => injector.as_ref() as &dyn Configurable,
})
}
fn find_data_source<'a>(
app: &'a PluginApp,
component_id: &str,
) -> Result<&'a dyn DataSource, JsonRpcError> {
match app.by_id.get(component_id) {
Some(ComponentEntry::DataSource(ds)) => Ok(ds.as_ref()),
Some(_) => Err(JsonRpcError::new(
codes::METHOD_NOT_FOUND,
format!("component is not a data source: {component_id}"),
)),
None => Err(JsonRpcError::new(
codes::METHOD_NOT_FOUND,
format!("data source not found: {component_id}"),
)),
}
}
fn find_executor<'a>(
app: &'a PluginApp,
component_id: &str,
) -> Result<&'a dyn ActionExecutor, JsonRpcError> {
match app.by_id.get(component_id) {
Some(ComponentEntry::Executor(ex)) => Ok(ex.as_ref()),
Some(_) => Err(JsonRpcError::new(
codes::METHOD_NOT_FOUND,
format!("component is not an executor: {component_id}"),
)),
None => Err(JsonRpcError::new(
codes::METHOD_NOT_FOUND,
format!("executor not found: {component_id}"),
)),
}
}
fn find_search_engine<'a>(
app: &'a PluginApp,
component_id: &str,
) -> Result<&'a dyn SearchEngine, JsonRpcError> {
match app.by_id.get(component_id) {
Some(ComponentEntry::SearchEngine(engine)) => Ok(engine.as_ref()),
Some(_) => Err(JsonRpcError::new(
codes::METHOD_NOT_FOUND,
format!("component is not a search engine: {component_id}"),
)),
None => Err(JsonRpcError::new(
codes::METHOD_NOT_FOUND,
format!("search engine not found: {component_id}"),
)),
}
}
fn find_score_booster<'a>(
app: &'a PluginApp,
component_id: &str,
) -> Result<&'a dyn ScoreBooster, JsonRpcError> {
match app.by_id.get(component_id) {
Some(ComponentEntry::ScoreBooster(booster)) => Ok(booster.as_ref()),
Some(_) => Err(JsonRpcError::new(
codes::METHOD_NOT_FOUND,
format!("component is not a score booster: {component_id}"),
)),
None => Err(JsonRpcError::new(
codes::METHOD_NOT_FOUND,
format!("score booster not found: {component_id}"),
)),
}
}
fn find_keyword_optimizer<'a>(
app: &'a PluginApp,
component_id: &str,
) -> Result<&'a dyn KeywordOptimizer, JsonRpcError> {
match app.by_id.get(component_id) {
Some(ComponentEntry::KeywordOptimizer(optimizer)) => Ok(optimizer.as_ref()),
Some(_) => Err(JsonRpcError::new(
codes::METHOD_NOT_FOUND,
format!("component is not a keyword optimizer: {component_id}"),
)),
None => Err(JsonRpcError::new(
codes::METHOD_NOT_FOUND,
format!("keyword optimizer not found: {component_id}"),
)),
}
}
fn find_keyword_injector<'a>(
app: &'a PluginApp,
component_id: &str,
) -> Result<&'a dyn KeywordInjector, JsonRpcError> {
match app.by_id.get(component_id) {
Some(ComponentEntry::KeywordInjector(injector)) => Ok(injector.as_ref()),
Some(_) => Err(JsonRpcError::new(
codes::METHOD_NOT_FOUND,
format!("component is not a keyword injector: {component_id}"),
)),
None => Err(JsonRpcError::new(
codes::METHOD_NOT_FOUND,
format!("keyword injector not found: {component_id}"),
)),
}
}