Class: Wurk::Capsule

Inherits:
Object
  • Object
show all
Defined in:
lib/wurk/capsule.rb

Overview

One processing unit: a set of threads + queues sharing a fetcher and a Redis pool. Configurations can hold many capsules; each maps to its own Manager and Processors.

Spec: docs/target/sidekiq-free.md §5 (Sidekiq::Capsule).

Constant Summary collapse

MODES =
%i[strict weighted random].freeze
POOL_HEADROOM =

Headroom above concurrency for the main pool; the whole pool is then floored at MIN_POOL_SIZE. Blocking BLMOVE fetch has its own pool (#fetch_redis_pool), so the main pool serves only the background loops — heartbeat, scheduled poller, leader election, cron, the two metrics rollups, reaper, history, health probe — plus the host's own job-code checkouts. concurrency + 5 (floor 10) gives each an unstarvable slot; the old concurrency + 2 starved them once fetch also drew from here — the #101 0/N pool-exhaustion incident. Override via config.redis[:size].

5
MIN_POOL_SIZE =
10

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(name, config) ⇒ Capsule

Returns a new instance of Capsule.



25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
# File 'lib/wurk/capsule.rb', line 25

def initialize(name, config)
  @name = name.to_s
  @config = config
  @concurrency = config[:concurrency] || 5
  @queues = ['default']
  @mode = :strict
  @weights = { 'default' => 0 }
  @fetcher = nil
  # One mutable Hash rather than two ivars: the pools are the only part of a
  # capsule that legitimately changes after Configuration#freeze! — fork
  # closes them, Launcher#stop releases them, an embedded host that boots
  # again rebuilds them. `Object#freeze` is shallow, so the Hash stays
  # writable and freezing a capsule keeps meaning "no more configuration"
  # instead of "these sockets are yours forever". Before this, a reset on a
  # frozen capsule disconnected the pool and then raised FrozenError on the
  # memo, leaving `redis_pool` answering with a shut-down pool for good.
  @pools = {}
  @client_chain = nil
  @server_chain = nil
end

Instance Attribute Details

#concurrencyObject

Returns the value of attribute concurrency.



18
19
20
# File 'lib/wurk/capsule.rb', line 18

def concurrency
  @concurrency
end

#configObject (readonly)

Returns the value of attribute config.



17
18
19
# File 'lib/wurk/capsule.rb', line 17

def config
  @config
end

#fetcherObject

Returns the value of attribute fetcher.



18
19
20
# File 'lib/wurk/capsule.rb', line 18

def fetcher
  @fetcher
end

#modeObject (readonly)

Returns the value of attribute mode.



17
18
19
# File 'lib/wurk/capsule.rb', line 17

def mode
  @mode
end

#nameObject (readonly)

Returns the value of attribute name.



17
18
19
# File 'lib/wurk/capsule.rb', line 17

def name
  @name
end

#queuesObject

Returns the value of attribute queues.



17
18
19
# File 'lib/wurk/capsule.rb', line 17

def queues
  @queues
end

#weightsObject (readonly)

Returns the value of attribute weights.



17
18
19
# File 'lib/wurk/capsule.rb', line 17

def weights
  @weights
end

Instance Method Details

#client_middleware {|chain| ... } ⇒ Object

Yields:

  • (chain)


104
105
106
107
108
# File 'lib/wurk/capsule.rb', line 104

def client_middleware
  chain = (@client_chain ||= @config.client_middleware.copy_for(self))
  yield chain if block_given?
  chain
end

#fetch_poll_intervalObject

Empty-poll BLMOVE backoff for this capsule's reliable fetcher (Pro super_fetch §3.3). nil → the fetcher falls back to its TIMEOUT default.



168
169
170
# File 'lib/wurk/capsule.rb', line 168

def fetch_poll_interval
  @config[:fetch_poll_interval]
end

#fetch_redis(idempotent: false) ⇒ Object

Checkout from the dedicated fetch pool. Only the reliable fetcher's blocking BLMOVE uses this, so a parked fetch never holds a main-pool slot.



158
159
160
# File 'lib/wurk/capsule.rb', line 158

def fetch_redis(idempotent: false, &)
  PoolCheckout.with(fetch_redis_pool, idempotent, &)
end

#fetch_redis_poolObject

Dedicated pool for the reliable fetcher's blocking BLMOVE: one slot per processor thread (concurrency), since at most that many threads park in fetch at once. Keeping fetch off the main pool is what lets an idle worker hold zero main-pool connections again.



138
139
140
# File 'lib/wurk/capsule.rb', line 138

def fetch_redis_pool
  @pools[:fetch] ||= build_pool(size: @concurrency, name: "#{@name}-fetch")
end

#loggerObject



172
173
174
# File 'lib/wurk/capsule.rb', line 172

def logger
  @config.logger
end

#lookup(name) ⇒ Object



162
163
164
# File 'lib/wurk/capsule.rb', line 162

def lookup(name)
  @config.lookup(name)
end

#prepare!Object

Materialize everything that lazy-inits via ||= and default the fetcher, BEFORE Configuration#freeze! freezes the capsule — otherwise the first post-freeze access (a fetch tick, a middleware call) hits a nil fetcher or FrozenErrors building a pool. The swarm's ChildBoot used to do this by hand; centralizing it here covers the standalone CLI and embedded paths too (the bug behind a nil fetcher in exe/wurk). Idempotent.



78
79
80
81
82
83
84
# File 'lib/wurk/capsule.rb', line 78

def prepare!
  prepare_shared!
  @fetcher ||= build_fetcher
  redis_pool
  fetch_redis_pool
  self
end

#prepare_shared!Object

The half of prepare! a forking parent can run on every child's behalf: the chains are a pure function of this capsule's identity, not of the slot (queues + concurrency) a swarm child is assigned later, and copy_for opens nothing. Run before the fork, the entries are allocated once and inherited copy-on-write instead of rebuilt in every child.

The rest of prepare! deliberately stays post-fork: both pools are sized off the slot's concurrency and one built here would hand every child an inherited socket, and build_fetcher fires the host's config[:fetch_setup] hook, which is per-child — running it in the parent would let a custom fetcher snapshot the wrong queues, or leak whatever the hook opened across the fork. Idempotent.



98
99
100
101
102
# File 'lib/wurk/capsule.rb', line 98

def prepare_shared!
  client_middleware
  server_middleware
  self
end

#queue_specsObject

Lossless name[,weight] specs (unlike queues, which is the weight-expanded list). Lets Configuration#topology rebuild a slot that round-trips back through queues= without flattening weights.



68
69
70
# File 'lib/wurk/capsule.rb', line 68

def queue_specs
  @weights.map { |q, w| w.positive? ? "#{q},#{w}" : q }
end

#redis(idempotent: false) ⇒ Object



152
153
154
# File 'lib/wurk/capsule.rb', line 152

def redis(idempotent: false, &)
  PoolCheckout.with(redis_pool, idempotent, &)
end

#redis_poolObject



130
131
132
# File 'lib/wurk/capsule.rb', line 130

def redis_pool
  @pools[:main] ||= build_pool(size: main_pool_size, name: "#{@name}-main")
end

#reset_redis_pools!Object

Disconnect and drop cached pools. Called by Wurk::Swarm just before fork (parent side: close inherited sockets), just after fork (child side: rebuild lazily), and by Launcher#stop (release what this process held). Connection_pool#shutdown is terminal, so dropping the reference is required — redis_pool will rebuild.



147
148
149
150
# File 'lib/wurk/capsule.rb', line 147

def reset_redis_pools!
  @pools.each_value(&:disconnect!)
  @pools.clear
end

#server_middleware {|chain| ... } ⇒ Object

Yields:

  • (chain)


110
111
112
113
114
115
116
117
# File 'lib/wurk/capsule.rb', line 110

def server_middleware
  # copy_for(self) — not dup — binds the chain's `@config` to this capsule,
  # so middleware that reach for `redis_pool`/`redis`/`logger` resolve them
  # instead of hitting `nil` (a plain dup leaves @config nil).
  chain = (@server_chain ||= @config.server_middleware.copy_for(self))
  yield chain if block_given?
  chain
end

#stopObject

No-op in OSS. Reserved for Pro/Ent hooks that need to flush per-capsule state on shutdown. Manager#stop invokes this in ensure so the contract is honored regardless of how stop unwinds.

Spec: docs/target/sidekiq-free.md §5.



181
# File 'lib/wurk/capsule.rb', line 181

def stop; end

#thread_priorityObject

Capsule-hosted components (Manager, Processor, Fetcher) hand their capsule to Component as config, and safe_thread reads the priority off it. Sidekiq delegates the same accessor (capsule.rb:30).



23
# File 'lib/wurk/capsule.rb', line 23

def thread_priority = @config.thread_priority

#to_hObject



46
47
48
# File 'lib/wurk/capsule.rb', line 46

def to_h
  { concurrency: @concurrency, mode: @mode, weights: @weights }
end