Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions .github/workflows/_test-integrations.yml
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,7 @@ jobs:
pip install -e '.[test]'
- name: Run Integration Testing
run: |
pytest --cov mindee -m integration
pytest --cov mindee -n 2 -m integration

- name: Notify Slack Action on Failure
uses: ravsamhq/notify-slack-action@2.5.0
Expand All @@ -71,7 +71,7 @@ jobs:
SLACK_WEBHOOK_URL: ${{ secrets.PRODUCTION_ISSUES_SLACK_HOOK_URL }}

pytest-lite:
name: Run Integration Tests
name: Run Integration Tests (Lite)
timeout-minutes: 30
strategy:
matrix:
Expand Down Expand Up @@ -107,4 +107,4 @@ jobs:
shell: bash
- name: Run Integration Testing
run: |
pytest -m "integration and not pypdfium2 and not pillow"
pytest -n 2 -m "integration and not pypdfium2 and not pillow"
2 changes: 1 addition & 1 deletion .github/workflows/_test-regressions.yml
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ jobs:
env:
MINDEE_API_KEY: ${{ secrets.MINDEE_API_KEY_SE_TESTS }}
run: |
pytest --cov mindee -m regression
pytest --cov mindee -n auto -m regression

- name: Notify Slack Action on Failure
uses: ravsamhq/notify-slack-action@2.5.0
Expand Down
8 changes: 8 additions & 0 deletions .github/workflows/_test-units.yml
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,14 @@ name: Test
on:
workflow_call:

env:
MINDEE_V2_SE_TESTS_FINDOC_MODEL_ID: ${{ secrets.MINDEE_V2_SE_TESTS_FINDOC_MODEL_ID }}
MINDEE_V2_SE_TESTS_BLANK_PDF_URL: ${{ secrets.MINDEE_V2_SE_TESTS_BLANK_PDF_URL }}
MINDEE_V2_SE_TESTS_CLASSIFICATION_MODEL_ID: ${{ secrets.MINDEE_V2_SE_TESTS_CLASSIFICATION_MODEL_ID }}
MINDEE_V2_SE_TESTS_CROP_MODEL_ID: ${{ secrets.MINDEE_V2_SE_TESTS_CROP_MODEL_ID }}
MINDEE_V2_SE_TESTS_SPLIT_MODEL_ID: ${{ secrets.MINDEE_V2_SE_TESTS_SPLIT_MODEL_ID }}
MINDEE_V2_SE_TESTS_OCR_MODEL_ID: ${{ secrets.MINDEE_V2_SE_TESTS_OCR_MODEL_ID }}

jobs:
pytest:
name: Run Unit Tests
Expand Down
9 changes: 0 additions & 9 deletions mindee/parsing/common/common_response.py
Original file line number Diff line number Diff line change
@@ -1,18 +1,9 @@
import json
from enum import Enum

from mindee.logger import logger
from mindee.parsing.common.string_dict import StringDict


class CommonStatus(str, Enum):
"""Response status."""

PROCESSING = "Processing"
FAILED = "Failed"
PROCESSED = "Processed"


class CommonResponse:
"""Base class for V1 & V2 responses."""

Expand Down
147 changes: 110 additions & 37 deletions mindee/v2/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,6 @@
from mindee.input.local_input_source import LocalInputSource
from mindee.logger import logger
from mindee.mindee_http.cancellation_token import CancellationToken
from mindee.parsing.common.common_response import CommonStatus
from mindee.v2.client_options.base_annotation_parameters import BaseAnnotationParameters
from mindee.v2.client_options.base_product_parameters import BaseProductParameters
from mindee.v2.client_options.base_rag_document_upload_parameters import (
Expand Down Expand Up @@ -128,50 +127,112 @@ def enqueue_and_get_result(

:return: A valid inference response.
"""
if not params.polling_options:
params.polling_options = PollingOptions()
params.polling_options.validate_settings()
if params.polling_options:
polling_options = params.polling_options
polling_options.validate_settings()
else:
polling_options = PollingOptions()

enqueue_response = self.enqueue(input_source, params)
logger.debug(
"Successfully enqueued document with job ID: %s", enqueue_response.job.id
)

return self._poll_for_result(
initial_response=enqueue_response,
response_class=response_type,
polling_options=polling_options,
cancellation_token=cancellation_token,
)

@staticmethod
def __check_webhooks_done(job_response: JobResponse) -> bool:
"""
Checks if all webhooks associated with a job have finished processing.
"""
are_webhooks_done = all(
webhook.status != "Processing" for webhook in job_response.job.webhooks
)
if are_webhooks_done:
logger.debug("All webhooks are completed.")
return True

logger.debug("Not all webhooks are completed.")
return False

def _poll_on_job(
self,
initial_response: JobResponse,
polling_options: PollingOptions,
wait_for_webhooks: bool,
Comment thread
ianardee marked this conversation as resolved.
cancellation_token: CancellationToken | None = None,
) -> JobResponse:
"""Polls a job until it is processed or the maximum number of tries is reached."""
logger.debug(
"Waiting %s seconds before attempting to retrieve the result...",
polling_options.initial_delay_sec,
)

if cancellation_token and cancellation_token.is_canceled:
raise MindeeError("Request canceled through cancellation token.")
sleep(params.polling_options.initial_delay_sec)

sleep(polling_options.initial_delay_sec)
try_counter = 0
while try_counter < params.polling_options.max_retries:

while try_counter < polling_options.max_retries:
if cancellation_token and cancellation_token.is_canceled:
raise MindeeError("Request canceled through cancellation token.")
job_response = self.get_job(enqueue_response.job.id)

logger.debug(
"Poll attempt %s of %s", try_counter + 1, polling_options.max_retries
)

job_response = self.get_job(initial_response.job.id)
assert isinstance(job_response, JobResponse)
if job_response.job.status == CommonStatus.FAILED.value:

if job_response.job.status == "Processed":
Comment thread
ianardee marked this conversation as resolved.
logger.debug(
"Job ID %s completed processing at: %s",
job_response.job.id,
job_response.job.completed_at,
)
if not wait_for_webhooks or self.__check_webhooks_done(job_response):
return job_response

# normally the mindee_api will throw on error, this is a fallback
if job_response.job.status == "Failed":
if job_response.job.error:
detail = job_response.job.error.detail
else:
detail = "No error detail available."
raise MindeeError(
f"Parsing failed for job {job_response.job.id}: {detail}"
)
if (
job_response.job.status == CommonStatus.PROCESSED.value
and job_response.job.result_url
):
logger.debug(
"Job ID %s completed processing at: %s",
job_response.job.id,
job_response.job.completed_at,
)
result = self.get_result_from_url(
response_type, job_response.job.result_url
)
assert isinstance(result, response_type), (
f'Invalid response type "{type(result)}"'
)
return result

try_counter += 1
sleep(params.polling_options.delay_sec)
sleep(polling_options.delay_sec)

raise MindeeError(f"Couldn't retrieve document after {try_counter} tries.")

raise MindeeError(f"Couldn't retrieve document after {try_counter + 1} tries.")
def _poll_for_result(
self,
initial_response: JobResponse,
response_class: type[TypeBaseInferenceResponse],
polling_options: PollingOptions,
wait_for_webhooks: bool = False,
cancellation_token: CancellationToken | None = None,
) -> TypeBaseInferenceResponse:
"""
Poll until the inference is finished processing or the max number of attempts is reached.
"""
job_response = self._poll_on_job(
initial_response, polling_options, wait_for_webhooks, cancellation_token
)
if not job_response.job.result_url:
raise MindeeError(
"The result URL is undefined. This is a server error, try again later or contact support."
)
return self.get_result_from_url(response_class, job_response.job.result_url)

def upload_rag_document(
self,
Expand All @@ -183,6 +244,7 @@ def upload_rag_document(
You will need to poll until the document is ready for use.
Add a document to the RAG database.
"""
logger.debug("Adding a document to the RAG database")
return self.mindee_api.req_post_rag_document(input_source, parameters)

def upload_and_get_rag_document(
Expand All @@ -195,11 +257,14 @@ def upload_and_get_rag_document(
"""
Add a document to the RAG database and return the initial annotation.
"""
if polling_options is None:
polling_options = PollingOptions()
else:
polling_options.validate_settings()

initial_response = self.upload_rag_document(input_source, parameters)
if initial_response.status != "Processing":
return initial_response
if polling_options is None:
polling_options = PollingOptions()
return self._poll_for_rag_document(
initial_response, polling_options, cancellation_token
)
Expand All @@ -212,6 +277,7 @@ def get_rag_document(
You will need to poll until the document is ready for use.
Get a document's info and annotations from the RAG database.
"""
logger.debug("Getting RAG document ID: %s", document_id)
return self.mindee_api.req_get_rag_annotation(response_type, document_id)

def get_ready_rag_document(
Expand All @@ -224,11 +290,14 @@ def get_ready_rag_document(
"""
Get a document's info and annotations from the RAG database.
"""
if polling_options is None:
polling_options = PollingOptions()
else:
polling_options.validate_settings()

initial_response = self.get_rag_document(response_type, document_id)
if initial_response.status != "Processing":
return initial_response
if polling_options is None:
polling_options = PollingOptions()
return self._poll_for_rag_document(
initial_response, polling_options, cancellation_token
)
Expand All @@ -241,6 +310,7 @@ def update_rag_annotations(
You will need to poll until the document is ready for use.
Update a document's annotations in the RAG database.
"""
logger.debug("Updating RAG document ID: %s", parameters.document_id)
return self.mindee_api.req_patch_rag_annotation(parameters)

def update_and_get_rag_annotations(
Expand All @@ -252,11 +322,14 @@ def update_and_get_rag_annotations(
"""
Update a document's annotations in the RAG database.
"""
if polling_options is None:
polling_options = PollingOptions()
else:
polling_options.validate_settings()

initial_response = self.update_rag_annotations(parameters)
if initial_response.status != "Processing":
return initial_response
if polling_options is None:
polling_options = PollingOptions()
return self._poll_for_rag_document(
initial_response, polling_options, cancellation_token
)
Expand All @@ -266,6 +339,7 @@ def delete_extraction_rag_document(self, document_id: str) -> bool:
Delete a document from the RAG database.
For extraction models only.
"""
logger.debug("Deleting RAG document ID: %s", document_id)
return self.mindee_api.req_delete_extraction_rag_document(document_id)

def _poll_for_rag_document(
Expand All @@ -277,12 +351,11 @@ def _poll_for_rag_document(
"""
Poll until the document is finished processing or the max number of attempts is reached.
"""
logger.info("Polling for RAG document ID: %s", initial_response.id)
polling_options.validate_settings()
logger.debug("Polling for RAG document ID: %s", initial_response.id)
max_retries = polling_options.max_retries + 1

logger.debug(
"Waiting %s seconds before attempting to retrieve the result...",
"Waiting %s seconds before attempting to retrieve the document...",
polling_options.initial_delay_sec,
)

Expand All @@ -296,7 +369,7 @@ def _poll_for_rag_document(
while retry_count < max_retries:
if cancellation_token and cancellation_token.is_canceled:
raise MindeeError("Request canceled through cancellation token.")
logger.info("Poll attempt %s of %s", retry_count, max_retries)
logger.debug("Poll attempt %s of %s", retry_count, max_retries)

response = self.get_rag_document(type(initial_response), document_id)
retry_count += 1
Expand All @@ -305,7 +378,7 @@ def _poll_for_rag_document(
sleep(polling_options.delay_sec)
continue
if response.status == "Failed":
raise MindeeError("Job failed without an error payload.")
raise MindeeError("RAG failed without an error payload.")
return response

raise MindeeError(f"RAG polling not complete after {retry_count - 1} attempts.")
Expand Down
2 changes: 2 additions & 0 deletions mindee/v2/parsing/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
from mindee.v2.parsing.inference.inference_active_options import InferenceActiveOptions
from mindee.v2.parsing.inference.inference_file import InferenceFile
from mindee.v2.parsing.inference.inference_model import InferenceModel
from mindee.v2.parsing.job.job import Job
from mindee.v2.parsing.job.job_response import JobResponse
from mindee.v2.product.extraction.extraction_inference import ExtractionInference
from mindee.v2.product.extraction.extraction_response import ExtractionResponse
Expand All @@ -25,5 +26,6 @@
"InferenceActiveOptions",
"InferenceFile",
"InferenceModel",
"Job",
"JobResponse",
]
1 change: 1 addition & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@ lint = [
test = [
"tomli~=2.4.1",
"pytest>=9.0.3,<9.2.0",
"pytest-xdist~=3.8.0",
"pytest-cov~=7.1.0",
"respx~=0.23.1",
]
Expand Down
16 changes: 0 additions & 16 deletions tests/conftest.py

This file was deleted.

This file was deleted.

Loading
Loading