xiaowei-system/skills/devops/devops-umbrella/scripts/agent-comm/redis-message-bus.py

120 lines
4.2 KiB
Python
Executable File

#!/usr/bin/env python3
"""
Redis-based message bus for multi-agent communication.
Usage:
Listener: python3 redis-message-bus.py listen <my_agent>
Sender: python3 redis-message-bus.py send <to_agent> <message>
Status: python3 redis-message-bus.py status
Queue: python3 redis-message-bus.py queue <my_agent> # pop pending messages
"""
import redis, json, time, sys, os, uuid
from datetime import datetime
REDIS_HOST = os.environ.get("REDIS_HOST", "localhost")
REDIS_PORT = int(os.environ.get("REDIS_PORT", 6379))
AGENT_HUB_DIR = os.path.expanduser("~/.agent-hub")
r = redis.Redis(host=REDIS_HOST, port=REDIS_PORT, decode_responses=True)
def get_config():
config_path = os.path.join(AGENT_HUB_DIR, "config", "agents.json")
if os.path.exists(config_path):
return json.load(open(config_path))
return None
def make_msg(from_agent, to_agent, msg_type, payload):
return {
"id": f"msg-{uuid.uuid4().hex[:12]}",
"from": from_agent,
"to": to_agent,
"type": msg_type,
"payload": payload,
"timestamp": time.time(),
"created_at": datetime.now().isoformat()
}
def send(to_agent, from_agent, msg_type, payload):
msg = make_msg(from_agent, to_agent, msg_type, payload)
channel = f"agent:{to_agent}"
r.publish(channel, json.dumps(msg))
# Also queue it for durability
r.hset(f"agent:queue:{to_agent}", msg["id"], json.dumps(msg))
return msg["id"]
def listen(agent, callback=None):
"""Listen for messages directed to this agent."""
pubsub = r.pubsub()
pubsub.subscribe(f"agent:{agent}")
print(f"[{agent}] Listening on agent:{agent}...", flush=True)
for msg in pubsub.listen():
if msg["type"] == "message":
data = json.loads(msg["data"])
if data.get("to") == agent:
print(f"[{agent}] {data['from']}{data['to']}: {data['type']} - {data['payload']}", flush=True)
if callback:
callback(data)
else:
# Mark as processed
r.hset(f"agent:queue:{agent}", data["id"], json.dumps({**data, "processed_at": time.time()}))
def queue(agent):
"""Pop pending messages for this agent."""
msgs = r.hgetall(f"agent:queue:{agent}")
pending = []
for msg_id, data in msgs.items():
d = json.loads(data)
if "processed_at" not in d:
pending.append(d)
return pending
def status():
"""Show system status."""
info = r.info()
agents_config = get_config()
if agents_config:
for agent_name in agents_config.get("agents", {}):
heartbeat = r.get(f"agent:{agent_name}:heartbeat")
qsize = r.hlen(f"agent:queue:{agent_name}")
last_seen = float(heartbeat) if heartbeat else None
status_str = "online" if last_seen and (time.time() - last_seen < 30) else "offline"
print(f" {agent_name}: {status_str} (queue: {qsize} pending)")
print(f" Redis: {info['connected_clients']} clients, {info['used_memory_human']} used")
print(f" Subscriptions: {r.pubsub_numsub('agent:*')}")
def heartbeat(agent, interval=10):
"""Send heartbeat every `interval` seconds. Run in background."""
while True:
r.setex(f"agent:{agent}:heartbeat", interval * 3, time.time())
time.sleep(interval)
if __name__ == "__main__":
if len(sys.argv) < 2:
print(__doc__)
sys.exit(1)
cmd = sys.argv[1]
if cmd == "listen" and len(sys.argv) >= 3:
listen(sys.argv[2])
elif cmd == "send" and len(sys.argv) >= 5:
_, send, to_agent, *rest = sys.argv
msg_text = " ".join(rest)
# Default from "xiao-wei" — callers can override this in code
msg_id = send(to_agent, "xiao-wei", "task", {"text": msg_text})
print(f"Sent {msg_id}{to_agent}")
elif cmd == "queue" and len(sys.argv) >= 3:
for msg in queue(sys.argv[2]):
print(json.dumps(msg, indent=2))
elif cmd == "status":
status()
elif cmd == "heartbeat" and len(sys.argv) >= 3:
interval = int(sys.argv[3]) if len(sys.argv) >= 4 else 10
heartbeat(sys.argv[2], interval)
else:
print(__doc__)