sadpandajoe commented on code in PR #43134:
URL: https://github.com/apache/superset/pull/43134#discussion_r3817875454


##########
superset/ai/api.py:
##########
@@ -0,0 +1,989 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+REST API for the AI assistant.
+
+Every route carries ``@protect()`` and is reached through ``@expose`` on a
+``BaseSupersetApi`` subclass, which is what makes Flask-AppBuilder's
+authorization actually run. Ownership is enforced a second time in the command
+and DAO layers, so a conversation identifier is never on its own a capability.
+"""
+
+from __future__ import annotations
+
+import logging
+import time
+from collections.abc import Generator
+from typing import Any, cast
+
+from flask import current_app, request, Response, stream_with_context
+from flask_appbuilder.api import expose, permission_name, protect, safe
+from marshmallow import ValidationError
+
+from superset.ai.events import (
+    error_event,
+    KEEPALIVE_FRAME,
+    KEEPALIVE_INTERVAL_SECONDS,
+)
+from superset.ai.schemas import (
+    AgentResponseSchema,
+    CancelPostSchema,
+    FeedbackPostSchema,
+    MessagePostSchema,
+    RunAcceptedResponseSchema,
+    SuggestedPromptsPostSchema,
+    ThreadDetailResponseSchema,
+    ThreadPostSchema,
+    ThreadPutSchema,
+    ThreadResponseSchema,
+)
+from superset.ai.types import MessageRole, MessageStatus
+from superset.commands.ai.exceptions import (
+    AIChatMessageInvalidError,
+    AIChatMessageNotFoundError,
+    AIChatThreadInvalidError,
+    AIChatThreadNotFoundError,
+)
+from superset.extensions import event_logger
+from superset.utils.core import get_user_id
+from superset.utils.decorators import transaction
+from superset.views.base_api import BaseSupersetApi, statsd_metrics
+
+logger = logging.getLogger(__name__)
+
+#: Upper bound on how long a client may hold a stream open, so an abandoned
+#: browser tab cannot pin a worker indefinitely.
+_STREAM_TIMEOUT_SECONDS = 900
+
+#: How often a reader checks the event bus for new frames.
+#:
+#: Deliberately separate from ``KEEPALIVE_INTERVAL_SECONDS``. Passing the
+#: keep-alive interval as the poll interval made the reader sleep fifteen 
seconds
+#: between checks and then deliver everything that had accumulated in one 
batch —
+#: so a worker-mode run showed no streaming at all: the answer and every tool 
call
+#: appeared in fifteen-second lumps. One controls responsiveness, the other how
+#: often an idle connection is reassured; they are not the same number.
+_EVENT_POLL_SECONDS = 0.1
+
+
+class AIRestApi(BaseSupersetApi):
+    """Conversations with the AI assistant."""
+
+    resource_name = "ai"
+    openapi_spec_tag = "AI Assistant"
+    allow_browser_login = True
+    class_permission_name = "AIAssistant"
+
+    openapi_spec_component_schemas = (
+        AgentResponseSchema,
+        CancelPostSchema,
+        FeedbackPostSchema,
+        MessagePostSchema,
+        RunAcceptedResponseSchema,
+        SuggestedPromptsPostSchema,
+        ThreadDetailResponseSchema,
+        ThreadPostSchema,
+        ThreadPutSchema,
+        ThreadResponseSchema,
+    )
+
+    @expose("/agent/", methods=("GET",))
+    @protect()
+    @safe
+    @statsd_metrics
+    @permission_name("read")
+    def agents(self) -> Response:
+        """List agent profiles the current user may select.
+        ---
+        get:
+          summary: List available agent profiles
+          responses:
+            200:
+              description: Available profiles
+              content:
+                application/json:
+                  schema:
+                    type: object
+                    properties:
+                      result:
+                        type: array
+                        items:
+                          $ref: '#/components/schemas/AgentResponseSchema'
+            401:
+              $ref: '#/components/responses/401'
+            403:
+              $ref: '#/components/responses/403'
+            404:
+              $ref: '#/components/responses/404'
+        """
+        if (unavailable := self._reject_if_unconfigured()) is not None:
+            return unavailable
+
+        from superset.ai.factories import get_profiles
+
+        profiles = get_profiles().visible_to_current_user()
+        return self.response(200, result=[p.to_public_dict() for p in 
profiles])
+
+    @expose("/model/", methods=("GET",))
+    @protect()
+    @safe
+    @statsd_metrics
+    @permission_name("read")
+    def models(self) -> Response:
+        """List models this deployment has configured.
+        ---
+        get:
+          summary: List selectable models
+          responses:
+            200:
+              description: Configured model identifiers
+              content:
+                application/json:
+                  schema:
+                    type: object
+                    properties:
+                      result:
+                        type: array
+                        items:
+                          type: string
+            401:
+              $ref: '#/components/responses/401'
+            404:
+              $ref: '#/components/responses/404'
+        """
+        if (unavailable := self._reject_if_unconfigured()) is not None:
+            return unavailable
+
+        from superset.ai.factories import get_provider
+
+        return self.response(200, result=get_provider().available_models())
+
+    @expose("/thread/", methods=("POST",))
+    @protect()
+    @safe
+    @statsd_metrics
+    @permission_name("write")
+    @event_logger.log_this_with_context(
+        action=lambda self, *args, **kwargs: 
f"{self.__class__.__name__}.post_thread",
+        log_to_statsd=False,
+    )
+    def post_thread(self) -> Response:
+        """Create a conversation.
+        ---
+        post:
+          summary: Create a conversation
+          requestBody:
+            content:
+              application/json:
+                schema:
+                  $ref: '#/components/schemas/ThreadPostSchema'
+          responses:
+            201:
+              description: Conversation created
+              content:
+                application/json:
+                  schema:
+                    type: object
+                    properties:
+                      result:
+                        $ref: '#/components/schemas/ThreadResponseSchema'
+            400:
+              $ref: '#/components/responses/400'
+            401:
+              $ref: '#/components/responses/401'
+            404:
+              $ref: '#/components/responses/404'
+        """
+        if (unavailable := self._reject_if_unconfigured()) is not None:
+            return unavailable
+
+        from superset.commands.ai import CreateAIChatThreadCommand
+
+        try:
+            payload = ThreadPostSchema().load(request.json or {})
+        except ValidationError as error:
+            return self.response_400(message=error.messages)
+        try:
+            thread = CreateAIChatThreadCommand(
+                user_id=self._user_id(),
+                title=payload.get("title"),
+                agent_key=payload.get("agent_key"),
+            ).run()
+        except AIChatThreadInvalidError as ex:
+            return self.response_422(message=str(ex))
+        return self.response(201, result=_thread_dict(thread))
+
+    @expose("/thread/", methods=("GET",))
+    @protect()
+    @safe
+    @statsd_metrics
+    @permission_name("read")
+    def get_threads(self) -> Response:
+        """List the current user's conversations.
+        ---
+        get:
+          summary: List conversations
+          parameters:
+          - in: query
+            name: limit
+            schema:
+              type: integer
+          - in: query
+            name: offset
+            schema:
+              type: integer
+          responses:
+            200:
+              description: Conversations
+              content:
+                application/json:
+                  schema:
+                    type: object
+                    properties:
+                      count:
+                        type: integer
+                      result:
+                        type: array
+                        items:
+                          $ref: '#/components/schemas/ThreadResponseSchema'
+            401:
+              $ref: '#/components/responses/401'
+            404:
+              $ref: '#/components/responses/404'
+        """
+        if (unavailable := self._reject_if_unconfigured()) is not None:
+            return unavailable
+
+        from superset.daos.ai import AIChatThreadDAO
+
+        limit = request.args.get("limit", type=int) or 50
+        offset = request.args.get("offset", type=int) or 0
+        threads = AIChatThreadDAO.find_all_for_user(
+            self._user_id(), limit=limit, offset=offset
+        )
+        return self.response(
+            200,
+            count=len(threads),
+            result=[_thread_dict(thread) for thread in threads],
+        )
+
+    @expose("/thread/<thread_uuid>", methods=("GET",))
+    @protect()
+    @safe
+    @statsd_metrics
+    @permission_name("read")
+    def get_thread(self, thread_uuid: str) -> Response:
+        """Fetch a conversation and its messages.
+        ---
+        get:
+          summary: Get a conversation
+          parameters:
+          - in: path
+            name: thread_uuid
+            required: true
+            schema:
+              type: string
+              format: uuid
+          responses:
+            200:
+              description: Conversation with messages
+              content:
+                application/json:
+                  schema:
+                    type: object
+                    properties:
+                      result:
+                        $ref: '#/components/schemas/ThreadDetailResponseSchema'
+            401:
+              $ref: '#/components/responses/401'
+            404:
+              $ref: '#/components/responses/404'
+        """
+        if (unavailable := self._reject_if_unconfigured()) is not None:
+            return unavailable
+
+        from superset.daos.ai import (
+            AIChatFeedbackDAO,
+            AIChatMessageDAO,
+            AIChatThreadDAO,
+        )
+
+        user_id = self._user_id()
+        thread = AIChatThreadDAO.find_by_uuid_for_user(thread_uuid, user_id)
+        if thread is None:
+            return self.response_404()
+
+        messages = AIChatMessageDAO.find_for_thread(thread)
+        # Resolved for the whole transcript at once so the panel can show which
+        # replies this user already rated; without it a reload loses the 
verdict
+        # and the message looks unrated.
+        verdicts = AIChatFeedbackDAO.find_verdicts_for_user(
+            [message.id for message in messages], user_id
+        )
+        detail = _thread_dict(thread)
+        detail["messages"] = [
+            _message_dict(message, liked=verdicts.get(message.id))
+            for message in messages
+        ]
+        return self.response(200, result=detail)
+
+    @expose("/thread/<thread_uuid>", methods=("PUT",))
+    @protect()
+    @safe
+    @statsd_metrics
+    @permission_name("write")
+    @event_logger.log_this_with_context(
+        action=lambda self, *args, **kwargs: 
f"{self.__class__.__name__}.put_thread",
+        log_to_statsd=False,
+    )
+    def put_thread(self, thread_uuid: str) -> Response:
+        """Rename or archive a conversation.
+        ---
+        put:
+          summary: Update a conversation
+          parameters:
+          - in: path
+            name: thread_uuid
+            required: true
+            schema:
+              type: string
+              format: uuid
+          requestBody:
+            content:
+              application/json:
+                schema:
+                  $ref: '#/components/schemas/ThreadPutSchema'
+          responses:
+            200:
+              description: Conversation updated
+            401:
+              $ref: '#/components/responses/401'
+            404:
+              $ref: '#/components/responses/404'
+            422:
+              $ref: '#/components/responses/422'
+        """
+        if (unavailable := self._reject_if_unconfigured()) is not None:
+            return unavailable
+
+        from superset.commands.ai import UpdateAIChatThreadCommand
+
+        try:
+            payload = ThreadPutSchema().load(request.json or {})
+        except ValidationError as error:
+            return self.response_400(message=error.messages)
+        try:
+            thread = UpdateAIChatThreadCommand(
+                thread_uuid,
+                self._user_id(),
+                title=payload.get("title"),
+                status=payload.get("status"),
+            ).run()
+        except AIChatThreadNotFoundError:
+            return self.response_404()
+        except AIChatThreadInvalidError as ex:
+            return self.response_422(message=str(ex))
+        return self.response(200, result=_thread_dict(thread))
+
+    @expose("/thread/<thread_uuid>", methods=("DELETE",))
+    @protect()
+    @safe
+    @statsd_metrics
+    @permission_name("write")
+    @event_logger.log_this_with_context(
+        action=lambda self, *args, **kwargs: 
f"{self.__class__.__name__}.delete_thread",
+        log_to_statsd=False,
+    )
+    def delete_thread(self, thread_uuid: str) -> Response:
+        """Delete a conversation and its messages.
+        ---
+        delete:
+          summary: Delete a conversation
+          parameters:
+          - in: path
+            name: thread_uuid
+            required: true
+            schema:
+              type: string
+              format: uuid
+          responses:
+            200:
+              description: Conversation deleted
+            401:
+              $ref: '#/components/responses/401'
+            404:
+              $ref: '#/components/responses/404'
+        """
+        if (unavailable := self._reject_if_unconfigured()) is not None:
+            return unavailable
+
+        from superset.commands.ai import DeleteAIChatThreadCommand
+
+        try:
+            DeleteAIChatThreadCommand(thread_uuid, self._user_id()).run()
+        except AIChatThreadNotFoundError:
+            return self.response_404()
+        return self.response(200, message="OK")
+
+    @expose("/thread/<thread_uuid>/message", methods=("POST",))
+    @protect()
+    @safe
+    @statsd_metrics
+    @permission_name("write")
+    @event_logger.log_this_with_context(
+        action=lambda self, *args, **kwargs: 
f"{self.__class__.__name__}.post_message",
+        log_to_statsd=False,
+    )
+    def post_message(self, thread_uuid: str) -> Response:
+        """Post a user message and start a run.
+        ---
+        post:
+          summary: Post a message
+          description: >
+            Stores the user's message, creates a placeholder assistant message,
+            and starts a run. Returns immediately; consume the answer from the
+            stream endpoint using the returned run identifier.
+          parameters:
+          - in: path
+            name: thread_uuid
+            required: true
+            schema:
+              type: string
+              format: uuid
+          requestBody:
+            content:
+              application/json:
+                schema:
+                  $ref: '#/components/schemas/MessagePostSchema'
+          responses:
+            202:
+              description: Run accepted
+              content:
+                application/json:
+                  schema:
+                    type: object
+                    properties:
+                      result:
+                        $ref: '#/components/schemas/RunAcceptedResponseSchema'
+            400:
+              $ref: '#/components/responses/400'
+            401:
+              $ref: '#/components/responses/401'
+            404:
+              $ref: '#/components/responses/404'
+            422:
+              $ref: '#/components/responses/422'
+        """
+        if (unavailable := self._reject_if_unconfigured()) is not None:
+            return unavailable
+
+        from superset.ai.orchestrator import new_run_id
+        from superset.commands.ai import AppendAIChatMessageCommand
+
+        try:
+            payload = MessagePostSchema().load(request.json or {})
+        except ValidationError as error:
+            return self.response_400(message=error.messages)
+        user_id = self._user_id()
+
+        try:
+            user_message = AppendAIChatMessageCommand(
+                thread_uuid,
+                user_id,
+                MessageRole.USER,
+                payload["content"],
+                request_id=payload.get("request_id"),
+            ).run()
+            # Created up front so a client that reconnects before any token
+            # arrives still has a row to attach its stream to.
+            assistant_message = AppendAIChatMessageCommand(
+                thread_uuid,
+                user_id,
+                MessageRole.ASSISTANT,
+                "",
+                request_id=payload.get("request_id"),
+                status=MessageStatus.PENDING,
+            ).run()
+        except AIChatThreadNotFoundError:
+            return self.response_404()
+        except (AIChatMessageInvalidError, AIChatThreadInvalidError) as ex:
+            return self.response_422(message=str(ex))
+
+        run_id = new_run_id()

Review Comment:
   A retry with the same `request_id` reuses the existing messages but still 
creates a new run ID and starts another run. If the first 202 response is lost, 
the retry can launch a second inference/tool execution and overwrite the 
original run context. Should this return the existing run instead when the 
idempotent append did not create a new assistant message?



##########
superset-frontend/src/features/ai/components/ChatTabsMenu.tsx:
##########
@@ -0,0 +1,356 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+/**
+ * @fileoverview The conversation list.
+ *
+ * Conversations live behind one menu rather than a tab strip: the panel is 
narrow
+ * enough in floating mode that a strip would truncate every name, and the list
+ * doubles as the history of past conversations, which a strip cannot be.
+ */
+
+import { useCallback, useState } from 'react';
+import type { MouseEvent as ReactMouseEvent } from 'react';
+import { styled } from '@apache-superset/core/theme';
+import { t } from '@apache-superset/core/translation';
+import { Button, Dropdown, Popconfirm } from '@superset-ui/core/components';
+import { Icons } from '@superset-ui/core/components/Icons';
+import type { ChatTab } from '../types';
+
+const MenuContainer = styled.div`
+  background: ${({ theme }) => theme.colorBgElevated};
+  border-radius: ${({ theme }) => theme.borderRadius}px;
+  box-shadow: ${({ theme }) => theme.boxShadowSecondary};
+  min-width: ${({ theme }) => theme.sizeUnit * 65}px;
+  max-height: ${({ theme }) => theme.sizeUnit * 100}px;
+  overflow-y: auto;
+  border: 1px solid ${({ theme }) => theme.colorBorderSecondary};
+`;
+
+const MenuHeader = styled.div`
+  padding: ${({ theme }) => theme.sizeUnit * 3}px
+    ${({ theme }) => theme.sizeUnit * 4}px;
+  border-bottom: 1px solid ${({ theme }) => theme.colorBorderSecondary};
+  font-weight: ${({ theme }) => theme.fontWeightStrong};
+  font-size: ${({ theme }) => theme.fontSizeSM}px;
+  color: ${({ theme }) => theme.colorTextSecondary};
+  text-transform: uppercase;
+  letter-spacing: 0.5px;
+`;
+
+const NewChatButton = styled.button`
+  display: flex;
+  align-items: center;
+  gap: ${({ theme }) => theme.sizeUnit * 2}px;
+  width: 100%;
+  padding: ${({ theme }) => theme.sizeUnit * 2.5}px
+    ${({ theme }) => theme.sizeUnit * 4}px;
+  cursor: pointer;
+  color: ${({ theme }) => theme.colorPrimary};
+  font-weight: ${({ theme }) => theme.fontWeightStrong};
+  background: none;
+  border: none;
+  text-align: left;
+  transition: background ${({ theme }) => theme.motionDurationMid};
+
+  &:hover {
+    background: ${({ theme }) => theme.colorFillTertiary};
+  }
+`;
+
+const TabItem = styled.div<{ isActive: boolean }>`
+  display: flex;
+  align-items: center;
+  justify-content: space-between;
+  padding: ${({ theme }) => theme.sizeUnit * 2.5}px
+    ${({ theme }) => theme.sizeUnit * 4}px;
+  cursor: pointer;
+  background: ${({ theme, isActive }) =>
+    isActive ? theme.colorFillSecondary : 'transparent'};
+  border-left: 3px solid
+    ${({ theme, isActive }) => (isActive ? theme.colorPrimary : 
'transparent')};
+  transition: background ${({ theme }) => theme.motionDurationMid};
+
+  &:hover {
+    background: ${({ theme }) => theme.colorFillTertiary};
+
+    .action-btn {
+      opacity: 1;
+    }
+  }
+`;
+
+const TabInfo = styled.div`
+  display: flex;
+  align-items: center;
+  gap: ${({ theme }) => theme.sizeUnit * 2}px;
+  flex: 1;
+  overflow: hidden;
+`;
+
+const TabName = styled.span`
+  font-size: ${({ theme }) => theme.fontSize}px;
+  color: ${({ theme }) => theme.colorText};
+  white-space: nowrap;
+  overflow: hidden;
+  text-overflow: ellipsis;
+  max-width: ${({ theme }) => theme.sizeUnit * 35}px;
+`;
+
+const TabNameInput = styled.input`
+  width: 100%;
+  max-width: ${({ theme }) => theme.sizeUnit * 40}px;
+  font-size: ${({ theme }) => theme.fontSize}px;
+  color: ${({ theme }) => theme.colorText};
+  background: ${({ theme }) => theme.colorBgContainer};
+  border: 1px solid ${({ theme }) => theme.colorBorder};
+  border-radius: ${({ theme }) => theme.borderRadius}px;
+  padding: 2px ${({ theme }) => theme.sizeUnit * 1.5}px;
+`;
+
+const TabTimestamp = styled.span`
+  font-size: ${({ theme }) => theme.fontSizeSM}px;
+  color: ${({ theme }) => theme.colorTextQuaternary};
+  white-space: nowrap;
+  flex-shrink: 0;
+`;
+
+const ActionButtons = styled.div`
+  display: flex;
+  align-items: center;
+  gap: 2px;
+`;
+
+const ActionButton = styled.button`
+  background: none;
+  border: none;
+  padding: ${({ theme }) => theme.sizeUnit}px;
+  cursor: pointer;
+  color: ${({ theme }) => theme.colorTextSecondary};
+  opacity: 0;
+  transition: all ${({ theme }) => theme.motionDurationMid};
+  display: flex;
+  align-items: center;
+  justify-content: center;
+  border-radius: ${({ theme }) => theme.borderRadius}px;
+
+  &:hover,
+  &:focus-visible {
+    opacity: 1;
+    color: ${({ theme }) => theme.colorError};
+    background: ${({ theme }) => theme.colorErrorBg};
+  }
+`;
+
+const Divider = styled.div`
+  height: 1px;
+  background: ${({ theme }) => theme.colorBorderSecondary};
+  margin: ${({ theme }) => theme.sizeUnit}px 0;
+`;
+
+const EmptyState = styled.div`
+  padding: ${({ theme }) => theme.sizeUnit * 5}px
+    ${({ theme }) => theme.sizeUnit * 4}px;
+  text-align: center;
+  color: ${({ theme }) => theme.colorTextSecondary};
+  font-size: ${({ theme }) => theme.fontSizeSM}px;
+`;
+
+const MINUTE_SECONDS = 60;
+const HOUR_MINUTES = 60;
+const DAY_HOURS = 24;
+const WEEK_DAYS = 7;
+
+export const formatRelativeTime = (timestamp: number): string => {
+  const seconds = Math.floor((Date.now() - timestamp) / 1000);
+  if (seconds < MINUTE_SECONDS) {
+    return t('just now');
+  }
+  const minutes = Math.floor(seconds / MINUTE_SECONDS);
+  if (minutes < HOUR_MINUTES) {
+    return t('%sm', String(minutes));
+  }
+  const hours = Math.floor(minutes / HOUR_MINUTES);
+  if (hours < DAY_HOURS) {
+    return t('%sh', String(hours));
+  }
+  const days = Math.floor(hours / DAY_HOURS);
+  if (days < WEEK_DAYS) {
+    return t('%sd', String(days));
+  }
+  return new Date(timestamp).toLocaleDateString(undefined, {
+    month: 'short',
+    day: 'numeric',
+  });
+};
+
+interface ChatTabsMenuProps {
+  tabs: ChatTab[];
+  activeTabId: string;
+  onSelectTab: (tabId: string) => void;
+  onNewChat: () => void;
+  onDeleteTab: (tabId: string) => void;
+  onRenameTab: (tabId: string, name: string) => void;
+}
+
+export const ChatTabsMenu = ({
+  tabs,
+  activeTabId,
+  onSelectTab,
+  onNewChat,
+  onDeleteTab,
+  onRenameTab,
+}: ChatTabsMenuProps) => {
+  const [editingTabId, setEditingTabId] = useState<string | null>(null);
+  const [editingName, setEditingName] = useState('');
+
+  const startEditing = useCallback((event: ReactMouseEvent, tab: ChatTab) => {
+    event.stopPropagation();
+    setEditingTabId(tab.id);
+    setEditingName(tab.name);
+  }, []);
+
+  const cancelEditing = useCallback(() => {
+    setEditingTabId(null);
+    setEditingName('');
+  }, []);
+
+  const commitRename = useCallback(
+    (tabId: string) => {
+      const trimmedName = editingName.trim();
+      if (trimmedName) {
+        onRenameTab(tabId, trimmedName);
+      }
+      cancelEditing();
+    },
+    [cancelEditing, editingName, onRenameTab],
+  );
+
+  const menuContent = (
+    <MenuContainer data-test="chat-tabs-menu">
+      <MenuHeader>{t('Conversations')}</MenuHeader>
+      <NewChatButton type="button" onClick={onNewChat}>
+        <Icons.PlusOutlined iconSize="s" />
+        <span>{t('New Chat')}</span>
+      </NewChatButton>
+      <Divider />
+      {tabs.length === 0 ? (
+        <EmptyState>{t('No conversations yet')}</EmptyState>
+      ) : (
+        tabs.map(tab => (
+          <TabItem
+            key={tab.id}
+            isActive={tab.id === activeTabId}
+            onClick={() => onSelectTab(tab.id)}
+          >
+            <TabInfo>
+              <Icons.MessageOutlined iconSize="s" />
+              {editingTabId === tab.id ? (
+                <TabNameInput
+                  autoFocus
+                  value={editingName}
+                  onChange={event => setEditingName(event.target.value)}
+                  onClick={event => event.stopPropagation()}
+                  onBlur={() => commitRename(tab.id)}
+                  onKeyDown={event => {
+                    event.stopPropagation();
+                    if (event.key === 'Enter') {
+                      commitRename(tab.id);
+                    } else if (event.key === 'Escape') {
+                      cancelEditing();
+                    }
+                  }}
+                  aria-label={t('Conversation name')}
+                />
+              ) : (
+                <TabName>{tab.name}</TabName>
+              )}
+              {tab.updatedAt !== undefined && (
+                
<TabTimestamp>{formatRelativeTime(tab.updatedAt)}</TabTimestamp>
+              )}
+            </TabInfo>
+            <ActionButtons>
+              <ActionButton
+                type="button"
+                className="action-btn"
+                onClick={event => startEditing(event, tab)}
+                title={t('Rename conversation')}
+                aria-label={t('Rename conversation')}
+              >
+                <Icons.EditOutlined iconSize="s" />
+              </ActionButton>
+              {/* A conversation with messages is confirmed before deletion; an
+                  empty one is discarded without a prompt. */}
+              {tab.messages.length > 0 ? (

Review Comment:
   History tabs are created without loading their messages, so an unopened 
conversation has `messages.length === 0` even when it contains persisted 
messages. Deleting one bypasses the irreversible-delete confirmation. Could 
this use the list response's message count, or always confirm non-new threads?



##########
superset/ai/eventbus.py:
##########
@@ -0,0 +1,314 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+Carries streamed events from whatever produced them to the HTTP response.
+
+Two implementations, matching the two execution modes. Inline execution needs
+nothing more than an in-process queue. Worker execution needs a shared,
+*replayable* channel — replayable because a browser that loses its connection
+must be able to rejoin a run already in progress, which rules out
+publish/subscribe: a subscriber that was absent when an event was published
+never sees it.
+
+The Redis implementation therefore uses streams, and reuses the cache backend
+that Superset's async-query channel already configures rather than introducing
+a second Redis client to operate.
+"""
+
+from __future__ import annotations
+
+import logging
+import queue
+from abc import ABC, abstractmethod
+from collections.abc import Iterator
+from typing import Any
+
+from superset.ai.events import StreamEvent
+from superset.ai.types import StreamEventType
+from superset.utils import json
+
+logger = logging.getLogger(__name__)
+
+#: Yielded by :meth:`BaseEventBus.consume` when nothing arrived within the poll
+#: interval, so a caller can emit a keep-alive rather than block indefinitely.
+IDLE = None
+
+#: Terminal event types. Seeing one ends consumption, so a reader does not hang
+#: waiting for a producer that has already finished.
+_TERMINAL = frozenset(
+    {StreamEventType.DONE, StreamEventType.ERROR, StreamEventType.CANCELLED}
+)
+
+
+class BaseEventBus(ABC):
+    """A per-run channel of events."""
+
+    @abstractmethod
+    def publish(self, run_id: str, event: StreamEvent) -> None:
+        """Append an event to a run's channel."""
+
+    @abstractmethod
+    def consume(
+        self,
+        run_id: str,
+        timeout_seconds: float,
+        poll_seconds: float = 1.0,
+    ) -> Iterator[StreamEvent | None]:
+        """
+        Yield a run's events until a terminal one arrives or time runs out.
+
+        Yields :data:`IDLE` when a poll interval passes with nothing new, which
+        is the caller's cue to send a keep-alive frame.
+        """
+
+    @abstractmethod
+    def close(self, run_id: str) -> None:
+        """Release any resources held for a run."""
+
+
+class MemoryEventBus(BaseEventBus):
+    """
+    An in-process queue per run.
+
+    Correct only when the producer and the streaming request share a process.
+    Selecting this alongside worker execution would leave every stream silent,
+    which :func:`get_event_bus` refuses to allow.
+    """
+
+    def __init__(self) -> None:
+        self._queues: dict[str, queue.SimpleQueue[StreamEvent]] = {}
+
+    def _queue_for(self, run_id: str) -> queue.SimpleQueue[StreamEvent]:
+        return self._queues.setdefault(run_id, queue.SimpleQueue())
+
+    def publish(self, run_id: str, event: StreamEvent) -> None:
+        self._queue_for(run_id).put(event)
+
+    def consume(
+        self,
+        run_id: str,
+        timeout_seconds: float,
+        poll_seconds: float = 1.0,
+    ) -> Iterator[StreamEvent | None]:
+        import time
+
+        # Deliberately not ``_queue_for``: reading must not create a channel.
+        # This bus lives for the life of the process, so a client polling
+        # unknown run identifiers would otherwise grow the dict without bound.
+        channel = self._queues.get(run_id)
+        deadline = time.monotonic() + timeout_seconds
+
+        while True:
+            remaining = deadline - time.monotonic()
+            if remaining <= 0:
+                return
+            if channel is None:
+                # The producer may not have published yet; look again rather
+                # than deciding the run does not exist. Only report idle if it
+                # is still absent, so a channel that appeared during the wait
+                # is drained on this pass instead of costing an extra tick.
+                channel = self._queues.get(run_id)
+                if channel is None:
+                    yield IDLE
+                    time.sleep(min(poll_seconds, remaining))
+                continue
+            try:
+                # Bounded by whichever is sooner, so a generous poll interval
+                # cannot overshoot the caller's deadline.
+                event = channel.get(timeout=min(poll_seconds, remaining))
+            except queue.Empty:
+                yield IDLE
+                continue
+            yield event
+            if event.type in _TERMINAL:
+                return
+
+    def close(self, run_id: str) -> None:
+        self._queues.pop(run_id, None)
+
+
+class RedisStreamEventBus(BaseEventBus):
+    """
+    A Redis stream per run.
+
+    Replayable by construction: a reconnecting reader starts from the beginning
+    of the stream and catches up, which is what makes worker execution usable
+    from a browser on a flaky connection.
+    """
+
+    def __init__(
+        self,
+        cache: Any,
+        prefix: str = "ai-events-",
+        ttl_seconds: int = 900,
+    ) -> None:
+        self._cache = cache
+        self._prefix = prefix
+        self._ttl = ttl_seconds
+
+    def _stream(self, run_id: str) -> str:
+        return f"{self._prefix}{run_id}"
+
+    def publish(self, run_id: str, event: StreamEvent) -> None:
+        payload = {
+            "data": json.dumps({"type": event.type.value, "payload": 
event.payload})
+        }
+        # A failure to publish must not kill the run that is producing useful
+        # work; the reader will time out and the answer is still persisted.
+        try:
+            self._cache.xadd(self._stream(run_id), payload, "*", 10_000)

Review Comment:
   The TTL is only set by `close()`, so a worker-mode run whose client never 
opens the stream leaves its Redis key without an expiry. Repeated abandoned 
requests will accumulate `ai-events-*` streams indefinitely. Could publishing a 
terminal event set or refresh the configured TTL?



##########
docs/admin_docs/configuration/ai-assistant.mdx:
##########
@@ -0,0 +1,489 @@
+---
+title: AI Assistant
+hide_title: true
+sidebar_position: 17
+version: 1
+---
+
+# AI Assistant
+
+The AI Assistant is a conversational interface for exploring your data. A user
+asks a question in plain language; the assistant finds relevant datasets,
+inspects their schema, writes and runs read-only SQL, and answers with both the
+result and the query it used.
+
+Superset ships **no model provider and talks to no model vendor by default**.
+The feature is disabled, and even when enabled it returns `404` until you point
+it at a provider you control. Nothing is sent anywhere until you configure it.
+
+## Enabling it
+
+Two things are required: the feature flag, and a provider.
+
+```python
+# superset_config.py
+FEATURE_FLAGS = {
+    "AI_ASSISTANT": True,
+}
+
+AI_LLM_PROVIDER_CLASS = "superset.ai.llm.anthropic.AnthropicProvider"
+AI_LLM_PROVIDER_CONFIG = {
+    "api_key": os.environ["ANTHROPIC_API_KEY"],
+    "models": {
+        "default": "claude-sonnet-4-5",
+        "fast": "claude-haiku-4-5",
+        "reasoning": "claude-opus-4-1",
+    },
+}
+```
+
+Install the matching extra:
+
+```bash
+pip install "apache-superset[ai-anthropic]"   # or [ai-openai]
+```
+
+Then run `superset init` so the assistant's permissions are created and 
assigned
+to roles. Without this the endpoints return `403`.
+
+Conversations are stored in Superset's metadata database, so no extra
+infrastructure is needed for the default configuration.
+
+### Which roles get access
+
+`superset init` grants `can_read`/`can_write` on `AIAssistant` to **Admin** and
+**Alpha** only. "Write" here means writing one's own conversation — the
+assistant's tools are read-only and it cannot create or modify assets.
+
+**Gamma does not get it by default.** The assistant runs queries and costs
+money per question, so it is granted deliberately rather than inherited. To
+give it to Gamma users, add `can_read`/`can_write` on `AIAssistant` to Gamma or
+to a custom role.
+
+Every query the assistant runs is subject to the *user's own* database and
+dataset permissions. It cannot read anything the person chatting with it could
+not read themselves.
+
+Because it is not in Gamma, it is also not inherited by the Public role when
+`PUBLIC_ROLE_LIKE = "Gamma"` — an anonymous visitor cannot reach the assistant
+unless you grant it explicitly.
+
+## Choosing a provider
+
+`AI_LLM_PROVIDER_CLASS` is a dotted path to a
+`superset.ai.llm.base.BaseLLMProvider` subclass. Two are bundled:
+
+| Class | Use for |
+| --- | --- |
+| `superset.ai.llm.anthropic.AnthropicProvider` | The Anthropic Messages API |
+| `superset.ai.llm.openai_compatible.OpenAICompatibleProvider` | OpenAI, and 
anything exposing an OpenAI-compatible endpoint — vLLM, Ollama, a private 
gateway |
+
+`AI_LLM_PROVIDER_CONFIG` is passed to the provider's constructor and its
+contents are provider-defined. For the OpenAI-compatible provider, `base_url`
+points it anywhere:
+
+```python
+AI_LLM_PROVIDER_CLASS = 
"superset.ai.llm.openai_compatible.OpenAICompatibleProvider"
+AI_LLM_PROVIDER_CONFIG = {
+    "base_url": "https://llm.internal.example.com/v1";,
+    "api_key": os.environ["MY_GATEWAY_KEY"],
+    "models": {"default": "our-hosted-model"},
+}
+```
+
+Everything vendor-specific — URLs, authentication, model naming — lives in the
+provider. Superset core contains none of it, so a self-hosted model or a 
private
+gateway needs configuration rather than a fork.
+
+### Model tiers and selection
+
+Profiles and prompts refer to capability *tiers* (`default`, `fast`,
+`reasoning`), never to a vendor's model names. The provider maps tiers to
+concrete models via the `models` dict. A tier you do not configure is an error
+when requested, never a silent substitution — so cost and answer quality stay
+attributable to the model actually used.
+
+Users may also pin a specific model per turn. Only models present in your
+`models` mapping are accepted; anything else is rejected.
+
+## Agent profiles
+
+A profile bundles the decisions that differ between a quick answer and a 
careful
+investigation: which tools are available, which model tier, and how many steps.
+Two ship by default — `default` and `analyst`.
+
+**Which tools a model may invoke is a decision each deployment makes**, so
+profiles are fully configurable. `AI_AGENT_PROFILES` maps a profile key to the
+fields you want to override, leaving the rest alone:
+
+```python
+AI_AGENT_PROFILES = {
+    # Let the assistant search and inspect, but never run SQL.
+    "default": {"tools": ["search_assets", "list_databases", "get_schema"]},
+
+    # Let the analyst profile think harder and longer.
+    "analyst": {"model_alias": "reasoning", "max_turns": 60},
+
+    # Add a profile only some users may select.
+    "deep": {
+        "name": "Deep analysis",
+        "description": "Slow, thorough, multi-step.",
+        "tools": ["search_assets", "get_schema", "execute_sql"],
+        "required_permission": ("can_write", "AIAssistant"),
+    },
+}
+```
+
+A tool name that does not exist is an error naming the typo and listing the
+valid names, rather than an assistant that quietly lacks a capability. An empty
+`tools` list is valid and means conversation with no data access.
+
+`required_permission` is enforced on both the listing *and* the run path, so a
+profile a user cannot see is also one they cannot invoke by posting its key.
+
+### Available tools
+
+| Tool | What it does |
+| --- | --- |
+| `search_assets` | Finds datasets, charts and dashboards the user can see |
+| `list_databases` | Lists database connections exposed to SQL Lab |
+| `get_schema` | Lists schemas, tables and columns |
+| `execute_sql` | Runs a **read-only** query |
+| `validate_sql` | Checks a query without running it |
+| `get_chart_context` | Reads a chart's definition |
+| `get_dashboard_context` | Reads a dashboard's definition |
+
+## Customising the prompt
+
+Three levers, in increasing order of bluntness.
+
+**Add to it.** `AI_EXTRA_PROMPT_SECTIONS` appends your own sections. This is
+where deployment-specific knowledge belongs — your table conventions, your
+warehouse's dialect quirks, how your business defines a metric. The shipped
+prompt is deliberately generic and mentions no particular database engine.
+
+**Remove from it.** `AI_DISABLED_PROMPT_SECTIONS` drops a shipped section by
+key, for when you disagree with one. The safety section cannot be disabled.
+
+**Replace it.** `AI_SYSTEM_PROMPT` substitutes the whole thing.
+
+:::warning
+Setting `AI_SYSTEM_PROMPT` discards the shipped safety and prompt-injection
+rules along with everything else. Your deployment then owns them.
+:::
+
+`AI_SYSTEM_PROMPT_MUTATOR` is a last-mile callable applied after assembly,
+mirroring `SQL_QUERY_MUTATOR`.
+
+## Where turns execute
+
+`AI_ASSISTANT_EXECUTION_MODE` decides where the work happens.
+
+**`"inline"`** (default) runs the turn in the web process. Nothing extra to
+deploy.
+
+**`"worker"`** hands it to Celery. Web workers stay free, and a browser that
+loses its connection can rejoin a run in progress. It requires Celery and a
+Redis event bus:
+
+```python
+AI_ASSISTANT_EXECUTION_MODE = "worker"
+AI_ASSISTANT_EVENT_BUS = "redis"
+AI_ASSISTANT_EVENT_BUS_CACHE_CONFIG = {
+    "CACHE_TYPE": "RedisCache",
+    "CACHE_REDIS_HOST": "redis",
+    "CACHE_REDIS_PORT": 6379,
+    "CACHE_REDIS_DB": 0,
+}
+
+class CeleryConfig:
+    imports = (
+        # ... your existing imports ...
+        "superset.ai.tasks",
+    )
+```
+
+Streams need Redis commands the general-purpose cache client does not expose,
+which is why the bus is configured separately rather than reusing 
`CACHE_CONFIG`.
+
+Selecting `"worker"` with the in-memory event bus raises rather than leaving
+every stream silently empty, and so does selecting the Redis bus without a
+usable connection.
+
+A turn is deliberately **not** retried after a worker crash: inference costs
+money, and re-running a turn the user may already have partly seen would charge
+twice. The message records that it failed and the user can ask again.
+
+## Safety and limits
+
+Guards are applied before any tool runs, configured via
+`AI_AGENT_TOOL_POLICIES`:
+
+- **Read-only SQL.** Enforced using Superset's own SQL parser, not pattern
+  matching — so a write hidden behind a comment, a CTE, a second statement, or
+  an unparseable construct is refused. `EXPLAIN`, `SHOW` and `DESCRIBE` are
+  permitted; everything the parser cannot vouch for is not.
+- **Identifier safety.** Table and column names are resolved against metadata
+  the user may see rather than interpolated into SQL.
+
+These bound blast radius; they do not replace authorization. Every tool that
+touches a data-bearing object performs the same permission check the REST API
+does.
+
+Result sizes are capped by `AI_AGENT_MAX_RESULT_ROWS` and
+`AI_AGENT_MAX_RESULT_BYTES`, and truncation is reported rather than hidden. 
Turn
+length is bounded by `AI_AGENT_MAX_TURNS` and `AI_AGENT_TIMEOUT_SECONDS`; a run
+that exhausts either answers with what it has.
+
+Content that arrives from your warehouse or asset metadata — table comments,
+chart titles, column labels — is marked as untrusted in the prompt, because a
+value in a database is data and not an instruction.
+
+### Cancellation
+
+Cancellation is cooperative: a run stops at its next step boundary. A run 
inside
+a single long model call or a single long query will not stop until that call
+returns.
+
+## Monitoring and tracing
+
+Superset bundles **no integration with any AI monitoring product**. Instead it
+exposes a small sink interface, `AITelemetry`, and calls it once per run, once
+per model round trip and once per tool call. Whatever you already use —
+Braintrust, LangSmith, Langfuse, Arize Phoenix, an OpenTelemetry collector, a
+self-hosted alternative, or a table in your own warehouse — you connect by
+implementing that interface and listing it in `AI_TELEMETRY`.
+
+Entries are instances or dotted paths, exactly as for `EVENT_LOGGER` and
+`STATS_LOGGER`. Two sinks ship in-tree and depend on nothing external:
+
+```python
+# superset_config.py
+import logging
+
+from superset.ai.telemetry import LoggingAITelemetry, StatsLoggerAITelemetry
+
+AI_TELEMETRY = [
+    # One structured line per span, at the level you choose.
+    LoggingAITelemetry(level=logging.INFO),
+    # Counters and timings through your configured STATS_LOGGER.
+    StatsLoggerAITelemetry(),
+]
+```
+
+`StatsLoggerAITelemetry` emits under a `superset.ai.` prefix: `run.start`,
+`run.end`, `run.outcome.<outcome>`, `run.duration_ms`, `run.turns`,
+`run.tokens.input`, `run.tokens.output`, `model_call`,
+`model_call.duration_ms`, `model_call.error`, `error`, and per tool
+`tool_call.<tool>`, `tool_call.<tool>.duration_ms`, `tool_call.<tool>.error`
+and `tool_call.<tool>.truncated`. User, run and thread identifiers deliberately
+never appear in a metric name — a metric per user is how a metrics backend gets
+brought down. That detail belongs in a trace, which is what a custom sink is
+for.
+
+### The content trade-off
+
+`AI_TELEMETRY_REDACT_CONTENT` defaults to `True`, and telemetry then carries
+**structure and measurements only**: durations, token counts, model names, tool
+names, outcomes, error classes, and the run, thread and user identifiers. No
+question, no answer, no SQL, no row of data. Redaction is applied where the
+trace is built, so a sink cannot receive content by accident even if it looks
+for it.
+
+Setting it to `False` is what makes a trace genuinely useful for debugging
+answer quality — you can read the prompt that produced a wrong answer and the
+statement it ran. It also means the text of business questions and values from
+your warehouse leave Superset for whichever service your sinks talk to. In many
+organisations that is a decision for someone other than the person editing the
+config file. `AI_TELEMETRY_MAX_CONTENT_CHARS` (default 10,000) caps any single
+content field so one large result cannot dominate a payload.
+
+### A custom sink
+
+Every method has a no-op default, so implement only the ones you need — a sink
+that only wants token counts overrides `on_model_call` and nothing else.
+
+```python
+from superset.ai.telemetry import AITelemetry, ModelCallTrace, RunTrace
+
+
+class TracingServiceTelemetry(AITelemetry):
+    """Forwards runs to an external tracing service."""
+
+    def __init__(self, client):
+        self._client = client
+
+    def on_run_start(self, run: RunTrace) -> None:
+        self._client.start_span(run.run_id, name="superset.ai.run", 
attributes={
+            "thread": run.thread_uuid,
+            "user": run.user_id,
+        })
+
+    def on_model_call(self, run: RunTrace, call: ModelCallTrace) -> None:
+        self._client.event(run.run_id, "model_call", {
+            "turn": call.turn,
+            "model": call.model,
+            "input_tokens": call.input_tokens,
+            "output_tokens": call.output_tokens,
+            # None unless you have turned redaction off.
+            "prompt": call.system_prompt,
+        })
+
+    def on_run_end(self, run: RunTrace) -> None:
+        self._client.end_span(run.run_id, status=str(run.outcome), attributes={
+            "duration_ms": run.duration_ms,
+            "turns": run.turns,
+            "usage": run.usage,
+        })
+
+
+AI_TELEMETRY = [TracingServiceTelemetry(client=my_tracing_client)]
+```
+
+Three things to know before you write one:
+
+- **Sinks are called on the thread answering the user.** Anything that makes a
+  network call should hand off to a queue or a background thread; otherwise a
+  slow monitoring backend becomes slow answers.
+- **A sink that raises cannot break a run.** Failures are logged once and
+  ignored, and the other configured sinks still receive everything. The same
+  applies to a dotted path that will not import: it is skipped with a warning
+  rather than taking the assistant down, because a missing observer loses the
+  record of a run and not the run itself.
+- **`agent_key`, `model` and `question` are resolved after the run starts**, so
+  a `RunTrace` passed to `on_run_start` may carry less than the one passed to
+  the later hooks. Read those on `on_run_end`.
+
+## Connecting your own MCP servers
+
+The assistant's built-in tools cover Superset itself. To let it reach anything
+else — your data catalog, a metrics service, a ticketing system — attach an
+[MCP](https://modelcontextprotocol.io) server. Superset bundles no third-party
+integration and connects to nothing by default; you name the servers.
+
+```bash
+pip install "apache-superset[ai-mcp]"
+```
+
+```python
+AI_AGENT_MCP_SERVERS = {
+    "acme_catalog": {
+        "url": "https://mcp.acme.internal/mcp";,
+        "transport": "streamable_http",       # or "sse"
+        "headers": {"Authorization": f"Bearer {os.environ['ACME_MCP_TOKEN']}"},
+        "timeout_seconds": 30,
+        "tool_allowlist": ["search_tables"],  # omit to offer every tool
+    },
+}
+
+# Then let a profile use it.
+AI_AGENT_PROFILES = {
+    "default": {"mcp_servers": ["acme_catalog"]},
+}
+```
+
+Its tools appear to the model as `mcp__acme_catalog__search_tables`. The
+namespace means a foreign tool can never shadow a built-in one, and it is the
+name to use in `tool_allowlist` and `tool_denylist`.
+
+### What Superset does to keep a foreign server contained
+
+A third-party server is untrusted input, and possibly untrusted intent:
+
+- **Everything it returns is marked as untrusted** before the model sees it, so
+  text in a tool result is treated as data rather than instructions. Tool
+  *descriptions* get the same treatment, since they enter the prompt every 
turn.
+- **No Superset credential is ever forwarded.** Only the headers you configured
+  for that server are sent — never the user's session cookie, CSRF token, or an
+  inbound authorization header.
+- **SQL execution through a foreign server is refused by default.** Superset's
+  read-only enforcement and per-dataset authorization cannot apply to a query
+  another system runs, so allowing it would silently bypass both. Set
+  `AI_AGENT_MCP_DENY_FOREIGN_SQL = False` to accept that trade deliberately.
+- **Results obey the same size cap** as built-in tools, and the cap is applied
+  while reading, so a hostile server cannot exhaust memory before truncation.
+- **A server being down does not break the assistant.** Discovery failure means
+  that server contributes no tools for the turn; the built-ins keep working.
+
+A profile naming a server you have not configured is an error, because a typo
+there is indistinguishable at runtime from an agent that has quietly lost a
+capability. Note that discovery happens per turn, so a slow server adds its
+latency to every turn that uses it.
+
+## Retention
+
+Conversations are kept for `AI_ASSISTANT_MESSAGE_RETENTION_DAYS` (default 30).
+Pruning is not automatic — schedule it if you want it enforced.
+
+## Trying it locally
+
+The development `docker compose` stack can bring the assistant up against the
+example data. Put the settings in `docker/.env-local`, which is untracked:
+
+```bash
+# docker/.env-local
+
+# Point at any OpenAI-compatible endpoint, including a private gateway.
+SUPERSET_AI_LLM_BASE_URL=https://your-gateway/v1
+SUPERSET_AI_LLM_API_KEY=your-token
+SUPERSET_AI_MODEL_DEFAULT=your-model-name
+```
+
+```bash
+docker compose up
+```
+
+The assistant appears once both a URL and a key are present; with neither set 
the
+stack behaves exactly as it did before. The model providers are optional 
extras,
+so add whichever one you need to `docker/requirements-local.txt`:
+
+```
+openai>=1.60.0,<2
+```
+
+The local stack also logs traces to the container output and, unlike the
+production default, includes prompts and SQL in them.
+
+## Full configuration reference
+
+| Setting | Default | Purpose |
+| --- | --- | --- |
+| `AI_LLM_PROVIDER_CLASS` | `None` | Provider class path. Unset means the 
feature is off. |

Review Comment:
   The configuration reference omits the `AI_SUGGESTED_PROMPTS_*` settings even 
though enabling them triggers configurable model calls. Operators cannot 
discover how to enable or bound this feature and will remain on fallback 
prompts. Could the reference include those settings and their opt-in/cost 
behavior?



##########
superset/ai/tools/base.py:
##########
@@ -0,0 +1,612 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+The tool contract and the registry that dispatches to it.
+
+A tool is a small, self-authorizing unit of work. "Self-authorizing" is the
+important half: the registry does not check permissions on a tool's behalf, and
+neither does the policy chain in :mod:`superset.ai.policy`, which answers the
+coarser question of whether a *shape* of call should be attempted at all. Any
+tool that returns or mutates a data-bearing object performs its own
+``security_manager`` check, because that is the only place with enough context 
to
+know which object is being touched.
+
+A tool returns two things. :attr:`ToolOutput.content` is what the model reads.
+:attr:`ToolOutput.display` is a summary for the UI, so a user can expand what 
the
+assistant did and see the SQL it ran and the rows it got back. Both are
+size-bounded here rather than in each tool: ``display`` is persisted on the
+message and shipped to the browser, so an unbounded one would be a second way 
to
+blow up a response.
+"""
+
+from __future__ import annotations
+
+import logging
+import time
+from abc import ABC, abstractmethod
+from dataclasses import dataclass, field
+from typing import Any, ClassVar
+
+from superset.ai.llm.base import ToolCall, ToolDefinition, ToolResult
+from superset.mcp_service.utils.sanitization import (
+    LLM_CONTEXT_CLOSE_DELIMITER,
+    LLM_CONTEXT_ESCAPED_CLOSE_DELIMITER,
+    LLM_CONTEXT_ESCAPED_OPEN_DELIMITER,
+    LLM_CONTEXT_OPEN_DELIMITER,
+)
+from superset.utils import json
+
+logger = logging.getLogger(__name__)
+
+#: Keys added to a payload whose budget was exceeded. Phrased for the model: it
+#: says what was lost and what to do differently, because a model told only
+#: "truncated" reissues the identical call.
+TRUNCATION_KEY = "_truncated"
+TRUNCATION_NOTE_KEY = "_truncation_note"
+
+#: Share of the response budget the UI summary may use. The model's copy is the
+#: one that has to be complete enough to reason over; ``display`` only has to 
be
+#: enough to render, and it is persisted, so it gets the smaller share.
+_DISPLAY_BUDGET_FRACTION = 0.5
+
+#: Fallback response budget for use outside an application context, matching 
the
+#: shipped ``AI_AGENT_MAX_RESULT_BYTES`` default.
+_DEFAULT_MAX_BYTES = 256 * 1024
+
+
+class ToolError(Exception):
+    """
+    A failure that should be shown to the model rather than raised at the user.
+
+    Tools raise this for conditions the model can act on — a database it may 
not
+    read, a column that does not resolve, SQL that will not parse. The registry
+    turns it into a :class:`~superset.ai.llm.base.ToolResult` with
+    ``is_error=True`` so the turn continues and the model can correct itself.
+
+    The message is model-visible, so it must never carry a driver exception, a
+    connection string, or a stack trace.
+    """
+
+
+@dataclass(frozen=True)
+class ToolOutput:
+    """
+    What a tool returns.
+
+    Build one with :meth:`of` rather than by hand, so that ``content`` and
+    ``payload`` cannot disagree.
+    """
+
+    #: Model-facing text. Becomes :attr:`ToolResult.content`.
+    content: str
+
+    #: JSON-serialisable summary for the UI, or ``None`` when there is nothing
+    #: worth rendering. Must never carry credentials, a connection string, or a
+    #: full result set.
+    display: dict[str, Any] | None = None
+
+    #: The structure ``content`` was serialised from. Retained so the registry
+    #: can shrink an oversized result by dropping rows rather than cutting JSON
+    #: mid-token. ``None`` when the tool supplied text directly.
+    payload: Any = None
+
+    @classmethod
+    def of(
+        cls,
+        payload: Any,
+        display: dict[str, Any] | None = None,
+    ) -> ToolOutput:
+        """Serialise ``payload`` as the model-facing content."""
+        return cls(
+            content=json.dumps(payload, default=str),
+            display=display,
+            payload=payload,
+        )
+
+
+@dataclass
+class ToolInvocation:
+    """
+    One completed dispatch, with everything a caller might need.
+
+    Exists because two consumers want different things from the same call: the
+    provider needs a :class:`~superset.ai.llm.base.ToolResult`, while the event
+    stream and message persistence want the UI summary and the timing.
+    """
+
+    call_id: str
+    tool_name: str
+    result: ToolResult
+    display: dict[str, Any] | None = None
+    duration_ms: int = 0
+    truncated: bool = False
+    arguments: dict[str, Any] = field(default_factory=dict)
+    error_type: str | None = None
+
+    @property
+    def is_error(self) -> bool:
+        """Whether the call failed."""
+        return self.result.is_error
+
+    def to_tool_result(self) -> ToolResult:
+        """The provider-neutral result to feed back to the model."""
+        return self.result
+
+
+class AITool(ABC):
+    """
+    One capability offered to the model.
+
+    Subclasses set :attr:`name`, :attr:`description` and :attr:`input_schema`,
+    and implement :meth:`run`. Everything else — size capping, error
+    translation, timing — is the registry's job.
+    """
+
+    #: Stable identifier the model calls, and the key an operator types when
+    #: configuring which tools an agent profile may use. Renaming one is a
+    #: breaking change: it appears in stored conversation history and in
+    #: deployment configuration.
+    name: ClassVar[str] = ""
+
+    #: Shown to the model verbatim. This is the tool's entire user manual, so 
it
+    #: should say when to reach for the tool and what it returns, not merely
+    #: what it is called.
+    description: ClassVar[str] = ""
+
+    #: JSON Schema for the arguments object. Providers translate it into
+    #: whatever their API expects.
+    input_schema: ClassVar[dict[str, Any]] = {"type": "object", "properties": 
{}}
+
+    @abstractmethod
+    def run(self, **kwargs: Any) -> ToolOutput:
+        """
+        Perform the work.
+
+        Raise :class:`ToolError` for anything the model should see and be able
+        to recover from. Any other exception is treated as a defect: it is
+        logged with a traceback and reported to the model as a generic failure,
+        so that an unexpected driver error cannot leak its message.
+        """
+
+    def definition(self) -> ToolDefinition:
+        """Provider-neutral description of this tool."""
+        return ToolDefinition(
+            name=self.name,
+            description=self.description,
+            input_schema=self.input_schema,
+        )
+
+
+def truncate_payload(payload: Any, max_bytes: int) -> tuple[str, bool]:
+    """
+    Serialise ``payload`` and bound it to ``max_bytes``.
+
+    Returns the JSON text and whether anything was dropped. Truncation is
+    applied to the largest list in a mapping payload — result rows in practice 
—
+    because halving a row count is comprehensible to the model whereas cutting 
a
+    JSON string mid-token is not. When there is no list to shrink, the text is
+    cut and the marker says so.
+    """
+    text = json.dumps(payload, default=str)
+    if len(text.encode("utf-8")) <= max_bytes:
+        return text, False
+
+    if isinstance(payload, dict):
+        list_keys = [
+            key for key, value in payload.items() if isinstance(value, list) 
and value
+        ]
+        if list_keys:
+            # Shrink the longest list first; it is the one carrying the bulk.
+            key = max(list_keys, key=lambda item: len(payload[item]))
+            rows = payload[key]
+            kept = len(rows)
+            # Halve until it fits, always leaving one element so the model can
+            # still see the shape of what it asked for.
+            while kept > 1:
+                kept //= 2
+                candidate = dict(payload)
+                candidate[key] = rows[:kept]
+                candidate[TRUNCATION_KEY] = True
+                candidate[TRUNCATION_NOTE_KEY] = (
+                    f"Returned {kept} of {len(rows)} {key} — the full result "
+                    f"exceeded the {max_bytes} byte response budget. Narrow 
the "
+                    f"request (fewer columns, a tighter filter, a smaller 
limit) "
+                    f"to see the rest."
+                )
+                text = json.dumps(candidate, default=str)
+                if len(text.encode("utf-8")) <= max_bytes:
+                    return text, True
+
+    # Nothing structural to shrink: cut the text and say so plainly.
+    note = f'…[truncated to {max_bytes} bytes]"}}'
+    keep = max(0, max_bytes - len(note.encode("utf-8")))
+    return text.encode("utf-8")[:keep].decode("utf-8", "ignore") + note, True
+
+
+def strip_prompt_framing(value: Any) -> Any:
+    """
+    Remove the model-facing untrusted-content framing from a value.
+
+    ``superset.mcp_service`` wraps user-authored text in 
``<UNTRUSTED-CONTENT>``
+    delimiters so a model can tell data from instruction. That framing is
+    meaningless to a person and appeared verbatim in the panel's tool log, so 
it
+    is removed on the way to the browser — and only there. The model's copy 
keeps
+    the delimiters, which is the whole point of them.
+
+    Text that was escaped because the author had literally typed a delimiter is
+    restored, since for a reader that literal text is the honest rendering.
+    """
+    if isinstance(value, str):
+        return (
+            value.replace(LLM_CONTEXT_OPEN_DELIMITER, "")
+            .replace(LLM_CONTEXT_CLOSE_DELIMITER, "")
+            .replace(LLM_CONTEXT_ESCAPED_OPEN_DELIMITER, 
LLM_CONTEXT_OPEN_DELIMITER)
+            .replace(LLM_CONTEXT_ESCAPED_CLOSE_DELIMITER, 
LLM_CONTEXT_CLOSE_DELIMITER)
+            .strip()
+        )
+    if isinstance(value, dict):
+        return {
+            strip_prompt_framing(key): strip_prompt_framing(nested)
+            for key, nested in value.items()
+        }
+    if isinstance(value, list):
+        return [strip_prompt_framing(item) for item in value]
+    return value
+
+
+def bound_display(
+    display: dict[str, Any] | None,
+    max_bytes: int,
+) -> dict[str, Any] | None:
+    """
+    Bound the UI summary to ``max_bytes``.
+
+    The summary is persisted on the message and sent to the browser, so it 
needs
+    its own ceiling rather than riding on the model-facing budget. Reuses the
+    row-dropping strategy so a sample stays a valid sample.
+    """
+    if display is None:
+        return None
+    # Stripped before bounding so the ceiling applies to what is actually sent,
+    # and so a payload is not spent on framing that gets removed anyway.
+    display = strip_prompt_framing(display)
+    text, truncated = truncate_payload(display, max_bytes)
+    if not truncated:
+        return display
+    try:
+        bounded = json.loads(text)
+    except Exception:  # pylint: disable=broad-except
+        # The text form was cut mid-structure and will not parse. A summary is
+        # a nicety, so drop it rather than shipping something unrenderable.
+        logger.info("Dropping an oversized tool display payload")
+        return {TRUNCATION_KEY: True, TRUNCATION_NOTE_KEY: "Summary was too 
large."}
+    return bounded if isinstance(bounded, dict) else None
+
+
+class ToolRegistry:
+    """
+    Holds tools by name, exports their definitions, and dispatches calls.
+
+    Registration is explicit — a caller constructs the tools it wants and hands
+    them over. There is no discovery, no entry points and no dynamic import: 
the
+    set of capabilities a deployment exposes to a model should be readable in 
one
+    place rather than assembled by import side effects.
+    """
+
+    def __init__(self, tools: list[AITool] | None = None) -> None:
+        self._tools: dict[str, AITool] = {}
+        for tool in tools or []:
+            self.register(tool)
+
+    def register(self, tool: AITool) -> None:
+        """
+        Add a tool.
+
+        A duplicate name is an error rather than an overwrite: silently
+        replacing a tool would change what the model can do depending on
+        registration order.
+        """
+        if not tool.name:
+            raise ValueError(f"{type(tool).__name__} has no name")
+        if tool.name in self._tools:
+            raise ValueError(f"Tool {tool.name!r} is already registered")
+        self._tools[tool.name] = tool
+
+    def __contains__(self, name: object) -> bool:
+        return name in self._tools
+
+    def __len__(self) -> int:
+        return len(self._tools)
+
+    def names(self) -> list[str]:
+        """Registered tool names, in registration order."""
+        return list(self._tools)
+
+    def get(self, name: str) -> AITool | None:
+        """The tool registered under ``name``, or ``None``."""
+        return self._tools.get(name)
+
+    def definitions(self) -> list[ToolDefinition]:
+        """Every tool's definition, for the provider's tool list."""
+        return [tool.definition() for tool in self._tools.values()]
+
+    def subset(self, names: list[str] | tuple[str, ...] | set[str]) -> 
ToolRegistry:
+        """
+        A registry containing only the named tools.
+
+        This is how a deployment restricts what one agent profile may do. An
+        unknown name raises rather than being skipped: a typo in configuration
+        that silently removed a capability would look like a model that had
+        simply stopped using the tool, which is close to undebuggable.
+        """
+        requested = list(names)
+        unknown = [name for name in requested if name not in self._tools]
+        if unknown:
+            raise ValueError(
+                f"Unknown tool(s) {', '.join(sorted(unknown))}. "
+                f"Available: {', '.join(sorted(self._tools))}."
+            )
+        # Registration order is preserved rather than the caller's, so two
+        # profiles listing the same tools present them to the model 
identically.
+        return ToolRegistry(
+            [tool for name, tool in self._tools.items() if name in 
set(requested)]
+        )
+
+    def invoke(
+        self,
+        call: ToolCall,
+        max_bytes: int | None = None,
+    ) -> ToolInvocation:
+        """
+        Run one tool call and return it with its UI summary and timing.
+
+        Never raises for a tool-level failure: an unknown name, a refused call
+        and a crash all come back as ``is_error=True`` results so the model can
+        react within the same turn. The alternative — propagating — ends the 
turn
+        on a mistake the model could have corrected.
+        """
+        started = time.monotonic()
+
+        def elapsed() -> int:
+            return int((time.monotonic() - started) * 1000)
+
+        def failure(message: str, error_type: str) -> ToolInvocation:
+            return ToolInvocation(
+                call_id=call.id,
+                tool_name=call.name,
+                result=ToolResult(call_id=call.id, content=message, 
is_error=True),
+                duration_ms=elapsed(),
+                arguments=dict(call.arguments),
+                error_type=error_type,
+            )
+
+        tool = self._tools.get(call.name)
+        if tool is None:
+            available = ", ".join(sorted(self._tools)) or "none"
+            return failure(
+                f"No tool named {call.name!r}. Available tools: {available}.",
+                "ToolUnavailable",
+            )
+
+        # Every tool here reads a permissioned resource, so all of them need a
+        # principal to check against. Gated once, here, rather than in each 
tool:
+        # a tool that forgot would otherwise fall through to whatever the
+        # security manager does with an absent user, which is not a decision
+        # worth leaving to chance.
+        if _current_user() is None:
+            return failure(
+                f"{call.name} requires an authenticated user and none is "
+                f"available for this request.",
+                "AuthenticationRequired",
+            )
+
+        if max_bytes is None:
+            max_bytes = _configured_max_bytes()
+
+        try:

Review Comment:
   The telemetry test hand-constructs an invocation with 
`DetachedInstanceError`, while the registry test only asserts the base 
`ToolError`. A regression that drops or hard-codes the generic-exception 
`error_type` would still pass. Could the test drive a concrete subclass through 
`ToolRegistry` and assert the telemetry span receives that class?



##########
superset/ai/runtime/messages.py:
##########
@@ -0,0 +1,582 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+The default runtime: a plain tool-use loop over the provider's message API.
+
+Chosen as the default because it needs nothing beyond an HTTP call — no agent
+engine subprocess, no working directory, no bundled binary — so it works with
+whatever provider a deployment configures.
+"""
+
+from __future__ import annotations
+
+import logging
+import time
+from collections.abc import AsyncIterator
+from typing import Any
+
+from superset.ai.events import (
+    assistant_delta_event,
+    checkpoint_event,
+    error_event,
+    final_event,
+    GENERIC_ERROR_MESSAGE,
+    StreamEvent,
+    thinking_event,
+    thoughts_event,
+)
+from superset.ai.llm.base import (
+    CompletionRequest,
+    LLMError,
+    LLMResponse,
+    Message,
+    StreamEventKind,
+    ToolCall,
+    ToolResult,
+)
+from superset.ai.runtime.base import BaseAgentRuntime, RunRequest, RunResult
+from superset.ai.telemetry import (
+    current_run,
+    POLICY_DENIED,
+    RunRecorder,
+    TOOL_UNAVAILABLE,
+)
+from superset.ai.types import MessageRole, ProgressStage, TokenUsage
+
+logger = logging.getLogger(__name__)
+
+#: How much of a tool's output is kept on the persisted message. The model
+#: still sees the whole thing; this is the audit copy.
+_RECORDED_OUTPUT_LIMIT = 2_000
+
+#: Size of the chunks the finished answer is delivered in.
+_DELIVERY_CHUNK_SIZE = 512
+
+#: How much reasoning is kept on the result. Reasoning can run several times
+#: longer than the answer, and this is persisted next to it.
+_RECORDED_THOUGHTS_LIMIT = 8_000
+
+_NO_ANSWER = (
+    "I wasn't able to reach an answer for that. Try narrowing the question, "
+    "or naming the dataset you have in mind."
+)
+
+
+class MessagesApiRuntime(BaseAgentRuntime):
+    """
+    Alternates model calls and tool calls until the model stops asking.
+
+    Two behaviours are worth understanding before changing this class.
+
+    First, prose the model emits *before* a tool call is treated as reasoning,
+    not answer: it becomes a ``thoughts`` event and is dropped from the answer.
+    A model narrating "the orders table looks right, let me check" is stating a
+    hypothesis it may abandon, and appending that to the answer produces a
+    reply that contradicts itself.
+
+    Second, the loop always terminates and never raises for an operational
+    failure. By the time it runs, response headers have been flushed and an
+    exception can no longer become an HTTP status, so every failure is an 
event.
+    """
+
+    def __init__(self, provider: Any) -> None:
+        super().__init__(provider)
+        self._result = RunResult()
+        #: Set when the model signals it has finished answering.
+        self._finished = False
+        #: The most recent round trip's response, or ``None`` if it failed. The
+        #: turn methods are generators and cannot return a value.
+        self._last_response: LLMResponse | None = None
+        #: Whether any answer text has already been sent as it was generated. 
The
+        #: finished answer is only replayed in chunks when it has not.
+        self._streamed_text = False
+
+    @property
+    def result(self) -> RunResult:
+        return self._result
+
+    async def run(self, request: RunRequest) -> AsyncIterator[StreamEvent]:
+        self._result = RunResult()
+        self._finished = False
+        self._last_response = None
+        self._streamed_text = False
+        answer_parts: list[str] = []
+
+        yield thinking_event(ProgressStage.START, "Working on your question")
+
+        # The provider's connection pool belongs to the loop this run is driven
+        # on, and the caller closes that loop as soon as the run ends. Closing
+        # here — inside the loop, however the run finishes, including when the
+        # generator is abandoned mid-way by a user pressing stop — is what 
keeps
+        # a client from being finalised against a dead loop.
+        try:
+            async for event in self._turn_loop(request, answer_parts):
+                yield event
+
+            # A run that failed or was abandoned has already said so; emitting 
an
+            # answer as well would contradict it.
+            if self._result.error is not None or self._result.cancelled:
+                return
+
+            answer = "\n\n".join(part for part in answer_parts if part).strip()
+            self._result.answer = answer or _NO_ANSWER
+
+            # Only replayed when nothing was streamed — a provider without
+            # streaming support still gets to deliver its answer progressively.
+            # Replaying after live text would show the answer twice.
+            if not self._streamed_text:
+                for chunk in _chunk(self._result.answer):
+                    yield assistant_delta_event(chunk)
+            yield final_event(self._result.answer)
+        finally:
+            await self.provider.aclose()
+
+    async def _turn_loop(
+        self,
+        request: RunRequest,
+        answer_parts: list[str],
+    ) -> AsyncIterator[StreamEvent]:
+        """
+        Alternate model and tool calls until the model stops or a budget runs 
out.
+
+        Appends to ``answer_parts`` rather than returning the answer, because 
an
+        async generator cannot both yield events and return a value.
+        """
+        deadline = time.monotonic() + request.timeout_seconds
+        conversation = list(request.messages)
+
+        for turn in range(1, request.max_turns + 1):
+            self._result.turns = turn
+
+            if self._should_stop(request, deadline):
+                if self._result.timed_out:
+                    yield thinking_event(
+                        ProgressStage.FALLBACK,
+                        "Taking longer than expected — answering with what I 
have",
+                    )
+                return
+
+            async for event in self._safe_turn(request, conversation, turn):
+                yield event
+            response = self._last_response
+            if response is None:
+                yield error_event()
+                return
+
+            async for event in self._consume(
+                request, response, conversation, answer_parts
+            ):
+                yield event
+
+            if self._finished or self._result.cancelled:
+                return
+
+        # Budget exhausted without the model choosing to stop.
+        yield thinking_event(
+            ProgressStage.FALLBACK,
+            "Reached the step limit — answering with what I have",
+        )
+
+    async def _consume(
+        self,
+        request: RunRequest,
+        response: LLMResponse,
+        conversation: list[Message],
+        answer_parts: list[str],
+    ) -> AsyncIterator[StreamEvent]:
+        """Act on one model response, running any tools it asked for."""
+        if response.thinking:
+            self._record_thoughts(response.thinking)
+            yield thoughts_event(response.thinking)
+
+        if not response.wants_tools:
+            self._finished = True
+            if response.text:
+                answer_parts.append(response.text)
+                # Recorded as it arrives, not just at the end, so a run stopped
+                # after this point still persists what the user already saw.
+                self._result.answer = "\n\n".join(
+                    part for part in answer_parts if part
+                ).strip()
+            return
+
+        # Prose accompanying a tool call is reasoning, not answer.
+        if response.text:
+            self._record_thoughts(response.text)
+            yield thoughts_event(response.text)
+
+        conversation.append(
+            Message(
+                role=MessageRole.ASSISTANT,
+                content=response.text,
+                tool_calls=list(response.tool_calls),
+            )
+        )
+
+        results: list[ToolResult] = []
+        async for event in self._run_tools(request, response.tool_calls, 
results):
+            yield event
+
+        conversation.append(Message(role=MessageRole.USER, 
tool_results=results))
+
+    async def _run_tools(
+        self,
+        request: RunRequest,
+        calls: list[ToolCall],
+        results: list[ToolResult],
+    ) -> AsyncIterator[StreamEvent]:
+        """Execute this turn's tool calls, appending outcomes to 
``results``."""
+        for call in calls:
+            if self._cancelled(request):
+                self._result.cancelled = True
+                return
+
+            yield thinking_event(
+                ProgressStage.TOOL,
+                f"Running {call.name}",
+                {"tool_name": call.name},
+            )
+            result, detail = self._invoke_tool(request, call)
+            results.append(result)
+            record = self._record_call(call, result, detail)
+
+            # The frame carries the same record that is persisted, rather than 
a
+            # subset assembled separately. The subset was missing the arguments
+            # and the output, so a step expanded during a run showed nothing at
+            # all unless its tool happened to supply a display — and then 
filled
+            # itself in on reload, which looked like the detail arrived late.
+            # Sharing one record makes that class of drift impossible.
+            yield checkpoint_event(
+                f"{'Failed' if result.is_error else 'Finished'} {call.name}",
+                # ``tool_name`` as well as ``name``: the progress frames use 
that
+                # key, so a consumer reading either finds what it expects.
+                {"tool_name": call.name, **record},
+            )
+
+    async def _safe_turn(
+        self,
+        request: RunRequest,
+        conversation: list[Message],
+        turn: int,
+    ) -> AsyncIterator[StreamEvent]:
+        """
+        One model round trip, converting failure into a ``None`` response.
+
+        A generator rather than a coroutine so the answer can reach the client 
as
+        the model produces it. The response is handed back on
+        :attr:`_last_response` because an async generator cannot both yield 
events
+        and return a value — the same reason ``_turn_loop`` writes into
+        ``answer_parts``.
+
+        The failure detail goes to the log; the caller emits a message that 
cannot
+        leak a URL, a credential or a fragment of someone else's query.
+        """
+        recorder = current_run()
+        started = time.monotonic()
+        self._last_response = None
+        try:
+            async for event in self._one_turn(request, conversation):
+                yield event
+        except LLMError as ex:
+            logger.warning("AI provider error on turn %s: %s", turn, ex)
+            self._result.error = str(ex)
+            self._trace_model_call(recorder, request, turn, started, error=ex)
+            self._last_response = None
+            return
+        except Exception as ex:  # pylint: disable=broad-except
+            logger.exception("Unexpected error in AI runtime on turn %s", turn)
+            self._result.error = GENERIC_ERROR_MESSAGE
+            self._trace_model_call(recorder, request, turn, started, error=ex)
+            self._last_response = None
+            return
+        self._trace_model_call(
+            recorder, request, turn, started, response=self._last_response
+        )
+
+    def _trace_model_call(
+        self,
+        recorder: RunRecorder,
+        request: RunRequest,
+        turn: int,
+        started: float,
+        response: LLMResponse | None = None,
+        error: BaseException | None = None,
+    ) -> None:
+        """
+        Report one round trip to telemetry.
+
+        Content is passed as-is; whether any of it survives into a trace is the
+        redaction policy's decision, made in one place rather than here.
+        """
+        if not recorder.enabled:
+            return
+        usage = response.usage if response is not None else TokenUsage()
+        recorder.model_call(
+            turn=turn,
+            # The concrete identifier when the provider reported one, and the
+            # capability tier otherwise, so a trace can always be grouped by
+            # what the run asked for.
+            model=usage.get("model") or request.model_alias.value,
+            duration_ms=int((time.monotonic() - started) * 1000),
+            input_tokens=usage.get("input_tokens"),
+            output_tokens=usage.get("output_tokens"),
+            stop_reason=response.stop_reason if response is not None else None,
+            error_type=type(error).__name__ if error is not None else None,
+            system_prompt=request.system_prompt,
+            response_text=response.text if response is not None else None,
+        )
+        if error is not None:
+            recorder.error(error)
+
+    async def _one_turn(
+        self,
+        request: RunRequest,
+        conversation: list[Message],
+    ) -> AsyncIterator[StreamEvent]:
+        """
+        Call the model once, yielding answer text as the model produces it.
+
+        Streaming is used when the provider supports it. The assembled response
+        is left on :attr:`_last_response` rather than returned, because a
+        generator cannot do both; it has the same shape either way, so callers 
do
+        not branch on which path ran.
+        """
+        completion = CompletionRequest(
+            messages=conversation,
+            system=request.system_prompt,
+            model_alias=request.model_alias,
+            tools=tuple(request.tools.definitions()) if request.tools else (),
+        )
+
+        if not self.provider.supports_streaming:
+            self._last_response = await self.provider.complete(completion)
+            return
+
+        text_parts: list[str] = []
+        thinking_parts: list[str] = []
+        tool_calls: list[ToolCall] = []
+        usage = None
+
+        async for event in self.provider.stream(completion):
+            if event.kind is StreamEventKind.TEXT:
+                text_parts.append(event.text)
+                # Forwarded as it arrives. Text is buffered as well, because 
the
+                # turn is only known to be an answer once the model stops 
without
+                # asking for a tool — prose before a tool call is reasoning, 
and
+                # is re-routed as such in ``_consume``. A reader that has 
already
+                # seen it replaces its copy on the ``final`` frame.
+                if event.text:
+                    self._streamed_text = True
+                    yield assistant_delta_event(event.text)
+            elif event.kind is StreamEventKind.THINKING:
+                thinking_parts.append(event.text)
+            elif event.kind is StreamEventKind.TOOL_USE and event.tool_call:
+                tool_calls.append(event.tool_call)
+            elif event.kind is StreamEventKind.USAGE:
+                usage = event.usage
+
+        self._last_response = LLMResponse(
+            text="".join(text_parts),
+            thinking="".join(thinking_parts),
+            tool_calls=tool_calls,
+            usage=usage or {},
+            stop_reason="tool_use" if tool_calls else "end_turn",
+        )
+
+    def _invoke_tool(
+        self,
+        request: RunRequest,
+        call: ToolCall,
+    ) -> tuple[ToolResult, dict[str, Any]]:
+        """
+        Run one tool call, through the policy chain first.
+
+        Returns the model-facing result plus a display dict for the UI. A 
denial
+        is handed back to the model as an error result rather than raised,
+        because the reason text is what steers it towards an acceptable
+        alternative. A tool that raises is treated the same way: one broken 
tool
+        should not end an otherwise productive turn.
+
+        Every exit reports a telemetry span, so a refusal and a failure are as
+        visible to a monitoring system as a success is.
+        """
+        if request.policies is not None:
+            denial = request.policies.check(call.name, call.arguments)
+            if denial is not None:
+                return self._traced(
+                    call,
+                    ToolResult(call_id=call.id, content=denial.reason, 
is_error=True),
+                    {"denied": True},
+                    error_type=POLICY_DENIED,
+                )
+
+        if request.tools is None:
+            return self._traced(
+                call,
+                ToolResult(
+                    call_id=call.id,
+                    content=f"No tool named {call.name} is available.",
+                    is_error=True,
+                ),
+                {},
+                error_type=TOOL_UNAVAILABLE,
+            )
+
+        try:
+            result, detail = self._dispatch(request.tools, call)
+        except Exception as ex:  # pylint: disable=broad-except
+            logger.exception("AI tool %s failed", call.name)
+            # The model is told the tool failed but not why, for the same 
reason
+            # the user is not: exception text is not a safe channel.
+            return self._traced(
+                call,
+                ToolResult(
+                    call_id=call.id,
+                    content=f"{call.name} failed and returned no result.",
+                    is_error=True,
+                ),
+                {},
+                error_type=type(ex).__name__,
+            )
+        error_type = detail.get("error_type")
+        return self._traced(
+            call,
+            result,
+            detail,
+            error_type=str(error_type) if error_type else None,
+        )
+
+    def _traced(
+        self,
+        call: ToolCall,
+        result: ToolResult,
+        detail: dict[str, Any],
+        error_type: str | None = None,
+    ) -> tuple[ToolResult, dict[str, Any]]:
+        """
+        Report one tool invocation and hand the outcome back unchanged.
+
+        Placed on the return path of :meth:`_invoke_tool` so that a denial, an
+        unknown tool and a raising tool are each reported once, with the 
failure
+        class the caller could not have reconstructed from the result alone.
+        """
+        recorder = current_run()
+        if recorder.enabled:
+            recorder.tool_call(
+                tool_name=call.name,
+                # The dispatcher already timed the work. Timing it again here
+                # would report the same milliseconds twice under two names.
+                duration_ms=int(detail.get("duration_ms") or 0),
+                ok=not result.is_error,
+                error_type=error_type,
+                truncated=bool(detail.get("truncated")),
+                arguments=call.arguments,
+                output=result.content,
+            )
+        return result, detail
+
+    def _dispatch(
+        self,
+        tools: Any,
+        call: ToolCall,
+    ) -> tuple[ToolResult, dict[str, Any]]:
+        """
+        Call the dispatcher, preferring the richer interface when offered.
+
+        A dispatcher that only implements ``dispatch`` still works — it simply
+        contributes no display detail — which keeps the minimum a test stub has
+        to implement small.
+        """
+        invoke = getattr(tools, "invoke", None)
+        if invoke is None:
+            return tools.dispatch(call), {}
+
+        invocation = invoke(call)
+        detail: dict[str, Any] = {"duration_ms": getattr(invocation, 
"duration_ms", 0)}
+        if getattr(invocation, "truncated", False):
+            detail["truncated"] = True
+        if error_type := getattr(invocation, "error_type", None):
+            detail["error_type"] = error_type

Review Comment:
   This adds the exception class to `detail`, but `_record_call()` persists 
that same dictionary and emits it on the stream. A tool failure such as 
`OperationalError` will therefore expose its internal class through the thread 
API despite this change being telemetry-only. Could `error_type` be passed 
directly to `_traced()` instead of putting it in the shared detail record?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to