saddle-framework 0.3.25

The single business-facing facade for Saddle applications
//! Same production factory/body/cleanup and write guard; no replacement executor.
use super::*;
use saddle_runtime::profusegw::*;
use saddle_runtime::request_task::reserved_set::ReservedCollectionJoin;
use std::{
    alloc::Layout,
    pin::pin,
    task::{Context, Poll, Waker},
};

fn budget() -> saddle_admission::DeploymentResourceBudget {
    let pending = saddle_admission::freeze_deployment_resource_budget(
        1, 32, 5000, 1, 1, 1_000_000, 1_000_000, 1_000_000, 1_000_000,
    )
    .unwrap();
    let (app, listener) = saddle_core::BootstrapRendezvousIssuer::issue()
        .freeze_application(saddle_core::GeneratedApplicationFreezeSource::new(
            "app",
            b"descriptor",
            &["route"],
        ))
        .unwrap();
    let listener = listener
        .freeze_listener(saddle_core::ListenerStartupFreezeSource::new(
            "app",
            "127.0.0.1:8000".parse().unwrap(),
            "127.0.0.1:9000".parse().unwrap(),
            Duration::from_millis(5000),
        ))
        .ok()
        .unwrap();
    let (whole, receipt) = saddle_core::pair_bootstrap_rendezvous(app, listener)
        .ok()
        .unwrap();
    saddle_admission::bind_deployment_resource_budget_bootstrap(pending, whole, receipt)
        .ok()
        .unwrap()
}
fn dispatch(
    _: saddle_boundary::ingress::AcceptedIngress,
    _: (),
    _: BusinessConfig<()>,
    _: crate::database_capability::DatabaseRequest,
) -> std::future::Ready<Result<Vec<u8>>> {
    std::future::ready(Ok(b"{}".to_vec()))
}
type Dispatch = fn(
    saddle_boundary::ingress::AcceptedIngress,
    (),
    BusinessConfig<()>,
    crate::database_capability::DatabaseRequest,
) -> std::future::Ready<Result<Vec<u8>>>;

#[test]
fn formal_factory_matching_join_last_reference_and_partial_write() {
    let mut coordinator = pin!(coordinate_profusegw_app_run(budget(), |process| {
        std::future::ready(run_profusegw_owned_application(
            process,
            |lease| async move {
                let startup = lease.take_database_startup_half().unwrap();
                let observer = saddle_observability::Observer::with_writer(
                    Default::default(),
                    std::io::sink(),
                )
                .unwrap();
                let application = saddle_core::ContextLabel::checked("app").unwrap();
                let dispatch: Dispatch = dispatch;
                let prepared = prepare::<(), (), _, _>("app", &lease, &dispatch).unwrap();
                let Prepared {
                    mut normal,
                    mut rejected,
                    indirect,
                    ..
                } = prepared;
                for mode in 0..3 {
                    let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
                    let mut client = TcpStream::connect(listener.local_addr().unwrap())
                        .await
                        .unwrap();
                    let (socket, _) = listener.accept().await.unwrap();
                    let input = NormalInput {
                        factory: Some(consumer_storage::NormalFactory {
                            socket,
                            adapter: saddle_boundary::ingress::ProfuseGwListenerAdapter::new("app")
                                .unwrap(),
                            deployment: Some(()),
                            business: Some(BusinessConfig::unit()),
                            ingress_token: None,
                            dispatch: Arc::new(dispatch),
                            observer: observer.clone(),
                            admission: None,
                            database: None,
                            diagnostic_handle: None,
                        }),
                        retained: EntryOwner::new(),
                    };
                    let outcome = lease.try_reserved_dispatch(
                        application.clone(),
                        None,
                        body_layout::<(), (), Dispatch, std::future::Ready<Result<Vec<u8>>>>(),
                        &indirect,
                        input,
                        factory::<(), (), _, _>,
                    );
                    let ReservedDispatchOutcome::Ready {
                        root,
                        future,
                        ticket,
                    } = outcome
                    else {
                        panic!("actual factory must fit")
                    };
                    let held = root.view(saddle_core::RequestViewPhase::Reading).unwrap();
                    let abort = normal
                        .spawn(future, ticket)
                        .unwrap_or_else(|_| panic!("original task slot"));
                    if mode == 0 {
                        abort.abort();
                    } else {
                        client.write_all(b"POST / HTTP/1.1\r\n").await.unwrap();
                        for _ in 0..3 {
                            tokio::task::yield_now().await;
                        }
                        if mode == 1 {
                            abort.abort();
                        } else {
                            client.shutdown().await.unwrap();
                        }
                    }
                    // Completed/aborted tasks still own the slot until their matching join.
                    while !abort.is_finished() {
                        tokio::task::yield_now().await;
                    }
                    assert_eq!(normal.len(), 1);
                    assert!(matches!(
                        lease.try_admit(),
                        ProfuseGwCoordinatorAdmissionOutcome::CapacityRejected(_)
                    ));
                    let ReservedCollectionJoin::Matched(joined) = normal.join_next().await.unwrap()
                    else {
                        panic!("matching ticket")
                    };
                    finish_normal(joined, &observer, None).await;
                    drop(root);
                    assert!(matches!(
                        lease.try_admit(),
                        ProfuseGwCoordinatorAdmissionOutcome::CapacityRejected(_)
                    ));
                    if mode == 1 {
                        // Exactly the production write/select/Drop owner, with 4-byte IO capacity.
                        let (mut write, mut read) = tokio::io::duplex(4);
                        let mut retained = None;
                        let mut future = Box::pin(deliver(
                            &mut write,
                            b"abcdefgh",
                            i64::MAX,
                            held.clone(),
                            None,
                            &mut retained,
                        ));
                        assert!(matches!(
                            future
                                .as_mut()
                                .poll(&mut Context::from_waker(Waker::noop())),
                            Poll::Pending
                        ));
                        let mut bytes = [0; 4];
                        read.read_exact(&mut bytes).await.unwrap();
                        assert_eq!(&bytes, b"abcd");
                        drop(future);
                        let ReservedDeliveryOutcome::Failed { axes, failure } =
                            retained.take().unwrap()
                        else {
                            panic!("cancel retains failure")
                        };
                        assert_eq!(axes.bytes_written, Some(4));
                        assert!(matches!(
                            axes.operation,
                            saddle_core::OperationOutcome::Cancelled
                        ));
                        let (_, source) = failure.into_parts();
                        source
                            .finish(
                                &held,
                                None,
                                RootOutcomeFacts {
                                    axes,
                                    ..Default::default()
                                },
                            )
                            .ok()
                            .unwrap();
                    }
                    drop(held);
                    let ProfuseGwCoordinatorAdmissionOutcome::Ready(owner, _) = lease.try_admit()
                    else {
                        panic!("last view must refund")
                    };
                    owner.cancel();
                    println!(
                        "S_REAL_FACTORY mode={mode} no_refund_before_join no_refund_before_last_view refund_after_last_view PASS"
                    );
                }
                println!(
                    "S_REAL_COLLECTION capacity={} ticket_entry={} rejection_body={}",
                    normal.capacity(),
                    Layout::new::<(
                        tokio::task::Id,
                        saddle_runtime::request_task::reserved::ReservedTaskTicket<
                            Owner<
                                (),
                                (),
                                fn(
                                    saddle_boundary::ingress::AcceptedIngress,
                                    (),
                                    BusinessConfig<()>,
                                    crate::database_capability::DatabaseRequest,
                                )
                                    -> std::future::Ready<Result<Vec<u8>>>,
                            >,
                        >
                    )>()
                    .size(),
                    rejected_layout().size()
                );
                // Actual rejection factory and body, up to the unchanged slot
                // bound. Byte capacity may reject before sixteen complete tasks.
                let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
                let mut slots = Vec::new();
                let mut peers = Vec::new();
                for _ in 0..16 {
                    let peer = TcpStream::connect(listener.local_addr().unwrap())
                        .await
                        .unwrap();
                    let (socket, _) = listener.accept().await.unwrap();
                    let input = RejectedInput {
                        socket: Some(socket),
                        code: 503,
                        admission: None,
                        retained: EntryOwner::new(),
                        observer: observer.clone(),
                        output: None,
                    };
                    match lease.try_reserved_rejection(
                        application.clone(),
                        None,
                        rejected_layout(),
                        &[],
                        input,
                        rejected_factory,
                    ) {
                        Ok(pair) => {
                            slots.push(pair);
                            peers.push(peer);
                        }
                        Err((input, _, error)) => {
                            let _facts = record_unrooted_failure(
                                &application,
                                None,
                                UnrootedEntryError::Rejection(&error),
                            );
                            drop(input);
                            break;
                        }
                    }
                }
                assert!(!slots.is_empty());
                let fitted = slots.len();
                // Grow the actual process table while every complete task is
                // retained. The table's old backing remains charged during grow.
                rejected.try_capacity(fitted).unwrap();
                let held = slots
                    .last()
                    .unwrap()
                    .0
                    .view(saddle_core::RequestViewPhase::Reading)
                    .unwrap();
                for (root, future, ticket) in slots {
                    let handle = rejected
                        .spawn(future, ticket)
                        .unwrap_or_else(|_| panic!("reserved rejection collection"));
                    handle.abort();
                    drop(root);
                }
                while let Some(join) = rejected.join_next().await {
                    let ReservedCollectionJoin::Matched(joined) = join else {
                        panic!("matching rejected ticket")
                    };
                    finish_rejected(joined).await;
                }
                assert_eq!(rejected.len(), 0);
                drop((held, peers));
                println!(
                    "S_REAL_REJECTION fitted={fitted} slot_limit=16 retained_capacity={} owner={} body={} matching_abort_cleanup PASS",
                    rejected.capacity(),
                    Layout::new::<RejectedInput>().size(),
                    rejected_layout().size()
                );
                drop((normal, rejected, startup, lease));
                assert!(
                    std::process::Command::new("kill")
                        .args(["-TERM", &std::process::id().to_string()])
                        .status()
                        .unwrap()
                        .success()
                );
                Ok(saddle_runtime::Application::new())
            },
        ))
    }));
    let Poll::Ready(Ok(())) = coordinator
        .as_mut()
        .poll(&mut Context::from_waker(Waker::noop()))
    else {
        panic!("original finalization must reach ZERO")
    };
}