atman-runtime 1.11.0

atman flow execution runtime: evaluator, tool dispatch, provider dispatch, executor, memory stores
Documentation
use crate::error::RuntimeError;
use crate::message::{Message, MessageRole};
use crate::tool::{BoxFut, Tier, Tool, ToolArgs, ToolCtx, ToolResult};
use crate::value::Value;

pub struct SessionPush;

impl Tool for SessionPush {
    fn name(&self) -> &str {
        "session.push"
    }

    fn tier(&self) -> Tier {
        Tier::Zero
    }

    fn description(&self) -> Option<&str> {
        Some(
            "Push a Message value into the current session's message history. \
             Use after dispatch_all to persist tool results so the next \
             llm.call(context: \"session\") call can see them. The message role \
             (user/assistant/tool/system) is preserved. Returns unit.",
        )
    }

    fn input_schema(&self) -> serde_json::Value {
        serde_json::json!({
            "type": "object",
            "properties": {
                "message": {
                    "type": "object",
                    "description": "The Message value to push (e.g. a tool_result from dispatch_all). Pass the value returned by dispatch_all directly."
                }
            },
            "required": ["message"]
        })
    }

    fn call<'a>(&'a self, args: ToolArgs, ctx: &'a ToolCtx) -> BoxFut<'a, ToolResult> {
        Box::pin(async move {
            let val = match args.named("message").or_else(|| args.positional(0).ok()) {
                Some(v) => v.clone(),
                None => {
                    return Err(RuntimeError::MissingArg("session.push: message".into()));
                }
            };
            let msgs = match val {
                Value::Message(m) => vec![m],
                Value::List(items) => items
                    .into_iter()
                    .filter_map(|v| match v {
                        Value::Message(m) => Some(m),
                        _ => None,
                    })
                    .collect(),
                other => {
                    return Err(RuntimeError::TypeMismatch {
                        expected: "message or list of message".into(),
                        actual: other.kind_name().into(),
                    });
                }
            };
            let Some(_handle) = &ctx.session_messages_handle else {
                return Err(RuntimeError::ToolFailed(
                    "session.push: no session messages handle available".into(),
                ));
            };
            let _compact_guard = match &ctx.compact_lock_handle {
                Some(lock) => Some(lock.lock().await),
                None => None,
            };
            for msg in msgs {
                let msg = crate::tools::tool_output::maybe_truncate_tool_message_with_budget(
                    &msg,
                    ctx.output_store.as_deref(),
                    ctx.tool_output_budget,
                );
                append_message_to_context(ctx, msg)?;
            }
            Ok(Value::Unit)
        })
    }
}

pub(crate) fn append_message_to_context(ctx: &ToolCtx, msg: Message) -> Result<(), RuntimeError> {
    if matches!(ctx.history_segment, crate::tool::HistorySegment::Root)
        && let Some(session) = ctx.session_runtime.as_ref()
    {
        let stream_flow_run_id = match msg.role {
            MessageRole::Assistant | MessageRole::Tool => ctx.flow_run_id.clone(),
            MessageRole::User | MessageRole::System => None,
        };
        session.append_message_with_stream_scope(
            msg,
            ctx.message_flow_run_id(),
            stream_flow_run_id,
        );
        return Ok(());
    }
    let Some(handle) = &ctx.session_messages_handle else {
        return Err(RuntimeError::ToolFailed(
            "session message context is unavailable".into(),
        ));
    };
    emit_message_event(ctx, &msg);
    let flow_run_id = match msg.role {
        MessageRole::Assistant | MessageRole::Tool => {
            ctx.flow_run_id.as_ref().map(|run_id| run_id.0.to_string())
        }
        MessageRole::User => ctx.message_flow_run_id().map(|run_id| run_id.0.to_string()),
        MessageRole::System => None,
    };
    if msg.origin != crate::message::MessageOrigin::Internal
        && msg.role != MessageRole::System
        && let Some(tx) = &ctx.stream_tx
    {
        let frame = match msg.role {
            MessageRole::Assistant => crate::stream::StreamFrame::AssistantMsg {
                flow_run_id,
                message: msg.clone(),
            },
            MessageRole::User | MessageRole::Tool => crate::stream::StreamFrame::ToolResultMsg {
                flow_run_id,
                message: msg.clone(),
            },
            MessageRole::System => unreachable!(),
        };
        let _ = tx.send(frame);
    }
    let mut messages = handle.lock().unwrap();
    messages.push(msg);
    Ok(())
}

fn emit_message_event(ctx: &ToolCtx, msg: &Message) {
    use crate::event::{Event, TurnId};
    let Some(sink) = &ctx.events else {
        return;
    };
    let turn_id = ctx.turn_id.clone().unwrap_or_else(TurnId::now);
    let flow_run_id = ctx.message_flow_run_id();
    let event = match msg.role {
        MessageRole::User => Event::UserMsg {
            turn_id,
            flow_run_id,
            message: msg.clone(),
        },
        MessageRole::Assistant => Event::AssistantMsg {
            turn_id,
            flow_run_id,
            message: msg.clone(),
        },
        MessageRole::Tool => Event::ToolResultMsg {
            turn_id,
            flow_run_id,
            message: msg.clone(),
        },
        MessageRole::System => Event::SystemMsg {
            turn_id,
            flow_run_id: ctx.message_flow_run_id(),
            message: msg.clone(),
        },
    };
    sink.emit(event);
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn spawned_context_does_not_rewrite_messages_outside_compaction() {
        let messages = std::sync::Arc::new(std::sync::Mutex::new(
            (0..100)
                .map(|index| {
                    Message::assistant_text(crate::event::TurnId::now(), format!("old-{index}"))
                })
                .collect(),
        ));
        let mut ctx = ToolCtx::new()
            .with_history_segment(crate::tool::HistorySegment::Spawned)
            .with_session_messages_handle(std::sync::Arc::clone(&messages));
        ctx.context_epoch_handle = Some(std::sync::Arc::new(std::sync::atomic::AtomicU64::new(0)));

        append_message_to_context(
            &ctx,
            Message::assistant_text(crate::event::TurnId::now(), "new"),
        )
        .unwrap();

        assert_eq!(messages.lock().unwrap().len(), 101);
        assert_eq!(messages.lock().unwrap()[0].text_concat(), "old-0");
        assert_eq!(ctx.context_epoch_seed().as_deref(), Some("generation:0"));
    }

    #[test]
    fn session_push_name_and_tier() {
        let tool = SessionPush;
        assert_eq!(tool.name(), "session.push");
        assert_eq!(tool.tier(), Tier::Zero);
        assert!(tool.description().is_some());
    }

    #[tokio::test]
    async fn session_push_scopes_live_tool_result_without_scoping_session_history() {
        use crate::event::{Event, FlowRunId, TurnId};
        use crate::message::{MessageOrigin, MessagePart};

        let session = std::sync::Arc::new(crate::session::Session::open_ephemeral());
        let run_id = FlowRunId::now();
        let mut stream_rx = session.stream_subscribe();
        let ctx = ToolCtx::new()
            .with_anchors(Some(TurnId::now()), Some(run_id.clone()), None)
            .with_events(session.sink().clone())
            .with_session_messages_handle(session.messages_handle())
            .with_session_runtime(session.clone());
        let message = Message {
            role: MessageRole::Tool,
            parts: vec![MessagePart::ToolResult {
                tool_use_id: "tu_1".into(),
                content: "done".into(),
                is_error: false,
            }],
            turn_id: TurnId::now(),
            origin: MessageOrigin::User,
        };

        SessionPush
            .call(
                ToolArgs {
                    positional: vec![Value::Message(message)],
                    named: Vec::new(),
                },
                &ctx,
            )
            .await
            .unwrap();

        let crate::stream::StreamFrame::ToolResultMsg { flow_run_id, .. } =
            stream_rx.recv().await.unwrap()
        else {
            panic!("tool result frame");
        };
        assert_eq!(flow_run_id.as_deref(), Some(run_id.0.to_string().as_str()));
        assert!(session.sink().snapshot().iter().any(|event| {
            matches!(
                event,
                Event::ToolResultMsg {
                    flow_run_id: None,
                    ..
                }
            )
        }));
        assert!(
            session
                .messages()
                .iter()
                .any(|message| message.role == MessageRole::Tool)
        );
    }

    #[test]
    fn internal_child_system_message_is_audited_without_becoming_a_live_transcript_frame() {
        use crate::context_plan::{
            ContextRecord, ContextRecordAuthority, ContextRecordBody, ContextRecordRetention,
        };
        use crate::event::{Event, FlowRunId, TurnId};

        let session = std::sync::Arc::new(crate::session::Session::open_ephemeral());
        let run_id = FlowRunId::now();
        let (stream_tx, mut stream_rx) = tokio::sync::broadcast::channel(8);
        let messages = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
        let ctx = ToolCtx::new()
            .with_anchors(Some(TurnId::now()), Some(run_id.clone()), None)
            .with_history_segment(crate::tool::HistorySegment::Spawned)
            .with_events(session.sink().clone())
            .with_session_messages_handle(messages)
            .with_stream_tx(stream_tx);
        let message = Message::context_record(
            TurnId::now(),
            ContextRecord::new(
                "handoff.parent",
                1,
                ContextRecordAuthority::Runtime,
                ContextRecordRetention::Latest,
                ContextRecordBody::text("delegated"),
            ),
        );

        append_message_to_context(&ctx, message).unwrap();

        assert!(matches!(
            stream_rx.try_recv(),
            Err(tokio::sync::broadcast::error::TryRecvError::Empty)
        ));
        assert!(session.sink().snapshot().iter().any(|event| {
            matches!(
                event,
                Event::SystemMsg {
                    flow_run_id: Some(owner),
                    ..
                } if owner == &run_id
            )
        }));
    }
}