mirror of
https://github.com/n8n-io/n8n.git
synced 2025-01-14 22:38:06 -08:00
88dea330b9
* ⚡ Update `lintfix` script * ⚡ Run baseline `lintfix` * 🔥 Remove unneeded exceptions (#3538) * 🔥 Remove exceptions for `node-param-default-wrong-for-simplify` * 🔥 Remove exceptions for `node-param-placeholder-miscased-id` * ⚡ Update version * 👕 Apply `node-param-placeholder-missing` (#3542) * 👕 Apply `filesystem-wrong-cred-filename` (#3543) * 👕 Apply `node-param-description-missing-from-dynamic-options` (#3545) Co-authored-by: Iván Ovejero <ivov.src@gmail.com> * 👕 Apply `node-class-description-empty-string` (#3546) * 👕 Apply `node-class-description-icon-not-svg` (#3548) * 👕 Apply `filesystem-wrong-node-filename` (#3549) Co-authored-by: Iván Ovejero <ivov.src@gmail.com> * 👕 Expand lintings to credentials (#3550) * 👕 Apply `node-param-multi-options-type-unsorted-items` (#3552) * ⚡ fix * ⚡ Minor fixes Co-authored-by: Michael Kret <michael.k@radency.com> * 👕 Apply `node-param-description-wrong-for-dynamic-multi-options` (#3541) * ⚡ Add new lint rule, node-param-description-wrong-for-dynamic-multi-options * ⚡ Fix with updated linting rules * ⚡ Minor fixes Co-authored-by: Iván Ovejero <ivov.src@gmail.com> * 👕 Apply `node-param-description-boolean-without-whether` (#3553) * ⚡ fix * Update packages/nodes-base/nodes/Clockify/ProjectDescription.ts Co-authored-by: Iván Ovejero <ivov.src@gmail.com> * 👕 Apply node-param-display-name-wrong-for-dynamic-multi-options (#3537) * 👕 Add exceptions * 👕 Add exception * ✏️ Alphabetize rules * ⚡ Restore `lintfix` command Co-authored-by: agobrech <45268029+agobrech@users.noreply.github.com> Co-authored-by: Omar Ajoue <krynble@gmail.com> Co-authored-by: Michael Kret <michael.k@radency.com> Co-authored-by: brianinoa <54530642+brianinoa@users.noreply.github.com> Co-authored-by: Michael Kret <88898367+michael-radency@users.noreply.github.com>
197 lines
4.7 KiB
TypeScript
197 lines
4.7 KiB
TypeScript
import {
|
|
IExecuteFunctions,
|
|
} from 'n8n-core';
|
|
|
|
import {
|
|
IDataObject,
|
|
INodeExecutionData,
|
|
INodeType,
|
|
INodeTypeDescription,
|
|
} from 'n8n-workflow';
|
|
|
|
import mqtt from 'mqtt';
|
|
|
|
import {
|
|
IClientOptions,
|
|
} from 'mqtt';
|
|
|
|
export class Mqtt implements INodeType {
|
|
description: INodeTypeDescription = {
|
|
displayName: 'MQTT',
|
|
name: 'mqtt',
|
|
icon: 'file:mqtt.svg',
|
|
group: ['input'],
|
|
version: 1,
|
|
description: 'Push messages to MQTT',
|
|
defaults: {
|
|
name: 'MQTT',
|
|
},
|
|
inputs: ['main'],
|
|
outputs: ['main'],
|
|
credentials: [
|
|
{
|
|
name: 'mqtt',
|
|
required: true,
|
|
},
|
|
],
|
|
properties: [
|
|
{
|
|
displayName: 'Topic',
|
|
name: 'topic',
|
|
type: 'string',
|
|
required: true,
|
|
default: '',
|
|
description: 'The topic to publish to',
|
|
},
|
|
{
|
|
displayName: 'Send Input Data',
|
|
name: 'sendInputData',
|
|
type: 'boolean',
|
|
default: true,
|
|
description: 'Whether to send the the data the node receives as JSON',
|
|
},
|
|
{
|
|
displayName: 'Message',
|
|
name: 'message',
|
|
type: 'string',
|
|
required: true,
|
|
displayOptions: {
|
|
show: {
|
|
sendInputData: [
|
|
false,
|
|
],
|
|
},
|
|
},
|
|
default: '',
|
|
description: 'The message to publish',
|
|
},
|
|
{
|
|
displayName: 'Options',
|
|
name: 'options',
|
|
type: 'collection',
|
|
placeholder: 'Add Option',
|
|
default: {},
|
|
options: [
|
|
{
|
|
displayName: 'QoS',
|
|
name: 'qos',
|
|
type: 'options',
|
|
options: [
|
|
{
|
|
name: 'Received at Most Once',
|
|
value: 0,
|
|
},
|
|
{
|
|
name: 'Received at Least Once',
|
|
value: 1,
|
|
},
|
|
{
|
|
name: 'Exactly Once',
|
|
value: 2,
|
|
},
|
|
],
|
|
default: 0,
|
|
description: 'QoS subscription level',
|
|
},
|
|
{
|
|
displayName: 'Retain',
|
|
name: 'retain',
|
|
type: 'boolean',
|
|
default: false,
|
|
// eslint-disable-next-line n8n-nodes-base/node-param-description-boolean-without-whether
|
|
description: 'Normally if a publisher publishes a message to a topic, and no one is subscribed to that topic the message is simply discarded by the broker. However the publisher can tell the broker to keep the last message on that topic by setting the retain flag to true.',
|
|
},
|
|
],
|
|
},
|
|
],
|
|
};
|
|
|
|
async execute(this: IExecuteFunctions): Promise<INodeExecutionData[][]> {
|
|
const items = this.getInputData();
|
|
const length = items.length;
|
|
const credentials = await this.getCredentials('mqtt');
|
|
|
|
const protocol = credentials.protocol as string || 'mqtt';
|
|
const host = credentials.host as string;
|
|
const brokerUrl = `${protocol}://${host}`;
|
|
const port = credentials.port as number || 1883;
|
|
const clientId = credentials.clientId as string || `mqttjs_${Math.random().toString(16).substr(2, 8)}`;
|
|
const clean = credentials.clean as boolean;
|
|
const ssl = credentials.ssl as boolean;
|
|
const ca = credentials.ca as string;
|
|
const cert = credentials.cert as string;
|
|
const key = credentials.key as string;
|
|
const rejectUnauthorized = credentials.rejectUnauthorized as boolean;
|
|
|
|
let client: mqtt.MqttClient;
|
|
|
|
if (ssl === false) {
|
|
const clientOptions: IClientOptions = {
|
|
port,
|
|
clean,
|
|
clientId,
|
|
};
|
|
|
|
if (credentials.username && credentials.password) {
|
|
clientOptions.username = credentials.username as string;
|
|
clientOptions.password = credentials.password as string;
|
|
}
|
|
|
|
client = mqtt.connect(brokerUrl, clientOptions);
|
|
}
|
|
else {
|
|
const clientOptions: IClientOptions = {
|
|
port,
|
|
clean,
|
|
clientId,
|
|
ca,
|
|
cert,
|
|
key,
|
|
rejectUnauthorized,
|
|
};
|
|
if (credentials.username && credentials.password) {
|
|
clientOptions.username = credentials.username as string;
|
|
clientOptions.password = credentials.password as string;
|
|
}
|
|
|
|
client = mqtt.connect(brokerUrl, clientOptions);
|
|
}
|
|
|
|
const sendInputData = this.getNodeParameter('sendInputData', 0) as boolean;
|
|
|
|
// tslint:disable-next-line: no-any
|
|
const data = await new Promise((resolve, reject): any => {
|
|
client.on('connect', () => {
|
|
for (let i = 0; i < length; i++) {
|
|
|
|
let message;
|
|
const topic = (this.getNodeParameter('topic', i) as string);
|
|
const options = (this.getNodeParameter('options', i) as IDataObject);
|
|
|
|
try {
|
|
if (sendInputData === true) {
|
|
message = JSON.stringify(items[i].json);
|
|
} else {
|
|
message = this.getNodeParameter('message', i) as string;
|
|
}
|
|
client.publish(topic, message, options);
|
|
} catch (e) {
|
|
reject(e);
|
|
}
|
|
}
|
|
//wait for the in-flight messages to be acked.
|
|
//needed for messages with QoS 1 & 2
|
|
client.end(false, {}, () => {
|
|
resolve([items]);
|
|
});
|
|
|
|
client.on('error', (e: string | undefined) => {
|
|
reject(e);
|
|
});
|
|
});
|
|
});
|
|
|
|
return data as INodeExecutionData[][];
|
|
}
|
|
}
|