summaryrefslogtreecommitdiffstats
path: root/plugins/omazureeventhubs
diff options
context:
space:
mode:
authorDaniel Baumann <daniel.baumann@progress-linux.org>2024-04-15 16:28:20 +0000
committerDaniel Baumann <daniel.baumann@progress-linux.org>2024-04-15 16:28:20 +0000
commitdcc721a95bef6f0d8e6d8775b8efe33e5aecd562 (patch)
tree66a2774cd0ee294d019efd71d2544c70f42b2842 /plugins/omazureeventhubs
parentInitial commit. (diff)
downloadrsyslog-dcc721a95bef6f0d8e6d8775b8efe33e5aecd562.tar.xz
rsyslog-dcc721a95bef6f0d8e6d8775b8efe33e5aecd562.zip
Adding upstream version 8.2402.0.upstream/8.2402.0
Signed-off-by: Daniel Baumann <daniel.baumann@progress-linux.org>
Diffstat (limited to '')
-rw-r--r--plugins/omazureeventhubs/Makefile.am13
-rw-r--r--plugins/omazureeventhubs/Makefile.in802
-rw-r--r--plugins/omazureeventhubs/omazureeventhubs.c1354
3 files changed, 2169 insertions, 0 deletions
diff --git a/plugins/omazureeventhubs/Makefile.am b/plugins/omazureeventhubs/Makefile.am
new file mode 100644
index 0000000..356bb88
--- /dev/null
+++ b/plugins/omazureeventhubs/Makefile.am
@@ -0,0 +1,13 @@
+pkglib_LTLIBRARIES = omazureeventhubs.la
+
+omazureeventhubs_la_SOURCES = omazureeventhubs.c
+if ENABLE_QPIDPROTON_STATIC
+omazureeventhubs_la_LDFLAGS = -module -avoid-version $(PROTON_PROACTOR_LIBS) $(PTHREADS_LIBS) $(OPENSSL_LIBS) -lm
+omazureeventhubs_la_LDFLAGS = -module -avoid-version -Wl,-whole-archive -l:libqpid-proton-proactor-static.a -l:libqpid-proton-core-static.a -Wl,--no-whole-archive $(PTHREADS_LIBS) $(OPENSSL_LIBS) ${RT_LIBS} -lsasl2
+omazureeventhubs_la_LIBADD =
+else
+omazureeventhubs_la_LDFLAGS = -module -avoid-version $(PROTON_PROACTOR_LIBS) $(PTHREADS_LIBS) $(OPENSSL_LIBS) -lm
+omazureeventhubs_la_LIBADD =
+endif
+omazureeventhubs_la_CPPFLAGS = $(RSRT_CFLAGS) $(PTHREADS_CFLAGS) $(CURL_CFLAGS) $(PROTON_PROACTOR_CFLAGS) -Wno-error=switch
+EXTRA_DIST =
diff --git a/plugins/omazureeventhubs/Makefile.in b/plugins/omazureeventhubs/Makefile.in
new file mode 100644
index 0000000..0fd0502
--- /dev/null
+++ b/plugins/omazureeventhubs/Makefile.in
@@ -0,0 +1,802 @@
+# Makefile.in generated by automake 1.16.1 from Makefile.am.
+# @configure_input@
+
+# Copyright (C) 1994-2018 Free Software Foundation, Inc.
+
+# This Makefile.in is free software; the Free Software Foundation
+# gives unlimited permission to copy and/or distribute it,
+# with or without modifications, as long as this notice is preserved.
+
+# This program is distributed in the hope that it will be useful,
+# but WITHOUT ANY WARRANTY, to the extent permitted by law; without
+# even the implied warranty of MERCHANTABILITY or FITNESS FOR A
+# PARTICULAR PURPOSE.
+
+@SET_MAKE@
+
+VPATH = @srcdir@
+am__is_gnu_make = { \
+ if test -z '$(MAKELEVEL)'; then \
+ false; \
+ elif test -n '$(MAKE_HOST)'; then \
+ true; \
+ elif test -n '$(MAKE_VERSION)' && test -n '$(CURDIR)'; then \
+ true; \
+ else \
+ false; \
+ fi; \
+}
+am__make_running_with_option = \
+ case $${target_option-} in \
+ ?) ;; \
+ *) echo "am__make_running_with_option: internal error: invalid" \
+ "target option '$${target_option-}' specified" >&2; \
+ exit 1;; \
+ esac; \
+ has_opt=no; \
+ sane_makeflags=$$MAKEFLAGS; \
+ if $(am__is_gnu_make); then \
+ sane_makeflags=$$MFLAGS; \
+ else \
+ case $$MAKEFLAGS in \
+ *\\[\ \ ]*) \
+ bs=\\; \
+ sane_makeflags=`printf '%s\n' "$$MAKEFLAGS" \
+ | sed "s/$$bs$$bs[$$bs $$bs ]*//g"`;; \
+ esac; \
+ fi; \
+ skip_next=no; \
+ strip_trailopt () \
+ { \
+ flg=`printf '%s\n' "$$flg" | sed "s/$$1.*$$//"`; \
+ }; \
+ for flg in $$sane_makeflags; do \
+ test $$skip_next = yes && { skip_next=no; continue; }; \
+ case $$flg in \
+ *=*|--*) continue;; \
+ -*I) strip_trailopt 'I'; skip_next=yes;; \
+ -*I?*) strip_trailopt 'I';; \
+ -*O) strip_trailopt 'O'; skip_next=yes;; \
+ -*O?*) strip_trailopt 'O';; \
+ -*l) strip_trailopt 'l'; skip_next=yes;; \
+ -*l?*) strip_trailopt 'l';; \
+ -[dEDm]) skip_next=yes;; \
+ -[JT]) skip_next=yes;; \
+ esac; \
+ case $$flg in \
+ *$$target_option*) has_opt=yes; break;; \
+ esac; \
+ done; \
+ test $$has_opt = yes
+am__make_dryrun = (target_option=n; $(am__make_running_with_option))
+am__make_keepgoing = (target_option=k; $(am__make_running_with_option))
+pkgdatadir = $(datadir)/@PACKAGE@
+pkgincludedir = $(includedir)/@PACKAGE@
+pkglibdir = $(libdir)/@PACKAGE@
+pkglibexecdir = $(libexecdir)/@PACKAGE@
+am__cd = CDPATH="$${ZSH_VERSION+.}$(PATH_SEPARATOR)" && cd
+install_sh_DATA = $(install_sh) -c -m 644
+install_sh_PROGRAM = $(install_sh) -c
+install_sh_SCRIPT = $(install_sh) -c
+INSTALL_HEADER = $(INSTALL_DATA)
+transform = $(program_transform_name)
+NORMAL_INSTALL = :
+PRE_INSTALL = :
+POST_INSTALL = :
+NORMAL_UNINSTALL = :
+PRE_UNINSTALL = :
+POST_UNINSTALL = :
+build_triplet = @build@
+host_triplet = @host@
+subdir = plugins/omazureeventhubs
+ACLOCAL_M4 = $(top_srcdir)/aclocal.m4
+am__aclocal_m4_deps = $(top_srcdir)/m4/ac_check_define.m4 \
+ $(top_srcdir)/m4/atomic_operations.m4 \
+ $(top_srcdir)/m4/atomic_operations_64bit.m4 \
+ $(top_srcdir)/m4/libtool.m4 $(top_srcdir)/m4/ltoptions.m4 \
+ $(top_srcdir)/m4/ltsugar.m4 $(top_srcdir)/m4/ltversion.m4 \
+ $(top_srcdir)/m4/lt~obsolete.m4 $(top_srcdir)/configure.ac
+am__configure_deps = $(am__aclocal_m4_deps) $(CONFIGURE_DEPENDENCIES) \
+ $(ACLOCAL_M4)
+DIST_COMMON = $(srcdir)/Makefile.am $(am__DIST_COMMON)
+mkinstalldirs = $(install_sh) -d
+CONFIG_HEADER = $(top_builddir)/config.h
+CONFIG_CLEAN_FILES =
+CONFIG_CLEAN_VPATH_FILES =
+am__vpath_adj_setup = srcdirstrip=`echo "$(srcdir)" | sed 's|.|.|g'`;
+am__vpath_adj = case $$p in \
+ $(srcdir)/*) f=`echo "$$p" | sed "s|^$$srcdirstrip/||"`;; \
+ *) f=$$p;; \
+ esac;
+am__strip_dir = f=`echo $$p | sed -e 's|^.*/||'`;
+am__install_max = 40
+am__nobase_strip_setup = \
+ srcdirstrip=`echo "$(srcdir)" | sed 's/[].[^$$\\*|]/\\\\&/g'`
+am__nobase_strip = \
+ for p in $$list; do echo "$$p"; done | sed -e "s|$$srcdirstrip/||"
+am__nobase_list = $(am__nobase_strip_setup); \
+ for p in $$list; do echo "$$p $$p"; done | \
+ sed "s| $$srcdirstrip/| |;"' / .*\//!s/ .*/ ./; s,\( .*\)/[^/]*$$,\1,' | \
+ $(AWK) 'BEGIN { files["."] = "" } { files[$$2] = files[$$2] " " $$1; \
+ if (++n[$$2] == $(am__install_max)) \
+ { print $$2, files[$$2]; n[$$2] = 0; files[$$2] = "" } } \
+ END { for (dir in files) print dir, files[dir] }'
+am__base_list = \
+ sed '$$!N;$$!N;$$!N;$$!N;$$!N;$$!N;$$!N;s/\n/ /g' | \
+ sed '$$!N;$$!N;$$!N;$$!N;s/\n/ /g'
+am__uninstall_files_from_dir = { \
+ test -z "$$files" \
+ || { test ! -d "$$dir" && test ! -f "$$dir" && test ! -r "$$dir"; } \
+ || { echo " ( cd '$$dir' && rm -f" $$files ")"; \
+ $(am__cd) "$$dir" && rm -f $$files; }; \
+ }
+am__installdirs = "$(DESTDIR)$(pkglibdir)"
+LTLIBRARIES = $(pkglib_LTLIBRARIES)
+omazureeventhubs_la_DEPENDENCIES =
+am_omazureeventhubs_la_OBJECTS = \
+ omazureeventhubs_la-omazureeventhubs.lo
+omazureeventhubs_la_OBJECTS = $(am_omazureeventhubs_la_OBJECTS)
+AM_V_lt = $(am__v_lt_@AM_V@)
+am__v_lt_ = $(am__v_lt_@AM_DEFAULT_V@)
+am__v_lt_0 = --silent
+am__v_lt_1 =
+omazureeventhubs_la_LINK = $(LIBTOOL) $(AM_V_lt) --tag=CC \
+ $(AM_LIBTOOLFLAGS) $(LIBTOOLFLAGS) --mode=link $(CCLD) \
+ $(AM_CFLAGS) $(CFLAGS) $(omazureeventhubs_la_LDFLAGS) \
+ $(LDFLAGS) -o $@
+AM_V_P = $(am__v_P_@AM_V@)
+am__v_P_ = $(am__v_P_@AM_DEFAULT_V@)
+am__v_P_0 = false
+am__v_P_1 = :
+AM_V_GEN = $(am__v_GEN_@AM_V@)
+am__v_GEN_ = $(am__v_GEN_@AM_DEFAULT_V@)
+am__v_GEN_0 = @echo " GEN " $@;
+am__v_GEN_1 =
+AM_V_at = $(am__v_at_@AM_V@)
+am__v_at_ = $(am__v_at_@AM_DEFAULT_V@)
+am__v_at_0 = @
+am__v_at_1 =
+DEFAULT_INCLUDES = -I.@am__isrc@ -I$(top_builddir)
+depcomp = $(SHELL) $(top_srcdir)/depcomp
+am__maybe_remake_depfiles = depfiles
+am__depfiles_remade = \
+ ./$(DEPDIR)/omazureeventhubs_la-omazureeventhubs.Plo
+am__mv = mv -f
+COMPILE = $(CC) $(DEFS) $(DEFAULT_INCLUDES) $(INCLUDES) $(AM_CPPFLAGS) \
+ $(CPPFLAGS) $(AM_CFLAGS) $(CFLAGS)
+LTCOMPILE = $(LIBTOOL) $(AM_V_lt) --tag=CC $(AM_LIBTOOLFLAGS) \
+ $(LIBTOOLFLAGS) --mode=compile $(CC) $(DEFS) \
+ $(DEFAULT_INCLUDES) $(INCLUDES) $(AM_CPPFLAGS) $(CPPFLAGS) \
+ $(AM_CFLAGS) $(CFLAGS)
+AM_V_CC = $(am__v_CC_@AM_V@)
+am__v_CC_ = $(am__v_CC_@AM_DEFAULT_V@)
+am__v_CC_0 = @echo " CC " $@;
+am__v_CC_1 =
+CCLD = $(CC)
+LINK = $(LIBTOOL) $(AM_V_lt) --tag=CC $(AM_LIBTOOLFLAGS) \
+ $(LIBTOOLFLAGS) --mode=link $(CCLD) $(AM_CFLAGS) $(CFLAGS) \
+ $(AM_LDFLAGS) $(LDFLAGS) -o $@
+AM_V_CCLD = $(am__v_CCLD_@AM_V@)
+am__v_CCLD_ = $(am__v_CCLD_@AM_DEFAULT_V@)
+am__v_CCLD_0 = @echo " CCLD " $@;
+am__v_CCLD_1 =
+SOURCES = $(omazureeventhubs_la_SOURCES)
+DIST_SOURCES = $(omazureeventhubs_la_SOURCES)
+am__can_run_installinfo = \
+ case $$AM_UPDATE_INFO_DIR in \
+ n|no|NO) false;; \
+ *) (install-info --version) >/dev/null 2>&1;; \
+ esac
+am__tagged_files = $(HEADERS) $(SOURCES) $(TAGS_FILES) $(LISP)
+# Read a list of newline-separated strings from the standard input,
+# and print each of them once, without duplicates. Input order is
+# *not* preserved.
+am__uniquify_input = $(AWK) '\
+ BEGIN { nonempty = 0; } \
+ { items[$$0] = 1; nonempty = 1; } \
+ END { if (nonempty) { for (i in items) print i; }; } \
+'
+# Make sure the list of sources is unique. This is necessary because,
+# e.g., the same source file might be shared among _SOURCES variables
+# for different programs/libraries.
+am__define_uniq_tagged_files = \
+ list='$(am__tagged_files)'; \
+ unique=`for i in $$list; do \
+ if test -f "$$i"; then echo $$i; else echo $(srcdir)/$$i; fi; \
+ done | $(am__uniquify_input)`
+ETAGS = etags
+CTAGS = ctags
+am__DIST_COMMON = $(srcdir)/Makefile.in $(top_srcdir)/depcomp
+DISTFILES = $(DIST_COMMON) $(DIST_SOURCES) $(TEXINFOS) $(EXTRA_DIST)
+ACLOCAL = @ACLOCAL@
+AMTAR = @AMTAR@
+AM_DEFAULT_VERBOSITY = @AM_DEFAULT_VERBOSITY@
+APU_CFLAGS = @APU_CFLAGS@
+APU_LIBS = @APU_LIBS@
+AR = @AR@
+AUTOCONF = @AUTOCONF@
+AUTOHEADER = @AUTOHEADER@
+AUTOMAKE = @AUTOMAKE@
+AWK = @AWK@
+CC = @CC@
+CCDEPMODE = @CCDEPMODE@
+CFLAGS = @CFLAGS@
+CIVETWEB_LIBS = @CIVETWEB_LIBS@
+CONF_FILE_PATH = @CONF_FILE_PATH@
+CPP = @CPP@
+CPPFLAGS = @CPPFLAGS@
+CURL_CFLAGS = @CURL_CFLAGS@
+CURL_LIBS = @CURL_LIBS@
+CYGPATH_W = @CYGPATH_W@
+CZMQ_CFLAGS = @CZMQ_CFLAGS@
+CZMQ_LIBS = @CZMQ_LIBS@
+DEFS = @DEFS@
+DEPDIR = @DEPDIR@
+DLLTOOL = @DLLTOOL@
+DL_LIBS = @DL_LIBS@
+DSYMUTIL = @DSYMUTIL@
+DUMPBIN = @DUMPBIN@
+ECHO_C = @ECHO_C@
+ECHO_N = @ECHO_N@
+ECHO_T = @ECHO_T@
+EGREP = @EGREP@
+EXEEXT = @EXEEXT@
+FAUP_LIBS = @FAUP_LIBS@
+FGREP = @FGREP@
+GLIB_CFLAGS = @GLIB_CFLAGS@
+GLIB_LIBS = @GLIB_LIBS@
+GNUTLS_CFLAGS = @GNUTLS_CFLAGS@
+GNUTLS_LIBS = @GNUTLS_LIBS@
+GREP = @GREP@
+GSS_LIBS = @GSS_LIBS@
+GT_KSI_LS12_CFLAGS = @GT_KSI_LS12_CFLAGS@
+GT_KSI_LS12_LIBS = @GT_KSI_LS12_LIBS@
+HASH_XXHASH_LIBS = @HASH_XXHASH_LIBS@
+HIREDIS_CFLAGS = @HIREDIS_CFLAGS@
+HIREDIS_LIBS = @HIREDIS_LIBS@
+IMUDP_LIBS = @IMUDP_LIBS@
+INSTALL = @INSTALL@
+INSTALL_DATA = @INSTALL_DATA@
+INSTALL_PROGRAM = @INSTALL_PROGRAM@
+INSTALL_SCRIPT = @INSTALL_SCRIPT@
+INSTALL_STRIP_PROGRAM = @INSTALL_STRIP_PROGRAM@
+IP = @IP@
+JAVA = @JAVA@
+JAVAC = @JAVAC@
+LD = @LD@
+LDFLAGS = @LDFLAGS@
+LEX = @LEX@
+LEXLIB = @LEXLIB@
+LEX_OUTPUT_ROOT = @LEX_OUTPUT_ROOT@
+LIBCAPNG_CFLAGS = @LIBCAPNG_CFLAGS@
+LIBCAPNG_LIBS = @LIBCAPNG_LIBS@
+LIBCAPNG_PRESENT_CFLAGS = @LIBCAPNG_PRESENT_CFLAGS@
+LIBCAPNG_PRESENT_LIBS = @LIBCAPNG_PRESENT_LIBS@
+LIBDBI_CFLAGS = @LIBDBI_CFLAGS@
+LIBDBI_LIBS = @LIBDBI_LIBS@
+LIBESTR_CFLAGS = @LIBESTR_CFLAGS@
+LIBESTR_LIBS = @LIBESTR_LIBS@
+LIBEVENT_CFLAGS = @LIBEVENT_CFLAGS@
+LIBEVENT_LIBS = @LIBEVENT_LIBS@
+LIBFASTJSON_CFLAGS = @LIBFASTJSON_CFLAGS@
+LIBFASTJSON_LIBS = @LIBFASTJSON_LIBS@
+LIBGCRYPT_CFLAGS = @LIBGCRYPT_CFLAGS@
+LIBGCRYPT_CONFIG = @LIBGCRYPT_CONFIG@
+LIBGCRYPT_LIBS = @LIBGCRYPT_LIBS@
+LIBLOGGING_CFLAGS = @LIBLOGGING_CFLAGS@
+LIBLOGGING_LIBS = @LIBLOGGING_LIBS@
+LIBLOGGING_STDLOG_CFLAGS = @LIBLOGGING_STDLOG_CFLAGS@
+LIBLOGGING_STDLOG_LIBS = @LIBLOGGING_STDLOG_LIBS@
+LIBLOGNORM_CFLAGS = @LIBLOGNORM_CFLAGS@
+LIBLOGNORM_LIBS = @LIBLOGNORM_LIBS@
+LIBLZ4_CFLAGS = @LIBLZ4_CFLAGS@
+LIBLZ4_LIBS = @LIBLZ4_LIBS@
+LIBM = @LIBM@
+LIBMONGOC_CFLAGS = @LIBMONGOC_CFLAGS@
+LIBMONGOC_LIBS = @LIBMONGOC_LIBS@
+LIBOBJS = @LIBOBJS@
+LIBRDKAFKA_CFLAGS = @LIBRDKAFKA_CFLAGS@
+LIBRDKAFKA_LIBS = @LIBRDKAFKA_LIBS@
+LIBS = @LIBS@
+LIBSYSTEMD_CFLAGS = @LIBSYSTEMD_CFLAGS@
+LIBSYSTEMD_JOURNAL_CFLAGS = @LIBSYSTEMD_JOURNAL_CFLAGS@
+LIBSYSTEMD_JOURNAL_LIBS = @LIBSYSTEMD_JOURNAL_LIBS@
+LIBSYSTEMD_LIBS = @LIBSYSTEMD_LIBS@
+LIBTOOL = @LIBTOOL@
+LIBUUID_CFLAGS = @LIBUUID_CFLAGS@
+LIBUUID_LIBS = @LIBUUID_LIBS@
+LIPO = @LIPO@
+LN_S = @LN_S@
+LTLIBOBJS = @LTLIBOBJS@
+LT_SYS_LIBRARY_PATH = @LT_SYS_LIBRARY_PATH@
+MAKEINFO = @MAKEINFO@
+MANIFEST_TOOL = @MANIFEST_TOOL@
+MKDIR_P = @MKDIR_P@
+MYSQL_CFLAGS = @MYSQL_CFLAGS@
+MYSQL_CONFIG = @MYSQL_CONFIG@
+MYSQL_LIBS = @MYSQL_LIBS@
+NM = @NM@
+NMEDIT = @NMEDIT@
+OBJDUMP = @OBJDUMP@
+OBJEXT = @OBJEXT@
+OPENSSL_CFLAGS = @OPENSSL_CFLAGS@
+OPENSSL_LIBS = @OPENSSL_LIBS@
+OTOOL = @OTOOL@
+OTOOL64 = @OTOOL64@
+PACKAGE = @PACKAGE@
+PACKAGE_BUGREPORT = @PACKAGE_BUGREPORT@
+PACKAGE_NAME = @PACKAGE_NAME@
+PACKAGE_STRING = @PACKAGE_STRING@
+PACKAGE_TARNAME = @PACKAGE_TARNAME@
+PACKAGE_URL = @PACKAGE_URL@
+PACKAGE_VERSION = @PACKAGE_VERSION@
+PATH_SEPARATOR = @PATH_SEPARATOR@
+PGSQL_CFLAGS = @PGSQL_CFLAGS@
+PGSQL_LIBS = @PGSQL_LIBS@
+PG_CONFIG = @PG_CONFIG@
+PID_FILE_PATH = @PID_FILE_PATH@
+PKG_CONFIG = @PKG_CONFIG@
+PKG_CONFIG_LIBDIR = @PKG_CONFIG_LIBDIR@
+PKG_CONFIG_PATH = @PKG_CONFIG_PATH@
+PROTON_CFLAGS = @PROTON_CFLAGS@
+PROTON_LIBS = @PROTON_LIBS@
+PROTON_PROACTOR_CFLAGS = @PROTON_PROACTOR_CFLAGS@
+PROTON_PROACTOR_LIBS = @PROTON_PROACTOR_LIBS@
+PTHREADS_CFLAGS = @PTHREADS_CFLAGS@
+PTHREADS_LIBS = @PTHREADS_LIBS@
+PYTHON = @PYTHON@
+PYTHON_EXEC_PREFIX = @PYTHON_EXEC_PREFIX@
+PYTHON_PLATFORM = @PYTHON_PLATFORM@
+PYTHON_PREFIX = @PYTHON_PREFIX@
+PYTHON_VERSION = @PYTHON_VERSION@
+RABBITMQ_CFLAGS = @RABBITMQ_CFLAGS@
+RABBITMQ_LIBS = @RABBITMQ_LIBS@
+RANLIB = @RANLIB@
+READLINK = @READLINK@
+REDIS = @REDIS@
+RELP_CFLAGS = @RELP_CFLAGS@
+RELP_LIBS = @RELP_LIBS@
+RSRT_CFLAGS = @RSRT_CFLAGS@
+RSRT_CFLAGS1 = @RSRT_CFLAGS1@
+RSRT_LIBS = @RSRT_LIBS@
+RSRT_LIBS1 = @RSRT_LIBS1@
+RST2MAN = @RST2MAN@
+RT_LIBS = @RT_LIBS@
+SED = @SED@
+SET_MAKE = @SET_MAKE@
+SHELL = @SHELL@
+SNMP_CFLAGS = @SNMP_CFLAGS@
+SNMP_LIBS = @SNMP_LIBS@
+SOL_LIBS = @SOL_LIBS@
+STRIP = @STRIP@
+TCL_BIN_DIR = @TCL_BIN_DIR@
+TCL_INCLUDE_SPEC = @TCL_INCLUDE_SPEC@
+TCL_LIB_FILE = @TCL_LIB_FILE@
+TCL_LIB_FLAG = @TCL_LIB_FLAG@
+TCL_LIB_SPEC = @TCL_LIB_SPEC@
+TCL_PATCH_LEVEL = @TCL_PATCH_LEVEL@
+TCL_SRC_DIR = @TCL_SRC_DIR@
+TCL_STUB_LIB_FILE = @TCL_STUB_LIB_FILE@
+TCL_STUB_LIB_FLAG = @TCL_STUB_LIB_FLAG@
+TCL_STUB_LIB_SPEC = @TCL_STUB_LIB_SPEC@
+TCL_VERSION = @TCL_VERSION@
+UDPSPOOF_CFLAGS = @UDPSPOOF_CFLAGS@
+UDPSPOOF_LIBS = @UDPSPOOF_LIBS@
+VALGRIND = @VALGRIND@
+VERSION = @VERSION@
+WARN_CFLAGS = @WARN_CFLAGS@
+WARN_LDFLAGS = @WARN_LDFLAGS@
+WARN_SCANNERFLAGS = @WARN_SCANNERFLAGS@
+WGET = @WGET@
+YACC = @YACC@
+YACC_FOUND = @YACC_FOUND@
+YFLAGS = @YFLAGS@
+ZLIB_CFLAGS = @ZLIB_CFLAGS@
+ZLIB_LIBS = @ZLIB_LIBS@
+ZSTD_CFLAGS = @ZSTD_CFLAGS@
+ZSTD_LIBS = @ZSTD_LIBS@
+abs_builddir = @abs_builddir@
+abs_srcdir = @abs_srcdir@
+abs_top_builddir = @abs_top_builddir@
+abs_top_srcdir = @abs_top_srcdir@
+ac_ct_AR = @ac_ct_AR@
+ac_ct_CC = @ac_ct_CC@
+ac_ct_DUMPBIN = @ac_ct_DUMPBIN@
+am__include = @am__include@
+am__leading_dot = @am__leading_dot@
+am__quote = @am__quote@
+am__tar = @am__tar@
+am__untar = @am__untar@
+bindir = @bindir@
+build = @build@
+build_alias = @build_alias@
+build_cpu = @build_cpu@
+build_os = @build_os@
+build_vendor = @build_vendor@
+builddir = @builddir@
+datadir = @datadir@
+datarootdir = @datarootdir@
+docdir = @docdir@
+dvidir = @dvidir@
+exec_prefix = @exec_prefix@
+host = @host@
+host_alias = @host_alias@
+host_cpu = @host_cpu@
+host_os = @host_os@
+host_vendor = @host_vendor@
+htmldir = @htmldir@
+includedir = @includedir@
+infodir = @infodir@
+install_sh = @install_sh@
+libdir = @libdir@
+libexecdir = @libexecdir@
+localedir = @localedir@
+localstatedir = @localstatedir@
+mandir = @mandir@
+mkdir_p = @mkdir_p@
+moddirs = @moddirs@
+oldincludedir = @oldincludedir@
+pdfdir = @pdfdir@
+pkgpyexecdir = @pkgpyexecdir@
+pkgpythondir = @pkgpythondir@
+prefix = @prefix@
+program_transform_name = @program_transform_name@
+psdir = @psdir@
+pyexecdir = @pyexecdir@
+pythondir = @pythondir@
+runstatedir = @runstatedir@
+sbindir = @sbindir@
+sharedstatedir = @sharedstatedir@
+srcdir = @srcdir@
+sysconfdir = @sysconfdir@
+target_alias = @target_alias@
+top_build_prefix = @top_build_prefix@
+top_builddir = @top_builddir@
+top_srcdir = @top_srcdir@
+pkglib_LTLIBRARIES = omazureeventhubs.la
+omazureeventhubs_la_SOURCES = omazureeventhubs.c
+@ENABLE_QPIDPROTON_STATIC_FALSE@omazureeventhubs_la_LDFLAGS = -module -avoid-version $(PROTON_PROACTOR_LIBS) $(PTHREADS_LIBS) $(OPENSSL_LIBS) -lm
+@ENABLE_QPIDPROTON_STATIC_TRUE@omazureeventhubs_la_LDFLAGS = -module -avoid-version -Wl,-whole-archive -l:libqpid-proton-proactor-static.a -l:libqpid-proton-core-static.a -Wl,--no-whole-archive $(PTHREADS_LIBS) $(OPENSSL_LIBS) ${RT_LIBS} -lsasl2
+@ENABLE_QPIDPROTON_STATIC_FALSE@omazureeventhubs_la_LIBADD =
+@ENABLE_QPIDPROTON_STATIC_TRUE@omazureeventhubs_la_LIBADD =
+omazureeventhubs_la_CPPFLAGS = $(RSRT_CFLAGS) $(PTHREADS_CFLAGS) $(CURL_CFLAGS) $(PROTON_PROACTOR_CFLAGS) -Wno-error=switch
+EXTRA_DIST =
+all: all-am
+
+.SUFFIXES:
+.SUFFIXES: .c .lo .o .obj
+$(srcdir)/Makefile.in: $(srcdir)/Makefile.am $(am__configure_deps)
+ @for dep in $?; do \
+ case '$(am__configure_deps)' in \
+ *$$dep*) \
+ ( cd $(top_builddir) && $(MAKE) $(AM_MAKEFLAGS) am--refresh ) \
+ && { if test -f $@; then exit 0; else break; fi; }; \
+ exit 1;; \
+ esac; \
+ done; \
+ echo ' cd $(top_srcdir) && $(AUTOMAKE) --gnu plugins/omazureeventhubs/Makefile'; \
+ $(am__cd) $(top_srcdir) && \
+ $(AUTOMAKE) --gnu plugins/omazureeventhubs/Makefile
+Makefile: $(srcdir)/Makefile.in $(top_builddir)/config.status
+ @case '$?' in \
+ *config.status*) \
+ cd $(top_builddir) && $(MAKE) $(AM_MAKEFLAGS) am--refresh;; \
+ *) \
+ echo ' cd $(top_builddir) && $(SHELL) ./config.status $(subdir)/$@ $(am__maybe_remake_depfiles)'; \
+ cd $(top_builddir) && $(SHELL) ./config.status $(subdir)/$@ $(am__maybe_remake_depfiles);; \
+ esac;
+
+$(top_builddir)/config.status: $(top_srcdir)/configure $(CONFIG_STATUS_DEPENDENCIES)
+ cd $(top_builddir) && $(MAKE) $(AM_MAKEFLAGS) am--refresh
+
+$(top_srcdir)/configure: $(am__configure_deps)
+ cd $(top_builddir) && $(MAKE) $(AM_MAKEFLAGS) am--refresh
+$(ACLOCAL_M4): $(am__aclocal_m4_deps)
+ cd $(top_builddir) && $(MAKE) $(AM_MAKEFLAGS) am--refresh
+$(am__aclocal_m4_deps):
+
+install-pkglibLTLIBRARIES: $(pkglib_LTLIBRARIES)
+ @$(NORMAL_INSTALL)
+ @list='$(pkglib_LTLIBRARIES)'; test -n "$(pkglibdir)" || list=; \
+ list2=; for p in $$list; do \
+ if test -f $$p; then \
+ list2="$$list2 $$p"; \
+ else :; fi; \
+ done; \
+ test -z "$$list2" || { \
+ echo " $(MKDIR_P) '$(DESTDIR)$(pkglibdir)'"; \
+ $(MKDIR_P) "$(DESTDIR)$(pkglibdir)" || exit 1; \
+ echo " $(LIBTOOL) $(AM_LIBTOOLFLAGS) $(LIBTOOLFLAGS) --mode=install $(INSTALL) $(INSTALL_STRIP_FLAG) $$list2 '$(DESTDIR)$(pkglibdir)'"; \
+ $(LIBTOOL) $(AM_LIBTOOLFLAGS) $(LIBTOOLFLAGS) --mode=install $(INSTALL) $(INSTALL_STRIP_FLAG) $$list2 "$(DESTDIR)$(pkglibdir)"; \
+ }
+
+uninstall-pkglibLTLIBRARIES:
+ @$(NORMAL_UNINSTALL)
+ @list='$(pkglib_LTLIBRARIES)'; test -n "$(pkglibdir)" || list=; \
+ for p in $$list; do \
+ $(am__strip_dir) \
+ echo " $(LIBTOOL) $(AM_LIBTOOLFLAGS) $(LIBTOOLFLAGS) --mode=uninstall rm -f '$(DESTDIR)$(pkglibdir)/$$f'"; \
+ $(LIBTOOL) $(AM_LIBTOOLFLAGS) $(LIBTOOLFLAGS) --mode=uninstall rm -f "$(DESTDIR)$(pkglibdir)/$$f"; \
+ done
+
+clean-pkglibLTLIBRARIES:
+ -test -z "$(pkglib_LTLIBRARIES)" || rm -f $(pkglib_LTLIBRARIES)
+ @list='$(pkglib_LTLIBRARIES)'; \
+ locs=`for p in $$list; do echo $$p; done | \
+ sed 's|^[^/]*$$|.|; s|/[^/]*$$||; s|$$|/so_locations|' | \
+ sort -u`; \
+ test -z "$$locs" || { \
+ echo rm -f $${locs}; \
+ rm -f $${locs}; \
+ }
+
+omazureeventhubs.la: $(omazureeventhubs_la_OBJECTS) $(omazureeventhubs_la_DEPENDENCIES) $(EXTRA_omazureeventhubs_la_DEPENDENCIES)
+ $(AM_V_CCLD)$(omazureeventhubs_la_LINK) -rpath $(pkglibdir) $(omazureeventhubs_la_OBJECTS) $(omazureeventhubs_la_LIBADD) $(LIBS)
+
+mostlyclean-compile:
+ -rm -f *.$(OBJEXT)
+
+distclean-compile:
+ -rm -f *.tab.c
+
+@AMDEP_TRUE@@am__include@ @am__quote@./$(DEPDIR)/omazureeventhubs_la-omazureeventhubs.Plo@am__quote@ # am--include-marker
+
+$(am__depfiles_remade):
+ @$(MKDIR_P) $(@D)
+ @echo '# dummy' >$@-t && $(am__mv) $@-t $@
+
+am--depfiles: $(am__depfiles_remade)
+
+.c.o:
+@am__fastdepCC_TRUE@ $(AM_V_CC)depbase=`echo $@ | sed 's|[^/]*$$|$(DEPDIR)/&|;s|\.o$$||'`;\
+@am__fastdepCC_TRUE@ $(COMPILE) -MT $@ -MD -MP -MF $$depbase.Tpo -c -o $@ $< &&\
+@am__fastdepCC_TRUE@ $(am__mv) $$depbase.Tpo $$depbase.Po
+@AMDEP_TRUE@@am__fastdepCC_FALSE@ $(AM_V_CC)source='$<' object='$@' libtool=no @AMDEPBACKSLASH@
+@AMDEP_TRUE@@am__fastdepCC_FALSE@ DEPDIR=$(DEPDIR) $(CCDEPMODE) $(depcomp) @AMDEPBACKSLASH@
+@am__fastdepCC_FALSE@ $(AM_V_CC@am__nodep@)$(COMPILE) -c -o $@ $<
+
+.c.obj:
+@am__fastdepCC_TRUE@ $(AM_V_CC)depbase=`echo $@ | sed 's|[^/]*$$|$(DEPDIR)/&|;s|\.obj$$||'`;\
+@am__fastdepCC_TRUE@ $(COMPILE) -MT $@ -MD -MP -MF $$depbase.Tpo -c -o $@ `$(CYGPATH_W) '$<'` &&\
+@am__fastdepCC_TRUE@ $(am__mv) $$depbase.Tpo $$depbase.Po
+@AMDEP_TRUE@@am__fastdepCC_FALSE@ $(AM_V_CC)source='$<' object='$@' libtool=no @AMDEPBACKSLASH@
+@AMDEP_TRUE@@am__fastdepCC_FALSE@ DEPDIR=$(DEPDIR) $(CCDEPMODE) $(depcomp) @AMDEPBACKSLASH@
+@am__fastdepCC_FALSE@ $(AM_V_CC@am__nodep@)$(COMPILE) -c -o $@ `$(CYGPATH_W) '$<'`
+
+.c.lo:
+@am__fastdepCC_TRUE@ $(AM_V_CC)depbase=`echo $@ | sed 's|[^/]*$$|$(DEPDIR)/&|;s|\.lo$$||'`;\
+@am__fastdepCC_TRUE@ $(LTCOMPILE) -MT $@ -MD -MP -MF $$depbase.Tpo -c -o $@ $< &&\
+@am__fastdepCC_TRUE@ $(am__mv) $$depbase.Tpo $$depbase.Plo
+@AMDEP_TRUE@@am__fastdepCC_FALSE@ $(AM_V_CC)source='$<' object='$@' libtool=yes @AMDEPBACKSLASH@
+@AMDEP_TRUE@@am__fastdepCC_FALSE@ DEPDIR=$(DEPDIR) $(CCDEPMODE) $(depcomp) @AMDEPBACKSLASH@
+@am__fastdepCC_FALSE@ $(AM_V_CC@am__nodep@)$(LTCOMPILE) -c -o $@ $<
+
+omazureeventhubs_la-omazureeventhubs.lo: omazureeventhubs.c
+@am__fastdepCC_TRUE@ $(AM_V_CC)$(LIBTOOL) $(AM_V_lt) --tag=CC $(AM_LIBTOOLFLAGS) $(LIBTOOLFLAGS) --mode=compile $(CC) $(DEFS) $(DEFAULT_INCLUDES) $(INCLUDES) $(omazureeventhubs_la_CPPFLAGS) $(CPPFLAGS) $(AM_CFLAGS) $(CFLAGS) -MT omazureeventhubs_la-omazureeventhubs.lo -MD -MP -MF $(DEPDIR)/omazureeventhubs_la-omazureeventhubs.Tpo -c -o omazureeventhubs_la-omazureeventhubs.lo `test -f 'omazureeventhubs.c' || echo '$(srcdir)/'`omazureeventhubs.c
+@am__fastdepCC_TRUE@ $(AM_V_at)$(am__mv) $(DEPDIR)/omazureeventhubs_la-omazureeventhubs.Tpo $(DEPDIR)/omazureeventhubs_la-omazureeventhubs.Plo
+@AMDEP_TRUE@@am__fastdepCC_FALSE@ $(AM_V_CC)source='omazureeventhubs.c' object='omazureeventhubs_la-omazureeventhubs.lo' libtool=yes @AMDEPBACKSLASH@
+@AMDEP_TRUE@@am__fastdepCC_FALSE@ DEPDIR=$(DEPDIR) $(CCDEPMODE) $(depcomp) @AMDEPBACKSLASH@
+@am__fastdepCC_FALSE@ $(AM_V_CC@am__nodep@)$(LIBTOOL) $(AM_V_lt) --tag=CC $(AM_LIBTOOLFLAGS) $(LIBTOOLFLAGS) --mode=compile $(CC) $(DEFS) $(DEFAULT_INCLUDES) $(INCLUDES) $(omazureeventhubs_la_CPPFLAGS) $(CPPFLAGS) $(AM_CFLAGS) $(CFLAGS) -c -o omazureeventhubs_la-omazureeventhubs.lo `test -f 'omazureeventhubs.c' || echo '$(srcdir)/'`omazureeventhubs.c
+
+mostlyclean-libtool:
+ -rm -f *.lo
+
+clean-libtool:
+ -rm -rf .libs _libs
+
+ID: $(am__tagged_files)
+ $(am__define_uniq_tagged_files); mkid -fID $$unique
+tags: tags-am
+TAGS: tags
+
+tags-am: $(TAGS_DEPENDENCIES) $(am__tagged_files)
+ set x; \
+ here=`pwd`; \
+ $(am__define_uniq_tagged_files); \
+ shift; \
+ if test -z "$(ETAGS_ARGS)$$*$$unique"; then :; else \
+ test -n "$$unique" || unique=$$empty_fix; \
+ if test $$# -gt 0; then \
+ $(ETAGS) $(ETAGSFLAGS) $(AM_ETAGSFLAGS) $(ETAGS_ARGS) \
+ "$$@" $$unique; \
+ else \
+ $(ETAGS) $(ETAGSFLAGS) $(AM_ETAGSFLAGS) $(ETAGS_ARGS) \
+ $$unique; \
+ fi; \
+ fi
+ctags: ctags-am
+
+CTAGS: ctags
+ctags-am: $(TAGS_DEPENDENCIES) $(am__tagged_files)
+ $(am__define_uniq_tagged_files); \
+ test -z "$(CTAGS_ARGS)$$unique" \
+ || $(CTAGS) $(CTAGSFLAGS) $(AM_CTAGSFLAGS) $(CTAGS_ARGS) \
+ $$unique
+
+GTAGS:
+ here=`$(am__cd) $(top_builddir) && pwd` \
+ && $(am__cd) $(top_srcdir) \
+ && gtags -i $(GTAGS_ARGS) "$$here"
+cscopelist: cscopelist-am
+
+cscopelist-am: $(am__tagged_files)
+ list='$(am__tagged_files)'; \
+ case "$(srcdir)" in \
+ [\\/]* | ?:[\\/]*) sdir="$(srcdir)" ;; \
+ *) sdir=$(subdir)/$(srcdir) ;; \
+ esac; \
+ for i in $$list; do \
+ if test -f "$$i"; then \
+ echo "$(subdir)/$$i"; \
+ else \
+ echo "$$sdir/$$i"; \
+ fi; \
+ done >> $(top_builddir)/cscope.files
+
+distclean-tags:
+ -rm -f TAGS ID GTAGS GRTAGS GSYMS GPATH tags
+
+distdir: $(BUILT_SOURCES)
+ $(MAKE) $(AM_MAKEFLAGS) distdir-am
+
+distdir-am: $(DISTFILES)
+ @srcdirstrip=`echo "$(srcdir)" | sed 's/[].[^$$\\*]/\\\\&/g'`; \
+ topsrcdirstrip=`echo "$(top_srcdir)" | sed 's/[].[^$$\\*]/\\\\&/g'`; \
+ list='$(DISTFILES)'; \
+ dist_files=`for file in $$list; do echo $$file; done | \
+ sed -e "s|^$$srcdirstrip/||;t" \
+ -e "s|^$$topsrcdirstrip/|$(top_builddir)/|;t"`; \
+ case $$dist_files in \
+ */*) $(MKDIR_P) `echo "$$dist_files" | \
+ sed '/\//!d;s|^|$(distdir)/|;s,/[^/]*$$,,' | \
+ sort -u` ;; \
+ esac; \
+ for file in $$dist_files; do \
+ if test -f $$file || test -d $$file; then d=.; else d=$(srcdir); fi; \
+ if test -d $$d/$$file; then \
+ dir=`echo "/$$file" | sed -e 's,/[^/]*$$,,'`; \
+ if test -d "$(distdir)/$$file"; then \
+ find "$(distdir)/$$file" -type d ! -perm -700 -exec chmod u+rwx {} \;; \
+ fi; \
+ if test -d $(srcdir)/$$file && test $$d != $(srcdir); then \
+ cp -fpR $(srcdir)/$$file "$(distdir)$$dir" || exit 1; \
+ find "$(distdir)/$$file" -type d ! -perm -700 -exec chmod u+rwx {} \;; \
+ fi; \
+ cp -fpR $$d/$$file "$(distdir)$$dir" || exit 1; \
+ else \
+ test -f "$(distdir)/$$file" \
+ || cp -p $$d/$$file "$(distdir)/$$file" \
+ || exit 1; \
+ fi; \
+ done
+check-am: all-am
+check: check-am
+all-am: Makefile $(LTLIBRARIES)
+installdirs:
+ for dir in "$(DESTDIR)$(pkglibdir)"; do \
+ test -z "$$dir" || $(MKDIR_P) "$$dir"; \
+ done
+install: install-am
+install-exec: install-exec-am
+install-data: install-data-am
+uninstall: uninstall-am
+
+install-am: all-am
+ @$(MAKE) $(AM_MAKEFLAGS) install-exec-am install-data-am
+
+installcheck: installcheck-am
+install-strip:
+ if test -z '$(STRIP)'; then \
+ $(MAKE) $(AM_MAKEFLAGS) INSTALL_PROGRAM="$(INSTALL_STRIP_PROGRAM)" \
+ install_sh_PROGRAM="$(INSTALL_STRIP_PROGRAM)" INSTALL_STRIP_FLAG=-s \
+ install; \
+ else \
+ $(MAKE) $(AM_MAKEFLAGS) INSTALL_PROGRAM="$(INSTALL_STRIP_PROGRAM)" \
+ install_sh_PROGRAM="$(INSTALL_STRIP_PROGRAM)" INSTALL_STRIP_FLAG=-s \
+ "INSTALL_PROGRAM_ENV=STRIPPROG='$(STRIP)'" install; \
+ fi
+mostlyclean-generic:
+
+clean-generic:
+
+distclean-generic:
+ -test -z "$(CONFIG_CLEAN_FILES)" || rm -f $(CONFIG_CLEAN_FILES)
+ -test . = "$(srcdir)" || test -z "$(CONFIG_CLEAN_VPATH_FILES)" || rm -f $(CONFIG_CLEAN_VPATH_FILES)
+
+maintainer-clean-generic:
+ @echo "This command is intended for maintainers to use"
+ @echo "it deletes files that may require special tools to rebuild."
+clean: clean-am
+
+clean-am: clean-generic clean-libtool clean-pkglibLTLIBRARIES \
+ mostlyclean-am
+
+distclean: distclean-am
+ -rm -f ./$(DEPDIR)/omazureeventhubs_la-omazureeventhubs.Plo
+ -rm -f Makefile
+distclean-am: clean-am distclean-compile distclean-generic \
+ distclean-tags
+
+dvi: dvi-am
+
+dvi-am:
+
+html: html-am
+
+html-am:
+
+info: info-am
+
+info-am:
+
+install-data-am:
+
+install-dvi: install-dvi-am
+
+install-dvi-am:
+
+install-exec-am: install-pkglibLTLIBRARIES
+
+install-html: install-html-am
+
+install-html-am:
+
+install-info: install-info-am
+
+install-info-am:
+
+install-man:
+
+install-pdf: install-pdf-am
+
+install-pdf-am:
+
+install-ps: install-ps-am
+
+install-ps-am:
+
+installcheck-am:
+
+maintainer-clean: maintainer-clean-am
+ -rm -f ./$(DEPDIR)/omazureeventhubs_la-omazureeventhubs.Plo
+ -rm -f Makefile
+maintainer-clean-am: distclean-am maintainer-clean-generic
+
+mostlyclean: mostlyclean-am
+
+mostlyclean-am: mostlyclean-compile mostlyclean-generic \
+ mostlyclean-libtool
+
+pdf: pdf-am
+
+pdf-am:
+
+ps: ps-am
+
+ps-am:
+
+uninstall-am: uninstall-pkglibLTLIBRARIES
+
+.MAKE: install-am install-strip
+
+.PHONY: CTAGS GTAGS TAGS all all-am am--depfiles check check-am clean \
+ clean-generic clean-libtool clean-pkglibLTLIBRARIES \
+ cscopelist-am ctags ctags-am distclean distclean-compile \
+ distclean-generic distclean-libtool distclean-tags distdir dvi \
+ dvi-am html html-am info info-am install install-am \
+ install-data install-data-am install-dvi install-dvi-am \
+ install-exec install-exec-am install-html install-html-am \
+ install-info install-info-am install-man install-pdf \
+ install-pdf-am install-pkglibLTLIBRARIES install-ps \
+ install-ps-am install-strip installcheck installcheck-am \
+ installdirs maintainer-clean maintainer-clean-generic \
+ mostlyclean mostlyclean-compile mostlyclean-generic \
+ mostlyclean-libtool pdf pdf-am ps ps-am tags tags-am uninstall \
+ uninstall-am uninstall-pkglibLTLIBRARIES
+
+.PRECIOUS: Makefile
+
+
+# Tell versions [3.59,3.63) of GNU make to not export all variables.
+# Otherwise a system limit (for SysV at least) may be exceeded.
+.NOEXPORT:
diff --git a/plugins/omazureeventhubs/omazureeventhubs.c b/plugins/omazureeventhubs/omazureeventhubs.c
new file mode 100644
index 0000000..a2bd8b3
--- /dev/null
+++ b/plugins/omazureeventhubs/omazureeventhubs.c
@@ -0,0 +1,1354 @@
+/* omazureeventhubs.c
+ * This output plugin make rsyslog talk to Azure EventHubs.
+ *
+ * Copyright 2014-2017 by Adiscon GmbH.
+ *
+ * This file is part of rsyslog.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ * -or-
+ * see COPYING.ASL20 in the source distribution
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+#include "config.h"
+#include <stdio.h>
+#include <stdarg.h>
+#include <stdlib.h>
+#include <string.h>
+#include <strings.h>
+#include <assert.h>
+#include <errno.h>
+#include <fcntl.h>
+#include <pthread.h>
+#include <sys/uio.h>
+#include <sys/queue.h>
+#include <sys/types.h>
+#include <math.h>
+#ifdef HAVE_SYS_STAT_H
+# include <sys/stat.h>
+#endif
+#include <sys/time.h>
+#include <time.h>
+
+// Include Proton headers
+#include <proton/version.h>
+#include <proton/condition.h>
+#include <proton/connection.h>
+#include <proton/delivery.h>
+#include <proton/link.h>
+#include <proton/listener.h>
+#include <proton/netaddr.h>
+#include <proton/message.h>
+#include <proton/object.h>
+#include <proton/proactor.h>
+#include <proton/sasl.h>
+#include <proton/session.h>
+#include <proton/transport.h>
+#include <proton/ssl.h>
+
+// Include rsyslog headers
+#include "rsyslog.h"
+#include "syslogd-types.h"
+#include "srUtils.h"
+#include "template.h"
+#include "module-template.h"
+#include "errmsg.h"
+#include "atomic.h"
+#include "statsobj.h"
+#include "unicode-helper.h"
+#include "datetime.h"
+#include "glbl.h"
+
+MODULE_TYPE_OUTPUT
+MODULE_TYPE_NOKEEP
+MODULE_CNFNAME("omazureeventhubs")
+
+/* internal structures
+ */
+DEF_OMOD_STATIC_DATA
+DEFobjCurrIf(glbl)
+DEFobjCurrIf(datetime)
+DEFobjCurrIf(strm)
+DEFobjCurrIf(statsobj)
+
+statsobj_t *azureStats;
+STATSCOUNTER_DEF(ctrMessageSubmit, mutCtrMessageSubmit);
+STATSCOUNTER_DEF(ctrAzureFail, mutCtrAzureFail);
+STATSCOUNTER_DEF(ctrCacheMiss, mutCtrCacheMiss);
+STATSCOUNTER_DEF(ctrCacheEvict, mutCtrCacheEvict);
+STATSCOUNTER_DEF(ctrCacheSkip, mutCtrCacheSkip);
+STATSCOUNTER_DEF(ctrAzureAck, mutCtrAzureAck);
+STATSCOUNTER_DEF(ctrAzureMsgTooLarge, mutCtrAzureMsgTooLarge);
+STATSCOUNTER_DEF(ctrAzureQueueFull, mutCtrAzureQueueFull);
+STATSCOUNTER_DEF(ctrAzureOtherErrors, mutCtrAzureOtherErrors);
+STATSCOUNTER_DEF(ctrAzureRespTimedOut, mutCtrAzureRespTimedOut);
+STATSCOUNTER_DEF(ctrAzureRespTransport, mutCtrAzureRespTransport);
+STATSCOUNTER_DEF(ctrAzureRespBrokerDown, mutCtrAzureRespBrokerDown);
+STATSCOUNTER_DEF(ctrAzureRespAuth, mutCtrAzureRespAuth);
+STATSCOUNTER_DEF(ctrAzureRespSSL, mutCtrAzureRespSSL);
+STATSCOUNTER_DEF(ctrAzureRespOther, mutCtrAzureRespOther);
+
+#define MAX_ERRMSG 1024 /* max size of error messages that we support */
+#define MAX_DEFAULTMSGS 1024 /* Initial max size of the proton message helper array */
+
+#define SETUP_PROTON_NONE 0
+#define SETUP_PROTON_AUTOCLOSE 1
+
+/* flags for writeAzure: shall we resubmit a failed message? */
+#define RESUBMIT 1
+#define NO_RESUBMIT 0
+
+/* flags for transaction Handling */
+enum proton_submission_status
+{
+ PROTON_UNSUBMITTED = 0, // Message not submitted yet
+ PROTON_SUBMITTED, // Message submitted to proton sender instance
+ PROTON_ACCEPTED, // Message accepted from remote target
+ PROTON_REJECTED, // Message rejected from remote target (zero credit?)
+};
+
+// event_property NEEDED?
+struct event_property {
+ const char *key;
+ const char *val;
+};
+
+static pn_timestamp_t time_now(void);
+
+/* Struct for Proton Messages Listitems */
+struct s_protonmsg_entry {
+ uchar* payload;
+ size_t payload_len;
+ uchar* MsgID;
+ size_t MsgID_len;
+
+ uchar* address;
+ char status;
+};
+typedef struct s_protonmsg_entry protonmsg_entry;
+
+/* Struct for module InstanceData */
+typedef struct _instanceData {
+ uchar *amqp_address;
+ uchar *azurehost;
+ uchar *azureport;
+ uchar *azure_key_name;
+ uchar *azure_key;
+ uchar *container;
+ uchar *tplName; /* assigned output template */
+
+ int nEventProperties;
+ struct event_property *eventProperties;
+
+ uchar *statsName;
+ statsobj_t *stats;
+ STATSCOUNTER_DEF(ctrMessageSubmit, mutCtrMessageSubmit);
+ STATSCOUNTER_DEF(ctrAzureFail, mutCtrAzureFail);
+ STATSCOUNTER_DEF(ctrAzureAck, mutCtrAzureAck);
+ STATSCOUNTER_DEF(ctrAzureOtherErrors, mutCtrAzureOtherErrors);
+} instanceData;
+
+/* Struct for module workerInstanceData */
+typedef struct wrkrInstanceData {
+ instanceData *pData;
+
+ protonmsg_entry **aProtonMsgs; /* dynamically sized array for transactional outputs */
+ unsigned int nProtonMsgs; /* current used proton msgs */
+ unsigned int nMaxProtonMsgs; /* current max */
+
+ int bIsConnecting; /* 1 if connecting, 0 if disconnected */
+ int bIsConnected; /* 1 if connected, 0 if disconnected */
+ int bIsSuspended; /* when broker fail, we need to suspend the action */
+ pthread_rwlock_t pnLock;
+
+ // PROTON Handles
+ pn_proactor_t *pnProactor;
+ pn_transport_t *pnTransport;
+ pn_connection_t *pnConn;
+ pn_link_t* pnSender;
+ pn_rwbytes_t pnMessageBuffer; /* Buffer for messages */
+ int pnStatus;
+
+ // Message Counters for sender link in worker instance
+ unsigned int iMsgSeq;
+ unsigned int iMaxMsgSeq;
+
+ /* The following structure controls the proton handling threads. Passes necessary pointers
+ * needed for their access.
+ */
+ sbool bThreadRunning;
+ pthread_t tid; /* the worker's thread ID */
+
+} wrkrInstanceData_t;
+
+#define INST_STATSCOUNTER_INC(inst, ctr, mut) \
+ do { \
+ if (inst->stats) { STATSCOUNTER_INC(ctr, mut); } \
+ } while(0);
+
+// QPID Proton Handler functions
+static rsRetVal proton_run_thread(wrkrInstanceData_t *pWrkrData);
+static rsRetVal proton_shutdown_thread(wrkrInstanceData_t *pWrkrData);
+static void * proton_thread(void *myInfo);
+static void handleProtonDelivery(wrkrInstanceData_t *const pWrkrData);
+static void handleProton(wrkrInstanceData_t *const pWrkrData, pn_event_t *event);
+static rsRetVal writeProton(wrkrInstanceData_t *__restrict__ const pWrkrData,
+ const actWrkrIParams_t *__restrict__ const pParam,
+ const int iMsg);
+
+/* tables for interfacing with the v6 config system */
+/* action (instance) parameters */
+static struct cnfparamdescr actpdescr[] = {
+ { "azurehost", eCmdHdlrString, CNFPARAM_REQUIRED },
+ { "azureport", eCmdHdlrString, CNFPARAM_REQUIRED },
+ { "azure_key_name", eCmdHdlrString, CNFPARAM_REQUIRED },
+ { "azure_key", eCmdHdlrString, CNFPARAM_REQUIRED },
+ { "amqp_address", eCmdHdlrString, 0 },
+ { "container", eCmdHdlrString, 0 },
+ { "eventproperties", eCmdHdlrArray, 0 },
+ { "template", eCmdHdlrGetWord, 0 },
+ { "statsname", eCmdHdlrGetWord, 0 }
+};
+static struct cnfparamblk actpblk =
+ { CNFPARAMBLK_VERSION,
+ sizeof(actpdescr)/sizeof(struct cnfparamdescr),
+ actpdescr
+ };
+
+BEGINinitConfVars /* (re)set config variables to default values */
+CODESTARTinitConfVars
+ENDinitConfVars
+
+static void ATTR_NONNULL(1)
+protonmsg_entry_destruct(protonmsg_entry *const __restrict__ fmsgEntry) {
+ free(fmsgEntry->MsgID);
+ free(fmsgEntry->payload);
+ free(fmsgEntry->address);
+ free(fmsgEntry);
+}
+
+/* note: we need the length of message as we need to deal with
+ * non-NUL terminated strings under some circumstances.
+ */
+static protonmsg_entry * ATTR_NONNULL(1,3)
+protonmsg_entry_construct( const char *const MsgID, const size_t msgidlen,
+ const char *const msg, const size_t msglen,
+ const char *const address)
+{
+ protonmsg_entry *etry = NULL;
+
+ if((etry = malloc(sizeof(struct s_protonmsg_entry))) == NULL) {
+ return NULL;
+ }
+ etry->status = PROTON_UNSUBMITTED; // Unsubmitted default */
+
+ etry->MsgID_len = msgidlen;
+ if((etry->MsgID = (uchar*)malloc(msgidlen+1)) == NULL) {
+ free(etry);
+ return NULL;
+ }
+ memcpy(etry->MsgID, MsgID, msgidlen);
+ etry->MsgID[msgidlen] = '\0';
+
+ etry->payload_len = msglen;
+ if((etry->payload = (uchar*)malloc(msglen+1)) == NULL) {
+ free(etry->MsgID);
+ free(etry);
+ return NULL;
+ }
+ memcpy(etry->payload, msg, msglen);
+ etry->payload[msglen] = '\0';
+
+ if((etry->address = (uchar*)strdup(address)) == NULL) {
+ free(etry->MsgID);
+ free(etry->payload);
+ free(etry);
+ return NULL;
+ }
+ return etry;
+}
+
+/* Create a message with a map { "sequence" : number } encode it and return the encoded buffer. */
+static pn_message_t* proton_encode_message(wrkrInstanceData_t *const pWrkrData, protonmsg_entry* pMsgEntry) {
+ instanceData *const pData = (instanceData *const) pWrkrData->pData;
+ /* Construct a message with the map */
+ pn_message_t* message = pn_message();
+ // Optionally include Address?
+ // pn_message_set_address(message, (char *) pWrkrData->amqp_address);
+
+ // Send in BINARY MODE ( as Stream )
+ pn_message_set_content_type(message, (char*) "application/octect-stream");
+ pn_message_set_creation_time(message, time_now());
+ pn_message_set_inferred(message, true);
+ // Set message ID
+ pn_message_set_id(message, (pn_atom_t){
+ .type=PN_STRING,
+ .u.as_bytes.start = (char*)pMsgEntry->MsgID,
+ .u.as_bytes.size = pMsgEntry->MsgID_len
+ });
+
+ if (pData->nEventProperties > 0) {
+ // Add Event properties
+ pn_data_t *props = pn_message_properties(message);
+ pn_data_put_map(props);
+ pn_data_enter(props);
+
+ for(int i = 0 ; i < pData->nEventProperties ; ++i) {
+ DBGPRINTF("proton_encode_message: add eventproperty %s:%s\n",
+ pData->eventProperties[i].key,
+ pData->eventProperties[i].val);
+ pn_data_put_string(props, pn_bytes(strlen(pData->eventProperties[i].key),
+ pData->eventProperties[i].key));
+ pn_data_put_string(props, pn_bytes(strlen(pData->eventProperties[i].val),
+ pData->eventProperties[i].val));
+ }
+ pn_data_exit(props);
+ }
+
+ // Set message BODY
+ pn_data_t* body = pn_message_body(message);
+ pn_data_enter(body);
+ pn_data_put_binary(body, pn_bytes(pMsgEntry->payload_len, (char*)pMsgEntry->payload));
+ pn_data_exit(body);
+
+ DBGPRINTF("proton_encode_message: created message id '%s': '%.*s'\n",
+ (char*)pMsgEntry->MsgID,
+ (pMsgEntry->payload_len > 0 ? (int)pMsgEntry->payload_len-1 : 0),
+ (char*)pMsgEntry->payload);
+
+ return message;
+}
+
+static rsRetVal
+closeProton(wrkrInstanceData_t *const __restrict__ pWrkrData)
+{
+ DEFiRet;
+ instanceData *const pData = (instanceData *const) pWrkrData->pData;
+#ifndef NDEBUG
+ DBGPRINTF("closeProton[%p]: ENTER\n", pWrkrData);
+#endif
+ if (pWrkrData->pnSender) {
+ pn_link_close(pWrkrData->pnSender);
+ DBGPRINTF("closeProton[%p]: pn_link_close\n", pWrkrData);
+ pn_session_close(pn_link_session(pWrkrData->pnSender));
+ DBGPRINTF("closeProton[%p]: pn_session_close\n", pWrkrData);
+ }
+ if (pWrkrData->pnConn) {
+ DBGPRINTF("closeProton[%p]: pn_connection_close connection\n", pWrkrData);
+ pn_connection_close(pWrkrData->pnConn);
+ }
+
+ pWrkrData->bIsConnecting = 0;
+ pWrkrData->bIsConnected = 0;
+ pWrkrData->pnStatus = PN_EVENT_NONE;
+
+ pWrkrData->pnSender = NULL;
+ pWrkrData->pnConn = NULL;
+ pWrkrData->iMsgSeq = 0;
+ pWrkrData->iMaxMsgSeq = 0;
+
+ // Mark all remaining entries as REJECTED
+ if(pWrkrData->aProtonMsgs != NULL) {
+ for(unsigned int i = 0 ; i < pWrkrData->nProtonMsgs ; ++i) {
+ if (pWrkrData->aProtonMsgs[i] != NULL && (
+ pWrkrData->aProtonMsgs[i]->status == PROTON_UNSUBMITTED ||
+ pWrkrData->aProtonMsgs[i]->status == PROTON_SUBMITTED)
+ ) {
+ pWrkrData->aProtonMsgs[i]->status = PROTON_REJECTED;
+ DBGPRINTF("closeProton[%p]: Setting ProtonMsg %s to PROTON_REJECTED \n",
+ pWrkrData, pWrkrData->aProtonMsgs[i]->MsgID);
+ // Increment Stats Counter
+ STATSCOUNTER_INC(ctrAzureFail, mutCtrAzureFail);
+ INST_STATSCOUNTER_INC(pData, pData->ctrAzureFail, pData->mutCtrAzureFail);
+ }
+ }
+ }
+
+ FINALIZE;
+finalize_it:
+ RETiRet;
+
+}
+
+static rsRetVal
+openProton(wrkrInstanceData_t *const __restrict__ pWrkrData)
+{
+ DEFiRet;
+ instanceData *const pData = (instanceData *const) pWrkrData->pData;
+ int pnErr = PN_OK;
+ char szAddr[PN_MAX_ADDR];
+ pn_ssl_t* pnSsl;
+#ifndef NDEBUG
+ DBGPRINTF("openProton[%p]: ENTER\n", pWrkrData);
+#endif
+ if(pWrkrData->bIsConnecting == 1 || pWrkrData->bIsConnected == 1)
+ FINALIZE;
+ pWrkrData->pnStatus = PN_EVENT_NONE;
+
+ pn_proactor_addr(szAddr, sizeof(szAddr),
+ (const char *) pData->azurehost,
+ (const char *) pData->azureport);
+
+ // Configure a transport for SSL. The transport will be freed by the proactor.
+ pWrkrData->pnTransport = pn_transport();
+ DBGPRINTF("openProton[%p]: create transport to '%s:%s'\n",
+ pWrkrData, pData->azurehost, pData->azureport);
+ pnSsl = pn_ssl(pWrkrData->pnTransport);
+ if (pnSsl != NULL) {
+ pn_ssl_domain_t* pnDomain = pn_ssl_domain(PN_SSL_MODE_CLIENT);
+ if (pData->azure_key_name != NULL && pData->azure_key != NULL) {
+ pnErr = pn_ssl_init(pnSsl, pnDomain, NULL);
+ if (pnErr) {
+ DBGPRINTF("openProton[%p]: pn_ssl_init failed for '%s:%s' with error %d: %s\n",
+ pWrkrData, pData->azurehost, pData->azureport,
+ pnErr, pn_code(pnErr));
+ }
+ pn_sasl_allowed_mechs(pn_sasl(pWrkrData->pnTransport), "PLAIN");
+ } else {
+ pnErr = pn_ssl_domain_set_peer_authentication(pnDomain, PN_SSL_ANONYMOUS_PEER, NULL);
+ if (!pnErr) {
+ pnErr = pn_ssl_init(pnSsl, pnDomain, NULL);
+ } else {
+ DBGPRINTF(
+ "openProton[%p]: pn_ssl_domain_set_peer_authentication failed with '%d'\n",
+ pWrkrData, pnErr);
+ }
+ }
+ pn_ssl_domain_free(pnDomain);
+ } else {
+ LogError(0, RS_RET_ERR, "openProton[%p]: openProton pn_ssl_init NULL", pWrkrData);
+ }
+
+ // Handle ERROR Output
+ if (pnErr) {
+ LogError(0, RS_RET_IO_ERROR, "openProton[%p]: creating transport to '%s:%s' "
+ "failed with error %d: %s\n",
+ pWrkrData, pData->azurehost, pData->azureport,
+ pnErr, pn_code(pnErr));
+ ABORT_FINALIZE(RS_RET_IO_ERROR);
+ }
+
+ // Connect to Azure Event Hubs
+ pn_proactor_connect2(pWrkrData->pnProactor, NULL, pWrkrData->pnTransport, szAddr);
+
+ // Successfully connecting
+ pWrkrData->bIsConnecting = 1;
+ pWrkrData->bIsSuspended = 0;
+finalize_it:
+ if(iRet != RS_RET_OK) {
+ closeProton(pWrkrData); // Make sure to free ressources
+ }
+ RETiRet;
+}
+
+static sbool
+proton_check_condition( pn_event_t *event,
+ wrkrInstanceData_t *const __restrict__ pWrkrData,
+ pn_condition_t *cond,
+ const char * pszReason) {
+ if (pn_condition_is_set(cond)) {
+ DBGPRINTF("proton_check_condition: %s %s: %s: %s",
+ pszReason,
+ pn_event_type_name(pn_event_type(event)),
+ pn_condition_get_name(cond),
+ pn_condition_get_description(cond));
+ LogError(0, RS_RET_ERR, "omazureeventhubs: %s %s: %s: %s",
+ pszReason,
+ pn_event_type_name(pn_event_type(event)),
+ pn_condition_get_name(cond),
+ pn_condition_get_description(cond));
+
+ // Connection can be closed
+ closeProton(pWrkrData);
+
+ // Set Worker to suspended state!
+ pWrkrData->bIsSuspended = 1;
+
+ return 0;
+ } else {
+ return 1;
+ }
+}
+
+static rsRetVal
+setupProtonHandle(wrkrInstanceData_t *const __restrict__ pWrkrData, int autoclose)
+{
+ DEFiRet;
+ DBGPRINTF("omazureeventhubs[%p]: setupProtonHandle ENTER\n", pWrkrData);
+
+ pthread_rwlock_wrlock(&pWrkrData->pnLock);
+ if (autoclose == SETUP_PROTON_AUTOCLOSE && (pWrkrData->bIsConnected == 1)) {
+ DBGPRINTF("omazureeventhubs[%p]: setupProtonHandle closeProton\n", pWrkrData);
+ closeProton(pWrkrData);
+ }
+ CHKiRet(openProton(pWrkrData));
+finalize_it:
+ if (iRet != RS_RET_OK) {
+ /* Parameter Error's cannot be resumed, so we need to disable the action */
+ if (iRet == RS_RET_PARAM_ERROR) {
+ iRet = RS_RET_DISABLE_ACTION;
+ LogError(0, iRet, "omazureeventhubs: action will be disabled due invalid "
+ "configuration parameters\n");
+ }
+ }
+ pthread_rwlock_unlock(&pWrkrData->pnLock);
+ RETiRet;
+}
+
+static rsRetVal
+writeProton(wrkrInstanceData_t *__restrict__ const pWrkrData,
+ const actWrkrIParams_t *__restrict__ const pParam,
+ const int iMsg)
+{
+ DEFiRet;
+ instanceData *const pData = (instanceData *const) pWrkrData->pData;
+ protonmsg_entry* fmsgEntry;
+
+ // Create Unqiue Message ID
+ char szMsgID[64];
+ sprintf(szMsgID, "%d", pWrkrData->iMsgSeq);
+
+ const char* pszParamStr = (const char*)actParam(pParam, 1 /*pData->iNumTpls*/, iMsg, 0).param;
+ size_t tzParamStrLen = actParam(pParam, 1 /*pData->iNumTpls*/, iMsg, 0).lenStr;
+
+ DBGPRINTF("omazureeventhubs[%p]: writeProton for msg %d (seq %d) msg:'%.*s%s'\n",
+ pWrkrData,
+ iMsg, pWrkrData->iMsgSeq,
+ (int)(tzParamStrLen > 0 ? (tzParamStrLen > 64 ? 64 : tzParamStrLen-1) : 0),
+ pszParamStr,
+ (tzParamStrLen > 64 ? "..." : ""));
+ // Increment Message sequence number
+ pWrkrData->iMsgSeq++;
+
+ // Add message to LIST for sending
+ CHKmalloc(fmsgEntry = protonmsg_entry_construct(
+ szMsgID, sizeof(szMsgID),
+ pszParamStr,
+ tzParamStrLen,
+ (const char*)pData->amqp_address));
+ // Add to helper Array
+ pWrkrData->aProtonMsgs[iMsg] = fmsgEntry;
+finalize_it:
+ RETiRet;
+}
+
+BEGINcreateInstance
+CODESTARTcreateInstance
+ DBGPRINTF("createInstance[%p]: ENTER\n", pData);
+ pData->amqp_address = NULL;
+ pData->azurehost = NULL;
+ pData->azureport = NULL;
+ pData->azure_key_name = NULL;
+ pData->azure_key = NULL;
+ pData->container = NULL;
+ENDcreateInstance
+
+
+BEGINcreateWrkrInstance
+CODESTARTcreateWrkrInstance
+ DBGPRINTF("createWrkrInstance[%p]: ENTER\n", pWrkrData);
+ pWrkrData->bIsConnecting = 0;
+ pWrkrData->bIsConnected = 0;
+ pWrkrData->bIsSuspended = 0;
+
+ // Create Proton proActor in Worker Instance
+ pWrkrData->pnProactor = pn_proactor();
+ pWrkrData->pnConn = NULL;
+ pWrkrData->pnTransport = NULL;
+ pWrkrData->pnSender = NULL;
+
+ pWrkrData->iMsgSeq = 0;
+ pWrkrData->iMaxMsgSeq = 0;
+ pWrkrData->pnMessageBuffer.start = NULL;
+
+ pWrkrData->nProtonMsgs = 0;
+ pWrkrData->nMaxProtonMsgs = MAX_DEFAULTMSGS;
+ CHKmalloc(pWrkrData->aProtonMsgs = calloc(MAX_DEFAULTMSGS, sizeof(struct s_protonmsg_entry)));
+
+ CHKiRet(pthread_rwlock_init(&pWrkrData->pnLock, NULL));
+
+ pWrkrData->bThreadRunning = 0;
+ pWrkrData->tid = 0;
+
+ // Run Proton Worker Thread now
+ proton_run_thread(pWrkrData);
+finalize_it:
+ DBGPRINTF("createWrkrInstance[%p] returned %d\n", pWrkrData, iRet);
+ENDcreateWrkrInstance
+
+
+BEGINisCompatibleWithFeature
+CODESTARTisCompatibleWithFeature
+ENDisCompatibleWithFeature
+
+
+BEGINfreeInstance
+CODESTARTfreeInstance
+ DBGPRINTF("freeInstance[%p]: ENTER\n", pData);
+
+ if (pData->stats) {
+ statsobj.Destruct(&pData->stats);
+ }
+
+ /* Free other mem */
+ free(pData->amqp_address);
+ free(pData->azurehost);
+ free(pData->azureport);
+ free(pData->azure_key_name);
+ free(pData->azure_key);
+ free(pData->container);
+
+ free(pData->tplName);
+ free(pData->statsName);
+ for(int i = 0 ; i < pData->nEventProperties ; ++i) {
+ free((void*) pData->eventProperties[i].key);
+ free((void*) pData->eventProperties[i].val);
+ }
+ free(pData->eventProperties);
+ENDfreeInstance
+
+BEGINfreeWrkrInstance
+CODESTARTfreeWrkrInstance
+ DBGPRINTF("freeWrkrInstance[%p]: ENTER\n", pWrkrData);
+
+ /* Closing azure first! */
+ pthread_rwlock_wrlock(&pWrkrData->pnLock);
+
+ // Close Proton Connection
+ closeProton(pWrkrData);
+
+ // Stop Proton Handle Thread
+ proton_shutdown_thread(pWrkrData);
+
+ // Free Proton Ressources
+ if (pWrkrData->pnProactor != NULL) {
+ DBGPRINTF("freeWrkrInstance[%p]: FREE proactor\n", pWrkrData);
+ pn_proactor_free(pWrkrData->pnProactor);
+ pWrkrData->pnProactor = NULL;
+ }
+ free(pWrkrData->pnMessageBuffer.start);
+
+ pthread_rwlock_unlock(&pWrkrData->pnLock);
+
+ // Free our proton helper array
+ if(pWrkrData->aProtonMsgs != NULL) {
+ for(unsigned int i = 0 ; i < pWrkrData->nProtonMsgs ; ++i) {
+ if (pWrkrData->aProtonMsgs[i] != NULL) {
+ protonmsg_entry_destruct(pWrkrData->aProtonMsgs[i]);
+ }
+ }
+ free(pWrkrData->aProtonMsgs);
+ }
+
+ pthread_rwlock_destroy(&pWrkrData->pnLock);
+ENDfreeWrkrInstance
+
+
+BEGINdbgPrintInstInfo
+CODESTARTdbgPrintInstInfo
+ENDdbgPrintInstInfo
+
+BEGINtryResume
+CODESTARTtryResume
+#ifndef NDEBUG
+ DBGPRINTF("omazureeventhubs[%p]: tryResume ENTER\n", pWrkrData);
+#endif
+ if (pWrkrData->bIsConnecting == 0 && pWrkrData->bIsConnected == 0) {
+ DBGPRINTF("omazureeventhubs[%p]: tryResume setupProtonHandle\n", pWrkrData);
+ CHKiRet(setupProtonHandle(pWrkrData, SETUP_PROTON_AUTOCLOSE));
+ }
+finalize_it:
+ DBGPRINTF("omazureeventhubs[%p]: tryResume returned %d\n", pWrkrData, iRet);
+ENDtryResume
+
+BEGINbeginTransaction
+CODESTARTbeginTransaction
+ /* we have nothing to do to begin a transaction */
+ DBGPRINTF("omazureeventhubs[%p]: beginTransaction ENTER\n", pWrkrData);
+ if (pWrkrData->bIsConnecting == 0 && pWrkrData->bIsConnected == 0) {
+ CHKiRet(setupProtonHandle(pWrkrData, SETUP_PROTON_NONE));
+ }
+finalize_it:
+ RETiRet;
+ENDbeginTransaction
+
+/*
+ * New Transaction action interface
+*/
+BEGINcommitTransaction
+ // instanceData *__restrict__ const pData = pWrkrData->pData;
+ unsigned i;
+ unsigned iNeedSubmission;
+ sbool bDone = 0;
+ protonmsg_entry* pMsgEntry = NULL;
+CODESTARTcommitTransaction
+#ifndef NDEBUG
+ DBGPRINTF("omazureeventhubs[%p]: commitTransaction [%d msgs] ENTER\n", pWrkrData, nParams);
+#endif
+
+ // Handle/Expand our proton helper array
+ if (nParams > pWrkrData->nMaxProtonMsgs) {
+ // Free old Array
+ if(pWrkrData->aProtonMsgs != NULL) {
+ free(pWrkrData->aProtonMsgs);
+ }
+ // Expand our proton helper array
+ DBGPRINTF("omazureeventhubs[%p]: commitTransaction expand helper array from %d to %d\n",
+ pWrkrData, pWrkrData->nMaxProtonMsgs, nParams);
+ pWrkrData->nMaxProtonMsgs = nParams; // Set new MAX
+ CHKmalloc(pWrkrData->aProtonMsgs = calloc(pWrkrData->nMaxProtonMsgs, sizeof(struct s_protonmsg_entry)));
+ }
+ // Copy count of New Messages and increase MaxMsgSeq
+ pWrkrData->nProtonMsgs = nParams;
+ pWrkrData->iMaxMsgSeq += nParams;
+
+ do {
+ iNeedSubmission = 0;
+ // Put unsubmitted messages into Proton
+ for(i = 0 ; i < nParams ; ++i) {
+ // Get reference to Proton Array Helper
+ pMsgEntry = ((protonmsg_entry*)pWrkrData->aProtonMsgs[i]);
+ if ( pMsgEntry == NULL) {
+ // Send Message to Proton
+ writeProton(pWrkrData, pParams, i);
+ } else if ( pMsgEntry->status == PROTON_REJECTED) {
+ // Reset Message Entry, it will be retried
+ pMsgEntry->status = PROTON_UNSUBMITTED;
+ }
+ }
+ bDone = 1;
+
+ // Wait 100 microseconds
+ srSleep(0, 100000);
+
+ // Verify if messages have been submitted successfully
+ for(i = 0 ; i < nParams ; ++i) {
+ // Get reference to Proton Array Helper
+ pMsgEntry = ((protonmsg_entry*)pWrkrData->aProtonMsgs[i]);
+ if (pMsgEntry != NULL) {
+ if ( pMsgEntry->status == PROTON_UNSUBMITTED) {
+ iNeedSubmission++;
+ // we need to retry check later
+ bDone = 0;
+ } else if ( pMsgEntry->status == PROTON_SUBMITTED) {
+ // we need to retry check later
+ bDone = 0;
+ }
+ }
+ }
+
+ if (iNeedSubmission > 0) {
+ if ( pWrkrData->bIsConnected == 1) {
+ int credits = pn_link_credit(pWrkrData->pnSender);
+ if (pn_link_credit(pWrkrData->pnSender) > 0) {
+ DBGPRINTF("omazureeventhubs[%p]: trigger pn_connection_wake\n",
+ pWrkrData);
+ pn_connection_wake(pWrkrData->pnConn);
+ } else {
+ DBGPRINTF("omazureeventhubs[%p]: warning pn_link_credit returned %d\n",
+ pWrkrData, credits);
+ }
+ } else {
+ DBGPRINTF("omazureeventhubs[%p]: commitTransaction Suspended=%s Connecting=%s\n",
+ pWrkrData,
+ pWrkrData->bIsSuspended == 1 ? "YES" : "NO",
+ pWrkrData->bIsConnecting == 1 ? "YES" : "NO");
+ if (pWrkrData->bIsSuspended == 1 && pWrkrData->bIsConnecting == 0) {
+ ABORT_FINALIZE(RS_RET_SUSPENDED);
+ }
+ }
+ }
+ } while (bDone == 0);
+finalize_it:
+ // Free Proton Message Helpers
+ if (pWrkrData->aProtonMsgs != NULL) {
+ for(i = 0 ; i < nParams ; ++i) {
+ if (pWrkrData->aProtonMsgs[i] != NULL) {
+ // Destroy
+ protonmsg_entry_destruct(pWrkrData->aProtonMsgs[i]);
+ pWrkrData->aProtonMsgs[i] = NULL;
+ }
+ }
+ }
+
+ /* TODO: Suspend Action if broker problems were reported in error callback */
+ if (pWrkrData->bIsSuspended == 1) {
+ DBGPRINTF("omazureeventhubs[%p]: commitTransaction failed to send messages, suspending action\n",
+ pWrkrData);
+ iRet = RS_RET_SUSPENDED;
+ }
+ if(iRet != RS_RET_OK) {
+ DBGPRINTF("omazureeventhubs[%p]: commitTransaction failed with status %d\n", pWrkrData, iRet);
+ }
+ DBGPRINTF("omazureeventhubs[%p]: commitTransaction [%d msgs] EXIT\n", pWrkrData, nParams);
+ENDcommitTransaction
+
+static void
+setInstParamDefaults(instanceData *pData) {
+ DBGPRINTF("setInstParamDefaults[%p]: ENTER\n", pData);
+ pData->amqp_address = NULL;
+ pData->azurehost = NULL;
+ pData->azureport = NULL;
+ pData->azure_key_name = NULL;
+ pData->azure_key = NULL;
+ pData->container = NULL;
+ pData->nEventProperties = 0;
+ pData->eventProperties = NULL;
+}
+
+static rsRetVal
+processEventProperty(char *const param,
+ const char **const name,
+ const char **const paramval) {
+ DEFiRet;
+ char *val = strstr(param, "=");
+ if(val == NULL) {
+ LogError(0, RS_RET_PARAM_ERROR, "missing equal sign in "
+ "parameter '%s'", param);
+ ABORT_FINALIZE(RS_RET_PARAM_ERROR);
+ }
+ *val = '\0'; /* terminates name */
+ ++val; /* now points to begin of value */
+ CHKmalloc(*name = strdup(param));
+ CHKmalloc(*paramval = strdup(val));
+finalize_it:
+ RETiRet;
+}
+
+BEGINnewActInst
+ struct cnfparamvals *pvals;
+ int i;
+ int iNumTpls;
+ DBGPRINTF("newActInst: ENTER\n");
+CODESTARTnewActInst
+ if((pvals = nvlstGetParams(lst, &actpblk, NULL)) == NULL) {
+ ABORT_FINALIZE(RS_RET_MISSING_CNFPARAMS);
+ }
+
+ CHKiRet(createInstance(&pData));
+ setInstParamDefaults(pData);
+
+ for(i = 0 ; i < actpblk.nParams ; ++i) {
+ if(!pvals[i].bUsed)
+ continue;
+ if(!strcmp(actpblk.descr[i].name, "amqp_address")) {
+ pData->amqp_address = (uchar*)es_str2cstr(pvals[i].val.d.estr, NULL);
+ } else if(!strcmp(actpblk.descr[i].name, "azurehost")) {
+ pData->azurehost = (uchar*)es_str2cstr(pvals[i].val.d.estr, NULL);
+ } else if(!strcmp(actpblk.descr[i].name, "azureport")) {
+ pData->azureport = (uchar*)es_str2cstr(pvals[i].val.d.estr, NULL);
+ } else if(!strcmp(actpblk.descr[i].name, "azure_key_name")) {
+ pData->azure_key_name = (uchar*)es_str2cstr(pvals[i].val.d.estr, NULL);
+ } else if(!strcmp(actpblk.descr[i].name, "azure_key")) {
+ pData->azure_key = (uchar*)es_str2cstr(pvals[i].val.d.estr, NULL);
+ } else if(!strcmp(actpblk.descr[i].name, "container")) {
+ pData->container = (uchar*)es_str2cstr(pvals[i].val.d.estr, NULL);
+ } else if(!strcmp(actpblk.descr[i].name, "eventproperties")) {
+ pData->nEventProperties = pvals[i].val.d.ar->nmemb;
+ CHKmalloc(pData->eventProperties = malloc(sizeof(struct event_property) *
+ pvals[i].val.d.ar->nmemb ));
+ for(int j = 0 ; j < pvals[i].val.d.ar->nmemb ; ++j) {
+ char *cstr = es_str2cstr(pvals[i].val.d.ar->arr[j], NULL);
+ CHKiRet(processEventProperty(cstr, &pData->eventProperties[j].key,
+ &pData->eventProperties[j].val));
+ free(cstr);
+ }
+ } else if(!strcmp(actpblk.descr[i].name, "template")) {
+ pData->tplName = (uchar*)es_str2cstr(pvals[i].val.d.estr, NULL);
+ } else if(!strcmp(actpblk.descr[i].name, "statsname")) {
+ pData->statsName = (uchar*)es_str2cstr(pvals[i].val.d.estr, NULL);
+ } else {
+ LogError(0, RS_RET_INTERNAL_ERROR,
+ "omazureeventhubs: program error, non-handled param '%s'\n", actpblk.descr[i].name);
+ }
+ }
+
+ if(pData->azure_key_name == NULL || pData->azure_key == NULL) {
+ LogError(0, RS_RET_CONFIG_ERROR,
+ "omazureeventhubs: azure_key_name and azure_key are requires to access azure eventhubs"
+ " - action definition invalid");
+ ABORT_FINALIZE(RS_RET_CONFIG_ERROR);
+ }
+
+ if(pData->container == NULL) {
+ LogError(0, RS_RET_CONFIG_ERROR,
+ "omazureeventhubs: Event Hubs \"container\" parameter (which is instance) not specified "
+ " - action definition invalid");
+ ABORT_FINALIZE(RS_RET_CONFIG_ERROR);
+ }
+
+ if(pData->amqp_address == NULL) {
+ if(pData->azurehost == NULL) {
+ LogMsg(0, NO_ERRCODE, LOG_INFO, "omazureeventhubs: \"azurehost\" parameter not specified "
+ "(youreventhubinstance.servicebus.windows.net- action definition invalid!");
+ ABORT_FINALIZE(RS_RET_CONFIG_ERROR);
+ }
+ if(pData->azureport== NULL) {
+ // Set default
+ CHKmalloc(pData->azureport = (uchar *) strdup("5671"));
+ }
+
+ // Create amqps URL from parameters
+ char szAddress[1024];
+ sprintf(szAddress, "amqps://%s:%s@%s:%s/%s",
+ pData->azure_key_name,
+ pData->azure_key,
+ pData->azurehost,
+ pData->azureport,
+ pData->container);
+ CHKmalloc(pData->amqp_address = (uchar*) strdup(szAddress));
+ }
+
+ iNumTpls = 1;
+ CODE_STD_STRING_REQUESTnewActInst(iNumTpls);
+
+ CHKiRet(OMSRsetEntry(*ppOMSR, 0, (uchar*)strdup((pData->tplName == NULL) ?
+ "RSYSLOG_FileFormat" : (char*)pData->tplName),
+ OMSR_NO_RQD_TPL_OPTS));
+
+ if (pData->statsName) {
+ CHKiRet(statsobj.Construct(&pData->stats));
+ CHKiRet(statsobj.SetName(pData->stats, (uchar *)pData->statsName));
+ CHKiRet(statsobj.SetOrigin(pData->stats, (uchar *)"omazureeventhubs"));
+
+ /* Track following stats */
+ STATSCOUNTER_INIT(pData->ctrMessageSubmit, pData->mutCtrMessageSubmit);
+ CHKiRet(statsobj.AddCounter(pData->stats, (uchar *)"submitted",
+ ctrType_IntCtr, CTR_FLAG_RESETTABLE, &pData->ctrMessageSubmit));
+ STATSCOUNTER_INIT(pData->ctrAzureFail, pData->mutCtrAzureFail);
+ CHKiRet(statsobj.AddCounter(pData->stats, (uchar *)"failures",
+ ctrType_IntCtr, CTR_FLAG_RESETTABLE, &pData->ctrAzureFail));
+ STATSCOUNTER_INIT(pData->ctrAzureAck, pData->mutCtrAzureAck);
+ CHKiRet(statsobj.AddCounter(pData->stats, (uchar *)"accepted",
+ ctrType_IntCtr, CTR_FLAG_RESETTABLE, &pData->ctrAzureAck));
+ STATSCOUNTER_INIT(pData->ctrAzureOtherErrors, pData->mutCtrAzureOtherErrors);
+ CHKiRet(statsobj.AddCounter(pData->stats, (uchar *)"othererrors",
+ ctrType_IntCtr, CTR_FLAG_RESETTABLE, &pData->ctrAzureOtherErrors));
+ CHKiRet(statsobj.ConstructFinalize(pData->stats));
+ }
+
+CODE_STD_FINALIZERnewActInst
+ cnfparamvalsDestruct(pvals, &actpblk);
+ENDnewActInst
+
+
+BEGINmodExit
+CODESTARTmodExit
+ DBGPRINTF("modExit: ENTER\n");
+ statsobj.Destruct(&azureStats);
+ CHKiRet(objRelease(statsobj, CORE_COMPONENT));
+ DESTROY_ATOMIC_HELPER_MUT(mutClock);
+
+ objRelease(glbl, CORE_COMPONENT);
+finalize_it:
+ENDmodExit
+
+NO_LEGACY_CONF_parseSelectorAct
+BEGINqueryEtryPt
+CODESTARTqueryEtryPt
+CODEqueryEtryPt_STD_OMODTX_QUERIES
+CODEqueryEtryPt_STD_OMOD8_QUERIES
+CODEqueryEtryPt_STD_CONF2_CNFNAME_QUERIES
+CODEqueryEtryPt_STD_CONF2_OMOD_QUERIES
+ENDqueryEtryPt
+
+
+BEGINmodInit()
+CODESTARTmodInit
+// uchar *pTmp;
+ DBGPRINTF("modInit: ENTER\n");
+INITLegCnfVars
+ *ipIFVersProvided = CURR_MOD_IF_VERSION;
+CODEmodInit_QueryRegCFSLineHdlr
+ /* request objects we use */
+ CHKiRet(objUse(glbl, CORE_COMPONENT));
+ CHKiRet(objUse(datetime, CORE_COMPONENT));
+ CHKiRet(objUse(strm, CORE_COMPONENT));
+ CHKiRet(objUse(statsobj, CORE_COMPONENT));
+
+// INIT_ATOMIC_HELPER_MUT(mutClock);
+ DBGPRINTF("omazureeventhubs %s using qpid-proton library %d.%d.%d\n",
+ VERSION, PN_VERSION_MAJOR, PN_VERSION_MINOR, PN_VERSION_POINT);
+
+ CHKiRet(statsobj.Construct(&azureStats));
+ CHKiRet(statsobj.SetName(azureStats, (uchar *)"omazureeventhubs"));
+ CHKiRet(statsobj.SetOrigin(azureStats, (uchar*)"omazureeventhubs"));
+ STATSCOUNTER_INIT(ctrMessageSubmit, mutCtrMessageSubmit);
+ CHKiRet(statsobj.AddCounter(azureStats, (uchar *)"submitted",
+ ctrType_IntCtr, CTR_FLAG_RESETTABLE, &ctrMessageSubmit));
+ STATSCOUNTER_INIT(ctrAzureFail, mutCtrAzureFail);
+ CHKiRet(statsobj.AddCounter(azureStats, (uchar *)"failures",
+ ctrType_IntCtr, CTR_FLAG_RESETTABLE, &ctrAzureFail));
+ STATSCOUNTER_INIT(ctrAzureAck, mutCtrAzureAck);
+ CHKiRet(statsobj.AddCounter(azureStats, (uchar *)"accepted",
+ ctrType_IntCtr, CTR_FLAG_RESETTABLE, &ctrAzureAck));
+ STATSCOUNTER_INIT(ctrAzureOtherErrors, mutCtrAzureOtherErrors);
+ CHKiRet(statsobj.AddCounter(azureStats, (uchar *)"failures_other",
+ ctrType_IntCtr, CTR_FLAG_RESETTABLE, &ctrAzureOtherErrors));
+ CHKiRet(statsobj.ConstructFinalize(azureStats));
+ENDmodInit
+
+pn_timestamp_t time_now(void)
+{
+ struct timeval now;
+ if (gettimeofday(&now, NULL)) return 0;
+ return ((pn_timestamp_t)now.tv_sec) * 1000 + (now.tv_usec / 1000);
+}
+
+/*
+* Start PROTON Handling Thread
+*/
+static rsRetVal proton_run_thread(wrkrInstanceData_t *pWrkrData)
+{
+ DEFiRet;
+ int iErr = 0;
+ if ( !pWrkrData->bThreadRunning) {
+ DBGPRINTF("omazureeventhubs[%p]: proton_run_thread\n", pWrkrData);
+
+ do {
+ iErr = pthread_create(&pWrkrData->tid,
+ NULL,
+ proton_thread,
+ pWrkrData);
+ if (!iErr) {
+ pWrkrData->bThreadRunning = 1;
+ DBGPRINTF("omazureeventhubs[%p]: proton_run_thread (tid %lx) created\n",
+ pWrkrData, pWrkrData->tid);
+ FINALIZE;
+ }
+ } while (iErr == EAGAIN);
+ } else {
+ DBGPRINTF("omazureeventhubs[%p]: proton_run_thread (tid %ld) already running\n",
+ pWrkrData, pWrkrData->tid);
+ }
+finalize_it:
+ if (iRet != RS_RET_OK) {
+ LogError(0, RS_RET_SYS_ERR, "omazureeventhubs: proton_run_thread thread create failed with error: %d",
+ iErr);
+ }
+ RETiRet;
+}
+/* Stop PROTON Handling Thread
+*/
+static rsRetVal
+proton_shutdown_thread(wrkrInstanceData_t *pWrkrData)
+{
+ DEFiRet;
+ if ( pWrkrData->bThreadRunning) {
+ DBGPRINTF("omazureeventhubs[%p]: STOPPING Thread\n", pWrkrData);
+ int r = pthread_cancel(pWrkrData->tid);
+ if(r == 0) {
+ pthread_join(pWrkrData->tid, NULL);
+ }
+ DBGPRINTF("omazureeventhubs[%p]: STOPPED Thread\n", pWrkrData);
+ pWrkrData->bThreadRunning = 0;
+ }
+ FINALIZE;
+finalize_it:
+ RETiRet;
+}
+
+/*
+* Workerthread function for a single ProActor Handler
+ */
+static void *
+proton_thread(void __attribute__((unused)) *pVoidWrkrData)
+{
+ wrkrInstanceData_t *const pWrkrData = (wrkrInstanceData_t *const) pVoidWrkrData;
+ instanceData *const pData = (instanceData *const) pWrkrData->pData;
+
+ DBGPRINTF("proton_thread[%p]: started protocol workerthread(%lx) for %s:%s/%s\n",
+ pWrkrData, pthread_self(), pData->azurehost, pData->azureport, pData->container);
+
+ do {
+ if ( pWrkrData->pnProactor != NULL) {
+ // Process Protocol events
+ pn_event_batch_t *events = pn_proactor_wait(pWrkrData->pnProactor);
+ pn_event_t *event;
+ while ((event = pn_event_batch_next(events))) {
+ handleProton(pWrkrData, event);
+ }
+ pn_proactor_done(pWrkrData->pnProactor, events);
+ } else {
+ // Wait 10 microseconds
+ srSleep(0, 10000);
+ }
+ } while(glbl.GetGlobalInputTermState() == 0);
+
+ DBGPRINTF("proton_thread[%p]: stopped protocol workerthread\n", pWrkrData);
+ return NULL;
+}
+
+static void
+handleProtonDelivery(wrkrInstanceData_t *const pWrkrData) {
+ instanceData *const pData = (instanceData *const) pWrkrData->pData;
+
+ /* Process messages from ARRAY */
+ for(unsigned int i = 0 ; i < pWrkrData->nProtonMsgs ; ++i) {
+ protonmsg_entry* pMsgEntry = (protonmsg_entry*) pWrkrData->aProtonMsgs[i];
+ // Process Unsubmitted messages only
+ if (pMsgEntry != NULL) {
+ if (pMsgEntry->status == PROTON_UNSUBMITTED) {
+ int iCreditBalance = pn_link_credit(pWrkrData->pnSender);
+ if (iCreditBalance > 0) {
+ DBGPRINTF(
+ "handleProtonDelivery: PN_LINK_FLOW deliver '%s' @ %p:%s:%s/%s, msg:'%.*s'\n",
+ pMsgEntry->MsgID,
+ pWrkrData,
+ pData->azurehost,
+ pData->azureport,
+ pData->container,
+ (pMsgEntry->payload_len > 0 ?
+ (int)pMsgEntry->payload_len-1 : 0),
+ pMsgEntry->payload);
+
+ /* Use sent counter as unique delivery tag. */
+ pn_delivery(pWrkrData->pnSender, pn_dtag((const char *)pMsgEntry->MsgID,
+ pMsgEntry->MsgID_len));
+
+ /* Construct a message with the map */
+ pn_message_t* message = proton_encode_message(pWrkrData, pMsgEntry);
+ /*
+ * pn_message_send does performs the following steps:
+ * - call pn_message_encode2() to encode the message to a buffer
+ * - call pn_link_send() to send the encoded message bytes
+ * - call pn_link_advance() to indicate the message is complete
+ */
+ if (pn_message_send( message,
+ pWrkrData->pnSender,
+ &pWrkrData->pnMessageBuffer) < 0) {
+ LogMsg(0, NO_ERRCODE, LOG_INFO,
+ "handleProtonDelivery: PN_LINK_FLOW deliver SEND ERROR %s\n",
+ pn_error_text(pn_message_error(message)));
+ pn_message_free(message);
+ break;
+ } else {
+ DBGPRINTF("handleProtonDelivery: PN_LINK_FLOW deliver SUCCESS\n");
+ pn_message_free(message);
+ }
+
+ STATSCOUNTER_INC(ctrMessageSubmit, mutCtrMessageSubmit);
+ INST_STATSCOUNTER_INC( pData,
+ pData->ctrMessageSubmit,
+ pData->mutCtrMessageSubmit);
+ pMsgEntry->status = PROTON_SUBMITTED;
+ } else {
+ DBGPRINTF("handleProtonDelivery: sender credit balance reached %d. "
+ "extend credit for %d\n", iCreditBalance, pWrkrData->nProtonMsgs);
+ pn_link_flow(pWrkrData->pnSender, pWrkrData->nProtonMsgs);
+
+ // TODO: MAKE CONFIGUREABLE
+ // Wait 10 microseconds
+ // srSleep(0, 10000);
+ break;
+ }
+ }
+ } else {
+ // Abort further processing if pMsgEntry is NULL
+ break;
+ }
+ }
+}
+
+/* Handles PROTON Communication in this Function
+*
+*/
+#pragma GCC diagnostic ignored "-Wswitch"
+static void
+handleProton(wrkrInstanceData_t *const pWrkrData, pn_event_t *event) {
+ instanceData *const pData = (instanceData *const) pWrkrData->pData;
+ /* Handle Proton Events */
+ switch (pn_event_type(event)) {
+ case PN_CONNECTION_BOUND: {
+ DBGPRINTF("handleProton: PN_CONNECTION_BOUND to %p:%s:%s/%s\n",
+ pWrkrData, pData->azurehost, pData->azureport, pData->container);
+ break;
+ }
+ case PN_SESSION_INIT : {
+ DBGPRINTF("handleProton: PN_SESSION_INIT to %p:%s:%s/%s\n",
+ pWrkrData, pData->azurehost, pData->azureport, pData->container);
+ break;
+ }
+ case PN_LINK_INIT: {
+ DBGPRINTF("handleProton: PN_LINK_INIT to %p:%s:%s/%s\n",
+ pWrkrData, pData->azurehost, pData->azureport, pData->container);
+ break;
+ }
+ case PN_LINK_REMOTE_OPEN: {
+ DBGPRINTF("handleProton: PN_LINK_REMOTE_OPEN to %p:%s:%s/%s\n",
+ pWrkrData, pData->azurehost, pData->azureport, pData->container);
+ pWrkrData->bIsConnected = 1;
+ pWrkrData->bIsConnecting = 0;
+ break;
+ }
+ case PN_CONNECTION_INIT: {
+ DBGPRINTF("handleProton: PN_CONNECTION_INIT to %p:%s:%s/%s\n",
+ pWrkrData, pData->azurehost, pData->azureport, pData->container);
+ pWrkrData->pnStatus = PN_CONNECTION_INIT;
+
+ // Get Connection
+ pWrkrData->pnConn = pn_event_connection(event);
+ // Set AMQP Properties
+ pn_connection_set_container(pWrkrData->pnConn, (const char *) pData->container);
+ pn_connection_set_hostname(pWrkrData->pnConn, (const char *) pData->azurehost);
+ pn_connection_set_user(pWrkrData->pnConn, (const char *) pData->azure_key_name);
+ pn_connection_set_password(pWrkrData->pnConn, (const char *) pData->azure_key);
+
+ pn_connection_open(pWrkrData->pnConn); // Open Connection
+ pn_session_t* pnSession = pn_session(pWrkrData->pnConn); // Create Session
+ pn_session_open(pnSession); // Open Session
+ pWrkrData->pnSender = pn_sender(pnSession, (char *)pData->azure_key_name); // Create Link
+
+ DBGPRINTF("handleProton: PN_CONNECTION_INIT with amqp address: %s\n",
+ pData->amqp_address);
+ pn_terminus_set_address(pn_link_target(pWrkrData->pnSender),
+ (const char *) pData->amqp_address);
+
+ pn_link_open(pWrkrData->pnSender);
+ pn_link_flow(pWrkrData->pnSender, pWrkrData->nProtonMsgs);
+ break;
+ }
+ case PN_CONNECTION_REMOTE_OPEN: {
+ DBGPRINTF("handleProton: PN_CONNECTION_REMOTE_OPEN to %p:%s:%s/%s\n",
+ pWrkrData, pData->azurehost, pData->azureport, pData->container);
+ pWrkrData->pnStatus = PN_CONNECTION_REMOTE_OPEN;
+ pn_ssl_t *ssl = pn_ssl(pn_event_transport(event));
+ if (ssl) {
+ char name[1024];
+ pn_ssl_get_protocol_name(ssl, name, sizeof(name));
+ {
+ const char *subject = pn_ssl_get_remote_subject(ssl);
+ if (subject) {
+ DBGPRINTF(
+ "handleProton: handleProton secure connection: to %s using %s\n",
+ subject, name);
+ } else {
+ DBGPRINTF(
+ "handleProton: handleProton anonymous connection: using %s\n",
+ name);
+ }
+ fflush(stdout);
+ }
+ }
+ break;
+ }
+ case PN_CONNECTION_WAKE: {
+ DBGPRINTF("handleProton: PN_CONNECTION_WAKE (%d) to %p:%s:%s/%s\n",
+ pWrkrData->nProtonMsgs,
+ pWrkrData, pData->azurehost, pData->azureport, pData->container);
+ /* Process messages */
+ handleProtonDelivery(pWrkrData);
+ break;
+ }
+ case PN_LINK_FLOW: {
+ pWrkrData->pnStatus = PN_LINK_FLOW;
+ /* The peer has given us some credit, now we can send messages */
+ pWrkrData->pnSender = pn_event_link(event);
+
+ /* Process messages */
+ handleProtonDelivery(pWrkrData);
+ break;
+ }
+ case PN_DELIVERY: {
+ pWrkrData->pnStatus = PN_DELIVERY;
+ /* We received acknowledgement from the peer that a message was delivered. */
+ pn_delivery_t* pDeliveryStatus = pn_event_delivery(event);
+ pn_delivery_tag_t pnTag = pn_delivery_tag(pDeliveryStatus);
+
+ if (pnTag.start != NULL) {
+ // Convert Tag into Number!
+ unsigned int iTagNum = (unsigned int) atoi((pnTag.start != NULL ? pnTag.start : ""));
+ // Calc QueueNumber from Tagnum
+ unsigned int iQueueNum = pWrkrData->nProtonMsgs - (pWrkrData->iMaxMsgSeq - iTagNum);
+ // Get proton Msg Helper (checks for out of bound array access)
+ protonmsg_entry* pMsgEntry = NULL;
+ if (pWrkrData->nMaxProtonMsgs > iQueueNum) {
+ pMsgEntry = (protonmsg_entry*) pWrkrData->aProtonMsgs[iQueueNum];
+ }
+ // Process if found
+ if (pMsgEntry != NULL) {
+ if (pn_delivery_remote_state(pDeliveryStatus) == PN_ACCEPTED) {
+ DBGPRINTF(
+ "handleProton: PN_DELIVERY SUCCESS for MSG '%s(Q %d, MAX %d)' @ %p:%s:%s/%s\n",
+ (pnTag.start != NULL ? (char*) pnTag.start : "NULL"),
+ iQueueNum,
+ pWrkrData->nMaxProtonMsgs,
+ pWrkrData, pData->azurehost, pData->azureport, pData->container);
+ pMsgEntry->status = PROTON_ACCEPTED;
+
+ // Increment Stats Counter
+ STATSCOUNTER_INC(ctrAzureAck, mutCtrAzureAck);
+ INST_STATSCOUNTER_INC(pData, pData->ctrAzureAck, pData->mutCtrAzureAck);
+ } else if (pn_delivery_remote_state(pDeliveryStatus) == PN_REJECTED) {
+ LogError(0, RS_RET_ERR,
+ "omazureeventhubs: PN_DELIVERY REJECTED for MSG '%s'"
+ " - @ %p:%s:%s/%s\n",
+ (pnTag.start != NULL ? (char*) pnTag.start : "NULL"),
+ pWrkrData, pData->azurehost, pData->azureport, pData->container);
+ pMsgEntry->status = PROTON_REJECTED;
+
+ // Increment Stats Counter
+ STATSCOUNTER_INC(ctrAzureFail, mutCtrAzureFail);
+ INST_STATSCOUNTER_INC(pData, pData->ctrAzureFail, pData->mutCtrAzureFail);
+ }
+ } else {
+ DBGPRINTF("handleProton: PN_DELIVERY MISSING MSG '%d(Q %d,MAX %d)' - @ %p:%s:%s/%s\n",
+ iTagNum,
+ iQueueNum,
+ pWrkrData->nMaxProtonMsgs,
+ pWrkrData, pData->azurehost, pData->azureport, pData->container);
+ STATSCOUNTER_INC(ctrAzureOtherErrors, mutCtrAzureOtherErrors);
+ INST_STATSCOUNTER_INC(pData, pData->ctrAzureOtherErrors, pData->mutCtrAzureOtherErrors);
+ }
+ } else {
+ LogError(0, RS_RET_ERR,"handleProton: PN_DELIVERY HELPER ARRAY is NULL - @ %p:%s:%s/%s\n",
+ pWrkrData, pData->azurehost, pData->azureport, pData->container);
+ STATSCOUNTER_INC(ctrAzureOtherErrors, mutCtrAzureOtherErrors);
+ INST_STATSCOUNTER_INC(pData, pData->ctrAzureOtherErrors, pData->mutCtrAzureOtherErrors);
+ }
+ break;
+ }
+ case PN_TRANSPORT_CLOSED:
+ DBGPRINTF("handleProton: transport closed for %p:%s\n", pWrkrData, pData->azurehost);
+ proton_check_condition(event, pWrkrData, pn_transport_condition(pn_event_transport(event)),
+ "transport closed");
+ break;
+ case PN_CONNECTION_REMOTE_CLOSE:
+ DBGPRINTF("handleProton: connection closed for %p:%s\n", pWrkrData, pData->azurehost);
+ proton_check_condition(event, pWrkrData, pn_connection_remote_condition(pn_event_connection(event)),
+ "connection closed");
+ break;
+ case PN_SESSION_REMOTE_CLOSE:
+ DBGPRINTF("handleProton: remote session closed for %p:%s\n", pWrkrData, pData->azurehost);
+ proton_check_condition(event, pWrkrData, pn_session_remote_condition(pn_event_session(event)),
+ "remote session closed");
+ break;
+ case PN_LINK_REMOTE_CLOSE:
+ case PN_LINK_REMOTE_DETACH:
+ DBGPRINTF("handleProton: remote link closed for %p:%s\n", pWrkrData, pData->azurehost);
+ proton_check_condition(event, pWrkrData, pn_link_remote_condition(pn_event_link(event)),
+ "remote link closed");
+ break;
+ case PN_PROACTOR_INACTIVE:
+ DBGPRINTF("handleProton: INAKTIVE for %p:%s\n", pWrkrData, pData->azurehost);
+ break;
+/* Workarround compiler warning:
+* error: enumeration value '...' not handled in switch [-Werror=switch-enum]
+*/
+#ifdef __GNU_C
+ default:
+ DBGPRINTF("handleProton: UNHANDELED EVENT %d for %p:%s\n", pn_event_type(event),
+ pWrkrData, pData->azurehost);
+ break;
+#endif
+ }
+}