Skip to main content

zenoh_plugin_trait/
manager.rs

1// Copyright (c) 2023 ZettaScale Technology
2//
3// This program and the accompanying materials are made available under the
4// terms of the Eclipse Public License 2.0 which is available at
5// http://www.eclipse.org/legal/epl-2.0, or the Apache License, Version 2.0
6// which is available at https://www.apache.org/licenses/LICENSE-2.0.
7//
8// SPDX-License-Identifier: EPL-2.0 OR Apache-2.0
9//
10// Contributors:
11//   ZettaScale Zenoh Team, <zenoh@zettascale.tech>
12//
13mod dynamic_plugin;
14mod static_plugin;
15
16use std::fmt;
17
18use zenoh_keyexpr::keyexpr;
19use zenoh_result::ZResult;
20use zenoh_util::LibLoader;
21
22use self::{
23    dynamic_plugin::{DynamicPlugin, DynamicPluginSource},
24    static_plugin::StaticPlugin,
25};
26use crate::*;
27
28pub trait DeclaredPlugin<StartArgs, Instance>: PluginStatus {
29    fn as_status(&self) -> &dyn PluginStatus;
30    fn load(&mut self) -> ZResult<Option<&mut dyn LoadedPlugin<StartArgs, Instance>>>;
31    fn loaded(&self) -> Option<&dyn LoadedPlugin<StartArgs, Instance>>;
32    fn loaded_mut(&mut self) -> Option<&mut dyn LoadedPlugin<StartArgs, Instance>>;
33}
34pub trait LoadedPlugin<StartArgs, Instance>: PluginStatus {
35    fn as_status(&self) -> &dyn PluginStatus;
36    fn required(&self) -> bool;
37    fn start(&mut self, args: &StartArgs) -> ZResult<&mut dyn StartedPlugin<StartArgs, Instance>>;
38    fn started(&self) -> Option<&dyn StartedPlugin<StartArgs, Instance>>;
39    fn started_mut(&mut self) -> Option<&mut dyn StartedPlugin<StartArgs, Instance>>;
40}
41
42pub trait StartedPlugin<StartArgs, Instance>: PluginStatus {
43    fn as_status(&self) -> &dyn PluginStatus;
44    fn stop(&mut self);
45    fn instance(&self) -> &Instance;
46    fn instance_mut(&mut self) -> &mut Instance;
47}
48
49struct PluginRecord<StartArgs: PluginStartArgs, Instance: PluginInstance>(
50    Box<dyn DeclaredPlugin<StartArgs, Instance> + Send + Sync>,
51);
52
53impl<StartArgs: PluginStartArgs, Instance: PluginInstance> PluginRecord<StartArgs, Instance> {
54    fn new<P: DeclaredPlugin<StartArgs, Instance> + Send + Sync + 'static>(plugin: P) -> Self {
55        Self(Box::new(plugin))
56    }
57}
58
59impl<StartArgs: PluginStartArgs, Instance: PluginInstance> PluginStatus
60    for PluginRecord<StartArgs, Instance>
61{
62    fn name(&self) -> &str {
63        self.0.name()
64    }
65
66    fn id(&self) -> &str {
67        self.0.id()
68    }
69
70    fn version(&self) -> Option<&str> {
71        self.0.version()
72    }
73    fn long_version(&self) -> Option<&str> {
74        self.0.long_version()
75    }
76    fn path(&self) -> &str {
77        self.0.path()
78    }
79    fn state(&self) -> PluginState {
80        self.0.state()
81    }
82    fn report(&self) -> PluginReport {
83        self.0.report()
84    }
85}
86
87impl<StartArgs: PluginStartArgs, Instance: PluginInstance> DeclaredPlugin<StartArgs, Instance>
88    for PluginRecord<StartArgs, Instance>
89{
90    fn as_status(&self) -> &dyn PluginStatus {
91        self
92    }
93    fn load(&mut self) -> ZResult<Option<&mut dyn LoadedPlugin<StartArgs, Instance>>> {
94        self.0.load()
95    }
96    fn loaded(&self) -> Option<&dyn LoadedPlugin<StartArgs, Instance>> {
97        self.0.loaded()
98    }
99    fn loaded_mut(&mut self) -> Option<&mut dyn LoadedPlugin<StartArgs, Instance>> {
100        self.0.loaded_mut()
101    }
102}
103
104/// A plugins manager that handles starting and stopping plugins.
105/// Plugins can be loaded from shared libraries using [`Self::declare_dynamic_plugin_by_name`] or [`Self::declare_dynamic_plugin_by_paths`], or added directly from the binary if available using [`Self::declare_static_plugin`].
106pub struct PluginsManager<StartArgs: PluginStartArgs, Instance: PluginInstance> {
107    default_lib_prefix: String,
108    loader: Option<LibLoader>,
109    plugins: Vec<PluginRecord<StartArgs, Instance>>,
110}
111
112impl<StartArgs: PluginStartArgs, Instance: PluginInstance> fmt::Debug
113    for PluginsManager<StartArgs, Instance>
114{
115    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
116        f.debug_struct("PluginsManager")
117            .field("default_lib_prefix", &self.default_lib_prefix)
118            .field("loader", &self.loader.as_ref().map(|_| ".."))
119            .field("plugins_len", &self.plugins.len())
120            .finish()
121    }
122}
123
124impl<StartArgs: PluginStartArgs + 'static, Instance: PluginInstance + 'static>
125    PluginsManager<StartArgs, Instance>
126{
127    /// Constructs a new plugin manager with dynamic library loading enabled.
128    pub fn dynamic<S: Into<String>>(loader: LibLoader, default_lib_prefix: S) -> Self {
129        PluginsManager {
130            default_lib_prefix: default_lib_prefix.into(),
131            loader: Some(loader),
132            plugins: Default::default(),
133        }
134    }
135    /// Constructs a new plugin manager with dynamic library loading disabled.
136    pub fn static_plugins_only() -> Self {
137        PluginsManager {
138            default_lib_prefix: String::new(),
139            loader: None,
140            plugins: Default::default(),
141        }
142    }
143
144    /// Adds a statically linked plugin to the manager.
145    pub fn declare_static_plugin<
146        P: Plugin<StartArgs = StartArgs, Instance = Instance> + Send + Sync,
147        S: Into<String>,
148    >(
149        &mut self,
150        id: S,
151        required: bool,
152    ) {
153        let id = id.into();
154        let plugin_loader: StaticPlugin<StartArgs, Instance, P> =
155            StaticPlugin::new(id.clone(), required);
156
157        if self.get_plugin_index(&id).is_some() {
158            tracing::warn!(
159                "Duplicate plugin with ID: {id}, only the last declared one will be loaded"
160            )
161        }
162
163        self.plugins.push(PluginRecord::new(plugin_loader));
164        tracing::debug!(
165            "Declared static plugin Id:{} - Name:{}",
166            self.plugins.last().unwrap().id(),
167            self.plugins.last().unwrap().name()
168        );
169    }
170
171    /// Add dynamic plugin to the manager by name, automatically prepending the default library prefix
172    pub fn declare_dynamic_plugin_by_name<S: Into<String>>(
173        &mut self,
174        id: S,
175        plugin_name: S,
176        required: bool,
177    ) -> ZResult<&mut dyn DeclaredPlugin<StartArgs, Instance>> {
178        let plugin_name = plugin_name.into();
179        let id = id.into();
180        let libplugin_name = format!("{}{}", self.default_lib_prefix, plugin_name);
181        let libloader = self
182            .loader
183            .as_ref()
184            .ok_or("Dynamic plugin loading is disabled")?
185            .clone();
186        tracing::debug!(
187            "Declared dynamic plugin {} by name {}",
188            &id,
189            &libplugin_name
190        );
191        let loader = DynamicPlugin::new(
192            plugin_name,
193            id.clone(),
194            DynamicPluginSource::ByName((libloader, libplugin_name)),
195            required,
196        );
197
198        if self.get_plugin_index(&id).is_some() {
199            tracing::warn!(
200                "Duplicate plugin with ID: {id}, only the last declared one will be loaded"
201            )
202        }
203        self.plugins.push(PluginRecord::new(loader));
204        Ok(self.plugins.last_mut().unwrap())
205    }
206
207    /// Add first available dynamic plugin from the list of paths to the plugin files
208    pub fn declare_dynamic_plugin_by_paths<S: Into<String>, P: AsRef<str> + std::fmt::Debug>(
209        &mut self,
210        name: S,
211        id: S,
212        paths: &[P],
213        required: bool,
214    ) -> ZResult<&mut dyn DeclaredPlugin<StartArgs, Instance>> {
215        let name = name.into();
216        let id = id.into();
217        let paths = paths.iter().map(|p| p.as_ref().into()).collect();
218        tracing::debug!("Declared dynamic plugin {} by paths {:?}", &id, &paths);
219        let loader = DynamicPlugin::new(
220            name,
221            id.clone(),
222            DynamicPluginSource::ByPaths(paths),
223            required,
224        );
225
226        if self.get_plugin_index(&id).is_some() {
227            tracing::warn!(
228                "Duplicate plugin with ID: {id}, only the last declared one will be loaded"
229            )
230        }
231
232        self.plugins.push(PluginRecord::new(loader));
233        Ok(self.plugins.last_mut().unwrap())
234    }
235
236    fn get_plugin_index(&self, id: &str) -> Option<usize> {
237        self.plugins.iter().position(|p| p.id() == id)
238    }
239
240    /// Lists all plugins
241    pub fn declared_plugins_iter(
242        &self,
243    ) -> impl Iterator<Item = &dyn DeclaredPlugin<StartArgs, Instance>> + '_ {
244        self.plugins
245            .iter()
246            .map(|p| p as &dyn DeclaredPlugin<StartArgs, Instance>)
247    }
248
249    /// Lists all plugins mutable
250    pub fn declared_plugins_iter_mut(
251        &mut self,
252    ) -> impl Iterator<Item = &mut dyn DeclaredPlugin<StartArgs, Instance>> + '_ {
253        self.plugins
254            .iter_mut()
255            .map(|p| p as &mut dyn DeclaredPlugin<StartArgs, Instance>)
256    }
257
258    /// Lists the loaded plugins
259    pub fn loaded_plugins_iter(
260        &self,
261    ) -> impl Iterator<Item = &dyn LoadedPlugin<StartArgs, Instance>> + '_ {
262        self.declared_plugins_iter().filter_map(|p| p.loaded())
263    }
264
265    /// Lists the loaded plugins mutable
266    pub fn loaded_plugins_iter_mut(
267        &mut self,
268    ) -> impl Iterator<Item = &mut dyn LoadedPlugin<StartArgs, Instance>> + '_ {
269        self.declared_plugins_iter_mut()
270            .filter_map(|p| p.loaded_mut())
271    }
272
273    /// Lists the started plugins
274    pub fn started_plugins_iter(
275        &self,
276    ) -> impl Iterator<Item = &dyn StartedPlugin<StartArgs, Instance>> + '_ {
277        self.loaded_plugins_iter().filter_map(|p| p.started())
278    }
279
280    /// Lists the started plugins mutable
281    pub fn started_plugins_iter_mut(
282        &mut self,
283    ) -> impl Iterator<Item = &mut dyn StartedPlugin<StartArgs, Instance>> + '_ {
284        self.loaded_plugins_iter_mut()
285            .filter_map(|p| p.started_mut())
286    }
287
288    /// Returns single plugin record by id
289    pub fn plugin(&self, id: &str) -> Option<&dyn DeclaredPlugin<StartArgs, Instance>> {
290        let index = self.get_plugin_index(id)?;
291        Some(&self.plugins[index])
292    }
293
294    /// Returns mutable plugin record by id
295    pub fn plugin_mut(&mut self, id: &str) -> Option<&mut dyn DeclaredPlugin<StartArgs, Instance>> {
296        let index = self.get_plugin_index(id)?;
297        Some(&mut self.plugins[index])
298    }
299
300    /// Returns loaded plugin record by id
301    pub fn loaded_plugin(&self, id: &str) -> Option<&dyn LoadedPlugin<StartArgs, Instance>> {
302        self.plugin(id)?.loaded()
303    }
304
305    /// Returns mutable loaded plugin record by id
306    pub fn loaded_plugin_mut(
307        &mut self,
308        id: &str,
309    ) -> Option<&mut dyn LoadedPlugin<StartArgs, Instance>> {
310        self.plugin_mut(id)?.loaded_mut()
311    }
312
313    /// Returns started plugin record by id
314    pub fn started_plugin(&self, id: &str) -> Option<&dyn StartedPlugin<StartArgs, Instance>> {
315        self.loaded_plugin(id)?.started()
316    }
317
318    /// Returns mutable started plugin record by id
319    pub fn started_plugin_mut(
320        &mut self,
321        id: &str,
322    ) -> Option<&mut dyn StartedPlugin<StartArgs, Instance>> {
323        self.loaded_plugin_mut(id)?.started_mut()
324    }
325}
326
327impl<StartArgs: PluginStartArgs + 'static, Instance: PluginInstance + 'static> PluginControl
328    for PluginsManager<StartArgs, Instance>
329{
330    fn plugins_status(&self, names: &keyexpr) -> Vec<PluginStatusRec<'_>> {
331        tracing::debug!(
332            "Plugin manager with prefix `{}` : requested plugins_status {:?}",
333            self.default_lib_prefix,
334            names
335        );
336        let mut plugins = Vec::new();
337        for plugin in self.declared_plugins_iter() {
338            let id = unsafe { keyexpr::from_str_unchecked(plugin.id()) };
339            if names.includes(id) {
340                let status = PluginStatusRec::new(plugin.as_status());
341                plugins.push(status);
342            }
343            // for running plugins append their subplugins prepended with the running plugin name
344            if let Some(plugin) = plugin.loaded() {
345                if let Some(plugin) = plugin.started() {
346                    if let [names, ..] = names.strip_prefix(id)[..] {
347                        plugins.append(
348                            &mut plugin
349                                .instance()
350                                .plugins_status(names)
351                                .into_iter()
352                                .map(|s| s.prepend_name(id))
353                                .collect(),
354                        );
355                    }
356                }
357            }
358        }
359        plugins
360    }
361}