| name | event-publisher |
| description | Generate Redis pub/sub event publisher and subscriber patterns.
Use when: (1) Broadcasting events between services, (2) Implementing event-driven patterns,
(3) Real-time notifications, (4) Decoupling service communication.
Builds on redis-pattern skill with domain-specific event types and handlers.
NOT for task queues (use redis-pattern task-queue instead).
|
Event Publisher Generator
Generate domain-specific event publishing and subscribing patterns.
Quick Reference
/event-publisher chat # Chat message events
/event-publisher stream # Streaming status events
/event-publisher tool # Tool execution events
/event-publisher <domain> # Custom domain events
Generated Structure
events/
├── base.py # Event base class
├── publisher.py # Event publisher
├── subscriber.py # Event subscriber
├── handlers.py # Event handlers
└── types/
├── chat.py # Chat events
├── stream.py # Stream events
└── tool.py # Tool events
Implementation
1. Event Types
from pydantic import BaseModel
from datetime import datetime
from typing import Literal
class ChatEvent(BaseModel):
"""Base chat event."""
event_id: str
conversation_id: str
user_id: str
timestamp: datetime
class MessageCreated(ChatEvent):
event_type: Literal["chat.message.created"] = "chat.message.created"
message_id: str
role: str
content: str
class MessageUpdated(ChatEvent):
event_type: Literal["chat.message.updated"] = "chat.message.updated"
message_id: str
content: str
class TypingStarted(ChatEvent):
event_type: Literal["chat.typing.started"] = "chat.typing.started"
class TypingStopped(ChatEvent):
event_type: Literal["chat.typing.stopped"] = "chat.typing.stopped"
ChatEventTypes = MessageCreated | MessageUpdated | TypingStarted | TypingStopped
"""
Stream events - synchronized with streaming-llm-responses skill.
These are the producer-side definitions; Stream Broker consumes them.
"""
from pydantic import BaseModel
from typing import Literal
from datetime import datetime
import uuid
class StreamEvent(BaseModel):
"""Base for all stream events."""
id: str = None
timestamp: datetime = None
user_id: str
def __init__(self, **data):
if "id" not in data or data["id"] is None:
data["id"] = str(uuid.uuid4())
if "timestamp" not in data or data["timestamp"] is None:
data["timestamp"] = datetime.utcnow()
super().__init__(**data)
class StatusEvent(StreamEvent):
"""Processing state updates (thinking, searching, etc.)."""
type: Literal["status"] = "status"
status: Literal["thinking", "searching", "generating"]
message: str
class DeltaEvent(StreamEvent):
"""Token-by-token content streaming."""
type: Literal["delta"] = "delta"
content: str
message_id: str
class ToolStartEvent(StreamEvent):
"""Tool execution has begun."""
type: Literal["tool_start"] = "tool_start"
tool_name: str
tool_call_id: str
class ToolEndEvent(StreamEvent):
"""Tool execution completed."""
type: Literal["tool_end"] = "tool_end"
tool_name: str
tool_call_id: str
result: dict | None = None
class ErrorEvent(StreamEvent):
"""Error notification."""
type: Literal["error"] = "error"
code: str
message: str
recoverable: bool = True
StreamEventTypes = StatusEvent | DeltaEvent | ToolStartEvent | ToolEndEvent | ErrorEvent
2. Publisher
"""
Event publisher - produces events for Stream Broker consumption.
Channel naming: events:stream:{user_id} (matches streaming-llm-responses).
"""
from redis.asyncio import Redis
from pydantic import BaseModel
from typing import Literal
from .types.stream import (
StatusEvent, DeltaEvent, ToolStartEvent, ToolEndEvent, ErrorEvent
)
class EventPublisher:
def __init__(self, redis: Redis, prefix: str = "events"):
self.redis = redis
self.prefix = prefix
def _channel(self, *parts: str) -> str:
return ":".join([self.prefix, *parts])
async def publish(self, channel: str, event: BaseModel) -> str:
"""Publish event to channel."""
await self.redis.publish(
self._channel(channel),
event.model_dump_json(),
)
return getattr(event, "id", "")
async def status(
self,
user_id: str,
status: Literal["thinking", "searching", "generating"],
message: str,
) -> str:
"""Publish status update to user's stream."""
return await self.publish(
f"stream:{user_id}",
StatusEvent(user_id=user_id, status=status, message=message),
)
async def delta(self, user_id: str, content: str, message_id: str) -> str:
"""Publish content delta to user's stream."""
return await self.publish(
f"stream:{user_id}",
DeltaEvent(user_id=user_id, content=content, message_id=message_id),
)
async def tool_start(self, user_id: str, tool_name: str, tool_call_id: str) -> str:
"""Publish tool start event."""
return await self.publish(
f"stream:{user_id}",
ToolStartEvent(user_id=user_id, tool_name=tool_name, tool_call_id=tool_call_id),
)
async def tool_end(
self,
user_id: str,
tool_name: str,
tool_call_id: str,
result: dict | None = None,
) -> str:
"""Publish tool end event."""
return await self.publish(
f"stream:{user_id}",
ToolEndEvent(
user_id=user_id,
tool_name=tool_name,
tool_call_id=tool_call_id,
result=result,
),
)
async def error(
self,
user_id: str,
code: str,
message: str,
recoverable: bool = True,
) -> str:
"""Publish error event."""
return await self.publish(
f"stream:{user_id}",
ErrorEvent(
user_id=user_id,
code=code,
message=message,
recoverable=recoverable,
),
)
async def to_conversation(self, conversation_id: str, event: BaseModel) -> str:
return await self.publish(f"conversation:{conversation_id}", event)
async def broadcast(self, event: BaseModel) -> str:
return await self.publish("broadcast", event)
3. Subscriber with Handlers
from redis.asyncio import Redis
from typing import Callable, Awaitable, Type
from pydantic import BaseModel
import asyncio
import json
Handler = Callable[[BaseModel], Awaitable[None]]
class EventSubscriber:
def __init__(self, redis: Redis, prefix: str = "events"):
self.redis = redis
self.prefix = prefix
self._handlers: dict[str, list[tuple[Type[BaseModel], Handler]]] = {}
self._running = False
def on(
self,
event_type: str,
model: Type[BaseModel],
) -> Callable[[Handler], Handler]:
"""Decorator to register event handler."""
def decorator(handler: Handler) -> Handler:
if event_type not in self._handlers:
self._handlers[event_type] = []
self._handlers[event_type].append((model, handler))
return handler
return decorator
async def subscribe(self, *channels: str) -> None:
"""Subscribe to channels and dispatch events."""
pubsub = self.redis.pubsub()
full_channels = [f"{self.prefix}:{ch}" for ch in channels]
await pubsub.subscribe(*full_channels)
self._running = True
async for message in pubsub.listen():
if not self._running:
break
if message["type"] != "message":
continue
data = json.loads(message["data"])
event_type = data.get("event_type")
for model, handler in self._handlers.get(event_type, []):
event = model.model_validate(data)
asyncio.create_task(handler(event))
await pubsub.close()
def stop(self) -> None:
self._running = False
4. Handler Registration
from events.subscriber import EventSubscriber
from events.types.chat import MessageCreated, MessageUpdated
from events.types.stream import StreamCompleted
def register_handlers(subscriber: EventSubscriber) -> None:
"""Register all event handlers."""
@subscriber.on("chat.message.created", MessageCreated)
async def on_message_created(event: MessageCreated):
pass
@subscriber.on("stream.completed", StreamCompleted)
async def on_stream_completed(event: StreamCompleted):
pass
Usage
Publishing Stream Events (Orchestrator/Tool Executor)
publisher = EventPublisher(redis)
await publisher.status("user123", "thinking", "Analyzing your request...")
await publisher.status("user123", "searching", "Searching the web...")
for token in llm_response.tokens:
await publisher.delta("user123", token, message_id="msg456")
await publisher.tool_start("user123", "web_search", "call_789")
result = await execute_tool(...)
await publisher.tool_end("user123", "web_search", "call_789", result)
await publisher.error("user123", "LLM_TIMEOUT", "Request timed out", recoverable=True)
Publishing Chat Events (Persistence notifications)
await publisher.to_conversation(
"conv123",
MessageCreated(
event_id=str(uuid.uuid4()),
conversation_id="conv123",
user_id="user456",
timestamp=datetime.utcnow(),
message_id="msg789",
role="assistant",
content="Hello!",
),
)
Subscribing to Events
subscriber = EventSubscriber(redis)
register_handlers(subscriber)
await subscriber.subscribe("conversation:*", "stream:*")
Skill Relationship
| Skill | Role | Channel Pattern |
|---|
event-publisher | Producer (Orchestrator, Tools) | Publishes to events:stream:{user_id} |
streaming-llm-responses | Consumer (Stream Broker) | Subscribes to events:stream:{user_id}, delivers via SSE |
Both skills define the same event types (status, delta, tool_start, tool_end, error) to ensure contract compatibility.