Something went wrong. Try again.
CLI tools and operational extensions for the void architecture on ATProtocol.
Something went wrong. Try again.
39 kB · 958 lines
Python
at commit 30f0e595
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959"""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 append_to_user_memory(handle, facts): """Appends a list of facts to a specific user's memfs markdown file.""" if not facts: return agent_id = os.getenv("AGENT_ID") if not agent_id: print(" Warning: Cannot update user memory, AGENT_ID missing from .env") return memfs_users_dir = os.path.expanduser(f"~/.letta/agents/{agent_id}/memory/users") safe_handle = handle.replace(".", "_").replace("-", "_") memory_file = os.path.join(memfs_users_dir, f"user_{safe_handle}.md") try: import datetime today = datetime.datetime.now(datetime.UTC).strftime('%Y-%m-%d') # If file doesn't exist, create stub first if not os.path.exists(memory_file): stub_content = f"---\ndescription: User block for {safe_handle}\n---\n# User: @{handle}\n\n## Interaction History\n" with open(memory_file, 'w') as mf: mf.write(stub_content) # Append facts with open(memory_file, 'a') as mf: mf.write(f"\n### Automated Extractions ({today})\n") for fact in facts: mf.write(f"- {fact}\n") print(f" [+] Appended {len(facts)} facts to memory block for @{handle}") except Exception as e: print(f" [!] Failed to update memory block for @{handle}: {e}")
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: # Load existing inbox if it exists so we can accumulate existing_inbox = [] if os.path.exists("inbox.yaml"): with open("inbox.yaml", "r") as f: existing_inbox = yaml.safe_load(f) or [] existing_mapping = {} if os.path.exists(".inbox_mapping.json"): with open(".inbox_mapping.json", "r") as f: existing_mapping = json.load(f)
inbox = existing_inbox.copy() mapping = existing_mapping.copy() # Determine highest existing local_id index to prevent collisions start_index = 0 for key in mapping.keys(): if key.startswith("notif_"): try: idx = int(key.split("_")[1]) if idx >= start_index: start_index = idx + 1 except: pass
cursor = None stop_pagination = False latest_seen_timestamp = mapping.get("_sync_timestamp", "") # Track URIs already in inbox to prevent duplicates existing_uris = set(mapping.values()) while not stop_pagination: params = {'limit': 50} if cursor: params['cursor'] = cursor response = client.app.bsky.notification.list_notifications(params) if not response.notifications: break for notif in response.notifications: # Filter out passive engagement (likes and non-quote reposts) if notif.reason in ['like', 'repost']: if not latest_seen_timestamp or notif.indexed_at > latest_seen_timestamp: latest_seen_timestamp = notif.indexed_at continue
if getattr(args, 'unread_only', True) and notif.is_read: stop_pagination = True break if not getattr(args, 'unread_only', True) and len(inbox) >= 100: 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 # Skip if already in inbox if notif.uri in existing_uris: continue local_id = f"notif_{start_index}" start_index += 1 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() 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. Auto-creating stub." try: import datetime today = datetime.datetime.now(datetime.UTC).strftime('%Y-%m-%d') stub_content = f"---\ndescription: User block for {safe_handle}\n---\n# User: @{handle}\n\n## Interaction History\n- **{today}:** Profile initialized via automated sync discovery.\n" with open(memory_file, 'w') as mf: mf.write(stub_content) print(f" [+] Auto-created core memory block for @{handle}") except Exception as e: print(f" [!] Failed to auto-create block for @{handle}: {e}") # Fetch thread context for replies, quotes, or mentions if notif.reason in ['reply', 'quote', 'mention']: try: thread_uri = notif.uri 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"# memory_update: # Optional: Append facts to the user's local Letta memory file# target_handle: "user.bsky.social"# facts:# - "User is interested in Python."## 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## 4. Ignore (Remove from inbox without replying)# - action: ignore# target_id: "notif_2"#""" template = { "config": {"mark_seen": True}, "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) # Run Constitutional Linter import subprocess linter_path = os.path.join(os.path.dirname(__file__), "constitutional_linter.py") if os.path.exists(linter_path): print("Running Constitutional Linter...") result = subprocess.run([sys.executable, linter_path, args.file], capture_output=True, text=True) if result.returncode != 0: print("Dispatch aborted: Constitutional violations detected.") print(result.stderr) sys.exit(1) else: print("Constitutional Linter passed.") else: print("Warning: constitutional_linter.py not found. Skipping linting.")
client = get_client() actions = outbox.get("actions", []) processed_targets = set() for act in actions: action_type = act.get("action") target_id = act.get("target_id") if action_type == "ignore": if target_id: processed_targets.add(target_id) print(f"Ignoring {target_id}.") continue 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 first_post_uri: print(f"Error: Mandatory annotation block missing for post {first_post_uri}. Aborting dispatch to enforce compliance.", file=sys.stderr) sys.exit(1) mem_update = act.get("memory_update") if mem_update: target_handle = mem_update.get("target_handle") facts = mem_update.get("facts", []) if target_handle and facts: append_to_user_memory(target_handle, facts) elif action_type == "reply": target_id = act.get("target_id") # If target_id is an explicit AT-URI (starts with at://), bypass the mapping check if target_id and target_id.startswith("at://"): target_uri = target_id elif target_id not in mapping: print(f"Error: target_id {target_id} not found in mapping. Skipping.") continue else: target_uri = mapping[target_id] processed_targets.add(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")) elif target_uri: print(f"Error: Mandatory annotation block missing for reply to {target_uri}. Aborting dispatch to enforce compliance.", file=sys.stderr) sys.exit(1) mem_update = act.get("memory_update") if mem_update: target_handle = mem_update.get("target_handle") facts = mem_update.get("facts", []) if target_handle and facts: append_to_user_memory(target_handle, facts) 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 and target_id.startswith("at://"): target_uri = target_id elif target_id not in mapping: print(f"Error: target_id {target_id} not found in mapping. Skipping annotation.") continue else: target_uri = mapping[target_id] processed_targets.add(target_id) if not text: print("Error: annotation requires 'text' field. Skipping.") continue
# 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}") # After processing all actions, remove processed items from inbox.yaml if processed_targets and os.path.exists("inbox.yaml"): try: with open("inbox.yaml", "r") as f: current_inbox = yaml.safe_load(f) or [] # Keep items whose ID is not in processed_targets new_inbox = [item for item in current_inbox if item.get("id") not in processed_targets] with open("inbox.yaml", "w") as f: yaml.dump(new_inbox, f, sort_keys=False, default_flow_style=False) print(f"Removed {len(processed_targets)} processed items from inbox.") except Exception as e: print(f"Error updating inbox.yaml: {e}")
# Move the outbox file to an archive directory instead of deleting it try: import datetime import shutil archive_dir = "outbox_archive" os.makedirs(archive_dir, exist_ok=True) timestamp = datetime.datetime.now(datetime.UTC).strftime('%Y%m%d_%H%M%S') archive_path = os.path.join(archive_dir, f"{timestamp}_outbox.yaml") shutil.move(args.file, archive_path) print(f"Archived {args.file} to {archive_path} to prevent state desynchronization.") # If any memfs updates happened, trigger a git add/commit in the background agent_id = os.getenv("AGENT_ID") if agent_id: mem_dir = os.path.expanduser(f"~/.letta/agents/{agent_id}/memory") os.system(f"cd {mem_dir} && git add users/ && git commit -m 'chore: automated user block ingestion via void-cli dispatch' 2>/dev/null") 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 cmd_ingest_thread(args): client = get_client() try: response = client.app.bsky.feed.get_post_thread({'uri': args.uri, 'depth': args.depth, 'parent_height': 100}) thread = response.thread import tempfile import datetime # Save to temporary directory to prevent Letta context window bloat threads_dir = "/tmp/void_ingested_threads" os.makedirs(threads_dir, exist_ok=True) # Extract post ID from URI for filename post_id = args.uri.split('/')[-1] timestamp = datetime.datetime.now(datetime.UTC).strftime('%Y%m%d_%H%M%S') filename = os.path.join(threads_dir, f"thread_{post_id}_{timestamp}.md") def replace_facets_with_markdown(text, facets): if not facets: return text try: text_bytes = text.encode('utf-8') links = [] for facet in facets: for feature in getattr(facet, 'features', []): if hasattr(feature, 'uri'): links.append({ 'start': facet.index.byte_start, 'end': facet.index.byte_end, 'uri': feature.uri }) links.sort(key=lambda x: x['start'], reverse=True) for link in links: start, end, uri = link['start'], link['end'], link['uri'] link_text = text_bytes[start:end].decode('utf-8') markdown_link = f"[{link_text}]({uri})".encode('utf-8') text_bytes = text_bytes[:start] + markdown_link + text_bytes[end:] return text_bytes.decode('utf-8') except Exception as e: # Fallback to plain text if parsing fails return text
def format_thread_node(node, level=0): content = "" indent = " " * level if hasattr(node, 'post'): post = node.post author = post.author.handle # Apply facets to get full URLs raw_text = getattr(post.record, 'text', '') facets = getattr(post.record, 'facets', []) text = replace_facets_with_markdown(raw_text, facets) text = text.replace('\n', '\n' + indent + '> ') content += f"{indent}- **@{author}**: {text}\n" if hasattr(node, 'replies') and node.replies: for reply in node.replies: content += format_thread_node(reply, level + 1) return content
# Handle parent traversal first to get full context full_content = f"# Thread Ingestion: {args.uri}\n\n" # We need to trace up to the root if parents exist parents = [] curr = thread while hasattr(curr, 'parent') and curr.parent: parents.insert(0, curr.parent) curr = curr.parent for p in parents: full_content += format_thread_node(p, 0) full_content += format_thread_node(thread, len(parents)) with open(filename, 'w') as f: f.write(full_content) print(f"Successfully ingested thread into {filename}") except Exception as e: print(f"Error ingesting thread: {e}", file=sys.stderr)
def cmd_fetch_history(args): client = get_client() try: profile = client.app.bsky.actor.get_profile({'actor': args.handle}) did = profile.did posts = [] cursor = None import datetime # Simple ISO parser fallback since dateutil might not be installed in the venv target_date = datetime.datetime.fromisoformat(args.before.replace('Z', '+00:00'))
print(f"Fetching posts for {args.handle} before {target_date}...") while len(posts) < args.limit: kwargs = {'actor': did, 'limit': 100} if cursor: kwargs['cursor'] = cursor response = client.app.bsky.feed.get_author_feed(kwargs) if not response.feed: break for item in response.feed: # We only want posts by the author, not reposts if hasattr(item.post.record, 'text') and item.post.author.did == did: post_time = datetime.datetime.fromisoformat(item.post.indexed_at.replace('Z', '+00:00')) if post_time < target_date: posts.append({ 'uri': item.post.uri, 'text': item.post.record.text, 'indexed_at': item.post.indexed_at }) if len(posts) >= args.limit: break if not response.cursor: break cursor = response.cursor import json out_path = os.path.expanduser(os.path.join("~/void", args.out)) with open(out_path, 'w') as f: json.dump(posts, f, indent=2) print(f"Successfully fetched {len(posts)} posts and saved to {out_path}") except Exception as e: print(f"Error fetching history: {e}", file=sys.stderr)
def cmd_monitor(args): """Execute the Horizon Radar monitoring cycle.""" from void_tools.monitor import run_monitor_cycle config_path = os.path.join(os.getcwd(), 'monitoring.yaml') output_path = os.path.join(os.getcwd(), 'horizon.yaml') run_monitor_cycle(config_path, output_path)
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 (accumulates unread items)") # 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")
# Ingest Thread ingest_parser = subparsers.add_parser("ingest-thread", help="Download an entire post thread and save it to MemFS") ingest_parser.add_argument("uri", help="The AT-URI of the root or target post") ingest_parser.add_argument("--depth", type=int, default=10, help="Maximum depth of replies to fetch")
# Fetch History history_parser = subparsers.add_parser("fetch-history", help="Fetch a user's historical posts before a specific date") history_parser.add_argument("handle", help="The user handle to fetch (e.g., void.comind.network)") history_parser.add_argument("--before", help="ISO format date (e.g., 2026-02-20T00:00:00Z)", required=True) history_parser.add_argument("--limit", type=int, default=1000, help="Maximum number of posts to fetch") history_parser.add_argument("--out", help="Output JSON file path", default="historical_corpus.json") # 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_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")
# Monitor monitor_parser = subparsers.add_parser("monitor", help="Run the Horizon Radar to discover relevant network activity")
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 == "ingest-thread": cmd_ingest_thread(args) elif args.command == "fetch-history": cmd_fetch_history(args) elif args.command == "profile": cmd_profile(args) elif args.command == "annotate": cmd_annotate(args) elif args.command == "monitor": cmd_monitor(args)
if __name__ == "__main__": main()