1use serde::{Deserialize, Serialize};
8use std::path::{Path, PathBuf};
9
10#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
12pub struct CronInstructions {
13 pub project: String,
15 pub phase: u32,
17 pub status: String,
19 pub retry_after: String,
21 pub resume: ResumeCommand,
23 pub hermes_cron: HermesCronJob,
25}
26
27#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
29pub struct ResumeCommand {
30 pub command: String,
32 pub args: Vec<String>,
34}
35
36#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
38pub struct HermesCronJob {
39 pub schedule: String,
41 pub name: String,
43 pub command: String,
45 pub once: bool,
47}
48
49#[derive(Debug, thiserror::Error)]
51pub enum ShipError {
52 #[error("ship I/O failed: {0}")]
54 Io(#[from] std::io::Error),
55 #[error("ship JSON failed: {0}")]
57 Json(#[from] serde_json::Error),
58 #[error("no last-ship record found — nothing to confirm or reject")]
60 Missing,
61}
62
63pub fn cron_instructions_path(project_root: &Path, phase: u32) -> PathBuf {
67 project_root
68 .join(".devflow")
69 .join(format!("cron-instructions-{phase:02}.json"))
70}
71
72pub(crate) fn legacy_cron_instructions_path(project_root: &Path) -> PathBuf {
75 project_root.join(".devflow").join("cron-instructions.json")
76}
77
78pub fn write_cron_instructions(
80 project_root: &Path,
81 instructions: &CronInstructions,
82) -> Result<(), ShipError> {
83 let path = cron_instructions_path(project_root, instructions.phase);
84 if let Some(parent) = path.parent() {
85 crate::workflow::ensure_devflow_dir(parent)?;
86 }
87 std::fs::write(&path, serde_json::to_string_pretty(instructions)?)?;
88 Ok(())
89}
90
91pub fn load_cron_instructions(
94 project_root: &Path,
95 phase: u32,
96) -> Result<CronInstructions, ShipError> {
97 let path = cron_instructions_path(project_root, phase);
98 if path.exists() {
99 return Ok(serde_json::from_str(&std::fs::read_to_string(&path)?)?);
100 }
101 let legacy = legacy_cron_instructions_path(project_root);
102 if legacy.exists() {
103 let instructions: CronInstructions =
104 serde_json::from_str(&std::fs::read_to_string(&legacy)?)?;
105 if instructions.phase == phase {
106 return Ok(instructions);
107 }
108 }
109 Err(ShipError::Missing)
110}
111
112pub fn list_cron_instructions(project_root: &Path) -> Vec<CronInstructions> {
115 let mut found = Vec::new();
116 if let Ok(entries) = std::fs::read_dir(project_root.join(".devflow")) {
117 for entry in entries.flatten() {
118 let name = entry.file_name();
119 let Some(name) = name.to_str() else { continue };
120 if !name.starts_with("cron-instructions") || !name.ends_with(".json") {
121 continue;
122 }
123 if let Ok(contents) = std::fs::read_to_string(entry.path())
124 && let Ok(instructions) = serde_json::from_str::<CronInstructions>(&contents)
125 {
126 found.push(instructions);
127 }
128 }
129 }
130 found.sort_by_key(|i| i.phase);
131 found.dedup_by_key(|i| i.phase);
132 found
133}
134
135pub fn delete_cron_instructions(project_root: &Path, phase: u32) -> Result<(), ShipError> {
138 let path = cron_instructions_path(project_root, phase);
139 if path.exists() {
140 std::fs::remove_file(path)?;
141 }
142 let legacy = legacy_cron_instructions_path(project_root);
143 if legacy.exists()
144 && let Ok(contents) = std::fs::read_to_string(&legacy)
145 && serde_json::from_str::<CronInstructions>(&contents)
146 .map(|i| i.phase == phase)
147 .unwrap_or(true)
148 {
149 std::fs::remove_file(&legacy)?;
150 }
151 Ok(())
152}
153
154pub fn build_single_agent_cron_instructions(
160 project_root: &Path,
161 phase: u32,
162 retry_after: &str,
163) -> CronInstructions {
164 let project = project_root.display().to_string();
165 let args = vec![
166 "resume".to_string(),
167 "--phase".to_string(),
168 phase.to_string(),
169 ];
170 CronInstructions {
171 project: project.clone(),
172 phase,
173 status: "rate_limited".to_string(),
174 retry_after: retry_after.to_string(),
175 resume: ResumeCommand {
176 command: "devflow".to_string(),
177 args,
178 },
179 hermes_cron: HermesCronJob {
180 schedule: cron_schedule_from_retry_after(retry_after).unwrap_or_default(),
181 name: format!("devflow-phase-{phase:02}-resume"),
182 command: format!(
183 "cd {} && devflow resume --phase {phase}",
184 shell_quote(&project)
185 ),
186 once: true,
187 },
188 }
189}
190
191pub fn cron_schedule_from_retry_after(retry_after: &str) -> Option<String> {
194 parse_retry_timestamp(retry_after).map(|ts| ts.round_up_minute().to_cron())
196}
197
198#[derive(Debug, Clone, Copy, PartialEq, Eq)]
199struct RetryTimestamp {
200 year: i32,
201 month: u32,
202 day: u32,
203 hour: u32,
204 minute: u32,
205 second: u32,
206}
207
208impl RetryTimestamp {
209 fn round_up_minute(self) -> Self {
210 if self.second == 0 {
211 return self;
212 }
213 Self::from_epoch_minutes(self.to_epoch_minutes() + 1)
214 }
215
216 fn to_cron(self) -> String {
217 format!(
218 "{} {} {} {} *",
219 self.minute, self.hour, self.day, self.month
220 )
221 }
222
223 fn to_epoch_minutes(self) -> i64 {
224 let days = days_from_civil(self.year, self.month, self.day);
225 days * 24 * 60 + i64::from(self.hour) * 60 + i64::from(self.minute)
226 }
227
228 fn from_epoch_minutes(minutes: i64) -> Self {
229 let days = minutes.div_euclid(24 * 60);
230 let minute_of_day = minutes.rem_euclid(24 * 60);
231 let (year, month, day) = civil_from_days(days);
232 Self {
233 year,
234 month,
235 day,
236 hour: (minute_of_day / 60) as u32,
237 minute: (minute_of_day % 60) as u32,
238 second: 0,
239 }
240 }
241}
242
243fn parse_retry_timestamp(input: &str) -> Option<RetryTimestamp> {
244 parse_unix_seconds(input).or_else(|| parse_rfc3339ish(input))
245}
246
247fn parse_unix_seconds(input: &str) -> Option<RetryTimestamp> {
248 let seconds = input.trim().parse::<i64>().ok()?;
249 let minutes = seconds.div_euclid(60) + i64::from(seconds.rem_euclid(60) > 0);
250 Some(RetryTimestamp::from_epoch_minutes(minutes))
251}
252
253fn parse_rfc3339ish(input: &str) -> Option<RetryTimestamp> {
254 let input = input.trim();
255 let split_at = input.find('T').or_else(|| input.find(' '))?;
256 let (date, rest) = input.split_at(split_at);
257 let time = rest.get(1..)?;
258 let mut date_parts = date.split('-');
259 let year = date_parts.next()?.parse::<i32>().ok()?;
260 let month = date_parts.next()?.parse::<u32>().ok()?;
261 let day = date_parts.next()?.parse::<u32>().ok()?;
262 if date_parts.next().is_some() {
263 return None;
264 }
265
266 let (time, offset_minutes) = split_time_and_offset(time);
267 let mut time_parts = time.split(':');
268 let hour = time_parts.next()?.parse::<u32>().ok()?;
269 let minute = time_parts.next()?.parse::<u32>().ok()?;
270 let second = time_parts
271 .next()
272 .map(|s| s.split('.').next().unwrap_or_default().parse::<u32>().ok())
273 .unwrap_or(Some(0))?;
274 if month == 0 || month > 12 || day == 0 || day > 31 || hour > 23 || minute > 59 || second > 60 {
275 return None;
276 }
277
278 let ts = RetryTimestamp {
279 year,
280 month,
281 day,
282 hour,
283 minute,
284 second,
285 };
286 let utc_minutes = ts.to_epoch_minutes() - i64::from(offset_minutes);
287 let mut normalized = RetryTimestamp::from_epoch_minutes(utc_minutes);
288 normalized.second = second;
295 Some(normalized)
296}
297
298fn split_time_and_offset(time: &str) -> (&str, i32) {
299 let trimmed = time.trim_end_matches('Z');
300 if trimmed.len() > 6 {
301 if let Some(idx) = trimmed.rfind('+') {
302 return (
303 &trimmed[..idx],
304 parse_offset_minutes(&trimmed[idx..]).unwrap_or(0),
305 );
306 }
307 if let Some(idx) = trimmed.rfind('-')
308 && idx > 0
309 {
310 return (
311 &trimmed[..idx],
312 parse_offset_minutes(&trimmed[idx..]).unwrap_or(0),
313 );
314 }
315 }
316 (trimmed, 0)
317}
318
319fn parse_offset_minutes(offset: &str) -> Option<i32> {
320 const MAX_OFFSET_HOURS: i32 = 23;
329 const MAX_OFFSET_MINUTES: i32 = 59;
330 let sign = if offset.starts_with('-') { -1 } else { 1 };
331 let rest = offset.get(1..)?;
332 let (hours_part, minutes_part) = match rest.split_once(':') {
333 Some((hours, minutes)) => (hours, minutes),
334 None => match rest.len() {
335 2 => (rest, "0"), 4 => (&rest[..2], &rest[2..]), _ => return None,
338 },
339 };
340 let hours = hours_part.parse::<i32>().ok()?;
341 let minutes = minutes_part.parse::<i32>().ok()?;
342 if !(0..=MAX_OFFSET_HOURS).contains(&hours) || !(0..=MAX_OFFSET_MINUTES).contains(&minutes) {
343 return None;
344 }
345 Some(sign * (hours * 60 + minutes))
346}
347
348fn days_from_civil(year: i32, month: u32, day: u32) -> i64 {
349 let year = year - i32::from(month <= 2);
350 let era = i64::from(year).div_euclid(400);
351 let yoe = i64::from(year) - era * 400;
352 let month = i64::from(month);
353 let doy = (153 * (month + if month > 2 { -3 } else { 9 }) + 2) / 5 + i64::from(day) - 1;
354 let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy;
355 era * 146_097 + doe - 719_468
356}
357
358fn civil_from_days(days: i64) -> (i32, u32, u32) {
359 let z = days + 719_468;
360 let era = z.div_euclid(146_097);
361 let doe = z - era * 146_097;
362 let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096).div_euclid(365);
363 let year = yoe + era * 400;
364 let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
365 let mp = (5 * doy + 2).div_euclid(153);
366 let day = doy - (153 * mp + 2).div_euclid(5) + 1;
367 let month = mp + if mp < 10 { 3 } else { -9 };
368 let year = year + i64::from(month <= 2);
369 (year as i32, month as u32, day as u32)
370}
371
372fn shell_quote(value: &str) -> String {
373 if value.chars().all(|c| {
380 c.is_ascii_alphanumeric()
381 || matches!(c, '/' | '.' | '_' | '-' | '~' | ':' | '@' | '+' | '=' | '%')
382 }) {
383 value.to_string()
384 } else {
385 format!("'{}'", value.replace('\'', "'\\''"))
386 }
387}
388
389pub fn prepend_changelog(existing: &str, version: &str, date: &str) -> String {
392 const HEADER: &str = "# Changelog\n\n\
393 All notable changes to this project are documented here.\n";
394 let entry = format!("## {version} — {date}\n\n- Released phase via DevFlow.\n");
395
396 if existing.trim().is_empty() {
397 return format!("{HEADER}\n{entry}");
398 }
399 if let Some(idx) = existing.find("\n\n") {
402 let (head, tail) = existing.split_at(idx + 2);
403 format!("{head}{entry}\n{tail}")
404 } else {
405 format!("{entry}\n{existing}")
406 }
407}
408
409#[cfg(test)]
410mod tests {
411 use super::*;
412
413 #[test]
414 fn cron_instructions_save_load_round_trips() {
415 let dir = tempfile::tempdir().unwrap();
416 let record = build_single_agent_cron_instructions(dir.path(), 7, "2026-06-18T15:45:30Z");
417
418 write_cron_instructions(dir.path(), &record).unwrap();
419
420 assert_eq!(load_cron_instructions(dir.path(), 7).unwrap(), record);
421 }
422
423 #[test]
424 fn delete_cron_instructions_is_idempotent() {
425 let dir = tempfile::tempdir().unwrap();
426 let record = build_single_agent_cron_instructions(dir.path(), 7, "2026-06-18T15:45:30Z");
427 write_cron_instructions(dir.path(), &record).unwrap();
428
429 delete_cron_instructions(dir.path(), 7).unwrap();
430 assert!(!cron_instructions_path(dir.path(), 7).exists());
431 delete_cron_instructions(dir.path(), 7).unwrap();
432 }
433
434 #[test]
437 fn cron_instructions_are_per_phase() {
438 let dir = tempfile::tempdir().unwrap();
439 let a = build_single_agent_cron_instructions(dir.path(), 7, "2026-06-18T15:45:30Z");
440 let b = build_single_agent_cron_instructions(dir.path(), 8, "2026-06-18T16:45:30Z");
441 write_cron_instructions(dir.path(), &a).unwrap();
442 write_cron_instructions(dir.path(), &b).unwrap();
443
444 assert_eq!(load_cron_instructions(dir.path(), 7).unwrap(), a);
445 assert_eq!(load_cron_instructions(dir.path(), 8).unwrap(), b);
446 let listed = list_cron_instructions(dir.path());
447 assert_eq!(listed.iter().map(|i| i.phase).collect::<Vec<_>>(), [7, 8]);
448
449 delete_cron_instructions(dir.path(), 7).unwrap();
450 assert!(load_cron_instructions(dir.path(), 7).is_err());
451 assert_eq!(load_cron_instructions(dir.path(), 8).unwrap(), b);
452 }
453
454 #[test]
457 fn legacy_cron_instructions_are_read_and_deleted() {
458 let dir = tempfile::tempdir().unwrap();
459 let record = build_single_agent_cron_instructions(dir.path(), 5, "2026-06-18T15:45:30Z");
460 let legacy = legacy_cron_instructions_path(dir.path());
461 std::fs::create_dir_all(legacy.parent().unwrap()).unwrap();
462 std::fs::write(&legacy, serde_json::to_string_pretty(&record).unwrap()).unwrap();
463
464 assert_eq!(load_cron_instructions(dir.path(), 5).unwrap(), record);
465 assert!(load_cron_instructions(dir.path(), 6).is_err());
466 assert_eq!(list_cron_instructions(dir.path()).len(), 1);
467
468 delete_cron_instructions(dir.path(), 5).unwrap();
469 assert!(!legacy.exists());
470 }
471
472 #[test]
473 fn cron_schedule_rounds_up_to_nearest_minute() {
474 assert_eq!(
475 cron_schedule_from_retry_after("2026-06-18T15:45:30Z"),
476 Some("46 15 18 6 *".to_string())
477 );
478 assert_eq!(
479 cron_schedule_from_retry_after("2026-06-18T15:45:00Z"),
480 Some("45 15 18 6 *".to_string())
481 );
482 }
483
484 #[test]
485 fn cron_schedule_normalizes_negative_offset() {
486 assert_eq!(
488 cron_schedule_from_retry_after("2026-06-18T15:45:30-05:00"),
489 Some("46 20 18 6 *".to_string())
490 );
491 assert_eq!(
493 cron_schedule_from_retry_after("2026-06-18T15:45:00-05:30"),
494 Some("15 21 18 6 *".to_string())
495 );
496 }
497
498 #[test]
505 fn cron_schedule_parses_all_iso8601_offset_forms() {
506 assert_eq!(
508 cron_schedule_from_retry_after("2026-06-18T15:45:30+0530"),
509 cron_schedule_from_retry_after("2026-06-18T15:45:30+05:30"),
510 );
511 assert_eq!(
513 cron_schedule_from_retry_after("2026-06-18T15:45:30-05"),
514 Some("46 20 18 6 *".to_string())
515 );
516 }
517
518 #[test]
519 fn parse_offset_minutes_bounds_and_forms() {
520 assert_eq!(parse_offset_minutes("+05:30"), Some(330));
521 assert_eq!(parse_offset_minutes("+0530"), Some(330));
522 assert_eq!(parse_offset_minutes("-0530"), Some(-330));
523 assert_eq!(parse_offset_minutes("+05"), Some(300));
524 assert_eq!(parse_offset_minutes("-05"), Some(-300));
525 assert_eq!(parse_offset_minutes("+24"), None);
527 assert_eq!(parse_offset_minutes("+05:60"), None);
528 assert_eq!(parse_offset_minutes("+5"), None);
529 assert_eq!(parse_offset_minutes("+530"), None);
530 assert_eq!(parse_offset_minutes("+abcd"), None);
531 }
532
533 #[test]
534 fn cron_schedule_formats_unix_seconds() {
535 assert_eq!(
536 cron_schedule_from_retry_after("1766678401"),
537 Some("1 16 25 12 *".to_string())
538 );
539 }
540
541 #[test]
542 fn shell_quote_leaves_common_safe_chars_unquoted() {
543 assert_eq!(
544 shell_quote("user@host:1.2.3+build"),
545 "user@host:1.2.3+build"
546 );
547 assert_eq!(shell_quote("~/proj/build=1_2%3"), "~/proj/build=1_2%3");
548 }
549
550 #[test]
551 fn shell_quote_quotes_unsafe_input() {
552 assert_eq!(shell_quote("a b"), "'a b'");
553 assert_eq!(shell_quote("it's"), "'it'\\''s'");
554 }
555
556 #[test]
561 fn single_agent_cron_instructions_resume_command_is_devflow_resume() {
562 let dir = tempfile::tempdir().unwrap();
563 let record = build_single_agent_cron_instructions(dir.path(), 9, "2026-06-18T15:45:30Z");
564
565 assert_eq!(record.resume.command, "devflow");
566 assert_eq!(record.resume.args, ["resume", "--phase", "9"]);
567 assert!(
568 record
569 .hermes_cron
570 .command
571 .contains("devflow resume --phase 9")
572 );
573 assert!(!record.hermes_cron.command.contains("sequentagent"));
574 assert!(!record.hermes_cron.command.contains(" start"));
575 assert!(record.hermes_cron.once);
576 }
577
578 #[test]
579 fn cron_instructions_reject_unparseable_retry_time() {
580 let dir = tempfile::tempdir().unwrap();
581 let record = build_single_agent_cron_instructions(dir.path(), 7, "unknown");
582
583 assert_ne!(record.hermes_cron.schedule, "* * * * *");
584 assert!(record.hermes_cron.schedule.is_empty());
585 }
586
587 #[test]
588 fn prepend_changelog_creates_header_when_empty() {
589 let out = prepend_changelog("", "0.5.2", "2026-06-18");
590 assert!(out.starts_with("# Changelog"));
591 assert!(out.contains("## 0.5.2 — 2026-06-18"));
592 }
593
594 #[test]
595 fn prepend_changelog_inserts_after_header() {
596 let existing = "# Changelog\n\n## 0.5.1 — 2026-06-17\n\n- old\n";
597 let out = prepend_changelog(existing, "0.5.2", "2026-06-18");
598 let new_idx = out.find("0.5.2").unwrap();
599 let old_idx = out.find("0.5.1").unwrap();
600 assert!(new_idx < old_idx, "new entry should come before old");
601 assert!(out.starts_with("# Changelog"));
602 }
603}