diff --git a/src/samples/stdinpub.c b/src/samples/stdinpub.c index 98ea59ee..9330fe02 100644 --- a/src/samples/stdinpub.c +++ b/src/samples/stdinpub.c @@ -63,7 +63,7 @@ void usage() printf(" --port (default is 1883)\n"); printf(" --qos (default is 0)\n"); printf(" --retained (default is off)\n"); - printf(" --delimiter (default is \n)"); + printf(" --delimiter (default is \\n)"); printf(" --clientid (default is hostname+timestamp)"); printf(" --maxdatalen 100\n"); printf(" --username none\n"); @@ -93,7 +93,7 @@ void cfinish(int sig) struct { char* clientid; - char delimiter; + char* delimiter; int maxdatalen; int qos; int retained; @@ -104,7 +104,7 @@ struct int verbose; } opts = { - "publisher", '\n', 100, 0, 0, NULL, NULL, "localhost", "1883", 0 + "publisher", "\n", 100, 0, 0, NULL, NULL, "localhost", "1883", 0 }; void getopts(int argc, char** argv); @@ -156,11 +156,19 @@ int main(int argc, char** argv) while (!toStop) { int data_len = 0; + int delim_len = 0; + delim_len = strlen(opts.delimiter); do { buffer[data_len++] = getchar(); - } while (buffer[data_len-1] != opts.delimiter && data_len < opts.maxdatalen); + if (data_len > delim_len) + { + //printf("comparing %s %s\n", opts.delimiter, &buffer[data_len - delim_len]); + if (strncmp(opts.delimiter, &buffer[data_len - delim_len], delim_len) == 0) + break; + } + } while (data_len < opts.maxdatalen); if (opts.verbose) printf("Publishing data of length %d\n", data_len); @@ -256,7 +264,7 @@ void getopts(int argc, char** argv) else if (strcmp(argv[count], "--delimiter") == 0) { if (++count < argc) - opts.delimiter = argv[count][0]; + opts.delimiter = argv[count]; else usage(); } diff --git a/src/samples/stdoutsub.c b/src/samples/stdoutsub.c index 1a5f0e99..9330fe02 100644 --- a/src/samples/stdoutsub.c +++ b/src/samples/stdoutsub.c @@ -13,27 +13,28 @@ * Contributors: * Ian Craggs - initial contribution *******************************************************************************/ - -/* - stdout subscriber + /* + stdin publisher compulsory parameters: - --topic topic to subscribe to + --topic topic to publish on defaulted parameters: --host localhost --port 1883 - --qos 2 - --delimiter \n - --clientid stdout_subscriber + --qos 0 + --delimiters \n + --clientid stdin_publisher + --maxdatalen 100 --userid none --password none */ + #include "MQTTClient.h" #include "MQTTClientPersistence.h" @@ -56,26 +57,27 @@ volatile int toStop = 0; void usage() { - printf("MQTT stdout subscriber\n"); - printf("Usage: stdoutsub topicname , where options are:\n"); + printf("MQTT stdin publisher\n"); + printf("Usage: stdinpub topicname , where options are:\n"); printf(" --host (default is localhost)\n"); printf(" --port (default is 1883)\n"); - printf(" --qos (default is 2)\n"); - printf(" --delimiter (default is no delimiter)\n"); - printf(" --clientid (default is hostname+timestamp)\n"); + printf(" --qos (default is 0)\n"); + printf(" --retained (default is off)\n"); + printf(" --delimiter (default is \\n)"); + printf(" --clientid (default is hostname+timestamp)"); + printf(" --maxdatalen 100\n"); printf(" --username none\n"); printf(" --password none\n"); - printf(" --showtopics (default is on if the topic has a wildcard, else off)\n"); exit(-1); } void myconnect(MQTTClient* client, MQTTClient_connectOptions* opts) { - int rc = 0; - if ((rc = MQTTClient_connect(*client, opts)) != 0) + printf("Connecting\n"); + if (MQTTClient_connect(*client, opts) != 0) { - printf("Failed to connect, return code %d\n", rc); + printf("Failed to connect\n"); exit(-1); } } @@ -88,50 +90,59 @@ void cfinish(int sig) } -struct opts_struct +struct { char* clientid; - int nodelimiter; - char delimiter; + char* delimiter; + int maxdatalen; int qos; + int retained; char* username; char* password; char* host; char* port; - int showtopics; + int verbose; } opts = { - "stdout-subscriber", 1, '\n', 2, NULL, NULL, "localhost", "1883", 0 + "publisher", "\n", 100, 0, 0, NULL, NULL, "localhost", "1883", 0 }; void getopts(int argc, char** argv); +int messageArrived(void* context, char* topicName, int topicLen, MQTTClient_message* m) +{ + /* not expecting any messages */ + return 1; +} + int main(int argc, char** argv) { MQTTClient client; MQTTClient_connectOptions conn_opts = MQTTClient_connectOptions_initializer; char* topic = NULL; + char* buffer = NULL; int rc = 0; char url[100]; - + if (argc < 2) usage(); - topic = argv[1]; - - if (strchr(topic, '#') || strchr(topic, '+')) - opts.showtopics = 1; - if (opts.showtopics) - printf("topic is %s\n", topic); - - getopts(argc, argv); + getopts(argc, argv); + sprintf(url, "%s:%s", opts.host, opts.port); + if (opts.verbose) + printf("URL is %s\n", url); + + topic = argv[1]; + printf("Using topic %s\n", topic); rc = MQTTClient_create(&client, url, opts.clientid, MQTTCLIENT_PERSISTENCE_NONE, NULL); signal(SIGINT, cfinish); signal(SIGTERM, cfinish); + rc = MQTTClient_setCallbacks(client, NULL, NULL, messageArrived, NULL); + conn_opts.keepAliveInterval = 10; conn_opts.reliable = 0; conn_opts.cleansession = 1; @@ -139,33 +150,41 @@ int main(int argc, char** argv) conn_opts.password = opts.password; myconnect(&client, &conn_opts); - - rc = MQTTClient_subscribe(client, topic, opts.qos); + buffer = malloc(opts.maxdatalen); + while (!toStop) { - char* topicName = NULL; - int topicLen; - MQTTClient_message* message = NULL; + int data_len = 0; + int delim_len = 0; - rc = MQTTClient_receive(client, &topicName, &topicLen, &message, 1000); - if (message) + delim_len = strlen(opts.delimiter); + do { - if (opts.showtopics) - printf("%s\t", topicName); - if (opts.nodelimiter) - printf("%.*s", message->payloadlen, (char*)message->payload); - else - printf("%.*s%c", message->payloadlen, (char*)message->payload, opts.delimiter); - fflush(stdout); - MQTTClient_freeMessage(&message); - MQTTClient_free(topicName); - } + buffer[data_len++] = getchar(); + if (data_len > delim_len) + { + //printf("comparing %s %s\n", opts.delimiter, &buffer[data_len - delim_len]); + if (strncmp(opts.delimiter, &buffer[data_len - delim_len], delim_len) == 0) + break; + } + } while (data_len < opts.maxdatalen); + + if (opts.verbose) + printf("Publishing data of length %d\n", data_len); + rc = MQTTClient_publish(client, topic, data_len, buffer, opts.qos, opts.retained, NULL); if (rc != 0) + { myconnect(&client, &conn_opts); + rc = MQTTClient_publish(client, topic, data_len, buffer, opts.qos, opts.retained, NULL); + } + if (opts.qos > 0) + MQTTClient_yield(); } printf("Stopping\n"); + + free(buffer); MQTTClient_disconnect(client, 0); @@ -180,7 +199,11 @@ void getopts(int argc, char** argv) while (count < argc) { - if (strcmp(argv[count], "--qos") == 0) + if (strcmp(argv[count], "--retained") == 0) + opts.retained = 1; + if (strcmp(argv[count], "--verbose") == 0) + opts.verbose = 1; + else if (strcmp(argv[count], "--qos") == 0) { if (++count < argc) { @@ -231,30 +254,17 @@ void getopts(int argc, char** argv) else usage(); } - else if (strcmp(argv[count], "--delimiter") == 0) + else if (strcmp(argv[count], "--maxdatalen") == 0) { if (++count < argc) - { - if (strcmp("newline", argv[count]) == 0) - opts.delimiter = '\n'; - else - opts.delimiter = argv[count][0]; - opts.nodelimiter = 0; - } + opts.maxdatalen = atoi(argv[count]); else usage(); } - else if (strcmp(argv[count], "--showtopics") == 0) + else if (strcmp(argv[count], "--delimiter") == 0) { if (++count < argc) - { - if (strcmp(argv[count], "on") == 0) - opts.showtopics = 1; - else if (strcmp(argv[count], "off") == 0) - opts.showtopics = 0; - else - usage(); - } + opts.delimiter = argv[count]; else usage(); } @@ -262,3 +272,4 @@ void getopts(int argc, char** argv) } } +