1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386238723882389239023912392239323942395239623972398239924002401240224032404240524062407240824092410241124122413241424152416241724182419242024212422242324242425242624272428242924302431243224332434243524362437243824392440244124422443244424452446244724482449245024512452245324542455245624572458245924602461246224632464246524662467246824692470247124722473247424752476247724782479248024812482248324842485248624872488248924902491249224932494249524962497249824992500250125022503250425052506250725082509251025112512251325142515251625172518251925202521252225232524252525262527252825292530253125322533253425352536253725382539254025412542254325442545254625472548254925502551255225532554255525562557255825592560256125622563256425652566256725682569257025712572257325742575257625772578257925802581258225832584258525862587258825892590259125922593259425952596259725982599260026012602260326042605260626072608260926102611261226132614261526162617261826192620262126222623262426252626262726282629263026312632263326342635263626372638 |
- using System;
- using System.Net;
- #if !(WINDOWS_APP || WINDOWS_PHONE_APP || (!UNITY_EDITOR&&UNITY_WSA_10_0&&!ENABLE_IL2CPP))
- using System.Net.Sockets;
- using System.Security.Cryptography.X509Certificates;
- #endif
- using System.Threading;
- using uPLibrary.Networking.M2Mqtt.Exceptions;
- using uPLibrary.Networking.M2Mqtt.Messages;
- using uPLibrary.Networking.M2Mqtt.Session;
- using uPLibrary.Networking.M2Mqtt.Utility;
- using uPLibrary.Networking.M2Mqtt.Internal;
- #if (MF_FRAMEWORK_VERSION_V4_2 || MF_FRAMEWORK_VERSION_V4_3)
- using Microsoft.SPOT;
- #if SSL
- using Microsoft.SPOT.Net.Security;
- #endif
- #else
- using System.Collections.Generic;
- #if (SSL && !(WINDOWS_APP || WINDOWS_PHONE_APP || (!UNITY_EDITOR&&UNITY_WSA_10_0&&!ENABLE_IL2CPP)))
- using System.Security.Authentication;
- using System.Net.Security;
- #endif
- #endif
- #if (WINDOWS_APP || WINDOWS_PHONE_APP || (!UNITY_EDITOR && UNITY_WSA_10_0 && !ENABLE_IL2CPP))
- using Windows.Networking.Sockets;
- #endif
- using System.Collections;
- using MqttUtility = uPLibrary.Networking.M2Mqtt.Utility;
- using System.IO;
- using System.Net.Security;
- using UnityEngine;
- namespace uPLibrary.Networking.M2Mqtt
- {
-
-
-
- public class MqttClient
- {
- #if BROKER
- #region Constants ...
-
- private const string RECEIVE_THREAD_NAME = "ReceiveThread";
- private const string RECEIVE_EVENT_THREAD_NAME = "DispatchEventThread";
- private const string PROCESS_INFLIGHT_THREAD_NAME = "ProcessInflightThread";
- private const string KEEP_ALIVE_THREAD = "KeepAliveThread";
- #endregion
- #endif
-
-
-
- public delegate void MqttMsgPublishEventHandler(object sender, MqttMsgPublishEventArgs e);
-
-
-
- public delegate void MqttMsgPublishedEventHandler(object sender, MqttMsgPublishedEventArgs e);
-
-
-
- public delegate void MqttMsgSubscribedEventHandler(object sender, MqttMsgSubscribedEventArgs e);
-
-
-
- public delegate void MqttMsgUnsubscribedEventHandler(object sender, MqttMsgUnsubscribedEventArgs e);
- #if BROKER
-
-
-
- public delegate void MqttMsgSubscribeEventHandler(object sender, MqttMsgSubscribeEventArgs e);
-
-
-
- public delegate void MqttMsgUnsubscribeEventHandler(object sender, MqttMsgUnsubscribeEventArgs e);
-
-
-
- public delegate void MqttMsgConnectEventHandler(object sender, MqttMsgConnectEventArgs e);
-
-
-
- public delegate void MqttMsgDisconnectEventHandler(object sender, EventArgs e);
- #endif
-
-
-
- public delegate void ConnectionClosedEventHandler(object sender, EventArgs e);
-
- private string brokerHostName;
- private int brokerPort;
-
- private bool isRunning;
-
- private AutoResetEvent receiveEventWaitHandle;
-
- private AutoResetEvent inflightWaitHandle;
-
- AutoResetEvent syncEndReceiving;
-
- MqttMsgBase msgReceived;
-
- Exception exReceiving;
-
- private int keepAlivePeriod;
-
- private AutoResetEvent keepAliveEvent;
- private AutoResetEvent keepAliveEventEnd;
-
- private int lastCommTime;
-
- public event MqttMsgPublishEventHandler MqttMsgPublishReceived;
-
- public event MqttMsgPublishedEventHandler MqttMsgPublished;
-
- public event MqttMsgSubscribedEventHandler MqttMsgSubscribed;
-
- public event MqttMsgUnsubscribedEventHandler MqttMsgUnsubscribed;
- #if BROKER
-
- public event MqttMsgSubscribeEventHandler MqttMsgSubscribeReceived;
-
- public event MqttMsgUnsubscribeEventHandler MqttMsgUnsubscribeReceived;
-
- public event MqttMsgConnectEventHandler MqttMsgConnected;
-
- public event MqttMsgDisconnectEventHandler MqttMsgDisconnected;
- #endif
-
- public event ConnectionClosedEventHandler ConnectionClosed;
-
-
- private IMqttNetworkChannel channel;
-
- private Queue inflightQueue;
-
- private Queue internalQueue;
-
- private Queue eventQueue;
-
- private MqttClientSession session;
-
- private MqttSettings settings;
-
- private ushort messageIdCounter = 0;
-
- private bool isConnectionClosing;
-
-
-
- public bool IsConnected { get; private set; }
-
-
-
- public string ClientId { get; private set; }
-
-
-
- public bool CleanSession { get; private set; }
-
-
-
- public bool WillFlag { get; private set; }
-
-
-
- public byte WillQosLevel { get; private set; }
-
-
-
- public string WillTopic { get; private set; }
-
-
-
- public string WillMessage { get; private set; }
-
-
-
- public MqttProtocolVersion ProtocolVersion { get; set; }
- #if BROKER
-
-
-
- public MqttClientSession Session
- {
- get { return this.session; }
- set { this.session = value; }
- }
- #endif
-
-
-
- public MqttSettings Settings
- {
- get { return this.settings; }
- }
- #if !(WINDOWS_APP || WINDOWS_PHONE_APP || (!UNITY_EDITOR&&UNITY_WSA_10_0&&!ENABLE_IL2CPP))
-
-
-
-
- [Obsolete("Use this ctor MqttClient(string brokerHostName) insted")]
- public MqttClient(IPAddress brokerIpAddress) :
- this(brokerIpAddress, MqttSettings.MQTT_BROKER_DEFAULT_PORT, false, null, null, MqttSslProtocols.None)
- {
- }
-
-
-
-
-
-
-
-
-
- [Obsolete("Use this ctor MqttClient(string brokerHostName, int brokerPort, bool secure, X509Certificate caCert) insted")]
- public MqttClient(IPAddress brokerIpAddress, int brokerPort, bool secure, X509Certificate caCert, X509Certificate clientCert, MqttSslProtocols sslProtocol)
- {
- #if !(MF_FRAMEWORK_VERSION_V4_2 || MF_FRAMEWORK_VERSION_V4_3 || COMPACT_FRAMEWORK)
- this.Init(brokerIpAddress.ToString(), brokerPort, secure, caCert, clientCert, sslProtocol, null, null);
- #else
- this.Init(brokerIpAddress.ToString(), brokerPort, secure, caCert, clientCert, sslProtocol);
- #endif
- }
- #endif
-
-
-
-
- public MqttClient(string brokerHostName) :
- #if !(WINDOWS_APP || WINDOWS_PHONE_APP || (!UNITY_EDITOR&&UNITY_WSA_10_0&&!ENABLE_IL2CPP))
- this(brokerHostName, MqttSettings.MQTT_BROKER_DEFAULT_PORT, false, null, null, MqttSslProtocols.None)
- #else
- this(brokerHostName, MqttSettings.MQTT_BROKER_DEFAULT_PORT, false, MqttSslProtocols.None)
- #endif
- {
- }
-
-
-
-
-
-
-
- #if !(WINDOWS_APP || WINDOWS_PHONE_APP || (!UNITY_EDITOR && UNITY_WSA_10_0 && !ENABLE_IL2CPP))
-
-
- public MqttClient(string brokerHostName, int brokerPort, bool secure, X509Certificate caCert, X509Certificate clientCert, MqttSslProtocols sslProtocol)
- #else
- public MqttClient(string brokerHostName, int brokerPort, bool secure, MqttSslProtocols sslProtocol)
- #endif
- {
- #if !(MF_FRAMEWORK_VERSION_V4_2 || MF_FRAMEWORK_VERSION_V4_3 || COMPACT_FRAMEWORK || WINDOWS_APP || WINDOWS_PHONE_APP || (!UNITY_EDITOR&&UNITY_WSA_10_0&&!ENABLE_IL2CPP))
- this.Init(brokerHostName, brokerPort, secure, caCert, clientCert, sslProtocol, null, null);
- #elif (WINDOWS_APP || WINDOWS_PHONE_APP || (!UNITY_EDITOR&&UNITY_WSA_10_0&&!ENABLE_IL2CPP))
- this.Init(brokerHostName, brokerPort, secure, sslProtocol);
- #else
- this.Init(brokerHostName, brokerPort, secure, caCert, clientCert, sslProtocol);
- #endif
- }
- #if !(MF_FRAMEWORK_VERSION_V4_2 || MF_FRAMEWORK_VERSION_V4_3 || COMPACT_FRAMEWORK || WINDOWS_APP || WINDOWS_PHONE_APP || (!UNITY_EDITOR&&UNITY_WSA_10_0&&!ENABLE_IL2CPP))
-
-
-
-
-
-
-
-
-
-
- public MqttClient(string brokerHostName, int brokerPort, bool secure, X509Certificate caCert, X509Certificate clientCert, MqttSslProtocols sslProtocol,
- RemoteCertificateValidationCallback userCertificateValidationCallback)
- : this(brokerHostName, brokerPort, secure, caCert, clientCert, sslProtocol, userCertificateValidationCallback, null)
- {
- }
-
-
-
-
-
-
-
-
-
- public MqttClient(string brokerHostName, int brokerPort, bool secure, MqttSslProtocols sslProtocol,
- RemoteCertificateValidationCallback userCertificateValidationCallback,
- LocalCertificateSelectionCallback userCertificateSelectionCallback)
- : this(brokerHostName, brokerPort, secure, null, null, sslProtocol, userCertificateValidationCallback, userCertificateSelectionCallback)
- {
- }
-
-
-
-
-
-
-
-
-
-
-
- public MqttClient(string brokerHostName, int brokerPort, bool secure, X509Certificate caCert, X509Certificate clientCert, MqttSslProtocols sslProtocol,
- RemoteCertificateValidationCallback userCertificateValidationCallback,
- LocalCertificateSelectionCallback userCertificateSelectionCallback)
- {
- this.Init(brokerHostName, brokerPort, secure, caCert, clientCert, sslProtocol, userCertificateValidationCallback, userCertificateSelectionCallback);
- }
- #endif
- #if BROKER
-
-
-
-
- public MqttClient(IMqttNetworkChannel channel)
- {
-
- this.ProtocolVersion = MqttProtocolVersion.Version_3_1_1;
- this.channel = channel;
-
- this.settings = MqttSettings.Instance;
-
- this.IsConnected = false;
- this.ClientId = null;
- this.CleanSession = true;
- this.keepAliveEvent = new AutoResetEvent(false);
-
- this.inflightWaitHandle = new AutoResetEvent(false);
- this.inflightQueue = new Queue();
-
- this.receiveEventWaitHandle = new AutoResetEvent(false);
- this.eventQueue = new Queue();
- this.internalQueue = new Queue();
-
- this.session = null;
- }
- #endif
-
-
-
-
-
-
-
-
-
- #if !(MF_FRAMEWORK_VERSION_V4_2 || MF_FRAMEWORK_VERSION_V4_3 || COMPACT_FRAMEWORK || WINDOWS_APP || WINDOWS_PHONE_APP || ((!UNITY_EDITOR && UNITY_WSA_10_0 && !ENABLE_IL2CPP)))
-
-
- private void Init(string brokerHostName, int brokerPort, bool secure, X509Certificate caCert, X509Certificate clientCert, MqttSslProtocols sslProtocol,
- RemoteCertificateValidationCallback userCertificateValidationCallback,
- LocalCertificateSelectionCallback userCertificateSelectionCallback)
- #elif (WINDOWS_APP || WINDOWS_PHONE_APP || (!UNITY_EDITOR&&UNITY_WSA_10_0&&!ENABLE_IL2CPP))
- private void Init(string brokerHostName, int brokerPort, bool secure, MqttSslProtocols sslProtocol)
- #else
- private void Init(string brokerHostName, int brokerPort, bool secure, X509Certificate caCert, X509Certificate clientCert, MqttSslProtocols sslProtocol)
- #endif
- {
-
- this.ProtocolVersion = MqttProtocolVersion.Version_3_1_1;
- #if !SSL
-
- if (secure)
- throw new ArgumentException("Library compiled without SSL support");
- #endif
- this.brokerHostName = brokerHostName;
- this.brokerPort = brokerPort;
-
- this.settings = MqttSettings.Instance;
-
- if (!secure)
- this.settings.Port = this.brokerPort;
- else
- this.settings.SslPort = this.brokerPort;
- this.syncEndReceiving = new AutoResetEvent(false);
- this.keepAliveEvent = new AutoResetEvent(false);
-
- this.inflightWaitHandle = new AutoResetEvent(false);
- this.inflightQueue = new Queue();
-
- this.receiveEventWaitHandle = new AutoResetEvent(false);
- this.eventQueue = new Queue();
- this.internalQueue = new Queue();
-
- this.session = null;
-
- #if !(MF_FRAMEWORK_VERSION_V4_2 || MF_FRAMEWORK_VERSION_V4_3 || COMPACT_FRAMEWORK || WINDOWS_APP || WINDOWS_PHONE_APP || (!UNITY_EDITOR&&UNITY_WSA_10_0&&!ENABLE_IL2CPP))
- this.channel = new MqttNetworkChannel(this.brokerHostName, this.brokerPort, secure, caCert, clientCert, sslProtocol, userCertificateValidationCallback, userCertificateSelectionCallback);
- #elif (WINDOWS_APP || WINDOWS_PHONE_APP || (!UNITY_EDITOR&&UNITY_WSA_10_0&&!ENABLE_IL2CPP))
- this.channel = new MqttNetworkChannel(this.brokerHostName, this.brokerPort, secure, sslProtocol);
- #else
- this.channel = new MqttNetworkChannel(this.brokerHostName, this.brokerPort, secure, caCert, clientCert, sslProtocol);
- #endif
- }
-
-
-
-
-
- public byte Connect(string clientId)
- {
- return this.Connect(clientId, null, null, false, MqttMsgConnect.QOS_LEVEL_AT_MOST_ONCE, false, null, null, true, MqttMsgConnect.KEEP_ALIVE_PERIOD_DEFAULT);
- }
-
-
-
-
-
-
-
- public byte Connect(string clientId,
- string username,
- string password)
- {
- return this.Connect(clientId, username, password, false, MqttMsgConnect.QOS_LEVEL_AT_MOST_ONCE, false, null, null, true, MqttMsgConnect.KEEP_ALIVE_PERIOD_DEFAULT);
- }
-
-
-
-
-
-
-
-
-
- public byte Connect(string clientId,
- string username,
- string password,
- bool cleanSession,
- ushort keepAlivePeriod)
- {
- return this.Connect(clientId, username, password, false, MqttMsgConnect.QOS_LEVEL_AT_MOST_ONCE, false, null, null, cleanSession, keepAlivePeriod);
- }
-
-
-
-
-
-
-
-
-
-
-
-
-
-
- public byte Connect(string clientId,
- string username,
- string password,
- bool willRetain,
- byte willQosLevel,
- bool willFlag,
- string willTopic,
- string willMessage,
- bool cleanSession,
- ushort keepAlivePeriod)
- {
-
- MqttMsgConnect connect = new MqttMsgConnect(clientId,
- username,
- password,
- willRetain,
- willQosLevel,
- willFlag,
- willTopic,
- willMessage,
- cleanSession,
- keepAlivePeriod,
- (byte)this.ProtocolVersion);
- try
- {
-
- this.channel.Connect();
- }
- catch (Exception ex)
- {
- throw new MqttConnectionException("Exception connecting to the broker", ex);
- }
- this.lastCommTime = 0;
- this.isRunning = true;
- this.isConnectionClosing = false;
-
- Fx.StartThread(this.ReceiveThread);
-
- MqttMsgConnack connack = (MqttMsgConnack)this.SendReceive(connect);
-
- if (connack.ReturnCode == MqttMsgConnack.CONN_ACCEPTED)
- {
-
- this.ClientId = clientId;
- this.CleanSession = cleanSession;
- this.WillFlag = willFlag;
- this.WillTopic = willTopic;
- this.WillMessage = willMessage;
- this.WillQosLevel = willQosLevel;
- this.keepAlivePeriod = keepAlivePeriod * 1000;
-
- this.RestoreSession();
-
- if (this.keepAlivePeriod != 0)
- {
-
- Fx.StartThread(this.KeepAliveThread);
- }
-
- Fx.StartThread(this.DispatchEventThread);
-
-
- Fx.StartThread(this.ProcessInflightThread);
- this.IsConnected = true;
- }
- return connack.ReturnCode;
- }
-
-
-
- public void Disconnect()
- {
- MqttMsgDisconnect disconnect = new MqttMsgDisconnect();
- this.Send(disconnect);
-
- this.OnConnectionClosing();
- }
- #if BROKER
-
-
-
- public void Open()
- {
- this.isRunning = true;
-
- Fx.StartThread(this.ReceiveThread);
-
- Fx.StartThread(this.DispatchEventThread);
-
- Fx.StartThread(this.ProcessInflightThread);
- }
- #endif
-
-
-
- #if BROKER
- public void Close()
- #else
- private void Close()
- #endif
- {
-
- this.isRunning = false;
-
- if (this.receiveEventWaitHandle != null)
- this.receiveEventWaitHandle.Set();
-
- if (this.inflightWaitHandle != null)
- this.inflightWaitHandle.Set();
- #if BROKER
-
- this.keepAliveEvent.Set();
- #else
-
- this.keepAliveEvent.Set();
- if (this.keepAliveEventEnd != null)
- this.keepAliveEventEnd.WaitOne();
- #endif
-
- this.inflightQueue.Clear();
- this.internalQueue.Clear();
- this.eventQueue.Clear();
-
- this.channel.Close();
- this.IsConnected = false;
- }
-
-
-
-
- private MqttMsgPingResp Ping()
- {
- MqttMsgPingReq pingreq = new MqttMsgPingReq();
- try
- {
-
- return (MqttMsgPingResp)this.SendReceive(pingreq, this.keepAlivePeriod);
- }
- catch (Exception e)
- {
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Error, "Exception occurred: {0}", e.ToString());
- #endif
-
- this.OnConnectionClosing();
- return null;
- }
- }
- #if BROKER
-
-
-
-
-
-
-
- public void Connack(MqttMsgConnect connect, byte returnCode, string clientId, bool sessionPresent)
- {
- this.lastCommTime = 0;
-
- MqttMsgConnack connack = new MqttMsgConnack();
- connack.ReturnCode = returnCode;
-
- if (this.ProtocolVersion == MqttProtocolVersion.Version_3_1_1)
- connack.SessionPresent = sessionPresent;
-
- this.Send(connack);
-
- if (connack.ReturnCode == MqttMsgConnack.CONN_ACCEPTED)
- {
-
-
- this.ClientId = (clientId == null) ? connect.ClientId : clientId;
- this.CleanSession = connect.CleanSession;
- this.WillFlag = connect.WillFlag;
- this.WillTopic = connect.WillTopic;
- this.WillMessage = connect.WillMessage;
- this.WillQosLevel = connect.WillQosLevel;
- this.keepAlivePeriod = connect.KeepAlivePeriod * 1000;
-
- this.keepAlivePeriod += (this.keepAlivePeriod / 2);
-
- Fx.StartThread(this.KeepAliveThread);
- this.isConnectionClosing = false;
- this.IsConnected = true;
- }
-
- else
- {
- this.Close();
- }
- }
-
-
-
-
-
- public void Suback(ushort messageId, byte[] grantedQosLevels)
- {
- MqttMsgSuback suback = new MqttMsgSuback();
- suback.MessageId = messageId;
- suback.GrantedQoSLevels = grantedQosLevels;
- this.Send(suback);
- }
-
-
-
-
- public void Unsuback(ushort messageId)
- {
- MqttMsgUnsuback unsuback = new MqttMsgUnsuback();
- unsuback.MessageId = messageId;
- this.Send(unsuback);
- }
- #endif
-
-
-
-
-
-
- public ushort Subscribe(string[] topics, byte[] qosLevels)
- {
- MqttMsgSubscribe subscribe = new MqttMsgSubscribe(topics, qosLevels);
- subscribe.MessageId = this.GetMessageId();
-
- this.EnqueueInflight(subscribe, MqttMsgFlow.ToPublish);
-
- return subscribe.MessageId;
- }
-
-
-
-
-
- public ushort Unsubscribe(string[] topics)
- {
- MqttMsgUnsubscribe unsubscribe =
- new MqttMsgUnsubscribe(topics);
- unsubscribe.MessageId = this.GetMessageId();
-
- this.EnqueueInflight(unsubscribe, MqttMsgFlow.ToPublish);
- return unsubscribe.MessageId;
- }
-
-
-
-
-
-
-
- public ushort Publish(string topic, byte[] message)
- {
- return this.Publish(topic, message, MqttMsgBase.QOS_LEVEL_AT_MOST_ONCE, false);
- }
-
-
-
-
-
-
-
-
- public ushort Publish(string topic, byte[] message, byte qosLevel, bool retain)
- {
- MqttMsgPublish publish =
- new MqttMsgPublish(topic, message, false, qosLevel, retain);
- publish.MessageId = this.GetMessageId();
-
- bool enqueue = this.EnqueueInflight(publish, MqttMsgFlow.ToPublish);
-
- if (enqueue)
- return publish.MessageId;
-
- else
- throw new MqttClientException(MqttClientErrorCode.InflightQueueFull);
- }
-
-
-
-
- private void OnInternalEvent(InternalEvent internalEvent)
- {
- lock (this.eventQueue)
- {
- this.eventQueue.Enqueue(internalEvent);
- }
- this.receiveEventWaitHandle.Set();
- }
-
-
-
- private void OnConnectionClosing()
- {
- if (!this.isConnectionClosing)
- {
- this.isConnectionClosing = true;
- this.receiveEventWaitHandle.Set();
- }
- }
-
-
-
-
- private void OnMqttMsgPublishReceived(MqttMsgPublish publish)
- {
- if (this.MqttMsgPublishReceived != null)
- {
- this.MqttMsgPublishReceived(this,
- new MqttMsgPublishEventArgs(publish.Topic, publish.Message, publish.DupFlag, publish.QosLevel, publish.Retain));
- }
- }
-
-
-
-
-
- private void OnMqttMsgPublished(ushort messageId, bool isPublished)
- {
- if (this.MqttMsgPublished != null)
- {
- this.MqttMsgPublished(this,
- new MqttMsgPublishedEventArgs(messageId, isPublished));
- }
- }
-
-
-
-
- private void OnMqttMsgSubscribed(MqttMsgSuback suback)
- {
- if (this.MqttMsgSubscribed != null)
- {
- this.MqttMsgSubscribed(this,
- new MqttMsgSubscribedEventArgs(suback.MessageId, suback.GrantedQoSLevels));
- }
- }
-
-
-
-
- private void OnMqttMsgUnsubscribed(ushort messageId)
- {
- if (this.MqttMsgUnsubscribed != null)
- {
- this.MqttMsgUnsubscribed(this,
- new MqttMsgUnsubscribedEventArgs(messageId));
- }
- }
- #if BROKER
-
-
-
-
-
-
- private void OnMqttMsgSubscribeReceived(ushort messageId, string[] topics, byte[] qosLevels)
- {
- if (this.MqttMsgSubscribeReceived != null)
- {
- this.MqttMsgSubscribeReceived(this,
- new MqttMsgSubscribeEventArgs(messageId, topics, qosLevels));
- }
- }
-
-
-
-
-
- private void OnMqttMsgUnsubscribeReceived(ushort messageId, string[] topics)
- {
- if (this.MqttMsgUnsubscribeReceived != null)
- {
- this.MqttMsgUnsubscribeReceived(this,
- new MqttMsgUnsubscribeEventArgs(messageId, topics));
- }
- }
-
-
-
- private void OnMqttMsgConnected(MqttMsgConnect connect)
- {
- if (this.MqttMsgConnected != null)
- {
- this.ProtocolVersion = (MqttProtocolVersion)connect.ProtocolVersion;
- this.MqttMsgConnected(this, new MqttMsgConnectEventArgs(connect));
- }
- }
-
-
-
- private void OnMqttMsgDisconnected()
- {
- if (this.MqttMsgDisconnected != null)
- {
- this.MqttMsgDisconnected(this, EventArgs.Empty);
- }
- }
- #endif
-
-
-
- private void OnConnectionClosed()
- {
- if (this.ConnectionClosed != null)
- {
- this.ConnectionClosed(this, EventArgs.Empty);
- }
- }
-
-
-
-
- private void Send(byte[] msgBytes)
- {
- try
- {
-
- this.channel.Send(msgBytes);
- #if !BROKER
-
- this.lastCommTime = Environment.TickCount;
- #endif
- }
- catch (Exception e)
- {
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Error, "Exception occurred: {0}", e.ToString());
- #endif
- throw new MqttCommunicationException(e);
- }
- }
-
-
-
-
- private void Send(MqttMsgBase msg)
- {
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Frame, "SEND {0}", msg);
- #endif
- this.Send(msg.GetBytes((byte)this.ProtocolVersion));
- }
-
-
-
-
-
- private MqttMsgBase SendReceive(byte[] msgBytes)
- {
- return this.SendReceive(msgBytes, MqttSettings.MQTT_DEFAULT_TIMEOUT);
- }
-
-
-
-
-
-
- private MqttMsgBase SendReceive(byte[] msgBytes, int timeout)
- {
-
- this.syncEndReceiving.Reset();
- try
- {
-
- this.channel.Send(msgBytes);
-
- this.lastCommTime = Environment.TickCount;
- }
- catch (Exception e)
- {
- #if !(MF_FRAMEWORK_VERSION_V4_2 || MF_FRAMEWORK_VERSION_V4_3 || COMPACT_FRAMEWORK || WINDOWS_APP || WINDOWS_PHONE_APP || (!UNITY_EDITOR&&UNITY_WSA_10_0&&!ENABLE_IL2CPP))
- if (typeof(SocketException) == e.GetType())
- {
-
- if (((SocketException)e).SocketErrorCode == SocketError.ConnectionReset)
- this.IsConnected = false;
- }
- #endif
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Error, "Exception occurred: {0}", e.ToString());
- #endif
- throw new MqttCommunicationException(e);
- }
- #if (MF_FRAMEWORK_VERSION_V4_2 || MF_FRAMEWORK_VERSION_V4_3 || COMPACT_FRAMEWORK)
-
- if (this.syncEndReceiving.WaitOne(timeout, false))
- #else
-
- if (this.syncEndReceiving.WaitOne(timeout))
- #endif
- {
-
- if (this.exReceiving == null)
- return this.msgReceived;
-
- else
- throw this.exReceiving;
- }
- else
- {
-
- throw new MqttCommunicationException();
- }
- }
-
-
-
-
-
- private MqttMsgBase SendReceive(MqttMsgBase msg)
- {
- return this.SendReceive(msg, MqttSettings.MQTT_DEFAULT_TIMEOUT);
- }
-
-
-
-
-
-
- private MqttMsgBase SendReceive(MqttMsgBase msg, int timeout)
- {
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Frame, "SEND {0}", msg);
- #endif
- return this.SendReceive(msg.GetBytes((byte)this.ProtocolVersion), timeout);
- }
-
-
-
-
-
-
- private bool EnqueueInflight(MqttMsgBase msg, MqttMsgFlow flow)
- {
-
- bool enqueue = true;
-
- if ((msg.Type == MqttMsgBase.MQTT_MSG_PUBLISH_TYPE) &&
- (msg.QosLevel == MqttMsgBase.QOS_LEVEL_EXACTLY_ONCE))
- {
- lock (this.inflightQueue)
- {
-
-
-
-
- MqttMsgContextFinder msgCtxFinder = new MqttMsgContextFinder(msg.MessageId, MqttMsgFlow.ToAcknowledge);
- MqttMsgContext msgCtx = (MqttMsgContext)this.inflightQueue.Get(msgCtxFinder.Find);
-
-
- if (msgCtx != null)
- {
- msgCtx.State = MqttMsgState.QueuedQos2;
- msgCtx.Flow = MqttMsgFlow.ToAcknowledge;
- enqueue = false;
- }
- }
- }
- if (enqueue)
- {
-
- MqttMsgState state = MqttMsgState.QueuedQos0;
-
- switch (msg.QosLevel)
- {
-
- case MqttMsgBase.QOS_LEVEL_AT_MOST_ONCE:
- state = MqttMsgState.QueuedQos0;
- break;
-
- case MqttMsgBase.QOS_LEVEL_AT_LEAST_ONCE:
- state = MqttMsgState.QueuedQos1;
- break;
-
- case MqttMsgBase.QOS_LEVEL_EXACTLY_ONCE:
- state = MqttMsgState.QueuedQos2;
- break;
- }
-
-
- if (msg.Type == MqttMsgBase.MQTT_MSG_SUBSCRIBE_TYPE)
- state = MqttMsgState.SendSubscribe;
- else if (msg.Type == MqttMsgBase.MQTT_MSG_UNSUBSCRIBE_TYPE)
- state = MqttMsgState.SendUnsubscribe;
-
- MqttMsgContext msgContext = new MqttMsgContext()
- {
- Message = msg,
- State = state,
- Flow = flow,
- Attempt = 0
- };
- lock (this.inflightQueue)
- {
-
- enqueue = (this.inflightQueue.Count < this.settings.InflightQueueSize);
- if (enqueue)
- {
-
- this.inflightQueue.Enqueue(msgContext);
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Queuing, "enqueued {0}", msg);
- #endif
-
- if (msg.Type == MqttMsgBase.MQTT_MSG_PUBLISH_TYPE)
- {
-
- if ((msgContext.Flow == MqttMsgFlow.ToPublish) &&
- ((msg.QosLevel == MqttMsgBase.QOS_LEVEL_AT_LEAST_ONCE) ||
- (msg.QosLevel == MqttMsgBase.QOS_LEVEL_EXACTLY_ONCE)))
- {
- if (this.session != null)
- this.session.InflightMessages.Add(msgContext.Key, msgContext);
- }
-
- else if ((msgContext.Flow == MqttMsgFlow.ToAcknowledge) &&
- (msg.QosLevel == MqttMsgBase.QOS_LEVEL_EXACTLY_ONCE))
- {
- if (this.session != null)
- this.session.InflightMessages.Add(msgContext.Key, msgContext);
- }
- }
- }
- }
- }
- this.inflightWaitHandle.Set();
- return enqueue;
- }
-
-
-
-
- private void EnqueueInternal(MqttMsgBase msg)
- {
-
- bool enqueue = true;
-
- if (msg.Type == MqttMsgBase.MQTT_MSG_PUBREL_TYPE)
- {
- lock (this.inflightQueue)
- {
-
-
-
-
-
- MqttMsgContextFinder msgCtxFinder = new MqttMsgContextFinder(msg.MessageId, MqttMsgFlow.ToAcknowledge);
- MqttMsgContext msgCtx = (MqttMsgContext)this.inflightQueue.Get(msgCtxFinder.Find);
-
-
- if (msgCtx == null)
- {
- MqttMsgPubcomp pubcomp = new MqttMsgPubcomp();
- pubcomp.MessageId = msg.MessageId;
- this.Send(pubcomp);
- enqueue = false;
- }
- }
- }
-
- else if (msg.Type == MqttMsgBase.MQTT_MSG_PUBCOMP_TYPE)
- {
- lock (this.inflightQueue)
- {
-
-
-
-
-
- MqttMsgContextFinder msgCtxFinder = new MqttMsgContextFinder(msg.MessageId, MqttMsgFlow.ToPublish);
- MqttMsgContext msgCtx = (MqttMsgContext)this.inflightQueue.Get(msgCtxFinder.Find);
-
- if (msgCtx == null)
- {
- enqueue = false;
- }
- }
- }
-
- else if (msg.Type == MqttMsgBase.MQTT_MSG_PUBREC_TYPE)
- {
- lock (this.inflightQueue)
- {
-
-
-
-
-
- MqttMsgContextFinder msgCtxFinder = new MqttMsgContextFinder(msg.MessageId, MqttMsgFlow.ToPublish);
- MqttMsgContext msgCtx = (MqttMsgContext)this.inflightQueue.Get(msgCtxFinder.Find);
-
- if (msgCtx == null)
- {
- enqueue = false;
- }
- }
- }
- if (enqueue)
- {
- lock (this.internalQueue)
- {
- this.internalQueue.Enqueue(msg);
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Queuing, "enqueued {0}", msg);
- #endif
- this.inflightWaitHandle.Set();
- }
- }
- }
-
-
-
- private void ReceiveThread()
- {
- int readBytes = 0;
- byte[] fixedHeaderFirstByte = new byte[1];
- byte msgType;
- while (this.isRunning)
- {
- try
- {
-
- readBytes = this.channel.Receive(fixedHeaderFirstByte);
- if (readBytes > 0)
- {
- #if BROKER
-
- this.lastCommTime = Environment.TickCount;
- #endif
-
- msgType = (byte)((fixedHeaderFirstByte[0] & MqttMsgBase.MSG_TYPE_MASK) >> MqttMsgBase.MSG_TYPE_OFFSET);
- switch (msgType)
- {
-
- case MqttMsgBase.MQTT_MSG_CONNECT_TYPE:
- #if BROKER
- MqttMsgConnect connect = MqttMsgConnect.Parse(fixedHeaderFirstByte[0], (byte)this.ProtocolVersion, this.channel);
- #if TRACE
- Trace.WriteLine(TraceLevel.Frame, "RECV {0}", connect);
- #endif
-
- this.OnInternalEvent(new MsgInternalEvent(connect));
- break;
- #else
- throw new MqttClientException(MqttClientErrorCode.WrongBrokerMessage);
- #endif
-
-
- case MqttMsgBase.MQTT_MSG_CONNACK_TYPE:
- #if BROKER
- throw new MqttClientException(MqttClientErrorCode.WrongBrokerMessage);
- #else
- this.msgReceived = MqttMsgConnack.Parse(fixedHeaderFirstByte[0], (byte)this.ProtocolVersion, this.channel);
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Frame, "RECV {0}", this.msgReceived);
- #endif
- this.syncEndReceiving.Set();
- break;
- #endif
-
- case MqttMsgBase.MQTT_MSG_PINGREQ_TYPE:
- #if BROKER
- this.msgReceived = MqttMsgPingReq.Parse(fixedHeaderFirstByte[0], (byte)this.ProtocolVersion, this.channel);
- #if TRACE
- Trace.WriteLine(TraceLevel.Frame, "RECV {0}", this.msgReceived);
- #endif
- MqttMsgPingResp pingresp = new MqttMsgPingResp();
- this.Send(pingresp);
- break;
- #else
- throw new MqttClientException(MqttClientErrorCode.WrongBrokerMessage);
- #endif
-
- case MqttMsgBase.MQTT_MSG_PINGRESP_TYPE:
- #if BROKER
- throw new MqttClientException(MqttClientErrorCode.WrongBrokerMessage);
- #else
- this.msgReceived = MqttMsgPingResp.Parse(fixedHeaderFirstByte[0], (byte)this.ProtocolVersion, this.channel);
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Frame, "RECV {0}", this.msgReceived);
- #endif
- this.syncEndReceiving.Set();
- break;
- #endif
-
- case MqttMsgBase.MQTT_MSG_SUBSCRIBE_TYPE:
- #if BROKER
- MqttMsgSubscribe subscribe = MqttMsgSubscribe.Parse(fixedHeaderFirstByte[0], (byte)this.ProtocolVersion, this.channel);
- #if TRACE
- Trace.WriteLine(TraceLevel.Frame, "RECV {0}", subscribe);
- #endif
-
- this.OnInternalEvent(new MsgInternalEvent(subscribe));
- break;
- #else
- throw new MqttClientException(MqttClientErrorCode.WrongBrokerMessage);
- #endif
-
- case MqttMsgBase.MQTT_MSG_SUBACK_TYPE:
- #if BROKER
- throw new MqttClientException(MqttClientErrorCode.WrongBrokerMessage);
- #else
-
- MqttMsgSuback suback = MqttMsgSuback.Parse(fixedHeaderFirstByte[0], (byte)this.ProtocolVersion, this.channel);
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Frame, "RECV {0}", suback);
- #endif
-
- this.EnqueueInternal(suback);
- break;
- #endif
-
- case MqttMsgBase.MQTT_MSG_PUBLISH_TYPE:
- MqttMsgPublish publish = MqttMsgPublish.Parse(fixedHeaderFirstByte[0], (byte)this.ProtocolVersion, this.channel);
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Frame, "RECV {0}", publish);
- #endif
-
- this.EnqueueInflight(publish, MqttMsgFlow.ToAcknowledge);
- break;
-
- case MqttMsgBase.MQTT_MSG_PUBACK_TYPE:
-
- MqttMsgPuback puback = MqttMsgPuback.Parse(fixedHeaderFirstByte[0], (byte)this.ProtocolVersion, this.channel);
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Frame, "RECV {0}", puback);
- #endif
-
- this.EnqueueInternal(puback);
- break;
-
- case MqttMsgBase.MQTT_MSG_PUBREC_TYPE:
-
- MqttMsgPubrec pubrec = MqttMsgPubrec.Parse(fixedHeaderFirstByte[0], (byte)this.ProtocolVersion, this.channel);
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Frame, "RECV {0}", pubrec);
- #endif
-
- this.EnqueueInternal(pubrec);
- break;
-
- case MqttMsgBase.MQTT_MSG_PUBREL_TYPE:
-
- MqttMsgPubrel pubrel = MqttMsgPubrel.Parse(fixedHeaderFirstByte[0], (byte)this.ProtocolVersion, this.channel);
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Frame, "RECV {0}", pubrel);
- #endif
-
- this.EnqueueInternal(pubrel);
- break;
-
-
- case MqttMsgBase.MQTT_MSG_PUBCOMP_TYPE:
-
- MqttMsgPubcomp pubcomp = MqttMsgPubcomp.Parse(fixedHeaderFirstByte[0], (byte)this.ProtocolVersion, this.channel);
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Frame, "RECV {0}", pubcomp);
- #endif
-
- this.EnqueueInternal(pubcomp);
- break;
-
- case MqttMsgBase.MQTT_MSG_UNSUBSCRIBE_TYPE:
- #if BROKER
- MqttMsgUnsubscribe unsubscribe = MqttMsgUnsubscribe.Parse(fixedHeaderFirstByte[0], (byte)this.ProtocolVersion, this.channel);
- #if TRACE
- Trace.WriteLine(TraceLevel.Frame, "RECV {0}", unsubscribe);
- #endif
-
- this.OnInternalEvent(new MsgInternalEvent(unsubscribe));
- break;
- #else
- throw new MqttClientException(MqttClientErrorCode.WrongBrokerMessage);
- #endif
-
- case MqttMsgBase.MQTT_MSG_UNSUBACK_TYPE:
- #if BROKER
- throw new MqttClientException(MqttClientErrorCode.WrongBrokerMessage);
- #else
-
- MqttMsgUnsuback unsuback = MqttMsgUnsuback.Parse(fixedHeaderFirstByte[0], (byte)this.ProtocolVersion, this.channel);
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Frame, "RECV {0}", unsuback);
- #endif
-
- this.EnqueueInternal(unsuback);
- break;
- #endif
-
- case MqttMsgDisconnect.MQTT_MSG_DISCONNECT_TYPE:
- #if BROKER
- MqttMsgDisconnect disconnect = MqttMsgDisconnect.Parse(fixedHeaderFirstByte[0], (byte)this.ProtocolVersion, this.channel);
- #if TRACE
- Trace.WriteLine(TraceLevel.Frame, "RECV {0}", disconnect);
- #endif
-
- this.OnInternalEvent(new MsgInternalEvent(disconnect));
- break;
- #else
- throw new MqttClientException(MqttClientErrorCode.WrongBrokerMessage);
- #endif
- default:
- throw new MqttClientException(MqttClientErrorCode.WrongBrokerMessage);
- }
- this.exReceiving = null;
- }
-
- else
- {
-
- this.OnConnectionClosing();
- }
- }
- catch (Exception e)
- {
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Error, "Exception occurred: {0}", e.ToString());
- #endif
- this.exReceiving = new MqttCommunicationException(e);
- bool close = false;
- if (e.GetType() == typeof(MqttClientException))
- {
-
- MqttClientException ex = e as MqttClientException;
- close = ((ex.ErrorCode == MqttClientErrorCode.InvalidFlagBits) ||
- (ex.ErrorCode == MqttClientErrorCode.InvalidProtocolName) ||
- (ex.ErrorCode == MqttClientErrorCode.InvalidConnectFlags));
- }
- #if !(WINDOWS_APP || WINDOWS_PHONE_APP || (!UNITY_EDITOR&&UNITY_WSA_10_0&&!ENABLE_IL2CPP))
- else if ((e.GetType() == typeof(IOException)) || (e.GetType() == typeof(SocketException)) ||
- ((e.InnerException != null) && (e.InnerException.GetType() == typeof(SocketException))))
- {
- close = true;
- }
- #endif
-
- if (close)
- {
-
- this.OnConnectionClosing();
- }
- }
- }
- }
-
-
-
- private void KeepAliveThread()
- {
- int delta = 0;
- int wait = this.keepAlivePeriod;
-
-
- this.keepAliveEventEnd = new AutoResetEvent(false);
- while (this.isRunning)
- {
- #if (MF_FRAMEWORK_VERSION_V4_2 || MF_FRAMEWORK_VERSION_V4_3 || COMPACT_FRAMEWORK)
-
- this.keepAliveEvent.WaitOne(wait, false);
- #else
-
- this.keepAliveEvent.WaitOne(wait);
- #endif
- if (this.isRunning)
- {
- delta = Environment.TickCount - this.lastCommTime;
-
- if (delta >= this.keepAlivePeriod)
- {
- #if BROKER
-
- this.OnConnectionClosing();
- #else
-
- this.Ping();
- wait = this.keepAlivePeriod;
- #endif
- }
- else
- {
-
- wait = this.keepAlivePeriod - delta;
- }
- }
- }
-
- this.keepAliveEventEnd.Set();
- }
-
-
-
- private void DispatchEventThread()
- {
- while (this.isRunning)
- {
- #if BROKER
- if ((this.eventQueue.Count == 0) && !this.isConnectionClosing)
- {
-
-
- if (!this.IsConnected)
- {
-
- if (!this.receiveEventWaitHandle.WaitOne(this.settings.TimeoutOnConnection))
- {
-
- this.Close();
-
- this.OnConnectionClosed();
- }
- }
- else
- {
-
- this.receiveEventWaitHandle.WaitOne();
- }
- }
- #else
- if ((this.eventQueue.Count == 0) && !this.isConnectionClosing)
-
- this.receiveEventWaitHandle.WaitOne();
- #endif
-
- if (this.isRunning)
- {
-
- InternalEvent internalEvent = null;
- lock (this.eventQueue)
- {
- if (this.eventQueue.Count > 0)
- internalEvent = (InternalEvent)this.eventQueue.Dequeue();
- }
-
- if (internalEvent != null)
- {
- MqttMsgBase msg = ((MsgInternalEvent)internalEvent).Message;
- if (msg != null)
- {
- switch (msg.Type)
- {
-
- case MqttMsgBase.MQTT_MSG_CONNECT_TYPE:
- #if BROKER
-
- this.OnMqttMsgConnected((MqttMsgConnect)msg);
- break;
- #else
- throw new MqttClientException(MqttClientErrorCode.WrongBrokerMessage);
- #endif
-
- case MqttMsgBase.MQTT_MSG_SUBSCRIBE_TYPE:
- #if BROKER
- MqttMsgSubscribe subscribe = (MqttMsgSubscribe)msg;
-
- this.OnMqttMsgSubscribeReceived(subscribe.MessageId, subscribe.Topics, subscribe.QoSLevels);
- break;
- #else
- throw new MqttClientException(MqttClientErrorCode.WrongBrokerMessage);
- #endif
-
- case MqttMsgBase.MQTT_MSG_SUBACK_TYPE:
-
- this.OnMqttMsgSubscribed((MqttMsgSuback)msg);
- break;
-
- case MqttMsgBase.MQTT_MSG_PUBLISH_TYPE:
-
- if (internalEvent.GetType() == typeof(MsgPublishedInternalEvent))
- this.OnMqttMsgPublished(msg.MessageId, false);
- else
-
- this.OnMqttMsgPublishReceived((MqttMsgPublish)msg);
- break;
-
- case MqttMsgBase.MQTT_MSG_PUBACK_TYPE:
-
-
- this.OnMqttMsgPublished(msg.MessageId, true);
- break;
-
- case MqttMsgBase.MQTT_MSG_PUBREL_TYPE:
-
-
- this.OnMqttMsgPublishReceived((MqttMsgPublish)msg);
- break;
-
- case MqttMsgBase.MQTT_MSG_PUBCOMP_TYPE:
-
-
- this.OnMqttMsgPublished(msg.MessageId, true);
- break;
-
- case MqttMsgBase.MQTT_MSG_UNSUBSCRIBE_TYPE:
- #if BROKER
- MqttMsgUnsubscribe unsubscribe = (MqttMsgUnsubscribe)msg;
-
- this.OnMqttMsgUnsubscribeReceived(unsubscribe.MessageId, unsubscribe.Topics);
- break;
- #else
- throw new MqttClientException(MqttClientErrorCode.WrongBrokerMessage);
- #endif
-
- case MqttMsgBase.MQTT_MSG_UNSUBACK_TYPE:
-
- this.OnMqttMsgUnsubscribed(msg.MessageId);
- break;
-
- case MqttMsgDisconnect.MQTT_MSG_DISCONNECT_TYPE:
- #if BROKER
-
- this.OnMqttMsgDisconnected();
- break;
- #else
- throw new MqttClientException(MqttClientErrorCode.WrongBrokerMessage);
- #endif
- }
- }
- }
-
-
- if ((this.eventQueue.Count == 0) && this.isConnectionClosing)
- {
-
- this.Close();
-
- this.OnConnectionClosed();
- }
- }
- }
- }
-
-
-
- private void ProcessInflightThread()
- {
- MqttMsgContext msgContext = null;
- MqttMsgBase msgInflight = null;
- MqttMsgBase msgReceived = null;
- InternalEvent internalEvent = null;
- bool acknowledge = false;
- int timeout = Timeout.Infinite;
- int delta;
- bool msgReceivedProcessed = false;
- try
- {
- while (this.isRunning)
- {
- #if (MF_FRAMEWORK_VERSION_V4_2 || MF_FRAMEWORK_VERSION_V4_3 || COMPACT_FRAMEWORK)
-
- this.inflightWaitHandle.WaitOne(timeout, false);
- #else
-
- this.inflightWaitHandle.WaitOne(timeout);
- #endif
-
- if (this.isRunning)
- {
- lock (this.inflightQueue)
- {
-
-
-
-
- msgReceivedProcessed = false;
- acknowledge = false;
- msgReceived = null;
-
-
- timeout = Int32.MaxValue;
-
-
- int count = this.inflightQueue.Count;
-
- while (count > 0)
- {
- count--;
- acknowledge = false;
- msgReceived = null;
-
- if (!this.isRunning)
- break;
-
- msgContext = (MqttMsgContext)this.inflightQueue.Dequeue();
-
- msgInflight = (MqttMsgBase)msgContext.Message;
- switch (msgContext.State)
- {
- case MqttMsgState.QueuedQos0:
-
- if (msgContext.Flow == MqttMsgFlow.ToPublish)
- {
- this.Send(msgInflight);
- }
-
- else if (msgContext.Flow == MqttMsgFlow.ToAcknowledge)
- {
- internalEvent = new MsgInternalEvent(msgInflight);
-
- this.OnInternalEvent(internalEvent);
- }
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Queuing, "processed {0}", msgInflight);
- #endif
- break;
- case MqttMsgState.QueuedQos1:
-
- case MqttMsgState.SendSubscribe:
- case MqttMsgState.SendUnsubscribe:
-
- if (msgContext.Flow == MqttMsgFlow.ToPublish)
- {
- msgContext.Timestamp = Environment.TickCount;
- msgContext.Attempt++;
- if (msgInflight.Type == MqttMsgBase.MQTT_MSG_PUBLISH_TYPE)
- {
-
- msgContext.State = MqttMsgState.WaitForPuback;
-
- if (msgContext.Attempt > 1)
- msgInflight.DupFlag = true;
- }
- else if (msgInflight.Type == MqttMsgBase.MQTT_MSG_SUBSCRIBE_TYPE)
-
- msgContext.State = MqttMsgState.WaitForSuback;
- else if (msgInflight.Type == MqttMsgBase.MQTT_MSG_UNSUBSCRIBE_TYPE)
-
- msgContext.State = MqttMsgState.WaitForUnsuback;
- this.Send(msgInflight);
-
- timeout = (this.settings.DelayOnRetry < timeout) ? this.settings.DelayOnRetry : timeout;
-
- this.inflightQueue.Enqueue(msgContext);
- }
-
- else if (msgContext.Flow == MqttMsgFlow.ToAcknowledge)
- {
- MqttMsgPuback puback = new MqttMsgPuback();
- puback.MessageId = msgInflight.MessageId;
- this.Send(puback);
- internalEvent = new MsgInternalEvent(msgInflight);
-
- this.OnInternalEvent(internalEvent);
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Queuing, "processed {0}", msgInflight);
- #endif
- }
- break;
- case MqttMsgState.QueuedQos2:
-
- if (msgContext.Flow == MqttMsgFlow.ToPublish)
- {
- msgContext.Timestamp = Environment.TickCount;
- msgContext.Attempt++;
- msgContext.State = MqttMsgState.WaitForPubrec;
-
- if (msgContext.Attempt > 1)
- msgInflight.DupFlag = true;
- this.Send(msgInflight);
-
- timeout = (this.settings.DelayOnRetry < timeout) ? this.settings.DelayOnRetry : timeout;
-
- this.inflightQueue.Enqueue(msgContext);
- }
-
- else if (msgContext.Flow == MqttMsgFlow.ToAcknowledge)
- {
- MqttMsgPubrec pubrec = new MqttMsgPubrec();
- pubrec.MessageId = msgInflight.MessageId;
- msgContext.State = MqttMsgState.WaitForPubrel;
- this.Send(pubrec);
-
- this.inflightQueue.Enqueue(msgContext);
- }
- break;
- case MqttMsgState.WaitForPuback:
- case MqttMsgState.WaitForSuback:
- case MqttMsgState.WaitForUnsuback:
-
-
-
- if (msgContext.Flow == MqttMsgFlow.ToPublish)
- {
- acknowledge = false;
- lock (this.internalQueue)
- {
- if (this.internalQueue.Count > 0)
- msgReceived = (MqttMsgBase)this.internalQueue.Peek();
- }
-
- if (msgReceived != null)
- {
-
- if (((msgReceived.Type == MqttMsgBase.MQTT_MSG_PUBACK_TYPE) && (msgInflight.Type == MqttMsgBase.MQTT_MSG_PUBLISH_TYPE) && (msgReceived.MessageId == msgInflight.MessageId)) ||
- ((msgReceived.Type == MqttMsgBase.MQTT_MSG_SUBACK_TYPE) && (msgInflight.Type == MqttMsgBase.MQTT_MSG_SUBSCRIBE_TYPE) && (msgReceived.MessageId == msgInflight.MessageId)) ||
- ((msgReceived.Type == MqttMsgBase.MQTT_MSG_UNSUBACK_TYPE) && (msgInflight.Type == MqttMsgBase.MQTT_MSG_UNSUBSCRIBE_TYPE) && (msgReceived.MessageId == msgInflight.MessageId)))
- {
- lock (this.internalQueue)
- {
-
- this.internalQueue.Dequeue();
- acknowledge = true;
- msgReceivedProcessed = true;
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Queuing, "dequeued {0}", msgReceived);
- #endif
- }
-
- if (msgReceived.Type == MqttMsgBase.MQTT_MSG_PUBACK_TYPE)
- internalEvent = new MsgPublishedInternalEvent(msgReceived, true);
- else
- internalEvent = new MsgInternalEvent(msgReceived);
-
- this.OnInternalEvent(internalEvent);
-
- if ((msgInflight.Type == MqttMsgBase.MQTT_MSG_PUBLISH_TYPE) &&
- (this.session != null) &&
- #if (MF_FRAMEWORK_VERSION_V4_2 || MF_FRAMEWORK_VERSION_V4_3 || COMPACT_FRAMEWORK)
- (this.session.InflightMessages.Contains(msgContext.Key)))
- #else
- (this.session.InflightMessages.ContainsKey(msgContext.Key)))
- #endif
- {
- this.session.InflightMessages.Remove(msgContext.Key);
- }
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Queuing, "processed {0}", msgInflight);
- #endif
- }
- }
-
- if (!acknowledge)
- {
- delta = Environment.TickCount - msgContext.Timestamp;
-
-
-
- if (delta >= this.settings.DelayOnRetry)
- {
-
- if (msgContext.Attempt < this.settings.AttemptsOnRetry)
- {
- msgContext.State = MqttMsgState.QueuedQos1;
-
- this.inflightQueue.Enqueue(msgContext);
-
- timeout = 0;
- }
- else
- {
-
- if (msgInflight.Type == MqttMsgBase.MQTT_MSG_PUBLISH_TYPE)
- {
-
- if ((this.session != null) &&
- #if (MF_FRAMEWORK_VERSION_V4_2 || MF_FRAMEWORK_VERSION_V4_3 || COMPACT_FRAMEWORK)
- (this.session.InflightMessages.Contains(msgContext.Key)))
- #else
- (this.session.InflightMessages.ContainsKey(msgContext.Key)))
- #endif
- {
- this.session.InflightMessages.Remove(msgContext.Key);
- }
- internalEvent = new MsgPublishedInternalEvent(msgInflight, false);
-
- this.OnInternalEvent(internalEvent);
- }
-
-
- }
- }
- else
- {
-
- this.inflightQueue.Enqueue(msgContext);
-
- int msgTimeout = (this.settings.DelayOnRetry - delta);
- timeout = (msgTimeout < timeout) ? msgTimeout : timeout;
- }
- }
- }
- break;
- case MqttMsgState.WaitForPubrec:
-
- if (msgContext.Flow == MqttMsgFlow.ToPublish)
- {
- acknowledge = false;
- lock (this.internalQueue)
- {
- if (this.internalQueue.Count > 0)
- msgReceived = (MqttMsgBase)this.internalQueue.Peek();
- }
-
- if ((msgReceived != null) && (msgReceived.Type == MqttMsgBase.MQTT_MSG_PUBREC_TYPE))
- {
-
- if (msgReceived.MessageId == msgInflight.MessageId)
- {
- lock (this.internalQueue)
- {
-
- this.internalQueue.Dequeue();
- acknowledge = true;
- msgReceivedProcessed = true;
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Queuing, "dequeued {0}", msgReceived);
- #endif
- }
- MqttMsgPubrel pubrel = new MqttMsgPubrel();
- pubrel.MessageId = msgInflight.MessageId;
- msgContext.State = MqttMsgState.WaitForPubcomp;
- msgContext.Timestamp = Environment.TickCount;
- msgContext.Attempt = 1;
- this.Send(pubrel);
-
- timeout = (this.settings.DelayOnRetry < timeout) ? this.settings.DelayOnRetry : timeout;
-
- this.inflightQueue.Enqueue(msgContext);
- }
- }
-
- if (!acknowledge)
- {
- delta = Environment.TickCount - msgContext.Timestamp;
-
- if (delta >= this.settings.DelayOnRetry)
- {
-
- if (msgContext.Attempt < this.settings.AttemptsOnRetry)
- {
- msgContext.State = MqttMsgState.QueuedQos2;
-
- this.inflightQueue.Enqueue(msgContext);
-
- timeout = 0;
- }
- else
- {
-
- if ((this.session != null) &&
- #if (MF_FRAMEWORK_VERSION_V4_2 || MF_FRAMEWORK_VERSION_V4_3 || COMPACT_FRAMEWORK)
- (this.session.InflightMessages.Contains(msgContext.Key)))
- #else
- (this.session.InflightMessages.ContainsKey(msgContext.Key)))
- #endif
- {
- this.session.InflightMessages.Remove(msgContext.Key);
- }
-
- internalEvent = new MsgPublishedInternalEvent(msgInflight, false);
-
- this.OnInternalEvent(internalEvent);
- }
- }
- else
- {
-
- this.inflightQueue.Enqueue(msgContext);
-
- int msgTimeout = (this.settings.DelayOnRetry - delta);
- timeout = (msgTimeout < timeout) ? msgTimeout : timeout;
- }
- }
- }
- break;
- case MqttMsgState.WaitForPubrel:
-
- if (msgContext.Flow == MqttMsgFlow.ToAcknowledge)
- {
- lock (this.internalQueue)
- {
- if (this.internalQueue.Count > 0)
- msgReceived = (MqttMsgBase)this.internalQueue.Peek();
- }
-
- if ((msgReceived != null) && (msgReceived.Type == MqttMsgBase.MQTT_MSG_PUBREL_TYPE))
- {
-
- if (msgReceived.MessageId == msgInflight.MessageId)
- {
- lock (this.internalQueue)
- {
-
- this.internalQueue.Dequeue();
- msgReceivedProcessed = true;
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Queuing, "dequeued {0}", msgReceived);
- #endif
- }
- MqttMsgPubcomp pubcomp = new MqttMsgPubcomp();
- pubcomp.MessageId = msgInflight.MessageId;
- this.Send(pubcomp);
- internalEvent = new MsgInternalEvent(msgInflight);
-
- this.OnInternalEvent(internalEvent);
-
- if ((msgInflight.Type == MqttMsgBase.MQTT_MSG_PUBLISH_TYPE) &&
- (this.session != null) &&
- #if (MF_FRAMEWORK_VERSION_V4_2 || MF_FRAMEWORK_VERSION_V4_3 || COMPACT_FRAMEWORK)
- (this.session.InflightMessages.Contains(msgContext.Key)))
- #else
- (this.session.InflightMessages.ContainsKey(msgContext.Key)))
- #endif
- {
- this.session.InflightMessages.Remove(msgContext.Key);
- }
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Queuing, "processed {0}", msgInflight);
- #endif
- }
- else
- {
-
- this.inflightQueue.Enqueue(msgContext);
- }
- }
- else
- {
-
- this.inflightQueue.Enqueue(msgContext);
- }
- }
- break;
- case MqttMsgState.WaitForPubcomp:
-
- if (msgContext.Flow == MqttMsgFlow.ToPublish)
- {
- acknowledge = false;
- lock (this.internalQueue)
- {
- if (this.internalQueue.Count > 0)
- msgReceived = (MqttMsgBase)this.internalQueue.Peek();
- }
-
- if ((msgReceived != null) && (msgReceived.Type == MqttMsgBase.MQTT_MSG_PUBCOMP_TYPE))
- {
-
- if (msgReceived.MessageId == msgInflight.MessageId)
- {
- lock (this.internalQueue)
- {
-
- this.internalQueue.Dequeue();
- acknowledge = true;
- msgReceivedProcessed = true;
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Queuing, "dequeued {0}", msgReceived);
- #endif
- }
- internalEvent = new MsgPublishedInternalEvent(msgReceived, true);
-
- this.OnInternalEvent(internalEvent);
-
- if ((msgInflight.Type == MqttMsgBase.MQTT_MSG_PUBLISH_TYPE) &&
- (this.session != null) &&
- #if (MF_FRAMEWORK_VERSION_V4_2 || MF_FRAMEWORK_VERSION_V4_3 || COMPACT_FRAMEWORK)
- (this.session.InflightMessages.Contains(msgContext.Key)))
- #else
- (this.session.InflightMessages.ContainsKey(msgContext.Key)))
- #endif
- {
- this.session.InflightMessages.Remove(msgContext.Key);
- }
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Queuing, "processed {0}", msgInflight);
- #endif
- }
- }
-
- else if ((msgReceived != null) && (msgReceived.Type == MqttMsgBase.MQTT_MSG_PUBREC_TYPE))
- {
-
-
- if (msgReceived.MessageId == msgInflight.MessageId)
- {
- lock (this.internalQueue)
- {
-
- this.internalQueue.Dequeue();
- acknowledge = true;
- msgReceivedProcessed = true;
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Queuing, "dequeued {0}", msgReceived);
- #endif
-
- this.inflightQueue.Enqueue(msgContext);
- }
- }
- }
-
- if (!acknowledge)
- {
- delta = Environment.TickCount - msgContext.Timestamp;
-
- if (delta >= this.settings.DelayOnRetry)
- {
-
- if (msgContext.Attempt < this.settings.AttemptsOnRetry)
- {
- msgContext.State = MqttMsgState.SendPubrel;
-
- this.inflightQueue.Enqueue(msgContext);
-
- timeout = 0;
- }
- else
- {
-
- if ((this.session != null) &&
- #if (MF_FRAMEWORK_VERSION_V4_2 || MF_FRAMEWORK_VERSION_V4_3 || COMPACT_FRAMEWORK)
- (this.session.InflightMessages.Contains(msgContext.Key)))
- #else
- (this.session.InflightMessages.ContainsKey(msgContext.Key)))
- #endif
- {
- this.session.InflightMessages.Remove(msgContext.Key);
- }
-
- internalEvent = new MsgPublishedInternalEvent(msgInflight, false);
-
- this.OnInternalEvent(internalEvent);
- }
- }
- else
- {
-
- this.inflightQueue.Enqueue(msgContext);
-
- int msgTimeout = (this.settings.DelayOnRetry - delta);
- timeout = (msgTimeout < timeout) ? msgTimeout : timeout;
- }
- }
- }
- break;
- case MqttMsgState.SendPubrec:
-
- break;
- case MqttMsgState.SendPubrel:
-
- if (msgContext.Flow == MqttMsgFlow.ToPublish)
- {
- MqttMsgPubrel pubrel = new MqttMsgPubrel();
- pubrel.MessageId = msgInflight.MessageId;
- msgContext.State = MqttMsgState.WaitForPubcomp;
- msgContext.Timestamp = Environment.TickCount;
- msgContext.Attempt++;
-
- if (this.ProtocolVersion == MqttProtocolVersion.Version_3_1)
- {
- if (msgContext.Attempt > 1)
- pubrel.DupFlag = true;
- }
-
- this.Send(pubrel);
-
- timeout = (this.settings.DelayOnRetry < timeout) ? this.settings.DelayOnRetry : timeout;
-
- this.inflightQueue.Enqueue(msgContext);
- }
- break;
- case MqttMsgState.SendPubcomp:
-
- break;
- case MqttMsgState.SendPuback:
-
- break;
- default:
- break;
- }
- }
-
- if (timeout == Int32.MaxValue)
- timeout = Timeout.Infinite;
-
-
- if ((msgReceived != null) && !msgReceivedProcessed)
- {
- this.internalQueue.Dequeue();
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Queuing, "dequeued {0} orphan", msgReceived);
- #endif
- }
- }
- }
- }
- }
- catch (MqttCommunicationException e)
- {
-
- if (msgContext != null)
-
- this.inflightQueue.Enqueue(msgContext);
- #if TRACE
- MqttUtility.Trace.WriteLine(TraceLevel.Error, "Exception occurred: {0}", e.ToString());
- #endif
-
- this.OnConnectionClosing();
- }
- }
-
-
-
- private void RestoreSession()
- {
-
- if (!this.CleanSession)
- {
-
- if (this.session != null)
- {
- lock (this.inflightQueue)
- {
- foreach (MqttMsgContext msgContext in this.session.InflightMessages.Values)
- {
- this.inflightQueue.Enqueue(msgContext);
-
- if ((msgContext.Message.Type == MqttMsgBase.MQTT_MSG_PUBLISH_TYPE) &&
- (msgContext.Flow == MqttMsgFlow.ToPublish))
- {
-
- if ((msgContext.Message.QosLevel == MqttMsgBase.QOS_LEVEL_AT_LEAST_ONCE) &&
- (msgContext.State == MqttMsgState.WaitForPuback))
- {
-
- msgContext.State = MqttMsgState.QueuedQos1;
- }
-
- else if (msgContext.Message.QosLevel == MqttMsgBase.QOS_LEVEL_EXACTLY_ONCE)
- {
-
- if (msgContext.State == MqttMsgState.WaitForPubrec)
- {
- msgContext.State = MqttMsgState.QueuedQos2;
- }
-
- else if (msgContext.State == MqttMsgState.WaitForPubcomp)
- {
- msgContext.State = MqttMsgState.SendPubrel;
- }
- }
- }
- }
- }
-
- this.inflightWaitHandle.Set();
- }
- else
- {
-
- this.session = new MqttClientSession(this.ClientId);
- }
- }
-
- else
- {
- if (this.session != null)
- this.session.Clear();
- }
- }
- #if BROKER
-
-
-
-
- public void LoadSession(MqttClientSession session)
- {
-
- if (!this.CleanSession)
- {
-
- this.session = session;
-
- this.RestoreSession();
- }
- }
- #endif
-
-
-
-
- private ushort GetMessageId()
- {
-
- this.messageIdCounter = ((this.messageIdCounter % UInt16.MaxValue) != 0) ? (ushort)(this.messageIdCounter + 1) : (ushort)1;
- return this.messageIdCounter;
- }
-
-
-
- internal class MqttMsgContextFinder
- {
-
- internal ushort MessageId { get; set; }
-
- internal MqttMsgFlow Flow { get; set; }
-
-
-
-
-
- internal MqttMsgContextFinder(ushort messageId, MqttMsgFlow flow)
- {
- this.MessageId = messageId;
- this.Flow = flow;
- }
- internal bool Find(object item)
- {
- MqttMsgContext msgCtx = (MqttMsgContext)item;
- return ((msgCtx.Message.Type == MqttMsgBase.MQTT_MSG_PUBLISH_TYPE) &&
- (msgCtx.Message.MessageId == this.MessageId) &&
- msgCtx.Flow == this.Flow);
- }
- }
- }
-
-
-
- public enum MqttProtocolVersion
- {
- Version_3_1 = MqttMsgConnect.PROTOCOL_VERSION_V3_1,
- Version_3_1_1 = MqttMsgConnect.PROTOCOL_VERSION_V3_1_1
- }
- }
|