kcode-kennedy-stepped-turn-runtime 0.1.0

Watermark-bounded mailbox sequencing and stepped-turn session driving
Documentation
use super::*;

fn queued(key: &str, delivery: usize) -> QueuedAdmission<usize> {
    QueuedAdmission {
        key: key.to_owned(),
        recorded_at: "2026-08-17T00:00:00Z".to_owned(),
        admission: PendingTurnAdmission::User {
            text: format!("text-{key}"),
            metadata: serde_json::json!({}),
        },
        delivery,
    }
}

#[test]
fn fifo_and_queued_duplicate_keep_original_sequence() {
    let mailbox = Mailbox::new();
    let sender = mailbox.sender();
    assert_eq!(
        sender.push(queued("a", 10)),
        PushResult::Queued { sequence: 1 }
    );
    assert_eq!(
        sender.push(queued("b", 20)),
        PushResult::Queued { sequence: 2 }
    );
    assert_eq!(
        sender.push(queued("a", 99)),
        PushResult::Duplicate { sequence: 1 }
    );

    let drained = mailbox.drain_through(mailbox.watermark());
    let observed: Vec<_> = drained
        .iter()
        .map(|entry| (entry.sequence, entry.item.key.as_str(), entry.item.delivery))
        .collect();
    assert_eq!(observed, [(1, "a", 10), (2, "b", 20)]);
}

#[test]
fn captured_watermark_excludes_later_push() {
    let mailbox = Mailbox::new();
    let sender = mailbox.sender();
    assert_eq!(
        sender.push(queued("a", 1)),
        PushResult::Queued { sequence: 1 }
    );
    let watermark = mailbox.watermark();
    assert_eq!(
        sender.push(queued("b", 2)),
        PushResult::Queued { sequence: 2 }
    );

    let first = mailbox.drain_through(watermark);
    assert_eq!(
        first
            .iter()
            .map(|entry| (entry.sequence, entry.item.delivery))
            .collect::<Vec<_>>(),
        [(1, 1)]
    );
    let later = mailbox.drain_through(mailbox.watermark());
    assert_eq!(
        later
            .iter()
            .map(|entry| (entry.sequence, entry.item.delivery))
            .collect::<Vec<_>>(),
        [(2, 2)]
    );
}

#[test]
fn privately_drained_duplicate_and_restoration_preserve_order() {
    let mailbox = Mailbox::new();
    let sender = mailbox.sender();
    sender.push(queued("a", 1));
    sender.push(queued("b", 2));
    let drained = mailbox.drain_through(mailbox.watermark());

    assert_eq!(
        sender.push(queued("a", 99)),
        PushResult::Duplicate { sequence: 1 }
    );
    assert_eq!(
        sender.push(queued("c", 3)),
        PushResult::Queued { sequence: 3 }
    );
    mailbox.restore_front(drained);

    let restored = mailbox.drain_through(mailbox.watermark());
    let observed: Vec<_> = restored
        .iter()
        .map(|entry| (entry.sequence, entry.item.key.as_str(), entry.item.delivery))
        .collect();
    assert_eq!(observed, [(1, "a", 1), (2, "b", 2), (3, "c", 3)]);
}

#[test]
fn exact_acknowledgement_permits_key_at_later_sequence() {
    let mailbox = Mailbox::new();
    let sender = mailbox.sender();
    assert_eq!(
        sender.push(queued("key", 1)),
        PushResult::Queued { sequence: 1 }
    );
    let first = mailbox.drain_through(mailbox.watermark()).pop().unwrap();
    mailbox.acknowledge(first.sequence, &first.item.key);

    assert_eq!(
        sender.push(queued("key", 2)),
        PushResult::Queued { sequence: 2 }
    );
    let second = mailbox.drain_through(mailbox.watermark()).pop().unwrap();
    assert_eq!((second.sequence, second.item.delivery), (2, 2));
    mailbox.acknowledge(second.sequence, &second.item.key);
}

#[tokio::test(flavor = "multi_thread")]
async fn coalesced_notify_and_concurrent_producers_preserve_sequences_and_payloads() {
    let mailbox = Mailbox::new();
    let sender = mailbox.sender();
    sender.push(queued("first", 100));
    sender.push(queued("second", 200));
    mailbox.notified().await;

    const N: usize = 32;
    let mut tasks = Vec::new();
    for delivery in 0..N {
        let producer = sender.clone();
        tasks.push(tokio::spawn(async move {
            let result = producer.push(queued(&format!("key-{delivery}"), delivery));
            let PushResult::Queued { sequence } = result else {
                panic!("unique producer key was duplicated");
            };
            (sequence, delivery)
        }));
    }

    let mut assigned = Vec::new();
    for task in tasks {
        assigned.push(task.await.unwrap());
    }
    assigned.sort_unstable();
    assert_eq!(
        assigned.iter().map(|entry| entry.0).collect::<Vec<_>>(),
        (3..=N as u64 + 2).collect::<Vec<_>>()
    );

    let drained = mailbox.drain_through(mailbox.watermark());
    let observed: Vec<_> = drained
        .iter()
        .map(|entry| (entry.sequence, entry.item.delivery))
        .collect();
    let mut expected = vec![(1, 100), (2, 200)];
    expected.extend(assigned);
    assert_eq!(observed, expected);
}