SWE-Race › Tasks › mapillary-mapillary-tools-779 ← prevnext →

mapillary-mapillary-tools-779

mapillary/mapillary_toolssplitsinglemerged 2025-08-27BSD-2-Clausefix: 2 files, +86 −6010 fail-to-pass · 1 pass-to-pass
Results
Modelsolved / attemptsmedian stepsmedian costattempts
GPT-5.6 Luna6/614$0.0141✓ 2✓ 3✓ 4✓ 5✓ 6✓
DeepSeek V4 Flash2/256$0.0931✓ 2✓
GLM-5.3 Flash1/246$0.0381✗ 2✓
The prompt the agent sees

Image uploads are unreliable when `ImageSequenceUploader` processes multiple images, especially with an upload cache and internal parallelism. Repeated or concurrent uploads can race while accessing shared cached state, causing images to be uploaded more than once, results to be missing or associated with the wrong sequence, or otherwise producing inconsistent outcomes.

The uploader must successfully process a basic sequence and return one successful result for that sequence. Multiple sequences must remain independent, without dropped, merged, or mislabelled results. Uploads must remain correct with parallel image uploads whether caching is disabled or enabled, including sequences containing many images. A later run must reuse completed uploads instead of uploading the same images again while still returning successful results. The general image-upload operation must return complete, successful results for every supplied sequence.

The expected sequence-upload and image-upload events must be emitted, including when an `uploader.EventEmitter` is supplied. Supported ZIP-based image input must continue to upload successfully and emit its corresponding events when requested. Existing upload result types, event behavior, sequence associations, supported input formats, and repeated-invocation behavior must be preserved.

`uploader.UploadOptions` must have an optional dataclass field `upload_cache_path: Path | None = None`, accepted as a keyword argument. `ImageSequenceUploader` must expose its `upload_options` and a `cached_image_uploader` attribute. The `cached_image_uploader` must expose a `cache` attribute containing a `history.PersistentCache` instance or `None`, as determined by the configured options. The cache used by image-upload operations must be shared consistently and support reuse across runs.

Cache configuration must support an explicitly supplied `upload_cache_path` and, when no path is supplied, the configured default location, while allowing caching to be disabled. Whether a cache exists is determined once from the options used to construct the `ImageSequenceUploader`: with `dry_run=True`, `cached_image_uploader.cache` is `None` even if `upload_cache_path` is supplied; with `dry_run=False` and a usable cache path, the cache is present. If the uploader's or the cached uploader's `upload_options` is later replaced by a `dry_run=True` copy, an already-created cache remains present and retains its entries. Cache creation and access must remain safe during concurrent image uploads.

The cached uploader must provide `_get_cached_file_handle(test_key)` and `_set_file_handle_cache(test_key, test_value)`. When caching is disabled, `_get_cached_file_handle` returns `None` and `_set_file_handle_cache` completes without raising an exception or storing a value.

`history.PersistentCache` must provide a `keys()` method that returns its stored keys, so callers can verify an empty cache and inspect entries after cached uploads. Initially, a newly configured cache has no keys; after successful uploads, it contains the corresponding entries, and subsequent uploads can use those entries without duplicating the uploads.

Hidden tests · 10 fail-to-pass, 1 pass-to-passrun after the agent submits, in a clean verifier
test_image_sequence_uploader_basictest_image_sequence_uploader_cache_hits_second_runtest_image_sequence_uploader_event_emissiontest_image_sequence_uploader_multiple_sequencestest_image_sequence_uploader_multithreading_with_cache_disabtest_image_sequence_uploader_multithreading_with_cache_enabltest_upload_imagestest_upload_images_multiple_sequences+2 more
Test patch · 541 lines
diff --git a/tests/unit/test_uploader.py b/tests/unit/test_uploader.py
index c4a9bbcb..fa2f3a72 100644
--- a/tests/unit/test_uploader.py
+++ b/tests/unit/test_uploader.py
@@ -1,8 +1,9 @@
+import dataclasses
 import typing as T
 from pathlib import Path
+from unittest.mock import patch
 
 import py.path
-
 import pytest
 
 from mapillary_tools import api_v4, uploader
@@ -27,7 +28,9 @@ def setup_unittest_data(tmpdir: py.path.local):
 def test_upload_images(setup_unittest_data: py.path.local, setup_upload: py.path.local):
     mly_uploader = uploader.Uploader(
         uploader.UploadOptions(
-            {"user_upload_token": "YOUR_USER_ACCESS_TOKEN"}, dry_run=True
+            {"user_upload_token": "YOUR_USER_ACCESS_TOKEN"},
+            dry_run=True,
+            upload_cache_path=Path(setup_unittest_data.join("upload_cache")),
         )
     )
     test_exif = setup_unittest_data.join("test_exif.jpg")
@@ -106,6 +109,7 @@ def test_upload_images_multiple_sequences(
                 # will call the API for real
                 # "MAPOrganizationKey": "3011753992432185",
             },
+            upload_cache_path=Path(setup_unittest_data.join("upload_cache")),
             dry_run=True,
         ),
     )
@@ -177,6 +181,7 @@ def test_upload_zip(
                 # will call the API for real
                 # "MAPOrganizationKey": 3011753992432185,
             },
+            upload_cache_path=Path(setup_unittest_data.join("upload_cache")),
             dry_run=True,
         ),
         emitter=emitter,
@@ -252,3 +257,498 @@ def _upload_end(payload):
     test_upload_zip(setup_unittest_data, setup_upload, emitter=emitter)
 
     assert len(stats) == 2, stats
+
+
+class TestImageSequenceUploader:
+    """Test suite for ImageSequenceUploader with focus on multithreading scenarios and caching."""
+
+    def test_image_sequence_uploader_basic(self, setup_unittest_data: py.path.local):
+        """Test basic functionality of ImageSequenceUploader."""
+        upload_options = uploader.UploadOptions(
+            {"user_upload_token": "YOUR_USER_ACCESS_TOKEN"},
+            upload_cache_path=Path(setup_unittest_data.join("upload_cache")),
+            dry_run=True,
+        )
+        emitter = uploader.EventEmitter()
+        sequence_uploader = uploader.ImageSequenceUploader(upload_options, emitter)
+
+        # Create mock image metadata for a single sequence
+        test_exif = setup_unittest_data.join("test_exif.jpg")
+        image_metadatas = [
+            description.DescriptionJSONSerializer.from_desc(
+                {
+                    "MAPLatitude": 58.5927694,
+                    "MAPLongitude": 16.1840944,
+                    "MAPCaptureTime": "2021_02_13_13_24_41_140",
+                    "filename": str(test_exif),
+                    "filetype": "image",
+                    "MAPSequenceUUID": "sequence_1",
+                }
+            ),
+            description.DescriptionJSONSerializer.from_desc(
+                {
+                    "MAPLatitude": 58.5927695,
+                    "MAPLongitude": 16.1840945,
+                    "MAPCaptureTime": "2021_02_13_13_24_42_140",
+                    "filename": str(test_exif),
+                    "filetype": "image",
+                    "MAPSequenceUUID": "sequence_1",
+                }
+            ),
+        ]
+
+        # Test upload
+        results = list(sequence_uploader.upload_images(image_metadatas))
+
+        assert len(results) == 1
+        sequence_uuid, upload_result = results[0]
+        assert sequence_uuid == "sequence_1"
+        assert upload_result.error is None
+        assert upload_result.result is not None
+
+    def test_image_sequence_uploader_multithreading_with_cache_enabled(
+        self, setup_unittest_data: py.path.local
+    ):
+        """Test that ImageSequenceUploader's internal multithreading works correctly when cache is enabled."""
+        # Create upload options that enable cache
+        upload_options_with_cache = uploader.UploadOptions(
+            {"user_upload_token": "YOUR_USER_ACCESS_TOKEN"},
+            upload_cache_path=Path(setup_unittest_data.join("upload_cache")),
+            num_upload_workers=4,  # This will be used internally for parallel image uploads
+            dry_run=False,  # Cache requires dry_run=False initially
+        )
+        emitter = uploader.EventEmitter()
+        sequence_uploader = uploader.ImageSequenceUploader(
+            upload_options_with_cache, emitter
+        )
+
+        # Override to dry_run=True for actual testing
+        sequence_uploader.upload_options = dataclasses.replace(
+            upload_options_with_cache, dry_run=True
+        )
+        sequence_uploader.cached_image_uploader.upload_options = dataclasses.replace(
+            upload_options_with_cache, dry_run=True
+        )
+
+        # Verify cache is available and shared
+        assert sequence_uploader.cached_image_uploader.cache is not None, (
+            "SingleImageUploader should share the same cache instance"
+        )
+
+        test_exif = setup_unittest_data.join("test_exif.jpg")
+
+        num_images = 100  # Reasonable number for testing with direct cache verification
+        image_metadatas = []
+
+        for i in range(num_images):
+            image_metadatas.append(
+                description.DescriptionJSONSerializer.from_desc(
+                    {
+                        "MAPLatitude": 58.5927694 + i * 0.0001,
+                        "MAPLongitude": 16.1840944 + i * 0.0001,
+                        "MAPCaptureTime": f"2021_02_13_13_{(24 + i) % 60:02d}_{(41 + i) % 60:02d}_140",
+                        "filename": str(test_exif),
+                        "filetype": "image",
+                        "MAPSequenceUUID": "multi_thread_sequence",
+                    }
+                )
+            )
+
+        # Test upload - this will internally use multithreading via _upload_images_parallel
+        results = list(sequence_uploader.upload_images(image_metadatas))
+
+        assert len(results) == 1, f"Expected 1 sequence result, got {len(results)}"
+        sequence_uuid, upload_result = results[0]
+        assert sequence_uuid == "multi_thread_sequence", (
+            f"Got wrong sequence UUID: {sequence_uuid}"
+        )
+        assert upload_result.error is None, (
+            f"Upload failed with error: {upload_result.error}"
+        )
+        assert upload_result.result is not None, "Upload should return a cluster ID"
+
+    def test_image_sequence_uploader_multithreading_with_cache_disabled(
+        self, setup_unittest_data: py.path.local
+    ):
+        """Test that ImageSequenceUploader's internal multithreading works correctly when cache is disabled."""
+        # Test with cache disabled via constants patch
+        with patch("mapillary_tools.constants.UPLOAD_CACHE_DIR", None):
+            upload_options = uploader.UploadOptions(
+                {"user_upload_token": "YOUR_USER_ACCESS_TOKEN"},
+                upload_cache_path=Path(setup_unittest_data.join("upload_cache")),
+                num_upload_workers=4,  # This will be used internally for parallel image uploads
+                dry_run=True,
+            )
+            emitter = uploader.EventEmitter()
+            sequence_uploader = uploader.ImageSequenceUploader(upload_options, emitter)
+
+            # Verify cache is disabled for both instances
+            assert sequence_uploader.cached_image_uploader.cache is None, (
+                "Should have cache disabled"
+            )
+
+            test_exif = setup_unittest_data.join("test_exif.jpg")
+
+            num_images = 100
+            image_metadatas = []
+
+            for i in range(num_images):
+                image_metadatas.append(
+                    description.DescriptionJSONSerializer.from_desc(
+                        {
+                            "MAPLatitude": 59.5927694 + i * 0.0001,
+                            "MAPLongi
… [15073 more characters]
Reference fix · 2 files, +86 −60the upstream merge, used only for grading calibration

The agent could not see this: the repository holds one commit and the sandbox has no network. Leak audit.

mapillary_tools/history.py, mapillary_tools/uploader.py

diff --git a/mapillary_tools/history.py b/mapillary_tools/history.py
index 8b7251ec..a0cf1311 100644
--- a/mapillary_tools/history.py
+++ b/mapillary_tools/history.py
@@ -162,6 +162,11 @@ def clear_expired(self) -> list[str]:
 
         return expired_keys
 
+    def keys(self):
+        with self._lock:
+            with dbm.open(self._file, flag="c") as db:
+                return db.keys()
+
     def _is_expired(self, payload: JSONDict) -> bool:
         expires_at = payload.get("expires_at")
         if isinstance(expires_at, (int, float)):
diff --git a/mapillary_tools/uploader.py b/mapillary_tools/uploader.py
index 8b2e2ad3..c81c91c4 100644
--- a/mapillary_tools/uploader.py
+++ b/mapillary_tools/uploader.py
@@ -2,6 +2,7 @@
 
 import concurrent.futures
 import dataclasses
+import hashlib
 import io
 import json
 import logging
@@ -56,6 +57,9 @@ class UploadOptions:
     user_items: config.UserItem
     chunk_size: int = int(constants.UPLOAD_CHUNK_SIZE_MB * 1024 * 1024)
     num_upload_workers: int = constants.MAX_IMAGE_UPLOAD_WORKERS
+    # When set, upload cache will be read/write there
+    # This option is exposed for testing purpose. In PROD, the path is calculated based on envvar and user_items
+    upload_cache_path: Path | None = None
     dry_run: bool = False
     nofinish: bool = False
     noresume: bool = False
@@ -471,7 +475,7 @@ def _zip_sequence_fp(
                 # Arcname should be unique, the name does not matter
                 arcname = f"{idx}.jpg"
                 zipinfo = zipfile.ZipInfo(arcname, date_time=(1980, 1, 1, 0, 0, 0))
-                zipf.writestr(zipinfo, SingleImageUploader.dump_image_bytes(metadata))
+                zipf.writestr(zipinfo, CachedImageUploader.dump_image_bytes(metadata))
             assert len(sequence) == len(set(zipf.namelist()))
             zipf.comment = json.dumps(
                 {"sequence_md5sum": sequence_md5sum},
@@ -537,6 +541,13 @@ class ImageSequenceUploader:
     def __init__(self, upload_options: UploadOptions, emitter: EventEmitter):
         self.upload_options = upload_options
         self.emitter = emitter
+        # Create a single shared SingleImageUploader instance that will be used across all uploads
+        cache = _maybe_create_persistent_cache_instance(self.upload_options)
+        if cache:
+            cache.clear_expired()
+        self.cached_image_uploader = CachedImageUploader(
+            self.upload_options, cache=cache
+        )
 
     def upload_images(
         self, image_metadatas: T.Sequence[types.ImageMetadata]
@@ -688,10 +699,6 @@ def _upload_images_from_queue(
         with api_v4.create_user_session(
             self.upload_options.user_items["user_upload_token"]
         ) as user_session:
-            single_image_uploader = SingleImageUploader(
-                self.upload_options, user_session=user_session
-            )
-
             while True:
                 # Assert that all images are already pushed into the queue
                 try:
@@ -710,8 +717,8 @@ def _upload_images_from_queue(
                 }
 
                 # image_progress will be updated during uploading
-                file_handle = single_image_uploader.upload(
-                    image_metadata, image_progress
+                file_handle = self.cached_image_uploader.upload(
+                    user_session, image_metadata, image_progress
                 )
 
                 # Update chunk_size (it was constant if set)
@@ -731,24 +738,27 @@ def _upload_images_from_queue(
         return indexed_file_handles
 
 
-class SingleImageUploader:
+class CachedImageUploader:
     def __init__(
         self,
         upload_options: UploadOptions,
-        user_session: requests.Session | None = None,
+        cache: history.PersistentCache | None = None,
     ):
         self.upload_options = upload_options
-        self.user_session = user_session
-        self.cache = self._maybe_create_persistent_cache_instance(
-            self.upload_options.user_items, upload_options
-        )
+        self.cache = cache
+        if self.cache:
+            self.cache.clear_expired()
 
+    # Thread-safe
     def upload(
-        self, image_metadata: types.ImageMetadata, image_progress: dict[str, T.Any]
+        self,
+        user_session: requests.Session,
+        image_metadata: types.ImageMetadata,
+        image_progress: dict[str, T.Any],
     ) -> str:
         image_bytes = self.dump_image_bytes(image_metadata)
 
-        uploader = Uploader(self.upload_options, user_session=self.user_session)
+        uploader = Uploader(self.upload_options, user_session=user_session)
 
         session_key = uploader._gen_session_key(io.BytesIO(image_bytes), image_progress)
 
@@ -786,51 +796,7 @@ def dump_image_bytes(cls, metadata: types.ImageMetadata) -> bytes:
                 f"Failed to dump EXIF bytes: {ex}", metadata.filename
             ) from ex
 
-    @classmethod
-    def _maybe_create_persistent_cache_instance(
-        cls, user_items: config.UserItem, upload_options: UploadOptions
-    ) -> history.PersistentCache | None:
-        if not constants.UPLOAD_CACHE_DIR:
-            LOG.debug(
-                "Upload cache directory is set empty, skipping caching upload file handles"
-            )
-            return None
-
-        if upload_options.dry_run:
-            LOG.debug("Dry-run mode enabled, skipping caching upload file handles")
-            return None
-
-        # Different python/CLI versions use different cache (dbm) formats.
-        # Separate them to avoid conflicts
-        py_version_parts = [str(part) for part in sys.version_info[:3]]
-        version = f"py_{'_'.join(py_version_parts)}_{VERSION}"
-
-        cache_path_dir = (
-            Path(constants.UPLOAD_CACHE_DIR)
-            .joinpath(version)
-            .joinpath(api_v4.MAPILLARY_CLIENT_TOKEN.replace("|", "_"))
-            .joinpath(
-                user_items.get("MAPSettingsUserKey", user_items["user_upload_token"])
-            )
-        )
-        cache_path_dir.mkdir(parents=True, exist_ok=True)
-        cache_path = cache_path_dir.joinpath("cached_file_handles")
-
-        # Sanitize sensitive segments for logging
-        sanitized_cache_path = (
-            Path(constants.UPLOAD_CACHE_DIR)
-            .joinpath(version)
-            .joinpath("***")
-            .joinpath("***")
-            .joinpath("cached_file_handles")
-        )
-        LOG.debug(f"File handle cache path: {sanitized_cache_path}")
-
-        cache = history.PersistentCache(str(cache_path.resolve()))
-        cache.clear_expired()
-
-        return cache
-
+    # Thread-safe
     def _get_cached_file_handle(self, key: str) -> str | None:
         if self.cache is None:
             return None
@@ -840,6 +806,7 @@ def _get_cached_file_handle(self, key: str) -> str | None:
 
         return self.cache.get(key)
 
+    # Thread-safe
     def _set_file_handle_cache(self, key: str, value: str) -> None:
         if self.cache is None:
             return
@@ -1168,3 +1135,57 @@ def _prefixed_uuid4():
 
 def _is_uuid(key: str) -> bool:
     return key.startswith("uuid_") or key.startswith("mly_tools_uuid_")
+
+
+def _build_upload_cache_path(upload_options: UploadOptions) -> Path:
+    # Different python/CLI versions use different cache (dbm) formats.
+    # Separate them to avoid conflicts
+    py_version_parts = [str(part) for part in sys.version_info[:3]]
+    version = f"py_{'_'.join(py_version_parts)}_{VERSION}"
+    # File handles are not sharable between different users
+    user_id = str(
+        upload_options.user_items.get(
+            "MAPSettingsUserKey", upload_options.user_items["user_upload_token"]
+        )
+    )
+    # Use hash to avoid log sensitive data
+    user_fingerprint = utils.md5sum_fp(
+        io.BytesIO((api_v4.MAPILLARY_CLIENT_TOKEN + user_id).encode("utf-8")),
+        md5=hashlib.sha256(),
+    ).hexdigest()[:24]
+
+    cache_path = (
+        Path(constants.UPLOA
… [1048 more characters]

What each model did

Submitted diff first, then every command. Outputs show 12 lines; expand for the rest.
GPT-5.6 Luna · attempt 3 · solved · 12 steps · $0.012
Submitted diff · 176 lines
diff --git a/mapillary_tools/history.py b/mapillary_tools/history.py
index 8b7251e..4011f0a 100644
--- a/mapillary_tools/history.py
+++ b/mapillary_tools/history.py
@@ -137,6 +137,15 @@ class PersistentCache:
             f"Cached file handle for {key} ({(time.perf_counter() - s) * 1000:.0f} ms)"
         )
 
+    def keys(self) -> list[str]:
+        """Return the keys currently stored in the cache."""
+        with self._lock:
+            with dbm.open(self._file, flag="c") as db:
+                return [
+                    key.decode("utf-8") if isinstance(key, bytes) else key
+                    for key in db.keys()
+                ]
+
     def clear_expired(self) -> list[str]:
         s = time.perf_counter()
 
diff --git a/mapillary_tools/uploader.py b/mapillary_tools/uploader.py
index 8b2e2ad..f1d79cf 100644
--- a/mapillary_tools/uploader.py
+++ b/mapillary_tools/uploader.py
@@ -59,6 +59,7 @@ class UploadOptions:
     dry_run: bool = False
     nofinish: bool = False
     noresume: bool = False
+    upload_cache_path: Path | None = None
 
     def __post_init__(self):
         if self.num_upload_workers <= 0:
@@ -534,13 +535,22 @@ class ZipUploader:
 
 
 class ImageSequenceUploader:
-    def __init__(self, upload_options: UploadOptions, emitter: EventEmitter):
+    def __init__(
+        self, upload_options: UploadOptions, emitter: EventEmitter | None = None
+    ):
         self.upload_options = upload_options
-        self.emitter = emitter
+        self.emitter = emitter if emitter is not None else EventEmitter()
+        # A single cache owner is important: dbm implementations are not safe when
+        # each worker opens the same database independently.  It also lets the
+        # uploader serialize the cache miss/upload/store operation.
+        self.cached_image_uploader = SingleImageUploader(upload_options)
 
     def upload_images(
         self, image_metadatas: T.Sequence[types.ImageMetadata]
     ) -> T.Generator[tuple[str, UploadResult], None, None]:
+        # Cache creation is deliberately done only in __init__.  Replacing the
+        # options with dry_run must not discard an existing cache.
+        self.cached_image_uploader.upload_options = self.upload_options
         sequences = types.group_and_sort_images(image_metadatas)
 
         for sequence_idx, (sequence_uuid, sequence) in enumerate(sequences.items()):
@@ -688,9 +698,7 @@ class ImageSequenceUploader:
         with api_v4.create_user_session(
             self.upload_options.user_items["user_upload_token"]
         ) as user_session:
-            single_image_uploader = SingleImageUploader(
-                self.upload_options, user_session=user_session
-            )
+            single_image_uploader = self.cached_image_uploader
 
             while True:
                 # Assert that all images are already pushed into the queue
@@ -711,7 +719,7 @@ class ImageSequenceUploader:
 
                 # image_progress will be updated during uploading
                 file_handle = single_image_uploader.upload(
-                    image_metadata, image_progress
+                    image_metadata, image_progress, user_session=user_session
                 )
 
                 # Update chunk_size (it was constant if set)
@@ -739,29 +747,39 @@ class SingleImageUploader:
     ):
         self.upload_options = upload_options
         self.user_session = user_session
+        self._cache_lock = threading.Lock()
         self.cache = self._maybe_create_persistent_cache_instance(
             self.upload_options.user_items, upload_options
         )
 
     def upload(
-        self, image_metadata: types.ImageMetadata, image_progress: dict[str, T.Any]
+        self,
+        image_metadata: types.ImageMetadata,
+        image_progress: dict[str, T.Any],
+        user_session: requests.Session | None = None,
     ) -> str:
         image_bytes = self.dump_image_bytes(image_metadata)
 
-        uploader = Uploader(self.upload_options, user_session=self.user_session)
+        uploader = Uploader(
+            self.upload_options,
+            user_session=self.user_session if user_session is None else user_session,
+        )
 
         session_key = uploader._gen_session_key(io.BytesIO(image_bytes), image_progress)
 
-        file_handle = self._get_cached_file_handle(session_key)
+        # Checking and populating the cache must be one operation.  Otherwise
+        # parallel workers can all observe a miss and upload the same image.
+        with self._cache_lock:
+            file_handle = self._get_cached_file_handle(session_key)
 
-        if file_handle is None:
-            # image_progress will be updated during uploading
-            file_handle = uploader.upload_stream(
-                io.BytesIO(image_bytes),
-                session_key=session_key,
-                progress=image_progress,
-            )
-            self._set_file_handle_cache(session_key, file_handle)
+            if file_handle is None:
+                # image_progress will be updated during uploading
+                file_handle = uploader.upload_stream(
+                    io.BytesIO(image_bytes),
+                    session_key=session_key,
+                    progress=image_progress,
+                )
+                self._set_file_handle_cache(session_key, file_handle)
 
         return file_handle
 
@@ -790,12 +808,6 @@ class SingleImageUploader:
     def _maybe_create_persistent_cache_instance(
         cls, user_items: config.UserItem, upload_options: UploadOptions
     ) -> history.PersistentCache | None:
-        if not constants.UPLOAD_CACHE_DIR:
-            LOG.debug(
-                "Upload cache directory is set empty, skipping caching upload file handles"
-            )
-            return None
-
         if upload_options.dry_run:
             LOG.debug("Dry-run mode enabled, skipping caching upload file handles")
             return None
@@ -805,16 +817,25 @@ class SingleImageUploader:
         py_
… [1549 more characters]
step 1

I’ll inspect the uploader, cache, and existing tests first, then reproduce the concurrency behavior. I’ll make the smallest source-only fix, run focused tests (and any broader relevant suite), and review the final diff.

$ cd /repo && ls && find . -maxdepth 3 -type f | head -80 && grep -R "class ImageSequenceUploader\|class UploadOptions\|PersistentCache\|cached_image_uploader" -n . --exclude-dir=.git
output · 46 lines
CHANGELOG.md
CODE_OF_CONDUCT.md
CONTRIBUTING.md
LICENSE
README.md
README_osx_package.txt
docs
mapillary_tools
mapillary_tools.egg-info
pyinstaller
pyproject.toml
schema
34 more lines
script
tests
./.git/COMMIT_EDITMSG
./.git/config
./.git/description
./.git/HEAD
./.git/hooks/applypatch-msg.sample
./.git/hooks/commit-msg.sample
./.git/hooks/fsmonitor-watchman.sample
./.git/hooks/post-update.sample
./.git/hooks/pre-applypatch.sample
./.git/hooks/pre-commit.sample
./.git/hooks/pre-merge-commit.sample
./.git/hooks/pre-push.sample
./.git/hooks/pre-rebase.sample
./.git/hooks/pre-receive.sample
./.git/hooks/prepare-commit-msg.sample
./.git/hooks/push-to-checkout.sample
./.git/hooks/sendemail-validate.sample
./.git/hooks/update.sample
./.git/index
./.git/info/exclude
./.git/logs/HEAD
./.github/workflows/codeql-analysis.yml
./.github/workflows/python-package.yml
./.github/workflows/release.yml
./.gitignore

…[3286 chars omitted]…
ts/unit/test_persistent_cache.py:190:    cache = PersistentCache(cache_file)
./tests/unit/test_persistent_cache.py:203:    cache = PersistentCache(cache_file)
./tests/unit/test_persistent_cache.py:228:    cache = PersistentCache(cache_file)
./tests/unit/test_persistent_cache.py:272:    cache = PersistentCache(cache_file)
./tests/unit/test_persistent_cache.py:286:    cache = PersistentCache(cache_file)
step 2
$ cd /repo && sed -n '1,240p' mapillary_tools/history.py && grep -R "ImageSequenceUploader\|upload_images\|upload_sequence" -n tests | head -100
output · 44 lines
from __future__ import annotations

import contextlib
import dbm
import json
import logging
import string
import threading
import time
import typing as T
from pathlib import Path
32 more lines
# dbm modules are dynamically imported, so here we explicitly import dbm.sqlite3 to make sure pyinstaller include it
# Otherwise you will see: ImportError: no dbm clone found; tried ['dbm.sqlite3', 'dbm.gnu', 'dbm.ndbm', 'dbm.dumb']
try:
    import dbm.sqlite3  # type: ignore
except ImportError:
    pass


from . import constants, types
from .serializer.description import DescriptionJSONSerializer

JSONDict = T.Dict[str, T.Union[str, int, float, None]]

LOG = logging.getLogger(__name__)


def _validate_hexdigits(md5sum: str):
    try:
        assert set(md5sum).issubset(string.hexdigits)
        assert 4 <= len(md5sum)
        _ = int(md5sum, 16)
    except Exception:
        raise ValueError(f"Invalid md5su
…[4326 chars omitted]…
, dict):
            LOG.warning(f"Invalid cache value format: {payload}")
            return {}

        return payload
tests/integration/test_process_and_upload.py:152:def test_process_and_upload_images_only(
tests/unit/test_uploader.py:27:def test_upload_images(setup_unittest_data: py.path.local, setup_upload: py.path.local):
tests/unit/test_uploader.py:68:def test_upload_images_multiple_sequences(
step 3
$ cd /repo && sed -n '520,710p' mapillary_tools/uploader.py && sed -n '1,130p' tests/unit/test_uploader.py && grep -n "def upload_images\|class EventEmitter\|def upload_sequence" mapillary_tools/uploader.py
output · 40 lines
done_path = wip_path.parent.joinpath(
                _suffix_session_key(upload_md5sum, api_v4.ClusterFileType.ZIP)
            )

            try:
                os.remove(done_path)
            except FileNotFoundError:
                pass
            wip_path.rename(done_path)
        finally:
            try:
                os.remove(wip_path)
28 more lines
            except FileNotFoundError:
                pass


class ImageSequenceUploader:
    def __init__(self, upload_options: UploadOptions, emitter: EventEmitter):
        self.upload_options = upload_options
        self.emitter = emitter

    def upload_images(
        self, image_metadatas: T.Sequence[types.ImageMetadata]
    ) -> T.Generator[tuple[str, UploadResult], None, None]:
        sequences = types.group_and_sort_images(image_metadatas)

        for sequence_idx, (sequence_uuid, sequence) in enumerate(sequences
…[8628 chars omitted]…
rt len(results) == 2
    actual_descs = sum(extract_all_uploaded_descs(Path(setup_upload)), [])
    assert 2 == len(actual_descs)


def test_upload_zip(
    setup_unittest_data: py.path.local, setup_upload: py.path.local, emitter=None
):
    test_exif = setup_unittest_data.join("test_exif.jpg")
    setup_unittest_data.join("another_directory").mkdir()
197:class EventEmitter:
541:    def upload_images(
step 4
$ cd /repo && sed -n '630,790p' mapillary_tools/uploader.py && sed -n '180,240p' mapillary_tools/uploader.py && grep -R "ImageSequenceUploader(" -n mapillary_tools tests | head -30 && sed -n '1,130p' mapillary_tools/config.py | grep -n "CACHE" -C2; grep -n "UPLOAD_CACHE" mapillary_tools/constants.py
output · 38 lines
# Lock is used to synchronize event emission
        lock = threading.Lock()

        # Push all images into the queue
        image_queue: queue.Queue[tuple[int, types.ImageMetadata]] = queue.Queue()
        for idx, image_metadata in enumerate(sequence):
            image_queue.put((idx, image_metadata))

        upload_interrupted = threading.Event()

        with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor:
            futures = [
26 more lines
                executor.submit(
                    self._upload_images_from_queue,
                    image_queue,
                    lock,
                    upload_interrupted,
                    sequence_progress,
                )
                for _ in range(max_workers)
            ]

            indexed_image_file_handles = []

            try:
                for future in futures:
                 
…[6177 chars omitted]…
except Exception as ex:
                yield video_metadata, UploadResult(error=ex)
                continue

            assert isinstance(video_metadata.md5sum, str), "md5sum should be updated"

            progress: SequenceProgress = {
mapillary_tools/upload.py:559:    image_uploader = uploader.ImageSequenceUploader(
150:UPLOAD_CACHE_DIR: str = os.getenv(
151:    _ENV_PREFIX + "UPLOAD_CACHE_DIR",
step 5
$ cd /repo && python - <<'PY'
output · 33 lines
mapillary_tools/history.py  |  9 +++++
 mapillary_tools/uploader.py | 86 ++++++++++++++++++++++++++++-----------------
 2 files changed, 63 insertions(+), 32 deletions(-)
diff --git a/mapillary_tools/uploader.py b/mapillary_tools/uploader.py
index 8b2e2ad..d995ead 100644
--- a/mapillary_tools/uploader.py
+++ b/mapillary_tools/uploader.py
@@ -59,6 +59,7 @@ class UploadOptions:
     dry_run: bool = False
     nofinish: bool = False
     noresume: bool = False
+    upload_cache_path: Path | None = None
21 more lines
 
     def __post_init__(self):
         if self.num_upload_workers <= 0:
@@ -367,6 +368,10 @@ class ZipUploader:
     def zip_images_and_upload(
         cls, uploader: Uploader, image_metadatas: T.Sequence[types.ImageMetadata]
     ) -> T.Generator[tuple[str, UploadResult], None, None]:
+        # The options may be replaced by callers between invocations, but cache
+        # creation 
…[5925 chars omitted]…
token"])
+                )
             )
-        )
-        cache_path_dir.mkdir(parents=True, exist_ok=True)
-        cache_path = cache_path_dir.joinpath("cached_file_handles")
+            cache_path_dir.mkdir(parents=True, exist_ok=True)
+            cache_path = cache_path_dir.joinpath("cached_file_handles")
 
         # Sanitize sensitive segments for logging
         sanitized_cache_path = (
step 6
$ cd /repo && sed -n '350,380p' mapillary_tools/uploader.py && grep -n "def _gen_session_key\|def _is_uuid" mapillary_tools/uploader.py && python -m compileall -q mapillary_tools && pytest -q tests/unit/test_persistent_cache.py tests/unit/test_uploader.py --disable-warnings --maxfail=1
output · 30 lines
) -> None:
        """
        Group images into sequences and zip each sequence into a zipfile.
        """
        sequences = types.group_and_sort_images(metadatas)
        os.makedirs(zip_dir, exist_ok=True)

        for sequence_uuid, sequence in sequences.items():
            _validate_metadatas(sequence)
            # For atomicity we write into a WIP file and then rename to the final file
            wip_zip_filename = zip_dir.joinpath(
                f".mly_zip_{uuid.uuid4()}_{sequence_uuid}_{os.getpid()}_{int(time.time())}"
18 more lines
            )
            with cls._wip_file_context(wip_zip_filename) as wip_path:
                with wip_path.open("wb") as wip_fp:
                    cls._zip_sequence_fp(sequence, wip_fp)

    @classmethod
    def zip_images_and_upload(
        cls, uploader: Uploader, image_metadatas: T.Sequence[types.ImageMetadata]
    ) -> T.Generator[tuple
…[3409 chars omitted]…
                                     ^^^^
E       NameError: name 'self' is not defined

mapillary_tools/uploader.py:374: NameError
=========================== short test summary info ============================
FAILED tests/unit/test_uploader.py::test_upload_images - NameError: name 'sel...
!!!!!!!!!!!!!!!!!!!!!!!!!! stopping after 1 failures !!!!!!!!!!!!!!!!!!!!!!!!!!!
1 failed, 60 passed in 40.51s
step 7
$ cd /repo && python - <<'PY'
output · 2 lines
.................................................................        [100%]
65 passed in 40.49s
step 8
$ cd /repo && pytest -q --disable-warnings --maxfail=1 && git diff --check && git diff -- mapillary_tools/history.py mapillary_tools/uploader.py | tail -100
output · 31 lines
sss....s.s.s.s...............ssssss....ss....sssssss.................... [ 25%]
...................................................ssssssssss........... [ 50%]
........................................................................ [ 75%]
.......................................................................  [100%]
255 passed, 32 skipped in 96.29s (0:01:36)
                 # Update chunk_size (it was constant if set)
@@ -739,29 +747,39 @@ class SingleImageUploader:
     ):
         self.upload_options = upload_options
         self.user_session = user_session
+        self._cache_lock = threading.Lock()
         self.cache = self._maybe_create_persistent_cache_instance(
19 more lines
             self.upload_options.user_items, upload_options
         )
 
     def upload(
-        self, image_metadata: types.ImageMetadata, image_progress: dict[str, T.Any]
+        self,
+        image_metadata:
…[3403 chars omitted]…
token"])
+                )
             )
-        )
-        cache_path_dir.mkdir(parents=True, exist_ok=True)
-        cache_path = cache_path_dir.joinpath("cached_file_handles")
+            cache_path_dir.mkdir(parents=True, exist_ok=True)
+            cache_path = cache_path_dir.joinpath("cached_file_handles")
 
         # Sanitize sensitive segments for logging
         sanitized_cache_path = (
step 9
$ cd /repo && git diff -- mapillary_tools/uploader.py | head -75 && git status --short
output · 38 lines
diff --git a/mapillary_tools/uploader.py b/mapillary_tools/uploader.py
index 8b2e2ad..f1d79cf 100644
--- a/mapillary_tools/uploader.py
+++ b/mapillary_tools/uploader.py
@@ -59,6 +59,7 @@ class UploadOptions:
     dry_run: bool = False
     nofinish: bool = False
     noresume: bool = False
+    upload_cache_path: Path | None = None
 
     def __post_init__(self):
         if self.num_upload_workers <= 0:
26 more lines
@@ -534,13 +535,22 @@ class ZipUploader:
 
 
 class ImageSequenceUploader:
-    def __init__(self, upload_options: UploadOptions, emitter: EventEmitter):
+    def __init__(
+        self, upload_options: UploadOptions, emitter: EventEmitter | None = None
+    ):
         self.upload_options = upload_options
-        self.emitter = emitter
+        self.emitter = emitter if emitter is not None else EventEmitter()
+        # A single cache owner is important: dbm implementations are not
…[1967 chars omitted]…
 )
 
     def upload(
-        self, image_metadata: types.ImageMetadata, image_progress: dict[str, T.Any]
+        self,
+        image_metadata: types.ImageMetadata,
+        image_progress: dict[str, T.Any],
+        user_session: requests.Session | None = None,
     ) -> str:
         image_bytes = self.dump_image_bytes(image_metadata)
 
 M mapillary_tools/history.py
 M mapillary_tools/uploader.py
step 10
$ cd /repo && sed -n '545,570p' mapillary_tools/upload.py && sed -n '835,860p' mapillary_tools/uploader.py && git diff -- mapillary_tools/history.py
output · 35 lines
return [d for d in metadatas if d.filename.resolve() in resolved_image_paths]


def _gen_upload_everything(
    mly_uploader: uploader.Uploader,
    metadatas: T.Sequence[types.Metadata],
    import_paths: T.Sequence[Path],
    skip_subfolders: bool,
):
    # Upload images
    image_metadatas = _find_metadata_with_filename_existed_in(
        (m for m in metadatas if isinstance(m, types.ImageMetadata)),
23 more lines
        utils.find_images(import_paths, skip_subfolders=skip_subfolders),
    )
    image_uploader = uploader.ImageSequenceUploader(
        mly_uploader.upload_options, emitter=mly_uploader.emitter
    )
    yield from image_uploader.upload_images(image_metadatas)

    # Upload videos
    video_metadatas = _find_metadata_with_filename_existed_in(
        (m for m in metadatas if isinstance(m, types.VideoMetadata)),
        utils.find_videos(import_paths, skip_subfolders=skip_subfo
…[1223 chars omitted]…
lf) -> list[str]:
+        """Return the keys currently stored in the cache."""
+        with self._lock:
+            with dbm.open(self._file, flag="c") as db:
+                return [
+                    key.decode("utf-8") if isinstance(key, bytes) else key
+                    for key in db.keys()
+                ]
+
     def clear_expired(self) -> list[str]:
         s = time.perf_counter()
step 11
$ cd /repo && sed -n '1095,1125p' mapillary_tools/uploader.py && sed -n '1185,1205p' mapillary_tools/uploader.py
output · 39 lines
# Start uploading
        return upload_service.upload_shifted_chunks(
            shifted_chunks, begin_offset, read_timeout=read_timeout
        )

    def _gen_session_key(self, fp: T.IO[bytes], progress: dict[str, T.Any]) -> str:
        if self.upload_options.noresume:
            # Generate a unique UUID for session_key when noresume is True
            # to prevent resuming from previous uploads
            session_key = f"{_prefixed_uuid4()}"
        else:
            fp.seek(0, io.SEEK_SET)
27 more lines
            session_key = utils.md5sum_fp(fp).hexdigest()

        filetype = progress.get("file_type")
        if filetype is not None:
            session_key = _suffix_session_key(session_key, types.FileType(filetype))

        return session_key


def _validate_metadatas(metadatas: T.Sequence[types.ImageMetadata]):
    for metadata in metadatas:
        validate_image_desc(Des
…[29 chars omitted]…
c(metadata))
        if not metadata.filename.is_file():
            raise FileNotFoundError(f"No such file {metadata.filename}")


def _is_immediate_retriable_exception(ex: BaseException) -> bool:
    if (
    prefixed = f"uuid_{uuid.uuid4().hex}"
    assert _is_uuid(prefixed)
    return prefixed


def _is_uuid(key: str) -> bool:
    return key.startswith("uuid_") or key.startswith("mly_tools_uuid_")
step 13
$ done: {"notes": "Implemented shared, concurrency-safe image upload caching, configurable cache paths, PersistentCache.keys(), optional emitter support, and validated with 255 passing tests (32 skipped)."}
output · 1 lines
submitted.
GPT-5.6 Luna · attempt 4 · solved · 12 steps · $0.011
Submitted diff · 195 lines
diff --git a/mapillary_tools/history.py b/mapillary_tools/history.py
index 8b7251e..3b02d80 100644
--- a/mapillary_tools/history.py
+++ b/mapillary_tools/history.py
@@ -95,6 +95,14 @@ class PersistentCache:
         else:
             self._lock = threading.Lock()
         self._file = file
+        # Public for users that need to coordinate cache operations.
+        self.file = file
+
+    def keys(self) -> list[str]:
+        """Return the currently stored cache keys."""
+        with self._lock:
+            with dbm.open(self._file, flag="c") as db:
+                return [key.decode() if isinstance(key, bytes) else key for key in db.keys()]
 
     def get(self, key: str) -> str | None:
         s = time.perf_counter()
diff --git a/mapillary_tools/uploader.py b/mapillary_tools/uploader.py
index 8b2e2ad..864f733 100644
--- a/mapillary_tools/uploader.py
+++ b/mapillary_tools/uploader.py
@@ -59,6 +59,7 @@ class UploadOptions:
     dry_run: bool = False
     nofinish: bool = False
     noresume: bool = False
+    upload_cache_path: Path | None = None
 
     def __post_init__(self):
         if self.num_upload_workers <= 0:
@@ -368,7 +369,6 @@ class ZipUploader:
         cls, uploader: Uploader, image_metadatas: T.Sequence[types.ImageMetadata]
     ) -> T.Generator[tuple[str, UploadResult], None, None]:
         sequences = types.group_and_sort_images(image_metadatas)
-
         for sequence_idx, (sequence_uuid, sequence) in enumerate(sequences.items()):
             try:
                 _validate_metadatas(sequence)
@@ -534,14 +534,24 @@ class ZipUploader:
 
 
 class ImageSequenceUploader:
-    def __init__(self, upload_options: UploadOptions, emitter: EventEmitter):
+    def __init__(self, upload_options: UploadOptions, emitter: EventEmitter | None = None):
         self.upload_options = upload_options
-        self.emitter = emitter
+        self.emitter = emitter if emitter is not None else EventEmitter()
+        # One cached uploader is shared by all image worker threads.  Apart from
+        # avoiding needless database handles, this makes the cache check and the
+        # upload it guards one operation from the caller's point of view.
+        self.cached_image_uploader = SingleImageUploader(
+            upload_options, emitter=self.emitter
+        )
 
     def upload_images(
         self, image_metadatas: T.Sequence[types.ImageMetadata]
     ) -> T.Generator[tuple[str, UploadResult], None, None]:
         sequences = types.group_and_sort_images(image_metadatas)
+        # Keep the operational options current while deliberately retaining the
+        # cache created at construction time (including when callers replace the
+        # options with a dry-run copy).
+        self.cached_image_uploader.upload_options = self.upload_options
 
         for sequence_idx, (sequence_uuid, sequence) in enumerate(sequences.items()):
             LOG.debug(f"Checksum for image sequence {sequence_uuid}...")
@@ -688,9 +698,7 @@ class ImageSequenceUploader:
         with api_v4.create_user_session(
             self.upload_options.user_items["user_upload_token"]
         ) as user_session:
-            single_image_uploader = SingleImageUploader(
-                self.upload_options, user_session=user_session
-            )
+            single_image_uploader = self.cached_image_uploader
 
             while True:
                 # Assert that all images are already pushed into the queue
@@ -732,13 +740,18 @@ class ImageSequenceUploader:
 
 
 class SingleImageUploader:
+    _cache_locks: dict[str, threading.Lock] = {}
+    _cache_locks_guard = threading.Lock()
+
     def __init__(
         self,
         upload_options: UploadOptions,
         user_session: requests.Session | None = None,
+        emitter: EventEmitter | None = None,
     ):
         self.upload_options = upload_options
         self.user_session = user_session
+        self.emitter = emitter if emitter is not None else EventEmitter()
         self.cache = self._maybe_create_persistent_cache_instance(
             self.upload_options.user_items, upload_options
         )
@@ -748,23 +761,39 @@ class SingleImageUploader:
     ) -> str:
         image_bytes = self.dump_image_bytes(image_metadata)
 
-        uploader = Uploader(self.upload_options, user_session=self.user_session)
+        uploader = Uploader(
+            self.upload_options, user_session=self.user_session, emitter=self.emitter
+        )
 
         session_key = uploader._gen_session_key(io.BytesIO(image_bytes), image_progress)
 
-        file_handle = self._get_cached_file_handle(session_key)
-
-        if file_handle is None:
-            # image_progress will be updated during uploading
-            file_handle = uploader.upload_stream(
-                io.BytesIO(image_bytes),
-                session_key=session_key,
-                progress=image_progress,
-            )
-            self._set_file_handle_cache(session_key, file_handle)
+        # A cache lookup followed by an upload must be serialized.  Otherwise
+        # two workers can both observe a miss and upload the same image.
+        lock_key = str(getattr(self.cache, "file", "")) if self.cache else ""
+        lock_key += "\0" + session_key
+        with self._cache_lock(lock_key):
+            file_handle = self._get_cached_file_handle(session_key)
+            if file_handle is None:
+                file_handle = uploader.upload_stream(
+                    io.BytesIO(image_bytes),
+                    session_key=session_key,
+                    progress=image_progress,
+                )
+                self._set_file_handle_cache(session_key, file_handle)
 
         return file_handle
 
+    @classmethod
+    @contextmanager
+    def _cache_lock(cls, key: str):
+        if not key:
+            yield
+            return
+        with cls._cache_locks_guard:
+            lock = cls._cache_locks.setdefault(key, threading.Lock())
+        with lock:
+            yield
+
     @classmeth
… [2261 more characters]
step 1

Plan: inspect the uploader, cache, and existing tests to understand the concurrency and API contracts; reproduce the likely race with focused scripts/tests; implement synchronization and option/cache semantics in non-test code; then run the affected test suite and review the final diff.

$ ls && grep -R "class ImageSequenceUploader\|class UploadOptions\|class PersistentCache\|cached_image_uploader" -n . --exclude-dir=.git | head -100
output · 17 lines
CHANGELOG.md
CODE_OF_CONDUCT.md
CONTRIBUTING.md
LICENSE
README.md
README_osx_package.txt
docs
mapillary_tools
mapillary_tools.egg-info
pyinstaller
pyproject.toml
schema
5 more lines
script
tests
./mapillary_tools/history.py:88:class PersistentCache:
./mapillary_tools/uploader.py:55:class UploadOptions:
./mapillary_tools/uploader.py:536:class ImageSequenceUploader:
step 2
$ grep -R "EventEmitter\|SequenceUpload\|ImageUpload" -n mapillary_tools/uploader.py tests | head -160
output · 9 lines
mapillary_tools/uploader.py:197:class EventEmitter:
mapillary_tools/uploader.py:474:                zipf.writestr(zipinfo, SingleImageUploader.dump_image_bytes(metadata))
mapillary_tools/uploader.py:536:class ImageSequenceUploader:
mapillary_tools/uploader.py:537:    def __init__(self, upload_options: UploadOptions, emitter: EventEmitter):
mapillary_tools/uploader.py:691:            single_image_uploader = SingleImageUploader(
mapillary_tools/uploader.py:734:class SingleImageUploader:
mapillary_tools/uploader.py:858:        emitter: EventEmitter | None = None,
mapillary_tools/uploader.py:864:            self.emitter = EventEmitter()
tests/unit/test_uploader.py:222:    emitter = uploader.EventEmitter()
step 3
$ sed -n '180,430p' mapillary_tools/uploader.py; sed -n '536,860p' mapillary_tools/uploader.py; sed -n '1,130p' tests/unit/test_uploader.py; sed -n '130,280p' tests/unit/test_uploader.py; grep -R "Cached\|cache" -n mapillary_tools/uploader.py mapillary_tools/config.py mapillary_tools/constants.py | head -100
output · 49 lines
#                 raise ex
#         else:
#             break
#     emit("upload_end")
#     finish_upload(data)
#     emit("upload_finished")
EventName = T.Literal[
    "upload_start",
    "upload_fetch_offset",
    "upload_progress",
    "upload_retrying",
    "upload_end",
37 more lines
    "upload_failed",
    "upload_finished",
]


class EventEmitter:
    events: dict[EventName, list]

    def __init__(self):
        self.events = {}

    def on(self, event: EventName):
        def _wrap(callback):
            self.events.setdefault(event, []).append(callback)
            return callback

        return _wrap

    def emit(self, event: EventName, *args, **kwargs):
        for callback in self.events.get(event, []):
            callback(*args, **kwargs)


@dataclasses.dataclass
class UploadResult:
    result: str | None = None
    error: Exception | None = None


class VideoUploader:
    @clas
…[8630 chars omitted]…
apillary_tools/uploader.py:841:        return self.cache.get(key)
mapillary_tools/uploader.py:843:    def _set_file_handle_cache(self, key: str, value: str) -> None:
mapillary_tools/uploader.py:844:        if self.cache is None:
mapillary_tools/uploader.py:850:        self.cache.set(key, value)
mapillary_tools/constants.py:152:    os.path.join(tempfile.gettempdir(), "mapillary_tools", "upload_cache"),
step 4
$ sed -n '536,740p' mapillary_tools/uploader.py; sed -n '734,855p' mapillary_tools/uploader.py; grep -R "upload_images(" -n mapillary_tools tests | head -50; sed -n '1,110p' tests/unit/test_uploader.py
output · 35 lines
class ImageSequenceUploader:
    def __init__(self, upload_options: UploadOptions, emitter: EventEmitter):
        self.upload_options = upload_options
        self.emitter = emitter

    def upload_images(
        self, image_metadatas: T.Sequence[types.ImageMetadata]
    ) -> T.Generator[tuple[str, UploadResult], None, None]:
        sequences = types.group_and_sort_images(image_metadatas)

        for sequence_idx, (sequence_uuid, sequence) in enumerate(sequences.items()):
            LOG.debug(f"Checksum for image sequence {sequence_uuid}...")
23 more lines
            sequence_md5sum = types.update_sequence_md5sum(sequence)

            sequence_progress: SequenceProgress = {
                "sequence_idx": sequence_idx,
                "total_sequence_count": len(sequences),
                "sequence_image_count": len(sequence),
                "sequence_uuid": sequence_uuid,
                
…[8629 chars omitted]…
None,
            "filetype": "image",
            "MAPSequenceUUID": "sequence_2",
        },
    ]
    mly_uploader = uploader.Uploader(
        uploader.UploadOptions(
            {
                "user_upload_token": "YOUR_USER_ACCESS_TOKEN",
                # will call the API for real
                # "MAPOrganizationKey": "3011753992432185",
            },
            dry_run=True,
        ),
step 5
$ sed -n '720,850p' mapillary_tools/uploader.py; sed -n '850,950p' mapillary_tools/uploader.py; grep -R "UploadOptions(" -n mapillary_tools | head -30; grep -n "def upload" mapillary_tools/uploader.py | tail -10
output · 43 lines
# Main thread will handle the interruption
                if upload_interrupted.is_set():
                    break

                with lock:
                    self.emitter.emit("upload_progress", image_progress)

                indexed_file_handles.append((idx, file_handle))

                image_queue.task_done()

        return indexed_file_handles
31 more lines

class SingleImageUploader:
    def __init__(
        self,
        upload_options: UploadOptions,
        user_session: requests.Session | None = None,
    ):
        self.upload_options = upload_options
        self.user_session = user_session
        self.cache = self._maybe_create_persistent_cache_instance(
            self.upload_options.user_items, upload_options
        )

    def upload(
        self, image_metadata: types.ImageMetadata, image_progress: dict[str, T.Any]
    ) -> str:
        image_bytes = 
…[6569 chars omitted]…
se(resp)
            # TODO: Validate cluster_id
            cluster_id = body.get("cluster_id")

        progress["cluster_id"] = cluster_id
        self.emitter.emit("upload_finished", progress)
mapillary_tools/upload.py:89:        upload_options = uploader.UploadOptions(
223:    def upload_videos(
320:    def upload_zipfiles(
541:    def upload_images(
746:    def upload(
868:    def upload_stream(
step 6
$ python - <<'PY'
step 7
$ git diff --stat && git diff -- mapillary_tools/uploader.py mapillary_tools/history.py | sed -n '1,260p'
output · 34 lines
mapillary_tools/history.py  |  8 ++++
 mapillary_tools/uploader.py | 91 ++++++++++++++++++++++++++++++---------------
 2 files changed, 68 insertions(+), 31 deletions(-)
diff --git a/mapillary_tools/history.py b/mapillary_tools/history.py
index 8b7251e..3b02d80 100644
--- a/mapillary_tools/history.py
+++ b/mapillary_tools/history.py
@@ -95,6 +95,14 @@ class PersistentCache:
         else:
             self._lock = threading.Lock()
         self._file = file
+        # Public for users that need to coordinate cache operations.
22 more lines
+        self.file = file
+
+    def keys(self) -> list[str]:
+        """Return the currently stored cache keys."""
+        with self._lock:
+            with dbm.open(self._file, flag="c") as db:
+                return [key.decode() if isinstance(key, bytes) else key for key in db.keys()]
 
     def get(self, key: str) -> str | None:
         s = time.perf_c
…[6162 chars omitted]…
token"])
+                )
             )
-        )
-        cache_path_dir.mkdir(parents=True, exist_ok=True)
-        cache_path = cache_path_dir.joinpath("cached_file_handles")
+            cache_path_dir.mkdir(parents=True, exist_ok=True)
+            cache_path = cache_path_dir.joinpath("cached_file_handles")
 
         # Sanitize sensitive segments for logging
         sanitized_cache_path = (
step 8
$ python - <<'PY'
output · 28 lines
F
=================================== FAILURES ===================================
______________________________ test_upload_images ______________________________

setup_unittest_data = local('/tmp/pytest-of-root/pytest-1/test_upload_images0/data')
setup_upload = local('/tmp/pytest-of-root/pytest-1/test_upload_images0/mapillary_public_uploads')

    def test_upload_images(setup_unittest_data: py.path.local, setup_upload: py.path.local):
        mly_uploader = uploader.Uploader(
            uploader.UploadOptions(
                {"user_upload_token": "YOUR_USER_ACCESS_TOKEN"}, dry_run=True
            )
16 more lines
        )
        test_exif = setup_unittest_data.join("test_exif.jpg")
        descs: T.List[description.DescriptionOrError] = [
            {
                "MAPLatitude": 58.5927694,
                "MAPLongitude": 16.1840944,
                "MAPCaptureTime": "2021_02_13_13_24_41
…[1768 chars omitted]…
                                                 ^^^^
E       NameError: name 'self' is not defined

mapillary_tools/uploader.py:375: NameError
=========================== short test summary info ============================
FAILED tests/unit/test_uploader.py::test_upload_images - NameError: name 'sel...
!!!!!!!!!!!!!!!!!!!!!!!!!! stopping after 1 failures !!!!!!!!!!!!!!!!!!!!!!!!!!!
1 failed in 0.73s
step 9
$ sed -n '360,385p' mapillary_tools/uploader.py; sed -n '545,565p' mapillary_tools/uploader.py
output · 28 lines
wip_zip_filename = zip_dir.joinpath(
                f".mly_zip_{uuid.uuid4()}_{sequence_uuid}_{os.getpid()}_{int(time.time())}"
            )
            with cls._wip_file_context(wip_zip_filename) as wip_path:
                with wip_path.open("wb") as wip_fp:
                    cls._zip_sequence_fp(sequence, wip_fp)

    @classmethod
    def zip_images_and_upload(
        cls, uploader: Uploader, image_metadatas: T.Sequence[types.ImageMetadata]
    ) -> T.Generator[tuple[str, UploadResult], None, None]:
        sequences = types.group_and_sort_images(image_metadatas)
16 more lines
        # Keep the operational options current while deliberately retaining the
        # cache created at construction time (including when callers replace the
        # options with a dry-run copy).
        self.cached_image_uploader.upload_options = self.upload_options

        for sequence_idx, (sequ
…[1028 chars omitted]…
replace the
        # options with a dry-run copy).
        self.cached_image_uploader.upload_options = self.upload_options

        for sequence_idx, (sequence_uuid, sequence) in enumerate(sequences.items()):
            LOG.debug(f"Checksum for image sequence {sequence_uuid}...")
            sequence_md5sum = types.update_sequence_md5sum(sequence)

            sequence_progress: SequenceProgress = {
step 10
$ python - <<'PY'
output · 2 lines
.................................................................        [100%]
65 passed in 41.22s
step 11
$ git diff --check; git diff --stat; grep -n "Keep the operational" -B2 -A4 mapillary_tools/uploader.py; python - <<'PY'
output · 12 lines
mapillary_tools/history.py  |  8 ++++
 mapillary_tools/uploader.py | 96 ++++++++++++++++++++++++++++++---------------
 2 files changed, 72 insertions(+), 32 deletions(-)
549-    ) -> T.Generator[tuple[str, UploadResult], None, None]:
550-        sequences = types.group_and_sort_images(image_metadatas)
551:        # Keep the operational options current while deliberately retaining the
552-        # cache created at construction time (including when callers replace the
553-        # options with a dry-run copy).
554-        self.cached_image_uploader.upload_options = self.upload_options
555-
<class 'mapillary_tools.history.PersistentCache'> []
True
step 12
$ done: {"notes": "Implemented shared cached image uploader with per-cache-key synchronization, cache reuse across parallel workers, optional cache path and stable cache lifecycle semantics, PersistentCache.keys(), and event propagation. Affected uploader/cache tests pass (65)."}
output · 1 lines
submitted.
GPT-5.6 Luna · attempt 5 · solved · 17 steps · $0.016
Submitted diff · 142 lines
diff --git a/mapillary_tools/history.py b/mapillary_tools/history.py
index 8b7251e..2c3d36b 100644
--- a/mapillary_tools/history.py
+++ b/mapillary_tools/history.py
@@ -137,6 +137,11 @@ class PersistentCache:
             f"Cached file handle for {key} ({(time.perf_counter() - s) * 1000:.0f} ms)"
         )
 
+    def keys(self) -> list[str]:
+        with self._lock:
+            with dbm.open(self._file, flag="c") as db:
+                return [key.decode("utf-8") if isinstance(key, bytes) else key for key in db.keys()]
+
     def clear_expired(self) -> list[str]:
         s = time.perf_counter()
 
diff --git a/mapillary_tools/uploader.py b/mapillary_tools/uploader.py
index 8b2e2ad..03dcbda 100644
--- a/mapillary_tools/uploader.py
+++ b/mapillary_tools/uploader.py
@@ -59,6 +59,7 @@ class UploadOptions:
     dry_run: bool = False
     nofinish: bool = False
     noresume: bool = False
+    upload_cache_path: Path | None = None
 
     def __post_init__(self):
         if self.num_upload_workers <= 0:
@@ -537,6 +538,7 @@ class ImageSequenceUploader:
     def __init__(self, upload_options: UploadOptions, emitter: EventEmitter):
         self.upload_options = upload_options
         self.emitter = emitter
+        self.cached_image_uploader = SingleImageUploader(upload_options)
 
     def upload_images(
         self, image_metadatas: T.Sequence[types.ImageMetadata]
@@ -691,6 +693,11 @@ class ImageSequenceUploader:
             single_image_uploader = SingleImageUploader(
                 self.upload_options, user_session=user_session
             )
+            # All workers use the same cache and per-key locks, while retaining
+            # their independent HTTP sessions.
+            single_image_uploader.cache = self.cached_image_uploader.cache
+            single_image_uploader._cache_locks = self.cached_image_uploader._cache_locks
+            single_image_uploader._cache_locks_guard = self.cached_image_uploader._cache_locks_guard
 
             while True:
                 # Assert that all images are already pushed into the queue
@@ -701,6 +708,7 @@ class ImageSequenceUploader:
 
                 # Main thread will handle the interruption
                 if upload_interrupted.is_set():
+                    image_queue.task_done()
                     break
 
                 # Create a new mutatble progress to keep the sequence_progress immutable
@@ -742,6 +750,8 @@ class SingleImageUploader:
         self.cache = self._maybe_create_persistent_cache_instance(
             self.upload_options.user_items, upload_options
         )
+        self._cache_locks: dict[str, threading.Lock] = {}
+        self._cache_locks_guard = threading.Lock()
 
     def upload(
         self, image_metadata: types.ImageMetadata, image_progress: dict[str, T.Any]
@@ -752,18 +762,21 @@ class SingleImageUploader:
 
         session_key = uploader._gen_session_key(io.BytesIO(image_bytes), image_progress)
 
-        file_handle = self._get_cached_file_handle(session_key)
-
-        if file_handle is None:
-            # image_progress will be updated during uploading
-            file_handle = uploader.upload_stream(
-                io.BytesIO(image_bytes),
-                session_key=session_key,
-                progress=image_progress,
+        if self.cache is None:
+            return uploader.upload_stream(
+                io.BytesIO(image_bytes), session_key=session_key, progress=image_progress
             )
-            self._set_file_handle_cache(session_key, file_handle)
 
-        return file_handle
+        with self._cache_locks_guard:
+            key_lock = self._cache_locks.setdefault(session_key, threading.Lock())
+        with key_lock:
+            file_handle = self._get_cached_file_handle(session_key)
+            if file_handle is None:
+                file_handle = uploader.upload_stream(
+                    io.BytesIO(image_bytes), session_key=session_key, progress=image_progress
+                )
+                self._set_file_handle_cache(session_key, file_handle)
+            return file_handle
 
     @classmethod
     def dump_image_bytes(cls, metadata: types.ImageMetadata) -> bytes:
@@ -790,31 +803,27 @@ class SingleImageUploader:
     def _maybe_create_persistent_cache_instance(
         cls, user_items: config.UserItem, upload_options: UploadOptions
     ) -> history.PersistentCache | None:
-        if not constants.UPLOAD_CACHE_DIR:
-            LOG.debug(
-                "Upload cache directory is set empty, skipping caching upload file handles"
-            )
-            return None
-
         if upload_options.dry_run:
             LOG.debug("Dry-run mode enabled, skipping caching upload file handles")
             return None
 
-        # Different python/CLI versions use different cache (dbm) formats.
-        # Separate them to avoid conflicts
         py_version_parts = [str(part) for part in sys.version_info[:3]]
         version = f"py_{'_'.join(py_version_parts)}_{VERSION}"
 
-        cache_path_dir = (
-            Path(constants.UPLOAD_CACHE_DIR)
-            .joinpath(version)
-            .joinpath(api_v4.MAPILLARY_CLIENT_TOKEN.replace("|", "_"))
-            .joinpath(
-                user_items.get("MAPSettingsUserKey", user_items["user_upload_token"])
+        if upload_options.upload_cache_path is not None:
+            cache_path = upload_options.upload_cache_path
+            cache_path.parent.mkdir(parents=True, exist_ok=True)
+        else:
+            if not constants.UPLOAD_CACHE_DIR:
+                LOG.debug("Upload cache directory is set empty, skipping caching upload file handles")
+                return None
+            cache_path_dir = (
+                Path(constants.UPLOAD_CACHE_DIR).joinpath(version)
+                .joinpath(api_v4.MAPILLARY_CLIENT_TOKEN.replace("|", "_"))
+                .joinpath(user_items.get("MAPSettingsUserKey", user_items["user_upload_token"]))
             )
-        )
-        cache_path_dir.mkdir(
… [320 more characters]
step 1

I’ll inspect the uploader, cache, and existing tests first, then reproduce the concurrency/cache behavior before making a focused synchronization and API fix.

$ cd /repo && ls && grep -R "class ImageSequenceUploader\|class UploadOptions\|class PersistentCache\|cached_image_uploader" -n . --exclude-dir=.git | head -100
output · 17 lines
CHANGELOG.md
CODE_OF_CONDUCT.md
CONTRIBUTING.md
LICENSE
README.md
README_osx_package.txt
docs
mapillary_tools
mapillary_tools.egg-info
pyinstaller
pyproject.toml
schema
5 more lines
script
tests
./mapillary_tools/history.py:88:class PersistentCache:
./mapillary_tools/uploader.py:55:class UploadOptions:
./mapillary_tools/uploader.py:536:class ImageSequenceUploader:
step 2
$ cd /repo && grep -n "def .*upload\|EventEmitter\|ThreadPool\|parallel" mapillary_tools/uploader.py
output · 23 lines
160:# def pseudo_upload(metadata):
197:class EventEmitter:
223:    def upload_videos(
320:    def upload_zipfiles(
367:    def zip_images_and_upload(
411:    def _upload_zipfile(
537:    def __init__(self, upload_options: UploadOptions, emitter: EventEmitter):
541:    def upload_images(
569:    def _upload_sequence_and_finish(
581:            image_file_handles = self._upload_images_parallel(
601:    def _upload_manifest(self, image_file_handles: T.Sequence[str]) -> str:
621:    def _upload_images_parallel(
11 more lines
641:        with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor:
679:    def _upload_images_from_queue(
746:    def upload(
858:        emitter: EventEmitter | None = None,
864:            self.emitter = EventEmitter()
868:    def upload_stream(
920:    def finish_upload(
954:    def _create_upload_service(
975:    def _handle_upload_exception(
1005:    def _upload_name(cls, progress: UploaderProgress):
1036:    def _upload_stream_retryable(
step 3
$ cd /repo && sed -n '190,430p' mapillary_tools/uploader.py; sed -n '530,760p' mapillary_tools/uploader.py; sed -n '760,840p' mapillary_tools/uploader.py; grep -R "_maybe_create_persistent\|MAPILLARY_TOOLS_UPLOAD_CACHE" -n mapillary_tools tests
output · 47 lines
"upload_retrying",
    "upload_end",
    "upload_failed",
    "upload_finished",
]


class EventEmitter:
    events: dict[EventName, list]

    def __init__(self):
        self.events = {}
35 more lines
    def on(self, event: EventName):
        def _wrap(callback):
            self.events.setdefault(event, []).append(callback)
            return callback

        return _wrap

    def emit(self, event: EventName, *args, **kwargs):
        for callback in self.events.get(event, []):
            callback(*args, **kwargs)


@dataclasses.dataclass
class UploadResult:
    result: str | None = None
    error: Exception | None = None


class VideoUploader:
    @classmethod
    def upload_videos(
        cls, mly_uploader: Uploader, video_metadatas: T.Sequence[types.VideoMetadata]
    ) -> T.Generator[tuple[types.VideoMetadata, UploadResult], None, None]:
        # If upload in a random order, then 
…[8630 chars omitted]…
      if self.cache is None:
            return None

        if _is_uuid(key):
            return None

mapillary_tools/uploader.py:742:        self.cache = self._maybe_create_persistent_cache_instance(
mapillary_tools/uploader.py:790:    def _maybe_create_persistent_cache_instance(
tests/integration/fixtures.py:69:    os.environ["MAPILLARY_TOOLS_UPLOAD_CACHE_DIR"] = str(tmpdir.mkdir("upload_cache"))
step 4
$ cd /repo && sed -n '536,625p' mapillary_tools/uploader.py; sed -n '780,830p' mapillary_tools/uploader.py; grep -R "upload_images(" -n mapillary_tools tests | head -40; grep -n "UPLOAD_CACHE" mapillary_tools/constants.py mapillary_tools/config.py
output · 26 lines
class ImageSequenceUploader:
    def __init__(self, upload_options: UploadOptions, emitter: EventEmitter):
        self.upload_options = upload_options
        self.emitter = emitter

    def upload_images(
        self, image_metadatas: T.Sequence[types.ImageMetadata]
    ) -> T.Generator[tuple[str, UploadResult], None, None]:
        sequences = types.group_and_sort_images(image_metadatas)

        for sequence_idx, (sequence_uuid, sequence) in enumerate(sequences.items()):
            LOG.debug(f"Checksum for image sequence {sequence_uuid}...")
14 more lines
            sequence_md5sum = types.update_sequence_md5sum(sequence)

            sequence_progress: SequenceProgress = {
                "sequence_idx": sequence_idx,
                "total_sequence_count": len(sequences),
                "sequence_image_count": len(sequence),
                "sequence_uuid": sequence_uuid,
                
…[4329 chars omitted]…
mapillary_tools/upload.py:562:    yield from image_uploader.upload_images(image_metadatas)
mapillary_tools/uploader.py:541:    def upload_images(
tests/unit/test_uploader.py:27:def test_upload_images(setup_unittest_data: py.path.local, setup_upload: py.path.local):
mapillary_tools/constants.py:150:UPLOAD_CACHE_DIR: str = os.getenv(
mapillary_tools/constants.py:151:    _ENV_PREFIX + "UPLOAD_CACHE_DIR",
step 5
$ cd /repo && sed -n '840,950p' mapillary_tools/uploader.py; sed -n '145,160p' mapillary_tools/constants.py; sed -n '1,80p' tests/unit/test_uploader.py; sed -n '1,100p' mapillary_tools/upload.py
output · 49 lines
return self.cache.get(key)

    def _set_file_handle_cache(self, key: str, value: str) -> None:
        if self.cache is None:
            return

        if _is_uuid(key):
            return

        self.cache.set(key, value)

37 more lines
class Uploader:
    def __init__(
        self,
        upload_options: UploadOptions,
        user_session: requests.Session | None = None,
        emitter: EventEmitter | None = None,
    ):
        self.upload_options = upload_options
        self.user_session = user_session
        if emitter is None:
            # An empty event emitter that does nothing
            self.emitter = EventEmitter()
        else:
            self.emitter = emitter

    def upload_stream(
        self,
        fp: T.IO[bytes],
        session_key: str | None = None,
        progress: dict[str, T.Any] | None = None,
    ) -> str:
        if progress is None:
         
…[7954 chars omitted]…

    try:
        upload_options = uploader.UploadOptions(
            user_items,
            dry_run=dry_run,
            nofinish=nofinish,
            noresume=noresume,
            num_upload_workers=num_upload_workers,
        )
    except ValueError as ex:
        raise exceptions.MapillaryBadParameterError(str(ex)) from ex

    mly_uploader = uploader.Uploader(upload_options, emitter=emitter)
step 6
$ cd /repo && sed -n '100,260p' tests/unit/test_uploader.py; grep -R "ImageSequenceUploader" -n tests mapillary_tools | cat; sed -n '530,580p' mapillary_tools/upload.py
output · 44 lines
},
    ]
    mly_uploader = uploader.Uploader(
        uploader.UploadOptions(
            {
                "user_upload_token": "YOUR_USER_ACCESS_TOKEN",
                # will call the API for real
                # "MAPOrganizationKey": "3011753992432185",
            },
            dry_run=True,
        ),
    )
32 more lines
    results = list(
        uploader.ZipUploader.zip_images_and_upload(
            mly_uploader,
            [
                description.DescriptionJSONSerializer.from_desc(T.cast(T.Any, desc))
                for desc in descs
            ],
        )
    )
    assert len(results) == 2
    actual_descs = sum(extract_all_uploaded_descs(Path(setup_upload)), [])
    assert 2 == len(actual_descs)


def test_upload_zip(
    setup_unittest_data: py.path.local, setup_upload: py.path.local, emitter=None
):
    test_exif = setup_unittest_data.join("test_exif.jpg")
    
…[5754 chars omitted]…
s(mly_uploader, video_metadatas)

    # Upload zip files
    zip_paths = utils.find_zipfiles(import_paths, skip_subfolders=skip_subfolders)
    yield from uploader.ZipUploader.upload_zipfiles(mly_uploader, zip_paths)


def _normalize_import_paths(import_path: Path | T.Sequence[Path]) -> list[Path]:
    import_paths: list[Path]

    if isinstance(import_path, Path):
        import_paths = [import_path]
step 7

The main race is clear: every worker creates its own `SingleImageUploader` and cache object, so concurrent cache misses can upload the same content and dbm access is not coordinated. I’ll create one cached uploader per image-sequence uploader, serialize cache check/upload/store as one operation, add explicit cache-path handling, and make cache inspection thread-safe. I’ll also fix queue bookkeeping on worker failures.

$ cd /repo && python - <<'PY'
output · 3 lines
cache block not found
 mapillary_tools/history.py | 9 +++++++++
 1 file changed, 9 insertions(+)
step 8
$ cd /repo && sed -n '735,780p' mapillary_tools/uploader.py; git checkout -- mapillary_tools/history.py
output · 37 lines
def __init__(
        self,
        upload_options: UploadOptions,
        user_session: requests.Session | None = None,
    ):
        self.upload_options = upload_options
        self.user_session = user_session
        self.cache = self._maybe_create_persistent_cache_instance(
            self.upload_options.user_items, upload_options
        )

    def upload(
25 more lines
        self, image_metadata: types.ImageMetadata, image_progress: dict[str, T.Any]
    ) -> str:
        image_bytes = self.dump_image_bytes(image_metadata)

        uploader = Uploader(self.upload_options, user_session=self.user_session)

        session_key = uploader._gen_session_key(io.BytesIO(image_bytes), image_progress)

        file_handle = self._get_cached_file_handle(session_key)

        if file_handle is None:
            # image_progress will be updated during uploading
            file_handle = uploader.u
…[340 chars omitted]…
  try:
            edit = exif_write.ExifEdit(metadata.filename)
        except struct.error as ex:
            raise ExifError(f"Failed to load EXIF: {ex}", metadata.filename) from ex

        # The cast is to fix the type checker error
        edit.add_image_description(
            T.cast(
                T.Dict, desc_file_to_exif(DescriptionJSONSerializer.as_desc(metadata))
            )
        )
step 9
$ cd /repo && python - <<'PY'
output · 37 lines
@classmethod
    def _maybe_create_persistent_cache_instance(
        cls, user_items: config.UserItem, upload_options: UploadOptions
    ) -> history.PersistentCache | None:
        if upload_options.upload_cache_path is not None:
            cache_path = upload_options.upload_cache_path
            cache_path.parent.mkdir(parents=True, exist_ok=True)
        else:
            if not constants.UPLOAD_CACHE_DIR:
                LOG.debug("Upload cache directory is set empty, skipping caching upload file handles")
                return None
            cache_path_dir = (
25 more lines
                Path(constants.UPLOAD_CACHE_DIR).joinpath(version)
                .joinpath(api_v4.MAPILLARY_CLIENT_TOKEN.replace("|", "_"))
                .joinpath(user_items.get("MAPSettingsUserKey", user_items["user_upload_token"]))
            )
            cache_path_dir.mkdir(parents=True, exist_ok=True)
…[502 chars omitted]…
_expired()

        return cache

    def _get_cached_file_handle(self, key: str) -> str | None:
        if self.cache is None:
            return None

        if _is_uuid(key):
            return None

        return self.cache.get(key)

    def _set_file_handle_cache(self, key: str, value: str) -> None:
        if self.cache is None:
            return

        if _is_uuid(key):
            return
step 10
$ cd /repo && sed -n '785,815p' mapillary_tools/uploader.py; git diff --stat; pytest -q tests/unit/test_persistent_cache.py tests/unit/test_uploader.py --disable-warnings --maxfail=1 p 2>/dev/null || true
output · 34 lines
raise ExifError(f"Failed to load EXIF: {ex}", metadata.filename) from ex

        # The cast is to fix the type checker error
        edit.add_image_description(
            T.cast(
                T.Dict, desc_file_to_exif(DescriptionJSONSerializer.as_desc(metadata))
            )
        )

        try:
            return edit.dump_image_bytes()
        except struct.error as ex:
22 more lines
            raise ExifError(
                f"Failed to dump EXIF bytes: {ex}", metadata.filename
            ) from ex

    @classmethod
    def _maybe_create_persistent_cache_instance(
        cls, user_items: config.UserItem, upload_options: UploadOptions
    ) -> history.PersistentCache | None:
        if upload_options.upload_cache_path is not None:
            cache_path = upload_options.upload_cache_path
            cache_path.parent.mkdir(parents=True, exist_ok=True)
        else:
     
…[232 chars omitted]…
ts.UPLOAD_CACHE_DIR).joinpath(version)
                .joinpath(api_v4.MAPILLARY_CLIENT_TOKEN.replace("|", "_"))
                .joinpath(user_items.get("MAPSettingsUserKey", user_items["user_upload_token"]))
 mapillary_tools/history.py  |  5 ++++
 mapillary_tools/uploader.py | 69 +++++++++++++++++++++++----------------------
 2 files changed, 40 insertions(+), 34 deletions(-)

no tests ran in 0.02s
step 11
$ cd /repo && sed -n '760,810p' mapillary_tools/uploader.py; grep -n 'version =' mapillary_tools/uploader.py
output · 31 lines
uploader = Uploader(self.upload_options, user_session=self.user_session)

        session_key = uploader._gen_session_key(io.BytesIO(image_bytes), image_progress)

        if self.cache is None:
            return uploader.upload_stream(
                io.BytesIO(image_bytes), session_key=session_key, progress=image_progress
            )

        with self._cache_locks_guard:
            key_lock = self._cache_locks.setdefault(session_key, threading.Lock())
        with key_lock:
19 more lines
            file_handle = self._get_cached_file_handle(session_key)
            if file_handle is None:
                file_handle = uploader.upload_stream(
                    io.BytesIO(image_bytes), session_key=session_key, progress=image_progress
                )
                self._set_file_handle_cache(session_key, file_handle)
            return file_handle

    @classmethod
    def dump_
…[801 chars omitted]…
ions: UploadOptions
    ) -> history.PersistentCache | None:
        if upload_options.upload_cache_path is not None:
            cache_path = upload_options.upload_cache_path
            cache_path.parent.mkdir(parents=True, exist_ok=True)
        else:
            if not constants.UPLOAD_CACHE_DIR:
                LOG.debug("Upload cache directory is set empty, skipping caching upload file handles")
step 12
$ cd /repo && python - <<'PY'
output · 34 lines
............................................................             [100%]
60 passed in 40.03s
diff --git a/mapillary_tools/uploader.py b/mapillary_tools/uploader.py
index 8b2e2ad..6ee791e 100644
--- a/mapillary_tools/uploader.py
+++ b/mapillary_tools/uploader.py
@@ -59,6 +59,7 @@ class UploadOptions:
     dry_run: bool = False
     nofinish: bool = False
     noresume: bool = False
+    upload_cache_path: Path | None = None
 
22 more lines
     def __post_init__(self):
         if self.num_upload_workers <= 0:
@@ -537,6 +538,7 @@ class ImageSequenceUploader:
     def __init__(self, upload_options: UploadOptions, emitter: EventEmitter):
         self.upload_options = upload_options
         self.emitter = emitter
+        self.cached_image_uploader = SingleImageUploader(upload_options)
 
     def upload_images(
         self, image_metadatas: T.Sequence[types.ImageMetadata]
@@ -691,6 +693,11 @
…[4186 chars omitted]…
tems["user_upload_token"]))
             )
-        )
-        cache_path_dir.mkdir(parents=True, exist_ok=True)
-        cache_path = cache_path_dir.joinpath("cached_file_handles")
+            cache_path_dir.mkdir(parents=True, exist_ok=True)
+            cache_path = cache_path_dir.joinpath("cached_file_handles")
 
         # Sanitize sensitive segments for logging
         sanitized_cache_path = (
step 13
$ cd /repo && python - <<'PY'
output · 6 lines
Traceback (most recent call last):
  File "<stdin>", line 12, in <module>
AssertionError
 mapillary_tools/history.py  |  5 ++++
 mapillary_tools/uploader.py | 62 +++++++++++++++++++++++++--------------------
 2 files changed, 40 insertions(+), 27 deletions(-)
step 14
$ cd /repo && sed -n '700,735p' mapillary_tools/uploader.py; git status --short
output · 36 lines
single_image_uploader._cache_locks_guard = self.cached_image_uploader._cache_locks_guard

            while True:
                # Assert that all images are already pushed into the queue
                try:
                    idx, image_metadata = image_queue.get_nowait()
                except queue.Empty:
                    break

                # Main thread will handle the interruption
                if upload_interrupted.is_set():
                    break
24 more lines
                # Create a new mutatble progress to keep the sequence_progress immutable
                image_progress = {
                    **sequence_progress,
                    "import_path": str(image_metadata.filename),
                }

                # image_progress will be updated during uploading
                file_handle = single_image_uploader.upload(
                    image_metadata, 
…[123 chars omitted]…
ss["chunk_size"] = image_metadata.filesize

                # Main thread will handle the interruption
                if upload_interrupted.is_set():
                    break

                with lock:
                    self.emitter.emit("upload_progress", image_progress)

                indexed_file_handles.append((idx, file_handle))

 M mapillary_tools/history.py
 M mapillary_tools/uploader.py
step 15
$ cd /repo && python - <<'PY'
output · 3 lines
mapillary_tools/history.py  |  5 ++++
 mapillary_tools/uploader.py | 63 ++++++++++++++++++++++++++-------------------
 2 files changed, 41 insertions(+), 27 deletions(-)
step 16
$ cd /repo && pytest -q tests/unit/test_uploader.py --disable-warnings --maxfail=1 >/tmp/test.out && tail -5 /tmp/test.out || (cat /tmp/test.out; exit 1)
output · 2 lines
.....                                                                    [100%]
5 passed in 0.73s
step 18
$ done: {"notes": "Implemented synchronized shared image caching, explicit cache-path support, cache inspection, and validated with persistent-cache and uploader unit tests."}
output · 1 lines
submitted.
GPT-5.6 Luna · attempt 6 · solved · 14 steps · $0.013
Submitted diff · 120 lines
diff --git a/mapillary_tools/history.py b/mapillary_tools/history.py
index 8b7251e..4011f0a 100644
--- a/mapillary_tools/history.py
+++ b/mapillary_tools/history.py
@@ -137,6 +137,15 @@ class PersistentCache:
             f"Cached file handle for {key} ({(time.perf_counter() - s) * 1000:.0f} ms)"
         )
 
+    def keys(self) -> list[str]:
+        """Return the keys currently stored in the cache."""
+        with self._lock:
+            with dbm.open(self._file, flag="c") as db:
+                return [
+                    key.decode("utf-8") if isinstance(key, bytes) else key
+                    for key in db.keys()
+                ]
+
     def clear_expired(self) -> list[str]:
         s = time.perf_counter()
 
diff --git a/mapillary_tools/uploader.py b/mapillary_tools/uploader.py
index 8b2e2ad..9c3bd57 100644
--- a/mapillary_tools/uploader.py
+++ b/mapillary_tools/uploader.py
@@ -59,6 +59,7 @@ class UploadOptions:
     dry_run: bool = False
     nofinish: bool = False
     noresume: bool = False
+    upload_cache_path: Path | None = None
 
     def __post_init__(self):
         if self.num_upload_workers <= 0:
@@ -534,9 +535,10 @@ class ZipUploader:
 
 
 class ImageSequenceUploader:
-    def __init__(self, upload_options: UploadOptions, emitter: EventEmitter):
+    def __init__(self, upload_options: UploadOptions, emitter: EventEmitter | None = None):
         self.upload_options = upload_options
-        self.emitter = emitter
+        self.emitter = emitter if emitter is not None else EventEmitter()
+        self.cached_image_uploader = SingleImageUploader(upload_options, emitter=self.emitter)
 
     def upload_images(
         self, image_metadatas: T.Sequence[types.ImageMetadata]
@@ -689,7 +691,10 @@ class ImageSequenceUploader:
             self.upload_options.user_items["user_upload_token"]
         ) as user_session:
             single_image_uploader = SingleImageUploader(
-                self.upload_options, user_session=user_session
+                self.upload_options, user_session=user_session,
+                cache=self.cached_image_uploader.cache,
+                cache_lock=self.cached_image_uploader._cache_lock,
+                emitter=self.emitter,
             )
 
             while True:
@@ -736,32 +741,36 @@ class SingleImageUploader:
         self,
         upload_options: UploadOptions,
         user_session: requests.Session | None = None,
+        cache: history.PersistentCache | None = None,
+        cache_lock: threading.RLock | None = None,
+        emitter: EventEmitter | None = None,
     ):
         self.upload_options = upload_options
         self.user_session = user_session
-        self.cache = self._maybe_create_persistent_cache_instance(
+        self.emitter = emitter if emitter is not None else EventEmitter()
+        self.cache = cache if cache is not None else self._maybe_create_persistent_cache_instance(
             self.upload_options.user_items, upload_options
         )
+        self._cache_lock = cache_lock if cache_lock is not None else threading.RLock()
 
     def upload(
         self, image_metadata: types.ImageMetadata, image_progress: dict[str, T.Any]
     ) -> str:
         image_bytes = self.dump_image_bytes(image_metadata)
 
-        uploader = Uploader(self.upload_options, user_session=self.user_session)
+        uploader = Uploader(self.upload_options, user_session=self.user_session, emitter=self.emitter)
 
         session_key = uploader._gen_session_key(io.BytesIO(image_bytes), image_progress)
 
-        file_handle = self._get_cached_file_handle(session_key)
-
-        if file_handle is None:
-            # image_progress will be updated during uploading
-            file_handle = uploader.upload_stream(
-                io.BytesIO(image_bytes),
-                session_key=session_key,
-                progress=image_progress,
-            )
-            self._set_file_handle_cache(session_key, file_handle)
+        with self._cache_lock:
+            file_handle = self._get_cached_file_handle(session_key)
+            if file_handle is None:
+                file_handle = uploader.upload_stream(
+                    io.BytesIO(image_bytes),
+                    session_key=session_key,
+                    progress=image_progress,
+                )
+                self._set_file_handle_cache(session_key, file_handle)
 
         return file_handle
 
@@ -800,6 +809,13 @@ class SingleImageUploader:
             LOG.debug("Dry-run mode enabled, skipping caching upload file handles")
             return None
 
+        if upload_options.upload_cache_path is not None:
+            cache_path = upload_options.upload_cache_path
+            cache_path.parent.mkdir(parents=True, exist_ok=True)
+            cache = history.PersistentCache(str(cache_path.resolve()))
+            cache.clear_expired()
+            return cache
+
         # Different python/CLI versions use different cache (dbm) formats.
         # Separate them to avoid conflicts
         py_version_parts = [str(part) for part in sys.version_info[:3]]
step 1

Plan: inspect the uploader, cache, and existing tests; reproduce the concurrency/cache behavior; identify shared-state and cache initialization races; implement a focused synchronization and API fix; then run targeted and broader tests and review the diff.

$ ls && find . -maxdepth 3 -type f | head -80 && grep -R "class ImageSequenceUploader\|class UploadOptions\|PersistentCache\|cached_image_uploader" -n . --exclude-dir=.git
output · 46 lines
CHANGELOG.md
CODE_OF_CONDUCT.md
CONTRIBUTING.md
LICENSE
README.md
README_osx_package.txt
docs
mapillary_tools
mapillary_tools.egg-info
pyinstaller
pyproject.toml
schema
34 more lines
script
tests
./.git/COMMIT_EDITMSG
./.git/config
./.git/description
./.git/HEAD
./.git/hooks/applypatch-msg.sample
./.git/hooks/commit-msg.sample
./.git/hooks/fsmonitor-watchman.sample
./.git/hooks/post-update.sample
./.git/hooks/pre-applypatch.sample
./.git/hooks/pre-commit.sample
./.git/hooks/pre-merge-commit.sample
./.git/hooks/pre-push.sample
./.git/hooks/pre-rebase.sample
./.git/hooks/pre-receive.sample
./.git/hooks/prepare-commit-msg.sample
./.git/hooks/push-to-checkout.sample
./.git/hooks/sendemail-validate.sample
./.git/hooks/update.sample
./.git/index
./.git/info/exclude
./.git/logs/HEAD
./.github/workflows/codeql-analysis.yml
./.github/workflows/python-package.yml
./.github/workflows/release.yml
./.gitignore

…[3286 chars omitted]…
ts/unit/test_persistent_cache.py:190:    cache = PersistentCache(cache_file)
./tests/unit/test_persistent_cache.py:203:    cache = PersistentCache(cache_file)
./tests/unit/test_persistent_cache.py:228:    cache = PersistentCache(cache_file)
./tests/unit/test_persistent_cache.py:272:    cache = PersistentCache(cache_file)
./tests/unit/test_persistent_cache.py:286:    cache = PersistentCache(cache_file)
step 2
$ sed -n '1,180p' mapillary_tools/uploader.py; sed -n '480,880p' mapillary_tools/uploader.py; sed -n '1,220p' mapillary_tools/history.py; grep -R "ImageSequenceUploader\|upload_images" -n tests | head -80
output · 63 lines
from __future__ import annotations

import concurrent.futures
import dataclasses
import io
import json
import logging
import os
import queue
import struct
import sys
import tempfile
51 more lines
import threading
import time
import typing as T
import uuid
import zipfile
from contextlib import contextmanager
from pathlib import Path

if sys.version_info >= (3, 11):
    from typing import Required
else:
    from typing_extensions import Required

import requests

from . import (
    api_v4,
    config,
    constants,
    exif_write,
    geo,
    history,
    telemetry,
    types,
    upload_api_v4,
    utils,
    VERSION,
)
from .camm import camm_builder, camm_parser
from .gpmf import gpmf_parser
from .mp4 import simple_mp4_builder
from .serializer.description import (
    desc_file_to_exif,
    DescriptionJSONSerializer,
    validate_image_desc,
)


LOG = logging.getLogger(__name__)


@dataclasses.d
…[8630 chars omitted]…
, dict):
            LOG.warning(f"Invalid cache value format: {payload}")
            return {}

        return payload
tests/integration/test_process_and_upload.py:152:def test_process_and_upload_images_only(
tests/unit/test_uploader.py:27:def test_upload_images(setup_unittest_data: py.path.local, setup_upload: py.path.local):
tests/unit/test_uploader.py:68:def test_upload_images_multiple_sequences(
step 3
$ sed -n '500,880p' mapillary_tools/uploader.py; sed -n '1,130p' tests/unit/test_uploader.py; grep -n "class .*Uploader\|def upload_images\|_get_cached\|cache" mapillary_tools/uploader.py
output · 42 lines
sequence_md5sum = zip_metadata.get("sequence_md5sum")

        if not sequence_md5sum and not isinstance(sequence_md5sum, str):
            raise InvalidMapillaryZipFileError("No sequence_md5sum found")

        return sequence_md5sum

    @classmethod
    @contextmanager
    def _wip_file_context(cls, wip_path: Path):
        try:
            os.remove(wip_path)
30 more lines
        except FileNotFoundError:
            pass
        try:
            yield wip_path

            with wip_path.open("rb") as fp:
                upload_md5sum = utils.md5sum_fp(fp).hexdigest()

            done_path = wip_path.parent.joinpath(
                _suffix_session_key(upload_md5sum, api_v4.ClusterFileType.ZIP)
            )

            try:
                os.remove(done_path)
            except FileNotFoundError:
                pass
            wip_path.rename(done_path)
        finally:
         
…[8629 chars omitted]…
istentCache(str(cache_path.resolve()))
830:        cache.clear_expired()
832:        return cache
834:    def _get_cached_file_handle(self, key: str) -> str | None:
835:        if self.cache is None:
841:        return self.cache.get(key)
843:    def _set_file_handle_cache(self, key: str, value: str) -> None:
844:        if self.cache is None:
850:        self.cache.set(key, value)
853:class Uploader:
step 4
$ sed -n '620,860p' mapillary_tools/uploader.py; sed -n '1,100p' mapillary_tools/uploader.py | tail -20; grep -R "upload_cache\|SingleImageUploader" -n tests mapillary_tools | head -80; sed -n '220,360p' mapillary_tools/uploader.py
output · 35 lines
def _upload_images_parallel(
        self,
        sequence: T.Sequence[types.ImageMetadata],
        sequence_progress: dict[str, T.Any],
    ) -> list[str]:
        if not sequence:
            return []

        max_workers = min(self.upload_options.num_upload_workers, len(sequence))

        # Lock is used to synchronize event emission
        lock = threading.Lock()
23 more lines
        # Push all images into the queue
        image_queue: queue.Queue[tuple[int, types.ImageMetadata]] = queue.Queue()
        for idx, image_metadata in enumerate(sequence):
            image_queue.put((idx, image_metadata))

        upload_interrupted = threading.Event()

        with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor:
            futures = [
                executor.submit(
                    self._upload_images_from_queue,
                    image_queue,
       
…[8629 chars omitted]…
roup_and_sort_images(metadatas)
        os.makedirs(zip_dir, exist_ok=True)

        for sequence_uuid, sequence in sequences.items():
            _validate_metadatas(sequence)
            # For atomicity we write into a WIP file and then rename to the final file
            wip_zip_filename = zip_dir.joinpath(
                f".mly_zip_{uuid.uuid4()}_{sequence_uuid}_{os.getpid()}_{int(time.time())}"
step 5
$ sed -n '730,850p' mapillary_tools/uploader.py; grep -R "UPLOAD_CACHE\|CACHE" -n mapillary_tools/config.py mapillary_tools/constants.py mapillary_tools/commands; sed -n '1,80p' mapillary_tools/history.py; sed -n '1,90p' tests/unit/test_persistent_cache.py
output · 43 lines
return indexed_file_handles


class SingleImageUploader:
    def __init__(
        self,
        upload_options: UploadOptions,
        user_session: requests.Session | None = None,
    ):
        self.upload_options = upload_options
        self.user_session = user_session
        self.cache = self._maybe_create_persistent_cache_instance(
31 more lines
            self.upload_options.user_items, upload_options
        )

    def upload(
        self, image_metadata: types.ImageMetadata, image_progress: dict[str, T.Any]
    ) -> str:
        image_bytes = self.dump_image_bytes(image_metadata)

        uploader = Uploader(self.upload_options, user_session=self.user_session)

        session_key = uploader._gen_session_key(io.BytesIO(image_bytes), image_progress)

        file_handle = self._get_cached_file_handle(session_key)

        if file_handle is None:
            # image_progress wil
…[7828 chars omitted]…
me, should_exist
):
    """Test expiration with different timing combinations."""
    cache_file = os.path.join(
        tmpdir, f"cache_param_exp_{dbm_backend}_{expire_time}_{sleep_time}"
    )
    cache = PersistentCache(cache_file)

    key = f"key_expires_in_{expire_time}_sleeps_{sleep_time}"
    cache.set(key, "test_value", expires_in=expire_time)

    time.sleep(sleep_time)

    if should_exist:
step 6
$ sed -n '850,1030p' mapillary_tools/uploader.py; sed -n '140,165p' mapillary_tools/constants.py; grep -R "ImageSequenceUploader(" -n mapillary_tools tests
output · 42 lines
self.cache.set(key, value)


class Uploader:
    def __init__(
        self,
        upload_options: UploadOptions,
        user_session: requests.Session | None = None,
        emitter: EventEmitter | None = None,
    ):
        self.upload_options = upload_options
        self.user_session = user_session
30 more lines
        if emitter is None:
            # An empty event emitter that does nothing
            self.emitter = EventEmitter()
        else:
            self.emitter = emitter

    def upload_stream(
        self,
        fp: T.IO[bytes],
        session_key: str | None = None,
        progress: dict[str, T.Any] | None = None,
    ) -> str:
        if progress is None:
            progress = {}

        if session_key is None:
            session_key = self._gen_session_key(fp, progress)

        fp.seek(0, io.SEEK_END)
        entity_size = fp.tell()

        progress["entity
…[6276 chars omitted]…
TE: Sequences themselves are uploaded sequentially, not in parallel.
MAX_IMAGE_UPLOAD_WORKERS: int = int(
    os.getenv(_ENV_PREFIX + "MAX_IMAGE_UPLOAD_WORKERS", 4)
)
# The chunk size in MB (see chunked transfer encoding https://en.wikipedia.org/wiki/Chunked_transfer_encoding)
# for uploading data to MLY upload service.
mapillary_tools/upload.py:559:    image_uploader = uploader.ImageSequenceUploader(
step 7
$ sed -n '530,580p' mapillary_tools/upload.py; grep -n "def test_.*cache\|cache" tests/unit/test_uploader.py tests/unit/test_persistent_cache.py | tail -30; grep -n "class EventEmitter" mapillary_tools/uploader.py
output · 33 lines
except requests.HTTPError as exc:
            LOG.warning(
                f"HTTPError from logging action {action}: {http.readable_http_error(exc)}"
            )
        except Exception:
            LOG.warning(f"Error from logging action {action}", exc_info=True)


_M = T.TypeVar("_M", bound=types.Metadata)


def _find_metadata_with_filename_existed_in(
21 more lines
    metadatas: T.Iterable[_M], paths: T.Iterable[Path]
) -> list[_M]:
    resolved_image_paths = set(p.resolve() for p in paths)
    return [d for d in metadatas if d.filename.resolve() in resolved_image_paths]


def _gen_upload_everything(
    mly_uploader: uploader.Uploader,
    metadatas: T.Sequence[types.Metadata],
    import_paths: T.Sequence[Path],
    skip_subfolders: bool,
):
    # Upload images
    image_metadatas = _find_metadata_with_filename_existed_in(
        (m for m in metadatas if isinstance(m, types.ImageM
…[3310 chars omitted]…
ent_cache.py:290:    assert cache._is_expired({"expires_at": time.time() + 10}) is False
tests/unit/test_persistent_cache.py:291:    assert cache._is_expired({}) is False
tests/unit/test_persistent_cache.py:292:    assert cache._is_expired({"expires_at": "not a number"}) is False
tests/unit/test_persistent_cache.py:293:    assert cache._is_expired({"expires_at": None}) is False
197:class EventEmitter:
step 8
$ python - <<'PY'
output · 5 lines
Traceback (most recent call last):
  File "<stdin>", line 121, in <module>
AssertionError
 mapillary_tools/history.py | 9 +++++++++
 1 file changed, 9 insertions(+)
step 9
$ grep -n "Different python\|def __init__\|file_handle = self._get" mapillary_tools/uploader.py | tail -15; git diff -- mapillary_tools/uploader.py
output · 7 lines
148:    def __init__(self, message: str, image_path: Path):
200:    def __init__(self):
537:    def __init__(self, upload_options: UploadOptions, emitter: EventEmitter):
735:    def __init__(
755:        file_handle = self._get_cached_file_handle(session_key)
803:        # Different python/CLI versions use different cache (dbm) formats.
854:    def __init__(
step 10
$ python - <<'PY'
output · 2 lines
.................................................................        [100%]
65 passed in 41.69s
step 11
$ git diff --check; git diff --stat; python -m pytest tests/unit -q --disable-warnings --maxfail=1. 2>/dev/null || true; git diff -- mapillary_tools/uploader.py mapillary_tools/history.py | sed -n '1,240p'
output · 32 lines
mapillary_tools/history.py  |  9 +++++++++
 mapillary_tools/uploader.py | 46 ++++++++++++++++++++++++++++++---------------
 2 files changed, 40 insertions(+), 15 deletions(-)
diff --git a/mapillary_tools/history.py b/mapillary_tools/history.py
index 8b7251e..4011f0a 100644
--- a/mapillary_tools/history.py
+++ b/mapillary_tools/history.py
@@ -137,6 +137,15 @@ class PersistentCache:
             f"Cached file handle for {key} ({(time.perf_counter() - s) * 1000:.0f} ms)"
         )
 
+    def keys(self) -> list[str]:
20 more lines
+        """Return the keys currently stored in the cache."""
+        with self._lock:
+            with dbm.open(self._file, flag="c") as db:
+                return [
+                    key.decode("utf-8") if isinstance(key, bytes) else key
+                    for key in db.keys()
+                ]
+
     def clear_expired(self) -> list[str]:
         s = time.perf_c
…[3953 chars omitted]…
cache_path
+            cache_path.parent.mkdir(parents=True, exist_ok=True)
+            cache = history.PersistentCache(str(cache_path.resolve()))
+            cache.clear_expired()
+            return cache
+
         # Different python/CLI versions use different cache (dbm) formats.
         # Separate them to avoid conflicts
         py_version_parts = [str(part) for part in sys.version_info[:3]]
step 12
$ python -m pytest tests/unit -q --disable-warnings --maxfail=1; git status --short
output · 7 lines
.......................................................................s [ 30%]
sssssssss............................................................... [ 61%]
........................................................................ [ 91%]
...................                                                      [100%]
225 passed, 10 skipped in 46.90s
 M mapillary_tools/history.py
 M mapillary_tools/uploader.py
step 13
$ python -m black mapillary_tools/uploader.py mapillary_tools/history.py >/dev/null 2>&1 || true; git diff --check; git diff --stat
output · 3 lines
mapillary_tools/history.py  |  9 +++++++++
 mapillary_tools/uploader.py | 46 ++++++++++++++++++++++++++++++---------------
 2 files changed, 40 insertions(+), 15 deletions(-)
step 14
$ done: {"notes": "Implemented shared, transactionally locked image cache access across parallel workers; added configurable upload_cache_path, exposed cached_image_uploader, preserved cache lifetime semantics, forwarded emitters for image uploads, and added PersistentCache.keys(). Unit suite passes: 225 passed, 10 skipped."}
output · 1 lines
submitted.
GLM-5.3 Flash · attempt 2 · solved · 64 steps · $0.059
Submitted diff · 253 lines
diff --git a/mapillary_tools/history.py b/mapillary_tools/history.py
index 8b7251e..27990e7 100644
--- a/mapillary_tools/history.py
+++ b/mapillary_tools/history.py
@@ -137,6 +137,20 @@ class PersistentCache:
             f"Cached file handle for {key} ({(time.perf_counter() - s) * 1000:.0f} ms)"
         )
 
+    def keys(self) -> list[str]:
+        """Return the (non-expired) keys stored in the cache"""
+
+        keys: list[str] = []
+
+        with self._lock:
+            with dbm.open(self._file, flag="c") as db:
+                for key in db.keys():
+                    if isinstance(key, bytes):
+                        key = key.decode("utf-8", errors="replace")
+                    keys.append(T.cast(str, key))
+
+        return keys
+
     def clear_expired(self) -> list[str]:
         s = time.perf_counter()
 
diff --git a/mapillary_tools/uploader.py b/mapillary_tools/uploader.py
index 8b2e2ad..100f4bd 100644
--- a/mapillary_tools/uploader.py
+++ b/mapillary_tools/uploader.py
@@ -59,6 +59,7 @@ class UploadOptions:
     dry_run: bool = False
     nofinish: bool = False
     noresume: bool = False
+    upload_cache_path: Path | None = None
 
     def __post_init__(self):
         if self.num_upload_workers <= 0:
@@ -537,6 +538,10 @@ class ImageSequenceUploader:
     def __init__(self, upload_options: UploadOptions, emitter: EventEmitter):
         self.upload_options = upload_options
         self.emitter = emitter
+        # Create the cached uploader (and thus the persistent cache) only once,
+        # from the options used to construct this uploader, so that the cache is
+        # shared consistently by all the concurrent image uploads
+        self.cached_image_uploader = CachedUploader(upload_options)
 
     def upload_images(
         self, image_metadatas: T.Sequence[types.ImageMetadata]
@@ -689,7 +694,9 @@ class ImageSequenceUploader:
             self.upload_options.user_items["user_upload_token"]
         ) as user_session:
             single_image_uploader = SingleImageUploader(
-                self.upload_options, user_session=user_session
+                self.upload_options,
+                user_session=user_session,
+                cache=self.cached_image_uploader.cache,
             )
 
             while True:
@@ -701,6 +708,7 @@ class ImageSequenceUploader:
 
                 # Main thread will handle the interruption
                 if upload_interrupted.is_set():
+                    image_queue.task_done()
                     break
 
                 # Create a new mutatble progress to keep the sequence_progress immutable
@@ -719,6 +727,7 @@ class ImageSequenceUploader:
 
                 # Main thread will handle the interruption
                 if upload_interrupted.is_set():
+                    image_queue.task_done()
                     break
 
                 with lock:
@@ -731,18 +740,102 @@ class ImageSequenceUploader:
         return indexed_file_handles
 
 
-class SingleImageUploader:
+class CachedUploader:
+    """
+    Uploader that caches upload file handles in a persistent cache, so that
+    repeated uploads of the same content can reuse the previously uploaded
+    file handles instead of uploading them again.
+
+    The cache instance is created once (from the upload options given to the
+    constructor) and shared by all the image uploads, so concurrent uploads
+    won't race while accessing the cached state.
+    """
+
     def __init__(
         self,
         upload_options: UploadOptions,
         user_session: requests.Session | None = None,
+        cache: history.PersistentCache | None = None,
     ):
         self.upload_options = upload_options
         self.user_session = user_session
-        self.cache = self._maybe_create_persistent_cache_instance(
-            self.upload_options.user_items, upload_options
-        )
+        if cache is None:
+            self.cache = self._maybe_create_persistent_cache_instance(upload_options)
+        else:
+            self.cache = cache
+
+    @classmethod
+    def _maybe_create_persistent_cache_instance(
+        cls, upload_options: UploadOptions
+    ) -> history.PersistentCache | None:
+        if upload_options.dry_run:
+            LOG.debug("Dry-run mode enabled, skipping caching upload file handles")
+            return None
+
+        if upload_options.upload_cache_path is not None:
+            cache_path = upload_options.upload_cache_path
+            cache_path.parent.mkdir(parents=True, exist_ok=True)
+        else:
+            if not constants.UPLOAD_CACHE_DIR:
+                LOG.debug(
+                    "Upload cache directory is set empty, skipping caching upload file handles"
+                )
+                return None
+
+            # Different python/CLI versions use different cache (dbm) formats.
+            # Separate them to avoid conflicts
+            py_version_parts = [str(part) for part in sys.version_info[:3]]
+            version = f"py_{'_'.join(py_version_parts)}_{VERSION}"
+
+            cache_path_dir = (
+                Path(constants.UPLOAD_CACHE_DIR)
+                .joinpath(version)
+                .joinpath(api_v4.MAPILLARY_CLIENT_TOKEN.replace("|", "_"))
+                .joinpath(
+                    upload_options.user_items.get(
+                        "MAPSettingsUserKey",
+                        upload_options.user_items["user_upload_token"],
+                    )
+                )
+            )
+            cache_path_dir.mkdir(parents=True, exist_ok=True)
+            cache_path = cache_path_dir.joinpath("cached_file_handles")
+
+            # Sanitize sensitive segments for logging
+            sanitized_cache_path = (
+                Path(constants.UPLOAD_CACHE_DIR)
+                .joinpath(version)
+                .joinpath("***")
+                .joinpath("***")
+                .joinpath("cached_file_handles")
+            )
+            LOG.debug(f"File handle cache path: {sanitized_cache_path}")
… [3052 more characters]
step 1
$ find /repo -type f -name "*.py" | head -50; echo ---; ls /repo
output · 44 lines
/repo/mapillary_tools/api_v4.py
/repo/mapillary_tools/authenticate.py
/repo/mapillary_tools/blackvue_parser.py
/repo/mapillary_tools/camm/camm_builder.py
/repo/mapillary_tools/camm/camm_parser.py
/repo/mapillary_tools/commands/authenticate.py
/repo/mapillary_tools/commands/process.py
/repo/mapillary_tools/commands/process_and_upload.py
/repo/mapillary_tools/commands/sample_video.py
/repo/mapillary_tools/commands/upload.py
/repo/mapillary_tools/commands/video_process.py
/repo/mapillary_tools/commands/video_process_and_upload.py
32 more lines
/repo/mapillary_tools/commands/zip.py
/repo/mapillary_tools/commands/__init__.py
/repo/mapillary_tools/commands/__main__.py
/repo/mapillary_tools/config.py
/repo/mapillary_tools/constants.py
/repo/mapillary_tools/exceptions.py
/repo/mapillary_tools/exiftool_read.py
/repo/mapillary_tools/exiftool_read_video.py
/repo/mapillary_tools/exiftool_runner.py
/repo/mapill
…[1167 chars omitted]…
ctors/native.py
/repo/mapillary_tools/geotag/__init__.py
/repo/mapillary_tools/gpmf/gpmf_gps_filter.py
/repo/mapillary_tools/gpmf/gpmf_parser.py
/repo/mapillary_tools/gpmf/gps_filter.py
/repo/mapillary_tools/history.py
---
CHANGELOG.md
CODE_OF_CONDUCT.md
CONTRIBUTING.md
LICENSE
README.md
README_osx_package.txt
docs
mapillary_tools
mapillary_tools.egg-info
pyinstaller
pyproject.toml
schema
script
tests
step 2
$ ls /repo/mapillary_tools; echo ---; grep -rn "ImageSequenceUploader\|cached_image_uploader\|upload_cache_path" /repo/mapillary_tools /repo/tests | head -40
output · 35 lines
__init__.py
api_v4.py
authenticate.py
blackvue_parser.py
camm
commands
config.py
constants.py
exceptions.py
exif_read.py
exif_write.py
exiftool_read.py
23 more lines
exiftool_read_video.py
exiftool_runner.py
ffmpeg.py
geo.py
geotag
gpmf
history.py
http.py
ipc.py
mp4
process_geotag_properties.py
process_sequence_properties.py
sample_video.py
serializer
telemetry.py
types.py
upload.py
upload_api_v4.py
uploader.py
utils.py
---
/repo/mapillary_tools/upload.py:559:    image_uploader = uploader.ImageSequenceUploader(
/repo/mapillary_tools/uploader.py:536:class ImageSequenceUploader:
step 3
$ sed -n 1,120p /repo/mapillary_tools/uploader.py; echo =====; sed -n 480,760p /repo/mapillary_tools/uploader.py
output · 65 lines
from __future__ import annotations

import concurrent.futures
import dataclasses
import io
import json
import logging
import os
import queue
import struct
import sys
import tempfile
53 more lines
import threading
import time
import typing as T
import uuid
import zipfile
from contextlib import contextmanager
from pathlib import Path

if sys.version_info >= (3, 11):
    from typing import Required
else:
    from typing_extensions import Required

import requests

from . import (
    api_v4,
    config,
    constants,
    exif_write,
    geo,
    history,
    telemetry,
    types,
    upload_api_v4,
    utils,
    VERSION,
)
from .camm import camm_builder, camm_parser
from .gpmf import gpmf_parser
from .mp4 import simple_mp4_builder
from .serializer.description import (
    desc_file_to_exif,
    DescriptionJSONSerializer,
    validate_image_desc,
)


LOG = logging.getLogger(__name__)


@dataclasses.d
…[8629 chars omitted]…
= Uploader(self.upload_options, user_session=self.user_session)

        session_key = uploader._gen_session_key(io.BytesIO(image_bytes), image_progress)

        file_handle = self._get_cached_file_handle(session_key)

        if file_handle is None:
            # image_progress will be updated during uploading
            file_handle = uploader.upload_stream(
                io.BytesIO(image_bytes),
step 4
$ sed -n 760,1000p /repo/mapillary_tools/uploader.py
output · 40 lines
io.BytesIO(image_bytes),
                session_key=session_key,
                progress=image_progress,
            )
            self._set_file_handle_cache(session_key, file_handle)

        return file_handle

    @classmethod
    def dump_image_bytes(cls, metadata: types.ImageMetadata) -> bytes:
        try:
            edit = exif_write.ExifEdit(metadata.filename)
28 more lines
        except struct.error as ex:
            raise ExifError(f"Failed to load EXIF: {ex}", metadata.filename) from ex

        # The cast is to fix the type checker error
        edit.add_image_description(
            T.cast(
                T.Dict, desc_file_to_exif(DescriptionJSONSerializer.as_desc(metadata))
            )
        )

        try:
            return edit.dump_image_bytes()
        except struct.error as ex:
            raise ExifError(
                f"Failed to dump EXIF bytes: 
…[6860 chars omitted]…
retries in the caller
            retries += 1
            if _is_immediate_retriable_exception(ex):
                sleep_for = 0
            else:
                sleep_for = min(2**retries, 16)
            LOG.info(
                f"Retrying in {sleep_for} seconds ({retries}/{constants.MAX_UPLOAD_RETRIES})"
            )
            if sleep_for:
                time.sleep(sleep_for)
        else:
step 5
$ sed -n 536,660p /repo/mapillary_tools/uploader.py; echo ====; sed -n 1,80p /repo/mapillary_tools/history.py
output · 31 lines
class ImageSequenceUploader:
    def __init__(self, upload_options: UploadOptions, emitter: EventEmitter):
        self.upload_options = upload_options
        self.emitter = emitter

    def upload_images(
        self, image_metadatas: T.Sequence[types.ImageMetadata]
    ) -> T.Generator[tuple[str, UploadResult], None, None]:
        sequences = types.group_and_sort_images(image_metadatas)

        for sequence_idx, (sequence_uuid, sequence) in enumerate(sequences.items()):
            LOG.debug(f"Checksum for image sequence {sequence_uuid}...")
19 more lines
            sequence_md5sum = types.update_sequence_md5sum(sequence)

            sequence_progress: SequenceProgress = {
                "sequence_idx": sequence_idx,
                "total_sequence_count": len(sequences),
                "sequence_image_count": len(sequence),
                "sequence_uuid": sequence_uuid,
                
…[5434 chars omitted]…
  summary: JSONDict,
    metadatas: T.Sequence[types.Metadata] | None = None,
) -> None:
    if not constants.MAPILLARY_UPLOAD_HISTORY_PATH:
        return
    path = history_desc_path(md5sum)
    LOG.debug("Writing upload history: %s", path)
    path.resolve().parent.mkdir(parents=True, exist_ok=True)
    history: dict[str, T.Any] = {"params": params, "summary": summary}
    if metadatas is not None:
step 6
$ grep -n "class PersistentCache" -A 60 /repo/mapillary_tools/history.py; echo ===; grep -n "UPLOAD_CACHE_DIR" /repo/mapillary_tools/constants.py
output · 43 lines
88:class PersistentCache:
89-    _lock: contextlib.nullcontext | threading.Lock
90-
91-    def __init__(self, file: str):
92-        # SQLite3 backend supports concurrent access without a lock
93-        if dbm.whichdb(file) == "dbm.sqlite3":
94-            self._lock = contextlib.nullcontext()
95-        else:
96-            self._lock = threading.Lock()
97-        self._file = file
98-
99-    def get(self, key: str) -> str | None:
31 more lines
100-        s = time.perf_counter()
101-
102-        with self._lock:
103-            with dbm.open(self._file, flag="c") as db:
104-                value: bytes | None = db.get(key)
105-
106-        if value is None:
107-            return None
108-
109-        payload = self._decode(value)
110-
111-        if self._is_expired(payload):
112-            return None
113-
114-        file_handle = payload.get("file_handle")
115-
116-        LOG.debug(
117-  
…[797 chars omitted]…
lf) -> list[str]:
141-        s = time.perf_counter()
142-
143-        expired_keys: list[str] = []
144-
145-        with self._lock:
146-            with dbm.open(self._file, flag="c") as db:
147-                if hasattr(db, "items"):
148-                    items: T.Iterable[tuple[str | bytes, bytes]] = db.items()
===
150:UPLOAD_CACHE_DIR: str = os.getenv(
151:    _ENV_PREFIX + "UPLOAD_CACHE_DIR",
step 7
$ sed -n 148,200p /repo/mapillary_tools/history.py; echo ===; sed -n 145,160p /repo/mapillary_tools/constants.py; echo ===; ls /repo/tests; grep -rn "ImageSequenceUploader\|PersistentCache\|upload_cache" /repo/tests | head
output · 30 lines
items: T.Iterable[tuple[str | bytes, bytes]] = db.items()
                else:
                    items = ((key, db[key]) for key in db.keys())

                for key, value in items:
                    payload = self._decode(value)
                    if self._is_expired(payload):
                        del db[key]
                        expired_keys.append(T.cast(str, key))

        if expired_keys:
            LOG.debug(
18 more lines
                f"Cleared {len(expired_keys)} expired entries from the cache ({(time.perf_counter() - s) * 1000:.0f} ms)"
            )

        return expired_keys

    def _is_expired(self, payload: JSONDict) -> bool:
        expires_at = payload.get("expires_at")
        if isinstance(expires_at, (int, float)):
            return expires_at is None or expires_at <= time.time()
        return False

    def _decode(self, value: bytes) -
…[1691 chars omitted]…
ersistent_cache.py:48:    cache = PersistentCache(cache_file)
/repo/tests/unit/test_persistent_cache.py:83:    cache = PersistentCache(cache_file)
/repo/tests/unit/test_persistent_cache.py:100:    cache = PersistentCache(cache_file)
/repo/tests/unit/test_persistent_cache.py:123:    cache = PersistentCache(cache_file)
/repo/tests/unit/test_persistent_cache.py:149:    cache = PersistentCache(cache_file)
step 8
$ cat /repo/tests/unit/test_persistent_cache.py; echo ===; ls /repo/tests/integration /repo/tests/unit
output · 53 lines
import dbm
import os
import threading
import time

import pytest

from mapillary_tools.history import PersistentCache


# DBM backends to test with
DBM_BACKENDS = ["dbm.sqlite3", "dbm.gnu", "dbm.ndbm", "dbm.dumb"]
41 more lines

@pytest.mark.parametrize("dbm_backend", DBM_BACKENDS)
def test_basic_operations_with_backend(tmpdir, dbm_backend):
    """Test basic operations with different DBM backends.

    Note: This is a demonstration of pytest's parametrize feature.
    The actual PersistentCache class might not support specifying backends.
    """
    cache_file = os.path.join(tmpdir, dbm_backend)
    # Here you would use the backend if the cache implementation supported it
    cache = PersistentCache(cache_file)

    # Perform basic operations
    cache.set("test_key", "test_value")
    assert cache.get("test_key") == "test_value"

    # Add specific test logic for different backends if needed
  
…[8628 chars omitted]…
amm_parser.py
test_config.py
test_description.py
test_exceptions.py
test_exifedit.py
test_exifread.py
test_ffmpeg.py
test_geo.py
test_gpmf_parser.py
test_gps_filter.py
test_io_utils.py
test_mp4_sample_parser.py
test_persistent_cache.py
test_sample_video.py
test_sequence_processing.py
test_simple_mp4_builder.py
test_simple_mp4_parser.py
test_types.py
test_upload_api_v4.py
test_uploader.py
test_utils.py
step 9
$ sed -n 1,200p /repo/tests/unit/test_uploader.py
output · 48 lines
import typing as T
from pathlib import Path

import py.path

import pytest

from mapillary_tools import api_v4, uploader
from mapillary_tools.serializer import description

from ..integration.fixtures import extract_all_uploaded_descs, setup_upload
36 more lines
IMPORT_PATH = "tests/unit/data"


@pytest.fixture
def setup_unittest_data(tmpdir: py.path.local):
    data_path = tmpdir.mkdir("data")
    source = py.path.local(IMPORT_PATH)
    source.copy(data_path)
    yield data_path
    if tmpdir.check():
        tmpdir.remove(ignore_errors=True)


def test_upload_images(setup_unittest_data: py.path.local, setup_upload: py.path.local):
    mly_uploader = uploader.Uploader(
        uploader.UploadOptions(
            {"user_upload_token": "YOUR_USER_ACCESS_TOKEN"}, dry_run=True
        )
    )
    test_exif = setup_unittest_data.join("test_exif.jpg")
    descs: T.List[description.DescriptionOrError] 
…[5134 chars omitted]…
 == len(actual_descs)


def test_upload_blackvue(tmpdir: py.path.local, setup_upload: py.path.local):
    mly_uploader = uploader.Uploader(
        uploader.UploadOptions(
            {
                "user_upload_token": "YOUR_USER_ACCESS_TOKEN",
                # will call the API for real
                # "MAPOrganizationKey": "3011753992432185",
            },
            dry_run=True,
        )
step 10
$ sed -n 200,400p /repo/tests/unit/test_uploader.py; echo ===; sed -n 1,120p /repo/tests/integration/fixtures.py
output · 51 lines
)
    )
    blackvue_path = tmpdir.join("blackvue.mp4")
    with open(blackvue_path, "wb") as fp:
        fp.write(b"this is a fake video")
    with Path(blackvue_path).open("rb") as fp:
        file_handle = mly_uploader.upload_stream(
            fp,
            session_key="this_is_a_blackvue.mp4",
        )
    cluster_id = mly_uploader.finish_upload(
        file_handle, api_v4.ClusterFileType.BLACKVUE
39 more lines
    )
    assert cluster_id == "0"
    assert setup_upload.join("this_is_a_blackvue.mp4").exists()
    with open(setup_upload.join("this_is_a_blackvue.mp4"), "rb") as fp:
        assert fp.read() == b"this is a fake video"


def test_upload_zip_with_emitter(
    setup_unittest_data: py.path.local, setup_upload: py.path.local
):
    emitter = uploader.EventEmitter()

    stats = {}

    @emitter.on("upload_start")
    def _upload_start(payload):
        assert payload["entit
…[3990 chars omitted]…
      return False
    return True


IS_FFMPEG_INSTALLED = _ffmpeg_installed()


def _exiftool_installed():
    try:
        subprocess.run(
            [EXIFTOOL_EXECUTABLE, "-ver"],
            stderr=subprocess.PIPE,
            stdout=subprocess.PIPE,
            check=True,
        )
    except FileNotFoundError:
        return False
    return True


IS_EXIFTOOL_INSTALLED = _exiftool_installed()
step 11
$ git log --oneline | head; git diff HEAD~1 --stat 2>/dev/null | tail -5
output · 1 lines
115aa82 base
step 12
$ grep -n "_gen_session_key\|_suffix_session_key\|_prefixed_uuid4\|_is_uuid" -A 20 /repo/mapillary_tools/uploader.py | sed -n 1,80p
output · 46 lines
521:                _suffix_session_key(upload_md5sum, api_v4.ClusterFileType.ZIP)
522-            )
523-
524-            try:
525-                os.remove(done_path)
526-            except FileNotFoundError:
527-                pass
528-            wip_path.rename(done_path)
529-        finally:
530-            try:
531-                os.remove(wip_path)
532-            except FileNotFoundError:
34 more lines
533-                pass
534-
535-
536-class ImageSequenceUploader:
537-    def __init__(self, upload_options: UploadOptions, emitter: EventEmitter):
538-        self.upload_options = upload_options
539-        self.emitter = emitter
540-
541-    def upload_images(
--
618:                manifest_fp, session_key=f"{_prefixed_uuid4()}.json"
619-            )
620-
621-    def _upload_images_parallel(
622-        self,
623-        sequence: T.Sequence[types.ImageMetadata],
624-        sequence
…[1391 chars omitted]…
(f"Failed to load EXIF: {ex}", metadata.filename) from ex
--
838:        if _is_uuid(key):
839-            return None
840-
841-        return self.cache.get(key)
842-
843-    def _set_file_handle_cache(self, key: str, value: str) -> None:
844-        if self.cache is None:
845-            return
846-
847:        if _is_uuid(key):
848-            return
849-
850-        self.cache.set(key, value)
851-
step 13
$ grep -n "def _gen_session_key" -A 25 /repo/mapillary_tools/uploader.py; grep -n "def _is_uuid\|def _prefixed_uuid4\|def _suffix_session_key" -A 10 /repo/mapillary_tools/uploader.py
output · 36 lines
1080:    def _gen_session_key(self, fp: T.IO[bytes], progress: dict[str, T.Any]) -> str:
1081-        if self.upload_options.noresume:
1082-            # Generate a unique UUID for session_key when noresume is True
1083-            # to prevent resuming from previous uploads
1084-            session_key = f"{_prefixed_uuid4()}"
1085-        else:
1086-            fp.seek(0, io.SEEK_SET)
1087-            session_key = utils.md5sum_fp(fp).hexdigest()
1088-
1089-        filetype = progress.get("file_type")
1090-        if filetype is not None:
1091-            session_key = _suffix_session_key(session_key, types.FileType(filetype))
24 more lines
1092-
1093-        return session_key
1094-
1095-
1096-def _validate_metadatas(metadatas: T.Sequence[types.ImageMetadata]):
1097-    for metadata in metadatas:
1098-        validate_image_desc(DescriptionJSONSerializer.as_desc(metadata))
1099-        if not me
…[412 chars omitted]…
56-    key = f"mly_tools_{key}{_SUFFIX_MAP[filetype]}"
1157-
1158-    assert _is_uuid(key) is is_uuid_before
1159-
1160-    return key
1161-
--
1163:def _prefixed_uuid4():
1164-    prefixed = f"uuid_{uuid.uuid4().hex}"
1165-    assert _is_uuid(prefixed)
1166-    return prefixed
1167-
1168-
1169:def _is_uuid(key: str) -> bool:
1170-    return key.startswith("uuid_") or key.startswith("mly_tools_uuid_")
step 14
$ sed -n 520,600p /repo/mapillary_tools/upload.py; grep -rn "SingleImageUploader" /repo/mapillary_tools
output · 33 lines
if constants.MAPILLARY_DISABLE_API_LOGGING:
        return

    payload_with_reason = {**payload, "reason": exc.__class__.__name__}
    action: api_v4.ActionType = "upload_failed_upload"

    with api_v4.create_client_session(disable_logging=True) as client_session:
        try:
            api_v4.log_event(client_session, action, payload_with_reason)
        except requests.HTTPError as exc:
            LOG.warning(
                f"HTTPError from logging action {action}: {http.readable_http_error(exc)}"
21 more lines
            )
        except Exception:
            LOG.warning(f"Error from logging action {action}", exc_info=True)


_M = T.TypeVar("_M", bound=types.Metadata)


def _find_metadata_with_filename_existed_in(
    metadatas: T.Iterable[_M], paths: T.Iterable[Path]
) -> list[_M]:
    resolved_image_paths = set(p.resolve() for p in paths)
    return [d for d in metadatas if d.fil
…[1793 chars omitted]…
eption:
    """
    Wrap the exception, or re-raise if it is a fatal error (i.e. there is no point to continue)
    """
/repo/mapillary_tools/uploader.py:474:                zipf.writestr(zipinfo, SingleImageUploader.dump_image_bytes(metadata))
/repo/mapillary_tools/uploader.py:691:            single_image_uploader = SingleImageUploader(
/repo/mapillary_tools/uploader.py:734:class SingleImageUploader:
step 15

Now I'll implement. Plan: add `upload_cache_path` to `UploadOptions`, add a `CachedUploader` class owning a single shared `PersistentCache`, have `ImageSequenceUploader` create one `cached_image_uploader` in `__init__` and share its cache with per-worker `SingleImageUploader`s, add `PersistentCache.keys()`, and fix the missing `task_done` on early break.

$ cd /repo && python - <<'PY'
output · 2 lines
mapillary_tools/uploader.py | 11 ++++++++++-
 1 file changed, 10 insertions(+), 1 deletion(-)
step 16
$ cd /repo && python - <<'PY'
output · 2 lines
767:    def upload(
862:    def _get_cached_file_handle(self, key: str) -> str | None:
step 17
$ cd /repo && sed -n 755,775p mapillary_tools/uploader.py; echo ...; sed -n 845,880p mapillary_tools/uploader.py
output · 44 lines
self,
        upload_options: UploadOptions,
        user_session: requests.Session | None = None,
        cache: history.PersistentCache | None = None,
    ):
        self.upload_options = upload_options
        self.user_session = user_session
        if cache is None:
            self.cache = self._maybe_create_persistent_cache_instance(upload_options)
        else:
            self.cache = cache
32 more lines
    def upload(
        self, image_metadata: types.ImageMetadata, image_progress: dict[str, T.Any]
    ) -> str:
        image_bytes = self.dump_image_bytes(image_metadata)

        uploader = Uploader(self.upload_options, user_session=self.user_session)

        session_key = uploader._gen_session_key(io.BytesIO(image_bytes), image_progress)

...
            cache_path = cache_path_dir.joinpath("cached_file_handles")

            # Sanitize sensitive segments for logging
       
…[440 chars omitted]…
 def _get_cached_file_handle(self, key: str) -> str | None:
        if self.cache is None:
            return None

        if _is_uuid(key):
            return None

        return self.cache.get(key)

    def _set_file_handle_cache(self, key: str, value: str) -> None:
        if self.cache is None:
            return

        if _is_uuid(key):
            return

        self.cache.set(key, value)
step 18
$ cd /repo && python - <<'PY'
output · 1 lines
ok
step 19
$ cd /repo && grep -n "class FakeUploadService" -A 40 mapillary_tools/upload_api_v4.py | head -60; grep -n "def upload\|def fetch\|class UploadService" mapillary_tools/upload_api_v4.py
output · 38 lines
173:class FakeUploadService(UploadService):
174-    """
175-    A mock upload service that simulates the upload process for testing purposes.
176-    It writes the uploaded data to a file in a temporary directory and generates a fake file handle.
177-    """
178-
179-    FILE_HANDLE_DIR: str = "file_handles"
180-
181-    def __init__(
182-        self,
183-        *args,
184-        upload_path: Path | None = None,
26 more lines
185-        transient_error_ratio: float = 0.0,
186-        **kwargs,
187-    ):
188-        super().__init__(*args, **kwargs)
189-        if upload_path is None:
190-            upload_path = Path(tempfile.gettempdir()).joinpath(
191-                "mapillary_public_uploads"
192-            )
193-        self._upload_path = upload_path
194-        self._transient_error_ratio = transient_error_ratio
195-
196-    @override
197-    def upload_shifted_chunks(
198-        self
…[529 chars omitted]…
with filename.open("ab") as fp:
212-            for chunk in shifted_chunks:
213-                self._randomly_raise_transient_error()
27:class UploadService:
39:    def fetch_offset(self) -> int:
113:    def upload_byte_stream(
128:    def upload_chunks(
141:    def upload_shifted_chunks(
197:    def upload_shifted_chunks(
227:    def fetch_offset(self) -> int:
237:    def upload_path(self) -> Path:
step 20
$ cd /repo && sed -n 213,260p mapillary_tools/upload_api_v4.py; grep -n "finish_upload\|create_user_session" mapillary_tools/api_v4.py | head
output · 35 lines
self._randomly_raise_transient_error()
                fp.write(chunk)
                self._randomly_raise_transient_error()

        file_handle_dir = self._upload_path.joinpath(self.FILE_HANDLE_DIR)
        file_handle_path = file_handle_dir.joinpath(self.session_key)
        if not file_handle_path.exists():
            os.makedirs(file_handle_dir, exist_ok=True)
            random_file_handle = uuid.uuid4().hex
            file_handle_path.write_text(random_file_handle)

        return file_handle_path.read_text()
23 more lines
    @override
    def fetch_offset(self) -> int:
        self._randomly_raise_transient_error()
        filename = self._upload_path.joinpath(self.session_key)
        if not filename.exists():
            return 0
        with open(filename, "rb") as fp:
            fp.seek(0, io.SEEK_END)
            return fp.tell()

    @property
    def upload_path
…[200 chars omitted]…
is for testing purposes only.
        """
        if random.random() <= self._transient_error_ratio:
            raise requests.ConnectionError(
                f"[TEST ONLY]: Transient error with ratio {self._transient_error_ratio}"
            )
40:def create_user_session(user_access_token: str) -> requests.Session:
149:def finish_upload(
155:    url = f"{MAPILLARY_GRAPH_API_ENDPOINT}/finish_upload"
step 21
$ cd /repo && cat > /tmp/repro.py <<'PY'
output · 11 lines
Traceback (most recent call last):
  File "/tmp/repro.py", line 47, in <module>
    u = uploader.ImageSequenceUploader(opts, emitter)
        ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/repo/mapillary_tools/uploader.py", line 544, in __init__
    self.cached_image_uploader = CachedUploader(upload_options)
                                 ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/repo/mapillary_tools/uploader.py", line 763, in __init__
    self.cache = self._maybe_create_persistent_cache_instance(upload_options)
                 ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
AttributeError: 'CachedUploader' object has no attribute '_maybe_create_persistent_cache_instance'
step 22
$ cd /repo && python - <<'PY'
output · 11 lines
Traceback (most recent call last):
  File "/tmp/repro.py", line 47, in <module>
    u = uploader.ImageSequenceUploader(opts, emitter)
        ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/repo/mapillary_tools/uploader.py", line 544, in __init__
    self.cached_image_uploader = CachedUploader(upload_options)
                                 ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/repo/mapillary_tools/uploader.py", line 763, in __init__
    self.cache = self._maybe_create_persistent_cache_instance(upload_options)
                 ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
AttributeError: 'CachedUploader' object has no attribute '_maybe_create_persistent_cache_instance'
step 23
$ cd /repo && sed -n 740,800p mapillary_tools/uploader.py; grep -n "class \|    def \|    @classmethod" mapillary_tools/uploader.py | sed -n '/73[0-9]/,$p' | awk 'NR>0' | sed -n 1,60p
output · 36 lines
return indexed_file_handles


class CachedUploader:
    """
    Uploader that caches upload file handles in a persistent cache, so that
    repeated uploads of the same content can reuse the previously uploaded
    file handles instead of uploading them again.

    The cache instance is created once (from the upload options given to the
    constructor) and shared by all the image uploads, so concurrent uploads
    won't race while accessing the cached state.
24 more lines
    """

    def __init__(
        self,
        upload_options: UploadOptions,
        user_session: requests.Session | None = None,
        cache: history.PersistentCache | None = None,
    ):
        self.upload_options = upload_options
        self.user_session = user_session
        if cache is None:
            self.cache = self._maybe_create_persistent_cache_instance(upload_options)
        else:
            self.c
…[1036 chars omitted]…
         cache_path_dir = (
                Path(constants.UPLOAD_CACHE_DIR)
                .joinpath(version)
                .joinpath(api_v4.MAPILLARY_CLIENT_TOKEN.replace("|", "_"))
                .joinpath(
                    upload_options.user_items.get(
                        "MAPSettingsUserKey",
                        upload_options.user_items["user_upload_token"],
                    )
step 24
$ cd /repo && python - <<'PY'
output · 14 lines
743:class CachedUploader:
768:class SingleImageUploader(CachedUploader):
822:class SingleImageUploader(CachedUploader):
Traceback (most recent call last):
  File "/tmp/repro.py", line 47, in <module>
    u = uploader.ImageSequenceUploader(opts, emitter)
        ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/repo/mapillary_tools/uploader.py", line 544, in __init__
    self.cached_image_uploader = CachedUploader(upload_options)
                                 ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/repo/mapillary_tools/uploader.py", line 763, in __init__
    self.cache = self._maybe_create_persistent_cache_instance(upload_options)
2 more lines
                 ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
AttributeError: 'CachedUploader' object has no attribute '_maybe_create_persistent_cache_instance'
step 25
$ cd /repo && sed -n 762,830p mapillary_tools/uploader.py
output · 37 lines
if cache is None:
            self.cache = self._maybe_create_persistent_cache_instance(upload_options)
        else:
            self.cache = cache


class SingleImageUploader(CachedUploader):
    @classmethod
    def _maybe_create_persistent_cache_instance(
        cls, upload_options: UploadOptions
    ) -> history.PersistentCache | None:
        if upload_options.dry_run:
25 more lines
            LOG.debug("Dry-run mode enabled, skipping caching upload file handles")
            return None

        if upload_options.upload_cache_path is not None:
            cache_path = upload_options.upload_cache_path
            cache_path.parent.mkdir(parents=True, exist_ok=True)
        else:
            if not constants.UPLOAD_CACHE_DIR:
                LOG.debug(
                    "Upload cache directory is set empty, skipping caching upload file handles"
                )
                ret
…[1338 chars omitted]…
 return cache


class SingleImageUploader(CachedUploader):
    def upload(
        self, image_metadata: types.ImageMetadata, image_progress: dict[str, T.Any]
    ) -> str:
        image_bytes = self.dump_image_bytes(image_metadata)

        uploader = Uploader(self.upload_options, user_session=self.user_session)

        session_key = uploader._gen_session_key(io.BytesIO(image_bytes), image_progress)
step 26
$ cd /repo && python - <<'PY'
output · 6 lines
results: [('seq1', UploadResult(result=None, error=NameError("name '_validate_metadatas' is not defined"))), ('seq2', UploadResult(result=None, error=NameError("name '_validate_metadatas' is not defined")))]
Traceback (most recent call last):
  File "/tmp/repro.py", line 57, in <module>
    assert r["result"], r
           ~^^^^^^^^^^
TypeError: 'UploadResult' object is not subscriptable
step 27
$ cd /repo && grep -n "_validate_metadatas\|def dump_image_bytes\|_is_immediate_retriable" mapillary_tools/uploader.py | head; git diff | grep "^-" | grep "_validate_metadatas"
output · 5 lines
358:            _validate_metadatas(sequence)
375:                _validate_metadatas(sequence)
579:        _validate_metadatas(sequence)
844:    def dump_image_bytes(cls, metadata: types.ImageMetadata) -> bytes:
-def _validate_metadatas(metadatas: T.Sequence[types.ImageMetadata]):
step 28
$ cd /repo && git show HEAD:mapillary_tools/uploader.py > /tmp/orig.py && python - <<'PY'
output · 7 lines
appended
results: [('seq1', UploadResult(result=None, error=NameError("name 'Uploader' is not defined"))), ('seq2', UploadResult(result=None, error=NameError("name 'Uploader' is not defined")))]
Traceback (most recent call last):
  File "/tmp/repro.py", line 57, in <module>
    assert r["result"], r
           ~^^^^^^^^^^
TypeError: 'UploadResult' object is not subscriptable
step 29
$ cd /repo && grep -n "^class \|^def " mapillary_tools/uploader.py
output · 20 lines
55:class UploadOptions:
74:class UploaderProgress(T.TypedDict, total=True):
105:class SequenceProgress(T.TypedDict, total=False):
134:class Progress(SequenceProgress, UploaderProgress):
138:class SequenceError(Exception):
148:class ExifError(SequenceError):
154:class InvalidMapillaryZipFileError(SequenceError):
198:class EventEmitter:
217:class UploadResult:
222:class VideoUploader:
319:class ZipUploader:
537:class ImageSequenceUploader:
8 more lines
743:class CachedUploader:
820:class SingleImageUploader(CachedUploader):
864:def _validate_metadatas(metadatas: T.Sequence[types.ImageMetadata]):
871:def _is_immediate_retriable_exception(ex: BaseException) -> bool:
887:def _is_retriable_exception(ex: BaseException) -> bool:
919:def _suffix_session_key(
931:def _prefixed_uuid4():
937:def _is_uuid(key: str) -> bool:
step 30
$ cd /repo && grep -n "^class \|^def " /tmp/orig.py
output · 20 lines
55:class UploadOptions:
73:class UploaderProgress(T.TypedDict, total=True):
104:class SequenceProgress(T.TypedDict, total=False):
133:class Progress(SequenceProgress, UploaderProgress):
137:class SequenceError(Exception):
147:class ExifError(SequenceError):
153:class InvalidMapillaryZipFileError(SequenceError):
197:class EventEmitter:
216:class UploadResult:
221:class VideoUploader:
318:class ZipUploader:
536:class ImageSequenceUploader:
8 more lines
734:class SingleImageUploader:
853:class Uploader:
1096:def _validate_metadatas(metadatas: T.Sequence[types.ImageMetadata]):
1103:def _is_immediate_retriable_exception(ex: BaseException) -> bool:
1119:def _is_retriable_exception(ex: BaseException) -> bool:
1151:def _suffix_session_key(
1163:def _prefixed_uuid4():
1169:def _is_uuid(key: str) -> bool:
step 31
$ cd /repo && python - <<'PY'
output · 6 lines
results: [('seq1', UploadResult(result=None, error=AttributeError("'SingleImageUploader' object has no attribute '_get_cached_file_handle'"))), ('seq2', UploadResult(result=None, error=AttributeError("'SingleImageUploader' object has no attribute '_get_cached_file_handle'")))]
Traceback (most recent call last):
  File "/tmp/repro.py", line 57, in <module>
    assert r["result"], r
           ~^^^^^^^^^^
TypeError: 'UploadResult' object is not subscriptable
step 32
$ cd /repo && sed -n 743,865p mapillary_tools/uploader.py
output · 42 lines
class CachedUploader:
    """
    Uploader that caches upload file handles in a persistent cache, so that
    repeated uploads of the same content can reuse the previously uploaded
    file handles instead of uploading them again.

    The cache instance is created once (from the upload options given to the
    constructor) and shared by all the image uploads, so concurrent uploads
    won't race while accessing the cached state.
    """

    def __init__(
30 more lines
        self,
        upload_options: UploadOptions,
        user_session: requests.Session | None = None,
        cache: history.PersistentCache | None = None,
    ):
        self.upload_options = upload_options
        self.user_session = user_session
        if cache is None:
            self.cache = self._maybe_create_persistent_cache_instance(upload_options)
        else:
            self.cache = cache

    @classmethod
    def
…[3198 chars omitted]…
it.add_image_description(
            T.cast(
                T.Dict, desc_file_to_exif(DescriptionJSONSerializer.as_desc(metadata))
            )
        )

        try:
            return edit.dump_image_bytes()
        except struct.error as ex:
            raise ExifError(
                f"Failed to dump EXIF bytes: {ex}", metadata.filename
            ) from ex

class Uploader:
    def __init__(
step 33
$ cd /repo && python - <<'PY'
output · 6 lines
results: [('seq1', UploadResult(result='123', error=None)), ('seq2', UploadResult(result='123', error=None))]
Traceback (most recent call last):
  File "/tmp/repro.py", line 57, in <module>
    assert r["result"], r
           ~^^^^^^^^^^
TypeError: 'UploadResult' object is not subscriptable
step 34
$ cd /repo && sed -i 's/r\["result"\]/r.result/g' /tmp/repro.py && python /tmp/repro.py
output · 6 lines
results: [('seq1', UploadResult(result='123', error=None)), ('seq2', UploadResult(result='123', error=None))]
Traceback (most recent call last):
  File "/tmp/repro.py", line 59, in <module>
    assert uploads["count"] == 14, uploads["count"]
           ^^^^^^^^^^^^^^^^^^^^^^
AssertionError: 4
step 35
$ cd /repo && python - <<'PY'
output · 3 lines
/tmp/tmpgaldyofu/data/test_exif.jpg 49411
/tmp/tmpgaldyofu/data/test_exif.jpg 49412
/tmp/tmpgaldyofu/data/test_exif.jpg 49412
step 36
$ cd /repo && python - <<'PY'
output · 3 lines
[('seq1', UploadResult(result='1', error=None)), ('seq2', UploadResult(result='1', error=None))]
keys: [None, None, None, None]
cache keys: ['mly_tools_04c92c8cd61b587ee68609dac32c941c.jpg', 'mly_tools_f92469558c0d52150a49b6ab0a7a05f3.jpg']
step 37
$ cd /repo && python - <<'PY'
output · 3 lines
58.5 2ba124e69c19a9a3b5bfa774e63a6aa9
58.51 a003718bdb608a1991c455d8624b7b07
58.52 d8488a7aa2aba1459f81df70b8ffc60c
step 38
$ cd /repo && python - <<'PY'
output · 7 lines
Error uploading sequence_seq1/test_exif.jpg at offset=None since begin_offset=None: ConnectionError: HTTPSConnectionPool(host='rupload.facebook.com', port=443): Max retries exceeded with url: /mapillary_public_uploads/mly_tools_d8488a7aa2aba1459f81df70b8ffc60c.jpg (Caused by NameResolutionError("HTTPSConnection(host='rupload.facebook.com', port=443): Failed to resolve 'rupload.facebook.com' ([Errno -3] Temporary failure in name resolution)"))
Error uploading sequence_seq1/test_exif.jpg at offset=None since begin_offset=None: ConnectionError: HTTPSConnectionPool(host='rupload.facebook.com', port=443): Max retries exceeded with url: /mapillary_public_uploads/mly_tools_d8488a7aa2aba1459f81df70b8ffc60c.jpg (Caused by NameResolutionError("HTTPSConnection(host='rupload.facebook.com', port=443): Failed to resolve 'rupload.facebook.com' ([Errno -3] Temporary failure in name resolution)"))
E
…[8631 chars omitted]…
begin_offset=None: ConnectionError: HTTPSConnectionPool(host='rupload.facebook.com', port=443): Max retries exceeded with url: /mapillary_public_uploads/mly_tools_d8488a7aa2aba1459f81df70b8ffc60c.jpg (Caused by NameResolutionError("HTTPSConnection(host='rupload.facebook.com', port=443): Failed to resolve 'rupload.facebook.com' ([Errno -3] Temporary failure in name resolution)"))

(command exceeded 600s)
step 39
$ cd /repo && timeout 120 python - <<'PY'
output · 4 lines
58.5 92ed607c130b4a768740176c7a78911f
58.51 dbc80f8d57d94dd78062e11ffb0f4bfa
58.52 94cbdf80aca34d8a8d5cb0752ee7d95f
['mly_tools_2ba124e69c19a9a3b5bfa774e63a6aa9.jpg', 'mly_tools_d8488a7aa2aba1459f81df70b8ffc60c.jpg', 'mly_tools_a003718bdb608a1991c455d8624b7b07.jpg']
step 40
$ cd /repo && timeout 300 python - <<'PY'
output · 2 lines
[('seq1', UploadResult(result=None, error=AttributeError("'C' object has no attribute 'session_key'")))]
[]
step 41
$ cd /repo && sed -n 27,60p mapillary_tools/upload_api_v4.py
output · 33 lines
class UploadService:
    """
    Upload byte streams to the Upload Service.
    """

    user_access_token: str
    session_key: str

    def __init__(self, user_session: requests.Session, session_key: str):
        self.user_session = user_session
        self.session_key = session_key
21 more lines
    def fetch_offset(self) -> int:
        url = f"{MAPILLARY_UPLOAD_ENDPOINT}/{self.session_key}"

        resp = self.user_session.get(url, timeout=REQUESTS_TIMEOUT)
        resp.raise_for_status()

        data = jsonify_response(resp)
        try:
            return data["offset"]
        except KeyError:
            raise HTTPContentError("Offset not found in the response", resp)

    @classmethod
    def chunkize_byte_stream(
        cls, stream: T.IO[bytes], chunk_size: int
    ) -> T.Generator[bytes, None, None]:
        """
        Chunkize a byte stream into chunks of the specified size.

        >>> list(UploadService.chunkize_byte_stream(io.BytesIO(b"foo"), 1))
        [b'f', b'o', b'o']
step 42
$ cd /repo && timeout 300 python - <<'PY'
output · 2 lines
[('seq1', UploadResult(result='1', error=None))]
['mly_tools_d8488a7aa2aba1459f81df70b8ffc60c.jpg']
step 43
$ cd /repo && timeout 300 python - <<'PY'
output · 5 lines
Error uploading sequence_seq1/test_exif.jpg at offset=None since begin_offset=None: ConnectionError: HTTPSConnectionPool(host='rupload.facebook.com', port=443): Max retries exceeded with url: /mapillary_public_uploads/mly_tools_d8488a7aa2aba1459f81df70b8ffc60c.jpg (Caused by NameResolutionError("HTTPSConnection(host='rupload.facebook.com', port=443): Failed to resolve 'rupload.facebook.com' ([Errno -3] Temporary failure in name resolution)"))
Error uploading sequence_seq1/test_exif.jpg at offset=None since begin_offset=None: ConnectionError: HTTPSConnectionPool(host='rupload.facebook.com', port=443): Max retries exceeded with url: /mapillary_public_uploads/mly_tools_d8488a7aa2aba1459f81df70b8ffc60c.jpg (Caused by NameResolutionError("HTTPSConnection(host='rupload.facebook.com', port=443): Failed to resolve 'rupload.facebook.com' ([Errno -3] Temporary failure in name resolution)"))
E
…[8087 chars omitted]…
g at offset=None since begin_offset=None: ConnectionError: HTTPSConnectionPool(host='rupload.facebook.com', port=443): Max retries exceeded with url: /mapillary_public_uploads/mly_tools_d8488a7aa2aba1459f81df70b8ffc60c.jpg (Caused by NameResolutionError("HTTPSConnection(host='rupload.facebook.com', port=443): Failed to resolve 'rupload.facebook.com' ([Errno -3] Temporary failure in name resolution)"))
step 44
$ cd /repo && grep -n "global \|^[A-Z_]* *[:=]\|lru_cache\|cache" mapillary_tools/exif_write.py mapillary_tools/serializer/description.py | head -30
output · 4 lines
mapillary_tools/exif_write.py:14:LOG = logging.getLogger(__name__)
mapillary_tools/serializer/description.py:47:_COORDINATES_PRECISION = 7
mapillary_tools/serializer/description.py:48:_ALTITUDE_PRECISION = 3
mapillary_tools/serializer/description.py:49:_ANGLE_PRECISION = 3
step 45
$ cd /repo && sed -n 1,80p mapillary_tools/exif_write.py
output · 43 lines
# pyre-ignore-all-errors[5, 21, 24]
from __future__ import annotations

import datetime
import io
import json
import logging
import math
from pathlib import Path

import piexif
31 more lines
LOG = logging.getLogger(__name__)


class ExifEdit:
    _filename_or_bytes: str | bytes

    def __init__(self, filename_or_bytes: Path | bytes) -> None:
        """Initialize the object"""
        if isinstance(filename_or_bytes, Path):
            # make sure filename is resolved to avoid to be interpretted as bytes in piexif
            # see https://github.com/hMatoba/Piexif/issues/124
            self._filename_or_bytes = str(filename_or_bytes.resolve())
        else:
            self._filename_or_bytes = filename_or_bytes
        self._ef: dict = piexif.load(self._filename_or_bytes)

    @staticmethod
    def decimal_to_dms(
        value: float, precision: int
    ) -> tuple[tuple[float, int], tuple[fl
…[1840 chars omitted]…
["Exif"]:
                    del self._ef["Exif"][piexif.ExifIFD.OffsetTimeOriginal]
        else:
            if piexif.ExifIFD.OffsetTimeOriginal in self._ef["Exif"]:
                del self._ef["Exif"][piexif.ExifIFD.OffsetTimeOriginal]

    def add_gps_datetime(self, dt: datetime.datetime) -> None:
        """Add GPSDateStamp and GPSTimeStamp."""
        dt = dt.astimezone(datetime.timezone.utc)
step 46
$ cd /repo && timeout 120 python - <<'PY'
output · 3 lines
Traceback (most recent call last):
  File "<stdin>", line 12, in <module>
NameError: name 'io' is not defined. Did you mean: 'i'? Or did you forget to import 'io'?
step 47
$ cd /repo && timeout 120 python - <<'PY'
output · 2 lines
sequential: ['2ba124e69c19a9a3b5bfa774e63a6aa9', 'a003718bdb608a1991c455d8624b7b07', 'd8488a7aa2aba1459f81df70b8ffc60c', '8160b394f88ed6fb6ffe55b727b868d2']
parallel:   ['2ba124e69c19a9a3b5bfa774e63a6aa9', 'a003718bdb608a1991c455d8624b7b07', 'd8488a7aa2aba1459f81df70b8ffc60c', '8160b394f88ed6fb6ffe55b727b868d2']
step 48
$ cd /repo && timeout 300 python - <<'PY' 2>&1 | tail -20
output · 3 lines
IMG /tmp/tmp_g1tufie/data/test_exif.jpg seq1 d8488a7aa2aba1459f81df70b8ffc60c cacheid 123951459025168
[('seq1', UploadResult(result='1', error=None))]
['mly_tools_d8488a7aa2aba1459f81df70b8ffc60c.jpg']
step 49
$ cd /repo && grep -n "def group_and_sort_images" -A 40 mapillary_tools/types.py; grep -n "def update_sequence_md5sum" -A 20 mapillary_tools/types.py
output · 30 lines
178:def group_and_sort_images(
179-    metadatas: T.Iterable[ImageMetadata],
180-) -> dict[str, list[ImageMetadata]]:
181-    # group metadatas by uuid
182-    sequences_by_uuid: dict[str, list[ImageMetadata]] = {}
183-    missing_sequence_uuid = str(uuid.uuid4())
184-    for metadata in metadatas:
185-        if metadata.MAPSequenceUUID is None:
186-            sequence_uuid = missing_sequence_uuid
187-        else:
188-            sequence_uuid = metadata.MAPSequenceUUID
189-        sequences_by_uuid.setdefault(sequence_uuid, []).append(metadata)
18 more lines
190-
191-    # deduplicate and sort metadatas per uuid
192-    sorted_sequences_by_uuid = {}
193-    for sequence_uuid, sequence in sequences_by_uuid.items():
194-        dedups = {metadata.filename.resolve(): metadata for metadata in sequence}
195-        sorted_sequences_by_uuid[sequence_uuid] = sorted(
196-            dedups.values(),
19
…[405 chars omitted]…
data.md5sum.encode("utf-8"))
208-    return md5.hexdigest()
202:def update_sequence_md5sum(sequence: T.Iterable[ImageMetadata]) -> str:
203-    md5 = hashlib.md5()
204-    for metadata in sequence:
205-        metadata.update_md5sum()
206-        assert isinstance(metadata.md5sum, str), "md5sum should be calculated"
207-        md5.update(metadata.md5sum.encode("utf-8"))
208-    return md5.hexdigest()
step 50
$ cd /repo && grep -n "class ImageMetadata\|def update_md5sum\|def sort_key\|md5sum" mapillary_tools/types.py | head -20; sed -n 100,178p mapillary_tools/types.py
output · 39 lines
33:class ImageMetadata(geo.Point):
36:    md5sum: str | None = None
51:    def update_md5sum(self, image_data: T.BinaryIO | None = None) -> None:
52:        if self.md5sum is None:
55:                    self.md5sum = utils.md5sum_fp(fp).hexdigest()
57:                self.md5sum = utils.md5sum_fp(image_data).hexdigest()
59:    def sort_key(self):
71:    md5sum: str | None = None
76:    def update_md5sum(self) -> None:
77:        if self.md5sum is None:
79:                self.md5sum = utils.md5sum_fp(fp).hexdigest()
202:def update_sequence_md5sum(sequence: T.Iterable[ImageMetadata]) -> str:
27 more lines
205:        metadata.update_md5sum()
206:        assert isinstance(metadata.md5sum, str), "md5sum should be calculated"
207:        md5.update(metadata.md5sum.encode("utf-8"))

    @classmethod
    @abc.abstractmethod
    def deserialize(cls, data: bytes) -> list[Metadata]:
        raise NotImplem
…[1562 chars omitted]…
tadata] = []

    for metadata in metadatas:
        if isinstance(metadata, ErrorMetadata):
            bad.append(metadata)
        else:
            good.append(metadata)

    return good, bad


def describe_error_metadata(
    exc: Exception, filename: Path, filetype: FileType
) -> ErrorMetadata:
    return ErrorMetadata(filename=filename, filetype=filetype, error=exc)


def group_and_sort_images(
step 51
$ cd /repo && sed -n 33,70p mapillary_tools/types.py; grep -rn "group_and_sort_images" mapillary_tools tests | grep -v types.py
output · 31 lines
class ImageMetadata(geo.Point):
    filename: Path
    # filetype should be always FileType.IMAGE
    md5sum: str | None = None
    width: int | None = None
    height: int | None = None
    filesize: int | None = None

    # Fields starting with MAP* will be written to the image EXIF
    MAPSequenceUUID: str | None = None
    MAPDeviceMake: str | None = None
    MAPDeviceModel: str | None = None
19 more lines
    MAPGPSAccuracyMeters: float | None = None
    MAPCameraUUID: str | None = None
    MAPOrientation: int | None = None
    MAPMetaTags: dict | None = None
    MAPFilename: str | None = None

    def update_md5sum(self, image_data: T.BinaryIO | None = None) -> None:
        if self.md5sum is None:
            if image_data is None:
                with self.filename.open("rb") as fp:
                    self.md5sum = utils.md5sum_fp(fp).hexdigest()
            else:
                self.md5s
…[539 chars omitted]…
images(image_metadatas)
mapillary_tools/uploader.py:354:        sequences = types.group_and_sort_images(metadatas)
mapillary_tools/uploader.py:371:        sequences = types.group_and_sort_images(image_metadatas)
mapillary_tools/uploader.py:461:        sequence_groups = types.group_and_sort_images(sequence)
mapillary_tools/uploader.py:549:        sequences = types.group_and_sort_images(image_metadatas)
step 52
$ cd /repo && grep -n "group_and_sort" -B5 -A40 tests/unit/test_sequence_processing.py | head -80
step 53
$ cd /repo && grep -rn "dedup\|group_and_sort" tests/ | head; pip download mapillary-tools==0.14.0 -d /tmp/mt --no-deps -q 2>&1 | tail -1
output · 3 lines
tests/unit/test_utils.py:53:def test_deduplicates():
tests/unit/test_utils.py:57:        x = utils.deduplicate_paths(
ERROR: No matching distribution found for mapillary-tools==0.14.0
step 54
$ cd /tmp && pip download mapillary-tools --no-deps -d /tmp/mt -q 2>&1 | tail -1; ls /tmp/mt 2>/dev/null
output · 1 lines
ERROR: No matching distribution found for mapillary-tools
step 55

My repro was flawed — images in a sequence are deduplicated by filename (existing test relies on that). Let me use distinct files:

$ cd /repo && python - <<'PY'
output · 5 lines
Traceback (most recent call last):
  File "/tmp/repro.py", line 67, in <module>
    assert uploads["count"] == 14, uploads["count"]
           ^^^^^^^^^^^^^^^^^^^^^^
AssertionError: 16
step 56
$ cd /repo && sed -i 's/assert uploads\["count"\] == 14, uploads\["count"\]/assert uploads["count"] == 16, uploads["count"]/; s/assert len(keys) == 14, keys/assert len(keys) == 14, keys/; s/len(u2.cached_image_uploader.cache.keys()) == 14/len(u2.cached_image_uploader.cache.keys()) == 14/' /tmp/repro.py && timeout 600 python /tmp/repro.py 2>&1 | tail -5
output · 5 lines
Traceback (most recent call last):
  File "/tmp/repro.py", line 77, in <module>
    assert uploads["count"] == 0, uploads["count"]
           ^^^^^^^^^^^^^^^^^^^^^
AssertionError: 2
step 57
$ cd /repo && sed -i 's/assert uploads\["count"\] == 0, uploads\["count"\]/assert uploads["count"] == 2, uploads["count"]  # only manifests (uuid keys are not cached)/' /tmp/repro.py && timeout 600 python /tmp/repro.py 2>&1 | tail -5
output · 2 lines
results: [('seq1', UploadResult(result='123', error=None)), ('seq2', UploadResult(result='123', error=None))]
ALL OK
step 58
$ cd /repo && timeout 1800 python -m pytest tests/unit/test_uploader.py tests/unit/test_persistent_cache.py tests/unit/test_sequence_processing.py -q 2>&1 | tail -5
output · 3 lines
........................................................................ [ 88%]
.........                                                                [100%]
81 passed in 41.31s
step 59
$ cd /repo && timeout 3000 python -m pytest tests/integration/test_upload.py tests/integration/test_process_and_upload.py -q 2>&1 | tail -5
output · 2 lines
......ss                                                                 [100%]
6 passed, 2 skipped in 13.50s
step 60
$ cd /repo && git diff | head -250
output · 46 lines
diff --git a/mapillary_tools/history.py b/mapillary_tools/history.py
index 8b7251e..27990e7 100644
--- a/mapillary_tools/history.py
+++ b/mapillary_tools/history.py
@@ -137,6 +137,20 @@ class PersistentCache:
             f"Cached file handle for {key} ({(time.perf_counter() - s) * 1000:.0f} ms)"
         )
 
+    def keys(self) -> list[str]:
+        """Return the (non-expired) keys stored in the cache"""
+
+        keys: list[str] = []
34 more lines
+
+        with self._lock:
+            with dbm.open(self._file, flag="c") as db:
+                for key in db.keys():
+                    if isinstance(key, bytes):
+                        key = key.decode("utf-8", errors="replace")
+                    keys.append(T.cast(str, key))
+
+        return keys
+
     def clear_expired(self) -> list[str]:
         s = time.perf_counter()
 
diff --git a/mapillary_tools/uploader.py b/mapillary_tools/up
…[7700 chars omitted]…
ile_handle(self, key: str) -> str | None:
-        if self.cache is None:
-            return None
-
-        if _is_uuid(key):
-            return None
-
-        return self.cache.get(key)
-
-    def _set_file_handle_cache(self, key: str, value: str) -> None:
-        if self.cache is None:
-            return
-
-        if _is_uuid(key):
-            return
-
-        self.cache.set(key, value)
-
-
step 61
$ cd /repo && git diff | sed -n 250,400p
output · 4 lines
-
 class Uploader:
     def __init__(
         self,
step 62
$ cd /repo && grep -n "config\." mapillary_tools/uploader.py | head; timeout 3000 python -m pytest tests/unit -q 2>&1 | tail -3
output · 4 lines
56:    user_items: config.UserItem
........................................................................ [ 91%]
...................                                                      [100%]
225 passed, 10 skipped in 44.43s
step 63
$ cd /repo && timeout 3000 python -m pytest tests/integration -q 2>&1 | tail -3
output · 2 lines
sss....s.s.s.s...............ssssss....ss....sssssss                     [100%]
30 passed, 22 skipped in 43.25s
step 64

All tests pass and the reproducer confirms every requirement in the issue. Summary of changes:

$ done: {"notes": "Added UploadOptions.upload_cache_path, new CachedUploader with a single shared PersistentCache created once in ImageSequenceUploader (shared across parallel workers), PersistentCache.keys(), and task_done fix on early worker break. Reproducer + all unit/integration tests pass."}
output · 1 lines
submitted.