3 changed files with 354 additions and 14 deletions
@ -0,0 +1,173 @@ |
|||
using System; |
|||
using System.Collections.Generic; |
|||
using System.IO; |
|||
using System.Runtime.CompilerServices; |
|||
using Avalonia.Rendering.Composition.Server; |
|||
|
|||
namespace Avalonia.Rendering.Composition.Transport; |
|||
|
|||
internal class BatchStreamData |
|||
{ |
|||
public Queue<BatchStreamSegment<ServerObject[]>> Objects { get; } = new(); |
|||
public Queue<BatchStreamSegment<IntPtr>> Structs { get; } = new(); |
|||
} |
|||
|
|||
public struct BatchStreamSegment<TData> |
|||
{ |
|||
public TData Data { get; set; } |
|||
public int ElementCount { get; set; } |
|||
} |
|||
|
|||
internal class BatchStreamWriter : IDisposable |
|||
{ |
|||
private readonly BatchStreamData _output; |
|||
private readonly BatchStreamMemoryPool _memoryPool; |
|||
private readonly BatchStreamObjectPool<ServerObject> _objectPool; |
|||
|
|||
private BatchStreamSegment<ServerObject[]?> _currentObjectSegment; |
|||
private BatchStreamSegment<IntPtr> _currentDataSegment; |
|||
|
|||
public BatchStreamWriter(BatchStreamData output, BatchStreamMemoryPool memoryPool, BatchStreamObjectPool<ServerObject> objectPool) |
|||
{ |
|||
_output = output; |
|||
_memoryPool = memoryPool; |
|||
_objectPool = objectPool; |
|||
} |
|||
|
|||
void CommitDataSegment() |
|||
{ |
|||
if (_currentDataSegment.Data != IntPtr.Zero) |
|||
_output.Structs.Enqueue(_currentDataSegment); |
|||
_currentDataSegment = new (); |
|||
} |
|||
|
|||
void NextDataSegment() |
|||
{ |
|||
CommitDataSegment(); |
|||
_currentDataSegment.Data = _memoryPool.Get(); |
|||
} |
|||
|
|||
void CommitObjectSegment() |
|||
{ |
|||
if (_currentObjectSegment.Data != null) |
|||
_output.Objects.Enqueue(_currentObjectSegment!); |
|||
_currentObjectSegment = new(); |
|||
} |
|||
|
|||
void NextObjectSegment() |
|||
{ |
|||
CommitObjectSegment(); |
|||
_currentObjectSegment.Data = _objectPool.Get(); |
|||
} |
|||
|
|||
public unsafe void Write<T>(T item) where T : unmanaged |
|||
{ |
|||
var size = Unsafe.SizeOf<T>(); |
|||
if (_currentDataSegment.Data == IntPtr.Zero || _currentDataSegment.ElementCount + size > _memoryPool.BufferSize) |
|||
NextDataSegment(); |
|||
*(T*)((byte*)_currentDataSegment.Data + _currentDataSegment.ElementCount) = item; |
|||
_currentDataSegment.ElementCount += size; |
|||
} |
|||
|
|||
public void Write(ServerObject item) |
|||
{ |
|||
if (_currentObjectSegment.Data == null || |
|||
_currentObjectSegment.ElementCount >= _currentObjectSegment.Data.Length) |
|||
NextObjectSegment(); |
|||
_currentObjectSegment.Data![_currentObjectSegment.ElementCount] = item; |
|||
_currentObjectSegment.ElementCount++; |
|||
} |
|||
|
|||
public void Dispose() |
|||
{ |
|||
CommitDataSegment(); |
|||
CommitObjectSegment(); |
|||
} |
|||
} |
|||
|
|||
internal class BatchStreamReader : IDisposable |
|||
{ |
|||
private readonly BatchStreamData _input; |
|||
private readonly BatchStreamMemoryPool _memoryPool; |
|||
private readonly BatchStreamObjectPool<ServerObject> _objectPool; |
|||
|
|||
private BatchStreamSegment<ServerObject[]?> _currentObjectSegment; |
|||
private BatchStreamSegment<IntPtr> _currentDataSegment; |
|||
private int _memoryOffset, _objectOffset; |
|||
|
|||
public BatchStreamReader(BatchStreamData _input, BatchStreamMemoryPool memoryPool, BatchStreamObjectPool<ServerObject> objectPool) |
|||
{ |
|||
this._input = _input; |
|||
_memoryPool = memoryPool; |
|||
_objectPool = objectPool; |
|||
} |
|||
|
|||
public unsafe T Read<T>() where T : unmanaged |
|||
{ |
|||
var size = Unsafe.SizeOf<T>(); |
|||
if (_currentDataSegment.Data == IntPtr.Zero) |
|||
{ |
|||
if (_input.Structs.Count == 0) |
|||
throw new EndOfStreamException(); |
|||
_currentDataSegment = _input.Structs.Dequeue(); |
|||
_memoryOffset = 0; |
|||
} |
|||
|
|||
if (_memoryOffset + size > _currentDataSegment.ElementCount) |
|||
throw new InvalidOperationException("Attempted to read more memory then left in the current segment"); |
|||
|
|||
var rv = *(T*)((byte*)_currentDataSegment.Data + size); |
|||
_memoryOffset += size; |
|||
if (_memoryOffset == _currentDataSegment.ElementCount) |
|||
{ |
|||
_memoryPool.Return(_currentDataSegment.Data); |
|||
_currentDataSegment = new(); |
|||
} |
|||
|
|||
return rv; |
|||
} |
|||
|
|||
public ServerObject ReadObject() |
|||
{ |
|||
if (_currentObjectSegment.Data == null) |
|||
{ |
|||
if (_input.Objects.Count == 0) |
|||
throw new EndOfStreamException(); |
|||
_currentObjectSegment = _input.Objects.Dequeue()!; |
|||
_objectOffset = 0; |
|||
} |
|||
|
|||
var rv = _currentObjectSegment.Data![_objectOffset]; |
|||
_objectOffset++; |
|||
if (_objectOffset == _currentObjectSegment.ElementCount) |
|||
{ |
|||
_objectPool.Return(_currentObjectSegment.Data); |
|||
_currentObjectSegment = new(); |
|||
} |
|||
|
|||
return rv; |
|||
} |
|||
|
|||
public bool IsStructEof => _currentDataSegment.Data == IntPtr.Zero && _input.Structs.Count == 0; |
|||
|
|||
public void Dispose() |
|||
{ |
|||
if (_currentDataSegment.Data != IntPtr.Zero) |
|||
{ |
|||
_memoryPool.Return(_currentDataSegment.Data); |
|||
_currentDataSegment = new(); |
|||
} |
|||
|
|||
while (_input.Structs.Count > 0) |
|||
_memoryPool.Return(_input.Structs.Dequeue().Data); |
|||
|
|||
if (_currentObjectSegment.Data != null) |
|||
{ |
|||
_objectPool.Return(_currentObjectSegment.Data); |
|||
_currentObjectSegment = new(); |
|||
} |
|||
|
|||
while (_input.Objects.Count > 0) |
|||
_objectPool.Return(_input.Objects.Dequeue().Data); |
|||
} |
|||
} |
|||
@ -0,0 +1,144 @@ |
|||
using System; |
|||
using System.Collections.Generic; |
|||
using System.Linq; |
|||
using System.Runtime.ConstrainedExecution; |
|||
using System.Runtime.InteropServices; |
|||
using Avalonia.Threading; |
|||
|
|||
namespace Avalonia.Rendering.Composition.Transport; |
|||
|
|||
/// <summary>
|
|||
/// A pool that keeps a number of elements that was used in the last 10 seconds
|
|||
/// </summary>
|
|||
internal abstract class BatchStreamPoolBase<T> : IDisposable |
|||
{ |
|||
readonly Stack<T> _pool = new(); |
|||
bool _disposed; |
|||
int _usage; |
|||
readonly int[] _usageStatistics = new int[10]; |
|||
int _usageStatisticsSlot; |
|||
|
|||
public BatchStreamPoolBase(bool needsFinalize = false) |
|||
{ |
|||
if(!needsFinalize) |
|||
GC.SuppressFinalize(needsFinalize); |
|||
|
|||
var updateRef = new WeakReference<BatchStreamPoolBase<T>>(this); |
|||
StartUpdateTimer(updateRef); |
|||
} |
|||
|
|||
static void StartUpdateTimer(WeakReference<BatchStreamPoolBase<T>> updateRef) |
|||
{ |
|||
DispatcherTimer.Run(() => |
|||
{ |
|||
if (updateRef.TryGetTarget(out var target)) |
|||
{ |
|||
target.UpdateStatistics(); |
|||
return true; |
|||
} |
|||
return false; |
|||
|
|||
}, TimeSpan.FromSeconds(1)); |
|||
} |
|||
|
|||
private void UpdateStatistics() |
|||
{ |
|||
lock (_pool) |
|||
{ |
|||
var maximumUsage = _usageStatistics.Max(); |
|||
var recentlyUsedPooledSlots = maximumUsage - _usage; |
|||
while (recentlyUsedPooledSlots < _pool.Count) |
|||
DestroyItem(_pool.Pop()); |
|||
|
|||
_usageStatistics[_usage] = 0; |
|||
_usageStatisticsSlot = (_usageStatisticsSlot + 1) % _usageStatistics.Length; |
|||
} |
|||
} |
|||
|
|||
protected abstract T CreateItem(); |
|||
|
|||
protected virtual void DestroyItem(T item) |
|||
{ |
|||
|
|||
} |
|||
|
|||
public T Get() |
|||
{ |
|||
lock (_pool) |
|||
{ |
|||
_usage++; |
|||
if (_usageStatistics[_usageStatisticsSlot] < _usage) |
|||
_usageStatistics[_usageStatisticsSlot] = _usage; |
|||
|
|||
if (_pool.Count != 0) |
|||
return _pool.Pop(); |
|||
} |
|||
|
|||
return CreateItem(); |
|||
} |
|||
|
|||
public void Return(T item) |
|||
{ |
|||
lock (_pool) |
|||
{ |
|||
_usage--; |
|||
if (!_disposed) |
|||
{ |
|||
_pool.Push(item); |
|||
return; |
|||
} |
|||
} |
|||
|
|||
DestroyItem(item); |
|||
} |
|||
|
|||
public void Dispose() |
|||
{ |
|||
lock (_pool) |
|||
{ |
|||
_disposed = true; |
|||
foreach (var item in _pool) |
|||
DestroyItem(item); |
|||
_pool.Clear(); |
|||
} |
|||
} |
|||
|
|||
~BatchStreamPoolBase() |
|||
{ |
|||
Dispose(); |
|||
} |
|||
} |
|||
|
|||
internal sealed class BatchStreamObjectPool<T> : BatchStreamPoolBase<T[]> where T : class |
|||
{ |
|||
private readonly int _arraySize; |
|||
|
|||
public BatchStreamObjectPool(int arraySize = 1024) |
|||
{ |
|||
_arraySize = arraySize; |
|||
} |
|||
|
|||
protected override T[] CreateItem() |
|||
{ |
|||
return new T[_arraySize]; |
|||
} |
|||
|
|||
protected override void DestroyItem(T[] item) |
|||
{ |
|||
Array.Clear(item, 0, item.Length); |
|||
} |
|||
} |
|||
|
|||
internal sealed class BatchStreamMemoryPool : BatchStreamPoolBase<IntPtr> |
|||
{ |
|||
public int BufferSize { get; } |
|||
|
|||
public BatchStreamMemoryPool(int bufferSize = 16384) |
|||
{ |
|||
BufferSize = bufferSize; |
|||
} |
|||
|
|||
protected override IntPtr CreateItem() => Marshal.AllocHGlobal(BufferSize); |
|||
|
|||
protected override void DestroyItem(IntPtr item) => Marshal.FreeHGlobal(item); |
|||
} |
|||
Loading…
Reference in new issue