26 #ifndef _MEMORYBOOKMARKSTORE_H_ 27 #define _MEMORYBOOKMARKSTORE_H_ 29 #include <amps/BookmarkStore.hpp> 42 #define AMPS_MIN_BOOKMARK_LEN 3 43 #define AMPS_INITIAL_MEMORY_BOOKMARK_SIZE 16384UL 64 typedef std::map<Message::Field, size_t, Message::Field::FieldHash> RecoveryMap;
65 typedef std::map<amps_uint64_t, amps_uint64_t> PublisherMap;
66 typedef std::map<Message::Field, size_t, Message::Field::FieldHash>::iterator
68 typedef std::map<amps_uint64_t, amps_uint64_t>::iterator PublisherIterator;
74 : _current(1), _currentBase(0), _least(1), _leastBase(0)
75 , _recoveryMin(AMPS_UNSET_INDEX), _recoveryBase(AMPS_UNSET_INDEX)
76 , _recoveryMax(AMPS_UNSET_INDEX), _recoveryMaxBase(AMPS_UNSET_INDEX)
77 , _entriesLength(AMPS_INITIAL_MEMORY_BOOKMARK_SIZE), _entries(NULL)
82 _store->resize(_id, (
char**)&_entries,
83 sizeof(Entry)*AMPS_INITIAL_MEMORY_BOOKMARK_SIZE,
false);
84 setLastPersistedToEpoch();
89 Lock<Mutex> guard(_subLock);
92 for (
size_t i = 0; i < _entriesLength; ++i)
94 _entries[i]._val.clear();
97 _store->resize(_id, (
char**)&_entries, 0);
101 _lastPersisted.clear();
104 _recoveryTimestamp.clear();
113 Lock<Mutex> guard(_subLock);
115 size_t index = recover(bookmark_,
true);
116 if (index == AMPS_UNSET_INDEX)
119 if (_current >= _entriesLength)
122 _currentBase += _entriesLength;
126 if ((_current == _least && _leastBase < _currentBase) ||
127 (_current == _recoveryMin && _recoveryBase < _currentBase))
129 if (!_store->resize(_id, (
char**)&_entries,
130 sizeof(Entry) * _entriesLength * 2))
133 return log(bookmark_);
166 bool isRange = BookmarkRange::isRange(bookmark_);
167 bool startExclusive =
false;
170 _range.set(bookmark_);
172 if (!_range.isValid())
175 throw CommandException(
"Invalid bookmark range specified.");
177 startExclusive = !_range.isStartInclusive();
179 if (isRange || _publishers.empty())
182 Message::Field parseable = isRange ? _range.getStart() : bookmark_;
183 amps_uint64_t publisher, sequence;
184 std::vector<Field> bmList = Field::parseBookmarkList(parseable);
185 if (bmList.empty() && !Field::isTimestamp(parseable))
193 for (std::vector<Field>::iterator bkmk = bmList.begin(); bkmk != bmList.end(); ++bkmk)
195 parseBookmark(*bkmk, publisher, sequence);
196 if (publisher != (amps_uint64_t)0)
198 if (isRange && startExclusive)
201 PublisherIterator pub = _publishers.find(publisher);
202 if (pub == _publishers.end() || pub->second < sequence)
204 _publishers[publisher] = sequence;
208 else if (!Field::isTimestamp(*bkmk))
225 _entries[_current]._val.deepCopy(bookmark_);
229 Unlock<Mutex> unlock(_subLock);
230 _store->updateAdapter(
this);
235 _entries[_current]._active =
true;
238 return index + _currentBase;
243 Lock<Mutex> guard(_subLock);
244 return _discard(index_);
254 Lock<Mutex> guard(_subLock);
255 size_t search = _least;
256 size_t searchBase = _leastBase;
257 size_t searchMax = _current;
258 size_t searchMaxBase = _currentBase;
259 if (_least + _leastBase == _current + _currentBase)
261 if (_recoveryMin != AMPS_UNSET_INDEX)
263 search = _recoveryMin;
264 searchBase = _recoveryBase;
265 searchMax = _recoveryMax;
266 searchMaxBase = _recoveryMaxBase;
273 assert(searchMax != AMPS_UNSET_INDEX);
274 assert(searchMaxBase != AMPS_UNSET_INDEX);
275 assert(search != AMPS_UNSET_INDEX);
276 assert(searchBase != AMPS_UNSET_INDEX);
278 while (search + searchBase < searchMax + searchMaxBase)
280 if (_entries[search]._val == bookmark_)
282 return _discard(search + searchBase);
284 if (++search == _entriesLength)
287 searchBase += _entriesLength;
296 amps_uint64_t& publisherId_,
297 amps_uint64_t& sequenceNumber_)
299 Message::Field::parseBookmark(field_, publisherId_, sequenceNumber_);
306 Lock<Mutex> guard(_subLock);
307 if (BookmarkRange::isRange(bookmark_))
312 amps_uint64_t publisher, sequence;
313 parseBookmark(bookmark_, publisher, sequence);
317 if (publisher == 0 && !Field::isTimestamp(bookmark_))
322 size_t recoveredIdx = recover(bookmark_,
false);
324 PublisherIterator pub = _publishers.find(publisher);
325 if (pub == _publishers.end() || pub->second < sequence)
327 _publishers[publisher] = sequence;
328 if (recoveredIdx == AMPS_UNSET_INDEX)
333 if (recoveredIdx != AMPS_UNSET_INDEX)
335 if (!_entries[recoveredIdx]._active)
337 _recovered.erase(bookmark_);
345 if (_store->_recovering)
353 size_t base = _leastBase;
354 for (
size_t i = _least; i + base < _current + _currentBase; i++)
356 if ( i >= _entriesLength )
361 if (_entries[i]._val == bookmark_)
363 return !_entries[i]._active;
370 bool empty(
void)
const 372 if (_least == AMPS_UNSET_INDEX ||
373 ((_least + _leastBase) == (_current + _currentBase) &&
374 _recoveryMin == AMPS_UNSET_INDEX))
381 void updateMostRecent()
383 Lock<Mutex> guard(_subLock);
387 const BookmarkRange& getRange()
const 394 Lock<Mutex> guard(_subLock);
395 bool useLastPersisted = !_lastPersisted.empty() &&
396 _lastPersisted.len() > 1;
401 bool useRecent = !_recent.empty() && _recent.len() > 1;
402 amps_uint64_t lastPublisher = 0;
403 amps_uint64_t lastSeq = 0;
404 amps_uint64_t recentPublisher = 0;
405 amps_uint64_t recentSeq = 0;
406 if (useLastPersisted)
408 parseBookmark(_lastPersisted, lastPublisher, lastSeq);
412 parseBookmark(_recent, recentPublisher, recentSeq);
413 if (empty() && useLastPersisted)
419 if (useLastPersisted && lastPublisher == recentPublisher)
421 if (lastSeq <= recentSeq)
427 useLastPersisted =
false;
433 size_t totalLen = (useLastPersisted ? _lastPersisted.len() + 1 : 0);
436 totalLen += _recent.len() + 1;
441 if (usePublishersList_
442 && ((!useLastPersisted && !useRecent)
445 std::ostringstream os;
446 for (PublisherIterator pub = _publishers.begin();
447 pub != _publishers.end(); ++pub)
449 if (pub->first == 0 && pub->second == 0)
453 if (pub->first == recentPublisher && recentSeq < pub->second)
455 os << recentPublisher <<
'|' << recentSeq <<
"|,";
459 os << pub->first <<
'|' << pub->second <<
"|,";
462 std::string recent = os.str();
465 totalLen = recent.length();
466 if (!_recoveryTimestamp.empty())
468 totalLen += _recoveryTimestamp.len();
469 recent += std::string(_recoveryTimestamp);
474 recent.erase(--totalLen);
479 if (_range.isValid())
481 if (_range.getStart() != recent
484 _range.replaceStart(_recentList,
true);
486 else if (_range.isStartInclusive())
488 amps_uint64_t publisher, sequence;
489 parseBookmark(_range.getStart(), publisher,
491 PublisherIterator pub = _publishers.find(publisher);
492 if (pub != _publishers.end()
493 && pub->second >= sequence)
495 _range.makeStartExclusive();
498 return _range.deepCopy();
500 return _recentList.deepCopy();
502 if (_range.isValid())
504 return _range.deepCopy();
507 if (!_recoveryTimestamp.empty() && !_range.isValid())
509 totalLen += _recoveryTimestamp.len() + 1;
513 || (_recent.len() < 2 && !empty()))
515 if (_range.isValid())
517 return _range.deepCopy();
523 _setLastPersistedToEpoch();
524 return _lastPersisted.deepCopy();
528 char* field = (
char*)malloc(totalLen);
537 memcpy(field, _recent.data(), len);
543 if (useLastPersisted)
545 memcpy(field + len, _lastPersisted.data(), _lastPersisted.len());
546 len += _lastPersisted.len();
552 if (!_recoveryTimestamp.empty() && !_range.isValid())
554 memcpy(field + len, _recoveryTimestamp.data(),
555 _recoveryTimestamp.len());
562 _recentList.assign(field, totalLen);
563 if (_range.isValid())
567 if (_range.getStart() != _recentList)
569 _range.replaceStart(_recentList,
true);
571 else if (_range.isStartInclusive())
573 amps_uint64_t publisher, sequence;
574 parseBookmark(_range.getStart(), publisher,
576 PublisherIterator pub = _publishers.find(publisher);
577 if (pub != _publishers.end()
578 && pub->second >= sequence)
580 _range.makeStartExclusive();
584 return _range.deepCopy();
586 return _recentList.deepCopy();
591 Lock<Mutex> guard(_subLock);
594 if (update_ && _store->_recentChanged)
610 Lock<Mutex> guard(_subLock);
611 return _lastPersisted;
617 _recent.deepCopy(recent_);
620 void setRecoveryTimestamp(
const char* recoveryTimestamp_,
623 _recoveryTimestamp.clear();
624 size_t len = (len_ == 0) ? AMPS_TIMESTAMP_LEN : len_;
625 char* ts = (
char*)malloc(len);
630 memcpy((
void*)ts, (
const void*)recoveryTimestamp_, len);
631 _recoveryTimestamp.assign(ts, len);
634 void moveEntries(
char* old_,
char* new_,
size_t newSize_)
636 size_t least = _least;
637 size_t leastBase = _leastBase;
638 if (_recoveryMin != AMPS_UNSET_INDEX)
640 least = _recoveryMin;
641 leastBase = _recoveryBase;
646 if (newSize_ - (
sizeof(Entry)*_entriesLength) >
sizeof(Entry)*least)
648 memcpy(new_ + (
sizeof(Entry)*_entriesLength),
649 old_, (
sizeof(Entry)*least));
651 memset(old_, 0,
sizeof(Entry)*least);
655 Entry* buffer =
new Entry[least];
656 memcpy((
void*)buffer, (
void*)old_,
sizeof(Entry)*least);
658 memcpy((
void*)new_, (
void*)((
char*)old_ + (
sizeof(Entry)*least)),
659 (_entriesLength - least)*
sizeof(Entry));
661 memcpy((
void*)((
char*)new_ + ((_entriesLength - least)*
sizeof(Entry))),
662 (
void*)buffer, least *
sizeof(Entry));
672 memcpy((
void*)new_, (
void*)((
char*)old_ + (
sizeof(Entry)*least)),
673 (_entriesLength - least)*
sizeof(Entry));
675 memcpy((
void*)((
char*)new_ + ((_entriesLength - least)*
sizeof(Entry))),
676 (
void*)old_, least *
sizeof(Entry));
681 if (_recoveryMin != AMPS_UNSET_INDEX)
683 _least = least + (_least + _leastBase) - (_recoveryMin + _recoveryBase);
684 _recoveryMax = least + (_recoveryMax + _recoveryMaxBase) -
685 (_recoveryMin + _recoveryBase);
686 _recoveryMaxBase = leastBase;
687 _recoveryMin = least;
688 _recoveryBase = leastBase;
694 _leastBase = leastBase;
696 _currentBase = _leastBase;
697 _current = least + _entriesLength;
702 Lock<Mutex> guard(_subLock);
704 return ((_least + _leastBase) == (_current + _currentBase)) ? AMPS_UNSET_INDEX :
712 || BookmarkRange::isRange(bookmark_))
716 Lock<Mutex> guard(_subLock);
717 return _setLastPersisted(bookmark_);
722 if (!_lastPersisted.empty())
724 amps_uint64_t publisher, publisher_lastPersisted;
725 amps_uint64_t sequence, sequence_lastPersisted;
726 parseBookmark(bookmark_, publisher, sequence);
727 parseBookmark(_lastPersisted, publisher_lastPersisted,
728 sequence_lastPersisted);
729 if (publisher == publisher_lastPersisted &&
730 sequence <= sequence_lastPersisted)
736 _lastPersisted.deepCopy(bookmark_);
737 _store->_recentChanged =
true;
738 _recoveryTimestamp.clear();
744 Lock<Mutex> guard(_subLock);
748 || BookmarkRange::isRange(bookmark))
752 _setLastPersisted(bookmark);
760 size_t recover(
const Message::Field& bookmark_,
bool relogIfNotDiscarded)
762 size_t retVal = AMPS_UNSET_INDEX;
763 if (_recovered.empty() || _recoveryBase == AMPS_UNSET_INDEX)
769 RecoveryIterator item = _recovered.find(bookmark_);
770 if (item != _recovered.end())
772 size_t seqNo = item->second;
773 size_t index = (seqNo - _recoveryBase) % _entriesLength;
776 if (_least + _leastBase == _current + _currentBase &&
777 !_entries[index]._active)
779 _store->_recentChanged =
true;
781 _recent = _entries[index]._val.deepCopy();
782 retVal = moveEntry(index);
783 if (retVal == AMPS_UNSET_INDEX)
785 recover(bookmark_, relogIfNotDiscarded);
788 _leastBase = _currentBase;
790 else if (!_entries[index]._active || relogIfNotDiscarded)
792 retVal = moveEntry(index);
793 if (retVal == AMPS_UNSET_INDEX)
795 recover(bookmark_, relogIfNotDiscarded);
802 _recovered.erase(item);
803 if (_recovered.empty())
805 _recoveryMin = AMPS_UNSET_INDEX;
806 _recoveryBase = AMPS_UNSET_INDEX;
807 _recoveryMax = AMPS_UNSET_INDEX;
808 _recoveryMaxBase = AMPS_UNSET_INDEX;
810 else if (index == _recoveryMin)
812 while (_entries[_recoveryMin]._val.empty() &&
813 (_recoveryMin + _recoveryBase) < (_recoveryMax + _recoveryMaxBase))
815 if (++_recoveryMin == _entriesLength)
818 _recoveryBase += _entriesLength;
838 Entry() : _active(
false)
844 typedef std::vector<Entry*> EntryPtrList;
846 void getRecoveryEntries(EntryPtrList& list_)
848 if (_recoveryMin == AMPS_UNSET_INDEX ||
849 _recoveryMax == AMPS_UNSET_INDEX)
853 size_t base = _recoveryBase;
854 size_t max = _recoveryMax + _recoveryMaxBase;
855 for (
size_t i = _recoveryMin; i + base < max; ++i)
857 if (i == _entriesLength)
860 base = _recoveryMaxBase;
863 list_.push_back(&(_entries[i]));
868 void getActiveEntries(EntryPtrList& list_)
870 size_t base = _leastBase;
871 for (
size_t i = _least; i + base < _current + _currentBase; ++i)
873 if (i >= _entriesLength)
879 list_.push_back(&(_entries[i]));
884 Entry* getEntryByIndex(
size_t index_)
886 Lock<Mutex> guard(_subLock);
887 size_t base = (_recoveryBase == AMPS_UNSET_INDEX ||
888 index_ >= _least + _leastBase)
889 ? _leastBase : _recoveryBase;
891 size_t min = (_recoveryMin == AMPS_UNSET_INDEX ?
892 _least + _leastBase :
893 _recoveryMin + _recoveryBase);
894 if (index_ >= _current + _currentBase || index_ < min)
898 return &(_entries[(index_ - base) % _entriesLength]);
903 Lock<Mutex> guard(_subLock);
906 getRecoveryEntries(list);
907 setPublishersToDiscarded(&list, &_publishers);
910 void setPublishersToDiscarded(EntryPtrList* recovered_,
911 PublisherMap* publishers_)
917 for (EntryPtrList::iterator i = recovered_->begin();
918 i != recovered_->end(); ++i)
920 if ((*i)->_val.empty())
924 amps_uint64_t publisher = (amps_uint64_t)0;
925 amps_uint64_t sequence = (amps_uint64_t)0;
926 parseBookmark((*i)->_val, publisher, sequence);
927 if (publisher && sequence && (*i)->_active &&
928 (*publishers_)[publisher] >= sequence)
930 (*publishers_)[publisher] = sequence - 1;
935 void clearLastPersisted()
937 Lock<Mutex> guard(_subLock);
938 _lastPersisted.clear();
941 void setLastPersistedToEpoch()
943 Lock<Mutex> guard(_subLock);
944 _setLastPersistedToEpoch();
948 Subscription(
const Subscription&);
949 Subscription& operator=(
const Subscription&);
951 size_t moveEntry(
size_t index_)
954 if (_current >= _entriesLength)
957 _currentBase += _entriesLength;
961 if ((_current == _least % _entriesLength &&
962 _leastBase < _currentBase) ||
963 (_current == _recoveryMin && _recoveryBase < _currentBase))
965 if (!_store->resize(_id, (
char**)&_entries,
966 sizeof(Entry) * _entriesLength * 2))
968 return AMPS_UNSET_INDEX;
973 _entries[_current]._val = _entries[index_]._val;
974 _entries[_current]._active = _entries[index_]._active;
976 _entries[index_]._val.assign(NULL, 0);
977 _entries[index_]._active =
false;
981 void _setLastPersistedToEpoch()
984 char* field = (
char*)malloc(fieldLen);
990 _lastPersisted.clear();
991 _lastPersisted.assign(field, fieldLen);
994 bool _discard(
size_t index_)
998 assert((_recoveryBase == AMPS_UNSET_INDEX && _recoveryMin == AMPS_UNSET_INDEX) ||
999 (_recoveryBase != AMPS_UNSET_INDEX && _recoveryMin != AMPS_UNSET_INDEX));
1000 size_t base = (_recoveryBase == AMPS_UNSET_INDEX
1001 || index_ >= _least + _leastBase)
1002 ? _leastBase : _recoveryBase;
1004 size_t min = (_recoveryMin == AMPS_UNSET_INDEX ? _least + _leastBase :
1005 _recoveryMin + _recoveryBase);
1006 if (index_ >= _current + _currentBase || index_ < min)
1013 Entry& e = _entries[(index_ - base) % _entriesLength];
1016 size_t index = index_;
1017 if (_recoveryMin != AMPS_UNSET_INDEX &&
1018 index_ == _recoveryMin + _recoveryBase)
1021 size_t j = _recoveryMin;
1022 while (j + _recoveryBase < _recoveryMax + _recoveryMaxBase &&
1023 !_entries[j]._active)
1044 if (!bookmark.
empty())
1046 _recovered.erase(bookmark);
1048 amps_uint64_t publisher, sequence;
1049 parseBookmark(bookmark, publisher, sequence);
1050 PublisherIterator pub = _publishers.find(publisher);
1051 if (pub == _publishers.end() || pub->second < sequence)
1053 _publishers[publisher] = sequence;
1055 if (_least + _leastBase == _current + _currentBase ||
1056 ((_least + _leastBase) % _entriesLength) ==
1057 ((_recoveryMin + _recoveryBase + 1)) % _entriesLength)
1061 _store->_recentChanged =
true;
1062 _recoveryTimestamp.clear();
1065 bookmark.assign(NULL, 0);
1075 if (++j == _entriesLength)
1078 _recoveryBase += _entriesLength;
1082 assert(j + _recoveryBase != _recoveryMax + _recoveryMaxBase ||
1083 _recovered.empty());
1084 if (_recovered.empty())
1086 _recoveryMin = AMPS_UNSET_INDEX;
1087 _recoveryBase = AMPS_UNSET_INDEX;
1088 _recoveryMax = AMPS_UNSET_INDEX;
1089 _recoveryMaxBase = AMPS_UNSET_INDEX;
1091 index = _least + _leastBase;
1100 if (index == _least + _leastBase)
1104 while (j + _leastBase < _current + _currentBase &&
1105 !_entries[j]._active)
1109 _recent = _entries[j]._val;
1110 _entries[j]._val.assign(NULL, 0);
1111 _store->_recentChanged =
true;
1113 _recoveryTimestamp.clear();
1116 if (++j == _entriesLength)
1119 _leastBase += _entriesLength;
1128 void _updateMostRecent()
1132 assert((_recoveryBase == AMPS_UNSET_INDEX && _recoveryMin == AMPS_UNSET_INDEX) ||
1133 (_recoveryBase != AMPS_UNSET_INDEX && _recoveryMin != AMPS_UNSET_INDEX));
1134 size_t base = (_recoveryMin == AMPS_UNSET_INDEX) ? _leastBase : _recoveryBase;
1135 size_t start = (_recoveryMin == AMPS_UNSET_INDEX) ? _least : _recoveryMin;
1136 _recoveryMin = AMPS_UNSET_INDEX;
1137 _recoveryBase = AMPS_UNSET_INDEX;
1138 _recoveryMax = AMPS_UNSET_INDEX;
1139 _recoveryMaxBase = AMPS_UNSET_INDEX;
1140 for (
size_t i = start; i + base < _current + _currentBase; i++)
1142 if ( i >= _entriesLength )
1145 base = _currentBase;
1147 if (i >= _recoveryMax + _recoveryBase && i < _least + _leastBase)
1151 Entry& entry = _entries[i];
1152 if (!entry._val.empty())
1154 _recovered[entry._val] = i + base;
1155 if (_recoveryMin == AMPS_UNSET_INDEX)
1158 _recoveryBase = base;
1159 _recoveryMax = _current;
1160 _recoveryMaxBase = _currentBase;
1164 if (_current == _entriesLength)
1167 _currentBase += _entriesLength;
1170 _leastBase = _currentBase;
1177 BookmarkRange _range;
1180 size_t _currentBase;
1183 size_t _recoveryMin;
1184 size_t _recoveryBase;
1185 size_t _recoveryMax;
1186 size_t _recoveryMaxBase;
1187 size_t _entriesLength;
1191 RecoveryMap _recovered;
1193 PublisherMap _publishers;
1202 _serverVersion(AMPS_DEFAULT_MIN_VERSION),
1203 _recentChanged(true),
1205 _recoveryPointAdapter(NULL),
1206 _recoveryPointFactory(NULL),
1207 _adapterSequence((amps_uint64_t)0),
1208 _nextAdapterUpdate((amps_uint64_t)0)
1211 typedef RecoveryPointAdapter::iterator RecoveryIterator;
1220 RecoveryPointFactory factory_ = NULL)
1224 , _serverVersion(AMPS_DEFAULT_MIN_VERSION)
1225 , _recentChanged(true)
1227 , _recoveryPointAdapter(adapter_)
1228 , _recoveryPointFactory(factory_)
1229 , _adapterSequence((amps_uint64_t)0)
1230 , _nextAdapterUpdate((amps_uint64_t)0)
1233 if (!_recoveryPointFactory)
1237 for (RecoveryIterator recoveryPoint = _recoveryPointAdapter.begin();
1238 recoveryPoint != _recoveryPointAdapter.end();
1241 Field subId(recoveryPoint->getSubId());
1242 msg.setSubscriptionHandle(static_cast<amps_subscription_handle>(0));
1244 Field bookmark = recoveryPoint->getBookmark();
1245 if (BookmarkRange::isRange(bookmark))
1252 std::vector<Field> bmList = Field::parseBookmarkList(bookmark);
1253 for (std::vector<Field>::iterator bkmk = bmList.begin(); bkmk != bmList.end(); ++bkmk)
1255 if (Field::isTimestamp(*bkmk))
1257 find(subId)->setRecoveryTimestamp(bkmk->data(), bkmk->len());
1271 _recovering =
false;
1276 if (_recoveryPointAdapter.isValid())
1278 _recoveryPointAdapter.close();
1290 Lock<Mutex> guard(_lock);
1291 return _log(message_);
1301 Lock<Mutex> guard(_lock);
1302 (void)_discard(message_);
1314 Lock<Mutex> guard(_lock);
1315 (void)_discard(subId_, bookmarkSeqNo_);
1325 Lock<Mutex> guard(_lock);
1326 return _getMostRecent(subId_);
1339 Lock<Mutex> guard(_lock);
1340 return _isDiscarded(message_);
1350 Lock<Mutex> guard(_lock);
1361 Lock<Mutex> guard(_lock);
1371 Lock<Mutex> guard(_lock);
1372 return _getOldestBookmarkSeq(subId_);
1383 Lock<Mutex> guard(_lock);
1384 _persisted(find(subId_), bookmark_);
1396 Lock<Mutex> guard(_lock);
1397 return _persisted(find(subId_), bookmark_);
1415 Lock<Mutex> guard(_subsLock);
1416 _serverVersion = version_;
1419 inline bool isWritableBookmark(
size_t length)
1421 return length >= AMPS_MIN_BOOKMARK_LEN;
1424 typedef Subscription::EntryPtrList EntryPtrList;
1429 size_t _log(
Message& message_)
1432 Subscription* pSub = (Subscription*)(message_.getSubscriptionHandle());
1441 message_.setSubscriptionHandle(
1442 static_cast<amps_subscription_handle>(pSub));
1444 size_t retVal = pSub->log(bookmark);
1445 message_.setBookmarkSeqNo(retVal);
1450 bool _discard(
const Message& message_)
1452 size_t bookmarkSeqNo = message_.getBookmarkSeqNo();
1453 Subscription* pSub = (Subscription*)(message_.getSubscriptionHandle());
1463 bool retVal = pSub->discard(bookmarkSeqNo);
1466 updateAdapter(pSub);
1472 bool _discard(
const Message::Field& subId_,
size_t bookmarkSeqNo_)
1474 Subscription* pSub = find(subId_);
1475 bool retVal = pSub->discard(bookmarkSeqNo_);
1478 updateAdapter(pSub);
1485 bool usePublishersList_ =
true)
1487 Subscription* pSub = find(subId_);
1488 return pSub->getMostRecentList(usePublishersList_);
1492 bool _isDiscarded(
Message& message_)
1499 Subscription* pSub = find(subId);
1500 message_.setSubscriptionHandle(
1501 static_cast<amps_subscription_handle>(pSub));
1508 Subscription* pSub = find(subId_);
1509 return pSub->getOldestBookmarkSeq();
1513 virtual void _persisted(Subscription* pSub_,
1516 if (pSub_->lastPersisted(bookmark_))
1518 updateAdapter(pSub_);
1523 virtual Message::Field _persisted(Subscription* pSub_,
size_t bookmark_)
1525 return pSub_->lastPersisted(bookmark_);
1531 if (_recoveryPointAdapter.isValid())
1533 _recoveryPointAdapter.purge();
1542 while (!_subs.empty())
1544 SubscriptionMap::iterator iter = _subs.begin();
1548 delete (iter->second);
1557 if (_recoveryPointAdapter.isValid())
1559 _recoveryPointAdapter.purge(subId_);
1567 Lock<Mutex> guard(_subsLock);
1568 SubscriptionMap::iterator iter = _subs.find(subId_);
1569 if (iter == _subs.end())
1574 delete (iter->second);
1582 find(subId_)->setMostRecent(recent_);
1587 static const char ENTRY_BOOKMARK =
'b';
1588 static const char ENTRY_DISCARD =
'd';
1589 static const char ENTRY_PERSISTED =
'p';
1595 throw StoreException(
"A valid subscription ID must be provided to the Bookmark Store");
1597 Lock<Mutex> guard(_subsLock);
1598 if (_subs.count(subId_) == 0)
1603 _subs[id] =
new Subscription(
this,
id);
1606 return _subs[subId_];
1609 virtual bool resize(
const Message::Field& subId_,
char** newBuffer_,
size_t size_,
1610 bool callResizeHandler_ =
true)
1625 if (callResizeHandler_ && !callResizeHandler(subId_, size_))
1629 char* allocBuffer = (
char*)malloc(size_);
1634 memset(allocBuffer, 0, size_);
1637 find(subId_)->moveEntries(*newBuffer_, allocBuffer, size_);
1640 *newBuffer_ = allocBuffer;
1645 void updateAdapter(Subscription* pSub_)
1647 if (_recovering || !_recentChanged || !_recoveryPointAdapter.isValid())
1651 Field bookmark = pSub_->getMostRecentList(
false);
1652 RecoveryPoint update = _recoveryPointFactory(pSub_->id(), bookmark);
1654 amps_uint64_t seq = _adapterSequence.fetch_add(1);
1655 Unlock<Mutex> unlock(_lock);
1656 while (_nextAdapterUpdate.load() < seq)
1662 _recoveryPointAdapter.update(update);
1664 catch (
const std::exception&)
1666 _nextAdapterUpdate.fetch_add(1);
1671 _nextAdapterUpdate.fetch_add(1);
1676 typedef std::map<Message::Field, Subscription*, Message::Field::FieldHash> SubscriptionMap;
1677 SubscriptionMap _subs;
1678 size_t _serverVersion;
1679 bool _recentChanged;
1681 typedef std::set<Subscription*> SubscriptionSet;
1684 std::atomic<amps_uint64_t> _adapterSequence;
1685 std::atomic<amps_uint64_t> _nextAdapterUpdate;
1690 #endif //_MEMORYBOOKMARKSTORE_H_ Defines the AMPS::Message class and related classes.
Abstract base class for storing received bookmarks for HA clients.
Definition: BookmarkStore.hpp:77
Field getSubscriptionId() const
Retrieves the value of the SubscriptionId header of the Message as a Field which references the under...
Definition: Message.hpp:1469
virtual void purge()
Called to purge the contents of this store.
Definition: MemoryBookmarkStore.hpp:1348
virtual void discard(const Message::Field &subId_, size_t bookmarkSeqNo_)
Log a discard-bookmark entry to the persistent log based on a bookmark sequence number.
Definition: MemoryBookmarkStore.hpp:1312
virtual Message::Field getMostRecent(const Message::Field &subId_)
Returns the most recent bookmark from the log that ought to be used for (re-)subscriptions.
Definition: MemoryBookmarkStore.hpp:1323
virtual Message::Field persisted(const Message::Field &subId_, size_t bookmark_)
Mark the bookmark provided as replicated to all sync replication destinations for the given subscript...
Definition: MemoryBookmarkStore.hpp:1393
Message encapsulates a single message sent to or received from an AMPS server, and provides methods f...
Definition: Message.hpp:519
void clear()
Deletes the data associated with this Field, should only be used on Fields that were created as deepC...
Definition: Field.hpp:266
static Field stringCopy(const char *str_)
Makes a copy of str_ in a new Field.
Definition: Field.hpp:253
MemoryBookmarkStore(const RecoveryPointAdapter &adapter_, RecoveryPointFactory factory_=NULL)
Creates a MemoryBookmarkStore.
Definition: MemoryBookmarkStore.hpp:1219
RecoveryPointAdapter a handle class for implementing external storage of subscription recovery points...
Definition: RecoveryPointAdapter.hpp:77
MemoryBookmarkStore()
Creates a MemoryBookmarkStore.
Definition: MemoryBookmarkStore.hpp:1199
Provides access to the subId and bookmark needed to restart a subscription.
Definition: RecoveryPoint.hpp:67
virtual void persisted(const Message::Field &subId_, const Message::Field &bookmark_)
Mark the bookmark provided as replicated to all sync replication destinations for the given subscript...
Definition: MemoryBookmarkStore.hpp:1380
Base class for all exceptions in AMPS.
Definition: AMPSException.hpp:40
A memory error occurred.
Definition: amps.h:225
Message & assignBookmark(const std::string &v)
Assigns the value of the Bookmark header for this Message without copying.
Definition: Message.hpp:1236
void setServerVersion(size_t version_)
Internally used to set the server version so the store knows how to deal with persisted acks and call...
Definition: MemoryBookmarkStore.hpp:1413
Field getSubscriptionIds() const
Retrieves the value of the SubscriptionIds header of the Message as a Field which references the unde...
Definition: Message.hpp:1470
bool empty() const
Returns 'true' if empty, 'false' otherwise.
Definition: Field.hpp:129
Defines the AMPS::Field class, which represents the value of a field in a message.
#define AMPS_BOOKMARK_EPOCH
Start the subscription at the beginning of the journal.
Definition: BookmarkStore.hpp:51
virtual bool isDiscarded(Message &message_)
Called for each arriving message to determine if the application has already seen this bookmark and s...
Definition: MemoryBookmarkStore.hpp:1337
Provides AMPS::RecoveryPointAdapter, an iterface for implementing external storage of bookmark subscr...
void setServerVersion(const VersionInfo &version_)
Internally used to set the server version so the store knows how to deal with persisted acks and call...
Definition: MemoryBookmarkStore.hpp:1404
virtual void discard(const Message &message_)
Log a discard-bookmark entry to the persistent log based on a bookmark sequence number.
Definition: MemoryBookmarkStore.hpp:1299
#define AMPS_BOOKMARK_NOW
Start the subscription at the point in time when AMPS processes the subscription. ...
Definition: BookmarkStore.hpp:55
Message & setSubId(const std::string &v)
Sets the value of the SubscriptionId header for this Message.
Definition: Message.hpp:1469
RecoveryPoint(* RecoveryPointFactory)(const Field &subId_, const Field &bookmark_)
RecoveryPointFactory is a function type for producing a RecoveryPoint that is sent to a RecoveryPoint...
Definition: RecoveryPoint.hpp:126
Provides AMPS::RecoveryPoint, AMPS::RecoveryPointFactory, AMPS::FixedRecoveryPoint, and AMPS::DynamicRecoveryPoint.
A BookmarkStoreImpl implementation that stores bookmarks in memory.
Definition: MemoryBookmarkStore.hpp:58
virtual size_t log(Message &message_)
Log a bookmark to the persistent log and return the corresponding sequence number for this bookmark...
Definition: MemoryBookmarkStore.hpp:1288
Field represents the value of a single field in a Message.
Definition: Field.hpp:87
static RecoveryPoint create(const Field &subId_, const Field &bookmark_)
Use this function in BookmarkStore::setRecoveryPointFactory( std::bind(&FixedRecoveryPoint::create, std::placeholder::_1, std::placeholder::_2))
Definition: RecoveryPoint.hpp:139
virtual size_t getOldestBookmarkSeq(const Message::Field &subId_)
Called to find the oldest bookmark in the store.
Definition: MemoryBookmarkStore.hpp:1369
virtual void purge(const Message::Field &subId_)
Called to purge the contents of this store for particular subId.
Definition: MemoryBookmarkStore.hpp:1359
Message & setBookmark(const std::string &v)
Sets the value of the Bookmark header for this Message.
Definition: Message.hpp:1236
void deepCopy(const Field &orig_)
Makes self a deep copy of the original field.
Definition: Field.hpp:219
Definition: AMPSException.hpp:32
Field getBookmark() const
Retrieves the value of the Bookmark header of the Message as a Field which references the underlying ...
Definition: Message.hpp:1236