Skip to content
Closed
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
8 changes: 4 additions & 4 deletions e2e/test/hooks_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -162,7 +162,7 @@ var _ = Describe("Hooks E2E Tests", Label("hooks"), Ordered, func() {
"--retry-timeout", "0",
"--selector", "example.com/board=hooks", "j", "power", "on")
Expect(err).To(HaveOccurred())
Expect(out).To(MatchRegexp(`(beforeLease hook fail|Exporter shutting down|Connection to exporter lost)`))
Expect(out).To(MatchRegexp(`(beforeLease hook fail|Exporter shutting down|Connection to exporter lost|unreachable after)`))

WaitForExporter("test-exporter-hooks")
})
Expand All @@ -175,7 +175,7 @@ var _ = Describe("Hooks E2E Tests", Label("hooks"), Ordered, func() {
"--retry-timeout", "0",
"--selector", "example.com/board=hooks", "j", "power", "on")
Expect(err).To(HaveOccurred())
Expect(out).To(MatchRegexp(`(beforeLease hook fail|Connection to exporter lost)`))
Expect(out).To(MatchRegexp(`(beforeLease hook fail|Connection to exporter lost|unreachable after)`))

// The exporter should release the lease and return to Available
WaitForExporter("test-exporter-hooks")
Expand All @@ -186,7 +186,7 @@ var _ = Describe("Hooks E2E Tests", Label("hooks"), Ordered, func() {
"--retry-timeout", "0",
"--selector", "example.com/board=hooks", "j", "power", "on")
Expect(err2).To(HaveOccurred())
Expect(out2).To(MatchRegexp(`(beforeLease hook fail|Connection to exporter lost)`))
Expect(out2).To(MatchRegexp(`(beforeLease hook fail|Connection to exporter lost|unreachable after)`))

// Exporter should recover again
WaitForExporter("test-exporter-hooks")
Expand All @@ -212,7 +212,7 @@ var _ = Describe("Hooks E2E Tests", Label("hooks"), Ordered, func() {
"--retry-timeout", "0",
"--selector", "example.com/board=hooks", "j", "power", "on")
Expect(err).To(HaveOccurred())
Expect(out).To(MatchRegexp(`(beforeLease hook fail|Exporter shutting down|Connection to exporter lost)`))
Expect(out).To(MatchRegexp(`(beforeLease hook fail|Exporter shutting down|Connection to exporter lost|unreachable after)`))

// Exporter process should have exited (allow extra time on slower runners like ARM)
Eventually(func() bool {
Expand Down
7 changes: 5 additions & 2 deletions python/packages/jumpstarter-cli/jumpstarter_cli/shell.py
Original file line number Diff line number Diff line change
Expand Up @@ -504,6 +504,7 @@ async def _shell_with_signal_handling( # noqa: C901
try:
async with anyio.from_thread.BlockingPortal() as portal:
connect_deadline = None
connect_start = None
while True:
async with config.lease_async(
selector, exporter_name, lease_name, duration, portal, acquisition_timeout,
Expand Down Expand Up @@ -534,11 +535,13 @@ async def _shell_with_signal_handling( # noqa: C901
"Session is no longer valid."
) from unreachable
if connect_deadline is None:
connect_deadline = time.monotonic() + lease.retry_timeout
connect_start = time.monotonic()
connect_deadline = connect_start + lease.retry_timeout
if time.monotonic() >= connect_deadline:
elapsed = time.monotonic() - connect_start
raise ExporterUnreachableError(
f"Exporter {lease.exporter_name} unreachable after "
f"{lease.retry_timeout:.0f}s of retrying"
f"{elapsed:.0f}s of retrying: {unreachable}"
) from unreachable
logger.warning(
"Exporter %s is unreachable, releasing lease and retrying...",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1169,7 +1169,8 @@ async def fake_run(*_):

exc = find_exception_in_group(exc, ExporterUnreachableError)
assert exc is not None
assert "after 0s of retrying" in str(exc)
assert "test-exporter" in str(exc)
assert "unreachable" in str(exc).lower()
assert state["call_count"] >= 1

async def test_retries_when_wrapped_in_exception_group(self):
Expand Down
101 changes: 53 additions & 48 deletions python/packages/jumpstarter/jumpstarter/client/lease.py
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,7 @@ class Lease(ContextManagerMixin, AsyncContextManagerMixin):
) # Called when lease is ending
lease_ended: bool = field(default=False, init=False) # Set when lease expires naturally
lease_transferred: bool = field(default=False, init=False) # Set when lease is transferred to another client
_connected: bool = field(default=False, init=False) # True after first successful Dial

def __post_init__(self):
if hasattr(super(), "__post_init__"):
Expand Down Expand Up @@ -325,31 +326,36 @@ def __contextmanager__(self) -> Generator[Self]:
with self.portal.wrap_async_context_manager(self) as value:
yield value

async def _dial_with_retry(self):
"""Dial the controller with exponential backoff, waiting for the exporter to be ready.

Returns DialResponse on success.
Raises ExporterUnreachableError on timeout or unrecoverable error.
"""
logger.debug("Dialing controller for lease %s", self.name)
async def handle_async(self, stream): # noqa: C901
logger.debug("Connecting to Lease with name %s", self.name)
started = time.monotonic()
base_delay = 0.3
max_delay = 2.0
deadline = time.monotonic() + self.dial_timeout
dial_deadline = started + self.dial_timeout
# Short budget for initial connection (fast reassignment via shell's outer loop),
# full budget for mid-session reconnects (session worth preserving)
unavail_budget = self.retry_timeout if self._connected else self.dial_timeout
unavailable_deadline = started + unavail_budget if unavail_budget > 0 else None
attempt = 0
warned_unavailable = False
while True:
try:
return await self.controller.Dial(jumpstarter_pb2.DialRequest(lease_name=self.name))
response = await self.controller.Dial(jumpstarter_pb2.DialRequest(lease_name=self.name))
self._connected = True
break
except AioRpcError as e:
if e.code() == grpc.StatusCode.FAILED_PRECONDITION and "not ready" in str(e.details()):
remaining = deadline - time.monotonic()
remaining = dial_deadline - time.monotonic()
if remaining <= 0:
elapsed = time.monotonic() - started
logger.debug(
"Exporter not ready and dial timeout (%.1fs) exceeded after %d attempts",
"Exporter %s not ready and dial timeout (%.1fs) exceeded after %d attempts",
self.exporter_name,
self.dial_timeout,
attempt + 1,
)
raise ExporterUnreachableError(
f"Exporter {self.exporter_name} not ready after {self.dial_timeout:.0f}s"
f"Exporter {self.exporter_name} not ready after {elapsed:.0f}s: {e.details()}"
) from e
delay = min(base_delay * (2 ** min(attempt, 10)), max_delay, remaining)
logger.debug(
Expand All @@ -362,18 +368,33 @@ async def _dial_with_retry(self):
attempt += 1
continue
if e.code() == grpc.StatusCode.UNAVAILABLE:
remaining = deadline - time.monotonic()
if unavailable_deadline is None:
logger.warning("Exporter %s unavailable and retry disabled", self.exporter_name)
raise ExporterUnreachableError(
f"Exporter {self.exporter_name} unavailable (retry disabled): {e.details()}"
) from e
remaining = unavailable_deadline - time.monotonic()
if remaining <= 0:
elapsed = time.monotonic() - started
logger.warning(
"Exporter unavailable and dial timeout (%.1fs) exceeded after %d attempts",
self.dial_timeout,
"Exporter %s unavailable, retry budget (%.1fs) exceeded after %d attempts",
self.exporter_name,
unavail_budget,
attempt + 1,
)
raise ExporterUnreachableError(
f"Exporter {self.exporter_name} unavailable after {self.dial_timeout:.0f}s"
f"Exporter {self.exporter_name} unavailable after "
f"{elapsed:.0f}s of retrying: {e.details()}"
) from e
if not warned_unavailable:
warned_unavailable = True
logger.warning(
"Controller/exporter %s unavailable, retrying for %.0fs...",
self.exporter_name,
unavail_budget,
)
delay = min(base_delay * (2 ** min(attempt, 10)), max_delay, remaining)
logger.warning(
logger.info(
"Exporter unavailable, retrying Dial in %.1fs (attempt %d, %.1fs remaining)",
delay,
attempt + 1,
Expand All @@ -382,47 +403,31 @@ async def _dial_with_retry(self):
await sleep(delay)
attempt += 1
continue
# Exporter went offline or lease ended - raise immediately
if "permission denied" in str(e.details()).lower():
if e.code() == grpc.StatusCode.PERMISSION_DENIED:
self.lease_transferred = True
logger.warning(
"Lease %s has been transferred to another client. Your session is no longer valid.",
self.name,
)
raise ExporterUnreachableError(
f"Lease {self.name} transferred to another client"
f"Lease {self.name} has been transferred to another client"
) from e
logger.warning("Connection to exporter lost: %s", e.details())
raise ExporterUnreachableError(
f"Connection to exporter {self.exporter_name} lost: {e.details()}"
) from e

@asynccontextmanager
async def serve_unix_async(self):
# Wait for exporter readiness before accepting connections.
# The response is intentionally discarded — each connection needs
# its own Dial to get a unique router tunnel.
await self._dial_with_retry()

async def _tunnel_handler(stream):
try:
response = await self.controller.Dial(
jumpstarter_pb2.DialRequest(lease_name=self.name)
)
except AioRpcError as e:
raise ExporterUnreachableError(
f"Per-connection Dial failed for {self.exporter_name}: {e.details()}"
) from e
try:
async with connect_router_stream(
response.router_endpoint,
response.router_token,
stream,
self.tls_config,
self.grpc_options,
response.router_endpoint, response.router_token, stream, self.tls_config, self.grpc_options
):
pass
except grpc.aio.AioRpcError as e:
raise ExporterUnreachableError(
f"Router {response.router_endpoint} unreachable: {e.details()}"
) from e
except OSError as e:
raise ExporterUnreachableError(
f"Router {response.router_endpoint} connection failed: {e}"
) from e

async with TemporaryUnixListener(_tunnel_handler) as path:
@asynccontextmanager
async def serve_unix_async(self):
async with TemporaryUnixListener(self.handle_async) as path:
logger.debug("Serving Unix socket at %s", path)
yield path

Expand Down
Loading
Loading