|
|
|
@ -1,51 +1,27 @@
|
|
|
|
|
using System;
|
|
|
|
|
using System.Collections.Concurrent;
|
|
|
|
|
using System.Collections.Generic;
|
|
|
|
|
using System.Linq;
|
|
|
|
|
using System.Threading;
|
|
|
|
|
using ZeroLevel.Services.Pools;
|
|
|
|
|
|
|
|
|
|
namespace ZeroLevel.Network
|
|
|
|
|
{
|
|
|
|
|
internal sealed class RequestBuffer
|
|
|
|
|
{
|
|
|
|
|
private SpinLock _reqeust_lock = new SpinLock();
|
|
|
|
|
private Dictionary<long, RequestInfo> _requests = new Dictionary<long, RequestInfo>();
|
|
|
|
|
private ConcurrentDictionary<long, RequestInfo> _requests = new ConcurrentDictionary<long, RequestInfo>();
|
|
|
|
|
private static ObjectPool<RequestInfo> _ri_pool = new ObjectPool<RequestInfo>(() => new RequestInfo());
|
|
|
|
|
|
|
|
|
|
public void RegisterForFrame(int identity, Action<byte[]> callback, Action<string> fail = null)
|
|
|
|
|
{
|
|
|
|
|
var ri = _ri_pool.Allocate();
|
|
|
|
|
bool take = false;
|
|
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
_reqeust_lock.Enter(ref take);
|
|
|
|
|
ri.Reset(callback, fail);
|
|
|
|
|
_requests.Add(identity, ri);
|
|
|
|
|
}
|
|
|
|
|
finally
|
|
|
|
|
{
|
|
|
|
|
if (take) _reqeust_lock.Exit(false);
|
|
|
|
|
}
|
|
|
|
|
_requests[identity] = ri;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public void Fail(long frameId, string message)
|
|
|
|
|
{
|
|
|
|
|
RequestInfo ri = null;
|
|
|
|
|
bool take = false;
|
|
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
_reqeust_lock.Enter(ref take);
|
|
|
|
|
if (_requests.ContainsKey(frameId))
|
|
|
|
|
{
|
|
|
|
|
ri = _requests[frameId];
|
|
|
|
|
_requests.Remove(frameId);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
finally
|
|
|
|
|
{
|
|
|
|
|
if (take) _reqeust_lock.Exit(false);
|
|
|
|
|
}
|
|
|
|
|
if (ri != null)
|
|
|
|
|
RequestInfo ri;
|
|
|
|
|
if (_requests.TryRemove(frameId, out ri))
|
|
|
|
|
{
|
|
|
|
|
ri.Fail(message);
|
|
|
|
|
_ri_pool.Free(ri);
|
|
|
|
@ -54,22 +30,8 @@ namespace ZeroLevel.Network
|
|
|
|
|
|
|
|
|
|
public void Success(long frameId, byte[] data)
|
|
|
|
|
{
|
|
|
|
|
RequestInfo ri = null;
|
|
|
|
|
bool take = false;
|
|
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
_reqeust_lock.Enter(ref take);
|
|
|
|
|
if (_requests.ContainsKey(frameId))
|
|
|
|
|
{
|
|
|
|
|
ri = _requests[frameId];
|
|
|
|
|
_requests.Remove(frameId);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
finally
|
|
|
|
|
{
|
|
|
|
|
if (take) _reqeust_lock.Exit(false);
|
|
|
|
|
}
|
|
|
|
|
if (ri != null)
|
|
|
|
|
RequestInfo ri;
|
|
|
|
|
if (_requests.TryRemove(frameId, out ri))
|
|
|
|
|
{
|
|
|
|
|
ri.Success(data);
|
|
|
|
|
_ri_pool.Free(ri);
|
|
|
|
@ -78,21 +40,8 @@ namespace ZeroLevel.Network
|
|
|
|
|
|
|
|
|
|
public void StartSend(long frameId)
|
|
|
|
|
{
|
|
|
|
|
RequestInfo ri = null;
|
|
|
|
|
bool take = false;
|
|
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
_reqeust_lock.Enter(ref take);
|
|
|
|
|
if (_requests.ContainsKey(frameId))
|
|
|
|
|
{
|
|
|
|
|
ri = _requests[frameId];
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
finally
|
|
|
|
|
{
|
|
|
|
|
if (take) _reqeust_lock.Exit(false);
|
|
|
|
|
}
|
|
|
|
|
if (ri != null)
|
|
|
|
|
RequestInfo ri;
|
|
|
|
|
if (_requests.TryGetValue(frameId, out ri))
|
|
|
|
|
{
|
|
|
|
|
ri.StartSend();
|
|
|
|
|
}
|
|
|
|
@ -100,22 +49,13 @@ namespace ZeroLevel.Network
|
|
|
|
|
|
|
|
|
|
public void Timeout(List<long> frameIds)
|
|
|
|
|
{
|
|
|
|
|
bool take = false;
|
|
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
_reqeust_lock.Enter(ref take);
|
|
|
|
|
RequestInfo ri;
|
|
|
|
|
for (int i = 0; i < frameIds.Count; i++)
|
|
|
|
|
{
|
|
|
|
|
if (_requests.ContainsKey(frameIds[i]))
|
|
|
|
|
if (_requests.TryRemove(frameIds[i], out ri))
|
|
|
|
|
{
|
|
|
|
|
_ri_pool.Free(_requests[frameIds[i]]);
|
|
|
|
|
_requests.Remove(frameIds[i]);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
_ri_pool.Free(ri);
|
|
|
|
|
}
|
|
|
|
|
finally
|
|
|
|
|
{
|
|
|
|
|
if (take) _reqeust_lock.Exit(false);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
@ -123,24 +63,13 @@ namespace ZeroLevel.Network
|
|
|
|
|
{
|
|
|
|
|
var now_ticks = DateTime.UtcNow.Ticks;
|
|
|
|
|
var to_remove = new List<long>();
|
|
|
|
|
KeyValuePair<long, RequestInfo>[] to_proceed;
|
|
|
|
|
bool take = false;
|
|
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
_reqeust_lock.Enter(ref take);
|
|
|
|
|
to_proceed = _requests.Select(x => x).ToArray();
|
|
|
|
|
}
|
|
|
|
|
finally
|
|
|
|
|
{
|
|
|
|
|
if (take) _reqeust_lock.Exit(false);
|
|
|
|
|
}
|
|
|
|
|
for (int i = 0; i < to_proceed.Length; i++)
|
|
|
|
|
foreach (var pair in _requests)
|
|
|
|
|
{
|
|
|
|
|
if (to_proceed[i].Value.Sended == false) continue;
|
|
|
|
|
var diff = now_ticks - to_proceed[i].Value.Timestamp;
|
|
|
|
|
if (pair.Value.Sended == false) continue;
|
|
|
|
|
var diff = now_ticks - pair.Value.Timestamp;
|
|
|
|
|
if (diff > BaseSocket.MAX_REQUEST_TIME_TICKS)
|
|
|
|
|
{
|
|
|
|
|
to_remove.Add(to_proceed[i].Key);
|
|
|
|
|
to_remove.Add(pair.Key);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
Timeout(to_remove);
|
|
|
|
|