use crate::service::{EventFilter, ScryerError};
use async_trait::async_trait;
use observation::{Event, EventScope};
use serde::{Deserialize, Serialize};
use std::sync::Arc;
use thiserror::Error;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ScopedEvent {
pub scope: EventScope,
pub event: Event,
}
impl ScopedEvent {
pub fn new(scope: EventScope, event: Event) -> Self {
Self { scope, event }
}
pub fn tag_all(scope: &EventScope, events: Vec<Event>) -> Vec<Self> {
events
.into_iter()
.map(|event| Self::new(scope.clone(), event))
.collect()
}
pub fn into_events(rows: Vec<Self>) -> Vec<Event> {
rows.into_iter().map(|r| r.event).collect()
}
}
#[derive(Debug, Clone)]
pub enum FederationRule {
All,
Tag(String),
}
#[derive(Debug, Clone, Default)]
pub struct PeerIdentity {
pub tags: Vec<String>,
}
impl PeerIdentity {
pub fn with_tag(mut self, tag: impl Into<String>) -> Self {
self.tags.push(tag.into());
self
}
}
pub trait FederationAcl: Send + Sync {
fn is_authorized(&self, identity: &PeerIdentity) -> bool;
}
pub struct OperatorTagAcl;
impl FederationAcl for OperatorTagAcl {
fn is_authorized(&self, identity: &PeerIdentity) -> bool {
identity.tags.iter().any(|t| t == "tag:operator")
}
}
pub struct DenyAllAcl;
impl FederationAcl for DenyAllAcl {
fn is_authorized(&self, _: &PeerIdentity) -> bool {
false
}
}
#[derive(Debug, Error)]
pub enum FederationError {
#[error("rpc: {0}")]
Rpc(String),
#[error("unauthorized: operator tag required for federated queries")]
Unauthorized,
#[error("{0}")]
Scryer(#[from] ScryerError),
}
#[async_trait]
pub trait FederationPeer: Send + Sync {
fn name(&self) -> &str;
async fn events(&self, filter: &EventFilter) -> Result<Vec<ScopedEvent>, FederationError>;
}
pub trait TimeOrdered {
fn order_key(&self) -> (u32, u32);
}
impl TimeOrdered for Event {
fn order_key(&self) -> (u32, u32) {
(self.offset_ms, self.seq)
}
}
impl TimeOrdered for ScopedEvent {
fn order_key(&self) -> (u32, u32) {
(self.event.offset_ms, self.event.seq)
}
}
pub fn merge_ordered<T: TimeOrdered>(a: Vec<T>, b: Vec<T>) -> Vec<T> {
let mut result = Vec::with_capacity(a.len() + b.len());
let mut ia = a.into_iter().peekable();
let mut ib = b.into_iter().peekable();
loop {
match (ia.peek(), ib.peek()) {
(None, None) => break,
(Some(_), None) => result.push(ia.next().unwrap()),
(None, Some(_)) => result.push(ib.next().unwrap()),
(Some(ea), Some(eb)) => {
if ea.order_key() <= eb.order_key() {
result.push(ia.next().unwrap());
} else {
result.push(ib.next().unwrap());
}
}
}
}
result
}
pub fn merge_events(a: Vec<Event>, b: Vec<Event>) -> Vec<Event> {
merge_ordered(a, b)
}
fn peer_matches_rule(peer: &dyn FederationPeer, rule: &FederationRule) -> bool {
match rule {
FederationRule::All => true,
FederationRule::Tag(tag) => peer.name().contains(tag.as_str()),
}
}
pub async fn federated_events(
local_events: Vec<ScopedEvent>,
peers: &[Arc<dyn FederationPeer>],
filter: &EventFilter,
rule: &FederationRule,
) -> Vec<ScopedEvent> {
let mut result = local_events;
for peer in peers {
if peer_matches_rule(peer.as_ref(), rule) {
if let Ok(remote) = peer.events(filter).await {
result = merge_ordered(result, remote);
}
}
}
result
}
#[cfg(test)]
mod acl {
use super::*;
#[test]
fn unauthorized_peer_is_rejected() {
let acl = OperatorTagAcl;
let identity = PeerIdentity::default(); assert!(!acl.is_authorized(&identity));
}
#[test]
fn operator_tagged_peer_is_allowed() {
let acl = OperatorTagAcl;
let identity = PeerIdentity::default().with_tag("tag:operator");
assert!(acl.is_authorized(&identity));
}
#[test]
fn non_operator_tag_is_rejected() {
let acl = OperatorTagAcl;
let identity = PeerIdentity::default().with_tag("tag:workload");
assert!(!acl.is_authorized(&identity));
}
#[test]
fn deny_all_rejects_any_identity() {
let acl = DenyAllAcl;
let identity = PeerIdentity::default().with_tag("tag:operator");
assert!(!acl.is_authorized(&identity));
}
}
#[cfg(test)]
mod merge_tests {
use super::*;
use observation::{EventSource, Level, TaskRunId};
use serde_json::json;
fn ev(offset_ms: u32, seq: u32) -> Event {
Event {
run_id: TaskRunId::new(),
seq,
offset_ms,
level: Level::Info,
target: "t".to_string(),
msg: format!("e{offset_ms}"),
fields: json!({}),
anchor: None,
source: EventSource::Synth,
}
}
#[test]
fn merge_interleaves_by_offset() {
let a = vec![ev(1, 0), ev(3, 1), ev(5, 2)];
let b = vec![ev(2, 0), ev(4, 1), ev(6, 2)];
let result = merge_events(a, b);
let offsets: Vec<u32> = result.iter().map(|e| e.offset_ms).collect();
assert_eq!(offsets, vec![1, 2, 3, 4, 5, 6]);
}
#[test]
fn merge_empty_inputs() {
assert!(merge_events(vec![], vec![]).is_empty());
let a = vec![ev(1, 0)];
assert_eq!(merge_events(a.clone(), vec![]).len(), 1);
assert_eq!(merge_events(vec![], a).len(), 1);
}
#[test]
fn merge_tie_breaks_by_seq() {
let a = vec![ev(10, 1)];
let b = vec![ev(10, 0)];
let result = merge_events(a, b);
assert_eq!(result[0].seq, 0);
assert_eq!(result[1].seq, 1);
}
}