Skip to content
This repository
tag: v0.3.0.beta.1
Fetching contributors…

Octocat-spinner-32-eaf2f5

Cannot retrieve contributors at this time

file 84 lines (69 sloc) 2.435 kb
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83
module EventMachine
  module Synchrony

    class ConnectionPool
      def initialize(opts, &block)
        @reserved = {} # map of in-progress connections
        @available = [] # pool of free connections
        @pending = [] # pending reservations (FIFO)

        opts[:size].times do
          @available.push(block.call) if block_given?
        end
      end

      # Choose first available connection and pass it to the supplied
      # block. This will block indefinitely until there is an available
      # connection to service the request.
      def execute(async)
        f = Fiber.current

        begin
          conn = acquire(f)
          yield conn
        ensure
          release(f) if not async
        end
      end

      private

        # Acquire a lock on a connection and assign it to executing fiber
        # - if connection is available, pass it back to the calling block
        # - if pool is full, yield the current fiber until connection is available
        def acquire(fiber)

          if conn = @available.pop
            @reserved[fiber.object_id] = conn
            conn
          else
            Fiber.yield @pending.push fiber
            acquire(fiber)
          end
        end

        # Release connection assigned to the supplied fiber and
        # resume any other pending connections (which will
        # immediately try to run acquire on the pool)
        def release(fiber)
          @available.push(@reserved.delete(fiber.object_id))

          if pending = @pending.pop
            pending.resume
          end
        end

        # Allow the pool to behave as the underlying connection
        #
        # If the requesting method begins with "a" prefix, then
        # hijack the callbacks and errbacks to fire a connection
        # pool release whenever the request is complete. Otherwise
        # yield the connection within execute method and release
        # once it is complete (assumption: fiber will yield until
        # data is available, or request is complete)
        #
        def method_missing(method, *args)
          async = (method[0,1] == "a")

          execute(async) do |conn|
            df = conn.send(method, *args)

            if async
              fiber = Fiber.current
              df.callback { release(fiber) }
              df.errback { release(fiber) }
            end

            df
          end
        end
    end

  end
end
Something went wrong with that request. Please try again.