Skip to main content

objects/store/
writer_lease.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Exclusive writer leases for agent-controlled thread mutations.
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    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    pub fn lease_expires_at(&self) -> DateTime<Utc> {
61        self.heartbeat_at + crate::store::AGENT_LEASE_DURATION
62    }
63
64    pub fn liveness_at(&self, now: DateTime<Utc>) -> Liveness {
65        if self.status != WriterLeaseStatus::Active {
66            return Liveness::Dead;
67        }
68        reservation_liveness_at(
69            self.pid,
70            self.boot_id.as_deref(),
71            Some(self.heartbeat_at),
72            now,
73        )
74    }
75}
76
77#[derive(Debug, Clone)]
78pub struct WriterLeaseDraft {
79    pub thread: String,
80    pub actor_session_id: Option<String>,
81    pub task_assignment_id: Option<String>,
82    pub anchor_state: Option<String>,
83    pub anchor_root: Option<String>,
84    pub path: Option<PathBuf>,
85    pub pid: Option<u32>,
86    pub boot_id: Option<String>,
87}
88
89#[derive(Debug)]
90pub struct WriterLeaseGrant {
91    pub lease: WriterLease,
92    pub token: String,
93}
94
95#[derive(Debug)]
96pub enum WriterLeaseReserveOutcome {
97    Reserved(WriterLeaseGrant),
98    LiveOwner(WriterLease),
99}
100
101#[derive(Debug)]
102pub enum WriterLeaseAuthOutcome {
103    Authorized(WriterLease),
104    Missing,
105    TokenMismatch,
106    Inactive(WriterLease),
107}
108
109pub struct WriterLeaseStore {
110    leases_dir: PathBuf,
111}
112
113impl WriterLeaseStore {
114    pub fn new(heddle_dir: &Path) -> Self {
115        Self {
116            leases_dir: heddle_dir.join("writer-leases"),
117        }
118    }
119
120    fn lock_path(&self) -> PathBuf {
121        self.leases_dir.join(".lock")
122    }
123
124    fn write_lock(&self) -> Result<crate::lock::WriteLockGuard> {
125        RepoLock::at(self.lock_path()).write().map_err(|err| {
126            HeddleError::Config(format!("failed to acquire writer lease lock: {err}"))
127        })
128    }
129
130    fn lease_path(&self, lease_id: &str) -> Result<PathBuf> {
131        validate_lease_id(lease_id)?;
132        Ok(self.leases_dir.join(format!("{lease_id}.toml")))
133    }
134
135    fn load_path(&self, path: &Path) -> Result<Option<WriterLease>> {
136        if !path.exists() {
137            return Ok(None);
138        }
139        let content = std::fs::read_to_string(path)?;
140        toml::from_str(&content)
141            .map(Some)
142            .map_err(|err| HeddleError::Config(err.to_string()))
143    }
144
145    fn write_lease(&self, lease: &WriterLease) -> Result<()> {
146        crate::fs_atomic::create_dir_all_durable(&self.leases_dir)?;
147        let content =
148            toml::to_string_pretty(lease).map_err(|err| HeddleError::Config(err.to_string()))?;
149        Ok(write_file_atomic(
150            &self.lease_path(&lease.lease_id)?,
151            content.as_bytes(),
152        )?)
153    }
154
155    fn list_locked(&self) -> Result<Vec<WriterLease>> {
156        if !self.leases_dir.exists() {
157            return Ok(Vec::new());
158        }
159        let mut leases = Vec::new();
160        for entry in std::fs::read_dir(&self.leases_dir)? {
161            let path = entry?.path();
162            if path
163                .extension()
164                .is_some_and(|extension| extension == "toml")
165                && let Some(lease) = self.load_path(&path)?
166            {
167                leases.push(lease);
168            }
169        }
170        leases.sort_by_key(|lease| std::cmp::Reverse(lease.started_at));
171        Ok(leases)
172    }
173
174    fn reap_expired_locked(&self, now: DateTime<Utc>) -> Result<()> {
175        for mut lease in self.list_locked()? {
176            if lease.liveness_at(now) == Liveness::Dead && lease.status == WriterLeaseStatus::Active
177            {
178                lease.status = WriterLeaseStatus::Abandoned;
179                lease.completed_at = Some(now);
180                self.write_lease(&lease)?;
181            }
182        }
183        Ok(())
184    }
185
186    pub fn reserve(
187        &self,
188        draft: WriterLeaseDraft,
189        now: DateTime<Utc>,
190    ) -> Result<WriterLeaseReserveOutcome> {
191        let _lock = self.write_lock()?;
192        self.reap_expired_locked(now)?;
193        if let Some(owner) = self
194            .list_locked()?
195            .into_iter()
196            .find(|lease| lease.thread == draft.thread && lease.status == WriterLeaseStatus::Active)
197        {
198            return Ok(WriterLeaseReserveOutcome::LiveOwner(owner));
199        }
200
201        let lease_id = generate_writer_lease_id();
202        let token = generate_writer_lease_token();
203        let lease = WriterLease {
204            lease_id,
205            thread: draft.thread,
206            actor_session_id: draft.actor_session_id,
207            task_assignment_id: draft.task_assignment_id,
208            anchor_state: draft.anchor_state,
209            anchor_root: draft.anchor_root,
210            path: draft.path,
211            token_hash: token_hash(&token),
212            pid: draft.pid,
213            boot_id: draft.boot_id,
214            heartbeat_at: now,
215            started_at: now,
216            status: WriterLeaseStatus::Active,
217            completed_at: None,
218        };
219        self.write_lease(&lease)?;
220        Ok(WriterLeaseReserveOutcome::Reserved(WriterLeaseGrant {
221            lease,
222            token,
223        }))
224    }
225
226    pub fn authenticate_and_renew(
227        &self,
228        lease_id: &str,
229        token: &str,
230        now: DateTime<Utc>,
231    ) -> Result<WriterLeaseAuthOutcome> {
232        let _lock = self.write_lock()?;
233        let path = self.lease_path(lease_id)?;
234        let Some(mut lease) = self.load_path(&path)? else {
235            return Ok(WriterLeaseAuthOutcome::Missing);
236        };
237        if lease.status != WriterLeaseStatus::Active || lease.liveness_at(now) == Liveness::Dead {
238            if lease.status == WriterLeaseStatus::Active {
239                lease.status = WriterLeaseStatus::Abandoned;
240                lease.completed_at = Some(now);
241                self.write_lease(&lease)?;
242            }
243            return Ok(WriterLeaseAuthOutcome::Inactive(lease));
244        }
245        if token_hash(token) != lease.token_hash {
246            return Ok(WriterLeaseAuthOutcome::TokenMismatch);
247        }
248        lease.heartbeat_at = now;
249        self.write_lease(&lease)?;
250        Ok(WriterLeaseAuthOutcome::Authorized(lease))
251    }
252
253    pub fn release(
254        &self,
255        lease_id: &str,
256        token: &str,
257        status: WriterLeaseStatus,
258        now: DateTime<Utc>,
259    ) -> Result<WriterLeaseAuthOutcome> {
260        let _lock = self.write_lock()?;
261        let path = self.lease_path(lease_id)?;
262        let Some(mut lease) = self.load_path(&path)? else {
263            return Ok(WriterLeaseAuthOutcome::Missing);
264        };
265        if lease.status != WriterLeaseStatus::Active {
266            return Ok(WriterLeaseAuthOutcome::Inactive(lease));
267        }
268        if token_hash(token) != lease.token_hash {
269            return Ok(WriterLeaseAuthOutcome::TokenMismatch);
270        }
271        lease.status = status;
272        lease.completed_at = Some(now);
273        self.write_lease(&lease)?;
274        Ok(WriterLeaseAuthOutcome::Authorized(lease))
275    }
276
277    pub fn list(&self) -> Result<Vec<WriterLease>> {
278        let _lock = self.write_lock()?;
279        self.reap_expired_locked(Utc::now())?;
280        self.list_locked()
281    }
282
283    pub fn abandon_thread(&self, thread: &str, now: DateTime<Utc>) -> Result<()> {
284        let _lock = self.write_lock()?;
285        for mut lease in self.list_locked()? {
286            if lease.thread == thread && lease.status == WriterLeaseStatus::Active {
287                lease.status = WriterLeaseStatus::Abandoned;
288                lease.completed_at = Some(now);
289                self.write_lease(&lease)?;
290            }
291        }
292        Ok(())
293    }
294
295    pub fn load(&self, lease_id: &str) -> Result<Option<WriterLease>> {
296        self.load_path(&self.lease_path(lease_id)?)
297    }
298}
299
300pub fn generate_writer_lease_id() -> String {
301    format!("lease-{}", random_base32())
302}
303
304pub fn generate_writer_lease_token() -> String {
305    format!("hwl_{}", random_base32())
306}
307
308fn random_base32() -> String {
309    let random_bytes: [u8; 24] = rand::random();
310    base32::encode(base32::Alphabet::Rfc4648 { padding: false }, &random_bytes).to_lowercase()
311}
312
313fn token_hash(token: &str) -> String {
314    blake3::hash(token.as_bytes()).to_hex().to_string()
315}
316
317fn validate_lease_id(lease_id: &str) -> Result<()> {
318    if lease_id.starts_with("lease-")
319        && lease_id
320            .bytes()
321            .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'-')
322    {
323        return Ok(());
324    }
325    Err(HeddleError::Config(format!(
326        "invalid writer lease id '{lease_id}'"
327    )))
328}
329
330#[cfg(test)]
331mod tests {
332    use tempfile::TempDir;
333
334    use super::*;
335
336    fn draft(thread: &str) -> WriterLeaseDraft {
337        WriterLeaseDraft {
338            thread: thread.to_string(),
339            actor_session_id: Some("agent-one".to_string()),
340            task_assignment_id: None,
341            anchor_state: Some("hd-state".to_string()),
342            anchor_root: Some("root".to_string()),
343            path: None,
344            pid: None,
345            boot_id: None,
346        }
347    }
348
349    #[test]
350    fn token_is_required_to_renew_or_release() {
351        let temp = TempDir::new().unwrap();
352        let store = WriterLeaseStore::new(temp.path());
353        let now = Utc::now();
354        let WriterLeaseReserveOutcome::Reserved(grant) =
355            store.reserve(draft("feature/a"), now).unwrap()
356        else {
357            panic!("first lease should reserve");
358        };
359        assert!(matches!(
360            store
361                .authenticate_and_renew(&grant.lease.lease_id, "wrong", now)
362                .unwrap(),
363            WriterLeaseAuthOutcome::TokenMismatch
364        ));
365        assert!(matches!(
366            store
367                .authenticate_and_renew(&grant.lease.lease_id, &grant.token, now)
368                .unwrap(),
369            WriterLeaseAuthOutcome::Authorized(_)
370        ));
371    }
372
373    #[test]
374    fn expired_lease_does_not_block_a_new_owner() {
375        let temp = TempDir::new().unwrap();
376        let store = WriterLeaseStore::new(temp.path());
377        let now = Utc::now();
378        let WriterLeaseReserveOutcome::Reserved(first) =
379            store.reserve(draft("feature/a"), now).unwrap()
380        else {
381            panic!("first lease should reserve");
382        };
383        let later = now + crate::store::AGENT_LEASE_DURATION + chrono::Duration::seconds(1);
384        let WriterLeaseReserveOutcome::Reserved(second) =
385            store.reserve(draft("feature/a"), later).unwrap()
386        else {
387            panic!("expired lease should not block");
388        };
389        assert_ne!(first.lease.lease_id, second.lease.lease_id);
390    }
391
392    #[test]
393    fn stored_lease_does_not_contain_bearer_token() {
394        let temp = TempDir::new().unwrap();
395        let store = WriterLeaseStore::new(temp.path());
396        let WriterLeaseReserveOutcome::Reserved(grant) =
397            store.reserve(draft("feature/a"), Utc::now()).unwrap()
398        else {
399            panic!("first lease should reserve");
400        };
401        let persisted = std::fs::read_to_string(
402            temp.path()
403                .join("writer-leases")
404                .join(format!("{}.toml", grant.lease.lease_id)),
405        )
406        .unwrap();
407        assert!(!persisted.contains(&grant.token));
408        assert!(persisted.contains("token_hash"));
409    }
410}