refactor: reorganize imports and add lint allows
CI / ${{ matrix.crate }} (workshop-engine) (push) Failing after 13m13s
CI / ${{ matrix.crate }} (workshop-baker) (push) Failing after 28m54s

This commit is contained in:
Catty Steve
2026-04-27 14:31:11 +08:00
parent 50b42e654c
commit 20e2ab4224
34 changed files with 296 additions and 259 deletions
+8
View File
@@ -58,6 +58,14 @@ globset = "0.4.18"
[dev-dependencies] [dev-dependencies]
criterion = "0.5" criterion = "0.5"
[lints.rust]
dead_code = "allow"
unreachable_code = "allow"
[lints.clippy]
inherent_to_string = "allow"
non_canonical_partial_ord_impl = "allow"
[[bench]] [[bench]]
name = "parser_bench" name = "parser_bench"
harness = false harness = false
+15 -15
View File
@@ -1,26 +1,26 @@
use workshop_engine::{Engine, EventSender, ExecutionContext};
use crate::bake::parser::Function;
use crate::bake::constant::{DEFAULT_WORKSPACE, TRIVIAL_SCRIPT_NAME}; use crate::bake::constant::{DEFAULT_WORKSPACE, TRIVIAL_SCRIPT_NAME};
use crate::bake::error::BakeError; use crate::bake::error::BakeError;
use crate::bake::parser::Function;
use crate::cli::Cli; use crate::cli::Cli;
use crate::prebake;
use crate::prebake::PrebakeConfig; use crate::prebake::PrebakeConfig;
use crate::prebake::security::get_drop_after; use crate::prebake::security::get_drop_after;
use crate::prebake::stage::PrebakeStage; use crate::prebake::stage::PrebakeStage;
use crate::prebake;
use std::path::{Path, PathBuf}; use std::path::{Path, PathBuf};
use workshop_engine::{Engine, EventSender, ExecutionContext};
mod builder; mod builder;
pub mod constant; pub mod constant;
mod decorator; mod decorator;
pub mod parser;
pub mod error; pub mod error;
pub mod parser;
mod schedule; mod schedule;
pub async fn bake( pub async fn bake(
script_path: &Path, script_path: &Path,
bake_base_path: &Path, bake_base_path: &Path,
prebake_path: &Path, prebake_path: &Path,
cli: &Cli, _cli: &Cli,
ctx: &mut ExecutionContext, ctx: &mut ExecutionContext,
event_tx: EventSender, event_tx: EventSender,
) -> Result<(), BakeError> { ) -> Result<(), BakeError> {
@@ -43,10 +43,10 @@ pub async fn bake(
let use_template = true; let use_template = true;
let engine = Engine::new(); let engine = Engine::new();
if let Some(env) = &prebake.envvars { if let Some(env) = &prebake.envvars
if let Some(prebake_env) = &env.bake { && let Some(prebake_env) = &env.bake
ctx.env_vars.extend(prebake_env.clone()); {
} ctx.env_vars.extend(prebake_env.clone());
} }
if !has_pipeline || trivial { if !has_pipeline || trivial {
@@ -57,7 +57,7 @@ pub async fn bake(
}; };
let script = let script =
builder::build_script(&function, bake_base_path, &workspace, None, use_template)?; builder::build_script(&function, bake_base_path, &workspace, None, use_template)?;
let result = engine.execute_script(&script, &ctx, &event_tx).await; let result = engine.execute_script(&script, ctx, &event_tx).await;
result?; result?;
} else { } else {
let sorted_functions = schedule::sort_function(functions.pipeline_functions)?; let sorted_functions = schedule::sort_function(functions.pipeline_functions)?;
@@ -72,7 +72,7 @@ pub async fn bake(
None, None,
use_template, use_template,
)?; )?;
let result = engine.execute_script(&script, &ctx, &event_tx).await; let result = engine.execute_script(&script, ctx, &event_tx).await;
result?; result?;
} }
} }
@@ -96,10 +96,10 @@ fn get_workspace(prebake_config: &PrebakeConfig) -> PathBuf {
} }
fn privileged(prebake_config: &PrebakeConfig) -> bool { fn privileged(prebake_config: &PrebakeConfig) -> bool {
if let Some(bootstrap) = &prebake_config.bootstrap { if let Some(bootstrap) = &prebake_config.bootstrap
if bootstrap.user == "root" { && bootstrap.user == "root"
return true; {
} return true;
} }
let drop_after = get_drop_after(&prebake_config.security).unwrap_or_default(); let drop_after = get_drop_after(&prebake_config.security).unwrap_or_default();
PrebakeStage::Never == drop_after PrebakeStage::Never == drop_after
+2 -2
View File
@@ -1,7 +1,7 @@
use crate::bake::decorator::Decorator;
use crate::bake::parser::Function;
use crate::bake::constant::{PIPELINE_SKIP_ERRORCODE, TRIVIAL_SCRIPT_NAME}; use crate::bake::constant::{PIPELINE_SKIP_ERRORCODE, TRIVIAL_SCRIPT_NAME};
use crate::bake::decorator::Decorator;
use crate::bake::error::BakeError; use crate::bake::error::BakeError;
use crate::bake::parser::Function;
use minijinja::{Environment, context}; use minijinja::{Environment, context};
use std::fs::Permissions; use std::fs::Permissions;
use std::os::unix::fs::PermissionsExt; use std::os::unix::fs::PermissionsExt;
+1 -4
View File
@@ -2,10 +2,7 @@ use thiserror::Error;
use workshop_engine::{ use workshop_engine::{
ExecutionError, ExecutionError,
error::{ error::{EXITCODE_IO_ERROR, EXITCODE_PARSE_ERROR, HasExitCode},
EXITCODE_IO_ERROR, EXITCODE_PARSE_ERROR,
HasExitCode,
},
}; };
#[derive(Error, Debug)] #[derive(Error, Debug)]
+5 -5
View File
@@ -9,12 +9,12 @@ pub mod types;
use crate::cli::Cli; use crate::cli::Cli;
use crate::finalize::constant::FINALIZE_STANDALONE_PLUGIN_PATH; use crate::finalize::constant::FINALIZE_STANDALONE_PLUGIN_PATH;
use workshop_engine::EventSender;
use crate::finalize::plugin::fetch::{self, FetchArgument}; use crate::finalize::plugin::fetch::{self, FetchArgument};
use std::path::PathBuf;
pub use config::*; pub use config::*;
pub use error::FinalizeError; pub use error::FinalizeError;
use std::path::Path; use std::path::Path;
use std::path::PathBuf;
use workshop_engine::EventSender;
pub fn parse(path: &Path) -> Result<FinalizeConfig, FinalizeError> { pub fn parse(path: &Path) -> Result<FinalizeConfig, FinalizeError> {
let content = std::fs::read_to_string(path).map_err(FinalizeError::IoError)?; let content = std::fs::read_to_string(path).map_err(FinalizeError::IoError)?;
@@ -24,7 +24,7 @@ pub fn parse(path: &Path) -> Result<FinalizeConfig, FinalizeError> {
pub async fn finalize( pub async fn finalize(
finalize_path: &Path, finalize_path: &Path,
cli: &Cli, _cli: &Cli,
ctx: &mut workshop_engine::ExecutionContext, ctx: &mut workshop_engine::ExecutionContext,
event_tx: EventSender, event_tx: EventSender,
) -> Result<(), FinalizeError> { ) -> Result<(), FinalizeError> {
@@ -55,7 +55,7 @@ pub async fn finalize(
checksum: config.checksum.clone(), checksum: config.checksum.clone(),
shallow: config.shallow, shallow: config.shallow,
}; };
fetch::fetch_plugin(&source, &name, &plugin_path, fetch_argument).await?; fetch::fetch_plugin(source, name, &plugin_path, fetch_argument).await?;
config.runtime.fspath = Some(plugin_path.join(name)); config.runtime.fspath = Some(plugin_path.join(name));
} }
} }
@@ -75,7 +75,7 @@ pub async fn finalize(
// Determine if we should continue with artifacts/latehook // Determine if we should continue with artifacts/latehook
// If prior stages failed, we still send notification but skip artifacts // If prior stages failed, we still send notification but skip artifacts
let prior_failed = earlyhook_result.is_err(); let _prior_failed = earlyhook_result.is_err();
// Prepare artifacts // Prepare artifacts
// TODO: artifact stage implementation // TODO: artifact stage implementation
+1 -2
View File
@@ -3,13 +3,13 @@
//! This module defines the configuration schema for the finalize stage, //! This module defines the configuration schema for the finalize stage,
//! which handles artifact collection, publishing, and notifications. //! which handles artifact collection, publishing, and notifications.
use workshop_engine::{CustomCommand, deserialize_duration_ms};
use crate::finalize::types::compression::CompressionMethod; use crate::finalize::types::compression::CompressionMethod;
use semver::{Version, VersionReq}; use semver::{Version, VersionReq};
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use serde_yaml::Value; use serde_yaml::Value;
use std::collections::HashMap; use std::collections::HashMap;
use std::path::PathBuf; use std::path::PathBuf;
use workshop_engine::{CustomCommand, deserialize_duration_ms};
/// Version requirement for finalize.yml schema compatibility /// Version requirement for finalize.yml schema compatibility
pub const VERSION_REQUIREMENT: &str = "^0.0.1"; pub const VERSION_REQUIREMENT: &str = "^0.0.1";
@@ -253,7 +253,6 @@ impl Default for ArtifactDef {
/// Artifact /// Artifact
/// TODO: Under heavy refactor /// TODO: Under heavy refactor
/// Cleanup configuration /// Cleanup configuration
#[derive(Debug, Deserialize, Serialize, Clone, Default)] #[derive(Debug, Deserialize, Serialize, Clone, Default)]
pub struct CleanupConfig { pub struct CleanupConfig {
+4 -4
View File
@@ -1,5 +1,5 @@
use crate::finalize::constant::{FINALIZE_PLUGIN_MANIFEST, FINALIZE_PLUGIN_MANIFEST_EXTENSION};
use crate::finalize::RawPluginMap; use crate::finalize::RawPluginMap;
use crate::finalize::constant::{FINALIZE_PLUGIN_MANIFEST, FINALIZE_PLUGIN_MANIFEST_EXTENSION};
use crate::finalize::plugin::metadata::{PluginMetadata, PluginType}; use crate::finalize::plugin::metadata::{PluginMetadata, PluginType};
use crate::{EventSender, ExecutionContext, finalize::error::PluginError}; use crate::{EventSender, ExecutionContext, finalize::error::PluginError};
use async_trait::async_trait; use async_trait::async_trait;
@@ -134,7 +134,7 @@ fn load_manifest(path: &Path, name: &str, suffix: &str) -> Result<String, Plugin
reason: format!( reason: format!(
"Failed to read plugin manifest from {}: {}", "Failed to read plugin manifest from {}: {}",
manifest_path.display(), manifest_path.display(),
e.to_string() e
), ),
}) })
} }
@@ -189,10 +189,10 @@ fn get_entry(plugin_type: PluginType, name: &str) -> Result<Box<dyn AsyncPluginF
"Plugin {} has unknown plugin type and cannot be inferred from entrypoint.", "Plugin {} has unknown plugin type and cannot be inferred from entrypoint.",
name name
); );
return Err(PluginError::PluginUnavailable { Err(PluginError::PluginUnavailable {
name: name.to_string(), name: name.to_string(),
reason: "Unknown plugin type and cannot be inferred from entrypoint".to_string(), reason: "Unknown plugin type and cannot be inferred from entrypoint".to_string(),
}); })
} }
} }
} }
+5 -8
View File
@@ -1,9 +1,6 @@
// Code in this module PARTIALLY OR FULLY utilized AI Coding Agent. // Code in this module PARTIALLY OR FULLY utilized AI Coding Agent.
use crate::finalize::fetch; use crate::finalize::fetch;
use crate::finalize::{ use crate::finalize::plugin::{PluginError, RawPluginMap};
PluginConfig,
plugin::{PluginError, RawPluginMap},
};
use std::collections::HashSet; use std::collections::HashSet;
use std::path::Path; use std::path::Path;
@@ -66,14 +63,14 @@ pub async fn fetch_plugin(
// TODO: Support client mode // TODO: Support client mode
// NOTE: In standalone mode, we just assume path is well-defined and utilize // NOTE: In standalone mode, we just assume path is well-defined and utilize
log::info!("Fetching file plugin: {}", url); log::info!("Fetching file plugin: {}", url);
return Ok(()); Ok(())
} }
URLScheme::Unknown => { URLScheme::Unknown => {
log::error!("Unknown URL scheme in plugin URL: {}", url); log::error!("Unknown URL scheme in plugin URL: {}", url);
return Err(PluginError::GeneralError(format!( Err(PluginError::GeneralError(format!(
"Unknown URL scheme in plugin URL: {}", "Unknown URL scheme in plugin URL: {}",
url url
))); )))
} }
} }
} }
@@ -104,7 +101,7 @@ pub async fn fetch_plugins(
checksum: config.checksum.clone(), checksum: config.checksum.clone(),
shallow: config.shallow, shallow: config.shallow,
}; };
fetch::fetch_plugin(&source, &name, &plugin_path, fetch_argument).await?; fetch::fetch_plugin(source, name, plugin_path, fetch_argument).await?;
log::debug!("Storing plugin {} to {}", name, plugin_path.display()); log::debug!("Storing plugin {} to {}", name, plugin_path.display());
config.runtime.fspath = Some(plugin_path.join(name)); config.runtime.fspath = Some(plugin_path.join(name));
} }
@@ -1,6 +1,6 @@
// Code in this module PARTIALLY OR FULLY utilized AI Coding Agent. // Code in this module PARTIALLY OR FULLY utilized AI Coding Agent.
use crate::finalize::plugin::PluginError; use crate::finalize::plugin::PluginError;
use std::path::PathBuf; use std::path::{Path, PathBuf};
use tokio::task; use tokio::task;
/// Compression type for tar archives. /// Compression type for tar archives.
@@ -48,7 +48,7 @@ impl CompressionType {
/// Removes existing destination directory if present, creates a fresh directory, /// Removes existing destination directory if present, creates a fresh directory,
/// then unpacks the archive using the appropriate decompression based on `compression`. /// then unpacks the archive using the appropriate decompression based on `compression`.
pub async fn extract_tar( pub async fn extract_tar(
archive: &PathBuf, archive: &Path,
destdir: &PathBuf, destdir: &PathBuf,
compression: CompressionType, compression: CompressionType,
) -> Result<(), PluginError> { ) -> Result<(), PluginError> {
@@ -61,7 +61,7 @@ pub async fn extract_tar(
.await .await
.map_err(|e| PluginError::GeneralError(format!("Failed to create directory: {}", e)))?; .map_err(|e| PluginError::GeneralError(format!("Failed to create directory: {}", e)))?;
let archive = archive.clone(); let archive = archive.to_path_buf();
let destdir = destdir.clone(); let destdir = destdir.clone();
task::spawn_blocking(move || { task::spawn_blocking(move || {
@@ -69,7 +69,7 @@ fn extract_archive_extension(url: &str) -> &str {
} else if lower.ends_with(".tar") { } else if lower.ends_with(".tar") {
".tar" ".tar"
} else { } else {
".tar" ".unknown"
} }
} }
@@ -86,23 +86,21 @@ pub async fn fetch_http_plugin(
download_to_file(url, &temp_file, name).await?; download_to_file(url, &temp_file, name).await?;
if let Some(checksum) = checksum_str { if let Some(checksum) = checksum_str
if !checksum.is_empty() { && !checksum.is_empty()
let bytes = tokio::fs::read(&temp_file).await.map_err(|e| { {
PluginError::GeneralError(format!("Failed to read downloaded file: {}", e)) let bytes = tokio::fs::read(&temp_file).await.map_err(|e| {
})?; PluginError::GeneralError(format!("Failed to read downloaded file: {}", e))
})?;
let checksum_type = let checksum_type = detect_checksum_type(checksum).map_err(PluginError::GeneralError)?;
detect_checksum_type(checksum).map_err(|e| PluginError::GeneralError(e))?;
if let Some(ct) = checksum_type { if let Some(ct) = checksum_type {
verify_checksum(&bytes, checksum, ct).map_err(|e| PluginError::GeneralError(e))?; verify_checksum(&bytes, checksum, ct).map_err(PluginError::GeneralError)?;
}
} }
} }
let compression = let compression = CompressionType::from_path(&temp_file).map_err(PluginError::GeneralError)?;
CompressionType::from_path(&temp_file).map_err(|e| PluginError::GeneralError(e))?;
let destdir = basedir.join(name); let destdir = basedir.join(name);
extract_tar(&temp_file, &destdir, compression).await?; extract_tar(&temp_file, &destdir, compression).await?;
+18 -18
View File
@@ -10,22 +10,22 @@ pub struct InternalPlugin {
} }
pub fn get_internal_plugin() -> Vec<InternalPlugin> { pub fn get_internal_plugin() -> Vec<InternalPlugin> {
let mut internal_plugins: Vec<InternalPlugin> = Vec::new(); vec![
internal_plugins.push(InternalPlugin { InternalPlugin {
name: "in-site-notify", name: "in-site-notify",
entry: Box::new(insitenotify::InsiteNotify), entry: Box::new(insitenotify::InsiteNotify),
}); },
internal_plugins.push(InternalPlugin { InternalPlugin {
name: "mail", name: "mail",
entry: Box::new(mail::Mail), entry: Box::new(mail::Mail),
}); },
internal_plugins.push(InternalPlugin { InternalPlugin {
name: "satori", name: "satori",
entry: Box::new(satori::Satori), entry: Box::new(satori::Satori),
}); },
internal_plugins.push(InternalPlugin { InternalPlugin {
name: "webhook", name: "webhook",
entry: Box::new(webhook::Webhook), entry: Box::new(webhook::Webhook),
}); },
internal_plugins ]
} }
@@ -1,6 +1,6 @@
// Code in this module PARTIALLY OR FULLY utilized AI Coding Agent. // Code in this module PARTIALLY OR FULLY utilized AI Coding Agent.
use workshop_engine::deserialize_duration_ms;
use serde::Deserialize; use serde::Deserialize;
use workshop_engine::deserialize_duration_ms;
/// Email address representation supporting plain string or structured object. /// Email address representation supporting plain string or structured object.
#[derive(Debug, Deserialize, Clone)] #[derive(Debug, Deserialize, Clone)]
+4 -4
View File
@@ -94,7 +94,7 @@ fn convert_argument(
envvars.insert(prefix.to_string(), value.to_string()); envvars.insert(prefix.to_string(), value.to_string());
} }
serde_yaml::Value::String(value) => { serde_yaml::Value::String(value) => {
envvars.insert(prefix.to_string(), format!("{}", value)); envvars.insert(prefix.to_string(), value.to_string());
} }
serde_yaml::Value::Number(value) => { serde_yaml::Value::Number(value) => {
envvars.insert(prefix.to_string(), value.to_string()); envvars.insert(prefix.to_string(), value.to_string());
@@ -112,12 +112,12 @@ fn convert_argument(
return envvars; return envvars;
} }
let mut array = String::new(); let mut array = String::new();
array.push_str("("); array.push('(');
let mut has_mapping = false; let mut has_mapping = false;
let mut _envvars: HashMap<String, String> = HashMap::new(); let mut _envvars: HashMap<String, String> = HashMap::new();
for (index, item) in seq.into_iter().enumerate() { for (index, item) in seq.into_iter().enumerate() {
if index > 0 { if index > 0 {
array.push_str(" "); array.push(' ');
} }
let new_prefix = format!("{}_{}", prefix, index); let new_prefix = format!("{}_{}", prefix, index);
let subenvvars = convert_argument(new_prefix.as_str(), item.clone(), true); let subenvvars = convert_argument(new_prefix.as_str(), item.clone(), true);
@@ -136,7 +136,7 @@ fn convert_argument(
if has_mapping { if has_mapping {
envvars.extend(_envvars); envvars.extend(_envvars);
} else { } else {
array.push_str(")"); array.push(')');
envvars.insert(prefix.to_string(), array); envvars.insert(prefix.to_string(), array);
} }
} }
+1 -1
View File
@@ -1,7 +1,7 @@
use crate::ExecutionContext; use crate::ExecutionContext;
use workshop_engine::{DeltaExecutionContext, Engine, EventSender, ExecutionError, CustomCommand};
use crate::finalize::plugin::{PluginMap, call_named_plugin}; use crate::finalize::plugin::{PluginMap, call_named_plugin};
use std::path::PathBuf; use std::path::PathBuf;
use workshop_engine::{CustomCommand, DeltaExecutionContext, Engine, EventSender, ExecutionError};
/// Executes post-build finalize hooks, supporting both plugin and shell command hooks. /// Executes post-build finalize hooks, supporting both plugin and shell command hooks.
/// ///
+1 -1
View File
@@ -7,7 +7,7 @@ use workshop_baker::{
finalize::finalize, finalize::finalize,
prebake::prebake, prebake::prebake,
}; };
use workshop_engine::{error::HasExitCode, EventSender, ExecutionContext}; use workshop_engine::{EventSender, ExecutionContext, error::HasExitCode};
#[tokio::main] #[tokio::main]
async fn main() { async fn main() {
+105 -88
View File
@@ -1,10 +1,18 @@
use crate::finalize::constant::FINALIZE_STANDALONE_PLUGIN_PATH;
use crate::finalize::FinalizeConfig; use crate::finalize::FinalizeConfig;
use crate::finalize::NotificationConfig; use crate::finalize::NotificationConfig;
use crate::finalize::NotificationTemplate;
use crate::finalize::config::TriggerItemRaw;
use crate::finalize::constant::FINALIZE_STANDALONE_PLUGIN_PATH;
use crate::finalize::parse; use crate::finalize::parse;
use crate::finalize::plugin; use crate::finalize::plugin;
use crate::finalize::plugin::fetch; use crate::finalize::plugin::fetch;
use crate::notify::types::Matchable;
use crate::notify::types::Method; use crate::notify::types::Method;
use crate::notify::types::NOTIFICATION_DEFAULT_FALLIBLE;
use crate::notify::types::NOTIFICATION_DEFAULT_MAX_LIVES;
use crate::notify::types::NOTIFICATION_DEFAULT_RETRY_BACKOFF;
use crate::notify::types::NOTIFICATION_DEFAULT_TIMEOUT_MS;
use crate::notify::types::NOTIFICATION_NOTREADY_DELAY;
use crate::notify::types::NOTIFICATION_TEMPORARY_DELAY_FACTOR; use crate::notify::types::NOTIFICATION_TEMPORARY_DELAY_FACTOR;
use crate::notify::types::NOTIFICATION_TEMPORARY_DELAY_TIME; use crate::notify::types::NOTIFICATION_TEMPORARY_DELAY_TIME;
use crate::notify::types::NOTIFICATION_TEMPORARY_MAX_TIME; use crate::notify::types::NOTIFICATION_TEMPORARY_MAX_TIME;
@@ -15,20 +23,14 @@ use crate::notify::types::NotificationPluginMap;
use crate::notify::types::NotificationRenderer; use crate::notify::types::NotificationRenderer;
use crate::notify::types::NotificationRule; use crate::notify::types::NotificationRule;
use crate::notify::types::NotificationRulesMap; use crate::notify::types::NotificationRulesMap;
use crate::notify::types::ParsedTrigger;
use crate::{EventSender, ExecutionContext}; use crate::{EventSender, ExecutionContext};
use chrono::Utc;
use globset::{Glob, GlobSetBuilder};
use std::collections::HashMap; use std::collections::HashMap;
use std::collections::HashSet; use std::collections::HashSet;
use std::path::Path; use std::path::Path;
use std::path::PathBuf; use std::path::PathBuf;
use crate::notify::types::ParsedTrigger;
use crate::notify::types::Matchable;
use globset::{Glob, GlobSetBuilder};
use crate::finalize::config::TriggerItemRaw;
use crate::finalize::NotificationTemplate;
use crate::notify::types::NOTIFICATION_DEFAULT_FALLIBLE;
use crate::notify::types::NOTIFICATION_DEFAULT_MAX_LIVES;
use crate::notify::types::NOTIFICATION_DEFAULT_TIMEOUT_MS;
use crate::notify::types::NOTIFICATION_DEFAULT_RETRY_BACKOFF;
mod error; mod error;
pub mod types; pub mod types;
@@ -49,8 +51,8 @@ fn extract_plugin(finalize: &FinalizeConfig) -> HashSet<String> {
finalize finalize
.notification .notification
.groups .groups
.iter() .values()
.flat_map(|(_, methods)| methods.keys()) .flat_map(|methods| methods.keys())
.cloned() .cloned()
.collect() .collect()
} }
@@ -58,7 +60,7 @@ fn extract_plugin(finalize: &FinalizeConfig) -> HashSet<String> {
pub async fn notify( pub async fn notify(
finalize_path: &Path, finalize_path: &Path,
ctx: &mut ExecutionContext, ctx: &mut ExecutionContext,
event_tx: EventSender, _event_tx: EventSender,
notifyevent_rx: &mut NotificationEventReceiver, notifyevent_rx: &mut NotificationEventReceiver,
notifyevent_tx: NotificationEventSender, notifyevent_tx: NotificationEventSender,
) -> anyhow::Result<()> { ) -> anyhow::Result<()> {
@@ -84,48 +86,65 @@ pub async fn notify(
let notification_rules = parse_rules(&finalize.notification); let notification_rules = parse_rules(&finalize.notification);
// loop { loop {
// tokio::select! { tokio::select! {
// Some(events) = notifyevent_rx.recv() => { Some(event) = notifyevent_rx.recv() => {
// match notify_handler(&notify_plugins, &finalize.notification, &events){ handle_notification_event(
// Ok(_) => {}, // Notification succeeded / failed with fallible event,
// Err(e) => { &notify_plugins,
// log::debug!("Failed to send notification: {:?}", e); &notification_rules,
// match e { &notifyevent_tx,
// NotifyError::PluginNotAvailable(_) => { )
// // Should keep lives, add static delay and give less priority .await?;
// for mut event in events { }
// event.not_before = Utc::now() + NOTIFY_NOTREADY_DELAY; }
// event.effective_priority = event.effective_priority.saturating_add(1); }
// notifyevent_tx.send(event).await?; }
// }
// }, /// Process a single notification event.
// NotifyError::TemporaryError(_) => { ///
// // Should cost 1 life, add dynamic (exponential) delay /// Delegates to [`notification_handler`] and handles errors:
// for mut event in events { /// - [`NotifyError::PluginNotAvailable`] — reschedule with static delay, lower priority
// event.lives -= 1; /// - [`NotifyError::TemporaryError`] — consume one life, reschedule with exponential backoff
// event.not_before = Utc::now() + notify_delay(&event); /// - [`NotifyError::PermanentError`] — drop the event
// notifyevent_tx.send(event).await?; /// - other — log and drop
// } async fn handle_notification_event(
// }, event: NotificationEvent,
// NotifyError::PermanentError(_) => { plugins: &NotificationPluginMap,
// // Unexpected! rules: &NotificationRulesMap,
// // Should be dropped notifyevent_tx: &NotificationEventSender,
// log::error!("None of the configured notification plugins works."); ) -> anyhow::Result<()> {
// log::error!("Message will be dropped."); match notification_handler(plugins, rules, &event) {
// // TODO: Send audit/metric information Ok(_) => {}
// } Err(e) => {
// error => { log::debug!("Failed to send notification: {:?}", e);
// // Unexpected! match e {
// log::error!("Internal error: notify_handler returned unexpected error: {:?}", error); NotifyError::PluginNotAvailable(_) => {
// } let mut event = event;
// }; event.not_before = Utc::now() + NOTIFICATION_NOTREADY_DELAY;
// } event.effective_priority = event.effective_priority.saturating_add(1);
// }; notifyevent_tx.send(event).await?;
// } }
// } NotifyError::TemporaryError(_) => {
// } let mut event = event;
todo!(); event.lives -= 1;
event.not_before = Utc::now() + notify_delay(&event);
notifyevent_tx.send(event).await?;
}
NotifyError::PermanentError(_) => {
log::error!("None of the configured notification plugins works.");
log::error!("Message will be dropped.");
}
error => {
log::error!(
"Internal error: notify_handler returned unexpected error: {:?}",
error
);
}
}
}
}
Ok(())
} }
fn notification_handler( fn notification_handler(
@@ -133,7 +152,7 @@ fn notification_handler(
config: &NotificationRulesMap, config: &NotificationRulesMap,
event: &NotificationEvent, event: &NotificationEvent,
) -> Result<(), NotifyError> { ) -> Result<(), NotifyError> {
for (priority, ruleset) in config.iter() { for (_priority, ruleset) in config.iter() {
for (method, rules) in ruleset.iter() { for (method, rules) in ruleset.iter() {
if !event.target.contains(method) { if !event.target.contains(method) {
log::trace!("target does not contain method name: {}", method); log::trace!("target does not contain method name: {}", method);
@@ -173,23 +192,21 @@ fn notification_handler(
fn parse_trigger_item_raw(item: &TriggerItemRaw) -> Vec<ParsedTrigger> { fn parse_trigger_item_raw(item: &TriggerItemRaw) -> Vec<ParsedTrigger> {
match item { match item {
TriggerItemRaw::Simple(s) => { TriggerItemRaw::Simple(s) => match s.as_str() {
match s.as_str() { "on_success" => vec![ParsedTrigger::OnSuccess],
"on_success" => vec![ParsedTrigger::OnSuccess], "on_failure" => vec![ParsedTrigger::OnFailure],
"on_failure" => vec![ParsedTrigger::OnFailure], s if s.starts_with("on_finish_of:") => {
s if s.starts_with("on_finish_of:") => { let pattern = s.trim_start_matches("on_finish_of:").trim();
let pattern = s.trim_start_matches("on_finish_of:").trim(); let mut builder = GlobSetBuilder::new();
let mut builder = GlobSetBuilder::new(); builder.add(Glob::new(pattern).expect("Invalid glob pattern"));
builder.add(Glob::new(pattern).expect("Invalid glob pattern")); let glob_set = builder.build().expect("Failed to compile glob set");
let glob_set = builder.build().expect("Failed to compile glob set"); vec![ParsedTrigger::OnFinishOf(Matchable::Glob(glob_set))]
vec![ParsedTrigger::OnFinishOf(Matchable::Glob(glob_set))]
}
_ => {
log::warn!("Unknown trigger string: {}", s);
vec![]
}
} }
} _ => {
log::warn!("Unknown trigger string: {}", s);
vec![]
}
},
TriggerItemRaw::Map(map) => { TriggerItemRaw::Map(map) => {
let mut triggers = Vec::new(); let mut triggers = Vec::new();
for (key, patterns) in map { for (key, patterns) in map {
@@ -219,22 +236,21 @@ fn merge_templates(
policy_template: &Option<serde_yaml::Value>, policy_template: &Option<serde_yaml::Value>,
) -> NotificationTemplate { ) -> NotificationTemplate {
let mut merged = method_template.clone(); let mut merged = method_template.clone();
if let Some(policy_tmpl) = policy_template { if let Some(policy_tmpl) = policy_template
if let Ok(policy_template) = && let Ok(policy_template) =
serde_yaml::from_value::<NotificationTemplate>(policy_tmpl.clone()) serde_yaml::from_value::<NotificationTemplate>(policy_tmpl.clone())
{ {
if policy_template.schema.is_some() { if policy_template.schema.is_some() {
merged.schema = policy_template.schema; merged.schema = policy_template.schema;
} }
if policy_template.on_success.is_some() { if policy_template.on_success.is_some() {
merged.on_success = policy_template.on_success; merged.on_success = policy_template.on_success;
} }
if policy_template.on_failure.is_some() { if policy_template.on_failure.is_some() {
merged.on_failure = policy_template.on_failure; merged.on_failure = policy_template.on_failure;
} }
if policy_template.on_finish_of.is_some() { if policy_template.on_finish_of.is_some() {
merged.on_finish_of = policy_template.on_finish_of; merged.on_finish_of = policy_template.on_finish_of;
}
} }
} }
merged merged
@@ -306,7 +322,8 @@ fn parse_rules(config: &NotificationConfig) -> NotificationRulesMap {
.flat_map(parse_trigger_item_raw) .flat_map(parse_trigger_item_raw)
.collect() .collect()
}; };
let template = merge_templates(&method.template, &Some(policy.template.clone())); let template =
merge_templates(&method.template, &Some(policy.template.clone()));
let method_overrides = extract_overrides(&method.config); let method_overrides = extract_overrides(&method.config);
let policy_overrides = extract_overrides(&policy.overrides); let policy_overrides = extract_overrides(&policy.overrides);
let fallible = if policy.overrides.get("fallible").is_some() { let fallible = if policy.overrides.get("fallible").is_some() {
+1 -1
View File
@@ -14,7 +14,7 @@ use std::time::Duration;
use uuid::Uuid; use uuid::Uuid;
use crate::finalize::NotificationTemplate; use crate::finalize::NotificationTemplate;
use crate::finalize::config::Trigger;
/// Final runtime configuration for notification, restored and flattened from yaml /// Final runtime configuration for notification, restored and flattened from yaml
pub struct NotificationRule { pub struct NotificationRule {
pub enabled: bool, pub enabled: bool,
+18 -20
View File
@@ -34,9 +34,9 @@ use crate::{
cli::Cli, cli::Cli,
prebake::{error::PrebakeError, security::get_drop_after, stage::PrebakeStage}, prebake::{error::PrebakeError, security::get_drop_after, stage::PrebakeStage},
}; };
use workshop_engine::EventSender;
use chrono::Utc; use chrono::Utc;
use std::path::Path; use std::path::Path;
use workshop_engine::EventSender;
pub use config::PrebakeConfig; pub use config::PrebakeConfig;
pub use security::prebake_drop_privilege; pub use security::prebake_drop_privilege;
@@ -51,7 +51,6 @@ pub use security::prebake_drop_privilege;
/// ///
/// Returns `Ok(PrebakeConfig)` on successful parsing, or `Err(serde_yaml::Error)` /// Returns `Ok(PrebakeConfig)` on successful parsing, or `Err(serde_yaml::Error)`
/// if the YAML is invalid. /// if the YAML is invalid.
pub fn parse(yaml_content: &str) -> Result<PrebakeConfig, serde_yaml::Error> { pub fn parse(yaml_content: &str) -> Result<PrebakeConfig, serde_yaml::Error> {
serde_yaml::from_str(yaml_content) serde_yaml::from_str(yaml_content)
} }
@@ -92,7 +91,7 @@ pub async fn prebake(
prebake_path: &Path, prebake_path: &Path,
cli: &Cli, cli: &Cli,
stage: Option<PrebakeStage>, stage: Option<PrebakeStage>,
mut ctx: &mut ExecutionContext, ctx: &mut ExecutionContext,
event_tx: EventSender, event_tx: EventSender,
) -> anyhow::Result<()> { ) -> anyhow::Result<()> {
// Initialize // Initialize
@@ -125,26 +124,26 @@ pub async fn prebake(
// Bootstrap stage // Bootstrap stage
if stage <= PrebakeStage::Bootstrap { if stage <= PrebakeStage::Bootstrap {
let stage_start = Utc::now(); let _stage_start = Utc::now();
log::info!("Running PBStage: Bootstrap"); log::info!("Running PBStage: Bootstrap");
prebake_drop_privilege(PrebakeStage::Bootstrap, dropafter, target_user, &mut ctx)?; prebake_drop_privilege(PrebakeStage::Bootstrap, dropafter, target_user, ctx)?;
ctx.task_id = format!("{}-{}-bootstrap", cli.pipeline, cli.build_id); ctx.task_id = format!("{}-{}-bootstrap", cli.pipeline, cli.build_id);
let osinfo = os_info::get(); let osinfo = os_info::get();
let result = stage::bootstrap::bootstrap(&prebake, &osinfo, &ctx, &event_tx).await; let result = stage::bootstrap::bootstrap(&prebake, &osinfo, ctx, &event_tx).await;
result?; result?;
} }
// EarlyHook stage // EarlyHook stage
if stage <= PrebakeStage::EarlyHook { if stage <= PrebakeStage::EarlyHook {
let stage_start = Utc::now(); let _stage_start = Utc::now();
log::info!("Running PBStage: EarlyHook"); log::info!("Running PBStage: EarlyHook");
prebake_drop_privilege(PrebakeStage::EarlyHook, dropafter, target_user, &mut ctx)?; prebake_drop_privilege(PrebakeStage::EarlyHook, dropafter, target_user, ctx)?;
ctx.task_id = format!("{}-{}-earlyhook", cli.pipeline, cli.build_id); ctx.task_id = format!("{}-{}-earlyhook", cli.pipeline, cli.build_id);
let earlyhook = match &prebake.hooks { let earlyhook = match &prebake.hooks {
Some(hooks) => hooks.early.clone(), Some(hooks) => hooks.early.clone(),
None => None, None => None,
}; };
let result = stage::hook::hook(earlyhook, &ctx, event_tx.clone()).await; let result = stage::hook::hook(earlyhook, ctx, event_tx.clone()).await;
result?; result?;
} }
@@ -160,7 +159,7 @@ pub async fn prebake(
let upm_sys_deps = if let Some(user) = &dependencies.user { let upm_sys_deps = if let Some(user) = &dependencies.user {
if stage <= PrebakeStage::DepsSystem { if stage <= PrebakeStage::DepsSystem {
stage::depsuser::collect_upm_sysdeps(user, &ctx, &endpoint).await? stage::depsuser::collect_upm_sysdeps(user, ctx, &endpoint).await?
} else { } else {
stage::depsuser::CollectedUPMdeps::default() stage::depsuser::CollectedUPMdeps::default()
} }
@@ -208,38 +207,37 @@ pub async fn prebake(
} }
if !system_config.is_empty() && stage <= PrebakeStage::DepsSystem { if !system_config.is_empty() && stage <= PrebakeStage::DepsSystem {
let stage_start = Utc::now(); let _stage_start = Utc::now();
log::info!("Running PBStage: DepsSystem"); log::info!("Running PBStage: DepsSystem");
prebake_drop_privilege(PrebakeStage::DepsSystem, dropafter, target_user, &mut ctx)?; prebake_drop_privilege(PrebakeStage::DepsSystem, dropafter, target_user, ctx)?;
ctx.task_id = format!("{}-{}-depssystem", cli.pipeline, cli.build_id); ctx.task_id = format!("{}-{}-depssystem", cli.pipeline, cli.build_id);
let result = let result = stage::depssystem::depssystem(&system_config, ctx, event_tx.clone()).await;
stage::depssystem::depssystem(&system_config, &ctx, event_tx.clone()).await;
result?; result?;
} }
if let Some(user) = dependencies.user if let Some(user) = dependencies.user
&& stage <= PrebakeStage::DepsUser && stage <= PrebakeStage::DepsUser
{ {
let stage_start = Utc::now(); let _stage_start = Utc::now();
log::info!("Running PBStage: DepsUser"); log::info!("Running PBStage: DepsUser");
prebake_drop_privilege(PrebakeStage::DepsUser, dropafter, target_user, &mut ctx)?; prebake_drop_privilege(PrebakeStage::DepsUser, dropafter, target_user, ctx)?;
ctx.task_id = format!("{}-{}-depsuser", cli.pipeline, cli.build_id); ctx.task_id = format!("{}-{}-depsuser", cli.pipeline, cli.build_id);
let result = stage::depsuser::depsuser(&user, &ctx, event_tx.clone(), &endpoint).await; let result = stage::depsuser::depsuser(&user, ctx, event_tx.clone(), &endpoint).await;
result?; result?;
} }
} }
// LateHook stage // LateHook stage
if stage <= PrebakeStage::LateHook { if stage <= PrebakeStage::LateHook {
let stage_start = Utc::now(); let _stage_start = Utc::now();
log::info!("Running PBStage: LateHook"); log::info!("Running PBStage: LateHook");
prebake_drop_privilege(PrebakeStage::LateHook, dropafter, target_user, &mut ctx)?; prebake_drop_privilege(PrebakeStage::LateHook, dropafter, target_user, ctx)?;
ctx.task_id = format!("{}-{}-latehook", cli.pipeline, cli.build_id); ctx.task_id = format!("{}-{}-latehook", cli.pipeline, cli.build_id);
let latehook = match &prebake.hooks { let latehook = match &prebake.hooks {
Some(hooks) => hooks.late.clone(), Some(hooks) => hooks.late.clone(),
None => None, None => None,
}; };
let result = stage::hook::hook(latehook, &ctx, event_tx.clone()).await; let result = stage::hook::hook(latehook, ctx, event_tx.clone()).await;
result?; result?;
} }
+1 -4
View File
@@ -3,10 +3,7 @@ use thiserror::Error;
use workshop_engine::{ use workshop_engine::{
ExecutionError, ExecutionError,
error::{ error::{EXITCODE_IO_ERROR, EXITCODE_PARSE_ERROR, EXITCODE_PRIV_DROP_FAILED, HasExitCode},
EXITCODE_IO_ERROR, EXITCODE_PARSE_ERROR, EXITCODE_PRIV_DROP_FAILED,
HasExitCode,
},
}; };
#[derive(Error, Debug)] #[derive(Error, Debug)]
+1 -1
View File
@@ -1,6 +1,6 @@
use crate::ExecutionContext; use crate::ExecutionContext;
use crate::prebake::error::PrebakeError;
use crate::prebake::config::Security; use crate::prebake::config::Security;
use crate::prebake::error::PrebakeError;
use crate::prebake::stage::PrebakeStage; use crate::prebake::stage::PrebakeStage;
use privdrop::PrivDrop; use privdrop::PrivDrop;
use std::ffi::CString; use std::ffi::CString;
@@ -9,15 +9,15 @@
//! and the target operating system, supporting multiple package managers. //! and the target operating system, supporting multiple package managers.
use crate::ExecutionContext; use crate::ExecutionContext;
use crate::prebake::constant::BOOTSTRAP_SCRIPT_PATH;
use workshop_engine::{Engine, EventSender, pm};
use crate::prebake::error::BootstrapError;
use crate::prebake::config::Bootstrap; use crate::prebake::config::Bootstrap;
use crate::prebake::config::PrebakeConfig; use crate::prebake::config::PrebakeConfig;
use crate::prebake::constant::BOOTSTRAP_SCRIPT_PATH;
use crate::prebake::error::BootstrapError;
use std::fs::File; use std::fs::File;
use std::io::Write; use std::io::Write;
use std::path::{Path, PathBuf}; use std::path::{Path, PathBuf};
use which::which; use which::which;
use workshop_engine::{Engine, EventSender, pm};
/// Executes the bootstrap stage of the prebake workflow. /// Executes the bootstrap stage of the prebake workflow.
/// ///
@@ -96,7 +96,7 @@ fn write_bootstrap_content(bootstrap_path: &Path, content: String) -> Result<(),
let mut file = File::create(bootstrap_path)?; let mut file = File::create(bootstrap_path)?;
file.write_all(content.as_bytes())?; file.write_all(content.as_bytes())?;
std::fs::set_permissions( std::fs::set_permissions(
&bootstrap_path, bootstrap_path,
std::os::unix::fs::PermissionsExt::from_mode(0o755), std::os::unix::fs::PermissionsExt::from_mode(0o755),
)?; )?;
Ok(()) Ok(())
@@ -1,7 +1,7 @@
use crate::ExecutionContext; use crate::ExecutionContext;
use workshop_engine::{Engine, EventSender, pm, error::DependencyError};
use crate::prebake::config::SystemDependency; use crate::prebake::config::SystemDependency;
use std::collections::HashMap; use std::collections::HashMap;
use workshop_engine::{Engine, EventSender, error::DependencyError, pm};
/// Handles system dependency installation for the target platform. /// Handles system dependency installation for the target platform.
/// ///
+1 -1
View File
@@ -1,8 +1,8 @@
use serde_yaml::Value; use serde_yaml::Value;
use crate::ExecutionContext; use crate::ExecutionContext;
use workshop_engine::{EventSender, upm, error::DependencyError, RepologyEndpoint};
use std::collections::HashMap; use std::collections::HashMap;
use workshop_engine::{EventSender, RepologyEndpoint, error::DependencyError, upm};
#[derive(Default)] #[derive(Default)]
pub struct CollectedUPMdeps { pub struct CollectedUPMdeps {
+2 -2
View File
@@ -13,8 +13,8 @@
//! and timeout configurations. //! and timeout configurations.
use crate::ExecutionContext; use crate::ExecutionContext;
use workshop_engine::{DeltaExecutionContext, Engine, EventSender, ExecutionError, CustomCommand};
use std::path::PathBuf; use std::path::PathBuf;
use workshop_engine::{CustomCommand, DeltaExecutionContext, Engine, EventSender, ExecutionError};
/// Executes a sequence of hook commands. /// Executes a sequence of hook commands.
/// ///
@@ -24,7 +24,7 @@ use std::path::PathBuf;
/// # Arguments /// # Arguments
/// ///
/// * `hook` - Optional vector of hook commands to execute. If `None` or empty, /// * `hook` - Optional vector of hook commands to execute. If `None` or empty,
/// the function returns immediately with `Ok(())`. /// the function returns immediately with `Ok(())`.
/// * `ctx` - Base execution context containing pipeline parameters and settings /// * `ctx` - Base execution context containing pipeline parameters and settings
/// * `event_tx` - Event sender for reporting hook execution progress /// * `event_tx` - Event sender for reporting hook execution progress
/// ///
@@ -8,7 +8,7 @@ use std::time::Duration;
use workshop_baker::ExecutionContext; use workshop_baker::ExecutionContext;
use workshop_baker::finalize::plugin::PluginMap; use workshop_baker::finalize::plugin::PluginMap;
use workshop_baker::finalize::stage::hook::hook; use workshop_baker::finalize::stage::hook::hook;
use workshop_engine::{ExecutionEvent, CustomCommand}; use workshop_engine::{CustomCommand, ExecutionEvent};
fn create_test_context() -> ExecutionContext { fn create_test_context() -> ExecutionContext {
ExecutionContext { ExecutionContext {
@@ -5,10 +5,10 @@
use std::collections::HashMap; use std::collections::HashMap;
use std::time::Duration; use std::time::Duration;
use workshop_baker::cli::Cli;
use workshop_baker::ExecutionContext; use workshop_baker::ExecutionContext;
use workshop_engine::ExecutionEvent; use workshop_baker::cli::Cli;
use workshop_baker::finalize::finalize; use workshop_baker::finalize::finalize;
use workshop_engine::ExecutionEvent;
fn create_test_context() -> ExecutionContext { fn create_test_context() -> ExecutionContext {
ExecutionContext { ExecutionContext {
@@ -1,7 +1,7 @@
use workshop_baker::finalize::types::compression::CompressionMethod;
use workshop_baker::prebake::types::architecture::Architecture; use workshop_baker::prebake::types::architecture::Architecture;
use workshop_baker::prebake::types::builderconfig::BuilderConfig; use workshop_baker::prebake::types::builderconfig::BuilderConfig;
use workshop_baker::prebake::types::cache::{CacheDirectory, CacheMode, CacheStrategyType}; use workshop_baker::prebake::types::cache::{CacheDirectory, CacheMode, CacheStrategyType};
use workshop_baker::finalize::types::compression::CompressionMethod;
use workshop_engine::{CustomCommand, RepologyEndpoint}; use workshop_engine::{CustomCommand, RepologyEndpoint};
#[test] #[test]
+8 -2
View File
@@ -41,11 +41,17 @@ impl ManagedChild for TokioChild {
} }
fn take_stdout(&mut self) -> Option<Box<dyn AsyncRead + Unpin + Send>> { fn take_stdout(&mut self) -> Option<Box<dyn AsyncRead + Unpin + Send>> {
self.0.stdout.take().map(|s| Box::new(s) as Box<dyn AsyncRead + Unpin + Send>) self.0
.stdout
.take()
.map(|s| Box::new(s) as Box<dyn AsyncRead + Unpin + Send>)
} }
fn take_stderr(&mut self) -> Option<Box<dyn AsyncRead + Unpin + Send>> { fn take_stderr(&mut self) -> Option<Box<dyn AsyncRead + Unpin + Send>> {
self.0.stderr.take().map(|s| Box::new(s) as Box<dyn AsyncRead + Unpin + Send>) self.0
.stderr
.take()
.map(|s| Box::new(s) as Box<dyn AsyncRead + Unpin + Send>)
} }
async fn wait(&mut self) -> io::Result<i32> { async fn wait(&mut self) -> io::Result<i32> {
+48 -34
View File
@@ -9,8 +9,8 @@ use tokio::time::timeout;
use crate::child::{ManagedChild, TokioChild}; use crate::child::{ManagedChild, TokioChild};
use crate::pm::PackageManager; use crate::pm::PackageManager;
use crate::types::{ use crate::types::{
DeltaExecutionContext, EventSender, ExecutionContext, ExecutionError, DeltaExecutionContext, EventSender, ExecutionContext, ExecutionError, ExecutionEvent,
ExecutionEvent, ExecutionResult, StreamType, ExecutionResult, StreamType,
}; };
impl crate::types::Engine { impl crate::types::Engine {
@@ -92,36 +92,36 @@ impl crate::types::Engine {
.stderr(Stdio::piped()); .stderr(Stdio::piped());
} }
/// Execute a pre-configured command, managing the full lifecycle. /// Execute a pre-configured command, managing the full lifecycle.
/// ///
/// This method: /// This method:
/// 1. Configures the command with context (env vars, working dir, stdio) /// 1. Configures the command with context (env vars, working dir, stdio)
/// 2. Creates a new process group on Unix for process-tree isolation /// 2. Creates a new process group on Unix for process-tree isolation
/// 3. Returns early in dry_run mode /// 3. Returns early in dry_run mode
/// 4. Spawns the OS process /// 4. Spawns the OS process
/// 5. Delegates to `execute_child` for wait/event/output management /// 5. Delegates to `execute_child` for wait/event/output management
pub async fn execute( pub async fn execute(
&self, &self,
command: Command, command: Command,
ctx: &ExecutionContext, ctx: &ExecutionContext,
event_tx: &EventSender, event_tx: &EventSender,
) -> Result<ExecutionResult, ExecutionError> { ) -> Result<ExecutionResult, ExecutionError> {
let mut command = command; let mut command = command;
self.construct_command(&mut command, ctx); self.construct_command(&mut command, ctx);
if ctx.dry_run { if ctx.dry_run {
log::info!("[DRY_RUN] {:?}", command); log::info!("[DRY_RUN] {:?}", command);
return Ok(ExecutionResult { return Ok(ExecutionResult {
..Default::default() ..Default::default()
}); });
} }
// Create a new process group so that kill() can terminate the // Create a new process group so that kill() can terminate the
// entire process tree (bash + its children) on timeout. // entire process tree (bash + its children) on timeout.
#[cfg(unix)] #[cfg(unix)]
command.process_group(0); command.process_group(0);
let child = command.spawn().map_err(|e| { let child = command.spawn().map_err(|e| {
ExecutionError::InvalidCommand(format!("Failed to spawn command: {}", e)) ExecutionError::InvalidCommand(format!("Failed to spawn command: {}", e))
})?; })?;
@@ -352,8 +352,8 @@ impl Default for crate::types::Engine {
#[cfg(test)] #[cfg(test)]
mod tests { mod tests {
use super::*; use super::*;
use crate::child::MockChild;
use crate::Engine; use crate::Engine;
use crate::child::MockChild;
use std::path::PathBuf; use std::path::PathBuf;
use std::time::Duration; use std::time::Duration;
@@ -423,7 +423,10 @@ mod tests {
let result = engine.execute_child(child, &test_ctx(), &None).await; let result = engine.execute_child(child, &test_ctx(), &None).await;
assert!(result.is_err()); assert!(result.is_err());
assert!(matches!(result.unwrap_err(), ExecutionError::ExecutionFailed(_))); assert!(matches!(
result.unwrap_err(),
ExecutionError::ExecutionFailed(_)
));
} }
#[tokio::test] #[tokio::test]
@@ -511,11 +514,20 @@ mod tests {
let mut found_stdout = false; let mut found_stdout = false;
while let Ok(event) = rx.try_recv() { while let Ok(event) = rx.try_recv() {
if matches!(&event, ExecutionEvent::OutputChunk { stream: StreamType::Stdout, .. }) { if matches!(
&event,
ExecutionEvent::OutputChunk {
stream: StreamType::Stdout,
..
}
) {
found_stdout = true; found_stdout = true;
} }
} }
assert!(found_stdout, "Should have received at least one OutputChunk for stdout"); assert!(
found_stdout,
"Should have received at least one OutputChunk for stdout"
);
} }
// ------------------------------------------------------------------ // ------------------------------------------------------------------
@@ -538,7 +550,9 @@ mod tests {
// Should have received TaskFailed for timeout // Should have received TaskFailed for timeout
let _ = rx.recv().await; // skip TaskStarted let _ = rx.recv().await; // skip TaskStarted
let event = rx.recv().await; let event = rx.recv().await;
assert!(matches!(event, Some(ExecutionEvent::TaskFailed { error, .. }) if error == "Timeout")); assert!(
matches!(event, Some(ExecutionEvent::TaskFailed { error, .. }) if error == "Timeout")
);
} }
// ------------------------------------------------------------------ // ------------------------------------------------------------------
+11 -8
View File
@@ -12,14 +12,17 @@ pub mod upm;
pub use child::{ManagedChild, TokioChild}; pub use child::{ManagedChild, TokioChild};
// Convenience re-exports // Convenience re-exports
pub use error::{ExecutionError, HasExitCode,
EXITCODE_OK, EXITCODE_GENERAL_ERROR, EXITCODE_INVALID_ARGS,
EXITCODE_IO_ERROR, EXITCODE_PARSE_ERROR, EXITCODE_EXECUTION_ERROR,
EXITCODE_TIMEOUT, EXITCODE_PRIV_DROP_FAILED};
pub use types::{Engine, EventSender, ExecutionContext, ExecutionEvent, ExecutionResult,
StreamType, DeltaExecutionContext};
pub use command::CustomCommand; pub use command::CustomCommand;
pub use repology::RepologyEndpoint; pub use error::{
pub use time::{parse_duration_to_ms, deserialize_duration_ms, deserialize_option_duration_ms}; EXITCODE_EXECUTION_ERROR, EXITCODE_GENERAL_ERROR, EXITCODE_INVALID_ARGS, EXITCODE_IO_ERROR,
EXITCODE_OK, EXITCODE_PARSE_ERROR, EXITCODE_PRIV_DROP_FAILED, EXITCODE_TIMEOUT, ExecutionError,
HasExitCode,
};
pub use pm::PackageManager; pub use pm::PackageManager;
pub use repology::RepologyEndpoint;
pub use time::{deserialize_duration_ms, deserialize_option_duration_ms, parse_duration_to_ms};
pub use types::{
DeltaExecutionContext, Engine, EventSender, ExecutionContext, ExecutionEvent, ExecutionResult,
StreamType,
};
pub use upm::UserPackageManager; pub use upm::UserPackageManager;
+1 -1
View File
@@ -3,8 +3,8 @@ use async_trait::async_trait;
use serde_yaml::Value; use serde_yaml::Value;
mod custom; mod custom;
use crate::types::ExecutionContext;
use crate::error::DependencyError; use crate::error::DependencyError;
use crate::types::ExecutionContext;
/// System dependencies required by a user package manager. /// System dependencies required by a user package manager.
pub struct UPMSysDeps { pub struct UPMSysDeps {
+4 -4
View File
@@ -2,11 +2,11 @@ use async_trait::async_trait;
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use std::path::PathBuf; use std::path::PathBuf;
use crate::types::Engine;
use crate::types::DeltaExecutionContext;
use crate::upm::UPMSysDeps;
use crate::error::DependencyError;
use crate::command::CustomCommand; use crate::command::CustomCommand;
use crate::error::DependencyError;
use crate::types::DeltaExecutionContext;
use crate::types::Engine;
use crate::upm::UPMSysDeps;
use super::UserPackageManager; use super::UserPackageManager;
+5 -2
View File
@@ -213,9 +213,9 @@ async fn test_execute_captures_stderr() {
#[tokio::test] #[tokio::test]
async fn test_process_group_kill_kills_subprocesses() { async fn test_process_group_kill_kills_subprocesses() {
use std::os::unix::process::CommandExt; use std::os::unix::process::CommandExt;
use std::process::Stdio;
use tokio::io::AsyncBufReadExt; use tokio::io::AsyncBufReadExt;
use tokio::process::Command; use tokio::process::Command;
use std::process::Stdio;
let mut cmd = Command::new("bash"); let mut cmd = Command::new("bash");
cmd.args(["-c", "sleep 999 & echo $!"]); cmd.args(["-c", "sleep 999 & echo $!"]);
@@ -232,7 +232,10 @@ async fn test_process_group_kill_kills_subprocesses() {
let stdout = child.stdout.take().expect("stdout pipe"); let stdout = child.stdout.take().expect("stdout pipe");
let mut reader = tokio::io::BufReader::new(stdout); let mut reader = tokio::io::BufReader::new(stdout);
let mut pid_line = String::new(); let mut pid_line = String::new();
reader.read_line(&mut pid_line).await.expect("read pid line"); reader
.read_line(&mut pid_line)
.await
.expect("read pid line");
let sleep_pid: u32 = pid_line.trim().parse().expect("valid sleep PID"); let sleep_pid: u32 = pid_line.trim().parse().expect("valid sleep PID");
eprintln!("[test] sleep PID: {}", sleep_pid); eprintln!("[test] sleep PID: {}", sleep_pid);