| 179 | } |
| 180 | |
| 181 | std::shared_ptr<core::FlowFile> Connection::poll(std::set<std::shared_ptr<core::FlowFile>> &expiredFlowRecords) { |
| 182 | std::lock_guard<std::mutex> lock(mutex_); |
| 183 | |
| 184 | while (!queue_.empty()) { |
| 185 | std::shared_ptr<core::FlowFile> item = queue_.front(); |
| 186 | queue_.pop(); |
| 187 | queued_data_size_ -= item->getSize(); |
| 188 | |
| 189 | if (expired_duration_ > 0) { |
| 190 | // We need to check for flow expiration |
| 191 | if (utils::timeutils::getTimeMillis() > (item->getEntryDate() + expired_duration_)) { |
| 192 | // Flow record expired |
| 193 | expiredFlowRecords.insert(item); |
| 194 | logger_->log_debug("Delete flow file UUID %s from connection %s, because it expired", item->getUUIDStr(), name_); |
| 195 | } else { |
| 196 | // Flow record not expired |
| 197 | if (item->isPenalized()) { |
| 198 | // Flow record was penalized |
| 199 | queue_.push(item); |
| 200 | queued_data_size_ += item->getSize(); |
| 201 | break; |
| 202 | } |
| 203 | std::shared_ptr<Connectable> connectable = std::static_pointer_cast<Connectable>(shared_from_this()); |
| 204 | item->setConnection(connectable); |
| 205 | logger_->log_debug("Dequeue flow file UUID %s from connection %s", item->getUUIDStr(), name_); |
| 206 | return item; |
| 207 | } |
| 208 | } else { |
| 209 | // Flow record not expired |
| 210 | if (item->isPenalized()) { |
| 211 | // Flow record was penalized |
| 212 | queue_.push(item); |
| 213 | queued_data_size_ += item->getSize(); |
| 214 | break; |
| 215 | } |
| 216 | std::shared_ptr<Connectable> connectable = std::static_pointer_cast<Connectable>(shared_from_this()); |
| 217 | item->setConnection(connectable); |
| 218 | logger_->log_debug("Dequeue flow file UUID %s from connection %s", item->getUUIDStr(), name_); |
| 219 | return item; |
| 220 | } |
| 221 | } |
| 222 | |
| 223 | return NULL; |
| 224 | } |
| 225 | |
| 226 | void Connection::drain(bool delete_permanently) { |
| 227 | std::lock_guard<std::mutex> lock(mutex_); |
no test coverage detected