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}