add support for Base Eventing Mechanism for Store's ...

This commit is contained in:
2025-07-14 05:59:41 +03:30
parent 36765f55f0
commit 4f6a03e6c9
4 changed files with 787 additions and 506 deletions
+12
View File
@@ -0,0 +1,12 @@
namespace xModels.Base
{
public class XBaseStorableDtoEventModel<T>
{
public T Model { get; set; }
public XBaseStorableDtoEventModel(T model)
{
Model = model;
}
}
}
+84
View File
@@ -0,0 +1,84 @@
using System;
using System.Collections.Generic;
using System.Reactive.Linq;
using System.Reactive.Subjects;
using xModels.Interfaces;
namespace xModels.Base
{
public abstract class XBaseStoreEvent<T, TKey> : IXBaseStoreEvents<T, TKey>
where T : XBaseStorableDto<TKey>
{
//
private readonly ISubject<XBaseStorableDtoEventModel<T>> onAddSubject;
private readonly ISubject<XBaseStorableDtoEventModel<T>> onUpdateSubject;
private readonly ISubject<XBaseStorableDtoEventModel<T>> onRemoveSubject;
//
private readonly ISubject<XBaseStorableDtoEventModel<IEnumerable<T>>> onAddManySubject;
private readonly ISubject<XBaseStorableDtoEventModel<IEnumerable<T>>> onUpdateManySubject;
private readonly ISubject<XBaseStorableDtoEventModel<IEnumerable<T>>> onRemoveManySubject;
//
public IObservable<XBaseStorableDtoEventModel<T>> OnAddObservable => onAddSubject.AsObservable();
public IObservable<XBaseStorableDtoEventModel<T>> OnUpdateObservable => onUpdateSubject.AsObservable();
public IObservable<XBaseStorableDtoEventModel<T>> OnRemoveObservable => onRemoveSubject.AsObservable();
//
public IObservable<XBaseStorableDtoEventModel<IEnumerable<T>>> OnAddManyObservable => onAddManySubject.AsObservable();
public IObservable<XBaseStorableDtoEventModel<IEnumerable<T>>> OnUpdateManyObservable => onUpdateManySubject.AsObservable();
public IObservable<XBaseStorableDtoEventModel<IEnumerable<T>>> OnRemoveManyObservable => onRemoveManySubject.AsObservable();
//
protected XBaseStoreEvent()
{
//
onAddSubject = new ReplaySubject<XBaseStorableDtoEventModel<T>>(1);
onUpdateSubject = new ReplaySubject<XBaseStorableDtoEventModel<T>>(1);
onRemoveSubject = new ReplaySubject<XBaseStorableDtoEventModel<T>>(1);
//
onAddManySubject = new ReplaySubject<XBaseStorableDtoEventModel<IEnumerable<T>>>(1);
onUpdateManySubject = new ReplaySubject<XBaseStorableDtoEventModel<IEnumerable<T>>>(1);
onRemoveManySubject = new ReplaySubject<XBaseStorableDtoEventModel<IEnumerable<T>>>(1);
}
//
#region Base Actions ...
public void AddEvent(XBaseStorableDtoEventModel<T> model)
{
onAddSubject.OnNext(model);
}
public void AddManyEvent(XBaseStorableDtoEventModel<IEnumerable<T>> model)
{
onAddManySubject.OnNext(model);
}
public void RemoveEvent(XBaseStorableDtoEventModel<T> model)
{
onRemoveSubject.OnNext(model);
}
public void RemoveManyEvent(XBaseStorableDtoEventModel<IEnumerable<T>> model)
{
onRemoveManySubject.OnNext(model);
}
public void UpdateEvent(XBaseStorableDtoEventModel<T> model)
{
onUpdateSubject.OnNext(model);
}
public void UpdateManyEvent(XBaseStorableDtoEventModel<IEnumerable<T>> model)
{
onUpdateManySubject.OnNext(model);
}
#endregion
}
}
+30
View File
@@ -0,0 +1,30 @@
using System;
using System.Collections.Generic;
using xModels.Base;
namespace xModels.Interfaces
{
public interface IXBaseStoreEvents<T, TKey> where T : XBaseStorableDto<TKey>
{
//
void AddEvent(XBaseStorableDtoEventModel<T> model);
void UpdateEvent(XBaseStorableDtoEventModel<T> model);
void RemoveEvent(XBaseStorableDtoEventModel<T> model);
//
void AddManyEvent(XBaseStorableDtoEventModel<IEnumerable<T>> model);
void UpdateManyEvent(XBaseStorableDtoEventModel<IEnumerable<T>> model);
void RemoveManyEvent(XBaseStorableDtoEventModel<IEnumerable<T>> model);
//
IObservable<XBaseStorableDtoEventModel<T>> OnAddObservable { get; }
IObservable<XBaseStorableDtoEventModel<T>> OnUpdateObservable { get; }
IObservable<XBaseStorableDtoEventModel<T>> OnRemoveObservable { get; }
//
IObservable<XBaseStorableDtoEventModel<IEnumerable<T>>> OnAddManyObservable { get; }
IObservable<XBaseStorableDtoEventModel<IEnumerable<T>>> OnUpdateManyObservable { get; }
IObservable<XBaseStorableDtoEventModel<IEnumerable<T>>> OnRemoveManyObservable { get; }
}
}
+217 -62
View File
@@ -8,9 +8,11 @@ using xCommons.Extensions;
using xModels.Base; using xModels.Base;
using xModels.Interfaces; using xModels.Interfaces;
namespace xModels.Providers { namespace xModels.Providers
{
public abstract class XBaseInMemoryStore<T, TKey> : IXBaseStore<T, TKey> public abstract class XBaseInMemoryStore<T, TKey> : IXBaseStore<T, TKey>
where T : XBaseStorableDto<TKey> { where T : XBaseStorableDto<TKey>
{
/// <summary> /// <summary>
/// this is main store of items ... /// this is main store of items ...
/// </summary> /// </summary>
@@ -18,13 +20,27 @@ namespace xModels.Providers {
/// <returns></returns> /// <returns></returns>
private static ConcurrentBag<T> STORE = new ConcurrentBag<T>(); private static ConcurrentBag<T> STORE = new ConcurrentBag<T>();
//
private readonly IXBaseStoreEvents<T, TKey> events;
//
#region Constructor ...
protected XBaseInMemoryStore(
IXBaseStoreEvents<T, TKey> events = null
)
{
this.events = events;
}
#endregion
// //
#region Retrieve ... #region Retrieve ...
/// <summary> /// <summary>
/// Retrieve all Exists Items ... /// Retrieve all Exists Items ...
/// </summary> /// </summary>
/// <returns></returns> /// <returns></returns>
public Task<IEnumerable<T>> GetAll () { public Task<IEnumerable<T>> GetAll()
{
// //
var result = STORE.AsEnumerable(); var result = STORE.AsEnumerable();
return Task.FromResult(result); return Task.FromResult(result);
@@ -35,25 +51,31 @@ namespace xModels.Providers {
/// </summary> /// </summary>
/// <param name="key"></param> /// <param name="key"></param>
/// <returns></returns> /// <returns></returns>
public Task<T> Get (TKey key) { public Task<T> Get(TKey key)
{
// //
return IsExistsByKey(key) return IsExistsByKey(key)
.ContinueWith (isExistsTask => { .ContinueWith(isExistsTask =>
{
// //
var isExists = isExistsTask var isExists = isExistsTask
.RunTask(); .RunTask();
// //
if (!isExists) { if (!isExists)
{
return null; return null;
} }
// //
try { try
{
var result = STORE var result = STORE
.FirstOrDefault(i => GetKey(i).Equals(key)); .FirstOrDefault(i => GetKey(i).Equals(key));
return result; return result;
} catch { }
catch
{
return null; return null;
} }
}); });
@@ -64,10 +86,12 @@ namespace xModels.Providers {
/// </summary> /// </summary>
/// <param name="condition"></param> /// <param name="condition"></param>
/// <returns></returns> /// <returns></returns>
public Task<IEnumerable<T>> FindMany (Expression<Func<T, bool>> condition) { public Task<IEnumerable<T>> FindMany(Expression<Func<T, bool>> condition)
{
// //
// Validate Args ... // Validate Args ...
if (condition.IsNull ()) { if (condition.IsNull())
{
return Task.FromResult(new List<T>().AsEnumerable()); return Task.FromResult(new List<T>().AsEnumerable());
} }
@@ -84,10 +108,12 @@ namespace xModels.Providers {
/// </summary> /// </summary>
/// <param name="condition"></param> /// <param name="condition"></param>
/// <returns></returns> /// <returns></returns>
public Task<T> FindOne (Expression<Func<T, bool>> condition) { public Task<T> FindOne(Expression<Func<T, bool>> condition)
{
// //
// Validate Args ... // Validate Args ...
if (condition.IsNull ()) { if (condition.IsNull())
{
return null; return null;
} }
@@ -107,27 +133,32 @@ namespace xModels.Providers {
/// </summary> /// </summary>
/// <param name="item"></param> /// <param name="item"></param>
/// <returns></returns> /// <returns></returns>
public Task<bool> Add (T item) { public Task<bool> Add(T item)
{
// //
// Validate Args ... // Validate Args ...
if (item.IsNull () || GetKey (item).IsNull ()) { if (item.IsNull() || GetKey(item).IsNull())
{
return Task.FromResult(false); return Task.FromResult(false);
} }
// //
// Try to add item ... // Try to add item ...
return IsExistsByKey(GetKey(item)) return IsExistsByKey(GetKey(item))
.ContinueWith (isExistsTask => { .ContinueWith(isExistsTask =>
{
// //
// Check item exists in Store or not ... // Check item exists in Store or not ...
var isExists = isExistsTask.RunTask(); var isExists = isExistsTask.RunTask();
if (isExists) { if (isExists)
{
return false; return false;
} }
// //
// Add Item to Store ... // Add Item to Store ...
STORE.Add(item); STORE.Add(item);
OnAdd(item);
// //
return true; return true;
@@ -139,13 +170,15 @@ namespace xModels.Providers {
/// </summary> /// </summary>
/// <param name="items"></param> /// <param name="items"></param>
/// <returns></returns> /// <returns></returns>
public Task<bool> AddMany (XBaseRangeRequest<T> items) { public Task<bool> AddMany(XBaseRangeRequest<T> items)
{
// //
// Validate Args ... // Validate Args ...
if ( if (
items.IsNull() || items.IsNull() ||
!items.Items.HasChild() !items.Items.HasChild()
) { )
{
return Task.FromResult(false); return Task.FromResult(false);
} }
@@ -156,21 +189,28 @@ namespace xModels.Providers {
.IsNull() && !STORE .IsNull() && !STORE
.Any(si => GetKey(i).Equals(GetKey(si))) .Any(si => GetKey(i).Equals(GetKey(si)))
); );
if (!mustAddItems.HasChild ()) { if (!mustAddItems.HasChild())
{
return Task.FromResult(false); return Task.FromResult(false);
} }
// //
// Add Items ... // Add Items ...
try { try
{
// //
mustAddItems mustAddItems
.ToList() .ToList()
.ForEach(mi => STORE.Add(mi)); .ForEach(mi => STORE.Add(mi));
//
OnAddMany(mustAddItems);
// //
return Task.FromResult(true); return Task.FromResult(true);
} catch { }
catch
{
return Task.FromResult(false); return Task.FromResult(false);
} }
} }
@@ -183,7 +223,8 @@ namespace xModels.Providers {
/// </summary> /// </summary>
/// <param name="item"></param> /// <param name="item"></param>
/// <returns></returns> /// <returns></returns>
public Task<bool> Remove (T item) { public Task<bool> Remove(T item)
{
// //
// Validate Args ... // Validate Args ...
if ( if (
@@ -191,16 +232,20 @@ namespace xModels.Providers {
!STORE !STORE
.Any(i => GetKey(i) .Any(i => GetKey(i)
.Equals(GetKey(item))) .Equals(GetKey(item)))
) { )
{
return Task.FromResult(false); return Task.FromResult(false);
} }
// //
try { try
{
// //
RemoveItemFromStore(item); RemoveItemFromStore(item);
return Task.FromResult(true); return Task.FromResult(true);
} catch { }
catch
{
return Task.FromResult(false); return Task.FromResult(false);
} }
} }
@@ -210,9 +255,11 @@ namespace xModels.Providers {
/// </summary> /// </summary>
/// <param name="key"></param> /// <param name="key"></param>
/// <returns></returns> /// <returns></returns>
public Task<bool> RemoveByKey (TKey key) { public Task<bool> RemoveByKey(TKey key)
{
return Get(key) return Get(key)
.ContinueWith (getTask => { .ContinueWith(getTask =>
{
// //
// retrieve item by it's Key ... // retrieve item by it's Key ...
var item = getTask var item = getTask
@@ -220,18 +267,22 @@ namespace xModels.Providers {
// //
// Check item exists ... // Check item exists ...
if (item.IsNull ()) { if (item.IsNull())
{
return false; return false;
} }
// //
// Try to remove item ... // Try to remove item ...
var result = false; var result = false;
try { try
{
// //
result = Remove(item) result = Remove(item)
.RunTask(); .RunTask();
} catch { }
catch
{
result = false; result = false;
} }
@@ -245,13 +296,15 @@ namespace xModels.Providers {
/// </summary> /// </summary>
/// <param name="items"></param> /// <param name="items"></param>
/// <returns></returns> /// <returns></returns>
public Task<bool> RemoveMany (XBaseRangeRequest<T> items) { public Task<bool> RemoveMany(XBaseRangeRequest<T> items)
{
// //
// Validate Args ... // Validate Args ...
if ( if (
items.IsNull() || items.IsNull() ||
!items.Items.HasChild() !items.Items.HasChild()
) { )
{
return Task.FromResult(false); return Task.FromResult(false);
} }
@@ -260,16 +313,20 @@ namespace xModels.Providers {
.Where(i => STORE .Where(i => STORE
.Any(si => GetKey(si).Equals(GetKey(i))) .Any(si => GetKey(si).Equals(GetKey(i)))
); );
if (!mustRemovedItems.HasChild ()) { if (!mustRemovedItems.HasChild())
{
return Task.FromResult(false); return Task.FromResult(false);
} }
// //
try { try
{
// //
RemoveItemsFromStore(mustRemovedItems); RemoveItemsFromStore(mustRemovedItems);
return Task.FromResult(true); return Task.FromResult(true);
} catch { }
catch
{
return Task.FromResult(false); return Task.FromResult(false);
} }
} }
@@ -292,21 +349,25 @@ namespace xModels.Providers {
ICollection<string> propertyBlackList = null, ICollection<string> propertyBlackList = null,
ICollection<KeyValuePair<string, Func<T, object>>> propertyValueProviders = null, ICollection<KeyValuePair<string, Func<T, object>>> propertyValueProviders = null,
bool updateWithNullOrEmptyValues = false bool updateWithNullOrEmptyValues = false
) { )
{
return Get(GetKey(item)) return Get(GetKey(item))
.ContinueWith (getTask => { .ContinueWith(getTask =>
{
// //
var existsItem = getTask var existsItem = getTask
.RunTask(); .RunTask();
// //
// Validate item Exists ... // Validate item Exists ...
if (existsItem.IsNull ()) { if (existsItem.IsNull())
{
return null; return null;
} }
// //
try { try
{
// //
var updatedItem = existsItem; var updatedItem = existsItem;
updatedItem.UpdateData( updatedItem.UpdateData(
@@ -321,10 +382,13 @@ namespace xModels.Providers {
// //
RemoveItemFromStore(existsItem); RemoveItemFromStore(existsItem);
STORE.Add(updatedItem); STORE.Add(updatedItem);
OnUpdate(item);
// //
return updatedItem; return updatedItem;
} catch { }
catch
{
return null; return null;
} }
}); });
@@ -345,13 +409,15 @@ namespace xModels.Providers {
ICollection<string> propertyBlackList = null, ICollection<string> propertyBlackList = null,
ICollection<KeyValuePair<string, Func<T, object>>> propertyValueProviders = null, ICollection<KeyValuePair<string, Func<T, object>>> propertyValueProviders = null,
bool updateWithNullOrEmptyValues = false bool updateWithNullOrEmptyValues = false
) { )
{
// //
// Validate Args ... // Validate Args ...
if ( if (
items.IsNull() || items.IsNull() ||
!items.Items.HasChild() !items.Items.HasChild()
) { )
{
return Task.FromResult(false); return Task.FromResult(false);
} }
@@ -363,7 +429,8 @@ namespace xModels.Providers {
.Any( .Any(
si => GetKey(i).Equals(GetKey(si))) si => GetKey(i).Equals(GetKey(si)))
); );
if (!mustUpdateItems.HasChild ()) { if (!mustUpdateItems.HasChild())
{
return Task.FromResult(false); return Task.FromResult(false);
} }
@@ -377,15 +444,21 @@ namespace xModels.Providers {
updateWithNullOrEmptyValues: updateWithNullOrEmptyValues updateWithNullOrEmptyValues: updateWithNullOrEmptyValues
)); ));
return Task.WhenAll(mustUpdateTasks) return Task.WhenAll(mustUpdateTasks)
.ContinueWith (updateTasks => { .ContinueWith(updateTasks =>
{
// //
var updatedItems = updateTasks var updatedItems = updateTasks
.RunTask(); .RunTask();
// //
if (!updatedItems.HasChild ()) { if (!updatedItems.HasChild())
{
return false; return false;
} else { }
else
{
//
OnUpdateMany(mustUpdateItems);
return true; return true;
} }
}); });
@@ -399,10 +472,12 @@ namespace xModels.Providers {
/// </summary> /// </summary>
/// <param name="condition"></param> /// <param name="condition"></param>
/// <returns></returns> /// <returns></returns>
public Task<bool> IsExists (Expression<Func<T, bool>> condition) { public Task<bool> IsExists(Expression<Func<T, bool>> condition)
{
// //
// Validate Args ... // Validate Args ...
if (condition.IsNull ()) { if (condition.IsNull())
{
return Task.FromResult(false); return Task.FromResult(false);
} }
@@ -419,8 +494,10 @@ namespace xModels.Providers {
/// </summary> /// </summary>
/// <param name="key"></param> /// <param name="key"></param>
/// <returns></returns> /// <returns></returns>
public Task<bool> IsExistsByKey (TKey key) { public Task<bool> IsExistsByKey(TKey key)
return Task.Run (() => { {
return Task.Run(() =>
{
return STORE.Any(i => GetKey(i).Equals(key)); return STORE.Any(i => GetKey(i).Equals(key));
}); });
} }
@@ -432,7 +509,8 @@ namespace xModels.Providers {
/// Count Exists items ... /// Count Exists items ...
/// </summary> /// </summary>
/// <returns></returns> /// <returns></returns>
public Task<int> Count () { public Task<int> Count()
{
// //
var result = STORE.Count(); var result = STORE.Count();
return Task.FromResult(result); return Task.FromResult(result);
@@ -449,11 +527,13 @@ namespace xModels.Providers {
public void SetKey( public void SetKey(
ref T item, ref T item,
TKey id TKey id
) { )
{
// //
var props = item.GetType().GetProperties(); var props = item.GetType().GetProperties();
var keyProp = props.FirstOrDefault(p => p.Name == "Id"); var keyProp = props.FirstOrDefault(p => p.Name == "Id");
if (keyProp.IsNull ()) { if (keyProp.IsNull())
{
return; return;
} }
@@ -468,35 +548,44 @@ namespace xModels.Providers {
/// </summary> /// </summary>
/// <param name="item"></param> /// <param name="item"></param>
/// <returns></returns> /// <returns></returns>
public TKey GetKey (T item) { public TKey GetKey(T item)
{
// //
var props = item.GetType().GetProperties(); var props = item.GetType().GetProperties();
var keyProp = props.FirstOrDefault(p => p.Name == "Id"); var keyProp = props.FirstOrDefault(p => p.Name == "Id");
// //
var keyString = string.Empty; var keyString = string.Empty;
if (keyProp.IsNull ()) { if (keyProp.IsNull())
{
keyString = string.Empty; keyString = string.Empty;
} else { }
else
{
keyString = keyProp.GetValue(item).ToString(); keyString = keyProp.GetValue(item).ToString();
} }
// //
if (keyString.IsNullOrEmpty ()) { if (keyString.IsNullOrEmpty())
{
return default(TKey); return default(TKey);
} }
// //
// Prevent Deserializing issues throug JsonReader ... // Prevent Deserializing issues throug JsonReader ...
if (keyString.IsGuid () && typeof (TKey) == typeof (Guid)) { if (keyString.IsGuid() && typeof(TKey) == typeof(Guid))
{
return item.Id; return item.Id;
} }
// //
TKey result; TKey result;
try { try
{
result = keyString.FromJSON<TKey>(); result = keyString.FromJSON<TKey>();
} catch { }
catch
{
result = keyString.ConvertTo<TKey>(); result = keyString.ConvertTo<TKey>();
} }
return result; return result;
@@ -505,17 +594,83 @@ namespace xModels.Providers {
// //
#region Private ... #region Private ...
private void RemoveItemFromStore (T item) { private void RemoveItemFromStore(T item)
{
//
STORE = new ConcurrentBag<T>( STORE = new ConcurrentBag<T>(
STORE.Except(new[] { item }) STORE.Except(new[] { item })
); );
//
OnRemove(item);
} }
private void RemoveItemsFromStore (IEnumerable<T> items) { private void RemoveItemsFromStore(IEnumerable<T> items)
{
//
STORE = new ConcurrentBag<T>( STORE = new ConcurrentBag<T>(
STORE.Except(items) STORE.Except(items)
); );
//
OnRemoveMany(items);
} }
//
#region Event Notifiers ...
private void OnAdd(T item)
{
//
if (!events.IsNull())
{
events.AddEvent(new XBaseStorableDtoEventModel<T>(item));
}
}
private void OnUpdate(T item)
{
//
if (!events.IsNull())
{
events.UpdateEvent(new XBaseStorableDtoEventModel<T>(item));
}
}
private void OnRemove(T item)
{
//
if (!events.IsNull())
{
events.RemoveEvent(new XBaseStorableDtoEventModel<T>(item));
}
}
private void OnAddMany(IEnumerable<T> items)
{
//
if (!events.IsNull())
{
events.AddManyEvent(new XBaseStorableDtoEventModel<IEnumerable<T>>(items));
}
}
private void OnUpdateMany(IEnumerable<T> items)
{
//
if (!events.IsNull())
{
events.UpdateManyEvent(new XBaseStorableDtoEventModel<IEnumerable<T>>(items));
}
}
private void OnRemoveMany(IEnumerable<T> items)
{
//
if (!events.IsNull())
{
events.RemoveManyEvent(new XBaseStorableDtoEventModel<IEnumerable<T>>(items));
}
}
#endregion
#endregion #endregion
} }
} }