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; } } } 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.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)) 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/src/Orleans.Transactions/Abstractions/Extensions/TransactionParticipantExtensionExtensions.cs b/src/Orleans.Transactions/Abstractions/Extensions/TransactionParticipantExtensionExtensions.cs index b8d71730824..ebd4ea890a7 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; } @@ -37,6 +34,7 @@ public static JsonSerializerSettings GetJsonSerializerSettings(ITypeResolver typ internal sealed class TransactionParticipantExtensionWrapper : ITransactionParticipant { 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..4f484f87475 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,21 +74,21 @@ internal class StorageBatch : ITransactionalStateStorageEvents private int confirm = 0; private int collect = 0; private int cancel = 0; - + private readonly 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 serializerSettings) { - MetaData = ReadMetaData(loadresponse); + 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,15 +97,16 @@ public StorageBatch(TransactionalStorageLoadResponse loadresponse) public StorageBatch(StorageBatch previous) { + this.serializerSettings = previous.serializerSettings; 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 +117,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 +176,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/TransactionParticipantExtension.cs b/src/Orleans.Transactions/State/TransactionParticipantExtension.cs index 842ff7c98b4..16056d2a5f4 100644 --- a/src/Orleans.Transactions/State/TransactionParticipantExtension.cs +++ b/src/Orleans.Transactions/State/TransactionParticipantExtension.cs @@ -3,6 +3,7 @@ using System; using System.Collections.Generic; using System.Threading.Tasks; +using Newtonsoft.Json; namespace Orleans.Transactions { diff --git a/src/Orleans.Transactions/State/TransactionalState.cs b/src/Orleans.Transactions/State/TransactionalState.cs index 8765c319e22..479729bfdf3 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 @@ -104,8 +105,7 @@ public TransactionalState( ITransactionAgent transactionAgent, IProviderRuntime runtime, ILoggerFactory loggerFactory, - ITypeResolver typeResolver, - IGrainFactory grainFactory, + JsonSerializerSettings serializerSettings, IClock clock ) { @@ -121,10 +121,7 @@ IClock clock storageWorker = new BatchWorkerFromDelegate(StorageWork); confirmationWorker = new BatchWorkerFromDelegate(ConfirmationWork); - if (MetaData.SerializerSettings == null) - { - MetaData.SerializerSettings = TransactionParticipantExtensionExtensions.GetJsonSerializerSettings(typeResolver, grainFactory); - } + this.serializerSettings = serializerSettings; } #region lifecycle @@ -190,7 +187,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 +206,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/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/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); 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..1250cbca0c6 100644 --- a/test/Transactions/Orleans.Transactions.Tests/Memory/SerializationTests.cs +++ b/test/Transactions/Orleans.Transactions.Tests/Memory/SerializationTests.cs @@ -1,46 +1,52 @@ 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; 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")} }); - MetaData.SerializerSettings = TransactionParticipantExtensionExtensions.GetJsonSerializerSettings( - this.environment.Client.ServiceProvider.GetService(), - this.environment.GrainFactory); + var serializerSettings = TransactionParticipantExtensionExtensions.GetJsonSerializerSettings( + this.fixture.Client.ServiceProvider.GetService(), + this.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); } }