From b16310ae736b8363927c8aa1dd3fc7fe396d3332 Mon Sep 17 00:00:00 2001 From: Eduardo Flores Date: Sat, 18 Jul 2026 02:14:11 -0600 Subject: [PATCH] feat(plugins): serialize lifecycle transitions --- .../loader/DynamicJobLoaderService.java | 34 +++-- .../loader/DynamicJobLoaderServiceTest.java | 139 +++++++++++++++++- 2 files changed, 161 insertions(+), 12 deletions(-) diff --git a/fr-batch-service/src/main/java/com/fronzec/frbatchservice/batchjobs/plugins/loader/DynamicJobLoaderService.java b/fr-batch-service/src/main/java/com/fronzec/frbatchservice/batchjobs/plugins/loader/DynamicJobLoaderService.java index fda0971..7fcba0f 100644 --- a/fr-batch-service/src/main/java/com/fronzec/frbatchservice/batchjobs/plugins/loader/DynamicJobLoaderService.java +++ b/fr-batch-service/src/main/java/com/fronzec/frbatchservice/batchjobs/plugins/loader/DynamicJobLoaderService.java @@ -53,12 +53,11 @@ public class DynamicJobLoaderService { private final Map classLoaders = new ConcurrentHashMap<>(); /** - * Per-definition load locks. Serializes concurrent {@link #loadJob(Long)} calls for the same - * definition ID so the already-loaded guard and the classloader/registry writes execute - * atomically. Without this, two concurrent loads of the same definition both pass the guard, each - * opens a JAR file-handle, and the loser's {@code classLoaders.put} silently overwrites the - * winner's entry — leaking the winner's handle. {@link #loadJob(Long)} allocates an entry only - * after confirming the definition exists, so unknown IDs cannot grow the map without bound. + * Per-definition lifecycle locks. Serialize concurrent load, unload, and reload calls for the + * same definition ID so registry and classloader transitions execute atomically. A holder is + * allocated only after confirming the definition exists, so unknown IDs cannot grow this map + * without bound. This coordination is limited to calls in this service instance; it does not + * provide distributed lifecycle locking. */ private final Map lifecycleLocks = new ConcurrentHashMap<>(); @@ -135,9 +134,7 @@ private static final class LifecycleLockHolder { } /** - * Performs the actual load. Always invoked while holding the per-definition lock acquired in - * {@link #loadJob(Long)}, so the already-loaded guard and the classloader/registry writes are - * atomic with respect to other loads of the same definition. + * Performs the actual load while the caller owns the per-definition lifecycle lock. */ private LoadResult doLoadJob(Long definitionId) { JobDefinitionEntity entity = @@ -349,6 +346,14 @@ private LoadResult doLoadJob(Long definitionId) { * {@code false} */ public LoadResult unloadJob(Long definitionId, boolean force) { + if (jobDefinitionRepository.findById(definitionId).isEmpty()) { + throw new NoSuchElementException("Job definition not found for id: " + definitionId); + } + return withDefinitionLock(definitionId, () -> doUnloadJob(definitionId, force)); + } + + /** Performs the actual unload while the caller owns the per-definition lifecycle lock. */ + private LoadResult doUnloadJob(Long definitionId, boolean force) { long unloadStartNanos = System.nanoTime(); JobDefinitionEntity entity = @@ -437,8 +442,15 @@ public LoadResult unloadJob(Long definitionId, boolean force) { * @throws JobLoadException if the load step fails */ public LoadResult reloadJob(Long definitionId) { - unloadJob(definitionId, true); - return loadJob(definitionId); + if (jobDefinitionRepository.findById(definitionId).isEmpty()) { + throw new NoSuchElementException("Job definition not found for id: " + definitionId); + } + return withDefinitionLock( + definitionId, + () -> { + doUnloadJob(definitionId, true); + return doLoadJob(definitionId); + }); } /** diff --git a/fr-batch-service/src/test/java/com/fronzec/frbatchservice/batchjobs/plugins/loader/DynamicJobLoaderServiceTest.java b/fr-batch-service/src/test/java/com/fronzec/frbatchservice/batchjobs/plugins/loader/DynamicJobLoaderServiceTest.java index 918ec8c..7b0bf3e 100644 --- a/fr-batch-service/src/test/java/com/fronzec/frbatchservice/batchjobs/plugins/loader/DynamicJobLoaderServiceTest.java +++ b/fr-batch-service/src/test/java/com/fronzec/frbatchservice/batchjobs/plugins/loader/DynamicJobLoaderServiceTest.java @@ -397,9 +397,11 @@ void loadJob_alreadyLoaded_throwsIllegalStateException() { // ── Unload ──────────────────────────────────────────────────────────────── @Test - void unloadJob_success_updatesStatus() { + void unloadJob_success_releasesClassLoaderOwnershipAndUpdatesStatus() throws Exception { JobDefinitionEntity entity = createEntity(8L, "unloadable-job", true, LoadResult.LOADED); + DynamicJobClassLoader classLoader = mock(DynamicJobClassLoader.class); + classLoadersOf(service).put(8L, classLoader); when(jobDefinitionRepository.findById(8L)).thenReturn(Optional.of(entity)); when(jobRepository.findRunningJobExecutions("unloadable-job")) @@ -412,6 +414,8 @@ void unloadJob_success_updatesStatus() { assertEquals(LoadResult.UNLOADED, result.status()); assertEquals(LoadResult.UNLOADED, entity.getLoadStatus()); verify(pluginRegistryService).unregisterDynamicPlugin("unloadable-job"); + verify(classLoader).cleanup(); + assertFalse(classLoadersOf(service).containsKey(8L)); } @Test @@ -480,6 +484,139 @@ void reloadJob_unloadsThenLoads_atomic() throws Exception { } } + @Test + void reloadJob_holdsOneLockAcrossUnloadAndReplacementLoad() throws Exception { + JobDefinitionEntity entity = createEntity(21L, "test-job", true, LoadResult.LOADED); + entity.setMainClassName(TestBatchJobPlugin.class.getName()); + when(jobDefinitionRepository.findById(21L)).thenReturn(Optional.of(entity)); + + AtomicInteger unregisterCalls = new AtomicInteger(); + CountDownLatch secondUnregisterEntered = new CountDownLatch(1); + doAnswer( + invocation -> { + if (unregisterCalls.incrementAndGet() == 2) { + secondUnregisterEntered.countDown(); + } + return null; + }) + .when(pluginRegistryService) + .unregisterDynamicPlugin("test-job"); + + CountDownLatch replacementRegisterEntered = new CountDownLatch(1); + CountDownLatch releaseReplacementRegister = new CountDownLatch(1); + doAnswer( + invocation -> { + replacementRegisterEntered.countDown(); + assertTrue(releaseReplacementRegister.await(5, TimeUnit.SECONDS)); + return null; + }) + .when(pluginRegistryService) + .registerDynamicPlugin(any(BatchJobPlugin.class)); + + ExecutorService pool = Executors.newFixedThreadPool(2); + try { + Future reload = pool.submit(() -> service.reloadJob(21L)); + assertTrue(replacementRegisterEntered.await(5, TimeUnit.SECONDS)); + + Future concurrentUnload = pool.submit(() -> service.unloadJob(21L, true)); + assertFalse( + secondUnregisterEntered.await(250, TimeUnit.MILLISECONDS), + "Concurrent unload must not enter while reload owns the lifecycle lock"); + + releaseReplacementRegister.countDown(); + reload.get(5, TimeUnit.SECONDS); + concurrentUnload.get(5, TimeUnit.SECONDS); + assertEquals(2, unregisterCalls.get()); + } finally { + releaseReplacementRegister.countDown(); + pool.shutdownNow(); + } + } + + @Test + void reloadJob_unregisterFailureAbortsReplacementAndPreservesLoadedState() throws Exception { + JobDefinitionEntity entity = createEntity(22L, "test-job", true, LoadResult.LOADED); + DynamicJobClassLoader oldLoader = mock(DynamicJobClassLoader.class); + classLoadersOf(service).put(22L, oldLoader); + when(jobDefinitionRepository.findById(22L)).thenReturn(Optional.of(entity)); + doThrow(new RuntimeException("unregister failed")) + .when(pluginRegistryService) + .unregisterDynamicPlugin("test-job"); + + assertThrows(RuntimeException.class, () -> service.reloadJob(22L)); + + verify(pluginRegistryService, never()).registerDynamicPlugin(any(BatchJobPlugin.class)); + verify(oldLoader, never()).cleanup(); + verify(jobDefinitionRepository, never()).save(entity); + assertSame(oldLoader, classLoadersOf(service).get(22L)); + assertEquals(LoadResult.LOADED, entity.getLoadStatus()); + } + + @Test + void reloadJob_replacementFailureCleansCandidateAndPersistsFailed() throws Exception { + JobDefinitionEntity entity = createEntity(23L, "test-job", true, LoadResult.LOADED); + DynamicJobClassLoader oldLoader = mock(DynamicJobClassLoader.class); + classLoadersOf(service).put(23L, oldLoader); + when(jobDefinitionRepository.findById(23L)).thenReturn(Optional.of(entity)); + + try (MockedConstruction mocked = + mockConstruction( + DynamicJobClassLoader.class, + (mock, context) -> + when(mock.loadClass("com.test.TestPlugin")) + .thenThrow(new ClassNotFoundException("replacement missing")))) { + assertThrows(JobLoadException.class, () -> service.reloadJob(23L)); + verify(oldLoader).cleanup(); + verify(mocked.constructed().getFirst()).cleanup(); + } + + verify(pluginRegistryService).unregisterDynamicPlugin("test-job"); + verify(pluginRegistryService, never()).registerDynamicPlugin(any(BatchJobPlugin.class)); + verify(jobDefinitionRepository, atLeastOnce()).save(entity); + assertTrue(classLoadersOf(service).isEmpty()); + assertEquals(LoadResult.FAILED, entity.getLoadStatus()); + } + + @Test + void loadJob_recoversFailedDefinitionWithOneRegistrationAndClassLoader() throws Exception { + JobDefinitionEntity entity = createEntity(24L, "test-job", true, LoadResult.FAILED); + when(jobDefinitionRepository.findById(24L)).thenReturn(Optional.of(entity)); + + try (MockedConstruction mocked = + mockConstruction( + DynamicJobClassLoader.class, + (mock, context) -> + when(mock.loadClass("com.test.TestPlugin")) + .thenAnswer(invocation -> TestBatchJobPlugin.class))) { + assertEquals(LoadResult.LOADED, service.loadJob(24L).status()); + assertEquals(1, mocked.constructed().size()); + } + + verify(pluginRegistryService, times(1)).registerDynamicPlugin(any(BatchJobPlugin.class)); + assertEquals(1, classLoadersOf(service).size()); + assertEquals(LoadResult.LOADED, entity.getLoadStatus()); + } + + @Test + void loadJob_failedRecoveryLeavesFailedWithoutLoaderOrRegistration() throws Exception { + JobDefinitionEntity entity = createEntity(25L, "test-job", true, LoadResult.FAILED); + when(jobDefinitionRepository.findById(25L)).thenReturn(Optional.of(entity)); + + try (MockedConstruction mocked = + mockConstruction( + DynamicJobClassLoader.class, + (mock, context) -> + when(mock.loadClass("com.test.TestPlugin")) + .thenThrow(new ClassNotFoundException("recovery missing")))) { + assertThrows(JobLoadException.class, () -> service.loadJob(25L)); + verify(mocked.constructed().getFirst()).cleanup(); + } + + verify(pluginRegistryService, never()).registerDynamicPlugin(any(BatchJobPlugin.class)); + assertTrue(classLoadersOf(service).isEmpty()); + assertEquals(LoadResult.FAILED, entity.getLoadStatus()); + } + // ── loadAllEnabled ──────────────────────────────────────────────────────── @Test