2
0
mirror of synced 2025-01-12 09:34:40 +00:00

97 lines
3.0 KiB
C#
Raw Normal View History

2023-08-15 11:01:18 +02:00
using k8s;
using Logging;
namespace KubernetesWorkflow
{
public class CrashWatcher
{
2023-09-12 10:31:55 +02:00
private readonly ILog log;
2023-08-15 11:01:18 +02:00
private readonly KubernetesClientConfiguration config;
private readonly string containerName;
private readonly string podName;
private readonly string recipeName;
2023-08-15 11:01:18 +02:00
private readonly string k8sNamespace;
private ILogHandler? logHandler;
private CancellationTokenSource cts;
private Task? worker;
private Exception? workerException;
public CrashWatcher(ILog log, KubernetesClientConfiguration config, string containerName, string podName, string recipeName, string k8sNamespace)
2023-08-15 11:01:18 +02:00
{
this.log = log;
this.config = config;
this.containerName = containerName;
this.podName = podName;
this.recipeName = recipeName;
2023-08-15 11:01:18 +02:00
this.k8sNamespace = k8sNamespace;
cts = new CancellationTokenSource();
}
public void Start(ILogHandler logHandler)
{
if (worker != null) throw new InvalidOperationException();
this.logHandler = logHandler;
cts = new CancellationTokenSource();
worker = Task.Run(Worker);
}
public void Stop()
{
if (worker == null) throw new InvalidOperationException();
cts.Cancel();
worker.Wait();
worker = null;
if (workerException != null) throw new Exception("Exception occurred in CrashWatcher worker thread.", workerException);
}
public bool HasContainerCrashed()
{
using var client = new Kubernetes(config);
return HasContainerBeenRestarted(client);
}
2023-08-15 11:01:18 +02:00
private void Worker()
{
try
{
MonitorContainer(cts.Token);
}
catch (Exception ex)
{
workerException = ex;
}
}
private void MonitorContainer(CancellationToken token)
{
using var client = new Kubernetes(config);
2023-08-15 11:01:18 +02:00
while (!token.IsCancellationRequested)
{
token.WaitHandle.WaitOne(TimeSpan.FromSeconds(10));
2023-08-15 11:01:18 +02:00
if (HasContainerBeenRestarted(client))
2023-08-15 11:01:18 +02:00
{
DownloadCrashedContainerLogs(client);
2023-08-15 11:01:18 +02:00
return;
}
}
}
private bool HasContainerBeenRestarted(Kubernetes client)
{
var podInfo = client.ReadNamespacedPod(podName, k8sNamespace);
return podInfo.Status.ContainerStatuses.Any(c => c.RestartCount > 0);
}
private void DownloadCrashedContainerLogs(Kubernetes client)
2023-08-15 11:01:18 +02:00
{
log.Log("Pod crash detected for " + containerName);
using var stream = client.ReadNamespacedPodLog(podName, k8sNamespace, recipeName, previous: true);
2023-08-15 11:01:18 +02:00
logHandler!.Log(stream);
}
}
}