use crate::session::{ActivityState, Session};
use std::collections::HashMap;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
const POST_TIMEOUT: Duration = Duration::from_secs(5);
pub struct Webhook {
target: String,
origin: String,
token: String,
warned: Arc<AtomicBool>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum State {
Busy,
Waiting,
Asking,
Stopped,
}
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,
}
}
impl Webhook {
pub fn new(target: String, origin: String, token: String) -> Webhook {
Webhook {
target,
origin,
token,
warned: Arc::new(AtomicBool::new(false)),
}
}
pub fn crossed<'a>(
&self,
before: &[Session],
after: &'a [Session],
) -> Vec<(&'static str, &'a Session)> {
let was: HashMap<String, State> = before.iter().map(|s| (s.key(), state_of(s))).collect();
let mut out = Vec::new();
for session in after {
let event = match state_of(session) {
State::Waiting => "waiting",
State::Asking => "asking",
State::Busy | State::Stopped => continue,
};
if was.get(&session.key()) == Some(&State::Busy) {
out.push((event, session));
}
}
out
}
pub fn post(&self, event: &'static str, session: &Session) {
let link = format!("{}/session/{}", self.origin, session.session_id);
let link = match self.token.is_empty() {
true => link,
false => format!("{link}?t={}", self.token),
};
let body = serde_json::json!({
"event": event,
"session_id": session.session_id,
"project": session.label_source,
"title": session.title,
"url": link,
})
.to_string();
let target = self.target.clone();
let warned = Arc::clone(&self.warned);
let session_id = session.session_id.clone();
std::thread::spawn(move || {
let agent: ureq::Agent = ureq::Agent::config_builder()
.timeout_global(Some(POST_TIMEOUT))
.build()
.into();
let sent = agent
.post(&target)
.header("Content-Type", "application/json")
.send(body.as_str());
crate::elog::event(
"notify",
"post",
serde_json::json!({
"event": event,
"session": session_id,
"ok": sent.is_ok(),
"status": sent.as_ref().map(|r| r.status().as_u16()).unwrap_or(0),
}),
);
if let Err(why) = sent
&& !warned.swap(true, Ordering::Relaxed)
{
eprintln!("cctop: notify POST to {target} failed: {why}");
}
});
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::pricing::Provider;
fn hook() -> Webhook {
Webhook::new(
"http://127.0.0.1:9/hook".into(),
"http://127.0.0.1:7777".into(),
"abc123".into(),
)
}
fn session(id: &str, running: bool, state: ActivityState) -> Session {
let mut s = Session::new(Provider::Claude, id.into());
s.activity_state = state;
if running {
s.process = Some(crate::proc::ProcInfo::default());
}
s
}
#[test]
fn only_the_crossing_into_waiting_posts() {
let hook = hook();
let busy = vec![session("a", true, ActivityState::Working)];
let waiting = vec![session("a", true, ActivityState::WaitingForInput)];
assert!(hook.crossed(&[], &busy).is_empty());
assert!(hook.crossed(&[], &waiting).is_empty());
let crossed = hook.crossed(&busy, &waiting);
assert_eq!(crossed.len(), 1);
assert_eq!(crossed[0].0, "waiting");
assert_eq!(crossed[0].1.session_id, "a");
assert!(hook.crossed(&waiting, &waiting).is_empty());
assert_eq!(hook.crossed(&waiting, &busy).len(), 0);
assert_eq!(hook.crossed(&busy, &waiting).len(), 1);
}
#[test]
fn asking_is_its_own_event_but_the_same_turn_does_not_fire_twice() {
let hook = hook();
let busy = vec![session("a", true, ActivityState::Working)];
let asking = vec![session("a", true, ActivityState::Asking)];
let waiting = vec![session("a", true, ActivityState::WaitingForInput)];
let crossed = hook.crossed(&busy, &asking);
assert_eq!(crossed[0].0, "asking");
assert!(hook.crossed(&asking, &waiting).is_empty());
}
#[test]
fn a_stopped_or_absent_previous_state_is_not_an_edge() {
let hook = hook();
let stopped = vec![session("a", false, ActivityState::Working)];
let waiting = vec![session("a", true, ActivityState::WaitingForInput)];
assert!(hook.crossed(&stopped, &waiting).is_empty());
assert!(hook.crossed(&[], &waiting).is_empty());
}
}