Compare commits
35 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
20292b7b7b | ||
|
|
26de3befa8 | ||
|
|
1098466055 | ||
|
|
2d3663e78c | ||
|
|
5e16282ad0 | ||
|
|
e35a43c4a4 | ||
|
|
087a203ba0 | ||
|
|
5d313bbf5e | ||
|
|
ce896f02c4 | ||
|
|
d3210c3c93 | ||
|
|
23f1207718 | ||
|
|
122ab88960 | ||
|
|
28b0ac1611 | ||
|
|
1cfb5cfab1 | ||
|
|
b023cd67a9 | ||
|
|
24ee6b5201 | ||
|
|
2e92a98db2 | ||
|
|
c59bddfd39 | ||
|
|
7bdb9cc0cd | ||
|
|
be62699702 | ||
|
|
77da47e1da | ||
|
|
88797bfafd | ||
|
|
1e3b37623d | ||
|
|
ba6a96976a | ||
|
|
6afd3939b3 | ||
|
|
2ffe0c13fa | ||
|
|
48eb0daf9a | ||
|
|
34c05bc37a | ||
|
|
7c96c4a5cc | ||
|
|
b280196395 | ||
|
|
c75f4893e8 | ||
|
|
d666f6a53b | ||
|
|
7ef18de755 | ||
|
|
838df3a34a | ||
|
|
8a25155fd8 |
18
.github/workflows/ci.yml
vendored
Normal file
18
.github/workflows/ci.yml
vendored
Normal file
@@ -0,0 +1,18 @@
|
||||
name: "CI"
|
||||
on:
|
||||
|
||||
jobs:
|
||||
ci:
|
||||
runs-on: ubuntu-20.04
|
||||
steps:
|
||||
- name: Checkout this repository
|
||||
uses: actions/checkout@v2.3.4
|
||||
- name: Cache for arduino-ci
|
||||
uses: actions/cache@v2.1.3
|
||||
with:
|
||||
path: |
|
||||
~/.arduino15
|
||||
key: ${{ runner.os }}-arduino
|
||||
- name: Install nix
|
||||
uses: cachix/install-nix-action@v12
|
||||
- run: nix-shell -I nixpkgs=channel:nixpkgs-unstable -p arduino-ci --run "arduino-ci"
|
||||
@@ -1,17 +1,19 @@
|
||||
# TinyMqtt
|
||||
|
||||

|
||||
[](https://github.com/hsaturn/TinyMqtt/actions/workflows/aunit.yml)
|
||||

|
||||

|
||||

|
||||

|
||||

|
||||
|
||||
ESP 8266 is a small, fast and capable Mqtt Broker and Client
|
||||
TinyMqtt is a small, fast and capable Mqtt Broker and Client for Esp8266 / Esp32 / Esp WROOM
|
||||
|
||||
## Features
|
||||
|
||||
- Very (very !!) fast broker I saw it re-sent 1000 topics per second for two
|
||||
clients that had subscribed (payload ~15 bytes). No topic lost.
|
||||
clients that had subscribed (payload ~15 bytes ESP8266). No topic lost.
|
||||
The max I've seen was 2k msg/s (1 client 1 subscription)
|
||||
- Act as as a mqtt broker and/or a mqtt client
|
||||
- Mqtt 3.1.1 / Qos 0 supported
|
||||
@@ -48,7 +50,7 @@ After a while one ESP naturally becomes a 'master' and all ESP are connected tog
|
||||
No problem if the master dies, a new master will be choosen soon.
|
||||
|
||||
## TODO List
|
||||
* Use [Async library](https://github.com/me-no-dev/ESPAsyncTCP)
|
||||
* ~~Use [Async library](https://github.com/me-no-dev/ESPAsyncTCP)~~
|
||||
* Implement zeroconf mode (needs async)
|
||||
* Add a max_clients in MqttBroker. Used with zeroconf, there will be
|
||||
no need for having tons of clients (also RAM is the problem with many clients)
|
||||
|
||||
@@ -3,6 +3,23 @@
|
||||
/**
|
||||
* Local broker that accept connections and two local clients
|
||||
*
|
||||
*
|
||||
* +-----------------------------+
|
||||
* | ESP |
|
||||
* | +--------+ | 1883 <--- External client/s
|
||||
* | +-------->| broker | | 1883 <--- External client/s
|
||||
* | | +--------+ |
|
||||
* | | ^ |
|
||||
* | | | |
|
||||
* | | | | -----
|
||||
* | v v | ---
|
||||
* | +----------+ +----------+ | -
|
||||
* | | internal | | internal | +-------* Wifi
|
||||
* | | client | | client | |
|
||||
* | +----------+ +----------+ |
|
||||
* | |
|
||||
* +-----------------------------+
|
||||
*
|
||||
* pros - Reduces internal latency (when publish is received by the same ESP)
|
||||
* - Reduces wifi traffic
|
||||
* - No need to have an external broker
|
||||
@@ -12,11 +29,6 @@
|
||||
* cons - Takes more memory
|
||||
* - a bit hard to understand
|
||||
*
|
||||
* This sounds crazy: a mqtt mqtt that do not need a broker !
|
||||
* The use case arise when one ESP wants to publish topics and subscribe to them at the same time.
|
||||
* Without broker, the ESP won't react to its own topics.
|
||||
*
|
||||
* TinyMqtt mqtt allows this use case to work.
|
||||
*/
|
||||
|
||||
#include <my_credentials.h>
|
||||
@@ -58,7 +70,7 @@ void setup()
|
||||
|
||||
void loop()
|
||||
{
|
||||
broker.loop();
|
||||
broker.loop(); // Don't forget to add loop for every broker and clients
|
||||
|
||||
mqtt_a.loop();
|
||||
mqtt_b.loop();
|
||||
|
||||
@@ -1,11 +1,29 @@
|
||||
#include <TinyMqtt.h> // https://github.com/hsaturn/TinyMqtt
|
||||
|
||||
/** TinyMQTT allows a disconnected mode:
|
||||
*
|
||||
* +-----------------------------+
|
||||
* | ESP |
|
||||
* | +--------+ |
|
||||
* | +-------->| broker | |
|
||||
* | | +--------+ |
|
||||
* | | ^ |
|
||||
* | | | |
|
||||
* | v v |
|
||||
* | +----------+ +----------+ |
|
||||
* | | internal | | internal | |
|
||||
* | | client | | client | |
|
||||
* | +----------+ +----------+ |
|
||||
* | |
|
||||
* +-----------------------------+
|
||||
*
|
||||
* In this example, local clients A and B are talking together, no need to be connected.
|
||||
*
|
||||
* A single ESP can use this to be able to comunicate with itself with the power
|
||||
* of MQTT, and once connected still continue to work with others.
|
||||
*
|
||||
* The broker may still be conected if wifi is on.
|
||||
*
|
||||
*/
|
||||
|
||||
std::string topic="sensor/temperature";
|
||||
|
||||
@@ -5,6 +5,16 @@
|
||||
#define PORT 1883
|
||||
MqttBroker broker(PORT);
|
||||
|
||||
/** Basic Mqtt Broker
|
||||
*
|
||||
* +-----------------------------+
|
||||
* | ESP |
|
||||
* | +--------+ |
|
||||
* | | broker | | 1883 <--- External client/s
|
||||
* | +--------+ |
|
||||
* | |
|
||||
* +-----------------------------+
|
||||
*/
|
||||
void setup()
|
||||
{
|
||||
Serial.begin(115200);
|
||||
|
||||
@@ -1,9 +1,23 @@
|
||||
#include <ESP8266WiFi.h>
|
||||
#include "TinyMqtt.h" // https://github.com/hsaturn/TinyMqtt
|
||||
|
||||
/** Simple Client
|
||||
/** Simple Client (The simplest configuration)
|
||||
*
|
||||
* This is the simplest Mqtt client configuration
|
||||
*
|
||||
* +--------+
|
||||
* +------>| broker |<--- < Other client
|
||||
* | +--------+
|
||||
* |
|
||||
* +-----------------+
|
||||
* | ESP | |
|
||||
* | +----------+ |
|
||||
* | | internal | |
|
||||
* | | client | |
|
||||
* | +----------+ |
|
||||
* | |
|
||||
* +-----------------+
|
||||
*
|
||||
* 1 - edit my_credentials.h to setup wifi essid/password
|
||||
* 2 - change BROKER values (or keep emqx.io test broker)
|
||||
*
|
||||
* pro - small memory footprint (both ram and flash)
|
||||
* - very simple to setup and use
|
||||
@@ -13,6 +27,9 @@
|
||||
* - local publishes takes more time (because they go outside)
|
||||
*/
|
||||
|
||||
const char* BROKER = "broker.emqx.io";
|
||||
const uint16_t BROKER_PORT = 1883;
|
||||
|
||||
#include <my_credentials.h>
|
||||
|
||||
static float temp=19;
|
||||
@@ -32,8 +49,7 @@ void setup()
|
||||
|
||||
Serial << "Connected to " << ssid << "IP address: " << WiFi.localIP() << endl;
|
||||
|
||||
client.connect("192.168.1.40", 1883); // Put here your broker ip / port
|
||||
|
||||
client.connect(BROKER, BROKER_PORT); // Put here your broker ip / port
|
||||
}
|
||||
|
||||
void loop()
|
||||
|
||||
@@ -1,22 +1,25 @@
|
||||
#include <TinyMqtt.h> // https://github.com/hsaturn/TinyMqtt
|
||||
#include <MqttStreaming.h>
|
||||
#if defined(ESP8266)
|
||||
#include <ESP8266mDNS.h>
|
||||
#elif defined(ESP32)
|
||||
#include <ESPmDNS.h>
|
||||
#else
|
||||
#error Unsupported platform
|
||||
#endif
|
||||
#include <ESPmDNS.h>
|
||||
|
||||
#include <sstream>
|
||||
#include <map>
|
||||
|
||||
/**
|
||||
/** Very complex example
|
||||
* Console allowing to make any kind of test.
|
||||
*
|
||||
* pros - Reduces internal latency (when publish is received by the same ESP)
|
||||
* - Reduces wifi traffic
|
||||
* - No need to have an external broker
|
||||
* - can still report to a 'main' broker (TODO see documentation that have to be written)
|
||||
* - accepts external clients
|
||||
*
|
||||
* cons - Takes more memory
|
||||
* - a bit hard to understand
|
||||
* Upload the sketch, the use the terminal.
|
||||
* Press H for mini help.
|
||||
*
|
||||
* tested with mqtt-spy-0.5.4
|
||||
* TODO examples of scripts
|
||||
*/
|
||||
|
||||
#include <my_credentials.h>
|
||||
@@ -138,7 +141,7 @@ std::set<std::string> commands = {
|
||||
"auto", "broker", "blink", "client", "connect",
|
||||
"create", "delete", "help", "interval",
|
||||
"ls", "ip", "off", "on", "set",
|
||||
"publish", "reset", "subscribe", "unsubscribe", "view"
|
||||
"publish", "reset", "subscribe", "unsubscribe", "view", "every"
|
||||
};
|
||||
|
||||
void getCommand(std::string& search)
|
||||
@@ -316,74 +319,44 @@ std::map<MqttClient*, automatic*> automatic::autos;
|
||||
|
||||
bool compare(std::string s, const char* cmd)
|
||||
{
|
||||
if (s.length()==0 or s.length()>strlen(cmd)) return false;
|
||||
return strncmp(cmd, s.c_str(), s.length())==0;
|
||||
uint8_t p=0;
|
||||
while(s[p++]==*cmd++)
|
||||
{
|
||||
if (*cmd==0 or s[p]==0) return true;
|
||||
if (s[p]==' ') return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
using ClientFunction = void(*)(std::string& cmd, MqttClient* publish);
|
||||
|
||||
struct Every
|
||||
{
|
||||
std::string cmd;
|
||||
uint32_t ms;
|
||||
uint32_t next;
|
||||
|
||||
void dump()
|
||||
{
|
||||
auto mill=millis();
|
||||
Serial << ms << "ms [" << cmd << "] next in ";
|
||||
if (mill > next)
|
||||
Serial << "now";
|
||||
else
|
||||
Serial << next-mill << "ms";
|
||||
}
|
||||
};
|
||||
|
||||
uint32_t blink_ms_on[16];
|
||||
uint32_t blink_ms_off[16];
|
||||
uint32_t blink_next[16];
|
||||
bool blink_state[16];
|
||||
int16_t blink;
|
||||
void loop()
|
||||
{
|
||||
auto ms=millis();
|
||||
int8_t out=1;
|
||||
int16_t blink_bits = blink;
|
||||
while(blink_bits)
|
||||
{
|
||||
if (blink_ms_on[out] and ms > blink_next[out])
|
||||
{
|
||||
if (blink_state[out])
|
||||
{
|
||||
blink_next[out] += blink_ms_on[out];
|
||||
digitalWrite(out, LOW);
|
||||
}
|
||||
else
|
||||
{
|
||||
blink_next[out] += blink_ms_off[out];
|
||||
digitalWrite(abs(out), HIGH);
|
||||
}
|
||||
blink_state[out] = not blink_state[out];
|
||||
}
|
||||
blink_bits >>=1;
|
||||
out++;
|
||||
}
|
||||
|
||||
static long count;
|
||||
MDNS.update();
|
||||
std::vector<Every> everies;
|
||||
|
||||
if (MqttClient::counter != count)
|
||||
void eval(std::string& cmd)
|
||||
{
|
||||
Serial << "# " << MqttClient::counter << endl;
|
||||
count = MqttClient::counter;
|
||||
}
|
||||
for(auto it: brokers)
|
||||
it.second->loop();
|
||||
|
||||
for(auto it: clients)
|
||||
it.second->loop();
|
||||
|
||||
automatic::loop();
|
||||
|
||||
if (Serial.available())
|
||||
{
|
||||
static std::string cmd;
|
||||
char c=Serial.read();
|
||||
|
||||
if (c==10 or c==14)
|
||||
{
|
||||
|
||||
Serial << "----------------[ " << cmd.c_str() << " ]--------------" << endl;
|
||||
static std::string last_cmd;
|
||||
if (cmd=="!")
|
||||
cmd=last_cmd;
|
||||
else
|
||||
last_cmd=cmd;
|
||||
|
||||
if (cmd.substr(0,3)!="set") replaceVars(cmd);
|
||||
while(cmd.length())
|
||||
{
|
||||
MqttError retval = MqttOk;
|
||||
@@ -501,10 +474,64 @@ void loop()
|
||||
client->dump();
|
||||
}
|
||||
}
|
||||
else if (compare(s, "on"))
|
||||
{
|
||||
uint8_t pin=getint(cmd, 2);
|
||||
pinMode(pin, OUTPUT);
|
||||
digitalWrite(pin, HIGH);
|
||||
}
|
||||
else if (compare(s, "off"))
|
||||
{
|
||||
uint8_t pin=getint(cmd, 2);
|
||||
pinMode(pin, OUTPUT);
|
||||
digitalWrite(pin, LOW);
|
||||
}
|
||||
else if (compare(s, "every"))
|
||||
{
|
||||
uint32_t ms = getint(cmd, 0);
|
||||
if (ms and cmd.length())
|
||||
{
|
||||
Every every;
|
||||
every.ms=ms;
|
||||
every.cmd=cmd;
|
||||
every.next=millis()+ms;
|
||||
everies.push_back(every);
|
||||
every.dump();
|
||||
Serial << endl;
|
||||
cmd="";
|
||||
}
|
||||
else if (ms==0 and compare(cmd, "list"))
|
||||
{
|
||||
getword(cmd);
|
||||
Serial << "List of everies (ms=" << millis() << ")" << endl;
|
||||
uint8_t count=0;
|
||||
for(auto& every: everies)
|
||||
{
|
||||
Serial << count << ": ";
|
||||
every.dump();
|
||||
Serial << endl;
|
||||
count++;
|
||||
}
|
||||
}
|
||||
else if (ms==0 and compare(cmd, "remove"))
|
||||
{
|
||||
getword(cmd);
|
||||
int8_t every=getint(cmd, -1);
|
||||
if (every==-1 and compare(cmd, "all"))
|
||||
{
|
||||
getword(cmd);
|
||||
everies.clear();
|
||||
}
|
||||
else if (everies.size() > (uint8_t)every)
|
||||
{
|
||||
everies.erase(everies.begin()+every);
|
||||
}
|
||||
}
|
||||
}
|
||||
else if (compare(s, "blink"))
|
||||
{
|
||||
uint8_t blink_nr = getint(cmd, 0);
|
||||
if (blink_nr)
|
||||
int8_t blink_nr = getint(cmd, -1);
|
||||
if (blink_nr >= 0)
|
||||
{
|
||||
blink_ms_on[blink_nr]=getint(cmd, blink_ms_on[blink_nr]);
|
||||
blink_ms_off[blink_nr]=getint(cmd, blink_ms_on[blink_nr]);
|
||||
@@ -512,10 +539,10 @@ void loop()
|
||||
blink_next[blink_nr] = millis();
|
||||
Serial << "Blink " << blink_nr << ' ' << (blink_ms_on[blink_nr] ? "on" : "off") << endl;
|
||||
if (blink_ms_on[blink_nr])
|
||||
blink |= 1<< (blink_nr-1);
|
||||
blink |= 1<< blink_nr;
|
||||
else
|
||||
{
|
||||
blink &= ~(1<<(blink_nr-1));
|
||||
blink &= ~(1<< blink_nr);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -639,6 +666,8 @@ void loop()
|
||||
Serial << " set [name][value]" << endl;
|
||||
Serial << " ! repeat last command" << endl;
|
||||
Serial << endl;
|
||||
Serial << " every ms [command]; every list; every remove [nr|all]" << endl;
|
||||
Serial << " on {output}; off {output}" << endl;
|
||||
Serial << " $id : name of the client." << endl;
|
||||
Serial << " default topic is '" << topic.c_str() << "'" << endl;
|
||||
Serial << endl;
|
||||
@@ -656,6 +685,78 @@ void loop()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void loop()
|
||||
{
|
||||
auto ms=millis();
|
||||
int8_t out=0;
|
||||
int16_t blink_bits = blink;
|
||||
|
||||
for(auto& every: everies)
|
||||
{
|
||||
if (every.ms && every.cmd.length() && ms > every.next)
|
||||
{
|
||||
std::string cmd(every.cmd);
|
||||
eval(cmd);
|
||||
every.next += every.ms;
|
||||
}
|
||||
}
|
||||
|
||||
while(blink_bits)
|
||||
{
|
||||
if (blink_ms_on[out] and ms > blink_next[out])
|
||||
{
|
||||
if (blink_state[out])
|
||||
{
|
||||
blink_next[out] += blink_ms_on[out];
|
||||
digitalWrite(out, LOW);
|
||||
}
|
||||
else
|
||||
{
|
||||
blink_next[out] += blink_ms_off[out];
|
||||
digitalWrite(abs(out), HIGH);
|
||||
}
|
||||
blink_state[out] = not blink_state[out];
|
||||
}
|
||||
blink_bits >>=1;
|
||||
out++;
|
||||
}
|
||||
|
||||
static long count;
|
||||
#if defined(ESP9266)
|
||||
MDNS.update();
|
||||
#endif
|
||||
|
||||
if (MqttClient::counter != count)
|
||||
{
|
||||
Serial << "# " << MqttClient::counter << endl;
|
||||
count = MqttClient::counter;
|
||||
}
|
||||
for(auto it: brokers)
|
||||
it.second->loop();
|
||||
|
||||
for(auto it: clients)
|
||||
it.second->loop();
|
||||
|
||||
automatic::loop();
|
||||
|
||||
if (Serial.available())
|
||||
{
|
||||
static std::string cmd;
|
||||
char c=Serial.read();
|
||||
|
||||
if (c==10 or c==14)
|
||||
{
|
||||
Serial << "----------------[ " << cmd.c_str() << " ]--------------" << endl;
|
||||
static std::string last_cmd;
|
||||
if (cmd=="!")
|
||||
cmd=last_cmd;
|
||||
else
|
||||
last_cmd=cmd;
|
||||
|
||||
if (cmd.substr(0,3)!="set") replaceVars(cmd);
|
||||
eval(cmd);
|
||||
}
|
||||
else
|
||||
{
|
||||
cmd=cmd+c;
|
||||
|
||||
@@ -1,12 +1,12 @@
|
||||
{
|
||||
"name": "TinyMqtt",
|
||||
"keywords": "ethernet, mqtt, m2m, iot",
|
||||
"description": "MQTT is a lightweight messaging protocol ideal for small devices. This library allows to send and receive MQTT messages. It does support MQTT 3.1.1 with QOS=0.",
|
||||
"description": "MQTT is a lightweight messaging protocol ideal for small devices. This library allows to send and receive and host a broker for MQTT. It does support MQTT 3.1.1 with QOS=0 on ESP8266 and ESP32 WROOM platfrms.",
|
||||
"repository": {
|
||||
"type": "git",
|
||||
"url": "https://github.com/hsaturn/TinyMqtt.git"
|
||||
},
|
||||
"version": "0.7.3",
|
||||
"version": "0.7.4",
|
||||
"exclude": "",
|
||||
"examples": "examples/*/*.ino",
|
||||
"frameworks": "arduino",
|
||||
|
||||
@@ -1,9 +1,9 @@
|
||||
name=TinyMqtt
|
||||
version=0.7.3
|
||||
version=0.7.4
|
||||
author=Francois BIOT, HSaturn, <hsaturn@gmail.com>
|
||||
maintainer=Francois BIOT, HSaturn, <hsaturn@gmail.com>
|
||||
sentence=A tiny broker and client library for MQTT messaging.
|
||||
paragraph=MQTT is a lightweight messaging protocol ideal for small devices. This library allows to send and receive MQTT messages and to host a broker in your ESP. It does support MQTT 3.1.1 with QoS=0.
|
||||
paragraph=MQTT is a lightweight messaging protocol ideal for small devices. This library allows to send and receive MQTT messages and to host a broker in your ESP 8266 and 32 WROOM. It does support MQTT 3.1.1 with QoS=0.
|
||||
category=Communication
|
||||
url=https://github.com/hsaturn/TinyMqtt
|
||||
architectures=*
|
||||
|
||||
102
src/TinyMqtt.cpp
102
src/TinyMqtt.cpp
@@ -11,8 +11,10 @@ void outstring(const char* prefix, const char*p, uint16_t len)
|
||||
|
||||
MqttBroker::MqttBroker(uint16_t port)
|
||||
{
|
||||
server = new AsyncServer(port);
|
||||
server = new TcpServer(port);
|
||||
#ifdef TCP_ASYNC
|
||||
server->onClient(onClient, this);
|
||||
#endif
|
||||
}
|
||||
|
||||
MqttBroker::~MqttBroker()
|
||||
@@ -25,12 +27,17 @@ MqttBroker::~MqttBroker()
|
||||
}
|
||||
|
||||
// private constructor used by broker only
|
||||
MqttClient::MqttClient(MqttBroker* parent, AsyncClient* new_client)
|
||||
: parent(parent), client(new_client)
|
||||
MqttClient::MqttClient(MqttBroker* parent, TcpClient* new_client)
|
||||
: parent(parent)
|
||||
{
|
||||
#ifdef TCP_ASYNC
|
||||
client = new_client;
|
||||
client->onData(onData, this);
|
||||
// client->onConnect() TODO
|
||||
// client->onDisconnect() TODO
|
||||
#else
|
||||
client = new WiFiClient(*new_client);
|
||||
#endif
|
||||
alive = millis()+5000; // client expires after 5s if no CONNECT msg
|
||||
}
|
||||
|
||||
@@ -72,33 +79,22 @@ void MqttClient::close(bool bSendDisconnect)
|
||||
void MqttClient::connect(std::string broker, uint16_t port, uint16_t ka)
|
||||
{
|
||||
debug("cnx: closing");
|
||||
keep_alive = ka;
|
||||
close();
|
||||
if (client) delete client;
|
||||
client = new AsyncClient;
|
||||
client = new TcpClient;
|
||||
|
||||
debug("Trying to connect to " << broker.c_str() << ':' << port);
|
||||
// TODO This may return immediately !!!
|
||||
// TODO so I have to add onConnect and move this code to onConnect
|
||||
// TODO also, as this is async now, I must take care of
|
||||
// TODO the broker that may disconnect and delete the client immediately
|
||||
#ifdef TCP_ASYNC
|
||||
client->onData(onData, this);
|
||||
client->onConnect(onConnect, this);
|
||||
client->connect(broker.c_str(), port);
|
||||
#else
|
||||
if (client->connect(broker.c_str(), port))
|
||||
{
|
||||
debug("cnx: connecting");
|
||||
MqttMessage msg(MqttMessage::Type::Connect);
|
||||
msg.add("MQTT",4);
|
||||
msg.add(0x4); // Mqtt protocol version 3.1.1
|
||||
msg.add(0x0); // Connect flags TODO user / name
|
||||
|
||||
keep_alive = ka;
|
||||
msg.add(0x00); // keep_alive
|
||||
msg.add((char)keep_alive);
|
||||
msg.add(clientId);
|
||||
debug("cnx: mqtt connecting");
|
||||
msg.sendTo(this);
|
||||
msg.reset();
|
||||
debug("cnx: mqtt sent " << (int32_t)parent);
|
||||
|
||||
clientAlive(0);
|
||||
onConnect(this, client);
|
||||
}
|
||||
#endif
|
||||
}
|
||||
|
||||
void MqttBroker::addClient(MqttClient* client)
|
||||
@@ -134,16 +130,24 @@ void MqttBroker::removeClient(MqttClient* remove)
|
||||
debug("Error cannot remove client"); // TODO should not occur
|
||||
}
|
||||
|
||||
void MqttBroker::onClient(void* broker_ptr, AsyncClient* client)
|
||||
void MqttBroker::onClient(void* broker_ptr, TcpClient* client)
|
||||
{
|
||||
MqttBroker* broker = static_cast<MqttBroker*>(broker_ptr);
|
||||
|
||||
broker->addClient(new MqttClient(broker, client));
|
||||
debug("New client #" << broker->clients.size());
|
||||
debug("New client");
|
||||
}
|
||||
|
||||
void MqttBroker::loop()
|
||||
{
|
||||
#ifndef TCP_ASYNC
|
||||
WiFiClient client = server->available();
|
||||
|
||||
if (client)
|
||||
{
|
||||
onClient(this, &client);
|
||||
}
|
||||
#endif
|
||||
if (broker)
|
||||
{
|
||||
// TODO should monitor broker's activity.
|
||||
@@ -163,7 +167,7 @@ void MqttBroker::loop()
|
||||
}
|
||||
else
|
||||
{
|
||||
debug("Client " << client->id().c_str() << " Disconnected, parent=" << (int32_t)client->parent);
|
||||
debug("Client " << client->id().c_str() << " Disconnected, parent=" << (dbg_ptr)client->parent);
|
||||
// Note: deleting a client not added by the broker itself will probably crash later.
|
||||
delete client;
|
||||
break;
|
||||
@@ -263,13 +267,45 @@ void MqttClient::loop()
|
||||
client->write((const char*)(&pingreq), 2);
|
||||
clientAlive(0);
|
||||
|
||||
// TODO when many MqttClient passes through a local browser
|
||||
// TODO when many MqttClient passes through a local broker
|
||||
// there is no need to send one PingReq per instance.
|
||||
}
|
||||
}
|
||||
#ifndef TCP_ASYNC
|
||||
while(client && client->available()>0)
|
||||
{
|
||||
message.incoming(client->read());
|
||||
if (message.type())
|
||||
{
|
||||
processMessage(&message);
|
||||
message.reset();
|
||||
}
|
||||
}
|
||||
#endif
|
||||
}
|
||||
|
||||
void MqttClient::onData(void* client_ptr, AsyncClient*, void* data, size_t len)
|
||||
void MqttClient::onConnect(void *mqttclient_ptr, TcpClient*)
|
||||
{
|
||||
MqttClient* mqtt = static_cast<MqttClient*>(mqttclient_ptr);
|
||||
debug("cnx: connecting");
|
||||
MqttMessage msg(MqttMessage::Type::Connect);
|
||||
msg.add("MQTT",4);
|
||||
msg.add(0x4); // Mqtt protocol version 3.1.1
|
||||
msg.add(0x0); // Connect flags TODO user / name
|
||||
|
||||
msg.add(0x00); // keep_alive
|
||||
msg.add((char)mqtt->keep_alive);
|
||||
msg.add(mqtt->clientId);
|
||||
debug("cnx: mqtt connecting");
|
||||
msg.sendTo(mqtt);
|
||||
msg.reset();
|
||||
debug("cnx: mqtt sent " << (dbg_ptr)mqtt->parent);
|
||||
|
||||
mqtt->clientAlive(0);
|
||||
}
|
||||
|
||||
#ifdef TCP_ASYNC
|
||||
void MqttClient::onData(void* client_ptr, TcpClient*, void* data, size_t len)
|
||||
{
|
||||
char* char_ptr = static_cast<char*>(data);
|
||||
MqttClient* client=static_cast<MqttClient*>(client_ptr);
|
||||
@@ -284,6 +320,7 @@ void MqttClient::onData(void* client_ptr, AsyncClient*, void* data, size_t len)
|
||||
len--;
|
||||
}
|
||||
}
|
||||
#endif
|
||||
|
||||
void MqttClient::resubscribe()
|
||||
{
|
||||
@@ -360,7 +397,7 @@ void MqttClient::processMessage(const MqttMessage* mesg)
|
||||
#ifdef TINY_MQTT_DEBUG
|
||||
if (mesg->type() != MqttMessage::Type::PingReq && mesg->type() != MqttMessage::Type::PingResp)
|
||||
{
|
||||
Serial << "---> INCOMING " << _HEX(mesg->type()) << " client(" << (int)client << ':' << clientId << ") mem=" << ESP.getFreeHeap() << endl;
|
||||
Serial << "---> INCOMING " << _HEX(mesg->type()) << " client(" << (dbg_ptr)client << ':' << clientId << ") mem=" << ESP.getFreeHeap() << endl;
|
||||
// mesg->hexdump("Incoming");
|
||||
}
|
||||
#endif
|
||||
@@ -497,6 +534,11 @@ if (mesg->type() != MqttMessage::Type::PingReq && mesg->type() != MqttMessage::T
|
||||
}
|
||||
break;
|
||||
|
||||
case MqttMessage::Type::UnSuback:
|
||||
if (!mqtt_connected) break;
|
||||
bclose = false;
|
||||
break;
|
||||
|
||||
case MqttMessage::Type::Publish:
|
||||
if (mqtt_connected or client == nullptr)
|
||||
{
|
||||
|
||||
@@ -1,6 +1,26 @@
|
||||
#pragma once
|
||||
#include <ESP8266WiFi.h>
|
||||
|
||||
// TODO Should add a AUnit with both TCP_ASYNC and not TCP_ASYNC
|
||||
// #define TCP_ASYNC // Uncomment this to use ESPAsyncTCP instead of normal cnx
|
||||
|
||||
#if defined(ESP8266) || defined(EPOXY_DUINO)
|
||||
#ifdef TCP_ASYNC
|
||||
#include <ESPAsyncTCP.h>
|
||||
#else
|
||||
#include <ESP8266WiFi.h>
|
||||
#endif
|
||||
#elif defined(ESP32)
|
||||
#ifdef TCP_ASYNC
|
||||
#include <AsyncTCP.h> // https://github.com/me-no-dev/AsyncTCP
|
||||
#else
|
||||
#include <WiFi.h>
|
||||
#endif
|
||||
#endif
|
||||
#ifdef EPOXY_DUINO
|
||||
#define dbg_ptr uint64_t
|
||||
#else
|
||||
#define dbg_ptr uint32_t
|
||||
#endif
|
||||
#include <vector>
|
||||
#include <set>
|
||||
#include <string>
|
||||
@@ -15,6 +35,14 @@
|
||||
#define debug(what) {}
|
||||
#endif
|
||||
|
||||
#ifdef TCP_ASYNC
|
||||
using TcpClient = AsyncClient;
|
||||
using TcpServer = AsyncServer;
|
||||
#else
|
||||
using TcpClient = WiFiClient;
|
||||
using TcpServer = WiFiServer;
|
||||
#endif
|
||||
|
||||
enum MqttError
|
||||
{
|
||||
MqttOk = 0,
|
||||
@@ -49,6 +77,7 @@ class MqttMessage
|
||||
Subscribe = 0x80,
|
||||
SubAck = 0x90,
|
||||
UnSubscribe = 0xA0,
|
||||
UnSuback = 0xB0,
|
||||
PingReq = 0xC0,
|
||||
PingResp = 0xD0,
|
||||
Disconnect = 0xE0
|
||||
@@ -183,12 +212,15 @@ class MqttClient
|
||||
static long counter;
|
||||
|
||||
private:
|
||||
static void onData(void* client_ptr, AsyncClient*, void* data, size_t len);
|
||||
static void onConnect(void * client_ptr, TcpClient*);
|
||||
#ifdef TCP_ASYNC
|
||||
static void onData(void* client_ptr, TcpClient*, void* data, size_t len);
|
||||
#endif
|
||||
MqttError sendTopic(const Topic& topic, MqttMessage::Type type, uint8_t qos);
|
||||
void resubscribe();
|
||||
|
||||
friend class MqttBroker;
|
||||
MqttClient(MqttBroker* parent, AsyncClient* client);
|
||||
MqttClient(MqttBroker* parent, TcpClient* client);
|
||||
// republish a received publish if topic matches any in subscriptions
|
||||
MqttError publishIfSubscribed(const Topic& topic, const MqttMessage& msg);
|
||||
|
||||
@@ -206,7 +238,7 @@ class MqttClient
|
||||
// (this is the case when MqttBroker isn't used except here)
|
||||
MqttBroker* parent=nullptr; // connection to local broker
|
||||
|
||||
AsyncClient* client=nullptr; // connection to mqtt client or to remote broker
|
||||
TcpClient* client=nullptr; // connection to mqtt client or to remote broker
|
||||
std::set<Topic> subscriptions;
|
||||
std::string clientId;
|
||||
CallBack callback = nullptr;
|
||||
@@ -246,7 +278,7 @@ class MqttBroker
|
||||
private:
|
||||
friend class MqttClient;
|
||||
|
||||
static void onClient(void*, AsyncClient*);
|
||||
static void onClient(void*, TcpClient*);
|
||||
bool checkUser(const char* user, uint8_t len) const
|
||||
{ return compareString(auth_user, user, len); }
|
||||
|
||||
@@ -264,7 +296,7 @@ class MqttBroker
|
||||
|
||||
bool compareString(const char* good, const char* str, uint8_t str_len) const;
|
||||
std::vector<MqttClient*> clients;
|
||||
AsyncServer* server;
|
||||
TcpServer* server;
|
||||
|
||||
const char* auth_user = "guest";
|
||||
const char* auth_password = "guest";
|
||||
|
||||
Reference in New Issue
Block a user