193 lines
8.5 KiB
Python
193 lines
8.5 KiB
Python
"""Stream clean MinIO artwork to an archive and verify it via an isolated prefix."""
|
|
import hashlib
|
|
import io
|
|
import json
|
|
import os
|
|
import sys
|
|
import tarfile
|
|
from datetime import datetime, timezone
|
|
from uuid import uuid4
|
|
|
|
from app.adapters import LocalS3Storage, require_local
|
|
from app.core.db import connect
|
|
|
|
CHUNK = 8 * 1024 * 1024
|
|
MAX_OBJECT = min(128 * 1024 * 1024, int(os.environ.get('SCAN_MAX_BYTES', '134217728')))
|
|
MAX_TOTAL = int(os.environ.get('STORAGE_QUOTA_BYTES', '53687091200'))
|
|
MAX_OBJECTS = 100000
|
|
|
|
|
|
class DigestReader:
|
|
def __init__(self, stream):
|
|
self.stream = stream
|
|
self.digest = hashlib.sha256()
|
|
self.size = 0
|
|
|
|
def read(self, size=-1):
|
|
block = self.stream.read(size)
|
|
if block:
|
|
self.digest.update(block)
|
|
self.size += len(block)
|
|
return block
|
|
|
|
|
|
def add_bytes(archive, name, content):
|
|
info = tarfile.TarInfo(name)
|
|
info.size = len(content)
|
|
info.mtime = int(datetime.now(timezone.utc).timestamp())
|
|
info.mode = 0o600
|
|
archive.addfile(info, io.BytesIO(content))
|
|
|
|
|
|
def export_archive():
|
|
require_local()
|
|
storage = LocalS3Storage()
|
|
with connect() as c:
|
|
rows = c.execute("""SELECT id,object_key,size,created_at,expires_at
|
|
FROM dtf_local.uploads WHERE complete AND purged_at IS NULL
|
|
AND expires_at>now() AND scan_state='clean' ORDER BY object_key""").fetchall()
|
|
if len(rows) > MAX_OBJECTS:
|
|
raise RuntimeError('Object count exceeds the backup safety limit')
|
|
expected_total = sum(row['size'] for row in rows)
|
|
if expected_total > MAX_TOTAL:
|
|
raise RuntimeError('Object bytes exceed the backup safety limit')
|
|
|
|
objects = []
|
|
with tarfile.open(fileobj=sys.stdout.buffer, mode='w|gz', compresslevel=6) as archive:
|
|
for index, row in enumerate(rows):
|
|
if row['size'] <= 0 or row['size'] > MAX_OBJECT:
|
|
raise RuntimeError('A clean object has an invalid backup size')
|
|
response = storage.client.get_object(Bucket=storage.bucket, Key=row['object_key'])
|
|
if response['ContentLength'] != row['size']:
|
|
response['Body'].close()
|
|
raise RuntimeError('Stored object size differs from database metadata')
|
|
reader = DigestReader(response['Body'])
|
|
name = f'objects/{index:08d}'
|
|
info = tarfile.TarInfo(name)
|
|
info.size = row['size']
|
|
info.mtime = int(row['created_at'].timestamp())
|
|
info.mode = 0o600
|
|
try:
|
|
archive.addfile(info, reader)
|
|
finally:
|
|
response['Body'].close()
|
|
if reader.size != row['size']:
|
|
raise RuntimeError('Object changed while the backup was streaming')
|
|
objects.append({'archive_name': name, 'upload_id': str(row['id']),
|
|
'object_key': row['object_key'], 'size': row['size'],
|
|
'sha256': reader.digest.hexdigest(),
|
|
'expires_at': row['expires_at'].isoformat()})
|
|
manifest = {'format': 'dtf-object-backup-v1',
|
|
'created_at': datetime.now(timezone.utc).isoformat(),
|
|
'source_bucket': storage.bucket, 'objects': objects,
|
|
'object_count': len(objects), 'total_bytes': expected_total}
|
|
add_bytes(archive, 'manifest.json', json.dumps(manifest, sort_keys=True).encode())
|
|
print('DTF_BACKUP_SUMMARY '+json.dumps({'object_count': len(objects),
|
|
'total_bytes': expected_total}), file=sys.stderr)
|
|
|
|
|
|
def upload_member(storage, source, size, key):
|
|
upload = storage.client.create_multipart_upload(
|
|
Bucket=storage.bucket, Key=key, ContentType='application/octet-stream')['UploadId']
|
|
parts = []
|
|
digest = hashlib.sha256()
|
|
remaining = size
|
|
try:
|
|
number = 1
|
|
while remaining:
|
|
block = source.read(min(CHUNK, remaining))
|
|
if not block:
|
|
raise ValueError('Archive object ended before its declared size')
|
|
digest.update(block)
|
|
result = storage.client.upload_part(Bucket=storage.bucket, Key=key,
|
|
UploadId=upload, PartNumber=number, Body=block, ContentLength=len(block))
|
|
parts.append({'PartNumber': number, 'ETag': result['ETag']})
|
|
remaining -= len(block)
|
|
number += 1
|
|
storage.client.complete_multipart_upload(Bucket=storage.bucket, Key=key,
|
|
UploadId=upload, MultipartUpload={'Parts': parts})
|
|
except Exception:
|
|
try:
|
|
storage.client.abort_multipart_upload(
|
|
Bucket=storage.bucket, Key=key, UploadId=upload)
|
|
except Exception:
|
|
# Preserve the original upload/read error; the verification-prefix
|
|
# cleanup below still removes any completed temporary object.
|
|
pass
|
|
raise
|
|
return digest.hexdigest()
|
|
|
|
|
|
def object_digest(storage, key):
|
|
response = storage.client.get_object(Bucket=storage.bucket, Key=key)
|
|
digest = hashlib.sha256()
|
|
size = 0
|
|
try:
|
|
for block in response['Body'].iter_chunks(chunk_size=1024 * 1024):
|
|
digest.update(block)
|
|
size += len(block)
|
|
finally:
|
|
response['Body'].close()
|
|
return size, digest.hexdigest()
|
|
|
|
|
|
def verify_archive():
|
|
require_local()
|
|
storage = LocalS3Storage()
|
|
verification = uuid4().hex
|
|
restored = []
|
|
observed = []
|
|
manifest = None
|
|
total = 0
|
|
try:
|
|
with tarfile.open(fileobj=sys.stdin.buffer, mode='r|gz') as archive:
|
|
for member in archive:
|
|
if not member.isfile():
|
|
raise ValueError('Backup archive contains a non-file member')
|
|
source = archive.extractfile(member)
|
|
if source is None:
|
|
raise ValueError('Backup archive member cannot be read')
|
|
if member.name == 'manifest.json':
|
|
if manifest is not None or member.size > 1024 * 1024:
|
|
raise ValueError('Invalid object-backup manifest')
|
|
manifest = json.loads(source.read().decode())
|
|
continue
|
|
expected_name = f'objects/{len(observed):08d}'
|
|
if member.name != expected_name or len(observed) >= MAX_OBJECTS:
|
|
raise ValueError('Unexpected object-backup member')
|
|
if member.size <= 0 or member.size > MAX_OBJECT:
|
|
raise ValueError('Object-backup member exceeds safety limits')
|
|
total += member.size
|
|
if total > MAX_TOTAL:
|
|
raise ValueError('Object-backup total exceeds safety limits')
|
|
key = f'originals/restore-verification/{verification}/{len(observed):08d}'
|
|
digest = upload_member(storage, source, member.size, key)
|
|
restored.append(key)
|
|
observed.append({'archive_name': member.name, 'size': member.size,
|
|
'sha256': digest})
|
|
if not manifest or manifest.get('format') != 'dtf-object-backup-v1':
|
|
raise ValueError('Object-backup manifest is missing or unsupported')
|
|
expected = manifest.get('objects')
|
|
if manifest.get('object_count') != len(observed) or manifest.get('total_bytes') != total:
|
|
raise ValueError('Object-backup summary does not match archive contents')
|
|
if not isinstance(expected, list) or len(expected) != len(observed):
|
|
raise ValueError('Object-backup manifest count does not match')
|
|
for actual, recorded, key in zip(observed, expected, restored):
|
|
for field in ('archive_name', 'size', 'sha256'):
|
|
if actual[field] != recorded.get(field):
|
|
raise ValueError('Object-backup checksum manifest does not match')
|
|
restored_size, restored_hash = object_digest(storage, key)
|
|
if restored_size != actual['size'] or restored_hash != actual['sha256']:
|
|
raise ValueError('Restored verification object differs from archive')
|
|
print(f'PASS: {len(observed)} clean objects ({total} bytes) restored and hashed in an isolated prefix.')
|
|
finally:
|
|
for key in restored:
|
|
storage.client.delete_object(Bucket=storage.bucket, Key=key)
|
|
print('Removed temporary verification objects. Active artwork was untouched.')
|
|
|
|
|
|
if __name__ == '__main__':
|
|
if len(sys.argv) != 2 or sys.argv[1] not in ('export', 'verify'):
|
|
raise SystemExit('usage: python -m ops.storage_backup export|verify')
|
|
export_archive() if sys.argv[1] == 'export' else verify_archive()
|