some changes

This commit is contained in:
Wolfgang Hottgenroth
2017-07-27 23:38:17 +02:00
parent 8685bdde87
commit 7e236974c2
4 changed files with 29 additions and 17 deletions

12
dist/main.js vendored
View File

@ -1,12 +1,16 @@
"use strict";
Object.defineProperty(exports, "__esModule", { value: true });
const log = require("./log");
const MqttClient = require("./mqttclient");
const mqtt = require("./mqttclient");
class Dispatcher {
constructor() {
this._mqttClient = new MqttClient.MqttClient();
this._mqttClient.register('IoT/test', null);
this._mqttClient.register('IoT/Device/#', null);
this._mqttClient = new mqtt.MqttClient();
this._mqttClient.register('IoT/test', (message) => {
log.info("Callback for IoT/test");
log.info(`message is ${message}`);
return `<<${message}>>`;
});
this._mqttClient.register('IoT/Device/#', mqtt.passThrough);
}
exec() {
log.info("Dispatcher starting");

8
dist/mqttclient.js vendored
View File

@ -3,6 +3,10 @@ Object.defineProperty(exports, "__esModule", { value: true });
const Mqtt = require("mqtt");
const log = require("./log");
const MQTT_BROKER_DEFAULT_URL = "mqtt://localhost";
function passThrough(message) {
return message;
}
exports.passThrough = passThrough;
class MqttClient {
constructor(mqttBrokerUrl) {
this._mqttBrokerUrl = (mqttBrokerUrl) ? mqttBrokerUrl : MQTT_BROKER_DEFAULT_URL;
@ -26,9 +30,7 @@ class MqttClient {
for (let topicHandler of this._topicHandlers) {
if (this.topicMatch(topicHandler.topic, topic)) {
log.info(`received topic ${topic} matches registered topic ${topicHandler.topic}`);
if (topicHandler.callback != null) {
topicHandler.callback(payload);
}
topicHandler.callback(payload);
}
}
});

View File

@ -1,14 +1,18 @@
import * as log from './log'
import * as MqttClient from './mqttclient'
import * as mqtt from './mqttclient'
class Dispatcher {
private _mqttClient: MqttClient.MqttClient
private _mqttClient: mqtt.MqttClient
constructor() {
this._mqttClient = new MqttClient.MqttClient()
this._mqttClient = new mqtt.MqttClient()
this._mqttClient.register('IoT/test', null)
this._mqttClient.register('IoT/Device/#', null)
this._mqttClient.register('IoT/test', (message: any) : any => {
log.info("Callback for IoT/test")
log.info(`message is ${message}`)
return `<<${message}>>`
})
this._mqttClient.register('IoT/Device/#', mqtt.passThrough)
}
exec() : void {

View File

@ -7,7 +7,11 @@ type TopicHandlerCallback = (message: any) => any
interface TopicHandler {
topic: string,
callback: TopicHandlerCallback|null
callback: TopicHandlerCallback
}
export function passThrough(message: any) {
return message
}
export class MqttClient {
@ -20,7 +24,7 @@ export class MqttClient {
this._topicHandlers = []
}
register(topic: string, callback: TopicHandlerCallback|null) : void {
register(topic: string, callback: TopicHandlerCallback) : void {
this._topicHandlers.push({topic, callback})
log.info(`handler registered for topic ${topic}`)
}
@ -39,9 +43,7 @@ export class MqttClient {
for (let topicHandler of this._topicHandlers) {
if (this.topicMatch(topicHandler.topic, topic)) {
log.info(`received topic ${topic} matches registered topic ${topicHandler.topic}`)
if (topicHandler.callback != null) {
topicHandler.callback(payload)
}
topicHandler.callback(payload)
}
}
})