2 Commits
Author SHA1 Message Date
BlubbFish 47cdef0959 Thread Abort in netcore is not working 2019-12-01 17:53:57 +01:00
BlubbFish 35ba2b8135 change to netcore 2019-11-29 14:45:10 +01:00
4 changed files with 37 additions and 11 deletions
View File
@@ -35,6 +35,7 @@
<Compile Include="Properties\AssemblyInfo.cs" />
</ItemGroup>
<ItemGroup>
<Content Include="..\CHANGELOG.md" />
<Content Include="..\CONTRIBUTING.md" />
<Content Include="..\LICENSE" />
<Content Include="..\README.md" />
@@ -15,6 +15,7 @@
<RepositoryType>git</RepositoryType>
<PackageReleaseNotes>1.1.0 Rewrite Module to reconnect itselfs, so you dont need to watch over the the state of the connection</PackageReleaseNotes>
<NeutralLanguage>de-DE</NeutralLanguage>
<PackageId>Mqtt.Data.Connector.IoT.Utils.BlubbFish</PackageId>
</PropertyGroup>
<ItemGroup>
@@ -23,6 +24,7 @@
</ItemGroup>
<ItemGroup>
<Content Include="../CHANGELOG.md" />
<Content Include="../CONTRIBUTING.md" />
<Content Include="../LICENSE" />
<Content Include="../README.md" />
+34 -11
View File
@@ -11,9 +11,10 @@ namespace BlubbFish.Utils.IoT.Connector.Data {
public class Mqtt : ADataBackend, IDisposable {
private MqttClient client;
private Thread connectionWatcher;
private Boolean connectionWatcherRunning;
public Mqtt(Dictionary<String, String> settings) : base(settings) {
Console.WriteLine("BlubbFish.Utils.IoT.Connector.Data.Mqtt.Mqtt()");
Console.WriteLine("BlubbFish.Utils.IoT.Connector.Data.Mqtt(" + this.ToString()+")");
Int32 port = 1883;
if(this.settings.ContainsKey("port")) {
port = Int32.Parse(this.settings["port"]);
@@ -22,25 +23,45 @@ namespace BlubbFish.Utils.IoT.Connector.Data {
this.ConnectionWatcher();
}
public override String ToString() {
String ret = "mqtt://";
if (this.settings.ContainsKey("user")) {
ret += this.settings["user"];
if (this.settings.ContainsKey("pass")) {
ret += ":" + this.settings["pass"];
}
ret += "@";
}
ret += this.settings["server"];
if (this.settings.ContainsKey("port")) {
ret += ":" + this.settings["port"];
}
if (this.settings.ContainsKey("topic")) {
ret += "/" + this.settings["topic"];
}
return ret;
}
#region ConectionManage
private void ConnectionWatcher() {
this.connectionWatcher = new Thread(this.ConnectionWatcherRunner);
this.connectionWatcherRunning = true;
this.connectionWatcher.Start();
}
private void ConnectionWatcherRunner() {
while(true) {
while(this.connectionWatcherRunning) {
try {
if(!this.IsConnected) {
this.Reconnect();
Thread.Sleep(1000);
}
Thread.Sleep(500);
Thread.Sleep(10);
} catch(Exception) { }
}
}
private void Reconnect() {
Console.WriteLine("BlubbFish.Utils.IoT.Connector.Data.Mqtt.Reconnect()");
if(this.IsConnected) {
this.Disconnect(true);
} else {
@@ -50,7 +71,7 @@ namespace BlubbFish.Utils.IoT.Connector.Data {
}
private void Disconnect(Boolean complete) {
Console.WriteLine("BlubbFish.Utils.IoT.Connector.Data.Mqtt.Disconnect()");
Console.WriteLine("BlubbFish.Utils.IoT.Connector.Data.Disconnect(" + this.ToString() + ")");
this.client.MqttMsgPublishReceived -= this.Client_MqttMsgPublishReceived;
this.Unsubscripe();
if(complete) {
@@ -59,7 +80,7 @@ namespace BlubbFish.Utils.IoT.Connector.Data {
}
private void Connect() {
Console.WriteLine("BlubbFish.Utils.IoT.Connector.Data.Mqtt.Connect()");
Console.WriteLine("BlubbFish.Utils.IoT.Connector.Data.Connect(" + this.ToString() + ")");
this.client.MqttMsgPublishReceived += this.Client_MqttMsgPublishReceived;
_ = this.settings.ContainsKey("user") && this.settings.ContainsKey("pass")
? this.client.Connect(Guid.NewGuid().ToString(), this.settings["user"], this.settings["pass"])
@@ -106,11 +127,13 @@ namespace BlubbFish.Utils.IoT.Connector.Data {
protected virtual void Dispose(Boolean disposing) {
if(!this.disposedValue) {
if(disposing) {try {
try {
this.connectionWatcher.Abort();
this.connectionWatcher = null;
} catch { }
if(disposing) {
try {
this.connectionWatcherRunning = false;
while(this.connectionWatcher.IsAlive) {
Thread.Sleep(10);
}
this.connectionWatcher = null;
this.Disconnect(true);
} catch (Exception) { }
}