airweave / rules
airweave-ai/airweave/.cursor/rules/monke.mdc
End-to-end testing framework for Airweave connectors using real external APIs Monke validates Airweave's data synchronization pipeline by creating real test data in external systems, triggering sync jobs, verifying data in the vector database, and testing update/deletion scenarios. The framework separates configuration from runtime state, uses pluggable authentication, and supports parallel execution. Test Runner (runner.py) - Entry point that orchestrates parallel test execution across connectors Test Flow (core/flow.py) - Executes complete lifecycle for a connector: setup → test steps → cleanup
- Reads credentials
---
globs: **/monke/**
alwaysApply: false
---
# Monke - Airweave Connector Testing Framework
**End-to-end testing framework for Airweave connectors using real external APIs**
## Overview
Monke validates Airweave's data synchronization pipeline by creating real test data in external systems, triggering sync jobs, verifying data in the vector database, and testing update/deletion scenarios. The framework separates configuration from runtime state, uses pluggable authentication, and supports parallel execution.
---
## Core Architecture
**Test Runner** (`runner.py`) - Entry point that orchestrates parallel test execution across connectors
**Test Flow** (`core/flow.py`) - Executes complete lifecycle for a connector: setup → test steps → cleanup
**Test Steps** (`core/steps.py`) - Individual operations (create, sync, verify, update, delete) with retry logic
**Configuration** (`core/config.py`) - YAML-based config with classes: `TestConfig`, `ConnectorConfig`, `TestFlowConfig`, `DeletionConfig`
**Test Context** (`core/context.py`) - Runtime state separate from config: tracks entities, infrastructure IDs, metrics
**Infrastructure** (`core/infrastructure.py`) - Creates/tears down test collections and source connections
**Authentication** (`auth/`) - Pluggable auth brokers (`BaseAuthBroker`, `ComposioBroker`) for credential resolution
---
## Test Flow
Default test steps execute in sequence:
```
1. collection_cleanup # Clean leftover test collections
2. cleanup # Clean entities in external system
3. create # Create test entities
4. sync # Trigger Airweave sync
5. verify # Verify entities in vector DB
6. update # Update entities
7. sync → verify # Verify updates
8. partial_delete # Delete subset of entities
9. sync → verify_partial_deletion
10. verify_remaining_entities
11. complete_delete # Delete all entities
12. sync → verify_complete_deletion
13. cleanup # Final cleanup
14. collection_cleanup # Delete test collection
```
Steps are customizable per connector via YAML config.
---
## Bongos
**Bongos** create, update, and delete real test data in external systems.
### Structure
Each bongo inherits from `BaseBongo` and implements:
- `create_entities()` - Create test data
- `update_entities()` - Modify test data
- `delete_entities()` - Remove all test data
- `delete_specific_entities(entities)` - Remove specific entities
- `cleanup()` - Force cleanup remaining artifacts
Bongos are auto-discovered by `BongoRegistry` via the `connector_type` class attribute.
### Example
```python
class GitHubBongo(BaseBongo):
def __init__(self, credentials: Dict[str, Any], config: Dict[str, Any]):
super().__init__(credentials)
self.repo_name = config["repo_name"] # Now from config_fields
self.token = credentials["personal_access_token"]
async def create_entities(self) -> List[Dict[str, Any]]:
# Create files in GitHub, return metadata
async def update_entities(self) -> List[Dict[str, Any]]:
# Update files
async def delete_entities(self) -> List[str]:
# Delete all tracked entities
```
---
## Configuration
Configs live in `configs/{connector_type}.yaml`:
```yaml
name: "GitHub Connector Test"
description: "End-to-end test of GitHub connector"
connector:
name: "github_test"
type: "github"
auth_mode: direct # or "provider"
auth_fields:
personal_access_token: MONKE_GITHUB_PERSONAL_ACCESS_TOKEN
config_fields:
repo_name: MONKE_GITHUB_REPO_NAME # Moved from auth_fields
branch: "main"
post_create_sleep_seconds: 15
test_flow:
steps: [cleanup, create, sync, verify, update, sync, verify, complete_delete, sync, verify_complete_deletion, cleanup]
entity_count: 3
deletion:
partial_delete_count: 1
verify_partial_deletion: true
verification:
score_threshold: 0.1
max_retries: 3
retry_delay_seconds: 5
```
Environment variables are substituted using `${VAR_NAME}` syntax.
---
## Authentication
### Direct Mode
Credentials specified via `auth_fields` with env var names:
```yaml
auth_mode: direct
auth_fields:
personal_access_token: MONKE_GITHUB_PERSONAL_ACCESS_TOKEN
```
### Provider Mode
External auth providers (e.g., Composio) fetch credentials:
```bash
DM_AUTH_PROVIDER=composio
DM_AUTH_PROVIDER_ID=provider_id
DM_AUTH_PROVIDER_API_KEY=api_key
GITHUB_AUTH_PROVIDER_AUTH_CONFIG_ID=ac_xxx
GITHUB_AUTH_PROVIDER_ACCOUNT_ID=ca_xxx
```
Auth resolution:
1. `credentials_resolver` determines mode
2. Direct: resolves env vars from `auth_fields`
3. Provider: uses `ComposioBroker` to fetch credentials
4. Credentials passed to bongo and Airweave source connection
---
## Content Generation
Generators in `generation/{connector}.py` create realistic test data with embedded tracking tokens:
```python
from monke.client.llm import LLMClient
async def generate_github_content(token: str) -> str:
llm = LLMClient()
prompt = f"Generate Python code. Include '{token}' in a comment."
return await llm.generate(prompt)
```
Pydantic schemas in `generation/schemas/{connector}.py` define entity structures.
---
## Verification
After each sync, verification steps:
1. Search Qdrant vector DB for tracking tokens
2. Check relevance scores meet threshold (default: 0.1)
3. Verify expected result count
4. Retry with exponential backoff on failures
Types:
- **Standard**: All entities exist with good scores
- **Partial Deletion**: Specific entities are gone
- **Remaining Entities**: Non-deleted entities still exist
- **Complete Deletion**: Zero results for all tokens
---
Each has: bongo implementation, config file, content generator, schema definitions.
---
## Usage
```bash
# Test changed connectors (auto-detect from git diff)
./monke.sh
# Test specific connector
./monke.sh github
# Test multiple in parallel
./monke.sh github asana notion
# Test all connectors
./monke.sh --all
# Python runner directly
python monke/runner.py github --max-concurrency 5
```
### Environment Setup
```bash
AIRWEAVE_API_URL=http://localhost:8001
AIRWEAVE_API_KEY=optional_api_key
OPENAI_API_KEY=sk-...
# Direct auth mode
MONKE_GITHUB_PERSONAL_ACCESS_TOKEN=ghp_...
MONKE_GITHUB_REPO_NAME=owner/repo
```
---
## Adding a New Connector
### 1. Create Bongo (`bongos/myapp.py`)
```python
from monke.bongos.base_bongo import BaseBongo
class MyAppBongo(BaseBongo):
connector_type = "myapp"
async def create_entities(self): pass
async def update_entities(self): pass
async def delete_entities(self): pass
async def delete_specific_entities(self, entities): pass
async def cleanup(self): pass
```
### 2. Create Generator (`generation/myapp.py`)
```python
async def generate_myapp_content(token: str):
llm = LLMClient()
return await llm.generate(f"Generate content with token: {token}")
```
### 3. Create Schema (`generation/schemas/myapp.py`)
```python
from pydantic import BaseModel
class MyAppEntity(BaseModel):
id: str
content: str
token: str
```
### 4. Create or Update Config (`configs/myapp.yaml`)
```yaml
name: "MyApp Test"
connector:
type: "myapp"
auth_mode: direct
auth_fields:
api_key: MONKE_MYAPP_API_KEY
config_fields: {}
entity_count: 5
test_flow:
steps: [cleanup, create, sync, verify, update, sync, verify, complete_delete, sync, verify_complete_deletion, cleanup]
verification:
score_threshold: 0.1
```
### 5. Test
```bash
./monke.sh myapp
```
---
## Key Design Principles
1. **Separation**: Configuration vs runtime state (TestConfig vs TestContext)
2. **Pluggable Auth**: Multiple credential sources via auth brokers
3. **Atomic Steps**: Independent, retryable test steps
4. **Auto Cleanup**: Cleanup runs even on failure
5. **Real APIs**: Tests use real external services
6. **Parallel**: Multiple connectors tested simultaneously
7. **Event-Driven**: Real-time progress via event bus
8. **Configurable**: Test flows customizable per connector
---
## Common Pitfalls & Troubleshooting
### Entity ID Collisions Across Types
**Problem**: When a source has multiple entity types (e.g., Pipedrive has persons, deals, notes), different types can share the same numeric ID from the external API (e.g., `deal/75` and `note/75` both exist).
**Symptoms**:
- `verify_partial_deletion` fails because wrong entities are matched
- Entities "found" in Qdrant that should be deleted
- UUID collisions in vector storage
**Solution**: Use **type-prefixed IDs** for internal tracking:
```python
ent = {
"type": "deal",
"id": f"deal_{pipedrive_id}", # Unique internal ID: "deal_75"
"pipedrive_id": pipedrive_id, # Raw API ID: 75
"token": token,
# ...
}
```
- Use `id` for Monke tracking, comparison, and Qdrant paths
- Use `pipedrive_id` (or equivalent) for actual API calls
### Deletion Order for Linked Entities
**Problem**: APIs often prevent deleting parent entities while children reference them.
**Solution**: Delete in dependency order (children first):
```python
# Notes → Activities → Leads → Deals → Persons → Organizations → Products
all_entities = (
self._notes # Depend on persons/deals
+ self._activities # Depend on persons/deals/orgs
+ self._leads # Depend on persons/orgs
+ self._deals # Depend on persons/orgs
+ self._persons # Depend on orgs
+ self._organizations # Independent
+ self._products # Independent
)
```
### Token Embedding in Searchable Fields
**Problem**: Qdrant verification searches for tokens. If the token isn't in an embeddable field, verification fails.
**Solution**:
1. Schema: Mark searchable fields with `embeddable=True` in Pydantic schemas
2. Generator: Embed token in the searchable field (e.g., `title`, `name`, `content`)
3. Bongo: Explicitly ensure token presence when creating:
```python
# Ensure token is in searchable field even if generator omits it
title_with_token = f"{deal.title} [{token}]"
payload = {"title": title_with_token, ...}
```
### Debugging Entity Verification
Add diagnostic logging to trace what's being created vs what's searched:
```python
self.logger.info(f"📝 DIAGNOSTIC: Creating deal with title='{title_with_token}'")
# After API response
actual_title = deal_data.get("title", "")
if token not in actual_title:
self.logger.warning(f"⚠️ Token {token} not in returned title: {actual_title}")
```
### Syntax Errors After Editing Bongos
**Problem**: Indentation errors after multi-line edits to Python files.
**Solution**: Always verify syntax after editing:
```bash
python3 -m py_compile monke/bongos/{connector}.py && echo "Syntax OK"
```
---
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.
No one has posted yet. Be the first.

