fix: restore orchestrator companion notifications
This commit is contained in:
@@ -10,6 +10,7 @@ use session_store::Store;
|
||||
use ticket::LocalTicketBackend;
|
||||
use ticket::config::TicketConfig;
|
||||
use tokio::sync::{broadcast, mpsc, oneshot};
|
||||
use tracing::{debug, warn};
|
||||
|
||||
use crate::discovery::{PodDiscovery, list_pods_tool, restore_pod_tool, send_to_peer_pod_tool};
|
||||
use crate::feature::FeatureRegistryBuilder;
|
||||
@@ -546,6 +547,30 @@ fn install_ticket_event_companion_notify_hook<C, St>(
|
||||
pod.cwd().to_path_buf(),
|
||||
spawned_registry,
|
||||
);
|
||||
match discovery.ensure_existing_peer(&companion_pod_name) {
|
||||
Ok(Some(_)) => {
|
||||
debug!(
|
||||
companion = %companion_pod_name,
|
||||
orchestrator = %pod.manifest().pod.name,
|
||||
"ensured Companion peer relationship for Orchestrator Ticket event notifications"
|
||||
);
|
||||
}
|
||||
Ok(None) => {
|
||||
debug!(
|
||||
companion = %companion_pod_name,
|
||||
orchestrator = %pod.manifest().pod.name,
|
||||
"Companion metadata is missing; Ticket event notifications will skip until Companion exists"
|
||||
);
|
||||
}
|
||||
Err(error) => {
|
||||
warn!(
|
||||
companion = %companion_pod_name,
|
||||
orchestrator = %pod.manifest().pod.name,
|
||||
error = %error,
|
||||
"failed to ensure Companion peer relationship for Orchestrator Ticket event notifications"
|
||||
);
|
||||
}
|
||||
}
|
||||
pod.add_post_tool_call_hook(TicketEventCompanionNotifyHook::new(
|
||||
LocalTicketBackend::new(backend_root),
|
||||
discovery,
|
||||
|
||||
+104
-21
@@ -7,6 +7,7 @@
|
||||
//! state that exists but is outside that visibility set.
|
||||
|
||||
use std::collections::BTreeMap;
|
||||
use std::fmt;
|
||||
use std::io;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::process::Stdio;
|
||||
@@ -177,6 +178,16 @@ where
|
||||
&self,
|
||||
peer_name: &str,
|
||||
) -> Result<PeerRegistrationResult, PodDiscoveryError> {
|
||||
self.ensure_existing_peer(peer_name)?
|
||||
.ok_or_else(|| PodDiscoveryError::MissingPod {
|
||||
pod_name: peer_name.to_string(),
|
||||
})
|
||||
}
|
||||
|
||||
pub fn ensure_existing_peer(
|
||||
&self,
|
||||
peer_name: &str,
|
||||
) -> Result<Option<PeerRegistrationResult>, PodDiscoveryError> {
|
||||
validate_pod_name(peer_name)?;
|
||||
if peer_name == self.self_pod_name {
|
||||
return Err(PodDiscoveryError::SelfPeer {
|
||||
@@ -191,9 +202,7 @@ where
|
||||
})?;
|
||||
let prior_self_peers = self_metadata.peers.clone();
|
||||
if self.store.read_by_name(peer_name)?.is_none() {
|
||||
return Err(PodDiscoveryError::MissingPod {
|
||||
pod_name: peer_name.to_string(),
|
||||
});
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
self.store.add_peer(&self.self_pod_name, peer_name)?;
|
||||
@@ -202,10 +211,10 @@ where
|
||||
return Err(PodDiscoveryError::PodStore(error));
|
||||
}
|
||||
|
||||
Ok(PeerRegistrationResult {
|
||||
Ok(Some(PeerRegistrationResult {
|
||||
source: self.self_pod_name.clone(),
|
||||
peer: peer_name.to_string(),
|
||||
})
|
||||
}))
|
||||
}
|
||||
|
||||
async fn visibility(&self) -> Result<VisibilitySet, PodDiscoveryError> {
|
||||
@@ -354,16 +363,41 @@ where
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn send_weak_notify_to_live_peer(&self, peer_name: &str, message: String) -> bool {
|
||||
let Ok(detail) = self.inspect(peer_name).await else {
|
||||
return false;
|
||||
pub async fn send_weak_notify_to_live_peer(
|
||||
&self,
|
||||
peer_name: &str,
|
||||
message: String,
|
||||
) -> WeakNotifyDelivery {
|
||||
let detail = match self.inspect(peer_name).await {
|
||||
Ok(detail) => detail,
|
||||
Err(PodDiscoveryError::StateMissing { .. } | PodDiscoveryError::MissingPod { .. }) => {
|
||||
return WeakNotifyDelivery::SkippedMissing;
|
||||
}
|
||||
Err(PodDiscoveryError::NotVisible { .. }) => {
|
||||
return WeakNotifyDelivery::SkippedNotVisible;
|
||||
}
|
||||
Err(error) => {
|
||||
return WeakNotifyDelivery::SendFailed {
|
||||
error: error.to_string(),
|
||||
};
|
||||
}
|
||||
};
|
||||
if detail.visibility != VisibilityReason::Peer || !detail.live.reachable {
|
||||
return false;
|
||||
if detail.visibility != VisibilityReason::Peer {
|
||||
return WeakNotifyDelivery::SkippedNotPeer {
|
||||
visibility: detail.visibility,
|
||||
};
|
||||
}
|
||||
if !detail.live.reachable {
|
||||
return WeakNotifyDelivery::SkippedNotLive {
|
||||
reason: detail.live.error,
|
||||
};
|
||||
}
|
||||
match send_notify(&detail.live.socket_path, message, false).await {
|
||||
Ok(()) => WeakNotifyDelivery::Delivered,
|
||||
Err(error) => WeakNotifyDelivery::SendFailed {
|
||||
error: error.to_string(),
|
||||
},
|
||||
}
|
||||
send_notify(&detail.live.socket_path, message, false)
|
||||
.await
|
||||
.is_ok()
|
||||
}
|
||||
|
||||
async fn live_for_name(&self, pod_name: &str, socket_override: Option<&Path>) -> LiveInfo {
|
||||
@@ -585,6 +619,50 @@ pub struct PeerRegistrationResult {
|
||||
pub peer: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub enum WeakNotifyDelivery {
|
||||
Delivered,
|
||||
SkippedMissing,
|
||||
SkippedNotVisible,
|
||||
SkippedNotPeer { visibility: VisibilityReason },
|
||||
SkippedNotLive { reason: Option<String> },
|
||||
SendFailed { error: String },
|
||||
}
|
||||
|
||||
impl WeakNotifyDelivery {
|
||||
pub fn delivered(&self) -> bool {
|
||||
matches!(self, WeakNotifyDelivery::Delivered)
|
||||
}
|
||||
}
|
||||
|
||||
impl fmt::Display for WeakNotifyDelivery {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
match self {
|
||||
WeakNotifyDelivery::Delivered => write!(f, "delivered"),
|
||||
WeakNotifyDelivery::SkippedMissing => {
|
||||
write!(f, "skipped: target pod metadata is missing")
|
||||
}
|
||||
WeakNotifyDelivery::SkippedNotVisible => {
|
||||
write!(f, "skipped: target pod is not visible")
|
||||
}
|
||||
WeakNotifyDelivery::SkippedNotPeer { visibility } => {
|
||||
write!(
|
||||
f,
|
||||
"skipped: target pod is visible as {visibility:?}, not peer"
|
||||
)
|
||||
}
|
||||
WeakNotifyDelivery::SkippedNotLive { reason } => {
|
||||
if let Some(reason) = reason {
|
||||
write!(f, "skipped: target peer is not live/reachable ({reason})")
|
||||
} else {
|
||||
write!(f, "skipped: target peer is not live/reachable")
|
||||
}
|
||||
}
|
||||
WeakNotifyDelivery::SendFailed { error } => write!(f, "send failed: {error}"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, thiserror::Error)]
|
||||
pub enum PodDiscoveryError {
|
||||
#[error("pod state missing for `{pod_name}`")]
|
||||
@@ -1524,18 +1602,20 @@ mod tests {
|
||||
}
|
||||
});
|
||||
|
||||
assert!(
|
||||
assert_eq!(
|
||||
discovery
|
||||
.send_weak_notify_to_live_peer("target", "weak event".into())
|
||||
.await
|
||||
.await,
|
||||
WeakNotifyDelivery::Delivered
|
||||
);
|
||||
assert_eq!(rx.recv().await.unwrap(), "weak event");
|
||||
target.await.unwrap();
|
||||
|
||||
assert!(
|
||||
!discovery
|
||||
assert_eq!(
|
||||
discovery
|
||||
.send_weak_notify_to_live_peer("missing", "no-op".into())
|
||||
.await
|
||||
.await,
|
||||
WeakNotifyDelivery::SkippedMissing
|
||||
);
|
||||
}
|
||||
|
||||
@@ -1567,10 +1647,13 @@ mod tests {
|
||||
SpawnedPodRegistry::new(runtime_dir),
|
||||
);
|
||||
|
||||
assert!(
|
||||
!discovery
|
||||
assert_eq!(
|
||||
discovery
|
||||
.send_weak_notify_to_live_peer("target", "must not send".into())
|
||||
.await
|
||||
.await,
|
||||
WeakNotifyDelivery::SkippedNotPeer {
|
||||
visibility: VisibilityReason::SpawnedChild
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
@@ -5,9 +5,9 @@ use minijinja::Value as TemplateValue;
|
||||
use serde_json::Value;
|
||||
use std::collections::BTreeMap;
|
||||
use ticket::{LocalTicketBackend, TicketBackend, TicketIdOrSlug};
|
||||
use tracing::debug;
|
||||
use tracing::{debug, warn};
|
||||
|
||||
use crate::discovery::PodDiscovery;
|
||||
use crate::discovery::{PodDiscovery, WeakNotifyDelivery};
|
||||
use crate::hook::{Hook, HookPostToolAction, PostToolCall, ToolResultSummary};
|
||||
use crate::prompt::catalog::{PodPrompt, PromptCatalog};
|
||||
use pod_store::PodMetadataStore;
|
||||
@@ -48,17 +48,58 @@ impl<St: PodMetadataStore + Clone + Send + Sync + 'static> Hook<PostToolCall>
|
||||
let Some(notice) = build_ticket_event_notice(&self.backend, summary) else {
|
||||
return HookPostToolAction::Continue;
|
||||
};
|
||||
let delivered = self
|
||||
match self
|
||||
.discovery
|
||||
.ensure_existing_peer(&self.companion_pod_name)
|
||||
{
|
||||
Ok(Some(_)) => {
|
||||
debug!(
|
||||
ticket = %notice.ticket_id,
|
||||
event_kind = %notice.event_kind,
|
||||
companion = %self.companion_pod_name,
|
||||
"ensured Companion peer relationship before Ticket event notification"
|
||||
);
|
||||
}
|
||||
Ok(None) => {
|
||||
debug!(
|
||||
ticket = %notice.ticket_id,
|
||||
event_kind = %notice.event_kind,
|
||||
companion = %self.companion_pod_name,
|
||||
"skipping Companion peer registration because Companion metadata is missing"
|
||||
);
|
||||
}
|
||||
Err(error) => {
|
||||
warn!(
|
||||
ticket = %notice.ticket_id,
|
||||
event_kind = %notice.event_kind,
|
||||
companion = %self.companion_pod_name,
|
||||
error = %error,
|
||||
"failed to ensure Companion peer relationship before Ticket event notification"
|
||||
);
|
||||
}
|
||||
}
|
||||
let delivery = self
|
||||
.discovery
|
||||
.send_weak_notify_to_live_peer(&self.companion_pod_name, notice.message)
|
||||
.await;
|
||||
if delivered {
|
||||
debug!(
|
||||
ticket = %notice.ticket_id,
|
||||
event_kind = %notice.event_kind,
|
||||
companion = %self.companion_pod_name,
|
||||
"delivered weak Ticket event notification to Companion peer"
|
||||
);
|
||||
match delivery {
|
||||
WeakNotifyDelivery::Delivered => {
|
||||
debug!(
|
||||
ticket = %notice.ticket_id,
|
||||
event_kind = %notice.event_kind,
|
||||
companion = %self.companion_pod_name,
|
||||
"delivered weak Ticket event notification to Companion peer"
|
||||
);
|
||||
}
|
||||
skipped => {
|
||||
warn!(
|
||||
ticket = %notice.ticket_id,
|
||||
event_kind = %notice.event_kind,
|
||||
companion = %self.companion_pod_name,
|
||||
delivery = %skipped,
|
||||
"skipped weak Ticket event notification to Companion peer"
|
||||
);
|
||||
}
|
||||
}
|
||||
HookPostToolAction::Continue
|
||||
}
|
||||
@@ -327,7 +368,7 @@ mod tests {
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn ticket_event_hook_delivers_weak_companion_notification() {
|
||||
async fn ticket_event_hook_ensures_peer_and_delivers_weak_companion_notification() {
|
||||
let root = tempdir().expect("tempdir");
|
||||
let runtime_base = root.path().join("runtime");
|
||||
let store_dir = root.path().join("store");
|
||||
@@ -339,9 +380,7 @@ mod tests {
|
||||
active: None,
|
||||
spawned_children: Vec::new(),
|
||||
reclaimed_children: Vec::new(),
|
||||
peers: vec![pod_store::PodPeer {
|
||||
pod_name: "companion".into(),
|
||||
}],
|
||||
peers: Vec::new(),
|
||||
resolved_manifest_snapshot: None,
|
||||
})
|
||||
.unwrap();
|
||||
@@ -351,9 +390,7 @@ mod tests {
|
||||
active: None,
|
||||
spawned_children: Vec::new(),
|
||||
reclaimed_children: Vec::new(),
|
||||
peers: vec![pod_store::PodPeer {
|
||||
pod_name: "orchestrator".into(),
|
||||
}],
|
||||
peers: Vec::new(),
|
||||
resolved_manifest_snapshot: None,
|
||||
})
|
||||
.unwrap();
|
||||
@@ -363,6 +400,7 @@ mod tests {
|
||||
.await
|
||||
.unwrap(),
|
||||
);
|
||||
let store_for_assert = store.clone();
|
||||
let hook = TicketEventCompanionNotifyHook::new(
|
||||
backend,
|
||||
PodDiscovery::new(
|
||||
@@ -448,6 +486,15 @@ mod tests {
|
||||
let message = rx.recv().await.unwrap();
|
||||
assert!(message.contains("event: state/queued->inprogress"));
|
||||
assert!(message.contains("title: Companion event hook"));
|
||||
let orchestrator = store_for_assert
|
||||
.read_by_name("orchestrator")
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(orchestrator.peers.len(), 1);
|
||||
assert_eq!(orchestrator.peers[0].pod_name, "companion");
|
||||
let companion_metadata = store_for_assert.read_by_name("companion").unwrap().unwrap();
|
||||
assert_eq!(companion_metadata.peers.len(), 1);
|
||||
assert_eq!(companion_metadata.peers[0].pod_name, "orchestrator");
|
||||
companion.await.unwrap();
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user