ericm-db commented on code in PR #58247:
URL: https://github.com/apache/spark/pull/58247#discussion_r3882775885
##########
python/pyspark/sql/tests/connect/test_connect_local_server_pool.py:
##########
@@ -837,6 +870,155 @@ def
test_claim_skips_malformed_and_incompatible_members(self) -> None:
for uid in records:
self.assertEqual(set(self._states(uid)), {"server"})
+ def test_reap_server_unreachable_and_idle(self) -> None:
+ with self.subTest("unreachable member is retired"):
+ gone = self._live_process()
+ with _non_listening_socket() as port:
+ self._write_state(
+ self._directory.server_path("dead"),
self._server_data(port, gone.pid)
+ )
+ with self._directory:
+ self._pool.reap("dead")
+ self.assertEqual(set(self._states("dead")), {"retired"})
+ self.assertTrue(_wait_proc_dead(gone))
+ with self.subTest("member idle past the timeout is retired"):
+ os.environ["SPARK_LOCAL_CONNECT_POOL_IDLE_TIMEOUT"] = "10"
+ with _listening_socket() as port:
+ idle = self._live_process()
+ self._write_state(
+ self._directory.server_path("1d1e"),
+ self._server_data(port, idle.pid, created=time.time() -
60),
+ )
+ fresh = self._live_process()
+ self._write_state(
+ self._directory.server_path("f2e5"),
self._server_data(port, fresh.pid)
+ )
+ with self._directory:
+ self._pool.janitor()
+ self.assertEqual(set(self._states("1d1e")), {"retired"})
+ self.assertEqual(set(self._states("f2e5")), {"server"})
+ self.assertTrue(_wait_proc_dead(idle))
+ self.assertIsNone(fresh.poll())
+
+ def test_reap_does_not_signal_reused_server_pid(self) -> None:
+ unrelated = subprocess.Popen(
+ [sys.executable, "-c", "import sys; sys.stdin.buffer.read()"],
+ stdin=subprocess.PIPE,
+ stdout=subprocess.DEVNULL,
+ stderr=subprocess.DEVNULL,
+ )
+ self._procs.append(unrelated)
+ with _non_listening_socket() as port:
+ self._write_state(
+ self._directory.server_path("bad6"), self._server_data(port,
unrelated.pid)
+ )
+ with self._directory:
+ self._pool.reap("bad6")
+
+ self.assertEqual(set(self._states("bad6")), {"retired"})
+ self.assertIsNone(unrelated.poll())
+
+ def test_reap_malformed_member_recovers_server_pid(self) -> None:
+ # An out-of-range created value makes the full member invalid, but its
independently
+ # valid pid must still be retired and tracked rather than leaving a
live JVM orphaned.
+ invalid_created = self._live_process()
+ with _non_listening_socket() as port:
+ self._write_state(
+ self._directory.server_path("bad0"),
+ self._server_data(port, invalid_created.pid, created=2**100),
+ )
+ with self._directory:
+ self._pool.reap("bad0")
+ states = self._states("bad0")
+ self.assertEqual(set(states), {"retired"})
+ with self._directory as directory:
+ self.assertEqual(directory.read_json(states["retired"])["pid"],
invalid_created.pid)
+ self.assertTrue(_wait_proc_dead(invalid_created))
+
+ # Claimed records use the same recovery when their client has
disappeared.
+ claimed = self._live_process()
+ with _non_listening_socket() as port:
+ self._write_state(
+ self._directory.claimed_path(2**31 - 1, "bad2"),
+ self._server_data(port, claimed.pid, created=2**100),
+ )
+ with self._directory:
+ self._pool.reap("bad2")
+ states = self._states("bad2")
+ self.assertEqual(set(states), {"retired"})
+ with self._directory as directory:
+ self.assertEqual(directory.read_json(states["retired"])["pid"],
claimed.pid)
+ self.assertTrue(_wait_proc_dead(claimed))
+
+ # If the record itself has no usable pid, fall back to
spark-daemon.sh's pid file. The
+ # fallback keeps the shutdown tracked but lacks the process identity
needed to signal it.
+ unreadable_pid = self._live_process()
+ self._write_daemon_pid("bad1", unreadable_pid.pid)
+ self._write_state(self._directory.server_path("bad1"), {"malformed":
True})
+ with self._directory:
+ self._pool.reap("bad1")
+ states = self._states("bad1")
+ self.assertEqual(set(states), {"member", "retired"})
+ with self._directory as directory:
+ retired = directory.read_json(states["retired"])
+ assert retired is not None
+ self.assertEqual(retired["pid"], unreadable_pid.pid)
+ self.assertNotIn("process_start_id", retired)
+ self.assertIsNone(unreadable_pid.poll())
+
+ def test_reap_claimed_of_dead_client(self) -> None:
Review Comment:
added!
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]