1use std::collections::HashSet;
2use std::path::{Path, PathBuf};
3use std::process::Command;
4use std::time::{Duration, Instant};
5
6use anyhow::{Context, Result};
7use fs2::FileExt;
8use serde::{Deserialize, Serialize};
9
10const OPENCODE_REGISTRY_FILE: &str = "runtime/opencode-servers.json";
11const REAP_TERM_GRACE: Duration = Duration::from_secs(2);
12
13#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
18pub struct OpenCodeReapReport {
19 pub reaped: u32,
20 pub errors: u32,
21}
22
23#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
24pub(crate) struct OpenCodeServerEntry {
25 pub opencode_pid: u32,
26 pub owner_loopflow_pid: u32,
27}
28
29pub(crate) fn registered_opencode_servers_at(lf_home: &Path) -> Result<Vec<OpenCodeServerEntry>> {
30 let path = lf_home.join(OPENCODE_REGISTRY_FILE);
31 let _lock = lock_registry_for_read(&path)?;
32 read_registry_entries(&path)
33}
34
35pub(crate) fn register_opencode_server(opencode_pid: u32) -> Result<()> {
36 register_opencode_server_at_path(®istry_path(), opencode_pid, std::process::id())
37}
38
39pub(crate) fn unregister_opencode_server(opencode_pid: u32) -> Result<()> {
40 unregister_opencode_server_at_path(®istry_path(), opencode_pid)
41}
42
43pub fn reap_orphaned_opencode_servers() -> OpenCodeReapReport {
44 reap_orphaned_opencode_servers_at(&crate::store::lf_home_dir())
45}
46
47pub(crate) fn reap_orphaned_opencode_servers_at(lf_home: &Path) -> OpenCodeReapReport {
48 reap_orphaned_opencode_servers_at_path(
49 &lf_home.join(OPENCODE_REGISTRY_FILE),
50 |_| true,
51 pid_is_alive,
52 classify_leader,
53 process_group_alive,
54 terminate_process_group,
55 )
56}
57
58pub(crate) fn reap_selected_orphaned_opencode_servers_at(
59 lf_home: &Path,
60 process_groups: &HashSet<u32>,
61) -> OpenCodeReapReport {
62 reap_orphaned_opencode_servers_at_path(
63 &lf_home.join(OPENCODE_REGISTRY_FILE),
64 |pid| process_groups.contains(&pid),
65 pid_is_alive,
66 classify_leader,
67 process_group_alive,
68 terminate_process_group,
69 )
70}
71
72fn registry_path() -> PathBuf {
73 crate::store::lf_home_dir().join(OPENCODE_REGISTRY_FILE)
74}
75
76fn register_opencode_server_at_path(
77 path: &Path,
78 opencode_pid: u32,
79 owner_loopflow_pid: u32,
80) -> Result<()> {
81 let _lock = lock_registry(path)?;
82 let mut entries = read_registry_entries(path)?;
83 entries.retain(|entry| entry.opencode_pid != opencode_pid);
84 entries.push(OpenCodeServerEntry {
85 opencode_pid,
86 owner_loopflow_pid,
87 });
88 write_registry_entries(path, &entries)
89}
90
91fn unregister_opencode_server_at_path(path: &Path, opencode_pid: u32) -> Result<()> {
92 let _lock = lock_registry(path)?;
93 let mut entries = read_registry_entries(path)?;
94 let original_len = entries.len();
95 entries.retain(|entry| entry.opencode_pid != opencode_pid);
96 if entries.len() == original_len {
97 return Ok(());
98 }
99 write_registry_entries(path, &entries)
100}
101
102#[derive(Debug, Clone, Copy, PartialEq, Eq)]
107enum LeaderState {
108 Dead,
110 Opencode,
112 Other,
114}
115
116fn reap_orphaned_opencode_servers_at_path(
117 path: &Path,
118 eligible: impl Fn(u32) -> bool,
119 owner_pid_alive: impl Fn(u32) -> bool,
120 leader: impl Fn(u32) -> LeaderState,
121 group_alive: impl Fn(u32) -> bool,
122 terminate_group: impl Fn(u32) -> bool,
123) -> OpenCodeReapReport {
124 let mut report = OpenCodeReapReport::default();
125 let _lock = match lock_registry(path) {
126 Ok(lock) => lock,
127 Err(err) => {
128 tracing::warn!(path = %path.display(), error = %err, "failed to lock OpenCode registry");
129 report.errors += 1;
130 return report;
131 }
132 };
133 let entries = match read_registry_entries(path) {
134 Ok(entries) => entries,
135 Err(err) => {
136 tracing::warn!(path = %path.display(), error = %err, "failed to read OpenCode registry");
137 report.errors += 1;
138 return report;
139 }
140 };
141
142 let mut retained = Vec::with_capacity(entries.len());
143 for entry in entries {
144 if !eligible(entry.opencode_pid) {
145 retained.push(entry);
146 continue;
147 }
148 if owner_pid_alive(entry.owner_loopflow_pid) {
149 retained.push(entry);
150 continue;
151 }
152
153 let reap_group = match leader(entry.opencode_pid) {
158 LeaderState::Opencode => true,
159 LeaderState::Dead => group_alive(entry.opencode_pid),
160 LeaderState::Other => {
161 tracing::info!(
162 opencode_pid = entry.opencode_pid,
163 "orphaned OpenCode pid reused by an unrelated process; leaving it"
164 );
165 false
166 }
167 };
168
169 if reap_group {
170 if terminate_group(entry.opencode_pid) {
171 report.reaped += 1;
172 } else {
173 tracing::warn!(
174 opencode_pid = entry.opencode_pid,
175 owner_loopflow_pid = entry.owner_loopflow_pid,
176 "failed to terminate orphaned OpenCode process group"
177 );
178 report.errors += 1;
179 retained.push(entry);
180 }
181 }
182 }
185
186 if let Err(err) = write_registry_entries(path, &retained) {
187 tracing::warn!(
188 path = %path.display(),
189 error = %err,
190 "failed to update OpenCode registry after orphan cleanup"
191 );
192 report.errors += 1;
193 }
194
195 report
196}
197
198fn lock_registry(path: &Path) -> Result<std::fs::File> {
199 let lock_path = path.with_extension("json.lock");
200 if let Some(parent) = lock_path.parent() {
201 std::fs::create_dir_all(parent)
202 .with_context(|| format!("failed creating runtime dir {}", parent.display()))?;
203 }
204 let lock = std::fs::OpenOptions::new()
205 .create(true)
206 .read(true)
207 .write(true)
208 .truncate(false)
209 .open(&lock_path)
210 .with_context(|| {
211 format!(
212 "failed opening OpenCode registry lock {}",
213 lock_path.display()
214 )
215 })?;
216 FileExt::lock_exclusive(&lock)
217 .with_context(|| format!("failed locking OpenCode registry {}", lock_path.display()))?;
218 Ok(lock)
219}
220
221fn lock_registry_for_read(path: &Path) -> Result<Option<std::fs::File>> {
222 let lock_path = path.with_extension("json.lock");
223 let lock = match std::fs::OpenOptions::new().read(true).open(&lock_path) {
224 Ok(lock) => lock,
225 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
226 Err(error) => {
227 return Err(error).with_context(|| {
228 format!(
229 "failed opening OpenCode registry lock {}",
230 lock_path.display()
231 )
232 })
233 }
234 };
235 FileExt::lock_shared(&lock)
236 .with_context(|| format!("failed locking OpenCode registry {}", lock_path.display()))?;
237 Ok(Some(lock))
238}
239
240fn read_registry_entries(path: &Path) -> Result<Vec<OpenCodeServerEntry>> {
241 let content = match std::fs::read_to_string(path) {
242 Ok(content) => content,
243 Err(err) if err.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
244 Err(err) => return Err(err.into()),
245 };
246
247 if content.trim().is_empty() {
248 return Ok(Vec::new());
249 }
250
251 serde_json::from_str(&content)
252 .with_context(|| format!("failed parsing OpenCode registry at {}", path.display()))
253}
254
255fn write_registry_entries(path: &Path, entries: &[OpenCodeServerEntry]) -> Result<()> {
256 if let Some(parent) = path.parent() {
257 std::fs::create_dir_all(parent)
258 .with_context(|| format!("failed creating runtime dir {}", parent.display()))?;
259 }
260
261 let json = serde_json::to_string_pretty(entries)
262 .context("failed serializing OpenCode server registry")?;
263 std::fs::write(path, json)
264 .with_context(|| format!("failed writing OpenCode registry at {}", path.display()))?;
265 Ok(())
266}
267
268fn classify_leader(pid: u32) -> LeaderState {
269 if !pid_is_alive(pid) {
270 return LeaderState::Dead;
271 }
272 if process_looks_like_opencode_serve(pid) {
273 LeaderState::Opencode
274 } else {
275 LeaderState::Other
276 }
277}
278
279fn process_looks_like_opencode_serve(pid: u32) -> bool {
280 let output = match Command::new("ps")
281 .arg("-o")
282 .arg("command=")
283 .arg("-p")
284 .arg(pid.to_string())
285 .output()
286 {
287 Ok(output) => output,
288 Err(_) => return false,
289 };
290
291 if !output.status.success() {
292 return false;
293 }
294
295 let command = String::from_utf8_lossy(&output.stdout).to_ascii_lowercase();
296 command.contains("opencode") && command.contains("serve")
297}
298
299#[cfg(unix)]
300fn pid_is_alive(pid: u32) -> bool {
301 if pid == 0 {
302 return false;
303 }
304 let Ok(raw) = i32::try_from(pid) else {
305 return false;
306 };
307 let result = unsafe { libc::kill(raw, 0) };
309 result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM)
310}
311
312#[cfg(not(unix))]
313fn pid_is_alive(pid: u32) -> bool {
314 if pid == 0 {
315 return false;
316 }
317 Command::new("kill")
318 .arg("-0")
319 .arg(pid.to_string())
320 .status()
321 .is_ok_and(|status| status.success())
322}
323
324#[cfg(unix)]
326fn process_group_alive(pgid: u32) -> bool {
327 if pgid == 0 {
328 return false;
329 }
330 let Ok(raw) = i32::try_from(pgid) else {
331 return false;
332 };
333 let result = unsafe { libc::kill(-raw, 0) };
335 result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM)
336}
337
338#[cfg(not(unix))]
339fn process_group_alive(_pgid: u32) -> bool {
340 false
341}
342
343#[cfg(unix)]
347fn terminate_process_group(pgid: u32) -> bool {
348 if pgid == 0 {
349 return false;
350 }
351 if crate::engine::process::current_process_group_id() == Some(pgid) {
352 tracing::warn!(pgid, "refusing to reap current process group");
353 return false;
354 }
355 let Ok(raw) = i32::try_from(pgid) else {
356 return false;
357 };
358
359 if signal_group(raw, libc::SIGTERM) {
360 return true;
361 }
362 if group_gone(raw, REAP_TERM_GRACE) {
363 return true;
364 }
365 signal_group(raw, libc::SIGKILL);
366 group_gone(raw, REAP_TERM_GRACE)
367}
368
369#[cfg(not(unix))]
370fn terminate_process_group(_pgid: u32) -> bool {
371 false
372}
373
374#[cfg(unix)]
375fn signal_group(pgid: i32, signal: libc::c_int) -> bool {
376 let result = unsafe { libc::kill(-pgid, signal) };
378 if result == 0 {
379 return false;
380 }
381 let err = std::io::Error::last_os_error();
382 if err.raw_os_error() == Some(libc::ESRCH) {
383 return true;
384 }
385 tracing::warn!(pgid, signal, error = %err, "failed to signal OpenCode process group");
386 false
387}
388
389#[cfg(unix)]
390fn group_gone(pgid: i32, grace: Duration) -> bool {
391 let deadline = Instant::now() + grace;
392 loop {
393 let result = unsafe { libc::kill(-pgid, 0) };
395 if result != 0 && std::io::Error::last_os_error().raw_os_error() == Some(libc::ESRCH) {
396 return true;
397 }
398 if Instant::now() >= deadline {
399 return false;
400 }
401 std::thread::sleep(Duration::from_millis(50));
402 }
403}
404
405#[cfg(test)]
406mod tests {
407 use std::collections::HashSet;
408 use std::sync::Mutex;
409
410 use tempfile::tempdir;
411
412 use super::*;
413
414 fn registry_path(root: &Path) -> PathBuf {
415 root.join("runtime").join("opencode-servers.json")
416 }
417
418 fn entry(opencode_pid: u32, owner_loopflow_pid: u32) -> OpenCodeServerEntry {
419 OpenCodeServerEntry {
420 opencode_pid,
421 owner_loopflow_pid,
422 }
423 }
424
425 #[test]
426 fn register_and_unregister_opencode_server_updates_registry() {
427 let tmp = tempdir().expect("tempdir");
428 let path = registry_path(tmp.path());
429
430 register_opencode_server_at_path(&path, 111, 222).expect("register pid");
431 let entries = read_registry_entries(&path).expect("read entries");
432 assert_eq!(entries, vec![entry(111, 222)]);
433
434 register_opencode_server_at_path(&path, 111, 444).expect("overwrite existing pid");
435 let entries = read_registry_entries(&path).expect("read entries");
436 assert_eq!(entries, vec![entry(111, 444)]);
437
438 unregister_opencode_server_at_path(&path, 111).expect("unregister pid");
439 let entries = read_registry_entries(&path).expect("read entries");
440 assert!(entries.is_empty());
441 }
442
443 #[test]
444 fn reading_an_absent_registry_creates_no_runtime_state() {
445 let tmp = tempdir().expect("tempdir");
446
447 assert!(registered_opencode_servers_at(tmp.path())
448 .expect("read absent registry")
449 .is_empty());
450 assert!(!tmp.path().join("runtime").exists());
451 }
452
453 #[test]
454 fn reap_kills_the_process_group_of_an_orphaned_opencode_server() {
455 let tmp = tempdir().expect("tempdir");
456 let path = registry_path(tmp.path());
457 write_registry_entries(&path, &[entry(10, 1), entry(11, 2), entry(12, 2)])
459 .expect("write registry");
460
461 let owner_alive: HashSet<u32> = [1].into_iter().collect();
462 let opencode_pids: HashSet<u32> = [11].into_iter().collect();
463 let killed = Mutex::new(Vec::new());
464
465 let report = reap_orphaned_opencode_servers_at_path(
466 &path,
467 |_| true,
468 |pid| owner_alive.contains(&pid),
469 |pid| {
470 if opencode_pids.contains(&pid) {
471 LeaderState::Opencode
472 } else {
473 LeaderState::Dead
474 }
475 },
476 |_| false,
477 |pid| {
478 killed.lock().expect("lock killed list").push(pid);
479 true
480 },
481 );
482
483 assert_eq!(
484 report,
485 OpenCodeReapReport {
486 reaped: 1,
487 errors: 0
488 }
489 );
490 assert_eq!(*killed.lock().expect("lock killed list"), vec![11]);
491 assert_eq!(
493 read_registry_entries(&path).expect("read entries"),
494 vec![entry(10, 1)]
495 );
496 }
497
498 #[test]
499 fn reap_kills_surviving_children_when_the_leader_is_dead() {
500 let tmp = tempdir().expect("tempdir");
501 let path = registry_path(tmp.path());
502 write_registry_entries(&path, &[entry(21, 2)]).expect("write registry");
504
505 let alive_groups: HashSet<u32> = [21].into_iter().collect();
506 let killed = Mutex::new(Vec::new());
507
508 let report = reap_orphaned_opencode_servers_at_path(
509 &path,
510 |_| true,
511 |_| false,
512 |_| LeaderState::Dead,
513 |pgid| alive_groups.contains(&pgid),
514 |pid| {
515 killed.lock().expect("lock killed list").push(pid);
516 true
517 },
518 );
519
520 assert_eq!(
521 report,
522 OpenCodeReapReport {
523 reaped: 1,
524 errors: 0
525 }
526 );
527 assert_eq!(*killed.lock().expect("lock killed list"), vec![21]);
528 assert!(read_registry_entries(&path)
529 .expect("read entries")
530 .is_empty());
531 }
532
533 #[test]
534 fn reap_prunes_a_fully_dead_tree_without_signalling() {
535 let tmp = tempdir().expect("tempdir");
536 let path = registry_path(tmp.path());
537 write_registry_entries(&path, &[entry(31, 2)]).expect("write registry");
538
539 let killed = Mutex::new(Vec::new());
540 let report = reap_orphaned_opencode_servers_at_path(
541 &path,
542 |_| true,
543 |_| false,
544 |_| LeaderState::Dead,
545 |_| false,
546 |pid| {
547 killed.lock().expect("lock killed list").push(pid);
548 true
549 },
550 );
551
552 assert_eq!(report, OpenCodeReapReport::default());
553 assert!(killed.lock().expect("lock killed list").is_empty());
554 assert!(read_registry_entries(&path)
555 .expect("read entries")
556 .is_empty());
557 }
558
559 #[test]
560 fn reap_leaves_a_reused_pid_alone() {
561 let tmp = tempdir().expect("tempdir");
562 let path = registry_path(tmp.path());
563 write_registry_entries(&path, &[entry(41, 2)]).expect("write registry");
565
566 let killed = Mutex::new(Vec::new());
567 let report = reap_orphaned_opencode_servers_at_path(
568 &path,
569 |_| true,
570 |_| false,
571 |_| LeaderState::Other,
572 |_| true,
573 |pid| {
574 killed.lock().expect("lock killed list").push(pid);
575 true
576 },
577 );
578
579 assert_eq!(report, OpenCodeReapReport::default());
581 assert!(killed.lock().expect("lock killed list").is_empty());
582 assert!(read_registry_entries(&path)
583 .expect("read entries")
584 .is_empty());
585 }
586
587 #[test]
588 fn reap_retains_an_entry_when_termination_fails() {
589 let tmp = tempdir().expect("tempdir");
590 let path = registry_path(tmp.path());
591 write_registry_entries(&path, &[entry(51, 2)]).expect("write registry");
592
593 let report = reap_orphaned_opencode_servers_at_path(
594 &path,
595 |_| true,
596 |_| false,
597 |_| LeaderState::Opencode,
598 |_| true,
599 |_| false,
600 );
601
602 assert_eq!(
603 report,
604 OpenCodeReapReport {
605 reaped: 0,
606 errors: 1
607 }
608 );
609 assert_eq!(
610 read_registry_entries(&path).expect("read entries"),
611 vec![entry(51, 2)]
612 );
613 }
614
615 #[test]
616 fn selected_reap_preserves_unlisted_orphans() {
617 let tmp = tempdir().expect("tempdir");
618 let path = registry_path(tmp.path());
619 write_registry_entries(&path, &[entry(60, 2), entry(61, 2)]).expect("write registry");
620
621 let report = reap_orphaned_opencode_servers_at_path(
622 &path,
623 |pid| pid == 60,
624 |_| false,
625 |_| LeaderState::Opencode,
626 |_| true,
627 |_| true,
628 );
629
630 assert_eq!(report.reaped, 1);
631 assert_eq!(
632 read_registry_entries(&path).expect("read entries"),
633 vec![entry(61, 2)]
634 );
635 }
636
637 #[test]
638 fn reap_is_idempotent() {
639 let tmp = tempdir().expect("tempdir");
640 let path = registry_path(tmp.path());
641 write_registry_entries(&path, &[entry(20, 2)]).expect("write registry");
642
643 let first = reap_orphaned_opencode_servers_at_path(
644 &path,
645 |_| true,
646 |_| false,
647 |pid| {
648 if pid == 20 {
649 LeaderState::Opencode
650 } else {
651 LeaderState::Dead
652 }
653 },
654 |_| false,
655 |_| true,
656 );
657 assert_eq!(
658 first,
659 OpenCodeReapReport {
660 reaped: 1,
661 errors: 0
662 }
663 );
664
665 let second = reap_orphaned_opencode_servers_at_path(
666 &path,
667 |_| true,
668 |_| false,
669 |pid| {
670 if pid == 20 {
671 LeaderState::Opencode
672 } else {
673 LeaderState::Dead
674 }
675 },
676 |_| false,
677 |_| true,
678 );
679 assert_eq!(second, OpenCodeReapReport::default());
680 }
681}