Skip to main content

nnrp_runtime/
server.rs

1use std::collections::BTreeMap;
2use std::fmt;
3use std::sync::{Arc, Mutex, MutexGuard};
4use std::time::{SystemTime, UNIX_EPOCH};
5
6use nnrp_core::{
7    validate_control_request_semantics, validate_partial_result_semantics,
8    validate_pressure_semantics, validate_profile_assignment, validate_progress_semantics,
9    validate_result_drop_header, validate_result_drop_reason_semantics,
10    validate_scheduling_semantics, validate_trace_context_semantics, BudgetMetadata,
11    CacheInvalidateMetadata, CacheMissMetadata, CacheObjectId, CacheObjectKind,
12    CacheReferenceMetadata, CapabilityMetadata, CommonHeader, ConnectionLifecycle,
13    ControlRequestMetadata, FlowUpdateMetadata, FrameSubmitMetadata, MessageType,
14    ObjectDeltaMetadata, ObjectDescriptorMetadata, ObjectReferenceMetadata, ObjectReleaseMetadata,
15    OperationCancelRequest, OperationDescriptor, OperationRegistry, PartialResultMetadata,
16    PressureMetadata, ProgressMetadata, RecoverableErrorMetadata, ResultDropReasonMetadata,
17    ResultHintMetadata, ResultPushMetadata, RetryAfterMetadata, RouteHintMetadata, RuntimeRole,
18    SchedulingMetadata, SchemaRegistry, SessionCloseAckMetadata, SessionCloseMetadata,
19    SessionCloseStatus, SessionMigrateAckMetadata, SessionMigrateMetadata, SessionOpenAckMetadata,
20    SessionOpenMetadata, SessionPatchAckMetadata, SessionPatchMetadata, SessionStatus,
21    SupersedeMetadata, TraceContextMetadata, TransportProbeAckMetadata, TransportProbeMetadata,
22    BUDGET_METADATA_LEN, CACHE_INVALIDATE_METADATA_LEN, CACHE_MISS_METADATA_LEN,
23    CACHE_REFERENCE_METADATA_LEN, CAPABILITY_METADATA_LEN, CONTROL_REQUEST_METADATA_LEN,
24    FLOW_UPDATE_METADATA_LEN, FRAME_SUBMIT_METADATA_LEN, OBJECT_DELTA_METADATA_LEN,
25    OBJECT_DESCRIPTOR_METADATA_LEN, OBJECT_REFERENCE_METADATA_LEN, OBJECT_RELEASE_METADATA_LEN,
26    PARTIAL_RESULT_METADATA_LEN, PRESSURE_METADATA_LEN, PROGRESS_METADATA_LEN,
27    RECOVERABLE_ERROR_METADATA_LEN, RESULT_DROP_REASON_DEADLINE_EXPIRED,
28    RESULT_DROP_REASON_METADATA_LEN, RESULT_PUSH_METADATA_LEN, RETRY_AFTER_METADATA_LEN,
29    ROUTE_HINT_METADATA_LEN, SCHEDULING_FLAG_EMIT_DROP_REASON, SCHEDULING_METADATA_LEN,
30    SESSION_ACK_FLAG_RESUME_ENABLED, SESSION_CLOSE_ACK_METADATA_LEN, SESSION_ERROR_LIMIT_REACHED,
31    SESSION_ERROR_NONE, SESSION_ERROR_PROFILE_UNSUPPORTED, SESSION_ERROR_RESUME_REJECTED,
32    SESSION_ERROR_SCHEMA_UNSUPPORTED, SESSION_FLAG_ALLOW_RESUME, SESSION_MIGRATE_ACK_METADATA_LEN,
33    SESSION_MIGRATE_METADATA_LEN, SESSION_OPEN_ACK_METADATA_LEN, SESSION_PATCH_ACK_METADATA_LEN,
34    SESSION_PATCH_METADATA_LEN, SUPERSEDE_METADATA_LEN, TRACE_CONTEXT_METADATA_LEN,
35};
36use tokio::net::TcpListener;
37
38use crate::{
39    BoxedFramedListener, BoxedFramedTransport, FramedListener, RuntimeError, RuntimePacket,
40    RuntimePressureState, RuntimeTransportKind, TcpFramedListener,
41};
42
43#[derive(Clone)]
44pub struct NnrpServerConfig {
45    pub transport: RuntimeTransportKind,
46    pub supported_profiles: Vec<u16>,
47    pub supported_cache_objects: Vec<CacheObjectKind>,
48    pub max_cache_objects: usize,
49    pub max_cache_object_bytes: u32,
50    pub schema_registry: SchemaRegistry,
51    pub resume_token_bytes: u32,
52    pub max_in_flight_operations: u16,
53    pub granted_operation_credit: u16,
54    pub lease_ttl_ms: u32,
55    pub resume_window_ms: u32,
56    pub application_policy: Arc<dyn NnrpServerPolicy>,
57}
58
59impl fmt::Debug for NnrpServerConfig {
60    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
61        formatter
62            .debug_struct("NnrpServerConfig")
63            .field("transport", &self.transport)
64            .field("supported_profiles", &self.supported_profiles)
65            .field("supported_cache_objects", &self.supported_cache_objects)
66            .field("max_cache_objects", &self.max_cache_objects)
67            .field("max_cache_object_bytes", &self.max_cache_object_bytes)
68            .field("schema_registry", &self.schema_registry)
69            .field("resume_token_bytes", &self.resume_token_bytes)
70            .field("max_in_flight_operations", &self.max_in_flight_operations)
71            .field("granted_operation_credit", &self.granted_operation_credit)
72            .field("lease_ttl_ms", &self.lease_ttl_ms)
73            .field("resume_window_ms", &self.resume_window_ms)
74            .field("application_policy", &"<dyn NnrpServerPolicy>")
75            .finish()
76    }
77}
78
79pub trait NnrpServerPolicy: Send + Sync {
80    fn validate_session_open(&self, open: &SessionOpenMetadata) -> Result<(), u32>;
81}
82
83#[derive(Debug, Default)]
84pub struct AllowAllServerPolicy;
85
86impl NnrpServerPolicy for AllowAllServerPolicy {
87    fn validate_session_open(&self, _open: &SessionOpenMetadata) -> Result<(), u32> {
88        Ok(())
89    }
90}
91
92impl Default for NnrpServerConfig {
93    fn default() -> Self {
94        Self {
95            transport: RuntimeTransportKind::Tcp,
96            supported_profiles: vec![nnrp_core::PROFILE_TOKEN],
97            supported_cache_objects: Vec::new(),
98            max_cache_objects: 0,
99            max_cache_object_bytes: 0,
100            schema_registry: SchemaRegistry::with_standard_preview3_profiles(),
101            resume_token_bytes: 24,
102            max_in_flight_operations: 4,
103            granted_operation_credit: 2,
104            lease_ttl_ms: 30_000,
105            resume_window_ms: 120_000,
106            application_policy: Arc::new(AllowAllServerPolicy),
107        }
108    }
109}
110
111impl NnrpServerConfig {
112    pub fn with_transport(mut self, transport: RuntimeTransportKind) -> Self {
113        self.transport = transport;
114        self
115    }
116
117    pub fn with_supported_profiles(mut self, profiles: impl Into<Vec<u16>>) -> Self {
118        self.supported_profiles = profiles.into();
119        self
120    }
121
122    pub fn with_supported_cache_objects(
123        mut self,
124        objects: impl Into<Vec<CacheObjectKind>>,
125    ) -> Self {
126        self.supported_cache_objects = objects.into();
127        self
128    }
129
130    pub fn with_cache_limits(mut self, max_objects: usize, max_object_bytes: u32) -> Self {
131        self.max_cache_objects = max_objects;
132        self.max_cache_object_bytes = max_object_bytes;
133        self
134    }
135
136    pub fn with_schema_registry(mut self, schema_registry: SchemaRegistry) -> Self {
137        self.schema_registry = schema_registry;
138        self
139    }
140
141    pub fn with_resume_token_bytes(mut self, resume_token_bytes: u32) -> Self {
142        self.resume_token_bytes = resume_token_bytes;
143        self
144    }
145
146    pub fn with_application_policy<P>(mut self, policy: P) -> Self
147    where
148        P: NnrpServerPolicy + 'static,
149    {
150        self.application_policy = Arc::new(policy);
151        self
152    }
153
154    fn validate_client_open(&self, open: &SessionOpenMetadata) -> Result<(), u32> {
155        if !self.supported_profiles.contains(&open.profile_id)
156            || validate_profile_assignment(open.profile_id).is_err()
157        {
158            return Err(SESSION_ERROR_PROFILE_UNSUPPORTED);
159        }
160
161        if self
162            .schema_registry
163            .get(open.schema_id, open.schema_version)
164            .is_none()
165        {
166            return Err(SESSION_ERROR_SCHEMA_UNSUPPORTED);
167        }
168
169        if open.max_in_flight_operations > self.max_in_flight_operations {
170            return Err(SESSION_ERROR_LIMIT_REACHED);
171        }
172
173        self.application_policy.validate_session_open(open)?;
174
175        Ok(())
176    }
177}
178
179pub struct NnrpServer {
180    listener: BoxedFramedListener,
181    config: NnrpServerConfig,
182    sessions: SharedSessionRegistry,
183}
184
185pub struct NnrpServerSession {
186    session_id: u32,
187    client_open: SessionOpenMetadata,
188    transport: BoxedFramedTransport,
189    lifecycle: ConnectionLifecycle,
190    operations: OperationRegistry,
191    frame_operations: BTreeMap<u32, u64>,
192    operation_frames: BTreeMap<u64, u32>,
193    pressure: RuntimePressureState,
194    cache_objects: Vec<CacheObjectId>,
195    max_cache_objects: usize,
196    sessions: SharedSessionRegistry,
197    pending_close: Option<SessionCloseMetadata>,
198}
199
200#[derive(Debug, Clone, PartialEq, Eq)]
201pub struct RuntimeSessionRecord {
202    pub session_id: u32,
203    pub profile_id: u16,
204    pub schema_id: u32,
205    pub schema_version: u32,
206    pub resume_enabled: bool,
207    pub resume_token_bytes: u32,
208    pub last_operation_id: u64,
209}
210
211type SharedSessionRegistry = Arc<Mutex<BTreeMap<u32, RuntimeSessionRecord>>>;
212
213#[derive(Debug, Clone, PartialEq, Eq)]
214pub struct NnrpSubmit {
215    pub operation_id: u64,
216    pub frame_id: u32,
217    pub metadata: FrameSubmitMetadata,
218    pub body: Vec<u8>,
219}
220
221#[derive(Debug, Clone, Copy, PartialEq, Eq)]
222pub struct NnrpCancel {
223    pub frame_id: u32,
224}
225
226#[derive(Debug, Clone, Copy, PartialEq, Eq)]
227pub struct NnrpMigration {
228    pub metadata: SessionMigrateMetadata,
229}
230
231#[derive(Debug, Clone, PartialEq, Eq)]
232pub struct NnrpRuntimeControl {
233    pub message_type: MessageType,
234    pub metadata: ControlRequestMetadata,
235    pub body: Vec<u8>,
236}
237
238#[derive(Debug, Clone, Copy, PartialEq, Eq)]
239pub struct NnrpSchedulingUpdate {
240    pub message_type: MessageType,
241    pub metadata: SchedulingMetadata,
242}
243
244#[derive(Debug, Clone, Copy, PartialEq, Eq)]
245pub struct NnrpPressureUpdate {
246    pub message_type: MessageType,
247    pub metadata: PressureMetadata,
248}
249
250#[derive(Debug, Clone, PartialEq, Eq)]
251pub enum NnrpServerEvent {
252    Submit(NnrpSubmit),
253    FrameCancel(NnrpCancel),
254    PartialResult {
255        metadata: PartialResultMetadata,
256        body: Vec<u8>,
257    },
258    Progress {
259        metadata: ProgressMetadata,
260        body: Vec<u8>,
261    },
262    ResultDropReason {
263        metadata: ResultDropReasonMetadata,
264        body: Vec<u8>,
265    },
266    Control(NnrpRuntimeControl),
267    Scheduling(NnrpSchedulingUpdate),
268    Supersede {
269        metadata: SupersedeMetadata,
270        body: Vec<u8>,
271    },
272    Budget(BudgetMetadata),
273    FlowUpdate(FlowUpdateMetadata),
274    Pressure(NnrpPressureUpdate),
275    Capability {
276        message_type: MessageType,
277        metadata: CapabilityMetadata,
278        body: Vec<u8>,
279    },
280    RouteHint {
281        message_type: MessageType,
282        metadata: RouteHintMetadata,
283        body: Vec<u8>,
284    },
285    TraceContext {
286        frame_id: u32,
287        metadata: TraceContextMetadata,
288        body: Vec<u8>,
289    },
290    RecoverableError {
291        metadata: RecoverableErrorMetadata,
292        body: Vec<u8>,
293    },
294    RetryAfter {
295        metadata: RetryAfterMetadata,
296        body: Vec<u8>,
297    },
298    ObjectDeclare {
299        metadata: ObjectDescriptorMetadata,
300        body: Vec<u8>,
301    },
302    ObjectRef {
303        metadata: ObjectReferenceMetadata,
304        body: Vec<u8>,
305    },
306    ObjectRelease {
307        metadata: ObjectReleaseMetadata,
308        body: Vec<u8>,
309    },
310    ObjectDelta {
311        message_type: MessageType,
312        metadata: ObjectDeltaMetadata,
313        body: Vec<u8>,
314    },
315    CacheReference {
316        metadata: CacheReferenceMetadata,
317        body: Vec<u8>,
318    },
319    CacheMiss {
320        metadata: CacheMissMetadata,
321        body: Vec<u8>,
322    },
323    CacheInvalidate(CacheInvalidateMetadata),
324    Close(SessionCloseMetadata),
325}
326
327impl NnrpServer {
328    pub async fn bind_tcp(
329        addr: impl tokio::net::ToSocketAddrs,
330        config: NnrpServerConfig,
331    ) -> Result<Self, RuntimeError> {
332        if config.transport != RuntimeTransportKind::Tcp {
333            return Err(RuntimeError::UnsupportedTransport(
334                "server config selected a non-TCP transport for bind_tcp",
335            ));
336        }
337        Self::from_listener(
338            TcpFramedListener::new(TcpListener::bind(addr).await?),
339            config,
340        )
341    }
342
343    pub async fn bind_quic(
344        _endpoint: &str,
345        config: NnrpServerConfig,
346    ) -> Result<Self, RuntimeError> {
347        if config.transport != RuntimeTransportKind::Quic {
348            return Err(RuntimeError::UnsupportedTransport(
349                "server config selected a non-QUIC transport for bind_quic",
350            ));
351        }
352        Err(RuntimeError::UnsupportedTransport(
353            "QUIC provider is not installed; use from_listener with a QUIC FramedListener",
354        ))
355    }
356
357    pub fn from_listener<L>(listener: L, config: NnrpServerConfig) -> Result<Self, RuntimeError>
358    where
359        L: FramedListener + 'static,
360    {
361        Self::from_boxed_listener(Box::new(listener), config)
362    }
363
364    pub fn from_boxed_listener(
365        listener: BoxedFramedListener,
366        config: NnrpServerConfig,
367    ) -> Result<Self, RuntimeError> {
368        if listener.transport_kind() != config.transport {
369            return Err(RuntimeError::UnsupportedTransport(
370                "server config transport does not match the provided listener slot",
371            ));
372        }
373        Ok(Self {
374            listener,
375            config,
376            sessions: Arc::new(Mutex::new(BTreeMap::new())),
377        })
378    }
379
380    pub fn local_addr(&self) -> Result<std::net::SocketAddr, RuntimeError> {
381        self.listener.local_addr()
382    }
383
384    pub fn session_count(&self) -> Result<usize, RuntimeError> {
385        Ok(self.session_registry()?.len())
386    }
387
388    pub async fn accept(&self) -> Result<NnrpServerSession, RuntimeError> {
389        let mut transport = self.listener.accept().await?;
390        let packet = loop {
391            let packet = transport.read_packet().await?;
392            if packet.header.message_type == MessageType::TransportProbe {
393                respond_to_transport_probe(&mut transport, packet).await?;
394                continue;
395            }
396            if packet.header.message_type != MessageType::SessionOpen {
397                return Err(RuntimeError::UnexpectedMessage(
398                    "server expected TRANSPORT_PROBE or SESSION_OPEN",
399                ));
400            }
401            break packet;
402        };
403
404        let open = SessionOpenMetadata::parse(&packet.metadata)?;
405        nnrp_core::validate_session_recovery_request(&open)?;
406        let ack = self.accept_ack(&open);
407        let mut ack_bytes = vec![0u8; SESSION_OPEN_ACK_METADATA_LEN];
408        ack.write(&mut ack_bytes)?;
409
410        let mut ack_header = CommonHeader::new(
411            MessageType::SessionOpenAck,
412            SESSION_OPEN_ACK_METADATA_LEN as u32,
413            0,
414        );
415        ack_header.session_id = ack.session_id;
416        transport
417            .write_packet(&RuntimePacket::new(ack_header, ack_bytes, Vec::new())?)
418            .await?;
419
420        if !matches!(
421            ack.session_status,
422            SessionStatus::Opened | SessionStatus::Resumed
423        ) {
424            return Err(RuntimeError::UnexpectedMessage(
425                "server rejected SESSION_OPEN",
426            ));
427        }
428
429        let mut lifecycle = ConnectionLifecycle::new();
430        lifecycle.apply_session_open_ack(&ack)?;
431        self.session_registry()?.insert(
432            ack.session_id,
433            RuntimeSessionRecord {
434                session_id: ack.session_id,
435                profile_id: ack.accepted_profile_id,
436                schema_id: ack.schema_id,
437                schema_version: ack.schema_version,
438                resume_enabled: ack.session_flags_ack & SESSION_ACK_FLAG_RESUME_ENABLED != 0,
439                resume_token_bytes: ack.resume_token_bytes,
440                last_operation_id: 0,
441            },
442        );
443
444        Ok(NnrpServerSession {
445            session_id: ack.session_id,
446            client_open: open,
447            transport,
448            lifecycle,
449            operations: OperationRegistry::new(),
450            frame_operations: BTreeMap::new(),
451            operation_frames: BTreeMap::new(),
452            pressure: RuntimePressureState::default(),
453            cache_objects: Vec::new(),
454            max_cache_objects: self.config.max_cache_objects,
455            sessions: Arc::clone(&self.sessions),
456            pending_close: None,
457        })
458    }
459
460    fn accept_ack(&self, open: &SessionOpenMetadata) -> SessionOpenAckMetadata {
461        let validation_error = self.config.validate_client_open(open).err();
462        let resume_attempt = open.resume_token_bytes > 0;
463        let existing_session = self
464            .session_registry()
465            .ok()
466            .and_then(|registry| registry.get(&open.requested_session_id).cloned());
467        let known_resume = resume_attempt
468            && existing_session
469                .as_ref()
470                .filter(|record| record.resume_enabled)
471                .is_some();
472        let recovery_error = if resume_attempt && !known_resume {
473            Some(SESSION_ERROR_RESUME_REJECTED)
474        } else if !resume_attempt && existing_session.is_some() {
475            Some(SESSION_ERROR_LIMIT_REACHED)
476        } else {
477            None
478        };
479        let accepted = validation_error.is_none() && recovery_error.is_none();
480        let session_id = if accepted {
481            open.requested_session_id.max(1)
482        } else {
483            0
484        };
485        let resume_enabled = open.session_flags & SESSION_FLAG_ALLOW_RESUME != 0;
486        let ack_resume_token_bytes = if accepted && resume_enabled {
487            self.config.resume_token_bytes
488        } else {
489            0
490        };
491        SessionOpenAckMetadata {
492            session_id,
493            accepted_profile_id: open.profile_id,
494            accepted_priority_class: open.priority_class,
495            session_status: if !accepted {
496                SessionStatus::Rejected
497            } else if resume_attempt {
498                SessionStatus::Resumed
499            } else {
500                SessionStatus::Opened
501            },
502            schema_id: open.schema_id,
503            schema_version: open.schema_version,
504            granted_operation_credit: self.config.granted_operation_credit,
505            max_in_flight_operations: self.config.max_in_flight_operations,
506            lease_ttl_ms: self.config.lease_ttl_ms,
507            resume_window_ms: self.config.resume_window_ms,
508            resume_token_bytes: ack_resume_token_bytes,
509            session_extension_bytes: 0,
510            server_session_tag: session_id as u64,
511            route_scope_id: 0,
512            session_error_code: validation_error
513                .or(recovery_error)
514                .unwrap_or(SESSION_ERROR_NONE),
515            session_flags_ack: if ack_resume_token_bytes > 0 {
516                SESSION_ACK_FLAG_RESUME_ENABLED
517            } else {
518                0
519            },
520        }
521    }
522
523    fn session_registry(
524        &self,
525    ) -> Result<MutexGuard<'_, BTreeMap<u32, RuntimeSessionRecord>>, RuntimeError> {
526        self.sessions
527            .lock()
528            .map_err(|_| RuntimeError::Internal("server session registry lock poisoned"))
529    }
530}
531
532async fn respond_to_transport_probe(
533    transport: &mut BoxedFramedTransport,
534    packet: RuntimePacket,
535) -> Result<(), RuntimeError> {
536    let probe = TransportProbeMetadata::parse(&packet.metadata)?;
537    if packet.body.len() != probe.probe_payload_bytes as usize {
538        return Err(nnrp_core::NnrpError::DeclaredLengthMismatch {
539            field: "transport_probe.probe_payload_bytes",
540            declared: probe.probe_payload_bytes as usize,
541            actual: packet.body.len(),
542        }
543        .into());
544    }
545    let ack = TransportProbeAckMetadata {
546        probe_id: probe.probe_id,
547        server_recv_ts_us: unix_time_us(),
548    };
549    transport
550        .write_packet(&RuntimePacket::new(
551            CommonHeader::new(MessageType::TransportProbeAck, 0, 0),
552            ack.to_bytes()?.to_vec(),
553            Vec::new(),
554        )?)
555        .await
556}
557
558fn unix_time_us() -> u64 {
559    SystemTime::now()
560        .duration_since(UNIX_EPOCH)
561        .unwrap_or_default()
562        .as_micros()
563        .try_into()
564        .unwrap_or(u64::MAX)
565}
566
567impl fmt::Debug for NnrpServer {
568    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
569        formatter
570            .debug_struct("NnrpServer")
571            .field("transport", &self.listener.transport_kind())
572            .field("config", &self.config)
573            .finish_non_exhaustive()
574    }
575}
576
577impl fmt::Debug for NnrpServerSession {
578    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
579        formatter
580            .debug_struct("NnrpServerSession")
581            .field("session_id", &self.session_id)
582            .field("client_open", &self.client_open)
583            .field("transport", &self.transport.transport_kind())
584            .field("lifecycle", &self.lifecycle)
585            .field("operations", &self.operations)
586            .field("pressure", &self.pressure)
587            .field("cache_objects", &self.cache_objects)
588            .field("max_cache_objects", &self.max_cache_objects)
589            .finish_non_exhaustive()
590    }
591}
592
593impl NnrpServerSession {
594    pub fn session_id(&self) -> u32 {
595        self.session_id
596    }
597
598    pub fn client_open(&self) -> &SessionOpenMetadata {
599        &self.client_open
600    }
601
602    pub fn lifecycle(&self) -> &ConnectionLifecycle {
603        &self.lifecycle
604    }
605
606    pub fn operations(&self) -> &OperationRegistry {
607        &self.operations
608    }
609
610    pub fn pressure_state(&self) -> RuntimePressureState {
611        self.pressure
612    }
613
614    pub fn cache_object_count(&self) -> usize {
615        self.cache_objects.len()
616    }
617
618    pub async fn receive_submit(&mut self) -> Result<NnrpSubmit, RuntimeError> {
619        let packet = self.transport.read_packet().await?;
620        self.handle_frame_submit_packet(packet)
621    }
622
623    fn handle_frame_submit_packet(
624        &mut self,
625        packet: RuntimePacket,
626    ) -> Result<NnrpSubmit, RuntimeError> {
627        if packet.header.message_type != MessageType::FrameSubmit {
628            return Err(RuntimeError::UnexpectedMessage(
629                "server expected FRAME_SUBMIT",
630            ));
631        }
632        if packet.header.session_id != self.session_id {
633            return Err(RuntimeError::UnexpectedMessage(
634                "server received submit for another session",
635            ));
636        }
637        if packet.metadata.len() != FRAME_SUBMIT_METADATA_LEN {
638            return Err(RuntimeError::UnexpectedMessage(
639                "server received malformed FRAME_SUBMIT metadata length",
640            ));
641        }
642
643        let metadata = FrameSubmitMetadata::parse(&packet.metadata)?;
644        if self.frame_operations.contains_key(&packet.header.frame_id) {
645            return Err(RuntimeError::UnexpectedMessage(
646                "server received duplicate FRAME_SUBMIT frame id",
647            ));
648        }
649        self.operations.register(OperationDescriptor::new(
650            self.session_id,
651            metadata.operation_id,
652        ))?;
653        self.frame_operations
654            .insert(packet.header.frame_id, metadata.operation_id);
655        self.operation_frames
656            .insert(metadata.operation_id, packet.header.frame_id);
657        self.update_registry_last_operation(metadata.operation_id)?;
658
659        Ok(NnrpSubmit {
660            operation_id: metadata.operation_id,
661            frame_id: packet.header.frame_id,
662            metadata,
663            body: packet.body,
664        })
665    }
666
667    pub async fn await_event(&mut self) -> Result<NnrpServerEvent, RuntimeError> {
668        let packet = self.transport.read_packet().await?;
669        match packet.header.message_type {
670            MessageType::FrameSubmit => self
671                .handle_frame_submit_packet(packet)
672                .map(NnrpServerEvent::Submit),
673            MessageType::FrameCancel => {
674                self.require_session_packet(&packet, "server received cancel for another session")?;
675                if !packet.metadata.is_empty() || !packet.body.is_empty() {
676                    return Err(RuntimeError::UnexpectedMessage(
677                        "server received malformed FRAME_CANCEL lengths",
678                    ));
679                }
680                let operation_id = self.operation_id_for_frame(packet.header.frame_id)?;
681                self.operations.cancel(OperationCancelRequest {
682                    session_id: self.session_id,
683                    operation_id,
684                    cancel_scope: nnrp_core::CancelScope::Operation,
685                })?;
686                Ok(NnrpServerEvent::FrameCancel(NnrpCancel {
687                    frame_id: packet.header.frame_id,
688                }))
689            }
690            MessageType::PartialResult => {
691                self.require_session_packet(
692                    &packet,
693                    "server received partial result for another session",
694                )?;
695                if packet.metadata.len() != PARTIAL_RESULT_METADATA_LEN {
696                    return Err(RuntimeError::UnexpectedMessage(
697                        "server received malformed PARTIAL_RESULT metadata length",
698                    ));
699                }
700                let metadata = PartialResultMetadata::parse(&packet.metadata)?;
701                validate_partial_result_semantics(&metadata)?;
702                require_body_len(
703                    packet.body.len(),
704                    metadata.body_bytes as usize,
705                    "server received PARTIAL_RESULT body length mismatch",
706                )?;
707                self.require_operation_frame(metadata.operation_id, packet.header.frame_id)?;
708                Ok(NnrpServerEvent::PartialResult {
709                    metadata,
710                    body: packet.body,
711                })
712            }
713            MessageType::Progress => {
714                self.require_session_packet(
715                    &packet,
716                    "server received progress for another session",
717                )?;
718                if packet.metadata.len() != PROGRESS_METADATA_LEN {
719                    return Err(RuntimeError::UnexpectedMessage(
720                        "server received malformed PROGRESS metadata length",
721                    ));
722                }
723                let metadata = ProgressMetadata::parse(&packet.metadata)?;
724                validate_progress_semantics(&metadata)?;
725                require_body_len(
726                    packet.body.len(),
727                    metadata.body_bytes as usize,
728                    "server received PROGRESS body length mismatch",
729                )?;
730                self.require_operation_frame(metadata.operation_id, packet.header.frame_id)?;
731                Ok(NnrpServerEvent::Progress {
732                    metadata,
733                    body: packet.body,
734                })
735            }
736            MessageType::ResultDropReason => {
737                self.require_session_packet(
738                    &packet,
739                    "server received drop reason for another session",
740                )?;
741                if packet.metadata.len() != RESULT_DROP_REASON_METADATA_LEN {
742                    return Err(RuntimeError::UnexpectedMessage(
743                        "server received malformed RESULT_DROP_REASON metadata length",
744                    ));
745                }
746                let metadata = ResultDropReasonMetadata::parse(&packet.metadata)?;
747                validate_result_drop_reason_semantics(&metadata)?;
748                require_body_len(
749                    packet.body.len(),
750                    metadata.diagnostic_bytes as usize,
751                    "server received RESULT_DROP_REASON body length mismatch",
752                )?;
753                self.require_operation_frame(metadata.operation_id, packet.header.frame_id)?;
754                Ok(NnrpServerEvent::ResultDropReason {
755                    metadata,
756                    body: packet.body,
757                })
758            }
759            MessageType::Cancel | MessageType::Abort => {
760                self.require_session_packet(
761                    &packet,
762                    "server received control for another session",
763                )?;
764                if packet.metadata.len() != CONTROL_REQUEST_METADATA_LEN {
765                    return Err(RuntimeError::UnexpectedMessage(
766                        "server received malformed runtime control lengths",
767                    ));
768                }
769                let metadata = ControlRequestMetadata::parse(&packet.metadata)?;
770                validate_control_request_semantics(packet.header.message_type, &metadata)?;
771                require_body_len(
772                    packet.body.len(),
773                    metadata.diagnostic_bytes as usize,
774                    "server received runtime control diagnostic body length mismatch",
775                )?;
776                self.require_operation_frame(metadata.operation_id, packet.header.frame_id)?;
777                match packet.header.message_type {
778                    MessageType::Cancel => {
779                        self.operations.cancel(OperationCancelRequest {
780                            session_id: self.session_id,
781                            operation_id: metadata.operation_id,
782                            cancel_scope: nnrp_core::CancelScope::Operation,
783                        })?;
784                    }
785                    MessageType::Abort => self.operations.abort(metadata.operation_id)?,
786                    _ => unreachable!("runtime control message type was matched earlier"),
787                }
788                Ok(NnrpServerEvent::Control(NnrpRuntimeControl {
789                    message_type: packet.header.message_type,
790                    metadata,
791                    body: packet.body,
792                }))
793            }
794            MessageType::PriorityUpdate | MessageType::Deadline | MessageType::ExpireAt => {
795                self.require_session_packet(
796                    &packet,
797                    "server received scheduling update for another session",
798                )?;
799                if packet.metadata.len() != SCHEDULING_METADATA_LEN || !packet.body.is_empty() {
800                    return Err(RuntimeError::UnexpectedMessage(
801                        "server received malformed scheduling metadata length",
802                    ));
803                }
804                let metadata = SchedulingMetadata::parse(&packet.metadata)?;
805                validate_scheduling_semantics(packet.header.message_type, &metadata)?;
806                self.require_operation_frame(metadata.operation_id, packet.header.frame_id)?;
807                self.operations.apply_scheduling_update(
808                    self.session_id,
809                    packet.header.message_type,
810                    metadata,
811                )?;
812                Ok(NnrpServerEvent::Scheduling(NnrpSchedulingUpdate {
813                    message_type: packet.header.message_type,
814                    metadata,
815                }))
816            }
817            MessageType::Supersede => {
818                self.require_session_packet(
819                    &packet,
820                    "server received supersede for another session",
821                )?;
822                if packet.metadata.len() != SUPERSEDE_METADATA_LEN {
823                    return Err(RuntimeError::UnexpectedMessage(
824                        "server received malformed SUPERSEDE metadata length",
825                    ));
826                }
827                let metadata = SupersedeMetadata::parse(&packet.metadata)?;
828                require_body_len(
829                    packet.body.len(),
830                    metadata.diagnostic_bytes as usize,
831                    "server received SUPERSEDE diagnostic body length mismatch",
832                )?;
833                self.require_operation_frame(metadata.old_operation_id, packet.header.frame_id)?;
834                Ok(NnrpServerEvent::Supersede {
835                    metadata,
836                    body: packet.body,
837                })
838            }
839            MessageType::BudgetUpdate => {
840                self.require_session_packet(
841                    &packet,
842                    "server received budget update for another session",
843                )?;
844                if packet.metadata.len() != BUDGET_METADATA_LEN || !packet.body.is_empty() {
845                    return Err(RuntimeError::UnexpectedMessage(
846                        "server received malformed BUDGET_UPDATE lengths",
847                    ));
848                }
849                let metadata = BudgetMetadata::parse(&packet.metadata)?;
850                self.require_operation_frame(metadata.operation_id, packet.header.frame_id)?;
851                Ok(NnrpServerEvent::Budget(metadata))
852            }
853            MessageType::FlowUpdate => {
854                if packet.metadata.len() != FLOW_UPDATE_METADATA_LEN || !packet.body.is_empty() {
855                    return Err(RuntimeError::UnexpectedMessage(
856                        "server received malformed FLOW_UPDATE lengths",
857                    ));
858                }
859                let metadata = FlowUpdateMetadata::parse(&packet.metadata)?;
860                self.lifecycle
861                    .validate_flow_update(&packet.header, &metadata)?;
862                Ok(NnrpServerEvent::FlowUpdate(metadata))
863            }
864            MessageType::Backpressure | MessageType::CreditUpdate => {
865                self.require_optional_session_packet(
866                    &packet,
867                    "server received pressure update for another session",
868                )?;
869                if packet.metadata.len() != PRESSURE_METADATA_LEN || !packet.body.is_empty() {
870                    return Err(RuntimeError::UnexpectedMessage(
871                        "server received malformed pressure metadata length",
872                    ));
873                }
874                let metadata = PressureMetadata::parse(&packet.metadata)?;
875                validate_pressure_semantics(packet.header.message_type, &metadata)?;
876                self.pressure
877                    .apply_inbound(packet.header.message_type, metadata)?;
878                Ok(NnrpServerEvent::Pressure(NnrpPressureUpdate {
879                    message_type: packet.header.message_type,
880                    metadata,
881                }))
882            }
883            MessageType::CapabilityNegotiation | MessageType::DegradeProfile => {
884                self.require_optional_session_packet(
885                    &packet,
886                    "server received capability update for another session",
887                )?;
888                if packet.metadata.len() != CAPABILITY_METADATA_LEN {
889                    return Err(RuntimeError::UnexpectedMessage(
890                        "server received malformed capability metadata length",
891                    ));
892                }
893                let metadata = CapabilityMetadata::parse(&packet.metadata)?;
894                require_body_len(
895                    packet.body.len(),
896                    metadata.body_bytes as usize,
897                    "server received capability body length mismatch",
898                )?;
899                Ok(NnrpServerEvent::Capability {
900                    message_type: packet.header.message_type,
901                    metadata,
902                    body: packet.body,
903                })
904            }
905            MessageType::RouteHint | MessageType::ExecutionHint => {
906                self.require_optional_session_packet(
907                    &packet,
908                    "server received route hint for another session",
909                )?;
910                if packet.metadata.len() != ROUTE_HINT_METADATA_LEN {
911                    return Err(RuntimeError::UnexpectedMessage(
912                        "server received malformed route hint metadata length",
913                    ));
914                }
915                let metadata = RouteHintMetadata::parse(&packet.metadata)?;
916                require_body_len(
917                    packet.body.len(),
918                    metadata.body_bytes as usize,
919                    "server received route hint body length mismatch",
920                )?;
921                self.require_operation_frame(metadata.operation_id, packet.header.frame_id)?;
922                Ok(NnrpServerEvent::RouteHint {
923                    message_type: packet.header.message_type,
924                    metadata,
925                    body: packet.body,
926                })
927            }
928            MessageType::TraceContext => {
929                self.require_optional_session_packet(
930                    &packet,
931                    "server received trace context for another session",
932                )?;
933                if packet.metadata.len() != TRACE_CONTEXT_METADATA_LEN {
934                    return Err(RuntimeError::UnexpectedMessage(
935                        "server received malformed TRACE_CONTEXT metadata length",
936                    ));
937                }
938                let metadata = TraceContextMetadata::parse(&packet.metadata)?;
939                validate_trace_context_semantics(&metadata)?;
940                require_body_len(
941                    packet.body.len(),
942                    metadata.body_bytes as usize,
943                    "server received TRACE_CONTEXT body length mismatch",
944                )?;
945                Ok(NnrpServerEvent::TraceContext {
946                    frame_id: packet.header.frame_id,
947                    metadata,
948                    body: packet.body,
949                })
950            }
951            MessageType::ErrorRecoverable => {
952                self.require_optional_session_packet(
953                    &packet,
954                    "server received recoverable error for another session",
955                )?;
956                if packet.metadata.len() != RECOVERABLE_ERROR_METADATA_LEN {
957                    return Err(RuntimeError::UnexpectedMessage(
958                        "server received malformed ERROR_RECOVERABLE metadata length",
959                    ));
960                }
961                let metadata = RecoverableErrorMetadata::parse(&packet.metadata)?;
962                require_body_len(
963                    packet.body.len(),
964                    metadata.diagnostic_bytes as usize,
965                    "server received ERROR_RECOVERABLE diagnostic body length mismatch",
966                )?;
967                Ok(NnrpServerEvent::RecoverableError {
968                    metadata,
969                    body: packet.body,
970                })
971            }
972            MessageType::RetryAfter => {
973                self.require_optional_session_packet(
974                    &packet,
975                    "server received retry-after for another session",
976                )?;
977                if packet.metadata.len() != RETRY_AFTER_METADATA_LEN {
978                    return Err(RuntimeError::UnexpectedMessage(
979                        "server received malformed RETRY_AFTER metadata length",
980                    ));
981                }
982                let metadata = RetryAfterMetadata::parse(&packet.metadata)?;
983                require_body_len(
984                    packet.body.len(),
985                    metadata.diagnostic_bytes as usize,
986                    "server received RETRY_AFTER diagnostic body length mismatch",
987                )?;
988                Ok(NnrpServerEvent::RetryAfter {
989                    metadata,
990                    body: packet.body,
991                })
992            }
993            MessageType::ObjectDeclare => {
994                self.require_session_packet(
995                    &packet,
996                    "server received object declaration for another session",
997                )?;
998                if packet.metadata.len() != OBJECT_DESCRIPTOR_METADATA_LEN {
999                    return Err(RuntimeError::UnexpectedMessage(
1000                        "server received malformed OBJECT_DECLARE metadata length",
1001                    ));
1002                }
1003                let metadata = ObjectDescriptorMetadata::parse(&packet.metadata)?;
1004                require_body_len(
1005                    packet.body.len(),
1006                    metadata.metadata_bytes as usize,
1007                    "server received OBJECT_DECLARE body length mismatch",
1008                )?;
1009                Ok(NnrpServerEvent::ObjectDeclare {
1010                    metadata,
1011                    body: packet.body,
1012                })
1013            }
1014            MessageType::ObjectRef => {
1015                self.require_session_packet(
1016                    &packet,
1017                    "server received object reference for another session",
1018                )?;
1019                if packet.metadata.len() != OBJECT_REFERENCE_METADATA_LEN {
1020                    return Err(RuntimeError::UnexpectedMessage(
1021                        "server received malformed OBJECT_REF metadata length",
1022                    ));
1023                }
1024                let metadata = ObjectReferenceMetadata::parse(&packet.metadata)?;
1025                require_body_len(
1026                    packet.body.len(),
1027                    metadata.metadata_bytes as usize,
1028                    "server received OBJECT_REF body length mismatch",
1029                )?;
1030                self.require_operation_frame(metadata.operation_id, packet.header.frame_id)?;
1031                Ok(NnrpServerEvent::ObjectRef {
1032                    metadata,
1033                    body: packet.body,
1034                })
1035            }
1036            MessageType::ObjectRelease => {
1037                self.require_session_packet(
1038                    &packet,
1039                    "server received object release for another session",
1040                )?;
1041                if packet.metadata.len() != OBJECT_RELEASE_METADATA_LEN {
1042                    return Err(RuntimeError::UnexpectedMessage(
1043                        "server received malformed OBJECT_RELEASE metadata length",
1044                    ));
1045                }
1046                let metadata = ObjectReleaseMetadata::parse(&packet.metadata)?;
1047                require_body_len(
1048                    packet.body.len(),
1049                    metadata.diagnostic_bytes as usize,
1050                    "server received OBJECT_RELEASE body length mismatch",
1051                )?;
1052                self.require_operation_frame(metadata.operation_id, packet.header.frame_id)?;
1053                Ok(NnrpServerEvent::ObjectRelease {
1054                    metadata,
1055                    body: packet.body,
1056                })
1057            }
1058            MessageType::ObjectPatch | MessageType::ObjectDelta => {
1059                self.require_session_packet(
1060                    &packet,
1061                    "server received object delta for another session",
1062                )?;
1063                if packet.metadata.len() != OBJECT_DELTA_METADATA_LEN {
1064                    return Err(RuntimeError::UnexpectedMessage(
1065                        "server received malformed object delta metadata length",
1066                    ));
1067                }
1068                let metadata = ObjectDeltaMetadata::parse(&packet.metadata)?;
1069                require_body_len(
1070                    packet.body.len(),
1071                    metadata.metadata_bytes.saturating_add(metadata.delta_bytes) as usize,
1072                    "server received object delta body length mismatch",
1073                )?;
1074                Ok(NnrpServerEvent::ObjectDelta {
1075                    message_type: packet.header.message_type,
1076                    metadata,
1077                    body: packet.body,
1078                })
1079            }
1080            MessageType::CacheReference => {
1081                self.require_session_packet(
1082                    &packet,
1083                    "server received cache reference for another session",
1084                )?;
1085                if packet.metadata.len() != CACHE_REFERENCE_METADATA_LEN {
1086                    return Err(RuntimeError::UnexpectedMessage(
1087                        "server received malformed CACHE_REFERENCE metadata length",
1088                    ));
1089                }
1090                let metadata = CacheReferenceMetadata::parse(&packet.metadata)?;
1091                require_body_len(
1092                    packet.body.len(),
1093                    metadata.metadata_bytes as usize,
1094                    "server received CACHE_REFERENCE body length mismatch",
1095                )?;
1096                Ok(NnrpServerEvent::CacheReference {
1097                    metadata,
1098                    body: packet.body,
1099                })
1100            }
1101            MessageType::CacheMiss => {
1102                self.require_session_packet(
1103                    &packet,
1104                    "server received cache miss for another session",
1105                )?;
1106                if packet.metadata.len() != CACHE_MISS_METADATA_LEN {
1107                    return Err(RuntimeError::UnexpectedMessage(
1108                        "server received malformed CACHE_MISS metadata length",
1109                    ));
1110                }
1111                let metadata = CacheMissMetadata::parse(&packet.metadata)?;
1112                require_body_len(
1113                    packet.body.len(),
1114                    metadata.diagnostic_bytes as usize,
1115                    "server received CACHE_MISS diagnostic body length mismatch",
1116                )?;
1117                Ok(NnrpServerEvent::CacheMiss {
1118                    metadata,
1119                    body: packet.body,
1120                })
1121            }
1122            MessageType::CacheInvalidate => {
1123                self.require_session_packet(
1124                    &packet,
1125                    "server received cache invalidate for another session",
1126                )?;
1127                if packet.metadata.len() != CACHE_INVALIDATE_METADATA_LEN || !packet.body.is_empty()
1128                {
1129                    return Err(RuntimeError::UnexpectedMessage(
1130                        "server received malformed CACHE_INVALIDATE lengths",
1131                    ));
1132                }
1133                Ok(NnrpServerEvent::CacheInvalidate(
1134                    CacheInvalidateMetadata::parse(&packet.metadata)?,
1135                ))
1136            }
1137            MessageType::SessionClose => {
1138                self.require_session_packet(&packet, "server received close for another session")?;
1139                let metadata = SessionCloseMetadata::parse(&packet.metadata)?;
1140                self.lifecycle
1141                    .begin_session_close(&packet.header, &metadata)?;
1142                self.pending_close = Some(metadata);
1143                Ok(NnrpServerEvent::Close(metadata))
1144            }
1145            _ => Err(RuntimeError::UnexpectedMessage(
1146                "server expected a submit, control, object, cache, or close event",
1147            )),
1148        }
1149    }
1150
1151    pub async fn send_result(
1152        &mut self,
1153        frame_id: u32,
1154        metadata: ResultPushMetadata,
1155        body: Vec<u8>,
1156    ) -> Result<(), RuntimeError> {
1157        let operation_id = self.operation_id_for_frame(frame_id)?;
1158        if let Some(schedule) = self
1159            .operations
1160            .expire_if_stale(operation_id, current_unix_ms())?
1161        {
1162            if schedule.flags & SCHEDULING_FLAG_EMIT_DROP_REASON != 0 {
1163                self.send_result_drop_reason(ResultDropReasonMetadata {
1164                    operation_id,
1165                    result_sequence: schedule.update_sequence,
1166                    drop_reason_code: RESULT_DROP_REASON_DEADLINE_EXPIRED,
1167                    source_role: RuntimeRole::Server as u8,
1168                    flags: 0,
1169                    diagnostic_bytes: 0,
1170                })
1171                .await?;
1172            }
1173            return Err(nnrp_core::NnrpError::InvalidOperationTransition {
1174                from: nnrp_core::OperationState::Superseded,
1175                to: nnrp_core::OperationState::Completed,
1176            }
1177            .into());
1178        }
1179        self.operations.complete(operation_id)?;
1180        let mut header = CommonHeader::new(
1181            MessageType::ResultPush,
1182            RESULT_PUSH_METADATA_LEN as u32,
1183            body.len() as u32,
1184        );
1185        header.session_id = self.session_id;
1186        header.frame_id = frame_id;
1187        self.transport
1188            .write_packet(&RuntimePacket::new(
1189                header,
1190                metadata.to_bytes()?.to_vec(),
1191                body,
1192            )?)
1193            .await?;
1194        Ok(())
1195    }
1196
1197    pub async fn send_result_drop(&mut self, frame_id: u32) -> Result<(), RuntimeError> {
1198        self.operation_id_for_frame(frame_id)?;
1199        let mut header = CommonHeader::new(MessageType::ResultDrop, 0, 0);
1200        header.session_id = self.session_id;
1201        header.frame_id = frame_id;
1202        validate_result_drop_header(&header)?;
1203        self.transport
1204            .write_packet(&RuntimePacket::new(header, Vec::new(), Vec::new())?)
1205            .await?;
1206        Ok(())
1207    }
1208
1209    pub async fn send_partial_result(
1210        &mut self,
1211        metadata: PartialResultMetadata,
1212        body: Vec<u8>,
1213    ) -> Result<(), RuntimeError> {
1214        validate_partial_result_semantics(&metadata)?;
1215        if metadata.body_bytes as usize != body.len() {
1216            return Err(RuntimeError::UnexpectedMessage(
1217                "server PARTIAL_RESULT body length mismatch",
1218            ));
1219        }
1220        let mut header = CommonHeader::new(
1221            MessageType::PartialResult,
1222            PARTIAL_RESULT_METADATA_LEN as u32,
1223            body.len() as u32,
1224        );
1225        header.session_id = self.session_id;
1226        header.frame_id = self.correlated_frame_id(metadata.operation_id)?;
1227        self.transport
1228            .write_packet(&RuntimePacket::new(
1229                header,
1230                metadata.to_bytes()?.to_vec(),
1231                body,
1232            )?)
1233            .await
1234    }
1235
1236    pub async fn send_progress(
1237        &mut self,
1238        metadata: ProgressMetadata,
1239        body: Vec<u8>,
1240    ) -> Result<(), RuntimeError> {
1241        validate_progress_semantics(&metadata)?;
1242        if metadata.body_bytes as usize != body.len() {
1243            return Err(RuntimeError::UnexpectedMessage(
1244                "server PROGRESS body length mismatch",
1245            ));
1246        }
1247        let mut header = CommonHeader::new(
1248            MessageType::Progress,
1249            PROGRESS_METADATA_LEN as u32,
1250            body.len() as u32,
1251        );
1252        header.session_id = self.session_id;
1253        header.frame_id = self.correlated_frame_id(metadata.operation_id)?;
1254        self.transport
1255            .write_packet(&RuntimePacket::new(
1256                header,
1257                metadata.to_bytes()?.to_vec(),
1258                body,
1259            )?)
1260            .await
1261    }
1262
1263    pub async fn send_result_drop_reason(
1264        &mut self,
1265        metadata: ResultDropReasonMetadata,
1266    ) -> Result<(), RuntimeError> {
1267        self.send_result_drop_reason_with_diagnostics(metadata, Vec::new())
1268            .await
1269    }
1270
1271    pub async fn send_result_drop_reason_with_diagnostics(
1272        &mut self,
1273        metadata: ResultDropReasonMetadata,
1274        diagnostics: Vec<u8>,
1275    ) -> Result<(), RuntimeError> {
1276        validate_result_drop_reason_semantics(&metadata)?;
1277        require_body_len(
1278            diagnostics.len(),
1279            metadata.diagnostic_bytes as usize,
1280            "server RESULT_DROP_REASON diagnostic body length mismatch",
1281        )?;
1282        let mut header = CommonHeader::new(
1283            MessageType::ResultDropReason,
1284            RESULT_DROP_REASON_METADATA_LEN as u32,
1285            diagnostics.len() as u32,
1286        );
1287        header.session_id = self.session_id;
1288        header.frame_id = self.correlated_frame_id(metadata.operation_id)?;
1289        self.transport
1290            .write_packet(&RuntimePacket::new(
1291                header,
1292                metadata.to_bytes()?.to_vec(),
1293                diagnostics,
1294            )?)
1295            .await
1296    }
1297
1298    pub async fn send_control_request(
1299        &mut self,
1300        message_type: MessageType,
1301        metadata: ControlRequestMetadata,
1302        diagnostics: Vec<u8>,
1303    ) -> Result<(), RuntimeError> {
1304        validate_control_request_semantics(message_type, &metadata)?;
1305        require_body_len(
1306            diagnostics.len(),
1307            metadata.diagnostic_bytes as usize,
1308            "server runtime control diagnostic body length mismatch",
1309        )?;
1310        let frame_id = self.correlated_frame_id(metadata.operation_id)?;
1311        self.write_runtime_packet(
1312            message_type,
1313            frame_id,
1314            metadata.to_bytes()?.to_vec(),
1315            diagnostics,
1316        )
1317        .await
1318    }
1319
1320    pub async fn send_scheduling_update(
1321        &mut self,
1322        message_type: MessageType,
1323        metadata: SchedulingMetadata,
1324    ) -> Result<(), RuntimeError> {
1325        validate_scheduling_semantics(message_type, &metadata)?;
1326        let frame_id = self.correlated_frame_id(metadata.operation_id)?;
1327        self.write_runtime_packet(
1328            message_type,
1329            frame_id,
1330            metadata.to_bytes()?.to_vec(),
1331            Vec::new(),
1332        )
1333        .await
1334    }
1335
1336    pub async fn supersede_operation(
1337        &mut self,
1338        metadata: SupersedeMetadata,
1339        diagnostics: Vec<u8>,
1340    ) -> Result<(), RuntimeError> {
1341        require_body_len(
1342            diagnostics.len(),
1343            metadata.diagnostic_bytes as usize,
1344            "server SUPERSEDE diagnostic body length mismatch",
1345        )?;
1346        let frame_id = self.correlated_frame_id(metadata.old_operation_id)?;
1347        self.write_runtime_packet(
1348            MessageType::Supersede,
1349            frame_id,
1350            metadata.to_bytes()?.to_vec(),
1351            diagnostics,
1352        )
1353        .await
1354    }
1355
1356    pub async fn update_budget(&mut self, metadata: BudgetMetadata) -> Result<(), RuntimeError> {
1357        let frame_id = self.correlated_frame_id(metadata.operation_id)?;
1358        self.write_runtime_packet(
1359            MessageType::BudgetUpdate,
1360            frame_id,
1361            metadata.to_bytes()?.to_vec(),
1362            Vec::new(),
1363        )
1364        .await
1365    }
1366
1367    pub async fn send_capability(
1368        &mut self,
1369        message_type: MessageType,
1370        metadata: CapabilityMetadata,
1371        body: Vec<u8>,
1372    ) -> Result<(), RuntimeError> {
1373        if !matches!(
1374            message_type,
1375            MessageType::CapabilityNegotiation | MessageType::DegradeProfile
1376        ) {
1377            return Err(RuntimeError::UnexpectedMessage(
1378                "server capability send requires CAPABILITY_NEGOTIATION or DEGRADE_PROFILE",
1379            ));
1380        }
1381        require_body_len(
1382            body.len(),
1383            metadata.body_bytes as usize,
1384            "server capability body length mismatch",
1385        )?;
1386        let mut header = CommonHeader::new(
1387            message_type,
1388            CAPABILITY_METADATA_LEN as u32,
1389            body.len() as u32,
1390        );
1391        header.session_id = self.session_id;
1392        self.transport
1393            .write_packet(&RuntimePacket::new(
1394                header,
1395                metadata.to_bytes()?.to_vec(),
1396                body,
1397            )?)
1398            .await
1399    }
1400
1401    pub async fn send_route_hint(
1402        &mut self,
1403        message_type: MessageType,
1404        metadata: RouteHintMetadata,
1405        body: Vec<u8>,
1406    ) -> Result<(), RuntimeError> {
1407        if !matches!(
1408            message_type,
1409            MessageType::RouteHint | MessageType::ExecutionHint
1410        ) {
1411            return Err(RuntimeError::UnexpectedMessage(
1412                "server route hint send requires ROUTE_HINT or EXECUTION_HINT",
1413            ));
1414        }
1415        require_body_len(
1416            body.len(),
1417            metadata.body_bytes as usize,
1418            "server route hint body length mismatch",
1419        )?;
1420        let mut header = CommonHeader::new(
1421            message_type,
1422            ROUTE_HINT_METADATA_LEN as u32,
1423            body.len() as u32,
1424        );
1425        header.session_id = self.session_id;
1426        header.frame_id = self.correlated_frame_id(metadata.operation_id)?;
1427        self.transport
1428            .write_packet(&RuntimePacket::new(
1429                header,
1430                metadata.to_bytes()?.to_vec(),
1431                body,
1432            )?)
1433            .await
1434    }
1435
1436    pub async fn send_object_declare(
1437        &mut self,
1438        metadata: ObjectDescriptorMetadata,
1439        body: Vec<u8>,
1440    ) -> Result<(), RuntimeError> {
1441        require_body_len(
1442            body.len(),
1443            metadata.metadata_bytes as usize,
1444            "server OBJECT_DECLARE body length mismatch",
1445        )?;
1446        let mut header = CommonHeader::new(
1447            MessageType::ObjectDeclare,
1448            OBJECT_DESCRIPTOR_METADATA_LEN as u32,
1449            body.len() as u32,
1450        );
1451        header.session_id = self.session_id;
1452        self.transport
1453            .write_packet(&RuntimePacket::new(
1454                header,
1455                metadata.to_bytes()?.to_vec(),
1456                body,
1457            )?)
1458            .await
1459    }
1460
1461    pub async fn send_object_ref(
1462        &mut self,
1463        metadata: ObjectReferenceMetadata,
1464        body: Vec<u8>,
1465    ) -> Result<(), RuntimeError> {
1466        require_body_len(
1467            body.len(),
1468            metadata.metadata_bytes as usize,
1469            "server OBJECT_REF body length mismatch",
1470        )?;
1471        let mut header = CommonHeader::new(
1472            MessageType::ObjectRef,
1473            OBJECT_REFERENCE_METADATA_LEN as u32,
1474            body.len() as u32,
1475        );
1476        header.session_id = self.session_id;
1477        header.frame_id = self.correlated_frame_id(metadata.operation_id)?;
1478        self.transport
1479            .write_packet(&RuntimePacket::new(
1480                header,
1481                metadata.to_bytes()?.to_vec(),
1482                body,
1483            )?)
1484            .await
1485    }
1486
1487    pub async fn send_object_release(
1488        &mut self,
1489        metadata: ObjectReleaseMetadata,
1490        body: Vec<u8>,
1491    ) -> Result<(), RuntimeError> {
1492        require_body_len(
1493            body.len(),
1494            metadata.diagnostic_bytes as usize,
1495            "server OBJECT_RELEASE body length mismatch",
1496        )?;
1497        let mut header = CommonHeader::new(
1498            MessageType::ObjectRelease,
1499            OBJECT_RELEASE_METADATA_LEN as u32,
1500            body.len() as u32,
1501        );
1502        header.session_id = self.session_id;
1503        header.frame_id = self.correlated_frame_id(metadata.operation_id)?;
1504        self.transport
1505            .write_packet(&RuntimePacket::new(
1506                header,
1507                metadata.to_bytes()?.to_vec(),
1508                body,
1509            )?)
1510            .await
1511    }
1512
1513    pub async fn send_object_delta(
1514        &mut self,
1515        message_type: MessageType,
1516        metadata: ObjectDeltaMetadata,
1517        body: Vec<u8>,
1518    ) -> Result<(), RuntimeError> {
1519        if !matches!(
1520            message_type,
1521            MessageType::ObjectPatch | MessageType::ObjectDelta
1522        ) {
1523            return Err(RuntimeError::UnexpectedMessage(
1524                "server object delta send requires OBJECT_PATCH or OBJECT_DELTA",
1525            ));
1526        }
1527        let expected_body_len =
1528            metadata.metadata_bytes.saturating_add(metadata.delta_bytes) as usize;
1529        require_body_len(
1530            body.len(),
1531            expected_body_len,
1532            "server object delta body length mismatch",
1533        )?;
1534        let mut header = CommonHeader::new(
1535            message_type,
1536            OBJECT_DELTA_METADATA_LEN as u32,
1537            body.len() as u32,
1538        );
1539        header.session_id = self.session_id;
1540        self.transport
1541            .write_packet(&RuntimePacket::new(
1542                header,
1543                metadata.to_bytes()?.to_vec(),
1544                body,
1545            )?)
1546            .await
1547    }
1548
1549    pub async fn send_cache_reference(
1550        &mut self,
1551        metadata: CacheReferenceMetadata,
1552        body: Vec<u8>,
1553    ) -> Result<(), RuntimeError> {
1554        require_body_len(
1555            body.len(),
1556            metadata.metadata_bytes as usize,
1557            "server CACHE_REFERENCE body length mismatch",
1558        )?;
1559        let mut header = CommonHeader::new(
1560            MessageType::CacheReference,
1561            CACHE_REFERENCE_METADATA_LEN as u32,
1562            body.len() as u32,
1563        );
1564        header.session_id = self.session_id;
1565        self.transport
1566            .write_packet(&RuntimePacket::new(
1567                header,
1568                metadata.to_bytes()?.to_vec(),
1569                body,
1570            )?)
1571            .await
1572    }
1573
1574    pub async fn send_cache_miss(
1575        &mut self,
1576        metadata: CacheMissMetadata,
1577        body: Vec<u8>,
1578    ) -> Result<(), RuntimeError> {
1579        require_body_len(
1580            body.len(),
1581            metadata.diagnostic_bytes as usize,
1582            "server CACHE_MISS body length mismatch",
1583        )?;
1584        let mut header = CommonHeader::new(
1585            MessageType::CacheMiss,
1586            CACHE_MISS_METADATA_LEN as u32,
1587            body.len() as u32,
1588        );
1589        header.session_id = self.session_id;
1590        self.transport
1591            .write_packet(&RuntimePacket::new(
1592                header,
1593                metadata.to_bytes()?.to_vec(),
1594                body,
1595            )?)
1596            .await
1597    }
1598
1599    pub async fn send_cache_invalidate(
1600        &mut self,
1601        metadata: CacheInvalidateMetadata,
1602    ) -> Result<(), RuntimeError> {
1603        let mut header = CommonHeader::new(
1604            MessageType::CacheInvalidate,
1605            CACHE_INVALIDATE_METADATA_LEN as u32,
1606            0,
1607        );
1608        header.session_id = self.session_id;
1609        self.transport
1610            .write_packet(&RuntimePacket::new(
1611                header,
1612                metadata.to_bytes()?.to_vec(),
1613                Vec::new(),
1614            )?)
1615            .await
1616    }
1617
1618    pub async fn receive_cancel(&mut self) -> Result<NnrpCancel, RuntimeError> {
1619        let packet = self.transport.read_packet().await?;
1620        if packet.header.message_type != MessageType::FrameCancel {
1621            return Err(RuntimeError::UnexpectedMessage(
1622                "server expected FRAME_CANCEL",
1623            ));
1624        }
1625        self.require_session_packet(&packet, "server received cancel for another session")?;
1626        if packet.header.meta_len != 0 || packet.header.body_len != 0 {
1627            return Err(RuntimeError::UnexpectedMessage(
1628                "server received malformed FRAME_CANCEL lengths",
1629            ));
1630        }
1631        let operation_id = self.operation_id_for_frame(packet.header.frame_id)?;
1632        self.operations.cancel(OperationCancelRequest {
1633            session_id: self.session_id,
1634            operation_id,
1635            cancel_scope: nnrp_core::CancelScope::Operation,
1636        })?;
1637        Ok(NnrpCancel {
1638            frame_id: packet.header.frame_id,
1639        })
1640    }
1641
1642    pub async fn receive_runtime_control(&mut self) -> Result<NnrpRuntimeControl, RuntimeError> {
1643        let packet = self.transport.read_packet().await?;
1644        if !matches!(
1645            packet.header.message_type,
1646            MessageType::Cancel | MessageType::Abort
1647        ) {
1648            return Err(RuntimeError::UnexpectedMessage(
1649                "server expected CANCEL or ABORT",
1650            ));
1651        }
1652        self.require_session_packet(&packet, "server received control for another session")?;
1653        if packet.metadata.len() != CONTROL_REQUEST_METADATA_LEN {
1654            return Err(RuntimeError::UnexpectedMessage(
1655                "server received malformed runtime control lengths",
1656            ));
1657        }
1658
1659        let metadata = ControlRequestMetadata::parse(&packet.metadata)?;
1660        validate_control_request_semantics(packet.header.message_type, &metadata)?;
1661        require_body_len(
1662            packet.body.len(),
1663            metadata.diagnostic_bytes as usize,
1664            "server received runtime control diagnostic body length mismatch",
1665        )?;
1666        self.require_operation_frame(metadata.operation_id, packet.header.frame_id)?;
1667        match packet.header.message_type {
1668            MessageType::Cancel => {
1669                self.operations.cancel(OperationCancelRequest {
1670                    session_id: self.session_id,
1671                    operation_id: metadata.operation_id,
1672                    cancel_scope: nnrp_core::CancelScope::Operation,
1673                })?;
1674            }
1675            MessageType::Abort => {
1676                self.operations.abort(metadata.operation_id)?;
1677            }
1678            _ => unreachable!("runtime control message type was validated earlier"),
1679        }
1680        Ok(NnrpRuntimeControl {
1681            message_type: packet.header.message_type,
1682            metadata,
1683            body: packet.body,
1684        })
1685    }
1686
1687    pub async fn receive_scheduling_update(
1688        &mut self,
1689    ) -> Result<NnrpSchedulingUpdate, RuntimeError> {
1690        let packet = self.transport.read_packet().await?;
1691        if !matches!(
1692            packet.header.message_type,
1693            MessageType::PriorityUpdate | MessageType::Deadline | MessageType::ExpireAt
1694        ) {
1695            return Err(RuntimeError::UnexpectedMessage(
1696                "server expected PRIORITY_UPDATE, DEADLINE, or EXPIRE_AT",
1697            ));
1698        }
1699        self.require_session_packet(
1700            &packet,
1701            "server received scheduling update for another session",
1702        )?;
1703        if packet.metadata.len() != SCHEDULING_METADATA_LEN || !packet.body.is_empty() {
1704            return Err(RuntimeError::UnexpectedMessage(
1705                "server received malformed scheduling metadata length",
1706            ));
1707        }
1708
1709        let metadata = SchedulingMetadata::parse(&packet.metadata)?;
1710        validate_scheduling_semantics(packet.header.message_type, &metadata)?;
1711        self.require_operation_frame(metadata.operation_id, packet.header.frame_id)?;
1712        self.operations.apply_scheduling_update(
1713            self.session_id,
1714            packet.header.message_type,
1715            metadata,
1716        )?;
1717        Ok(NnrpSchedulingUpdate {
1718            message_type: packet.header.message_type,
1719            metadata,
1720        })
1721    }
1722
1723    pub async fn receive_pressure_update(&mut self) -> Result<NnrpPressureUpdate, RuntimeError> {
1724        let packet = self.transport.read_packet().await?;
1725        if !matches!(
1726            packet.header.message_type,
1727            MessageType::Backpressure | MessageType::CreditUpdate
1728        ) {
1729            return Err(RuntimeError::UnexpectedMessage(
1730                "server expected BACKPRESSURE or CREDIT_UPDATE",
1731            ));
1732        }
1733        self.require_optional_session_packet(
1734            &packet,
1735            "server received pressure update for another session",
1736        )?;
1737        if packet.metadata.len() != PRESSURE_METADATA_LEN || !packet.body.is_empty() {
1738            return Err(RuntimeError::UnexpectedMessage(
1739                "server received malformed pressure metadata length",
1740            ));
1741        }
1742
1743        let metadata = PressureMetadata::parse(&packet.metadata)?;
1744        validate_pressure_semantics(packet.header.message_type, &metadata)?;
1745        self.pressure
1746            .apply_inbound(packet.header.message_type, metadata)?;
1747        Ok(NnrpPressureUpdate {
1748            message_type: packet.header.message_type,
1749            metadata,
1750        })
1751    }
1752
1753    pub async fn send_backpressure(
1754        &mut self,
1755        metadata: PressureMetadata,
1756    ) -> Result<(), RuntimeError> {
1757        validate_pressure_semantics(MessageType::Backpressure, &metadata)?;
1758        self.pressure
1759            .apply_outbound(MessageType::Backpressure, metadata)?;
1760        let mut header =
1761            CommonHeader::new(MessageType::Backpressure, PRESSURE_METADATA_LEN as u32, 0);
1762        header.session_id = self.session_id;
1763        self.transport
1764            .write_packet(&RuntimePacket::new(
1765                header,
1766                metadata.to_bytes()?.to_vec(),
1767                Vec::new(),
1768            )?)
1769            .await
1770    }
1771
1772    pub async fn send_credit_update(
1773        &mut self,
1774        metadata: PressureMetadata,
1775    ) -> Result<(), RuntimeError> {
1776        validate_pressure_semantics(MessageType::CreditUpdate, &metadata)?;
1777        self.pressure
1778            .apply_outbound(MessageType::CreditUpdate, metadata)?;
1779        self.write_runtime_packet(
1780            MessageType::CreditUpdate,
1781            0,
1782            metadata.to_bytes()?.to_vec(),
1783            Vec::new(),
1784        )
1785        .await
1786    }
1787
1788    pub async fn send_trace_context(
1789        &mut self,
1790        frame_id: u32,
1791        metadata: TraceContextMetadata,
1792        body: Vec<u8>,
1793    ) -> Result<(), RuntimeError> {
1794        validate_trace_context_semantics(&metadata)?;
1795        require_body_len(
1796            body.len(),
1797            metadata.body_bytes as usize,
1798            "server trace context body length mismatch",
1799        )?;
1800        self.write_runtime_packet(
1801            MessageType::TraceContext,
1802            frame_id,
1803            metadata.to_bytes()?.to_vec(),
1804            body,
1805        )
1806        .await
1807    }
1808
1809    pub async fn send_recoverable_error(
1810        &mut self,
1811        metadata: RecoverableErrorMetadata,
1812        diagnostics: Vec<u8>,
1813    ) -> Result<(), RuntimeError> {
1814        require_body_len(
1815            diagnostics.len(),
1816            metadata.diagnostic_bytes as usize,
1817            "server recoverable error diagnostic body length mismatch",
1818        )?;
1819        self.write_runtime_packet(
1820            MessageType::ErrorRecoverable,
1821            metadata.related_frame_id,
1822            metadata.to_bytes()?.to_vec(),
1823            diagnostics,
1824        )
1825        .await
1826    }
1827
1828    pub async fn send_retry_after(
1829        &mut self,
1830        metadata: RetryAfterMetadata,
1831        diagnostics: Vec<u8>,
1832    ) -> Result<(), RuntimeError> {
1833        require_body_len(
1834            diagnostics.len(),
1835            metadata.diagnostic_bytes as usize,
1836            "server retry-after diagnostic body length mismatch",
1837        )?;
1838        self.write_runtime_packet(
1839            MessageType::RetryAfter,
1840            0,
1841            metadata.to_bytes()?.to_vec(),
1842            diagnostics,
1843        )
1844        .await
1845    }
1846
1847    pub async fn send_result_hint(
1848        &mut self,
1849        metadata: ResultHintMetadata,
1850    ) -> Result<(), RuntimeError> {
1851        self.write_runtime_packet(
1852            MessageType::ResultHint,
1853            0,
1854            metadata.to_bytes()?.to_vec(),
1855            Vec::new(),
1856        )
1857        .await
1858    }
1859
1860    async fn write_runtime_packet(
1861        &mut self,
1862        message_type: MessageType,
1863        frame_id: u32,
1864        metadata: Vec<u8>,
1865        body: Vec<u8>,
1866    ) -> Result<(), RuntimeError> {
1867        let mut header = CommonHeader::new(message_type, metadata.len() as u32, body.len() as u32);
1868        header.session_id = self.session_id;
1869        header.frame_id = frame_id;
1870        self.transport
1871            .write_packet(&RuntimePacket::new(header, metadata, body)?)
1872            .await
1873    }
1874
1875    pub fn track_cache_object(&mut self, object_id: CacheObjectId) -> Result<(), RuntimeError> {
1876        if self.cache_objects.contains(&object_id) {
1877            return Ok(());
1878        }
1879        if self.max_cache_objects != 0 && self.cache_objects.len() >= self.max_cache_objects {
1880            return Err(RuntimeError::UnexpectedMessage(
1881                "server cache object limit reached",
1882            ));
1883        }
1884        self.cache_objects.push(object_id);
1885        Ok(())
1886    }
1887
1888    pub async fn receive_patch(&mut self) -> Result<SessionPatchMetadata, RuntimeError> {
1889        let packet = self.transport.read_packet().await?;
1890        if packet.header.message_type != MessageType::SessionPatch {
1891            return Err(RuntimeError::UnexpectedMessage(
1892                "server expected SESSION_PATCH",
1893            ));
1894        }
1895        self.require_session_packet(&packet, "server received patch for another session")?;
1896        if packet.metadata.len() != SESSION_PATCH_METADATA_LEN {
1897            return Err(RuntimeError::UnexpectedMessage(
1898                "server received malformed SESSION_PATCH metadata length",
1899            ));
1900        }
1901        Ok(SessionPatchMetadata::parse(&packet.metadata)?)
1902    }
1903
1904    pub async fn send_patch_ack(
1905        &mut self,
1906        ack: SessionPatchAckMetadata,
1907    ) -> Result<(), RuntimeError> {
1908        let mut header = CommonHeader::new(
1909            MessageType::SessionPatchAck,
1910            SESSION_PATCH_ACK_METADATA_LEN as u32,
1911            ack.profile_patch_ack_bytes,
1912        );
1913        header.session_id = self.session_id;
1914        self.transport
1915            .write_packet(&RuntimePacket::new(
1916                header,
1917                ack.to_bytes()?.to_vec(),
1918                Vec::new(),
1919            )?)
1920            .await
1921    }
1922
1923    pub async fn send_flow_update(
1924        &mut self,
1925        metadata: FlowUpdateMetadata,
1926    ) -> Result<(), RuntimeError> {
1927        let mut header =
1928            CommonHeader::new(MessageType::FlowUpdate, FLOW_UPDATE_METADATA_LEN as u32, 0);
1929        if !matches!(metadata.scope_kind, nnrp_core::FlowScopeKind::Connection) {
1930            header.session_id = self.session_id;
1931        }
1932        metadata.validate_routing(&header)?;
1933        self.transport
1934            .write_packet(&RuntimePacket::new(
1935                header,
1936                metadata.to_bytes()?.to_vec(),
1937                Vec::new(),
1938            )?)
1939            .await
1940    }
1941
1942    pub async fn receive_migrate(&mut self) -> Result<NnrpMigration, RuntimeError> {
1943        let packet = self.transport.read_packet().await?;
1944        if packet.header.message_type != MessageType::SessionMigrate {
1945            return Err(RuntimeError::UnexpectedMessage(
1946                "server expected SESSION_MIGRATE",
1947            ));
1948        }
1949        self.require_session_packet(&packet, "server received migrate for another session")?;
1950        if packet.metadata.len() != SESSION_MIGRATE_METADATA_LEN {
1951            return Err(RuntimeError::UnexpectedMessage(
1952                "server received malformed SESSION_MIGRATE metadata length",
1953            ));
1954        }
1955        Ok(NnrpMigration {
1956            metadata: SessionMigrateMetadata::parse(&packet.metadata)?,
1957        })
1958    }
1959
1960    pub async fn send_migrate_ack(
1961        &mut self,
1962        request: &SessionMigrateMetadata,
1963        ack: SessionMigrateAckMetadata,
1964    ) -> Result<(), RuntimeError> {
1965        nnrp_core::validate_migration_recovery(request, &ack)?;
1966        let mut header = CommonHeader::new(
1967            MessageType::SessionMigrateAck,
1968            SESSION_MIGRATE_ACK_METADATA_LEN as u32,
1969            0,
1970        );
1971        header.session_id = self.session_id;
1972        self.transport
1973            .write_packet(&RuntimePacket::new(
1974                header,
1975                ack.to_bytes()?.to_vec(),
1976                Vec::new(),
1977            )?)
1978            .await
1979    }
1980
1981    pub async fn receive_close(&mut self) -> Result<SessionCloseMetadata, RuntimeError> {
1982        let packet = self.transport.read_packet().await?;
1983        if packet.header.message_type != MessageType::SessionClose {
1984            return Err(RuntimeError::UnexpectedMessage(
1985                "server expected SESSION_CLOSE",
1986            ));
1987        }
1988        if packet.header.session_id != self.session_id {
1989            return Err(RuntimeError::UnexpectedMessage(
1990                "server received close for another session",
1991            ));
1992        }
1993        let close = SessionCloseMetadata::parse(&packet.metadata)?;
1994        self.lifecycle.begin_session_close(&packet.header, &close)?;
1995        self.pending_close = Some(close);
1996        Ok(close)
1997    }
1998
1999    pub async fn ack_close(&mut self, close: &SessionCloseMetadata) -> Result<(), RuntimeError> {
2000        let ack = SessionCloseAckMetadata {
2001            close_status: SessionCloseStatus::Closed,
2002            last_operation_id: close.last_operation_id,
2003            session_error_code: SESSION_ERROR_NONE,
2004        };
2005        let mut header = CommonHeader::new(
2006            MessageType::SessionCloseAck,
2007            SESSION_CLOSE_ACK_METADATA_LEN as u32,
2008            0,
2009        );
2010        header.session_id = self.session_id;
2011        self.lifecycle.apply_session_close_ack(&header, &ack)?;
2012        self.transport
2013            .write_packet(&RuntimePacket::new(
2014                header,
2015                ack.to_bytes()?.to_vec(),
2016                Vec::new(),
2017            )?)
2018            .await?;
2019        if self.pending_close == Some(*close) {
2020            self.pending_close = None;
2021        }
2022        Ok(())
2023    }
2024
2025    pub async fn close(mut self) -> Result<(), RuntimeError> {
2026        self.close_in_place().await
2027    }
2028
2029    pub async fn close_in_place(&mut self) -> Result<(), RuntimeError> {
2030        if let Some(close) = self.pending_close.take() {
2031            self.ack_close(&close).await?;
2032        }
2033        self.remove_from_registry()?;
2034        self.transport.close().await
2035    }
2036
2037    fn require_session_packet(
2038        &self,
2039        packet: &RuntimePacket,
2040        message: &'static str,
2041    ) -> Result<(), RuntimeError> {
2042        if packet.header.session_id != self.session_id {
2043            return Err(RuntimeError::UnexpectedMessage(message));
2044        }
2045        Ok(())
2046    }
2047
2048    fn require_optional_session_packet(
2049        &self,
2050        packet: &RuntimePacket,
2051        message: &'static str,
2052    ) -> Result<(), RuntimeError> {
2053        if packet.header.session_id != 0 && packet.header.session_id != self.session_id {
2054            return Err(RuntimeError::UnexpectedMessage(message));
2055        }
2056        Ok(())
2057    }
2058
2059    fn correlated_frame_id(&self, operation_id: u64) -> Result<u32, RuntimeError> {
2060        if operation_id == 0 {
2061            return Ok(0);
2062        }
2063        self.operation_frames
2064            .get(&operation_id)
2065            .copied()
2066            .ok_or(nnrp_core::NnrpError::UnknownOperation(operation_id).into())
2067    }
2068
2069    fn require_operation_frame(
2070        &self,
2071        operation_id: u64,
2072        frame_id: u32,
2073    ) -> Result<(), RuntimeError> {
2074        if self.correlated_frame_id(operation_id)? != frame_id {
2075            return Err(RuntimeError::UnexpectedMessage(
2076                "server runtime event frame id does not match its operation",
2077            ));
2078        }
2079        Ok(())
2080    }
2081
2082    fn operation_id_for_frame(&self, frame_id: u32) -> Result<u64, RuntimeError> {
2083        self.frame_operations
2084            .get(&frame_id)
2085            .copied()
2086            .ok_or(RuntimeError::UnexpectedMessage(
2087                "server frame id is not bound to an operation",
2088            ))
2089    }
2090
2091    fn update_registry_last_operation(&self, operation_id: u64) -> Result<(), RuntimeError> {
2092        let mut sessions = self.session_registry()?;
2093        if let Some(record) = sessions.get_mut(&self.session_id) {
2094            record.last_operation_id = record.last_operation_id.max(operation_id);
2095        }
2096        Ok(())
2097    }
2098
2099    fn remove_from_registry(&self) -> Result<(), RuntimeError> {
2100        self.session_registry()?.remove(&self.session_id);
2101        Ok(())
2102    }
2103
2104    fn session_registry(
2105        &self,
2106    ) -> Result<MutexGuard<'_, BTreeMap<u32, RuntimeSessionRecord>>, RuntimeError> {
2107        self.sessions
2108            .lock()
2109            .map_err(|_| RuntimeError::Internal("server session registry lock poisoned"))
2110    }
2111}
2112
2113fn current_unix_ms() -> u64 {
2114    SystemTime::now()
2115        .duration_since(UNIX_EPOCH)
2116        .map(|duration| duration.as_millis().min(u128::from(u64::MAX)) as u64)
2117        .unwrap_or(0)
2118}
2119
2120fn require_body_len(
2121    actual: usize,
2122    expected: usize,
2123    message: &'static str,
2124) -> Result<(), RuntimeError> {
2125    if actual != expected {
2126        return Err(RuntimeError::UnexpectedMessage(message));
2127    }
2128    Ok(())
2129}