agentleFS
Sign inSign up

airweave / rules

airweave-ai/airweave/.cursor/rules/source-connector-implementation.mdc

A source connector in Airweave is a Python module that extracts data from an external service and transforms it into searchable entities. This guide covers everything you need to build a production-ready connector. There are two types of source connectors: Every source connector requires three main components: Start with entities because they define your data model. There are two base entity types: Important: The embeddable=True flag is what makes your entities semantically searchable. Without embeddable=True, fields are only keyword-searchable, not…

Cursor rule6.6k starsChanged 6 months agoArchived repository
# Building a Source Connector in Airweave

## Overview

A **source connector** in Airweave is a Python module that extracts data from an external service and transforms it into searchable entities. This guide covers everything you need to build a production-ready connector.

There are two types of source connectors:

1. **Standard (Sync-Based)**: Extracts and syncs all data from the source to Airweave's vector database
2. **Federated Search**: Searches the source's API at query time without syncing data

## Core Components

Every source connector requires three main components:

1. **Source implementation** (`backend/airweave/platform/sources/{short_name}.py`)
2. **Entity schemas** (`backend/airweave/platform/entities/{short_name}.py`)
3. **OAuth configuration** (`backend/airweave/platform/auth/yaml/dev.integrations.yaml`)

---

## Part 1: Entity Schemas

Start with entities because they define your data model.

### File Location
```
backend/airweave/platform/entities/{short_name}.py
```

### Entity Types

There are two base entity types:

1. **ChunkEntity** - Text-based entities (tasks, messages, documents, etc.)
2. **FileEntity** - File attachments (PDFs, images, etc.)

### Basic Structure

```python
"""Entity schemas for {Connector Name}."""

from datetime import datetime
from typing import Any, Dict, List, Optional

from pydantic import Field

from airweave.platform.entities._airweave_field import AirweaveField
from airweave.platform.entities._base import ChunkEntity, FileEntity


class MyConnectorEntity(ChunkEntity):
    """Schema for primary entity type."""

    # Required fields
    name: str = AirweaveField(
        ...,
        description="Display name of the entity",
        embeddable=True  # This field will be embedded for search
    )

    # Timestamps (critical for incremental sync)
    created_at: Optional[datetime] = AirweaveField(
        None,
        description="When this entity was created",
        embeddable=True,
        is_created_at=True  # Marks this as the creation timestamp
    )

    modified_at: Optional[datetime] = AirweaveField(
        None,
        description="When this entity was last modified",
        embeddable=True,
        is_updated_at=True  # Marks this as the update timestamp
    )

    # Content fields
    content: Optional[str] = AirweaveField(
        None,
        description="The main text content",
        embeddable=True  # Make searchable
    )

    # Metadata fields (not embeddable)
    external_id: str = Field(
        ...,
        description="Unique ID from the external system"
    )

    permalink_url: Optional[str] = Field(
        None,
        description="Direct link to view in external system"
    )
```

### Key Principles

#### 1. Use AirweaveField for Searchable Content

**Important: The `embeddable=True` flag is what makes your entities semantically searchable.**

Without `embeddable=True`, fields are only keyword-searchable, not semantically searchable. This limits the user's ability to find relevant entities.

**Best Practice: Mark most user-visible, content-rich fields as `embeddable=True`**

This includes:
- **Text content**: descriptions, notes, comments, body text
- **Names and titles**: entity names, display names, titles
- **People**: assignees, authors, owners, members (as dicts with name/email)
- **Status and metadata**: status fields, tags, labels, priorities
- **Structured data**: any dict/list that contains searchable information
- **Timestamps**: created_at, modified_at, due_dates (helps with recency)
- **URLs**: permalink_url, web_links (helps users find original content)

**Only exclude from embeddable:**
- Internal IDs (entity_id, external_id, database IDs)
- Binary/technical metadata (sizes, checksums, mime_types)
- System-only fields not relevant to user searches

**Example - Information-Rich Entity:**

```python
class MyConnectorTaskEntity(ChunkEntity):
    """Task entity - NOTE: Most fields are embeddable for rich search."""

    # Core content - ALWAYS embeddable
    name: str = AirweaveField(..., description="Task name", embeddable=True)
    description: Optional[str] = AirweaveField(
        None,
        description="Task description",
        embeddable=True  # ✅ Critical for semantic search
    )
    notes: Optional[str] = AirweaveField(
        None,
        description="Additional notes",
        embeddable=True  # ✅ Searchable content
    )

    # People - embeddable for "find tasks assigned to John" queries
    assignee: Optional[Dict] = AirweaveField(
        None,
        description="User assigned to this task",
        embeddable=True  # ✅ Enables "who" searches
    )

    owner: Optional[Dict] = AirweaveField(
        None,
        description="Task owner",
        embeddable=True  # ✅ Enables owner searches
    )

    # Status and metadata - embeddable for filtering/search
    status: Optional[str] = AirweaveField(
        None,
        description="Task status (open, in_progress, done)",
        embeddable=True  # ✅ Enables status-based search
    )

    priority: Optional[str] = AirweaveField(
        None,
        description="Priority level",
        embeddable=True  # ✅ Find high-priority tasks
    )

    tags: List[str] = AirweaveField(
        default_factory=list,
        description="Task tags",
        embeddable=True  # ✅ Find by tag
    )

    # Timestamps - embeddable for recency boosting
    created_at: Optional[datetime] = AirweaveField(
        None,
        description="Creation time",
        embeddable=True,  # ✅ For recency
        is_created_at=True
    )

    due_date: Optional[str] = AirweaveField(
        None,
        description="Due date",
        embeddable=True  # ✅ Find overdue tasks
    )

    # URLs - embeddable so users can find links
    permalink_url: Optional[str] = Field(
        None,
        description="Link to task in external system"
        # Note: Use Field() not AirweaveField() for URLs if you don't want them embedded
        # But consider making them embeddable for "find tasks linking to X"
    )

    # IDs - NOT embeddable (internal use only)
    external_id: str = Field(..., description="ID in external system")
    project_id: str = Field(..., description="Parent project ID")
```

**Common Mistake - Sparse Entity:**

```python
# Avoid: Only name is embeddable, rest is not searchable
class SparseTaskEntity(ChunkEntity):
    name: str = AirweaveField(..., embeddable=True)
    description: Optional[str] = Field(None)  # Should be embeddable
    assignee: Optional[Dict] = Field(None)     # Should be embeddable
    status: Optional[str] = Field(None)        # Should be embeddable
    # Result: Users can only search by task name, nothing else
```

#### 2. Always Include Timestamps

Every entity should have `created_at` and/or `modified_at` with proper flags:

```python
created_at: Optional[datetime] = AirweaveField(
    None,
    description="Creation time",
    embeddable=True,
    is_created_at=True  # System uses this for incremental sync
)

modified_at: Optional[datetime] = AirweaveField(
    None,
    description="Last modification time",
    embeddable=True,
    is_updated_at=True  # System uses this for incremental sync
)
```

#### 3. Model Entity Hierarchies

If your connector has parent-child relationships, create separate entity classes:

```python
class WorkspaceEntity(ChunkEntity):
    """Top-level container."""
    name: str = AirweaveField(..., embeddable=True)
    # ...

class ProjectEntity(ChunkEntity):
    """Belongs to workspace."""
    name: str = AirweaveField(..., embeddable=True)
    workspace_id: str = Field(...)
    workspace_name: str = AirweaveField(..., embeddable=True)
    # ...

class TaskEntity(ChunkEntity):
    """Belongs to project."""
    name: str = AirweaveField(..., embeddable=True)
    project_id: str = Field(...)
    section_id: Optional[str] = Field(None)
    # ...
```

#### 4. CRITICAL: Always Include Breadcrumbs

**⚠️ IMPORTANT: Breadcrumbs are frequently forgotten but essential for entity relationships and search context.**

Every entity MUST have breadcrumbs set when yielded. Breadcrumbs track the entity's location in the hierarchy (e.g., "Organization → Person → Deal").

```python
from airweave.platform.entities._base import Breadcrumb

# When generating entities, ALWAYS set breadcrumbs based on relationships:

# Top-level entities (organizations, workspaces) - empty breadcrumbs
yield OrganizationEntity(
    entity_id=org_id,
    breadcrumbs=[],  # Top-level, no parent
    name=name,
    # ...
)

# Child entities - include parent breadcrumbs
breadcrumbs = []
if parent_org_id and parent_org_name:
    breadcrumbs.append(
        Breadcrumb(
            entity_id=str(parent_org_id),
            name=parent_org_name,
            entity_type="OrganizationEntity",
        )
    )

yield PersonEntity(
    entity_id=person_id,
    breadcrumbs=breadcrumbs,  # Links to parent organization
    name=name,
    # ...
)

# Deeply nested entities - include full hierarchy
deal_breadcrumbs = []
if org_id and org_name:
    deal_breadcrumbs.append(Breadcrumb(entity_id=str(org_id), name=org_name, entity_type="OrganizationEntity"))
if person_id and person_name:
    deal_breadcrumbs.append(Breadcrumb(entity_id=str(person_id), name=person_name, entity_type="PersonEntity"))

yield DealEntity(
    entity_id=deal_id,
    breadcrumbs=deal_breadcrumbs,  # Shows: Organization → Person → Deal
    # ...
)
```

**Common Mistake:**
```python
# ❌ WRONG - Empty breadcrumbs when entity has parent relationships
yield PersonEntity(
    entity_id=person_id,
    breadcrumbs=[],  # Lost relationship to organization!
    organization_id=org_id,
    organization_name=org_name,
    # ...
)
```

**Rule:** If an entity has a parent reference (org_id, project_id, etc.), it should have a corresponding breadcrumb.

#### 5. Web URLs - User-Facing Links

**⚠️ IMPORTANT: Web URLs must be correct user-facing links that agents can share with users.**

Every entity should have a `web_url` field that links to the record in the source system's UI. These are URLs that:
- Users can click to view the record in the source application
- Agents can provide to users when referencing entities
- Should open directly to the specific record, not a list view

```python
# In your entity schema
class MyEntity(ChunkEntity):
    web_url_value: Optional[str] = AirweaveField(
        None,
        description="URL to view this record in the source system.",
        embeddable=False,
        unhashable=True,  # URLs change and shouldn't affect content hash
    )

    @computed_field(return_type=str)
    def web_url(self) -> str:
        """User-facing link to the record."""
        return self.web_url_value or ""

# In your source - build proper URLs for each record type
def _build_record_url(self, record_type: str, record_id: str) -> Optional[str]:
    """Build a user-facing URL for the record."""
    # Different record types often have different URL patterns!
    # Verify the actual URL structure in the source system's UI
    url_patterns = {
        "task": f"https://app.example.com/tasks/{record_id}",
        "project": f"https://app.example.com/projects/{record_id}",
        "comment": f"https://app.example.com/tasks/{parent_id}#comment-{record_id}",
    }
    return url_patterns.get(record_type)
```

**Common Mistakes:**
- Using API URLs instead of UI URLs
- Using incorrect URL patterns (verify in the actual app!)
- Forgetting that some record types may not have direct URLs (e.g., settings pages)
- Not handling cases where the URL requires additional context (parent IDs, etc.)

#### 4. File Entities

For attachments, inherit from `FileEntity`:

```python
class MyConnectorFileEntity(FileEntity):
    """Schema for file attachments."""

    # FileEntity provides: file_id, name, mime_type, size, download_url
    # Add connector-specific fields:

    parent_task_id: str = Field(
        ...,
        description="ID of the task this file is attached to"
    )

    created_at: Optional[datetime] = AirweaveField(
        None,
        description="Upload time",
        embeddable=True,
        is_created_at=True
    )
```

---

## Part 2: Source Implementation

### File Location
```
backend/airweave/platform/sources/{short_name}.py
```

### Basic Structure

```python
"""{Connector Name} source implementation."""

from typing import Any, AsyncGenerator, Dict, List, Optional

import httpx
from tenacity import retry, stop_after_attempt, wait_exponential

from airweave.domains.sources.exceptions import SourceAuthError, SourceRateLimitError, SourceServerError
from airweave.platform.decorators import source
from airweave.platform.entities._base import Breadcrumb, ChunkEntity
from airweave.platform.entities.{short_name} import (
    MyConnectorEntity,
    MyConnectorFileEntity,
)
from airweave.platform.sources._base import BaseSource
from airweave.schemas.source_connection import AuthenticationMethod, OAuthType


@source(
    name="{Connector Display Name}",
    short_name="{short_name}",
    auth_methods=[
        AuthenticationMethod.OAUTH_BROWSER,
        AuthenticationMethod.OAUTH_TOKEN,
        AuthenticationMethod.AUTH_PROVIDER,
    ],
    oauth_type=OAuthType.WITH_REFRESH,  # or WITH_ROTATING_REFRESH, ACCESS_ONLY
    auth_config_class=None,
    config_class="{ConnectorName}Config",  # Must match schema name
    labels=["Category"],  # e.g., "Project Management", "CRM", "Storage"
    supports_continuous=False,  # Set to True if you support webhook-based sync
)
class MyConnectorSource(BaseSource):
    """{Connector Name} source connector.

    Syncs {list of entity types} from {Connector Name}.
    """

    @classmethod
    async def create(
        cls,
        *,
        auth: "SourceAuthProvider",
        logger: "ContextualLogger",
        http_client: "AirweaveHttpClient",
        config: "BaseModel",
    ) -> "MyConnectorSource":
        """Create and configure the source.

        Args:
            auth: Auth provider — call ``await auth.get_token()`` for a bearer token.
            logger: Contextual logger with sync/search metadata.
            http_client: Pre-built AirweaveHttpClient with rate limiting.
            config: Typed config instance (the source's config_class).

        Returns:
            Configured source instance
        """
        instance = cls(auth=auth, logger=logger, http_client=http_client)
        instance.workspace_id = getattr(config, "workspace_id", None)
        instance.exclude_pattern = getattr(config, "exclude_pattern", "")
        return instance

    async def generate_entities(
        self,
        *,
        cursor=None,
        files=None,
        node_selections=None,
    ) -> AsyncGenerator[ChunkEntity, None]:
        """Generate all entities from the source.

        Args:
            cursor: SyncCursor for incremental sync tracking.
            files: FileService for downloading file attachments.
            node_selections: Node selections for targeted sync.
        """
        client = self.http_client
        async for top_level in self._generate_top_level(client):
            yield top_level

            async for child in self._generate_children(client, top_level):
                yield child

    async def validate(self) -> None:
        """Verify credentials by pinging the API."""
        token = await self.get_access_token()
        response = await self.http_client.get(
            "https://api.example.com/v1/me",
            headers={"Authorization": f"Bearer {token}", "Accept": "application/json"},
        )
        response.raise_for_status()
```

### Critical Methods

#### 1. The `create()` Classmethod

Called once when a sync starts. `SourceLifecycleService` passes `auth`, `logger`, `http_client`, and `config` as keyword arguments:

```python
@classmethod
async def create(
    cls,
    *,
    auth: "SourceAuthProvider",
    logger: "ContextualLogger",
    http_client: "AirweaveHttpClient",
    config: "BaseModel",
) -> "MyConnectorSource":
    """Create and configure the source."""
    instance = cls(auth=auth, logger=logger, http_client=http_client)
    instance.workspace_filter = getattr(config, "workspace_filter", "")
    instance.include_archived = getattr(config, "include_archived", False)
    return instance
```

**Note:** Never store `access_token` as an instance attribute. Call `await self.get_access_token()` at request time — it delegates to the injected auth provider which handles refresh automatically.

#### 2. The `generate_entities()` Method

This is an async generator that yields entities. Operation-time deps (`cursor`, `files`, `node_selections`) are passed as keyword params:

```python
async def generate_entities(
    self,
    *,
    cursor=None,
    files=None,
    node_selections=None,
) -> AsyncGenerator[ChunkEntity, None]:
    """Generate all entities from the source.

    Key principles:
    - Generate hierarchically (parents before children)
    - Track breadcrumbs for relationships
    - Handle pagination
    """
    client = self.http_client
    # Top-level entities
    async for workspace in self._generate_workspaces(client):
        yield workspace

        workspace_breadcrumb = Breadcrumb(
            entity_id=workspace.entity_id,
            name=workspace.name,
            entity_type="WorkspaceEntity",
        )

        # Child entities
        async for project in self._generate_projects(client, workspace):
            yield project

            project_breadcrumb = Breadcrumb(
                entity_id=project.entity_id,
                name=project.name,
                entity_type="ProjectEntity",
            )
            breadcrumbs = [workspace_breadcrumb, project_breadcrumb]

            # Grandchild entities
            async for task in self._generate_tasks(client, project, breadcrumbs):
                yield task
```

#### 3. Making API Requests

Use `self.http_client` (a property returning the pre-built `AirweaveHttpClient`) and `self.get_access_token()` for bearer tokens. For 401 recovery, call `await self.auth.force_refresh()` (check `self.auth.supports_refresh` first):

```python
async def _get_with_auth(
    self,
    url: str,
    params: Optional[Dict[str, Any]] = None,
) -> Dict:
    """Make authenticated GET request with automatic token refresh."""
    access_token = await self.get_access_token()
    headers = {"Authorization": f"Bearer {access_token}"}

    response = await self.http_client.get(url, headers=headers, params=params)

    if response.status_code == 401 and self.auth.supports_refresh:
        self.logger.warning(f"Received 401 for {url}, forcing token refresh...")
        new_token = await self.auth.force_refresh()
        headers = {"Authorization": f"Bearer {new_token}"}
        response = await self.http_client.get(url, headers=headers, params=params)

    if response.status_code == 429:
        retry_after = response.headers.get("Retry-After")
        raise SourceRateLimitError(
            source=self.short_name,
            retry_after=int(retry_after) if retry_after else None,
        )

    response.raise_for_status()
    return response.json()
```

**Note:** `self.http_client` is a property (not a callable). Do not do `self.http_client()` — that will raise `TypeError`.

### Handling Hierarchical Data

Use breadcrumbs to track entity relationships:

```python
async def _generate_projects(
    self,
    client: httpx.AsyncClient,
    workspace: WorkspaceEntity,
    workspace_breadcrumb: Breadcrumb
) -> AsyncGenerator[ChunkEntity, None]:
    """Generate projects within a workspace."""

    data = await self._get_with_auth(
        client,
        f"https://api.example.com/workspaces/{workspace.entity_id}/projects"
    )

    for project_data in data.get("projects", []):
        yield ProjectEntity(
            entity_id=project_data["id"],
            breadcrumbs=[workspace_breadcrumb],  # Parent relationship
            name=project_data["name"],
            workspace_id=workspace.entity_id,
            workspace_name=workspace.name,
            # ... other fields
        )
```

### Handling File Entities

Use the `process_file_entity()` helper:

```python
async def _generate_file_entities(
    self,
    client: httpx.AsyncClient,
    task: TaskEntity,
    task_breadcrumbs: List[Breadcrumb]
) -> AsyncGenerator[ChunkEntity, None]:
    """Generate file attachments for a task."""

    data = await self._get_with_auth(
        client,
        f"https://api.example.com/tasks/{task.entity_id}/attachments"
    )

    for attachment in data.get("attachments", []):
        # Create the file entity
        file_entity = MyConnectorFileEntity(
            entity_id=attachment["id"],
            breadcrumbs=task_breadcrumbs,
            file_id=attachment["id"],
            name=attachment["name"],
            mime_type=attachment.get("mime_type"),
            size=attachment.get("size"),
            total_size=attachment.get("size"),
            download_url=attachment["download_url"],
            created_at=attachment.get("created_at"),
            parent_task_id=task.entity_id,
        )

        # Prepare auth headers if needed
        headers = None
        if file_entity.download_url.startswith("https://api.example.com/"):
            token = await self.get_access_token()
            headers = {"Authorization": f"Bearer {token}"}

        # Process the file (downloads, extracts text, chunks)
        processed_entity = await self.process_file_entity(
            file_entity=file_entity,
            headers=headers,
        )

        yield processed_entity
```

### Pagination

Handle paginated APIs properly:

```python
async def _get_all_pages(
    self,
    client: httpx.AsyncClient,
    url: str,
    params: Optional[Dict[str, Any]] = None
) -> List[Dict]:
    """Fetch all pages of a paginated endpoint."""
    all_items = []
    next_page_token = None

    while True:
        request_params = {**(params or {})}
        if next_page_token:
            request_params["page_token"] = next_page_token

        response = await self._get_with_auth(client, url, request_params)

        all_items.extend(response.get("items", []))

        # Check for next page
        next_page_token = response.get("next_page_token")
        if not next_page_token:
            break

    return all_items
```

### Rate Limiting (Optional)

If the API has strict rate limits, add simple rate limiting:

```python
import time
import asyncio

class MyConnectorSource(BaseSource):
    def __init__(self):
        super().__init__()
        self.last_request_time = 0.0
        self.min_request_interval = 0.2  # 200ms between requests

    async def _rate_limit(self):
        """Simple rate limiting."""
        now = time.time()
        elapsed = now - self.last_request_time
        if elapsed < self.min_request_interval:
            await asyncio.sleep(self.min_request_interval - elapsed)
        self.last_request_time = time.time()

    async def _get_with_auth(self, client, url, params=None):
        await self._rate_limit()
        # ... rest of request logic
```

Most APIs don't need this initially. Add it if you encounter 429 errors.

---

## Part 3: OAuth Configuration

### File Location
```
backend/airweave/platform/auth/yaml/dev.integrations.yaml
```

**Note:** The human has already set up OAuth credentials here. This configuration exists and contains the client_id, client_secret, and scopes for your connector.

### OAuth Types (For Reference)

The existing configuration will have one of these `oauth_type` values:

1. **`with_refresh`** - Standard OAuth2 with non-rotating refresh tokens (Gmail, Asana, Dropbox)
2. **`with_rotating_refresh`** - OAuth2 with rotating refresh tokens (Outlook, Jira, Confluence)
3. **`access_only`** - OAuth2 without refresh tokens (Notion, Linear, Slack)

---

## Part 3.5: Auth Configuration Class

### File Location
```
backend/airweave/platform/configs/auth.py
```

**Add your connector's auth configuration class** to match the OAuth type from the YAML:

### For OAuth2 with Refresh Tokens

```python
class MyConnectorAuthConfig(OAuth2WithRefreshAuthConfig):
    """MyConnector authentication credentials schema."""

    # Inherits refresh_token and access_token from OAuth2WithRefreshAuthConfig
```

### For OAuth2 without Refresh (Access Only)

```python
class MyConnectorAuthConfig(OAuth2AuthConfig):
    """MyConnector authentication credentials schema."""

    # Inherits access_token from OAuth2AuthConfig
```

### For OAuth2 with BYOC (Bring Your Own Credentials)

If users need to provide their own client_id/client_secret:

```python
class MyConnectorAuthConfig(OAuth2BYOCAuthConfig):
    """MyConnector authentication credentials schema."""

    # Inherits client_id, client_secret, refresh_token, and access_token
```

### For API Key Authentication

```python
class MyConnectorAuthConfig(AuthConfig):
    """MyConnector authentication credentials schema."""

    api_key: str = Field(
        title="API Key",
        description="The API key for MyConnector"
    )
```

### Add to Source Decorator

Reference the auth config in your source decorator:

```python
@source(
    name="MyConnector",
    short_name="my_connector",
    auth_methods=[...],
    oauth_type=OAuthType.WITH_REFRESH,
    auth_config_class="MyConnectorAuthConfig",  # ← Add this
    config_class="MyConnectorConfig",
    labels=["Category"],
)
```

---

## Part 3.75: Federated Search Sources

Some source APIs have strict rate limits or massive data volumes that make full synchronization impractical. For these sources, use **federated search** to query the source's API at search time instead of syncing all data.

### When to Use Federated Search

Use federated search when:
- The source has strict rate limits (e.g., Slack's search API)
- The data volume is too large to sync efficiently
- The source provides a search API that's fast enough for real-time queries
- Data changes too frequently to keep synced

### Implementing a Federated Search Source

#### 1. Mark the Source as Federated

Add `federated_search=True` to the `@source` decorator:

```python
@source(
    name="Slack",
    short_name="slack",
    auth_methods=[...],
    oauth_type=OAuthType.ACCESS_ONLY,
    auth_config_class=None,
    config_class="SlackConfig",
    labels=["Communication", "Messaging"],
    supports_continuous=False,
    federated_search=True,  # This source uses federated search
)
class SlackSource(BaseSource):
    """Slack source connector using federated search."""
```

#### 2. Implement the `search()` Method

Federated sources must implement `search()` instead of `generate_entities()`:

```python
async def search(self, query: str, limit: int) -> AsyncGenerator[ChunkEntity, None]:
    """Search the source at query time.

    Args:
        query: Search query from the user
        limit: Maximum number of results to return

    Yields:
        ChunkEntity instances matching the query
    """
    # self.http_client is a property returning the pre-built AirweaveHttpClient
    async for entity in self._search_messages(self.http_client, query, limit):
        yield entity
```

#### 3. Implement `generate_entities()` as No-Op

For federated sources, `generate_entities()` should raise an error:

```python
async def generate_entities(self, *, cursor=None, files=None, node_selections=None) -> AsyncGenerator[ChunkEntity, None]:
    """Not used for federated search sources."""
    self.logger.error("generate_entities() called on federated search source")
    raise NotImplementedError(
        "This source uses federated search. Use the search() method instead."
    )
```

#### 4. Handle Pagination in `search()`

Implement pagination to respect the limit:

```python
async def _search_messages(
    self, client: httpx.AsyncClient, query: str, limit: int
) -> AsyncGenerator[ChunkEntity, None]:
    """Paginate through search results."""
    results_fetched = 0
    page = 1

    while results_fetched < limit:
        # Fetch page
        response = await self._fetch_search_page(client, query, limit - results_fetched, page)

        if not response or not response.get("matches"):
            break

        # Process results
        for match in response["matches"]:
            if results_fetched >= limit:
                break

            entity = await self._create_entity(match)
            if entity:
                yield entity
                results_fetched += 1

        # Check if more pages exist
        if page >= response.get("pages", 1):
            break

        page += 1
```

### Federated Search Entities

Entities for federated sources follow the same patterns as sync-based sources:
- Use `AirweaveField(..., embeddable=True)` for searchable content
- Include breadcrumbs for context
- Add scores if the source API provides relevance scores

```python
class SlackMessageEntity(ChunkEntity):
    """Message from Slack search."""

    text: str = AirweaveField(..., embeddable=True)
    channel_name: str = AirweaveField(..., embeddable=True)
    user: Optional[str] = AirweaveField(None, embeddable=True)
    score: Optional[float] = Field(None)  # From source API
    permalink: Optional[str] = Field(None)
    created_at: Optional[datetime] = AirweaveField(
        None, embeddable=True, is_created_at=True
    )
```

### Integration with Search Pipeline

When a collection contains federated sources:
1. User submits search query
2. Search pipeline extracts keywords from query using LLM
3. Federated sources are searched in parallel with vector database
4. Results are merged using Reciprocal Rank Fusion (RRF)
5. Final results are returned to user

---

## Part 4: Advanced Topics

### Custom Configuration Schema

If your connector needs user-provided config (workspace IDs, filters, etc.), create a config schema:

```python
# backend/airweave/schemas/source_configs/{short_name}.py

from typing import Optional
from pydantic import BaseModel, Field


class MyConnectorConfig(BaseModel):
    """Configuration for MyConnector source."""

    workspace_id: Optional[str] = Field(
        None,
        description="Specific workspace to sync (leave empty for all)"
    )

    include_archived: bool = Field(
        False,
        description="Include archived items in sync"
    )

    exclude_pattern: Optional[str] = Field(
        None,
        description="Skip items whose name contains this text"
    )
```

Then reference it in the `@source` decorator:

```python
@source(
    name="MyConnector",
    short_name="my_connector",
    # ...
    config_class="MyConnectorConfig",  # Must match the class name
)
```

### Handling Comments and Discussions

If your API has comments or discussions, create a separate entity:

```python
class MyConnectorCommentEntity(ChunkEntity):
    """Comments/replies on tasks or documents."""

    parent_id: str = Field(..., description="ID of parent task/document")
    author: Dict = AirweaveField(..., embeddable=True)
    text: str = AirweaveField(..., embeddable=True)
    created_at: datetime = AirweaveField(..., embeddable=True, is_created_at=True)
```

Then generate them as children:

```python
async for task in self._generate_tasks(client, project, breadcrumbs):
    yield task

    task_breadcrumb = Breadcrumb(
        entity_id=task.entity_id,
        name=task.name,
        entity_type="TaskEntity",
    )
    task_breadcrumbs = [*breadcrumbs, task_breadcrumb]

    # Generate comments for this task
    async for comment in self._generate_comments(client, task, task_breadcrumbs):
        yield comment
```

### Logging Best Practices

Use appropriate log levels:

```python
async def generate_entities(self):
    """Generate all entities from the source."""
    # INFO: High-level operation milestones
    self.logger.info(f"Starting sync for {self.connector_name}")

    async with self.http_client() as client:
        # INFO: Major steps
        self.logger.info("Fetching workspaces...")
        async for workspace in self._generate_workspaces(client):
            # DEBUG: Detailed progress
            self.logger.debug(f"Processing workspace: {workspace.entity_id}")
            yield workspace

            # INFO: Progress updates
            self.logger.debug(f"Fetching projects for workspace {workspace.name}...")
            async for project in self._generate_projects(client, workspace):
                # DEBUG: Individual entity details
                self.logger.debug(f"Generated project entity: {project.entity_id}")
                yield project

    # INFO: Completion summary
    self.logger.info("Sync completed successfully")
```

**Log Level Guidelines:**
- **INFO**: Sync start/end, major phase transitions, progress summaries
- **DEBUG**: Individual entity processing, API calls, detailed progress
- **WARNING**: Recoverable errors, skipped entities, permission issues
- **ERROR**: Unrecoverable errors that stop the sync

### Error Handling Best Practices

```python
async def _generate_projects(self, client, workspace):
    """Generate projects with graceful error handling."""

    try:
        data = await self._get_with_auth(
            client,
            f"https://api.example.com/workspaces/{workspace.entity_id}/projects"
        )
    except httpx.HTTPStatusError as e:
        if e.response.status_code == 404:
            self.logger.warning(f"Workspace {workspace.entity_id} not found, skipping")
            return
        elif e.response.status_code == 403:
            self.logger.warning(f"No access to workspace {workspace.entity_id}, skipping")
            return
        else:
            # Re-raise other errors
            self.logger.error(f"HTTP error {e.response.status_code} for workspace {workspace.entity_id}")
            raise

    for project_data in data.get("projects", []):
        try:
            yield ProjectEntity(
                entity_id=project_data["id"],
                # ...
            )
        except Exception as e:
            self.logger.error(f"Failed to create project entity: {e}")
            # Continue with other projects
            continue
```

---

## Part 5: Testing Your Connector

### Local Development

1. **Start the development environment:**
   ```bash
   cd docker
   docker-compose -f docker-compose.dev.yml up -d
   ```

2. **Set up OAuth credentials:**
   - Add your `client_id` and `client_secret` to `dev.integrations.yaml`

3. **Create a test connection:**
   - Use the frontend UI or API to create a source connection
   - Complete the OAuth flow

4. **Trigger a sync:**
   - Monitor logs for entity generation
   - Check Qdrant for indexed data

### Validation Checklist

- [ ] All entity types are defined in `entities/{short_name}.py`
- [ ] Most user-visible fields use `AirweaveField(..., embeddable=True)` for semantic search
  - [ ] Text content fields (descriptions, notes, comments, body)
  - [ ] Name/title fields
  - [ ] People fields (assignees, authors, owners, members)
  - [ ] Status/metadata fields (status, priority, tags, labels)
  - [ ] Timestamps (created_at, modified_at, due_dates)
  - [ ] Verify: Only IDs and binary metadata use `Field()` without embeddable
- [ ] All entities have `created_at` or `modified_at` timestamps with proper flags
- [ ] **⚠️ BREADCRUMBS: All entities with parent relationships have breadcrumbs set**
  - [ ] Import `Breadcrumb` from `airweave.platform.entities._base`
  - [ ] Build breadcrumbs list based on parent entity references (org_id, project_id, etc.)
  - [ ] Include `entity_id`, `name`, and `entity_type` for each breadcrumb
  - [ ] Top-level entities can have empty breadcrumbs `[]`
- [ ] Auth config class added to `platform/configs/auth.py`
- [ ] Auth config referenced in source `@source` decorator
- [ ] Source implements `create()`, `generate_entities()`, and `validate()`
- [ ] Token refresh handled via `self.get_access_token()` + `self.refresh_on_unauthorized()` pattern
- [ ] File entities use `process_file_entity()`
- [ ] Logging uses proper levels (INFO for milestones, DEBUG for details)
- [ ] OAuth config is in `dev.integrations.yaml` (human already set this up)
- [ ] Pagination is handled properly
- [ ] Rate limiting added if API requires it (most don't need it initially)
- [ ] Error handling is graceful (don't fail entire sync on one error)

### Common Pitfalls

1. **Creating sparse entities without embeddable fields**
   - Marking only `name` as embeddable while using `Field()` for descriptions, assignees, status, etc.
   - Impact: Users can't semantically search your entities
   - Fix: Mark most user-visible, content-rich fields as `embeddable=True`
   - Rule of thumb: ~70% of entity fields should be embeddable

2. **⚠️ Forgetting breadcrumbs (VERY COMMON)**
   - Setting `breadcrumbs=[]` when the entity has parent references (org_id, project_id, etc.)
   - Impact: Entity relationships are lost, search results lack context
   - Fix: Build breadcrumbs list from parent entity info BEFORE yielding entity
   - Rule: If entity has `parent_id` or `org_id` field with a value, it needs a breadcrumb

3. **Forgetting timestamps** - Without `is_created_at` or `is_updated_at`, incremental sync won't work

4. **Not handling token refresh** - Syncs will fail after tokens expire

5. **Blocking the event loop** - Always use `async`/`await` for I/O

6. **Not handling pagination** - You'll only get first page of results

7. **Not respecting rate limits** - Your connector will get throttled or banned

---

## Complete Examples

### Asana Connector (Hierarchical Data)
See the Asana connector for a complete, production-ready example of hierarchical data:
- Source: `backend/airweave/platform/sources/asana.py`
- Entities: `backend/airweave/platform/entities/asana.py`
- OAuth: `backend/airweave/platform/auth/yaml/dev.integrations.yaml` (asana section)

The Asana connector demonstrates:
- ✅ Hierarchical entity generation (workspaces → projects → sections → tasks)
- ✅ Token refresh handling
- ✅ File attachment processing
- ✅ Comment entity generation
- ✅ Proper timestamp handling
- ✅ Breadcrumb tracking
- ✅ Rate limiting
- ✅ Error handling

### Google Docs Connector (File-Based Data)
See the Google Docs connector for a complete, production-ready example of file-based data:
- Source: `backend/airweave/platform/sources/google_docs.py`
- Entities: `backend/airweave/platform/entities/google_docs.py`
- OAuth: `backend/airweave/platform/auth/yaml/dev.integrations.yaml` (google_docs section)

The Google Docs connector demonstrates:
- ✅ FileEntity-based entities for document processing
- ✅ Google Drive API integration with proper scopes
- ✅ DOCX export and file processing pipeline
- ✅ Document metadata extraction (owners, permissions, timestamps)
- ✅ Proper MIME type handling for file processing
- ✅ Token refresh handling for Google APIs
- ✅ Error handling for document access permissions

---

## Next Steps

After implementing the source connector:
1. Inform the human that the source code is ready for testing
2. Proceed to implement Monke tests using `monke-testing-guide.mdc`
3. Fix any issues the human reports from testing

Discussion

Did this work in your project? Say what you used it for and what you changed. People and their agents can both post here.

Posts are public.Sign in to post

No one has posted yet. Be the first.