MCPcopy Create free account
hub / github.com/atomicdotdev/atomic / wait_turn_publication_lock

Method wait_turn_publication_lock

atomic-agent/src/turn/orchestrator/turn.rs:477–519  ·  view source on GitHub ↗
(
        &self,
        session_id: &str,
    )

Source from the content-addressed store, hash-verified

475 }
476
477 fn wait_turn_publication_lock(
478 &self,
479 session_id: &str,
480 ) -> AgentResult<Option<TurnEndLockGuard>> {
481 use fs2::FileExt;
482 let failure = |reason: String| AgentError::ProvenanceJournalFailed {
483 session_id: session_id.to_owned(),
484 reason,
485 };
486 let canonical = match atomic_repository::Repository::canonical_dot_dir(&self.repo_root) {
487 Ok(path) => path,
488 // Legacy session-only orchestrators can run without a repository.
489 // A journal-backed Stop must always have canonical coordination.
490 Err(
491 atomic_repository::RepositoryError::NotFound { .. }
492 | atomic_repository::RepositoryError::NotInRepository,
493 ) if self.journal_sink.is_none() => return Ok(None),
494 Err(error) => return Err(failure(error.to_string())),
495 };
496 // Never unlink a lock file: waiters must all lock the same inode.
497 let file = std::fs::OpenOptions::new()
498 .create(true)
499 .write(true)
500 .truncate(false)
501 .open(canonical.join("turn-publication.lock"))
502 .map_err(|error| failure(error.to_string()))?;
503 let start = std::time::Instant::now();
504 let timeout = std::time::Duration::from_secs(10);
505 loop {
506 match file.try_lock_exclusive() {
507 Ok(()) => return Ok(Some(TurnEndLockGuard { file })),
508 Err(error) if super::is_lock_contended(&error) => {
509 if start.elapsed() >= timeout {
510 return Err(failure(
511 "timed out waiting for another Stop to publish; retry this Stop".into(),
512 ));
513 }
514 std::thread::sleep(std::time::Duration::from_millis(10));
515 }
516 Err(error) => return Err(failure(error.to_string())),
517 }
518 }
519 }
520
521 fn try_turn_end_lock(&self, session_id: &str) -> TurnEndLock {
522 use fs2::FileExt;

Callers 1

handle_turn_endMethod · 0.80

Calls 6

is_noneMethod · 0.80
truncateMethod · 0.80
is_lock_contendedFunction · 0.70
openMethod · 0.45
writeMethod · 0.45
createMethod · 0.45

Tested by

no test coverage detected