Skip to content

Commit ccf5896

Browse files
committed
Preserve memory-only events during flush
Partition DoNotStoreOnDisk records out of both batched and per-record persistence so flush returns them to memory instead of writing them to SQLite. Use the repository exception macros so exception-disabled builds retain their supported control flow, and correct the C# sample sequence property target. Files changed: - lib/offline/OfflineStorageHandler.cpp - tests/unittests/OfflineStorageTests.cpp - examples/cs/SampleCsNet48/Program.cs Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: c4cfcad8-1637-4e46-86cf-6bf200244b04
1 parent 0165474 commit ccf5896

3 files changed

Lines changed: 129 additions & 38 deletions

File tree

‎examples/cs/SampleCsNet48/Program.cs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -55,7 +55,7 @@ static void Main(string[] args)
5555
for (int i = 0; i < 999; i++)
5656
{
5757
EventProperties props2 = new EventProperties("EventSimpleFromCSharpApp");
58-
props.SetProperty("EventSeqNum", Convert.ToString(i));
58+
props2.SetProperty("EventSeqNum", Convert.ToString(i));
5959
logger.LogEvent(props2);
6060
}
6161

‎lib/offline/OfflineStorageHandler.cpp‎

Lines changed: 49 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@
1414
#include <algorithm>
1515
#include <cstdio>
1616
#include <exception>
17+
#include <iterator>
1718
#include <limits>
1819
#include <numeric>
1920
#include <set>
@@ -36,8 +37,7 @@ namespace MAT_NS_BEGIN {
3637
{
3738
}
3839

39-
OfflineStorageHandler::OfflineStorageHandler(ILogManager& logManager, IRuntimeConfig& runtimeConfig,
40-
ITaskDispatcher& taskDispatcher, std::shared_ptr<IOfflineStorageProvider> storageProvider) :
40+
OfflineStorageHandler::OfflineStorageHandler(ILogManager& logManager, IRuntimeConfig& runtimeConfig, ITaskDispatcher& taskDispatcher, std::shared_ptr<IOfflineStorageProvider> storageProvider) :
4141
m_observer(nullptr),
4242
m_logManager(logManager),
4343
m_config(runtimeConfig),
@@ -58,7 +58,7 @@ namespace MAT_NS_BEGIN {
5858
{
5959
if (!m_storageProvider)
6060
{
61-
throw std::invalid_argument("OfflineStorageHandler requires a storage provider");
61+
MATSDK_THROW(std::invalid_argument("OfflineStorageHandler requires a storage provider"));
6262
}
6363

6464
// TODO: [MG] - OfflineStorage_SQLite.cpp is performing similar checks
@@ -107,15 +107,11 @@ namespace MAT_NS_BEGIN {
107107
{
108108
if (m_active)
109109
{
110-
try
110+
MATSDK_TRY
111111
{
112112
m_logManager.EndActivity();
113113
}
114-
catch (const std::exception& e)
115-
{
116-
std::fprintf(stderr, "Failed to end telemetry activity: %s\n", e.what());
117-
}
118-
catch (...)
114+
MATSDK_CATCH(...)
119115
{
120116
std::fputs("Failed to end telemetry activity\n", stderr);
121117
}
@@ -262,7 +258,7 @@ namespace MAT_NS_BEGIN {
262258
return;
263259
}
264260
std::vector<StorageRecord> recordsToRecover;
265-
try
261+
MATSDK_TRY
266262
{
267263
// Flush could be executed from context of worker thread, as well as from TPM and
268264
// after HTTP callback. Make sure it is atomic / thread-safe.
@@ -293,14 +289,29 @@ namespace MAT_NS_BEGIN {
293289

294290
const size_t drainedBatchSize = recordsToRecover.size();
295291
recordsRemaining -= std::min(recordsRemaining, drainedBatchSize);
296-
const size_t batchSaved = m_offlineStorageDisk->StoreRecords(recordsToRecover);
292+
293+
auto memoryOnlyBegin = std::partition(
294+
recordsToRecover.begin(), recordsToRecover.end(),
295+
[](StorageRecord const& record)
296+
{
297+
return record.persistence != EventPersistence_DoNotStoreOnDisk;
298+
});
299+
std::vector<StorageRecord> memoryOnlyRecords(
300+
std::make_move_iterator(memoryOnlyBegin),
301+
std::make_move_iterator(recordsToRecover.end()));
302+
recordsToRecover.erase(memoryOnlyBegin, recordsToRecover.end());
303+
ReturnRecordsToMemory(memoryOnlyRecords);
304+
305+
const size_t batchSaved = recordsToRecover.empty()
306+
? 0
307+
: m_offlineStorageDisk->StoreRecords(recordsToRecover);
297308
// StoreRecords() removes permanently-invalid records before
298309
// returning, so compare against the remaining valid records.
299310
const size_t validBatchSize = recordsToRecover.size();
300311
if (batchSaved != validBatchSize)
301312
{
302313
LOG_WARN("Flush: disk store failed for the batch of %zu records; returning it to the queue for retry",
303-
validBatchSize);
314+
validBatchSize);
304315
ReturnRecordsToMemory(recordsToRecover);
305316
recordsToRecover.clear();
306317
break;
@@ -342,28 +353,26 @@ namespace MAT_NS_BEGIN {
342353
m_flushComplete.post();
343354
m_flushPending = false;
344355
}
345-
catch (...)
356+
MATSDK_CATCH(...)
346357
{
358+
#if HAVE_EXCEPTIONS
347359
std::exception_ptr failure = std::current_exception();
348-
try
360+
MATSDK_TRY
349361
{
350362
if (m_offlineStorageMemory && !recordsToRecover.empty())
351363
{
352364
ReturnRecordsToMemory(recordsToRecover);
353365
}
354366
}
355-
catch (const std::exception& e)
356-
{
357-
std::fprintf(stderr, "Failed to recover records after flush failure: %s\n", e.what());
358-
}
359-
catch (...)
367+
MATSDK_CATCH(...)
360368
{
361369
std::fputs("Failed to recover records after flush failure\n", stderr);
362370
}
363371
LOCKGUARD(m_flushLock);
364372
m_flushComplete.post();
365373
m_flushPending = false;
366374
std::rethrow_exception(failure);
375+
#endif
367376
}
368377
}
369378

@@ -461,21 +470,28 @@ namespace MAT_NS_BEGIN {
461470
{
462471
(void)record;
463472
LOG_ERROR("Flush: dropping event %s:%s: Invalid parameters",
464-
tenantTokenToId(record.tenantToken).c_str(), record.id.c_str());
473+
tenantTokenToId(record.tenantToken).c_str(), record.id.c_str());
465474
OnStorageFailed("Invalid parameters");
466475
}
467476

468477
size_t OfflineStorageHandler::StoreRecordsIndividually(std::vector<StorageRecord>& records)
469478
{
470479
size_t totalSaved = 0;
471480
std::vector<StorageRecord> recordsToRetry;
481+
std::vector<StorageRecord> memoryOnlyRecords;
472482
size_t nextRecord = 0;
473483

474-
try
484+
MATSDK_TRY
475485
{
476486
for (; nextRecord < records.size(); ++nextRecord)
477487
{
478488
auto const& record = records[nextRecord];
489+
if (record.persistence == EventPersistence_DoNotStoreOnDisk)
490+
{
491+
memoryOnlyRecords.push_back(record);
492+
continue;
493+
}
494+
479495
if (!IsValidDiskStorageRecord(record))
480496
{
481497
ReportInvalidDiskRecord(record);
@@ -503,8 +519,9 @@ namespace MAT_NS_BEGIN {
503519
break;
504520
}
505521
}
506-
catch (...)
522+
MATSDK_CATCH(...)
507523
{
524+
#if HAVE_EXCEPTIONS
508525
recordsToRetry.clear();
509526
for (size_t retryIndex = nextRecord; retryIndex < records.size(); ++retryIndex)
510527
{
@@ -514,14 +531,17 @@ namespace MAT_NS_BEGIN {
514531
}
515532
}
516533
records.clear();
534+
ReturnRecordsToMemory(memoryOnlyRecords);
517535
ReturnRecordsToMemory(recordsToRetry);
518-
throw;
536+
std::rethrow_exception(std::current_exception());
537+
#endif
519538
}
520539

540+
ReturnRecordsToMemory(memoryOnlyRecords);
521541
if (!recordsToRetry.empty())
522542
{
523543
LOG_WARN("Flush: per-record disk store failed after saving %zu of %zu records; returning %zu records to the queue for retry",
524-
totalSaved, records.size(), recordsToRetry.size());
544+
totalSaved, records.size(), recordsToRetry.size());
525545
ReturnRecordsToMemory(recordsToRetry);
526546
}
527547

@@ -535,38 +555,30 @@ namespace MAT_NS_BEGIN {
535555

536556
for (auto const& record : records)
537557
{
538-
try
558+
MATSDK_TRY
539559
{
540560
if (m_offlineStorageMemory && m_offlineStorageMemory->StoreRecord(record))
541561
{
542562
++returned;
543563
continue;
544564
}
545565
LOG_ERROR("Flush: failed to return event %s:%s to memory queue after disk store failure; dropping record",
546-
tenantTokenToId(record.tenantToken).c_str(), record.id.c_str());
566+
tenantTokenToId(record.tenantToken).c_str(), record.id.c_str());
547567
dropped[record.tenantToken]++;
548568
}
549-
catch (const std::exception& e)
550-
{
551-
std::fprintf(stderr, "Failed to recover a record after flush failure: %s\n", e.what());
552-
}
553-
catch (...)
569+
MATSDK_CATCH(...)
554570
{
555571
std::fputs("Failed to recover a record after flush failure\n", stderr);
556572
}
557573
}
558574

559575
if (!dropped.empty())
560576
{
561-
try
577+
MATSDK_TRY
562578
{
563579
OnStorageRecordsDropped(dropped);
564580
}
565-
catch (const std::exception& e)
566-
{
567-
std::fprintf(stderr, "Failed to report dropped records after flush failure: %s\n", e.what());
568-
}
569-
catch (...)
581+
MATSDK_CATCH(...)
570582
{
571583
std::fputs("Failed to report dropped records after flush failure\n", stderr);
572584
}

‎tests/unittests/OfflineStorageTests.cpp‎

Lines changed: 79 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -343,6 +343,85 @@ TEST(OfflineStorageHandlerFlushTests, BatchingOptOutUsesPerRecordDiskStores)
343343
handler.Flush();
344344
}
345345

346+
TEST(OfflineStorageHandlerFlushTests, BatchedFlushKeepsMemoryOnlyRecordsOffDisk)
347+
{
348+
NullLogManager logManager;
349+
NiceMock<MockIRuntimeConfig> config;
350+
NoopTaskDispatcher dispatcher;
351+
StrictMock<MockIOfflineStorageObserver> observer;
352+
353+
config[CFG_BOOL_ENABLE_BATCHED_STORAGE_FLUSH] = true;
354+
config[CFG_INT_RAM_QUEUE_SIZE] = 4096 * 20;
355+
356+
auto memory = std::make_shared<StrictMock<MockIOfflineStorage>>();
357+
auto disk = std::make_shared<StrictMock<MockIOfflineStorage>>();
358+
auto provider = std::make_shared<MockOfflineStorageProvider>(memory, disk);
359+
OfflineStorageHandler handler(logManager, config, dispatcher, provider);
360+
EXPECT_CALL(*memory, Initialize(Ref(handler))).WillOnce(Return());
361+
EXPECT_CALL(*disk, Initialize(Ref(handler))).WillOnce(Return());
362+
handler.Initialize(observer);
363+
364+
std::vector<StorageRecord> records;
365+
records.push_back(StorageRecord("persisted", "tenant-token",
366+
EventLatency_Normal, EventPersistence_Normal, 1, std::vector<uint8_t>{'x'}));
367+
records.push_back(StorageRecord("memory-only", "tenant-token",
368+
EventLatency_Normal, EventPersistence_DoNotStoreOnDisk, 1, std::vector<uint8_t>{'y'}));
369+
370+
EXPECT_CALL(*memory, GetSize())
371+
.WillOnce(Return(records.size()))
372+
.WillOnce(Return(static_cast<size_t>(1)));
373+
EXPECT_CALL(*memory, GetRecordCount(EventLatency_Unspecified))
374+
.WillOnce(Return(records.size()));
375+
EXPECT_CALL(*memory, GetRecords(false, EventLatency_Unspecified, 2000))
376+
.WillOnce(Return(records));
377+
EXPECT_CALL(*memory, StoreRecord(Field(&StorageRecord::id, "memory-only")))
378+
.WillOnce(Return(true));
379+
EXPECT_CALL(*disk, StoreRecord(_)).Times(0);
380+
EXPECT_CALL(*disk, StoreRecords(_)).WillOnce(Return(1));
381+
EXPECT_CALL(observer, OnStorageRecordsSaved(1));
382+
383+
handler.Flush();
384+
}
385+
386+
TEST(OfflineStorageHandlerFlushTests, PerRecordFlushKeepsMemoryOnlyRecordsOffDisk)
387+
{
388+
NullLogManager logManager;
389+
NiceMock<MockIRuntimeConfig> config;
390+
NoopTaskDispatcher dispatcher;
391+
StrictMock<MockIOfflineStorageObserver> observer;
392+
393+
config[CFG_BOOL_ENABLE_BATCHED_STORAGE_FLUSH] = false;
394+
config[CFG_INT_RAM_QUEUE_SIZE] = 4096 * 20;
395+
396+
auto memory = std::make_shared<StrictMock<MockIOfflineStorage>>();
397+
auto disk = std::make_shared<StrictMock<MockIOfflineStorage>>();
398+
auto provider = std::make_shared<MockOfflineStorageProvider>(memory, disk);
399+
OfflineStorageHandler handler(logManager, config, dispatcher, provider);
400+
EXPECT_CALL(*memory, Initialize(Ref(handler))).WillOnce(Return());
401+
EXPECT_CALL(*disk, Initialize(Ref(handler))).WillOnce(Return());
402+
handler.Initialize(observer);
403+
404+
std::vector<StorageRecord> records;
405+
records.push_back(StorageRecord("persisted", "tenant-token",
406+
EventLatency_Normal, EventPersistence_Normal, 1, std::vector<uint8_t>{'x'}));
407+
records.push_back(StorageRecord("memory-only", "tenant-token",
408+
EventLatency_Normal, EventPersistence_DoNotStoreOnDisk, 1, std::vector<uint8_t>{'y'}));
409+
410+
EXPECT_CALL(*memory, GetSize())
411+
.WillOnce(Return(records.size()))
412+
.WillOnce(Return(static_cast<size_t>(1)));
413+
EXPECT_CALL(*memory, GetRecords(false, EventLatency_Unspecified, 0))
414+
.WillOnce(Return(records));
415+
EXPECT_CALL(*memory, StoreRecord(Field(&StorageRecord::id, "memory-only")))
416+
.WillOnce(Return(true));
417+
EXPECT_CALL(*disk, StoreRecords(_)).Times(0);
418+
EXPECT_CALL(*disk, StoreRecord(Field(&StorageRecord::id, "persisted")))
419+
.WillOnce(Return(true));
420+
EXPECT_CALL(observer, OnStorageRecordsSaved(1));
421+
422+
handler.Flush();
423+
}
424+
346425
TEST(OfflineStorageHandlerFlushTests, BatchedFlushLimitsEachDiskWrite)
347426
{
348427
NullLogManager logManager;

0 commit comments

Comments
 (0)