use super::*;
#[derive(Clone)]
pub struct TurnReviewHost {
pub(super) events: mpsc::UnboundedSender<HostEvent>,
pub(super) shared: Arc<HostShared>,
}
pub(super) struct HostShared {
pub(super) views: Mutex<BTreeMap<String, RuntimeReviewView>>,
pub(super) changed: Arc<dyn Fn() + Send + Sync>,
pub(super) shutdown: tokio::sync::OnceCell<Result<(), String>>,
}
impl std::fmt::Debug for TurnReviewHost {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str("TurnReviewHost")
}
}
impl TurnReviewHost {
#[must_use]
pub fn spawn(control: SessionManagerControl, config: ReviewConfigSource) -> Self {
Self::spawn_notifying(control, config, Arc::new(|| {}))
}
#[must_use]
pub fn spawn_notifying(
control: SessionManagerControl,
config: ReviewConfigSource,
changed: Arc<dyn Fn() + Send + Sync>,
) -> Self {
Self::spawn_in_notifying(control, config, Arc::new(ControllerEnvironment), changed)
}
#[must_use]
pub fn spawn_in(
control: SessionManagerControl,
config: ReviewConfigSource,
environment: Arc<dyn ReviewEnvironment>,
) -> Self {
Self::spawn_in_notifying(control, config, environment, Arc::new(|| {}))
}
#[must_use]
pub(super) fn spawn_in_notifying(
control: SessionManagerControl,
config: ReviewConfigSource,
environment: Arc<dyn ReviewEnvironment>,
changed: Arc<dyn Fn() + Send + Sync>,
) -> Self {
let (events, receiver) = mpsc::unbounded_channel();
let (persistence, persistence_receiver) = mpsc::unbounded_channel();
let shared = Arc::new(HostShared {
views: Mutex::default(),
changed,
shutdown: tokio::sync::OnceCell::new(),
});
let host = Self {
events: events.clone(),
shared: shared.clone(),
};
let persistence_task = tokio::spawn(persistence_loop(
environment.clone(),
events.clone(),
persistence_receiver,
));
persistence
.send(PersistenceRequest::SweepInterrupted)
.expect("new review persistence lane accepts its initial sweep");
tokio::spawn(host_loop(
HostState {
control,
config,
environment,
shared,
events,
persistence: Some(persistence),
persistence_task: Some(persistence_task),
reviews: BTreeMap::new(),
preparing: BTreeSet::new(),
pending_open: BTreeMap::new(),
closing: BTreeSet::new(),
awaiting_forward_persistence: BTreeMap::new(),
next_epoch: 0,
sessions: BTreeMap::new(),
missing_reviewer_reported: BTreeSet::new(),
recovery_candidates: BTreeSet::new(),
recovery_in_flight: BTreeSet::new(),
},
receiver,
));
host
}
pub fn retain_sessions(&self, live: std::collections::BTreeSet<String>) {
let _ = self.events.send(HostEvent::Retain { live });
}
pub fn observe(&self, session_id: &str, view: &ManagedSessionView) {
let _ = self.events.send(HostEvent::View {
session_id: session_id.to_owned(),
snapshot: view
.snapshot
.as_ref()
.map(|snapshot| Box::new(snapshot.materialized.clone())),
prompt_driven: view
.snapshot
.as_ref()
.is_some_and(|snapshot| snapshot.operational.active_prompt.is_some()),
});
}
pub async fn start(&self, session_id: &str, manual: bool) -> Result<(), StartRefusal> {
let (reply, answer) = oneshot::channel();
self.events
.send(HostEvent::Start {
session_id: session_id.to_owned(),
manual,
reply: Some(reply),
})
.map_err(|_| StartRefusal("the review host stopped".to_owned()))?;
answer
.await
.map_err(|_| StartRefusal("the review host stopped".to_owned()))?
}
pub async fn resolve(&self, session_id: &str, resolution: Resolution) -> Result<(), String> {
let (reply, answer) = oneshot::channel();
self.events
.send(HostEvent::Resolve {
session_id: session_id.to_owned(),
resolution,
reply,
})
.map_err(|_| "the review host stopped".to_owned())?;
answer
.await
.map_err(|_| "the review host stopped".to_owned())?
}
pub async fn shutdown(&self) -> Result<(), String> {
self.shared
.shutdown
.get_or_init(|| async {
let (reply, answer) = oneshot::channel();
self.events
.send(HostEvent::Shutdown { reply })
.map_err(|_| "the review host stopped".to_owned())?;
answer
.await
.map_err(|_| "the review host stopped during shutdown".to_owned())?
})
.await
.clone()
}
#[must_use]
pub fn views(&self) -> Vec<RuntimeReviewView> {
self.shared
.views
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.values()
.cloned()
.collect()
}
#[must_use]
pub fn view(&self, session_id: &str) -> Option<RuntimeReviewView> {
self.shared
.views
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(session_id)
.cloned()
}
#[must_use]
pub fn refuses_prompt(&self, session_id: &str) -> bool {
prompt_refusal(session_id).is_some()
}
}