From ac6fd6ac0d45bc11f0183bf604f68042dfe9c812 Mon Sep 17 00:00:00 2001 From: Koushik Dutta Date: Mon, 20 Sep 2021 14:45:34 -0700 Subject: [PATCH] mqtt: broker logging --- plugins/mqtt/package-lock.json | 4 ++-- plugins/mqtt/package.json | 2 +- plugins/mqtt/src/main.ts | 27 +++++++++++++++++---------- 3 files changed, 20 insertions(+), 13 deletions(-) diff --git a/plugins/mqtt/package-lock.json b/plugins/mqtt/package-lock.json index 6c9ce8756..ddc5bed15 100644 --- a/plugins/mqtt/package-lock.json +++ b/plugins/mqtt/package-lock.json @@ -1,12 +1,12 @@ { "name": "@scrypted/mqtt", - "version": "0.0.11", + "version": "0.0.14", "lockfileVersion": 2, "requires": true, "packages": { "": { "name": "@scrypted/mqtt", - "version": "0.0.11", + "version": "0.0.14", "dependencies": { "@types/node": "^16.6.1", "aedes": "^0.46.1", diff --git a/plugins/mqtt/package.json b/plugins/mqtt/package.json index 06e7d4cc8..bb8f1279a 100644 --- a/plugins/mqtt/package.json +++ b/plugins/mqtt/package.json @@ -35,5 +35,5 @@ "devDependencies": { "@scrypted/sdk": "file:../../sdk" }, - "version": "0.0.11" + "version": "0.0.14" } diff --git a/plugins/mqtt/src/main.ts b/plugins/mqtt/src/main.ts index f985872c4..5cee195d3 100644 --- a/plugins/mqtt/src/main.ts +++ b/plugins/mqtt/src/main.ts @@ -28,7 +28,7 @@ class MqttDevice extends ScryptedDeviceBase implements Scriptable, Settings { client: Client; handler: any; pathname: string; - + constructor(nativeId: string) { super(nativeId); @@ -246,15 +246,22 @@ class MqttProvider extends ScryptedDeviceBase implements DeviceProvider, Setting this.netServer?.close(); if (this.storage.getItem('enableBroker') !== 'true') - return; - const instance = aedes(); - this.netServer = net.createServer(instance.handle); - const tcpPort = parseInt(this.storage.getItem('tcpPort')) || 1883; - const httpPort = parseInt(this.storage.getItem('httpPort')) || 8888; - this.netServer.listen(tcpPort); - this.httpServer = http.createServer(); - ws.createServer({ server: this.httpServer}).on('connection', instance.handle); - this.httpServer.listen(httpPort); + return; + const instance = aedes(); + this.netServer = net.createServer(instance.handle); + const tcpPort = parseInt(this.storage.getItem('tcpPort')) || 1883; + const httpPort = parseInt(this.storage.getItem('httpPort')) || 8888; + this.netServer.listen(tcpPort); + this.httpServer = http.createServer(); + ws.createServer({ server: this.httpServer }).on('connection', instance.handle); + this.httpServer.listen(httpPort); + + instance.on('publish', packet => { + if (!packet.payload) + return; + const preview = packet.payload.length > 512 ? '[large payload suppressed]' : packet.payload.toString(); + this.console.log('mqtt message', packet.topic, preview); + }); } async putSetting(key: string, value: string | number) {