lean_ctx/
daemon_client.rs1use anyhow::{Context, Result};
2use tokio::io::{AsyncReadExt, AsyncWriteExt};
3
4use crate::daemon;
5use crate::ipc;
6
7pub async fn daemon_request(method: &str, path: &str, body: &str) -> Result<String> {
10 use std::time::Duration;
11 use tokio::time::timeout;
12
13 const CONNECT_TIMEOUT: Duration = Duration::from_secs(3);
14 const IO_TIMEOUT: Duration = Duration::from_secs(10);
15
16 let addr = daemon::daemon_addr();
17 if !addr.is_listening() {
18 anyhow::bail!(
19 "Daemon endpoint not found at {}. Is the daemon running?",
20 addr.display()
21 );
22 }
23
24 let request = format_http_request(method, path, body);
25
26 #[cfg(unix)]
27 {
28 let mut stream = timeout(CONNECT_TIMEOUT, ipc::connect(&addr))
29 .await
30 .with_context(|| {
31 format!(
32 "connect to daemon timed out ({}s)",
33 CONNECT_TIMEOUT.as_secs()
34 )
35 })?
36 .with_context(|| format!("cannot connect to daemon at {}", addr.display()))?;
37
38 timeout(IO_TIMEOUT, stream.write_all(request.as_bytes()))
39 .await
40 .context("write to daemon timed out")?
41 .context("failed to write request to daemon")?;
42
43 let mut response_buf = Vec::with_capacity(4096);
44 timeout(IO_TIMEOUT, stream.read_to_end(&mut response_buf))
45 .await
46 .context("read from daemon timed out")?
47 .context("failed to read response from daemon")?;
48
49 parse_http_response(&response_buf)
50 }
51
52 #[cfg(windows)]
53 {
54 let mut stream = timeout(CONNECT_TIMEOUT, ipc::connect(&addr))
55 .await
56 .with_context(|| {
57 format!(
58 "connect to daemon timed out ({}s)",
59 CONNECT_TIMEOUT.as_secs()
60 )
61 })?
62 .with_context(|| format!("cannot connect to daemon at {}", addr.display()))?;
63
64 timeout(IO_TIMEOUT, stream.write_all(request.as_bytes()))
65 .await
66 .context("write to daemon timed out")?
67 .context("failed to write request to daemon")?;
68
69 let mut response_buf = Vec::with_capacity(4096);
70 timeout(IO_TIMEOUT, stream.read_to_end(&mut response_buf))
71 .await
72 .context("read from daemon timed out")?
73 .context("failed to read response from daemon")?;
74
75 parse_http_response(&response_buf)
76 }
77}
78
79pub async fn daemon_health_check() -> bool {
81 match daemon_request("GET", "/health", "").await {
82 Ok(body) => body.trim() == "ok",
83 Err(_) => false,
84 }
85}
86
87pub async fn daemon_tool_call(name: &str, arguments: Option<&serde_json::Value>) -> Result<String> {
89 let body = serde_json::json!({
90 "name": name,
91 "arguments": arguments,
92 });
93 daemon_request("POST", "/v1/tools/call", &body.to_string()).await
94}
95
96fn format_http_request(method: &str, path: &str, body: &str) -> String {
97 if body.is_empty() {
98 format!("{method} {path} HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n")
99 } else {
100 let content_length = body.len();
101 format!(
102 "{method} {path} HTTP/1.1\r\nHost: localhost\r\nContent-Type: application/json\r\nContent-Length: {content_length}\r\nConnection: close\r\n\r\n{body}"
103 )
104 }
105}
106
107fn parse_http_response(raw: &[u8]) -> Result<String> {
108 let response_str = std::str::from_utf8(raw).context("daemon response is not valid UTF-8")?;
109
110 let Some(header_end) = response_str.find("\r\n\r\n") else {
111 anyhow::bail!("malformed HTTP response from daemon (no header boundary)");
112 };
113
114 let headers = &response_str[..header_end];
115 let body = &response_str[header_end + 4..];
116
117 let status_line = headers.lines().next().unwrap_or("");
118 let status_code = status_line
119 .split_whitespace()
120 .nth(1)
121 .and_then(|s| s.parse::<u16>().ok())
122 .unwrap_or(0);
123
124 if status_code >= 400 {
125 anyhow::bail!("daemon returned HTTP {status_code}: {body}");
126 }
127
128 Ok(body.to_string())
129}
130
131pub async fn try_daemon_request(method: &str, path: &str, body: &str) -> Option<String> {
133 if !daemon::is_daemon_running() {
134 return None;
135 }
136 daemon_request(method, path, body).await.ok()
137}
138
139pub fn notify_cache_clear() -> bool {
145 if !daemon::is_daemon_running() {
146 return false;
147 }
148 let Ok(rt) = tokio::runtime::Runtime::new() else {
149 return false;
150 };
151 let body = serde_json::json!({
152 "name": "ctx_cache",
153 "arguments": { "action": "clear" },
154 });
155 rt.block_on(async {
156 try_daemon_request("POST", "/v1/tools/call", &body.to_string())
157 .await
158 .is_some()
159 })
160}
161
162#[allow(clippy::needless_pass_by_value)]
166pub fn try_daemon_tool_call_blocking(
167 name: &str,
168 arguments: Option<serde_json::Value>,
169) -> Option<String> {
170 use std::time::Duration;
171
172 let rt = tokio::runtime::Runtime::new().ok()?;
173
174 let addr = daemon::daemon_addr();
175 let mut ready = addr.is_listening() && rt.block_on(async { daemon_health_check().await });
176
177 if !ready {
178 if std::env::var("LEAN_CTX_HOOK_CHILD").is_ok() {
179 return None;
180 }
181
182 let lock = crate::core::startup_guard::try_acquire_lock(
183 "daemon-start",
184 Duration::from_millis(1200),
185 Duration::from_secs(5),
186 );
187
188 if let Some(g) = lock {
189 g.touch();
190 let mut did_start = false;
191
192 if !daemon::is_daemon_running() {
193 if daemon::start_daemon(&[]).is_ok() {
194 did_start = true;
195 } else {
196 return None;
197 }
198 }
199
200 for _ in 0..60 {
201 if addr.is_listening() && rt.block_on(async { daemon_health_check().await }) {
202 ready = true;
203 break;
204 }
205 std::thread::sleep(Duration::from_millis(50));
206 }
207
208 if ready && did_start && crate::core::protocol::meta_visible() {
209 eprintln!("\x1b[2m▸ daemon auto-started\x1b[0m");
210 }
211 } else {
212 for _ in 0..60 {
213 if addr.is_listening() && rt.block_on(async { daemon_health_check().await }) {
214 ready = true;
215 break;
216 }
217 std::thread::sleep(Duration::from_millis(50));
218 }
219 }
220 }
221
222 if !ready {
223 return None;
224 }
225
226 if let Some(out) = rt.block_on(async { daemon_tool_call(name, arguments.as_ref()).await.ok() })
227 {
228 return Some(out);
229 }
230
231 for _ in 0..5 {
232 std::thread::sleep(Duration::from_millis(50));
233 if let Some(out) =
234 rt.block_on(async { daemon_tool_call(name, arguments.as_ref()).await.ok() })
235 {
236 return Some(out);
237 }
238 }
239
240 None
241}
242
243fn unwrap_mcp_tool_text(body: &str) -> Option<String> {
244 let v: serde_json::Value = serde_json::from_str(body).ok()?;
245 let result = v.get("result")?;
246
247 if let Some(content) = result.get("content").and_then(|c| c.as_array()) {
248 let mut texts: Vec<String> = Vec::new();
249 for item in content {
250 if let Some(text) = item.get("text").and_then(|t| t.as_str())
251 && !text.is_empty()
252 {
253 texts.push(text.to_string());
254 }
255 }
256 if !texts.is_empty() {
257 return Some(texts.join("\n"));
258 }
259 }
260
261 if let Some(text) = result.get("text").and_then(|t| t.as_str()) {
262 return Some(text.to_string());
263 }
264
265 result.as_str().map(std::string::ToString::to_string)
266}
267
268pub fn try_daemon_tool_call_blocking_text(
270 name: &str,
271 arguments: Option<serde_json::Value>,
272) -> Option<String> {
273 let body = try_daemon_tool_call_blocking(name, arguments)?;
274 let trimmed = body.trim_start();
275 if !trimmed.starts_with('{') {
276 return Some(body);
277 }
278 Some(unwrap_mcp_tool_text(&body).unwrap_or(body))
279}