2019-11-09 06:37:05 +00:00
|
|
|
|
using System;
|
|
|
|
|
using System.Collections.Generic;
|
|
|
|
|
using System.Linq;
|
2019-11-09 21:29:55 +00:00
|
|
|
|
using System.Reactive.Subjects;
|
2019-11-09 06:37:05 +00:00
|
|
|
|
using System.Text;
|
2019-11-11 03:47:25 +00:00
|
|
|
|
using System.Threading;
|
2019-11-09 06:37:05 +00:00
|
|
|
|
using System.Threading.Tasks;
|
2019-11-10 00:22:28 +00:00
|
|
|
|
using System.Windows.Forms.VisualStyles;
|
2019-11-09 06:37:05 +00:00
|
|
|
|
|
|
|
|
|
namespace Wabbajack.Common.CSP
|
|
|
|
|
{
|
|
|
|
|
public static class CSPExtensions
|
|
|
|
|
{
|
|
|
|
|
public static async Task OntoChannel<TIn, TOut>(this IEnumerable<TIn> coll, IChannel<TIn, TOut> chan)
|
|
|
|
|
{
|
|
|
|
|
foreach (var val in coll)
|
|
|
|
|
{
|
|
|
|
|
if (!await chan.Put(val)) break;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// <summary>
|
|
|
|
|
/// Turns a IEnumerable collection into a channel. Note, computation of the enumerable will happen inside
|
|
|
|
|
/// the lock of the channel, so try to keep the work of the enumerable light.
|
|
|
|
|
/// </summary>
|
|
|
|
|
/// <typeparam name="T"></typeparam>
|
|
|
|
|
/// <param name="coll">Collection to spool out of the channel.</param>
|
|
|
|
|
/// <returns></returns>
|
|
|
|
|
public static IChannel<T, T> ToChannel<T>(this IEnumerable<T> coll)
|
|
|
|
|
{
|
2019-11-09 21:29:55 +00:00
|
|
|
|
var chan = Channel.Create(coll.GetEnumerator());
|
|
|
|
|
chan.Close();
|
|
|
|
|
return chan;
|
2019-11-09 06:37:05 +00:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
2019-11-09 21:29:55 +00:00
|
|
|
|
/// <summary>
|
|
|
|
|
/// Takes all the values from chan, once the channel closes returns a List of the values taken.
|
|
|
|
|
/// </summary>
|
|
|
|
|
/// <typeparam name="TOut"></typeparam>
|
|
|
|
|
/// <typeparam name="TIn"></typeparam>
|
|
|
|
|
/// <param name="chan"></param>
|
|
|
|
|
/// <returns></returns>
|
2019-11-09 06:37:05 +00:00
|
|
|
|
public static async Task<List<TOut>> TakeAll<TOut, TIn>(this IChannel<TIn, TOut> chan)
|
|
|
|
|
{
|
|
|
|
|
List<TOut> acc = new List<TOut>();
|
|
|
|
|
while (true)
|
|
|
|
|
{
|
2019-11-09 14:49:00 +00:00
|
|
|
|
var (open, val) = await chan.Take();
|
|
|
|
|
|
|
|
|
|
if (!open) break;
|
|
|
|
|
|
|
|
|
|
acc.Add(val);
|
2019-11-09 06:37:05 +00:00
|
|
|
|
}
|
|
|
|
|
return acc;
|
|
|
|
|
}
|
2019-11-09 21:29:55 +00:00
|
|
|
|
|
|
|
|
|
|
|
|
|
|
/// <summary>
|
|
|
|
|
/// Pipes values from `from` into `to`
|
|
|
|
|
/// </summary>
|
|
|
|
|
/// <typeparam name="TIn"></typeparam>
|
|
|
|
|
/// <typeparam name="TMid"></typeparam>
|
|
|
|
|
/// <typeparam name="TOut"></typeparam>
|
|
|
|
|
/// <param name="from">source channel</param>
|
|
|
|
|
/// <param name="to">destination channel</param>
|
|
|
|
|
/// <param name="closeOnFinished">Tf true, will close the other channel when one channel closes</param>
|
|
|
|
|
/// <returns></returns>
|
|
|
|
|
public static async Task Pipe<TIn, TMid, TOut>(this IChannel<TIn, TMid> from, IChannel<TMid, TOut> to, bool closeOnFinished = true)
|
|
|
|
|
{
|
|
|
|
|
while (true)
|
|
|
|
|
{
|
|
|
|
|
var (isFromOpen, val) = await from.Take();
|
|
|
|
|
if (isFromOpen)
|
|
|
|
|
{
|
|
|
|
|
var isToOpen = await to.Put(val);
|
|
|
|
|
if (isToOpen) continue;
|
|
|
|
|
if (closeOnFinished)
|
|
|
|
|
@from.Close();
|
|
|
|
|
break;
|
|
|
|
|
}
|
|
|
|
|
if (closeOnFinished)
|
|
|
|
|
to.Close();
|
|
|
|
|
break;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2019-11-11 03:47:25 +00:00
|
|
|
|
public static Task<T> ThreadedTask<T>(Func<T> action)
|
2019-11-09 21:29:55 +00:00
|
|
|
|
{
|
2019-11-11 03:47:25 +00:00
|
|
|
|
var src = new TaskCompletionSource<T>();
|
|
|
|
|
var th = new Thread(() =>
|
2019-11-09 21:29:55 +00:00
|
|
|
|
{
|
2019-11-11 03:47:25 +00:00
|
|
|
|
try
|
2019-11-09 21:29:55 +00:00
|
|
|
|
{
|
2019-11-11 03:47:25 +00:00
|
|
|
|
src.SetResult(action());
|
2019-11-09 21:29:55 +00:00
|
|
|
|
}
|
2019-11-11 03:47:25 +00:00
|
|
|
|
catch (Exception ex)
|
|
|
|
|
{
|
|
|
|
|
src.SetException(ex);
|
|
|
|
|
}
|
|
|
|
|
}) {Priority = ThreadPriority.BelowNormal};
|
|
|
|
|
th.Start();
|
|
|
|
|
return src.Task;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public static Task ThreadedTask<T>(Action action)
|
|
|
|
|
{
|
|
|
|
|
var src = new TaskCompletionSource<bool>();
|
|
|
|
|
var th = new Thread(() =>
|
|
|
|
|
{
|
|
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
action();
|
|
|
|
|
src.SetResult(true);
|
|
|
|
|
}
|
|
|
|
|
catch (Exception ex)
|
|
|
|
|
{
|
|
|
|
|
src.SetException(ex);
|
|
|
|
|
}
|
|
|
|
|
})
|
|
|
|
|
{ Priority = ThreadPriority.BelowNormal };
|
|
|
|
|
th.Start();
|
|
|
|
|
return src.Task;
|
|
|
|
|
}
|
2019-11-09 21:29:55 +00:00
|
|
|
|
|
2019-11-09 06:37:05 +00:00
|
|
|
|
}
|
2019-11-10 00:22:28 +00:00
|
|
|
|
|
2019-11-09 06:37:05 +00:00
|
|
|
|
}
|