1use std::{
5 path::{Path, PathBuf},
6 sync::{
7 Arc, Mutex,
8 atomic::{AtomicU64, Ordering},
9 },
10 time::SystemTime,
11};
12
13use heddle_object_model::op_record::OpRecord;
14use objects::{
15 error::{HeddleError, Result},
16 object::{MarkerName, StateId, ThreadName},
17};
18
19use super::{
20 Head, RefExpectation, RefUpdate,
21 backend::CoreRefBackend,
22 format_state_id_text,
23 packed_refs::PackedRefs,
24 reconcile::{LoadRequest, Loaded, RefClass, RefCommitter, RefReconciler},
25 ref_backend::RefBackend,
26 refs_storage::RefsLock,
27 resolve_refspec,
28};
29use crate::fs_atomic::{create_dir_all_durable, sync_directory};
30
31const WATERMARK_UNSET: u64 = u64::MAX;
34
35fn watermark_covers(cached: u64, tip: u64) -> bool {
36 cached != WATERMARK_UNSET && cached >= tip
37}
38
39const RECONCILE_WATERMARK_LOCAL: &str = "RECONCILE_WATERMARK_LOCAL";
43
44const RECONCILE_WATERMARK_SHARED: &str = "RECONCILE_WATERMARK_SHARED";
50
51const SNAPSHOT_WITNESS_LOCAL: &str = "SNAPSHOT_WITNESS_LOCAL";
55
56const SNAPSHOT_WITNESS_SHARED: &str = "SNAPSHOT_WITNESS_SHARED";
58
59pub const UNDO_RECOVERY_HANDLE: &str = ".undo-recovery";
76
77struct CachedPackedRefs {
83 stamp: Option<(SystemTime, u64)>,
84 packed: Arc<PackedRefs>,
85}
86
87pub struct RefManager {
89 pub(crate) root: PathBuf,
90 pub(crate) local_head: Option<PathBuf>,
91 reconciler: Option<Arc<dyn RefReconciler>>,
95 committer: Option<Arc<dyn RefCommitter>>,
99 cached_local_generation: AtomicU64,
102 cached_shared_generation: AtomicU64,
105 packed_refs_cache: Mutex<Option<CachedPackedRefs>>,
107}
108
109impl RefManager {
110 pub fn new(heddle_dir: impl AsRef<Path>) -> Self {
111 Self {
112 root: heddle_dir.as_ref().to_path_buf(),
113 local_head: None,
114 reconciler: None,
115 committer: None,
116 cached_local_generation: AtomicU64::new(WATERMARK_UNSET),
117 cached_shared_generation: AtomicU64::new(WATERMARK_UNSET),
118 packed_refs_cache: Mutex::new(None),
119 }
120 }
121
122 pub fn with_local_head(mut self, path: PathBuf) -> Self {
123 self.local_head = Some(path);
124 self
125 }
126
127 pub fn with_reconciler(mut self, reconciler: Arc<dyn RefReconciler>) -> Self {
143 let generation = reconciler.generation().unwrap_or(WATERMARK_UNSET);
148 self.cached_local_generation
149 .store(generation, Ordering::Release);
150 self.cached_shared_generation
151 .store(generation, Ordering::Release);
152 self.reconciler = Some(reconciler);
153 self
154 }
155
156 pub fn with_committer(mut self, committer: Arc<dyn RefCommitter>) -> Self {
160 self.committer = Some(committer);
161 self
162 }
163
164 pub fn commit_and_publish(
185 &self,
186 records: &[OpRecord],
187 ref_updates: &[RefUpdate],
188 scope: Option<&str>,
189 ) -> Result<()> {
190 self.write_chokepoint(|lock| {
191 self.validate_commit_publish(ref_updates, lock, || {
192 let committed_for_reconcile = self.committer.is_some() && !records.is_empty();
195 if let Some(committer) = self.committer.as_ref() {
196 committer.commit_records(records, scope)?;
197 } else if !records.is_empty() {
198 return Err(HeddleError::Config(format!(
204 "commit_and_publish was handed {} record(s) but this RefManager has no \
205 committer; refusing to publish and silently drop committed data",
206 records.len()
207 )));
208 }
209 Ok(committed_for_reconcile)
210 })
211 })
212 }
213
214 pub(super) fn write_chokepoint<T>(
241 &self,
242 body: impl FnOnce(&RefsLock) -> Result<T>,
243 ) -> Result<T> {
244 let lock = self.lock_refs()?;
245 self.materialize_committed_tail(&lock)?;
246 body(&lock)
247 }
248
249 fn materialize_committed_tail(&self, lock: &RefsLock) -> Result<()> {
255 self.materialize_class(RefClass::Local, lock)?;
256 self.materialize_class(RefClass::Shared, lock)?;
257 Ok(())
258 }
259
260 pub fn materialize_snapshot_thread_after_commit(
267 &self,
268 thread: &ThreadName,
269 state: StateId,
270 tip: u64,
271 ) -> Result<()> {
272 let lock = self.lock_refs()?;
273 let outcome = super::reconcile::ReconcileOutcome {
274 loaded: Loaded::Point(Some(state)),
275 republish: vec![RefUpdate::Thread {
276 name: thread.clone(),
277 expected: RefExpectation::Any,
278 new: Some(state),
279 }],
280 remote_updates: Vec::new(),
281 undo_recovery: None,
282 };
283 self.materialize_with_ref_durability(&outcome, &lock, false)?;
284 self.cached_shared_generation.store(tip, Ordering::Release);
285 self.cached_local_generation.store(tip, Ordering::Release);
288 let _ = std::fs::write(
289 self.snapshot_witness_local_path(),
290 format!("{tip}\nattached\n{thread}\n"),
291 );
292 let _ = std::fs::write(
293 self.snapshot_witness_shared_path(),
294 format!("{tip}\nthread\n{thread}\n{}\n", state.to_string_full()),
295 );
296 Ok(())
297 }
298
299 pub fn materialize_snapshot_head_after_commit(&self, state: StateId, tip: u64) -> Result<()> {
302 let lock = self.lock_refs()?;
303 let outcome = super::reconcile::ReconcileOutcome {
304 loaded: Loaded::Head(Head::Detached { state }),
305 republish: vec![RefUpdate::Head {
306 expected: RefExpectation::Any,
307 new: Head::Detached { state },
308 }],
309 remote_updates: Vec::new(),
310 undo_recovery: None,
311 };
312 self.materialize_with_ref_durability(&outcome, &lock, false)?;
313 self.cached_local_generation.store(tip, Ordering::Release);
314 self.cached_shared_generation.store(tip, Ordering::Release);
316 let _ = std::fs::write(
317 self.snapshot_witness_local_path(),
318 format!("{tip}\ndetached\n{}\n", state.to_string_full()),
319 );
320 Ok(())
321 }
322
323 fn materialize_class(&self, class: RefClass, lock: &RefsLock) -> Result<()> {
330 let Some(reconciler) = self.reconciler.as_ref() else {
331 return Ok(());
332 };
333 let watermark = self.class_watermark(class);
334 let tip = reconciler.generation()?;
335 if watermark_covers(watermark.load(Ordering::Acquire), tip) {
336 return Ok(());
337 }
338 self.refresh_persisted_watermark(class, lock)?;
342 let cached = watermark.load(Ordering::Acquire);
343 if watermark_covers(cached, tip) {
344 return Ok(());
345 }
346 let req = Self::class_probe(class);
347 let raw = self.raw_load(&req)?;
348 let since = if cached == WATERMARK_UNSET { 0 } else { cached };
349 let outcome = reconciler.reconcile(&req, raw, since)?;
350 self.materialize(&outcome, lock)?;
351 watermark.store(tip, Ordering::Release);
352 let _ = self.persist_reconcile_watermark(lock);
353 Ok(())
354 }
355
356 fn class_watermark(&self, class: RefClass) -> &AtomicU64 {
358 match class {
359 RefClass::Local => &self.cached_local_generation,
360 RefClass::Shared => &self.cached_shared_generation,
361 }
362 }
363
364 fn class_probe(class: RefClass) -> LoadRequest {
369 match class {
370 RefClass::Local => LoadRequest::Head,
371 RefClass::Shared => LoadRequest::MarkerList,
372 }
373 }
374
375 fn refresh_persisted_watermark(&self, class: RefClass, _lock: &RefsLock) -> Result<()> {
387 let path = match class {
388 RefClass::Local => self.reconcile_watermark_local_path(),
389 RefClass::Shared => self.reconcile_watermark_shared_path(),
390 };
391 let Some(persisted) = self.read_single_watermark(&path)? else {
392 return Ok(());
393 };
394 let watermark = self.class_watermark(class);
395 let cached = watermark.load(Ordering::Acquire);
396 let next = if cached == WATERMARK_UNSET {
399 persisted
400 } else {
401 cached.max(persisted)
402 };
403 if next != cached {
404 watermark.store(next, Ordering::Release);
405 }
406 Ok(())
407 }
408
409 fn reconciled_load(&self, req: LoadRequest) -> Result<Loaded> {
416 heddle_perf_contract::record_ref_read();
417 let Some(reconciler) = self.reconciler.as_ref() else {
418 return self.raw_load(&req);
419 };
420
421 let watermark = match req.ref_class() {
422 RefClass::Local => &self.cached_local_generation,
423 RefClass::Shared => &self.cached_shared_generation,
424 };
425
426 let tip = reconciler.generation()?;
432 if watermark_covers(watermark.load(Ordering::Acquire), tip) {
433 return self.raw_load(&req);
434 }
435 if let Some(loaded) = self.snapshot_witnessed_load(&req, tip)? {
436 return Ok(loaded);
437 }
438
439 let lock = self.lock_refs()?;
447 let tip = reconciler.generation()?;
448 self.refresh_persisted_watermark(req.ref_class(), &lock)?;
455 let cached = watermark.load(Ordering::Acquire);
456 let raw = self.raw_load(&req)?;
457 if watermark_covers(cached, tip) {
458 return Ok(raw);
461 }
462
463 let since = if cached == WATERMARK_UNSET { 0 } else { cached };
468 let outcome = reconciler.reconcile(&req, raw, since)?;
469 self.materialize(&outcome, &lock)?;
470 watermark.store(tip, Ordering::Release);
471 let _ = self.persist_reconcile_watermark(&lock);
477 Ok(outcome.loaded)
478 }
479
480 pub(super) fn reconciled_value_under_lock(&self, req: &LoadRequest) -> Result<Loaded> {
488 let raw = self.raw_load(req)?;
489 let Some(reconciler) = self.reconciler.as_ref() else {
490 return Ok(raw);
491 };
492 let tip = reconciler.generation()?;
493 let watermark = match req.ref_class() {
494 RefClass::Local => &self.cached_local_generation,
495 RefClass::Shared => &self.cached_shared_generation,
496 };
497 let cached = watermark.load(Ordering::Acquire);
498 if watermark_covers(cached, tip) {
499 return Ok(raw);
500 }
501 let since = if cached == WATERMARK_UNSET { 0 } else { cached };
502 Ok(reconciler.reconcile(req, raw, since)?.loaded)
503 }
504
505 pub fn init_reconcile_watermark(&self) -> Result<()> {
532 if self.reconciler.is_none() {
533 return Ok(());
534 }
535 let (local, shared) = self.read_persisted_reconcile_watermark()?;
536 if let Some(local) = local {
537 self.cached_local_generation.store(local, Ordering::Release);
538 }
539 if let Some(shared) = shared {
540 self.cached_shared_generation
541 .store(shared, Ordering::Release);
542 }
543 if local.is_none() || shared.is_none() {
548 let lock = self.lock_refs()?;
549 self.persist_reconcile_watermark(&lock)?;
550 }
551 Ok(())
552 }
553
554 fn reconcile_watermark_local_path(&self) -> PathBuf {
557 self.head_path()
558 .parent()
559 .map(|dir| dir.join(RECONCILE_WATERMARK_LOCAL))
560 .unwrap_or_else(|| self.root.join(RECONCILE_WATERMARK_LOCAL))
561 }
562
563 fn reconcile_watermark_shared_path(&self) -> PathBuf {
569 self.root.join(RECONCILE_WATERMARK_SHARED)
570 }
571
572 fn snapshot_witness_local_path(&self) -> PathBuf {
573 self.head_path()
574 .parent()
575 .map(|dir| dir.join(SNAPSHOT_WITNESS_LOCAL))
576 .unwrap_or_else(|| self.root.join(SNAPSHOT_WITNESS_LOCAL))
577 }
578
579 fn snapshot_witness_shared_path(&self) -> PathBuf {
580 self.root.join(SNAPSHOT_WITNESS_SHARED)
581 }
582
583 fn snapshot_witnessed_load(&self, request: &LoadRequest, tip: u64) -> Result<Option<Loaded>> {
587 let parse_tip = |lines: &mut std::str::Lines<'_>| {
588 lines
589 .next()
590 .and_then(|value| value.parse::<u64>().ok())
591 .filter(|value| *value == tip)
592 };
593 match request {
594 LoadRequest::Head => {
595 let Some(contents) =
596 self.read_optional_string(&self.snapshot_witness_local_path())?
597 else {
598 return Ok(None);
599 };
600 let mut lines = contents.lines();
601 if parse_tip(&mut lines).is_none() {
602 return Ok(None);
603 }
604 let witnessed = match lines.next() {
605 Some("attached") => lines.next().map(|thread| Head::Attached {
606 thread: ThreadName::new(thread),
607 }),
608 Some("detached") => lines
609 .next()
610 .and_then(|state| StateId::parse(state).ok())
611 .map(|state| Head::Detached { state }),
612 _ => None,
613 };
614 let Some(witnessed) = witnessed else {
615 return Ok(None);
616 };
617 let raw = self.read_head_state()?.head;
618 Ok((raw == witnessed).then_some(Loaded::Head(raw)))
619 }
620 LoadRequest::Thread(requested) => {
621 let Some(contents) =
622 self.read_optional_string(&self.snapshot_witness_shared_path())?
623 else {
624 return Ok(None);
625 };
626 let mut lines = contents.lines();
627 if parse_tip(&mut lines).is_none() || lines.next() != Some("thread") {
628 return Ok(None);
629 }
630 let Some(thread) = lines.next() else {
631 return Ok(None);
632 };
633 let Some(state) = lines.next().and_then(|value| StateId::parse(value).ok()) else {
634 return Ok(None);
635 };
636 if thread != requested.as_str() {
637 return Ok(None);
638 }
639 let raw = self.raw_get_thread(requested)?;
640 Ok((raw == Some(state)).then_some(Loaded::Point(raw)))
641 }
642 _ => Ok(None),
643 }
644 }
645
646 fn read_persisted_reconcile_watermark(&self) -> Result<(Option<u64>, Option<u64>)> {
650 let local = self.read_single_watermark(&self.reconcile_watermark_local_path())?;
651 let shared = self.read_single_watermark(&self.reconcile_watermark_shared_path())?;
652 Ok((local, shared))
653 }
654
655 fn read_single_watermark(&self, path: &Path) -> Result<Option<u64>> {
658 let Some(contents) = self.read_optional_string(path)? else {
659 return Ok(None);
660 };
661 Ok(contents
662 .split_whitespace()
663 .next()
664 .and_then(|s| s.parse::<u64>().ok()))
665 }
666
667 fn persist_reconcile_watermark(&self, _lock: &RefsLock) -> Result<()> {
671 let local = self.cached_local_generation.load(Ordering::Acquire);
672 let shared = self.cached_shared_generation.load(Ordering::Acquire);
673 self.persist_watermark_file(&self.reconcile_watermark_local_path(), local)?;
674 self.persist_watermark_file(&self.reconcile_watermark_shared_path(), shared)?;
675 Ok(())
676 }
677
678 fn persist_watermark_file(&self, path: &Path, value: u64) -> Result<()> {
686 if value == WATERMARK_UNSET {
687 return Ok(());
688 }
689 let on_disk = self.read_single_watermark(path)?.unwrap_or(0);
690 let next = value.max(on_disk);
691 self.write_string(path, &format!("{next}\n"))
692 }
693
694 fn materialize(
710 &self,
711 outcome: &super::reconcile::ReconcileOutcome,
712 lock: &RefsLock,
713 ) -> Result<()> {
714 self.materialize_with_ref_durability(outcome, lock, true)
715 }
716
717 fn materialize_with_ref_durability(
718 &self,
719 outcome: &super::reconcile::ReconcileOutcome,
720 lock: &RefsLock,
721 durable_refs: bool,
722 ) -> Result<()> {
723 let plans = self.plan_materialization(&outcome.republish)?;
729 if !plans.is_empty() {
730 if durable_refs {
731 self.publish_ref_plans(plans, lock)?;
732 } else {
733 self.publish_ref_plans_reconstructible(plans, lock)?;
734 }
735 }
736 for (remote, thread, value) in &outcome.remote_updates {
737 if self.raw_get_remote_thread(remote, thread)? != *value {
738 match value {
739 Some(state) => self.set_remote_thread_locked(remote, thread, state, lock)?,
740 None => {
741 self.delete_remote_thread_locked(remote, thread, lock)?;
742 }
743 }
744 }
745 }
746 if let Some(state) = &outcome.undo_recovery {
747 let current = self.read_state_id_at(
748 &self.undo_recovery_path(),
749 "undo recovery",
750 UNDO_RECOVERY_HANDLE,
751 )?;
752 if current.as_ref() != Some(state) {
753 self.set_undo_recovery_locked(state, lock)?;
754 }
755 }
756 Ok(())
757 }
758
759 fn raw_load(&self, req: &LoadRequest) -> Result<Loaded> {
763 Ok(match req {
764 LoadRequest::Head => Loaded::Head(self.read_head_state()?.head),
765 LoadRequest::Thread(name) => Loaded::Point(self.raw_get_thread(name)?),
766 LoadRequest::Marker(name) => Loaded::Point(self.raw_get_marker(name)?),
767 LoadRequest::UndoRecovery => Loaded::Point(self.read_state_id_at(
768 &self.undo_recovery_path(),
769 "undo recovery",
770 UNDO_RECOVERY_HANDLE,
771 )?),
772 LoadRequest::RemoteThread { remote, thread } => {
773 Loaded::Point(self.raw_get_remote_thread(remote, thread)?)
774 }
775 LoadRequest::ThreadList => Loaded::ThreadList(self.raw_list_threads()?),
776 LoadRequest::MarkerList => Loaded::MarkerList(self.raw_list_markers()?),
777 LoadRequest::RemoteList => Loaded::RemoteList(self.raw_list_remotes()?),
778 LoadRequest::RemoteThreadList { remote } => {
779 Loaded::RemoteThreadList(self.raw_list_remote_threads(remote)?)
780 }
781 })
782 }
783
784 fn raw_get_thread(&self, name: &ThreadName) -> Result<Option<StateId>> {
785 let path = self.thread_path(name)?;
786 if let Some(id) = self.read_state_id_at(&path, "thread", name)? {
787 return Ok(Some(id));
788 }
789 Ok(self.load_packed_refs_cached()?.get_thread(name))
790 }
791
792 fn raw_get_marker(&self, name: &MarkerName) -> Result<Option<StateId>> {
793 let path = self.marker_path(name)?;
794 if let Some(id) = self.read_state_id_at(&path, "marker", name)? {
795 return Ok(Some(id));
796 }
797 Ok(self.load_packed_refs_cached()?.get_marker(name))
798 }
799
800 fn packed_refs_stamp(path: &Path) -> Option<(SystemTime, u64)> {
804 let meta = std::fs::metadata(path).ok()?;
805 let modified = meta.modified().ok()?;
806 Some((modified, meta.len()))
807 }
808
809 pub(super) fn load_packed_refs_cached(&self) -> Result<Arc<PackedRefs>> {
813 let path = self.packed_refs_path();
814 let stamp = Self::packed_refs_stamp(&path);
815 let mut guard = self.packed_refs_cache.lock().map_err(|_| {
816 HeddleError::Config("Failed to acquire packed-refs cache lock".to_string())
817 })?;
818 if let Some(cached) = guard.as_ref()
819 && cached.stamp == stamp
820 {
821 return Ok(cached.packed.clone());
822 }
823 let packed = Arc::new(PackedRefs::load(&path)?);
824 *guard = Some(CachedPackedRefs {
825 stamp,
826 packed: packed.clone(),
827 });
828 Ok(packed)
829 }
830
831 pub(super) fn invalidate_packed_refs_cache(&self) {
834 if let Ok(mut guard) = self.packed_refs_cache.lock() {
835 *guard = None;
836 }
837 }
838
839 pub(super) fn raw_get_remote_thread(
840 &self,
841 remote: &str,
842 thread: &ThreadName,
843 ) -> Result<Option<StateId>> {
844 let path = self.remote_thread_path(remote, thread)?;
845 self.read_state_id_at(&path, "remote thread", &format!("{}/{}", remote, thread))
846 }
847
848 fn raw_list_threads(&self) -> Result<Vec<ThreadName>> {
849 if let Some(summary) = self.try_read_ref_summary_index() {
850 return Ok(summary.thread_names());
851 }
852 self.list_threads_from_storage()
853 }
854
855 fn raw_list_markers(&self) -> Result<Vec<MarkerName>> {
856 if let Some(summary) = self.try_read_ref_summary_index() {
857 return Ok(summary.marker_names());
858 }
859 self.list_markers_from_storage()
860 }
861
862 fn raw_list_remotes(&self) -> Result<Vec<String>> {
863 if let Some(summary) = self.try_read_ref_summary_index() {
864 return Ok(summary.remote_names());
865 }
866 self.list_remotes_from_storage()
867 }
868
869 fn raw_list_remote_threads(&self, remote: &str) -> Result<Vec<ThreadName>> {
870 if let Some(summary) = self.try_read_ref_summary_index() {
871 return Ok(summary.remote_thread_names(remote));
872 }
873 self.list_remote_threads_from_storage(remote)
874 }
875
876 pub fn init(&self) -> Result<()> {
877 create_dir_all_durable(&self.threads_dir())?;
878 create_dir_all_durable(&self.markers_dir())?;
879 create_dir_all_durable(&self.remotes_dir())?;
880 Ok(())
881 }
882
883 pub fn cleanup_stale_temps(&self) {
884 let refs_dir = self.refs_dir();
885 if let Ok(entries) = std::fs::read_dir(&refs_dir) {
886 for entry in entries.flatten() {
887 let path = entry.path();
888 if path
889 .extension()
890 .and_then(|e| e.to_str())
891 .map(|e| e.starts_with("tmp-"))
892 .unwrap_or(false)
893 {
894 let _ = std::fs::remove_file(&path);
895 }
896 }
897 }
898 }
899
900 pub fn read_head(&self) -> Result<Head> {
901 match self.reconciled_load(LoadRequest::Head)? {
902 Loaded::Head(head) => Ok(head),
903 _ => unreachable!("Head request yields Head"),
904 }
905 }
906
907 pub fn write_head(&self, head: &Head) -> Result<()> {
908 self.write_head_cas(RefExpectation::Any, head)
909 }
910
911 pub fn write_head_cas(&self, expected: RefExpectation<Head>, head: &Head) -> Result<()> {
912 self.update_refs(&[RefUpdate::Head {
913 expected,
914 new: head.clone(),
915 }])
916 }
917
918 fn reconciled_point(&self, request: LoadRequest) -> Result<Option<StateId>> {
923 match self.reconciled_load(request)? {
924 Loaded::Point(id) => Ok(id),
925 _ => unreachable!("point request yields Point"),
926 }
927 }
928
929 pub fn get_thread(&self, name: &ThreadName) -> Result<Option<StateId>> {
930 self.reconciled_point(LoadRequest::Thread(name.clone()))
931 }
932
933 pub fn set_thread(&self, name: &ThreadName, state: &StateId) -> Result<()> {
934 self.set_thread_cas(name, RefExpectation::Any, state)
935 }
936
937 pub fn set_thread_cas(
938 &self,
939 name: &ThreadName,
940 expected: RefExpectation<StateId>,
941 state: &StateId,
942 ) -> Result<()> {
943 self.update_refs(&[RefUpdate::Thread {
944 name: name.clone(),
945 expected,
946 new: Some(*state),
947 }])
948 }
949
950 pub fn delete_thread(&self, name: &ThreadName) -> Result<Option<StateId>> {
951 let state = self.get_thread(name)?;
952 if state.is_some() {
953 self.update_refs(&[RefUpdate::Thread {
954 name: name.clone(),
955 expected: RefExpectation::Any,
956 new: None,
957 }])?;
958 }
959 Ok(state)
960 }
961
962 pub fn delete_thread_cas(
963 &self,
964 name: &ThreadName,
965 expected: RefExpectation<StateId>,
966 ) -> Result<()> {
967 self.update_refs(&[RefUpdate::Thread {
968 name: name.clone(),
969 expected,
970 new: None,
971 }])
972 }
973
974 pub fn list_threads(&self) -> Result<Vec<ThreadName>> {
975 match self.reconciled_load(LoadRequest::ThreadList)? {
976 Loaded::ThreadList(names) => Ok(names),
977 _ => unreachable!("ThreadList request yields ThreadList"),
978 }
979 }
980
981 pub fn list_threads_with_states(&self) -> Result<Vec<(ThreadName, StateId)>> {
985 let names = self.list_threads()?;
986 if let Some(summary) = self.try_read_ref_summary_index() {
987 let states = summary.thread_states();
988 if states.len() == names.len()
989 && states
990 .iter()
991 .zip(&names)
992 .all(|((state_name, _), name)| state_name == name)
993 {
994 return Ok(states);
995 }
996 }
997 names
998 .into_iter()
999 .filter_map(|name| match self.raw_get_thread(&name) {
1000 Ok(Some(state)) => Some(Ok((name, state))),
1001 Ok(None) => None,
1002 Err(error) => Some(Err(error)),
1003 })
1004 .collect()
1005 }
1006
1007 pub fn get_marker(&self, name: &MarkerName) -> Result<Option<StateId>> {
1008 self.reconciled_point(LoadRequest::Marker(name.clone()))
1009 }
1010
1011 pub fn create_marker(&self, name: &MarkerName, state: &StateId) -> Result<()> {
1012 self.set_marker_cas(name, RefExpectation::Missing, state)
1013 }
1014
1015 pub fn set_marker_cas(
1016 &self,
1017 name: &MarkerName,
1018 expected: RefExpectation<StateId>,
1019 state: &StateId,
1020 ) -> Result<()> {
1021 self.update_refs(&[RefUpdate::Marker {
1022 name: name.clone(),
1023 expected,
1024 new: Some(*state),
1025 }])
1026 }
1027
1028 pub fn delete_marker(&self, name: &MarkerName) -> Result<Option<StateId>> {
1029 let state = self.get_marker(name)?;
1030 if state.is_some() {
1031 self.delete_marker_cas(name, RefExpectation::Any)?;
1032 }
1033 Ok(state)
1034 }
1035
1036 pub fn delete_marker_cas(
1037 &self,
1038 name: &MarkerName,
1039 expected: RefExpectation<StateId>,
1040 ) -> Result<()> {
1041 self.update_refs(&[RefUpdate::Marker {
1042 name: name.clone(),
1043 expected,
1044 new: None,
1045 }])
1046 }
1047
1048 pub fn list_markers(&self) -> Result<Vec<MarkerName>> {
1049 match self.reconciled_load(LoadRequest::MarkerList)? {
1050 Loaded::MarkerList(names) => Ok(names),
1051 _ => unreachable!("MarkerList request yields MarkerList"),
1052 }
1053 }
1054
1055 pub fn set_undo_recovery(&self, state: &StateId) -> Result<()> {
1061 self.set_undo_recovery_raw(state)
1062 }
1063
1064 fn set_undo_recovery_raw(&self, state: &StateId) -> Result<()> {
1071 self.write_chokepoint(|lock| self.set_undo_recovery_locked(state, lock))
1072 }
1073
1074 fn set_undo_recovery_locked(&self, state: &StateId, _lock: &RefsLock) -> Result<()> {
1078 self.write_string(
1079 &self.undo_recovery_path(),
1080 &super::format_state_id_text(state),
1081 )
1082 }
1083
1084 pub fn clear_undo_recovery(&self) -> Result<()> {
1093 self.write_chokepoint(|lock| self.clear_undo_recovery_locked(lock))
1094 }
1095
1096 fn clear_undo_recovery_locked(&self, _lock: &RefsLock) -> Result<()> {
1097 let path = self.undo_recovery_path();
1098 match std::fs::remove_file(&path) {
1099 Ok(()) => Ok(()),
1100 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
1101 Err(e) => Err(HeddleError::Io(e)),
1102 }
1103 }
1104
1105 pub fn get_undo_recovery(&self) -> Result<Option<StateId>> {
1108 self.reconciled_point(LoadRequest::UndoRecovery)
1109 }
1110
1111 pub fn get_remote_thread(&self, remote: &str, thread: &ThreadName) -> Result<Option<StateId>> {
1112 self.reconciled_point(LoadRequest::RemoteThread {
1113 remote: remote.to_string(),
1114 thread: thread.clone(),
1115 })
1116 }
1117
1118 pub fn set_remote_thread(
1119 &self,
1120 remote: &str,
1121 thread: &ThreadName,
1122 state: &StateId,
1123 ) -> Result<()> {
1124 self.set_remote_thread_raw(remote, thread, state)
1125 }
1126
1127 fn set_remote_thread_raw(
1131 &self,
1132 remote: &str,
1133 thread: &ThreadName,
1134 state: &StateId,
1135 ) -> Result<()> {
1136 self.write_chokepoint(|lock| self.set_remote_thread_locked(remote, thread, state, lock))
1137 }
1138
1139 fn set_remote_thread_locked(
1142 &self,
1143 remote: &str,
1144 thread: &ThreadName,
1145 state: &StateId,
1146 lock: &RefsLock,
1147 ) -> Result<()> {
1148 let path = self.remote_thread_path(remote, thread)?;
1149 let content = format_state_id_text(state);
1150 let parent = path.parent().ok_or_else(|| {
1151 HeddleError::Config(format!(
1152 "invalid remote thread path for {}/{}",
1153 remote, thread
1154 ))
1155 })?;
1156 create_dir_all_durable(parent)?;
1157 self.write_string(&path, &content)?;
1158 if self.rebuild_ref_summary_index_with_lock(lock).is_err() {
1159 self.invalidate_ref_summary_index();
1160 }
1161 Ok(())
1162 }
1163
1164 pub fn delete_remote_thread(
1165 &self,
1166 remote: &str,
1167 thread: &ThreadName,
1168 ) -> Result<Option<StateId>> {
1169 self.delete_remote_thread_raw(remote, thread)
1170 }
1171
1172 fn delete_remote_thread_raw(
1176 &self,
1177 remote: &str,
1178 thread: &ThreadName,
1179 ) -> Result<Option<StateId>> {
1180 self.write_chokepoint(|lock| self.delete_remote_thread_locked(remote, thread, lock))
1181 }
1182
1183 fn delete_remote_thread_locked(
1186 &self,
1187 remote: &str,
1188 thread: &ThreadName,
1189 lock: &RefsLock,
1190 ) -> Result<Option<StateId>> {
1191 let state = self.raw_get_remote_thread(remote, thread)?;
1192 if state.is_some() {
1193 let path = self.remote_thread_path(remote, thread)?;
1194 match std::fs::remove_file(&path) {
1195 Ok(()) => {}
1196 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
1197 Err(e) => return Err(HeddleError::from(e)),
1198 }
1199 }
1200 if self.rebuild_ref_summary_index_with_lock(lock).is_err() {
1201 self.invalidate_ref_summary_index();
1202 }
1203 Ok(state)
1204 }
1205
1206 pub fn list_remotes(&self) -> Result<Vec<String>> {
1207 match self.reconciled_load(LoadRequest::RemoteList)? {
1208 Loaded::RemoteList(names) => Ok(names),
1209 _ => unreachable!("RemoteList request yields RemoteList"),
1210 }
1211 }
1212
1213 pub fn list_remote_threads(&self, remote: &str) -> Result<Vec<ThreadName>> {
1214 match self.reconciled_load(LoadRequest::RemoteThreadList {
1215 remote: remote.to_string(),
1216 })? {
1217 Loaded::RemoteThreadList(names) => Ok(names),
1218 _ => unreachable!("RemoteThreadList request yields RemoteThreadList"),
1219 }
1220 }
1221
1222 pub fn update_refs(&self, updates: &[RefUpdate]) -> Result<()> {
1223 if updates.is_empty() {
1224 return Ok(());
1225 }
1226 self.write_chokepoint(|lock| self.update_refs_with_lock(updates, lock))
1233 }
1234
1235 pub fn resolve(&self, refspec: &str) -> Result<Option<StateId>> {
1236 resolve_refspec(
1237 refspec,
1238 || self.read_head(),
1239 |name| self.get_thread(&ThreadName::new(name)),
1240 |name| self.get_marker(&MarkerName::new(name)),
1241 || self.get_undo_recovery(),
1242 )
1243 }
1244
1245 pub fn pack_refs(&self) -> Result<()> {
1246 let lock = self.lock_refs()?;
1247 let packed_path = self.packed_refs_path();
1248 let mut packed = (*self.load_packed_refs_cached()?).clone();
1249
1250 let threads = self.list_threads_from_storage()?;
1251 for name in &threads {
1252 let path = self.thread_path(name)?;
1253 if let Some(id) = self.read_state_id_at(&path, "thread", name)? {
1254 packed.set_thread(name, id);
1255 }
1256 }
1257 let markers = self.list_markers_from_storage()?;
1258 for name in &markers {
1259 let path = self.marker_path(name)?;
1260 if let Some(id) = self.read_state_id_at(&path, "marker", name)? {
1261 packed.set_marker(name, id);
1262 }
1263 }
1264 if !packed.is_empty() {
1265 packed.save(&packed_path)?;
1266 self.invalidate_packed_refs_cache();
1267 let packed_parent = packed_path
1268 .parent()
1269 .ok_or_else(|| HeddleError::Config("invalid packed-refs path".to_string()))?;
1270 sync_directory(packed_parent)?;
1271 for name in &threads {
1272 let path = self.thread_path(name)?;
1273 if path.exists() {
1274 std::fs::remove_file(&path)?;
1275 }
1276 }
1277 for name in &markers {
1278 let path = self.marker_path(name)?;
1279 if path.exists() {
1280 std::fs::remove_file(&path)?;
1281 }
1282 }
1283 }
1284 if self.rebuild_ref_summary_index_with_lock(&lock).is_err() {
1285 self.invalidate_ref_summary_index();
1286 }
1287 drop(lock);
1288 Ok(())
1289 }
1290}
1291
1292impl CoreRefBackend for RefManager {
1293 type Error = HeddleError;
1294
1295 fn read_head(&self) -> Result<Head> {
1296 RefManager::read_head(self)
1297 }
1298 fn write_head(&self, head: &Head) -> Result<()> {
1299 RefManager::write_head(self, head)
1300 }
1301 fn write_head_cas(&self, expected: RefExpectation<Head>, head: &Head) -> Result<()> {
1302 RefManager::write_head_cas(self, expected, head)
1303 }
1304 async fn get_thread(&self, name: &ThreadName) -> Result<Option<StateId>> {
1305 RefManager::get_thread(self, name)
1306 }
1307 fn set_thread(&self, name: &ThreadName, state: &StateId) -> Result<()> {
1308 RefManager::set_thread(self, name, state)
1309 }
1310 fn set_thread_cas(
1311 &self,
1312 name: &ThreadName,
1313 expected: RefExpectation<StateId>,
1314 state: &StateId,
1315 ) -> Result<()> {
1316 RefManager::set_thread_cas(self, name, expected, state)
1317 }
1318 fn delete_thread(&self, name: &ThreadName) -> Result<Option<StateId>> {
1319 RefManager::delete_thread(self, name)
1320 }
1321 fn delete_thread_cas(
1322 &self,
1323 name: &ThreadName,
1324 expected: RefExpectation<StateId>,
1325 ) -> Result<()> {
1326 RefManager::delete_thread_cas(self, name, expected)
1327 }
1328 fn list_threads(&self) -> Result<Vec<ThreadName>> {
1329 RefManager::list_threads(self)
1330 }
1331 async fn get_marker(&self, name: &MarkerName) -> Result<Option<StateId>> {
1332 RefManager::get_marker(self, name)
1333 }
1334 async fn create_marker(&self, name: &MarkerName, state: &StateId) -> Result<()> {
1335 RefManager::create_marker(self, name, state)
1336 }
1337 fn set_marker_cas(
1338 &self,
1339 name: &MarkerName,
1340 expected: RefExpectation<StateId>,
1341 state: &StateId,
1342 ) -> Result<()> {
1343 RefManager::set_marker_cas(self, name, expected, state)
1344 }
1345 fn delete_marker(&self, name: &MarkerName) -> Result<Option<StateId>> {
1346 RefManager::delete_marker(self, name)
1347 }
1348 fn delete_marker_cas(
1349 &self,
1350 name: &MarkerName,
1351 expected: RefExpectation<StateId>,
1352 ) -> Result<()> {
1353 RefManager::delete_marker_cas(self, name, expected)
1354 }
1355 fn list_markers(&self) -> Result<Vec<MarkerName>> {
1356 RefManager::list_markers(self)
1357 }
1358 fn update_refs(&self, updates: &[RefUpdate]) -> Result<()> {
1359 RefManager::update_refs(self, updates)
1360 }
1361 async fn resolve(&self, refspec: &str) -> Result<Option<StateId>> {
1362 RefManager::resolve(self, refspec)
1363 }
1364}
1365
1366impl RefBackend for RefManager {
1367 fn can_commit_records(&self) -> bool {
1368 self.committer.is_some()
1369 }
1370
1371 fn get_remote_thread(&self, remote: &str, thread: &ThreadName) -> Result<Option<StateId>> {
1372 RefManager::get_remote_thread(self, remote, thread)
1373 }
1374 fn set_remote_thread(&self, remote: &str, thread: &ThreadName, state: &StateId) -> Result<()> {
1375 RefManager::set_remote_thread(self, remote, thread, state)
1376 }
1377 fn delete_remote_thread(&self, remote: &str, thread: &ThreadName) -> Result<Option<StateId>> {
1378 RefManager::delete_remote_thread(self, remote, thread)
1379 }
1380 fn list_remotes(&self) -> Result<Vec<String>> {
1381 RefManager::list_remotes(self)
1382 }
1383 fn list_remote_threads(&self, remote: &str) -> Result<Vec<ThreadName>> {
1384 RefManager::list_remote_threads(self, remote)
1385 }
1386 fn commit_and_publish(
1387 &self,
1388 records: &[OpRecord],
1389 ref_updates: &[RefUpdate],
1390 scope: Option<&str>,
1391 ) -> Result<()> {
1392 RefManager::commit_and_publish(self, records, ref_updates, scope)
1393 }
1394 fn inspect_ref_summary_index(&self) -> Result<super::RefSummaryIndexInspection> {
1395 RefManager::inspect_ref_summary_index(self)
1396 }
1397 fn rebuild_ref_summary_index(&self) -> Result<super::RefSummaryIndexInspection> {
1398 RefManager::rebuild_ref_summary_index(self)
1399 }
1400 fn pack_refs(&self) -> Result<()> {
1401 RefManager::pack_refs(self)
1402 }
1403 fn cleanup_stale_temps(&self) {
1404 RefManager::cleanup_stale_temps(self)
1405 }
1406}