Merge pull request 'Merge upstream v0.16.25' (#155) from merge/upstream-v0.16.25 into main
ci / fork-checks (push) Skipped
ci / build (push) Skipped
github/ci (branch) GitHub Actions
ci / github (push) Successful in 42m47s

This commit was merged in pull request #155.
This commit is contained in:
jcoffey-dev committed 2026-10-06 07:30:12 +00:00
commit b64f6690f6
60 files changed
+1593 -490

No files matched your search

+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "common"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
build = "build.rs"
+1 -1
View File
@@ -225,7 +225,7 @@ pub struct Caches {
pub dns_ipv6: CacheWithTtl<Box<str>, RecordSet<Ipv6Addr>>,
pub dns_tlsa: CacheWithTtl<Box<str>, Arc<Tlsa>>,
pub dns_mta_sts: CacheWithTtl<Box<str>, Arc<Policy>>,
pub dns_rbl: CacheWithTtl<Box<str>, Option<Arc<IpResolver>>>,
pub dns_rbl: CacheWithTtl<Box<str>, Option<Arc<[IpResolver]>>>,
pub negative_cache_ttl: Duration,
}
+14 -2
View File
@@ -227,10 +227,22 @@ impl WebApplicationManager {
let cached = if force_refresh {
None
} else {
server
match server
.blob_store()
.get_blob(self.blob_key.as_slice(), 0..usize::MAX)
.await?
.await
{
Ok(cached) => cached,
Err(err) => {
trc::event!(
Resource(trc::ResourceEvent::Error),
Reason = err,
Url = self.url.clone(),
Details = "Failed to read cached application bundle, downloading it again"
);
None
}
}
};
let is_cached = cached.is_some();
let bundle = match cached {
@@ -11,7 +11,7 @@ use quick_xml::Reader;
use quick_xml::XmlVersion;
use quick_xml::events::Event;
use registry::schema::{enums::ServiceProtocol, structs::Service};
use std::fmt::Write;
use std::{borrow::Cow, fmt::Write};
use utils::map::vec_map::VecMap;
impl Server {
@@ -20,32 +20,78 @@ impl Server {
body: Option<Vec<u8>>,
) -> trc::Result<Resource<Vec<u8>>> {
// Obtain parameters
let emailaddress = parse_autodiscover_request(body.as_deref().unwrap_or_default())
.map_err(|err| {
let request =
parse_autodiscover_request(body.as_deref().unwrap_or_default()).map_err(|err| {
trc::ResourceEvent::BadParameters
.into_err()
.details("Failed to parse autodiscover request")
.ctx(trc::Key::Reason, err)
})?;
// inbuxa: legacy-protocols LP-7, LP-14a
let legacy_off = match emailaddress.rsplit_once('@') {
let legacy_off = match request.email.rsplit_once('@') {
Some((_, domain)) => self.legacy_off_for(domain).await?,
None => self.legacy_off_for("").await?,
};
Ok(Resource::new(
"application/xml; charset=utf-8",
build_autodiscover_response(
&emailaddress,
let response = match request.response_schema {
ResponseSchema::Outlook => build_autodiscover_response(
&request.email,
&self.core.network.server_name,
&self.core.network.info.services,
|protocol| legacy_off.service(protocol),
)
.into_bytes(),
))
ResponseSchema::Unsupported => PROVIDER_NOT_AVAILABLE_RESPONSE.as_bytes().to_vec(),
};
Ok(Resource::new("application/xml; charset=utf-8", response))
}
}
const OUTLOOK_RESPONSE_SCHEMA: &str =
"http://schemas.microsoft.com/exchange/autodiscover/outlook/responseschema/2006a";
const PROVIDER_NOT_AVAILABLE_RESPONSE: &str = concat!(
"<?xml version=\"1.0\" encoding=\"UTF-8\"?>\n",
"<Autodiscover xmlns=\"http://schemas.microsoft.com/exchange/autodiscover/responseschema/2006\">\n",
"\t<Response>\n",
"\t\t<Error>\n",
"\t\t\t<ErrorCode>601</ErrorCode>\n",
"\t\t\t<Message>Provider is not available</Message>\n",
"\t\t\t<DebugData />\n",
"\t\t</Error>\n",
"\t</Response>\n",
"</Autodiscover>\n",
);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ResponseSchema {
Outlook,
Unsupported,
}
impl ResponseSchema {
fn parse(value: &str) -> Self {
if value.trim().eq_ignore_ascii_case(OUTLOOK_RESPONSE_SCHEMA) {
ResponseSchema::Outlook
} else {
ResponseSchema::Unsupported
}
}
}
#[derive(Debug, PartialEq, Eq)]
struct AutodiscoverRequest {
email: String,
response_schema: ResponseSchema,
}
#[derive(Clone, Copy)]
enum RequestField {
EmailAddress,
ResponseSchema,
}
fn build_autodiscover_response(
emailaddress: &str,
default_host: &str,
@@ -124,7 +170,7 @@ fn build_autodiscover_response(
config
}
fn parse_autodiscover_request(bytes: &[u8]) -> Result<String, String> {
fn parse_autodiscover_request(bytes: &[u8]) -> Result<AutodiscoverRequest, String> {
if bytes.is_empty() {
return Err("Empty request body".to_string());
}
@@ -132,8 +178,9 @@ fn parse_autodiscover_request(bytes: &[u8]) -> Result<String, String> {
let mut reader = Reader::from_reader(bytes);
reader.config_mut().trim_text(true);
let mut buf = Vec::with_capacity(128);
let mut value_buf = Vec::with_capacity(128);
'outer: for tag_name in ["Autodiscover", "Request", "EMailAddress"] {
'outer: for tag_name in ["Autodiscover", "Request"] {
loop {
match reader.read_event_into(&mut buf) {
Ok(Event::Start(e)) => {
@@ -143,30 +190,6 @@ fn parse_autodiscover_request(bytes: &[u8]) -> Result<String, String> {
.eq_ignore_ascii_case(found_tag_name.as_ref())
{
continue 'outer;
} else if tag_name == "EMailAddress" {
// Skip unsupported tags under Request, such as AcceptableResponseSchema
let mut tag_count = 0;
loop {
match reader.read_event_into(&mut buf) {
Ok(Event::End(_)) => {
if tag_count == 0 {
break;
} else {
tag_count -= 1;
}
}
Ok(Event::Start(_)) => {
tag_count += 1;
}
Ok(Event::Eof) => {
return Err(format!(
"Expected value, found unexpected EOF at position {}.",
reader.buffer_position()
));
}
_ => (),
}
}
} else {
return Err(format!(
"Expected tag {}, found unexpected tag {} at position {}.",
@@ -195,36 +218,170 @@ fn parse_autodiscover_request(bytes: &[u8]) -> Result<String, String> {
}
}
if let Ok(Event::Text(text)) = reader.read_event_into(&mut buf)
&& let Ok(text) = text.xml_content(XmlVersion::Implicit1_0)
&& text.contains('@')
{
return Ok(text.trim().to_lowercase());
let mut email = None;
let mut response_schema = ResponseSchema::Outlook;
loop {
match reader.read_event_into(&mut buf) {
Ok(Event::Start(e)) => {
let local_name = e.local_name();
let field = hashify::tiny_map_ignore_case!(local_name.as_ref(),
b"EMailAddress" => RequestField::EmailAddress,
b"AcceptableResponseSchema" => RequestField::ResponseSchema,
);
let value = match reader.read_event_into(&mut value_buf) {
Ok(Event::End(_)) => None,
Ok(event) => {
let value = match event {
Event::Text(text) => text
.xml_content(XmlVersion::Implicit1_0)
.ok()
.map(Cow::into_owned),
_ => None,
};
reader
.read_to_end_into(e.name(), &mut value_buf)
.map_err(|err| {
format!("Error at position {}: {:?}", reader.buffer_position(), err)
})?;
value
}
Err(err) => {
return Err(format!(
"Error at position {}: {:?}",
reader.buffer_position(),
err
));
}
};
match (field, value) {
(Some(RequestField::EmailAddress), Some(value)) => {
email = Some(value);
}
(Some(RequestField::ResponseSchema), Some(value)) => {
response_schema = ResponseSchema::parse(&value);
}
_ => (),
}
}
Ok(Event::End(_) | Event::Eof) => break,
Ok(_) => (),
Err(e) => {
return Err(format!(
"Error at position {}: {:?}",
reader.buffer_position(),
e
));
}
}
}
Err(format!(
"Expected email address, found unexpected value at position {}.",
reader.buffer_position()
))
match email {
Some(email) if email.contains('@') => Ok(AutodiscoverRequest {
email: email.trim().to_lowercase(),
response_schema,
}),
_ => Err(format!(
"Expected email address, found unexpected value at position {}.",
reader.buffer_position()
)),
}
}
#[cfg(test)]
mod tests {
use super::{AutodiscoverRequest, ResponseSchema, parse_autodiscover_request};
#[test]
fn parse_autodiscover() {
let r = r#"<?xml version="1.0" encoding="utf-8"?>
const OUTLOOK: &str =
"http://schemas.microsoft.com/exchange/autodiscover/outlook/responseschema/2006a";
const MOBILESYNC: &str =
"http://schemas.microsoft.com/exchange/autodiscover/mobilesync/responseschema/2006";
for (request, expected) in [
(
format!(
r#"<?xml version="1.0" encoding="utf-8"?>
<Autodiscover xmlns="http://schemas.microsoft.com/exchange/autodiscover/outlook/requestschema/2006">
<Request>
<EMailAddress>email@example.com</EMailAddress>
<AcceptableResponseSchema>http://schemas.microsoft.com/exchange/autodiscover/outlook/responseschema/2006a</AcceptableResponseSchema>
<EMailAddress>Email@Example.com</EMailAddress>
<AcceptableResponseSchema>{OUTLOOK}</AcceptableResponseSchema>
</Request>
</Autodiscover>"#;
</Autodiscover>"#
),
ResponseSchema::Outlook,
),
(
format!(
r#"<Autodiscover xmlns="http://schemas.microsoft.com/exchange/autodiscover/outlook/requestschema/2006">
<Request>
<AcceptableResponseSchema>{OUTLOOK}</AcceptableResponseSchema>
<EMailAddress>[email protected]</EMailAddress>
</Request>
</Autodiscover>"#
),
ResponseSchema::Outlook,
),
(
r#"<Autodiscover>
<Request>
<EMailAddress>[email protected]</EMailAddress>
</Request>
</Autodiscover>"#
.to_string(),
ResponseSchema::Outlook,
),
(
format!(
r#"<?xml version="1.0" encoding="utf-8"?>
<Autodiscover xmlns="http://schemas.microsoft.com/exchange/autodiscover/mobilesync/requestschema/2006">
<Request>
<EMailAddress>[email protected]</EMailAddress>
<AcceptableResponseSchema>{MOBILESYNC}</AcceptableResponseSchema>
</Request>
</Autodiscover>"#
),
ResponseSchema::Unsupported,
),
(
format!(
r#"<Autodiscover>
<Request>
<LegacyDN>/o=Example/ou=Users/cn=email</LegacyDN>
<Unknown><Nested>value</Nested><Empty/></Unknown>
<AcceptableResponseSchema>{MOBILESYNC}</AcceptableResponseSchema>
<EMailAddress>[email protected]</EMailAddress>
</Request>
</Autodiscover>"#
),
ResponseSchema::Unsupported,
),
] {
assert_eq!(
parse_autodiscover_request(request.as_bytes()).expect("valid request"),
AutodiscoverRequest {
email: "[email protected]".to_string(),
response_schema: expected,
},
"{request}"
);
}
assert_eq!(
super::parse_autodiscover_request(r.as_bytes()).unwrap(),
"[email protected]"
);
for request in [
"",
"<Autodiscover><Request></Request></Autodiscover>",
"<Autodiscover><Request><EMailAddress>no-domain</EMailAddress></Request></Autodiscover>",
"<Autodiscover><Request><EMailAddress>[email protected]</Request></Autodiscover>",
"<Request><EMailAddress>[email protected]</EMailAddress></Request>",
] {
assert!(
parse_autodiscover_request(request.as_bytes()).is_err(),
"{request}"
);
}
}
#[test]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "coordinator"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "dav-proto"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "dav"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "directory"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "email"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+2 -6
View File
@@ -20,7 +20,7 @@ use groupware::{
scheduling::{ItipError, ItipMessages},
};
use mail_parser::{
DateTime, Header, HeaderName, HeaderValue, Message, MessageParser, MimeHeaders, PartType,
Header, HeaderName, HeaderValue, Message, MessageParser, MimeHeaders, PartType,
parsers::fields::thread::thread_name,
};
use registry::{
@@ -924,11 +924,7 @@ impl EmailIngest for Server {
span_id: u64,
) {
if let Some(config) = &self.core.spam.classifier {
let mut dt = DateTime::from_timestamp(now() as i64);
dt.hour = 0;
dt.minute = 0;
dt.second = 0;
let until = dt.to_timestamp() as u64 + config.hold_samples_for;
let until = now() + config.hold_samples_for;
let sample = SpamTrainingSample {
account_id: Some(Id::from(account_id)),
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "groupware"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "http_proto"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "http"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "imap_proto"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+50 -2
View File
@@ -677,7 +677,13 @@ pub trait SerializeResponse {
impl SerializeResponse for trc::Error {
fn serialize(&self) -> Vec<u8> {
let mut buf = Vec::with_capacity(128);
if let Some(tag) = self.value_as_str(trc::Key::Id) {
if let Some(tag) = self
.keys()
.iter()
.rev()
.find_map(|(key, value)| (*key == trc::Key::Id).then_some(value))
.and_then(|value| value.as_str())
{
buf.extend_from_slice(tag.as_bytes());
} else {
buf.push(b'*');
@@ -813,9 +819,51 @@ impl Display for Command {
#[cfg(test)]
mod tests {
use crate::parser::parse_sequence_set;
use crate::protocol::ObjectId;
use crate::protocol::{ObjectId, SerializeResponse};
use types::id::Id;
#[test]
fn serialize_error_uses_command_tag() {
for (error, expected) in [
(
trc::AuthEvent::Failed.into_err().id("a1"),
"a1 NO [AUTHENTICATIONFAILED] ",
),
(
trc::AuthEvent::Failed
.into_err()
.ctx(trc::Key::Id, 7u32)
.id("a1"),
"a1 NO [AUTHENTICATIONFAILED] ",
),
(
trc::AuthEvent::Error
.into_err()
.ctx(trc::Key::Id, "12")
.id("a2"),
"a2 NO [AUTHENTICATIONFAILED] ",
),
(
trc::AuthEvent::TooManyAttempts
.into_err()
.caused_by(trc::AuthEvent::Failed.into_err().ctx(trc::Key::Id, 7u32))
.id("a3"),
"a3 NO [AUTHENTICATIONFAILED] ",
),
(
trc::AuthEvent::Failed.into_err().ctx(trc::Key::Id, 7u32),
"* NO [AUTHENTICATIONFAILED] ",
),
] {
let response = error.serialize();
assert!(
response.starts_with(expected.as_bytes()),
"{:?} does not start with {expected:?}",
String::from_utf8_lossy(&response)
);
}
}
#[test]
fn serialize_objectid_compound() {
// Empty compound
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "imap"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+4 -1
View File
@@ -92,7 +92,10 @@ impl<T: SessionStream> Session<T> {
auth_failures: auth_failures + 1,
};
} else {
return trc::AuthEvent::TooManyAttempts.into_err().caused_by(err);
return trc::AuthEvent::TooManyAttempts
.into_err()
.caused_by(err)
.id(tag.clone());
}
}
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "jmap_proto"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "jmap"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
@@ -80,9 +80,9 @@ pub(crate) async fn bootstrap_set(
mut set: RegistrySetResponse<'_>,
) -> trc::Result<RegistrySetResponse<'_>> {
if !set.server.registry().is_bootstrap_mode() {
set.fail_all_create("This operation is only allowed bootstrap mode");
set.fail_all_update("This operation is only allowed bootstrap mode");
set.fail_all_destroy("This operation is only allowed bootstrap mode");
set.fail_all_create("This operation is only allowed in bootstrap mode");
set.fail_all_update("This operation is only allowed in bootstrap mode");
set.fail_all_destroy("This operation is only allowed in bootstrap mode");
return Ok(set);
}
+20 -1
View File
@@ -102,6 +102,10 @@ pub(crate) async fn validate_domain(
let will_trigger_dkim = matches!(domain.dkim_management, DkimManagement::Automatic(_))
&& old_domain
.is_none_or(|old| !matches!(old.dkim_management, DkimManagement::Automatic(_)));
let will_schedule_dkim = !will_trigger_dkim
&& matches!(domain.dkim_management, DkimManagement::Automatic(_))
&& publishes_dkim(domain)
&& old_domain.is_some_and(|old| !publishes_dkim(old));
let will_trigger_acme = if let DnsManagement::Automatic(details) = &domain.dns_management
&& old_domain.is_none_or(|old| !matches!(old.dns_management, DnsManagement::Automatic(_)))
{
@@ -125,11 +129,19 @@ pub(crate) async fn validate_domain(
}));
on_success_renew_certificate
} else {
if will_schedule_dkim {
tasks.push(Task::DnsManagement(TaskDnsManagement {
domain_id: Id::default(),
update_records: Map::new(vec![DnsRecordType::Dkim]),
on_success_renew_certificate: false,
status: TaskStatus::now(),
}));
}
false
};
// Schedule DKIM key rotation task
if will_trigger_dkim {
if will_trigger_dkim || will_schedule_dkim {
tasks.push(Task::DkimManagement(TaskDomainManagement {
domain_id: Id::default(),
status: TaskStatus::now(),
@@ -176,6 +188,13 @@ pub(crate) async fn validate_domain(
Ok(Ok(response))
}
fn publishes_dkim(domain: &Domain) -> bool {
matches!(
&domain.dns_management,
DnsManagement::Automatic(details) if details.publish_records.contains(&DnsRecordType::Dkim)
)
}
pub(crate) async fn validate_dns_server(
set: &RegistrySetResponse<'_>,
dns: &mut DnsServer,
+1 -1
View File
@@ -7,7 +7,7 @@ keywords = ["imap", "jmap", "smtp", "email", "mail", "webdav", "server"]
categories = ["email"]
# Upstream offers AGPL-3.0-only OR LicenseRef-SEL; inbuxa takes the AGPL only.
license = "AGPL-3.0-only"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[[bin]]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "managesieve"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "migration"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "nlp"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "pop3"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "registry"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "scim-proto"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "scim"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "services"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+158 -44
View File
@@ -68,12 +68,14 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
"Domain is not set to automatic DKIM management".to_string(),
));
};
let mut create_signatures = dkim.algorithms.into_inner();
if create_signatures.is_empty() {
let configured = dkim.algorithms.into_inner();
if configured.is_empty() {
return Ok(TaskResult::permanent(
"No DKIM algorithms configured for domain".to_string(),
));
}
let mut create_signatures = configured.clone();
let mut active_types: Vec<DkimSignatureType> = Vec::with_capacity(configured.len());
let dns_updater = match domain.dns_management {
DnsManagement::Automatic(props) if props.publish_records.contains(&DnsRecordType::Dkim) => {
@@ -95,6 +97,7 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
let mut retire_signatures = Vec::new();
let mut retiring_signatures = Vec::new();
let mut delete_signatures = Vec::new();
let mut schedule_signatures = Vec::new();
let mut next_transition = None;
let signature_ids = server
@@ -123,16 +126,35 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
create_signatures.retain(|algo| algo != &key_algo);
publish_signatures.push(key)
}
DkimRotationStage::Active => retiring_signatures.push(key),
DkimRotationStage::Active if dns_updater.is_some() => retiring_signatures.push(key),
DkimRotationStage::Active => {
create_signatures.retain(|algo| algo != &key_algo);
active_types.push(key_algo);
}
DkimRotationStage::Retiring => retire_signatures.push(key),
DkimRotationStage::Retired => delete_signatures.push(key),
}
} else {
if key.object.is_active() {
create_signatures.retain(|algo| algo != &key_algo);
let transition = key.object.next_transition();
match key.object.stage() {
DkimRotationStage::Active => {
create_signatures.retain(|algo| algo != &key_algo);
active_types.push(key_algo);
if transition.is_none() && dns_updater.is_some() {
schedule_signatures.push(key);
}
}
DkimRotationStage::Pending => {
create_signatures.retain(|algo| algo != &key_algo);
if transition.is_none() || dns_updater.is_none() {
publish_signatures.push(key);
continue;
}
}
DkimRotationStage::Retiring | DkimRotationStage::Retired => {}
}
if let Some(transition) = key.object.next_transition()
if let Some(transition) = transition
&& next_transition.is_none_or(|next| transition < next)
{
next_transition = Some(transition);
@@ -142,6 +164,7 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
let now = now();
let mut do_refresh = false;
let mut temporary_errors = String::new();
for algorithm in create_signatures {
#[cfg(feature = "test_mode")]
@@ -221,7 +244,7 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
));
};
let propagation_target = txt_value.clone();
let published = updater
let published = match updater
.set_rrset(
origin,
&record.name,
@@ -229,12 +252,43 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
vec![record.record.clone()],
)
.await
.is_ok();
let signature_transition = if published
&& updater
.wait_for_txt_propagation(&record.name, origin, &propagation_target)
.await
{
Ok(_) => {
let propagated = updater
.wait_for_txt_propagation(&record.name, origin, &propagation_target)
.await;
if !propagated {
if !temporary_errors.is_empty() {
temporary_errors.push_str("; ");
}
let _ = write!(
&mut temporary_errors,
"DKIM record {} did not propagate, will retry.",
record.name
);
}
propagated
}
Err(err) => {
if !temporary_errors.is_empty() {
temporary_errors.push_str("; ");
}
let _ = write!(
&mut temporary_errors,
"Failed to publish DKIM record {}: {err}.",
record.name
);
trc::event!(
Dns(DnsEvent::RecordCreationFailed),
Hostname = record.name.clone(),
Details = origin.clone(),
Type = "TXT",
Reason = err,
);
false
}
};
let signature_transition = if published {
trc::event!(
Dkim(DkimEvent::SignaturePublished),
Id = selector.clone(),
@@ -246,7 +300,7 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
} else {
// Something went wrong, reschedule.
signature.set_stage(DkimRotationStage::Pending);
UTCDateTime::from_timestamp((now + 60) as i64) // Retry after 1 minute
UTCDateTime::from_timestamp(now as i64)
};
if next_transition.is_none_or(|next| signature_transition < next) {
@@ -257,12 +311,16 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
}
// Write key
let is_active = signature.is_active();
match server
.registry()
.write(RegistryWrite::insert(&signature.into()))
.await?
{
RegistryWriteResult::Success(_) => {
if is_active {
active_types.push(algorithm);
}
trc::event!(
Dkim(DkimEvent::SignatureCreated),
Id = selector,
@@ -277,11 +335,35 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
}
}
for signature in schedule_signatures {
let record = generate_dkim_dns_record_name(&signature.object, &domain.name);
let signature_transition =
UTCDateTime::from_timestamp((now + dkim.rotate_after.as_secs()) as i64);
if next_transition.is_none_or(|next| signature_transition < next) {
next_transition = Some(signature_transition);
}
let mut new_signature = signature.object.clone();
new_signature.set_next_transition(signature_transition);
if let SignatureUpdate::Failed(task_result) = update_signature(
server,
signature,
new_signature,
&record,
&mut temporary_errors,
)
.await?
{
return Ok(task_result);
}
}
// Publish signatures
let mut temporary_errors = String::new();
for signature in publish_signatures {
let record = generate_dkim_dns_record(&signature.object, &domain.name).await?;
if let Some((updater, origin)) = &dns_updater {
if let Some((updater, origin)) = &dns_updater {
for signature in publish_signatures {
let record = generate_dkim_dns_record(&signature.object, &domain.name).await?;
let dns_update::DnsRecord::TXT(txt_value) = &record.record else {
return Ok(TaskResult::permanent(
"DKIM record must be a TXT record".to_string(),
@@ -311,6 +393,7 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
next_transition = Some(signature_transition);
}
let signature_type = signature.object.object_type();
let mut new_signature = signature.object.clone();
new_signature.set_next_transition(signature_transition);
@@ -323,7 +406,7 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
);
// Write key
if let Some(task_result) = update_signature(
match update_signature(
server,
signature,
new_signature,
@@ -332,7 +415,9 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
)
.await?
{
return Ok(task_result);
SignatureUpdate::Written => active_types.push(signature_type),
SignatureUpdate::Conflict => {}
SignatureUpdate::Failed(task_result) => return Ok(task_result),
}
do_refresh = true;
}
@@ -357,20 +442,52 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
);
}
}
} else {
if !temporary_errors.is_empty() {
temporary_errors.push_str("; ");
}
} else {
for signature in publish_signatures {
let signature_type = signature.object.object_type();
if active_types.contains(&signature_type) || !configured.contains(&signature_type) {
continue;
}
let _ = write!(
let record = generate_dkim_dns_record_name(&signature.object, &domain.name);
let mut new_signature = signature.object.clone();
new_signature.set_stage(DkimRotationStage::Active);
match &mut new_signature {
DkimSignature::Dkim1Ed25519Sha256(sign) | DkimSignature::Dkim1RsaSha256(sign) => {
sign.next_transition_at = None
}
DkimSignature::Dkim2Ed25519Sha256(sign) | DkimSignature::Dkim2RsaSha256(sign) => {
sign.next_transition_at = None
}
}
match update_signature(
server,
signature,
new_signature,
&record,
&mut temporary_errors,
"No DNS server configured, cannot publish DKIM record {}.",
record.name
);
)
.await?
{
SignatureUpdate::Written => {
active_types.push(signature_type);
do_refresh = true;
}
SignatureUpdate::Conflict => active_types.push(signature_type),
SignatureUpdate::Failed(task_result) => return Ok(task_result),
}
}
}
// Retiring signatures
for signature in retiring_signatures {
let signature_type = signature.object.object_type();
if configured.contains(&signature_type) && !active_types.contains(&signature_type) {
continue;
}
let record = generate_dkim_dns_record_name(&signature.object, &domain.name);
let signature_transition =
UTCDateTime::from_timestamp((now + dkim.retire_after.as_secs()) as i64);
@@ -391,7 +508,7 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
);
// Write key
if let Some(task_result) = update_signature(
if let SignatureUpdate::Failed(task_result) = update_signature(
server,
signature,
new_signature,
@@ -406,9 +523,9 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
}
// Retire signatures
for signature in retire_signatures {
let record = generate_dkim_dns_record_name(&signature.object, &domain.name);
if let Some((updater, origin)) = &dns_updater {
if let Some((updater, origin)) = &dns_updater {
for signature in retire_signatures {
let record = generate_dkim_dns_record_name(&signature.object, &domain.name);
match updater
.set_rrset(origin, &record, dns_update::DnsRecordType::TXT, Vec::new())
.await
@@ -433,7 +550,7 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
);
// Write key
if let Some(task_result) = update_signature(
if let SignatureUpdate::Failed(task_result) = update_signature(
server,
signature,
new_signature,
@@ -458,15 +575,6 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
);
}
}
} else {
if !temporary_errors.is_empty() {
temporary_errors.push_str("; ");
}
let _ = write!(
&mut temporary_errors,
"No DNS server configured, cannot retire DKIM record {}.",
record
);
}
}
@@ -556,13 +664,19 @@ async fn dkim_management(server: &Server, task: &TaskDomainManagement) -> trc::R
}
}
enum SignatureUpdate {
Written,
Conflict,
Failed(TaskResult),
}
async fn update_signature(
server: &Server,
signature: RegistryObject<DkimSignature>,
new_signature: DkimSignature,
name: &str,
temporary_errors: &mut String,
) -> trc::Result<Option<TaskResult>> {
) -> trc::Result<SignatureUpdate> {
match server
.registry()
.write(RegistryWrite::update(
@@ -575,8 +689,8 @@ async fn update_signature(
))
.await
{
Ok(RegistryWriteResult::Success(_)) => Ok(None),
Ok(err) => Ok(Some(TaskResult::permanent(format!(
Ok(RegistryWriteResult::Success(_)) => Ok(SignatureUpdate::Written),
Ok(err) => Ok(SignatureUpdate::Failed(TaskResult::permanent(format!(
"Failed to write DKIM signature for record {name}: {err}"
)))),
Err(err) => {
@@ -588,7 +702,7 @@ async fn update_signature(
temporary_errors,
"Failed to write DKIM signature for record {name} due to concurrent modification, will retry.",
);
Ok(None)
Ok(SignatureUpdate::Conflict)
} else {
Err(err)
}
+1 -1
View File
@@ -6,7 +6,7 @@ homepage = "https://inbuxa.org"
keywords = ["smtp", "email", "mail", "server"]
categories = ["email"]
license = "AGPL-3.0-only OR LicenseRef-SEL"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+2 -2
View File
@@ -127,8 +127,8 @@ impl HasQueueQuota for Server {
refs: &mut Vec<Metadata>,
session_id: u64,
) -> bool {
if !quota.expr.is_empty()
&& self
if quota.expr.is_empty()
|| self
.eval_if(&quota.expr, envelope, session_id)
.await
.unwrap_or(false)
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "spam-filter"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+91 -65
View File
@@ -37,7 +37,7 @@ use std::{
hash::{Hash, RandomState},
sync::Arc,
};
use store::ahash::AHashSet;
use store::ahash::AHashMap;
use store::rand::seq::SliceRandom;
use store::write::{BlobLink, RegistryClass, now};
use store::{
@@ -92,7 +92,14 @@ struct TrainingTask {
sample: TrainingSample,
is_spam: bool,
is_replay: bool,
remove: Option<u64>,
remove: SampleRemoval,
}
#[derive(Debug, Clone, Copy)]
enum SampleRemoval {
Keep,
Item,
ItemAndLink { until: u64 },
}
#[derive(rkyv::Archive, rkyv::Deserialize, rkyv::Serialize, Debug)]
@@ -199,7 +206,7 @@ impl SpamClassifier for Server {
object_id,
item_id: u64::MAX,
}));
let mut seen_samples = AHashSet::new();
let mut seen_samples = AHashMap::new();
let mut spam_count = 0;
let mut ham_count = 0;
self.store()
@@ -220,43 +227,57 @@ impl SpamClassifier for Server {
.unwrap_or(u32::MAX),
};
if seen_samples.insert(sample.clone()) {
// Add to reservoir
if !do_remove {
trainer.reservoir.update_reservoir(
&sample,
match seen_samples.entry(sample.clone()) {
Entry::Vacant(entry) => {
entry.insert((!do_remove).then_some(until));
// Add to reservoir
if !do_remove {
trainer.reservoir.update_reservoir(
&sample,
is_spam,
config.reservoir_capacity,
);
} else {
trainer.reservoir.update_counts(is_spam);
}
samples.push(TrainingTask {
id,
sample,
is_spam,
config.reservoir_capacity,
);
} else {
trainer.reservoir.update_counts(is_spam);
is_replay: false,
remove: if do_remove {
SampleRemoval::ItemAndLink { until }
} else {
SampleRemoval::Keep
},
});
remove_entries |= do_remove;
// Update trainer stats
if is_spam {
spam_count += 1;
} else {
ham_count += 1;
}
}
samples.push(TrainingTask {
id,
sample,
is_spam,
is_replay: false,
remove: do_remove.then_some(until),
});
remove_entries |= do_remove;
// Update trainer stats
if is_spam {
spam_count += 1;
} else {
ham_count += 1;
Entry::Occupied(entry) => {
let remove = if *entry.get() == Some(until) {
SampleRemoval::Item
} else {
SampleRemoval::ItemAndLink { until }
};
duplicate_samples.push(TrainingTask {
id,
sample,
is_spam,
is_replay: false,
remove,
});
remove_entries = true;
}
} else {
duplicate_samples.push(TrainingTask {
id,
sample,
is_spam,
is_replay: false,
remove: Some(until),
});
remove_entries = true;
}
trainer.last_id = trainer.last_id.max(id);
@@ -315,7 +336,7 @@ impl SpamClassifier for Server {
sample: sample.clone(),
is_spam: false,
is_replay: true,
remove: None,
remove: SampleRemoval::Keep,
}),
);
} else if ham_count > spam_count {
@@ -329,7 +350,7 @@ impl SpamClassifier for Server {
sample: sample.clone(),
is_spam: true,
is_replay: true,
remove: None,
remove: SampleRemoval::Keep,
}),
);
}
@@ -884,33 +905,38 @@ async fn delete_samples(
let object_id = ObjectType::SpamTrainingSample.to_id();
let mut batch = BatchBuilder::new();
for sample in samples.into_iter().chain(duplicate_samples) {
if let Some(until) = sample.remove {
batch
.with_account_id(sample.sample.account_id)
.clear(BlobOp::Link {
hash: sample.sample.hash,
to: BlobLink::Temporary { until },
})
.clear(ValueClass::Registry(RegistryClass::Item {
object_id,
item_id: sample.id,
}))
.clear(ValueClass::Registry(RegistryClass::Index {
index_id: Property::AccountId.to_id(),
object_id,
item_id: sample.id,
key: (sample.sample.account_id as u64).serialize(),
}));
let until = match sample.remove {
SampleRemoval::Keep => continue,
SampleRemoval::Item => None,
SampleRemoval::ItemAndLink { until } => Some(until),
};
batch.with_account_id(sample.sample.account_id);
if let Some(until) = until {
batch.clear(BlobOp::Link {
hash: sample.sample.hash,
to: BlobLink::Temporary { until },
});
}
batch
.clear(ValueClass::Registry(RegistryClass::Item {
object_id,
item_id: sample.id,
}))
.clear(ValueClass::Registry(RegistryClass::Index {
index_id: Property::AccountId.to_id(),
object_id,
item_id: sample.id,
key: (sample.sample.account_id as u64).serialize(),
}));
if batch.is_large_batch() {
server
.store()
.write(batch.build_all())
.await
.caused_by(trc::location!())?;
batch = BatchBuilder::new();
batch.with_account_id(sample.sample.account_id);
}
if batch.is_large_batch() {
server
.store()
.write(batch.build_all())
.await
.caused_by(trc::location!())?;
batch = BatchBuilder::new();
batch.with_account_id(sample.sample.account_id);
}
}
if !batch.is_empty() {
+143 -95
View File
@@ -11,7 +11,14 @@ use common::{
config::mailstore::spamfilter::{DnsBlServer, Element, IpResolver, Location},
expr::functions::ResolveVariable,
};
use mail_auth::{Error, common::resolver::ToFqdn};
use mail_auth::common::resolver::ToFqdn;
#[cfg(not(feature = "test_mode"))]
use mail_auth::hickory_resolver::{
net::{DnsError, NetError},
proto::rr::{Name, RData},
};
#[cfg(feature = "test_mode")]
use mail_auth::{DnsError, Error};
use std::{
net::Ipv4Addr,
sync::Arc,
@@ -19,6 +26,18 @@ use std::{
};
use trc::SpamEvent;
const MAX_NEGATIVE_TTL: u32 = 3600;
enum DnsblAnswer {
Listed {
ips: Vec<Ipv4Addr>,
expires: Instant,
},
NotListed {
expires: Option<Instant>,
},
}
pub(crate) async fn check_dnsbl(
server: &Server,
ctx: &mut SpamFilterContext<'_>,
@@ -49,7 +68,7 @@ pub(crate) async fn check_dnsbl(
for dnsbl in &server.core.spam.dnsbl.servers {
if dnsbl.scope == scope
&& checks < max_checks
&& let Some(tag) = is_dnsbl(
&& let Some(codes) = dnsbl_codes(
server,
dnsbl,
SpamFilterResolver::new(ctx, resolver, location),
@@ -58,7 +77,19 @@ pub(crate) async fn check_dnsbl(
)
.await
{
ctx.result.add_tag(tag);
for code in codes.iter() {
let tag = server
.eval_if::<String, _>(
&dnsbl.tags,
&SpamFilterResolver::new(ctx, code, location),
ctx.input.span_id,
)
.await;
if let Some(tag) = tag {
ctx.result.add_tag(tag);
}
}
}
}
@@ -71,13 +102,13 @@ pub(crate) async fn check_dnsbl(
}
}
async fn is_dnsbl(
async fn dnsbl_codes(
server: &Server,
config: &DnsBlServer,
resolver: SpamFilterResolver<'_, impl ResolveVariable>,
element: Element,
checks: &mut usize,
) -> Option<String> {
) -> Option<Arc<[IpResolver]>> {
let time = Instant::now();
let zone = server
.eval_if::<String, _>(&config.zone, &resolver, resolver.ctx.input.span_id)
@@ -92,105 +123,122 @@ async fn is_dnsbl(
{
None
} else {
server
.eval_if(
&config.tags,
&SpamFilterResolver::new(
resolver.ctx,
&IpResolver::new(
format!("127.0.{}.{}", parts[1], parts[0]).parse().unwrap(),
),
resolver.location,
),
resolver.ctx.input.span_id,
)
.await
Some(Arc::from([IpResolver::new(
format!("127.0.{}.{}", parts[1], parts[0]).parse().unwrap(),
)]))
};
}
}
let result = match server.inner.cache.dns_rbl.get(zone.as_str()) {
Some(Some(result)) => result,
Some(None) => return None,
None => {
*checks += 1;
if let Some(codes) = server.inner.cache.dns_rbl.get(zone.as_str()) {
return codes;
}
match server
.core
.smtp
.resolvers
.dns
.ipv4_lookup_raw(zone.to_fqdn().as_ref())
.await
{
Ok(result) => {
trc::event!(
Spam(SpamEvent::Dnsbl),
Hostname = zone.clone(),
Result = result
.entry
.iter()
.map(|ip| trc::Value::from(ip.to_string()))
.collect::<Vec<_>>(),
Details = element.as_str(),
Elapsed = time.elapsed()
);
*checks += 1;
let entry = Arc::new(IpResolver::new(
result
.entry
.iter()
.copied()
.next()
.unwrap_or(Ipv4Addr::BROADCAST)
.into(),
));
match resolve_zone(server, zone.to_fqdn().as_ref()).await {
Ok(DnsblAnswer::Listed { ips, expires }) => {
trc::event!(
Spam(SpamEvent::Dnsbl),
Hostname = zone.clone(),
Result = ips
.iter()
.map(|ip| trc::Value::from(ip.to_string()))
.collect::<Vec<_>>(),
Details = element.as_str(),
Elapsed = time.elapsed()
);
server.inner.cache.dns_rbl.insert_with_expiry(
zone.into(),
Some(entry.clone()),
result.expires,
);
let codes: Arc<[IpResolver]> = ips
.into_iter()
.map(|ip| IpResolver::new(ip.into()))
.collect();
entry
}
Err(Error::Dns(mail_auth::DnsError::RecordNotFound(_))) => {
trc::event!(
Spam(SpamEvent::Dnsbl),
Hostname = zone.clone(),
Result = trc::Value::None,
Details = element.as_str(),
Elapsed = time.elapsed()
);
server.inner.cache.dns_rbl.insert_with_expiry(
zone.into(),
Some(codes.clone()),
expires,
);
server.inner.cache.dns_rbl.insert(
zone.into(),
None,
Duration::from_secs(86400),
);
return None;
}
Err(err) => {
trc::event!(
Spam(SpamEvent::DnsblError),
Hostname = zone,
Elapsed = time.elapsed(),
Details = element.as_str(),
CausedBy = err.to_string()
);
return None;
}
}
Some(codes)
}
};
Ok(DnsblAnswer::NotListed { expires }) => {
trc::event!(
Spam(SpamEvent::Dnsbl),
Hostname = zone.clone(),
Result = trc::Value::None,
Details = element.as_str(),
Elapsed = time.elapsed()
);
server
.eval_if(
&config.tags,
&SpamFilterResolver::new(resolver.ctx, result.as_ref(), resolver.location),
resolver.ctx.input.span_id,
)
.await
if let Some(expires) = expires {
server
.inner
.cache
.dns_rbl
.insert_with_expiry(zone.into(), None, expires);
}
None
}
Err(err) => {
trc::event!(
Spam(SpamEvent::DnsblError),
Hostname = zone,
Elapsed = time.elapsed(),
Details = element.as_str(),
CausedBy = err
);
None
}
}
}
#[cfg(not(feature = "test_mode"))]
async fn resolve_zone(server: &Server, zone: &str) -> Result<DnsblAnswer, String> {
let name = Name::from_str_relaxed(zone).map_err(|err| err.to_string())?;
match server.core.smtp.resolvers.dns.0.ipv4_lookup(name).await {
Ok(lookup) => {
let expires = lookup.valid_until();
let ips = lookup
.answers()
.iter()
.filter_map(|record| match &record.data {
RData::A(a) => Some(a.0),
_ => None,
})
.collect::<Vec<_>>();
Ok(if !ips.is_empty() {
DnsblAnswer::Listed { ips, expires }
} else {
DnsblAnswer::NotListed {
expires: Some(expires),
}
})
}
Err(NetError::Dns(DnsError::NoRecordsFound(no_records))) => Ok(DnsblAnswer::NotListed {
expires: no_records
.negative_ttl
.filter(|ttl| *ttl > 0)
.map(|ttl| Instant::now() + Duration::from_secs(ttl.min(MAX_NEGATIVE_TTL).into())),
}),
Err(err) => Err(err.to_string()),
}
}
#[cfg(feature = "test_mode")]
async fn resolve_zone(server: &Server, zone: &str) -> Result<DnsblAnswer, String> {
match server.core.smtp.resolvers.dns.ipv4_lookup_raw(zone).await {
Ok(result) => Ok(DnsblAnswer::Listed {
ips: result.entry.to_vec(),
expires: result.expires,
}),
Err(Error::Dns(DnsError::RecordNotFound(_))) => Ok(DnsblAnswer::NotListed {
expires: Some(Instant::now() + Duration::from_secs(MAX_NEGATIVE_TTL.into())),
}),
Err(err) => Err(err.to_string()),
}
}
+39 -35
View File
@@ -33,17 +33,9 @@ pub(crate) async fn pyzor_check(
message: &Message<'_>,
config: &PyzorConfig,
) -> trc::Result<Option<PyzorResponse>> {
// Make sure there is at least one text part
if !message
.parts
.iter()
.any(|p| matches!(p.body, PartType::Text(_) | PartType::Html(_)))
{
let Some(request) = message.pyzor_check_message() else {
return Ok(None);
}
// Hash message
let request = message.pyzor_check_message();
};
// Send message to address. inbuxa: in tests, a fixed table answers
// instead of a public server (test_response).
@@ -167,15 +159,15 @@ impl PyzorWrite for Sha1 {
}
trait PyzorDigest<W: PyzorWrite> {
fn pyzor_digest(&self, writer: W) -> W;
fn pyzor_digest(&self, writer: W) -> Option<W>;
}
pub trait PyzorCheck {
fn pyzor_check_message(&self) -> String;
fn pyzor_check_message(&self) -> Option<String>;
}
impl<W: PyzorWrite> PyzorDigest<W> for Message<'_> {
fn pyzor_digest(&self, writer: W) -> W {
fn pyzor_digest(&self, writer: W) -> Option<W> {
let parts = self
.parts
.iter()
@@ -191,7 +183,7 @@ impl<W: PyzorWrite> PyzorDigest<W> for Message<'_> {
}
impl PyzorCheck for Message<'_> {
fn pyzor_check_message(&self) -> String {
fn pyzor_check_message(&self) -> Option<String> {
let time = SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.map_or(0, |d| d.as_secs());
@@ -204,9 +196,9 @@ impl PyzorCheck for Message<'_> {
}
}
fn pyzor_create_message(message: &Message<'_>, time: u64, thread: u16) -> String {
fn pyzor_create_message(message: &Message<'_>, time: u64, thread: u16) -> Option<String> {
// Hash message
let hash = message.pyzor_digest(Sha1::new()).finalize().hex_encode();
let hash = message.pyzor_digest(Sha1::new())?.finalize().hex_encode();
// Hash key
let mut hash_key = Sha1::new();
hash_key.update("anonymous:".as_bytes());
@@ -226,10 +218,10 @@ fn pyzor_create_message(message: &Message<'_>, time: u64, thread: u16) -> String
sig.update(format!(":{time}:{hash_key}"));
let sig = sig.finalize().hex_encode();
format!("{message}\nSig: {sig}\n")
Some(format!("{message}\nSig: {sig}\n"))
}
fn pyzor_digest<'x, I, W>(mut writer: W, lines: I) -> W
fn pyzor_digest<'x, I, W>(mut writer: W, lines: I) -> Option<W>
where
I: Iterator<Item = &'x str>,
W: PyzorWrite,
@@ -291,6 +283,10 @@ where
}
}
if result.is_empty() {
return None;
}
if result.len() > ATOMIC_NUM_LINES {
for (offset, length) in DIGEST_SPEC {
for i in 0..*length {
@@ -305,7 +301,7 @@ where
}
}
writer
Some(writer)
}
fn html_to_text(input: &str) -> String {
@@ -489,7 +485,8 @@ mod test {
&MessageParser::new().parse(HTML_TEXT_STYLE_SCRIPT).unwrap(),
1697468672,
49005,
);
)
.unwrap();
assert_eq!(
message,
@@ -522,10 +519,9 @@ mod test {
"http://spammer.com/special-offers?buy=now",
] {
assert_eq!(
String::from_utf8(pyzor_digest(
Vec::new(),
format!("Test {strip_me} Test2").lines(),
))
String::from_utf8(
pyzor_digest(Vec::new(), format!("Test {strip_me} Test2").lines()).unwrap()
)
.unwrap(),
"TestTest2"
);
@@ -533,20 +529,26 @@ mod test {
// Test short lines
assert_eq!(
String::from_utf8(pyzor_digest(
Vec::new(),
concat!("This line is included\n", "not this\n", "This also").lines(),
))
String::from_utf8(
pyzor_digest(
Vec::new(),
concat!("This line is included\n", "not this\n", "This also").lines(),
)
.unwrap()
)
.unwrap(),
"ThislineisincludedThisalso"
);
// Test atomic
assert_eq!(
String::from_utf8(pyzor_digest(
Vec::new(),
"All this message\nShould be included\nIn the digest".lines(),
))
String::from_utf8(
pyzor_digest(
Vec::new(),
"All this message\nShould be included\nIn the digest".lines(),
)
.unwrap()
)
.unwrap(),
"AllthismessageShouldbeincludedInthedigest"
);
@@ -561,7 +563,7 @@ mod test {
expected += format!("Line{i}testtesttest").as_str();
}
assert_eq!(
String::from_utf8(pyzor_digest(Vec::new(), text.lines(),)).unwrap(),
String::from_utf8(pyzor_digest(Vec::new(), text.lines()).unwrap()).unwrap(),
expected
);
@@ -593,7 +595,8 @@ mod test {
MessageParser::new()
.parse(input)
.unwrap()
.pyzor_digest(Vec::new(),)
.pyzor_digest(Vec::new())
.unwrap()
)
.unwrap(),
expected,
@@ -606,7 +609,8 @@ mod test {
MessageParser::new()
.parse(HTML_TEXT_STYLE_SCRIPT)
.unwrap()
.pyzor_digest(Sha1::new(),)
.pyzor_digest(Sha1::new())
.unwrap()
.finalize()
.hex_encode(),
"b2c27325a034c581df0c9ef37e4a0d63208a3e7e",
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "store"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+4
View File
@@ -26,6 +26,8 @@ const CHURN_DELETION_WINDOW: usize = 4096;
const CHURN_DELETION_TRIGGER: usize = 1024;
const CHURN_DELETION_RATIO: f64 = 0.5;
const BYTES_PER_SYNC: u64 = 1024 * 1024;
const MAX_LOG_FILE_SIZE: usize = 10 * 1024 * 1024;
const KEEP_LOG_FILE_NUM: usize = 5;
#[derive(Clone, Copy)]
enum CfProfile {
@@ -118,6 +120,8 @@ impl RocksDbStore {
.set_db_write_buffer_size((config.buffer_size as usize).max(MIN_DB_WRITE_BUFFER_SIZE));
db_opts.set_bytes_per_sync(BYTES_PER_SYNC);
db_opts.set_wal_bytes_per_sync(BYTES_PER_SYNC);
db_opts.set_max_log_file_size(MAX_LOG_FILE_SIZE);
db_opts.set_keep_log_file_num(KEEP_LOG_FILE_NUM);
Ok(Store::RocksDb(Arc::new(RocksDbStore {
db: OptimisticTransactionDB::open_cf_descriptors(&db_opts, idx_path, cfs)
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "trc"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "event_macro"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[lib]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "types"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "utils"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "proc_macros"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[lib]