#![cfg(feature = "session")]
use std::collections::BTreeMap;
use std::net::SocketAddr;
use std::sync::Arc;
use std::time::Duration;
use axum::body::{to_bytes, Body};
use axum::extract::ConnectInfo;
use axum::http::{Request, StatusCode};
use base64::Engine;
use boatramp_core::config::{DeployConfig, HandlersSiteConfig, SessionConfig, SiteConfig};
use boatramp_core::deploy::{sha256_hex, DeployStore, FileEntry, Manifest};
use boatramp_core::kv::{KvStore, MemoryKv};
use boatramp_core::project::ProjectRef;
use boatramp_core::ByteStream;
use boatramp_handlers::{HandlerEngine, Limits};
use boatramp_server::{router, Auth, HandlerRuntime};
use boatramp_storage::FsStorage;
use futures::StreamExt;
use tower::ServiceExt;
const SESSION_ECHO: &[u8] =
include_bytes!("../../boatramp-handlers/tests/fixtures/session-echo.wasm");
const SITE: &str = "agentsite";
const ROUTE: &str = "/agent";
type App = axum::Router;
async fn live_app() -> (App, Arc<MemoryKv>) {
let kv = Arc::new(MemoryKv::new());
let dir = std::env::temp_dir().join(format!("br-session-live-{}", std::process::id()));
let storage = Arc::new(FsStorage::new(dir));
let deploy = DeployStore::new(storage.clone(), kv.clone());
let hash = sha256_hex(SESSION_ECHO);
let stream: ByteStream =
futures::stream::once(async move { Ok(bytes::Bytes::from_static(SESSION_ECHO)) }).boxed();
deploy.put_blob(&hash, stream).await.unwrap();
let mut files = BTreeMap::new();
files.insert(
"session/echo.wasm".to_string(),
FileEntry {
hash,
size: SESSION_ECHO.len() as u64,
content_type: None,
variants: BTreeMap::new(),
},
);
let config = DeployConfig {
sessions: vec![SessionConfig {
route: ROUTE.to_string(),
component: "session/echo.wasm".to_string(),
imports: vec!["session".to_string()],
..Default::default()
}],
..Default::default()
};
let manifest = Manifest {
files,
config,
..Default::default()
};
let id = deploy.put_manifest(&manifest).await.unwrap();
deploy
.activate(ProjectRef::DEFAULT, SITE, &id)
.await
.unwrap();
deploy
.set_site_config(
ProjectRef::DEFAULT,
SITE,
&SiteConfig {
handlers: Some(HandlersSiteConfig {
enabled: true,
allow_imports: vec!["session".to_string()],
..Default::default()
}),
..Default::default()
},
)
.await
.unwrap();
let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
let runtime = HandlerRuntime::new(engine, kv.clone(), storage, None, None);
(router(deploy, Auth::disabled(), runtime), kv)
}
fn req(method: &str, uri: &str, body: &'static [u8]) -> Request<Body> {
let mut r = Request::builder()
.method(method)
.uri(uri)
.body(Body::from(body))
.unwrap();
r.extensions_mut()
.insert(ConnectInfo(SocketAddr::from(([127, 0, 0, 1], 40000))));
r
}
async fn post(app: &App, uri: &str, body: &'static [u8]) -> (StatusCode, String) {
let resp = app.clone().oneshot(req("POST", uri, body)).await.unwrap();
let status = resp.status();
let bytes = to_bytes(resp.into_body(), usize::MAX).await.unwrap();
(status, String::from_utf8_lossy(&bytes).into_owned())
}
async fn get_sse(app: &App, uri: &str, last_event_id: Option<&str>) -> String {
let mut builder = Request::builder().method("GET").uri(uri);
if let Some(id) = last_event_id {
builder = builder.header("last-event-id", id);
}
let mut r = builder.body(Body::empty()).unwrap();
r.extensions_mut()
.insert(ConnectInfo(SocketAddr::from(([127, 0, 0, 1], 40001))));
let resp = app.clone().oneshot(r).await.unwrap();
assert_eq!(resp.status(), StatusCode::OK, "SSE open should be 200");
let mut stream = resp.into_body().into_data_stream();
let mut acc: Vec<u8> = Vec::new();
let _ = tokio::time::timeout(Duration::from_millis(1500), async {
while let Some(Ok(chunk)) = stream.next().await {
acc.extend_from_slice(&chunk);
}
})
.await;
String::from_utf8_lossy(&acc).into_owned()
}
fn b64(s: &str) -> String {
base64::engine::general_purpose::STANDARD.encode(s.as_bytes())
}
async fn record(kv: &MemoryKv, id: &str) -> Option<serde_json::Value> {
let key = format!("session/{}/{id}", ProjectRef::DEFAULT.as_str());
kv.get(&key)
.await
.unwrap()
.map(|b| serde_json::from_slice(&b).unwrap())
}
#[tokio::test]
#[ignore = "capability gate: run via capability.yml or with --ignored"]
async fn session_capability_duplex_resume_dedup_and_close() {
let (app, kv) = live_app().await;
let base = format!("/_sites/{SITE}{ROUTE}");
let s1 = format!("{base}?id=S1");
assert_eq!(post(&app, &s1, b"hello").await.0, StatusCode::ACCEPTED);
let rec = record(&kv, "S1").await.expect("session S1 persisted");
let outbound = rec["state"]["outbound"].as_array().unwrap();
assert_eq!(outbound.len(), 1, "one outbound frame buffered");
assert_eq!(outbound[0]["cursor"].as_u64(), Some(1));
let cp = rec["state"]["checkpoint"].as_array().unwrap();
assert_eq!(
cp.first().and_then(serde_json::Value::as_u64),
Some(1),
"count=1 checkpointed"
);
let sse = get_sse(&app, &s1, None).await;
assert!(
sse.contains("event: frame"),
"sse had no frame event: {sse}"
);
assert!(sse.contains("id: 1"), "sse missing cursor id: {sse}");
assert!(
sse.contains(&b64("echo#1:hello")),
"sse missing echoed payload: {sse}"
);
assert_eq!(post(&app, &s1, b"world").await.0, StatusCode::ACCEPTED);
let rec = record(&kv, "S1").await.unwrap();
let outbound = rec["state"]["outbound"].as_array().unwrap();
let cursor2 = outbound
.iter()
.find(|f| f["cursor"].as_u64() == Some(2))
.unwrap();
let payload2: Vec<u8> = cursor2["payload"]
.as_array()
.unwrap()
.iter()
.map(|v| v.as_u64().unwrap() as u8)
.collect();
assert_eq!(
payload2,
b"echo#2:world",
"count must resume from the checkpoint (got {:?})",
String::from_utf8_lossy(&payload2)
);
let sse = get_sse(&app, &s1, Some("1")).await;
assert!(sse.contains("id: 2"), "resume missing cursor 2: {sse}");
assert!(
sse.contains(&b64("echo#2:world")),
"resume missing payload: {sse}"
);
assert!(
!sse.contains(&b64("echo#1:hello")),
"resume must NOT redeliver acked cursor 1: {sse}"
);
let s1_idem = format!("{base}?id=S1&idem=K1");
assert_eq!(post(&app, &s1_idem, b"dupe").await.0, StatusCode::ACCEPTED);
let after_first = record(&kv, "S1").await.unwrap()["state"]["out_cursor"]
.as_u64()
.unwrap();
let (status, body) = post(&app, &s1_idem, b"dupe").await;
assert_eq!(
status,
StatusCode::OK,
"retried POST should be deduped: {body}"
);
let after_dup = record(&kv, "S1").await.unwrap()["state"]["out_cursor"]
.as_u64()
.unwrap();
assert_eq!(
after_first, after_dup,
"a deduped POST must not re-run the guest"
);
assert_eq!(post(&app, &s1, b"cancel").await.0, StatusCode::ACCEPTED);
let sse = get_sse(&app, &s1, None).await;
assert!(
sse.contains("event: close"),
"no close event after cancel: {sse}"
);
assert!(sse.contains("client cancel"), "close reason missing: {sse}");
assert_eq!(
post(&app, &s1, b"again").await.0,
StatusCode::GONE,
"a POST to a closed session must be 410"
);
let bad = format!("{base}?id=has%2Fslash");
assert_eq!(
post(&app, &bad, b"x").await.0,
StatusCode::BAD_REQUEST,
"an invalid session id must be rejected 400"
);
let s2 = format!("{base}?id=S2");
assert_eq!(post(&app, &s2, b"solo").await.0, StatusCode::ACCEPTED);
assert!(record(&kv, "S2").await.is_some());
println!(
"SESSION CAPABILITY GATE OK: duplex send/recv, checkpoint-resume, cursor-resume, \
idempotent redelivery, guest close, and id hardening all held end to end."
);
}