Record stream first byte on upstream event

This commit is contained in:
fawney19
2026-05-27 10:19:49 +08:00
parent 433a4d3c7d
commit 0ee6e393ce
@@ -605,209 +605,6 @@ fn append_stream_capture_bytes(
}
}
const CLIENT_VISIBLE_TEXT_TRACKER_MAX_BUFFER_BYTES: usize = 256 * 1024;
struct ClientVisibleTextTracker {
client_api_format: String,
buffered: Vec<u8>,
observed: bool,
}
impl ClientVisibleTextTracker {
fn new(client_api_format: &str) -> Self {
Self {
client_api_format: client_api_format.trim().to_ascii_lowercase(),
buffered: Vec::new(),
observed: false,
}
}
fn observe_chunk(&mut self, chunk: &[u8]) -> bool {
if self.observed || chunk.is_empty() {
return false;
}
self.buffered.extend_from_slice(chunk);
while let Some((block_end, separator_len)) = find_sse_block_boundary(&self.buffered) {
let block_len = block_end + separator_len;
let block = self.buffered.drain(..block_len).collect::<Vec<_>>();
if sse_block_has_client_visible_text(&self.client_api_format, &block) {
self.observed = true;
self.buffered.clear();
return true;
}
}
if self.buffered.len() > CLIENT_VISIBLE_TEXT_TRACKER_MAX_BUFFER_BYTES {
self.buffered.clear();
}
false
}
}
fn sse_block_has_client_visible_text(client_api_format: &str, block: &[u8]) -> bool {
let text = String::from_utf8_lossy(block);
let mut data_lines = Vec::new();
for raw_line in text.lines() {
let line = raw_line.trim_end_matches('\r');
let Some(data) = line.strip_prefix("data:") else {
continue;
};
data_lines.push(data.strip_prefix(' ').unwrap_or(data));
}
if data_lines.is_empty() {
return false;
}
let data = data_lines.join("\n");
let data = data.trim();
if data.is_empty() || data == "[DONE]" {
return false;
}
let Ok(value) = serde_json::from_str::<Value>(data) else {
return false;
};
stream_json_has_client_visible_text(client_api_format, &value)
}
fn stream_json_has_client_visible_text(client_api_format: &str, value: &Value) -> bool {
if client_api_format.contains("openai:chat") || client_api_format.contains("openai_chat") {
return openai_chat_stream_json_has_visible_text(value);
}
if client_api_format.contains("openai:responses")
|| client_api_format.contains("openai_responses")
{
return openai_responses_stream_json_has_visible_text(value);
}
if client_api_format.contains("claude") {
return claude_stream_json_has_visible_text(value);
}
if client_api_format.contains("gemini") {
return gemini_stream_json_has_visible_text(value);
}
openai_chat_stream_json_has_visible_text(value)
|| openai_responses_stream_json_has_visible_text(value)
|| claude_stream_json_has_visible_text(value)
|| gemini_stream_json_has_visible_text(value)
}
fn visible_text_value(value: Option<&Value>) -> bool {
value
.and_then(Value::as_str)
.is_some_and(|text| text.chars().any(|ch| !ch.is_whitespace()))
}
fn openai_chat_stream_json_has_visible_text(value: &Value) -> bool {
value
.get("choices")
.and_then(Value::as_array)
.is_some_and(|choices| {
choices.iter().any(|choice| {
choice
.get("delta")
.and_then(Value::as_object)
.is_some_and(|delta| {
visible_text_value(delta.get("content"))
|| visible_text_value(delta.get("refusal"))
})
|| visible_text_value(choice.get("text"))
})
})
}
fn openai_responses_stream_json_has_visible_text(value: &Value) -> bool {
match value
.get("type")
.and_then(Value::as_str)
.unwrap_or_default()
{
"response.output_text.delta" => visible_text_value(value.get("delta")),
"response.output_text.done" => visible_text_value(value.get("text")),
"response.completed" => value
.get("response")
.is_some_and(openai_response_body_has_visible_text),
_ => false,
}
}
fn openai_response_body_has_visible_text(value: &Value) -> bool {
value
.get("output")
.and_then(Value::as_array)
.is_some_and(|items| {
items.iter().any(|item| {
item.get("content")
.and_then(Value::as_array)
.is_some_and(|content| {
content.iter().any(|part| {
part.get("type")
.and_then(Value::as_str)
.is_some_and(|kind| matches!(kind, "output_text" | "text"))
&& visible_text_value(part.get("text"))
})
})
})
})
}
fn claude_stream_json_has_visible_text(value: &Value) -> bool {
match value
.get("type")
.and_then(Value::as_str)
.unwrap_or_default()
{
"content_block_delta" => {
value
.get("delta")
.and_then(Value::as_object)
.is_some_and(|delta| {
delta.get("type").and_then(Value::as_str) == Some("text_delta")
&& visible_text_value(delta.get("text"))
})
}
"content_block_start" => value
.get("content_block")
.and_then(Value::as_object)
.is_some_and(|block| {
block.get("type").and_then(Value::as_str) == Some("text")
&& visible_text_value(block.get("text"))
}),
_ => false,
}
}
fn gemini_stream_json_has_visible_text(value: &Value) -> bool {
let event = value
.get("response")
.filter(|response| response.get("candidates").is_some())
.unwrap_or(value);
event
.get("candidates")
.and_then(Value::as_array)
.is_some_and(|candidates| {
candidates.iter().any(|candidate| {
candidate
.get("content")
.and_then(|content| content.get("parts"))
.and_then(Value::as_array)
.is_some_and(|parts| {
parts.iter().any(|part| {
!part
.get("thought")
.and_then(Value::as_bool)
.unwrap_or(false)
&& visible_text_value(part.get("text"))
})
})
})
})
}
fn observe_stream_usage_bytes(
observer: &mut StreamingStandardTerminalObserver,
report_context: &Value,
@@ -1913,17 +1710,36 @@ fn stream_chunk_contains_sse_done(chunk: &[u8]) -> bool {
tracker.observe_chunk(chunk)
}
async fn next_stream_frame<R>(
buffered_frames: &mut VecDeque<StreamFrame>,
struct ObservedStreamFrame {
frame: StreamFrame,
observed_at: Instant,
}
async fn read_next_observed_stream_frame<R>(
lines: &mut FramedRead<R, LinesCodec>,
) -> Result<Option<StreamFrame>, GatewayError>
) -> Result<Option<ObservedStreamFrame>, GatewayError>
where
R: tokio::io::AsyncRead + Unpin,
{
Ok(read_next_frame(lines)
.await?
.map(|frame| ObservedStreamFrame {
frame,
observed_at: Instant::now(),
}))
}
async fn next_stream_frame<R>(
buffered_frames: &mut VecDeque<ObservedStreamFrame>,
lines: &mut FramedRead<R, LinesCodec>,
) -> Result<Option<ObservedStreamFrame>, GatewayError>
where
R: tokio::io::AsyncRead + Unpin,
{
if let Some(frame) = buffered_frames.pop_front() {
return Ok(Some(frame));
}
read_next_frame(lines).await
read_next_observed_stream_frame(lines).await
}
fn should_refresh_stream_usage_telemetry(
@@ -1943,11 +1759,19 @@ fn stream_elapsed_ms_since(started_at: Instant) -> u64 {
started_at.elapsed().as_millis().min(u128::from(u64::MAX)) as u64
}
fn first_visible_text_telemetry(
fn stream_elapsed_ms_at(started_at: Instant, observed_at: Instant) -> u64 {
observed_at
.saturating_duration_since(started_at)
.as_millis()
.min(u128::from(u64::MAX)) as u64
}
fn first_stream_event_telemetry(
stream_started_at: Instant,
event_observed_at: Instant,
upstream_telemetry: Option<&ExecutionTelemetry>,
) -> ExecutionTelemetry {
let elapsed_ms = stream_elapsed_ms_since(stream_started_at);
let elapsed_ms = stream_elapsed_ms_at(stream_started_at, event_observed_at);
ExecutionTelemetry {
ttfb_ms: Some(elapsed_ms),
elapsed_ms: Some(elapsed_ms),
@@ -1955,6 +1779,28 @@ fn first_visible_text_telemetry(
}
}
fn maybe_capture_first_stream_event_telemetry(
stream_started_at: Instant,
event_observed_at: Instant,
upstream_telemetry: Option<&ExecutionTelemetry>,
usage_stream_telemetry: &mut Option<ExecutionTelemetry>,
) -> bool {
if usage_stream_telemetry
.as_ref()
.and_then(|telemetry| telemetry.ttfb_ms)
.is_some()
{
return false;
}
*usage_stream_telemetry = Some(first_stream_event_telemetry(
stream_started_at,
event_observed_at,
upstream_telemetry,
));
true
}
fn usage_refresh_telemetry(
upstream_telemetry: &ExecutionTelemetry,
usage_stream_telemetry: Option<&ExecutionTelemetry>,
@@ -1966,35 +1812,32 @@ fn usage_refresh_telemetry(
}
}
fn maybe_record_first_visible_text_started(
fn maybe_record_first_stream_event_started(
state: &AppState,
lifecycle_seed: &LifecycleUsageSeed,
status_code: u16,
tracker: &mut ClientVisibleTextTracker,
chunk: &[u8],
stream_started_at: Instant,
event_observed_at: Instant,
upstream_telemetry: Option<&ExecutionTelemetry>,
usage_stream_telemetry: &mut Option<ExecutionTelemetry>,
) {
if usage_stream_telemetry
.as_ref()
.and_then(|telemetry| telemetry.ttfb_ms)
.is_some()
{
if !maybe_capture_first_stream_event_telemetry(
stream_started_at,
event_observed_at,
upstream_telemetry,
usage_stream_telemetry,
) {
return;
}
if !tracker.observe_chunk(chunk) {
let Some(telemetry) = usage_stream_telemetry.as_ref() else {
return;
}
let telemetry = first_visible_text_telemetry(stream_started_at, upstream_telemetry);
};
state.usage_runtime.record_stream_started(
state.data.as_ref(),
lifecycle_seed,
status_code,
Some(&telemetry),
Some(telemetry),
);
*usage_stream_telemetry = Some(telemetry);
}
fn build_terminal_stream_telemetry(
@@ -2068,14 +1911,14 @@ fn should_probe_success_failover_before_stream(headers: &BTreeMap<String, String
}
async fn probe_local_stream_success_failover_text<R>(
buffered_frames: &mut VecDeque<StreamFrame>,
buffered_frames: &mut VecDeque<ObservedStreamFrame>,
lines: &mut FramedRead<R, LinesCodec>,
) -> Result<Option<String>, GatewayError>
where
R: tokio::io::AsyncRead + Unpin,
{
while let Some(frame) = read_next_frame(lines).await? {
let probe_text = match &frame.payload {
while let Some(observed_frame) = read_next_observed_stream_frame(lines).await? {
let probe_text = match &observed_frame.frame.payload {
StreamFramePayload::Data { chunk_b64, text } => {
match decode_stream_data_chunk(chunk_b64.as_deref(), text.as_deref()) {
Ok(chunk) if !chunk.is_empty() => {
@@ -2087,7 +1930,7 @@ where
StreamFramePayload::Error { .. } | StreamFramePayload::Eof { .. } => None,
StreamFramePayload::Headers { .. } | StreamFramePayload::Telemetry { .. } => None,
};
buffered_frames.push_back(frame);
buffered_frames.push_back(observed_frame);
if probe_text.is_some() {
return Ok(probe_text);
}
@@ -2536,8 +2379,6 @@ async fn execute_stream_from_frame_stream(
let mut prefetched_inspection_body = Vec::new();
let mut prefetched_telemetry: Option<ExecutionTelemetry> = None;
let mut prefetched_usage_telemetry: Option<ExecutionTelemetry> = None;
let mut prefetched_client_visible_text_tracker =
ClientVisibleTextTracker::new(plan.client_api_format.as_str());
let mut reached_eof = false;
let mut sync_json_stream_bridge_active = false;
if skip_direct_finalize_prefetch {
@@ -2597,7 +2438,7 @@ async fn execute_stream_from_frame_stream(
} else {
next_stream_frame(&mut buffered_frames, &mut lines).await
};
let Some(frame) = (match next_frame_result {
let Some(observed_frame) = (match next_frame_result {
Ok(frame) => frame,
Err(err) => {
let failure = build_stream_failure_report(
@@ -2625,8 +2466,15 @@ async fn execute_stream_from_frame_stream(
reached_eof = true;
break;
};
match frame.payload {
let frame_observed_at = observed_frame.observed_at;
match observed_frame.frame.payload {
StreamFramePayload::Data { chunk_b64, text } => {
maybe_capture_first_stream_event_telemetry(
stream_started_at,
frame_observed_at,
prefetched_telemetry.as_ref(),
&mut prefetched_usage_telemetry,
);
let chunk =
match decode_stream_data_chunk(chunk_b64.as_deref(), text.as_deref()) {
Ok(chunk) => chunk,
@@ -2759,19 +2607,6 @@ async fn execute_stream_from_frame_stream(
"text/event-stream".to_string(),
);
stream_terminal_summary = outcome.terminal_summary;
if prefetched_usage_telemetry
.as_ref()
.and_then(|telemetry| telemetry.ttfb_ms)
.is_none()
&& prefetched_client_visible_text_tracker
.observe_chunk(&outcome.sse_body)
{
prefetched_usage_telemetry =
Some(first_visible_text_telemetry(
stream_started_at,
prefetched_telemetry.as_ref(),
));
}
prefetched_body.extend_from_slice(&outcome.sse_body);
prefetched_chunks.push(Bytes::from(outcome.sse_body));
sync_json_stream_bridge_active = true;
@@ -2871,18 +2706,6 @@ async fn execute_stream_from_frame_stream(
normalized_chunk
};
if !rewritten_chunk.is_empty() {
if prefetched_usage_telemetry
.as_ref()
.and_then(|telemetry| telemetry.ttfb_ms)
.is_none()
&& prefetched_client_visible_text_tracker
.observe_chunk(&rewritten_chunk)
{
prefetched_usage_telemetry = Some(first_visible_text_telemetry(
stream_started_at,
prefetched_telemetry.as_ref(),
));
}
prefetched_body.extend_from_slice(&rewritten_chunk);
prefetched_chunks.push(Bytes::from(rewritten_chunk));
}
@@ -2987,7 +2810,6 @@ async fn execute_stream_from_frame_stream(
let prefetched_chunks_for_body = prefetched_chunks;
let sync_json_stream_bridge_active_for_report = sync_json_stream_bridge_active;
let initial_telemetry = prefetched_telemetry;
let initial_client_visible_text_tracker = prefetched_client_visible_text_tracker;
let initial_reached_eof = reached_eof;
let direct_stream_finalize_kind_owned = direct_stream_finalize_kind;
let candidate_started_unix_secs_for_report = candidate_started_unix_secs;
@@ -3077,7 +2899,6 @@ async fn execute_stream_from_frame_stream(
let mut client_stream_completion_tracker = ClientVisibleStreamCompletionTracker::default();
let mut client_visible_stream_completed =
client_stream_completion_tracker.observe_chunk(&prefetched_body_for_report);
let mut client_visible_text_tracker = initial_client_visible_text_tracker;
let mut usage_stream_telemetry: Option<ExecutionTelemetry> = initial_usage_telemetry;
let mut telemetry: Option<ExecutionTelemetry> = initial_telemetry;
let reached_eof = initial_reached_eof;
@@ -3289,19 +3110,27 @@ async fn execute_stream_from_frame_stream(
break;
}
};
let Some(frame) = next_frame else {
let Some(observed_frame) = next_frame else {
if tx.is_closed() {
downstream_dropped = true;
}
break;
};
let frame_elapsed_ms = stream_started_at_for_report
.elapsed()
.as_millis()
.min(u128::from(u64::MAX)) as u64;
let frame_observed_at = observed_frame.observed_at;
let frame_elapsed_ms =
stream_elapsed_ms_at(stream_started_at_for_report, frame_observed_at);
last_upstream_frame_elapsed_ms.store(frame_elapsed_ms, Ordering::Relaxed);
match frame.payload {
match observed_frame.frame.payload {
StreamFramePayload::Data { chunk_b64, text } => {
maybe_record_first_stream_event_started(
&state_for_report,
&lifecycle_seed_for_report,
status_code,
stream_started_at_for_report,
frame_observed_at,
telemetry.as_ref(),
&mut usage_stream_telemetry,
);
if sync_json_stream_bridge_active_for_report {
continue;
}
@@ -3450,16 +3279,6 @@ async fn execute_stream_from_frame_stream(
);
downstream_dropped = true;
} else {
maybe_record_first_visible_text_started(
&state_for_report,
&lifecycle_seed_for_report,
status_code,
&mut client_visible_text_tracker,
rewritten_chunk.as_ref(),
stream_started_at_for_report,
telemetry.as_ref(),
&mut usage_stream_telemetry,
);
client_visible_stream_completed |= client_stream_completion_tracker
.observe_chunk(rewritten_chunk.as_ref());
client_stream_bytes.fetch_add(rewritten_chunk_len, Ordering::Relaxed);
@@ -3606,16 +3425,6 @@ async fn execute_stream_from_frame_stream(
);
downstream_dropped = true;
} else {
maybe_record_first_visible_text_started(
&state_for_report,
&lifecycle_seed_for_report,
status_code,
&mut client_visible_text_tracker,
rewritten_chunk.as_ref(),
stream_started_at_for_report,
telemetry.as_ref(),
&mut usage_stream_telemetry,
);
client_visible_stream_completed |= client_stream_completion_tracker
.observe_chunk(rewritten_chunk.as_ref());
client_stream_bytes
@@ -3687,16 +3496,6 @@ async fn execute_stream_from_frame_stream(
);
downstream_dropped = true;
} else {
maybe_record_first_visible_text_started(
&state_for_report,
&lifecycle_seed_for_report,
status_code,
&mut client_visible_text_tracker,
flushed_chunk.as_ref(),
stream_started_at_for_report,
telemetry.as_ref(),
&mut usage_stream_telemetry,
);
client_visible_stream_completed |= client_stream_completion_tracker
.observe_chunk(flushed_chunk.as_ref());
client_stream_bytes.fetch_add(flushed_chunk_len, Ordering::Relaxed);
@@ -4130,7 +3929,7 @@ mod tests {
stream_requires_observed_terminal_event, stream_terminal_summary_missing_observed_finish,
stream_terminal_summary_missing_observed_finish_with_requirement,
stream_terminal_summary_represents_failure_with_requirement,
ClientVisibleStreamCompletionTracker, ClientVisibleTextTracker,
ClientVisibleStreamCompletionTracker,
};
use crate::control::GatewayControlDecision;
use crate::tunnel::{tunnel_protocol, TunnelProxyConn};
@@ -5025,35 +4824,6 @@ mod tests {
));
}
#[test]
fn client_visible_text_tracker_ignores_control_role_and_tool_events_until_text() {
let mut tracker = ClientVisibleTextTracker::new("openai:chat");
assert!(!tracker.observe_chunk(
b": keepalive\n\n\
data: {\"choices\":[{\"delta\":{\"role\":\"assistant\"}}]}\n\n"
));
assert!(!tracker.observe_chunk(
b"data: {\"choices\":[{\"delta\":{\"tool_calls\":[{\"index\":0,\"function\":{\"arguments\":\"{}\"}}]}}]}\n\n"
));
assert!(
tracker.observe_chunk(b"data: {\"choices\":[{\"delta\":{\"content\":\"Hello\"}}]}\n\n")
);
assert!(!tracker
.observe_chunk(b"data: {\"choices\":[{\"delta\":{\"content\":\"again\"}}]}\n\n"));
}
#[test]
fn client_visible_text_tracker_detects_responses_output_text_across_chunks() {
let mut tracker = ClientVisibleTextTracker::new("openai:responses");
assert!(!tracker.observe_chunk(
b"event: response.created\n\
data: {\"type\":\"response.created\"}\n\n\
event: response.output_text.delta\n\
data: {\"type\":\"response.output_text.delta\",\"delta\":\""
));
assert!(tracker.observe_chunk(b"Hi\"}\n\n"));
}
#[tokio::test]
async fn sse_body_stream_emits_initial_and_periodic_keepalive_without_business_chunks() {
let (_tx, rx) = mpsc::channel::<Result<Bytes, std::io::Error>>(1);
@@ -6199,16 +5969,16 @@ mod tests {
}
#[tokio::test]
async fn execute_execution_runtime_stream_waits_for_visible_text_before_recording_first_byte() {
async fn execute_execution_runtime_stream_records_first_stream_event_before_visible_text() {
let listener = crate::test_support::bind_loopback_listener()
.await
.expect("listener should bind");
let addr = listener.local_addr().expect("local addr should resolve");
let role_only_seen = Arc::new(Notify::new());
let first_event_seen = Arc::new(Notify::new());
let release_text = Arc::new(Notify::new());
let text_seen = Arc::new(Notify::new());
let release_terminal = Arc::new(Notify::new());
let role_only_seen_for_route = Arc::clone(&role_only_seen);
let first_event_seen_for_route = Arc::clone(&first_event_seen);
let release_text_for_route = Arc::clone(&release_text);
let text_seen_for_route = Arc::clone(&text_seen);
let release_terminal_for_route = Arc::clone(&release_terminal);
@@ -6216,7 +5986,7 @@ mod tests {
let app = Router::new().route(
"/v1/execute/stream",
any(move |_request: Request| {
let role_only_seen = Arc::clone(&role_only_seen_for_route);
let first_event_seen = Arc::clone(&first_event_seen_for_route);
let release_text = Arc::clone(&release_text_for_route);
let text_seen = Arc::clone(&text_seen_for_route);
let release_terminal = Arc::clone(&release_terminal_for_route);
@@ -6226,12 +5996,12 @@ mod tests {
b"{\"type\":\"headers\",\"payload\":{\"kind\":\"headers\",\"status_code\":200,\"headers\":{\"content-type\":\"text/event-stream\"}}}\n",
));
yield Ok::<Bytes, Infallible>(Bytes::from_static(
b"{\"type\":\"data\",\"payload\":{\"kind\":\"data\",\"text\":\"data: {\\\"choices\\\":[{\\\"delta\\\":{\\\"role\\\":\\\"assistant\\\"}}]}\\n\\n\"}}\n",
b"{\"type\":\"data\",\"payload\":{\"kind\":\"data\",\"text\":\"\"}}\n",
));
yield Ok::<Bytes, Infallible>(Bytes::from_static(
b"{\"type\":\"telemetry\",\"payload\":{\"kind\":\"telemetry\",\"telemetry\":{\"ttfb_ms\":11,\"elapsed_ms\":12}}}\n",
));
role_only_seen.notify_one();
first_event_seen.notify_one();
release_text.notified().await;
yield Ok::<Bytes, Infallible>(Bytes::from_static(
b"{\"type\":\"data\",\"payload\":{\"kind\":\"data\",\"text\":\"data: {\\\"choices\\\":[{\\\"delta\\\":{\\\"content\\\":\\\"hello\\\"}}]}\\n\\n\"}}\n",
@@ -6275,8 +6045,8 @@ mod tests {
})
.with_execution_runtime_override_base_url(format!("http://{addr}"));
let plan = ExecutionPlan {
request_id: "req-live-stream-visible-text".into(),
candidate_id: Some("cand-live-stream-visible-text".into()),
request_id: "req-live-stream-first-event".into(),
candidate_id: Some("cand-live-stream-first-event".into()),
provider_name: Some("openai".into()),
provider_id: "prov-1".into(),
endpoint_id: "ep-1".into(),
@@ -6318,7 +6088,7 @@ mod tests {
let response = execute_execution_runtime_stream(
&state,
plan,
"trace-live-stream-visible-text",
"trace-live-stream-first-event",
&decision,
"openai_chat_stream",
None,
@@ -6331,55 +6101,37 @@ mod tests {
.expect("execution should succeed")
.expect("execution should return a client response");
role_only_seen.notified().await;
first_event_seen.notified().await;
let deadline = tokio::time::Instant::now() + Duration::from_secs(2);
let role_only_usage = loop {
let first_event_usage = loop {
let usage = usage_repository
.find_by_request_id("req-live-stream-visible-text")
.find_by_request_id("req-live-stream-first-event")
.await
.expect("usage should read");
if usage.as_ref().is_some_and(|usage| {
usage.status == "streaming" && usage.response_time_ms == Some(12)
usage.status == "streaming"
&& usage.response_time_ms == Some(12)
&& usage.first_byte_time_ms.is_some()
}) {
break usage.expect("streaming usage should exist");
}
assert!(
tokio::time::Instant::now() < deadline,
"usage should record streaming status"
"usage should record first byte on the first upstream stream event"
);
tokio::time::sleep(Duration::from_millis(10)).await;
};
assert_eq!(role_only_usage.first_byte_time_ms, None);
assert_eq!(role_only_usage.response_time_ms, Some(12));
assert!(first_event_usage.first_byte_time_ms.is_some());
assert_eq!(first_event_usage.response_time_ms, Some(12));
release_text.notify_one();
text_seen.notified().await;
let deadline = tokio::time::Instant::now() + Duration::from_secs(2);
let visible_text_usage = loop {
let usage = usage_repository
.find_by_request_id("req-live-stream-visible-text")
.await
.expect("usage should read");
if usage
.as_ref()
.is_some_and(|usage| usage.first_byte_time_ms.is_some())
{
break usage.expect("visible text usage should exist");
}
assert!(
tokio::time::Instant::now() < deadline,
"usage should record first byte only after visible text"
);
tokio::time::sleep(Duration::from_millis(10)).await;
};
assert!(visible_text_usage.first_byte_time_ms.unwrap_or(0) >= 12);
release_terminal.notify_one();
let body = to_bytes(response.into_body(), usize::MAX)
.await
.expect("response body should read");
let text = String::from_utf8(body.to_vec()).expect("response body should be utf8");
assert!(text.contains("\"role\":\"assistant\""));
assert!(text.contains("\"content\":\"hello\""));
server.abort();