1use dashmap::DashMap;
14use std::collections::HashMap;
15use std::sync::Arc;
16use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader};
17use tokio::sync::{mpsc, oneshot};
18
19use zerolaunch_plugin_api::config::Configurable;
20use zerolaunch_plugin_api::{
21 ActionExecutor, CachedCandidateData, DataSource, KeywordInjector, KeywordOptimizer, Plugin,
22 ScoreBooster, SearchEngine,
23};
24use zerolaunch_plugin_protocol::codec::{
25 encode_frame, summarize_message, MAX_FRAME_SIZE, MAX_HEADER_SIZE,
26};
27use zerolaunch_plugin_protocol::jsonrpc::{Message, Request, Response};
28use zerolaunch_plugin_protocol::messages::*;
29use zerolaunch_plugin_protocol::methods::plugin as plugin_methods;
30use zerolaunch_plugin_protocol::{codes, JsonRpcError, PROTOCOL_VERSION};
31
32use crate::host_proxy::HostProxy;
33use crate::logging;
34
35use std::sync::OnceLock;
36
37static HOST_PROXY: OnceLock<Arc<HostProxy>> = OnceLock::new();
39
40pub fn host() -> Arc<HostProxy> {
43 HOST_PROXY
44 .get()
45 .expect("host() 必须在 run() 之后调用")
46 .clone()
47}
48
49struct IncomingRequest {
51 id: u64,
52 method: String,
53 params: serde_json::Value,
54}
55
56pub struct PluginApp {
62 pub(crate) plugin: Arc<dyn Plugin>,
64 pub(crate) data_sources: Vec<Arc<dyn DataSource>>,
65 pub(crate) executors: Vec<Arc<dyn ActionExecutor>>,
66 pub(crate) search_engines: Vec<Arc<dyn SearchEngine>>,
68 pub(crate) score_boosters: Vec<Arc<dyn ScoreBooster>>,
70 pub(crate) keyword_optimizers: Vec<Arc<dyn KeywordOptimizer>>,
72 pub(crate) keyword_injectors: Vec<Arc<dyn KeywordInjector>>,
74 by_id: HashMap<String, ComponentEntry>,
77}
78
79enum ComponentEntry {
81 Plugin(Arc<dyn Plugin>),
82 DataSource(Arc<dyn DataSource>),
83 Executor(Arc<dyn ActionExecutor>),
84 SearchEngine(Arc<dyn SearchEngine>),
86 ScoreBooster(Arc<dyn ScoreBooster>),
88 KeywordOptimizer(Arc<dyn KeywordOptimizer>),
90 KeywordInjector(Arc<dyn KeywordInjector>),
92}
93
94impl PluginApp {
95 pub fn new(plugin: impl Plugin + 'static) -> Self {
97 let plugin = Arc::new(plugin);
98 let mut by_id = HashMap::new();
99 by_id.insert(
100 plugin.component_id().to_string(),
101 ComponentEntry::Plugin(plugin.clone()),
102 );
103 Self {
104 plugin,
105 data_sources: Vec::new(),
106 executors: Vec::new(),
107 search_engines: Vec::new(),
108 score_boosters: Vec::new(),
109 keyword_optimizers: Vec::new(),
110 keyword_injectors: Vec::new(),
111 by_id,
112 }
113 }
114
115 pub fn with_data_source(mut self, ds: impl DataSource + 'static) -> Self {
118 let ds = Arc::new(ds);
119 let component_id = ds.component_id().to_string();
120 assert!(
121 self.by_id
122 .insert(component_id, ComponentEntry::DataSource(ds.clone()))
123 .is_none(),
124 "DataSource 组件 id 重复:{}",
125 ds.component_id()
126 );
127 self.data_sources.push(ds);
128 self
129 }
130
131 pub fn with_executor(mut self, ex: impl ActionExecutor + 'static) -> Self {
134 let ex = Arc::new(ex);
135 let component_id = ex.component_id().to_string();
136 assert!(
137 self.by_id
138 .insert(component_id, ComponentEntry::Executor(ex.clone()))
139 .is_none(),
140 "ActionExecutor 组件 id 重复:{}",
141 ex.component_id()
142 );
143 self.executors.push(ex);
144 self
145 }
146 pub fn with_search_engine(mut self, engine: impl SearchEngine + 'static) -> Self {
149 let engine = Arc::new(engine);
150 let component_id = engine.component_id().to_string();
151 assert!(
152 self.by_id
153 .insert(component_id, ComponentEntry::SearchEngine(engine.clone()))
154 .is_none(),
155 "SearchEngine 组件 id 重复:{}",
156 engine.component_id()
157 );
158 self.search_engines.push(engine);
159 self
160 }
161
162 pub fn with_score_booster(mut self, booster: impl ScoreBooster + 'static) -> Self {
165 let booster = Arc::new(booster);
166 let component_id = booster.component_id().to_string();
167 assert!(
168 self.by_id
169 .insert(component_id, ComponentEntry::ScoreBooster(booster.clone()))
170 .is_none(),
171 "ScoreBooster 组件 id 重复:{}",
172 booster.component_id()
173 );
174 self.score_boosters.push(booster);
175 self
176 }
177
178 pub fn with_keyword_optimizer(mut self, optimizer: impl KeywordOptimizer + 'static) -> Self {
181 let optimizer = Arc::new(optimizer);
182 let component_id = optimizer.component_id().to_string();
183 assert!(
184 self.by_id
185 .insert(
186 component_id,
187 ComponentEntry::KeywordOptimizer(optimizer.clone()),
188 )
189 .is_none(),
190 "KeywordOptimizer 组件 id 重复:{}",
191 optimizer.component_id()
192 );
193 self.keyword_optimizers.push(optimizer);
194 self
195 }
196
197 pub fn with_keyword_injector(mut self, injector: impl KeywordInjector + 'static) -> Self {
200 let injector = Arc::new(injector);
201 let component_id = injector.component_id().to_string();
202 assert!(
203 self.by_id
204 .insert(
205 component_id,
206 ComponentEntry::KeywordInjector(injector.clone()),
207 )
208 .is_none(),
209 "KeywordInjector 组件 id 重复:{}",
210 injector.component_id()
211 );
212 self.keyword_injectors.push(injector);
213 self
214 }
215
216 pub fn run(self) {
218 let rt = tokio::runtime::Builder::new_multi_thread()
219 .worker_threads(2)
220 .enable_all()
221 .build()
222 .expect("failed to build tokio runtime");
223
224 rt.block_on(async move {
225 run_async(self).await;
226 });
227 }
228}
229
230pub fn run(plugin: impl Plugin + 'static) {
233 PluginApp::new(plugin).run()
234}
235
236async fn run_async(mut app: PluginApp) {
237 let mut log_rx = logging::init_logging();
239 let stdin = tokio::io::stdin();
240 let stdout = tokio::io::stdout();
241
242 let (request_tx, mut request_rx) = mpsc::channel::<IncomingRequest>(64);
244 let (outbound_tx, mut outbound_rx) = mpsc::channel::<Vec<u8>>(64);
245 let pending: Arc<DashMap<u64, oneshot::Sender<Result<serde_json::Value, JsonRpcError>>>> =
246 Arc::new(DashMap::new());
247 let host_proxy = Arc::new(HostProxy::new(pending.clone(), outbound_tx.clone()));
249 let _ = HOST_PROXY.set(host_proxy);
250 let hp_for_logs = HOST_PROXY.get().expect("HOST_PROXY 已设置").clone();
251
252 let mut plugin_context: Option<zerolaunch_plugin_api::PluginContext> = None;
254
255 tokio::spawn(async move {
257 while let Some(entry) = log_rx.recv().await {
258 hp_for_logs.log_no_wait(&entry.level, &entry.message);
259 }
260 });
261
262 let pending_r = pending.clone();
264 let request_tx_clone = request_tx.clone();
265 let mut read_handle = tokio::spawn(async move {
266 let reader = BufReader::new(stdin);
267 let mut stdin = reader;
268 while let Ok(body) = read_frame(&mut stdin).await {
269 tracing::debug!(
270 "收到宿主原始帧 ({} bytes): {:?}",
271 body.len(),
272 String::from_utf8_lossy(&body)
273 );
274 let msg: Message = match serde_json::from_slice(&body) {
275 Ok(m) => m,
276 Err(e) => {
277 tracing::debug!(
278 "解析宿主消息失败: {}; raw frame: {:?}",
279 e,
280 String::from_utf8_lossy(&body)
281 );
282 continue;
283 }
284 };
285 match msg {
286 Message::Response(resp) => {
287 if let Some((_, tx)) = pending_r.remove(&resp.id) {
288 let result = if let Some(err) = resp.error {
292 Err(err)
293 } else {
294 Ok(resp.result.unwrap_or(serde_json::Value::Null))
295 };
296 let _ = tx.send(result);
297 } else {
298 tracing::warn!(
300 "收到未知的响应 id={},可能已超时或被取消: {:?}",
301 resp.id,
302 resp
303 );
304 }
305 }
306 Message::Request(req) => {
307 let ret = request_tx_clone
308 .send(IncomingRequest {
309 id: req.id,
310 method: req.method,
311 params: req.params,
312 })
313 .await;
314 if ret.is_err() {
316 tracing::warn!(
317 "无法将请求发送到 dispatch task,可能 dispatch task 已退出: {:?}",
318 ret
319 );
320 }
321 }
322 Message::Notification(_) => {
323 tracing::trace!("忽略通知");
324 }
325 }
326 }
327 });
328
329 let outbound_dispatch = outbound_tx.clone();
331 let mut dispatch_handle = tokio::spawn(async move {
332 while let Some(incoming) = request_rx.recv().await {
333 let req = Request::new(incoming.id, &incoming.method, incoming.params);
334 tracing::debug!(
335 "dispatch 收到请求: {:?}",
336 summarize_message(&Message::Request(req.clone()))
337 );
338 let result = handle_request(&mut app, &req, &mut plugin_context).await;
340 if let Ok(payload) = serde_json::to_vec(&result) {
341 if outbound_dispatch.send(payload).await.is_err() {
342 tracing::warn!(
343 "dispatch 发送响应失败(outbound channel 关闭),dispatch task 退出"
344 );
345 break;
346 }
347 } else {
348 tracing::error!("dispatch 序列化响应失败,请求被丢弃");
349 }
350 }
351 tracing::warn!("dispatch task 退出(request channel 关闭或发送失败)");
352 });
353
354 let write_handle = tokio::spawn(async move {
356 let mut writer = stdout;
357 while let Some(payload) = outbound_rx.recv().await {
358 match serde_json::from_slice::<Message>(&payload) {
359 Ok(msg) => tracing::debug!("发送给宿主的消息: {:?}", summarize_message(&msg)),
360 Err(e) => tracing::debug!(
361 "发送给宿主的原始负载 ({} bytes, 解析失败: {}): {:?}",
362 payload.len(),
363 e,
364 String::from_utf8_lossy(&payload)
365 ),
366 }
367 let frame = encode_frame(&payload);
368 if writer.write_all(&frame).await.is_err() {
369 break;
370 }
371 if writer.flush().await.is_err() {
372 break;
373 }
374 }
375 });
376
377 tokio::select! {
379 _ = &mut read_handle => {
381 drop(request_tx);
384 write_handle.abort();
387 let _ = (&mut dispatch_handle).await;
388 }
389 r = &mut dispatch_handle => {
391 if r.is_err() {
392 tracing::error!("dispatch task 因 panic 退出,插件进程终止(宿主将按 auto_restart 重启)");
393 }
394 }
395 }
396}
397
398async fn read_frame<R: tokio::io::AsyncBufRead + Unpin>(reader: &mut R) -> Result<Vec<u8>, String> {
401 let mut content_length: Option<usize> = None;
402 let mut total_header_len = 0usize;
403 loop {
404 let mut line = String::new();
405 let n = reader
406 .read_line(&mut line)
407 .await
408 .map_err(|e| format!("read error: {}", e))?;
409 if n == 0 {
410 return Err("transport closed".into());
411 }
412 total_header_len += n;
413 if total_header_len > MAX_HEADER_SIZE {
414 return Err("header too long".into());
415 }
416 let trimmed = line.trim();
417 if trimmed.is_empty() {
418 break;
419 }
420 if let Some(value) = trimmed.strip_prefix("Content-Length:") {
421 content_length = Some(
422 value
423 .trim()
424 .parse::<usize>()
425 .map_err(|e| format!("bad Content-Length: {}", e))?,
426 );
427 }
428 }
429 let len = content_length.ok_or("missing Content-Length")?;
430 if len > MAX_FRAME_SIZE {
431 return Err(format!("Content-Length too large: {}", len));
432 }
433 let mut body = vec![0u8; len];
434 reader
435 .read_exact(&mut body)
436 .await
437 .map_err(|e| format!("read body: {}", e))?;
438 Ok(body)
439}
440
441async fn handle_request(
443 app: &mut PluginApp,
444 req: &Request,
445 plugin_ctx: &mut Option<zerolaunch_plugin_api::PluginContext>,
446) -> Message {
447 let id = req.id;
448 let result = dispatch(app, &req.method, &req.params, plugin_ctx).await;
449 match result {
450 Ok(value) => Message::Response(Response::ok(id, value)),
451 Err(err) => Message::Response(Response::err(id, err)),
452 }
453}
454
455async fn dispatch(
456 app: &mut PluginApp,
457 method: &str,
458 params: &serde_json::Value,
459 plugin_ctx: &mut Option<zerolaunch_plugin_api::PluginContext>,
460) -> Result<serde_json::Value, JsonRpcError> {
461 match method {
462 plugin_methods::INITIALIZE => {
464 let p: InitializeParams = serde_json::from_value(params.clone())
465 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
466 crate::set_plugin_id(&p.plugin_id);
468 *plugin_ctx = Some(zerolaunch_plugin_api::PluginContext {
469 trace_id: "init".into(),
470 query_id: None,
471 plugin_id: Some(p.plugin_id),
472 query_revision_gate: None,
474 query_channel: zerolaunch_plugin_api::QueryChannel::Ui,
476 locale: p.locale,
478 });
479 let result = InitializeResult {
480 protocol_version: PROTOCOL_VERSION.to_string(),
481 };
482 Ok(serde_json::to_value(result).unwrap_or_default())
483 }
484 plugin_methods::GET_COMPONENTS => {
486 let mut components = vec![ComponentDescriptor {
487 component_id: app.plugin.component_id().to_string(),
488 component_name: app.plugin.component_name().to_string(),
489 component_description: app.plugin.component_description().to_string(),
490 component_type: app.plugin.component_type(),
491 kind: ComponentKind::Plugin,
492 priority: app.plugin.priority(),
493 }];
494 for ds in &app.data_sources {
495 components.push(ComponentDescriptor {
496 component_id: ds.component_id().to_string(),
497 component_name: ds.component_name().to_string(),
498 component_description: ds.component_description().to_string(),
499 component_type: ds.component_type(),
500 kind: ComponentKind::DataSource,
501 priority: ds.priority(),
502 });
503 }
504 for ex in &app.executors {
505 components.push(ComponentDescriptor {
506 component_id: ex.component_id().to_string(),
507 component_name: ex.component_name().to_string(),
508 component_description: ex.component_description().to_string(),
509 component_type: ex.component_type(),
510 kind: ComponentKind::ActionExecutor {
511 target_types: ex.supported_target_types(),
512 },
513 priority: ex.priority(),
514 });
515 }
516 for engine in &app.search_engines {
517 components.push(ComponentDescriptor {
518 component_id: engine.component_id().to_string(),
519 component_name: engine.component_name().to_string(),
520 component_description: engine.component_description().to_string(),
521 component_type: engine.component_type(),
522 kind: ComponentKind::SearchEngine,
523 priority: engine.priority(),
524 });
525 }
526 for booster in &app.score_boosters {
527 components.push(ComponentDescriptor {
528 component_id: booster.component_id().to_string(),
529 component_name: booster.component_name().to_string(),
530 component_description: booster.component_description().to_string(),
531 component_type: booster.component_type(),
532 kind: ComponentKind::ScoreBooster,
533 priority: booster.priority(),
534 });
535 }
536 for optimizer in &app.keyword_optimizers {
537 components.push(ComponentDescriptor {
538 component_id: optimizer.component_id().to_string(),
539 component_name: optimizer.component_name().to_string(),
540 component_description: optimizer.component_description().to_string(),
541 component_type: optimizer.component_type(),
542 kind: ComponentKind::KeywordOptimizer,
543 priority: optimizer.priority(),
544 });
545 }
546 for injector in &app.keyword_injectors {
547 components.push(ComponentDescriptor {
548 component_id: injector.component_id().to_string(),
549 component_name: injector.component_name().to_string(),
550 component_description: injector.component_description().to_string(),
551 component_type: injector.component_type(),
552 kind: ComponentKind::KeywordInjector,
553 priority: injector.priority(),
554 });
555 }
556 Ok(serde_json::to_value(components).unwrap_or_default())
557 }
558 plugin_methods::GET_SETTINGS_SCHEMA => {
560 let p: GetSettingsSchemaParams = serde_json::from_value(params.clone())
561 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
562 let conf = find_configurable(app, &p.component_id)?;
563 Ok(serde_json::to_value(conf.setting_schema()).unwrap_or(serde_json::Value::Null))
564 }
565 plugin_methods::GET_SETTINGS => {
567 let p: GetSettingsParams = serde_json::from_value(params.clone())
568 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
569 let conf = find_configurable(app, &p.component_id)?;
570 Ok(conf.get_settings())
571 }
572 plugin_methods::GET_DEFAULT_ENABLED => {
574 let p: GetDefaultEnabledParams = serde_json::from_value(params.clone())
575 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
576 let conf = find_configurable(app, &p.component_id)?;
577 Ok(serde_json::to_value(conf.default_enabled()).unwrap_or_default())
578 }
579 plugin_methods::APPLY_SETTINGS => {
581 let p: ApplySettingsParams = serde_json::from_value(params.clone())
582 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
583 let conf = find_configurable(app, &p.component_id)?;
584 conf.apply_settings(p.settings)
585 .await
586 .map_err(|e| JsonRpcError::new(codes::PLUGIN_ERROR, e.to_string()))?;
587 Ok(serde_json::Value::Null)
588 }
589 plugin_methods::VALIDATE_SETTINGS => {
591 let p: ValidateSettingsParams = serde_json::from_value(params.clone())
592 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
593 let conf = find_configurable(app, &p.component_id)?;
594 let result = match conf.validate_settings(&p.settings).await {
595 Ok(()) => ValidateSettingsResult { error: None },
596 Err(e) => ValidateSettingsResult {
597 error: Some(e.to_string()),
598 },
599 };
600 Ok(serde_json::to_value(result).unwrap_or_default())
601 }
602 plugin_methods::CONFIG_ACTIONS => {
605 let p: ConfigActionsParams = serde_json::from_value(params.clone())
606 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
607 let conf = find_configurable(app, &p.component_id)?;
608 Ok(serde_json::to_value(conf.config_actions()).unwrap_or_default())
609 }
610 plugin_methods::EXECUTE_CONFIG_ACTION => {
611 let p: ExecuteConfigActionParams = serde_json::from_value(params.clone())
612 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
613 let conf = find_configurable(app, &p.component_id)?;
614 conf.execute_config_action(&p.action, &p.params)
615 .await
616 .map_err(|e| JsonRpcError::new(codes::PLUGIN_ERROR, e.to_string()))
617 }
618 plugin_methods::QUERY => {
619 let p: QueryParams = serde_json::from_value(params.clone())
620 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
621 let response = app
622 .plugin
623 .query(&p.ctx, &p.query)
624 .await
625 .map_err(|e| JsonRpcError::new(codes::PLUGIN_ERROR, e.to_string()))?;
626 Ok(serde_json::to_value(response).unwrap_or_default())
627 }
628 plugin_methods::MATCH_QUERY => {
631 let p: MatchQueryParams = serde_json::from_value(params.clone())
632 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
633 Ok(serde_json::Value::Bool(
634 app.plugin
635 .match_query(&p.raw_query, &p.trigger_keywords)
636 .await,
637 ))
638 }
639 plugin_methods::EXECUTE_ACTION => {
640 let p: ExecuteActionParams = serde_json::from_value(params.clone())
641 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
642 app.plugin
643 .execute_action(&p.ctx, &p.action_id, p.payload)
644 .await
645 .map_err(|e| JsonRpcError::new(codes::PLUGIN_ERROR, e.to_string()))?;
646 Ok(serde_json::Value::Null)
647 }
648 plugin_methods::INIT => {
652 let p: InitParams = serde_json::from_value(params.clone())
653 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
654 app.plugin
655 .init(&p.ctx, None)
656 .await
657 .map_err(|e| JsonRpcError::new(codes::PLUGIN_ERROR, e.to_string()))?;
658 Ok(serde_json::Value::Null)
659 }
660 plugin_methods::INTERACTION_POLICY => {
662 let _p: InteractionPolicyParams = serde_json::from_value(params.clone())
663 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
664 Ok(serde_json::to_value(app.plugin.interaction_policy()).unwrap_or_default())
665 }
666 plugin_methods::FETCH_CANDIDATES => {
668 let p: FetchCandidatesParams = serde_json::from_value(params.clone())
669 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
670 let ds = find_data_source(app, &p.component_id)?;
671 let cache = ds.fetch_candidates().await;
672 Ok(serde_json::to_value(FetchCandidatesResult {
673 candidates: cache.get_candidates().clone(),
674 })
675 .unwrap_or_default())
676 }
677 plugin_methods::SUPPORTED_TARGET_TYPES => {
679 let p: SupportedTargetTypesParams = serde_json::from_value(params.clone())
680 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
681 let ex = find_executor(app, &p.component_id)?;
682 Ok(serde_json::to_value(ex.supported_target_types()).unwrap_or_default())
683 }
684 plugin_methods::SUPPORTED_ACTIONS => {
687 let p: SupportedActionsParams = serde_json::from_value(params.clone())
688 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
689 let ex = find_executor(app, &p.component_id)?;
690 Ok(serde_json::to_value(ex.supported_actions()).unwrap_or_default())
691 }
692 plugin_methods::EXECUTOR_EXECUTE => {
694 let p: ExecutorExecuteParams = serde_json::from_value(params.clone())
695 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
696 let ex = find_executor(app, &p.component_id)?;
697 let result = match ex.execute(&p.execution_ctx, &p.action_id).await {
698 Ok(()) => ExecutorExecuteResult { error: None },
699 Err(e) => ExecutorExecuteResult {
700 error: Some(e.to_string()),
701 },
702 };
703 Ok(serde_json::to_value(result).unwrap_or_default())
704 }
705 plugin_methods::CALCULATE_SCORES => {
708 let p: CalculateScoresParams = serde_json::from_value(params.clone())
709 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
710 let engine = find_search_engine(app, &p.component_id)?;
711 let cache = CachedCandidateData::from_data(p.candidates);
712 let scored = engine.calculate_scores(&cache, &p.query).await;
713 Ok(serde_json::to_value(scored).unwrap_or_default())
714 }
715 plugin_methods::BOOSTER_BOOST => {
717 let p: BoosterBoostParams = serde_json::from_value(params.clone())
718 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
719 let booster = find_score_booster(app, &p.component_id)?;
720 let cache = CachedCandidateData::from_data(p.candidates);
721 let mut scored = p.scored;
722 booster.boost(&mut scored, &cache, &p.query).await;
723 Ok(serde_json::to_value(scored).unwrap_or_default())
724 }
725 plugin_methods::BOOSTER_RECORD => {
727 let p: BoosterRecordParams = serde_json::from_value(params.clone())
728 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
729 let booster = find_score_booster(app, &p.component_id)?;
730 let cache = CachedCandidateData::from_data(p.candidates);
731 booster.record(p.candidate_id, &cache, &p.query).await;
732 Ok(serde_json::Value::Null)
733 }
734 plugin_methods::KEYWORD_OPTIMIZER_INFO => {
736 let p: KeywordOptimizerInfoParams = serde_json::from_value(params.clone())
737 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
738 let optimizer = find_keyword_optimizer(app, &p.component_id)?;
739 Ok(serde_json::to_value(KeywordOptimizerInfo {
740 input_source: optimizer.input_source(),
741 priority: optimizer.get_priority(),
742 })
743 .unwrap_or_default())
744 }
745 plugin_methods::KEYWORD_OPTIMIZE => {
747 let p: KeywordOptimizeParams = serde_json::from_value(params.clone())
748 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
749 let optimizer = find_keyword_optimizer(app, &p.component_id)?;
750 Ok(serde_json::to_value(optimizer.optimize(&p.keyword).await).unwrap_or_default())
751 }
752 plugin_methods::KEYWORD_INJECT => {
754 let p: KeywordInjectParams = serde_json::from_value(params.clone())
755 .map_err(|e| JsonRpcError::new(codes::INVALID_PARAMS, e.to_string()))?;
756 let injector = find_keyword_injector(app, &p.component_id)?;
757 Ok(
758 serde_json::to_value(injector.inject_keywords(&p.candidate).await)
759 .unwrap_or_default(),
760 )
761 }
762 _ => Err(JsonRpcError::new(
763 codes::METHOD_NOT_FOUND,
764 format!("method not found: {}", method),
765 )),
766 }
767}
768
769fn find_configurable<'a>(
772 app: &'a PluginApp,
773 component_id: &str,
774) -> Result<&'a dyn Configurable, JsonRpcError> {
775 let entry = app.by_id.get(component_id).ok_or_else(|| {
776 JsonRpcError::new(
777 codes::METHOD_NOT_FOUND,
778 format!("component not found: {component_id}"),
779 )
780 })?;
781 Ok(match entry {
782 ComponentEntry::Plugin(p) => p.as_ref() as &dyn Configurable,
783 ComponentEntry::DataSource(ds) => ds.as_ref() as &dyn Configurable,
784 ComponentEntry::Executor(ex) => ex.as_ref() as &dyn Configurable,
785 ComponentEntry::SearchEngine(engine) => engine.as_ref() as &dyn Configurable,
786 ComponentEntry::ScoreBooster(booster) => booster.as_ref() as &dyn Configurable,
787 ComponentEntry::KeywordOptimizer(optimizer) => optimizer.as_ref() as &dyn Configurable,
788 ComponentEntry::KeywordInjector(injector) => injector.as_ref() as &dyn Configurable,
789 })
790}
791
792fn find_data_source<'a>(
794 app: &'a PluginApp,
795 component_id: &str,
796) -> Result<&'a dyn DataSource, JsonRpcError> {
797 match app.by_id.get(component_id) {
798 Some(ComponentEntry::DataSource(ds)) => Ok(ds.as_ref()),
799 Some(_) => Err(JsonRpcError::new(
800 codes::METHOD_NOT_FOUND,
801 format!("component is not a data source: {component_id}"),
802 )),
803 None => Err(JsonRpcError::new(
804 codes::METHOD_NOT_FOUND,
805 format!("data source not found: {component_id}"),
806 )),
807 }
808}
809
810fn find_executor<'a>(
812 app: &'a PluginApp,
813 component_id: &str,
814) -> Result<&'a dyn ActionExecutor, JsonRpcError> {
815 match app.by_id.get(component_id) {
816 Some(ComponentEntry::Executor(ex)) => Ok(ex.as_ref()),
817 Some(_) => Err(JsonRpcError::new(
818 codes::METHOD_NOT_FOUND,
819 format!("component is not an executor: {component_id}"),
820 )),
821 None => Err(JsonRpcError::new(
822 codes::METHOD_NOT_FOUND,
823 format!("executor not found: {component_id}"),
824 )),
825 }
826}
827fn find_search_engine<'a>(
829 app: &'a PluginApp,
830 component_id: &str,
831) -> Result<&'a dyn SearchEngine, JsonRpcError> {
832 match app.by_id.get(component_id) {
833 Some(ComponentEntry::SearchEngine(engine)) => Ok(engine.as_ref()),
834 Some(_) => Err(JsonRpcError::new(
835 codes::METHOD_NOT_FOUND,
836 format!("component is not a search engine: {component_id}"),
837 )),
838 None => Err(JsonRpcError::new(
839 codes::METHOD_NOT_FOUND,
840 format!("search engine not found: {component_id}"),
841 )),
842 }
843}
844
845fn find_score_booster<'a>(
847 app: &'a PluginApp,
848 component_id: &str,
849) -> Result<&'a dyn ScoreBooster, JsonRpcError> {
850 match app.by_id.get(component_id) {
851 Some(ComponentEntry::ScoreBooster(booster)) => Ok(booster.as_ref()),
852 Some(_) => Err(JsonRpcError::new(
853 codes::METHOD_NOT_FOUND,
854 format!("component is not a score booster: {component_id}"),
855 )),
856 None => Err(JsonRpcError::new(
857 codes::METHOD_NOT_FOUND,
858 format!("score booster not found: {component_id}"),
859 )),
860 }
861}
862
863fn find_keyword_optimizer<'a>(
865 app: &'a PluginApp,
866 component_id: &str,
867) -> Result<&'a dyn KeywordOptimizer, JsonRpcError> {
868 match app.by_id.get(component_id) {
869 Some(ComponentEntry::KeywordOptimizer(optimizer)) => Ok(optimizer.as_ref()),
870 Some(_) => Err(JsonRpcError::new(
871 codes::METHOD_NOT_FOUND,
872 format!("component is not a keyword optimizer: {component_id}"),
873 )),
874 None => Err(JsonRpcError::new(
875 codes::METHOD_NOT_FOUND,
876 format!("keyword optimizer not found: {component_id}"),
877 )),
878 }
879}
880
881fn find_keyword_injector<'a>(
883 app: &'a PluginApp,
884 component_id: &str,
885) -> Result<&'a dyn KeywordInjector, JsonRpcError> {
886 match app.by_id.get(component_id) {
887 Some(ComponentEntry::KeywordInjector(injector)) => Ok(injector.as_ref()),
888 Some(_) => Err(JsonRpcError::new(
889 codes::METHOD_NOT_FOUND,
890 format!("component is not a keyword injector: {component_id}"),
891 )),
892 None => Err(JsonRpcError::new(
893 codes::METHOD_NOT_FOUND,
894 format!("keyword injector not found: {component_id}"),
895 )),
896 }
897}