|
|
|
@ -35,7 +35,7 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
/// <summary>
|
|
|
|
|
/// 最大队列长度
|
|
|
|
|
/// </summary>
|
|
|
|
|
public int MaxOperatingQueueSize = 5;
|
|
|
|
|
public int MaxOperatingQueueSize = 4;
|
|
|
|
|
|
|
|
|
|
/// <summary>
|
|
|
|
|
/// 同时上传的数量
|
|
|
|
@ -58,8 +58,6 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
throw new System.ArgumentNullException();
|
|
|
|
|
this.mTaskIdentify = ID;
|
|
|
|
|
this.mOperator = op;
|
|
|
|
|
//sLog = LoggerFactory.ForContext("Uploader:" + this.mTaskIdentify);
|
|
|
|
|
// 入口
|
|
|
|
|
mDatasWorkerThread = new System.Threading.Thread(uploadWorker)
|
|
|
|
|
{
|
|
|
|
|
Name = "upload worker",
|
|
|
|
@ -67,7 +65,6 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
};
|
|
|
|
|
mDatasWorkerThread.Start();
|
|
|
|
|
|
|
|
|
|
// Swap 处理: 超过4s保存至数据库 Caches 表
|
|
|
|
|
mSwapClearTimer = new Timer((x) =>
|
|
|
|
|
{
|
|
|
|
|
var toDBs = new List<UploadModel<T>>();
|
|
|
|
@ -160,15 +157,16 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
|
|
|
|
|
if (reqModel == null)
|
|
|
|
|
{
|
|
|
|
|
//sLog.Formf($"no request found:{request.RequestData}");
|
|
|
|
|
Console.WriteLine("no request found:{request.RequestData}");
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
reqModel.ErrorInfo = reason;
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
//sLog.Formf($"{reqModel.Request} finished:{success}");
|
|
|
|
|
string cancelRetryReason = null;
|
|
|
|
|
if (success)
|
|
|
|
|
{
|
|
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
using (var db = new CodeFirstDbContext())
|
|
|
|
|
{
|
|
|
|
|
db.UploadFinishs.Add(new UploadFinish
|
|
|
|
@ -177,10 +175,16 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
RetryAt = reqModel.RetryAt,
|
|
|
|
|
RetryCount = reqModel.RetryCount,
|
|
|
|
|
RequestData = mOperator.ConvertRequestToCachel(reqModel.Request),
|
|
|
|
|
Tag = reqModel.Tag,
|
|
|
|
|
UploaderID = mTaskIdentify
|
|
|
|
|
});
|
|
|
|
|
db.SaveChanges();
|
|
|
|
|
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
catch (Exception ex)
|
|
|
|
|
{
|
|
|
|
|
Console.WriteLine(ex.Message);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
else if (mOperator.ShouldRetryRequest(reqModel, out cancelRetryReason))
|
|
|
|
|
{
|
|
|
|
@ -195,6 +199,8 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
else
|
|
|
|
|
{
|
|
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
using (var db = new CodeFirstDbContext())
|
|
|
|
|
{
|
|
|
|
@ -205,14 +211,18 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
RetryCount = reqModel.RetryCount,
|
|
|
|
|
RequestData = mOperator.ConvertRequestToCachel(reqModel.Request),
|
|
|
|
|
ErrorInfo = "Task cancel Retry:" + cancelRetryReason ?? "No Reason",
|
|
|
|
|
Tag = reqModel.Tag,
|
|
|
|
|
UploaderID = mTaskIdentify
|
|
|
|
|
});;
|
|
|
|
|
});
|
|
|
|
|
db.SaveChanges();
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
Console.WriteLine($"{DateTime.UtcNow.ToString("mm:ss.fff")} task {x} result:{success}");
|
|
|
|
|
//throw new NotImplementedException();
|
|
|
|
|
catch (Exception ex)
|
|
|
|
|
{
|
|
|
|
|
Console.WriteLine(ex.Message);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
Console.WriteLine($"{DateTime.Now.ToString("HH:mm:ss.fff")} task {x} result:{success}");
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// <summary>
|
|
|
|
@ -228,7 +238,7 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
if (reqCount >= MaxCocurrentTaskCount)
|
|
|
|
|
{
|
|
|
|
|
System.Threading.Thread.Sleep(10);
|
|
|
|
|
Console.WriteLine($"{DateTime.UtcNow.ToString("mm:ss.fff")} task exceed {MaxCocurrentTaskCount}, waiting");
|
|
|
|
|
Console.WriteLine($@"{DateTime.Now.ToString("HH:mm:ss.fff")} task exceed {MaxCocurrentTaskCount}, waiting");
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
|
int newRequestCount = MaxCocurrentTaskCount - reqCount;
|
|
|
|
@ -242,6 +252,8 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
{
|
|
|
|
|
string donotStartReason;
|
|
|
|
|
if (!mOperator.ShouldStartRequest(d, out donotStartReason))
|
|
|
|
|
{
|
|
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
using (var db = new CodeFirstDbContext())
|
|
|
|
|
{
|
|
|
|
@ -252,10 +264,16 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
RetryCount = d.RetryCount,
|
|
|
|
|
RequestData = mOperator.ConvertRequestToCachel(d.Request),
|
|
|
|
|
ErrorInfo = "Task cancel Start:" + donotStartReason ?? "No Reason",
|
|
|
|
|
Tag = d.Tag,
|
|
|
|
|
UploaderID = mTaskIdentify
|
|
|
|
|
}); ;
|
|
|
|
|
db.SaveChanges();
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
catch (Exception ex)
|
|
|
|
|
{
|
|
|
|
|
Console.WriteLine(ex.Message);
|
|
|
|
|
}
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
|
d.RetryAt = DateTime.Now;
|
|
|
|
@ -263,8 +281,6 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
var req = mOperator.UploadData(d.Request);
|
|
|
|
|
enqueRequest(req, d);
|
|
|
|
|
ThreadPool.QueueUserWorkItem(new WaitCallback((o) => req.Start()));
|
|
|
|
|
//new System.Threading.Thread(req.Start).Start();
|
|
|
|
|
//req.Start();
|
|
|
|
|
i++;
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
@ -290,7 +306,7 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
{
|
|
|
|
|
using (var db = new CodeFirstDbContext())
|
|
|
|
|
{
|
|
|
|
|
//sLog.Formf($@"cache requests,{datas.Count()}");
|
|
|
|
|
Console.WriteLine($@"cache requests,{datas.Count()}");
|
|
|
|
|
var caches = datas.Select((data) =>
|
|
|
|
|
new UploadCache
|
|
|
|
|
{
|
|
|
|
@ -307,9 +323,9 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
db.SaveChanges();
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
catch (Exception)
|
|
|
|
|
catch (Exception ex)
|
|
|
|
|
{
|
|
|
|
|
//sLog.FormErrorf(ex.Message);
|
|
|
|
|
Console.WriteLine(ex.Message);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
@ -319,7 +335,8 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
/// <param name="datas"></param>
|
|
|
|
|
private void cacheRequest(List<UploadModel<T>> datas)
|
|
|
|
|
{
|
|
|
|
|
//sLog.Formf($"cache requests count,{datas.Count()}");
|
|
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
using (var db = new CodeFirstDbContext())
|
|
|
|
|
{
|
|
|
|
|
var caches = datas.Select((data) =>
|
|
|
|
@ -330,6 +347,7 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
RetryAt = data.RetryAt,
|
|
|
|
|
RetryCount = data.RetryCount,
|
|
|
|
|
ErrorInfo = data.ErrorInfo,
|
|
|
|
|
Tag = data.Tag,
|
|
|
|
|
RequestData = mOperator.ConvertRequestToCachel(data.Request)
|
|
|
|
|
}
|
|
|
|
|
);
|
|
|
|
@ -337,6 +355,11 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
db.SaveChanges();
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
catch (Exception ex)
|
|
|
|
|
{
|
|
|
|
|
Console.WriteLine(ex.Message);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// <summary>
|
|
|
|
|
/// 从数据库中提取指定数量的请求数据
|
|
|
|
@ -355,6 +378,9 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
toGet--;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
using (var db = new CodeFirstDbContext())
|
|
|
|
|
{
|
|
|
|
|
var caches = from c in db.UploadCaches
|
|
|
|
@ -372,15 +398,19 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
RetryAt = data.RetryAt,
|
|
|
|
|
RetryCount = data.RetryCount,
|
|
|
|
|
ErrorInfo = data.ErrorInfo,
|
|
|
|
|
Tag = data.Tag,
|
|
|
|
|
Request = mOperator.ConvertCacheToRequest(data.RequestData)
|
|
|
|
|
}
|
|
|
|
|
).ToList();
|
|
|
|
|
|
|
|
|
|
ret.AddRange(modelsFromDB);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
catch (Exception ex)
|
|
|
|
|
{
|
|
|
|
|
Console.WriteLine(ex.Message);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
//if (ret.Count() != 0)
|
|
|
|
|
//sLog.Formf($"get {number} cache return {ret.Count()}");
|
|
|
|
|
return ret;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
@ -391,7 +421,6 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
{
|
|
|
|
|
if (!Interlocked.Equals(mIsCanceled, 1))
|
|
|
|
|
{
|
|
|
|
|
|
|
|
|
|
var requestData = new UploadModel<T>
|
|
|
|
|
{
|
|
|
|
|
CreateAt = DateTime.Now,
|
|
|
|
@ -403,18 +432,18 @@ namespace Ksat.Supplyment.Library.Uploader
|
|
|
|
|
};
|
|
|
|
|
if (mDatas.Count >= MaxOperatingQueueSize)
|
|
|
|
|
{
|
|
|
|
|
//sLog.Formf($"add request to swap,{data}");
|
|
|
|
|
Console.WriteLine("add request to swap");
|
|
|
|
|
mSwapDatas.Enqueue(requestData);
|
|
|
|
|
}
|
|
|
|
|
else
|
|
|
|
|
{
|
|
|
|
|
//sLog.Formf($"add request to queue,{data}");
|
|
|
|
|
Console.WriteLine("add request to queue");
|
|
|
|
|
mDatas.Enqueue(requestData);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
else
|
|
|
|
|
{
|
|
|
|
|
//sLog.Formf("task canceled");
|
|
|
|
|
Console.WriteLine("task canceled");
|
|
|
|
|
cacheRequest(new List<Tuple<T, string>>() { new Tuple<T, string>(data, tag) });
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|