use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use frame_conv::id::CorrelationId;
use super::{BarrierEngine, DevReply, DevRequest, DevStatusEvent, NodeControl};
pub struct InboundManagement {
pub request: DevRequest,
pub correlation: CorrelationId,
}
pub trait ManagementConversation {
fn next_request(&mut self, wait: Duration) -> Result<Option<InboundManagement>, String>;
fn reply(&mut self, correlation: CorrelationId, reply: &DevReply) -> Result<(), String>;
fn publish_status(&mut self, event: &DevStatusEvent) -> Result<(), String>;
}
#[derive(Debug, PartialEq, Eq)]
pub enum ServeExit {
Closed,
ConnectionLost {
detail: String,
},
}
pub fn serve<C: ManagementConversation>(
conversation: &mut C,
engine: &mut BarrierEngine,
control: &mut dyn NodeControl,
close: &AtomicBool,
quantum: Duration,
) -> ServeExit {
loop {
if close.load(Ordering::Acquire) {
return ServeExit::Closed;
}
match conversation.next_request(quantum) {
Ok(None) => {}
Ok(Some(inbound)) => {
let mut publish_failure: Option<String> = None;
let reply = {
let failure = &mut publish_failure;
engine.handle(inbound.request, control, &mut |event| {
if failure.is_none()
&& let Err(detail) = conversation_publish(conversation, &event)
{
*failure = Some(detail);
}
})
};
if let Err(detail) = conversation.reply(inbound.correlation, &reply) {
return ServeExit::ConnectionLost { detail };
}
if let Some(detail) = publish_failure {
return ServeExit::ConnectionLost { detail };
}
}
Err(detail) => return ServeExit::ConnectionLost { detail },
}
}
}
fn conversation_publish<C: ManagementConversation>(
conversation: &mut C,
event: &DevStatusEvent,
) -> Result<(), String> {
conversation.publish_status(event)
}
#[cfg(test)]
mod tests {
#![allow(clippy::expect_used, clippy::panic)]
use std::sync::atomic::AtomicBool;
use std::time::Duration;
use frame_conv::id::CorrelationId;
use super::{InboundManagement, ManagementConversation, ServeExit, serve};
use crate::dev::{
BarrierEngine, CandidateBytes, DevReply, DevRequest, DevStatusEvent, GenerationReport,
NodeControl, NodeMode, StageRefusal, StartFailure,
};
const QUANTUM: Duration = Duration::from_millis(1);
struct FakeConversation {
script: Vec<Option<InboundManagement>>,
replies: Vec<DevReply>,
published: Vec<DevStatusEvent>,
pump_ticks: usize,
}
impl FakeConversation {
fn new(mut script: Vec<Option<InboundManagement>>) -> Self {
script.reverse();
Self {
script,
replies: Vec::new(),
published: Vec::new(),
pump_ticks: 0,
}
}
}
impl ManagementConversation for FakeConversation {
fn next_request(&mut self, _wait: Duration) -> Result<Option<InboundManagement>, String> {
self.pump_ticks += 1;
match self.script.pop() {
Some(entry) => Ok(entry),
None => Err("script exhausted: connection closed".to_owned()),
}
}
fn reply(&mut self, _correlation: CorrelationId, reply: &DevReply) -> Result<(), String> {
self.replies.push(reply.clone());
Ok(())
}
fn publish_status(&mut self, event: &DevStatusEvent) -> Result<(), String> {
self.published.push(event.clone());
Ok(())
}
}
#[derive(Default)]
struct CountingNode {
operations: usize,
}
impl NodeControl for CountingNode {
fn stop_component(&mut self) -> Result<(), String> {
self.operations += 1;
Ok(())
}
fn stage(&mut self, _candidate: &CandidateBytes) -> Result<(), String> {
self.operations += 1;
Ok(())
}
fn start_component(&mut self) -> Result<(), StartFailure> {
self.operations += 1;
Ok(())
}
fn witness_mailbox(&mut self) -> Result<(), String> {
self.operations += 1;
Ok(())
}
fn witness_content(&mut self) -> Result<(), String> {
self.operations += 1;
Ok(())
}
fn report(
&mut self,
build_generation: u64,
content_digest: [u8; 32],
) -> Result<GenerationReport, String> {
self.operations += 1;
Ok(GenerationReport {
build_generation,
content_digest,
component_incarnation: 1,
module_generation: 1,
})
}
}
fn candidate(generation: u64) -> CandidateBytes {
CandidateBytes {
build_generation: generation,
content_digest: [0; 32],
component_beam: vec![1],
ffi_beam: vec![2],
}
}
fn inbound(generation: u64) -> InboundManagement {
InboundManagement {
request: DevRequest::StageAndActivate(candidate(generation)),
correlation: CorrelationId::mint(),
}
}
#[test]
fn quiet_ticks_do_no_work() {
let mut conversation = FakeConversation::new(vec![None, None, None, None, None]);
let mut engine = BarrierEngine::new(NodeMode::Dev, candidate(0));
let mut node = CountingNode::default();
let close = AtomicBool::new(false);
let exit = serve(&mut conversation, &mut engine, &mut node, &close, QUANTUM);
assert!(matches!(exit, ServeExit::ConnectionLost { .. }));
assert_eq!(conversation.pump_ticks, 6, "five quiet ticks + the fate");
assert_eq!(node.operations, 0, "quiet ticks touch the node ZERO times");
assert!(
conversation.published.is_empty(),
"quiet ticks publish nothing"
);
assert!(
conversation.replies.is_empty(),
"quiet ticks reply to nothing"
);
}
#[test]
fn requests_are_served_between_quiet_ticks() {
let mut conversation = FakeConversation::new(vec![None, Some(inbound(1)), None]);
let mut engine = BarrierEngine::new(NodeMode::Dev, candidate(0));
let mut node = CountingNode::default();
let close = AtomicBool::new(false);
let exit = serve(&mut conversation, &mut engine, &mut node, &close, QUANTUM);
assert!(matches!(exit, ServeExit::ConnectionLost { .. }));
assert_eq!(conversation.replies.len(), 1, "exactly one reply");
assert!(matches!(conversation.replies[0], DevReply::Activated(_)));
assert!(
matches!(
conversation.published.last(),
Some(DevStatusEvent::RunningCurrent(_))
),
"the barrier's status transitions were published"
);
}
#[test]
fn close_flag_exits_ordered() {
let mut conversation = FakeConversation::new(vec![None]);
let mut engine = BarrierEngine::new(NodeMode::Dev, candidate(0));
let mut node = CountingNode::default();
let close = AtomicBool::new(true);
let exit = serve(&mut conversation, &mut engine, &mut node, &close, QUANTUM);
assert_eq!(exit, ServeExit::Closed);
assert_eq!(conversation.pump_ticks, 0, "no pump arm after close");
}
#[test]
fn production_refusal_holds_through_the_adapter() {
let mut conversation = FakeConversation::new(vec![Some(inbound(1))]);
let mut engine = BarrierEngine::new(NodeMode::Production, candidate(0));
let mut node = CountingNode::default();
let close = AtomicBool::new(false);
let _exit = serve(&mut conversation, &mut engine, &mut node, &close, QUANTUM);
assert_eq!(
conversation.replies[0],
DevReply::Refused(StageRefusal::NotADevNode)
);
assert_eq!(node.operations, 0);
assert!(conversation.published.is_empty());
}
}