frame_host/dev/
conversation.rs1use std::time::Duration;
13
14use frame_conv::error::CallError;
15use frame_conv::handle::ConversationHandle;
16use frame_conv::id::{ConversationId, CorrelationId};
17use frame_conv::outcome::InboundRequest;
18use frame_conv::store::ResumeStore;
19use frame_conv::{AttachError, BusAttachment};
20
21use super::adapter::{InboundManagement, ManagementConversation};
22use super::{DevReply, DevRequest, DevStatusEvent, StageRefusal};
23
24pub const MANAGEMENT_LINE_PREFIX: &str = "dev management at ";
27
28#[must_use]
30pub fn format_management_line(endpoint: &str, conversation: ConversationId) -> String {
31 format!("{MANAGEMENT_LINE_PREFIX}{endpoint} conversation {conversation}")
32}
33
34#[must_use]
37pub fn parse_management_line(line: &str) -> Option<(String, ConversationId)> {
38 let after = line.split(MANAGEMENT_LINE_PREFIX).nth(1)?;
39 let (endpoint, id) = after.split_once(" conversation ")?;
40 let conversation = id.trim().parse().ok()?;
41 Some((endpoint.to_owned(), conversation))
42}
43
44#[derive(Default)]
49pub struct DevResumeStore(Vec<u8>);
50
51impl ResumeStore for DevResumeStore {
52 fn persist(&mut self, canonical_state: &[u8]) -> Result<(), frame_conv::error::StoreError> {
53 self.0 = canonical_state.to_vec();
54 Ok(())
55 }
56}
57
58pub struct DevManagementConversation<S: ResumeStore> {
61 handle: ConversationHandle<S>,
62}
63
64impl<S: ResumeStore> DevManagementConversation<S> {
65 pub fn open(
72 attachment: &BusAttachment,
73 store: S,
74 ) -> Result<(Self, ConversationId), AttachError> {
75 let (handle, grant) = ConversationHandle::open(attachment, store)?;
76 Ok((Self { handle }, grant.conversation))
77 }
78
79 pub fn join(
85 attachment: &BusAttachment,
86 conversation: ConversationId,
87 store: S,
88 ) -> Result<Self, AttachError> {
89 let (handle, _grant) = ConversationHandle::join(attachment, conversation, store)?;
90 Ok(Self { handle })
91 }
92
93 pub fn handle_mut(&mut self) -> &mut ConversationHandle<S> {
97 &mut self.handle
98 }
99}
100
101impl<S: ResumeStore> ManagementConversation for DevManagementConversation<S> {
102 fn next_request(&mut self, wait: Duration) -> Result<Option<InboundManagement>, String> {
103 match self.handle.next_request::<DevRequest>(wait) {
104 Ok(None) => Ok(None),
105 Ok(Some(InboundRequest::Valid(incoming))) => Ok(Some(InboundManagement {
106 request: incoming.body,
107 correlation: incoming.correlation,
108 })),
109 Ok(Some(InboundRequest::SchemaInvalid {
110 correlation,
111 detail,
112 ..
113 })) => {
114 let refusal = DevReply::Refused(StageRefusal::SchemaInvalid {
118 detail: detail.clone(),
119 });
120 self.handle
121 .reply(correlation, &refusal)
122 .map_err(|error| format!("schema-invalid refusal reply: {error}"))?;
123 tracing::error!(%detail, "management request failed schema validation; refused typed");
124 Ok(None)
125 }
126 Err(CallError::ConnectionLost { detail }) => Err(detail),
127 Err(other) => Err(other.to_string()),
128 }
129 }
130
131 fn reply(&mut self, correlation: CorrelationId, reply: &DevReply) -> Result<(), String> {
132 self.handle
133 .reply(correlation, reply)
134 .map(|_| ())
135 .map_err(|error| error.to_string())
136 }
137
138 fn publish_status(&mut self, event: &DevStatusEvent) -> Result<(), String> {
139 self.handle
140 .publish_event(event)
141 .map(|_| ())
142 .map_err(|error| error.to_string())
143 }
144}
145
146#[cfg(test)]
147mod tests {
148 #![allow(clippy::expect_used)]
149
150 use super::{format_management_line, parse_management_line};
151 use frame_conv::id::ConversationId;
152
153 #[test]
156 fn handoff_line_round_trips_and_rejects_noise() {
157 let id: ConversationId = "12345".parse().expect("id parses");
158 let line = format_management_line("127.0.0.1:9410", id);
159 assert_eq!(
160 parse_management_line(&line),
161 Some(("127.0.0.1:9410".to_owned(), id))
162 );
163 assert_eq!(
164 parse_management_line(&format!("prefix {line}")),
165 Some(("127.0.0.1:9410".to_owned(), id)),
166 "a logger prefix does not defeat the parse"
167 );
168 assert!(parse_management_line("demo serving at http://127.0.0.1:4191").is_none());
169 assert!(parse_management_line("dev management at half a line").is_none());
170 assert!(
171 parse_management_line("dev management at 1.2.3.4:1 conversation notanumber").is_none()
172 );
173 }
174}