1use 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
16pub 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 _ => 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 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 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 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 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 pub fn live_owner(&self, thread: &str, path: Option<&Path>) -> Result<Option<WriterLease>> {
439 self.live_owner_inner(thread, path, None)
440 }
441
442 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 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 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 pub fn list_without_reaping(&self) -> Result<Vec<WriterLease>> {
593 self.list_locked()
594 }
595
596 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}