Skip to main content

zerolaunch_plugin_sdk_rust/
runtime.rs

1//! 第三方 Rust 插件的 JSON-RPC 运行时。
2//!
3//! 运行时采用三任务异步架构:
4//! - `read_task`:唯一 stdin 读取者,解析 LSP 帧,路由响应到 pending_map,
5//!   转发请求到 dispatch_task。
6//! - `write_task`:唯一 stdout 写入者,将所有出站消息编码为 LSP 帧。
7//! - `dispatch_task`:处理 plugin/* 请求,调用用户 Plugin trait 实现,
8//!   将响应发到 write_task。
9//!
10//! HostProxy 通过共享的 pending_map 和 outbound_tx 发送 host/* 请求,
11//! 避免了同步 BufReader 造成的死锁问题。
12
13use 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
37// 全局 HostProxy,由 `run_async` 初始化。
38static HOST_PROXY: OnceLock<Arc<HostProxy>> = OnceLock::new();
39
40/// 返回当前运行时的 `HostProxy`。
41/// 在 `run()` 之前调用会 panic。
42pub fn host() -> Arc<HostProxy> {
43    HOST_PROXY
44        .get()
45        .expect("host() 必须在 run() 之后调用")
46        .clone()
47}
48
49/// 从 read task 路由到 dispatch task 的入站 JSON-RPC 请求。
50struct IncomingRequest {
51    id: u64,
52    method: String,
53    params: serde_json::Value,
54}
55
56/// SDK 组件集合:一个 Plugin 主组件 + 任意个 DataSource / ActionExecutor /
57/// SearchEngine / ScoreBooster / KeywordOptimizer / KeywordInjector 附加组件。
58///
59/// 与内置插件完全对等:每个组件都是独立的 Configurable(各自的 component_id、
60/// schema、设置),进程向宿主声明全部组件(GET_COMPONENTS 返回全部)。
61pub struct PluginApp {
62    /// 插件主组件:进程级元数据、触发词、面板查询。
63    pub(crate) plugin: Arc<dyn Plugin>,
64    pub(crate) data_sources: Vec<Arc<dyn DataSource>>,
65    pub(crate) executors: Vec<Arc<dyn ActionExecutor>>,
66    /// 搜索引擎组件(与内置引擎互斥,仅允许一个启用)。
67    pub(crate) search_engines: Vec<Arc<dyn SearchEngine>>,
68    /// 分数增强器组件(引擎打分后追加修正)。
69    pub(crate) score_boosters: Vec<Arc<dyn ScoreBooster>>,
70    /// 关键词优化器组件(扩展候选关键词)。
71    pub(crate) keyword_optimizers: Vec<Arc<dyn KeywordOptimizer>>,
72    /// 关键词注入器组件(基于候选上下文注入关键词)。
73    pub(crate) keyword_injectors: Vec<Arc<dyn KeywordInjector>>,
74    /// component_id → 组件统一索引(全部组件类别皆入)。
75    /// dispatch 按 component_id 以 O(1) 路由组件方法。
76    by_id: HashMap<String, ComponentEntry>,
77}
78
79/// 组件索引条目:统一持有各 trait 对象,按需向上转型。
80enum ComponentEntry {
81    Plugin(Arc<dyn Plugin>),
82    DataSource(Arc<dyn DataSource>),
83    Executor(Arc<dyn ActionExecutor>),
84    /// 搜索引擎组件(与内置引擎互斥,仅允许一个启用)。
85    SearchEngine(Arc<dyn SearchEngine>),
86    /// 分数增强器组件。
87    ScoreBooster(Arc<dyn ScoreBooster>),
88    /// 关键词优化器组件。
89    KeywordOptimizer(Arc<dyn KeywordOptimizer>),
90    /// 关键词注入器组件。
91    KeywordInjector(Arc<dyn KeywordInjector>),
92}
93
94impl PluginApp {
95    /// 以 Plugin 主组件构建应用(必须存在:进程级 metadata/触发词/面板查询属于它)。
96    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    /// 附加 DataSource 组件(候选采集;与内置数据源完全对等)。
116    /// 组件 id 重复属于插件编码错误,直接 panic 暴露。
117    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    /// 附加 ActionExecutor 组件(候选执行;与内置执行器完全对等)。
132    /// 组件 id 重复属于插件编码错误,直接 panic 暴露。
133    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    /// 附加 SearchEngine 组件(搜索打分;与内置搜索引擎完全对等)。
147    /// 组件 id 重复属于插件编码错误,直接 panic 暴露。
148    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    /// 附加 ScoreBooster 组件(分数增强;与内置分数增强器完全对等)。
163    /// 组件 id 重复属于插件编码错误,直接 panic 暴露。
164    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    /// 附加 KeywordOptimizer 组件(关键词扩展;与内置优化器完全对等)。
179    /// 组件 id 重复属于插件编码错误,直接 panic 暴露。
180    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    /// 附加 KeywordInjector 组件(基于候选上下文注入关键词;与内置注入器完全对等)。
198    /// 组件 id 重复属于插件编码错误,直接 panic 暴露。
199    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    /// 运行 JSON-RPC stdio 循环(阻塞当前线程直到进程退出)。
217    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
230/// 使用给定的 Plugin 实现运行 JSON-RPC stdio 循环。
231/// 等价于 `PluginApp::new(plugin).run()`。
232pub fn run(plugin: impl Plugin + 'static) {
233    PluginApp::new(plugin).run()
234}
235
236async fn run_async(mut app: PluginApp) {
237    // 初始化日志系统(双写:stderr → 文件 + WARN/ERROR → host/log 转发)
238    let mut log_rx = logging::init_logging();
239    let stdin = tokio::io::stdin();
240    let stdout = tokio::io::stdout();
241
242    // 通道
243    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    // 创建 HostProxy 并写入全局,供 dispatch/read/write 三任务共享访问。
248    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    // 插件状态
253    let mut plugin_context: Option<zerolaunch_plugin_api::PluginContext> = None;
254
255    // --- 日志转发后台任务:将 WARN/ERROR 非阻塞转发到宿主 ---
256    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    // --- 读任务:stdin → pending_map(响应)或 request_tx(新请求)---
263    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                        // 错误走 Err 分支携带结构化 JsonRpcError,不再伪装成字符串结果:
289                        // 宿主拒绝/模型错误等可被调用方区分,模型方法 from_value 不再被
290                        // "invalid type: string" 淹没真实错误。
291                        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                        // 如果没有对应的 pending channel,说明响应已经超时或被取消,忽略。同时打印一下被忽略的信息
299                        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                    // 如果 dispatch task 已退出,说明插件可能已经崩溃或被关闭,无法处理请求。打印警告信息。
315                    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    // --- 分发任务:plugin/* 请求 → 用户 Plugin → 响应到 outbound_tx ---
330    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            // 调用用户实现的 Plugin trait 处理,并将响应发送到 outbound_tx。
339            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    // --- 写任务:outbound_rx → stdout ---
355    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    // 等待读任务结束(传输层关闭),或 dispatch task 崩溃(插件代码 panic)。
378    tokio::select! {
379        // 正常退出:stdin 关闭(宿主 shutdown / EOF)→ 读任务结束。
380        _ = &mut read_handle => {
381            // 释放 request_tx → dispatch_task 在当前请求处理完后
382            // 通过 channel 关闭优雅退出。
383            drop(request_tx);
384            // 写任务无法通过 channel 关闭退出,因为 HOST_PROXY 全局
385            // 仍持有 outbound_tx 的 clone,因此直接 abort。
386            write_handle.abort();
387            let _ = (&mut dispatch_handle).await;
388        }
389        // 插件代码 panic:dispatch task 死亡,进程退出由宿主按 auto_restart 重启。
390        r = &mut dispatch_handle => {
391            if r.is_err() {
392                tracing::error!("dispatch task 因 panic 退出,插件进程终止(宿主将按 auto_restart 重启)");
393            }
394        }
395    }
396}
397
398/// 读取单条 LSP 风格 Content-Length 帧消息。
399/// 返回原始 JSON 字节,或在解析/大小/IO 失败时返回错误字符串。
400async 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
441// / 处理单条 plugin/* 请求,返回响应 Message。
442async 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_ctx 并返回插件版本信息。
463        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            // 记录插件 id:t_key() 据此自动补 `plugin.<id>.` 前缀
467            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                // 远端插件无宿主查询版本门控,恒视为最新。
473                query_revision_gate: None,
474                // 远端插件会话由宿主经 RPC 下发通道,未收到时缺省视为 GUI 通道。
475                query_channel: zerolaunch_plugin_api::QueryChannel::Ui,
476                // 宿主语言在握手时下发(InitializeParams.locale),写入会话上下文。
477                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 + 附加 DataSource / ActionExecutor)
485        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        // 返回指定组件的注册配置项
559        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        // 返回指定组件当前的配置值
566        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        // 返回指定组件的默认启用状态
573        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        // 宿主下发新的配置值,指定组件据此更新自身行为
580        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        // 验证一组配置值是否合法(不会实际应用),返回验证结果或错误信息
590        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        // 返回裸数组:宿主 discover 流程以 `Vec<ConfigActionDef>` 反序列化
603        // (与 get_settings_schema 的裸数组约定一致),包装结构会解析失败。
604        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        // 查询匹配裁决(插件级语义):宿主在路由阶段询问本插件是否接管当前输入。
629        // 请求携带该插件声明的触发词:默认实现据此做框架关键词判定(覆盖者可忽略)。
630        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        // 插件初始化钩子:宿主在注册完成后调用(内置插件在启动期统一 init)。
649        // 远端进程无宿主 PluginHandle(跨进程不可序列化),以 None 传入,
650        // 平台能力经 host() 的 host/* RPC 访问——与内置插件语义对等。
651        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        // 插件交互策略(查询触发方式/防抖/按键绑定):宿主在会话推送时读取。
661        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        // DataSource 组件:采集候选项(与内置数据源对等)
667        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        // ActionExecutor 组件:支持的目标类型列表
678        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        // ActionExecutor 组件:支持的动作列表(与内置 supported_actions() 语义一致,
685        // 不区分 target_type——宿主按 (id, label) 去重合并)
686        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        // ActionExecutor 组件:执行动作(完整 ExecutionContext 原样下发)
693        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        // SearchEngine 组件:对缓存候选计算分数(候选快照保真还原,
706        // 插件进程内按原 id / 原顺序可回查)
707        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        // ScoreBooster 组件:增强分数(返回增强后的完整列表,索引与传入一致)
716        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        // ScoreBooster 组件:记录用户确认(选中候选 → 学习用户习惯)
726        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        // KeywordOptimizer 组件:声明属性(input_source / priority 经 info RPC 上报)
735        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        // KeywordOptimizer 组件:优化单个关键词
746        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        // KeywordInjector 组件:对单个候选注入关键词
753        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
769/// 按 component_id 查找 Configurable 组件(Plugin / DataSource / ActionExecutor 皆可)。
770/// 组件未注册时返回 METHOD_NOT_FOUND。
771fn 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
792/// 按 component_id 查找 DataSource 组件。
793fn 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
810/// 按 component_id 查找 ActionExecutor 组件。
811fn 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}
827/// 按 component_id 查找 SearchEngine 组件。
828fn 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
845/// 按 component_id 查找 ScoreBooster 组件。
846fn 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
863/// 按 component_id 查找 KeywordOptimizer 组件。
864fn 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
881/// 按 component_id 查找 KeywordInjector 组件。
882fn 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}