use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use alkcall::channels::operations::{
ChannelCore, ChannelPlan, Establishment, EstablishmentError, OpenEstablisher, OpenHandler,
};
use alkcall::core::auth::AuthContext;
use alkcall::core::ownership::OwnershipProvider;
use alkcall::core::Connection;
use alkcall::registry::spec::{
AccessControl, ChannelOpenSpec, ErrorDefinition, OperationSpec, OperationType, Visibility,
};
use serde_json::{json, Value};
use tracing::debug;
use crate::adapter::{drive_session_pre_negotiated, TTY_OPEN_SCOPE};
use crate::backend::TtyBackend;
use crate::negotiation::{NegotiateRequest, NegotiationWriter};
pub const OP_TTY_OPEN: &str = "channels/tty/sub";
pub const TTY_ALPN: &str = "alk/tty";
pub fn register_openable(
core: &ChannelCore,
backends: Arc<HashMap<String, Arc<dyn TtyBackend>>>,
ownership: Option<Arc<dyn OwnershipProvider>>,
registry: &mut alkcall::registry::registration::OperationRegistry,
auth: AuthContext,
) -> Result<(), String> {
let spec = tty_open_spec();
let establisher = make_tty_establisher(Arc::clone(&backends), ownership);
let open_handler = make_tty_open_handler(backends);
core.register_openable_with_establisher(
spec,
Some(establisher),
open_handler,
registry,
auth,
None,
)
}
pub fn tty_open_spec() -> OperationSpec {
let open_failed_schema = json!({
"type": "object",
"properties": {
"reason": {
"type": "string",
"enum": ["dial_failed", "unknown_resource", "handler_error", "timeout"]
},
"message": { "type": "string" }
},
"required": ["reason", "message"]
});
OperationSpec::new(
OP_TTY_OPEN,
OperationType::Sub,
Visibility::External,
json!({
"type": "object",
"properties": {
"carriage": { "type": "string" },
"backend": { "type": "string" },
"cmd": {
"type": "array",
"items": { "type": "string" }
}
},
"required": ["carriage", "backend", "cmd"]
}),
json!({
"type": "object",
"properties": {
"channel_id": { "type": "integer", "minimum": 0 }
}
}),
vec![ErrorDefinition {
code: "channel:open_failed".to_string(),
description: "Establishment failed after channel allocation: the establisher \
rejected the negotiation (malformed negotiation, unknown backend, \
ownership denial), the backend allocation failed (dial_failed), or \
the establishment deadline expired. Details: { reason, message }. \
No channel_id is returned."
.to_string(),
schema: open_failed_schema,
http_status: None,
}],
AccessControl {
required_scopes: vec![TTY_OPEN_SCOPE.to_string()],
required_scopes_any: None,
resource_type: None,
resource_action: None,
},
None,
)
.with_description(
"Open a terminal session (alk/tty). The open op's input is the \
NegotiateRequest (ADR-009); the data plane is TTY's 5-byte chunk \
format on the allocated channel.",
)
.with_channel_open(ChannelOpenSpec::new(TTY_ALPN))
}
fn make_tty_establisher(
backends: Arc<HashMap<String, Arc<dyn TtyBackend>>>,
ownership: Option<Arc<dyn OwnershipProvider>>,
) -> OpenEstablisher {
Arc::new(move |input: Value, auth: AuthContext| {
let backends = Arc::clone(&backends);
let ownership = ownership.clone();
let identity = auth.identity.clone();
Box::pin(async move {
let req: NegotiateRequest = match serde_json::from_value(input) {
Ok(r) => r,
Err(e) => {
debug!("tty: channels establisher: invalid negotiate params: {e}");
return Err(EstablishmentError::HandlerError {
message: format!("malformed negotiation: {e}"),
});
}
};
if req.carriage != "raw" {
return Err(EstablishmentError::HandlerError {
message: "malformed negotiation: carriage must be 'raw'".to_string(),
});
}
if req.cmd.is_empty() {
return Err(EstablishmentError::HandlerError {
message: "malformed negotiation: cmd must be non-empty".to_string(),
});
}
let backend = match backends.get(&req.backend) {
Some(b) => Arc::clone(b),
None => {
return Err(EstablishmentError::UnknownResource {
message: format!("unknown backend: {}", req.backend),
})
}
};
let params = crate::backend::TtyParams::from(req);
if let Some(provider) = ownership {
if let Some((kind, id)) = backend.resource_id(¶ms) {
let owns = identity
.as_ref()
.map(|id_ref| provider.owns(id_ref, kind, &id, "tty"))
.unwrap_or(false);
if !owns {
debug!("tty: channels establisher: ownership denied");
return Err(EstablishmentError::HandlerError {
message: "forbidden: caller does not own the requested tty \
resource"
.to_string(),
});
}
}
}
let handle = match backend.allocate(¶ms).await {
Ok(h) => h,
Err(e) => {
debug!("tty: channels establisher: allocation failed: {e}");
return Err(EstablishmentError::DialFailed {
message: e.to_string(),
});
}
};
Ok(Establishment::new(Arc::new(AllocatedHandle {
handle: Mutex::new(Some(handle)),
}) as ChannelPlan))
})
})
}
struct AllocatedHandle {
handle: Mutex<Option<crate::backend::TtyHandle>>,
}
impl AllocatedHandle {
fn take(&self) -> Option<crate::backend::TtyHandle> {
self.handle.lock().unwrap_or_else(|e| e.into_inner()).take()
}
}
fn make_tty_open_handler(backends: Arc<HashMap<String, Arc<dyn TtyBackend>>>) -> OpenHandler {
Arc::new(
move |input: Value,
plan: Option<ChannelPlan>,
channel_conn: Connection,
auth: AuthContext| {
let backends = Arc::clone(&backends);
let identity = auth.identity.clone();
tokio::spawn(async move {
let req: NegotiateRequest = match serde_json::from_value(input) {
Ok(r) => r,
Err(e) => {
debug!("tty: channels open: invalid negotiate params: {e}");
let stream = match channel_conn.accept_bi().await {
Ok(s) => s,
Err(e) => {
debug!("tty: channels open: accept_bi failed: {e}");
return;
}
};
let (_, mut client_write) = tokio::io::split(stream);
crate::adapter::send_negotiation_error(
NegotiationWriter::new(&mut client_write),
"malformed_negotiation",
&[("message", &e.to_string())],
)
.await;
return;
}
};
let stream = match channel_conn.accept_bi().await {
Ok(s) => s,
Err(e) => {
debug!("tty: channels open: accept_bi failed: {e}");
return;
}
};
let (client_read, client_write) = tokio::io::split(stream);
let allocated = plan
.as_ref()
.and_then(|p| p.downcast_ref::<AllocatedHandle>())
.and_then(|slot| slot.take());
match allocated {
Some(handle) => {
crate::adapter::drive_session_pre_allocated(
client_write,
client_read,
handle,
)
.await;
}
None => {
drive_session_pre_negotiated(
client_write,
client_read,
req,
backends,
None,
identity,
)
.await;
}
}
})
},
)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::backend::{MockBackend, TtyError};
use crate::testing::wire_client_and_server;
use alkcall::channels::operations::ChannelCore;
use alkcall::core::auth::Identity;
use alkcall::registry::registration::OperationRegistry;
use std::collections::HashMap as StdHashMap;
use tokio::io::duplex;
#[tokio::test]
async fn open_handler_writes_error_frame_on_schema_bypassing_input() {
use alkcall::core::types::{BiStream, BidiStreamSource, StreamError};
use tokio::io::AsyncReadExt;
struct YieldOnce {
stream: tokio::sync::Mutex<Option<BiStream>>,
}
#[async_trait::async_trait]
impl BidiStreamSource for YieldOnce {
async fn accept_bi(&self) -> Result<BiStream, StreamError> {
self.stream
.lock()
.await
.take()
.ok_or(StreamError::ConnectionClosed)
}
async fn open_bi(&self) -> Result<BiStream, StreamError> {
Err(StreamError::StreamClosed)
}
fn remote_addr(&self) -> Option<std::net::SocketAddr> {
None
}
fn close(&self, _code: u32, _reason: &str) {
}
}
let (client_end, server_end) = duplex(64 * 1024);
let (server_read, server_write) = tokio::io::split(server_end);
let bidi = BiStream::from_joined(server_read, server_write);
let channel_conn = Connection::from_source(
YieldOnce {
stream: tokio::sync::Mutex::new(Some(bidi)),
},
TTY_ALPN.as_bytes().to_vec(),
);
let mut backends: HashMap<String, Arc<dyn TtyBackend>> = HashMap::new();
backends.insert("mock".to_string(), Arc::new(MockBackend::with_exit_code(0)));
let handler = make_tty_open_handler(Arc::new(backends));
let task = handler(
serde_json::json!({
"carriage": "raw",
"backend": "mock",
"cmd": ["true"],
"cwd": 42
}),
None,
channel_conn,
AuthContext::anonymous(b"test"),
);
task.await.expect("handler task");
let mut read = client_end;
let mut first = [0u8; 1];
tokio::time::timeout(
std::time::Duration::from_secs(5),
read.read_exact(&mut first),
)
.await
.expect("no byte from handler")
.expect("read first byte");
assert_eq!(first[0], 0x00, "error frame length prefix starts with 0x00");
let mut len_rest = [0u8; 3];
read.read_exact(&mut len_rest).await.expect("read len rest");
let len = u32::from_be_bytes([first[0], len_rest[0], len_rest[1], len_rest[2]]) as usize;
let mut body = vec![0u8; len];
read.read_exact(&mut body).await.expect("read error body");
let v: serde_json::Value = serde_json::from_slice(&body).expect("parse error frame");
assert_eq!(v["error"], "malformed_negotiation");
assert!(v["message"].as_str().is_some_and(|m| !m.is_empty()));
}
#[test]
fn op_tty_open_is_channels_tty_sub() {
assert_eq!(OP_TTY_OPEN, "channels/tty/sub");
}
fn mock_backends_one() -> Arc<HashMap<String, Arc<dyn TtyBackend>>> {
let mut backends: HashMap<String, Arc<dyn TtyBackend>> = HashMap::new();
backends.insert("mock".to_string(), Arc::new(MockBackend::with_exit_code(0)));
Arc::new(backends)
}
async fn establisher_result(
backends: Arc<HashMap<String, Arc<dyn TtyBackend>>>,
ownership: Option<Arc<dyn OwnershipProvider>>,
input: Value,
) -> Result<Establishment, EstablishmentError> {
let establisher = make_tty_establisher(backends, ownership);
establisher(input, AuthContext::anonymous(b"test")).await
}
#[tokio::test]
async fn establisher_ok_on_valid_params() {
let result = establisher_result(
mock_backends_one(),
None,
json!({ "carriage": "raw", "backend": "mock", "cmd": ["true"] }),
)
.await;
assert!(result.is_ok(), "valid params establish, got {result:?}");
}
#[tokio::test]
async fn establisher_rejects_schema_bypassing_parse() {
let result = establisher_result(
mock_backends_one(),
None,
json!({ "carriage": "raw", "backend": "mock", "cmd": ["true"], "cwd": 42 }),
)
.await;
match result {
Err(e) => {
assert_eq!(e.reason(), "handler_error");
assert!(e.message().contains("malformed negotiation"));
}
Ok(_) => panic!("schema-bypassing input must be rejected by the establisher"),
}
}
#[tokio::test]
async fn establisher_rejects_non_raw_carriage() {
let result = establisher_result(
mock_backends_one(),
None,
json!({ "carriage": "line", "backend": "mock", "cmd": ["true"] }),
)
.await;
match result {
Err(e) => {
assert_eq!(e.reason(), "handler_error");
assert!(e.message().contains("carriage must be 'raw'"));
}
Ok(_) => panic!("non-raw carriage must be rejected"),
}
}
#[tokio::test]
async fn establisher_rejects_empty_cmd() {
let result = establisher_result(
mock_backends_one(),
None,
json!({ "carriage": "raw", "backend": "mock", "cmd": [] }),
)
.await;
match result {
Err(e) => {
assert_eq!(e.reason(), "handler_error");
assert!(e.message().contains("cmd must be non-empty"));
}
Ok(_) => panic!("empty cmd must be rejected"),
}
}
#[tokio::test]
async fn establisher_rejects_unknown_backend_as_unknown_resource() {
let result = establisher_result(
mock_backends_one(),
None,
json!({ "carriage": "raw", "backend": "nope", "cmd": ["true"] }),
)
.await;
match result {
Err(e) => {
assert_eq!(e.reason(), "unknown_resource");
assert_eq!(e.message(), "unknown backend: nope");
}
Ok(_) => panic!("unknown backend must be rejected"),
}
}
#[tokio::test]
async fn establisher_rejects_ownership_denial() {
use alkcall::core::OwnershipStore;
let store = alkcall::core::ownership::InMemoryOwnershipStore::new();
store
.record(
&Identity {
id: "mallory".to_string(),
scopes: vec![],
resources: HashMap::new(),
},
"session",
"s-42",
)
.await
.expect("record");
let provider: Arc<dyn OwnershipProvider> = Arc::new(store);
struct OwnedBackend;
#[async_trait::async_trait]
impl TtyBackend for OwnedBackend {
async fn allocate(
&self,
_params: &crate::backend::TtyParams,
) -> Result<crate::backend::TtyHandle, crate::backend::TtyError> {
panic!("allocate must not run during establishment");
}
fn resource_id(
&self,
_params: &crate::backend::TtyParams,
) -> Option<(&'static str, String)> {
Some(("session", "s-42".to_string()))
}
}
let mut backends: HashMap<String, Arc<dyn TtyBackend>> = HashMap::new();
backends.insert("mock".to_string(), Arc::new(OwnedBackend));
let result = establisher_result(
Arc::new(backends),
Some(provider),
json!({ "carriage": "raw", "backend": "mock", "cmd": ["true"] }),
)
.await;
match result {
Err(e) => {
assert_eq!(e.reason(), "handler_error");
assert!(e.message().contains("forbidden"));
}
Ok(_) => panic!("ownership denial must be rejected by the establisher"),
}
}
#[tokio::test]
async fn establisher_allocates_and_maps_failure_to_dial_failed() {
struct AllocFailBackend;
#[async_trait::async_trait]
impl TtyBackend for AllocFailBackend {
async fn allocate(
&self,
_params: &crate::backend::TtyParams,
) -> Result<crate::backend::TtyHandle, crate::backend::TtyError> {
Err(crate::backend::TtyError::AllocFailed {
message: "out of ptys".to_string(),
})
}
}
let mut failing: HashMap<String, Arc<dyn TtyBackend>> = HashMap::new();
failing.insert("mock".to_string(), Arc::new(AllocFailBackend));
let result = establisher_result(
Arc::new(failing),
None,
json!({ "carriage": "raw", "backend": "mock", "cmd": ["true"] }),
)
.await;
match result {
Err(e) => {
assert_eq!(e.reason(), "dial_failed");
assert!(e.message().contains("out of ptys"));
}
Ok(_) => panic!("allocation failure must be rejected by the establisher"),
}
let result = establisher_result(
mock_backends_one(),
None,
json!({ "carriage": "raw", "backend": "mock", "cmd": ["true"] }),
)
.await;
let establishment = result.expect("successful allocate carries the handle");
let slot = establishment
.plan
.as_ref()
.and_then(|p| p.downcast_ref::<AllocatedHandle>())
.expect("plan is the AllocatedHandle slot")
.take()
.expect("the allocated handle is in the slot");
let code = slot.exit_code.await.expect("exit_code resolves");
assert_eq!(code, 0);
}
#[test]
fn tty_open_spec_declares_description_and_open_failed_error() {
let spec = tty_open_spec();
assert!(
spec.description.is_some(),
"description set for services/list disclosure (review 006 E-02)"
);
let def = spec
.error_schemas
.iter()
.find(|e| e.code == "channel:open_failed")
.expect("channel:open_failed ErrorDefinition declared");
let enum_values = def.schema["properties"]["reason"]["enum"]
.as_array()
.expect("reason enum");
let reasons: Vec<&str> = enum_values.iter().filter_map(|v| v.as_str()).collect();
assert_eq!(
reasons,
vec![
"dial_failed",
"unknown_resource",
"handler_error",
"timeout"
]
);
}
#[test]
fn tty_alpn_is_alk_tty() {
assert_eq!(TTY_ALPN, "alk/tty");
}
#[test]
fn tty_open_spec_has_channel_open_marker_for_alk_tty() {
let spec = tty_open_spec();
assert_eq!(spec.name, OP_TTY_OPEN);
assert_eq!(spec.op_type, OperationType::Sub);
assert_eq!(spec.visibility, Visibility::External);
let marker = spec.channel_open.expect("channel_open marker set");
assert_eq!(marker.alpn, TTY_ALPN);
}
#[test]
fn tty_open_spec_requires_tty_open_scope() {
let spec = tty_open_spec();
assert_eq!(spec.access_control.required_scopes, vec![TTY_OPEN_SCOPE]);
assert!(spec.access_control.has_restrictions());
}
#[test]
fn tty_open_spec_input_schema_has_required_fields() {
let spec = tty_open_spec();
let schema = spec.input_schema;
let required = schema
.get("required")
.and_then(|v| v.as_array())
.expect("required array");
let required_names: Vec<&str> = required.iter().filter_map(|v| v.as_str()).collect();
assert!(required_names.contains(&"carriage"));
assert!(required_names.contains(&"backend"));
assert!(required_names.contains(&"cmd"));
}
#[tokio::test]
async fn register_openable_registers_op_on_registry() {
use alkcall::channels::policy::default_policy;
let (_client, server) = duplex(1024);
let (_reader, writer) = tokio::io::split(server);
let (handle, _runner) = alkcall::channels::mux::MuxRunner::new(Box::new(writer));
let manager = alkcall::channels::manager::ChannelManager::with_defaults(handle, None);
let core = ChannelCore::new(manager, default_policy());
let backends: Arc<HashMap<String, Arc<dyn TtyBackend>>> = Arc::new(HashMap::new());
let mut registry = OperationRegistry::new();
register_openable(
&core,
backends,
None,
&mut registry,
AuthContext::anonymous(b"alk/channels"),
)
.expect("register_openable");
assert!(registry.registration(OP_TTY_OPEN).is_some());
}
#[tokio::test]
async fn end_to_end_open_op_returns_channel_id() {
let mut backends: HashMap<String, Arc<dyn TtyBackend>> = HashMap::new();
backends.insert("mock".to_string(), Arc::new(MockBackend::with_exit_code(0)));
let backends = Arc::new(backends);
let identity = Identity {
id: "alice".to_string(),
scopes: vec![TTY_OPEN_SCOPE.to_string()],
resources: StdHashMap::new(),
};
let client = wire_client_and_server(backends, None, Some(identity)).await;
let response = tokio::time::timeout(
std::time::Duration::from_secs(10),
client.call_open_op(
OP_TTY_OPEN,
json!({
"carriage": "raw",
"backend": "mock",
"cmd": ["true"],
}),
),
)
.await
.expect("open op timed out");
let out = response.result.expect("open op should succeed");
let channel_id = out
.get("channel_id")
.and_then(|v| v.as_u64())
.expect("channel_id in response");
assert!(
channel_id > 0,
"channel_id should be non-zero (channel 0 is the call channel), got {channel_id}"
);
}
#[tokio::test]
async fn end_to_end_open_op_denies_without_tty_open_scope() {
let mut backends: HashMap<String, Arc<dyn TtyBackend>> = HashMap::new();
backends.insert("mock".to_string(), Arc::new(MockBackend::with_exit_code(0)));
let backends = Arc::new(backends);
let identity = Identity {
id: "alice".to_string(),
scopes: vec![],
resources: StdHashMap::new(),
};
let client = wire_client_and_server(backends, None, Some(identity)).await;
let response = tokio::time::timeout(
std::time::Duration::from_secs(10),
client.call_open_op(
OP_TTY_OPEN,
json!({
"carriage": "raw",
"backend": "mock",
"cmd": ["true"],
}),
),
)
.await
.expect("open op timed out");
let err = response.result.expect_err("open op should be denied");
assert!(
err.message.contains("scope")
|| err.code.contains("FORBIDDEN")
|| err.code.contains("AUTH"),
"error should mention scope/forbidden/auth, got code={} message={}",
err.code,
err.message
);
}
#[tokio::test]
async fn mock_backend_resolves_configured_exit_code() {
let backend = MockBackend::with_exit_code(7);
let params = crate::backend::TtyParams {
terminal: None,
cmd: vec!["true".to_string()],
cwd: None,
env: HashMap::new(),
backend_params: serde_json::Map::new(),
};
let handle = backend.allocate(¶ms).await.expect("allocate");
let code = handle.exit_code.await.expect("exit_code resolves");
assert_eq!(code, 7);
let _: Result<i32, TtyError> = Ok(code);
}
}