Skip to main content

telelog_server/
lib.rs

1//! telelog-server as a library, so integration tests can run it in-process.
2
3pub mod archive;
4pub mod auth;
5mod service;
6
7use std::future::Future;
8use std::path::PathBuf;
9use std::sync::Arc;
10
11use anyhow::{Context as _, Result};
12use telelog_proto::LogServiceServer;
13use telelog_sources::Sources;
14use tokio::net::TcpListener;
15use tokio_stream::wrappers::TcpListenerStream;
16use tonic::transport::{Identity, Server, ServerTlsConfig};
17
18pub struct TlsFiles {
19    pub cert: PathBuf,
20    pub key: PathBuf,
21}
22
23pub struct Config {
24    pub token: Option<String>,
25    pub tls: Option<TlsFiles>,
26    /// When set, every source's output is archived to this bucket and can be queried.
27    pub archive: Option<Arc<archive::Archive>>,
28}
29
30/// Serves the gRPC API on `listener` until `shutdown` resolves.
31pub async fn serve(
32    listener: TcpListener,
33    config: Config,
34    sources: Sources,
35    shutdown: impl Future<Output = ()>,
36) -> Result<()> {
37    // object_store and tonic enable different rustls crypto backends; pick one explicitly.
38    let _ = rustls::crypto::ring::default_provider().install_default();
39    let mut builder = Server::builder();
40    if let Some(tls) = &config.tls {
41        let cert = std::fs::read(&tls.cert).with_context(|| format!("reading {}", tls.cert.display()))?;
42        let key = std::fs::read(&tls.key).with_context(|| format!("reading {}", tls.key.display()))?;
43        builder = builder
44            .tls_config(ServerTlsConfig::new().identity(Identity::from_pem(cert, key)))
45            .context("configuring TLS")?;
46    }
47    // Archiving runs for as long as the server does, whether or not any app is connected.
48    let background: Vec<_> = match &config.archive {
49        Some(archive) => vec![
50            tokio::spawn(archive::ingest::run(archive.clone(), sources.clone())),
51            tokio::spawn(archive::ingest::sweep_forever(archive.clone())),
52        ],
53        None => Vec::new(),
54    };
55    let service = LogServiceServer::with_interceptor(
56        service::Logs::new(sources, config.archive.clone()),
57        auth::TokenAuth::new(config.token),
58    );
59    builder
60        .add_service(service)
61        .serve_with_incoming_shutdown(TcpListenerStream::new(listener), shutdown)
62        .await?;
63    for task in background {
64        task.abort();
65    }
66    Ok(())
67}