durable_workflow/
worker_session.rs

1//! Worker-held session routing. Session memory is never durable workflow state.
2
3use 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/// One worker's acknowledged session lease.
21///
22/// Clones share lifecycle state. `active()` is a local lease hint, not authority
23/// for an external side effect. Server fences activity completion. Replacement
24/// holders must rebuild process-local resources rather than restore this handle.
25#[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    /// Last validated session snapshot, including holder and original TTL.
58    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    /// Renew the holder lease without extending the original TTL.
87    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    /// Close the session after its activities drain. Duplicate close reuses its receipt.
107    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    /// Opt in to remote activity sessions. Sticky execution stays unsupported.
295    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    /// Additional resource requirements this worker can actually satisfy.
307    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    /// Obtain a shared handle for this registered worker and immutable session options.
322    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            // Never issue a close for a handle that was not admitted to this holder.
410            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    /// Session affinity and holder-local lifecycle for this remote activity.
424    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/// Routing and lifetime options for one worker-held session.
458///
459/// A holder replacement must rebuild its local resources. TTL expiry and explicit
460/// close are terminal for the session identity, even when lease reacquisition is
461/// enabled. Options do not create a session until sent to Server.
462#[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    /// Create, reuse or reacquire a session for this registered worker.
625    ///
626    /// Server owns admission, requirements, capacity and holder authority.
627    /// An admitted reacquisition requires rebuilding worker-local resources.
628    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    /// Renew only the current session holder's lease. This does not extend TTL.
646    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    /// Close one holder's session. A closed identity cannot be reacquired.
672    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    /// Route this remote activity through a durable worker-session identity.
695    ///
696    /// The matching worker needs the session capability and its requirements.
697    /// Options are validated before scheduling and checked during cold replay.
698    pub fn in_worker_session(mut self, options: WorkerSessionOptions) -> Self {
699        self.worker_session = Some(options);
700        self
701    }
702
703    /// Decode this activity result directly from the lossless Avro value.
704    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    /// Return the lossless Avro result of this activity call.
712    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    /// Route every activity leaf in this parallel group through one session.
719    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}