1use crate::address::HeaderAddressList;
2#[cfg(feature = "impl")]
3use crate::dkim::Signer;
4#[cfg(feature = "impl")]
5use crate::dkim::SIGN_POOL;
6pub use crate::queue_name::QueueNameComponents;
7use crate::scheduling::Scheduling;
8use anyhow::Context;
9use bstr::{BString, ByteSlice, ByteVec};
10use chrono::{DateTime, Utc};
11#[cfg(feature = "impl")]
12use config::{any_err, from_lua_value, serialize_options, SerdeWrappedValue};
13use futures::{FutureExt, TryFutureExt};
14use intrusive_collections::{intrusive_adapter, LinkedList, LinkedListAtomicLink};
15use kumo_chrono_helper::*;
16#[cfg(feature = "impl")]
17use kumo_dkim::arc::ARC;
18use kumo_log_types::rfc3464::Report;
19use kumo_log_types::rfc5965::ARFReport;
20use kumo_prometheus::declare_metric;
21#[cfg(feature = "impl")]
22use mailparsing::{AuthenticationResult, AuthenticationResults, EncodeHeaderValue};
23use mailparsing::{
24 CheckFixSettings, DecodedBody, Header, HeaderParseResult, MessageConformance, MimePart,
25};
26#[cfg(feature = "impl")]
27use mlua::{IntoLua, LuaSerdeExt, UserData, UserDataMethods};
28#[cfg(feature = "impl")]
29use mod_dns_resolver::get_resolver_instance;
30use parking_lot::Mutex;
31use rfc5321::parser::EnvelopeAddress;
32use serde::{Deserialize, Serialize};
33use serde_with::formats::PreferOne;
34use serde_with::{serde_as, OneOrMany};
35use spool::{get_data_spool, get_meta_spool, Spool, SpoolId};
36use std::collections::{BTreeMap, HashSet};
37use std::hash::Hash;
38use std::sync::{Arc, LazyLock, Weak};
39use std::time::{Duration, Instant};
40use timeq::TimerEntryWithDelay;
41
42bitflags::bitflags! {
43 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
44 struct MessageFlags: u8 {
45 const META_DIRTY = 1;
47 const DATA_DIRTY = 2;
49 const SCHEDULED = 4;
51 const FORCE_SYNC = 8;
53 }
54}
55
56declare_metric! {
57static MESSAGE_COUNT: IntGauge("message_count");
64}
65
66declare_metric! {
67static META_COUNT: IntGauge("message_meta_resident_count");
74}
75
76declare_metric! {
77static DATA_COUNT: IntGauge("message_data_resident_count");
87}
88
89static NO_DATA: LazyLock<Arc<Box<[u8]>>> = LazyLock::new(|| Arc::new(vec![].into_boxed_slice()));
90
91declare_metric! {
92static SAVE_HIST: Histogram("message_save_latency");
102}
103
104declare_metric! {
105static LOAD_DATA_HIST: Histogram("message_data_load_latency");
116}
117
118declare_metric! {
119static LOAD_META_HIST: Histogram("message_meta_load_latency");
125}
126
127#[derive(Debug)]
128struct MessageInner {
129 metadata: Option<Box<MetaData>>,
130 data: Arc<Box<[u8]>>,
131 flags: MessageFlags,
132 num_attempts: u16,
133 due: Option<DateTime<Utc>>,
134}
135
136#[derive(Debug)]
137pub(crate) struct MessageWithId {
138 id: SpoolId,
139 inner: Mutex<MessageInner>,
140 link: LinkedListAtomicLink,
141}
142
143intrusive_adapter!(
144 pub(crate) MessageWithIdAdapter = Arc<MessageWithId>: MessageWithId { link: LinkedListAtomicLink }
145);
146
147pub struct MessageList {
152 list: LinkedList<MessageWithIdAdapter>,
153 len: usize,
154}
155
156impl Default for MessageList {
157 fn default() -> Self {
158 Self::new()
159 }
160}
161
162impl MessageList {
163 pub fn new() -> Self {
165 Self {
166 list: LinkedList::new(MessageWithIdAdapter::default()),
167 len: 0,
168 }
169 }
170
171 pub fn len(&self) -> usize {
173 self.len
174 }
175
176 pub fn is_empty(&self) -> bool {
177 self.len == 0
178 }
179
180 pub fn take(&mut self) -> Self {
183 let new_list = Self {
184 list: self.list.take(),
185 len: self.len,
186 };
187 self.len = 0;
188 new_list
189 }
190
191 pub fn push_back(&mut self, message: Message) {
193 self.list.push_back(message.msg_and_id);
194 self.len += 1;
195 }
196
197 pub fn pop_front(&mut self) -> Option<Message> {
199 self.list.pop_front().map(|msg_and_id| {
200 self.len -= 1;
201 Message { msg_and_id }
202 })
203 }
204
205 pub fn pop_back(&mut self) -> Option<Message> {
207 self.list.pop_back().map(|msg_and_id| {
208 self.len -= 1;
209 Message { msg_and_id }
210 })
211 }
212
213 pub fn drain(&mut self) -> Vec<Message> {
217 let mut messages = Vec::with_capacity(self.len);
218 while let Some(msg) = self.pop_front() {
219 messages.push(msg);
220 }
221 messages
222 }
223
224 pub fn extend_from_iter<I>(&mut self, iter: I)
225 where
226 I: Iterator<Item = Message>,
227 {
228 for msg in iter {
229 self.push_back(msg)
230 }
231 }
232}
233
234impl IntoIterator for MessageList {
235 type Item = Message;
236 type IntoIter = MessageListIter;
237 fn into_iter(self) -> MessageListIter {
238 MessageListIter { list: self.list }
239 }
240}
241
242pub struct MessageListIter {
243 list: LinkedList<MessageWithIdAdapter>,
244}
245
246impl Iterator for MessageListIter {
247 type Item = Message;
248 fn next(&mut self) -> Option<Message> {
249 self.list
250 .pop_front()
251 .map(|msg_and_id| Message { msg_and_id })
252 }
253}
254
255#[derive(Clone, Debug)]
256#[cfg_attr(feature = "impl", derive(mlua::FromLua))]
257pub struct Message {
258 pub(crate) msg_and_id: Arc<MessageWithId>,
259}
260
261impl PartialEq for Message {
262 fn eq(&self, other: &Self) -> bool {
263 self.id() == other.id()
264 }
265}
266impl Eq for Message {}
267
268impl Hash for Message {
269 fn hash<H>(&self, hasher: &mut H)
270 where
271 H: std::hash::Hasher,
272 {
273 self.id().hash(hasher)
274 }
275}
276
277#[derive(Clone, Debug)]
278pub struct WeakMessage {
279 weak: Weak<MessageWithId>,
280}
281
282impl WeakMessage {
283 pub fn upgrade(&self) -> Option<Message> {
284 Some(Message {
285 msg_and_id: self.weak.upgrade()?,
286 })
287 }
288}
289
290#[serde_as]
291#[derive(Debug, Serialize, Deserialize, Clone)]
292pub(crate) struct MetaData {
293 pub sender: EnvelopeAddress,
294 #[serde_as(as = "OneOrMany<_, PreferOne>")]
295 pub recipient: Vec<EnvelopeAddress>,
296 pub meta: serde_json::Value,
297 #[serde(default)]
298 pub schedule: Option<Scheduling>,
299}
300
301impl Drop for MessageInner {
302 fn drop(&mut self) {
303 if self.metadata.is_some() {
304 META_COUNT.dec();
305 }
306 if !self.data.is_empty() {
307 DATA_COUNT.dec();
308 }
309 MESSAGE_COUNT.dec();
310 }
311}
312
313impl Message {
314 pub fn new_dirty(
317 id: SpoolId,
318 sender: EnvelopeAddress,
319 recipient: Vec<EnvelopeAddress>,
320 meta: serde_json::Value,
321 data: Arc<Box<[u8]>>,
322 ) -> anyhow::Result<Self> {
323 anyhow::ensure!(meta.is_object(), "metadata must be a json object");
324 MESSAGE_COUNT.inc();
325 DATA_COUNT.inc();
326 META_COUNT.inc();
327 Ok(Self {
328 msg_and_id: Arc::new(MessageWithId {
329 id,
330 inner: Mutex::new(MessageInner {
331 metadata: Some(Box::new(MetaData {
332 sender,
333 recipient,
334 meta,
335 schedule: None,
336 })),
337 data,
338 flags: MessageFlags::META_DIRTY | MessageFlags::DATA_DIRTY,
339 num_attempts: 0,
340 due: None,
341 }),
342 link: LinkedListAtomicLink::default(),
343 }),
344 })
345 }
346
347 pub fn weak(&self) -> WeakMessage {
348 WeakMessage {
349 weak: Arc::downgrade(&self.msg_and_id),
350 }
351 }
352
353 pub fn new_from_spool(id: SpoolId, metadata: Vec<u8>) -> anyhow::Result<Self> {
357 let metadata: MetaData = serde_json::from_slice(&metadata)?;
358 MESSAGE_COUNT.inc();
359 META_COUNT.inc();
360
361 let flags = if metadata.schedule.is_some() {
362 MessageFlags::SCHEDULED
363 } else {
364 MessageFlags::empty()
365 };
366
367 Ok(Self {
368 msg_and_id: Arc::new(MessageWithId {
369 id,
370 inner: Mutex::new(MessageInner {
371 metadata: Some(Box::new(metadata)),
372 data: NO_DATA.clone(),
373 flags,
374 num_attempts: 0,
375 due: None,
376 }),
377 link: LinkedListAtomicLink::default(),
378 }),
379 })
380 }
381
382 pub(crate) fn new_from_parts(id: SpoolId, metadata: MetaData, data: Arc<Box<[u8]>>) -> Self {
383 MESSAGE_COUNT.inc();
384 META_COUNT.inc();
385
386 let flags = if metadata.schedule.is_some() {
387 MessageFlags::SCHEDULED
388 } else {
389 MessageFlags::empty()
390 };
391
392 Self {
393 msg_and_id: Arc::new(MessageWithId {
394 id,
395 inner: Mutex::new(MessageInner {
396 metadata: Some(Box::new(metadata)),
397 data,
398 flags: flags | MessageFlags::META_DIRTY | MessageFlags::DATA_DIRTY,
399 num_attempts: 0,
400 due: None,
401 }),
402 link: LinkedListAtomicLink::default(),
403 }),
404 }
405 }
406
407 pub async fn new_with_id(id: SpoolId) -> anyhow::Result<Self> {
408 let meta_spool = get_meta_spool();
409 let data = meta_spool.load(id).await?;
410 Self::new_from_spool(id, data)
411 }
412
413 pub fn get_num_attempts(&self) -> u16 {
414 let inner = self.msg_and_id.inner.lock();
415 inner.num_attempts
416 }
417
418 pub fn set_num_attempts(&self, num_attempts: u16) {
419 let mut inner = self.msg_and_id.inner.lock();
420 inner.num_attempts = num_attempts;
421 }
422
423 pub fn increment_num_attempts(&self) {
424 let mut inner = self.msg_and_id.inner.lock();
425 inner.num_attempts += 1;
426 }
427
428 pub async fn set_scheduling(
429 &self,
430 scheduling: Option<Scheduling>,
431 ) -> anyhow::Result<Option<Scheduling>> {
432 self.load_meta_if_needed().await?;
433 let mut inner = self.msg_and_id.inner.lock();
434 match &mut inner.metadata {
435 None => anyhow::bail!("set_scheduling: metadata must be loaded first"),
436 Some(meta) => {
437 meta.schedule = scheduling;
438 inner.flags.set(MessageFlags::META_DIRTY, true);
439 inner
440 .flags
441 .set(MessageFlags::SCHEDULED, scheduling.is_some());
442 if let Some(sched) = scheduling {
443 let due = inner.due.unwrap_or_else(Utc::now);
444 inner.due = Some(sched.adjust_for_schedule(due));
445 }
446 Ok(scheduling)
447 }
448 }
449 }
450
451 pub async fn get_scheduling(&self) -> anyhow::Result<Option<Scheduling>> {
452 self.load_meta_if_needed().await?;
453 let inner = self.msg_and_id.inner.lock();
454 Ok(inner.metadata.as_ref().and_then(|meta| meta.schedule))
455 }
456
457 pub fn get_due(&self) -> Option<DateTime<Utc>> {
458 let inner = self.msg_and_id.inner.lock();
459 inner.due
460 }
461
462 pub async fn delay_with_jitter(&self, limit: i64) -> anyhow::Result<Option<DateTime<Utc>>> {
463 let scale = rand::random::<f32>();
464 let value = (scale * limit as f32) as i64;
465 self.delay_by(seconds(value)?).await
466 }
467
468 pub async fn delay_by(
469 &self,
470 duration: chrono::Duration,
471 ) -> anyhow::Result<Option<DateTime<Utc>>> {
472 let due = Utc::now() + duration;
473 self.set_due(Some(due)).await
474 }
475
476 pub async fn delay_by_and_jitter(
478 &self,
479 duration: chrono::Duration,
480 ) -> anyhow::Result<Option<DateTime<Utc>>> {
481 let scale = rand::random::<f32>();
482 let value = (scale * 60.) as i64;
483 let due = Utc::now() + duration + seconds(value)?;
484 self.set_due(Some(due)).await
485 }
486
487 pub async fn set_due(
488 &self,
489 due: Option<DateTime<Utc>>,
490 ) -> anyhow::Result<Option<DateTime<Utc>>> {
491 let due = {
492 let mut inner = self.msg_and_id.inner.lock();
493
494 if !inner.flags.contains(MessageFlags::SCHEDULED) {
495 inner.due = due;
497 return Ok(inner.due);
498 }
499
500 let due = due.unwrap_or_else(Utc::now);
501
502 if let Some(meta) = &inner.metadata {
503 inner.due = match &meta.schedule {
504 Some(sched) => Some(sched.adjust_for_schedule(due)),
505 None => Some(due),
506 };
507 return Ok(inner.due);
508 }
509
510 due
513 };
514
515 self.load_meta().await?;
516
517 {
518 let mut inner = self.msg_and_id.inner.lock();
519 match &inner.metadata {
520 Some(meta) => {
521 inner.due = match &meta.schedule {
522 Some(sched) => Some(sched.adjust_for_schedule(due)),
523 None => Some(due),
524 };
525 Ok(inner.due)
526 }
527 None => anyhow::bail!("loaded metadata, but metadata is not set!?"),
528 }
529 }
530 }
531
532 fn get_data_if_dirty(&self) -> Option<Arc<Box<[u8]>>> {
533 let inner = self.msg_and_id.inner.lock();
534 if inner.flags.contains(MessageFlags::DATA_DIRTY) {
535 Some(Arc::clone(&inner.data))
536 } else {
537 None
538 }
539 }
540
541 fn get_meta_if_dirty(&self) -> Option<MetaData> {
542 let inner = self.msg_and_id.inner.lock();
543 if inner.flags.contains(MessageFlags::META_DIRTY) {
544 inner.metadata.as_ref().map(|md| (**md).clone())
545 } else {
546 None
547 }
548 }
549
550 pub(crate) async fn clone_meta_data(&self) -> anyhow::Result<MetaData> {
551 self.load_meta_if_needed().await?;
552 let inner = self.msg_and_id.inner.lock();
553 inner
554 .metadata
555 .as_ref()
556 .ok_or_else(|| anyhow::anyhow!("metadata not loaded even though we just loaded it"))
557 .map(|md| (**md).clone())
558 }
559
560 pub fn set_force_sync(&self, force: bool) {
561 let mut inner = self.msg_and_id.inner.lock();
562 inner.flags.set(MessageFlags::FORCE_SYNC, force);
563 }
564
565 pub fn needs_save(&self) -> bool {
566 let inner = self.msg_and_id.inner.lock();
567 inner.flags.contains(MessageFlags::META_DIRTY)
568 || inner.flags.contains(MessageFlags::DATA_DIRTY)
569 }
570
571 pub async fn save(&self, deadline: Option<Instant>) -> anyhow::Result<()> {
572 let _timer = SAVE_HIST.start_timer();
573 self.save_to(&**get_meta_spool(), &**get_data_spool(), deadline)
574 .await
575 }
576
577 pub async fn save_to(
578 &self,
579 meta_spool: &(dyn Spool + Send + Sync),
580 data_spool: &(dyn Spool + Send + Sync),
581 deadline: Option<Instant>,
582 ) -> anyhow::Result<()> {
583 let force_sync = self
584 .msg_and_id
585 .inner
586 .lock()
587 .flags
588 .contains(MessageFlags::FORCE_SYNC);
589
590 let data_fut: futures::future::BoxFuture<'_, anyhow::Result<bool>> =
591 if let Some(data) = self.get_data_if_dirty() {
592 anyhow::ensure!(!data.is_empty(), "message data must not be empty");
593 data_spool
594 .store(self.msg_and_id.id, data, force_sync, deadline)
595 .map_ok(|_| true)
596 .boxed()
597 } else {
598 futures::future::ready(Ok(false)).boxed()
599 };
600 let meta_fut: futures::future::BoxFuture<'_, anyhow::Result<bool>> =
601 if let Some(meta) = self.get_meta_if_dirty() {
602 let meta = Arc::new(serde_json::to_vec(&meta)?.into_boxed_slice());
603 meta_spool
604 .store(self.msg_and_id.id, meta, force_sync, deadline)
605 .map_ok(|_| true)
606 .boxed()
607 } else {
608 futures::future::ready(Ok(false)).boxed()
609 };
610
611 let (data_res, meta_res) = tokio::join!(data_fut, meta_fut);
617
618 if matches!(data_res, Ok(true)) {
624 self.msg_and_id
625 .inner
626 .lock()
627 .flags
628 .remove(MessageFlags::DATA_DIRTY);
629 }
630 if matches!(meta_res, Ok(true)) {
631 self.msg_and_id
632 .inner
633 .lock()
634 .flags
635 .remove(MessageFlags::META_DIRTY);
636 }
637
638 match (data_res, meta_res) {
652 (Ok(_), Ok(_)) => Ok(()),
653 (Err(data_err), Ok(_)) => Err(data_err.context("data spool store")),
654 (Ok(_), Err(meta_err)) => Err(meta_err.context("meta spool store")),
655 (Err(data_err), Err(meta_err)) => {
656 let data_unhealthy = data_err.root_cause().is::<::spool::SpoolUnhealthyError>();
657 let meta_unhealthy = meta_err.root_cause().is::<::spool::SpoolUnhealthyError>();
658 match (data_unhealthy, meta_unhealthy) {
659 (true, _) => Err(data_err.context(format!(
660 "data spool store; meta spool store also failed: {meta_err:#}"
661 ))),
662 (false, true) => Err(meta_err.context(format!(
663 "meta spool store; data spool store also failed: {data_err:#}"
664 ))),
665 (false, false) => Err(anyhow::anyhow!(
666 "data spool store: {data_err:#}; meta spool store: {meta_err:#}"
667 )),
668 }
669 }
670 }
671 }
672
673 pub fn id(&self) -> &SpoolId {
674 &self.msg_and_id.id
675 }
676
677 pub async fn save_and_shrink(&self) -> anyhow::Result<bool> {
679 self.save(None).await?;
680 self.shrink()
681 }
682
683 pub async fn save_and_shrink_data(&self) -> anyhow::Result<bool> {
685 self.save(None).await?;
686 self.shrink_data()
687 }
688
689 pub fn shrink_data(&self) -> anyhow::Result<bool> {
690 let mut inner = self.msg_and_id.inner.lock();
691 let mut did_shrink = false;
692 if inner.flags.contains(MessageFlags::DATA_DIRTY) {
693 anyhow::bail!("Cannot shrink message: DATA_DIRTY");
694 }
695 if !inner.data.is_empty() {
696 DATA_COUNT.dec();
697 did_shrink = true;
698 }
699 if !inner.data.is_empty() {
700 inner.data = NO_DATA.clone();
701 did_shrink = true;
702 }
703 Ok(did_shrink)
704 }
705
706 pub fn shrink(&self) -> anyhow::Result<bool> {
707 let mut inner = self.msg_and_id.inner.lock();
708 let mut did_shrink = false;
709 if inner.flags.contains(MessageFlags::DATA_DIRTY) {
710 anyhow::bail!("Cannot shrink message: DATA_DIRTY");
711 }
712 if inner.flags.contains(MessageFlags::META_DIRTY) {
713 anyhow::bail!("Cannot shrink message: META_DIRTY");
714 }
715 if inner.metadata.take().is_some() {
716 META_COUNT.dec();
717 did_shrink = true;
718 }
719 if !inner.data.is_empty() {
720 DATA_COUNT.dec();
721 did_shrink = true;
722 }
723 if !inner.data.is_empty() {
724 inner.data = NO_DATA.clone();
725 did_shrink = true;
726 }
727 Ok(did_shrink)
728 }
729
730 pub async fn sender(&self) -> anyhow::Result<EnvelopeAddress> {
731 self.load_meta_if_needed().await?;
732 let inner = self.msg_and_id.inner.lock();
733 match &inner.metadata {
734 Some(meta) => Ok(meta.sender.clone()),
735 None => anyhow::bail!("Message::sender metadata is not loaded"),
736 }
737 }
738
739 pub async fn set_sender(&self, sender: EnvelopeAddress) -> anyhow::Result<()> {
740 self.load_meta_if_needed().await?;
741 let mut inner = self.msg_and_id.inner.lock();
742 match &mut inner.metadata {
743 Some(meta) => {
744 meta.sender = sender;
745 inner.flags.set(MessageFlags::META_DIRTY, true);
746 Ok(())
747 }
748 None => anyhow::bail!("Message::set_sender: metadata is not loaded"),
749 }
750 }
751
752 #[deprecated = "use recipient_list or first_recipient instead"]
753 pub async fn recipient(&self) -> anyhow::Result<EnvelopeAddress> {
754 self.first_recipient().await
755 }
756
757 pub async fn first_recipient(&self) -> anyhow::Result<EnvelopeAddress> {
758 self.load_meta_if_needed().await?;
759 let inner = self.msg_and_id.inner.lock();
760 match &inner.metadata {
761 Some(meta) => match meta.recipient.first() {
762 Some(recip) => Ok(recip.clone()),
763 None => anyhow::bail!("recipient list is empty!?"),
764 },
765 None => anyhow::bail!("Message::first_recipient: metadata is not loaded"),
766 }
767 }
768
769 pub async fn recipient_list(&self) -> anyhow::Result<Vec<EnvelopeAddress>> {
770 self.load_meta_if_needed().await?;
771 let inner = self.msg_and_id.inner.lock();
772 match &inner.metadata {
773 Some(meta) => Ok(meta.recipient.clone()),
774 None => anyhow::bail!("Message::recipient_list: metadata is not loaded"),
775 }
776 }
777
778 pub async fn recipient_list_string(&self) -> anyhow::Result<Vec<String>> {
779 self.load_meta_if_needed().await?;
780 let inner = self.msg_and_id.inner.lock();
781 match &inner.metadata {
782 Some(meta) => Ok(meta.recipient.iter().map(|a| a.to_string()).collect()),
783 None => anyhow::bail!("Message::recipient_list_string: metadata is not loaded"),
784 }
785 }
786
787 #[deprecated = "use set_recipient_list instead"]
788 pub async fn set_recipient(&self, recipient: EnvelopeAddress) -> anyhow::Result<()> {
789 self.set_recipient_list(vec![recipient]).await
790 }
791
792 pub async fn set_recipient_list(&self, recipient: Vec<EnvelopeAddress>) -> anyhow::Result<()> {
793 self.load_meta_if_needed().await?;
794
795 let mut inner = self.msg_and_id.inner.lock();
796 match &mut inner.metadata {
797 Some(meta) => {
798 meta.recipient = recipient;
799 inner.flags.set(MessageFlags::META_DIRTY, true);
800 Ok(())
801 }
802 None => anyhow::bail!("Message::set_recipient_list: metadata is not loaded"),
803 }
804 }
805
806 pub fn is_meta_loaded(&self) -> bool {
807 self.msg_and_id.inner.lock().metadata.is_some()
808 }
809
810 pub fn is_data_loaded(&self) -> bool {
811 !self.msg_and_id.inner.lock().data.is_empty()
812 }
813
814 pub async fn load_meta_if_needed(&self) -> anyhow::Result<()> {
815 if self.is_meta_loaded() {
816 return Ok(());
817 }
818 self.load_meta().await
819 }
820
821 pub async fn data(&self) -> anyhow::Result<Arc<Box<[u8]>>> {
822 self.load_data_if_needed().await
823 }
824
825 async fn load_data_if_needed(&self) -> anyhow::Result<Arc<Box<[u8]>>> {
826 if self.is_data_loaded() {
827 return Ok(self.get_data_maybe_not_loaded());
828 }
829 self.load_data().await
830 }
831
832 pub async fn load_meta(&self) -> anyhow::Result<()> {
833 let _timer = LOAD_META_HIST.start_timer();
834 self.load_meta_from(&**get_meta_spool()).await
835 }
836
837 async fn load_meta_from(&self, meta_spool: &(dyn Spool + Send + Sync)) -> anyhow::Result<()> {
838 let id = self.id();
839 let data = meta_spool.load(*id).await?;
840 let mut inner = self.msg_and_id.inner.lock();
841 let was_not_loaded = inner.metadata.is_none();
842 let metadata: MetaData = serde_json::from_slice(&data)?;
843 inner.metadata.replace(Box::new(metadata));
844 if was_not_loaded {
845 META_COUNT.inc();
846 }
847 Ok(())
848 }
849
850 pub async fn load_data(&self) -> anyhow::Result<Arc<Box<[u8]>>> {
851 let _timer = LOAD_DATA_HIST.start_timer();
852 self.load_data_from(&**get_data_spool()).await
853 }
854
855 async fn load_data_from(
856 &self,
857 data_spool: &(dyn Spool + Send + Sync),
858 ) -> anyhow::Result<Arc<Box<[u8]>>> {
859 let data = data_spool.load(*self.id()).await?;
860 let mut inner = self.msg_and_id.inner.lock();
861 let was_empty = inner.data.is_empty();
862 inner.data = Arc::new(data.into_boxed_slice());
863 if was_empty {
864 DATA_COUNT.inc();
865 }
866 Ok(inner.data.clone())
867 }
868
869 pub fn assign_data(&self, data: Vec<u8>) {
870 let mut inner = self.msg_and_id.inner.lock();
871 let was_empty = inner.data.is_empty();
872 inner.data = Arc::new(data.into_boxed_slice());
873 inner.flags.set(MessageFlags::DATA_DIRTY, true);
874 if was_empty {
875 DATA_COUNT.inc();
876 }
877 }
878
879 pub fn get_data_maybe_not_loaded(&self) -> Arc<Box<[u8]>> {
880 let inner = self.msg_and_id.inner.lock();
881 inner.data.clone()
882 }
883
884 pub async fn set_meta<S: AsRef<str>, V: Into<serde_json::Value>>(
885 &self,
886 key: S,
887 value: V,
888 ) -> anyhow::Result<()> {
889 self.load_meta_if_needed().await?;
890 let mut inner = self.msg_and_id.inner.lock();
891 match &mut inner.metadata {
892 None => anyhow::bail!("set_meta: metadata must be loaded first"),
893 Some(meta) => {
894 let key = key.as_ref();
895 let value = value.into();
896
897 match &mut meta.meta {
898 serde_json::Value::Object(map) => {
899 map.insert(key.to_string(), value);
900 }
901 _ => anyhow::bail!("metadata is somehow not a json object"),
902 }
903
904 inner.flags.set(MessageFlags::META_DIRTY, true);
905 Ok(())
906 }
907 }
908 }
909
910 pub async fn unset_meta<S: AsRef<str>>(&self, key: S) -> anyhow::Result<()> {
911 self.load_meta_if_needed().await?;
912 let mut inner = self.msg_and_id.inner.lock();
913 match &mut inner.metadata {
914 None => anyhow::bail!("set_meta: metadata must be loaded first"),
915 Some(meta) => {
916 let key = key.as_ref();
917
918 match &mut meta.meta {
919 serde_json::Value::Object(map) => {
920 map.remove(key);
921 }
922 _ => anyhow::bail!("metadata is somehow not a json object"),
923 }
924
925 inner.flags.set(MessageFlags::META_DIRTY, true);
926 Ok(())
927 }
928 }
929 }
930
931 pub async fn get_meta_string<S: serde_json::value::Index + std::fmt::Display + Copy>(
933 &self,
934 key: S,
935 ) -> anyhow::Result<Option<String>> {
936 match self.get_meta(key).await {
937 Ok(serde_json::Value::String(value)) => Ok(Some(value.to_string())),
938 Ok(serde_json::Value::Null) => Ok(None),
939 hmm => {
940 anyhow::bail!("expected '{key}' to be a string value, got {hmm:?}");
941 }
942 }
943 }
944
945 pub async fn get_meta_obj(&self) -> anyhow::Result<serde_json::Value> {
946 self.load_meta_if_needed().await?;
947 let inner = self.msg_and_id.inner.lock();
948 match &inner.metadata {
949 None => anyhow::bail!("get_meta_obj: metadata must be loaded first"),
950 Some(meta) => Ok(meta.meta.clone()),
951 }
952 }
953
954 pub async fn get_meta<S: serde_json::value::Index>(
955 &self,
956 key: S,
957 ) -> anyhow::Result<serde_json::Value> {
958 self.load_meta_if_needed().await?;
959 let inner = self.msg_and_id.inner.lock();
960 match &inner.metadata {
961 None => anyhow::bail!("get_meta: metadata must be loaded first"),
962 Some(meta) => match meta.meta.get(key) {
963 Some(value) => Ok(value.clone()),
964 None => Ok(serde_json::Value::Null),
965 },
966 }
967 }
968
969 pub fn age(&self, now: DateTime<Utc>) -> chrono::Duration {
970 self.msg_and_id.id.age(now)
971 }
972
973 pub async fn get_queue_name(&self) -> anyhow::Result<String> {
974 Ok(match self.get_meta_string("queue").await? {
975 Some(name) => name,
976 None => {
977 let name = QueueNameComponents::format(
978 self.get_meta_string("campaign").await?,
979 self.get_meta_string("tenant").await?,
980 self.first_recipient()
981 .await?
982 .domain()
983 .to_string()
984 .to_lowercase(),
985 self.get_meta_string("routing_domain").await?,
986 );
987 name.to_string()
988 }
989 })
990 }
991
992 #[cfg(feature = "impl")]
993 pub async fn arc_verify(
994 &self,
995 opt_resolver_name: Option<String>,
996 ) -> anyhow::Result<AuthenticationResult> {
997 let resolver = get_resolver_instance(&opt_resolver_name)?;
998 let data = self.data().await?;
999 let bytes = mailparsing::SharedString::try_from(data.as_ref().as_ref())?;
1000
1001 let parsed = mailparsing::Header::parse_headers(bytes.clone())?;
1002 let message = kumo_dkim::ParsedEmail::HeaderOnlyParse { bytes, parsed };
1003
1004 let arc = ARC::verify(&message, &**resolver).await;
1005 Ok(arc.authentication_result())
1006 }
1007
1008 #[cfg(feature = "impl")]
1009 pub async fn arc_seal(
1010 &self,
1011 signer: Signer,
1012 auth_results: AuthenticationResults,
1013 opt_resolver_name: Option<String>,
1014 ) -> anyhow::Result<()> {
1015 let resolver = get_resolver_instance(&opt_resolver_name)?;
1016 let data = self.data().await?;
1017 let bytes = mailparsing::SharedString::try_from(data.as_ref().as_ref())?;
1018 let parsed = mailparsing::Header::parse_headers(bytes.clone())?;
1019 let message = kumo_dkim::ParsedEmail::HeaderOnlyParse { bytes, parsed };
1020 let arc = ARC::verify(&message, &**resolver).await;
1021
1022 let headers = arc.seal(&message, auth_results, signer.signer())?;
1023 if !headers.is_empty() {
1024 let mut new_data = Vec::<u8>::with_capacity(data.len() + 1024);
1025
1026 for hdr in headers {
1027 hdr.write_header(&mut new_data).ok();
1028 }
1029 new_data.extend_from_slice(&data);
1030 self.assign_data(new_data);
1031 }
1032
1033 Ok(())
1034 }
1035
1036 #[cfg(feature = "impl")]
1037 pub async fn dkim_verify(
1038 &self,
1039 opt_resolver_name: Option<String>,
1040 ) -> anyhow::Result<Vec<AuthenticationResult>> {
1041 let resolver = get_resolver_instance(&opt_resolver_name)?;
1042 let data = self.data().await?;
1043 let bytes = mailparsing::SharedString::try_from(data.as_ref().as_ref())?;
1044
1045 let parsed = mailparsing::Header::parse_headers(bytes.clone())?;
1046 if parsed
1047 .overall_conformance
1048 .contains(MessageConformance::NON_CANONICAL_LINE_ENDINGS)
1049 {
1050 return Ok(vec![AuthenticationResult {
1051 method: "dkim".into(),
1052 method_version: None,
1053 result: "permerror".into(),
1054 reason: Some("message has non-canonical line endings".into()),
1055 props: Default::default(),
1056 }]);
1057 }
1058 let message = kumo_dkim::ParsedEmail::HeaderOnlyParse { bytes, parsed };
1059
1060 let results = kumo_dkim::verify_email_with_resolver(&message, &**resolver).await?;
1061 Ok(results)
1062 }
1063
1064 pub async fn parse(&self) -> anyhow::Result<MimePart<'static>> {
1069 let data = self.data().await?;
1070 let owned_data = String::from_utf8_lossy(data.as_ref().as_ref()).to_string();
1071 Ok(MimePart::parse(owned_data)?)
1072 }
1073
1074 pub async fn parse_rfc3464(&self) -> anyhow::Result<Option<Report>> {
1075 let data = self.data().await?;
1076 Report::parse(&data)
1077 }
1078
1079 pub async fn parse_rfc5965(&self) -> anyhow::Result<Option<ARFReport>> {
1080 let data = self.data().await?;
1081 ARFReport::parse(&data)
1082 }
1083
1084 pub async fn prepend_header(&self, name: Option<&str>, value: &[u8]) -> anyhow::Result<()> {
1085 let data = self.data().await?;
1086 let mut new_data = Vec::with_capacity(size_header(name, value) + 2 + data.len());
1087 emit_header(&mut new_data, name, value);
1088 new_data.extend_from_slice(&data);
1089 self.assign_data(new_data);
1090 Ok(())
1091 }
1092
1093 pub async fn append_header(&self, name: Option<&str>, value: &[u8]) -> anyhow::Result<()> {
1094 let data = self.data().await?;
1095 let mut new_data = Vec::with_capacity(size_header(name, value.as_bytes()) + 2 + data.len());
1096 for (idx, window) in data.windows(4).enumerate() {
1097 if window == b"\r\n\r\n" {
1098 let headers = &data[0..idx + 2];
1099 let body = &data[idx + 2..];
1100
1101 new_data.extend_from_slice(headers);
1102 emit_header(&mut new_data, name, value);
1103 new_data.extend_from_slice(body);
1104 self.assign_data(new_data);
1105 return Ok(());
1106 }
1107 }
1108
1109 anyhow::bail!("append_header could not find the end of the header block");
1110 }
1111
1112 pub async fn get_address_header(
1113 &self,
1114 header_name: &str,
1115 ) -> anyhow::Result<Option<HeaderAddressList>> {
1116 let data = self.data().await?;
1117 let HeaderParseResult { headers, .. } =
1118 mailparsing::Header::parse_headers(data.as_ref().as_ref())?;
1119
1120 match headers.get_first(header_name) {
1121 Some(hdr) => {
1122 let list = hdr.as_address_list()?;
1123 let result: HeaderAddressList = list.into();
1124 Ok(Some(result))
1125 }
1126 None => Ok(None),
1127 }
1128 }
1129
1130 pub async fn get_first_named_header_value(
1131 &self,
1132 name: &str,
1133 ) -> anyhow::Result<Option<BString>> {
1134 let data = self.data().await?;
1135 let HeaderParseResult { headers, .. } = Header::parse_headers(data.as_ref().as_ref())?;
1136
1137 match headers.get_first(name) {
1138 Some(hdr) => Ok(Some(
1139 hdr.as_unstructured()
1140 .unwrap_or_else(|_| hdr.get_raw_value().into()),
1141 )),
1142 None => Ok(None),
1143 }
1144 }
1145
1146 pub async fn get_all_named_header_values(&self, name: &str) -> anyhow::Result<Vec<BString>> {
1147 let data = self.data().await?;
1148 let HeaderParseResult { headers, .. } = Header::parse_headers(data.as_ref().as_ref())?;
1149
1150 let mut values = vec![];
1151 for hdr in headers.iter_named(name) {
1152 values.push(
1153 hdr.as_unstructured()
1154 .unwrap_or_else(|_| hdr.get_raw_value().into()),
1155 );
1156 }
1157 Ok(values)
1158 }
1159
1160 pub async fn get_all_headers(&self) -> anyhow::Result<Vec<(BString, BString)>> {
1161 let data = self.data().await?;
1162 let HeaderParseResult { headers, .. } = Header::parse_headers(data.as_ref().as_ref())?;
1163
1164 let mut values = vec![];
1165 for hdr in headers.iter() {
1166 values.push((
1167 hdr.get_name().into(),
1168 hdr.as_unstructured()
1169 .unwrap_or_else(|_| hdr.get_raw_value().into()),
1170 ));
1171 }
1172 Ok(values)
1173 }
1174
1175 pub async fn retain_headers<F: FnMut(usize, &Header) -> bool>(
1176 &self,
1177 mut func: F,
1178 ) -> anyhow::Result<()> {
1179 let data = self.data().await?;
1180 let mut new_data = Vec::with_capacity(data.len());
1181 let HeaderParseResult {
1182 headers,
1183 body_offset,
1184 ..
1185 } = Header::parse_headers(data.as_ref().as_ref())?;
1186 for (idx, hdr) in headers.iter().enumerate() {
1187 let retain = (func)(idx, hdr);
1188 if !retain {
1189 continue;
1190 }
1191 hdr.write_header(&mut new_data)?;
1192 }
1193 new_data.extend_from_slice(b"\r\n");
1194 new_data.extend_from_slice(&data[body_offset..]);
1195 self.assign_data(new_data);
1196 Ok(())
1197 }
1198
1199 pub async fn remove_first_named_header(&self, name: &str) -> anyhow::Result<()> {
1200 let mut removed = false;
1201 self.retain_headers(|_, hdr| {
1202 if hdr.get_name().eq_ignore_ascii_case(name.as_bytes()) && !removed {
1203 removed = true;
1204 false
1205 } else {
1206 true
1207 }
1208 })
1209 .await
1210 }
1211
1212 pub async fn import_x_headers(&self, names: Vec<String>) -> anyhow::Result<()> {
1213 let specs = if names.is_empty() {
1214 vec![ImportHeaderSpec {
1215 name: "X-*".to_string(),
1216 ..ImportHeaderSpec::default()
1217 }]
1218 } else {
1219 names
1220 .into_iter()
1221 .map(|name| ImportHeaderSpec {
1222 name,
1223 ..ImportHeaderSpec::default()
1224 })
1225 .collect()
1226 };
1227 self.import_headers(specs).await
1228 }
1229
1230 pub async fn import_headers(&self, specs: Vec<ImportHeaderSpec>) -> anyhow::Result<()> {
1231 let compiled: Vec<CompiledImportHeaderSpec> = specs
1232 .into_iter()
1233 .map(CompiledImportHeaderSpec::compile)
1234 .collect::<anyhow::Result<_>>()?;
1235
1236 let data = self.data().await?;
1237 let HeaderParseResult { headers, .. } = Header::parse_headers(data.as_ref().as_ref())?;
1238
1239 let mut accumulators: Vec<PerSpecAccumulator> = compiled
1240 .iter()
1241 .map(|s| PerSpecAccumulator::new(s.match_mode))
1242 .collect();
1243 let mut indices_to_remove: HashSet<usize> = HashSet::new();
1244
1245 for (idx, hdr) in headers.iter().enumerate() {
1246 let hdr_name = hdr.get_name();
1247 for (spec_idx, spec) in compiled.iter().enumerate() {
1248 if spec.matches(hdr_name) {
1249 let key = spec.target_key(&hdr.get_name_lossy());
1250 let value = hdr.as_unstructured()?.to_str_lossy().to_string();
1251 accumulators[spec_idx].record(key, value);
1252 if spec.remove {
1253 indices_to_remove.insert(idx);
1254 }
1255 break;
1256 }
1257 }
1258 }
1259
1260 for acc in accumulators {
1261 acc.write_to(self).await?;
1262 }
1263
1264 if !indices_to_remove.is_empty() {
1265 self.retain_headers(|idx, _| !indices_to_remove.contains(&idx))
1266 .await?;
1267 }
1268
1269 Ok(())
1270 }
1271
1272 pub async fn remove_x_headers(&self, names: Vec<String>) -> anyhow::Result<()> {
1273 self.retain_headers(|_, hdr| {
1274 if names.is_empty() {
1275 !is_x_header(hdr.get_name())
1276 } else {
1277 !names
1278 .iter()
1279 .any(|n| hdr.get_name().eq_ignore_ascii_case(n.as_bytes()))
1280 }
1281 })
1282 .await
1283 }
1284
1285 pub async fn remove_all_named_headers(&self, name: &str) -> anyhow::Result<()> {
1286 self.retain_headers(|_, hdr| !hdr.get_name().eq_ignore_ascii_case(name.as_bytes()))
1287 .await
1288 }
1289
1290 #[cfg(feature = "impl")]
1291 pub async fn dkim_sign(&self, signer: Signer) -> anyhow::Result<()> {
1292 let data = self.data().await?;
1293 let header = if let Some(runtime) = SIGN_POOL.get() {
1294 runtime.spawn_blocking(move || signer.sign(&data)).await??
1295 } else {
1296 signer.sign(&data)?
1297 };
1298 self.prepend_header(None, header.as_bytes()).await?;
1299 Ok(())
1300 }
1301
1302 pub async fn import_scheduling_header(
1303 &self,
1304 header_name: &str,
1305 remove: bool,
1306 ) -> anyhow::Result<Option<Scheduling>> {
1307 if let Some(value) = self.get_first_named_header_value(header_name).await? {
1308 let sched: Scheduling = serde_json::from_slice(&value).with_context(|| {
1309 format!("{value} from header {header_name} is not a valid Scheduling header")
1310 })?;
1311 let result = self.set_scheduling(Some(sched)).await?;
1312
1313 if remove {
1314 self.remove_all_named_headers(header_name).await?;
1315 }
1316
1317 Ok(result)
1318 } else {
1319 Ok(None)
1320 }
1321 }
1322
1323 pub async fn append_text_plain(&self, content: &str) -> anyhow::Result<bool> {
1324 let data = self.data().await?;
1325 let mut msg = MimePart::parse(data.as_ref().as_ref())?;
1326 let parts = msg.simplified_structure_pointers()?;
1327 if let Some(p) = parts.text_part.and_then(|p| msg.resolve_ptr_mut(p)) {
1328 match p.body()? {
1329 DecodedBody::Text(text) => {
1330 let mut text = text.as_bytes().to_vec();
1331 text.push_str("\r\n");
1332 text.push_str(content);
1333 p.replace_text_body("text/plain", &*text)?;
1334
1335 let new_data = msg.to_message_bytes();
1336 self.assign_data(new_data);
1337 Ok(true)
1338 }
1339 DecodedBody::Binary(_) => {
1340 anyhow::bail!("expected text/plain part to be text, but it is binary");
1341 }
1342 }
1343 } else {
1344 Ok(false)
1345 }
1346 }
1347
1348 pub async fn append_text_html(&self, content: &str) -> anyhow::Result<bool> {
1349 let data = self.data().await?;
1350 let mut msg = MimePart::parse(data.as_ref().as_ref())?;
1351 let parts = msg.simplified_structure_pointers()?;
1352 if let Some(p) = parts.html_part.and_then(|p| msg.resolve_ptr_mut(p)) {
1353 match p.body()? {
1354 DecodedBody::Text(text) => {
1355 let mut text = text.as_bytes().to_vec();
1356
1357 match text.rfind("</body>").or_else(|| text.rfind("</BODY>")) {
1358 Some(idx) => {
1359 text.insert_str(idx, content);
1360 text.insert_str(idx, "\r\n");
1361 }
1362 None => {
1363 text.push_str("\r\n");
1365 text.push_str(content);
1366 }
1367 }
1368
1369 p.replace_text_body("text/html", &*text)?;
1370
1371 let new_data = msg.to_message_bytes();
1372 self.assign_data(new_data);
1373 Ok(true)
1374 }
1375 DecodedBody::Binary(_) => {
1376 anyhow::bail!("expected text/html part to be text, but it is binary");
1377 }
1378 }
1379 } else {
1380 Ok(false)
1381 }
1382 }
1383
1384 pub async fn check_fix_conformance(
1385 &self,
1386 check: MessageConformance,
1387 fix: MessageConformance,
1388 settings: Option<&CheckFixSettings>,
1389 ) -> anyhow::Result<()> {
1390 let data = self.data().await?;
1391 let data_bytes = data.as_ref().as_ref();
1392 let msg = MimePart::parse(data_bytes)?;
1393
1394 let mut settings = settings.map(Clone::clone).unwrap_or_default();
1395
1396 if fix.contains(MessageConformance::MISSING_MESSAGE_ID_HEADER)
1397 && settings.message_id.is_none()
1398 && matches!(msg.headers().message_id(), Err(_) | Ok(None))
1399 {
1400 let sender = self.sender().await?;
1401 let domain = sender.domain();
1402 let id = *self.id();
1403 settings.message_id.replace(format!("{id}@{domain}"));
1404 }
1405
1406 if settings.detect_encoding {
1407 settings.data_bytes.replace(data.clone());
1408 }
1409
1410 let opt_msg = msg.check_fix_conformance(check, fix, settings)?;
1411
1412 if let Some(msg) = opt_msg {
1413 let new_data = msg.to_message_bytes();
1414 self.assign_data(new_data);
1415 }
1416
1417 Ok(())
1418 }
1419}
1420
1421fn imported_header_name(name: &str) -> String {
1422 name.chars()
1423 .map(|c| match c.to_ascii_lowercase() {
1424 '-' => '_',
1425 c => c,
1426 })
1427 .collect()
1428}
1429
1430fn is_x_header(name: &[u8]) -> bool {
1431 name.starts_with_str("X-") || name.starts_with_str("x-")
1432}
1433
1434#[derive(Default, Debug, Clone, Copy, Deserialize, Serialize, PartialEq, Eq)]
1435#[serde(rename_all = "snake_case")]
1436pub enum MatchMode {
1437 First,
1438 #[default]
1439 Last,
1440 All,
1441}
1442
1443#[derive(Default, Debug, Clone, Copy, Deserialize, Serialize, PartialEq, Eq)]
1444#[serde(rename_all = "snake_case")]
1445pub enum NameTransform {
1446 #[default]
1447 SnakeCase,
1448 KebabCase,
1449 CamelCase,
1450 PascalCase,
1451}
1452
1453#[derive(Debug, Clone, Deserialize)]
1454#[serde(deny_unknown_fields)]
1455pub struct ImportHeaderSpec {
1456 pub name: String,
1457 #[serde(default, rename = "match")]
1458 pub match_mode: MatchMode,
1459 #[serde(default)]
1460 pub transform: NameTransform,
1461 #[serde(default)]
1462 pub target: Option<String>,
1463 #[serde(default)]
1464 pub remove: bool,
1465}
1466
1467impl Default for ImportHeaderSpec {
1468 fn default() -> Self {
1469 Self {
1470 name: String::new(),
1471 match_mode: MatchMode::default(),
1472 transform: NameTransform::default(),
1473 target: None,
1474 remove: false,
1475 }
1476 }
1477}
1478
1479#[derive(Debug, Clone)]
1480enum HeaderPattern {
1481 Exact(String),
1483 Prefix(String),
1486}
1487
1488#[derive(Debug, Clone)]
1489struct CompiledImportHeaderSpec {
1490 pattern: HeaderPattern,
1491 match_mode: MatchMode,
1492 transform: NameTransform,
1493 target: Option<String>,
1494 remove: bool,
1495}
1496
1497impl CompiledImportHeaderSpec {
1498 fn compile(spec: ImportHeaderSpec) -> anyhow::Result<Self> {
1499 let pattern = compile_header_pattern(&spec.name)?;
1500 if spec.target.is_some() && matches!(pattern, HeaderPattern::Prefix(_)) {
1501 anyhow::bail!(
1502 "import_headers: `target` cannot be used with wildcard pattern {:?}",
1503 spec.name
1504 );
1505 }
1506 Ok(Self {
1507 pattern,
1508 match_mode: spec.match_mode,
1509 transform: spec.transform,
1510 target: spec.target,
1511 remove: spec.remove,
1512 })
1513 }
1514
1515 fn matches(&self, hdr_name: &[u8]) -> bool {
1516 match &self.pattern {
1517 HeaderPattern::Exact(s) => hdr_name.eq_ignore_ascii_case(s.as_bytes()),
1518 HeaderPattern::Prefix(p) => {
1519 hdr_name.len() >= p.len() && hdr_name[..p.len()].eq_ignore_ascii_case(p.as_bytes())
1520 }
1521 }
1522 }
1523
1524 fn target_key(&self, matched_name: &str) -> String {
1525 if let Some(target) = &self.target {
1526 return target.clone();
1527 }
1528 apply_name_transform(matched_name, self.transform)
1529 }
1530}
1531
1532fn compile_header_pattern(pat: &str) -> anyhow::Result<HeaderPattern> {
1533 if pat.is_empty() {
1534 anyhow::bail!("import_headers: header name pattern must not be empty");
1535 }
1536 let star_count = pat.chars().filter(|c| *c == '*').count();
1537 if star_count == 0 {
1538 return Ok(HeaderPattern::Exact(pat.to_string()));
1539 }
1540 let Some(prefix) = pat.strip_suffix('*') else {
1541 anyhow::bail!(
1542 "import_headers: only a single trailing `*` is supported in pattern {:?}",
1543 pat
1544 );
1545 };
1546 if star_count > 1 {
1547 anyhow::bail!(
1548 "import_headers: only a single trailing `*` is supported in pattern {:?}",
1549 pat
1550 );
1551 }
1552 if prefix.is_empty() {
1553 anyhow::bail!("import_headers: bare `*` patterns are not supported");
1554 }
1555 Ok(HeaderPattern::Prefix(prefix.to_string()))
1556}
1557
1558fn apply_name_transform(name: &str, transform: NameTransform) -> String {
1559 match transform {
1560 NameTransform::SnakeCase => imported_header_name(name),
1561 NameTransform::KebabCase => name
1562 .split('-')
1563 .map(|p| p.to_ascii_lowercase())
1564 .collect::<Vec<_>>()
1565 .join("-"),
1566 NameTransform::CamelCase => name
1567 .split('-')
1568 .enumerate()
1569 .map(|(i, p)| {
1570 if i == 0 {
1571 p.to_ascii_lowercase()
1572 } else {
1573 titlecase_ascii(p)
1574 }
1575 })
1576 .collect::<Vec<_>>()
1577 .join(""),
1578 NameTransform::PascalCase => name
1579 .split('-')
1580 .map(titlecase_ascii)
1581 .collect::<Vec<_>>()
1582 .join(""),
1583 }
1584}
1585
1586fn titlecase_ascii(s: &str) -> String {
1587 let mut chars = s.chars();
1588 match chars.next() {
1589 None => String::new(),
1590 Some(first) => {
1591 let mut out = String::with_capacity(s.len());
1592 out.push(first.to_ascii_uppercase());
1593 for c in chars {
1594 out.push(c.to_ascii_lowercase());
1595 }
1596 out
1597 }
1598 }
1599}
1600
1601enum AccValue {
1602 Str(String),
1603 Arr(Vec<String>),
1604}
1605
1606struct PerSpecAccumulator {
1607 mode: MatchMode,
1608 by_key: BTreeMap<String, AccValue>,
1609}
1610
1611impl PerSpecAccumulator {
1612 fn new(mode: MatchMode) -> Self {
1613 Self {
1614 mode,
1615 by_key: BTreeMap::new(),
1616 }
1617 }
1618
1619 fn record(&mut self, key: String, value: String) {
1620 match self.mode {
1621 MatchMode::First => {
1622 self.by_key.entry(key).or_insert(AccValue::Str(value));
1623 }
1624 MatchMode::Last => {
1625 self.by_key.insert(key, AccValue::Str(value));
1626 }
1627 MatchMode::All => {
1628 let entry = self
1629 .by_key
1630 .entry(key)
1631 .or_insert_with(|| AccValue::Arr(Vec::new()));
1632 if let AccValue::Arr(v) = entry {
1633 v.push(value);
1634 }
1635 }
1636 }
1637 }
1638
1639 async fn write_to(self, msg: &Message) -> anyhow::Result<()> {
1640 for (key, value) in self.by_key {
1641 match value {
1642 AccValue::Str(s) => msg.set_meta(key, s).await?,
1643 AccValue::Arr(arr) => {
1644 let json = serde_json::Value::Array(
1645 arr.into_iter().map(serde_json::Value::String).collect(),
1646 );
1647 msg.set_meta(key, json).await?;
1648 }
1649 }
1650 }
1651 Ok(())
1652 }
1653}
1654
1655fn size_header(name: Option<&str>, value: &[u8]) -> usize {
1656 name.map(|name| name.len() + 2).unwrap_or(0) + value.len()
1657}
1658
1659fn emit_header(dest: &mut Vec<u8>, name: Option<&str>, value: &[u8]) {
1660 if let Some(name) = name {
1661 dest.extend_from_slice(name.as_bytes());
1662 dest.extend_from_slice(b": ");
1663 }
1664 dest.extend_from_slice(value);
1665 if !value.ends_with_str("\r\n") {
1666 dest.extend_from_slice(b"\r\n");
1667 }
1668}
1669
1670#[cfg(feature = "impl")]
1671impl UserData for Message {
1672 fn add_methods<M: UserDataMethods<Self>>(methods: &mut M) {
1673 methods.add_async_method(
1674 "set_meta",
1675 move |_, this, (name, value): (String, mlua::Value)| async move {
1676 let value = serde_json::value::to_value(value).map_err(any_err)?;
1677 this.set_meta(name, value).await.map_err(any_err)?;
1678 Ok(())
1679 },
1680 );
1681 methods.add_async_method("get_meta", move |lua, this, name: String| async move {
1682 let value = this.get_meta(name).await.map_err(any_err)?;
1683 Ok(Some(lua.to_value_with(&value, serialize_options())?))
1684 });
1685 methods.add_async_method("get_data", |lua, this, _: ()| async move {
1686 let data = this.data().await.map_err(any_err)?;
1687 lua.create_string(&*data)
1688 });
1689 methods.add_method("set_data", move |_lua, this, data: mlua::String| {
1690 this.assign_data(data.as_bytes().to_vec());
1691 Ok(())
1692 });
1693
1694 methods.add_async_method("parse_mime", |_lua, this, _: ()| async move {
1695 let data = this.data().await.map_err(any_err)?;
1696 let owned_data = BString::new(data.as_ref().to_vec());
1697 let part = MimePart::parse(owned_data).map_err(any_err)?;
1698 Ok(mod_mimepart::PartRef::new(part))
1699 });
1700
1701 methods.add_async_method("append_text_plain", |_lua, this, data: String| async move {
1702 this.append_text_plain(&data).await.map_err(any_err)
1703 });
1704
1705 methods.add_async_method("append_text_html", |_lua, this, data: String| async move {
1706 this.append_text_html(&data).await.map_err(any_err)
1707 });
1708
1709 methods.add_method("id", move |_, this, _: ()| Ok(this.id().to_string()));
1710 methods.add_async_method("sender", move |_, this, _: ()| async move {
1711 this.sender().await.map_err(any_err)
1712 });
1713
1714 methods.add_method("num_attempts", move |_, this, _: ()| {
1715 Ok(this.get_num_attempts())
1716 });
1717 methods.add_method("increment_num_attempts", move |_, this, _: ()| {
1718 Ok(this.increment_num_attempts())
1719 });
1720
1721 methods.add_async_method("queue_name", move |_, this, _: ()| async move {
1722 this.get_queue_name().await.map_err(any_err)
1723 });
1724
1725 methods.add_async_method("set_due", move |lua, this, due: mlua::Value| async move {
1726 let due: Option<DateTime<Utc>> = lua.from_value(due)?;
1727 let revised_due = this.set_due(due).await.map_err(any_err)?;
1728 lua.to_value(&revised_due)
1729 });
1730
1731 methods.add_async_method(
1732 "set_sender",
1733 move |lua, this, value: mlua::Value| async move {
1734 let sender = match value {
1735 mlua::Value::String(s) => {
1736 let s = s.to_str()?;
1737 EnvelopeAddress::parse(&s).map_err(any_err)?
1738 }
1739 _ => lua.from_value::<EnvelopeAddress>(value.clone())?,
1740 };
1741 this.set_sender(sender).await.map_err(any_err)
1742 },
1743 );
1744
1745 methods.add_async_method("recipient", move |lua, this, _: ()| async move {
1746 let mut recipients = this.recipient_list().await.map_err(any_err)?;
1747 match recipients.len() {
1748 0 => Ok(mlua::Value::Nil),
1749 1 => {
1750 let recip: EnvelopeAddress = recipients.pop().expect("have 1");
1751 recip.into_lua(&lua)
1752 }
1753 _ => recipients.into_lua(&lua),
1754 }
1755 });
1756
1757 methods.add_async_method("recipient_list", move |lua, this, _: ()| async move {
1758 let recipients = this.recipient_list().await.map_err(any_err)?;
1759 recipients.into_lua(&lua)
1760 });
1761
1762 methods.add_async_method(
1763 "set_recipient",
1764 move |lua, this, value: mlua::Value| async move {
1765 let recipients = match value {
1766 mlua::Value::String(s) => {
1767 let s = s.to_str()?;
1768 vec![EnvelopeAddress::parse(&s).map_err(any_err)?]
1769 }
1770 _ => {
1771 if let Ok(recips) = lua.from_value::<Vec<EnvelopeAddress>>(value.clone()) {
1772 recips
1773 } else {
1774 vec![lua.from_value::<EnvelopeAddress>(value.clone())?]
1775 }
1776 }
1777 };
1778 this.set_recipient_list(recipients).await.map_err(any_err)
1779 },
1780 );
1781
1782 #[cfg(feature = "impl")]
1783 methods.add_async_method("dkim_sign", |_, this, signer: Signer| async move {
1784 this.dkim_sign(signer).await.map_err(any_err)
1785 });
1786
1787 methods.add_async_method("shrink", |_, this, _: ()| async move {
1788 if this.needs_save() {
1789 this.save(None).await.map_err(any_err)?;
1790 }
1791 this.shrink().map_err(any_err)
1792 });
1793
1794 methods.add_async_method("shrink_data", |_, this, _: ()| async move {
1795 if this.needs_save() {
1796 this.save(None).await.map_err(any_err)?;
1797 }
1798 this.shrink_data().map_err(any_err)
1799 });
1800
1801 methods.add_async_method(
1802 "add_authentication_results",
1803 |lua, this, (serv_id, results): (mlua::String, mlua::Value)| async move {
1804 let results: Vec<AuthenticationResult> = lua.from_value(results)?;
1805 let results = AuthenticationResults {
1806 serv_id: serv_id.as_bytes().as_ref().into(),
1807 version: None,
1808 results,
1809 };
1810
1811 this.prepend_header(
1812 Some("Authentication-Results"),
1813 results.encode_value().as_bytes(),
1814 )
1815 .await
1816 .map_err(any_err)?;
1817
1818 Ok(())
1819 },
1820 );
1821
1822 #[cfg(feature = "impl")]
1823 methods.add_async_method(
1824 "arc_verify",
1825 |lua, this, opt_resolver_name: Option<String>| async move {
1826 let results = this.arc_verify(opt_resolver_name).await.map_err(any_err)?;
1827 lua.to_value_with(&results, serialize_options())
1828 },
1829 );
1830
1831 #[cfg(feature = "impl")]
1832 methods.add_async_method(
1833 "arc_seal",
1834 |lua,
1835 this,
1836 (signer, serv_id, auth_res, opt_resolver_name): (
1837 Signer,
1838 mlua::String,
1839 mlua::Value,
1840 Option<String>,
1841 )| async move {
1842 let results: Vec<AuthenticationResult> = lua.from_value(auth_res)?;
1843 this.arc_seal(
1844 signer,
1845 AuthenticationResults {
1846 serv_id: serv_id.as_bytes().as_ref().into(),
1847 version: None,
1848 results,
1849 },
1850 opt_resolver_name,
1851 )
1852 .await
1853 .map_err(any_err)
1854 },
1855 );
1856
1857 #[cfg(feature = "impl")]
1858 methods.add_async_method(
1859 "dkim_verify",
1860 |lua, this, opt_resolver_name: Option<String>| async move {
1861 let results = this.dkim_verify(opt_resolver_name).await.map_err(any_err)?;
1862 lua.to_value_with(&results, serialize_options())
1863 },
1864 );
1865
1866 methods.add_async_method(
1867 "prepend_header",
1868 |_, this, (name, value, encode): (String, String, Option<bool>)| async move {
1869 let encode = encode.unwrap_or(false);
1870 if encode {
1871 let header = Header::new_unstructured(name.clone(), value);
1872 this.prepend_header(Some(&name), header.get_raw_value())
1873 .await
1874 .map_err(any_err)?;
1875 } else {
1876 this.prepend_header(Some(&name), value.as_bytes())
1877 .await
1878 .map_err(any_err)?;
1879 }
1880 Ok(())
1881 },
1882 );
1883 methods.add_async_method(
1884 "append_header",
1885 |_, this, (name, value, encode): (String, String, Option<bool>)| async move {
1886 let encode = encode.unwrap_or(false);
1887 if encode {
1888 let header = Header::new_unstructured(name.clone(), value);
1889 this.append_header(Some(&name), header.get_raw_value())
1890 .await
1891 .map_err(any_err)?;
1892 } else {
1893 this.append_header(Some(&name), value.as_bytes())
1894 .await
1895 .map_err(any_err)?;
1896 }
1897 Ok(())
1898 },
1899 );
1900 methods.add_async_method("get_address_header", |_, this, name: String| async move {
1901 this.get_address_header(&name).await.map_err(any_err)
1902 });
1903 methods.add_async_method("from_header", |_, this, ()| async move {
1904 this.get_address_header("From").await.map_err(any_err)
1905 });
1906 methods.add_async_method("to_header", |_, this, ()| async move {
1907 this.get_address_header("To").await.map_err(any_err)
1908 });
1909
1910 methods.add_async_method(
1911 "get_first_named_header_value",
1912 |_, this, name: String| async move {
1913 this.get_first_named_header_value(&name)
1914 .await
1915 .map_err(any_err)
1916 },
1917 );
1918 methods.add_async_method(
1919 "get_all_named_header_values",
1920 |_, this, name: String| async move {
1921 this.get_all_named_header_values(&name)
1922 .await
1923 .map_err(any_err)
1924 },
1925 );
1926 methods.add_async_method("get_all_headers", |_, this, _: ()| async move {
1927 Ok(this
1928 .get_all_headers()
1929 .await
1930 .map_err(any_err)?
1931 .into_iter()
1932 .map(|(name, value)| vec![name, value])
1933 .collect::<Vec<Vec<BString>>>())
1934 });
1935 methods.add_async_method(
1936 "import_x_headers",
1937 |_, this, names: Option<Vec<String>>| async move {
1938 this.import_x_headers(names.unwrap_or_default())
1939 .await
1940 .map_err(any_err)
1941 },
1942 );
1943 methods.add_async_method(
1944 "import_headers",
1945 |_, this, specs: SerdeWrappedValue<Vec<ImportHeaderSpec>>| async move {
1946 this.import_headers(specs.0).await.map_err(any_err)
1947 },
1948 );
1949
1950 methods.add_async_method(
1951 "remove_x_headers",
1952 |_, this, names: Option<Vec<String>>| async move {
1953 this.remove_x_headers(names.unwrap_or_default())
1954 .await
1955 .map_err(any_err)
1956 },
1957 );
1958 methods.add_async_method(
1959 "remove_all_named_headers",
1960 |_, this, name: String| async move {
1961 this.remove_all_named_headers(&name).await.map_err(any_err)
1962 },
1963 );
1964
1965 methods.add_async_method(
1966 "import_scheduling_header",
1967 |lua, this, (header_name, remove): (String, bool)| async move {
1968 let opt_schedule = this
1969 .import_scheduling_header(&header_name, remove)
1970 .await
1971 .map_err(any_err)?;
1972 lua.to_value(&opt_schedule)
1973 },
1974 );
1975
1976 methods.add_async_method(
1977 "set_scheduling",
1978 move |lua, this, params: mlua::Value| async move {
1979 let sched: Option<Scheduling> = from_lua_value(&lua, params)?;
1980 let opt_schedule = this.set_scheduling(sched).await.map_err(any_err)?;
1981 lua.to_value(&opt_schedule)
1982 },
1983 );
1984
1985 methods.add_async_method("parse_rfc3464", |lua, this, _: ()| async move {
1986 let report = this.parse_rfc3464().await.map_err(any_err)?;
1987 match report {
1988 Some(report) => lua.to_value_with(&report, serialize_options()),
1989 None => Ok(mlua::Value::Nil),
1990 }
1991 });
1992
1993 methods.add_async_method("parse_rfc5965", |lua, this, _: ()| async move {
1994 let report = this.parse_rfc5965().await.map_err(any_err)?;
1995 match report {
1996 Some(report) => lua.to_value_with(&report, serialize_options()),
1997 None => Ok(mlua::Value::Nil),
1998 }
1999 });
2000
2001 methods.add_async_method("save", |_, this, ()| async move {
2002 this.save(None).await.map_err(any_err)
2003 });
2004
2005 methods.add_method("set_force_sync", move |_, this, force: bool| {
2006 this.set_force_sync(force);
2007 Ok(())
2008 });
2009
2010 methods.add_async_method(
2011 "check_fix_conformance",
2012 |lua, this, (check, fix, settings): (String, String, Option<mlua::Value>)| async move {
2013 use std::str::FromStr;
2014 let check = MessageConformance::from_str(&check).map_err(any_err)?;
2015 let fix = MessageConformance::from_str(&fix).map_err(any_err)?;
2016
2017 let settings = match settings {
2018 Some(v) => Some(lua.from_value(v).map_err(any_err)?),
2019 None => None,
2020 };
2021
2022 match this
2023 .check_fix_conformance(check, fix, settings.as_ref())
2024 .await
2025 {
2026 Ok(_) => Ok(None),
2027 Err(err) => Ok(Some(format!("{err:#}"))),
2028 }
2029 },
2030 );
2031 }
2032}
2033
2034impl TimerEntryWithDelay for WeakMessage {
2035 fn delay(&self) -> Duration {
2036 match self.upgrade() {
2037 None => {
2038 Duration::from_millis(0)
2040 }
2041 Some(msg) => msg.delay(),
2042 }
2043 }
2044}
2045
2046impl TimerEntryWithDelay for Message {
2047 fn delay(&self) -> Duration {
2048 let inner = self.msg_and_id.inner.lock();
2049 match inner.due {
2050 Some(time) => {
2051 let now = Utc::now();
2052 let delta = time - now;
2053 delta.to_std().unwrap_or(Duration::from_millis(0))
2054 }
2055 None => Duration::from_millis(0),
2056 }
2057 }
2058}
2059
2060#[cfg(test)]
2061pub(crate) mod test {
2062 use super::*;
2063 use serde_json::json;
2064
2065 pub fn new_msg_body<S: AsRef<[u8]>>(s: S) -> Message {
2066 Message::new_dirty(
2067 SpoolId::new(),
2068 EnvelopeAddress::parse("sender@example.com").unwrap(),
2069 vec![EnvelopeAddress::parse("recip@example.com").unwrap()],
2070 serde_json::json!({}),
2071 Arc::new(s.as_ref().to_vec().into_boxed_slice()),
2072 )
2073 .unwrap()
2074 }
2075
2076 fn data_as_string(msg: &Message) -> String {
2077 String::from_utf8(msg.get_data_maybe_not_loaded().to_vec()).unwrap()
2078 }
2079
2080 const X_HDR_CONTENT: &str =
2081 "X-Hello: there\r\nX-Header: value\r\nSubject: Hello\r\nFrom :Someone\r\n\r\nBody";
2082
2083 #[tokio::test]
2084 async fn import_all_x_headers() {
2085 let msg = new_msg_body(X_HDR_CONTENT);
2086
2087 msg.import_x_headers(vec![]).await.unwrap();
2088 k9::assert_equal!(
2089 msg.get_meta_obj().await.unwrap(),
2090 json!({
2091 "x_hello": "there",
2092 "x_header": "value",
2093 })
2094 );
2095 }
2096
2097 #[tokio::test]
2098 async fn meta_and_nil() {
2099 let msg = new_msg_body(X_HDR_CONTENT);
2100 msg.set_meta("test", serde_json::Value::Null).await.unwrap();
2102 k9::assert_equal!(msg.get_meta("test").await.unwrap(), serde_json::Value::Null);
2103
2104 let lua = mlua::Lua::new();
2106 lua.globals().set("msg", msg).unwrap();
2107 lua.load("assert(msg:get_meta('test') == nil)")
2108 .exec()
2109 .unwrap();
2110 }
2111
2112 #[tokio::test]
2113 async fn set_sender_marks_meta_dirty() {
2114 let msg = new_msg_body(X_HDR_CONTENT);
2115 {
2118 let mut inner = msg.msg_and_id.inner.lock();
2119 inner
2120 .flags
2121 .remove(MessageFlags::DATA_DIRTY | MessageFlags::META_DIRTY);
2122 }
2123
2124 msg.set_sender(EnvelopeAddress::parse("new-sender@example.com").unwrap())
2125 .await
2126 .unwrap();
2127
2128 let inner = msg.msg_and_id.inner.lock();
2129 assert!(
2130 inner.flags.contains(MessageFlags::META_DIRTY),
2131 "changing the sender must mark the metadata dirty"
2132 );
2133 assert!(
2134 !inner.flags.contains(MessageFlags::DATA_DIRTY),
2135 "changing the sender must not mark the data dirty"
2136 );
2137 }
2138
2139 #[tokio::test]
2140 async fn set_recipient_list_marks_meta_dirty() {
2141 let msg = new_msg_body(X_HDR_CONTENT);
2142 {
2143 let mut inner = msg.msg_and_id.inner.lock();
2144 inner
2145 .flags
2146 .remove(MessageFlags::DATA_DIRTY | MessageFlags::META_DIRTY);
2147 }
2148
2149 msg.set_recipient_list(vec![
2150 EnvelopeAddress::parse("new-recip@example.com").unwrap()
2151 ])
2152 .await
2153 .unwrap();
2154
2155 let inner = msg.msg_and_id.inner.lock();
2156 assert!(
2157 inner.flags.contains(MessageFlags::META_DIRTY),
2158 "changing the recipient list must mark the metadata dirty"
2159 );
2160 assert!(
2161 !inner.flags.contains(MessageFlags::DATA_DIRTY),
2162 "changing the recipient list must not mark the data dirty"
2163 );
2164 }
2165
2166 #[tokio::test]
2167 async fn set_sender_after_shrink_does_not_flag_empty_data() {
2168 let msg = new_msg_body(X_HDR_CONTENT);
2169 {
2172 let mut inner = msg.msg_and_id.inner.lock();
2173 inner
2174 .flags
2175 .remove(MessageFlags::DATA_DIRTY | MessageFlags::META_DIRTY);
2176 inner.data = NO_DATA.clone();
2177 }
2178
2179 msg.set_sender(EnvelopeAddress::parse("new-sender@example.com").unwrap())
2180 .await
2181 .unwrap();
2182
2183 assert!(
2187 msg.get_data_if_dirty().is_none(),
2188 "an envelope change must not flag the empty data buffer for saving"
2189 );
2190 assert!(
2191 msg.needs_save(),
2192 "the metadata change must still require a save"
2193 );
2194 }
2195
2196 #[tokio::test]
2197 async fn set_scheduling_marks_meta_dirty() {
2198 let msg = new_msg_body(X_HDR_CONTENT);
2199 {
2200 let mut inner = msg.msg_and_id.inner.lock();
2201 inner
2202 .flags
2203 .remove(MessageFlags::DATA_DIRTY | MessageFlags::META_DIRTY);
2204 }
2205
2206 msg.set_scheduling(Some(Scheduling {
2207 restriction: None,
2208 first_attempt: None,
2209 expires: None,
2210 }))
2211 .await
2212 .unwrap();
2213
2214 let inner = msg.msg_and_id.inner.lock();
2215 assert!(
2216 inner.flags.contains(MessageFlags::META_DIRTY),
2217 "changing the schedule must mark the metadata dirty"
2218 );
2219 assert!(
2220 !inner.flags.contains(MessageFlags::DATA_DIRTY),
2221 "changing the schedule must not mark the data dirty"
2222 );
2223 }
2224
2225 #[tokio::test]
2226 async fn import_some_x_headers() {
2227 let msg = new_msg_body(X_HDR_CONTENT);
2228
2229 msg.import_x_headers(vec!["x-hello".to_string()])
2230 .await
2231 .unwrap();
2232 k9::assert_equal!(
2233 msg.get_meta_obj().await.unwrap(),
2234 json!({
2235 "x_hello": "there",
2236 })
2237 );
2238 }
2239
2240 #[tokio::test]
2241 async fn import_headers_wildcard_remove() {
2242 let msg = new_msg_body(X_HDR_CONTENT);
2243
2244 msg.import_headers(vec![ImportHeaderSpec {
2245 name: "X-*".to_string(),
2246 remove: true,
2247 ..ImportHeaderSpec::default()
2248 }])
2249 .await
2250 .unwrap();
2251 k9::assert_equal!(
2252 msg.get_meta_obj().await.unwrap(),
2253 json!({
2254 "x_hello": "there",
2255 "x_header": "value",
2256 })
2257 );
2258 k9::assert_equal!(
2259 data_as_string(&msg),
2260 "Subject: Hello\r\nFrom :Someone\r\n\r\nBody"
2261 );
2262 }
2263
2264 #[tokio::test]
2265 async fn import_headers_match_modes() {
2266 let body =
2267 "Received: from a\r\nReceived: from b\r\nReceived: from c\r\nSubject: hi\r\n\r\nBody";
2268 for (mode, expected) in [
2269 (MatchMode::First, json!("from a")),
2270 (MatchMode::Last, json!("from c")),
2271 (MatchMode::All, json!(["from a", "from b", "from c"])),
2272 ] {
2273 let msg = new_msg_body(body);
2274 msg.import_headers(vec![ImportHeaderSpec {
2275 name: "Received".to_string(),
2276 match_mode: mode,
2277 ..ImportHeaderSpec::default()
2278 }])
2279 .await
2280 .unwrap();
2281 k9::assert_equal!(
2282 msg.get_meta_obj().await.unwrap(),
2283 json!({ "received": expected })
2284 );
2285 }
2286 }
2287
2288 #[tokio::test]
2289 async fn import_headers_no_match_skips_meta() {
2290 let msg = new_msg_body(X_HDR_CONTENT);
2291 msg.import_headers(vec![ImportHeaderSpec {
2292 name: "Nonexistent".to_string(),
2293 match_mode: MatchMode::All,
2294 ..ImportHeaderSpec::default()
2295 }])
2296 .await
2297 .unwrap();
2298 k9::assert_equal!(msg.get_meta_obj().await.unwrap(), json!({}));
2299 }
2300
2301 #[tokio::test]
2302 async fn import_headers_specific_before_wildcard() {
2303 let body = "X-Campaign-Id: 42\r\nX-Mailer: foo\r\n\r\nBody";
2304 let msg = new_msg_body(body);
2305 msg.import_headers(vec![
2306 ImportHeaderSpec {
2307 name: "X-Campaign-Id".to_string(),
2308 target: Some("campaign".to_string()),
2309 ..ImportHeaderSpec::default()
2310 },
2311 ImportHeaderSpec {
2312 name: "X-*".to_string(),
2313 ..ImportHeaderSpec::default()
2314 },
2315 ])
2316 .await
2317 .unwrap();
2318 k9::assert_equal!(
2319 msg.get_meta_obj().await.unwrap(),
2320 json!({
2321 "campaign": "42",
2322 "x_mailer": "foo",
2323 })
2324 );
2325 }
2326
2327 #[test]
2328 fn name_transforms() {
2329 let n = "X-Campaign-Id";
2330 k9::assert_equal!(
2331 apply_name_transform(n, NameTransform::SnakeCase),
2332 "x_campaign_id"
2333 );
2334 k9::assert_equal!(
2335 apply_name_transform(n, NameTransform::KebabCase),
2336 "x-campaign-id"
2337 );
2338 k9::assert_equal!(
2339 apply_name_transform(n, NameTransform::CamelCase),
2340 "xCampaignId"
2341 );
2342 k9::assert_equal!(
2343 apply_name_transform(n, NameTransform::PascalCase),
2344 "XCampaignId"
2345 );
2346 }
2347
2348 #[test]
2349 fn pattern_compile_rejects_bad_patterns() {
2350 compile_header_pattern("").unwrap_err();
2351 compile_header_pattern("*").unwrap_err();
2352 compile_header_pattern("*-Id").unwrap_err();
2353 compile_header_pattern("X-*-Id").unwrap_err();
2354 compile_header_pattern("X-**").unwrap_err();
2355 assert!(matches!(
2356 compile_header_pattern("Subject").unwrap(),
2357 HeaderPattern::Exact(_)
2358 ));
2359 assert!(matches!(
2360 compile_header_pattern("X-*").unwrap(),
2361 HeaderPattern::Prefix(_)
2362 ));
2363 }
2364
2365 #[tokio::test]
2366 async fn import_headers_target_with_wildcard_rejected() {
2367 let msg = new_msg_body(X_HDR_CONTENT);
2368 msg.import_headers(vec![ImportHeaderSpec {
2369 name: "X-*".to_string(),
2370 target: Some("oops".to_string()),
2371 ..ImportHeaderSpec::default()
2372 }])
2373 .await
2374 .unwrap_err();
2375 }
2376
2377 #[tokio::test]
2378 async fn remove_all_x_headers() {
2379 let msg = new_msg_body(X_HDR_CONTENT);
2380
2381 msg.remove_x_headers(vec![]).await.unwrap();
2382 k9::assert_equal!(
2383 data_as_string(&msg),
2384 "Subject: Hello\r\nFrom :Someone\r\n\r\nBody"
2385 );
2386 }
2387
2388 #[tokio::test]
2389 async fn prepend_header_2_params() {
2390 let msg = new_msg_body(X_HDR_CONTENT);
2391
2392 msg.prepend_header(Some("Date"), b"Today").await.unwrap();
2393 k9::assert_equal!(
2394 data_as_string(&msg),
2395 "Date: Today\r\nX-Hello: there\r\nX-Header: value\r\nSubject: Hello\r\nFrom :Someone\r\n\r\nBody"
2396 );
2397 }
2398
2399 #[tokio::test]
2400 async fn prepend_header_1_params() {
2401 let msg = new_msg_body(X_HDR_CONTENT);
2402
2403 msg.prepend_header(None, b"Date: Today").await.unwrap();
2404 k9::assert_equal!(
2405 data_as_string(&msg),
2406 "Date: Today\r\nX-Hello: there\r\nX-Header: value\r\nSubject: Hello\r\nFrom :Someone\r\n\r\nBody"
2407 );
2408 }
2409
2410 #[tokio::test]
2411 async fn append_header_2_params() {
2412 let msg = new_msg_body(X_HDR_CONTENT);
2413
2414 msg.append_header(Some("Date"), b"Today").await.unwrap();
2415 k9::assert_equal!(
2416 data_as_string(&msg),
2417 "X-Hello: there\r\nX-Header: value\r\nSubject: Hello\r\nFrom :Someone\r\nDate: Today\r\n\r\nBody"
2418 );
2419 }
2420
2421 #[tokio::test]
2422 async fn append_header_1_params() {
2423 let msg = new_msg_body(X_HDR_CONTENT);
2424
2425 msg.append_header(None, b"Date: Today").await.unwrap();
2426 k9::assert_equal!(
2427 data_as_string(&msg),
2428 "X-Hello: there\r\nX-Header: value\r\nSubject: Hello\r\nFrom :Someone\r\nDate: Today\r\n\r\nBody"
2429 );
2430 }
2431
2432 const MULTI_HEADER_CONTENT: &str =
2433 "X-Hello: there\r\nX-Header: value\r\nSubject: Hello\r\nX-Header: another value\r\nFrom :Someone@somewhere\r\n\r\nBody";
2434
2435 #[tokio::test]
2436 async fn get_first_header() {
2437 let msg = new_msg_body(MULTI_HEADER_CONTENT);
2438 k9::assert_equal!(
2439 msg.get_first_named_header_value("X-header")
2440 .await
2441 .unwrap()
2442 .unwrap(),
2443 "value"
2444 );
2445 }
2446
2447 #[tokio::test]
2448 async fn get_all_header() {
2449 let msg = new_msg_body(MULTI_HEADER_CONTENT);
2450 k9::assert_equal!(
2451 msg.get_all_named_header_values("X-header").await.unwrap(),
2452 vec!["value".to_string(), "another value".to_string()]
2453 );
2454 }
2455
2456 #[tokio::test]
2457 async fn remove_first() {
2458 let msg = new_msg_body(MULTI_HEADER_CONTENT);
2459 msg.remove_first_named_header("X-header").await.unwrap();
2460 k9::assert_equal!(
2461 data_as_string(&msg),
2462 "X-Hello: there\r\nSubject: Hello\r\nX-Header: another value\r\nFrom :Someone@somewhere\r\n\r\nBody"
2463 );
2464 }
2465
2466 #[tokio::test]
2467 async fn remove_all() {
2468 let msg = new_msg_body(MULTI_HEADER_CONTENT);
2469 msg.remove_all_named_headers("X-header").await.unwrap();
2470 k9::assert_equal!(
2471 data_as_string(&msg),
2472 "X-Hello: there\r\nSubject: Hello\r\nFrom :Someone@somewhere\r\n\r\nBody"
2473 );
2474 }
2475
2476 #[tokio::test]
2477 async fn append_text_plain() {
2478 let msg = new_msg_body(MULTI_HEADER_CONTENT);
2479 msg.append_text_plain("I am at the bottom").await.unwrap();
2480 k9::assert_equal!(
2481 data_as_string(&msg),
2482 "X-Hello: there\r\n\
2483 X-Header: value\r\n\
2484 Subject: Hello\r\n\
2485 X-Header: another value\r\n\
2486 From :Someone@somewhere\r\n\
2487 Content-Type: text/plain;\r\n\
2488 \tcharset=\"us-ascii\"\r\n\
2489 \r\n\
2490 Body\r\n\
2491 I am at the bottom\r\n"
2492 );
2493 }
2494
2495 const MIXED_CONTENT: &str = "Content-Type: multipart/mixed;\r\n\
2496\tboundary=\"my-boundary\"\r\n\
2497\r\n\
2498--my-boundary\r\n\
2499Content-Type: text/plain;\r\n\
2500\tcharset=\"us-ascii\"\r\n\
2501\r\n\
2502plain text\r\n\
2503--my-boundary\r\n\
2504Content-Type: text/html;\r\n\
2505\tcharset=\"us-ascii\"\r\n\
2506\r\n\
2507<b>rich</b> text\r\n\
2508--my-boundary\r\n\
2509Content-Type: application/octet-stream\r\n\
2510Content-Transfer-Encoding: base64\r\n\
2511Content-Disposition: attachment;\r\n\
2512\tfilename=\"woot.bin\"\r\n\
2513Content-ID: <woot.id@somewhere>\r\n\
2514\r\n\
2515AAECAw==\r\n\
2516--my-boundary--\r\n\
2517\r\n";
2518
2519 const MIXED_CONTENT_ENCLOSING_BODY: &str = "Content-Type: multipart/mixed;\r\n\
2520\tboundary=\"my-boundary\"\r\n\
2521\r\n\
2522--my-boundary\r\n\
2523Content-Type: text/plain;\r\n\
2524\tcharset=\"us-ascii\"\r\n\
2525\r\n\
2526plain text\r\n\
2527--my-boundary\r\n\
2528Content-Type: text/html;\r\n\
2529\tcharset=\"us-ascii\"\r\n\
2530\r\n\
2531<BODY>\r\n\
2532<b>rich</b> text\r\n\
2533</BODY>\r\n\
2534--my-boundary\r\n\
2535Content-Type: application/octet-stream\r\n\
2536Content-Transfer-Encoding: base64\r\n\
2537Content-Disposition: attachment;\r\n\
2538\tfilename=\"woot.bin\"\r\n\
2539Content-ID: <woot.id>\r\n\
2540\r\n\
2541AAECAw==\r\n\
2542--my-boundary--\r\n\
2543\r\n";
2544
2545 #[tokio::test]
2546 async fn append_text_html() {
2547 let msg = new_msg_body(MIXED_CONTENT);
2548 msg.append_text_html("bottom html").await.unwrap();
2549 k9::snapshot!(
2550 data_as_string(&msg),
2551 r#"
2552Content-Type: multipart/mixed;\r
2553\tboundary="my-boundary"\r
2554\r
2555--my-boundary\r
2556Content-Type: text/plain;\r
2557\tcharset="us-ascii"\r
2558\r
2559plain text\r
2560--my-boundary\r
2561Content-Type: text/html;\r
2562\tcharset="us-ascii"\r
2563\r
2564<b>rich</b> text\r
2565\r
2566bottom html\r
2567--my-boundary\r
2568Content-Type: application/octet-stream\r
2569Content-Transfer-Encoding: base64\r
2570Content-Disposition: attachment;\r
2571\tfilename="woot.bin"\r
2572Content-ID: <woot.id@somewhere>\r
2573\r
2574AAECAw==\r
2575--my-boundary--\r
2576\r
2577
2578"#
2579 );
2580
2581 let msg = new_msg_body(MIXED_CONTENT_ENCLOSING_BODY);
2582 msg.append_text_html("bottom html 👻").await.unwrap();
2583 k9::snapshot!(
2584 data_as_string(&msg),
2585 r#"
2586Content-Type: multipart/mixed;\r
2587\tboundary="my-boundary"\r
2588\r
2589--my-boundary\r
2590Content-Type: text/plain;\r
2591\tcharset="us-ascii"\r
2592\r
2593plain text\r
2594--my-boundary\r
2595Content-Type: text/html;\r
2596\tcharset="utf-8"\r
2597Content-Transfer-Encoding: quoted-printable\r
2598\r
2599<BODY>\r
2600<b>rich</b> text\r
2601\r
2602bottom html =F0=9F=91=BB</BODY>\r
2603--my-boundary\r
2604Content-Type: application/octet-stream\r
2605Content-Transfer-Encoding: base64\r
2606Content-Disposition: attachment;\r
2607\tfilename="woot.bin"\r
2608Content-ID: <woot.id>\r
2609\r
2610AAECAw==\r
2611--my-boundary--\r
2612\r
2613
2614"#
2615 );
2616 }
2617
2618 #[tokio::test]
2619 async fn append_text_plain_mixed() {
2620 let msg = new_msg_body(MIXED_CONTENT);
2621 msg.append_text_plain("bottom text 👾").await.unwrap();
2622 k9::snapshot!(
2623 data_as_string(&msg),
2624 r#"
2625Content-Type: multipart/mixed;\r
2626\tboundary="my-boundary"\r
2627\r
2628--my-boundary\r
2629Content-Type: text/plain;\r
2630\tcharset="utf-8"\r
2631Content-Transfer-Encoding: quoted-printable\r
2632\r
2633plain text\r
2634\r
2635bottom text =F0=9F=91=BE\r
2636--my-boundary\r
2637Content-Type: text/html;\r
2638\tcharset="us-ascii"\r
2639\r
2640<b>rich</b> text\r
2641--my-boundary\r
2642Content-Type: application/octet-stream\r
2643Content-Transfer-Encoding: base64\r
2644Content-Disposition: attachment;\r
2645\tfilename="woot.bin"\r
2646Content-ID: <woot.id@somewhere>\r
2647\r
2648AAECAw==\r
2649--my-boundary--\r
2650\r
2651
2652"#
2653 );
2654 }
2655
2656 #[tokio::test]
2657 async fn check_conformance_angle_msg_id() {
2658 const DOUBLE_ANGLE_ONLY: &str = "Subject: hello\r
2659Message-ID: <<1234@example.com>>\r
2660\r
2661Hello";
2662 let msg = new_msg_body(DOUBLE_ANGLE_ONLY);
2663 k9::snapshot!(
2664 msg.check_fix_conformance(
2665 MessageConformance::MISSING_MESSAGE_ID_HEADER,
2666 MessageConformance::empty(),
2667 None,
2668 )
2669 .await
2670 .unwrap_err()
2671 .to_string(),
2672 "Message has conformance issues: MISSING_MESSAGE_ID_HEADER"
2673 );
2674
2675 msg.check_fix_conformance(
2676 MessageConformance::MISSING_MESSAGE_ID_HEADER,
2677 MessageConformance::MISSING_MESSAGE_ID_HEADER,
2678 None,
2679 )
2680 .await
2681 .unwrap();
2682
2683 const DOUBLE_ANGLE_AND_LONG_LINE: &str = "Subject: hello\r
2698Message-ID: <<1234@example.com>>\r
2699\r
2700Hello this is a really long line Hello this is a really long line \
2701Hello this is a really long line Hello this is a really long line \
2702Hello this is a really long line Hello this is a really long line \
2703Hello this is a really long line Hello this is a really long line \
2704Hello this is a really long line Hello this is a really long line \
2705Hello this is a really long line Hello this is a really long line \
2706Hello this is a really long line Hello this is a really long line
2707";
2708 let msg = new_msg_body(DOUBLE_ANGLE_AND_LONG_LINE);
2709 msg.check_fix_conformance(
2710 MessageConformance::MISSING_COLON_VALUE,
2711 MessageConformance::MISSING_MESSAGE_ID_HEADER | MessageConformance::LINE_TOO_LONG,
2712 None,
2713 )
2714 .await
2715 .unwrap();
2716
2717 }
2741
2742 #[tokio::test]
2743 async fn check_conformance() {
2744 let msg = new_msg_body(MULTI_HEADER_CONTENT);
2745 msg.check_fix_conformance(
2746 MessageConformance::default(),
2747 MessageConformance::MISSING_MIME_VERSION,
2748 None,
2749 )
2750 .await
2751 .unwrap();
2752 k9::snapshot!(
2753 data_as_string(&msg),
2754 r#"
2755X-Hello: there\r
2756X-Header: value\r
2757Subject: Hello\r
2758X-Header: another value\r
2759From :Someone@somewhere\r
2760Mime-Version: 1.0\r
2761\r
2762Body
2763"#
2764 );
2765
2766 let msg = new_msg_body(MULTI_HEADER_CONTENT);
2767 msg.check_fix_conformance(
2768 MessageConformance::default(),
2769 MessageConformance::MISSING_MIME_VERSION | MessageConformance::NAME_ENDS_WITH_SPACE,
2770 None,
2771 )
2772 .await
2773 .unwrap();
2774 k9::snapshot!(
2775 data_as_string(&msg),
2776 r#"
2777Content-Type: text/plain;\r
2778\tcharset="us-ascii"\r
2779X-Hello: there\r
2780X-Header: value\r
2781Subject: Hello\r
2782X-Header: another value\r
2783From: <Someone@somewhere>\r
2784Mime-Version: 1.0\r
2785\r
2786Body\r
2787
2788"#
2789 );
2790 }
2791
2792 #[tokio::test]
2793 async fn check_fix_latin_input() {
2794 const POUNDS: &[u8] = b"Subject: \xa3\r\n\r\nGBP\r\n";
2795 let msg = new_msg_body(&*POUNDS);
2796 msg.check_fix_conformance(
2797 MessageConformance::default(),
2798 MessageConformance::NEEDS_TRANSFER_ENCODING,
2799 Some(&CheckFixSettings {
2800 detect_encoding: true,
2801 include_encodings: vec!["iso-8859-1".to_string()],
2802 ..Default::default()
2803 }),
2804 )
2805 .await
2806 .unwrap();
2807
2808 let subject = msg
2809 .get_first_named_header_value("subject")
2810 .await
2811 .unwrap()
2812 .unwrap();
2813 assert_eq!(subject, "£");
2814 }
2815
2816 #[tokio::test]
2817 async fn set_scheduling() -> anyhow::Result<()> {
2818 let msg = new_msg_body(MULTI_HEADER_CONTENT);
2819 assert!(msg.get_due().is_none(), "due is implicitly now");
2820
2821 let now = Utc::now();
2822 let one_day = chrono::Duration::try_days(1).expect("1 day to be valid");
2823
2824 msg.set_scheduling(Some(Scheduling {
2825 restriction: None,
2826 first_attempt: Some((now + one_day).into()),
2827 expires: None,
2828 }))
2829 .await?;
2830
2831 let due = msg.get_due().expect("due to now be set");
2832 assert!(due - now >= one_day, "due time is at least 1 day away");
2833
2834 Ok(())
2835 }
2836
2837 #[cfg(all(test, target_pointer_width = "64"))]
2838 #[test]
2839 fn sizes() {
2840 assert_eq!(std::mem::size_of::<Message>(), 8);
2841 assert_eq!(std::mem::size_of::<MessageInner>(), 32);
2842 assert_eq!(std::mem::size_of::<MessageWithId>(), 72);
2843 }
2844}