Spaces:
Sleeping
Sleeping
Claude Code Claude Opus 4.6 commited on
Commit ·
57701c5
1
Parent(s): 45a307e
Claude Code: Add import-time worker bootstrap for guaranteed startup
Browse filesDual-trigger mechanism:
1. Import-time bootstrap: Worker starts immediately when app.py loads
2. Startup event failsafe: Redundant backup in case bootstrap fails
- Run asyncio polling in daemon thread with its own event loop
- Hard-write worker_active=True immediately at import time
- Startup event now a failsafe wrapper that checks bootstrap status
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
app.py
CHANGED
|
@@ -35,6 +35,94 @@ _worker_stop_event = None
|
|
| 35 |
_polling_task = None
|
| 36 |
_polling_lock = threading.Lock()
|
| 37 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 38 |
|
| 39 |
def _run_worker_thread():
|
| 40 |
"""
|
|
@@ -134,39 +222,47 @@ async def _poll_a2a_ready_state():
|
|
| 134 |
|
| 135 |
@app.on_event("startup")
|
| 136 |
async def startup_event():
|
| 137 |
-
"""
|
| 138 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 139 |
|
| 140 |
-
# 1. Verbose startup logging - prove handler ran
|
| 141 |
logger.info("=" * 60)
|
| 142 |
-
logger.info("[STARTUP] ========== STARTUP EVENT
|
| 143 |
logger.info(f"[STARTUP] Timestamp: {datetime.utcnow().isoformat()}")
|
| 144 |
logger.info(f"[STARTUP] PID: {os.getpid()}")
|
|
|
|
| 145 |
logger.info("=" * 60)
|
| 146 |
|
| 147 |
-
#
|
| 148 |
-
|
| 149 |
-
|
| 150 |
-
logger.info("[Cain] Worker disabled by WORKER_MODE=disabled")
|
| 151 |
else:
|
| 152 |
-
logger.info("[
|
| 153 |
-
|
| 154 |
-
|
| 155 |
-
|
| 156 |
-
|
| 157 |
-
|
| 158 |
-
|
| 159 |
-
|
| 160 |
-
|
| 161 |
-
|
| 162 |
-
|
| 163 |
-
|
| 164 |
-
|
| 165 |
-
|
| 166 |
-
|
| 167 |
-
|
| 168 |
-
|
| 169 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
| 170 |
|
| 171 |
|
| 172 |
@app.on_event("shutdown")
|
|
|
|
| 35 |
_polling_task = None
|
| 36 |
_polling_lock = threading.Lock()
|
| 37 |
|
| 38 |
+
# ========== IMMEDIATE BOOTSTRAP (Module Import Time) ==========
|
| 39 |
+
# This fires the moment app.py is imported, BEFORE FastAPI starts
|
| 40 |
+
# Bypasses unreliable startup event by starting worker immediately
|
| 41 |
+
_bootstrap_worker_active = False
|
| 42 |
+
|
| 43 |
+
|
| 44 |
+
def _run_asyncio_polling_in_thread():
|
| 45 |
+
"""
|
| 46 |
+
Run the asyncio polling task in its own event loop within a daemon thread.
|
| 47 |
+
This ensures the polling loop starts immediately at module import time.
|
| 48 |
+
"""
|
| 49 |
+
global _bootstrap_worker_active
|
| 50 |
+
|
| 51 |
+
logger.info("=" * 60)
|
| 52 |
+
logger.info("[BOOTSTRAP] ========== IMPORT-TIME WORKER START ==========")
|
| 53 |
+
logger.info(f"[BOOTSTRAP] Timestamp: {datetime.utcnow().isoformat()}")
|
| 54 |
+
logger.info(f"[BOOTSTRAP] PID: {os.getpid()}")
|
| 55 |
+
logger.info(f"[BOOTSTRAP] Thread: {threading.current_thread().name}")
|
| 56 |
+
logger.info("=" * 60)
|
| 57 |
+
|
| 58 |
+
# Create new event loop for this thread
|
| 59 |
+
loop = asyncio.new_event_loop()
|
| 60 |
+
asyncio.set_event_loop(loop)
|
| 61 |
+
|
| 62 |
+
try:
|
| 63 |
+
# Run the polling task - this blocks forever in while True
|
| 64 |
+
loop.run_until_complete(_poll_a2a_ready_state())
|
| 65 |
+
except Exception as e:
|
| 66 |
+
logger.error(f"[BOOTSTRAP] Polling loop failed: {e}")
|
| 67 |
+
finally:
|
| 68 |
+
loop.close()
|
| 69 |
+
|
| 70 |
+
|
| 71 |
+
def _bootstrap_immediate_worker():
|
| 72 |
+
"""Immediately bootstrap the worker at module import time."""
|
| 73 |
+
global _bootstrap_worker_active
|
| 74 |
+
|
| 75 |
+
# Check WORKER_MODE before bootstrapping
|
| 76 |
+
worker_mode = os.environ.get("WORKER_MODE", "auto").lower()
|
| 77 |
+
if worker_mode == "disabled":
|
| 78 |
+
logger.info("[BOOTSTRAP] Worker disabled by WORKER_MODE=disabled")
|
| 79 |
+
return
|
| 80 |
+
|
| 81 |
+
logger.info("[BOOTSTRAP] Starting immediate daemon thread for asyncio polling...")
|
| 82 |
+
|
| 83 |
+
# Start daemon thread with its own asyncio event loop
|
| 84 |
+
bootstrap_thread = threading.Thread(
|
| 85 |
+
target=_run_asyncio_polling_in_thread,
|
| 86 |
+
daemon=True,
|
| 87 |
+
name="bootstrap-asyncio-polling"
|
| 88 |
+
)
|
| 89 |
+
bootstrap_thread.start()
|
| 90 |
+
|
| 91 |
+
_bootstrap_worker_active = True
|
| 92 |
+
logger.info(f"[BOOTSTRAP] SUCCESS: Daemon thread started: {bootstrap_thread.name}")
|
| 93 |
+
|
| 94 |
+
# 3. Hard-Write State: Immediately set worker_active in status file
|
| 95 |
+
data_dir = Path(os.environ.get("OPENCLAW_DATA_DIR", "/data"))
|
| 96 |
+
data_dir.mkdir(parents=True, exist_ok=True)
|
| 97 |
+
status_file = data_dir / "cain_status.json"
|
| 98 |
+
|
| 99 |
+
try:
|
| 100 |
+
# Update or create status file with worker_active=True
|
| 101 |
+
if status_file.exists():
|
| 102 |
+
with open(status_file, "r") as f:
|
| 103 |
+
data = json.load(f)
|
| 104 |
+
else:
|
| 105 |
+
data = {}
|
| 106 |
+
|
| 107 |
+
data["worker_active"] = True
|
| 108 |
+
data["_worker_bootstrap_pid"] = os.getpid()
|
| 109 |
+
data["_worker_bootstrap_time"] = datetime.utcnow().isoformat() + "+00:00"
|
| 110 |
+
data["_worker_bootstrap_mode"] = "import_time_daemon"
|
| 111 |
+
|
| 112 |
+
with open(status_file, "w") as f:
|
| 113 |
+
json.dump(data, f, indent=2)
|
| 114 |
+
|
| 115 |
+
logger.info(f"[BOOTSTRAP] Hard-wrote worker_active=True to {status_file}")
|
| 116 |
+
except Exception as e:
|
| 117 |
+
logger.warning(f"[BOOTSTRAP] Failed to hard-write status: {e}")
|
| 118 |
+
|
| 119 |
+
|
| 120 |
+
# ========== EXECUTE IMMEDIATE BOOTSTRAP ==========
|
| 121 |
+
# This runs IMMEDIATELY when app.py is imported, before FastAPI even starts
|
| 122 |
+
_bootstrap_immediate_worker()
|
| 123 |
+
logger.info("[BOOTSTRAP] Import-time bootstrap complete")
|
| 124 |
+
# ========== END IMMEDIATE BOOTSTRAP ==========
|
| 125 |
+
|
| 126 |
|
| 127 |
def _run_worker_thread():
|
| 128 |
"""
|
|
|
|
| 222 |
|
| 223 |
@app.on_event("startup")
|
| 224 |
async def startup_event():
|
| 225 |
+
"""
|
| 226 |
+
FAILSAFE STARTUP: Backup worker initialization.
|
| 227 |
+
|
| 228 |
+
The worker should already be running from import-time bootstrap.
|
| 229 |
+
This event is a redundant failsafe in case the bootstrap failed.
|
| 230 |
+
"""
|
| 231 |
+
global _worker_thread, _worker_stop_event, _polling_task, _bootstrap_worker_active
|
| 232 |
|
|
|
|
| 233 |
logger.info("=" * 60)
|
| 234 |
+
logger.info("[STARTUP] ========== FAILSAFE STARTUP EVENT ==========")
|
| 235 |
logger.info(f"[STARTUP] Timestamp: {datetime.utcnow().isoformat()}")
|
| 236 |
logger.info(f"[STARTUP] PID: {os.getpid()}")
|
| 237 |
+
logger.info(f"[STARTUP] Bootstrap worker active? {_bootstrap_worker_active}")
|
| 238 |
logger.info("=" * 60)
|
| 239 |
|
| 240 |
+
# 2. Failsafe Wrapper: Try to start worker if bootstrap didn't
|
| 241 |
+
if _bootstrap_worker_active:
|
| 242 |
+
logger.info("[STARTUP] Bootstrap already started worker - this event is redundant")
|
|
|
|
| 243 |
else:
|
| 244 |
+
logger.info("[STARTUP] Bootstrap did NOT start worker - attempting failsafe startup...")
|
| 245 |
+
|
| 246 |
+
# Check WORKER_MODE before starting worker
|
| 247 |
+
worker_mode = os.environ.get("WORKER_MODE", "auto").lower()
|
| 248 |
+
if worker_mode == "disabled":
|
| 249 |
+
logger.info("[Cain] Worker disabled by WORKER_MODE=disabled")
|
| 250 |
+
return
|
| 251 |
+
|
| 252 |
+
try:
|
| 253 |
+
logger.info("[STARTUP] Starting persistent background worker thread...")
|
| 254 |
+
_worker_stop_event = threading.Event()
|
| 255 |
+
_worker_thread = threading.Thread(target=_run_worker_thread, daemon=True)
|
| 256 |
+
_worker_thread.start()
|
| 257 |
+
logger.info(f"[STARTUP] Worker thread started: {_worker_thread.name}")
|
| 258 |
+
|
| 259 |
+
logger.info("[STARTUP] About to create asyncio task for A2A polling...")
|
| 260 |
+
_polling_task = asyncio.create_task(_poll_a2a_ready_state())
|
| 261 |
+
logger.info(f"[STARTUP] SUCCESS: asyncio.create_task() returned task object: {_polling_task}")
|
| 262 |
+
except Exception as e:
|
| 263 |
+
logger.error(f"[STARTUP] FAILED: Failsafe startup raised exception: {e}")
|
| 264 |
+
import traceback
|
| 265 |
+
traceback.print_exc()
|
| 266 |
|
| 267 |
|
| 268 |
@app.on_event("shutdown")
|