diff --git a/app/tasks/embed_metadata_into_pdf.py b/app/tasks/embed_metadata_into_pdf.py index d5d59254..9e998195 100644 --- a/app/tasks/embed_metadata_into_pdf.py +++ b/app/tasks/embed_metadata_into_pdf.py @@ -7,6 +7,7 @@ import json from app.config import settings from app.tasks.retry_config import BaseTaskWithRetry from app.tasks.finalize_document_storage import finalize_document_storage +from app.utils import task_logger, log_task # Import the shared Celery instance from app.celery_app import celery @@ -39,6 +40,7 @@ def persist_metadata(metadata, final_pdf_path): return json_path @celery.task(base=BaseTaskWithRetry) +@log_task("embed_metadata") def embed_metadata_into_pdf(local_file_path: str, extracted_text: str, metadata: dict): """ Embeds extracted metadata into the PDF's standard metadata fields. @@ -59,8 +61,10 @@ def embed_metadata_into_pdf(local_file_path: str, extracted_text: str, metadata: alt_path = os.path.join(settings.workdir, "tmp", os.path.basename(local_file_path)) if os.path.exists(alt_path): local_file_path = alt_path + task_logger(f"Using alternative path: {local_file_path}", step_name="embed_metadata") else: - print(f"[ERROR] Local file {local_file_path} not found, cannot embed metadata.") + task_logger(f"Local file {local_file_path} not found, cannot embed metadata.", + level="error", step_name="embed_metadata") return {"error": "File not found"} # Work on a safe copy in /tmp @@ -70,9 +74,10 @@ def embed_metadata_into_pdf(local_file_path: str, extracted_text: str, metadata: # Create a safe copy to work on shutil.copy(original_file, processed_file) + task_logger(f"Created working copy at {processed_file}", step_name="embed_metadata") try: - print(f"[DEBUG] Embedding metadata into {processed_file}...") + task_logger(f"Embedding metadata into {processed_file}", step_name="embed_metadata") # Open the PDF doc = fitz.open(processed_file) @@ -87,7 +92,7 @@ def embed_metadata_into_pdf(local_file_path: str, extracted_text: str, metadata: doc.save(processed_file, incremental=True, encryption=fitz.PDF_ENCRYPT_KEEP) doc.close() - print(f"[INFO] Metadata embedded successfully in {processed_file}") + task_logger("Metadata embedded successfully", step_name="embed_metadata") # Use the suggested filename from metadata; if not provided, use the original basename. suggested_filename = metadata.get("filename", os.path.splitext(os.path.basename(local_file_path))[0]) @@ -101,28 +106,34 @@ def embed_metadata_into_pdf(local_file_path: str, extracted_text: str, metadata: # Move the processed file using shutil.move to handle cross-device moves. shutil.move(processed_file, final_file_path) + task_logger(f"Moved processed file to {final_file_path}", step_name="embed_metadata") + # Ensure the temporary file is deleted if it still exists. if os.path.exists(processed_file): os.remove(processed_file) # Persist the metadata into a JSON file with the same base name. json_path = persist_metadata(metadata, final_file_path) - print(f"[INFO] Metadata persisted to {json_path}") + task_logger(f"Metadata persisted to {json_path}", step_name="embed_metadata") # Trigger the next step: final storage. - finalize_document_storage.delay(original_file, final_file_path, metadata) + finalize_doc_task = finalize_document_storage.delay(original_file, final_file_path, metadata) + task_logger(f"Triggered final document storage with task ID: {finalize_doc_task.id}", + step_name="embed_metadata") # After triggering final storage, delete the original file if it is in workdir/tmp. workdir_tmp = os.path.join(settings.workdir, "tmp") if original_file.startswith(workdir_tmp) and os.path.exists(original_file): try: os.remove(original_file) - print(f"[INFO] Deleted original file from {original_file}") + task_logger(f"Deleted original file from {original_file}", step_name="embed_metadata") except Exception as e: - print(f"[ERROR] Could not delete original file {original_file}: {e}") + task_logger(f"Could not delete original file {original_file}: {e}", + level="warning", step_name="embed_metadata") return {"file": final_file_path, "metadata_file": json_path, "status": "Metadata embedded"} except Exception as e: - print(f"[ERROR] Failed to embed metadata into {processed_file}: {e}") + task_logger(f"Failed to embed metadata into {processed_file}: {e}", + level="error", step_name="embed_metadata") return {"error": str(e)} diff --git a/app/tasks/extract_metadata_with_gpt.py b/app/tasks/extract_metadata_with_gpt.py index 18c44aee..dcb4c795 100644 --- a/app/tasks/extract_metadata_with_gpt.py +++ b/app/tasks/extract_metadata_with_gpt.py @@ -2,9 +2,13 @@ import json import re +import os from app.config import settings from app.tasks.retry_config import BaseTaskWithRetry from app.tasks.embed_metadata_into_pdf import embed_metadata_into_pdf +from app.utils import log_task_progress, task_step_logging +from app.database import SessionLocal +from app.models import FileRecord # Import the shared Celery instance from app.celery_app import celery @@ -34,9 +38,13 @@ def extract_json_from_text(text): return None @celery.task(base=BaseTaskWithRetry) +@task_step_logging def extract_metadata_with_gpt(s3_filename: str, cleaned_text: str): """Uses OpenAI to classify document metadata.""" - prompt = f""" + session = SessionLocal() + try: + log_task_progress(session, s3_filename, "Starting metadata extraction") + prompt = f""" You are a specialized document analyzer trained to extract structured metadata from documents. Your task is to analyze the given text and return a well-structured JSON object. @@ -68,7 +76,6 @@ Extracted text: Return only valid JSON with no additional commentary. """ - try: print(f"[DEBUG] Sending classification request for {s3_filename}...") completion = client.chat.completions.create( model=settings.openai_model, @@ -85,6 +92,7 @@ Return only valid JSON with no additional commentary. json_text = extract_json_from_text(content) if not json_text: print(f"[ERROR] Could not find valid JSON in GPT response for {s3_filename}.") + log_task_progress(session, s3_filename, "Failed to extract valid JSON") return {} metadata = json.loads(json_text) @@ -92,9 +100,19 @@ Return only valid JSON with no additional commentary. # Trigger the next step: embedding metadata into the PDF embed_metadata_into_pdf.delay(s3_filename, cleaned_text, metadata) + log_task_progress(session, s3_filename, "Metadata extraction completed") + + # Update database record + file_record = session.query(FileRecord).filter(FileRecord.s3_filename == s3_filename).first() + if file_record: + file_record.metadata = metadata + session.commit() return {"s3_file": s3_filename, "metadata": metadata} except Exception as e: print(f"[ERROR] OpenAI classification failed for {s3_filename}: {e}") + log_task_progress(session, s3_filename, f"Error: {e}") return {} + finally: + session.close() diff --git a/app/tasks/finalize_document_storage.py b/app/tasks/finalize_document_storage.py index 6d5216d0..b1a5b826 100644 --- a/app/tasks/finalize_document_storage.py +++ b/app/tasks/finalize_document_storage.py @@ -2,25 +2,29 @@ from app.config import settings from app.tasks.retry_config import BaseTaskWithRetry -# Import the shared Celery instance from app.celery_app import celery +from app.utils import task_logger, log_task -# 1) Import the aggregator task +# Import the aggregator task from app.tasks.send_to_all import send_to_all_destinations - @celery.task(base=BaseTaskWithRetry) +@log_task("finalize_storage") def finalize_document_storage(original_file: str, processed_file: str, metadata: dict): """ Final storage step after embedding metadata. We will now call 'send_to_all_destinations' to push the final PDF to Dropbox/Nextcloud/Paperless. """ - print(f"[INFO] Finalizing document storage for {processed_file}") + task_logger(f"Finalizing document storage for {processed_file}", step_name="finalize_storage") - # 2) Enqueue uploads to all destinations (Dropbox, Nextcloud, Paperless) - send_to_all_destinations.delay(processed_file) + # Enqueue uploads to all destinations (Dropbox, Nextcloud, Paperless) + send_task = send_to_all_destinations.delay(processed_file) + + task_logger(f"Triggered send to all destinations with task ID: {send_task.id}", + step_name="finalize_storage", status="success") return { "status": "Completed", - "file": processed_file + "file": processed_file, + "send_task_id": send_task.id } diff --git a/app/tasks/process_document.py b/app/tasks/process_document.py index d31e4199..9d4e9407 100644 --- a/app/tasks/process_document.py +++ b/app/tasks/process_document.py @@ -13,7 +13,7 @@ from app.tasks.extract_metadata_with_gpt import extract_metadata_with_gpt from app.celery_app import celery from app.database import SessionLocal from app.models import FileRecord -from app.utils import hash_file +from app.utils import hash_file, log_task_progress, task_step_logging @celery.task(base=BaseTaskWithRetry) @@ -28,77 +28,96 @@ def process_document(original_local_file: str): - Check for embedded text. If present, run local GPT extraction - Otherwise, queue Textract-based OCR """ + task_id = process_document.request.id + log_task_progress(task_id, "process_document", "pending", f"Processing {original_local_file}", file_path=original_local_file) if not os.path.exists(original_local_file): - print(f"[ERROR] File {original_local_file} not found.") + log_task_progress(task_id, "process_document", "failure", f"File {original_local_file} not found.", file_path=original_local_file) return {"error": "File not found"} # 0. Compute the file hash and check for duplicates - filehash = hash_file(original_local_file) - original_filename = os.path.basename(original_local_file) - file_size = os.path.getsize(original_local_file) - mime_type, _ = mimetypes.guess_type(original_local_file) - if not mime_type: - mime_type = "application/octet-stream" + with task_step_logging(task_id, "compute_hash", file_path=original_local_file): + filehash = hash_file(original_local_file) + original_filename = os.path.basename(original_local_file) + file_size = os.path.getsize(original_local_file) + mime_type, _ = mimetypes.guess_type(original_local_file) + if not mime_type: + mime_type = "application/octet-stream" # Acquire DB session in the task + new_record = None with SessionLocal() as db: - existing = db.query(FileRecord).filter_by(filehash=filehash).one_or_none() - if existing: - print(f"[INFO] Duplicate file detected (hash={filehash[:10]}...) Skipping processing.") - return { - "status": "duplicate_file", - "file_id": existing.id, - "detail": "File already processed." - } + with task_step_logging(task_id, "check_duplicates", file_path=original_local_file): + existing = db.query(FileRecord).filter_by(filehash=filehash).one_or_none() + if existing: + log_task_progress(task_id, "process_document", "success", + f"Duplicate file detected (hash={filehash[:10]}...). Skipping processing.", + file_id=existing.id) + return { + "status": "duplicate_file", + "file_id": existing.id, + "detail": "File already processed." + } # Not a duplicate -> insert a new record - new_record = FileRecord( - filehash=filehash, - original_filename=original_filename, - local_filename="", # Will fill in after we move it - file_size=file_size, - mime_type=mime_type, - ) - db.add(new_record) - db.commit() - db.refresh(new_record) + with task_step_logging(task_id, "create_file_record", file_path=original_local_file): + new_record = FileRecord( + filehash=filehash, + original_filename=original_filename, + local_filename="", # Will fill in after we move it + file_size=file_size, + mime_type=mime_type, + ) + db.add(new_record) + db.commit() + db.refresh(new_record) # 1. Generate a UUID-based filename and place it in /workdir/tmp - file_ext = os.path.splitext(original_local_file)[1] - file_uuid = str(uuid.uuid4()) - new_filename = f"{file_uuid}{file_ext}" + with task_step_logging(task_id, "copy_to_workdir", file_id=new_record.id, file_path=original_local_file): + file_ext = os.path.splitext(original_local_file)[1] + file_uuid = str(uuid.uuid4()) + new_filename = f"{file_uuid}{file_ext}" - tmp_dir = os.path.join(settings.workdir, "tmp") - os.makedirs(tmp_dir, exist_ok=True) - new_local_path = os.path.join(tmp_dir, new_filename) + tmp_dir = os.path.join(settings.workdir, "tmp") + os.makedirs(tmp_dir, exist_ok=True) + new_local_path = os.path.join(tmp_dir, new_filename) - # Copy the file instead of moving it - shutil.copy(original_local_file, new_local_path) + # Copy the file instead of moving it + shutil.copy(original_local_file, new_local_path) - # Update the DB with final local filename - new_record.local_filename = new_local_path - db.commit() + # Update the DB with final local filename + new_record.local_filename = new_local_path + db.commit() # 2. Check for embedded text (outside the DB session to avoid long open transactions) - pdf_doc = fitz.open(new_local_path) - has_text = any(page.get_text() for page in pdf_doc) - pdf_doc.close() + with task_step_logging(task_id, "check_embedded_text", file_id=new_record.id, file_path=new_local_path): + pdf_doc = fitz.open(new_local_path) + has_text = any(page.get_text() for page in pdf_doc) + pdf_doc.close() if has_text: - print(f"[INFO] PDF {original_local_file} contains embedded text. Processing locally.") + log_task_progress(task_id, "process_document", "in_progress", + f"PDF {original_local_file} contains embedded text. Processing locally.", + file_id=new_record.id) # Extract text locally extracted_text = "" - pdf_doc = fitz.open(new_local_path) - for page in pdf_doc: - extracted_text += page.get_text("text") + "\n" - pdf_doc.close() + with task_step_logging(task_id, "extract_text_locally", file_id=new_record.id, file_path=new_local_path): + pdf_doc = fitz.open(new_local_path) + for page in pdf_doc: + extracted_text += page.get_text("text") + "\n" + pdf_doc.close() # Call metadata extraction directly + log_task_progress(task_id, "process_document", "success", + "Text extracted locally. Queuing for metadata extraction.", + file_id=new_record.id) extract_metadata_with_gpt.delay(new_filename, extracted_text) - return {"file": new_local_path, "status": "Text extracted locally"} + return {"file": new_local_path, "status": "Text extracted locally", "file_id": new_record.id} # 3. If no embedded text, queue Textract processing + log_task_progress(task_id, "process_document", "success", + "No embedded text found. Queuing for OCR.", + file_id=new_record.id) process_with_textract.delay(new_filename) - return {"file": new_local_path, "status": "Queued for OCR"} + return {"file": new_local_path, "status": "Queued for OCR", "file_id": new_record.id} diff --git a/app/tasks/process_with_textract.py b/app/tasks/process_with_textract.py index d3319545..d15defcc 100644 --- a/app/tasks/process_with_textract.py +++ b/app/tasks/process_with_textract.py @@ -8,6 +8,9 @@ from app.config import settings from app.tasks.retry_config import BaseTaskWithRetry from app.tasks.extract_metadata_with_gpt import extract_metadata_with_gpt from app.celery_app import celery +from app.utils import log_task_progress, task_step_logging +from app.database import SessionLocal +from app.models import FileRecord logger = logging.getLogger(__name__) @@ -30,38 +33,58 @@ def process_with_textract(s3_filename: str): 4. Extracts the text content for metadata processing. 5. Triggers downstream metadata extraction by calling extract_metadata_with_gpt. """ + task_id = process_with_textract.request.id + tmp_file_path = os.path.join(settings.workdir, "tmp", s3_filename) + + # Get the file_id from the database + file_id = None + with SessionLocal() as db: + file_record = db.query(FileRecord).filter( + FileRecord.local_filename == tmp_file_path + ).first() + if file_record: + file_id = file_record.id + + log_task_progress(task_id, "process_with_textract", "pending", + f"Starting OCR for {s3_filename}", file_id, tmp_file_path) + + if not os.path.exists(tmp_file_path): + log_task_progress(task_id, "process_with_textract", "failure", + f"Local file not found: {tmp_file_path}", file_id, tmp_file_path) + raise FileNotFoundError(f"Local file not found: {tmp_file_path}") + try: - tmp_file_path = os.path.join(settings.workdir, "tmp", s3_filename) - if not os.path.exists(tmp_file_path): - raise FileNotFoundError(f"Local file not found: {tmp_file_path}") + with task_step_logging(task_id, "azure_document_intelligence", file_id, tmp_file_path): + # Open and send the document for processing + with open(tmp_file_path, "rb") as f: + poller = document_intelligence_client.begin_analyze_document( + "prebuilt-read", body=f, output=[AnalyzeOutputOption.PDF] + ) + result: AnalyzeResult = poller.result() + operation_id = poller.details["operation_id"] - logger.info(f"Processing {s3_filename} with Azure Document Intelligence OCR.") - - # Open and send the document for processing - with open(tmp_file_path, "rb") as f: - poller = document_intelligence_client.begin_analyze_document( - "prebuilt-read", body=f, output=[AnalyzeOutputOption.PDF] + with task_step_logging(task_id, "retrieve_and_save_searchable_pdf", file_id, tmp_file_path): + # Retrieve the processed searchable PDF + response = document_intelligence_client.get_analyze_result_pdf( + model_id=result.model_id, result_id=operation_id ) - result: AnalyzeResult = poller.result() - operation_id = poller.details["operation_id"] - - # Retrieve the processed searchable PDF - response = document_intelligence_client.get_analyze_result_pdf( - model_id=result.model_id, result_id=operation_id - ) - searchable_pdf_path = tmp_file_path # Overwrite the original PDF location - with open(searchable_pdf_path, "wb") as writer: - writer.writelines(response) - logger.info(f"Searchable PDF saved at: {searchable_pdf_path}") - - # Extract raw text content from the result - extracted_text = result.content if result.content else "" - logger.info(f"Extracted text for {s3_filename}: {len(extracted_text)} characters") + searchable_pdf_path = tmp_file_path # Overwrite the original PDF location + with open(searchable_pdf_path, "wb") as writer: + writer.writelines(response) + + # Extract raw text content from the result + extracted_text = result.content if result.content else "" + log_task_progress(task_id, "process_with_textract", "in_progress", + f"Extracted {len(extracted_text)} characters of text", file_id, tmp_file_path) # Trigger downstream metadata extraction + log_task_progress(task_id, "process_with_textract", "success", + "OCR completed. Queueing metadata extraction.", file_id, tmp_file_path) extract_metadata_with_gpt.delay(s3_filename, extracted_text) return {"s3_file": s3_filename, "searchable_pdf": searchable_pdf_path, "cleaned_text": extracted_text} except Exception as e: + log_task_progress(task_id, "process_with_textract", "failure", + f"Error processing with Azure Document Intelligence: {e}", file_id, tmp_file_path) logger.error(f"Error processing {s3_filename} with Azure Document Intelligence: {e}") raise diff --git a/app/tasks/send_to_all.py b/app/tasks/send_to_all.py index 15f78d42..8c8a3886 100644 --- a/app/tasks/send_to_all.py +++ b/app/tasks/send_to_all.py @@ -4,18 +4,31 @@ from app.celery_app import celery from app.tasks.upload_to_dropbox import upload_to_dropbox from app.tasks.upload_to_nextcloud import upload_to_nextcloud from app.tasks.upload_to_paperless import upload_to_paperless +from app.utils import task_logger, log_task @celery.task +@log_task("send_to_all_destinations") def send_to_all_destinations(file_path: str): """ Fires off tasks to upload a single file to Dropbox, Nextcloud, and Paperless. These tasks run in parallel (Celery returns immediately from each .delay()). """ - upload_to_dropbox.delay(file_path) - upload_to_nextcloud.delay(file_path) - upload_to_paperless.delay(file_path) + task_logger(f"Sending {file_path} to all destinations", step_name="send_to_all") + + dropbox_task = upload_to_dropbox.delay(file_path) + nextcloud_task = upload_to_nextcloud.delay(file_path) + paperless_task = upload_to_paperless.delay(file_path) + + task_logger(f"Enqueued file for all destinations: Dropbox (task: {dropbox_task.id}), " + f"Nextcloud (task: {nextcloud_task.id}), Paperless (task: {paperless_task.id})", + step_name="send_to_all", status="success") return { "status": "All upload tasks enqueued", - "file_path": file_path + "file_path": file_path, + "task_ids": { + "dropbox": dropbox_task.id, + "nextcloud": nextcloud_task.id, + "paperless": paperless_task.id + } } diff --git a/app/tasks/upload_to_dropbox.py b/app/tasks/upload_to_dropbox.py index 448360a6..ad605703 100644 --- a/app/tasks/upload_to_dropbox.py +++ b/app/tasks/upload_to_dropbox.py @@ -6,6 +6,7 @@ import dropbox from app.config import settings from app.tasks.retry_config import BaseTaskWithRetry from app.celery_app import celery +from app.utils import task_logger, log_task def get_dropbox_access_token(): """Refresh the Dropbox access token using the stored refresh token from ENV.""" @@ -25,14 +26,16 @@ def get_dropbox_access_token(): return response.json()["access_token"] else: error_msg = f"Failed to refresh Dropbox token: {response.status_code} - {response.text}" - print(f"[ERROR] {error_msg}") + task_logger(error_msg, level="error", step_name="dropbox_auth") raise Exception(error_msg) @celery.task(base=BaseTaskWithRetry) +@log_task("upload_to_dropbox") def upload_to_dropbox(file_path: str): """Uploads a file to Dropbox using the API.""" if not os.path.exists(file_path): + task_logger(f"File not found: {file_path}", level="error", step_name="dropbox_upload") raise FileNotFoundError(f"File not found: {file_path}") # Extract filename and set target path @@ -41,16 +44,20 @@ def upload_to_dropbox(file_path: str): try: # Get fresh access token + task_logger(f"Getting Dropbox access token", step_name="dropbox_auth") access_token = get_dropbox_access_token() dbx = dropbox.Dropbox(access_token) file_size = os.path.getsize(file_path) chunk_size = 4 * 1024 * 1024 # 4MB chunk size + task_logger(f"Starting upload of {filename} ({file_size} bytes) to Dropbox", step_name="dropbox_upload") + with open(file_path, "rb") as file_data: if file_size <= chunk_size: dbx.files_upload(file_data.read(), dropbox_path) else: + task_logger(f"Using chunked upload for {filename}", step_name="dropbox_upload") upload_session_start_result = dbx.files_upload_session_start(file_data.read(chunk_size)) cursor = dropbox.files.UploadSessionCursor( session_id=upload_session_start_result.session_id, @@ -65,10 +72,10 @@ def upload_to_dropbox(file_path: str): dbx.files_upload_session_append_v2(file_data.read(chunk_size), cursor) cursor.offset = file_data.tell() - print(f"[INFO] Successfully uploaded {filename} to Dropbox at {dropbox_path}.") + task_logger(f"Successfully uploaded {filename} to Dropbox at {dropbox_path}", step_name="dropbox_upload", status="success") return {"status": "Completed", "file": file_path} except Exception as e: - error_msg = f"[ERROR] Failed to upload {filename} to Dropbox: {str(e)}" - print(error_msg) + error_msg = f"Failed to upload {filename} to Dropbox: {str(e)}" + task_logger(error_msg, level="error", step_name="dropbox_upload", status="failure") raise Exception(error_msg) diff --git a/app/tasks/upload_to_nextcloud.py b/app/tasks/upload_to_nextcloud.py index b64c2972..513758ed 100644 --- a/app/tasks/upload_to_nextcloud.py +++ b/app/tasks/upload_to_nextcloud.py @@ -5,12 +5,15 @@ import requests from app.config import settings from app.tasks.retry_config import BaseTaskWithRetry from app.celery_app import celery +from app.utils import task_logger, log_task @celery.task(base=BaseTaskWithRetry) +@log_task("upload_to_nextcloud") def upload_to_nextcloud(file_path: str): """Uploads a file to Nextcloud in the configured folder.""" if not os.path.exists(file_path): + task_logger(f"File not found: {file_path}", level="error", step_name="nextcloud_upload") raise FileNotFoundError(f"File not found: {file_path}") # Extract filename @@ -18,6 +21,8 @@ def upload_to_nextcloud(file_path: str): # Construct the full upload URL nextcloud_url = f"{settings.nextcloud_upload_url}/{settings.nextcloud_folder}/{filename}" + + task_logger(f"Starting upload of {filename} to Nextcloud", step_name="nextcloud_upload") # Read file content with open(file_path, "rb") as file_data: @@ -29,9 +34,10 @@ def upload_to_nextcloud(file_path: str): # Check if upload was successful if response.status_code in (200, 201): - print(f"[INFO] Successfully uploaded {filename} to Nextcloud at {nextcloud_url}.") + task_logger(f"Successfully uploaded {filename} to Nextcloud at {nextcloud_url}", + step_name="nextcloud_upload", status="success") return {"status": "Completed", "file": file_path} else: - error_msg = f"[ERROR] Failed to upload {filename} to Nextcloud: {response.status_code} - {response.text}" - print(error_msg) + error_msg = f"Failed to upload {filename} to Nextcloud: {response.status_code} - {response.text}" + task_logger(error_msg, level="error", step_name="nextcloud_upload", status="failure") raise Exception(error_msg) diff --git a/app/tasks/upload_to_paperless.py b/app/tasks/upload_to_paperless.py index 135c886b..676561a3 100644 --- a/app/tasks/upload_to_paperless.py +++ b/app/tasks/upload_to_paperless.py @@ -10,6 +10,7 @@ from typing import Dict, Any from app.config import settings from app.tasks.retry_config import BaseTaskWithRetry from app.celery_app import celery +from app.utils import task_logger, log_task logger = logging.getLogger(__name__) @@ -46,13 +47,15 @@ def poll_task_for_document_id(task_id: str) -> int: while attempts < POLL_MAX_ATTEMPTS: try: + task_logger(f"Polling Paperless for task {task_id}, attempt {attempts+1}/{POLL_MAX_ATTEMPTS}", + step_name="paperless_poll") resp = requests.get(url, headers=_get_headers(), params={"task_id": task_id}) resp.raise_for_status() tasks_data = resp.json() except requests.exceptions.RequestException as exc: - logger.warning( - "Failed to poll for task_id='%s'. Attempt=%d Error=%s", - task_id, attempts + 1, exc + task_logger( + f"Failed to poll for task_id='{task_id}'. Attempt={attempts + 1}/{POLL_MAX_ATTEMPTS} Error={exc}", + level="warning", step_name="paperless_poll" ) time.sleep(POLL_INTERVAL_SEC) attempts += 1 @@ -67,21 +70,29 @@ def poll_task_for_document_id(task_id: str) -> int: if status == "SUCCESS": doc_str = task_info.get("related_document") if doc_str: + task_logger(f"Task {task_id} completed successfully with document ID: {doc_str}", + step_name="paperless_poll", status="success") return int(doc_str) raise RuntimeError( f"Task {task_id} completed but no doc ID found. Task info: {task_info}" ) elif status == "FAILURE": - raise RuntimeError(f"Task {task_id} failed: {task_info.get('result')}") + error_msg = f"Task {task_id} failed: {task_info.get('result')}" + task_logger(error_msg, level="error", step_name="paperless_poll", status="failure") + raise RuntimeError(error_msg) + else: + task_logger(f"Task {task_id} status: {status}, waiting {POLL_INTERVAL_SEC}s", + step_name="paperless_poll") attempts += 1 time.sleep(POLL_INTERVAL_SEC) - raise TimeoutError( - f"Task {task_id} didn't reach SUCCESS within {POLL_MAX_ATTEMPTS} attempts." - ) + timeout_msg = f"Task {task_id} didn't reach SUCCESS within {POLL_MAX_ATTEMPTS} attempts." + task_logger(timeout_msg, level="error", step_name="paperless_poll", status="failure") + raise TimeoutError(timeout_msg) @celery.task(base=BaseTaskWithRetry) +@log_task("upload_to_paperless") def upload_to_paperless(file_path: str) -> Dict[str, Any]: """ Uploads a PDF to Paperless with minimal metadata (filename and date only). @@ -93,9 +104,11 @@ def upload_to_paperless(file_path: str) -> Dict[str, Any]: Returns a dict with status, the paperless_task_id, paperless_document_id, and file_path. """ if not os.path.exists(file_path): + task_logger(f"File not found: {file_path}", level="error", step_name="paperless_upload") raise FileNotFoundError(f"File not found: {file_path}") base_name = os.path.basename(file_path) + task_logger(f"Starting upload of {base_name} to Paperless", step_name="paperless_upload") # Upload the PDF post_url = _paperless_api_url("/api/documents/post_document/") @@ -106,22 +119,21 @@ def upload_to_paperless(file_path: str) -> Dict[str, Any]: data = {"title": base_name} # Title = Filename (no additional metadata) try: - logger.debug("Posting document to Paperless: file=%s", base_name) + task_logger(f"Posting document to Paperless: file={base_name}", step_name="paperless_upload") resp = requests.post(post_url, headers=_get_headers(), files=files, data=data) resp.raise_for_status() except requests.exceptions.RequestException as exc: - logger.error( - "Failed to upload document '%s' to Paperless. Error: %s. Response=%s", - file_path, exc, getattr(exc.response, "text", "") - ) + error_msg = f"Failed to upload document '{file_path}' to Paperless. Error: {exc}. Response={getattr(exc.response, 'text', '')}" + task_logger(error_msg, level="error", step_name="paperless_upload", status="failure") raise raw_task_id = resp.text.strip().strip('"').strip("'") - logger.info(f"Received Paperless task ID: {raw_task_id}") + task_logger(f"Received Paperless task ID: {raw_task_id}", step_name="paperless_upload") # Poll tasks until success/fail => get doc_id doc_id = poll_task_for_document_id(raw_task_id) - logger.info(f"Document {file_path} successfully ingested => ID={doc_id}") + task_logger(f"Document {file_path} successfully ingested => ID={doc_id}", + step_name="paperless_upload", status="success") return { "status": "Completed", diff --git a/app/utils.py b/app/utils.py index 887d780f..1b960063 100644 --- a/app/utils.py +++ b/app/utils.py @@ -1,7 +1,14 @@ # app/utils.py import hashlib +import logging +import contextlib +from functools import wraps +from typing import Optional, Callable +from celery import Task from app.database import SessionLocal -from app.models import ProcessingLog +from app.models import ProcessingLog, FileRecord + +logger = logging.getLogger(__name__) def hash_file(filepath, chunk_size=65536): """ @@ -18,18 +25,153 @@ def hash_file(filepath, chunk_size=65536): return sha256.hexdigest() - -def log_task_progress(task_id, step_name, status, message=None, file_id=None): +def log_task_progress(task_id: str, step_name: str, status: str, message: Optional[str] = None, + file_id: Optional[int] = None, file_path: Optional[str] = None): """ Logs the progress of a Celery task to the database. + + Parameters: + task_id (str): The Celery task ID + step_name (str): Name of the processing step + status (str): Status of the step ("pending", "in_progress", "success", "failure") + message (str, optional): Additional message or error details + file_id (int, optional): ID of associated FileRecord + file_path (str, optional): Path to file - will attempt to find file_id from path """ - with SessionLocal() as db: - log_entry = ProcessingLog( - task_id=task_id, - step_name=step_name, - status=status, - message=message, - file_id=file_id, - ) - db.add(log_entry) - db.commit() + try: + with SessionLocal() as db: + # If file_path is provided but not file_id, try to look up the file_id + if not file_id and file_path: + file_record = db.query(FileRecord).filter( + FileRecord.local_filename == file_path + ).first() + if file_record: + file_id = file_record.id + + log_entry = ProcessingLog( + task_id=task_id, + step_name=step_name, + status=status, + message=message, + file_id=file_id, + ) + db.add(log_entry) + db.commit() + logger.info(f"Task {task_id} - {step_name}: {status} {message or ''}") + return log_entry.id + except Exception as e: + logger.error(f"Failed to log task progress: {e}") + return None + + +@contextlib.contextmanager +def task_step_logging(task_id: str, step_name: str, file_id: Optional[int] = None, file_path: Optional[str] = None): + """ + Context manager for logging the beginning and end of a task step. + + Example: + with task_step_logging(task.request.id, "extract_text", file_path=pdf_path): + # Do the actual work + text = extract_text_from_pdf(pdf_path) + """ + log_id = log_task_progress(task_id, step_name, "in_progress", + "Starting processing step", file_id, file_path) + try: + yield + log_task_progress(task_id, step_name, "success", + "Successfully completed", file_id, file_path) + except Exception as e: + log_task_progress(task_id, step_name, "failure", + f"Error: {str(e)}", file_id, file_path) + raise # Re-raise the exception after logging + + +def log_task(step_name: str): + """ + Decorator for Celery tasks to automatically log progress. + + Example: + @celery.task + @log_task("process_pdf") + def process_pdf(file_path): + # Task implementation + """ + def decorator(func: Callable): + @wraps(func) + def wrapper(*args, **kwargs): + # Get task_id from Celery's current task + task = wrapper.request if hasattr(wrapper, 'request') else None + task_id = task.id if task else "unknown_task" + + # Try to determine file_id or file_path from arguments + file_path = None + if args and isinstance(args[0], str): + file_path = args[0] # Assume first arg is file path + + # Log start + log_task_progress(task_id, step_name, "pending", "Task queued", file_path=file_path) + + try: + # Log in_progress + log_task_progress(task_id, step_name, "in_progress", "Task started", file_path=file_path) + + # Execute the task + result = func(*args, **kwargs) + + # Log success + log_task_progress(task_id, step_name, "success", "Task completed", file_path=file_path) + + return result + except Exception as e: + # Log failure + log_task_progress(task_id, step_name, "failure", f"Error: {str(e)}", file_path=file_path) + raise # Re-raise the exception + + return wrapper + return decorator + + +def task_logger(message: str, level: str = "info", task_id: str = None, step_name: str = None, + status: str = None, file_path: Optional[str] = None, file_id: Optional[int] = None): + """ + Unified logging function that logs to both console and database. + This replaces print() statements in tasks with proper logging. + + Parameters: + message: The log message + level: Log level (info, error, warning, debug) + task_id: Celery task ID (tries to get from current task if None) + step_name: Name of the processing step + status: Status for database logging (pending, in_progress, success, failure) + file_path: Path to the file being processed + file_id: ID of the FileRecord + + Usage: + task_logger("Processing file", task_id=task.request.id, step_name="process_pdf") + task_logger("Error processing file", level="error") + """ + # Get task_id from current task if not provided + if task_id is None: + from celery._state import get_current_task + current_task = get_current_task() + task_id = current_task.request.id if current_task else "unknown_task" + + # Default step name if not provided + if step_name is None: + step_name = "general" + + # Default status if not provided + if status is None: + if level == "error": + status = "failure" + elif level == "warning": + status = "warning" + else: + status = "in_progress" + + # Log to console + log_method = getattr(logger, level.lower(), logger.info) + log_method(f"[{step_name}] {message}") + + # Log to database + return log_task_progress(task_id, step_name, status, message, file_id, file_path)