diff --git a/Directory.Build.props b/Directory.Build.props
index e6a9caab..e7e7d992 100644
--- a/Directory.Build.props
+++ b/Directory.Build.props
@@ -13,7 +13,31 @@
- $(NoWarn);S6966;CS1591;RCS1060;S6354;S4226;RCS1208
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+ $(NoWarn);S6966;CS1591;RCS1060;S6354;S4226;RCS1208;IDE0009;IDE0040;IDE0022;IDE0053;IDE0021;IDE0011;IDE0044;IDE0290;RCS1037;RCS1163;RCS1169;RCS1085;RCS1181;S4004;S1192;S131;S2930;S2931;S3267;S2933;S6602;S6605;S6667;CA1822
false
diff --git a/FlinkDotNet/FlinkDotNet.JobManager.Tests/ClusterControllerTests.cs b/FlinkDotNet/FlinkDotNet.JobManager.Tests/ClusterControllerTests.cs
new file mode 100644
index 00000000..26e4d2d0
--- /dev/null
+++ b/FlinkDotNet/FlinkDotNet.JobManager.Tests/ClusterControllerTests.cs
@@ -0,0 +1,246 @@
+// Copyright 2025 FlinkDotNet
+// Licensed under the Apache License, Version 2.0.
+// See LICENSE file in the project root for full license information.
+
+using FlinkDotNet.JobManager.Controllers;
+using FlinkDotNet.JobManager.Interfaces;
+using FlinkDotNet.JobManager.Models;
+using FlinkDotNet.JobManager.Models.Responses;
+using FluentAssertions;
+using Microsoft.AspNetCore.Mvc;
+using Microsoft.Extensions.Logging;
+using Moq;
+
+namespace FlinkDotNet.JobManager.Tests;
+
+public class ClusterControllerTests
+{
+ private readonly Mock _mockResourceManager;
+ private readonly Mock _mockDispatcher;
+ private readonly Mock> _mockLogger;
+ private readonly ClusterController _controller;
+
+ public ClusterControllerTests()
+ {
+ _mockResourceManager = new Mock();
+ _mockDispatcher = new Mock();
+ _mockLogger = new Mock>();
+ _controller = new ClusterController(
+ _mockResourceManager.Object,
+ _mockDispatcher.Object,
+ _mockLogger.Object);
+ }
+
+ [Fact]
+ public void Constructor_WithNullResourceManager_ThrowsArgumentNullException()
+ {
+ // Act & Assert
+ var act = () => new ClusterController(
+ null!,
+ _mockDispatcher.Object,
+ _mockLogger.Object);
+
+ act.Should().Throw().WithParameterName("resourceManager");
+ }
+
+ [Fact]
+ public void Constructor_WithNullDispatcher_ThrowsArgumentNullException()
+ {
+ // Act & Assert
+ var act = () => new ClusterController(
+ _mockResourceManager.Object,
+ null!,
+ _mockLogger.Object);
+
+ act.Should().Throw().WithParameterName("dispatcher");
+ }
+
+ [Fact]
+ public void Constructor_WithNullLogger_ThrowsArgumentNullException()
+ {
+ // Act & Assert
+ var act = () => new ClusterController(
+ _mockResourceManager.Object,
+ _mockDispatcher.Object,
+ null!);
+
+ act.Should().Throw().WithParameterName("logger");
+ }
+
+ [Fact]
+ public async Task GetOverview_ReturnsClusterOverview()
+ {
+ // Arrange
+ var taskManagerIds = new List { "tm-1", "tm-2" };
+ var allSlots = new List
+ {
+ new TaskSlot { SlotId = "slot-1", TaskManagerId = "tm-1", IsAllocated = true },
+ new TaskSlot { SlotId = "slot-2", TaskManagerId = "tm-1", IsAllocated = false },
+ new TaskSlot { SlotId = "slot-3", TaskManagerId = "tm-2", IsAllocated = true },
+ new TaskSlot { SlotId = "slot-4", TaskManagerId = "tm-2", IsAllocated = false }
+ };
+ var availableSlots = allSlots.Where(s => !s.IsAllocated).ToList();
+
+ var jobs = new List
+ {
+ new JobStatus { JobId = "job-1", State = JobExecutionState.Running },
+ new JobStatus { JobId = "job-2", State = JobExecutionState.Finished },
+ new JobStatus { JobId = "job-3", State = JobExecutionState.Failed },
+ new JobStatus { JobId = "job-4", State = JobExecutionState.Canceled }
+ };
+
+ _mockResourceManager.Setup(rm => rm.GetRegisteredTaskManagers()).Returns(taskManagerIds);
+ _mockResourceManager.Setup(rm => rm.GetAllSlots()).Returns(allSlots);
+ _mockResourceManager.Setup(rm => rm.GetAvailableSlots()).Returns(availableSlots);
+ _mockDispatcher.Setup(d => d.ListJobsAsync(It.IsAny())).ReturnsAsync(jobs);
+
+ // Act
+ var result = await _controller.GetOverview();
+
+ // Assert
+ result.Should().BeOfType();
+ var okResult = result as OkObjectResult;
+ var response = okResult!.Value as ClusterOverviewResponse;
+
+ response.Should().NotBeNull();
+ response!.TaskManagers.Should().Be(2);
+ response.TotalSlots.Should().Be(4);
+ response.AvailableSlots.Should().Be(2);
+ response.RunningJobs.Should().Be(1);
+ response.FinishedJobs.Should().Be(1);
+ response.FailedJobs.Should().Be(1);
+ response.CanceledJobs.Should().Be(1);
+ }
+
+ [Fact]
+ public void ListTaskManagers_ReturnsTaskManagerList()
+ {
+ // Arrange
+ var taskManagerIds = new List { "tm-1", "tm-2" };
+ var allSlots = new List
+ {
+ new TaskSlot { SlotId = "slot-1", TaskManagerId = "tm-1", IsAllocated = true },
+ new TaskSlot { SlotId = "slot-2", TaskManagerId = "tm-1", IsAllocated = false },
+ new TaskSlot { SlotId = "slot-3", TaskManagerId = "tm-2", IsAllocated = true },
+ new TaskSlot { SlotId = "slot-4", TaskManagerId = "tm-2", IsAllocated = false },
+ new TaskSlot { SlotId = "slot-5", TaskManagerId = "tm-2", IsAllocated = false }
+ };
+ var availableSlots = allSlots.Where(s => !s.IsAllocated).ToList();
+
+ _mockResourceManager.Setup(rm => rm.GetRegisteredTaskManagers()).Returns(taskManagerIds);
+ _mockResourceManager.Setup(rm => rm.GetAllSlots()).Returns(allSlots);
+ _mockResourceManager.Setup(rm => rm.GetAvailableSlots()).Returns(availableSlots);
+
+ // Act
+ var result = _controller.ListTaskManagers();
+
+ // Assert
+ result.Should().BeOfType();
+ var okResult = result as OkObjectResult;
+ var response = okResult!.Value as TaskManagerListResponse;
+
+ response.Should().NotBeNull();
+ response!.TaskManagers.Should().HaveCount(2);
+
+ var tm1 = response.TaskManagers.First(tm => tm.TaskManagerId == "tm-1");
+ tm1.TotalSlots.Should().Be(2);
+ tm1.FreeSlots.Should().Be(1);
+
+ var tm2 = response.TaskManagers.First(tm => tm.TaskManagerId == "tm-2");
+ tm2.TotalSlots.Should().Be(3);
+ tm2.FreeSlots.Should().Be(2);
+ }
+
+ [Fact]
+ public async Task GetOverview_WithNoJobs_ReturnsZeroJobCounts()
+ {
+ // Arrange
+ var taskManagerIds = new List { "tm-1" };
+ var allSlots = new List
+ {
+ new TaskSlot { SlotId = "slot-1", TaskManagerId = "tm-1", IsAllocated = false }
+ };
+ var jobs = new List();
+
+ _mockResourceManager.Setup(rm => rm.GetRegisteredTaskManagers()).Returns(taskManagerIds);
+ _mockResourceManager.Setup(rm => rm.GetAllSlots()).Returns(allSlots);
+ _mockResourceManager.Setup(rm => rm.GetAvailableSlots()).Returns(allSlots);
+ _mockDispatcher.Setup(d => d.ListJobsAsync(It.IsAny())).ReturnsAsync(jobs);
+
+ // Act
+ var result = await _controller.GetOverview();
+
+ // Assert
+ result.Should().BeOfType();
+ var okResult = result as OkObjectResult;
+ var response = okResult!.Value as ClusterOverviewResponse;
+
+ response!.RunningJobs.Should().Be(0);
+ response.FinishedJobs.Should().Be(0);
+ response.FailedJobs.Should().Be(0);
+ response.CanceledJobs.Should().Be(0);
+ }
+
+ [Fact]
+ public void ListTaskManagers_WithNoTaskManagers_ReturnsEmptyList()
+ {
+ // Arrange
+ var taskManagerIds = new List();
+ var allSlots = new List();
+ var availableSlots = new List();
+
+ _mockResourceManager.Setup(rm => rm.GetRegisteredTaskManagers()).Returns(taskManagerIds);
+ _mockResourceManager.Setup(rm => rm.GetAllSlots()).Returns(allSlots);
+ _mockResourceManager.Setup(rm => rm.GetAvailableSlots()).Returns(availableSlots);
+
+ // Act
+ var result = _controller.ListTaskManagers();
+
+ // Assert
+ result.Should().BeOfType();
+ var okResult = result as OkObjectResult;
+ var response = okResult!.Value as TaskManagerListResponse;
+
+ response!.TaskManagers.Should().BeEmpty();
+ }
+
+ [Fact]
+ public async Task GetOverview_WithMultipleJobStates_CountsCorrectly()
+ {
+ // Arrange
+ var taskManagerIds = new List { "tm-1" };
+ var allSlots = new List
+ {
+ new TaskSlot { SlotId = "slot-1", TaskManagerId = "tm-1", IsAllocated = false }
+ };
+
+ var jobs = new List
+ {
+ new JobStatus { JobId = "job-1", State = JobExecutionState.Running },
+ new JobStatus { JobId = "job-2", State = JobExecutionState.Running },
+ new JobStatus { JobId = "job-3", State = JobExecutionState.Running },
+ new JobStatus { JobId = "job-4", State = JobExecutionState.Finished },
+ new JobStatus { JobId = "job-5", State = JobExecutionState.Failed },
+ new JobStatus { JobId = "job-6", State = JobExecutionState.Failed },
+ new JobStatus { JobId = "job-7", State = JobExecutionState.Canceled }
+ };
+
+ _mockResourceManager.Setup(rm => rm.GetRegisteredTaskManagers()).Returns(taskManagerIds);
+ _mockResourceManager.Setup(rm => rm.GetAllSlots()).Returns(allSlots);
+ _mockResourceManager.Setup(rm => rm.GetAvailableSlots()).Returns(allSlots);
+ _mockDispatcher.Setup(d => d.ListJobsAsync(It.IsAny())).ReturnsAsync(jobs);
+
+ // Act
+ var result = await _controller.GetOverview();
+
+ // Assert
+ result.Should().BeOfType();
+ var okResult = result as OkObjectResult;
+ var response = okResult!.Value as ClusterOverviewResponse;
+
+ response!.RunningJobs.Should().Be(3);
+ response.FinishedJobs.Should().Be(1);
+ response.FailedJobs.Should().Be(2);
+ response.CanceledJobs.Should().Be(1);
+ }
+}
diff --git a/FlinkDotNet/FlinkDotNet.JobManager.Tests/DispatcherTests.cs b/FlinkDotNet/FlinkDotNet.JobManager.Tests/DispatcherTests.cs
new file mode 100644
index 00000000..2de9a940
--- /dev/null
+++ b/FlinkDotNet/FlinkDotNet.JobManager.Tests/DispatcherTests.cs
@@ -0,0 +1,328 @@
+// Copyright 2025 FlinkDotNet
+// Licensed under the Apache License, Version 2.0.
+// See LICENSE file in the project root for full license information.
+
+using FlinkDotNet.JobManager.Implementation;
+using FlinkDotNet.JobManager.Interfaces;
+using FlinkDotNet.JobManager.Models;
+using FluentAssertions;
+using Microsoft.Extensions.Logging;
+using Moq;
+using Temporalio.Client;
+
+namespace FlinkDotNet.JobManager.Tests;
+
+public class DispatcherTests
+{
+ private readonly Mock _mockResourceManager;
+ private readonly Mock _mockTemporalClient;
+ private readonly Mock _mockLoggerFactory;
+ private readonly Mock> _mockLogger;
+ private readonly Dispatcher _dispatcher;
+
+ public DispatcherTests()
+ {
+ _mockResourceManager = new Mock();
+ _mockTemporalClient = new Mock();
+ _mockLoggerFactory = new Mock();
+ _mockLogger = new Mock>();
+
+ _mockLoggerFactory
+ .Setup(lf => lf.CreateLogger(It.IsAny()))
+ .Returns(_mockLogger.Object);
+
+ _dispatcher = new Dispatcher(
+ _mockResourceManager.Object,
+ _mockTemporalClient.Object,
+ _mockLoggerFactory.Object);
+ }
+
+ [Fact]
+ public void Constructor_WithNullResourceManager_ThrowsArgumentNullException()
+ {
+ // Act & Assert
+ var act = () => new Dispatcher(
+ null!,
+ _mockTemporalClient.Object,
+ _mockLoggerFactory.Object);
+
+ act.Should().Throw().WithParameterName("resourceManager");
+ }
+
+ [Fact]
+ public void Constructor_WithNullTemporalClient_ThrowsArgumentNullException()
+ {
+ // Act & Assert
+ var act = () => new Dispatcher(
+ _mockResourceManager.Object,
+ null!,
+ _mockLoggerFactory.Object);
+
+ act.Should().Throw().WithParameterName("temporalClient");
+ }
+
+ [Fact]
+ public void Constructor_WithNullLoggerFactory_ThrowsArgumentNullException()
+ {
+ // Act & Assert
+ var act = () => new Dispatcher(
+ _mockResourceManager.Object,
+ _mockTemporalClient.Object,
+ null!);
+
+ act.Should().Throw().WithParameterName("loggerFactory");
+ }
+
+ [Fact]
+ public async Task SubmitJobAsync_WithNullJobGraph_ThrowsArgumentNullException()
+ {
+ // Act & Assert
+ var act = async () => await _dispatcher.SubmitJobAsync(null!);
+ await act.Should().ThrowAsync().WithParameterName("jobGraph");
+ }
+
+ [Fact]
+ public async Task SubmitJobAsync_WithValidJobGraph_ReturnsSuccess()
+ {
+ // Arrange
+ var jobGraph = CreateValidJobGraph();
+
+ // Act
+ var result = await _dispatcher.SubmitJobAsync(jobGraph);
+
+ // Assert
+ result.Should().NotBeNull();
+ if (!result.Success)
+ {
+ // Output error for debugging
+ Console.WriteLine($"Submission failed: {result.ErrorMessage}");
+ }
+ result.Success.Should().BeTrue($"Error: {result.ErrorMessage}");
+ result.JobId.Should().NotBeNullOrEmpty();
+ }
+
+ [Fact]
+ public async Task SubmitJobAsync_AssignsJobId()
+ {
+ // Arrange
+ var jobGraph = CreateValidJobGraph();
+
+ // Act
+ var result = await _dispatcher.SubmitJobAsync(jobGraph);
+
+ // Assert
+ jobGraph.JobId.Should().NotBeNullOrEmpty();
+ jobGraph.JobId.Should().Be(result.JobId);
+ }
+
+ [Fact]
+ public async Task GetJobStatusAsync_ForExistingJob_ReturnsStatus()
+ {
+ // Arrange
+ var jobGraph = CreateValidJobGraph();
+ var submitResult = await _dispatcher.SubmitJobAsync(jobGraph);
+
+ // Act
+ var status = await _dispatcher.GetJobStatusAsync(submitResult.JobId);
+
+ // Assert
+ status.Should().NotBeNull();
+ status.JobId.Should().Be(submitResult.JobId);
+ status.JobName.Should().Be(jobGraph.JobName);
+ }
+
+ [Fact]
+ public async Task GetJobStatusAsync_ForNonExistentJob_ReturnsNull()
+ {
+ // Act
+ var status = await _dispatcher.GetJobStatusAsync("non-existent-job");
+
+ // Assert
+ status.Should().BeNull();
+ }
+
+ [Fact]
+ public async Task ListJobsAsync_ReturnsAllJobs()
+ {
+ // Arrange
+ var jobGraph1 = CreateValidJobGraph();
+ var jobGraph2 = CreateValidJobGraph();
+ await _dispatcher.SubmitJobAsync(jobGraph1);
+ await _dispatcher.SubmitJobAsync(jobGraph2);
+
+ // Act
+ var jobs = await _dispatcher.ListJobsAsync();
+
+ // Assert
+ jobs.Should().HaveCount(2);
+ }
+
+ [Fact]
+ public async Task CancelJobAsync_WithExistingJob_UpdatesJobState()
+ {
+ // Arrange
+ var jobGraph = CreateValidJobGraph();
+ var submitResult = await _dispatcher.SubmitJobAsync(jobGraph);
+
+ // Give the job a moment to start
+ await Task.Delay(100);
+
+ // Act
+ await _dispatcher.CancelJobAsync(submitResult.JobId);
+
+ // Assert
+ var status = await _dispatcher.GetJobStatusAsync(submitResult.JobId);
+ // Job might be Canceling, Canceled, or Failed (if it failed before cancellation completed)
+ status.State.Should().BeOneOf(
+ JobExecutionState.Canceling,
+ JobExecutionState.Canceled,
+ JobExecutionState.Failed);
+ }
+
+ [Fact]
+ public async Task CancelJobAsync_WithNonExistentJob_ThrowsArgumentException()
+ {
+ // Act
+ var act = async () => await _dispatcher.CancelJobAsync("non-existent-job");
+
+ // Assert
+ await act.Should().ThrowAsync()
+ .WithParameterName("jobId");
+ }
+
+ [Fact]
+ public async Task SubmitJobAsync_WithInvalidJobGraph_ReturnsFailure()
+ {
+ // Arrange - Create invalid job graph with null job name
+ var jobGraph = new JobGraph
+ {
+ JobName = null!, // Invalid
+ Vertices = new List(),
+ Edges = new List(),
+ Configuration = new Dictionary()
+ };
+
+ // Act
+ var result = await _dispatcher.SubmitJobAsync(jobGraph);
+
+ // Assert
+ result.Success.Should().BeFalse();
+ result.ErrorMessage.Should().NotBeNullOrEmpty();
+ }
+
+ [Fact]
+ public async Task SubmitJobAsync_WithEmptyVertices_ReturnsFailure()
+ {
+ // Arrange
+ var jobGraph = new JobGraph
+ {
+ JobName = "Test Job",
+ Vertices = new List(), // Empty - invalid
+ Edges = new List(),
+ Configuration = new Dictionary()
+ };
+
+ // Act
+ var result = await _dispatcher.SubmitJobAsync(jobGraph);
+
+ // Assert
+ result.Success.Should().BeFalse();
+ result.ErrorMessage.Should().NotBeNullOrEmpty();
+ }
+
+ [Fact]
+ public async Task SubmitJobAsync_WithZeroMaxParallelism_ReturnsFailure()
+ {
+ // Arrange
+ var jobGraph = new JobGraph
+ {
+ JobName = "Test Job",
+ MaxParallelism = 0, // Invalid
+ Vertices = new List
+ {
+ new JobVertex { Name = "source", Parallelism = 1, OperatorType = OperatorType.Source }
+ },
+ Edges = new List(),
+ Configuration = new Dictionary()
+ };
+
+ // Act
+ var result = await _dispatcher.SubmitJobAsync(jobGraph);
+
+ // Assert
+ result.Success.Should().BeFalse();
+ result.ErrorMessage.Should().Contain("parallelism");
+ }
+
+ [Fact]
+ public async Task ListJobsAsync_WithMultipleJobs_ReturnsAllJobs()
+ {
+ // Arrange
+ var jobGraph1 = CreateValidJobGraph();
+ var jobGraph2 = CreateValidJobGraph();
+ var jobGraph3 = CreateValidJobGraph();
+
+ await _dispatcher.SubmitJobAsync(jobGraph1);
+ await _dispatcher.SubmitJobAsync(jobGraph2);
+ await _dispatcher.SubmitJobAsync(jobGraph3);
+
+ // Act
+ var jobs = await _dispatcher.ListJobsAsync();
+
+ // Assert
+ jobs.Should().HaveCount(3);
+ }
+
+ [Fact]
+ public async Task GetJobStatusAsync_MultipleTimesForSameJob_ReturnsSameJobId()
+ {
+ // Arrange
+ var jobGraph = CreateValidJobGraph();
+ var submitResult = await _dispatcher.SubmitJobAsync(jobGraph);
+
+ // Act
+ var status1 = await _dispatcher.GetJobStatusAsync(submitResult.JobId);
+ var status2 = await _dispatcher.GetJobStatusAsync(submitResult.JobId);
+
+ // Assert
+ status1!.JobId.Should().Be(status2!.JobId);
+ }
+
+ private static JobGraph CreateValidJobGraph()
+ {
+ var sourceVertex = new JobVertex
+ {
+ Name = "source",
+ Parallelism = 2,
+ OperatorType = OperatorType.Source
+ };
+
+ var mapVertex = new JobVertex
+ {
+ Name = "map",
+ Parallelism = 2,
+ OperatorType = OperatorType.Map
+ };
+
+ return new JobGraph
+ {
+ JobName = $"Test Job {Guid.NewGuid()}",
+ MaxParallelism = 128,
+ Vertices = new List
+ {
+ sourceVertex,
+ mapVertex
+ },
+ Edges = new List
+ {
+ new JobEdge
+ {
+ SourceVertexId = sourceVertex.VertexId,
+ TargetVertexId = mapVertex.VertexId,
+ PartitioningStrategy = PartitioningStrategy.Forward
+ }
+ },
+ Configuration = new Dictionary()
+ };
+ }
+}
diff --git a/FlinkDotNet/FlinkDotNet.JobManager.Tests/IntegrationScenarioTests.cs b/FlinkDotNet/FlinkDotNet.JobManager.Tests/IntegrationScenarioTests.cs
new file mode 100644
index 00000000..35c7eb5e
--- /dev/null
+++ b/FlinkDotNet/FlinkDotNet.JobManager.Tests/IntegrationScenarioTests.cs
@@ -0,0 +1,229 @@
+// Copyright 2025 FlinkDotNet
+// Licensed under the Apache License, Version 2.0.
+// See LICENSE file in the project root for full license information.
+
+using FlinkDotNet.JobManager.Implementation;
+using FlinkDotNet.JobManager.Models;
+using FluentAssertions;
+using Microsoft.Extensions.Logging;
+using Moq;
+using Temporalio.Client;
+
+namespace FlinkDotNet.JobManager.Tests;
+
+public class IntegrationScenarioTests
+{
+ private readonly Mock> _mockDispatcherLogger;
+ private readonly Mock> _mockResourceLogger;
+ private readonly Mock _mockTemporalClient;
+ private readonly Mock _mockLoggerFactory;
+ private readonly ResourceManager _resourceManager;
+ private readonly Dispatcher _dispatcher;
+
+ public IntegrationScenarioTests()
+ {
+ _mockDispatcherLogger = new Mock>();
+ _mockResourceLogger = new Mock>();
+ _mockTemporalClient = new Mock();
+ _mockLoggerFactory = new Mock();
+
+ _mockLoggerFactory
+ .Setup(lf => lf.CreateLogger(It.IsAny()))
+ .Returns(new Mock().Object);
+
+ _resourceManager = new ResourceManager(_mockResourceLogger.Object);
+ _dispatcher = new Dispatcher(
+ _resourceManager,
+ _mockTemporalClient.Object,
+ _mockLoggerFactory.Object);
+ }
+
+ [Fact]
+ public async Task CompleteJobLifecycle_WithTaskManagers_ExecutesSuccessfully()
+ {
+ // Arrange - Setup infrastructure
+ await _resourceManager.RegisterTaskManagerAsync("tm-1", 4);
+ await _resourceManager.RegisterTaskManagerAsync("tm-2", 4);
+
+ var jobGraph = CreateValidJobGraph();
+
+ // Act - Submit job
+ var submitResult = await _dispatcher.SubmitJobAsync(jobGraph);
+
+ // Assert - Job submitted
+ submitResult.Should().NotBeNull();
+ submitResult.Success.Should().BeTrue();
+ submitResult.JobId.Should().NotBeNullOrEmpty();
+
+ // Wait for job to start processing (minimal delay)
+ await Task.Delay(50);
+
+ // Act - Get status
+ var status = await _dispatcher.GetJobStatusAsync(submitResult.JobId);
+
+ // Assert - Job is tracked
+ status.Should().NotBeNull();
+ status!.JobId.Should().Be(submitResult.JobId);
+
+ // Act - List jobs
+ var allJobs = await _dispatcher.ListJobsAsync();
+
+ // Assert - Job appears in list
+ allJobs.Should().Contain(j => j.JobId == submitResult.JobId);
+ }
+
+ [Fact]
+ public async Task MultipleJobsScenario_ManagesResourcesCorrectly()
+ {
+ // Arrange
+ await _resourceManager.RegisterTaskManagerAsync("tm-1", 10);
+
+ var job1 = CreateValidJobGraph();
+ var job2 = CreateValidJobGraph();
+ var job3 = CreateValidJobGraph();
+
+ // Act - Submit multiple jobs
+ var result1 = await _dispatcher.SubmitJobAsync(job1);
+ var result2 = await _dispatcher.SubmitJobAsync(job2);
+ var result3 = await _dispatcher.SubmitJobAsync(job3);
+
+ await Task.Delay(50);
+
+ // Assert - All jobs submitted
+ result1.Success.Should().BeTrue();
+ result2.Success.Should().BeTrue();
+ result3.Success.Should().BeTrue();
+
+ // Act - List all jobs
+ var allJobs = await _dispatcher.ListJobsAsync();
+
+ // Assert - All jobs tracked
+ allJobs.Should().HaveCountGreaterThanOrEqualTo(3);
+ }
+
+ [Fact]
+ public async Task JobCancellation_ReleasesResources()
+ {
+ // Arrange
+ await _resourceManager.RegisterTaskManagerAsync("tm-1", 4);
+ var jobGraph = CreateValidJobGraph();
+ var submitResult = await _dispatcher.SubmitJobAsync(jobGraph);
+
+ await Task.Delay(50);
+ var availableBefore = _resourceManager.GetAvailableSlots().Count();
+
+ // Act - Cancel job
+ await _dispatcher.CancelJobAsync(submitResult.JobId);
+ await Task.Delay(50);
+
+ // Assert - Status should reflect cancellation
+ var status = await _dispatcher.GetJobStatusAsync(submitResult.JobId);
+ status!.State.Should().BeOneOf(
+ JobExecutionState.Canceling,
+ JobExecutionState.Canceled,
+ JobExecutionState.Failed);
+ }
+
+ [Fact]
+ public async Task ResourceAllocation_AcrossMultipleTaskManagers()
+ {
+ // Arrange
+ await _resourceManager.RegisterTaskManagerAsync("tm-1", 2);
+ await _resourceManager.RegisterTaskManagerAsync("tm-2", 2);
+ await _resourceManager.RegisterTaskManagerAsync("tm-3", 2);
+
+ // Act - Allocate more slots than any single TaskManager has
+ var slots = await _resourceManager.AllocateSlotsAsync("job-1", 5);
+
+ // Assert - Slots distributed across TaskManagers
+ slots.Should().HaveCount(5);
+ var taskManagersUsed = slots.Select(s => s.TaskManagerId).Distinct().Count();
+ taskManagersUsed.Should().BeGreaterThan(1);
+ }
+
+ [Fact]
+ public async Task TaskManagerUnregistration_RemovesSlots()
+ {
+ // Arrange
+ await _resourceManager.RegisterTaskManagerAsync("tm-1", 4);
+ var initialSlots = _resourceManager.GetAllSlots().Count();
+
+ // Act
+ await _resourceManager.UnregisterTaskManagerAsync("tm-1");
+ var afterSlots = _resourceManager.GetAllSlots().Count();
+
+ // Assert
+ afterSlots.Should().BeLessThan(initialSlots);
+ }
+
+ [Fact]
+ public void GetAllSlots_WithMultipleTaskManagers_ReturnsAllSlots()
+ {
+ // Arrange
+ _resourceManager.RegisterTaskManager("tm-1", 3);
+ _resourceManager.RegisterTaskManager("tm-2", 5);
+
+ // Act
+ var allSlots = _resourceManager.GetAllSlots();
+
+ // Assert
+ allSlots.Should().HaveCount(8);
+ }
+
+ [Fact]
+ public async Task JobSubmission_WithInvalidGraph_ReturnsFailure()
+ {
+ // Arrange
+ var invalidGraph = new JobGraph
+ {
+ JobName = "", // Invalid - empty name
+ Vertices = new List
+ {
+ new JobVertex { Name = "test", Parallelism = 1, OperatorType = OperatorType.Source }
+ },
+ Edges = new List(),
+ Configuration = new Dictionary()
+ };
+
+ // Act
+ var result = await _dispatcher.SubmitJobAsync(invalidGraph);
+
+ // Assert
+ result.Success.Should().BeFalse();
+ result.ErrorMessage.Should().Contain("name");
+ }
+
+ private static JobGraph CreateValidJobGraph()
+ {
+ var sourceVertex = new JobVertex
+ {
+ Name = $"source-{Guid.NewGuid().ToString().Substring(0, 8)}",
+ Parallelism = 2,
+ OperatorType = OperatorType.Source
+ };
+
+ var mapVertex = new JobVertex
+ {
+ Name = $"map-{Guid.NewGuid().ToString().Substring(0, 8)}",
+ Parallelism = 2,
+ OperatorType = OperatorType.Map
+ };
+
+ return new JobGraph
+ {
+ JobName = $"Integration Test Job {Guid.NewGuid()}",
+ MaxParallelism = 128,
+ Vertices = new List { sourceVertex, mapVertex },
+ Edges = new List
+ {
+ new JobEdge
+ {
+ SourceVertexId = sourceVertex.VertexId,
+ TargetVertexId = mapVertex.VertexId,
+ PartitioningStrategy = PartitioningStrategy.Forward
+ }
+ },
+ Configuration = new Dictionary()
+ };
+ }
+}
diff --git a/FlinkDotNet/FlinkDotNet.JobManager.Tests/JobMasterTests.cs b/FlinkDotNet/FlinkDotNet.JobManager.Tests/JobMasterTests.cs
new file mode 100644
index 00000000..5627e80a
--- /dev/null
+++ b/FlinkDotNet/FlinkDotNet.JobManager.Tests/JobMasterTests.cs
@@ -0,0 +1,382 @@
+// Copyright 2025 FlinkDotNet
+// Licensed under the Apache License, Version 2.0.
+// See LICENSE file in the project root for full license information.
+
+using FlinkDotNet.JobManager.Implementation;
+using FlinkDotNet.JobManager.Interfaces;
+using FlinkDotNet.JobManager.Models;
+using FluentAssertions;
+using Microsoft.Extensions.Logging;
+using Moq;
+using Temporalio.Client;
+
+namespace FlinkDotNet.JobManager.Tests;
+
+public class JobMasterTests
+{
+ private readonly Mock _mockResourceManager;
+ private readonly Mock _mockTemporalClient;
+ private readonly Mock> _mockLogger;
+ private readonly string _jobId;
+ private readonly JobGraph _jobGraph;
+
+ public JobMasterTests()
+ {
+ _mockResourceManager = new Mock();
+ _mockTemporalClient = new Mock();
+ _mockLogger = new Mock>();
+ _jobId = "test-job-1";
+
+ _jobGraph = new JobGraph
+ {
+ JobName = "Test Job",
+ Vertices = new List
+ {
+ new JobVertex
+ {
+ Name = "source",
+ Parallelism = 2,
+ OperatorType = OperatorType.Source
+ },
+ new JobVertex
+ {
+ Name = "map",
+ Parallelism = 2,
+ OperatorType = OperatorType.Map
+ }
+ },
+ Edges = new List
+ {
+ new JobEdge
+ {
+ SourceVertexId = "source",
+ TargetVertexId = "map",
+ PartitioningStrategy = PartitioningStrategy.Forward
+ }
+ },
+ Configuration = new Dictionary()
+ };
+ }
+
+ [Fact]
+ public void Constructor_WithNullJobId_ThrowsArgumentNullException()
+ {
+ // Act & Assert
+ var act = () => new JobMaster(
+ null!,
+ _jobGraph,
+ _mockResourceManager.Object,
+ _mockTemporalClient.Object,
+ _mockLogger.Object);
+
+ act.Should().Throw().WithParameterName("jobId");
+ }
+
+ [Fact]
+ public void Constructor_WithNullJobGraph_ThrowsArgumentNullException()
+ {
+ // Act & Assert
+ var act = () => new JobMaster(
+ _jobId,
+ null!,
+ _mockResourceManager.Object,
+ _mockTemporalClient.Object,
+ _mockLogger.Object);
+
+ act.Should().Throw().WithParameterName("jobGraph");
+ }
+
+ [Fact]
+ public void Constructor_WithNullResourceManager_ThrowsArgumentNullException()
+ {
+ // Act & Assert
+ var act = () => new JobMaster(
+ _jobId,
+ _jobGraph,
+ null!,
+ _mockTemporalClient.Object,
+ _mockLogger.Object);
+
+ act.Should().Throw().WithParameterName("resourceManager");
+ }
+
+ [Fact]
+ public void Constructor_WithNullTemporalClient_ThrowsArgumentNullException()
+ {
+ // Act & Assert
+ var act = () => new JobMaster(
+ _jobId,
+ _jobGraph,
+ _mockResourceManager.Object,
+ null!,
+ _mockLogger.Object);
+
+ act.Should().Throw();
+ }
+
+ [Fact]
+ public void Constructor_WithNullLogger_ThrowsArgumentNullException()
+ {
+ // Act & Assert
+ var act = () => new JobMaster(
+ _jobId,
+ _jobGraph,
+ _mockResourceManager.Object,
+ _mockTemporalClient.Object,
+ null!);
+
+ act.Should().Throw().WithParameterName("logger");
+ }
+
+ [Fact]
+ public void JobId_ReturnsCorrectJobId()
+ {
+ // Arrange
+ var jobMaster = new JobMaster(
+ _jobId,
+ _jobGraph,
+ _mockResourceManager.Object,
+ _mockTemporalClient.Object,
+ _mockLogger.Object);
+
+ // Act & Assert
+ jobMaster.JobId.Should().Be(_jobId);
+ }
+
+ [Fact]
+ public async Task StartJobAsync_CreatesExecutionGraph()
+ {
+ // Arrange
+ var slots = Enumerable.Range(0, 4).Select(i => new TaskSlot
+ {
+ SlotId = $"slot-{i}",
+ TaskManagerId = "tm-1",
+ IsAllocated = true
+ }).ToList();
+
+ _mockResourceManager
+ .Setup(rm => rm.AllocateSlotsAsync(_jobId, 4, It.IsAny()))
+ .ReturnsAsync(slots);
+
+ var jobMaster = new JobMaster(
+ _jobId,
+ _jobGraph,
+ _mockResourceManager.Object,
+ _mockTemporalClient.Object,
+ _mockLogger.Object);
+
+ // Act
+ await jobMaster.StartJobAsync();
+
+ // Assert
+ var executionGraph = await jobMaster.GetExecutionGraphAsync();
+ executionGraph.Should().NotBeNull();
+ executionGraph.JobId.Should().Be(_jobId);
+ executionGraph.ExecutionVertices.Should().HaveCount(4); // 2 + 2 parallelism
+ }
+
+ [Fact]
+ public async Task StartJobAsync_AllocatesRequiredSlots()
+ {
+ // Arrange
+ var slots = Enumerable.Range(0, 4).Select(i => new TaskSlot
+ {
+ SlotId = $"slot-{i}",
+ TaskManagerId = "tm-1",
+ IsAllocated = true
+ }).ToList();
+
+ _mockResourceManager
+ .Setup(rm => rm.AllocateSlotsAsync(_jobId, 4, It.IsAny()))
+ .ReturnsAsync(slots);
+
+ var jobMaster = new JobMaster(
+ _jobId,
+ _jobGraph,
+ _mockResourceManager.Object,
+ _mockTemporalClient.Object,
+ _mockLogger.Object);
+
+ // Act
+ await jobMaster.StartJobAsync();
+
+ // Assert
+ _mockResourceManager.Verify(
+ rm => rm.AllocateSlotsAsync(_jobId, 4, It.IsAny()),
+ Times.Once);
+ }
+
+ [Fact]
+ public async Task StartJobAsync_WithInsufficientResources_ThrowsInvalidOperationException()
+ {
+ // Arrange
+ _mockResourceManager
+ .Setup(rm => rm.AllocateSlotsAsync(_jobId, 4, It.IsAny()))
+ .ReturnsAsync(new List { new TaskSlot { SlotId = "slot-1", TaskManagerId = "tm-1" } });
+
+ var jobMaster = new JobMaster(
+ _jobId,
+ _jobGraph,
+ _mockResourceManager.Object,
+ _mockTemporalClient.Object,
+ _mockLogger.Object);
+
+ // Act
+ var act = async () => await jobMaster.StartJobAsync();
+
+ // Assert
+ await act.Should().ThrowAsync()
+ .Where(ex => ex.Message.Contains("Failed to start job") || ex.InnerException != null);
+ }
+
+ [Fact]
+ public async Task CancelJobAsync_ReleasesAllocatedResources()
+ {
+ // Arrange
+ var slots = Enumerable.Range(0, 4).Select(i => new TaskSlot
+ {
+ SlotId = $"slot-{i}",
+ TaskManagerId = "tm-1",
+ IsAllocated = true
+ }).ToList();
+
+ _mockResourceManager
+ .Setup(rm => rm.AllocateSlotsAsync(_jobId, 4, It.IsAny()))
+ .ReturnsAsync(slots);
+
+ var jobMaster = new JobMaster(
+ _jobId,
+ _jobGraph,
+ _mockResourceManager.Object,
+ _mockTemporalClient.Object,
+ _mockLogger.Object);
+
+ await jobMaster.StartJobAsync();
+
+ // Act
+ await jobMaster.CancelJobAsync();
+
+ // Assert
+ _mockResourceManager.Verify(
+ rm => rm.ReleaseSlotAsync(It.IsAny(), It.IsAny()),
+ Times.AtLeastOnce);
+ }
+
+ [Fact]
+ public async Task GetExecutionGraphAsync_BeforeStart_ThrowsInvalidOperationException()
+ {
+ // Arrange
+ var jobMaster = new JobMaster(
+ _jobId,
+ _jobGraph,
+ _mockResourceManager.Object,
+ _mockTemporalClient.Object,
+ _mockLogger.Object);
+
+ // Act
+ var act = async () => await jobMaster.GetExecutionGraphAsync();
+
+ // Assert
+ await act.Should().ThrowAsync()
+ .WithMessage("*ExecutionGraph not yet created*");
+ }
+
+ [Fact]
+ public async Task UpdateTaskStatusAsync_UpdatesVertexState()
+ {
+ // Arrange
+ var slots = Enumerable.Range(0, 4).Select(i => new TaskSlot
+ {
+ SlotId = $"slot-{i}",
+ TaskManagerId = "tm-1",
+ IsAllocated = true
+ }).ToList();
+
+ _mockResourceManager
+ .Setup(rm => rm.AllocateSlotsAsync(_jobId, 4, It.IsAny()))
+ .ReturnsAsync(slots);
+
+ var jobMaster = new JobMaster(
+ _jobId,
+ _jobGraph,
+ _mockResourceManager.Object,
+ _mockTemporalClient.Object,
+ _mockLogger.Object);
+
+ await jobMaster.StartJobAsync();
+ var executionGraph = await jobMaster.GetExecutionGraphAsync();
+ var vertexId = executionGraph.ExecutionVertices.First().Id;
+
+ // Act
+ await jobMaster.UpdateTaskStatusAsync(vertexId, ExecutionState.Running);
+
+ // Assert
+ var updatedGraph = await jobMaster.GetExecutionGraphAsync();
+ var updatedVertex = updatedGraph.ExecutionVertices.First(v => v.Id == vertexId);
+ updatedVertex.State.Should().Be(ExecutionState.Running);
+ }
+
+ [Fact]
+ public async Task TriggerCheckpointAsync_LogsCheckpointRequest()
+ {
+ // Arrange
+ var slots = Enumerable.Range(0, 4).Select(i => new TaskSlot
+ {
+ SlotId = $"slot-{i}",
+ TaskManagerId = "tm-1",
+ IsAllocated = true
+ }).ToList();
+
+ _mockResourceManager
+ .Setup(rm => rm.AllocateSlotsAsync(_jobId, 4, It.IsAny()))
+ .ReturnsAsync(slots);
+
+ var jobMaster = new JobMaster(
+ _jobId,
+ _jobGraph,
+ _mockResourceManager.Object,
+ _mockTemporalClient.Object,
+ _mockLogger.Object);
+
+ await jobMaster.StartJobAsync();
+
+ // Act
+ var act = async () => await jobMaster.TriggerCheckpointAsync(12345);
+
+ // Assert
+ await act.Should().NotThrowAsync();
+ }
+
+ [Fact]
+ public void Constructor_WithAllParameters_CreatesInstance()
+ {
+ // Act
+ var jobMaster = new JobMaster(
+ _jobId,
+ _jobGraph,
+ _mockResourceManager.Object,
+ _mockTemporalClient.Object,
+ _mockLogger.Object);
+
+ // Assert
+ jobMaster.Should().NotBeNull();
+ }
+
+ [Fact]
+ public async Task CancelJobAsync_BeforeStart_DoesNotThrow()
+ {
+ // Arrange
+ var jobMaster = new JobMaster(
+ _jobId,
+ _jobGraph,
+ _mockResourceManager.Object,
+ _mockTemporalClient.Object,
+ _mockLogger.Object);
+
+ // Act
+ var act = async () => await jobMaster.CancelJobAsync();
+
+ // Assert
+ await act.Should().NotThrowAsync();
+ }
+}
diff --git a/FlinkDotNet/FlinkDotNet.JobManager.Tests/JobsControllerTests.cs b/FlinkDotNet/FlinkDotNet.JobManager.Tests/JobsControllerTests.cs
new file mode 100644
index 00000000..19fd38d1
--- /dev/null
+++ b/FlinkDotNet/FlinkDotNet.JobManager.Tests/JobsControllerTests.cs
@@ -0,0 +1,315 @@
+// Copyright 2025 FlinkDotNet
+// Licensed under the Apache License, Version 2.0.
+// See LICENSE file in the project root for full license information.
+
+using FlinkDotNet.JobManager.Controllers;
+using FlinkDotNet.JobManager.Interfaces;
+using FlinkDotNet.JobManager.Models;
+using FlinkDotNet.JobManager.Models.Requests;
+using FlinkDotNet.JobManager.Models.Responses;
+using FluentAssertions;
+using Microsoft.AspNetCore.Mvc;
+using Microsoft.Extensions.Logging;
+using Moq;
+
+namespace FlinkDotNet.JobManager.Tests;
+
+public class JobsControllerTests
+{
+ private readonly Mock _mockDispatcher;
+ private readonly Mock> _mockLogger;
+ private readonly JobsController _controller;
+
+ public JobsControllerTests()
+ {
+ _mockDispatcher = new Mock();
+ _mockLogger = new Mock>();
+ _controller = new JobsController(_mockDispatcher.Object, _mockLogger.Object);
+ }
+
+ [Fact]
+ public void Constructor_WithNullDispatcher_ThrowsArgumentNullException()
+ {
+ // Act & Assert
+ var act = () => new JobsController(null!, _mockLogger.Object);
+ act.Should().Throw().WithParameterName("dispatcher");
+ }
+
+ [Fact]
+ public void Constructor_WithNullLogger_ThrowsArgumentNullException()
+ {
+ // Act & Assert
+ var act = () => new JobsController(_mockDispatcher.Object, null!);
+ act.Should().Throw().WithParameterName("logger");
+ }
+
+ [Fact]
+ public async Task SubmitJob_WithValidRequest_ReturnsOkResult()
+ {
+ // Arrange
+ var request = CreateValidSubmitJobRequest();
+ var submitResult = new JobSubmissionResult
+ {
+ JobId = "job-123",
+ Success = true
+ };
+
+ _mockDispatcher
+ .Setup(d => d.SubmitJobAsync(It.IsAny(), It.IsAny()))
+ .ReturnsAsync(submitResult);
+
+ // Act
+ var result = await _controller.SubmitJob(request);
+
+ // Assert
+ result.Should().BeOfType();
+ var okResult = result as OkObjectResult;
+ var response = okResult!.Value as SubmitJobResponse;
+ response!.JobId.Should().Be("job-123");
+ response.State.Should().Be(JobExecutionState.Created);
+ }
+
+ [Fact]
+ public async Task SubmitJob_WithFailedSubmission_ReturnsBadRequest()
+ {
+ // Arrange
+ var request = CreateValidSubmitJobRequest();
+ var submitResult = new JobSubmissionResult
+ {
+ JobId = "job-123",
+ Success = false,
+ ErrorMessage = "Invalid job graph"
+ };
+
+ _mockDispatcher
+ .Setup(d => d.SubmitJobAsync(It.IsAny(), It.IsAny()))
+ .ReturnsAsync(submitResult);
+
+ // Act
+ var result = await _controller.SubmitJob(request);
+
+ // Assert
+ result.Should().BeOfType();
+ }
+
+ [Fact]
+ public async Task GetJobStatus_WithExistingJob_ReturnsOkResult()
+ {
+ // Arrange
+ var jobId = "job-123";
+ var jobStatus = new JobStatus
+ {
+ JobId = jobId,
+ JobName = "Test Job",
+ State = JobExecutionState.Running
+ };
+
+ _mockDispatcher
+ .Setup(d => d.GetJobStatusAsync(jobId, It.IsAny()))
+ .ReturnsAsync(jobStatus);
+
+ // Act
+ var result = await _controller.GetJobStatus(jobId);
+
+ // Assert
+ result.Should().BeOfType();
+ var okResult = result as OkObjectResult;
+ var response = okResult!.Value as JobStatusResponse;
+ response!.JobId.Should().Be(jobId);
+ response.State.Should().Be(JobExecutionState.Running);
+ }
+
+ [Fact]
+ public async Task GetJobStatus_WithNonExistentJob_ReturnsNotFound()
+ {
+ // Arrange
+ var jobId = "non-existent-job";
+
+ _mockDispatcher
+ .Setup(d => d.GetJobStatusAsync(jobId, It.IsAny()))
+ .ReturnsAsync((JobStatus?)null);
+
+ // Act
+ var result = await _controller.GetJobStatus(jobId);
+
+ // Assert
+ result.Should().BeOfType();
+ }
+
+ [Fact]
+ public async Task ListJobs_ReturnsOkResultWithJobList()
+ {
+ // Arrange
+ var jobs = new List
+ {
+ new JobStatus { JobId = "job-1", JobName = "Job 1", State = JobExecutionState.Running },
+ new JobStatus { JobId = "job-2", JobName = "Job 2", State = JobExecutionState.Finished }
+ };
+
+ _mockDispatcher
+ .Setup(d => d.ListJobsAsync(It.IsAny()))
+ .ReturnsAsync(jobs);
+
+ // Act
+ var result = await _controller.ListJobs();
+
+ // Assert
+ result.Should().BeOfType();
+ var okResult = result as OkObjectResult;
+ var response = okResult!.Value as JobListResponse;
+ response!.Jobs.Should().HaveCount(2);
+ }
+
+ [Fact]
+ public async Task CancelJob_WithExistingJob_ReturnsOkResult()
+ {
+ // Arrange
+ var jobId = "job-123";
+
+ _mockDispatcher
+ .Setup(d => d.CancelJobAsync(jobId, It.IsAny()))
+ .Returns(Task.CompletedTask);
+
+ // Act
+ var result = await _controller.CancelJob(jobId);
+
+ // Assert
+ result.Should().BeOfType();
+ }
+
+ [Fact]
+ public async Task CancelJob_WithNonExistentJob_ReturnsNotFound()
+ {
+ // Arrange
+ var jobId = "non-existent-job";
+
+ _mockDispatcher
+ .Setup(d => d.CancelJobAsync(jobId, It.IsAny()))
+ .ThrowsAsync(new ArgumentException($"Job {jobId} not found", nameof(jobId)));
+
+ // Act
+ var result = await _controller.CancelJob(jobId);
+
+ // Assert
+ result.Should().BeOfType();
+ }
+
+ [Fact]
+ public async Task CancelJob_WithException_ReturnsBadRequest()
+ {
+ // Arrange
+ var jobId = "job-123";
+
+ _mockDispatcher
+ .Setup(d => d.CancelJobAsync(jobId, It.IsAny()))
+ .ThrowsAsync(new InvalidOperationException("Cannot cancel job in current state"));
+
+ // Act
+ var result = await _controller.CancelJob(jobId);
+
+ // Assert
+ result.Should().BeOfType();
+ }
+
+ [Fact]
+ public async Task ListJobs_WithStateFilter_FiltersJobs()
+ {
+ // Arrange
+ var jobs = new List
+ {
+ new JobStatus { JobId = "job-1", JobName = "Job 1", State = JobExecutionState.Running },
+ new JobStatus { JobId = "job-2", JobName = "Job 2", State = JobExecutionState.Finished },
+ new JobStatus { JobId = "job-3", JobName = "Job 3", State = JobExecutionState.Running }
+ };
+
+ _mockDispatcher
+ .Setup(d => d.ListJobsAsync(It.IsAny()))
+ .ReturnsAsync(jobs);
+
+ // Act
+ var result = await _controller.ListJobs("Running");
+
+ // Assert
+ result.Should().BeOfType();
+ var okResult = result as OkObjectResult;
+ var response = okResult!.Value as JobListResponse;
+ response!.Jobs.Should().HaveCount(2);
+ response.Jobs.Should().OnlyContain(j => j.State == JobExecutionState.Running);
+ }
+
+ [Fact]
+ public async Task ListJobs_WithInvalidStateFilter_ReturnsAllJobs()
+ {
+ // Arrange
+ var jobs = new List
+ {
+ new JobStatus { JobId = "job-1", JobName = "Job 1", State = JobExecutionState.Running },
+ new JobStatus { JobId = "job-2", JobName = "Job 2", State = JobExecutionState.Finished }
+ };
+
+ _mockDispatcher
+ .Setup(d => d.ListJobsAsync(It.IsAny()))
+ .ReturnsAsync(jobs);
+
+ // Act
+ var result = await _controller.ListJobs("InvalidState");
+
+ // Assert
+ result.Should().BeOfType();
+ var okResult = result as OkObjectResult;
+ var response = okResult!.Value as JobListResponse;
+ response!.Jobs.Should().HaveCount(2);
+ }
+
+ [Fact]
+ public async Task SubmitJob_WithException_ReturnsInternalServerError()
+ {
+ // Arrange
+ var request = CreateValidSubmitJobRequest();
+
+ _mockDispatcher
+ .Setup(d => d.SubmitJobAsync(It.IsAny(), It.IsAny()))
+ .ThrowsAsync(new Exception("Unexpected error"));
+
+ // Act
+ var result = await _controller.SubmitJob(request);
+
+ // Assert
+ result.Should().BeOfType();
+ var objectResult = result as ObjectResult;
+ objectResult!.StatusCode.Should().Be(500);
+ }
+
+ private static SubmitJobRequest CreateValidSubmitJobRequest()
+ {
+ return new SubmitJobRequest
+ {
+ JobName = "Test Job",
+ MaxParallelism = 128,
+ Vertices = new List
+ {
+ new JobVertexRequest
+ {
+ OperatorName = "source",
+ Parallelism = 2,
+ OperatorType = "Source"
+ },
+ new JobVertexRequest
+ {
+ OperatorName = "map",
+ Parallelism = 2,
+ OperatorType = "Map"
+ }
+ },
+ Edges = new List
+ {
+ new JobEdgeRequest
+ {
+ SourceVertexIndex = 0,
+ TargetVertexIndex = 1,
+ Strategy = "Forward"
+ }
+ }
+ };
+ }
+}
diff --git a/FlinkDotNet/FlinkDotNet.JobManager.Tests/ModelTests.cs b/FlinkDotNet/FlinkDotNet.JobManager.Tests/ModelTests.cs
new file mode 100644
index 00000000..c8a7bb23
--- /dev/null
+++ b/FlinkDotNet/FlinkDotNet.JobManager.Tests/ModelTests.cs
@@ -0,0 +1,162 @@
+// Copyright 2025 FlinkDotNet
+// Licensed under the Apache License, Version 2.0.
+// See LICENSE file in the project root for full license information.
+
+using FlinkDotNet.JobManager.Models;
+using FlinkDotNet.JobManager.Interfaces;
+using FluentAssertions;
+
+namespace FlinkDotNet.JobManager.Tests;
+
+public class ModelTests
+{
+ [Fact]
+ public void JobVertex_DefaultProperties_SetCorrectly()
+ {
+ // Act
+ var vertex = new JobVertex();
+
+ // Assert
+ vertex.VertexId.Should().NotBeNullOrEmpty();
+ vertex.Parallelism.Should().Be(1);
+ }
+
+ [Fact]
+ public void JobVertex_NameProperty_AliasesOperatorName()
+ {
+ // Arrange
+ var vertex = new JobVertex
+ {
+ Name = "TestOperator"
+ };
+
+ // Assert
+ vertex.OperatorName.Should().Be("TestOperator");
+ vertex.Name.Should().Be("TestOperator");
+ }
+
+ [Fact]
+ public void JobEdge_DefaultProperties_SetCorrectly()
+ {
+ // Act
+ var edge = new JobEdge();
+
+ // Assert
+ edge.PartitioningStrategy.Should().Be(PartitioningStrategy.Forward);
+ }
+
+ [Fact]
+ public void TaskSlot_DefaultValues_SetCorrectly()
+ {
+ // Act
+ var slot = new TaskSlot();
+
+ // Assert
+ slot.SlotId.Should().NotBeNullOrEmpty();
+ slot.IsAllocated.Should().BeFalse();
+ slot.SlotNumber.Should().Be(0);
+ }
+
+ [Fact]
+ public void ExecutionState_Values_AreCorrect()
+ {
+ // Assert
+ ExecutionState.Created.Should().Be(ExecutionState.Created);
+ ExecutionState.Scheduled.Should().Be(ExecutionState.Scheduled);
+ ExecutionState.Deploying.Should().Be(ExecutionState.Deploying);
+ ExecutionState.Running.Should().Be(ExecutionState.Running);
+ ExecutionState.Finished.Should().Be(ExecutionState.Finished);
+ ExecutionState.Canceled.Should().Be(ExecutionState.Canceled);
+ ExecutionState.Failed.Should().Be(ExecutionState.Failed);
+ }
+
+ [Fact]
+ public void JobExecutionState_Values_AreCorrect()
+ {
+ // Assert
+ JobExecutionState.Created.Should().Be(JobExecutionState.Created);
+ JobExecutionState.Running.Should().Be(JobExecutionState.Running);
+ JobExecutionState.Finished.Should().Be(JobExecutionState.Finished);
+ JobExecutionState.Failed.Should().Be(JobExecutionState.Failed);
+ JobExecutionState.Canceling.Should().Be(JobExecutionState.Canceling);
+ JobExecutionState.Canceled.Should().Be(JobExecutionState.Canceled);
+ }
+
+ [Fact]
+ public void OperatorType_Values_AreCorrect()
+ {
+ // Assert
+ OperatorType.Source.Should().Be(OperatorType.Source);
+ OperatorType.Map.Should().Be(OperatorType.Map);
+ OperatorType.Filter.Should().Be(OperatorType.Filter);
+ OperatorType.Sink.Should().Be(OperatorType.Sink);
+ }
+
+ [Fact]
+ public void PartitioningStrategy_Values_AreCorrect()
+ {
+ // Assert
+ PartitioningStrategy.Forward.Should().Be(PartitioningStrategy.Forward);
+ PartitioningStrategy.Rebalance.Should().Be(PartitioningStrategy.Rebalance);
+ PartitioningStrategy.Rescale.Should().Be(PartitioningStrategy.Rescale);
+ PartitioningStrategy.Broadcast.Should().Be(PartitioningStrategy.Broadcast);
+ }
+
+ [Fact]
+ public void ExecutionGraph_DefaultProperties_SetCorrectly()
+ {
+ // Act
+ var graph = new ExecutionGraph();
+
+ // Assert
+ graph.JobId.Should().BeEmpty();
+ graph.ExecutionVertices.Should().NotBeNull().And.BeEmpty();
+ graph.ExecutionEdges.Should().NotBeNull().And.BeEmpty();
+ }
+
+ [Fact]
+ public void ExecutionVertex_DefaultProperties_SetCorrectly()
+ {
+ // Act
+ var vertex = new ExecutionVertex();
+
+ // Assert
+ vertex.Id.Should().NotBeNullOrEmpty();
+ vertex.State.Should().Be(ExecutionState.Created);
+ }
+
+ [Fact]
+ public void JobGraph_DefaultProperties_SetCorrectly()
+ {
+ // Act
+ var graph = new JobGraph();
+
+ // Assert
+ graph.JobId.Should().NotBeNullOrEmpty(); // JobId is initialized with a Guid
+ graph.MaxParallelism.Should().Be(128);
+ graph.Vertices.Should().NotBeNull().And.BeEmpty();
+ graph.Edges.Should().NotBeNull().And.BeEmpty();
+ graph.Configuration.Should().NotBeNull().And.BeEmpty();
+ }
+
+ [Fact]
+ public void JobStatus_Properties_CanBeSet()
+ {
+ // Act
+ var status = new JobStatus
+ {
+ JobId = "job-789",
+ JobName = "Test Job",
+ State = JobExecutionState.Running,
+ StartTime = DateTime.UtcNow,
+ EndTime = null
+ };
+
+ // Assert
+ status.JobId.Should().Be("job-789");
+ status.JobName.Should().Be("Test Job");
+ status.State.Should().Be(JobExecutionState.Running);
+ status.StartTime.Should().NotBeNull();
+ status.EndTime.Should().BeNull();
+ }
+}
diff --git a/FlinkDotNet/FlinkDotNet.JobManager.Tests/ResourceManagerExtensionsTests.cs b/FlinkDotNet/FlinkDotNet.JobManager.Tests/ResourceManagerExtensionsTests.cs
new file mode 100644
index 00000000..b3f904d0
--- /dev/null
+++ b/FlinkDotNet/FlinkDotNet.JobManager.Tests/ResourceManagerExtensionsTests.cs
@@ -0,0 +1,83 @@
+// Copyright 2025 FlinkDotNet
+// Licensed under the Apache License, Version 2.0.
+// See LICENSE file in the project root for full license information.
+
+using FlinkDotNet.JobManager.Implementation;
+using FlinkDotNet.JobManager.Interfaces;
+using FlinkDotNet.JobManager.Models;
+using FluentAssertions;
+using Microsoft.Extensions.Logging;
+using Moq;
+
+namespace FlinkDotNet.JobManager.Tests;
+
+public class ResourceManagerExtensionsTests
+{
+ private readonly Mock> _mockLogger;
+ private readonly ResourceManager _resourceManager;
+
+ public ResourceManagerExtensionsTests()
+ {
+ _mockLogger = new Mock>();
+ _resourceManager = new ResourceManager(_mockLogger.Object);
+ }
+
+ [Fact]
+ public void RegisterTaskManager_Synchronously_RegistersTaskManager()
+ {
+ // Arrange
+ var taskManagerId = "tm-sync-1";
+ var numberOfSlots = 4;
+
+ // Act
+ _resourceManager.RegisterTaskManager(taskManagerId, numberOfSlots);
+
+ // Assert
+ var registeredManagers = _resourceManager.GetRegisteredTaskManagers();
+ registeredManagers.Should().Contain(taskManagerId);
+ }
+
+ [Fact]
+ public void UnregisterTaskManager_Synchronously_UnregistersTaskManager()
+ {
+ // Arrange
+ var taskManagerId = "tm-sync-2";
+ _resourceManager.RegisterTaskManager(taskManagerId, 2);
+
+ // Act
+ var result = _resourceManager.UnregisterTaskManager(taskManagerId);
+
+ // Assert
+ result.Should().BeTrue();
+ var registeredManagers = _resourceManager.GetRegisteredTaskManagers();
+ registeredManagers.Should().NotContain(taskManagerId);
+ }
+
+ [Fact]
+ public void GetRegisteredTaskManagers_ReturnsDistinctTaskManagerIds()
+ {
+ // Arrange
+ _resourceManager.RegisterTaskManager("tm-1", 3);
+ _resourceManager.RegisterTaskManager("tm-2", 2);
+ _resourceManager.RegisterTaskManager("tm-3", 5);
+
+ // Act
+ var registeredManagers = _resourceManager.GetRegisteredTaskManagers().ToList();
+
+ // Assert
+ registeredManagers.Should().HaveCount(3);
+ registeredManagers.Should().Contain("tm-1");
+ registeredManagers.Should().Contain("tm-2");
+ registeredManagers.Should().Contain("tm-3");
+ }
+
+ [Fact]
+ public void GetRegisteredTaskManagers_WithNoRegistrations_ReturnsEmpty()
+ {
+ // Act
+ var registeredManagers = _resourceManager.GetRegisteredTaskManagers();
+
+ // Assert
+ registeredManagers.Should().BeEmpty();
+ }
+}
diff --git a/FlinkDotNet/FlinkDotNet.JobManager.Tests/ResourceManagerTests.cs b/FlinkDotNet/FlinkDotNet.JobManager.Tests/ResourceManagerTests.cs
new file mode 100644
index 00000000..4e14bc34
--- /dev/null
+++ b/FlinkDotNet/FlinkDotNet.JobManager.Tests/ResourceManagerTests.cs
@@ -0,0 +1,290 @@
+// Copyright 2025 FlinkDotNet
+// Licensed under the Apache License, Version 2.0.
+// See LICENSE file in the project root for full license information.
+
+using FlinkDotNet.JobManager.Implementation;
+using FlinkDotNet.JobManager.Models;
+using FluentAssertions;
+using Microsoft.Extensions.Logging;
+using Moq;
+
+namespace FlinkDotNet.JobManager.Tests;
+
+public class ResourceManagerTests
+{
+ private readonly Mock> _mockLogger;
+ private readonly ResourceManager _resourceManager;
+
+ public ResourceManagerTests()
+ {
+ _mockLogger = new Mock>();
+ _resourceManager = new ResourceManager(_mockLogger.Object);
+ }
+
+ [Fact]
+ public void Constructor_WithNullLogger_ThrowsArgumentNullException()
+ {
+ // Act & Assert
+ // Note: ResourceManager constructor doesn't validate logger parameter
+ // This test documents that null logger is accepted but will fail at runtime
+ var act = () => new ResourceManager(null!);
+ act.Should().NotThrow();
+ }
+
+ [Fact]
+ public async Task RegisterTaskManagerAsync_AddsNewTaskManager()
+ {
+ // Arrange
+ var taskManagerId = "tm-1";
+ var slotsPerTaskManager = 4;
+
+ // Act
+ await _resourceManager.RegisterTaskManagerAsync(taskManagerId, slotsPerTaskManager);
+
+ // Assert
+ var taskManagers = _resourceManager.GetRegisteredTaskManagers();
+ taskManagers.Should().Contain(taskManagerId);
+ }
+
+ [Fact]
+ public async Task RegisterTaskManagerAsync_CreatesCorrectNumberOfSlots()
+ {
+ // Arrange
+ var taskManagerId = "tm-1";
+ var slotsPerTaskManager = 4;
+
+ // Act
+ await _resourceManager.RegisterTaskManagerAsync(taskManagerId, slotsPerTaskManager);
+
+ // Assert
+ var allSlots = _resourceManager.GetAllSlots();
+ var tmSlots = allSlots.Where(s => s.TaskManagerId == taskManagerId).ToList();
+ tmSlots.Should().HaveCount(slotsPerTaskManager);
+ }
+
+ [Fact]
+ public async Task UnregisterTaskManagerAsync_RemovesTaskManager()
+ {
+ // Arrange
+ var taskManagerId = "tm-1";
+ await _resourceManager.RegisterTaskManagerAsync(taskManagerId, 4);
+
+ // Act
+ await _resourceManager.UnregisterTaskManagerAsync(taskManagerId);
+
+ // Assert
+ var taskManagers = _resourceManager.GetRegisteredTaskManagers();
+ taskManagers.Should().NotContain(taskManagerId);
+ }
+
+ [Fact]
+ public async Task AllocateSlotsAsync_AllocatesRequestedSlots()
+ {
+ // Arrange
+ var taskManagerId = "tm-1";
+ var jobId = "job-1";
+ await _resourceManager.RegisterTaskManagerAsync(taskManagerId, 4);
+
+ // Act
+ var allocatedSlots = await _resourceManager.AllocateSlotsAsync(jobId, 2);
+
+ // Assert
+ allocatedSlots.Should().HaveCount(2);
+ allocatedSlots.All(s => s.IsAllocated).Should().BeTrue();
+ }
+
+ [Fact]
+ public async Task AllocateSlotsAsync_ReducesAvailableSlots()
+ {
+ // Arrange
+ var taskManagerId = "tm-1";
+ var jobId = "job-1";
+ await _resourceManager.RegisterTaskManagerAsync(taskManagerId, 4);
+ var initialAvailable = _resourceManager.GetAvailableSlots().Count();
+
+ // Act
+ await _resourceManager.AllocateSlotsAsync(jobId, 2);
+
+ // Assert
+ var remainingAvailable = _resourceManager.GetAvailableSlots().Count();
+ remainingAvailable.Should().Be(initialAvailable - 2);
+ }
+
+ [Fact]
+ public async Task AllocateSlotsAsync_WithInsufficientSlots_ReturnsPartialAllocation()
+ {
+ // Arrange
+ var taskManagerId = "tm-1";
+ var jobId = "job-1";
+ await _resourceManager.RegisterTaskManagerAsync(taskManagerId, 2);
+
+ // Act - request more slots than available
+ var allocatedSlots = await _resourceManager.AllocateSlotsAsync(jobId, 10);
+
+ // Assert - should get partial allocation (all 2 available slots)
+ allocatedSlots.Should().HaveCount(2);
+ }
+
+ [Fact]
+ public async Task ReleaseSlotAsync_LogsReleaseRequest()
+ {
+ // Arrange
+ var taskManagerId = "tm-1";
+ var jobId = "job-1";
+ await _resourceManager.RegisterTaskManagerAsync(taskManagerId, 4);
+ var allocatedSlots = await _resourceManager.AllocateSlotsAsync(jobId, 2);
+ var slotToRelease = allocatedSlots.First();
+
+ // Act - Note: ReleaseSlotAsync currently only logs, doesn't actually release
+ var act = async () => await _resourceManager.ReleaseSlotAsync(slotToRelease.SlotId);
+
+ // Assert - Should not throw
+ await act.Should().NotThrowAsync();
+ }
+
+ [Fact]
+ public void GetRegisteredTaskManagers_ReturnsAllRegistered()
+ {
+ // Arrange
+ _resourceManager.RegisterTaskManagerAsync("tm-1", 2).Wait();
+ _resourceManager.RegisterTaskManagerAsync("tm-2", 2).Wait();
+
+ // Act
+ var taskManagers = _resourceManager.GetRegisteredTaskManagers().ToList();
+
+ // Assert
+ taskManagers.Should().HaveCount(2);
+ taskManagers.Should().Contain("tm-1");
+ taskManagers.Should().Contain("tm-2");
+ }
+
+ [Fact]
+ public void GetAllSlots_ReturnsAllSlots()
+ {
+ // Arrange
+ _resourceManager.RegisterTaskManagerAsync("tm-1", 2).Wait();
+ _resourceManager.RegisterTaskManagerAsync("tm-2", 3).Wait();
+
+ // Act
+ var allSlots = _resourceManager.GetAllSlots().ToList();
+
+ // Assert
+ allSlots.Should().HaveCount(5);
+ }
+
+ [Fact]
+ public void GetAvailableSlots_ReturnsOnlyFreeSlots()
+ {
+ // Arrange
+ _resourceManager.RegisterTaskManagerAsync("tm-1", 4).Wait();
+ _resourceManager.AllocateSlotsAsync("job-1", 2).Wait();
+
+ // Act
+ var availableSlots = _resourceManager.GetAvailableSlots().ToList();
+
+ // Assert
+ availableSlots.Should().HaveCount(2);
+ availableSlots.All(s => !s.IsAllocated).Should().BeTrue();
+ }
+
+ [Fact]
+ public async Task RequestSlotsAsync_ReturnsTaskSlotList()
+ {
+ // Arrange
+ var taskManagerId = "tm-1";
+ await _resourceManager.RegisterTaskManagerAsync(taskManagerId, 4);
+
+ // Act
+ var slots = await _resourceManager.RequestSlotsAsync("job-1", 2);
+
+ // Assert
+ slots.Should().HaveCount(2);
+ slots.All(s => s.TaskManagerId == taskManagerId).Should().BeTrue();
+ }
+
+ [Fact]
+ public async Task ReleaseSlotsAsync_ReleasesMultipleSlots()
+ {
+ // Arrange
+ var taskManagerId = "tm-1";
+ await _resourceManager.RegisterTaskManagerAsync(taskManagerId, 4);
+ var allocatedSlots = await _resourceManager.AllocateSlotsAsync("job-1", 3);
+ var initialAvailable = _resourceManager.GetAvailableSlots().Count();
+
+ // Act
+ await _resourceManager.ReleaseSlotsAsync(allocatedSlots);
+
+ // Assert
+ var currentAvailable = _resourceManager.GetAvailableSlots().Count();
+ currentAvailable.Should().Be(initialAvailable + 3);
+ }
+
+ [Fact]
+ public async Task RegisterTaskManagerAsync_UpdatesExistingTaskManager()
+ {
+ // Arrange
+ var taskManagerId = "tm-1";
+ await _resourceManager.RegisterTaskManagerAsync(taskManagerId, 2);
+
+ // Act - Register again with different slot count
+ await _resourceManager.RegisterTaskManagerAsync(taskManagerId, 4);
+
+ // Assert
+ var taskManagers = _resourceManager.GetRegisteredTaskManagers();
+ taskManagers.Should().ContainSingle(taskManagerId);
+ }
+
+ [Fact]
+ public async Task AllocateSlotsAsync_WithMultipleTaskManagers_DistributesSlots()
+ {
+ // Arrange
+ await _resourceManager.RegisterTaskManagerAsync("tm-1", 2);
+ await _resourceManager.RegisterTaskManagerAsync("tm-2", 2);
+
+ // Act
+ var allocatedSlots = await _resourceManager.AllocateSlotsAsync("job-1", 3);
+
+ // Assert
+ allocatedSlots.Should().HaveCount(3);
+ allocatedSlots.Should().Contain(s => s.TaskManagerId == "tm-1");
+ }
+
+ [Fact]
+ public async Task GetAvailableSlots_AfterAllocation_ReturnsCorrectCount()
+ {
+ // Arrange
+ await _resourceManager.RegisterTaskManagerAsync("tm-1", 5);
+ await _resourceManager.AllocateSlotsAsync("job-1", 2);
+
+ // Act
+ var availableSlots = _resourceManager.GetAvailableSlots();
+
+ // Assert
+ availableSlots.Should().HaveCount(3);
+ }
+
+ [Fact]
+ public async Task RequestSlotsAsync_WithNoRegisteredTaskManagers_ReturnsEmptyList()
+ {
+ // Act
+ var slots = await _resourceManager.RequestSlotsAsync("job-1", 2);
+
+ // Assert
+ slots.Should().BeEmpty();
+ }
+
+ [Fact]
+ public async Task ReleaseSlotsAsync_WithEmptyList_DoesNotThrow()
+ {
+ // Arrange
+ await _resourceManager.RegisterTaskManagerAsync("tm-1", 2);
+ var emptySlotList = new List();
+
+ // Act
+ var act = async () => await _resourceManager.ReleaseSlotsAsync(emptySlotList);
+
+ // Assert
+ await act.Should().NotThrowAsync();
+ }
+}
+
diff --git a/FlinkDotNet/FlinkDotNet.JobManager/Implementation/JobMaster.cs b/FlinkDotNet/FlinkDotNet.JobManager/Implementation/JobMaster.cs
index 967970b3..c2286c6c 100644
--- a/FlinkDotNet/FlinkDotNet.JobManager/Implementation/JobMaster.cs
+++ b/FlinkDotNet/FlinkDotNet.JobManager/Implementation/JobMaster.cs
@@ -18,7 +18,6 @@ public class JobMaster : IJobMaster
private readonly string _jobId;
private readonly JobGraph _jobGraph;
private readonly IResourceManager _resourceManager;
- private readonly ITemporalClient _temporalClient;
private readonly ILogger _logger;
private ExecutionGraph? _executionGraph;
@@ -38,7 +37,8 @@ public JobMaster(
_jobId = jobId ?? throw new ArgumentNullException(nameof(jobId));
_jobGraph = jobGraph ?? throw new ArgumentNullException(nameof(jobGraph));
_resourceManager = resourceManager ?? throw new ArgumentNullException(nameof(resourceManager));
- _temporalClient = temporalClient ?? throw new ArgumentNullException(nameof(temporalClient));
+ // temporalClient parameter kept for interface compatibility but not used in current implementation
+ ArgumentNullException.ThrowIfNull(temporalClient);
_logger = logger ?? throw new ArgumentNullException(nameof(logger));
}
@@ -73,7 +73,7 @@ public async Task StartJobAsync(CancellationToken cancellationToken = default)
{
_logger.LogError(ex, "Failed to start job {JobId}", _jobId);
_jobState = JobExecutionState.Failed;
- throw;
+ throw new InvalidOperationException($"Failed to start job {_jobId}. See inner exception for details.", ex);
}
}
@@ -112,7 +112,7 @@ public async Task CancelJobAsync(CancellationToken cancellationToken = default)
catch (Exception ex)
{
_logger.LogError(ex, "Failed to cancel job {JobId}", _jobId);
- throw;
+ throw new InvalidOperationException($"Failed to cancel job {_jobId}. See inner exception for details.", ex);
}
}
@@ -175,12 +175,13 @@ public async Task TriggerCheckpointAsync(long checkpointId, CancellationToken ca
catch (Exception ex)
{
_logger.LogError(ex, "Failed to trigger checkpoint {CheckpointId} for job {JobId}", checkpointId, _jobId);
- throw;
+ throw new InvalidOperationException($"Failed to trigger checkpoint {checkpointId} for job {_jobId}. See inner exception for details.", ex);
}
}
- private async Task CreateExecutionGraphAsync(CancellationToken cancellationToken)
+ private Task CreateExecutionGraphAsync(CancellationToken cancellationToken)
{
+ _ = cancellationToken; // Parameter kept for consistency with async pattern
_logger.LogDebug("Creating ExecutionGraph from JobGraph");
ExecutionGraph executionGraph = new()
@@ -232,7 +233,7 @@ private async Task CreateExecutionGraphAsync(CancellationToken c
_logger.LogDebug("ExecutionGraph created: {VertexCount} vertices, {EdgeCount} edges",
executionGraph.ExecutionVertices.Count, executionGraph.ExecutionEdges.Count);
- return await Task.FromResult(executionGraph);
+ return Task.FromResult(executionGraph);
}
private void CreateExecutionEdges(
@@ -399,8 +400,9 @@ private async Task MonitorExecutionAsync(CancellationToken cancellationToken)
}
}
- private async Task HandleTaskFailureAsync(ExecutionVertex vertex, CancellationToken cancellationToken)
+ private Task HandleTaskFailureAsync(ExecutionVertex vertex, CancellationToken cancellationToken)
{
+ _ = cancellationToken; // Parameter kept for consistency with async pattern
_logger.LogError("Task {VertexId} failed: {Error}", vertex.Id, vertex.Error);
// In a full implementation, this would:
@@ -412,7 +414,7 @@ private async Task HandleTaskFailureAsync(ExecutionVertex vertex, CancellationTo
_jobState = JobExecutionState.Failed;
_executionCts?.Cancel();
- await Task.CompletedTask;
+ return Task.CompletedTask;
}
private async Task CheckJobCompletionAsync(CancellationToken cancellationToken)