Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
16 commits
Select commit Hold shift + click to select a range
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
Original file line number Diff line number Diff line change
Expand Up @@ -1601,6 +1601,17 @@ all subsequent phases have E2E coverage from the start.
- document tested pairs; W3C Trace Context in gRPC remains
best-effort across Python and Go (no OTel SDK requirement to
propagate `traceparent` where needed).
- **`operation` as a Prometheus label (Phase 2):** exporter series use
`@export` method names (`on`, `off`, `flash`, …) as the `operation`
label. The set is finite per process (loaded drivers), and unknown
methods are not recorded. If cardinality grows in the field, evaluate
moving `operation` from a series label to an exemplar key (see
*Cardinality guidelines*).
- **Per-chunk `jumpstarter_stream_bytes_total` increments (Phase 2):**
`copy_stream` calls `Counter.inc()` once per chunk. That is in-process
and uses independent tx/rx label sets, so Phase 2 accepts it. If flash
or storage throughput regresses, batch byte counts and flush
periodically instead of incrementing on every chunk.

## Rejected Alternatives

Expand Down Expand Up @@ -1669,6 +1680,11 @@ all subsequent phases have E2E coverage from the start.
- Event retention: Loki retention policy (per-tenant, per-stream retention
classes) for annotated log events (**DD-2**); whether Jumpstarter should
document recommended retention defaults or leave this to operators.
- Whether the exporter `operation` metric label should stay a bounded
series dimension or move to exemplars if cardinality grows (see
*Risks*).
- Whether `jumpstarter_stream_bytes_total` should batch per-chunk
increments if flash performance degrades (see *Risks*).

## Future Possibilities

Expand Down
136 changes: 92 additions & 44 deletions python/packages/jumpstarter-cli/jumpstarter_cli/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,8 +11,15 @@
from jumpstarter_cli_common.config import opt_config
from jumpstarter_cli_common.exceptions import handle_exceptions

from jumpstarter.metrics import start_metrics_server

logger = logging.getLogger(__name__)

# Phase 2 interim: always expose local HTTP /metrics on ephemeral loopback.
# Phase 3 replaces this with Telemetry reverse-scrape (unix/memory or in-process);
# bind address is intentionally not a user-facing CLI option.
_METRICS_BIND_ADDRESS = ":0"


def _parse_listener_bind(value: str) -> tuple[str, int]:
"""Parse '[host:]port' into (host, port). Default host is 0.0.0.0."""
Expand Down Expand Up @@ -70,7 +77,14 @@ def _reap_zombie_processes(capture_child=None):
logger.warning(f"PARENT: Error during zombie reaping: {e}")


def _handle_child(config, parsed_bind=None, tls_insecure=False, tls_cert=None, tls_key=None, passphrase=None): # noqa: C901
def _handle_child( # noqa: C901
config,
parsed_bind=None,
tls_insecure=False,
tls_cert=None,
tls_key=None,
passphrase=None,
):
"""Handle child process with graceful shutdown."""
async def serve_with_graceful_shutdown(): # noqa: C901
received_signal = 0
Expand Down Expand Up @@ -99,46 +113,53 @@ async def signal_handler():
# Start signal handler immediately
signal_tg.start_soon(signal_handler)

if parsed_bind is not None:
host, port = parsed_bind
tls_credentials = None
if tls_insecure:
if passphrase:
click.echo(
"WARNING: --passphrase has no effect without TLS; "
"the passphrase will be transmitted in plaintext",
err=True,
)
elif tls_cert and tls_key:
tls_credentials = _tls_server_credentials(tls_cert, tls_key)

interceptors = None
if passphrase:
from jumpstarter.exporter.auth import PassphraseInterceptor
interceptors = [PassphraseInterceptor(passphrase)]

exporter_exit_code = None
async with config.create_exporter(standalone=True) as exporter:
try:
await exporter.serve_standalone_tcp(
host, port,
tls_credentials=tls_credentials,
interceptors=interceptors,
)
except* Exception as excgroup:
_handle_exporter_exceptions(excgroup)
exporter_exit_code = exporter.exit_code
else:
# Create exporter and run it (controller mode)
exporter_exit_code = None
async with config.create_exporter() as exporter:
try:
await exporter.serve()
except* Exception as excgroup:
_handle_exporter_exceptions(excgroup)
listen_addr, shutdown_metrics = start_metrics_server(_METRICS_BIND_ADDRESS)
logger.info("Serving metrics server at http://%s/metrics", listen_addr)

# Check if exporter set an exit code (e.g., from hook failure with on_failure='exit')
exporter_exit_code = exporter.exit_code
try:
if parsed_bind is not None:
host, port = parsed_bind
tls_credentials = None
if tls_insecure:
if passphrase:
click.echo(
"WARNING: --passphrase has no effect without TLS; "
"the passphrase will be transmitted in plaintext",
err=True,
)
elif tls_cert and tls_key:
tls_credentials = _tls_server_credentials(tls_cert, tls_key)

interceptors = None
if passphrase:
from jumpstarter.exporter.auth import PassphraseInterceptor
interceptors = [PassphraseInterceptor(passphrase)]

exporter_exit_code = None
async with config.create_exporter(standalone=True) as exporter:
try:
await exporter.serve_standalone_tcp(
host, port,
tls_credentials=tls_credentials,
interceptors=interceptors,
)
except* Exception as excgroup:
_handle_exporter_exceptions(excgroup)
exporter_exit_code = exporter.exit_code
else:
# Create exporter and run it (controller mode)
exporter_exit_code = None
async with config.create_exporter() as exporter:
try:
await exporter.serve()
except* Exception as excgroup:
_handle_exporter_exceptions(excgroup)

# Check if exporter set an exit code (e.g., from hook failure with on_failure='exit')
exporter_exit_code = exporter.exit_code
finally:
if shutdown_metrics is not None:
shutdown_metrics()

# Cancel the signal handler after exporter completes
signal_tg.cancel_scope.cancel()
Expand Down Expand Up @@ -204,7 +225,12 @@ def parent_signal_handler(signum, _):


def _serve_with_exc_handling(
config, parsed_bind=None, tls_insecure=False, tls_cert=None, tls_key=None, passphrase=None
config,
parsed_bind=None,
tls_insecure=False,
tls_cert=None,
tls_key=None,
passphrase=None,
):
max_rapid_failures = config.failure_detection.max_rapid_failures
rapid_failure_window = config.failure_detection.rapid_failure_window
Expand Down Expand Up @@ -253,7 +279,14 @@ def _serve_with_exc_handling(
rapid_failure_count = 0
else:
os.setsid() # Become group leader so all spawned subprocesses are reached by parent's signals
_handle_child(config, parsed_bind, tls_insecure, tls_cert, tls_key, passphrase)
_handle_child(
config,
parsed_bind,
tls_insecure,
tls_cert,
tls_key,
passphrase,
)
sys.exit(1) # should never happen


Expand Down Expand Up @@ -295,7 +328,15 @@ def _serve_with_exc_handling(
help="Exit after the current lease ends instead of waiting for a new one.",
)
@handle_exceptions
def run(config, listener_bind, tls_insecure, tls_cert, tls_key, passphrase, exit_on_lease_end):
def run(
config,
listener_bind,
tls_insecure,
tls_cert,
tls_key,
passphrase,
exit_on_lease_end,
):
"""Run an exporter locally."""
if listener_bind is not None and config is None:
raise click.UsageError("--exporter-config (or --exporter) is required when using --tls-grpc-listener")
Expand All @@ -313,4 +354,11 @@ def run(config, listener_bind, tls_insecure, tls_cert, tls_key, passphrase, exit
if exit_on_lease_end:
config.exit_on_lease_end = True
parsed_bind = _parse_listener_bind(listener_bind) if listener_bind is not None else None
return _serve_with_exc_handling(config, parsed_bind, tls_insecure, tls_cert, tls_key, passphrase)
return _serve_with_exc_handling(
config,
parsed_bind,
tls_insecure,
tls_cert,
tls_key,
passphrase,
)
Loading
Loading