(
&self,
session_id: &str,
)
| 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; |
no test coverage detected