Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_datime/src/schedule/action.rs

15.7 KiB, 63 runs

created by r1870400018:8488, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! Actions a scheduled task can carry out.
2//!
3//! An action is anything implementing the Action trait; a callback, a system
4//! command, a log line, an HTTP request, or a composite of several.
5//!
6//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
7//! Anthropic Claude
8
9use oxedyne_fe2o3_core::prelude::*;
10use crate::time::CalClock;
11use std::{
12 fmt,
13 sync::Arc,
14};
15
16#[derive(Debug, Clone)]
17pub struct ActionContext {
18 pub scheduled_time: CalClock,
19 pub execution_time: CalClock, // when it actually ran
20 pub task_name: String,
21 pub execution_count: u32, // previous runs
22 pub is_retry: bool,
23 pub retry_count: u32, // zero on the first attempt
24}
25
26#[derive(Debug, Clone)]
27pub enum ActionResult {
28 Success,
29 Warning(String), // completed, with something to say
30 Error(ActionError),
31}
32
33#[derive(Debug, Clone)]
34pub enum ActionError {
35 Timeout,
36 InvalidConfig(String),
37 ExternalFailure(String),
38 Cancelled,
39 ExecutionError(String),
40}
41
42impl fmt::Display for ActionError {
43 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
44 match self {
45 ActionError::Timeout => write!(f, "Action execution timed out"),
46 ActionError::InvalidConfig(msg) => write!(f, "Invalid configuration: {}", msg),
47 ActionError::ExternalFailure(msg) => write!(f, "External failure: {}", msg),
48 ActionError::Cancelled => write!(f, "Action was cancelled"),
49 ActionError::ExecutionError(msg) => write!(f, "Execution error: {}", msg),
50 }
51 }
52}
53
54pub trait Action: Send + Sync + fmt::Debug {
55 fn execute(&self, context: &ActionContext) -> Outcome<ActionResult>;
56
57 fn description(&self) -> &str;
58
59 /// Action is not Clone, so cloning a task goes through this.
60 fn box_clone(&self) -> Box<dyn Action>;
61
62 // validate runs when the task is built, prepare before each execution and
63 // cleanup after it. All three do nothing unless overridden.
64 fn validate(&self) -> Outcome<()> {
65 Ok(())
66 }
67
68 fn prepare(&self, _context: &ActionContext) -> Outcome<()> {
69 Ok(())
70 }
71
72 fn cleanup(&self, _context: &ActionContext, _result: &ActionResult) -> Outcome<()> {
73 Ok(())
74 }
75}
76
77#[derive(Clone)]
78pub struct CallbackAction {
79 callback: Arc<dyn Fn() -> Outcome<()> + Send + Sync>,
80 description: String,
81}
82
83impl fmt::Debug for CallbackAction {
84 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
85 f.debug_struct("CallbackAction")
86 .field("description", &self.description)
87 .field("callback", &"<function>")
88 .finish()
89 }
90}
91
92impl CallbackAction {
93 pub fn new<F>(callback: F) -> Self
94 where
95 F: Fn() -> Outcome<()> + Send + Sync + 'static
96 {
97 CallbackAction {
98 callback: Arc::new(callback),
99 description: "Callback action".to_string(),
100 }
101 }
102
103 pub fn with_description<F, S>(callback: F, description: S) -> Self
104 where
105 F: Fn() -> Outcome<()> + Send + Sync + 'static,
106 S: Into<String>
107 {
108 CallbackAction {
109 callback: Arc::new(callback),
110 description: description.into(),
111 }
112 }
113}
114
115impl Action for CallbackAction {
116 fn execute(&self, _context: &ActionContext) -> Outcome<ActionResult> {
117 match (self.callback)() {
118 Ok(()) => Ok(ActionResult::Success),
119 Err(e) => Ok(ActionResult::Error(ActionError::ExecutionError(e.to_string()))),
120 }
121 }
122
123 fn description(&self) -> &str {
124 &self.description
125 }
126
127 fn box_clone(&self) -> Box<dyn Action> {
128 Box::new(self.clone())
129 }
130}
131
132#[derive(Debug, Clone)]
133pub struct CommandAction {
134 command: String,
135 args: Vec<String>,
136 working_dir: Option<String>,
137 env_vars: Vec<(String, String)>,
138 description: String,
139}
140
141impl CommandAction {
142 pub fn new<S: Into<String>>(command: S) -> Self {
143 let cmd = command.into();
144 CommandAction {
145 description: format!("Command: {}", cmd),
146 command: cmd,
147 args: Vec::new(),
148 working_dir: None,
149 env_vars: Vec::new(),
150 }
151 }
152
153 pub fn args<I, S>(mut self, args: I) -> Self
154 where
155 I: IntoIterator<Item = S>,
156 S: Into<String>
157 {
158 self.args.extend(args.into_iter().map(Into::into));
159 self
160 }
161
162 pub fn working_dir<S: Into<String>>(mut self, dir: S) -> Self {
163 self.working_dir = Some(dir.into());
164 self
165 }
166
167 pub fn env<K, V>(mut self, key: K, value: V) -> Self
168 where
169 K: Into<String>,
170 V: Into<String>
171 {
172 self.env_vars.push((key.into(), value.into()));
173 self
174 }
175}
176
177impl Action for CommandAction {
178 fn execute(&self, _context: &ActionContext) -> Outcome<ActionResult> {
179 use std::process::Command;
180
181 let mut cmd = Command::new(&self.command);
182 cmd.args(&self.args);
183
184 if let Some(ref dir) = self.working_dir {
185 cmd.current_dir(dir);
186 }
187
188 for (key, value) in &self.env_vars {
189 cmd.env(key, value);
190 }
191
192 match cmd.output() {
193 Ok(output) => {
194 if output.status.success() {
195 Ok(ActionResult::Success)
196 } else {
197 let stderr = String::from_utf8_lossy(&output.stderr);
198 Ok(ActionResult::Error(ActionError::ExecutionError(
199 format!("Command failed with exit code {:?}: {}", output.status.code(), stderr)
200 )))
201 }
202 },
203 Err(e) => Ok(ActionResult::Error(ActionError::ExecutionError(
204 format!("Failed to execute command: {}", e)
205 ))),
206 }
207 }
208
209 fn description(&self) -> &str {
210 &self.description
211 }
212
213 fn validate(&self) -> Outcome<()> {
214 // Basic validation - check if command exists
215 if self.command.is_empty() {
216 return Err(err!("Command cannot be empty"; Invalid, Input));
217 }
218 Ok(())
219 }
220
221 fn box_clone(&self) -> Box<dyn Action> {
222 Box::new(self.clone())
223 }
224}
225
226#[derive(Debug, Clone)]
227pub struct LogAction {
228 message: String,
229 log_file: Option<String>,
230 description: String,
231}
232
233impl LogAction {
234 /// Writes to stdout.
235 pub fn new<S: Into<String>>(message: S) -> Self {
236 LogAction {
237 message: message.into(),
238 log_file: None,
239 description: "Log message".to_string(),
240 }
241 }
242
243 pub fn to_file<M, F>(message: M, file_path: F) -> Self
244 where
245 M: Into<String>,
246 F: Into<String>
247 {
248 let file = file_path.into();
249 LogAction {
250 message: message.into(),
251 description: format!("Log to file: {}", file),
252 log_file: Some(file),
253 }
254 }
255}
256
257impl Action for LogAction {
258 fn execute(&self, context: &ActionContext) -> Outcome<ActionResult> {
259 let timestamp = context.execution_time.to_string();
260 let log_entry = format!("[{}] Task '{}': {}", timestamp, context.task_name, self.message);
261
262 if let Some(ref file_path) = self.log_file {
263 match std::fs::OpenOptions::new()
264 .create(true)
265 .append(true)
266 .open(file_path)
267 {
268 Ok(mut file) => {
269 use std::io::Write;
270 if let Err(e) = writeln!(file, "{}", log_entry) {
271 return Ok(ActionResult::Error(ActionError::ExecutionError(
272 format!("Failed to write to log file: {}", e)
273 )));
274 }
275 },
276 Err(e) => {
277 return Ok(ActionResult::Error(ActionError::ExecutionError(
278 format!("Failed to open log file: {}", e)
279 )));
280 }
281 }
282 } else {
283 println!("{}", log_entry);
284 }
285
286 Ok(ActionResult::Success)
287 }
288
289 fn description(&self) -> &str {
290 &self.description
291 }
292
293 fn box_clone(&self) -> Box<dyn Action> {
294 Box::new(self.clone())
295 }
296}
297
298#[derive(Debug, Clone)]
299pub struct HttpAction {
300 url: String,
301 #[allow(dead_code)]
302 method: String,
303 headers: Vec<(String, String)>,
304 body: Option<String>,
305 description: String,
306}
307
308impl HttpAction {
309 pub fn get<S: Into<String>>(url: S) -> Self {
310 let url_str = url.into();
311 HttpAction {
312 description: format!("HTTP GET: {}", url_str),
313 url: url_str,
314 method: "GET".to_string(),
315 headers: Vec::new(),
316 body: None,
317 }
318 }
319
320 pub fn post<S: Into<String>>(url: S) -> Self {
321 let url_str = url.into();
322 HttpAction {
323 description: format!("HTTP POST: {}", url_str),
324 url: url_str,
325 method: "POST".to_string(),
326 headers: Vec::new(),
327 body: None,
328 }
329 }
330
331 pub fn header<K, V>(mut self, key: K, value: V) -> Self
332 where
333 K: Into<String>,
334 V: Into<String>
335 {
336 self.headers.push((key.into(), value.into()));
337 self
338 }
339
340 pub fn body<S: Into<String>>(mut self, body: S) -> Self {
341 self.body = Some(body.into());
342 self
343 }
344}
345
346impl Action for HttpAction {
347 fn execute(&self, _context: &ActionContext) -> Outcome<ActionResult> {
348 // Note: In a real implementation, you'd use a proper HTTP client like reqwest
349 // For now, this is a placeholder that simulates HTTP requests
350
351 if self.url.is_empty() {
352 return Ok(ActionResult::Error(ActionError::InvalidConfig(
353 "URL cannot be empty".to_string()
354 )));
355 }
356
357 // Simulate HTTP request
358 if self.url.starts_with("https://") || self.url.starts_with("http://") {
359 Ok(ActionResult::Success)
360 } else {
361 Ok(ActionResult::Error(ActionError::InvalidConfig(
362 "Invalid URL format".to_string()
363 )))
364 }
365 }
366
367 fn description(&self) -> &str {
368 &self.description
369 }
370
371 fn validate(&self) -> Outcome<()> {
372 if self.url.is_empty() {
373 return Err(err!("URL cannot be empty"; Invalid, Input));
374 }
375
376 if !self.url.starts_with("http://") && !self.url.starts_with("https://") {
377 return Err(err!("URL must start with http:// or https://"; Invalid, Input));
378 }
379
380 Ok(())
381 }
382
383 fn box_clone(&self) -> Box<dyn Action> {
384 Box::new(self.clone())
385 }
386}
387
388/// Runs its actions in the order they were added.
389#[derive(Debug)]
390pub struct CompositeAction {
391 actions: Vec<Box<dyn Action>>,
392 description: String,
393 stop_on_failure: bool,
394}
395
396impl Clone for CompositeAction {
397 fn clone(&self) -> Self {
398 CompositeAction {
399 actions: self.actions.iter().map(|a| a.box_clone()).collect(),
400 description: self.description.clone(),
401 stop_on_failure: self.stop_on_failure,
402 }
403 }
404}
405
406impl CompositeAction {
407 pub fn new<S: Into<String>>(description: S) -> Self {
408 CompositeAction {
409 actions: Vec::new(),
410 description: description.into(),
411 stop_on_failure: true,
412 }
413 }
414
415 pub fn add_action<A: Action + 'static>(mut self, action: A) -> Self {
416 self.actions.push(Box::new(action));
417 self
418 }
419
420 pub fn stop_on_failure(mut self, stop: bool) -> Self {
421 self.stop_on_failure = stop;
422 self
423 }
424}
425
426impl Action for CompositeAction {
427 fn execute(&self, context: &ActionContext) -> Outcome<ActionResult> {
428 let mut warnings = Vec::new();
429
430 for (i, action) in self.actions.iter().enumerate() {
431 match res!(action.execute(context)) {
432 ActionResult::Success => continue,
433 ActionResult::Warning(msg) => {
434 warnings.push(format!("Action {}: {}", i + 1, msg));
435 continue;
436 },
437 ActionResult::Error(err) => {
438 if self.stop_on_failure {
439 return Ok(ActionResult::Error(err));
440 } else {
441 warnings.push(format!("Action {} failed: {}", i + 1, err));
442 continue;
443 }
444 }
445 }
446 }
447
448 if warnings.is_empty() {
449 Ok(ActionResult::Success)
450 } else {
451 Ok(ActionResult::Warning(warnings.join("; ")))
452 }
453 }
454
455 fn description(&self) -> &str {
456 &self.description
457 }
458
459 fn validate(&self) -> Outcome<()> {
460 for action in &self.actions {
461 res!(action.validate());
462 }
463 Ok(())
464 }
465
466 fn box_clone(&self) -> Box<dyn Action> {
467 Box::new(self.clone())
468 }
469}
470
471#[cfg(test)]
472mod tests {
473 use super::*;
474 use crate::time::CalClockZone;
475
476 #[test]
477 fn test_callback_action() {
478 let action = CallbackAction::new(|| {
479 println!("Test callback executed");
480 Ok(())
481 });
482
483 let context = ActionContext {
484 scheduled_time: CalClock::new(2024, 1, 1, 12, 0, 0, 0, CalClockZone::utc()).unwrap(),
485 execution_time: CalClock::new(2024, 1, 1, 12, 0, 1, 0, CalClockZone::utc()).unwrap(),
486 task_name: "test".to_string(),
487 execution_count: 1,
488 is_retry: false,
489 retry_count: 0,
490 };
491
492 let result = action.execute(&context).unwrap();
493 assert!(matches!(result, ActionResult::Success));
494 }
495
496 #[test]
497 fn test_log_action() {
498 let action = LogAction::new("Test log message");
499
500 let context = ActionContext {
501 scheduled_time: CalClock::new(2024, 1, 1, 12, 0, 0, 0, CalClockZone::utc()).unwrap(),
502 execution_time: CalClock::new(2024, 1, 1, 12, 0, 1, 0, CalClockZone::utc()).unwrap(),
503 task_name: "test_task".to_string(),
504 execution_count: 1,
505 is_retry: false,
506 retry_count: 0,
507 };
508
509 let result = action.execute(&context).unwrap();
510 assert!(matches!(result, ActionResult::Success));
511 }
512
513 #[test]
514 fn test_command_action_validation() {
515 let action = CommandAction::new("ls").args(["-la"]);
516 assert!(action.validate().is_ok());
517
518 let empty_action = CommandAction::new("");
519 assert!(empty_action.validate().is_err());
520 }
521
522 #[test]
523 fn test_http_action_validation() {
524 let valid_action = HttpAction::get("https://example.com");
525 assert!(valid_action.validate().is_ok());
526
527 let invalid_action = HttpAction::get("invalid-url");
528 assert!(invalid_action.validate().is_err());
529
530 let empty_action = HttpAction::get("");
531 assert!(empty_action.validate().is_err());
532 }
533
534 #[test]
535 fn test_composite_action() {
536 let composite = CompositeAction::new("Test composite")
537 .add_action(LogAction::new("First action"))
538 .add_action(LogAction::new("Second action"))
539 .stop_on_failure(false);
540
541 assert!(composite.validate().is_ok());
542
543 let context = ActionContext {
544 scheduled_time: CalClock::new(2024, 1, 1, 12, 0, 0, 0, CalClockZone::utc()).unwrap(),
545 execution_time: CalClock::new(2024, 1, 1, 12, 0, 1, 0, CalClockZone::utc()).unwrap(),
546 task_name: "composite_test".to_string(),
547 execution_count: 1,
548 is_retry: false,
549 retry_count: 0,
550 };
551
552 let result = composite.execute(&context).unwrap();
553 assert!(matches!(result, ActionResult::Success));
554 }
555}