lenso-native-adapter 0.3.21

Native Rust Execution Adapter for Lenso.
Documentation
use std::{
    cell::{Cell, RefCell},
    rc::Rc,
    time::Duration,
};

use futures::channel::oneshot;
use lenso_app_plan::{PluginInstancePlan, ResolvedAppPlan};
use lenso_kernel::{
    DeterministicDriver, Kernel, ManagedTaskError, RuntimeDriver, RuntimeFailure, ShutdownOutcome,
};
use lenso_native_adapter::{
    CompleteObjectLifecycle, LifecycleContext, ManagedTasksError, NativePluginFactory,
    NativePluginFactoryContext, NativePluginInstance, NativePluginRegistry, PluginObject,
};

#[derive(Debug, Default)]
struct TaskProbe {
    ready_runs: Cell<u32>,
    cancellation_runs: Cell<u32>,
    task_drops: Cell<u32>,
    stop_runs: Cell<u32>,
}

struct TaskLifetime(Rc<TaskProbe>);

impl Drop for TaskLifetime {
    fn drop(&mut self) {
        self.0.task_drops.set(self.0.task_drops.get() + 1);
    }
}

#[lenso_native_adapter::plugin]
#[derive(Debug)]
struct ConstructingTaskPlugin {
    lifecycle: LifecycleContext,
    probe: Rc<TaskProbe>,
    ready_signal: RefCell<Option<oneshot::Receiver<()>>>,
}

#[lenso_native_adapter::plugin_impl]
impl ConstructingTaskPlugin {
    #[create]
    fn create(#[lifecycle] lifecycle: LifecycleContext) -> Result<Self, &'static str> {
        let ready = lifecycle
            .readiness()
            .map_err(|_| "no construction readiness")?;
        assert!(!ready.is_open());
        let probe = Rc::new(TaskProbe::default());
        let task_probe = probe.clone();
        let lifetime = TaskLifetime(probe.clone());
        let (started, ready_signal) = oneshot::channel();
        lifecycle
            .spawn_local(async move {
                let _lifetime = lifetime;
                ready.wait().await;
                task_probe.ready_runs.set(task_probe.ready_runs.get() + 1);
                let _ = started.send(());
                ready.cancellation().cancelled().await;
                task_probe
                    .cancellation_runs
                    .set(task_probe.cancellation_runs.get() + 1);
            })
            .map_err(|_| "construction task was rejected")?;

        assert_eq!(probe.ready_runs.get(), 0);
        Ok(Self {
            lifecycle,
            probe,
            ready_signal: RefCell::new(Some(ready_signal)),
        })
    }

    #[stop]
    fn stop(&self, #[lifecycle] lifecycle: LifecycleContext) {
        assert!(matches!(
            lifecycle.readiness(),
            Err(RuntimeFailure::AdmissionClosed)
        ));
        assert!(matches!(
            lifecycle.spawn_local(async {}),
            Err(ManagedTasksError::Inactive)
        ));
        self.probe.stop_runs.set(self.probe.stop_runs.get() + 1);
        drop(lifecycle);
    }
}

#[derive(Debug)]
struct TaskFactory {
    object: PluginObject<ConstructingTaskPlugin>,
}

impl NativePluginFactory for TaskFactory {
    fn package_id(&self) -> &'static str {
        "test.constructor-tasks"
    }

    fn runtime_profile(&self) -> &'static str {
        "lenso.native-authoring@2"
    }

    fn instantiate(
        &self,
        context: NativePluginFactoryContext<'_>,
    ) -> Result<NativePluginInstance, RuntimeFailure> {
        let lifecycle =
            CompleteObjectLifecycle::linked(self.object.clone(), context.configuration())?;
        Ok(NativePluginInstance::with_lifecycle(Vec::new(), lifecycle))
    }
}

#[test]
fn constructor_work_waits_for_ready_and_is_cancelled_with_its_generation() {
    let object = PluginObject::empty();
    let plan = ResolvedAppPlan::new(
        vec![
            PluginInstancePlan::new("tasks", "test.constructor-tasks")
                .with_authoring(2, "lenso.native-authoring@2"),
        ],
        vec![],
    );
    let driver = DeterministicDriver::new();
    let app = driver
        .run(Kernel::start_native(
            plan,
            driver.clone(),
            NativePluginRegistry::new().with_factory(TaskFactory {
                object: object.clone(),
            }),
        ))
        .expect("authoring-2 construction should register generation-owned work");
    let plugin = object.get().expect("constructor should install the object");
    let ready_signal = plugin
        .ready_signal
        .borrow_mut()
        .take()
        .expect("one readiness observation");
    driver
        .run(ready_signal)
        .expect("managed work should observe the same Ready Gate");
    assert_eq!(plugin.probe.ready_runs.get(), 1);
    assert_eq!(plugin.probe.task_drops.get(), 0);
    assert_eq!(plugin.probe.cancellation_runs.get(), 0);

    assert_eq!(
        driver.run(app.shutdown(Duration::from_secs(1))),
        ShutdownOutcome::Clean
    );
    assert_eq!(plugin.probe.stop_runs.get(), 1);
    assert_eq!(plugin.probe.task_drops.get(), 1);
    assert_eq!(plugin.probe.cancellation_runs.get(), 1);
    assert!(matches!(
        plugin.lifecycle.spawn_local(async {}),
        Err(ManagedTasksError::Scope(ManagedTaskError::ScopeClosed))
    ));
    driver.run(driver.yield_now());
    assert_eq!(plugin.probe.ready_runs.get(), 1);
    assert_eq!(plugin.probe.task_drops.get(), 1);
}