Skip to main content

objects/store/
writer_lease.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Exclusive writers for checkouts; separate checkouts may share a Thread.
3
4use std::path::{Path, PathBuf};
5
6use chrono::{DateTime, Utc};
7use serde::{Deserialize, Serialize};
8
9use crate::{
10    fs_atomic::write_file_atomic,
11    lock::RepoLock,
12    object::ContentHash,
13    store::{HeddleError, Liveness, Result, reservation_liveness_at},
14};
15
16/// The physical checkout lock shared by capture and lease handoff.
17pub fn checkout_writer_lock(heddle_dir: &Path, root: &Path) -> Result<RepoLock> {
18    Ok(checkout_writer_lock_for_path(
19        heddle_dir,
20        &root.canonicalize()?,
21    ))
22}
23
24fn checkout_writer_lock_for_path(heddle_dir: &Path, path: &Path) -> RepoLock {
25    let path_key = ContentHash::compute_typed(
26        "checkout-writer-path-v2",
27        path.as_os_str().as_encoded_bytes(),
28    );
29    RepoLock::at(
30        heddle_dir
31            .join("locks")
32            .join(format!("checkout-{}.lock", path_key.to_hex())),
33    )
34}
35
36#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
37#[serde(rename_all = "snake_case")]
38pub enum WriterLeaseStatus {
39    Active,
40    Complete,
41    Abandoned,
42}
43
44#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
45pub struct PidNamespace {
46    pub device: u64,
47    pub inode: u64,
48}
49
50impl std::fmt::Display for WriterLeaseStatus {
51    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
52        match self {
53            Self::Active => write!(f, "active"),
54            Self::Complete => write!(f, "complete"),
55            Self::Abandoned => write!(f, "abandoned"),
56        }
57    }
58}
59
60#[derive(Debug, Clone, Serialize, Deserialize)]
61pub struct WriterLease {
62    pub lease_id: String,
63    pub thread: String,
64    #[serde(default)]
65    pub actor_session_id: Option<String>,
66    #[serde(default)]
67    pub task_assignment_id: Option<String>,
68    #[serde(default)]
69    pub anchor_state: Option<String>,
70    #[serde(default)]
71    pub anchor_root: Option<String>,
72    #[serde(default)]
73    pub path: Option<PathBuf>,
74    pub token_hash: String,
75    #[serde(default)]
76    pub pid: Option<u32>,
77    #[serde(default)]
78    pub boot_id: Option<String>,
79    #[serde(default)]
80    pub pid_birth: Option<String>,
81    #[serde(default)]
82    pub pid_namespace: Option<PidNamespace>,
83    #[serde(default)]
84    pub harness_session_id: Option<String>,
85    pub heartbeat_at: DateTime<Utc>,
86    pub started_at: DateTime<Utc>,
87    pub status: WriterLeaseStatus,
88    #[serde(default)]
89    pub completed_at: Option<DateTime<Utc>>,
90}
91
92impl WriterLease {
93    pub fn from_draft(
94        draft: WriterLeaseDraft,
95        lease_id: String,
96        token_hash: String,
97        now: DateTime<Utc>,
98    ) -> Self {
99        Self {
100            lease_id,
101            thread: draft.thread,
102            actor_session_id: draft.actor_session_id,
103            task_assignment_id: draft.task_assignment_id,
104            anchor_state: draft.anchor_state,
105            anchor_root: draft.anchor_root,
106            path: draft.path,
107            token_hash,
108            pid: draft.pid,
109            boot_id: draft.boot_id,
110            pid_birth: None,
111            pid_namespace: draft.pid.and_then(|_| current_pid_namespace()),
112            harness_session_id: None,
113            heartbeat_at: now,
114            started_at: now,
115            status: WriterLeaseStatus::Active,
116            completed_at: None,
117        }
118    }
119
120    pub fn matches_current_pid_namespace(&self) -> bool {
121        if cfg!(target_os = "linux") {
122            self.pid_namespace
123                .is_some_and(|recorded| Some(recorded) == current_pid_namespace())
124        } else {
125            true
126        }
127    }
128
129    fn conflicts_with(&self, thread: &str, path: Option<&Path>) -> bool {
130        self.status == WriterLeaseStatus::Active
131            && match (self.path.as_deref(), path) {
132                (Some(owned), Some(requested)) => owned == requested,
133                // An unmaterialized reservation has no narrower checkout scope.
134                _ => self.thread == thread,
135            }
136    }
137
138    pub fn lease_expires_at(&self) -> DateTime<Utc> {
139        self.heartbeat_at + crate::store::AGENT_LEASE_DURATION
140    }
141
142    pub fn liveness_at(&self, now: DateTime<Utc>) -> Liveness {
143        self.liveness_at_in_namespace(now, current_pid_namespace())
144    }
145
146    fn liveness_at_in_namespace(
147        &self,
148        now: DateTime<Utc>,
149        observed_namespace: Option<PidNamespace>,
150    ) -> Liveness {
151        if self.status != WriterLeaseStatus::Active {
152            return Liveness::Dead;
153        }
154        // A PID in another namespace identifies a different process here.
155        // Without comparable namespace identities, only the heartbeat can expire it.
156        let same_namespace = if cfg!(target_os = "linux") {
157            self.pid_namespace
158                .is_some_and(|recorded| Some(recorded) == observed_namespace)
159        } else {
160            true
161        };
162        if !same_namespace {
163            return reservation_liveness_at(None, None, Some(self.heartbeat_at), now);
164        }
165        if let (Some(pid), Some(birth)) = (self.pid, self.pid_birth.as_deref())
166            && super::process_birth(pid).as_deref() != Some(birth)
167        {
168            return Liveness::Dead;
169        }
170        reservation_liveness_at(
171            self.pid,
172            self.boot_id.as_deref(),
173            Some(self.heartbeat_at),
174            now,
175        )
176    }
177}
178
179#[cfg(target_os = "linux")]
180fn current_pid_namespace() -> Option<PidNamespace> {
181    use std::os::unix::fs::MetadataExt;
182
183    let metadata = std::fs::metadata("/proc/self/ns/pid").ok()?;
184    Some(PidNamespace {
185        device: metadata.dev(),
186        inode: metadata.ino(),
187    })
188}
189
190#[cfg(not(target_os = "linux"))]
191fn current_pid_namespace() -> Option<PidNamespace> {
192    None
193}
194
195#[derive(Debug, Clone)]
196pub struct WriterLeaseDraft {
197    pub thread: String,
198    pub actor_session_id: Option<String>,
199    pub task_assignment_id: Option<String>,
200    pub anchor_state: Option<String>,
201    pub anchor_root: Option<String>,
202    pub path: Option<PathBuf>,
203    pub pid: Option<u32>,
204    pub boot_id: Option<String>,
205}
206
207#[derive(Debug)]
208pub struct WriterLeaseGrant {
209    pub lease: WriterLease,
210    pub token: String,
211}
212
213#[derive(Debug)]
214pub enum WriterLeaseReserveOutcome {
215    Reserved(WriterLeaseGrant),
216    LiveOwner(WriterLease),
217}
218
219#[derive(Debug)]
220pub enum WriterLeaseAuthOutcome {
221    Authorized(WriterLease),
222    Missing,
223    TokenMismatch,
224    Inactive(WriterLease),
225}
226
227pub struct WriterLeaseStore {
228    leases_dir: PathBuf,
229}
230
231impl WriterLeaseStore {
232    pub fn new(heddle_dir: &Path) -> Self {
233        Self {
234            leases_dir: heddle_dir.join("writer-leases"),
235        }
236    }
237
238    fn lock_path(&self) -> PathBuf {
239        self.leases_dir.join(".lock")
240    }
241
242    fn checkout_lock(&self, path: &Path) -> Result<RepoLock> {
243        let heddle_dir = self
244            .leases_dir
245            .parent()
246            .ok_or_else(|| HeddleError::Config("writer lease directory has no parent".into()))?;
247        Ok(checkout_writer_lock_for_path(heddle_dir, path))
248    }
249
250    fn write_lock(&self) -> Result<crate::lock::WriteLockGuard> {
251        RepoLock::at(self.lock_path()).write().map_err(|err| {
252            HeddleError::Config(format!("failed to acquire writer lease lock: {err}"))
253        })
254    }
255
256    fn lease_path(&self, lease_id: &str) -> Result<PathBuf> {
257        validate_lease_id(lease_id)?;
258        Ok(self.leases_dir.join(format!("{lease_id}.toml")))
259    }
260
261    fn load_path(&self, path: &Path) -> Result<Option<WriterLease>> {
262        if !path.exists() {
263            return Ok(None);
264        }
265        let content = std::fs::read_to_string(path)?;
266        toml::from_str(&content)
267            .map(Some)
268            .map_err(|err| HeddleError::Config(err.to_string()))
269    }
270
271    fn write_lease(&self, lease: &WriterLease) -> Result<()> {
272        crate::fs_atomic::create_dir_all_durable(&self.leases_dir)?;
273        let content =
274            toml::to_string_pretty(lease).map_err(|err| HeddleError::Config(err.to_string()))?;
275        Ok(write_file_atomic(
276            &self.lease_path(&lease.lease_id)?,
277            content.as_bytes(),
278        )?)
279    }
280
281    fn list_locked(&self) -> Result<Vec<WriterLease>> {
282        if !self.leases_dir.exists() {
283            return Ok(Vec::new());
284        }
285        let mut leases = Vec::new();
286        for entry in std::fs::read_dir(&self.leases_dir)? {
287            let path = entry?.path();
288            if path
289                .extension()
290                .is_some_and(|extension| extension == "toml")
291                && let Some(lease) = self.load_path(&path)?
292            {
293                leases.push(lease);
294            }
295        }
296        leases.sort_by_key(|lease| std::cmp::Reverse(lease.started_at));
297        Ok(leases)
298    }
299
300    fn reap_expired_locked(&self, now: DateTime<Utc>, held_checkout: Option<&Path>) -> Result<()> {
301        for mut lease in self.list_locked()? {
302            if lease.liveness_at(now) == Liveness::Dead && lease.status == WriterLeaseStatus::Active
303            {
304                // Never retire authority while a capture holds this checkout.
305                // try_write avoids reversing the checkout -> lease-store lock
306                // order used by capture and handoff.
307                let checkout_guard = if let Some(path) = lease.path.as_deref() {
308                    if held_checkout == Some(path) {
309                        None
310                    } else {
311                        let lock = self.checkout_lock(path)?;
312                        let Some(guard) = lock.try_write().map_err(|error| {
313                            HeddleError::Config(format!(
314                                "failed to acquire checkout writer lock: {error}"
315                            ))
316                        })?
317                        else {
318                            continue;
319                        };
320                        Some(guard)
321                    }
322                } else {
323                    None
324                };
325                lease.status = WriterLeaseStatus::Abandoned;
326                lease.completed_at = Some(now);
327                self.write_lease(&lease)?;
328                drop(checkout_guard);
329            }
330        }
331        Ok(())
332    }
333
334    pub fn reserve(
335        &self,
336        draft: WriterLeaseDraft,
337        now: DateTime<Utc>,
338    ) -> Result<WriterLeaseReserveOutcome> {
339        self.reserve_prepared(
340            draft,
341            generate_writer_lease_id(),
342            generate_writer_lease_token(),
343            now,
344        )
345    }
346
347    /// The caller already holds the checkout mutation lock for `draft.path`.
348    pub fn reserve_with_checkout_lock(
349        &self,
350        draft: WriterLeaseDraft,
351        now: DateTime<Utc>,
352    ) -> Result<WriterLeaseReserveOutcome> {
353        self.reserve_prepared_inner(
354            draft,
355            generate_writer_lease_id(),
356            generate_writer_lease_token(),
357            now,
358            true,
359        )
360    }
361
362    /// A command journal persists these random credentials before reservation.
363    /// Retrying the exact live reservation recovers its token; a different actor,
364    /// path, token or an expired/released reservation can never revive it.
365    pub fn reserve_prepared(
366        &self,
367        draft: WriterLeaseDraft,
368        lease_id: String,
369        token: String,
370        now: DateTime<Utc>,
371    ) -> Result<WriterLeaseReserveOutcome> {
372        self.reserve_prepared_inner(draft, lease_id, token, now, false)
373    }
374
375    fn reserve_prepared_inner(
376        &self,
377        mut draft: WriterLeaseDraft,
378        lease_id: String,
379        token: String,
380        now: DateTime<Utc>,
381        checkout_lock_held: bool,
382    ) -> Result<WriterLeaseReserveOutcome> {
383        validate_lease_id(&lease_id)?;
384        if token.len() < 32 || token.len() > 256 {
385            return Err(HeddleError::Config("invalid prepared writer token".into()));
386        }
387        draft.path = draft.path.map(std::fs::canonicalize).transpose()?;
388        let _checkout_guard = if !checkout_lock_held {
389            draft
390                .path
391                .as_deref()
392                .map(|path| {
393                    self.checkout_lock(path)?.write().map_err(|error| {
394                        HeddleError::Config(format!(
395                            "failed to acquire checkout writer lock: {error}"
396                        ))
397                    })
398                })
399                .transpose()?
400        } else {
401            None
402        };
403        let _lock = self.write_lock()?;
404        self.reap_expired_locked(now, draft.path.as_deref())?;
405        if let Some(old) = self.load_path(&self.lease_path(&lease_id)?)? {
406            if old.status == WriterLeaseStatus::Active
407                && old.liveness_at(now) != Liveness::Dead
408                && old.thread == draft.thread
409                && old.path == draft.path
410                && old.actor_session_id == draft.actor_session_id
411                && old.token_hash == token_hash(&token)
412            {
413                return Ok(WriterLeaseReserveOutcome::Reserved(WriterLeaseGrant {
414                    lease: old,
415                    token,
416                }));
417            }
418            return Err(HeddleError::Config(
419                "prepared writer reservation changed or expired".into(),
420            ));
421        }
422        if let Some(owner) = self
423            .list_locked()?
424            .into_iter()
425            .find(|lease| lease.conflicts_with(&draft.thread, draft.path.as_deref()))
426        {
427            return Ok(WriterLeaseReserveOutcome::LiveOwner(owner));
428        }
429        let lease = WriterLease::from_draft(draft, lease_id, token_hash(&token), now);
430        self.write_lease(&lease)?;
431        Ok(WriterLeaseReserveOutcome::Reserved(WriterLeaseGrant {
432            lease,
433            token,
434        }))
435    }
436
437    /// Advisory preflight. `reserve` repeats this check under the write lock.
438    pub fn live_owner(&self, thread: &str, path: Option<&Path>) -> Result<Option<WriterLease>> {
439        self.live_owner_inner(thread, path, None)
440    }
441
442    /// Advisory preflight when the caller holds this checkout's mutation lock.
443    pub fn live_owner_with_checkout_lock(
444        &self,
445        thread: &str,
446        path: &Path,
447        _checkout_guard: &crate::lock::WriteLockGuard,
448    ) -> Result<Option<WriterLease>> {
449        self.live_owner_inner(thread, Some(path), Some(path))
450    }
451
452    fn live_owner_inner(
453        &self,
454        thread: &str,
455        path: Option<&Path>,
456        held_checkout: Option<&Path>,
457    ) -> Result<Option<WriterLease>> {
458        let path = path.map(std::fs::canonicalize).transpose()?;
459        let held_checkout = held_checkout.map(std::fs::canonicalize).transpose()?;
460        let _lock = self.write_lock()?;
461        self.reap_expired_locked(Utc::now(), held_checkout.as_deref())?;
462        Ok(self
463            .list_locked()?
464            .into_iter()
465            .find(|lease| lease.conflicts_with(thread, path.as_deref())))
466    }
467
468    pub fn authenticate_and_renew(
469        &self,
470        lease_id: &str,
471        token: &str,
472        now: DateTime<Utc>,
473    ) -> Result<WriterLeaseAuthOutcome> {
474        let _lock = self.write_lock()?;
475        let path = self.lease_path(lease_id)?;
476        let Some(mut lease) = self.load_path(&path)? else {
477            return Ok(WriterLeaseAuthOutcome::Missing);
478        };
479        if lease.status != WriterLeaseStatus::Active || lease.liveness_at(now) == Liveness::Dead {
480            return Ok(WriterLeaseAuthOutcome::Inactive(lease));
481        }
482        if token_hash(token) != lease.token_hash {
483            return Ok(WriterLeaseAuthOutcome::TokenMismatch);
484        }
485        lease.heartbeat_at = now;
486        self.write_lease(&lease)?;
487        Ok(WriterLeaseAuthOutcome::Authorized(lease))
488    }
489
490    /// A hook proves both possession of the lane credential and its process
491    /// ancestry. Subsequent hook calls may only renew the same live session.
492    pub fn bind_hook_session(
493        &self,
494        lease_id: &str,
495        token: &str,
496        session: &str,
497        pid: u32,
498        birth: &str,
499        now: DateTime<Utc>,
500    ) -> Result<WriterLeaseAuthOutcome> {
501        let _lock = self.write_lock()?;
502        let path = self.lease_path(lease_id)?;
503        let Some(mut lease) = self.load_path(&path)? else {
504            return Ok(WriterLeaseAuthOutcome::Missing);
505        };
506        if lease.status != WriterLeaseStatus::Active || lease.liveness_at(now) == Liveness::Dead {
507            return Ok(WriterLeaseAuthOutcome::Inactive(lease));
508        }
509        if token_hash(token) != lease.token_hash {
510            return Ok(WriterLeaseAuthOutcome::TokenMismatch);
511        }
512        if (lease.pid.is_some() || lease.pid_namespace.is_some())
513            && !lease.matches_current_pid_namespace()
514        {
515            return Ok(WriterLeaseAuthOutcome::TokenMismatch);
516        }
517        if let Some(existing) = lease.harness_session_id.as_deref()
518            && existing != session
519        {
520            return Ok(WriterLeaseAuthOutcome::TokenMismatch);
521        }
522        if let Some(existing) = lease.pid_birth.as_deref()
523            && (lease.pid != Some(pid) || existing != birth)
524        {
525            return Ok(WriterLeaseAuthOutcome::TokenMismatch);
526        }
527        lease.harness_session_id = Some(session.to_owned());
528        lease.pid = Some(pid);
529        lease.pid_birth = Some(birth.to_owned());
530        lease.pid_namespace = current_pid_namespace();
531        lease.boot_id = super::current_boot_id();
532        lease.heartbeat_at = now;
533        self.write_lease(&lease)?;
534        Ok(WriterLeaseAuthOutcome::Authorized(lease))
535    }
536
537    pub fn release(
538        &self,
539        lease_id: &str,
540        token: &str,
541        status: WriterLeaseStatus,
542        now: DateTime<Utc>,
543    ) -> Result<WriterLeaseAuthOutcome> {
544        let current = self.load(lease_id)?;
545        let _checkout_guard = current
546            .as_ref()
547            .and_then(|lease| lease.path.as_deref())
548            .map(|path| {
549                self.checkout_lock(path)?.write().map_err(|error| {
550                    HeddleError::Config(format!("failed to acquire checkout writer lock: {error}"))
551                })
552            })
553            .transpose()?;
554        self.release_with_checkout_lock(lease_id, token, status, now)
555    }
556
557    /// The caller holds the checkout mutation lock through credential cleanup.
558    pub fn release_with_checkout_lock(
559        &self,
560        lease_id: &str,
561        token: &str,
562        status: WriterLeaseStatus,
563        now: DateTime<Utc>,
564    ) -> Result<WriterLeaseAuthOutcome> {
565        let _lock = self.write_lock()?;
566        let path = self.lease_path(lease_id)?;
567        let Some(mut lease) = self.load_path(&path)? else {
568            return Ok(WriterLeaseAuthOutcome::Missing);
569        };
570        if token_hash(token) != lease.token_hash {
571            return Ok(WriterLeaseAuthOutcome::TokenMismatch);
572        }
573        if lease.status != WriterLeaseStatus::Active {
574            return Ok(WriterLeaseAuthOutcome::Inactive(lease));
575        }
576        lease.status = status;
577        lease.completed_at = Some(now);
578        self.write_lease(&lease)?;
579        Ok(WriterLeaseAuthOutcome::Authorized(lease))
580    }
581
582    pub fn list(&self) -> Result<Vec<WriterLease>> {
583        let _lock = self.write_lock()?;
584        self.reap_expired_locked(Utc::now(), None)?;
585        self.list_locked()
586    }
587
588    /// List persisted leases without reaping expired active records.
589    ///
590    /// Cleanup previews use this read-only view so residue detection cannot
591    /// mutate lease storage as a side effect.
592    pub fn list_without_reaping(&self) -> Result<Vec<WriterLease>> {
593        self.list_locked()
594    }
595
596    /// Remove reservations belonging to task assignments rolled back before launch.
597    pub fn delete_for_tasks(&self, task_ids: &[String]) -> Result<()> {
598        let _lock = self.write_lock()?;
599        for lease in self.list_locked()? {
600            if lease
601                .task_assignment_id
602                .as_ref()
603                .is_some_and(|task_id| task_ids.contains(task_id))
604            {
605                let path = self.lease_path(&lease.lease_id)?;
606                if path.exists() {
607                    std::fs::remove_file(path)?;
608                }
609            }
610        }
611        Ok(())
612    }
613
614    pub fn abandon_thread(&self, thread: &str, now: DateTime<Utc>) -> Result<()> {
615        let _lock = self.write_lock()?;
616        let mut checkout_guards = Vec::new();
617        for lease in self
618            .list_locked()?
619            .into_iter()
620            .filter(|lease| lease.thread == thread && lease.status == WriterLeaseStatus::Active)
621        {
622            if let Some(path) = lease.path.as_deref() {
623                let lock = self.checkout_lock(path)?;
624                let guard = lock
625                    .try_write()
626                    .map_err(|error| {
627                        HeddleError::Config(format!(
628                            "failed to acquire checkout writer lock: {error}"
629                        ))
630                    })?
631                    .ok_or_else(|| HeddleError::Config("checkout has an active mutation".into()))?;
632                checkout_guards.push(guard);
633            }
634        }
635        for mut lease in self.list_locked()? {
636            if lease.thread == thread && lease.status == WriterLeaseStatus::Active {
637                lease.status = WriterLeaseStatus::Abandoned;
638                lease.completed_at = Some(now);
639                self.write_lease(&lease)?;
640            }
641        }
642        Ok(())
643    }
644
645    pub fn load(&self, lease_id: &str) -> Result<Option<WriterLease>> {
646        self.load_path(&self.lease_path(lease_id)?)
647    }
648}
649
650pub fn generate_writer_lease_id() -> String {
651    format!("lease-{}", random_base32())
652}
653
654pub fn generate_writer_lease_token() -> String {
655    format!("hwl_{}", random_base32())
656}
657
658fn random_base32() -> String {
659    let random_bytes: [u8; 24] = rand::random();
660    base32::encode(base32::Alphabet::Rfc4648 { padding: false }, &random_bytes).to_lowercase()
661}
662
663fn token_hash(token: &str) -> String {
664    blake3::hash(token.as_bytes()).to_hex().to_string()
665}
666
667fn validate_lease_id(lease_id: &str) -> Result<()> {
668    if lease_id.starts_with("lease-")
669        && lease_id
670            .bytes()
671            .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'-')
672    {
673        return Ok(());
674    }
675    Err(HeddleError::Config(format!(
676        "invalid writer lease id '{lease_id}'"
677    )))
678}
679
680#[cfg(test)]
681mod tests {
682    use tempfile::TempDir;
683
684    use super::*;
685
686    fn draft(thread: &str) -> WriterLeaseDraft {
687        WriterLeaseDraft {
688            thread: thread.to_string(),
689            actor_session_id: Some("agent-one".to_string()),
690            task_assignment_id: None,
691            anchor_state: Some("hd-state".to_string()),
692            anchor_root: Some("root".to_string()),
693            path: None,
694            pid: None,
695            boot_id: None,
696        }
697    }
698
699    #[cfg(target_os = "linux")]
700    #[test]
701    fn constructing_a_pid_bound_lease_records_its_namespace() {
702        let now = Utc::now();
703        let mut bound = draft("lane");
704        bound.pid = Some(std::process::id());
705        let lease = WriterLease::from_draft(bound, "lease-one".into(), "hash".into(), now);
706        assert_eq!(lease.pid_namespace, current_pid_namespace());
707        assert!(lease.matches_current_pid_namespace());
708    }
709
710    #[cfg(target_os = "linux")]
711    #[test]
712    fn recycled_pid_does_not_keep_bound_lease_alive() {
713        let now = Utc::now();
714        let temp = TempDir::new().expect("lease store");
715        let store = WriterLeaseStore::new(temp.path());
716        let grant = match store.reserve(draft("lane"), now).expect("reserve") {
717            WriterLeaseReserveOutcome::Reserved(grant) => grant,
718            WriterLeaseReserveOutcome::LiveOwner(_) => panic!("new lane has an owner"),
719        };
720        let outcome = store
721            .bind_hook_session(
722                &grant.lease.lease_id,
723                &grant.token,
724                "codex:session-a",
725                std::process::id(),
726                "wrong-birth-tick",
727                now,
728            )
729            .expect("bind");
730        assert!(matches!(outcome, WriterLeaseAuthOutcome::Authorized(_)));
731        assert_eq!(
732            store
733                .load(&grant.lease.lease_id)
734                .expect("load")
735                .expect("lease")
736                .liveness_at(now),
737            Liveness::Dead,
738        );
739        assert!(matches!(
740            store.reserve(draft("lane"), now).expect("reacquire"),
741            WriterLeaseReserveOutcome::Reserved(_)
742        ));
743    }
744
745    #[cfg(target_os = "linux")]
746    #[test]
747    fn pid_namespace_records_device_and_inode() {
748        use std::os::unix::fs::MetadataExt;
749
750        let metadata = std::fs::metadata("/proc/self/ns/pid").expect("PID namespace metadata");
751        assert_eq!(
752            current_pid_namespace(),
753            Some(PidNamespace {
754                device: metadata.dev(),
755                inode: metadata.ino(),
756            })
757        );
758    }
759
760    #[cfg(target_os = "linux")]
761    #[test]
762    fn foreign_pid_namespace_cannot_reap_fresh_bound_lease() {
763        let now = Utc::now();
764        let temp = TempDir::new().expect("lease store");
765        let store = WriterLeaseStore::new(temp.path());
766        let grant = match store.reserve(draft("lane"), now).expect("reserve") {
767            WriterLeaseReserveOutcome::Reserved(grant) => grant,
768            WriterLeaseReserveOutcome::LiveOwner(_) => panic!("new lane has an owner"),
769        };
770        let outcome = store
771            .bind_hook_session(
772                &grant.lease.lease_id,
773                &grant.token,
774                "opencode:session-a",
775                std::process::id(),
776                "wrong-birth-tick",
777                now,
778            )
779            .expect("bind");
780        let WriterLeaseAuthOutcome::Authorized(mut lease) = outcome else {
781            panic!("hook should bind");
782        };
783        assert_eq!(lease.pid_namespace, current_pid_namespace());
784        assert_eq!(
785            lease.liveness_at_in_namespace(now, lease.pid_namespace),
786            Liveness::Dead,
787            "a recycled PID in the same namespace must be rejected"
788        );
789        let namespace = lease.pid_namespace.expect("PID namespace");
790        lease.pid_namespace = None;
791        assert_eq!(lease.liveness_at(now), Liveness::Alive);
792        lease.pid_namespace = Some(PidNamespace {
793            device: namespace.device.wrapping_add(1),
794            ..namespace
795        });
796        store
797            .write_lease(&lease)
798            .expect("simulate foreign namespace");
799        assert_eq!(lease.liveness_at(now), Liveness::Alive);
800        assert!(matches!(
801            store
802                .bind_hook_session(
803                    &grant.lease.lease_id,
804                    &grant.token,
805                    "opencode:session-a",
806                    std::process::id(),
807                    "wrong-birth-tick",
808                    now,
809                )
810                .expect("foreign namespace bind"),
811            WriterLeaseAuthOutcome::TokenMismatch
812        ));
813        assert_eq!(
814            store
815                .load(&grant.lease.lease_id)
816                .expect("load")
817                .expect("lease")
818                .pid_namespace,
819            lease.pid_namespace,
820            "failed bind must preserve the recorded namespace"
821        );
822        assert!(matches!(
823            store.reserve(draft("lane"), now).expect("competing writer"),
824            WriterLeaseReserveOutcome::LiveOwner(_)
825        ));
826        assert!(matches!(
827            store
828                .reserve(draft("lane"), now + chrono::Duration::minutes(6))
829                .expect("expired writer"),
830            WriterLeaseReserveOutcome::Reserved(_)
831        ));
832    }
833
834    #[test]
835    fn separate_checkouts_of_one_thread_have_independent_writers() {
836        let temp = TempDir::new().expect("lease directory");
837        let store = WriterLeaseStore::new(temp.path());
838        for name in ["agent-a", "agent-b"] {
839            let path = temp.path().join(name);
840            std::fs::create_dir(&path).expect("checkout directory");
841            let mut request = draft("shared-thread");
842            request.path = Some(path);
843            assert!(
844                matches!(
845                    store
846                        .reserve(request, Utc::now())
847                        .expect("reserve checkout"),
848                    WriterLeaseReserveOutcome::Reserved(_)
849                ),
850                "different checkouts must not contend on their Thread"
851            );
852        }
853        assert_eq!(store.list().expect("leases").len(), 2);
854    }
855
856    #[test]
857    fn one_checkout_cannot_have_two_writers_even_under_different_thread_names() {
858        let temp = TempDir::new().expect("lease directory");
859        let store = WriterLeaseStore::new(temp.path());
860        let path = temp.path().join("checkout");
861        std::fs::create_dir(&path).expect("checkout directory");
862        let mut first = draft("thread-a");
863        first.path = Some(path.clone());
864        let mut second = draft("thread-b");
865        second.path = Some(path.join("."));
866        assert!(matches!(
867            store.reserve(first, Utc::now()).expect("first writer"),
868            WriterLeaseReserveOutcome::Reserved(_)
869        ));
870        assert!(matches!(
871            store.reserve(second, Utc::now()).expect("competing writer"),
872            WriterLeaseReserveOutcome::LiveOwner(_)
873        ));
874    }
875
876    #[test]
877    fn token_is_required_to_renew_or_release() {
878        let temp = TempDir::new().unwrap();
879        let store = WriterLeaseStore::new(temp.path());
880        let now = Utc::now();
881        let WriterLeaseReserveOutcome::Reserved(grant) =
882            store.reserve(draft("feature/a"), now).unwrap()
883        else {
884            panic!("first lease should reserve");
885        };
886        assert!(matches!(
887            store
888                .authenticate_and_renew(&grant.lease.lease_id, "wrong", now)
889                .unwrap(),
890            WriterLeaseAuthOutcome::TokenMismatch
891        ));
892        assert!(matches!(
893            store
894                .authenticate_and_renew(&grant.lease.lease_id, &grant.token, now)
895                .unwrap(),
896            WriterLeaseAuthOutcome::Authorized(_)
897        ));
898    }
899
900    #[test]
901    fn expired_lease_does_not_block_a_new_owner() {
902        let temp = TempDir::new().unwrap();
903        let store = WriterLeaseStore::new(temp.path());
904        let now = Utc::now();
905        let WriterLeaseReserveOutcome::Reserved(first) =
906            store.reserve(draft("feature/a"), now).unwrap()
907        else {
908            panic!("first lease should reserve");
909        };
910        let later = now + crate::store::AGENT_LEASE_DURATION + chrono::Duration::seconds(1);
911        let WriterLeaseReserveOutcome::Reserved(second) =
912            store.reserve(draft("feature/a"), later).unwrap()
913        else {
914            panic!("expired lease should not block");
915        };
916        assert_ne!(first.lease.lease_id, second.lease.lease_id);
917    }
918
919    #[test]
920    fn stored_lease_does_not_contain_bearer_token() {
921        let temp = TempDir::new().unwrap();
922        let store = WriterLeaseStore::new(temp.path());
923        let WriterLeaseReserveOutcome::Reserved(grant) =
924            store.reserve(draft("feature/a"), Utc::now()).unwrap()
925        else {
926            panic!("first lease should reserve");
927        };
928        let persisted = std::fs::read_to_string(
929            temp.path()
930                .join("writer-leases")
931                .join(format!("{}.toml", grant.lease.lease_id)),
932        )
933        .unwrap();
934        assert!(!persisted.contains(&grant.token));
935        assert!(persisted.contains("token_hash"));
936    }
937}