16#include "git-revision.h"
22#include <sys/resource.h>
62#ifndef DOXYGEN_SHOULD_SKIP_THIS
161extern unsigned _stklen = 60000U;
274 "\"exptab\" file not found and MIDAS_DIR or MIDAS_EXPTAB environment variable is not defined"},
333 f = fopen(
"mem.txt",
"w");
351 memset(adr, 0, size *
count);
369 f = fopen(
"mem.txt",
"w");
381static std::vector<std::string>
split(
const char* sep,
const std::string& s)
383 unsigned sep_len = strlen(sep);
384 std::vector<std::string> v;
385 std::string::size_type pos = 0;
387 std::string::size_type next = s.find(sep, pos);
388 if (next == std::string::npos) {
389 v.push_back(s.substr(pos));
392 v.push_back(s.substr(pos, next-pos));
398static std::string
join(
const char* sep,
const std::vector<std::string>& v)
402 for (
unsigned i=0;
i<v.size();
i++) {
416 return s[s.length()-1] ==
c;
420 assert(format != NULL);
422 va_start(ap, format);
424 int length = vsnprintf(
nullptr, 0, format, ap1);
430 size_t size = (size_t) length + 1;
431 char *buffer = (
char *)malloc(size);
436 vsnprintf(buffer, size, format, ap);
438 std::string s(buffer);
481 return msprintf(
"unlisted status code %d", code);
529 if (pos != std::string::npos) {
541 for (
size_t i = 0;
i < flist.size();
i++) {
542 const char *p = flist[
i].c_str();
543 if (strchr(p,
'_') == NULL && !(p[0] >=
'0' && p[0] <=
'9')) {
544 size_t pos = flist[
i].rfind(
'.');
545 if (pos != std::string::npos) {
546 flist[
i].resize(pos);
548 list->push_back(flist[
i]);
557void cm_msg_get_logfile(
const char *
fac, time_t t, std::string* filename, std::string* linkname, std::string* linktarget) {
567 *filename = std::string(
fac) +
".log";
582 std::string facility;
588 std::string message_format;
590 if (message_format.find(
'%') != std::string::npos) {
597 localtime_r(&t, &tms);
601 strftime(
de + 1,
sizeof(
de)-1, strchr(message_format.c_str(),
'%'), &tms);
605 std::string message_dir;
607 if (message_dir.empty()) {
609 if (message_dir.empty()) {
611 if (message_dir.empty()) {
625 *filename = message_dir + facility + message_format +
".log";
626 if (!message_format.empty()) {
628 *linkname = message_dir + facility +
".log";
630 *linktarget = facility + message_format +
".log";
689 fprintf(stderr,
"cm_msg_log: Message \"%s\" not written to midas.log because rpc_call(RPC_CM_MSG_LOG) failed with status %d\n",
message,
status);
693 fprintf(stderr,
"cm_msg_log: Message \"%s\" not written to midas.log, no connection to mserver\n",
message);
699 std::string filename, linkname, linktarget;
704 if (!linkname.empty()) {
712 unlink(linkname.c_str());
713 status = symlink(linktarget.c_str(), linkname.c_str());
716 "cm_msg_log: Error: Cannot symlink message log file \'%s' to \'%s\', symlink() errno: %d (%s)\n",
717 linktarget.c_str(), linkname.c_str(), errno, strerror(errno));
722 int fh = open(filename.c_str(), O_WRONLY | O_CREAT | O_APPEND |
O_LARGEFILE, 0644);
725 "cm_msg_log: Message \"%s\" not written to midas.log because open(%s) failed with errno %d (%s)\n",
726 message, filename.c_str(), errno, strerror(errno));
734 localtime_r(&
tv.tv_sec, &tms);
737 strftime(
str,
sizeof(
str),
"%H:%M:%S", &tms);
738 sprintf(
str + strlen(
str),
".%03d ", (
int) (
tv.tv_usec / 1000));
739 strftime(
str + strlen(
str),
sizeof(
str),
"%Y/%m/%d", &tms);
750 ssize_t len = msg.length();
753 ssize_t wr =
write(fh, msg.c_str(), len);
756 fprintf(stderr,
"cm_msg_log: Message \"%s\" not written to \"%s\", write() error, errno %d (%s)\n",
message, filename.c_str(), errno, strerror(errno));
757 }
else if (wr != len) {
758 fprintf(stderr,
"cm_msg_log: Message \"%s\" not written to \"%s\", short write() wrote %d instead of %d bytes\n",
message, filename.c_str(), (
int)wr, (
int)len);
769static std::string
cm_msg_format(
INT message_type,
const char *filename,
INT line,
const char *routine,
const char *format, va_list *argptr)
772 const char* pc = filename + strlen(filename);
773 while (*pc !=
'\\' && *pc !=
'/' && pc != filename)
779 std::string type_str;
788 if (message_type &
MT_LOG)
800 if (
name.length() > 0)
808 message +=
msprintf(
"[%s:%d:%s,%s] ", pc, line, routine, type_str.c_str());
809 }
else if (message_type ==
MT_USER) {
814 char* buf = (
char*)malloc(bufsize);
817 for (
int i=0;
i<10;
i++) {
819 va_copy(ap, *argptr);
822 int n = vsnprintf(buf, bufsize, format, ap);
838 buf = (
char*)realloc(buf, bufsize);
852 if (message_type !=
MT_LOG) {
855 size_t len = strlen(send_message);
857 char event[event_length];
860 memcpy(event +
sizeof(
EVENT_HEADER), send_message, len + 1);
894 for (
i = 0;
i < 100;
i++) {
939INT cm_msg(
INT message_type,
const char *filename,
INT line,
const char *routine,
const char *format, ...)
946 va_start(argptr, format);
955 if (message_type !=
MT_LOG) {
998 const char *facility,
const char *routine,
const char *format, ...) {
1010 va_start(argptr, format);
1089static void add_message(
char **messages,
int *length,
int *allocated, time_t tstamp,
const char *new_message) {
1090 int new_message_length = strlen(new_message);
1091 int new_allocated = 1024 + 2 * ((*allocated) + new_message_length);
1097 if (*length + new_message_length + 100 > *allocated) {
1098 *messages = (
char *) realloc(*messages, new_allocated);
1099 assert(*messages != NULL);
1100 *allocated = new_allocated;
1104 if ((*messages)[(*length) - 1] !=
'\n') {
1105 (*messages)[*length] =
'\n';
1109 sprintf(buf,
"%ld ", tstamp);
1110 buf_length = strlen(buf);
1111 memcpy(&((*messages)[*length]), buf, buf_length);
1112 (*length) += buf_length;
1114 memcpy(&((*messages)[*length]), new_message, new_message_length);
1115 (*length) += new_message_length;
1116 (*messages)[*length] = 0;
1120static int cm_msg_retrieve1(
const char *filename, time_t t,
INT n_messages,
char **messages,
int *length,
int *allocated,
1121 int *num_messages) {
1125 struct stat stat_buf;
1126 time_t tstamp, tstamp_valid, tstamp_last;
1132 fh = open(filename, O_RDONLY |
O_TEXT, 0644);
1134 cm_msg(
MERROR,
"cm_msg_retrieve1",
"Cannot open log file \"%s\", errno %d (%s)", filename, errno,
1140 if (fstat(fh, &stat_buf) != 0) {
1141 cm_msg(
MERROR,
"cm_msg_retrieve1",
"Cannot stat log file \"%s\", errno %d (%s)", filename, errno,
1146 ssize_t size = stat_buf.st_size;
1149 ssize_t maxsize = 10 * 1024 * 1024;
1150 if (size > maxsize) {
1151 lseek(fh, -maxsize, SEEK_END);
1156 char *buffer = (
char *) malloc(size + 1);
1158 if (buffer == NULL) {
1159 cm_msg(
MERROR,
"cm_msg_retrieve1",
"Cannot malloc %d bytes to read log file \"%s\", errno %d (%s)", (
int) size,
1160 filename, errno, strerror(errno));
1165 ssize_t rd =
read(fh, buffer, size);
1168 cm_msg(
MERROR,
"cm_msg_retrieve1",
"Cannot read %d bytes from log file \"%s\", read() returned %d, errno %d (%s)",
1169 (
int) size, filename, (
int) rd, errno, strerror(errno));
1178 tstamp_last = tstamp_valid = 0;
1186 p = buffer + size - 1;
1188 while (p > buffer && (*p ==
'\n' || *p ==
'\r'))
1191 if (p == buffer && (*p ==
'\n' || *p ==
'\r')) {
1197 for (
n = 0; !stop;) {
1201 for (
i = 0; p != buffer && (*p !=
'\n' && *p !=
'\r');
i++)
1205 if (
i >= (
int)
sizeof(
str))
1206 i =
sizeof(
str) - 1;
1212 memcpy(
str, p + 1,
i);
1214 if (strchr(
str,
'\n'))
1215 *strchr(
str,
'\n') = 0;
1216 if (strchr(
str,
'\r'))
1217 *strchr(
str,
'\r') = 0;
1218 mstrlcat(
str,
"\n",
sizeof(
str));
1225 localtime_r(&now, &tms);
1227 if (
str[0] >=
'0' &&
str[0] <=
'9') {
1229 tms.tm_hour = atoi(
str);
1230 tms.tm_min = atoi(
str + 3);
1231 tms.tm_sec = atoi(
str + 6);
1232 tms.tm_year = atoi(
str + 13) - 1900;
1233 tms.tm_mon = atoi(
str + 18) - 1;
1234 tms.tm_mday = atoi(
str + 21);
1237 tms.tm_hour = atoi(
str + 11);
1238 tms.tm_min = atoi(
str + 14);
1239 tms.tm_sec = atoi(
str + 17);
1240 tms.tm_year = atoi(
str + 20) - 1900;
1241 for (
i = 0;
i < 12;
i++)
1245 tms.tm_mday = atoi(
str + 8);
1249 tstamp_valid = tstamp;
1252 if (n_messages == 0) {
1253 if (tstamp_valid < t)
1258 if (n_messages != 0) {
1259 if (tstamp_last > 0 && tstamp_valid < tstamp_last)
1263 if (t == 0 || tstamp == -1 ||
1264 (n_messages > 0 && tstamp <= t) ||
1265 (n_messages == 0 && tstamp >= t)) {
1275 while (p > buffer && (*p ==
'\n' || *p ==
'\r'))
1278 if (p == buffer && (*p ==
'\n' || *p ==
'\r'))
1281 if (n_messages == 1)
1283 else if (n_messages > 1) {
1285 if (
n == n_messages)
1286 tstamp_last = tstamp_valid;
1289 if (
n == n_messages && tstamp_valid == 0)
1312 std::string filename, linkname;
1324 if (!linkname.empty()) {
1326 filename = linkname;
1330 cm_msg_retrieve1(filename.c_str(), t, n_message, messages, &length, &allocated, &
n);
1336 if (linkname.empty()) {
1344 while (
n < n_message) {
1345 filedate -= 3600 * 24;
1352 cm_msg_retrieve1(filename.c_str(), t, n_message -
n, messages, &length, &allocated, &
i);
1383 char *messages = NULL;
1384 int num_messages = 0;
1395 mstrlcpy(
message, messages, buf_size);
1396 int len = strlen(messages);
1397 if (len >= buf_size)
1432 if (seconds != NULL) {
1535 return GIT_REVISION;
1549 assert(path[0] != 0);
1572 assert(path_size !=
sizeof(
char *));
1576 mstrlcpy(path,
_path_name.c_str(), path_size);
1596 assert(path != NULL);
1642#ifdef LOCAL_ROUTINES
1669 if (getenv(
"MIDAS_DIR")) {
1674 if (getenv(
"MIDAS_EXPT_NAME")) {
1675 e.
name = getenv(
"MIDAS_EXPT_NAME");
1678 cm_msg(
MERROR,
"cm_read_exptab",
"Experiments that use MIDAS_DIR must also set MIDAS_EXPT_NAME to the name of the experiment! Using experiment name \"%s\"",
e.name.c_str());
1681 e.directory = getenv(
"MIDAS_DIR");
1690#if defined (OS_WINNT)
1692 if (getenv(
"SystemRoot"))
1693 str = getenv(
"SystemRoot");
1694 else if (getenv(
"windir"))
1695 str = getenv(
"windir");
1699 std::string alt_str =
str;
1700 str +=
"\\system32\\exptab";
1701 alt_str +=
"\\system\\exptab";
1702#elif defined (OS_UNIX)
1703 std::string
str =
"/etc/exptab";
1704 std::string alt_str =
"/exptab";
1706 std::strint
str =
"exptab";
1707 std::string alt_str =
"exptab";
1711 if (getenv(
"MIDAS_EXPTAB")) {
1712 str = getenv(
"MIDAS_EXPTAB");
1713 alt_str = getenv(
"MIDAS_EXPTAB");
1719 FILE* f = fopen(
str.c_str(),
"r");
1721 f = fopen(alt_str.c_str(),
"r");
1730 memset(buf, 0,
sizeof(buf));
1731 char*
str = fgets(buf,
sizeof(buf)-1, f);
1734 if (
str[0] == 0)
continue;
1735 if (
str[0] ==
'#')
continue;
1743 while (*
str && isspace(*
str))
1749 while (*p2 && !isspace(*p2))
1752 ssize_t len = p2-p1;
1759 e.
name = std::string(p1, len);
1767 while (*
str && isspace(*
str))
1773 while (*p2 && !isspace(*p2))
1783 e.directory = std::string(p1, len);
1791 while (*
str && isspace(*
str))
1797 while (*p2 && !isspace(*p2))
1804 e.user = std::string(p1, len);
1818 for (
unsigned j=0;
j<exptab->
exptab.size();
j++) {
1819 cm_msg(
MINFO,
"cm_read_exptab",
"entry %d, experiment \"%s\", directory \"%s\", user \"%s\"",
j, exptab->
exptab[
j].name.c_str(), exptab->
exptab[
j].directory.c_str(), exptab->
exptab[
j].user.c_str());
1880int cm_get_exptab(
const char *expname,
char *dir,
int dir_size,
char *user,
int user_size) {
1881 std::string sdir, suser;
1885 mstrlcpy(dir, sdir.c_str(), dir_size);
1887 mstrlcpy(user, suser.c_str(), user_size);
1923#ifdef LOCAL_ROUTINES
1944 const char *program_name,
INT hw_type,
const char *password,
DWORD watchdog_timeout) {
1947 host_name, program_name, hw_type, password, watchdog_timeout);
1949#ifdef LOCAL_ROUTINES
1954 BOOL call_watchdog, allow;
1956 std::string start_command(255,
'\0');
1957 std::string alarm_class(31,
'\0');
1959 {
"Required",
false},
1960 {
"Watchdog timeout", 10000},
1961 {
"Check interval", (uint32_t)180000},
1962 {
"Start command", start_command},
1963 {
"Auto start",
false},
1964 {
"Auto stop",
false},
1965 {
"Auto restart",
false},
1966 {
"Alarm class", alarm_class},
1967 {
"First failed", (uint32_t)0}
1971 cm_msg(
MINFO,
"cm_set_client_info",
"Client name \"%s\" length %zu is longer than NAME_LENGTH %d", program_name, strlen(program_name),
NAME_LENGTH);
1994 if (!allow && strcmp(password,
pwd) != 0) {
2007 std::string
str =
msprintf(
"System/Clients/%0d", pid);
2017 std::string client_name = program_name;
2039 client_name =
msprintf(
"%s%d", program_name, idx);
2051 cm_msg(
MERROR,
"cm_set_client_info",
"cannot set client name, db_set_value(%s) status %d",
str.c_str(),
status);
2099 size =
sizeof(watchdog_timeout);
2100 str =
msprintf(
"/Programs/%s/Watchdog Timeout", program_name);
2214 if (
host_name && getenv(
"MIDAS_SERVER_HOST"))
2215 mstrlcpy(
host_name, getenv(
"MIDAS_SERVER_HOST"), host_name_size);
2217 if (
exp_name && getenv(
"MIDAS_EXPT_NAME"))
2218 mstrlcpy(
exp_name, getenv(
"MIDAS_EXPT_NAME"), exp_name_size);
2229 if (
host_name && getenv(
"MIDAS_SERVER_HOST"))
2230 *
host_name = getenv(
"MIDAS_SERVER_HOST");
2232 if (
exp_name && getenv(
"MIDAS_EXPT_NAME"))
2233 *
exp_name = getenv(
"MIDAS_EXPT_NAME");
2238#ifdef LOCAL_ROUTINES
2242 std::string exp_name1;
2252 std::string expdir, expuser;
2262 cm_msg(
MERROR,
"cm_set_experiment_local",
"Experiment \"%s\" directory \"%s\" does not exist", exp_name1.c_str(), expdir.c_str());
2277 cm_msg(
MERROR,
"cm_check_connect",
"cm_disconnect_experiment not called at end of program");
2372 const char *client_name,
void (*func)(
char *),
INT odb_size,
DWORD watchdog_timeout) {
2399 if (WSAStartup(MAKEWORD(1, 1), &WSAData) != 0)
2404 std::string default_exp_name1;
2405 if (default_exp_name)
2406 default_exp_name1 = default_exp_name;
2410 if (default_exp_name1.length() == 0) {
2429#ifdef LOCAL_ROUTINES
2438 INT semaphore_elog, semaphore_alarm, semaphore_history, semaphore_msg;
2443 cm_msg(
MERROR,
"cm_connect_experiment",
"Cannot create alarm semaphore");
2448 cm_msg(
MERROR,
"cm_connect_experiment",
"Cannot create elog semaphore");
2453 cm_msg(
MERROR,
"cm_connect_experiment",
"Cannot create history semaphore");
2458 cm_msg(
MERROR,
"cm_connect_experiment",
"Cannot create message semaphore");
2477 cm_msg(
MERROR,
"cm_connect_experiment1",
"cannot open database, db_open_database() status %d",
status);
2485 size =
sizeof(odb_timeout);
2488 cm_msg(
MERROR,
"cm_connect_experiment1",
"cannot get ODB /Experiment/ODB timeout, status %d",
status);
2491 if (odb_timeout > 0) {
2496 size =
sizeof(protect_odb);
2499 cm_msg(
MERROR,
"cm_connect_experiment1",
"cannot get ODB /Experiment/Protect ODB, status %d",
status);
2507 size =
sizeof(enable_core_dumps);
2510 cm_msg(
MERROR,
"cm_connect_experiment1",
"cannot get ODB /Experiment/Enable core dumps, status %d",
status);
2513 if (enable_core_dumps) {
2515 struct rlimit limit;
2516 limit.rlim_cur = RLIM_INFINITY;
2517 limit.rlim_max = RLIM_INFINITY;
2518 status = setrlimit(RLIMIT_CORE, &limit);
2520 cm_msg(
MERROR,
"cm_connect_experiment",
"Cannot setrlimit(RLIMIT_CORE, RLIM_INFINITY), errno %d (%s)", errno,
2524#warning setrlimit(RLIMIT_CORE) is not available
2533 "cannot get ODB /Experiment/Security/Enable non-localhost RPC, status %d",
status);
2536 std::string local_host_name;
2540 local_host_name =
"localhost";
2545 if (watchdog_timeout == 0)
2584 cm_msg(
MERROR,
"cm_connect_experiment1",
"cannot open message buffer, cm_msg_open_buffer() status %d",
status);
2592 std::string current_name;
2594 if (current_name.length() == 0 || current_name ==
"Default") {
2608 cm_msg(
MERROR,
"cm_connect_experiment",
"Cannot register RPC server, cm_register_server() status %d",
status);
2616 size =
sizeof(watchdog_timeout);
2617 sprintf(
str,
"/Programs/%s/Watchdog Timeout", client_name);
2623 std::string path =
"/Programs/" + std::string(client_name);
2626 prog[
"Start command"] == std::string(
""))
2633 cm_msg(
MLOG,
"cm_connect_experiment",
"Program %s on host %s started", xclient_name.c_str(), local_host_name.c_str());
2650#ifdef LOCAL_ROUTINES
2659 assert(exp_names != NULL);
2691 assert(exp_names != NULL);
2695 mstrlcpy(hname,
host_name,
sizeof(hname));
2696 s = strchr(hname,
':');
2699 port = strtoul(s + 1, NULL, 0);
2707 cm_msg(
MERROR,
"cm_list_experiments_remote",
"Cannot connect to \"%s\" port %d: %s", hname, port, errmsg.c_str());
2712 send(sock,
"I", 2, 0);
2725 exp_names->push_back(
str);
2733#ifdef LOCAL_ROUTINES
2753 if (expts.size() == 1) {
2755 }
else if (expts.size() > 1) {
2756 printf(
"Available experiments on local computer:\n");
2758 for (
unsigned i = 0;
i < expts.size();
i++) {
2759 printf(
"%d : %s\n",
i, expts[
i].c_str());
2763 printf(
"Select number from 0 to %d: ", ((
int)expts.size())-1);
2766 int isel = atoi(
str);
2769 if (isel >= (
int)expts.size())
2802 if (expts.size() > 1) {
2803 printf(
"Available experiments on server %s:\n",
host_name);
2805 for (
unsigned i = 0;
i < expts.size();
i++) {
2806 printf(
"%d : %s\n",
i, expts[
i].c_str());
2810 printf(
"Select number from 0 to %d: ", ((
int)expts.size())-1);
2813 int isel = atoi(
str);
2816 if (isel >= (
int)expts.size())
2871 length =
sizeof(
INT);
2926 printf(
"Waiting for transition to finish...\n");
2938 std::string local_host_name;
2941 local_host_name =
"localhost";
2949 cm_msg(
MLOG,
"cm_disconnect_experiment",
"Program %s on host %s stopped", client_name.c_str(), local_host_name.c_str());
3027#ifndef DOXYGEN_SHOULD_SKIP_THIS
3088 if (hKeyClient != NULL)
3095 if (hKeyClient != NULL)
3102#ifndef DOXYGEN_SHOULD_SKIP_THIS
3126 if (semaphore_alarm)
3130 if (semaphore_history)
3135 *semaphore_msg = -1;
3143#ifdef LOCAL_ROUTINES
3159 assert(pbuf != NULL);
3163 printf(
"lock_buffer_guard(%s) ctor without lock\n",
fBuf->
buffer_name);
3182 printf(
"lock_buffer_guard(invalid) dtor\n");
3184 assert(
fBuf != NULL);
3202 assert(
fBuf != NULL);
3213 assert(
fBuf != NULL);
3230 assert(
fBuf != NULL);
3318#ifdef LOCAL_ROUTINES
3321 std::vector<BUFFER*> mybuffers;
3328 for (
BUFFER* pbuf : mybuffers) {
3330 if (!pbuf || !pbuf->attached)
3399 *call_watchdog =
FALSE;
3420#ifdef LOCAL_ROUTINES
3429#ifndef DOXYGEN_SHOULD_SKIP_THIS
3452 str = (
char *) malloc(max_size);
3456 int size = max_size;
3461 if (strlen(
str) < 1)
3473 int new_size = last + 10;
3477 "Cannot resize the RPC hosts access control list, db_set_num_values(%d) status %d", new_size,
status);
3490 strcpy(buf,
"localhost");
3496 cm_msg(
MERROR,
"init_rpc_hosts",
"Cannot create the RPC hosts access control list, db_get_value() status %d",
3506 cm_msg(
MERROR,
"init_rpc_hosts",
"Cannot create \"Disable RPC hosts check\", db_get_value() status %d",
status);
3516 cm_msg(
MERROR,
"init_rpc_hosts",
"Cannot find the RPC hosts access control list, db_find_key() status %d",
3526 cm_msg(
MERROR,
"init_rpc_hosts",
"Cannot watch the RPC hosts access control list, db_watch() status %d",
status);
3561 size =
sizeof(
name);
3565 cm_msg(
MERROR,
"cm_register_server",
"cannot get client name, db_get_value() status %d",
status);
3569 mstrlcpy(
str,
"/Experiment/Security/RPC ports/",
sizeof(
str));
3572 size =
sizeof(port);
3576 cm_msg(
MERROR,
"cm_register_server",
"cannot get RPC port number, db_get_value(%s) status %d",
str,
status);
3584 cm_msg(
MERROR,
"cm_register_server",
"error, rpc_register_server(port=%d) status %d", port,
status);
3597 cm_msg(
MERROR,
"cm_register_server",
"error, db_find_key(\"Server Port\") status %d",
status);
3607 cm_msg(
MERROR,
"cm_register_server",
"error, db_set_data(\"Server Port\"=%d) status %d", port,
status);
3828 }
else if (
count > 1) {
3875 "Cannot set client run state, client hKey %d into /System/Clients is not valid, maybe this client was removed by a watchdog timeout",
3896#ifndef DOXYGEN_SHOULD_SKIP_THIS
3919 char tr_key_name[256];
3951 cm_msg(
MERROR,
"cm_register_deferred_transition",
"Cannot hotlink /Runinfo/Requested Transition");
3986 cm_msg(
MERROR,
"cm_check_deferred_transition",
"Cannot perform deferred transition: %s",
str);
4002#ifndef DOXYGEN_SHOULD_SKIP_THIS
4048 printf(
", wait for:");
4056static bool tr_compare(
const std::unique_ptr<TrClient>& arg1,
const std::unique_ptr<TrClient>& arg2) {
4057 return arg1->sequence_number < arg2->sequence_number;
4087 const char *buf =
"Success";
4091 sprintf(buf,
"status %d",
status);
4165 const char *args[100];
4167 char debug_arg[256];
4168 char start_arg[256];
4170 std::string mserver_hostname;
4176 const char *midassys = getenv(
"MIDASSYS");
4183 path +=
"mtransition";
4185 args[iarg++] = path.c_str();
4190 args[iarg++] =
"-h";
4191 args[iarg++] = mserver_hostname.c_str();
4198 args[iarg++] =
"-e";
4203 args[iarg++] =
"-d";
4205 sprintf(debug_arg,
"%d", debug_flag);
4206 args[iarg++] = debug_arg;
4210 args[iarg++] =
"STOP";
4212 args[iarg++] =
"PAUSE";
4214 args[iarg++] =
"RESUME";
4216 args[iarg++] =
"START";
4219 args[iarg++] = start_arg;
4222 args[iarg++] = NULL;
4225 for (iarg = 0; args[iarg] != NULL; iarg++) {
4226 printf(
"arg[%d] [%s]\n", iarg, args[iarg]);
4233 if (errstr != NULL) {
4234 sprintf(errstr,
"Cannot execute mtransition, ss_spawnv() returned %d",
status);
4249 int connect_timeout = 10000;
4250 int timeout = 120000;
4278 assert(wait_for_index >= 0);
4279 assert(wait_for_index < (
int)s->
clients.size());
4301 if (wait_for == NULL)
4308 printf(
"Client \"%s\" waits for client \"%s\"\n", tr_client->
client_name.c_str(), wait_for->
client_name.c_str());
4315 cm_msg(
MERROR,
"cm_transition_call",
"Client \"%s\" transition %d aborted while waiting for client \"%s\": \"/Runinfo/Transition in progress\" was cleared", tr_client->
client_name.c_str(), tr_client->
transition, wait_for->
client_name.c_str());
4331 printf(
"Connecting to client \"%s\" on host %s...\n", tr_client->
client_name.c_str(), tr_client->
host_name.c_str());
4333 cm_msg(
MINFO,
"cm_transition_call",
"cm_transition_call: Connecting to client \"%s\" on host %s...", tr_client->
client_name.c_str(), tr_client->
host_name.c_str());
4336 size =
sizeof(timeout);
4339 if (connect_timeout < 1000)
4340 connect_timeout = 1000;
4343 size =
sizeof(timeout);
4368 "cannot connect to client \"%s\" on host %s, port %d, status %d",
4392 printf(
"Connection established to client \"%s\" on host %s\n", tr_client->
client_name.c_str(), tr_client->
host_name.c_str());
4395 "cm_transition: Connection established to client \"%s\" on host %s",
4407 printf(
"Executing RPC transition client \"%s\" on host %s...\n",
4411 "cm_transition: Executing RPC transition client \"%s\" on host %s...",
4439 printf(
"RPC transition finished client \"%s\" on host \"%s\" in %d ms with status %d\n",
4443 "cm_transition: RPC transition finished client \"%s\" on host \"%s\" in %d ms with status %d",
4463 printf(
"hconn %d cm_transition_call(%s) finished init %d connect %d end %d rpc %d end %d xxx %d end %d\n",
4510 for (
size_t i = 0;
i <
n;
i++) {
4518 printf(
"Calling local transition callback\n");
4520 cm_msg(
MINFO,
"cm_transition_call_direct",
"cm_transition: Calling local transition callback");
4536 printf(
"Local transition callback finished, status %d\n",
int(tr_client->
status));
4538 cm_msg(
MINFO,
"cm_transition_call_direct",
"cm_transition: Local transition callback finished, status %d",
int(tr_client->
status));
4546 return tr_client->
status;
4608 char tr_key_name[256];
4618 errstr_size =
sizeof(xerrstr);
4634 mstrlcpy(errstr,
"Invalid transition request", errstr_size);
4645 sprintf(errstr,
"Start/Stop transition %d already in progress, please try again later\n",
i);
4646 mstrlcat(errstr,
"or set \"/Runinfo/Transition in progress\" manually to zero.\n", errstr_size);
4648 cm_msg(
MERROR,
"cm_transition",
"another transition is already in progress");
4701 cm_msg(
MERROR,
"cm_transition",
"Run start abort due to alarms: %s", alarms.c_str());
4702 mstrlcpy(errstr,
"Cannot start run due to alarms: ", errstr_size);
4703 mstrlcat(errstr, alarms.c_str(), errstr_size);
4714 HNDLE hkeyroot, hkey;
4731 size =
sizeof(program_info_required);
4734 cm_msg(
MERROR,
"cm_transition",
"Cannot get program info required, status %d",
status);
4738 if (program_info_required) {
4743 cm_msg(
MERROR,
"cm_transition",
"Run start abort due to program \"%s\" not running",
key.
name);
4744 std::string serrstr =
msprintf(
"Run start abort due to program \"%s\" not running",
key.
name);
4745 mstrlcpy(errstr, serrstr.c_str(), errstr_size);
4759 mstrlcpy(errstr,
"Unknown error", errstr_size);
4761 if (debug_flag == 0) {
4788 if (debug_flag == 1)
4789 printf(
"Setting run number %d in ODB\n",
run_number);
4790 if (debug_flag == 2)
4795 cm_msg(
MERROR,
"cm_transition",
"cannot set Runinfo/Run number in database, status %d",
status);
4801 if (debug_flag == 1)
4802 printf(
"Clearing /Runinfo/Requested transition\n");
4803 if (debug_flag == 2)
4804 cm_msg(
MINFO,
"cm_transition",
"cm_transition: Clearing /Runinfo/Requested transition");
4812 cm_msg(
MERROR,
"cm_transition",
"cannot find System/Clients entry in database");
4814 mstrlcpy(errstr,
"Cannot find /System/Clients in ODB", errstr_size);
4823 mstrlcpy(errstr,
"Deferred transition already in progress", errstr_size);
4824 mstrlcat(errstr,
", to cancel, set \"/Runinfo/Requested transition\" to zero", errstr_size);
4831 sprintf(tr_key_name,
"Transition %s DEFERRED", trname.c_str());
4840 size =
sizeof(sequence_number);
4849 if (debug_flag == 1)
4850 printf(
"---- Transition %s deferred by client \"%s\" ----\n", trname.c_str(),
str);
4851 if (debug_flag == 2)
4852 cm_msg(
MINFO,
"cm_transition",
"cm_transition: ---- Transition %s deferred by client \"%s\" ----", trname.c_str(),
str);
4854 if (debug_flag == 1)
4855 printf(
"Setting /Runinfo/Requested transition\n");
4856 if (debug_flag == 2)
4857 cm_msg(
MINFO,
"cm_transition",
"cm_transition: Setting /Runinfo/Requested transition");
4868 sprintf(errstr,
"Transition %s deferred by client \"%s\"", trname.c_str(),
str);
4899 size =
sizeof(program_info_auto_start);
4902 cm_msg(
MERROR,
"cm_transition",
"Cannot get program info auto start, status %d",
status);
4906 if (program_info_auto_start) {
4908 start_command[0] = 0;
4910 size =
sizeof(start_command);
4913 cm_msg(
MERROR,
"cm_transition",
"Cannot get program info start command, status %d",
status);
4917 if (start_command[0]) {
4918 cm_msg(
MINFO,
"cm_transition",
"Auto Starting program \"%s\", command \"%s\"",
key.
name,
4953 size =
sizeof(
state);
4959 cm_msg(
MERROR,
"cm_transition",
"cannot get Runinfo/State in database");
4966 cm_msg(
MERROR,
"cm_transition",
"cannot set \"Runinfo/Stop Time binary\" in database");
4973 cm_msg(
MERROR,
"cm_transition",
"cannot set \"Runinfo/Stop Time\" in database");
4979 cm_msg(
MERROR,
"cm_transition",
"cannot find System/Clients entry in database");
4981 mstrlcpy(errstr,
"Cannot find /System/Clients in ODB", errstr_size);
5011 if (debug_flag == 1)
5012 printf(
"---- Transition %s started ----\n", trname.c_str());
5013 if (debug_flag == 2)
5014 cm_msg(
MINFO,
"cm_transition",
"cm_transition: ---- Transition %s started ----", trname.c_str());
5016 sprintf(tr_key_name,
"Transition %s", trname.c_str());
5020 for (
int i = 0,
status = 0;;
i++) {
5037 size =
sizeof(sequence_number);
5046 c->async_flag = async_flag;
5047 c->debug_flag = debug_flag;
5048 c->sequence_number = sequence_number;
5050 c->key_name = subkey.
name;
5054 size =
sizeof(client_name);
5056 c->client_name = client_name;
5069 size =
sizeof(port);
5081 if (cc->
port ==
c->port)
5087 s.
clients.push_back(std::unique_ptr<TrClient>(
c));
5090 cm_msg(
MERROR,
"cm_transition",
"transition %s: client \"%s\" is registered with sequence number %d more than once", trname.c_str(),
c->client_name.c_str(),
c->sequence_number);
5102 for (
size_t idx = 0; idx < s.
clients.size(); idx++) {
5103 if (s.
clients[idx]->sequence_number == 0) {
5108 for (
size_t i = idx - 1; ;
i--) {
5109 if (s.
clients[
i]->sequence_number < s.
clients[idx]->sequence_number) {
5110 if (s.
clients[
i]->sequence_number > 0) {
5111 s.
clients[idx]->wait_for_index.push_back(
i);
5121 for (
size_t idx = 0; idx < s.
clients.size(); idx++) {
5126 for (
size_t idx = 0; idx < s.
clients.size(); idx++) {
5127 printf(
"TrClient[%d]: ",
int(idx));
5135 for (
size_t idx = 0; idx < s.
clients.size(); idx++) {
5136 if (debug_flag == 1)
5137 printf(
"\n==== Found client \"%s\" with sequence number %d\n",
5138 s.
clients[idx]->client_name.c_str(), s.
clients[idx]->sequence_number);
5139 if (debug_flag == 2)
5141 "cm_transition: ==== Found client \"%s\" with sequence number %d",
5142 s.
clients[idx]->client_name.c_str(), s.
clients[idx]->sequence_number);
5146 assert(s.
clients[idx]->thread == NULL);
5149 if (s.
clients[idx]->port == 0) {
5159 cm_msg(
MERROR,
"cm_transition",
"transition %s aborted: client \"%s\" returned status %d", trname.c_str(),
5160 s.
clients[idx]->client_name.c_str(),
int(s.
clients[idx]->status));
5170 for (
size_t idx = 0; idx < s.
clients.size(); idx++) {
5173 s.
clients[idx]->thread->join();
5174 delete s.
clients[idx]->thread;
5175 s.
clients[idx]->thread = NULL;
5186 cm_msg(
MERROR,
"cm_transition",
"transition %s aborted: \"/Runinfo/Transition in progress\" was cleared", trname.c_str());
5189 mstrlcpy(errstr,
"Canceled", errstr_size);
5195 for (
size_t idx = 0; idx < s.
clients.size(); idx++)
5199 mstrlcpy(errstr, s.
clients[idx]->errorstr.c_str(), errstr_size);
5214 if (debug_flag == 1)
5215 printf(
"\n---- Transition %s finished ----\n", trname.c_str());
5216 if (debug_flag == 2)
5217 cm_msg(
MINFO,
"cm_transition",
"cm_transition: ---- Transition %s finished ----", trname.c_str());
5232 size =
sizeof(
state);
5235 cm_msg(
MERROR,
"cm_transition",
"cannot set Runinfo/State in database, db_set_value() status %d",
status);
5283 size =
sizeof(program_info_auto_stop);
5286 cm_msg(
MERROR,
"cm_transition",
"Cannot get program info auto stop, status %d",
status);
5290 if (program_info_auto_stop) {
5304 mstrlcpy(errstr,
"Success", errstr_size);
5318 cm_msg(
MERROR,
"cm_transition",
"Could not start a run: cm_transition() status %d, message \'%s\'",
status,
5362 int sflag = async_flag &
TR_SYNC;
5367 cm_msg(
MERROR,
"cm_transition",
"previous transition did not finish yet");
5382 mstrlcpy(errstr,
"Invalid transition request", errstr_size);
5390 int size =
sizeof(
i);
5394 sprintf(errstr,
"Start/Stop transition %d already in progress, please try again later\n",
i);
5395 mstrlcat(errstr,
"or set \"/Runinfo/Transition in progress\" manually to zero.\n", errstr_size);
5397 cm_msg(
MERROR,
"cm_transition",
"another transition is already in progress");
5456#ifndef DOXYGEN_SHOULD_SKIP_THIS
5484 if (client_socket) {
5498 if (strchr(
str,
' '))
5499 *strchr(
str,
' ') = 0;
5517 printf(
"Received 2nd Ctrl-C, hard abort\n");
5520 printf(
"Received Ctrl-C, aborting...\n");
5580 std::string command;
5585 cm_msg(
MERROR,
"cm_exec_script",
"Script ODB \"%s\" of type TID_STRING, db_get_value_string() error %d",
5586 odb_path_to_script,
status);
5590 for (
int i = 0;;
i++) {
5602 cm_msg(
MERROR,
"cm_exec_script",
"Script ODB \"%s/%s\" should not be TID_KEY", odb_path_to_script,
5607 char *buf = (
char *) malloc(size);
5608 assert(buf != NULL);
5611 cm_msg(
MERROR,
"cm_exec_script",
"Script ODB \"%s/%s\" of type %d, db_get_data() error %d",
5625 cm_msg(
MERROR,
"cm_exec_script",
"Script ODB \"%s\" has invalid type %d, should be TID_STRING or TID_KEY",
5626 odb_path_to_script,
key.
type);
5632 if (command.length() > 0) {
5633 cm_msg(
MINFO,
"cm_exec_script",
"Executing script \"%s\" from ODB \"%s\"", command.c_str(), odb_path_to_script);
5654 static DWORD alarm_last_checked_sec = 0;
5658 static DWORD last_millitime = 0;
5659 DWORD tdiff_millitime = now_millitime - last_millitime;
5660 const DWORD kPeriod = 1000;
5661 if (last_millitime == 0) {
5662 last_millitime = now_millitime;
5663 tdiff_millitime = kPeriod;
5673 if (now_sec - alarm_last_checked_sec > 10) {
5675 alarm_last_checked_sec = now_sec;
5680 if (tdiff_millitime >= kPeriod) {
5682 if (tdiff_millitime > 60000)
5683 wrong_interval =
TRUE;
5687 bm_cleanup(
"cm_periodic_tasks", now_millitime, wrong_interval);
5688 db_cleanup(
"cm_periodic_tasks", now_millitime, wrong_interval);
5692 last_millitime = now_millitime;
5801 static int check_cm_execute = 1;
5802 static int enable_cm_execute = 0;
5807 if (check_cm_execute) {
5811 check_cm_execute = 0;
5816 size =
sizeof(enable_cm_execute);
5823 if (!enable_cm_execute) {
5825 mstrlcpy(buf, command,
sizeof(buf));
5826 cm_msg(
MERROR,
"cm_execute",
"cm_execute(%s...) is disabled by ODB \"/Experiment/Enable cm_execute\"", buf);
5832 std::string
str =
msprintf(
"%s > %s", command, filename.c_str());
5836 fh = open(filename.c_str(), O_RDONLY, 0644);
5839 n =
read(fh, result, bufsize - 1);
5840 result[
MAX(0,
n)] = 0;
5843 remove(filename.c_str());
5845 status = system(command);
5859#ifndef DOXYGEN_SHOULD_SKIP_THIS
5895 sprintf(
str,
"RPC/%d",
id);
5923 if (history_channel && (strlen(history_channel) > 0)) {
5925 p +=
"/Logger/History/";
5926 p += history_channel;
5927 p +=
"/History dir";
5988#ifdef LOCAL_ROUTINES
5998 bool badindex =
false;
5999 bool badclient =
false;
6009 if (pclient->
name[0] == 0)
6024 printf(
"bm_validate_client_index: pbuf=%p, buf_name \"%s\", client_index=%d, max_client_index=%d, badindex %d, pid=%d\n",
6027 }
else if (badclient) {
6028 printf(
"bm_validate_client_index: pbuf=%p, buf_name \"%s\", client_index=%d, max_client_index=%d, client_name=\'%s\', client_pid=%d, pid=%d, badclient %d\n",
6033 printf(
"bm_validate_client_index: pbuf=%p, buf_name \"%s\", client_index=%d, max_client_index=%d, client_name=\'%s\', client_pid=%d, pid=%d, goodclient\n",
6040 if (badindex || badclient) {
6041 static int prevent_recursion = 1;
6043 if (prevent_recursion) {
6044 prevent_recursion = 0;
6052 cm_msg(
MERROR,
"bm_validate_client_index",
"Maybe this client was removed by a timeout. See midas.log. Cannot continue, aborting...");
6061 fprintf(stderr,
"bm_validate_client_index: Maybe this client was removed by a timeout. See midas.log. Cannot continue, aborting...\n");
6101#ifdef LOCAL_ROUTINES
6127 pbctmp = pheader->
client;
6144 pbclient = pheader->
client;
6148 if (pbclient->
pid) {
6151 "Client \'%s\' on buffer \'%s\' removed by %s because process pid %d does not exist", pbclient->
name,
6152 pheader->
name, who, pbclient->
pid);
6163 printf(
"buffer [%s] client [%-32s] times 0x%08x 0x%08x, diff 0x%08x %5d, timeout %d\n",
6175 cm_msg(
MINFO,
"bm_cleanup",
"Client \'%s\' on buffer \'%s\' removed by %s (idle %1.1lfs, timeout %1.0lfs)",
6176 pbclient->
name, pheader->
name, who,
6192 std::vector<BUFFER*> mybuffers;
6198 for (
BUFFER* pbuf : mybuffers) {
6201 if (pbuf->attached) {
6211 if (pclient->
pid == pid) {
6226#ifdef LOCAL_ROUTINES
6230 std::vector<BUFFER*> mybuffers;
6237 for (
BUFFER* pbuf : mybuffers) {
6240 if (pbuf->attached) {
6252 if (!wrong_interval)
6259#ifdef LOCAL_ROUTINES
6262 if (rp < 0 || rp > pheader->
size) {
6264 "error: buffer \"%s\" is corrupted: rp %d is invalid. buffer read_pointer %d, write_pointer %d, size %d, called from %s",
6277 "error: buffer \"%s\" is corrupted: rp %d plus event header point beyond the end of buffer by %d bytes. buffer read_pointer %d, write_pointer %d, size %d, called from %s",
6292static FILE* gRpLog = NULL;
6298 if (gRpLog == NULL) {
6299 gRpLog = fopen(
"rp.log",
"a");
6301 if (gRpLog && (total_size < 16)) {
6302 const char *pdata = (
const char *) (pheader + 1);
6303 const DWORD *pevent = (
const DWORD*) (pdata + rp);
6304 fprintf(gRpLog,
"%s: rp %d, total_size %d, at rp 0x%08x 0x%08x 0x%08x 0x%08x 0x%08x 0x%08x\n", pheader->
name, rp, total_size,
6305 pevent[0], pevent[1], pevent[2], pevent[3], pevent[4], pevent[5]);
6311 assert(total_size > 0);
6315 if (rp >= pheader->
size) {
6316 rp -= pheader->
size;
6333 if (pevent->
data_size <= 0 || total_size <= 0 || total_size > pheader->
size) {
6335 "error: buffer \"%s\" is corrupted: rp %d points to an invalid event: data_size %d, event size %d, total_size %d, buffer read_pointer %d, write_pointer %d, size %d, called from %s",
6349 if (rp < pheader->write_pointer) {
6352 remaining = pheader->
size - rp;
6358 if (total_size > remaining) {
6360 "error: buffer \"%s\" is corrupted: rp %d points to an invalid event: data_size %d, event size %d, total_size %d, buffer read_pointer %d, write_pointer %d, size %d, remaining %d, called from %s",
6381 const char *pdata = (
const char *) (pheader + 1);
6391 "buffer \"%s\" is corrupted: invalid read pointer %d. Size %d, write pointer %d", pheader->
name,
6398 "buffer \"%s\" is corrupted: invalid write pointer %d. Size %d, read pointer %d", pheader->
name,
6404 cm_msg(
MERROR,
"bm_validate_buffer",
"buffer \"%s\" is corrupted: read pointer %d is invalid", pheader->
name,
6413 cm_msg(
MERROR,
"bm_validate_buffer",
"buffer \"%s\" is corrupted: invalid rp %d, last good event at rp %d",
6414 pheader->
name, rp, rp0);
6418 int rp1 =
bm_next_rp(
"bm_validate_buffer_locked", pheader, pdata, rp);
6421 "buffer \"%s\" is corrupted: invalid event at rp %d, last good event at rp %d", pheader->
name, rp, rp0);
6440 get_all = (get_all || xget_all);
6444 int rp =
c->read_pointer;
6448 int rp1 =
bm_next_rp(
"bm_validate_buffer_locked", pheader, pdata, rp);
6451 "buffer \"%s\" is corrupted for client \"%s\" rp %d: invalid event at rp %d, last good event at rp %d",
6452 pheader->
name,
c->name,
c->read_pointer, rp, rp0);
6553 double buf_size = pheader->
size;
6557 double buf_fill = 0;
6558 double buf_cptr = 0;
6559 double buf_cused = 0;
6560 double buf_cused_pct = 0;
6562 if (client_index >= 0 && client_index <= pheader->max_client_index) {
6565 if (buf_wptr == buf_cptr) {
6567 }
else if (buf_wptr > buf_cptr) {
6568 buf_cused = buf_wptr - buf_cptr;
6570 buf_cused = (buf_size - buf_cptr) + buf_wptr;
6573 buf_cused_pct = buf_cused / buf_size * 100.0;
6582 if (buf_wptr == buf_rptr) {
6584 }
else if (buf_wptr > buf_rptr) {
6585 buf_fill = buf_wptr - buf_rptr;
6587 buf_fill = (buf_size - buf_rptr) + buf_wptr;
6590 double buf_fill_pct = buf_fill / buf_size * 100.0;
6640 sprintf(
str,
"writes_blocked_by/%s/count_write_wait", pheader->
client[
i].
name);
6643 sprintf(
str,
"writes_blocked_by/%s/time_write_wait", pheader->
client[
i].
name);
6669 if ((strlen(
buffer_name.c_str()) < 1) || (strlen(client_name.c_str()) < 1)) {
6672 cm_msg(
MERROR,
"bm_write_buffer_statistics_to_odb",
"Invalid empty buffer name \"%s\" or client name \"%s\"",
buffer_name.c_str(), client_name.c_str());
6689 size_t sbuffer_handle = buffer_handle;
6697 if (buffer_handle >=1 && sbuffer_handle <= nbuf) {
6703 if (sbuffer_handle > nbuf || buffer_handle <= 0) {
6705 cm_msg(
MERROR, who,
"invalid buffer handle %d: out of range [1..%d]", buffer_handle, (
int)nbuf);
6713 cm_msg(
MERROR, who,
"invalid buffer handle %d: empty slot", buffer_handle);
6721 cm_msg(
MERROR, who,
"invalid buffer handle %d: not attached", buffer_handle);
6797 int size =
sizeof(
INT);
6801 cm_msg(
MERROR,
"bm_open_buffer",
"Cannot get ODB /Experiment/MAX_EVENT_SIZE, db_get_value() status %d",
6808#ifdef LOCAL_ROUTINES
6813 const int max_buffer_size = 2 * 1000 * 1024 * 1024;
6818 cm_msg(
MERROR,
"bm_open_buffer",
"cannot open buffer with zero name");
6835 std::string odb_path;
6836 odb_path +=
"/Experiment/Buffer sizes/";
6839 int size =
sizeof(
INT);
6842 if (buffer_size <= 0 || buffer_size > max_buffer_size) {
6844 "Cannot open buffer \"%s\", invalid buffer size %d in ODB \"%s\", maximum buffer size is %d",
6845 buffer_name, buffer_size, odb_path.c_str(), max_buffer_size);
6853 "Will use default SYSMSG buffer size %d to open buffer \"%s\"",
6864 cm_msg(
MERROR,
"bm_open_buffer",
"Cannot get ODB /Experiment/MAX_EVENT_SIZE, db_get_value() status %d",
6874 *buffer_handle =
i + 1;
6883 static std::mutex gNewBufferMutex;
6884 std::lock_guard<std::mutex> guard(gNewBufferMutex);
6894 *buffer_handle =
i + 1;
6967 pheader->
size = buffer_size;
6977 "Buffer \"%s\" is corrupted, mismatch of buffer name in shared memory \"%s\"",
buffer_name,
6988 cm_msg(
MERROR,
"bm_open_buffer",
"Buffer \"%s\" is corrupted, num_clients %d exceeds MAX_CLIENTS %d",
6999 cm_msg(
MERROR,
"bm_open_buffer",
"Buffer \"%s\" is corrupted, max_client_index %d exceeds MAX_CLIENTS %d",
7007 if (pheader->
size != buffer_size) {
7008 cm_msg(
MINFO,
"bm_open_buffer",
"Buffer \"%s\" requested size %d differs from existing size %d",
7011 buffer_size = pheader->
size;
7043 "buffer \'%s\' is corrupted, bm_validate_buffer() status %d, calling bm_reset_buffer()...",
buffer_name,
7082 mstrlcpy(pclient->
name, client_name.c_str(),
sizeof(pclient->
name));
7113 *buffer_handle =
i+1;
7118 *buffer_handle =
gBuffers.size() + 1;
7155 *buffer_handle =
i + 1;
7176#ifdef LOCAL_ROUTINES
7191 std::vector<EventRequest> request_list_copy =
_request_list;
7193 for (
size_t i = 0;
i < request_list_copy.size();
i++) {
7194 if (request_list_copy[
i].buffer_handle == buffer_handle) {
7321#ifdef LOCAL_ROUTINES
7329 for (
size_t i = nbuf;
i > 0;
i--) {
7355#ifdef LOCAL_ROUTINES
7367 std::vector<BUFFER*> mybuffers;
7373 for (
BUFFER* pbuf : mybuffers) {
7374 if (!pbuf || !pbuf->attached)
7393#ifdef LOCAL_ROUTINES
7412 for (
i = 0;
i < 20;
i++) {
7434#ifdef LOCAL_ROUTINES
7449#ifdef LOCAL_ROUTINES
7494 for (
int i = 0;;
i++) {
7504 std::string client_name;
7505 std::string remote_host;
7512 size =
sizeof(client_name);
7521 size =
sizeof(port);
7537 int client_pid = atoi(
key.
name);
7539 cm_msg(
MERROR,
"cm_shutdown",
"Cannot connect to client \'%s\' on host \'%s\', port %d", client_name.c_str(), remote_host.c_str(), port);
7541 cm_msg(
MERROR,
"cm_shutdown",
"Killing and Deleting client \'%s\' pid %d", client_name.c_str(), client_pid);
7542 kill(client_pid, SIGKILL);
7546 cm_msg(
MERROR,
"cm_shutdown",
"Cannot delete client info for client \'%s\', pid %d, status %d",
name, client_pid,
status);
7560 int client_pid = atoi(
key.
name);
7562 cm_msg(
MERROR,
"cm_shutdown",
"Client \'%s\' not responding to shutdown command", client_name.c_str());
7564 cm_msg(
MERROR,
"cm_shutdown",
"Killing and Deleting client \'%s\' pid %d", client_name.c_str(), client_pid);
7565 kill(client_pid, SIGKILL);
7568 cm_msg(
MERROR,
"cm_shutdown",
"Cannot delete client info for client \'%s\', pid %d, status %d",
name, client_pid,
status);
7582 return return_status;
7612 for (
int i = 0;;
i++) {
7621 std::string client_name;
7648 std::string program_name;
7664 size_t name_len = strlen(
name);
7666 if (client_name.length() < name_len) {
7671 std::string truncated_client_name = client_name;
7673 truncated_client_name.resize(name_len);
7727#ifdef LOCAL_ROUTINES
7732 std::vector<BUFFER*> mybuffers;
7739 for (
BUFFER* pbuf : mybuffers) {
7742 if (pbuf->attached) {
7758 if (
j != pbuf->client_index && pbclient->
pid &&
7759 (client_name == NULL || client_name[0] == 0
7760 || strncmp(pbclient->
name, client_name, strlen(client_name)) == 0)) {
7774 "Client \'%s\' on \'%s\' removed by cm_cleanup (idle %1.1lfs, timeout %1.0lfs)",
7797 db_cleanup2(client_name, ignore_timeout, now,
"cm_cleanup");
7824 const char *s =
str;
7829 std::string envname;
7836 const char *
e = getenv(envname.c_str());
7857 printf(
"test_expand_env: [%s] -> [%s] expected [%s]",
7861 if (s != expected) {
7862 printf(
", MISMATCH!\n");
7871 printf(
"Test expand_end()\n");
7872 setenv(
"FOO",
"foo", 1);
7873 setenv(
"BAR",
"bar", 1);
7874 setenv(
"EMPTY",
"", 1);
7890 printf(
"test_expand_env: all tests passed!\n");
7892 printf(
"test_expand_env: test FAILED!\n");
7901#ifndef DOXYGEN_SHOULD_SKIP_THIS
7930#ifdef LOCAL_ROUTINES
7974#ifdef LOCAL_ROUTINES
7994 *n_bytes += pheader->
size;
8014#ifdef LOCAL_ROUTINES
8022 fprintf(stderr,
"bm_lock_buffer_read_cache: Error: Cannot lock read cache of buffer \"%s\", ss_timed_mutex_wait_for_sec() timeout, aborting...\n", pbuf->
buffer_name);
8023 cm_msg(
MERROR,
"bm_lock_buffer_read_cache",
"Cannot lock read cache of buffer \"%s\", ss_timed_mutex_wait_for_sec() timeout, aborting...", pbuf->
buffer_name);
8030 fprintf(stderr,
"bm_lock_buffer_read_cache: Error: Cannot lock read cache of buffer \"%s\", buffer was closed while we waited for the buffer_mutex\n", pbuf->
buffer_name);
8043 fprintf(stderr,
"bm_lock_buffer_write_cache: Error: Cannot lock write cache of buffer \"%s\", ss_timed_mutex_wait_for_sec() timeout, aborting...\n", pbuf->
buffer_name);
8044 cm_msg(
MERROR,
"bm_lock_buffer_write_cache",
"Cannot lock write cache of buffer \"%s\", ss_timed_mutex_wait_for_sec() timeout, aborting...", pbuf->
buffer_name);
8051 fprintf(stderr,
"bm_lock_buffer_write_cache: Error: Cannot lock write cache of buffer \"%s\", buffer was closed while we waited for the buffer_mutex\n", pbuf->
buffer_name);
8066 fprintf(stderr,
"bm_lock_buffer_mutex: Error: Cannot lock buffer \"%s\", ss_timed_mutex_wait_for_sec() timeout, aborting...\n", pbuf->
buffer_name);
8067 cm_msg(
MERROR,
"bm_lock_buffer_mutex",
"Cannot lock buffer \"%s\", ss_timed_mutex_wait_for_sec() timeout, aborting...", pbuf->
buffer_name);
8074 fprintf(stderr,
"bm_lock_buffer_mutex: Error: Cannot lock buffer \"%s\", buffer was closed while we waited for the buffer_mutex\n", pbuf->
buffer_name);
8107 fprintf(stderr,
"bm_lock_buffer: Lock buffer \"%s\" is taking longer than 1 second!\n", pbuf->
buffer_name);
8112 fprintf(stderr,
"bm_lock_buffer: Lock buffer \"%s\" is taking longer than 10 seconds, buffer semaphore is probably stuck, delete %s.SHM and try again!\n", pbuf->
buffer_name, pbuf->
buffer_name);
8123 fprintf(stderr,
"bm_lock_buffer: Error: Cannot lock buffer \"%s\", ss_semaphore_wait_for() status %d, aborting...\n", pbuf->
buffer_name,
status);
8161 printf(
"unlock [??????] unused1 ????? pid %d\n", getpid());
8199#ifdef LOCAL_ROUTINES
8259#ifdef LOCAL_ROUTINES
8278 if (write_size > 0) {
8287 if (write_size > max_write_size) {
8288 size_t new_write_size = max_write_size;
8289 cm_msg(
MERROR,
"bm_set_cache_size",
"requested write cache size %zu on buffer \"%s\" is too big: buffer size is %d, write cache size will be %zu bytes", write_size, pbuf->
buffer_name, pbuf->
buffer_header->
size, new_write_size);
8290 write_size = new_write_size;
8308 if (read_size > 0) {
8309 pbuf->
read_cache = (
char *) malloc(read_size);
8315 cm_msg(
MERROR,
"bm_set_cache_size",
"not enough memory to allocate read cache for buffer \"%s\", malloc(%zu) failed", pbuf->
buffer_name, read_size);
8344 if (write_size > 0) {
8351 cm_msg(
MERROR,
"bm_set_cache_size",
"not enough memory to allocate write cache for buffer \"%s\", malloc(%zu) failed", pbuf->
buffer_name, write_size);
8407 static std::mutex mutex;
8414 std::lock_guard<std::mutex> lock(mutex);
8424#ifndef DOXYGEN_SHOULD_SKIP_THIS
8476#ifdef LOCAL_ROUTINES
8492 if (func == NULL && pbuf->
callback) {
8494 cm_msg(
MERROR,
"bm_add_event_request",
"mixing callback/non callback requests not possible");
8501 cm_msg(
MERROR,
"bm_add_event_request",
"GET_RECENT request not possible if read cache is enabled");
8580 INT sampling_type,
HNDLE *request_id,
8583 assert(request_id != NULL);
8635#ifdef LOCAL_ROUTINES
8701 if (request_id < 0 ||
size_t(request_id) >=
_request_list.size()) {
8706 int buffer_handle =
_request_list[request_id].buffer_handle;
8722 pclient = pheader->
client;
8724 printf(
"buffer \'%s\', rptr: %d, wptr: %d, size: %d\n", pheader->
name, pheader->
read_pointer,
8727 if (pclient[
i].pid) {
8728 printf(
"pointers: client %d \'%s\', rptr %d\n",
i, pclient[
i].
name, pclient[
i].read_pointer);
8743 "Corrected read pointer for client \'%s\' on buffer \'%s\' from %d to %d, write pointer %d, size %d",
8752 "Corrected read pointer for client \'%s\' on buffer \'%s\' from %d to %d, read pointer %d, size %d",
8763 "Corrected read pointer for client \'%s\' on buffer \'%s\' from %d to %d, write pointer %d, size %d",
8772 "Corrected read pointer for client \'%s\' on buffer \'%s\' from %d to %d, write pointer %d, size %d",
8781 "Corrected read pointer for client \'%s\' on buffer \'%s\' from %d to %d, write pointer %d, size %d",
8797 if (pclient[
i].pid) {
8798 bm_validate_client_pointers(pheader, &pclient[
i]);
8834 assert(caller_name);
8846 printf(
"bm_update_read_pointer: [%s] rp %d, wp %d, size %d, min_rp %d, client [%s] rp %d\n",
8877 min_rp += pheader->
size;
8879 assert(min_rp >= 0);
8880 assert(min_rp < pheader->size);
8887 printf(
"bm_update_read_pointer: [%s] rp %d, wp %d, size %d, new_rp %d, moved\n",
8902 int have_get_all_requests = 0;
8909 if (!have_get_all_requests)
8919 if (free_space <= 0)
8920 free_space += pheader->
size;
8922 if (free_space >= pheader->
size * 0.5) {
8942 for (
size_t i = 0;
i <
n;
i++) {
8959 r.
dispatcher(buffer_handle,
i, pevent, (
void *) (pevent + 1));
8966#ifdef LOCAL_ROUTINES
8992 *ptotal_size = total_size;
9021 cm_msg(
MERROR,
"bm_peek_buffer_locked",
"event buffer \"%s\" is corrupted: client \"%s\" read pointer %d is invalid. buffer read pointer %d, write pointer %d, size %d", pheader->
name, pc->
name, pc->
read_pointer, pheader->
read_pointer, pheader->
write_pointer, pheader->
size);
9025 char *pdata = (
char *) (pheader + 1);
9031 if ((total_size <= 0) || (total_size > pheader->
size)) {
9032 cm_msg(
MERROR,
"bm_peek_buffer_locked",
"event buffer \"%s\" is corrupted: client \"%s\" read pointer %d points to invalid event: data_size %d, event_size %d, total_size %d. buffer size: %d, read_pointer: %d, write_pointer: %d", pheader->
name, pc->
name, pc->
read_pointer, pevent->
data_size,
event_size, total_size, pheader->
size, pheader->
read_pointer, pheader->
write_pointer);
9036 assert(total_size > 0);
9037 assert(total_size <= pheader->size);
9044 *ptotal_size = total_size;
9051 const char *pdata = (
const char *) (pheader + 1);
9053 if (rp + event_size <= pheader->size) {
9058 int size = pheader->
size - rp;
9059 memcpy(buf, pdata + rp, size);
9060 memcpy(buf + size, pdata,
event_size - size);
9066 const char *pdata = (
const char *) (pheader + 1);
9068 if (rp + event_size <= pheader->size) {
9070 vecptr->assign(pdata + rp, pdata + rp +
event_size);
9073 int size = pheader->
size - rp;
9074 vecptr->assign(pdata + rp, pdata + rp + size);
9075 vecptr->insert(vecptr->end(), pdata, pdata +
event_size - size);
9085 if (prequest->
valid) {
9095 is_requested =
TRUE;
9100 return is_requested;
9184 if (convert_flags) {
9206 char *pdata = (
char *) (pheader + 1);
9213 requested_space += 100;
9215 if (requested_space >= pheader->
size)
9219 DWORD time_end = time_start + timeout_msec;
9223 int blocking_client_index = -1;
9225 blocking_client_name[0] = 0;
9233 free += pheader->
size;
9237 if (requested_space < free) {
9265 "error: buffer \"%s\" is corrupted: read_pointer %d, write_pointer %d, size %d, free %d, waiting for %d bytes: read pointer is invalid",
9280 printf(
"bm_wait_for_free_space: buffer pointers: read: %d, write: %d, free space: %d, bufsize: %d, event size: %d, blocking event size %d/%d\n", pheader->
read_pointer, pheader->
write_pointer, free, pheader->
size, requested_space,
event_size, total_size);
9283 if (pevent->
data_size <= 0 || total_size <= 0 || total_size > pheader->
size) {
9285 "error: buffer \"%s\" is corrupted: read_pointer %d, write_pointer %d, size %d, free %d, waiting for %d bytes: read pointer points to an invalid event: data_size %d, event size %d, total_size %d",
9298 int blocking_client = -1;
9328 blocking_client =
i;
9337 if (blocking_client >= 0) {
9338 blocking_client_index = blocking_client;
9339 mstrlcpy(blocking_client_name, pheader->
client[blocking_client].
name,
sizeof(blocking_client_name));
9356 "error: buffer \"%s\" is corrupted: read_pointer %d, write_pointer %d, size %d, free %d, waiting for %d bytes: read pointer did not move as expected",
9388 int sleep_time_msec = 1000;
9390 if (timeout_msec ==
BM_WAIT) {
9397 if (now >= time_end) {
9402 sleep_time_msec = time_end - now;
9404 if (sleep_time_msec <= 0) {
9405 sleep_time_msec = 10;
9406 }
else if (sleep_time_msec > 1000) {
9407 sleep_time_msec = 1000;
9417 if (unlock_write_cache)
9466 if (unlock_write_cache) {
9475 if (!pbuf_guard.
relock()) {
9476 if (unlock_write_cache) {
9527 DWORD time_wait = time_start + timeout_msec;
9528 DWORD sleep_time = 1000;
9531 }
else if (timeout_msec ==
BM_WAIT) {
9534 if (sleep_time > (
DWORD)timeout_msec)
9535 sleep_time = timeout_msec;
9556 if (unlock_read_cache)
9563 }
else if (timeout_msec ==
BM_WAIT) {
9568 if (now >= time_wait) {
9571 sleep_time = time_wait - now;
9572 if (sleep_time > 1000)
9580 if (unlock_read_cache) {
9588 if (!pbuf_guard.
relock()) {
9589 if (unlock_read_cache) {
9619 char *pdata = (
char *) (pheader + 1);
9627 for (
int i=0;
i<sg_n;
i++) {
9629 memcpy(wptr, sg_ptr[
i], sg_len[
i]);
9657 for (;
i<sg_n;
i++) {
9658 if (
count + sg_len[
i] > size)
9660 memcpy(wptr, sg_ptr[
i], sg_len[
i]);
9669 size_t first = size -
count;
9670 size_t second = sg_len[
i] - first;
9671 assert(first + second == sg_len[
i]);
9672 assert(
count + first == size);
9676 memcpy(wptr, sg_ptr[
i], first);
9679 memcpy(wptr, sg_ptr[
i] + first, second);
9686 for (;
i<sg_n;
i++) {
9687 memcpy(wptr, sg_ptr[
i], sg_len[
i]);
9708 return prequest->
id;
9717 if (request_id >= 0) {
9721 sprintf(
str,
"B %s %d", pheader->
name, request_id);
9738 DWORD time_end = time_start + timeout_msec;
9740 int xtimeout_msec = timeout_msec;
9743 if (timeout_msec ==
BM_WAIT) {
9744 xtimeout_msec = 1000;
9748 if (xtimeout_msec > 1000) {
9749 xtimeout_msec = 1000;
9758 if (timeout_msec ==
BM_WAIT) {
9766 if (now >= time_end) {
9771 DWORD remain = time_end - now;
9773 if (remain < xtimeout_msec) {
9774 xtimeout_msec = remain;
9793 const DWORD MAX_DATA_SIZE = (0x7FFFFFF0 - 16);
9796 if (data_size == 0) {
9797 cm_msg(
MERROR,
"bm_send_event",
"invalid event data size zero");
9801 if (data_size > MAX_DATA_SIZE) {
9802 cm_msg(
MERROR,
"bm_send_event",
"invalid event data size %d (0x%x) maximum is %d (0x%x)", data_size, data_size, MAX_DATA_SIZE, MAX_DATA_SIZE);
9820 const char* cptr =
event.data();
9821 size_t clen =
event.size();
9825int bm_send_event_vec(
int buffer_handle,
const std::vector<std::vector<char>>& event,
int timeout_msec)
9827 int sg_n =
event.size();
9828 const char* sg_ptr[sg_n];
9829 size_t sg_len[sg_n];
9830 for (
int i=0;
i<sg_n;
i++) {
9831 sg_ptr[
i] =
event[
i].data();
9832 sg_len[
i] =
event[
i].size();
9834 return bm_send_event_sg(buffer_handle, sg_n, sg_ptr, sg_len, timeout_msec);
9837#ifdef LOCAL_ROUTINES
9891int bm_send_event_sg(
int buffer_handle,
int sg_n,
const char*
const sg_ptr[],
const size_t sg_len[],
int timeout_msec)
9897 cm_msg(
MERROR,
"bm_send_event",
"invalid sg_n %d", sg_n);
9901 if (sg_ptr[0] == NULL) {
9902 cm_msg(
MERROR,
"bm_send_event",
"invalid sg_ptr[0] is NULL");
9907 cm_msg(
MERROR,
"bm_send_event",
"invalid sg_len[0] value %d is smaller than event header size %d", (
int)sg_len[0], (
int)
sizeof(
EVENT_HEADER));
9913 const DWORD MAX_DATA_SIZE = (0x7FFFFFF0 - 16);
9916 if (data_size == 0) {
9917 cm_msg(
MERROR,
"bm_send_event",
"invalid event data size zero");
9921 if (data_size > MAX_DATA_SIZE) {
9922 cm_msg(
MERROR,
"bm_send_event",
"invalid event data size %d (0x%x) maximum is %d (0x%x)", data_size, data_size, MAX_DATA_SIZE, MAX_DATA_SIZE);
9929 for (
int i=0;
i<sg_n;
i++) {
9934 cm_msg(
MERROR,
"bm_send_event",
"data size mismatch: event data_size %d, event_size %d not same as sum of sg_len %d", (
int)data_size, (
int)
event_size, (
int)
count);
9940#ifdef LOCAL_ROUTINES
9989 cm_msg(
MERROR,
"bm_send_event",
"write cache size is bigger than buffer size");
9998 if (!too_big && pbuf->
write_cache_wp + total_size <= pbuf->write_cache_size) {
10003 for (
int i=0;
i<sg_n;
i++) {
10004 memcpy(wptr, sg_ptr[
i], sg_len[
i]);
10034 printf(
"bm_send_event: corrupted 111!\n");
10040 if (total_size >= (
size_t)pheader->
size) {
10042 cm_msg(
MERROR,
"bm_send_event",
"total event size (%d) larger than size (%d) of buffer \'%s\'", (
int)total_size, pheader->
size, pheader->
name);
10056 printf(
"bm_send_event: corrupted 222!\n");
10085 printf(
"bm_send_event: corrupted 333!\n");
10105 DWORD time_end = time_start + timeout_msec;
10106 DWORD time_bombout = time_end;
10108 if (timeout_msec < 10000)
10109 time_bombout = time_start + 10000;
10111 int xtimeout_msec = timeout_msec;
10114 if (timeout_msec ==
BM_WAIT) {
10115 xtimeout_msec = 1000;
10119 if (xtimeout_msec > 1000) {
10120 xtimeout_msec = 1000;
10129 if (timeout_msec ==
BM_WAIT) {
10131 if (now >= time_bombout) {
10143 if (now >= time_end) {
10148 DWORD remain = time_end - now;
10150 if (remain < (
DWORD)xtimeout_msec) {
10151 xtimeout_msec = remain;
10154 if (now >= time_bombout) {
10190#ifdef LOCAL_ROUTINES
10209 request_id[
i] = -1;
10219 if (ask_rp == ask_wp) {
10223 assert(ask_rp < ask_wp);
10225 size_t ask_free =
ALIGN8(ask_wp - ask_rp);
10227 if (ask_free == 0) {
10234 printf(
"bm_flush_cache: corrupted 111!\n");
10271 printf(
"bm_flush_cache: cache size %d, wp %d, rp %d, event data_size %d, event_size %d, total_size %d, free %d, written %d\n",
10284 assert(total_size <= (
size_t)pheader->
size);
10341#ifdef LOCAL_ROUTINES
10390#ifdef LOCAL_ROUTINES
10397 max_size = *buf_size;
10419 if (!pbuf_guard.
relock()) {
10460 if (convert_flags) {
10463 }
else if (bufptr) {
10467 }
else if (vecptr) {
10469 char* cptr = (
char*)pevent;
10490 if (!pbuf_guard.
relock())
10527 if (is_requested) {
10535 "buffer size %d is smaller than event size %d, event truncated. buffer \"%s\"", max_size,
10547 if (convert_flags) {
10553 }
else if (dispatch || bufptr) {
10559 }
else if (vecptr) {
10622 xbuf_size = *buf_size;
10623 }
else if (ppevent) {
10628 xbuf = pvec->data();
10629 xbuf_size = pvec->size();
10631 assert(!
"incorrect call to bm_receivent_event_rpc()");
10636 DWORD time_end = time_start + timeout_msec;
10638 int xtimeout_msec = timeout_msec;
10640 int zbuf_size = xbuf_size;
10643 if (timeout_msec ==
BM_WAIT) {
10644 xtimeout_msec = 1000;
10648 if (xtimeout_msec > 1000) {
10649 xtimeout_msec = 1000;
10653 zbuf_size = xbuf_size;
10660 if (timeout_msec ==
BM_WAIT) {
10668 if (now >= time_end) {
10673 DWORD remain = time_end - now;
10675 if (remain < (
DWORD)xtimeout_msec) {
10676 xtimeout_msec = remain;
10691 }
else if (ppevent) {
10697 assert(!
"incorrect call to bm_receivent_event_rpc()");
10706 *buf_size = zbuf_size;
10707 }
else if (ppevent) {
10711 pvec->resize(zbuf_size);
10713 assert(!
"incorrect call to bm_receivent_event_rpc()");
10723 std::vector<char> *pv;
10726 pv =
new std::vector<char>;
10734 DWORD time_end = time_start + timeout_msec;
10736 int xtimeout_msec = timeout_msec;
10739 if (timeout_msec ==
BM_WAIT) {
10740 xtimeout_msec = 1000;
10744 if (xtimeout_msec > 1000) {
10745 xtimeout_msec = 1000;
10751 printf(
"bm_receive_event_rpc_cxx: handle %d, timeout %d, status %d, size %zu, via RPC_BM_RECEIVE_EVENT_CXX\n", buffer_handle, xtimeout_msec,
status, pv->size());
10754 if (timeout_msec ==
BM_WAIT) {
10762 if (now >= time_end) {
10767 DWORD remain = time_end - now;
10769 if (remain < (
DWORD)xtimeout_msec) {
10770 xtimeout_msec = remain;
10785 }
else if (ppevent) {
10791 assert(!
"incorrect call to bm_receivent_event_rpc_cxx()");
10803 if (pv->size() > (
size_t)*buf_size) {
10805 memcpy(buf, pv->data(), *buf_size);
10807 *buf_size = pv->size();
10808 memcpy(buf, pv->data(), *buf_size);
10810 }
else if (ppevent) {
10811 if (*ppevent == NULL) {
10813 assert(*ppevent != NULL);
10814 memcpy(*ppevent, pv->data(), pv->size());
10816 *ppevent = (
EVENT_HEADER*)realloc(*ppevent, pv->size());
10817 assert(*ppevent != NULL);
10818 memcpy(*ppevent, pv->data(), pv->size());
10823 assert(!
"incorrect call to bm_receivent_event_rpc()");
10896#ifdef LOCAL_ROUTINES
10976#ifdef LOCAL_ROUTINES
10987 return bm_read_buffer(pbuf, buffer_handle, (
void **) ppevent, NULL, NULL, NULL, timeout_msec, convert_flags,
FALSE);
11054#ifdef LOCAL_ROUTINES
11065 return bm_read_buffer(pbuf, buffer_handle, NULL, NULL, NULL, pvec, timeout_msec, convert_flags,
FALSE);
11072#ifdef LOCAL_ROUTINES
11119#ifdef LOCAL_ROUTINES
11135#ifdef LOCAL_ROUTINES
11161 std::vector<BUFFER*> mybuffers;
11167 for (
size_t i = 0;
i < mybuffers.size();
i++) {
11196#ifdef LOCAL_ROUTINES
11212 std::vector<BUFFER*> mybuffers;
11219 for (
size_t idx = 0; idx < mybuffers.size(); idx++) {
11220 BUFFER* pbuf = mybuffers[idx];
11353 if (convert_flags) {
11388 std::vector<char>
vec;
11392 bool locked =
true;
11394 for (
size_t i = 0;
i <
n;
i++) {
11418 dispatched_something =
TRUE;
11427 cm_msg(
MERROR,
"bm_poll_event",
"received event was truncated, buffer size %d is too small, see messages and increase /Experiment/MAX_EVENT_SIZE in ODB", (
int)
vec.size());
11450 if (dispatched_something)
11485#ifdef LOCAL_ROUTINES
11487 std::vector<BUFFER*> mybuffers;
11494 for (
BUFFER* pbuf : mybuffers) {
11497 if (!pbuf->attached)
11511#ifndef DOXYGEN_SHOULD_SKIP_THIS
11513#define MAX_DEFRAG_EVENTS 10
11569 "Received new event with ID %d while old fragments were not completed",
11570 (pevent->event_id & 0x0FFF));
11580 "Not enough defragment buffers, please increase MAX_DEFRAG_EVENTS and recompile");
11587 "Received first event fragment with %d bytes instead of %d bytes, event ignored",
11600 cm_msg(
MERROR,
"bm_defragement_event",
"Not enough memory to allocate event defragment buffer");
11622 "Received fragment without first fragment (ID %d) Ser#:%d",
11633 "Received fragments with more data (%d) than event size (%d)",
11686 printf(
"index %d, client \"%s\", host \"%s\", port %d, socket %d, connected %d, timeout %d",
11817 *convert_flags = 0;
11845 unsigned short int lo, hi;
11848 lo = *((
short int *) (var) + 1);
11849 hi = *((
short int *) (var));
11855 *((
short int *) (var) + 1) = hi;
11856 *((
short int *) (var)) = lo;
11860 unsigned short int lo, hi;
11863 lo = *((
short int *) (var) + 1);
11864 hi = *((
short int *) (var));
11870 *((
short int *) (var) + 1) = hi;
11871 *((
short int *) (var)) = lo;
11876 unsigned short int i1, i2, i3, i4;
11879 i1 = *((
short int *) (var) + 3);
11880 i2 = *((
short int *) (var) + 2);
11881 i3 = *((
short int *) (var) + 1);
11882 i4 = *((
short int *) (var));
11888 *((
short int *) (var) + 3) = i4;
11889 *((
short int *) (var) + 2) = i3;
11890 *((
short int *) (var) + 1) = i2;
11891 *((
short int *) (var)) = i1;
11895 unsigned short int i1, i2, i3, i4;
11898 i1 = *((
short int *) (var) + 3);
11899 i2 = *((
short int *) (var) + 2);
11900 i3 = *((
short int *) (var) + 1);
11901 i4 = *((
short int *) (var));
11907 *((
short int *) (var) + 3) = i4;
11908 *((
short int *) (var) + 2) = i3;
11909 *((
short int *) (var) + 1) = i2;
11910 *((
short int *) (var)) = i1;
11972 if (single_size == 0)
11975 int n = total_size / single_size;
11977 for (
int i = 0;
i <
n;
i++) {
11978 char* p = (
char *)
data + (
i * single_size);
12001 return "<unknown>";
12008 return "<unknown>";
12062 for (
int i = 0; new_list[
i].
id != 0;
i++) {
12067 cm_msg(
MERROR,
"rpc_register_functions",
"registered RPC function with invalid ID %d", new_list[
i].
id);
12074 for (
int i = 0; new_list[
i].
id != 0;
i++) {
12083 for (
int i = 0; new_list[
i].
id != 0;
i++) {
12087 if (
e.dispatch == NULL) {
12100#ifndef DOXYGEN_SHOULD_SKIP_THIS
12188 char net_buffer[256];
12190 int n =
recv_tcp(sock, net_buffer,
sizeof(net_buffer), 0);
12205 struct timeval timeout;
12212 FD_SET(sock, &readfds);
12214 timeout.tv_sec = 0;
12215 timeout.tv_usec = 0;
12217 select(
FD_SETSIZE, &readfds, NULL, NULL, &timeout);
12219 if (FD_ISSET(sock, &readfds)) {
12220 n =
recv_tcp(sock, net_buffer,
sizeof(net_buffer), 0);
12234 }
while (FD_ISSET(sock, &readfds));
12270 bool debug =
false;
12274 cm_msg(
MERROR,
"rpc_client_connect",
"cm_connect_experiment/rpc_set_name not called");
12280 cm_msg(
MERROR,
"rpc_client_connect",
"invalid port %d", port);
12286 static std::mutex gHostnameMutex;
12292 printf(
"rpc_client_connect: host \"%s\", port %d, client \"%s\"\n",
host_name, port, client_name);
12295 printf(
"client connection %d: ", (
int)
i);
12308 bool hostname_locked =
false;
12313 if (
c &&
c->connected) {
12315 if (!hostname_locked) {
12316 gHostnameMutex.lock();
12317 hostname_locked =
true;
12320 if ((
c->host_name ==
host_name) && (
c->port == port)) {
12324 gHostnameMutex.unlock();
12325 hostname_locked =
false;
12326 std::lock_guard<std::mutex> cguard(
c->mutex);
12328 if (
c->connected) {
12333 *hConnection =
c->index;
12335 printf(
"already connected: ");
12351 if (hostname_locked) {
12352 gHostnameMutex.unlock();
12353 hostname_locked =
false;
12359 static int last_reused = 0;
12362 for (
int j = 1;
j < size;
j++) {
12363 int i = (last_reused +
j) % size;
12367 printf(
"last reused %d, reusing slot %d: ", last_reused, (
int)
i);
12386 printf(
"new connection appended to array: ");
12393 c->connected =
true;
12402 std::string errmsg;
12407 cm_msg(
MERROR,
"rpc_client_connect",
"cannot connect to \"%s\" port %d: %s",
host_name, port, errmsg.c_str());
12412 gHostnameMutex.lock();
12417 gHostnameMutex.unlock();
12419 c->client_name = client_name;
12424 setsockopt(
c->send_sock, IPPROTO_TCP, TCP_NODELAY, (
char *) &
i,
sizeof(
i));
12432 std::string cstr =
msprintf(
"%d %s %s %s", hw_type,
cm_get_version(), local_prog_name.c_str(), local_host_name.c_str());
12434 int size = cstr.length() + 1;
12435 i = send(
c->send_sock, cstr.c_str(), size, 0);
12436 if (
i < 0 ||
i != size) {
12437 cm_msg(
MERROR,
"rpc_client_connect",
"cannot send %d bytes, send() returned %d, errno %d (%s)", size,
i, errno, strerror(errno));
12442 bool restore_watchdog_timeout =
false;
12443 BOOL watchdog_call;
12444 DWORD watchdog_timeout;
12450 restore_watchdog_timeout =
true;
12459 if (restore_watchdog_timeout) {
12464 cm_msg(
MERROR,
"rpc_client_connect",
"timeout waiting for server reply");
12470 int remote_hw_type = 0;
12471 char remote_version[32];
12472 remote_version[0] = 0;
12473 sscanf(
str,
"%d %s", &remote_hw_type, remote_version);
12475 c->remote_hw_type = remote_hw_type;
12479 mstrlcpy(v1, remote_version,
sizeof(v1));
12480 if (strchr(v1,
'.'))
12481 if (strchr(strchr(v1,
'.') + 1,
'.'))
12482 *strchr(strchr(v1,
'.') + 1,
'.') = 0;
12485 if (strchr(
str,
'.'))
12486 if (strchr(strchr(
str,
'.') + 1,
'.'))
12487 *strchr(strchr(
str,
'.') + 1,
'.') = 0;
12489 if (strcmp(v1,
str) != 0) {
12490 cm_msg(
MERROR,
"rpc_client_connect",
"remote MIDAS version \'%s\' differs from local version \'%s\'", remote_version,
cm_get_version());
12493 c->connected =
true;
12495 *hConnection =
c->index;
12513 for (
i = 0;
i < MAX_RPC_CONNECTION;
i++)
12514 if (_client_connection[
i].send_sock != 0)
12515 printf(
"slot %d, checking client %s socket %d, connected %d\n",
i, _client_connection[
i].client_name, _client_connection[
i].send_sock, _client_connection[
i].connected);
12523 if (
c &&
c->connected) {
12524 std::lock_guard<std::mutex> cguard(
c->mutex);
12526 if (!
c->connected) {
12539 FD_SET(
c->send_sock, &readfds);
12541 struct timeval timeout;
12542 timeout.tv_sec = 0;
12543 timeout.tv_usec = 0;
12552 }
while (
status == -1 && errno == EINTR);
12555 if (!FD_ISSET(
c->send_sock, &readfds)) {
12562 status = recv(
c->send_sock, (
char *) buffer,
sizeof(buffer), MSG_PEEK);
12567 if (errno == EAGAIN) {
12574 "RPC client connection to \"%s\" on host \"%s\" is broken, recv() errno %d (%s)",
12575 c->client_name.c_str(),
12576 c->host_name.c_str(),
12577 errno, strerror(errno));
12580 }
else if (
status == 0) {
12585 cm_msg(
MINFO,
"rpc_client_check",
"RPC client connection to \"%s\" on host \"%s\" unexpectedly closed",
c->client_name.c_str(),
c->host_name.c_str());
12644 INT remote_hw_type, hw_type;
12645 char str[200], version[32], v1[32];
12647 struct timeval timeout;
12656 if (WSAStartup(MAKEWORD(1, 1), &WSAData) != 0)
12670 cm_msg(
MERROR,
"rpc_server_connect",
"cm_connect_experiment/rpc_set_name not called");
12682 bool listen_localhost =
false;
12684 if (strcmp(
host_name,
"localhost") == 0)
12685 listen_localhost =
true;
12687 int lsock1, lport1;
12688 int lsock2, lport2;
12689 int lsock3, lport3;
12691 std::string errmsg;
12696 cm_msg(
MERROR,
"rpc_server_connect",
"cannot create listener socket: %s", errmsg.c_str());
12703 cm_msg(
MERROR,
"rpc_server_connect",
"cannot create listener socket: %s", errmsg.c_str());
12710 cm_msg(
MERROR,
"rpc_server_connect",
"cannot create listener socket: %s", errmsg.c_str());
12716 s = strchr(
str,
':');
12719 port = strtoul(s + 1, NULL, 0);
12727 cm_msg(
MERROR,
"rpc_server_connect",
"cannot connect to mserver on host \"%s\" port %d: %s",
str, port, errmsg.c_str());
12733 sprintf(
str,
"C %d %d %d %s Default", lport1, lport2, lport3,
cm_get_version());
12737 send(sock,
str, strlen(
str) + 1, 0);
12741 cm_msg(
MERROR,
"rpc_server_connect",
"timeout on receive status from server");
12745 status = version[0] = 0;
12746 sscanf(
str,
"%d %s", &
status, version);
12754 strcpy(v1, version);
12755 if (strchr(v1,
'.'))
12756 if (strchr(strchr(v1,
'.') + 1,
'.'))
12757 *strchr(strchr(v1,
'.') + 1,
'.') = 0;
12760 if (strchr(
str,
'.'))
12761 if (strchr(strchr(
str,
'.') + 1,
'.'))
12762 *strchr(strchr(
str,
'.') + 1,
'.') = 0;
12764 if (strcmp(v1,
str) != 0) {
12765 cm_msg(
MERROR,
"rpc_server_connect",
"remote MIDAS version \'%s\' differs from local version \'%s\'", version,
12771 FD_SET(lsock1, &readfds);
12772 FD_SET(lsock2, &readfds);
12773 FD_SET(lsock3, &readfds);
12776 timeout.tv_usec = 0;
12787 if (!FD_ISSET(lsock1, &readfds)) {
12788 cm_msg(
MERROR,
"rpc_server_connect",
"mserver subprocess could not be started (check path)");
12800 cm_msg(
MERROR,
"rpc_server_connect",
"accept() failed");
12814 flag = 2 * 1024 * 1024;
12817 cm_msg(
MERROR,
"rpc_server_connect",
"cannot setsockopt(SOL_SOCKET, SO_SNDBUF), errno %d (%s)", errno, strerror(errno));
12822 sprintf(
str,
"%d %s", hw_type, local_prog_name.c_str());
12829 cm_msg(
MERROR,
"rpc_server_connect",
"timeout on receive remote computer info");
12833 sscanf(
str,
"%d", &remote_hw_type);
12850 if (
c &&
c->connected) {
12853 if (!
c->connected) {
12873 if (
c &&
c->connected) {
12890 if (!
c->connected) {
12959 static int rpc_server_disconnect_recursion_level = 0;
12961 if (rpc_server_disconnect_recursion_level)
12964 rpc_server_disconnect_recursion_level = 1;
12989 rpc_server_disconnect_recursion_level = 0;
13081 INT tmp_type, size;
13099 dummy = 0x12345678;
13100 p = (
unsigned char *) &dummy;
13103 else if (*p == 0x12)
13106 cm_msg(
MERROR,
"rpc_get_option",
"unknown byte order format");
13109 f = (float) 1.2345;
13111 memcpy(&dummy, &f,
sizeof(f));
13112 if ((dummy & 0xFF) == 0x19 &&
13113 ((dummy >> 8) & 0xFF) == 0x04 && ((dummy >> 16) & 0xFF) == 0x9E
13114 && ((dummy >> 24) & 0xFF) == 0x3F)
13116 else if ((dummy & 0xFF) == 0x9E &&
13117 ((dummy >> 8) & 0xFF) == 0x40 && ((dummy >> 16) & 0xFF) == 0x19
13118 && ((dummy >> 24) & 0xFF) == 0x04)
13121 cm_msg(
MERROR,
"rpc_get_option",
"unknown floating point format");
13123 d = (double) 1.2345;
13125 memcpy(&dummy, &
d,
sizeof(f));
13126 if ((dummy & 0xFF) == 0x8D &&
13127 ((dummy >> 8) & 0xFF) == 0x97 && ((dummy >> 16) & 0xFF) == 0x6E
13128 && ((dummy >> 24) & 0xFF) == 0x12)
13130 else if ((dummy & 0xFF) == 0x83 &&
13131 ((dummy >> 8) & 0xFF) == 0xC0 && ((dummy >> 16) & 0xFF) == 0xF3
13132 && ((dummy >> 24) & 0xFF) == 0x3F)
13134 else if ((dummy & 0xFF) == 0x13 &&
13135 ((dummy >> 8) & 0xFF) == 0x40 && ((dummy >> 16) & 0xFF) == 0x83
13136 && ((dummy >> 24) & 0xFF) == 0xC0)
13138 else if ((dummy & 0xFF) == 0x9E &&
13139 ((dummy >> 8) & 0xFF) == 0x40 && ((dummy >> 16) & 0xFF) == 0x18
13140 && ((dummy >> 24) & 0xFF) == 0x04)
13142 "MIDAS cannot handle VAX D FLOAT format. Please compile with the /g_float flag");
13144 cm_msg(
MERROR,
"rpc_get_option",
"unknown floating point format");
13168 else if (hConn == -2)
13185 setsockopt(
c->send_sock, IPPROTO_TCP, TCP_NODELAY, (
char *) &
value,
sizeof(
value));
13215 int timeout =
c->rpc_timeout;
13236 if (old_timeout_msec)
13240 if (old_timeout_msec)
13246 if (old_timeout_msec)
13247 *old_timeout_msec =
c->rpc_timeout;
13248 c->rpc_timeout = timeout_msec;
13251 if (old_timeout_msec)
13252 *old_timeout_msec = 0;
13260#ifndef DOXYGEN_SHOULD_SKIP_THIS
13413 va_start(argptr, format);
13414 vsprintf(
str, (
char *) format, argptr);
13427 switch (arg_type) {
13436 *((
int *) arg) = va_arg(*arg_ptr,
int);
13441 *((
INT *) arg) = va_arg(*arg_ptr,
INT);
13450 *((
float *) arg) = (float) va_arg(*arg_ptr,
double);
13454 *((
double *) arg) = va_arg(*arg_ptr,
double);
13458 *((
char **) arg) = va_arg(*arg_ptr,
char *);
13468 bool debug =
false;
13471 printf(
"encode rpc_id %d \"%s\"\n", rl.
id, rl.
name);
13476 printf(
"i=%d, tid %d, flags 0x%x, n %d\n",
i, tid, flags,
n);
13505 size_t buf_size =
sizeof(
NET_COMMAND) + 4 * 1024;
13506 char* buf = (
char *)malloc(buf_size);
13514 char* param_ptr = (*nc)->param;
13538 char* arg = args[
i];
13561 arg_size = 1 + strlen((
char *) *((
char **) arg));
13574 const char* arg_tmp = args[
i+1];
13578 arg_size = *((
INT *) *((
void **) arg_tmp));
13580 arg_size = *((
INT *) arg_tmp);
13582 *((
INT *) param_ptr) =
ALIGN8(arg_size);
13592 int param_size =
ALIGN8(arg_size);
13595 size_t param_offset = (
char *) param_ptr - (
char *)(*nc);
13597 if (param_offset + param_size + 16 > buf_size) {
13598 size_t new_size = param_offset + param_size + 1024;
13600 buf = (
char *) realloc(buf, new_size);
13602 buf_size = new_size;
13604 param_ptr = buf + param_offset;
13610 printf(
"encode param %d, flags 0x%x, tid %d, arg_type %d, arg_size %d, param_size %d, memcpy pointer %d\n",
i, flags, tid, arg_type, arg_size, param_size, arg_size);
13612 memcpy(param_ptr, (
void *) *((
void **) arg), arg_size);
13615 printf(
"encode param %d, flags 0x%x, tid %d, arg_type %d, arg_size %d, param_size %d, double->float\n",
i, flags, tid, arg_type, arg_size, param_size);
13618 *((
float *) param_ptr) = (float) *((
double *) arg);
13621 printf(
"encode param %d, flags 0x%x, tid %d, arg_type %d, arg_size %d, param_size %d, memcpy %d\n",
i, flags, tid, arg_type, arg_size, param_size, arg_size);
13623 memcpy(param_ptr, arg, arg_size);
13626 param_ptr += param_size;
13633 printf(
"encode rpc_id %d \"%s\" buf_size %d, param_size %d\n", rl.
id, rl.
name, (
int)buf_size, (*nc)->header.param_size);
13641 bool debug =
false;
13644 printf(
"decode reply to rpc_id %d \"%s\" has %d bytes\n", rl.
id, rl.
name, (
int)buf_size);
13648 const char* param_ptr = buf;
13672 if (param_ptr == NULL) {
13673 cm_msg(
MERROR,
"rpc_call_decode",
"routine \"%s\": no data in RPC reply, needed to decode an RPC_OUT parameter. param_ptr is NULL", rl.
name);
13681 arg_size = strlen((
char *) (param_ptr)) + 1;
13684 arg_size = *((
INT *) param_ptr);
13692 int param_size =
ALIGN8(arg_size);
13695 if (*((
char **) arg)) {
13697 printf(
"decode param %d, flags 0x%x, tid %d, arg_type %d, arg_size %d, param_size %d, memcpy %d\n",
i, flags, tid, arg_type, arg_size, param_size, arg_size);
13698 memcpy((
void *) *((
char **) arg), param_ptr, arg_size);
13701 param_ptr += param_size;
13713 bool debug =
false;
13722 printf(
"encode rpc_id %d \"%s\"\n", rl.
id, rl.
name);
13727 printf(
"param %2d, tid %2d, flags 0x%02x, n %3d\n",
i, tid, flags,
n);
13756 size_t buf_size =
sizeof(
NET_COMMAND) + 4 * 1024;
13757 char* buf = (
char *)malloc(buf_size);
13765 char* param_ptr = (*nc)->param;
13789 char* arg = args[
i];
13810 void* parg = (
void *) *((
void **) arg);
13814 std::string* s = (std::string*)parg;
13815 arg_size = 1 + s->length();
13816 parg = (
void*)s->c_str();
13820 arg_size = 1 + strlen((
char *) *((
char **) arg));
13834 std::vector<char>* pv = (std::vector<char>*)parg;
13835 arg_size = pv->size();
13836 parg = (
void*)pv->data();
13840 *((
INT *) param_ptr) = arg_size;
13843 const char* arg_tmp = args[
i+1];
13847 arg_size = *((
INT *) *((
void **) arg_tmp));
13849 arg_size = *((
INT *) arg_tmp);
13852 *((
INT *) param_ptr) =
ALIGN8(arg_size);
13864 int param_size =
ALIGN8(arg_size);
13867 size_t param_offset = (
char *) param_ptr - (
char *)(*nc);
13869 if (param_offset + param_size + 16 > buf_size) {
13870 size_t new_size = param_offset + param_size + 1024;
13872 buf = (
char *) realloc(buf, new_size);
13874 buf_size = new_size;
13876 param_ptr = buf + param_offset;
13882 printf(
"encode param %2d, flags 0x%02x, tid %2d, arg_type %2d, arg_size %3d, param_size %3d, memcpy pointer %3d\n",
i, flags, tid, arg_type, arg_size, param_size, arg_size);
13884 memcpy(param_ptr, parg, arg_size);
13887 printf(
"encode param %2d, flags 0x%02x, tid %2d, arg_type %2d, arg_size %3d, param_size %3d, double->float\n",
i, flags, tid, arg_type, arg_size, param_size);
13890 *((
float *) param_ptr) = (float) *((
double *) arg);
13893 printf(
"encode param %2d, flags 0x%02x, tid %2d, arg_type %2d, arg_size %3d, param_size %3d, memcpy %3d\n",
i, flags, tid, arg_type, arg_size, param_size, arg_size);
13895 memcpy(param_ptr, arg, arg_size);
13898 param_ptr += param_size;
13905 printf(
"encode rpc_id %d \"%s\" buf_size %d, param_size %d\n", rl.
id, rl.
name, (
int)buf_size, (*nc)->header.param_size);
13913 bool debug =
false;
13922 printf(
"decode reply to rpc_id %d \"%s\" has %d bytes\n", rl.
id, rl.
name, (
int)buf_size);
13926 const char* param_ptr = buf;
13945 char arg[
sizeof(double)+
sizeof(uint64_t)+
sizeof(
char*)];
13950 if (param_ptr == NULL) {
13951 cm_msg(
MERROR,
"rpc_call_decode_cxx",
"routine \"%s\": no data in RPC reply, needed to decode an RPC_OUT parameter. param_ptr is NULL", rl.
name);
13958 arg_size = strlen((
char *) (param_ptr)) + 1;
13961 printf(
"decode param %2d, flags 0x%02x, tid %2d, arg_type %2d, arg_size %3d, string [%s]\n",
i, flags, tid, arg_type, arg_size, (
char *) (param_ptr));
13965 arg_size = *((
INT *) param_ptr);
13973 int param_size =
ALIGN8(arg_size);
13976 void* parg = *(
char**) arg;
13980 printf(
"decode param %2d, flags 0x%02x, tid %2d, arg_type %2d, arg_size %3d, param_size %3d, assign %3d to std::string at %p, offset %zu, string [%s]\n",
i, flags, tid, arg_type, arg_size, param_size, arg_size, parg, param_ptr - buf, param_ptr);
13981 *(std::string*)parg = param_ptr;
13984 printf(
"decode param %2d, flags 0x%02x, tid %2d, arg_type %2d, arg_size %3d, param_size %3d, assign %3d to std::vector at %p, offset %zu\n",
i, flags, tid, arg_type, arg_size, param_size, arg_size, parg, param_ptr - buf);
13985 std::vector<char>* pvec = (std::vector<char>*)parg;
13987 pvec->insert(pvec->end(), param_ptr, param_ptr + arg_size);
13990 printf(
"decode param %2d, flags 0x%02x, tid %2d, arg_type %2d, arg_size %3d, param_size %3d, memcpy %3d to %p, offset %zu\n",
i, flags, tid, arg_type, arg_size, param_size, arg_size, parg, param_ptr - buf);
13991 memcpy(parg, param_ptr, arg_size);
13995 param_ptr += param_size;
14057 cm_msg(
MERROR,
"rpc_client_call",
"invalid rpc connection handle %d", hConn);
14068 routine_id &= ~RPC_NO_REPLY;
14078 bool rpc_cxx =
false;
14083 cm_msg(
MERROR,
"rpc_client_call",
"call to \"%s\" on \"%s\" with invalid RPC ID %d",
c->client_name.c_str(),
c->host_name.c_str(), routine_id);
14088 const char *rpc_name = rpc_entry.
name;
14094 va_start(ap, routine_id);
14111 if (rpc_no_reply) {
14112 i =
send_tcp(
c->send_sock, (
char *) nc, send_size, 0);
14114 if (
i != send_size) {
14115 cm_msg(
MERROR,
"rpc_client_call",
"call to \"%s\" on \"%s\" RPC \"%s\": send_tcp() failed",
c->client_name.c_str(),
c->host_name.c_str(), rpc_name);
14135 i =
send_tcp(
c->send_sock, (
char *) nc, send_size, 0);
14136 if (
i != send_size) {
14137 cm_msg(
MERROR,
"rpc_client_call",
"call to \"%s\" on \"%s\" RPC \"%s\": send_tcp() failed",
c->client_name.c_str(),
c->host_name.c_str(), rpc_name);
14145 bool restore_watchdog_timeout =
false;
14146 BOOL watchdog_call;
14147 DWORD watchdog_timeout;
14152 if (
c->rpc_timeout >= (
int) watchdog_timeout) {
14153 restore_watchdog_timeout =
true;
14157 DWORD rpc_status = 0;
14158 DWORD buf_size = 0;
14164 if (restore_watchdog_timeout) {
14169 cm_msg(
MERROR,
"rpc_client_call",
"call to \"%s\" on \"%s\" RPC \"%s\": timeout waiting for reply",
c->client_name.c_str(),
c->host_name.c_str(), rpc_name);
14177 cm_msg(
MERROR,
"rpc_client_call",
"call to \"%s\" on \"%s\" RPC \"%s\": error, ss_recv_net_command() status %d",
c->client_name.c_str(),
c->host_name.c_str(), rpc_name,
status);
14187 cm_msg(
MERROR,
"rpc_client_call",
"call to \"%s\" on \"%s\" RPC \"%s\": error, unknown RPC, status %d",
c->client_name.c_str(),
c->host_name.c_str(), rpc_name, rpc_status);
14195 va_start(ap, routine_id);
14246 routine_id &= ~RPC_NO_REPLY;
14255 fprintf(stderr,
"rpc_call(routine_id=%d) failed, no connection to mserver.\n", routine_id);
14273 bool rpc_cxx =
false;
14279 cm_msg(
MERROR,
"rpc_call",
"invalid rpc ID (%d)", routine_id);
14283 const char* rpc_name = rpc_entry.
name;
14290 va_start(ap, routine_id);
14307 if (rpc_no_reply) {
14308 i =
send_tcp(send_sock, (
char *) nc, send_size, 0);
14310 if (
i != send_size) {
14312 cm_msg(
MERROR,
"rpc_call",
"rpc \"%s\" error: send_tcp() failed", rpc_name);
14323 i =
send_tcp(send_sock, (
char *) nc, send_size, 0);
14324 if (
i != send_size) {
14326 cm_msg(
MERROR,
"rpc_call",
"rpc \"%s\" error: send_tcp() failed", rpc_name);
14334 bool restore_watchdog_timeout =
false;
14335 BOOL watchdog_call;
14336 DWORD watchdog_timeout;
14346 if (rpc_timeout >= (
int) watchdog_timeout) {
14347 restore_watchdog_timeout =
true;
14352 DWORD rpc_status = 0;
14353 DWORD buf_size = 0;
14358 if (restore_watchdog_timeout) {
14369 cm_msg(
MERROR,
"rpc_call",
"routine \"%s\": timeout waiting for reply, program abort", rpc_name);
14377 cm_msg(
MERROR,
"rpc_call",
"routine \"%s\": error, ss_recv_net_command() status %d, program abort", rpc_name,
status);
14385 cm_msg(
MERROR,
"rpc_call",
"routine \"%s\": error, unknown RPC, status %d", rpc_name, rpc_status);
14393 va_start(ap, routine_id);
14456 return bm_send_event(buffer_handle, pevent, unused, async_flag);
14478 cm_msg(
MERROR,
"rpc_send_event_sg",
"invalid sg_n %d", sg_n);
14482 if (sg_ptr[0] == NULL) {
14483 cm_msg(
MERROR,
"rpc_send_event_sg",
"invalid sg_ptr[0] is NULL");
14488 cm_msg(
MERROR,
"rpc_send_event_sg",
"invalid sg_len[0] value %d is smaller than event header size %d", (
int)sg_len[0], (
int)
sizeof(
EVENT_HEADER));
14494 const DWORD MAX_DATA_SIZE = (0x7FFFFFF0 - 16);
14497 if (data_size == 0) {
14498 cm_msg(
MERROR,
"rpc_send_event_sg",
"invalid event data size zero");
14502 if (data_size > MAX_DATA_SIZE) {
14503 cm_msg(
MERROR,
"rpc_send_event_sg",
"invalid event data size %d (0x%x) maximum is %d (0x%x)", data_size, data_size, MAX_DATA_SIZE, MAX_DATA_SIZE);
14511 for (
int i=0;
i<sg_n;
i++) {
14516 cm_msg(
MERROR,
"rpc_send_event_sg",
"data size mismatch: event data_size %d, event_size %d not same as sum of sg_len %d", (
int)data_size, (
int)
event_size, (
int)
count);
14542 assert(
sizeof(
DWORD) == 4);
14543 DWORD bh_buf = buffer_handle;
14548 cm_msg(
MERROR,
"rpc_send_event_sg",
"ss_write_tcp(buffer handle) failed, event socket is now closed");
14554 for (
int i=0;
i<sg_n;
i++) {
14558 cm_msg(
MERROR,
"rpc_send_event_sg",
"ss_write_tcp(event data) failed, event socket is now closed");
14565 if (
count < total_size) {
14566 char padding[8] = { 0,0,0,0,0,0,0,0 };
14567 size_t padlen = total_size -
count;
14568 assert(padlen < 8);
14572 cm_msg(
MERROR,
"rpc_send_event_sg",
"ss_write_tcp(padding) failed, event socket is now closed");
14638 for (
size_t i = 0;
i <
n;
i++) {
14661 cm_msg(
MERROR,
"rpc_transition_dispatch",
"no handler for transition %d with sequence number %d",
CINT(0),
CINT(4));
14664 cm_msg(
MERROR,
"rpc_transition_dispatch",
"received unrecognized command %d", idx);
14720void debug_dump(
unsigned char *p,
int size)
14725 for (
i = 0;
i < (size - 1) / 16 + 1;
i++) {
14726 printf(
"%p ", p +
i * 16);
14727 for (
j = 0;
j < 16;
j++)
14728 if (
i * 16 +
j < size)
14729 printf(
"%02X ", p[
i * 16 +
j]);
14734 for (
j = 0;
j < 16;
j++) {
14736 if (
i * 16 +
j < size)
14737 printf(
"%c", (
c >= 32 &&
c < 128) ? p[
i * 16 +
j] :
'.');
14776 char *buffer = NULL;
14798 int param_size = -1;
14807 if (param_size == -1) {
14827 char *p = (
char *) realloc(*pbuf, new_size);
14830 cm_msg(
MERROR,
"recv_net_command_realloc",
"cannot reallocate buffer from %d bytes to %d bytes", *pbufsize, new_size);
14836 *pbufsize = new_size;
14847 int size = write_ptr - read_ptr;
14850 memcpy(buffer + copied, net_buffer + read_ptr, size);
14852 read_ptr = write_ptr;
14856 write_ptr = recv(sock, net_buffer + misalign, sa->
net_buffer_size - 8, 0);
14859 }
while (write_ptr == -1 && errno == EINTR);
14861 write_ptr = recv(sock, net_buffer + misalign, sa->
net_buffer_size - 8, 0);
14865 if (write_ptr <= 0) {
14866 if (write_ptr == 0)
14867 cm_msg(
MERROR,
"recv_net_command_realloc",
"rpc connection from \'%s\' on \'%s\' unexpectedly closed", sa->
prog_name.c_str(), sa->
host_name.c_str());
14869 cm_msg(
MERROR,
"recv_net_command_realloc",
"recv() returned %d, errno: %d (%s)", write_ptr, errno, strerror(errno));
14877 read_ptr = misalign;
14878 write_ptr += misalign;
14880 misalign = write_ptr % 8;
14885 memcpy(buffer + copied, net_buffer + read_ptr, size);
14890 if (write_ptr - read_ptr < param_size)
14893 *remaining = write_ptr - read_ptr;
14900 return size + copied;
14964 char header_buf[header_size];
14975 int hrd =
recv_tcp2(sock, header_buf, header_size, 1);
14984 cm_msg(
MERROR,
"recv_event_server",
"recv_tcp2(header) returned %d", hrd);
14988 if (hrd < (
int) header_size) {
14989 int hrd1 =
recv_tcp2(sock, header_buf + hrd, header_size - hrd, 0);
14993 cm_msg(
MERROR,
"recv_event_server",
"recv_tcp2(more header) returned %d", hrd1);
15001 if (hrd != (
int) header_size) {
15002 cm_msg(
MERROR,
"recv_event_server",
"recv_tcp2(header) returned %d instead of %d", hrd, (
int) header_size);
15006 INT *pbh = (
INT *) header_buf;
15025 for (
int i=0;
i<5;
i++) {
15026 printf(
"recv_event_server: header[%d]: 0x%08x\n",
i, pbh[
i]);
15034 "received event header with invalid data_size %d: event_size %d, total_size %d", pevent->
data_size,
15042 int bufsize =
sizeof(
INT) + total_size;
15047 if (*pbuffer_size < bufsize) {
15048 int newsize = 1024 +
ALIGN8(bufsize);
15052 char *newbuf = (
char *) realloc(*pbuffer, newsize);
15053 if (newbuf == NULL) {
15054 cm_msg(
MERROR,
"recv_event_server",
"cannot realloc() event buffer from %d to %d bytes", *pbuffer_size,
15059 *pbuffer_size = newsize;
15064 memcpy(*pbuffer, header_buf, header_size);
15068 int to_read =
sizeof(
INT) + total_size - header_size;
15069 int rptr = header_size;
15072 int drd =
recv_tcp2(sock, (*pbuffer) + rptr, to_read, 0);
15076 cm_msg(
MERROR,
"recv_event_server",
"recv_tcp2(data) returned %d instead of %d", drd, to_read);
15157 std::string errmsg;
15162 cm_msg(
MERROR,
"rpc_register_server",
"cannot listen to tcp port %d: %s", port, errmsg.c_str());
15167#if defined(F_SETFD) && defined(FD_CLOEXEC)
15168 status = fcntl(lsock, F_SETFD, fcntl(lsock, F_GETFD) | FD_CLOEXEC);
15170 cm_msg(
MERROR,
"rpc_register_server",
"fcntl(F_SETFD, FD_CLOEXEC) failed, errno %d (%s)", errno, strerror(errno));
15227 char *in_param_ptr, *out_param_ptr, *last_param_ptr;
15230 INT param_size, max_size;
15231 void *prpc_param[20];
15232 char debug_line[1024], *return_buffer;
15233 int return_buffer_size;
15234 int return_buffer_tls;
15238 int initial_buffer_size = 1024;
15261 return_buffer_tls =
i;
15264 assert(return_buffer);
15274 if (convert_flags) {
15290 assert(xroutine_id == routine_id);
15294 in_param_ptr = nc_in->
param;
15297 out_param_ptr = nc_out->
param;
15299 sprintf(debug_line,
"%s(", rl.
name);
15309 param_size =
ALIGN8(1 + strlen((
char *) (in_param_ptr)));
15313 param_size = *((
INT *) in_param_ptr);
15316 param_size =
ALIGN8(param_size);
15324 prpc_param[
i] = in_param_ptr;
15327 if (convert_flags) {
15338 if (strlen(debug_line) +
str.length() + 2 <
sizeof(debug_line)) {
15339 strcat(debug_line,
"\"");
15340 strcat(debug_line,
str.c_str());
15341 strcat(debug_line,
"\"");
15343 strcat(debug_line,
"...");
15345 strcat(debug_line,
str.c_str());
15347 in_param_ptr += param_size;
15359 max_size = *((
INT *) in_param_ptr);
15363 max_size =
ALIGN8(max_size);
15365 *((
INT *) out_param_ptr) = max_size;
15371 param_size = max_size;
15377 if ((
POINTER_T) out_param_ptr - (
POINTER_T) nc_out + param_size > return_buffer_size) {
15380 "return parameters (%d) too large for network buffer (%d)",
15390 "rpc_execute: return parameters (%d) too large for network buffer (%d), new buffer size (%d)",
15391 (
int)((
POINTER_T) out_param_ptr - (
POINTER_T) nc_out + param_size), return_buffer_size, new_size);
15394 itls = return_buffer_tls;
15400 cm_msg(
MERROR,
"rpc_execute",
"Cannot allocate return buffer of size %d", new_size);
15406 assert(return_buffer);
15414 memcpy(out_param_ptr, prpc_param[
i], param_size);
15417 strcat(debug_line,
"-");
15419 prpc_param[
i] = out_param_ptr;
15420 out_param_ptr += param_size;
15424 strcat(debug_line,
", ");
15429 strcat(debug_line,
")");
15432 last_param_ptr = out_param_ptr;
15461 out_param_ptr = nc_out->
param;
15469 max_size = *((
INT *) out_param_ptr);
15474 const char* param_ptr = ((
char *) out_param_ptr) +
ALIGN8(
sizeof(
INT));
15476 param_size = strlen(param_ptr) + 1;
15477 param_size =
ALIGN8(param_size);
15480 memmove(out_param_ptr, out_param_ptr +
ALIGN8(
sizeof(
INT)), param_size);
15483 memmove(out_param_ptr + param_size,
15484 out_param_ptr + max_size +
ALIGN8(
sizeof(
INT)),
15490 max_size = *((
INT *) out_param_ptr);
15498 param_size = *((
INT *) prpc_param[
i + 1]);
15499 *((
INT *) out_param_ptr) = param_size;
15505 param_size =
ALIGN8(param_size);
15508 memmove(out_param_ptr + param_size,
15509 out_param_ptr + max_size,
15517 if (convert_flags) {
15527 out_param_ptr += param_size;
15538 if (convert_flags) {
15586 std::vector<char>*
pv = NULL;
15635 bool debug =
false;
15641 if (convert_flags) {
15656 assert(xroutine_id == routine_id);
15674 char* in_param_ptr = (
char*)nc_in->
param;
15677 printf(
"rpc_execute_cxx: routine_id %d, name \"%s\"\n", routine_id, rl.
name);
15684 std::vector<RPE> params;
15691 params.resize(nparams);
15693 size_t in_offset = 0;
15695 for (
int i = 0;
i < nparams;
i++) {
15696 in_param_size[
i] = 0;
15697 in_param_offset[
i] = 0;
15706 arg_size = 1 + strlen((
char *) (in_param_ptr));
15711 int arg_size_align8 = *((
INT *) in_param_ptr);
15718 arg_size = arg_size_align8;
15721 arg_size = *((
INT *) (((
char*)in_param_ptr) +
ALIGN8(arg_size_align8)));
15728 cm_msg(
MERROR,
"rpc_execute_cxx",
"RPC %d, param %d tid %d flags 0x%x size mismatch: header %d vs next param %d", routine_id,
i, tid, flags, arg_size_align8, arg_size);
15737 int param_size =
ALIGN8(arg_size);
15739 in_param_size[
i] = param_size;
15740 in_param_offset[
i] = in_offset;
15742 params[
i].offset = in_offset;
15743 params[
i].arg_size = arg_size;
15744 params[
i].param_size = param_size;
15747 if (convert_flags) {
15755 in_param_ptr += param_size;
15756 in_offset += param_size;
15763 params[
i].out_max_size = 0;
15770 params[
i].out_max_size_offset = in_offset;
15772 INT max_size = *((
INT *) in_param_ptr);
15777 if (max_size < 0 || (tid ==
TID_STRING && max_size == 0)) {
15778 cm_msg(
MERROR,
"rpc_execute_cxx",
"RPC %d, param %d tid %d flags 0x%x invalid maximum output size %d", routine_id,
i, tid, flags, max_size);
15782 params[
i].out_max_size = max_size;
15786 params[
i].out_max_size = rl.
param[
i].
n;
15792 params[
i].ps =
new std::string;
15794 *(params[
i].ps) = (
char*)nc_in->
param + in_param_offset[
i];
15797 prpc_param[
i] = (
void*) params[
i].ps;
15799 params[
i].pv =
new std::vector<char>;
15801 params[
i].pv->insert(params[
i].pv->end(), (
char*)nc_in->
param + in_param_offset[
i], (
char*)nc_in->
param + in_param_offset[
i] + params[
i].arg_size);
15804 prpc_param[
i] = (
void*) params[
i].pv;
15806 cm_msg(
MERROR,
"rpc_execute_cxx",
"RPC %d: param %d tid %d flags 0x%x, TID not compatible with flag RPC_CXX", routine_id,
i, tid, flags);
15811 params[
i].pv =
new std::vector<char>;
15812 params[
i].pv->insert(params[
i].pv->end(), (
char*)nc_in->
param + in_param_offset[
i], (
char*)nc_in->
param + in_param_offset[
i] + params[
i].arg_size);
15813 size_t want_size = params[
i].out_max_size;
15815 if (params[
i].pv->size() < want_size)
15816 params[
i].pv->resize(want_size);
15817 prpc_param[
i] = params[
i].pv->data();
15818 }
else if (flags &
RPC_IN) {
15819 prpc_param[
i] = (
char*)nc_in->
param + in_param_offset[
i];
15821 }
else if (flags &
RPC_OUT) {
15822 params[
i].pv =
new std::vector<char>;
15823 params[
i].pv->resize(params[
i].out_max_size);
15824 prpc_param[
i] = params[
i].pv->data();
15830 printf(
"rpc_execute_cxx: param %2d, tid %2d, flags 0x%04x, in %3zu+%-3zu+%-3zu, out max size %3zu at %3zu, ptr %p\n",
i, tid, flags, params[
i].
offset, params[
i].arg_size, params[
i].param_size, params[
i].out_max_size, params[
i].out_max_size_offset, prpc_param[
i]);
15834 printf(
"rpc_execute_cxx: nc_in size %d, in_offset %zu\n", nc_in->
header.
param_size, in_offset);
15839 ok &= in_param_offset[ 0] == 0; ok &= in_param_size[ 0] == 8;
15840 ok &= in_param_offset[ 1] == 0; ok &= in_param_size[ 1] == 0;
15841 ok &= in_param_offset[ 2] == 8; ok &= in_param_size[ 2] == 8;
15842 ok &= in_param_offset[ 3] == 16; ok &= in_param_size[ 3] == 16;
15843 ok &= in_param_offset[ 4] == 0; ok &= in_param_size[ 4] == 0;
15844 ok &= in_param_offset[ 5] == 32; ok &= in_param_size[ 5] == 8;
15845 ok &= in_param_offset[ 6] == 0; ok &= in_param_size[ 6] == 0;
15846 ok &= in_param_offset[ 7] == 40; ok &= in_param_size[ 7] == 8;
15847 ok &= in_param_offset[ 8] == 48; ok &= in_param_size[ 8] == 16;
15848 ok &= in_param_offset[ 9] == 64; ok &= in_param_size[ 9] == 8;
15849 ok &= in_param_offset[10] == 72; ok &= in_param_size[10] == 72;
15850 ok &= in_param_offset[11] == 0; ok &= in_param_size[11] == 0;
15851 ok &= in_param_offset[12] == 144; ok &= in_param_size[12] == 72;
15852 ok &= in_param_offset[13] == 224; ok &= in_param_size[13] == 40;
15853 ok &= in_param_offset[14] == 264; ok &= in_param_size[14] == 8;
15854 ok &= in_param_offset[15] == 280; ok &= in_param_size[15] == 16;
15855 ok &= in_param_offset[16] == 296; ok &= in_param_size[16] == 8;
15856 ok &= in_param_offset[17] == 0; ok &= in_param_size[17] == 0;
15857 ok &= in_param_offset[18] == 304; ok &= in_param_size[18] == 8;
15858 ok &= in_offset == 312;
15861 cm_msg(
MERROR,
"rpc_execute_cxx",
"RPC_TEST2 parameters encoding error!");
15867 printf(
"rpc_execute_cxx: calling dispatch()\n");
15879 printf(
"rpc_execute_cxx: dispatch() status %d\n",
status);
15904 std::vector<char> v_out;
15908 for (
int i = 0;
i < nparams;
i++) {
15915 size_t arg_size = 1 + params[
i].ps->length();
15916 size_t param_size =
ALIGN8(arg_size);
15919 printf(
"rpc_execute_cxx: param %2d, std::string arg_size %zu, param_size %zu, string [%s]\n",
i, arg_size, param_size, params[
i].ps->c_str());
15921 v_out.insert(v_out.end(), params[
i].ps->c_str(), params[
i].ps->c_str() + arg_size);
15922 v_out.resize(v_out.size() + param_size - arg_size);
15924 size_t arg_size = params[
i].pv->size();
15925 size_t param_size =
ALIGN8(arg_size);
15928 printf(
"rpc_execute_cxx: param %2d, std::vector arg_size %zu, param_size %zu\n",
i, arg_size, param_size);
15931 *((
INT *) buf) = arg_size;
15934 v_out.insert(v_out.end(), buf, buf +
ALIGN8(
sizeof(
INT)));
15935 v_out.insert(v_out.end(), params[
i].pv->data(), params[
i].pv->data() + arg_size);
15936 v_out.resize(v_out.size() + param_size - arg_size);
15938 cm_msg(
MERROR,
"rpc_execute_cxx",
"RPC %d: param %d tid %d flags 0x%x, TID not compatible with flag RPC_CXX", routine_id,
i, tid, flags);
15942 size_t convert_offset = 0;
15943 size_t convert_size = 0;
15946 size_t max_size = params[
i].out_max_size;
15947 char* param_ptr = (
char *) prpc_param[
i];
15949 char* string_end = (
char*)memchr(param_ptr, 0, max_size);
15950 size_t arg_size = string_end ? string_end - param_ptr : max_size;
15952 param_ptr[max_size - 1] = 0;
15953 arg_size = max_size - 1;
15956 size_t param_size =
ALIGN8(arg_size);
15959 printf(
"rpc_execute_cxx: param %2d, string max_size %zu, string_size %zu, param_size %zu\n",
i, max_size, arg_size, param_size);
15961 v_out.insert(v_out.end(), param_ptr, param_ptr + arg_size);
15962 v_out.resize(v_out.size() + param_size - arg_size);
15964 size_t max_size = params[
i].out_max_size;
15965 INT arg_size_int = *((
INT *) prpc_param[
i + 1]);
15966 if (arg_size_int < 0 || (
size_t) arg_size_int > max_size) {
15967 cm_msg(
MERROR,
"rpc_execute_cxx",
"RPC %d, param %d output array size %d exceeds maximum size %zu", routine_id,
i, arg_size_int, max_size);
15970 size_t arg_size = arg_size_int;
15971 char* param_ptr = (
char*)prpc_param[
i];
15972 size_t param_size =
ALIGN8(arg_size);
15975 printf(
"rpc_execute_cxx: param %2d, array max_size %zu, param_size %zu\n",
i, max_size, param_size);
15978 *((
INT *) buf) = arg_size;
15981 v_out.insert(v_out.end(), buf, buf +
ALIGN8(
sizeof(
INT)));
15982 convert_offset = v_out.size();
15983 convert_size = arg_size;
15984 v_out.insert(v_out.end(), param_ptr, param_ptr + arg_size);
15985 v_out.resize(v_out.size() + param_size - arg_size);
15987 char* param_ptr = (
char*)prpc_param[
i];
15991 size_t param_size =
ALIGN8(arg_size);
15995 printf(
"rpc_execute_cxx: param %2d, tid %2d, arg_size %zu, param_size %zu, value %d\n",
i, tid, arg_size, param_size, *(
int*)param_ptr);
15997 printf(
"rpc_execute_cxx: param %2d, tid %2d, arg_size %zu, param_size %zu\n",
i, tid, arg_size, param_size);
16001 convert_offset = v_out.size();
16002 convert_size = arg_size;
16003 v_out.insert(v_out.end(), param_ptr, param_ptr + arg_size);
16004 v_out.resize(v_out.size() + param_size - arg_size);
16008 if (convert_flags) {
16026 if (convert_flags) {
16039 printf(
"rpc_execute_cxx: send_tcp() sent %d bytes\n",
status);
16072 printf(
"rpc_test_rpc_test2!\n");
16075 int int_inout = 456;
16077 char string_out[33];
16078 char string2_out[49];
16080 char string_inout[25];
16081 strcpy(string_inout,
"string_inout");
16085 struct_in.
type = 111;
16087 strcpy(struct_in.
name,
"name");
16093 struct_inout.
type = 111111;
16095 strcpy(struct_inout.
name,
"name_name");
16098 uint32_t dwordarray_inout[9];
16099 size_t dwordarray_inout_size =
sizeof(dwordarray_inout);
16101 for (
int i=0;
i<9;
i++) {
16102 dwordarray_inout[
i] =
i*10;
16107 for (
size_t i=0;
i<
sizeof(array_in);
i++) {
16108 array_in[
i] =
'a' +
i;
16111 char array_out[16];
16112 size_t array_out_size =
sizeof(array_out);
16114 for (
size_t i=0;
i<
sizeof(array_out);
i++) {
16115 array_out[
i] =
'Z';
16123 string_out,
sizeof(string_out),
16124 string2_out,
sizeof(string2_out),
16125 string_inout,
sizeof(string_inout),
16129 dwordarray_inout, &dwordarray_inout_size,
16130 array_in,
sizeof(array_in),
16131 array_out, &array_out_size
16135 printf(
"rpc_call(RPC_TEST2) status %d\n",
status);
16139 if (int_out != 789) {
16140 printf(
"int_out mismatch!\n");
16144 if (int_inout != 456*2) {
16145 printf(
"int_inout mismatch!\n");
16149 if (strcmp(string_out,
"string_out") != 0) {
16150 printf(
"string_out mismatch [%s] vs [%s]\n", string_out,
"string_out");
16154 if (strcmp(string2_out,
"second string_out") != 0) {
16155 printf(
"string2_out mismatch [%s] vs [%s]\n", string2_out,
"second string_out");
16159 if (strcmp(string_inout,
"return string_inout") != 0) {
16160 printf(
"string_inout mismatch [%s] vs [%s]\n", string_inout,
"return string_inout");
16170 pkey = &struct_out;
16173 printf(
"struct_out mismatch: type %d, num_values %d, name [%s], last_written %d\n", pkey->
type, pkey->
num_values, pkey->
name, pkey->
last_written);
16177 pkey = &struct_inout;
16180 printf(
"struct_inout mismatch: type %d, num_values %d, name [%s], last_written %d\n", pkey->
type, pkey->
num_values, pkey->
name, pkey->
last_written);
16184 if (dwordarray_inout_size != 4*5) {
16185 printf(
"dwordarray_inout_size mismatch %d vs %d\n", (
int)dwordarray_inout_size, 4*5);
16188 for (
size_t i=0;
i<dwordarray_inout_size/
sizeof(uint32_t);
i++) {
16189 if (dwordarray_inout[
i] !=
i*10+
i) {
16190 printf(
"dwordarray_inout[%d] data mismatch %d vs %zu\n", (
int)
i, dwordarray_inout[
i],
i*10+
i);
16201 if (array_out_size != 15) {
16202 printf(
"array_out_size mismatch %d vs %d\n", (
int)array_out_size, 15);
16205 if (strcmp(array_out,
"test test test") != 0) {
16206 printf(
"array_out data mismatch\n");
16234 printf(
"rpc_test_rpc_test2_cxx!\n");
16237 int int_inout = 456;
16239 char string_out[33];
16240 std::string string2_out;
16241 std::string string_inout =
"string_inout";
16245 struct_in.
type = 111;
16247 strcpy(struct_in.
name,
"name");
16253 struct_inout.
type = 111111;
16255 strcpy(struct_inout.
name,
"name_name");
16258 uint32_t dwordarray_inout[9];
16259 size_t dwordarray_inout_size =
sizeof(dwordarray_inout);
16261 for (
int i=0;
i<9;
i++) {
16262 dwordarray_inout[
i] =
i*10;
16265 std::vector<char> array_in;
16266 int array_in_size = 10;
16268 for (
int i=0;
i<array_in_size;
i++) {
16269 array_in.push_back(
'a' +
i);
16272 std::vector<char> array_out;
16273 size_t array_out_size = 16;
16280 string_out,
sizeof(string_out),
16286 dwordarray_inout, &dwordarray_inout_size,
16287 &array_in, array_in_size,
16288 &array_out, &array_out_size
16292 printf(
"rpc_call(RPC_TEST2_CXX) status %d\n",
status);
16296 if (int_out != 789) {
16297 printf(
"int_out mismatch!\n");
16301 if (int_inout != 456*2) {
16302 printf(
"int_inout mismatch!\n");
16306 if (strcmp(string_out,
"string_out") != 0) {
16307 printf(
"string_out mismatch [%s] vs [%s]\n", string_out,
"string_out");
16311 if (string2_out !=
"second string_out") {
16312 printf(
"string2_out mismatch [%s] vs [%s]\n", string2_out.c_str(),
"second string_out");
16316 if (string_inout !=
"return string_inout") {
16317 printf(
"string_inout mismatch [%s] vs [%s]\n", string_inout.c_str(),
"return string_inout");
16327 pkey = &struct_out;
16330 printf(
"struct_out mismatch: type %d, num_values %d, name [%s], last_written %d\n", pkey->
type, pkey->
num_values, pkey->
name, pkey->
last_written);
16334 pkey = &struct_inout;
16337 printf(
"struct_inout mismatch: type %d, num_values %d, name [%s], last_written %d\n", pkey->
type, pkey->
num_values, pkey->
name, pkey->
last_written);
16341 if (dwordarray_inout_size != 4*5) {
16342 printf(
"dwordarray_inout_size mismatch %d vs %d\n", (
int)dwordarray_inout_size, 4*5);
16345 for (
size_t i=0;
i<dwordarray_inout_size/
sizeof(uint32_t);
i++) {
16346 if (dwordarray_inout[
i] !=
i*10+
i) {
16347 printf(
"dwordarray_inout[%d] data mismatch %d vs %zu\n", (
int)
i, dwordarray_inout[
i],
i*10+
i);
16358 if (array_out_size != 15) {
16359 printf(
"array_out_size mismatch %d vs %d\n", (
int)array_out_size, 15);
16361 }
else if (array_out.size() != 15) {
16362 printf(
"array_out.size() mismatch %d vs %d\n", (
int)array_out.size(), 15);
16365 if (strcmp(array_out.data(),
"test test test") != 0) {
16366 printf(
"array_out data mismatch\n");
16394 printf(
"rpc_test_rpc_test3_cxx!\n");
16397 int int_inout = 456;
16399 char string_out[33];
16400 std::string string2_out;
16401 std::string string_inout =
"string_inout";
16405 struct_in.
type = 111;
16407 strcpy(struct_in.
name,
"name");
16413 struct_inout.
type = 111111;
16415 strcpy(struct_inout.
name,
"name_name");
16418 uint32_t dwordarray_inout[9];
16419 size_t dwordarray_inout_size =
sizeof(dwordarray_inout);
16421 for (
int i=0;
i<9;
i++) {
16422 dwordarray_inout[
i] =
i*10;
16425 std::vector<char> array_in;
16426 int array_in_size = 10;
16428 for (
int i=0;
i<array_in_size;
i++) {
16429 array_in.push_back(
'a' +
i);
16432 std::vector<char> array_out;
16439 string_out,
sizeof(string_out),
16445 dwordarray_inout, &dwordarray_inout_size,
16451 printf(
"rpc_call(RPC_TEST3_CXX) status %d\n",
status);
16455 if (int_out != 789) {
16456 printf(
"int_out mismatch!\n");
16460 if (int_inout != 456*2) {
16461 printf(
"int_inout mismatch!\n");
16465 if (strcmp(string_out,
"string_out") != 0) {
16466 printf(
"string_out mismatch [%s] vs [%s]\n", string_out,
"string_out");
16470 if (string2_out !=
"second string_out") {
16471 printf(
"string2_out mismatch [%s] vs [%s]\n", string2_out.c_str(),
"second string_out");
16475 if (string_inout !=
"return string_inout") {
16476 printf(
"string_inout mismatch [%s] vs [%s]\n", string_inout.c_str(),
"return string_inout");
16486 pkey = &struct_out;
16489 printf(
"struct_out mismatch: type %d, num_values %d, name [%s], last_written %d\n", pkey->
type, pkey->
num_values, pkey->
name, pkey->
last_written);
16493 pkey = &struct_inout;
16496 printf(
"struct_inout mismatch: type %d, num_values %d, name [%s], last_written %d\n", pkey->
type, pkey->
num_values, pkey->
name, pkey->
last_written);
16500 if (dwordarray_inout_size != 4*5) {
16501 printf(
"dwordarray_inout_size mismatch %d vs %d\n", (
int)dwordarray_inout_size, 4*5);
16504 for (
size_t i=0;
i<dwordarray_inout_size/
sizeof(uint32_t);
i++) {
16505 if (dwordarray_inout[
i] !=
i*10+
i) {
16506 printf(
"dwordarray_inout[%d] data mismatch %d vs %zu\n", (
int)
i, dwordarray_inout[
i],
i*10+
i);
16517 if (array_out.size() != 15) {
16518 printf(
"array_out.size() mismatch %d vs %d\n", (
int)array_out.size(), 15);
16521 if (strcmp(array_out.data(),
"test test test") != 0) {
16522 printf(
"array_out data mismatch\n");
16550 printf(
"rpc_test_rpc_test4_cxx!\n");
16553 int int_inout = 456;
16555 std::string string_in =
"test string";
16556 std::string string_out;
16557 std::string string_inout =
"string_inout";
16559 std::vector<char> array_in;
16560 int array_in_size = 10;
16562 for (
int i=0;
i<array_in_size;
i++) {
16563 array_in.push_back(
'a' +
i);
16566 std::vector<char> array_out;
16568 std::vector<char> array_inout;
16569 int array_inout_size = 6;
16571 for (
int i=0;
i<array_inout_size;
i++) {
16572 array_inout.push_back(
'0' +
i);
16588 printf(
"rpc_call(RPC_TEST4_CXX) status %d\n",
status);
16592 if (int_out != 789) {
16593 printf(
"int_out mismatch!\n");
16597 if (int_inout != 456*2) {
16598 printf(
"int_inout mismatch!\n");
16602 if (string_out !=
"return string_out") {
16603 printf(
"string_out mismatch [%s] vs [%s]\n", string_out.c_str(),
"return string_out");
16607 if (string_inout !=
"return string_inout") {
16608 printf(
"string_inout mismatch [%s] vs [%s]\n", string_inout.c_str(),
"return string_inout");
16612 if (array_out.size() != 15) {
16613 printf(
"array_out.size() mismatch %d vs %d\n", (
int)array_out.size(), 15);
16616 if (strcmp(array_out.data(),
"test test test") != 0) {
16617 printf(
"array_out data mismatch\n");
16622 if (array_inout.size() != 12) {
16623 printf(
"array_inout.size() mismatch %d vs %d\n", (
int)array_inout.size(), 12);
16626 for (
int i=0;
i<6;
i++) {
16627 if (array_inout[
i] !=
'0' +
i) {
16628 printf(
"array_inout data mismatch, index %d, value %d should be %d\n",
i, array_inout[
i], (
'0'+
i));
16632 for (
int i=6;
i<12;
i++) {
16633 if (array_inout[
i] != 2*(
'0' + (
i-6))) {
16634 printf(
"array_inout data mismatch, index %d, value %d should be %d\n",
i, array_inout[
i], 2*(
'0'+
i));
16753 if (strcmp(hostname,
"localhost") == 0)
16756 if (strcmp(hostname,
"localhost.localdomain") == 0)
16759 if (strcmp(hostname,
"localhost6") == 0)
16762 if (strcmp(hostname,
"ip6-localhost") == 0)
16770 if (h == hostname) {
16787 std::string hostname;
16799 static std::atomic_int max_report(10);
16800 if (max_report > 0) {
16802 if (max_report == 0) {
16803 cm_msg(
MERROR,
"rpc_socket_check_allowed_host",
"rejecting connection from unallowed host \'%s\', this message will no longer be reported", hostname.c_str());
16805 cm_msg(
MERROR,
"rpc_socket_check_allowed_host",
"rejecting connection from unallowed host \'%s\'. Add this host to \"/Experiment/Security/RPC hosts/Allowed hosts\"", hostname.c_str());
16839 INT port1, port2, port3;
16841 char net_buffer[256];
16842 struct linger ling;
16847 sock = accept(lsock, NULL, NULL);
16872 char command = (char) toupper(net_buffer[0]);
16886#ifdef LOCAL_ROUTINES
16889 for (
unsigned i=0;
i<exptab.
exptab.size();
i++) {
16891 const char*
str = exptab.
exptab[
i].name.c_str();
16892 send(sock,
str, strlen(
str) + 1, 0);
16894 send(sock,
"", 1, 0);
16905 port1 = port2 = version[0] = 0;
16912 port1 = strtoul(net_buffer + 2, &ptr, 0);
16913 port2 = strtoul(ptr, &ptr, 0);
16914 port3 = strtoul(ptr, &ptr, 0);
16916 while (*ptr ==
' ')
16920 for (; *ptr != 0 && *ptr !=
' ' &&
i < (int)
sizeof(version) - 1;)
16921 version[
i++] = *ptr++;
16924 assert(
i < (
int)
sizeof(version));
16928 for (; *ptr != 0 && *ptr !=
' ';)
16931 while (*ptr ==
' ')
16935 for (; *ptr != 0 && *ptr !=
' ' && *ptr !=
'\n' && *ptr !=
'\r' &&
i < (int)
sizeof(
experiment) - 1;)
16945 mstrlcpy(v1, version,
sizeof(v1));
16946 if (strchr(v1,
'.'))
16947 if (strchr(strchr(v1,
'.') + 1,
'.'))
16948 *strchr(strchr(v1,
'.') + 1,
'.') = 0;
16952 if (strchr(
str,
'.'))
16953 if (strchr(strchr(
str,
'.') + 1,
'.'))
16954 *strchr(strchr(
str,
'.') + 1,
'.') = 0;
16956 if (strcmp(v1,
str) != 0) {
16958 cm_msg(
MERROR,
"rpc_server_accept",
"received string: %s", net_buffer + 2);
16973#ifdef LOCAL_ROUTINES
16979 bool found =
false;
16985 for (idx = 0; idx < exptab.
exptab.size(); idx++) {
16998 send(sock,
"2", 2, 0);
17007 char host_port1_str[30], host_port2_str[30], host_port3_str[30];
17008 char debug_str[30];
17017 const char *argv[10];
17018 argv[0] = mserver_path;
17020 argv[2] = host_port1_str;
17021 argv[3] = host_port2_str;
17022 argv[4] = host_port3_str;
17023 argv[5] = debug_str;
17030 argv[0], argv[1], argv[2], argv[3], argv[4], argv[5], argv[6], argv[7], argv[8],
17039 send(sock,
str, strlen(
str) + 1, 0);
17045 send(sock,
str, strlen(
str) + 1, 0);
17052 cm_msg(
MERROR,
"rpc_server_accept",
"received unknown command '%c' code %d", command, command);
17062 setsockopt(sock, SOL_SOCKET, SO_LINGER, (
char *) &ling,
sizeof(ling));
17093 INT client_hw_type = 0, hw_type;
17094 std::string client_program;
17097 char net_buffer[256], *p;
17099 int sock = accept(lsock, NULL, NULL);
17115 client_program =
"(unknown)";
17118 i =
recv_string(sock, net_buffer,
sizeof(net_buffer), 10000);
17125 p = strtok(net_buffer,
" ");
17127 client_hw_type = atoi(p);
17128 p = strtok(NULL,
" ");
17132 p = strtok(NULL,
" ");
17135 client_program = p;
17136 p = strtok(NULL,
" ");
17140 p = strtok(NULL,
" ");
17161 status = send(sock,
str.c_str(),
str.length() + 1, 0);
17195 int recv_sock, send_sock, event_sock;
17197 std::string client_program;
17198 INT client_hw_type, hw_type;
17200 char net_buffer[256];
17208 std::string errmsg;
17242 flag = 2 * 1024 * 1024;
17243 status = setsockopt(event_sock, SOL_SOCKET, SO_RCVBUF, (
char *) &flag,
sizeof(
INT));
17245 cm_msg(
MERROR,
"rpc_server_callback",
"cannot setsockopt(SOL_SOCKET, SO_RCVBUF), errno %d (%s)", errno,
17250 cm_msg(
MERROR,
"rpc_server_callback",
"timeout on receive remote computer info");
17259 client_hw_type = strtoul(net_buffer, &p, 0);
17264 client_program = p;
17298 sprintf(
str,
"%d", hw_type);
17299 send(recv_sock,
str, strlen(
str) + 1, 0);
17368 if (n_received <= 0) {
17370 cm_msg(
MERROR,
"rpc_server_receive_rpc",
"recv_net_command() returned %d", n_received);
17377 memcpy(&nc_in, buf,
sizeof(nc_in));
17388 bool rpc_cxx =
false;
17393 cm_msg(
MERROR,
"rpc_server_receive_rpc",
"Unknown RPC routine_id %d", routine_id);
17403 cm_msg(
MERROR,
"rpc_server_receive_rpc",
"rpc_execute() returned %d, abort",
status);
17413 }
while (remaining);
17428 if (strchr(
str,
'.'))
17429 *strchr(
str,
'.') = 0;
17430 cm_msg(
MTALK,
"rpc_server_receive_rpc",
"Program \'%s\' on host \'%s\' aborted", sa->
prog_name.c_str(),
str);
17447 cm_msg(
MERROR,
"rpc_server_receive_rpc",
"mserver unexpected shutdown, status %d",
status);
17509 static char *xbuf = NULL;
17510 static int xbufsize = 0;
17511 static bool xbufempty =
true;
17514 if (sa == NULL && xbufempty)
17517 static bool recurse =
false;
17520 cm_msg(
MERROR,
"rpc_server_receive_event",
"internal error: called recursively");
17532 if (xbufempty && sa) {
17535 if (n_received < 0) {
17537 cm_msg(
MERROR,
"rpc_server_receive_event",
"recv_event_server_realloc() returned %d, abort", n_received);
17541 if (n_received == 0) {
17557 INT *pbh = (
INT *) xbuf;
17565 cm_msg(
MERROR,
"rpc_server_receive_event",
"bm_send_event() error %d (SS_ABORT), abort",
status);
17576 cm_msg(
MERROR,
"rpc_server_receive_event",
"bm_send_event() error %d, mserver dropped this event",
status);
17592 if (strchr(
str,
'.'))
17593 *strchr(
str,
'.') = 0;
17594 cm_msg(
MTALK,
"rpc_server_receive_event",
"Program \'%s\' on host \'%s\' aborted", sa->
prog_name.c_str(),
str);
17660 }
else if (timeout_msec ==
BM_WAIT) {
17705 struct linger ling;
17714 setsockopt(sa->
recv_sock, SOL_SOCKET, SO_LINGER, (
char *) &ling,
sizeof(ling));
17718 setsockopt(sa->
send_sock, SOL_SOCKET, SO_LINGER, (
char *) &ling,
sizeof(ling));
17723 setsockopt(sa->
event_sock, SOL_SOCKET, SO_LINGER, (
char *) &ling,
sizeof(ling));
17778 struct timeval timeout;
17805 if (convert_flags) {
17814 cm_msg(
MINFO,
"rpc_check_channels",
"client \"%s\" on host \"%s\" failed watchdog test after %d sec, send_tcp() returned %d",
17837 timeout.tv_sec = 1;
17838 timeout.tv_usec = 0;
17846 if (now > timeout_end_ms)
17859 if (!FD_ISSET(sa->
send_sock, &readfds) &&
17862 cm_msg(
MERROR,
"rpc_check_channels",
"client \"%s\" on host \"%s\" failed watchdog test after %d sec",
17878 if (FD_ISSET(sa->
send_sock, &readfds)) {
17881 cm_msg(
MERROR,
"rpc_check_channels",
"client \"%s\" on host \"%s\" failed watchdog test after %d sec, recv_tcp() returned %d",
17929#ifndef DOXYGEN_SHOULD_SKIP_THIS
18080 if (((
PTYPE) event & 0x07) != 0) {
18081 cm_msg(
MERROR,
"bk_create",
"Bank %s created with unaligned event pointer",
name);
18090 *pdata = pbk32a + 1;
18098 *pdata = pbk32 + 1;
18111#ifndef DOXYGEN_SHOULD_SKIP_THIS
18124 DWORD bklen, bktype, bksze;
18150 memmove(pdest, (
char *) psbkh32a,
ALIGN8(bksze) +
sizeof(
BANK32A));
18164 memmove(pdest, (
char *) psbkh32,
ALIGN8(bksze) +
sizeof(
BANK32));
18171 psbkh = ((
BANK *) psdata - 1);
18178 memmove(pdest, (
char *) psbkh,
ALIGN8(bksze) +
sizeof(
BANK));
18217 if (*((
DWORD *) pbk32a->
name) == dname) {
18219 remaining = ((
char *) event + ((
BANK_HEADER *) event)->data_size +
18227 memmove(pbk32a, (
char *) (pbk32a + 1) +
ALIGN8(pbk32a->
data_size), remaining);
18232 }
while ((
DWORD) ((
char *) pbk32a - (
char *) event) <
18239 if (*((
DWORD *) pbk32->
name) == dname) {
18241 remaining = ((
char *) event + ((
BANK_HEADER *) event)->data_size +
18249 memmove(pbk32, (
char *) (pbk32 + 1) +
ALIGN8(pbk32->
data_size), remaining);
18254 }
while ((
DWORD) ((
char *) pbk32 - (
char *) event) <
18263 remaining = ((
char *) event + ((
BANK_HEADER *) event)->data_size +
18271 memmove(pbk, (
char *) (pbk + 1) +
ALIGN8(pbk->
data_size), remaining);
18276 }
while ((
DWORD) ((
char *) pbk - (
char *) event) <
18300 pbk32a->
data_size = (
DWORD) ((
char *) pdata - (
char *) (pbk32a + 1));
18302 printf(
"Warning: TID_STRUCT bank %c%c%c%c has zero size\n", pbk32a->
name[0], pbk32a->
name[1], pbk32a->
name[2], pbk32a->
name[3]);
18307 pbk32->
data_size = (
DWORD) ((
char *) pdata - (
char *) (pbk32 + 1));
18309 printf(
"Warning: TID_STRUCT bank %c%c%c%c has zero size\n", pbk32->
name[0], pbk32->
name[1], pbk32->
name[2], pbk32->
name[3]);
18314 uint32_t size = (uint32_t) ((
char *) pdata - (
char *) (pbk + 1));
18315 if (size > 0xFFFF) {
18316 printf(
"Error: Bank size %d exceeds 16-bit limit of 65526, please use bk_init32() to create a 32-bit bank\n", size);
18321 printf(
"Warning: TID_STRUCT bank %c%c%c%c has zero size\n", pbk->
name[0], pbk->
name[1], pbk->
name[2], pbk->
name[3]);
18323 if (size > 0xFFFF) {
18324 printf(
"Error: Bank size %d exceeds 16-bit limit of 65526, please use bk_init32() to create a 32-bit bank\n", size);
18371 if (pmbk32a == NULL)
18375 if (pmbk32 == NULL)
18389 strncat(bklist, (
char *) pmbk32a->
name, 4);
18391 strncat(bklist, (
char *) pmbk32->
name, 4);
18393 strncat(bklist, (
char *) pmbk->
name, 4);
18414 auto range_fits = [event,
event_size](
const void *ptr,
size_t size) {
18415 ptrdiff_t
offset = (
const char *) ptr - (
const char *) event;
18419 auto report_corrupt_bank = [pdata]() {
18420 cm_msg(
MERROR,
"bk_locate",
"corrupted bank in event");
18421 *((
void **) pdata) = NULL;
18428 while (range_fits(pbk32a, 1)) {
18429 if (!range_fits(pbk32a,
sizeof(
BANK32A)))
18430 return report_corrupt_bank();
18431 size_t remaining =
event_size - ((
char *) pbk32a - (
char *) event) -
sizeof(
BANK32A);
18433 size_t padding = (8 - bank_size % 8) % 8;
18434 if (bank_size > remaining || padding > remaining - bank_size)
18435 return report_corrupt_bank();
18436 if (*((
DWORD *) pbk32a->
name) == dname) {
18437 int tid = pbk32a->
type & 0xFF;
18439 return report_corrupt_bank();
18440 *((
void **) pdata) = pbk32a + 1;
18445 pbk32a = (
BANK32A *) ((
char *) (pbk32a + 1) + bank_size + padding);
18450 while (range_fits(pbk32, 1)) {
18451 if (!range_fits(pbk32,
sizeof(
BANK32)))
18452 return report_corrupt_bank();
18453 size_t remaining =
event_size - ((
char *) pbk32 - (
char *) event) -
sizeof(
BANK32);
18455 size_t padding = (8 - bank_size % 8) % 8;
18456 if (bank_size > remaining || padding > remaining - bank_size)
18457 return report_corrupt_bank();
18458 if (*((
DWORD *) pbk32->
name) == dname) {
18459 int tid = pbk32->
type & 0xFF;
18461 return report_corrupt_bank();
18462 *((
void **) pdata) = pbk32 + 1;
18467 pbk32 = (
BANK32 *) ((
char *) (pbk32 + 1) + bank_size + padding);
18472 while (range_fits(pbk, 1)) {
18473 if (!range_fits(pbk,
sizeof(
BANK)))
18474 return report_corrupt_bank();
18475 size_t remaining =
event_size - ((
char *) pbk - (
char *) event) -
sizeof(
BANK);
18477 size_t padding = (8 - bank_size % 8) % 8;
18478 if (bank_size > remaining || padding > remaining - bank_size)
18479 return report_corrupt_bank();
18481 int tid = pbk->
type & 0xFF;
18483 return report_corrupt_bank();
18484 *((
void **) pdata) = pbk + 1;
18489 pbk = (
BANK *) ((
char *) (pbk + 1) + bank_size + padding);
18495 *((
void **) pdata) = NULL;
18516 if (*((
DWORD *) pbk32a->
name) == dname) {
18517 int tid = pbk32a->
type & 0xFF;
18519 *((
void **) pdata) = NULL;
18522 *((
void **) pdata) = pbk32a + 1;
18528 *bktype = pbk32a->
type;
18537 if (*((
DWORD *) pbk32->
name) == dname) {
18538 int tid = pbk32->
type & 0xFF;
18540 *((
void **) pdata) = NULL;
18543 *((
void **) pdata) = pbk32 + 1;
18549 *bktype = pbk32->
type;
18559 int tid = pbk->
type & 0xFF;
18561 *((
void **) pdata) = NULL;
18564 *((
void **) pdata) = pbk + 1;
18570 *bktype = pbk->
type;
18578 *((
void **) pdata) = NULL;
18622 *pbk = (
BANK *) ((
char *) (*pbk + 1) +
ALIGN8((*pbk)->data_size));
18624 *((
void **) pdata) = (*pbk) + 1;
18627 *pbk = *((
BANK **) pdata) = NULL;
18636#ifndef DOXYGEN_SHOULD_SKIP_THIS
18662 *pbk = (
BANK32 *) ((
char *) (*pbk + 1) +
ALIGN8((*pbk)->data_size));
18664 *((
void **) pdata) = (*pbk) + 1;
18695 if (*pbk32a == NULL)
18698 *pbk32a = (
BANK32A *) ((
char *) (*pbk32a + 1) +
ALIGN8((*pbk32a)->data_size));
18700 *((
void **) pdata) = (*pbk32a) + 1;
18740 if (pbh->
flags < 0x10000 && !force)
18747 pbk = (
BANK *) (pbh + 1);
18757 pdata = pbk32a + 1;
18774 pbk = (
BANK *) pbk32a;
18777 pbk = (
BANK *) pbk32;
18786 while ((
char *) pdata < (
char *) pbk) {
18788 pdata = (
void *) (((
WORD *) pdata) + 1);
18796 while ((
char *) pdata < (
char *) pbk) {
18798 pdata = (
void *) (((
DWORD *) pdata) + 1);
18805 while ((
char *) pdata < (
char *) pbk) {
18807 pdata = (
void *) (((
double *) pdata) + 1);
18827#ifndef DOXYGEN_SHOULD_SKIP_THIS
18852#define MAX_RING_BUFFER 100
18932 if (handle == NULL || size <= 0 || max_event_size <= 0 || max_event_size > size / 2)
18936 if (
rb[
i].buffer == NULL)
18944 if (
rb[
i].buffer == NULL)
19029 for (
i = 0;
i <= millisec / 10;
i++) {
19042 rp >
rb[h].buffer) {
19096 unsigned char *new_wp;
19104 cm_msg(
MERROR,
"rb_increment_wp",
"invalid event size of %d bytes, max_event_size is %u bytes",
19109 new_wp =
rb[h].
wp + size;
19115 assert(
rb[h].rp !=
rb[h].buffer);
19117 if (new_wp >
rb[h].ep)
19172 for (
i = 0;
i <= millisec / 10;
i++) {
19174 if (
rb[h].wp !=
rb[h].rp) {
19176 *p =
rb[handle - 1].
rp;
19227 unsigned char *new_rp;
19238 new_rp =
rb[h].
rp + size;
19242 if (new_rp >= ep &&
rb[h].wp < ep)
19245 rb[handle - 1].
rp = new_rp;
19284 if (
rb[h].wp >=
rb[h].rp)
19302 cm_msg(
MERROR,
"cm_write_event_to_odb",
"event %d ODB record size mismatch, db_set_record() status %d", pevent->
event_id,
status);
19309 char *pdata, *pdata0;
19318 HNDLE hKeyRoot, hKeyl, *hKeys;
19327 for (
n=0 ; ;
n++) {
19330 if (pbk32a == NULL)
19353 if (pbk32a == NULL)
19391 cm_msg(
MERROR,
"cm_write_event_to_odb",
"please define bank \"%s\" in BANK_LIST in frontend",
name);
19396 for (
i = 0;;
i++) {
19409 cm_msg(
MERROR,
"cm_write_event_to_odb",
"cannot write bank \"%s\" to ODB, db_set_data1() status %d",
name,
status);
19412 hKeys[
n++] = hKeyl;
19423 cm_msg(
MERROR,
"cm_write_event_to_odb",
"cannot create key for bank \"%s\" with tid %d in ODB, db_create_key() status %d",
name, bktype,
status);
19428 cm_msg(
MERROR,
"cm_write_event_to_odb",
"cannot find key for bank \"%s\" in ODB, after db_create_key(), db_find_key() status %d",
name,
status);
19435 cm_msg(
MERROR,
"cm_write_event_to_odb",
"cannot write bank \"%s\" to ODB, db_set_data1() status %d",
name,
status);
19437 hKeys[
n++] = hKeyRoot;
19449 cm_msg(
MERROR,
"cm_write_event_to_odb",
"event format %d is not supported (see midas.h definitions of FORMAT_xxx)", format);
std::atomic_bool connected
size_t out_max_size_offset
BUFFER * get_pbuf() const
bm_lock_buffer_guard(BUFFER *pbuf, bool do_not_lock=false)
bm_lock_buffer_guard & operator=(const bm_lock_buffer_guard &)=delete
bm_lock_buffer_guard(const bm_lock_buffer_guard &)=delete
void set_string_size(std::string s, int size)
static bool exists(const std::string &name)
void connect(const std::string &path, const std::string &name, bool write_defaults, bool delete_keys_not_in_defaults=false)
INT transition(INT run_number, char *error)
INT al_get_alarms(std::string *presult)
void bk_init32a(void *event)
INT bk_close(void *event, void *pdata)
INT bk_iterate32a(const void *event, BANK32A **pbk32a, void *pdata)
static void copy_bk_name(char *dst, const char *src)
INT bk_swap(void *event, BOOL force)
BOOL bk_is32a(const void *event)
int bk_delete(void *event, const char *name)
BOOL bk_is32(const void *event)
INT bk_iterate32(const void *event, BANK32 **pbk, void *pdata)
INT bk_locate(const void *event, const char *name, void *pdata)
void bk_init(void *event)
INT bk_list(const void *event, char *bklist)
INT bk_copy(char *pevent, char *psrce, const char *bkname)
INT bk_iterate(const void *event, BANK **pbk, void *pdata)
void bk_init32(void *event)
void bk_create(void *event, const char *name, WORD type, void **pdata)
INT bk_find(const BANK_HEADER *pbkh, const char *name, DWORD *bklen, DWORD *bktype, void **pdata)
INT bk_size(const void *event)
static void bm_wakeup_producers_locked(const BUFFER_HEADER *pheader, const BUFFER_CLIENT *pc)
static INT bm_receive_event_rpc(INT buffer_handle, void *buf, int *buf_size, EVENT_HEADER **ppevent, std::vector< char > *pvec, int timeout_msec)
INT bm_open_buffer(const char *buffer_name, INT buffer_size, INT *buffer_handle)
static BOOL bm_validate_rp(const char *who, const BUFFER_HEADER *pheader, int rp)
INT bm_send_event(INT buffer_handle, const EVENT_HEADER *pevent, int unused, int timeout_msec)
static int bm_flush_cache_rpc(int buffer_handle, int timeout_msec)
static void bm_write_buffer_statistics_to_odb_copy(HNDLE hDB, const char *buffer_name, const char *client_name, int client_index, BUFFER_INFO *pbuf, BUFFER_HEADER *pheader)
static int bm_skip_event(BUFFER *pbuf)
static INT bm_flush_cache_locked(bm_lock_buffer_guard &pbuf_guard, int timeout_msec)
static INT bm_receive_event_rpc_cxx(INT buffer_handle, void *buf, int *buf_size, EVENT_HEADER **ppevent, std::vector< char > *pvec, int timeout_msec)
INT bm_write_statistics_to_odb(void)
#define MAX_DEFRAG_EVENTS
INT bm_delete_request(INT request_id)
INT bm_close_all_buffers(void)
INT bm_add_event_request(INT buffer_handle, short int event_id, short int trigger_mask, INT sampling_type, EVENT_HANDLER *func, INT request_id)
static void bm_write_buffer_statistics_to_odb(HNDLE hDB, BUFFER *pbuf, BOOL force)
static int bm_incr_rp_no_check(const BUFFER_HEADER *pheader, int rp, int total_size)
INT bm_receive_event_vec(INT buffer_handle, std::vector< char > *pvec, int timeout_msec)
static void bm_notify_reader_locked(BUFFER_HEADER *pheader, BUFFER_CLIENT *pc, int old_write_pointer, int request_id)
static int bm_find_first_request_locked(BUFFER_CLIENT *pc, const EVENT_HEADER *pevent)
static BOOL bm_check_requests(const BUFFER_CLIENT *pc, const EVENT_HEADER *pevent)
static void bm_cleanup_buffer_locked(BUFFER *pbuf, const char *who, DWORD actual_time)
static BOOL bm_update_read_pointer_locked(const char *caller_name, BUFFER_HEADER *pheader)
static void bm_convert_event_header(EVENT_HEADER *pevent, int convert_flags)
INT bm_request_event(HNDLE buffer_handle, short int event_id, short int trigger_mask, INT sampling_type, HNDLE *request_id, EVENT_HANDLER *func)
static int bm_validate_client_index_locked(bm_lock_buffer_guard &pbuf_guard)
INT bm_set_cache_size(INT buffer_handle, size_t read_size, size_t write_size)
static int bm_validate_buffer_locked(const BUFFER *pbuf)
INT bm_receive_event(INT buffer_handle, void *destination, INT *buf_size, int timeout_msec)
static int bm_fill_read_cache_locked(bm_lock_buffer_guard &pbuf_guard, int timeout_msec)
INT bm_compose_event_threadsafe(EVENT_HEADER *event_header, short int event_id, short int trigger_mask, DWORD data_size, DWORD *serial)
INT bm_remove_event_request(INT buffer_handle, INT request_id)
static double _bm_mutex_timeout_sec
INT bm_close_buffer(INT buffer_handle)
int bm_send_event_sg(int buffer_handle, int sg_n, const char *const sg_ptr[], const size_t sg_len[], int timeout_msec)
INT bm_compose_event(EVENT_HEADER *event_header, short int event_id, short int trigger_mask, DWORD data_size, DWORD serial)
static void bm_write_to_buffer_locked(BUFFER_HEADER *pheader, int sg_n, const char *const sg_ptr[], const size_t sg_len[], size_t total_size)
static EVENT_DEFRAG_BUFFER defrag_buffer[MAX_DEFRAG_EVENTS]
static int bm_next_rp(const char *who, const BUFFER_HEADER *pheader, const char *pdata, int rp)
INT bm_match_event(short int event_id, short int trigger_mask, const EVENT_HEADER *pevent)
static void bm_read_from_buffer_locked(const BUFFER_HEADER *pheader, int rp, char *buf, int event_size)
static BOOL bm_peek_read_cache_locked(BUFFER *pbuf, EVENT_HEADER **ppevent, int *pevent_size, int *ptotal_size)
static void bm_update_last_activity(DWORD millitime)
static void bm_incr_read_cache_locked(BUFFER *pbuf, int total_size)
static DWORD _bm_max_event_size
static void bm_clear_buffer_statistics(HNDLE hDB, BUFFER *pbuf)
INT bm_flush_cache(int buffer_handle, int timeout_msec)
static void bm_dispatch_event(int buffer_handle, EVENT_HEADER *pevent)
static void bm_validate_client_pointers_locked(const BUFFER_HEADER *pheader, BUFFER_CLIENT *pclient)
static INT bm_read_buffer(BUFFER *pbuf, INT buffer_handle, void **bufptr, void *buf, INT *buf_size, std::vector< char > *vecptr, int timeout_msec, int convert_flags, BOOL dispatch)
static int bm_wait_for_free_space_locked(bm_lock_buffer_guard &pbuf_guard, int timeout_msec, int requested_space, bool unlock_write_cache)
INT bm_get_buffer_handle(const char *buffer_name, INT *buffer_handle)
int bm_send_event_vec(int buffer_handle, const std::vector< char > &event, int timeout_msec)
static int _bm_lock_timeout
static int bm_peek_buffer_locked(BUFFER *pbuf, BUFFER_HEADER *pheader, BUFFER_CLIENT *pc, EVENT_HEADER **ppevent, int *pevent_size, int *ptotal_size)
INT bm_receive_event_alloc(INT buffer_handle, EVENT_HEADER **ppevent, int timeout_msec)
static void bm_reset_buffer_locked(BUFFER *pbuf)
void bm_remove_client_locked(BUFFER_HEADER *pheader, int j)
static INT bm_push_buffer(BUFFER *pbuf, int buffer_handle)
static int bm_wait_for_more_events_locked(bm_lock_buffer_guard &pbuf_guard, BUFFER_CLIENT *pc, int timeout_msec, BOOL unlock_read_cache)
INT cm_set_path(const char *path)
INT cm_register_transition(INT transition, INT(*func)(INT, char *), INT sequence_number)
INT cm_shutdown(const char *name, BOOL bUnique)
static int cm_transition_call(TrState *s, int idx)
INT cm_disconnect_client(HNDLE hConn, BOOL bShutdown)
static void load_rpc_hosts(HNDLE hDB, HNDLE hKey, int index, void *info)
static std::atomic< std::thread * > _watchdog_thread
static int cm_transition_detach(INT transition, INT run_number, char *errstr, INT errstr_size, INT async_flag, INT debug_flag)
INT cm_yield(INT millisec)
INT cm_get_experiment_database(HNDLE *hDB, HNDLE *hKeyClient)
INT cm_list_experiments_remote(const char *host_name, STRING_LIST *exp_names)
INT cm_get_watchdog_params(BOOL *call_watchdog, DWORD *timeout)
static void bm_defragment_event(HNDLE buffer_handle, HNDLE request_id, EVENT_HEADER *pevent, void *pdata, EVENT_HANDLER *dispatcher)
INT cm_connect_client(const char *client_name, HNDLE *hConn)
INT cm_connect_experiment(const char *host_name, const char *exp_name, const char *client_name, void(*func)(char *))
static BUFFER_CLIENT * bm_get_my_client_locked(bm_lock_buffer_guard &pbuf_guard)
INT cm_list_experiments_local(STRING_LIST *exp_names)
INT cm_start_watchdog_thread()
static INT bm_push_event(const char *buffer_name)
INT cm_get_experiment_semaphore(INT *semaphore_alarm, INT *semaphore_elog, INT *semaphore_history, INT *semaphore_msg)
static int tr_finish(HNDLE hDB, TrState *tr, int transition, int status, const char *errorstr)
INT cm_set_client_run_state(INT state)
static int bm_lock_buffer_read_cache(BUFFER *pbuf)
static void write_tr_client_to_odb(HNDLE hDB, const TrClient *tr_client)
INT cm_transition(INT transition, INT run_number, char *errstr, INT errstr_size, INT async_flag, INT debug_flag)
INT cm_stop_watchdog_thread()
INT cm_register_function(INT id, INT(*func)(INT, void **))
INT cm_exist(const char *name, BOOL bClientName)
INT cm_connect_experiment1(const char *host_name, const char *default_exp_name, const char *client_name, void(*func)(char *), INT odb_size, DWORD watchdog_timeout)
INT cm_check_client(HNDLE hDB, HNDLE hKeyClient)
static int xbm_lock_buffer(BUFFER *pbuf)
static void xbm_unlock_buffer(BUFFER *pbuf)
INT cm_dispatch_ipc(const char *message, int message_size, int client_socket)
static INT cm_transition1(INT transition, INT run_number, char *errstr, INT errstr_size, INT async_flag, INT debug_flag)
INT cm_select_experiment_remote(const char *host_name, std::string *exp_name)
INT cm_register_server(void)
static BOOL _ctrlc_pressed
static void init_rpc_hosts(HNDLE hDB)
void cm_ack_ctrlc_pressed()
INT cm_execute(const char *command, char *result, INT bufsize)
INT cm_get_watchdog_info(HNDLE hDB, const char *client_name, DWORD *timeout, DWORD *last)
INT cm_cleanup(const char *client_name, BOOL ignore_timeout)
void cm_check_connect(void)
static BUFFER * bm_get_buffer(const char *who, INT buffer_handle, int *pstatus)
INT cm_select_experiment_local(std::string *exp_name)
std::string cm_expand_env(const char *str)
std::string cm_get_client_name()
static int bm_lock_buffer_mutex(BUFFER *pbuf)
int cm_exec_script(const char *odb_path_to_script)
INT EXPRT cm_get_path_string(std::string *path)
INT cm_set_client_info(HNDLE hDB, HNDLE *hKeyClient, const char *host_name, const char *program_name, INT hw_type, const char *password, DWORD watchdog_timeout)
static void rpc_client_shutdown()
static std::atomic< bool > _watchdog_thread_is_running
static exptab_struct _exptab
INT cm_disconnect_experiment(void)
static void bm_cleanup(const char *who, DWORD actual_time, BOOL wrong_interval)
static INT bm_notify_client(const char *buffer_name, int s)
int cm_get_exptab(const char *expname, std::string *dir, std::string *user)
static void xcm_watchdog_thread()
static bool test_cm_expand_env1(const char *str, const char *expected)
INT cm_synchronize(DWORD *seconds)
std::string cm_get_exptab_filename()
std::string cm_get_path()
static DWORD _deferred_transition_mask
std::string cm_get_history_path(const char *history_channel)
void cm_test_expand_env()
INT cm_register_deferred_transition(INT transition, BOOL(*func)(INT, BOOL))
int cm_set_experiment_local(const char *exp_name)
static INT _requested_transition
INT cm_get_environment(char *host_name, int host_name_size, char *exp_name, int exp_name_size)
const char * cm_get_version()
INT cm_read_exptab(exptab_struct *exptab)
static int bm_lock_buffer_write_cache(BUFFER *pbuf)
INT cm_deregister_transition(INT transition)
INT cm_check_deferred_transition()
std::string cm_get_experiment_name()
INT cm_set_transition_sequence(INT transition, INT sequence_number)
static bool tr_compare(const std::unique_ptr< TrClient > &arg1, const std::unique_ptr< TrClient > &arg2)
INT cm_delete_client_info(HNDLE hDB, INT pid)
static std::atomic< bool > _watchdog_thread_run
const char * cm_get_revision()
INT cm_watchdog_thread(void *unused)
INT cm_set_experiment_database(HNDLE hDB, HNDLE hKeyClient)
BOOL cm_is_ctrlc_pressed()
static INT tr_main_thread(void *param)
void cm_ctrlc_handler(int sig)
INT cm_set_watchdog_params_local(BOOL call_watchdog, DWORD timeout)
INT cm_transition_cleanup()
INT cm_set_watchdog_params(BOOL call_watchdog, DWORD timeout)
static int cm_transition_call_direct(TrClient *tr_client)
INT cm_set_experiment_semaphore(INT semaphore_alarm, INT semaphore_elog, INT semaphore_history, INT semaphore_msg)
static INT cm_transition2(INT transition, INT run_number, char *errstr, INT errstr_size, INT async_flag, INT debug_flag)
INT cm_set_experiment_name(const char *name)
#define CM_INVALID_TRANSITION
#define CM_DEFERRED_TRANSITION
#define CM_TRANSITION_IN_PROGRESS
#define CM_WRONG_PASSWORD
#define CM_TRANSITION_CANCELED
#define CM_VERSION_MISMATCH
#define BM_INVALID_MIXING
#define BM_INVALID_HANDLE
#define DB_INVALID_HANDLE
#define DB_NO_MORE_SUBKEYS
#define RPC_EXCEED_BUFFER
#define RPC_DOUBLE_DEFINED
#define RPC_NOT_REGISTERED
#define RPC_MUTEX_TIMEOUT
#define RPC_NO_CONNECTION
#define RPC_HNDLE_CONNECT
#define RPC_HNDLE_MSERVER
#define VALIGN(adr, align)
RPC_LIST * rpc_get_internal_list(INT flag)
#define MESSAGE_BUFFER_NAME
#define DRI_LITTLE_ENDIAN
#define MAX_STRING_LENGTH
#define MESSAGE_BUFFER_SIZE
void() EVENT_HANDLER(HNDLE buffer_handler, HNDLE request_id, EVENT_HEADER *event_header, void *event_data)
INT() RPC_HANDLER(INT index, void *prpc_param[])
std::string ss_gethostname()
INT ss_suspend(INT millisec, INT msg)
INT ss_get_struct_align()
INT ss_mutex_release(MUTEX_T *mutex)
INT ss_suspend_init_odb_port()
bool ss_event_socket_has_data()
time_t ss_mktime(struct tm *tms)
int ss_file_exist(const char *path)
INT ss_semaphore_create(const char *name, HNDLE *semaphore_handle)
INT recv_tcp2(int sock, char *net_buffer, int buffer_size, int timeout_ms)
INT ss_suspend_set_client_listener(int listen_socket)
INT ss_socket_get_peer_name(int sock, std::string *hostp, int *portp)
int ss_file_link_exist(const char *path)
INT ss_mutex_delete(MUTEX_T *mutex)
int ss_socket_wait(int sock, INT millisec)
INT ss_suspend_set_server_acceptions(RPC_SERVER_ACCEPTION_LIST *acceptions)
DWORD ss_settime(DWORD seconds)
char * ss_getpass(const char *prompt)
INT ss_suspend_set_client_connection(RPC_SERVER_CONNECTION *connection)
INT ss_mutex_create(MUTEX_T **mutex, BOOL recursive)
INT ss_shm_open(const char *name, INT size, void **adr, size_t *shm_size, HNDLE *handle, BOOL get_size)
INT recv_string(int sock, char *buffer, DWORD buffer_size, INT millisec)
INT ss_write_tcp(int sock, const char *buffer, size_t buffer_size)
int ss_dir_exist(const char *path)
bool ss_timed_mutex_wait_for_sec(std::timed_mutex &mutex, const char *mutex_name, double timeout_sec)
INT ss_semaphore_release(HNDLE semaphore_handle)
std::string ss_get_cmdline(void)
int ss_file_copy(const char *src, const char *dst, bool append)
INT recv_tcp(int sock, char *net_buffer, DWORD buffer_size, INT flags)
INT ss_resume(INT port, const char *message)
midas_thread_t ss_gettid(void)
INT ss_semaphore_delete(HNDLE semaphore_handle, INT destroy_flag)
INT ss_sleep(INT millisec)
INT ss_socket_connect_tcp(const char *hostname, int tcp_port, int *sockp, std::string *error_msg_p)
INT ss_semaphore_wait_for(HNDLE semaphore_handle, DWORD timeout_millisec)
INT ss_socket_listen_tcp(bool listen_localhost, int tcp_port, int *sockp, int *tcp_port_p, std::string *error_msg_p)
char * ss_crypt(const char *buf, const char *salt)
INT ss_spawnv(INT mode, const char *cmdname, const char *const argv[])
INT ss_suspend_get_buffer_port(midas_thread_t thread_id, INT *port)
INT ss_socket_close(int *sockp)
char * ss_gets(char *string, int size)
INT ss_recv_net_command(int sock, DWORD *routine_id, DWORD *param_size, char **param_ptr, int timeout_ms)
INT ss_shm_close(const char *name, void *adr, size_t shm_size, HNDLE handle, INT destroy_flag)
INT send_tcp(int sock, char *buffer, DWORD buffer_size, INT flags)
void * ss_ctrlc_handler(void(*func)(int))
BOOL ss_pid_exists(int pid)
INT ss_system(const char *command)
INT ss_mutex_wait_for(MUTEX_T *mutex, INT timeout)
INT ss_file_find(const char *path, const char *pattern, char **plist)
static int cm_msg_retrieve1(const char *filename, time_t t, INT n_messages, char **messages, int *length, int *allocated, int *num_messages)
INT cm_msg1(INT message_type, const char *filename, INT line, const char *facility, const char *routine, const char *format,...)
int cm_msg_early_init(void)
INT EXPRT cm_msg_facilities(STRING_LIST *list)
int cm_msg_open_buffer(void)
int cm_msg_close_buffer(void)
static std::mutex gMsgBufMutex
static void add_message(char **messages, int *length, int *allocated, time_t tstamp, const char *new_message)
INT cm_msg_register(EVENT_HANDLER *func)
INT cm_msg_log(INT message_type, const char *facility, const char *message)
INT cm_msg_flush_buffer()
static std::deque< msg_buffer_entry > gMsgBuf
static INT cm_msg_send_event(DWORD ts, INT message_type, const char *send_message)
std::string cm_get_error(INT code)
INT cm_msg(INT message_type, const char *filename, INT line, const char *routine, const char *format,...)
static std::string cm_msg_format(INT message_type, const char *filename, INT line, const char *routine, const char *format, va_list *argptr)
INT cm_msg_retrieve(INT n_message, char *message, INT buf_size)
INT cm_msg_retrieve2(const char *facility, time_t t, INT n_message, char **messages, int *num_messages)
void cm_msg_get_logfile(const char *fac, time_t t, std::string *filename, std::string *linkname, std::string *linktarget)
INT cm_set_msg_print(INT system_mask, INT user_mask, int(*func)(const char *))
struct rpc_server_acception_struct RPC_SERVER_ACCEPTION
BOOL equal_ustring(const char *str1, const char *str2)
INT db_flush_database(HNDLE hDB)
INT db_get_data_index(HNDLE hDB, HNDLE hKey, void *data, INT *buf_size, INT idx, DWORD type)
INT db_delete_key(HNDLE hDB, HNDLE hKey, BOOL follow_links)
INT db_check_client(HNDLE hDB, HNDLE hKeyClient)
INT db_get_value(HNDLE hDB, HNDLE hKeyRoot, const char *key_name, void *data, INT *buf_size, DWORD type, BOOL create)
INT db_open_record(HNDLE hDB, HNDLE hKey, void *ptr, INT rec_size, WORD access_mode, void(*dispatcher)(INT, INT, void *), void *info)
INT db_open_database(const char *xdatabase_name, INT database_size, HNDLE *hDB, const char *client_name)
void db_cleanup(const char *who, DWORD actual_time, BOOL wrong_interval)
INT db_get_data(HNDLE hDB, HNDLE hKey, void *data, INT *buf_size, DWORD type)
INT db_create_key(HNDLE hDB, HNDLE hKey, const char *key_name, DWORD type)
INT db_set_mode(HNDLE hDB, HNDLE hKey, WORD mode, BOOL recurse)
INT db_get_key(HNDLE hDB, HNDLE hKey, KEY *key)
INT EXPRT db_get_value_string(HNDLE hdb, HNDLE hKeyRoot, const char *key_name, int index, std::string *s, BOOL create, int create_string_length)
INT db_get_watchdog_info(HNDLE hDB, const char *client_name, DWORD *timeout, DWORD *last)
INT db_set_data_index(HNDLE hDB, HNDLE hKey, const void *data, INT data_size, INT idx, DWORD type)
INT db_close_all_records()
INT db_watch(HNDLE hDB, HNDLE hKey, void(*dispatcher)(INT, INT, INT, void *), void *info)
INT db_close_all_databases(void)
INT db_set_data(HNDLE hDB, HNDLE hKey, const void *data, INT buf_size, INT num_values, DWORD type)
INT db_delete(HNDLE hDB, HNDLE hKeyRoot, const char *odb_path)
INT db_sprintf(char *string, const void *data, INT data_size, INT idx, DWORD type)
INT db_update_last_activity(DWORD millitime)
INT db_set_data1(HNDLE hDB, HNDLE hKey, const void *data, INT buf_size, INT num_values, DWORD type)
INT db_set_value(HNDLE hDB, HNDLE hKeyRoot, const char *key_name, const void *data, INT data_size, INT num_values, DWORD type)
INT db_find_key(HNDLE hDB, HNDLE hKey, const char *key_name, HNDLE *subhKey)
INT db_update_record_local(INT hDB, INT hKeyRoot, INT hKey, int index)
void db_set_watchdog_params(DWORD timeout)
INT db_update_record_mserver(INT hDB, INT hKeyRoot, INT hKey, int index, int client_socket)
void db_cleanup2(const char *client_name, int ignore_timeout, DWORD actual_time, const char *who)
int db_delete_client_info(HNDLE hDB, int pid)
static DATABASE * db_lock_database(HNDLE hDB, int *pstatus, const char *caller, bool check_attached=true)
INT db_set_client_name(HNDLE hDB, const char *client_name)
INT db_notify_clients_array(HNDLE hDB, HNDLE hKeys[], INT size)
INT db_set_record(HNDLE hDB, HNDLE hKey, void *data, INT buf_size, INT align)
INT db_enum_key(HNDLE hDB, HNDLE hKey, INT idx, HNDLE *subkey_handle)
static void db_unlock_database(DATABASE *pdb, const char *caller)
INT EXPRT db_set_value_string(HNDLE hDB, HNDLE hKeyRoot, const char *key_name, const std::string *s)
INT db_set_lock_timeout(HNDLE hDB, int timeout_millisec)
INT db_set_num_values(HNDLE hDB, HNDLE hKey, INT num_values)
INT db_protect_database(HNDLE hDB)
int rb_get_rp(int handle, void **p, int millisec)
int rb_delete(int handle)
int rb_get_wp(int handle, void **p, int millisec)
int rb_increment_rp(int handle, int size)
static volatile int _rb_nonblocking
int rb_increment_wp(int handle, int size)
int rb_create(int size, int max_event_size, int *handle)
int rb_get_buffer_level(int handle, int *n_bytes)
static RING_BUFFER rb[MAX_RING_BUFFER]
INT rpc_add_allowed_host(const char *hostname)
void rpc_convert_data(void *data, INT tid, INT flags, INT total_size, INT convert_flags)
INT rpc_client_connect(const char *host_name, INT port, const char *client_name, HNDLE *hConnection)
#define RPC_BM_ADD_EVENT_REQUEST
INT rpc_register_server(int port, int *plsock, int *pport)
#define RPC_CM_CHECK_CLIENT
static int recv_event_server_realloc(INT idx, RPC_SERVER_ACCEPTION *psa, char **pbuffer, int *pbuffer_size)
INT rpc_get_opt_tcp_size()
INT rpc_client_disconnect(HNDLE hConn, BOOL bShutdown)
#define RPC_BM_SEND_EVENT
INT rpc_client_call(HNDLE hConn, DWORD routine_id,...)
INT rpc_register_functions(const RPC_LIST *new_list, RPC_HANDLER func)
static std::atomic_bool gAllowedHostsEnabled(false)
INT rpc_server_callback(struct callback_addr *pcallback)
static std::mutex _client_connections_mutex
INT rpc_set_timeout(HNDLE hConn, int timeout_msec, int *old_timeout_msec)
#define RPC_CM_SYNCHRONIZE
static std::mutex gAllowedHostsMutex
INT recv_tcp_check(int sock)
#define RPC_BM_GET_BUFFER_INFO
#define RPC_CM_SET_CLIENT_INFO
const char * rpc_get_mserver_path()
RPC_SERVER_ACCEPTION * rpc_get_mserver_acception()
void rpc_calc_convert_flags(INT hw_type, INT remote_hw_type, INT *convert_flags)
INT rpc_server_connect(const char *host_name, const char *exp_name)
std::string rpc_get_name()
int rpc_test_rpc_test2_cxx()
static RPC_SERVER_ACCEPTION * rpc_new_server_acception()
#define RPC_BM_REMOVE_EVENT_REQUEST
INT rpc_server_receive_rpc(RPC_SERVER_ACCEPTION *sa)
void rpc_debug_printf(const char *format,...)
const char * rpc_tid_name_old(INT id)
int cm_query_transition(int *transition, int *run_number, int *trans_time)
#define RPC_RC_TRANSITION
int rpc_test_rpc_test3_cxx()
void rpc_va_arg(va_list *arg_ptr, INT arg_type, void *arg)
static INT rpc_execute_old(INT sock, int xroutine_id, const RPC_LIST &rl, char *buffer, INT convert_flags)
INT rpc_server_loop(void)
INT rpc_clear_allowed_hosts()
std::string rpc_get_mserver_hostname(void)
INT rpc_deregister_functions()
bool rpc_is_connected(void)
INT rpc_set_mserver_path(const char *path)
static std::vector< RPC_LIST > rpc_list
#define RPC_BM_CLOSE_BUFFER
static TLS_POINTER * tls_buffer
#define RPC_BM_SET_CACHE_SIZE
static std::vector< RPC_CLIENT_CONNECTION * > _client_connections
static void rpc_call_encode(va_list &ap, const RPC_LIST &rl, NET_COMMAND **nc)
static RPC_CLIENT_CONNECTION * rpc_get_locked_client_connection(HNDLE hConn)
static std::vector< std::string > gAllowedHosts
INT rpc_call(DWORD routine_id,...)
static std::mutex rpc_list_mutex
const char * rpc_tid_name(INT id)
static void rpc_call_encode_cxx(va_list &ap, const RPC_LIST &rl, NET_COMMAND **nc)
INT rpc_register_client(const char *name, RPC_LIST *list)
static std::vector< RPC_SERVER_ACCEPTION * > _server_acceptions
static INT rpc_socket_check_allowed_host(int sock)
#define RPC_BM_CLOSE_ALL_BUFFERS
INT rpc_server_shutdown(void)
int rpc_flush_event_socket(int timeout_msec)
INT rpc_send_event(INT buffer_handle, const EVENT_HEADER *pevent, int unused, INT async_flag, INT mode)
static TR_FIFO _tr_fifo[10]
static std::string _mserver_path
INT rpc_get_timeout(HNDLE hConn)
INT rpc_server_disconnect()
INT rpc_set_debug(void(*func)(const char *), INT mode)
INT rpc_client_accept(int lsock)
void rpc_vax2ieee_float(float *var)
static int rpc_find_rpc(int routine_id, RPC_LIST *pentry, bool *prpc_cxx)
INT rpc_set_opt_tcp_size(INT tcp_size)
#define RPC_BM_RECEIVE_EVENT_CXX
#define RPC_BM_GET_BUFFER_LEVEL
int rpc_name_tid(const char *name)
INT rpc_register_listener(int port, RPC_HANDLER func, int *plsock, int *pport)
INT rpc_server_receive_event(int idx, RPC_SERVER_ACCEPTION *sa, int timeout_msec)
#define RPC_BM_OPEN_BUFFER
#define RPC_BM_EMPTY_BUFFERS
INT rpc_server_accept(int lsock)
static std::mutex _tr_fifo_mutex
bool rpc_is_mserver(void)
static RPC_SERVER_ACCEPTION * _mserver_acception
void rpc_ieee2vax_float(float *var)
#define RPC_CM_SET_WATCHDOG_PARAMS
INT rpc_send_event1(INT buffer_handle, const EVENT_HEADER *pevent)
#define RPC_BM_RECEIVE_EVENT
static bool _rpc_is_remote
static int rpc_call_decode(va_list &ap, const RPC_LIST &rl, const char *buf, size_t buf_size)
INT rpc_send_event_sg(INT buffer_handle, int sg_n, const char *const sg_ptr[], const size_t sg_len[])
INT rpc_set_name(const char *name)
INT rpc_register_function(INT id, INT(*func)(INT, void **))
#define RPC_CM_MSG_RETRIEVE
#define RPC_BM_SKIP_EVENT
static INT rpc_execute_cxx(INT sock, int xroutine_id, const RPC_LIST &rl, char *buffer, INT convert_flags)
INT rpc_client_dispatch(int sock)
#define RPC_BM_INIT_BUFFER_COUNTERS
static RPC_SERVER_CONNECTION _server_connection
INT rpc_check_channels(void)
static INT rpc_transition_dispatch(INT idx, void *prpc_param[])
int rpc_test_rpc_test4_cxx()
static int handle_msg_odb(int n, const NET_COMMAND *nc)
#define RPC_BM_FLUSH_CACHE
void rpc_ieee2vax_double(double *var)
static int rpc_call_decode_cxx(va_list &ap, const RPC_LIST &rl, const char *buf, size_t buf_size)
void rpc_vax2ieee_double(double *var)
#define RPC_CM_GET_WATCHDOG_INFO
static int recv_net_command_realloc(RPC_SERVER_ACCEPTION *sa, char **pbuf, int *pbufsize, INT *remaining)
INT rpc_get_convert_flags(void)
INT rpc_check_allowed_host(const char *hostname)
void rpc_convert_single(void *data, INT tid, INT flags, INT convert_flags)
char exp_name[NAME_LENGTH]
BOOL debug
debug printouts
char host_name[HOST_NAME_LENGTH]
char expt_name[NAME_LENGTH]
char buffer_name[NAME_LENGTH]
static const int tid_size[]
static std::string join(const char *sep, const std::vector< std::string > &v)
static std::vector< TRANS_TABLE > _trans_table
static std::atomic_int _message_mask_system
static int disable_bind_rpc_to_localhost
static std::mutex gBuffersMutex
int(* MessagePrintCallback)(const char *)
std::string cm_transition_name(int transition)
static DBG_MEM_LOC * _mem_loc
static std::vector< BUFFER * > gBuffers
static void(* _debug_print)(const char *)
static std::mutex _trans_table_mutex
static std::string _experiment_name
static std::string _path_name
static std::string _client_name
static std::vector< EventRequest > _request_list
static int _rpc_connect_timeout
static const char * tid_name[]
static const ERROR_TABLE _error_table[]
static INT _watchdog_timeout
INT bm_get_buffer_info(INT buffer_handle, BUFFER_HEADER *buffer_header)
static MUTEX_T * _mutex_rpc
static EVENT_HANDLER * _msg_dispatch
void * dbg_calloc(unsigned int size, unsigned int count, char *file, int line)
static BOOL _rpc_registered
static std::atomic< MessagePrintCallback > _message_print
bool ends_with_char(const std::string &s, char c)
void dbg_free(void *adr, char *file, int line)
static std::atomic_int _message_mask_user
static std::mutex _request_list_mutex
INT bm_get_buffer_level(INT buffer_handle, INT *n_bytes)
static TRANS_TABLE _deferred_trans_table[]
static int _rpc_listen_socket
std::string msprintf(const char *format,...)
void * dbg_malloc(unsigned int size, char *file, int line)
static const char * tid_name_old[]
static std::vector< std::string > split(const char *sep, const std::string &s)
int cm_write_event_to_odb(HNDLE hDB, HNDLE hKey, const EVENT_HEADER *pevent, INT format)
INT bm_init_buffer_counters(INT buffer_handle)
#define DIR_SEPARATOR_STR
#define DEFAULT_WATCHDOG_TIMEOUT
#define DEFAULT_RPC_TIMEOUT
#define MIN_WRITE_CACHE_SIZE
#define MAX_WRITE_CACHE_EVENT_SIZE_DIV
#define DEFAULT_MAX_EVENT_SIZE
#define WATCHDOG_INTERVAL
#define MAX_EVENT_REQUESTS
#define BANK_FORMAT_64BIT_ALIGNED
#define BANK_FORMAT_32BIT
std::vector< std::string > STRING_LIST
#define MAX_WRITE_CACHE_SIZE_DIV
#define TRANSITION_ERROR_STRING_LENGTH
#define BANK_FORMAT_VERSION
#define message(type, str)
#define write(n, a, f, d)
static std::string remove(const std::string s, char c)
int gettimeofday(struct timeval *tp, void *tzp)
struct callback_addr callback
EVENT_REQUEST event_request[MAX_EVENT_REQUESTS]
BUFFER_CLIENT client[MAX_CLIENTS]
BUFFER_INFO(BUFFER *pbuf)
int client_count_write_wait[MAX_CLIENTS]
DWORD client_time_write_wait[MAX_CLIENTS]
std::timed_mutex buffer_mutex
std::timed_mutex read_cache_mutex
char client_name[NAME_LENGTH]
std::timed_mutex write_cache_mutex
int client_count_write_wait[MAX_CLIENTS]
BUFFER_HEADER * buffer_header
std::atomic< size_t > read_cache_size
std::atomic< size_t > write_cache_size
char buffer_name[NAME_LENGTH]
std::atomic_bool attached
DWORD client_time_write_wait[MAX_CLIENTS]
EVENT_HANDLER * dispatcher
NET_COMMAND_HEADER header
unsigned int max_event_size
RPC_PARAM param[MAX_RPC_PARAMS]
std::atomic< std::thread * > thread
std::atomic_bool finished
std::vector< int > wait_for_index
std::string waiting_for_client
std::vector< std::unique_ptr< TrClient > > clients
unsigned short host_port1
unsigned short host_port2
unsigned short host_port3
std::vector< exptab_entry > exptab
std::mutex event_sock_mutex
static double fac(double a)
static te_expr * list(state *s)