use std::sync::{Arc, RwLock};
use clawless_core::event::{Event, EventReceiver};
use tokio::task::JoinHandle;
pub use self::entry::Entry;
use self::state::ProjectionState;
mod entry;
mod state;
#[derive(Debug)]
pub struct Projection {
state: Arc<RwLock<ProjectionState>>,
_drain_handle: JoinHandle<()>,
}
impl Projection {
pub fn new(receiver: EventReceiver) -> Self {
let state = Arc::new(RwLock::new(ProjectionState::default()));
let drain_state = Arc::clone(&state);
let handle = tokio::spawn(drain(receiver, drain_state));
Self {
state,
_drain_handle: handle,
}
}
#[allow(clippy::expect_used)]
pub fn entries(&self) -> Vec<Entry> {
self.state.read().expect("lock poisoned").entries()
}
#[allow(clippy::expect_used)]
pub fn messages(&self) -> Vec<Entry> {
self.state.read().expect("lock poisoned").messages()
}
#[allow(clippy::expect_used)]
pub fn details(&self) -> Vec<Entry> {
self.state.read().expect("lock poisoned").details()
}
#[allow(clippy::expect_used)]
pub fn artifacts(&self) -> Vec<Entry> {
self.state.read().expect("lock poisoned").artifacts()
}
#[allow(clippy::expect_used)]
pub fn is_complete(&self) -> bool {
self.state.read().expect("lock poisoned").is_complete()
}
}
#[allow(clippy::expect_used)]
async fn drain(mut receiver: EventReceiver, state: Arc<RwLock<ProjectionState>>) {
while let Some(event) = receiver.recv().await {
let entry = match event {
Event::Message(text) => Entry::Message(text),
Event::Detail(text) => Entry::Detail(text),
Event::Artifact(artifact) => Entry::Artifact(Arc::from(artifact)),
};
state.write().expect("lock poisoned").push(entry);
}
state.write().expect("lock poisoned").set_complete();
}
#[cfg(test)]
mod tests {
#![allow(clippy::missing_panics_doc)]
use std::fmt;
use clawless_core::event::event_channel;
use serde::Serialize;
use super::*;
#[derive(Clone, Debug, Serialize)]
struct TestArtifact {
value: String,
}
impl fmt::Display for TestArtifact {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "{}", self.value)
}
}
#[tokio::test]
async fn artifacts_returns_artifact_entries_only() {
let (sender, receiver) = event_channel();
let projection = Projection::new(receiver);
sender
.send(Event::Message("msg".to_string()))
.await
.expect("should send");
sender
.send(Event::Artifact(Box::new(TestArtifact {
value: "art".to_string(),
})))
.await
.expect("should send");
drop(sender);
tokio::task::yield_now().await;
let artifacts = projection.artifacts();
assert_eq!(artifacts.len(), 1);
let Entry::Artifact(a) = &artifacts[0] else {
panic!("expected Entry::Artifact");
};
assert_eq!(a.to_string(), "art");
}
#[tokio::test]
async fn details_returns_detail_entries_only() {
let (sender, receiver) = event_channel();
let projection = Projection::new(receiver);
sender
.send(Event::Detail("dtl".to_string()))
.await
.expect("should send");
sender
.send(Event::Message("msg".to_string()))
.await
.expect("should send");
sender
.send(Event::Detail("dtl2".to_string()))
.await
.expect("should send");
drop(sender);
tokio::task::yield_now().await;
let details = projection.details();
assert_eq!(details.len(), 2);
let Entry::Detail(s) = &details[0] else {
panic!("expected Entry::Detail");
};
assert_eq!(s, "dtl");
let Entry::Detail(s) = &details[1] else {
panic!("expected Entry::Detail");
};
assert_eq!(s, "dtl2");
}
#[tokio::test]
async fn drain_processes_buffered_events_before_completing() {
let (sender, receiver) = event_channel();
sender
.send(Event::Message("buffered".to_string()))
.await
.expect("should send");
drop(sender);
let projection = Projection::new(receiver);
tokio::task::yield_now().await;
assert!(projection.is_complete());
let entries = projection.entries();
assert_eq!(entries.len(), 1);
let Entry::Message(s) = &entries[0] else {
panic!("expected Entry::Message");
};
assert_eq!(s, "buffered");
}
#[tokio::test]
async fn entries_returns_artifact_events_as_entries() {
let (sender, receiver) = event_channel();
let projection = Projection::new(receiver);
sender
.send(Event::Artifact(Box::new(TestArtifact {
value: "result".to_string(),
})))
.await
.expect("should send");
drop(sender);
tokio::task::yield_now().await;
let entries = projection.entries();
assert_eq!(entries.len(), 1);
let Entry::Artifact(a) = &entries[0] else {
panic!("expected Entry::Artifact");
};
assert_eq!(a.to_string(), "result");
}
#[tokio::test]
async fn entries_returns_detail_events_as_entries() {
let (sender, receiver) = event_channel();
let projection = Projection::new(receiver);
sender
.send(Event::Detail("info".to_string()))
.await
.expect("should send");
drop(sender);
tokio::task::yield_now().await;
let entries = projection.entries();
assert_eq!(entries.len(), 1);
let Entry::Detail(s) = &entries[0] else {
panic!("expected Entry::Detail");
};
assert_eq!(s, "info");
}
#[tokio::test]
async fn entries_returns_message_events_as_entries() {
let (sender, receiver) = event_channel();
let projection = Projection::new(receiver);
sender
.send(Event::Message("hello".to_string()))
.await
.expect("should send");
drop(sender);
tokio::task::yield_now().await;
let entries = projection.entries();
assert_eq!(entries.len(), 1);
let Entry::Message(s) = &entries[0] else {
panic!("expected Entry::Message");
};
assert_eq!(s, "hello");
}
#[tokio::test]
async fn entries_preserves_receive_order() {
let (sender, receiver) = event_channel();
let projection = Projection::new(receiver);
sender
.send(Event::Message("first".to_string()))
.await
.expect("should send");
sender
.send(Event::Detail("second".to_string()))
.await
.expect("should send");
sender
.send(Event::Message("third".to_string()))
.await
.expect("should send");
drop(sender);
tokio::task::yield_now().await;
let entries = projection.entries();
assert_eq!(entries.len(), 3);
let Entry::Message(s) = &entries[0] else {
panic!("expected Entry::Message");
};
assert_eq!(s, "first");
let Entry::Detail(s) = &entries[1] else {
panic!("expected Entry::Detail");
};
assert_eq!(s, "second");
let Entry::Message(s) = &entries[2] else {
panic!("expected Entry::Message");
};
assert_eq!(s, "third");
}
#[tokio::test]
async fn is_complete_returns_false_while_channel_open() {
let (_sender, receiver) = event_channel();
let projection = Projection::new(receiver);
tokio::task::yield_now().await;
assert!(!projection.is_complete());
}
#[tokio::test]
async fn messages_returns_message_entries_only() {
let (sender, receiver) = event_channel();
let projection = Projection::new(receiver);
sender
.send(Event::Message("msg".to_string()))
.await
.expect("should send");
sender
.send(Event::Detail("dtl".to_string()))
.await
.expect("should send");
sender
.send(Event::Message("msg2".to_string()))
.await
.expect("should send");
drop(sender);
tokio::task::yield_now().await;
let messages = projection.messages();
assert_eq!(messages.len(), 2);
let Entry::Message(s) = &messages[0] else {
panic!("expected Entry::Message");
};
assert_eq!(s, "msg");
let Entry::Message(s) = &messages[1] else {
panic!("expected Entry::Message");
};
assert_eq!(s, "msg2");
}
#[tokio::test]
async fn new_returns_empty_projection() {
let (_sender, receiver) = event_channel();
let projection = Projection::new(receiver);
assert!(projection.entries().is_empty());
assert!(!projection.is_complete());
}
#[tokio::test]
async fn new_starts_drain_that_completes_when_channel_closes() {
let (sender, receiver) = event_channel();
let projection = Projection::new(receiver);
drop(sender);
tokio::task::yield_now().await;
assert!(projection.is_complete());
}
#[test]
fn trait_send() {
fn assert_send<T: Send>() {}
assert_send::<Projection>();
}
#[test]
fn trait_sync() {
fn assert_sync<T: Sync>() {}
assert_sync::<Projection>();
}
#[test]
fn trait_unpin() {
fn assert_unpin<T: Unpin>() {}
assert_unpin::<Projection>();
}
}