diff --git a/handlers/src/config.ts b/handlers/src/config.ts index 1b1da74..a37ac11 100644 --- a/handlers/src/config.ts +++ b/handlers/src/config.ts @@ -93,3 +93,4 @@ export const REVIEW_DRAFTS = `${DRAFTS_DIR}/review`; export const REJECTED_DIR = `${DRAFTS_DIR}/rejected`; export const NOTES_DIR = `${DRAFTS_DIR}/notes`; export const PUBLISHED_DIR = `${DRAFTS_DIR}/published`; +export const PROCESSED_DIR = `${DRAFTS_DIR}/processed`; // Marker files for items sent to LLM (prevents re-processing) diff --git a/handlers/src/notification-handler.ts b/handlers/src/notification-handler.ts index 7d21007..4e0dac4 100644 --- a/handlers/src/notification-handler.ts +++ b/handlers/src/notification-handler.ts @@ -20,6 +20,7 @@ import { BLUESKY_DRAFTS, REVIEW_DRAFTS, PUBLISHED_DIR, + PROCESSED_DIR, INDEXER_URL, } from "./config.js"; @@ -40,6 +41,9 @@ interface QueueItem { * Check if we've already processed this notification */ function alreadyProcessed(id: string): boolean { + // Check if already sent to LLM (marker from previous run) + if (fs.existsSync(path.join(PROCESSED_DIR, `${id}.marker`))) return true; + const patterns = [ path.join(BLUESKY_DRAFTS, `reply-${id}.txt`), path.join(REVIEW_DRAFTS, `bluesky-reply-${id}.txt`), @@ -56,6 +60,17 @@ function alreadyProcessed(id: string): boolean { return false; } +/** + * Mark items as processed (sent to LLM). Prevents re-processing + * even if LLM chose not to draft a response. + */ +function markProcessed(ids: string[]): void { + fs.mkdirSync(PROCESSED_DIR, { recursive: true }); + for (const id of ids) { + fs.writeFileSync(path.join(PROCESSED_DIR, `${id}.marker`), ""); + } +} + /** * Chunk array into batches */ @@ -301,16 +316,22 @@ Write each draft file. Skip what should be skipped.`; await Promise.race([processPromise(), timeoutPromise]); session.close(); + + // Mark ALL items in this batch as processed (even if LLM chose not to respond) + const batchIds = batch.map(item => item.uri?.split("/").pop() || item.cid); + markProcessed(batchIds); + console.log(`Marked ${batchIds.length} items as processed`); } catch (error) { if (error instanceof Error && error.message.includes("timeout")) { console.error(`\n⚠ Batch ${i + 1} timed out. Items will retry next run.`); } else { console.error("Error invoking Central:", error); } + // Don't mark as processed on error - let them retry } } - // Clean up queue: remove items that have been processed (have draft files) + // Clean up queue: remove items that have been processed or are SKIP const remainingQueue = queue.filter(item => { const id = item.uri?.split("/").pop() || item.cid; return !alreadyProcessed(id) && item.priority !== "SKIP"; diff --git a/handlers/src/x-handler.ts b/handlers/src/x-handler.ts index 7a62915..94ca41a 100644 --- a/handlers/src/x-handler.ts +++ b/handlers/src/x-handler.ts @@ -17,6 +17,7 @@ import { X_DRAFTS, REVIEW_DRAFTS, PUBLISHED_DIR, + PROCESSED_DIR, CENTRAL_AGENT_ID, INDEXER_URL, } from "./config.js"; @@ -36,6 +37,9 @@ interface XQueueItem { * Check if we've already processed this notification */ function alreadyProcessed(id: string): boolean { + // Check if already sent to LLM (marker from previous run) + if (fs.existsSync(path.join(PROCESSED_DIR, `${id}.marker`))) return true; + const patterns = [ path.join(X_DRAFTS, `reply-${id}.txt`), path.join(REVIEW_DRAFTS, `x-reply-${id}.txt`), @@ -56,6 +60,17 @@ function alreadyProcessed(id: string): boolean { return false; } +/** + * Mark items as processed (sent to LLM). Prevents re-processing + * even if LLM chose not to draft a response. + */ +function markProcessed(ids: string[]): void { + fs.mkdirSync(PROCESSED_DIR, { recursive: true }); + for (const id of ids) { + fs.writeFileSync(path.join(PROCESSED_DIR, `${id}.marker`), ""); + } +} + /** * Chunk array into batches */ @@ -186,16 +201,21 @@ Write each draft file. Skip what should be skipped.`; await Promise.race([processPromise(), timeoutPromise]); session.close(); + + // Mark ALL items in this batch as processed (even if LLM chose not to respond) + markProcessed(batch.map(item => item.id)); + console.log(`Marked ${batch.length} items as processed`); } catch (error) { if (error instanceof Error && error.message.includes("timeout")) { console.error(`\n⚠ Batch ${i + 1} timed out. Items will retry next run.`); } else { console.error("Error invoking Central:", error); } + // Don't mark as processed on error - let them retry } } - // Clean up queue: remove items that have been processed (have draft files) + // Clean up queue: remove items that have been processed or are SKIP const remainingQueue = queue.filter(item => { return !alreadyProcessed(item.id) && item.priority !== "SKIP"; }); diff --git a/tools/responder.py b/tools/responder.py index 331b4e7..9f9bc18 100644 --- a/tools/responder.py +++ b/tools/responder.py @@ -20,6 +20,7 @@ console = Console() DRAFTS_FILE = Path("drafts/queue.yaml") SENT_FILE = Path("drafts/sent.txt") # Track URIs we've replied to MENTIONS_LOG = Path("logs/mentions.jsonl") +SEEN_AT_FILE = Path("drafts/notification_seen_at.txt") # Track newest notification timestamp QUEUE_TTL_HOURS = 24 # Auto-remove items older than this # Priority system (matches daemon.py) @@ -184,6 +185,16 @@ async def queue_notifications(limit=200): return notifications = resp.json().get("notifications", []) + + # Filter out notifications older than our cursor (avoids reprocessing) + seen_at_cursor = SEEN_AT_FILE.read_text().strip() if SEEN_AT_FILE.exists() else None + if seen_at_cursor: + before_count = len(notifications) + notifications = [n for n in notifications if n.get("indexedAt", "") > seen_at_cursor] + skipped = before_count - len(notifications) + if skipped > 0: + console.print(f"[dim]Skipped {skipped} notifications older than cursor[/dim]") + # Sort: Cameron first, then by time (ensures CRITICAL never gets buried by volume) notifications.sort(key=lambda n: (0 if n.get("author", {}).get("did") == CAMERON_DID else 1)) queue = [] @@ -211,19 +222,21 @@ async def queue_notifications(limit=200): count = 0 + # Track newest notification timestamp for incremental fetching + newest_indexed_at = None + for n in notifications: + # Track newest for seenAt cursor + indexed_at = n.get("indexedAt") + if indexed_at and (newest_indexed_at is None or indexed_at > newest_indexed_at): + newest_indexed_at = indexed_at + if n["reason"] not in ["mention", "reply"]: continue if n["uri"] in existing_uris: continue if n["uri"] in sent_uris: continue # Already replied to this - if n.get("isRead", False): # Only fetch unread? Maybe allow fetching recent read ones too? - # For now, let's include even read ones if they aren't in the queue, - # but usually we want to clear the queue. - # Actually, let's stick to unread to avoid noise, OR allow a --all flag. - # Defaulting to unread only for safety. - pass # Fetch context (parent/root) for threading # We need the post record to get reply refs @@ -283,6 +296,10 @@ async def queue_notifications(limit=200): with open(DRAFTS_FILE, "w") as f: yaml.dump(queue, f, sort_keys=False, indent=2) + # Update seenAt cursor for next run + if newest_indexed_at: + SEEN_AT_FILE.write_text(newest_indexed_at) + console.print(f"[green]Queued {count} new notifications.[/green]") console.print(f"Edit {DRAFTS_FILE} to draft responses.") diff --git a/tools/x_responder.py b/tools/x_responder.py index d21ac24..f57e327 100644 --- a/tools/x_responder.py +++ b/tools/x_responder.py @@ -53,6 +53,7 @@ def is_low_effort(text: str) -> bool: QUEUE_PATH = Path("drafts/x_queue.yaml") SENT_PATH = Path("drafts/x_sent.txt") MENTIONS_LOG = Path("logs/x_mentions.jsonl") +SINCE_ID_PATH = Path("drafts/x_since_id.txt") # Track newest mention seen def _log_mention(entry: dict): @@ -119,16 +120,36 @@ def save_sent_id(tweet_id: str): f.write(f"{tweet_id}\n") +def load_since_id() -> str | None: + """Load the most recent mention ID we've seen.""" + if SINCE_ID_PATH.exists(): + return SINCE_ID_PATH.read_text().strip() + return None + + +def save_since_id(tweet_id: str): + """Save the most recent mention ID we've seen.""" + SINCE_ID_PATH.write_text(tweet_id) + + def fetch_mentions(limit: int = 20) -> list[dict]: - """Fetch recent mentions.""" + """Fetch recent mentions newer than since_id.""" client = get_client() me = client.get_me() - + + # Only fetch mentions newer than last seen + since_id = load_since_id() + params = { + "max_results": min(limit, 100), + "tweet_fields": ["created_at", "author_id", "conversation_id"], + "expansions": ["author_id"], + } + if since_id: + params["since_id"] = since_id + response = client.get_users_mentions( me.data.id, - max_results=min(limit, 100), - tweet_fields=["created_at", "author_id", "conversation_id"], - expansions=["author_id"], + **params ) if not response.data: @@ -150,7 +171,11 @@ def fetch_mentions(limit: int = 20) -> list[dict]: "created_at": str(tweet.created_at) if tweet.created_at else None, "conversation_id": str(tweet.conversation_id) if tweet.conversation_id else None, }) - + + # Save newest mention ID for next run (mentions are newest-first) + if mentions: + save_since_id(mentions[0]["id"]) + return mentions