use std::any::TypeId;
use headers::ContentType;
use http::StatusCode;
use tokio::sync::{mpsc, oneshot};
use utoipa::openapi::{RefOr, Schema};
use super::collectors::Collectors;
use super::operation::CalledOperation;
use super::schema::{SchemaEntry, Schemas};
const CHANNEL_BUFFER_SIZE: usize = 256;
pub(in crate::client) enum CollectorMessage {
AddSchemas(Schemas),
AddSchemaEntry(SchemaEntry),
AddExample {
type_id: TypeId,
type_name: &'static str,
example: serde_json::Value,
},
RegisterOperation(CalledOperation),
RegisterResponse {
operation_id: String,
status: StatusCode,
content_type: Option<ContentType>,
schema: Option<RefOr<Schema>>,
description: String,
},
#[cfg(feature = "redaction")]
RegisterResponseWithExample {
operation_id: String,
status: StatusCode,
content_type: Option<ContentType>,
schema: RefOr<Schema>,
example: serde_json::Value,
},
GetCollectors(oneshot::Sender<Collectors>),
}
#[derive(Debug, Clone)]
pub(in crate::client) struct CollectorSender {
inner: mpsc::Sender<CollectorMessage>,
}
impl CollectorSender {
pub(in crate::client) fn dummy() -> Self {
let (sender, _) = mpsc::channel::<CollectorMessage>(1);
Self { inner: sender }
}
pub(in crate::client) async fn send(&self, msg: CollectorMessage) {
self.inner
.send(msg)
.await
.expect("Collector task should be running");
}
}
#[derive(Debug, Clone)]
pub(in crate::client) struct CollectorHandle {
sender: CollectorSender,
}
impl CollectorHandle {
pub(in crate::client) fn spawn() -> Self {
let (sender, receiver) = mpsc::channel::<CollectorMessage>(CHANNEL_BUFFER_SIZE);
tokio::spawn(collector_task(receiver));
Self {
sender: CollectorSender { inner: sender },
}
}
pub(in crate::client) fn sender(&self) -> CollectorSender {
self.sender.clone()
}
pub(in crate::client) async fn get_collectors(&self) -> Collectors {
let (tx, rx) = oneshot::channel();
self.sender.send(CollectorMessage::GetCollectors(tx)).await;
rx.await.expect("Collector task should respond")
}
}
async fn collector_task(mut receiver: mpsc::Receiver<CollectorMessage>) {
let mut collectors = Collectors::default();
while let Some(msg) = receiver.recv().await {
match msg {
CollectorMessage::AddSchemas(schemas) => {
collectors.collect_schemas(schemas);
}
CollectorMessage::AddSchemaEntry(entry) => {
collectors.collect_schema_entry(entry);
}
CollectorMessage::AddExample {
type_id,
type_name,
example,
} => {
collectors
.schemas
.add_example_by_id(type_id, type_name, example);
}
CollectorMessage::RegisterOperation(operation) => {
collectors.collect_operation(operation);
}
CollectorMessage::RegisterResponse {
operation_id,
status,
content_type,
schema,
description,
} => {
collectors.register_response(
&operation_id,
status,
content_type.as_ref(),
schema,
description,
);
}
#[cfg(feature = "redaction")]
CollectorMessage::RegisterResponseWithExample {
operation_id,
status,
content_type,
schema,
example,
} => {
collectors.register_response_with_example(
&operation_id,
status,
content_type.as_ref(),
schema,
example,
);
}
CollectorMessage::GetCollectors(responder) => {
let _ = responder.send(collectors.clone());
}
}
}
}