AkurAI Build
Menu

akurai-tasks

public

Latest 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)
}