caelix_core/
microservice.rs1use crate::{BoxFuture, Container, Injectable, ProviderDef, Result};
4use serde_json::Value;
5use std::{collections::BTreeMap, sync::Arc, time::SystemTime};
6use tokio_util::sync::CancellationToken;
7
8#[derive(Clone, Copy, Debug, Eq, PartialEq)]
10pub enum MessageHandlerKind {
11 Command,
13 Event,
15}
16
17#[derive(Clone, Debug, Default)]
19pub struct MessageDelivery {
20 pub stream: Option<String>,
22 pub consumer: Option<String>,
24 pub attempt: u64,
26}
27
28#[derive(Clone, Debug)]
30pub struct MessageContext {
31 subject: String,
32 headers: BTreeMap<String, String>,
33 correlation_id: Option<String>,
34 deadline: Option<SystemTime>,
35 cancellation: CancellationToken,
36 delivery: Option<MessageDelivery>,
37 event_id: Option<String>,
38}
39
40impl MessageContext {
41 pub fn new(
43 subject: impl Into<String>,
44 headers: BTreeMap<String, String>,
45 correlation_id: Option<String>,
46 deadline: Option<SystemTime>,
47 cancellation: CancellationToken,
48 delivery: Option<MessageDelivery>,
49 event_id: Option<String>,
50 ) -> Self {
51 Self {
52 subject: subject.into(),
53 headers,
54 correlation_id,
55 deadline,
56 cancellation,
57 delivery,
58 event_id,
59 }
60 }
61
62 pub fn subject(&self) -> &str {
64 &self.subject
65 }
66
67 pub fn headers(&self) -> &BTreeMap<String, String> {
69 &self.headers
70 }
71
72 pub fn correlation_id(&self) -> Option<&str> {
74 self.correlation_id.as_deref()
75 }
76
77 pub fn deadline(&self) -> Option<SystemTime> {
79 self.deadline
80 }
81
82 pub fn cancellation_token(&self) -> &CancellationToken {
84 &self.cancellation
85 }
86
87 pub fn delivery(&self) -> Option<&MessageDelivery> {
89 self.delivery.as_ref()
90 }
91
92 pub fn delivery_attempt(&self) -> u64 {
94 self.delivery
95 .as_ref()
96 .map_or(0, |delivery| delivery.attempt)
97 }
98
99 pub fn event_id(&self) -> Option<&str> {
101 self.event_id.as_deref()
102 }
103}
104
105type InvokeMessageFn = Arc<
106 dyn for<'a> Fn(&'a Container, MessageContext, Value) -> BoxFuture<'a, Result<Option<Value>>>
107 + Send
108 + Sync,
109>;
110
111#[derive(Clone)]
113pub struct MessageHandlerDef {
114 pub kind: MessageHandlerKind,
116 pub pattern: &'static str,
118 invoke: InvokeMessageFn,
119}
120
121impl MessageHandlerDef {
122 pub fn new(
124 kind: MessageHandlerKind,
125 pattern: &'static str,
126 invoke: impl for<'a> Fn(
127 &'a Container,
128 MessageContext,
129 Value,
130 ) -> BoxFuture<'a, Result<Option<Value>>>
131 + Send
132 + Sync
133 + 'static,
134 ) -> Self {
135 Self {
136 kind,
137 pattern,
138 invoke: Arc::new(invoke),
139 }
140 }
141
142 pub fn invoke<'a>(
144 &self,
145 container: &'a Container,
146 context: MessageContext,
147 payload: Value,
148 ) -> BoxFuture<'a, Result<Option<Value>>> {
149 (self.invoke)(container, context, payload)
150 }
151}
152
153pub struct MicroserviceDef {
155 pub provider: ProviderDef,
157 handlers_fn: fn() -> Vec<MessageHandlerDef>,
158}
159
160impl MicroserviceDef {
161 pub fn of<T: Injectable>(handlers_fn: fn() -> Vec<MessageHandlerDef>) -> Self {
163 Self {
164 provider: ProviderDef::of::<T>(),
165 handlers_fn,
166 }
167 }
168
169 pub fn handlers(&self) -> Vec<MessageHandlerDef> {
171 (self.handlers_fn)()
172 }
173}
174
175pub trait Microservice: Injectable {
177 fn definition() -> MicroserviceDef;
179}
180
181#[doc(hidden)]
186pub fn _assert_microservice_send_sync<T: Send + Sync + 'static>() {
187 let _ = std::marker::PhantomData::<Arc<T>>;
188}