mirror of
https://github.com/open-jarvis/OpenJarvis.git
synced 2026-07-30 02:42:16 +00:00
feat: wire AgentScheduler into server, register agents on create
This commit is contained in:
@@ -251,12 +251,41 @@ def serve(
|
||||
except Exception as exc:
|
||||
logger.debug("Agent manager init failed: %s", exc)
|
||||
|
||||
# Set up agent scheduler for cron/interval agents
|
||||
agent_scheduler = None
|
||||
if agent_manager is not None:
|
||||
try:
|
||||
from openjarvis.agents.executor import AgentExecutor
|
||||
from openjarvis.agents.scheduler import AgentScheduler
|
||||
|
||||
executor = AgentExecutor(manager=agent_manager, event_bus=bus)
|
||||
from openjarvis.system import SystemBuilder
|
||||
system = SystemBuilder(config).build()
|
||||
executor.set_system(system)
|
||||
|
||||
agent_scheduler = AgentScheduler(
|
||||
manager=agent_manager,
|
||||
executor=executor,
|
||||
event_bus=bus,
|
||||
)
|
||||
for ag in agent_manager.list_agents():
|
||||
sched_type = ag.get("config", {}).get("schedule_type", "manual")
|
||||
if sched_type in ("cron", "interval") and ag["status"] not in (
|
||||
"archived", "error",
|
||||
):
|
||||
agent_scheduler.register_agent(ag["id"])
|
||||
agent_scheduler.start()
|
||||
console.print(" Scheduler: [cyan]active[/cyan]")
|
||||
except Exception as exc:
|
||||
logger.debug("Agent scheduler init failed: %s", exc)
|
||||
|
||||
app = create_app(
|
||||
engine, model_name, agent=agent, bus=bus,
|
||||
engine_name=engine_name, agent_name=agent_key or "",
|
||||
channel_bridge=channel_bridge, config=config,
|
||||
speech_backend=speech_backend,
|
||||
agent_manager=agent_manager,
|
||||
agent_scheduler=agent_scheduler,
|
||||
)
|
||||
|
||||
console.print(
|
||||
|
||||
@@ -7,7 +7,7 @@ from typing import Any, Dict, List, Optional, Tuple
|
||||
from openjarvis.agents.manager import AgentManager
|
||||
|
||||
try:
|
||||
from fastapi import APIRouter, HTTPException
|
||||
from fastapi import APIRouter, HTTPException, Request
|
||||
from pydantic import BaseModel
|
||||
except ImportError:
|
||||
raise ImportError("fastapi and pydantic are required for server routes")
|
||||
@@ -70,14 +70,23 @@ def create_agent_manager_router(
|
||||
return {"agents": manager.list_agents()}
|
||||
|
||||
@agents_router.post("")
|
||||
async def create_agent(req: CreateAgentRequest):
|
||||
async def create_agent(req: CreateAgentRequest, request: Request):
|
||||
if req.template_id:
|
||||
return manager.create_from_template(
|
||||
agent = manager.create_from_template(
|
||||
req.template_id, req.name, overrides=req.config
|
||||
)
|
||||
return manager.create_agent(
|
||||
name=req.name, agent_type=req.agent_type, config=req.config
|
||||
)
|
||||
else:
|
||||
agent = manager.create_agent(
|
||||
name=req.name, agent_type=req.agent_type, config=req.config
|
||||
)
|
||||
|
||||
# Register with scheduler if cron/interval
|
||||
scheduler = getattr(request.app.state, "agent_scheduler", None)
|
||||
sched_type = (req.config or {}).get("schedule_type", "manual")
|
||||
if scheduler and sched_type in ("cron", "interval"):
|
||||
scheduler.register_agent(agent["id"])
|
||||
|
||||
return agent
|
||||
|
||||
@agents_router.get("/{agent_id}")
|
||||
async def get_agent(agent_id: str):
|
||||
|
||||
@@ -58,6 +58,7 @@ def create_app(
|
||||
config=None,
|
||||
speech_backend=None,
|
||||
agent_manager=None,
|
||||
agent_scheduler=None,
|
||||
) -> FastAPI:
|
||||
"""Create and configure the FastAPI application.
|
||||
|
||||
@@ -133,6 +134,7 @@ def create_app(
|
||||
app.state.channel_bridge = channel_bridge
|
||||
app.state.speech_backend = speech_backend
|
||||
app.state.agent_manager = agent_manager
|
||||
app.state.agent_scheduler = agent_scheduler
|
||||
app.state.session_start = time.time()
|
||||
|
||||
app.include_router(router)
|
||||
|
||||
Reference in New Issue
Block a user