use std::sync::RwLock;
use std::time::{Duration, SystemTime};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum GateDecision {
Allow,
Review,
Block,
}
#[derive(Debug, Clone)]
struct ApprovalEntry {
worker_id: String,
decision: GateDecision,
approved_at: SystemTime,
}
pub struct ApprovalRegistry {
entries: RwLock<Vec<ApprovalEntry>>,
max_entries: usize,
ttl: Duration,
}
impl ApprovalRegistry {
pub fn new(max_entries: usize, ttl: Duration) -> Self {
Self {
entries: RwLock::new(Vec::new()),
max_entries,
ttl,
}
}
pub fn approve(&self, _proposal_id: &str, worker_id: &str, decision: GateDecision) {
if decision == GateDecision::Block {
return;
}
let mut entries = self.entries.write().unwrap();
entries.retain(|e| e.is_valid(self.ttl));
if entries.len() >= self.max_entries {
entries.remove(0);
}
entries.push(ApprovalEntry {
worker_id: worker_id.into(),
decision,
approved_at: SystemTime::now(),
});
}
pub fn is_approved(&self, _proposal_id: &str, worker_id: &str) -> bool {
let entries = self.entries.read().unwrap();
entries.iter().any(|e| {
e.worker_id == worker_id
&& e.decision != GateDecision::Block
&& e.is_valid(self.ttl)
})
}
pub fn purge(&self) {
let mut entries = self.entries.write().unwrap();
entries.retain(|e| e.is_valid(self.ttl));
}
}
impl Default for ApprovalRegistry {
fn default() -> Self {
Self::new(100, Duration::from_secs(300))
}
}
impl ApprovalEntry {
fn is_valid(&self, ttl: Duration) -> bool {
SystemTime::now()
.duration_since(self.approved_at)
.map(|elapsed| elapsed < ttl)
.unwrap_or(false)
}
}
pub fn validate_drain_request(
worker_id: &str,
proposal_id: &str,
drain_budget_secs: Option<u64>,
registry: &ApprovalRegistry,
) -> Result<(), String> {
if proposal_id.is_empty() {
return Err("Drain request rejected: empty proposal_id".into());
}
if worker_id.is_empty() {
return Err("Drain request rejected: empty worker_id".into());
}
if !registry.is_approved(proposal_id, worker_id) {
return Err(format!(
"Drain request rejected: proposal_id '{}' is not approved for worker '{}'",
proposal_id, worker_id
));
}
if let Some(budget) = drain_budget_secs {
if budget == 0 {
return Err("Drain request rejected: drain_budget_secs must be > 0".into());
}
if budget > 300 {
return Err(format!(
"Drain request rejected: drain_budget_secs {} exceeds max 300s",
budget
));
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
fn registry() -> ApprovalRegistry {
ApprovalRegistry::new(10, Duration::from_secs(300))
}
#[test]
fn test_approved_proposal_is_valid() {
let reg = registry();
reg.approve("prop-abc", "chest", GateDecision::Allow);
assert!(validate_drain_request("chest", "prop-abc", Some(30), ®).is_ok());
}
#[test]
fn test_unapproved_proposal_is_rejected() {
let reg = registry();
assert!(validate_drain_request("chest", "prop-unknown", None, ®).is_err());
}
#[test]
fn test_empty_proposal_rejected() {
let reg = registry();
assert!(validate_drain_request("chest", "", None, ®).is_err());
}
#[test]
fn test_empty_worker_rejected() {
let reg = registry();
assert!(validate_drain_request("", "prop-abc", None, ®).is_err());
}
#[test]
fn test_blocked_proposal_not_allowed() {
let reg = registry();
reg.approve("prop-blocked", "chest", GateDecision::Block);
assert!(validate_drain_request("chest", "prop-blocked", None, ®).is_err());
}
#[test]
fn test_zero_budget_rejected() {
let reg = registry();
reg.approve("p1", "test", GateDecision::Allow);
assert!(validate_drain_request("test", "p1", Some(0), ®).is_err());
}
#[test]
fn test_excessive_budget_rejected() {
let reg = registry();
reg.approve("p2", "test", GateDecision::Allow);
assert!(validate_drain_request("test", "p2", Some(301), ®).is_err());
}
#[test]
fn test_max_budget_accepted() {
let reg = registry();
reg.approve("p3", "test", GateDecision::Allow);
assert!(validate_drain_request("test", "p3", Some(300), ®).is_ok());
}
#[test]
fn test_wrong_worker_rejected() {
let reg = registry();
reg.approve("p4", "chest", GateDecision::Allow);
assert!(validate_drain_request("evernight", "p4", None, ®).is_err());
}
#[test]
fn test_review_decision_is_valid() {
let reg = registry();
reg.approve("p5", "chest", GateDecision::Review);
assert!(validate_drain_request("chest", "p5", None, ®).is_ok());
}
#[test]
fn test_purge_removes_expired() {
let reg = ApprovalRegistry::new(10, Duration::ZERO);
reg.approve("p6", "chest", GateDecision::Allow);
std::thread::sleep(Duration::from_millis(10));
let reg2 = ApprovalRegistry::new(10, Duration::ZERO);
assert!(validate_drain_request("chest", "p6", None, ®).is_err());
}
}