From d47b7615dda6dd8ddecc29476b18be3cb6fd3040 Mon Sep 17 00:00:00 2001 From: Agent Zero Date: Thu, 20 Aug 2026 01:12:32 +0200 Subject: [PATCH] =?UTF-8?q?docs:=20architecture=20plan=20for=20all=20open?= =?UTF-8?q?=20tasks=20=E2=80=94=2060=20tasks=20across=20Phase=20B/F/G/H/I/?= =?UTF-8?q?J/K,=2015=20new=20migrations,=205=20new=20plugins,=209=20new=20?= =?UTF-8?q?CommMessageBlock=20types,=20all=20building=20on=20existing=20sy?= =?UTF-8?q?stems?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- ARCHITECTURE_PLAN.md | 2433 ++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 2433 insertions(+) create mode 100644 ARCHITECTURE_PLAN.md diff --git a/ARCHITECTURE_PLAN.md b/ARCHITECTURE_PLAN.md new file mode 100644 index 0000000..71568e7 --- /dev/null +++ b/ARCHITECTURE_PLAN.md @@ -0,0 +1,2433 @@ +# 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. \ No newline at end of file