Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
101 changes: 53 additions & 48 deletions src/app.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1498,70 +1498,55 @@ impl App {
}

fn completion_notification_stats(&self) -> Option<String> {
let message = self.chat_state.chat.messages.iter().rev().find(|msg| {
msg.role == crate::session::types::MessageRole::Assistant && msg.is_complete
})?;

if let (Some(t0), Some(t1), Some(tn)) = (message.t0_ms, message.t1_ms, message.tn_ms) {
let output_tokens = message.output_tokens.or(message.token_count).unwrap_or(0);
let ttft_ms = t1.saturating_sub(t0);
let decode_ms = message.duration_ms.unwrap_or_else(|| tn.saturating_sub(t1));
let total_ms = ttft_ms.saturating_add(decode_ms);

let total_sec = total_ms as f64 / 1000.0;
let tokens_per_sec = if decode_ms > 0 && output_tokens > 0 {
(output_tokens as f64) / (decode_ms as f64 / 1000.0)
} else {
0.0
};

return Some(format!("{:.1}s | {:.0}t/s", total_sec, tokens_per_sec));
}

if let (Some(token_count), Some(duration_ms)) = (message.token_count, message.duration_ms) {
let duration_sec = duration_ms as f64 / 1000.0;
let tokens_per_sec = if duration_ms > 0 {
(token_count as f64) / (duration_ms as f64 / 1000.0)
} else {
0.0
};

return Some(format!("{:.1}s | {:.0}t/s", duration_sec, tokens_per_sec));
}

None
Self::completion_notification_stats_for_chat(&self.chat_state.chat)
}

fn completion_notification_stats_for_chat(chat: &Chat) -> Option<String> {
let message = chat.messages.iter().rev().find(|msg| {
msg.role == crate::session::types::MessageRole::Assistant && msg.is_complete
})?;

let format_tps = |precomputed: Option<f64>, tokens: usize, decode_ms: u64| -> Option<f64> {
if let Some(tps) = precomputed {
if tps.is_finite() && tps > 0.0 {
return Some(tps);
}
}
// OpenCode inter-token: (n - 1) / duration; need >1 token.
if decode_ms == 0 || tokens < 2 {
return None;
}
let tps = ((tokens - 1) as f64) / (decode_ms as f64 / 1000.0);
if tps.is_finite() && tps > 0.0 {
Some(tps)
} else {
None
}
};

if let (Some(t0), Some(t1), Some(tn)) = (message.t0_ms, message.t1_ms, message.tn_ms) {
let output_tokens = message.output_tokens.or(message.token_count).unwrap_or(0);
let ttft_ms = t1.saturating_sub(t0);
let decode_ms = message.duration_ms.unwrap_or_else(|| tn.saturating_sub(t1));
let total_ms = ttft_ms.saturating_add(decode_ms);

let total_sec = total_ms as f64 / 1000.0;
let tokens_per_sec = if decode_ms > 0 && output_tokens > 0 {
(output_tokens as f64) / (decode_ms as f64 / 1000.0)
} else {
0.0
};

return Some(format!("{:.1}s | {:.0}t/s", total_sec, tokens_per_sec));
if let Some(tokens_per_sec) =
format_tps(message.tokens_per_sec, output_tokens, decode_ms)
{
return Some(format!("{:.1}s | {:.0}t/s", total_sec, tokens_per_sec));
}
return Some(format!("{:.1}s", total_sec));
}

if let (Some(token_count), Some(duration_ms)) = (message.token_count, message.duration_ms) {
let duration_sec = duration_ms as f64 / 1000.0;
let tokens_per_sec = if duration_ms > 0 {
(token_count as f64) / (duration_ms as f64 / 1000.0)
} else {
0.0
};

return Some(format!("{:.1}s | {:.0}t/s", duration_sec, tokens_per_sec));
if let Some(tokens_per_sec) =
format_tps(message.tokens_per_sec, token_count, duration_ms)
{
return Some(format!("{:.1}s | {:.0}t/s", duration_sec, tokens_per_sec));
}
return Some(format!("{:.1}s", duration_sec));
}

None
Expand Down Expand Up @@ -2117,6 +2102,9 @@ impl App {
fn set_session_retry_status(&mut self, session_id: &str, status: Option<StreamingRetryStatus>) {
self.ensure_session_view_state(session_id);
if let Some(state) = self.session_view_states.get_mut(session_id) {
if state.retry_status == status {
return;
}
state.retry_status = status;
}
}
Expand Down Expand Up @@ -8940,6 +8928,11 @@ impl App {
crate::llm::ChunkMessage::Metrics { .. } => true,
crate::llm::ChunkMessage::ToolCalls(tool_calls) => {
self.set_session_retry_status(session_id, None);
// Close the generation sample as a tool-calls finish (excluded from
// TPS) and pause timing for the duration of tool execution.
if let Some(chat) = self.chat_for_session_mut(session_id) {
chat.end_generation_for_tool_calls();
}
self.add_tool_calls_to_session(session_id, tool_calls);
true
}
Expand Down Expand Up @@ -13015,6 +13008,9 @@ mod tests {
.chat
.add_message(crate::session::types::Message::incomplete(""));
app.chat_state.chat.begin_streaming_turn();
app.chat_state
.chat
.prepare_streaming_token_counter("test-model");

let (sender, receiver) = tokio::sync::mpsc::unbounded_channel();
sender
Expand Down Expand Up @@ -13139,14 +13135,23 @@ mod tests {
app.overlay_focus = OverlayFocus::SessionsDialog;
app.sessions_dialog_state.dialog.show();
app.refresh_sessions_dialog();
app.sessions_dialog_live_dirty = false;
let probe_before = app.last_sessions_dialog_metadata_probe;

app.base_focus = BaseFocus::Chat;
app.chat_state
.chat
.add_message(crate::session::types::Message::incomplete(""));
app.chat_state.chat.begin_streaming_turn();
// Warm tiktoken before probing so load cost is not charged to the
// sessions-dialog probe interval during process_streaming_chunks.
app.chat_state
.chat
.prepare_streaming_token_counter("test-model");

// Reset probe after any setup cost (e.g. tiktoken warm-up) so the
// assertion only covers process_streaming_chunks itself.
app.sessions_dialog_live_dirty = false;
app.last_sessions_dialog_metadata_probe = std::time::Instant::now();
let probe_before = app.last_sessions_dialog_metadata_probe;

let (sender, receiver) = tokio::sync::mpsc::unbounded_channel();
sender
Expand Down
1 change: 1 addition & 0 deletions src/persistence/conversions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -200,6 +200,7 @@ impl TryFrom<Message> for SessionMessage {
output_tokens: msg
.output_tokens
.and_then(|v| if v > 0 { Some(v as usize) } else { None }),
tokens_per_sec: None,

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Persist the precomputed TPS value

Every saved and reloaded assistant message receives None here because the persistence message/schema does not carry tokens_per_sec. For multi-step turns, the fallback cannot reconstruct the same value: output_tokens includes reasoning and tool-call-ending steps, while the aggregate TPS excludes those contributions. Consequently reopening a session changes the displayed throughput, so this field needs to be round-tripped through persistence.

Useful? React with 👍 / 👎.

model: msg.model.clone(),
provider: msg.provider.clone(),
local_image_paths,
Expand Down
5 changes: 5 additions & 0 deletions src/session/types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -174,6 +174,9 @@ pub struct Message {
pub t1_ms: Option<u64>,
pub tn_ms: Option<u64>,
pub output_tokens: Option<usize>,
/// Precomputed tokens/s (OpenCode inter-token aggregate). Prefer over
/// recomputing `output_tokens / duration_ms`.
pub tokens_per_sec: Option<f64>,
pub model: Option<String>,
pub provider: Option<String>,
pub local_image_paths: Vec<String>,
Expand Down Expand Up @@ -205,6 +208,7 @@ impl Message {
t1_ms: None,
tn_ms: None,
output_tokens: None,
tokens_per_sec: None,
model: None,
provider: None,
local_image_paths: Vec::new(),
Expand Down Expand Up @@ -252,6 +256,7 @@ impl Message {
t1_ms: None,
tn_ms: None,
output_tokens: None,
tokens_per_sec: None,
model: None,
provider: None,
local_image_paths: Vec::new(),
Expand Down
Loading
Loading