gRPC Push and Fan-Out

Viewed 6134

Is it possible and if so good practice to use gRPC as a push service with fan-out functionality? In the examples given from Google, there is following code for the server side (C#):

    public override async Task ListFeatures(Rectangle request, IServerStreamWriter<Feature> responseStream, ServerCallContext context)
    {
        var responses = features.FindAll( (feature) => feature.Exists() && request.Contains(feature.Location) );
        foreach (var response in responses)
        {
            await responseStream.WriteAsync(response);
        }
    }

The problem is here that:

  1. Data is only generated and written if the client explicitly asks for it.
  2. Only the client that asked will get the new data.

I think what I need is:

  1. Keep all the IServerStreamWriters for each client asked (subscribed).
  2. Trigger the write by an external event when new data is available.
  3. Write to all streamWriters

Edit: Based on Carl's suggestions I now have the following:

Proto:

service PubSub {
 rpc Subscribe(Subscription) returns (stream Event) {}
 rpc Unsubscribe(Subscription) returns (Unsubscription) {}
}

message Event
{
   string Value = 1;
}
message Subscription
{
  string Id = 1;
}
message Unsubscription
{
  string Id = 1;
}

PubSubImpl:

public class PubSubImpl : PubSub.PubSubBase
{
    private readonly BufferBlock<Event> _buffer = new BufferBlock<Event>();

    private Dictionary<string, IServerStreamWriter<Event>> _subscriberWritersMap =
        new Dictionary<string, IServerStreamWriter<Event>>();

     public override async Task Subscribe(Subscription subscription, IServerStreamWriter<Event> responseStream, ServerCallContext context)
    {
        //Dict to hold a streamWriter for each subscriber.
        _subscriberWritersMap[subscription.Id] = responseStream;

        while (_subscriberWritersMap.ContainsKey(subscription.Id))
        {
            //Wait on BufferBlock from MS Dataflow package.
            var @event = await _buffer.ReceiveAsync();
            foreach (var serverStreamWriter in _subscriberWritersMap.Values)
            {
                await serverStreamWriter.WriteAsync(@event);
            }
        }
    }

    public override Task<Unsubscription> Unsubscribe(Subscription request, ServerCallContext context)
    {
        _subscriberWritersMap.Remove(request.Id);
        return Task.FromResult(new Unsubscription() { Id = request.Id });
    }

    public void Publish(string input)
    {
        _buffer.Post(new Event() { Value = input });
    }
}

"Push"s can now be sent like this:

   while ((input = Console.ReadLine()) != "q")
    {
      pubsubImp.Publish(input);
    }

On client side I have:

public async Task Subscribe()
{
    _subscription = new Subscription() { Id = Guid.NewGuid().ToString() };
    using (var call = _pubSubClient.Subscribe(_subscription))
    {
        //Receive
        var responseReaderTask = Task.Run(async () =>
        {
            while (await call.ResponseStream.MoveNext())
            {
                Console.WriteLine("Event received: " + call.ResponseStream.Current);
            }
        });

        await responseReaderTask;
    }
}

public void Unsubscribe()
{
    _pubSubClient.Unsubscribe(_subscription);
}

Client-Main works like this:

static void Main(string[] args)
{
    var channel = new Channel("127.0.0.1:50052", 
                                ChannelCredentials.Insecure);
    var subscriber = new Subsriber(new PubSub.PubSubClient(channel));

    Task.Run(async () =>
    {
        await subscriber.Subscribe();
    }).GetAwaiter();

    Console.WriteLine("Hit key to unsubscribe");
    Console.ReadLine();

    subscriber.Unsubscribe();

    Console.WriteLine("Unsubscribed...");

    Console.WriteLine("Hit key to exit...");
    Console.ReadLine();

}

At the moment it looks like it works. Is that how it should/could be done? Test solution can be found at: https://github.com/KingKnecht/gRPC-PubSub

2 Answers

I don't think you need to keep track of the clients in your server-app. Instead, for each "push" method, create an EventWaitHandle and continuously wait on it. Then in another context, signal the EventWaitHandle so that the waiting "push" method can WriteAsync to the client.

e.g.

public class PubSubImpl : PubSub.PubSubBase
{
    EventWaitHandle evwStatus = new EventWaitHandle(false, EventResetMode.ManualReset);
    string status;

    public override async Task Subscribe(Subscription subscription, IServerStreamWriter<Event> responseStream, ServerCallContext context)
    {
        // respond with the current status
        await serverStreamWriter.WriteAsync(new Event { Value = status });
        // wait until we're signaled with a different status
        while(evwStatus.WaitOne())
        {
            await serverStreamWriter.WriteAsync(new Event { Value = status });
        }
    }   

    public void Publish(string input)
    {
        status = input;
        // let the waiting threads respond
        evwStatus.Set();
        // halt the waiting threads
        evwStatus.Reset();
    }
}
Related