mirror of
https://github.com/MindWorkAI/AI-Studio.git
synced 2026-10-04 23:29:40 +00:00
Fixed attached files silently reaching the AI as empty documents (#904)
This commit is contained in:
1 parent
b9ec13edcf
commit
8f8f788896
31 files changed
+1730
-285
No files matched your search
+651
-107
@@ -8,10 +8,12 @@ use axum::extract::Query;
|
||||
use axum::extract::rejection::QueryRejection;
|
||||
use axum::response::sse::{Event, Sse};
|
||||
use base64::{engine::general_purpose, Engine as _};
|
||||
use calamine::{open_workbook_auto, Reader};
|
||||
use calamine::{open_workbook_auto, Error as CalamineError, Reader};
|
||||
use chardetng::{EncodingDetector, Iso2022JpDetection, Utf8Detection};
|
||||
use encoding_rs::Encoding;
|
||||
use file_format::{FileFormat, Kind};
|
||||
use futures::{Stream, StreamExt};
|
||||
use pdfium_render::prelude::Pdfium;
|
||||
use pdfium_render::prelude::{Pdfium, PdfiumError, PdfiumInternalError};
|
||||
use pptx_to_md::{DiagnosticSeverity, ImageHandlingMode, MarkdownOptions, ParserConfig, PresentationContainer, PresentationFormat, PresentationMetadata, ReadingOrder};
|
||||
use serde::{Deserialize, Deserializer, Serialize};
|
||||
use serde::de::{Error as SerdeError, Visitor};
|
||||
@@ -19,7 +21,7 @@ use std::path::Path;
|
||||
use std::pin::Pin;
|
||||
use std::fmt;
|
||||
use log::{debug, error, warn};
|
||||
use tokio::io::AsyncBufReadExt;
|
||||
use tokio::io::AsyncReadExt;
|
||||
use tokio::sync::mpsc;
|
||||
use tokio_stream::wrappers::ReceiverStream;
|
||||
|
||||
@@ -34,7 +36,23 @@ impl Chunk {
|
||||
pub fn new(content: String, metadata: Metadata) -> Self {
|
||||
Chunk { content, stream_id: String::new(), metadata }
|
||||
}
|
||||
|
||||
|
||||
/// Creates a chunk which reports a failed extraction. Errors travel through the same
|
||||
/// schema as content chunks, so the .NET app is able to deserialize and surface them
|
||||
/// instead of silently treating a failure as empty file content.
|
||||
pub fn from_error(error: &ExtractionError) -> Self {
|
||||
Chunk {
|
||||
content: String::new(),
|
||||
stream_id: String::new(),
|
||||
metadata: Metadata::Error {
|
||||
code: error.code,
|
||||
message: error.message.clone(),
|
||||
page_number: error.page_number,
|
||||
detected_format: error.detected_format.clone(),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
pub fn set_stream_id(&mut self, stream_id: &str) { self.stream_id = stream_id.to_string(); }
|
||||
}
|
||||
|
||||
@@ -60,6 +78,135 @@ pub enum Metadata {
|
||||
slide_number: u32,
|
||||
image: Option<Base64Image>,
|
||||
},
|
||||
|
||||
Error {
|
||||
code: ExtractionErrorCode,
|
||||
message: String,
|
||||
page_number: Option<usize>,
|
||||
detected_format: Option<String>,
|
||||
},
|
||||
}
|
||||
|
||||
/// Classifies why an extraction failed, so the .NET app can tell the user what happened
|
||||
/// instead of showing an empty document.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
|
||||
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
|
||||
pub enum ExtractionErrorCode {
|
||||
/// The request itself was malformed, e.g. a missing query parameter.
|
||||
InvalidRequest,
|
||||
|
||||
FileNotFound,
|
||||
FileNotReadable,
|
||||
|
||||
/// Another process holds the file open and denies us reading it.
|
||||
FileLocked,
|
||||
|
||||
FormatDetectionFailed,
|
||||
NotAValidPdf,
|
||||
NotAValidSpreadsheet,
|
||||
PdfiumUnavailable,
|
||||
PdfEncrypted,
|
||||
PageExtractionFailed,
|
||||
NoTextExtracted,
|
||||
|
||||
/// The content does not match the file extension. This is a notice, not a failure: we read
|
||||
/// the file according to its content and only tell the user about the wrong extension.
|
||||
ExtensionMismatch,
|
||||
|
||||
/// The file was read as text, but its bytes are not text.
|
||||
NotTextContent,
|
||||
|
||||
/// The file is an executable, no matter what its extension claims.
|
||||
ExecutableRejected,
|
||||
|
||||
Unsupported,
|
||||
|
||||
/// Any failure which does not carry a code of its own yet.
|
||||
Internal,
|
||||
}
|
||||
|
||||
/// An extraction failure with a machine-readable code. It implements `std::error::Error`,
|
||||
/// so it travels through the existing boxed error channel and `?` keeps working for the
|
||||
/// error types of the underlying crates.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct ExtractionError {
|
||||
pub code: ExtractionErrorCode,
|
||||
pub message: String,
|
||||
pub page_number: Option<usize>,
|
||||
|
||||
/// The format we identified by looking at the content, e.g. when it contradicts the file
|
||||
/// extension. The app names it so the user learns what the file really is.
|
||||
pub detected_format: Option<String>,
|
||||
}
|
||||
|
||||
impl ExtractionError {
|
||||
pub fn new(code: ExtractionErrorCode, message: impl Into<String>) -> Self {
|
||||
Self { code, message: message.into(), page_number: None, detected_format: None }
|
||||
}
|
||||
|
||||
pub fn on_page(code: ExtractionErrorCode, message: impl Into<String>, page_number: usize) -> Self {
|
||||
Self { code, message: message.into(), page_number: Some(page_number), detected_format: None }
|
||||
}
|
||||
|
||||
/// Creates an error which names the format we identified by looking at the content.
|
||||
pub fn with_detected_format(code: ExtractionErrorCode, message: impl Into<String>, detected_format: &FileFormat) -> Self {
|
||||
Self { code, message: message.into(), page_number: None, detected_format: Some(detected_format.name().to_string()) }
|
||||
}
|
||||
|
||||
/// Recovers the structured error from a boxed error. Errors which do not carry a code
|
||||
/// yet are reported as `Internal`, so every failure reaches the .NET app through the
|
||||
/// same schema.
|
||||
fn from_boxed(error: &(dyn std::error::Error + Send + Sync + 'static)) -> Self {
|
||||
match error.downcast_ref::<ExtractionError>() {
|
||||
Some(extraction_error) => extraction_error.clone(),
|
||||
None => Self::new(ExtractionErrorCode::Internal, error.to_string()),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl fmt::Display for ExtractionError {
|
||||
fn fmt(&self, formatter: &mut fmt::Formatter) -> fmt::Result {
|
||||
match self.page_number {
|
||||
Some(page_number) => write!(formatter, "[{:?}] page {page_number}: {}", self.code, self.message),
|
||||
None => write!(formatter, "[{:?}] {}", self.code, self.message),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl std::error::Error for ExtractionError {}
|
||||
|
||||
/// Detects whether a file system error means that another process holds the file open.
|
||||
///
|
||||
/// Windows answers with `ERROR_SHARING_VIOLATION` (32) or `ERROR_LOCK_VIOLATION` (33). This also
|
||||
/// covers files on a network drive, because the SMB server enforces the lock and the client
|
||||
/// surfaces the very same codes.
|
||||
#[cfg(windows)]
|
||||
fn is_locked_error(error: &std::io::Error) -> bool {
|
||||
matches!(error.raw_os_error(), Some(32) | Some(33))
|
||||
}
|
||||
|
||||
/// Detects whether a file system error means that another process holds the file open.
|
||||
///
|
||||
/// Unix has no distinct error for this. A lock held through an SMB share surfaces as a permission
|
||||
/// problem, which we cannot tell apart from an actual permission problem, so we never claim a file
|
||||
/// is locked here.
|
||||
#[cfg(not(windows))]
|
||||
fn is_locked_error(_error: &std::io::Error) -> bool {
|
||||
false
|
||||
}
|
||||
|
||||
/// Classifies a file system error, so a file which another program holds open is reported as such
|
||||
/// instead of collapsing into a generic read failure.
|
||||
fn classify_io_error(error: &std::io::Error) -> ExtractionErrorCode {
|
||||
if is_locked_error(error) {
|
||||
return ExtractionErrorCode::FileLocked;
|
||||
}
|
||||
|
||||
match error.kind() {
|
||||
std::io::ErrorKind::NotFound => ExtractionErrorCode::FileNotFound,
|
||||
std::io::ErrorKind::InvalidData => ExtractionErrorCode::FormatDetectionFailed,
|
||||
_ => ExtractionErrorCode::FileNotReadable,
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Serialize)]
|
||||
@@ -77,10 +224,27 @@ impl Base64Image {
|
||||
}
|
||||
|
||||
const TO_MARKDOWN: &str = "markdown";
|
||||
|
||||
/// Pandoc's markup-free output format. We do not use it as content, only to find out whether a
|
||||
/// conversion produced any readable text at all.
|
||||
const PANDOC_PLAIN: &str = "plain";
|
||||
|
||||
const DOCX: &str = "docx";
|
||||
const ODT: &str = "odt";
|
||||
const HTML: &str = "html";
|
||||
const IMAGE_SEGMENT_SIZE_IN_CHARS: usize = 8_192; // equivalent to ~ 5500 token
|
||||
|
||||
/// Every PDF file starts with this signature.
|
||||
const PDF_MAGIC: &[u8] = b"%PDF-";
|
||||
|
||||
/// How many bytes we probe to verify the PDF signature. The few extra bytes beyond the
|
||||
/// signature itself make the diagnostics useful when the signature does not match.
|
||||
const PDF_HEADER_PROBE_SIZE: u64 = 8;
|
||||
|
||||
/// Last-resort payload used when even an error event cannot be serialized. It keeps the
|
||||
/// chunk schema intact, so the .NET app never has to parse a bare string.
|
||||
const FALLBACK_ERROR_EVENT_JSON: &str = r#"{"content":"","stream_id":"","metadata":{"Error":{"code":"INTERNAL","message":"The extraction error could not be serialized.","page_number":null,"detected_format":null}}}"#;
|
||||
|
||||
type Result<T> = std::result::Result<T, Box<dyn std::error::Error + Send + Sync>>;
|
||||
type ChunkStream = Pin<Box<dyn Stream<Item = Result<Chunk>> + Send>>;
|
||||
|
||||
@@ -124,6 +288,20 @@ where
|
||||
deserializer.deserialize_any(BoolVisitor)
|
||||
}
|
||||
|
||||
/// Reports an extraction failure as a schema-conformant SSE event, so the .NET app is able
|
||||
/// to deserialize it like any other chunk.
|
||||
fn error_event(error: &ExtractionError, stream_id: Option<&str>) -> Event {
|
||||
let mut chunk = Chunk::from_error(error);
|
||||
if let Some(stream_id) = stream_id {
|
||||
chunk.set_stream_id(stream_id);
|
||||
}
|
||||
|
||||
Event::default().json_data(&chunk).unwrap_or_else(|serialization_error| {
|
||||
error!("Failed to serialize an extraction error event: {serialization_error}");
|
||||
Event::default().data(FALLBACK_ERROR_EVENT_JSON)
|
||||
})
|
||||
}
|
||||
|
||||
pub async fn extract_data(
|
||||
_token: APIToken,
|
||||
query: std::result::Result<Query<ExtractDataQuery>, QueryRejection>,
|
||||
@@ -133,7 +311,7 @@ pub async fn extract_data(
|
||||
Err(e) => {
|
||||
let message = format!("Invalid query for '/retrieval/fs/extract': {e}");
|
||||
warn!("{message}");
|
||||
Err(message)
|
||||
Err(ExtractionError::new(ExtractionErrorCode::InvalidRequest, message))
|
||||
},
|
||||
};
|
||||
|
||||
@@ -142,6 +320,7 @@ pub async fn extract_data(
|
||||
Ok(query) => {
|
||||
let stream_result = stream_data(&query.path, query.extract_images).await;
|
||||
let id_ref = &query.stream_id;
|
||||
let path_ref = &query.path;
|
||||
|
||||
match stream_result {
|
||||
Ok(mut stream) => {
|
||||
@@ -149,11 +328,16 @@ pub async fn extract_data(
|
||||
match chunk {
|
||||
Ok(mut chunk) => {
|
||||
chunk.set_stream_id(id_ref);
|
||||
yield Ok(Event::default().json_data(&chunk).unwrap_or_else(|e| Event::default().data(format!("Error: {e}"))));
|
||||
yield Ok(Event::default().json_data(&chunk).unwrap_or_else(|e| {
|
||||
error!("Failed to serialize a content chunk for '{path_ref}': {e}");
|
||||
error_event(&ExtractionError::new(ExtractionErrorCode::Internal, format!("Failed to serialize a content chunk: {e}")), Some(id_ref))
|
||||
}));
|
||||
},
|
||||
|
||||
Err(e) => {
|
||||
yield Ok(Event::default().json_data(format!("Error: {e}")).unwrap_or_else(|_| Event::default().data(format!("Error: {e}"))));
|
||||
let extraction_error = ExtractionError::from_boxed(e.as_ref());
|
||||
error!("Extraction failed for '{path_ref}': {extraction_error}");
|
||||
yield Ok(error_event(&extraction_error, Some(id_ref)));
|
||||
break;
|
||||
},
|
||||
}
|
||||
@@ -161,13 +345,15 @@ pub async fn extract_data(
|
||||
},
|
||||
|
||||
Err(e) => {
|
||||
yield Ok(Event::default().json_data(format!("Error starting stream: {e}")).unwrap_or_else(|_| Event::default().data(format!("Error starting stream: {e}"))));
|
||||
let extraction_error = ExtractionError::from_boxed(e.as_ref());
|
||||
error!("Could not start the extraction stream for '{path_ref}': {extraction_error}");
|
||||
yield Ok(error_event(&extraction_error, Some(id_ref)));
|
||||
}
|
||||
};
|
||||
},
|
||||
|
||||
Err(e) => {
|
||||
yield Ok(Event::default().json_data(format!("Error starting stream: {e}")).unwrap_or_else(|_| Event::default().data(format!("Error starting stream: {e}"))));
|
||||
Err(extraction_error) => {
|
||||
yield Ok(error_event(&extraction_error, None));
|
||||
},
|
||||
}
|
||||
};
|
||||
@@ -175,18 +361,107 @@ pub async fn extract_data(
|
||||
Sse::new(stream)
|
||||
}
|
||||
|
||||
/// How a file is read.
|
||||
///
|
||||
/// Deriving the route from the extension and from the content separately is what lets us notice
|
||||
/// when the two disagree, instead of trusting a possibly wrong extension blindly.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
enum ExtractionRoute {
|
||||
Pdf,
|
||||
PandocDocx,
|
||||
PandocOdt,
|
||||
PandocHtml,
|
||||
PresentationPptx,
|
||||
PresentationOdp,
|
||||
Spreadsheet,
|
||||
Csv,
|
||||
Text,
|
||||
Image,
|
||||
|
||||
/// The file is an executable and is never read.
|
||||
Executable,
|
||||
|
||||
/// A format we recognize but have no reader for, e.g. the legacy binary Office formats.
|
||||
Unsupported,
|
||||
}
|
||||
|
||||
/// Derives the route from the file extension.
|
||||
fn route_from_extension(ext: &str) -> Option<ExtractionRoute> {
|
||||
match ext {
|
||||
"pdf" => Some(ExtractionRoute::Pdf),
|
||||
DOCX => Some(ExtractionRoute::PandocDocx),
|
||||
ODT => Some(ExtractionRoute::PandocOdt),
|
||||
HTML | "htm" => Some(ExtractionRoute::PandocHtml),
|
||||
"csv" | "tsv" => Some(ExtractionRoute::Csv),
|
||||
"pptx" => Some(ExtractionRoute::PresentationPptx),
|
||||
"odp" => Some(ExtractionRoute::PresentationOdp),
|
||||
"xlsx" | "ods" | "xls" | "xlsm" | "xlsb" | "xla" | "xlam" => Some(ExtractionRoute::Spreadsheet),
|
||||
"jpg" | "jpeg" | "png" | "gif" | "bmp" | "tiff" | "svg" | "webp" | "heic" => Some(ExtractionRoute::Image),
|
||||
|
||||
//
|
||||
// Everything else claims nothing in particular. Text formats end up here on purpose:
|
||||
// their content cannot be identified beyond "this is text", so there is nothing to
|
||||
// contradict. Every extension which does have a reader must be listed above, otherwise
|
||||
// a correctly named file looks like a mismatch.
|
||||
//
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Derives the route from the content we identified.
|
||||
///
|
||||
/// `None` means the content does not point at any particular reader. Such a file keeps whatever
|
||||
/// its extension asks for, and the text reader decides whether the bytes are readable at all.
|
||||
fn route_from_content(fmt: FileFormat) -> Option<ExtractionRoute> {
|
||||
match fmt {
|
||||
FileFormat::PortableDocumentFormat => Some(ExtractionRoute::Pdf),
|
||||
FileFormat::OfficeOpenXmlDocument => Some(ExtractionRoute::PandocDocx),
|
||||
FileFormat::OpendocumentText => Some(ExtractionRoute::PandocOdt),
|
||||
FileFormat::HypertextMarkupLanguage => Some(ExtractionRoute::PandocHtml),
|
||||
FileFormat::OfficeOpenXmlPresentation => Some(ExtractionRoute::PresentationPptx),
|
||||
FileFormat::OpendocumentPresentation => Some(ExtractionRoute::PresentationOdp),
|
||||
|
||||
// Calamine reads the legacy binary spreadsheet format as well:
|
||||
FileFormat::OfficeOpenXmlSpreadsheet
|
||||
| FileFormat::OpendocumentSpreadsheet
|
||||
| FileFormat::MicrosoftExcelSpreadsheet => Some(ExtractionRoute::Spreadsheet),
|
||||
|
||||
FileFormat::PlainText => Some(ExtractionRoute::Text),
|
||||
|
||||
//
|
||||
// The legacy binary Word and PowerPoint formats have no reader here: pptx_to_md only
|
||||
// handles PPTX and ODP, and Pandoc cannot read the binary .doc format at all. Saying so
|
||||
// is better than handing the file to a reader which is bound to fail.
|
||||
//
|
||||
FileFormat::MicrosoftWordDocument | FileFormat::MicrosoftPowerpointPresentation => Some(ExtractionRoute::Unsupported),
|
||||
|
||||
_ => match fmt.kind() {
|
||||
Kind::Executable => Some(ExtractionRoute::Executable),
|
||||
Kind::Image => Some(ExtractionRoute::Image),
|
||||
Kind::Ebook | Kind::Archive | Kind::Compressed => Some(ExtractionRoute::Unsupported),
|
||||
_ => None,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
async fn stream_data(file_path: &str, extract_images: bool) -> Result<ChunkStream> {
|
||||
if !Path::new(file_path).exists() {
|
||||
error!("File does not exist: '{file_path}'");
|
||||
return Err("File does not exist.".into());
|
||||
return Err(ExtractionError::new(ExtractionErrorCode::FileNotFound, format!("The file does not exist: '{file_path}'.")).into());
|
||||
}
|
||||
|
||||
let file_path_clone = file_path.to_owned();
|
||||
let fmt = match FileFormat::from_file(&file_path_clone) {
|
||||
Ok(format) => format,
|
||||
Err(error) => {
|
||||
error!("Failed to determine file format for '{file_path}': {error}");
|
||||
return Err(format!("Failed to determine file format for '{file_path}': {error}").into());
|
||||
//
|
||||
// Detecting the format opens the file, so this is the first place a file which another
|
||||
// program holds open fails. Reporting that as a format problem would send the user
|
||||
// looking in the wrong direction, hence we classify the error instead.
|
||||
//
|
||||
let code = classify_io_error(&error);
|
||||
error!("Failed to read '{file_path}' while determining its file format ({code:?}): {error}");
|
||||
return Err(ExtractionError::new(code, format!("The file could not be read: {error}")).into());
|
||||
},
|
||||
};
|
||||
|
||||
@@ -195,82 +470,155 @@ async fn stream_data(file_path: &str, extract_images: bool) -> Result<ChunkStrea
|
||||
.and_then(|extension| extension.to_str())
|
||||
.map(str::to_ascii_lowercase)
|
||||
.unwrap_or_default();
|
||||
debug!("Extracting data from file: '{file_path}', format: '{fmt:?}', extension: '{ext}'");
|
||||
|
||||
let stream = match ext.as_str() {
|
||||
DOCX | ODT => {
|
||||
let from = if ext == DOCX { "docx" } else { "odt" };
|
||||
convert_with_pandoc(file_path, from, TO_MARKDOWN).await?
|
||||
}
|
||||
|
||||
"csv" | "tsv" => {
|
||||
stream_text_file(file_path, true, Some("csv".to_string())).await?
|
||||
},
|
||||
|
||||
"pptx" => stream_presentation(file_path, extract_images, PresentationFormat::Pptx).await?,
|
||||
"odp" => stream_presentation(file_path, extract_images, PresentationFormat::Odp).await?,
|
||||
|
||||
"xlsx" | "ods" | "xls" | "xlsm" | "xlsb" | "xla" | "xlam" => {
|
||||
stream_spreadsheet_as_csv(file_path).await?
|
||||
}
|
||||
|
||||
_ => match fmt.kind() {
|
||||
Kind::Document => match fmt {
|
||||
FileFormat::PortableDocumentFormat => stream_pdf(file_path).await?,
|
||||
|
||||
FileFormat::MicrosoftWordDocument => {
|
||||
convert_with_pandoc(file_path, "docx", TO_MARKDOWN).await?
|
||||
},
|
||||
|
||||
FileFormat::OfficeOpenXmlDocument => {
|
||||
convert_with_pandoc(file_path, fmt.extension(), TO_MARKDOWN).await?
|
||||
},
|
||||
|
||||
_ => stream_text_file(file_path, false, None).await?,
|
||||
},
|
||||
|
||||
Kind::Ebook => return Err("Ebooks not yet supported".into()),
|
||||
|
||||
Kind::Image => {
|
||||
if !extract_images {
|
||||
return Err("Image extraction is disabled.".into());
|
||||
}
|
||||
|
||||
chunk_image(file_path).await?
|
||||
},
|
||||
|
||||
Kind::Other => match fmt {
|
||||
FileFormat::HypertextMarkupLanguage => {
|
||||
convert_with_pandoc(file_path, fmt.extension(), TO_MARKDOWN).await?
|
||||
},
|
||||
|
||||
_ => stream_text_file(file_path, false, None).await?,
|
||||
},
|
||||
|
||||
Kind::Presentation => match fmt {
|
||||
FileFormat::OfficeOpenXmlPresentation => {
|
||||
stream_presentation(file_path, extract_images, PresentationFormat::Pptx).await?
|
||||
},
|
||||
FileFormat::OpendocumentPresentation => {
|
||||
stream_presentation(file_path, extract_images, PresentationFormat::Odp).await?
|
||||
}
|
||||
|
||||
_ => stream_text_file(file_path, false, None).await?,
|
||||
},
|
||||
|
||||
Kind::Spreadsheet => stream_spreadsheet_as_csv(file_path).await?,
|
||||
|
||||
_ => stream_text_file(file_path, false, None).await?,
|
||||
},
|
||||
|
||||
// The size is part of the diagnostics: a truncated or not-yet-available file on a network
|
||||
// share is what tells a broken extraction apart from a document without text.
|
||||
let file_size = match tokio::fs::metadata(file_path).await {
|
||||
Ok(metadata) => format!("{} bytes", metadata.len()),
|
||||
Err(error) => format!("unknown size ({error})"),
|
||||
};
|
||||
|
||||
debug!("Extracting data from file: '{file_path}', {file_size}, format: '{fmt:?}', extension: '{ext}'");
|
||||
|
||||
let extension_route = route_from_extension(ext.as_str());
|
||||
let content_route = route_from_content(fmt);
|
||||
|
||||
//
|
||||
// The content decides whenever it points at a specific reader and contradicts the extension.
|
||||
// `Text` is excluded on purpose: it is the least specific answer, and letting it win would
|
||||
// cost a `.csv` its CSV fence. When the content says nothing, the extension keeps its say and
|
||||
// the text reader decides whether the bytes are readable at all.
|
||||
//
|
||||
let content_is_specific = matches!(content_route, Some(route) if route != ExtractionRoute::Text);
|
||||
let content_contradicts_extension = content_is_specific && content_route != extension_route;
|
||||
|
||||
let route = match (extension_route, content_route) {
|
||||
_ if content_contradicts_extension => content_route.unwrap(),
|
||||
(Some(from_extension), _) => from_extension,
|
||||
(None, Some(from_content)) => from_content,
|
||||
(None, None) => ExtractionRoute::Text,
|
||||
};
|
||||
|
||||
debug!("Reading '{file_path}' via {route:?} (extension: {extension_route:?}, content: {content_route:?}).");
|
||||
|
||||
match route {
|
||||
ExtractionRoute::Executable => {
|
||||
error!("Refused to read '{file_path}': its content is an executable ({name}).", name = fmt.name());
|
||||
return Err(ExtractionError::with_detected_format(
|
||||
ExtractionErrorCode::ExecutableRejected,
|
||||
format!("The file is an executable ({name}), which is never read.", name = fmt.name()),
|
||||
&fmt,
|
||||
).into());
|
||||
},
|
||||
|
||||
ExtractionRoute::Unsupported => {
|
||||
return Err(ExtractionError::with_detected_format(
|
||||
ExtractionErrorCode::Unsupported,
|
||||
format!("The format '{name}' is not supported.", name = fmt.name()),
|
||||
&fmt,
|
||||
).into());
|
||||
},
|
||||
|
||||
ExtractionRoute::Image if !extract_images => {
|
||||
return Err(ExtractionError::new(ExtractionErrorCode::Unsupported, "Image extraction is disabled.").into());
|
||||
},
|
||||
|
||||
_ => {},
|
||||
}
|
||||
|
||||
let stream = match route {
|
||||
ExtractionRoute::Pdf => stream_pdf(file_path).await?,
|
||||
ExtractionRoute::PandocDocx => convert_with_pandoc(file_path, DOCX, TO_MARKDOWN).await?,
|
||||
ExtractionRoute::PandocOdt => convert_with_pandoc(file_path, ODT, TO_MARKDOWN).await?,
|
||||
ExtractionRoute::PandocHtml => convert_with_pandoc(file_path, HTML, TO_MARKDOWN).await?,
|
||||
ExtractionRoute::PresentationPptx => stream_presentation(file_path, extract_images, PresentationFormat::Pptx).await?,
|
||||
ExtractionRoute::PresentationOdp => stream_presentation(file_path, extract_images, PresentationFormat::Odp).await?,
|
||||
ExtractionRoute::Spreadsheet => stream_spreadsheet_as_csv(file_path).await?,
|
||||
ExtractionRoute::Csv => stream_text_file(file_path, true, Some("csv".to_string())).await?,
|
||||
ExtractionRoute::Text => stream_text_file(file_path, false, None).await?,
|
||||
ExtractionRoute::Image => chunk_image(file_path).await?,
|
||||
|
||||
// Handled above, before any reader was chosen:
|
||||
ExtractionRoute::Executable | ExtractionRoute::Unsupported => unreachable!(),
|
||||
};
|
||||
|
||||
//
|
||||
// The file was readable, but not as its extension claims. We prepend a notice so the user
|
||||
// learns what the file really is, while the content itself is read correctly.
|
||||
//
|
||||
if content_contradicts_extension {
|
||||
warn!("The content of '{file_path}' is '{name}', which does not match its extension '{ext}'.", name = fmt.name());
|
||||
|
||||
let notice = Chunk::from_error(&ExtractionError::with_detected_format(
|
||||
ExtractionErrorCode::ExtensionMismatch,
|
||||
format!("The content is '{name}', which does not match the file extension '{ext}'.", name = fmt.name()),
|
||||
&fmt,
|
||||
));
|
||||
|
||||
let notice_stream = stream! { yield Ok(notice); };
|
||||
return Ok(Box::pin(notice_stream.chain(stream)));
|
||||
}
|
||||
|
||||
Ok(Box::pin(stream))
|
||||
}
|
||||
|
||||
/// How many bytes we inspect for NUL bytes to tell binary content from text.
|
||||
const BINARY_PROBE_SIZE: usize = 8_192;
|
||||
|
||||
/// Reads a text file and decodes it, no matter which encoding it uses.
|
||||
///
|
||||
/// Insisting on UTF-8 is not enough in practice: text files written on Windows are frequently
|
||||
/// encoded in Windows-1252, where umlauts are single bytes which UTF-8 rejects. Such a file used
|
||||
/// to look like it was not text at all.
|
||||
async fn read_text_file(file_path: &str) -> Result<String> {
|
||||
let bytes = tokio::fs::read(file_path).await.map_err(|error| ExtractionError::new(
|
||||
classify_io_error(&error),
|
||||
format!("The file could not be read: {error}"),
|
||||
))?;
|
||||
|
||||
//
|
||||
// A byte order mark is authoritative and also covers UTF-16, which the detector below does not
|
||||
// recognize. We therefore check it first and let `decode` act on it.
|
||||
//
|
||||
if let Some((encoding, _)) = Encoding::for_bom(&bytes) {
|
||||
let (text, _, _) = encoding.decode(&bytes);
|
||||
debug!("Decoded '{file_path}' as {name}, chosen by its byte order mark.", name = encoding.name());
|
||||
return Ok(text.into_owned());
|
||||
}
|
||||
|
||||
//
|
||||
// Without a byte order mark, every byte sequence decodes into *something*, so the decoder can
|
||||
// no longer tell us that a file is binary. NUL bytes do: they do not occur in text, and after
|
||||
// the check above no UTF-16 file can reach this point.
|
||||
//
|
||||
let probe_length = min(bytes.len(), BINARY_PROBE_SIZE);
|
||||
if bytes[..probe_length].contains(&0) {
|
||||
return Err(ExtractionError::new(
|
||||
ExtractionErrorCode::NotTextContent,
|
||||
"The file contains binary data and is not a text file.",
|
||||
).into());
|
||||
}
|
||||
|
||||
//
|
||||
// Both options are about untrusted web content which may run scripts, which is not what we
|
||||
// read here: these are local files the user picked, so allowing both guesses gives the better
|
||||
// detection.
|
||||
//
|
||||
let mut detector = EncodingDetector::new(Iso2022JpDetection::Allow);
|
||||
detector.feed(&bytes, true);
|
||||
|
||||
let (text, encoding, had_errors) = detector.guess(None, Utf8Detection::Allow).decode(&bytes);
|
||||
if had_errors {
|
||||
warn!("Decoding '{file_path}' as {name} replaced malformed sequences.", name = encoding.name());
|
||||
} else {
|
||||
debug!("Decoded '{file_path}' as {name}.", name = encoding.name());
|
||||
}
|
||||
|
||||
Ok(text.into_owned())
|
||||
}
|
||||
|
||||
async fn stream_text_file(file_path: &str, use_md_fences: bool, fence_language: Option<String>) -> Result<ChunkStream> {
|
||||
let file = tokio::fs::File::open(file_path).await?;
|
||||
let reader = tokio::io::BufReader::new(file);
|
||||
let mut lines = reader.lines();
|
||||
let text = read_text_file(file_path).await?;
|
||||
let mut line_number = 0;
|
||||
|
||||
let stream = stream! {
|
||||
@@ -291,10 +639,10 @@ async fn stream_text_file(file_path: &str, use_md_fences: bool, fence_language:
|
||||
};
|
||||
}
|
||||
|
||||
while let Ok(Some(line)) = lines.next_line().await {
|
||||
for line in text.lines() {
|
||||
line_number += 1;
|
||||
yield Ok(Chunk::new(
|
||||
line,
|
||||
line.to_string(),
|
||||
Metadata::Text { line_number }
|
||||
));
|
||||
}
|
||||
@@ -307,7 +655,62 @@ async fn stream_text_file(file_path: &str, use_md_fences: bool, fence_language:
|
||||
Ok(Box::pin(stream))
|
||||
}
|
||||
|
||||
/// Verifies the file really is a PDF before handing it to PDFium. Without this check, a file
|
||||
/// which only carries the `.pdf` extension, or whose bytes are not available, would end up in
|
||||
/// the text branch and silently produce empty content.
|
||||
async fn ensure_pdf_header(file_path: &str) -> Result<()> {
|
||||
let file = tokio::fs::File::open(file_path).await.map_err(|error| ExtractionError::new(
|
||||
classify_io_error(&error),
|
||||
format!("The file could not be opened: {error}"),
|
||||
))?;
|
||||
|
||||
let file_size = file.metadata().await.map_err(|error| ExtractionError::new(
|
||||
classify_io_error(&error),
|
||||
format!("The file size could not be read: {error}"),
|
||||
))?.len();
|
||||
|
||||
let mut header = Vec::with_capacity(PDF_HEADER_PROBE_SIZE as usize);
|
||||
file.take(PDF_HEADER_PROBE_SIZE).read_to_end(&mut header).await.map_err(|error| ExtractionError::new(
|
||||
classify_io_error(&error),
|
||||
format!("The first bytes of the file could not be read: {error}"),
|
||||
))?;
|
||||
|
||||
if header.starts_with(PDF_MAGIC) {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let header_hex = header.iter().map(|byte| format!("{byte:02x}")).collect::<Vec<_>>().join(" ");
|
||||
error!("The file '{file_path}' does not start with the PDF signature; size: {file_size} bytes, first bytes: [{header_hex}].");
|
||||
|
||||
Err(ExtractionError::new(
|
||||
ExtractionErrorCode::NotAValidPdf,
|
||||
format!("The file does not start with the PDF signature. Size: {file_size} bytes, first bytes: [{header_hex}]."),
|
||||
).into())
|
||||
}
|
||||
|
||||
/// Classifies why PDFium refused to open a document, so the cause reaches the user instead of
|
||||
/// collapsing into a generic failure.
|
||||
fn classify_pdf_load_error(error: &PdfiumError) -> ExtractionError {
|
||||
let code = match error {
|
||||
PdfiumError::PdfiumLibraryInternalError(internal_error) => match internal_error {
|
||||
// The document is encrypted or its security settings forbid access:
|
||||
PdfiumInternalError::PasswordError | PdfiumInternalError::SecurityError => ExtractionErrorCode::PdfEncrypted,
|
||||
|
||||
// Pdfium could not read the file itself, e.g. because a network share went away:
|
||||
PdfiumInternalError::FileError => ExtractionErrorCode::FileNotReadable,
|
||||
|
||||
_ => ExtractionErrorCode::NotAValidPdf,
|
||||
},
|
||||
|
||||
_ => ExtractionErrorCode::NotAValidPdf,
|
||||
};
|
||||
|
||||
ExtractionError::new(code, format!("The PDF could not be opened: {error}"))
|
||||
}
|
||||
|
||||
async fn stream_pdf(file_path: &str) -> Result<ChunkStream> {
|
||||
ensure_pdf_header(file_path).await?;
|
||||
|
||||
let path = file_path.to_owned();
|
||||
let (tx, rx) = mpsc::channel(10);
|
||||
|
||||
@@ -315,39 +718,98 @@ async fn stream_pdf(file_path: &str) -> Result<ChunkStream> {
|
||||
let pdfium = match Pdfium::ai_studio_init() {
|
||||
Ok(pdfium) => pdfium,
|
||||
Err(e) => {
|
||||
let _ = tx.blocking_send(Err(e));
|
||||
let _ = tx.blocking_send(Err(ExtractionError::new(
|
||||
ExtractionErrorCode::PdfiumUnavailable,
|
||||
format!("The PDF engine could not be initialized: {e}"),
|
||||
).into()));
|
||||
return;
|
||||
}
|
||||
};
|
||||
let doc = match pdfium.load_pdf_from_file(&path, None) {
|
||||
Ok(document) => document,
|
||||
Err(e) => {
|
||||
let _ = tx.blocking_send(Err(e.into()));
|
||||
let _ = tx.blocking_send(Err(classify_pdf_load_error(&e).into()));
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
let mut number_of_pages = 0;
|
||||
let mut number_of_characters = 0;
|
||||
let mut number_of_failed_pages = 0;
|
||||
let mut receiver_gone = false;
|
||||
|
||||
for (num_page, page) in doc.pages().iter().enumerate() {
|
||||
let page_number = num_page + 1;
|
||||
number_of_pages = page_number;
|
||||
|
||||
let content = match page.text().map(|t| t.all()) {
|
||||
Ok(text_content) => text_content,
|
||||
Err(e) => {
|
||||
let _ = tx.blocking_send(Err(e.into()));
|
||||
//
|
||||
// A single unreadable page must not end the document: we report it as a
|
||||
// non-fatal error chunk and continue with the next page. Sending it as an
|
||||
// `Err` would stop the consumer and silently truncate everything after it.
|
||||
//
|
||||
number_of_failed_pages += 1;
|
||||
warn!("The text of page {page_number} of '{path}' could not be extracted: {e}");
|
||||
|
||||
if tx.blocking_send(Ok(Chunk::from_error(&ExtractionError::on_page(
|
||||
ExtractionErrorCode::PageExtractionFailed,
|
||||
format!("The text of page {page_number} could not be extracted: {e}"),
|
||||
page_number,
|
||||
)))).is_err() {
|
||||
receiver_gone = true;
|
||||
break;
|
||||
}
|
||||
|
||||
continue;
|
||||
}
|
||||
};
|
||||
|
||||
number_of_characters += content.chars().count();
|
||||
|
||||
if tx.blocking_send(Ok(Chunk::new(
|
||||
content,
|
||||
Metadata::Pdf { page_number: num_page + 1 }
|
||||
content,
|
||||
Metadata::Pdf { page_number }
|
||||
))).is_err() {
|
||||
receiver_gone = true;
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
if receiver_gone {
|
||||
debug!("The consumer stopped reading the PDF stream of '{path}' after {number_of_pages} page(s).");
|
||||
return;
|
||||
}
|
||||
|
||||
debug!("Extracted {number_of_characters} character(s) from {number_of_pages} page(s) of '{path}'; failed pages: {number_of_failed_pages}.");
|
||||
|
||||
//
|
||||
// Without this marker, a PDF without a text layer and a broken extraction both arrive as
|
||||
// an empty document, and the AI would answer as if the file had no content at all.
|
||||
//
|
||||
if number_of_characters == 0 {
|
||||
warn!("No text could be extracted from '{path}': {number_of_pages} page(s), {number_of_failed_pages} failed page(s). The PDF may consist of scanned images without a text layer.");
|
||||
|
||||
let _ = tx.blocking_send(Ok(Chunk::from_error(&ExtractionError::new(
|
||||
ExtractionErrorCode::NoTextExtracted,
|
||||
format!("No text could be extracted from {number_of_pages} page(s). The PDF may consist of scanned images without a text layer."),
|
||||
))));
|
||||
}
|
||||
});
|
||||
|
||||
Ok(Box::pin(ReceiverStream::new(rx)))
|
||||
}
|
||||
|
||||
/// Classifies a spreadsheet failure, so an unreadable file, e.g. on a network share which went
|
||||
/// away, is not reported as a corrupt workbook.
|
||||
fn classify_spreadsheet_error_code(error: &CalamineError) -> ExtractionErrorCode {
|
||||
match error {
|
||||
CalamineError::Io(io_error) => classify_io_error(io_error),
|
||||
_ => ExtractionErrorCode::NotAValidSpreadsheet,
|
||||
}
|
||||
}
|
||||
|
||||
async fn stream_spreadsheet_as_csv(file_path: &str) -> Result<ChunkStream> {
|
||||
let path = file_path.to_owned();
|
||||
let (tx, rx) = mpsc::channel(10);
|
||||
@@ -356,7 +818,10 @@ async fn stream_spreadsheet_as_csv(file_path: &str) -> Result<ChunkStream> {
|
||||
let mut workbook = match open_workbook_auto(&path) {
|
||||
Ok(w) => w,
|
||||
Err(e) => {
|
||||
let _ = tx.blocking_send(Err(e.into()));
|
||||
let _ = tx.blocking_send(Err(ExtractionError::new(
|
||||
classify_spreadsheet_error_code(&e),
|
||||
format!("The spreadsheet could not be opened: {e}"),
|
||||
).into()));
|
||||
return;
|
||||
}
|
||||
};
|
||||
@@ -365,7 +830,20 @@ async fn stream_spreadsheet_as_csv(file_path: &str) -> Result<ChunkStream> {
|
||||
let range = match workbook.worksheet_range(&sheet_name) {
|
||||
Ok(r) => r,
|
||||
Err(e) => {
|
||||
let _ = tx.blocking_send(Err(e.into()));
|
||||
//
|
||||
// One unreadable sheet must not end the workbook: we report it as a non-fatal
|
||||
// error chunk and continue with the next sheet. Sending it as an `Err` would
|
||||
// stop the consumer and silently drop all remaining sheets.
|
||||
//
|
||||
warn!("The sheet '{sheet_name}' of '{path}' could not be read: {e}");
|
||||
|
||||
if tx.blocking_send(Ok(Chunk::from_error(&ExtractionError::new(
|
||||
classify_spreadsheet_error_code(&e),
|
||||
format!("The sheet '{sheet_name}' could not be read: {e}"),
|
||||
)))).is_err() {
|
||||
return;
|
||||
}
|
||||
|
||||
continue;
|
||||
}
|
||||
};
|
||||
@@ -421,27 +899,93 @@ async fn convert_with_pandoc(
|
||||
.with_output_format(to)
|
||||
.build()
|
||||
.command.output().await?;
|
||||
|
||||
|
||||
let exit_code = output.status.code();
|
||||
let stderr_text = String::from_utf8_lossy(&output.stderr).trim().to_string();
|
||||
debug!("Pandoc converted '{file_path}' from '{from}' to '{to}': exit={exit_code:?}, {stdout_length} byte(s) of output.", stdout_length = output.stdout.len());
|
||||
|
||||
if !stderr_text.is_empty() {
|
||||
warn!("Pandoc reported while converting '{file_path}': {stderr_text}");
|
||||
}
|
||||
|
||||
if !output.status.success() {
|
||||
return Err(ExtractionError::new(
|
||||
ExtractionErrorCode::Internal,
|
||||
format!("Pandoc failed with exit code {exit_code:?}: {stderr_text}"),
|
||||
).into());
|
||||
}
|
||||
|
||||
let content = String::from_utf8(output.stdout).map_err(|e| ExtractionError::new(
|
||||
ExtractionErrorCode::Internal,
|
||||
format!("The output of Pandoc was not valid UTF-8: {e}"),
|
||||
))?;
|
||||
|
||||
//
|
||||
// Pandoc succeeded, yet nothing came out. Passing that on as content would hand an empty
|
||||
// document to the AI, which is exactly what this whole path must not do.
|
||||
//
|
||||
if content.trim().is_empty() || !pandoc_found_readable_text(file_path, from, &content).await {
|
||||
return Err(ExtractionError::new(
|
||||
ExtractionErrorCode::NoTextExtracted,
|
||||
format!("Pandoc read the file without finding any readable text{separator}{stderr_text}", separator = if stderr_text.is_empty() { "." } else { ": " }),
|
||||
).into());
|
||||
}
|
||||
|
||||
let stream = stream! {
|
||||
if output.status.success() {
|
||||
match String::from_utf8(output.stdout.clone()) {
|
||||
Ok(content) => yield Ok(Chunk::new(
|
||||
content,
|
||||
Metadata::Document {}
|
||||
)),
|
||||
Err(e) => yield Err(e.into()),
|
||||
}
|
||||
} else {
|
||||
yield Err(format!(
|
||||
"Pandoc error: {}",
|
||||
String::from_utf8_lossy(&output.stderr)
|
||||
).into());
|
||||
}
|
||||
yield Ok(Chunk::new(
|
||||
content,
|
||||
Metadata::Document {}
|
||||
));
|
||||
};
|
||||
|
||||
Ok(Box::pin(stream))
|
||||
}
|
||||
|
||||
/// Decides whether a conversion produced actual text rather than just structure.
|
||||
///
|
||||
/// HTML is the one input where markup can masquerade as content: a page which builds its text with
|
||||
/// scripts converts into nothing but fenced divs and class names. That looks like content, yet it
|
||||
/// says nothing, and the AI would be asked to work with it. Pandoc's plain output settles the
|
||||
/// question, because it carries no markup at all. Documents such as `.docx` carry their text
|
||||
/// statically, so the check above is enough for them and they are spared the extra conversion.
|
||||
async fn pandoc_found_readable_text(file_path: &str, from: &str, content: &str) -> bool {
|
||||
if from != HTML {
|
||||
return true;
|
||||
}
|
||||
|
||||
let output = PandocProcessBuilder::new()
|
||||
.with_input_file(file_path)
|
||||
.with_input_format(from)
|
||||
.with_output_format(PANDOC_PLAIN)
|
||||
.build()
|
||||
.command.output().await;
|
||||
|
||||
match output {
|
||||
Ok(output) if output.status.success() => {
|
||||
let has_text = !String::from_utf8_lossy(&output.stdout).trim().is_empty();
|
||||
if !has_text {
|
||||
warn!("'{file_path}' converted into {length} character(s) of pure structure without any readable text.", length = content.trim().len());
|
||||
}
|
||||
|
||||
has_text
|
||||
},
|
||||
|
||||
//
|
||||
// We could not find out, so we do not claim the file is empty. The content we already have
|
||||
// is the better answer than an error we cannot justify.
|
||||
//
|
||||
Ok(output) => {
|
||||
warn!("Could not check '{file_path}' for readable text, Pandoc exited with {code:?}.", code = output.status.code());
|
||||
true
|
||||
},
|
||||
|
||||
Err(e) => {
|
||||
warn!("Could not check '{file_path}' for readable text: {e}");
|
||||
true
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
async fn chunk_image(file_path: &str) -> Result<ChunkStream> {
|
||||
let data = tokio::fs::read(file_path).await?;
|
||||
let base64 = general_purpose::STANDARD.encode(&data);
|
||||
|
||||
Reference in new issue
Block a user