3 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
BlubbFish 8381edd133 Refactoring and add netcore 2019-11-28 17:42:21 +01:00
6 changed files with 150 additions and 41 deletions
View File
@@ -34,6 +34,12 @@
<Compile Include="Mqtt.cs" /> <Compile Include="Mqtt.cs" />
<Compile Include="Properties\AssemblyInfo.cs" /> <Compile Include="Properties\AssemblyInfo.cs" />
</ItemGroup> </ItemGroup>
<ItemGroup>
<Content Include="..\CHANGELOG.md" />
<Content Include="..\CONTRIBUTING.md" />
<Content Include="..\LICENSE" />
<Content Include="..\README.md" />
</ItemGroup>
<ItemGroup> <ItemGroup>
<ProjectReference Include="..\..\..\Librarys\mqtt\M2Mqtt\M2Mqtt_4.7.1.csproj"> <ProjectReference Include="..\..\..\Librarys\mqtt\M2Mqtt\M2Mqtt_4.7.1.csproj">
<Project>{a11aef5a-b246-4fe8-8330-06db73cc8074}</Project> <Project>{a11aef5a-b246-4fe8-8330-06db73cc8074}</Project>
@@ -0,0 +1,40 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<TargetFramework>netcoreapp3.0</TargetFramework>
<RootNamespace>BlubbFish.Utils.IoT.Connector.Data</RootNamespace>
<AssemblyName>ConnectorDataMqtt</AssemblyName>
<Company>BlubbFish</Company>
<Authors>BlubbFish</Authors>
<Version>1.1.0</Version>
<Copyright>Copyright © BlubbFish 2017 - 27.05.2019</Copyright>
<Description>ADataBackend Connector that connects to mqtt using M2Mqtt</Description>
<PackageLicenseFile>LICENSE</PackageLicenseFile>
<PackageProjectUrl>http://git.blubbfish.net/vs_utils/ConnectorDataMqtt</PackageProjectUrl>
<RepositoryUrl>http://git.blubbfish.net/vs_utils/ConnectorDataMqtt.git</RepositoryUrl>
<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>
<ProjectReference Include="..\..\..\Librarys\mqtt\M2Mqtt\M2Mqtt_Core.csproj" />
<ProjectReference Include="..\..\Utils-IoT\Utils-IoT\Utils-IoT_Core.csproj" />
</ItemGroup>
<ItemGroup>
<Content Include="../CHANGELOG.md" />
<Content Include="../CONTRIBUTING.md" />
<Content Include="../LICENSE" />
<Content Include="../README.md" />
</ItemGroup>
<ItemGroup>
<None Include="..\LICENSE">
<Pack>True</Pack>
<PackagePath></PackagePath>
</None>
</ItemGroup>
</Project>
+50 -38
View File
@@ -11,9 +11,10 @@ namespace BlubbFish.Utils.IoT.Connector.Data {
public class Mqtt : ADataBackend, IDisposable { public class Mqtt : ADataBackend, IDisposable {
private MqttClient client; private MqttClient client;
private Thread connectionWatcher; private Thread connectionWatcher;
private Boolean connectionWatcherRunning;
public Mqtt(Dictionary<String, String> settings) : base(settings) { 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; Int32 port = 1883;
if(this.settings.ContainsKey("port")) { if(this.settings.ContainsKey("port")) {
port = Int32.Parse(this.settings["port"]); port = Int32.Parse(this.settings["port"]);
@@ -22,25 +23,45 @@ namespace BlubbFish.Utils.IoT.Connector.Data {
this.ConnectionWatcher(); 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 #region ConectionManage
private void ConnectionWatcher() { private void ConnectionWatcher() {
this.connectionWatcher = new Thread(this.ConnectionWatcherRunner); this.connectionWatcher = new Thread(this.ConnectionWatcherRunner);
this.connectionWatcherRunning = true;
this.connectionWatcher.Start(); this.connectionWatcher.Start();
} }
private void ConnectionWatcherRunner() { private void ConnectionWatcherRunner() {
while(true) { while(this.connectionWatcherRunning) {
try { try {
if(!this.IsConnected) { if(!this.IsConnected) {
this.Reconnect(); this.Reconnect();
Thread.Sleep(1000);
} }
Thread.Sleep(500); Thread.Sleep(10);
} catch(Exception) { } } catch(Exception) { }
} }
} }
private void Reconnect() { private void Reconnect() {
Console.WriteLine("BlubbFish.Utils.IoT.Connector.Data.Mqtt.Reconnect()");
if(this.IsConnected) { if(this.IsConnected) {
this.Disconnect(true); this.Disconnect(true);
} else { } else {
@@ -50,7 +71,7 @@ namespace BlubbFish.Utils.IoT.Connector.Data {
} }
private void Disconnect(Boolean complete) { 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.client.MqttMsgPublishReceived -= this.Client_MqttMsgPublishReceived;
this.Unsubscripe(); this.Unsubscripe();
if(complete) { if(complete) {
@@ -59,25 +80,19 @@ namespace BlubbFish.Utils.IoT.Connector.Data {
} }
private void Connect() { 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.client.MqttMsgPublishReceived += this.Client_MqttMsgPublishReceived;
if (this.settings.ContainsKey("user") && this.settings.ContainsKey("pass")) { _ = this.settings.ContainsKey("user") && this.settings.ContainsKey("pass")
this.client.Connect(Guid.NewGuid().ToString(), this.settings["user"], this.settings["pass"]); ? this.client.Connect(Guid.NewGuid().ToString(), this.settings["user"], this.settings["pass"])
} else { : this.client.Connect(Guid.NewGuid().ToString());
this.client.Connect(Guid.NewGuid().ToString());
}
this.Subscripe(); this.Subscripe();
} }
#endregion #endregion
#region Subscription #region Subscription
private void Unsubscripe() { private void Unsubscripe() => _ = this.settings.ContainsKey("topic")
if(this.settings.ContainsKey("topic")) { ? this.client.Unsubscribe(this.settings["topic"].Split(';'))
this.client.Unsubscribe(this.settings["topic"].Split(';')); : this.client.Unsubscribe(new String[] { "#" });
} else {
this.client.Unsubscribe(new String[] { "#" });
}
}
private void Subscripe() { private void Subscripe() {
if(this.settings.ContainsKey("topic")) { if(this.settings.ContainsKey("topic")) {
@@ -86,42 +101,39 @@ namespace BlubbFish.Utils.IoT.Connector.Data {
for(Int32 i = 0; i < qos.Length; i++) { for(Int32 i = 0; i < qos.Length; i++) {
qos[i] = MqttMsgBase.QOS_LEVEL_AT_LEAST_ONCE; qos[i] = MqttMsgBase.QOS_LEVEL_AT_LEAST_ONCE;
} }
this.client.Subscribe(this.settings["topic"].Split(';'), qos); _ = this.client.Subscribe(this.settings["topic"].Split(';'), qos);
} else { } else {
this.client.Subscribe(new String[] { "#" }, new Byte[] { MqttMsgBase.QOS_LEVEL_AT_LEAST_ONCE }); _ = this.client.Subscribe(new String[] { "#" }, new Byte[] { MqttMsgBase.QOS_LEVEL_AT_LEAST_ONCE });
} }
} }
#endregion #endregion
private async void Client_MqttMsgPublishReceived(Object sender, MqttMsgPublishEventArgs e) => await Task.Run(() => { private async void Client_MqttMsgPublishReceived(Object sender, MqttMsgPublishEventArgs e) => await Task.Run(() => this.NotifyClientIncomming(new DataEvent(Encoding.UTF8.GetString(e.Message), e.Topic, DateTime.Now)));
this.NotifyClientIncomming(new DataEvent(Encoding.UTF8.GetString(e.Message), e.Topic, DateTime.Now));
});
public override void Send(String topic, String data) { public override void Send(String topic, String data) {
this.client.Publish(topic, Encoding.UTF8.GetBytes(data)); _ = this.client.Publish(topic, Encoding.UTF8.GetBytes(data));
this.NotifyClientSending(new DataEvent(data, topic, DateTime.Now)); this.NotifyClientSending(new DataEvent(data, topic, DateTime.Now));
} }
public void Send(String topic, Byte[] data) {
_ = this.client.Publish(topic, data);
this.NotifyClientSending(new DataEvent(Encoding.UTF8.GetString(data), topic, DateTime.Now));
}
#region IDisposable Support #region IDisposable Support
private Boolean disposedValue = false; private Boolean disposedValue = false;
public override Boolean IsConnected { public override Boolean IsConnected => this.client != null ? this.client.IsConnected : false;
get {
if(this.client != null) {
return this.client.IsConnected;
} else {
return false;
}
}
}
protected virtual void Dispose(Boolean disposing) { protected virtual void Dispose(Boolean disposing) {
if(!this.disposedValue) { if(!this.disposedValue) {
if(disposing) {try { if(disposing) {
try { try {
this.connectionWatcher.Abort(); this.connectionWatcherRunning = false;
this.connectionWatcher = null; while(this.connectionWatcher.IsAlive) {
} catch { } Thread.Sleep(10);
}
this.connectionWatcher = null;
this.Disconnect(true); this.Disconnect(true);
} catch (Exception) { } } catch (Exception) { }
} }
+5 -3
View File
@@ -1,4 +1,5 @@
using System.Reflection; #if !NETCOREAPP
using System.Reflection;
using System.Resources; using System.Resources;
using System.Runtime.InteropServices; using System.Runtime.InteropServices;
@@ -10,7 +11,7 @@ using System.Runtime.InteropServices;
[assembly: AssemblyConfiguration("")] [assembly: AssemblyConfiguration("")]
[assembly: AssemblyCompany("BlubbFish")] [assembly: AssemblyCompany("BlubbFish")]
[assembly: AssemblyProduct("ConnectorDataMqtt")] [assembly: AssemblyProduct("ConnectorDataMqtt")]
[assembly: AssemblyCopyright("Copyright © 2017 - 27.05.2019")] [assembly: AssemblyCopyright("Copyright © BlubbFish 2017 - 27.05.2019")]
[assembly: AssemblyTrademark("© BlubbFish")] [assembly: AssemblyTrademark("© BlubbFish")]
[assembly: AssemblyCulture("")] [assembly: AssemblyCulture("")]
[assembly: NeutralResourcesLanguage("de-DE")] [assembly: NeutralResourcesLanguage("de-DE")]
@@ -35,7 +36,8 @@ using System.Runtime.InteropServices;
// [assembly: AssemblyVersion("1.0.*")] // [assembly: AssemblyVersion("1.0.*")]
[assembly: AssemblyVersion("1.1.0")] [assembly: AssemblyVersion("1.1.0")]
[assembly: AssemblyFileVersion("1.1.0")] [assembly: AssemblyFileVersion("1.1.0")]
#endif
/* /*
* 1.1.0 Rewrite Module to reconnect itselfs, so you dont need to watch over the the state of the connection * 1.1.0 Rewrite Module to reconnect itselfs, so you dont need to watch over the the state of the connection
*/ */
+49
View File
@@ -0,0 +1,49 @@
Microsoft Visual Studio Solution File, Format Version 12.00
# Visual Studio Version 16
VisualStudioVersion = 16.0.29519.87
MinimumVisualStudioVersion = 10.0.40219.1
Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "ConnectorDataMqtt_Core", "ConnectorDataMqtt\ConnectorDataMqtt_Core.csproj", "{5D7A8892-6687-4FE2-85D1-12B49782EDF0}"
EndProject
Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Utils-IoT_Core", "..\Utils-IoT\Utils-IoT\Utils-IoT_Core.csproj", "{45837342-5906-4906-AB49-971767D9E180}"
EndProject
Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "M2Mqtt_Core", "..\..\Librarys\mqtt\M2Mqtt\M2Mqtt_Core.csproj", "{05B829EF-6F4F-496D-8936-35D7F2179DF0}"
EndProject
Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "litjson_Core", "..\..\Librarys\litjson\litjson\litjson_Core.csproj", "{CA22D3E9-4195-453D-9CE0-3EBC672E6D06}"
EndProject
Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Utils_Core", "..\Utils\Utils\Utils_Core.csproj", "{19424C4A-D024-484A-BE65-09C7C529A2C4}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
Release|Any CPU = Release|Any CPU
EndGlobalSection
GlobalSection(ProjectConfigurationPlatforms) = postSolution
{5D7A8892-6687-4FE2-85D1-12B49782EDF0}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{5D7A8892-6687-4FE2-85D1-12B49782EDF0}.Debug|Any CPU.Build.0 = Debug|Any CPU
{5D7A8892-6687-4FE2-85D1-12B49782EDF0}.Release|Any CPU.ActiveCfg = Release|Any CPU
{5D7A8892-6687-4FE2-85D1-12B49782EDF0}.Release|Any CPU.Build.0 = Release|Any CPU
{45837342-5906-4906-AB49-971767D9E180}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{45837342-5906-4906-AB49-971767D9E180}.Debug|Any CPU.Build.0 = Debug|Any CPU
{45837342-5906-4906-AB49-971767D9E180}.Release|Any CPU.ActiveCfg = Release|Any CPU
{45837342-5906-4906-AB49-971767D9E180}.Release|Any CPU.Build.0 = Release|Any CPU
{05B829EF-6F4F-496D-8936-35D7F2179DF0}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{05B829EF-6F4F-496D-8936-35D7F2179DF0}.Debug|Any CPU.Build.0 = Debug|Any CPU
{05B829EF-6F4F-496D-8936-35D7F2179DF0}.Release|Any CPU.ActiveCfg = Release|Any CPU
{05B829EF-6F4F-496D-8936-35D7F2179DF0}.Release|Any CPU.Build.0 = Release|Any CPU
{CA22D3E9-4195-453D-9CE0-3EBC672E6D06}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{CA22D3E9-4195-453D-9CE0-3EBC672E6D06}.Debug|Any CPU.Build.0 = Debug|Any CPU
{CA22D3E9-4195-453D-9CE0-3EBC672E6D06}.Release|Any CPU.ActiveCfg = Release|Any CPU
{CA22D3E9-4195-453D-9CE0-3EBC672E6D06}.Release|Any CPU.Build.0 = Release|Any CPU
{19424C4A-D024-484A-BE65-09C7C529A2C4}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{19424C4A-D024-484A-BE65-09C7C529A2C4}.Debug|Any CPU.Build.0 = Debug|Any CPU
{19424C4A-D024-484A-BE65-09C7C529A2C4}.Release|Any CPU.ActiveCfg = Release|Any CPU
{19424C4A-D024-484A-BE65-09C7C529A2C4}.Release|Any CPU.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
EndGlobalSection
GlobalSection(ExtensibilityGlobals) = postSolution
SolutionGuid = {92985FD1-A967-4AE1-B2E0-5E938C0A7D6A}
EndGlobalSection
EndGlobal