Skip to content
Open
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
69 changes: 69 additions & 0 deletions client/platform/web-girder/api/waitForFolderDatasetReady.spec.ts
Original file line number Diff line number Diff line change
@@ -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/);
});
});
9 changes: 8 additions & 1 deletion client/platform/web-girder/api/waitForFolderDatasetReady.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -73,7 +74,8 @@ export async function waitForFolderDatasetReady(
): Promise<void> {
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<boolean> {
const { data: folder } = await getFolder(folderId);
Expand Down Expand Up @@ -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) {
Expand Down
4 changes: 4 additions & 0 deletions server/dive_server/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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)
Expand Down
37 changes: 37 additions & 0 deletions server/dive_server/crud_rpc.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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 = {
Expand Down Expand Up @@ -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)
Expand Down
66 changes: 47 additions & 19 deletions server/dive_tasks/convert_images.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
"""
Expand All @@ -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,
Expand Down
2 changes: 2 additions & 0 deletions server/dive_tasks/tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -40,6 +41,7 @@
'convert_images',
'convert_large_images',
'convert_video',
'create_large_image_tiles',
'download_google_drive_zip',
'export_trained_pipeline',
'extract_zip',
Expand Down
Loading