sadpandajoe commented on code in PR #43136: URL: https://github.com/apache/superset/pull/43136#discussion_r4050749065
########## 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) + except Exception: # pylint: disable=broad-except + logger.warning("Could not publish AI event for run %s", run_id) + + def consume( + self, + run_id: str, + timeout_seconds: float, + poll_seconds: float = 1.0, + ) -> Iterator[StreamEvent | None]: + import time + + stream = self._stream(run_id) + deadline = time.monotonic() + timeout_seconds + last_id = "-" + + while time.monotonic() < deadline: + try: + entries = self._cache.xrange(stream, last_id, "+", 100) + except Exception: # pylint: disable=broad-except + logger.warning("Could not read AI events for run %s", run_id) + yield IDLE + time.sleep(poll_seconds) + continue + + fresh = [entry for entry in entries if _entry_id(entry) != last_id] + if not fresh: + yield IDLE + time.sleep(poll_seconds) + continue + + for entry in fresh: + last_id = _entry_id(entry) + event = _decode(entry) + if event is None: + continue + yield event + if event.type in _TERMINAL: + return + + def close(self, run_id: str) -> None: + # The stream is left to expire rather than deleted, so a reader that is + # still catching up is not cut off mid-replay. + try: + self._cache.expire(self._stream(run_id), self._ttl) + except Exception: # pylint: disable=broad-except + logger.debug("Could not set TTL on AI event stream for %s", run_id) + + +def _entry_id(entry: Any) -> str: + """Stream entry id, tolerating bytes from the Redis client.""" + raw = entry[0] + return raw.decode() if isinstance(raw, bytes) else str(raw) + + +def _decode(entry: Any) -> StreamEvent | None: + """Rebuild an event from a stream entry, skipping anything malformed.""" + fields = entry[1] + raw = fields.get(b"data") or fields.get("data") + if raw is None: + return None + if isinstance(raw, bytes): + raw = raw.decode() + try: + decoded = json.loads(raw) + return StreamEvent(StreamEventType(decoded["type"]), decoded["payload"]) + except (json.JSONDecodeError, KeyError, ValueError, TypeError): + logger.warning("Discarding malformed AI event") + return None + + +def get_event_bus() -> BaseEventBus: + """ + Build the configured bus, refusing combinations that cannot work. + + An in-memory bus with worker execution is a silent failure — every stream + would sit empty while the run completed elsewhere — so it is rejected at + construction rather than discovered in production. + """ + from flask import current_app + + mode = current_app.config.get("AI_ASSISTANT_EXECUTION_MODE", "inline") + kind = current_app.config.get("AI_ASSISTANT_EVENT_BUS", "memory") + + if mode == "worker" and kind == "memory": + raise RuntimeError( + "AI_ASSISTANT_EXECUTION_MODE='worker' requires " + "AI_ASSISTANT_EVENT_BUS='redis': an in-process bus cannot carry " + "events from a Celery worker to the web process." + ) + + if kind == "memory": + return _memory_bus() + + return RedisStreamEventBus( + cache=_stream_backend(), + prefix=current_app.config.get("AI_ASSISTANT_EVENT_STREAM_PREFIX", "ai-events-"), + ttl_seconds=current_app.config.get("AI_ASSISTANT_EVENT_TTL_SECONDS", 900), + ) + + +def _stream_backend() -> Any: + """ + Build a cache client that can speak Redis streams. + + Deliberately not ``cache_manager.cache``: the general-purpose cache is a + Flask-Caching client with no stream commands, so publishing through it would + fail with an ``AttributeError`` on the first event. The stream methods live + on the same backend classes the async-query channel uses, and those are + constructed from a config dict rather than taken from the extension. + """ + from flask import current_app + + from superset.async_events.cache_backend import ( Review Comment: Worker mode imports a module that no longer exists at this head (`superset.async_events` moved to `superset.coordination`), so queued turns fail before `stream_turn` can finalize the assistant message and leave it permanently pending. Could this use the current backend import and add a worker-mode regression that actually constructs the Redis event bus? ########## superset/migrations/versions/2026-09-14_00-00_e84f26b1c903_merge_ai_chat_with_master.py: ########## @@ -0,0 +1,34 @@ +# 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. +"""Merge AI chat tables with the main migration chain. + +Revision ID: e84f26b1c903 +Revises: a1c4f7e29b31, 7e2c9a4f1b83 +Create Date: 2026-09-14 00:00:00.000000 + +""" + +revision = "e84f26b1c903" +down_revision = ("a1c4f7e29b31", "7e2c9a4f1b83") Review Comment: Merging this revision with the current `master` migration graph leaves both `e84f26b1c903` and `60f94cd6cd11` as heads, so `superset db upgrade` cannot select a single upgrade target. Could this be rebased and joined to the current migration head before merge? ########## superset/ai/orchestrator.py: ########## @@ -0,0 +1,647 @@ +# 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. +""" +Runs one assistant turn end to end. + +Sits between the HTTP layer and the runtime: loads the conversation, assembles +the prompt, resolves the tools the chosen profile allows, drives the runtime, +publishes every event to the bus, and records the outcome on the assistant +message. + +Deliberately independent of *where* it runs. The same function body serves the +inline path and the Celery path, which is what makes the execution mode a +configuration choice rather than two implementations that drift apart. +""" + +from __future__ import annotations + +import asyncio +import logging +import uuid as uuid_module +from collections.abc import AsyncIterator, Iterator +from dataclasses import dataclass +from typing import Any + +from superset.ai.events import ( + cancelled_event, + done_event, + error_event, + GENERIC_ERROR_MESSAGE, + session_event, + StreamEvent, +) +from superset.ai.llm.base import Message +from superset.ai.telemetry import bind_run, current_run, start_run +from superset.ai.types import MessageRole, MessageStatus, RunOutcome, StreamEventType +from superset.utils.decorators import transaction + +logger = logging.getLogger(__name__) + +#: Cache key prefix for a run's cancellation flag. A flag rather than a signal +#: because a worker cannot be interrupted mid-call reliably; the runtime checks +#: this between steps. +_CANCEL_PREFIX = "ai-cancel-" + +#: How long a cancellation request stays meaningful. +_CANCEL_TTL_SECONDS = 900 + +#: Stored when a run is stopped before it produced any answer, so the +#: transcript still records that the turn happened. +_STOPPED_WITHOUT_ANSWER = "_Stopped before an answer was produced._" + +#: Stored when a run exhausted its time budget without saying anything. Phrased +#: as something the user can act on, because retrying is usually the right move. +_TIMED_OUT_WITHOUT_ANSWER = ( + "The assistant ran out of time before it could answer. Please try again." +) + +#: Ceiling on the page context recorded on a message. Well below the prompt's own +#: limit: this is stored per turn and read back with the whole transcript. +_RECORDED_CONTEXT_LIMIT = 4_000 + + +@dataclass +class TurnRequest: + """One unit of work: answer the latest message on a thread.""" + + thread_uuid: str + user_id: int + run_id: str + #: Assistant message row to fill in. Created before the run starts so a + #: client that reconnects has something to attach to. + assistant_message_uuid: str + profile_key: str | None = None + #: Concrete model to pin, overriding the profile's tier. + model: str | None = None + #: What the user had on screen when they asked. Supplied by the client, + #: which is the only party that knows which tab is open, what is typed in + #: the editor and which filters are applied. + page_context: dict[str, Any] | None = None + + def to_payload(self) -> dict[str, Any]: + """Serialise for the task broker.""" + return { + "thread_uuid": self.thread_uuid, + "user_id": self.user_id, + "run_id": self.run_id, + "assistant_message_uuid": self.assistant_message_uuid, + "profile_key": self.profile_key, + "model": self.model, + "page_context": self.page_context, + } + + @classmethod + def from_payload(cls, payload: dict[str, Any]) -> TurnRequest: + """Rebuild from a broker payload.""" + return cls(**payload) + + +def new_run_id() -> str: + """Identifier for one run, used as the event-stream key.""" + return str(uuid_module.uuid4()) + + +#: Runs cancelled in this process. +#: +#: Held alongside the cache rather than instead of it. Superset's default cache +#: is a null cache, which accepts a write and discards it — so a cache-only +#: implementation would leave cancellation silently broken on a default install, +#: with the button appearing to work and nothing stopping. This set makes inline +#: execution correct with no cache at all; the cache is what carries a +#: cancellation across processes for worker execution. +_CANCELLED_LOCALLY: set[str] = set() + + +def request_cancel(run_id: str) -> None: + """ + Ask a run to stop. + + Cooperative by design: the flag is recorded here and observed by the runtime + between steps. A run blocked inside a single long model call or query will + not notice until that call returns, which is a real limit worth documenting + rather than hiding. + """ + from superset.extensions import cache_manager + + _CANCELLED_LOCALLY.add(run_id) + try: + cache_manager.cache.set( + f"{_CANCEL_PREFIX}{run_id}", True, timeout=_CANCEL_TTL_SECONDS + ) + except Exception: # pylint: disable=broad-except + logger.warning("Could not record cancellation for AI run %s", run_id) + + +def is_cancelled(run_id: str) -> bool: + """Whether a stop has been requested for this run.""" + from superset.extensions import cache_manager + + if run_id in _CANCELLED_LOCALLY: + return True + try: + return bool(cache_manager.cache.get(f"{_CANCEL_PREFIX}{run_id}")) + except Exception: # pylint: disable=broad-except + # A cache that cannot be read must not make every run appear cancelled; + # that would stop all inference the moment the cache went away. + return False + + +def clear_cancel(run_id: str) -> None: + """Drop a run's cancellation flag.""" + from superset.extensions import cache_manager + + _CANCELLED_LOCALLY.discard(run_id) + try: + cache_manager.cache.delete(f"{_CANCEL_PREFIX}{run_id}") + except Exception: # pylint: disable=broad-except + logger.debug("Could not clear cancellation flag for AI run %s", run_id) + + +def stream_turn(request: TurnRequest) -> Iterator[StreamEvent]: + """ + Answer a turn, yielding events as they happen. + + This is the primary entry point. Inline execution consumes it directly from + inside the streaming response, which means the producer and the reader are + the same process by construction — important because Superset runs several + web workers, and a turn that published to one process's in-memory queue + while the browser's stream landed on another would appear to hang forever. + + Never raises for an operational failure: a failure is an ``error`` event and + an ``error`` message status, because the caller may already have flushed + response headers or may be a worker with no one to report to. + """ + recorder = start_run( + run_id=request.run_id, + thread_uuid=request.thread_uuid, + user_id=request.user_id, + ) + # Shared with ``_run`` so the ``finally`` below can see the runtime's + # partial result and whether the message was already written. + state: dict[str, Any] = {} + try: + # Bound here rather than inside ``_run`` so that a run which fails before + # it has resolved a profile still produces a start and an end, and so + # that the runtime can report its own spans without the runtime contract + # growing a telemetry parameter. + with bind_run(recorder): + recorder.run_started() + yield from _run(request, state) + except Exception as ex: # pylint: disable=broad-except + logger.exception("AI turn failed for run %s", request.run_id) + recorder.error(ex) + recorder.run_ended(outcome=RunOutcome.ERROR) + answer, extra = _partial_from_state(state) + extra["outcome"] = RunOutcome.ERROR.value + _finalise_message( + request.assistant_message_uuid, + # The generic text rather than the exception: this is persisted and + # served back to the browser, so it must not carry internals. The + # detail is in the log line above, keyed by run id. + content=answer or GENERIC_ERROR_MESSAGE, + status=MessageStatus.ERROR, + extra=extra, + ) + state["finalised"] = True + yield error_event() + yield done_event(ok=False) + finally: + clear_cancel(request.run_id) + # A client that stops the run, or simply navigates away, abandons this + # generator part-way through. Nothing above will have written the + # message, so it would otherwise sit in ``streaming`` with no content + # for ever — the user loses both the partial answer and any record that + # the turn happened. Persist whatever was produced. + _abandon_message(request.assistant_message_uuid, state) + # Idempotent, so the ordinary paths above win. + recorder.run_ended(outcome=RunOutcome.CANCELLED) + + +def execute_turn(request: TurnRequest) -> RunOutcome: + """ + Answer a turn, publishing events to the event bus. + + Used by worker execution, where the reader is in another process. Shares its + whole body with :func:`stream_turn` so the two execution modes cannot drift + apart in behaviour. + """ + from superset.ai.eventbus import get_event_bus + + bus = get_event_bus() + outcome = RunOutcome.SUCCESS + + for event in stream_turn(request): + bus.publish(request.run_id, event) + if event.type is StreamEventType.ERROR: + outcome = RunOutcome.ERROR + elif event.type is StreamEventType.CANCELLED: + outcome = RunOutcome.CANCELLED + elif event.type is StreamEventType.DONE and not event.payload.get("ok"): + # A run that ended un-ok without an explicit error frame timed out. + if outcome is RunOutcome.SUCCESS: + outcome = RunOutcome.TIMEOUT + + return outcome + + +def _run(request: TurnRequest, state: dict[str, Any]) -> Iterator[StreamEvent]: + """Assemble and drive the run. See :func:`stream_turn` for error policy.""" + from superset.ai.factories import ( + get_profiles, + get_provider, + get_runtime, + get_tools_for_profile, + ) + from superset.ai.policy import load_policy_chain + from superset.ai.runtime.base import RunRequest + from superset.daos.ai import AIChatMessageDAO, AIChatThreadDAO + + recorder = current_run() + + thread = AIChatThreadDAO.find_by_uuid_for_user(request.thread_uuid, request.user_id) + if thread is None: + # The thread vanished between accepting the message and running it. + recorder.run_ended(outcome=RunOutcome.ERROR) + yield error_event("That conversation is no longer available.") + yield done_event(ok=False) + return + + profile = get_profiles().get(request.profile_key) + tools = get_tools_for_profile(profile) + provider = get_provider() + runtime = get_runtime(provider) + state["runtime"] = runtime + + yield session_event(request.thread_uuid, request.assistant_message_uuid) + + from superset.ai.page_context import render_page_context + + history = _build_history(AIChatMessageDAO.find_for_thread(thread)) + # Recorded as well as prompted with, so the transcript can show what the + # assistant was told about the user's screen. An answer that looks wrong is + # usually an answer to a different question than the reader assumed, and the + # page context is where that difference lives. + rendered_context = render_page_context(request.page_context) + state["page_context"] = rendered_context + system_prompt = _build_system_prompt(tools, rendered_context) + model = _resolved_model(provider, request.model, profile) + + recorder.describe( + agent_key=profile.key, + model=model, + question=_latest_question(history), + ) + + run_request = RunRequest( + messages=history, + system_prompt=system_prompt, + tools=tools, + policies=load_policy_chain(), + model_alias=profile.model_alias, + max_turns=profile.max_turns or _config("AI_AGENT_MAX_TURNS", 20), + timeout_seconds=profile.timeout_seconds + or _config("AI_AGENT_TIMEOUT_SECONDS", 300), + should_cancel=lambda: is_cancelled(request.run_id), + ) + + _mark_streaming(request.assistant_message_uuid) + + # The runtime is async and this is a synchronous generator, so the async + # events are drained into a list per batch rather than bridged with a + # thread. Collecting the whole run before yielding would defeat streaming, + # so the loop pulls one event at a time from a dedicated event loop. + yield from _drain(runtime.run(run_request)) + + result = runtime.result + outcome = _outcome_of(result) + if result.error is not None: + # The only place the provider's own words are recorded. They do not go on + # the message: that is served back to the browser, and a transport error + # can name internal hosts. + logger.warning( + "AI run %s failed: %s", + request.run_id, + result.error, + ) + _finalise_message( + request.assistant_message_uuid, + content=_terminal_content(result, outcome), + status=_status_of(outcome), + extra={ + "outcome": outcome.value, + "agent_key": profile.key, + "model": model, + "tool_calls": result.tool_calls, + "turns": result.turns, + **_recorded_context(rendered_context), + }, + ) + + recorder.run_ended( + outcome=outcome, + turns=result.turns, + answer=result.answer, + ) + + state["finalised"] = True + + if outcome is RunOutcome.CANCELLED: + yield cancelled_event() + yield done_event(ok=outcome is RunOutcome.SUCCESS) + + +def _drain(source: AsyncIterator[StreamEvent]) -> Iterator[StreamEvent]: + """ + Pull an async iterator one item at a time from a synchronous caller. + + A single event loop is kept for the whole run and stepped with + ``__anext__``, so each event reaches the client as it is produced rather + than after the run completes. + """ + loop = asyncio.new_event_loop() + try: + iterator = source.__aiter__() + while True: + try: + yield loop.run_until_complete(iterator.__anext__()) + except StopAsyncIteration: + return + finally: + loop.close() + + +def _build_history(messages: list[Any]) -> list[Message]: + """ + Convert stored rows into provider messages, trimmed to the configured budget. + + Trimming is newest-first by count and then by total characters, because an + old turn is less useful than a recent one and an oversized request is + rejected outright by every provider. + """ + max_messages = _config("AI_ASSISTANT_MAX_HISTORY_MESSAGES", 25) + max_chars = _config("AI_ASSISTANT_MAX_HISTORY_CHARS", 100_000) + + usable = [ + message + for message in messages + if message.content and message.role != MessageRole.SYSTEM.value + ] + window = usable[-max_messages:] + + # Trimmed to the budget, but never to nothing: a single over-budget message + # is still sent, because the provider's own error about it is more useful + # than a request with no question in it. + total = sum(len(message.content) for message in window) + while len(window) > 1 and total > max_chars: + dropped = window.pop(0) + total -= len(dropped.content) + + return [ + Message(role=MessageRole(message.role), content=message.content) + for message in window + ] + + +def _latest_question(history: list[Message]) -> str | None: + """ + The question this turn is answering. + + Offered to telemetry, which drops it unless a deployment has turned + redaction off. ``None`` when the turn was somehow queued with no user + message, which is a state worth being able to see rather than crash on. + """ + for message in reversed(history): + if message.role is MessageRole.USER and message.content: + return message.content + return None + + +def _build_system_prompt(tools: Any, rendered_context: str) -> str: + """ + Assemble the system prompt for the tools actually on offer. + + The page context arrives already rendered, and is appended after assembly + rather than joining the section list, because it is per-request data rather + than a configured section — and because the layering rules deliberately + refuse content from anywhere but ``superset.core`` in the prompt's own + sections. Rendering happens in the caller so the same text can be recorded on + the message without rendering it twice. + """ + from flask import current_app + + from superset.ai.prompts import assemble_system_prompt + from superset.ai.prompts.core import core_sections + + prompt = assemble_system_prompt( + core_sections(), Review Comment: Configured `AI_KNOWLEDGE_PROVIDERS` never reach this production prompt path, and no knowledge tool consumes their domains, so deployments silently run without the catalog or dialect guidance the setting promises. Could this resolve and include the configured providers here, or remove the unwired public contract? ########## superset-frontend/src/features/ai/hooks/useChatBot.ts: ########## @@ -0,0 +1,1323 @@ +/** + * 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 Conversation state and the send loop. + * + * Runs are tracked per conversation, not globally. That is the point of the + * structure: a user can start something slow in one conversation, switch to + * another and keep working, and come back to find the first still going. A single + * `isLoading` flag would have made switching away cancel or corrupt the run. + * + * The server owns the transcript. A finished run is re-read from it rather than + * assembled from the frames, so the tool calls persisted on the message are what + * the user sees, and what they see survives a reload. + */ + +import { useCallback, useEffect, useRef, useState } from 'react'; +import type { TextAreaRef } from 'antd/es/input/TextArea'; +import { logging } from '@apache-superset/core/utils'; +import { t } from '@apache-superset/core/translation'; +import { + type AiAgent, + type AiToolCall, + type ChatMessageWithMeta, + type ChatTab, + type CheckpointPayload, +} from '../types'; +import { + AGENT_STORAGE_KEY, + ChatRequestAbortedError, + ChatStreamEventError, + ChatStreamTimeoutError, + DEFAULT_AGENT_KEY, + DEFAULT_CHAT_AGENT, + cancelChatRun, + describeRequestError, + fetchAgents, + fetchSuggestedPrompts, + loadStoredAgentKey, + normalizeChatAgents, + startRun, + streamRun, + submitFeedback, +} from './chatRequest'; +import { + NEW_CHAT_NAME, + createThread, + deleteThread as deleteThreadApi, + getThread, + listThreads, + threadToTab, + updateThread, +} from './chatThreadsApi'; +import { buildQuickPrompts } from './quickPrompts'; +import { + buildPageContextPayload, + usePageContext, + type PageContext, +} from './usePageContext'; + +/** Cache of the conversation list, so the menu renders before the list arrives. */ +export const CHAT_TABS_STORAGE_KEY = 'superset-chat-tabs'; + +/** Which conversation was last open. */ +export const ACTIVE_TAB_STORAGE_KEY = 'superset-chat-active-tab'; + +/** Recent inputs, recalled with the arrow keys. */ +export const HISTORY_STORAGE_KEY = 'superset-chat-history'; + +export { AGENT_STORAGE_KEY } from './chatRequest'; + +/** How many inputs the arrow-key history keeps. */ +const MAX_INPUT_HISTORY = 50; + +/** A conversation title derived from a message is clipped to this. */ +const MAX_TAB_NAME_LENGTH = 30; + +export type ChatRunStatus = 'running' | 'cancelling'; + +/** Shared empty list, so a render with no steps yet keeps a stable identity. */ +const EMPTY_TOOL_CALLS: AiToolCall[] = []; + +interface ActiveChatRun { + requestId: string; + tabId: string; + threadId: string; + runId?: string; + controller: AbortController; + isStreaming: boolean; + liveThoughts: string; + liveToolLog: string; + /** + * Steps taken so far, as structured records rather than log lines. + * + * Carried alongside `liveToolLog` so a run in flight can be rendered the same + * way a finished one is — expandable per step, with the SQL and the rows it + * returned — instead of as a wall of text that only becomes legible once the + * transcript is re-read from the server. + */ + liveToolCalls: AiToolCall[]; + /** The page context this run was given, so the live view can show it too. */ + livePageContext?: string; + /** + * The answer so far, as the model produces it. + * + * Rendered directly: the deltas used to be folded into `liveThinking`, which + * nothing displayed, so an answer appeared in one piece the moment the run + * ended however long it had taken to generate. + */ + liveAnswer: string; + liveThinking: string; + status: ChatRunStatus; + startedAt: number; + checkpoint: CheckpointPayload | null; +} + +/** + * An identifier for a turn. + * + * Drawn from `crypto`, not `Math.random`. These become the idempotency key on a + * turn and the handle used to cancel one, so a value another session could guess + * is a correctness and a security problem rather than merely a collision risk. + */ +const generateId = (): string => { + if (typeof crypto.randomUUID === 'function') { + return crypto.randomUUID(); + } + // Older engines expose the entropy source without the convenience wrapper. + const bytes = new Uint8Array(16); + crypto.getRandomValues(bytes); + return Array.from(bytes, byte => byte.toString(16).padStart(2, '0')).join(''); +}; + +const createNewTab = (name: string = NEW_CHAT_NAME): ChatTab => ({ + id: generateId(), + name, + messages: [], + createdAt: Date.now(), +}); + +const truncateTabName = ( + name: string, + maxLength: number = MAX_TAB_NAME_LENGTH, +): string => + name.length <= maxLength ? name : `${name.substring(0, maxLength)}...`; + +const readJson = <T>(key: string, fallback: T): T => { + try { + const stored = localStorage.getItem(key); + return stored ? (JSON.parse(stored) as T) : fallback; + } catch (caught) { + logging.warn(`[ai] could not read ${key}`, caught); + return fallback; + } +}; + +const writeJson = (key: string, value: unknown): void => { + try { + localStorage.setItem(key, JSON.stringify(value)); + } catch (caught) { + logging.warn(`[ai] could not write ${key}`, caught); + } +}; + +/** + * Reconciles the server's transcript with what is already on screen. + * + * The server's copy is authoritative — it carries the tool calls — but it is not + * necessarily complete the moment a run ends, and replacing outright would then + * erase an answer the user has just read. So anything local that the server has + * not accounted for is kept, matched by identity first and by role and content + * second, which is how a locally-appended turn is recognised once the server + * returns its own copy of it under a real uuid. + */ +export const mergeMessages = ( + fromServer: ChatMessageWithMeta[], + local: ChatMessageWithMeta[], +): ChatMessageWithMeta[] => { + const serverIds = new Set(fromServer.map(message => message.id)); + const serverTurns = new Set( + fromServer.map(message => `${message.role}:${message.content}`), + ); + const unaccounted = local.filter( + message => + !serverIds.has(message.id) && + !serverTurns.has(`${message.role}:${message.content}`), + ); + return [...fromServer, ...unaccounted]; +}; + +/** + * The `page_context` body for one turn. + * + * Returns undefined when there is nothing to send, so an omitted field is + * distinguishable from an empty one. + */ +export const buildRequestPageContext = ( + context: PageContext | undefined, + directive?: string, +): Record<string, unknown> | undefined => { + const payload = context ? buildPageContextPayload(context) : undefined; + if (!directive) { + return payload; + } + const existing = payload?.helper_directives; + return { + ...payload, + helper_directives: [ + directive, + ...(Array.isArray(existing) ? existing : []), + ], + }; +}; + +export interface UseChatBotReturn { + // Conversations + chatTabs: ChatTab[]; + activeTabId: string; + activeTab: ChatTab | undefined; + threadsLoaded: boolean; + handleNewChat: () => Promise<string>; + handleSelectTab: (tabId: string) => Promise<void>; + handleDeleteTab: (tabId: string) => Promise<void>; + handleRenameTab: (tabId: string, newName: string) => void; + // Messages of the active conversation + messages: ChatMessageWithMeta[]; + // Input + inputValue: string; + setInputValue: (value: string) => void; + handleKeyDown: (event: React.KeyboardEvent) => void; + inputRef: React.RefObject<TextAreaRef>; + messagesEndRef: React.RefObject<HTMLDivElement>; + // The run in flight, if any, for the active conversation + isLoading: boolean; + isStreamingResponse: boolean; + liveThoughts: string; + liveToolLog: string; + /** Steps taken so far in the run in flight, for the structured live view. */ + liveToolCalls: AiToolCall[]; + /** The page context the run in flight was given. */ + livePageContext?: string; + /** The answer so far for the run in flight. */ + liveAnswer: string; + checkpoint: CheckpointPayload | null; + activeRunStatus: ChatRunStatus | null; + error?: string; + // Actions + sendMessage: ( + messageOverride?: string, + systemPromptOverride?: string, + ) => Promise<void>; + handleCancelRun: () => Promise<void>; + handleCheckpointContinue: () => void; + handleFeedback: (messageId: string, feedback: 'like' | 'dislike') => void; + messageFeedback: Record<string, 'like' | 'dislike'>; + // Suggestions + /** The message whose run just ended; its thought process stays open. */ + justCompletedId?: string; + quickPrompts: string[]; + loadQuickPrompts: () => void; + applyQuickPrompt: (prompt: string) => Promise<void>; + // Agent profiles + agents: AiAgent[]; + selectedAgent: string; + setSelectedAgent: (key: string) => void; + // Page context + pageContext: PageContext; + includePageContext: boolean; + toggleIncludePageContext: () => void; +} + +export const useChatBot = (): UseChatBotReturn => { + const [chatTabs, setChatTabs] = useState<ChatTab[]>(() => + readJson<ChatTab[]>(CHAT_TABS_STORAGE_KEY, []).map(tab => ({ + // The cache is a placeholder for the menu; message bodies are re-read from + // the server so a stale cache cannot show a conversation that has moved on. + ...tab, + messages: [], + })), + ); + const [activeTabId, setActiveTabId] = useState<string>(() => { + try { + return localStorage.getItem(ACTIVE_TAB_STORAGE_KEY) ?? ''; + } catch { + return ''; + } + }); + const [threadsLoaded, setThreadsLoaded] = useState(false); + const [error, setError] = useState<string | undefined>(undefined); + + const [inputValue, setInputValue] = useState(''); + const [activeRunsByTab, setActiveRunsByTab] = useState< + Record<string, ActiveChatRun> + >({}); + const [quickPrompts, setQuickPrompts] = useState<string[]>([]); + const [messageFeedback, setMessageFeedback] = useState< + Record<string, 'like' | 'dislike'> + >({}); + const [includePageContext, setIncludePageContext] = useState(true); + /** + * The assistant message whose run has only just ended. + * + * Its thought process stays open, because collapsing it the instant the answer + * lands moves everything below it — the answer the user is mid-sentence through + * jumps up the panel. Older messages start closed. + */ + const [justCompletedId, setJustCompletedId] = useState<string | undefined>(); + const [agents, setAgents] = useState<AiAgent[]>([DEFAULT_CHAT_AGENT]); + const [selectedAgent, setSelectedAgent] = useState<string>(() => + loadStoredAgentKey(AGENT_STORAGE_KEY), + ); + + const [messageHistory, setMessageHistory] = useState<string[]>(() => + readJson<string[]>(HISTORY_STORAGE_KEY, []), + ); + const [historyIndex, setHistoryIndex] = useState(-1); + const [currentDraft, setCurrentDraft] = useState(''); + + const messagesEndRef = useRef<HTMLDivElement>(null); + const inputRef = useRef<TextAreaRef>(null); + + /** + * The run map and the conversation list are also held in refs, and the refs are + * the authority. + * + * The send loop has to ask "is this still my run?" between awaits, and it cannot + * ask React: a run that starts and fails inside one batch never causes a render, + * so a ref synced at render time would still be empty and the loop would discard + * its own result as stale. Writing the ref at the point of mutation removes that + * window. The callbacks read the refs rather than the state so their identities + * do not churn on every streamed frame, which would restart effects mid-run. + */ + const activeRunsByTabRef = useRef<Record<string, ActiveChatRun>>({}); + const chatTabsRef = useRef<ChatTab[]>(chatTabs); + const activeTabIdRef = useRef(activeTabId); + activeTabIdRef.current = activeTabId; + + const updateRuns = useCallback( + ( + updater: ( + previous: Record<string, ActiveChatRun>, + ) => Record<string, ActiveChatRun>, + ) => { + activeRunsByTabRef.current = updater(activeRunsByTabRef.current); + setActiveRunsByTab(activeRunsByTabRef.current); + }, + [], + ); + + const updateTabs = useCallback( + (updater: (previous: ChatTab[]) => ChatTab[]) => { + chatTabsRef.current = updater(chatTabsRef.current); + setChatTabs(chatTabsRef.current); + }, + [], + ); + + /** Resolved when the user answers a checkpoint; see `streamRun`. */ + const checkpointGateRef = useRef<{ resolve: () => void } | null>(null); + const mountedRef = useRef(true); + useEffect( + () => () => { + mountedRef.current = false; + }, + [], + ); + + const activeTab = chatTabs.find(tab => tab.id === activeTabId); + const messages = activeTab?.messages ?? []; Review Comment: Deleting the last conversation, or failing the initial thread load, leaves `activeTab` undefined, so this expression creates a new array every render. The quick-prompt effect depends on that array and writes another new array to state, producing an update loop that crashes the panel. Could the empty-message value be stable and this zero-tab path get a regression test? ########## 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. Review Comment: The current head still has no pruning task, command, or DAO operation that reads `AI_ASSISTANT_MESSAGE_RETENTION_DAYS`, so the documented instruction to schedule retention cannot be carried out. Could this add the schedulable pruning path, or remove the inert setting and scheduling claim? -- 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]
