| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -126,7 +126,17 @@ async def state_lookup(self) -> int: | |||
| 126 | 126 | ||
| 127 | 127 | async def open(self) -> None: | |
| 128 | 128 | """Opens the underlying bidi-gRPC stream.""" | |
| 129 | - raise NotImplementedError("open is not implemented yet.") | ||
| 129 | + if self._is_stream_open: | ||
| 130 | + raise ValueError("Underlying bidi-gRPC stream is already open") | ||
| 131 | + | ||
| 132 | + await self.write_obj_stream.open() | ||
| 133 | + self._is_stream_open = True | ||
| 134 | + if self.generation is None: | ||
| 135 | + self.generation = self.write_obj_stream.generation_number | ||
| 136 | + self.write_handle = self.write_obj_stream.write_handle | ||
| 137 | + | ||
| 138 | + # Update self.persisted_size | ||
| 139 | + _ = await self.state_lookup() | ||
| 130 | 140 | ||
| 131 | 141 | async def append(self, data: bytes): | |
| 132 | 142 | raise NotImplementedError("append is not implemented yet.") | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -112,13 +112,62 @@ async def test_state_lookup(mock_write_object_stream, mock_client): | |||
| 112 | 112 | ||
| 113 | 113 | ||
| 114 | 114 | @pytest.mark.asyncio | |
| 115 | - async def test_unimplemented_methods_raise_error(mock_client): | ||
| 116 | - """Test that all currently unimplemented methods raise NotImplementedError.""" | ||
| 115 | + @mock.patch( | ||
| 116 | + "google.cloud.storage._experimental.asyncio.async_appendable_object_writer._AsyncWriteObjectStream" | ||
| 117 | + ) | ||
| 118 | + async def test_open_appendable_object_writer(mock_write_object_stream, mock_client): | ||
| 119 | + """Test the open method.""" | ||
| 120 | + # Arrange | ||
| 117 | 121 | writer = AsyncAppendableObjectWriter(mock_client, BUCKET, OBJECT) | |
| 122 | + mock_stream = mock_write_object_stream.return_value | ||
| 123 | + mock_stream.open = mock.AsyncMock() | ||
| 124 | + mock_stream.send = mock.AsyncMock() | ||
| 125 | + mock_stream.recv = mock.AsyncMock() | ||
| 118 | 126 | ||
| 119 | - with pytest.raises(NotImplementedError): | ||
| 127 | + mock_state_response = mock.MagicMock() | ||
| 128 | + mock_state_response.persisted_size = 1024 | ||
| 129 | + mock_stream.recv.return_value = mock_state_response | ||
| 130 | + | ||
| 131 | + mock_stream.generation_number = GENERATION | ||
| 132 | + mock_stream.write_handle = WRITE_HANDLE | ||
| 133 | + | ||
| 134 | + # Act | ||
| 135 | + await writer.open() | ||
| 136 | + | ||
| 137 | + # Assert | ||
| 138 | + mock_stream.open.assert_awaited_once() | ||
| 139 | + assert writer._is_stream_open | ||
| 140 | + assert writer.generation == GENERATION | ||
| 141 | + assert writer.write_handle == WRITE_HANDLE | ||
| 142 | + | ||
| 143 | + expected_request = _storage_v2.BidiWriteObjectRequest(state_lookup=True) | ||
| 144 | + mock_stream.send.assert_awaited_once_with(expected_request) | ||
| 145 | + mock_stream.recv.assert_awaited_once() | ||
| 146 | + assert writer.persisted_size == 1024 | ||
| 147 | + | ||
| 148 | + | ||
| 149 | + @pytest.mark.asyncio | ||
| 150 | + @mock.patch( | ||
| 151 | + "google.cloud.storage._experimental.asyncio.async_appendable_object_writer._AsyncWriteObjectStream" | ||
| 152 | + ) | ||
| 153 | + async def test_open_when_already_open_raises_error( | ||
| 154 | + mock_write_object_stream, mock_client | ||
| 155 | + ): | ||
| 156 | + """Test that opening an already open writer raises a ValueError.""" | ||
| 157 | + # Arrange | ||
| 158 | + writer = AsyncAppendableObjectWriter(mock_client, BUCKET, OBJECT) | ||
| 159 | + writer._is_stream_open = True # Manually set to open | ||
| 160 | + | ||
| 161 | + # Act & Assert | ||
| 162 | + with pytest.raises(ValueError, match="Underlying bidi-gRPC stream is already open"): | ||
| 120 | 163 | await writer.open() | |
| 121 | 164 | ||
| 165 | + | ||
| 166 | + @pytest.mark.asyncio | ||
| 167 | + async def test_unimplemented_methods_raise_error(mock_client): | ||
| 168 | + """Test that all currently unimplemented methods raise NotImplementedError.""" | ||
| 169 | + writer = AsyncAppendableObjectWriter(mock_client, BUCKET, OBJECT) | ||
| 170 | + | ||
| 122 | 171 | with pytest.raises(NotImplementedError): | |
| 123 | 172 | await writer.append(b"data") | |
| 124 | 173 | ||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -2350,7 +2350,7 @@ def test_move_blob_needs_url_encoding(self): | |||
| 2350 | 2350 | timeout=30, | |
| 2351 | 2351 | retry=None, | |
| 2352 | 2352 | _target_object=new_blob, | |
| 2353 | - ) | ||
| 2353 | + ) | ||
| 2354 | 2354 | ||
| 2355 | 2355 | def test_move_blob_w_user_project(self): | |
| 2356 | 2356 | source_name = "source" | |
| Back | FazBrowse Home | New Git URL |
0 commit comments