use std::collections::BTreeMap;
use wip_client::{
CallOutcome, Client, ClientErrorKind, ClientLimits, Completion, DispatchedTransportFailure,
ObservationState, OutcomeUnknownReason, SecurityContext,
};
use wip_http::http::Response;
use wip_http::{
EncodeHttpError, Limits, encode_call_operation_response, encode_fetch_interface_response,
encode_observe_response, encode_protocol_error_response,
};
use wip_protocol::{
CallOperationResponse, FetchInterfaceRequest, FetchInterfaceResponse, INTERFACE_FORMAT_V1,
InterfaceDescriptor, Object, ObjectObservation, ObserveRequest, OperationDeclaration,
ParameterDeclaration, ProtocolError, ProtocolErrorCode, ProtocolInteraction, ReturnDeclaration,
TypeExpr, Value,
};
fn wire_limits() -> Limits {
Limits::new(64 * 1024, 64 * 1024, 64).unwrap()
}
#[derive(Clone)]
struct TestObservation {
object: Object,
children: Vec<TestObservation>,
}
struct TestObservationResponse {
root: TestObservation,
}
fn protocol_observation(node: &TestObservation, depth: u32) -> ObjectObservation {
ObjectObservation {
object: node.object.clone(),
children: (depth > 0).then(|| {
node.children
.iter()
.map(|child| protocol_observation(child, depth - 1))
.collect()
}),
}
}
fn encode_test_observation_response(
request: &ObserveRequest,
value: &TestObservationResponse,
limits: Limits,
) -> Result<Response<Vec<u8>>, EncodeHttpError> {
encode_observe_response(
request,
&protocol_observation(&value.root, request.depth),
limits,
)
}
fn client(objects: usize, interfaces: usize, in_flight: usize, history: usize) -> Client {
Client::new(
ClientLimits::new(8, objects, interfaces, in_flight, history).unwrap(),
wire_limits(),
)
}
fn object(name: &str, reference: Option<&str>, validator: &[u8]) -> Object {
Object {
name: name.into(),
description: Some(format!("object {name}")),
interfaces: vec!["math".into()],
r#ref: reference.map(str::to_owned),
validator: Some(validator.to_vec()),
}
}
fn descriptor(return_type: TypeExpr) -> InterfaceDescriptor {
InterfaceDescriptor {
format: INTERFACE_FORMAT_V1.into(),
documentation: None,
types: vec![],
operations: vec![OperationDeclaration {
name: "calculate".into(),
documentation: None,
parameters: vec![ParameterDeclaration {
name: "value".into(),
required: true,
documentation: None,
r#type: TypeExpr::Integer,
}],
returns: ReturnDeclaration {
documentation: None,
r#type: return_type,
},
}],
}
}
fn observation_response(path: &str, value: &Object) -> Response<Vec<u8>> {
encode_observe_response(
&ObserveRequest {
path: path.into(),
depth: 0,
},
&ObjectObservation {
object: value.clone(),
children: None,
},
wire_limits(),
)
.unwrap()
}
fn observe_object(client: &mut Client, session: &wip_client::SessionId, path: &str, value: Object) {
let request = client.refresh_observed(session, path, 0).unwrap();
client
.complete(request.id, observation_response(path, &value))
.unwrap();
}
fn observe_interface(
client: &mut Client,
session: &wip_client::SessionId,
value: FetchInterfaceResponse,
) {
let request = client
.prepare_interface(session, value.interface.clone())
.unwrap();
let response = encode_fetch_interface_response(
&FetchInterfaceRequest {
interface: value.interface.clone(),
},
&value,
wire_limits(),
)
.unwrap();
client.complete(request.id, response).unwrap();
}
fn args(value: i64) -> BTreeMap<String, Value> {
BTreeMap::from([("value".into(), Value::Integer(value))])
}
#[test]
fn sessions_partition_known_space_by_canonical_endpoint_and_security_context() {
let mut client = client(8, 8, 8, 8);
let alice = client
.open_session(
"https://example.test/wip",
SecurityContext::new("alice-token"),
)
.unwrap();
let alice_canonical = client
.open_session(
"https://example.test/wip/",
SecurityContext::new("alice-token"),
)
.unwrap();
let bob = client
.open_session(
"https://example.test/wip",
SecurityContext::new("bob-token"),
)
.unwrap();
let other = client
.open_session(
"https://other.test/wip",
SecurityContext::new("alice-token"),
)
.unwrap();
assert_eq!(alice, alice_canonical);
observe_object(&mut client, &alice, "/item", object("item", None, b"a"));
assert!(client.object(&alice, "/item").is_some());
assert!(client.object(&bob, "/item").is_none());
assert!(client.object(&other, "/item").is_none());
}
#[test]
fn ensure_observed_deduplicates_loading_and_fresh_coverage() {
let mut client = client(8, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
let request = ObserveRequest {
path: "/".into(),
depth: 1,
};
let value = TestObservationResponse {
root: TestObservation {
object: object("", None, b"root"),
children: vec![TestObservation {
object: object("child", None, b"child"),
children: vec![],
}],
},
};
let root = client
.ensure_observed(&session, "/", 1)
.unwrap()
.expect("missing coverage starts an observation");
assert!(client.ensure_observed(&session, "/", 1).unwrap().is_none());
client
.complete(
root.id,
encode_test_observation_response(&request, &value, wire_limits()).unwrap(),
)
.unwrap();
assert!(client.ensure_observed(&session, "/", 1).unwrap().is_none());
let child = client
.ensure_observed(&session, "/child", 1)
.unwrap()
.expect("depth-boundary children remain missing coverage");
assert!(
client
.ensure_observed(&session, "/child", 1)
.unwrap()
.is_none()
);
client
.complete(
child.id,
encode_test_observation_response(
&ObserveRequest {
path: "/child".into(),
depth: 1,
},
&TestObservationResponse {
root: TestObservation {
object: object("child", None, b"child"),
children: vec![],
},
},
wire_limits(),
)
.unwrap(),
)
.unwrap();
assert!(
client
.ensure_observed(&session, "/child", 1)
.unwrap()
.is_none()
);
}
#[test]
fn ensure_observed_respects_capacity_and_does_not_loop_on_failure() {
let mut client = client(8, 8, 1, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
let first = client
.ensure_observed(&session, "/one", 0)
.unwrap()
.expect("first observation fits");
assert_eq!(
client
.ensure_observed(&session, "relative", 0)
.unwrap_err()
.kind,
ClientErrorKind::InvalidRequest
);
assert!(
client
.ensure_observed(&session, "/two", 0)
.unwrap()
.is_none()
);
assert_eq!(
client
.fail_transport(first.id, "scripted failure")
.unwrap_err()
.kind,
ClientErrorKind::Transport
);
assert!(
client
.ensure_observed(&session, "/one", 0)
.unwrap()
.is_none(),
"retained failures require an explicit refresh"
);
assert!(
client
.ensure_observed(&session, "/two", 0)
.unwrap()
.is_some()
);
}
#[test]
fn ensure_interface_deduplicates_loading_fresh_and_retained_failure() {
let mut client = client(8, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
let request = client
.ensure_interface(&session, "math")
.unwrap()
.expect("missing descriptor starts a fetch");
assert!(client.ensure_interface(&session, "math").unwrap().is_none());
let value = FetchInterfaceResponse {
interface: "math".into(),
descriptor: descriptor(TypeExpr::Integer),
validator: Some(b"math-v1".to_vec()),
};
client
.complete(
request.id,
encode_fetch_interface_response(
&FetchInterfaceRequest {
interface: "math".into(),
},
&value,
wire_limits(),
)
.unwrap(),
)
.unwrap();
assert!(client.ensure_interface(&session, "math").unwrap().is_none());
let failed = client
.ensure_interface(&session, "broken")
.unwrap()
.expect("new descriptor starts a fetch");
assert_eq!(
client
.fail_transport(failed.id, "scripted failure")
.unwrap_err()
.kind,
ClientErrorKind::Transport
);
assert!(
client
.ensure_interface(&session, "broken")
.unwrap()
.is_none(),
"retained interface failures require an explicit refresh"
);
}
#[test]
fn ensure_interface_waits_for_request_and_cache_capacity_without_evicting() {
let mut client = client(8, 1, 1, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
let pending_object = client.refresh_object(&session, "/busy").unwrap();
assert!(
client.ensure_interface(&session, "math").unwrap().is_none(),
"lookahead waits while the in-flight bound is full"
);
client
.complete(
pending_object.id,
observation_response("/busy", &object("busy", None, b"busy")),
)
.unwrap();
let request = client
.ensure_interface(&session, "math")
.unwrap()
.expect("lookahead resumes after capacity becomes available");
let value = FetchInterfaceResponse {
interface: "math".into(),
descriptor: descriptor(TypeExpr::Integer),
validator: None,
};
client
.complete(
request.id,
encode_fetch_interface_response(
&FetchInterfaceRequest {
interface: "math".into(),
},
&value,
wire_limits(),
)
.unwrap(),
)
.unwrap();
assert!(
client
.ensure_interface(&session, "other")
.unwrap()
.is_none(),
"lookahead does not churn a full descriptor cache"
);
assert!(client.interface(&session, "math").is_some());
assert!(client.interface(&session, "other").is_none());
let explicit = client.prepare_interface(&session, "other").unwrap();
assert!(client.interface(&session, "math").is_none());
assert!(matches!(
client.interface(&session, "other").unwrap().state,
ObservationState::Loading(id) if id == explicit.id
));
}
#[test]
fn independent_retrievals_run_concurrently_and_complete_out_of_order() {
let mut client = client(8, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
let first = client.refresh_observed(&session, "/one", 0).unwrap();
let second = client.refresh_observed(&session, "/two", 0).unwrap();
let interface = client.prepare_interface(&session, "math").unwrap();
assert_eq!(first.request.uri().path(), "/wip/v1/observe");
assert_eq!(second.request.uri().path(), "/wip/v1/observe");
assert_eq!(interface.request.uri().path(), "/wip/v1/fetch_interface");
assert_eq!(client.in_flight_len(), 3);
client
.complete(
second.id,
observation_response("/two", &object("two", None, b"two")),
)
.unwrap();
client
.complete(
interface.id,
encode_fetch_interface_response(
&FetchInterfaceRequest {
interface: "math".into(),
},
&FetchInterfaceResponse {
interface: "math".into(),
descriptor: descriptor(TypeExpr::Integer),
validator: Some(b"math-v1".to_vec()),
},
wire_limits(),
)
.unwrap(),
)
.unwrap();
client
.complete(
first.id,
observation_response("/one", &object("one", None, b"one")),
)
.unwrap();
assert_eq!(
client.object(&session, "/one").unwrap().state,
ObservationState::Fresh
);
assert_eq!(
client.object(&session, "/two").unwrap().state,
ObservationState::Fresh
);
assert_eq!(
client.interface(&session, "math").unwrap().state,
ObservationState::Fresh
);
assert_eq!(client.in_flight_len(), 0);
}
#[test]
fn newer_refresh_rejects_out_of_order_response_for_the_same_subject() {
let mut client = client(8, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
let old = client.refresh_observed(&session, "/item", 0).unwrap();
let new = client.refresh_observed(&session, "/item", 0).unwrap();
assert_eq!(
client
.complete(
old.id,
observation_response("/item", &object("item", None, b"old"))
)
.unwrap(),
Completion::StaleResponseRejected(old.id)
);
assert!(matches!(
client.object(&session, "/item").unwrap().state,
ObservationState::Loading(id) if id == new.id
));
client
.complete(
new.id,
observation_response("/item", &object("item", None, b"new")),
)
.unwrap();
assert_eq!(
client.object(&session, "/item").unwrap().validator,
Some(b"new".to_vec())
);
}
#[test]
fn normalized_object_cache_rejects_conflicts_and_evicts_to_its_bound() {
let mut client = client(2, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
observe_object(
&mut client,
&session,
"/a/shared",
object("shared", Some("shared-ref"), b"v1"),
);
observe_object(
&mut client,
&session,
"/b/shared",
object("shared", Some("shared-ref"), b"v1"),
);
let conflict = client.refresh_observed(&session, "/b/shared", 0).unwrap();
let mut changed = object("shared", Some("shared-ref"), b"v1");
changed.description = Some("contradiction".into());
let error = client
.complete(conflict.id, observation_response("/b/shared", &changed))
.unwrap_err();
assert_eq!(error.kind, ClientErrorKind::InvalidResponse);
observe_object(
&mut client,
&session,
"/second",
object("second", None, b"2"),
);
observe_object(&mut client, &session, "/third", object("third", None, b"3"));
assert!(client.object(&session, "/a/shared").is_none());
assert!(client.object(&session, "/second").is_some());
assert!(client.object(&session, "/third").is_some());
}
#[test]
fn normalized_aliases_preserve_newer_request_ownership() {
let mut client = client(8, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
observe_object(
&mut client,
&session,
"/a/shared",
object("shared", Some("shared-ref"), b"initial"),
);
observe_object(
&mut client,
&session,
"/b/shared",
object("shared", Some("shared-ref"), b"initial"),
);
let request = ObserveRequest {
path: "/a".into(),
depth: 1,
};
let older_tree_response = || TestObservationResponse {
root: TestObservation {
object: object("a", None, b"root"),
children: vec![TestObservation {
object: object("shared", Some("shared-ref"), b"older"),
children: vec![],
}],
},
};
let older_tree = client.refresh_observed(&session, "/a", 1).unwrap();
let newer_alias = client.refresh_observed(&session, "/b/shared", 0).unwrap();
client
.complete(
newer_alias.id,
observation_response("/b/shared", &object("shared", Some("shared-ref"), b"newer")),
)
.unwrap();
client
.complete(
older_tree.id,
encode_test_observation_response(&request, &older_tree_response(), wire_limits())
.unwrap(),
)
.unwrap();
assert_eq!(
client.object(&session, "/a/shared").unwrap().validator,
Some(b"newer".to_vec())
);
assert_eq!(
client.object(&session, "/b/shared").unwrap().validator,
Some(b"newer".to_vec())
);
let older_tree = client.refresh_observed(&session, "/a", 1).unwrap();
let newer_alias = client.refresh_observed(&session, "/b/shared", 0).unwrap();
client
.complete(
older_tree.id,
encode_test_observation_response(&request, &older_tree_response(), wire_limits())
.unwrap(),
)
.unwrap();
assert!(matches!(
client.object(&session, "/b/shared").unwrap().state,
ObservationState::Loading(id) if id == newer_alias.id
));
assert_eq!(
client.object(&session, "/b/shared").unwrap().validator,
Some(b"newer".to_vec())
);
client
.complete(
newer_alias.id,
observation_response(
"/b/shared",
&object("shared", Some("shared-ref"), b"newest"),
),
)
.unwrap();
assert_eq!(
client.object(&session, "/a/shared").unwrap().validator,
Some(b"newest".to_vec())
);
let older_tree = client.refresh_observed(&session, "/a", 1).unwrap();
let newer_alias = client.refresh_observed(&session, "/b/shared", 0).unwrap();
client
.complete(
newer_alias.id,
observation_response(
"/b/shared",
&object("shared", Some("shared-ref"), b"newest"),
),
)
.unwrap();
let mut contradictory = object("shared", Some("shared-ref"), b"newest");
contradictory.description = Some("contradictory representation".into());
let error = client
.complete(
older_tree.id,
encode_test_observation_response(
&request,
&TestObservationResponse {
root: TestObservation {
object: object("a", None, b"root"),
children: vec![TestObservation {
object: contradictory,
children: vec![],
}],
},
},
wire_limits(),
)
.unwrap(),
)
.unwrap_err();
assert_eq!(error.kind, ClientErrorKind::InvalidResponse);
assert_eq!(
client.object(&session, "/b/shared").unwrap().validator,
Some(b"newest".to_vec())
);
observe_object(
&mut client,
&session,
"/child",
object("child", Some("child-ref"), b"stable"),
);
let older_tree = client.refresh_observed(&session, "/", 1).unwrap();
let newer_child = client.refresh_observed(&session, "/child", 0).unwrap();
client
.complete(
newer_child.id,
observation_response("/child", &object("child", Some("child-ref"), b"stable")),
)
.unwrap();
let mut contradictory = object("child", Some("child-ref"), b"stable");
contradictory.description = Some("same-path contradiction".into());
let root_request = ObserveRequest {
path: "/".into(),
depth: 1,
};
let error = client
.complete(
older_tree.id,
encode_test_observation_response(
&root_request,
&TestObservationResponse {
root: TestObservation {
object: object("", None, b"root"),
children: vec![TestObservation {
object: contradictory,
children: vec![],
}],
},
},
wire_limits(),
)
.unwrap(),
)
.unwrap_err();
assert_eq!(error.kind, ClientErrorKind::InvalidResponse);
assert_eq!(
client.object(&session, "/child").unwrap().state,
ObservationState::Fresh
);
}
#[test]
fn observation_expansion_is_complete_and_refreshes_edges() {
let mut client = client(16, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
assert!(client.tree(&session, "/").is_none());
let request = client.refresh_observed(&session, "/", 1).unwrap();
let logical_request = ObserveRequest {
path: "/".into(),
depth: 1,
};
let tree = TestObservationResponse {
root: TestObservation {
object: object("", None, b"root"),
children: vec![TestObservation {
object: object("items", None, b"items"),
children: vec![],
}],
},
};
client
.complete(
request.id,
encode_test_observation_response(&logical_request, &tree, wire_limits()).unwrap(),
)
.unwrap();
assert_eq!(client.tree(&session, "/").unwrap().children, ["/items"]);
assert!(client.object(&session, "/items").is_some());
let refresh = client.refresh_observed(&session, "/", 1).unwrap();
let refreshed = TestObservationResponse {
root: TestObservation {
object: object("", None, b"root-2"),
children: vec![],
},
};
client
.complete(
refresh.id,
encode_test_observation_response(&logical_request, &refreshed, wire_limits()).unwrap(),
)
.unwrap();
assert!(client.tree(&session, "/").unwrap().children.is_empty());
assert_eq!(
client.object(&session, "/items").unwrap().state,
ObservationState::Stale
);
}
#[test]
fn depth_zero_observation_preserves_previously_observed_children() {
let mut client = client(16, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
let deep_request = ObserveRequest {
path: "/".into(),
depth: 1,
};
let deep = client.refresh_observed(&session, "/", 1).unwrap();
client
.complete(
deep.id,
encode_test_observation_response(
&deep_request,
&TestObservationResponse {
root: TestObservation {
object: object("", None, b"root-v1"),
children: vec![TestObservation {
object: object("child", None, b"child-v1"),
children: vec![],
}],
},
},
wire_limits(),
)
.unwrap(),
)
.unwrap();
let shallow = client.refresh_observed(&session, "/", 0).unwrap();
client
.complete(
shallow.id,
observation_response("/", &object("", None, b"root-v2")),
)
.unwrap();
assert_eq!(client.tree(&session, "/").unwrap().children, ["/child"]);
assert_eq!(
client.object(&session, "/child").unwrap().state,
ObservationState::Fresh
);
assert_eq!(
client.object(&session, "/").unwrap().validator,
Some(b"root-v2".to_vec())
);
}
#[test]
fn shallow_observation_supersedes_loading_deep_observation_without_leaking_loading_state() {
let mut client = client(16, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
let deep_request = ObserveRequest {
path: "/".into(),
depth: 1,
};
let deep_value = TestObservationResponse {
root: TestObservation {
object: object("", None, b"old-root"),
children: vec![],
},
};
let deep = client.refresh_observed(&session, "/", 1).unwrap();
let shallow = client.refresh_observed(&session, "/", 0).unwrap();
assert_eq!(
client.tree(&session, "/").unwrap().state,
ObservationState::Stale
);
assert_eq!(
client
.complete(
deep.id,
encode_test_observation_response(&deep_request, &deep_value, wire_limits())
.unwrap(),
)
.unwrap(),
Completion::StaleResponseRejected(deep.id)
);
client
.complete(
shallow.id,
observation_response("/", &object("", None, b"new-root")),
)
.unwrap();
assert_eq!(
client.tree(&session, "/").unwrap().state,
ObservationState::Stale
);
}
#[test]
fn newer_parent_boundary_observation_replans_superseded_child_lookahead() {
let mut client = client(16, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
let root_request = ObserveRequest {
path: "/".into(),
depth: 1,
};
let root_value = |validator: &'static [u8]| TestObservationResponse {
root: TestObservation {
object: object("", None, validator),
children: vec![TestObservation {
object: object("child", None, validator),
children: vec![],
}],
},
};
let initial = client.refresh_observed(&session, "/", 1).unwrap();
client
.complete(
initial.id,
encode_test_observation_response(&root_request, &root_value(b"initial"), wire_limits())
.unwrap(),
)
.unwrap();
let older_child = client.refresh_observed(&session, "/child", 1).unwrap();
let newer_parent = client.refresh_observed(&session, "/", 1).unwrap();
client
.complete(
newer_parent.id,
encode_test_observation_response(&root_request, &root_value(b"newer"), wire_limits())
.unwrap(),
)
.unwrap();
assert_eq!(
client.tree(&session, "/child").unwrap().state,
ObservationState::Stale
);
let replanned = client
.ensure_observed(&session, "/child", 1)
.unwrap()
.expect("superseded lookahead is missing rather than permanently loading");
assert_eq!(
client
.complete(
older_child.id,
encode_test_observation_response(
&ObserveRequest {
path: "/child".into(),
depth: 1,
},
&TestObservationResponse {
root: TestObservation {
object: object("child", None, b"older"),
children: vec![],
},
},
wire_limits(),
)
.unwrap(),
)
.unwrap(),
Completion::StaleResponseRejected(older_child.id)
);
assert!(matches!(
client.object(&session, "/child").unwrap().state,
ObservationState::Loading(id) if id == replanned.id
));
}
#[test]
fn call_context_fixes_membership_validators_descriptor_and_result_validation() {
let mut client = client(8, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
observe_object(
&mut client,
&session,
"/item",
object("item", Some("item-ref"), b"object-v1"),
);
observe_interface(
&mut client,
&session,
FetchInterfaceResponse {
interface: "math".into(),
descriptor: descriptor(TypeExpr::Integer),
validator: Some(b"interface-v1".to_vec()),
},
);
let call = client
.prepare_call(&session, "/item", "math", "calculate", args(7))
.unwrap();
assert_eq!(
client
.pending_call_context(call.id)
.unwrap()
.operation()
.name,
"calculate"
);
client.mark_dispatched(call.id).unwrap();
observe_interface(
&mut client,
&session,
FetchInterfaceResponse {
interface: "math".into(),
descriptor: descriptor(TypeExpr::String),
validator: Some(b"interface-v2".to_vec()),
},
);
let history_context = match client
.complete(
call.id,
encode_call_operation_response(
&wip_protocol::CallOperationRequest {
target: wip_protocol::Target {
path: "/item".into(),
validator: Some(b"object-v1".to_vec()),
},
interface: wip_protocol::InterfaceTarget {
reference: "math".into(),
validator: Some(b"interface-v1".to_vec()),
},
operation: "calculate".into(),
arguments: args(7),
},
&descriptor(TypeExpr::Integer),
&CallOperationResponse {
result: Value::Integer(14),
validator: Some(b"object-v2".to_vec()),
},
wire_limits(),
)
.unwrap(),
)
.unwrap()
{
Completion::Call(record) => record,
other => panic!("unexpected completion: {other:?}"),
};
assert_eq!(
history_context.context.descriptor(),
&descriptor(TypeExpr::Integer)
);
assert_eq!(
history_context.context.object().validator,
Some(b"object-v1".to_vec())
);
assert_eq!(
history_context.context.interface().validator,
Some(b"interface-v1".to_vec())
);
assert!(matches!(history_context.outcome, CallOutcome::Success(_)));
assert_eq!(
client.interface(&session, "math").unwrap().descriptor,
Some(descriptor(TypeExpr::String))
);
let observation = client.object(&session, "/item").unwrap();
assert_eq!(observation.state, ObservationState::Fresh);
assert_eq!(observation.validator, Some(b"object-v2".to_vec()));
}
#[test]
fn successful_call_without_validator_preserves_observed_validator() {
let mut client = client(8, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
observe_object(
&mut client,
&session,
"/item",
object("item", None, b"object-v1"),
);
observe_interface(
&mut client,
&session,
FetchInterfaceResponse {
interface: "math".into(),
descriptor: descriptor(TypeExpr::Integer),
validator: None,
},
);
let call = client
.prepare_call(&session, "/item", "math", "calculate", args(1))
.unwrap();
let request = client
.pending_call_context(call.id)
.unwrap()
.request()
.clone();
let fixed_descriptor = client
.pending_call_context(call.id)
.unwrap()
.descriptor()
.clone();
client.mark_dispatched(call.id).unwrap();
client
.complete(
call.id,
encode_call_operation_response(
&request,
&fixed_descriptor,
&CallOperationResponse {
result: Value::Integer(2),
validator: None,
},
wire_limits(),
)
.unwrap(),
)
.unwrap();
let observation = client.object(&session, "/item").unwrap();
assert_eq!(observation.state, ObservationState::Fresh);
assert_eq!(observation.validator, Some(b"object-v1".to_vec()));
let next = client
.prepare_call(&session, "/item", "math", "calculate", args(2))
.unwrap();
assert_eq!(
client
.pending_call_context(next.id)
.unwrap()
.request()
.target
.validator,
Some(b"object-v1".to_vec())
);
}
#[test]
fn newer_observation_prevents_older_success_validator_overwrite() {
let mut client = client(8, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
observe_object(
&mut client,
&session,
"/item",
object("item", None, b"object-v1"),
);
observe_interface(
&mut client,
&session,
FetchInterfaceResponse {
interface: "math".into(),
descriptor: descriptor(TypeExpr::Integer),
validator: None,
},
);
let call = client
.prepare_call(&session, "/item", "math", "calculate", args(1))
.unwrap();
let request = client
.pending_call_context(call.id)
.unwrap()
.request()
.clone();
let fixed_descriptor = client
.pending_call_context(call.id)
.unwrap()
.descriptor()
.clone();
client.mark_dispatched(call.id).unwrap();
observe_object(
&mut client,
&session,
"/item",
object("item", None, b"object-v3"),
);
client
.complete(
call.id,
encode_call_operation_response(
&request,
&fixed_descriptor,
&CallOperationResponse {
result: Value::Integer(2),
validator: Some(b"object-v2".to_vec()),
},
wire_limits(),
)
.unwrap(),
)
.unwrap();
let observation = client.object(&session, "/item").unwrap();
assert_eq!(observation.state, ObservationState::Fresh);
assert_eq!(observation.validator, Some(b"object-v3".to_vec()));
}
#[test]
fn older_retrieval_cannot_overwrite_success_validator() {
let mut client = client(8, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
observe_object(
&mut client,
&session,
"/item",
object("item", Some("shared-ref"), b"object-v1"),
);
observe_interface(
&mut client,
&session,
FetchInterfaceResponse {
interface: "math".into(),
descriptor: descriptor(TypeExpr::Integer),
validator: Some(b"interface-v1".to_vec()),
},
);
let older_alias = client.refresh_observed(&session, "/alias/item", 0).unwrap();
let call = client
.prepare_call(&session, "/item", "math", "calculate", args(7))
.unwrap();
let request = client
.pending_call_context(call.id)
.unwrap()
.request()
.clone();
let fixed_descriptor = client
.pending_call_context(call.id)
.unwrap()
.descriptor()
.clone();
client.mark_dispatched(call.id).unwrap();
client
.complete(
call.id,
encode_call_operation_response(
&request,
&fixed_descriptor,
&CallOperationResponse {
result: Value::Integer(14),
validator: Some(b"object-v2".to_vec()),
},
wire_limits(),
)
.unwrap(),
)
.unwrap();
let current = client.object(&session, "/item").unwrap();
assert_eq!(current.state, ObservationState::Fresh);
assert_eq!(current.validator, Some(b"object-v2".to_vec()));
client
.complete(
older_alias.id,
observation_response(
"/alias/item",
&object("item", Some("shared-ref"), b"object-v1"),
),
)
.unwrap();
let alias = client.object(&session, "/alias/item").unwrap();
assert_eq!(alias.state, ObservationState::Stale);
assert_eq!(alias.validator, Some(b"object-v2".to_vec()));
assert_eq!(
client
.prepare_call(&session, "/alias/item", "math", "calculate", args(1))
.unwrap_err()
.kind,
ClientErrorKind::MissingObservation
);
observe_object(
&mut client,
&session,
"/item",
object("item", Some("shared-ref"), b"object-v2"),
);
let older_alias = client
.refresh_observed(&session, "/second/item", 0)
.unwrap();
let call = client
.prepare_call(&session, "/item", "math", "calculate", args(8))
.unwrap();
client.mark_dispatched(call.id).unwrap();
let unknown = ProtocolError {
code: ProtocolErrorCode::OperationOutcomeUnknown,
message: "commit uncertain".into(),
};
client
.complete(
call.id,
encode_protocol_error_response(
ProtocolInteraction::CallOperation,
&unknown,
wire_limits(),
)
.unwrap(),
)
.unwrap();
client
.complete(
older_alias.id,
observation_response(
"/second/item",
&object("item", Some("shared-ref"), b"object-v2"),
),
)
.unwrap();
assert_eq!(
client.object(&session, "/second/item").unwrap().state,
ObservationState::Stale
);
}
#[test]
fn in_flight_calls_pin_target_identity_until_staleness_is_applied() {
let mut client = client(1, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
observe_object(
&mut client,
&session,
"/item",
object("item", Some("shared-ref"), b"object-v1"),
);
observe_interface(
&mut client,
&session,
FetchInterfaceResponse {
interface: "math".into(),
descriptor: descriptor(TypeExpr::Integer),
validator: Some(b"interface-v1".to_vec()),
},
);
let call = client
.prepare_call(&session, "/item", "math", "calculate", args(7))
.unwrap();
client.mark_dispatched(call.id).unwrap();
assert_eq!(
client
.refresh_observed(&session, "/other", 0)
.unwrap_err()
.kind,
ClientErrorKind::Capacity
);
let unknown = ProtocolError {
code: ProtocolErrorCode::OperationOutcomeUnknown,
message: "commit uncertain".into(),
};
client
.complete(
call.id,
encode_protocol_error_response(
ProtocolInteraction::CallOperation,
&unknown,
wire_limits(),
)
.unwrap(),
)
.unwrap();
assert_eq!(
client.object(&session, "/item").unwrap().state,
ObservationState::Stale
);
observe_object(
&mut client,
&session,
"/other",
object("other", None, b"other"),
);
assert!(client.object(&session, "/item").is_none());
assert!(client.object(&session, "/other").is_some());
}
#[test]
fn validator_mismatch_stales_only_its_observation_and_never_retries() {
let mut client = client(8, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
observe_object(&mut client, &session, "/item", object("item", None, b"v1"));
observe_interface(
&mut client,
&session,
FetchInterfaceResponse {
interface: "math".into(),
descriptor: descriptor(TypeExpr::Integer),
validator: Some(b"i1".to_vec()),
},
);
let call = client
.prepare_call(&session, "/item", "math", "calculate", args(1))
.unwrap();
client.mark_dispatched(call.id).unwrap();
let error = ProtocolError {
code: ProtocolErrorCode::InterfaceValidatorMismatch,
message: "stale descriptor".into(),
};
let response =
encode_protocol_error_response(ProtocolInteraction::CallOperation, &error, wire_limits())
.unwrap();
client.complete(call.id, response).unwrap();
assert_eq!(
client.interface(&session, "math").unwrap().state,
ObservationState::Stale
);
assert_eq!(
client.object(&session, "/item").unwrap().state,
ObservationState::Fresh
);
assert_eq!(client.in_flight_len(), 0);
assert_eq!(client.call_history(&session).unwrap().len(), 1);
observe_interface(
&mut client,
&session,
FetchInterfaceResponse {
interface: "math".into(),
descriptor: descriptor(TypeExpr::Integer),
validator: Some(b"i2".to_vec()),
},
);
assert_eq!(client.in_flight_len(), 0);
assert_eq!(client.call_history(&session).unwrap().len(), 1);
}
#[test]
fn protocol_and_transport_unknown_outcomes_stale_without_retry() {
let mut client = client(8, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
observe_object(&mut client, &session, "/item", object("item", None, b"v1"));
observe_interface(
&mut client,
&session,
FetchInterfaceResponse {
interface: "math".into(),
descriptor: descriptor(TypeExpr::Integer),
validator: None,
},
);
let call = client
.prepare_call(&session, "/item", "math", "calculate", args(1))
.unwrap();
client.mark_dispatched(call.id).unwrap();
let unknown = ProtocolError {
code: ProtocolErrorCode::OperationOutcomeUnknown,
message: "commit uncertain".into(),
};
let completion = client
.complete(
call.id,
encode_protocol_error_response(
ProtocolInteraction::CallOperation,
&unknown,
wire_limits(),
)
.unwrap(),
)
.unwrap();
assert!(matches!(
completion,
Completion::Call(record)
if matches!(
record.outcome,
CallOutcome::Unknown {
reason: OutcomeUnknownReason::Protocol,
..
}
)
));
assert_eq!(
client.object(&session, "/item").unwrap().state,
ObservationState::Stale
);
observe_object(&mut client, &session, "/item", object("item", None, b"v2"));
let second = client
.prepare_call(&session, "/item", "math", "calculate", args(2))
.unwrap();
client.mark_dispatched(second.id).unwrap();
let completion = client
.fail_dispatched_call(
second.id,
DispatchedTransportFailure::Timeout,
"deadline elapsed",
)
.unwrap();
assert!(matches!(
completion,
Completion::Call(record)
if matches!(
record.outcome,
CallOutcome::Unknown {
reason: OutcomeUnknownReason::Timeout,
..
}
)
));
assert_eq!(client.in_flight_len(), 0);
assert_eq!(client.call_history(&session).unwrap().len(), 2);
}
#[test]
fn overlapping_tree_work_preserves_newer_subject_request_ownership() {
let mut client = client(8, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
let tree_request = ObserveRequest {
path: "/".into(),
depth: 1,
};
let tree_value = TestObservationResponse {
root: TestObservation {
object: object("", None, b"root"),
children: vec![TestObservation {
object: object("child", None, b"tree-child"),
children: vec![],
}],
},
};
let old_tree = client.refresh_observed(&session, "/", 1).unwrap();
let newer_observation = client.refresh_observed(&session, "/child", 0).unwrap();
client
.complete(
old_tree.id,
encode_test_observation_response(&tree_request, &tree_value, wire_limits()).unwrap(),
)
.unwrap();
assert!(matches!(
client.object(&session, "/child").unwrap().state,
ObservationState::Loading(id) if id == newer_observation.id
));
client
.complete(
newer_observation.id,
observation_response("/child", &object("child", None, b"newer-observation")),
)
.unwrap();
assert_eq!(
client.object(&session, "/child").unwrap().validator,
Some(b"newer-observation".to_vec())
);
let older_tree = client.refresh_observed(&session, "/", 1).unwrap();
let completed_newer_observation = client.refresh_observed(&session, "/child", 0).unwrap();
client
.complete(
completed_newer_observation.id,
observation_response("/child", &object("child", None, b"completed-newer")),
)
.unwrap();
client
.complete(
older_tree.id,
encode_test_observation_response(&tree_request, &tree_value, wire_limits()).unwrap(),
)
.unwrap();
assert_eq!(
client.object(&session, "/child").unwrap().validator,
Some(b"completed-newer".to_vec())
);
let older_omitting_tree = client.refresh_observed(&session, "/", 1).unwrap();
let newer_observation = client.refresh_observed(&session, "/child", 0).unwrap();
client
.complete(
newer_observation.id,
observation_response("/child", &object("child", None, b"newer-than-omission")),
)
.unwrap();
client
.complete(
older_omitting_tree.id,
encode_test_observation_response(
&tree_request,
&TestObservationResponse {
root: TestObservation {
object: object("", None, b"older-root"),
children: vec![],
},
},
wire_limits(),
)
.unwrap(),
)
.unwrap();
let child = client.object(&session, "/child").unwrap();
assert_eq!(child.state, ObservationState::Fresh);
assert_eq!(child.validator, Some(b"newer-than-omission".to_vec()));
let old_tree = client.refresh_observed(&session, "/", 1).unwrap();
let newer_child_tree = client.refresh_observed(&session, "/child", 1).unwrap();
client
.complete(
old_tree.id,
encode_test_observation_response(&tree_request, &tree_value, wire_limits()).unwrap(),
)
.unwrap();
assert!(matches!(
client.tree(&session, "/child").unwrap().state,
ObservationState::Loading(id) if id == newer_child_tree.id
));
let child_observation = client.object(&session, "/child").unwrap();
assert!(matches!(
child_observation.state,
ObservationState::Loading(id) if id == newer_child_tree.id
));
assert_eq!(
child_observation.validator,
Some(b"newer-than-omission".to_vec())
);
let child_request = ObserveRequest {
path: "/child".into(),
depth: 1,
};
let child_value = TestObservationResponse {
root: TestObservation {
object: object("child", None, b"newer-tree"),
children: vec![],
},
};
client
.complete(
newer_child_tree.id,
encode_test_observation_response(&child_request, &child_value, wire_limits()).unwrap(),
)
.unwrap();
assert_eq!(
client.tree(&session, "/child").unwrap().state,
ObservationState::Fresh
);
let older_parent = client.refresh_observed(&session, "/", 1).unwrap();
let newer_child = client.refresh_observed(&session, "/child", 1).unwrap();
client
.complete(
newer_child.id,
encode_test_observation_response(
&child_request,
&TestObservationResponse {
root: TestObservation {
object: object("child", None, b"completed-newer-tree"),
children: vec![],
},
},
wire_limits(),
)
.unwrap(),
)
.unwrap();
client
.complete(
older_parent.id,
encode_test_observation_response(&tree_request, &tree_value, wire_limits()).unwrap(),
)
.unwrap();
assert_eq!(
client.object(&session, "/child").unwrap().validator,
Some(b"completed-newer-tree".to_vec())
);
assert_eq!(
client.tree(&session, "/child").unwrap().revision,
newer_child.id
);
}
#[test]
fn superseded_transport_failure_does_not_poison_newer_refresh() {
let mut client = client(8, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
let old = client.refresh_observed(&session, "/item", 0).unwrap();
let newer = client.refresh_observed(&session, "/item", 0).unwrap();
assert_eq!(
client
.fail_transport(old.id, "old connection failed")
.unwrap(),
Completion::StaleResponseRejected(old.id)
);
assert!(matches!(
client.object(&session, "/item").unwrap().state,
ObservationState::Loading(id) if id == newer.id
));
client
.complete(
newer.id,
observation_response("/item", &object("item", None, b"new")),
)
.unwrap();
let old = client.refresh_observed(&session, "/", 1).unwrap();
let newer_tree = client.refresh_observed(&session, "/", 1).unwrap();
assert_eq!(
client
.fail_transport(old.id, "old tree connection failed")
.unwrap(),
Completion::StaleResponseRejected(old.id)
);
assert!(matches!(
client.tree(&session, "/").unwrap().state,
ObservationState::Loading(id) if id == newer_tree.id
));
client
.complete(
newer_tree.id,
encode_test_observation_response(
&ObserveRequest {
path: "/".into(),
depth: 1,
},
&TestObservationResponse {
root: TestObservation {
object: object("", None, b"root"),
children: vec![],
},
},
wire_limits(),
)
.unwrap(),
)
.unwrap();
let old = client.prepare_interface(&session, "math").unwrap();
let newer = client.prepare_interface(&session, "math").unwrap();
assert_eq!(
client
.fail_transport(old.id, "old connection failed")
.unwrap(),
Completion::StaleResponseRejected(old.id)
);
assert!(matches!(
client.interface(&session, "math").unwrap().state,
ObservationState::Loading(id) if id == newer.id
));
let value = FetchInterfaceResponse {
interface: "math".into(),
descriptor: descriptor(TypeExpr::Integer),
validator: Some(b"current".to_vec()),
};
client
.complete(
newer.id,
encode_fetch_interface_response(
&FetchInterfaceRequest {
interface: "math".into(),
},
&value,
wire_limits(),
)
.unwrap(),
)
.unwrap();
}
#[test]
fn validator_free_interface_observation_remains_immutable() {
let mut client = client(8, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
observe_interface(
&mut client,
&session,
FetchInterfaceResponse {
interface: "math".into(),
descriptor: descriptor(TypeExpr::Integer),
validator: None,
},
);
let refresh = client.prepare_interface(&session, "math").unwrap();
let contradictory = FetchInterfaceResponse {
interface: "math".into(),
descriptor: descriptor(TypeExpr::Integer),
validator: Some(b"became-mutable".to_vec()),
};
let error = client
.complete(
refresh.id,
encode_fetch_interface_response(
&FetchInterfaceRequest {
interface: "math".into(),
},
&contradictory,
wire_limits(),
)
.unwrap(),
)
.unwrap_err();
assert_eq!(error.kind, ClientErrorKind::InvalidResponse);
assert_eq!(client.interface(&session, "math").unwrap().validator, None);
}
#[test]
fn entry_eviction_never_discards_a_loading_tree_subject() {
let mut client = client(1, 8, 8, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
observe_object(&mut client, &session, "/", object("", None, b"root"));
let tree = client.refresh_observed(&session, "/", 1).unwrap();
assert_eq!(
client
.refresh_observed(&session, "/other", 0)
.unwrap_err()
.kind,
ClientErrorKind::Capacity
);
assert!(matches!(
client.tree(&session, "/").unwrap().state,
ObservationState::Loading(id) if id == tree.id
));
client
.complete(
tree.id,
encode_test_observation_response(
&ObserveRequest {
path: "/".into(),
depth: 1,
},
&TestObservationResponse {
root: TestObservation {
object: object("", None, b"tree"),
children: vec![],
},
},
wire_limits(),
)
.unwrap(),
)
.unwrap();
assert_eq!(
client.tree(&session, "/").unwrap().state,
ObservationState::Fresh
);
}
#[test]
fn failed_entry_and_tree_observation_metadata_is_bounded() {
let mut client = client(2, 8, 2, 8);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
for path in ["/one", "/two", "/three"] {
let request = client.refresh_observed(&session, path, 0).unwrap();
assert_eq!(
client
.fail_transport(request.id, "retrieval failed")
.unwrap_err()
.kind,
ClientErrorKind::Transport
);
}
assert!(client.object(&session, "/one").is_none());
assert!(client.object(&session, "/two").is_some());
assert!(client.object(&session, "/three").is_some());
for path in ["/alpha", "/beta", "/gamma"] {
let request = client.refresh_observed(&session, path, 1).unwrap();
assert_eq!(
client
.fail_transport(request.id, "tree failed")
.unwrap_err()
.kind,
ClientErrorKind::Transport
);
}
assert!(client.tree(&session, "/alpha").is_none());
assert!(client.tree(&session, "/gamma").is_some());
}
#[test]
fn in_flight_and_history_are_bounded() {
let mut client = client(8, 8, 1, 1);
let session = client
.open_session("https://example.test/wip", SecurityContext::new("subject"))
.unwrap();
let first = client.refresh_observed(&session, "/one", 0).unwrap();
let error = client.refresh_observed(&session, "/two", 0).unwrap_err();
assert_eq!(error.kind, ClientErrorKind::Capacity);
client
.complete(
first.id,
observation_response("/one", &object("one", None, b"1")),
)
.unwrap();
let stale_old = client.refresh_observed(&session, "/one", 0).unwrap();
let stale_new = client.refresh_observed(&session, "/one", 0).unwrap_err();
assert_eq!(stale_new.kind, ClientErrorKind::Capacity);
client
.complete(
stale_old.id,
Response::builder()
.status(200)
.header("content-type", "application/json")
.body(b"not json".to_vec())
.unwrap(),
)
.unwrap_err();
assert_eq!(client.diagnostics(&session).unwrap().len(), 1);
}