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
16 changes: 15 additions & 1 deletion src/remote_ops.py
Original file line number Diff line number Diff line change
Expand Up @@ -250,7 +250,21 @@ def terminate(self) -> None:
if self._read_remote_rc() is not None:
return

self.send_signal(os_signal.SIGTERM)
try:
self.send_signal(os_signal.SIGTERM)
except ExecUtilException as e1:
if e1.exit_code != 1:
raise

assert type(e1.exit_code) is int
assert e1.exit_code == 1 # No such process

try:
self.wait(
OsOperationStaticConfig.remote_process_controller__terminate_wait_timeout,
)
except ExecTimeoutException:
raise e1
return

def poll(self) -> typing.Optional[int]:
Expand Down
7 changes: 7 additions & 0 deletions src/static_config.py
Original file line number Diff line number Diff line change
Expand Up @@ -87,4 +87,11 @@ class OsOperationStaticConfig:
4 * 3600.0,
)

remote_process_controller__terminate_wait_timeout = _get_opt_float(
10.0,
"TESTGRES_OS_OPS_CFG__REMOTE_PROCESS_CONTROLLER__TERMINATE_WAIT_TIMEOUT",
0.0,
60.0,
)

# //////////////////////////////////////////////////////////////////////////////
95 changes: 94 additions & 1 deletion tests/test_os_ops_common.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@
from tests.helpers.local_check import LocalCheck
from tests.helpers.local_check import OsOpsHelpers

from tests.conftest_helpers import TestServices

from src.os_ops import OsProcessController
from src.os_ops import OsCommandResult
from src.os_ops import T_OS_EXEC_ENV
Expand Down Expand Up @@ -4710,7 +4712,98 @@ def test_popen_terminate(self, os_ops_descr: OsOpsDescr):
type(os_ops).__name__,
))
pass
pass
return

def test_popen_terminate_mt(
self,
os_ops_descr: OsOpsDescr,
):
assert type(os_ops_descr) is OsOpsDescr
assert isinstance(os_ops_descr.os_ops, OsOperations)

RunConditions.skip_if_windows()
os_ops = os_ops_descr.os_ops

controller: typing.Optional[OsProcessController] = None

try:
N_WORKERS = 100

logging.info("Process is creating ...")
cmd1 = ["sleep", "100"]
controller = os_ops.popen(cmd1)
assert isinstance(controller, OsProcessController)

logging.info("Worker are creating ...")
threadPool = ThreadPoolExecutor(
max_workers=N_WORKERS,
thread_name_prefix="ex_creator",
)

class tadWorkerData:
future: ThreadFuture

workerDatas: typing.List[tadWorkerData] = list()

nErrors = 0

try:
for n in range(N_WORKERS):
logging.info("worker #{} is creating ...".format(n))

workerDatas.append(tadWorkerData())

workerDatas[n].future = threadPool.submit(
controller.terminate,
)

assert workerDatas[n].future is not None

logging.info("OK. All the workers were created!")
except BaseException as e:
nErrors += 1
logging.error("A problem is detected ({}): {}".format(
type(e).__name__,
TestServices.ExceptionToHumanText(e),
))

logging.info("Will wait for stop of all the workers...")

nWorkers = 0

assert type(workerDatas) is list

for i in range(len(workerDatas)):
worker = workerDatas[i].future

if worker is None:
break

nWorkers += 1

assert isinstance(worker, ThreadFuture)

try:
logging.info("Wait for worker #{}".format(i))
worker.result()
except BaseException as e:
nErrors += 1
logging.error("Worker #{} finished with error ({}): {}".format(
i,
type(e).__name__,
TestServices.ExceptionToHumanText(e),
))
continue

assert nWorkers == N_WORKERS

if nErrors != 0:
raise RuntimeError("Some problems were detected. Please examine the log messages.")

finally:
if controller is not None:
controller.close()

return

def test_popen_kill(self, os_ops_descr: OsOpsDescr):
Expand Down
Loading