Forge documentation
Library referenceRust

forge-collab

ANVIL Collaboration Contract: roles, sessions, task delegation, shared context, and interrupts for the Forge SDK

ANVIL Collaboration Contract: roles, sessions, task delegation, shared context, and interrupts for the Forge SDK

Package contract

FieldValue
Languagerust
Source version0.2.0
Manifestforge-rs/crates/forge-collab/Cargo.toml
Source files9
EvidenceSource reference; registry publication and runtime conformance are separate checks

Import boundary

use forge_collab;

Use a source checkout or your verified private registry. Manifest coordinates identify the package; they do not establish that a public registry release exists.

Crate boundary

The following entries are taken from src/lib.rs. Feature conditions in the exact source still apply.

pub mod capability;

pub mod context;

pub mod delegation;

pub mod error;

pub mod interrupt;

pub mod roles;

pub mod session;

pub mod types;

pub mod prelude;

pub use crate::capability::{match_task_to_agents, CapabilityAdvertiser};

pub use crate::context::{InMemorySharedContext, SharedContextContract};

pub use crate::delegation::{create_delegated_task, validate_task_constraints};

pub use crate::error::{CollabError, CollabResult};

pub use crate::interrupt::{create_interrupt, InterruptHandler};

pub use crate::roles::{CoordinatorContract, PeerContract, WorkerContract};

pub use crate::session::{SessionContract, SessionManager};

pub use crate::types::{
        AgentCapabilityProfile, CollaborationRole, CollaborationSession, ContextEntry,
        ContextVisibility, DelegatedTask, Interrupt, InterruptResponse, InterruptType,
        InterruptedState, SessionParticipant, SessionState, SessionTransition, TaskAcknowledgment,
        TaskConstraints, TaskPriority, TaskProgress, TaskResult, TaskStatus,
    };

Source reference

Download package reference JSON. Each original source file and generated declaration artifact has its own SHA-256 digest. Function bodies and constant values are omitted from downloads. These are source declaration inventories, not compiler-resolved rustdoc, TypeDoc, DocC, or Dokka output. Private modules can contain public declarations that are not reachable through the package boundary; consult the entry point before importing.

capability.rs

Read declaration text · 2 declaration entries

pub trait CapabilityAdvertiser: Send + Sync {
    /// Returns the agent's full capability profile.
    ///
    /// The profile includes the agent's supported roles, task types,
    /// available tools, current load, and maximum concurrency.
    ///
    /// # Returns
    ///
    /// An [`AgentCapabilityProfile`] describing the agent's capabilities.
    fn advertise_capabilities(&self) -> AgentCapabilityProfile;

    /// Returns `true` if this agent can handle the given task type.
    ///
    /// This is a convenience check that avoids constructing the full
    /// capability profile for simple task type matching.
    ///
    /// # Arguments
    ///
    /// * `task_type` - The task type to check (e.g., "code-review").
    ///
    /// # Returns
    ///
    /// `true` if the agent supports this task type.
    fn can_handle(&self, task_type: &str) -> bool;

    /// Returns the agent's current availability as a factor between 0.0
    /// (fully loaded) and 1.0 (completely idle).
    ///
    /// This is the inverse of `current_load` in the capability profile:
    /// `availability = 1.0 - current_load`.
    ///
    /// # Returns
    ///
    /// A value between 0.0 and 1.0 representing available capacity.
    fn current_availability(&self) -> f64;
}

pub fn match_task_to_agents(
    task: &DelegatedTask,
    agents: &[AgentCapabilityProfile],
) -> Vec<String>;

context.rs

Read declaration text · 8 declaration entries

#[async_trait::async_trait]
pub trait SharedContextContract: Send + Sync {
    /// Reads a context entry by key from the specified session.
    ///
    /// # Arguments
    ///
    /// * `session_id` - The session whose context to read from.
    /// * `key` - The key of the entry to read.
    ///
    /// # Returns
    ///
    /// The [`ContextEntry`] for the given key.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError::ContextKeyNotFound`] if the key does not
    /// exist in the session's context.
    ///
    /// Returns [`CollabError::SessionNotFound`] if the session does not
    /// exist.
    // ANVIL Spec section 11.4 -- Shared Context: context_read
    async fn context_read(&self, session_id: &str, key: &str) -> Result<ContextEntry, CollabError>;

    /// Writes a context entry to the specified session.
    ///
    /// If the key already exists, the entry is updated with the new value
    /// and its version is incremented.
    ///
    /// # Arguments
    ///
    /// * `session_id` - The session whose context to write to.
    /// * `entry` - The context entry to write.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError::SessionNotFound`] if the session does not
    /// exist.
    // ANVIL Spec section 11.4 -- Shared Context: context_write
    async fn context_write(&self, session_id: &str, entry: ContextEntry)
        -> Result<(), CollabError>;

    /// Lists all keys in the specified session's context.
    ///
    /// # Arguments
    ///
    /// * `session_id` - The session whose context keys to list.
    ///
    /// # Returns
    ///
    /// A `Vec` of key strings.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError::SessionNotFound`] if the session does not
    /// exist.
    // ANVIL Spec section 11.4 -- Shared Context: context_keys
    async fn context_keys(&self, session_id: &str) -> Result<Vec<String>, CollabError>;
}

#[derive(Debug, Clone)]
pub struct InMemorySharedContext {

}

pub fn new() -> Self;

pub fn create_session(&mut self, session_id: &str);

pub fn session_count(&self) -> usize;

pub fn write_entry(
        &mut self,
        session_id: &str,
        entry: ContextEntry,
    ) -> Result<(), CollabError>;

pub fn read_entry(&self, session_id: &str, key: &str) -> Result<ContextEntry, CollabError>;

pub fn list_keys(&self, session_id: &str) -> Result<Vec<String>, CollabError>;

delegation.rs

Read declaration text · 2 declaration entries

pub fn create_delegated_task(
    delegator_did: &str,
    task_type: &str,
    description: &str,
    input: serde_json::Value,
    priority: TaskPriority,
    constraints: TaskConstraints,
) -> DelegatedTask;

pub fn validate_task_constraints(constraints: &TaskConstraints) -> Result<(), CollabError>;

error.rs

Read declaration text · 2 declaration entries

#[derive(Debug, Error)]
pub enum CollabError {
    /// The specified collaboration session does not exist.
    ///
    /// This is returned when an operation references a session ID that has
    /// not been created or has already been garbage-collected after reaching
    /// a terminal state.
    ///
    /// See ANVIL Spec section 11.2 -- Collaboration Session.
    #[error("collaboration session '{session_id}' not found (see ANVIL Spec section 11.2)")]
    SessionNotFound {
        /// The session ID that was not found.
        session_id: String,
    },

    /// A session with the given ID is already in the Active state.
    ///
    /// Duplicate session activation is not permitted. Sessions must be
    /// dissolved before a new session with the same logical purpose is created.
    ///
    /// See ANVIL Spec section 11.2 -- Collaboration Session.
    #[error(
        "collaboration session '{session_id}' is already active (see ANVIL Spec section 11.2)"
    )]
    SessionAlreadyActive {
        /// The session ID that is already active.
        session_id: String,
    },

    /// The session exceeded its configured timeout duration.
    ///
    /// Sessions with a `timeout_seconds` value will automatically transition
    /// to the `TimedOut` terminal state when the deadline passes.
    ///
    /// See ANVIL Spec section 11.2 -- Collaboration Session.
    #[error("collaboration session '{session_id}' timed out (see ANVIL Spec section 11.2)")]
    SessionTimedOut {
        /// The session ID that timed out.
        session_id: String,
    },

    /// An invalid session state transition was attempted.
    ///
    /// The collaboration session state machine defines exactly which transitions
    /// are valid from each state. This error is returned when a transition
    /// violates the state machine rules.
    ///
    /// See ANVIL Spec section 11.2 -- Session State Machine.
    #[error("invalid session transition from '{from}' to '{to}' (see ANVIL Spec section 11.2)")]
    InvalidSessionTransition {
        /// The current session state.
        from: String,
        /// The target state that was attempted.
        to: String,
    },

    /// The agent is not a participant in the specified session.
    ///
    /// Only agents that have joined a session can perform operations within
    /// that session (read/write context, receive tasks, send interrupts).
    ///
    /// See ANVIL Spec section 11.2 -- Session Participation.
    #[error("agent '{agent_did}' is not a participant in session '{session_id}' (see ANVIL Spec section 11.2)")]
    NotAParticipant {
        /// The agent DID that is not a participant.
        agent_did: String,
        /// The session ID the agent attempted to interact with.
        session_id: String,
    },

    /// The specified delegated task does not exist.
    ///
    /// This is returned when an operation references a task ID that has
    /// not been created or has been completed and removed.
    ///
    /// See ANVIL Spec section 11.3 -- Delegated Task.
    #[error("delegated task '{task_id}' not found (see ANVIL Spec section 11.3)")]
    TaskNotFound {
        /// The task ID that was not found.
        task_id: String,
    },

    /// The task is already assigned to another agent.
    ///
    /// A delegated task can only be assigned to one worker at a time.
    /// The existing assignment must be cancelled before re-assignment.
    ///
    /// See ANVIL Spec section 11.3 -- Delegated Task.
    #[error("delegated task '{task_id}' is already assigned to '{assignee}' (see ANVIL Spec section 11.3)")]
    TaskAlreadyAssigned {
        /// The task ID that is already assigned.
        task_id: String,
        /// The DID of the agent currently assigned to the task.
        assignee: String,
    },

    /// The agent's capabilities do not match the task requirements.
    ///
    /// Task assignment validates that the target agent supports the required
    /// task type, tools, and has sufficient capacity.
    ///
    /// See ANVIL Spec section 11.6 -- Agent Capability Profile.
    #[error("capability mismatch for task type '{task_type}' and agent '{agent_did}': {reason} (see ANVIL Spec section 11.6)")]
    CapabilityMismatch {
        /// The task type that was not supported.
        task_type: String,
        /// The agent DID that lacks the capability.
        agent_did: String,
        /// A human-readable explanation of the mismatch.
        reason: String,
    },

    /// The specified key was not found in the shared context.
    ///
    /// Context reads return this when the key has not been written to
    /// the session's shared context store.
    ///
    /// See ANVIL Spec section 11.4 -- Shared Context.
    #[error("context key '{key}' not found (see ANVIL Spec section 11.4)")]
    ContextKeyNotFound {
        /// The key that was not found.
        key: String,
    },

    /// The agent does not have permission to access the specified context key.
    ///
    /// Context entries can have visibility restrictions based on session scope,
    /// role, or specific agent DID.
    ///
    /// See ANVIL Spec section 11.4 -- Shared Context.
    #[error("agent '{agent_did}' does not have permission to access context key '{key}' (see ANVIL Spec section 11.4)")]
    ContextPermissionDenied {
        /// The context key the agent attempted to access.
        key: String,
        /// The agent DID that was denied access.
        agent_did: String,
    },

    /// The interrupt was rejected by the target agent.
    ///
    /// Agents may reject interrupts if they are in a state that does not
    /// permit interruption (e.g., in a critical section) or if the interrupt
    /// type is not supported.
    ///
    /// See ANVIL Spec section 11.5 -- Interrupts.
    #[error("interrupt '{interrupt_id}' rejected: {reason} (see ANVIL Spec section 11.5)")]
    InterruptRejected {
        /// The interrupt ID that was rejected.
        interrupt_id: String,
        /// A human-readable explanation of the rejection.
        reason: String,
    },

    /// A task delegation operation failed.
    ///
    /// This is a general delegation failure covering scenarios such as
    /// no eligible workers, constraint violations, or internal errors
    /// during the delegation flow.
    ///
    /// See ANVIL Spec section 11.3 -- Delegated Task.
    #[error("delegation failed: {reason} (see ANVIL Spec section 11.3)")]
    DelegationFailed {
        /// A human-readable explanation of the failure.
        reason: String,
    },

    /// The agent attempted an action that violates its assigned role.
    ///
    /// For example, a Worker attempting to decompose tasks (a Coordinator
    /// action) or a Peer attempting to assign tasks.
    ///
    /// See ANVIL Spec section 11.1 -- Collaboration Roles.
    #[error("role violation: agent '{agent_did}' with role '{role}' cannot perform action '{action}' (see ANVIL Spec section 11.1)")]
    RoleViolation {
        /// The agent DID that violated its role.
        agent_did: String,
        /// The role the agent holds.
        role: String,
        /// The action the agent attempted.
        action: String,
    },
}

pub type CollabResult<T> = Result<T, CollabError>;

interrupt.rs

Read declaration text · 2 declaration entries

#[async_trait::async_trait]
pub trait InterruptHandler: Send + Sync {
    /// Handles an incoming interrupt.
    ///
    /// The handler inspects the interrupt type and priority, then decides
    /// whether to acknowledge it. The response includes the agent's current
    /// state if applicable (e.g., for suspendable tasks).
    ///
    /// # Arguments
    ///
    /// * `interrupt` - The interrupt to handle.
    ///
    /// # Returns
    ///
    /// An [`InterruptResponse`] indicating acknowledgment and optional state.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError::InterruptRejected`] if the interrupt cannot
    /// be handled (e.g., the agent is in a critical section).
    // ANVIL Spec section 11.5 -- Interrupt Handler: on_interrupt
    async fn on_interrupt(&self, interrupt: &Interrupt) -> Result<InterruptResponse, CollabError>;

    /// Called when the agent must be preempted by a higher-priority task.
    ///
    /// The agent should save its current state and prepare for the current
    /// task to be suspended or aborted.
    ///
    /// # Arguments
    ///
    /// * `interrupt` - The preemption interrupt.
    ///
    /// # Returns
    ///
    /// An [`InterruptedState`] snapshot of the agent's state at the moment
    /// of preemption.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError::InterruptRejected`] if preemption is not
    /// possible.
    // ANVIL Spec section 11.5 -- Interrupt Handler: on_preempt
    async fn on_preempt(&self, interrupt: &Interrupt) -> Result<InterruptedState, CollabError>;
}

pub fn create_interrupt(
    interrupt_type: InterruptType,
    source_did: &str,
    priority: TaskPriority,
    payload: serde_json::Value,
) -> Interrupt;

lib.rs

Read declaration text · 17 declaration entries

pub mod capability;

pub mod context;

pub mod delegation;

pub mod error;

pub mod interrupt;

pub mod roles;

pub mod session;

pub mod types;

pub mod prelude;

pub use crate::capability::{match_task_to_agents, CapabilityAdvertiser};

pub use crate::context::{InMemorySharedContext, SharedContextContract};

pub use crate::delegation::{create_delegated_task, validate_task_constraints};

pub use crate::error::{CollabError, CollabResult};

pub use crate::interrupt::{create_interrupt, InterruptHandler};

pub use crate::roles::{CoordinatorContract, PeerContract, WorkerContract};

pub use crate::session::{SessionContract, SessionManager};

pub use crate::types::{
        AgentCapabilityProfile, CollaborationRole, CollaborationSession, ContextEntry,
        ContextVisibility, DelegatedTask, Interrupt, InterruptResponse, InterruptType,
        InterruptedState, SessionParticipant, SessionState, SessionTransition, TaskAcknowledgment,
        TaskConstraints, TaskPriority, TaskProgress, TaskResult, TaskStatus,
    };

roles.rs

Read declaration text · 3 declaration entries

#[async_trait::async_trait]
pub trait CoordinatorContract: Send + Sync {
    /// Decomposes a high-level task into smaller sub-tasks.
    ///
    /// The coordinator analyzes the input task and produces a list of
    /// sub-tasks that can be individually assigned to workers. The sub-tasks
    /// should collectively cover the scope of the original task.
    ///
    /// # Arguments
    ///
    /// * `task` - The task to decompose.
    ///
    /// # Returns
    ///
    /// A `Vec` of sub-tasks derived from the original task.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError::DelegationFailed`] if the task cannot be
    /// decomposed (e.g., it is already atomic or invalid).
    // ANVIL Spec section 11.1 -- Coordinator: decompose_task
    async fn decompose_task(&self, task: &DelegatedTask)
        -> Result<Vec<DelegatedTask>, CollabError>;

    /// Assigns a task to a specific worker agent.
    ///
    /// The coordinator sends the task to the identified worker and receives
    /// an acknowledgment indicating acceptance or rejection.
    ///
    /// # Arguments
    ///
    /// * `task` - The task to assign.
    /// * `worker_did` - The OAS DID of the target worker.
    ///
    /// # Returns
    ///
    /// A [`TaskAcknowledgment`] from the worker.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError::CapabilityMismatch`] if the worker cannot
    /// handle the task type, or [`CollabError::DelegationFailed`] if the
    /// assignment fails for other reasons.
    // ANVIL Spec section 11.1 -- Coordinator: assign_task
    async fn assign_task(
        &self,
        task: &DelegatedTask,
        worker_did: &str,
    ) -> Result<TaskAcknowledgment, CollabError>;

    /// Aggregates results from multiple worker tasks into a single result.
    ///
    /// After all sub-tasks have completed, the coordinator merges their
    /// results into a unified response for the original task.
    ///
    /// # Arguments
    ///
    /// * `results` - The results from all completed sub-tasks.
    ///
    /// # Returns
    ///
    /// A single aggregated [`TaskResult`].
    ///
    /// # Errors
    ///
    /// Returns [`CollabError::DelegationFailed`] if the results cannot be
    /// meaningfully aggregated (e.g., no results, conflicting outputs).
    // ANVIL Spec section 11.1 -- Coordinator: aggregate_results
    async fn aggregate_results(&self, results: &[TaskResult]) -> Result<TaskResult, CollabError>;

    /// Handles a failure reported by a worker.
    ///
    /// When a worker fails to complete a task, the coordinator decides how
    /// to recover: reassign the task, mark it as failed, or abort the
    /// session.
    ///
    /// # Arguments
    ///
    /// * `task_id` - The ID of the failed task.
    /// * `worker_did` - The DID of the worker that failed.
    /// * `error` - A description of the failure.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError::DelegationFailed`] if recovery is not possible.
    // ANVIL Spec section 11.1 -- Coordinator: handle_worker_failure
    async fn handle_worker_failure(
        &self,
        task_id: &str,
        worker_did: &str,
        error: &str,
    ) -> Result<(), CollabError>;
}

#[async_trait::async_trait]
pub trait WorkerContract: Send + Sync {
    /// Called when a task is delegated to this worker.
    ///
    /// The worker inspects the task and decides whether to accept or reject
    /// it based on its capabilities and current load.
    ///
    /// # Arguments
    ///
    /// * `task` - The delegated task to evaluate.
    ///
    /// # Returns
    ///
    /// A [`TaskAcknowledgment`] indicating acceptance or rejection.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError`] if the acknowledgment cannot be produced.
    // ANVIL Spec section 11.1 -- Worker: on_task_delegated
    async fn on_task_delegated(
        &self,
        task: &DelegatedTask,
    ) -> Result<TaskAcknowledgment, CollabError>;

    /// Reports progress on an in-flight task.
    ///
    /// Workers should report progress periodically so the coordinator can
    /// monitor execution and detect stalls.
    ///
    /// # Arguments
    ///
    /// * `progress` - The progress report.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError::TaskNotFound`] if the task ID in the progress
    /// report does not match an active task.
    // ANVIL Spec section 11.1 -- Worker: report_progress
    async fn report_progress(&self, progress: &TaskProgress) -> Result<(), CollabError>;

    /// Submits the result of a completed task.
    ///
    /// Called when the worker has finished executing the task. The result
    /// includes the output data, status, and optional confidence score.
    ///
    /// # Arguments
    ///
    /// * `result` - The task result to submit.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError::TaskNotFound`] if the task ID in the result
    /// does not match an active task.
    // ANVIL Spec section 11.1 -- Worker: submit_result
    async fn submit_result(&self, result: &TaskResult) -> Result<(), CollabError>;

    /// Called when a task is cancelled by the coordinator.
    ///
    /// The worker should clean up any in-progress work for the specified
    /// task and release resources.
    ///
    /// # Arguments
    ///
    /// * `task_id` - The ID of the cancelled task.
    /// * `reason` - A human-readable cancellation reason.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError::TaskNotFound`] if the task is not known.
    // ANVIL Spec section 11.1 -- Worker: on_task_cancelled
    async fn on_task_cancelled(&self, task_id: &str, reason: &str) -> Result<(), CollabError>;
}

#[async_trait::async_trait]
pub trait PeerContract: Send + Sync {
    /// Submits a proposal for peer consensus.
    ///
    /// # Arguments
    ///
    /// * `proposal` - The proposal data to submit for voting.
    ///
    /// # Returns
    ///
    /// A unique proposal ID string.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError::DelegationFailed`] if the proposal cannot be
    /// submitted (e.g., session is not in an active state).
    // ANVIL Spec section 11.1 -- Peer: propose
    async fn propose(&self, proposal: &serde_json::Value) -> Result<String, CollabError>;

    /// Votes on an existing proposal.
    ///
    /// # Arguments
    ///
    /// * `proposal_id` - The ID of the proposal to vote on.
    /// * `approve` - `true` to approve, `false` to reject.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError::DelegationFailed`] if the proposal is not
    /// found or voting has closed.
    // ANVIL Spec section 11.1 -- Peer: vote
    async fn vote(&self, proposal_id: &str, approve: bool) -> Result<(), CollabError>;

    /// Called when consensus is reached on a proposal.
    ///
    /// # Arguments
    ///
    /// * `proposal_id` - The ID of the proposal that reached consensus.
    /// * `result` - The consensus result data.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError`] if the agent cannot act on the consensus.
    // ANVIL Spec section 11.1 -- Peer: on_consensus
    async fn on_consensus(
        &self,
        proposal_id: &str,
        result: &serde_json::Value,
    ) -> Result<(), CollabError>;
}

session.rs

Read declaration text · 8 declaration entries

#[async_trait::async_trait]
pub trait SessionContract: Send + Sync {
    /// Called when the agent joins a collaboration session.
    ///
    /// # Arguments
    ///
    /// * `session` - The session being joined.
    /// * `role` - The role assigned to this agent in the session.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError`] if the agent cannot join the session.
    // ANVIL Spec section 11.2 -- Session Lifecycle: on_session_join
    async fn on_session_join(
        &self,
        session: &CollaborationSession,
        role: CollaborationRole,
    ) -> Result<(), CollabError>;

    /// Called when the session transitions between states.
    ///
    /// # Arguments
    ///
    /// * `session_id` - The ID of the session.
    /// * `from` - The previous state.
    /// * `to` - The new state.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError`] if the agent cannot handle the transition.
    // ANVIL Spec section 11.2 -- Session Lifecycle: on_session_transition
    async fn on_session_transition(
        &self,
        session_id: &str,
        from: SessionState,
        to: SessionState,
    ) -> Result<(), CollabError>;

    /// Called when the agent leaves a collaboration session.
    ///
    /// # Arguments
    ///
    /// * `session_id` - The ID of the session being left.
    /// * `reason` - A human-readable reason for leaving.
    ///
    /// # Errors
    ///
    /// Returns [`CollabError`] if cleanup fails.
    // ANVIL Spec section 11.2 -- Session Lifecycle: on_session_leave
    async fn on_session_leave(&self, session_id: &str, reason: &str) -> Result<(), CollabError>;
}

#[derive(Debug, Clone)]
pub struct SessionManager {

}

pub fn new() -> Self;

pub fn state(&self) -> SessionState;

pub fn transition(&mut self, target: SessionState) -> Result<SessionTransition, CollabError>;

pub fn can_transition_to(&self, target: SessionState) -> bool;

pub fn valid_transitions(&self) -> Vec<SessionState>;

pub fn history(&self) -> &[SessionTransition];

types.rs

Read declaration text · 21 declaration entries

#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CollaborationRole {
    /// Coordinator: decomposes tasks, assigns them to workers, and aggregates
    /// results. There is at most one coordinator per session.
    Coordinator,
    /// Worker: executes delegated tasks, reports progress, and submits results
    /// back to the coordinator.
    Worker,
    /// Peer: participates in consensus-based collaboration where all agents
    /// have equal standing.
    Peer,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CollaborationSession {
/// Unique identifier for this session.

pub session_id: String,
/// The type of collaboration (e.g., "code-review", "data-analysis").

pub session_type: String,
/// The agents participating in this session.

pub participants: Vec<SessionParticipant>,
/// The DID of the coordinator agent, if any.

pub coordinator: Option<String>,
/// The ID of the shared context store for this session.

pub shared_context_id: Option<String>,
/// ISO 8601 timestamp when the session was created.

pub created_at: String,
/// Optional session timeout in seconds. When elapsed, the session

/// automatically transitions to the `TimedOut` terminal state.

pub timeout_seconds: Option<u64>,
/// Arbitrary metadata for application-specific session properties.

pub metadata: serde_json::Value
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SessionState {
    /// The session has been proposed but not yet accepted by all participants.
    Proposed,
    /// The session is active and participants can collaborate.
    Active,
    /// The session is in the process of completing (aggregating results).
    Completing,
    /// The session has completed successfully. Terminal state.
    Completed,
    /// The session exceeded its timeout. Terminal state.
    TimedOut,
    /// The session was dissolved before completion. Terminal state.
    Dissolved,
}

pub fn valid_transitions(&self) -> &'static [SessionState];

pub fn is_terminal(&self) -> bool;

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SessionParticipant {
/// The OAS DID of the participating agent.

pub agent_did: String,
/// The role this agent plays in the session.

pub role: CollaborationRole,
/// ISO 8601 timestamp when the agent joined the session.

pub joined_at: String,
/// Current participation status (e.g., "active", "left", "disconnected").

pub status: String
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DelegatedTask {
/// Unique identifier for this task.

pub task_id: String,
/// The type of task (e.g., "code-review", "summarize", "translate").

pub task_type: String,
/// Human-readable description of the task.

pub description: String,
/// Input data for the task.

pub input: serde_json::Value,
/// Optional JSON Schema describing the expected output format.

pub output_schema: Option<serde_json::Value>,
/// Constraints bounding the task execution.

pub constraints: TaskConstraints,
/// The OAS DID of the agent that delegated this task.

pub delegator: String,
/// Priority level for task scheduling.

pub priority: TaskPriority,
/// Optional ISO 8601 deadline for task completion.

pub deadline: Option<String>,
/// Keys in the shared context that this task may read.

pub context_keys: Vec<String>
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TaskConstraints {
/// Maximum number of reasoning steps the worker may take.

pub max_steps: Option<u32>,
/// Maximum number of tokens the worker may consume.

pub max_tokens: Option<u64>,
/// Maximum duration in seconds for task execution.

pub max_duration_seconds: Option<u64>,
/// Allowed tool names. If empty, no tool restrictions apply.

pub allowed_tools: Vec<String>,
/// Minimum required confidence score (0.0 to 1.0) for the result.

pub required_confidence: Option<f64>
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TaskResult {
/// The ID of the task this result corresponds to.

pub task_id: String,
/// The completion status of the task.

pub status: TaskStatus,
/// The output data produced by the worker.

pub output: serde_json::Value,
/// Optional confidence score (0.0 to 1.0) for the result.

pub confidence: Option<f64>,
/// Arbitrary metadata about the task execution.

pub metadata: serde_json::Value,
/// Optional Ed25519 signature over the result for verification.

pub signature: Option<String>
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum TaskStatus {
    /// The task completed successfully with full output.
    Completed,
    /// The task failed and could not produce output.
    Failed,
    /// The task produced partial output but could not fully complete.
    Partial,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TaskAcknowledgment {
/// The ID of the task being acknowledged.

pub task_id: String,
/// Whether the worker accepts the task.

pub accepted: bool,
/// Reason for rejection, if `accepted` is `false`.

pub rejection_reason: Option<String>,
/// Estimated time to completion in seconds, if accepted.

pub estimated_completion_seconds: Option<u64>
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TaskProgress {
/// The ID of the task this progress report corresponds to.

pub task_id: String,
/// Completion percentage (0.0 to 100.0).

pub percentage: f64,
/// Human-readable status message.

pub message: String
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum TaskPriority {
    /// Highest priority -- task must be handled immediately.
    Critical,
    /// High priority -- task should be handled soon.
    High,
    /// Normal priority -- default scheduling.
    Normal,
    /// Low priority -- task can wait.
    Low,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Interrupt {
/// Unique identifier for this interrupt.

pub interrupt_id: String,
/// The type of interrupt.

pub interrupt_type: InterruptType,
/// The OAS DID of the agent that sent the interrupt.

pub source: String,
/// Priority of the interrupt.

pub priority: TaskPriority,
/// Arbitrary payload data for the interrupt.

pub payload: serde_json::Value,
/// ISO 8601 timestamp when the interrupt was created.

pub timestamp: String
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum InterruptType {
    /// Override the current task with a higher-priority task.
    PriorityOverride,
    /// Suspend the current task (can be resumed later).
    Suspend,
    /// Resume a previously suspended task.
    Resume,
    /// Abort the current task entirely.
    Abort,
    /// Redirect the agent to a different task or session.
    Redirect,
    /// A human operator is interjecting into the agent's workflow.
    HumanInterjection,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct InterruptResponse {
/// The ID of the interrupt being responded to.

pub interrupt_id: String,
/// Whether the agent acknowledged and handled the interrupt.

pub acknowledged: bool,
/// The agent's state at the time of interruption, if applicable.

pub current_state: Option<InterruptedState>
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct InterruptedState {
/// The ID of the task that was interrupted, if any.

pub task_id: Option<String>,
/// The number of steps completed before interruption.

pub step_count: u32,
/// Progress percentage at the time of interruption (0.0 to 100.0).

pub progress_percentage: f64,
/// Whether the task can be resumed from this state.

pub can_resume: bool
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ContextEntry {
/// The key identifying this context entry.

pub key: String,
/// The value stored in this entry.

pub value: serde_json::Value,
/// The type of the value (e.g., "json", "text", "binary").

pub value_type: String,
/// The OAS DID of the agent that wrote this entry.

pub author: String,
/// Monotonically increasing version number for this key.

pub version: u64,
/// ISO 8601 timestamp when this version was written.

pub timestamp: String,
/// Visibility control for this entry.

pub visibility: ContextVisibility
}

#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ContextVisibility {
    /// Visible to all participants in the session.
    Session,
    /// Visible only to agents with the specified role.
    Role(CollaborationRole),
    /// Visible only to the specific agent identified by DID.
    Agent(String),
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AgentCapabilityProfile {
/// The OAS DID of the agent.

pub agent_did: String,
/// Collaboration roles this agent supports.

pub supported_roles: Vec<CollaborationRole>,
/// Task types this agent can handle (e.g., "code-review", "translate").

pub supported_task_types: Vec<String>,
/// Tool names available to this agent.

pub available_tools: Vec<String>,
/// Current load factor (0.0 = idle, 1.0 = fully loaded).

pub current_load: f64,
/// Maximum number of tasks this agent can execute concurrently.

pub max_concurrent_tasks: u32
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SessionTransition {
/// The state before the transition.

pub from: SessionState,
/// The state after the transition.

pub to: SessionState,
/// ISO 8601 timestamp when the transition occurred.

pub timestamp: String
}

Continue

On this page