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
41 changes: 41 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
# Agent guidelines

Coriolis migrates VMs between clouds. Cloud support lives in separate
provider plugins (often private). Target Python 3.10 and 3.12.

## Coriolis terminology

- **Minion**: temporary worker VM for disk transfer and os-morphing
(guest prep: network, packages, drivers). **Minion pools** reuse them.
Some source providers skip minions and read disks directly.
- **Transfer**: creates destination volumes and copies disk data.
Re-run a transfer as a new **execution** to pick up later changes.
Most providers can do that incrementally.
- **Deployment**: creates the destination VM from a completed transfer.
- **Replica** vs **migration** (`transfer.scenario`: `replica` /
`live_migration`) is a licensing split, not two engines. Replicas can
be re-executed and re-deployed even after license fulfillment.

Users configure source/destination **endpoints** (credentials) and
**environment options** (transfer and resulting VM settings).

## Working in this repo

- Use `.tox/py3/bin/` (`stestr`, `ruff`); it has project deps. Ignore
`.mypy_cache`, `.ruff_cache`, `.tox`.
- Follow ruff (`tox.ini`, `ruff.toml`).
- Public methods need docstrings (subclasses may inherit). Use type
hints when the type is known.
- Do not add helpers for trivial checks such as
`server.power_status == "RUNNING"`; keep those inline.
- Do not strip still-relevant inline comments.
- If regenerating a file, replace its contents; do not append duplicates.
- Empty `__init__.py` files must not have license headers.

## Tests

- Prefer `@mock.patch` / `@mock.patch.object` over `with mock.patch`.
- For multiple mock calls, use `assert_has_calls`. For a single call,
`assert_called_once_with` is fine. Match nearby tests.
- Integration tests use Docker providers; external cloud providers are
optional.
79 changes: 74 additions & 5 deletions coriolis/providers/backup_writers.py
Original file line number Diff line number Diff line change
Expand Up @@ -257,7 +257,16 @@ def truncate(self, size):
pass

@abc.abstractmethod
def write(self, data):
def write(self, data, encoding=None, uncompressed_size=None):
"""Write the given data chunk.

Use the `seek` method to specify the write offset.

:param data: data chunk to write (raw bytes)
:param encoding: a data encoding format supported by the backup writer.
:param uncompressed_size: the uncompressed size in case of
compressed chunks.
"""
pass

@abc.abstractmethod
Expand Down Expand Up @@ -319,7 +328,19 @@ def seek(self, pos):
def truncate(self, size):
self._file.truncate(size)

def write(self, data):
def write(self, data, encoding=None, uncompressed_size=None):
"""Write the given data chunk.

Use the `seek` method to specify the write offset.

:param data: data chunk to write (raw bytes)
:param encoding: unsupported, this backup writer only accepts raw chunks.
:param uncompressed_size: unsupported.
"""
if encoding:
raise exception.InvalidInput(
f"The file backup writer does not support {encoding} encoded chunks."
)
self._file.write(data)

def close(self):
Expand Down Expand Up @@ -463,7 +484,15 @@ def _encoder(self):
self._enc_q.task_done()
LOG.debug("Backup encoder stopped.")

def write(self, data):
def write(self, data, encoding=None, uncompressed_size=None):
"""Write the given data chunk.

Use the `seek` method to specify the write offset.

:param data: data chunk to write (raw bytes)
:param encoding: unsupported, this backup writer only accepts raw chunks.
:param uncompressed_size: unsupported.
"""
if self._closing:
raise exception.CoriolisException("Attempted to write to a closed writer.")

Expand All @@ -472,6 +501,11 @@ def write(self, data):
"Failed to write data. See log for details."
) from self._exception

if encoding:
raise exception.InvalidInput(
f"The file backup writer does not support {encoding} encoded chunks."
)

payload = {
"offset": self._offset,
"data": data,
Expand Down Expand Up @@ -785,6 +819,10 @@ def _sender(self):
if payload.get("encoding", None):
enc = copy.copy(payload["encoding"])
headers["content-encoding"] = enc
if payload.get("uncompressed_size") is not None:
headers["X-Uncompressed-Content-Length"] = str(
payload["uncompressed_size"]
)

@utils.retry_on_error()
def send():
Expand Down Expand Up @@ -831,7 +869,18 @@ def send():
LOG.debug("Backup sender stopped.")

@utils.retry_on_error()
def write(self, data):
def write(self, data, encoding=None, uncompressed_size=None):
"""Write the given data chunk.

Use the `seek` method to specify the write offset.

:param data: data chunk to write (raw bytes)
:param encoding: the encoding of the data chunk, one of the following:
* compression algorithms: fastlz, deflate, gzip, zlib
* None: raw data, no encoding (default)
:param uncompressed_size: the uncompressed size in case of
compressed chunks. Required for fastlz.
"""
if self._closing:
raise exception.CoriolisException("Attempted to write to a closed writer.")
if self._exception:
Expand All @@ -841,7 +890,27 @@ def write(self, data):
"offset": self._offset,
"data": data,
}
self._comp_q.put(payload)
if encoding is None:
self._comp_q.put(payload)
elif encoding in ("fastlz", "gzip", "zlib", "deflate"):
Comment thread
petrutlucian94 marked this conversation as resolved.
# The payload is already compressed, skip the compressor
# queue, use the sender queue directly.
payload["encoding"] = encoding
payload["chunk"] = data
if encoding == "fastlz":
if uncompressed_size is None:
raise exception.InvalidInput(
"fastlz without explicit uncompressed size."
)
payload["uncompressed_size"] = uncompressed_size
self._sender_q.put(payload)
elif encoding == "incompressible":

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

What are the usual reasons for a chunk to be incompressible?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Repeated sequences (e.g. zero blocks, text, etc) can easily be compressed. On the other hand, encrypted data usually cannot be compressed. Same applies to already compressed chunks (e.g. rotated log files).

High entropy leads to low compression rates.

# The caller determined that the chunk is incompressible,
# skip the compression queue.
payload["chunk"] = data
self._sender_q.put(payload)
else:
raise exception.InvalidInput("Unsupported write encoding: %s" % encoding)
self._offset += len(data)

def _wait_for_queues(self):
Expand Down
Binary file modified coriolis/resources/bin/coriolis-writer
Binary file not shown.
169 changes: 169 additions & 0 deletions coriolis/tests/providers/test_backup_writers.py
Original file line number Diff line number Diff line change
Expand Up @@ -282,6 +282,15 @@ def test_write(self):
self.writer.write(mock.sentinel.data)
self.writer._file.write.assert_called_once_with(mock.sentinel.data)

def test_write_with_encoding(self):
self.assertRaises(
exception.InvalidInput,
self.writer.write,
mock.sentinel.data,
encoding='gzip',
)
self.writer._file.write.assert_not_called()

@mock.patch('os.system')
def test_close(self, mock_system):
self.writer.close()
Expand Down Expand Up @@ -479,6 +488,19 @@ def test_write(self):
self.assertEqual(self.writer._offset, len('test_data'))
self.assertEqual(self.writer._msg_id, 1)

def test_write_with_encoding(self):
self.writer._closing = False
self.writer._exception = None
self.writer._enc_q = self._enc_queue

self.assertRaises(
exception.InvalidInput,
self.writer.write,
'test_data',
encoding='gzip',
)
self._enc_queue.put.assert_not_called()

def test_write_with_closing(self):
self.writer._closing = True
self.assertRaises(
Expand Down Expand Up @@ -1055,6 +1077,44 @@ def test__sender(self, mock_conf, mock_uri, mock_ensure_session):
self.assertEqual(self.writer._write_error, False)
self.assertIsInstance(self.writer._exception, BaseException)

@mock.patch.object(backup_writers.HTTPBackupWriterImpl, '_ensure_session')
@mock.patch.object(backup_writers.HTTPBackupWriterImpl, '_uri')
@mock.patch.object(backup_writers, 'CONF')
def test__sender_with_uncompressed_size(
self, mock_conf, mock_uri, mock_ensure_session
):
self.writer._session = mock.MagicMock()
self.writer._sender_q = mock.MagicMock()
self.writer._sender_q.get.return_value = {
"offset": mock.sentinel.offset,
"chunk": 'compressed_data',
"encoding": "fastlz",
"uncompressed_size": 4096,
}

mock_response = mock.MagicMock()
mock_response.status_code = 200
mock_response.content = "OK"
mock_response.raise_for_status.side_effect = [None, BaseException()]
self.writer._session.post.return_value = mock_response

with self.assertLogs('coriolis.providers.backup_writers', level=logging.ERROR):
self.assertRaises(BaseException, self.writer._sender)

expected_headers = {
"X-Write-Offset": str(mock.sentinel.offset),
"X-Client-Token": self.writer._id,
"content-encoding": "fastlz",
"X-Uncompressed-Content-Length": "4096",
}

self.writer._session.post.assert_called_with(
mock_uri,
headers=expected_headers,
data='compressed_data',
timeout=mock_conf.default_requests_timeout,
)

@mock.patch("time.sleep")
@mock.patch.object(backup_writers.HTTPBackupWriterImpl, '_ensure_session')
@mock.patch.object(backup_writers.HTTPBackupWriterImpl, '_uri')
Expand Down Expand Up @@ -1112,8 +1172,117 @@ def test_write(self):
"data": 'test_data',
}
self.writer._comp_q.put.assert_called_once_with(expected_payload)
self.writer._sender_q.put.assert_not_called()
self.assertEqual(self.writer._offset, len('test_data'))

def test_write_with_precompressed_encoding(self):
self.writer._closing = False
self.writer._exception = None
self.writer._comp_q = mock.MagicMock()
self.writer._sender_q = mock.MagicMock()
original_write = testutils.get_wrapped_function(self.writer.write)

for encoding in ('gzip', 'zlib', 'deflate'):
self.writer._offset = 0
self.writer._comp_q.reset_mock()
self.writer._sender_q.reset_mock()

original_write(self.writer, 'compressed', encoding=encoding)

expected_payload = {
"offset": 0,
"data": 'compressed',
"encoding": encoding,
"chunk": 'compressed',
}
self.writer._sender_q.put.assert_called_once_with(expected_payload)
self.writer._comp_q.put.assert_not_called()
self.assertEqual(self.writer._offset, len('compressed'))

def test_write_with_fastlz_encoding(self):
self.writer._closing = False
self.writer._exception = None
self.writer._comp_q = mock.MagicMock()
self.writer._sender_q = mock.MagicMock()
self.writer._offset = 0

original_write = testutils.get_wrapped_function(self.writer.write)

original_write(
self.writer, 'compressed', encoding='fastlz', uncompressed_size=4096
)

expected_payload = {
"offset": 0,
"data": 'compressed',
"encoding": 'fastlz',
"chunk": 'compressed',
"uncompressed_size": 4096,
}
self.writer._sender_q.put.assert_called_once_with(expected_payload)
self.writer._comp_q.put.assert_not_called()
self.assertEqual(self.writer._offset, len('compressed'))

def test_write_with_fastlz_encoding_missing_uncompressed_size(self):
self.writer._closing = False
self.writer._exception = None
self.writer._comp_q = mock.MagicMock()
self.writer._sender_q = mock.MagicMock()
self.writer._offset = 0

original_write = testutils.get_wrapped_function(self.writer.write)

self.assertRaises(
exception.InvalidInput,
original_write,
self.writer,
'compressed',
encoding='fastlz',
)
self.writer._sender_q.put.assert_not_called()
self.writer._comp_q.put.assert_not_called()
self.assertEqual(self.writer._offset, 0)

def test_write_with_incompressible_encoding(self):
self.writer._closing = False
self.writer._exception = None
self.writer._comp_q = mock.MagicMock()
self.writer._sender_q = mock.MagicMock()
self.writer._offset = 0

original_write = testutils.get_wrapped_function(self.writer.write)

original_write(self.writer, 'test_data', encoding='incompressible')

expected_payload = {
"offset": 0,
"data": 'test_data',
"chunk": 'test_data',
}
self.writer._sender_q.put.assert_called_once_with(expected_payload)
self.writer._comp_q.put.assert_not_called()
self.assertEqual(self.writer._offset, len('test_data'))

def test_write_with_unsupported_encoding(self):
self.writer._closing = False
self.writer._exception = None
self.writer._comp_q = mock.MagicMock()
self.writer._sender_q = mock.MagicMock()
self.writer._offset = 0

original_write = testutils.get_wrapped_function(self.writer.write)

self.assertRaises(
exception.InvalidInput,
original_write,
self.writer,
'test_data',
encoding='lz4',
)
self.writer._sender_q.put.assert_not_called()
self.writer._comp_q.put.assert_not_called()
self.assertEqual(self.writer._offset, 0)

def test_write_with_closing(self):
self.writer._closing = True

Expand Down
Loading