Skip to content

Instantly share code, notes, and snippets.

@georgebearden
Created July 24, 2015 15:10
Show Gist options
  • Select an option

  • Save georgebearden/78de0c1c9f94a1bf97f0 to your computer and use it in GitHub Desktop.

Select an option

Save georgebearden/78de0c1c9f94a1bf97f0 to your computer and use it in GitHub Desktop.
A sample implementation of using Rx.NET on top of a C# HttpListener
using System;
using System.Collections.Generic;
using System.Linq;
using System.Net;
using System.Reactive.Linq;
using System.Threading;
using System.Threading.Tasks;
namespace Tests
{
public class ReactiveServer : IObservable<HttpListenerContext>, IDisposable
{
private readonly IObservable<HttpListenerContext> _observable;
private readonly CancellationTokenSource _cancelTokenSource;
public ReactiveServer(IEnumerable<string> prefixes)
{
_cancelTokenSource = new CancellationTokenSource();
_observable = Observable.Create<HttpListenerContext>(async observer =>
{
var server = new HttpListener();
prefixes.ForEach(server.Prefixes.Add);
server.Start();
while (!_cancelTokenSource.IsCancellationRequested)
{
var context = await Task.Run(() => server.GetContext(), _cancelTokenSource.Token);
observer.OnNext(context);
}
});
}
public void Dispose()
{
if (!_cancelTokenSource.IsCancellationRequested)
_cancelTokenSource.Cancel();
}
public IDisposable Subscribe(IObserver<HttpListenerContext> observer)
{
return _observable.Subscribe(observer);
}
}
}
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment