Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 27 additions & 0 deletions src/borg/repoobj.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
147 changes: 145 additions & 2 deletions src/borg/repository.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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):
Expand Down Expand Up @@ -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).

Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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 <pack_id> 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.

Expand Down
3 changes: 2 additions & 1 deletion src/borg/testsuite/archiver/check_cmd_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
62 changes: 62 additions & 0 deletions src/borg/testsuite/repoobj_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
OBJ_VERSION,
RepoObj,
object_validator,
whole_object_authenticator,
)
from ..legacy.repoobj import RepoObj1
from ..compress import LZ4
Expand Down Expand Up @@ -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)
Loading
Loading