| 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 |