From fc7b12b2bb2b607b9c2cf9dacb4402e90390388d Mon Sep 17 00:00:00 2001
From: Harsh Soni <64592571+harshsoni2024@users.noreply.github.com>
Date: Tue, 8 Sep 2026 05:26:02 +0000
Subject: [PATCH 1/3] fix(mode): Ingest Mode report query metadata (#32533)
* Fixes #22559: Ingest Mode report query metadata
* Address Mode connector review feedback
* Refine Mode pagination loop guard
* Harden Mode connector edge cases
---
ingestion/.basedpyright/baseline.json | 128 ------
.../ingestion/source/dashboard/mode/client.py | 73 ++--
.../source/dashboard/mode/metadata.py | 323 +++++++++++----
.../unit/source/dashboard/mode/test_client.py | 117 ++++++
.../unit/topology/dashboard/test_mode.py | 368 ++++++++++++++++++
.../entity/data/dashboardDataModel.json | 4 +
.../ui/public/locales/en-US/Dashboard/Mode.md | 18 +-
.../api/data/createDashboardDataModel.ts | 1 +
.../entity/data/dashboardDataModel.ts | 1 +
9 files changed, 802 insertions(+), 231 deletions(-)
create mode 100644 ingestion/tests/unit/source/dashboard/mode/test_client.py
create mode 100644 ingestion/tests/unit/topology/dashboard/test_mode.py
diff --git a/ingestion/.basedpyright/baseline.json b/ingestion/.basedpyright/baseline.json
index 464858669139..62f2b7c57a34 100644
--- a/ingestion/.basedpyright/baseline.json
+++ b/ingestion/.basedpyright/baseline.json
@@ -24813,30 +24813,6 @@
"lineCount": 1
}
},
- {
- "code": "reportIndexIssue",
- "range": {
- "startColumn": 22,
- "endColumn": 42,
- "lineCount": 1
- }
- },
- {
- "code": "reportOptionalSubscript",
- "range": {
- "startColumn": 22,
- "endColumn": 42,
- "lineCount": 1
- }
- },
- {
- "code": "reportReturnType",
- "range": {
- "startColumn": 19,
- "endColumn": 27,
- "lineCount": 1
- }
- },
{
"code": "reportReturnType",
"range": {
@@ -27049,30 +27025,6 @@
"lineCount": 3
}
},
- {
- "code": "reportReturnType",
- "range": {
- "startColumn": 15,
- "endColumn": 41,
- "lineCount": 1
- }
- },
- {
- "code": "reportArgumentType",
- "range": {
- "startColumn": 28,
- "endColumn": 63,
- "lineCount": 1
- }
- },
- {
- "code": "reportArgumentType",
- "range": {
- "startColumn": 25,
- "endColumn": 66,
- "lineCount": 1
- }
- },
{
"code": "reportArgumentType",
"range": {
@@ -27081,38 +27033,6 @@
"lineCount": 6
}
},
- {
- "code": "reportAttributeAccessIssue",
- "range": {
- "startColumn": 56,
- "endColumn": 73,
- "lineCount": 1
- }
- },
- {
- "code": "reportAttributeAccessIssue",
- "range": {
- "startColumn": 48,
- "endColumn": 54,
- "lineCount": 1
- }
- },
- {
- "code": "reportAttributeAccessIssue",
- "range": {
- "startColumn": 39,
- "endColumn": 56,
- "lineCount": 1
- }
- },
- {
- "code": "reportCallIssue",
- "range": {
- "startColumn": 14,
- "endColumn": 45,
- "lineCount": 1
- }
- },
{
"code": "reportGeneralTypeIssues",
"range": {
@@ -27121,54 +27041,6 @@
"lineCount": 1
}
},
- {
- "code": "reportArgumentType",
- "range": {
- "startColumn": 28,
- "endColumn": 25,
- "lineCount": 6
- }
- },
- {
- "code": "reportReturnType",
- "range": {
- "startColumn": 30,
- "endColumn": 25,
- "lineCount": 3
- }
- },
- {
- "code": "reportArgumentType",
- "range": {
- "startColumn": 38,
- "endColumn": 47,
- "lineCount": 1
- }
- },
- {
- "code": "reportArgumentType",
- "range": {
- "startColumn": 61,
- "endColumn": 72,
- "lineCount": 1
- }
- },
- {
- "code": "reportCallIssue",
- "range": {
- "startColumn": 18,
- "endColumn": 13,
- "lineCount": 7
- }
- },
- {
- "code": "reportAttributeAccessIssue",
- "range": {
- "startColumn": 51,
- "endColumn": 68,
- "lineCount": 1
- }
- },
{
"code": "reportCallIssue",
"range": {
diff --git a/ingestion/src/metadata/ingestion/source/dashboard/mode/client.py b/ingestion/src/metadata/ingestion/source/dashboard/mode/client.py
index 080637976f43..395a110b71eb 100644
--- a/ingestion/src/metadata/ingestion/source/dashboard/mode/client.py
+++ b/ingestion/src/metadata/ingestion/source/dashboard/mode/client.py
@@ -13,7 +13,7 @@
"""
import traceback
from base64 import b64encode
-from typing import Optional
+from typing import Any, Dict, List, Optional, cast
from metadata.ingestion.connections.source_api_client import TrackedREST
from metadata.ingestion.ometa.client import ClientConfig
@@ -24,7 +24,7 @@
EMBEDDED = "_embedded"
-COLLECTIONS = "collections"
+SPACES = "spaces"
TOKEN = "token"
REPORTS = "reports"
QUERIES = "queries"
@@ -38,6 +38,7 @@
LINKS = "_links"
SHARE = "share"
HREF = "href"
+REPORTS_PAGE_SIZE = 30
class ModeApiClient:
@@ -66,7 +67,7 @@ def __init__(self, config):
def fetch_all_reports(
self, workspace_name: str, filter: Optional[str] = "all"
- ) -> Optional[list]:
+ ) -> List[Dict[str, Any]]:
"""Method to fetch all reports for Mode
Args:
workspace_name:
@@ -75,47 +76,57 @@ def fetch_all_reports(
dict
"""
if filter not in ["custom", "all"]:
- logger.warning(
- "Invalid value for filter. Should be one of ['custom', 'all']"
+ raise ValueError(
+ f"Invalid Mode filter [{filter}]. Expected one of ['custom', 'all']"
)
- return
- all_reports = []
+ all_reports: List[Dict[str, Any]] = []
filter_param = f"?filter={filter}"
- response_collections = self.client.get(
- f"/{workspace_name}/{COLLECTIONS}{filter_param}"
+ response_spaces = cast(
+ Dict[str, Any],
+ self.client.get(f"/{workspace_name}/{SPACES}{filter_param}"),
)
- collections = response_collections[EMBEDDED]["spaces"]
- for collection in collections:
- response_reports = self.get_all_reports_for_collection(
- workspace_name=workspace_name,
- collection_token=collection.get(TOKEN),
- )
- if response_reports:
+ spaces = response_spaces[EMBEDDED][SPACES]
+ for space in spaces:
+ page = 1
+ previous_reports = None
+ while True:
+ response_reports = self.get_reports_for_space(
+ workspace_name=workspace_name,
+ space_token=space[TOKEN],
+ page=page,
+ )
reports = response_reports[EMBEDDED][REPORTS]
+ if reports and reports == previous_reports:
+ raise RuntimeError(
+ f"Mode returned the same report page twice for space "
+ f"[{space[TOKEN]}] at page [{page}]"
+ )
all_reports.extend(reports)
+ if len(reports) < REPORTS_PAGE_SIZE:
+ break
+ previous_reports = reports
+ page += 1
return all_reports
- def get_all_reports_for_collection(
- self, workspace_name: str, collection_token: str
- ) -> Optional[dict]:
- """Method to fetch all reports for a collection
+ def get_reports_for_space(
+ self, workspace_name: str, space_token: str, page: int
+ ) -> Dict[str, Any]:
+ """Fetch one page of reports for a space.
+
Args:
workspace_name:
- collection_token:
+ space_token:
+ page:
Returns:
dict
"""
- try:
- response = self.client.get(
- f"/{workspace_name}/{COLLECTIONS}/{collection_token}/{REPORTS}"
- )
- return response
- except Exception as exc: # pylint: disable=broad-except
- logger.debug(traceback.format_exc())
- logger.warning(f"Error fetching charts: {exc}")
-
- return None
+ return cast(
+ Dict[str, Any],
+ self.client.get(
+ f"/{workspace_name}/{SPACES}/{space_token}/{REPORTS}?page={page}"
+ ),
+ )
def get_all_queries(self, workspace_name: str, report_token: str) -> Optional[dict]:
"""Method to fetch all queries
diff --git a/ingestion/src/metadata/ingestion/source/dashboard/mode/metadata.py b/ingestion/src/metadata/ingestion/source/dashboard/mode/metadata.py
index 77ebe7f744c0..b199bfaab3d2 100644
--- a/ingestion/src/metadata/ingestion/source/dashboard/mode/metadata.py
+++ b/ingestion/src/metadata/ingestion/source/dashboard/mode/metadata.py
@@ -11,14 +11,20 @@
"""Mode source module"""
import traceback
-from typing import Iterable, List, Optional
+from dataclasses import dataclass
+from typing import Any, Dict, Iterable, List, Optional, Union, cast
from metadata.generated.schema.api.data.createChart import CreateChartRequest
from metadata.generated.schema.api.data.createDashboard import CreateDashboardRequest
+from metadata.generated.schema.api.data.createDashboardDataModel import (
+ CreateDashboardDataModelRequest,
+)
from metadata.generated.schema.api.lineage.addLineage import AddLineageRequest
from metadata.generated.schema.entity.data.chart import Chart, ChartType
-from metadata.generated.schema.entity.data.dashboard import (
- Dashboard as Lineage_Dashboard,
+from metadata.generated.schema.entity.data.dashboard import Dashboard
+from metadata.generated.schema.entity.data.dashboardDataModel import (
+ DashboardDataModel,
+ DataModelType,
)
from metadata.generated.schema.entity.data.table import Table
from metadata.generated.schema.entity.services.connections.dashboard.modeConnection import (
@@ -35,6 +41,7 @@
FullyQualifiedEntityName,
Markdown,
SourceUrl,
+ SqlQuery,
)
from metadata.ingestion.api.models import Either
from metadata.ingestion.api.steps import InvalidSourceException
@@ -43,13 +50,23 @@
from metadata.ingestion.source.dashboard.dashboard_service import DashboardServiceSource
from metadata.ingestion.source.dashboard.mode import client
from metadata.utils import fqn
-from metadata.utils.filters import filter_by_chart
+from metadata.utils.filters import filter_by_chart, filter_by_datamodel
from metadata.utils.fqn import build_es_fqn_search_string
from metadata.utils.helpers import clean_uri
from metadata.utils.logger import ingestion_logger
logger = ingestion_logger()
+ModeRecord = Dict[str, Any]
+
+
+@dataclass(frozen=True)
+class ModeDashboardDetails:
+ """Mode report and the queries fetched for it."""
+
+ report: ModeRecord
+ queries: List[ModeRecord]
+
class ModeSource(DashboardServiceSource):
"""
@@ -64,7 +81,10 @@ def __init__(
super().__init__(config, metadata)
self.workspace_name = config.serviceConnection.root.config.workspaceName
self.filter_query_param = config.serviceConnection.root.config.filterQueryParam
- self.data_sources = self.client.get_all_data_sources(self.workspace_name)
+ self.data_sources = cast(
+ Dict[str, ModeRecord],
+ self.client.get_all_data_sources(self.workspace_name) or {},
+ )
@classmethod
def create(
@@ -78,7 +98,7 @@ def create(
)
return cls(config, metadata)
- def get_dashboards_list(self) -> Optional[List[dict]]:
+ def get_dashboards_list(self) -> List[ModeRecord]:
"""
Get List of all dashboards
"""
@@ -86,56 +106,152 @@ def get_dashboards_list(self) -> Optional[List[dict]]:
filter_param = "all" if not self.filter_query_param else self.filter_query_param
return self.client.fetch_all_reports(self.workspace_name, filter_param)
- def get_dashboard_name(self, dashboard: dict) -> str:
+ def get_dashboard_name(self, dashboard: ModeRecord) -> str:
"""
Get Dashboard Name
"""
- return dashboard.get(client.NAME)
+ return cast(str, dashboard.get(client.NAME) or dashboard[client.TOKEN])
+
+ def _dashboard_service_name(self) -> str:
+ return cast(str, cast(Any, self.context.get()).dashboard_service)
- def get_dashboard_details(self, dashboard: dict) -> dict:
+ def get_dashboard_details(self, dashboard: ModeRecord) -> ModeDashboardDetails:
"""
Get Dashboard Details
"""
- return dashboard
+ response = self.client.get_all_queries(
+ workspace_name=self.workspace_name,
+ report_token=cast(str, dashboard[client.TOKEN]),
+ )
+ queries = cast(
+ List[ModeRecord],
+ response.get(client.EMBEDDED, {}).get(client.QUERIES, [])
+ if response
+ else [],
+ )
+ return ModeDashboardDetails(report=dashboard, queries=queries or [])
def yield_dashboard(
- self, dashboard_details: dict
+ self, dashboard_details: ModeDashboardDetails
) -> Iterable[Either[CreateDashboardRequest]]:
"""
Method to Get Dashboard Entity
"""
- dashboard_path = dashboard_details[client.LINKS][client.SHARE][client.HREF]
+ report = dashboard_details.report
+ dashboard_service = self._dashboard_service_name()
+ charts = cast(List[str], getattr(self.context.get(), "charts", []) or [])
+ data_models = cast(
+ List[str], getattr(self.context.get(), "dataModels", []) or []
+ )
+ dashboard_path = cast(str, report[client.LINKS][client.SHARE][client.HREF])
+ report_token = cast(str, report[client.TOKEN])
+ description = cast(Optional[str], report.get(client.DESCRIPTION))
dashboard_url = f"{clean_uri(self.service_connection.hostPort)}{dashboard_path}"
dashboard_request = CreateDashboardRequest(
- name=EntityName(dashboard_details.get(client.TOKEN)),
+ name=EntityName(report_token),
sourceUrl=SourceUrl(dashboard_url),
- displayName=dashboard_details.get(client.NAME),
- description=(
- Markdown(dashboard_details.get(client.DESCRIPTION))
- if dashboard_details.get(client.DESCRIPTION)
- else None
- ),
+ displayName=cast(Optional[str], report.get(client.NAME)),
+ description=Markdown(description) if description else None,
charts=[
FullyQualifiedEntityName(
fqn.build(
self.metadata,
entity_type=Chart,
- service_name=self.context.get().dashboard_service,
+ service_name=dashboard_service,
chart_name=chart,
)
)
- for chart in self.context.get().charts or []
+ for chart in charts
],
- service=self.context.get().dashboard_service,
- owners=self.get_owner_ref(dashboard_details=dashboard_details),
+ dataModels=(
+ [
+ FullyQualifiedEntityName(
+ cast(
+ str,
+ fqn.build(
+ self.metadata,
+ entity_type=DashboardDataModel,
+ service_name=dashboard_service,
+ data_model_name=data_model,
+ ),
+ )
+ )
+ for data_model in data_models
+ ]
+ if self.source_config.includeDataModels
+ else None
+ ),
+ service=FullyQualifiedEntityName(dashboard_service),
+ owners=self.get_owner_ref(dashboard_details=report),
)
- yield Either(right=dashboard_request)
+ yield Either(left=None, right=dashboard_request)
self.register_record(dashboard_request=dashboard_request)
- # pylint: disable=too-many-locals
+ def yield_datamodel( # pyright: ignore[reportIncompatibleMethodOverride]
+ self, dashboard_details: ModeDashboardDetails
+ ) -> Iterable[Either[CreateDashboardDataModelRequest]]:
+ """Yield each Mode query as a dashboard data model."""
+ if not self.source_config.includeDataModels:
+ return
+
+ dashboard_service = self._dashboard_service_name()
+ report_token = cast(str, dashboard_details.report[client.TOKEN])
+ for query in dashboard_details.queries:
+ query_token = cast(Optional[str], query.get(client.TOKEN))
+ if not query_token:
+ yield Either(
+ left=StackTraceError(
+ name=cast(Optional[str], query.get(client.NAME)) or "",
+ error="Mode query is missing its token",
+ stackTrace="",
+ ),
+ right=None,
+ )
+ continue
+ query_name = cast(Optional[str], query.get(client.NAME)) or query_token
+ try:
+ if filter_by_datamodel(
+ self.source_config.dataModelFilterPattern,
+ query_name,
+ ):
+ self.status.filter(query_name, "Data model filtered out.")
+ continue
+
+ raw_query = cast(Optional[str], query.get("raw_query"))
+ datamodel_request = CreateDashboardDataModelRequest(
+ name=EntityName(self._data_model_name(report_token, query_token)),
+ displayName=query_name,
+ service=FullyQualifiedEntityName(dashboard_service),
+ serviceType=self.service_connection.type.value,
+ dataModelType=DataModelType.ModeDataModel,
+ sourceUrl=SourceUrl(
+ f"{clean_uri(self.service_connection.hostPort)}/"
+ f"{self.workspace_name}/{client.REPORTS}/{report_token}/"
+ f"{client.QUERIES}/{query_token}"
+ ),
+ sql=SqlQuery(raw_query) if raw_query else None,
+ columns=[],
+ )
+ yield Either(left=None, right=datamodel_request)
+ self.register_record_datamodel(datamodel_request=datamodel_request)
+ except Exception as exc:
+ yield Either(
+ left=StackTraceError(
+ name=query_name or "",
+ error=f"Error yielding Mode query data model [{query_name}]: {exc}",
+ stackTrace=traceback.format_exc(),
+ ),
+ right=None,
+ )
+
+ @staticmethod
+ def _data_model_name(report_token: str, query_token: str) -> str:
+ return f"{report_token}.{query_token}"
+
+ # pylint: disable=too-many-locals,too-many-branches
def yield_dashboard_lineage_details(
self,
- dashboard_details: dict,
+ dashboard_details: ModeDashboardDetails,
db_service_prefix: Optional[str] = None,
) -> Iterable[Either[AddLineageRequest]]:
"""Get lineage method"""
@@ -147,20 +263,19 @@ def yield_dashboard_lineage_details(
) = self.parse_db_service_prefix(db_service_prefix)
try:
- response_queries = self.client.get_all_queries(
- workspace_name=self.workspace_name,
- report_token=dashboard_details[client.TOKEN],
- )
- queries = response_queries[client.EMBEDDED][client.QUERIES]
-
- for query in queries:
- if not query.get("data_source_id"):
+ for query in dashboard_details.queries:
+ data_source_id = cast(Optional[str], query.get("data_source_id"))
+ if not data_source_id:
continue
- data_source = self.data_sources.get(query.get("data_source_id"))
+ data_source = self.data_sources.get(data_source_id)
if not data_source:
continue
- database_name = data_source.get(client.DATABASE)
+ raw_query = cast(Optional[str], query.get("raw_query"))
+ if not raw_query:
+ continue
+
+ database_name = cast(Optional[str], data_source.get(client.DATABASE))
if (
prefix_database_name
and database_name
@@ -171,11 +286,25 @@ def yield_dashboard_lineage_details(
)
continue
+ search_database_name = prefix_database_name or database_name
+ if not search_database_name:
+ logger.warning(
+ "Skipping Mode query lineage because its data source does not "
+ "provide a database name"
+ )
+ continue
+
lineage_parser = LineageParser(
- query.get("raw_query"),
+ raw_query,
parser_type=self.get_query_parser_type(),
)
query_hash = lineage_parser.query_hash
+ to_entity = self._resolve_lineage_target(
+ dashboard_details=dashboard_details,
+ query=query,
+ )
+ if not to_entity:
+ continue
for table in lineage_parser.source_tables:
database_schema_name, table = fqn.split(str(table))[-2:]
database_schema_name = self.check_database_schema_name(
@@ -204,55 +333,117 @@ def yield_dashboard_lineage_details(
continue
fqn_search_string = build_es_fqn_search_string(
- database_name=prefix_database_name or database_name,
+ database_name=search_database_name,
schema_name=prefix_schema_name or database_schema_name,
service_name=prefix_service_name or "*",
table_name=prefix_table_name or table,
)
- from_entities = self.metadata.search_in_any_service(
- entity_type=Table,
- fqn_search_string=fqn_search_string,
- fetch_multiple_entities=True,
- )
-
- to_entity = self.metadata.get_by_name(
- entity=Lineage_Dashboard,
- fqn=fqn.build(
- self.metadata,
- Lineage_Dashboard,
- service_name=self.config.serviceName,
- dashboard_name=dashboard_details.get(client.TOKEN),
+ from_entities = cast(
+ Optional[List[Table]],
+ self.metadata.search_in_any_service(
+ entity_type=Table,
+ fqn_search_string=fqn_search_string,
+ fetch_multiple_entities=True,
),
)
+
for from_entity in from_entities or []:
- yield self._get_add_lineage_request(
- to_entity=to_entity, from_entity=from_entity
+ lineage = self._get_add_lineage_request(
+ to_entity=to_entity,
+ from_entity=from_entity,
+ sql=raw_query,
)
+ if lineage:
+ yield lineage
except Exception as exc: # pylint: disable=broad-except
yield Either(
left=StackTraceError(
name="Lineage",
error=f"Error to yield dashboard lineage details for service name [{prefix_service_name}]: {exc}",
stackTrace=traceback.format_exc(),
- )
+ ),
+ right=None,
)
+ def _resolve_lineage_target(
+ self,
+ dashboard_details: ModeDashboardDetails,
+ query: ModeRecord,
+ ) -> Optional[Union[DashboardDataModel, Dashboard]]:
+ """Resolve the query model, falling back to its report dashboard."""
+ dashboard_service = self._dashboard_service_name()
+ report_token = cast(str, dashboard_details.report[client.TOKEN])
+ query_token = cast(Optional[str], query.get(client.TOKEN))
+ datamodel_name = (
+ self._data_model_name(report_token, query_token) if query_token else None
+ )
+ data_models = cast(
+ List[str], getattr(self.context.get(), "dataModels", None) or []
+ )
+ if self.source_config.includeDataModels and datamodel_name in data_models:
+ try:
+ datamodel_fqn = cast(
+ str,
+ fqn.build(
+ metadata=self.metadata,
+ entity_type=DashboardDataModel,
+ service_name=dashboard_service,
+ data_model_name=datamodel_name,
+ ),
+ )
+ datamodel = self.metadata.get_by_name(
+ entity=DashboardDataModel,
+ fqn=datamodel_fqn,
+ )
+ if datamodel:
+ return datamodel
+ except Exception as exc:
+ logger.debug(
+ "Could not resolve Mode query data model [%s]: %s",
+ datamodel_name,
+ exc,
+ )
+
+ dashboard_fqn = cast(
+ str,
+ fqn.build(
+ metadata=self.metadata,
+ entity_type=Dashboard,
+ service_name=dashboard_service,
+ dashboard_name=report_token,
+ ),
+ )
+ return self.metadata.get_by_name(entity=Dashboard, fqn=dashboard_fqn)
+
def yield_dashboard_chart(
- self, dashboard_details: dict
+ self, dashboard_details: ModeDashboardDetails
) -> Iterable[Either[CreateChartRequest]]:
"""Get chart method"""
- response_queries = self.client.get_all_queries(
- workspace_name=self.workspace_name,
- report_token=dashboard_details.get(client.TOKEN),
- )
- queries = response_queries[client.EMBEDDED][client.QUERIES]
- for query in queries:
+ report_token = cast(str, dashboard_details.report[client.TOKEN])
+ dashboard_service = self._dashboard_service_name()
+ for query in dashboard_details.queries:
+ query_token = cast(Optional[str], query.get(client.TOKEN))
+ if not query_token:
+ yield Either(
+ left=StackTraceError(
+ name=cast(Optional[str], query.get(client.NAME)) or "",
+ error="Mode query is missing its token; charts could not be fetched",
+ stackTrace="",
+ ),
+ right=None,
+ )
+ continue
response_charts = self.client.get_all_charts(
workspace_name=self.workspace_name,
- report_token=dashboard_details.get(client.TOKEN),
- query_token=query.get(client.TOKEN),
+ report_token=report_token,
+ query_token=query_token,
+ )
+ charts = cast(
+ List[ModeRecord],
+ response_charts.get(client.EMBEDDED, {}).get(client.CHARTS, [])
+ if response_charts
+ else [],
)
- charts = response_charts[client.EMBEDDED][client.CHARTS]
for chart in charts:
chart_name = chart[client.VIEW_VEGAS].get(client.TITLE)
try:
@@ -270,11 +461,11 @@ def yield_dashboard_chart(
f"{clean_uri(self.service_connection.hostPort)}{chart_path}"
)
chart_request = CreateChartRequest(
- name=EntityName(chart.get(client.TOKEN)),
+ name=EntityName(cast(str, chart.get(client.TOKEN))),
displayName=chart_name,
chartType=ChartType.Other,
sourceUrl=SourceUrl(chart_url),
- service=self.context.get().dashboard_service,
+ service=FullyQualifiedEntityName(dashboard_service),
)
yield Either(right=chart_request)
self.register_record_chart(chart_request=chart_request)
diff --git a/ingestion/tests/unit/source/dashboard/mode/test_client.py b/ingestion/tests/unit/source/dashboard/mode/test_client.py
new file mode 100644
index 000000000000..9cbe359df122
--- /dev/null
+++ b/ingestion/tests/unit/source/dashboard/mode/test_client.py
@@ -0,0 +1,117 @@
+# Copyright 2025 Collate
+# Licensed under the Collate Community License, Version 1.0 (the "License");
+# you may not use this file except in compliance with the License.
+# You may obtain a copy of the License at
+# https://github.com/open-metadata/OpenMetadata/blob/main/ingestion/LICENSE
+# 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.
+"""Unit tests for the Mode API client."""
+
+from unittest.mock import MagicMock, call
+
+import pytest
+
+from metadata.ingestion.source.dashboard.mode.client import ModeApiClient
+
+
+def _reports(prefix: str, count: int) -> list[dict]:
+ return [{"token": f"{prefix}-{index}"} for index in range(count)]
+
+
+def _embedded(name: str, values: list[dict]) -> dict:
+ return {"_embedded": {name: values}}
+
+
+@pytest.fixture
+def mode_client() -> ModeApiClient:
+ api_client = ModeApiClient.__new__(ModeApiClient)
+ api_client.client = MagicMock()
+ return api_client
+
+
+def test_fetch_all_reports_paginates_every_space(mode_client):
+ first_space_page = _reports("finance", 30)
+ second_space_page = _reports("finance-extra", 2)
+ operations_page = _reports("operations", 1)
+ mode_client.client.get.side_effect = [
+ _embedded("spaces", [{"token": "finance"}, {"token": "operations"}]),
+ _embedded("reports", first_space_page),
+ _embedded("reports", second_space_page),
+ _embedded("reports", operations_page),
+ ]
+
+ reports = mode_client.fetch_all_reports("acme", "custom")
+
+ assert reports == first_space_page + second_space_page + operations_page
+ assert mode_client.client.get.call_args_list == [
+ call("/acme/spaces?filter=custom"),
+ call("/acme/spaces/finance/reports?page=1"),
+ call("/acme/spaces/finance/reports?page=2"),
+ call("/acme/spaces/operations/reports?page=1"),
+ ]
+
+
+def test_fetch_all_reports_requests_page_after_exactly_thirty_results(mode_client):
+ first_page = _reports("report", 30)
+ mode_client.client.get.side_effect = [
+ _embedded("spaces", [{"token": "space-token"}]),
+ _embedded("reports", first_page),
+ _embedded("reports", []),
+ ]
+
+ reports = mode_client.fetch_all_reports("acme")
+
+ assert reports == first_page
+ assert mode_client.client.get.call_args_list[-1] == call(
+ "/acme/spaces/space-token/reports?page=2"
+ )
+
+
+def test_fetch_all_reports_propagates_later_page_failure(mode_client):
+ mode_client.client.get.side_effect = [
+ _embedded("spaces", [{"token": "space-token"}]),
+ _embedded("reports", _reports("report", 30)),
+ RuntimeError("Mode page request failed"),
+ ]
+
+ with pytest.raises(RuntimeError, match="Mode page request failed"):
+ mode_client.fetch_all_reports("acme")
+
+
+def test_fetch_all_reports_rejects_invalid_filter(mode_client):
+ with pytest.raises(ValueError, match="Invalid Mode filter"):
+ mode_client.fetch_all_reports("acme", "invalid")
+
+ mode_client.client.get.assert_not_called()
+
+
+def test_fetch_all_reports_rejects_repeated_full_page(mode_client):
+ repeated_page = _reports("report", 30)
+ mode_client.client.get.side_effect = [
+ _embedded("spaces", [{"token": "space-token"}]),
+ _embedded("reports", repeated_page),
+ _embedded("reports", repeated_page),
+ ]
+
+ with pytest.raises(RuntimeError, match="same report page twice"):
+ mode_client.fetch_all_reports("acme")
+
+ assert mode_client.client.get.call_args_list[-1] == call(
+ "/acme/spaces/space-token/reports?page=2"
+ )
+
+
+def test_fetch_all_reports_allows_distinct_tokenless_pages(mode_client):
+ first_page = [{"name": f"first-{index}"} for index in range(30)]
+ second_page = [{"name": f"second-{index}"} for index in range(30)]
+ mode_client.client.get.side_effect = [
+ _embedded("spaces", [{"token": "space-token"}]),
+ _embedded("reports", first_page),
+ _embedded("reports", second_page),
+ _embedded("reports", []),
+ ]
+
+ assert mode_client.fetch_all_reports("acme") == first_page + second_page
diff --git a/ingestion/tests/unit/topology/dashboard/test_mode.py b/ingestion/tests/unit/topology/dashboard/test_mode.py
new file mode 100644
index 000000000000..eba45a6555f9
--- /dev/null
+++ b/ingestion/tests/unit/topology/dashboard/test_mode.py
@@ -0,0 +1,368 @@
+# Copyright 2025 Collate
+# Licensed under the Collate Community License, Version 1.0 (the "License");
+# you may not use this file except in compliance with the License.
+# You may obtain a copy of the License at
+# https://github.com/open-metadata/OpenMetadata/blob/main/ingestion/LICENSE
+# 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.
+"""Unit tests for the Mode dashboard source."""
+
+from types import SimpleNamespace
+from unittest.mock import MagicMock, patch
+
+import pytest
+
+from metadata.generated.schema.entity.data.dashboard import Dashboard
+from metadata.generated.schema.entity.data.dashboardDataModel import DashboardDataModel
+from metadata.generated.schema.entity.data.table import Table
+from metadata.generated.schema.metadataIngestion.workflow import (
+ OpenMetadataWorkflowConfig,
+)
+from metadata.generated.schema.type.basic import FullyQualifiedEntityName, Uuid
+from metadata.generated.schema.type.filterPattern import FilterPattern
+from metadata.ingestion.ometa.ometa_api import OpenMetadata
+from metadata.ingestion.source.dashboard.mode.metadata import ModeSource
+
+MOCK_CONFIG = {
+ "source": {
+ "type": "mode",
+ "serviceName": "mock_mode",
+ "serviceConnection": {
+ "config": {
+ "type": "Mode",
+ "hostPort": "https://app.mode.com/",
+ "accessToken": "token",
+ "accessTokenPassword": "password",
+ "workspaceName": "acme",
+ }
+ },
+ "sourceConfig": {
+ "config": {
+ "type": "DashboardMetadata",
+ "dashboardFilterPattern": {},
+ "chartFilterPattern": {},
+ "dataModelFilterPattern": {},
+ }
+ },
+ },
+ "sink": {"type": "metadata-rest", "config": {}},
+ "workflowConfig": {
+ "loggerLevel": "DEBUG",
+ "openMetadataServerConfig": {
+ "hostPort": "http://localhost:8585/api",
+ "authProvider": "openmetadata",
+ "securityConfig": {"jwtToken": "test-token"},
+ },
+ },
+}
+
+REPORT = {
+ "token": "report-token",
+ "name": "Revenue Report",
+ "description": "Revenue by account",
+ "_links": {"share": {"href": "/acme/reports/report-token"}},
+}
+
+QUERY = {
+ "token": "query-token",
+ "name": "Revenue Query",
+ "raw_query": "SELECT * FROM analytics.orders",
+ "data_source_id": "source-id",
+}
+
+
+def _embedded(name: str, values: list[dict]) -> dict:
+ return {"_embedded": {name: values}}
+
+
+@pytest.fixture
+def mode_source():
+ mode_api = MagicMock()
+ mode_api.get_all_data_sources.return_value = {
+ "source-id": {
+ "token": "warehouse-token",
+ "name": "Warehouse",
+ "database": "analytics_db",
+ }
+ }
+ mode_api.get_all_queries.return_value = _embedded("queries", [QUERY])
+ mode_api.get_all_charts.return_value = _embedded("charts", [])
+
+ with (
+ patch(
+ "metadata.ingestion.source.dashboard.dashboard_service."
+ "DashboardServiceSource.test_connection"
+ ),
+ patch(
+ "metadata.ingestion.source.dashboard.dashboard_service.get_connection"
+ ) as get_connection,
+ ):
+ get_connection.return_value = mode_api
+ config = OpenMetadataWorkflowConfig.model_validate(MOCK_CONFIG)
+ source = ModeSource.create(
+ MOCK_CONFIG["source"],
+ OpenMetadata(config.workflowConfig.openMetadataServerConfig),
+ )
+
+ source.client = mode_api
+ source.data_sources = mode_api.get_all_data_sources.return_value
+ source.dashboard_source_state = set()
+ source.datamodel_source_state = set()
+ source.chart_source_state = set()
+ source.context.get().__dict__["dashboard_service"] = "mock_mode"
+ source.context.get().__dict__["charts"] = []
+ source.context.get().__dict__["dataModels"] = []
+ source.status = MagicMock()
+ return source
+
+
+def _details(mode_source, queries=None):
+ if queries is not None:
+ mode_source.client.get_all_queries.return_value = _embedded("queries", queries)
+ return mode_source.get_dashboard_details(REPORT)
+
+
+def _lineage_entities():
+ data_model = DashboardDataModel.model_construct(
+ id=Uuid("6e781e63-e30f-4c6e-891a-389f1f982cab"),
+ fullyQualifiedName=FullyQualifiedEntityName(
+ 'mock_mode."report-token.query-token"'
+ ),
+ )
+ dashboard = Dashboard.model_construct(
+ id=Uuid("4248eaa4-2183-4bc4-980a-26893311676f"),
+ fullyQualifiedName=FullyQualifiedEntityName("mock_mode.report-token"),
+ )
+ table = Table.model_construct(
+ id=Uuid("b9553fd0-408d-45aa-b38a-43e72ec731ee"),
+ fullyQualifiedName=FullyQualifiedEntityName(
+ "warehouse.analytics_db.analytics.orders"
+ ),
+ )
+ return data_model, dashboard, table
+
+
+class TestModeQueryMetadata:
+ def test_dashboard_name_falls_back_to_token(self, mode_source):
+ assert (
+ mode_source.get_dashboard_name({**REPORT, "name": None}) == "report-token"
+ )
+
+ def test_query_becomes_registered_dashboard_data_model(self, mode_source):
+ result = list(mode_source.yield_datamodel(_details(mode_source)))
+
+ assert len(result) == 1
+ data_model = result[0].right
+ assert data_model.name.root == "report-token.query-token"
+ assert data_model.displayName == "Revenue Query"
+ assert data_model.dataModelType.value == "ModeDataModel"
+ assert data_model.service.root == "mock_mode"
+ assert data_model.serviceType.value == "Mode"
+ assert data_model.sql.root == QUERY["raw_query"]
+ assert data_model.columns == []
+ assert data_model.sourceUrl.root == (
+ "https://app.mode.com/acme/reports/report-token/queries/query-token"
+ )
+ assert len(mode_source.datamodel_source_state) == 1
+
+ def test_query_name_falls_back_to_token(self, mode_source):
+ unnamed_query = {**QUERY, "name": None}
+
+ result = list(
+ mode_source.yield_datamodel(_details(mode_source, [unnamed_query]))
+ )
+
+ assert result[0].right.displayName == "query-token"
+
+ def test_data_model_filter_uses_query_name(self, mode_source):
+ mode_source.source_config.dataModelFilterPattern = FilterPattern(
+ excludes=["Revenue Query"]
+ )
+
+ result = list(mode_source.yield_datamodel(_details(mode_source)))
+
+ assert result == []
+ mode_source.status.filter.assert_called_once_with(
+ "Revenue Query", "Data model filtered out."
+ )
+
+ def test_include_data_models_false_skips_queries(self, mode_source):
+ mode_source.source_config.includeDataModels = False
+
+ assert list(mode_source.yield_datamodel(_details(mode_source))) == []
+ assert mode_source.datamodel_source_state == set()
+
+ def test_dashboard_contains_only_emitted_query_models(self, mode_source):
+ details = _details(mode_source)
+ mode_source.context.get().__dict__["dataModels"] = ["report-token.query-token"]
+
+ dashboard = next(iter(mode_source.yield_dashboard(details))).right
+
+ assert [value.root for value in dashboard.dataModels] == [
+ 'mock_mode.model."report-token.query-token"'
+ ]
+
+ def test_dashboard_omits_data_models_when_disabled(self, mode_source):
+ mode_source.source_config.includeDataModels = False
+ details = _details(mode_source)
+ mode_source.context.get().__dict__["dataModels"] = ["report-token.query-token"]
+
+ dashboard = next(iter(mode_source.yield_dashboard(details))).right
+
+ assert dashboard.dataModels is None
+
+ def test_report_without_queries_is_still_ingested(self, mode_source):
+ details = _details(mode_source, [])
+
+ assert list(mode_source.yield_datamodel(details)) == []
+ assert list(mode_source.yield_dashboard_chart(details)) == []
+ dashboard = next(iter(mode_source.yield_dashboard(details))).right
+ assert dashboard.name.root == "report-token"
+ assert dashboard.dataModels == []
+
+ def test_charts_use_the_queries_already_loaded_for_the_report(self, mode_source):
+ mode_source.client.get_all_charts.return_value = _embedded(
+ "charts",
+ [
+ {
+ "token": "chart-token",
+ "view_vegas": {"title": "Revenue by Account"},
+ "_links": {
+ "report_viz_web": {
+ "href": "/acme/reports/report-token/charts/chart-token"
+ }
+ },
+ }
+ ],
+ )
+ details = _details(mode_source)
+
+ chart = next(iter(mode_source.yield_dashboard_chart(details))).right
+
+ assert chart.name.root == "chart-token"
+ assert chart.displayName == "Revenue by Account"
+ assert chart.sourceUrl.root.endswith(
+ "/acme/reports/report-token/charts/chart-token"
+ )
+ mode_source.client.get_all_queries.assert_called_once()
+
+ def test_chart_stage_reports_missing_query_token_and_continues(self, mode_source):
+ query_without_token = {**QUERY, "token": None}
+ mode_source.client.get_all_charts.return_value = _embedded("charts", [])
+
+ result = list(
+ mode_source.yield_dashboard_chart(
+ _details(mode_source, [query_without_token, QUERY])
+ )
+ )
+
+ assert result[0].left.name == "Revenue Query"
+ assert (
+ result[0].left.error
+ == "Mode query is missing its token; charts could not be fetched"
+ )
+ mode_source.client.get_all_charts.assert_called_once_with(
+ workspace_name="acme",
+ report_token="report-token",
+ query_token="query-token",
+ )
+
+ def test_chart_stage_handles_failed_charts_request(self, mode_source):
+ mode_source.client.get_all_charts.return_value = None
+
+ assert list(mode_source.yield_dashboard_chart(_details(mode_source))) == []
+
+ def test_queries_are_fetched_once_for_all_report_stages(self, mode_source):
+ details = _details(mode_source)
+
+ list(mode_source.yield_dashboard_chart(details))
+ list(mode_source.yield_datamodel(details))
+ list(mode_source.yield_dashboard_lineage_details(details))
+
+ mode_source.client.get_all_queries.assert_called_once_with(
+ workspace_name="acme", report_token="report-token"
+ )
+
+
+class TestModeQueryLineage:
+ @staticmethod
+ def _prepare_metadata(mode_source, data_model_available=True):
+ data_model, dashboard, table = _lineage_entities()
+
+ def get_by_name(entity, **_):
+ if entity is DashboardDataModel:
+ return data_model if data_model_available else None
+ if entity is Dashboard:
+ return dashboard
+ return None
+
+ mode_source.metadata = MagicMock()
+ mode_source.metadata.get_by_name = MagicMock(side_effect=get_by_name)
+ mode_source.metadata.search_in_any_service = MagicMock(return_value=[table])
+ return data_model, dashboard, table
+
+ @staticmethod
+ def _run_lineage(mode_source):
+ parser = SimpleNamespace(
+ query_hash="query-hash",
+ source_tables=["analytics.orders"],
+ )
+ with patch(
+ "metadata.ingestion.source.dashboard.mode.metadata.LineageParser",
+ return_value=parser,
+ ):
+ return list(
+ mode_source.yield_dashboard_lineage_details(_details(mode_source))
+ )
+
+ def test_table_lineage_targets_query_data_model_with_sql(self, mode_source):
+ data_model, _, table = self._prepare_metadata(mode_source)
+ mode_source.context.get().__dict__["dataModels"] = ["report-token.query-token"]
+
+ result = self._run_lineage(mode_source)
+
+ assert len(result) == 1
+ edge = result[0].right.edge
+ assert edge.fromEntity.id == table.id
+ assert edge.toEntity.id == data_model.id
+ assert edge.toEntity.type == "dashboardDataModel"
+ assert edge.lineageDetails.sqlQuery.root == QUERY["raw_query"]
+
+ def test_lineage_targets_dashboard_when_data_models_disabled(self, mode_source):
+ _, dashboard, _ = self._prepare_metadata(mode_source)
+ mode_source.source_config.includeDataModels = False
+
+ result = self._run_lineage(mode_source)
+
+ assert result[0].right.edge.toEntity.id == dashboard.id
+ assert result[0].right.edge.toEntity.type == "dashboard"
+
+ def test_lineage_targets_dashboard_when_query_model_was_filtered(self, mode_source):
+ _, dashboard, _ = self._prepare_metadata(mode_source)
+ mode_source.context.get().__dict__["dataModels"] = []
+
+ result = self._run_lineage(mode_source)
+
+ assert result[0].right.edge.toEntity.id == dashboard.id
+
+ def test_lineage_falls_back_when_query_model_is_unavailable(self, mode_source):
+ _, dashboard, _ = self._prepare_metadata(
+ mode_source, data_model_available=False
+ )
+ mode_source.context.get().__dict__["dataModels"] = ["report-token.query-token"]
+
+ result = self._run_lineage(mode_source)
+
+ assert result[0].right.edge.toEntity.id == dashboard.id
+
+ def test_lineage_skips_query_when_database_is_unknown(self, mode_source):
+ self._prepare_metadata(mode_source)
+ mode_source.data_sources["source-id"].pop("database")
+
+ result = self._run_lineage(mode_source)
+
+ assert result == []
+ mode_source.metadata.search_in_any_service.assert_not_called()
diff --git a/openmetadata-spec/src/main/resources/json/schema/entity/data/dashboardDataModel.json b/openmetadata-spec/src/main/resources/json/schema/entity/data/dashboardDataModel.json
index 21a361d6a646..015bde1b5d92 100644
--- a/openmetadata-spec/src/main/resources/json/schema/entity/data/dashboardDataModel.json
+++ b/openmetadata-spec/src/main/resources/json/schema/entity/data/dashboardDataModel.json
@@ -22,6 +22,7 @@
"TableauEmbeddedDatasource",
"SupersetDataModel",
"MetabaseDataModel",
+ "ModeDataModel",
"LookMlView",
"LookMlExplore",
"PowerBIDataModel",
@@ -51,6 +52,9 @@
{
"name": "MetabaseDataModel"
},
+ {
+ "name": "ModeDataModel"
+ },
{
"name": "LookMlView"
},
diff --git a/openmetadata-ui/src/main/resources/ui/public/locales/en-US/Dashboard/Mode.md b/openmetadata-ui/src/main/resources/ui/public/locales/en-US/Dashboard/Mode.md
index 5fa0d757b8af..1a89a4d2ea28 100644
--- a/openmetadata-ui/src/main/resources/ui/public/locales/en-US/Dashboard/Mode.md
+++ b/openmetadata-ui/src/main/resources/ui/public/locales/en-US/Dashboard/Mode.md
@@ -8,6 +8,15 @@ OpenMetadata relies on Mode's API, which is exclusive to members of the Mode Bus
You can find further information on the Mode connector in the docs.
+## Metadata Mapping
+
+- Mode reports are ingested as dashboards.
+- Report visualizations are ingested as charts.
+- Each query associated with a report is ingested as a dashboard data model. The data model includes the query name, SQL text, and a link to the query in Mode. Mode's query response does not include result-column metadata, so query data models have an empty column list.
+- When the query SQL can be parsed and its source tables can be resolved, lineage is created from the tables to the query data model and then to the report dashboard. If data model ingestion is disabled or a query data model is filtered out, lineage is created directly from the tables to the dashboard.
+
+The connector can only ingest spaces, reports, queries, charts, and data sources visible to the API token. Ensure the token's workspace member has access to every space and report that should be cataloged. See Mode's report API guide and query API reference for the source API behavior.
+
## Connection Details
$$section
@@ -46,10 +55,7 @@ $$
$$section
### Filter Query Param $(id="filterQueryParam")
-This value is the `filter` query parameter that is passed to the Mode API. Different API
-calls use different types of acceptable values. Currently this parameter is only implemented
-to list all collections.
-The valid values that is currently supported are: `all`
-and `custom`. If this field is left empty, `all` will be used.
+This value is the `filter` query parameter passed to Mode's spaces API when discovering reports.
+The supported values are `all` and `custom`. If this field is left empty, `all` will be used.
-$$
\ No newline at end of file
+$$
diff --git a/openmetadata-ui/src/main/resources/ui/src/generated/api/data/createDashboardDataModel.ts b/openmetadata-ui/src/main/resources/ui/src/generated/api/data/createDashboardDataModel.ts
index 21966000bf6f..ebd5f516cac1 100644
--- a/openmetadata-ui/src/main/resources/ui/src/generated/api/data/createDashboardDataModel.ts
+++ b/openmetadata-ui/src/main/resources/ui/src/generated/api/data/createDashboardDataModel.ts
@@ -752,6 +752,7 @@ export enum DataModelType {
LookMlView = "LookMlView",
MetabaseDataModel = "MetabaseDataModel",
MicroStrategyDataset = "MicroStrategyDataset",
+ ModeDataModel = "ModeDataModel",
PowerBIDataFlow = "PowerBIDataFlow",
PowerBIDataModel = "PowerBIDataModel",
PowerBIDatamart = "PowerBIDatamart",
diff --git a/openmetadata-ui/src/main/resources/ui/src/generated/entity/data/dashboardDataModel.ts b/openmetadata-ui/src/main/resources/ui/src/generated/entity/data/dashboardDataModel.ts
index fbf780dbcb0a..2f143bc56ffe 100644
--- a/openmetadata-ui/src/main/resources/ui/src/generated/entity/data/dashboardDataModel.ts
+++ b/openmetadata-ui/src/main/resources/ui/src/generated/entity/data/dashboardDataModel.ts
@@ -896,6 +896,7 @@ export enum DataModelType {
LookMlView = "LookMlView",
MetabaseDataModel = "MetabaseDataModel",
MicroStrategyDataset = "MicroStrategyDataset",
+ ModeDataModel = "ModeDataModel",
PowerBIDataFlow = "PowerBIDataFlow",
PowerBIDataModel = "PowerBIDataModel",
PowerBIDatamart = "PowerBIDatamart",
From ad466b81c63da3e80b5f73c48fb26a0af0d1ad3f Mon Sep 17 00:00:00 2001
From: harshsoni2024
Date: Thu, 10 Sep 2026 20:26:33 +0530
Subject: [PATCH 2/3] Address Mode report pagination review feedback
Report pagination terminated on `len(reports) < REPORTS_PAGE_SIZE`, which made
two assumptions about Mode's API that it does not guarantee.
If a page ever holds fewer than 30 records, the first page satisfies the
condition and the rest of the space is dropped with no warning - reproducing
the missing-reports symptom this change set exists to fix. If a space holds an
exact multiple of 30 records and Mode clamps the out-of-range page back to the
last one, the repeated-page guard raised and killed the whole source, dropping
every dashboard of every space.
Pagination now stops once a page carries no report that was not already seen in
that space, which terminates on an empty page, on a clamped repeat, and on an
ignored page parameter, without assuming how many records a full page holds.
Reports are de-duplicated by token, falling back to the record itself when Mode
omits one.
Co-Authored-By: Claude Opus 5 (1M context)
---
.../ingestion/source/dashboard/mode/client.py | 38 +++++++---
.../unit/source/dashboard/mode/test_client.py | 71 +++++++++++++------
2 files changed, 76 insertions(+), 33 deletions(-)
diff --git a/ingestion/src/metadata/ingestion/source/dashboard/mode/client.py b/ingestion/src/metadata/ingestion/source/dashboard/mode/client.py
index 395a110b71eb..1a08cd6c48e1 100644
--- a/ingestion/src/metadata/ingestion/source/dashboard/mode/client.py
+++ b/ingestion/src/metadata/ingestion/source/dashboard/mode/client.py
@@ -11,6 +11,8 @@
"""
REST Auth & Client for Mode
"""
+
+import json
import traceback
from base64 import b64encode
from typing import Any, Dict, List, Optional, cast
@@ -38,7 +40,14 @@
LINKS = "_links"
SHARE = "share"
HREF = "href"
-REPORTS_PAGE_SIZE = 30
+
+
+def _report_key(report: Dict[str, Any]) -> str:
+ """Identify a report for de-duplication, tolerating a missing token."""
+ token = report.get(TOKEN)
+ if token:
+ return str(token)
+ return json.dumps(report, sort_keys=True, default=str)
class ModeApiClient:
@@ -69,11 +78,18 @@ def fetch_all_reports(
self, workspace_name: str, filter: Optional[str] = "all"
) -> List[Dict[str, Any]]:
"""Method to fetch all reports for Mode
+
+ Mode neither documents a stable report page size nor guarantees an empty page
+ past the last one: an out-of-range page may be clamped back to the last page,
+ and a page parameter the API ignores repeats page one indefinitely. Pagination
+ therefore stops once a page carries no unseen report, which terminates in all
+ of those cases without assuming how many records a full page holds.
+
Args:
workspace_name:
filter:
Returns:
- dict
+ the report records of every visible space
"""
if filter not in ["custom", "all"]:
raise ValueError(
@@ -88,8 +104,8 @@ def fetch_all_reports(
)
spaces = response_spaces[EMBEDDED][SPACES]
for space in spaces:
+ seen_reports = set()
page = 1
- previous_reports = None
while True:
response_reports = self.get_reports_for_space(
workspace_name=workspace_name,
@@ -97,15 +113,15 @@ def fetch_all_reports(
page=page,
)
reports = response_reports[EMBEDDED][REPORTS]
- if reports and reports == previous_reports:
- raise RuntimeError(
- f"Mode returned the same report page twice for space "
- f"[{space[TOKEN]}] at page [{page}]"
- )
- all_reports.extend(reports)
- if len(reports) < REPORTS_PAGE_SIZE:
+ new_reports = [
+ report
+ for report in reports
+ if _report_key(report) not in seen_reports
+ ]
+ if not new_reports:
break
- previous_reports = reports
+ seen_reports.update(_report_key(report) for report in new_reports)
+ all_reports.extend(new_reports)
page += 1
return all_reports
diff --git a/ingestion/tests/unit/source/dashboard/mode/test_client.py b/ingestion/tests/unit/source/dashboard/mode/test_client.py
index 9cbe359df122..87d9167c93fe 100644
--- a/ingestion/tests/unit/source/dashboard/mode/test_client.py
+++ b/ingestion/tests/unit/source/dashboard/mode/test_client.py
@@ -17,8 +17,8 @@
from metadata.ingestion.source.dashboard.mode.client import ModeApiClient
-def _reports(prefix: str, count: int) -> list[dict]:
- return [{"token": f"{prefix}-{index}"} for index in range(count)]
+def _reports(prefix: str, count: int, start: int = 0) -> list[dict]:
+ return [{"token": f"{prefix}-{index}"} for index in range(start, start + count)]
def _embedded(name: str, values: list[dict]) -> dict:
@@ -40,7 +40,9 @@ def test_fetch_all_reports_paginates_every_space(mode_client):
_embedded("spaces", [{"token": "finance"}, {"token": "operations"}]),
_embedded("reports", first_space_page),
_embedded("reports", second_space_page),
+ _embedded("reports", []),
_embedded("reports", operations_page),
+ _embedded("reports", []),
]
reports = mode_client.fetch_all_reports("acme", "custom")
@@ -50,26 +52,56 @@ def test_fetch_all_reports_paginates_every_space(mode_client):
call("/acme/spaces?filter=custom"),
call("/acme/spaces/finance/reports?page=1"),
call("/acme/spaces/finance/reports?page=2"),
+ call("/acme/spaces/finance/reports?page=3"),
call("/acme/spaces/operations/reports?page=1"),
+ call("/acme/spaces/operations/reports?page=2"),
]
-def test_fetch_all_reports_requests_page_after_exactly_thirty_results(mode_client):
- first_page = _reports("report", 30)
+def test_fetch_all_reports_keeps_paging_after_a_page_smaller_than_thirty(mode_client):
+ """A page size below 30 must not be mistaken for the end of the space."""
+ first_page = _reports("report", 25)
+ second_page = _reports("report", 25, start=25)
mode_client.client.get.side_effect = [
_embedded("spaces", [{"token": "space-token"}]),
_embedded("reports", first_page),
+ _embedded("reports", second_page),
_embedded("reports", []),
]
- reports = mode_client.fetch_all_reports("acme")
+ assert mode_client.fetch_all_reports("acme") == first_page + second_page
+
+
+def test_fetch_all_reports_stops_when_a_page_repeats(mode_client):
+ """Mode may clamp an out-of-range page, or ignore the page parameter entirely."""
+ repeated_page = _reports("report", 30)
+ mode_client.client.get.side_effect = [
+ _embedded("spaces", [{"token": "space-token"}]),
+ _embedded("reports", repeated_page),
+ _embedded("reports", repeated_page),
+ ]
- assert reports == first_page
+ assert mode_client.fetch_all_reports("acme") == repeated_page
assert mode_client.client.get.call_args_list[-1] == call(
"/acme/spaces/space-token/reports?page=2"
)
+def test_fetch_all_reports_keeps_only_the_first_copy_of_an_overlapping_report(
+ mode_client,
+):
+ first_page = _reports("report", 30)
+ overlapping_page = _reports("report", 4, start=28)
+ mode_client.client.get.side_effect = [
+ _embedded("spaces", [{"token": "space-token"}]),
+ _embedded("reports", first_page),
+ _embedded("reports", overlapping_page),
+ _embedded("reports", []),
+ ]
+
+ assert mode_client.fetch_all_reports("acme") == first_page + overlapping_page[2:]
+
+
def test_fetch_all_reports_propagates_later_page_failure(mode_client):
mode_client.client.get.side_effect = [
_embedded("spaces", [{"token": "space-token"}]),
@@ -88,22 +120,6 @@ def test_fetch_all_reports_rejects_invalid_filter(mode_client):
mode_client.client.get.assert_not_called()
-def test_fetch_all_reports_rejects_repeated_full_page(mode_client):
- repeated_page = _reports("report", 30)
- mode_client.client.get.side_effect = [
- _embedded("spaces", [{"token": "space-token"}]),
- _embedded("reports", repeated_page),
- _embedded("reports", repeated_page),
- ]
-
- with pytest.raises(RuntimeError, match="same report page twice"):
- mode_client.fetch_all_reports("acme")
-
- assert mode_client.client.get.call_args_list[-1] == call(
- "/acme/spaces/space-token/reports?page=2"
- )
-
-
def test_fetch_all_reports_allows_distinct_tokenless_pages(mode_client):
first_page = [{"name": f"first-{index}"} for index in range(30)]
second_page = [{"name": f"second-{index}"} for index in range(30)]
@@ -115,3 +131,14 @@ def test_fetch_all_reports_allows_distinct_tokenless_pages(mode_client):
]
assert mode_client.fetch_all_reports("acme") == first_page + second_page
+
+
+def test_fetch_all_reports_stops_when_a_tokenless_page_repeats(mode_client):
+ repeated_page = [{"name": f"report-{index}"} for index in range(30)]
+ mode_client.client.get.side_effect = [
+ _embedded("spaces", [{"token": "space-token"}]),
+ _embedded("reports", repeated_page),
+ _embedded("reports", repeated_page),
+ ]
+
+ assert mode_client.fetch_all_reports("acme") == repeated_page
From 052b4309f6c7bdd7930b86161d75093837d29cf4 Mon Sep 17 00:00:00 2001
From: harshsoni2024
Date: Fri, 11 Sep 2026 10:19:25 +0530
Subject: [PATCH 3/3] Make the Mode unit test directory a package
Without `__init__.py`, pytest puts the test file's own directory on sys.path
and imports it under its bare basename, so the new Mode test claimed the
top-level module name `test_client` and collection of the pre-existing
`tests/unit/source/database/burstiq/test_client.py` failed with an import file
mismatch. `tests/unit/source/dashboard/qlikcloud` already carries the same
marker for the same reason, and the 1.13 backport dropped the one the change
carries on main.
Co-Authored-By: Claude Opus 5 (1M context)
---
ingestion/tests/unit/source/dashboard/mode/__init__.py | 10 ++++++++++
1 file changed, 10 insertions(+)
create mode 100644 ingestion/tests/unit/source/dashboard/mode/__init__.py
diff --git a/ingestion/tests/unit/source/dashboard/mode/__init__.py b/ingestion/tests/unit/source/dashboard/mode/__init__.py
new file mode 100644
index 000000000000..c87f25e5c64e
--- /dev/null
+++ b/ingestion/tests/unit/source/dashboard/mode/__init__.py
@@ -0,0 +1,10 @@
+# Copyright 2025 Collate
+# Licensed under the Collate Community License, Version 1.0 (the "License");
+# you may not use this file except in compliance with the License.
+# You may obtain a copy of the License at
+# https://github.com/open-metadata/OpenMetadata/blob/main/ingestion/LICENSE
+# 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.