Menu
akurai-tasks
publicLatest change 1c1f5692cd8d8026234139f030e7432c298e73b9 - Log Tasks webhook delivery outcomes by Ólafur Búi Ólafsson
use super::{field, object, string, Value};
use std::io::{Read, Write};
use std::net::{SocketAddr, TcpStream};
use std::sync::mpsc::{self, Sender};
use std::thread;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
#[derive(Clone)]
pub struct TransitionWebhook {
tx: Sender<WebhookEvent>,
url: String,
token: String,
}
struct WebhookEvent {
url: String,
token: String,
payload: String,
}
impl TransitionWebhook {
pub fn from_env() -> Option<Self> {
let url = std::env::var("TASKS_WEBHOOK_URL")
.unwrap_or_else(|_| "http://100.88.0.9:4173/api/ingest-tasks-webhook".into());
let token = std::env::var("TASKS_WEBHOOK_TOKEN").ok()?;
if url.trim().is_empty() || token.trim().is_empty() || parse_http_url(&url).is_none() {
return None;
}
let (tx, rx) = mpsc::channel();
thread::Builder::new()
.name("tasks-webhook".into())
.spawn(move || {
while let Ok(event) = rx.recv() {
for attempt in 0..4 {
match post_webhook(&event) {
Ok(status) if (200..300).contains(&status) => {
eprintln!("tasks-webhook delivered status={status}");
break;
}
Ok(status) => {
eprintln!(
"tasks-webhook rejected status={status} attempt={}",
attempt + 1
);
}
Err(()) => {
eprintln!("tasks-webhook failed attempt={}", attempt + 1);
}
}
if attempt == 3 {
break;
}
thread::sleep(Duration::from_millis(100u64 << attempt));
}
}
})
.ok()?;
Some(Self { tx, url, token })
}
pub fn enqueue(&self, previous: &Value, result: &Value, actor: &str) {
let Some(id) = field(result, "id") else {
return;
};
let Some(project) = field(result, "project") else {
return;
};
let Some(title) = field(result, "title") else {
return;
};
let Some(previous_state) = field(previous, "state") else {
return;
};
let Some(state) = field(result, "state") else {
return;
};
let Some(revision) = result.get("revision").and_then(Value::as_i64) else {
return;
};
let payload = object(vec![
("event", string("work_transition")),
("id", string(id)),
("project", string(project)),
("title", string(title)),
("previousState", string(previous_state)),
("state", string(state)),
("actor", string(actor)),
("revision", Value::Int(revision)),
("at", Value::Int(unix_now())),
])
.to_json();
let _ = self.tx.send(WebhookEvent {
url: self.url.clone(),
token: self.token.clone(),
payload,
});
}
}
fn unix_now() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|duration| duration.as_secs() as i64)
.unwrap_or_default()
}
fn parse_http_url(url: &str) -> Option<(String, u16, String)> {
let rest = url.strip_prefix("http://")?;
let (authority, path) = rest.split_once('/').unwrap_or((rest, ""));
let (host, port) = if let Some((host, port)) = authority.rsplit_once(':') {
(host, port.parse().ok()?)
} else {
(authority, 80)
};
(!host.is_empty()).then(|| (host.to_string(), port, format!("/{path}")))
}
fn post_webhook(event: &WebhookEvent) -> Result<u16, ()> {
let (host, port, path) = parse_http_url(&event.url).ok_or(())?;
let address: SocketAddr = format!("{host}:{port}").parse().map_err(|_| ())?;
let mut stream =
TcpStream::connect_timeout(&address, Duration::from_secs(5)).map_err(|_| ())?;
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.map_err(|_| ())?;
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.map_err(|_| ())?;
let request = format!(
"POST {path} HTTP/1.1\r\nHost: {host}\r\nAuthorization: Bearer {}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
event.token,
event.payload.len(),
event.payload
);
stream.write_all(request.as_bytes()).map_err(|_| ())?;
let mut response = [0u8; 1024];
let length = stream.read(&mut response).map_err(|_| ())?;
let status = std::str::from_utf8(&response[..length])
.ok()
.and_then(|value| value.lines().next())
.and_then(|line| line.split_whitespace().nth(1))
.and_then(|value| value.parse::<u16>().ok())
.ok_or(())?;
Ok(status)
}