diff --git a/src/borg/repoobj.py b/src/borg/repoobj.py index 6f5960582f..5236f9146b 100644 --- a/src/borg/repoobj.py +++ b/src/borg/repoobj.py @@ -292,5 +292,32 @@ def validate(chunk_id, obj): return validate +def whole_object_authenticator(repo_objs): + """Return authenticate(chunk_id, obj): True if obj is the whole repo object with id chunk_id. + + obj is an object's header, metadata slot and data slot. Parsing it checks that the header's sizes + add up to len(obj) and verifies the tag of each slot. A slot's tag is computed over the slot and + over header_aad + slot_tag + chunk_id as AAD (additional authenticated data: bytes the tag covers + without being part of the ciphertext), see OBJ_VERSION_HEADER_AAD. Only the tags are verified, not + that the plaintext hashes to chunk_id. + + Raises Error for an "authenticated-*" key with the authenticated_no_key workaround, which skips + the tag verification. + """ + from .crypto.key import MACKeyBase # crypto.key imports this module + + if AUTHENTICATED_NO_KEY and isinstance(repo_objs.key, MACKeyBase): + raise Error("Objects can not be authenticated with BORG_WORKAROUNDS=authenticated_no_key.") + + def authenticate(chunk_id, obj): + try: + repo_objs.parse(chunk_id, obj, decompress=False, want_compressed=True, ro_type=ROBJ_DONTCARE) + except (IntegrityErrorBase, msgpack.UnpackException): + return False + return True + + return authenticate + + # Backward compatibility: RepoObj1 has moved to borg.legacy.repoobj from .legacy.repoobj import RepoObj1 # noqa: F401 diff --git a/src/borg/repository.py b/src/borg/repository.py index 7e95cdbf9a..2745a9d9f1 100644 --- a/src/borg/repository.py +++ b/src/borg/repository.py @@ -14,11 +14,12 @@ from borgstore.backends.rest import REST, ssh_cmd from borgstore.store import ObjectNotFound as StoreObjectNotFound, ReadRangeError from borgstore.backends.errors import BackendError as StoreBackendError +from borgstore.backends.errors import BackendConnectionError as StoreBackendConnectionError from borgstore.backends.errors import BackendDoesNotExist as StoreBackendDoesNotExist from borgstore.backends.errors import BackendAlreadyExists as StoreBackendAlreadyExists from .constants import * # NOQA -from .hashindex import ChunkIndex +from .hashindex import ChunkIndex, ChunkIndexEntry from .helpers import Error, ErrorWithTraceback, IntegrityError from .helpers import Location from .helpers import bin_to_hex, hex_to_bin @@ -472,7 +473,7 @@ def _validation_problem(self, hdr, offset, buf, buf_offset, validate): end = start + size obj = buf[start:end] if end <= len(buf) else self.read(offset, size) if not validate(hdr.chunk_id, obj): - return "object does not authenticate" + return "object header or metadata does not authenticate" return None def _find_header(self, offset, pack_size, validate): @@ -685,6 +686,21 @@ def remove_missing_pack_entries(chunks, missing_pack_ids): return len(stale_ids) +# Repository.salvage_pack outcomes. +SALVAGE_INTACT = "intact" # the pack's store hash matches its name +SALVAGE_DONE = "salvaged" # the pack was replaced by one holding only its authenticated objects +SALVAGE_NOTHING_AUTHENTICATES = "nothing authenticates" # no object in the pack authenticates +SALVAGE_READS_DIFFER = "reads differ" # two reads of the pack disagree +SALVAGE_READ_ERROR = "read error" # reading the pack failed + +# status: one of the SALVAGE_* outcomes. +# new_pack_id: id of the replacement pack (SALVAGE_DONE), else None. +# kept: (chunk_id, obj_offset, obj_size) of each object in the replacement pack, else []. +# dropped_bytes: number of bytes of the old pack left out of the replacement pack, else 0. +# removed_ids: chunk ids whose chunk index entries were removed, else []. +SalvageResult = namedtuple("SalvageResult", "status new_pack_id kept dropped_bytes removed_ids") + + class PackTracker: """Pack verification results, mapping pack_id -> (timestamp, result). @@ -961,6 +977,8 @@ def __init__( if cache_size: ns_config["packs/"]["size"] = int(cache_size) cache_url = cache_dir.as_uri() + # True if packs are cached locally (BORG_STORE_CACHE): store.load() of a pack may return the cached copy. + self.uses_pack_store_cache = cache_url is not None propagate_rsh() # borgstore shall use the same remote shell command as borg @@ -2363,6 +2381,131 @@ def transform_pack(self, pack_id, ids, transform, *, validate, chunks=None, befo self.store_delete(pack_key) return new_pack_id, len(pack_data) + def salvage_pack(self, pack_id, *, validate, authenticate, chunks=None, before_old_pack_delete=None): + """Replace pack by a pack holding only the objects in it that authenticate. + + A pack is named by the store hash (STORE_HASH_NAME) of its content. + + validate: validate(chunk_id, obj) -> bool, True if obj (an object's header and metadata slot) + is the repo object with id chunk_id, see repoobj.object_validator. The pack is walked with + PackReader.iter_headers(validate). + authenticate: authenticate(chunk_id, obj) -> bool, True if obj (a whole object: header, metadata + slot and data slot) is the repo object with id chunk_id, see repoobj.whole_object_authenticator. + It must verify the tags of both slots. + chunks: the ChunkIndex to update. Default: self.chunks. + before_old_pack_delete: callable without arguments, called once after the replacement pack is + stored and before chunks is updated and the old pack is deleted; use it to invalidate stored + chunk indexes for crash safety (see #9748). Not called if the old pack is not deleted, see + SALVAGE_DONE. + + Every object the walk yields and authenticate accepts is kept, whether the chunk index lists + it or not. Everything else is dropped: objects authenticate rejects, byte ranges the walk + skips, and trailing bytes too few for an object header. The replacement pack is the kept + objects' bytes in their old order. + + Returns a SalvageResult. The store and the chunk index change only for SALVAGE_DONE: + - SALVAGE_READS_DIFFER: the loaded bytes hash to the pack's name although the store hash did + not, or a second load of the pack differs from the first, e.g. due to corruption in memory + or in transfer. Objects are dropped only if both loads return the same bytes. + - SALVAGE_READ_ERROR: reading the pack raised OSError, StoreBackendConnectionError or + ReadRangeError. All reads happen before the first store change. + - SALVAGE_DONE: the replacement pack is stored, before_old_pack_delete is called, chunks is + updated, then the old pack is deleted. If the replacement pack has the old pack's name + (the kept bytes are the undamaged pack), storing it overwrites the old pack, and + before_old_pack_delete and the delete are skipped. An exception in one of these steps leaves + the steps before it done. + + The chunk index update: an entry of this pack is pointed at the kept object at its offset, + or else at a kept copy of the same chunk id, or else removed. A kept object whose chunk id + has no entry gets one with flags F_USED and size 0 (the plaintext size is unknown). Entries + of other packs and F_PENDING entries stay as they are. + + Raises Error before any store access if uses_pack_store_cache is set (then two loads can + return the same cached copy). Raises Repository.PermissionDenied unless the repo permissions + allow compaction (see assert_writable), and StoreObjectNotFound if the pack is missing. Other + store backend errors propagate. + + Updates the in-memory chunk index only; the caller holds the exclusive lock and writes the + index back to the store afterwards. + """ + if self.uses_pack_store_cache: + raise Error("Pack salvage refused: with BORG_STORE_CACHE, pack reads may return the cached copy.") + self._lock_refresh() + if chunks is None: + chunks = self.chunks + self.assert_writable() + pack_hex = bin_to_hex(pack_id) + pack_key = "packs/" + pack_hex + + def unchanged(status): + return SalvageResult(status, None, [], 0, []) + + try: + if self.store.hash(pack_key, algorithm=STORE_HASH_NAME) == pack_hex: + return unchanged(SALVAGE_INTACT) + pack_contents = self.store.load(pack_key) + if store_hash(pack_contents).digest() == pack_id: + return unchanged(SALVAGE_READS_DIFFER) + reader = PackReader(pack_id=pack_id, pack_contents=pack_contents) + kept_old = [] # (chunk_id, old offset, size) of each kept object, offset-ordered + for chunk_id, offset, size in reader.iter_headers(validate=validate): + if authenticate(chunk_id, reader.read(offset, size)): + kept_old.append((chunk_id, offset, size)) + if not kept_old: + return unchanged(SALVAGE_NOTHING_AUTHENTICATES) + if self.store.load(pack_key) != pack_contents: + return unchanged(SALVAGE_READS_DIFFER) + except (OSError, StoreBackendConnectionError, ReadRangeError) as exc: + logger.warning(f"pack {pack_hex}: {exc}, not salvaging it.") + return unchanged(SALVAGE_READ_ERROR) + + new_pack_data = b"".join(pack_contents[offset : offset + size] for _, offset, size in kept_old) + # equal to pack_id if the kept bytes are the undamaged pack, e.g. when the damage is appended bytes. + new_pack_id = store_hash(new_pack_data).digest() + kept = [] # (chunk_id, new offset, size) + new_offset = 0 + for chunk_id, _, size in kept_old: + kept.append((chunk_id, new_offset, size)) + new_offset += size + dropped_bytes = len(pack_contents) - len(new_pack_data) + + self.store_store("packs/" + bin_to_hex(new_pack_id), new_pack_data) + if before_old_pack_delete is not None and new_pack_id != pack_id: + before_old_pack_delete() + + new_by_old_offset = {old[1]: new for old, new in zip(kept_old, kept)} + new_by_id = {} # chunk_id -> its first kept copy + for new in kept: + new_by_id.setdefault(new[0], new) + # collect first: the index must not be mutated while iterating it. + listed = [ + (chunk_id, entry.obj_offset) + for chunk_id, entry in chunks.iteritems() + if entry.pack_id == pack_id and not (entry.flags & ChunkIndex.F_PENDING) + ] + new_locations = [] + removed_ids = [] + for chunk_id, old_offset in listed: + new = new_by_old_offset.get(old_offset) + if new is None or new[0] != chunk_id: + new = new_by_id.get(chunk_id) + if new is None: + del chunks[chunk_id] + removed_ids.append(chunk_id) + else: + new_locations.append((chunk_id, new_pack_id, new[1], new[2])) + chunks.update_pack_info(new_locations) + for chunk_id, (_, offset, size) in new_by_id.items(): + if chunk_id not in chunks: + chunks[chunk_id] = ChunkIndexEntry( + flags=ChunkIndex.F_USED, size=0, pack_id=new_pack_id, obj_offset=offset, obj_size=size + ) + + if new_pack_id != pack_id: # else storing the replacement pack overwrote the old one + self.store_delete(pack_key) + self._pack_cache.pop(pack_id, None) + return SalvageResult(SALVAGE_DONE, new_pack_id, kept, dropped_bytes, removed_ids) + def acquire_lock(self): """Lock the repository (as requested by open()), loading its key first if no key was set yet. diff --git a/src/borg/testsuite/archiver/check_cmd_test.py b/src/borg/testsuite/archiver/check_cmd_test.py index 3f5417dc59..90a03f873f 100644 --- a/src/borg/testsuite/archiver/check_cmd_test.py +++ b/src/borg/testsuite/archiver/check_cmd_test.py @@ -1101,7 +1101,8 @@ def test_repair_resyncs_pack_with_corrupt_object_header(archivers, request, dama repository.store_store(key, pack) output = cmd(archiver, "check", "--repair", "--debug", exit_code=0) - problem = {"magic": "no object header", "data_size": "object does not authenticate"}[damaged_field] + problems = {"magic": "no object header", "data_size": "object header or metadata does not authenticate"} + problem = problems[damaged_field] assert f"{problem} at offset {damaged_offset}" in output assert f"continuing at the object at offset {next_offset}" in output # the rebuild resumed at the next object # the resync dropped an object, so the summary reports a problem. diff --git a/src/borg/testsuite/repoobj_test.py b/src/borg/testsuite/repoobj_test.py index b05363ba9d..7ca95f0658 100644 --- a/src/borg/testsuite/repoobj_test.py +++ b/src/borg/testsuite/repoobj_test.py @@ -13,6 +13,7 @@ OBJ_VERSION, RepoObj, object_validator, + whole_object_authenticator, ) from ..legacy.repoobj import RepoObj1 from ..compress import LZ4 @@ -520,3 +521,64 @@ def test_object_validator_accepts_every_compression(compression): repo_objs.compressor = CompressionSpec(compression).compressor chunk_id, head = validator_input(repo_objs, b"payload" * 100) assert object_validator(repo_objs)(chunk_id, head) + + +@pytest.mark.parametrize("key_class", [AuthenticatedKey, CHPOKey, AESOCBKey]) +def test_whole_object_authenticator_checks_both_slots(key_class): + # a flipped byte in the metadata slot or in the data slot, a wrong chunk id or a truncated object + # fails authentication, for each key family. + key = key_class(None) + key.init_from_random_data() + key.init_ciphers() + repo_objs = RepoObj(key) + data = b"payload" * 100 + chunk_id = repo_objs.id_hash(data) + obj = repo_objs.format(chunk_id, {}, data, ro_type=ROBJ_FILE_STREAM) + authenticate = whole_object_authenticator(repo_objs) + assert authenticate(chunk_id, obj) + assert authenticate(chunk_id, memoryview(obj)) + # the first byte after the metadata slot's envelope header and tag, and the last byte of the data slot: + # each is covered by its slot's tag. + for pos in (RepoObj.obj_header.size + key.PAYLOAD_OVERHEAD, len(obj) - 1): + bad = bytearray(obj) + bad[pos] ^= 0x01 + assert not authenticate(chunk_id, bytes(bad)) + assert not authenticate(repo_objs.id_hash(b"other"), obj) + assert not authenticate(chunk_id, obj[:-1]) + + +@pytest.mark.parametrize("key_class", [AuthenticatedKey, CHPOKey, AESOCBKey]) +def test_whole_object_authenticator_refuses_authenticated_no_key(key_class, monkeypatch): + # the workaround skips the tag verification of the "authenticated-*" keys only. + key = key_class(None) + key.init_from_random_data() + key.init_ciphers() + monkeypatch.setattr("borg.repoobj.AUTHENTICATED_NO_KEY", True) + if key_class is AuthenticatedKey: + with pytest.raises(Error, match="authenticated_no_key"): + whole_object_authenticator(RepoObj(key)) + else: + assert callable(whole_object_authenticator(RepoObj(key))) + + +def test_whole_object_authenticator_does_not_check_the_id(aead_key): + # only the tags are checked: an object whose slots authenticate but whose plaintext does not match + # its id is accepted. + repo_objs = RepoObj(aead_key) + chunk_id = repo_objs.id_hash(b"foobar" * 10) + assert whole_object_authenticator(repo_objs)(chunk_id, wrong_content_object(repo_objs, chunk_id)) + + +def test_whole_object_authenticator_propagates_an_unexpected_exception(monkeypatch): + # only a failed authentication or unpacking gives False, any other exception propagates. + repo_objs = RepoObj(make_test_key()) + data = b"payload" * 100 + chunk_id = repo_objs.id_hash(data) + obj = repo_objs.format(chunk_id, {}, data, ro_type=ROBJ_FILE_STREAM) + + def parse_raising_valueerror(*args, **kwargs): + raise ValueError("not a failure to authenticate") + + monkeypatch.setattr(repo_objs, "parse", parse_raising_valueerror) + with pytest.raises(ValueError): + whole_object_authenticator(repo_objs)(chunk_id, obj) diff --git a/src/borg/testsuite/repository_test.py b/src/borg/testsuite/repository_test.py index 4d7d0a0e18..55019ab20c 100644 --- a/src/borg/testsuite/repository_test.py +++ b/src/borg/testsuite/repository_test.py @@ -9,6 +9,7 @@ import pytest from borghash import HashTableNT +from borgstore.backends.errors import BackendConnectionError, BackendMustBeOpen, PermissionDenied, ReadRangeError from borgstore.backends.rest import REST from ..crypto.key import store_hash @@ -22,7 +23,9 @@ from ..platform import get_process_id from ..repository import Repository, MAX_DATA_SIZE, MAX_VALIDATED_META_SIZE, propagate_rsh, rest_serve_command from ..repository import PackWriter, PackReader, PackTracker, superseded_gap_ranges -from ..repoobj import RepoObj, OBJ_MAGIC, OBJ_VERSION, object_validator +from ..repository import SALVAGE_DONE, SALVAGE_INTACT, SALVAGE_NOTHING_AUTHENTICATES +from ..repository import SALVAGE_READ_ERROR, SALVAGE_READS_DIFFER, StoreObjectNotFound +from ..repoobj import RepoObj, OBJ_MAGIC, OBJ_VERSION, object_validator, whole_object_authenticator from . import make_test_key, set_test_key_on_open from .hashindex_test import H from .repoobj_test import CHUNK_ID_OFFSET, DATA_SIZE_OFFSET, META_SIZE_OFFSET @@ -3270,7 +3273,10 @@ def test_superseded_gap_ranges_warns_where_it_keeps_bytes(tmp_path, caplog): assert gap_ranges(bytes(damaged) + garbage, chunks, validate) == [] assert gap_ranges(b"\0" * 3, chunks, validate) == [] - assert f"pack {pack_hex}: object does not authenticate at offset 0 in a gap, keeping its bytes." in caplog.text + assert ( + f"pack {pack_hex}: object header or metadata does not authenticate at offset 0 in a gap, keeping its bytes." + in caplog.text + ) assert ( f"pack {pack_hex}: no object header at offset {len(rejected)} in a gap, " f"keeping the remaining {len(garbage)} bytes of the gap." in caplog.text @@ -3438,3 +3444,380 @@ def failing_save_config(self, key=None): with Repository(location, exclusive=True, create=True): # and creating it afterwards works pass assert os.path.exists(os.path.join(location, "config", "config")) + + +def store_damaged_pack(repository, objs, *, listed, flip=(), tail=b"", size=100): + """Store objs as one pack named by the store hash of their bytes, then damage it. + + The byte at each position in flip is flipped and tail is appended, so a flip or a tail makes the + pack fail its store hash. The objects at the indexes in listed get chunk index entries, with + plaintext size (the tests' chunks hold 100 bytes). + Returns pack_id. + """ + pack = b"".join(obj for _, obj in objs) + pack_id = store_hash(pack).digest() + offsets = [] + offset = 0 + for _, obj in objs: + offsets.append(offset) + offset += len(obj) + damaged = bytearray(pack) + for pos in flip: + damaged[pos] ^= 0xFF + repository.store_store("packs/" + bin_to_hex(pack_id), bytes(damaged) + tail) + for i in listed: + chunk_id, obj = objs[i] + repository.chunks[chunk_id] = ChunkIndexEntry( + flags=ChunkIndex.F_USED, size=size, pack_id=pack_id, obj_offset=offsets[i], obj_size=len(obj) + ) + return pack_id + + +def last_byte_offset(objs, i): + # position of the last byte of objs[i] in a pack of objs: the last byte of its data slot. + return sum(len(obj) for _, obj in objs[: i + 1]) - 1 + + +def salvage(repository, repo_objs, pack_id, **kwargs): + kwargs.setdefault("authenticate", whole_object_authenticator(repo_objs)) + return repository.salvage_pack(pack_id, validate=object_validator(repo_objs), **kwargs) + + +def store_contents(repository): + # {name: content} of every pack in the store. + return {info.name: repository.store_load("packs/" + info.name) for info in repository.store_list("packs")} + + +def index_contents(chunks): + return dict(chunks.iteritems()) + + +@pytest.fixture() +def salvage_repository(tmp_path): + with Repository(os.fspath(tmp_path / "repo"), exclusive=True, create=True) as repository: + repository.chunks = ChunkIndex() # the tests add their entries to an empty chunk index + yield repository + + +def three_objects(repo_objs): + return [real_chunk(repo_objs, bytes([i]) * 100) for i in range(3)] + + +def test_salvage_pack_leaves_an_intact_pack_alone(salvage_repository): + repo_objs = plain_repo_objs() + objs = three_objects(repo_objs) + pack_id = store_damaged_pack(salvage_repository, objs, listed=range(3)) + packs_before, index_before = store_contents(salvage_repository), index_contents(salvage_repository.chunks) + called = [] + + result = salvage(salvage_repository, repo_objs, pack_id, before_old_pack_delete=lambda: called.append(1)) + + assert result.status == SALVAGE_INTACT + assert (result.new_pack_id, result.kept, result.dropped_bytes, result.removed_ids) == (None, [], 0, []) + assert store_contents(salvage_repository) == packs_before + assert index_contents(salvage_repository.chunks) == index_before + assert not called + + +def test_salvage_pack_drops_an_object_failing_authentication(salvage_repository): + # a flipped byte in the data slot of the middle object: the header and metadata slot still validate, + # so the walk yields it, and authenticate rejects it. + repo_objs = plain_repo_objs() + objs = three_objects(repo_objs) + (id0, obj0), (id1, obj1), (id2, obj2) = objs + pack_id = store_damaged_pack(salvage_repository, objs, listed=range(3), flip=[last_byte_offset(objs, 1)]) + list(salvage_repository.get_many([id0])) # loads the pack into the pack cache + + result = salvage(salvage_repository, repo_objs, pack_id) + + assert result.status == SALVAGE_DONE + assert result.new_pack_id == store_hash(obj0 + obj2).digest() + assert result.kept == [(id0, 0, len(obj0)), (id2, len(obj0), len(obj2))] + assert result.dropped_bytes == len(obj1) + assert result.removed_ids == [id1] + assert store_contents(salvage_repository) == {bin_to_hex(result.new_pack_id): obj0 + obj2} + assert pack_id not in salvage_repository._pack_cache + assert id1 not in salvage_repository.chunks + assert bytes(salvage_repository.get(id0)) == obj0 + assert bytes(salvage_repository.get(id2)) == obj2 + assert salvage_repository.chunks[id2].size == 100 # a repointed entry keeps its size + + +@pytest.mark.parametrize("field", ["magic", "metadata"]) +def test_salvage_pack_drops_an_object_the_walk_skips(salvage_repository, field): + # a flipped byte in the header magic or in the metadata slot of the middle object: the walk rejects + # its header, resyncs at the next object and never yields the damaged one to authenticate. + repo_objs = plain_repo_objs() + objs = three_objects(repo_objs) + (id0, obj0), (id1, obj1), (id2, obj2) = objs + start = last_byte_offset(objs, 0) + 1 # the first byte of objs[1] + pos = {"magic": start, "metadata": start + RepoObj.obj_header.size + repo_objs.key.PAYLOAD_OVERHEAD}[field] + pack_id = store_damaged_pack(salvage_repository, objs, listed=range(3), flip=[pos]) + authenticate = whole_object_authenticator(repo_objs) + authenticated = [] + + def spy_authenticate(chunk_id, obj): + authenticated.append(chunk_id) + return authenticate(chunk_id, obj) + + result = salvage(salvage_repository, repo_objs, pack_id, authenticate=spy_authenticate) + + assert result.status == SALVAGE_DONE + assert authenticated == [id0, id2] + assert result.kept == [(id0, 0, len(obj0)), (id2, len(obj0), len(obj2))] + assert result.dropped_bytes == len(obj1) + assert result.removed_ids == [id1] + assert store_contents(salvage_repository) == {bin_to_hex(result.new_pack_id): obj0 + obj2} + + +def test_salvage_pack_keeps_an_object_the_index_does_not_list(salvage_repository): + repo_objs = plain_repo_objs() + objs = three_objects(repo_objs) + (id0, obj0), (id1, obj1), (id2, obj2) = objs + pack_id = store_damaged_pack(salvage_repository, objs, listed=[0, 2], flip=[last_byte_offset(objs, 2)]) + + result = salvage(salvage_repository, repo_objs, pack_id) + + assert result.status == SALVAGE_DONE + assert result.kept == [(id0, 0, len(obj0)), (id1, len(obj0), len(obj1))] + assert result.removed_ids == [id2] + entry = salvage_repository.chunks[id1] + assert (entry.flags, entry.size) == (ChunkIndex.F_USED, 0) # the plaintext size is unknown + assert (entry.pack_id, entry.obj_offset, entry.obj_size) == (result.new_pack_id, len(obj0), len(obj1)) + assert bytes(salvage_repository.get(id1)) == obj1 + + +def test_salvage_pack_leaves_a_pack_alone_if_nothing_authenticates(salvage_repository): + repo_objs = plain_repo_objs() + objs = three_objects(repo_objs) + flip = [last_byte_offset(objs, i) for i in range(3)] + pack_id = store_damaged_pack(salvage_repository, objs, listed=range(3), flip=flip) + packs_before, index_before = store_contents(salvage_repository), index_contents(salvage_repository.chunks) + + result = salvage(salvage_repository, repo_objs, pack_id) + + assert result.status == SALVAGE_NOTHING_AUTHENTICATES + assert result.new_pack_id is None + assert store_contents(salvage_repository) == packs_before + assert index_contents(salvage_repository.chunks) == index_before + + +@pytest.mark.parametrize("tail", [b"x" * 10, b"junk" * 100], ids=["shorter-than-a-header", "longer"]) +def test_salvage_pack_drops_uncovered_trailing_bytes(salvage_repository, tail): + repo_objs = plain_repo_objs() + objs = three_objects(repo_objs) + pack_id = store_damaged_pack(salvage_repository, objs, listed=range(3), tail=tail) + index_before = index_contents(salvage_repository.chunks) + called = [] + + result = salvage(salvage_repository, repo_objs, pack_id, before_old_pack_delete=lambda: called.append(1)) + + assert result.status == SALVAGE_DONE + assert result.dropped_bytes == len(tail) + assert result.removed_ids == [] + # the kept bytes are the original pack, so the replacement has the store hash name of the undamaged pack + # and the old pack is not deleted. + assert result.new_pack_id == pack_id + assert called == [] + assert store_contents(salvage_repository) == {bin_to_hex(pack_id): b"".join(obj for _, obj in objs)} + assert index_contents(salvage_repository.chunks) == index_before + for chunk_id, obj in objs: + assert bytes(salvage_repository.get(chunk_id)) == obj + + +def test_salvage_pack_refuses_with_a_pack_store_cache(tmp_path, monkeypatch): + # with BORG_STORE_CACHE, both loads of a pack can return the same cached copy. + monkeypatch.setenv("BORG_STORE_CACHE", os.fspath(tmp_path / "cache")) + with Repository(os.fspath(tmp_path / "repo"), exclusive=True, create=True) as repository: + repository.chunks = ChunkIndex() + repo_objs = plain_repo_objs() + objs = three_objects(repo_objs) + pack_id = store_damaged_pack(repository, objs, listed=range(3), flip=[last_byte_offset(objs, 1)]) + with pytest.raises(Error, match="BORG_STORE_CACHE"): + salvage(repository, repo_objs, pack_id) + assert bin_to_hex(pack_id) in store_contents(repository) + + +def test_salvage_pack_refuses_without_write_permission(salvage_repository): + repo_objs = plain_repo_objs() + objs = three_objects(repo_objs) + pack_id = store_damaged_pack(salvage_repository, objs, listed=range(3), flip=[last_byte_offset(objs, 1)]) + packs_before, index_before = store_contents(salvage_repository), index_contents(salvage_repository.chunks) + salvage_repository.permissions = {"packs": "lrw", "index": "lrwWD"} # packs/ has no delete + with pytest.raises(Repository.PermissionDenied): + salvage(salvage_repository, repo_objs, pack_id) + assert store_contents(salvage_repository) == packs_before + assert index_contents(salvage_repository.chunks) == index_before + + +def test_salvage_pack_leaves_a_pack_alone_if_two_reads_differ(salvage_repository, monkeypatch): + # the first load has an extra flipped byte in object 0, the second load does not. + repo_objs = plain_repo_objs() + objs = three_objects(repo_objs) + pack_id = store_damaged_pack(salvage_repository, objs, listed=range(3), flip=[last_byte_offset(objs, 1)]) + packs_before, index_before = store_contents(salvage_repository), index_contents(salvage_repository.chunks) + load = salvage_repository.store.load + loads = [] + + def flaky_load(name, **kwargs): + data = load(name, **kwargs) + loads.append(name) + if len(loads) == 1: + data = bytearray(data) + data[last_byte_offset(objs, 0)] ^= 0xFF + data = bytes(data) + return data + + monkeypatch.setattr(salvage_repository.store, "load", flaky_load) + + result = salvage(salvage_repository, repo_objs, pack_id) + + assert result.status == SALVAGE_READS_DIFFER + assert len(loads) == 2 + monkeypatch.undo() + assert store_contents(salvage_repository) == packs_before + assert index_contents(salvage_repository.chunks) == index_before + + +def test_salvage_pack_leaves_a_pack_that_reads_intact_alone(salvage_repository, monkeypatch): + # the store hash does not match the name, but the loaded bytes do: the reads disagree, and no object + # is authenticated. + repo_objs = plain_repo_objs() + objs = three_objects(repo_objs) + pack_id = store_damaged_pack(salvage_repository, objs, listed=range(3)) + packs_before, index_before = store_contents(salvage_repository), index_contents(salvage_repository.chunks) + monkeypatch.setattr(salvage_repository.store, "hash", lambda name, algorithm: "0" * 64) + authenticated = [] + + def authenticate(chunk_id, obj): + authenticated.append(chunk_id) + return False + + result = salvage(salvage_repository, repo_objs, pack_id, authenticate=authenticate) + + assert result.status == SALVAGE_READS_DIFFER + assert authenticated == [] + assert store_contents(salvage_repository) == packs_before + assert index_contents(salvage_repository.chunks) == index_before + + +@pytest.mark.parametrize("method", ["hash", "load"]) +@pytest.mark.parametrize( + "error", + [OSError(5, "Input/output error"), BackendConnectionError("connection lost"), ReadRangeError("short read")], + ids=["OSError", "BackendConnectionError", "ReadRangeError"], +) +def test_salvage_pack_leaves_a_pack_alone_on_a_read_error(salvage_repository, monkeypatch, method, error): + repo_objs = plain_repo_objs() + objs = three_objects(repo_objs) + pack_id = store_damaged_pack(salvage_repository, objs, listed=range(3), flip=[last_byte_offset(objs, 1)]) + packs_before, index_before = store_contents(salvage_repository), index_contents(salvage_repository.chunks) + + def failing(*args, **kwargs): + raise error + + monkeypatch.setattr(salvage_repository.store, method, failing) + + result = salvage(salvage_repository, repo_objs, pack_id) + + assert result.status == SALVAGE_READ_ERROR + monkeypatch.undo() + assert store_contents(salvage_repository) == packs_before + assert index_contents(salvage_repository.chunks) == index_before + + +@pytest.mark.parametrize("method", ["hash", "load"]) +@pytest.mark.parametrize( + "error", [BackendMustBeOpen("not open"), PermissionDenied("denied")], ids=["BackendMustBeOpen", "PermissionDenied"] +) +def test_salvage_pack_raises_other_backend_errors(salvage_repository, monkeypatch, method, error): + repo_objs = plain_repo_objs() + objs = three_objects(repo_objs) + pack_id = store_damaged_pack(salvage_repository, objs, listed=range(3), flip=[last_byte_offset(objs, 1)]) + packs_before, index_before = store_contents(salvage_repository), index_contents(salvage_repository.chunks) + + def failing(*args, **kwargs): + raise error + + monkeypatch.setattr(salvage_repository.store, method, failing) + + with pytest.raises(type(error)): + salvage(salvage_repository, repo_objs, pack_id) + monkeypatch.undo() + assert store_contents(salvage_repository) == packs_before + assert index_contents(salvage_repository.chunks) == index_before + + +def test_salvage_pack_raises_if_the_pack_is_missing(salvage_repository): + repo_objs = plain_repo_objs() + with pytest.raises(StoreObjectNotFound): + salvage(salvage_repository, repo_objs, bytes(32)) + + +def test_salvage_pack_updates_the_given_chunk_index(salvage_repository): + repo_objs = plain_repo_objs() + objs = three_objects(repo_objs) + pack_id = store_damaged_pack(salvage_repository, objs, listed=range(3), flip=[last_byte_offset(objs, 1)]) + chunks = salvage_repository.chunks + salvage_repository.chunks = ChunkIndex() + + result = salvage(salvage_repository, repo_objs, pack_id, chunks=chunks) + + assert result.status == SALVAGE_DONE + assert result.removed_ids == [objs[1][0]] + assert chunks[objs[0][0]].pack_id == result.new_pack_id + assert index_contents(salvage_repository.chunks) == {} + + +def test_salvage_pack_indexes_duplicates(salvage_repository): + # a chunk id indexed in another pack keeps its entry. The entry of a dropped object is pointed at a + # kept copy of the same chunk id in the pack. + repo_objs = plain_repo_objs() + id_a, obj_a = real_chunk(repo_objs, b"A" * 100) + id_b, obj_b = real_chunk(repo_objs, b"B" * 100) + id_c, obj_c = real_chunk(repo_objs, b"C" * 100) + other_pack_id = store_damaged_pack(salvage_repository, [(id_a, obj_a)], listed=[0]) + objs = [(id_a, obj_a), (id_b, obj_b), (id_b, obj_b), (id_c, obj_c)] + # id_a is unlisted here, id_b is listed at its damaged first copy, id_c is listed and its only copy is damaged. + pack_id = store_damaged_pack( + salvage_repository, objs, listed=[1, 3], flip=[last_byte_offset(objs, 1), last_byte_offset(objs, 3)] + ) + + result = salvage(salvage_repository, repo_objs, pack_id) + + assert result.status == SALVAGE_DONE + assert result.kept == [(id_a, 0, len(obj_a)), (id_b, len(obj_a), len(obj_b))] + assert result.removed_ids == [id_c] + assert salvage_repository.chunks[id_a].pack_id == other_pack_id + entry = salvage_repository.chunks[id_b] + assert (entry.pack_id, entry.obj_offset) == (result.new_pack_id, len(obj_a)) + + +def test_salvage_pack_order_of_changes(salvage_repository, monkeypatch): + # the replacement pack is stored first, then before_old_pack_delete runs while the index still points + # at the old pack, then the index is updated, and the old pack is deleted last. + repo_objs = plain_repo_objs() + objs = three_objects(repo_objs) + id0 = objs[0][0] + pack_id = store_damaged_pack(salvage_repository, objs, listed=range(3), flip=[last_byte_offset(objs, 1)]) + events = [] + store_store, store_delete = salvage_repository.store_store, salvage_repository.store_delete + + def spy_store(name, value): + events.append(("store", name, salvage_repository.chunks[id0].pack_id)) + return store_store(name, value) + + def spy_delete(name, **kwargs): + events.append(("delete", name, salvage_repository.chunks[id0].pack_id)) + return store_delete(name, **kwargs) + + monkeypatch.setattr(salvage_repository, "store_store", spy_store) + monkeypatch.setattr(salvage_repository, "store_delete", spy_delete) + + def before_old_pack_delete(): + events.append(("marker", None, salvage_repository.chunks[id0].pack_id)) + + result = salvage(salvage_repository, repo_objs, pack_id, before_old_pack_delete=before_old_pack_delete) + + new_name, old_name = "packs/" + bin_to_hex(result.new_pack_id), "packs/" + bin_to_hex(pack_id) + assert events == [("store", new_name, pack_id), ("marker", None, pack_id), ("delete", old_name, result.new_pack_id)]