| 34 | } |
| 35 | |
| 36 | int |
| 37 | Puller::pull(const ACE_Time_Value& /*duration*/) |
| 38 | { |
| 39 | // Block until Publisher completes |
| 40 | DDS::StatusCondition_var condition = reader_->get_statuscondition(); |
| 41 | condition->set_enabled_statuses(DDS::SUBSCRIPTION_MATCHED_STATUS); |
| 42 | |
| 43 | DDS::WaitSet_var ws = new DDS::WaitSet; |
| 44 | ws->attach_condition(condition); |
| 45 | |
| 46 | // DDS::Duration_t timeout = {duration.sec(), 0}; |
| 47 | DDS::Duration_t timeout = {DDS::DURATION_INFINITE_SEC, DDS::DURATION_INFINITE_NSEC}; |
| 48 | |
| 49 | DDS::ConditionSeq conditions; |
| 50 | DDS::SubscriptionMatchedStatus matches = {0, 0, 0, 0, 0}; |
| 51 | |
| 52 | do |
| 53 | { |
| 54 | if (ws->wait(conditions, timeout) != DDS::RETCODE_OK) |
| 55 | { |
| 56 | ACE_ERROR_RETURN((LM_ERROR, |
| 57 | ACE_TEXT("%N:%l pull()") |
| 58 | ACE_TEXT(" ERROR: wait() failed: %p\n")), -1); |
| 59 | } |
| 60 | |
| 61 | if (reader_->get_subscription_matched_status(matches) != DDS::RETCODE_OK) |
| 62 | { |
| 63 | ACE_ERROR_RETURN((LM_ERROR, |
| 64 | ACE_TEXT("%N:%l pull()") |
| 65 | ACE_TEXT(" ERROR: get_subscription_matched_status() failed: %p\n")), -1); |
| 66 | } |
| 67 | } |
| 68 | while (matches.current_count > 0); |
| 69 | |
| 70 | ws->detach_condition(condition); |
| 71 | |
| 72 | return 0; |
| 73 | } |
no test coverage detected