Coverage for agentos/server/daemon.py: 21%
346 statements
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-06 10:59 +0800
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-06 10:59 +0800
1"""
2AgentOS Server Daemon — Independent production server wrapper.
4Provides:
5 - PID file management (start/stop/status/restart)
6 - Signal handling (SIGTERM/SIGINT) with graceful shutdown
7 - Background asyncio task queue
8 - Structured file logging with rotation
9 - Health check endpoint (/healthz auto-mounted)
10 - Configuration via environment variables
11 - Memory persistence (v1.14.9): auto-save on shutdown, auto-load on startup
13Usage:
14 agentos-daemon start # start as daemon
15 agentos-daemon stop # stop running daemon
16 agentos-daemon status # check if daemon is running
17 agentos-daemon restart # stop + start
18 agentos-daemon run # run in foreground (debug)
20Environment:
21 AGENTOS_DAEMON_HOST=0.0.0.0
22 AGENTOS_DAEMON_PORT=8910
23 AGENTOS_DAEMON_PIDFILE=~/.agentos/daemon.pid
24 AGENTOS_DAEMON_LOGFILE=~/.agentos/daemon.log
25 AGENTOS_DAEMON_WORKERS=4
26 AGENTOS_DAEMON_TIMEOUT=30 # graceful shutdown timeout seconds
27 AGENTOS_DAEMON_LOG_LEVEL=info
28 AGENTOS_DAEMON_MEMORY_DIR=~/.agentos/memory # memory persistence dir
29"""
31from __future__ import annotations
33import asyncio
34import atexit
35import logging
36import os
37import signal
38import sys
39import time
40from contextlib import asynccontextmanager
41from dataclasses import dataclass, field
42from pathlib import Path
43from typing import Callable, Any
45import uvicorn
46from fastapi import FastAPI, APIRouter
47from fastapi.middleware.cors import CORSMiddleware
49from agentos.memory.persistence import MemoryPersistenceManager
51__all__ = [
52 "ServerDaemon",
53 "DaemonConfig",
54 "BackgroundTask",
55 "get_daemon",
56 "create_daemon_app",
57 "daemon_main",
58]
60# ── Config ──────────────────────────────────────────────
63@dataclass
64class DaemonConfig:
65 """Server daemon configuration."""
67 host: str = "0.0.0.0"
68 port: int = 8910
69 pidfile: str = "~/.agentos/daemon.pid"
70 logfile: str = "~/.agentos/daemon.log"
71 workers: int = 1
72 shutdown_timeout: float = 30.0
73 log_level: str = "info"
74 log_max_bytes: int = 10 * 1024 * 1024 # 10 MB
75 log_backup_count: int = 5
76 memory_dir: str = "~/.agentos/memory"
77 persist_memory: bool = True
79 @classmethod
80 def from_env(cls) -> "DaemonConfig":
81 """Load config from environment variables."""
82 return cls(
83 host=os.getenv("AGENTOS_DAEMON_HOST", "0.0.0.0"),
84 port=int(os.getenv("AGENTOS_DAEMON_PORT", "8910")),
85 pidfile=os.path.expanduser(
86 os.getenv("AGENTOS_DAEMON_PIDFILE", "~/.agentos/daemon.pid")
87 ),
88 logfile=os.path.expanduser(
89 os.getenv("AGENTOS_DAEMON_LOGFILE", "~/.agentos/daemon.log")
90 ),
91 workers=int(os.getenv("AGENTOS_DAEMON_WORKERS", "1")),
92 shutdown_timeout=float(os.getenv("AGENTOS_DAEMON_TIMEOUT", "30")),
93 log_level=os.getenv("AGENTOS_DAEMON_LOG_LEVEL", "info"),
94 memory_dir=os.path.expanduser(
95 os.getenv("AGENTOS_DAEMON_MEMORY_DIR", "~/.agentos/memory")
96 ),
97 persist_memory=os.getenv("AGENTOS_DAEMON_PERSIST_MEMORY", "true").lower()
98 not in ("0", "false", "no"),
99 )
102# ── Background Task Queue ──────────────────────────────
105@dataclass
106class BackgroundTask:
107 """A tracked background task."""
109 task_id: str
110 name: str
111 created_at: float = field(default_factory=time.time)
112 status: str = "pending"
113 result: Any = None
114 error: str | None = None
116 def to_dict(self) -> dict:
117 return {
118 "task_id": self.task_id,
119 "name": self.name,
120 "created_at": self.created_at,
121 "status": self.status,
122 "result": str(self.result)[:500] if self.result is not None else None,
123 "error": self.error,
124 }
127class TaskQueue:
128 """Minimal async background task queue (in-process, no external broker)."""
130 def __init__(self, max_history: int = 1000) -> None:
131 self._tasks: dict[str, BackgroundTask] = {}
132 self._max_history = max_history
133 self._asyncio_tasks: dict[str, asyncio.Task] = {}
134 self._counter = 0
136 def submit(self, name: str, coro) -> str:
137 """Submit a coroutine for background execution. Returns task_id."""
138 self._counter += 1
139 tid = f"task_{self._counter}_{int(time.time())}"
140 bt = BackgroundTask(task_id=tid, name=name, status="running")
141 bt.created_at = time.time()
142 self._tasks[tid] = bt
144 async def _runner():
145 try:
146 bt.result = await coro
147 bt.status = "completed"
148 except Exception as exc:
149 bt.error = str(exc)
150 bt.status = "failed"
152 self._asyncio_tasks[tid] = asyncio.create_task(_runner())
153 self._prune_history()
154 return tid
156 def get(self, task_id: str) -> BackgroundTask | None:
157 return self._tasks.get(task_id)
159 def list_tasks(self, limit: int = 50) -> list[BackgroundTask]:
160 items = sorted(
161 self._tasks.values(), key=lambda t: t.created_at, reverse=True
162 )
163 return items[:limit]
165 def active_count(self) -> int:
166 return sum(1 for t in self._tasks.values() if t.status == "running")
168 def _prune_history(self) -> None:
169 if len(self._tasks) > self._max_history:
170 completed = [
171 tid for tid, t in self._tasks.items() if t.status in ("completed", "failed")
172 ]
173 overflow = len(self._tasks) - self._max_history
174 for tid in completed[:overflow]:
175 del self._tasks[tid]
176 self._asyncio_tasks.pop(tid, None)
178 async def shutdown(self, timeout: float = 10.0) -> None:
179 """Cancel all running tasks with a timeout."""
180 if not self._asyncio_tasks:
181 return
182 for t in list(self._asyncio_tasks.values()):
183 if not t.done():
184 t.cancel()
185 try:
186 await asyncio.wait_for(
187 asyncio.gather(*self._asyncio_tasks.values(), return_exceptions=True),
188 timeout=timeout,
189 )
190 except asyncio.TimeoutError:
191 pass
194# ── Daemon Core ────────────────────────────────────────
197def _setup_file_logging(config: DaemonConfig) -> logging.Logger:
198 """Configure rotating file logger."""
199 from logging.handlers import RotatingFileHandler
201 log_dir = Path(config.logfile).parent
202 log_dir.mkdir(parents=True, exist_ok=True)
204 logger = logging.getLogger("agentos.daemon")
205 logger.setLevel(getattr(logging, config.log_level.upper(), logging.INFO))
206 logger.propagate = False
208 # Remove existing handlers
209 for h in list(logger.handlers):
210 logger.removeHandler(h)
212 handler = RotatingFileHandler(
213 config.logfile,
214 maxBytes=config.log_max_bytes,
215 backupCount=config.log_backup_count,
216 encoding="utf-8",
217 )
218 handler.setFormatter(
219 logging.Formatter(
220 "%(asctime)s [%(levelname)s] %(name)s: %(message)s",
221 datefmt="%Y-%m-%d %H:%M:%S",
222 )
223 )
224 logger.addHandler(handler)
225 return logger
228class ServerDaemon:
229 """Independent server daemon wrapping any ASGI app.
231 Lifecycle:
232 daemon = ServerDaemon(config, app_factory)
233 daemon.start() # daemonize
234 daemon.stop() # send SIGTERM
235 daemon.status() # read PID file
236 daemon.restart() # stop + start
237 """
239 def __init__(
240 self,
241 config: DaemonConfig | None = None,
242 app_factory: Callable[[], FastAPI] | None = None,
243 ) -> None:
244 self.config = config or DaemonConfig.from_env()
245 self._app_factory = app_factory or self._default_app_factory
246 self._logger = _setup_file_logging(self.config)
247 self.task_queue = TaskQueue()
248 self._started_at: float | None = None
249 self._shutdown_event: asyncio.Event | None = None
251 # Memory persistence (v1.14.9)
252 self._persistence_mgr = MemoryPersistenceManager(
253 base_dir=self.config.memory_dir,
254 compress=True,
255 )
256 self._memory_objects: dict[str, Any] = {
257 "pyramid": None,
258 "working": None,
259 "conversation": None,
260 "long_term": None,
261 "reflection_engine": None,
262 "consolidation_pipeline": None,
263 }
265 # ── Memory Persistence API (v1.14.9) ──────
267 def register_memory(
268 self,
269 *,
270 pyramid: Any = None,
271 working: Any = None,
272 conversation: Any = None,
273 long_term: Any = None,
274 reflection_engine: Any = None,
275 consolidation_pipeline: Any = None,
276 ) -> None:
277 """Register memory objects for crash-safe persistence.
279 Objects must implement get_state() and restore_state().
280 On daemon shutdown, all registered objects are automatically saved
281 to disk. On daemon startup, they are automatically restored.
282 """
283 if pyramid is not None:
284 self._memory_objects["pyramid"] = pyramid
285 if working is not None:
286 self._memory_objects["working"] = working
287 if conversation is not None:
288 self._memory_objects["conversation"] = conversation
289 if long_term is not None:
290 self._memory_objects["long_term"] = long_term
291 if reflection_engine is not None:
292 self._memory_objects["reflection_engine"] = reflection_engine
293 if consolidation_pipeline is not None:
294 self._memory_objects["consolidation_pipeline"] = consolidation_pipeline
296 self._logger.info(
297 f"Registered {sum(1 for v in self._memory_objects.values() if v is not None)} "
298 f"memory subsystems for persistence"
299 )
301 async def save_memory_snapshot(self) -> str | None:
302 """Save current memory state to disk. Returns snapshot path or None."""
303 if not self.config.persist_memory:
304 return None
306 mo = self._memory_objects
307 try:
308 path = await self._persistence_mgr.save_all(
309 pyramid=mo["pyramid"],
310 working=mo["working"],
311 conversation=mo["conversation"],
312 long_term=mo["long_term"],
313 reflection_engine=mo["reflection_engine"],
314 consolidation_pipeline=mo["consolidation_pipeline"],
315 )
316 self._logger.info(f"Memory snapshot saved: {path}")
317 return path
318 except Exception as exc:
319 self._logger.error(f"Failed to save memory snapshot: {exc}")
320 return None
322 async def load_memory_snapshot(self) -> int:
323 """Load memory state from disk and restore into registered objects.
324 Returns count of subsystems restored.
325 """
326 if not self.config.persist_memory:
327 return 0
329 mo = self._memory_objects
330 try:
331 restored = await self._persistence_mgr.restore_all(
332 pyramid=mo["pyramid"],
333 working=mo["working"],
334 conversation=mo["conversation"],
335 long_term=mo["long_term"],
336 reflection_engine=mo["reflection_engine"],
337 consolidation_pipeline=mo["consolidation_pipeline"],
338 )
339 if restored > 0:
340 self._logger.info(f"Memory snapshot loaded: {restored} subsystems restored")
341 return restored
342 except Exception as exc:
343 self._logger.error(f"Failed to load memory snapshot: {exc}")
344 return 0
346 def memory_snapshot_info(self) -> dict[str, Any]:
347 """Return metadata about the current memory snapshot on disk."""
348 return self._persistence_mgr.snapshot_info()
350 # ── PID file helpers ──────────────────────
352 def _read_pid(self) -> int | None:
353 """Read PID from pidfile. Returns None if not running."""
354 path = Path(self.config.pidfile)
355 if not path.exists():
356 return None
357 try:
358 pid = int(path.read_text().strip())
359 except (ValueError, OSError):
360 return None
361 # Check if process is actually running
362 try:
363 os.kill(pid, 0)
364 return pid
365 except (ProcessLookupError, PermissionError):
366 return None
368 def _write_pid(self, pid: int) -> None:
369 """Write PID to pidfile."""
370 path = Path(self.config.pidfile)
371 path.parent.mkdir(parents=True, exist_ok=True)
372 path.write_text(f"{pid}\n")
374 def _remove_pid(self) -> None:
375 """Remove pidfile."""
376 path = Path(self.config.pidfile)
377 if path.exists():
378 path.unlink()
380 # ── Status ────────────────────────────────
382 def status(self) -> dict:
383 """Get daemon status as a dict."""
384 pid = self._read_pid()
385 running = pid is not None
386 uptime = time.time() - self._started_at if self._started_at and running else 0
387 return {
388 "running": running,
389 "pid": pid,
390 "host": self.config.host,
391 "port": self.config.port,
392 "pidfile": self.config.pidfile,
393 "logfile": self.config.logfile,
394 "uptime_seconds": round(uptime, 1),
395 "active_tasks": self.task_queue.active_count() if running else 0,
396 }
398 # ── Start ─────────────────────────────────
400 def start(self, daemonize: bool = True) -> int:
401 """Start the server. Returns PID if daemonized, 0 if foreground."""
402 if self._read_pid():
403 self._logger.warning("Daemon is already running.")
404 print(
405 f"Daemon already running (pid={self._read_pid()}) on "
406 f"http://{self.config.host}:{self.config.port}"
407 )
408 return self._read_pid() or 0
410 if daemonize:
411 return self._daemonize()
412 else:
413 return self._run_foreground()
415 def _daemonize(self) -> int:
416 """Fork into background daemon."""
417 pid = os.fork()
418 if pid > 0:
419 # Parent: wait briefly for child to start, then return
420 time.sleep(0.5)
421 child_pid = self._read_pid()
422 if child_pid:
423 self._logger.info(f"Daemon started (pid={child_pid})")
424 print(
425 f"Daemon started (pid={child_pid})\n"
426 f" http://{self.config.host}:{self.config.port}\n"
427 f" health: http://{self.config.host}:{self.config.port}/healthz\n"
428 f" logs: {self.config.logfile}"
429 )
430 return child_pid
431 else:
432 print("Failed to start daemon — check logs.")
433 return 1
435 # Child process
436 os.setsid()
437 # Second fork to detach from session
438 pid2 = os.fork()
439 if pid2 > 0:
440 os._exit(0)
442 # Grandchild: the actual daemon
443 self._write_pid(os.getpid())
444 atexit.register(self._remove_pid)
446 # Redirect stdin/stdout/stderr
447 devnull = os.open(os.devnull, os.O_RDWR)
448 os.dup2(devnull, sys.stdin.fileno())
449 os.dup2(devnull, sys.stdout.fileno())
450 os.dup2(devnull, sys.stderr.fileno())
451 if devnull > 2:
452 os.close(devnull)
454 self._started_at = time.time()
455 self._logger.info(
456 f"Daemon starting on http://{self.config.host}:{self.config.port}"
457 )
458 self._run_server()
459 return 0
461 def _run_foreground(self) -> int:
462 """Run in foreground (debug mode)."""
463 self._write_pid(os.getpid())
464 atexit.register(self._remove_pid)
465 self._started_at = time.time()
466 print(
467 f"Running in foreground on http://{self.config.host}:{self.config.port}\n"
468 f" health: http://{self.config.host}:{self.config.port}/healthz\n"
469 f" press Ctrl+C to stop"
470 )
471 self._run_server()
472 return 0
474 def _run_server(self) -> None:
475 """Run the uvicorn server (blocking)."""
476 app = self._build_app()
477 uvicorn.run(
478 app,
479 host=self.config.host,
480 port=self.config.port,
481 log_level=self.config.log_level,
482 workers=self.config.workers if self.config.workers > 1 else None,
483 timeout_graceful_shutdown=self.config.shutdown_timeout,
484 )
486 # ── Stop ──────────────────────────────────
488 def stop(self) -> bool:
489 """Stop the running daemon via SIGTERM."""
490 pid = self._read_pid()
491 if not pid:
492 print("No running daemon found.")
493 return False
495 self._logger.info(f"Stopping daemon (pid={pid})")
496 print(f"Stopping daemon (pid={pid})...")
497 try:
498 os.kill(pid, signal.SIGTERM)
499 except ProcessLookupError:
500 self._remove_pid()
501 print("Daemon already stopped.")
502 return True
504 # Wait for graceful shutdown
505 timeout = self.config.shutdown_timeout
506 for _ in range(int(timeout * 2)):
507 time.sleep(0.5)
508 try:
509 os.kill(pid, 0)
510 except ProcessLookupError:
511 self._remove_pid()
512 print("Daemon stopped.")
513 return True
515 # Force kill
516 print(f"Daemon did not stop within {timeout}s, sending SIGKILL...")
517 try:
518 os.kill(pid, signal.SIGKILL)
519 except ProcessLookupError:
520 pass
521 self._remove_pid()
522 print("Daemon force-stopped.")
523 return True
525 # ── Restart ───────────────────────────────
527 def restart(self, daemonize: bool = True) -> int:
528 """Stop then start the daemon."""
529 self.stop()
530 time.sleep(1)
531 return self.start(daemonize=daemonize)
533 # ── App building ──────────────────────────
535 def _default_app_factory(self) -> FastAPI:
536 """Default app: minimal standalone with health check."""
537 return create_daemon_app(self.task_queue)
539 def _build_app(self) -> FastAPI:
540 """Build the FastAPI application with memory persistence hooks."""
541 app = self._app_factory()
543 # Inject daemon state into app
544 app.state.daemon = self
545 app.state.task_queue = self.task_queue
547 # Patch lifespan to add memory persistence hooks (v1.14.9)
548 _original_lifespan = getattr(app.router, "lifespan_context", None)
550 @asynccontextmanager
551 async def memory_lifespan(app: FastAPI):
552 """Wrap existing lifespan with memory save/load."""
553 # Load on startup
554 restored = await self.load_memory_snapshot()
555 if restored > 0:
556 self._logger.info(f"Loaded {restored} memory subsystems from snapshot")
558 # Execute original lifespan
559 if _original_lifespan is not None:
560 async with _original_lifespan(app):
561 pass
563 # Yield to the app
564 yield
566 # Save on shutdown
567 await self.save_memory_snapshot()
568 self._logger.info("Memory snapshot saved on shutdown")
570 app.router.lifespan_context = memory_lifespan
572 # Ensure /healthz endpoint exists
573 has_healthz = any(
574 any(r.path == "/healthz" for r in router.routes)
575 for router in app.router.routes
576 if hasattr(router, "routes")
577 )
578 # Also check top-level routes
579 has_healthz = has_healthz or any(
580 getattr(r, "path", "") == "/healthz" for r in app.routes
581 )
583 if not has_healthz:
584 health_router = _make_health_router(self, self.task_queue)
585 app.include_router(health_router)
587 return app
590# ── Health / API routes ─────────────────────────────────
593def _make_health_router(daemon: ServerDaemon, tq: TaskQueue) -> APIRouter:
594 """Create the health check and management router."""
595 router = APIRouter(tags=["daemon"])
597 @router.get("/healthz")
598 async def healthz():
599 """Kubernetes-style health check."""
600 return {
601 "status": "healthy",
602 "uptime_seconds": round(time.time() - (daemon._started_at or time.time()), 1),
603 "active_tasks": tq.active_count(),
604 }
606 @router.get("/healthz/ready")
607 async def ready():
608 """Readiness check."""
609 return {
610 "status": "ready",
611 "active_tasks": tq.active_count(),
612 }
614 @router.get("/api/daemon/status")
615 async def daemon_status():
616 """Full daemon status."""
617 return daemon.status()
619 @router.get("/api/daemon/tasks")
620 async def list_tasks(limit: int = 50):
621 """List background tasks."""
622 tasks = [t.to_dict() for t in tq.list_tasks(limit)]
623 return {"count": len(tasks), "active": tq.active_count(), "tasks": tasks}
625 @router.get("/api/daemon/memory")
626 async def memory_status():
627 """Memory persistence status."""
628 info = daemon.memory_snapshot_info()
629 return {
630 "persistence_enabled": daemon.config.persist_memory,
631 "memory_dir": daemon.config.memory_dir,
632 "snapshot": info,
633 }
635 return router
638# ── Convenience factory ────────────────────────────────
641def create_daemon_app(task_queue: TaskQueue | None = None) -> FastAPI:
642 """Create a minimal standalone daemon FastAPI app."""
643 tq = task_queue or TaskQueue()
645 @asynccontextmanager
646 async def lifespan(app: FastAPI):
647 """Handle startup and shutdown."""
648 yield
649 await tq.shutdown(timeout=10.0)
651 app = FastAPI(
652 title="AgentOS Daemon",
653 version="1.14.9",
654 lifespan=lifespan,
655 )
656 app.add_middleware(
657 CORSMiddleware,
658 allow_origins=["*"],
659 allow_credentials=True,
660 allow_methods=["*"],
661 allow_headers=["*"],
662 )
664 @app.get("/")
665 async def root():
666 return {
667 "service": "AgentOS Daemon",
668 "version": "1.14.9",
669 "endpoints": {
670 "health": "/healthz",
671 "ready": "/healthz/ready",
672 "status": "/api/daemon/status",
673 "tasks": "/api/daemon/tasks",
674 },
675 }
677 return app
680# ── Module-level singleton ─────────────────────────────
682_daemon_instance: ServerDaemon | None = None
685def get_daemon(config: DaemonConfig | None = None) -> ServerDaemon:
686 """Get or create the module-level daemon singleton."""
687 global _daemon_instance
688 if _daemon_instance is None:
689 _daemon_instance = ServerDaemon(config=config)
690 return _daemon_instance
693# ── Main entry point ──────────────────────────────────
696def daemon_main(args: list[str] | None = None) -> int:
697 """CLI entry point for daemon commands."""
698 import argparse
700 parser = argparse.ArgumentParser(
701 prog="agentos-daemon",
702 description="AgentOS Independent Server Daemon",
703 )
704 sub = parser.add_subparsers(dest="command", required=True)
706 sub.add_parser("start", help="Start daemon in background")
707 sub.add_parser("run", help="Run in foreground")
708 sub.add_parser("stop", help="Stop running daemon")
709 sub.add_parser("restart", help="Stop then start")
710 sub.add_parser("status", help="Show daemon status")
712 ns = parser.parse_args(args)
714 daemon = get_daemon()
716 if ns.command == "start":
717 return daemon.start(daemonize=True)
718 elif ns.command == "run":
719 return daemon.start(daemonize=False)
720 elif ns.command == "stop":
721 return 0 if daemon.stop() else 1
722 elif ns.command == "restart":
723 return daemon.restart()
724 elif ns.command == "status":
725 s = daemon.status()
726 if s["running"]:
727 print(
728 f"RUNNING (pid={s['pid']})\n"
729 f" url: http://{s['host']}:{s['port']}\n"
730 f" uptime: {s['uptime_seconds']}s\n"
731 f" tasks: {s['active_tasks']} active"
732 )
733 else:
734 print("STOPPED")
735 return 0
737 return 1
740if __name__ == "__main__":
741 sys.exit(daemon_main())