# LeoCRM — Architektur-Plan für Phasen B–K > **Erstellt:** 2026-08-20 > **Basis:** Code-Analyse aller vorhandenen Systeme (kommunikation, graph_rag, unified_search, automation, workflows, ai, approval, storage, outbox, wiki) > **Leitlinie:** Alles aufbauend auf vorhandenem Code. Keine parallelen Systeme. > **Letzte Migration:** 0128_ai_decision_records.py → nächste: 0129+ --- ## Inhaltsverzeichnis 1. [Phase B Lücken (3 Tasks)](#phase-b-lücken) 2. [Phase F Lücken (3 Tasks)](#phase-f-lücken) 3. [Phase G Lücke (1 Task)](#phase-g-lücke) 4. [Phase H Rest (12 Tasks)](#phase-h-rest) 5. [Phase I (25 Tasks)](#phase-i) 6. [Phase J (10 Tasks)](#phase-j) 7. [Phase K (6 Tasks)](#phase-k) --- ## Vorhandene Systeme (Basis für alle Phasen) ### Kommunikation Plugin - **Pfad:** `app/plugins/builtins/kommunikation/` - **Models:** `CommConversation` (is_system, is_direct, is_archived, metadata_), `CommParticipant` (participant_type: user/ai/system/gateway), `CommMessage` (sender_type, content, content_format, metadata_, reply_to_id), `CommMessageBlock` (block_type, block_data JSONB, sort_order), `CommMessageAttachment`, `CommMessageReaction`, `CommMessageRead`, `CommConversationMute`, `CommConversationPin` - **Services:** `send_message()`, `get_conversation()`, `get_messages()`, `parse_mentions()`, `create_plugin_room()`, `post_system_message()`, `get_or_create_system_channel()` - **Contracts:** `KommunikationContract` registriert in `ContractRegistry` unter `"kommunikation"` - **MiniAppRegistry:** `register()`, `unregister()`, `unregister_plugin()`, `list_apps()`, `get_app()` — Singleton via `get_miniapp_registry()` - **ContentTypes:** `BLOCK_TYPES` dict: text, markdown, html, image, audio, video, file, action_card, contact_card, miniapp - **WebSocket:** `WebSocketManager` für Real-time - **ParticipantRegistry:** `get_participant_registry()` mit `ParticipantHandler` ### GraphRAG Plugin - **Pfad:** `app/plugins/builtins/graph_rag/` - **Models:** `EntityRelationship` (source_type, source_id, target_type, target_id, relationship_type, meta JSONB) - **Services:** `create_relationship()`, `traverse_graph()` (BFS, max_hops) - **Provider:** `GraphRAGSearchProvider` registriert in `SearchProviderRegistry` - **Routes:** `/api/v1/graph` - **Dependencies:** `unified_search` ### Unified Search Plugin - **Pfad:** `app/plugins/builtins/unified_search/` - **BaseSearchProvider:** `search_fts()`, `search_vector()` mit visibility filtering, `supports_fts/vector/rag/graph` flags - **SearchProviderRegistry:** `register()`, `unregister()`, `get()`, `get_all()` — Singleton via `get_search_registry()` - **Providers:** contact, company, task, tag, mail, ai_chat, etc. (14+) - **Embedding:** `llm_embed()` integration, `EMBEDDING_DIMENSIONS = 768` - **Contracts:** `get_search_registry()` in `unified_search.contracts` ### Automation Plugin - **Pfad:** `app/plugins/builtins/automation/` - **Models:** `AgentDefinition`, `AutomationDefinition` - **Agent Runner:** `app/plugins/builtins/automation/agent_runner.py` — ruft `run_react_loop`, `build_agent_context`, `resolve_agent_permissions`, `get_agent_tools`, `enforce_data_policy`, `mark_as_ai_generated`, `create_decision_record` - **Agent Loop:** `app/ai/agent_loop.py` — ReAct-Loop mit Tool-Execution, Approval-Pause - **Agent Comm:** `agent_comm.py` — `register_agent_comm_tool()` registriert Komm-Tool für Agents - **Prebuilt:** `prebuilt/email_triage_agent.py`, `prebuilt/follow_up_agent.py`, `prebuilt/contact_enrichment_agent.py`, `prebuilt/report_agent.py` — **NICHT in on_activate registriert** - **on_activate:** Registriert agent_comm_tool, agent_coordinator_tools, MiniApps, cron jobs — **aber KEINE prebuilt agents** ### Workflow Engine - **Pfad:** `app/workflows/engine.py` - **Models:** `Workflow` (steps JSONB, trigger_event), `WorkflowInstance` (status, current_step_index, context, resume_at, step_state, lock_owner, retry_count), `WorkflowStepHistory` - **Step Types:** action, approval, notification, condition, wait, http, mail, calendar, dms, search, agent, crm, event, webhook - **Decision Guard:** `app/workflows/decision_guard.py` — erstellt ApprovalRequest bei High-Risk - **Step Handlers:** `app/workflows/step_handlers.py` — `get_step_handler()`, `StepResult` ### Approval System - **Pfad:** `app/core/approval.py` - **Model:** `ApprovalRequest` (entity_type, entity_id, action, requested_by, requested_by_type, status: pending→approved/rejected/expired, request_metadata JSONB) - **Functions:** `create_approval_request()`, resolve functions - **Routes:** `app/routes/approvals.py` ### LLM Client - **Pfad:** `app/ai/llm_client.py` - **Functions:** `llm_complete()` (chat completion mit retry, cost tracking), `llm_embed()` (text embedding, 768 dims) - **Provider:** LiteLLM (100+ providers), OpenRouter für embeddings ### Outbox + Event Bus - **Outbox:** `app/core/outbox.py` — `enqueue_outbox_event()`, DLQ, replay, monitoring - **Event Bus:** `app/core/event_bus.py` — `get_event_bus()` - **Hooks:** `app/core/hooks.py` — `do_action()`, `apply_filters()`, `register_action()`, `register_filter()` ### Storage - **Pfad:** `app/core/storage.py` - **Classes:** `StorageBackend` (ABC), `LocalStorage`, `S3Storage` - **Methods:** `save()`, `save_stream()`, `read()`, `delete()` - **Config:** `STORAGE_BACKEND` env (local/s3), `STORAGE_PATH`, S3 settings ### Wiki Plugin - **Pfad:** `app/plugins/builtins/wiki/` - **Models:** WikiArticle, WikiCategory, WikiArticleVersion (versioning with restore) - **Routes:** `/api/v1/wiki` - **Plugin:** `WikiPlugin` — **kein on_activate** (kein search provider, keine AI tools) ### Notifications (zu deprecieren) - **Model:** `Notification` (app/models/notification.py) — noch vorhanden - **Routes:** `app/routes/notifications.py` — noch in main.py aktiv (line 545) - **Service:** `app/core/notifications.py` — `post_system_message()` delegiert an kommunikation, `create_notification()` ist deprecated wrapper - **Frontend:** `NotificationDropdown.tsx`, `NotificationItem.tsx` — noch vorhanden --- ## Phase B Lücken ### B-VEC-IVF: IVFFlat Index-Strategie implementieren **Basis:** `app/core/storage.py` (nein), `app/plugins/builtins/unified_search/embedding.py` + pgvector **Was existiert:** - HNSW ist konfiguriert (`settings.hnsw_ef_search` in `base_provider.py`) - `test_vector_performance.py` existiert - Keine IVFFlat-Konfiguration im Code **Was neu gebaut wird:** 1. **Config-Erweiterung:** `app/config.py` ```python # Neue Settings vector_index_type: str = "hnsw" # "hnsw" or "ivfflat" ivfflat_lists: int = 100 # number of lists for IVFFlat ivfflat_probes: int = 10 # number of probes at query time ``` 2. **Index-Manager:** `app/plugins/builtins/unified_search/index_manager.py` (NEU, ~200 Zeilen) ```python async def create_ivfflat_index(db, table: str, column: str, lists: int): """Create IVFFlat index on vector column.""" await db.execute(text( f"CREATE INDEX IF NOT EXISTS idx_{table}_{column}_ivfflat " f"ON {table} USING ivfflat ({column} vector_cosine_ops) WITH (lists = {lists})" )) async def set_ivfflat_probes(db, probes: int): await db.execute(text(f"SET LOCAL ivfflat.probes = {probes}")) async def switch_index_strategy(db, table: str, column: str, strategy: str): """Switch between HNSW and IVFFlat.""" # Drop old, create new ``` 3. **Migration:** `alembic/versions/0129_ivfflat_index_strategy.py` - Fügt `vector_index_type` zu system_settings hinzu - Erstellt IVFFlat-Index alternativ zu HNSW (nicht beide gleichzeitig) - Downgrade: Drop IVFFlat, restore HNSW 4. **BaseSearchProvider Anpassung:** `app/plugins/builtins/unified_search/base_provider.py` - In `search_vector()`: Wenn `settings.vector_index_type == "ivfflat"`, set `ivfflat.probes` statt `hnsw.ef_search` 5. **Admin API:** `app/plugins/builtins/unified_search/routes.py` - `POST /api/v1/search/index/switch` — Switch index strategy (admin only) - `GET /api/v1/search/index/status` — Current index info **Verbindungen:** - `base_provider.py` importiert `index_manager` für probe/ef_search setting - `config.py` erweitert mit vector_index_type **Migrationen:** - `0129_ivfflat_index_strategy.py` **Tests:** - `tests/test_ivfflat_index.py` — IVFFlat index creation, query with probes, performance comparison - `tests/test_vector_performance.py` — erweitert um IVFFlat benchmarks **Frontend:** - Settings-Seite: Toggle HNSW ↔ IVFFlat in Admin-Settings --- ### B-STOR-EXT: External Storage Plugin System (WebDAV, Nextcloud) **Basis:** `app/core/storage.py` (StorageBackend ABC, LocalStorage, S3Storage) **Was existiert:** - `StorageBackend` ABC mit `save()`, `save_stream()`, `read()`, `delete()` - `LocalStorage`, `S3Storage` implementiert - Kein Plugin-Interface für externe Storage-Provider **Was neu gebaut wird:** 1. **Storage Provider Registry:** `app/core/storage_registry.py` (NEU, ~150 Zeilen) ```python class StorageProviderRegistry: """Registry for pluggable storage backends.""" def register(self, name: str, backend_class: type[StorageBackend]) -> None: ... def unregister(self, name: str) -> None: ... def get(self, name: str) -> type[StorageBackend] | None: ... def list_providers(self) -> list[str]: ... _registry: StorageProviderRegistry | None = None def get_storage_registry() -> StorageProviderRegistry: ... ``` 2. **WebDAV Storage Backend:** `app/plugins/builtins/storage_webdav/` (NEU, komplettes Plugin, ~400 Zeilen) - `plugin.py` — `WebDAVStoragePlugin(BasePlugin)` mit Manifest - `backend.py` — `WebDAVStorage(StorageBackend)` implementiert save/read/delete via HTTP (PUT/GET/DELETE) - `routes.py` — `POST /api/v1/storage/webdav/test` — Test connection - `schemas.py` — WebDAVConfig (url, username, password, base_path) - `migrations/0001_initial.sql` — storage_provider_configs table 3. **Nextcloud Storage Backend:** `app/plugins/builtins/storage_nextcloud/` (NEU, ~300 Zeilen) - Baut auf WebDAV auf (Nextcloud hat WebDAV-API) - `plugin.py` — `NextcloudStoragePlugin(BasePlugin)` - `backend.py` — `NextcloudStorage(WebDAVStorage)` mit Nextcloud-spezifischen Erweiterungen (sharing, OCS API) - `routes.py` — `POST /api/v1/storage/nextcloud/test`, `GET /api/v1/storage/nextcloud/shares` 4. **Storage Config Model:** `app/models/storage_config.py` (NEU) ```python class StorageProviderConfig(Base, TenantMixin): __tablename__ = "storage_provider_configs" id: UUID PK provider_name: str # "local", "s3", "webdav", "nextcloud" config: dict[str, Any] = JSONB # provider-specific config (encrypted secrets) is_active: bool = True priority: int = 0 # for fallback ordering created_at: TIMESTAMPTZ ``` 5. **Storage Factory:** `app/core/storage.py` erweitern ```python async def get_storage_backend(db, tenant_id) -> StorageBackend: """Get configured storage backend for tenant.""" # Query StorageProviderConfig, instantiate via registry ``` 6. **Admin API:** `app/routes/storage.py` (NEU) - `GET /api/v1/storage/providers` — List available providers - `POST /api/v1/storage/providers` — Configure provider for tenant - `PUT /api/v1/storage/providers/{id}` — Update config - `DELETE /api/v1/storage/providers/{id}` — Remove provider - `POST /api/v1/storage/providers/{id}/test` — Test connection **Verbindungen:** - `storage.py` importiert `storage_registry` für dynamische Provider-Auflösung - DMS/Mail Plugins nutzen `get_storage_backend()` statt direkter LocalStorage/S3Storage - WebDAV/Nextcloud Plugins registrieren sich in `StorageProviderRegistry` via `on_activate` **Migrationen:** - `0130_storage_provider_configs.py` — storage_provider_configs table - Plugin-Migrations: `storage_webdav/migrations/0001_initial.sql`, `storage_nextcloud/migrations/0001_initial.sql` **Tests:** - `tests/test_storage_webdav.py` — WebDAV save/read/delete mit Mock-HTTP-Server - `tests/test_storage_nextcloud.py` — Nextcloud-spezifische Tests - `tests/test_storage_registry.py` — Registry register/unregister/get - `tests/test_storage_provider_config.py` — CRUD API tests **Frontend:** - `frontend/src/pages/StorageSettings.tsx` — Provider-Konfiguration - `frontend/src/components/storage/ProviderConfigForm.tsx` — Config form per provider type - `frontend/src/components/storage/ConnectionTest.tsx` — Test connection button --- ### B-NOTIF-DEPREC: Notification System zurückbauen **Basis:** `app/core/notifications.py`, `app/models/notification.py`, `app/routes/notifications.py` **Was existiert:** - `Notification` Model in `app/models/notification.py` - `NotificationType` Model - `/api/v1/notifications` Routes in `app/routes/notifications.py` (aktiv in main.py:545) - `NotificationDropdown.tsx`, `NotificationItem.tsx` im Frontend - `post_system_message()` in `app/core/notifications.py` delegiert bereits an kommunikation - `create_notification()` ist deprecated wrapper - `NotificationPreference` Model existiert (wird behalten für preferences) **Was gemacht wird:** 1. **Routes stilllegen:** `app/routes/notifications.py` - Alle Endpoints markieren als `@deprecated` in OpenAPI - `GET /api/v1/notifications` → redirect zu `GET /api/v1/comm/system-channel/messages` - `POST /api/v1/notifications/read` → redirect zu comm equivalent - `GET /api/v1/notifications/unread-count` → redirect zu comm equivalent - Nach 1 Release: Routes aus main.py entfernen (line 545) 2. **Model als deprecated markieren:** `app/models/notification.py` - `Notification` und `NotificationType` erhalten Docstring `DEPRECATED — use CommConversation is_system=True` - Keine neuen Writes, nur noch Reads für Migration 3. **Frontend entfernen:** - `NotificationDropdown.tsx` → ersetzen durch `CommSystemChannel.tsx` (neu, liest aus comm system channel) - `NotificationItem.tsx` → ersetzen durch `CommMessageItem.tsx` (existiert bereits in kommunikation Frontend) - Header-Komponente anpassen: Notification-Bell → Comm-Message-Bell 4. **Daten-Migration:** `alembic/versions/0131_migrate_notifications_to_comm.py` - Liest alle `Notification` records - Erstellt entsprechende `CommMessage` im system channel - Markiert Notifications als migrated (neue Spalte `migrated_at`) - Nach erfolgreichem Migration-Run: Drop `notifications` und `notification_types` tables 5. **Service cleanup:** `app/core/notifications.py` - `create_notification()` entfernen (deprecated wrapper) - `post_system_message()` behalten (delegiert an kommunikation) - Direkte Notification-Queries entfernen 6. **main.py cleanup:** - `notifications` Router entfernen (line 545) - `notifications` aus taginfo entfernen (line 432) - `notifications` aus imports entfernen (line 57) **Verbindungen:** - Alle ehemaligen Notification-Consumer nutzen `post_system_message()` aus `app/core/notifications.py` - Frontend nutzt Comm-System-Channel API **Migrationen:** - `0131_migrate_notifications_to_comm.py` — Daten-Migration + Drop tables **Tests:** - `tests/test_notification_deprecation.py` — Verify deprecated routes redirect/404 - `tests/test_notification_migration.py` — Verify data migration correctness - Existierende `tests/test_notifications.py` anpassen (nur noch post_system_message testen) **Frontend:** - `frontend/src/components/comm/SystemChannelBell.tsx` (NEU) — ersetzt NotificationDropdown - Header-Komponente aktualisieren --- ## Phase F Lücken ### F-PREBUILT: Pre-built Agents in plugin.py on_activate registrieren **Basis:** `app/plugins/builtins/automation/plugin.py` on_activate, `app/plugins/builtins/automation/prebuilt/` **Was existiert:** - 4 Pre-built Agent Definitionen: `email_triage_agent.py`, `follow_up_agent.py`, `contact_enrichment_agent.py`, `report_agent.py` - Jede Datei hat eine `create_*_agent()` Funktion - `on_activate` in `plugin.py` registriert agent_comm_tool, agent_coordinator_tools, MiniApps, cron jobs — **aber NICHT die prebuilt agents** **Was neu gebaut wird:** 1. **Prebuilt Agent Registration:** `app/plugins/builtins/automation/plugin.py` on_activate erweitern ```python async def on_activate(self, db, service_container, event_bus) -> None: await super().on_activate(db, service_container, event_bus) # ... existing registrations ... # Register pre-built agents try: from app.plugins.builtins.automation.prebuilt.email_triage_agent import create_email_triage_agent from app.plugins.builtins.automation.prebuilt.follow_up_agent import create_follow_up_agent from app.plugins.builtins.automation.prebuilt.contact_enrichment_agent import create_contact_enrichment_agent from app.plugins.builtins.automation.prebuilt.report_agent import create_report_agent for create_fn in [create_email_triage_agent, create_follow_up_agent, create_contact_enrichment_agent, create_report_agent]: agent = await create_fn(db) logger.info(f"Pre-built agent registered: {agent.name}") except Exception: logger.exception("Failed to register pre-built agents") ``` 2. **Idempotency:** Jede `create_*_agent()` Funktion prüft ob Agent bereits existiert (by name + tenant) - Wenn ja: skip (kein Duplikat) - Wenn nein: erstelle AgentDefinition 3. **on_deactivate cleanup:** Pre-built agents werden bei Deaktivierung entfernt (oder als system-markiert belassen) **Verbindungen:** - `plugin.py on_activate` → `prebuilt/*.create_*_agent()` - AgentDefinition in DB → verfügbar in Agent Dashboard und Agent Runner **Migrationen:** - Keine neue Migration (AgentDefinition Tabelle existiert bereits) **Tests:** - `tests/test_prebuilt_agent_registration.py` — Verify 4 agents exist after on_activate - `tests/test_prebuilt_agent_idempotency.py` — Verify no duplicates on re-activate **Frontend:** - Keine Änderung — Agents erscheinen automatisch im Agent Dashboard --- ### F-AGENT-COMM: Agent Run-Results in Kommunikationszentrale posten **Basis:** `app/plugins/builtins/automation/agent_runner.py`, `app/plugins/builtins/kommunikation/contracts.py` **Was existiert:** - `agent_comm.py` registriert ein `send_message` Tool für Agents (postet in comm conversations) - `agent_runner.py` führt Agent Loop aus und sammelt Results - KommunikationContract: `send_message()`, `post_system_message()`, `create_plugin_room()` - **Aber:** Agent Run-Results (Zusammenfassung, Steps, Cost) werden NICHT automatisch in Kommunikationszentrale gepostet **Was neu gebaut wird:** 1. **Agent Result Poster:** `app/plugins/builtins/automation/agent_result_poster.py` (NEU, ~200 Zeilen) ```python async def post_agent_run_result( db: AsyncSession, tenant_id: uuid.UUID, agent_run_id: uuid.UUID, agent_name: str, result: AgentRunResult, user_id: uuid.UUID, ) -> None: """Post agent run result to kommunikation system channel.""" from app.plugins.builtins.contracts import get_contract komm = get_contract("kommunikation") if not komm: return # Create rich message with blocks blocks = [ {"block_type": "markdown", "block_data": {"markdown": f"## Agent: {agent_name}\n\n{result.final_content}"}}, {"block_type": "action_card", "block_data": { "title": "Agent Run Summary", "body": f"Steps: {result.steps_taken} | Cost: ${result.total_cost_usd:.4f} | Status: {result.status}", "actions": [ {"label": "View Details", "action": f"/agents/runs/{agent_run_id}", "type": "link"}, ], }}, ] await komm.send_message( db=db, tenant_id=tenant_id, conversation_id=system_channel_id, # get_or_create_system_channel sender_id=agent_run_id, sender_type="ai", content=result.final_content or "Agent run completed", blocks=blocks, metadata={"agent_run_id": str(agent_run_id), "agent_name": agent_name}, ) ``` 2. **Agent Runner Integration:** `app/plugins/builtins/automation/agent_runner.py` - Nach `run_react_loop()` completion: rufe `post_agent_run_result()` auf - Bei Error: poste error summary in system channel - Bei Approval-Pause: poste approval request in system channel ```python # In agent_runner.py nach run_react_loop: from app.plugins.builtins.automation.agent_result_poster import post_agent_run_result await post_agent_run_result(db, tenant_id, agent_run_id, agent_name, result, user_id) ``` 3. **CommMessageBlock Types erweitern:** `app/plugins/builtins/kommunikation/content_types.py` - Neuer block_type: `"agent_result"` mit fields: `agent_run_id`, `agent_name`, `status`, `steps`, `cost` - Neuer block_type: `"approval_request"` mit fields: `approval_id`, `action`, `entity_type`, `entity_id` **Verbindungen:** - `agent_runner.py` → `agent_result_poster.py` → `kommunikation.contracts.send_message()` - `content_types.py` erweitert um agent_result und approval_request blocks **Migrationen:** - Keine (CommMessageBlock nutzt JSONB, schema-flexibel) **Tests:** - `tests/test_agent_result_posting.py` — Verify agent run results appear in system channel - `tests/test_agent_result_blocks.py` — Verify block structure and types **Frontend:** - `frontend/src/components/comm/AgentResultBlock.tsx` (NEU) — Rendert agent_result block_type - `frontend/src/components/comm/ApprovalRequestBlock.tsx` (NEU) — Rendert approval_request block_type - Block-Renderer in ChatView erweitern --- ### F-WORK: Agent Workstream auf kommunikation Plugin aufbauen **Basis:** `app/plugins/builtins/kommunikation/` (CommConversation, CommMessage, CommMessageBlock, MiniAppRegistry) **Was existiert:** - KommunikationPlugin mit CommConversation (is_system, is_direct, metadata_) - CommMessageBlock (block_type, block_data JSONB) — unterstützt action_card, contact_card, miniapp - MiniAppRegistry mit register/unregister - `create_plugin_room()` für plugin-spezifische Conversations - Agent Comm Tool (send_message) bereits registriert **Was neu gebaut wird:** 1. **Agent Workstream Service:** `app/plugins/builtins/automation/agent_workstream.py` (NEU, ~300 Zeilen) ```python async def create_agent_workstream_room( db: AsyncSession, tenant_id: uuid.UUID, agent_id: uuid.UUID, user_id: uuid.UUID ) -> dict: """Create a dedicated workstream conversation for an agent.""" from app.plugins.builtins.contracts import get_contract komm = get_contract("kommunikation") # Create plugin room with agent as participant room = await komm.create_plugin_room( db=db, tenant_id=tenant_id, plugin_name="automation", room_key=f"agent_{agent_id}", title=f"Agent Workstream", metadata={"agent_id": str(agent_id), "type": "agent_workstream"}, ) # Add agent and user as participants return room async def post_workstream_update( db, tenant_id, conversation_id, sender_type, sender_id, content, blocks=None ) -> dict: """Post a workstream update with rich blocks.""" from app.plugins.builtins.contracts import get_contract komm = get_contract("kommunikation") return await komm.send_message( db=db, tenant_id=tenant_id, conversation_id=conversation_id, sender_id=sender_id, sender_type=sender_type, content=content, blocks=blocks or [], ) async def get_agent_workstream(db, tenant_id, agent_id, user_id) -> dict | None: """Get or create workstream room for agent.""" # Lookup by metadata.agent_id, create if not exists ``` 2. **Workstream Block Types:** `app/plugins/builtins/kommunikation/content_types.py` erweitern - `"task_card"` — fields: task_id, title, status, assignee, due_date - `"workflow_card"` — fields: workflow_id, instance_id, status, current_step - `"knowledge_card"` — fields: entity_type, entity_id, title, source - `"progress_card"` — fields: current, total, label, percentage 3. **Agent Runner Integration:** `agent_runner.py` - Bei Agent-Start: `create_agent_workstream_room()` → post "Agent started" message - Bei jedem Step: `post_workstream_update()` mit progress_card block - Bei Completion: post result summary (wie F-AGENT-COMM) - Bei Approval: post approval_request block 4. **Workstream API Routes:** `app/plugins/builtins/automation/agent_routes.py` erweitern - `GET /api/v1/agents/{id}/workstream` — Get workstream conversation for agent - `POST /api/v1/agents/{id}/workstream/message` — Post message to workstream - `GET /api/v1/agents/{id}/workstream/messages` — List workstream messages **Verbindungen:** - `agent_runner.py` → `agent_workstream.py` → `kommunikation.contracts.create_plugin_room()` + `send_message()` - `content_types.py` erweitert um task_card, workflow_card, knowledge_card, progress_card - `agent_routes.py` erweitert um workstream endpoints **Migrationen:** - Keine (nutzt vorhandene comm_conversations, comm_messages, comm_message_blocks) **Tests:** - `tests/test_agent_workstream.py` — Workstream room creation, message posting, block rendering - `tests/test_agent_workstream_integration.py` — Agent run → workstream messages appear **Frontend:** - `frontend/src/components/comm/TaskCardBlock.tsx` (NEU) - `frontend/src/components/comm/WorkflowCardBlock.tsx` (NEU) - `frontend/src/components/comm/KnowledgeCardBlock.tsx` (NEU) - `frontend/src/components/comm/ProgressCardBlock.tsx` (NEU) - `frontend/src/components/comm/WorkstreamBlockRenderer.tsx` (NEU) — Dispatch block_type → component - AgentChat-Seite erweitert um Workstream-View --- ## Phase G Lücke ### G-WORK: Workflow Workstream auf kommunikation Plugin aufbauen **Basis:** `app/workflows/engine.py`, `app/plugins/builtins/kommunikation/` **Was existiert:** - WorkflowEngine verarbeitet Steps (action, approval, condition, wait, etc.) - WorkflowInstance hat status, current_step_index, context - `post_system_message()` in `app/core/notifications.py` delegiert an kommunikation - engine.py importiert bereits `post_system_message` aus `app.core.notifications` **Was neu gebaut wird:** 1. **Workflow Workstream Service:** `app/workflows/workstream.py` (NEU, ~250 Zeilen) ```python async def create_workflow_workstream_room( db: AsyncSession, tenant_id: uuid.UUID, workflow_instance_id: uuid.UUID, user_id: uuid.UUID ) -> dict: """Create workstream conversation for a workflow instance.""" from app.plugins.builtins.contracts import get_contract komm = get_contract("kommunikation") room = await komm.create_plugin_room( db=db, tenant_id=tenant_id, plugin_name="workflows", room_key=f"wf_{workflow_instance_id}", title=f"Workflow: {workflow_name}", metadata={"workflow_instance_id": str(workflow_instance_id), "type": "workflow_workstream"}, ) return room async def post_workflow_step_update( db, tenant_id, conversation_id, step_index, step_type, status, result=None ) -> dict: """Post workflow step progress to workstream.""" blocks = [ {"block_type": "progress_card", "block_data": { "current": step_index + 1, "total": total_steps, "label": f"Step {step_index + 1}: {step_type}", "percentage": int((step_index + 1) / total_steps * 100), }}, {"block_type": "workflow_card", "block_data": { "workflow_id": str(workflow_id), "instance_id": str(instance_id), "status": status, "current_step": step_index, }}, ] await komm.send_message(db, tenant_id, conversation_id, ...) async def post_workflow_approval_request( db, tenant_id, conversation_id, approval_id, step_description ) -> dict: """Post approval request to workstream.""" blocks = [{"block_type": "approval_request", "block_data": { "approval_id": str(approval_id), "action": "approve_workflow_step", "entity_type": "workflow_instance", "entity_id": str(instance_id), }}] await komm.send_message(...) async def post_workflow_completion(db, tenant_id, conversation_id, status, summary): """Post workflow completion summary.""" ``` 2. **Workflow Engine Integration:** `app/workflows/engine.py` erweitern - Bei Instance-Start: `create_workflow_workstream_room()` → post "Workflow started" message - Bei jedem Step-Übergang: `post_workflow_step_update()` - Bei Approval-Step: `post_workflow_approval_request()` - Bei Completion: `post_workflow_completion()` ```python # In engine.py process_step(): from app.workflows.workstream import post_workflow_step_update await post_workflow_step_update(db, tenant_id, conversation_id, step_index, step_type, status, result) ``` 3. **Workflow Routes erweitern:** `app/routes/workflows.py` - `GET /api/v1/workflows/instances/{id}/workstream` — Get workstream conversation - `POST /api/v1/workflows/instances/{id}/workstream/message` — Post message **Verbindungen:** - `engine.py` → `workstream.py` → `kommunikation.contracts.create_plugin_room()` + `send_message()` - `workstream.py` nutzt workflow_card, progress_card, approval_request block types (aus F-WORK definiert) **Migrationen:** - Keine (nutzt vorhandene comm Tabellen) **Tests:** - `tests/test_workflow_workstream.py` — Workstream room creation, step updates, approval posts - `tests/test_workflow_workstream_integration.py` — Full workflow run → workstream messages **Frontend:** - `frontend/src/pages/WorkflowDetail.tsx` erweitert um Workstream-Tab - Nutzt WorkstreamBlockRenderer aus F-WORK --- ## Phase H Rest ### H-SRC: Knowledge Source Adapter — auf unified_search providers aufbauen **Basis:** `app/plugins/builtins/unified_search/` (BaseSearchProvider, SearchProviderRegistry) **Was existiert:** - 14+ Search Providers (contact, company, task, mail, etc.) - BaseSearchProvider mit search_fts, search_vector, get_embedding_text - SearchProviderRegistry mit register/unregister/get_all **Was neu gebaut wird:** 1. **Knowledge Source Adapter:** `app/plugins/builtins/knowledge/source_adapter.py` (NEU, ~300 Zeilen) ```python class KnowledgeSourceAdapter: """Adapts unified_search providers as knowledge sources for extraction.""" async def fetch_source_content( self, db: AsyncSession, tenant_id: uuid.UUID, entity_type: str, entity_id: uuid.UUID, ) -> dict[str, Any] | None: """Fetch full content for an entity via its search provider.""" from app.plugins.builtins.unified_search.contracts import get_search_registry registry = get_search_registry() provider = registry.get(entity_type) if not provider: return None embedding_text = await provider.get_embedding_text(db, entity_id, tenant_id) return {"entity_type": entity_type, "entity_id": str(entity_id), "content": embedding_text} async def batch_fetch( self, db, tenant_id, items: list[dict], ) -> list[dict]: """Batch fetch content for multiple entities.""" async def discover_sources( self, db, tenant_id, since: datetime | None = None, ) -> list[dict]: """Discover all entities that could be knowledge sources.""" # Query all providers for recently updated entities ``` 2. **Knowledge Source Model:** `app/plugins/builtins/knowledge/models.py` (NEU) ```python class KnowledgeSource(Base, TenantMixin, OwnedMixin): __tablename__ = "knowledge_sources" id: UUID PK entity_type: str # "contact", "company", "mail", "file", "wiki_article" entity_id: UUID content_hash: str # SHA256 of content for change detection last_extracted_at: TIMESTAMPTZ | None extraction_status: str # "pending", "extracted", "failed", "stale" metadata_: dict = JSONB ``` 3. **Knowledge Plugin:** `app/plugins/builtins/knowledge/` (NEU, komplettes Plugin) - `plugin.py` — `KnowledgePlugin(BasePlugin)` mit Manifest, dependencies: ["unified_search", "graph_rag"] - `routes.py` — `/api/v1/knowledge/sources`, `/api/v1/knowledge/extraction` - `services.py` — Extraction orchestration - `schemas.py` — Pydantic schemas **Verbindungen:** - `source_adapter.py` → `unified_search.contracts.get_search_registry()` → provider.get_embedding_text() - Knowledge Plugin dependencies: unified_search, graph_rag **Migrationen:** - `0132_knowledge_sources.py` — knowledge_sources table - Plugin-Migration: `knowledge/migrations/0001_initial.sql` **Tests:** - `tests/test_knowledge_source_adapter.py` — Fetch content via providers, batch fetch, discover - `tests/test_knowledge_sources.py` — CRUD API, change detection **Frontend:** - `frontend/src/pages/KnowledgeDashboard.tsx` (NEU) — Source overview, extraction status --- ### H-LLM-REL: LLM Relationship Extraction — auf graph_rag aufbauen **Basis:** `app/plugins/builtins/graph_rag/services.py` (create_relationship), `app/ai/llm_client.py` (llm_complete) **Was neu gebaut wird:** 1. **Relationship Extractor:** `app/plugins/builtins/knowledge/relationship_extractor.py` (NEU, ~250 Zeilen) ```python async def extract_relationships( db: AsyncSession, tenant_id: uuid.UUID, entity_type: str, entity_id: uuid.UUID, content: str, ) -> list[dict]: """Use LLM to extract relationships from entity content.""" from app.ai.llm_client import llm_complete prompt = f"""Analyze the following content and extract relationships. Return JSON array of {{source_type, source_id, target_type, target_id, relationship_type, confidence}}. Content: {content[:4000]} Entity: {entity_type} {entity_id} """ response = await llm_complete(messages=[{"role": "user", "content": prompt}], ...) relationships = parse_llm_relationships(response, entity_type, entity_id) # Create relationships via graph_rag from app.plugins.builtins.graph_rag.services import create_relationship for rel in relationships: await create_relationship(db, tenant_id, **rel) return relationships ``` 2. **LLM Prompt Template:** Strukturiertes Prompt für Relationship-Extraction - Input: Entity content + context - Output: JSON array of relationships with confidence scores - System prompt: Domain-specific relationship types (works_for, has_email, related_to, etc.) **Verbindungen:** - `relationship_extractor.py` → `llm_client.llm_complete()` → `graph_rag.services.create_relationship()` **Migrationen:** - Keine (EntityRelationship existiert bereits) **Tests:** - `tests/test_relationship_extraction.py` — LLM mock, verify relationships created in graph_rag **Frontend:** - Knowledge Dashboard: Extracted relationships view --- ### H-ENT: Entity Extraction — auf graph_rag aufbauen **Basis:** `app/plugins/builtins/graph_rag/models.py`, `app/ai/llm_client.py` **Was neu gebaut wird:** 1. **Entity Extractor:** `app/plugins/builtins/knowledge/entity_extractor.py` (NEU, ~250 Zeilen) ```python async def extract_entities( db: AsyncSession, tenant_id: uuid.UUID, content: str, source_type: str, source_id: uuid.UUID, ) -> list[dict]: """Use LLM to extract named entities from content.""" from app.ai.llm_client import llm_complete prompt = f"""Extract named entities from the following content. Return JSON array of {{entity_type, name, attributes, mentions: [{{start, end}}]}}. Entity types: person, organization, email, phone, date, location, project Content: {content[:4000]} """ response = await llm_complete(...) entities = parse_llm_entities(response) # Link entities to existing CRM records or create new ones for entity in entities: matched = await match_entity_to_crm(db, tenant_id, entity) if matched: # Create relationship: source → matched entity await create_relationship(db, tenant_id, source_type, source_id, matched.type, matched.id, "mentions") return entities ``` 2. **Entity Matching Service:** `app/plugins/builtins/knowledge/entity_matcher.py` (NEU, ~150 Zeilen) - Match extracted entities against existing contacts, companies, etc. - Fuzzy matching by name, email, phone - Returns matched entity or None **Verbindungen:** - `entity_extractor.py` → `llm_client.llm_complete()` → `graph_rag.services.create_relationship()` - `entity_matcher.py` → Contact/Company models for matching **Tests:** - `tests/test_entity_extraction.py` — LLM mock, verify entity extraction and matching **Frontend:** - Knowledge Dashboard: Extracted entities view --- ### H-AUTO: Auto-Relationship Creation — auf graph_rag aufbauen **Basis:** `app/plugins/builtins/graph_rag/services.py`, `app/core/hooks.py`, `app/core/outbox.py` **Was neu gebaut wird:** 1. **Auto-Relationship Engine:** `app/plugins/builtins/knowledge/auto_relationship.py` (NEU, ~200 Zeilen) ```python async def auto_create_relationships( db: AsyncSession, tenant_id: uuid.UUID, entity_type: str, entity_id: uuid.UUID, ) -> list[dict]: """Automatically create relationships based on entity data.""" # Rule-based: contact → company (works_for), mail → contact (sent_by), etc. rules = get_relationship_rules(entity_type) relationships = [] for rule in rules: targets = await rule.find_targets(db, tenant_id, entity_type, entity_id) for target in targets: result = await create_relationship(db, tenant_id, entity_type, entity_id, target.type, target.id, rule.relationship_type) if "error" not in result: relationships.append(result) return relationships ``` 2. **Hook Integration:** `app/plugins/builtins/knowledge/plugin.py` on_activate - Register hooks: `contact.after_create`, `contact.after_update`, `mail.received`, `company.after_create` - On hook fire: `auto_create_relationships()` 3. **Rule Definitions:** `app/plugins/builtins/knowledge/rules.py` (NEU, ~200 Zeilen) - Contact → Company: if contact.company_id exists, create "works_for" relationship - Mail → Contact: if mail.from_address matches contact.email, create "sent_by" relationship - Task → Contact: if task.assigned_to matches contact, create "assigned_to" relationship - File → Contact: if file.metadata has contact reference, create "belongs_to" relationship **Verbindungen:** - `auto_relationship.py` → `graph_rag.services.create_relationship()` - `plugin.py on_activate` → `hooks.register_action()` - Hook callbacks → `auto_create_relationships()` **Tests:** - `tests/test_auto_relationships.py` — Create contact with company → verify "works_for" relationship - `tests/test_relationship_rules.py` — Each rule tested independently **Frontend:** - Knowledge Dashboard: Auto-created relationships view --- ### H-CONF: Confidence Score + Review Queue — auf graph_rag aufbauen **Basis:** `app/plugins/builtins/graph_rag/models.py` (EntityRelationship) **Was neu gebaut wird:** 1. **EntityRelationship erweitern:** Neue Migration fügt confidence und review_status hinzu ```python # In EntityRelationship (via migration): confidence: float = 0.0 # 0.0-1.0 review_status: str = "auto" # "auto", "pending_review", "approved", "rejected" reviewed_by: UUID | None reviewed_at: TIMESTAMPTZ | None extraction_method: str = "auto" # "auto", "llm", "manual" ``` 2. **Review Queue Service:** `app/plugins/builtins/knowledge/review_queue.py` (NEU, ~200 Zeilen) ```python async def get_pending_reviews(db, tenant_id, limit=50) -> list[EntityRelationship]: ... async def approve_relationship(db, tenant_id, rel_id, user_id) -> dict: ... async def reject_relationship(db, tenant_id, rel_id, user_id, reason) -> dict: ... async def batch_approve(db, tenant_id, rel_ids, user_id) -> dict: ... ``` 3. **Review API:** `app/plugins/builtins/knowledge/routes.py` erweitern - `GET /api/v1/knowledge/review-queue` — List pending relationships - `POST /api/v1/knowledge/review/{id}/approve` — Approve - `POST /api/v1/knowledge/review/{id}/reject` — Reject - `POST /api/v1/knowledge/review/batch-approve` — Batch approve **Verbindungen:** - EntityRelationship erweitert um confidence/review_status - Review Queue nutzt graph_rag models **Migrationen:** - `0133_entity_relationship_confidence.py` — Add confidence, review_status, reviewed_by, reviewed_at, extraction_method **Tests:** - `tests/test_review_queue.py` — Pending reviews, approve, reject, batch approve - `tests/test_confidence_scoring.py` — Confidence threshold filtering **Frontend:** - `frontend/src/pages/KnowledgeReviewQueue.tsx` (NEU) — Review queue UI - `frontend/src/components/knowledge/RelationshipReviewCard.tsx` (NEU) --- ### H-EVT: Event-Driven Extraction — auf Event Bus + graph_rag aufbauen **Basis:** `app/core/event_bus.py`, `app/core/outbox.py`, `app/core/hooks.py` **Was neu gebaut wird:** 1. **Event-Driven Extraction Handler:** `app/plugins/builtins/knowledge/event_handler.py` (NEU, ~200 Zeilen) ```python async def handle_entity_created(event_name: str, payload: dict): """Handle entity.created event — trigger extraction.""" entity_type = payload.get("entity_type") entity_id = uuid.UUID(payload.get("entity_id")) tenant_id = uuid.UUID(payload.get("tenant_id")) # Enqueue extraction job from app.core.jobs import enqueue_job await enqueue_job("knowledge_extract", { "entity_type": entity_type, "entity_id": str(entity_id), "tenant_id": str(tenant_id), }) async def handle_entity_updated(event_name: str, payload: dict): """Handle entity.updated — mark source as stale, re-extract.""" ``` 2. **Extraction Job:** `app/plugins/builtins/knowledge/jobs.py` (NEU, ~150 Zeilen) - ARQ job function: `async def knowledge_extract(ctx, entity_type, entity_id, tenant_id)` - Calls source_adapter → entity_extractor → relationship_extractor → auto_relationship - Updates KnowledgeSource.extraction_status 3. **Event Registration:** `app/plugins/builtins/knowledge/plugin.py` on_activate ```python event_bus.subscribe("contact.created", handle_entity_created) event_bus.subscribe("contact.updated", handle_entity_updated) event_bus.subscribe("mail.received", handle_entity_created) event_bus.subscribe("company.created", handle_entity_created) event_bus.subscribe("company.updated", handle_entity_updated) ``` **Verbindungen:** - `event_handler.py` → `event_bus.subscribe()` → `jobs.enqueue_job()` → extraction pipeline - Extraction pipeline: source_adapter → entity_extractor → relationship_extractor → graph_rag **Tests:** - `tests/test_event_driven_extraction.py` — Fire event → verify extraction job enqueued and executed **Frontend:** - Knowledge Dashboard: Extraction status per source --- ### H-CITE: Evidence/Source References — auf unified_search aufbauen **Basis:** `app/plugins/builtins/unified_search/` (search results with entity references) **Was neu gebaut wird:** 1. **Evidence Model:** `app/plugins/builtins/knowledge/models.py` erweitern ```python class KnowledgeEvidence(Base, TenantMixin): __tablename__ = "knowledge_evidence" id: UUID PK relationship_id: UUID FK → entity_relationships.id source_type: str # "contact", "mail", "file", "wiki_article" source_id: UUID source_snippet: str # Text snippet that supports the relationship confidence: float created_at: TIMESTAMPTZ ``` 2. **Evidence Service:** `app/plugins/builtins/knowledge/evidence_service.py` (NEU, ~150 Zeilen) ```python async def add_evidence(db, tenant_id, relationship_id, source_type, source_id, snippet, confidence): ... async def get_evidence_for_relationship(db, tenant_id, relationship_id) -> list[dict]: ... async def search_evidence(db, tenant_id, query) -> list[dict]: ... ``` 3. **Evidence API:** `app/plugins/builtins/knowledge/routes.py` erweitern - `GET /api/v1/knowledge/relationships/{id}/evidence` — List evidence for relationship - `POST /api/v1/knowledge/relationships/{id}/evidence` — Add evidence **Verbindungen:** - Evidence linked to EntityRelationship (graph_rag) - Source references use unified_search entity types **Migrationen:** - `0134_knowledge_evidence.py` — knowledge_evidence table **Tests:** - `tests/test_knowledge_evidence.py` — Add, list, search evidence **Frontend:** - Relationship Detail: Evidence section --- ### H-WIKI-SEARCH: Wiki Search Provider — Wiki als unified_search provider registrieren **Basis:** `app/plugins/builtins/wiki/` (WikiPlugin, models, routes), `app/plugins/builtins/unified_search/` (BaseSearchProvider) **Was existiert:** - WikiPlugin hat kein on_activate (kein search provider registriert) - WikiArticle Model existiert - BaseSearchProvider mit search_fts, search_vector **Was neu gebaut wird:** 1. **Wiki Search Provider:** `app/plugins/builtins/wiki/search_provider.py` (NEU, ~200 Zeilen) ```python class WikiSearchProvider(BaseSearchProvider): entity_type = "wiki_article" supports_fts = True supports_vector = True supports_rag = True async def _search_fts_filtered(self, db, tsquery, tenant_id, limit, visible_ids): # Search wiki_articles.title and content via FTS async def _search_vector_filtered(self, db, embedding, tenant_id, limit, visible_ids): # Search wiki_articles.embedding via vector cosine async def get_embedding_text(self, db, entity_id, tenant_id) -> str: # Return article title + content def to_search_result(self, article) -> dict: return {"entity_type": "wiki_article", "entity_id": str(article.id), "title": article.title, ...} ``` 2. **Wiki Plugin on_activate:** `app/plugins/builtins/wiki/plugin.py` erweitern ```python async def on_activate(self, db, service_container, event_bus) -> None: from app.plugins.builtins.unified_search.contracts import get_search_registry from app.plugins.builtins.wiki.search_provider import WikiSearchProvider registry = get_search_registry() registry.register(WikiSearchProvider()) await super().on_activate(db, service_container, event_bus) async def on_deactivate(self, db, service_container, event_bus) -> None: from app.plugins.builtins.unified_search.contracts import get_search_registry get_search_registry().unregister("wiki_article") await super().on_deactivate(db, service_container, event_bus) ``` 3. **Wiki Embedding Index:** `app/plugins/builtins/wiki/models.py` erweitern - WikiArticle braucht `embedding` vector(768) column (via migration) - `lifecycle.py` hook: On article save, generate embedding via `llm_embed()` **Verbindungen:** - `wiki/plugin.py on_activate` → `unified_search.contracts.get_search_registry().register(WikiSearchProvider())` - `wiki/search_provider.py` extends `unified_search.base_provider.BaseSearchProvider` **Migrationen:** - `0135_wiki_embedding_column.py` — Add embedding vector(768) to wiki_articles **Tests:** - `tests/test_wiki_search_provider.py` — FTS and vector search on wiki articles - `tests/test_wiki_search_integration.py` — Wiki results appear in unified search **Frontend:** - Keine Änderung — Wiki results appear in global search --- ### H-WIKI-EMBED: Wiki Embeddings — auf llm_client.llm_embed aufbauen **Basis:** `app/ai/llm_client.py` (llm_embed), `app/plugins/builtins/wiki/models.py` **Was neu gebaut wird:** 1. **Wiki Embedding Service:** `app/plugins/builtins/wiki/embedding_service.py` (NEU, ~150 Zeilen) ```python async def generate_wiki_embedding(db, tenant_id, article_id) -> None: """Generate and store embedding for wiki article.""" from app.ai.llm_client import llm_embed article = await get_article(db, tenant_id, article_id) text = f"{article.title}\n\n{article.content}" embedding = await llm_embed(text) article.embedding = embedding await db.commit() async def batch_generate_embeddings(db, tenant_id, batch_size=50) -> int: """Generate embeddings for all articles without one.""" ``` 2. **Wiki Lifecycle Hook:** `app/plugins/builtins/wiki/plugin.py` on_activate - Register hook: `wiki.article.after_create` → `generate_wiki_embedding()` - Register hook: `wiki.article.after_update` → `generate_wiki_embedding()` (re-embed) 3. **Batch Job:** `app/plugins/builtins/wiki/jobs.py` (NEU) - ARQ job: `async def wiki_embed_batch(ctx, tenant_id)` - Called via cron or manual trigger **Verbindungen:** - `embedding_service.py` → `llm_client.llm_embed()` → update WikiArticle.embedding - `plugin.py on_activate` → `hooks.register_action()` **Migrationen:** - Siehe H-WIKI-SEARCH (0135 fügt embedding column hinzu) **Tests:** - `tests/test_wiki_embeddings.py` — Generate embedding, verify vector stored, batch job **Frontend:** - Wiki Settings: "Re-generate embeddings" button --- ### H-WIKI-LINK: Auto-Linking — Wiki → Entity Links **Basis:** `app/plugins/builtins/wiki/models.py` (entity links), `app/plugins/builtins/graph_rag/services.py` **Was neu gebaut wird:** 1. **Wiki Auto-Linker:** `app/plugins/builtins/wiki/auto_linker.py` (NEU, ~200 Zeilen) ```python async def auto_link_wiki_to_entities( db: AsyncSession, tenant_id: uuid.UUID, article_id: uuid.UUID, ) -> list[dict]: """Analyze wiki article content and auto-link to CRM entities.""" article = await get_article(db, tenant_id, article_id) # 1. Extract entities from article content from app.plugins.builtins.knowledge.entity_extractor import extract_entities entities = await extract_entities(db, tenant_id, article.content, "wiki_article", article_id) # 2. Create entity links for entity in entities: if entity.get("matched_id"): await create_entity_link(db, tenant_id, "wiki_article", article_id, entity["type"], entity["matched_id"]) # Also create graph_rag relationship await create_relationship(db, tenant_id, "wiki_article", article_id, entity["type"], entity["matched_id"], "references") return entities ``` 2. **Wiki Hook Integration:** `app/plugins/builtins/wiki/plugin.py` on_activate - Register hook: `wiki.article.after_create` → `auto_link_wiki_to_entities()` - Register hook: `wiki.article.after_update` → `auto_link_wiki_to_entities()` (re-link) **Verbindungen:** - `auto_linker.py` → `knowledge.entity_extractor.extract_entities()` → `graph_rag.services.create_relationship()` - `plugin.py on_activate` → `hooks.register_action()` **Tests:** - `tests/test_wiki_auto_linking.py` — Create article mentioning contact → verify link created **Frontend:** - Wiki Article Detail: "Linked Entities" section (auto-linked + manual) --- ### H-DATA-LIFE: Derived-Data Lifecycle — auf Outbox + Event Bus aufbauen **Basis:** `app/core/outbox.py` (enqueue_outbox_event), `app/core/event_bus.py`, `app/core/hooks.py` **Was neu gebaut wird:** 1. **Derived Data Tracker:** `app/plugins/builtins/knowledge/derived_data.py` (NEU, ~250 Zeilen) ```python class DerivedDataRegistry: """Tracks which derived data depends on which source data.""" def register_dependency(self, source_type, source_id, derived_type, derived_id, plugin_name): ... def get_dependencies(self, source_type, source_id) -> list[dict]: ... async def invalidate_dependents(self, db, tenant_id, source_type, source_id): ... ``` 2. **Derived Data Model:** `app/plugins/builtins/knowledge/models.py` erweitern ```python class DerivedDataDependency(Base, TenantMixin): __tablename__ = "derived_data_dependencies" id: UUID PK source_type: str source_id: UUID derived_type: str # "embedding", "graph_relationship", "search_index", "rag_chunk" derived_id: UUID plugin_name: str is_valid: bool = True invalidated_at: TIMESTAMPTZ | None created_at: TIMESTAMPTZ ``` 3. **Event Handler:** `app/plugins/builtins/knowledge/lifecycle_handler.py` (NEU, ~150 Zeilen) ```python async def handle_source_updated(event_name, payload): """When source data changes, invalidate derived data.""" source_type = payload["entity_type"] source_id = uuid.UUID(payload["entity_id"]) tenant_id = uuid.UUID(payload["tenant_id"]) await registry.invalidate_dependents(db, tenant_id, source_type, source_id) # Enqueue re-computation job await enqueue_job("recompute_derived_data", { "source_type": source_type, "source_id": str(source_id), "tenant_id": str(tenant_id), }) ``` 4. **Plugin Registration:** `app/plugins/builtins/knowledge/plugin.py` on_activate - Subscribe to: `contact.updated`, `company.updated`, `mail.updated`, `file.updated`, `wiki.article.updated` - On event: `handle_source_updated()` **Verbindungen:** - `lifecycle_handler.py` → `outbox.enqueue_outbox_event()` → `event_bus` - `derived_data.py` → tracks dependencies across plugins **Migrationen:** - `0136_derived_data_dependencies.py` — derived_data_dependencies table **Tests:** - `tests/test_derived_data_lifecycle.py` — Update contact → verify embeddings invalidated and recomputed **Frontend:** - Admin: Derived data status dashboard --- ### H-RET: Knowledge Retention — auf bestehende Retention Patterns aufbauen **Basis:** `app/core/outbox.py` (DLQ, retention), `app/models/audit.py` (audit log retention) **Was neu gebaut wird:** 1. **Retention Policy Model:** `app/plugins/builtins/knowledge/models.py` erweitern ```python class KnowledgeRetentionPolicy(Base, TenantMixin): __tablename__ = "knowledge_retention_policies" id: UUID PK entity_type: str # "knowledge_source", "entity_relationship", "knowledge_evidence" max_age_days: int | None # None = unlimited max_count: int | None # None = unlimited action: str # "archive", "delete", "anonymize" is_active: bool = True created_at: TIMESTAMPTZ ``` 2. **Retention Service:** `app/plugins/builtins/knowledge/retention.py` (NEU, ~200 Zeilen) ```python async def apply_retention_policies(db, tenant_id) -> dict: """Apply all active retention policies.""" policies = await get_active_policies(db, tenant_id) results = {} for policy in policies: if policy.action == "archive": count = await archive_old_records(db, tenant_id, policy) elif policy.action == "delete": count = await delete_old_records(db, tenant_id, policy) elif policy.action == "anonymize": count = await anonymize_old_records(db, tenant_id, policy) results[policy.entity_type] = count return results ``` 3. **Retention Job:** `app/plugins/builtins/knowledge/jobs.py` erweitern - ARQ job: `async def apply_retention(ctx, tenant_id)` - Cron: täglich um 3 Uhr 4. **Retention API:** `app/plugins/builtins/knowledge/routes.py` erweitern - `GET /api/v1/knowledge/retention/policies` — List policies - `POST /api/v1/knowledge/retention/policies` — Create policy - `PUT /api/v1/knowledge/retention/policies/{id}` — Update - `DELETE /api/v1/knowledge/retention/policies/{id}` — Delete - `POST /api/v1/knowledge/retention/apply` — Manual apply **Verbindungen:** - `retention.py` → KnowledgeSource, EntityRelationship, KnowledgeEvidence models - Cron job registriert via automation plugin cron_jobs contribution **Migrationen:** - `0137_knowledge_retention.py` — knowledge_retention_policies table **Tests:** - `tests/test_knowledge_retention.py` — Create old records, apply policy, verify archived/deleted **Frontend:** - `frontend/src/pages/KnowledgeRetentionSettings.tsx` (NEU) — Policy management --- ## Phase I ### I-1: Cross-System Integration — Agent→Workflow, Agent→Knowledge, MCP **Basis:** `app/ai/agent_loop.py`, `app/workflows/engine.py`, `app/plugins/builtins/knowledge/` **Was neu gebaut wird:** 1. **Agent → Workflow Bridge:** `app/plugins/builtins/automation/agent_workflow_bridge.py` (NEU, ~200 Zeilen) ```python async def agent_trigger_workflow( db: AsyncSession, tenant_id: uuid.UUID, agent_id: uuid.UUID, workflow_id: uuid.UUID, context: dict, ) -> dict: """Agent triggers a workflow execution.""" from app.services.workflow_service import create_instance instance = await create_instance(db, tenant_id, workflow_id, context, initiated_by=agent_id) # Post to agent workstream await post_workstream_update(...) return {"workflow_instance_id": str(instance.id)} ``` - Register as agent tool: `trigger_workflow` 2. **Agent → Knowledge Bridge:** `app/plugins/builtins/automation/agent_knowledge_bridge.py` (NEU, ~200 Zeilen) ```python async def agent_query_knowledge( db: AsyncSession, tenant_id: uuid.UUID, query: str, entity_type: str | None = None, ) -> dict: """Agent queries knowledge graph and unified search.""" from app.plugins.builtins.unified_search.contracts import get_search_registry from app.plugins.builtins.graph_rag.services import traverse_graph # 1. Unified search for relevant entities # 2. Graph traversal for related entities # 3. Combine and return ``` - Register as agent tool: `query_knowledge` 3. **MCP Exposure:** `app/plugins/builtins/mcp/` (NEU, komplettes Plugin, ~400 Zeilen) - `plugin.py` — `MCPPlugin(BasePlugin)` mit Manifest - `server.py` — MCP Server mit tool definitions für CRM entities - `routes.py` — `/api/v1/mcp/tools`, `/api/v1/mcp/call` - Exposes CRM operations as MCP tools for external AI clients **Verbindungen:** - `agent_workflow_bridge.py` → `workflow_service.create_instance()` - `agent_knowledge_bridge.py` → `unified_search` + `graph_rag` - MCP Plugin → CRM routes (read-only initially) **Migrationen:** - Plugin-Migration: `mcp/migrations/0001_initial.sql` **Tests:** - `tests/test_agent_workflow_bridge.py` — Agent triggers workflow - `tests/test_agent_knowledge_bridge.py` — Agent queries knowledge - `tests/test_mcp_exposure.py` — MCP tool listing and calling **Frontend:** - Agent Editor: Available tools include `trigger_workflow`, `query_knowledge` - Settings: MCP configuration --- ### I-2 bis I-5: Human-AI Workstream — auf kommunikation Plugin aufbauen **Basis:** `app/plugins/builtins/kommunikation/` (CommConversation, CommMessageBlock, MiniAppRegistry) **Was neu gebaut wird:** 1. **Human-AI Workstream Service:** `app/plugins/builtins/kommunikation/human_ai_workstream.py` (NEU, ~300 Zeilen) ```python async def create_human_ai_session( db, tenant_id, user_id, agent_id, topic: str, ) -> dict: """Create a human-AI collaboration session (CommConversation with type=human_ai).""" conv = await create_plugin_room(db, tenant_id, "kommunikation", f"hai_{agent_id}_{user_id}", title=f"Human-AI: {topic}", metadata={"type": "human_ai", "agent_id": str(agent_id)}) # Add user and agent as participants return conv async def post_ai_proposal(db, tenant_id, conversation_id, proposal_type, content, actions): """Post an AI proposal with action_card block for human review.""" blocks = [{"block_type": "action_card", "block_data": { "title": f"AI Proposal: {proposal_type}", "body": content, "actions": actions, # [{label, action, type: "approve/reject/edit"}] }}] await send_message(db, tenant_id, conversation_id, sender_type="ai", ...) async def handle_human_response(db, tenant_id, conversation_id, message_id, response_type, edits): """Process human response to AI proposal.""" ``` 2. **Workstream Block Types:** `content_types.py` erweitern - `"ai_proposal"` — fields: proposal_type, content, actions, confidence - `"human_decision"` — fields: decision, rationale, decided_by - `"ai_explanation"` — fields: explanation, evidence_ids, confidence 3. **Workstream API:** `app/plugins/builtins/kommunikation/routes.py` erweitern - `POST /api/v1/comm/workstream/sessions` — Create human-AI session - `GET /api/v1/comm/workstream/sessions` — List sessions - `POST /api/v1/comm/workstream/sessions/{id}/proposals` — Post AI proposal - `POST /api/v1/comm/workstream/sessions/{id}/responses` — Human response **Verbindungen:** - `human_ai_workstream.py` → `kommunikation.services.send_message()` + `create_plugin_room()` - Agent Runner → `post_ai_proposal()` when agent needs human input **Tests:** - `tests/test_human_ai_workstream.py` — Session creation, proposal, response flow **Frontend:** - `frontend/src/components/comm/AIProposalBlock.tsx` (NEU) - `frontend/src/components/comm/HumanDecisionBlock.tsx` (NEU) - `frontend/src/pages/HumanAIWorkstream.tsx` (NEU) --- ### I-6 bis I-8: MiniApp Runtime — auf kommunikation/miniapp_registry aufbauen **Basis:** `app/plugins/builtins/kommunikation/miniapp_registry.py` (MiniAppRegistry, MiniAppDef) **Was existiert:** - MiniAppRegistry mit register/unregister/list_apps/get_app - MiniAppDef: app_id, name, icon, description, plugin_name, render_schema - CommMessageBlock block_type="miniapp" bereits definiert **Was neu gebaut wird:** 1. **MiniApp Runtime Service:** `app/plugins/builtins/kommunikation/miniapp_runtime.py` (NEU, ~300 Zeilen) ```python class MiniAppRuntime: """Executes mini-app actions and manages state.""" async def execute_action( self, db, tenant_id, app_id, action_name, params, user_id, ) -> dict: """Execute a mini-app action.""" app = get_miniapp_registry().get_app(app_id) if not app: raise NotFoundError(f"MiniApp {app_id} not found") # Execute action via plugin's action handler handler = self._get_action_handler(app.plugin_name, app_id) result = await handler(db, tenant_id, action_name, params, user_id) return result async def get_state(self, db, tenant_id, app_id, context) -> dict: """Get current mini-app state for rendering.""" ``` 2. **MiniApp Action API:** `app/plugins/builtins/kommunikation/routes.py` erweitern - `GET /api/v1/comm/miniapps` — List all registered mini-apps - `GET /api/v1/comm/miniapps/{app_id}` — Get mini-app definition - `POST /api/v1/comm/miniapps/{app_id}/actions` — Execute action - `GET /api/v1/comm/miniapps/{app_id}/state` — Get state 3. **Plugin Action Handler Convention:** Plugins die MiniApps registrieren, stellen action handlers bereit - In `plugin.py`: `def get_miniapp_action_handler(self, app_id) -> Callable` - Automation plugin: Agent control mini-app - DMS plugin: File preview mini-app - Tasks plugin: Task board mini-app **Verbindungen:** - `miniapp_runtime.py` → `miniapp_registry.get_app()` → plugin action handler - `routes.py` → `miniapp_runtime.execute_action()` **Tests:** - `tests/test_miniapp_runtime.py` — Register app, execute action, get state **Frontend:** - `frontend/src/components/comm/MiniAppBlockRenderer.tsx` (NEU) — Renders miniapp blocks - `frontend/src/components/miniapps/GenericMiniApp.tsx` (NEU) — Schema-based renderer --- ### I-9 bis I-12: Dashboard & Analytics — echte DB-Queries **Basis:** Alle vorhandenen Models (Contact, Company, Mail, Task, Calendar, DMS, AgentRun, WorkflowInstance, etc.) **Was neu gebaut wird:** 1. **Dashboard Service:** `app/services/dashboard_service.py` (NEU, ~400 Zeilen) ```python async def get_dashboard_metrics(db, tenant_id, user_id) -> dict: """Get real dashboard metrics from DB.""" return { "contacts": await count_contacts(db, tenant_id), "companies": await count_companies(db, tenant_id), "tasks": { "open": await count_tasks(db, tenant_id, status="open"), "overdue": await count_overdue_tasks(db, tenant_id), "completed_this_week": await count_completed_this_week(db, tenant_id), }, "mails": { "unread": await count_unread_mails(db, tenant_id, user_id), "total_today": await count_mails_today(db, tenant_id), }, "agents": { "active_runs": await count_active_agent_runs(db, tenant_id), "total_runs": await count_total_agent_runs(db, tenant_id), "success_rate": await calculate_agent_success_rate(db, tenant_id), }, "workflows": { "active_instances": await count_active_workflow_instances(db, tenant_id), "completed_this_month": await count_completed_workflows(db, tenant_id), }, "knowledge": { "relationships": await count_relationships(db, tenant_id), "sources_extracted": await count_extracted_sources(db, tenant_id), }, } async def get_activity_feed(db, tenant_id, user_id, limit=50) -> list[dict]: """Get recent activity across all systems.""" # Query audit_log, agent_runs, workflow_instances, comm_messages ``` 2. **Dashboard API:** `app/routes/dashboard.py` (NEU) - `GET /api/v1/dashboard/metrics` — Real-time metrics - `GET /api/v1/dashboard/activity-feed` — Activity feed - `GET /api/v1/dashboard/trends` — Trend data (7/30/90 days) **Verbindungen:** - `dashboard_service.py` → alle Core + Plugin Models (echte SQL queries) - Keine Mocks, keine hardcoded data **Tests:** - `tests/test_dashboard.py` — Verify metrics match real DB counts **Frontend:** - `frontend/src/pages/Dashboard.tsx` (NEU/überarbeitet) — Real metrics - `frontend/src/components/dashboard/MetricCard.tsx` (NEU) - `frontend/src/components/dashboard/ActivityFeed.tsx` (NEU) - `frontend/src/components/dashboard/TrendChart.tsx` (NEU) --- ### I-13 bis I-16: DSGVO Export — echte DB-Queries über alle Plugins **Basis:** Alle Models mit tenant_id + user_id reference **Was neu gebaut wird:** 1. **DSGVO Export Service:** `app/services/dsgvo_export.py` (NEU, ~400 Zeilen) ```python async def export_user_data(db, tenant_id, user_id) -> dict: """Export all data associated with a user for DSGVO compliance.""" data = { "user": await get_user_data(db, tenant_id, user_id), "contacts": await get_user_contacts(db, tenant_id, user_id), "companies": await get_user_companies(db, tenant_id, user_id), "tasks": await get_user_tasks(db, tenant_id, user_id), "calendar_events": await get_user_events(db, tenant_id, user_id), "mails": await get_user_mails(db, tenant_id, user_id), "files": await get_user_files(db, tenant_id, user_id), "audit_logs": await get_user_audit_logs(db, tenant_id, user_id), "agent_runs": await get_user_agent_runs(db, tenant_id, user_id), "workflow_instances": await get_user_workflows(db, tenant_id, user_id), "comm_messages": await get_user_comm_messages(db, tenant_id, user_id), "knowledge_relationships": await get_user_relationships(db, tenant_id, user_id), "approvals": await get_user_approvals(db, tenant_id, user_id), } return data async def export_to_zip(data: dict, output_path: str) -> str: """Export data as ZIP with JSON files.""" ``` 2. **DSGVO Export API:** `app/routes/dsgvo.py` (NEU) - `POST /api/v1/dsgvo/export` — Trigger export (async job) - `GET /api/v1/dsgvo/export/{job_id}` — Get export status - `GET /api/v1/dsgvo/export/{job_id}/download` — Download ZIP - `DELETE /api/v1/dsgvo/export/{job_id}` — Delete export file 3. **DSGVO Delete (Right to be forgotten):** - `POST /api/v1/dsgvo/delete` — Anonymize user data (soft-delete + PII removal) - `POST /api/v1/dsgvo/delete/{user_id}` — Admin: delete specific user data **Verbindungen:** - `dsgvo_export.py` → alle Plugin Models (echte SQL queries) - Export job via ARQ - Sensitive data excluded via `SENSITIVE_FIELDS` **Tests:** - `tests/test_dsgvo_export.py` — Export completeness, sensitive data exclusion - `tests/test_dsgvo_delete.py` — Anonymization correctness **Frontend:** - `frontend/src/pages/DSGVOExport.tsx` (NEU) — Export/Download/Delete UI --- ### I-17 bis I-20: Onboarding — Setup Wizard **Basis:** `app/config.py`, `app/models/system_settings.py`, alle Plugin on_activate **Was neu gebaut wird:** 1. **Onboarding Service:** `app/services/onboarding.py` (NEU, ~300 Zeilen) ```python async def get_onboarding_status(db, tenant_id, user_id) -> dict: """Check which onboarding steps are completed.""" steps = [ {"id": "profile", "label": "Complete your profile", "done": await is_profile_complete(db, user_id)}, {"id": "import_contacts", "label": "Import contacts", "done": await has_contacts(db, tenant_id)}, {"id": "mail_account", "label": "Connect mail account", "done": await has_mail_account(db, tenant_id, user_id)}, {"id": "first_agent", "label": "Create your first agent", "done": await has_agents(db, tenant_id)}, {"id": "first_workflow", "label": "Create a workflow", "done": await has_workflows(db, tenant_id)}, {"id": "knowledge_extraction", "label": "Run knowledge extraction", "done": await has_knowledge(db, tenant_id)}, ] return {"steps": steps, "completion": sum(s["done"] for s in steps) / len(steps)} async def complete_step(db, tenant_id, user_id, step_id, data) -> dict: """Mark onboarding step as complete.""" ``` 2. **Onboarding API:** `app/routes/onboarding.py` (NEU) - `GET /api/v1/onboarding/status` — Get onboarding status - `POST /api/v1/onboarding/steps/{step_id}/complete` — Complete step - `POST /api/v1/onboarding/skip` — Skip onboarding 3. **Onboarding Model:** `app/models/onboarding.py` (NEU) ```python class OnboardingProgress(Base, TenantMixin): __tablename__ = "onboarding_progress" id: UUID PK user_id: UUID FK → users.id step_id: str completed_at: TIMESTAMPTZ | None skipped: bool = False metadata_: dict = JSONB ``` **Verbindungen:** - `onboarding.py` → Contact, MailAccount, AgentDefinition, Workflow models - System settings for onboarding configuration **Migrationen:** - `0138_onboarding_progress.py` — onboarding_progress table **Tests:** - `tests/test_onboarding.py` — Status check, step completion, skip **Frontend:** - `frontend/src/pages/SetupWizard.tsx` (NEU) — Multi-step wizard - `frontend/src/components/onboarding/StepCard.tsx` (NEU) - `frontend/src/components/onboarding/ProgressBar.tsx` (NEU) --- ### I-21 bis I-23: Performance Optimization **Basis:** Alle Routes und Services, `app/core/redis.py` **Was neu gebaut wird:** 1. **Query Optimization:** - N+1 query detection und batch loading mit `selectinload()` / `joinedload()` - Pagination auf alle List-Endpoints (bereits teilweise vorhanden via `app/core/pagination.py`) - Redis caching für häufige Queries 2. **Cache Service:** `app/core/cache.py` erweitern ```python async def cached_query(key: str, ttl: int, query_func, *args, **kwargs): """Cache query result in Redis.""" redis = await get_redis() cached = await redis.get(key) if cached: return json.loads(cached) result = await query_func(*args, **kwargs) await redis.setex(key, ttl, json.dumps(result)) return result ``` 3. **Index Optimization:** - Audit aller DB-Indizes via `scripts/check_indexes.py` - Fehlende Indizes identifizieren und via Migration hinzufügen - Unused Indizes entfernen 4. **Frontend Performance:** - Code splitting: Lazy-load plugin pages - React Query: staleTime, cacheTime optimization - Bundle size analysis **Migrationen:** - `0139_performance_indexes.py` — Additional indexes based on query analysis **Tests:** - `tests/test_performance.py` — Query performance benchmarks - `tests/test_cache.py` — Cache hit/miss, TTL expiry **Frontend:** - `frontend/src/components/common/LazyPage.tsx` (NEU) — Lazy loading wrapper --- ### I-24 bis I-25: Final Polish **Basis:** Alle Systeme **Was gemacht wird:** 1. **Error Handling Polish:** - Alle Routes nutzen `build_error_response()` aus `error_codes.py` - Frontend Error Boundaries auf allen Plugin-Seiten - Partial-Failure-Semantik für Batch-Operationen 2. **Documentation Update:** - `docs/api-documentation.md` — Alle neuen Endpoints - `docs/plugin-development-guide.md` — Knowledge plugin, MCP plugin - `README.md` — Updated features list - `docs/test-strategy.md` — Updated test coverage 3. **PWA Re-activation:** - Service Worker in `frontend/src/main.tsx` re-aktivieren (aktuell deaktiviert) - Offline-first für kritische Views 4. **Accessibility Audit:** - ARIA attributes auf allen interaktiven Elementen - 44px touch targets - Keyboard navigation **Tests:** - `tests/test_error_handling.py` — Verify structured error responses - Frontend: `frontend/src/__tests__/accessibility.test.tsx` **Frontend:** - Error Boundaries: `frontend/src/components/common/ErrorBoundary.tsx` (NEU) - PWA: `frontend/src/sw.ts` (NEU/überarbeitet) --- ## Phase J ### J-1: Controlled Self-Improvement — auf echte DB-Tabellen + Services aufbauen **Basis:** `app/models/audit.py` (AuditLog), `app/plugins/builtins/automation/models.py` (AgentDefinition, AgentRun), `app/models/workflow.py` (WorkflowInstance) **Was neu gebaut wird:** 1. **Self-Improvement Plugin:** `app/plugins/builtins/self_improvement/` (NEU, komplettes Plugin) - `plugin.py` — `SelfImprovementPlugin(BasePlugin)` mit Manifest - `models.py` — ImprovementSignal, ImprovementProposal, ProposalVersion, EvaluationResult - `services.py` — Signal collection, proposal generation, evaluation - `routes.py` — `/api/v1/improvements` - `schemas.py` — Pydantic schemas - `jobs.py` — ARQ jobs for pattern detection and evaluation 2. **ImprovementSignal Model:** ```python class ImprovementSignal(Base, TenantMixin): __tablename__ = "improvement_signals" id: UUID PK signal_type: str # "repeated_error", "low_success_rate", "slow_workflow", "manual_feedback" source_type: str # "agent_run", "workflow_instance", "audit_log", "user_feedback" source_id: UUID severity: str # "low", "medium", "high" description: str metadata_: dict = JSONB detected_at: TIMESTAMPTZ status: str = "new" # "new", "analyzed", "proposal_created", "dismissed" ``` 3. **ImprovementProposal Model:** ```python class ImprovementProposal(Base, TenantMixin, OwnedMixin): __tablename__ = "improvement_proposals" id: UUID PK signal_id: UUID FK → improvement_signals.id title: str description: str proposed_changes: dict = JSONB # What should change current_version: int = 1 status: str = "draft" # "draft", "in_review", "approved", "rejected", "implemented", "evaluated" risk_assessment: dict = JSONB expected_impact: str # "low", "medium", "high" created_at: TIMESTAMPTZ updated_at: TIMESTAMPTZ ``` 4. **ProposalVersion Model:** ```python class ProposalVersion(Base, TenantMixin): __tablename__ = "improvement_proposal_versions" id: UUID PK proposal_id: UUID FK → improvement_proposals.id version: int content: dict = JSONB # Full proposal content at this version created_by: UUID created_at: TIMESTAMPTZ ``` 5. **EvaluationResult Model:** ```python class EvaluationResult(Base, TenantMixin): __tablename__ = "improvement_evaluation_results" id: UUID PK proposal_id: UUID FK → improvement_proposals.id metric_type: str # "success_rate", "error_count", "execution_time", "user_satisfaction" before_value: float after_value: float improvement_pct: float evaluated_at: TIMESTAMPTZ metadata_: dict = JSONB ``` **Verbindungen:** - Signal collection → AuditLog, AgentRun, WorkflowInstance queries - Proposal → Approval system (`app/core/approval.py`) - Evaluation → DB queries before/after implementation **Migrationen:** - `0140_self_improvement_tables.py` — improvement_signals, improvement_proposals, improvement_proposal_versions, improvement_evaluation_results **Tests:** - `tests/test_self_improvement.py` — Signal creation, proposal lifecycle, evaluation **Frontend:** - `frontend/src/pages/ImprovementCenter.tsx` (NEU) — Overview dashboard - `frontend/src/components/improvement/ProposalCard.tsx` (NEU) - `frontend/src/components/improvement/SignalList.tsx` (NEU) --- ### J-2: Improvement Signals — auf AuditLog + AgentRun + WorkflowInstance aufbauen **Basis:** `app/models/audit.py`, `app/plugins/builtins/automation/models.py`, `app/models/workflow.py` **Was neu gebaut wird:** 1. **Signal Detector:** `app/plugins/builtins/self_improvement/signal_detector.py` (NEU, ~300 Zeilen) ```python async def detect_repeated_errors(db, tenant_id, time_window_hours=24) -> list[dict]: """Detect repeated error patterns in audit logs.""" # Query audit_log for error entries, group by error_type + entity_type # If count > threshold, create ImprovementSignal async def detect_low_agent_success_rate(db, tenant_id) -> list[dict]: """Detect agents with low success rates.""" # Query agent_runs, calculate success rate per agent # If below threshold, create signal async def detect_slow_workflows(db, tenant_id) -> list[dict]: """Detect workflows with long execution times.""" # Query workflow_instances, calculate avg duration per workflow # If above threshold, create signal async def detect_manual_feedback(db, tenant_id) -> list[dict]: """Collect manual user feedback signals.""" # Query feedback entries (if feedback system exists) ``` 2. **Signal Collection Job:** `app/plugins/builtins/self_improvement/jobs.py` - ARQ job: `async def collect_signals(ctx, tenant_id)` - Cron: stündlich - Calls all detect_* functions **Verbindungen:** - `signal_detector.py` → AuditLog, AgentRun, WorkflowInstance models (echte SQL queries) - Cron job via automation plugin **Tests:** - `tests/test_signal_detection.py` — Create errors, run detection, verify signals --- ### J-3: Pattern Detection — auf echten Signalen aufbauen **Basis:** `app/plugins/builtins/self_improvement/models.py` (ImprovementSignal) **Was neu gebaut wird:** 1. **Pattern Detector:** `app/plugins/builtins/self_improvement/pattern_detector.py` (NEU, ~250 Zeilen) ```python async def detect_patterns(db, tenant_id, signal_ids: list[uuid.UUID]) -> list[dict]: """Analyze signals and detect patterns.""" signals = await get_signals(db, tenant_id, signal_ids) # Group signals by: # - Entity type (all errors on contacts) # - Agent (all failures from one agent) # - Workflow (all slow workflows) # - Time clustering (errors spike at certain times) patterns = [] # 1. Frequency analysis # 2. Correlation analysis # 3. Time-series analysis return patterns ``` 2. **LLM-assisted Pattern Analysis:** ```python async def llm_analyze_patterns(patterns: list[dict]) -> list[dict]: """Use LLM to suggest root causes and improvements.""" from app.ai.llm_client import llm_complete prompt = f"Analyze these patterns and suggest root causes: {json.dumps(patterns)}" response = await llm_complete(messages=[{"role": "user", "content": prompt}]) return parse_llm_analysis(response) ``` **Verbindungen:** - `pattern_detector.py` → ImprovementSignal queries → `llm_client.llm_complete()` **Tests:** - `tests/test_pattern_detection.py` — Create signals, detect patterns --- ### J-4: ImprovementProposal — echte SQLAlchemy Modelle + Migration **Basis:** `app/plugins/builtins/self_improvement/models.py` (aus J-1) **Was neu gebaut wird:** 1. **Proposal Generator:** `app/plugins/builtins/self_improvement/proposal_generator.py` (NEU, ~250 Zeilen) ```python async def generate_proposal(db, tenant_id, signal_id, pattern_analysis) -> dict: """Generate improvement proposal from signal and pattern analysis.""" from app.ai.llm_client import llm_complete prompt = f"""Based on the following signal and pattern analysis, generate a concrete improvement proposal. Signal: {signal_description} Pattern: {pattern_analysis} Return JSON with: title, description, proposed_changes, risk_assessment, expected_impact """ response = await llm_complete(...) proposal_data = parse_llm_proposal(response) # Create proposal in DB proposal = await create_proposal(db, tenant_id, signal_id, proposal_data) return proposal ``` 2. **Proposal API:** `app/plugins/builtins/self_improvement/routes.py` - `GET /api/v1/improvements/proposals` — List proposals - `GET /api/v1/improvements/proposals/{id}` — Get proposal - `POST /api/v1/improvements/proposals` — Create proposal - `PUT /api/v1/improvements/proposals/{id}` — Update - `POST /api/v1/improvements/proposals/{id}/submit` — Submit for review **Tests:** - `tests/test_proposal_generation.py` — Signal → proposal generation --- ### J-5: Versioned Drafts — auf bestehende Versionierung aufbauen **Basis:** `app/plugins/builtins/wiki/models.py` (WikiArticleVersion), `app/plugins/builtins/self_improvement/models.py` (ProposalVersion) **Was neu gebaut wird:** 1. **Version Service:** `app/plugins/builtins/self_improvement/version_service.py` (NEU, ~200 Zeilen) ```python async def create_version(db, tenant_id, proposal_id, content, user_id) -> dict: """Create a new version of a proposal.""" current_version = await get_current_version(db, tenant_id, proposal_id) new_version = current_version + 1 version = ProposalVersion(proposal_id=proposal_id, version=new_version, content=content, ...) db.add(version) await db.commit() # Update proposal.current_version return version async def get_version_history(db, tenant_id, proposal_id) -> list[dict]: ... async def restore_version(db, tenant_id, proposal_id, version_id) -> dict: ... async def diff_versions(db, tenant_id, proposal_id, v1, v2) -> dict: ... ``` 2. **Version API:** `app/plugins/builtins/self_improvement/routes.py` erweitern - `GET /api/v1/improvements/proposals/{id}/versions` — Version history - `GET /api/v1/improvements/proposals/{id}/versions/{v}` — Get specific version - `POST /api/v1/improvements/proposals/{id}/versions/{v}/restore` — Restore - `GET /api/v1/improvements/proposals/{id}/diff?v1=1&v2=2` — Diff **Tests:** - `tests/test_proposal_versioning.py` — Create, restore, diff versions --- ### J-6: Evaluation/Sandbox — auf Test-DB aufbauen **Basis:** `tests/conftest.py` (test DB setup), `app/plugins/builtins/self_improvement/models.py` **Was neu gebaut wird:** 1. **Sandbox Evaluator:** `app/plugins/builtins/self_improvement/sandbox.py` (NEU, ~300 Zeilen) ```python async def evaluate_proposal_in_sandbox( proposal_id: uuid.UUID, tenant_id: uuid.UUID, ) -> dict: """Evaluate proposal in isolated sandbox environment.""" # 1. Create sandbox DB transaction (savepoint) # 2. Apply proposed changes # 3. Run test scenarios # 4. Measure metrics (success rate, error count, execution time) # 5. Rollback transaction # 6. Return before/after comparison ``` 2. **Test Scenario Runner:** `app/plugins/builtins/self_improvement/test_runner.py` (NEU, ~200 Zeilen) ```python async def run_test_scenarios(db, tenant_id, scenarios: list[dict]) -> dict: """Run predefined test scenarios and collect metrics.""" results = [] for scenario in scenarios: result = await execute_scenario(db, tenant_id, scenario) results.append(result) return aggregate_results(results) ``` 3. **Evaluation API:** `app/plugins/builtins/self_improvement/routes.py` erweitern - `POST /api/v1/improvements/proposals/{id}/evaluate` — Run sandbox evaluation - `GET /api/v1/improvements/proposals/{id}/evaluation` — Get evaluation results **Tests:** - `tests/test_sandbox_evaluation.py` — Evaluate proposal, verify metrics --- ### J-7: Human Approval — auf approval.py aufbauen **Basis:** `app/core/approval.py` (ApprovalRequest, create_approval_request) **Was neu gebaut wird:** 1. **Proposal Approval Integration:** `app/plugins/builtins/self_improvement/approval_handler.py` (NEU, ~200 Zeilen) ```python async def request_proposal_approval( db, tenant_id, proposal_id, requested_by, user_id, ) -> dict: """Create approval request for improvement proposal.""" from app.core.approval import create_approval_request approval = await create_approval_request( db=db, tenant_id=tenant_id, entity_type="improvement_proposal", entity_id=proposal_id, action="implement_proposal", requested_by=requested_by, requested_by_type="system", approver_id=user_id, metadata={"proposal_id": str(proposal_id), "risk": proposal.risk_assessment}, ) # Post to workstream await post_workflow_approval_request(...) return approval async def handle_approval_resolved(db, tenant_id, approval_id, status, user_id): """Handle approval resolution — implement or reject proposal.""" if status == "approved": await implement_proposal(db, tenant_id, proposal_id) else: await reject_proposal(db, tenant_id, proposal_id) ``` 2. **Approval Hook:** Register callback for `improvement_proposal` entity type - When ApprovalRequest resolved → `handle_approval_resolved()` **Verbindungen:** - `approval_handler.py` → `app.core.approval.create_approval_request()` - Approval resolution → `implement_proposal()` or `reject_proposal()` - Workstream notification via `post_workflow_approval_request()` **Tests:** - `tests/test_proposal_approval.py` — Request approval, approve, implement --- ### J-8: Impact Measurement — auf echte DB-Queries aufbauen **Basis:** `app/plugins/builtins/self_improvement/models.py` (EvaluationResult), `app/models/audit.py` **Was neu gebaut wird:** 1. **Impact Measurement Service:** `app/plugins/builtins/self_improvement/impact_measurement.py` (NEU, ~250 Zeilen) ```python async def measure_impact_before(db, tenant_id, proposal_id) -> dict: """Measure metrics before proposal implementation.""" return { "error_count": await count_errors(db, tenant_id, since=proposal.created_at), "agent_success_rate": await calculate_agent_success_rate(db, tenant_id), "workflow_avg_duration": await calculate_avg_workflow_duration(db, tenant_id), "user_satisfaction": await get_satisfaction_score(db, tenant_id), } async def measure_impact_after(db, tenant_id, proposal_id, days=7) -> dict: """Measure metrics after proposal implementation.""" # Same metrics, but after implementation date async def calculate_improvement(before: dict, after: dict) -> dict: """Calculate improvement percentages.""" improvements = {} for key in before: if before[key] != 0: improvements[key] = ((after[key] - before[key]) / before[key]) * 100 else: improvements[key] = 0.0 return improvements ``` 2. **Impact API:** `app/plugins/builtins/self_improvement/routes.py` erweitern - `GET /api/v1/improvements/proposals/{id}/impact` — Get impact measurement - `POST /api/v1/improvements/proposals/{id}/measure-impact` — Trigger measurement **Verbindungen:** - `impact_measurement.py` → AuditLog, AgentRun, WorkflowInstance queries (echte SQL) - Results stored in EvaluationResult model **Tests:** - `tests/test_impact_measurement.py` — Before/after metrics, improvement calculation --- ### J-9: Pattern Insight Frontend **Frontend:** - `frontend/src/components/improvement/PatternInsight.tsx` (NEU) — Pattern visualization - `frontend/src/components/improvement/ImpactChart.tsx` (NEU) — Before/after comparison - `frontend/src/components/improvement/SignalTimeline.tsx` (NEU) — Signal timeline --- ### J-10: Self-Improvement Documentation - `docs/self-improvement-guide.md` (NEU) — How the self-improvement system works - `docs/api-documentation.md` — Update with improvement endpoints - `PROGRESS.md` — Update with J-phase status --- ## Phase K ### K-1: EU Compliance Finalization — AI Use-Case Registry **Basis:** `app/ai/ai_use_case.py` (existing), `app/ai/data_policy.py`, `app/ai/transparency.py` **Was neu gebaut wird:** 1. **AI Use-Case Registry erweitern:** `app/ai/ai_use_case.py` - Vollständige Use-Case-Registration für alle AI-Features - Pflichtfelder: intended_purpose, owner, agents, models, provider, data_categories, allowed_actions, human_oversight_policy, risk_class - API: `GET /api/v1/compliance/ai-use-cases`, `POST /api/v1/compliance/ai-use-cases` 2. **Compliance Dashboard Service:** `app/services/compliance_service.py` (NEU, ~300 Zeilen) ```python async def get_compliance_status(db, tenant_id) -> dict: return { "ai_use_cases": await get_all_use_cases(db, tenant_id), "data_processing_activities": await get_processing_activities(db, tenant_id), "retention_policies": await get_retention_policies(db, tenant_id), "data_subject_requests": await get_dsr_status(db, tenant_id), "risk_assessments": await get_risk_assessments(db, tenant_id), } ``` **Tests:** - `tests/test_compliance_ai_use_cases.py` — Use-case registration, validation **Frontend:** - `frontend/src/pages/ComplianceDashboard.tsx` (NEU) --- ### K-2: Data Processing Registry **Basis:** `app/models/audit.py`, `app/core/sensitive_data.py` **Was neu gebaut wird:** 1. **Processing Activity Model:** `app/models/processing_activity.py` (NEU) ```python class ProcessingActivity(Base, TenantMixin): __tablename__ = "processing_activities" id: UUID PK name: str purpose: str legal_basis: str # DSGVO Art. 6 basis data_categories: list[str] = JSONB recipients: list[str] = JSONB retention_period_days: int | None dpia_required: bool = False dpia_completed: bool = False is_active: bool = True ``` 2. **Processing Activity API:** `app/routes/compliance.py` (NEU) - CRUD endpoints for processing activities **Migrationen:** - `0141_processing_activities.py` **Tests:** - `tests/test_processing_activities.py` --- ### K-3: Data Subject Request (DSR) Automation **Basis:** `app/services/dsgvo_export.py` (aus Phase I) **Was neu gebaut wird:** 1. **DSR Model:** `app/models/data_subject_request.py` (NEU) ```python class DataSubjectRequest(Base, TenantMixin): __tablename__ = "data_subject_requests" id: UUID PK request_type: str # "access", "rectification", "erasure", "portability", "restriction", "objection" requested_by: UUID FK → users.id status: str = "new" # "new", "processing", "completed", "rejected" requested_at: TIMESTAMPTZ completed_at: TIMESTAMPTZ | None result_data: dict = JSONB # Export data or action result metadata_: dict = JSONB ``` 2. **DSR Service:** `app/services/dsr_service.py` (NEU, ~250 Zeilen) - `create_request()`, `process_request()`, `complete_request()` - Access: triggers DSGVO export - Erasure: triggers anonymization - Rectification: triggers data update workflow 3. **DSR API:** `app/routes/compliance.py` erweitern - `POST /api/v1/compliance/dsr` — Create request - `GET /api/v1/compliance/dsr/{id}` — Status - `GET /api/v1/compliance/dsr` — List requests **Migrationen:** - `0142_data_subject_requests.py` **Tests:** - `tests/test_dsr.py` — All request types, processing flow **Frontend:** - `frontend/src/pages/DataSubjectRequests.tsx` (NEU) --- ### K-4: DPIA (Data Protection Impact Assessment) **Basis:** `app/models/processing_activity.py` (aus K-2) **Was neu gebaut wird:** 1. **DPIA Model:** `app/models/dpia.py` (NEU) ```python class DPIA(Base, TenantMixin): __tablename__ = "dpias" id: UUID PK processing_activity_id: UUID FK → processing_activities.id risk_level: str # "low", "medium", "high" assessment: dict = JSONB mitigation_measures: list[dict] = JSONB approved_by: UUID | None approved_at: TIMESTAMPTZ | None status: str = "draft" ``` 2. **DPIA API:** `app/routes/compliance.py` erweitern - CRUD for DPIAs **Migrationen:** - `0143_dpias.py` **Tests:** - `tests/test_dpia.py` --- ### K-5: AI Act Compliance **Basis:** `app/ai/ai_use_case.py`, `app/ai/transparency.py`, `app/ai/oversight.py` **Was neu gebaut wird:** 1. **AI Act Risk Classification:** `app/ai/ai_act_compliance.py` (NEU, ~200 Zeilen) ```python class AIActRiskClass: MINIMAL = "minimal" LIMITED = "limited" HIGH = "high" UNACCEPTABLE = "unacceptable" async def classify_ai_use_case(use_case) -> str: """Classify AI use case according to EU AI Act risk levels.""" # Based on: purpose, data categories, automation level, human oversight async def get_ai_act_requirements(risk_class: str) -> dict: """Get required compliance measures for risk class.""" ``` 2. **Transparency Requirements:** - AI-generierte Inhalte müssen gekennzeichnet sein (transparency.py bereits vorhanden) - AI-Akteure im Workstream als AI gekennzeichnet (kommunikation sender_type="ai") - Deepfake-Kennzeichnung für AI-generierte Medien 3. **Human Oversight Requirements:** - High-Risk AI-Use-Cases MÜSSEN Human Approval haben (approval.py bereits vorhanden) - Oversight-Records für alle High-Risk-Entscheidungen (oversight.py bereits vorhanden) **Tests:** - `tests/test_ai_act_compliance.py` — Risk classification, requirements --- ### K-6: Compliance Documentation & Audit Trail **Basis:** `app/models/audit.py`, `app/services/compliance_service.py` **Was neu gebaut wird:** 1. **Compliance Audit Trail:** - Alle Compliance-relevanten Aktionen werden im AuditLog protokolliert - DSR requests, DPIA approvals, AI use-case changes, data exports 2. **Compliance Report Generator:** `app/services/compliance_report.py` (NEU, ~200 Zeilen) ```python async def generate_compliance_report(db, tenant_id, period_start, period_end) -> dict: """Generate comprehensive compliance report for a period.""" return { "period": {"start": period_start, "end": period_end}, "ai_use_cases": await get_use_cases_summary(db, tenant_id, period_start, period_end), "data_exports": await get_exports_summary(db, tenant_id, period_start, period_end), "dsr_requests": await get_dsr_summary(db, tenant_id, period_start, period_end), "audit_trail": await get_audit_summary(db, tenant_id, period_start, period_end), "retention_actions": await get_retention_summary(db, tenant_id, period_start, period_end), } ``` 3. **Compliance Report API:** `app/routes/compliance.py` erweitern - `GET /api/v1/compliance/report?start=...&end=...` — Generate report - `POST /api/v1/compliance/report/export` — Export as PDF/JSON **Tests:** - `tests/test_compliance_report.py` — Report generation, completeness **Frontend:** - `frontend/src/pages/ComplianceReport.tsx` (NEU) --- ## Migrations-Übersicht | Migration | Phase | Beschreibung | |-----------|-------|--------------| | 0129 | B | IVFFlat index strategy config | | 0130 | B | storage_provider_configs table | | 0131 | B | Migrate notifications to comm + drop notification tables | | 0132 | H | knowledge_sources table | | 0133 | H | entity_relationships: confidence, review_status, reviewed_by | | 0134 | H | knowledge_evidence table | | 0135 | H | wiki_articles: embedding vector(768) column | | 0136 | H | derived_data_dependencies table | | 0137 | H | knowledge_retention_policies table | | 0138 | I | onboarding_progress table | | 0139 | I | Performance indexes | | 0140 | J | Self-improvement tables (signals, proposals, versions, evaluations) | | 0141 | K | processing_activities table | | 0142 | K | data_subject_requests table | | 0143 | K | dpias table | ## Neue Plugins | Plugin | Phase | Dependencies | |--------|-------|-------------| | storage_webdav | B | [] | | storage_nextcloud | B | [storage_webdav] | | knowledge | H | [unified_search, graph_rag] | | mcp | I | [] | | self_improvement | J | [automation] | ## Neue Frontend-Seiten | Seite | Phase | Pfad | |-------|-------|------| | StorageSettings | B | /settings/storage | | KnowledgeDashboard | H | /knowledge | | KnowledgeReviewQueue | H | /knowledge/review | | HumanAIWorkstream | I | /workstream | | Dashboard | I | /dashboard | | DSGVOExport | I | /settings/dsgvo | | SetupWizard | I | /onboarding | | ImprovementCenter | J | /improvements | | ComplianceDashboard | K | /compliance | | DataSubjectRequests | K | /compliance/dsr | | ComplianceReport | K | /compliance/report | ## Neue CommMessageBlock Types | Block Type | Phase | Verwendung | |-----------|-------|-----------| | agent_result | F | Agent run results | | approval_request | F | Approval requests in workstream | | task_card | F | Task references in workstream | | workflow_card | F | Workflow references in workstream | | knowledge_card | F | Knowledge entity references | | progress_card | F | Progress indicators | | ai_proposal | I | AI proposals for human review | | human_decision | I | Human decision records | | ai_explanation | I | AI explanations with evidence | ## Verbindungs-Matrix (Wichtigste) | Von | Nach | Mechanismus | |-----|------|-----------| | Agent Runner | Kommunikation | `kommunikation.contracts.send_message()` | | Agent Runner | Workstream | `agent_workstream.py` → `create_plugin_room()` | | Workflow Engine | Kommunikation | `workstream.py` → `send_message()` | | Knowledge Plugin | Unified Search | `get_search_registry().get()` → `get_embedding_text()` | | Knowledge Plugin | GraphRAG | `graph_rag.services.create_relationship()` | | Knowledge Plugin | LLM Client | `llm_complete()` für extraction | | Knowledge Plugin | Event Bus | `event_bus.subscribe()` für event-driven extraction | | Wiki Plugin | Unified Search | `get_search_registry().register(WikiSearchProvider())` | | Wiki Plugin | LLM Client | `llm_embed()` für embeddings | | Self-Improvement | Approval | `create_approval_request()` für proposal approval | | Self-Improvement | AuditLog | SQL queries für signal detection | | MCP Plugin | CRM Routes | Read-only CRM operations als MCP tools | | Storage Plugins | Storage Registry | `get_storage_registry().register()` | | Automation Plugin | Prebuilt Agents | `create_*_agent()` in `on_activate()` | --- ## Implementierungs-Reihenfolge 1. **Phase B Lücken** (3 Tasks) — Fundament: Storage, Index, Notification cleanup 2. **Phase F Lücken** (3 Tasks) — Agent Integration: Prebuilt registration, Agent→Comm, Workstream 3. **Phase G Lücke** (1 Task) — Workflow Workstream 4. **Phase H Rest** (12 Tasks) — Knowledge System auf graph_rag + unified_search 5. **Phase I** (25 Tasks) — Cross-System Integration, Human-AI, MiniApps, Dashboard, DSGVO, Onboarding 6. **Phase J** (10 Tasks) — Self-Improvement auf echte DB + AuditLog 7. **Phase K** (6 Tasks) — EU Compliance Finalization **Total: 60 Tasks** --- > Dieser Plan ist so detailliert, dass direkt mit der Implementierung begonnen werden kann. Jeder Task hat konkrete Dateipfade, Model-Definitionen, Route-Definitionen, und Verbindungs-Punkte zu vorhandenen Systemen.