async_factory.txt (7099B)
1 Async Factory Pattern 2 3 A constructor method that constructs the Future function without calling 4 it. This is used when we don't know about the lifetime of the function. 5 6 The factory pattern splits execution into two phases: 7 8 1. Synchronous setup, runs now: request_body(messages) and api_key.clone() execute immediately, while &mut self is still held. Everything the request needs is copied out of self. 9 10 2. The actual I/O, runs whenever: the async move block captures only those owned values (body, api_key, chunks), so the returned future is 'static — it has no ties to self's lifetime. Hence the doc's "we don't know about the lifetime of the function": the caller might spawn it on another task, join it with other futures, or drop it, and none of that has to be coordinated with the borrow of self. 11 12 (Rust futures are always lazy — even a plain async fn doesn't run until polled. What the factory buys you isn't laziness, it's the lifetime decoupling: setup borrows self, the future doesn't.) 13 14 Box::pin earns its keep when the concrete type must be erased — e.g. different provider implementations behind one trait, or storing the future in a struct field. 15 16 17 Async Factory: 18 19 impl LlmInference for LlmClaudeContext { 20 /// Sends the call and pushes each body frame into `chunks`. 21 /// 22 /// The body — settings folded together with `messages` — and the key are both 23 /// taken before the returned future is built, so the future borrows nothing 24 /// from `self` and can be spawned or joined freely. 25 /// 26 /// A non-success status is rejected before any frame is sent, so a consumer 27 /// never sees part of a failed response — the provider's own error text comes 28 /// back on the error instead. 29 /// 30 /// # Arguments 31 /// - `chunks`: Each body frame, in arrival order. Closed when this returns. 32 /// 33 /// # Returns 34 /// - `Ok(())` at the end of the body, or an [`LlmError`] if the request 35 /// failed, the status was rejected, the connection dropped mid-body, or the 36 /// receiver went away while frames were still arriving. 37 fn http_infer( 38 &mut self, 39 chunks: Sender<Vec<u8>>, 40 messages: Vec<LlmClaudeMessage>, 41 ) -> Pin<Box<dyn Future<Output = Result<(), LlmError>> + Send + 'static>> { 42 let body = self.request_body(messages); 43 let api_key = self.api_key.clone(); 44 45 Box::pin(async move { 46 let mut response = Client::new() 47 .post(REQUEST_URL) 48 .header("x-api-key", api_key) 49 .header("anthropic-version", PROVIDER_VERSION) 50 .json(&body) 51 .send() 52 .await 53 .map_err(LlmError::Http)?; 54 55 // `send` resolves on the response headers, so the status is known 56 // while the body is still in flight. Rejecting a bad one here means 57 // the body never reaches the consumer. 58 let status = response.status(); 59 if !status.is_success() { 60 let detail = response 61 .text() 62 .await 63 .map_err(|e| LlmError::Inference { status, message: e.to_string() })?; 64 return Err(LlmError::Inference { status, message: detail }); 65 } 66 67 // `chunk` resolves as soon as hyper has a frame and yields `None` at 68 // the end of the body, so nothing collects the reply. `send` waits 69 // when the channel is full — that wait is the backpressure. 70 while let Some(frame) = response.chunk().await.map_err(LlmError::Http)? { 71 chunks.send(frame.to_vec()).await.map_err(|_| { 72 LlmError::ClaudeResponseStreaming { 73 message: "the chunk receiver was dropped while the response body was \ 74 still arriving" 75 .to_string(), 76 } 77 })?; 78 } 79 Ok(()) 80 }) 81 } 82 } 83 84 Call site: 85 86 In the call site it just spawns that constructed Future function, and then 87 polls the Future's inner channel that is streaming the responses from 88 claude, until all messages are polled, and then calls await on the handle 89 itself to make sure the Future is finished. 90 91 impl LlmInferenceContext { 92 /// Runs one inference pass: says `message`, then turns the streamed reply into 93 /// the files it asked for. 94 /// 95 /// The turn is appended to [`messages`](LlmInferenceContext::messages) before 96 /// the call and the model's own reply is appended after it, so a later pass 97 /// sees the whole conversation. The provider is handed a *snapshot*, which it 98 /// consumes; the context's own history is untouched by that. 99 /// 100 /// Reading and lexing run concurrently with the request: the provider is 101 /// spawned and its frames are consumed here as they arrive, so a file is 102 /// available as soon as its fence closes rather than at the end of the reply. 103 /// 104 /// # Arguments 105 /// - `root_path`: The session directory every fenced path is resolved against. 106 /// - `message`: What to say to the model this pass. 107 /// 108 /// # Returns 109 /// - Every file the reply carried, in the order the model emitted them. An 110 /// [`LlmError`] if the request failed, a frame could not be lexed, or a 111 /// fenced path was unsafe to write. 112 pub async fn infer_llm( 113 &mut self, 114 root_path: impl Into<String>, 115 message: impl Into<String>, 116 ) -> Result<Vec<AgentFileEvent>, LlmError> { 117 // Append, then snapshot: the provider consumes its copy, so the clone is 118 // what keeps this context's history intact across calls. 119 self.messages.push(LlmClaudeMessage::user(message)); 120 let history = self.messages.clone(); 121 122 let (tx, mut rx) = mpsc::channel(10); 123 let req_fut = self.http_ai_provider.http_infer(tx, history); 124 125 let handle = tokio::spawn(req_fut); 126 // Right now we're using the claude buffer 127 let mut buffer = Buffer::new(root_path); 128 let mut events_buffer = Vec::new(); 129 130 while let Some(chunk) = rx.recv().await { 131 match buffer.ingest_chunk(chunk)? { 132 BufferResponse::Events(events) => events_buffer.extend(events), 133 BufferResponse::End { events, end_of_stream } => { 134 events_buffer.extend(events); 135 // The model's turn, verbatim, so the next pass can see what it 136 // already wrote rather than only what we asked for. 137 self.messages.push(LlmClaudeMessage::assistant(end_of_stream.reply)); 138 }, 139 BufferResponse::None => {}, 140 } 141 } 142 // Two layers: the outer says whether the task ran at all (it panicked or 143 // was cancelled), the inner whether the request itself succeeded. 144 handle.await.map_err(|e| LlmError::ClaudeResponseStreaming { 145 message: format!("the inference task did not finish: {e}"), 146 })??; 147 Ok(events_buffer) 148 } 149 }