AI models

This commit is contained in:
mdecimus
2024-10-05 19:05:04 +02:00
parent d5c2dcb817
commit d0ce2b1a96
41 changed files with 1171 additions and 406 deletions

View File

@@ -10,23 +10,18 @@
use std::sync::Arc;
use ahash::AHashMap;
use base64::{engine::general_purpose, Engine};
use common::{config::server::Listeners, listener::SessionData, Core, Data, Inner};
use directory::{backend::internal::PrincipalField, QueryBy};
use hyper::{body, server::conn::http1, service::service_fn, Method, StatusCode, Uri};
use hyper_util::rt::TokioIo;
use jmap::api::{
http::{fetch_body, ToHttpResponse},
HttpResponse, JsonResponse,
};
use hyper::{Method, StatusCode};
use jmap::api::{http::ToHttpResponse, JsonResponse};
use mail_send::Credentials;
use serde_json::json;
use tokio::sync::watch;
use trc::{AuthEvent, EventType};
use utils::config::Config;
use crate::{add_test_certs, directory::DirectoryTest, AssertConfig};
use crate::{
directory::DirectoryTest,
http_server::{spawn_mock_http_server, HttpMessage},
};
static TEST_TOKEN: &str = "eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJzdWIiOiIxMjM0NTY3ODkwIiwibmFtZSI6IkpvaG4gRG9lIiwiaWF0IjoxNTE2MjM5MDIyfQ";
@@ -143,121 +138,3 @@ async fn oidc_directory() {
assert_eq!(principal.description(), Some("John Doe"));
}
}
const MOCK_HTTP_SERVER: &str = r#"
[server]
hostname = "'oidc.example.org'"
http.url = "'https://127.0.0.1:9090'"
[server.listener.jmap]
bind = ['127.0.0.1:9090']
protocol = 'http'
tls.implicit = true
[server.socket]
reuse-addr = true
[certificate.default]
cert = '%{file:{CERT}}%'
private-key = '%{file:{PK}}%'
default = true
"#;
#[derive(Clone)]
pub struct HttpSessionManager {
inner: HttpRequestHandler,
}
pub type HttpRequestHandler = Arc<dyn Fn(HttpMessage) -> HttpResponse + Sync + Send>;
#[derive(Debug)]
pub struct HttpMessage {
method: Method,
headers: AHashMap<String, String>,
uri: Uri,
body: Option<Vec<u8>>,
}
impl HttpMessage {
pub fn get_url_encoded(&self, key: &str) -> Option<String> {
form_urlencoded::parse(self.body.as_ref()?.as_slice())
.find(|(k, _)| k == key)
.map(|(_, v)| v.into_owned())
}
}
pub async fn spawn_mock_http_server(
handler: HttpRequestHandler,
) -> (watch::Sender<bool>, watch::Receiver<bool>) {
// Start mock push server
let mut settings = Config::new(add_test_certs(MOCK_HTTP_SERVER)).unwrap();
settings.resolve_all_macros().await;
let mock_inner = Arc::new(Inner {
shared_core: Core::parse(&mut settings, Default::default(), Default::default())
.await
.into_shared(),
data: Data::parse(&mut settings),
..Default::default()
});
settings.errors.clear();
settings.warnings.clear();
let mut servers = Listeners::parse(&mut settings);
servers.parse_tcp_acceptors(&mut settings, mock_inner.clone());
// Start JMAP server
servers.bind_and_drop_priv(&mut settings);
settings.assert_no_errors();
servers.spawn(|server, acceptor, shutdown_rx| {
server.spawn(
HttpSessionManager {
inner: handler.clone(),
},
mock_inner.clone(),
acceptor,
shutdown_rx,
);
})
}
impl common::listener::SessionManager for HttpSessionManager {
#[allow(clippy::manual_async_fn)]
fn handle<T: common::listener::SessionStream>(
self,
session: SessionData<T>,
) -> impl std::future::Future<Output = ()> + Send {
async move {
let sender = self.inner;
let _ = http1::Builder::new()
.keep_alive(false)
.serve_connection(
TokioIo::new(session.stream),
service_fn(|mut req: hyper::Request<body::Incoming>| {
let sender = sender.clone();
async move {
let response = sender(HttpMessage {
method: req.method().clone(),
uri: req.uri().clone(),
headers: req
.headers()
.iter()
.map(|(k, v)| {
(k.as_str().to_lowercase(), v.to_str().unwrap().to_string())
})
.collect(),
body: fetch_body(&mut req, 1024 * 1024, 0).await,
});
Ok::<_, hyper::Error>(response.build())
}
}),
)
.await;
}
}
#[allow(clippy::manual_async_fn)]
fn shutdown(&self) -> impl std::future::Future<Output = ()> + Send {
async {}
}
}