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