soma_codemode/pool/
checkout.rs1use std::sync::Arc;
2
3use tokio::sync::{Mutex, Semaphore};
4
5use crate::ToolError;
6
7use super::config::PoolConfig;
8use super::disposition::RunnerDisposition;
9use super::runner_handle::{RunnerHandle, RunnerSpawn};
10
11pub struct RunnerPool {
12 config: PoolConfig,
13 spawn: RunnerSpawn,
14 overflow: Arc<Semaphore>,
15 available: Mutex<Vec<RunnerHandle>>,
16}
17
18impl RunnerPool {
19 pub fn new(config: PoolConfig, spawn: RunnerSpawn) -> Self {
20 Self {
21 overflow: Arc::new(Semaphore::new(
22 config.size.saturating_add(config.max_overflow).max(1),
23 )),
24 config,
25 spawn,
26 available: Mutex::new(Vec::new()),
27 }
28 }
29
30 pub async fn checkout(&self) -> Result<RunnerLease, ToolError> {
31 let permit = self
32 .overflow
33 .clone()
34 .acquire_owned()
35 .await
36 .map_err(|_| ToolError::internal_message("runner pool semaphore closed"))?;
37 let handle = if self.config.is_disabled() {
38 RunnerHandle::spawn(&self.spawn)?
39 } else {
40 match self.available.lock().await.pop() {
41 Some(handle) => handle,
42 None => RunnerHandle::spawn(&self.spawn)?,
43 }
44 };
45 Ok(RunnerLease {
46 handle: Some(handle),
47 _permit: permit,
48 })
49 }
50
51 pub async fn release(&self, mut lease: RunnerLease, disposition: RunnerDisposition) {
52 let Some(handle) = lease.handle.take() else {
53 return;
54 };
55 if self.config.is_disabled() || !matches!(disposition, RunnerDisposition::Reuse) {
56 return;
57 }
58 let mut available = self.available.lock().await;
59 if available.len() < self.config.size {
60 available.push(handle);
61 }
62 }
63
64 pub fn config(&self) -> PoolConfig {
65 self.config
66 }
67
68 pub fn spawn(&self) -> &RunnerSpawn {
69 &self.spawn
70 }
71}
72
73pub struct RunnerLease {
74 pub handle: Option<RunnerHandle>,
75 _permit: tokio::sync::OwnedSemaphorePermit,
76}
77
78impl RunnerLease {
79 pub fn handle_mut(&mut self) -> Result<&mut RunnerHandle, ToolError> {
80 self.handle
81 .as_mut()
82 .ok_or_else(|| ToolError::internal_message("runner lease has no handle"))
83 }
84}