From ee5a9f07f3996e5f798b9d59433c79f38ce45fe6 Mon Sep 17 00:00:00 2001 From: Soyoung Kim Date: Mon, 5 Oct 2026 14:57:06 +0200 Subject: [PATCH 1/2] feat: add sha media path --- synapse/media/filepath.py | 33 +++++++++++++++++++++++++++++++++ tests/media/test_filepath.py | 20 ++++++++++++++++++++ 2 files changed, 53 insertions(+) diff --git a/synapse/media/filepath.py b/synapse/media/filepath.py index df637f3be3..ab3cef1097 100644 --- a/synapse/media/filepath.py +++ b/synapse/media/filepath.py @@ -134,6 +134,16 @@ def _wrapped( } +def _content_hash_rel(directory: str, sha256: str) -> str: + """Relative path of a content-addressed blob inside `directory`.""" + return os.path.join( + directory, + _validate_path_component(sha256[0:2]), + _validate_path_component(sha256[2:4]), + _validate_path_component(sha256[4:]), + ) + + def _validate_path_component(name: str) -> str: """Checks that the given string can be safely used as a path component @@ -184,6 +194,17 @@ def local_media_filepath_rel(self, media_id: str) -> str: local_media_filepath = _wrap_in_base_path(local_media_filepath_rel) + @_wrap_with_jail_check(relative=True) + def local_media_blob_rel(self, sha256: str) -> str: + """Relative path of the local content-hash blob. + + Downloads keep using `local_media_filepath`. This path is an extra + hardlink so identical local uploads share one inode. + """ + return _content_hash_rel("local_blobs", sha256) + + local_media_blob_filepath = _wrap_in_base_path(local_media_blob_rel) + @_wrap_with_jail_check(relative=True) def local_media_thumbnail_rel( self, media_id: str, width: int, height: int, content_type: str, method: str @@ -230,6 +251,18 @@ def remote_media_filepath_rel(self, server_name: str, file_id: str) -> str: remote_media_filepath = _wrap_in_base_path(remote_media_filepath_rel) + @_wrap_with_jail_check(relative=True) + def remote_media_blob_rel(self, sha256: str) -> str: + """Relative path of the remote content-hash blob. + + Remote downloads keep using `remote_media_filepath`. This path is a + separate tree from `local_blobs`, so local and remote media never + share an inode. + """ + return _content_hash_rel("remote_blobs", sha256) + + remote_media_blob_filepath = _wrap_in_base_path(remote_media_blob_rel) + @_wrap_with_jail_check(relative=True) def remote_media_thumbnail_rel( self, diff --git a/tests/media/test_filepath.py b/tests/media/test_filepath.py index cd21b369b4..d44ba8b815 100644 --- a/tests/media/test_filepath.py +++ b/tests/media/test_filepath.py @@ -66,6 +66,26 @@ def test_local_media_thumbnail_dir(self) -> None: "/media_store/local_thumbnails/Ge/rZ/NDnDZVjsOtardLuwfIBg", ) + def test_content_hash_blob_paths(self) -> None: + """Local and remote blobs are separate trees, sharded like media ids.""" + sha256 = "ebf4f635a17d10d6eb46ba680b70142419aa3220f228001a036d311a22ee9d2a" + self.assertEqual( + self.filepaths.local_media_blob_rel(sha256), + "local_blobs/eb/f4/f635a17d10d6eb46ba680b70142419aa3220f228001a036d311a22ee9d2a", + ) + self.assertEqual( + self.filepaths.local_media_blob_filepath(sha256), + "/media_store/local_blobs/eb/f4/f635a17d10d6eb46ba680b70142419aa3220f228001a036d311a22ee9d2a", + ) + self.assertEqual( + self.filepaths.remote_media_blob_rel(sha256), + "remote_blobs/eb/f4/f635a17d10d6eb46ba680b70142419aa3220f228001a036d311a22ee9d2a", + ) + self.assertEqual( + self.filepaths.remote_media_blob_filepath(sha256), + "/media_store/remote_blobs/eb/f4/f635a17d10d6eb46ba680b70142419aa3220f228001a036d311a22ee9d2a", + ) + def test_remote_media_filepath(self) -> None: """Test remote media paths""" self.assertEqual( From 7dc74c3ef2500a6044e640f6e0cbb66b5680b7da Mon Sep 17 00:00:00 2001 From: Soyoung Kim Date: Tue, 6 Oct 2026 12:08:25 +0200 Subject: [PATCH 2/2] feat: save file in sha256 path blob --- synapse/media/media_repository.py | 4 + synapse/media/media_storage.py | 67 ++++++++ tests/media/test_media_storage.py | 259 +++++++++++++++++++++++++++++- 3 files changed, 329 insertions(+), 1 deletion(-) diff --git a/synapse/media/media_repository.py b/synapse/media/media_repository.py index 31bb401e65..49881c1c00 100644 --- a/synapse/media/media_repository.py +++ b/synapse/media/media_repository.py @@ -439,6 +439,8 @@ async def create_or_update_content( quarantined_by="system" if should_quarantine else None, ) + await self.media_storage.link_id_path_to_blob(fname, sha256, remote=False) + try: await self._generate_thumbnails(None, media_id, media_id, media_type) except Exception as e: @@ -880,6 +882,8 @@ async def _store_remote_media_with_cleanup( # Re-raise so the caller can handle it raise + await self.media_storage.link_id_path_to_blob(fname, sha256, remote=True) + async def _download_remote_file( self, server_name: str, diff --git a/synapse/media/media_storage.py b/synapse/media/media_storage.py index 4c942f07f8..ee40900c85 100644 --- a/synapse/media/media_storage.py +++ b/synapse/media/media_storage.py @@ -332,6 +332,33 @@ async def fetch_media(self, file_info: FileInfo) -> Responder | None: return None + def _blob_path(self, sha256: str, remote: bool) -> str | None: + try: + if remote: + return self.filepaths.remote_media_blob_filepath(sha256) + return self.filepaths.local_media_blob_filepath(sha256) + except ValueError: + logger.warning("Ignoring invalid media sha256 %r", sha256) + return None + + async def link_id_path_to_blob( + self, id_path: str, sha256: str, remote: bool + ) -> None: + """Hardlink `id_path` to the content-hash blob depending on the remote flag. + Args: + id_path: Absolute path of the media id file that was just stored. + sha256: Hex digest of that file. + remote: True when the file is remote media. + """ + if not sha256 or self.local_media_directory is None: + return + + blob_path = self._blob_path(sha256, remote) + if blob_path is None: + return + + _link_id_path_to_blob(id_path, blob_path) + @trace @contextlib.asynccontextmanager async def ensure_media_is_in_local_cache( @@ -464,6 +491,46 @@ def _file_info_to_path(self, file_info: FileInfo) -> str: return self.filepaths.local_media_filepath_rel(file_info.file_id) +def _link_id_path_to_blob(id_path: str, blob_path: str) -> None: + """Point `id_path` at the inode named by `blob_path`, creating the blob if needed. + + A failed link leaves `id_path` in place. + """ + if not os.path.isfile(id_path): + return + + tmp_link = id_path + ".dedup-link" + try: + os.makedirs(os.path.dirname(blob_path), exist_ok=True) + if not os.path.exists(blob_path): + try: + os.link(id_path, blob_path) + return + except FileExistsError: + pass + if os.path.samefile(id_path, blob_path): + return + + try: + os.remove(tmp_link) + except FileNotFoundError: + pass + # Link the existing blob to a temporary name first, replacing the id path only happens after that link succeeds. + os.link(blob_path, tmp_link) + os.replace(tmp_link, id_path) + except OSError: + logger.warning( + "Failed to hardlink media file %s to %s", + id_path, + blob_path, + exc_info=True, + ) + try: + os.remove(tmp_link) + except OSError: + pass + + @trace def _write_file_synchronously(source: IO, dest: IO) -> None: """Write `source` to the file like `dest` synchronously. Should be called diff --git a/tests/media/test_media_storage.py b/tests/media/test_media_storage.py index da5fa8a140..d8069425b4 100644 --- a/tests/media/test_media_storage.py +++ b/tests/media/test_media_storage.py @@ -18,6 +18,7 @@ # [This file includes modifications made by New Vector Limited] # # +import hashlib import os import shutil import tempfile @@ -47,7 +48,12 @@ from synapse.logging.context import make_deferred_yieldable from synapse.media._base import FileInfo, ThumbnailInfo from synapse.media.filepath import MediaFilePaths -from synapse.media.media_storage import MediaStorage, ReadableFileWrapper +from synapse.media.media_storage import ( + FileResponder, + MediaStorage, + ReadableFileWrapper, + _link_id_path_to_blob, +) from synapse.media.storage_provider import ( FileStorageProviderBackend, StorageProviderWrapper, @@ -1374,6 +1380,257 @@ async def _send_request(*args: Any, **kwargs: Any) -> IResponse: ) +class MediaHashDedupTestCase(unittest.HomeserverTestCase): + servlets = [ + admin.register_servlets, + login.register_servlets, + media.register_servlets, + ] + + def prepare(self, reactor: MemoryReactor, clock: Clock, hs: HomeServer) -> None: + self.user = self.register_user("user", "pass") + self.tok = self.login("user", "pass") + self.store = hs.get_datastores().main + self.media_repo = hs.get_media_repository() + self.client = hs.get_federation_http_client() + + def create_resource_dict(self) -> dict[str, Resource]: + resources = super().create_resource_dict() + resources["/_matrix/media"] = self.hs.get_media_repository_resource() + return resources + + def _upload(self, data: bytes) -> str: + media = self.helper.upload_media(data, tok=self.tok, expect_code=200) + return media["content_uri"].rsplit("/", 1)[-1] + + def _local_path(self, media_id: str) -> str: + return self.media_repo.filepaths.local_media_filepath(media_id) + + def _blob_path(self, sha256: str) -> str: + return self.media_repo.filepaths.local_media_blob_filepath(sha256) + + def test_link_id_path_to_blob(self) -> None: + """ + Ensure that `_link_id_path_to_blob` correctly links the media ID path to the + blob (via hardlink) and that the filesystem reflects the expected state. + + This test uploads a file, then links its media ID location to the blob location + by content hash, checks that they are the same inode. + """ + media_id = self._upload(SMALL_PNG) + path = self._local_path(media_id) + blob = self._blob_path(SMALL_PNG_SHA256) + + # Create a hard link from the media ID path to the blob file path. + _link_id_path_to_blob(path, blob) + + # Assert that both the media ID file path and the blob path point to the same inode. + self.assertTrue(os.path.samefile(path, blob)) + + # Check that the blob now has exactly two hard links: one for the blob location and one for the original media ID path. + self.assertEqual(os.stat(blob).st_nlink, 2) + + def test_duplicate_local_uploads_share_one_blob(self) -> None: + """ + Test that multiple uploads of identical file content share a single blob. + + This test uploads the same file twice, then checks that the two media ID paths + and the blob path all point to the same file on disk. + """ + + # Upload the same file twice, get each media ID. + media_id_a = self._upload(SMALL_PNG) + media_id_b = self._upload(SMALL_PNG) + + # Compute the absolute paths of each upload and the shared blob path. + path_a = self._local_path(media_id_a) + path_b = self._local_path(media_id_b) + blob = self.media_repo.filepaths.local_media_blob_filepath(SMALL_PNG_SHA256) + + # Ensure both media IDs and the blob path all point to the same file on disk. + self.assertTrue(os.path.samefile(path_a, blob)) + self.assertTrue(os.path.samefile(path_b, blob)) + + # The hard link count should be three: blob, media_a, and media_b. + self.assertEqual(os.stat(blob).st_nlink, 3) + + # File contents for both media are identical to expected bytes. + with open(path_a, "rb") as f: + self.assertEqual(f.read(), SMALL_PNG) + with open(path_b, "rb") as f: + self.assertEqual(f.read(), SMALL_PNG) + + # Cleanup by removing both uploaded media entries. + self.get_success( + self.media_repo.delete_local_media_ids([media_id_a, media_id_b]) + ) + + def test_legacy_file_without_blob_still_downloads_and_deletes(self) -> None: + """ + Test that legacy local media files that predate the deduplication scheme: + - A media object without a sha256 + - A media object with a sha256 but with the blob missing + are still readable and deletable. + """ + user = UserID.from_string(self.user) + + # 1. Legacy file case without sha256 + legacy_body = b"legacy-file" + legacy_id = "LegacyMediaIdNoBlob0001" + # Store the media metadata in the database, omitting sha256 so that dedup logic does not run. + self.get_success( + self.store.store_local_media( + media_id=legacy_id, + media_type="text/plain", + time_now_ms=self.hs.get_clock().time_msec(), + upload_name=None, + media_length=len(legacy_body), + user_id=user, + sha256=None, + ) + ) + # Write the actual file contents to the file location (old uploads) + legacy_path = self._local_path(legacy_id) + os.makedirs(os.path.dirname(legacy_path), exist_ok=True) + with open(legacy_path, "wb") as f: + f.write(legacy_body) + + # 2. Legacy file case with sha256 but no blob file + hashed_body = b"legacy-hashed" + hashed_id = "LegacyMediaIdHashed0001" + digest = hashlib.sha256(hashed_body).hexdigest() + # Store the media with a sha256, but do not create the blob file. + self.get_success( + self.store.store_local_media( + media_id=hashed_id, + media_type="text/plain", + time_now_ms=self.hs.get_clock().time_msec(), + upload_name=None, + media_length=len(hashed_body), + user_id=user, + sha256=digest, + ) + ) + # Create only the file at the media_id location, not the blob location. + hashed_path = self._local_path(hashed_id) + os.makedirs(os.path.dirname(hashed_path), exist_ok=True) + with open(hashed_path, "wb") as f: + f.write(hashed_body) + blob = self.media_repo.filepaths.local_media_blob_filepath(digest) + self.assertFalse( + os.path.exists(blob), + "No deduplicated blob file should exist for this sha256", + ) + + # 3. Attempt to download the files + # For both cases, verify that fetch_media works and the content matches. + for media_id, body in ((legacy_id, legacy_body), (hashed_id, hashed_body)): + responder = self.get_success( + self.media_repo.media_storage.fetch_media(FileInfo(None, media_id)) + ) + self.assertIsInstance(responder, FileResponder) + assert isinstance(responder, FileResponder) + try: + # Check contents. + self.assertEqual(responder.open_file.read(), body) + finally: + responder.__exit__(None, None, None) + + # 4. Attempt to delete both files + # They should delete cleanly, without regard to missing blobs. + self.get_success(self.media_repo.delete_local_media_ids([legacy_id, hashed_id])) + self.assertFalse(os.path.exists(legacy_path)) + self.assertFalse(os.path.exists(hashed_path)) + self.assertFalse(os.path.exists(blob)) + + @override_config({"enable_authenticated_media": False}) + @patch( + "synapse.http.matrixfederationclient.read_body_with_max_size", + read_body, + ) + def test_local_and_remote_blobs_stay_separate(self) -> None: + """ + Test that local and remote media blobs are stored as separate inodes. + """ + + # Patch the federation file fetching with a successful response + async def _send_request(*args: Any, **kwargs: Any) -> IResponse: + resp = MagicMock(spec=IResponse) + resp.code = 200 + resp.length = len(SMALL_PNG) + resp.headers = Headers({"Content-Type": ["image/png"]}) + resp.phrase = b"OK" + return resp + + self.client._send_request = _send_request # type: ignore[method-assign] + + # Upload to the local media repository and get its local media_id. + local_id = self._upload(SMALL_PNG) + + # Download the same bytes from two different "remote" server to cache them remotely. + # This triggers federation download and caching logic, but for different servers. + for server_name, media_id in (("remote.org", "abc"), ("other.org", "def")): + channel = self.make_request( + "GET", + f"/_matrix/media/v3/download/{server_name}/{media_id}", + shorthand=False, + access_token=self.tok, + ) + self.assertEqual(channel.code, 200) + + # Retrieve metadata for each remote cached copy. + remote_a = self.get_success( + self.store.get_cached_remote_media("remote.org", "abc") + ) + remote_b = self.get_success( + self.store.get_cached_remote_media("other.org", "def") + ) + assert remote_a is not None + assert remote_b is not None + + # Compute absolute paths for the cached files and blobs. + remote_path_a = self.media_repo.filepaths.remote_media_filepath( + "remote.org", remote_a.filesystem_id + ) + remote_path_b = self.media_repo.filepaths.remote_media_filepath( + "other.org", remote_b.filesystem_id + ) + remote_blob = self.media_repo.filepaths.remote_media_blob_filepath( + SMALL_PNG_SHA256 + ) + local_blob = self.media_repo.filepaths.local_media_blob_filepath( + SMALL_PNG_SHA256 + ) + + # Assert the uploaded local file is the same inode as the local blob. + self.assertTrue(os.path.samefile(self._local_path(local_id), local_blob)) + + # Assert that both remote cached files are the same inode as the remote blob. + self.assertTrue(os.path.samefile(remote_path_a, remote_blob)) + self.assertTrue(os.path.samefile(remote_path_b, remote_blob)) + + # Assert the local blob is a different inode to the remote blob. + self.assertNotEqual(os.stat(local_blob).st_ino, os.stat(remote_blob).st_ino) + + # Update access timestamps so remote deletions target only one of the remote servers. + self.get_success( + self.store.update_cached_last_access_time([], [("remote.org", "abc")], 1) + ) + self.get_success( + self.store.update_cached_last_access_time( + [], [("other.org", "def")], 5_000_000 + ) + ) + + # Delete old remote media + self.get_success(self.media_repo.delete_old_remote_media(1000)) + self.assertFalse(os.path.exists(remote_path_a)) + # Ensure that the remote blob still exists with the other remote media. + self.assertTrue(os.path.samefile(remote_path_b, remote_blob)) + # Local blob should still exist. + self.assertTrue(os.path.exists(local_blob)) + + class MediaRepoSizeModuleCallbackTestCase(unittest.HomeserverTestCase): servlets = [ login.register_servlets,