harn-vm 0.10.29

Async bytecode virtual machine for the Harn programming language
Documentation
use std::sync::{Arc, Mutex};

use crate::agent_events::sinks::{event_log_flush_schedule, EventLogFlushSchedule};
use crate::agent_events::{
    clear_session_sinks, emit_event, flush_and_clear_session_sinks, flush_session_sinks,
    register_sink, register_wildcard_sink, session_external_sink_count, unregister_wildcard_sink,
    AgentEvent, AgentEventSink, AgentEventSinkError, AgentEventSinkFlush, EventLogSink, MultiSink,
    ToolCallStatus, ToolMutationStatus,
};
use crate::event_log::{AnyEventLog, EventLog, MemoryEventLog, SqliteEventLog, Topic};

struct FlushProbe {
    name: &'static str,
    failure: Option<AgentEventSinkError>,
    flushed: Arc<Mutex<Vec<&'static str>>>,
}

impl AgentEventSink for FlushProbe {
    fn handle_event(&self, _event: &AgentEvent) {}

    fn flush(&self) -> AgentEventSinkFlush<'_> {
        Box::pin(async move {
            self.flushed.lock().unwrap().push(self.name);
            self.failure.clone().map_or(Ok(()), Err)
        })
    }
}

fn flush_probe(
    name: &'static str,
    failure: Option<AgentEventSinkError>,
    flushed: &Arc<Mutex<Vec<&'static str>>>,
) -> Arc<dyn AgentEventSink> {
    Arc::new(FlushProbe {
        name,
        failure,
        flushed: flushed.clone(),
    })
}

#[test]
fn redacts_tool_payloads_before_append() {
    let log = Arc::new(AnyEventLog::Memory(MemoryEventLog::new(8)));
    let sink = EventLogSink::new(log.clone(), "s");
    sink.handle_event(&AgentEvent::ToolCall {
        session_id: "s".into(),
        tool_call_id: "tc-1".into(),
        tool_name: "http".into(),
        kind: None,
        status: ToolCallStatus::Pending,
        raw_input: serde_json::json!({
            "authorization": "Bearer raw-bearer-value",
            "url": "https://user:password@example.com/items?sig=raw-signature&ok=1"
        }),
        parsing: None,
        audit: None,
    });

    let topic = Topic::new("observability.agent_events.s").unwrap();
    let events = futures::executor::block_on(log.read_range(&topic, None, 8)).unwrap();
    assert_eq!(events.len(), 1);
    let persisted = serde_json::to_string(&events[0].1).unwrap();
    assert!(persisted.contains("[redacted]") || persisted.contains("%5Bredacted%5D"));
    for secret in ["raw-bearer-value", "user:password", "raw-signature"] {
        assert!(
            !persisted.contains(secret),
            "event-log sink appended secret {secret}: {persisted}"
        );
    }
}

#[test]
fn skips_text_parsing_candidates() {
    let log = Arc::new(AnyEventLog::Memory(MemoryEventLog::new(8)));
    let sink = EventLogSink::new(log.clone(), "s");
    sink.handle_event(&AgentEvent::ToolCall {
        session_id: "s".into(),
        tool_call_id: "text-cand-0".into(),
        tool_name: "read".into(),
        kind: None,
        status: ToolCallStatus::Pending,
        raw_input: serde_json::json!({}),
        parsing: Some(true),
        audit: None,
    });
    sink.handle_event(&AgentEvent::ToolCallUpdate {
        session_id: "s".into(),
        tool_call_id: "call-real".into(),
        tool_name: "read".into(),
        status: ToolCallStatus::Completed,
        raw_output: Some(serde_json::json!({"text": "ok"})),
        error: None,
        duration_ms: None,
        execution_duration_ms: None,
        error_category: None,
        mutation_status: ToolMutationStatus::Unknown,
        changed_paths: None,
        executor: None,
        parsing: None,
        raw_input: None,
        raw_input_partial: None,
        audit: None,
    });

    let topic = Topic::new("observability.agent_events.s").unwrap();
    let events = futures::executor::block_on(log.read_range(&topic, None, 8)).unwrap();
    assert_eq!(events.len(), 1);
    let persisted = serde_json::to_string(&events[0].1).unwrap();
    assert!(persisted.contains("call-real"));
    assert!(!persisted.contains("text-cand-0"));
}

#[tokio::test(flavor = "current_thread")]
async fn flush_waits_for_queued_appends_without_polling() {
    let log = Arc::new(AnyEventLog::Memory(MemoryEventLog::new(8)));
    let sink = EventLogSink::new(log.clone(), "flush-session");
    sink.handle_event(&AgentEvent::AgentMessageChunk {
        session_id: "flush-session".into(),
        content: "persist before replay".into(),
    });

    sink.flush().await.expect("flush queued event-log append");

    let topic = Topic::new("observability.agent_events.flush-session").unwrap();
    let events = log.read_range(&topic, None, 8).await.unwrap();
    assert_eq!(events.len(), 1);
    assert_eq!(
        events[0].1.payload["event"]["content"],
        "persist before replay"
    );
}

#[tokio::test(flavor = "current_thread")]
async fn session_flush_is_a_causal_registry_barrier() {
    let session_id = "registry-flush-session";
    let log = Arc::new(AnyEventLog::Memory(MemoryEventLog::new(8)));
    register_sink(session_id, EventLogSink::new(log.clone(), session_id));
    emit_event(&AgentEvent::AgentMessageChunk {
        session_id: session_id.into(),
        content: "causally durable".into(),
    });

    flush_session_sinks(session_id)
        .await
        .expect("flush registered sinks");

    let topic = Topic::new("observability.agent_events.registry-flush-session").unwrap();
    let events = log.read_range(&topic, None, 8).await.unwrap();
    assert_eq!(events.len(), 1);
    assert_eq!(events[0].1.payload["event"]["content"], "causally durable");
    clear_session_sinks(session_id);
}

#[tokio::test(flavor = "current_thread")]
async fn session_flush_drains_later_sinks_after_first_failure() {
    let session_id = "drain-all-session";
    let flushed = Arc::new(Mutex::new(Vec::new()));
    register_sink(
        session_id,
        flush_probe(
            "first",
            Some(AgentEventSinkError::new("first", "expected failure")),
            &flushed,
        ),
    );
    register_sink(session_id, flush_probe("second", None, &flushed));

    let error = flush_session_sinks(session_id)
        .await
        .expect_err("first sink failure must be reported after all sinks drain");

    assert_eq!(error.sink(), "first");
    assert_eq!(error.message(), "expected failure");
    assert_eq!(*flushed.lock().unwrap(), ["first", "second"]);
    clear_session_sinks(session_id);
}

#[tokio::test(flavor = "current_thread")]
async fn multi_sink_flush_drains_later_sinks_after_first_failure() {
    let flushed = Arc::new(Mutex::new(Vec::new()));
    let sinks = MultiSink::new();
    sinks.push(flush_probe(
        "first",
        Some(AgentEventSinkError::new("first", "expected failure")),
        &flushed,
    ));
    sinks.push(flush_probe("second", None, &flushed));

    let error = sinks
        .flush()
        .await
        .expect_err("first sink failure must be reported after all sinks drain");

    assert_eq!(error.sink(), "first");
    assert_eq!(*flushed.lock().unwrap(), ["first", "second"]);
}

#[tokio::test(flavor = "current_thread")]
async fn cancelled_flush_does_not_consume_sticky_append_error() {
    let temp = tempfile::tempdir().expect("event-log tempdir");
    let path = temp.path().join("events.sqlite");
    drop(SqliteEventLog::open(path.clone(), 8).expect("initialize writable event log"));
    let read_only = SqliteEventLog::open_read_only(path, 8).expect("open read-only event log");
    let sink = EventLogSink::new(
        Arc::new(AnyEventLog::Sqlite(read_only)),
        "cancelled-flush-session",
    );
    sink.handle_event(&AgentEvent::AgentMessageChunk {
        session_id: "cancelled-flush-session".into(),
        content: "append must fail".into(),
    });

    let cancelled = sink.enqueue_flush_for_test();
    drop(cancelled);
    let error = sink
        .flush()
        .await
        .expect_err("append error must remain visible after a cancelled flush waiter");

    assert_eq!(error.sink(), "event_log");
    assert!(error.message().contains("readonly") || error.message().contains("read-only"));
}

#[test]
fn sqlite_full_checkpoint_uses_the_blocking_pool() {
    let temp = tempfile::tempdir().expect("event-log tempdir");
    let sqlite = AnyEventLog::Sqlite(
        SqliteEventLog::open(temp.path().join("events.sqlite"), 8).expect("open sqlite event log"),
    );
    let memory = AnyEventLog::Memory(MemoryEventLog::new(8));

    assert_eq!(
        event_log_flush_schedule(&sqlite),
        EventLogFlushSchedule::BlockingPool
    );
    assert_eq!(
        event_log_flush_schedule(&memory),
        EventLogFlushSchedule::AsyncExecutor
    );
}

#[tokio::test(flavor = "current_thread")]
async fn flush_and_clear_removes_sinks_after_append_failure() {
    let session_id = "failed-append-clear-session";
    let temp = tempfile::tempdir().expect("event-log tempdir");
    let path = temp.path().join("events.sqlite");
    drop(SqliteEventLog::open(path.clone(), 8).expect("initialize writable event log"));
    let read_only = SqliteEventLog::open_read_only(path, 8).expect("open read-only event log");
    register_sink(
        session_id,
        EventLogSink::new(Arc::new(AnyEventLog::Sqlite(read_only)), session_id),
    );
    emit_event(&AgentEvent::AgentMessageChunk {
        session_id: session_id.into(),
        content: "append must fail before cleanup".into(),
    });

    let error = flush_and_clear_session_sinks(session_id)
        .await
        .expect_err("append failure must cross the cleanup barrier");

    assert_eq!(error.sink(), "event_log");
    assert!(error.message().contains("readonly") || error.message().contains("read-only"));
    assert_eq!(session_external_sink_count(session_id), 0);
}

#[tokio::test(flavor = "current_thread")]
async fn flush_drains_concurrent_producers_without_scheduler_polling() {
    const PRODUCERS: usize = 8;
    const EVENTS_PER_PRODUCER: usize = 32;

    let session_id = "concurrent-flush-session";
    let log = Arc::new(AnyEventLog::Memory(MemoryEventLog::new(
        PRODUCERS * EVENTS_PER_PRODUCER,
    )));
    let sink = EventLogSink::new(log.clone(), session_id);
    let mut producers = Vec::with_capacity(PRODUCERS);
    for producer in 0..PRODUCERS {
        let sink = sink.clone();
        producers.push(std::thread::spawn(move || {
            for event in 0..EVENTS_PER_PRODUCER {
                sink.handle_event(&AgentEvent::AgentMessageChunk {
                    session_id: session_id.into(),
                    content: format!("producer-{producer}-event-{event}"),
                });
            }
        }));
    }
    for producer in producers {
        producer.join().expect("event producer");
    }

    sink.flush().await.expect("flush all accepted events");

    let topic = Topic::new("observability.agent_events.concurrent-flush-session").unwrap();
    let events = log
        .read_range(&topic, None, PRODUCERS * EVENTS_PER_PRODUCER)
        .await
        .unwrap();
    let mut actual = events
        .iter()
        .map(|(_, record)| {
            record.payload["event"]["content"]
                .as_str()
                .expect("agent message content")
                .to_string()
        })
        .collect::<Vec<_>>();
    actual.sort();
    let mut expected = (0..PRODUCERS)
        .flat_map(|producer| {
            (0..EVENTS_PER_PRODUCER).map(move |event| format!("producer-{producer}-event-{event}"))
        })
        .collect::<Vec<_>>();
    expected.sort();
    assert_eq!(actual, expected);
}

#[tokio::test(flavor = "current_thread")]
async fn session_flush_includes_wildcard_sinks() {
    let session_id = "wildcard-flush-session";
    let log = Arc::new(AnyEventLog::Memory(MemoryEventLog::new(8)));
    let handle = register_wildcard_sink(EventLogSink::new(log.clone(), session_id));
    emit_event(&AgentEvent::AgentMessageChunk {
        session_id: session_id.into(),
        content: "persisted by wildcard".into(),
    });

    flush_session_sinks(session_id)
        .await
        .expect("flush wildcard event sink");

    let topic = Topic::new("observability.agent_events.wildcard-flush-session").unwrap();
    let events = log.read_range(&topic, None, 8).await.unwrap();
    assert_eq!(events.len(), 1);
    assert_eq!(
        events[0].1.payload["event"]["content"],
        "persisted by wildcard"
    );
    unregister_wildcard_sink(handle);
}