1+ import logging
12import queue
23from unittest .mock import Mock
34
45import pytest
56
6- from influxdb_client_3 .write_client .client .util .multiprocessing_helper import MultiprocessingWriter
7+ from influxdb_client_3 .write_client .client .util .multiprocessing_helper import MultiprocessingWriter , _PoisonPill
78
89
910def make_writer (** kwargs ):
@@ -17,8 +18,9 @@ def make_writer(**kwargs):
1718
1819
1920class RecordingWriteApi :
20- def __init__ (self , error = None ):
21+ def __init__ (self , error = None , close_error = None ):
2122 self .error = error
23+ self .close_error = close_error
2224 self .writes = []
2325 self .close_calls = 0
2426
@@ -29,6 +31,8 @@ def write(self, **record):
2931
3032 def close (self ):
3133 self .close_calls += 1
34+ if self .close_error is not None :
35+ raise self .close_error
3236
3337
3438class FakeProcess :
@@ -85,6 +89,38 @@ def test_write_after_close_raises_runtime_error():
8589 writer .write (record = "value" )
8690
8791
92+ def test_write_after_start_queues_record ():
93+ writer = make_writer ()
94+ writer .process = FakeProcess ()
95+ writer .queue_ = FakeQueue ()
96+ writer .start ()
97+ writer .write (record = "value" )
98+
99+ assert writer .queue_ .items == [{"record" : "value" }]
100+
101+
102+ def test_start_after_close_raises_runtime_error ():
103+ writer = make_writer ()
104+ writer .close ()
105+
106+ with pytest .raises (RuntimeError , match = "after it has been closed" ):
107+ writer .start ()
108+
109+
110+ def test_start_twice_raises_runtime_error ():
111+ writer = make_writer ()
112+ writer .process = FakeProcess ()
113+ writer .queue_ = FakeQueue ()
114+ writer .start ()
115+
116+ with pytest .raises (RuntimeError , match = "already started" ):
117+ writer .start ()
118+
119+
120+ def test_get_start_processing_method_returns_context_method ():
121+ assert make_writer ().get_start_processing_method () == "fork"
122+
123+
88124def test_worker_timeout_closes_api_and_invokes_callback_once ():
89125 writer = make_writer (process_ttl = 0.01 )
90126 writer .queue_ = Mock ()
@@ -102,6 +138,48 @@ def test_worker_timeout_closes_api_and_invokes_callback_once():
102138 writer .write (record = "after-timeout" )
103139
104140
141+ def test_worker_timeout_continues_when_write_api_close_fails ():
142+ writer = make_writer ()
143+ writer .queue_ = Mock ()
144+ writer .queue_ .get .side_effect = queue .Empty
145+ write_api = RecordingWriteApi (close_error = ValueError ("close failed" ))
146+
147+ writer .run (write_api , writer .disposed , 0.01 , None )
148+
149+ assert writer .disposed .value == 1
150+ assert write_api .close_calls == 1
151+
152+
153+ def test_worker_poison_pill_closes_write_api ():
154+ writer = make_writer ()
155+ writer .queue_ .put (_PoisonPill ())
156+ write_api = RecordingWriteApi ()
157+
158+ writer .run (write_api , writer .disposed , 0.01 , None )
159+ writer .queue_ .join ()
160+
161+ assert write_api .close_calls == 1
162+
163+
164+ def test_shutdown_callback_is_called_once_when_invoked_repeatedly ():
165+ writer = make_writer ()
166+ callback = Mock ()
167+
168+ writer ._call_on_shutdown (callback )
169+ writer ._call_on_shutdown (callback )
170+
171+ callback .assert_called_once_with ()
172+
173+
174+ def test_shutdown_callback_failure_is_logged (caplog ):
175+ writer = make_writer ()
176+ callback = Mock (side_effect = ValueError ("callback failed" ))
177+
178+ writer ._call_on_shutdown (callback )
179+
180+ assert "shutdown callback failed" in caplog .text
181+
182+
105183def test_worker_write_failure_marks_queue_item_done ():
106184 writer = make_writer ()
107185 writer .queue_ .put ({"record" : "value" })
@@ -148,3 +226,13 @@ def test_context_manager_uses_close():
148226
149227 writer .start .assert_called_once_with ()
150228 writer .close .assert_called_once_with ()
229+
230+
231+ def test_destructor_swallows_cleanup_failure (caplog ):
232+ caplog .set_level (logging .DEBUG )
233+ writer = make_writer ()
234+ writer .close = Mock (side_effect = RuntimeError ("cleanup failed" ))
235+
236+ writer .__del__ ()
237+
238+ assert "cleanup failed" in caplog .text
0 commit comments