From 622ced5b556616c9cb4d2b450e401f7394f8e952 Mon Sep 17 00:00:00 2001 From: sunlishuo Date: Sat, 3 Oct 2026 10:33:54 +0800 Subject: [PATCH 1/2] [4444][fix] scope job monitor tenant context --- .../main/java/org/dinky/job/FlinkJobTask.java | 24 +- .../java/org/dinky/job/FlinkJobTaskTest.java | 356 ++++++++++++++++++ 2 files changed, 374 insertions(+), 6 deletions(-) create mode 100644 dinky-admin/src/test/java/org/dinky/job/FlinkJobTaskTest.java diff --git a/dinky-admin/src/main/java/org/dinky/job/FlinkJobTask.java b/dinky-admin/src/main/java/org/dinky/job/FlinkJobTask.java index 979d7d6ca7..92957b7a59 100644 --- a/dinky-admin/src/main/java/org/dinky/job/FlinkJobTask.java +++ b/dinky-admin/src/main/java/org/dinky/job/FlinkJobTask.java @@ -21,6 +21,7 @@ import org.dinky.assertion.Asserts; import org.dinky.context.SpringContextUtils; +import org.dinky.context.TenantContextHolder; import org.dinky.daemon.constant.FlinkTaskConstant; import org.dinky.daemon.task.DaemonTask; import org.dinky.daemon.task.DaemonTaskConfig; @@ -101,14 +102,25 @@ public DaemonTaskConfig getConfig() { public boolean dealTask() { volatilityBalance(); - boolean isDone = JobRefreshHandler.refreshJob(jobInfoDetail, isNeedSave()); - if (Asserts.isAllNotNull(jobInfoDetail.getClusterInstance())) { - JobAlertHandler.getInstance().check(jobInfoDetail); - if (SystemConfiguration.getInstances().getMetricsSysEnable().getValue()) { - JobMetricsHandler.refreshAndWriteFlinkMetrics(jobInfoDetail, verticesAndMetricsMap); + Object previousTenant = TenantContextHolder.get(); + try { + // A reused monitor worker may have inherited another job's tenant. + TenantContextHolder.set(jobInfoDetail.getInstance().getTenantId()); + boolean isDone = JobRefreshHandler.refreshJob(jobInfoDetail, isNeedSave()); + if (Asserts.isAllNotNull(jobInfoDetail.getClusterInstance())) { + JobAlertHandler.getInstance().check(jobInfoDetail); + if (SystemConfiguration.getInstances().getMetricsSysEnable().getValue()) { + JobMetricsHandler.refreshAndWriteFlinkMetrics(jobInfoDetail, verticesAndMetricsMap); + } + } + return isDone; + } finally { + if (previousTenant == null) { + TenantContextHolder.clear(); + } else { + TenantContextHolder.set(previousTenant); } } - return isDone; } /** diff --git a/dinky-admin/src/test/java/org/dinky/job/FlinkJobTaskTest.java b/dinky-admin/src/test/java/org/dinky/job/FlinkJobTaskTest.java new file mode 100644 index 0000000000..8e041d4dde --- /dev/null +++ b/dinky-admin/src/test/java/org/dinky/job/FlinkJobTaskTest.java @@ -0,0 +1,356 @@ +/* + * + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.dinky.job; + +import static org.junit.jupiter.api.Assertions.assertAll; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertSame; +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.anyInt; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.reset; +import static org.mockito.Mockito.when; + +import org.dinky.configure.MybatisPlusConfig; +import org.dinky.context.SpringContextUtils; +import org.dinky.context.TenantContextHolder; +import org.dinky.data.dto.JobDataDto; +import org.dinky.data.enums.GatewayType; +import org.dinky.data.enums.JobStatus; +import org.dinky.data.model.ClusterInstance; +import org.dinky.data.model.SystemConfiguration; +import org.dinky.data.model.ext.JobInfoDetail; +import org.dinky.data.model.job.JobInstance; +import org.dinky.mapper.AlertRulesMapper; +import org.dinky.mybatis.properties.MybatisPlusFillProperties; +import org.dinky.service.AlertHistoryService; +import org.dinky.service.ClusterInstanceService; +import org.dinky.service.HistoryService; +import org.dinky.service.JobHistoryService; +import org.dinky.service.JobInstanceService; +import org.dinky.service.MonitorService; +import org.dinky.service.TaskService; +import org.dinky.service.UserService; +import org.dinky.service.impl.AlertRuleServiceImpl; + +import org.apache.ibatis.annotations.Update; +import org.apache.ibatis.datasource.unpooled.UnpooledDataSource; +import org.apache.ibatis.mapping.Environment; +import org.apache.ibatis.session.Configuration; +import org.apache.ibatis.session.SqlSession; +import org.apache.ibatis.session.SqlSessionFactoryBuilder; +import org.apache.ibatis.transaction.jdbc.JdbcTransactionFactory; + +import java.io.IOException; +import java.io.OutputStream; +import java.net.InetSocketAddress; +import java.nio.charset.StandardCharsets; +import java.sql.Connection; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.sql.Statement; +import java.util.Collections; +import java.util.HashMap; +import java.util.Map; +import java.util.UUID; + +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.context.ApplicationContext; +import org.springframework.context.support.StaticApplicationContext; + +import com.sun.net.httpserver.HttpServer; + +class FlinkJobTaskTest { + + private static final JobInstanceService JOB_INSTANCE_SERVICE = mock(JobInstanceService.class); + + private static StaticApplicationContext applicationContext; + private static ApplicationContext previousApplicationContext; + + private final Map jobTenants = new HashMap<>(); + private final Map persistedStatuses = new HashMap<>(); + + @BeforeAll + static void registerServices() { + previousApplicationContext = SpringContextUtils.applicationContext; + applicationContext = new StaticApplicationContext(); + applicationContext.getBeanFactory().registerSingleton("jobInstanceServiceImpl", JOB_INSTANCE_SERVICE); + applicationContext.getBeanFactory().registerSingleton("monitorServiceImpl", mock(MonitorService.class)); + applicationContext.getBeanFactory().registerSingleton("jobHistoryServiceImpl", mock(JobHistoryService.class)); + applicationContext + .getBeanFactory() + .registerSingleton("clusterInstanceServiceImpl", mock(ClusterInstanceService.class)); + applicationContext.getBeanFactory().registerSingleton("historyServiceImpl", mock(HistoryService.class)); + applicationContext.getBeanFactory().registerSingleton("taskServiceImpl", mock(TaskService.class)); + applicationContext + .getBeanFactory() + .registerSingleton("alertHistoryServiceImpl", mock(AlertHistoryService.class)); + applicationContext.getBeanFactory().registerSingleton("userServiceImpl", mock(UserService.class)); + AlertRuleServiceImpl alertRuleService = mock(AlertRuleServiceImpl.class); + AlertRulesMapper alertRulesMapper = mock(AlertRulesMapper.class); + when(alertRuleService.getBaseMapper()).thenReturn(alertRulesMapper); + when(alertRulesMapper.selectWithTemplate()).thenReturn(Collections.emptyList()); + applicationContext.getBeanFactory().registerSingleton("alertRuleServiceImpl", alertRuleService); + SpringContextUtils.applicationContext = applicationContext; + } + + @AfterAll + static void restoreApplicationContext() { + SpringContextUtils.applicationContext = previousApplicationContext; + applicationContext.close(); + } + + @BeforeEach + void prepareTenantScopedPersistence() { + TenantContextHolder.clear(); + reset(JOB_INSTANCE_SERVICE); + doAnswer(invocation -> { + TenantContextHolder.set(jobTenants.get(invocation.getArgument(0))); + return null; + }) + .when(JOB_INSTANCE_SERVICE) + .initTenantByJobInstanceId(anyInt()); + when(JOB_INSTANCE_SERVICE.updateById(any(JobInstance.class))).thenAnswer(invocation -> { + JobInstance instance = invocation.getArgument(0); + // A tenant-filtered update matches no row when a worker keeps another job's tenant. + if (!jobTenants.get(instance.getId()).equals(TenantContextHolder.get())) { + return false; + } + persistedStatuses.put(instance.getId(), instance.getStatus()); + return true; + }); + } + + @AfterEach + void clearTenantContext() { + TenantContextHolder.clear(); + } + + @Test + void persistsTerminalStatusUsingTheJobTenant() { + TenantContextHolder.set(1); + FlinkJobTask task = task(101, 2); + + assertTrue(task.dealTask()); + + assertAll( + () -> assertEquals(JobStatus.UNKNOWN.getValue(), persistedStatuses.get(101)), + () -> assertEquals(1, TenantContextHolder.get())); + } + + @Test + void clearsTenantContextWhenTheWorkerHadNoTenant() { + FlinkJobTask task = task(101, 2); + + assertTrue(task.dealTask()); + + assertAll( + () -> assertEquals(JobStatus.UNKNOWN.getValue(), persistedStatuses.get(101)), + () -> assertNull(TenantContextHolder.get())); + } + + @Test + void refreshesDifferentTenantsOnTheSameWorker() { + TenantContextHolder.set(1); + FlinkJobTask firstTask = task(101, 2); + FlinkJobTask secondTask = task(102, 3); + + assertTrue(firstTask.dealTask()); + assertTrue(secondTask.dealTask()); + + assertAll( + () -> assertEquals(JobStatus.UNKNOWN.getValue(), persistedStatuses.get(101)), + () -> assertEquals(JobStatus.UNKNOWN.getValue(), persistedStatuses.get(102)), + () -> assertEquals(1, TenantContextHolder.get())); + } + + @Test + void restoresPreviousTenantWhenRefreshFails() { + TenantContextHolder.set(1); + FlinkJobTask task = task(101, 2); + IllegalStateException failure = new IllegalStateException("Unable to persist status"); + doAnswer(invocation -> { + assertEquals(2, TenantContextHolder.get()); + throw failure; + }) + .when(JOB_INSTANCE_SERVICE) + .updateById(any(JobInstance.class)); + + assertSame(failure, assertThrows(IllegalStateException.class, task::dealTask)); + assertEquals(1, TenantContextHolder.get()); + } + + @Test + void clearsTenantContextWhenRefreshFailsWithoutPreviousTenant() { + FlinkJobTask task = task(101, 2); + IllegalStateException failure = new IllegalStateException("Unable to persist status"); + doThrow(failure).when(JOB_INSTANCE_SERVICE).updateById(any(JobInstance.class)); + + assertSame(failure, assertThrows(IllegalStateException.class, task::dealTask)); + assertNull(TenantContextHolder.get()); + } + + @Test + void preservesTheTenantForManualRefresh() { + TenantContextHolder.set(2); + FlinkJobTask task = task(101, 2); + + assertTrue(task.dealTask()); + + assertAll( + () -> assertEquals(JobStatus.UNKNOWN.getValue(), persistedStatuses.get(101)), + () -> assertEquals(2, TenantContextHolder.get())); + } + + @Test + void persistsFailedYarnSessionJobWithTenantFilteringEnabled() throws Exception { + HttpServer flink = failedJobServer(); + Boolean metricsEnabled = + SystemConfiguration.getInstances().getMetricsSysEnable().getValue(); + Integer resendInterval = + SystemConfiguration.getInstances().getJobReSendDiffSecond().getValue(); + SystemConfiguration.getInstances().getMetricsSysEnable().setValue(false); + SystemConfiguration.getInstances().getJobReSendDiffSecond().setValue(60); + UnpooledDataSource dataSource = + new UnpooledDataSource("org.h2.Driver", "jdbc:h2:mem:" + UUID.randomUUID(), "sa", ""); + + try (Connection connection = dataSource.getConnection(); + Statement statement = connection.createStatement()) { + statement.execute( + "CREATE TABLE dinky_job_instance (id INT PRIMARY KEY, tenant_id INT, status VARCHAR(32))"); + statement.execute("INSERT INTO dinky_job_instance VALUES (101, 2, 'RUNNING'), (102, 1, 'RUNNING')"); + Configuration configuration = + new Configuration(new Environment("test", new JdbcTransactionFactory(), dataSource)); + configuration.addInterceptor( + new MybatisPlusConfig(new MybatisPlusFillProperties()).mybatisPlusInterceptor()); + configuration.addMapper(JobStatusMapper.class); + + try (SqlSession session = + new SqlSessionFactoryBuilder().build(configuration).openSession(true)) { + JobStatusMapper mapper = session.getMapper(JobStatusMapper.class); + TenantContextHolder.set(1); + assertFalse(TenantContextHolder.isIgnoreTenant()); + FlinkJobTask task = task(101, 2); + JobInfoDetail detail = task.getJobInfoDetail(); + JobInstance instance = detail.getInstance(); + instance.setTaskId(10); + instance.setName("failed-job"); + instance.setJid("job-101"); + instance.setStatus(JobStatus.FAILED.getValue()); + assertEquals(0, mapper.updateStatus(instance), "Another tenant must not be able to update this row"); + instance.setStatus(JobStatus.RUNNING.getValue()); + detail.setJobDataDto(JobDataDto.builder().id(101).tenantId(2).build()); + ClusterInstance cluster = new ClusterInstance(); + cluster.setName("yarn-session"); + cluster.setType(GatewayType.YARN_SESSION.getLongValue()); + cluster.setJobManagerHost("127.0.0.1:" + flink.getAddress().getPort()); + cluster.setHosts(cluster.getJobManagerHost()); + detail.setClusterInstance(cluster); + doAnswer(invocation -> mapper.updateStatus(invocation.getArgument(0)) == 1) + .when(JOB_INSTANCE_SERVICE) + .updateById(any(JobInstance.class)); + + assertTrue(task.dealTask()); + + assertAll( + () -> assertEquals(JobStatus.FAILED.getValue(), persistedStatus(connection, 101)), + () -> assertEquals(JobStatus.RUNNING.getValue(), persistedStatus(connection, 102)), + () -> assertEquals(1, TenantContextHolder.get()), + () -> assertFalse(TenantContextHolder.isIgnoreTenant())); + } + } finally { + flink.stop(0); + SystemConfiguration.getInstances().getMetricsSysEnable().setValue(metricsEnabled); + SystemConfiguration.getInstances().getJobReSendDiffSecond().setValue(resendInterval); + } + } + + private static HttpServer failedJobServer() throws IOException { + Map responses = new HashMap<>(); + responses.put( + "/jobs/job-101", + "{\"jid\":\"job-101\",\"name\":\"failed-job\",\"state\":\"FAILED\"," + + "\"start-time\":1000,\"end-time\":2000,\"duration\":1000,\"vertices\":[],\"plan\":{\"nodes\":[]}}"); + responses.put("/jobs/job-101/config", "{\"jid\":\"job-101\",\"name\":\"failed-job\",\"execution-config\":{}}"); + responses.put("/jobs/job-101/checkpoints", "{\"errors\":[]}"); + responses.put("/jobs/job-101/checkpoints/config", "{\"errors\":[]}"); + responses.put( + "/jobs/job-101/exceptions", + "{\"all-exceptions\":[],\"root-exception\":\"\",\"timestamp\":2000,\"truncated\":false," + + "\"exceptionHistory\":{\"entries\":[],\"truncated\":false}}"); + HttpServer server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); + server.createContext("/jobs/job-101", exchange -> { + String response = responses.get(exchange.getRequestURI().getPath()); + boolean validRequest = "GET".equals(exchange.getRequestMethod()) && response != null; + byte[] body = (validRequest ? response : "{\"errors\":[\"Unexpected request\"]}") + .getBytes(StandardCharsets.UTF_8); + exchange.getResponseHeaders().set("Content-Type", "application/json"); + exchange.sendResponseHeaders(validRequest ? 200 : 404, body.length); + try (OutputStream output = exchange.getResponseBody()) { + output.write(body); + } + }); + server.start(); + return server; + } + + private static String persistedStatus(Connection connection, int id) throws SQLException { + try (PreparedStatement statement = + connection.prepareStatement("SELECT status FROM dinky_job_instance WHERE id = ?")) { + statement.setInt(1, id); + try (ResultSet rows = statement.executeQuery()) { + assertTrue(rows.next()); + return rows.getString("status"); + } + } + } + + interface JobStatusMapper { + @Update("UPDATE dinky_job_instance SET status = #{status} WHERE id = #{id}") + int updateStatus(JobInstance instance); + } + + private FlinkJobTask task(int id, int tenantId) { + JobInstance instance = new JobInstance(); + instance.setId(id); + instance.setTenantId(tenantId); + instance.setStatus(JobStatus.RUNNING.getValue()); + jobTenants.put(id, tenantId); + persistedStatuses.put(id, JobStatus.RUNNING.getValue()); + + // Missing cluster metadata is a terminal refresh path that needs no Flink or YARN server. + JobInfoDetail detail = new JobInfoDetail(id); + detail.setInstance(instance); + FlinkJobTask task = new FlinkJobTask(); + task.setJobInfoDetail(detail); + return task; + } +} From b1290f24f19e158c5ee737b464fd5370e991d2e6 Mon Sep 17 00:00:00 2001 From: sunlishuo Date: Sat, 3 Oct 2026 11:06:54 +0800 Subject: [PATCH 2/2] [4444][test] isolate job monitor regression state --- .../java/org/dinky/job/FlinkJobTaskTest.java | 569 +++++++++++------- 1 file changed, 336 insertions(+), 233 deletions(-) diff --git a/dinky-admin/src/test/java/org/dinky/job/FlinkJobTaskTest.java b/dinky-admin/src/test/java/org/dinky/job/FlinkJobTaskTest.java index 8e041d4dde..4a3ae9ac47 100644 --- a/dinky-admin/src/test/java/org/dinky/job/FlinkJobTaskTest.java +++ b/dinky-admin/src/test/java/org/dinky/job/FlinkJobTaskTest.java @@ -65,7 +65,11 @@ import org.apache.ibatis.transaction.jdbc.JdbcTransactionFactory; import java.io.IOException; +import java.io.InputStream; import java.io.OutputStream; +import java.lang.reflect.Constructor; +import java.lang.reflect.InvocationTargetException; +import java.lang.reflect.Method; import java.net.InetSocketAddress; import java.nio.charset.StandardCharsets; import java.sql.Connection; @@ -78,279 +82,378 @@ import java.util.Map; import java.util.UUID; -import org.junit.jupiter.api.AfterAll; -import org.junit.jupiter.api.AfterEach; -import org.junit.jupiter.api.BeforeAll; -import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.springframework.context.ApplicationContext; import org.springframework.context.support.StaticApplicationContext; +import com.google.common.io.ByteStreams; import com.sun.net.httpserver.HttpServer; class FlinkJobTaskTest { - private static final JobInstanceService JOB_INSTANCE_SERVICE = mock(JobInstanceService.class); - - private static StaticApplicationContext applicationContext; - private static ApplicationContext previousApplicationContext; - - private final Map jobTenants = new HashMap<>(); - private final Map persistedStatuses = new HashMap<>(); - - @BeforeAll - static void registerServices() { - previousApplicationContext = SpringContextUtils.applicationContext; - applicationContext = new StaticApplicationContext(); - applicationContext.getBeanFactory().registerSingleton("jobInstanceServiceImpl", JOB_INSTANCE_SERVICE); - applicationContext.getBeanFactory().registerSingleton("monitorServiceImpl", mock(MonitorService.class)); - applicationContext.getBeanFactory().registerSingleton("jobHistoryServiceImpl", mock(JobHistoryService.class)); - applicationContext - .getBeanFactory() - .registerSingleton("clusterInstanceServiceImpl", mock(ClusterInstanceService.class)); - applicationContext.getBeanFactory().registerSingleton("historyServiceImpl", mock(HistoryService.class)); - applicationContext.getBeanFactory().registerSingleton("taskServiceImpl", mock(TaskService.class)); - applicationContext - .getBeanFactory() - .registerSingleton("alertHistoryServiceImpl", mock(AlertHistoryService.class)); - applicationContext.getBeanFactory().registerSingleton("userServiceImpl", mock(UserService.class)); - AlertRuleServiceImpl alertRuleService = mock(AlertRuleServiceImpl.class); - AlertRulesMapper alertRulesMapper = mock(AlertRulesMapper.class); - when(alertRuleService.getBaseMapper()).thenReturn(alertRulesMapper); - when(alertRulesMapper.selectWithTemplate()).thenReturn(Collections.emptyList()); - applicationContext.getBeanFactory().registerSingleton("alertRuleServiceImpl", alertRuleService); - SpringContextUtils.applicationContext = applicationContext; - } - - @AfterAll - static void restoreApplicationContext() { - SpringContextUtils.applicationContext = previousApplicationContext; - applicationContext.close(); - } - - @BeforeEach - void prepareTenantScopedPersistence() { - TenantContextHolder.clear(); - reset(JOB_INSTANCE_SERVICE); - doAnswer(invocation -> { - TenantContextHolder.set(jobTenants.get(invocation.getArgument(0))); - return null; - }) - .when(JOB_INSTANCE_SERVICE) - .initTenantByJobInstanceId(anyInt()); - when(JOB_INSTANCE_SERVICE.updateById(any(JobInstance.class))).thenAnswer(invocation -> { - JobInstance instance = invocation.getArgument(0); - // A tenant-filtered update matches no row when a worker keeps another job's tenant. - if (!jobTenants.get(instance.getId()).equals(TenantContextHolder.get())) { - return false; - } - persistedStatuses.put(instance.getId(), instance.getStatus()); - return true; - }); + @Test + void persistsTerminalStatusUsingTheJobTenant() throws Throwable { + runIsolated("persistsTerminalStatusUsingTheJobTenant"); } - @AfterEach - void clearTenantContext() { - TenantContextHolder.clear(); + @Test + void clearsTenantContextWhenTheWorkerHadNoTenant() throws Throwable { + runIsolated("clearsTenantContextWhenTheWorkerHadNoTenant"); } @Test - void persistsTerminalStatusUsingTheJobTenant() { - TenantContextHolder.set(1); - FlinkJobTask task = task(101, 2); - - assertTrue(task.dealTask()); - - assertAll( - () -> assertEquals(JobStatus.UNKNOWN.getValue(), persistedStatuses.get(101)), - () -> assertEquals(1, TenantContextHolder.get())); + void refreshesDifferentTenantsOnTheSameWorker() throws Throwable { + runIsolated("refreshesDifferentTenantsOnTheSameWorker"); } @Test - void clearsTenantContextWhenTheWorkerHadNoTenant() { - FlinkJobTask task = task(101, 2); - - assertTrue(task.dealTask()); - - assertAll( - () -> assertEquals(JobStatus.UNKNOWN.getValue(), persistedStatuses.get(101)), - () -> assertNull(TenantContextHolder.get())); + void restoresPreviousTenantWhenRefreshFails() throws Throwable { + runIsolated("restoresPreviousTenantWhenRefreshFails"); } @Test - void refreshesDifferentTenantsOnTheSameWorker() { - TenantContextHolder.set(1); - FlinkJobTask firstTask = task(101, 2); - FlinkJobTask secondTask = task(102, 3); - - assertTrue(firstTask.dealTask()); - assertTrue(secondTask.dealTask()); - - assertAll( - () -> assertEquals(JobStatus.UNKNOWN.getValue(), persistedStatuses.get(101)), - () -> assertEquals(JobStatus.UNKNOWN.getValue(), persistedStatuses.get(102)), - () -> assertEquals(1, TenantContextHolder.get())); + void clearsTenantContextWhenRefreshFailsWithoutPreviousTenant() throws Throwable { + runIsolated("clearsTenantContextWhenRefreshFailsWithoutPreviousTenant"); } @Test - void restoresPreviousTenantWhenRefreshFails() { - TenantContextHolder.set(1); - FlinkJobTask task = task(101, 2); - IllegalStateException failure = new IllegalStateException("Unable to persist status"); - doAnswer(invocation -> { - assertEquals(2, TenantContextHolder.get()); - throw failure; - }) - .when(JOB_INSTANCE_SERVICE) - .updateById(any(JobInstance.class)); - - assertSame(failure, assertThrows(IllegalStateException.class, task::dealTask)); - assertEquals(1, TenantContextHolder.get()); + void preservesTheTenantForManualRefresh() throws Throwable { + runIsolated("preservesTheTenantForManualRefresh"); } @Test - void clearsTenantContextWhenRefreshFailsWithoutPreviousTenant() { - FlinkJobTask task = task(101, 2); - IllegalStateException failure = new IllegalStateException("Unable to persist status"); - doThrow(failure).when(JOB_INSTANCE_SERVICE).updateById(any(JobInstance.class)); + void persistsFailedYarnSessionJobWithTenantFilteringEnabled() throws Throwable { + runIsolated("persistsFailedYarnSessionJobWithTenantFilteringEnabled"); + } - assertSame(failure, assertThrows(IllegalStateException.class, task::dealTask)); - assertNull(TenantContextHolder.get()); + private void runIsolated(String testMethod) throws Throwable { + ClassLoader previousLoader = Thread.currentThread().getContextClassLoader(); + ClassLoader isolatedLoader = new FixtureClassLoader(getClass().getClassLoader()); + try { + Thread.currentThread().setContextClassLoader(isolatedLoader); + Class fixtureClass = isolatedLoader.loadClass(getClass().getName() + "$Fixture"); + Constructor constructor = fixtureClass.getDeclaredConstructor(); + constructor.setAccessible(true); + Object fixture = constructor.newInstance(); + invoke(fixture, "registerServices"); + try { + invoke(fixture, "prepareTenantScopedPersistence"); + invoke(fixture, testMethod); + } finally { + try { + invoke(fixture, "clearTenantContext"); + } finally { + invoke(fixture, "restoreApplicationContext"); + } + } + } finally { + Thread.currentThread().setContextClassLoader(previousLoader); + } } - @Test - void preservesTheTenantForManualRefresh() { - TenantContextHolder.set(2); - FlinkJobTask task = task(101, 2); + private static void invoke(Object fixture, String methodName) throws Throwable { + Method method = fixture.getClass().getDeclaredMethod(methodName); + method.setAccessible(true); + try { + method.invoke(fixture); + } catch (InvocationTargetException e) { + throw e.getCause(); + } + } - assertTrue(task.dealTask()); + // The handlers retain static service references and configuration listeners. Load the complete + // Dinky fixture in its own namespace so those references cannot affect other tests or IDE reruns. + private static class FixtureClassLoader extends ClassLoader { - assertAll( - () -> assertEquals(JobStatus.UNKNOWN.getValue(), persistedStatuses.get(101)), - () -> assertEquals(2, TenantContextHolder.get())); - } + FixtureClassLoader(ClassLoader parent) { + super(parent); + } - @Test - void persistsFailedYarnSessionJobWithTenantFilteringEnabled() throws Exception { - HttpServer flink = failedJobServer(); - Boolean metricsEnabled = - SystemConfiguration.getInstances().getMetricsSysEnable().getValue(); - Integer resendInterval = - SystemConfiguration.getInstances().getJobReSendDiffSecond().getValue(); - SystemConfiguration.getInstances().getMetricsSysEnable().setValue(false); - SystemConfiguration.getInstances().getJobReSendDiffSecond().setValue(60); - UnpooledDataSource dataSource = - new UnpooledDataSource("org.h2.Driver", "jdbc:h2:mem:" + UUID.randomUUID(), "sa", ""); - - try (Connection connection = dataSource.getConnection(); - Statement statement = connection.createStatement()) { - statement.execute( - "CREATE TABLE dinky_job_instance (id INT PRIMARY KEY, tenant_id INT, status VARCHAR(32))"); - statement.execute("INSERT INTO dinky_job_instance VALUES (101, 2, 'RUNNING'), (102, 1, 'RUNNING')"); - Configuration configuration = - new Configuration(new Environment("test", new JdbcTransactionFactory(), dataSource)); - configuration.addInterceptor( - new MybatisPlusConfig(new MybatisPlusFillProperties()).mybatisPlusInterceptor()); - configuration.addMapper(JobStatusMapper.class); - - try (SqlSession session = - new SqlSessionFactoryBuilder().build(configuration).openSession(true)) { - JobStatusMapper mapper = session.getMapper(JobStatusMapper.class); - TenantContextHolder.set(1); - assertFalse(TenantContextHolder.isIgnoreTenant()); - FlinkJobTask task = task(101, 2); - JobInfoDetail detail = task.getJobInfoDetail(); - JobInstance instance = detail.getInstance(); - instance.setTaskId(10); - instance.setName("failed-job"); - instance.setJid("job-101"); - instance.setStatus(JobStatus.FAILED.getValue()); - assertEquals(0, mapper.updateStatus(instance), "Another tenant must not be able to update this row"); - instance.setStatus(JobStatus.RUNNING.getValue()); - detail.setJobDataDto(JobDataDto.builder().id(101).tenantId(2).build()); - ClusterInstance cluster = new ClusterInstance(); - cluster.setName("yarn-session"); - cluster.setType(GatewayType.YARN_SESSION.getLongValue()); - cluster.setJobManagerHost("127.0.0.1:" + flink.getAddress().getPort()); - cluster.setHosts(cluster.getJobManagerHost()); - detail.setClusterInstance(cluster); - doAnswer(invocation -> mapper.updateStatus(invocation.getArgument(0)) == 1) - .when(JOB_INSTANCE_SERVICE) - .updateById(any(JobInstance.class)); - - assertTrue(task.dealTask()); - - assertAll( - () -> assertEquals(JobStatus.FAILED.getValue(), persistedStatus(connection, 101)), - () -> assertEquals(JobStatus.RUNNING.getValue(), persistedStatus(connection, 102)), - () -> assertEquals(1, TenantContextHolder.get()), - () -> assertFalse(TenantContextHolder.isIgnoreTenant())); + @Override + protected Class loadClass(String name, boolean resolve) throws ClassNotFoundException { + if (!name.startsWith("org.dinky.")) { + return super.loadClass(name, resolve); + } + synchronized (getClassLoadingLock(name)) { + Class loaded = findLoadedClass(name); + if (loaded == null) { + String resource = name.replace('.', '/') + ".class"; + try (InputStream input = getParent().getResourceAsStream(resource)) { + if (input == null) { + throw new ClassNotFoundException(name); + } + byte[] bytes = ByteStreams.toByteArray(input); + loaded = defineClass(name, bytes, 0, bytes.length); + } catch (IOException e) { + throw new ClassNotFoundException(name, e); + } + } + if (resolve) { + resolveClass(loaded); + } + return loaded; } - } finally { - flink.stop(0); - SystemConfiguration.getInstances().getMetricsSysEnable().setValue(metricsEnabled); - SystemConfiguration.getInstances().getJobReSendDiffSecond().setValue(resendInterval); } } - private static HttpServer failedJobServer() throws IOException { - Map responses = new HashMap<>(); - responses.put( - "/jobs/job-101", - "{\"jid\":\"job-101\",\"name\":\"failed-job\",\"state\":\"FAILED\"," - + "\"start-time\":1000,\"end-time\":2000,\"duration\":1000,\"vertices\":[],\"plan\":{\"nodes\":[]}}"); - responses.put("/jobs/job-101/config", "{\"jid\":\"job-101\",\"name\":\"failed-job\",\"execution-config\":{}}"); - responses.put("/jobs/job-101/checkpoints", "{\"errors\":[]}"); - responses.put("/jobs/job-101/checkpoints/config", "{\"errors\":[]}"); - responses.put( - "/jobs/job-101/exceptions", - "{\"all-exceptions\":[],\"root-exception\":\"\",\"timestamp\":2000,\"truncated\":false," - + "\"exceptionHistory\":{\"entries\":[],\"truncated\":false}}"); - HttpServer server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); - server.createContext("/jobs/job-101", exchange -> { - String response = responses.get(exchange.getRequestURI().getPath()); - boolean validRequest = "GET".equals(exchange.getRequestMethod()) && response != null; - byte[] body = (validRequest ? response : "{\"errors\":[\"Unexpected request\"]}") - .getBytes(StandardCharsets.UTF_8); - exchange.getResponseHeaders().set("Content-Type", "application/json"); - exchange.sendResponseHeaders(validRequest ? 200 : 404, body.length); - try (OutputStream output = exchange.getResponseBody()) { - output.write(body); + private static class Fixture { + + private static final JobInstanceService JOB_INSTANCE_SERVICE = mock(JobInstanceService.class); + + private static StaticApplicationContext applicationContext; + private static ApplicationContext previousApplicationContext; + + private final Map jobTenants = new HashMap<>(); + private final Map persistedStatuses = new HashMap<>(); + + static void registerServices() { + previousApplicationContext = SpringContextUtils.applicationContext; + applicationContext = new StaticApplicationContext(); + applicationContext.getBeanFactory().registerSingleton("jobInstanceServiceImpl", JOB_INSTANCE_SERVICE); + applicationContext.getBeanFactory().registerSingleton("monitorServiceImpl", mock(MonitorService.class)); + applicationContext + .getBeanFactory() + .registerSingleton("jobHistoryServiceImpl", mock(JobHistoryService.class)); + applicationContext + .getBeanFactory() + .registerSingleton("clusterInstanceServiceImpl", mock(ClusterInstanceService.class)); + applicationContext.getBeanFactory().registerSingleton("historyServiceImpl", mock(HistoryService.class)); + applicationContext.getBeanFactory().registerSingleton("taskServiceImpl", mock(TaskService.class)); + applicationContext + .getBeanFactory() + .registerSingleton("alertHistoryServiceImpl", mock(AlertHistoryService.class)); + applicationContext.getBeanFactory().registerSingleton("userServiceImpl", mock(UserService.class)); + AlertRuleServiceImpl alertRuleService = mock(AlertRuleServiceImpl.class); + AlertRulesMapper alertRulesMapper = mock(AlertRulesMapper.class); + when(alertRuleService.getBaseMapper()).thenReturn(alertRulesMapper); + when(alertRulesMapper.selectWithTemplate()).thenReturn(Collections.emptyList()); + applicationContext.getBeanFactory().registerSingleton("alertRuleServiceImpl", alertRuleService); + SpringContextUtils.applicationContext = applicationContext; + } + + static void restoreApplicationContext() { + SpringContextUtils.applicationContext = previousApplicationContext; + applicationContext.close(); + } + + void prepareTenantScopedPersistence() { + TenantContextHolder.clear(); + reset(JOB_INSTANCE_SERVICE); + doAnswer(invocation -> { + TenantContextHolder.set(jobTenants.get(invocation.getArgument(0))); + return null; + }) + .when(JOB_INSTANCE_SERVICE) + .initTenantByJobInstanceId(anyInt()); + when(JOB_INSTANCE_SERVICE.updateById(any(JobInstance.class))).thenAnswer(invocation -> { + JobInstance instance = invocation.getArgument(0); + // A tenant-filtered update matches no row when a worker keeps another job's tenant. + if (!jobTenants.get(instance.getId()).equals(TenantContextHolder.get())) { + return false; + } + persistedStatuses.put(instance.getId(), instance.getStatus()); + return true; + }); + } + + void clearTenantContext() { + TenantContextHolder.clear(); + } + + void persistsTerminalStatusUsingTheJobTenant() { + TenantContextHolder.set(1); + FlinkJobTask task = task(101, 2); + + assertTrue(task.dealTask()); + + assertAll( + () -> assertEquals(JobStatus.UNKNOWN.getValue(), persistedStatuses.get(101)), + () -> assertEquals(1, TenantContextHolder.get())); + } + + void clearsTenantContextWhenTheWorkerHadNoTenant() { + FlinkJobTask task = task(101, 2); + + assertTrue(task.dealTask()); + + assertAll( + () -> assertEquals(JobStatus.UNKNOWN.getValue(), persistedStatuses.get(101)), + () -> assertNull(TenantContextHolder.get())); + } + + void refreshesDifferentTenantsOnTheSameWorker() { + TenantContextHolder.set(1); + FlinkJobTask firstTask = task(101, 2); + FlinkJobTask secondTask = task(102, 3); + + assertTrue(firstTask.dealTask()); + assertTrue(secondTask.dealTask()); + + assertAll( + () -> assertEquals(JobStatus.UNKNOWN.getValue(), persistedStatuses.get(101)), + () -> assertEquals(JobStatus.UNKNOWN.getValue(), persistedStatuses.get(102)), + () -> assertEquals(1, TenantContextHolder.get())); + } + + void restoresPreviousTenantWhenRefreshFails() { + TenantContextHolder.set(1); + FlinkJobTask task = task(101, 2); + IllegalStateException failure = new IllegalStateException("Unable to persist status"); + doAnswer(invocation -> { + assertEquals(2, TenantContextHolder.get()); + throw failure; + }) + .when(JOB_INSTANCE_SERVICE) + .updateById(any(JobInstance.class)); + + assertSame(failure, assertThrows(IllegalStateException.class, task::dealTask)); + assertEquals(1, TenantContextHolder.get()); + } + + void clearsTenantContextWhenRefreshFailsWithoutPreviousTenant() { + FlinkJobTask task = task(101, 2); + IllegalStateException failure = new IllegalStateException("Unable to persist status"); + doThrow(failure).when(JOB_INSTANCE_SERVICE).updateById(any(JobInstance.class)); + + assertSame(failure, assertThrows(IllegalStateException.class, task::dealTask)); + assertNull(TenantContextHolder.get()); + } + + void preservesTheTenantForManualRefresh() { + TenantContextHolder.set(2); + FlinkJobTask task = task(101, 2); + + assertTrue(task.dealTask()); + + assertAll( + () -> assertEquals(JobStatus.UNKNOWN.getValue(), persistedStatuses.get(101)), + () -> assertEquals(2, TenantContextHolder.get())); + } + + void persistsFailedYarnSessionJobWithTenantFilteringEnabled() throws Exception { + HttpServer flink = failedJobServer(); + Boolean metricsEnabled = + SystemConfiguration.getInstances().getMetricsSysEnable().getValue(); + Integer resendInterval = + SystemConfiguration.getInstances().getJobReSendDiffSecond().getValue(); + SystemConfiguration.getInstances().getMetricsSysEnable().setValue(false); + SystemConfiguration.getInstances().getJobReSendDiffSecond().setValue(60); + UnpooledDataSource dataSource = + new UnpooledDataSource("org.h2.Driver", "jdbc:h2:mem:" + UUID.randomUUID(), "sa", ""); + + try (Connection connection = dataSource.getConnection(); + Statement statement = connection.createStatement()) { + statement.execute( + "CREATE TABLE dinky_job_instance (id INT PRIMARY KEY, tenant_id INT, status VARCHAR(32))"); + statement.execute("INSERT INTO dinky_job_instance VALUES (101, 2, 'RUNNING'), (102, 1, 'RUNNING')"); + Configuration configuration = + new Configuration(new Environment("test", new JdbcTransactionFactory(), dataSource)); + configuration.addInterceptor( + new MybatisPlusConfig(new MybatisPlusFillProperties()).mybatisPlusInterceptor()); + configuration.addMapper(JobStatusMapper.class); + + try (SqlSession session = + new SqlSessionFactoryBuilder().build(configuration).openSession(true)) { + JobStatusMapper mapper = session.getMapper(JobStatusMapper.class); + TenantContextHolder.set(1); + assertFalse(TenantContextHolder.isIgnoreTenant()); + FlinkJobTask task = task(101, 2); + JobInfoDetail detail = task.getJobInfoDetail(); + JobInstance instance = detail.getInstance(); + instance.setTaskId(10); + instance.setName("failed-job"); + instance.setJid("job-101"); + instance.setStatus(JobStatus.FAILED.getValue()); + assertEquals( + 0, mapper.updateStatus(instance), "Another tenant must not be able to update this row"); + instance.setStatus(JobStatus.RUNNING.getValue()); + detail.setJobDataDto( + JobDataDto.builder().id(101).tenantId(2).build()); + ClusterInstance cluster = new ClusterInstance(); + cluster.setName("yarn-session"); + cluster.setType(GatewayType.YARN_SESSION.getLongValue()); + cluster.setJobManagerHost("127.0.0.1:" + flink.getAddress().getPort()); + cluster.setHosts(cluster.getJobManagerHost()); + detail.setClusterInstance(cluster); + doAnswer(invocation -> mapper.updateStatus(invocation.getArgument(0)) == 1) + .when(JOB_INSTANCE_SERVICE) + .updateById(any(JobInstance.class)); + + assertTrue(task.dealTask()); + + assertAll( + () -> assertEquals(JobStatus.FAILED.getValue(), persistedStatus(connection, 101)), + () -> assertEquals(JobStatus.RUNNING.getValue(), persistedStatus(connection, 102)), + () -> assertEquals(1, TenantContextHolder.get()), + () -> assertFalse(TenantContextHolder.isIgnoreTenant())); + } + } finally { + flink.stop(0); + SystemConfiguration.getInstances().getMetricsSysEnable().setValue(metricsEnabled); + SystemConfiguration.getInstances().getJobReSendDiffSecond().setValue(resendInterval); } - }); - server.start(); - return server; - } + } + + private static HttpServer failedJobServer() throws IOException { + Map responses = new HashMap<>(); + responses.put( + "/jobs/job-101", + "{\"jid\":\"job-101\",\"name\":\"failed-job\",\"state\":\"FAILED\"," + + "\"start-time\":1000,\"end-time\":2000,\"duration\":1000,\"vertices\":[],\"plan\":{\"nodes\":[]}}"); + responses.put( + "/jobs/job-101/config", "{\"jid\":\"job-101\",\"name\":\"failed-job\",\"execution-config\":{}}"); + responses.put("/jobs/job-101/checkpoints", "{\"errors\":[]}"); + responses.put("/jobs/job-101/checkpoints/config", "{\"errors\":[]}"); + responses.put( + "/jobs/job-101/exceptions", + "{\"all-exceptions\":[],\"root-exception\":\"\",\"timestamp\":2000,\"truncated\":false," + + "\"exceptionHistory\":{\"entries\":[],\"truncated\":false}}"); + HttpServer server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); + server.createContext("/jobs/job-101", exchange -> { + String response = responses.get(exchange.getRequestURI().getPath()); + boolean validRequest = "GET".equals(exchange.getRequestMethod()) && response != null; + byte[] body = (validRequest ? response : "{\"errors\":[\"Unexpected request\"]}") + .getBytes(StandardCharsets.UTF_8); + exchange.getResponseHeaders().set("Content-Type", "application/json"); + exchange.sendResponseHeaders(validRequest ? 200 : 404, body.length); + try (OutputStream output = exchange.getResponseBody()) { + output.write(body); + } + }); + server.start(); + return server; + } - private static String persistedStatus(Connection connection, int id) throws SQLException { - try (PreparedStatement statement = - connection.prepareStatement("SELECT status FROM dinky_job_instance WHERE id = ?")) { - statement.setInt(1, id); - try (ResultSet rows = statement.executeQuery()) { - assertTrue(rows.next()); - return rows.getString("status"); + private static String persistedStatus(Connection connection, int id) throws SQLException { + try (PreparedStatement statement = + connection.prepareStatement("SELECT status FROM dinky_job_instance WHERE id = ?")) { + statement.setInt(1, id); + try (ResultSet rows = statement.executeQuery()) { + assertTrue(rows.next()); + return rows.getString("status"); + } } } - } - interface JobStatusMapper { - @Update("UPDATE dinky_job_instance SET status = #{status} WHERE id = #{id}") - int updateStatus(JobInstance instance); - } + interface JobStatusMapper { + @Update("UPDATE dinky_job_instance SET status = #{status} WHERE id = #{id}") + int updateStatus(JobInstance instance); + } - private FlinkJobTask task(int id, int tenantId) { - JobInstance instance = new JobInstance(); - instance.setId(id); - instance.setTenantId(tenantId); - instance.setStatus(JobStatus.RUNNING.getValue()); - jobTenants.put(id, tenantId); - persistedStatuses.put(id, JobStatus.RUNNING.getValue()); - - // Missing cluster metadata is a terminal refresh path that needs no Flink or YARN server. - JobInfoDetail detail = new JobInfoDetail(id); - detail.setInstance(instance); - FlinkJobTask task = new FlinkJobTask(); - task.setJobInfoDetail(detail); - return task; + private FlinkJobTask task(int id, int tenantId) { + JobInstance instance = new JobInstance(); + instance.setId(id); + instance.setTenantId(tenantId); + instance.setStatus(JobStatus.RUNNING.getValue()); + jobTenants.put(id, tenantId); + persistedStatuses.put(id, JobStatus.RUNNING.getValue()); + + // Missing cluster metadata is a terminal refresh path that needs no Flink or YARN server. + JobInfoDetail detail = new JobInfoDetail(id); + detail.setInstance(instance); + FlinkJobTask task = new FlinkJobTask(); + task.setJobInfoDetail(detail); + return task; + } } }