Skip to content

Commit 063bd77

Browse files
committed
feat: Various minor improvements
- Exit early if no plugin initializes successfully. - Moved the log CLI arguments to the bottom and simplified argument help messages. - Further `TaskTracker` adoption. - Moved plugin `Uuid` into `AvailablePlugin`. - Moved service registration functions into service sub modules. - Temporarily removed the old deregistration implementation. - Naming and ordering consistency.
1 parent 83c53a7 commit 063bd77

17 files changed

Lines changed: 614 additions & 666 deletions

File tree

src/cli.rs

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -24,13 +24,10 @@ impl From<String> for CliLogParametersFileRotation {
2424
#[derive(Parser)]
2525
#[command(about, long_about = None, version, author)]
2626
pub struct Cli {
27-
#[command(flatten)]
28-
pub log_parameters: CliLogParameters,
29-
3027
#[arg(default_value = "./config.yaml", short, long, value_name = "FILE PATH", help = "The path to the program its configuration file", long_help = None)]
3128
pub config_file: PathBuf,
3229

33-
#[arg(default_value = "./.env", short, long, value_name = "FILE PATH", help = "The path to an env file, used by the program its configuration file for env var interpolation", long_help = None)]
30+
#[arg(default_value = "./.env", short, long, value_name = "FILE PATH", help = "The path to the program its env file", long_help = None)]
3431
pub env_file: PathBuf,
3532

3633
#[arg(default_value = "./plugins", short, long, value_name = "DIRECTORY PATH", help = "The path to the program its plugin directory", long_help = None)]
@@ -41,6 +38,9 @@ pub struct Cli {
4138

4239
#[arg(default_value_t = false, short, long, help = "Run in restricted mode, in this case plugin permissions are opt in", long_help = None)]
4340
pub restricted: bool,
41+
42+
#[command(flatten)]
43+
pub log_parameters: CliLogParameters,
4444
}
4545

4646
#[derive(Args)]

src/config.rs

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ use std::{collections::HashMap, fs, path::Path};
66
use anyhow::Result;
77
use serde::Deserialize;
88
use tracing::info;
9+
use uuid::Uuid;
910

1011
use crate::config::{plugins::ConfigPlugin, services::ConfigServices};
1112

@@ -14,13 +15,18 @@ pub mod services;
1415

1516
#[derive(Deserialize)]
1617
pub struct Config {
18+
#[serde(default = "Config::default_name")]
1719
pub name: String,
1820
#[serde(default)]
1921
pub services: ConfigServices,
2022
pub plugins: HashMap<String, ConfigPlugin>,
2123
}
2224

2325
impl Config {
26+
fn default_name() -> String {
27+
Uuid::new_v4().to_string()
28+
}
29+
2430
pub fn new(file_path: &Path, restricted: bool) -> Result<Self> {
2531
info!("Loading and parsing the config file");
2632

src/main.rs

Lines changed: 21 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -101,21 +101,21 @@ async fn main() -> Result<ExitCode> {
101101
let database = database::new(&cli.database_directory)?;
102102

103103
let message_handler = message_handler(
104-
database.clone(),
105104
Arc::new(RwLock::new(Some(channels.core.runtime_tx))),
106105
Arc::new(RwLock::new(channels.core.job_scheduler_tx)),
107106
Arc::new(RwLock::new(channels.core.discord_tx)),
108-
Arc::new(shutdown_signal_listener),
109107
channels.core.rx,
108+
database.clone(),
109+
Arc::new(shutdown_signal_listener),
110110
);
111111

112112
let setup_result = setup(
113113
cli.plugin_directory,
114-
database,
115-
channels.services,
116-
channels.runtime,
117114
config,
118115
secrets,
116+
channels.services,
117+
channels.runtime,
118+
database,
119119
)
120120
.await;
121121

@@ -129,18 +129,23 @@ async fn main() -> Result<ExitCode> {
129129
}
130130

131131
fn message_handler(
132-
database: Database,
133132
runtime_tx: Arc<RwLock<Option<UnboundedSender<RuntimeMessages>>>>,
134133
job_scheduler_tx: Arc<RwLock<Option<UnboundedSender<JobSchedulerMessages>>>>,
135134
discord_tx: Arc<RwLock<Option<UnboundedSender<DiscordMessages>>>>,
136-
shutdown_signal_listener: Arc<JoinHandle<()>>,
137135
mut rx: UnboundedReceiver<CoreMessages>,
136+
database: Database,
137+
shutdown_signal_listener: Arc<JoinHandle<()>>,
138138
) -> JoinHandle<Result<()>> {
139139
debug!("Starting the message handler");
140140

141141
tokio::spawn(async move {
142142
while let Some(core_message) = rx.recv().await {
143143
match core_message {
144+
CoreMessages::Runtime(runtime_message) => {
145+
if let Some(runtime_tx) = runtime_tx.read().await.as_ref() {
146+
runtime_tx.send(runtime_message).unwrap();
147+
}
148+
}
144149
CoreMessages::JobScheduler(job_scheduler_message) => {
145150
if let Some(job_scheduler_tx) = job_scheduler_tx.read().await.as_ref() {
146151
job_scheduler_tx.send(job_scheduler_message).unwrap();
@@ -151,18 +156,13 @@ fn message_handler(
151156
discord_tx.send(discord_message).unwrap();
152157
}
153158
}
154-
CoreMessages::Runtime(runtime_message) => {
155-
if let Some(runtime_tx) = runtime_tx.read().await.as_ref() {
156-
runtime_tx.send(runtime_message).unwrap();
157-
}
158-
}
159159
CoreMessages::Shutdown(shutdown_kind) => {
160160
tokio::spawn(shutdown(
161-
shutdown_kind,
162161
runtime_tx.clone(),
163162
job_scheduler_tx.clone(),
164163
discord_tx.clone(),
165164
shutdown_signal_listener.clone(),
165+
shutdown_kind,
166166
));
167167
}
168168
}
@@ -174,27 +174,27 @@ fn message_handler(
174174

175175
async fn setup(
176176
plugin_directory_path: PathBuf,
177-
database: Database,
178-
service_channels: ChannelsServices,
179-
runtime_channels: ChannelsRuntime,
180177
config: Config,
181178
secrets: Secrets,
179+
service_channels: ChannelsServices,
180+
runtime_channels: ChannelsRuntime,
181+
database: Database,
182182
) -> Result<()> {
183183
let config_name = Arc::new(config.name);
184184

185185
let available_plugins = registry::get_plugins(
186186
&plugin_directory_path,
187-
database.clone(),
188187
config_name.clone(),
189188
config.plugins,
189+
database.clone(),
190190
)
191191
.await?;
192192

193193
services::setup(
194194
config.services,
195195
secrets.services,
196-
database.clone(),
197196
service_channels,
197+
database.clone(),
198198
)
199199
.await?;
200200

@@ -204,9 +204,9 @@ async fn setup(
204204
.initialize_plugins(
205205
plugin_directory_path,
206206
config_name,
207-
available_plugins,
208-
database,
209207
runtime_channels.core_tx,
208+
database,
209+
available_plugins,
210210
)
211211
.await?;
212212

@@ -293,11 +293,11 @@ fn shutdown_signal_listener(core_tx: UnboundedSender<CoreMessages>) -> JoinHandl
293293
}
294294

295295
async fn shutdown(
296-
shutdown_kind: Shutdown,
297296
runtime_tx: Arc<RwLock<Option<UnboundedSender<RuntimeMessages>>>>,
298297
job_scheduler_tx: Arc<RwLock<Option<UnboundedSender<JobSchedulerMessages>>>>,
299298
discord_tx: Arc<RwLock<Option<UnboundedSender<DiscordMessages>>>>,
300299
shutdown_signal_listener: Arc<JoinHandle<()>>,
300+
shutdown_kind: Shutdown,
301301
) {
302302
let mut shutdown_guard = SHUTDOWN.write().await;
303303

src/registry.rs

Lines changed: 28 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -27,18 +27,18 @@ static DEFAULT_NAMESPACE_ID: &str = "wpbs-rs";
2727
#[hotpath::measure]
2828
pub async fn get_plugins(
2929
plugin_directory_path: &Path,
30-
database: Database,
3130
config_name: Arc<String>,
3231
config_plugins: HashMap<String, ConfigPlugin>,
33-
) -> Result<Vec<(Uuid, AvailablePlugin)>> {
32+
database: Database,
33+
) -> Result<Vec<AvailablePlugin>> {
3434
info!("Getting all plugins from their respective registries");
3535

3636
let caching_client =
3737
create_registry_client(&plugin_directory_path.join("binaries").join("remote")).await?;
3838

3939
let mut available_plugins = Vec::new();
4040

41-
let mut plugin_tasks: Vec<JoinHandle<Result<(Uuid, AvailablePlugin)>>> = Vec::new();
41+
let mut plugin_tasks: Vec<JoinHandle<Result<AvailablePlugin>>> = Vec::new();
4242

4343
let plugins_keyspace = database.keyspace("plugins", KeyspaceCreateOptions::default)?;
4444

@@ -57,41 +57,38 @@ pub async fn get_plugins(
5757
if namespace_id == "local" {
5858
get_local_plugin(plugin_directory_path, &plugin_id, &plugin_version).await?;
5959

60-
let uuid = get_plugin_uuid(&plugins_keyspace, &config_name, &plugin_user_id)?;
61-
62-
return Ok((
63-
uuid,
64-
AvailablePlugin {
65-
namespace_id,
66-
plugin_id,
67-
version: plugin_version,
68-
content_digest: None,
69-
user_id: plugin_user_id,
70-
environment: plugin_config.environment,
71-
permissions: plugin_config.permissions,
72-
settings: plugin_config.settings,
73-
},
74-
));
75-
}
76-
77-
let release =
78-
fetch_plugin(caching_client, &namespace_id, &plugin_id, &plugin_version).await?;
79-
80-
let uuid = get_plugin_uuid(&plugins_keyspace, &config_name, &plugin_user_id)?;
60+
let plugin_uuid =
61+
get_plugin_uuid(&plugins_keyspace, &config_name, &plugin_user_id)?;
8162

82-
Ok((
83-
uuid,
84-
AvailablePlugin {
63+
return Ok(AvailablePlugin {
64+
plugin_uuid,
8565
namespace_id,
8666
plugin_id,
87-
version: release.version,
88-
content_digest: Some(release.content_digest),
67+
version: plugin_version,
68+
content_digest: None,
8969
user_id: plugin_user_id,
9070
permissions: plugin_config.permissions,
9171
environment: plugin_config.environment,
9272
settings: plugin_config.settings,
93-
},
94-
))
73+
});
74+
}
75+
76+
let release =
77+
fetch_plugin(caching_client, &namespace_id, &plugin_id, &plugin_version).await?;
78+
79+
let plugin_uuid = get_plugin_uuid(&plugins_keyspace, &config_name, &plugin_user_id)?;
80+
81+
Ok(AvailablePlugin {
82+
plugin_uuid,
83+
namespace_id,
84+
plugin_id,
85+
version: release.version,
86+
content_digest: Some(release.content_digest),
87+
user_id: plugin_user_id,
88+
permissions: plugin_config.permissions,
89+
environment: plugin_config.environment,
90+
settings: plugin_config.settings,
91+
})
9592
}));
9693
}
9794

src/registry/plugins.rs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,12 +4,14 @@
44
use std::collections::HashMap;
55

66
use semver::Version;
7+
use uuid::Uuid;
78
use wasm_pkg_client::ContentDigest;
89
use yaml_serde::Value;
910

1011
use crate::config::plugins::permissions::PluginPermissions;
1112

1213
pub struct AvailablePlugin {
14+
pub plugin_uuid: Uuid,
1315
pub namespace_id: String,
1416
pub plugin_id: String,
1517
pub version: Version,

0 commit comments

Comments
 (0)