pub struct NativeAsyncCoordinator { /* private fields */ }Implementations§
Source§impl NativeAsyncCoordinator
impl NativeAsyncCoordinator
Sourcepub async fn open(
journal: Box<dyn NativeAsyncJournal>,
executor: Arc<dyn NativeAsyncExecutor>,
max_concurrency: usize,
parallel_tool_calls: bool,
) -> Result<NativeAsyncCoordinator, AgentLoopError>
pub async fn open( journal: Box<dyn NativeAsyncJournal>, executor: Arc<dyn NativeAsyncExecutor>, max_concurrency: usize, parallel_tool_calls: bool, ) -> Result<NativeAsyncCoordinator, AgentLoopError>
Recover only after acquiring the journal’s exclusive ownership fence. Safe calls are reauthorized and restarted; unsafe calls become explicit interrupted outputs. Ambiguous HTTP delivery requires receipt recovery.
pub fn checkpoint(&self) -> &NativeAsyncCheckpoint
pub async fn persist_host_outcome( &mut self, outcome: Value, ) -> Result<(), AgentLoopError>
pub async fn release(&mut self) -> Result<(), AgentLoopError>
pub async fn register( &mut self, call: NativeToolCall, ) -> Result<(), AgentLoopError>
pub async fn begin_transcript_response( &mut self, message_id: String, ) -> Result<(), AgentLoopError>
pub async fn stage_transcript_result( &mut self, result: Value, ) -> Result<(), AgentLoopError>
pub async fn transcript_committed( &mut self, message_id: &str, ) -> Result<(), AgentLoopError>
Sourcepub async fn begin_response(&mut self) -> Result<(), AgentLoopError>
pub async fn begin_response(&mut self) -> Result<(), AgentLoopError>
Record request intent before a host opens its provider HTTP stream.
Sourcepub async fn next_response_event(
&mut self,
stream: &mut Pin<Box<dyn Stream<Item = Result<LlmStreamEvent, AgentLoopError>> + Send>>,
) -> Result<LlmStreamEvent, AgentLoopError>
pub async fn next_response_event( &mut self, stream: &mut Pin<Box<dyn Stream<Item = Result<LlmStreamEvent, AgentLoopError>> + Send>>, ) -> Result<LlmStreamEvent, AgentLoopError>
Drive jobs alongside one stream event, retaining call events for the host’s normal transcript pipeline. A completed response is not a turn end.
Sourcepub async fn pump(
&mut self,
stream: Pin<Box<dyn Stream<Item = Result<LlmStreamEvent, AgentLoopError>> + Send>>,
observe: impl FnMut(LlmStreamEvent),
) -> Result<(), AgentLoopError>
pub async fn pump( &mut self, stream: Pin<Box<dyn Stream<Item = Result<LlmStreamEvent, AgentLoopError>> + Send>>, observe: impl FnMut(LlmStreamEvent), ) -> Result<(), AgentLoopError>
Consume a response while jobs execute. Returns once the provider response
is complete, even if async work remains. Independent follow-up responses
may be pumped before waiting for jobs. observe receives prose/reasoning
unchanged; only executable call events are consumed by the coordinator.
Sourcepub async fn wait_next(&mut self) -> Result<bool, AgentLoopError>
pub async fn wait_next(&mut self) -> Result<bool, AgentLoopError>
Wait for one completion; outputs may be delivered out of launch order.
Sourcepub async fn prepare_delivery(
&mut self,
) -> Result<Option<Delivery>, AgentLoopError>
pub async fn prepare_delivery( &mut self, ) -> Result<Option<Delivery>, AgentLoopError>
Persist the delivery intent before the caller submits its HTTP request. Synchronous calls must all finish before the provider can continue.
Sourcepub async fn cancel(&mut self) -> Result<(), AgentLoopError>
pub async fn cancel(&mut self) -> Result<(), AgentLoopError>
Dropping owned futures cancels local execution. Persist cancellation outputs; the conversation remains incomplete until they are delivered.
Source§impl NativeAsyncCoordinator
impl NativeAsyncCoordinator
Sourcepub async fn run<Request, RequestFuture>(
&mut self,
max_responses: usize,
request: Request,
observe: impl FnMut(LlmStreamEvent),
) -> Result<(), AgentLoopError>where
Request: FnMut(Option<Delivery>, Option<String>) -> RequestFuture,
RequestFuture: Future<Output = Result<Pin<Box<dyn Stream<Item = Result<LlmStreamEvent, AgentLoopError>> + Send>>, AgentLoopError>>,
pub async fn run<Request, RequestFuture>(
&mut self,
max_responses: usize,
request: Request,
observe: impl FnMut(LlmStreamEvent),
) -> Result<(), AgentLoopError>where
Request: FnMut(Option<Delivery>, Option<String>) -> RequestFuture,
RequestFuture: Future<Output = Result<Pin<Box<dyn Stream<Item = Result<LlmStreamEvent, AgentLoopError>> + Send>>, AgentLoopError>>,
Drive HTTP response continuations through completion, returning only when every accepted call’s output has a provider receipt. The caller supplies request construction; no transport or model defaults are chosen here.