项目文件夹

文件
wehub-resource-sync bf2343b7e4
Integration Tests - MySQL + Elasticsearch / Detect Changes (push) Has been cancelled
Integration Tests - MySQL + Elasticsearch / integration-tests-mysql-elasticsearch (push) Has been cancelled
Integration Tests - PostgreSQL + Elasticsearch + Redis / Detect Changes (push) Has been cancelled
Integration Tests - PostgreSQL + Elasticsearch + Redis / integration-tests-postgres-elasticsearch-redis (push) Has been cancelled
Integration Tests - PostgreSQL + OpenSearch / Detect Changes (push) Has been cancelled
Integration Tests - PostgreSQL + OpenSearch / integration-tests-postgres-opensearch (push) Has been cancelled
Java Checkstyle / java-checkstyle (push) Has been cancelled
Maven Collate Tests / maven-collate-ci (push) Has been cancelled
OpenMetadata Service Unit Tests / openmetadata-service-unit-tests-status (push) Has been cancelled
Publish Package to Maven Central Repository / publish-maven-packages (push) Has been cancelled
OpenMetadata Service Unit Tests / Detect Changes (push) Has been cancelled
OpenMetadata Service Unit Tests / openmetadata-service-unit-tests (push) Has been cancelled
OpenMetadata Service Unit Tests / k8s_operator-unit-tests (push) Has been cancelled
chore: import upstream snapshot with attribution
2026-07-13 13:35:45 +08:00

85 行
3.3 KiB
Python

# 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.
"""Tests for MetadataRestSink.write_barrier dispatcher.
Guards the contract that a Barrier record flushes the bulk buffer
synchronously so subsequent records in the same stream see committed
entities.
"""
from unittest.mock import MagicMock, Mock
import pytest
from metadata.generated.schema.api.data.createDashboardDataModel import (
CreateDashboardDataModelRequest,
)
from metadata.generated.schema.entity.data.dashboardDataModel import DataModelType
from metadata.generated.schema.entity.data.table import Column, DataType
from metadata.generated.schema.type.basic import EntityName, FullyQualifiedEntityName
from metadata.ingestion.models.barrier import Barrier
from metadata.ingestion.sink.metadata_rest import MetadataRestSink, MetadataRestSinkConfig
def _make_data_model(name: str) -> CreateDashboardDataModelRequest:
return CreateDashboardDataModelRequest(
name=EntityName(name),
displayName=name,
service=FullyQualifiedEntityName("test_service"),
dataModelType=DataModelType.QuickSightDataModel,
columns=[Column(name="col1", dataType=DataType.STRING)],
)
def _mock_bulk_success(entities, use_async=False, **kwargs):
result = MagicMock()
result.status.value = "success"
result.numberOfRowsProcessed.root = len(entities)
result.numberOfRowsFailed.root = 0
result.successRequest = entities
result.failedRequest = []
return result
@pytest.fixture
def sink():
mock_metadata = Mock()
mock_metadata.bulk_create_or_update = Mock(side_effect=_mock_bulk_success)
config = MetadataRestSinkConfig(bulk_sink_batch_size=10)
return MetadataRestSink(config, mock_metadata)
class TestBarrierDispatcher:
"""write_barrier must flush the buffer when non-empty and be a no-op when empty."""
def test_barrier_flushes_non_empty_buffer(self, sink):
"""A Barrier on a non-empty buffer triggers bulk_create_or_update."""
sink.write_create_request(_make_data_model("dm-1"))
sink.write_create_request(_make_data_model("dm-2"))
assert len(sink.buffer) == 2
sink.write_barrier(Barrier(reason="test"))
sink.metadata.bulk_create_or_update.assert_called_once()
# Buffer should be empty after flush
assert len(sink.buffer) == 0
def test_barrier_on_empty_buffer_is_noop(self, sink):
"""A Barrier on an empty buffer must not call bulk_create_or_update."""
assert len(sink.buffer) == 0
result = sink.write_barrier(Barrier(reason="empty"))
sink.metadata.bulk_create_or_update.assert_not_called()
# Returns Either(right=None) — protocol-conformant
assert result is not None
assert result.right is None