refactor(cron): inline the multiplex-context span into _launch_external_cron_worker
The `_inner` wrapper existed only to set and reset the multiplex contextvar around the whole launch. Every consumer of that context (the payload flag, `hydrate_profile_secret_sources`, `strip_launch_profile_env`, the scrubbed `build_subprocess_env`) runs before the worker is spawned, so the span now covers exactly the handoff-environment build and ends before `Popen` and the acknowledgement wait. One function, one try/finally, same behaviour.
This commit is contained in:
@@ -3162,18 +3162,6 @@ def _launch_external_cron_worker(job: dict) -> bool:
|
||||
and built the worker environment with the launch profile's residue and no scrub (review on
|
||||
f5f88d5058). Enable it for exactly this span; the worker then re-establishes it from the payload.
|
||||
"""
|
||||
from agent.secret_scope import is_multiplex_active, reset_multiplex_context, set_multiplex_context
|
||||
from cron.scheduler_provider import routed_profile_fire
|
||||
|
||||
context_token = set_multiplex_context(True) if routed_profile_fire() and not is_multiplex_active() else None
|
||||
try:
|
||||
return _launch_external_cron_worker_inner(job)
|
||||
finally:
|
||||
if context_token is not None:
|
||||
reset_multiplex_context(context_token)
|
||||
|
||||
|
||||
def _launch_external_cron_worker_inner(job: dict) -> bool:
|
||||
execution_id = str(job["execution_id"])
|
||||
job_id = str(job["id"])
|
||||
handoff_dir = _get_hermes_home() / "cron" / "external-workers"
|
||||
@@ -3192,9 +3180,12 @@ def _launch_external_cron_worker_inner(job: dict) -> bool:
|
||||
from agent.secret_scope import (
|
||||
build_profile_secret_scope,
|
||||
is_multiplex_active,
|
||||
reset_multiplex_context,
|
||||
reset_secret_scope,
|
||||
set_multiplex_context,
|
||||
set_secret_scope,
|
||||
)
|
||||
from cron.scheduler_provider import routed_profile_fire
|
||||
from hermes_cli.env_loader import hydrate_profile_secret_sources
|
||||
from tools.environments.local import build_subprocess_env, restore_managed_env, strip_launch_profile_env
|
||||
from tools.process_registry import (
|
||||
@@ -3202,59 +3193,66 @@ def _launch_external_cron_worker_inner(job: dict) -> bool:
|
||||
systemd_user_bus_env,
|
||||
)
|
||||
|
||||
try:
|
||||
require_restart_safe_scope = bool(
|
||||
(load_config_readonly().get("cron") or {}).get("require_restart_safe_scope", False)
|
||||
)
|
||||
except Exception:
|
||||
require_restart_safe_scope = False
|
||||
multiplex_active = is_multiplex_active()
|
||||
dispatch = restart_safe_gateway_child_argv(
|
||||
command,
|
||||
unit_suffix=f"cron-{job_id}-exec-{execution_id}",
|
||||
require_restart_safe_scope=require_restart_safe_scope,
|
||||
context_token = (
|
||||
set_multiplex_context(True) if routed_profile_fire() and not is_multiplex_active() else None
|
||||
)
|
||||
if dispatch.mode == "in_process":
|
||||
return False
|
||||
|
||||
if mark_execution_handoff_pending(execution_id) is None:
|
||||
raise RuntimeError(
|
||||
"cron execution claim changed before external worker handoff"
|
||||
)
|
||||
|
||||
_ensure_cron_dir(handoff_dir)
|
||||
try:
|
||||
handoff_dir.chmod(0o700)
|
||||
except OSError:
|
||||
pass
|
||||
fd = os.open(payload_path, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600)
|
||||
try:
|
||||
with os.fdopen(fd, "w", encoding="utf-8") as payload_file:
|
||||
json.dump(
|
||||
{
|
||||
"job": job,
|
||||
"profile_home": str(_get_hermes_home().resolve()),
|
||||
"multiplex_active": multiplex_active,
|
||||
},
|
||||
payload_file,
|
||||
try:
|
||||
require_restart_safe_scope = bool(
|
||||
(load_config_readonly().get("cron") or {}).get("require_restart_safe_scope", False)
|
||||
)
|
||||
payload_file.flush()
|
||||
os.fsync(payload_file.fileno())
|
||||
except BaseException:
|
||||
payload_path.unlink(missing_ok=True)
|
||||
raise
|
||||
except Exception:
|
||||
require_restart_safe_scope = False
|
||||
multiplex_active = is_multiplex_active()
|
||||
dispatch = restart_safe_gateway_child_argv(
|
||||
command,
|
||||
unit_suffix=f"cron-{job_id}-exec-{execution_id}",
|
||||
require_restart_safe_scope=require_restart_safe_scope,
|
||||
)
|
||||
if dispatch.mode == "in_process":
|
||||
return False
|
||||
|
||||
profile_home = _get_hermes_home().resolve()
|
||||
hydrate_profile_secret_sources(profile_home)
|
||||
secret_token = set_secret_scope(build_profile_secret_scope(profile_home))
|
||||
try:
|
||||
worker_env = restore_managed_env(strip_launch_profile_env(build_subprocess_env(
|
||||
scrub_secrets=multiplex_active,
|
||||
inherit_profile_home=True,
|
||||
extra={"HERMES_HOME": str(profile_home)},
|
||||
)))
|
||||
if mark_execution_handoff_pending(execution_id) is None:
|
||||
raise RuntimeError(
|
||||
"cron execution claim changed before external worker handoff"
|
||||
)
|
||||
|
||||
_ensure_cron_dir(handoff_dir)
|
||||
try:
|
||||
handoff_dir.chmod(0o700)
|
||||
except OSError:
|
||||
pass
|
||||
fd = os.open(payload_path, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600)
|
||||
try:
|
||||
with os.fdopen(fd, "w", encoding="utf-8") as payload_file:
|
||||
json.dump(
|
||||
{
|
||||
"job": job,
|
||||
"profile_home": str(_get_hermes_home().resolve()),
|
||||
"multiplex_active": multiplex_active,
|
||||
},
|
||||
payload_file,
|
||||
)
|
||||
payload_file.flush()
|
||||
os.fsync(payload_file.fileno())
|
||||
except BaseException:
|
||||
payload_path.unlink(missing_ok=True)
|
||||
raise
|
||||
|
||||
profile_home = _get_hermes_home().resolve()
|
||||
hydrate_profile_secret_sources(profile_home)
|
||||
secret_token = set_secret_scope(build_profile_secret_scope(profile_home))
|
||||
try:
|
||||
worker_env = restore_managed_env(strip_launch_profile_env(build_subprocess_env(
|
||||
scrub_secrets=multiplex_active,
|
||||
inherit_profile_home=True,
|
||||
extra={"HERMES_HOME": str(profile_home)},
|
||||
)))
|
||||
finally:
|
||||
reset_secret_scope(secret_token)
|
||||
finally:
|
||||
reset_secret_scope(secret_token)
|
||||
if context_token is not None:
|
||||
reset_multiplex_context(context_token)
|
||||
worker_env = systemd_user_bus_env(worker_env)
|
||||
try:
|
||||
process = subprocess.Popen(
|
||||
|
||||
Reference in New Issue
Block a user