use crate::session::{ActivityState, Session};
use std::collections::{HashMap, VecDeque};
use std::time::{Duration, Instant};
pub const MARK_FOR: Duration = Duration::from_secs(30);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum State {
Busy,
Waiting,
Asking,
Stopped,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Reason {
NeedsInput,
Asking,
Stopped,
}
#[derive(Debug, Clone)]
pub struct Rang {
pub key: String,
pub label: String,
pub reason: Reason,
pub at: Instant,
}
const MAX_RANG: usize = 16;
#[derive(Default)]
pub struct Notifier {
pub enabled: bool,
pub webhook: Option<String>,
pub link_base: Option<String>,
watched: HashMap<String, (State, String)>,
recent: VecDeque<Rang>,
quota_limited: HashMap<String, bool>,
}
impl Notifier {
pub fn new(enabled: bool) -> Self {
Notifier {
enabled,
webhook: std::env::var("CCTOP_NOTIFY_URL")
.ok()
.filter(|url| !url.is_empty()),
..Default::default()
}
}
pub fn observe(&mut self, sessions: &[Session]) {
self.recent.retain(|r| r.at.elapsed() < MARK_FOR);
let mut crossed: Vec<Rang> = Vec::new();
let mut posts: Vec<(&'static str, &Session)> = Vec::new();
let mut next = HashMap::with_capacity(self.watched.len());
for session in sessions.iter().filter(|s| s.is_running()) {
let key = session.key();
let state = state_of(session);
let label = session.display_label().to_string();
if let Some((State::Busy, _)) = self.watched.remove(&key)
&& let Some(reason) = reason_for(state)
{
crossed.push(Rang {
key: key.clone(),
label: label.clone(),
reason,
at: Instant::now(),
});
let event = match reason {
Reason::NeedsInput => Some("waiting"),
Reason::Asking => Some("asking"),
Reason::Stopped => None,
};
if let Some(event) = event {
posts.push((event, session));
}
}
if state == State::Busy {
self.recent.retain(|r| r.key != key);
}
next.insert(key, (state, label));
}
for (key, (state, label)) in self.watched.drain() {
if state == State::Busy {
crossed.push(Rang {
key,
label,
reason: Reason::Stopped,
at: Instant::now(),
});
}
}
self.watched = next;
if let Some(target) = &self.webhook {
for (event, session) in posts {
self.post_session(target.clone(), event, session);
}
}
if !self.enabled || crossed.is_empty() {
return;
}
let extra = crossed.len() - 1;
let rang = crossed.swap_remove(0);
ring(&desktop_text(&rang, extra));
self.recent.push_back(rang);
self.recent.extend(crossed);
while self.recent.len() > MAX_RANG {
self.recent.pop_front();
}
}
#[cfg(test)]
pub(crate) fn record_for_test(&mut self, rang: Rang) {
self.recent.push_back(rang);
}
pub fn unanswered(&self, selected: Option<&str>) -> Option<&Rang> {
self.recent
.iter()
.find(|r| r.at.elapsed() < MARK_FOR && Some(r.key.as_str()) != selected)
.or_else(|| self.recent.iter().rev().find(|r| r.at.elapsed() < MARK_FOR))
}
pub fn rang_recently(&self, key: &str) -> bool {
self.recent
.iter()
.any(|r| r.key == key && r.at.elapsed() < MARK_FOR)
}
pub fn footer(&self, selected: Option<&str>) -> Option<String> {
let rang = self.unanswered(selected)?;
if Some(rang.key.as_str()) == selected {
return None;
}
let secs = rang.at.elapsed().as_secs();
let ago = if secs < 60 {
format!("{secs}s")
} else {
format!("{}m", secs / 60)
};
let extra = self
.recent
.iter()
.filter(|r| r.at.elapsed() < MARK_FOR && r.key != rang.key)
.count();
let more = match extra {
0 => String::new(),
n => format!(" · +{n} more"),
};
Some(format!(
"Bell: ◉ {} · {} · {ago} ago{more} · b jumps to it",
rang.label,
match rang.reason {
Reason::NeedsInput => "waiting for input",
Reason::Asking => "needs permission",
Reason::Stopped => "stopped",
}
))
}
pub fn observe_quota(&mut self, quota: &crate::quota::Quota) -> Vec<String> {
let mut freed: Vec<String> = Vec::new();
for (provider, profiles) in [("claude", "a.claude), ("codex", "a.codex)] {
for profile in profiles {
let crate::quota::ProviderStatus::Ok(q) = &profile.status else {
continue;
};
let mut watch = |key: String, limited: bool, what: String| {
if self.quota_limited.insert(key, limited) == Some(true) && !limited {
freed.push(what);
}
};
watch(
format!("{provider}/{}", profile.profile),
q.limit_reached,
format!(
"{provider} · {}: the rate limit has lifted",
profile.profile
),
);
for window in &q.windows {
watch(
format!("{provider}/{}/{}", profile.profile, window.label),
window.pct >= 100,
format!(
"{provider} · {}: the {} window is open again",
profile.profile, window.label
),
);
}
}
}
freed
}
fn post_session(&self, target: String, event: &'static str, session: &Session) {
let mut body = serde_json::json!({
"event": event,
"session_id": session.session_id,
"project": session.label_source,
"title": session.title,
});
if let Some(base) = &self.link_base {
body["url"] = session_link(base, &session.session_id).into();
}
crate::serve::notify::post(
target,
body.to_string(),
serde_json::json!({ "event": event, "session": session.session_id }),
);
}
pub fn post_event(&self, event: &'static str, text: &str) {
let Some(target) = &self.webhook else {
return;
};
crate::serve::notify::post(
target.clone(),
serde_json::json!({ "event": event, "text": text }).to_string(),
serde_json::json!({ "event": event }),
);
}
}
fn session_link(base: &str, id: &str) -> String {
match base.split_once("/?") {
Some((origin, query)) => format!("{origin}/session/{id}?{query}"),
None => format!("{base}session/{id}"),
}
}
fn reason_for(state: State) -> Option<Reason> {
match state {
State::Asking => Some(Reason::Asking),
State::Waiting => Some(Reason::NeedsInput),
State::Busy | State::Stopped => None,
}
}
fn state_of(session: &Session) -> State {
if !session.is_running() {
return State::Stopped;
}
match session.activity_state {
ActivityState::Asking => State::Asking,
ActivityState::WaitingForInput => State::Waiting,
ActivityState::Working | ActivityState::ApiError => State::Busy,
}
}
fn desktop_text(rang: &Rang, extra: usize) -> String {
let what = match rang.reason {
Reason::NeedsInput => "is waiting for input",
Reason::Asking => "needs permission",
Reason::Stopped => "stopped",
};
let more = if extra > 0 {
format!(" (+{extra} more)")
} else {
String::new()
};
format!("cctop: {} {what}{more}", rang.label)
}
pub(crate) fn ring(text: &str) {
use std::io::Write;
if cfg!(test) {
return;
}
let mut out = std::io::stdout();
let _ = write!(out, "\x07\x1b]9;{}\x07", sanitize(text));
let _ = out.flush();
}
fn sanitize(text: &str) -> String {
text.chars().filter(|c| !c.is_control()).collect()
}
#[cfg(test)]
mod tests {
use super::*;
use crate::pricing::Provider;
fn session(id: &str, running: bool, state: ActivityState) -> Session {
let mut s = Session::new(Provider::Claude, id.into());
s.label_source = format!("/home/x/{id}");
s.abbrev_label = id.into();
s.activity_state = state;
if running {
s.process = Some(crate::proc::ProcInfo::default());
}
s
}
#[test]
fn a_held_permission_prompt_is_not_a_finished_turn() {
let mut n = Notifier {
enabled: true,
..Default::default()
};
n.observe(&[session("a", true, ActivityState::Working)]);
n.observe(&[session("a", true, ActivityState::Asking)]);
let rang = n.recent.back().expect("a blocked agent rings");
assert_eq!(rang.reason, Reason::Asking);
assert!(
desktop_text(rang, 0).contains("needs permission"),
"the bell says which kind of waiting it is: {}",
desktop_text(rang, 0)
);
let mut n = Notifier {
enabled: true,
..Default::default()
};
n.observe(&[session("a", true, ActivityState::Working)]);
n.observe(&[session("a", true, ActivityState::WaitingForInput)]);
assert_eq!(n.recent.back().map(|r| r.reason), Some(Reason::NeedsInput));
}
#[test]
fn answering_a_prompt_does_not_ring_again_when_the_turn_ends() {
let mut n = Notifier {
enabled: true,
..Default::default()
};
n.observe(&[session("a", true, ActivityState::Working)]);
n.observe(&[session("a", true, ActivityState::Asking)]);
n.recent.clear();
n.observe(&[session("a", true, ActivityState::WaitingForInput)]);
assert!(n.recent.is_empty(), "the same turn, reported twice");
}
#[test]
fn a_session_that_stops_working_rings_once_and_not_again() {
let mut n = Notifier {
enabled: true,
..Default::default()
};
let busy = vec![session("a", true, ActivityState::Working)];
n.observe(&busy);
assert!(n.recent.is_empty(), "still working, nothing to say");
let waiting = vec![session("a", true, ActivityState::WaitingForInput)];
n.observe(&waiting);
let first = n.recent.back().cloned().expect("the crossing rings");
assert_eq!(first.reason, Reason::NeedsInput);
n.observe(&waiting);
assert_eq!(
n.recent.back().map(|r| r.at),
Some(first.at),
"a session that is still waiting must not ring again"
);
}
#[test]
fn a_busy_session_that_disappears_rings_as_stopped() {
let mut n = Notifier {
enabled: true,
..Default::default()
};
n.observe(&[session("a", true, ActivityState::Working)]);
n.observe(&[session("a", false, ActivityState::Working)]);
let rang = n.recent.back().expect("an agent that exits is news");
assert_eq!(rang.reason, Reason::Stopped);
assert_eq!(rang.label, "a");
n.recent.clear();
n.observe(&[]);
assert!(n.recent.is_empty());
}
#[test]
fn a_session_first_seen_idle_never_rings_for_being_idle() {
let mut n = Notifier {
enabled: true,
..Default::default()
};
let waiting = vec![session("a", true, ActivityState::WaitingForInput)];
n.observe(&waiting);
n.observe(&waiting);
assert!(n.recent.is_empty());
}
#[test]
fn the_state_machine_stays_warm_while_notifications_are_off() {
let mut n = Notifier::default();
n.observe(&[session("a", true, ActivityState::Working)]);
n.enabled = true;
n.observe(&[session("a", true, ActivityState::WaitingForInput)]);
assert!(
!n.recent.is_empty(),
"the busy state seen before the toggle still counts"
);
}
#[test]
fn going_back_to_work_answers_the_bell() {
let mut n = Notifier {
enabled: true,
..Default::default()
};
n.observe(&[session("a", true, ActivityState::Working)]);
n.observe(&[session("a", true, ActivityState::WaitingForInput)]);
assert!(!n.recent.is_empty());
n.observe(&[session("a", true, ActivityState::Working)]);
assert!(n.recent.is_empty(), "answered, so stop naming it");
}
#[test]
fn the_footer_names_the_session_until_it_is_selected() {
let mut n = Notifier {
enabled: true,
..Default::default()
};
n.observe(&[session("a", true, ActivityState::Working)]);
n.observe(&[session("a", true, ActivityState::WaitingForInput)]);
let key = n.recent.back().unwrap().key.clone();
assert!(n.footer(None).is_some_and(|t| t.contains('a')));
assert!(n.footer(Some("claude:other")).is_some());
assert!(n.footer(Some(&key)).is_none());
assert!(n.rang_recently(&key));
assert!(!n.rang_recently("claude:other"));
}
#[test]
fn a_label_cannot_break_out_of_the_notification() {
let text = sanitize("proj\x07\x1b]0;pwned\x07");
assert!(!text.contains('\x07'));
assert!(!text.contains('\x1b'));
assert_eq!(text, "proj]0;pwned");
}
#[test]
fn one_bell_for_a_refresh_that_finishes_several_sessions() {
let rang = Rang {
key: "claude:a".into(),
label: "alpha".into(),
reason: Reason::NeedsInput,
at: Instant::now(),
};
assert_eq!(desktop_text(&rang, 0), "cctop: alpha is waiting for input");
assert_eq!(
desktop_text(&rang, 2),
"cctop: alpha is waiting for input (+2 more)"
);
}
#[test]
fn every_session_that_crossed_stays_findable() {
let mut n = Notifier {
enabled: true,
..Default::default()
};
n.observe(&[
session("a", true, ActivityState::Working),
session("b", true, ActivityState::Working),
]);
n.observe(&[
session("a", true, ActivityState::WaitingForInput),
session("b", true, ActivityState::Asking),
]);
assert!(n.rang_recently("claude:a"));
assert!(
n.rang_recently("claude:b"),
"the crossing the bell did not name lost its marker"
);
assert_eq!(n.unanswered(None).map(|r| r.key.as_str()), Some("claude:a"));
assert_eq!(
n.unanswered(Some("claude:a")).map(|r| r.key.as_str()),
Some("claude:b")
);
assert_eq!(
n.unanswered(Some("claude:b")).map(|r| r.key.as_str()),
Some("claude:a")
);
let footer = n.footer(None).expect("someone rang");
assert!(footer.contains('a'), "{footer}");
assert!(footer.contains("+1 more"), "{footer}");
}
#[test]
fn a_quota_window_rings_once_on_the_way_back() {
let mut n = Notifier::default();
let at = |pct: u32| crate::quota::Quota {
fetched: true,
claude: vec![crate::quota::ProfileQuota {
profile: "default".into(),
status: crate::quota::ProviderStatus::Ok(crate::quota::ProviderQuota {
plan: None,
windows: vec![crate::quota::Window {
label: "5h",
pct,
duration: None,
resets_at: None,
}],
limit_reached: false,
}),
source: crate::config::AccountSource::Directory,
}],
codex: Vec::new(),
};
assert!(n.observe_quota(&at(40)).is_empty());
assert!(
n.observe_quota(&at(100)).is_empty(),
"filling is not the edge"
);
let freed = n.observe_quota(&at(60));
assert_eq!(freed.len(), 1, "the crossing is the announcement");
assert!(freed[0].contains("5h"), "{freed:?}");
assert!(
n.observe_quota(&at(60)).is_empty(),
"a window that stays open stays quiet"
);
let mut gone = at(60);
gone.claude[0].status = crate::quota::ProviderStatus::RateLimited { retry_at: None };
assert!(n.observe_quota(&gone).is_empty());
assert!(
n.observe_quota(&at(60)).is_empty(),
"still nothing new to say"
);
}
#[test]
fn a_provider_level_limit_rings_on_lift() {
let mut n = Notifier::default();
let quota = |held: bool| crate::quota::Quota {
fetched: true,
claude: Vec::new(),
codex: vec![crate::quota::ProfileQuota {
profile: "default".into(),
status: crate::quota::ProviderStatus::Ok(crate::quota::ProviderQuota {
plan: None,
windows: Vec::new(),
limit_reached: held,
}),
source: crate::config::AccountSource::Directory,
}],
};
assert!(n.observe_quota("a(true)).is_empty());
let freed = n.observe_quota("a(false));
assert_eq!(freed.len(), 1);
assert!(freed[0].contains("limit has lifted"), "{freed:?}");
}
#[test]
fn a_session_link_keeps_the_token_in_the_query() {
assert_eq!(
session_link("https://x.trycloudflare.com/?t=abc", "s1"),
"https://x.trycloudflare.com/session/s1?t=abc"
);
assert_eq!(
session_link("http://127.0.0.1:7777/", "s1"),
"http://127.0.0.1:7777/session/s1"
);
}
}