Browse Source

Adding health checks (#432)

jkotalik/porter
areller 6 years ago
committed by GitHub
parent
commit
f27905b09a
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 27
      src/Microsoft.Tye.Core/ApplicationFactory.cs
  2. 15
      src/Microsoft.Tye.Core/ConfigModel/ConfigApplication.cs
  3. 18
      src/Microsoft.Tye.Core/ConfigModel/ConfigHttpProber.cs
  4. 16
      src/Microsoft.Tye.Core/ConfigModel/ConfigProbe.cs
  5. 2
      src/Microsoft.Tye.Core/ConfigModel/ConfigService.cs
  6. 4
      src/Microsoft.Tye.Core/ContainerServiceBuilder.cs
  7. 11
      src/Microsoft.Tye.Core/CoreStrings.resx
  8. 4
      src/Microsoft.Tye.Core/ExecutableServiceBuilder.cs
  9. 16
      src/Microsoft.Tye.Core/HttpProberBuilder.cs
  10. 68
      src/Microsoft.Tye.Core/KubernetesManifestGenerator.cs
  11. 16
      src/Microsoft.Tye.Core/ProbeBuilder.cs
  12. 4
      src/Microsoft.Tye.Core/ProjectServiceBuilder.cs
  13. 173
      src/Microsoft.Tye.Core/Serialization/ConfigServiceParser.cs
  14. 85
      src/Microsoft.Tye.Hosting/DockerRunner.cs
  15. 69
      src/Microsoft.Tye.Hosting/HttpProxyService.cs
  16. 17
      src/Microsoft.Tye.Hosting/Model/HttpProber.cs
  17. 18
      src/Microsoft.Tye.Hosting/Model/Probe.cs
  18. 13
      src/Microsoft.Tye.Hosting/Model/ReplicaBinding.cs
  19. 2
      src/Microsoft.Tye.Hosting/Model/ReplicaState.cs
  20. 8
      src/Microsoft.Tye.Hosting/Model/ReplicaStatus.cs
  21. 5
      src/Microsoft.Tye.Hosting/Model/Service.cs
  22. 2
      src/Microsoft.Tye.Hosting/Model/ServiceDescription.cs
  23. 1
      src/Microsoft.Tye.Hosting/Model/V1/V1ReplicaStatus.cs
  24. 2
      src/Microsoft.Tye.Hosting/PortAssigner.cs
  25. 9
      src/Microsoft.Tye.Hosting/ProcessRunner.cs
  26. 53
      src/Microsoft.Tye.Hosting/ProxyService.cs
  27. 455
      src/Microsoft.Tye.Hosting/ReplicaMonitor.cs
  28. 3
      src/Microsoft.Tye.Hosting/TyeDashboardApi.cs
  29. 1
      src/Microsoft.Tye.Hosting/TyeHost.cs
  30. 28
      src/tye/ApplicationBuilderExtensions.cs
  31. 453
      test/E2ETest/HealthCheckTests.cs
  32. 3
      test/E2ETest/Microsoft.Tye.E2ETests.csproj
  33. 143
      test/E2ETest/ReplicaStoppingTests.cs
  34. 35
      test/E2ETest/TyeGenerateTests.cs
  35. 76
      test/E2ETest/testassets/generate/health-checks.yaml
  36. 32
      test/E2ETest/testassets/projects/health-checks/api/Program.cs
  37. 30
      test/E2ETest/testassets/projects/health-checks/api/Properties/launchSettings.json
  38. 166
      test/E2ETest/testassets/projects/health-checks/api/Startup.cs
  39. 8
      test/E2ETest/testassets/projects/health-checks/api/api.csproj
  40. 9
      test/E2ETest/testassets/projects/health-checks/api/appsettings.Development.json
  41. 10
      test/E2ETest/testassets/projects/health-checks/api/appsettings.json
  42. 34
      test/E2ETest/testassets/projects/health-checks/health-checks.sln
  43. 36
      test/E2ETest/testassets/projects/health-checks/tye-all.yaml
  44. 30
      test/E2ETest/testassets/projects/health-checks/tye-ingress.yaml
  45. 17
      test/E2ETest/testassets/projects/health-checks/tye-liveness.yaml
  46. 10
      test/E2ETest/testassets/projects/health-checks/tye-none.yaml
  47. 26
      test/E2ETest/testassets/projects/health-checks/tye-proxy.yaml
  48. 17
      test/E2ETest/testassets/projects/health-checks/tye-readiness.yaml
  49. 127
      test/Test.Infrastructure/TestHelpers.cs

27
src/Microsoft.Tye.Core/ApplicationFactory.cs

@ -97,6 +97,9 @@ namespace Microsoft.Tye
}
project.Replicas = configService.Replicas ?? 1;
project.Liveness = configService.Liveness != null ? GetProbeBuilder(configService.Liveness) : null;
project.Readiness = configService.Readiness != null ? GetProbeBuilder(configService.Readiness) : null;
// We don't apply more container defaults here because we might need
// to prompt for the registry name.
project.ContainerInfo = new ContainerInfo() { UseMultiphaseDockerfile = false, };
@ -117,6 +120,9 @@ namespace Microsoft.Tye
DockerFileContext = GetDockerFileContext(source, configService)
};
service = container;
container.Liveness = configService.Liveness != null ? GetProbeBuilder(configService.Liveness) : null;
container.Readiness = configService.Readiness != null ? GetProbeBuilder(configService.Readiness) : null;
}
else if (!string.IsNullOrEmpty(configService.Executable))
{
@ -139,6 +145,9 @@ namespace Microsoft.Tye
Replicas = configService.Replicas ?? 1
};
service = executable;
executable.Liveness = configService.Liveness != null ? GetProbeBuilder(configService.Liveness) : null;
executable.Readiness = configService.Readiness != null ? GetProbeBuilder(configService.Readiness) : null;
}
else if (!string.IsNullOrEmpty(configService.Include))
{
@ -379,5 +388,23 @@ namespace Microsoft.Tye
}
}
}
private static ProbeBuilder GetProbeBuilder(ConfigProbe config) => new ProbeBuilder()
{
Http = config.Http != null ? GetHttpProberBuilder(config.Http) : null,
InitialDelay = config.InitialDelay,
Period = config.Period,
Timeout = config.Timeout,
SuccessThreshold = config.SuccessThreshold,
FailureThreshold = config.FailureThreshold
};
private static HttpProberBuilder GetHttpProberBuilder(ConfigHttpProber config) => new HttpProberBuilder()
{
Path = config.Path,
Headers = config.Headers,
Port = config.Port,
Protocol = config.Protocol
};
}
}

15
src/Microsoft.Tye.Core/ConfigModel/ConfigApplication.cs

@ -135,6 +135,21 @@ namespace Microsoft.Tye.ConfigModel
string.Join(Environment.NewLine, results.Select(r => r.ErrorMessage)));
}
}
var probes = new[] { (Name: "liveness", Probe: service.Liveness), (Name: "readiness", Probe: service.Readiness) }.Where(p => p.Probe != null).ToArray();
foreach (var probe in probes)
{
if (probe.Name == "liveness" && probe.Probe.SuccessThreshold != 1)
{
throw new TyeYamlException(CoreStrings.FormatSuccessThresholdMustBeOne(probe.Name));
}
// right now only http is supported, so it must be set
if (probe.Probe!.Http == null)
{
throw new TyeYamlException(CoreStrings.FormatProberRequired(probe.Name));
}
}
}
foreach (var ingress in config.Ingress)

18
src/Microsoft.Tye.Core/ConfigModel/ConfigHttpProber.cs

@ -0,0 +1,18 @@
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.
// See the LICENSE file in the project root for more information.
using System.Collections.Generic;
using System.ComponentModel.DataAnnotations;
using System.Net.Http.Headers;
namespace Microsoft.Tye.ConfigModel
{
public class ConfigHttpProber
{
[Required] public string Path { get; set; } = default!;
public int? Port { get; set; }
public string? Protocol { get; set; }
public List<KeyValuePair<string, object>> Headers { get; set; } = new List<KeyValuePair<string, object>>();
}
}

16
src/Microsoft.Tye.Core/ConfigModel/ConfigProbe.cs

@ -0,0 +1,16 @@
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.
// See the LICENSE file in the project root for more information.
namespace Microsoft.Tye.ConfigModel
{
public class ConfigProbe
{
public ConfigHttpProber? Http { get; set; }
public int InitialDelay { get; set; } = 0;
public int Period { get; set; } = 1;
public int Timeout { get; set; } = 1;
public int SuccessThreshold { get; set; } = 1;
public int FailureThreshold { get; set; } = 3;
}
}

2
src/Microsoft.Tye.Core/ConfigModel/ConfigService.cs

@ -37,5 +37,7 @@ namespace Microsoft.Tye.ConfigModel
[YamlMember(Alias = "env")]
public List<ConfigConfigurationSource> Configuration { get; set; } = new List<ConfigConfigurationSource>();
public List<BuildProperty> BuildProperties { get; set; } = new List<BuildProperty>();
public ConfigProbe? Liveness { get; set; }
public ConfigProbe? Readiness { get; set; }
}
}

4
src/Microsoft.Tye.Core/ContainerServiceBuilder.cs

@ -29,5 +29,9 @@ namespace Microsoft.Tye
public List<EnvironmentVariableBuilder> EnvironmentVariables { get; } = new List<EnvironmentVariableBuilder>();
public List<VolumeBuilder> Volumes { get; } = new List<VolumeBuilder>();
public ProbeBuilder? Liveness { get; set; }
public ProbeBuilder? Readiness { get; set; }
}
}

11
src/Microsoft.Tye.Core/CoreStrings.resx

@ -141,6 +141,12 @@
<data name="MultipleBindingWithSamePort" xml:space="preserve">
<value>Cannot have multiple {0} bindings with the same port.</value>
</data>
<data name="ProberRequired" xml:space="preserve">
<value>A prober must be configured for the {0} probe.</value>
</data>
<data name="SuccessThresholdMustBeOne" xml:space="preserve">
<value>"successThreshold" for {0} probe must be set to "1".</value>
</data>
<data name="MustBeABoolean" xml:space="preserve">
<value>"{value}" must be a boolean value (true/false).</value>
</data>
@ -150,6 +156,9 @@
<data name="MustBePositive" xml:space="preserve">
<value>"{value}" value cannot be negative.</value>
</data>
<data name="MustBeGreaterThanZero" xml:space="preserve">
<value>"{value}" value must be greater than zero.</value>
</data>
<data name="ProjectImageExecutableExclusive" xml:space="preserve">
<value>Cannot have both "{0}" and "{1}" set for a service. Only one of project, image, and executable can be set for a given service.</value>
</data>
@ -162,4 +171,4 @@
<data name="UnrecognizedKey" xml:space="preserve">
<value>Unexpected key "{key}" in tye.yaml.</value>
</data>
</root>
</root>

4
src/Microsoft.Tye.Core/ExecutableServiceBuilder.cs

@ -23,5 +23,9 @@ namespace Microsoft.Tye
public int Replicas { get; set; } = 1;
public List<EnvironmentVariableBuilder> EnvironmentVariables { get; } = new List<EnvironmentVariableBuilder>();
public ProbeBuilder? Liveness { get; set; }
public ProbeBuilder? Readiness { get; set; }
}
}

16
src/Microsoft.Tye.Core/HttpProberBuilder.cs

@ -0,0 +1,16 @@
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.
// See the LICENSE file in the project root for more information.
using System.Collections.Generic;
namespace Microsoft.Tye
{
public class HttpProberBuilder
{
public string Path { get; set; } = default!;
public int? Port { get; set; }
public string? Protocol { get; set; }
public List<KeyValuePair<string, object>> Headers { get; set; } = new List<KeyValuePair<string, object>>();
}
}

68
src/Microsoft.Tye.Core/KubernetesManifestGenerator.cs

@ -378,6 +378,16 @@ namespace Microsoft.Tye
}
}
}
if (project.Liveness != null)
{
AddProbe(project, container, project.Liveness!, "livenessProbe");
}
if (project.Readiness != null)
{
AddProbe(project, container, project.Readiness!, "readinessProbe");
}
}
foreach (var sidecar in project.Sidecars)
@ -452,6 +462,64 @@ namespace Microsoft.Tye
return new KubernetesDeploymentOutput(project.Name, new YamlDocument(root));
}
private static void AddProbe(ServiceBuilder service, YamlMappingNode container, ProbeBuilder builder, string name)
{
var probe = new YamlMappingNode();
container.Add(name, probe);
if (builder.Http != null)
{
var builderHttp = builder.Http;
var httpGet = new YamlMappingNode();
probe.Add("httpGet", httpGet);
httpGet.Add("path", builderHttp.Path);
if (builderHttp.Protocol != null)
{
httpGet.Add("scheme", builderHttp.Protocol.ToUpper());
}
if (builderHttp.Port != null)
{
httpGet.Add("port", builderHttp.Port.ToString()!);
}
else
{
// If port is not given, we pull default port
var binding = service.Bindings.First(b => builderHttp.Protocol == null || b.Protocol == builderHttp.Protocol);
if (binding.Port != null)
{
httpGet.Add("port", binding.Port.Value.ToString());
}
if (builderHttp.Protocol == null && binding.Protocol != null)
{
httpGet.Add("scheme", binding.Protocol.ToUpper());
}
}
if (builderHttp.Headers.Count > 0)
{
var headers = new YamlSequenceNode();
httpGet.Add("httpHeaders", headers);
foreach (var builderHeader in builderHttp.Headers)
{
var header = new YamlMappingNode();
header.Add("name", builderHeader.Key);
header.Add("value", builderHeader.Value.ToString()!);
headers.Add(header);
}
}
}
probe.Add("initialDelaySeconds", builder.InitialDelay.ToString());
probe.Add("periodSeconds", builder.Period.ToString()!);
probe.Add("successThreshold", builder.SuccessThreshold.ToString()!);
probe.Add("failureThreshold", builder.FailureThreshold.ToString()!);
}
private static void AddEnvironmentVariablesForComputedBindings(YamlSequenceNode env, ComputedBindings bindings)
{
foreach (var binding in bindings.Bindings.OfType<EnvironmentVariableInputBinding>())

16
src/Microsoft.Tye.Core/ProbeBuilder.cs

@ -0,0 +1,16 @@
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.
// See the LICENSE file in the project root for more information.
namespace Microsoft.Tye
{
public class ProbeBuilder
{
public HttpProberBuilder? Http { get; set; }
public int InitialDelay { get; set; }
public int Period { get; set; }
public int Timeout { get; set; }
public int SuccessThreshold { get; set; }
public int FailureThreshold { get; set; }
}
}

4
src/Microsoft.Tye.Core/ProjectServiceBuilder.cs

@ -54,5 +54,9 @@ namespace Microsoft.Tye
public Dictionary<string, string> BuildProperties { get; } = new Dictionary<string, string>();
public List<SidecarBuilder> Sidecars { get; } = new List<SidecarBuilder>();
public ProbeBuilder? Liveness { get; set; }
public ProbeBuilder? Readiness { get; set; }
}
}

173
src/Microsoft.Tye.Core/Serialization/ConfigServiceParser.cs

@ -3,6 +3,8 @@
// See the LICENSE file in the project root for more information.
using System.Collections.Generic;
using System.Linq;
using Microsoft.Build.Evaluation;
using Microsoft.Tye.ConfigModel;
using YamlDotNet.RepresentationModel;
@ -37,6 +39,7 @@ namespace Tye.Serialization
{
throw new TyeYamlException(child.Value.Start, CoreStrings.FormatMustBeABoolean(key));
}
service.External = external;
break;
case "image":
@ -56,6 +59,7 @@ namespace Tye.Serialization
{
throw new TyeYamlException(child.Value.Start, CoreStrings.FormatExpectedYamlSequence(key));
}
HandleBuildProperties((child.Value as YamlSequenceNode)!, service.BuildProperties);
break;
case "include":
@ -69,6 +73,7 @@ namespace Tye.Serialization
{
throw new TyeYamlException(child.Value.Start, CoreStrings.FormatMustBeABoolean(key));
}
service.Build = build;
break;
case "executable":
@ -115,8 +120,17 @@ namespace Tye.Serialization
{
throw new TyeYamlException(child.Value.Start, CoreStrings.FormatExpectedYamlSequence(key));
}
HandleServiceConfiguration((child.Value as YamlSequenceNode)!, service.Configuration);
break;
case "liveness":
service.Liveness = new ConfigProbe();
HandleServiceProbe((YamlMappingNode)child.Value, service.Liveness!);
break;
case "readiness":
service.Readiness = new ConfigProbe();
HandleServiceProbe((YamlMappingNode)child.Value, service.Readiness!);
break;
default:
throw new TyeYamlException(child.Key.Start, CoreStrings.FormatUnrecognizedKey(key));
}
@ -187,6 +201,165 @@ namespace Tye.Serialization
}
}
private static void HandleServiceProbe(YamlMappingNode yamlMappingNode, ConfigProbe probe)
{
foreach (var child in yamlMappingNode.Children)
{
var key = YamlParser.GetScalarValue(child.Key);
switch (key)
{
case "http":
probe.Http = new ConfigHttpProber();
HandleServiceHttpProber((YamlMappingNode)child.Value, probe.Http!);
break;
case "initialDelay":
if (!int.TryParse(YamlParser.GetScalarValue(key, child.Value), out var initialDelay))
{
throw new TyeYamlException(child.Value.Start, CoreStrings.FormatMustBeAnInteger(key));
}
if (initialDelay < 0)
{
throw new TyeYamlException(child.Value.Start, CoreStrings.FormatMustBePositive(key));
}
probe.InitialDelay = initialDelay;
break;
case "period":
if (!int.TryParse(YamlParser.GetScalarValue(key, child.Value), out var period))
{
throw new TyeYamlException(child.Value.Start, CoreStrings.FormatMustBeAnInteger(key));
}
if (period < 1)
{
throw new TyeYamlException(child.Value.Start, CoreStrings.FormatMustBeGreaterThanZero(key));
}
probe.Period = period;
break;
case "timeout":
if (!int.TryParse(YamlParser.GetScalarValue(key, child.Value), out var timeout))
{
throw new TyeYamlException(child.Value.Start, CoreStrings.FormatMustBeAnInteger(key));
}
if (timeout < 1)
{
throw new TyeYamlException(child.Value.Start, CoreStrings.FormatMustBeGreaterThanZero(key));
}
probe.Timeout = timeout;
break;
case "successThreshold":
if (!int.TryParse(YamlParser.GetScalarValue(key, child.Value), out var successThreshold))
{
throw new TyeYamlException(child.Value.Start, CoreStrings.FormatMustBeAnInteger(key));
}
if (successThreshold < 1)
{
throw new TyeYamlException(child.Value.Start, CoreStrings.FormatMustBeGreaterThanZero(key));
}
probe.SuccessThreshold = successThreshold;
break;
case "failureThreshold":
if (!int.TryParse(YamlParser.GetScalarValue(key, child.Value), out var failureThreshold))
{
throw new TyeYamlException(child.Value.Start, CoreStrings.FormatMustBeAnInteger(key));
}
if (failureThreshold < 1)
{
throw new TyeYamlException(child.Value.Start, CoreStrings.FormatMustBeGreaterThanZero(key));
}
probe.FailureThreshold = failureThreshold;
break;
default:
throw new TyeYamlException(child.Key.Start, CoreStrings.FormatUnrecognizedKey(key));
}
}
}
private static void HandleServiceHttpProber(YamlMappingNode yamlMappingNode, ConfigHttpProber prober)
{
foreach (var child in yamlMappingNode.Children)
{
var key = YamlParser.GetScalarValue(child.Key);
switch (key)
{
case "path":
prober.Path = YamlParser.GetScalarValue("path", child.Value);
break;
case "port":
if (!int.TryParse(YamlParser.GetScalarValue(key, child.Value), out var port))
{
throw new TyeYamlException(child.Value.Start, CoreStrings.FormatMustBeAnInteger(key));
}
prober.Port = port;
break;
case "protocol":
prober.Path = YamlParser.GetScalarValue("protocol", child.Value);
break;
case "headers":
prober.Headers = new List<KeyValuePair<string, object>>();
var headersNode = child.Value as YamlSequenceNode;
if (headersNode is null)
{
throw new TyeYamlException(child.Value.Start, CoreStrings.FormatExpectedYamlSequence("headers"));
}
foreach (var header in headersNode.Children)
{
HandleServiceProbeHttpHeader((YamlMappingNode)header, prober.Headers);
}
break;
default:
throw new TyeYamlException(child.Key.Start, CoreStrings.FormatUnrecognizedKey(key));
}
}
}
private static void HandleServiceProbeHttpHeader(YamlMappingNode yamlMappingNode, List<KeyValuePair<string, object>> headers)
{
string? name = null;
object? value = null;
foreach (var child in yamlMappingNode.Children)
{
var key = YamlParser.GetScalarValue(child.Key);
switch (key)
{
case "name":
name = YamlParser.GetScalarValue("name", child.Value);
break;
case "value":
value = YamlParser.GetScalarValue("value", child.Value);
break;
default:
throw new TyeYamlException(child.Key.Start, CoreStrings.FormatUnrecognizedKey(key));
}
}
if (name is null)
{
throw new TyeYamlException(yamlMappingNode.Start, CoreStrings.FormatExpectedYamlScalar("name"));
}
else if (value is null)
{
throw new TyeYamlException(yamlMappingNode.Start, CoreStrings.FormatExpectedYamlScalar("value"));
}
headers.Add(new KeyValuePair<string, object>(name, value));
}
private static void HandleServiceVolumeNameMapping(YamlMappingNode yamlMappingNode, ConfigVolume volume)
{
foreach (var child in yamlMappingNode!.Children)

85
src/Microsoft.Tye.Hosting/DockerRunner.cs

@ -248,6 +248,7 @@ namespace Microsoft.Tye.Hosting
if (hasPorts)
{
status.Ports = ports.Select(p => p.Port);
status.Bindings = ports.Select(p => new ReplicaBinding() { Port = p.Port, ExternalPort = p.ExternalPort, Protocol = p.Protocol }).ToList();
// These are the ports that the application should use for binding
@ -327,7 +328,7 @@ namespace Microsoft.Tye.Hosting
if (result.ExitCode != 0)
{
_logger.LogError("docker run failed for {ServiceName} with exit code {ExitCode}:" + result.StandardError, service.Description.Name, result.ExitCode);
service.Replicas.TryRemove(replica, out _);
service.Replicas.TryRemove(replica, out var _);
service.ReplicaEvents.OnNext(new ReplicaEvent(ReplicaState.Removed, status));
PrintStdOutAndErr(service, replica, result);
@ -367,42 +368,73 @@ namespace Microsoft.Tye.Hosting
PrintStdOutAndErr(service, replica, result);
}
service.ReplicaEvents.OnNext(new ReplicaEvent(ReplicaState.Started, status));
_logger.LogInformation("Collecting docker logs for {ContainerName}.", replica);
var backOff = TimeSpan.FromSeconds(5);
var sentStartedEvent = false;
while (!dockerInfo.StoppingTokenSource.Token.IsCancellationRequested)
{
var logsRes = await ProcessUtil.RunAsync("docker", $"logs -f {containerId}",
outputDataReceived: data => service.Logs.OnNext($"[{replica}]: {data}"),
errorDataReceived: data => service.Logs.OnNext($"[{replica}]: {data}"),
throwOnError: false,
cancellationToken: dockerInfo.StoppingTokenSource.Token);
if (logsRes.ExitCode != 0)
if (sentStartedEvent)
{
break;
using var restartCts = new CancellationTokenSource(DockerStopTimeout);
result = await ProcessUtil.RunAsync("docker", $"restart {containerId}", throwOnError: false, cancellationToken: restartCts.Token);
if (restartCts.IsCancellationRequested)
{
_logger.LogWarning($"Failed to restart container after {DockerStopTimeout.Seconds} seconds.", replica, shortContainerId);
break; // implement retry mechanism?
}
else if (result.ExitCode != 0)
{
_logger.LogWarning($"Failed to restart container due to exit code {result.ExitCode}.", replica, shortContainerId);
break;
}
service.ReplicaEvents.OnNext(new ReplicaEvent(ReplicaState.Stopped, status));
}
if (!dockerInfo.StoppingTokenSource.IsCancellationRequested)
using var stoppingCts = new CancellationTokenSource();
status.StoppingTokenSource = stoppingCts;
service.ReplicaEvents.OnNext(new ReplicaEvent(ReplicaState.Started, status));
sentStartedEvent = true;
await using var _ = dockerInfo.StoppingTokenSource.Token.Register(() => status.StoppingTokenSource.Cancel());
_logger.LogInformation("Collecting docker logs for {ContainerName}.", replica);
var backOff = TimeSpan.FromSeconds(5);
while (!status.StoppingTokenSource.Token.IsCancellationRequested)
{
try
var logsRes = await ProcessUtil.RunAsync("docker", $"logs -f {containerId}",
outputDataReceived: data => service.Logs.OnNext($"[{replica}]: {data}"),
errorDataReceived: data => service.Logs.OnNext($"[{replica}]: {data}"),
throwOnError: false,
cancellationToken: status.StoppingTokenSource.Token);
if (logsRes.ExitCode != 0)
{
// Avoid spamming logs if restarts are happening
await Task.Delay(backOff, dockerInfo.StoppingTokenSource.Token);
break;
}
catch (OperationCanceledException)
if (!status.StoppingTokenSource.IsCancellationRequested)
{
break;
try
{
// Avoid spamming logs if restarts are happening
await Task.Delay(backOff, status.StoppingTokenSource.Token);
}
catch (OperationCanceledException)
{
break;
}
}
backOff *= 2;
}
backOff *= 2;
}
_logger.LogInformation("docker logs collection for {ContainerName} complete with exit code {ExitCode}", replica, result.ExitCode);
_logger.LogInformation("docker logs collection for {ContainerName} complete with exit code {ExitCode}", replica, result.ExitCode);
status.StoppingTokenSource = null;
}
// Docker has a tendency to get stuck so we're going to timeout this shutdown process
var timeoutCts = new CancellationTokenSource(DockerStopTimeout);
@ -418,7 +450,10 @@ namespace Microsoft.Tye.Hosting
PrintStdOutAndErr(service, replica, result);
service.ReplicaEvents.OnNext(new ReplicaEvent(ReplicaState.Stopped, status));
if (sentStartedEvent)
{
service.ReplicaEvents.OnNext(new ReplicaEvent(ReplicaState.Stopped, status));
}
_logger.LogInformation("Stopped container {ContainerName} with ID {ContainerId} exited with {ExitCode}", replica, shortContainerId, result.ExitCode);
@ -433,7 +468,7 @@ namespace Microsoft.Tye.Hosting
_logger.LogInformation("Removed container {ContainerName} with ID {ContainerId} exited with {ExitCode}", replica, shortContainerId, result.ExitCode);
service.Replicas.TryRemove(replica, out _);
service.Replicas.TryRemove(replica, out var _);
service.ReplicaEvents.OnNext(new ReplicaEvent(ReplicaState.Removed, status));
};

69
src/Microsoft.Tye.Hosting/HttpProxyService.cs

@ -3,6 +3,7 @@
// See the LICENSE file in the project root for more information.
using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Linq;
using System.Net;
@ -26,9 +27,12 @@ namespace Microsoft.Tye.Hosting
private List<WebApplication> _webApplications = new List<WebApplication>();
private readonly ILogger _logger;
private ConcurrentDictionary<int, bool> _readyPorts;
public HttpProxyService(ILogger logger)
{
_logger = logger;
_readyPorts = new ConcurrentDictionary<int, bool>();
}
public async Task StartAsync(Application application)
@ -98,8 +102,9 @@ namespace Microsoft.Tye.Hosting
_logger.LogInformation("Processing ingress rule: Path:{Path}, Host:{Host}, Service:{Service}", rule.Path, rule.Host, rule.Service);
var targetServiceDescription = target.Description;
RegisterListener(target);
var uris = new List<Uri>();
var uris = new List<(int Port, Uri Uri)>();
// HTTP before HTTPS (this might change once we figure out certs...)
var targetBinding = targetServiceDescription.Bindings.FirstOrDefault(b => b.Protocol == "http") ??
@ -117,23 +122,42 @@ namespace Microsoft.Tye.Hosting
{
var port = targetBinding.ReplicaPorts[i];
var url = $"{targetBinding.Protocol}://localhost:{port}";
uris.Add(new Uri(url));
uris.Add((port, new Uri(url)));
}
_logger.LogInformation("Service {ServiceName} is using {Urls}", targetServiceDescription.Name, string.Join(",", uris.Select(u => u.ToString())));
// The only load balancing strategy here is round robin
long count = 0;
RequestDelegate del = context =>
RequestDelegate del = async context =>
{
var next = (int)(Interlocked.Increment(ref count) % uris.Count);
var uri = new UriBuilder(uris[next])
// we find the first `Ready` port
for (int i = 0; i < uris.Count; i++)
{
if (_readyPorts.ContainsKey(uris[next].Port))
{
break;
}
next = (int)(Interlocked.Increment(ref count) % uris.Count);
}
// if we've looped through all the port and didn't find a single one that is `Ready`, we return HTTP BadGateway
if (!_readyPorts.ContainsKey(uris[next].Port))
{
context.Response.StatusCode = (int)HttpStatusCode.BadGateway;
await context.Response.WriteAsync("Bad gateway");
return;
}
var uri = new UriBuilder(uris[next].Uri)
{
Path = (string)context.Request.RouteValues["path"]
};
return context.ProxyRequest(invoker, uri.Uri);
await context.ProxyRequest(invoker, uri.Uri);
};
IEndpointConventionBuilder conventions = null!;
@ -167,6 +191,14 @@ namespace Microsoft.Tye.Hosting
public async Task StopAsync(Application application)
{
foreach (var service in application.Services.Values)
{
if (service.Items.TryGetValue(typeof(Subscription), out var item) && item is IDisposable disposable)
{
disposable.Dispose();
}
}
foreach (var webApp in _webApplications)
{
try
@ -183,5 +215,32 @@ namespace Microsoft.Tye.Hosting
}
}
}
private void RegisterListener(Service service)
{
if (!service.Items.ContainsKey(typeof(Subscription)))
{
service.Items[typeof(Subscription)] = service.ReplicaEvents.Subscribe(OnReplicaEvent);
}
}
private void OnReplicaEvent(ReplicaEvent replicaEvent)
{
foreach (var binding in replicaEvent.Replica.Bindings)
{
if (replicaEvent.State == ReplicaState.Ready)
{
_readyPorts.TryAdd(binding.Port, true);
}
else
{
_readyPorts.TryRemove(binding.Port, out _);
}
}
}
private class Subscription
{
}
}
}

17
src/Microsoft.Tye.Hosting/Model/HttpProber.cs

@ -0,0 +1,17 @@
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.
// See the LICENSE file in the project root for more information.
using System;
using System.Collections.Generic;
namespace Microsoft.Tye.Hosting.Model
{
public class HttpProber
{
public string Path { get; set; } = default!;
public int? Port { get; set; }
public string? Protocol { get; set; }
public List<KeyValuePair<string, object>> Headers { get; set; } = new List<KeyValuePair<string, object>>();
}
}

18
src/Microsoft.Tye.Hosting/Model/Probe.cs

@ -0,0 +1,18 @@
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.
// See the LICENSE file in the project root for more information.
using System;
namespace Microsoft.Tye.Hosting.Model
{
public class Probe
{
public HttpProber? Http { get; set; }
public TimeSpan InitialDelay { get; set; }
public TimeSpan Period { get; set; }
public TimeSpan Timeout { get; set; }
public int SuccessThreshold { get; set; }
public int FailureThreshold { get; set; }
}
}

13
src/Microsoft.Tye.Hosting/Model/ReplicaBinding.cs

@ -0,0 +1,13 @@
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.
// See the LICENSE file in the project root for more information.
namespace Microsoft.Tye.Hosting.Model
{
public class ReplicaBinding
{
public int Port { get; set; }
public int ExternalPort { get; set; }
public string? Protocol { get; set; }
}
}

2
src/Microsoft.Tye.Hosting/Model/ReplicaState.cs

@ -10,5 +10,7 @@ namespace Microsoft.Tye.Hosting.Model
Added,
Started,
Stopped,
Healthy,
Ready
}
}

8
src/Microsoft.Tye.Hosting/Model/ReplicaStatus.cs

@ -2,8 +2,10 @@
// The .NET Foundation licenses this file to you under the MIT license.
// See the LICENSE file in the project root for more information.
using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Threading;
namespace Microsoft.Tye.Hosting.Model
{
@ -26,5 +28,11 @@ namespace Microsoft.Tye.Hosting.Model
public ConcurrentDictionary<string, string> Metrics { get; set; } = new ConcurrentDictionary<string, string>();
public IDictionary<string, string>? Environment { get; set; }
public ReplicaState? State { get; set; }
public CancellationTokenSource? StoppingTokenSource { get; set; }
public List<ReplicaBinding> Bindings { get; set; } = new List<ReplicaBinding>();
}
}

5
src/Microsoft.Tye.Hosting/Model/Service.cs

@ -24,6 +24,11 @@ namespace Microsoft.Tye.Hosting.Model
CachedLogs.Enqueue(entry);
});
ReplicaEvents.Subscribe(entry =>
{
entry.Replica.State = entry.State;
});
}
public ServiceDescription Description { get; }

2
src/Microsoft.Tye.Hosting/Model/ServiceDescription.cs

@ -20,5 +20,7 @@ namespace Microsoft.Tye.Hosting.Model
public List<ServiceBinding> Bindings { get; } = new List<ServiceBinding>();
public List<EnvironmentVariable> Configuration { get; } = new List<EnvironmentVariable>();
public List<string> Dependencies { get; } = new List<string>();
public Probe? Liveness { get; set; }
public Probe? Readiness { get; set; }
}
}

1
src/Microsoft.Tye.Hosting/Model/V1/V1ReplicaStatus.cs

@ -18,5 +18,6 @@ namespace Microsoft.Tye.Hosting.Model.V1
public int? ExitCode { get; set; }
public int? Pid { get; set; }
public IDictionary<string, string>? Environment { get; set; }
public ReplicaState? State { get; set; }
}
}

2
src/Microsoft.Tye.Hosting/PortAssigner.cs

@ -50,7 +50,7 @@ namespace Microsoft.Tye.Hosting
binding.Port = GetNextPort();
}
if (service.Description.Replicas == 1)
if (service.Description.Readiness == null && service.Description.Replicas == 1)
{
// No need to proxy, the port maps to itself
binding.ReplicaPorts.Add(binding.Port.Value);

9
src/Microsoft.Tye.Hosting/ProcessRunner.cs

@ -238,6 +238,10 @@ namespace Microsoft.Tye.Hosting
var status = new ProcessStatus(service, replica);
service.Replicas[replica] = status;
using var stoppingCts = new CancellationTokenSource();
status.StoppingTokenSource = stoppingCts;
await using var _ = processInfo.StoppedTokenSource.Token.Register(() => status.StoppingTokenSource.Cancel());
service.ReplicaEvents.OnNext(new ReplicaEvent(ReplicaState.Added, status));
// This isn't your host name
@ -250,6 +254,7 @@ namespace Microsoft.Tye.Hosting
if (hasPorts)
{
status.Ports = ports.Select(p => p.Port);
status.Bindings = ports.Select(p => new ReplicaBinding() { Port = p.Port, ExternalPort = p.ExternalPort, Protocol = p.Protocol }).ToList();
}
// TODO clean this up.
@ -291,7 +296,7 @@ namespace Microsoft.Tye.Hosting
service.ReplicaEvents.OnNext(new ReplicaEvent(ReplicaState.Started, status));
},
throwOnError: false,
cancellationToken: processInfo.StoppedTokenSource.Token);
cancellationToken: status.StoppingTokenSource.Token);
status.ExitCode = result.ExitCode;
@ -324,7 +329,7 @@ namespace Microsoft.Tye.Hosting
}
// Remove the replica from the set
service.Replicas.TryRemove(replica, out _);
service.Replicas.TryRemove(replica, out var _);
service.ReplicaEvents.OnNext(new ReplicaEvent(ReplicaState.Removed, status));
}
}

53
src/Microsoft.Tye.Hosting/ProxyService.cs

@ -3,6 +3,7 @@
// See the LICENSE file in the project root for more information.
using System;
using System.Collections.Concurrent;
using System.IO;
using System.IO.Pipelines;
using System.Net;
@ -13,7 +14,6 @@ using System.Threading.Tasks;
using Bedrock.Framework;
using Microsoft.AspNetCore.Connections;
using Microsoft.AspNetCore.Connections.Features;
using Microsoft.AspNetCore.Hosting;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using Microsoft.Tye.Hosting.Model;
@ -25,9 +25,12 @@ namespace Microsoft.Tye.Hosting
private IHost? _host;
private readonly ILogger _logger;
private ConcurrentDictionary<int, CancellationTokenSource> _cancellationsByReplicaPort;
public ProxyService(ILogger logger)
{
_logger = logger;
_cancellationsByReplicaPort = new ConcurrentDictionary<int, CancellationTokenSource>();
}
public Task StartAsync(Application application)
@ -44,6 +47,8 @@ namespace Microsoft.Tye.Hosting
continue;
}
service.Items[typeof(Subscription)] = service.ReplicaEvents.Subscribe(OnReplicaEvent);
foreach (var binding in service.Description.Bindings)
{
if (binding.Port == null)
@ -52,7 +57,7 @@ namespace Microsoft.Tye.Hosting
continue;
}
if (service.Description.Replicas == 1)
if (service.Description.Readiness == null && service.Description.Replicas == 1)
{
// No need to proxy for a single replica, we may want to do this later but right now we skip it
continue;
@ -77,6 +82,15 @@ namespace Microsoft.Tye.Hosting
var next = (int)(Interlocked.Increment(ref count) % ports.Count);
if (!_cancellationsByReplicaPort.TryGetValue(ports[next], out var cts))
{
// replica in ready state <=> it's ports have cancellation tokens in the dictionary
// if replica is not in ready state, we don't forward traffic, but return instead
return;
}
using var _ = cts.Token.Register(() => notificationFeature.RequestClose());
NetworkStream? targetStream = null;
try
@ -134,6 +148,8 @@ namespace Microsoft.Tye.Hosting
{
_logger.LogDebug(0, ex, "Proxy error for service {ServiceName}", service.Description.Name);
}
_logger.LogDebug("Existing proxy {ServiceName} {ExternalPort}:{InternalPort}", service.Description.Name, binding.Port, ports[next]);
}
catch (Exception ex)
{
@ -159,6 +175,14 @@ namespace Microsoft.Tye.Hosting
public async Task StopAsync(Application application)
{
foreach (var service in application.Services.Values)
{
if (service.Items.TryGetValue(typeof(Subscription), out var item) && item is IDisposable disposable)
{
disposable.Dispose();
}
}
if (_host != null)
{
await _host.StopAsync();
@ -169,5 +193,30 @@ namespace Microsoft.Tye.Hosting
}
}
}
private void OnReplicaEvent(ReplicaEvent replicaEvent)
{
// when a replica becomes ready for the first time, it shouldn't have a cancellation token in the dictionary
// for any event other than ready, we want to cancel the token and remove it from the dictionary
foreach (var binding in replicaEvent.Replica.Bindings)
{
if (_cancellationsByReplicaPort.TryRemove(binding.Port, out var cts))
{
cts.Cancel();
}
}
if (replicaEvent.State == ReplicaState.Ready)
{
foreach (var binding in replicaEvent.Replica.Bindings)
{
_cancellationsByReplicaPort.TryAdd(binding.Port, new CancellationTokenSource());
}
}
}
private class Subscription
{
}
}
}

455
src/Microsoft.Tye.Hosting/ReplicaMonitor.cs

@ -0,0 +1,455 @@
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.
// See the LICENSE file in the project root for more information.
using System;
using System.Collections.Concurrent;
using System.Linq;
using System.Net.Http;
using System.Reactive.Subjects;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.Extensions.Logging;
using Microsoft.Tye.Hosting.Model;
namespace Microsoft.Tye.Hosting
{
public class ReplicaMonitor : IApplicationProcessor
{
private ILogger _logger;
private ConcurrentDictionary<string, ReplicaMonitorState> _states;
public ReplicaMonitor(ILogger logger)
{
_logger = logger;
_states = new ConcurrentDictionary<string, ReplicaMonitorState>();
}
public Task StartAsync(Application application)
{
foreach (var service in application.Services.Values)
{
service.Items[typeof(Subscription)] = service.ReplicaEvents.Subscribe(OnReplicaChanged);
}
return Task.CompletedTask;
}
public Task StopAsync(Application application)
{
foreach (var service in application.Services.Values)
{
if (service.Items.TryGetValue(typeof(Subscription), out var item) && item is IDisposable disposable)
{
disposable.Dispose();
}
}
return Task.CompletedTask;
}
private void OnReplicaChanged(ReplicaEvent replicaEvent)
{
switch (replicaEvent.State)
{
case ReplicaState.Started:
_states.TryAdd(replicaEvent.Replica.Name, new ReplicaMonitorState(replicaEvent.Replica, _logger));
break;
case ReplicaState.Stopped:
if (_states.TryRemove(replicaEvent.Replica.Name, out var stateToDispose))
{
stateToDispose.Dispose();
}
break;
default:
if (_states.TryGetValue(replicaEvent.Replica.Name, out var state))
{
state.Update(replicaEvent);
}
break;
}
}
private class Subscription
{
}
private class ReplicaMonitorState : IDisposable
{
private ReplicaStatus _replica;
private ILogger _logger;
private Prober? _livenessProber;
private Prober? _readinessProber;
private IDisposable? _livenessProberObserver;
private IDisposable? _readinessProberObserver;
private ReplicaState _currentState;
private DateTime _lastStateChange;
private object _stateChangeLocker;
public ReplicaMonitorState(ReplicaStatus replica, ILogger logger)
{
_replica = replica;
_logger = logger;
_currentState = ReplicaState.Started;
_lastStateChange = DateTime.Now;
_stateChangeLocker = new object();
Init();
}
private void Init()
{
var serviceDesc = _replica.Service.Description;
if (serviceDesc.Liveness == null && serviceDesc.Readiness == null)
{
MoveToReady();
}
else if (serviceDesc.Liveness == null)
{
MoveToHealthy(from: ReplicaState.Started);
StartReadinessProbe(serviceDesc.Readiness!);
}
else if (serviceDesc.Readiness == null)
{
StartLivenessProbe(serviceDesc.Liveness, moveToOnSuccess: ReplicaState.Ready);
}
else
{
StartLivenessProbe(serviceDesc.Liveness);
StartReadinessProbe(serviceDesc.Readiness);
}
}
private void StartLivenessProbe(Probe probe, ReplicaState moveToOnSuccess = ReplicaState.Healthy)
{
// currently only HTTP is available
if (probe.Http == null)
{
_logger.LogWarning("Cannot start probing replica {name} because probe configuration is not set", _replica.Name);
return;
}
_livenessProber = new HttpProber(_replica, "liveness", probe, probe.Http, _logger);
var failureThreshold = probe.FailureThreshold;
var failures = 0;
var dead = false;
_livenessProberObserver = _livenessProber.ProbeResults.Subscribe(entry =>
{
if (dead)
{
return;
}
if (entry)
{
// Reset failures count on success
failures = 0;
}
(var currentState, _) = ReadCurrentState();
var failuresPastThreshold = failures >= failureThreshold;
switch ((entry, currentState, moveToOnSuccess, failuresPastThreshold))
{
case (false, _, _, true):
dead = true;
Kill();
break;
case (false, _, _, false):
Interlocked.Increment(ref failures);
break;
case (true, ReplicaState.Started, ReplicaState.Ready, _):
case (true, ReplicaState.Healthy, ReplicaState.Ready, _):
MoveToReady();
break;
case (true, ReplicaState.Started, ReplicaState.Healthy, _):
MoveToHealthy(from: ReplicaState.Started);
break;
}
});
_livenessProber.Start();
}
private void StartReadinessProbe(Probe probe)
{
// currently only HTTP is available
if (probe.Http == null)
{
_logger.LogWarning("Cannot start probing replica {name} because probe configuration is not set", _replica.Name);
return;
}
_readinessProber = new HttpProber(_replica, "readiness", probe, probe.Http, _logger);
var successThreshold = probe.SuccessThreshold;
var failureThreshold = probe.FailureThreshold;
var successes = 0;
var failures = 0;
_readinessProberObserver = _readinessProber.ProbeResults.Subscribe(entry =>
{
if (entry)
{
// Reset failures count on success
failures = 0;
}
else
{
// Reset successes count on failure
successes = 0;
}
(var currentState, _) = ReadCurrentState();
var successesPastThreshold = successes >= successThreshold;
var failuresPastThreshold = failures >= failureThreshold;
switch ((entry, currentState, failuresPastThreshold, successesPastThreshold))
{
case (false, ReplicaState.Ready, true, _):
MoveToHealthy(from: ReplicaState.Ready);
break;
case (false, ReplicaState.Ready, false, _):
Interlocked.Increment(ref failures);
break;
case (true, ReplicaState.Healthy, _, true):
MoveToReady();
break;
case (true, ReplicaState.Healthy, _, false):
Interlocked.Increment(ref successes);
break;
}
});
_readinessProber.Start();
}
private void MoveToHealthy(ReplicaState from)
{
_logger.LogDebug("Replica {name} is moving to an healthy state", _replica.Name);
ChangeState(ReplicaState.Healthy);
}
private void MoveToReady()
{
_logger.LogDebug("Replica {name} is moving to a ready state", _replica.Name);
ChangeState(ReplicaState.Ready);
}
private void Kill()
{
_logger.LogDebug("Killing replica {name} because it has failed the liveness probe", _replica.Name);
// it is assumed that a `Started` replica should have an initialized stopping token source
_replica.StoppingTokenSource!.Cancel();
}
private void ChangeState(ReplicaState state)
{
_replica.Service.ReplicaEvents.OnNext(new ReplicaEvent(state, _replica));
lock (_stateChangeLocker)
{
_currentState = state;
_lastStateChange = DateTime.Now;
}
}
private (ReplicaState state, DateTime lastChanged) ReadCurrentState()
{
lock (_stateChangeLocker)
{
return (_currentState, _lastStateChange);
}
}
public void Update(ReplicaEvent replicaEvent)
{
}
public void Dispose()
{
_livenessProber?.Dispose();
_readinessProber?.Dispose();
_livenessProberObserver?.Dispose();
_readinessProberObserver?.Dispose();
}
}
private abstract class Prober : IDisposable
{
protected Prober()
{
ProbeResults = new Subject<bool>();
}
public Subject<bool> ProbeResults { get; }
public abstract void Start();
public abstract void Dispose();
}
private class HttpProber : Prober
{
private static HttpClient _httpClient;
static HttpProber()
{
_httpClient = new HttpClient();
}
private ReplicaStatus _replica;
private ReplicaBinding? _selectedBinding;
private string _probeName;
private Probe _probe;
private Model.HttpProber _httpProberSettings;
private Timer _probeTimer;
private CancellationTokenSource _cts;
private ILogger _logger;
private bool _lastStatus;
public HttpProber(ReplicaStatus replica, string probeName, Probe probe, Model.HttpProber httpProberSettings, ILogger logger)
: base()
{
_replica = replica;
_selectedBinding = null;
_probeName = probeName;
_probe = probe;
_httpProberSettings = httpProberSettings;
_probeTimer = new Timer(DoProbe, null, Timeout.Infinite, Timeout.Infinite);
_cts = new CancellationTokenSource();
_logger = logger;
_lastStatus = true;
}
private void DoProbe(object? state)
{
_probeTimer.Change(Timeout.Infinite, Timeout.Infinite);
_ = DoProbeAsync();
}
private async Task DoProbeAsync()
{
if (_cts.Token.IsCancellationRequested)
{
return;
}
try
{
var protocol = _selectedBinding!.Protocol;
var address = $"{protocol}://localhost:{_selectedBinding.Port}{_httpProberSettings.Path}";
using var timeoutCts = new CancellationTokenSource(_probe.Timeout);
var req = new HttpRequestMessage(HttpMethod.Get, address);
foreach (var header in _httpProberSettings.Headers)
{
req.Headers.Add(header.Key, header.Value.ToString());
}
var res = await _httpClient.SendAsync(req, timeoutCts.Token);
if (!res.IsSuccessStatusCode)
{
ShowWarning($"Replica {_replica.Name} failed http probe at address '{_httpProberSettings.Path}' due to a failed status ({res.StatusCode})");
Send(false);
return;
}
Send(true);
}
catch (HttpRequestException ex)
{
ShowWarning($"Replica {_replica.Name} failed http probe at address '{_httpProberSettings.Path}' due to an http exception", ex);
Send(false);
}
catch (TaskCanceledException)
{
ShowWarning($"Replica {_replica.Name} failed http probe at address '{_httpProberSettings.Path}' due to timeout");
Send(false);
}
finally
{
try
{
_probeTimer.Change(_probe.Period, Timeout.InfiniteTimeSpan);
}
catch (ObjectDisposedException)
{
}
}
}
private void Send(bool status)
{
ProbeResults.OnNext(status);
_lastStatus = status;
}
private void ShowWarning(string message, Exception? ex = null)
{
if (!_lastStatus)
{
return;
}
if (ex != null)
{
_logger.LogWarning(ex, message);
}
else
{
_logger.LogWarning(message);
}
}
public override void Start()
{
// the logic that selects the binding for the probing depends on the protocol and port fields that were provided in the probe
// if neither port nor protocol were provided, we select the first binding
// if just port was provided, we select the first binding with that port
// if just protocol was provided, we select the first binding with that protocol (http/https)
// if both port and protocol were provided, we select the first binding with that port and protocol
Func<ReplicaBinding, bool> bindingClosure = (_httpProberSettings.Port.HasValue, _httpProberSettings.Protocol != null) switch
{
(false, false) => _ => true,
(true, false) => r => r.ExternalPort == _httpProberSettings.Port!.Value,
(false, true) => r => r.Protocol == _httpProberSettings.Protocol!,
(true, true) => r => r.ExternalPort == _httpProberSettings.Port!.Value && r.Protocol == _httpProberSettings.Protocol!
};
var selectedBindings = _replica.Bindings.Where(bindingClosure);
if (selectedBindings.Count() == 0)
{
_logger.LogWarning($"No suitable binding was found for replica {_replica.Name} for probe '{_probeName}'");
return;
}
_selectedBinding = selectedBindings.First();
try
{
_probeTimer.Change(_probe.InitialDelay, Timeout.InfiniteTimeSpan);
}
catch (ObjectDisposedException)
{
}
}
public override void Dispose()
{
_cts.Cancel();
_probeTimer.Dispose();
}
}
}
}

3
src/Microsoft.Tye.Hosting/TyeDashboardApi.cs

@ -171,7 +171,8 @@ namespace Microsoft.Tye.Hosting
{
Name = replica.Value.Name,
Ports = replica.Value.Ports,
Environment = replica.Value.Environment
Environment = replica.Value.Environment,
State = replica.Value.State
};
replicateDictionary[replica.Key] = replicaStatus;

1
src/Microsoft.Tye.Hosting/TyeHost.cs

@ -276,6 +276,7 @@ namespace Microsoft.Tye.Hosting
new ProxyService(logger),
new HttpProxyService(logger),
new DockerImagePuller(logger),
new ReplicaMonitor(logger),
new DockerRunner(logger, replicaRegistry),
new ProcessRunner(logger, replicaRegistry, ProcessRunnerOptions.FromHostOptions(options))
};

28
src/tye/ApplicationBuilderExtensions.cs

@ -36,11 +36,15 @@ namespace Microsoft.Tye
foreach (var service in application.Services)
{
RunInfo? runInfo;
Probe? liveness;
Probe? readiness;
int replicas;
var env = new List<EnvironmentVariable>();
if (service is ExternalServiceBuilder)
{
runInfo = null;
liveness = null;
readiness = null;
replicas = 1;
}
else if (service is ContainerServiceBuilder container)
@ -70,6 +74,8 @@ namespace Microsoft.Tye
runInfo = dockerRunInfo;
replicas = container.Replicas;
liveness = container.Liveness != null ? GetProbeFromBuilder(container.Liveness) : null;
readiness = container.Readiness != null ? GetProbeFromBuilder(container.Readiness) : null;
foreach (var entry in container.EnvironmentVariables)
{
@ -80,6 +86,8 @@ namespace Microsoft.Tye
{
runInfo = new ExecutableRunInfo(executable.Executable, executable.WorkingDirectory, executable.Args);
replicas = executable.Replicas;
liveness = executable.Liveness != null ? GetProbeFromBuilder(executable.Liveness) : null;
readiness = executable.Readiness != null ? GetProbeFromBuilder(executable.Readiness) : null;
foreach (var entry in executable.EnvironmentVariables)
{
@ -107,6 +115,8 @@ namespace Microsoft.Tye
runInfo = projectInfo;
replicas = project.Replicas;
liveness = project.Liveness != null ? GetProbeFromBuilder(project.Liveness) : null;
readiness = project.Readiness != null ? GetProbeFromBuilder(project.Readiness) : null;
foreach (var entry in project.EnvironmentVariables)
{
@ -121,6 +131,8 @@ namespace Microsoft.Tye
var description = new ServiceDescription(service.Name, runInfo)
{
Replicas = replicas,
Liveness = liveness,
Readiness = readiness
};
description.Configuration.AddRange(env);
@ -187,5 +199,21 @@ namespace Microsoft.Tye
return env;
}
private static Probe GetProbeFromBuilder(ProbeBuilder builder) => new Probe()
{
Http = builder.Http != null ? new HttpProber()
{
Path = builder.Http.Path,
Headers = builder.Http.Headers,
Port = builder.Http.Port,
Protocol = builder.Http.Protocol
} : null,
InitialDelay = TimeSpan.FromSeconds(builder.InitialDelay),
Period = TimeSpan.FromSeconds(builder.Period),
Timeout = TimeSpan.FromSeconds(builder.Timeout),
SuccessThreshold = builder.SuccessThreshold,
FailureThreshold = builder.FailureThreshold
};
}
}

453
test/E2ETest/HealthCheckTests.cs

@ -0,0 +1,453 @@
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.
// See the LICENSE file in the project root for more information.
using System;
using System.Collections.Generic;
using System.IO;
using System.Linq;
using System.Net;
using System.Net.Http;
using System.Text.Json;
using System.Text.Json.Serialization;
using System.Threading.Tasks;
using Microsoft.Tye;
using Microsoft.Tye.Hosting;
using Microsoft.Tye.Hosting.Model;
using Test.Infrastructure;
using Xunit;
using Xunit.Abstractions;
using static Test.Infrastructure.TestHelpers;
namespace E2ETest
{
public class HealthCheckTests
{
private readonly ITestOutputHelper _output;
private readonly TestOutputLogEventSink _sink;
private readonly JsonSerializerOptions _options;
private static readonly ReplicaState?[] startedOrHigher = new ReplicaState?[] { ReplicaState.Started, ReplicaState.Healthy, ReplicaState.Ready };
private static readonly ReplicaState?[] stoppedOrLower = new ReplicaState?[] { ReplicaState.Stopped, ReplicaState.Removed };
private static HttpClient _client;
static HealthCheckTests()
{
var handler = new HttpClientHandler
{
ServerCertificateCustomValidationCallback = (a, b, c, d) => true,
AllowAutoRedirect = false
};
_client = new HttpClient(new RetryHandler(handler));
}
public HealthCheckTests(ITestOutputHelper output)
{
_output = output;
_sink = new TestOutputLogEventSink(output);
_options = new JsonSerializerOptions()
{
PropertyNameCaseInsensitive = true,
PropertyNamingPolicy = JsonNamingPolicy.CamelCase,
WriteIndented = true,
};
_options.Converters.Add(new JsonStringEnumConverter(JsonNamingPolicy.CamelCase));
}
[Fact]
public async Task ServiceWithoutLivenessReadinessShouldDefaultToReadyTests()
{
using var projectDirectory = CopyTestProjectDirectory("health-checks");
var projectFile = new FileInfo(Path.Combine(projectDirectory.DirectoryPath, "tye-none.yaml"));
var outputContext = new OutputContext(_sink, Verbosity.Debug);
var application = await ApplicationFactory.CreateAsync(outputContext, projectFile);
await using var host = new TyeHost(application.ToHostingApplication(), new HostOptions())
{
Sink = _sink,
};
await StartHostAndWaitForReplicasToStart(host, new[] { "health-none" }, ReplicaState.Ready);
}
[Fact]
public async Task ServiceWithoutLivenessShouldDefaultToHealthyTests()
{
using var projectDirectory = CopyTestProjectDirectory("health-checks");
var projectFile = new FileInfo(Path.Combine(projectDirectory.DirectoryPath, "tye-readiness.yaml"));
var outputContext = new OutputContext(_sink, Verbosity.Debug);
var application = await ApplicationFactory.CreateAsync(outputContext, projectFile);
await using var host = new TyeHost(application.ToHostingApplication(), new HostOptions())
{
Sink = _sink,
};
await StartHostAndWaitForReplicasToStart(host, new[] { "health-readiness" }, ReplicaState.Healthy);
}
[Fact]
public async Task ServicWithoutLivenessShouldBecomeReadyWhenReadyTests()
{
using var projectDirectory = CopyTestProjectDirectory("health-checks");
var projectFile = new FileInfo(Path.Combine(projectDirectory.DirectoryPath, "tye-readiness.yaml"));
var outputContext = new OutputContext(_sink, Verbosity.Debug);
var application = await ApplicationFactory.CreateAsync(outputContext, projectFile);
await using var host = new TyeHost(application.ToHostingApplication(), new HostOptions())
{
Sink = _sink,
};
await StartHostAndWaitForReplicasToStart(host, new[] { "health-readiness" }, ReplicaState.Healthy);
var replicas = host.Application.Services["health-readiness"].Replicas.Select(r => r.Value).ToList();
Assert.True(await DoOperationAndWaitForReplicasToChangeState(host, ReplicaState.Ready, replicas.Count, replicas.Select(r => r.Name).ToHashSet(), null, TimeSpan.Zero, async _ =>
{
await Task.WhenAll(replicas.Select(r => SetHealthyReadyInReplica(r, ready: true)));
}));
}
[Fact]
public async Task ServiceWithoutReadinessShouldDefaultToReadyWhenHealthyTests()
{
using var projectDirectory = CopyTestProjectDirectory("health-checks");
var projectFile = new FileInfo(Path.Combine(projectDirectory.DirectoryPath, "tye-liveness.yaml"));
var outputContext = new OutputContext(_sink, Verbosity.Debug);
var application = await ApplicationFactory.CreateAsync(outputContext, projectFile);
await using var host = new TyeHost(application.ToHostingApplication(), new HostOptions())
{
Sink = _sink,
};
await StartHostAndWaitForReplicasToStart(host, new[] { "health-liveness" }, ReplicaState.Started);
var replicasToBecomeReady = host.Application.Services["health-liveness"].Replicas.Select(r => r.Value);
var replicasNamesToBecomeReady = replicasToBecomeReady.Select(r => r.Name).ToHashSet();
Assert.True(await DoOperationAndWaitForReplicasToChangeState(host, ReplicaState.Ready, replicasNamesToBecomeReady.Count, replicasNamesToBecomeReady, null, TimeSpan.Zero, async _ =>
{
await Task.WhenAll(replicasToBecomeReady.Select(r => SetHealthyReadyInReplica(r, healthy: true)));
}));
}
[Fact]
public async Task ReadyServiceShouldBecomeHealthyWhenReadinessFailsTests()
{
using var projectDirectory = CopyTestProjectDirectory("health-checks");
var projectFile = new FileInfo(Path.Combine(projectDirectory.DirectoryPath, "tye-all.yaml"));
var outputContext = new OutputContext(_sink, Verbosity.Debug);
var application = await ApplicationFactory.CreateAsync(outputContext, projectFile);
await using var host = new TyeHost(application.ToHostingApplication(), new HostOptions())
{
Sink = _sink,
};
SetReplicasInitialState(host, true, true);
await StartHostAndWaitForReplicasToStart(host, new[] { "health-all" }, ReplicaState.Ready);
var replicasToBecomeReady = host.Application.Services["health-all"].Replicas.Select(r => r.Value).ToList();
var replicasNamesToBecomeReady = replicasToBecomeReady.Select(r => r.Name).ToHashSet();
var randomReplica = replicasToBecomeReady[new Random().Next(0, replicasToBecomeReady.Count)];
Assert.True(await DoOperationAndWaitForReplicasToChangeState(host, ReplicaState.Healthy, 1, new[] { randomReplica.Name }.ToHashSet(), replicasNamesToBecomeReady.Where(r => r != randomReplica.Name).ToHashSet(), TimeSpan.FromSeconds(1), async _ =>
{
await SetHealthyReadyInReplica(randomReplica, ready: false);
}));
}
[Fact]
public async Task ReadyServiceShouldRestartWhenLivenessFailsTests()
{
using var projectDirectory = CopyTestProjectDirectory("health-checks");
var projectFile = new FileInfo(Path.Combine(projectDirectory.DirectoryPath, "tye-all.yaml"));
var outputContext = new OutputContext(_sink, Verbosity.Debug);
var application = await ApplicationFactory.CreateAsync(outputContext, projectFile);
await using var host = new TyeHost(application.ToHostingApplication(), new HostOptions())
{
Sink = _sink,
};
SetReplicasInitialState(host, true, true);
await StartHostAndWaitForReplicasToStart(host, new[] { "health-all" }, ReplicaState.Ready);
var replicasToBecomeReady = host.Application.Services["health-all"].Replicas.Select(r => r.Value).ToList();
var replicasNamesToBecomeReady = replicasToBecomeReady.Select(r => r.Name).ToHashSet();
var randomReplica = replicasToBecomeReady[new Random().Next(0, replicasToBecomeReady.Count)];
Assert.True(await DoOperationAndWaitForReplicasToRestart(host, new[] { randomReplica.Name }.ToHashSet(), replicasNamesToBecomeReady.Where(r => r != randomReplica.Name).ToHashSet(), TimeSpan.FromSeconds(1), async _ =>
{
await SetHealthyReadyInReplica(randomReplica, healthy: false);
}));
}
[Fact]
public async Task ProbeShouldRespectTimeoutTests()
{
using var projectDirectory = CopyTestProjectDirectory("health-checks");
var projectFile = new FileInfo(Path.Combine(projectDirectory.DirectoryPath, "tye-all.yaml"));
var outputContext = new OutputContext(_sink, Verbosity.Debug);
var application = await ApplicationFactory.CreateAsync(outputContext, projectFile);
await using var host = new TyeHost(application.ToHostingApplication(), new HostOptions())
{
Sink = _sink,
};
SetReplicasInitialState(host, true, true);
await StartHostAndWaitForReplicasToStart(host, new[] { "health-all" }, ReplicaState.Ready);
var replicasToBecomeReady = host.Application.Services["health-all"].Replicas.Select(r => r.Value).ToList();
var replicasNamesToBecomeReady = replicasToBecomeReady.Select(r => r.Name).ToHashSet();
var randomNumber = new Random().Next(0, replicasToBecomeReady.Count);
var randomReplica1 = replicasToBecomeReady[randomNumber];
var randomReplica2 = replicasToBecomeReady[(randomNumber + 1) % replicasToBecomeReady.Count];
Assert.True(await DoOperationAndWaitForReplicasToChangeState(host, ReplicaState.Healthy, 1, new[] { randomReplica1.Name }.ToHashSet(), replicasNamesToBecomeReady.Where(r => r != randomReplica1.Name).ToHashSet(), TimeSpan.FromSeconds(2), async _ =>
{
await Task.WhenAll(new[]
{
SetHealthyReadyInReplica(randomReplica1, readyDelay: 2),
SetHealthyReadyInReplica(randomReplica2, readyDelay: 1)
});
}));
}
[Fact]
public async Task ProxyShouldNotProxyToNonReadyReplicasTests()
{
using var projectDirectory = CopyTestProjectDirectory("health-checks");
var projectFile = new FileInfo(Path.Combine(projectDirectory.DirectoryPath, "tye-proxy.yaml"));
var outputContext = new OutputContext(_sink, Verbosity.Debug);
var application = await ApplicationFactory.CreateAsync(outputContext, projectFile);
await using var host = new TyeHost(application.ToHostingApplication(), new HostOptions())
{
Sink = _sink,
};
SetReplicasInitialState(host, true, true);
await StartHostAndWaitForReplicasToStart(host, new[] { "health-proxy" }, ReplicaState.Ready);
var replicasToBecomeReady = host.Application.Services["health-proxy"].Replicas.Select(r => r.Value).ToList();
// we assume that proxy will continue sending http request to the same replica
var randomReplicaPortRes1 = await _client.GetAsync($"http://localhost:{host.Application.Services["health-proxy"].Description.Bindings.First().Port}/ports");
var randomReplicaPort1 = JsonSerializer.Deserialize<int[]>(await randomReplicaPortRes1.Content.ReadAsStringAsync())[0];
var randomReplica1 = replicasToBecomeReady.First(r => r.Bindings.Any(b => b.Port == randomReplicaPort1));
await DoOperationAndWaitForReplicasToChangeState(host, ReplicaState.Healthy, 1, new[] { randomReplica1.Name }.ToHashSet(), null, TimeSpan.Zero, async _ =>
{
await SetHealthyReadyInReplica(randomReplica1, ready: false);
});
var randomReplicaPortRes2 = await _client.GetAsync($"http://localhost:{host.Application.Services["health-proxy"].Description.Bindings.First().Port}/ports");
var randomReplicaPort2 = JsonSerializer.Deserialize<int[]>(await randomReplicaPortRes2.Content.ReadAsStringAsync())[0];
var randomReplica2 = replicasToBecomeReady.First(r => r.Bindings.Any(b => b.Port == randomReplicaPort2));
Assert.NotEqual(randomReplicaPort1, randomReplicaPort2);
await DoOperationAndWaitForReplicasToChangeState(host, ReplicaState.Healthy, 1, new[] { randomReplica2.Name }.ToHashSet(), null, TimeSpan.Zero, async _ =>
{
await SetHealthyReadyInReplica(randomReplica2, ready: false);
});
try
{
var resShouldFail = await _client.GetAsync($"http://localhost:{host.Application.Services["health-proxy"].Description.Bindings.First().Port}/ports");
Assert.False(resShouldFail.IsSuccessStatusCode);
}
catch (HttpRequestException)
{
}
await DoOperationAndWaitForReplicasToChangeState(host, ReplicaState.Ready, 1, new[] { randomReplica2.Name }.ToHashSet(), null, TimeSpan.Zero, async _ =>
{
await SetHealthyReadyInReplica(randomReplica2, ready: true);
});
var randomReplicaPortRes3 = await _client.GetAsync($"http://localhost:{host.Application.Services["health-proxy"].Description.Bindings.First().Port}/ports");
var randomReplicaPort3 = JsonSerializer.Deserialize<int[]>(await randomReplicaPortRes3.Content.ReadAsStringAsync())[0];
Assert.Equal(randomReplicaPort3, randomReplicaPort2);
}
[Fact]
public async Task IngressShouldNotProxyToNonReadyReplicasTests()
{
using var projectDirectory = CopyTestProjectDirectory("health-checks");
var projectFile = new FileInfo(Path.Combine(projectDirectory.DirectoryPath, "tye-ingress.yaml"));
var outputContext = new OutputContext(_sink, Verbosity.Debug);
var application = await ApplicationFactory.CreateAsync(outputContext, projectFile);
await using var host = new TyeHost(application.ToHostingApplication(), new HostOptions())
{
Sink = _sink,
};
SetReplicasInitialState(host, true, true);
await StartHostAndWaitForReplicasToStart(host, new[] { "health-ingress-svc" }, ReplicaState.Ready);
var replicasToBecomeReady = host.Application.Services["health-ingress-svc"].Replicas.Select(r => r.Value).ToList();
var ingressBinding = host.Application.Services.First(s => s.Value.Description.RunInfo is IngressRunInfo).Value.Description.Bindings.First();
var uniqueIdUrl = $"{ingressBinding.Protocol}://localhost:{ingressBinding.Port}/api/id";
var uniqueIds = await ProbeNumberOfUniqueReplicas(uniqueIdUrl);
Assert.Equal(2, uniqueIds);
var firstReplica = replicasToBecomeReady.First();
var secondReplica = replicasToBecomeReady.Skip(1).First();
await DoOperationAndWaitForReplicasToChangeState(host, ReplicaState.Healthy, 1, new[] { firstReplica.Name }.ToHashSet(), null, TimeSpan.Zero, async _ =>
{
await SetHealthyReadyInReplica(firstReplica, ready: false);
});
uniqueIds = await ProbeNumberOfUniqueReplicas(uniqueIdUrl);
Assert.Equal(1, uniqueIds);
await DoOperationAndWaitForReplicasToChangeState(host, ReplicaState.Healthy, 1, new[] { secondReplica.Name }.ToHashSet(), null, TimeSpan.Zero, async _ =>
{
await SetHealthyReadyInReplica(secondReplica, ready: false);
});
var res = await _client.GetAsync(uniqueIdUrl);
Assert.Equal(HttpStatusCode.BadGateway, res.StatusCode);
await DoOperationAndWaitForReplicasToChangeState(host, ReplicaState.Ready, 2, new[] { firstReplica.Name, secondReplica.Name }.ToHashSet(), null, TimeSpan.Zero, async _ =>
{
await SetHealthyReadyInReplica(firstReplica, ready: true);
await SetHealthyReadyInReplica(secondReplica, ready: true);
});
uniqueIds = await ProbeNumberOfUniqueReplicas(uniqueIdUrl);
Assert.Equal(2, uniqueIds);
}
[Fact]
public async Task HeadersTests()
{
using var projectDirectory = CopyTestProjectDirectory("health-checks");
var projectFile = new FileInfo(Path.Combine(projectDirectory.DirectoryPath, "tye-all.yaml"));
var outputContext = new OutputContext(_sink, Verbosity.Debug);
var application = await ApplicationFactory.CreateAsync(outputContext, projectFile);
await using var host = new TyeHost(application.ToHostingApplication(), new HostOptions())
{
Sink = _sink,
};
SetReplicasInitialState(host, true, true);
await StartHostAndWaitForReplicasToStart(host, new[] { "health-all" }, ReplicaState.Ready);
var res = await _client.GetAsync($"http://localhost:{host.Application.Services["health-all"].Description.Bindings.First().Port}/livenessHeaders");
Assert.True(res.IsSuccessStatusCode);
var headers = JsonSerializer.Deserialize<Dictionary<string, string>>(await res.Content.ReadAsStringAsync());
Assert.Equal("value1", headers["name1"]);
Assert.Equal("value2", headers["name2"]);
res = await _client.GetAsync($"http://localhost:{host.Application.Services["health-all"].Description.Bindings.First().Port}/readinessHeaders");
Assert.True(res.IsSuccessStatusCode);
headers = JsonSerializer.Deserialize<Dictionary<string, string>>(await res.Content.ReadAsStringAsync());
Assert.Equal("value3", headers["name3"]);
Assert.Equal("value4", headers["name4"]);
}
private async Task SetHealthyReadyInReplica(ReplicaStatus replica, bool? healthy = null, bool? ready = null, int? healthyDelay = null, int? readyDelay = null)
{
var query = new List<string>();
if (healthy.HasValue)
{
query.Add("healthy=" + healthy);
}
if (ready.HasValue)
{
query.Add("ready=" + ready);
}
if (healthyDelay.HasValue)
{
query.Add("healthyDelay=" + healthyDelay.Value);
}
if (readyDelay.HasValue)
{
query.Add("readyDelay=" + readyDelay.Value);
}
await _client.GetAsync($"http://localhost:{replica.Ports.First()}/set?" + string.Join("&", query));
}
private async Task<int> ProbeNumberOfUniqueReplicas(string url)
{
// this assumes roundrobin
var unique = new HashSet<string>();
string? id = null;
while (id == null || unique.Add(id))
{
var res = await _client.GetAsync(url);
id = await res.Content.ReadAsStringAsync();
}
return unique.Count;
}
private void SetReplicasInitialState(TyeHost host, bool? healthy, bool? ready, string[]? services = null)
{
if (services == null)
{
services = host.Application.Services.Select(s => s.Key).ToArray();
}
else
{
if (services.Any(s => !host.Application.Services.ContainsKey(s)))
{
throw new ArgumentException($"not all services given in {nameof(services)} exist");
}
}
foreach (var service in services)
{
if (healthy.HasValue)
{
host.Application.Services[service].Description.Configuration.Add(new EnvironmentVariable("healthy") { Value = "true" });
}
if (ready.HasValue)
{
host.Application.Services[service].Description.Configuration.Add(new EnvironmentVariable("ready") { Value = "true" });
}
}
}
}
}

3
test/E2ETest/Microsoft.Tye.E2ETests.csproj

@ -1,3 +1,4 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
@ -30,4 +31,4 @@
<Compile Remove="testassets\**\*" />
</ItemGroup>
</Project>
</Project>

143
test/E2ETest/ReplicaStoppingTests.cs

@ -0,0 +1,143 @@
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.
// See the LICENSE file in the project root for more information.
using System;
using System.IO;
using System.Linq;
using System.Net.Http;
using System.Text.Json;
using System.Text.Json.Serialization;
using System.Threading.Tasks;
using Microsoft.Tye;
using Microsoft.Tye.Hosting;
using Microsoft.Tye.Hosting.Model;
using Microsoft.Tye.Hosting.Model.V1;
using Test.Infrastructure;
using Xunit;
using Xunit.Abstractions;
using static Test.Infrastructure.TestHelpers;
namespace E2ETest
{
public class ReplicaStoppingTests
{
private readonly ITestOutputHelper _output;
private readonly TestOutputLogEventSink _sink;
private readonly JsonSerializerOptions _options;
private static readonly ReplicaState?[] startedOrHigher = new ReplicaState?[] { ReplicaState.Started, ReplicaState.Healthy, ReplicaState.Ready };
private static readonly ReplicaState?[] stoppedOrLower = new ReplicaState?[] { ReplicaState.Stopped, ReplicaState.Removed };
public ReplicaStoppingTests(ITestOutputHelper output)
{
_output = output;
_sink = new TestOutputLogEventSink(output);
_options = new JsonSerializerOptions()
{
PropertyNameCaseInsensitive = true,
PropertyNamingPolicy = JsonNamingPolicy.CamelCase,
WriteIndented = true,
};
_options.Converters.Add(new JsonStringEnumConverter(JsonNamingPolicy.CamelCase));
}
[Fact]
public async Task MultiProjectStoppingTests()
{
using var projectDirectory = CopyTestProjectDirectory("health-checks");
var projectFile = new FileInfo(Path.Combine(projectDirectory.DirectoryPath, "tye-none.yaml"));
var outputContext = new OutputContext(_sink, Verbosity.Debug);
var application = await ApplicationFactory.CreateAsync(outputContext, projectFile);
await RunHostingApplication(application, new HostOptions(), async (host, uri) =>
{
var replicaToStop = host.Application.Services["health-none"].Replicas.First();
Assert.Contains(replicaToStop.Value.State, startedOrHigher);
var replicasToRestart = new[] { replicaToStop.Key };
var restOfReplicas = host.Application.Services.SelectMany(s => s.Value.Replicas).Select(r => r.Value.Name).Where(r => r != replicaToStop.Key).ToArray();
Assert.True(await DoOperationAndWaitForReplicasToRestart(host, replicasToRestart.ToHashSet(), restOfReplicas.ToHashSet(), TimeSpan.FromSeconds(1), _ =>
{
replicaToStop.Value.StoppingTokenSource!.Cancel();
return Task.CompletedTask;
}));
Assert.Contains(replicaToStop.Value.State, stoppedOrLower);
Assert.True(host.Application.Services.SelectMany(s => s.Value.Replicas).All(r => startedOrHigher.Contains(r.Value.State)));
});
}
[ConditionalFact]
[SkipIfDockerNotRunning]
public async Task MultiProjectDockerStoppingTests()
{
using var projectDirectory = CopyTestProjectDirectory("multi-project");
var projectFile = new FileInfo(Path.Combine(projectDirectory.DirectoryPath, "tye.yaml"));
var outputContext = new OutputContext(_sink, Verbosity.Debug);
var application = await ApplicationFactory.CreateAsync(outputContext, projectFile);
await RunHostingApplication(application, new HostOptions() { Docker = true }, async (host, uri) =>
{
var replicaToStop = host.Application.Services["frontend"].Replicas.First();
Assert.Contains(replicaToStop.Value.State, startedOrHigher);
var replicasToRestart = new[] { replicaToStop.Key };
var restOfReplicas = host.Application.Services.SelectMany(s => s.Value.Replicas).Select(r => r.Value.Name).Where(r => r != replicaToStop.Key).ToArray();
Assert.True(await DoOperationAndWaitForReplicasToRestart(host, replicasToRestart.ToHashSet(), restOfReplicas.ToHashSet(), TimeSpan.FromSeconds(1), _ =>
{
replicaToStop.Value.StoppingTokenSource!.Cancel();
return Task.CompletedTask;
}));
Assert.Contains(replicaToStop.Value.State, startedOrHigher); // when a container restarts, it assumes the same replica
Assert.True(host.Application.Services.SelectMany(s => s.Value.Replicas).All(r => startedOrHigher.Contains(r.Value.State)));
});
}
private async Task RunHostingApplication(ApplicationBuilder application, HostOptions options, Func<TyeHost, Uri, Task> execute)
{
await using var host = new TyeHost(application.ToHostingApplication(), options)
{
Sink = _sink,
};
try
{
await StartHostAndWaitForReplicasToStart(host);
var uri = new Uri(host.DashboardWebApplication!.Addresses.First());
await execute(host, uri!);
}
finally
{
if (host.DashboardWebApplication != null)
{
var uri = new Uri(host.DashboardWebApplication!.Addresses.First());
using var client = new HttpClient();
foreach (var s in host.Application.Services.Values)
{
var logs = await client.GetStringAsync(new Uri(uri, $"/api/v1/logs/{s.Description.Name}"));
_output.WriteLine($"Logs for service: {s.Description.Name}");
_output.WriteLine(logs);
var description = await client.GetStringAsync(new Uri(uri, $"/api/v1/services/{s.Description.Name}"));
_output.WriteLine($"Service defintion: {s.Description.Name}");
_output.WriteLine(description);
}
}
}
}
}
}

35
test/E2ETest/TyeGenerateTests.cs

@ -392,5 +392,40 @@ namespace E2ETest
await DockerAssert.DeleteDockerImagesAsync(output, "appb");
}
}
[ConditionalFact]
[SkipIfDockerNotRunning]
public async Task Generate_HealthChecks()
{
var applicationName = "health-checks";
var environment = "production";
var projectName = "health-all";
await DockerAssert.DeleteDockerImagesAsync(output, projectName);
using var projectDirectory = TestHelpers.CopyTestProjectDirectory(applicationName);
var projectFile = new FileInfo(Path.Combine(projectDirectory.DirectoryPath, "tye-all.yaml"));
var outputContext = new OutputContext(sink, Verbosity.Debug);
var application = await ApplicationFactory.CreateAsync(outputContext, projectFile);
try
{
await GenerateHost.ExecuteGenerateAsync(outputContext, application, environment, interactive: false);
// name of application is the folder
var content = await File.ReadAllTextAsync(Path.Combine(projectDirectory.DirectoryPath, $"{applicationName}-generate-{environment}.yaml"));
var expectedContent = await File.ReadAllTextAsync($"testassets/generate/{applicationName}.yaml");
YamlAssert.Equals(expectedContent, content, output);
await DockerAssert.AssertImageExistsAsync(output, projectName);
}
finally
{
await DockerAssert.DeleteDockerImagesAsync(output, projectName);
}
}
}
}

76
test/E2ETest/testassets/generate/health-checks.yaml

@ -0,0 +1,76 @@
kind: Deployment
apiVersion: apps/v1
metadata:
name: health-all
labels:
app.kubernetes.io/name: 'health-all'
app.kubernetes.io/part-of: 'health-checks'
spec:
replicas: 3
selector:
matchLabels:
app.kubernetes.io/name: health-all
template:
metadata:
labels:
app.kubernetes.io/name: 'health-all'
app.kubernetes.io/part-of: 'health-checks'
spec:
containers:
- name: health-all
image: health-all:1.0.0
imagePullPolicy: Always
env:
- name: ASPNETCORE_URLS
value: 'http://*:8004'
- name: PORT
value: '8004'
ports:
- containerPort: 8004
livenessProbe:
httpGet:
path: /healthy
port: 8004
scheme: HTTP
httpHeaders:
- name: name1
value: value1
- name: name2
value: value2
initialDelaySeconds: 5
periodSeconds: 1
successThreshold: 1
failureThreshold: 1
readinessProbe:
httpGet:
path: /ready
port: 8004
scheme: HTTP
httpHeaders:
- name: name3
value: value3
- name: name4
value: value4
initialDelaySeconds: 5
periodSeconds: 1
successThreshold: 1
failureThreshold: 1
...
---
kind: Service
apiVersion: v1
metadata:
name: health-all
labels:
app.kubernetes.io/name: 'health-all'
app.kubernetes.io/part-of: 'health-checks'
spec:
selector:
app.kubernetes.io/name: health-all
type: ClusterIP
ports:
- name: http
protocol: TCP
port: 8004
targetPort: 8004
...

32
test/E2ETest/testassets/projects/health-checks/api/Program.cs

@ -0,0 +1,32 @@
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.
// See the LICENSE file in the project root for more information.
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading.Tasks;
using Microsoft.AspNetCore.Hosting;
using Microsoft.AspNetCore.Server.Kestrel.Core;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
namespace api
{
public class Program
{
public static void Main(string[] args)
{
CreateHostBuilder(args).Build().Run();
}
public static IHostBuilder CreateHostBuilder(string[] args) =>
Host.CreateDefaultBuilder(args)
.ConfigureWebHostDefaults(web =>
{
web.UseStartup<Startup>()
.ConfigureKestrel(options => {});
});
}
}

30
test/E2ETest/testassets/projects/health-checks/api/Properties/launchSettings.json

@ -0,0 +1,30 @@
{
"$schema": "http://json.schemastore.org/launchsettings.json",
"iisSettings": {
"windowsAuthentication": false,
"anonymousAuthentication": true,
"iisExpress": {
"applicationUrl": "http://localhost:39065",
"sslPort": 44337
}
},
"profiles": {
"IIS Express": {
"commandName": "IISExpress",
"launchBrowser": true,
"launchUrl": "weatherforecast",
"environmentVariables": {
"ASPNETCORE_ENVIRONMENT": "Development"
}
},
"api": {
"commandName": "Project",
"launchBrowser": true,
"launchUrl": "weatherforecast",
"applicationUrl": "https://localhost:5001;http://localhost:5000",
"environmentVariables": {
"ASPNETCORE_ENVIRONMENT": "Development"
}
}
}
}

166
test/E2ETest/testassets/projects/health-checks/api/Startup.cs

@ -0,0 +1,166 @@
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.
// See the LICENSE file in the project root for more information.
using System;
using System.Collections.Generic;
using System.Linq;
using System.Text;
using System.Text.Json;
using System.Threading.Tasks;
using Microsoft.AspNetCore.Builder;
using Microsoft.AspNetCore.Hosting;
using Microsoft.AspNetCore.Http;
using Microsoft.AspNetCore.HttpsPolicy;
using Microsoft.AspNetCore.Mvc;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
namespace api
{
public class Startup
{
private static string _randomId = Guid.NewGuid().ToString();
private static bool _healthy = false;
private static bool _ready = false;
private static int _healthyDelay = 0;
private static int _readyDelay = 0;
private static Dictionary<string, string> _livenessHeaders;
private static Dictionary<string, string> _readinessHeaders;
private static int[] _ports;
private static object _locker = new object();
public Startup(IConfiguration configuration)
{
Configuration = configuration;
var healthyEnv = Environment.GetEnvironmentVariable("healthy");
var readyEnv = Environment.GetEnvironmentVariable("ready");
if (!string.IsNullOrEmpty(Environment.GetEnvironmentVariable("PORT")))
{
var portParts = Environment.GetEnvironmentVariable("PORT").Split(';');
_ports = portParts.Select(p => int.Parse(p)).ToArray();
}
if (healthyEnv != null)
{
_healthy = bool.Parse(healthyEnv);
}
if (readyEnv != null)
{
_ready = bool.Parse(readyEnv);
}
}
public IConfiguration Configuration { get; }
// This method gets called by the runtime. Use this method to add services to the container.
public void ConfigureServices(IServiceCollection services)
{
services.AddControllers();
}
// This method gets called by the runtime. Use this method to configure the HTTP request pipeline.
public void Configure(IApplicationBuilder app, IWebHostEnvironment env)
{
if (env.IsDevelopment())
{
app.UseDeveloperExceptionPage();
}
app.UseRouting();
app.UseAuthorization();
app.UseEndpoints(endpoints =>
{
endpoints.MapGet("/", async ctx =>
{
await ctx.Response.WriteAsync("Hello");
});
endpoints.MapGet("/ports", async ctx =>
{
await ctx.Response.WriteAsync(JsonSerializer.Serialize(_ports));
});
endpoints.MapGet("/id", async ctx =>
{
await ctx.Response.WriteAsync(_randomId);
});
endpoints.MapGet("/healthy", async ctx =>
{
if (_healthyDelay != 0)
{
await Task.Delay(TimeSpan.FromSeconds(_healthyDelay));
}
_livenessHeaders = ctx.Request.Headers.ToDictionary(h => h.Key, h => h.Value.ToString());
ctx.Response.StatusCode = _healthy ? 200 : 500;
await ctx.Response.WriteAsync(ctx.Response.StatusCode.ToString());
});
endpoints.MapGet("/ready", async ctx =>
{
if (_readyDelay != 0)
{
await Task.Delay(TimeSpan.FromSeconds(_readyDelay));
}
_readinessHeaders = ctx.Request.Headers.ToDictionary(h => h.Key, h => h.Value.ToString());
ctx.Response.StatusCode = _ready ? 200 : 500;
await ctx.Response.WriteAsync(ctx.Response.StatusCode.ToString());
});
// Should be technically POST/PUT, but it's just for tests...
endpoints.MapGet("/set", async ctx =>
{
var query = ctx.Request.Query.ToDictionary(kv => kv.Key).ToDictionary(kv => kv.Key, kv => kv.Value.Value.First());
if (query.ContainsKey("healthy"))
{
_healthy = bool.Parse(query["healthy"]);
}
if (query.ContainsKey("ready"))
{
_ready = bool.Parse(query["ready"]);
}
if (query.ContainsKey("healthyDelay"))
{
_healthyDelay = int.Parse(query["healthyDelay"]);
}
if (query.ContainsKey("readyDelay"))
{
_readyDelay = int.Parse(query["readyDelay"]);
}
await ctx.Response.WriteAsync(_randomId);
});
endpoints.MapGet("/livenessHeaders", async ctx =>
{
await ctx.Response.WriteAsync(JsonSerializer.Serialize(_livenessHeaders));
});
endpoints.MapGet("/readinessHeaders", async ctx =>
{
await ctx.Response.WriteAsync(JsonSerializer.Serialize(_readinessHeaders));
});
});
}
}
}

8
test/E2ETest/testassets/projects/health-checks/api/api.csproj

@ -0,0 +1,8 @@
<Project Sdk="Microsoft.NET.Sdk.Web">
<PropertyGroup>
<TargetFramework>netcoreapp3.1</TargetFramework>
</PropertyGroup>
</Project>

9
test/E2ETest/testassets/projects/health-checks/api/appsettings.Development.json

@ -0,0 +1,9 @@
{
"Logging": {
"LogLevel": {
"Default": "Information",
"Microsoft": "Warning",
"Microsoft.Hosting.Lifetime": "Information"
}
}
}

10
test/E2ETest/testassets/projects/health-checks/api/appsettings.json

@ -0,0 +1,10 @@
{
"Logging": {
"LogLevel": {
"Default": "Information",
"Microsoft": "Warning",
"Microsoft.Hosting.Lifetime": "Information"
}
},
"AllowedHosts": "*"
}

34
test/E2ETest/testassets/projects/health-checks/health-checks.sln

@ -0,0 +1,34 @@

Microsoft Visual Studio Solution File, Format Version 12.00
# Visual Studio 15
VisualStudioVersion = 15.0.26124.0
MinimumVisualStudioVersion = 15.0.26124.0
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "api", "api\api.csproj", "{239E0475-3113-4C1C-A517-38AD9188B65D}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
Debug|x64 = Debug|x64
Debug|x86 = Debug|x86
Release|Any CPU = Release|Any CPU
Release|x64 = Release|x64
Release|x86 = Release|x86
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
EndGlobalSection
GlobalSection(ProjectConfigurationPlatforms) = postSolution
{239E0475-3113-4C1C-A517-38AD9188B65D}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{239E0475-3113-4C1C-A517-38AD9188B65D}.Debug|Any CPU.Build.0 = Debug|Any CPU
{239E0475-3113-4C1C-A517-38AD9188B65D}.Debug|x64.ActiveCfg = Debug|Any CPU
{239E0475-3113-4C1C-A517-38AD9188B65D}.Debug|x64.Build.0 = Debug|Any CPU
{239E0475-3113-4C1C-A517-38AD9188B65D}.Debug|x86.ActiveCfg = Debug|Any CPU
{239E0475-3113-4C1C-A517-38AD9188B65D}.Debug|x86.Build.0 = Debug|Any CPU
{239E0475-3113-4C1C-A517-38AD9188B65D}.Release|Any CPU.ActiveCfg = Release|Any CPU
{239E0475-3113-4C1C-A517-38AD9188B65D}.Release|Any CPU.Build.0 = Release|Any CPU
{239E0475-3113-4C1C-A517-38AD9188B65D}.Release|x64.ActiveCfg = Release|Any CPU
{239E0475-3113-4C1C-A517-38AD9188B65D}.Release|x64.Build.0 = Release|Any CPU
{239E0475-3113-4C1C-A517-38AD9188B65D}.Release|x86.ActiveCfg = Release|Any CPU
{239E0475-3113-4C1C-A517-38AD9188B65D}.Release|x86.Build.0 = Release|Any CPU
EndGlobalSection
EndGlobal

36
test/E2ETest/testassets/projects/health-checks/tye-all.yaml

@ -0,0 +1,36 @@
# tye application configuration file
# read all about it at https://github.com/dotnet/tye
name: health-checks
services:
- name: health-all
project: api/api.csproj
replicas: 3
bindings:
- port: 8004
liveness:
http:
path: /healthy
headers:
- name: name1
value: value1
- name: name2
value: value2
initialDelay: 5
period: 1
timeout: 1
successThreshold: 1
failureThreshold: 1
readiness:
http:
path: /ready
headers:
- name: name3
value: value3
- name: name4
value: value4
initialDelay: 5
period: 1
timeout: 2
successThreshold: 1
failureThreshold: 1

30
test/E2ETest/testassets/projects/health-checks/tye-ingress.yaml

@ -0,0 +1,30 @@
# tye application configuration file
# read all about it at https://github.com/dotnet/tye
name: health-checks
ingress:
- name: ingress
bindings:
- port: 8006
rules:
- path: /api
service: health-ingress-svc
services:
- name: health-ingress-svc
project: api/api.csproj
replicas: 2
liveness:
http:
path: /healthy
initialDelay: 5
period: 1
timeout: 1
successThreshold: 1
failureThreshold: 1
readiness:
http:
path: /ready
initialDelay: 5
period: 1
timeout: 1
successThreshold: 1
failureThreshold: 1

17
test/E2ETest/testassets/projects/health-checks/tye-liveness.yaml

@ -0,0 +1,17 @@
# tye application configuration file
# read all about it at https://github.com/dotnet/tye
name: health-checks
services:
- name: health-liveness
project: api/api.csproj
replicas: 3
bindings:
- port: 8002
liveness:
http:
path: /healthy
initialDelay: 5
period: 1
timeout: 1
successThreshold: 1
failureThreshold: 1

10
test/E2ETest/testassets/projects/health-checks/tye-none.yaml

@ -0,0 +1,10 @@
# tye application configuration file
# read all about it at https://github.com/dotnet/tye
name: health-checks
services:
- name: health-none
project: api/api.csproj
replicas: 3
bindings:
- port: 8001

26
test/E2ETest/testassets/projects/health-checks/tye-proxy.yaml

@ -0,0 +1,26 @@
# tye application configuration file
# read all about it at https://github.com/dotnet/tye
name: health-checks
services:
- name: health-proxy
project: api/api.csproj
replicas: 2
bindings:
- port: 8005
liveness:
http:
path: /healthy
initialDelay: 5
period: 1
timeout: 1
successThreshold: 1
failureThreshold: 1
readiness:
http:
path: /ready
initialDelay: 5
period: 1
timeout: 1
successThreshold: 1
failureThreshold: 1

17
test/E2ETest/testassets/projects/health-checks/tye-readiness.yaml

@ -0,0 +1,17 @@
# tye application configuration file
# read all about it at https://github.com/dotnet/tye
name: health-checks
services:
- name: health-readiness
project: api/api.csproj
replicas: 3
bindings:
- port: 8003
readiness:
http:
path: /ready
initialDelay: 5
period: 1
timeout: 1
successThreshold: 1
failureThreshold: 1

127
test/Test.Infrastructure/TestHelpers.cs

@ -101,38 +101,50 @@ namespace Test.Infrastructure
return temp;
}
public static async Task StartHostAndWaitForReplicasToStart(TyeHost host)
public static async Task<bool> DoOperationAndWaitForReplicasToChangeState(TyeHost host, ReplicaState desiredState, int n, HashSet<string>? toChange, HashSet<string>? rest, Func<ReplicaEvent, string> entitySelector, TimeSpan waitUntilSuccess, Func<TyeHost, Task> operation)
{
var startedTask = new TaskCompletionSource<bool>();
var alreadyStarted = 0;
var totalReplicas = host.Application.Services.Sum(s => s.Value.Description.Replicas);
if (toChange != null && rest != null && rest.Overlaps(toChange))
{
throw new ArgumentException($"{nameof(toChange)} and {nameof(rest)} can't overlap");
}
var changedTask = new TaskCompletionSource<bool>();
var remaining = n;
void OnReplicaChange(ReplicaEvent ev)
{
if (ev.State == ReplicaState.Started)
if (rest != null && rest.Contains(entitySelector(ev)))
{
Interlocked.Increment(ref alreadyStarted);
changedTask!.TrySetResult(false);
}
else if (ev.State == ReplicaState.Stopped)
else if ((toChange == null || toChange.Contains(entitySelector(ev))) && ev.State == desiredState)
{
Interlocked.Decrement(ref alreadyStarted);
Interlocked.Decrement(ref remaining);
}
if (alreadyStarted == totalReplicas)
if (remaining == 0)
{
startedTask!.TrySetResult(true);
Task.Delay(waitUntilSuccess)
.ContinueWith(_ =>
{
if (!changedTask!.Task.IsCompleted)
{
changedTask!.TrySetResult(remaining == 0);
}
});
}
}
var servicesStateObserver = host.Application.Services.Select(srv => srv.Value.ReplicaEvents.Subscribe(OnReplicaChange)).ToList();
await host.StartAsync();
await operation(host);
using var cancellation = new CancellationTokenSource(WaitForServicesTimeout);
try
{
await using (cancellation.Token.Register(() => startedTask.TrySetCanceled()))
await using (cancellation.Token.Register(() => changedTask.TrySetCanceled()))
{
await startedTask.Task;
return await changedTask.Task;
}
}
finally
@ -144,46 +156,58 @@ namespace Test.Infrastructure
}
}
public static async Task PurgeHostAndWaitForGivenReplicasToStop(TyeHost host, string[] replicas)
public static async Task<bool> DoOperationAndWaitForReplicasToRestart(TyeHost host, HashSet<string> toRestart, HashSet<string>? rest, Func<ReplicaEvent, string> entitySelector, TimeSpan waitUntilSuccess, Func<TyeHost, Task> operation)
{
static async Task Purge(TyeHost host)
if (rest != null && rest.Overlaps(toRestart))
{
var logger = host.DashboardWebApplication!.Logger;
var replicaRegistry = new ReplicaRegistry(host.Application.ContextDirectory, logger);
var processRunner = new ProcessRunner(logger, replicaRegistry, new ProcessRunnerOptions());
var dockerRunner = new DockerRunner(logger, replicaRegistry);
await processRunner.StartAsync(new Application(new FileInfo(host.Application.Source), new Dictionary<string, Service>()));
await dockerRunner.StartAsync(new Application(new FileInfo(host.Application.Source), new Dictionary<string, Service>()));
throw new ArgumentException($"{nameof(toRestart)} and {nameof(rest)} can't overlap");
}
var stoppedTask = new TaskCompletionSource<bool>();
var remaining = replicas.Length;
var restartedTask = new TaskCompletionSource<bool>();
var remaining = toRestart.Count;
var alreadyStarted = 0;
void OnReplicaChange(ReplicaEvent ev)
{
if (replicas.Contains(ev.Replica.Name) && ev.State == ReplicaState.Stopped)
if (ev.State == ReplicaState.Started)
{
Interlocked.Decrement(ref remaining);
Interlocked.Increment(ref alreadyStarted);
}
else if (ev.State == ReplicaState.Stopped)
{
if (toRestart.Contains(entitySelector(ev)))
{
Interlocked.Decrement(ref remaining);
}
else if (rest != null && rest.Contains(entitySelector(ev)))
{
restartedTask!.SetResult(false);
}
}
if (remaining == 0)
if (remaining == 0 && alreadyStarted == toRestart.Count)
{
stoppedTask!.TrySetResult(true);
Task.Delay(waitUntilSuccess)
.ContinueWith(_ =>
{
if (!restartedTask!.Task.IsCompleted)
{
restartedTask!.SetResult(remaining == 0 && alreadyStarted == toRestart.Count);
}
});
}
}
var servicesStateObserver = host.Application.Services.Select(srv => srv.Value.ReplicaEvents.Subscribe(OnReplicaChange)).ToList();
// We purge existing replicas by restarting the host which will initiate the purging process
await Purge(host);
await operation(host);
using var cancellation = new CancellationTokenSource(WaitForServicesTimeout);
try
{
await using (cancellation.Token.Register(() => stoppedTask.TrySetCanceled()))
await using (cancellation.Token.Register(() => restartedTask.TrySetCanceled()))
{
await stoppedTask.Task;
return await restartedTask.Task;
}
}
finally
@ -194,5 +218,44 @@ namespace Test.Infrastructure
}
}
}
public static Task<bool> DoOperationAndWaitForReplicasToChangeState(TyeHost host, ReplicaState desiredState, int n, HashSet<string>? toChange, HashSet<string>? rest, TimeSpan waitUntilSuccess, Func<TyeHost, Task> operation)
=> DoOperationAndWaitForReplicasToChangeState(host, desiredState, n, toChange, rest, ev => ev.Replica.Name, waitUntilSuccess, operation);
public static Task<bool> DoOperationAndWaitForReplicasToRestart(TyeHost host, HashSet<string> toRestart, HashSet<string>? rest, TimeSpan waitUntilSuccess, Func<TyeHost, Task> operation)
=> DoOperationAndWaitForReplicasToRestart(host, toRestart, rest, ev => ev.Replica.Name, waitUntilSuccess, operation);
public static async Task StartHostAndWaitForReplicasToStart(TyeHost host, string[]? services = null, ReplicaState desiredState = ReplicaState.Started)
{
if (services == null)
{
await DoOperationAndWaitForReplicasToChangeState(host, desiredState, host.Application.Services.Sum(s => s.Value.Description.Replicas), null, null, TimeSpan.Zero, h => h.StartAsync());
}
else
{
if (services.Any(s => !host.Application.Services.ContainsKey(s)))
{
throw new ArgumentException($"not all services given in {nameof(services)} exist");
}
await DoOperationAndWaitForReplicasToChangeState(host, desiredState, host.Application.Services.Where(s => services.Contains(s.Value.Description.Name)).Sum(s => s.Value.Description.Replicas), services.ToHashSet(), null, ev => ev.Replica.Service.Description.Name, TimeSpan.Zero, h => h.StartAsync());
}
}
public static async Task PurgeHostAndWaitForGivenReplicasToStop(TyeHost host, string[] replicas)
{
static async Task Purge(TyeHost host)
{
var logger = host.DashboardWebApplication!.Logger;
var replicaRegistry = new ReplicaRegistry(host.Application.ContextDirectory, logger);
var processRunner = new ProcessRunner(logger, replicaRegistry, new ProcessRunnerOptions());
var dockerRunner = new DockerRunner(logger, replicaRegistry);
await processRunner.StartAsync(new Application(new FileInfo(host.Application.Source), new Dictionary<string, Service>()));
await dockerRunner.StartAsync(new Application(new FileInfo(host.Application.Source), new Dictionary<string, Service>()));
}
await DoOperationAndWaitForReplicasToChangeState(host, ReplicaState.Stopped, replicas.Length, replicas.ToHashSet(), new HashSet<string>(), TimeSpan.Zero, Purge);
}
}
}

Loading…
Cancel
Save