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 store::{HeddleError, Liveness, Result, reservation_liveness_at},
13};
14
15#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
16#[serde(rename_all = "snake_case")]
17pub enum WriterLeaseStatus {
18 Active,
19 Complete,
20 Abandoned,
21}
22
23impl std::fmt::Display for WriterLeaseStatus {
24 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
25 match self {
26 Self::Active => write!(f, "active"),
27 Self::Complete => write!(f, "complete"),
28 Self::Abandoned => write!(f, "abandoned"),
29 }
30 }
31}
32
33#[derive(Debug, Clone, Serialize, Deserialize)]
34pub struct WriterLease {
35 pub lease_id: String,
36 pub thread: String,
37 #[serde(default)]
38 pub actor_session_id: Option<String>,
39 #[serde(default)]
40 pub task_assignment_id: Option<String>,
41 #[serde(default)]
42 pub anchor_state: Option<String>,
43 #[serde(default)]
44 pub anchor_root: Option<String>,
45 #[serde(default)]
46 pub path: Option<PathBuf>,
47 pub token_hash: String,
48 #[serde(default)]
49 pub pid: Option<u32>,
50 #[serde(default)]
51 pub boot_id: Option<String>,
52 pub heartbeat_at: DateTime<Utc>,
53 pub started_at: DateTime<Utc>,
54 pub status: WriterLeaseStatus,
55 #[serde(default)]
56 pub completed_at: Option<DateTime<Utc>>,
57}
58
59impl WriterLease {
60 fn conflicts_with(&self, thread: &str, path: Option<&Path>) -> bool {
61 self.status == WriterLeaseStatus::Active
62 && match (self.path.as_deref(), path) {
63 (Some(owned), Some(requested)) => owned == requested,
64 _ => self.thread == thread,
66 }
67 }
68
69 pub fn lease_expires_at(&self) -> DateTime<Utc> {
70 self.heartbeat_at + crate::store::AGENT_LEASE_DURATION
71 }
72
73 pub fn liveness_at(&self, now: DateTime<Utc>) -> Liveness {
74 if self.status != WriterLeaseStatus::Active {
75 return Liveness::Dead;
76 }
77 reservation_liveness_at(
78 self.pid,
79 self.boot_id.as_deref(),
80 Some(self.heartbeat_at),
81 now,
82 )
83 }
84}
85
86#[derive(Debug, Clone)]
87pub struct WriterLeaseDraft {
88 pub thread: String,
89 pub actor_session_id: Option<String>,
90 pub task_assignment_id: Option<String>,
91 pub anchor_state: Option<String>,
92 pub anchor_root: Option<String>,
93 pub path: Option<PathBuf>,
94 pub pid: Option<u32>,
95 pub boot_id: Option<String>,
96}
97
98#[derive(Debug)]
99pub struct WriterLeaseGrant {
100 pub lease: WriterLease,
101 pub token: String,
102}
103
104#[derive(Debug)]
105pub enum WriterLeaseReserveOutcome {
106 Reserved(WriterLeaseGrant),
107 LiveOwner(WriterLease),
108}
109
110#[derive(Debug)]
111pub enum WriterLeaseAuthOutcome {
112 Authorized(WriterLease),
113 Missing,
114 TokenMismatch,
115 Inactive(WriterLease),
116}
117
118pub struct WriterLeaseStore {
119 leases_dir: PathBuf,
120}
121
122impl WriterLeaseStore {
123 pub fn new(heddle_dir: &Path) -> Self {
124 Self {
125 leases_dir: heddle_dir.join("writer-leases"),
126 }
127 }
128
129 fn lock_path(&self) -> PathBuf {
130 self.leases_dir.join(".lock")
131 }
132
133 fn write_lock(&self) -> Result<crate::lock::WriteLockGuard> {
134 RepoLock::at(self.lock_path()).write().map_err(|err| {
135 HeddleError::Config(format!("failed to acquire writer lease lock: {err}"))
136 })
137 }
138
139 fn lease_path(&self, lease_id: &str) -> Result<PathBuf> {
140 validate_lease_id(lease_id)?;
141 Ok(self.leases_dir.join(format!("{lease_id}.toml")))
142 }
143
144 fn load_path(&self, path: &Path) -> Result<Option<WriterLease>> {
145 if !path.exists() {
146 return Ok(None);
147 }
148 let content = std::fs::read_to_string(path)?;
149 toml::from_str(&content)
150 .map(Some)
151 .map_err(|err| HeddleError::Config(err.to_string()))
152 }
153
154 fn write_lease(&self, lease: &WriterLease) -> Result<()> {
155 crate::fs_atomic::create_dir_all_durable(&self.leases_dir)?;
156 let content =
157 toml::to_string_pretty(lease).map_err(|err| HeddleError::Config(err.to_string()))?;
158 Ok(write_file_atomic(
159 &self.lease_path(&lease.lease_id)?,
160 content.as_bytes(),
161 )?)
162 }
163
164 fn list_locked(&self) -> Result<Vec<WriterLease>> {
165 if !self.leases_dir.exists() {
166 return Ok(Vec::new());
167 }
168 let mut leases = Vec::new();
169 for entry in std::fs::read_dir(&self.leases_dir)? {
170 let path = entry?.path();
171 if path
172 .extension()
173 .is_some_and(|extension| extension == "toml")
174 && let Some(lease) = self.load_path(&path)?
175 {
176 leases.push(lease);
177 }
178 }
179 leases.sort_by_key(|lease| std::cmp::Reverse(lease.started_at));
180 Ok(leases)
181 }
182
183 fn reap_expired_locked(&self, now: DateTime<Utc>) -> Result<()> {
184 for mut lease in self.list_locked()? {
185 if lease.liveness_at(now) == Liveness::Dead && lease.status == WriterLeaseStatus::Active
186 {
187 lease.status = WriterLeaseStatus::Abandoned;
188 lease.completed_at = Some(now);
189 self.write_lease(&lease)?;
190 }
191 }
192 Ok(())
193 }
194
195 pub fn reserve(
196 &self,
197 draft: WriterLeaseDraft,
198 now: DateTime<Utc>,
199 ) -> Result<WriterLeaseReserveOutcome> {
200 self.reserve_prepared(
201 draft,
202 generate_writer_lease_id(),
203 generate_writer_lease_token(),
204 now,
205 )
206 }
207
208 pub fn reserve_prepared(
212 &self,
213 mut draft: WriterLeaseDraft,
214 lease_id: String,
215 token: String,
216 now: DateTime<Utc>,
217 ) -> Result<WriterLeaseReserveOutcome> {
218 validate_lease_id(&lease_id)?;
219 if token.len() < 32 || token.len() > 256 {
220 return Err(HeddleError::Config("invalid prepared writer token".into()));
221 }
222 draft.path = draft.path.map(std::fs::canonicalize).transpose()?;
223 let _lock = self.write_lock()?;
224 self.reap_expired_locked(now)?;
225 if let Some(old) = self.load_path(&self.lease_path(&lease_id)?)? {
226 if old.status == WriterLeaseStatus::Active
227 && old.liveness_at(now) != Liveness::Dead
228 && old.thread == draft.thread
229 && old.path == draft.path
230 && old.actor_session_id == draft.actor_session_id
231 && old.token_hash == token_hash(&token)
232 {
233 return Ok(WriterLeaseReserveOutcome::Reserved(WriterLeaseGrant {
234 lease: old,
235 token,
236 }));
237 }
238 return Err(HeddleError::Config(
239 "prepared writer reservation changed or expired".into(),
240 ));
241 }
242 if let Some(owner) = self
243 .list_locked()?
244 .into_iter()
245 .find(|lease| lease.conflicts_with(&draft.thread, draft.path.as_deref()))
246 {
247 return Ok(WriterLeaseReserveOutcome::LiveOwner(owner));
248 }
249 let lease = WriterLease {
250 lease_id,
251 thread: draft.thread,
252 actor_session_id: draft.actor_session_id,
253 task_assignment_id: draft.task_assignment_id,
254 anchor_state: draft.anchor_state,
255 anchor_root: draft.anchor_root,
256 path: draft.path,
257 token_hash: token_hash(&token),
258 pid: draft.pid,
259 boot_id: draft.boot_id,
260 heartbeat_at: now,
261 started_at: now,
262 status: WriterLeaseStatus::Active,
263 completed_at: None,
264 };
265 self.write_lease(&lease)?;
266 Ok(WriterLeaseReserveOutcome::Reserved(WriterLeaseGrant {
267 lease,
268 token,
269 }))
270 }
271
272 pub fn live_owner(&self, thread: &str, path: Option<&Path>) -> Result<Option<WriterLease>> {
274 let path = path.map(std::fs::canonicalize).transpose()?;
275 Ok(self
276 .list()?
277 .into_iter()
278 .find(|lease| lease.conflicts_with(thread, path.as_deref())))
279 }
280
281 pub fn authenticate_and_renew(
282 &self,
283 lease_id: &str,
284 token: &str,
285 now: DateTime<Utc>,
286 ) -> Result<WriterLeaseAuthOutcome> {
287 let _lock = self.write_lock()?;
288 let path = self.lease_path(lease_id)?;
289 let Some(mut lease) = self.load_path(&path)? else {
290 return Ok(WriterLeaseAuthOutcome::Missing);
291 };
292 if lease.status != WriterLeaseStatus::Active || lease.liveness_at(now) == Liveness::Dead {
293 if lease.status == WriterLeaseStatus::Active {
294 lease.status = WriterLeaseStatus::Abandoned;
295 lease.completed_at = Some(now);
296 self.write_lease(&lease)?;
297 }
298 return Ok(WriterLeaseAuthOutcome::Inactive(lease));
299 }
300 if token_hash(token) != lease.token_hash {
301 return Ok(WriterLeaseAuthOutcome::TokenMismatch);
302 }
303 lease.heartbeat_at = now;
304 self.write_lease(&lease)?;
305 Ok(WriterLeaseAuthOutcome::Authorized(lease))
306 }
307
308 pub fn release(
309 &self,
310 lease_id: &str,
311 token: &str,
312 status: WriterLeaseStatus,
313 now: DateTime<Utc>,
314 ) -> Result<WriterLeaseAuthOutcome> {
315 let _lock = self.write_lock()?;
316 let path = self.lease_path(lease_id)?;
317 let Some(mut lease) = self.load_path(&path)? else {
318 return Ok(WriterLeaseAuthOutcome::Missing);
319 };
320 if lease.status != WriterLeaseStatus::Active {
321 return Ok(WriterLeaseAuthOutcome::Inactive(lease));
322 }
323 if token_hash(token) != lease.token_hash {
324 return Ok(WriterLeaseAuthOutcome::TokenMismatch);
325 }
326 lease.status = status;
327 lease.completed_at = Some(now);
328 self.write_lease(&lease)?;
329 Ok(WriterLeaseAuthOutcome::Authorized(lease))
330 }
331
332 pub fn list(&self) -> Result<Vec<WriterLease>> {
333 let _lock = self.write_lock()?;
334 self.reap_expired_locked(Utc::now())?;
335 self.list_locked()
336 }
337
338 pub fn list_without_reaping(&self) -> Result<Vec<WriterLease>> {
343 self.list_locked()
344 }
345
346 pub fn abandon_thread(&self, thread: &str, now: DateTime<Utc>) -> Result<()> {
347 let _lock = self.write_lock()?;
348 for mut lease in self.list_locked()? {
349 if lease.thread == thread && lease.status == WriterLeaseStatus::Active {
350 lease.status = WriterLeaseStatus::Abandoned;
351 lease.completed_at = Some(now);
352 self.write_lease(&lease)?;
353 }
354 }
355 Ok(())
356 }
357
358 pub fn load(&self, lease_id: &str) -> Result<Option<WriterLease>> {
359 self.load_path(&self.lease_path(lease_id)?)
360 }
361}
362
363pub fn generate_writer_lease_id() -> String {
364 format!("lease-{}", random_base32())
365}
366
367pub fn generate_writer_lease_token() -> String {
368 format!("hwl_{}", random_base32())
369}
370
371fn random_base32() -> String {
372 let random_bytes: [u8; 24] = rand::random();
373 base32::encode(base32::Alphabet::Rfc4648 { padding: false }, &random_bytes).to_lowercase()
374}
375
376fn token_hash(token: &str) -> String {
377 blake3::hash(token.as_bytes()).to_hex().to_string()
378}
379
380fn validate_lease_id(lease_id: &str) -> Result<()> {
381 if lease_id.starts_with("lease-")
382 && lease_id
383 .bytes()
384 .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'-')
385 {
386 return Ok(());
387 }
388 Err(HeddleError::Config(format!(
389 "invalid writer lease id '{lease_id}'"
390 )))
391}
392
393#[cfg(test)]
394mod tests {
395 use tempfile::TempDir;
396
397 use super::*;
398
399 fn draft(thread: &str) -> WriterLeaseDraft {
400 WriterLeaseDraft {
401 thread: thread.to_string(),
402 actor_session_id: Some("agent-one".to_string()),
403 task_assignment_id: None,
404 anchor_state: Some("hd-state".to_string()),
405 anchor_root: Some("root".to_string()),
406 path: None,
407 pid: None,
408 boot_id: None,
409 }
410 }
411
412 #[test]
413 fn separate_checkouts_of_one_thread_have_independent_writers() {
414 let temp = TempDir::new().expect("lease directory");
415 let store = WriterLeaseStore::new(temp.path());
416 for name in ["agent-a", "agent-b"] {
417 let path = temp.path().join(name);
418 std::fs::create_dir(&path).expect("checkout directory");
419 let mut request = draft("shared-thread");
420 request.path = Some(path);
421 assert!(
422 matches!(
423 store
424 .reserve(request, Utc::now())
425 .expect("reserve checkout"),
426 WriterLeaseReserveOutcome::Reserved(_)
427 ),
428 "different checkouts must not contend on their Thread"
429 );
430 }
431 assert_eq!(store.list().expect("leases").len(), 2);
432 }
433
434 #[test]
435 fn one_checkout_cannot_have_two_writers_even_under_different_thread_names() {
436 let temp = TempDir::new().expect("lease directory");
437 let store = WriterLeaseStore::new(temp.path());
438 let path = temp.path().join("checkout");
439 std::fs::create_dir(&path).expect("checkout directory");
440 let mut first = draft("thread-a");
441 first.path = Some(path.clone());
442 let mut second = draft("thread-b");
443 second.path = Some(path.join("."));
444 assert!(matches!(
445 store.reserve(first, Utc::now()).expect("first writer"),
446 WriterLeaseReserveOutcome::Reserved(_)
447 ));
448 assert!(matches!(
449 store.reserve(second, Utc::now()).expect("competing writer"),
450 WriterLeaseReserveOutcome::LiveOwner(_)
451 ));
452 }
453
454 #[test]
455 fn token_is_required_to_renew_or_release() {
456 let temp = TempDir::new().unwrap();
457 let store = WriterLeaseStore::new(temp.path());
458 let now = Utc::now();
459 let WriterLeaseReserveOutcome::Reserved(grant) =
460 store.reserve(draft("feature/a"), now).unwrap()
461 else {
462 panic!("first lease should reserve");
463 };
464 assert!(matches!(
465 store
466 .authenticate_and_renew(&grant.lease.lease_id, "wrong", now)
467 .unwrap(),
468 WriterLeaseAuthOutcome::TokenMismatch
469 ));
470 assert!(matches!(
471 store
472 .authenticate_and_renew(&grant.lease.lease_id, &grant.token, now)
473 .unwrap(),
474 WriterLeaseAuthOutcome::Authorized(_)
475 ));
476 }
477
478 #[test]
479 fn expired_lease_does_not_block_a_new_owner() {
480 let temp = TempDir::new().unwrap();
481 let store = WriterLeaseStore::new(temp.path());
482 let now = Utc::now();
483 let WriterLeaseReserveOutcome::Reserved(first) =
484 store.reserve(draft("feature/a"), now).unwrap()
485 else {
486 panic!("first lease should reserve");
487 };
488 let later = now + crate::store::AGENT_LEASE_DURATION + chrono::Duration::seconds(1);
489 let WriterLeaseReserveOutcome::Reserved(second) =
490 store.reserve(draft("feature/a"), later).unwrap()
491 else {
492 panic!("expired lease should not block");
493 };
494 assert_ne!(first.lease.lease_id, second.lease.lease_id);
495 }
496
497 #[test]
498 fn stored_lease_does_not_contain_bearer_token() {
499 let temp = TempDir::new().unwrap();
500 let store = WriterLeaseStore::new(temp.path());
501 let WriterLeaseReserveOutcome::Reserved(grant) =
502 store.reserve(draft("feature/a"), Utc::now()).unwrap()
503 else {
504 panic!("first lease should reserve");
505 };
506 let persisted = std::fs::read_to_string(
507 temp.path()
508 .join("writer-leases")
509 .join(format!("{}.toml", grant.lease.lease_id)),
510 )
511 .unwrap();
512 assert!(!persisted.contains(&grant.token));
513 assert!(persisted.contains("token_hash"));
514 }
515}