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
7 changes: 3 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -39,10 +39,9 @@ only the files an artifact asks for are read out of the image. The GUI picks
`raw` on its own for those extensions. See `admin/docs/raw_image_input.md`.

`tar` also reads an xz-compressed tar (`.tar.xz`), and the GUI picks `tar` for that
extension. A compressed tar, `.tar.gz` included, is read much more slowly than a plain one:
each time a file earlier in the archive is needed, the reader decompresses from the start
again. For a large extraction, decompress it first (`xz -dk` or `gunzip -k`) and give the
tool the `.tar`.
extension. A compressed tar, `.tar.gz` included, is decompressed once into the report folder
before any file is read, so the run needs free space there for the uncompressed tar. The
copy is deleted when the run ends, and the run log says how long the step took.
`iva` reads a Berla iVe export as it stands: the raw image inside it is read the
same way, and the export's `Vehicle.json` is reported beside the vehicle data.

Expand Down
134 changes: 134 additions & 0 deletions admin/test/scripts/test_seeker_tar_compressed.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,134 @@
"""Pin how the tar seeker reads a compressed tar: decompressed once, then read as a plain tar.

A compressed stream seeks backwards by decompressing again from its start, and artifacts
search one after another, so a compressed tar read in place rewinds over and over. On a
5.36 GB Android extraction that was 56 rewinds and 220 GB decompressed as a .tar.gz (347 s
against 33 s for the plain tar), and about two hours projected as a .tar.xz. The seeker
now decompresses a compressed tar once into the report folder, reads that copy, and deletes
it in cleanup(). When the copy cannot be written it reads the archive in place as before.

The archives are built here with tarfile, in every compressed format this Python can
write. Every expected value is one of the member contents defined below. The same file
runs in all five LEAPP cores.
"""
import errno
import hashlib
import io
import os
import pathlib
import shutil
import sys
import tarfile
import tempfile
import unittest
from unittest import mock

REPO_ROOT = pathlib.Path(__file__).resolve().parents[3]
sys.path.insert(0, str(REPO_ROOT))

import scripts.search_files as search_files # pylint: disable=wrong-import-position
from scripts.search_files import FileSeekerTar # pylint: disable=wrong-import-position

# 64 KiB that does not compress, so cutting a compressed archive in half always lands
# inside this member's data and never inside the first header.
MIDDLE = b''.join(hashlib.sha256(i.to_bytes(4, 'big')).digest() for i in range(2048))
MEMBERS = (
('a/first.txt', b'FIRST'),
('m/middle.bin', MIDDLE),
('z/last.txt', b'LAST'),
)
SPOOL_NAME = '_decompressed_input.tar'


def _compressed_modes():
modes = ['w:gz', 'w:bz2', 'w:xz']
try:
import compression.zstd # pylint: disable=import-outside-toplevel,unused-import
modes.append('w:zst') # Python 3.14 and later
except ImportError:
pass
return modes


class TestCompressedTarIsReadOnce(unittest.TestCase):

def setUp(self):
self.tmp = tempfile.mkdtemp(prefix='leapp_seeker_tar_compressed_')
self.addCleanup(shutil.rmtree, self.tmp, True)

def _archive(self, mode, name):
path = os.path.join(self.tmp, name)
with tarfile.open(path, mode) as archive:
for member, content in MEMBERS:
info = tarfile.TarInfo(member)
info.size = len(content)
info.mtime = 1000
archive.addfile(info, io.BytesIO(content))
return path

def _report(self, label):
report = os.path.join(self.tmp, label)
os.makedirs(report)
return report, os.path.join(report, 'data')

@staticmethod
def _read(path):
with open(path, 'rb') as handle:
return handle.read()

def test_a_compressed_tar_is_read_from_one_decompressed_copy(self):
for mode in _compressed_modes():
with self.subTest(mode=mode):
path = self._archive(mode, f'case.tar.{mode[2:]}')
report, data = self._report(f'report_{mode[2:]}')
seeker = FileSeekerTar(path, data)
spool = os.path.join(report, SPOOL_NAME)
self.assertTrue(os.path.isfile(spool))
self.assertNotIsInstance(seeker.tar_file.fileobj, search_files._compressed_tar_streams()) # pylint: disable=protected-access
# out of archive order: the last member first, then the first
self.assertEqual(self._read(seeker.search('*/last.txt', return_on_first_hit=True)), b'LAST')
self.assertEqual(self._read(seeker.search('*/first.txt', return_on_first_hit=True)), b'FIRST')
self.assertEqual(self._read(seeker.search('*/middle.bin', return_on_first_hit=True)), MIDDLE)
seeker.cleanup()
self.assertFalse(os.path.exists(spool))
self.assertEqual(sorted(os.listdir(report)), ['data'])

def test_a_plain_tar_is_read_in_place(self):
path = self._archive('w', 'case.tar')
report, data = self._report('report_plain')
seeker = FileSeekerTar(path, data)
self.assertIsNone(seeker._spool_path) # pylint: disable=protected-access
self.assertFalse(os.path.exists(os.path.join(report, SPOOL_NAME)))
self.assertEqual(self._read(seeker.search('*/last.txt', return_on_first_hit=True)), b'LAST')
seeker.cleanup()

def test_the_archive_is_read_in_place_when_the_copy_cannot_be_written(self):
path = self._archive('w:xz', 'full.tar.xz')
report, data = self._report('report_full')
messages = []
disk_full = OSError(errno.ENOSPC, 'No space left on device')
with mock.patch.object(search_files, 'copyfileobj', side_effect=disk_full), \
mock.patch.object(search_files, 'logfunc', side_effect=messages.append):
seeker = FileSeekerTar(path, data)
self.assertIsNone(seeker._spool_path) # pylint: disable=protected-access
self.assertFalse(os.path.exists(os.path.join(report, SPOOL_NAME)))
self.assertIsInstance(seeker.tar_file.fileobj, search_files._compressed_tar_streams()) # pylint: disable=protected-access
self.assertTrue(any('No space left on device' in m and 'in place' in m for m in messages), messages)
self.assertEqual(self._read(seeker.search('*/last.txt', return_on_first_hit=True)), b'LAST')
self.assertEqual(self._read(seeker.search('*/first.txt', return_on_first_hit=True)), b'FIRST')
seeker.cleanup()

def test_a_truncated_archive_still_fails_and_leaves_no_copy(self):
for mode in ('w:gz', 'w:xz'):
with self.subTest(mode=mode):
path = self._archive(mode, f'cut.tar.{mode[2:]}')
with open(path, 'r+b') as handle:
handle.truncate(os.path.getsize(path) // 2)
report, data = self._report(f'report_cut_{mode[2:]}')
with self.assertRaises(EOFError):
FileSeekerTar(path, data)
self.assertFalse(os.path.exists(os.path.join(report, SPOOL_NAME)))


if __name__ == '__main__':
unittest.main()
63 changes: 63 additions & 0 deletions scripts/search_files.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@
"""

import time as timex
import atexit
import importlib
import os
import shutil
import tarfile
Expand Down Expand Up @@ -291,6 +293,18 @@ def search(self, filepattern, return_on_first_hit=False, force=False):
self.searched[filepattern] = pathlist
return pathlist

def _compressed_tar_streams():
"""The reader classes tarfile wraps a compressed tar in, for the codecs this Python has."""
streams = []
for module, name in (('gzip', 'GzipFile'), ('bz2', 'BZ2File'), ('lzma', 'LZMAFile'),
('compression.zstd', 'ZstdFile')):
try:
streams.append(getattr(importlib.import_module(module), name))
except (ImportError, AttributeError):
pass
return tuple(streams)


class FileSeekerTar(FileSeekerBase):
"""
This is a class that extends FileSeekerBase to facilitate searching and extracting files
Expand Down Expand Up @@ -328,12 +342,51 @@ def __init__(self, tar_file_path, data_folder):
mode = 'r:gz' if self.is_gzip else 'r'
self.tar_file = tarfile.open(tar_file_path, mode)
self.data_folder = data_folder
self._spool_path = None
if isinstance(self.tar_file.fileobj, _compressed_tar_streams()):
self._read_from_a_decompressed_copy(tar_file_path, mode)
self.searched = {}
self.copied = {}
self.file_infos = {}
self._init_dest_guard(self.data_folder)
self._search_members, self._other_versions = self._index_member_names()

def _read_from_a_decompressed_copy(self, tar_file_path, mode):
"""Decompress a compressed tar once and read the plain copy instead.

A compressed stream can only seek backwards by decompressing again from its
start. Each search reads its matches in archive order, but artifacts search one
after another, so reading a compressed tar in place rewinds over and over. On a
5.36 GB Android extraction that was 56 rewinds and 220 GB decompressed as a
.tar.gz (347 s against 33 s for the plain tar), and about two hours projected
as a .tar.xz. The copy sits in the report folder and is deleted by cleanup(),
or at exit if a run never reaches it. When it cannot be written, the archive is
read in place as before.
"""
spool = os.path.join(os.path.dirname(os.path.normpath(self.data_folder)), '_decompressed_input.tar')
logfunc(f'Decompressing {os.path.basename(tar_file_path)} once so its files can be read in any order')
started = timex.time()
written = False
try:
self.tar_file.fileobj.seek(0)
with open(spool, 'wb') as out:
copyfileobj(self.tar_file.fileobj, out, 16 << 20)
written = True
except OSError as ex:
logfunc(f'Could not write the decompressed copy ({ex.strerror or type(ex).__name__}); '
'reading the compressed archive in place, which is much slower')
finally:
if not written and os.path.exists(spool):
os.remove(spool)
self.tar_file.close()
if not written:
self.tar_file = tarfile.open(tar_file_path, mode)
return
self._spool_path = spool
atexit.register(self._discard_spool)
self.tar_file = tarfile.open(spool, 'r:')
logfunc(f'Decompressed {os.path.getsize(spool):,} bytes in {timex.time() - started:.0f} s')

def _index_member_names(self):
"""One member per distinct name, in archive order, and the other versions of any name stored twice."""
by_name = {}
Expand Down Expand Up @@ -441,6 +494,16 @@ def _stage_other_versions(self, name):

def cleanup(self):
self.tar_file.close()
if self._spool_path:
atexit.unregister(self._discard_spool)
self._discard_spool()

def _discard_spool(self):
"""Close the archive and delete the decompressed copy made for it."""
self.tar_file.close()
if self._spool_path and os.path.exists(self._spool_path):
os.remove(self._spool_path)
self._spool_path = None


class FileSeekerZip(FileSeekerBase):
Expand Down
Loading