| 230 | // |
| 231 | // |
| 232 | void Queue::put(const Variant& vData) |
| 233 | { |
| 234 | // Apply filter |
| 235 | if (m_pFilter && m_pFilter->filterElem(vData)) |
| 236 | { |
| 237 | ++m_nFilteredCount; |
| 238 | return; |
| 239 | } |
| 240 | |
| 241 | // Simple mode processing |
| 242 | if (m_nBatchSize == 0) |
| 243 | { |
| 244 | putIntoData(vData, true); |
| 245 | return; |
| 246 | } |
| 247 | |
| 248 | // Batch mode processing |
| 249 | { |
| 250 | std::unique_lock sync(m_mtxData); |
| 251 | if (!testFlag(m_eMode, QueueMode::Put)) |
| 252 | error::OperationDeclined().throwException(); |
| 253 | } |
| 254 | |
| 255 | std::scoped_lock lock(m_mtxBatch); |
| 256 | |
| 257 | m_vBatch.push_back(vData); |
| 258 | auto nSize = m_vBatch.getSize(); |
| 259 | if (nSize >= m_nBatchSize) |
| 260 | { |
| 261 | if (m_pBatchTimer) |
| 262 | m_pBatchTimer = nullptr; |
| 263 | putIntoData(m_vBatch, true); |
| 264 | m_vBatch = Sequence(); |
| 265 | } |
| 266 | else if (m_nBatchTimeout != 0) |
| 267 | { |
| 268 | if (m_pBatchTimer == nullptr) |
| 269 | { |
| 270 | // avoid circular pointing timer <-> this |
| 271 | ObjWeakPtr<Queue> weakPtrToThis(getPtrFromThis(this)); |
| 272 | m_pBatchTimer = runWithDelay(m_nBatchTimeout, [weakPtrToThis]() |
| 273 | { |
| 274 | auto pThis = weakPtrToThis.lock(); |
| 275 | if (pThis == nullptr) |
| 276 | return false; |
| 277 | return pThis->flushBatch(); |
| 278 | }); |
| 279 | } |
| 280 | } |
| 281 | } |
| 282 | |
| 283 | // |
| 284 | // |
no test coverage detected