DESIGN-VINTAGE (March 2026) — retained for history. Neo4j references below describe the design substrate of that date, not the live backend: Neo4j was retired by owner decision 2026-08-09; the live KG is Firestore-backed, and KG writes go only through the governed Cell 3/24 writer path (
.claude/rules/03-canonical-architecture.md). Banner added by the truth-debt sweep, 2026-08-21.
Marketing Event Stream to Journey Graph Integration Plan
Status
As of March 2026, the repository contains most of the building blocks for a marketing-event-graph pipeline, but the implementation is still a composed system rather than one fully unified runtime path.
What exists today:
src/cells/cell02already ingests and normalizes multi-source marketing data.src/cells/cell02/enhanced_identity_resolver.pyalready performs deterministic identity resolution.src/cells/cell02/customer_journey_profiler.pyalready constructs unified journey profiles.services/cell03-intelligencealready builds knowledge-graph structures and attribution outputs.services/gemini-kg-pipelinealready performs Gemini-based graph enrichment.services/mizoki-journey-kgalready writes journey events to Neo4j and exposes prediction endpoints.- Boss already routes tools through
src/moe/mcp_boss_agent.pyandmcp/service_registry.yaml. - The marketing event graph runtime surfaces now exist under:
miz-oki-adk-agents/boss/marketing_event_graph_mcp.pymiz-oki-adk-agents/marketing_event_graph_service/main.py
The remaining gap is contract alignment and orchestration clarity rather than missing greenfield capability.
Why This Document Was Rewritten
An earlier version of this document mixed aspirational ownership with current implementation facts and then accumulated unresolved merge blocks. This version keeps the repo-accurate mapping and removes the conflicting alternative cell ownership model.
The important corrections are:
cell11is not the current audience-state owner.cell15is not the canonical attribution owner.cell20is not a dedicated LTV scoring service today.cell30is not the primary identity-resolution service.
Current Control Points In Repo
Boss orchestration
- Dynamic MCP discovery and routing:
src/moe/mcp_boss_agent.py - MCP service catalog:
mcp/service_registry.yaml
Ingestion and normalization
- Unified schema and source mappings:
src/cells/cell02/multi_source_config.py - Integration clients and DTOs:
src/cells/cell02/integrations - Broker adapters:
src/cells/cell02/integrations/broker.ts
Identity and journey assembly
- Identity resolution:
src/cells/cell02/enhanced_identity_resolver.py - Journey profiling:
src/cells/cell02/customer_journey_profiler.py - Journey stitching service:
src/cells/cell09/Cell9_Journey.py - Journey MCP service:
customer_journey_systeminmcp/service_registry.yaml
Knowledge graph and attribution
- KG build and semantic structuring:
services/cell03-intelligence/src/knowledge_graph - Attribution engine:
services/cell03-intelligence/src/attribution/engine.py - Boss marketing KG integration:
miz-oki-adk-agents/boss/marketing_kg_integration.py - Journey graph writer:
services/mizoki-journey-kg/graph-writer
Prediction and causal scoring
- Journey predictions API:
services/mizoki-journey-kg/predictions-api - Causal uplift estimation:
src/cells/cell26/Cell26.py
Semantic and provenance vocabulary
- Marketing semantic conventions:
config/semantic-conventions/marketing.yaml - UEES marketing schema fragments:
services/ekis/src/schema/uees.marketing.ts
Repo-Accurate Ownership Map
1. Raw extraction and first-pass normalization
Primary owner: cell02
Use:
src/cells/cell02/multi_source_config.pysrc/cells/cell02/integrations/clients/facebook.client.tssrc/cells/cell02/integrations/clients/googleads.client.tssrc/cells/cell02/integrations/clients/ga4.client.ts
Responsibilities:
- Pull or receive Meta, Google Ads, GA4, Shopify, Klaviyo, and related source data.
- Preserve original source payloads in raw fields.
- Normalize source records into
UnifiedSchema. - Emit broker-friendly fragments for downstream consumers.
2. Identity resolution
Primary owner: cell02
Use:
src/cells/cell02/enhanced_identity_resolver.py
Responsibilities:
- Resolve email, fbclid, gclid, GA4 client ID, Shopify, and Klaviyo identifiers.
- Produce stable
unified_customer_idvalues. - Carry forward confidence for deterministic versus weaker matches.
Current note:
- The existing implementation is deterministic-heavy and BigQuery-backed.
- There is no repo evidence that
cell30should own this path.
3. Journey profile construction
Primary owners: cell02 and journey services
Use:
src/cells/cell02/customer_journey_profiler.pysrc/cells/cell09/Cell9_Journey.pycustomer_journey_systemMCP service
Responsibilities:
- Build ordered touchpoint timelines per customer.
- Classify journey stages.
- Maintain ordered paths suitable for analytics and graph persistence.
Recommended split:
cell02remains the offline profile constructor.customer_journey_systemremains the runtime MCP entry point.cell09remains optional for stitching, anomaly detection, and golden-path analysis.
4. Knowledge graph persistence and semantic enrichment
Primary owners: services/cell03-intelligence, services/gemini-kg-pipeline, and graph_writer
Use:
services/cell03-intelligence/src/knowledge_graph/kg_builder.pyservices/cell03-intelligence/src/knowledge_graph/gemini_structurer.pyservices/gemini-kg-pipelineservices/mizoki-journey-kg/graph-writer
Responsibilities:
- Convert curated journey, customer, and campaign data into graph nodes and edges.
- Apply Gemini semantic enrichment where useful.
- Persist journey events and customer profiles into the graph store.
Recommended split:
cell03-intelligenceowns graph model construction from curated records.gemini-kg-pipelineowns LLM extraction and enrichment jobs.graph_writerowns direct event and profile graph persistence for runtime journey operations.
5. Attribution
Primary owners: services/cell03-intelligence and existing Boss KG tools
Use:
services/cell03-intelligence/src/attribution/engine.pymiz-oki-adk-agents/boss/ads_decision_integration.pymiz-oki-adk-agents/boss/marketing_kg_integration.py
Responsibilities:
- Compute first-touch, last-touch, linear, time-decay, and position-based attribution.
- Record touchpoints and conversion attribution in the temporal KG.
- Expose attributed journey context back to Boss and downstream services.
Current note:
- The repo already has attribution logic.
- What is missing is a single Boss-exposed MCP wrapper presenting attribution as one clean pipeline step.
6. Journey prediction and uplift
Primary owners: services/mizoki-journey-kg and cell26
Use:
services/mizoki-journey-kg/predictions-apisrc/cells/cell26/Cell26.py
Responsibilities:
- Predict next best action and completion probability from the journey graph.
- Estimate causal uplift on journey edges.
- Feed decision support back into Boss and optimization layers.
7. LTV and cohort scoring
Current state: no single service cleanly matches a dedicated score_customer_ltv runtime owner.
Recommendation:
- Treat LTV scoring as an explicit gap.
- Do not assign it to
cell20unless that contract is intentionally implemented.
Contract Strategy
Do not invent a third canonical event schema.
The repo already has two useful layers:
Layer A: curated internal record
Use UnifiedSchema from src/cells/cell02/multi_source_config.py as the curated record exchanged between Cell02-style ingestion and downstream graph processing.
Layer B: transport fragments
Use UEES-style marketing fragments from services/ekis/src/schema/uees.marketing.ts for broker or integration boundaries.
Cross-cutting metadata
Use config/semantic-conventions/marketing.yaml for:
marketing.*prov.*- activation and attribution context
Recommended wrapper
If a Boss handoff envelope is needed, keep it thin:
payload.unified_record->UnifiedSchemapayload.uees_fragment-> optional UEES transport fragmenttrace-> semantic marketing and provenance attributes
Boss MCP Mapping
The conceptual pipeline is:
- extract marketing signals
- normalize marketing events
- resolve customer identity
- build journey profile
- compute graph attribution
- score customer LTV
- upsert marketing knowledge graph
The repo already exposes or approximates much of this with existing runtime tools:
customer_journey_trackgemini_kg_ingest_eventgraph_writer_write_eventsgraph_writer_upsert_profilekg_record_touchpointkg_conversion_attributionjourney_kg_predict_next_stepjourney_kg_recommend
Gaps that still need an explicit MCP-safe wrapper:
- identity resolution
- attribution engine exposure
- LTV scoring
Recommended Implementation Sequence
Phase 1: align contracts
- Standardize Cell02 output around
UnifiedSchema. - Ensure raw payload retention.
- Ensure marketing and provenance attributes are attached consistently.
- Keep the Boss handoff shape thin.
Phase 2: expose identity and attribution as MCP tools
Target:
mcp/service_registry.yamlsrc/moe/mcp_boss_agent.py- Boss integration modules under
miz-oki-adk-agents/boss
Actions:
- Register identity resolution backed by
EnhancedIdentityResolver. - Register attribution backed by the Cell03 attribution engine or an adapter around it.
Phase 3: unify journey graph writes
- Decide whether runtime writes land in
graph_writerfirst and batch enrichment lands incell03-intelligenceandgemini-kg-pipeline, or vice versa. - Make that split explicit so events are not written twice through competing graph paths.
Phase 4: add missing LTV scoring layer
- Create a dedicated LTV or cohort scoring module.
- Only then register
score_customer_ltvin MCP.
Phase 5: add a Boss-level workflow wrapper
Create one higher-level workflow that sequences:
- identity resolution
- journey update
- graph write
- attribution
- prediction or recommendation
Runtime Implementation Status
There is now a concrete runtime surface for this integration:
Boss runtime files
miz-oki-adk-agents/boss/marketing_event_graph_mcp.pymiz-oki-adk-agents/boss/boss_agent_core.py
Registered tools include:
extract_marketing_signalsnormalize_marketing_eventsresolve_customer_identitybuild_journey_profilecompute_graph_attributionscore_customer_ltvupsert_marketing_knowledge_graphrun_marketing_event_graph_pipeline
Cloud Run service files
miz-oki-adk-agents/marketing_event_graph_service/main.pymiz-oki-adk-agents/marketing_event_graph_service/Dockerfilemiz-oki-adk-agents/marketing_event_graph_service/cloudbuild.yamldeployment/cloudrun_root/marketing-event-graph-service.prod.yamldeployment/cloudrun_root/marketing-event-graph-service.staging.yaml
HTTP endpoints
GET /healthPOST /api/v1/marketing-event-graph/extractPOST /api/v1/marketing-event-graph/normalizePOST /api/v1/marketing-event-graph/resolve-identityPOST /api/v1/marketing-event-graph/build-journeyPOST /api/v1/marketing-event-graph/attributionPOST /api/v1/marketing-event-graph/score-ltvPOST /api/v1/marketing-event-graph/upsert-kgPOST /api/v1/marketing-event-graph/run
Definition Of Done
This integration should be considered complete only when all of the following are true:
- Boss can invoke a real multi-step marketing journey pipeline through MCP.
- Identity resolution is exposed as a runtime tool, not only library code.
- Journey tracking and graph persistence use one explicit write path.
- Attribution is exposed as a first-class runtime capability.
- Provenance and marketing semantic attributes are carried consistently across ingestion, graph write, attribution, and prediction steps.
- LTV and cohort scoring has an explicit owner and contract.
Production Verification Commands
Run these only from an allowed network:
curl -sS https://boss-agent-adk-698171499447.us-central1.run.app/health
curl -sS https://miz-oki-cell2-698171499447.us-central1.run.app/health
curl -sS https://miz-oki-cell3-698171499447.us-central1.run.app/health
curl -sS https://miz-oki-cell11-698171499447.us-central1.run.app/health
curl -sS https://miz-oki-cell15-698171499447.us-central1.run.app/health
curl -sS https://miz-oki-cell20-698171499447.us-central1.run.app/health
curl -sS https://miz-oki-cell30-698171499447.us-central1.run.app/health
Expected:
- HTTP 200 from all services.
- Boss registry or list-tools includes the marketing-event-graph tool chain.
- End-to-end test payloads produce lineage and confidence fields at each step.
Practical Conclusion
The repo is closer to integration than the earlier aspirational mapping implied. The shortest path from current state to a clean production architecture is:
- standardize on existing contracts
- expose identity and attribution as MCP tools
- converge on one journey graph write path
- add a real scoring service for LTV and cohorts