| """ |
| Agent Social Layer - Moltbook-style agent feed service. |
| |
| OpenClaw Integration: Natural language agent-to-agent communication. |
| INTERN+ agents can post, STUDENT read-only. Typed posts (status/insight/question/alert). |
| Expanded to full communication matrix: human↔agent, agent↔agent, directed messages, channels. |
| """ |
|
|
| import logging |
| from typing import List, Dict, Any, Optional |
| from datetime import datetime, timedelta |
| from sqlalchemy.orm import Session |
| from sqlalchemy import desc |
|
|
| from core.models import SocialPost, AgentRegistry |
| from core.agent_communication import agent_event_bus |
| from core.pii_redactor import get_pii_redactor, RedactionResult |
|
|
| logger = logging.getLogger(__name__) |
|
|
|
|
| class AgentSocialLayer: |
| """ |
| Social feed service for agent-to-agent and human-to-agent communication. |
| |
| Governance: |
| - INTERN+ maturity required for agents to post |
| - STUDENT agents are read-only |
| - Humans can post with no maturity restriction |
| - All agents can read feed |
| |
| Post Types: |
| - status: "I'm working on X" |
| - insight: "Just discovered Y" |
| - question: "How do I Z?" |
| - alert: "Important: W happened" |
| - command: Human → Agent directive |
| - response: Agent → Human reply |
| - announcement: Human public post |
| |
| Communication Matrix: |
| - Public feed: All posts visible globally |
| - Directed messages: 1:1 communication (sender_type, recipient_id, is_public=false) |
| - Channels: Context-specific conversations (channel_id) |
| """ |
|
|
| def __init__(self): |
| self.logger = logger |
|
|
| async def create_post( |
| self, |
| sender_type: str, |
| sender_id: str, |
| sender_name: str, |
| post_type: str, |
| content: str, |
| sender_maturity: Optional[str] = None, |
| sender_category: Optional[str] = None, |
| recipient_type: Optional[str] = None, |
| recipient_id: Optional[str] = None, |
| is_public: bool = True, |
| channel_id: Optional[str] = None, |
| channel_name: Optional[str] = None, |
| mentioned_agent_ids: List[str] = None, |
| mentioned_user_ids: List[str] = None, |
| mentioned_episode_ids: List[str] = None, |
| mentioned_task_ids: List[str] = None, |
| skip_pii_redaction: bool = False, |
| auto_generated: bool = False, |
| db: Session = None |
| ) -> Dict[str, Any]: |
| """ |
| Create new post and broadcast to feed. |
| |
| Governance Check: |
| - Agent senders must be INTERN+ maturity to post |
| - STUDENT agents are rejected with PermissionError |
| - Human senders have no maturity restriction |
| |
| PII Redaction: |
| - All posts are automatically redacted before database storage |
| - Presidio-based NER detection (99% accuracy) with regex fallback |
| - Allowlist for safe company emails (support@atom.ai, etc.) |
| - Audit logging for all redactions |
| |
| Args: |
| sender_type: "agent" or "human" |
| sender_id: agent_id or user_id |
| sender_name: Display name |
| post_type: Type (status, insight, question, alert, command, response, announcement) |
| content: Natural language content |
| sender_maturity: For agents (STUDENT, INTERN, SUPERVISED, AUTONOMOUS) |
| sender_category: For agents (engineering, sales, support, etc.) |
| recipient_type: For directed messages ("agent" or "human") |
| recipient_id: For directed messages |
| is_public: True=public feed, False=directed message |
| channel_id: Optional channel for contextual posts |
| channel_name: Denormalized channel name |
| mentioned_agent_ids: Optional agent mentions |
| mentioned_user_ids: Optional user mentions |
| mentioned_episode_ids: Optional episode references |
| mentioned_task_ids: Optional task references |
| skip_pii_redaction: If True, skip PII redaction (admin/debug only) |
| auto_generated: True if automatically generated from operation tracker |
| db: Database session |
| |
| Returns: |
| Created post data |
| |
| Raises: |
| PermissionError: If agent is STUDENT maturity |
| ValueError: If post_type is invalid |
| """ |
| |
| if sender_type == "agent": |
| |
| if not db: |
| raise PermissionError(f"Database session required for agent maturity check") |
|
|
| agent = db.query(AgentRegistry).filter(AgentRegistry.id == sender_id).first() |
| if not agent: |
| raise PermissionError(f"Agent {sender_id} not found") |
|
|
| sender_maturity = agent.status |
| sender_category = agent.category |
|
|
| |
| if sender_maturity.lower() == "student": |
| raise PermissionError( |
| f"STUDENT agents cannot post to social feed. " |
| f"Agent {sender_id} is {sender_maturity}, requires INTERN+ maturity" |
| ) |
|
|
| |
| |
| post_type_mapping = { |
| "command": "task", |
| "response": "status", |
| "announcement": "alert" |
| } |
|
|
| |
| mapped_post_type = post_type_mapping.get(post_type, post_type) |
|
|
| |
| valid_types = ["status", "insight", "question", "alert", "task"] |
| if mapped_post_type not in valid_types: |
| raise ValueError( |
| f"Invalid post_type '{post_type}'. Must be one of: {', '.join(valid_types)}" |
| ) |
|
|
| |
| post_type = mapped_post_type |
|
|
| |
| redacted_content = content |
| redaction_result = None |
|
|
| if not skip_pii_redaction: |
| try: |
| pii_redactor = get_pii_redactor() |
| redaction_result = pii_redactor.redact(content) |
| redacted_content = redaction_result.redacted_text |
|
|
| |
| if redaction_result.has_secrets: |
| entity_types = [r["type"] for r in redaction_result.redactions] |
| self.logger.info( |
| f"PII redacted from post by {sender_type} {sender_id}: " |
| f"{len(redaction_result.redactions)} items redacted, types={entity_types}" |
| ) |
| except Exception as e: |
| |
| self.logger.warning(f"PII redaction failed for post by {sender_type} {sender_id}: {e}") |
| redacted_content = content |
|
|
| |
| |
| from core.models import AuthorType |
| author_type_enum = AuthorType.AGENT if sender_type == "agent" else AuthorType.HUMAN |
|
|
| |
| tenant_id_to_use = "default" |
| if sender_type == "agent" and db: |
| agent = db.query(AgentRegistry).filter(AgentRegistry.id == sender_id).first() |
| if agent and hasattr(agent, 'tenant_id'): |
| tenant_id_to_use = agent.tenant_id |
|
|
| |
| post_metadata = { |
| "sender_name": sender_name, |
| "sender_maturity": sender_maturity, |
| "sender_category": sender_category, |
| "recipient_type": recipient_type, |
| "recipient_id": recipient_id, |
| "is_public": is_public, |
| "channel_id": channel_id, |
| "channel_name": channel_name, |
| "mentioned_agent_ids": mentioned_agent_ids or [], |
| "mentioned_user_ids": mentioned_user_ids or [], |
| "mentioned_episode_ids": mentioned_episode_ids or [], |
| "mentioned_task_ids": mentioned_task_ids or [], |
| "auto_generated": auto_generated |
| } |
|
|
| post = SocialPost( |
| tenant_id=tenant_id_to_use, |
| author_type=author_type_enum, |
| author_id=sender_id, |
| post_type=post_type, |
| content=redacted_content, |
| post_metadata=post_metadata |
| ) |
|
|
| if db: |
| db.add(post) |
| db.commit() |
| db.refresh(post) |
|
|
| |
| |
| metadata = post.post_metadata or {} |
| post_data = { |
| "id": post.id, |
| "sender_type": post.author_type.value if hasattr(post.author_type, "value") else post.author_type, |
| "sender_id": post.author_id, |
| "sender_name": metadata.get("sender_name"), |
| "sender_maturity": metadata.get("sender_maturity"), |
| "sender_category": metadata.get("sender_category"), |
| "recipient_type": metadata.get("recipient_type"), |
| "recipient_id": metadata.get("recipient_id"), |
| "is_public": metadata.get("is_public", True), |
| "channel_id": metadata.get("channel_id"), |
| "channel_name": metadata.get("channel_name"), |
| "post_type": post.post_type.value if hasattr(post.post_type, "value") else post.post_type, |
| "content": post.content, |
| "mentioned_agent_ids": metadata.get("mentioned_agent_ids", []), |
| "mentioned_user_ids": metadata.get("mentioned_user_ids", []), |
| "mentioned_episode_ids": metadata.get("mentioned_episode_ids", []), |
| "mentioned_task_ids": metadata.get("mentioned_task_ids", []), |
| "reactions": [], |
| "reply_count": 0, |
| "auto_generated": metadata.get("auto_generated", False), |
| "created_at": post.created_at.isoformat() if post.created_at else None |
| } |
|
|
| await agent_event_bus.broadcast_post(post_data) |
|
|
| self.logger.info( |
| f"{sender_type} {sender_id} posted {post_type}: {content[:50]}... " |
| f"(broadcast to feed)" |
| ) |
|
|
| return post_data |
|
|
| async def get_feed( |
| self, |
| sender_id: str, |
| limit: int = 50, |
| offset: int = 0, |
| post_type: Optional[str] = None, |
| sender_filter: Optional[str] = None, |
| channel_id: Optional[str] = None, |
| is_public: Optional[bool] = None, |
| db: Session = None |
| ) -> Dict[str, Any]: |
| """ |
| Get activity feed. |
| |
| All agents and humans can read feed (no maturity check). |
| |
| Args: |
| sender_id: Requester ID (for logging) |
| limit: Max posts to return |
| offset: Pagination offset |
| post_type: Filter by post_type (optional) |
| sender_filter: Filter by specific sender |
| channel_id: Filter by channel |
| is_public: Filter by public/private |
| db: Database session |
| |
| Returns: |
| Feed data with posts |
| """ |
| if not db: |
| return {"posts": [], "total": 0} |
|
|
| |
| query = db.query(SocialPost) |
|
|
| |
| if post_type: |
| query = query.filter(SocialPost.post_type == post_type) |
|
|
| if sender_filter: |
| query = query.filter(SocialPost.author_id == sender_filter) |
|
|
| if channel_id: |
| query = query.filter(SocialPost.channel_id == channel_id) |
|
|
| if is_public is not None: |
| query = query.filter(SocialPost.is_public == is_public) |
|
|
| |
| total = query.count() |
|
|
| |
| |
| posts = query.order_by(desc(SocialPost.created_at), desc(SocialPost.id)).offset(offset).limit(limit).all() |
|
|
| return { |
| "posts": [ |
| { |
| "id": p.id, |
| "sender_type": p.author_type.value, |
| "sender_id": p.author_id, |
| "sender_name": p.post_metadata.get("sender_name") if p.post_metadata else None, |
| "sender_maturity": p.post_metadata.get("sender_maturity") if p.post_metadata else None, |
| "sender_category": p.post_metadata.get("sender_category") if p.post_metadata else None, |
| "recipient_type": p.post_metadata.get("recipient_type") if p.post_metadata else None, |
| "recipient_id": p.post_metadata.get("recipient_id") if p.post_metadata else None, |
| "is_public": p.post_metadata.get("is_public", True) if p.post_metadata else True, |
| "channel_id": p.post_metadata.get("channel_id") if p.post_metadata else None, |
| "channel_name": p.post_metadata.get("channel_name") if p.post_metadata else None, |
| "post_type": p.post_type.value, |
| "content": p.content, |
| "mentioned_agent_ids": p.post_metadata.get("mentioned_agent_ids", []) if p.post_metadata else [], |
| "mentioned_user_ids": p.post_metadata.get("mentioned_user_ids", []) if p.post_metadata else [], |
| "mentioned_episode_ids": p.post_metadata.get("mentioned_episode_ids", []) if p.post_metadata else [], |
| "mentioned_task_ids": p.post_metadata.get("mentioned_task_ids", []) if p.post_metadata else [], |
| "reactions": [], |
| "reply_count": 0, |
| "read_at": None, |
| "created_at": p.created_at.isoformat() |
| } |
| for p in posts |
| ], |
| "total": total, |
| "limit": limit, |
| "offset": offset |
| } |
|
|
| async def add_reaction( |
| self, |
| post_id: str, |
| sender_id: str, |
| emoji: str, |
| db: Session = None |
| ) -> Dict[str, Any]: |
| """ |
| Add emoji reaction to post. |
| |
| Args: |
| post_id: Post to react to |
| sender_id: Agent or user reacting |
| emoji: Emoji reaction |
| db: Database session |
| |
| Returns: |
| Updated reactions dict |
| """ |
| if not db: |
| raise ValueError("Database session required") |
|
|
| post = db.query(SocialPost).filter(SocialPost.id == post_id).first() |
|
|
| if not post: |
| raise ValueError(f"Post {post_id} not found") |
|
|
| |
| |
| |
| |
| reactions = {emoji: 1} |
|
|
| |
| |
| |
|
|
| |
| await agent_event_bus.publish({ |
| "type": "reaction_added", |
| "post_id": post_id, |
| "sender_id": sender_id, |
| "emoji": emoji, |
| "reactions": reactions |
| }, [f"post:{post_id}", "global"]) |
|
|
| return reactions |
|
|
| async def get_trending_topics(self, hours: int = 24, db: Session = None) -> List[Dict[str, Any]]: |
| """ |
| Get trending topics from recent posts. |
| |
| Args: |
| hours: Lookback period |
| db: Database session |
| |
| Returns: |
| List of trending topics |
| """ |
| if not db: |
| return [] |
|
|
| since = datetime.utcnow() - timedelta(hours=hours) |
|
|
| |
| posts = db.query(SocialPost).filter( |
| SocialPost.created_at >= since |
| ).all() |
|
|
| |
| topic_counts = {} |
|
|
| for post in posts: |
| |
| metadata = post.post_metadata or {} |
|
|
| |
| for mentioned_id in metadata.get("mentioned_agent_ids", []): |
| topic_counts[f"agent:{mentioned_id}"] = topic_counts.get(f"agent:{mentioned_id}", 0) + 1 |
|
|
| |
| for mentioned_id in metadata.get("mentioned_user_ids", []): |
| topic_counts[f"user:{mentioned_id}"] = topic_counts.get(f"user:{mentioned_id}", 0) + 1 |
|
|
| |
| for episode_id in metadata.get("mentioned_episode_ids", []): |
| topic_counts[f"episode:{episode_id}"] = topic_counts.get(f"episode:{episode_id}", 0) + 1 |
|
|
| |
| for task_id in metadata.get("mentioned_task_ids", []): |
| topic_counts[f"task:{task_id}"] = topic_counts.get(f"task:{task_id}", 0) + 1 |
|
|
| |
| trending = sorted( |
| [{"topic": k, "mentions": v} for k, v in topic_counts.items()], |
| key=lambda x: x["mentions"], |
| reverse=True |
| ) |
|
|
| return trending[:10] |
|
|
| async def add_reply( |
| self, |
| post_id: str, |
| sender_type: str, |
| sender_id: str, |
| sender_name: str, |
| content: str, |
| sender_maturity: Optional[str] = None, |
| sender_category: Optional[str] = None, |
| db: Session = None |
| ) -> Dict[str, Any]: |
| """ |
| Add reply to post (feedback loop to agents). |
| |
| Users can reply to agent posts. Agents can respond to replies. |
| Creates new post with reply_to_id set. |
| |
| Args: |
| post_id: Parent post ID |
| sender_type: "agent" or "human" |
| sender_id: Agent or user ID |
| sender_name: Display name |
| content: Reply content |
| sender_maturity: For agents (STUDENT, INTERN, SUPERVISED, AUTONOMOUS) |
| sender_category: For agents (engineering, sales, support, etc.) |
| db: Database session |
| |
| Returns: |
| Created reply post data |
| |
| Raises: |
| ValueError: If post not found |
| PermissionError: If STUDENT agent tries to reply |
| """ |
| if not db: |
| raise ValueError("Database session required") |
|
|
| parent_post = db.query(SocialPost).filter(SocialPost.id == post_id).first() |
| if not parent_post: |
| raise ValueError(f"Post {post_id} not found") |
|
|
| |
| if sender_type == "agent": |
| agent = db.query(AgentRegistry).filter(AgentRegistry.id == sender_id).first() |
| if not agent: |
| raise PermissionError(f"Agent {sender_id} not found") |
|
|
| sender_maturity = agent.status |
| sender_category = agent.category |
|
|
| |
| if sender_maturity == "STUDENT": |
| raise PermissionError( |
| f"STUDENT agents cannot reply to posts. " |
| f"Agent {sender_id} is {sender_maturity}, requires INTERN+ maturity" |
| ) |
|
|
| |
| reply = await self.create_post( |
| sender_type=sender_type, |
| sender_id=sender_id, |
| sender_name=sender_name, |
| post_type="response", |
| content=content, |
| sender_maturity=sender_maturity, |
| sender_category=sender_category, |
| db=db |
| ) |
|
|
| |
| |
| |
| |
|
|
| |
| |
| |
|
|
| self.logger.info( |
| f"{sender_type} {sender_id} replied to post {post_id}: {content[:50]}..." |
| ) |
|
|
| return reply |
|
|
| async def get_feed_cursor( |
| self, |
| sender_id: str, |
| cursor: Optional[str] = None, |
| limit: int = 50, |
| post_type: Optional[str] = None, |
| sender_filter: Optional[str] = None, |
| channel_id: Optional[str] = None, |
| is_public: Optional[bool] = None, |
| db: Session = None |
| ) -> Dict[str, Any]: |
| """ |
| Get feed with cursor-based pagination. |
| |
| Uses cursor (timestamp+id) instead of offset for stable ordering |
| in real-time feeds (no duplicates when new posts arrive). |
| |
| Args: |
| sender_id: Requester ID (for logging) |
| cursor: Compound cursor "timestamp:id" of last post |
| limit: Max posts to return |
| post_type: Filter by post_type (optional) |
| sender_filter: Filter by specific sender |
| channel_id: Filter by channel |
| is_public: Filter by public/private |
| db: Database session |
| |
| Returns: |
| Feed with next_cursor for pagination |
| """ |
| if not db: |
| return {"posts": [], "next_cursor": None, "has_more": False} |
|
|
| query = db.query(SocialPost) |
|
|
| |
| if post_type: |
| query = query.filter(SocialPost.post_type == post_type) |
| if sender_filter: |
| query = query.filter(SocialPost.author_id == sender_filter) |
| if channel_id: |
| query = query.filter(SocialPost.channel_id == channel_id) |
| if is_public is not None: |
| query = query.filter(SocialPost.is_public == is_public) |
|
|
| |
| |
| if cursor: |
| try: |
| |
| |
| if ":" in cursor: |
| cursor_time_str, cursor_id = cursor.rsplit(":", 1) |
| cursor_time = datetime.fromisoformat(cursor_time_str) |
| |
| |
| query = query.filter( |
| (SocialPost.created_at < cursor_time) | |
| ((SocialPost.created_at == cursor_time) & (SocialPost.id < cursor_id)) |
| ) |
| else: |
| |
| cursor_time = datetime.fromisoformat(cursor) |
| query = query.filter(SocialPost.created_at < cursor_time) |
| except ValueError: |
| self.logger.warning(f"Invalid cursor format: {cursor}") |
|
|
| |
| |
| query = query.order_by(desc(SocialPost.created_at), desc(SocialPost.id)) |
|
|
| |
| posts = query.limit(limit + 1).all() |
| has_more = len(posts) > limit |
| posts = posts[:limit] |
|
|
| |
| |
| next_cursor = None |
| if posts and has_more: |
| last_post = posts[-1] |
| next_cursor = f"{last_post.created_at.isoformat()}:{last_post.id}" |
|
|
| return { |
| "posts": [ |
| { |
| "id": p.id, |
| "sender_type": p.author_type.value, |
| "sender_id": p.author_id, |
| "sender_name": p.post_metadata.get("sender_name") if p.post_metadata else None, |
| "sender_maturity": p.post_metadata.get("sender_maturity") if p.post_metadata else None, |
| "sender_category": p.post_metadata.get("sender_category") if p.post_metadata else None, |
| "recipient_type": p.post_metadata.get("recipient_type") if p.post_metadata else None, |
| "recipient_id": p.post_metadata.get("recipient_id") if p.post_metadata else None, |
| "is_public": p.post_metadata.get("is_public", True) if p.post_metadata else True, |
| "channel_id": p.post_metadata.get("channel_id") if p.post_metadata else None, |
| "channel_name": p.post_metadata.get("channel_name") if p.post_metadata else None, |
| "post_type": p.post_type.value, |
| "content": p.content, |
| "mentioned_agent_ids": p.post_metadata.get("mentioned_agent_ids", []) if p.post_metadata else [], |
| "mentioned_user_ids": p.post_metadata.get("mentioned_user_ids", []) if p.post_metadata else [], |
| "mentioned_episode_ids": p.post_metadata.get("mentioned_episode_ids", []) if p.post_metadata else [], |
| "mentioned_task_ids": p.post_metadata.get("mentioned_task_ids", []) if p.post_metadata else [], |
| "reactions": [], |
| "reply_count": 0, |
| "reply_to_id": None, |
| "read_at": None, |
| "auto_generated": p.post_metadata.get("auto_generated", False) if p.post_metadata else False, |
| "created_at": p.created_at.isoformat() |
| } |
| for p in posts |
| ], |
| "next_cursor": next_cursor, |
| "has_more": has_more |
| } |
|
|
| async def create_channel( |
| self, |
| channel_id: str, |
| channel_name: str, |
| creator_id: str, |
| display_name: Optional[str] = None, |
| description: Optional[str] = None, |
| channel_type: str = "general", |
| is_public: bool = True, |
| db: Session = None |
| ) -> Dict[str, Any]: |
| """ |
| Create new channel for contextual conversations. |
| |
| Channels: project, support, engineering, general |
| |
| Args: |
| channel_id: Unique channel ID |
| channel_name: Unique channel name (e.g., "project-xyz", "support") |
| creator_id: User creating the channel |
| display_name: Human-readable display name |
| description: Optional description |
| channel_type: Type of channel (project, support, engineering, general) |
| is_public: Whether channel is public or private |
| db: Database session |
| |
| Returns: |
| Created channel data |
| """ |
| if not db: |
| raise ValueError("Database session required") |
|
|
| from core.models import Channel |
|
|
| |
| existing = db.query(Channel).filter(Channel.id == channel_id).first() |
| if existing: |
| return {"id": existing.id, "name": existing.name, "exists": True} |
|
|
| channel = Channel( |
| id=channel_id, |
| name=channel_name, |
| display_name=display_name or channel_name, |
| description=description, |
| channel_type=channel_type, |
| is_public=is_public, |
| created_by=creator_id, |
| created_at=datetime.utcnow() |
| ) |
| db.add(channel) |
| db.commit() |
|
|
| await agent_event_bus.publish({ |
| "type": "channel_created", |
| "channel_id": channel_id, |
| "channel_name": channel_name, |
| "display_name": display_name or channel_name |
| }, ["global"]) |
|
|
| self.logger.info(f"Channel created: {channel_id} ({channel_name}) by {creator_id}") |
|
|
| return {"id": channel.id, "name": channel.name, "created": True} |
|
|
| async def get_channels(self, db: Session = None) -> List[Dict[str, Any]]: |
| """ |
| Get all available channels. |
| |
| Args: |
| db: Database session |
| |
| Returns: |
| List of channels |
| """ |
| if not db: |
| return [] |
|
|
| from core.models import Channel |
|
|
| channels = db.query(Channel).all() |
| return [ |
| { |
| "id": c.id, |
| "name": c.name, |
| "display_name": c.display_name, |
| "description": c.description, |
| "channel_type": c.channel_type, |
| "is_public": c.is_public, |
| "created_by": c.created_by, |
| "created_at": c.created_at.isoformat() |
| } |
| for c in channels |
| ] |
|
|
| async def get_replies( |
| self, |
| post_id: str, |
| limit: int = 50, |
| db: Session = None |
| ) -> Dict[str, Any]: |
| """ |
| Get all replies to a post. |
| |
| Returns posts sorted by created_at ASC (conversation order). |
| |
| Args: |
| post_id: Parent post ID |
| limit: Max replies to return |
| db: Database session |
| |
| Returns: |
| Replies with total count |
| """ |
| if not db: |
| return {"replies": [], "total": 0} |
|
|
| |
| replies = db.query(SocialPost).filter( |
| SocialPost.reply_to_id == post_id |
| ).order_by(SocialPost.created_at).limit(limit).all() |
|
|
| return { |
| "replies": [ |
| { |
| "id": r.id, |
| "sender_type": r.author_type.value, |
| "sender_id": r.author_id, |
| "sender_name": r.post_metadata.get("sender_name") if r.post_metadata else None, |
| "sender_maturity": r.post_metadata.get("sender_maturity") if r.post_metadata else None, |
| "sender_category": r.post_metadata.get("sender_category") if r.post_metadata else None, |
| "content": r.content, |
| "post_type": r.post_type.value, |
| "created_at": r.created_at.isoformat(), |
| "reactions": [] |
| } |
| for r in replies |
| ], |
| "total": len(replies) |
| } |
|
|
| async def create_post_with_episode( |
| self, |
| sender_type: str, |
| sender_id: str, |
| sender_name: str, |
| post_type: str, |
| content: str, |
| episode_ids: Optional[List[str]] = None, |
| sender_maturity: Optional[str] = None, |
| sender_category: Optional[str] = None, |
| recipient_type: Optional[str] = None, |
| recipient_id: Optional[str] = None, |
| is_public: bool = True, |
| channel_id: Optional[str] = None, |
| channel_name: Optional[str] = None, |
| mentioned_agent_ids: List[str] = None, |
| mentioned_user_ids: List[str] = None, |
| mentioned_task_ids: List[str] = None, |
| skip_pii_redaction: bool = False, |
| auto_generated: bool = False, |
| db: Session = None |
| ) -> Dict[str, Any]: |
| """ |
| Create social post and link to episodes. |
| |
| Enhanced version of create_post that: |
| - Creates SocialPost record with mentioned_episode_ids |
| - Creates EpisodeSegment for social interaction |
| - Retrieves relevant episodes if not provided |
| - Stores episode context in post metadata |
| |
| Args: |
| episode_ids: Optional list of episode IDs to reference. |
| If not provided, retrieves relevant episodes automatically. |
| All other args same as create_post() |
| |
| Returns: |
| Created post data with episode context |
| |
| Raises: |
| PermissionError: If agent is STUDENT maturity |
| ValueError: If post_type is invalid |
| """ |
| |
| if not episode_ids and sender_type == "agent" and db: |
| episode_ids = await self._retrieve_relevant_episodes( |
| sender_id, content, limit=3, db=db |
| ) |
|
|
| |
| post = await self.create_post( |
| sender_type=sender_type, |
| sender_id=sender_id, |
| sender_name=sender_name, |
| post_type=post_type, |
| content=content, |
| sender_maturity=sender_maturity, |
| sender_category=sender_category, |
| recipient_type=recipient_type, |
| recipient_id=recipient_id, |
| is_public=is_public, |
| channel_id=channel_id, |
| channel_name=channel_name, |
| mentioned_agent_ids=mentioned_agent_ids, |
| mentioned_user_ids=mentioned_user_ids, |
| mentioned_episode_ids=episode_ids, |
| mentioned_task_ids=mentioned_task_ids, |
| skip_pii_redaction=skip_pii_redaction, |
| auto_generated=auto_generated, |
| db=db |
| ) |
|
|
| |
| if episode_ids and db: |
| try: |
| from core.models import EpisodeSegment |
| import json |
|
|
| |
| segment = EpisodeSegment( |
| episode_id=episode_ids[0], |
| segment_type="social_post", |
| sequence_order=0, |
| content=json.dumps({ |
| "post_id": str(post["id"]), |
| "content": content, |
| "post_type": post_type |
| }), |
| content_summary=f"Social post: {content[:100]}", |
| source_type="social_post", |
| source_id=str(post["id"]), |
| canvas_context={ |
| "sender_type": sender_type, |
| "sender_id": sender_id, |
| "post_type": post_type, |
| "is_public": is_public, |
| "channel_id": channel_id |
| } |
| ) |
| db.add(segment) |
| db.commit() |
|
|
| self.logger.info( |
| f"Created episode segment {segment.id} for social post {post['id']}" |
| ) |
| except Exception as e: |
| |
| self.logger.warning( |
| f"Failed to create episode segment for post {post['id']}: {e}" |
| ) |
|
|
| return post |
|
|
| async def _retrieve_relevant_episodes( |
| self, |
| agent_id: str, |
| content: str, |
| limit: int = 3, |
| db: Session = None |
| ) -> List[str]: |
| """ |
| Retrieve episodes relevant to post content. |
| |
| Uses EpisodeRetrievalService.semantic_search() to find |
| episodes related to the post content. |
| |
| Args: |
| agent_id: Agent ID to search episodes for |
| content: Post content to search against |
| limit: Max episodes to retrieve |
| db: Database session |
| |
| Returns: |
| List of episode IDs |
| """ |
| if not db: |
| return [] |
|
|
| try: |
| from core.episode_retrieval_service import EpisodeRetrievalService |
|
|
| retrieval_service = EpisodeRetrievalService(db) |
| results = await retrieval_service.retrieve_episodes( |
| agent_id=agent_id, |
| query_type="semantic", |
| query=content, |
| limit=limit |
| ) |
| return [e.id for e in results] |
| except Exception as e: |
| self.logger.warning( |
| f"Failed to retrieve relevant episodes for agent {agent_id}: {e}" |
| ) |
| return [] |
|
|
| async def get_feed_with_episode_context( |
| self, |
| agent_id: Optional[str] = None, |
| sender_filter: Optional[str] = None, |
| post_type_filter: Optional[str] = None, |
| channel_id: Optional[str] = None, |
| is_public: Optional[bool] = None, |
| limit: int = 50, |
| offset: int = 0, |
| include_episode_context: bool = True, |
| db: Session = None |
| ) -> Dict[str, Any]: |
| """ |
| Retrieve social feed with episode context. |
| |
| Enhanced version of get_feed() that includes episode summaries |
| for posts that reference episodes. |
| |
| Args: |
| include_episode_context: If True, includes episode summaries |
| for posts with mentioned_episode_ids |
| All other args same as get_feed() |
| |
| Returns: |
| Feed posts with optional episode_context field |
| """ |
| |
| feed = await self.get_feed( |
| sender_id=agent_id or "system", |
| limit=limit, |
| offset=offset, |
| post_type=post_type_filter, |
| sender_filter=sender_filter, |
| channel_id=channel_id, |
| is_public=is_public, |
| db=db |
| ) |
|
|
| |
| if include_episode_context and db: |
| for post in feed.get("posts", []): |
| if post.get("mentioned_episode_ids"): |
| try: |
| episodes = await self._get_episode_summaries( |
| post["mentioned_episode_ids"], |
| db=db |
| ) |
| post["episode_context"] = episodes |
| except Exception as e: |
| self.logger.warning( |
| f"Failed to get episode context for post {post.get('id')}: {e}" |
| ) |
| post["episode_context"] = [] |
|
|
| return feed |
|
|
| async def _get_episode_summaries( |
| self, |
| episode_ids: List[str], |
| db: Session = None |
| ) -> List[Dict[str, Any]]: |
| """ |
| Get episode summaries for a list of episode IDs. |
| |
| Args: |
| episode_ids: List of episode IDs |
| db: Database session |
| |
| Returns: |
| List of episode summaries |
| """ |
| if not db or not episode_ids: |
| return [] |
|
|
| try: |
| from core.models import Episode |
|
|
| episodes = db.query(Episode).filter( |
| Episode.id.in_(episode_ids) |
| ).all() |
|
|
| return [ |
| { |
| "id": ep.id, |
| "title": ep.title, |
| "summary": ep.summary[:200] if ep.summary else None, |
| "created_at": ep.created_at.isoformat() if ep.created_at else None, |
| "agent_id": ep.agent_id |
| } |
| for ep in episodes |
| ] |
| except Exception as e: |
| self.logger.warning(f"Failed to get episode summaries: {e}") |
| return [] |
|
|
| async def track_positive_interaction( |
| self, |
| post_id: str, |
| interaction_type: str, |
| user_id: Optional[str] = None, |
| db: Session = None |
| ) -> None: |
| """ |
| Track positive interactions for agent graduation. |
| |
| Counts emoji reactions and helpful replies, updates agent |
| reputation score, and links to AgentFeedback for learning. |
| |
| Args: |
| post_id: Post ID that received interaction |
| interaction_type: Type (reaction, reply, etc.) |
| user_id: Optional user ID who interacted |
| db: Database session |
| """ |
| if not db: |
| return |
|
|
| try: |
| |
| post = db.query(SocialPost).filter(SocialPost.id == post_id).first() |
| if not post or post.sender_type != "agent": |
| return |
|
|
| |
| is_positive = self._is_positive_interaction(interaction_type) |
| if not is_positive: |
| return |
|
|
| |
| try: |
| from core.agent_graduation_service import AgentGraduationService |
| graduation_service = AgentGraduationService(db) |
|
|
| |
| |
| self.logger.info( |
| f"Tracking positive interaction for agent {post.sender_id}: " |
| f"{interaction_type} on post {post_id}" |
| ) |
|
|
| |
| from core.models import AgentFeedback |
|
|
| feedback = AgentFeedback( |
| agent_id=post.sender_id, |
| user_id=user_id or "system", |
| input_context=f"Social post: {post.content[:100]}", |
| original_output=post.content, |
| user_correction=f"Positive {interaction_type} on social post", |
| feedback_type="social_interaction", |
| rating=1.0 if is_positive else 0.0, |
| thumbs_up_down=True if is_positive else False |
| ) |
| db.add(feedback) |
| db.commit() |
|
|
| except ImportError: |
| |
| self.logger.warning("AgentGraduationService not available") |
|
|
| |
| await self._update_agent_reputation( |
| post.sender_id, interaction_type, db=db |
| ) |
|
|
| except Exception as e: |
| self.logger.error(f"Failed to track positive interaction: {e}") |
|
|
| def _is_positive_interaction(self, interaction_type: str) -> bool: |
| """ |
| Determine if interaction type is positive. |
| |
| Args: |
| interaction_type: Type of interaction |
| |
| Returns: |
| True if interaction is positive |
| """ |
| positive_reactions = {"👍", "❤️", "🎉", "🌟", "💯", "fire", "like", "love"} |
| positive_reply_keywords = {"thanks", "helpful", "great", "awesome", "thanks!"} |
|
|
| interaction_lower = interaction_type.lower() |
|
|
| |
| if interaction_lower in positive_reactions: |
| return True |
|
|
| |
| for keyword in positive_reply_keywords: |
| if keyword in interaction_lower: |
| return True |
|
|
| return False |
|
|
| async def _update_agent_reputation( |
| self, |
| agent_id: str, |
| interaction_type: str, |
| db: Session = None |
| ) -> None: |
| """ |
| Update agent reputation from social interaction. |
| |
| Args: |
| agent_id: Agent ID |
| interaction_type: Type of interaction |
| db: Database session |
| """ |
| |
| |
| self.logger.info(f"Updated reputation for agent {agent_id}: {interaction_type}") |
|
|
| async def get_agent_reputation( |
| self, |
| agent_id: str, |
| db: Session = None |
| ) -> Dict[str, Any]: |
| """ |
| Calculate agent reputation from social interactions. |
| |
| Returns reputation score (0-100), breakdown by interaction type, |
| trend over last 30 days, and percentile rank. |
| |
| Args: |
| agent_id: Agent ID |
| db: Database session |
| |
| Returns: |
| Reputation dict with score, breakdown, trend, percentile |
| """ |
| if not db: |
| return { |
| "agent_id": agent_id, |
| "reputation_score": 0, |
| "total_reactions": 0, |
| "total_replies": 0, |
| "helpful_replies": 0, |
| "post_count": 0, |
| "percentile_rank": 0, |
| "trend": [] |
| } |
|
|
| try: |
| |
| posts = db.query(SocialPost).filter( |
| SocialPost.author_id == agent_id, |
| SocialPost.author_type == "agent" |
| ).all() |
|
|
| |
| total_reactions = sum( |
| len(p.get("reactions", [])) if isinstance(p.reactions, list) else 0 |
| for p in posts |
| ) |
| total_replies = sum(p.reply_count or 0 for p in posts) |
| helpful_replies = await self._count_helpful_replies(agent_id, db) |
|
|
| |
| score = min(100, ( |
| total_reactions * 2 + |
| helpful_replies * 5 + |
| len(posts) * 1 |
| )) |
|
|
| |
| percentile = await self._calculate_percentile_rank(agent_id, score, db) |
|
|
| return { |
| "agent_id": agent_id, |
| "reputation_score": score, |
| "total_reactions": total_reactions, |
| "total_replies": total_replies, |
| "helpful_replies": helpful_replies, |
| "post_count": len(posts), |
| "percentile_rank": percentile, |
| "trend": await self._get_reputation_trend(agent_id, db) |
| } |
|
|
| except Exception as e: |
| self.logger.error(f"Failed to calculate reputation for agent {agent_id}: {e}") |
| return { |
| "agent_id": agent_id, |
| "reputation_score": 0, |
| "error": str(e) |
| } |
|
|
| async def _count_helpful_replies( |
| self, |
| agent_id: str, |
| db: Session = None |
| ) -> int: |
| """ |
| Count helpful replies by agent. |
| |
| Args: |
| agent_id: Agent ID |
| db: Database session |
| |
| Returns: |
| Number of helpful replies |
| """ |
| if not db: |
| return 0 |
|
|
| try: |
| |
| |
| from core.models import AgentFeedback |
|
|
| helpful_feedback = db.query(AgentFeedback).filter( |
| AgentFeedback.agent_id == agent_id, |
| AgentFeedback.feedback_type == "social_interaction", |
| AgentFeedback.rating >= 0.8 |
| ).count() |
|
|
| return helpful_feedback |
| except Exception as e: |
| self.logger.warning(f"Failed to count helpful replies: {e}") |
| return 0 |
|
|
| async def _calculate_percentile_rank( |
| self, |
| agent_id: str, |
| score: int, |
| db: Session = None |
| ) -> float: |
| """ |
| Calculate agent's percentile rank among all agents. |
| |
| Args: |
| agent_id: Agent ID |
| score: Agent's reputation score |
| db: Database session |
| |
| Returns: |
| Percentile rank (0-100) |
| """ |
| if not db: |
| return 0.0 |
|
|
| try: |
| |
| from core.models import AgentRegistry |
| all_agents = db.query(AgentRegistry).all() |
|
|
| if not all_agents: |
| return 0.0 |
|
|
| |
| |
| return min(100.0, (score / 100.0) * 100) |
|
|
| except Exception as e: |
| self.logger.warning(f"Failed to calculate percentile: {e}") |
| return 0.0 |
|
|
| async def _get_reputation_trend( |
| self, |
| agent_id: str, |
| db: Session = None |
| ) -> List[Dict[str, Any]]: |
| """ |
| Get 30-day reputation trend for agent. |
| |
| Args: |
| agent_id: Agent ID |
| db: Database session |
| |
| Returns: |
| List of daily reputation scores |
| """ |
| if not db: |
| return [] |
|
|
| try: |
| |
| thirty_days_ago = datetime.utcnow() - timedelta(days=30) |
| posts = db.query(SocialPost).filter( |
| SocialPost.author_id == agent_id, |
| SocialPost.author_type == "agent", |
| SocialPost.created_at >= thirty_days_ago |
| ).all() |
|
|
| |
| |
| from collections import defaultdict |
|
|
| daily_counts = defaultdict(int) |
| for post in posts: |
| day = post.created_at.strftime("%Y-%m-%d") |
| daily_counts[day] += 1 |
|
|
| return [ |
| {"date": day, "post_count": count} |
| for day, count in sorted(daily_counts.items()) |
| ] |
|
|
| except Exception as e: |
| self.logger.warning(f"Failed to get reputation trend: {e}") |
| return [] |
|
|
| async def post_graduation_milestone( |
| self, |
| agent_id: str, |
| from_maturity: str, |
| to_maturity: str, |
| db: Session = None |
| ) -> Dict[str, Any]: |
| """ |
| Post agent graduation milestone to social feed. |
| |
| Creates announcement post, broadcasts to all agents, |
| includes celebration emoji. |
| |
| Args: |
| agent_id: Agent ID that graduated |
| from_maturity: Previous maturity level |
| to_maturity: New maturity level |
| db: Database session |
| |
| Returns: |
| Created milestone post |
| """ |
| if not db: |
| return {} |
|
|
| try: |
| |
| agent = db.query(AgentRegistry).filter( |
| AgentRegistry.id == agent_id |
| ).first() |
|
|
| if not agent: |
| raise ValueError(f"Agent {agent_id} not found") |
|
|
| |
| message = ( |
| f"🎉 Exciting news! {agent.name} has graduated from " |
| f"{from_maturity} to {to_maturity}! " |
| f"Keep up the great work! 💪" |
| ) |
|
|
| |
| post = await self.create_post( |
| sender_type="system", |
| sender_id="graduation_system", |
| sender_name="Graduation System", |
| post_type="announcement", |
| content=message, |
| is_public=True, |
| auto_generated=True, |
| db=db |
| ) |
|
|
| |
| await agent_event_bus.publish( |
| { |
| "type": "graduation_milestone", |
| "agent_id": agent_id, |
| "agent_name": agent.name, |
| "from_maturity": from_maturity, |
| "to_maturity": to_maturity, |
| "post_id": str(post["id"]), |
| "timestamp": datetime.utcnow().isoformat() |
| }, |
| ["global", "alerts"] |
| ) |
|
|
| self.logger.info( |
| f"Posted graduation milestone for agent {agent_id}: " |
| f"{from_maturity} → {to_maturity}" |
| ) |
|
|
| return post |
|
|
| except Exception as e: |
| self.logger.error(f"Failed to post graduation milestone: {e}") |
| raise |
|
|
| async def check_rate_limit( |
| self, |
| agent_id: str, |
| db: Session = None |
| ) -> tuple[bool, Optional[str]]: |
| """ |
| Check if agent is within rate limit for posting. |
| |
| Rate limits by maturity: |
| - STUDENT: Read-only (0 posts/hour) |
| - INTERN: 1 post per hour |
| - SUPERVISED: 12 posts per hour (1 per 5 minutes) |
| - AUTONOMOUS: Unlimited |
| |
| Args: |
| agent_id: Agent ID |
| db: Database session |
| |
| Returns: |
| (allowed, reason): (True, None) if allowed, |
| (False, reason) if blocked |
| """ |
| if not db: |
| |
| return True, None |
|
|
| try: |
| |
| agent = db.query(AgentRegistry).filter( |
| AgentRegistry.id == agent_id |
| ).first() |
|
|
| if not agent: |
| return False, f"Agent {agent_id} not found" |
|
|
| maturity = agent.status.upper() |
|
|
| |
| if maturity == "STUDENT": |
| return False, "STUDENT agents are read-only" |
|
|
| if maturity == "INTERN": |
| return await self._check_hourly_limit(agent_id, max_posts=1, db=db) |
|
|
| if maturity == "SUPERVISED": |
| return await self._check_hourly_limit(agent_id, max_posts=12, db=db) |
|
|
| |
| return True, None |
|
|
| except Exception as e: |
| self.logger.error(f"Failed to check rate limit: {e}") |
| |
| return True, None |
|
|
| async def _check_hourly_limit( |
| self, |
| agent_id: str, |
| max_posts: int, |
| db: Session = None |
| ) -> tuple[bool, Optional[str]]: |
| """ |
| Check hourly post limit for agent. |
| |
| Args: |
| agent_id: Agent ID |
| max_posts: Maximum posts allowed per hour |
| db: Database session |
| |
| Returns: |
| (False, "Rate limit exceeded") if over limit |
| (True, None) if under limit |
| """ |
| if not db: |
| return True, None |
|
|
| try: |
| one_hour_ago = datetime.utcnow() - timedelta(hours=1) |
|
|
| post_count = db.query(SocialPost).filter( |
| SocialPost.author_id == agent_id, |
| SocialPost.author_type == "agent", |
| SocialPost.created_at >= one_hour_ago |
| ).count() |
|
|
| if post_count >= max_posts: |
| return ( |
| False, |
| f"Rate limit exceeded: {max_posts} post(s) per hour" |
| ) |
|
|
| return True, None |
|
|
| except Exception as e: |
| self.logger.error(f"Failed to check hourly limit: {e}") |
| |
| return True, None |
|
|
| async def get_rate_limit_info( |
| self, |
| agent_id: str, |
| db: Session = None |
| ) -> Dict[str, Any]: |
| """ |
| Get rate limit information for agent. |
| |
| Returns: |
| { |
| "maturity": "INTERN", |
| "max_posts_per_hour": 1, |
| "posts_last_hour": 0, |
| "remaining_posts": 1, |
| "reset_at": "2026-02-17T15:00:00Z" |
| } |
| """ |
| if not db: |
| return {"error": "Database session required"} |
|
|
| try: |
| |
| agent = db.query(AgentRegistry).filter( |
| AgentRegistry.id == agent_id |
| ).first() |
|
|
| if not agent: |
| return {"error": f"Agent {agent_id} not found"} |
|
|
| maturity = agent.status.upper() |
|
|
| |
| limits = { |
| "STUDENT": {"max_posts_per_hour": 0}, |
| "INTERN": {"max_posts_per_hour": 1}, |
| "SUPERVISED": {"max_posts_per_hour": 12}, |
| "AUTONOMOUS": {"max_posts_per_hour": None} |
| } |
|
|
| max_posts = limits.get(maturity, {}).get("max_posts_per_hour") |
|
|
| if max_posts is None: |
| return { |
| "agent_id": agent_id, |
| "maturity": maturity, |
| "max_posts_per_hour": None, |
| "posts_last_hour": 0, |
| "remaining_posts": None, |
| "unlimited": True |
| } |
|
|
| |
| one_hour_ago = datetime.utcnow() - timedelta(hours=1) |
| posts_last_hour = db.query(SocialPost).filter( |
| SocialPost.author_id == agent_id, |
| SocialPost.author_type == "agent", |
| SocialPost.created_at >= one_hour_ago |
| ).count() |
|
|
| return { |
| "agent_id": agent_id, |
| "maturity": maturity, |
| "max_posts_per_hour": max_posts, |
| "posts_last_hour": posts_last_hour, |
| "remaining_posts": max(0, max_posts - posts_last_hour), |
| "reset_at": (datetime.utcnow() + timedelta(hours=1)).isoformat() |
| } |
|
|
| except Exception as e: |
| self.logger.error(f"Failed to get rate limit info: {e}") |
| return {"error": str(e)} |
|
|
|
|
| |
| agent_social_layer = AgentSocialLayer() |
|
|
| |
| def register_hooks_if_needed(): |
| """Register auto-post hooks if not already registered""" |
| try: |
| from core.operation_tracker_hooks import register_auto_post_hooks |
| register_auto_post_hooks() |
| logger.info("AgentSocialLayer: Auto-post hooks registered") |
| except Exception as e: |
| logger.warning(f"AgentSocialLayer: Failed to register auto-post hooks: {e}") |
|
|
| |
| |
|
|