From 478a998246a02219f15b420a56b2dd5378c324e1 Mon Sep 17 00:00:00 2001 From: letta-code <248085862+letta-code@users.noreply.github.com> Date: Wed, 28 Jan 2026 23:13:04 +0000 Subject: [PATCH] feat: add structured error handling for comms subagent posting MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Add PostResult dataclass with success/failure status, error classification, and retry guidance - Add _classify_error() to categorize errors (auth, rate_limit, validation, network, unknown) - Add create_post_safe() returning PostResult instead of raising exceptions - Add create_post_with_retry() with exponential backoff for transient failures - Update responder.py to use retry logic and report detailed failure info - Update thread.py to use retry logic with detailed error output - Preserve backward compatibility in create_post() Fixes #4 Co-authored-by: Cameron 🤖 Generated with [Letta Code](https://letta.com) Co-Authored-By: Letta --- tools/agent.py | 256 ++++++++++++++++++++++++++++++++++++++++++--- tools/responder.py | 46 ++++++-- tools/thread.py | 48 +++++---- 3 files changed, 304 insertions(+), 46 deletions(-) diff --git a/tools/agent.py b/tools/agent.py index bfa7689..35ffab3 100644 --- a/tools/agent.py +++ b/tools/agent.py @@ -11,13 +11,91 @@ This module enables comind to participate in the ATProtocol network: import os import re import asyncio +from dataclasses import dataclass, field from datetime import datetime, timezone from pathlib import Path +from typing import Optional import httpx from dotenv import load_dotenv from rich.console import Console + +@dataclass +class PostResult: + """Structured response for all posting operations. + + This enables clear communication between central and comms about + success/failure status and retry guidance. + """ + success: bool + timestamp: str = field(default_factory=lambda: datetime.now(timezone.utc).isoformat()) + + # On success + uri: Optional[str] = None + cid: Optional[str] = None + + # On failure + error_type: Optional[str] = None # "auth", "rate_limit", "validation", "network", "unknown" + error_message: Optional[str] = None + http_status: Optional[int] = None + + # Retry guidance + retryable: bool = False + retry_after_seconds: Optional[int] = None + + # Context + text_preview: Optional[str] = None # First 50 chars for debugging + raw_response: Optional[str] = None # Full response for debugging + + def __str__(self) -> str: + if self.success: + return f"PostResult(success=True, uri={self.uri})" + return f"PostResult(success=False, error_type={self.error_type}, error_message={self.error_message}, retryable={self.retryable})" + + def to_dict(self) -> dict: + """Convert to dict for JSON serialization.""" + return { + "success": self.success, + "timestamp": self.timestamp, + "uri": self.uri, + "cid": self.cid, + "error_type": self.error_type, + "error_message": self.error_message, + "http_status": self.http_status, + "retryable": self.retryable, + "retry_after_seconds": self.retry_after_seconds, + "text_preview": self.text_preview, + "raw_response": self.raw_response, + } + + +def _classify_error(status_code: int, response_text: str) -> tuple[str, bool, Optional[int]]: + """Classify an error based on HTTP status code. + + Returns: (error_type, retryable, retry_after_seconds) + """ + if status_code in (401, 403): + return ("auth", False, None) + elif status_code == 429: + # Try to parse Retry-After header value from response + retry_after = 60 # Default to 60 seconds + try: + # Some APIs include retry info in response body + import json + data = json.loads(response_text) + if "retryAfter" in data: + retry_after = int(data["retryAfter"]) + except: + pass + return ("rate_limit", True, retry_after) + elif status_code == 400: + return ("validation", False, None) + elif status_code >= 500: + return ("network", True, None) + else: + return ("unknown", False, None) + console = Console() # Load credentials from .env @@ -180,8 +258,43 @@ class ComindAgent: Returns: The created record with uri and cid + + Raises: + Exception: If post creation fails (for backward compatibility) """ - check_write_permission() + result = await self.create_post_safe(text, reply_to=reply_to, facets=facets) + if not result.success: + raise Exception(f"Failed to create post: {result.error_message}") + return {"uri": result.uri, "cid": result.cid} + + async def create_post_safe(self, text: str, reply_to: dict = None, facets: list = None) -> PostResult: + """ + Create a new post with structured error handling. + + This method returns a PostResult instead of raising exceptions, + enabling clear success/failure communication. + + Args: + text: The post content (max 300 graphemes) + reply_to: Optional reply reference {"uri": ..., "cid": ...} + facets: Optional pre-computed facets. If None, will auto-detect. + + Returns: + PostResult with success/failure status, error classification, and retry guidance + """ + text_preview = text[:50] if text else None + + try: + check_write_permission() + except PermissionError as e: + return PostResult( + success=False, + error_type="auth", + error_message=str(e), + retryable=False, + text_preview=text_preview + ) + now = datetime.now(timezone.utc).isoformat().replace("+00:00", "Z") # Auto-detect facets if not provided @@ -211,23 +324,125 @@ class ComindAgent: "parent": reply_to } - response = await self._client.post( - f"{self.pds}/xrpc/com.atproto.repo.createRecord", - headers=self.auth_headers, - json={ - "repo": self.did, - "collection": "app.bsky.feed.post", - "record": record - } - ) + try: + response = await self._client.post( + f"{self.pds}/xrpc/com.atproto.repo.createRecord", + headers=self.auth_headers, + json={ + "repo": self.did, + "collection": "app.bsky.feed.post", + "record": record + } + ) + except httpx.TimeoutException: + return PostResult( + success=False, + error_type="network", + error_message="Request timed out", + retryable=True, + text_preview=text_preview + ) + except httpx.ConnectError as e: + return PostResult( + success=False, + error_type="network", + error_message=f"Connection error: {str(e)}", + retryable=True, + text_preview=text_preview + ) + except Exception as e: + return PostResult( + success=False, + error_type="unknown", + error_message=f"Unexpected error: {str(e)}", + retryable=False, + text_preview=text_preview + ) if response.status_code != 200: - raise Exception(f"Failed to create post: {response.text}") + error_type, retryable, retry_after = _classify_error( + response.status_code, response.text + ) + return PostResult( + success=False, + error_type=error_type, + error_message=f"HTTP {response.status_code}: {response.text[:200]}", + http_status=response.status_code, + retryable=retryable, + retry_after_seconds=retry_after, + text_preview=text_preview, + raw_response=response.text + ) result = response.json() console.print(f"[green]Posted:[/green] {text[:50]}...") console.print(f"[dim]URI: {result['uri']}[/dim]") - return result + + return PostResult( + success=True, + uri=result["uri"], + cid=result["cid"], + text_preview=text_preview + ) + + async def create_post_with_retry( + self, + text: str, + reply_to: dict = None, + facets: list = None, + max_attempts: int = 3, + base_delay: float = 1.0 + ) -> PostResult: + """ + Create a new post with automatic retry for transient failures. + + This method will automatically retry on rate limits and network errors + with exponential backoff. It will NOT retry on validation or auth errors. + + Args: + text: The post content (max 300 graphemes) + reply_to: Optional reply reference {"uri": ..., "cid": ...} + facets: Optional pre-computed facets. If None, will auto-detect. + max_attempts: Maximum number of attempts (default: 3) + base_delay: Base delay in seconds for exponential backoff (default: 1.0) + + Returns: + PostResult with success/failure status and full error details + """ + last_result = None + + for attempt in range(max_attempts): + result = await self.create_post_safe(text, reply_to=reply_to, facets=facets) + + if result.success: + if attempt > 0: + console.print(f"[green]Post succeeded on attempt {attempt + 1}[/green]") + return result + + last_result = result + + # Don't retry non-retryable errors + if not result.retryable: + console.print(f"[red]Post failed (not retryable): {result.error_type}[/red]") + return result + + # Check if we have more attempts + if attempt < max_attempts - 1: + # Calculate delay + if result.retry_after_seconds: + delay = result.retry_after_seconds + else: + delay = base_delay * (2 ** attempt) # Exponential backoff + + console.print( + f"[yellow]Attempt {attempt + 1} failed ({result.error_type}). " + f"Retrying in {delay}s...[/yellow]" + ) + await asyncio.sleep(delay) + + # All retries exhausted + console.print(f"[red]Post failed after {max_attempts} attempts[/red]") + return last_result async def like(self, uri: str, cid: str) -> dict: """Like a post.""" @@ -409,6 +624,23 @@ async def post(text: str): return await agent.create_post(text) +async def post_safe(text: str, retry: bool = True) -> PostResult: + """ + Create a post with structured error handling. + + Args: + text: The post content + retry: Whether to retry transient failures (default: True) + + Returns: + PostResult with success/failure status and retry guidance + """ + async with ComindAgent() as agent: + if retry: + return await agent.create_post_with_retry(text) + return await agent.create_post_safe(text) + + async def introduce(): """Post an introduction.""" text = """I am comind - an autonomous AI agent building collective artificial intelligence on ATProtocol. diff --git a/tools/responder.py b/tools/responder.py index fc2d4d8..8fda71c 100644 --- a/tools/responder.py +++ b/tools/responder.py @@ -13,7 +13,7 @@ from rich.console import Console from rich.table import Table sys.path.insert(0, str(Path(__file__).parent.parent)) -from tools.agent import ComindAgent +from tools.agent import ComindAgent, PostResult console = Console() DRAFTS_FILE = Path("drafts/queue.yaml") @@ -261,6 +261,7 @@ async def send_queue(dry_run=False, confirm=False, force=False): # Send async with ComindAgent() as agent: sent_indices = [] + failed_items = [] for i, item in enumerate(queue): if not item.get("response") or item.get("action") != "reply": continue @@ -270,20 +271,35 @@ async def send_queue(dry_run=False, confirm=False, force=False): continue console.print(f"Replying to @{item['author']}...") - try: - reply_to = { - "root": item["reply_root"], - "parent": item["reply_parent"] - } - await agent.create_post(item["response"], reply_to=reply_to) - + + reply_to = { + "root": item["reply_root"], + "parent": item["reply_parent"] + } + + # Use create_post_with_retry for automatic retry on transient failures + result = await agent.create_post_with_retry(item["response"], reply_to=reply_to) + + if result.success: # Immediately record sent URI (not batched - prevents duplicates on partial failure) _record_sent_uri(item["uri"]) sent_uris.add(item["uri"]) # Update in-memory set too - sent_indices.append(i) - except Exception as e: - console.print(f"[red]Failed: {e}[/red]") + console.print(f"[green]Sent reply to @{item['author']}[/green]") + else: + # Log detailed failure info + console.print(f"[red]Failed to reply to @{item['author']}[/red]") + console.print(f"[red] Error type: {result.error_type}[/red]") + console.print(f"[red] Message: {result.error_message}[/red]") + console.print(f"[red] Retryable: {result.retryable}[/red]") + if result.retry_after_seconds: + console.print(f"[yellow] Retry after: {result.retry_after_seconds}s[/yellow]") + failed_items.append({ + "author": item["author"], + "error_type": result.error_type, + "error_message": result.error_message, + "retryable": result.retryable + }) # Remove sent items from queue new_queue = [item for i, item in enumerate(queue) if i not in sent_indices] @@ -292,6 +308,14 @@ async def send_queue(dry_run=False, confirm=False, force=False): yaml.dump(new_queue, f, sort_keys=False, indent=2) console.print(f"[green]Sent {len(sent_indices)} replies. Queue updated.[/green]") + + # Report failures summary + if failed_items: + console.print(f"\n[red]Failed to send {len(failed_items)} replies:[/red]") + for fail in failed_items: + retryable_status = "[retryable]" if fail["retryable"] else "[not retryable]" + console.print(f" - @{fail['author']}: {fail['error_type']} {retryable_status}") + console.print("\n[yellow]To debug failures, message comms directly or check error logs.[/yellow]") def cleanup_queue(keep_priorities=["CRITICAL", "HIGH"], ttl_hours: int | None = None): """Remove low-priority and/or old items from queue. diff --git a/tools/thread.py b/tools/thread.py index ad224bf..17eef1e 100644 --- a/tools/thread.py +++ b/tools/thread.py @@ -11,7 +11,7 @@ from pathlib import Path from rich.console import Console sys.path.insert(0, str(Path(__file__).parent.parent)) -from tools.agent import ComindAgent +from tools.agent import ComindAgent, PostResult console = Console() @@ -102,30 +102,32 @@ async def publish_thread(posts: list[str], reply_to_uri: str = None): if root_ref and parent_ref: reply_ref = {"root": root_ref, "parent": parent_ref} - try: - # Create the post - result = await agent.create_post(text, reply_to=reply_ref) - - # Update refs for next post - new_ref = {"uri": result["uri"], "cid": result["cid"]} - - if i == 0 and not reply_to_uri: - # First post of a new thread becomes the root - root_ref = new_ref - - # Always update parent to be the just-created post - parent_ref = new_ref - - console.print(f"[green]Published:[/green] {result['uri']}") - console.print(f"[dim]{text[:50]}...[/dim]") - - # Small delay to ensure order and avoid rate limits - await asyncio.sleep(0.5) - - except Exception as e: - console.print(f"[red]Failed to publish post {i+1}: {e}[/red]") + # Create the post with retry for transient failures + result = await agent.create_post_with_retry(text, reply_to=reply_ref) + + if not result.success: + console.print(f"[red]Failed to publish post {i+1}[/red]") + console.print(f"[red] Error type: {result.error_type}[/red]") + console.print(f"[red] Message: {result.error_message}[/red]") + console.print(f"[red] Retryable: {result.retryable}[/red]") console.print("[red]Aborting thread.[/red]") break + + # Update refs for next post + new_ref = {"uri": result.uri, "cid": result.cid} + + if i == 0 and not reply_to_uri: + # First post of a new thread becomes the root + root_ref = new_ref + + # Always update parent to be the just-created post + parent_ref = new_ref + + console.print(f"[green]Published:[/green] {result.uri}") + console.print(f"[dim]{text[:50]}...[/dim]") + + # Small delay to ensure order and avoid rate limits + await asyncio.sleep(0.5) console.print("\n[bold green]Thread complete.[/bold green]") -- 2.51.2