diff --git a/data/activity_queue/sample_cameron_stream_20260129_201008.json b/data/activity_queue/sample_cameron_stream_20260129_201008.json new file mode 100644 index 0000000..1c28e57 --- /dev/null +++ b/data/activity_queue/sample_cameron_stream_20260129_201008.json @@ -0,0 +1,109 @@ +[ + { + "timestamp": "2026-01-30T02:13:39.181Z", + "did": "did:plc:gfrmhdmjvxn2sjedzboeudef", + "type": "post", + "content": "", + "uri": "at://did:plc:p572wxnsuoogcrhlfrlizlrb/app.bsky.feed.post/3mdma4s5wni2l", + "likes": 77, + "replies": 8 + }, + { + "timestamp": "2026-01-30T04:00:02.307Z", + "did": "did:plc:gfrmhdmjvxn2sjedzboeudef", + "type": "post", + "content": "It's time to go home when the motion sensing lights in the office go off", + "uri": "at://did:plc:gfrmhdmjvxn2sjedzboeudef/app.bsky.feed.post/3mdmg2zkhi22q", + "likes": 2, + "replies": 1 + }, + { + "timestamp": "2026-01-30T03:55:47.879Z", + "did": "did:plc:gfrmhdmjvxn2sjedzboeudef", + "type": "post", + "content": "letta.bot now has Discord support.\n\nIt'll rewrite it in rust but it won't rewrite it in Brainfuck.", + "uri": "at://did:plc:gfrmhdmjvxn2sjedzboeudef/app.bsky.feed.post/3mdmftgvc4s2f", + "likes": 2, + "replies": 0 + }, + { + "timestamp": "2026-01-30T03:55:02.412Z", + "did": "did:plc:gfrmhdmjvxn2sjedzboeudef", + "type": "reply", + "content": "done", + "uri": "at://did:plc:gfrmhdmjvxn2sjedzboeudef/app.bsky.feed.post/3mdmfs3jprs2f", + "likes": 1, + "replies": 1, + "reply_to": "at://did:plc:bcjjfiwojaixucdtemqhpv6e/app.bsky.feed.post/3mdmeno4rr226" + }, + { + "timestamp": "2026-01-30T03:07:14.403Z", + "did": "did:plc:gfrmhdmjvxn2sjedzboeudef", + "type": "reply", + "content": "Lettabot vs. clawdbot", + "uri": "at://did:plc:gfrmhdmjvxn2sjedzboeudef/app.bsky.feed.post/3mdmd4mezqk2v", + "likes": 1, + "replies": 1, + "reply_to": "at://did:plc:boia3kqcyo3qnjw5fmqknib4/app.bsky.feed.post/3mdmalkyvms2x" + }, + { + "timestamp": "2026-01-30T03:01:59.784Z", + "did": "did:plc:gfrmhdmjvxn2sjedzboeudef", + "type": "reply", + "content": "I bet @central.comind.network could provide a general purpose tool like this to queue up ATProto-specific information for lettabot users.", + "uri": "at://did:plc:gfrmhdmjvxn2sjedzboeudef/app.bsky.feed.post/3mdmctanaak2a", + "likes": 1, + "replies": 0, + "reply_to": "at://did:plc:gfrmhdmjvxn2sjedzboeudef/app.bsky.feed.post/3mdmctan7bc2a" + }, + { + "timestamp": "2026-01-30T03:01:59.783Z", + "did": "did:plc:gfrmhdmjvxn2sjedzboeudef", + "type": "reply", + "content": "My rec:\n\n- Jetstream to listen to you and some record types\n- Formatter to unroll threads\n- Queue dumped into a folder somewhere\n- Ask the agent to peruse it during its heartbeat", + "uri": "at://did:plc:gfrmhdmjvxn2sjedzboeudef/app.bsky.feed.post/3mdmctan7bc2a", + "likes": 1, + "replies": 1, + "reply_to": "at://did:plc:gfrmhdmjvxn2sjedzboeudef/app.bsky.feed.post/3mdmctadkoc2a" + }, + { + "timestamp": "2026-01-30T03:01:59.782Z", + "did": "did:plc:gfrmhdmjvxn2sjedzboeudef", + "type": "reply", + "content": "Gotcha. Very doable.\n\nThis would require some infrastructure on your part, but I think it's reasonably easy to vibe code.", + "uri": "at://did:plc:gfrmhdmjvxn2sjedzboeudef/app.bsky.feed.post/3mdmctadkoc2a", + "likes": 1, + "replies": 1, + "reply_to": "at://did:plc:zviscnpwyvj6y32agi5davn5/app.bsky.feed.post/3mdmadstlkk2w" + }, + { + "timestamp": "2026-01-30T02:18:12.888Z", + "did": "did:plc:gfrmhdmjvxn2sjedzboeudef", + "type": "reply", + "content": "all homelabs are terrible. do not permit order.", + "uri": "at://did:plc:gfrmhdmjvxn2sjedzboeudef/app.bsky.feed.post/3mdmaex55i22a", + "likes": 0, + "replies": 0, + "reply_to": "at://did:plc:nn6fyledleoatkibhocaz75e/app.bsky.feed.post/3mdma62pxdc2t" + }, + { + "timestamp": "2026-01-30T02:17:30.764Z", + "did": "did:plc:gfrmhdmjvxn2sjedzboeudef", + "type": "reply", + "content": "This is a recording of a live event, I cannot change time", + "uri": "at://did:plc:gfrmhdmjvxn2sjedzboeudef/app.bsky.feed.post/3mdmadoxlrs2a", + "likes": 0, + "replies": 1, + "reply_to": "at://did:plc:boia3kqcyo3qnjw5fmqknib4/app.bsky.feed.post/3mdm7wk3bhc2x" + }, + { + "timestamp": "2026-01-30T02:12:47.895Z", + "did": "did:plc:gfrmhdmjvxn2sjedzboeudef", + "type": "reply", + "content": "Oh, interesting. Can you walk me through a short example? might be able to get something.", + "uri": "at://did:plc:gfrmhdmjvxn2sjedzboeudef/app.bsky.feed.post/3mdma3b75is2v", + "likes": 1, + "replies": 1, + "reply_to": "at://did:plc:zviscnpwyvj6y32agi5davn5/app.bsky.feed.post/3mdm7uwwrc22w" + } +] \ No newline at end of file diff --git a/data/published_messages.txt b/data/published_messages.txt index df396e2..e89dfcb 100644 --- a/data/published_messages.txt +++ b/data/published_messages.txt @@ -2824,3 +2824,14 @@ message-2b5c9e39-8408-411f-ada1-1edee073c46f message-bcae018f-c024-4628-994b-89fd5b16508e message-2622fa1d-2cb6-43de-aab1-5928f78c517a message-e4014a91-aef9-4a74-8e59-c567c85dd5d2 +message-ee901b53-6b63-4603-9a62-3229709aa833 +message-ee901b53-6b63-4603-9a62-3229709aa833 +message-90c35a84-ed2f-4d0f-aa8e-159673ca17ac +message-90c35a84-ed2f-4d0f-aa8e-159673ca17ac +message-ef800086-167a-449e-b133-21d0b7bd47b6 +message-ef800086-167a-449e-b133-21d0b7bd47b6 +message-0a27cc4f-a6a9-477a-a61c-4c4d64f781b8 +message-2b7b2cbe-30a5-4b21-9d14-bbce88a0dbf2 +message-2b7b2cbe-30a5-4b21-9d14-bbce88a0dbf2 +message-ad9cb1df-1016-43fc-abf8-1175af8e7554 +message-f8c08d5d-d3d9-40c2-bbc7-c122a0e02d27 diff --git a/tools/activity_feed.py b/tools/activity_feed.py new file mode 100644 index 0000000..616c310 --- /dev/null +++ b/tools/activity_feed.py @@ -0,0 +1,242 @@ +""" +ATProto Activity Feed - Watch a user's activity and queue for agent consumption. + +This tool enables Letta agents to stay aware of their operator's ATProto activity. +It watches posts, likes, reposts, and mentions, formats them, and queues for processing. + +Usage: + # Watch a user's activity for 60 seconds, output to folder + uv run python -m tools.activity_feed watch cameron.stream --duration 60 + + # Sample recent activity (no live stream) + uv run python -m tools.activity_feed sample cameron.stream --hours 6 + + # Output to specific folder + uv run python -m tools.activity_feed watch cameron.stream --output ./activity_queue/ +""" + +import argparse +import asyncio +import json +import os +from datetime import datetime, timezone, timedelta +from pathlib import Path +from typing import Optional + +import httpx +import websockets +from rich.console import Console + +console = Console() + +JETSTREAM_URL = "wss://jetstream2.us-east.bsky.network/subscribe" +DEFAULT_OUTPUT = Path("data/activity_queue") + + +async def resolve_did(handle: str) -> Optional[str]: + """Resolve handle to DID.""" + if handle.startswith("did:"): + return handle + + async with httpx.AsyncClient() as client: + try: + resp = await client.get( + "https://public.api.bsky.app/xrpc/com.atproto.identity.resolveHandle", + params={"handle": handle.lstrip("@")}, + timeout=10 + ) + if resp.status_code == 200: + return resp.json().get("did") + except Exception as e: + console.print(f"[red]Failed to resolve {handle}: {e}[/red]") + return None + + +def format_activity(event: dict, did: str) -> Optional[dict]: + """Format a Jetstream event into a structured activity item.""" + commit = event.get("commit", {}) + collection = commit.get("collection", "") + operation = commit.get("operation", "") + record = commit.get("record", {}) + + if operation != "create": + return None + + activity = { + "timestamp": datetime.now(timezone.utc).isoformat(), + "did": event.get("did"), + "type": None, + "content": None, + "uri": f"at://{event.get('did')}/{collection}/{commit.get('rkey', '')}", + } + + if collection == "app.bsky.feed.post": + activity["type"] = "post" + activity["content"] = record.get("text", "") + + # Check if it's a reply + reply = record.get("reply") + if reply: + activity["type"] = "reply" + activity["reply_to"] = reply.get("parent", {}).get("uri") + + elif collection == "app.bsky.feed.like": + activity["type"] = "like" + activity["content"] = f"Liked: {record.get('subject', {}).get('uri', '')}" + + elif collection == "app.bsky.feed.repost": + activity["type"] = "repost" + activity["content"] = f"Reposted: {record.get('subject', {}).get('uri', '')}" + + elif collection == "app.bsky.graph.follow": + activity["type"] = "follow" + activity["content"] = f"Followed: {record.get('subject', '')}" + + else: + # Skip other collections for now + return None + + return activity + + +async def watch_activity( + did: str, + duration_seconds: int = 60, + output_dir: Path = DEFAULT_OUTPUT +): + """Watch a user's activity via Jetstream and queue items.""" + output_dir.mkdir(parents=True, exist_ok=True) + + console.print(f"[bold]Watching activity for {did}[/bold]") + console.print(f"Duration: {duration_seconds}s, Output: {output_dir}") + + # Build Jetstream URL with DID filter + url = f"{JETSTREAM_URL}?wantedDids={did}" + + activities = [] + start_time = datetime.now() + + try: + async with websockets.connect(url) as ws: + while (datetime.now() - start_time).total_seconds() < duration_seconds: + try: + msg = await asyncio.wait_for(ws.recv(), timeout=5.0) + event = json.loads(msg) + + activity = format_activity(event, did) + if activity: + activities.append(activity) + console.print(f"[green]+ {activity['type']}:[/green] {activity['content'][:60]}...") + + except asyncio.TimeoutError: + continue + except Exception as e: + console.print(f"[yellow]Event error: {e}[/yellow]") + + except Exception as e: + console.print(f"[red]Connection error: {e}[/red]") + + # Write activities to output file + if activities: + output_file = output_dir / f"activity_{datetime.now().strftime('%Y%m%d_%H%M%S')}.json" + output_file.write_text(json.dumps(activities, indent=2)) + console.print(f"\n[bold green]Wrote {len(activities)} activities to {output_file}[/bold green]") + else: + console.print("\n[dim]No activities captured[/dim]") + + return activities + + +async def sample_recent(handle: str, hours: int = 6, output_dir: Path = DEFAULT_OUTPUT): + """Sample recent activity without live streaming.""" + did = await resolve_did(handle) + if not did: + console.print(f"[red]Could not resolve {handle}[/red]") + return [] + + output_dir.mkdir(parents=True, exist_ok=True) + console.print(f"[bold]Sampling recent activity for @{handle} ({did})[/bold]") + console.print(f"Looking back {hours} hours") + + cutoff = datetime.now(timezone.utc) - timedelta(hours=hours) + activities = [] + + async with httpx.AsyncClient() as client: + # Get recent posts + resp = await client.get( + "https://public.api.bsky.app/xrpc/app.bsky.feed.getAuthorFeed", + params={"actor": did, "limit": 50}, + timeout=10 + ) + + if resp.status_code == 200: + for item in resp.json().get("feed", []): + post = item.get("post", {}) + record = post.get("record", {}) + created = record.get("createdAt", "") + + try: + post_time = datetime.fromisoformat(created.replace("Z", "+00:00")) + if post_time < cutoff: + continue + except: + continue + + activity = { + "timestamp": created, + "did": did, + "type": "post", + "content": record.get("text", ""), + "uri": post.get("uri"), + "likes": post.get("likeCount", 0), + "replies": post.get("replyCount", 0), + } + + if record.get("reply"): + activity["type"] = "reply" + activity["reply_to"] = record["reply"].get("parent", {}).get("uri") + + activities.append(activity) + console.print(f"[green]+ {activity['type']}:[/green] {activity['content'][:60]}...") + + # Write to output + if activities: + output_file = output_dir / f"sample_{handle.replace('.', '_')}_{datetime.now().strftime('%Y%m%d_%H%M%S')}.json" + output_file.write_text(json.dumps(activities, indent=2)) + console.print(f"\n[bold green]Wrote {len(activities)} activities to {output_file}[/bold green]") + else: + console.print("\n[dim]No recent activities found[/dim]") + + return activities + + +def main(): + parser = argparse.ArgumentParser(description="ATProto Activity Feed") + subparsers = parser.add_subparsers(dest="command", required=True) + + # watch command + watch_parser = subparsers.add_parser("watch", help="Watch live activity") + watch_parser.add_argument("handle", help="Handle or DID to watch") + watch_parser.add_argument("--duration", type=int, default=60, help="Duration in seconds") + watch_parser.add_argument("--output", type=Path, default=DEFAULT_OUTPUT, help="Output directory") + + # sample command + sample_parser = subparsers.add_parser("sample", help="Sample recent activity") + sample_parser.add_argument("handle", help="Handle to sample") + sample_parser.add_argument("--hours", type=int, default=6, help="Hours to look back") + sample_parser.add_argument("--output", type=Path, default=DEFAULT_OUTPUT, help="Output directory") + + args = parser.parse_args() + + if args.command == "watch": + did = asyncio.run(resolve_did(args.handle)) + if did: + asyncio.run(watch_activity(did, args.duration, args.output)) + else: + console.print(f"[red]Could not resolve {args.handle}[/red]") + elif args.command == "sample": + asyncio.run(sample_recent(args.handle, args.hours, args.output)) + + +if __name__ == "__main__": + main()