acorn-lib 0.3.2

ACORN library
use crate::io::api::webhooks::store::{EnqueueStatus, OperationQueue, OperationState};
use crate::io::database::backend::params;
use crate::io::ApiResult;
use crate::util::constants::app::MAX_WEBHOOK_OPERATION_ATTEMPTS;
use acorn_core::util::to_rfc3339;
use color_eyre::eyre::eyre;
use jiff::{SignedDuration, Timestamp};
use nanoid::nanoid;
use std::env::temp_dir;

fn queue() -> OperationQueue {
    OperationQueue::from(temp_dir().join(format!("acorn-bot-{}.db", nanoid!())))
}

fn settlement_failure() -> ApiResult<()> {
    Err(eyre!("injected settlement failure"))
}

#[test]
fn test_cancel_and_explicit_replay_are_state_safe() {
    let queue = queue();
    let operation_key = "repository-process:1:2:3";
    queue.enqueue("delivery-1", operation_key, "{}").unwrap();
    assert!(queue.cancel(operation_key).unwrap());
    assert_eq!(queue.state(operation_key).unwrap(), Some(OperationState::Cancelled));
    assert_eq!(queue.counts().unwrap().cancelled, 1);
    assert!(!queue.cancel(operation_key).unwrap());
    assert!(queue.replay(operation_key).unwrap());
    assert_eq!(queue.state(operation_key).unwrap(), Some(OperationState::Queued));
    assert!(!queue.replay(operation_key).unwrap());
}

#[test]
fn test_claims_retries_and_retains_terminal_failure() {
    let queue = queue();
    queue.enqueue("delivery-1", "work-item-check:1:2:3", "{}").unwrap();
    let mut operation = Some(queue.claim_next(SignedDuration::from_mins(5)).unwrap().unwrap());
    for expected_attempt in 1..=MAX_WEBHOOK_OPERATION_ATTEMPTS {
        let claim = operation.take().unwrap();
        assert_eq!(claim.attempts(), expected_attempt);
        let operation_key = claim.operation_key().to_string();
        let state = queue.fail(claim, "temporary").unwrap();
        if expected_attempt < MAX_WEBHOOK_OPERATION_ATTEMPTS {
            assert_eq!(state, OperationState::RetryableFailure);
            queue
                .with_connection(|connection| {
                    connection
                        .execute(
                            "UPDATE webhook_operations SET available_at = ? WHERE operation_key = ?",
                            params![
                                to_rfc3339(Timestamp::now().checked_sub(SignedDuration::from_secs(1)).unwrap()),
                                operation_key
                            ],
                        )
                        .map(|_| ())
                        .map_err(|why| eyre!("{why}"))
                })
                .unwrap();
            operation = queue.claim_next(SignedDuration::from_mins(5)).unwrap();
        } else {
            assert_eq!(state, OperationState::TerminalFailure);
            assert_eq!(queue.state(&operation_key).unwrap(), Some(OperationState::TerminalFailure));
        }
    }
    assert_eq!(queue.counts().unwrap().failed, 1);
}

#[test]
fn test_duplicate_delivery_and_operation_are_idempotent() {
    let queue = queue();
    assert_eq!(
        queue.enqueue("delivery-1", "mr-check:1:2:abc", r#"{"event":"normalized"}"#).unwrap(),
        EnqueueStatus::Inserted
    );
    assert_eq!(
        queue.enqueue("delivery-1", "mr-check:1:2:abc", r#"{"event":"changed"}"#).unwrap(),
        EnqueueStatus::DuplicateDelivery
    );
    assert_eq!(
        queue.enqueue("delivery-2", "mr-check:1:2:abc", r#"{"event":"normalized"}"#).unwrap(),
        EnqueueStatus::DuplicateOperation
    );
    assert_eq!(queue.counts().unwrap().deduplicated, 1);
}

#[test]
fn test_older_attempt_cannot_settle_recovered_operation() {
    let queue = queue();
    queue.enqueue("delivery-1", "mr-check:1:2:abc", "{}").unwrap();
    let stale = queue.claim_next(SignedDuration::from_mins(5)).unwrap().unwrap();
    queue
        .with_connection(|connection| {
            connection
                .execute(
                    "UPDATE webhook_operations SET claimed_at = ? WHERE operation_key = ?",
                    params![
                        to_rfc3339(Timestamp::now().checked_sub(SignedDuration::from_mins(10)).unwrap()),
                        stale.operation_key()
                    ],
                )
                .map(|_| ())
                .map_err(|why| eyre!("{why}"))
        })
        .unwrap();
    let current = queue.claim_next(SignedDuration::from_mins(5)).unwrap().unwrap();
    let operation_key = current.operation_key().to_string();
    assert_eq!(current.attempts(), 2);
    assert!(queue.succeed(stale).unwrap_err().to_string().contains("stale or already settled"));
    assert_eq!(queue.state(&operation_key).unwrap(), Some(OperationState::Running));
    queue.succeed(current).unwrap();
    assert_eq!(queue.state(&operation_key).unwrap(), Some(OperationState::Succeeded));
}

#[test]
#[cfg(not(feature = "duckdb"))]
fn test_schema_does_not_store_raw_webhook_bodies() {
    let queue = queue();
    queue.migrate().unwrap();
    queue
        .with_connection(|connection| {
            connection
                .query_row("SELECT sql FROM sqlite_master WHERE name = 'webhook_deliveries'", params![], |row| {
                    row.get::<_, String>(0)
                })
                .map_err(|why| eyre!("{why}"))
                .map(|sql| assert!(!sql.contains("raw_body")))
        })
        .unwrap();
}

#[test]
fn test_settlement_failure_leaves_claim_running_for_recovery() {
    let path = temp_dir().join(format!("acorn-bot-{}.db", nanoid!()));
    let queue = OperationQueue::from(path.clone()).with_settlement_hook(settlement_failure);
    queue.enqueue("delivery-1", "mr-check:1:2:abc", "{}").unwrap();
    let claim = queue.claim_next(SignedDuration::from_mins(5)).unwrap().unwrap();
    let operation_key = claim.operation_key().to_string();
    assert_eq!(queue.succeed(claim).unwrap_err().to_string(), "injected settlement failure");
    assert_eq!(queue.state(&operation_key).unwrap(), Some(OperationState::Running));
    let recovered = OperationQueue::from(path).claim_next(SignedDuration::from_secs(-1)).unwrap().unwrap();
    assert_eq!(recovered.attempts(), 2);
}

#[test]
fn test_stale_running_operation_is_recovered_after_restart() {
    let path = temp_dir().join(format!("acorn-bot-{}.db", nanoid!()));
    let queue = OperationQueue::from(path.clone());
    queue.enqueue("delivery-1", "mr-check:1:2:abc", "{}").unwrap();
    let claimed = queue.claim_next(SignedDuration::from_mins(5)).unwrap().unwrap();
    queue
        .with_connection(|connection| {
            connection
                .execute(
                    "UPDATE webhook_operations SET claimed_at = ? WHERE operation_key = ?",
                    params![
                        to_rfc3339(Timestamp::now().checked_sub(SignedDuration::from_mins(10)).unwrap()),
                        claimed.operation_key()
                    ],
                )
                .map(|_| ())
                .map_err(|why| eyre!("{why}"))
        })
        .unwrap();
    let restarted = OperationQueue::from(path);
    assert_eq!(restarted.claim_next(SignedDuration::from_mins(5)).unwrap().unwrap().attempts(), 2);
}