durable-actors 0.7.10

Standalone regional durable-actors control plane, host, and durability runtime
use super::*;

pub(super) struct Segment {
    pub manifest: Manifest,
    writers: Vec<Box<dyn LogWriter>>,
    pub offset: u64,
    healthy: bool,
}

impl Segment {
    pub async fn open(
        storage: &LogStorage,
        stream: StateStream,
        first_version: u64,
    ) -> Result<Option<Self>> {
        let started = std::time::Instant::now();
        let actor = crate::storage_paths::actor_from_snapshot(&stream.object(first_version))?;
        let prefix = object_name(&crate::storage_paths::snapshots(&actor)?)?.replacen(
            "snapshots-",
            "logs-",
            1,
        );
        let prepared = storage.prepared.lock().unwrap().take();
        let prepared = match prepared {
            Some(prepared) => {
                ensure!(
                    prepared.prefix == prefix,
                    "prepared log belongs to another actor"
                );
                prepared
            }
            None => Prepared::new(prefix, storage.zones.clone()),
        };
        let mut manifest = Manifest {
            format: 1,
            id: prepared.id.clone(),
            stream,
            first_version,
            replicas: Vec::new(),
        };
        let Some(opened) = prepared.take().await? else {
            return Ok(None);
        };
        let mut writers = Vec::new();
        for (replica, writer) in opened {
            manifest.replicas.push(replica);
            writers.push(writer);
        }
        let key = manifest.key()?;
        let open_ms = started.elapsed().as_secs_f64() * 1000.0;
        ensure!(
            bounded(replace(
                storage.archive.as_ref(),
                &key,
                None,
                crate::payload::encode(&manifest)?
            ))
            .await?,
            "conflicting log manifest"
        );
        tracing::info!(
            event = "rapid_log_open",
            owner_epoch = manifest.stream.owner_epoch,
            open_ms,
            manifest_ms = started.elapsed().as_secs_f64() * 1000.0 - open_ms,
            duration_ms = started.elapsed().as_secs_f64() * 1000.0
        );
        Ok(Some(Self {
            manifest,
            writers,
            offset: 0,
            healthy: true,
        }))
    }

    pub async fn append(&mut self, version: u64, state: Bytes) -> Result<Bytes> {
        ensure!(
            self.healthy,
            "uncertain log append requires activation recovery"
        );
        let frame = Record { version, state }.encode()?;
        let expected = self
            .offset
            .checked_add(frame.len() as u64)
            .context("log offset overflow")?;
        // An interrupted flush may have persisted bytes; this pair must never be reused.
        self.healthy = false;
        let flushed = bounded(async {
            Ok(futures_util::future::join_all(
                self.writers
                    .iter_mut()
                    .map(|writer| writer.append_and_flush(frame.clone())),
            )
            .await)
        })
        .await?;
        for result in flushed {
            ensure!(result? == expected, "unexpected persisted log offset");
        }
        self.offset = expected;
        self.healthy = true;
        Ok(frame)
    }
}