Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,166 @@
// 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 FluentAssertions;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using Moq;

namespace FlinkDotNet.JobManager.Tests;

public class HeartbeatMonitoringServiceTests
{
private readonly Mock<IResourceManager> _mockResourceManager;
private readonly Mock<ILogger<HeartbeatMonitoringService>> _mockLogger;
private readonly HeartbeatConfiguration _configuration;

public HeartbeatMonitoringServiceTests()
{
_mockResourceManager = new Mock<IResourceManager>();
_mockLogger = new Mock<ILogger<HeartbeatMonitoringService>>();
_configuration = new HeartbeatConfiguration
{
TimeoutSeconds = 2, // Short timeout for testing
CheckIntervalSeconds = 1 // Short interval for testing
};
}

[Fact]
public void Constructor_WithNullResourceManager_ThrowsArgumentNullException()
{
// Arrange & Act & Assert
var act = () => new HeartbeatMonitoringService(
null!,
Options.Create(_configuration),
_mockLogger.Object);

act.Should().Throw<ArgumentNullException>()
.WithParameterName("resourceManager");
}

[Fact]
public void Constructor_WithNullConfiguration_ThrowsArgumentNullException()
{
// Arrange & Act & Assert
var act = () => new HeartbeatMonitoringService(
_mockResourceManager.Object,
null!,
_mockLogger.Object);

act.Should().Throw<ArgumentNullException>();
}

[Fact]
public void Constructor_WithNullLogger_ThrowsArgumentNullException()
{
// Arrange & Act & Assert
var act = () => new HeartbeatMonitoringService(
_mockResourceManager.Object,
Options.Create(_configuration),
null!);

act.Should().Throw<ArgumentNullException>()
.WithParameterName("logger");
}

[Fact]
public async Task HeartbeatMonitoring_DetectsTimeout_AndUnregistersTaskManager()
{
// Arrange
var taskManagerId = "tm-timeout";
var oldHeartbeat = DateTime.UtcNow.AddSeconds(-10); // Old heartbeat (10 seconds ago)

_mockResourceManager
.Setup(rm => rm.GetRegisteredTaskManagers())
.Returns(new[] { taskManagerId });

_mockResourceManager
.Setup(rm => rm.GetLastHeartbeat(taskManagerId))
.Returns(oldHeartbeat);

var service = new HeartbeatMonitoringService(
_mockResourceManager.Object,
Options.Create(_configuration),
_mockLogger.Object);

// Act
await service.StartAsync(CancellationToken.None);
await Task.Delay(TimeSpan.FromSeconds(2)); // Wait for check interval
await service.StopAsync(CancellationToken.None);

// Assert
_mockResourceManager.Verify(
rm => rm.UnregisterTaskManagerAsync(taskManagerId, It.IsAny<CancellationToken>()),
Times.AtLeastOnce());
}

[Fact]
public async Task HeartbeatMonitoring_WithRecentHeartbeat_DoesNotUnregister()
{
// Arrange
var taskManagerId = "tm-healthy";

_mockResourceManager
.Setup(rm => rm.GetRegisteredTaskManagers())
.Returns(new[] { taskManagerId });

// Return a fresh heartbeat each time it's queried
_mockResourceManager
.Setup(rm => rm.GetLastHeartbeat(taskManagerId))
.Returns(() => DateTime.UtcNow);

var service = new HeartbeatMonitoringService(
_mockResourceManager.Object,
Options.Create(_configuration),
_mockLogger.Object);

// Act
await service.StartAsync(CancellationToken.None);
await Task.Delay(TimeSpan.FromSeconds(2)); // Wait for check interval
await service.StopAsync(CancellationToken.None);

// Assert
_mockResourceManager.Verify(
rm => rm.UnregisterTaskManagerAsync(taskManagerId, It.IsAny<CancellationToken>()),
Times.Never());
}

[Fact]
public async Task HeartbeatMonitoring_WithNoTaskManagers_DoesNothing()
{
// Arrange
_mockResourceManager
.Setup(rm => rm.GetRegisteredTaskManagers())
.Returns(Array.Empty<string>());

var service = new HeartbeatMonitoringService(
_mockResourceManager.Object,
Options.Create(_configuration),
_mockLogger.Object);

// Act
await service.StartAsync(CancellationToken.None);
await Task.Delay(TimeSpan.FromSeconds(2)); // Wait for check interval
await service.StopAsync(CancellationToken.None);

// Assert
_mockResourceManager.Verify(
rm => rm.UnregisterTaskManagerAsync(It.IsAny<string>(), It.IsAny<CancellationToken>()),
Times.Never());
}

[Fact]
public void HeartbeatConfiguration_HasCorrectDefaults()
{
// Arrange & Act
var config = new HeartbeatConfiguration();

// Assert
config.TimeoutSeconds.Should().Be(30);
config.CheckIntervalSeconds.Should().Be(10);
HeartbeatConfiguration.SectionName.Should().Be("Heartbeat");
}
}
168 changes: 168 additions & 0 deletions FlinkDotNet/FlinkDotNet.JobManager.Tests/HeartbeatTests.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,168 @@
// 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 FluentAssertions;
using Microsoft.Extensions.Logging;
using Moq;

namespace FlinkDotNet.JobManager.Tests;

public class HeartbeatTests
{
private readonly Mock<ILogger<ResourceManager>> _mockLogger;
private readonly ResourceManager _resourceManager;

public HeartbeatTests()
{
_mockLogger = new Mock<ILogger<ResourceManager>>();
_resourceManager = new ResourceManager(_mockLogger.Object);
}

[Fact]
public async Task RecordHeartbeatAsync_UpdatesLastHeartbeatTimestamp()
{
// Arrange
var taskManagerId = "tm-heartbeat-1";
var numberOfSlots = 4;

await _resourceManager.RegisterTaskManagerAsync(taskManagerId, numberOfSlots);
DateTime? initialHeartbeat = _resourceManager.GetLastHeartbeat(taskManagerId);

// Wait a small amount to ensure timestamp difference
await Task.Delay(10);

// Act
await _resourceManager.RecordHeartbeatAsync(taskManagerId);

// Assert
DateTime? updatedHeartbeat = _resourceManager.GetLastHeartbeat(taskManagerId);

updatedHeartbeat.Should().NotBeNull();
initialHeartbeat.Should().NotBeNull();
updatedHeartbeat.Should().BeAfter(initialHeartbeat.Value);
}

[Fact]
public async Task RecordHeartbeatAsync_ForUnregisteredTaskManager_LogsWarning()
{
// Arrange
var unregisteredTaskManagerId = "tm-unregistered";

// Act
await _resourceManager.RecordHeartbeatAsync(unregisteredTaskManagerId);

// Assert
// Verify that a warning was logged (implementation logs warning)
DateTime? heartbeat = _resourceManager.GetLastHeartbeat(unregisteredTaskManagerId);
heartbeat.Should().BeNull();
}

[Fact]
public async Task GetLastHeartbeat_ForRegisteredTaskManager_ReturnsTimestamp()
{
// Arrange
var taskManagerId = "tm-heartbeat-2";
var numberOfSlots = 4;

await _resourceManager.RegisterTaskManagerAsync(taskManagerId, numberOfSlots);

// Act
DateTime? heartbeat = _resourceManager.GetLastHeartbeat(taskManagerId);

// Assert
heartbeat.Should().NotBeNull();
heartbeat.Should().BeCloseTo(DateTime.UtcNow, TimeSpan.FromSeconds(5));
}

[Fact]
public void GetLastHeartbeat_ForUnregisteredTaskManager_ReturnsNull()
{
// Arrange
var unregisteredTaskManagerId = "tm-not-registered";

// Act
DateTime? heartbeat = _resourceManager.GetLastHeartbeat(unregisteredTaskManagerId);

// Assert
heartbeat.Should().BeNull();
}

[Fact]
public async Task RegisterTaskManagerAsync_InitializesLastHeartbeat()
{
// Arrange
var taskManagerId = "tm-heartbeat-3";
var numberOfSlots = 4;

// Act
await _resourceManager.RegisterTaskManagerAsync(taskManagerId, numberOfSlots);

// Assert
DateTime? heartbeat = _resourceManager.GetLastHeartbeat(taskManagerId);
heartbeat.Should().NotBeNull();
heartbeat.Should().BeCloseTo(DateTime.UtcNow, TimeSpan.FromSeconds(5));
}

[Fact]
public async Task MultipleHeartbeats_UpdateTimestampSequentially()
{
// Arrange
var taskManagerId = "tm-heartbeat-4";
var numberOfSlots = 4;

await _resourceManager.RegisterTaskManagerAsync(taskManagerId, numberOfSlots);

// Act & Assert
DateTime? heartbeat1 = _resourceManager.GetLastHeartbeat(taskManagerId);
heartbeat1.Should().NotBeNull();

await Task.Delay(10);
await _resourceManager.RecordHeartbeatAsync(taskManagerId);
DateTime? heartbeat2 = _resourceManager.GetLastHeartbeat(taskManagerId);
heartbeat2.Should().BeAfter(heartbeat1.Value);

await Task.Delay(10);
await _resourceManager.RecordHeartbeatAsync(taskManagerId);
DateTime? heartbeat3 = _resourceManager.GetLastHeartbeat(taskManagerId);
heartbeat3.Should().BeAfter(heartbeat2.Value);
}

[Fact]
public async Task ConcurrentHeartbeats_AreThreadSafe()
{
// Arrange
var taskManagerId = "tm-concurrent";
var numberOfSlots = 4;

await _resourceManager.RegisterTaskManagerAsync(taskManagerId, numberOfSlots);

// Act - Send concurrent heartbeats
var tasks = Enumerable.Range(0, 10).Select(_ =>
Task.Run(async () => await _resourceManager.RecordHeartbeatAsync(taskManagerId))
);

await Task.WhenAll(tasks);

// Assert - Should not throw and should have a valid timestamp
DateTime? heartbeat = _resourceManager.GetLastHeartbeat(taskManagerId);
heartbeat.Should().NotBeNull();
}

[Fact]
public void SynchronousRegisterTaskManager_InitializesLastHeartbeat()
{
// Arrange
var taskManagerId = "tm-sync-heartbeat";
var numberOfSlots = 4;

// Act
_resourceManager.RegisterTaskManager(taskManagerId, numberOfSlots);

// Assert
DateTime? heartbeat = _resourceManager.GetLastHeartbeat(taskManagerId);
heartbeat.Should().NotBeNull();
heartbeat.Should().BeCloseTo(DateTime.UtcNow, TimeSpan.FromSeconds(5));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@
// Licensed under the Apache License, Version 2.0.
// See LICENSE file in the project root for full license information.

#nullable enable

using FlinkDotNet.JobManager.Controllers;
using FlinkDotNet.JobManager.Interfaces;
using FlinkDotNet.JobManager.Models;
Expand Down Expand Up @@ -125,9 +127,9 @@
// Arrange
var jobId = "non-existent-job";

_mockDispatcher

Check warning on line 130 in FlinkDotNet/FlinkDotNet.JobManager.Tests/JobsControllerTests.cs

View workflow job for this annotation

GitHub Actions / Run .NET Unit Tests

Argument of type 'ISetup<IDispatcher, Task<JobStatus>>' cannot be used for parameter 'mock' of type 'IReturns<IDispatcher, Task<JobStatus?>>' in 'IReturnsResult<IDispatcher> ReturnsExtensions.ReturnsAsync<IDispatcher, JobStatus?>(IReturns<IDispatcher, Task<JobStatus?>> mock, JobStatus? value)' due to differences in the nullability of reference types.

Check warning on line 130 in FlinkDotNet/FlinkDotNet.JobManager.Tests/JobsControllerTests.cs

View workflow job for this annotation

GitHub Actions / Run .NET Unit Tests

Argument of type 'ISetup<IDispatcher, Task<JobStatus>>' cannot be used for parameter 'mock' of type 'IReturns<IDispatcher, Task<JobStatus?>>' in 'IReturnsResult<IDispatcher> ReturnsExtensions.ReturnsAsync<IDispatcher, JobStatus?>(IReturns<IDispatcher, Task<JobStatus?>> mock, JobStatus? value)' due to differences in the nullability of reference types.
.Setup(d => d.GetJobStatusAsync(jobId, It.IsAny<CancellationToken>()))
.ReturnsAsync((JobStatus?) null);
.ReturnsAsync((JobStatus?)null);

// Act
var result = await _controller.GetJobStatus(jobId);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -154,6 +154,29 @@ public IActionResult UnregisterTaskManager(string taskManagerId)
message = $"TaskManager {taskManagerId} unregistered successfully"
});
}

/// <summary>
/// Record heartbeat from a TaskManager.
/// TaskManagers should call this endpoint periodically to indicate they are alive.
/// </summary>
/// <param name="taskManagerId">ID of the TaskManager sending the heartbeat.</param>
/// <returns>Heartbeat acknowledgement.</returns>
[HttpPost("taskmanagers/{taskManagerId}/heartbeat")]
[ProducesResponseType(typeof(object), StatusCodes.Status200OK)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
public async Task<IActionResult> RecordHeartbeat(string taskManagerId)
{
this._logger.LogDebug("Received heartbeat from TaskManager: {TaskManagerId}", taskManagerId);

await this._resourceManager.RecordHeartbeatAsync(taskManagerId);

return Ok(new
{
message = "Heartbeat recorded",
taskManagerId,
timestamp = DateTime.UtcNow
});
}
}

/// <summary>
Expand Down
Loading
Loading