Created
April 15, 2020 18:50
-
-
Save deepj/9d2626b76b3ea617eae3c3d7e7483ac0 to your computer and use it in GitHub Desktop.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
# frozen_string_literal: true | |
require 'async/io' | |
require 'async/io/stream' | |
require 'async/pool/controller' | |
require 'async/await' | |
module Async | |
module NATS | |
def self.local_endpoint | |
Async::IO::Endpoint.tcp('localhost', 4222) | |
end | |
class Client | |
include Async::Await | |
def initialize(endpoint = NATS.local_endpoint, **options) | |
@endpoint = endpoint | |
@pool = connect(**options) | |
end | |
def ping | |
call("PING") | |
end | |
def call(*arguments) | |
@pool.acquire do |connection| | |
connection.write_lines(arguments) | |
connection.flush | |
return connection.read_response | |
end | |
end | |
private | |
def connect(**options) | |
Async::Pool::Controller.wrap(**options) do | |
peer = @endpoint.connect | |
peer.setsockopt(Socket::IPPROTO_TCP, Socket::TCP_NODELAY, 1) | |
stream = IO::Stream.new(peer) | |
@protocol = Async::IO::Protocol::Line.new(stream, CLRF) | |
end | |
end | |
end | |
end | |
end | |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment