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}