mapillary-mapillary-tools-779
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 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]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]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)
$ 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($ 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($ 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.pyoutput · 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",$ 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 = None21 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 = ($ 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$ cd /repo && python - <<'PY'
output · 2 lines
................................................................. [100%] 65 passed in 40.49s
$ 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 = ($ 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$ 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()$ 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_")$ 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]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:
$ 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()
$ 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"),$ 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.pyoutput · 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,
),$ 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 -10output · 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_handles31 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($ python - <<'PY'
$ 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 = ($ 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$ 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 = {$ python - <<'PY'
output · 2 lines
................................................................. [100%] 65 passed in 41.22s
$ 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
$ 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]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:
$ 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(
$ 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"))$ 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.pyoutput · 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",$ 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)$ 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]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(+)
$ 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))
)
)$ 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$ 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$ 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")$ 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 = ($ 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(-)
$ 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():
break24 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$ cd /repo && python - <<'PY'
output · 3 lines
mapillary_tools/history.py | 5 ++++ mapillary_tools/uploader.py | 63 ++++++++++++++++++++++++++------------------- 2 files changed, 41 insertions(+), 27 deletions(-)
$ 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
$ 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]]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)
$ 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($ 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:$ 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())}"$ 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:$ sed -n '850,1030p' mapillary_tools/uploader.py; sed -n '140,165p' mapillary_tools/constants.py; grep -R "ImageSequenceUploader(" -n mapillary_tools testsoutput · 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_session30 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($ 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:$ 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(+)
$ 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__(
$ python - <<'PY'
output · 2 lines
................................................................. [100%] 65 passed in 41.69s
$ 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]]$ 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
$ 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(-)
$ 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]$ 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
$ 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:
$ 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),$ 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:$ 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:$ 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",$ 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)$ 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$ 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,
)$ 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.BLACKVUE39 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()$ git log --oneline | head; git diff HEAD~1 --stat 2>/dev/null | tail -5
output · 1 lines
115aa82 base
$ 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-$ 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_")$ 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: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(-)
$ cd /repo && python - <<'PY'
output · 2 lines
767: def upload( 862: def _get_cached_file_handle(self, key: str) -> str | None:
$ 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)$ cd /repo && python - <<'PY'
output · 1 lines
ok
$ 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:$ 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"$ 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'$ 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'$ 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"],
)$ 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'
$ 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)$ 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$ 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]):
$ 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$ 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:
$ 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:
$ 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$ 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__($ 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$ 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$ 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
$ 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']$ cd /repo && python - <<'PY'
output · 3 lines
58.5 2ba124e69c19a9a3b5bfa774e63a6aa9 58.51 a003718bdb608a1991c455d8624b7b07 58.52 d8488a7aa2aba1459f81df70b8ffc60c
$ 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)$ 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']
$ cd /repo && timeout 300 python - <<'PY'
output · 2 lines
[('seq1', UploadResult(result=None, error=AttributeError("'C' object has no attribute 'session_key'")))]
[]$ 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']$ cd /repo && timeout 300 python - <<'PY'
output · 2 lines
[('seq1', UploadResult(result='1', error=None))]
['mly_tools_d8488a7aa2aba1459f81df70b8ffc60c.jpg']$ 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)"))$ 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
$ 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)$ 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'?
$ cd /repo && timeout 120 python - <<'PY'
output · 2 lines
sequential: ['2ba124e69c19a9a3b5bfa774e63a6aa9', 'a003718bdb608a1991c455d8624b7b07', 'd8488a7aa2aba1459f81df70b8ffc60c', '8160b394f88ed6fb6ffe55b727b868d2'] parallel: ['2ba124e69c19a9a3b5bfa774e63a6aa9', 'a003718bdb608a1991c455d8624b7b07', 'd8488a7aa2aba1459f81df70b8ffc60c', '8160b394f88ed6fb6ffe55b727b868d2']
$ 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']$ 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()$ 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($ 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 = None19 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)$ cd /repo && grep -n "group_and_sort" -B5 -A40 tests/unit/test_sequence_processing.py | head -80
$ 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
$ 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
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$ 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$ 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$ 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
$ 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
$ 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)
-
-$ cd /repo && git diff | sed -n 250,400p
output · 4 lines
-
class Uploader:
def __init__(
self,$ 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
$ 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
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.