Skip to main content

onlyne_client/ops/local_cli/
plugin.rs

1//! Plugin package verbs: install or remove a vendored coding-agent package,
2//! then tell a running client that the package changed.
3
4use super::config::{
5    append_plugin_entry, config_lists_plugin, remove_plugin_entry, render_plugins,
6};
7use anyhow::{Context, Result, anyhow};
8use onlyne_adapter::AdapterIo;
9use onlyne_config::layout::RoleWorkspace;
10use onlyne_proto::{AdapterMsg, HelloArgs, MountKind, PROTOCOL_VERSION, PluginOp};
11use onlyne_wire::socket::connect_local;
12use std::path::{Path, PathBuf};
13use std::time::Duration;
14
15/// Wait bound for one handshake with a running client.
16pub const CLIENT_PROBE_TIMEOUT: Duration = Duration::from_secs(2);
17/// Line an operator sees when a plugin change needs a client restart.
18pub const RESTART_HINT: &str =
19    "onlyne: no running client; restart `onlyne-client run` to mount this plugin";
20
21/// Plugin id rule: the id becomes a path component under `agent/<id>/`.
22pub fn valid_plugin_id(id: &str) -> bool {
23    let bytes = id.as_bytes();
24    if bytes.is_empty() || bytes.len() > 32 {
25        return false;
26    }
27    let first = bytes[0];
28    let first_ok = first.is_ascii_lowercase() || first.is_ascii_digit();
29    if !first_ok {
30        return false;
31    }
32    bytes.iter().all(|byte| {
33        byte.is_ascii_lowercase() || byte.is_ascii_digit() || *byte == b'_' || *byte == b'-'
34    })
35}
36
37/// Directory holding one vendored coding-agent plugin package.
38pub fn agent_package_dir(workspace: &Path, plugin_id: &str) -> PathBuf {
39    RoleWorkspace::resolve(workspace).agent_dir(plugin_id)
40}
41
42/// Install a coding-agent plugin package into `agent/<id>/`.
43pub fn agent_install(
44    workspace: &Path,
45    package: &Path,
46    plugin_id: &str,
47    agent: Option<&str>,
48) -> Result<Vec<String>> {
49    if !valid_plugin_id(plugin_id) {
50        return Err(anyhow!("onlyne: invalid plugin id {plugin_id}"));
51    }
52    let target = agent_package_dir(workspace, plugin_id);
53    if target.exists() {
54        return Err(anyhow!("onlyne: plugin {plugin_id} already installed"));
55    }
56    // The config is the other half of "installed": a workspace that already
57    // lists the id is installed even when its package directory is gone, so
58    // the refusal lands before any filesystem change.
59    if config_lists_plugin(workspace, plugin_id)? {
60        return Err(anyhow!("onlyne: plugin {plugin_id} already installed"));
61    }
62    std::fs::create_dir_all(&target).with_context(|| format!("create {}", target.display()))?;
63    install_package_files(package, &target)?;
64    let binary = find_plugin_binary(&target)?;
65    set_executable(&target.join(&binary))?;
66    let agent_name = agent.unwrap_or(plugin_id).to_string();
67    let manifest = format!(
68        "id = {plugin_id:?}\nbinary = {binary:?}\nagent = {agent_name:?}\ncapabilities = []\n"
69    );
70    std::fs::write(target.join("plugin.toml"), manifest)?;
71    let ids = append_plugin_entry(workspace, plugin_id)?;
72    Ok(vec![
73        format!("installed {plugin_id}"),
74        format!("wrote {}", target.join("plugin.toml").display()),
75        format!("registered plugin {plugin_id} in {}", render_plugins(&ids)),
76    ])
77}
78
79/// Remove a plugin package and its config entry.
80pub fn agent_uninstall(workspace: &Path, plugin_id: &str) -> Result<Vec<String>> {
81    if !valid_plugin_id(plugin_id) {
82        return Err(anyhow!("onlyne: invalid plugin id {plugin_id}"));
83    }
84    let target = agent_package_dir(workspace, plugin_id);
85    if !target.exists() {
86        return Err(anyhow!("onlyne: no plugin {plugin_id} installed"));
87    }
88    remove_plugin_entry(workspace, plugin_id)?;
89    std::fs::remove_dir_all(&target).with_context(|| format!("remove {}", target.display()))?;
90    Ok(vec![
91        format!("deregistered plugin {plugin_id} from plugins"),
92        format!("removed {}", target.display()),
93    ])
94}
95
96/// Frame the client would send to hot-mount a plugin over its local socket.
97pub fn agent_mount_frame(plugin_id: &str) -> serde_json::Value {
98    serde_json::json!({"f": "req", "id": plugin_id, "op": "control", "args": {"plugin": plugin_id, "action": "mount"}})
99}
100
101/// Frame the client would send to stop a mounted plugin.
102pub fn agent_stop_frame(plugin_id: &str) -> serde_json::Value {
103    serde_json::json!({"f": "req", "id": plugin_id, "op": "control", "args": {"plugin": plugin_id, "action": "stop"}})
104}
105
106/// What one plugin verb did, split by the stream each line belongs to.
107#[derive(Debug, Clone, PartialEq, Eq)]
108pub struct PluginAction {
109    /// Lines for stdout, one per filesystem action.
110    pub lines: Vec<String>,
111    /// Line for stderr when no client answered the local socket.
112    pub hint: Option<String>,
113}
114
115/// Install a package, then tell a running client about it.
116pub async fn install_verb(
117    workspace: &Path,
118    package: &Path,
119    plugin_id: &str,
120    agent: Option<&str>,
121) -> Result<PluginAction> {
122    let lines = agent_install(workspace, package, plugin_id, agent)?;
123    let running = notify_client(workspace, plugin_id).await;
124    Ok(PluginAction {
125        lines,
126        hint: (!running).then(|| RESTART_HINT.to_string()),
127    })
128}
129
130/// Remove a package and its config entry, then tell a running client.
131pub async fn uninstall_verb(workspace: &Path, plugin_id: &str) -> Result<PluginAction> {
132    let lines = agent_uninstall(workspace, plugin_id)?;
133    let running = notify_client(workspace, plugin_id).await;
134    Ok(PluginAction {
135        lines,
136        hint: (!running).then(|| RESTART_HINT.to_string()),
137    })
138}
139
140/// Operator refusal for a plugin verb.
141pub fn plugin_exit_code(error: &anyhow::Error) -> i32 {
142    if error.to_string().starts_with("onlyne: ") {
143        2
144    } else {
145        1
146    }
147}
148
149/// Tell a running client that a plugin package changed.
150///
151/// The adapter socket has one host-side entry point, the `hello` handshake, so
152/// the notification is that handshake on an `admin` mount. `false` means no
153/// client is listening on the workspace socket.
154pub async fn notify_client(workspace: &Path, plugin_id: &str) -> bool {
155    let socket = RoleWorkspace::resolve(workspace).socket_path();
156    let Ok(stream) = connect_local(&socket).await else {
157        return false;
158    };
159    let io = AdapterIo::new(stream, CLIENT_PROBE_TIMEOUT, CLIENT_PROBE_TIMEOUT);
160    let hello = HelloArgs {
161        protocol: PROTOCOL_VERSION,
162        plugin: format!("onlyne-client-cli:{plugin_id}"),
163        version: env!("CARGO_PKG_VERSION").to_string(),
164        kind: MountKind::Admin,
165        capabilities: Vec::new(),
166        mount: None,
167    };
168    matches!(io.request(AdapterMsg::Plugin(PluginOp::Hello(hello))).await, Ok(body) if body.ok)
169}
170
171fn install_package_files(package: &Path, target: &Path) -> Result<()> {
172    if package.is_dir() {
173        let entries: Vec<_> = std::fs::read_dir(package)
174            .with_context(|| format!("read {}", package.display()))?
175            .collect::<std::result::Result<Vec<_>, _>>()?;
176        for entry in &entries {
177            let name = entry.file_name().to_string_lossy().into_owned();
178            let source = entry.path();
179            let destination = target.join(&name);
180            if source.is_dir() {
181                return Err(anyhow!("onlyne: package directory must be flat"));
182            }
183            std::fs::copy(&source, &destination).with_context(|| format!("copy {}", name))?;
184        }
185        return Ok(());
186    }
187    if package.extension().and_then(|ext| ext.to_str()) == Some("gz") {
188        let file =
189            std::fs::File::open(package).with_context(|| format!("open {}", package.display()))?;
190        let decoder = flate2::read::GzDecoder::new(file);
191        let mut archive = tar::Archive::new(decoder);
192        archive
193            .unpack(target)
194            .with_context(|| format!("unpack {}", package.display()))?;
195        return Ok(());
196    }
197    Err(anyhow!(
198        "onlyne: package must be a directory or a .tar.gz file"
199    ))
200}
201
202fn find_plugin_binary(target: &Path) -> Result<String> {
203    let mut binaries = Vec::new();
204    for entry in std::fs::read_dir(target).with_context(|| format!("read {}", target.display()))? {
205        let entry = entry?;
206        let name = entry.file_name().to_string_lossy().into_owned();
207        if name.starts_with("onlyne-agent-") && entry.file_type()?.is_file() {
208            binaries.push(name);
209        }
210    }
211    binaries.sort();
212    if binaries.len() == 1 {
213        Ok(binaries.remove(0))
214    } else {
215        Err(anyhow!(
216            "onlyne: package must contain exactly one onlyne-agent-* executable"
217        ))
218    }
219}
220
221#[cfg(unix)]
222fn set_executable(path: &Path) -> Result<()> {
223    use std::os::unix::fs::PermissionsExt;
224    let mut permissions = std::fs::metadata(path)
225        .with_context(|| format!("stat {}", path.display()))?
226        .permissions();
227    permissions.set_mode(0o755);
228    std::fs::set_permissions(path, permissions)
229        .with_context(|| format!("chmod {}", path.display()))?;
230    Ok(())
231}
232
233#[cfg(not(unix))]
234fn set_executable(_path: &Path) -> Result<()> {
235    Ok(())
236}