src/AMQPCommand.cpp
branchv_0
changeset 0 08cb319d7c3a
--- /dev/null	Thu Jan 01 00:00:00 1970 +0000
+++ b/src/AMQPCommand.cpp	Sun May 01 18:29:58 2022 +0200
@@ -0,0 +1,74 @@
+/**
+ * Relational pipes
+ * Copyright © 2022 František Kučera (Frantovo.cz, GlobalCode.info)
+ *
+ * This program is free software: you can redistribute it and/or modify
+ * it under the terms of the GNU General Public License as published by
+ * the Free Software Foundation, version 3 of the License.
+ *
+ * This program is distributed in the hope that it will be useful,
+ * but WITHOUT ANY WARRANTY; without even the implied warranty of
+ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
+ * GNU General Public License for more details.
+ *
+ * You should have received a copy of the GNU General Public License
+ * along with this program. If not, see <http://www.gnu.org/licenses/>.
+ */
+#include <cstdlib>
+#include <vector>
+#include <memory>
+#include <locale>
+#include <regex>
+#include <algorithm>
+#include <unistd.h>
+#include <sstream>
+#include <iomanip>
+
+#include <relpipe/writer/RelationalWriter.h>
+#include <relpipe/writer/RelpipeWriterException.h>
+#include <relpipe/writer/AttributeMetadata.h>
+#include <relpipe/writer/Factory.h>
+#include <relpipe/writer/TypeId.h>
+
+#include <relpipe/cli/CLI.h>
+
+#include "AMQPCommand.h"
+#include "AMQP.h"
+#include "Hex.h"
+
+using namespace std;
+using namespace relpipe::cli;
+using namespace relpipe::writer;
+
+namespace relpipe {
+namespace in {
+namespace amqp {
+
+void AMQPCommand::process(std::shared_ptr<writer::RelationalWriter> writer, Configuration& configuration) {
+	vector<AttributeMetadata> metadata;
+
+	std::shared_ptr<AMQP> mq(AMQP::open(convertor.to_bytes(configuration.queue), configuration.unlinkOnClose));
+
+	writer->startRelation(configuration.relation,{
+		{L"queue", TypeId::STRING},
+		{L"text", TypeId::STRING},
+		{L"data", TypeId::STRING}
+	}, true);
+
+	for (int i = configuration.messageCount; continueProcessing && i > 0; i--) {
+		// TODO: maybe rather call mq_timedreceive() inside and check continueProcessing (to be able to stop even when no messages are comming)
+		std::string message = mq->receive();
+
+		writer->writeAttribute(configuration.queue);
+		writer->writeAttribute(Hex::toTxt(message));
+		writer->writeAttribute(Hex::toHex(message));
+	}
+
+}
+
+AMQPCommand::~AMQPCommand() {
+}
+
+}
+}
+}