#![allow(non_snake_case)]
use async_trait::async_trait;
use observation::{Event, EventScope, EventSource, Level, TaskRunId};
use yah_scryer::federation::{
DenyAllAcl, FederationAcl, FederationError, FederationPeer, FederationRule, OperatorTagAcl,
PeerIdentity, ScopedEvent, federated_events,
};
use yah_scryer::service::{EventFilter, Scryer, ScryerConfig};
use serde_json::json;
use std::sync::Arc;
use tempfile::TempDir;
fn make_event(run_id: &TaskRunId, seq: u32, offset_ms: u32) -> Event {
Event {
run_id: run_id.clone(),
seq,
offset_ms,
level: Level::Info,
target: "federation::test".to_string(),
msg: format!("event at offset {offset_ms}"),
fields: json!({}),
anchor: None,
source: EventSource::Synth,
}
}
fn make_scoped(run_id: &TaskRunId, seq: u32, offset_ms: u32) -> ScopedEvent {
ScopedEvent::new(
EventScope::TaskRun(run_id.clone()),
make_event(run_id, seq, offset_ms),
)
}
struct MockPeer {
peer_name: String,
events: Vec<ScopedEvent>,
}
#[async_trait]
impl FederationPeer for MockPeer {
fn name(&self) -> &str {
&self.peer_name
}
async fn events(&self, _filter: &EventFilter) -> Result<Vec<ScopedEvent>, FederationError> {
Ok(self.events.clone())
}
}
struct FailingPeer;
#[async_trait]
impl FederationPeer for FailingPeer {
fn name(&self) -> &str {
"failing-peer"
}
async fn events(&self, _filter: &EventFilter) -> Result<Vec<ScopedEvent>, FederationError> {
Err(FederationError::Rpc("simulated gRPC error".to_string()))
}
}
#[tokio::test]
async fn scryer_federation__local() {
let run_id = TaskRunId::new();
let local = vec![
make_scoped(&run_id, 0, 1),
make_scoped(&run_id, 2, 3),
make_scoped(&run_id, 4, 5),
];
let peer_b: Arc<dyn FederationPeer> = Arc::new(MockPeer {
peer_name: "yubaba-pdx-2".to_string(),
events: vec![
make_scoped(&run_id, 1, 2),
make_scoped(&run_id, 3, 4),
make_scoped(&run_id, 5, 6),
],
});
let filter = EventFilter::default();
let result = federated_events(local, &[peer_b], &filter, &FederationRule::All).await;
assert_eq!(result.len(), 6, "3 local + 3 from peer B");
let offsets: Vec<u32> = result.iter().map(|r| r.event.offset_ms).collect();
assert_eq!(offsets, vec![1, 2, 3, 4, 5, 6], "events must be interleaved in time order");
assert!(
result.iter().all(|r| r.scope == EventScope::TaskRun(run_id.clone())),
"scope envelope must survive the federated merge",
);
}
#[tokio::test]
async fn scryer_federated_events_method() {
let dir = TempDir::new().unwrap();
let cfg = ScryerConfig::new(dir.path().join("events.db"));
let scryer = Scryer::new(cfg, None).unwrap();
let run_id = TaskRunId::new();
let scope = EventScope::TaskRun(run_id.clone());
for (i, offset) in [(0u32, 10u32), (2, 30), (4, 50)] {
scryer.push(scope.clone(), make_event(&run_id, i, offset)).unwrap();
}
scryer.flush_ring().unwrap();
let peer_b: Arc<dyn FederationPeer> = Arc::new(MockPeer {
peer_name: "peer-b".to_string(),
events: vec![
make_scoped(&run_id, 1, 20),
make_scoped(&run_id, 3, 40),
make_scoped(&run_id, 5, 60),
],
});
let filter = EventFilter::default();
let result = scryer
.federated_events(&scope, &filter, &FederationRule::All, &[peer_b])
.await
.unwrap();
assert_eq!(result.len(), 6);
let offsets: Vec<u32> = result.iter().map(|r| r.event.offset_ms).collect();
assert_eq!(offsets, vec![10, 20, 30, 40, 50, 60]);
assert!(result.iter().all(|r| r.scope == scope));
}
#[tokio::test]
async fn scryer_federation__failing_peer_is_skipped() {
let run_id = TaskRunId::new();
let local = vec![make_scoped(&run_id, 0, 1)];
let peer: Arc<dyn FederationPeer> = Arc::new(FailingPeer);
let filter = EventFilter::default();
let result = federated_events(local, &[peer], &filter, &FederationRule::All).await;
assert_eq!(result.len(), 1, "local event survives; failing peer is ignored");
}
#[tokio::test]
async fn scryer_federation__tag_rule_filters_peers() {
let run_id = TaskRunId::new();
let local = vec![make_scoped(&run_id, 0, 1)];
let matching_peer: Arc<dyn FederationPeer> = Arc::new(MockPeer {
peer_name: "yubaba-tier=public-1".to_string(),
events: vec![make_scoped(&run_id, 1, 2)],
});
let excluded_peer: Arc<dyn FederationPeer> = Arc::new(MockPeer {
peer_name: "yubaba-private-1".to_string(),
events: vec![make_scoped(&run_id, 2, 3)],
});
let filter = EventFilter::default();
let result = federated_events(
local,
&[matching_peer, excluded_peer],
&filter,
&FederationRule::Tag("tier=public".to_string()),
)
.await;
assert_eq!(result.len(), 2, "local + matching peer only");
assert_eq!(result[0].event.offset_ms, 1);
assert_eq!(result[1].event.offset_ms, 2);
}
#[test]
fn scryer_acl__unauthorized_peer_rejected() {
let acl = OperatorTagAcl;
let identity = PeerIdentity::default(); assert!(!acl.is_authorized(&identity), "untagged peer must be rejected");
}
#[test]
fn scryer_acl__operator_tag_is_allowed() {
let acl = OperatorTagAcl;
let identity = PeerIdentity::default().with_tag("tag:operator");
assert!(acl.is_authorized(&identity));
}
#[test]
fn scryer_acl__deny_all_blocks_operator() {
let acl = DenyAllAcl;
let identity = PeerIdentity::default().with_tag("tag:operator");
assert!(!acl.is_authorized(&identity), "DenyAllAcl must reject all");
}