17 changed files with 885 additions and 9 deletions
@ -1,13 +1,16 @@ |
|||
using Volo.Abp.Collections; |
|||
using System.Collections.Generic; |
|||
using Volo.Abp.Collections; |
|||
|
|||
namespace LINGYUN.Abp.AI.Tools; |
|||
public class AbpAIToolsOptions |
|||
{ |
|||
public ITypeList<IAIToolDefinitionProvider> DefinitionProviders { get; } |
|||
public ITypeList<IAIToolProvider> AIToolProviders { get; } |
|||
public HashSet<string> DeletedAITools { get; } |
|||
public AbpAIToolsOptions() |
|||
{ |
|||
DefinitionProviders = new TypeList<IAIToolDefinitionProvider>(); |
|||
AIToolProviders = new TypeList<IAIToolProvider>(); |
|||
DeletedAITools = new HashSet<string>(); |
|||
} |
|||
} |
|||
|
|||
@ -0,0 +1,67 @@ |
|||
using LINGYUN.Abp.AI.Tools; |
|||
using System.Collections.Generic; |
|||
using System.Globalization; |
|||
using System.Threading.Tasks; |
|||
using Volo.Abp.Data; |
|||
using Volo.Abp.DependencyInjection; |
|||
using Volo.Abp.Guids; |
|||
using Volo.Abp.Localization; |
|||
using Volo.Abp.SimpleStateChecking; |
|||
|
|||
namespace LINGYUN.Abp.AIManagement.Tools; |
|||
public class AIToolDefinitionSerializer : IAIToolDefinitionSerializer, ITransientDependency |
|||
{ |
|||
protected IGuidGenerator GuidGenerator { get; } |
|||
protected ISimpleStateCheckerSerializer StateCheckerSerializer { get; } |
|||
protected ILocalizableStringSerializer LocalizableStringSerializer { get; } |
|||
|
|||
public AIToolDefinitionSerializer( |
|||
IGuidGenerator guidGenerator, |
|||
ISimpleStateCheckerSerializer stateCheckerSerializer, |
|||
ILocalizableStringSerializer localizableStringSerializer) |
|||
{ |
|||
GuidGenerator = guidGenerator; |
|||
StateCheckerSerializer = stateCheckerSerializer; |
|||
LocalizableStringSerializer = localizableStringSerializer; |
|||
} |
|||
|
|||
public async virtual Task<AIToolDefinitionRecord[]> SerializeAsync(IEnumerable<AIToolDefinition> definitions) |
|||
{ |
|||
var records = new List<AIToolDefinitionRecord>(); |
|||
foreach (var aiToolDef in definitions) |
|||
{ |
|||
records.Add(await SerializeAsync(aiToolDef)); |
|||
} |
|||
|
|||
return records.ToArray(); |
|||
} |
|||
|
|||
public virtual Task<AIToolDefinitionRecord> SerializeAsync(AIToolDefinition definition) |
|||
{ |
|||
using (CultureHelper.Use(CultureInfo.InvariantCulture)) |
|||
{ |
|||
var aiToolRecord = new AIToolDefinitionRecord( |
|||
GuidGenerator.Create(), |
|||
definition.Name, |
|||
definition.Provider, |
|||
definition.Description != null ? LocalizableStringSerializer.Serialize(definition.Description) : null, |
|||
SerializeStateCheckers(definition.StateCheckers)); |
|||
|
|||
|
|||
foreach (var property in definition.Properties) |
|||
{ |
|||
aiToolRecord.SetProperty(property.Key, property.Value); |
|||
} |
|||
|
|||
aiToolRecord.IsEnabled = definition.IsEnabled; |
|||
aiToolRecord.IsSystem = true; |
|||
|
|||
return Task.FromResult(aiToolRecord); |
|||
} |
|||
} |
|||
|
|||
protected virtual string? SerializeStateCheckers(List<ISimpleStateChecker<AIToolDefinition>> stateCheckers) |
|||
{ |
|||
return StateCheckerSerializer.Serialize(stateCheckers); |
|||
} |
|||
} |
|||
@ -0,0 +1,138 @@ |
|||
using LINGYUN.Abp.AI.Tools; |
|||
using Microsoft.Extensions.DependencyInjection; |
|||
using Microsoft.Extensions.Hosting; |
|||
using Microsoft.Extensions.Logging; |
|||
using Microsoft.Extensions.Logging.Abstractions; |
|||
using Microsoft.Extensions.Options; |
|||
using Polly; |
|||
using System; |
|||
using System.Threading; |
|||
using System.Threading.Tasks; |
|||
using Volo.Abp.DependencyInjection; |
|||
using Volo.Abp.Threading; |
|||
|
|||
namespace LINGYUN.Abp.AIManagement.Tools; |
|||
public class AIToolDynamicInitializer : ITransientDependency |
|||
{ |
|||
public ILogger<AIToolDynamicInitializer> Logger { get; set; } |
|||
|
|||
protected IServiceProvider ServiceProvider { get; } |
|||
|
|||
public AIToolDynamicInitializer(IServiceProvider serviceProvider) |
|||
{ |
|||
Logger = NullLogger<AIToolDynamicInitializer>.Instance; |
|||
|
|||
ServiceProvider = serviceProvider; |
|||
} |
|||
|
|||
public virtual Task InitializeAsync(bool runInBackground, CancellationToken cancellationToken = default) |
|||
{ |
|||
var options = ServiceProvider.GetRequiredService<IOptions<AIManagementOptions>>().Value; |
|||
|
|||
if (!options.SaveStaticAIToolsToDatabase && !options.IsDynamicAIToolStoreEnabled) |
|||
{ |
|||
return Task.CompletedTask; |
|||
} |
|||
|
|||
if (runInBackground) |
|||
{ |
|||
var applicationLifetime = ServiceProvider.GetService<IHostApplicationLifetime>(); |
|||
Task.Run(async () => |
|||
{ |
|||
if (cancellationToken == default && applicationLifetime?.ApplicationStopping != null) |
|||
{ |
|||
cancellationToken = applicationLifetime.ApplicationStopping; |
|||
} |
|||
await ExecuteInitializationAsync(options, cancellationToken); |
|||
}, cancellationToken); |
|||
|
|||
return Task.CompletedTask; |
|||
} |
|||
|
|||
return ExecuteInitializationAsync(options, cancellationToken); |
|||
} |
|||
|
|||
protected virtual async Task ExecuteInitializationAsync(AIManagementOptions options, CancellationToken cancellationToken) |
|||
{ |
|||
try |
|||
{ |
|||
var cancellationTokenProvider = ServiceProvider.GetRequiredService<ICancellationTokenProvider>(); |
|||
using (cancellationTokenProvider.Use(cancellationToken)) |
|||
{ |
|||
if (cancellationTokenProvider.Token.IsCancellationRequested) |
|||
{ |
|||
return; |
|||
} |
|||
|
|||
await SaveStaticAIToolsToDatabaseAsync(options, cancellationToken); |
|||
|
|||
if (cancellationTokenProvider.Token.IsCancellationRequested) |
|||
{ |
|||
return; |
|||
} |
|||
|
|||
await PreCacheDynamicAIToolsAsync(options); |
|||
} |
|||
} |
|||
catch |
|||
{ |
|||
// No need to log here since inner calls log
|
|||
} |
|||
} |
|||
|
|||
protected virtual async Task SaveStaticAIToolsToDatabaseAsync( |
|||
AIManagementOptions options, |
|||
CancellationToken cancellationToken) |
|||
{ |
|||
if (!options.SaveStaticAIToolsToDatabase) |
|||
{ |
|||
return; |
|||
} |
|||
|
|||
var staticAIToolSaver = ServiceProvider.GetRequiredService<IStaticAIToolSaver>(); |
|||
|
|||
await Policy |
|||
.Handle<Exception>(ex => ex is not OperationCanceledException) |
|||
.WaitAndRetryAsync( |
|||
8, |
|||
retryAttempt => TimeSpan.FromSeconds( |
|||
Volo.Abp.RandomHelper.GetRandom( |
|||
(int)Math.Pow(2, retryAttempt) * 8, |
|||
(int)Math.Pow(2, retryAttempt) * 12) |
|||
) |
|||
) |
|||
.ExecuteAsync(async _ => |
|||
{ |
|||
try |
|||
{ |
|||
await staticAIToolSaver.SaveAsync(); |
|||
} |
|||
catch (Exception ex) |
|||
{ |
|||
Logger.LogException(ex); |
|||
throw; // Polly will catch it
|
|||
} |
|||
}, cancellationToken); |
|||
} |
|||
|
|||
protected virtual async Task PreCacheDynamicAIToolsAsync(AIManagementOptions options) |
|||
{ |
|||
if (!options.IsDynamicAIToolStoreEnabled) |
|||
{ |
|||
return; |
|||
} |
|||
|
|||
var dynamicAIToolDefinitionStore = ServiceProvider.GetRequiredService<IDynamicAIToolDefinitionStore>(); |
|||
|
|||
try |
|||
{ |
|||
// Pre-cache AITools, so first request doesn't wait
|
|||
await dynamicAIToolDefinitionStore.GetAllAsync(); |
|||
} |
|||
catch (Exception ex) |
|||
{ |
|||
Logger.LogException(ex); |
|||
throw; // It will be cached in Initialize()
|
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,160 @@ |
|||
using JetBrains.Annotations; |
|||
using LINGYUN.Abp.AI.Tools; |
|||
using Microsoft.Extensions.Caching.Distributed; |
|||
using Microsoft.Extensions.Options; |
|||
using System; |
|||
using System.Collections.Generic; |
|||
using System.Threading.Tasks; |
|||
using Volo.Abp; |
|||
using Volo.Abp.Caching; |
|||
using Volo.Abp.DependencyInjection; |
|||
using Volo.Abp.DistributedLocking; |
|||
using Volo.Abp.Threading; |
|||
|
|||
namespace LINGYUN.Abp.AIManagement.Tools; |
|||
|
|||
[Dependency(ReplaceServices = true)] |
|||
public class DynamicAIToolDefinitionStore : IDynamicAIToolDefinitionStore, ITransientDependency |
|||
{ |
|||
protected IAIToolDefinitionRecordRepository AIToolDefinitionRecordRepository { get; } |
|||
protected IAIToolDefinitionSerializer AIToolDefinitionSerializer { get; } |
|||
protected IDynamicAIToolDefinitionStoreInMemoryCache StoreCache { get; } |
|||
protected IDistributedCache DistributedCache { get; } |
|||
protected IAbpDistributedLock DistributedLock { get; } |
|||
public AIManagementOptions AIManagementOptions { get; } |
|||
protected AbpDistributedCacheOptions CacheOptions { get; } |
|||
|
|||
public DynamicAIToolDefinitionStore( |
|||
IAIToolDefinitionRecordRepository aiToolDefinitionRecordRepository, |
|||
IAIToolDefinitionSerializer aiToolDefinitionSerializer, |
|||
IDynamicAIToolDefinitionStoreInMemoryCache storeCache, |
|||
IDistributedCache distributedCache, |
|||
IOptions<AbpDistributedCacheOptions> cacheOptions, |
|||
IOptions<AIManagementOptions> aiManagementOptions, |
|||
IAbpDistributedLock distributedLock) |
|||
{ |
|||
AIToolDefinitionRecordRepository = aiToolDefinitionRecordRepository; |
|||
AIToolDefinitionSerializer = aiToolDefinitionSerializer; |
|||
StoreCache = storeCache; |
|||
DistributedCache = distributedCache; |
|||
DistributedLock = distributedLock; |
|||
AIManagementOptions = aiManagementOptions.Value; |
|||
CacheOptions = cacheOptions.Value; |
|||
} |
|||
|
|||
public async virtual Task<IReadOnlyList<AIToolDefinition>> GetAllAsync() |
|||
{ |
|||
if (!AIManagementOptions.IsDynamicAIToolStoreEnabled) |
|||
{ |
|||
return Array.Empty<AIToolDefinition>(); |
|||
} |
|||
|
|||
using (await StoreCache.SyncSemaphore.LockAsync()) |
|||
{ |
|||
await EnsureCacheIsUptoDateAsync(); |
|||
return StoreCache.GetAITools(); |
|||
} |
|||
} |
|||
|
|||
public async virtual Task<AIToolDefinition> GetAsync([NotNull] string name) |
|||
{ |
|||
Check.NotNull(name, nameof(name)); |
|||
|
|||
return await GetOrNullAsync(name) ?? throw new AbpException("Undefined AITool: " + name); |
|||
} |
|||
|
|||
public async virtual Task<AIToolDefinition?> GetOrNullAsync([NotNull] string name) |
|||
{ |
|||
Check.NotNull(name, nameof(name)); |
|||
|
|||
if (!AIManagementOptions.IsDynamicAIToolStoreEnabled) |
|||
{ |
|||
return null; |
|||
} |
|||
|
|||
using (await StoreCache.SyncSemaphore.LockAsync()) |
|||
{ |
|||
await EnsureCacheIsUptoDateAsync(); |
|||
return StoreCache.GetAIToolOrNull(name); |
|||
} |
|||
} |
|||
protected virtual async Task EnsureCacheIsUptoDateAsync() |
|||
{ |
|||
if (StoreCache.LastCheckTime.HasValue && |
|||
DateTime.Now.Subtract(StoreCache.LastCheckTime.Value).TotalSeconds < 30) |
|||
{ |
|||
return; |
|||
} |
|||
|
|||
var stampInDistributedCache = await GetOrSetStampInDistributedCache(); |
|||
|
|||
if (stampInDistributedCache == StoreCache.CacheStamp) |
|||
{ |
|||
StoreCache.LastCheckTime = DateTime.Now; |
|||
return; |
|||
} |
|||
|
|||
await UpdateInMemoryStoreCache(); |
|||
|
|||
StoreCache.CacheStamp = stampInDistributedCache; |
|||
StoreCache.LastCheckTime = DateTime.Now; |
|||
} |
|||
|
|||
protected virtual async Task UpdateInMemoryStoreCache() |
|||
{ |
|||
var workspaces = await AIToolDefinitionRecordRepository.GetListAsync(); |
|||
|
|||
await StoreCache.FillAsync(workspaces); |
|||
} |
|||
|
|||
protected virtual async Task<string> GetOrSetStampInDistributedCache() |
|||
{ |
|||
var cacheKey = GetCommonStampCacheKey(); |
|||
|
|||
var stampInDistributedCache = await DistributedCache.GetStringAsync(cacheKey); |
|||
if (stampInDistributedCache != null) |
|||
{ |
|||
return stampInDistributedCache; |
|||
} |
|||
|
|||
await using (var commonLockHandle = await DistributedLock |
|||
.TryAcquireAsync(GetCommonDistributedLockKey(), TimeSpan.FromMinutes(2))) |
|||
{ |
|||
if (commonLockHandle == null) |
|||
{ |
|||
throw new AbpException( |
|||
"Could not acquire distributed lock for AITool definition common stamp check!" |
|||
); |
|||
} |
|||
|
|||
stampInDistributedCache = await DistributedCache.GetStringAsync(cacheKey); |
|||
if (stampInDistributedCache != null) |
|||
{ |
|||
return stampInDistributedCache; |
|||
} |
|||
|
|||
stampInDistributedCache = Guid.NewGuid().ToString(); |
|||
|
|||
await DistributedCache.SetStringAsync( |
|||
cacheKey, |
|||
stampInDistributedCache, |
|||
new DistributedCacheEntryOptions |
|||
{ |
|||
SlidingExpiration = TimeSpan.FromDays(30) |
|||
} |
|||
); |
|||
} |
|||
|
|||
return stampInDistributedCache; |
|||
} |
|||
|
|||
protected virtual string GetCommonStampCacheKey() |
|||
{ |
|||
return $"{CacheOptions.KeyPrefix}_AbpInMemoryAIToolCacheStamp"; |
|||
} |
|||
|
|||
protected virtual string GetCommonDistributedLockKey() |
|||
{ |
|||
return $"{CacheOptions.KeyPrefix}_Common_AbpAIToolUpdateLock"; |
|||
} |
|||
} |
|||
@ -0,0 +1,57 @@ |
|||
using Microsoft.Extensions.Caching.Distributed; |
|||
using Microsoft.Extensions.Options; |
|||
using System; |
|||
using System.Threading.Tasks; |
|||
using Volo.Abp.Caching; |
|||
using Volo.Abp.DependencyInjection; |
|||
using Volo.Abp.Domain.Entities.Events; |
|||
using Volo.Abp.EventBus; |
|||
using Volo.Abp.Threading; |
|||
using Volo.Abp.Timing; |
|||
|
|||
namespace LINGYUN.Abp.AIManagement.Tools; |
|||
public class DynamicAIToolDefinitionStoreCacheInvalidator : |
|||
ILocalEventHandler<EntityChangedEventData<AIToolDefinitionRecord>>, |
|||
ITransientDependency |
|||
{ |
|||
private readonly IDynamicAIToolDefinitionStoreInMemoryCache _storeCache; |
|||
|
|||
private readonly IClock _clock; |
|||
private readonly IDistributedCache _distributedCache; |
|||
private readonly AbpDistributedCacheOptions _cacheOptions; |
|||
|
|||
public DynamicAIToolDefinitionStoreCacheInvalidator( |
|||
IClock clock, |
|||
IDistributedCache distributedCache, |
|||
IDynamicAIToolDefinitionStoreInMemoryCache storeCache, |
|||
IOptions<AbpDistributedCacheOptions> cacheOptions) |
|||
{ |
|||
_storeCache = storeCache; |
|||
_clock = clock; |
|||
_distributedCache = distributedCache; |
|||
_cacheOptions = cacheOptions.Value; |
|||
} |
|||
|
|||
public async virtual Task HandleEventAsync(EntityChangedEventData<AIToolDefinitionRecord> eventData) |
|||
{ |
|||
await RemoveStampInDistributedCacheAsync(); |
|||
} |
|||
|
|||
protected async virtual Task RemoveStampInDistributedCacheAsync() |
|||
{ |
|||
using (await _storeCache.SyncSemaphore.LockAsync()) |
|||
{ |
|||
var cacheKey = GetCommonStampCacheKey(); |
|||
|
|||
await _distributedCache.RemoveAsync(cacheKey); |
|||
|
|||
_storeCache.CacheStamp = Guid.NewGuid().ToString(); |
|||
_storeCache.LastCheckTime = _clock.Now.AddMinutes(-5); |
|||
} |
|||
} |
|||
|
|||
protected virtual string GetCommonStampCacheKey() |
|||
{ |
|||
return $"{_cacheOptions.KeyPrefix}_AbpInMemoryAIToolCacheStamp"; |
|||
} |
|||
} |
|||
@ -0,0 +1,79 @@ |
|||
using LINGYUN.Abp.AI.Tools; |
|||
using System; |
|||
using System.Collections.Generic; |
|||
using System.Linq; |
|||
using System.Threading; |
|||
using System.Threading.Tasks; |
|||
using Volo.Abp.DependencyInjection; |
|||
using Volo.Abp.Localization; |
|||
using Volo.Abp.SimpleStateChecking; |
|||
|
|||
namespace LINGYUN.Abp.AIManagement.Tools; |
|||
public class DynamicAIToolDefinitionStoreInMemoryCache : IDynamicAIToolDefinitionStoreInMemoryCache, ISingletonDependency |
|||
{ |
|||
public string CacheStamp { get; set; } |
|||
protected IDictionary<string, AIToolDefinition> AIToolDefinitions { get; } |
|||
protected ISimpleStateCheckerSerializer StateCheckerSerializer { get; } |
|||
protected ILocalizableStringSerializer LocalizableStringSerializer { get; } |
|||
|
|||
public SemaphoreSlim SyncSemaphore { get; } = new(1, 1); |
|||
|
|||
public DateTime? LastCheckTime { get; set; } |
|||
|
|||
public DynamicAIToolDefinitionStoreInMemoryCache( |
|||
ISimpleStateCheckerSerializer stateCheckerSerializer, |
|||
ILocalizableStringSerializer localizableStringSerializer) |
|||
{ |
|||
StateCheckerSerializer = stateCheckerSerializer; |
|||
LocalizableStringSerializer = localizableStringSerializer; |
|||
|
|||
AIToolDefinitions = new Dictionary<string, AIToolDefinition>(); |
|||
} |
|||
|
|||
public Task FillAsync(List<AIToolDefinitionRecord> tools) |
|||
{ |
|||
AIToolDefinitions.Clear(); |
|||
|
|||
foreach (var tool in tools) |
|||
{ |
|||
var toolDef = new AIToolDefinition( |
|||
tool.Name, |
|||
tool.Provider, |
|||
!tool.Description.IsNullOrWhiteSpace() ? LocalizableStringSerializer.Deserialize(tool.Description) : null); |
|||
|
|||
toolDef.IsEnabled = tool.IsEnabled; |
|||
|
|||
if (!tool.StateCheckers.IsNullOrWhiteSpace()) |
|||
{ |
|||
var checkers = StateCheckerSerializer |
|||
.DeserializeArray( |
|||
tool.StateCheckers, |
|||
toolDef |
|||
); |
|||
toolDef.StateCheckers.AddRange(checkers); |
|||
} |
|||
|
|||
foreach (var property in tool.ExtraProperties) |
|||
{ |
|||
if (property.Value != null) |
|||
{ |
|||
toolDef.WithProperty(property.Key, property.Value); |
|||
} |
|||
} |
|||
|
|||
AIToolDefinitions[tool.Name] = toolDef; |
|||
} |
|||
|
|||
return Task.CompletedTask; |
|||
} |
|||
|
|||
public AIToolDefinition? GetAIToolOrNull(string name) |
|||
{ |
|||
return AIToolDefinitions.GetOrDefault(name); |
|||
} |
|||
|
|||
public IReadOnlyList<AIToolDefinition> GetAITools() |
|||
{ |
|||
return AIToolDefinitions.Values.ToList(); |
|||
} |
|||
} |
|||
@ -0,0 +1,11 @@ |
|||
using LINGYUN.Abp.AI.Tools; |
|||
using System.Collections.Generic; |
|||
using System.Threading.Tasks; |
|||
|
|||
namespace LINGYUN.Abp.AIManagement.Tools; |
|||
public interface IAIToolDefinitionSerializer |
|||
{ |
|||
Task<AIToolDefinitionRecord[]> SerializeAsync(IEnumerable<AIToolDefinition> definitions); |
|||
|
|||
Task<AIToolDefinitionRecord> SerializeAsync(AIToolDefinition definition); |
|||
} |
|||
@ -0,0 +1,21 @@ |
|||
using LINGYUN.Abp.AI.Tools; |
|||
using System; |
|||
using System.Collections.Generic; |
|||
using System.Threading; |
|||
using System.Threading.Tasks; |
|||
|
|||
namespace LINGYUN.Abp.AIManagement.Tools; |
|||
public interface IDynamicAIToolDefinitionStoreInMemoryCache |
|||
{ |
|||
string CacheStamp { get; set; } |
|||
|
|||
SemaphoreSlim SyncSemaphore { get; } |
|||
|
|||
DateTime? LastCheckTime { get; set; } |
|||
|
|||
Task FillAsync(List<AIToolDefinitionRecord> tools); |
|||
|
|||
AIToolDefinition? GetAIToolOrNull(string name); |
|||
|
|||
IReadOnlyList<AIToolDefinition> GetAITools(); |
|||
} |
|||
@ -0,0 +1,7 @@ |
|||
using System.Threading.Tasks; |
|||
|
|||
namespace LINGYUN.Abp.AIManagement.Tools; |
|||
public interface IStaticAIToolSaver |
|||
{ |
|||
Task SaveAsync(); |
|||
} |
|||
@ -0,0 +1,241 @@ |
|||
using LINGYUN.Abp.AI.Tools; |
|||
using Microsoft.Extensions.Caching.Distributed; |
|||
using Microsoft.Extensions.Options; |
|||
using System; |
|||
using System.Collections.Generic; |
|||
using System.Linq; |
|||
using System.Text; |
|||
using System.Text.Json; |
|||
using System.Text.Json.Serialization.Metadata; |
|||
using System.Threading.Tasks; |
|||
using Volo.Abp; |
|||
using Volo.Abp.Caching; |
|||
using Volo.Abp.DependencyInjection; |
|||
using Volo.Abp.DistributedLocking; |
|||
using Volo.Abp.Guids; |
|||
using Volo.Abp.Json.SystemTextJson.Modifiers; |
|||
using Volo.Abp.Threading; |
|||
using Volo.Abp.Uow; |
|||
|
|||
namespace LINGYUN.Abp.AIManagement.Tools; |
|||
public class StaticAIToolSaver : IStaticAIToolSaver, ITransientDependency |
|||
{ |
|||
protected IStaticAIToolDefinitionStore StaticStore { get; } |
|||
protected IAIToolDefinitionRecordRepository AIToolDefinitionRecordRepository { get; } |
|||
protected IAIToolDefinitionSerializer AIToolDefinitionSerializer { get; } |
|||
protected IDistributedCache Cache { get; } |
|||
protected IApplicationInfoAccessor ApplicationInfoAccessor { get; } |
|||
protected IAbpDistributedLock DistributedLock { get; } |
|||
protected AbpAIToolsOptions AIToolOptions { get; } |
|||
protected ICancellationTokenProvider CancellationTokenProvider { get; } |
|||
protected AbpDistributedCacheOptions CacheOptions { get; } |
|||
protected IUnitOfWorkManager UnitOfWorkManager { get; } |
|||
protected IGuidGenerator GuidGenerator { get; } |
|||
|
|||
public StaticAIToolSaver( |
|||
IStaticAIToolDefinitionStore staticStore, |
|||
IAIToolDefinitionRecordRepository aiToolDefinitionRecordRepository, |
|||
IAIToolDefinitionSerializer aiToolDefinitionSerializer, |
|||
IDistributedCache cache, |
|||
IOptions<AbpDistributedCacheOptions> cacheOptions, |
|||
IApplicationInfoAccessor applicationInfoAccessor, |
|||
IAbpDistributedLock distributedLock, |
|||
IOptions<AbpAIToolsOptions> aiToolOptions, |
|||
ICancellationTokenProvider cancellationTokenProvider, |
|||
IUnitOfWorkManager unitOfWorkManager, |
|||
IGuidGenerator guidGenerator) |
|||
{ |
|||
StaticStore = staticStore; |
|||
AIToolDefinitionRecordRepository = aiToolDefinitionRecordRepository; |
|||
AIToolDefinitionSerializer = aiToolDefinitionSerializer; |
|||
Cache = cache; |
|||
ApplicationInfoAccessor = applicationInfoAccessor; |
|||
DistributedLock = distributedLock; |
|||
CancellationTokenProvider = cancellationTokenProvider; |
|||
AIToolOptions = aiToolOptions.Value; |
|||
CacheOptions = cacheOptions.Value; |
|||
UnitOfWorkManager = unitOfWorkManager; |
|||
GuidGenerator = guidGenerator; |
|||
} |
|||
|
|||
[UnitOfWork] |
|||
public async Task SaveAsync() |
|||
{ |
|||
await using var applicationLockHandle = await DistributedLock.TryAcquireAsync( |
|||
GetApplicationDistributedLockKey() |
|||
); |
|||
|
|||
if (applicationLockHandle == null) |
|||
{ |
|||
return; |
|||
} |
|||
|
|||
var cacheKey = GetApplicationHashCacheKey(); |
|||
var cachedHash = await Cache.GetStringAsync(cacheKey, CancellationTokenProvider.Token); |
|||
|
|||
var aiTools = await AIToolDefinitionSerializer.SerializeAsync(await StaticStore.GetAllAsync()); |
|||
var currentHash = CalculateHash(aiTools, AIToolOptions.DeletedAITools); |
|||
|
|||
if (cachedHash == currentHash) |
|||
{ |
|||
return; |
|||
} |
|||
|
|||
await using (var commonLockHandle = await DistributedLock.TryAcquireAsync( |
|||
GetCommonDistributedLockKey(), |
|||
TimeSpan.FromMinutes(5))) |
|||
{ |
|||
if (commonLockHandle == null) |
|||
{ |
|||
/* It will re-try */ |
|||
throw new AbpException("Could not acquire distributed lock for saving static AITool!"); |
|||
} |
|||
|
|||
using (var unitOfWork = UnitOfWorkManager.Begin(requiresNew: true, isTransactional: true)) |
|||
{ |
|||
try |
|||
{ |
|||
var hasChangesInAITools = await UpdateChangedAIToolsAsync(aiTools); |
|||
|
|||
if (hasChangesInAITools) |
|||
{ |
|||
await Cache.SetStringAsync( |
|||
GetCommonStampCacheKey(), |
|||
Guid.NewGuid().ToString(), |
|||
new DistributedCacheEntryOptions |
|||
{ |
|||
SlidingExpiration = TimeSpan.FromDays(30) |
|||
}, |
|||
CancellationTokenProvider.Token |
|||
); |
|||
} |
|||
} |
|||
catch |
|||
{ |
|||
try |
|||
{ |
|||
await unitOfWork.RollbackAsync(); |
|||
} |
|||
catch |
|||
{ |
|||
/* ignored */ |
|||
} |
|||
|
|||
throw; |
|||
} |
|||
|
|||
await unitOfWork.CompleteAsync(); |
|||
} |
|||
} |
|||
|
|||
await Cache.SetStringAsync( |
|||
cacheKey, |
|||
currentHash, |
|||
new DistributedCacheEntryOptions |
|||
{ |
|||
SlidingExpiration = TimeSpan.FromDays(30) |
|||
}, |
|||
CancellationTokenProvider.Token |
|||
); |
|||
} |
|||
|
|||
private async Task<bool> UpdateChangedAIToolsAsync(AIToolDefinitionRecord[] aiToolRecords) |
|||
{ |
|||
var newRecords = new List<AIToolDefinitionRecord>(); |
|||
var changedRecords = new List<AIToolDefinitionRecord>(); |
|||
|
|||
var aiToolRecordsInDatabase = (await AIToolDefinitionRecordRepository.GetListAsync()).ToDictionary(x => x.Name); |
|||
|
|||
foreach (var record in aiToolRecords) |
|||
{ |
|||
var aiToolRecordInDatabase = aiToolRecordsInDatabase.GetOrDefault(record.Name); |
|||
if (aiToolRecordInDatabase == null) |
|||
{ |
|||
/* New group */ |
|||
newRecords.Add(record); |
|||
continue; |
|||
} |
|||
|
|||
if (record.HasSameData(aiToolRecordInDatabase)) |
|||
{ |
|||
/* Not changed */ |
|||
continue; |
|||
} |
|||
|
|||
/* Changed */ |
|||
aiToolRecordInDatabase.Patch(record); |
|||
changedRecords.Add(aiToolRecordInDatabase); |
|||
} |
|||
|
|||
/* Deleted */ |
|||
var deletedRecords = new List<AIToolDefinitionRecord>(); |
|||
|
|||
if (AIToolOptions.DeletedAITools.Any()) |
|||
{ |
|||
deletedRecords.AddRange(aiToolRecordsInDatabase.Values.Where(x => AIToolOptions.DeletedAITools.Contains(x.Name))); |
|||
} |
|||
|
|||
if (newRecords.Any()) |
|||
{ |
|||
await AIToolDefinitionRecordRepository.InsertManyAsync(newRecords); |
|||
} |
|||
|
|||
if (changedRecords.Any()) |
|||
{ |
|||
await AIToolDefinitionRecordRepository.UpdateManyAsync(changedRecords); |
|||
} |
|||
|
|||
if (deletedRecords.Any()) |
|||
{ |
|||
await AIToolDefinitionRecordRepository.DeleteManyAsync(deletedRecords); |
|||
} |
|||
|
|||
return newRecords.Any() || changedRecords.Any() || deletedRecords.Any(); |
|||
} |
|||
|
|||
private string GetApplicationDistributedLockKey() |
|||
{ |
|||
return $"{CacheOptions.KeyPrefix}_{ApplicationInfoAccessor.ApplicationName}_AbpAIToolUpdateLock"; |
|||
} |
|||
|
|||
private string GetCommonDistributedLockKey() |
|||
{ |
|||
return $"{CacheOptions.KeyPrefix}_Common_AbpAIToolUpdateLock"; |
|||
} |
|||
|
|||
private string GetApplicationHashCacheKey() |
|||
{ |
|||
return $"{CacheOptions.KeyPrefix}_{ApplicationInfoAccessor.ApplicationName}_AbpAIToolsHash"; |
|||
} |
|||
|
|||
private string GetCommonStampCacheKey() |
|||
{ |
|||
return $"{CacheOptions.KeyPrefix}_AbpInMemoryAIToolCacheStamp"; |
|||
} |
|||
|
|||
private string CalculateHash(AIToolDefinitionRecord[] aiToolRecords, IEnumerable<string> deletedAITool) |
|||
{ |
|||
var jsonSerializerOptions = new JsonSerializerOptions |
|||
{ |
|||
TypeInfoResolver = new DefaultJsonTypeInfoResolver |
|||
{ |
|||
Modifiers = |
|||
{ |
|||
new AbpIgnorePropertiesModifiers<AIToolDefinitionRecord, Guid>().CreateModifyAction(x => x.Id), |
|||
} |
|||
} |
|||
}; |
|||
|
|||
var stringBuilder = new StringBuilder(); |
|||
|
|||
stringBuilder.Append("AITools:"); |
|||
stringBuilder.AppendLine(JsonSerializer.Serialize(aiToolRecords, jsonSerializerOptions)); |
|||
|
|||
stringBuilder.Append("DeletedAITool:"); |
|||
stringBuilder.Append(deletedAITool.JoinAsString(",")); |
|||
|
|||
return stringBuilder |
|||
.ToString() |
|||
.ToMd5(); |
|||
} |
|||
} |
|||
Loading…
Reference in new issue