BaseProtocolServer
Base class for protocol servers (FIX, SBE, etc.). Contains shared infrastructure: lifecycle, sessions, subscriptions, message channels, background tasks.
Inherits: BaseLogReceiver
Implements: IProtocolServer, ILogReceiver, ILogSource, IDisposable
Constructors
protected BaseProtocolServer(TSettings settings)
baseProtocolServer = BaseProtocolServer(settings)
Initializes a new instance.
- settings
- Server settings.
Properties
protected LoginRateLimiter LoginRateLimiter { get; private set; }
value = baseProtocolServer.LoginRateLimiter
baseProtocolServer.LoginRateLimiter = value
Rate limiter for login attempts (shared across all listeners).
public IEnumerable<IMessageListenerSession> Sessions { get; }
value = baseProtocolServer.Sessions
Snapshot of all client sessions currently held by the server. Useful for tests that need to verify server-side session bookkeeping independently of what the client adapter happens to see (a FIX adapter, for example, silently auto-reconnects after a server-side kick, so the client-side DisconnectMessage is not a reliable signal).
public ChannelStates State { get; protected set; }
value = baseProtocolServer.State
baseProtocolServer.State = value
Current server state.
Methods
public void AddSubscription(ServerSubscription subscription)
baseProtocolServer.AddSubscription(subscription)
Add subscription.
- subscription
- Subscription.
protected virtual void ClearState()
baseProtocolServer.ClearState()
Clear all state: listeners, sessions, subscriptions.
protected string Convert(string value)
result = baseProtocolServer.Convert(value)
Convert text to latin if configured.
protected abstract TSubscription CreateSubscription(TClientSession session, string requestId, ServerSubscription subscription, MessageTypes type)
result = baseProtocolServer.CreateSubscription(session, requestId, subscription, type)
Create a subscription info instance.
public void Disconnect(IMessageListenerSession session)
baseProtocolServer.Disconnect(session)
Disconnect session.
- session
- Session.
protected override void DisposeManaged()
baseProtocolServer.DisposeManaged()
Release resources.
protected void EnqueueMessage(TClientSession session, string requestId, Func<TSubscription, TClientSession, string, CancellationToken, ValueTask> handler, CancellationToken cancellationToken)
baseProtocolServer.EnqueueMessage(session, requestId, handler, cancellationToken)
Enqueue message for a session with request ID.
protected void EnqueueMessage(TSubscription subscription, TClientSession session, string requestId, Func<TSubscription, TClientSession, string, CancellationToken, ValueTask<bool>> handler, CancellationToken cancellationToken)
baseProtocolServer.EnqueueMessage(subscription, session, requestId, handler, cancellationToken)
Core enqueue: schedule work on a client session.
protected void EnqueueMessage(TSubscription[] subscribers, Func<TSubscription, TClientSession, string, CancellationToken, ValueTask> handler, CancellationToken cancellationToken)
baseProtocolServer.EnqueueMessage(subscribers, handler, cancellationToken)
Enqueue message for multiple subscribers.
protected void EnqueueMessage(TSubscription[] subscribers, Func<TSubscription, TClientSession, string, CancellationToken, ValueTask<bool>> handler, CancellationToken cancellationToken)
baseProtocolServer.EnqueueMessage(subscribers, handler, cancellationToken)
Enqueue message for multiple subscribers (with bool result).
protected void EnqueueMessage(TSubscription subscription, Func<TSubscription, TClientSession, string, CancellationToken, ValueTask> handler, CancellationToken cancellationToken)
baseProtocolServer.EnqueueMessage(subscription, handler, cancellationToken)
Enqueue message for a single subscription.
protected void EnqueueMessage(TSubscription subscription, Func<TSubscription, TClientSession, string, CancellationToken, ValueTask<bool>> handler, CancellationToken cancellationToken)
baseProtocolServer.EnqueueMessage(subscription, handler, cancellationToken)
Enqueue message for a single subscription (with bool result).
protected void EnqueueMessage(TClientSession session, Func<TSubscription, TClientSession, string, CancellationToken, ValueTask> handler, CancellationToken cancellationToken)
baseProtocolServer.EnqueueMessage(session, handler, cancellationToken)
Enqueue message for a session (no request ID).
protected abstract string GetSessionKey(TClientSession session)
result = baseProtocolServer.GetSessionKey(session)
Get the key used to group client sessions (e.g. SenderCompId).
public IEnumerable<ServerSubscription> GetSubscriptions(IMessageListenerSession session)
result = baseProtocolServer.GetSubscriptions(session)
Get subscription for the specified session.
- session
- Session.
Returns: Subscriptions.
public bool HasSubscriptions(DataType dataType, SecurityId securityId)
result = baseProtocolServer.HasSubscriptions(dataType, securityId)
Are there subscribers.
protected abstract bool IsStubSubscription(TSubscription subscription)
result = baseProtocolServer.IsStubSubscription(subscription)
Check if subscription is a stub (placeholder).
public override void Load(SettingsStorage storage)
baseProtocolServer.Load(storage)
Load settings.
- storage
- Settings storage.
protected abstract void OnClearListenerState()
baseProtocolServer.OnClearListenerState()
Clear protocol-specific listener state (stop listeners, etc.).
protected abstract void OnStartListeners(CancellationToken cancellationToken)
baseProtocolServer.OnStartListeners(cancellationToken)
Start protocol-specific listeners (TCP, etc.). Add tasks to _backgroundTasks.
protected abstract void ProcessInMessage(Message message, CancellationToken cancellationToken)
baseProtocolServer.ProcessInMessage(message, cancellationToken)
Process an incoming message and route to subscribed clients.
protected virtual void ProcessLogout(TClientSession session)
baseProtocolServer.ProcessLogout(session)
Process client logout/disconnect.
protected void RaiseNewOutMessage(IMessageListenerSession session, Message message)
baseProtocolServer.RaiseNewOutMessage(session, message)
Write outgoing message to channel.
protected void RaiseSessionConnected(IMessageListenerSession session)
baseProtocolServer.RaiseSessionConnected(session)
Raise session connected event.
protected void RaiseSessionDisconnected(IMessageListenerSession session)
baseProtocolServer.RaiseSessionDisconnected(session)
Raise session disconnected event.
public bool RemoveSubscription(ServerSubscription subscription)
result = baseProtocolServer.RemoveSubscription(subscription)
Remove subscription.
- subscription
- Subscription.
Returns: if subscription was found, otherwise .
public bool Resume(ServerSubscription subscription)
result = baseProtocolServer.Resume(subscription)
Resume subscription.
- subscription
- Subscription.
Returns: if subscription was found, otherwise .
public void Resume(IMessageListenerSession session)
baseProtocolServer.Resume(session)
Resume session.
- session
- Session.
public IAsyncEnumerable<ValueTuple<IMessageListenerSession, Message>> RunAsync(CancellationToken cancellationToken)
result = baseProtocolServer.RunAsync(cancellationToken)
Run async message processing and return outgoing messages as IAsyncEnumerable. Opens server at start and closes at end.
public override void Save(SettingsStorage storage)
baseProtocolServer.Save(storage)
Save settings.
- storage
- Settings storage.
protected abstract void ScheduleSessionWork(TClientSession session, ServerSubscription subscription, string requestId, Func<TClientSession, string, CancellationToken, ValueTask<bool>> handler)
baseProtocolServer.ScheduleSessionWork(session, subscription, requestId, handler)
Schedule work on a client session's outgoing queue.
public virtual bool SendInMessage(Message message)
result = baseProtocolServer.SendInMessage(message)
Send incoming message to server for processing and routing to clients.
protected void StartServer(CancellationToken cancellationToken)
baseProtocolServer.StartServer(cancellationToken)
Start the server: initialize state, start listeners and background tasks.
protected void StopServer()
baseProtocolServer.StopServer()
Stop the server: cancel tasks, complete channels, clear state.
public void Suspend(IMessageListenerSession session)
baseProtocolServer.Suspend(session)
Suspend session.
- session
- Session.
public bool Suspend(ServerSubscription subscription)
result = baseProtocolServer.Suspend(subscription)
Suspend subscription.
- subscription
- Subscription.
Returns: if subscription was found, otherwise .
Events
public event Action<IMessageListenerSession> SessionConnected
baseProtocolServer.SessionConnected += handler
Session connected event.
public event Action<IMessageListenerSession> SessionDisconnected
baseProtocolServer.SessionDisconnected += handler
Session disconnected event.
public event Action StateChanged
baseProtocolServer.StateChanged += handler
State change event.
public event Action<ServerSubscription> SubscriptionChanged
baseProtocolServer.SubscriptionChanged += handler
Client subscription changed event.
Fields
protected readonly List<Task> _backgroundTasks
value = baseProtocolServer._backgroundTasks
Background tasks (accept clients, process messages, etc.).
protected readonly SynchronizedSet<long> _cancellationRequests
value = baseProtocolServer._cancellationRequests
Pending order cancellation request transaction IDs.
protected readonly SynchronizedSet<ServerSubscription> _changedSubscriptions
value = baseProtocolServer._changedSubscriptions
Changed subscriptions pending notification.
protected readonly CachedSynchronizedDictionary<string, CachedSynchronizedList<TClientSession>> _clientSessions
value = baseProtocolServer._clientSessions
Client sessions grouped by session key (e.g. SenderCompId).
protected readonly Channel<Message> _incomingMessages
value = baseProtocolServer._incomingMessages
Incoming messages channel.
protected readonly Channel<ValueTuple<IMessageListenerSession, Message>> _outgoingMessages
value = baseProtocolServer._outgoingMessages
Outgoing messages channel (session + message pairs).
protected readonly SubscriptionHolder<TSubscription, TClientSession> _subscriptions
value = baseProtocolServer._subscriptions
Subscription holder.