Cell Integration with SRPVDAL Pipeline

Date: November 7, 2025
Objective: Ensure all cell operations flow through Boss Agent SRPVDAL orchestration


Architecture Overview

Current State

User Request → Cell 2 directly → Cell 3 KG → Cell 4 Ontology
             ↓
         Direct processing (no SRPVDAL orchestration)

Target State (SRPVDAL Integration)

User Request → Boss Agent SRPVDAL Orchestrator
             ↓
         SENSE → REASON → DECIDE → ACT → LEARN
             ↓           ↓         ↓      ↓
         Cell 2       Cell 3    Cell 4   Update KG
         Extract      Analyze   Validate  Store

Cell 2: Data Extraction & PII Masking

Current Capabilities

Current Endpoints

SRPVDAL Integration Required

Cell 2 operations map to SRPVDAL stages:

SENSE Stage: Data ingestion and schema inference - Load data from source - Detect data format - Sample for schema - Profile data quality

REASON Stage: PII detection and validation - Identify PII fields - Validate against ontology (Cell 4) - Assess data quality - Determine masking strategy

DECIDE Stage: Processing strategy - Choose masking method - Decide on forwarding to KG - Determine batch size - Select processing mode

ACT Stage: Execute transformation - Mask PII data - Transform to unified schema - Forward to Cell 3 KG - Store in BigQuery

LEARN Stage: Update patterns - Store successful schemas - Record PII patterns - Update validation rules - Optimize extraction strategy


Integration Approach

Flow:

Boss Agent receives: "Ingest Facebook data"
  ↓ SENSE: Classify as data extraction task
  ↓ REASON: Determine Cell 2 needed for Facebook ingestion
  ↓ DECIDE: Route to Cell 2 with parameters
  ↓ ACT: Call Cell 2 /api/v1/extract endpoint
  ↓ LEARN: Store ingestion metadata in KG

Implementation:

# In miz-oki-adk-agents/boss/app.py

async def handle_data_extraction_request(message: str, context: dict) -> dict:
    """
    Handle data extraction requests through SRPVDAL pipeline
    Routes to Cell 2 for processing
    """
    # SENSE: Classify request
    if "ingest" in message.lower() or "extract" in message.lower():
        source_type = detect_source_type(message)  # facebook, shopify, etc.

        # REASON: Determine parameters
        extraction_params = {
            "source": context.get("source_url") or "auto-detect",
            "format": context.get("format", "auto"),
            "ontology_id": "unified_marketing_v1",
            "metadata": {
                "requested_by": context.get("userId"),
                "request_type": "boss_orchestrated"
            }
        }

        # DECIDE: Choose Cell 2 endpoint
        cell2_url = "https://miz-oki-cell2-698171499447.us-central1.run.app"

        # ACT: Execute via Cell 2
        async with httpx.AsyncClient() as client:
            response = await client.post(
                f"{cell2_url}/api/v1/extract",
                json=extraction_params
            )
            result = response.json()

        # LEARN: Record in KG
        await record_extraction_event(result)

        return {
            "extraction_id": result["extraction_id"],
            "records_processed": result["records_processed"],
            "kg_ingestion_id": result["kg_ingestion_id"]
        }

Option 2: Cell 2 Reports to SRPVDAL (Hybrid)

Flow:

Boss Agent → Cell 2 /api/v1/extract
             ↓
         Cell 2 processes internally
             ↓
         Cell 2 sends SRPVDAL events to Boss Agent
             ↓
         Boss Agent updates orchestration state

Implementation:

# In src/cells/cell02/Cell2.py

class SRPVDALReporter:
    """Reports Cell 2 operations to Boss Agent SRPVDAL orchestrator"""

    def __init__(self, boss_agent_url: str):
        self.boss_agent_url = boss_agent_url
        self.http_client = httpx.AsyncClient()

    async def report_srpvdal_event(self, stage: str, operation: str, data: dict):
        """Report SRPVDAL stage completion to Boss Agent"""
        try:
            await self.http_client.post(
                f"{self.boss_agent_url}/srpvdal/event",
                json={
                    "cell_id": 2,
                    "stage": stage,
                    "operation": operation,
                    "data": data,
                    "timestamp": datetime.utcnow().isoformat()
                }
            )
        except Exception as e:
            logger.warning(f"Failed to report SRPVDAL event: {e}")

# In extract_and_process method:
async def extract_and_process(self, request: ExtractionRequest) -> ExtractionResult:
    # SENSE stage
    await srpvdal_reporter.report_srpvdal_event("SENSE", "data_loaded", {"source": request.source})
    df = await self._load_and_prepare_data(request)

    # REASON stage  
    await srpvdal_reporter.report_srpvdal_event("REASON", "pii_detected", {"fields": pii_summary})
    masked_df, pii_summary = await self._process_pii(df)

    # DECIDE stage
    await srpvdal_reporter.report_srpvdal_event("DECIDE", "validation_complete", {"valid": is_valid})
    is_valid, errors = await self._validate_schema(schema, ...)

    # ACT stage
    await srpvdal_reporter.report_srpvdal_event("ACT", "kg_forwarded", {"kg_id": kg_id})
    kg_id = await self._forward_to_kg_if_valid(...)

    # LEARN stage
    await srpvdal_reporter.report_srpvdal_event("LEARN", "extraction_complete", result)

    return result

Start with Option 1 (Boss Agent orchestrates Cell 2):

Advantages:

Steps to Implement:

  1. Add Cell 2 to Boss Agent's cell registry python # In miz-oki-adk-agents/boss/app.py CELL_URLS = { 1: "http://boss-agent:8000", 2: "https://miz-oki-cell2-698171499447.us-central1.run.app", # ADD THIS 3: "http://cell3-graph-engine:8003", # ... rest }

  2. Add data extraction intent to classifier python # In SRPVDALOrchestrator.classify_intent() intent_keywords = { "data_extraction": ["ingest", "extract", "load data", "import", "upload"], "pii_masking": ["mask", "anonymize", "pii", "privacy"], # ... existing intents }

  3. Add Cell 2 specialist agent python # In /api/agents response { "id": "data-extractor", "name": "Data Extraction & PII Masking Agent", "type": "specialist", "capabilities": ["data-ingestion", "pii-masking", "schema-validation"], "status": "active", "srpvdal_stage": "SENSE", "cell": 2 }

  4. Wire Cell 2 calls in orchestrator python # In SRPVDALOrchestrator.call_specialist_agent() elif agent_id == "data-extractor": async with httpx.AsyncClient() as client: response = await client.post( f"{self.cell_urls[2]}/api/v1/extract", json={ "source": context.get("data_source"), "format": "auto", "ontology_id": "unified_marketing_v1" } ) return response.json()


Next Steps

Phase 1: Cell 2 SRPVDAL Integration (2 hours)

  1. Add Cell 2 URL to Boss Agent configuration
  2. Add data extraction intents to classifier
  3. Wire Cell 2 calls in orchestrator
  4. Test end-to-end: Boss Agent → Cell 2 → Cell 3

Phase 2: Other Cell Integration (4-6 hours)

For each cell (3, 4, 7, 11, 21, 23-32): 1. Identify cell purpose and operations 2. Map operations to SRPVDAL stages 3. Add cell URL to Boss Agent registry 4. Add cell-specific intents 5. Wire cell calls in orchestrator 6. Test integration

Phase 3: Multi-Cell Workflows (2-3 hours)

Create complex workflows that coordinate multiple cells:

Example: "Analyze campaign attribution"
  SENSE: Cell 2 extracts campaign data
  REASON: Cell 3 KG analyzes relationships  
  REASON: Cell 7 performs causal inference
  DECIDE: Cell 11 optimizes attribution model
  ACT: Cell 15 generates recommendations
  LEARN: Cell 19 updates learning models

Cell Registry

Based on code inspection:

Cell Purpose SRPVDAL Stage URL
1 Discovery SENSE https://cell1-discovery-698171499447.us-central1.run.app
2 Data Extraction & PII SENSE https://miz-oki-cell2-698171499447.us-central1.run.app
3 Knowledge Graph SENSE/REASON https://miz-oki-cell3-698171499447.us-central1.run.app
4-6 (To be mapped) TBD TBD
7 Causal Inference REASON https://miz-oki-cell7-698171499447.us-central1.run.app
11 Attribution Optimizer DECIDE https://miz-oki-cell11-698171499447.us-central1.run.app
15 Campaign Executor ACT (Map from service list)
19 Learning Engine LEARN (Map from service list)
21 Smart Expert Router DECIDE https://miz-oki-cell21-698171499447.us-central1.run.app
23-25 Causal Discovery REASON https://miz-oki-cell23-698171499447.us-central1.run.app
26-32 Advanced Analytics Various (Map from service list)

Ready to Start?

Say "continue" and I'll: 1. Wire Cell 2 into Boss Agent SRPVDAL orchestrator 2. Add data extraction intents and routing 3. Test Cell 2 → Boss Agent → SRPVDAL flow 4. Then proceed with other cells systematically

← All docsView source on GitHub →