From dad62cc16d4d6de671563f6b1e54196f617926a5 Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Wed, 27 Jun 2018 03:57:34 +1000 Subject: [PATCH 1/6] Fix client connection logic (#4711) --- Samples/2.0/docker-aspnet-core/API/Startup.cs | 22 ++++++++----------- 1 file changed, 9 insertions(+), 13 deletions(-) diff --git a/Samples/2.0/docker-aspnet-core/API/Startup.cs b/Samples/2.0/docker-aspnet-core/API/Startup.cs index 3c79b057327..d7146d9bf49 100644 --- a/Samples/2.0/docker-aspnet-core/API/Startup.cs +++ b/Samples/2.0/docker-aspnet-core/API/Startup.cs @@ -5,6 +5,7 @@ using Microsoft.AspNetCore.Hosting; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; using Orleans; using Orleans.Configuration; using Orleans.Hosting; @@ -42,6 +43,8 @@ public void Configure(IApplicationBuilder app, Microsoft.AspNetCore.Hosting.IHos private IClusterClient CreateClusterClient(IServiceProvider serviceProvider) { + var log = serviceProvider.GetService>(); + // TODO replace with your connection string const string connectionString = "YOUR_CONNECTION_STRING_HERE"; var client = new ClientBuilder() @@ -53,22 +56,15 @@ private IClusterClient CreateClusterClient(IServiceProvider serviceProvider) }) .UseAzureStorageClustering(options => options.ConnectionString = connectionString) .Build(); - StartClientWithRetries(client).Wait(); + + client.Connect(RetryFilter).GetAwaiter().GetResult(); return client; - } - private static async Task StartClientWithRetries(IClusterClient client) - { - for (var i=0; i<5; i++) + async Task RetryFilter(Exception exception) { - try - { - await client.Connect(); - return; - } - catch(Exception) - { } - await Task.Delay(TimeSpan.FromSeconds(5)); + log?.LogWarning("Exception while attempting to connect to Orleans cluster: {Exception}", exception); + await Task.Delay(TimeSpan.FromSeconds(2)); + return true; } } } From 6b9d58cbdce3107731d36113e98ffa971732f39f Mon Sep 17 00:00:00 2001 From: Xiao Zeng Date: Tue, 26 Jun 2018 12:44:27 -0700 Subject: [PATCH 2/6] Make LifecycleSubject logging less verbose (#4660) * make LifecycleSubject logging less verbose * pr feedback * fix compile error * fix tets --- .../Lifecycle/LifecycleSubject.cs | 22 +++++++++++++---- .../Lifecycle/SiloLifecycleSubject.cs | 24 +++++++++++-------- .../StorageTests/Relational/CommonFixture.cs | 2 +- 3 files changed, 32 insertions(+), 16 deletions(-) diff --git a/src/Orleans.Core/Lifecycle/LifecycleSubject.cs b/src/Orleans.Core/Lifecycle/LifecycleSubject.cs index 9139b8d86b0..9d58625a6b4 100644 --- a/src/Orleans.Core/Lifecycle/LifecycleSubject.cs +++ b/src/Orleans.Core/Lifecycle/LifecycleSubject.cs @@ -30,7 +30,13 @@ public LifecycleSubject(ILogger logger) this.subscribers = new ConcurrentDictionary(); } - public async Task OnStart(CancellationToken ct) + protected virtual void PerfMeasureOnStart(int? stage, TimeSpan timelapsed) + { + if (this.logger != null && this.logger.IsEnabled(LogLevel.Trace)) + this.logger.Trace(ErrorCode.SiloStartPerfMeasure, $"Starting lifecycle stage {stage} took {timelapsed.TotalMilliseconds} Milliseconds"); + } + + public virtual async Task OnStart(CancellationToken ct) { if (this.highStage.HasValue) throw new InvalidOperationException("Lifecycle has already been started."); try @@ -47,7 +53,7 @@ public async Task OnStart(CancellationToken ct) Stopwatch stopWatch = Stopwatch.StartNew(); await Task.WhenAll(observerGroup.Select(orderedObserver => WrapExecution(ct, orderedObserver.Observer.OnStart))); stopWatch.Stop(); - this.logger?.Info(ErrorCode.SiloStartPerfMeasure, $"Starting lifecycle stage {this.highStage} took {stopWatch.ElapsedMilliseconds} Milliseconds"); + this.PerfMeasureOnStart(this.highStage, stopWatch.Elapsed); } } catch (Exception ex) @@ -58,7 +64,13 @@ public async Task OnStart(CancellationToken ct) } } - public async Task OnStop(CancellationToken ct) + protected virtual void PerfMeasureOnStop(int? stage, TimeSpan timelapsed) + { + if (this.logger != null && this.logger.IsEnabled(LogLevel.Trace)) + this.logger.Trace(ErrorCode.SiloStartPerfMeasure, $"Stopping lifecycle stage {stage} took {timelapsed.TotalMilliseconds} Milliseconds"); + } + + public virtual async Task OnStop(CancellationToken ct) { // if not started, do nothing if (!this.highStage.HasValue) return; @@ -74,7 +86,7 @@ public async Task OnStop(CancellationToken ct) Stopwatch stopWatch = Stopwatch.StartNew(); await Task.WhenAll(observerGroup.Select(orderedObserver => WrapExecution(ct, orderedObserver.Observer.OnStop))); stopWatch.Stop(); - this.logger?.Info(ErrorCode.SiloStartPerfMeasure, $"Stopping lifecycle stage {this.highStage} took {stopWatch.ElapsedMilliseconds} Milliseconds"); + this.PerfMeasureOnStop(this.highStage, stopWatch.Elapsed); } catch (Exception ex) { @@ -83,7 +95,7 @@ public async Task OnStop(CancellationToken ct) } } - public IDisposable Subscribe(string observerName, int stage, ILifecycleObserver observer) + public virtual IDisposable Subscribe(string observerName, int stage, ILifecycleObserver observer) { if (observer == null) throw new ArgumentNullException(nameof(observer)); if (this.highStage.HasValue) throw new InvalidOperationException("Lifecycle has already been started."); diff --git a/src/Orleans.Runtime/Lifecycle/SiloLifecycleSubject.cs b/src/Orleans.Runtime/Lifecycle/SiloLifecycleSubject.cs index fa0ac9d3417..ebb377cff82 100644 --- a/src/Orleans.Runtime/Lifecycle/SiloLifecycleSubject.cs +++ b/src/Orleans.Runtime/Lifecycle/SiloLifecycleSubject.cs @@ -12,38 +12,42 @@ namespace Orleans.Runtime /// /// Decorator over lifecycle subject for silo. Adds some logging and monitoring /// - public class SiloLifecycleSubject : ISiloLifecycleSubject + public class SiloLifecycleSubject : LifecycleSubject, ISiloLifecycleSubject { - private readonly ILifecycleSubject subject; private readonly ILogger logger; private readonly List observers; - public SiloLifecycleSubject(ILifecycleSubject subject, ILogger logger) + public SiloLifecycleSubject(ILogger logger) + :base(logger) { - this.subject = subject; this.logger = logger; this.observers = new List(); } - public Task OnStart(CancellationToken ct) + public override Task OnStart(CancellationToken ct) { foreach(IGrouping stage in this.observers.GroupBy(o => o.Stage).OrderBy(s => s.Key)) { this.logger?.Info(ErrorCode.LifecycleStagesReport, $"Stage {stage.Key}: {string.Join(", ", stage.Select(o => o.Name))}", stage.Key); } - return this.subject.OnStart(ct); + return base.OnStart(ct); } - public Task OnStop(CancellationToken ct) + protected override void PerfMeasureOnStop(int? stage, TimeSpan timelapsed) { - return this.subject.OnStop(ct); + this.logger?.Info(ErrorCode.SiloStartPerfMeasure, $"Stopping lifecycle stage {stage} took {timelapsed.TotalMilliseconds} Milliseconds"); } - public IDisposable Subscribe(string observerName, int stage, ILifecycleObserver observer) + protected override void PerfMeasureOnStart(int? stage, TimeSpan timelapsed) + { + this.logger?.Info(ErrorCode.SiloStartPerfMeasure, $"Starting lifecycle stage {stage} took {timelapsed.TotalMilliseconds} Milliseconds"); + } + + public override IDisposable Subscribe(string observerName, int stage, ILifecycleObserver observer) { var monitoredObserver = new MonitoredObserver(observerName, stage, observer, this.logger); this.observers.Add(monitoredObserver); - return this.subject.Subscribe(observerName, stage, monitoredObserver); + return base.Subscribe(observerName, stage, monitoredObserver); } private class MonitoredObserver : ILifecycleObserver diff --git a/test/Extensions/TesterAdoNet/StorageTests/Relational/CommonFixture.cs b/test/Extensions/TesterAdoNet/StorageTests/Relational/CommonFixture.cs index a1a4a99f0ee..9e4ac022e09 100644 --- a/test/Extensions/TesterAdoNet/StorageTests/Relational/CommonFixture.cs +++ b/test/Extensions/TesterAdoNet/StorageTests/Relational/CommonFixture.cs @@ -92,7 +92,7 @@ public async Task GetStorageProvider(string storageInvariant) ServiceId = Guid.NewGuid().ToString() }; var storageProvider = new AdoNetGrainStorage(DefaultProviderRuntime.ServiceProvider.GetService>(), DefaultProviderRuntime, Options.Create(options), Options.Create(clusterOptions), storageInvariant + "_StorageProvider"); - ISiloLifecycleSubject siloLifeCycle = new SiloLifecycleSubject(new LifecycleSubject(NullLoggerFactory.Instance.CreateLogger()), NullLoggerFactory.Instance.CreateLogger()); + ISiloLifecycleSubject siloLifeCycle = new SiloLifecycleSubject(NullLoggerFactory.Instance.CreateLogger()); storageProvider.Participate(siloLifeCycle); await siloLifeCycle.OnStart(CancellationToken.None); From 637830d8bc8eeeb5be84c2b808e4abd7029a7c0b Mon Sep 17 00:00:00 2001 From: Benjamin Petit Date: Tue, 26 Jun 2018 16:45:37 -0700 Subject: [PATCH 3/6] Do not use ip address from interface not operational" (#4713) --- src/Orleans.Core/Configuration/ConfigUtilities.cs | 1 + 1 file changed, 1 insertion(+) diff --git a/src/Orleans.Core/Configuration/ConfigUtilities.cs b/src/Orleans.Core/Configuration/ConfigUtilities.cs index 6e4f46b5359..8d5c211d023 100644 --- a/src/Orleans.Core/Configuration/ConfigUtilities.cs +++ b/src/Orleans.Core/Configuration/ConfigUtilities.cs @@ -349,6 +349,7 @@ internal static async Task ResolveIPAddress(string addrOrHost, byte[] if (string.IsNullOrEmpty(addrOrHost)) { nodeIps = NetworkInterface.GetAllNetworkInterfaces() + .Where(iface => iface.OperationalStatus == OperationalStatus.Up) .SelectMany(iface => iface.GetIPProperties().UnicastAddresses) .Select(addr => addr.Address) .Where(addr => addr.AddressFamily == family && !IPAddress.IsLoopback(addr)) From 7c2239e8ea50183c0c666ba9c62e7a43131614e5 Mon Sep 17 00:00:00 2001 From: xiazen Date: Fri, 22 Jun 2018 10:54:27 -0700 Subject: [PATCH 4/6] fix TransactionParticipantExtensionWrapper serialization issues --- ...reTableTransactionalStateStorageFactory.cs | 5 +++- ...ansactionParticipantExtensionExtensions.cs | 5 ++-- .../State/StorageBatch.cs | 30 +++++++++++-------- .../State/TransactionParticipant.cs | 2 +- .../State/TransactionParticipantExtension.cs | 3 ++ .../State/TransactionalState.cs | 14 ++++----- .../TestFixture.cs | 2 -- .../Memory/SerializationTests.cs | 8 +++-- 8 files changed, 38 insertions(+), 31 deletions(-) diff --git a/src/Azure/Orleans.Transactions.AzureStorage/TransactionalState/AzureTableTransactionalStateStorageFactory.cs b/src/Azure/Orleans.Transactions.AzureStorage/TransactionalState/AzureTableTransactionalStateStorageFactory.cs index a9587fe9283..a08e345f69b 100644 --- a/src/Azure/Orleans.Transactions.AzureStorage/TransactionalState/AzureTableTransactionalStateStorageFactory.cs +++ b/src/Azure/Orleans.Transactions.AzureStorage/TransactionalState/AzureTableTransactionalStateStorageFactory.cs @@ -10,6 +10,7 @@ using Orleans.Runtime; using Orleans.Serialization; using Orleans.Transactions.Abstractions; +using Orleans.Transactions.Abstractions.Extensions; namespace Orleans.Transactions.AzureStorage { @@ -33,7 +34,9 @@ public AzureTableTransactionalStateStorageFactory(string name, AzureTableTransac this.name = name; this.options = options; this.clusterOptions = clusterOptions.Value; - this.jsonSettings = OrleansJsonSerializer.GetDefaultSerializerSettings(typeResolver, grainFactory); + this.jsonSettings = TransactionParticipantExtensionExtensions.GetJsonSerializerSettings( + typeResolver, + grainFactory); this.loggerFactory = loggerFactory; } diff --git a/src/Orleans.Transactions/Abstractions/Extensions/TransactionParticipantExtensionExtensions.cs b/src/Orleans.Transactions/Abstractions/Extensions/TransactionParticipantExtensionExtensions.cs index b8d71730824..3f5bf9892c0 100644 --- a/src/Orleans.Transactions/Abstractions/Extensions/TransactionParticipantExtensionExtensions.cs +++ b/src/Orleans.Transactions/Abstractions/Extensions/TransactionParticipantExtensionExtensions.cs @@ -7,8 +7,6 @@ using Orleans.Runtime; using Orleans.Serialization; -[assembly: GenerateSerializer(typeof(Orleans.Transactions.Abstractions.Extensions.TransactionParticipantExtensionExtensions.TransactionParticipantExtensionWrapper))] - namespace Orleans.Transactions.Abstractions.Extensions { public static class TransactionParticipantExtensionExtensions @@ -27,7 +25,6 @@ public static ITransactionParticipant AsTransactionParticipant(this ITransaction public static JsonSerializerSettings GetJsonSerializerSettings(ITypeResolver typeResolver, IGrainFactory grainFactory) { var serializerSettings = OrleansJsonSerializer.GetDefaultSerializerSettings(typeResolver, grainFactory); - serializerSettings.TypeNameHandling = TypeNameHandling.Auto; serializerSettings.PreserveReferencesHandling = PreserveReferencesHandling.None; return serializerSettings; } @@ -36,7 +33,9 @@ public static JsonSerializerSettings GetJsonSerializerSettings(ITypeResolver typ [Immutable] internal sealed class TransactionParticipantExtensionWrapper : ITransactionParticipant { + [JsonProperty] private readonly ITransactionParticipantExtension extension; + [JsonProperty] private readonly string resourceId; public TransactionParticipantExtensionWrapper(ITransactionParticipantExtension transactionalExtension, string resourceId) diff --git a/src/Orleans.Transactions/State/StorageBatch.cs b/src/Orleans.Transactions/State/StorageBatch.cs index 2abc2502fa3..2feca30d0b7 100644 --- a/src/Orleans.Transactions/State/StorageBatch.cs +++ b/src/Orleans.Transactions/State/StorageBatch.cs @@ -1,10 +1,12 @@ using System; using System.Collections.Generic; using System.Linq; +using System.Runtime.CompilerServices; using System.Threading.Tasks; using Newtonsoft.Json; using Orleans.Concurrency; using Orleans.Transactions.Abstractions; +using Orleans.Transactions.Abstractions.Extensions; namespace Orleans.Transactions { @@ -35,8 +37,6 @@ public class MetaData public DateTime TimeStamp { get; set; } public Dictionary CommitRecords { get; set; } - - public static JsonSerializerSettings SerializerSettings { get; set; } } [Serializable] @@ -74,38 +74,43 @@ internal class StorageBatch : ITransactionalStateStorageEvents private int confirm = 0; private int collect = 0; private int cancel = 0; - + private JsonSerializerSettings serializerSettings; public MetaData MetaData { get; private set; } public string ETag { get; set; } public int BatchSize => total; - public override string ToString() { return $"batchsize={total} [{read}r {prepare}p {commit}c {confirm}cf {collect}cl {cancel}cc]"; } - public StorageBatch(TransactionalStorageLoadResponse loadresponse) + public StorageBatch(TransactionalStorageLoadResponse loadresponse, JsonSerializerSettings settings) { - MetaData = ReadMetaData(loadresponse); + if(settings == null) + throw new ArgumentNullException(nameof(settings)); + this.serializerSettings = settings; + MetaData = ReadMetaData(loadresponse, this.serializerSettings); ETag = loadresponse.ETag; confirmUpTo = loadresponse.CommittedSequenceId; cancelAbove = loadresponse.PendingStates.LastOrDefault()?.SequenceId ?? loadresponse.CommittedSequenceId; cancelAboveStart = cancelAbove; } - public StorageBatch(StorageBatch previous) + public StorageBatch(StorageBatch previous, JsonSerializerSettings settings) { + if (settings == null) + throw new ArgumentNullException(nameof(settings)); + this.serializerSettings = settings; MetaData = previous.MetaData; confirmUpTo = previous.confirmUpTo; cancelAbove = previous.cancelAbove; cancelAboveStart = cancelAbove; } - private static MetaData ReadMetaData(TransactionalStorageLoadResponse loadresponse) + private static MetaData ReadMetaData(TransactionalStorageLoadResponse loadresponse, JsonSerializerSettings serializerSettings) { - if (string.IsNullOrEmpty(loadresponse.Metadata)) + if (string.IsNullOrEmpty(loadresponse?.Metadata)) { // this thing is fresh... did not exist in storage yet return new MetaData() @@ -116,14 +121,13 @@ private static MetaData ReadMetaData(TransactionalStorageLoadResponse lo } else { - return JsonConvert.DeserializeObject(loadresponse.Metadata, MetaData.SerializerSettings); + return JsonConvert.DeserializeObject(loadresponse.Metadata, serializerSettings); } } public Task Store(ITransactionalStateStorage storage) { - var jsonMetaData = JsonConvert.SerializeObject(MetaData, MetaData.SerializerSettings); - + var jsonMetaData = JsonConvert.SerializeObject(MetaData, this.serializerSettings); var list = new List>(); if (prepares != null) @@ -176,7 +180,7 @@ public void Prepare(long sequenceNumber, Guid transactionId, DateTime timestamp, prepares = new SortedDictionary>(); var tmstring = (transactionManager == null) ? null : - JsonConvert.SerializeObject(transactionManager, MetaData.SerializerSettings); + JsonConvert.SerializeObject(transactionManager, this.serializerSettings); prepares[sequenceNumber] = new PendingTransactionState { diff --git a/src/Orleans.Transactions/State/TransactionParticipant.cs b/src/Orleans.Transactions/State/TransactionParticipant.cs index 52e3f9e0141..c0e805ba711 100644 --- a/src/Orleans.Transactions/State/TransactionParticipant.cs +++ b/src/Orleans.Transactions/State/TransactionParticipant.cs @@ -67,7 +67,7 @@ private async Task StorageWork() { // get the next batch in place so it can be filled while we store the old one batchBeingSentToStorage = storageBatch; - storageBatch = new StorageBatch(batchBeingSentToStorage); + storageBatch = new StorageBatch(batchBeingSentToStorage, this.serializerSettings); // perform the actual store, and record the e-tag storageBatch.ETag = await batchBeingSentToStorage.Store(storage); diff --git a/src/Orleans.Transactions/State/TransactionParticipantExtension.cs b/src/Orleans.Transactions/State/TransactionParticipantExtension.cs index 842ff7c98b4..f13dbefbe73 100644 --- a/src/Orleans.Transactions/State/TransactionParticipantExtension.cs +++ b/src/Orleans.Transactions/State/TransactionParticipantExtension.cs @@ -3,11 +3,14 @@ using System; using System.Collections.Generic; using System.Threading.Tasks; +using Newtonsoft.Json; namespace Orleans.Transactions { + [Serializable] public class TransactionParticipantExtension : ITransactionParticipantExtension { + [JsonProperty] private readonly Dictionary localParticipants = new Dictionary(); public void Register(string resourceId, ITransactionParticipant localTransactionParticipant) diff --git a/src/Orleans.Transactions/State/TransactionalState.cs b/src/Orleans.Transactions/State/TransactionalState.cs index 8765c319e22..fb80c44a6db 100644 --- a/src/Orleans.Transactions/State/TransactionalState.cs +++ b/src/Orleans.Transactions/State/TransactionalState.cs @@ -89,6 +89,7 @@ private Dictionary> confirmationTasks private CausalClock clock; + private JsonSerializerSettings serializerSettings; // collection tasks private Dictionary unprocessedPreparedMessages; private class PMessages @@ -120,11 +121,8 @@ IClock clock lockWorker = new BatchWorkerFromDelegate(LockWork); storageWorker = new BatchWorkerFromDelegate(StorageWork); confirmationWorker = new BatchWorkerFromDelegate(ConfirmationWork); - - if (MetaData.SerializerSettings == null) - { - MetaData.SerializerSettings = TransactionParticipantExtensionExtensions.GetJsonSerializerSettings(typeResolver, grainFactory); - } + + this.serializerSettings = TransactionParticipantExtensionExtensions.GetJsonSerializerSettings(typeResolver, grainFactory); } #region lifecycle @@ -190,7 +188,8 @@ private async Task Restore() var loadresponse = await loadtask; - storageBatch = new StorageBatch(loadresponse); + storageBatch = new StorageBatch(loadresponse, this.serializerSettings); + stableState = loadresponse.CommittedState; stableSequenceNumber = loadresponse.CommittedSequenceId; @@ -208,9 +207,8 @@ private async Task Restore() { if (logger.IsEnabled(LogLevel.Debug)) logger.Debug($"recover two-phase-commit {pr.TransactionId}"); - var tm = (pr.TransactionManager == null) ? null : - (ITransactionParticipant) JsonConvert.DeserializeObject(pr.TransactionManager, MetaData.SerializerSettings); + (ITransactionParticipant) JsonConvert.DeserializeObject(pr.TransactionManager, this.serializerSettings); commitQueue.Add(new TransactionRecord() { diff --git a/test/Transactions/Orleans.Transactions.Azure.Test/TestFixture.cs b/test/Transactions/Orleans.Transactions.Azure.Test/TestFixture.cs index 31958e2cd57..ab3bdcfa4bf 100644 --- a/test/Transactions/Orleans.Transactions.Azure.Test/TestFixture.cs +++ b/test/Transactions/Orleans.Transactions.Azure.Test/TestFixture.cs @@ -3,9 +3,7 @@ using Orleans.Hosting; using Orleans.TestingHost; using Orleans.Transactions.Tests; -using Orleans.TestingHost.Utils; using TestExtensions; -using Microsoft.Extensions.Logging; using Tester; namespace Orleans.Transactions.AzureStorage.Tests diff --git a/test/Transactions/Orleans.Transactions.Tests/Memory/SerializationTests.cs b/test/Transactions/Orleans.Transactions.Tests/Memory/SerializationTests.cs index b32bf87bad9..e13dadcae53 100644 --- a/test/Transactions/Orleans.Transactions.Tests/Memory/SerializationTests.cs +++ b/test/Transactions/Orleans.Transactions.Tests/Memory/SerializationTests.cs @@ -1,8 +1,10 @@ using System; using System.Collections.Generic; using System.Reflection; +using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; using Newtonsoft.Json; +using Orleans.Providers; using Xunit; using Orleans.Runtime; using Orleans.Runtime.Configuration; @@ -34,13 +36,13 @@ public void JsonConcertCanSerializeMetaData() Timestamp = DateTime.UtcNow, WriteParticipants = new List() { new TransactionParticipantExtension().AsTransactionParticipant("resourceId")} }); - MetaData.SerializerSettings = TransactionParticipantExtensionExtensions.GetJsonSerializerSettings( + var serializerSettings = TransactionParticipantExtensionExtensions.GetJsonSerializerSettings( this.environment.Client.ServiceProvider.GetService(), this.environment.GrainFactory); //should be able to serialize it - var jsonMetaData = JsonConvert.SerializeObject(metaData, MetaData.SerializerSettings); + var jsonMetaData = JsonConvert.SerializeObject(metaData, serializerSettings); - var deseriliazedMetaData = JsonConvert.DeserializeObject(jsonMetaData, MetaData.SerializerSettings); + var deseriliazedMetaData = JsonConvert.DeserializeObject(jsonMetaData, serializerSettings); Assert.Equal(metaData.TimeStamp, deseriliazedMetaData.TimeStamp); } } From be0c39ed89067ad75744a2b0bf0302b3abf1dd40 Mon Sep 17 00:00:00 2001 From: xiazen Date: Wed, 27 Jun 2018 14:20:38 -0700 Subject: [PATCH 5/6] pr feedback --- ...ansactionParticipantExtensionExtensions.cs | 1 - .../State/StorageBatch.cs | 20 ++++++++----------- .../State/TransactionParticipant.cs | 2 +- .../State/TransactionParticipantExtension.cs | 2 -- .../State/TransactionalState.cs | 7 +++---- .../State/TransactionalStateFactory.cs | 10 +++++++--- .../Memory/SerializationTests.cs | 20 +++++++++++-------- 7 files changed, 31 insertions(+), 31 deletions(-) diff --git a/src/Orleans.Transactions/Abstractions/Extensions/TransactionParticipantExtensionExtensions.cs b/src/Orleans.Transactions/Abstractions/Extensions/TransactionParticipantExtensionExtensions.cs index 3f5bf9892c0..ebd4ea890a7 100644 --- a/src/Orleans.Transactions/Abstractions/Extensions/TransactionParticipantExtensionExtensions.cs +++ b/src/Orleans.Transactions/Abstractions/Extensions/TransactionParticipantExtensionExtensions.cs @@ -33,7 +33,6 @@ public static JsonSerializerSettings GetJsonSerializerSettings(ITypeResolver typ [Immutable] internal sealed class TransactionParticipantExtensionWrapper : ITransactionParticipant { - [JsonProperty] private readonly ITransactionParticipantExtension extension; [JsonProperty] private readonly string resourceId; diff --git a/src/Orleans.Transactions/State/StorageBatch.cs b/src/Orleans.Transactions/State/StorageBatch.cs index 2feca30d0b7..6b1c1ae7263 100644 --- a/src/Orleans.Transactions/State/StorageBatch.cs +++ b/src/Orleans.Transactions/State/StorageBatch.cs @@ -74,7 +74,7 @@ internal class StorageBatch : ITransactionalStateStorageEvents private int confirm = 0; private int collect = 0; private int cancel = 0; - private JsonSerializerSettings serializerSettings; + public JsonSerializerSettings SerializerSettings { get; private set; } public MetaData MetaData { get; private set; } public string ETag { get; set; } @@ -85,23 +85,19 @@ public override string ToString() return $"batchsize={total} [{read}r {prepare}p {commit}c {confirm}cf {collect}cl {cancel}cc]"; } - public StorageBatch(TransactionalStorageLoadResponse loadresponse, JsonSerializerSettings settings) + public StorageBatch(TransactionalStorageLoadResponse loadresponse, JsonSerializerSettings serializerSettings) { - if(settings == null) - throw new ArgumentNullException(nameof(settings)); - this.serializerSettings = settings; - MetaData = ReadMetaData(loadresponse, this.serializerSettings); + this.SerializerSettings = serializerSettings??throw new ArgumentNullException(nameof(serializerSettings)); + MetaData = ReadMetaData(loadresponse, this.SerializerSettings); ETag = loadresponse.ETag; confirmUpTo = loadresponse.CommittedSequenceId; cancelAbove = loadresponse.PendingStates.LastOrDefault()?.SequenceId ?? loadresponse.CommittedSequenceId; cancelAboveStart = cancelAbove; } - public StorageBatch(StorageBatch previous, JsonSerializerSettings settings) + public StorageBatch(StorageBatch previous) { - if (settings == null) - throw new ArgumentNullException(nameof(settings)); - this.serializerSettings = settings; + this.SerializerSettings = previous.SerializerSettings; MetaData = previous.MetaData; confirmUpTo = previous.confirmUpTo; cancelAbove = previous.cancelAbove; @@ -127,7 +123,7 @@ private static MetaData ReadMetaData(TransactionalStorageLoadResponse lo public Task Store(ITransactionalStateStorage storage) { - var jsonMetaData = JsonConvert.SerializeObject(MetaData, this.serializerSettings); + var jsonMetaData = JsonConvert.SerializeObject(MetaData, this.SerializerSettings); var list = new List>(); if (prepares != null) @@ -180,7 +176,7 @@ public void Prepare(long sequenceNumber, Guid transactionId, DateTime timestamp, prepares = new SortedDictionary>(); var tmstring = (transactionManager == null) ? null : - JsonConvert.SerializeObject(transactionManager, this.serializerSettings); + JsonConvert.SerializeObject(transactionManager, this.SerializerSettings); prepares[sequenceNumber] = new PendingTransactionState { diff --git a/src/Orleans.Transactions/State/TransactionParticipant.cs b/src/Orleans.Transactions/State/TransactionParticipant.cs index c0e805ba711..52e3f9e0141 100644 --- a/src/Orleans.Transactions/State/TransactionParticipant.cs +++ b/src/Orleans.Transactions/State/TransactionParticipant.cs @@ -67,7 +67,7 @@ private async Task StorageWork() { // get the next batch in place so it can be filled while we store the old one batchBeingSentToStorage = storageBatch; - storageBatch = new StorageBatch(batchBeingSentToStorage, this.serializerSettings); + storageBatch = new StorageBatch(batchBeingSentToStorage); // perform the actual store, and record the e-tag storageBatch.ETag = await batchBeingSentToStorage.Store(storage); diff --git a/src/Orleans.Transactions/State/TransactionParticipantExtension.cs b/src/Orleans.Transactions/State/TransactionParticipantExtension.cs index f13dbefbe73..16056d2a5f4 100644 --- a/src/Orleans.Transactions/State/TransactionParticipantExtension.cs +++ b/src/Orleans.Transactions/State/TransactionParticipantExtension.cs @@ -7,10 +7,8 @@ namespace Orleans.Transactions { - [Serializable] public class TransactionParticipantExtension : ITransactionParticipantExtension { - [JsonProperty] private readonly Dictionary localParticipants = new Dictionary(); public void Register(string resourceId, ITransactionParticipant localTransactionParticipant) diff --git a/src/Orleans.Transactions/State/TransactionalState.cs b/src/Orleans.Transactions/State/TransactionalState.cs index fb80c44a6db..479729bfdf3 100644 --- a/src/Orleans.Transactions/State/TransactionalState.cs +++ b/src/Orleans.Transactions/State/TransactionalState.cs @@ -105,8 +105,7 @@ public TransactionalState( ITransactionAgent transactionAgent, IProviderRuntime runtime, ILoggerFactory loggerFactory, - ITypeResolver typeResolver, - IGrainFactory grainFactory, + JsonSerializerSettings serializerSettings, IClock clock ) { @@ -121,8 +120,8 @@ IClock clock lockWorker = new BatchWorkerFromDelegate(LockWork); storageWorker = new BatchWorkerFromDelegate(StorageWork); confirmationWorker = new BatchWorkerFromDelegate(ConfirmationWork); - - this.serializerSettings = TransactionParticipantExtensionExtensions.GetJsonSerializerSettings(typeResolver, grainFactory); + + this.serializerSettings = serializerSettings; } #region lifecycle diff --git a/src/Orleans.Transactions/State/TransactionalStateFactory.cs b/src/Orleans.Transactions/State/TransactionalStateFactory.cs index a4cf8a97e04..b9b79767380 100644 --- a/src/Orleans.Transactions/State/TransactionalStateFactory.cs +++ b/src/Orleans.Transactions/State/TransactionalStateFactory.cs @@ -1,21 +1,25 @@ using Microsoft.Extensions.DependencyInjection; +using Newtonsoft.Json; using Orleans.Runtime; using Orleans.Transactions.Abstractions; +using Orleans.Transactions.Abstractions.Extensions; namespace Orleans.Transactions { public class TransactionalStateFactory : ITransactionalStateFactory { private IGrainActivationContext context; - - public TransactionalStateFactory(IGrainActivationContext context) + private JsonSerializerSettings serializerSettings; + public TransactionalStateFactory(IGrainActivationContext context, ITypeResolver typeResolver, IGrainFactory grainFactory) { this.context = context; + this.serializerSettings = + TransactionParticipantExtensionExtensions.GetJsonSerializerSettings(typeResolver, grainFactory); } public ITransactionalState Create(ITransactionalStateConfiguration config) where TState : class, new() { - TransactionalState transactionalState = ActivatorUtilities.CreateInstance>(this.context.ActivationServices, config, this.context); + TransactionalState transactionalState = ActivatorUtilities.CreateInstance>(this.context.ActivationServices, config, this.serializerSettings, this.context); transactionalState.Participate(context.ObservableLifecycle); return transactionalState; } diff --git a/test/Transactions/Orleans.Transactions.Tests/Memory/SerializationTests.cs b/test/Transactions/Orleans.Transactions.Tests/Memory/SerializationTests.cs index e13dadcae53..1250cbca0c6 100644 --- a/test/Transactions/Orleans.Transactions.Tests/Memory/SerializationTests.cs +++ b/test/Transactions/Orleans.Transactions.Tests/Memory/SerializationTests.cs @@ -11,34 +11,38 @@ using Orleans.Serialization; using Orleans.Transactions.Abstractions; using Orleans.Transactions.Abstractions.Extensions; +using Orleans.Transactions.Tests.Correctness; using TestExtensions; +using Xunit.Abstractions; namespace Orleans.Transactions.Tests.Memory { [TestCategory("BVT"), TestCategory("Transactions")] - public class SerializationTests + public class SerializationTests: TransactionTestRunnerBase, IClassFixture { - private SerializationTestEnvironment environment; - public SerializationTests() + private MemoryTransactionsFixture fixture; + public SerializationTests(MemoryTransactionsFixture fixture, ITestOutputHelper output) + :base(fixture.GrainFactory, output) { - var config = new ClientConfiguration { SerializationProviders = { typeof(OrleansJsonSerializer).GetTypeInfo() } }; - this.environment = SerializationTestEnvironment.InitializeWithDefaults(config); + this.fixture = fixture; } [Fact] public void JsonConcertCanSerializeMetaData() { + var grainRef = this.RandomTestGrain(TransactionTestConstants.SingleStateTransactionalGrain); + var ext = grainRef.Cast(); var metaData = new MetaData(); metaData.TimeStamp = DateTime.UtcNow; metaData.CommitRecords = new Dictionary(); metaData.CommitRecords.Add(Guid.NewGuid(), new CommitRecord() { Timestamp = DateTime.UtcNow, - WriteParticipants = new List() { new TransactionParticipantExtension().AsTransactionParticipant("resourceId")} + WriteParticipants = new List() { ext.AsTransactionParticipant("resourceId")} }); var serializerSettings = TransactionParticipantExtensionExtensions.GetJsonSerializerSettings( - this.environment.Client.ServiceProvider.GetService(), - this.environment.GrainFactory); + this.fixture.Client.ServiceProvider.GetService(), + this.grainFactory); //should be able to serialize it var jsonMetaData = JsonConvert.SerializeObject(metaData, serializerSettings); From 73713877e5a8413d188ac68142598bb3ea945c06 Mon Sep 17 00:00:00 2001 From: xiazen Date: Wed, 27 Jun 2018 16:59:31 -0700 Subject: [PATCH 6/6] pr feedback --- src/Orleans.Transactions/State/StorageBatch.cs | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/src/Orleans.Transactions/State/StorageBatch.cs b/src/Orleans.Transactions/State/StorageBatch.cs index 6b1c1ae7263..4f484f87475 100644 --- a/src/Orleans.Transactions/State/StorageBatch.cs +++ b/src/Orleans.Transactions/State/StorageBatch.cs @@ -74,7 +74,7 @@ internal class StorageBatch : ITransactionalStateStorageEvents private int confirm = 0; private int collect = 0; private int cancel = 0; - public JsonSerializerSettings SerializerSettings { get; private set; } + private readonly JsonSerializerSettings serializerSettings; public MetaData MetaData { get; private set; } public string ETag { get; set; } @@ -87,8 +87,8 @@ public override string ToString() public StorageBatch(TransactionalStorageLoadResponse loadresponse, JsonSerializerSettings serializerSettings) { - this.SerializerSettings = serializerSettings??throw new ArgumentNullException(nameof(serializerSettings)); - MetaData = ReadMetaData(loadresponse, this.SerializerSettings); + this.serializerSettings = serializerSettings ?? throw new ArgumentNullException(nameof(serializerSettings)); + MetaData = ReadMetaData(loadresponse, this.serializerSettings); ETag = loadresponse.ETag; confirmUpTo = loadresponse.CommittedSequenceId; cancelAbove = loadresponse.PendingStates.LastOrDefault()?.SequenceId ?? loadresponse.CommittedSequenceId; @@ -97,7 +97,7 @@ public StorageBatch(TransactionalStorageLoadResponse loadresponse, JsonS public StorageBatch(StorageBatch previous) { - this.SerializerSettings = previous.SerializerSettings; + this.serializerSettings = previous.serializerSettings; MetaData = previous.MetaData; confirmUpTo = previous.confirmUpTo; cancelAbove = previous.cancelAbove; @@ -123,7 +123,7 @@ private static MetaData ReadMetaData(TransactionalStorageLoadResponse lo public Task Store(ITransactionalStateStorage storage) { - var jsonMetaData = JsonConvert.SerializeObject(MetaData, this.SerializerSettings); + var jsonMetaData = JsonConvert.SerializeObject(MetaData, this.serializerSettings); var list = new List>(); if (prepares != null) @@ -176,7 +176,7 @@ public void Prepare(long sequenceNumber, Guid transactionId, DateTime timestamp, prepares = new SortedDictionary>(); var tmstring = (transactionManager == null) ? null : - JsonConvert.SerializeObject(transactionManager, this.SerializerSettings); + JsonConvert.SerializeObject(transactionManager, this.serializerSettings); prepares[sequenceNumber] = new PendingTransactionState {