2019-11-17 16:27:58 +00:00
|
|
|
using System;
|
|
|
|
using System.Collections.Generic;
|
|
|
|
using System.Text;
|
|
|
|
using System.Threading;
|
|
|
|
using System.Threading.Tasks;
|
|
|
|
using Microsoft.Azure.Devices.Client;
|
|
|
|
using Microsoft.Extensions.Logging;
|
|
|
|
using Newtonsoft.Json;
|
|
|
|
|
|
|
|
namespace NucuCar.Domain.Telemetry
|
|
|
|
{
|
|
|
|
public class TelemetryPublisherAzure : TelemetryPublisher, IDisposable
|
|
|
|
{
|
|
|
|
// Needs to be configured via the Configure method or setup directly.
|
2019-11-23 14:09:44 +00:00
|
|
|
public string ConnectionString { get; set; }
|
|
|
|
public string TelemetrySource { private get; set; }
|
2019-11-17 16:27:58 +00:00
|
|
|
protected DeviceClient DeviceClient;
|
2019-11-23 14:09:44 +00:00
|
|
|
|
2019-11-17 16:27:58 +00:00
|
|
|
public override void Start()
|
|
|
|
{
|
|
|
|
try
|
|
|
|
{
|
2019-11-23 14:09:44 +00:00
|
|
|
DeviceClient = DeviceClient.CreateFromConnectionString(ConnectionString, TransportType.Mqtt);
|
2019-11-17 16:27:58 +00:00
|
|
|
}
|
|
|
|
catch (FormatException)
|
|
|
|
{
|
|
|
|
Logger.LogCritical("Can't start telemetry service! Malformed connection string!");
|
|
|
|
throw;
|
|
|
|
}
|
|
|
|
Logger.LogInformation("Started the AzureTelemetryPublisher!");
|
|
|
|
}
|
|
|
|
|
|
|
|
public override async Task PublishAsync(CancellationToken cancellationToken)
|
|
|
|
{
|
|
|
|
foreach (var telemeter in RegisteredTelemeters)
|
|
|
|
{
|
|
|
|
var data = telemeter.GetTelemetryData();
|
|
|
|
if (data == null)
|
|
|
|
{
|
|
|
|
Logger.LogWarning($"Warning! Data for {telemeter.GetIdentifier()} is null!");
|
|
|
|
continue;
|
|
|
|
}
|
2019-11-23 14:09:44 +00:00
|
|
|
var metadata = new Dictionary<string, object>
|
|
|
|
{
|
|
|
|
["source"] = TelemetrySource ?? nameof(TelemetryPublisherAzure),
|
|
|
|
["id"] = telemeter.GetIdentifier(),
|
|
|
|
["timestamp"] = DateTime.Now,
|
|
|
|
["data"] = data,
|
|
|
|
};
|
2019-11-17 16:27:58 +00:00
|
|
|
|
2019-11-23 14:09:44 +00:00
|
|
|
await PublishViaMqtt(metadata, cancellationToken);
|
2019-11-17 16:27:58 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
private async Task PublishViaMqtt(Dictionary<string, object> data, CancellationToken cancellationToken)
|
|
|
|
{
|
|
|
|
if (cancellationToken.IsCancellationRequested)
|
|
|
|
{
|
|
|
|
Logger.LogInformation("Stopping the AzureTelemetryPublisher, cancellation requested.");
|
|
|
|
await DeviceClient.CloseAsync(cancellationToken);
|
|
|
|
return;
|
|
|
|
}
|
|
|
|
var messageString = JsonConvert.SerializeObject(data);
|
|
|
|
var message = new Message(Encoding.ASCII.GetBytes(messageString));
|
|
|
|
Logger.LogDebug($"Telemetry message: {message}");
|
|
|
|
await DeviceClient.SendEventAsync(message, cancellationToken);
|
|
|
|
}
|
|
|
|
|
|
|
|
public void Dispose()
|
|
|
|
{
|
|
|
|
DeviceClient.CloseAsync().GetAwaiter().GetResult();
|
|
|
|
}
|
|
|
|
|
|
|
|
public override bool Publish(int timeout)
|
|
|
|
{
|
|
|
|
throw new NotImplementedException();
|
|
|
|
}
|
|
|
|
|
|
|
|
#pragma warning disable 1998
|
|
|
|
public override async Task StartAsync()
|
|
|
|
{
|
|
|
|
throw new NotImplementedException();
|
|
|
|
}
|
|
|
|
#pragma warning restore 1998
|
|
|
|
}
|
|
|
|
}
|