mirror of
https://gitlab.com/wirelos/sprocket-plugin-mqtt.git
synced 2025-12-16 14:25:04 +01:00
basic mqtt plugin
This commit is contained in:
150
src/MqttPlugin.cpp
Normal file
150
src/MqttPlugin.cpp
Normal file
@@ -0,0 +1,150 @@
|
||||
#ifndef __MQTT_PLUGIN__
|
||||
#define __MQTT_PLUGIN__
|
||||
|
||||
#define _TASK_SLEEP_ON_IDLE_RUN
|
||||
#define _TASK_STD_FUNCTION
|
||||
|
||||
#include <WiFiClient.h>
|
||||
#include <PubSubClient.h>
|
||||
#include <Plugin.h>
|
||||
|
||||
using namespace std;
|
||||
using namespace std::placeholders;
|
||||
|
||||
struct MqttConfig
|
||||
{
|
||||
const char *clientName;
|
||||
const char *brokerHost;
|
||||
int brokerPort;
|
||||
const char *inTopicRoot;
|
||||
const char *outTopicRoot;
|
||||
};
|
||||
|
||||
class MqttPlugin : public Plugin
|
||||
{
|
||||
public:
|
||||
PubSubClient *client;
|
||||
WiFiClient wifiClient;
|
||||
Task connectTask;
|
||||
Task processTask;
|
||||
MqttConfig mqttConfig;
|
||||
vector<String> topics;
|
||||
|
||||
MqttPlugin(MqttConfig cfg)
|
||||
{
|
||||
mqttConfig = cfg;
|
||||
//subscribe("mqtt/out", bind(&MqttPlugin::mqttPublish, this, _1));
|
||||
}
|
||||
|
||||
void activate(Scheduler *scheduler)
|
||||
{
|
||||
client = new PubSubClient(mqttConfig.brokerHost, mqttConfig.brokerPort, bind(&MqttPlugin::downstreamHandler, this, _1, _2, _3), wifiClient);
|
||||
enableConnectTask(scheduler);
|
||||
enableProcessTask(scheduler);
|
||||
Serial.println("MQTT plugin activated"); // FIXME use PRINT_MSG
|
||||
}
|
||||
|
||||
private:
|
||||
void enableConnectTask(Scheduler *scheduler)
|
||||
{
|
||||
connectTask.set(TASK_SECOND * 5, TASK_FOREVER, bind(&MqttPlugin::connect, this));
|
||||
scheduler->addTask(connectTask);
|
||||
connectTask.enable();
|
||||
}
|
||||
|
||||
void enableProcessTask(Scheduler *scheduler)
|
||||
{
|
||||
processTask.set(TASK_MILLISECOND * 5, TASK_FOREVER, bind(&MqttPlugin::process, this));
|
||||
scheduler->addTask(processTask);
|
||||
processTask.enable();
|
||||
}
|
||||
|
||||
void process()
|
||||
{
|
||||
client->loop();
|
||||
}
|
||||
|
||||
void connect()
|
||||
{
|
||||
if (!client->connected())
|
||||
{
|
||||
Serial.println("MQTT try connect");
|
||||
if (client->connect(mqttConfig.clientName))
|
||||
{
|
||||
Serial.println("MQTT connected");
|
||||
upstreamHandler("lobby", "join");
|
||||
// overkill to subscribe to all...?
|
||||
for (auto const &localSub : mediator->subscriptions)
|
||||
{
|
||||
subscribe(localSub.first.c_str(), bind(&MqttPlugin::upstreamHandler, this, localSub.first.c_str(), _1));
|
||||
client->subscribe(
|
||||
(String(mqttConfig.inTopicRoot) + String(localSub.first.c_str()))
|
||||
.c_str());
|
||||
}
|
||||
}
|
||||
}
|
||||
else
|
||||
{
|
||||
publish("local/topic", "Free heap: " + String(ESP.getFreeHeap()));
|
||||
}
|
||||
}
|
||||
|
||||
void mqttPublish(String msg)
|
||||
{
|
||||
client->publish(mqttConfig.outTopicRoot, msg.c_str());
|
||||
}
|
||||
|
||||
void upstreamHandler(String topic, String msg)
|
||||
{
|
||||
Serial.println(msg);
|
||||
// publish message on remote queue
|
||||
// TODO rootTopic + topic?
|
||||
client->publish((String(mqttConfig.outTopicRoot) + topic).c_str(), msg.c_str());
|
||||
}
|
||||
|
||||
void downstreamHandler(char *topic, uint8_t *payload, unsigned int length)
|
||||
{
|
||||
char *cleanPayload = (char *)malloc(length + 1);
|
||||
payload[length] = '\0';
|
||||
memcpy(cleanPayload, payload, length + 1);
|
||||
String msg = String(cleanPayload);
|
||||
free(cleanPayload);
|
||||
|
||||
Serial.println("MQTT topic " + String(topic));
|
||||
Serial.println("MQTT msg " + msg);
|
||||
|
||||
int topicRootLength = String(mqttConfig.inTopicRoot).length();
|
||||
String localTopic = String(topic).substring(topicRootLength);
|
||||
Serial.println("Local topic:" + localTopic);
|
||||
|
||||
// publish message to local subscribers on device
|
||||
publish(localTopic, msg);
|
||||
}
|
||||
/* void mqttCallback(char* topic, uint8_t* payload, unsigned int length) {
|
||||
char* cleanPayload = (char*)malloc(length+1);
|
||||
payload[length] = '\0';
|
||||
memcpy(cleanPayload, payload, length+1);
|
||||
String msg = String(cleanPayload);
|
||||
free(cleanPayload);
|
||||
|
||||
int topicRootLength = String(mqttConfig.topicRoot).length();
|
||||
String targetStr = String(topic).substring(topicRootLength + 4);
|
||||
|
||||
if(targetStr == "gateway"){
|
||||
if(msg == "getNodes") {
|
||||
client->publish(MQTT_TOPIC_FROM_GATEWAY, net->mesh.subConnectionJson().c_str());
|
||||
}
|
||||
} else if(targetStr == "broadcast") {
|
||||
net->mesh.sendBroadcast(msg);
|
||||
} else {
|
||||
uint32_t target = strtoul(targetStr.c_str(), NULL, 10);
|
||||
if(net->mesh.isConnected(target)){
|
||||
net->mesh.sendSingle(target, msg);
|
||||
} else {
|
||||
client->publish(MQTT_TOPIC_FROM_GATEWAY, "Client not connected!");
|
||||
}
|
||||
}
|
||||
} */
|
||||
};
|
||||
|
||||
#endif
|
||||
36
src/examples/basic/config.h
Normal file
36
src/examples/basic/config.h
Normal file
@@ -0,0 +1,36 @@
|
||||
#ifndef __DEVICE_CONFIG__
|
||||
#define __DEVICE_CONFIG__
|
||||
|
||||
// Scheduler config
|
||||
#define _TASK_SLEEP_ON_IDLE_RUN
|
||||
#define _TASK_STD_FUNCTION
|
||||
#define _TASK_PRIORITY
|
||||
|
||||
// Chip config
|
||||
#define SPROCKET_TYPE "SPROCKET"
|
||||
#define SERIAL_BAUD_RATE 115200
|
||||
#define STARTUP_DELAY 1000
|
||||
|
||||
// network config
|
||||
#define SPROCKET_MODE 1
|
||||
#define WIFI_CHANNEL 11
|
||||
#define MESH_PORT 5555
|
||||
#define AP_SSID "sprocket"
|
||||
#define AP_PASSWORD "th3r31sn0sp00n"
|
||||
#define MESH_PREFIX "sprocket-mesh"
|
||||
#define MESH_PASSWORD "th3r31sn0sp00n"
|
||||
#define STATION_SSID "MyAP"
|
||||
#define STATION_PASSWORD "th3r31sn0sp00n"
|
||||
#define HOSTNAME "sprocket"
|
||||
#define CONNECT_TIMEOUT 10000
|
||||
#define MESH_DEBUG_TYPES ERROR | STARTUP | CONNECTION
|
||||
//#define MESH_DEBUG_TYPES ERROR | MESH_STATUS | CONNECTION | SYNC | COMMUNICATION | GENERAL | MSG_TYPES | REMOTE
|
||||
|
||||
// WebServer
|
||||
#define WEB_CONTEXT_PATH "/"
|
||||
#define WEB_DOC_ROOT "/www"
|
||||
#define WEB_DEFAULT_FILE "index.html"
|
||||
#define WEB_PORT 80
|
||||
|
||||
|
||||
#endif
|
||||
36
src/examples/basic/main.cpp
Normal file
36
src/examples/basic/main.cpp
Normal file
@@ -0,0 +1,36 @@
|
||||
#include "config.h"
|
||||
#include "WiFiNet.h"
|
||||
#include "Sprocket.h"
|
||||
#include "MqttPlugin.cpp"
|
||||
|
||||
WiFiNet *network;
|
||||
Sprocket *sprocket;
|
||||
|
||||
void setup()
|
||||
{
|
||||
sprocket = new Sprocket(
|
||||
{STARTUP_DELAY, SERIAL_BAUD_RATE});
|
||||
sprocket->subscribe("local/topic", [](String msg){
|
||||
Serial.println("Received: " + msg);
|
||||
});
|
||||
sprocket->addPlugin(
|
||||
new MqttPlugin({"sprocket", "192.168.1.2", 1883, "wirelos/mqttSprocket-in/", "wirelos/mqttSprocket-out/"}));
|
||||
|
||||
network = new WiFiNet(
|
||||
SPROCKET_MODE,
|
||||
STATION_SSID,
|
||||
STATION_PASSWORD,
|
||||
AP_SSID,
|
||||
AP_PASSWORD,
|
||||
HOSTNAME,
|
||||
CONNECT_TIMEOUT);
|
||||
network->connect();
|
||||
|
||||
sprocket->activate();
|
||||
}
|
||||
|
||||
void loop()
|
||||
{
|
||||
sprocket->loop();
|
||||
yield();
|
||||
}
|
||||
Reference in New Issue
Block a user