diff --git a/AGENTS.md b/AGENTS.md new file mode 100644 index 000000000..ba5fe55bc --- /dev/null +++ b/AGENTS.md @@ -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. diff --git a/coriolis/providers/backup_writers.py b/coriolis/providers/backup_writers.py index 1c7e34737..8c20f4790 100644 --- a/coriolis/providers/backup_writers.py +++ b/coriolis/providers/backup_writers.py @@ -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 @@ -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): @@ -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.") @@ -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, @@ -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(): @@ -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: @@ -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"): + # 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": + # 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): diff --git a/coriolis/resources/bin/coriolis-writer b/coriolis/resources/bin/coriolis-writer index ea9f62458..7a6035fac 100755 Binary files a/coriolis/resources/bin/coriolis-writer and b/coriolis/resources/bin/coriolis-writer differ diff --git a/coriolis/tests/providers/test_backup_writers.py b/coriolis/tests/providers/test_backup_writers.py index 10c470ed8..913299e72 100644 --- a/coriolis/tests/providers/test_backup_writers.py +++ b/coriolis/tests/providers/test_backup_writers.py @@ -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() @@ -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( @@ -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') @@ -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