1use crate::{
4 percent_encode_path_segment, ActivityCall, AvroValue, Client, Error, HandlerKind, HistoryEvent,
5 RequestProtocol, Result, WORKER_PROTOCOL_VERSION,
6};
7use serde::{Deserialize, Serialize};
8use serde_json::{json, Value};
9use std::sync::{Arc, Mutex};
10use std::time::{Duration, Instant};
11
12#[derive(Debug, Default)]
13struct SessionState {
14 expires: Option<Instant>,
15 receipt: Option<Value>,
16 close_receipt: Option<Value>,
17 original_ttl: Option<(i64, u32)>,
18}
19
20#[derive(Clone, Debug)]
26pub struct WorkerSession {
27 client: Client,
28 worker_id: String,
29 options: WorkerSessionOptions,
30 state: Arc<Mutex<SessionState>>,
31 operation: Arc<tokio::sync::Mutex<()>>,
32}
33
34impl WorkerSession {
35 fn new(client: Client, worker_id: String, options: WorkerSessionOptions) -> Self {
36 Self {
37 client,
38 worker_id,
39 options,
40 state: Arc::new(Mutex::new(SessionState::default())),
41 operation: Arc::new(tokio::sync::Mutex::new(())),
42 }
43 }
44
45 pub fn options(&self) -> &WorkerSessionOptions {
46 &self.options
47 }
48
49 pub fn active(&self) -> bool {
50 self.state
51 .lock()
52 .ok()
53 .and_then(|state| state.expires)
54 .is_some_and(|expires| Instant::now() < expires)
55 }
56
57 pub fn snapshot(&self) -> Result<Option<Value>> {
59 Ok(self
60 .state
61 .lock()
62 .map_err(|_| Error::WorkflowStatePoisoned)?
63 .receipt
64 .clone())
65 }
66
67 pub async fn create(&self) -> Result<Value> {
68 let _operation = self.operation.lock().await;
69 if self
70 .state
71 .lock()
72 .map_err(|_| Error::WorkflowStatePoisoned)?
73 .close_receipt
74 .is_some()
75 {
76 return Err(invalid("a closed session identity cannot be recreated"));
77 }
78 let response = self
79 .client
80 .create_worker_session(&self.worker_id, &self.options)
81 .await;
82 let response = self.settle(response, &["created", "reused", "reacquired"])?;
83 Ok(response)
84 }
85
86 pub async fn renew(&self) -> Result<Value> {
88 let _operation = self.operation.lock().await;
89 if !self.active() {
90 return Err(invalid(
91 "cannot renew an uncreated, expired or closed local handle",
92 ));
93 }
94 let response = self
95 .client
96 .renew_worker_session(
97 &self.worker_id,
98 self.options.session_id(),
99 self.options.lease_seconds,
100 )
101 .await;
102 let response = self.settle(response, &["heartbeat_recorded"])?;
103 Ok(response)
104 }
105
106 pub async fn close(&self, reason: &str) -> Result<Value> {
108 let _operation = self.operation.lock().await;
109 if let Some(receipt) = self
110 .state
111 .lock()
112 .map_err(|_| Error::WorkflowStatePoisoned)?
113 .close_receipt
114 .clone()
115 {
116 return Ok(receipt);
117 }
118 self.invalidate()?;
119 let response = self
120 .client
121 .close_worker_session(&self.worker_id, self.options.session_id(), reason)
122 .await?;
123 self.validate_receipt(&response, &["closed", "already_closed"], "closed")?;
124 let mut state = self
125 .state
126 .lock()
127 .map_err(|_| Error::WorkflowStatePoisoned)?;
128 state.receipt = Some(response["session"].clone());
129 state.close_receipt = Some(response.clone());
130 Ok(response)
131 }
132
133 fn invalidate(&self) -> Result<()> {
134 self.state
135 .lock()
136 .map_err(|_| Error::WorkflowStatePoisoned)?
137 .expires = None;
138 Ok(())
139 }
140
141 fn validate_receipt(&self, receipt: &Value, outcomes: &[&str], status: &str) -> Result<()> {
142 let session = &receipt["session"];
143 if receipt["admitted"] != true
144 || !receipt["outcome"]
145 .as_str()
146 .is_some_and(|outcome| outcomes.contains(&outcome))
147 || session["session_id"].as_str() != Some(self.options.session_id())
148 || session["namespace"].as_str() != Some(self.client.namespace.as_str())
149 || session["lease_owner"].as_str() != Some(self.worker_id.as_str())
150 || session["status"].as_str() != Some(status)
151 {
152 return Err(invalid(
153 "Server did not acknowledge this session, namespace, holder and lifecycle outcome",
154 ));
155 }
156 Ok(())
157 }
158
159 fn accept(&self, receipt: &Value, outcomes: &[&str]) -> Result<()> {
160 self.validate_receipt(receipt, outcomes, "active")?;
161 self.track(&receipt["session"])?;
162 Ok(())
163 }
164
165 fn settle(&self, result: Result<Value>, outcomes: &[&str]) -> Result<Value> {
166 match result.and_then(|receipt| {
167 self.accept(&receipt, outcomes)?;
168 Ok(receipt)
169 }) {
170 Ok(receipt) => Ok(receipt),
171 Err(error) => {
172 self.invalidate()?;
173 Err(error)
174 }
175 }
176 }
177
178 pub(super) async fn wait_until_unavailable(&self) {
179 while self.active() {
180 let delay = self
181 .state
182 .lock()
183 .ok()
184 .and_then(|state| state.expires)
185 .map(|expires| expires.saturating_duration_since(Instant::now()))
186 .unwrap_or_default()
187 .min(Duration::from_millis(100));
188 tokio::time::sleep(delay).await;
189 }
190 }
191
192 pub(super) fn track(&self, affinity: &Value) -> Result<()> {
193 let received: WorkerSessionOptions = serde_json::from_value(affinity.clone())?;
194 if affinity["session_id"].as_str() != Some(self.options.session_id())
195 || affinity["status"] != "active"
196 || affinity["lease_owner"].as_str() != Some(self.worker_id.as_str())
197 || affinity["queue"].as_str() != self.options.queue.as_deref()
198 || received.to_wire()? != self.options.to_wire()?
199 || affinity
200 .get("namespace")
201 .is_some_and(|namespace| namespace.as_str() != Some(self.client.namespace.as_str()))
202 {
203 return Err(invalid(
204 "activity session affinity does not match the current holder and queue",
205 ));
206 }
207 let expires = lease_deadline(affinity)?;
208 let ttl =
209 chrono::DateTime::parse_from_rfc3339(affinity["ttl_expires_at"].as_str().unwrap())
210 .map_err(|_| invalid("invalid session TTL deadline"))?;
211 let ttl = (ttl.timestamp(), ttl.timestamp_subsec_nanos());
212 let mut state = self
213 .state
214 .lock()
215 .map_err(|_| Error::WorkflowStatePoisoned)?;
216 if state.original_ttl.is_some_and(|original| original != ttl) {
217 return Err(invalid("session renewal changed its original TTL deadline"));
218 }
219 state.original_ttl = Some(ttl);
220 state.expires = Some(expires);
221 state.receipt = Some(affinity.clone());
222 if let Some(receipt) = &mut state.receipt {
223 receipt["namespace"] = json!(self.client.namespace);
224 }
225 state.close_receipt = None;
226 Ok(())
227 }
228}
229
230fn lease_deadline(affinity: &Value) -> Result<Instant> {
231 let now = std::time::SystemTime::now();
232 let mut budget = Duration::MAX;
233 for field in ["lease_expires_at", "ttl_expires_at"] {
234 let deadline = affinity[field]
235 .as_str()
236 .and_then(|value| chrono::DateTime::parse_from_rfc3339(value).ok())
237 .ok_or_else(|| invalid("session receipt requires valid lease and TTL deadlines"))?;
238 let timestamp = deadline.timestamp();
239 if timestamp < 0 {
240 return Err(invalid("session lease or TTL already expired"));
241 }
242 let deadline = std::time::UNIX_EPOCH
243 + Duration::new(timestamp as u64, deadline.timestamp_subsec_nanos());
244 let remaining = deadline
245 .duration_since(now)
246 .map_err(|_| invalid("session lease or TTL already expired"))?;
247 if remaining.is_zero() {
248 return Err(invalid("session lease or TTL already expired"));
249 }
250 budget = budget.min(remaining);
251 }
252 Instant::now()
253 .checked_add(budget)
254 .ok_or_else(|| invalid("session deadline exceeds local clock range"))
255}
256
257#[derive(Clone, Debug, Deserialize)]
258pub(super) struct SessionActivityTask {
259 #[serde(flatten)]
260 pub(super) task: crate::ActivityTask,
261 #[serde(default)]
262 pub(super) worker_session: Option<Value>,
263}
264
265#[derive(Debug, Deserialize)]
266pub(super) struct SessionPollResponse {
267 #[serde(default)]
268 pub(super) task: Option<SessionActivityTask>,
269 #[serde(default)]
270 poll_status: Option<String>,
271 #[serde(default)]
272 reason: Option<String>,
273}
274
275impl SessionPollResponse {
276 pub(super) fn outcome(&self) -> crate::WorkerPollOutcome {
277 crate::worker_poll_outcome(
278 self.task.is_some(),
279 self.poll_status.as_deref(),
280 self.reason.as_deref(),
281 )
282 }
283
284 pub(super) fn ordinary(self) -> crate::PollActivityTaskResponse {
285 crate::PollActivityTaskResponse {
286 task: self.task.map(|task| task.task),
287 poll_status: self.poll_status,
288 reason: self.reason,
289 }
290 }
291}
292
293impl crate::Worker {
294 pub fn worker_sessions(mut self, enabled: bool) -> Self {
296 self.client.worker_sessions_enabled = enabled;
297 self.session_registration_confirmed = Arc::new(std::sync::atomic::AtomicBool::new(false));
298 self
299 }
300
301 pub fn max_concurrent_worker_sessions(mut self, count: usize) -> Self {
302 self.client.max_concurrent_worker_sessions = count.max(1);
303 self
304 }
305
306 pub fn capabilities<I, S>(mut self, capabilities: I) -> Self
308 where
309 I: IntoIterator<Item = S>,
310 S: Into<String>,
311 {
312 self.resource_capabilities = capabilities
313 .into_iter()
314 .map(|value| value.into().trim().to_owned())
315 .collect();
316 self.resource_capabilities.sort();
317 self.resource_capabilities.dedup();
318 self
319 }
320
321 pub fn worker_session(&self, mut options: WorkerSessionOptions) -> Result<WorkerSession> {
323 self.require_session_registration()?;
324 if options.queue.is_none() {
325 options.queue = Some(self.task_queue.clone());
326 }
327 options.to_wire()?;
328 if options.queue.as_deref() != Some(self.task_queue.as_str()) {
329 return Err(invalid("session queue must match this worker"));
330 }
331 let mut sessions = self
332 .sessions
333 .lock()
334 .map_err(|_| Error::WorkflowStatePoisoned)?;
335 if let Some(session) = sessions.get(options.session_id()) {
336 if session.options.to_wire()? != options.to_wire()? {
337 return Err(invalid("session identity already has different options"));
338 }
339 return Ok(session.clone());
340 }
341 sessions.retain(|_, session| {
342 session
343 .state
344 .lock()
345 .map(|state| {
346 state.close_receipt.is_none()
347 && (Arc::strong_count(&session.state) > 1
348 || state
349 .expires
350 .is_some_and(|expires| Instant::now() < expires))
351 })
352 .unwrap_or(true)
353 });
354 if sessions.len() >= self.client.max_concurrent_worker_sessions {
355 return Err(invalid("local worker session capacity exhausted"));
356 }
357 let session = WorkerSession::new(self.client.clone(), self.worker_id.clone(), options);
358 sessions.insert(session.options.session_id.clone(), session.clone());
359 Ok(session)
360 }
361
362 pub(super) fn require_session_registration(&self) -> Result<()> {
363 if !self.client.worker_sessions_enabled
364 || !self
365 .session_registration_confirmed
366 .load(std::sync::atomic::Ordering::SeqCst)
367 {
368 return Err(invalid(
369 "register a session-capable worker before acquiring session handles or tasks",
370 ));
371 }
372 Ok(())
373 }
374
375 pub(super) fn track_session_task(
376 &self,
377 affinity: Option<&Value>,
378 ) -> Result<Option<WorkerSession>> {
379 let Some(affinity) = affinity.filter(|value| !value.is_null()) else {
380 return Ok(None);
381 };
382 let options: WorkerSessionOptions = serde_json::from_value(affinity.clone())?;
383 let session = self.worker_session(options)?;
384 session.track(affinity)?;
385 Ok(Some(session))
386 }
387
388 pub(super) fn session_available(&self) -> usize {
389 self.sessions
390 .lock()
391 .map(|sessions| {
392 self.client
393 .max_concurrent_worker_sessions
394 .saturating_sub(sessions.values().filter(|session| session.active()).count())
395 })
396 .unwrap_or(0)
397 }
398
399 pub(super) async fn close_worker_sessions(&self) -> Result<()> {
400 let sessions: Vec<_> = self
401 .sessions
402 .lock()
403 .map_err(|_| Error::WorkflowStatePoisoned)?
404 .values()
405 .cloned()
406 .collect();
407 let mut first_error = None;
408 for session in sessions {
409 if session.snapshot()?.is_some() {
411 if let Err(error) = session.close("worker_shutdown").await {
412 if first_error.is_none() {
413 first_error = Some(error);
414 }
415 }
416 }
417 }
418 first_error.map_or(Ok(()), Err)
419 }
420}
421
422impl crate::ActivityContext {
423 pub fn worker_session(&self) -> Option<&WorkerSession> {
425 self.worker_session.as_ref()
426 }
427}
428
429pub(super) fn settle_activity_heartbeat(
430 session: &WorkerSession,
431 result: Result<Value>,
432 task_id: &str,
433 attempt_id: &str,
434 owner: &str,
435) -> Result<crate::ActivityHeartbeatResponse> {
436 let result = result.and_then(|value| {
437 if value["task_id"].as_str() != Some(task_id)
438 || value["activity_attempt_id"].as_str() != Some(attempt_id)
439 || value["lease_owner"].as_str() != Some(owner)
440 || value["can_continue"] != true
441 || value["heartbeat_recorded"] != true
442 || value["cancel_requested"] != false
443 {
444 return Err(invalid(
445 "activity heartbeat did not acknowledge this claim and session",
446 ));
447 }
448 session.track(&value["worker_session"])?;
449 serde_json::from_value(value).map_err(Error::from)
450 });
451 result.map_err(|error| {
452 let _ = session.invalidate();
453 Error::ActivityExecutionAbandoned(error.to_string())
454 })
455}
456
457#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
463pub struct WorkerSessionOptions {
464 session_id: String,
465 #[serde(default, skip_serializing_if = "Option::is_none")]
466 connection: Option<String>,
467 #[serde(default, skip_serializing_if = "Option::is_none")]
468 queue: Option<String>,
469 #[serde(default)]
470 requirements: Vec<String>,
471 #[serde(default = "default_lease")]
472 lease_seconds: u64,
473 #[serde(default = "default_ttl")]
474 ttl_seconds: u64,
475 #[serde(default = "default_concurrency")]
476 max_concurrent_activities: usize,
477 #[serde(default = "default_true")]
478 create_if_missing: bool,
479 #[serde(default = "default_true")]
480 allow_reacquire_after_failure: bool,
481}
482
483fn default_lease() -> u64 {
484 120
485}
486fn default_ttl() -> u64 {
487 1800
488}
489fn default_concurrency() -> usize {
490 1
491}
492fn default_true() -> bool {
493 true
494}
495
496impl WorkerSessionOptions {
497 pub fn new(session_id: impl Into<String>) -> Self {
498 Self {
499 session_id: session_id.into().trim().to_owned(),
500 connection: None,
501 queue: None,
502 requirements: Vec::new(),
503 lease_seconds: default_lease(),
504 ttl_seconds: default_ttl(),
505 max_concurrent_activities: 1,
506 create_if_missing: true,
507 allow_reacquire_after_failure: true,
508 }
509 }
510
511 pub fn session_id(&self) -> &str {
512 &self.session_id
513 }
514 pub fn task_queue(&self) -> Option<&str> {
515 self.queue.as_deref()
516 }
517 pub fn lease_duration(&self) -> Duration {
518 Duration::from_secs(self.lease_seconds)
519 }
520 pub fn ttl_duration(&self) -> Duration {
521 Duration::from_secs(self.ttl_seconds)
522 }
523 pub fn connection(mut self, connection: impl Into<String>) -> Self {
524 self.connection = Some(connection.into().trim().to_owned());
525 self
526 }
527 pub fn queue(mut self, queue: impl Into<String>) -> Self {
528 self.queue = Some(queue.into().trim().to_owned());
529 self
530 }
531 pub fn requirements<I, S>(mut self, requirements: I) -> Self
532 where
533 I: IntoIterator<Item = S>,
534 S: Into<String>,
535 {
536 self.requirements = requirements
537 .into_iter()
538 .map(|value| value.into().trim().to_owned())
539 .collect();
540 self.requirements.sort();
541 self.requirements.dedup();
542 self
543 }
544 pub fn lease_seconds(mut self, seconds: u64) -> Self {
545 self.lease_seconds = seconds;
546 self
547 }
548 pub fn ttl_seconds(mut self, seconds: u64) -> Self {
549 self.ttl_seconds = seconds;
550 self
551 }
552 pub fn max_concurrent_activities(mut self, count: usize) -> Self {
553 self.max_concurrent_activities = count;
554 self
555 }
556 pub fn create_if_missing(mut self, enabled: bool) -> Self {
557 self.create_if_missing = enabled;
558 self
559 }
560 pub fn allow_reacquire_after_failure(mut self, enabled: bool) -> Self {
561 self.allow_reacquire_after_failure = enabled;
562 self
563 }
564
565 pub fn to_wire(&self) -> Result<Value> {
566 identifier("session_id", &self.session_id)?;
567 for (field, value) in [("connection", &self.connection), ("queue", &self.queue)] {
568 if let Some(value) = value {
569 identifier(field, value)?;
570 }
571 }
572 for requirement in &self.requirements {
573 identifier("requirements", requirement)?;
574 }
575 if self.lease_seconds == 0
576 || self.ttl_seconds == 0
577 || self.max_concurrent_activities == 0
578 || self.lease_seconds > i64::MAX as u64
579 || self.ttl_seconds > i64::MAX as u64
580 || self.max_concurrent_activities as u128 > i64::MAX as u128
581 {
582 return Err(invalid("lease, TTL and concurrency must be positive"));
583 }
584 let mut canonical = self.clone();
585 canonical.requirements.sort();
586 canonical.requirements.dedup();
587 Ok(serde_json::to_value(canonical)?)
588 }
589}
590
591fn invalid(message: &str) -> Error {
592 Error::WorkerLoop(format!("invalid_worker_session: {message}"))
593}
594
595fn identifier(field: &str, value: &str) -> Result<()> {
596 if value.trim().is_empty() || value.trim() != value || value.chars().count() > 255 {
597 return Err(invalid(&format!(
598 "{field} must be a canonical non-empty string of at most 255 characters"
599 )));
600 }
601 Ok(())
602}
603
604pub(super) fn validate_resource_capability(value: &str) -> Result<()> {
605 identifier("capabilities", value)?;
606 if value.starts_with("prepared_local_")
607 || value.starts_with("cancellation_scope")
608 || matches!(
609 value,
610 "local_activities"
611 | "worker_sessions"
612 | "sticky_execution"
613 | "cooperative_cancellation"
614 )
615 {
616 return Err(invalid(
617 "built-in execution capabilities must use their explicit worker opt-in",
618 ));
619 }
620 Ok(())
621}
622
623impl Client {
624 pub async fn create_worker_session(
629 &self,
630 worker_id: &str,
631 options: &WorkerSessionOptions,
632 ) -> Result<Value> {
633 identifier("worker_id", worker_id)?;
634 let mut body = options.to_wire()?;
635 body["worker_id"] = json!(worker_id);
636 self.request_json(
637 reqwest::Method::POST,
638 "/worker/sessions",
639 RequestProtocol::Worker(WORKER_PROTOCOL_VERSION),
640 Some(&body),
641 )
642 .await
643 }
644
645 pub async fn renew_worker_session(
647 &self,
648 worker_id: &str,
649 session_id: &str,
650 lease_seconds: u64,
651 ) -> Result<Value> {
652 identifier("worker_id", worker_id)?;
653 identifier("session_id", session_id)?;
654 if lease_seconds == 0 || lease_seconds > i64::MAX as u64 {
655 return Err(invalid(
656 "lease_seconds must be a positive signed 64-bit integer",
657 ));
658 }
659 self.request_json(
660 reqwest::Method::POST,
661 &format!(
662 "/worker/sessions/{}/heartbeat",
663 percent_encode_path_segment(session_id)
664 ),
665 RequestProtocol::Worker(WORKER_PROTOCOL_VERSION),
666 Some(&json!({"worker_id":worker_id,"lease_seconds":lease_seconds})),
667 )
668 .await
669 }
670
671 pub async fn close_worker_session(
673 &self,
674 worker_id: &str,
675 session_id: &str,
676 reason: &str,
677 ) -> Result<Value> {
678 identifier("worker_id", worker_id)?;
679 identifier("session_id", session_id)?;
680 self.request_json(
681 reqwest::Method::DELETE,
682 &format!(
683 "/worker/sessions/{}",
684 percent_encode_path_segment(session_id)
685 ),
686 RequestProtocol::Worker(WORKER_PROTOCOL_VERSION),
687 Some(&json!({"worker_id":worker_id,"reason":reason})),
688 )
689 .await
690 }
691}
692
693impl ActivityCall {
694 pub fn in_worker_session(mut self, options: WorkerSessionOptions) -> Self {
699 self.worker_session = Some(options);
700 self
701 }
702
703 pub async fn typed<O: serde::de::DeserializeOwned>(mut self) -> Result<O> {
705 let activity_type = self.activity_type.clone();
706 let result =
707 std::future::poll_fn(|cx| std::pin::Pin::new(&mut self).poll_avro_value(cx)).await?;
708 crate::decode_handler_result(result, HandlerKind::Activity, &activity_type)
709 }
710
711 pub async fn avro_value(mut self) -> Result<AvroValue> {
713 std::future::poll_fn(|cx| std::pin::Pin::new(&mut self).poll_avro_value(cx)).await
714 }
715}
716
717impl crate::ParallelCall {
718 pub fn in_worker_session(mut self, options: WorkerSessionOptions) -> Self {
720 self.worker_session = Some(options);
721 self
722 }
723}
724
725pub(super) fn recorded_session(events: &[&HistoryEvent], sequence: u64) -> Result<Option<Value>> {
726 let mut original = None;
727 for event in events {
728 for source in [Some(&event.payload), event.payload.get("activity")]
729 .into_iter()
730 .flatten()
731 {
732 let Some(value) = source.get("worker_session") else {
733 continue;
734 };
735 let session = if value.is_null() {
736 None
737 } else {
738 let mut value = value.clone();
739 if let Some(object) = value.as_object_mut() {
740 for field in ["lease_seconds", "ttl_seconds", "max_concurrent_activities"] {
741 if object.get(field).is_some_and(Value::is_null) {
742 object.remove(field);
743 }
744 }
745 }
746 let options = serde_json::from_value::<WorkerSessionOptions>(value)
747 .and_then(|options| options.to_wire().map_err(serde::de::Error::custom))
748 .map_err(|error| {
749 crate::invalid_recorded_history(
750 "worker_session_invalid",
751 sequence,
752 "valid worker-session routing",
753 &error.to_string(),
754 "recorded worker-session metadata is invalid",
755 )
756 })?;
757 Some(options)
758 };
759 if original.as_ref().is_some_and(|old| old != &session) {
760 return Err(crate::invalid_recorded_history(
761 "worker_session_conflict",
762 sequence,
763 "one worker-session identity",
764 "conflicting session options",
765 "activity history changes worker-session routing at one command boundary",
766 ));
767 }
768 original = Some(session);
769 }
770 }
771 Ok(original.flatten())
772}
773
774#[cfg(test)]
775mod tests {
776 use super::*;
777
778 #[test]
779 fn invalid_routing_and_lifetime_are_refused_before_dispatch() {
780 for options in [
781 WorkerSessionOptions::new(" "),
782 WorkerSessionOptions::new("render").queue(" "),
783 WorkerSessionOptions::new("render").connection(""),
784 WorkerSessionOptions::new("render").requirements([""]),
785 WorkerSessionOptions::new("render").lease_seconds(0),
786 WorkerSessionOptions::new("render").ttl_seconds(0),
787 WorkerSessionOptions::new("render").max_concurrent_activities(0),
788 ] {
789 assert!(
790 matches!(options.to_wire(), Err(Error::WorkerLoop(message)) if message.starts_with("invalid_worker_session:"))
791 );
792 }
793 }
794
795 #[test]
796 fn session_requirements_have_one_canonical_wire_identity() {
797 let a =
798 WorkerSessionOptions::new(" render ").requirements(["gpu:l4", "codec:av1", "gpu:l4"]);
799 let b = WorkerSessionOptions::new("render").requirements(["codec:av1", "gpu:l4"]);
800 assert_eq!(a, b);
801 assert_eq!(a.to_wire().unwrap(), b.to_wire().unwrap());
802 assert_eq!(
803 a.to_wire().unwrap()["requirements"],
804 json!(["codec:av1", "gpu:l4"])
805 );
806 }
807}