From 4c14ddfeb8964d09b83f14ce7629e42863fa67ce Mon Sep 17 00:00:00 2001
From: cai <cai@nbcai.cc>
Date: Wed, 12 Aug 2026 15:58:53 +0800
Subject: [PATCH] feat(helper): distinguish current participant visibility

---
 src/main.rs    |  183 ++++++++++++++++++++++++++++++++++++++++-----
 src/service.rs |    2 
 2 files changed, 164 insertions(+), 21 deletions(-)

diff --git a/src/main.rs b/src/main.rs
index 8803bba..af61418 100644
--- a/src/main.rs
+++ b/src/main.rs
@@ -127,6 +127,7 @@
 struct ControlledFixtureVisibilityEvidence {
     first_visible_bucket: &'static str,
     visibility_source: &'static str,
+    visibility_result: &'static str,
     binding_matched: bool,
 }
 
@@ -368,48 +369,74 @@
     }
 }
 
-async fn observe_controlled_fixture_post_expiry<F>(
+async fn observe_controlled_fixture_post_expiry_views<F, G>(
     received_at: Instant,
     observation_deadline: Instant,
-    actual_participant: &str,
+    held_participant: &str,
     expected_participant: Option<&str>,
     requested_sequence: &str,
     lifecycle_active: Arc<AtomicBool>,
-    mut read_attributes: F,
+    mut read_held_attributes: F,
+    mut read_current_participant: G,
 ) -> Option<ControlledFixtureVisibilityEvidence>
 where
     F: FnMut() -> std::collections::HashMap<String, String>,
+    G: FnMut() -> Option<(String, std::collections::HashMap<String, String>)>,
 {
     loop {
         if !lifecycle_active.load(Ordering::Acquire) {
             return None;
         }
         let now = Instant::now();
-        let decision = classify_controlled_fixture_attributes(
-            actual_participant,
+        let held_decision = classify_controlled_fixture_attributes(
+            held_participant,
             expected_participant,
-            &read_attributes(),
+            &read_held_attributes(),
             requested_sequence,
         );
-        match decision {
+        match held_decision {
             Ok(()) => {
                 return Some(ControlledFixtureVisibilityEvidence {
                     first_visible_bucket: controlled_fixture_visibility_bucket(
                         now.saturating_duration_since(received_at),
                     ),
                     visibility_source: "participant_attributes_poll",
+                    visibility_result: "held_visible",
                     binding_matched: true,
                 });
             }
-            Err("missing_attributes") if now < observation_deadline => {}
-            Err("missing_attributes") => {
-                return Some(ControlledFixtureVisibilityEvidence {
-                    first_visible_bucket: "never_visible_within_observation_window",
-                    visibility_source: "participant_attributes_poll",
-                    binding_matched: true,
-                });
-            }
+            Err("missing_attributes") => {}
             Err(_) => return None,
+        }
+        let Some((current_identity, current_attributes)) = read_current_participant() else {
+            return None;
+        };
+        match classify_controlled_fixture_attributes(
+            &current_identity,
+            expected_participant,
+            &current_attributes,
+            requested_sequence,
+        ) {
+            Ok(()) => {
+                return Some(ControlledFixtureVisibilityEvidence {
+                    first_visible_bucket: controlled_fixture_visibility_bucket(
+                        now.saturating_duration_since(received_at),
+                    ),
+                    visibility_source: "current_room_lookup",
+                    visibility_result: "held_stale_current_visible",
+                    binding_matched: true,
+                });
+            }
+            Err("missing_attributes") => {}
+            Err(_) => return None,
+        }
+        if now >= observation_deadline {
+            return Some(ControlledFixtureVisibilityEvidence {
+                first_visible_bucket: "never_visible_within_observation_window",
+                visibility_source: "held_and_current_room_lookup",
+                visibility_result: "unavailable_both",
+                binding_matched: true,
+            });
         }
         sleep(CONTROLLED_FIXTURE_PROBE_RECHECK_DELAY).await;
     }
@@ -435,6 +462,7 @@
             "stage": "post_expiry_visibility",
             "first_visible_bucket": evidence.first_visible_bucket,
             "visibility_source": evidence.visibility_source,
+            "visibility_result": evidence.visibility_result,
             "binding_matched": evidence.binding_matched,
             "call_id_hash": probe.call_id_hash,
             "trace_id_hash": probe.call_trace_id_hash,
@@ -448,6 +476,7 @@
 fn spawn_controlled_fixture_post_expiry_observation(
     probe: PendingControlledFixtureProbe,
     participant: RemoteParticipant,
+    room: Arc<Room>,
     expected_participant: Option<String>,
     lifecycle_active: Arc<AtomicBool>,
     runtime_call_id: String,
@@ -455,7 +484,7 @@
 ) {
     tokio::spawn(async move {
         let participant_identity = participant.identity().to_string();
-        let evidence = observe_controlled_fixture_post_expiry(
+        let evidence = observe_controlled_fixture_post_expiry_views(
             probe.received_at,
             probe.received_at + CONTROLLED_FIXTURE_POST_EXPIRY_WINDOW,
             &participant_identity,
@@ -463,6 +492,11 @@
             &probe.sequence,
             lifecycle_active.clone(),
             || participant.attributes(),
+            || {
+                room.remote_participants()
+                    .get(&probe.sender)
+                    .map(|current| (current.identity().to_string(), current.attributes()))
+            },
         )
         .await;
         if lifecycle_active.load(Ordering::Acquire) {
@@ -1609,6 +1643,7 @@
         spawn_controlled_fixture_post_expiry_observation(
             probe,
             participant.clone(),
+            sink.room.clone(),
             expected_participant.map(str::to_string),
             lifecycle_active,
             runtime_call_id.to_string(),
@@ -5855,7 +5890,7 @@
                 published: true,
             }
         );
-        let evidence = observe_controlled_fixture_post_expiry(
+        let evidence = observe_controlled_fixture_post_expiry_views(
             started_at,
             started_at + CONTROLLED_FIXTURE_POST_EXPIRY_WINDOW,
             "user-1",
@@ -5869,6 +5904,7 @@
                     HashMap::new()
                 }
             },
+            || Some(("user-1".to_string(), HashMap::new())),
         )
         .await;
         assert_eq!(acknowledged.len(), 1);
@@ -5877,13 +5913,14 @@
             Some(ControlledFixtureVisibilityEvidence {
                 first_visible_bucket: "250_500ms",
                 visibility_source: "participant_attributes_poll",
+                visibility_result: "held_visible",
                 binding_matched: true,
             })
         );
 
         let never_started_at = Instant::now();
         assert_eq!(
-            observe_controlled_fixture_post_expiry(
+            observe_controlled_fixture_post_expiry_views(
                 never_started_at,
                 never_started_at + Duration::from_millis(40),
                 "user-1",
@@ -5891,17 +5928,19 @@
                 "fixture-01",
                 Arc::new(AtomicBool::new(true)),
                 HashMap::new,
+                || Some(("user-1".to_string(), HashMap::new())),
             )
             .await,
             Some(ControlledFixtureVisibilityEvidence {
                 first_visible_bucket: "never_visible_within_observation_window",
-                visibility_source: "participant_attributes_poll",
+                visibility_source: "held_and_current_room_lookup",
+                visibility_result: "unavailable_both",
                 binding_matched: true,
             })
         );
 
         assert_eq!(
-            observe_controlled_fixture_post_expiry(
+            observe_controlled_fixture_post_expiry_views(
                 Instant::now(),
                 Instant::now() + Duration::from_millis(50),
                 "cross-call-user",
@@ -5909,6 +5948,7 @@
                 "fixture-01",
                 Arc::new(AtomicBool::new(true)),
                 || expected.clone(),
+                || Some(("user-1".to_string(), expected.clone())),
             )
             .await,
             None
@@ -5916,7 +5956,7 @@
 
         let inactive = Arc::new(AtomicBool::new(false));
         assert_eq!(
-            observe_controlled_fixture_post_expiry(
+            observe_controlled_fixture_post_expiry_views(
                 Instant::now(),
                 Instant::now() + Duration::from_millis(50),
                 "user-1",
@@ -5924,6 +5964,7 @@
                 "fixture-01",
                 inactive,
                 HashMap::new,
+                || Some(("user-1".to_string(), HashMap::new())),
             )
             .await,
             None
@@ -5959,6 +6000,7 @@
                 "sequence_hash",
                 "stage",
                 "trace_id_hash",
+                "visibility_result",
                 "visibility_source",
             ]
         );
@@ -5974,6 +6016,105 @@
         }
     }
 
+    #[tokio::test]
+    async fn production_post_expiry_observation_distinguishes_held_stale_from_current_room_view() {
+        let started_at = Instant::now();
+        let current_attributes = HashMap::from([
+            (
+                "inputSourceCategory".to_string(),
+                "controlled_fixture".to_string(),
+            ),
+            (
+                "clientFixtureSequence".to_string(),
+                "fixture-01".to_string(),
+            ),
+        ]);
+        let mut acknowledged = HashSet::new();
+        let expired_ack = complete_controlled_fixture_ack_publish(
+            async { Ok::<(), ()>(()) },
+            "runtime-call-current-view",
+            "runtime-trace-current-view",
+            false,
+            "timeout",
+            Some("expired"),
+            "call-current-view",
+            "trace-current-view",
+            CONTROLLED_FIXTURE_GENERATION,
+            "fixture-01",
+            &mut acknowledged,
+        )
+        .await;
+        let evidence = observe_controlled_fixture_post_expiry_views(
+            started_at,
+            started_at + Duration::from_millis(100),
+            "user-1",
+            Some("user-1"),
+            "fixture-01",
+            Arc::new(AtomicBool::new(true)),
+            HashMap::new,
+            || Some(("user-1".to_string(), current_attributes.clone())),
+        )
+        .await;
+
+        assert_eq!(
+            expired_ack,
+            ControlledFixtureAckPublishOutcome {
+                observed: false,
+                published: true,
+            }
+        );
+        assert_eq!(acknowledged.len(), 1);
+        assert_eq!(
+            evidence,
+            Some(ControlledFixtureVisibilityEvidence {
+                first_visible_bucket: "lte_250ms",
+                visibility_source: "current_room_lookup",
+                visibility_result: "held_stale_current_visible",
+                binding_matched: true,
+            })
+        );
+        let mut observer_starts = 0;
+        assert!(!start_observer_after_controlled_fixture_probe(
+            true,
+            Some(expired_ack.observed),
+            || observer_starts += 1,
+        ));
+        assert_eq!(observer_starts, 0);
+
+        for current_view in [
+            None,
+            Some(("cross-call-user".to_string(), current_attributes.clone())),
+            Some((
+                "user-1".to_string(),
+                HashMap::from([
+                    (
+                        "inputSourceCategory".to_string(),
+                        "controlled_fixture".to_string(),
+                    ),
+                    (
+                        "clientFixtureSequence".to_string(),
+                        "fixture-old".to_string(),
+                    ),
+                ]),
+            )),
+        ] {
+            assert_eq!(
+                observe_controlled_fixture_post_expiry_views(
+                    Instant::now(),
+                    Instant::now() + Duration::from_millis(20),
+                    "user-1",
+                    Some("user-1"),
+                    "fixture-01",
+                    Arc::new(AtomicBool::new(true)),
+                    HashMap::new,
+                    || current_view.clone(),
+                )
+                .await,
+                None
+            );
+        }
+    }
+
     #[test]
     fn controlled_fixture_ack_payload_is_reliable_and_redacted() {
         let ack = ControlledFixtureAttributeAck {
diff --git a/src/service.rs b/src/service.rs
index 76e2627..2d0d3ab 100644
--- a/src/service.rs
+++ b/src/service.rs
@@ -1323,6 +1323,7 @@
             &ControlledFixtureVisibilityEvidence {
                 first_visible_bucket: "250_500ms",
                 visibility_source: "participant_attributes_poll",
+                visibility_result: "held_visible",
                 binding_matched: true,
             },
         );
@@ -1330,6 +1331,7 @@
             .expect("post-expiry evidence must enter the production stdout projection");
         assert_eq!(projected["extension"]["first_visible_bucket"], "250_500ms");
         assert_eq!(projected["extension"]["binding_matched"], true);
+        assert_eq!(projected["extension"]["visibility_result"], "held_visible");
     }
 
     #[test]

--
Gitblit v1.9.3