Class: Async::Aws::ClientCache

Inherits:
Object
  • Object
show all
Defined in:
lib/async/aws/client_cache.rb

Defined Under Namespace

Classes: Entry, ProxyClient

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize ⇒ void



55
56
57
58
59
# File 'lib/async/aws/client_cache.rb', line 55

def initialize
  @clients = {}
  @mutex = Mutex.new
  @access_count = 0
end

Class Method Details

.default_cert_store ⇒ Object



17
18
19
20
21
22
23
24
25
26
27
# File 'lib/async/aws/client_cache.rb', line 17

def self.default_cert_store
  return @default_cert_store if @default_cert_store

  @default_cert_store_mutex.synchronize do
    return @default_cert_store if @default_cert_store

    store = OpenSSL::X509::Store.new
    store.set_default_paths
    @default_cert_store = store
  end
end

Instance Method Details

#clear!(timeout: nil) ⇒ void

This method returns an undefined value.

Closes all cached clients and clears the cache. Intended for shutdown only.



185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
# File 'lib/async/aws/client_cache.rb', line 185

def clear!(timeout: nil)
  clients = @mutex.synchronize do
    values = @clients.values
    @clients.clear
    values
  end

  clients.each do |client_entry|
    current_reactor = Async::Task.current? ? Async::Task.current.reactor : nil
    owner_reactor = entry_reactor(client_entry)
    if timeout && current_reactor && owner_reactor == current_reactor
      begin
        client = extract_client(client_entry)
        Async::Task.current.with_timeout(timeout, Async::TimeoutError) { client.close if client.respond_to?(:close) }
      rescue Async::TimeoutError
        logger = logger_for
        logger&.warn('[aws-sdk-http-async] force-closing client (timeout)')
        close_entry(client_entry, force: true)
      rescue StandardError => e
        logger = logger_for
        logger&.warn("[aws-sdk-http-async] failed to close client: #{e.message}")
      end
    else
      close_entry(client_entry, force: true)
    end
  end
end

#client_for(endpoint, config) ⇒ Async::HTTP::Client

Parameters:

  • endpoint (URI::HTTP, URI::HTTPS)
  • config (Seahorse::Client::Configuration)

Returns:

  • (Async::HTTP::Client)

Raises:



64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
# File 'lib/async/aws/client_cache.rb', line 64

def client_for(endpoint, config)
  raise NoReactorError, 'Async reactor is required. Wrap calls in Sync { }.' unless Async::Task.current?

  reactor = Async::Task.current.reactor
  key = cache_key(endpoint, config, reactor)
  entry = nil
  stale_entry = nil

  entry = @mutex.synchronize do
    cached = @clients[key]
    if entry_valid_for?(cached, reactor)
      touch_lru!(key, cached)
      cached
    else
      stale_entry = @clients.delete(key) if cached
      nil
    end
  end

  close_entry(stale_entry) if stale_entry
  sweep_dead_entries_if_needed.each { |dead| close_entry(dead) }
  return entry.client if entry

  new_entry = nil
  evicted = []
  used_entry = nil
  stale_existing = nil

  new_entry = Entry.new(build_client(endpoint, config), WeakRef.new(reactor), 0)

  @mutex.synchronize do
    existing = @clients[key]
    if entry_valid_for?(existing, reactor)
      touch_lru!(key, existing)
      used_entry = existing
    else
      stale_existing = existing
      @clients[key] = new_entry
      touch_lru!(key, new_entry)
      used_entry = new_entry
      evicted = evict_entries_locked(config, reactor)
    end
  end

  close_entry(new_entry) if used_entry != new_entry
  close_entry(stale_existing) if stale_existing
  evicted.each { |entry_to_close| close_entry(entry_to_close) }

  used_entry.client
end

#close! ⇒ void

This method returns an undefined value.



214
215
216
# File 'lib/async/aws/client_cache.rb', line 214

def close!
  clear!
end

#with_client(endpoint, config) {|client| ... } ⇒ Object

Parameters:

  • endpoint (URI::HTTP, URI::HTTPS)
  • config (Seahorse::Client::Configuration)

Yield Parameters:

  • client (Async::HTTP::Client)

Returns:

  • (Object)

Raises:



119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
# File 'lib/async/aws/client_cache.rb', line 119

def with_client(endpoint, config)
  raise NoReactorError, 'Async reactor is required. Wrap calls in Sync { }.' unless Async::Task.current?

  reactor = Async::Task.current.reactor
  key = cache_key(endpoint, config, reactor)
  entry = nil
  stale_entry = nil

  entry = @mutex.synchronize do
    cached = @clients[key]
    if entry_valid_for?(cached, reactor)
      touch_lru!(key, cached)
      cached.inflight = cached.inflight.to_i + 1
      cached
    else
      stale_entry = @clients.delete(key) if cached
      nil
    end
  end

  close_entry(stale_entry) if stale_entry
  sweep_dead_entries_if_needed.each { |dead| close_entry(dead) }

  unless entry
    new_entry = Entry.new(build_client(endpoint, config), WeakRef.new(reactor), 0)
    evicted = []
    used_entry = nil
    stale_existing = nil

    @mutex.synchronize do
      existing = @clients[key]
      if entry_valid_for?(existing, reactor)
        touch_lru!(key, existing)
        existing.inflight = existing.inflight.to_i + 1
        used_entry = existing
      else
        stale_existing = existing
        @clients[key] = new_entry
        touch_lru!(key, new_entry)
        new_entry.inflight = new_entry.inflight.to_i + 1
        used_entry = new_entry
        evicted = evict_entries_locked(config, reactor)
      end
    end

    close_entry(new_entry) if used_entry != new_entry
    close_entry(stale_existing) if stale_existing
    evicted.each { |entry_to_close| close_entry(entry_to_close) }

    entry = used_entry
  end

  begin
    yield extract_client(entry)
  ensure
    @mutex.synchronize do
      if entry
        entry.inflight = [entry.inflight.to_i - 1, 0].max
      end
    end
  end
end