Skip to main content

core_api/
ingress.rs

1//! Transport-independent ingress rules. Identity always comes from the authenticated host session.
2use crate::{
3    ErrorResponse, ExceptionCode, MwsMessage, MwsMessageType, ServiceCoreInput, ServiceCoreResponse,
4};
5use std::time::Duration;
6
7pub const CORE_REQUEST_TIMEOUT: Duration = Duration::from_secs(60);
8pub const IPC_TIMEOUT_GRACE: Duration = Duration::from_secs(5);
9
10pub fn request_timeout(target: &str) -> Duration {
11    match target {
12        "/harness/test/session/open" => Duration::from_secs(150),
13        _ => CORE_REQUEST_TIMEOUT,
14    }
15}
16
17pub fn input_timeout(input: &ServiceCoreInput) -> Duration {
18    match input {
19        ServiceCoreInput::ClientRequest { target, .. }
20        | ServiceCoreInput::ConversationRequest { target, .. }
21        | ServiceCoreInput::NodeRequest { target, .. } => request_timeout(target),
22        _ => CORE_REQUEST_TIMEOUT,
23    }
24}
25
26pub struct NodeSession<'a> {
27    pub connection_key: &'a str,
28    pub tenant_id: &'a str,
29    pub scope_id: &'a str,
30    pub node_type: &'a str,
31    pub node_id: &'a str,
32}
33
34pub fn node_input(
35    request: &MwsMessage,
36    session: NodeSession<'_>,
37) -> Result<ServiceCoreInput, ErrorResponse> {
38    let target = request.target.as_deref().unwrap_or_default();
39    if target.is_empty() {
40        return Err(ErrorResponse::new(
41            ExceptionCode::BadRequest,
42            "target is required",
43        ));
44    }
45    let payload = request
46        .payload
47        .as_deref()
48        .filter(|p| !p.is_empty())
49        .unwrap_or("{}");
50    if target == crate::node::status::STATUS_CHANGED_TARGET {
51        let status = serde_json::from_str(payload).map_err(|error| {
52            ErrorResponse::new(
53                ExceptionCode::BadRequest,
54                format!("invalid Node status payload: {error}"),
55            )
56        })?;
57        return Ok(ServiceCoreInput::NodeStatusObserved {
58            connection_key: session.connection_key.into(),
59            status: Box::new(status),
60        });
61    }
62    Ok(ServiceCoreInput::NodeRequest {
63        target: target.into(),
64        payload: payload.into(),
65        control: request.control.clone(),
66        tenant_id: session.tenant_id.into(),
67        scope_id: session.scope_id.into(),
68        node_type: session.node_type.into(),
69        node_id: session.node_id.into(),
70    })
71}
72
73/// Apply at the outer boundary, including authorization and validation failures.
74pub fn node_reply<T>(message_type: &str, response: Option<T>) -> Option<T> {
75    if message_type == MwsMessageType::NODE_DATA {
76        None
77    } else {
78        response
79    }
80}
81
82pub fn node_response(
83    request: &MwsMessage,
84    scope_id: &str,
85    result: Result<Option<ServiceCoreResponse>, ErrorResponse>,
86) -> Option<MwsMessage> {
87    let result = result.and_then(|response| match response {
88        Some(ServiceCoreResponse::NodeRequest { response_payload }) => {
89            Ok(response_payload.unwrap_or_else(|| "{}".into()))
90        }
91        Some(ServiceCoreResponse::Ack { handled: true })
92            if request.target.as_deref() == Some(crate::node::status::STATUS_CHANGED_TARGET) =>
93        {
94            Ok("{}".into())
95        }
96        _ => Err(ErrorResponse::new(
97            ExceptionCode::Internal,
98            "unexpected Core response for Node request",
99        )),
100    });
101    let (payload, error) = match result {
102        Ok(payload) => (Some(payload), None),
103        Err(error) => (None, Some(error)),
104    };
105    node_reply(
106        request.r#type.as_deref().unwrap_or_default(),
107        Some(MwsMessage {
108            scope_id: Some(scope_id.into()),
109            from: None,
110            target: request.target.clone(),
111            sig: request.sig.clone(),
112            r#type: Some(MwsMessageType::NODE_RESP.into()),
113            payload,
114            error,
115            control: None,
116            client_info: None,
117        }),
118    )
119}
120
121#[cfg(test)]
122mod tests {
123    use super::*;
124    fn request(target: &str, kind: &str, payload: &str) -> MwsMessage {
125        serde_json::from_value(serde_json::json!({"target":target,"type":kind,"payload":payload,"scopeId":"untrusted","sig":"s"})).unwrap()
126    }
127    fn session() -> NodeSession<'static> {
128        NodeSession {
129            connection_key: "connection",
130            tenant_id: "tenant",
131            scope_id: "scope",
132            node_type: "camera",
133            node_id: "node",
134        }
135    }
136    #[test]
137    fn regular_node_input_uses_authenticated_identity() {
138        match node_input(&request("/x", MwsMessageType::NODE_REQ, ""), session()).unwrap() {
139            ServiceCoreInput::NodeRequest {
140                tenant_id,
141                scope_id,
142                node_id,
143                payload,
144                ..
145            } => {
146                assert_eq!(
147                    (
148                        tenant_id.as_str(),
149                        scope_id.as_str(),
150                        node_id.as_str(),
151                        payload.as_str()
152                    ),
153                    ("tenant", "scope", "node", "{}")
154                );
155            }
156            _ => panic!("wrong input"),
157        }
158    }
159    #[test]
160    fn invalid_status_is_not_forwarded_as_an_ordinary_request() {
161        let error = node_input(
162            &request(
163                crate::node::status::STATUS_CHANGED_TARGET,
164                MwsMessageType::NODE_DATA,
165                "{}",
166            ),
167            session(),
168        )
169        .unwrap_err();
170        assert_eq!(error.error, ExceptionCode::BadRequest.as_str());
171    }
172    #[test]
173    fn status_uses_the_transport_connection_not_payload_identity() {
174        let payload = serde_json::json!({
175            "contract":crate::node::STATUS_CONTRACT,"serviceId":"camera","nodeType":"camera",
176            "processGeneration":"process","revision":1,"process":"running","updatedAtMs":1,
177            "instance":{"nodeId":"node","runtimeId":"r","hubId":"h","tenantId":"tenant","scopeId":"scope",
178                "identity":"commissioned","hub":"resolved","connection":"connected","runtime":"ready",
179                "runtimeGeneration":1,"connectionGeneration":1,"revision":1,"effective":"online","updatedAtMs":1},
180            "runtime":{"runtimeId":"r","runtime":"ready","generation":1,"revision":1,"updatedAtMs":1},
181            "connectionKey":"forged"
182        });
183        let request = request(
184            crate::node::status::STATUS_CHANGED_TARGET,
185            MwsMessageType::NODE_REQ,
186            &payload.to_string(),
187        );
188        match node_input(&request, session()).unwrap() {
189            ServiceCoreInput::NodeStatusObserved {
190                connection_key,
191                status,
192            } => {
193                assert_eq!(connection_key, "connection");
194                assert_eq!(status.instance.node_id, "node");
195            }
196            _ => panic!("status must use the status manager input"),
197        }
198        let response = node_response(
199            &request,
200            "scope",
201            Ok(Some(ServiceCoreResponse::Ack { handled: true })),
202        )
203        .unwrap();
204        assert_eq!(response.payload.as_deref(), Some("{}"));
205        assert!(response.error.is_none());
206    }
207    #[test]
208    fn windows_pipe_names_are_bounded_and_workdir_specific() {
209        use crate::external_core_pipe_name as name;
210        use std::path::Path;
211        assert_eq!(
212            name(Path::new("C:/Meow/dev/core.sock")),
213            name(Path::new(r"c:\meow\dev\core.sock"))
214        );
215        assert_ne!(
216            name(Path::new("C:/Meow/dev/core.sock")),
217            name(Path::new("C:/Meow/prod/core.sock"))
218        );
219        assert!(name(Path::new(&"a".repeat(5000))).len() < 256);
220    }
221    #[test]
222    fn reports_never_reply_and_requests_preserve_error_details() {
223        let error = ErrorResponse::new(ExceptionCode::PreconditionFail, "expired")
224            .with_details(serde_json::json!({"preparedActionStatus":"expired"}));
225        for target in ["", "/x", crate::node::status::STATUS_CHANGED_TARGET] {
226            assert!(node_response(
227                &request(target, MwsMessageType::NODE_DATA, "{}"),
228                "scope",
229                Err(error.clone())
230            )
231            .is_none());
232        }
233        let response = node_response(
234            &request("/x", MwsMessageType::NODE_REQ, "{}"),
235            "scope",
236            Err(error.clone()),
237        )
238        .unwrap();
239        assert_eq!(
240            serde_json::to_value(response.error.unwrap()).unwrap(),
241            serde_json::to_value(error).unwrap()
242        );
243        assert!(node_reply(MwsMessageType::NODE_DATA, Some("unauthorized")).is_none());
244    }
245    #[test]
246    fn ipc_budget_covers_core_budget_including_long_requests() {
247        assert_eq!(
248            request_timeout("/app/plugin/camera/watch/set").as_secs(),
249            60
250        );
251        assert_eq!(request_timeout("/harness/test/session/open").as_secs(), 150);
252        assert!(IPC_TIMEOUT_GRACE > Duration::ZERO);
253    }
254}