Skip to content

Commit 823e931

Browse files
committed
Merge branch 'remove_client_eof'
2 parents ac6c0a0 + 9a42ef3 commit 823e931

5 files changed

Lines changed: 45 additions & 17 deletions

File tree

yggdrasil/communication/CommBase.py

Lines changed: 12 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -894,6 +894,16 @@ def get_status_message(self, nindent=0, extra_lines_before=None,
894894
lines += ['%s%s' % (prefix, x) for x in extra_lines_after]
895895
return lines, prefix
896896

897+
# Re-enable this once the environment is crystalized on initialization
898+
# @property
899+
# def print_name(self):
900+
# r"""str: Name of the class object."""
901+
# out = super(CommBase, self).print_name
902+
# model_name = self.full_model_name
903+
# if model_name:
904+
# out += '[%s]' % model_name
905+
# return out
906+
897907
def printStatus(self, *args, **kwargs):
898908
r"""Print status of the communicator."""
899909
nindent = kwargs.get('nindent', 0)
@@ -1986,8 +1996,8 @@ def prepare_message(self, *args, header_kwargs=None, skip_serialization=False,
19861996
elif isinstance(msg.args, bytes) and (msg.args == YGG_CLIENT_EOF):
19871997
once_per_partner = True
19881998
if once_per_partner and (self.partner_copies > 1):
1989-
self.info("Sending %s to %d model(s)", msg.args,
1990-
self.partner_copies)
1999+
self.debug("Sending %s to %d model(s)", msg.args,
2000+
self.partner_copies)
19912001
for i in range(self.partner_copies - 1):
19922002
msg.add_message(args=msg.args,
19932003
header=copy.deepcopy(msg.header))

yggdrasil/config.py

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -571,8 +571,10 @@ def resolve_config_parser(args):
571571
args.validate_components = True
572572
args.validate_messages = True
573573
else:
574-
args.loglevel = 'INFO'
575-
args.client_loglevel = 'INFO'
574+
if args.loglevel is None:
575+
args.loglevel = 'INFO'
576+
if args.client_loglevel is None:
577+
args.client_loglevel = 'INFO'
576578
if args.validate_messages in ['True', 'False']:
577579
args.validate_messages = (args.validate_messages == 'True')
578580
return args

yggdrasil/drivers/RPCRequestDriver.py

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -117,16 +117,16 @@ def remove_model(self, direction, name):
117117
self.send_eof()
118118
return out
119119

120-
def send_eof(self):
121-
r"""Send EOF message.
120+
# def send_eof(self):
121+
# r"""Send EOF message.
122122

123-
Returns:
124-
bool: Success or failure of send.
123+
# Returns:
124+
# bool: Success or failure of send.
125125

126-
"""
127-
if self.ocomm.partner_copies > 1:
128-
self.ocomm.partner_copies = len(self.servers_recvd)
129-
return super(RPCRequestDriver, self).send_eof()
126+
# """
127+
# if self.ocomm.partner_copies > 1:
128+
# self.ocomm.partner_copies = len(self.servers_recvd)
129+
# return super(RPCRequestDriver, self).send_eof()
130130

131131
def on_eof(self, msg):
132132
r"""On EOF, decrement number of clients. Only send EOF if the number

yggdrasil/drivers/RPCResponseDriver.py

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,16 @@ def response_address(self):
4343
r"""str: Address of response comm."""
4444
return self.icomm.address
4545

46+
def send_eof(self):
47+
r"""Send EOF message.
48+
49+
Returns:
50+
bool: Success or failure of send.
51+
52+
"""
53+
# Don't send EOF
54+
return False
55+
4656
def send_message(self, msg, **kwargs):
4757
r"""Propagate the request_id.
4858

yggdrasil/languages/C/communication/ZMQComm.h

Lines changed: 11 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -790,14 +790,20 @@ int free_zmq_comm(comm_t *x) {
790790
if (_ygg_error_flag == 0) {
791791
size_t data_len = 100;
792792
char *data = (char*)malloc(data_len);
793+
comm_head_t head;
794+
bool is_eof_flag = false;
793795
while (zmq_comm_nmsg(x) > 0) {
794796
ret = zmq_comm_recv(x, &data, data_len, 1);
795-
if (ret < 0) {
796-
if (ret == -2) {
797+
if (ret >= 0) {
798+
head = parse_comm_header(data, ret);
799+
if (strncmp(YGG_MSG_EOF, data + head.bodybeg, strlen(YGG_MSG_EOF)) == 0)
800+
is_eof_flag = true;
801+
destroy_header(&head);
802+
if ((head.flags & HEAD_FLAG_VALID) && is_eof_flag) {
797803
x->const_flags[0] = x->const_flags[0] | COMM_EOF_RECV;
798-
break;
799-
}
800-
}
804+
break;
805+
}
806+
}
801807
}
802808
free(data);
803809
}

0 commit comments

Comments
 (0)