This is an automated email from the ASF dual-hosted git repository.

terrymanu pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/shardingsphere.git


The following commit(s) were added to refs/heads/master by this push:
     new 996c854c2da Fix CDC selector parsing and job startup validation 
(#39612)
996c854c2da is described below

commit 996c854c2da02ce322050a37d36c8b260b0991ac
Author: Liang Zhang <[email protected]>
AuthorDate: Thu Aug 27 11:14:04 2026 +0800

    Fix CDC selector parsing and job startup validation (#39612)
    
    * Add code implementation Skill and non-regression policy gates
    
    * Add code implementation Skill and non-regression policy gates
    
    * Fix CDC selector parsing and job startup validation
---
 .../data/pipeline/cdc/api/CDCJobAPI.java           | 41 ++++++++--
 .../pipeline/cdc/handler/CDCBackendHandler.java    |  1 -
 .../pipeline/cdc/util/CDCSchemaTableUtils.java     | 12 +--
 .../data/pipeline/cdc/api/CDCJobAPITest.java       | 89 +++++++++++++++++++++-
 .../cdc/handler/CDCBackendHandlerTest.java         | 45 +++++++----
 .../pipeline/cdc/util/CDCSchemaTableUtilsTest.java | 22 ++++++
 .../frontend/netty/CDCChannelInboundHandler.java   | 16 ++--
 .../netty/CDCChannelInboundHandlerTest.java        | 49 +++++++++---
 8 files changed, 226 insertions(+), 49 deletions(-)

diff --git 
a/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/api/CDCJobAPI.java
 
b/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/api/CDCJobAPI.java
index 6d251ee5631..839f8ce8724 100644
--- 
a/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/api/CDCJobAPI.java
+++ 
b/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/api/CDCJobAPI.java
@@ -138,8 +138,7 @@ public final class CDCJobAPI implements TransmissionJobAPI {
         if 
(governanceFacade.getJobFacade().getConfiguration().isExisted(jobConfig.getJobId()))
 {
             log.warn("CDC job already exists in registry center, ignore, job 
id is `{}`", jobConfig.getJobId());
         } else {
-            checkDataSources(jobConfig);
-            checkSchemaTableNames(jobConfig.getSchemaTableNames());
+            checkJobConfiguration(jobConfig);
             
governanceFacade.getJobFacade().getJob().create(jobConfig.getJobId(), 
jobType.getOption().getJobClass());
             JobConfigurationPOJO jobConfigPOJO = 
jobConfigManager.convertToJobConfigurationPOJO(jobConfig);
             jobConfigPOJO.setDisabled(true);
@@ -151,6 +150,11 @@ public final class CDCJobAPI implements TransmissionJobAPI 
{
         return jobConfig.getJobId();
     }
     
+    private void checkJobConfiguration(final CDCJobConfiguration jobConfig) {
+        checkDataSources(jobConfig);
+        checkSchemaTableNames(jobConfig.getSchemaTableNames());
+    }
+    
     private void checkDataSources(final CDCJobConfiguration jobConfig) {
         Map<String, Map<String, Object>> dataSources = 
jobConfig.getDataSourceConfig().getRootConfig().getDataSources();
         for (DataNode each : getDataNodes(jobConfig)) {
@@ -255,17 +259,40 @@ public final class CDCJobAPI implements 
TransmissionJobAPI {
      * @param sink sink
      */
     public void start(final String jobId, final PipelineSink sink) {
+        CDCJobConfiguration jobConfig = 
jobConfigManager.getJobConfiguration(jobId);
+        try {
+            checkJobConfiguration(jobConfig);
+        } catch (final PipelineInvalidParameterException ex) {
+            try {
+                PipelineJobRegistry.stop(jobId);
+                // CHECKSTYLE:OFF
+            } catch (final RuntimeException cleanupException) {
+                // CHECKSTYLE:ON
+                ex.addSuppressed(cleanupException);
+            }
+            try {
+                JobConfigurationPOJO jobConfigPOJO = 
PipelineJobIdUtils.getElasticJobConfigurationPOJO(jobId);
+                if (!jobConfigPOJO.isDisabled()) {
+                    disable(jobConfigPOJO);
+                }
+                // CHECKSTYLE:OFF
+            } catch (final RuntimeException cleanupException) {
+                // CHECKSTYLE:ON
+                ex.addSuppressed(cleanupException);
+            }
+            throw ex;
+        }
+        PipelineJobRegistry.stop(jobId);
         CDCJob job = new CDCJob(sink);
         PipelineJobRegistry.add(jobId, job);
-        enable(jobId);
         JobConfigurationPOJO jobConfigPOJO = 
PipelineJobIdUtils.getElasticJobConfigurationPOJO(jobId);
+        enable(jobConfigPOJO);
         OneOffJobBootstrap oneOffJobBootstrap = new 
OneOffJobBootstrap(PipelineAPIFactory.getRegistryCenter(PipelineJobIdUtils.parseContextKey(jobId)),
 job, jobConfigPOJO.toJobConfiguration());
         job.getJobRunnerManager().setJobBootstrap(oneOffJobBootstrap);
         oneOffJobBootstrap.execute();
     }
     
-    private void enable(final String jobId) {
-        JobConfigurationPOJO jobConfigPOJO = 
PipelineJobIdUtils.getElasticJobConfigurationPOJO(jobId);
+    private void enable(final JobConfigurationPOJO jobConfigPOJO) {
         jobConfigPOJO.setDisabled(false);
         jobConfigPOJO.getProps().setProperty("start_time_millis", 
String.valueOf(System.currentTimeMillis()));
         jobConfigPOJO.getProps().remove("stop_time");
@@ -279,6 +306,10 @@ public final class CDCJobAPI implements TransmissionJobAPI 
{
      */
     public void disable(final String jobId) {
         JobConfigurationPOJO jobConfigPOJO = 
PipelineJobIdUtils.getElasticJobConfigurationPOJO(jobId);
+        disable(jobConfigPOJO);
+    }
+    
+    private void disable(final JobConfigurationPOJO jobConfigPOJO) {
         jobConfigPOJO.setDisabled(true);
         jobConfigPOJO.getProps().setProperty("stop_time", 
LocalDateTime.now().format(DateTimeFormatterFactory.getDatetimeFormatter()));
         
PipelineAPIFactory.getJobConfigurationAPI(PipelineJobIdUtils.parseContextKey(jobConfigPOJO.getJobName())).updateJobConfiguration(jobConfigPOJO);
diff --git 
a/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/handler/CDCBackendHandler.java
 
b/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/handler/CDCBackendHandler.java
index 7167f58dd79..c9a34c11a99 100644
--- 
a/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/handler/CDCBackendHandler.java
+++ 
b/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/handler/CDCBackendHandler.java
@@ -123,7 +123,6 @@ public final class CDCBackendHandler {
     public void startStreaming(final String jobId, final CDCConnectionContext 
connectionContext, final Channel channel) {
         CDCJobConfiguration cdcJobConfig = 
jobConfigManager.getJobConfiguration(jobId);
         ShardingSpherePreconditions.checkNotNull(cdcJobConfig, () -> new 
PipelineJobNotFoundException(jobId));
-        PipelineJobRegistry.stop(jobId);
         ShardingSphereDatabase database = 
PipelineContextManager.getProxyContext().getMetaDataContexts().getMetaData().getDatabase(cdcJobConfig.getDatabaseName());
         jobAPI.start(jobId, new PipelineCDCSocketSink(channel, database, 
cdcJobConfig.getSchemaTableNames()));
         connectionContext.setJobId(jobId);
diff --git 
a/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/util/CDCSchemaTableUtils.java
 
b/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/util/CDCSchemaTableUtils.java
index ba986b053bb..771985bfb70 100644
--- 
a/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/util/CDCSchemaTableUtils.java
+++ 
b/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/util/CDCSchemaTableUtils.java
@@ -69,15 +69,13 @@ public final class CDCSchemaTableUtils {
         Map<String, Set<String>> result = new HashMap<>();
         for (SchemaTable each : schemaTables) {
             if ("*".equals(each.getSchema())) {
-                result.putAll(parseTableExpressionWithAllSchema(database, 
systemSchemas, each));
+                mergeSchemaTables(result, 
parseTableExpressionWithAllSchema(database, systemSchemas, each));
             } else if ("*".equals(each.getTable())) {
                 result.putAll(parseTableExpressionWithAllTable(database, 
each));
             } else {
-                String schemaName = each.getSchema();
-                if (schemaName.isEmpty()) {
-                    schemaName = 
dialectDatabaseMetaData.getSchemaOption().getDefaultSchema().orElse("");
-                }
+                String schemaName = each.getSchema().isEmpty() ? 
dialectDatabaseMetaData.getSchemaOption().getDefaultSchema().orElse("") : 
each.getSchema();
                 ShardingSphereSchema schema = database.getSchema(new 
IdentifierValue(schemaName));
+                ShardingSpherePreconditions.checkNotNull(schema, () -> new 
SchemaNotFoundException(schemaName));
                 
ShardingSpherePreconditions.checkNotNull(schema.getTable(each.getTable()), () 
-> new TableNotFoundException(each.getTable()));
                 result.computeIfAbsent(schema.getName(), ignored -> new 
HashSet<>()).add(each.getTable());
             }
@@ -85,6 +83,10 @@ public final class CDCSchemaTableUtils {
         return result;
     }
     
+    private static void mergeSchemaTables(final Map<String, Set<String>> 
target, final Map<String, Set<String>> source) {
+        source.forEach((schemaName, tableNames) -> 
target.computeIfAbsent(schemaName, ignored -> new 
HashSet<>()).addAll(tableNames));
+    }
+    
     private static Map<String, Set<String>> 
parseTableExpressionWithAllTables(final ShardingSphereDatabase database, final 
Collection<String> systemSchemas) {
         Map<String, Set<String>> result = new 
HashMap<>(database.getAllSchemas().size(), 1F);
         for (ShardingSphereSchema schema : database.getAllSchemas()) {
diff --git 
a/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/api/CDCJobAPITest.java
 
b/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/api/CDCJobAPITest.java
index f23a98eb03d..0444a0e10e5 100644
--- 
a/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/api/CDCJobAPITest.java
+++ 
b/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/api/CDCJobAPITest.java
@@ -103,6 +103,7 @@ import static org.hamcrest.Matchers.containsString;
 import static org.hamcrest.Matchers.emptyString;
 import static org.hamcrest.Matchers.hasSize;
 import static org.hamcrest.Matchers.is;
+import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -121,7 +122,7 @@ import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
 @ExtendWith(AutoMockExtension.class)
-@StaticMockSettings({PipelineAPIFactory.class, PipelineJobIdUtils.class, 
PipelineJobRegistry.class})
+@StaticMockSettings({PipelineAPIFactory.class, PipelineJobIdUtils.class})
 @MockitoSettings(strictness = Strictness.LENIENT)
 class CDCJobAPITest {
     
@@ -427,19 +428,30 @@ class CDCJobAPITest {
     }
     
     @Test
-    void assertStartEnableDisableAndType() {
+    void assertStartEnableDisableAndType() throws ReflectiveOperationException 
{
+        PipelineJobConfigurationManager jobConfigManager = 
mockPersistedJobConfiguration(createJobConfiguration(1));
         JobConfigurationPOJO jobConfigPOJO = createJobConfigurationPOJO();
         jobConfigPOJO.setShardingTotalCount(1);
+        jobConfigPOJO.setDisabled(true);
+        jobConfigPOJO.getProps().setProperty("stop_time", "2026-08-26 
00:00:00");
         
when(PipelineJobIdUtils.getElasticJobConfigurationPOJO("foo_job")).thenReturn(jobConfigPOJO);
         JobConfigurationAPI jobConfigAPI = mock(JobConfigurationAPI.class);
         
when(PipelineAPIFactory.getJobConfigurationAPI(any())).thenReturn(jobConfigAPI);
         
when(PipelineAPIFactory.getRegistryCenter(any())).thenReturn(mock(CoordinatorRegistryCenter.class));
         PipelineSink sink = mock(PipelineSink.class);
-        try (MockedConstruction<OneOffJobBootstrap> jobBootstrapConstruction = 
mockConstruction(OneOffJobBootstrap.class)) {
+        try (
+                MockedStatic<PipelineJobRegistry> jobRegistryMocked = 
mockStatic(PipelineJobRegistry.class);
+                MockedConstruction<OneOffJobBootstrap> 
jobBootstrapConstruction = mockConstruction(OneOffJobBootstrap.class)) {
             jobAPI.start("foo_job", sink);
-            
assertThat(jobConfigPOJO.getProps().getProperty("start_time_millis"), 
is(jobConfigPOJO.getProps().getProperty("start_time_millis")));
+            verify(jobConfigManager).getJobConfiguration("foo_job");
+            jobRegistryMocked.verify(() -> 
PipelineJobRegistry.stop("foo_job"));
+            jobRegistryMocked.verify(() -> 
PipelineJobRegistry.add(eq("foo_job"), any()));
+            assertFalse(jobConfigPOJO.isDisabled());
+            
assertNotNull(jobConfigPOJO.getProps().getProperty("start_time_millis"));
+            assertFalse(jobConfigPOJO.getProps().containsKey("stop_time"));
             verify(jobConfigAPI).updateJobConfiguration(jobConfigPOJO);
             jobAPI.disable("foo_job");
+            assertTrue(jobConfigPOJO.isDisabled());
             assertNotNull(jobConfigPOJO.getProps().getProperty("stop_time"));
             jobAPI.commit("foo_job");
             jobAPI.rollback("foo_job");
@@ -448,6 +460,75 @@ class CDCJobAPITest {
         }
     }
     
+    @Test
+    void assertStartRejectsPersistedUnresolvedDataSource() throws 
ReflectiveOperationException {
+        
mockPersistedJobConfiguration(createJobConfiguration(Collections.singletonList("foo_missing_ds"),
 false));
+        JobConfigurationPOJO jobConfigPOJO = createJobConfigurationPOJO();
+        jobConfigPOJO.setDisabled(false);
+        
when(PipelineJobIdUtils.getElasticJobConfigurationPOJO("foo_job")).thenReturn(jobConfigPOJO);
+        JobConfigurationAPI jobConfigAPI = mock(JobConfigurationAPI.class);
+        
when(PipelineAPIFactory.getJobConfigurationAPI(any())).thenReturn(jobConfigAPI);
+        try (
+                MockedStatic<PipelineJobRegistry> jobRegistryMocked = 
mockStatic(PipelineJobRegistry.class);
+                MockedConstruction<OneOffJobBootstrap> 
jobBootstrapConstruction = mockConstruction(OneOffJobBootstrap.class)) {
+            PipelineInvalidParameterException actual = 
assertThrows(PipelineInvalidParameterException.class, () -> 
jobAPI.start("foo_job", mock(PipelineSink.class)));
+            assertThat(actual.getMessage(), containsString("foo_missing_ds"));
+            jobRegistryMocked.verify(() -> 
PipelineJobRegistry.stop("foo_job"));
+            jobRegistryMocked.verify(() -> 
PipelineJobRegistry.add(anyString(), any()), never());
+            assertTrue(jobConfigPOJO.isDisabled());
+            verify(jobConfigAPI).updateJobConfiguration(jobConfigPOJO);
+            assertTrue(jobBootstrapConstruction.constructed().isEmpty());
+        }
+    }
+    
+    @Test
+    void 
assertStartRejectsPersistedDuplicateBareTableWithoutRepeatedGovernanceWrite() 
throws ReflectiveOperationException {
+        
mockPersistedJobConfiguration(createJobConfiguration(Collections.singletonList("foo_ds"),
 false, Arrays.asList("foo_schema.foo_tbl", "bar_schema.FOO_TBL")));
+        JobConfigurationPOJO jobConfigPOJO = createJobConfigurationPOJO();
+        jobConfigPOJO.setDisabled(true);
+        
when(PipelineJobIdUtils.getElasticJobConfigurationPOJO("foo_job")).thenReturn(jobConfigPOJO);
+        JobConfigurationAPI jobConfigAPI = mock(JobConfigurationAPI.class);
+        
when(PipelineAPIFactory.getJobConfigurationAPI(any())).thenReturn(jobConfigAPI);
+        try (
+                MockedStatic<PipelineJobRegistry> jobRegistryMocked = 
mockStatic(PipelineJobRegistry.class);
+                MockedConstruction<OneOffJobBootstrap> 
jobBootstrapConstruction = mockConstruction(OneOffJobBootstrap.class)) {
+            PipelineInvalidParameterException actual = 
assertThrows(PipelineInvalidParameterException.class, () -> 
jobAPI.start("foo_job", mock(PipelineSink.class)));
+            assertThat(actual.getMessage(), is("There is invalid parameter 
value. More than one schema table has the same table name `FOO_TBL`."));
+            jobRegistryMocked.verify(() -> 
PipelineJobRegistry.stop("foo_job"));
+            jobRegistryMocked.verify(() -> 
PipelineJobRegistry.add(anyString(), any()), never());
+            verify(jobConfigAPI, never()).updateJobConfiguration(any());
+            assertTrue(jobBootstrapConstruction.constructed().isEmpty());
+        }
+    }
+    
+    @Test
+    void assertStartPreservesValidationExceptionWhenCleanupFails() throws 
ReflectiveOperationException {
+        
mockPersistedJobConfiguration(createJobConfiguration(Collections.singletonList("foo_missing_ds"),
 false));
+        JobConfigurationPOJO jobConfigPOJO = createJobConfigurationPOJO();
+        jobConfigPOJO.setDisabled(false);
+        
when(PipelineJobIdUtils.getElasticJobConfigurationPOJO("foo_job")).thenReturn(jobConfigPOJO);
+        JobConfigurationAPI jobConfigAPI = mock(JobConfigurationAPI.class);
+        
when(PipelineAPIFactory.getJobConfigurationAPI(any())).thenReturn(jobConfigAPI);
+        doThrow(new IllegalStateException("governance 
unavailable")).when(jobConfigAPI).updateJobConfiguration(jobConfigPOJO);
+        try (MockedStatic<PipelineJobRegistry> jobRegistryMocked = 
mockStatic(PipelineJobRegistry.class)) {
+            jobRegistryMocked.when(() -> 
PipelineJobRegistry.stop("foo_job")).thenThrow(new IllegalStateException("stop 
unavailable"));
+            PipelineInvalidParameterException actual = 
assertThrows(PipelineInvalidParameterException.class, () -> 
jobAPI.start("foo_job", mock(PipelineSink.class)));
+            assertThat(actual.getMessage(), containsString("foo_missing_ds"));
+            assertThat(actual.getSuppressed().length, is(2));
+            assertThat(actual.getSuppressed()[0].getMessage(), is("stop 
unavailable"));
+            assertThat(actual.getSuppressed()[1].getMessage(), is("governance 
unavailable"));
+            jobRegistryMocked.verify(() -> 
PipelineJobRegistry.stop("foo_job"));
+            assertTrue(jobConfigPOJO.isDisabled());
+        }
+    }
+    
+    private PipelineJobConfigurationManager 
mockPersistedJobConfiguration(final CDCJobConfiguration jobConfig) throws 
ReflectiveOperationException {
+        PipelineJobConfigurationManager result = 
mock(PipelineJobConfigurationManager.class);
+        when(result.getJobConfiguration("foo_job")).thenReturn(jobConfig);
+        
Plugins.getMemberAccessor().set(CDCJobAPI.class.getDeclaredField("jobConfigManager"),
 jobAPI, result);
+        return result;
+    }
+    
     private void putContext(final Map<String, StorageUnit> storageUnits) {
         MetaDataContexts metaDataContexts = mock(MetaDataContexts.class, 
RETURNS_DEEP_STUBS);
         
when(metaDataContexts.getMetaData().getDatabase("foo_db").getResourceMetaData().getStorageUnits()).thenReturn(storageUnits);
diff --git 
a/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/handler/CDCBackendHandlerTest.java
 
b/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/handler/CDCBackendHandlerTest.java
index 44126f1f566..16d31bf0ed9 100644
--- 
a/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/handler/CDCBackendHandlerTest.java
+++ 
b/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/handler/CDCBackendHandlerTest.java
@@ -63,6 +63,7 @@ import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.ExtendWith;
 import org.mockito.ArgumentCaptor;
+import org.mockito.MockedStatic;
 import org.mockito.internal.configuration.plugins.Plugins;
 
 import java.util.ArrayList;
@@ -81,13 +82,14 @@ import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.Mockito.RETURNS_DEEP_STUBS;
 import static org.mockito.Mockito.lenient;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.verifyNoInteractions;
 import static org.mockito.Mockito.when;
 
 @ExtendWith(AutoMockExtension.class)
-@StaticMockSettings({TypedSPILoader.class, DatabaseTypedSPILoader.class, 
PipelineContextManager.class, CDCSchemaTableUtils.class, 
PipelineDataNodeUtils.class, PipelineJobRegistry.class,
-        PipelineJobIdUtils.class, CDCImporterManager.class})
+@StaticMockSettings({TypedSPILoader.class, DatabaseTypedSPILoader.class, 
PipelineContextManager.class, CDCSchemaTableUtils.class, 
PipelineDataNodeUtils.class, PipelineJobIdUtils.class,
+        CDCImporterManager.class})
 class CDCBackendHandlerTest {
     
     private CDCJobAPI jobAPI;
@@ -183,11 +185,15 @@ class CDCBackendHandlerTest {
         mockProxyContext(mock(ShardingSphereDatabase.class));
         Channel channel = mock(Channel.class);
         CDCConnectionContext connectionContext = createConnectionContext();
-        backendHandler.startStreaming("foo_job", connectionContext, channel);
-        ArgumentCaptor<PipelineSink> sinkCaptor = 
ArgumentCaptor.forClass(PipelineSink.class);
-        verify(jobAPI).start(eq("foo_job"), sinkCaptor.capture());
-        assertThat(((PipelineCDCSocketSink) 
sinkCaptor.getValue()).getChannel(), is(channel));
-        assertThat(connectionContext.getJobId(), is("foo_job"));
+        try (MockedStatic<PipelineJobRegistry> jobRegistryMocked = 
mockStatic(PipelineJobRegistry.class)) {
+            jobRegistryMocked.when(() -> 
PipelineJobRegistry.stop("foo_job")).thenThrow(new IllegalStateException("stop 
should be owned by CDCJobAPI"));
+            backendHandler.startStreaming("foo_job", connectionContext, 
channel);
+            ArgumentCaptor<PipelineSink> sinkCaptor = 
ArgumentCaptor.forClass(PipelineSink.class);
+            verify(jobAPI).start(eq("foo_job"), sinkCaptor.capture());
+            assertThat(((PipelineCDCSocketSink) 
sinkCaptor.getValue()).getChannel(), is(channel));
+            assertThat(connectionContext.getJobId(), is("foo_job"));
+            jobRegistryMocked.verifyNoInteractions();
+        }
     }
     
     @Test
@@ -203,9 +209,11 @@ class CDCBackendHandlerTest {
     
     @Test
     void assertStopStreamingWhenJobMissing() {
-        when(PipelineJobRegistry.get("foo_job")).thenReturn(null);
-        backendHandler.stopStreaming("foo_job", 
DefaultChannelId.newInstance());
-        verifyNoInteractions(jobAPI);
+        try (MockedStatic<PipelineJobRegistry> jobRegistryMocked = 
mockStatic(PipelineJobRegistry.class)) {
+            jobRegistryMocked.when(() -> 
PipelineJobRegistry.get("foo_job")).thenReturn(null);
+            backendHandler.stopStreaming("foo_job", 
DefaultChannelId.newInstance());
+            verifyNoInteractions(jobAPI);
+        }
     }
     
     @Test
@@ -214,9 +222,11 @@ class CDCBackendHandlerTest {
         Channel channel = mockChannel(DefaultChannelId.newInstance());
         CDCJob job = mock(CDCJob.class);
         when(job.getSink()).thenReturn(new PipelineCDCSocketSink(channel, 
mock(ShardingSphereDatabase.class), Collections.emptyList()));
-        when(PipelineJobRegistry.get("foo_job")).thenReturn(job);
-        backendHandler.stopStreaming("foo_job", targetChannelId);
-        verifyNoInteractions(jobAPI);
+        try (MockedStatic<PipelineJobRegistry> jobRegistryMocked = 
mockStatic(PipelineJobRegistry.class)) {
+            jobRegistryMocked.when(() -> 
PipelineJobRegistry.get("foo_job")).thenReturn(job);
+            backendHandler.stopStreaming("foo_job", targetChannelId);
+            verifyNoInteractions(jobAPI);
+        }
     }
     
     @Test
@@ -225,9 +235,12 @@ class CDCBackendHandlerTest {
         Channel channel = mockChannel(targetChannelId);
         CDCJob job = mock(CDCJob.class);
         when(job.getSink()).thenReturn(new PipelineCDCSocketSink(channel, 
mock(ShardingSphereDatabase.class), Collections.emptyList()));
-        when(PipelineJobRegistry.get("foo_job")).thenReturn(job);
-        backendHandler.stopStreaming("foo_job", targetChannelId);
-        verify(jobAPI).disable("foo_job");
+        try (MockedStatic<PipelineJobRegistry> jobRegistryMocked = 
mockStatic(PipelineJobRegistry.class)) {
+            jobRegistryMocked.when(() -> 
PipelineJobRegistry.get("foo_job")).thenReturn(job);
+            backendHandler.stopStreaming("foo_job", targetChannelId);
+            jobRegistryMocked.verify(() -> 
PipelineJobRegistry.stop("foo_job"));
+            verify(jobAPI).disable("foo_job");
+        }
     }
     
     @Test
diff --git 
a/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/util/CDCSchemaTableUtilsTest.java
 
b/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/util/CDCSchemaTableUtilsTest.java
index 05500aeed79..e617a41f6b1 100644
--- 
a/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/util/CDCSchemaTableUtilsTest.java
+++ 
b/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/util/CDCSchemaTableUtilsTest.java
@@ -21,6 +21,7 @@ import 
org.apache.shardingsphere.infra.config.props.ConfigurationProperties;
 
 import 
org.apache.shardingsphere.data.pipeline.cdc.protocol.request.StreamDataRequestBody.SchemaTable;
 import org.apache.shardingsphere.database.connector.core.type.DatabaseType;
+import 
org.apache.shardingsphere.infra.exception.kernel.metadata.SchemaNotFoundException;
 import 
org.apache.shardingsphere.infra.exception.kernel.metadata.TableNotFoundException;
 import 
org.apache.shardingsphere.infra.metadata.database.ShardingSphereDatabase;
 import 
org.apache.shardingsphere.infra.metadata.database.schema.model.ShardingSphereSchema;
@@ -92,6 +93,19 @@ class CDCSchemaTableUtilsTest {
         assertThat(actualResult, is(expectedResult));
     }
     
+    @Test
+    void assertParseTableExpressionsMergesSchemaWildcardResults() {
+        ShardingSphereSchema publicSchema = mockSchema("public", "t_order", 
"t_order_item");
+        ShardingSphereDatabase database =
+                new ShardingSphereDatabase("sharding_db", databaseType, null, 
null, Collections.singleton(publicSchema), new ConfigurationProperties(new 
Properties()));
+        List<SchemaTable> schemaTables = Arrays.asList(
+                
SchemaTable.newBuilder().setSchema("*").setTable("t_order").build(),
+                
SchemaTable.newBuilder().setSchema("*").setTable("t_order_item").build());
+        Map<String, Set<String>> actualResult = 
CDCSchemaTableUtils.parseTableExpressions(database, schemaTables);
+        Map<String, Set<String>> expectedResult = 
Collections.singletonMap("public", new HashSet<>(Arrays.asList("t_order", 
"t_order_item")));
+        assertThat(actualResult, is(expectedResult));
+    }
+    
     @Test
     void assertParseTableExpressionsWithAllTablesInSchema() {
         ShardingSphereSchema quotedSchema = mockSchema("CaseSchema", 
"t_order", "t_order_item");
@@ -156,6 +170,14 @@ class CDCSchemaTableUtilsTest {
         assertThrows(TableNotFoundException.class, () -> 
CDCSchemaTableUtils.parseTableExpressions(database, schemaTables));
     }
     
+    @Test
+    void assertParseTableExpressionsWithMissingSchema() {
+        ShardingSphereSchema publicSchema = mockSchema("public", "t_order");
+        ShardingSphereDatabase database = new 
ShardingSphereDatabase("sharding_db", databaseType, null, null, 
Collections.singleton(publicSchema), new ConfigurationProperties(new 
Properties()));
+        List<SchemaTable> schemaTables = 
Collections.singletonList(SchemaTable.newBuilder().setSchema("missing").setTable("t_order").build());
+        assertThrows(SchemaNotFoundException.class, () -> 
CDCSchemaTableUtils.parseTableExpressions(database, schemaTables));
+    }
+    
     @Test
     void assertParseTableExpressionsWithoutSchemaWithWildcard() {
         ShardingSphereSchema publicSchema = mockSchema("public", "t_order", 
"t_order2");
diff --git 
a/proxy/frontend/core/src/main/java/org/apache/shardingsphere/proxy/frontend/netty/CDCChannelInboundHandler.java
 
b/proxy/frontend/core/src/main/java/org/apache/shardingsphere/proxy/frontend/netty/CDCChannelInboundHandler.java
index 206e1a587c4..0684c49b1d3 100644
--- 
a/proxy/frontend/core/src/main/java/org/apache/shardingsphere/proxy/frontend/netty/CDCChannelInboundHandler.java
+++ 
b/proxy/frontend/core/src/main/java/org/apache/shardingsphere/proxy/frontend/netty/CDCChannelInboundHandler.java
@@ -48,8 +48,8 @@ import 
org.apache.shardingsphere.database.exception.core.exception.connection.Ac
 import 
org.apache.shardingsphere.database.exception.core.exception.syntax.database.UnknownDatabaseException;
 import org.apache.shardingsphere.infra.version.ShardingSphereVersion;
 import org.apache.shardingsphere.infra.exception.ShardingSpherePreconditions;
+import 
org.apache.shardingsphere.infra.exception.external.sql.ShardingSphereSQLException;
 import 
org.apache.shardingsphere.infra.exception.external.sql.sqlstate.XOpenSQLState;
-import 
org.apache.shardingsphere.infra.exception.external.sql.type.kernel.category.PipelineSQLException;
 import 
org.apache.shardingsphere.infra.exception.kernel.metadata.rule.MissingRequiredRuleException;
 import org.apache.shardingsphere.infra.metadata.user.Grantee;
 import org.apache.shardingsphere.infra.metadata.user.ShardingSphereUser;
@@ -181,7 +181,7 @@ public final class CDCChannelInboundHandler extends 
ChannelInboundHandlerAdapter
         try {
             CDCResponse response = 
backendHandler.streamData(request.getRequestId(), requestBody, 
connectionContext, ctx.channel());
             ctx.writeAndFlush(response);
-        } catch (final PipelineSQLException ex) {
+        } catch (final ShardingSphereSQLException ex) {
             throw new CDCExceptionWrapper(request.getRequestId(), ex);
         }
     }
@@ -205,10 +205,14 @@ public final class CDCChannelInboundHandler extends 
ChannelInboundHandlerAdapter
         if (requestBody.getStreamingId().isEmpty()) {
             throw new CDCExceptionWrapper(request.getRequestId(), new 
PipelineInvalidParameterException("Streaming id is empty"));
         }
-        String database = 
backendHandler.getDatabaseNameByJobId(requestBody.getStreamingId());
-        checkPrivileges(request.getRequestId(), 
connectionContext.getCurrentUser().getGrantee(), database);
-        backendHandler.startStreaming(requestBody.getStreamingId(), 
connectionContext, ctx.channel());
-        ctx.writeAndFlush(CDCResponseUtils.succeed(request.getRequestId()));
+        try {
+            String database = 
backendHandler.getDatabaseNameByJobId(requestBody.getStreamingId());
+            checkPrivileges(request.getRequestId(), 
connectionContext.getCurrentUser().getGrantee(), database);
+            backendHandler.startStreaming(requestBody.getStreamingId(), 
connectionContext, ctx.channel());
+            
ctx.writeAndFlush(CDCResponseUtils.succeed(request.getRequestId()));
+        } catch (final ShardingSphereSQLException ex) {
+            throw new CDCExceptionWrapper(request.getRequestId(), ex);
+        }
     }
     
     private void processStopStreamingRequest(final ChannelHandlerContext ctx, 
final CDCRequest request, final CDCConnectionContext connectionContext) {
diff --git 
a/proxy/frontend/core/src/test/java/org/apache/shardingsphere/proxy/frontend/netty/CDCChannelInboundHandlerTest.java
 
b/proxy/frontend/core/src/test/java/org/apache/shardingsphere/proxy/frontend/netty/CDCChannelInboundHandlerTest.java
index 72d03d42da1..49825306517 100644
--- 
a/proxy/frontend/core/src/test/java/org/apache/shardingsphere/proxy/frontend/netty/CDCChannelInboundHandlerTest.java
+++ 
b/proxy/frontend/core/src/test/java/org/apache/shardingsphere/proxy/frontend/netty/CDCChannelInboundHandlerTest.java
@@ -53,7 +53,7 @@ import 
org.apache.shardingsphere.database.exception.core.exception.syntax.databa
 import org.apache.shardingsphere.infra.config.props.ConfigurationProperties;
 import org.apache.shardingsphere.infra.config.props.ConfigurationPropertyKey;
 import 
org.apache.shardingsphere.infra.exception.external.sql.sqlstate.XOpenSQLState;
-import 
org.apache.shardingsphere.infra.exception.external.sql.type.kernel.category.PipelineSQLException;
+import 
org.apache.shardingsphere.infra.exception.kernel.metadata.SchemaNotFoundException;
 import 
org.apache.shardingsphere.infra.exception.kernel.metadata.rule.MissingRequiredRuleException;
 import org.apache.shardingsphere.infra.metadata.database.rule.RuleMetaData;
 import org.apache.shardingsphere.infra.metadata.user.Grantee;
@@ -83,13 +83,13 @@ import static org.hamcrest.MatcherAssert.assertThat;
 import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertNull;
-import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyString;
 import static org.mockito.ArgumentMatchers.argThat;
 import static org.mockito.Mockito.RETURNS_DEEP_STUBS;
 import static org.mockito.Mockito.atLeastOnce;
+import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verify;
@@ -345,16 +345,27 @@ class CDCChannelInboundHandlerTest {
     }
     
     @Test
-    void assertStreamDataRequestWrapsPipelineSQLException() {
-        CDCConnectionContext connectionContext = new 
CDCConnectionContext(user);
-        channel.attr(CONNECTION_CONTEXT_KEY).set(connectionContext);
-        StreamDataRequestBody.Builder bodyBuilder = 
StreamDataRequestBody.newBuilder().setDatabase("logic_db");
-        
bodyBuilder.addSourceSchemaTable(StreamDataRequestBody.SchemaTable.newBuilder().setSchema("schema").setTable("table").build());
-        CDCRequest request = 
CDCRequest.newBuilder().setType(Type.STREAM_DATA).setRequestId("stream-request").setStreamDataRequestBody(bodyBuilder.build()).build();
-        when(backendHandler.streamData(any(), any(), any(), 
any())).thenThrow(mock(PipelineSQLException.class));
-        ChannelHandlerContext context = channel.pipeline().context(handler);
-        assertThrows(CDCExceptionWrapper.class, () -> 
handler.channelRead(context, request));
-        assertThat(channel.attr(CONNECTION_CONTEXT_KEY).get(), 
is(connectionContext));
+    void assertStreamDataRequestWithSchemaNotFoundException() {
+        channel.attr(CONNECTION_CONTEXT_KEY).set(new 
CDCConnectionContext(user));
+        CDCRequest request = createStreamDataRequest("logic_db");
+        when(backendHandler.streamData(any(), any(), any(), 
any())).thenThrow(new SchemaNotFoundException("foo_schema"));
+        channel.writeInbound(request);
+        CDCResponse response = readResponseSkippingGreeting();
+        assertThat(response.getStatus(), is(Status.FAILED));
+        assertThat(response.getRequestId(), is("stream-request"));
+        assertThat(response.getErrorCode(), 
is(XOpenSQLState.NOT_FOUND.getValue()));
+    }
+    
+    @Test
+    void assertStreamDataRequestWithInvalidParameterException() {
+        channel.attr(CONNECTION_CONTEXT_KEY).set(new 
CDCConnectionContext(user));
+        CDCRequest request = createStreamDataRequest("logic_db");
+        when(backendHandler.streamData(any(), any(), any(), 
any())).thenThrow(new PipelineInvalidParameterException("invalid job 
configuration"));
+        channel.writeInbound(request);
+        CDCResponse response = readResponseSkippingGreeting();
+        assertThat(response.getStatus(), is(Status.FAILED));
+        assertThat(response.getRequestId(), is("stream-request"));
+        assertThat(response.getErrorCode(), 
is(XOpenSQLState.INVALID_PARAMETER_VALUE.getValue()));
     }
     
     @Test
@@ -429,6 +440,20 @@ class CDCChannelInboundHandlerTest {
         assertThat(response.getStatus(), is(Status.SUCCEED));
     }
     
+    @Test
+    void assertStartStreamingRequestWithInvalidParameterException() {
+        channel.attr(CONNECTION_CONTEXT_KEY).set(new 
CDCConnectionContext(user));
+        
when(backendHandler.getDatabaseNameByJobId("job-1")).thenReturn("logic_db");
+        doThrow(new PipelineInvalidParameterException("invalid job 
configuration")).when(backendHandler).startStreaming(any(), any(), any());
+        StartStreamingRequestBody body = 
StartStreamingRequestBody.newBuilder().setStreamingId("job-1").build();
+        CDCRequest request = 
CDCRequest.newBuilder().setType(Type.START_STREAMING).setRequestId("start-request").setStartStreamingRequestBody(body).build();
+        channel.writeInbound(request);
+        CDCResponse response = readResponseSkippingGreeting();
+        assertThat(response.getStatus(), is(Status.FAILED));
+        assertThat(response.getRequestId(), is("start-request"));
+        assertThat(response.getErrorCode(), 
is(XOpenSQLState.INVALID_PARAMETER_VALUE.getValue()));
+    }
+    
     @Test
     void assertStopStreamingRequestSucceed() {
         CDCConnectionContext connectionContext = new 
CDCConnectionContext(user);

Reply via email to