From baf31733e906da9be6db0771e418cb194b2dfd4f Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 19:32:42 -0700 Subject: [PATCH] refactor(cron): AST-neutral bracket packing --- cron/blueprint_catalog.py | 4 +- cron/executions.py | 6 +-- cron/incidents.py | 6 +-- cron/lifecycle_guard.py | 10 +--- cron/monitor.py | 6 +-- cron/scheduler_provider.py | 83 +++++++--------------------------- cron/scripts/classify_items.py | 4 +- cron/suggestion_catalog.py | 3 +- cron/suggestions.py | 10 +--- 9 files changed, 28 insertions(+), 104 deletions(-) diff --git a/cron/blueprint_catalog.py b/cron/blueprint_catalog.py index 76a084d3f7..bd9a8a448c 100644 --- a/cron/blueprint_catalog.py +++ b/cron/blueprint_catalog.py @@ -719,9 +719,7 @@ def _resolve_schedule(blueprint: AutomationBlueprint, values: Dict[str, Any]) -> def fill_blueprint( - blueprint: AutomationBlueprint, - values: Dict[str, Any], - *, + blueprint: AutomationBlueprint, values: Dict[str, Any], *, origin: Optional[Dict[str, Any]] = None, ) -> Dict[str, Any]: """Validate ``values`` and return ``cron.jobs.create_job`` kwargs. diff --git a/cron/executions.py b/cron/executions.py index 054f995e1a..e0798bbb97 100644 --- a/cron/executions.py +++ b/cron/executions.py @@ -50,8 +50,7 @@ def prepare_ledger(conn: sqlite3.Connection, *, db_label: str, synchronous_full: @contextmanager def ledger_transaction( - lock: threading.RLock, - connect: Callable[[], sqlite3.Connection], + lock: threading.RLock, connect: Callable[[], sqlite3.Connection], initialize_schema: Callable[[sqlite3.Connection], None], ) -> Iterator[sqlite3.Connection]: """Open a connection, commit/rollback on exit, always close. ``sqlite3.Connection``'s own context @@ -249,8 +248,7 @@ def recover_interrupted_executions() -> int: def list_executions( - *, job_id: Optional[str] = None, limit: int = 50, - before_claimed_at: Optional[str] = None, + *, job_id: Optional[str] = None, limit: int = 50, before_claimed_at: Optional[str] = None, ) -> List[Dict[str, Any]]: """Return indexed, newest-first execution history with cursor pagination.""" clauses: List[str] = [] diff --git a/cron/incidents.py b/cron/incidents.py index b66f8c640b..e8ea8be4c0 100644 --- a/cron/incidents.py +++ b/cron/incidents.py @@ -137,11 +137,7 @@ def _classify_failure_type(error: str) -> str: def upsert_incident( - job_id: str, - error: str, - *, - job_name: Optional[str] = None, - failure_type: Optional[str] = None, + job_id: str, error: str, *, job_name: Optional[str] = None, failure_type: Optional[str] = None, output_file: Optional[str] = None, ) -> tuple[str, bool]: """Record (or refresh) the incident for ``job_id`` + ``error``; returns ``(incident_id, is_new)``. diff --git a/cron/lifecycle_guard.py b/cron/lifecycle_guard.py index 518b313aae..3b9add8b48 100644 --- a/cron/lifecycle_guard.py +++ b/cron/lifecycle_guard.py @@ -688,11 +688,7 @@ def _read_script_for_scanning(script_path: str) -> str: # --- recursive walk --------------------------------------------------------------------------- def _contains_unsafe_gateway_action( - command: str, - *, - cwd: Optional[str], - depth: int, - visited: set[Path], + command: str, *, cwd: Optional[str], depth: int, visited: set[Path], read_remote_script: Optional[_ReadRemoteScriptFn] = None, ) -> bool: if _direct_lifecycle_scan(command): @@ -735,9 +731,7 @@ def _contains_unsafe_gateway_action( def contains_gateway_lifecycle_command_or_referenced_script( - command: str, - *, - cwd: Optional[str] = None, + command: str, *, cwd: Optional[str] = None, read_remote_script: Optional[_ReadRemoteScriptFn] = None, ) -> bool: """Detect lifecycle/submit commands, including bounded nested scripts. diff --git a/cron/monitor.py b/cron/monitor.py index 9b82b6cc4c..7d63cd3bcb 100644 --- a/cron/monitor.py +++ b/cron/monitor.py @@ -49,11 +49,7 @@ def build_monitor_diff(old: str, new: str) -> str: """Unified diff of old vs new monitor output, capped at MAX_DIFF_CHARS.""" diff = "\n".join( difflib.unified_diff( - old.splitlines(), - new.splitlines(), - fromfile="previous", - tofile="current", - lineterm="", + old.splitlines(), new.splitlines(), fromfile="previous", tofile="current", lineterm="", ) ) if len(diff) > MAX_DIFF_CHARS: diff --git a/cron/scheduler_provider.py b/cron/scheduler_provider.py index 3b006d427f..6363c1d340 100644 --- a/cron/scheduler_provider.py +++ b/cron/scheduler_provider.py @@ -78,11 +78,7 @@ class CronScheduler(ABC): @abstractmethod def start( - self, - stop_event: threading.Event, - *, - adapters: Any = None, - loop: Any = None, + self, stop_event: threading.Event, *, adapters: Any = None, loop: Any = None, interval: int = 60, ) -> None: """Begin firing due jobs. Built-in BLOCKS until stop_event is set (run in a daemon thread); @@ -115,12 +111,7 @@ class CronScheduler(ABC): return provider_supports_force_fire(self) def fire_due( - self, - job_id: str, - *, - adapters: Any = None, - loop: Any = None, - force: bool = False, + self, job_id: str, *, adapters: Any = None, loop: Any = None, force: bool = False, ) -> bool: """Run one job NOW (inbound fire webhook entry). Store CAS claim (multi-machine at-most-once) then shared ``run_one_job``. True if THIS caller claimed and processed the @@ -144,8 +135,7 @@ class CronScheduler(ABC): claimed_job = claim_job_for_fire(job_id, **claim_kwargs) except BaseException as exc: finish_execution( - execution["id"], - success=False, + execution["id"], success=False, error=f"Fire claim failed before dispatch: {type(exc).__name__}: {exc}", ) raise @@ -156,11 +146,7 @@ class CronScheduler(ABC): return claimed_job def fire_claimed( - self, - claimed_job: dict, - *, - adapters: Any = None, - loop: Any = None, + self, claimed_job: dict, *, adapters: Any = None, loop: Any = None, cancel_event: Any = None, ) -> bool: """Run an exact ``claim_fire`` snapshot; ``cancel_event`` lets the transport stop it @@ -220,11 +206,7 @@ def _misfire_grace_minutes() -> float: def fire_overdue_jobs( - provider: "CronScheduler", - *, - adapters: Any = None, - loop: Any = None, - now: Any = None, + provider: "CronScheduler", *, adapters: Any = None, loop: Any = None, now: Any = None, ) -> int: """Misfire backstop (gateway housekeeping loop): fire jobs whose external HTTP fire never arrived, else ``next_run_at`` stays parked in the past forever. No-op for the built-in (its tick @@ -289,10 +271,8 @@ def fire_overdue_jobs( if claimed is None: continue threading.Thread( - target=provider.fire_claimed, - args=(claimed,), - kwargs={"adapters": adapters, "loop": loop}, - daemon=True, + target=provider.fire_claimed, args=(claimed,), + kwargs={"adapters": adapters, "loop": loop}, daemon=True, name=f"cron-misfire-{job_id[:12]}", ).start() fired += 1 @@ -356,17 +336,8 @@ class InProcessCronScheduler(CronScheduler): return "builtin" def start( - self, - stop_event, - *, - adapters=None, - loop=None, - interval=60, - can_dispatch=None, - profile_homes=None, - profile_adapters=None, - default_profile=None, - profile_gate=None, + self, stop_event, *, adapters=None, loop=None, interval=60, can_dispatch=None, + profile_homes=None, profile_adapters=None, default_profile=None, profile_gate=None, ): from cron.scheduler import CronTickYielded from cron.scheduler import tick as cron_tick @@ -377,15 +348,9 @@ class InProcessCronScheduler(CronScheduler): # Multiplex: tick EACH profile's store every cycle, heartbeats/recovery scoped per profile. if profile_homes: self._start_multiplex( - stop_event, - profile_homes=profile_homes, - adapters=adapters, - loop=loop, - interval=interval, - can_dispatch=can_dispatch, - profile_adapters=profile_adapters, - default_profile=default_profile, - profile_gate=profile_gate, + stop_event, profile_homes=profile_homes, adapters=adapters, loop=loop, + interval=interval, can_dispatch=can_dispatch, profile_adapters=profile_adapters, + default_profile=default_profile, profile_gate=profile_gate, ) return @@ -423,26 +388,15 @@ class InProcessCronScheduler(CronScheduler): stop_event.wait(_backoff_wait_seconds(interval, consecutive_failures)) def _start_multiplex( - self, - stop_event, - *, - profile_homes, - adapters=None, - loop=None, - interval=60, - can_dispatch=None, - profile_adapters=None, - default_profile=None, - profile_gate=None, + self, stop_event, *, profile_homes, adapters=None, loop=None, interval=60, + can_dispatch=None, profile_adapters=None, default_profile=None, profile_gate=None, ): """Tick every profile's store, each scoped via ``_profile_cron_scope``. ``profile_gate(name, home)``, when given, is consulted every cycle; a rejected profile is neither ticked nor heartbeated.""" from cron.scheduler import tick as cron_tick from cron.scheduler import ( - CronTickYielded, - SharedRouteAdapters, - _is_fd_exhaustion, + CronTickYielded, SharedRouteAdapters, _is_fd_exhaustion, _primary_profile_routes_for_current_home, ) from cron.jobs import clear_ticker_error, record_ticker_error, record_ticker_heartbeat @@ -496,11 +450,8 @@ class InProcessCronScheduler(CronScheduler): try: with _profile_cron_scope(home): cron_tick( - verbose=False, - adapters=tick_adapters_for(_pname), - loop=loop, - sync=False, - can_dispatch=can_dispatch, + verbose=False, adapters=tick_adapters_for(_pname), loop=loop, + sync=False, can_dispatch=can_dispatch, ) except CronTickYielded as e: # Yield for THIS profile only; one fresh gateway must not stop others. diff --git a/cron/scripts/classify_items.py b/cron/scripts/classify_items.py index 59c0db65cf..101cd128cc 100644 --- a/cron/scripts/classify_items.py +++ b/cron/scripts/classify_items.py @@ -130,9 +130,7 @@ def main() -> int: prompt = _build_prompt(items, args.criteria) try: resp = call_llm( - task="monitor", - messages=[{"role": "user", "content": prompt}], - max_tokens=1024, + task="monitor", messages=[{"role": "user", "content": prompt}], max_tokens=1024, temperature=0, ) content = resp.choices[0].message.content diff --git a/cron/suggestion_catalog.py b/cron/suggestion_catalog.py index 4cfeb8eea8..da1fbdc22b 100644 --- a/cron/suggestion_catalog.py +++ b/cron/suggestion_catalog.py @@ -116,8 +116,7 @@ CATALOG: List[CatalogEntry] = [ def seed_catalog_suggestions( - *, - add_fn: Optional[Callable[..., Optional[Dict[str, Any]]]] = None, + *, add_fn: Optional[Callable[..., Optional[Dict[str, Any]]]] = None, keys: Optional[List[str]] = None, ) -> List[Dict[str, Any]]: """Register catalog entries as pending suggestions. diff --git a/cron/suggestions.py b/cron/suggestions.py index 52352dac95..a7e84b32bc 100644 --- a/cron/suggestions.py +++ b/cron/suggestions.py @@ -111,12 +111,7 @@ def list_pending() -> List[Dict[str, Any]]: def add_suggestion( - *, - title: str, - description: str, - source: str, - job_spec: Dict[str, Any], - dedup_key: str, + *, title: str, description: str, source: str, job_spec: Dict[str, Any], dedup_key: str, ) -> Optional[Dict[str, Any]]: """Register a pending suggestion. Returns the record, or None when skipped: the same ``dedup_key`` was already decided on or is still pending (never re-offer, never duplicate), or the pending list @@ -197,8 +192,7 @@ def accept_suggestion(ref: str, *, origin: Optional[Dict[str, Any]] = None) -> O return None from cron.scheduler import ( - CronSchedulerRegistrationError, - create_job_with_scheduler_registration, + CronSchedulerRegistrationError, create_job_with_scheduler_registration, ) spec = dict(s.get("job_spec") or {})