Browse Source

Fix event consumer on documentdb

pull/1272/head
Sebastian Stehle 8 months ago
parent
commit
1fb76e4afa
  1. 5
      backend/global.json
  2. 14
      backend/src/Squidex.Data.EntityFramework/Squidex.Data.EntityFramework.csproj
  3. 12
      backend/src/Squidex.Data.MongoDb/Squidex.Data.MongoDb.csproj
  4. 2
      backend/src/Squidex.Domain.Apps.Core.Model/Squidex.Domain.Apps.Core.Model.csproj
  5. 4
      backend/src/Squidex.Domain.Apps.Core.Operations/Squidex.Domain.Apps.Core.Operations.csproj
  6. 29
      backend/src/Squidex.Infrastructure/EventSourcing/Consume/EventConsumerProcessor.cs
  7. 14
      backend/src/Squidex.Infrastructure/Squidex.Infrastructure.csproj
  8. 24
      backend/src/Squidex/Squidex.csproj
  9. 4
      frontend/src/app/shared/state/asset-scripts.state.ts

5
backend/global.json

@ -0,0 +1,5 @@
{
"sdk": {
"version": "8.0.416"
}
}

14
backend/src/Squidex.Data.EntityFramework/Squidex.Data.EntityFramework.csproj

@ -40,13 +40,13 @@
<PackageReference Include="Pomelo.EntityFrameworkCore.MySql.Json.Microsoft" Version="8.0.3" /> <PackageReference Include="Pomelo.EntityFrameworkCore.MySql.Json.Microsoft" Version="8.0.3" />
<PackageReference Include="Pomelo.EntityFrameworkCore.MySql.NetTopologySuite" Version="8.0.3" /> <PackageReference Include="Pomelo.EntityFrameworkCore.MySql.NetTopologySuite" Version="8.0.3" />
<PackageReference Include="RefactoringEssentials" Version="5.6.0" PrivateAssets="all" /> <PackageReference Include="RefactoringEssentials" Version="5.6.0" PrivateAssets="all" />
<PackageReference Include="Squidex.AI.EntityFramework" Version="7.33.0" /> <PackageReference Include="Squidex.AI.EntityFramework" Version="7.34.0" />
<PackageReference Include="Squidex.Assets.EntityFramework" Version="7.33.0" /> <PackageReference Include="Squidex.Assets.EntityFramework" Version="7.34.0" />
<PackageReference Include="Squidex.Assets.TusAdapter" Version="7.33.0" /> <PackageReference Include="Squidex.Assets.TusAdapter" Version="7.34.0" />
<PackageReference Include="Squidex.Events.EntityFramework" Version="7.33.0" /> <PackageReference Include="Squidex.Events.EntityFramework" Version="7.34.0" />
<PackageReference Include="Squidex.Flows.EntityFramework" Version="7.33.0" /> <PackageReference Include="Squidex.Flows.EntityFramework" Version="7.34.0" />
<PackageReference Include="Squidex.Hosting" Version="7.33.0" /> <PackageReference Include="Squidex.Hosting" Version="7.34.0" />
<PackageReference Include="Squidex.Messaging.EntityFramework" Version="7.33.0" /> <PackageReference Include="Squidex.Messaging.EntityFramework" Version="7.34.0" />
<PackageReference Include="Squidex.OpenIdDict.EntityFramework" Version="5.8.4" /> <PackageReference Include="Squidex.OpenIdDict.EntityFramework" Version="5.8.4" />
<PackageReference Include="StyleCop.Analyzers" Version="1.1.118" PrivateAssets="all" /> <PackageReference Include="StyleCop.Analyzers" Version="1.1.118" PrivateAssets="all" />
<PackageReference Include="System.ValueTuple" Version="4.5.0" /> <PackageReference Include="System.ValueTuple" Version="4.5.0" />

12
backend/src/Squidex.Data.MongoDb/Squidex.Data.MongoDb.csproj

@ -25,12 +25,12 @@
<PackageReference Include="MongoDB.Driver.GridFS" Version="2.30.0" /> <PackageReference Include="MongoDB.Driver.GridFS" Version="2.30.0" />
<PackageReference Include="NodaTime.Serialization.SystemTextJson" Version="1.3.0" /> <PackageReference Include="NodaTime.Serialization.SystemTextJson" Version="1.3.0" />
<PackageReference Include="RefactoringEssentials" Version="5.6.0" PrivateAssets="all" /> <PackageReference Include="RefactoringEssentials" Version="5.6.0" PrivateAssets="all" />
<PackageReference Include="Squidex.AI.Mongo" Version="7.33.0" /> <PackageReference Include="Squidex.AI.Mongo" Version="7.34.0" />
<PackageReference Include="Squidex.Assets.Mongo" Version="7.33.0" /> <PackageReference Include="Squidex.Assets.Mongo" Version="7.34.0" />
<PackageReference Include="Squidex.Events.Mongo" Version="7.33.0" /> <PackageReference Include="Squidex.Events.Mongo" Version="7.34.0" />
<PackageReference Include="Squidex.Flows.Mongo" Version="7.33.0" /> <PackageReference Include="Squidex.Flows.Mongo" Version="7.34.0" />
<PackageReference Include="Squidex.Hosting" Version="7.33.0" /> <PackageReference Include="Squidex.Hosting" Version="7.34.0" />
<PackageReference Include="Squidex.Messaging.Mongo" Version="7.33.0" /> <PackageReference Include="Squidex.Messaging.Mongo" Version="7.34.0" />
<PackageReference Include="Squidex.OpenIddict.MongoDb" Version="5.8.5" /> <PackageReference Include="Squidex.OpenIddict.MongoDb" Version="5.8.5" />
<PackageReference Include="StyleCop.Analyzers" Version="1.1.118" PrivateAssets="all" /> <PackageReference Include="StyleCop.Analyzers" Version="1.1.118" PrivateAssets="all" />
<PackageReference Include="System.ValueTuple" Version="4.5.0" /> <PackageReference Include="System.ValueTuple" Version="4.5.0" />

2
backend/src/Squidex.Domain.Apps.Core.Model/Squidex.Domain.Apps.Core.Model.csproj

@ -20,7 +20,7 @@
<PackageReference Include="NetTopologySuite" Version="2.5.0" /> <PackageReference Include="NetTopologySuite" Version="2.5.0" />
<PackageReference Include="NodaTime.Serialization.SystemTextJson" Version="1.3.0" /> <PackageReference Include="NodaTime.Serialization.SystemTextJson" Version="1.3.0" />
<PackageReference Include="RefactoringEssentials" Version="5.6.0" PrivateAssets="all" /> <PackageReference Include="RefactoringEssentials" Version="5.6.0" PrivateAssets="all" />
<PackageReference Include="Squidex.Flows" Version="7.33.0" /> <PackageReference Include="Squidex.Flows" Version="7.34.0" />
<PackageReference Include="StyleCop.Analyzers" Version="1.1.118" PrivateAssets="all" /> <PackageReference Include="StyleCop.Analyzers" Version="1.1.118" PrivateAssets="all" />
<PackageReference Include="System.Collections.Immutable" Version="8.0.0" /> <PackageReference Include="System.Collections.Immutable" Version="8.0.0" />
<PackageReference Include="System.ComponentModel.Annotations" Version="5.0.0" /> <PackageReference Include="System.ComponentModel.Annotations" Version="5.0.0" />

4
backend/src/Squidex.Domain.Apps.Core.Operations/Squidex.Domain.Apps.Core.Operations.csproj

@ -29,8 +29,8 @@
<PackageReference Include="NJsonSchema" Version="11.0.2" /> <PackageReference Include="NJsonSchema" Version="11.0.2" />
<PackageReference Include="NodaTime.Serialization.SystemTextJson" Version="1.3.0" /> <PackageReference Include="NodaTime.Serialization.SystemTextJson" Version="1.3.0" />
<PackageReference Include="RefactoringEssentials" Version="5.6.0" PrivateAssets="all" /> <PackageReference Include="RefactoringEssentials" Version="5.6.0" PrivateAssets="all" />
<PackageReference Include="Squidex.AI" Version="7.33.0" /> <PackageReference Include="Squidex.AI" Version="7.34.0" />
<PackageReference Include="Squidex.Messaging.Subscriptions" Version="7.33.0" /> <PackageReference Include="Squidex.Messaging.Subscriptions" Version="7.34.0" />
<PackageReference Include="StyleCop.Analyzers" Version="1.1.118" PrivateAssets="all" /> <PackageReference Include="StyleCop.Analyzers" Version="1.1.118" PrivateAssets="all" />
<PackageReference Include="System.Collections.Immutable" Version="8.0.0" /> <PackageReference Include="System.Collections.Immutable" Version="8.0.0" />
<PackageReference Include="System.Linq.Async" Version="6.0.1" /> <PackageReference Include="System.Linq.Async" Version="6.0.1" />

29
backend/src/Squidex.Infrastructure/EventSourcing/Consume/EventConsumerProcessor.cs

@ -21,6 +21,7 @@ public class EventConsumerProcessor : IEventSubscriber<ParsedEvents>
private readonly IEventStore eventStore; private readonly IEventStore eventStore;
private readonly ILogger<EventConsumerProcessor> log; private readonly ILogger<EventConsumerProcessor> log;
private readonly AsyncLock asyncLock = new AsyncLock(); private readonly AsyncLock asyncLock = new AsyncLock();
private readonly RetryWindow logWindow = new RetryWindow(TimeSpan.FromSeconds(10), 0);
private IEventSubscription? currentSubscription; private IEventSubscription? currentSubscription;
public EventConsumerState State public EventConsumerState State
@ -73,9 +74,9 @@ public class EventConsumerProcessor : IEventSubscriber<ParsedEvents>
} }
} }
public virtual async ValueTask OnNextAsync(IEventSubscription subscription, ParsedEvents @event) public virtual ValueTask OnNextAsync(IEventSubscription subscription, ParsedEvents @event)
{ {
await UpdateAsync(async () => return UpdateAsync(async () =>
{ {
if (!ReferenceEquals(subscription, currentSubscription)) if (!ReferenceEquals(subscription, currentSubscription))
{ {
@ -83,14 +84,13 @@ public class EventConsumerProcessor : IEventSubscriber<ParsedEvents>
} }
await DispatchAsync(@event.Events); await DispatchAsync(@event.Events);
State = State.Handled(@event.Position, @event.Events.Count); State = State.Handled(@event.Position, @event.Events.Count);
}, State.Position); }, State.Position);
} }
public virtual async ValueTask OnErrorAsync(IEventSubscription subscription, Exception exception) public virtual ValueTask OnErrorAsync(IEventSubscription subscription, Exception exception)
{ {
await UpdateAsync(() => return UpdateAsync(() =>
{ {
if (!ReferenceEquals(subscription, currentSubscription)) if (!ReferenceEquals(subscription, currentSubscription))
{ {
@ -98,12 +98,16 @@ public class EventConsumerProcessor : IEventSubscriber<ParsedEvents>
} }
Unsubscribe(); Unsubscribe();
State = State.Stopped(exception); State = State.Stopped(exception);
if (logWindow.CanRetryAfterFailure())
{
log.LogError(exception, "Failed to handle event.");
}
}, State.Position); }, State.Position);
} }
public virtual Task ActivateAsync() public virtual ValueTask ActivateAsync()
{ {
return UpdateAsync(() => return UpdateAsync(() =>
{ {
@ -120,7 +124,7 @@ public class EventConsumerProcessor : IEventSubscriber<ParsedEvents>
}, State.Position); }, State.Position);
} }
public virtual Task StartAsync() public virtual ValueTask StartAsync()
{ {
return UpdateAsync(() => return UpdateAsync(() =>
{ {
@ -130,12 +134,11 @@ public class EventConsumerProcessor : IEventSubscriber<ParsedEvents>
} }
Subscribe(); Subscribe();
State = State.Started(); State = State.Started();
}, State.Position); }, State.Position);
} }
public virtual Task StopAsync() public virtual ValueTask StopAsync()
{ {
return UpdateAsync(() => return UpdateAsync(() =>
{ {
@ -145,7 +148,6 @@ public class EventConsumerProcessor : IEventSubscriber<ParsedEvents>
} }
Unsubscribe(); Unsubscribe();
State = State.Stopped(); State = State.Stopped();
}, State.Position); }, State.Position);
} }
@ -164,7 +166,6 @@ public class EventConsumerProcessor : IEventSubscriber<ParsedEvents>
await ClearAsync(); await ClearAsync();
State = EventConsumerState.Initial; State = EventConsumerState.Initial;
Subscribe(); Subscribe();
}, State.Position); }, State.Position);
} }
@ -177,7 +178,7 @@ public class EventConsumerProcessor : IEventSubscriber<ParsedEvents>
} }
} }
private Task UpdateAsync(Action action, string? position, [CallerMemberName] string? caller = null) private ValueTask UpdateAsync(Action action, string? position, [CallerMemberName] string? caller = null)
{ {
return UpdateAsync(() => return UpdateAsync(() =>
{ {
@ -187,7 +188,7 @@ public class EventConsumerProcessor : IEventSubscriber<ParsedEvents>
}, position, caller); }, position, caller);
} }
private async Task UpdateAsync(Func<Task> action, string? position, [CallerMemberName] string? caller = null) private async ValueTask UpdateAsync(Func<Task> action, string? position, [CallerMemberName] string? caller = null)
{ {
// We do not want to deal with concurrency in this class, therefore we just use a lock. // We do not want to deal with concurrency in this class, therefore we just use a lock.
using (await asyncLock.EnterAsync()) using (await asyncLock.EnterAsync())

14
backend/src/Squidex.Infrastructure/Squidex.Infrastructure.csproj

@ -24,13 +24,13 @@
<PackageReference Include="NodaTime" Version="3.2.0" /> <PackageReference Include="NodaTime" Version="3.2.0" />
<PackageReference Include="OpenTelemetry.Api" Version="1.9.0" /> <PackageReference Include="OpenTelemetry.Api" Version="1.9.0" />
<PackageReference Include="RefactoringEssentials" Version="5.6.0" PrivateAssets="all" /> <PackageReference Include="RefactoringEssentials" Version="5.6.0" PrivateAssets="all" />
<PackageReference Include="Squidex.Assets" Version="7.33.0" /> <PackageReference Include="Squidex.Assets" Version="7.34.0" />
<PackageReference Include="Squidex.Caching" Version="7.33.0" /> <PackageReference Include="Squidex.Caching" Version="7.34.0" />
<PackageReference Include="Squidex.Events" Version="7.33.0" /> <PackageReference Include="Squidex.Events" Version="7.34.0" />
<PackageReference Include="Squidex.Hosting.Abstractions" Version="7.33.0" /> <PackageReference Include="Squidex.Hosting.Abstractions" Version="7.34.0" />
<PackageReference Include="Squidex.Log" Version="7.33.0" /> <PackageReference Include="Squidex.Log" Version="7.34.0" />
<PackageReference Include="Squidex.Messaging" Version="7.33.0" /> <PackageReference Include="Squidex.Messaging" Version="7.34.0" />
<PackageReference Include="Squidex.Text" Version="7.33.0" /> <PackageReference Include="Squidex.Text" Version="7.34.0" />
<PackageReference Include="StyleCop.Analyzers" Version="1.1.118" PrivateAssets="all" /> <PackageReference Include="StyleCop.Analyzers" Version="1.1.118" PrivateAssets="all" />
<PackageReference Include="System.Collections.Immutable" Version="8.0.0" /> <PackageReference Include="System.Collections.Immutable" Version="8.0.0" />
<PackageReference Include="System.ComponentModel.Annotations" Version="5.0.0" /> <PackageReference Include="System.ComponentModel.Annotations" Version="5.0.0" />

24
backend/src/Squidex/Squidex.csproj

@ -60,17 +60,17 @@
<PackageReference Include="OpenTelemetry.Instrumentation.Runtime" Version="1.9.0" /> <PackageReference Include="OpenTelemetry.Instrumentation.Runtime" Version="1.9.0" />
<PackageReference Include="RefactoringEssentials" Version="5.6.0" PrivateAssets="all" /> <PackageReference Include="RefactoringEssentials" Version="5.6.0" PrivateAssets="all" />
<PackageReference Include="ReportGenerator" Version="5.4.1" PrivateAssets="all" /> <PackageReference Include="ReportGenerator" Version="5.4.1" PrivateAssets="all" />
<PackageReference Include="Squidex.Assets.Azure" Version="7.33.0" /> <PackageReference Include="Squidex.Assets.Azure" Version="7.34.0" />
<PackageReference Include="Squidex.Assets.GoogleCloud" Version="7.33.0" /> <PackageReference Include="Squidex.Assets.GoogleCloud" Version="7.34.0" />
<PackageReference Include="Squidex.Assets.FTP" Version="7.33.0" /> <PackageReference Include="Squidex.Assets.FTP" Version="7.34.0" />
<PackageReference Include="Squidex.Assets.ImageSharp" Version="7.33.0" /> <PackageReference Include="Squidex.Assets.ImageSharp" Version="7.34.0" />
<PackageReference Include="Squidex.Assets.S3" Version="7.33.0" /> <PackageReference Include="Squidex.Assets.S3" Version="7.34.0" />
<PackageReference Include="Squidex.Assets.TusAdapter" Version="7.33.0" /> <PackageReference Include="Squidex.Assets.TusAdapter" Version="7.34.0" />
<PackageReference Include="Squidex.ClientLibrary" Version="21.8.0" /> <PackageReference Include="Squidex.ClientLibrary" Version="21.8.0" />
<PackageReference Include="Squidex.Events.GetEventStore" Version="7.33.0" /> <PackageReference Include="Squidex.Events.GetEventStore" Version="7.34.0" />
<PackageReference Include="Squidex.Hosting" Version="7.33.0" /> <PackageReference Include="Squidex.Hosting" Version="7.34.0" />
<PackageReference Include="Squidex.Messaging.All" Version="7.33.0" /> <PackageReference Include="Squidex.Messaging.All" Version="7.34.0" />
<PackageReference Include="Squidex.Messaging.Subscriptions" Version="7.33.0" /> <PackageReference Include="Squidex.Messaging.Subscriptions" Version="7.34.0" />
<PackageReference Include="StyleCop.Analyzers" Version="1.1.118" PrivateAssets="all" /> <PackageReference Include="StyleCop.Analyzers" Version="1.1.118" PrivateAssets="all" />
<PackageReference Include="YDotNet" Version="0.4.3" /> <PackageReference Include="YDotNet" Version="0.4.3" />
<PackageReference Include="YDotNet.Native" Version="0.4.3" /> <PackageReference Include="YDotNet.Native" Version="0.4.3" />
@ -84,11 +84,11 @@
</ItemGroup> </ItemGroup>
<ItemGroup Condition="'$(IncludeMagick)' == 'true'"> <ItemGroup Condition="'$(IncludeMagick)' == 'true'">
<PackageReference Include="Squidex.Assets.ImageMagick" Version="7.33.0" /> <PackageReference Include="Squidex.Assets.ImageMagick" Version="7.34.0" />
</ItemGroup> </ItemGroup>
<ItemGroup Condition="'$(IncludeKafka)' == 'true'"> <ItemGroup Condition="'$(IncludeKafka)' == 'true'">
<PackageReference Include="Squidex.Messaging.Kafka" Version="7.33.0" /> <PackageReference Include="Squidex.Messaging.Kafka" Version="7.34.0" />
</ItemGroup> </ItemGroup>
<PropertyGroup> <PropertyGroup>

4
frontend/src/app/shared/state/asset-scripts.state.ts

@ -99,10 +99,10 @@ export class AssetScriptsState extends State<Snapshot> {
} }
private replaceAssetScripts(payload: AssetScriptsDto, version: VersionTag) { private replaceAssetScripts(payload: AssetScriptsDto, version: VersionTag) {
const { canUpdate, _links: _, version: __, ...scripts } = payload.toJSON(); const { _links: _, version: __, ...scripts } = payload.toJSON();
this.next({ this.next({
canUpdate, canUpdate: payload.canUpdate,
scripts, scripts,
isLoaded: true, isLoaded: true,
isLoading: false, isLoading: false,

Loading…
Cancel
Save