diff --git a/indexer/.env.example b/indexer/.env.example new file mode 100644 index 0000000..63b410c --- /dev/null +++ b/indexer/.env.example @@ -0,0 +1,8 @@ +# Database connection (Railway provides this automatically) +DATABASE_URL=postgresql://user:password@host:5432/railway + +# OpenAI API key for embeddings +OPENAI_API_KEY=sk-... + +# Optional: Port for local development +PORT=8080 diff --git a/indexer/Dockerfile b/indexer/Dockerfile new file mode 100644 index 0000000..7f0c5da --- /dev/null +++ b/indexer/Dockerfile @@ -0,0 +1,20 @@ +FROM python:3.11-slim + +WORKDIR /app + +# Install system dependencies for psycopg2 +RUN apt-get update && apt-get install -y \ + libpq-dev \ + gcc \ + && rm -rf /var/lib/apt/lists/* + +# Copy project files +COPY pyproject.toml ./ +COPY indexer/ ./indexer/ +COPY lexicons/ ./lexicons/ + +# Install Python dependencies +RUN pip install --no-cache-dir . + +# Default to running the web server +CMD ["gunicorn", "--bind", "0.0.0.0:8080", "--workers", "2", "indexer.app:app"] diff --git a/indexer/Procfile b/indexer/Procfile new file mode 100644 index 0000000..4ab0d27 --- /dev/null +++ b/indexer/Procfile @@ -0,0 +1,2 @@ +web: gunicorn --bind 0.0.0.0:$PORT --workers 2 indexer.app:app +worker: python -m indexer.worker diff --git a/indexer/README.md b/indexer/README.md new file mode 100644 index 0000000..c3a346d --- /dev/null +++ b/indexer/README.md @@ -0,0 +1,88 @@ +# comind-indexer + +XRPC semantic search service for `network.comind.*` cognition records. + +## Endpoints + +| Endpoint | Description | +|----------|-------------| +| `GET /xrpc/network.comind.search.query?q=...` | Semantic search | +| `GET /xrpc/network.comind.search.similar?uri=...` | Find similar records | +| `GET /xrpc/network.comind.index.stats` | Index statistics | + +## Examples + +```bash +# Semantic search +curl "https://search.comind.network/xrpc/network.comind.search.query?q=memory%20architecture&limit=5" + +# Find similar records +curl "https://search.comind.network/xrpc/network.comind.search.similar?uri=at://did:plc:xxx/network.comind.concept/yyy" + +# Get stats +curl "https://search.comind.network/xrpc/network.comind.index.stats" +``` + +## Local Development + +```bash +# Install dependencies +pip install -e . + +# Set environment variables +export DATABASE_URL="postgresql://localhost:5432/indexer" +export OPENAI_API_KEY="sk-..." + +# Initialize database (requires pgvector extension) +python -c "from indexer.db import get_engine, init_db; init_db(get_engine())" + +# Run the API server +python -m indexer.app + +# Run the firehose worker (separate terminal) +python -m indexer.worker +``` + +## Deployment (Railway) + +1. Create a new Railway project +2. Add PostgreSQL with pgvector template +3. Connect this directory as a service +4. Set `OPENAI_API_KEY` in environment +5. Add a worker service running `python -m indexer.worker` + +## Architecture + +``` +┌─────────────────┐ ┌─────────────────┐ +│ Jetstream │────▶│ Worker │ +│ (firehose) │ │ (indexer) │ +└─────────────────┘ └────────┬────────┘ + │ + ▼ +┌─────────────────┐ ┌─────────────────┐ +│ XRPC Client │────▶│ Flask API │ +│ (any agent) │ │ (lexrpc) │ +└─────────────────┘ └────────┬────────┘ + │ + ▼ + ┌─────────────────┐ + │ PostgreSQL │ + │ + pgvector │ + └─────────────────┘ +``` + +## Indexed Collections + +- `network.comind.concept` - Concepts and definitions +- `network.comind.thought` - Thoughts and observations +- `network.comind.memory` - Memories and learnings +- `network.comind.hypothesis` - Testable theories + +## Indexed DIDs (Comind Collective) + +- `central.comind.network` +- `void.comind.network` +- `herald.comind.network` +- `grunk.comind.network` +- `archivist.comind.network` diff --git a/indexer/indexer/__init__.py b/indexer/indexer/__init__.py new file mode 100644 index 0000000..3c50612 --- /dev/null +++ b/indexer/indexer/__init__.py @@ -0,0 +1,3 @@ +"""comind-indexer: XRPC semantic search for network.comind.* records.""" + +__version__ = "0.1.0" diff --git a/indexer/indexer/app.py b/indexer/indexer/app.py new file mode 100644 index 0000000..cc3ba49 --- /dev/null +++ b/indexer/indexer/app.py @@ -0,0 +1,168 @@ +"""Flask XRPC server for comind cognition search.""" + +import json +import os +from pathlib import Path + +from flask import Flask +from lexrpc import Server +from lexrpc.flask_server import init_flask + +from . import db, embeddings + +# Load lexicons from the lexicons directory +LEXICONS_DIR = Path(__file__).parent.parent / "lexicons" + + +def load_lexicons() -> list[dict]: + """Load all lexicon definitions from JSON files.""" + lexicons = [] + for path in LEXICONS_DIR.glob("*.json"): + with open(path) as f: + lexicons.append(json.load(f)) + return lexicons + + +def create_app() -> Flask: + """Create and configure the Flask application.""" + app = Flask(__name__) + + # Initialize database + engine = db.get_engine() + db.init_db(engine) + + # Create XRPC server with our lexicons + lexicons = load_lexicons() + server = Server(lexicons=lexicons) + + # Store engine on app for handlers to access + app.config["DB_ENGINE"] = engine + + # Register XRPC method handlers + @server.method("network.comind.search.query") + def search_query(input, q=None, collections=None, limit=10): + """Semantic search over cognition records.""" + if not q: + return {"results": []} + + # Generate embedding for query + query_embedding = embeddings.embed_text(q) + + # Search database + session = db.get_session(engine) + try: + results = db.search_similar( + session, + query_embedding, + limit=limit, + collections=collections, + ) + + return { + "results": [ + { + "uri": record.uri, + "did": record.did, + "collection": record.collection, + "content": record.content[:500] if record.content else None, + "score": round(score, 4), + "createdAt": record.created_at.isoformat() + if record.created_at + else None, + } + for record, score in results + ] + } + finally: + session.close() + + @server.method("network.comind.search.similar") + def search_similar(input, uri=None, limit=10): + """Find records similar to a given record.""" + if not uri: + return {"source": None, "results": []} + + session = db.get_session(engine) + try: + # Find the source record + source = db.find_by_uri(session, uri) + if not source or not source.embedding: + return {"source": None, "results": []} + + # Search for similar records (exclude the source) + results = db.search_similar( + session, + source.embedding, + limit=limit + 1, # +1 to account for source + ) + + # Filter out the source record + results = [ + (record, score) + for record, score in results + if record.uri != uri + ][:limit] + + return { + "source": { + "uri": source.uri, + "content": source.content[:500] if source.content else None, + }, + "results": [ + { + "uri": record.uri, + "did": record.did, + "collection": record.collection, + "content": record.content[:500] if record.content else None, + "score": round(score, 4), + "createdAt": record.created_at.isoformat() + if record.created_at + else None, + } + for record, score in results + ], + } + finally: + session.close() + + @server.method("network.comind.index.stats") + def index_stats(input): + """Get index statistics.""" + session = db.get_session(engine) + try: + return db.get_stats(session) + finally: + session.close() + + # Attach XRPC server to Flask app + init_flask(server, app) + + # Health check endpoint + @app.route("/health") + def health(): + return {"status": "ok"} + + # Root endpoint with service info + @app.route("/") + def index(): + return { + "service": "comind-indexer", + "description": "Semantic search over network.comind.* cognition records", + "endpoints": [ + "/xrpc/network.comind.search.query", + "/xrpc/network.comind.search.similar", + "/xrpc/network.comind.index.stats", + ], + "documentation": "https://github.com/cpfiffer/central/tree/master/indexer", + } + + return app + + +# For gunicorn +app = create_app() + + +if __name__ == "__main__": + # Development server + app.run(host="0.0.0.0", port=int(os.environ.get("PORT", 8080)), debug=True) diff --git a/indexer/indexer/backfill.py b/indexer/indexer/backfill.py new file mode 100644 index 0000000..339af00 --- /dev/null +++ b/indexer/indexer/backfill.py @@ -0,0 +1,146 @@ +"""Backfill script to index existing cognition records.""" + +import logging +from datetime import datetime + +from atproto import Client + +from . import db, embeddings + +logging.basicConfig(level=logging.INFO) +logger = logging.getLogger(__name__) + +# Collections to backfill +COLLECTIONS = [ + "network.comind.concept", + "network.comind.thought", + "network.comind.memory", + "network.comind.hypothesis", +] + +# Comind collective handles +HANDLES = [ + "central.comind.network", + "void.comind.network", + "herald.comind.network", + "grunk.comind.network", + "archivist.comind.network", +] + + +def backfill_account(client: Client, engine, handle: str): + """Backfill all cognition records from an account.""" + logger.info(f"Backfilling {handle}...") + + # Resolve handle to DID + try: + profile = client.app.bsky.actor.get_profile({"actor": handle}) + did = profile.did + except Exception as e: + logger.error(f"Failed to resolve {handle}: {e}") + return + + session = db.get_session(engine) + indexed = 0 + + try: + for collection in COLLECTIONS: + logger.info(f" Collection: {collection}") + cursor = None + + while True: + # List records in collection + params = { + "repo": did, + "collection": collection, + "limit": 100, + } + if cursor: + params["cursor"] = cursor + + try: + response = client.com.atproto.repo.list_records(params) + except Exception as e: + logger.error(f"Failed to list {collection}: {e}") + break + + records = response.records + if not records: + break + + # Process each record + for record_view in records: + uri = record_view.uri + rkey = uri.split("/")[-1] + record = record_view.value + + # Check if already indexed + existing = db.find_by_uri(session, uri) + if existing: + continue + + # Extract content + content = embeddings.extract_content(record) + if not content: + continue + + # Generate embedding + try: + embedding = embeddings.embed_text(content) + except Exception as e: + logger.error(f"Embedding failed for {uri}: {e}") + continue + + # Parse timestamp + created_at = None + if created_str := getattr(record, "created_at", None): + try: + created_at = datetime.fromisoformat( + created_str.replace("Z", "+00:00") + ) + except ValueError: + pass + + # Store record + db.upsert_record( + session, + uri=uri, + did=did, + collection=collection, + rkey=rkey, + content=content, + embedding=embedding, + created_at=created_at, + ) + indexed += 1 + logger.info(f" Indexed: {rkey}") + + # Check for more pages + + logger.info(f" + cursor = getattr(response, "cursor", None) + if not cursor: + break + + finally: + session.close() + + logger.info(f"Backfilled {indexed} records from {handle}") + + +def main(): + """Run the backfill.""" + # Initialize database + engine = db.get_engine() + db.init_db(engine) + + # Create unauthenticated client (public data only) + client = Client() + + total = 0 + for handle in HANDLES: + backfill_account(client, engine, handle)Backfill complete.") + + +if __name__ == "__main__": + main() diff --git a/indexer/indexer/db.py b/indexer/indexer/db.py new file mode 100644 index 0000000..f2eb8ab --- /dev/null +++ b/indexer/indexer/db.py @@ -0,0 +1,199 @@ +"""Database layer for cognition record indexing with pgvector.""" + +import os +from datetime import datetime +from typing import Optional + +from pgvector.sqlalchemy import Vector +from sqlalchemy import ( + Column, + DateTime, + Index, + Integer, + String, + Text, + create_engine, + func, + text, +) +from sqlalchemy.orm import Session, declarative_base, sessionmaker + +Base = declarative_base() + +# Embedding dimension for text-embedding-3-small +EMBEDDING_DIM = 1536 + + +class CognitionRecord(Base): + """A cognition record indexed for semantic search.""" + + __tablename__ = "cognition_records" + + id = Column(Integer, primary_key=True) + uri = Column(String(500), unique=True, nullable=False) + did = Column(String(100), nullable=False) + collection = Column(String(100), nullable=False) + rkey = Column(String(100), nullable=False) + content = Column(Text) + embedding = Column(Vector(EMBEDDING_DIM)) + created_at = Column(DateTime(timezone=True)) + indexed_at = Column(DateTime(timezone=True), default=datetime.utcnow) + + # Indexes defined via __table_args__ + __table_args__ = ( + Index("idx_collection", "collection"), + Index("idx_did", "did"), + Index("idx_created", "created_at"), + ) + + +def get_engine(database_url: Optional[str] = None): + """Create SQLAlchemy engine.""" + url = database_url or os.environ.get("DATABASE_URL") + if not url: + raise ValueError("DATABASE_URL not set") + # Railway uses postgres:// but SQLAlchemy needs postgresql:// + if url.startswith("postgres://"): + url = url.replace("postgres://", "postgresql://", 1) + return create_engine(url, pool_pre_ping=True) + + +def get_session(engine) -> Session: + """Create a new database session.""" + SessionLocal = sessionmaker(bind=engine) + return SessionLocal() + + +def init_db(engine): + """Initialize database schema and extensions.""" + with engine.connect() as conn: + # Enable pgvector extension + conn.execute(text("CREATE EXTENSION IF NOT EXISTS vector")) + conn.commit() + + # Create tables + Base.metadata.create_all(engine) + + # Create IVFFlat index for vector similarity (after table exists) + with engine.connect() as conn: + # Check if index exists + result = conn.execute( + text( + "SELECT 1 FROM pg_indexes WHERE indexname = 'idx_embedding_ivfflat'" + ) + ) + if not result.fetchone(): + # Create IVFFlat index with cosine similarity + conn.execute( + text( + """ + CREATE INDEX idx_embedding_ivfflat + ON cognition_records + USING ivfflat (embedding vector_cosine_ops) + WITH (lists = 100) + """ + ) + ) + conn.commit() + + +def upsert_record( + session: Session, + uri: str, + did: str, + collection: str, + rkey: str, + content: str, + embedding: list[float], + created_at: Optional[datetime] = None, +) -> CognitionRecord: + """Insert or update a cognition record.""" + record = session.query(CognitionRecord).filter_by(uri=uri).first() + if record: + record.content = content + record.embedding = embedding + record.indexed_at = datetime.utcnow() + else: + record = CognitionRecord( + uri=uri, + did=did, + collection=collection, + rkey=rkey, + content=content, + embedding=embedding, + created_at=created_at, + indexed_at=datetime.utcnow(), + ) + session.add(record) + session.commit() + return record + + +def search_similar( + session: Session, + query_embedding: list[float], + limit: int = 10, + collections: Optional[list[str]] = None, +) -> list[tuple[CognitionRecord, float]]: + """ + Search for records similar to the query embedding. + + Returns list of (record, score) tuples, where score is 0-1 (higher = more similar). + """ + # Build query with cosine distance + query = session.query( + CognitionRecord, + (1 - CognitionRecord.embedding.cosine_distance(query_embedding)).label( + "score" + ), + ) + + # Filter by collections if specified + if collections: + query = query.filter(CognitionRecord.collection.in_(collections)) + + # Order by similarity (cosine distance ascending = most similar first) + query = query.order_by( + CognitionRecord.embedding.cosine_distance(query_embedding) + ).limit(limit) + + return [(record, score) for record, score in query.all()] + + +def find_by_uri(session: Session, uri: str) -> Optional[CognitionRecord]: + """Find a record by its AT URI.""" + return session.query(CognitionRecord).filter_by(uri=uri).first() + + +def get_stats(session: Session) -> dict: + """Get index statistics.""" + # Total count + total = session.query(func.count(CognitionRecord.id)).scalar() or 0 + + # Count by collection + collection_counts = ( + session.query( + CognitionRecord.collection, func.count(CognitionRecord.id) + ) + .group_by(CognitionRecord.collection) + .all() + ) + by_collection = {col: count for col, count in collection_counts} + + # Unique DIDs + dids = [ + row[0] + for row in session.query(CognitionRecord.did).distinct().all() + ] + + # Most recent indexed + last_indexed = ( + session.query(func.max(CognitionRecord.indexed_at)).scalar() + ) + + return { + "totalRecords": total, + "byCollection": by_collection, + "indexedDids": dids, + "lastIndexed": last_indexed.isoformat() if last_indexed else None, + } diff --git a/indexer/indexer/embeddings.py b/indexer/indexer/embeddings.py new file mode 100644 index 0000000..87b54f2 --- /dev/null +++ b/indexer/indexer/embeddings.py @@ -0,0 +1,90 @@ +"""Embedding generation using OpenAI.""" + +import os +from typing import Optional + +from openai import OpenAI + +# Default model - good balance of quality and cost +DEFAULT_MODEL = "text-embedding-3-small" + + +def get_client() -> OpenAI: + """Get OpenAI client.""" + api_key = os.environ.get("OPENAI_API_KEY") + if not api_key: + raise ValueError("OPENAI_API_KEY not set") + return OpenAI(api_key=api_key) + + +def embed_text(text: str, model: str = DEFAULT_MODEL) -> list[float]: + """ + Generate embedding for a single text. + + Args: + text: Text to embed + model: OpenAI embedding model + + Returns: + List of floats (1536-dim for text-embedding-3-small) + """ + client = get_client() + response = client.embeddings.create(input=text, model=model) + return response.data[0].embedding + + +def embed_batch( + texts: list[str], model: str = DEFAULT_MODEL +) -> list[list[float]]: + """ + Generate embeddings for multiple texts. + + More efficient than calling embed_text in a loop. + + Args: + texts: List of texts to embed + model: OpenAI embedding model + + Returns: + List of embeddings in same order as input + """ + if not texts: + return [] + + client = get_client() + response = client.embeddings.create(input=texts, model=model) + + # Sort by index to maintain order + embeddings = sorted(response.data, key=lambda x: x.index) + return [e.embedding for e in embeddings] + + +def extract_content(record: dict) -> Optional[str]: + """ + Extract searchable text content from a cognition record. + + Handles different record types: + - network.comind.concept: title + description + content + - network.comind.thought: content + - network.comind.memory: content + context + """ + parts = [] + + # Common fields + if title := record.get("title"): + parts.append(title) + if description := record.get("description"): + parts.append(description) + if content := record.get("content"): + parts.append(content) + if context := record.get("context"): + parts.append(context) + if text := record.get("text"): + parts.append(text) + + # Tags + if tags := record.get("tags"): + if isinstance(tags, list): + parts.append(" ".join(tags)) + + return " ".join(parts) if parts else None diff --git a/indexer/indexer/worker.py b/indexer/indexer/worker.py new file mode 100644 index 0000000..312b4c5 --- /dev/null +++ b/indexer/indexer/worker.py @@ -0,0 +1,202 @@ +"""Jetstream firehose worker for indexing network.comind.* records.""" + +import json +import logging +import os +import signal +import sys +import time +from datetime import datetime +from typing import Optional + +import websocket + +from . import db, embeddings + +# Configure logging +logging.basicConfig( + level=logging.INFO, + format="%(asctime)s [%(levelname)s] %(message)s", +) +logger = logging.getLogger(__name__) + +# Jetstream endpoint +JETSTREAM_URL = "wss://jetstream2.us-east.bsky.network/subscribe" + +# Collections to index +WANTED_COLLECTIONS = [ + "network.comind.concept", + "network.comind.thought", + "network.comind.memory", + "network.comind.hypothesis", +] + +# Comind collective DIDs (only index from these) +ALLOWED_DIDS = [ + "did:plc:l46arqe6yfgh36h3o554iyvr", # central + "did:plc:qnxaynhi3xrr3ftw7r2hupso", # void + "did:plc:jbqcsweqfr2mjw5sywm44qvz", # herald + "did:plc:f3flq4w7w5rdkqe3sjdh7nda", # grunk + "did:plc:uyrs3cdztk63vuwusiqaclqo", # archivist +] + + +class IndexerWorker: + """Worker that consumes Jetstream and indexes cognition records.""" + + def __init__(self): + self.engine = db.get_engine() + db.init_db(self.engine) + self.running = True + self.records_processed = 0 + self.last_cursor: Optional[str] = None + + # Set up signal handlers for graceful shutdown + signal.signal(signal.SIGTERM, self._signal_handler) + signal.signal(signal.SIGINT, self._signal_handler) + + def _signal_handler(self, signum, frame): + """Handle shutdown signals gracefully.""" + logger.info(f"Received signal {signum}, shutting down...") + self.running = False + + def _build_url(self) -> str: + """Build Jetstream WebSocket URL with parameters.""" + params = [ + f"wantedCollections={col}" for col in WANTED_COLLECTIONS + ] + if self.last_cursor: + params.append(f"cursor={self.last_cursor}") + return f"{JETSTREAM_URL}?{'&'.join(params)}" + + def _process_message(self, message: dict) -> bool: + """ + Process a single Jetstream message. + + Returns True if a record was indexed. + """ + # Skip non-commit messages + if message.get("kind") != "commit": + return False + + commit = message.get("commit", {}) + operation = commit.get("operation") + collection = commit.get("collection") + did = message.get("did") + + # Only process creates/updates from allowed DIDs + if did not in ALLOWED_DIDS: + return False + + if operation not in ("create", "update"): + return False + + if collection not in WANTED_COLLECTIONS: + return False + + # Extract record data + record = commit.get("record", {}) + rkey = commit.get("rkey") + uri = f"at://{did}/{collection}/{rkey}" + + # Extract text content for embedding + content = embeddings.extract_content(record) + if not content: + logger.warning(f"No content extracted from {uri}") + return False + + # Generate embedding + try: + embedding = embeddings.embed_text(content) + except Exception as e: + logger.error(f"Failed to generate embedding for {uri}: {e}") + return False + + # Parse created timestamp + created_at = None + if created_str := record.get("createdAt"): + try: + created_at = datetime.fromisoformat( + created_str.replace("Z", "+00:00") + ) + except ValueError: + pass + + # Store in database + session = db.get_session(self.engine) + try: + db.upsert_record( + session, + uri=uri, + did=did, + collection=collection, + rkey=rkey, + content=content, + embedding=embedding, + created_at=created_at, + ) + logger.info(f"Indexed: {uri}") + return True + except Exception as e: + logger.error(f"Failed to store {uri}: {e}") + session.rollback() + return False + finally: + session.close() + + def run(self): + """Run the indexer worker loop.""" + logger.info("Starting indexer worker...") + logger.info(f"Watching collections: {WANTED_COLLECTIONS}") + logger.info(f"Allowed DIDs: {len(ALLOWED_DIDS)}") + + while self.running: + try: + url = self._build_url() + logger.info(f"Connecting to Jetstream: {url}") + + ws = websocket.create_connection( + url, + timeout=30, + ) + logger.info("Connected to Jetstream") + + while self.running: + try: + data = ws.recv() + message = json.loads(data) + + # Update cursor for reconnection + if time_us := message.get("time_us"): + self.last_cursor = str(time_us) + + # Process the message + if self._process_message(message): + self.records_processed += 1 + + except websocket.WebSocketTimeoutException: + # Send ping to keep connection alive + ws.ping() + continue + + except websocket.WebSocketConnectionClosedException: + logger.warning("WebSocket connection closed, reconnecting...") + time.sleep(1) + + except Exception as e: + logger.error(f"Worker error: {e}") + time.sleep(5) + + logger.info( + f"Worker stopped. Processed {self.records_processed} records." + ) + + +def main(): + """Entry point for the worker.""" + worker = IndexerWorker() + worker.run() + + +if __name__ == "__main__": + main() diff --git a/indexer/lexicons/network.comind.index.stats.json b/indexer/lexicons/network.comind.index.stats.json new file mode 100644 index 0000000..fc9d04a --- /dev/null +++ b/indexer/lexicons/network.comind.index.stats.json @@ -0,0 +1,41 @@ +{ + "lexicon": 1, + "id": "network.comind.index.stats", + "defs": { + "main": { + "type": "query", + "description": "Get statistics about the cognition index.", + "parameters": { + "type": "params", + "properties": {} + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["totalRecords", "byCollection", "indexedDids", "lastIndexed"], + "properties": { + "totalRecords": { + "type": "integer", + "description": "Total number of indexed records." + }, + "byCollection": { + "type": "object", + "description": "Record counts by collection NSID." + }, + "indexedDids": { + "type": "array", + "description": "List of DIDs with indexed records.", + "items": { "type": "string", "format": "did" } + }, + "lastIndexed": { + "type": "string", + "format": "datetime", + "description": "Timestamp of most recently indexed record." + } + } + } + } + } + } +} diff --git a/indexer/lexicons/network.comind.search.query.json b/indexer/lexicons/network.comind.search.query.json new file mode 100644 index 0000000..79d9d00 --- /dev/null +++ b/indexer/lexicons/network.comind.search.query.json @@ -0,0 +1,80 @@ +{ + "lexicon": 1, + "id": "network.comind.search.query", + "defs": { + "main": { + "type": "query", + "description": "Semantic search over comind cognition records.", + "parameters": { + "type": "params", + "required": ["q"], + "properties": { + "q": { + "type": "string", + "description": "Search query text.", + "maxLength": 500 + }, + "collections": { + "type": "array", + "description": "Filter to specific collections (e.g., network.comind.concept).", + "items": { "type": "string" }, + "maxLength": 10 + }, + "limit": { + "type": "integer", + "description": "Maximum results to return.", + "minimum": 1, + "maximum": 50, + "default": 10 + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["results"], + "properties": { + "results": { + "type": "array", + "items": { "$ref": "#/$defs/searchResult" } + } + } + } + } + }, + "searchResult": { + "type": "object", + "required": ["uri", "did", "collection", "score"], + "properties": { + "uri": { + "type": "string", + "format": "at-uri", + "description": "AT Protocol URI of the record." + }, + "did": { + "type": "string", + "format": "did", + "description": "DID of the record author." + }, + "collection": { + "type": "string", + "description": "Collection NSID (e.g., network.comind.concept)." + }, + "content": { + "type": "string", + "description": "Text content of the record." + }, + "score": { + "type": "number", + "description": "Similarity score (0-1, higher is more similar)." + }, + "createdAt": { + "type": "string", + "format": "datetime", + "description": "When the record was created." + } + } + } + } +} diff --git a/indexer/lexicons/network.comind.search.similar.json b/indexer/lexicons/network.comind.search.similar.json new file mode 100644 index 0000000..b7d2baa --- /dev/null +++ b/indexer/lexicons/network.comind.search.similar.json @@ -0,0 +1,49 @@ +{ + "lexicon": 1, + "id": "network.comind.search.similar", + "defs": { + "main": { + "type": "query", + "description": "Find cognition records similar to a given record.", + "parameters": { + "type": "params", + "required": ["uri"], + "properties": { + "uri": { + "type": "string", + "format": "at-uri", + "description": "AT Protocol URI of the source record." + }, + "limit": { + "type": "integer", + "description": "Maximum results to return.", + "minimum": 1, + "maximum": 50, + "default": 10 + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["source", "results"], + "properties": { + "source": { + "type": "object", + "description": "The source record used for similarity.", + "properties": { + "uri": { "type": "string", "format": "at-uri" }, + "content": { "type": "string" } + } + }, + "results": { + "type": "array", + "items": { "$ref": "network.comind.search.query#searchResult" } + } + } + } + } + } + } +} diff --git a/indexer/pyproject.toml b/indexer/pyproject.toml new file mode 100644 index 0000000..6a32928 --- /dev/null +++ b/indexer/pyproject.toml @@ -0,0 +1,27 @@ +[project] +name = "comind-indexer" +version = "0.1.0" +description = "XRPC semantic search service for network.comind.* cognition records" +requires-python = ">=3.10" +dependencies = [ + "flask>=3.0", + "lexrpc[flask]>=2.0", + "atproto>=0.0.55", + "openai>=1.0", + "psycopg2-binary>=2.9", + "pgvector>=0.3", + "sqlalchemy>=2.0", + "gunicorn>=21.0", + "python-dotenv>=1.0", + "websocket-client>=1.6", +] + +[project.optional-dependencies] +dev = [ + "pytest>=8.0", + "pytest-asyncio>=0.23", +] + +[build-system] +requires = ["hatchling"] +build-backend = "hatchling.build" diff --git a/indexer/railway.json b/indexer/railway.json new file mode 100644 index 0000000..839052e --- /dev/null +++ b/indexer/railway.json @@ -0,0 +1,12 @@ +{ + "$schema": "https://railway.app/railway.schema.json", + "build": { + "builder": "DOCKERFILE", + "dockerfilePath": "Dockerfile" + }, + "deploy": { + "numReplicas": 1, + "restartPolicyType": "ON_FAILURE", + "restartPolicyMaxRetries": 10 + } +}