#![allow(clippy::panic)] use crate::{
triggers::mutation::{
AfterMutationTrigger, BeforeMutationTrigger, EntityEvent, EventKind, TriggerMatcher,
},
types::EventPayload,
};
#[test]
fn test_after_mutation_fires_on_insert() {
let trigger = AfterMutationTrigger {
function_name: "onUserCreated".to_string(),
entity_type: "User".to_string(),
event_filter: Some(EventKind::Insert),
predicates: Vec::new(),
};
let event = EntityEvent {
entity: "User".to_string(),
event_kind: EventKind::Insert,
old: None,
new: Some(serde_json::json!({ "id": 1, "name": "Alice" })),
timestamp: chrono::Utc::now(),
};
let payload = trigger.build_payload(&event);
assert_eq!(payload.trigger_type, "after:mutation:onUserCreated");
assert_eq!(payload.entity, "User");
assert_eq!(payload.event_kind, "insert");
assert_eq!(payload.data["old"], serde_json::Value::Null);
assert!(payload.data["new"].is_object());
}
#[test]
fn test_after_mutation_fires_on_update() {
let trigger = AfterMutationTrigger {
function_name: "onUserUpdated".to_string(),
entity_type: "User".to_string(),
event_filter: Some(EventKind::Update),
predicates: Vec::new(),
};
let event = EntityEvent {
entity: "User".to_string(),
event_kind: EventKind::Update,
old: Some(serde_json::json!({ "id": 1, "name": "Alice" })),
new: Some(serde_json::json!({ "id": 1, "name": "Alice Smith" })),
timestamp: chrono::Utc::now(),
};
let payload = trigger.build_payload(&event);
assert_eq!(payload.trigger_type, "after:mutation:onUserUpdated");
assert_eq!(payload.entity, "User");
assert_eq!(payload.event_kind, "update");
assert!(payload.data["old"].is_object());
assert!(payload.data["new"].is_object());
}
#[test]
fn test_after_mutation_fires_on_delete() {
let trigger = AfterMutationTrigger {
function_name: "onUserDeleted".to_string(),
entity_type: "User".to_string(),
event_filter: Some(EventKind::Delete),
predicates: Vec::new(),
};
let event = EntityEvent {
entity: "User".to_string(),
event_kind: EventKind::Delete,
old: Some(serde_json::json!({ "id": 1, "name": "Alice" })),
new: None,
timestamp: chrono::Utc::now(),
};
let payload = trigger.build_payload(&event);
assert_eq!(payload.trigger_type, "after:mutation:onUserDeleted");
assert_eq!(payload.entity, "User");
assert_eq!(payload.event_kind, "delete");
assert!(payload.data["old"].is_object());
assert_eq!(payload.data["new"], serde_json::Value::Null);
}
#[test]
fn test_after_mutation_receives_entity_type() {
let trigger = AfterMutationTrigger {
function_name: "onPostCreated".to_string(),
entity_type: "Post".to_string(),
event_filter: Some(EventKind::Insert),
predicates: Vec::new(),
};
let event = EntityEvent {
entity: "Post".to_string(),
event_kind: EventKind::Insert,
old: None,
new: Some(serde_json::json!({ "id": 1, "title": "Hello" })),
timestamp: chrono::Utc::now(),
};
let payload = trigger.build_payload(&event);
assert_eq!(payload.entity, "Post");
assert_eq!(trigger.entity_type, "Post");
}
#[test]
fn test_after_mutation_trigger_matching() {
let trigger_insert = AfterMutationTrigger {
function_name: "onUserCreated".to_string(),
entity_type: "User".to_string(),
event_filter: Some(EventKind::Insert),
predicates: Vec::new(),
};
let trigger_all = AfterMutationTrigger {
function_name: "onUserChanged".to_string(),
entity_type: "User".to_string(),
event_filter: None,
predicates: Vec::new(),
};
assert!(trigger_insert.matches("User", EventKind::Insert));
assert!(!trigger_insert.matches("User", EventKind::Update));
assert!(!trigger_insert.matches("Post", EventKind::Insert));
assert!(trigger_all.matches("User", EventKind::Insert));
assert!(trigger_all.matches("User", EventKind::Update));
assert!(trigger_all.matches("User", EventKind::Delete));
assert!(!trigger_all.matches("Post", EventKind::Insert));
}
#[test]
fn test_before_mutation_trigger_matching() {
let trigger = BeforeMutationTrigger {
function_name: "validateUserInput".to_string(),
mutation_name: "createUser".to_string(),
};
assert!(trigger.matches("createUser"));
assert!(!trigger.matches("updateUser"));
assert!(!trigger.matches("deleteUser"));
}
#[test]
fn test_before_mutation_multiple_triggers() {
let trigger_a = BeforeMutationTrigger {
function_name: "validateInput".to_string(),
mutation_name: "createUser".to_string(),
};
let trigger_b = BeforeMutationTrigger {
function_name: "checkDuplicates".to_string(),
mutation_name: "createUser".to_string(),
};
let trigger_c = BeforeMutationTrigger {
function_name: "auditLog".to_string(),
mutation_name: "createUser".to_string(),
};
assert!(trigger_a.matches("createUser"));
assert!(trigger_b.matches("createUser"));
assert!(trigger_c.matches("createUser"));
}
#[test]
fn test_after_mutation_payload_serialization() {
let trigger = AfterMutationTrigger {
function_name: "onUserCreated".to_string(),
entity_type: "User".to_string(),
event_filter: Some(EventKind::Insert),
predicates: Vec::new(),
};
let event = EntityEvent {
entity: "User".to_string(),
event_kind: EventKind::Insert,
old: None,
new: Some(serde_json::json!({ "id": 1, "name": "Alice" })),
timestamp: chrono::Utc::now(),
};
let payload = trigger.build_payload(&event);
let json = serde_json::to_string(&payload).expect("serialize");
let restored: EventPayload = serde_json::from_str(&json).expect("deserialize");
assert_eq!(restored.trigger_type, payload.trigger_type);
assert_eq!(restored.entity, payload.entity);
assert_eq!(restored.event_kind, payload.event_kind);
}
#[test]
fn test_trigger_dispatch_finds_matching_triggers() {
let mut matcher = TriggerMatcher::new();
matcher.add(AfterMutationTrigger {
function_name: "onUserCreated".to_string(),
entity_type: "User".to_string(),
event_filter: Some(EventKind::Insert),
predicates: Vec::new(),
});
matcher.add(AfterMutationTrigger {
function_name: "onUserChanged".to_string(),
entity_type: "User".to_string(),
event_filter: None, predicates: Vec::new(),
});
let triggers = matcher.find("User", EventKind::Insert);
assert_eq!(triggers.len(), 2);
let names: Vec<_> = triggers.iter().map(|t| t.function_name.as_str()).collect();
assert!(names.contains(&"onUserCreated"));
assert!(names.contains(&"onUserChanged"));
let triggers = matcher.find("User", EventKind::Update);
assert_eq!(triggers.len(), 1);
assert_eq!(triggers[0].function_name, "onUserChanged");
}
#[tokio::test]
async fn test_after_mutation_async_dispatch_nonblocking() {
let trigger = AfterMutationTrigger {
function_name: "onUserCreated".to_string(),
entity_type: "User".to_string(),
event_filter: Some(EventKind::Insert),
predicates: Vec::new(),
};
let event = EntityEvent {
entity: "User".to_string(),
event_kind: EventKind::Insert,
old: None,
new: Some(serde_json::json!({ "id": 1, "name": "Alice" })),
timestamp: chrono::Utc::now(),
};
let payload = trigger.build_payload(&event);
assert_eq!(payload.trigger_type, "after:mutation:onUserCreated");
}
#[test]
fn test_trigger_dispatch_multiple_mutations() {
let mut matcher = TriggerMatcher::new();
matcher.add(AfterMutationTrigger {
function_name: "onUserCreated".to_string(),
entity_type: "User".to_string(),
event_filter: Some(EventKind::Insert),
predicates: Vec::new(),
});
matcher.add(AfterMutationTrigger {
function_name: "onUserDeleted".to_string(),
entity_type: "User".to_string(),
event_filter: Some(EventKind::Delete),
predicates: Vec::new(),
});
matcher.add(AfterMutationTrigger {
function_name: "onPostCreated".to_string(),
entity_type: "Post".to_string(),
event_filter: Some(EventKind::Insert),
predicates: Vec::new(),
});
let triggers = matcher.find("User", EventKind::Insert);
assert_eq!(triggers.len(), 1);
assert_eq!(triggers[0].function_name, "onUserCreated");
let triggers = matcher.find("User", EventKind::Delete);
assert_eq!(triggers.len(), 1);
assert_eq!(triggers[0].function_name, "onUserDeleted");
let triggers = matcher.find("Post", EventKind::Insert);
assert_eq!(triggers.len(), 1);
assert_eq!(triggers[0].function_name, "onPostCreated");
let triggers = matcher.find("Post", EventKind::Delete);
assert!(triggers.is_empty());
}
use crate::triggers::mutation::BeforeMutationResult;
#[test]
fn test_before_mutation_receives_proposed_input() {
let input = serde_json::json!({
"name": "Alice",
"email": "alice@example.com"
});
assert!(input.is_object());
assert_eq!(input["name"], "Alice");
assert_eq!(input["email"], "alice@example.com");
}
#[test]
fn test_before_mutation_proceed_allows_mutation() {
let input = serde_json::json!({
"name": "Alice",
"email": "alice@example.com"
});
let result = BeforeMutationResult::Proceed(input);
match result {
BeforeMutationResult::Proceed(modified) => {
assert_eq!(modified["name"], "Alice");
assert_eq!(modified["email"], "alice@example.com");
},
BeforeMutationResult::Abort(_) => {
panic!("Expected Proceed, got Abort");
},
}
}
#[test]
fn test_before_mutation_proceed_with_modified_input() {
let modified = serde_json::json!({
"name": "ALICE",
"email": "alice@example.com"
});
let result = BeforeMutationResult::Proceed(modified);
match result {
BeforeMutationResult::Proceed(output) => {
assert_eq!(output["name"], "ALICE");
assert_ne!(output["name"], "alice");
},
BeforeMutationResult::Abort(_) => {
panic!("Expected Proceed, got Abort");
},
}
}
#[test]
fn test_before_mutation_abort_cancels_mutation() {
let result: BeforeMutationResult =
BeforeMutationResult::Abort("validation failed: name is required".to_string());
match result {
BeforeMutationResult::Proceed(_) => {
panic!("Expected Abort, got Proceed");
},
BeforeMutationResult::Abort(error) => {
assert_eq!(error, "validation failed: name is required");
},
}
}
#[test]
fn test_before_mutation_chain_order() {
let trigger_a = BeforeMutationTrigger {
function_name: "validateInput".to_string(),
mutation_name: "createUser".to_string(),
};
let trigger_b = BeforeMutationTrigger {
function_name: "checkDuplicates".to_string(),
mutation_name: "createUser".to_string(),
};
let trigger_c = BeforeMutationTrigger {
function_name: "auditLog".to_string(),
mutation_name: "createUser".to_string(),
};
let chain = crate::triggers::mutation::BeforeMutationChain {
triggers: vec![trigger_a, trigger_b, trigger_c],
};
assert_eq!(chain.triggers[0].function_name, "validateInput");
assert_eq!(chain.triggers[1].function_name, "checkDuplicates");
assert_eq!(chain.triggers[2].function_name, "auditLog");
}
#[test]
fn test_before_mutation_result_serialization() {
let proceed_result = BeforeMutationResult::Proceed(serde_json::json!({"name": "Alice"}));
let json = serde_json::to_string(&proceed_result).expect("serialize");
let restored: BeforeMutationResult = serde_json::from_str(&json).expect("deserialize");
match restored {
BeforeMutationResult::Proceed(value) => {
assert_eq!(value["name"], "Alice");
},
BeforeMutationResult::Abort(_) => {
panic!("Expected Proceed after deserialization");
},
}
}
#[test]
fn test_before_mutation_abort_serialization() {
let abort_result = BeforeMutationResult::Abort("validation error".to_string());
let json = serde_json::to_string(&abort_result).expect("serialize");
let restored: BeforeMutationResult = serde_json::from_str(&json).expect("deserialize");
match restored {
BeforeMutationResult::Proceed(_) => {
panic!("Expected Abort after deserialization");
},
BeforeMutationResult::Abort(error) => {
assert_eq!(error, "validation error");
},
}
}
#[test]
fn test_before_mutation_chain_execution_simulation() {
use crate::triggers::mutation::BeforeMutationChain;
let chain = BeforeMutationChain {
triggers: vec![
BeforeMutationTrigger {
function_name: "normalizeEmail".to_string(),
mutation_name: "createUser".to_string(),
},
BeforeMutationTrigger {
function_name: "validateName".to_string(),
mutation_name: "createUser".to_string(),
},
BeforeMutationTrigger {
function_name: "enrichProfile".to_string(),
mutation_name: "createUser".to_string(),
},
],
};
assert_eq!(chain.triggers.len(), 3);
let mut current_input = serde_json::json!({
"name": "alice smith",
"email": " ALICE@EXAMPLE.COM "
});
current_input["email"] = serde_json::Value::String("alice@example.com".to_string());
current_input["name"] = serde_json::Value::String("Alice Smith".to_string());
current_input["profile"] = serde_json::json!({"bio": "User"});
assert_eq!(current_input["email"], "alice@example.com");
assert_eq!(current_input["name"], "Alice Smith");
assert!(current_input["profile"].is_object());
assert_eq!(current_input["profile"]["bio"], "User");
}
#[test]
fn test_before_mutation_chain_abort_simulation() {
let chain = crate::triggers::mutation::BeforeMutationChain {
triggers: vec![
BeforeMutationTrigger {
function_name: "validateInput".to_string(),
mutation_name: "createUser".to_string(),
},
BeforeMutationTrigger {
function_name: "checkDuplicates".to_string(),
mutation_name: "createUser".to_string(),
},
BeforeMutationTrigger {
function_name: "auditLog".to_string(),
mutation_name: "createUser".to_string(),
},
],
};
assert_eq!(chain.triggers.len(), 3);
let result1 = BeforeMutationResult::Abort("name is required".to_string());
match result1 {
BeforeMutationResult::Abort(error) => {
assert_eq!(error, "name is required");
},
BeforeMutationResult::Proceed(_) => {
panic!("Expected abort");
},
}
}
use crate::triggers::storage::{StorageEventPayload, StorageOperation, StorageTrigger};
#[test]
fn test_after_storage_upload_fires() {
let trigger = StorageTrigger {
function_name: "onAvatarUpload".to_string(),
bucket: "avatars".to_string(),
operation: StorageOperation::Upload,
};
let storage_event = StorageEventPayload {
bucket: "avatars".to_string(),
key: "users/alice/avatar.jpg".to_string(),
size_bytes: 204_800,
content_type: "image/jpeg".to_string(),
owner_id: Some("user123".to_string()),
operation: StorageOperation::Upload,
};
let payload = trigger.build_payload(&storage_event);
assert_eq!(payload.trigger_type, "after:storage:avatars:upload");
assert_eq!(payload.entity, "avatars");
assert_eq!(payload.event_kind, "upload");
assert_eq!(payload.data["bucket"], "avatars");
assert_eq!(payload.data["key"], "users/alice/avatar.jpg");
assert_eq!(payload.data["size_bytes"], 204_800);
assert_eq!(payload.data["content_type"], "image/jpeg");
assert_eq!(payload.data["owner_id"], "user123");
}
#[test]
fn test_after_storage_delete_fires() {
let trigger = StorageTrigger {
function_name: "onDocumentDelete".to_string(),
bucket: "documents".to_string(),
operation: StorageOperation::Delete,
};
let storage_event = StorageEventPayload {
bucket: "documents".to_string(),
key: "reports/2024/report.pdf".to_string(),
size_bytes: 0,
content_type: "application/pdf".to_string(),
owner_id: Some("user456".to_string()),
operation: StorageOperation::Delete,
};
let payload = trigger.build_payload(&storage_event);
assert_eq!(payload.trigger_type, "after:storage:documents:delete");
assert_eq!(payload.entity, "documents");
assert_eq!(payload.event_kind, "delete");
assert_eq!(payload.data["bucket"], "documents");
assert_eq!(payload.data["key"], "reports/2024/report.pdf");
assert_eq!(payload.data["operation"], "delete");
}
#[test]
fn test_after_storage_matches_bucket() {
let avatar_trigger = StorageTrigger {
function_name: "onAvatarUpload".to_string(),
bucket: "avatars".to_string(),
operation: StorageOperation::Upload,
};
let avatar_event = StorageEventPayload {
bucket: "avatars".to_string(),
key: "user/avatar.jpg".to_string(),
size_bytes: 100_000,
content_type: "image/jpeg".to_string(),
owner_id: Some("user1".to_string()),
operation: StorageOperation::Upload,
};
let doc_event = StorageEventPayload {
bucket: "documents".to_string(),
key: "report.pdf".to_string(),
size_bytes: 500_000,
content_type: "application/pdf".to_string(),
owner_id: Some("user1".to_string()),
operation: StorageOperation::Upload,
};
assert!(avatar_trigger.matches(&avatar_event));
assert!(!avatar_trigger.matches(&doc_event));
}
#[test]
fn test_after_storage_matches_operation() {
let upload_trigger = StorageTrigger {
function_name: "onAvatarUpload".to_string(),
bucket: "avatars".to_string(),
operation: StorageOperation::Upload,
};
let upload_event = StorageEventPayload {
bucket: "avatars".to_string(),
key: "avatar.jpg".to_string(),
size_bytes: 100_000,
content_type: "image/jpeg".to_string(),
owner_id: Some("user1".to_string()),
operation: StorageOperation::Upload,
};
let delete_event = StorageEventPayload {
bucket: "avatars".to_string(),
key: "avatar.jpg".to_string(),
size_bytes: 0,
content_type: "image/jpeg".to_string(),
owner_id: Some("user1".to_string()),
operation: StorageOperation::Delete,
};
assert!(upload_trigger.matches(&upload_event));
assert!(!upload_trigger.matches(&delete_event));
}
#[test]
fn test_after_storage_matches_any_operation() {
let any_trigger = StorageTrigger {
function_name: "onStorageEvent".to_string(),
bucket: "avatars".to_string(),
operation: StorageOperation::Any,
};
let upload_event = StorageEventPayload {
bucket: "avatars".to_string(),
key: "avatar.jpg".to_string(),
size_bytes: 100_000,
content_type: "image/jpeg".to_string(),
owner_id: Some("user1".to_string()),
operation: StorageOperation::Upload,
};
let delete_event = StorageEventPayload {
bucket: "avatars".to_string(),
key: "avatar.jpg".to_string(),
size_bytes: 0,
content_type: "image/jpeg".to_string(),
owner_id: Some("user1".to_string()),
operation: StorageOperation::Delete,
};
assert!(any_trigger.matches(&upload_event));
assert!(any_trigger.matches(&delete_event));
}
#[test]
fn test_after_storage_ignores_transform_cache() {
let trigger = StorageTrigger {
function_name: "onAvatarUpload".to_string(),
bucket: "avatars".to_string(),
operation: StorageOperation::Upload,
};
let transform_event = StorageEventPayload {
bucket: "avatars".to_string(),
key: "_transforms/avatar-thumb.jpg".to_string(),
size_bytes: 50000,
content_type: "image/jpeg".to_string(),
owner_id: None,
operation: StorageOperation::Upload,
};
assert!(!trigger.should_fire(&transform_event));
}
#[test]
fn test_after_storage_payload_includes_metadata() {
let trigger = StorageTrigger {
function_name: "onUpload".to_string(),
bucket: "documents".to_string(),
operation: StorageOperation::Upload,
};
let storage_event = StorageEventPayload {
bucket: "documents".to_string(),
key: "invoices/INV-001.pdf".to_string(),
size_bytes: 1_024_000,
content_type: "application/pdf".to_string(),
owner_id: Some("company_789".to_string()),
operation: StorageOperation::Upload,
};
let payload = trigger.build_payload(&storage_event);
assert!(payload.data.is_object());
assert!(payload.data["bucket"].is_string());
assert!(payload.data["key"].is_string());
assert!(payload.data["size_bytes"].is_number());
assert!(payload.data["content_type"].is_string());
assert!(payload.data["owner_id"].is_string());
}
use crate::triggers::cron::{CronExecutionState, CronSchedule, CronTrigger};
#[test]
fn test_cron_trigger_parses_daily_expression() {
let trigger = CronTrigger {
function_name: "dailyCleanup".to_string(),
schedule: "0 2 * * *".to_string(), timezone: "UTC".to_string(),
};
assert_eq!(trigger.function_name, "dailyCleanup");
assert_eq!(trigger.schedule, "0 2 * * *");
assert_eq!(trigger.timezone, "UTC");
}
#[test]
fn test_cron_trigger_parses_hourly_expression() {
let trigger = CronTrigger {
function_name: "hourlySync".to_string(),
schedule: "0 * * * *".to_string(), timezone: "UTC".to_string(),
};
assert_eq!(trigger.function_name, "hourlySync");
assert_eq!(trigger.schedule, "0 * * * *");
}
#[test]
fn test_cron_trigger_parses_every_5_minutes() {
let trigger = CronTrigger {
function_name: "frequentCheck".to_string(),
schedule: "*/5 * * * *".to_string(), timezone: "UTC".to_string(),
};
assert_eq!(trigger.schedule, "*/5 * * * *");
}
#[test]
fn test_cron_schedule_matches_exact_time() {
let schedule = CronSchedule::parse("0 2 * * *").expect("parse cron");
let matching_time = chrono::DateTime::parse_from_rfc3339("2024-03-15T02:00:00+00:00")
.expect("parse datetime")
.with_timezone(&chrono::Utc);
assert!(schedule.matches(&matching_time));
}
#[test]
fn test_cron_schedule_does_not_match_wrong_hour() {
let schedule = CronSchedule::parse("0 2 * * *").expect("parse cron");
let non_matching_time = chrono::DateTime::parse_from_rfc3339("2024-03-15T03:00:00+00:00")
.expect("parse datetime")
.with_timezone(&chrono::Utc);
assert!(!schedule.matches(&non_matching_time));
}
#[test]
fn test_cron_schedule_matches_every_5_minutes() {
let schedule = CronSchedule::parse("*/5 * * * *").expect("parse cron");
let time_00 = chrono::DateTime::parse_from_rfc3339("2024-03-15T10:00:00+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
let time_05 = chrono::DateTime::parse_from_rfc3339("2024-03-15T10:05:00+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
let time_03 = chrono::DateTime::parse_from_rfc3339("2024-03-15T10:03:00+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
assert!(schedule.matches(&time_00));
assert!(schedule.matches(&time_05));
assert!(!schedule.matches(&time_03));
}
#[test]
fn test_cron_trigger_tracks_last_execution() {
let mut state = CronExecutionState::new();
assert!(state.last_executed.is_none());
let now = chrono::Utc::now();
state.record_execution(now);
assert_eq!(state.last_executed, Some(now));
}
#[test]
fn test_cron_trigger_should_execute_first_time() {
let schedule = CronSchedule::parse("0 2 * * *").expect("parse cron");
let state = CronExecutionState::new();
let exec_time = chrono::DateTime::parse_from_rfc3339("2024-03-15T02:00:00+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
assert!(state.should_execute(&schedule, &exec_time));
}
#[test]
fn test_cron_trigger_prevents_duplicate_in_window() {
use crate::triggers::cron::CronDecision;
let schedule = CronSchedule::parse("0 2 * * *").expect("parse cron");
let mut state = CronExecutionState::new();
let exec_time = chrono::DateTime::parse_from_rfc3339("2024-03-15T02:00:03+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
assert!(state.should_execute(&schedule, &exec_time));
state.record_execution(exec_time);
let within_window = chrono::DateTime::parse_from_rfc3339("2024-03-15T02:00:52+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
match state.decide(&schedule, &within_window) {
CronDecision::AlreadyFired {
window_start,
last_executed,
} => {
assert_eq!(window_start.to_rfc3339(), "2024-03-15T02:00:00+00:00");
assert_eq!(last_executed, exec_time);
},
other => panic!("expected AlreadyFired, got {other:?}"),
}
let not_due = chrono::DateTime::parse_from_rfc3339("2024-03-15T02:05:00+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
assert_eq!(state.decide(&schedule, ¬_due), CronDecision::NotDue);
}
#[test]
fn test_cron_trigger_allows_next_window() {
let schedule = CronSchedule::parse("0 * * * *").expect("parse cron");
let mut state = CronExecutionState::new();
let time_200 = chrono::DateTime::parse_from_rfc3339("2024-03-15T02:00:00+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
assert!(state.should_execute(&schedule, &time_200));
state.record_execution(time_200);
let time_300 = chrono::DateTime::parse_from_rfc3339("2024-03-15T03:00:00+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
assert!(state.should_execute(&schedule, &time_300));
}
fn next_jitter_micros(rng: &mut u64) -> i64 {
*rng ^= *rng << 13;
*rng ^= *rng >> 7;
*rng ^= *rng << 17;
i64::try_from(*rng % 1_000_000).expect("micros < 1e6 fits i64")
}
fn simulate_cron_fires(
expr: &str,
start: chrono::DateTime<chrono::Utc>,
ticks: usize,
seed: u64,
) -> (usize, usize) {
let schedule = CronSchedule::parse(expr).expect("parse cron");
let mut state = CronExecutionState::new();
let mut rng = seed.max(1);
let mut fires = 0;
let mut expected = 0;
for n in 0..ticks {
let minute = start + chrono::Duration::minutes(i64::try_from(n).expect("tick fits i64"));
if schedule.matches(&minute) {
expected += 1;
}
let now = minute + chrono::Duration::microseconds(next_jitter_micros(&mut rng));
if state.should_execute(&schedule, &now) {
state.record_execution(now);
fires += 1;
}
}
(fires, expected)
}
#[test]
fn test_cron_fires_exactly_once_per_window_under_wall_clock_jitter() {
let start = chrono::DateTime::parse_from_rfc3339("2026-01-01T00:00:00+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
for expr in [
"* * * * *",
"*/5 * * * *",
"0 * * * *",
"0 2 * * *",
"0 3 * * 1",
] {
for seed in [1_u64, 42, 20_260_101] {
let (fires, expected) = simulate_cron_fires(expr, start, 20_000, seed);
assert_eq!(
fires, expected,
"`{expr}` (seed {seed}): fired {fires} times across 20000 jittered ticks, \
expected exactly one fire per matching window = {expected}"
);
}
}
}
#[test]
fn test_cron_daily_schedule_fires_on_consecutive_days_with_jitter() {
let schedule = CronSchedule::parse("0 2 * * *").expect("parse cron");
let mut state = CronExecutionState::new();
let day1 = chrono::DateTime::parse_from_rfc3339("2026-03-01T02:00:07.312+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
assert!(state.should_execute(&schedule, &day1), "first firing");
state.record_execution(day1);
let day2 = chrono::DateTime::parse_from_rfc3339("2026-03-02T02:00:07.905+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
assert!(
state.should_execute(&schedule, &day2),
"a new day is a new window: yesterday's firing must not suppress today's"
);
}
#[test]
fn test_cron_minutely_schedule_fires_next_minute_when_jitter_increases() {
let schedule = CronSchedule::parse("* * * * *").expect("parse cron");
let mut state = CronExecutionState::new();
let tick1 = chrono::DateTime::parse_from_rfc3339("2026-03-01T12:00:00.100+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
assert!(state.should_execute(&schedule, &tick1));
state.record_execution(tick1);
let tick2 = chrono::DateTime::parse_from_rfc3339("2026-03-01T12:01:00.900+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
assert!(
state.should_execute(&schedule, &tick2),
"12:01 is a new window regardless of sub-second jitter ordering"
);
}
#[test]
fn test_cron_does_not_double_fire_within_one_window_with_nonzero_seconds() {
let schedule = CronSchedule::parse("0 2 * * *").expect("parse cron");
let mut state = CronExecutionState::new();
let first = chrono::DateTime::parse_from_rfc3339("2026-03-01T02:00:07.300+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
assert!(state.should_execute(&schedule, &first));
state.record_execution(first);
let again = chrono::DateTime::parse_from_rfc3339("2026-03-01T02:00:41.900+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
assert!(
!state.should_execute(&schedule, &again),
"02:00:41 is the same 02:00 window — already fired"
);
}
#[test]
fn test_cron_fires_after_gap_longer_than_an_hour() {
let schedule = CronSchedule::parse("0 */3 * * *").expect("parse cron");
let mut state = CronExecutionState::new();
let first = chrono::DateTime::parse_from_rfc3339("2026-03-01T03:00:30.500+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
assert!(state.should_execute(&schedule, &first));
state.record_execution(first);
let three_hours_later = chrono::DateTime::parse_from_rfc3339("2026-03-01T06:00:10.200+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
assert!(
state.should_execute(&schedule, &three_hours_later),
"06:00 is a new window three hours after the 03:00 firing"
);
}
fn august_2026_at_nine(day: u32) -> chrono::DateTime<chrono::Utc> {
chrono::DateTime::parse_from_rfc3339(&format!("2026-08-{day:02}T09:00:00+00:00"))
.expect("parse")
.with_timezone(&chrono::Utc)
}
#[test]
fn test_cron_weekday_tokens_match_posix_days() {
for (token, matching_day) in [(0, 2), (1, 3), (2, 4), (3, 5), (4, 6), (5, 7), (6, 8)] {
let schedule =
CronSchedule::parse(&format!("0 9 * * {token}")).expect("parse weekday cron");
for day in 2..=8_u32 {
let time = august_2026_at_nine(day);
assert_eq!(
schedule.matches(&time),
day == matching_day,
"`0 9 * * {token}` on 2026-08-{day:02} ({})",
time.format("%A")
);
}
}
}
#[test]
fn test_cron_weekday_seven_is_sunday() {
let schedule = CronSchedule::parse("0 9 * * 7").expect("parse weekday cron");
assert!(schedule.matches(&august_2026_at_nine(2)), "7 must match Sunday");
for day in 3..=8_u32 {
assert!(!schedule.matches(&august_2026_at_nine(day)), "7 must match only Sunday");
}
}
#[test]
fn test_cron_weekday_range_mon_to_fri() {
let schedule = CronSchedule::parse("0 9 * * 1-5").expect("parse weekday cron");
let expectations = [
(2_u32, false), (3, true), (4, true), (5, true), (6, true), (7, true), (8, false), ];
for (day, expected) in expectations {
let time = august_2026_at_nine(day);
assert_eq!(
schedule.matches(&time),
expected,
"`0 9 * * 1-5` on 2026-08-{day:02} ({})",
time.format("%A")
);
}
}
#[test]
fn test_cron_trigger_catches_up_missed_executions() {
let schedule = CronSchedule::parse("0 * * * *").expect("parse cron");
let state = CronExecutionState::new();
let last_known = chrono::DateTime::parse_from_rfc3339("2024-03-15T01:00:00+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
let now = chrono::DateTime::parse_from_rfc3339("2024-03-15T03:30:00+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
let missed = state.find_missed_executions(&schedule, &last_known, &now);
assert_eq!(missed.len(), 2);
}
#[test]
fn test_cron_trigger_payload_includes_schedule_info() {
let trigger = CronTrigger {
function_name: "dailyCleanup".to_string(),
schedule: "0 2 * * *".to_string(),
timezone: "UTC".to_string(),
};
let exec_time = chrono::DateTime::parse_from_rfc3339("2024-03-15T02:00:00+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
let payload = trigger.build_payload(&exec_time);
assert_eq!(payload.trigger_type, "cron:dailyCleanup");
assert_eq!(payload.entity, "cron");
assert_eq!(payload.event_kind, "scheduled");
assert_eq!(payload.data["schedule"], "0 2 * * *");
assert_eq!(payload.data["timezone"], "UTC");
}
#[test]
fn test_cron_trigger_payload_includes_execution_time() {
let trigger = CronTrigger {
function_name: "hourlySync".to_string(),
schedule: "0 * * * *".to_string(),
timezone: "UTC".to_string(),
};
let exec_time = chrono::DateTime::parse_from_rfc3339("2024-03-15T14:00:00+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
let payload = trigger.build_payload(&exec_time);
assert!(payload.data.is_object());
assert!(payload.data["executed_at"].is_string());
assert_eq!(payload.data["executed_at"], "2024-03-15T14:00:00Z");
}
#[test]
fn test_cron_trigger_with_specific_timezone() {
let trigger = CronTrigger {
function_name: "morningReport".to_string(),
schedule: "0 9 * * *".to_string(), timezone: "America/New_York".to_string(),
};
assert_eq!(trigger.timezone, "America/New_York");
}
#[test]
fn test_cron_trigger_serialization() {
let trigger = CronTrigger {
function_name: "dailyCleanup".to_string(),
schedule: "0 2 * * *".to_string(),
timezone: "UTC".to_string(),
};
let json = serde_json::to_string(&trigger).expect("serialize");
let restored: CronTrigger = serde_json::from_str(&json).expect("deserialize");
assert_eq!(restored.function_name, trigger.function_name);
assert_eq!(restored.schedule, trigger.schedule);
assert_eq!(restored.timezone, trigger.timezone);
}
#[test]
fn test_cron_execution_state_serialization() {
let mut state = CronExecutionState::new();
let exec_time = chrono::DateTime::parse_from_rfc3339("2024-03-15T02:00:00+00:00")
.expect("parse")
.with_timezone(&chrono::Utc);
state.record_execution(exec_time);
let json = serde_json::to_string(&state).expect("serialize");
let restored: CronExecutionState = serde_json::from_str(&json).expect("deserialize");
assert_eq!(restored.last_executed, Some(exec_time));
}
#[test]
fn test_http_trigger_get_route() {
use crate::triggers::http::HttpTriggerRoute;
let route = HttpTriggerRoute {
function_name: "helloWorld".to_string(),
method: "GET".to_string(),
path: "/functions/v1/hello".to_string(),
requires_auth: false,
};
assert_eq!(route.function_name, "helloWorld");
assert_eq!(route.method, "GET");
assert_eq!(route.path, "/functions/v1/hello");
assert!(!route.requires_auth);
}
#[test]
fn test_http_trigger_post_route_with_auth() {
use crate::triggers::http::HttpTriggerRoute;
let route = HttpTriggerRoute {
function_name: "processData".to_string(),
method: "POST".to_string(),
path: "/functions/v1/process".to_string(),
requires_auth: true,
};
assert_eq!(route.function_name, "processData");
assert_eq!(route.method, "POST");
assert!(route.requires_auth);
}
#[test]
fn test_http_trigger_request_payload() {
use crate::triggers::http::HttpTriggerPayload;
let payload = HttpTriggerPayload {
method: "POST".to_string(),
path: "/functions/v1/users".to_string(),
headers: serde_json::json!({
"content-type": "application/json",
"x-user-id": "123"
}),
query: serde_json::json!({}),
params: serde_json::json!({
"id": "user-123"
}),
body: Some(serde_json::json!({
"name": "Alice",
"email": "alice@example.com"
})),
};
assert_eq!(payload.method, "POST");
assert_eq!(payload.path, "/functions/v1/users");
assert!(payload.body.is_some());
assert_eq!(payload.body.expect("body exists")["name"], "Alice");
}
#[test]
fn test_http_trigger_path_params() {
use crate::triggers::http::HttpTriggerPayload;
let payload = HttpTriggerPayload {
method: "GET".to_string(),
path: "/functions/v1/users/123".to_string(),
headers: serde_json::json!({}),
query: serde_json::json!({}),
params: serde_json::json!({
"id": "123"
}),
body: None,
};
assert_eq!(payload.params["id"], "123");
}
#[test]
fn test_http_trigger_response_custom_status() {
use crate::triggers::http::HttpTriggerResponse;
let response = HttpTriggerResponse {
status: 201,
headers: serde_json::json!({
"x-custom-header": "value"
}),
body: serde_json::json!({
"id": "new-user-123",
"created": true
}),
};
assert_eq!(response.status, 201);
assert_eq!(response.headers["x-custom-header"], "value");
assert_eq!(response.body["id"], "new-user-123");
}
#[test]
fn test_http_trigger_response_default_status() {
use crate::triggers::http::HttpTriggerResponse;
let response = HttpTriggerResponse {
status: 200,
headers: serde_json::json!({}),
body: serde_json::json!({"message": "OK"}),
};
assert_eq!(response.status, 200);
assert_eq!(response.body["message"], "OK");
}
#[test]
fn test_http_trigger_method_parsing() {
use crate::triggers::http::HttpTriggerRoute;
for method in &["GET", "POST", "PUT", "DELETE", "PATCH", "HEAD", "OPTIONS"] {
let route = HttpTriggerRoute {
function_name: "test".to_string(),
method: method.to_string(),
path: "/test".to_string(),
requires_auth: false,
};
assert_eq!(route.method, *method);
}
}
#[test]
fn test_http_trigger_route_matching() {
use crate::triggers::http::{HttpTriggerMatcher, HttpTriggerRoute};
let mut matcher = HttpTriggerMatcher::new();
matcher.add(HttpTriggerRoute {
function_name: "getUser".to_string(),
method: "GET".to_string(),
path: "/users/:id".to_string(),
requires_auth: true,
});
matcher.add(HttpTriggerRoute {
function_name: "createUser".to_string(),
method: "POST".to_string(),
path: "/users".to_string(),
requires_auth: true,
});
let route = matcher.find("GET", "/users/123");
assert!(route.is_some());
assert_eq!(route.expect("route matched").function_name, "getUser");
let route = matcher.find("POST", "/users");
assert!(route.is_some());
assert_eq!(route.expect("route matched").function_name, "createUser");
let route = matcher.find("GET", "/posts");
assert!(route.is_none());
}
#[test]
fn test_http_trigger_query_parameters() {
use crate::triggers::http::HttpTriggerPayload;
let payload = HttpTriggerPayload {
method: "GET".to_string(),
path: "/functions/v1/search".to_string(),
headers: serde_json::json!({}),
query: serde_json::json!({
"q": "alice",
"limit": 10
}),
params: serde_json::json!({}),
body: None,
};
assert_eq!(payload.query["q"], "alice");
assert_eq!(payload.query["limit"], 10);
}
#[test]
fn test_http_trigger_event_payload() {
use crate::triggers::http::HttpTriggerRoute;
let route = HttpTriggerRoute {
function_name: "handleRequest".to_string(),
method: "POST".to_string(),
path: "/functions/v1/webhook".to_string(),
requires_auth: false,
};
let http_payload = serde_json::json!({
"method": "POST",
"path": "/functions/v1/webhook",
"body": {"event": "user.created"}
});
let trigger_type = format!("http:{}:{}", route.method, route.path);
assert_eq!(trigger_type, "http:POST:/functions/v1/webhook");
let event = EventPayload {
trigger_type,
entity: "HttpRequest".to_string(),
event_kind: "request".to_string(),
data: http_payload,
timestamp: chrono::Utc::now(),
};
assert_eq!(event.entity, "HttpRequest");
assert_eq!(event.event_kind, "request");
}
#[test]
fn test_function_definition_creation() {
use crate::FunctionDefinition;
let func = FunctionDefinition::new(
"onUserCreated",
"after:mutation:createUser",
crate::RuntimeType::Deno,
);
assert_eq!(func.name, "onUserCreated");
assert_eq!(func.trigger, "after:mutation:createUser");
assert!(func.is_after_mutation());
assert!(!func.is_cron());
}
#[test]
fn test_function_definition_trigger_detection() {
use crate::{FunctionDefinition, RuntimeType};
let after_mutation =
FunctionDefinition::new("test", "after:mutation:createUser", RuntimeType::Deno);
assert!(after_mutation.is_after_mutation());
assert!(!after_mutation.is_before_mutation());
let before_mutation =
FunctionDefinition::new("test", "before:mutation:validateUser", RuntimeType::Deno);
assert!(before_mutation.is_before_mutation());
assert!(!before_mutation.is_after_mutation());
let cron = FunctionDefinition::new("test", "cron:0 * * * *", RuntimeType::Deno);
assert!(cron.is_cron());
let http = FunctionDefinition::new("test", "http:GET:/hello", RuntimeType::Deno);
assert!(http.is_http());
let storage =
FunctionDefinition::new("test", "after:storage:avatars:upload", RuntimeType::Deno);
assert!(storage.is_after_storage());
}
#[test]
fn test_function_definition_effective_timeout() {
use std::time::Duration;
use crate::{FunctionDefinition, RuntimeType};
let before_mutation =
FunctionDefinition::new("test", "before:mutation:createUser", RuntimeType::Deno);
assert_eq!(before_mutation.effective_timeout(), Duration::from_millis(500));
let after_mutation =
FunctionDefinition::new("test", "after:mutation:createUser", RuntimeType::Deno);
assert_eq!(after_mutation.effective_timeout(), Duration::from_secs(5));
let custom =
FunctionDefinition::new("test", "http:GET:/hello", RuntimeType::Deno).with_timeout(1000);
assert_eq!(custom.effective_timeout(), Duration::from_secs(1));
}
#[test]
fn test_trigger_registry_loads_definitions() {
use crate::{FunctionDefinition, RuntimeType};
let functions = [
FunctionDefinition::new("onUserCreated", "after:mutation:createUser", RuntimeType::Deno),
FunctionDefinition::new(
"validateUserInput",
"before:mutation:createUser",
RuntimeType::Deno,
),
FunctionDefinition::new("getUser", "http:GET:/users/:id", RuntimeType::Deno),
FunctionDefinition::new("dailyReport", "cron:0 2 * * *", RuntimeType::Deno),
];
assert_eq!(functions.len(), 4);
assert_eq!(functions[0].name, "onUserCreated");
assert!(functions[0].is_after_mutation());
assert!(functions[1].is_before_mutation());
assert!(functions[2].is_http());
assert!(functions[3].is_cron());
}
#[test]
fn test_trigger_registry_validates_format() {
use crate::{FunctionDefinition, RuntimeType};
let valid_triggers = vec![
"after:mutation:createUser",
"before:mutation:deleteUser",
"after:storage:avatars:upload",
"cron:0 * * * *",
"http:GET:/users/:id",
"http:POST:/data",
];
for trigger in valid_triggers {
let func = FunctionDefinition::new("test", trigger, RuntimeType::Deno);
assert_eq!(func.trigger, trigger);
}
}
#[test]
fn test_trigger_registry_multiple_same_type() {
use crate::{FunctionDefinition, RuntimeType};
let functions = [
FunctionDefinition::new("onUserCreated", "after:mutation:createUser", RuntimeType::Deno),
FunctionDefinition::new("onUserUpdated", "after:mutation:updateUser", RuntimeType::Deno),
FunctionDefinition::new("onUserDeleted", "after:mutation:deleteUser", RuntimeType::Deno),
];
assert_eq!(functions.iter().filter(|f| f.is_after_mutation()).count(), 3);
}
#[test]
fn test_mutation_with_before_hook_validation() {
use crate::{FunctionDefinition, RuntimeType};
let validation_hook = FunctionDefinition::new(
"validateUserInput",
"before:mutation:createUser",
RuntimeType::Deno,
);
let _input = serde_json::json!({
"name": "Alice",
"email": "alice@example.com"
});
assert!(validation_hook.is_before_mutation());
assert_eq!(validation_hook.name, "validateUserInput");
assert!(validation_hook.trigger.contains("before:mutation:createUser"));
}
#[test]
fn test_mutation_with_before_and_after_hooks() {
use crate::{FunctionDefinition, RuntimeType};
let before_hook =
FunctionDefinition::new("validateCreate", "before:mutation:createUser", RuntimeType::Deno);
let after_hook =
FunctionDefinition::new("logCreated", "after:mutation:createUser", RuntimeType::Deno);
assert!(before_hook.is_before_mutation());
assert!(after_hook.is_after_mutation());
assert_eq!(before_hook.name, "validateCreate");
assert_eq!(after_hook.name, "logCreated");
assert!(before_hook.trigger.contains("createUser"));
assert!(after_hook.trigger.contains("createUser"));
}
#[test]
fn test_after_mutation_and_storage_trigger_cascade() {
use crate::{FunctionDefinition, RuntimeType};
let mutation_hook = FunctionDefinition::new(
"onAvatarUpload",
"after:mutation:updateUserAvatar",
RuntimeType::Deno,
);
let storage_hook =
FunctionDefinition::new("processAvatar", "after:storage:avatars:upload", RuntimeType::Deno);
assert!(mutation_hook.is_after_mutation());
assert!(storage_hook.is_after_storage());
assert_eq!(mutation_hook.name, "onAvatarUpload");
assert_eq!(storage_hook.name, "processAvatar");
}
#[test]
fn test_cron_and_http_trigger_coexist() {
use crate::{FunctionDefinition, RuntimeType};
let cron_job = FunctionDefinition::new("dailyReport", "cron:0 2 * * *", RuntimeType::Deno);
let http_endpoint =
FunctionDefinition::new("getMetrics", "http:GET:/metrics", RuntimeType::Deno);
assert!(cron_job.is_cron());
assert!(http_endpoint.is_http());
assert_ne!(cron_job.name, http_endpoint.name);
assert_ne!(cron_job.trigger, http_endpoint.trigger);
}
#[test]
fn test_before_mutation_timeout_during_cascade() {
use crate::{FunctionDefinition, RuntimeType};
let before_hook =
FunctionDefinition::new("slowValidation", "before:mutation:deleteUser", RuntimeType::Deno);
let effective_timeout = before_hook.effective_timeout();
assert_eq!(effective_timeout.as_millis(), 500);
assert!(before_hook.is_before_mutation());
}
#[test]
fn test_trigger_registry_startup_with_all_types() {
use crate::{FunctionDefinition, RuntimeType};
let functions = [
FunctionDefinition::new("onUserCreated", "after:mutation:createUser", RuntimeType::Deno),
FunctionDefinition::new("validateUser", "before:mutation:createUser", RuntimeType::Deno),
FunctionDefinition::new("processFile", "after:storage:uploads:upload", RuntimeType::Deno),
FunctionDefinition::new("hourlySync", "cron:0 * * * *", RuntimeType::Deno),
FunctionDefinition::new("apiHandler", "http:POST:/api/process", RuntimeType::Deno),
];
assert_eq!(functions.iter().filter(|f| f.is_after_mutation()).count(), 1);
assert_eq!(functions.iter().filter(|f| f.is_before_mutation()).count(), 1);
assert_eq!(functions.iter().filter(|f| f.is_after_storage()).count(), 1);
assert_eq!(functions.iter().filter(|f| f.is_cron()).count(), 1);
assert_eq!(functions.iter().filter(|f| f.is_http()).count(), 1);
}
#[test]
fn test_trigger_graceful_shutdown() {
use crate::{FunctionDefinition, RuntimeType};
let functions = [
FunctionDefinition::new("onUserCreated", "after:mutation:createUser", RuntimeType::Deno),
FunctionDefinition::new("validateUser", "before:mutation:createUser", RuntimeType::Deno),
FunctionDefinition::new("hourlySync", "cron:0 * * * *", RuntimeType::Deno),
];
assert_eq!(functions.len(), 3);
}
#[test]
fn test_trigger_error_recovery() {
use crate::{FunctionDefinition, RuntimeType};
let functions = [
FunctionDefinition::new("failingFunction", "after:mutation:deleteUser", RuntimeType::Deno),
FunctionDefinition::new("workingFunction", "after:mutation:createUser", RuntimeType::Deno),
];
assert_eq!(functions.len(), 2);
assert!(functions[0].is_after_mutation());
assert!(functions[1].is_after_mutation());
}
#[test]
fn test_http_trigger_with_auth_context() {
use crate::{FunctionDefinition, RuntimeType};
let public_endpoint =
FunctionDefinition::new("publicMetrics", "http:GET:/public/metrics", RuntimeType::Deno);
let protected_endpoint =
FunctionDefinition::new("adminPanel", "http:GET:/admin/dashboard", RuntimeType::Deno);
assert!(public_endpoint.is_http());
assert!(protected_endpoint.is_http());
}
#[test]
fn test_cron_with_storage_cascade() {
use crate::{FunctionDefinition, RuntimeType};
let cron_job = FunctionDefinition::new("backupDaily", "cron:0 3 * * *", RuntimeType::Deno);
let storage_trigger =
FunctionDefinition::new("archiveBackup", "after:storage:backups:upload", RuntimeType::Deno);
assert!(cron_job.is_cron());
assert!(storage_trigger.is_after_storage());
assert_ne!(cron_job.name, storage_trigger.name);
}
#[test]
fn test_registry_loads_cron_triggers() {
use crate::{FunctionDefinition, RuntimeType, triggers::registry::TriggerRegistry};
let functions = vec![
FunctionDefinition::new("dailyCleanup", "cron:0 2 * * *", RuntimeType::Deno),
FunctionDefinition::new("hourlySync", "cron:0 * * * *", RuntimeType::Deno),
FunctionDefinition::new("onUserCreated", "after:mutation:User:insert", RuntimeType::Deno),
];
let registry = TriggerRegistry::load_from_definitions(&functions).expect("load registry");
assert_eq!(registry.cron_trigger_count(), 2, "should have 2 cron triggers");
assert_eq!(registry.cron_triggers[0].function_name, "dailyCleanup");
assert_eq!(registry.cron_triggers[0].schedule, "0 2 * * *");
assert_eq!(registry.cron_triggers[1].function_name, "hourlySync");
assert_eq!(registry.cron_triggers[1].schedule, "0 * * * *");
}
#[test]
fn test_registry_cron_scheduler_returns_some_when_triggers_exist() {
use crate::{FunctionDefinition, RuntimeType, triggers::registry::TriggerRegistry};
let functions = vec![FunctionDefinition::new(
"dailyJob",
"cron:0 3 * * *",
RuntimeType::Deno,
)];
let registry = TriggerRegistry::load_from_definitions(&functions).expect("load registry");
let scheduler = registry.cron_scheduler();
assert!(scheduler.is_some(), "should return a scheduler when cron triggers exist");
assert_eq!(
scheduler
.expect("cron_scheduler should return Some when triggers exist")
.trigger_count(),
1
);
}
#[test]
fn test_registry_cron_scheduler_returns_none_when_no_triggers() {
use crate::{FunctionDefinition, RuntimeType, triggers::registry::TriggerRegistry};
let functions = vec![
FunctionDefinition::new("onUserCreated", "after:mutation:User:insert", RuntimeType::Deno),
FunctionDefinition::new("validate", "before:mutation:createUser", RuntimeType::Deno),
];
let registry = TriggerRegistry::load_from_definitions(&functions).expect("load registry");
assert!(
registry.cron_scheduler().is_none(),
"no cron triggers → cron_scheduler() should return None (fast path)"
);
}
#[test]
fn test_cron_scheduler_new_creates_with_triggers() {
use crate::triggers::cron::{CronScheduler, CronTrigger};
let triggers = vec![
CronTrigger {
function_name: "dailyCleanup".to_string(),
schedule: "0 2 * * *".to_string(),
timezone: "UTC".to_string(),
},
CronTrigger {
function_name: "hourlySync".to_string(),
schedule: "0 * * * *".to_string(),
timezone: "UTC".to_string(),
},
];
let scheduler = CronScheduler::new(triggers);
assert_eq!(scheduler.trigger_count(), 2);
}
#[tokio::test]
async fn test_cron_scheduler_starts_and_provides_handle() {
use std::{collections::HashMap, sync::Arc};
use crate::{
observer::FunctionObserver,
triggers::cron::{CronScheduler, CronTrigger},
};
let triggers = vec![CronTrigger {
function_name: "dailyCleanup".to_string(),
schedule: "0 2 * * *".to_string(), timezone: "UTC".to_string(),
}];
let observer = Arc::new(FunctionObserver::new());
let handle = CronScheduler::new(triggers).start(observer, HashMap::new());
handle.stop();
}
#[tokio::test]
async fn test_cron_scheduler_handle_stops_gracefully() {
use std::{collections::HashMap, sync::Arc};
use crate::{
observer::FunctionObserver,
triggers::cron::{CronScheduler, CronTrigger},
};
let triggers = vec![CronTrigger {
function_name: "neverFires".to_string(),
schedule: "0 0 31 2 *".to_string(), timezone: "UTC".to_string(),
}];
let observer = Arc::new(FunctionObserver::new());
let handle = CronScheduler::new(triggers).start(observer, HashMap::new());
handle.stop();
tokio::task::yield_now().await;
}
#[tokio::test]
async fn test_cron_scheduler_empty_starts_cleanly() {
use std::{collections::HashMap, sync::Arc};
use crate::{observer::FunctionObserver, triggers::cron::CronScheduler};
let scheduler = CronScheduler::new(vec![]);
assert_eq!(scheduler.trigger_count(), 0);
let observer = Arc::new(FunctionObserver::new());
let handle = scheduler.start(observer, HashMap::new());
handle.stop();
}