diff options
author | Daniel Baumann <daniel.baumann@progress-linux.org> | 2024-04-15 16:28:20 +0000 |
---|---|---|
committer | Daniel Baumann <daniel.baumann@progress-linux.org> | 2024-04-15 16:28:20 +0000 |
commit | dcc721a95bef6f0d8e6d8775b8efe33e5aecd562 (patch) | |
tree | 66a2774cd0ee294d019efd71d2544c70f42b2842 /plugins/omazureeventhubs | |
parent | Initial commit. (diff) | |
download | rsyslog-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 'plugins/omazureeventhubs')
-rw-r--r-- | plugins/omazureeventhubs/Makefile.am | 13 | ||||
-rw-r--r-- | plugins/omazureeventhubs/Makefile.in | 802 | ||||
-rw-r--r-- | plugins/omazureeventhubs/omazureeventhubs.c | 1354 |
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 + } +} |