diff --git a/client/platform/web-girder/api/waitForFolderDatasetReady.spec.ts b/client/platform/web-girder/api/waitForFolderDatasetReady.spec.ts new file mode 100644 index 000000000..daea0e606 --- /dev/null +++ b/client/platform/web-girder/api/waitForFolderDatasetReady.spec.ts @@ -0,0 +1,69 @@ +import { + beforeEach, describe, expect, it, vi, +} from 'vitest'; + +import girderRest from 'platform/web-girder/plugins/girder'; +import { getFolder } from './girder.service'; +import { waitForFolderDatasetReady } from './waitForFolderDatasetReady'; + +vi.mock('platform/web-girder/plugins/girder', () => ({ + default: { get: vi.fn() }, +})); +vi.mock('./girder.service', () => ({ + getFolder: vi.fn(), + getItemsInFolder: vi.fn(), +})); + +const RUNNING = 2; +const SUCCESS = 3; + +/** A job that reports one more unit of progress on every poll, reaching `total`. */ +function progressingJob(total: number) { + let current = 0; + return () => { + current = Math.min(current + 1, total); + return { + data: { + status: current < total ? RUNNING : SUCCESS, + progress: { current, total }, + }, + }; + }; +} + +describe('waitForFolderDatasetReady', () => { + beforeEach(() => { + vi.mocked(girderRest.get).mockReset(); + vi.mocked(getFolder).mockReset(); + }); + + it('keeps waiting past the timeout while a job is still making progress', async () => { + const job = progressingJob(40); + let polls = 0; + vi.mocked(girderRest.get).mockImplementation(async () => job()); + // Not ready until the job has run for far longer than timeoutMs. + vi.mocked(getFolder).mockImplementation(async () => { + polls += 1; + return { data: { meta: { annotate: polls > 40 } } } as never; + }); + + await expect(waitForFolderDatasetReady( + 'folder', + { pollIntervalMs: 5, timeoutMs: 50 }, + ['job'], + )).resolves.toBeUndefined(); + }); + + it('times out when the job stops making progress', async () => { + vi.mocked(girderRest.get).mockResolvedValue({ + data: { status: RUNNING, progress: { current: 3, total: 10 } }, + }); + vi.mocked(getFolder).mockResolvedValue({ data: { meta: {} } } as never); + + await expect(waitForFolderDatasetReady( + 'folder', + { pollIntervalMs: 5, timeoutMs: 50 }, + ['job'], + )).rejects.toThrow(/Timed out/); + }); +}); diff --git a/client/platform/web-girder/api/waitForFolderDatasetReady.ts b/client/platform/web-girder/api/waitForFolderDatasetReady.ts index 45843d8db..025565dc9 100644 --- a/client/platform/web-girder/api/waitForFolderDatasetReady.ts +++ b/client/platform/web-girder/api/waitForFolderDatasetReady.ts @@ -61,6 +61,7 @@ export async function waitForFolderDatasetReady( folderId: string, options?: { pollIntervalMs?: number; + /** Give up after this long without job progress. */ timeoutMs?: number; /** Called with average job completion fraction in [0, 1] when jobs report progress. */ onProgress?: (fraction: number) => void; @@ -73,7 +74,8 @@ export async function waitForFolderDatasetReady( ): Promise { const pollIntervalMs = options?.pollIntervalMs ?? 1000; const timeoutMs = options?.timeoutMs ?? 10 * 60 * 1000; - const deadline = Date.now() + timeoutMs; + let deadline = Date.now() + timeoutMs; + let lastProgress = ''; async function folderReady(): Promise { const { data: folder } = await getFolder(folderId); @@ -108,6 +110,11 @@ export async function waitForFolderDatasetReady( return job; }), ); + const progress = JSON.stringify(jobs.map((job) => [job.status, job.progress?.current])); + if (progress !== lastProgress) { + lastProgress = progress; + deadline = Date.now() + timeoutMs; + } if (options?.onProgress) { const jobsWithProgress = jobs.filter((job) => job.progress?.total); if (jobsWithProgress.length) { diff --git a/server/dive_server/__init__.py b/server/dive_server/__init__.py index 0cda82a9c..7f635fcf4 100644 --- a/server/dive_server/__init__.py +++ b/server/dive_server/__init__.py @@ -6,9 +6,11 @@ from girder.constants import AccessType from girder.models.user import User from girder.plugin import getPlugin +from girder.settings import SettingDefault from girder.utility import mail_utils from girder.utility.model_importer import ModelImporter from girder_jobs.models.job import Job +from girder_large_image.constants import PluginSettings as LargeImageSettings from dive_utils import constants @@ -42,6 +44,8 @@ def load(self, info): info["apiRoot"].dive_scoring = ScoringResource("dive_scoring") # required because girder doesn't load plugins in order so we need to manually load first. getPlugin('jobs').load(info) + # Tile in a postprocess job, not per file on upload; an admin setting still wins. + SettingDefault.defaults[LargeImageSettings.LARGE_IMAGE_AUTO_SET] = False # Setup route additions for exsting resources info['apiRoot'].job.route("GET", ("queued",), countJobs) info["apiRoot"].user.route("PUT", (":id", "use_private_queue"), use_private_queue) diff --git a/server/dive_server/crud_rpc.py b/server/dive_server/crud_rpc.py index 5a337f809..bfdbcd1cb 100644 --- a/server/dive_server/crud_rpc.py +++ b/server/dive_server/crud_rpc.py @@ -1606,6 +1606,7 @@ def _postprocess( When skipJobs=False, the following may run as jobs: Transcoding of Video Transcoding of Images + Tile creation for large images (TIFF, NITF, ...) Conversion of KPF annotations into track JSON Extraction and upload of zip files @@ -1737,6 +1738,13 @@ def _postprocess( largeImageItems = Folder().childItems( dsFolder, filters={"lowerName": {"$regex": constants.largeImageRegEx}} ) + untiledLargeImageItems = Folder().childItems( + dsFolder, + filters={ + "lowerName": {"$regex": constants.largeImageRegEx}, + "largeImage": {"$exists": False}, + }, + ) if imageItems.count() > safeImageItems.count(): convert_params = { @@ -1766,6 +1774,35 @@ def _postprocess( ) created_job_ids.append(job['_id']) + elif untiledLargeImageItems.count() > 0: + # auto_set is off, so large images get their tiles here. + tiles_params = { + 'user_id': str(user["_id"]), + 'user_login': str(user["login"]), + 'input_folder': str(dsFolder["_id"]), + } + newjob = tasks.create_large_image_tiles.apply_async( + queue=_get_queue_name(user), + kwargs=dict( + folderId=str(dsFolder["_id"]), + user_id=str(user["_id"]), + user_login=str(user["login"]), + girder_client_token=str(token["_id"]), + girder_job_title=f"Preparing {dsFolder['name']} large images for viewing", + girder_job_type="private" if job_is_private else "convert", + ), + ) + job = _persist_async_job_metadata( + newjob, + **{ + constants.JOBCONST_PRIVATE_QUEUE: job_is_private, + constants.JOBCONST_DATASET_ID: dsFolder["_id"], + constants.JOBCONST_PARAMS: tiles_params, + constants.JOBCONST_CREATOR: str(user['_id']), + }, + ) + created_job_ids.append(job['_id']) + elif imageItems.count() > 0 or largeImageItems.count() > 0: # Safe/web images need annotate now; convert_images sets it when it runs. crud.refresh_folder_document(dsFolder) diff --git a/server/dive_tasks/convert_images.py b/server/dive_tasks/convert_images.py index e52619c5c..dafb488f3 100644 --- a/server/dive_tasks/convert_images.py +++ b/server/dive_tasks/convert_images.py @@ -215,6 +215,28 @@ def convert_images(self: Task, folderId, user_id: str, user_login: str): ) +def _ensure_large_image(gc: GirderClient, manager: JobManager, item: GirderModel) -> None: + """Create the item's large-image tile metadata unless it already has it.""" + try: + gc.get(f'item/{item["_id"]}/tiles') + manager.write(f'Skipping {item["name"]}, already a large image\n') + return + except HttpError as e: + # Safely parse JSON if possible + message = "" + try: + message = e.response.json().get("message", "") + except Exception: + pass # non-JSON response, leave message empty + # This is the Girder message when no large image exists + if e.status == 400 and message == "No large image file in this item.": + manager.write(f'Converting {item["name"]} to large image\n') + gc.post(f'item/{item["_id"]}/tiles') + else: + # Re-raise unexpected errors to fail the job + raise + + @app.task(bind=True, acks_late=True) def convert_large_images(self: Task, folderId, user_id: str, user_login: str): """ @@ -236,31 +258,37 @@ def convert_large_images(self: Task, folderId, user_id: str, user_login: str): ] for item in items_to_convert: # Assumes 1 file per item - try: - # Does it already have tiles? - gc.get(f'item/{item["_id"]}/tiles') - manager.write(f'Skipping {item["name"]}, already a large image\n') - continue - except HttpError as e: - # Safely parse JSON if possible - message = "" - try: - message = e.response.json().get("message", "") - except Exception: - pass # non-JSON response, leave message empty - # This is the Girder message when no large image exists - if e.status == 400 and message == "No large image file in this item.": - manager.write(f'Converting {item["name"]} to large image\n') - gc.post(f'item/{item["_id"]}/tiles') - else: - # Re-raise unexpected errors to fail the job - raise + _ensure_large_image(gc, manager, item) gc.addMetadataToFolder( str(folderId), {"type": constants.LargeImageType}, # mark the parent folder as able to annotate. ) +@app.task(bind=True, acks_late=True, ignore_result=True) +def create_large_image_tiles(self: Task, folderId, user_id: str, user_login: str): + """Create tiles for every large-image file in the folder, then mark it a dataset.""" + context: dict = {} + gc: GirderClient = self.girder_client + manager: JobManager = patch_manager(self.job_manager) + if utils.check_canceled(self, context): + manager.updateStatus(JobStatus.CANCELED) + return + + items = [ + item for item in gc.listItem(folderId) if constants.largeImageRegEx.search(item["name"]) + ] + for index, item in enumerate(items): + if utils.check_canceled(self, context, force=False): + manager.updateStatus(JobStatus.CANCELED) + return + manager.updateProgress(total=len(items), current=index) + # Assumes 1 file per item + _ensure_large_image(gc, manager, item) + manager.updateProgress(total=len(items), current=len(items), forceFlush=True) + gc.addMetadataToFolder(str(folderId), {constants.DatasetMarker: True}) + + @app.task(bind=True, acks_late=True, ignore_result=True) def extract_zip( self: Task, diff --git a/server/dive_tasks/tasks.py b/server/dive_tasks/tasks.py index 931348ac5..315a6f195 100644 --- a/server/dive_tasks/tasks.py +++ b/server/dive_tasks/tasks.py @@ -9,6 +9,7 @@ convert_calibration, convert_images, convert_large_images, + create_large_image_tiles, extract_zip, ) from dive_tasks.convert_video import convert_video, resolve_annotation_fps @@ -40,6 +41,7 @@ 'convert_images', 'convert_large_images', 'convert_video', + 'create_large_image_tiles', 'download_google_drive_zip', 'export_trained_pipeline', 'extract_zip',