|
|
@ -1,4 +1,5 @@
|
|
|
|
using FASTER.core;
|
|
|
|
using FASTER.core;
|
|
|
|
|
|
|
|
using System;
|
|
|
|
using System.Collections.Generic;
|
|
|
|
using System.Collections.Generic;
|
|
|
|
using System.IO;
|
|
|
|
using System.IO;
|
|
|
|
using System.Threading.Tasks;
|
|
|
|
using System.Threading.Tasks;
|
|
|
@ -28,8 +29,15 @@ namespace ZeroLevel.Services.Microservices.Dump
|
|
|
|
public void Dump(T value)
|
|
|
|
public void Dump(T value)
|
|
|
|
{
|
|
|
|
{
|
|
|
|
var packet = MessageSerializer.SerializeCompatible(value);
|
|
|
|
var packet = MessageSerializer.SerializeCompatible(value);
|
|
|
|
while (!log.TryEnqueue(packet, out _)) ;
|
|
|
|
try
|
|
|
|
log.Commit();
|
|
|
|
{
|
|
|
|
|
|
|
|
while (!log.TryEnqueue(packet, out _)) ;
|
|
|
|
|
|
|
|
log.Commit();
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
catch (Exception ex)
|
|
|
|
|
|
|
|
{
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
public async Task DumpAsync(T value)
|
|
|
|
public async Task DumpAsync(T value)
|
|
|
@ -41,7 +49,7 @@ namespace ZeroLevel.Services.Microservices.Dump
|
|
|
|
public IEnumerable<T> ReadAndTruncate()
|
|
|
|
public IEnumerable<T> ReadAndTruncate()
|
|
|
|
{
|
|
|
|
{
|
|
|
|
byte[] result;
|
|
|
|
byte[] result;
|
|
|
|
using (var iter = log.Scan(log.BeginAddress, log.TailAddress))
|
|
|
|
using (var iter = log.Scan(log.BeginAddress, long.MaxValue))
|
|
|
|
{
|
|
|
|
{
|
|
|
|
while (iter.GetNext(out result, out int length))
|
|
|
|
while (iter.GetNext(out result, out int length))
|
|
|
|
{
|
|
|
|
{
|
|
|
|