Skip to main content

loopflow/harness/
opencode_runtime.rs

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/// Outcome of reaping orphaned opencode `serve` processes. Public so the
14/// per-wave runtime can call [`reap_orphaned_opencode_servers`] at startup to
15/// clear servers — and their descendant process trees — left behind by a
16/// crashed `lf wave`.
17#[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(&registry_path(), opencode_pid, std::process::id())
37}
38
39pub(crate) fn unregister_opencode_server(opencode_pid: u32) -> Result<()> {
40    unregister_opencode_server_at_path(&registry_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/// How the leader PID of a registered OpenCode server presents itself when its
103/// owner loopflow process is gone. The OpenCode harness spawns `opencode serve`
104/// in its own process group (`process_group(0)`), so the registered pid is both
105/// the leader and the process-group id; its descendants live in that group.
106#[derive(Debug, Clone, Copy, PartialEq, Eq)]
107enum LeaderState {
108    /// PID no longer exists.
109    Dead,
110    /// PID alive and its command is `opencode ... serve`.
111    Opencode,
112    /// PID alive but running something else — recycled into an unrelated process.
113    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        // Decide whether to reap this entry's process group. The group id is
154        // the registered pid (the harness makes the child its own group
155        // leader), so killing the group reaps the server AND any descendants
156        // it spawned — accurately, without touching unrelated processes.
157        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        // Falling through (Dead+empty group, or Other) prunes the entry: it is
183        // not pushed to `retained`, so the registry drops it.
184    }
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    // SAFETY: signal 0 is an existence/permission probe; no pointers are used.
308    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/// Probe whether a process group still has any live member.
325#[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    // SAFETY: a negative pid targets the process group; signal 0 only probes.
334    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/// Reap an orphaned OpenCode process group: SIGTERM (let the server close its
344/// port and clean up its own children), wait, then SIGKILL the stragglers.
345/// Returns true if the group has no live members afterwards.
346#[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    // SAFETY: kill with a negative pid signals the process group; no pointers.
377    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        // SAFETY: signal 0 probes the process group; no pointers are used.
394        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        // Owner 1 alive (retained); owner 2 dead, leader is a live opencode server.
458        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        // Owner-alive entry retained; the dead leader at pid 12 pruned (no group).
492        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        // Owner dead, leader pid gone, but children still live in its group.
503        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        // Owner dead, but the pid now runs an unrelated process.
564        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        // Not reaped (would risk killing an unrelated group), entry pruned.
580        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}