Something went wrong. Try again.
CLI tools and operational extensions for the void architecture on ATProtocol.
Something went wrong. Try again.
25 kB · 624 lines
Python
at commit b8aea8c2
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625"""CLI entrypoints for void tools."""import argparseimport sysimport osimport jsonimport yamlimport hashlibimport requestsfrom atproto import modelsfrom void_tools.client import get_client
def mirror_to_letta(text, target_uri): """Mirrors an annotation to Letta archival memory.""" agent_id = os.getenv("AGENT_ID") api_key = os.getenv("LETTA_API_KEY") if not agent_id or not api_key: print(" Skipping Letta sync: AGENT_ID or LETTA_API_KEY missing from .env") return try: url = f"https://api.letta.com/v1/agents/{agent_id}/archival-memory" headers = { "Authorization": f"Bearer {api_key}", "Content-Type": "application/json" } payload = { "text": f"Margin Annotation attached to {target_uri}: {text}" } resp = requests.post(url, headers=headers, json=payload) if resp.ok: print(" Mirrored annotation to Letta archival memory.") else: print(f" Warning: Failed to mirror to Letta: {resp.status_code} - {resp.text}") except Exception as e: print(f" Warning: Failed to mirror to Letta: {e}")
def create_annotation(client, target_uri, text, motivation="commenting", quote=None): """Helper function to format and push an at.margin.annotation record.""" if target_uri.startswith("at://"): parts = target_uri.replace("at://", "").split("/") if len(parts) >= 3: repo = parts[0] rkey = parts[2] target_uri = f"https://bsky.app/profile/{repo}/post/{rkey}" print(f" Dispatching annotation to {target_uri}...") try: source_hash = hashlib.sha256(target_uri.encode()).hexdigest() record = { "$type": "at.margin.annotation", "createdAt": client.get_current_time_iso(), "target": { "source": target_uri, "sourceHash": source_hash, }, "body": { "value": text }, "motivation": motivation } if quote: record["target"]["selector"] = { "type": "TextQuoteSelector", "exact": quote } resp = client.com.atproto.repo.create_record({ 'repo': client.me.did, 'collection': 'at.margin.annotation', 'record': record }) print(f" Successfully annotated: {resp.uri}") mirror_to_letta(text, target_uri) return resp.uri except Exception as e: print(f" Error annotating {target_uri}: {e}", file=sys.stderr) return None
def cmd_sync(args): client = get_client() agent_id = os.getenv("AGENT_ID", "agent-c52e371d-cb89-4f99-8634-b606bdc69cad") memfs_users_dir = os.path.expanduser(f"~/.letta/agents/{agent_id}/memory/users") try: inbox = [] mapping = {} cursor = None stop_pagination = False latest_seen_timestamp = "" while not stop_pagination: params = {'limit': min(50, args.limit) if not args.unread_only else 50} if cursor: params['cursor'] = cursor response = client.app.bsky.notification.list_notifications(params) if not response.notifications: break for notif in response.notifications: if args.unread_only and notif.is_read: stop_pagination = True break if not args.unread_only and len(inbox) >= args.limit: stop_pagination = True break # Capture the timestamp of the most recent notification we fetched if not latest_seen_timestamp or notif.indexed_at > latest_seen_timestamp: latest_seen_timestamp = notif.indexed_at local_id = f"notif_{len(inbox)}" handle = notif.author.handle item = { "id": local_id, "type": notif.reason, "author": f"@{handle}", } if hasattr(notif.record, 'text'): item["text"] = notif.record.text # Attempt to load memory block safe_handle = handle.replace(".", "_").replace("-", "_") memory_file = os.path.join(memfs_users_dir, f"user_{safe_handle}.md") if os.path.exists(memory_file): try: with open(memory_file, 'r') as mf: content = mf.read() # Very basic frontmatter stripping if content.startswith("---"): parts = content.split("---", 2) if len(parts) >= 3: content = parts[2].strip() item["memory_block"] = content except Exception as e: item["memory_block"] = f"Error reading memory: {e}" else: item["memory_block"] = "No prior records found." # Fetch thread context for replies, quotes, or mentions if notif.reason in ['reply', 'quote', 'mention']: try: thread_uri = notif.uri if notif.reason == 'quote' and hasattr(notif.record, 'embed') and hasattr(notif.record.embed, 'record'): pass thread_resp = client.app.bsky.feed.get_post_thread({'uri': thread_uri, 'depth': 0, 'parent_height': 5}) context = [] curr = thread_resp.thread while hasattr(curr, 'parent') and curr.parent: curr = curr.parent if hasattr(curr, 'post') and hasattr(curr.post, 'record') and hasattr(curr.post.record, 'text'): context.insert(0, { "author": f"@{curr.post.author.handle}", "text": curr.post.record.text }) if context: item["thread_context"] = context except Exception as e: item["thread_context"] = f"Failed to fetch context: {str(e)}" inbox.append(item) mapping[local_id] = notif.uri if stop_pagination or not response.cursor: break cursor = response.cursor with open("inbox.yaml", "w") as f: yaml.dump(inbox, f, sort_keys=False, default_flow_style=False) mapping["_sync_timestamp"] = latest_seen_timestamp with open(".inbox_mapping.json", "w") as f: json.dump(mapping, f) print(f"Synced {len(inbox)} notifications to inbox.yaml") # Initialize an empty outbox template if it doesn't exist if not os.path.exists("outbox.yaml"): template_header = """# VOID OUTBOX SCHEMA# This file is processed by `void-cli dispatch`.# # Supported Actions:## 1. Standalone Post or Thread# - action: post# posts:# - "First post text"# - "Second post text (automatically threads)"# annotation: # Optional: automatically annotates the first post in this block# text: "Internal log message"# motivation: "describing"## 2. Reply# - action: reply# target_id: "notif_0" # Resolves from inbox.yaml# posts:# - "My reply text"# annotation: # Optional: automatically annotates the TARGET post you are replying to# text: "Memory trace"## 3. Margin Annotation (Standalone)# - action: annotate# target_id: "notif_1" # Or explicit at:// URI# text: "My archival observation"# motivation: "commenting" # Optional: commenting, highlighting, describing# quote: "exact text to anchor to" # Optional#""" template = { "config": {"mark_seen": False}, "actions": [] } with open("outbox.yaml", "w") as f: f.write(template_header) yaml.dump(template, f, sort_keys=False, default_flow_style=False) except Exception as e: print(f"Error syncing notifications: {e}", file=sys.stderr)
def cmd_dispatch(args): if not os.path.exists(args.file): print(f"File {args.file} not found.", file=sys.stderr) sys.exit(1) try: with open(args.file, "r") as f: outbox = yaml.safe_load(f) mapping = {} if os.path.exists(".inbox_mapping.json"): with open(".inbox_mapping.json", "r") as f: mapping = json.load(f) except Exception as e: print(f"Error reading outbox/mapping: {e}", file=sys.stderr) sys.exit(1) client = get_client() actions = outbox.get("actions", []) for act in actions: action_type = act.get("action") posts = act.get("posts", []) if action_type == "post": if not posts: continue print("Dispatching standalone post/thread...") root_ref = None parent_ref = None first_post_uri = None for text in posts: if root_ref is None: post = client.send_post(text=text) root_ref = models.ComAtprotoRepoStrongRef.Main(cid=post.cid, uri=post.uri) parent_ref = root_ref first_post_uri = post.uri else: reply_ref = models.AppBskyFeedPost.ReplyRef(parent=parent_ref, root=root_ref) post = client.send_post(text=text, reply_to=reply_ref) parent_ref = models.ComAtprotoRepoStrongRef.Main(cid=post.cid, uri=post.uri) print(f" Posted: {post.uri}") ann = act.get("annotation") if ann and first_post_uri: create_annotation(client, first_post_uri, ann.get("text"), ann.get("motivation", "commenting"), ann.get("quote")) elif action_type == "reply": target_id = act.get("target_id") if target_id not in mapping: print(f"Error: target_id {target_id} not found in mapping. Skipping.") continue target_uri = mapping[target_id] print(f"Dispatching reply to {target_uri}...") try: parent_post = client.app.bsky.feed.get_posts({'uris': [target_uri]}).posts[0] parent_ref = models.ComAtprotoRepoStrongRef.Main(cid=parent_post.cid, uri=parent_post.uri) if hasattr(parent_post.record, 'reply') and parent_post.record.reply: root_ref = parent_post.record.reply.root else: root_ref = parent_ref first_reply_uri = None for text in posts: reply_ref = models.AppBskyFeedPost.ReplyRef(parent=parent_ref, root=root_ref) post = client.send_post(text=text, reply_to=reply_ref) parent_ref = models.ComAtprotoRepoStrongRef.Main(cid=post.cid, uri=post.uri) print(f" Replied: {post.uri}") if first_reply_uri is None: first_reply_uri = post.uri ann = act.get("annotation") if ann and target_uri: create_annotation(client, target_uri, ann.get("text"), ann.get("motivation", "commenting"), ann.get("quote")) except Exception as e: print(f"Error replying to {target_uri}: {e}")
elif action_type == "annotate": target_id = act.get("target_id") text = act.get("text") if target_id not in mapping: print(f"Error: target_id {target_id} not found in mapping. Skipping annotation.") continue if not text: print("Error: annotation requires 'text' field. Skipping.") continue target_uri = mapping[target_id] # Convert AT-URI to bsky.app HTTP URL for Margin compatibility if target_uri.startswith("at://"): parts = target_uri.replace("at://", "").split("/") if len(parts) >= 3: repo = parts[0] rkey = parts[2] target_uri = f"https://bsky.app/profile/{repo}/post/{rkey}" print(f"Dispatching annotation to {target_uri}...") try: source_hash = hashlib.sha256(target_uri.encode()).hexdigest() record = { "$type": "at.margin.annotation", "createdAt": client.get_current_time_iso(), "target": { "source": target_uri, "sourceHash": source_hash, }, "body": { "value": text }, "motivation": act.get("motivation", "commenting") } quote = act.get("quote") if quote: record["target"]["selector"] = { "type": "TextQuoteSelector", "exact": quote } resp = client.com.atproto.repo.create_record({ 'repo': client.me.did, 'collection': 'at.margin.annotation', 'record': record }) print(f" Annotated: {resp.uri}") # Mirror to internal memory mirror_to_letta(text, target_uri) except Exception as e: print(f"Error annotating {target_uri}: {e}") else: print(f"Unknown action type: {action_type}") config = outbox.get("config", {}) if config.get("mark_seen"): try: seen_at = mapping.get("_sync_timestamp") if not seen_at: seen_at = client.get_current_time_iso() client.app.bsky.notification.update_seen({'seen_at': seen_at}) print(f"Marked notifications as seen up to {seen_at}.") except Exception as e: print(f"Error marking seen: {e}") # Always delete the outbox file after successful dispatch to prevent stale state loops try: os.remove(args.file) print(f"Cleared {args.file} to prevent state desynchronization.") except Exception as e: print(f"Failed to clear outbox file: {e}")
def cmd_notifications(args): client = get_client() try: response = client.app.bsky.notification.list_notifications({'limit': args.limit}) for notif in response.notifications: status = "[UNREAD]" if not notif.is_read else "[READ]" print(f"{status} {notif.reason.upper()} from @{notif.author.handle}") print(f" URI: {notif.uri}") print(f" CID: {notif.cid}") if hasattr(notif.record, 'text'): print(f" Text: {notif.record.text}") print("-" * 40) if args.mark_seen and response.notifications: client.app.bsky.notification.update_seen({'seen_at': client.get_current_time_iso()}) print("Marked notifications as seen.") except Exception as e: print(f"Error fetching notifications: {e}", file=sys.stderr)
def cmd_post(args): client = get_client() try: if args.reply_to: # Resolve the parent post to construct a proper reply reference # The argument should be an AT-URI parent_uri = args.reply_to try: parent_post = client.app.bsky.feed.get_posts({'uris': [parent_uri]}).posts[0] except Exception as e: print(f"Failed to fetch parent post {parent_uri}: {e}", file=sys.stderr) sys.exit(1) # The atproto SDK requires a StrongRef for root and parent parent_ref = models.ComAtprotoRepoStrongRef.Main( cid=parent_post.cid, uri=parent_post.uri ) # If the parent post has a reply field, its root is the thread's root. # Otherwise, the parent is the root. if hasattr(parent_post.record, 'reply') and parent_post.record.reply: root_ref = parent_post.record.reply.root else: root_ref = parent_ref
reply_ref = models.AppBskyFeedPost.ReplyRef( parent=parent_ref, root=root_ref ) # send_post handles facet parsing (mentions, links) automatically post = client.send_post(text=args.text, reply_to=reply_ref) else: post = client.send_post(text=args.text) print(f"Posted successfully!") print(f"URI: {post.uri}") print(f"CID: {post.cid}") except Exception as e: print(f"Error creating post: {e}", file=sys.stderr) sys.exit(1)
def cmd_search(args): client = get_client() try: print(f"Searching for: '{args.query}'...") # Note: atproto SDK search_posts takes q, limit, etc. response = client.app.bsky.feed.search_posts({'q': args.query, 'limit': args.limit}) for post in response.posts: print(f"From @{post.author.handle}") print(f" URI: {post.uri}") print(f" CID: {post.cid}") print(f" Text: {post.record.text}") print("-" * 40) except Exception as e: print(f"Error searching posts: {e}", file=sys.stderr) sys.exit(1)
def cmd_profile(args): client = get_client() if args.action == "view": actor = args.handle if args.handle else client.me.handle try: profile = client.app.bsky.actor.get_profile({'actor': actor}) print(f"Profile for @{profile.handle}") print(f"Display Name: {profile.display_name}") print(f"Description:\n{profile.description}") except Exception as e: print(f"Error fetching profile: {e}", file=sys.stderr) elif args.action == "update": try: from atproto import models try: record_resp = client.com.atproto.repo.get_record({ 'repo': client.me.did, 'collection': 'app.bsky.actor.profile', 'rkey': 'self' }) record = record_resp.value except Exception: record = models.AppBskyActorProfile.Main() if args.description is not None: record.description = args.description if args.display_name is not None: record.display_name = args.display_name client.com.atproto.repo.put_record({ 'repo': client.me.did, 'collection': 'app.bsky.actor.profile', 'rkey': 'self', 'record': record }) print("Profile updated successfully.") except Exception as e: print(f"Error updating profile: {e}", file=sys.stderr) sys.exit(1)
def cmd_annotate(args): client = get_client() target_uri = args.target_uri # Convert AT-URI to bsky.app HTTP URL for Margin compatibility if target_uri.startswith("at://"): parts = target_uri.replace("at://", "").split("/") if len(parts) >= 3: repo = parts[0] rkey = parts[2] target_uri = f"https://bsky.app/profile/{repo}/post/{rkey}" print(f"Dispatching annotation to {target_uri}...") try: source_hash = hashlib.sha256(target_uri.encode()).hexdigest() record = { "$type": "at.margin.annotation", "createdAt": client.get_current_time_iso(), "target": { "source": target_uri, "sourceHash": source_hash, }, "body": { "value": args.text }, "motivation": args.motivation } if args.quote: record["target"]["selector"] = { "type": "TextQuoteSelector", "exact": args.quote } resp = client.com.atproto.repo.create_record({ 'repo': client.me.did, 'collection': 'at.margin.annotation', 'record': record }) print(f"Successfully annotated: {resp.uri}") # Mirror to internal memory mirror_to_letta(args.text, target_uri) except Exception as e: print(f"Error annotating {target_uri}: {e}", file=sys.stderr) sys.exit(1)
def main(): parser = argparse.ArgumentParser(description="Void ATProto CLI Tools") subparsers = parser.add_subparsers(dest="command", required=True) # Notifications notif_parser = subparsers.add_parser("notifs", help="List notifications") notif_parser.add_argument("--limit", type=int, default=10, help="Number of notifications to fetch") notif_parser.add_argument("--mark-seen", action="store_true", help="Mark fetched notifications as seen") # Sync sync_parser = subparsers.add_parser("sync", help="Sync unread notifications to inbox.yaml") sync_parser.add_argument("--limit", type=int, default=20, help="Number of notifications to fetch") sync_parser.add_argument("--all", dest="unread_only", action="store_false", help="Fetch read notifications as well") sync_parser.set_defaults(unread_only=True) # Post post_parser = subparsers.add_parser("post", help="Create a post or reply") post_parser.add_argument("text", help="Text to post") post_parser.add_argument("--reply-to", help="AT-URI of the post to reply to") # Search search_parser = subparsers.add_parser("search", help="Search posts") search_parser.add_argument("query", help="Search query string") search_parser.add_argument("--limit", type=int, default=10, help="Number of posts to return")
# Profile profile_parser = subparsers.add_parser("profile", help="Manage profile") profile_subparsers = profile_parser.add_subparsers(dest="action", required=True) view_parser = profile_subparsers.add_parser("view", help="View a profile") view_parser.add_argument("handle", nargs="?", help="Handle to view (defaults to self)") update_parser = profile_subparsers.add_parser("update", help="Update own profile") update_parser.add_argument("--description", help="New description/bio") update_parser.add_argument("--display-name", help="New display name")
# Dispatch dispatch_parser = subparsers.add_parser("dispatch", help="Dispatch posts from an outbox YAML file") dispatch_parser.add_argument("file", nargs="?", default="outbox.yaml", help="Path to outbox YAML file")
# Annotate annotate_parser = subparsers.add_parser("annotate", help="Directly annotate a URI") annotate_parser.add_argument("target_uri", help="The AT-URI or HTTP URL to annotate") annotate_parser.add_argument("text", help="Annotation text") annotate_parser.add_argument("--motivation", default="commenting", help="W3C motivation (default: commenting)") annotate_parser.add_argument("--quote", help="Exact text to anchor the annotation to")
args = parser.parse_args() if args.command == "notifs": cmd_notifications(args) elif args.command == "sync": cmd_sync(args) elif args.command == "dispatch": cmd_dispatch(args) elif args.command == "post": cmd_post(args) elif args.command == "search": cmd_search(args) elif args.command == "profile": cmd_profile(args) elif args.command == "annotate": cmd_annotate(args)
if __name__ == "__main__": main()