Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 9 additions & 13 deletions Samples/2.0/docker-aspnet-core/API/Startup.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -42,6 +43,8 @@ public void Configure(IApplicationBuilder app, Microsoft.AspNetCore.Hosting.IHos

private IClusterClient CreateClusterClient(IServiceProvider serviceProvider)
{
var log = serviceProvider.GetService<ILogger<Startup>>();

// TODO replace with your connection string
const string connectionString = "YOUR_CONNECTION_STRING_HERE";
var client = new ClientBuilder()
Expand All @@ -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<bool> 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;
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
using Orleans.Runtime;
using Orleans.Serialization;
using Orleans.Transactions.Abstractions;
using Orleans.Transactions.Abstractions.Extensions;

namespace Orleans.Transactions.AzureStorage
{
Expand All @@ -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;
}

Expand Down
1 change: 1 addition & 0 deletions src/Orleans.Core/Configuration/ConfigUtilities.cs
Original file line number Diff line number Diff line change
Expand Up @@ -349,6 +349,7 @@ internal static async Task<IPAddress> 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))
Expand Down
22 changes: 17 additions & 5 deletions src/Orleans.Core/Lifecycle/LifecycleSubject.cs
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,13 @@ public LifecycleSubject(ILogger<LifecycleSubject> logger)
this.subscribers = new ConcurrentDictionary<object, OrderedObserver>();
}

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
Expand All @@ -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)
Expand All @@ -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;
Expand All @@ -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)
{
Expand All @@ -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.");
Expand Down
24 changes: 14 additions & 10 deletions src/Orleans.Runtime/Lifecycle/SiloLifecycleSubject.cs
Original file line number Diff line number Diff line change
Expand Up @@ -12,38 +12,42 @@ namespace Orleans.Runtime
/// <summary>
/// Decorator over lifecycle subject for silo. Adds some logging and monitoring
/// </summary>
public class SiloLifecycleSubject : ISiloLifecycleSubject
public class SiloLifecycleSubject : LifecycleSubject, ISiloLifecycleSubject
{
private readonly ILifecycleSubject subject;
private readonly ILogger<SiloLifecycleSubject> logger;
private readonly List<MonitoredObserver> observers;

public SiloLifecycleSubject(ILifecycleSubject subject, ILogger<SiloLifecycleSubject> logger)
public SiloLifecycleSubject(ILogger<SiloLifecycleSubject> logger)
:base(logger)
{
this.subject = subject;
this.logger = logger;
this.observers = new List<MonitoredObserver>();
}

public Task OnStart(CancellationToken ct)
public override Task OnStart(CancellationToken ct)
{
foreach(IGrouping<int,MonitoredObserver> 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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;
}
Expand All @@ -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)
Expand Down
24 changes: 12 additions & 12 deletions src/Orleans.Transactions/State/StorageBatch.cs
Original file line number Diff line number Diff line change
@@ -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
{
Expand Down Expand Up @@ -35,8 +37,6 @@ public class MetaData
public DateTime TimeStamp { get; set; }

public Dictionary<Guid, CommitRecord> CommitRecords { get; set; }

public static JsonSerializerSettings SerializerSettings { get; set; }
}

[Serializable]
Expand Down Expand Up @@ -74,21 +74,21 @@ internal class StorageBatch<TState> : ITransactionalStateStorageEvents<TState>
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<TState> loadresponse)
public StorageBatch(TransactionalStorageLoadResponse<TState> 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;
Expand All @@ -97,15 +97,16 @@ public StorageBatch(TransactionalStorageLoadResponse<TState> loadresponse)

public StorageBatch(StorageBatch<TState> previous)
{
this.serializerSettings = previous.serializerSettings;
MetaData = previous.MetaData;
confirmUpTo = previous.confirmUpTo;
cancelAbove = previous.cancelAbove;
cancelAboveStart = cancelAbove;
}

private static MetaData ReadMetaData(TransactionalStorageLoadResponse<TState> loadresponse)
private static MetaData ReadMetaData(TransactionalStorageLoadResponse<TState> 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()
Expand All @@ -116,14 +117,13 @@ private static MetaData ReadMetaData(TransactionalStorageLoadResponse<TState> lo
}
else
{
return JsonConvert.DeserializeObject<MetaData>(loadresponse.Metadata, MetaData.SerializerSettings);
return JsonConvert.DeserializeObject<MetaData>(loadresponse.Metadata, serializerSettings);
}
}

public Task<string> Store(ITransactionalStateStorage<TState> storage)
{
var jsonMetaData = JsonConvert.SerializeObject(MetaData, MetaData.SerializerSettings);

var jsonMetaData = JsonConvert.SerializeObject(MetaData, this.serializerSettings);
var list = new List<PendingTransactionState<TState>>();

if (prepares != null)
Expand Down Expand Up @@ -176,7 +176,7 @@ public void Prepare(long sequenceNumber, Guid transactionId, DateTime timestamp,
prepares = new SortedDictionary<long, PendingTransactionState<TState>>();

var tmstring = (transactionManager == null) ? null :
JsonConvert.SerializeObject(transactionManager, MetaData.SerializerSettings);
JsonConvert.SerializeObject(transactionManager, this.serializerSettings);

prepares[sequenceNumber] = new PendingTransactionState<TState>
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
using System;
using System.Collections.Generic;
using System.Threading.Tasks;
using Newtonsoft.Json;

namespace Orleans.Transactions
{
Expand Down
15 changes: 6 additions & 9 deletions src/Orleans.Transactions/State/TransactionalState.cs
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,7 @@ private Dictionary<Guid, TransactionRecord<TState>> confirmationTasks

private CausalClock clock;

private JsonSerializerSettings serializerSettings;
// collection tasks
private Dictionary<DateTime, PMessages> unprocessedPreparedMessages;
private class PMessages
Expand All @@ -104,8 +105,7 @@ public TransactionalState(
ITransactionAgent transactionAgent,
IProviderRuntime runtime,
ILoggerFactory loggerFactory,
ITypeResolver typeResolver,
IGrainFactory grainFactory,
JsonSerializerSettings serializerSettings,
IClock clock
)
{
Expand All @@ -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
Expand Down Expand Up @@ -190,7 +187,8 @@ private async Task Restore()

var loadresponse = await loadtask;

storageBatch = new StorageBatch<TState>(loadresponse);
storageBatch = new StorageBatch<TState>(loadresponse, this.serializerSettings);


stableState = loadresponse.CommittedState;
stableSequenceNumber = loadresponse.CommittedSequenceId;
Expand All @@ -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<ITransactionParticipant>(pr.TransactionManager, MetaData.SerializerSettings);
(ITransactionParticipant) JsonConvert.DeserializeObject<ITransactionParticipant>(pr.TransactionManager, this.serializerSettings);

commitQueue.Add(new TransactionRecord<TState>()
{
Expand Down
10 changes: 7 additions & 3 deletions src/Orleans.Transactions/State/TransactionalStateFactory.cs
Original file line number Diff line number Diff line change
@@ -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<TState> Create<TState>(ITransactionalStateConfiguration config) where TState : class, new()
{
TransactionalState<TState> transactionalState = ActivatorUtilities.CreateInstance<TransactionalState<TState>>(this.context.ActivationServices, config, this.context);
TransactionalState<TState> transactionalState = ActivatorUtilities.CreateInstance<TransactionalState<TState>>(this.context.ActivationServices, config, this.serializerSettings, this.context);
transactionalState.Participate(context.ObservableLifecycle);
return transactionalState;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,7 @@ public async Task<IGrainStorage> GetStorageProvider(string storageInvariant)
ServiceId = Guid.NewGuid().ToString()
};
var storageProvider = new AdoNetGrainStorage(DefaultProviderRuntime.ServiceProvider.GetService<ILogger<AdoNetGrainStorage>>(), DefaultProviderRuntime, Options.Create(options), Options.Create(clusterOptions), storageInvariant + "_StorageProvider");
ISiloLifecycleSubject siloLifeCycle = new SiloLifecycleSubject(new LifecycleSubject(NullLoggerFactory.Instance.CreateLogger<LifecycleSubject>()), NullLoggerFactory.Instance.CreateLogger<SiloLifecycleSubject>());
ISiloLifecycleSubject siloLifeCycle = new SiloLifecycleSubject(NullLoggerFactory.Instance.CreateLogger<SiloLifecycleSubject>());
storageProvider.Participate(siloLifeCycle);
await siloLifeCycle.OnStart(CancellationToken.None);

Expand Down
Loading