a73x

66997675

Add a server-side lease store for issue claims

a73x   2026-09-06 07:41

Commit message
Add a server-side lease store for issue claims

Cargo.lock
Old New
@@ -1320,6 +1320,18 @@ dependencies = [
1320 ] 1320 ]
1321 1321
1322 [[package]] 1322 [[package]]
1323 name = "fallible-iterator"
1324 version = "0.3.0"
1325 source = "registry+https://github.com/rust-lang/crates.io-index"
1326 checksum = "2acce4a10f12dc2fb14a218589d4f1f62ef011b2d0cc4b3cb1bba8e94da14649"
1327
1328 [[package]]
1329 name = "fallible-streaming-iterator"
1330 version = "0.1.9"
1331 source = "registry+https://github.com/rust-lang/crates.io-index"
1332 checksum = "7360491ce676a36bf9bb3c56c1aa791658183a54d2744120f27285738d90465a"
1333
1334 [[package]]
1323 name = "fancy-regex" 1335 name = "fancy-regex"
1324 version = "0.11.0" 1336 version = "0.11.0"
1325 source = "registry+https://github.com/rust-lang/crates.io-index" 1337 source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -1619,6 +1631,7 @@ dependencies = [
1619 "rand_core 0.6.4", 1631 "rand_core 0.6.4",
1620 "ratatui", 1632 "ratatui",
1621 "regex", 1633 "regex",
1634 "rusqlite",
1622 "russh", 1635 "russh",
1623 "serde", 1636 "serde",
1624 "serde_json", 1637 "serde_json",
@@ -1692,6 +1705,24 @@ dependencies = [
1692 ] 1705 ]
1693 1706
1694 [[package]] 1707 [[package]]
1708 name = "hashbrown"
1709 version = "0.17.1"
1710 source = "registry+https://github.com/rust-lang/crates.io-index"
1711 checksum = "ed5909b6e89a2db4456e54cd5f673791d7eca6732202bbf2a9cc504fe2f9b84a"
1712 dependencies = [
1713 "foldhash 0.2.0",
1714 ]
1715
1716 [[package]]
1717 name = "hashlink"
1718 version = "0.12.1"
1719 source = "registry+https://github.com/rust-lang/crates.io-index"
1720 checksum = "32069d97bb81e38fa67eab65e3393bf804bb85969f2bc06bf13f64aef5aba248"
1721 dependencies = [
1722 "hashbrown 0.17.1",
1723 ]
1724
1725 [[package]]
1695 name = "heck" 1726 name = "heck"
1696 version = "0.5.0" 1727 version = "0.5.0"
1697 source = "registry+https://github.com/rust-lang/crates.io-index" 1728 source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -2173,6 +2204,17 @@ dependencies = [
2173 ] 2204 ]
2174 2205
2175 [[package]] 2206 [[package]]
2207 name = "libsqlite3-sys"
2208 version = "0.38.2"
2209 source = "registry+https://github.com/rust-lang/crates.io-index"
2210 checksum = "f1d20bef17f513b9b3004532233187769cd072d790971f4e4da0e346eb6401e8"
2211 dependencies = [
2212 "cc",
2213 "pkg-config",
2214 "vcpkg",
2215 ]
2216
2217 [[package]]
2176 name = "libssh2-sys" 2218 name = "libssh2-sys"
2177 version = "0.3.1" 2219 version = "0.3.1"
2178 source = "registry+https://github.com/rust-lang/crates.io-index" 2220 source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -3289,6 +3331,31 @@ dependencies = [
3289 ] 3331 ]
3290 3332
3291 [[package]] 3333 [[package]]
3334 name = "rsqlite-vfs"
3335 version = "0.1.1"
3336 source = "registry+https://github.com/rust-lang/crates.io-index"
3337 checksum = "c51c9ae4df8a7fba42103df5c621fa3c37eccf3a3c650879e90fc48b11cc192c"
3338 dependencies = [
3339 "hashbrown 0.16.1",
3340 "thiserror 2.0.18",
3341 ]
3342
3343 [[package]]
3344 name = "rusqlite"
3345 version = "0.40.2"
3346 source = "registry+https://github.com/rust-lang/crates.io-index"
3347 checksum = "23f2a97da3e3873c73cb2a2e71b35c40ff95e0b1eefa8d72d8499a6928c3b5b3"
3348 dependencies = [
3349 "bitflags 2.11.0",
3350 "fallible-iterator",
3351 "fallible-streaming-iterator",
3352 "hashlink",
3353 "libsqlite3-sys",
3354 "smallvec",
3355 "sqlite-wasm-rs",
3356 ]
3357
3358 [[package]]
3292 name = "russh" 3359 name = "russh"
3293 version = "0.62.5" 3360 version = "0.62.5"
3294 source = "registry+https://github.com/rust-lang/crates.io-index" 3361 source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -3758,6 +3825,18 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
3758 checksum = "3a0219bd7d979d58245a4f41f695e1ac9f8befdffadd7f61f1bae9e39abc6620" 3825 checksum = "3a0219bd7d979d58245a4f41f695e1ac9f8befdffadd7f61f1bae9e39abc6620"
3759 3826
3760 [[package]] 3827 [[package]]
3828 name = "sqlite-wasm-rs"
3829 version = "0.5.5"
3830 source = "registry+https://github.com/rust-lang/crates.io-index"
3831 checksum = "dc3efc0da82635d7e1ced0053bbbfa8c7ab9645d0bf36ceb4f7127bb85315d75"
3832 dependencies = [
3833 "cc",
3834 "js-sys",
3835 "rsqlite-vfs",
3836 "wasm-bindgen",
3837 ]
3838
3839 [[package]]
3761 name = "ssh-cipher" 3840 name = "ssh-cipher"
3762 version = "0.3.0" 3841 version = "0.3.0"
3763 source = "registry+https://github.com/rust-lang/crates.io-index" 3842 source = "registry+https://github.com/rust-lang/crates.io-index"
Cargo.toml
Old New
@@ -55,6 +55,7 @@ regex = "1"
55 unicode-width = "0.2" 55 unicode-width = "0.2"
56 tempfile = "3" 56 tempfile = "3"
57 tokio-util = { version = "0.7", features = ["io"] } 57 tokio-util = { version = "0.7", features = ["io"] }
58 rusqlite = { version = "0.40.2", features = ["bundled"] }
58 59
59 [build-dependencies] 60 [build-dependencies]
60 clap = { version = "4", features = ["derive"] } 61 clap = { version = "4", features = ["derive"] }
src/server/leases.rs
Old New
@@ -0,0 +1,471 @@
1 //! Work leases on issues: the one collaboration primitive git cannot
2 //! express (atomic claim with TTL). One SQLite database per server,
3 //! `(repo, issue_id)` primary key, lazy expiry. See
4 //! docs/superpowers/specs/2026-09-05-server-authoritative-collab-design.md.
5
6 use std::path::Path;
7
8 use rusqlite::{Connection, OptionalExtension, TransactionBehavior};
9
10 /// A live or expired lease row.
11 #[derive(Debug, Clone, PartialEq)]
12 pub struct Lease {
13 pub repo: String,
14 pub issue_id: String,
15 pub holder: String,
16 pub token: i64,
17 pub acquired_at: i64,
18 /// None = open-ended (human assignment).
19 pub expires_at: Option<i64>,
20 }
21
22 impl Lease {
23 pub fn live(&self, now: i64) -> bool {
24 self.expires_at.map(|e| e > now).unwrap_or(true)
25 }
26 }
27
28 #[derive(Debug, PartialEq)]
29 pub enum Acquire {
30 /// Caller now holds the lease (fresh tenure, or idempotent re-acquire).
31 Acquired { token: i64, expires_at: Option<i64> },
32 /// A different holder has a live lease.
33 Held {
34 holder: String,
35 expires_at: Option<i64>,
36 },
37 }
38
39 #[derive(Debug, PartialEq)]
40 pub enum Renew {
41 Renewed { token: i64, expires_at: Option<i64> },
42 /// No live lease held by caller (expired, released, or someone else's).
43 NotHolder { holder: Option<String> },
44 }
45
46 #[derive(Debug, PartialEq)]
47 pub enum Release {
48 /// Row deleted, or there was nothing to release (idempotent success).
49 Released,
50 /// A different holder has a live lease; refuse.
51 NotHolder { holder: String },
52 }
53
54 /// Open (creating if needed) the lease database and ensure the schema.
55 pub fn open(path: &Path) -> rusqlite::Result<Connection> {
56 let conn = Connection::open(path)?;
57 init(&conn)?;
58 Ok(conn)
59 }
60
61 /// Set pragmas and ensure the schema on an already-open connection.
62 /// Factored out of `open` so tests can run against an in-memory database.
63 fn init(conn: &Connection) -> rusqlite::Result<()> {
64 // WAL so a web-UI read never blocks behind an SSH-session write;
65 // busy_timeout so two SSH sessions serialize instead of erroring.
66 conn.pragma_update(None, "journal_mode", "WAL")?;
67 conn.pragma_update(None, "busy_timeout", 5000)?;
68 conn.execute_batch(
69 "CREATE TABLE IF NOT EXISTS leases (
70 repo TEXT NOT NULL,
71 issue_id TEXT NOT NULL,
72 holder TEXT NOT NULL,
73 token INTEGER NOT NULL,
74 acquired_at INTEGER NOT NULL,
75 expires_at INTEGER,
76 PRIMARY KEY (repo, issue_id)
77 )",
78 )
79 }
80
81 fn get_row(conn: &Connection, repo: &str, issue_id: &str) -> rusqlite::Result<Option<Lease>> {
82 conn.query_row(
83 "SELECT repo, issue_id, holder, token, acquired_at, expires_at
84 FROM leases WHERE repo = ?1 AND issue_id = ?2",
85 (repo, issue_id),
86 |row| {
87 Ok(Lease {
88 repo: row.get(0)?,
89 issue_id: row.get(1)?,
90 holder: row.get(2)?,
91 token: row.get(3)?,
92 acquired_at: row.get(4)?,
93 expires_at: row.get(5)?,
94 })
95 },
96 )
97 .optional()
98 }
99
100 pub fn acquire(
101 conn: &mut Connection,
102 repo: &str,
103 issue_id: &str,
104 holder: &str,
105 ttl_secs: Option<i64>,
106 now: i64,
107 ) -> rusqlite::Result<Acquire> {
108 let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
109 let existing = get_row(&tx, repo, issue_id)?;
110 let expires_at = ttl_secs.map(|t| now + t);
111 let outcome = match existing {
112 Some(lease) if lease.live(now) => {
113 if lease.holder == holder {
114 // Idempotent re-acquire: same tenure, same token, fresh expiry
115 // from THIS call's ttl. A retrying client must not deadlock
116 // against itself.
117 tx.execute(
118 "UPDATE leases SET expires_at = ?3 WHERE repo = ?1 AND issue_id = ?2",
119 (repo, issue_id, expires_at),
120 )?;
121 Acquire::Acquired {
122 token: lease.token,
123 expires_at,
124 }
125 } else {
126 Acquire::Held {
127 holder: lease.holder,
128 expires_at: lease.expires_at,
129 }
130 }
131 }
132 other => {
133 // Free, or expired: take tenure. The token increments per tenure
134 // change; it is the fencing token later phases will enforce.
135 let token = other.map(|l| l.token + 1).unwrap_or(1);
136 tx.execute(
137 "INSERT OR REPLACE INTO leases
138 (repo, issue_id, holder, token, acquired_at, expires_at)
139 VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
140 (repo, issue_id, holder, token, now, expires_at),
141 )?;
142 Acquire::Acquired { token, expires_at }
143 }
144 };
145 tx.commit()?;
146 Ok(outcome)
147 }
148
149 pub fn renew(
150 conn: &mut Connection,
151 repo: &str,
152 issue_id: &str,
153 holder: &str,
154 ttl_secs: Option<i64>,
155 now: i64,
156 ) -> rusqlite::Result<Renew> {
157 let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
158 let existing = get_row(&tx, repo, issue_id)?;
159 let outcome = match existing {
160 Some(lease) if lease.live(now) && lease.holder == holder => {
161 let expires_at = ttl_secs.map(|t| now + t);
162 tx.execute(
163 "UPDATE leases SET expires_at = ?3 WHERE repo = ?1 AND issue_id = ?2",
164 (repo, issue_id, expires_at),
165 )?;
166 Renew::Renewed {
167 token: lease.token,
168 expires_at,
169 }
170 }
171 Some(lease) if lease.live(now) => Renew::NotHolder {
172 holder: Some(lease.holder),
173 },
174 _ => Renew::NotHolder { holder: None },
175 };
176 tx.commit()?;
177 Ok(outcome)
178 }
179
180 pub fn release(
181 conn: &mut Connection,
182 repo: &str,
183 issue_id: &str,
184 holder: &str,
185 now: i64,
186 ) -> rusqlite::Result<Release> {
187 let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
188 let existing = get_row(&tx, repo, issue_id)?;
189 let outcome = match existing {
190 Some(lease) if lease.live(now) && lease.holder != holder => {
191 Release::NotHolder {
192 holder: lease.holder,
193 }
194 }
195 Some(_) => {
196 tx.execute(
197 "DELETE FROM leases WHERE repo = ?1 AND issue_id = ?2",
198 (repo, issue_id),
199 )?;
200 Release::Released
201 }
202 // Releasing nothing succeeds: a client retrying after a dropped
203 // connection must not fail.
204 None => Release::Released,
205 };
206 tx.commit()?;
207 Ok(outcome)
208 }
209
210 /// A live lease on the issue, if any. Reaps an expired row it encounters.
211 pub fn current(
212 conn: &mut Connection,
213 repo: &str,
214 issue_id: &str,
215 now: i64,
216 ) -> rusqlite::Result<Option<Lease>> {
217 match get_row(conn, repo, issue_id)? {
218 Some(lease) if lease.live(now) => Ok(Some(lease)),
219 Some(_) => {
220 conn.execute(
221 "DELETE FROM leases WHERE repo = ?1 AND issue_id = ?2 AND expires_at <= ?3",
222 (repo, issue_id, now),
223 )?;
224 Ok(None)
225 }
226 None => Ok(None),
227 }
228 }
229
230 /// All live leases in a repo, oldest first. Reaps expired rows it encounters.
231 pub fn list(conn: &mut Connection, repo: &str, now: i64) -> rusqlite::Result<Vec<Lease>> {
232 conn.execute(
233 "DELETE FROM leases WHERE repo = ?1 AND expires_at IS NOT NULL AND expires_at <= ?2",
234 (repo, now),
235 )?;
236 let mut stmt = conn.prepare(
237 "SELECT repo, issue_id, holder, token, acquired_at, expires_at
238 FROM leases WHERE repo = ?1 ORDER BY acquired_at, issue_id",
239 )?;
240 let rows = stmt.query_map((repo,), |row| {
241 Ok(Lease {
242 repo: row.get(0)?,
243 issue_id: row.get(1)?,
244 holder: row.get(2)?,
245 token: row.get(3)?,
246 acquired_at: row.get(4)?,
247 expires_at: row.get(5)?,
248 })
249 })?;
250 rows.collect()
251 }
252
253 #[cfg(test)]
254 mod tests {
255 use super::*;
256
257 fn mem() -> Connection {
258 let conn = Connection::open_in_memory().unwrap();
259 init(&conn).unwrap();
260 conn
261 }
262
263 #[test]
264 fn acquire_free_issue_returns_token_1() {
265 let mut c = mem();
266 let got = acquire(&mut c, "r", "i", "alice", None, 100).unwrap();
267 assert_eq!(
268 got,
269 Acquire::Acquired {
270 token: 1,
271 expires_at: None
272 }
273 );
274 }
275
276 #[test]
277 fn acquire_sets_expiry_from_ttl() {
278 let mut c = mem();
279 let got = acquire(&mut c, "r", "i", "alice", Some(300), 100).unwrap();
280 assert_eq!(
281 got,
282 Acquire::Acquired {
283 token: 1,
284 expires_at: Some(400)
285 }
286 );
287 }
288
289 #[test]
290 fn acquire_without_ttl_is_open_ended() {
291 let mut c = mem();
292 acquire(&mut c, "r", "i", "alice", None, 100).unwrap();
293 // Far future: still held.
294 let lease = current(&mut c, "r", "i", i64::MAX - 1).unwrap().unwrap();
295 assert_eq!(lease.holder, "alice");
296 assert_eq!(lease.expires_at, None);
297 }
298
299 #[test]
300 fn second_acquire_by_other_holder_returns_held() {
301 let mut c = mem();
302 acquire(&mut c, "r", "i", "alice", Some(300), 100).unwrap();
303 let got = acquire(&mut c, "r", "i", "bob", Some(300), 150).unwrap();
304 assert_eq!(
305 got,
306 Acquire::Held {
307 holder: "alice".into(),
308 expires_at: Some(400)
309 }
310 );
311 }
312
313 #[test]
314 fn reacquire_by_holder_is_idempotent_same_token_new_expiry() {
315 let mut c = mem();
316 acquire(&mut c, "r", "i", "alice", Some(300), 100).unwrap();
317 let got = acquire(&mut c, "r", "i", "alice", Some(300), 200).unwrap();
318 assert_eq!(
319 got,
320 Acquire::Acquired {
321 token: 1,
322 expires_at: Some(500)
323 }
324 );
325 }
326
327 #[test]
328 fn acquire_after_expiry_takes_over_and_bumps_token() {
329 let mut c = mem();
330 acquire(&mut c, "r", "i", "alice", Some(300), 100).unwrap();
331 let got = acquire(&mut c, "r", "i", "bob", Some(300), 500).unwrap();
332 assert_eq!(
333 got,
334 Acquire::Acquired {
335 token: 2,
336 expires_at: Some(800)
337 }
338 );
339 }
340
341 #[test]
342 fn renew_extends_expiry_keeps_token() {
343 let mut c = mem();
344 acquire(&mut c, "r", "i", "alice", Some(300), 100).unwrap();
345 let got = renew(&mut c, "r", "i", "alice", Some(300), 200).unwrap();
346 assert_eq!(
347 got,
348 Renew::Renewed {
349 token: 1,
350 expires_at: Some(500)
351 }
352 );
353 }
354
355 #[test]
356 fn renew_by_non_holder_returns_not_holder() {
357 let mut c = mem();
358 acquire(&mut c, "r", "i", "alice", Some(300), 100).unwrap();
359 let got = renew(&mut c, "r", "i", "bob", Some(300), 200).unwrap();
360 assert_eq!(
361 got,
362 Renew::NotHolder {
363 holder: Some("alice".into())
364 }
365 );
366 }
367
368 #[test]
369 fn renew_after_expiry_returns_not_holder() {
370 let mut c = mem();
371 acquire(&mut c, "r", "i", "alice", Some(300), 100).unwrap();
372 let got = renew(&mut c, "r", "i", "alice", Some(300), 500).unwrap();
373 assert_eq!(got, Renew::NotHolder { holder: None });
374 }
375
376 #[test]
377 fn release_by_holder_deletes() {
378 let mut c = mem();
379 acquire(&mut c, "r", "i", "alice", Some(300), 100).unwrap();
380 assert_eq!(
381 release(&mut c, "r", "i", "alice", 200).unwrap(),
382 Release::Released
383 );
384 assert_eq!(current(&mut c, "r", "i", 200).unwrap(), None);
385 }
386
387 #[test]
388 fn release_idempotent_when_absent() {
389 let mut c = mem();
390 assert_eq!(
391 release(&mut c, "r", "i", "alice", 200).unwrap(),
392 Release::Released
393 );
394 }
395
396 #[test]
397 fn release_by_non_holder_refused() {
398 let mut c = mem();
399 acquire(&mut c, "r", "i", "alice", Some(300), 100).unwrap();
400 assert_eq!(
401 release(&mut c, "r", "i", "bob", 200).unwrap(),
402 Release::NotHolder {
403 holder: "alice".into()
404 }
405 );
406 }
407
408 #[test]
409 fn release_of_expired_lease_by_anyone_succeeds() {
410 let mut c = mem();
411 acquire(&mut c, "r", "i", "alice", Some(300), 100).unwrap();
412 assert_eq!(
413 release(&mut c, "r", "i", "bob", 500).unwrap(),
414 Release::Released
415 );
416 }
417
418 #[test]
419 fn list_shows_only_live_leases() {
420 let mut c = mem();
421 acquire(&mut c, "r", "a", "alice", Some(300), 100).unwrap();
422 acquire(&mut c, "r", "b", "bob", Some(100), 100).unwrap();
423 acquire(&mut c, "r", "c", "carol", None, 150).unwrap();
424 acquire(&mut c, "other", "a", "dave", None, 100).unwrap();
425 let live = list(&mut c, "r", 250).unwrap();
426 let holders: Vec<&str> = live.iter().map(|l| l.holder.as_str()).collect();
427 assert_eq!(holders, vec!["alice", "carol"]);
428 }
429
430 #[test]
431 fn current_none_after_expiry() {
432 let mut c = mem();
433 acquire(&mut c, "r", "i", "alice", Some(300), 100).unwrap();
434 assert!(current(&mut c, "r", "i", 200).unwrap().is_some());
435 assert_eq!(current(&mut c, "r", "i", 400).unwrap(), None);
436 }
437
438 #[test]
439 fn tenure_token_monotonic_across_holders() {
440 let mut c = mem();
441 acquire(&mut c, "r", "i", "alice", Some(100), 0).unwrap();
442 let b = acquire(&mut c, "r", "i", "bob", Some(100), 200).unwrap();
443 assert_eq!(
444 b,
445 Acquire::Acquired {
446 token: 2,
447 expires_at: Some(300)
448 }
449 );
450 let a = acquire(&mut c, "r", "i", "alice", Some(100), 400).unwrap();
451 assert_eq!(
452 a,
453 Acquire::Acquired {
454 token: 3,
455 expires_at: Some(500)
456 }
457 );
458 }
459
460 #[test]
461 fn open_creates_file_and_schema() {
462 let dir = tempfile::tempdir().unwrap();
463 let path = dir.path().join("collab.db");
464 {
465 let mut conn = open(&path).unwrap();
466 acquire(&mut conn, "r", "i", "alice", None, 100).unwrap();
467 }
468 let mut conn = open(&path).unwrap();
469 assert!(current(&mut conn, "r", "i", 200).unwrap().is_some());
470 }
471 }
src/server/main.rs
Old New
@@ -6,6 +6,7 @@ use tracing::info;
6 mod config; 6 mod config;
7 mod governance; 7 mod governance;
8 mod http; 8 mod http;
9 mod leases;
9 mod refs; 10 mod refs;
10 mod releases; 11 mod releases;
11 mod repos; 12 mod repos;