Skip to content

Commit 5570b6a

Browse files
committed
Make synthetic worker provider reachable from Linux containers
1 parent 1e462b2 commit 5570b6a

1 file changed

Lines changed: 32 additions & 4 deletions

File tree

pipeline/tests/e2e/test_worker_durability.py

Lines changed: 32 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33
import importlib.util
44
import json
55
import os
6+
import re
67
import shutil
78
import subprocess
89
import threading
@@ -97,7 +98,11 @@ def do_GET(self): # noqa: N802 - BaseHTTPRequestHandler API
9798
def log_message(self, _format, *_args):
9899
return
99100

100-
server = ThreadingHTTPServer(("127.0.0.1", 0), Handler)
101+
# Linux Docker's host-gateway reaches the host interface rather than
102+
# loopback. This server is disposable and returns only synthetic data;
103+
# bind on all interfaces so both Docker Desktop and Linux Engine can
104+
# reach it through the test-only host-gateway route.
105+
server = ThreadingHTTPServer(("0.0.0.0", 0), Handler)
101106
thread = threading.Thread(target=server.serve_forever, daemon=True)
102107
thread.start()
103108
return server, thread, state
@@ -121,6 +126,22 @@ def _result_count(self, record_id, provider):
121126
(record_id, provider),
122127
).fetchone()[0]
123128

129+
@staticmethod
130+
def _redacted_diagnostics(text):
131+
text = text or ""
132+
text = re.sub(r"(?i)postgres(?:ql)?://[^\s]+", "postgresql://[redacted]", text)
133+
text = text.replace("synthetic query", "[query-redacted]")
134+
return text[-2000:]
135+
136+
def _container_diagnostics(self, container_name):
137+
result = subprocess.run(
138+
["docker", "logs", "--tail", "40", container_name],
139+
capture_output=True,
140+
text=True,
141+
check=False,
142+
)
143+
return self._redacted_diagnostics(result.stdout + result.stderr)
144+
124145
def _docker_worker_command(self, image, database_url, provider_url, *, limit=1, lease_timeout=900):
125146
harness = (ROOT / "tests" / "e2e" / "docker_worker_sitecustomize.py").resolve()
126147
volume = f"{str(harness).replace(chr(92), '/') }:/app/sitecustomize.py:ro"
@@ -434,7 +455,11 @@ def test_real_worker_image_drains_synthetic_queue_and_exits_without_private_logs
434455
server_thread.join(10)
435456
server.server_close()
436457
self.assertEqual(result.returncode, 0, result.stderr[-4000:])
437-
self.assertEqual(state["calls"], 1)
458+
self.assertEqual(
459+
state["calls"], 1,
460+
"synthetic provider request missing; "
461+
f"worker_exit={result.returncode} diagnostics={self._redacted_diagnostics(result.stdout + result.stderr)}",
462+
)
438463
self.assertNotIn("synthetic query", result.stdout)
439464
self.assertNotIn("GEOAPIFY_API_KEY", result.stdout + result.stderr)
440465
with psycopg.connect(self.env.database_url) as db:
@@ -459,8 +484,11 @@ def test_process_kill_preserves_completed_result_and_restart_reclaims_second_job
459484
try:
460485
started = subprocess.run(command, cwd=ROOT.parent, capture_output=True, text=True, timeout=30, check=False)
461486
self.assertEqual(started.returncode, 0, started.stderr[-2000:])
462-
self.assertTrue(self._wait_for(lambda: self._result_count(first_record, "synthetic-docker") == 1, 20))
463-
self.assertTrue(state["second_started"].wait(20))
487+
self.assertTrue(
488+
self._wait_for(lambda: self._result_count(first_record, "synthetic-docker") == 1, 20),
489+
self._container_diagnostics(container_name),
490+
)
491+
self.assertTrue(state["second_started"].wait(20), self._container_diagnostics(container_name))
464492
killed = subprocess.run(["docker", "kill", container_name], capture_output=True, text=True, check=False)
465493
self.assertEqual(killed.returncode, 0, killed.stderr[-2000:])
466494
with psycopg.connect(self.env.database_url) as db:

0 commit comments

Comments
 (0)