durable-actors 0.7.10

Standalone regional durable-actors control plane, host, and durability runtime
use super::storage::LtxFile;
use anyhow::{Context, Result, ensure};
use async_trait::async_trait;
use base64::engine::general_purpose::STANDARD;
use std::{
    io::{self, Read},
    time::{Duration, Instant},
};
use tokio_util::sync::CancellationToken;

#[async_trait]
pub(crate) trait LtxCompactor: Send + Sync {
    async fn compact(&self, files: &[LtxFile]) -> Result<LtxFile>;
}

pub(crate) struct RustCompactor;

#[async_trait]
impl LtxCompactor for RustCompactor {
    async fn compact(&self, files: &[LtxFile]) -> Result<LtxFile> {
        let files = files.to_vec();
        let cancel = CancellationToken::new();
        let _guard = cancel.clone().drop_guard();
        tokio::task::spawn_blocking(move || compact(files, cancel)).await?
    }
}

fn compact(files: Vec<LtxFile>, cancel: CancellationToken) -> Result<LtxFile> {
    let last = files.last().context("empty checkpoint")?.last;
    let deadline = Instant::now() + Duration::from_secs(120);
    let readers = files
        .into_iter()
        .map(|file| {
            Ok(CompactionInput {
                reader: base64::read::DecoderReader::new(
                    std::io::Cursor::new(file.data),
                    &STANDARD,
                ),
                cancel: cancel.clone(),
                deadline,
            })
        })
        .collect::<Result<Vec<_>>>()?;
    let mut data = crate::payload::Spool::new();
    let mut encoder = base64::write::EncoderWriter::new(&mut data, &STANDARD);
    let header = terse_ltx::compact(readers, &mut encoder).context("compact LTX checkpoint")?;
    ensure!(
        header.min_txid == 1 && header.max_txid == last,
        "compacted checkpoint transaction mismatch"
    );
    encoder.finish()?;
    drop(encoder);
    Ok(LtxFile {
        level: 9,
        first: header.min_txid,
        last: header.max_txid,
        data: crate::payload::Text::from_bytes(data.finish()?)?,
    })
}

struct CompactionInput {
    reader: base64::read::DecoderReader<
        'static,
        base64::engine::general_purpose::GeneralPurpose,
        std::io::Cursor<crate::payload::Text>,
    >,
    cancel: CancellationToken,
    deadline: Instant,
}

impl Read for CompactionInput {
    fn read(&mut self, bytes: &mut [u8]) -> io::Result<usize> {
        if self.cancel.is_cancelled() {
            return Err(io::Error::other("LTX compaction cancelled"));
        }
        if Instant::now() >= self.deadline {
            return Err(io::Error::new(
                io::ErrorKind::TimedOut,
                "LTX compaction timed out",
            ));
        }
        self.reader.read(bytes)
    }
}