Skip to main content

arcula/
operations.rs

1use std::fs;
2use std::path::PathBuf;
3
4use anyhow::{Context, Result};
5use serde::{Deserialize, Serialize};
6
7use crate::approvals::{self, ApprovalMode, ApprovalRecord};
8use crate::config::{get_connection_policy, MongoConfig};
9use crate::core::sync::{perform_sync, SyncReport};
10use crate::plans::{self, PlanStatus, SyncPlanRecord};
11use crate::storage;
12use crate::utils::mongodb;
13
14#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
15#[serde(rename_all = "snake_case")]
16pub enum OperationKind {
17    Sync,
18    Revert,
19}
20
21#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
22#[serde(rename_all = "snake_case")]
23pub enum OperationStatus {
24    Running,
25    Completed,
26    Failed,
27}
28
29#[derive(Debug, Clone, Serialize, Deserialize)]
30pub struct OperationStatusEvent {
31    pub at: String,
32    pub status: OperationStatus,
33    pub message: String,
34}
35
36#[derive(Debug, Clone, Serialize, Deserialize)]
37pub struct OperationRecord {
38    pub version: u8,
39    pub id: String,
40    pub kind: OperationKind,
41    pub plan_id: Option<String>,
42    pub plan_hash: Option<String>,
43    pub original_operation_id: Option<String>,
44    pub created_at: String,
45    pub updated_at: String,
46    pub status: OperationStatus,
47    pub status_history: Vec<OperationStatusEvent>,
48    pub approval: Option<ApprovalRecord>,
49    pub sync_report: Option<SyncReport>,
50    pub error: Option<String>,
51}
52
53#[derive(Debug, Clone, Serialize)]
54pub struct RevertPreview {
55    pub operation_id: String,
56    pub target_env: String,
57    pub target_db: String,
58    pub backup_path: String,
59    pub target_policy_requires_human_approval: bool,
60}
61
62pub async fn run_plan(plan_id: &str, agent: bool) -> Result<OperationRecord> {
63    let plan = plans::load_plan(plan_id)?;
64    if matches!(plan.status, PlanStatus::Running | PlanStatus::Completed) {
65        anyhow::bail!(
66            "Plan '{}' is {:?} and cannot be run again. Create a new plan instead.",
67            plan.id,
68            plan.status
69        );
70    }
71    let approval = ensure_plan_approval(&plan, agent)?;
72
73    let mut operation = OperationRecord::new_sync(&plan, approval.clone());
74    save_operation(&operation)?;
75    plans::update_plan_status(&plan.id, PlanStatus::Running)?;
76
77    match perform_sync(plan.config.clone()).await {
78        Ok(report) => {
79            operation.sync_report = Some(report);
80            operation.push_status(OperationStatus::Completed, "sync completed");
81            save_operation(&operation)?;
82            plans::update_plan_status(&plan.id, PlanStatus::Completed)?;
83            Ok(operation)
84        }
85        Err(err) => {
86            operation.error = Some(format!("{err:#}"));
87            operation.push_status(OperationStatus::Failed, "sync failed");
88            save_operation(&operation)?;
89            let _ = plans::update_plan_status(&plan.id, PlanStatus::Failed);
90            Err(err)
91        }
92    }
93}
94
95pub fn preview_revert(operation_id: &str) -> Result<RevertPreview> {
96    let operation = load_operation(operation_id)?;
97    let report = operation
98        .sync_report
99        .as_ref()
100        .context("Operation has no sync report and cannot be reverted")?;
101    let backup_path = report
102        .backup_path
103        .clone()
104        .context("Operation has no backup path and cannot be reverted")?;
105    if !std::path::Path::new(&backup_path).exists() {
106        anyhow::bail!("Backup path does not exist: {backup_path}");
107    }
108    let target_policy = get_connection_policy(&report.target_env);
109
110    Ok(RevertPreview {
111        operation_id: operation.id,
112        target_env: report.target_env.name().to_string(),
113        target_db: report.target_db.clone(),
114        backup_path,
115        target_policy_requires_human_approval: target_policy.human_approval_required,
116    })
117}
118
119pub async fn revert_operation(operation_id: &str) -> Result<OperationRecord> {
120    crate::config::check_mongodb_tools()
121        .map_err(|err| anyhow::anyhow!("MongoDB tools not found: {err}"))?;
122
123    let original = load_operation(operation_id)?;
124    let report = original
125        .sync_report
126        .as_ref()
127        .context("Operation has no sync report and cannot be reverted")?;
128    let backup_path = report
129        .backup_path
130        .clone()
131        .context("Operation has no backup path and cannot be reverted")?;
132    if !std::path::Path::new(&backup_path).exists() {
133        anyhow::bail!("Backup path does not exist: {backup_path}");
134    }
135
136    let target_policy = get_connection_policy(&report.target_env);
137    let approval = if target_policy.human_approval_required {
138        Some(approvals::create_human_presence_approval(
139            &format!("revert:{}", original.id),
140            &revert_hash(&original)?,
141        )?)
142    } else {
143        Some(approvals::create_agent_policy_approval(
144            &format!("revert:{}", original.id),
145            &revert_hash(&original)?,
146        )?)
147    };
148
149    let mut operation = OperationRecord::new_revert(&original, approval);
150    save_operation(&operation)?;
151
152    let target_config = MongoConfig::from_env(report.target_env.clone()).context(format!(
153        "Failed to get configuration for {}",
154        report.target_env
155    ))?;
156    match mongodb::restore_backup(
157        &target_config,
158        &report.target_db,
159        std::path::Path::new(&backup_path),
160    )
161    .await
162    {
163        Ok(()) => {
164            let mut revert_report = report.clone();
165            revert_report.restored_from_backup = true;
166            operation.sync_report = Some(revert_report);
167            operation.push_status(OperationStatus::Completed, "revert completed");
168            save_operation(&operation)?;
169            Ok(operation)
170        }
171        Err(err) => {
172            operation.error = Some(format!("{err:#}"));
173            operation.push_status(OperationStatus::Failed, "revert failed");
174            save_operation(&operation)?;
175            Err(err)
176        }
177    }
178}
179
180pub fn load_operation(id: &str) -> Result<OperationRecord> {
181    let path = operation_path(id);
182    let contents = fs::read_to_string(&path)
183        .with_context(|| format!("Failed to read operation {} from {}", id, path.display()))?;
184    serde_json::from_str(&contents).with_context(|| format!("Failed to parse operation {id}"))
185}
186
187pub fn list_operations() -> Result<Vec<OperationRecord>> {
188    let dir = storage::operations_dir();
189    if !dir.exists() {
190        return Ok(Vec::new());
191    }
192
193    let mut operations = Vec::new();
194    for entry in fs::read_dir(&dir).with_context(|| format!("Failed to read {}", dir.display()))? {
195        let entry = entry?;
196        let path = entry.path().join("operation.json");
197        if path.exists() {
198            let contents = fs::read_to_string(&path)?;
199            let operation: OperationRecord = serde_json::from_str(&contents)?;
200            operations.push(operation);
201        }
202    }
203    operations.sort_by(|a, b| b.created_at.cmp(&a.created_at));
204    Ok(operations)
205}
206
207fn ensure_plan_approval(plan: &SyncPlanRecord, agent: bool) -> Result<Option<ApprovalRecord>> {
208    if let Some(approval) = plans::load_approval(&plan.id)? {
209        approvals::verify_approval(&approval, &plan.id, &plan.hash)?;
210        if plan.requires_human_approval && approval.mode != ApprovalMode::HumanPresence {
211            anyhow::bail!("Plan '{}' requires human OS approval", plan.id);
212        }
213        return Ok(Some(approval));
214    }
215
216    if plan.requires_human_approval {
217        anyhow::bail!(
218            "Plan '{}' requires human approval. Run: arcula plan approve {}",
219            plan.id,
220            plan.id
221        );
222    }
223
224    if agent && !plan.target_policy.allow_agent_apply {
225        anyhow::bail!(
226            "Target policy does not allow agent execution for plan '{}'",
227            plan.id
228        );
229    }
230
231    let approval = approvals::create_agent_policy_approval(&plan.id, &plan.hash)?;
232    plans::save_approval(&plan.id, &approval)?;
233    plans::update_plan_status(&plan.id, PlanStatus::Approved)?;
234    Ok(Some(approval))
235}
236
237fn save_operation(operation: &OperationRecord) -> Result<()> {
238    storage::atomic_write_json(&operation_path(&operation.id), operation)
239}
240
241fn operation_dir(id: &str) -> PathBuf {
242    storage::operations_dir().join(id)
243}
244
245fn operation_path(id: &str) -> PathBuf {
246    operation_dir(id).join("operation.json")
247}
248
249fn revert_hash(operation: &OperationRecord) -> Result<String> {
250    let bytes = serde_json::to_vec(operation)?;
251    use sha2::{Digest, Sha256};
252    let mut hasher = Sha256::new();
253    hasher.update(bytes);
254    Ok(to_hex(&hasher.finalize()))
255}
256
257fn generate_operation_id(prefix: &str) -> String {
258    let timestamp = chrono::Utc::now().format("%Y%m%d%H%M%S");
259    let suffix: u32 = rand::random();
260    format!("{prefix}_{timestamp}_{suffix:08x}")
261}
262
263fn now() -> String {
264    chrono::Utc::now().to_rfc3339()
265}
266
267impl OperationRecord {
268    fn new_sync(plan: &SyncPlanRecord, approval: Option<ApprovalRecord>) -> Self {
269        let created_at = now();
270        let mut operation = Self {
271            version: 1,
272            id: generate_operation_id("op"),
273            kind: OperationKind::Sync,
274            plan_id: Some(plan.id.clone()),
275            plan_hash: Some(plan.hash.clone()),
276            original_operation_id: None,
277            created_at: created_at.clone(),
278            updated_at: created_at,
279            status: OperationStatus::Running,
280            status_history: Vec::new(),
281            approval,
282            sync_report: None,
283            error: None,
284        };
285        operation.push_status(OperationStatus::Running, "sync started");
286        operation
287    }
288
289    fn new_revert(original: &OperationRecord, approval: Option<ApprovalRecord>) -> Self {
290        let created_at = now();
291        let mut operation = Self {
292            version: 1,
293            id: generate_operation_id("revert"),
294            kind: OperationKind::Revert,
295            plan_id: original.plan_id.clone(),
296            plan_hash: original.plan_hash.clone(),
297            original_operation_id: Some(original.id.clone()),
298            created_at: created_at.clone(),
299            updated_at: created_at,
300            status: OperationStatus::Running,
301            status_history: Vec::new(),
302            approval,
303            sync_report: None,
304            error: None,
305        };
306        operation.push_status(OperationStatus::Running, "revert started");
307        operation
308    }
309
310    fn push_status(&mut self, status: OperationStatus, message: &str) {
311        self.status = status.clone();
312        self.updated_at = now();
313        self.status_history.push(OperationStatusEvent {
314            at: self.updated_at.clone(),
315            status,
316            message: message.to_string(),
317        });
318    }
319}
320
321fn to_hex(bytes: &[u8]) -> String {
322    const HEX: &[u8; 16] = b"0123456789abcdef";
323    let mut output = String::with_capacity(bytes.len() * 2);
324    for byte in bytes {
325        output.push(HEX[(byte >> 4) as usize] as char);
326        output.push(HEX[(byte & 0x0f) as usize] as char);
327    }
328    output
329}