lenso 0.5.27

Rust authoring and Host facade for Lenso applications and Plugins.
Documentation
#![cfg(feature = "host")]

use futures::FutureExt;
use lenso::host::{
    self, CanonicalDocument, ControlPlaneError, GenerationRuntime, HostBuilder,
    MemoryControlStateStore, ResolvedGeneration,
};
use lenso_app_plan::ResolvedAppPlan;
use lenso_plugin_control_plane::{
    AppGenerationSpec, AppGenerationTransitionSpec, EffectiveHostGrantSet, ReplacementMode,
    ResolvedArtifactSet, RolloutPolicy,
};
use lenso_runtime_codec::{ArtifactCatalog, InstanceResourceCatalog};
use serde_json::{Value, json};
use std::{
    collections::BTreeMap,
    sync::{
        Arc,
        atomic::{AtomicUsize, Ordering},
    },
    time::Duration,
};
use tokio::io::{AsyncReadExt, AsyncWriteExt};

#[derive(Debug)]
struct Runtime;
impl GenerationRuntime for Runtime {
    type Handle = ();
    type Route = ();
    fn stage<'a>(
        &'a mut self,
        _: &'a ResolvedGeneration,
        _: u64,
    ) -> futures::future::LocalBoxFuture<'a, Result<(), ControlPlaneError>> {
        async { Ok(()) }.boxed_local()
    }
    fn shutdown(
        &mut self,
        (): (),
        _: u64,
    ) -> futures::future::LocalBoxFuture<'_, Result<(), ControlPlaneError>> {
        async { Ok(()) }.boxed_local()
    }
    fn terminal_failure(&self, (): &()) -> Option<ControlPlaneError> {
        None
    }
    fn route(&self, (): &()) {}
}

fn generation() -> ResolvedGeneration {
    generation_with_authority("authority")
}

fn generation_with_authority(authority: &str) -> ResolvedGeneration {
    let artifacts = CanonicalDocument::from_value(
        "artifacts",
        ResolvedArtifactSet {
            schema_version: 3,
            resolution_authority_digest: authority.into(),
            host_execution_policy_digest: "policy".into(),
            artifacts: vec![],
            instance_resources: vec![],
        },
    )
    .unwrap();
    let grants = CanonicalDocument::from_value(
        "grants",
        EffectiveHostGrantSet {
            schema_version: 2,
            resolution_authority_digest: authority.into(),
            grants: vec![],
        },
    )
    .unwrap();
    let spec = CanonicalDocument::from_value(
        "generation",
        AppGenerationSpec {
            schema_version: 2,
            app_id: "example.app".into(),
            host_build_manifest_digest: "host".into(),
            host_execution_policy_digest: "policy".into(),
            resolved_plan_digest: "plan".into(),
            resolution_authority_digest: authority.into(),
            resolved_artifact_set_digest: artifacts.digest().into(),
            effective_host_grant_set_digest: grants.digest().into(),
        },
    )
    .unwrap();
    ResolvedGeneration {
        plan: ResolvedAppPlan::new(vec![], vec![]),
        artifacts: ArtifactCatalog::new(),
        resources: InstanceResourceCatalog::new(),
        stateful_instances: BTreeMap::new(),
        artifact_set: artifacts,
        grants,
        spec,
    }
}

async fn start_host(
    candidate: ResolvedGeneration,
) -> Result<(host::Host<()>, ResolvedGeneration), ControlPlaneError> {
    let app =
        HostBuilder::new("example.app", Runtime, MemoryControlStateStore::default()).build()?;
    let transition = CanonicalDocument::from_value(
        "transition",
        AppGenerationTransitionSpec {
            schema_version: 1,
            app_id: "example.app".into(),
            from_generation_spec_digest: None,
            to_generation_spec_digest: candidate.spec.digest().into(),
            replacement_mode: ReplacementMode::Initial,
            state_compatibility_receipt_digests: vec![],
            rollout_policy: RolloutPolicy {
                ready_timeout_nanos: "1000000000".into(),
                drain_timeout_nanos: "1000000000".into(),
                rollback_window_nanos: "0".into(),
                automatic_rollback_on_generation_failure: false,
            },
        },
    )
    .unwrap();
    app.transition(transition, candidate.clone(), BTreeMap::new())
        .await?;
    Ok((app, candidate))
}

async fn write(stream: &mut tokio::io::DuplexStream, value: Value) {
    let bytes = serde_json::to_vec(&value).unwrap();
    stream
        .write_u32(u32::try_from(bytes.len()).unwrap())
        .await
        .unwrap();
    stream.write_all(&bytes).await.unwrap();
}
async fn read(stream: &mut tokio::io::DuplexStream) -> Value {
    let len = stream.read_u32().await.unwrap();
    let mut bytes = vec![0; len as usize];
    stream.read_exact(&mut bytes).await.unwrap();
    serde_json::from_slice(&bytes).unwrap()
}

#[tokio::test(flavor = "current_thread")]
async fn handshake_ready_inspect_and_repeatable_stop_use_one_runtime_authority() {
    tokio::task::LocalSet::new().run_until(async {
    let (client,server)=tokio::io::duplex(4096); let (read_half,write_half)=tokio::io::split(server);
    let candidate=generation();
    let task=tokio::task::spawn_local(host::control::serve(host::control::ControlOptions { distribution:"dist-v1".into(), startup_timeout:Duration::from_secs(1), stop_timeout:Duration::from_secs(1) }, read_half, write_half, move || start_host(candidate)));
    let mut client=client;
    write(&mut client,json!({"op":"start","version":1,"id":1,"distribution":"dist-v1"})).await;
    assert_eq!(read(&mut client).await["kind"],"ready");
    let started=read(&mut client).await; assert_eq!(started["kind"],"started"); let revision=started["revision"].clone();
    write(&mut client,json!({"op":"inspect","version":1,"id":2,"revision":revision,"offset":0,"limit":64})).await;
    assert_eq!(read(&mut client).await["kind"],"inspected");
    write(&mut client,json!({"op":"stop","version":1,"id":3})).await;
    let terminal=read(&mut client).await; assert_eq!(terminal["shutdown"],"suspended"); assert_eq!(terminal["id"],3);
    task.await.unwrap().unwrap();
 }).await;
}

#[tokio::test(flavor = "current_thread")]
async fn invalid_start_and_queued_overflow_never_open_readiness() {
    tokio::task::LocalSet::new()
        .run_until(async {
            let (client, server) = tokio::io::duplex(8192);
            let (read_half, write_half) = tokio::io::split(server);
            let task = tokio::task::spawn_local(host::control::serve(
                host::control::ControlOptions {
                    distribution: "dist-v1".into(),
                    startup_timeout: Duration::from_secs(1),
                    stop_timeout: Duration::from_secs(1),
                },
                read_half,
                write_half,
                must_not_start,
            ));
            let mut client = client;
            write(
                &mut client,
                json!({"op":"start","version":2,"id":1,"distribution":"dist-v1"}),
            )
            .await;
            assert_eq!(read(&mut client).await["kind"], "start_failed");
            task.await.unwrap().unwrap();
        })
        .await;
}

async fn must_not_start() -> Result<(host::Host<()>, ResolvedGeneration), ControlPlaneError> {
    panic!("invalid handshake must not execute runtime")
}

#[tokio::test(flavor = "current_thread")]
async fn failed_activation_receipt_is_reported_and_retryable_without_replacing_route() {
    tokio::task::LocalSet::new()
        .run_until(async {
            let (client, server) = tokio::io::duplex(4096);
            let (reader, writer) = tokio::io::split(server);
            let candidate = generation();
            let retry_candidate = candidate.clone();
            let attempts = Arc::new(AtomicUsize::new(0));
            let receipt_attempts = attempts.clone();
            let task = tokio::task::spawn_local(host::control::serve_with_reconcile_and_receipt(
                host::control::ControlOptions {
                    distribution: "dist-v1".into(),
                    startup_timeout: Duration::from_secs(1),
                    stop_timeout: Duration::from_secs(1),
                },
                reader,
                writer,
                move || start_host(candidate),
                move || {
                    let candidate = retry_candidate.clone();
                    async move { Ok(candidate) }
                },
                move |_| {
                    let attempt = receipt_attempts.fetch_add(1, Ordering::SeqCst);
                    async move {
                        if attempt == 0 {
                            Err(ControlPlaneError::HostFailure {
                                detail: "receipt unavailable".into(),
                            })
                        } else {
                            Ok(())
                        }
                    }
                },
            ));
            let mut client = client;
            write(
                &mut client,
                json!({"op":"start","version":1,"id":1,"distribution":"dist-v1"}),
            )
            .await;
            assert_eq!(read(&mut client).await["kind"], "ready");
            let started = read(&mut client).await;
            assert_eq!(started["activation_recorded"], false);
            write(
                &mut client,
                json!({"op":"reconcile","version":1,"id":2,"revision":started["revision"]}),
            )
            .await;
            let reconciled = read(&mut client).await;
            assert_eq!(reconciled["changed"], false);
            assert_eq!(reconciled["activation_recorded"], true);
            assert_eq!(attempts.load(Ordering::SeqCst), 2);
            write(&mut client, json!({"op":"stop","version":1,"id":3})).await;
            assert_eq!(read(&mut client).await["shutdown"], "suspended");
            task.await.unwrap().unwrap();
        })
        .await;
}

#[tokio::test(flavor = "current_thread")]
#[allow(
    clippy::too_many_lines,
    reason = "one ordered control session proves stale, failed, switched, and idempotent outcomes"
)]
async fn reconcile_keeps_old_route_on_source_failure_then_switches_exact_candidate() {
    tokio::task::LocalSet::new()
        .run_until(async {
            let (client, server) = tokio::io::duplex(8192);
            let (read_half, write_half) = tokio::io::split(server);
            let initial = generation();
            let replacement = generation_with_authority("authority-v2");
            let replacement_digest = replacement.spec.digest().to_owned();
            let mut invalid = generation_with_authority("authority-invalid");
            let mut invalid_spec = invalid.spec.value().clone();
            invalid_spec.app_id = "other.app".into();
            invalid.spec = CanonicalDocument::from_value("generation", invalid_spec).unwrap();
            let mut attempts = 0;
            let task = tokio::task::spawn_local(host::control::serve_with_reconcile(
                host::control::ControlOptions {
                    distribution: "dist-v1".into(),
                    startup_timeout: Duration::from_secs(1),
                    stop_timeout: Duration::from_secs(1),
                },
                read_half,
                write_half,
                move || start_host(initial),
                move || {
                    attempts += 1;
                    let candidate = replacement.clone();
                    let invalid = invalid.clone();
                    async move {
                        match attempts {
                            1 => Err(ControlPlaneError::HostFailure {
                                detail: "source unavailable".into(),
                            }),
                            2 => Ok(invalid),
                            _ => Ok(candidate),
                        }
                    }
                },
            ));
            let mut client = client;
            write(
                &mut client,
                json!({"op":"start","version":1,"id":1,"distribution":"dist-v1"}),
            )
            .await;
            assert_eq!(read(&mut client).await["kind"], "ready");
            let started = read(&mut client).await;
            let revision = started["revision"].as_u64().unwrap();
            write(
                &mut client,
                json!({"op":"reconcile","version":1,"id":2,"revision":revision - 1}),
            )
            .await;
            assert_eq!(read(&mut client).await["code"], "snapshot_expired");
            write(
                &mut client,
                json!({"op":"reconcile","version":1,"id":3,"revision":revision}),
            )
            .await;
            assert_eq!(read(&mut client).await["code"], "candidate_unavailable");
            write(
                &mut client,
                json!({"op":"inspect","version":1,"id":4,"revision":revision,"offset":0,"limit":1}),
            )
            .await;
            let still_old = read(&mut client).await;
            assert_eq!(still_old["kind"], "inspected");
            assert_ne!(still_old["generation"], replacement_digest);
            write(
                &mut client,
                json!({"op":"reconcile","version":1,"id":5,"revision":revision}),
            )
            .await;
            assert_eq!(
                read(&mut client).await["code"],
                "candidate_not_ready_or_incompatible"
            );
            write(
                &mut client,
                json!({"op":"inspect","version":1,"id":6,"revision":null,"offset":0,"limit":1}),
            )
            .await;
            let still_old = read(&mut client).await;
            assert_eq!(still_old["kind"], "inspected");
            assert_ne!(still_old["generation"], replacement_digest);
            let revision = still_old["revision"].as_u64().unwrap();
            write(
                &mut client,
                json!({"op":"reconcile","version":1,"id":7,"revision":revision}),
            )
            .await;
            let switched = read(&mut client).await;
            assert_eq!(switched["kind"], "reconciled");
            assert_eq!(switched["changed"], true);
            assert_eq!(switched["generation"], replacement_digest);
            let new_revision = switched["revision"].as_u64().unwrap();
            assert!(new_revision > revision);
            write(
                &mut client,
                json!({"op":"reconcile","version":1,"id":8,"revision":revision}),
            )
            .await;
            assert_eq!(read(&mut client).await["code"], "snapshot_expired");
            write(
                &mut client,
                json!({"op":"inspect","version":1,"id":9,"revision":null,"offset":0,"limit":1}),
            )
            .await;
            let refreshed = read(&mut client).await;
            assert_eq!(refreshed["kind"], "inspected");
            let current_revision = refreshed["revision"].as_u64().unwrap();
            write(
                &mut client,
                json!({"op":"reconcile","version":1,"id":10,"revision":current_revision}),
            )
            .await;
            let unchanged = read(&mut client).await;
            assert_eq!(unchanged["changed"], false);
            assert_eq!(unchanged["revision"], current_revision);
            write(&mut client, json!({"op":"stop","version":1,"id":11})).await;
            assert_eq!(read(&mut client).await["shutdown"], "suspended");
            task.await.unwrap().unwrap();
        })
        .await;
}