Class: Parse::Agent::MCPRackApp::CancellationRegistry
- Inherits:
-
Object
- Object
- Parse::Agent::MCPRackApp::CancellationRegistry
- Defined in:
- lib/parse/agent/mcp_rack_app.rb
Instance Method Summary collapse
-
#cancel(correlation_id, request_id, reason: :notifications_cancelled) ⇒ Boolean
Trip the matching token.
-
#cancel_all_for(correlation_id, reason: :session_terminated) ⇒ Object
Trip every token registered under the given correlation_id.
-
#deregister(correlation_id, request_id, entry_id) ⇒ Boolean
Release a previously-registered entry.
-
#initialize ⇒ CancellationRegistry
constructor
A new instance of CancellationRegistry.
-
#register(correlation_id, request_id, token) ⇒ String?
Register a cancellation token for the given session and request id pair.
-
#size ⇒ Integer
Number of currently-registered tokens.
Constructor Details
#initialize ⇒ CancellationRegistry
Returns a new instance of CancellationRegistry.
2016 2017 2018 2019 |
# File 'lib/parse/agent/mcp_rack_app.rb', line 2016 def initialize @entries = {} @mutex = Mutex.new end |
Instance Method Details
#cancel(correlation_id, request_id, reason: :notifications_cancelled) ⇒ Boolean
Trip the matching token. Silent no-op when the entry is missing — by design, to avoid a probe oracle.
2071 2072 2073 2074 2075 2076 2077 |
# File 'lib/parse/agent/mcp_rack_app.rb', line 2071 def cancel(correlation_id, request_id, reason: :notifications_cancelled) return false if correlation_id.nil? || correlation_id.to_s.empty? entry = @mutex.synchronize { @entries[[correlation_id, request_id]] } return false unless entry entry[1].cancel!(reason: reason) true end |
#cancel_all_for(correlation_id, reason: :session_terminated) ⇒ Object
Trip every token registered under the given correlation_id.
Used by DELETE / session termination — when a client tears
down its session, any in-flight requests still running under
that correlation_id are cancelled so worker threads exit
promptly instead of carrying a doomed result to completion.
Silent no-op when no entries match (or correlation_id is blank). Returns the number of tokens tripped.
2087 2088 2089 2090 2091 2092 2093 2094 2095 |
# File 'lib/parse/agent/mcp_rack_app.rb', line 2087 def cancel_all_for(correlation_id, reason: :session_terminated) return 0 if correlation_id.nil? || correlation_id.to_s.empty? tokens = @mutex.synchronize do keys = @entries.keys.select { |(cid, _)| cid == correlation_id } keys.map { |k| @entries.delete(k)[1] } end tokens.each { |t| t.cancel!(reason: reason) } tokens.size end |
#deregister(correlation_id, request_id, entry_id) ⇒ Boolean
Release a previously-registered entry. Removes the slot only when the current owner matches the passed entry-id, so a stale on_close from a request whose slot was overwritten by a sibling registration cannot evict the sibling's token. Idempotent.
2053 2054 2055 2056 2057 2058 2059 2060 2061 2062 2063 2064 2065 |
# File 'lib/parse/agent/mcp_rack_app.rb', line 2053 def deregister(correlation_id, request_id, entry_id) return false if correlation_id.nil? || correlation_id.to_s.empty? return false if entry_id.nil? @mutex.synchronize do current = @entries[[correlation_id, request_id]] if current && current[0] == entry_id @entries.delete([correlation_id, request_id]) true else false end end end |
#register(correlation_id, request_id, token) ⇒ String?
Register a cancellation token for the given session and request id pair. Returns an opaque entry-id that the caller must pass to #deregister to release the slot. If multiple registrations land on the same key (legitimate id-reuse by the same session, or a request retry), only the latest registration is reachable for #cancel; older entries can still be safely released via their entry-id even though they no longer "own" the slot.
2037 2038 2039 2040 2041 2042 2043 2044 |
# File 'lib/parse/agent/mcp_rack_app.rb', line 2037 def register(correlation_id, request_id, token) return nil if correlation_id.nil? || correlation_id.to_s.empty? entry_id = SecureRandom.uuid @mutex.synchronize do @entries[[correlation_id, request_id]] = [entry_id, token] end entry_id end |
#size ⇒ Integer
Returns number of currently-registered tokens. Used by tests and operator dashboards.
2099 2100 2101 |
# File 'lib/parse/agent/mcp_rack_app.rb', line 2099 def size @mutex.synchronize { @entries.size } end |