Skip to content

Commit 8ee4a4c

Browse files
committed
fix(passport-worker): move PDF conversion config on the worker config side
1 parent 061462a commit 8ee4a4c

6 files changed

Lines changed: 43 additions & 48 deletions

File tree

workers/passport-worker/passport_worker/activities.py

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -38,7 +38,6 @@
3838
ImagePreprocessorConfig,
3939
PassportDetectionArgs,
4040
PassportDetectionResponse,
41-
PDFConverterConfig,
4241
PreprocessingBatches,
4342
)
4443
from .preprocessing import (
@@ -147,30 +146,29 @@ async def convert_to_pdfs(
147146
self,
148147
batch: Path,
149148
project: str,
150-
config: PDFConverterConfig,
151149
*,
152150
progress: Annotated[
153151
AsyncProgressRateHandler | None, Weight(value=_CONVERT_TO_PDF_WEIGHT)
154152
] = None,
155153
) -> tuple[Path, Path]:
156154
worker_config = cast(PassportWorkerConfig, lifespan_worker_config())
155+
config = worker_config.preprocessing.pdfs
157156
cache = lifespan_pdf_converter_cache()
158-
pdf_converter_cache_key = config_cache_key(config)
157+
pdf_converter_cache_key = config_cache_key(config.pdf_converter)
159158
pdf_converter_factory = async_enter_cm(
160-
partial(PDFConverter.from_config, config)
159+
partial(PDFConverter.from_config, config.pdf_converter)
161160
)
162161
pdf_converter = await cache.async_get_or_cache_resource(
163162
pdf_converter_cache_key, pdf_converter_factory
164163
)
165164
workdir = worker_config.paths.workdir
166165
pdfs_root = activity_workdir(workdir, project, act_context=False)
167166
pdfs_root.mkdir(parents=True, exist_ok=True)
168-
max_concurrency = worker_config.preprocessing.pdfs.max_concurrency
169167
successes, errors = await convert_to_pdfs_act(
170168
batch,
171169
pdf_converter,
172170
worker_config.paths,
173-
max_concurrency,
171+
config.max_concurrency,
174172
output_root=pdfs_root,
175173
progress=progress,
176174
)

workers/passport-worker/passport_worker/config.py

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,6 @@
11
from concurrent.futures import ProcessPoolExecutor
2+
from enum import StrEnum
3+
from typing import ClassVar
24

35
import datashare_python
46
from datashare_python.config import (
@@ -8,6 +10,7 @@
810
WorkerConfig,
911
)
1012
from datashare_python.objects import DatashareModel, WorkerPaths
13+
from icij_common.registrable import RegistrableConfig
1114
from pydantic import Field
1215

1316
_ALL_LOGGERS = [datashare_python.__name__, __name__, "__main__"]
@@ -31,9 +34,37 @@ def to_image_preprocessing_executor(self) -> ProcessPoolExecutor:
3134
return ProcessPoolExecutor(max_workers=self.n_processes)
3235

3336

37+
class PDFConverterType(StrEnum):
38+
GOTENBERG = "gotenberg"
39+
40+
41+
class PDFConverterConfigBase(DatashareModel, RegistrableConfig):
42+
registry_key: ClassVar[str] = Field(frozen=True, default="type")
43+
type: ClassVar[PDFConverterType]
44+
45+
46+
class GotenbergPDFConverterConfig(PDFConverterConfigBase):
47+
type: ClassVar[PDFConverterType] = Field(
48+
frozen=True, default=PDFConverterType.GOTENBERG
49+
)
50+
gotenberg_url: str = "http://localhost:3000"
51+
max_retries: int = 5
52+
min_retry_wait_s: float = 5.0
53+
max_retry_wait_s: float = 30.0
54+
max_retry_randomness_s: float = 2.0
55+
56+
57+
# TODO: use a tagged union here when we have more options
58+
PDFConverterConfig = GotenbergPDFConverterConfig
59+
60+
3461
class PDFConversionWorkerConfig(DatashareModel):
3562
max_concurrency: int = 10
3663

64+
pdf_converter: PDFConverterConfig = Field(
65+
default_factory=GotenbergPDFConverterConfig
66+
)
67+
3768

3869
class PreprocessingWorkerConfig(DatashareModel):
3970
target_n_pages_per_batch: int = 200

workers/passport-worker/passport_worker/objects.py

Lines changed: 0 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -60,36 +60,10 @@ class DefaultImagePreprocessorConfig(ImagePreprocessorConfigBase):
6060
ImagePreprocessorConfig = DefaultImagePreprocessorConfig
6161

6262

63-
class PDFConverterType(StrEnum):
64-
GOTENBERG = "gotenberg"
65-
66-
67-
class PDFConverterConfigBase(DatashareModel, RegistrableConfig):
68-
registry_key: ClassVar[str] = Field(frozen=True, default="type")
69-
type: ClassVar[PDFConverterType]
70-
71-
72-
class GotenbergPDFConverterConfig(PDFConverterConfigBase):
73-
type: ClassVar[PDFConverterType] = Field(
74-
frozen=True, default=PDFConverterType.GOTENBERG
75-
)
76-
77-
gotenberg_url: str = "http://localhost:3000"
78-
max_retries: int = 5
79-
min_retry_wait_s: float = 5.0
80-
max_retry_wait_s: float = 30.0
81-
max_retry_randomness_s: float = 2.0
82-
83-
84-
# TODO: use a tagged union here when we have more options
85-
PDFConverterConfig = GotenbergPDFConverterConfig
86-
87-
8863
class PreprocessingConfig(DatashareModel):
8964
images: ImagePreprocessorConfig = Field(
9065
default_factory=DefaultImagePreprocessorConfig
9166
)
92-
pdfs: PDFConverterConfig = Field(default_factory=GotenbergPDFConverterConfig)
9367

9468

9569
class PassportDetectorType(StrEnum):

workers/passport-worker/passport_worker/preprocessing.py

Lines changed: 3 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@
66
from inspect import iscoroutinefunction
77
from pathlib import Path
88
from types import TracebackType
9-
from typing import ClassVar, Protocol, Self, TypeVar
9+
from typing import Protocol, Self, TypeVar
1010

1111
from aiofile import async_open
1212
from datashare_python.objects import ProcessedPage, WorkerPaths
@@ -18,21 +18,19 @@
1818
to_raw_async_progress,
1919
to_raw_sync_progress,
2020
)
21-
from icij_common.registrable import RegistrableConfig, RegistrableFromConfig
21+
from icij_common.registrable import RegistrableFromConfig
2222
from passport_service import GotenbergClient
2323
from passport_service.constants import Colorspace
2424
from passport_service.core import process_image, process_pdf
2525
from passport_service.exceptions import UnsupportedDocExtension
2626
from passport_service.utils import run_with_concurrency
27-
from pydantic import Field
2827

28+
from passport_worker.config import GotenbergPDFConverterConfig, PDFConverterType
2929
from passport_worker.constants import pil_supported_extensions
3030
from passport_worker.objects import (
3131
DefaultImagePreprocessorConfig,
3232
FileProcessingError,
33-
GotenbergPDFConverterConfig,
3433
ImagePreprocessorType,
35-
PDFConverterType,
3634
ProcessedFile,
3735
)
3836

@@ -76,11 +74,6 @@ def _from_config(cls, config: DefaultImagePreprocessorConfig, **extras) -> Self:
7674
return cls(config)
7775

7876

79-
class PDFConverterConfig(RegistrableConfig):
80-
registry_key: ClassVar[str] = Field(frozen=True, default="type")
81-
type: ClassVar[PDFConverterType]
82-
83-
8477
class PDFConverter(RegistrableFromConfig):
8578
max_concurrency: int = 10
8679

workers/passport-worker/passport_worker/workflows.py

Lines changed: 3 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,6 @@
1616
ImagePreprocessorConfig,
1717
PassportDetectionArgs,
1818
PassportDetectionResponse,
19-
PDFConverterConfig,
2019
PreprocessingBatches,
2120
)
2221

@@ -99,7 +98,7 @@ async def preprocess(
9998
preprocessing_batches.images, args.project, args.config.preprocessing.images
10099
)
101100
convert_to_pdf_tasks = _convert_to_pdfs_tasks(
102-
preprocessing_batches.to_pdf, args.project, args.config.preprocessing.pdfs
101+
preprocessing_batches.to_pdf, args.project
103102
)
104103
im_preprocessing_tasks = asyncio.gather(*im_preprocessing_tasks)
105104
convert_to_pdf_tasks = asyncio.gather(*convert_to_pdf_tasks)
@@ -157,15 +156,13 @@ def _im_processing_tasks(
157156
return im_preprocessing_tasks
158157

159158

160-
def _convert_to_pdfs_tasks(
161-
batches: Batches, project: str, config: PDFConverterConfig
162-
) -> list[Coroutine]:
159+
def _convert_to_pdfs_tasks(batches: Batches, project: str) -> list[Coroutine]:
163160
all_tasks = []
164161
for b in batches:
165162
all_tasks.append(
166163
execute_activity(
167164
PassportDetectionActivities.convert_to_pdfs,
168-
args=(b, project, config),
165+
args=(b, project),
169166
task_queue=TaskQueue.IO,
170167
start_to_close_timeout=_CONVERT_TO_PDF_TIMEOUT,
171168
)

workers/passport-worker/uv.lock

Lines changed: 2 additions & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

0 commit comments

Comments
 (0)