Files
AI_engine/main.py
2026-06-26 16:08:31 +05:30

207 lines
6.6 KiB
Python

"""
LogiFlow AI - Autonomous Logistics Agent System
=================================================
Usage:
python main.py --production # Start all agents (production mode)
python main.py --dashboard # Launch admin dashboard
python main.py --portal # Launch customer portal
python main.py --help # Show this help
"""
import asyncio
import sys
import signal
from datetime import datetime
from typing import Dict, Any
# Force UTF-8 on Windows so emoji in output doesn't crash
if hasattr(sys.stdout, "reconfigure"):
sys.stdout.reconfigure(encoding="utf-8", errors="replace")
if hasattr(sys.stderr, "reconfigure"):
sys.stderr.reconfigure(encoding="utf-8", errors="replace")
from dotenv import load_dotenv
load_dotenv() # must happen before any config import reads os.getenv()
from core.logger import logger
from core.agent import MasterAgent, Agent
from core.message_bus import message_bus
from core.http_client import close_session
from core.types import AgentTask, MessageType, Priority
from agents.order_agent import OrderAgent
from agents.dispatch_agent import DispatchAgent
from agents.fleet_agent import FleetAgent
from agents.hub_agent import HubAgent
from agents.customer_agent import CustomerAgent
from agents.exception_agent import ExceptionAgent
from agents.route_optimizer_agent import RouteOptimizerAgent
class LogiFlowAI:
"""Main orchestration class — manages all agents and system lifecycle."""
def __init__(self):
logger.info("=" * 60)
logger.info("LOGIFLOW AI - AUTONOMOUS LOGISTICS SYSTEM")
logger.info("=" * 60)
self.master_agent = MasterAgent()
self.agents: Dict[str, Agent] = {}
self.running = False
def initialize_agents(self):
logger.info("Initializing AI Agents...")
self.agents["ORDER_AGENT"] = OrderAgent()
self.agents["DISPATCH_AGENT"] = DispatchAgent()
self.agents["FLEET_AGENT"] = FleetAgent()
self.agents["HUB_AGENT"] = HubAgent()
self.agents["CUSTOMER_AGENT"] = CustomerAgent()
self.agents["EXCEPTION_AGENT"] = ExceptionAgent()
self.agents["ROUTE_OPTIMIZER"] = RouteOptimizerAgent()
for agent in self.agents.values():
self.master_agent.register_sub_agent(agent)
logger.info(f"{len(self.agents) + 1} agents initialized (JARVIS + {len(self.agents)} specialized)")
async def start_agents(self):
logger.info("Starting Agent System...")
self.running = True
tasks = [asyncio.create_task(self.master_agent.start())]
tasks += [asyncio.create_task(a.start()) for a in self.agents.values()]
logger.info("All agents running — waiting for work")
try:
await asyncio.gather(*tasks, return_exceptions=True)
except asyncio.CancelledError:
pass
finally:
self.running = False
async def stop(self):
self.running = False
await self.master_agent.stop()
for agent in self.agents.values():
await agent.stop()
def get_system_status(self) -> Dict[str, Any]:
status = {
"timestamp": datetime.now().isoformat(),
"system": "operational" if self.running else "stopped",
"agents": {},
}
for agent_id, agent in self.agents.items():
status["agents"][agent_id] = {
"status": agent.state.status,
"tasks_completed": agent.state.tasks_completed,
"current_task": agent.state.current_task,
}
status["agents"]["JARVIS"] = {
"status": self.master_agent.state.status,
"sub_agents": len(self.agents),
}
return status
# ------------------------------------------------------------------ #
# Run modes #
# ------------------------------------------------------------------ #
async def production_mode():
logger.info("PRODUCTION MODE — connecting to infrastructure")
try:
await message_bus.connect()
except Exception as e:
logger.warning(f"Could not connect to NATS: {e} — using local in-process fallback")
system = LogiFlowAI()
system.initialize_agents()
logger.info("Infrastructure connected. Press Ctrl+C to stop.")
loop = asyncio.get_running_loop()
def _shutdown(signum, frame):
logger.info("Shutdown signal received — stopping agents...")
loop.create_task(_graceful_shutdown(system))
signal.signal(signal.SIGINT, _shutdown)
signal.signal(signal.SIGTERM, _shutdown)
try:
await system.start_agents()
finally:
await message_bus.disconnect()
await close_session()
logger.info("Clean shutdown complete")
async def _graceful_shutdown(system: LogiFlowAI):
await system.stop()
tasks = [t for t in asyncio.all_tasks() if t is not asyncio.current_task()]
for t in tasks:
t.cancel()
async def dashboard_mode():
import subprocess
logger.info("Launching Admin Dashboard at http://localhost:8501")
subprocess.run(["streamlit", "run", "dashboard/admin_dashboard.py"])
async def portal_mode():
import subprocess
logger.info("Launching Customer Portal at http://localhost:8502")
subprocess.run(["streamlit", "run", "customer_portal/portal.py", "--server.port", "8502"])
# ------------------------------------------------------------------ #
# Entry point #
# ------------------------------------------------------------------ #
def main():
if len(sys.argv) < 2 or sys.argv[1] == "--help":
print_help()
return
cmd = sys.argv[1]
if cmd == "--production":
asyncio.run(production_mode())
elif cmd == "--dashboard":
asyncio.run(dashboard_mode())
elif cmd == "--portal":
asyncio.run(portal_mode())
else:
logger.error(f"Unknown argument: {cmd}")
print_help()
sys.exit(1)
def print_help():
print("""
LogiFlow AI - Autonomous Logistics Agent System
===============================================
Usage:
python main.py --production Connect to NATS/Redis/Postgres and run forever
python main.py --dashboard Launch Streamlit admin dashboard
python main.py --portal Launch Streamlit customer portal
python main.py --help Show this help
Infrastructure (set in .env):
NATS nats://doormile@66.116.226.161:4223
Redis 66.116.226.255:6380
PG DB_HOST / DB_PORT / DB_NAME / DB_USER / DB_PASSWORD
API GO_API_BASE_URL
""")
if __name__ == "__main__":
main()