返回 CodeWhale
pdf.rs
根目录 / crates / tui / src / tools / pdf.rs
1 //! Shared PDF-to-text adapter.
2 //!
3 //! PDF parsing is intentionally delegated to the optional `pdftotext`
4 //! executable. Keeping the adapter here gives file and web tools one error
5 //! contract without carrying a second parser and font stack in Codewhale.
6
7 use std::ffi::{OsStr, OsString};
8 use std::fmt;
9 use std::io::Write;
10 use std::path::{Path, PathBuf};
11 use std::process::Stdio;
12 use std::time::Duration;
13
14 use serde_json::json;
15 use tokio::io::{AsyncRead, AsyncReadExt};
16 use tokio_util::sync::CancellationToken;
17
18 use super::spec::ToolError;
19
20 const PDF_TEXT_TIMEOUT: Duration = Duration::from_secs(30);
21 const PDF_PIPE_DRAIN_TIMEOUT: Duration = Duration::from_secs(1);
22 const MAX_PDF_STDOUT_BYTES: usize = 16 * 1024 * 1024;
23 const MAX_PDF_STDERR_BYTES: usize = 32 * 1024;
24
25 #[derive(Debug, Clone, PartialEq, Eq)]
26 pub(crate) enum PdfTextError {
27 BinaryUnavailable,
28 Cancelled,
29 TimedOut,
30 Execution(String),
31 }
32
33 impl fmt::Display for PdfTextError {
34 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
35 match self {
36 Self::BinaryUnavailable => formatter.write_str(
37 "PDF text extraction requires the optional `pdftotext` executable (Poppler)",
38 ),
39 Self::Cancelled => formatter.write_str("PDF text extraction was cancelled"),
40 Self::TimedOut => write!(
41 formatter,
42 "PDF text extraction timed out after {} seconds",
43 PDF_TEXT_TIMEOUT.as_secs()
44 ),
45 Self::Execution(message) => formatter.write_str(message),
46 }
47 }
48 }
49
50 /// One typed mapping shared by local-file and fetched-PDF consumers.
51 ///
52 /// The missing-binary message is deliberately a small JSON object. The
53 /// `NotAvailable` variant gives the runtime a failed terminal status while
54 /// callers that inspect the variant retain machine-readable recovery data.
55 pub(super) fn into_tool_error(error: PdfTextError) -> ToolError {
56 match error {
57 PdfTextError::BinaryUnavailable => ToolError::not_available(
58 json!({
59 "type": "binary_unavailable",
60 "kind": "pdf",
61 "binary": "pdftotext",
62 "reason": "optional pdftotext executable is not installed",
63 "hint": "install Poppler and ensure pdftotext is on PATH"
64 })
65 .to_string(),
66 ),
67 PdfTextError::Cancelled => ToolError::cancelled("PDF text extraction was cancelled"),
68 PdfTextError::TimedOut => ToolError::Timeout {
69 seconds: PDF_TEXT_TIMEOUT.as_secs(),
70 },
71 PdfTextError::Execution(message) => ToolError::execution_failed(message),
72 }
73 }
74
75 #[derive(Clone, Copy)]
76 pub(crate) struct PdfTextCommand<'a> {
77 binary: &'a OsStr,
78 timeout: Duration,
79 cancel: Option<&'a CancellationToken>,
80 context: Option<&'a super::spec::ToolContext>,
81 }
82
83 impl<'a> PdfTextCommand<'a> {
84 pub(super) fn context(self) -> Option<&'a super::spec::ToolContext> {
85 self.context
86 }
87 pub(super) fn system(context: Option<&'a super::spec::ToolContext>) -> Self {
88 Self {
89 binary: OsStr::new("pdftotext"),
90 timeout: PDF_TEXT_TIMEOUT,
91 cancel: context.and_then(|context| context.cancel_token.as_ref()),
92 context,
93 }
94 }
95
96 #[cfg(test)]
97 pub(super) fn test(
98 binary: &'a OsStr,
99 timeout: Duration,
100 cancel: Option<&'a CancellationToken>,
101 ) -> Self {
102 Self {
103 binary,
104 timeout,
105 cancel,
106 context: None,
107 }
108 }
109
110 #[cfg(all(test, unix))]
111 pub(super) fn with_context(mut self, context: &'a super::spec::ToolContext) -> Self {
112 self.context = Some(context);
113 if self.cancel.is_none() {
114 self.cancel = context.cancel_token.as_ref();
115 }
116 self
117 }
118 }
119
120 pub(super) async fn extract_path(
121 path: &Path,
122 page_range: Option<(u32, u32)>,
123 command: PdfTextCommand<'_>,
124 ) -> Result<String, PdfTextError> {
125 extract_captured(
126 CapturedPdf {
127 binary: command.binary.to_os_string(),
128 path: path.to_path_buf(),
129 page_range,
130 timeout: command.timeout,
131 _staged: None,
132 },
133 command,
134 )
135 .await
136 }
137
138 pub(super) async fn extract_bytes(
139 bytes: &[u8],
140 command: PdfTextCommand<'_>,
141 ) -> Result<String, PdfTextError> {
142 let mut input = tempfile::NamedTempFile::new().map_err(|error| {
143 PdfTextError::Execution(format!("failed to stage fetched PDF: {error}"))
144 })?;
145 input.write_all(bytes).map_err(|error| {
146 PdfTextError::Execution(format!("failed to stage fetched PDF: {error}"))
147 })?;
148 input.flush().map_err(|error| {
149 PdfTextError::Execution(format!("failed to stage fetched PDF: {error}"))
150 })?;
151 let path = input.path().to_path_buf();
152 extract_captured(
153 CapturedPdf {
154 binary: command.binary.to_os_string(),
155 path,
156 page_range: None,
157 timeout: command.timeout,
158 _staged: Some(input),
159 },
160 command,
161 )
162 .await
163 }
164
165 /// Private admitted parser request. No command, path or staged PDF crosses IPC.
166 pub(crate) struct CapturedPdf {
167 pub(crate) binary: OsString,
168 pub(crate) path: PathBuf,
169 pub(crate) page_range: Option<(u32, u32)>,
170 pub(crate) timeout: Duration,
171 _staged: Option<tempfile::NamedTempFile>,
172 }
173 impl CapturedPdf {
174 pub(crate) fn digest(&self) -> String {
175 let mut bytes = self.binary.as_encoded_bytes().to_vec();
176 bytes.push(0);
177 bytes.extend_from_slice(self.path.as_os_str().as_encoded_bytes());
178 bytes.push(0);
179 if let Some((start, end)) = self.page_range {
180 bytes.extend_from_slice(&start.to_le_bytes());
181 bytes.extend_from_slice(&end.to_le_bytes());
182 }
183 crate::hashing::sha256_hex(&bytes)
184 }
185 }
186
187 /// The existing bounded process result stays private in one Execution job.
188 pub(crate) struct PdfProcessOutcome {
189 output: Result<(BoundedOutput, BoundedOutput, std::process::ExitStatus), PdfTextError>,
190 }
191 impl PdfProcessOutcome {
192 pub(crate) fn projection(&self) -> serde_json::Value {
193 match &self.output {
194 Ok((stdout, stderr, status)) => json!({"kind":"pdf_process","state":"complete",
195 "success":status.success(),"exit_code":status.code(),
196 "stdout_truncated":stdout.truncated,"stderr":sanitized_text(&stderr.bytes),
197 "stderr_truncated":stderr.truncated}),
198 Err(error) => json!({"kind":"pdf_process","state":match error {
199 PdfTextError::BinaryUnavailable=>"binary_unavailable", PdfTextError::Cancelled=>"cancelled",
200 PdfTextError::TimedOut=>"timed_out", PdfTextError::Execution(_)=>"execution",
201 },"message":error.to_string()}),
202 }
203 }
204 fn into_rust_text(self) -> Result<String, PdfTextError> {
205 let (stdout, stderr, status) = self.output?;
206 if stdout.truncated {
207 return Err(PdfTextError::Execution(format!(
208 "pdftotext output exceeded the {} byte safety limit",
209 MAX_PDF_STDOUT_BYTES
210 )));
211 }
212 if !status.success() {
213 let suffix = if stderr.truncated { " [truncated]" } else { "" };
214 let stderr = sanitized_text(&stderr.bytes);
215 let stderr = if stderr.is_empty() {
216 "no diagnostic output".to_string()
217 } else {
218 stderr
219 };
220 return Err(PdfTextError::Execution(format!(
221 "pdftotext failed (exit {:?}): {stderr}{suffix}",
222 status.code()
223 )));
224 }
225 Ok(String::from_utf8_lossy(&stdout.bytes).into_owned())
226 }
227 fn into_host_text(self, result: super::spec::ToolResult) -> Result<String, PdfTextError> {
228 #[derive(serde::Deserialize)]
229 #[serde(deny_unknown_fields)]
230 struct Decision {
231 kind: String,
232 code: String,
233 message: Option<String>,
234 }
235 let invalid = || {
236 PdfTextError::Execution("PDF host returned a malformed or inconsistent decision; no Rust fallback was attempted".into())
237 };
238 if !result.success {
239 return Err(invalid());
240 }
241 let decision: Decision =
242 serde_json::from_value(result.metadata.ok_or_else(invalid)?).map_err(|_| invalid())?;
243 if decision.kind != "pdf_decision" {
244 return Err(invalid());
245 }
246 match self.output {
247 // Cancellation/unavailability keep their real typed terminal status.
248 Err(error) => {
249 let code = match &error {
250 PdfTextError::BinaryUnavailable => "binary_unavailable",
251 PdfTextError::Cancelled => "cancelled",
252 PdfTextError::TimedOut => "timed_out",
253 PdfTextError::Execution(_) => "execution",
254 };
255 if decision.code != code {
256 return Err(invalid());
257 }
258 Err(error)
259 }
260 Ok((stdout, _, status)) => match decision.code.as_str() {
261 // Mandatory data-loss/process-success guard remains Core-owned.
262 "success"
263 if !stdout.truncated && status.success() && decision.message.is_none() =>
264 {
265 Ok(String::from_utf8_lossy(&stdout.bytes).into_owned())
266 }
267 "execution" if stdout.truncated || !status.success() => Err(
268 PdfTextError::Execution(decision.message.ok_or_else(invalid)?),
269 ),
270 _ => Err(invalid()),
271 },
272 }
273 }
274 }
275
276 async fn extract_captured(
277 input: CapturedPdf,
278 request: PdfTextCommand<'_>,
279 ) -> Result<String, PdfTextError> {
280 if request.cancel.is_some_and(CancellationToken::is_cancelled) {
281 return Err(PdfTextError::Cancelled);
282 }
283 if let Some(context) = request
284 .context
285 .filter(|context| context.features.enabled(crate::features::Feature::PdfHost))
286 {
287 // Capture the existing Core deadline before entering the broker. An
288 // expired job's final receipt check can beat its timeout response; the
289 // generic RPC error must not erase the real Core terminal outcome.
290 let deadline = context
291 .turn_deadline
292 .map(|deadline| deadline.min(tokio::time::Instant::now() + input.timeout))
293 .unwrap_or_else(|| tokio::time::Instant::now() + input.timeout);
294 let result = crate::extension_host::manager()
295 .execute_pdf(input, context)
296 .await;
297 return match result {
298 Ok((outcome, decision)) => outcome.into_host_text(decision),
299 Err(error) => Err(host_terminal_error(error, deadline, request.cancel)),
300 };
301 }
302 run_pdf_driver(input, request.cancel).await.into_rust_text()
303 }
304
305 fn host_terminal_error(
306 error: ToolError,
307 deadline: tokio::time::Instant,
308 cancel: Option<&CancellationToken>,
309 ) -> PdfTextError {
310 if cancel.is_some_and(CancellationToken::is_cancelled)
311 || matches!(error, ToolError::Cancelled { .. })
312 {
313 PdfTextError::Cancelled
314 } else if tokio::time::Instant::now() >= deadline || matches!(error, ToolError::Timeout { .. })
315 {
316 PdfTextError::TimedOut
317 } else {
318 PdfTextError::Execution(format!(
319 "PDF host failed: {error}; no Rust fallback was attempted"
320 ))
321 }
322 }
323
324 /// One actual process driver shared by default Rust and the admitted Host job.
325 /// Host only receives the small process projection; stdout stays in this result.
326 pub(crate) async fn run_pdf_driver(
327 input: CapturedPdf,
328 cancel: Option<&CancellationToken>,
329 ) -> PdfProcessOutcome {
330 let output = async {
331 if cancel.is_some_and(CancellationToken::is_cancelled) { return Err(PdfTextError::Cancelled); }
332 let mut command = tokio::process::Command::new(&input.binary);
333 crate::utils::suppress_tokio_console_window(&mut command);
334 command.arg("-layout");
335 if let Some((start,end))=input.page_range { command.arg("-f").arg(start.to_string()).arg("-l").arg(end.to_string()); }
336 command.arg(&input.path).arg("-").stdin(Stdio::null()).stdout(Stdio::piped()).stderr(Stdio::piped()).kill_on_drop(true);
337 crate::child_env::apply_to_tokio_command(&mut command,std::iter::empty::<(&str,&str)>());
338 #[cfg(unix)] command.process_group(0);
339 let mut child = command.spawn().map_err(|error| if error.kind()==std::io::ErrorKind::NotFound { PdfTextError::BinaryUnavailable } else { PdfTextError::Execution(format!("failed to launch pdftotext: {error}")) })?;
340 let tree = crate::process_tree::ProcessTree::attach_tokio(&child).map_err(|error| PdfTextError::Execution(format!("failed to contain pdftotext: {error}")))?;
341 let stdout=child.stdout.take().ok_or_else(||PdfTextError::Execution("failed to capture pdftotext stdout".into()))?;
342 let stderr=child.stderr.take().ok_or_else(||PdfTextError::Execution("failed to capture pdftotext stderr".into()))?;
343 let stdout_task=tokio::spawn(read_bounded(stdout,MAX_PDF_STDOUT_BYTES));
344 let stderr_task=tokio::spawn(read_bounded(stderr,MAX_PDF_STDERR_BYTES));
345 let status = tokio::select! {
346 biased;
347 ()=wait_for_cancellation(cancel)=>{
348 let _=tree.kill(); terminate_child(&mut child).await;
349 finish_capture_tasks(stdout_task,stderr_task).await?; return Err(PdfTextError::Cancelled);
350 }
351 ()=tokio::time::sleep(input.timeout)=>{
352 let _=tree.kill(); terminate_child(&mut child).await;
353 finish_capture_tasks(stdout_task,stderr_task).await?; return Err(PdfTextError::TimedOut);
354 }
355 status=child.wait()=>status.map_err(|error|PdfTextError::Execution(format!("failed to wait for pdftotext: {error}")))?,
356 };
357 // No parser descendant may retain the capture pipes after the primary exits.
358 let _=tree.kill();
359 let (stdout,stderr)=finish_capture_tasks(stdout_task,stderr_task).await?;
360 Ok((stdout,stderr,status))
361 }.await;
362 PdfProcessOutcome { output }
363 }
364
365 async fn wait_for_cancellation(cancel: Option<&CancellationToken>) {
366 match cancel {
367 Some(cancel) => cancel.cancelled().await,
368 None => std::future::pending::<()>().await,
369 }
370 }
371
372 async fn terminate_child(child: &mut tokio::process::Child) {
373 let _ = child.kill().await;
374 let _ = child.wait().await;
375 }
376
377 struct BoundedOutput {
378 bytes: Vec<u8>,
379 truncated: bool,
380 }
381
382 async fn read_bounded(
383 mut reader: impl AsyncRead + Unpin,
384 max_bytes: usize,
385 ) -> std::io::Result<BoundedOutput> {
386 let mut bytes = Vec::with_capacity(max_bytes.min(8 * 1024));
387 let mut buffer = [0u8; 8 * 1024];
388 let mut truncated = false;
389 loop {
390 let read = reader.read(&mut buffer).await?;
391 if read == 0 {
392 break;
393 }
394 let remaining = max_bytes.saturating_sub(bytes.len());
395 let retained = read.min(remaining);
396 bytes.extend_from_slice(&buffer[..retained]);
397 truncated |= retained < read;
398 }
399 Ok(BoundedOutput { bytes, truncated })
400 }
401
402 async fn finish_capture_tasks(
403 mut stdout: tokio::task::JoinHandle<std::io::Result<BoundedOutput>>,
404 mut stderr: tokio::task::JoinHandle<std::io::Result<BoundedOutput>>,
405 ) -> Result<(BoundedOutput, BoundedOutput), PdfTextError> {
406 let joined = tokio::time::timeout(PDF_PIPE_DRAIN_TIMEOUT, async {
407 tokio::join!(&mut stdout, &mut stderr)
408 })
409 .await;
410 let (stdout, stderr) = match joined {
411 Ok(output) => output,
412 Err(_) => {
413 stdout.abort();
414 stderr.abort();
415 let _ = tokio::join!(stdout, stderr);
416 return Err(PdfTextError::Execution(
417 "pdftotext output pipes did not close after process termination".to_string(),
418 ));
419 }
420 };
421 let stdout = stdout
422 .map_err(|error| PdfTextError::Execution(format!("stdout reader failed: {error}")))?
423 .map_err(|error| PdfTextError::Execution(format!("stdout reader failed: {error}")))?;
424 let stderr = stderr
425 .map_err(|error| PdfTextError::Execution(format!("stderr reader failed: {error}")))?
426 .map_err(|error| PdfTextError::Execution(format!("stderr reader failed: {error}")))?;
427 Ok((stdout, stderr))
428 }
429
430 fn sanitized_text(bytes: &[u8]) -> String {
431 String::from_utf8_lossy(bytes)
432 .trim()
433 .chars()
434 .map(|character| match character {
435 '\n' | '\t' => character,
436 character if character.is_control() => '\u{fffd}',
437 character => character,
438 })
439 .collect()
440 }
441
442 #[cfg(test)]
443 mod tests;
444
445 #[cfg(all(test, unix))]
446 #[path = "pdf/host_tests.rs"]
447 mod host_tests;
448
448 lines RUST