Class: X::StreamingClient
- Inherits:
-
Object
- Object
- X::StreamingClient
- Defined in:
- x-streaming/lib/x/streaming/streaming_client.rb
Overview
A client for the streaming endpoints, which hold a connection open rather than answer a request
A stream reads until it is interrupted, so it reads with a short timeout and reconnects when it drops. The streaming endpoints require app-only authentication, so it streams with the bearer token of a client that has one, or fetches one with the API key and secret of a client that signs with OAuth 1.0a. It takes its credentials, base URL, parsing classes, and on_response hook from the client it was built from, and keeps the rest of its settings itself.
Each stream opens a connection of its own and closes it when it ends, rather than keep one open for the request that follows, as a client does between requests. So a streaming client has neither the keep_alive_timeout of a client, which says how long a connection is kept open, nor its close, which closes the connections it kept: a streaming client keeps none between streams. A stream is stopped by breaking or raising from the block that reads it, or from another thread, or the trap of a signal, with #stop, which stops a stream that delivers nothing as well. A streaming client that was stopped stays stopped, and Client#streaming builds a new one each time it is called, so a streaming client to be stopped is kept in a variable, rather than built again to stop.
A stream is opened with X::Client#get_stream, of an app-only copy of the client whose read_timeout is the stream's, so it carries the credentials, headers, and proxy of the client as any request does, and none of its credentials to another origin than the base URL of the client.
A streaming client keeps the settings it was built with for as long as it lives, as a client does, so a stream that runs for hours never reads a setting another thread is halfway through changing. Client#streaming builds one whose read_timeout, max_reconnects, or on_reconnect differ, the settings a streaming client keeps itself, and it takes the rest of its settings from the client it is built from, which #client reads, so a stream that connects differently is opened from a copy of that client: client.with(open_timeout: 2).streaming.
Constant Summary collapse
- DEFAULT_READ_TIMEOUT =
Default timeout for reading from a stream in seconds, half again the 20-second interval of the keep-alive X sends
30- DEFAULT_MAX_RECONNECTS =
Default maximum number of times in a row to reconnect a stream that drops without delivering an object or a keep-alive
ReconnectHandler::DEFAULT_MAX_RECONNECTS
Instance Attribute Summary collapse
-
#client ⇒ Client
readonly
The client the stream authenticates and parses with.
-
#on_reconnect ⇒ #call?
readonly
The callable passed the error that dropped a stream and its wait to reconnect.
Instance Method Summary collapse
-
#add_rules(rules, dry_run: false) {|problem| ... } ⇒ Array<StreamRule>
Add rules for the filtered stream to match posts against.
-
#as_json ⇒ void
Refuse to be read as JSON, which would write its credentials.
-
#delete_rules(rules, dry_run: false) {|problem| ... } ⇒ Integer
Delete rules of the filtered stream.
-
#encode_with(_coder) ⇒ void
Refuse to be written as YAML, which would write its credentials.
-
#initialize(client, read_timeout: DEFAULT_READ_TIMEOUT, max_reconnects: DEFAULT_MAX_RECONNECTS, on_reconnect: nil) ⇒ StreamingClient
constructor
Initialize a client for the streaming endpoints.
-
#inspect ⇒ String
Summarize the streaming client for the console without revealing credentials.
-
#marshal_dump ⇒ void
Refuse to be written with Marshal, which would write its credentials.
-
#max_reconnects ⇒ Integer, Float
The maximum number of times in a row to reconnect a stream.
-
#read_timeout ⇒ Integer, ...
The timeout for reading from a stream, in seconds.
-
#rules(params: nil) ⇒ Array<StreamRule>
The rules the filtered stream matches posts against.
-
#stop ⇒ nil
Stop every stream this streaming client runs, now and from then on.
-
#stopped? ⇒ Boolean
Whether #stop was called, after which each stream returns nil at once.
-
#stream(endpoint, params: nil, headers: {}, array_class: client.default_array_class, object_class: client.default_object_class) {|Hash, Array| ... } ⇒ Object?
Stream data from the X API.
-
#to_json(_state = nil) ⇒ void
Refuse to be written as JSON, which would write its credentials.
Constructor Details
#initialize(client, read_timeout: DEFAULT_READ_TIMEOUT, max_reconnects: DEFAULT_MAX_RECONNECTS, on_reconnect: nil) ⇒ StreamingClient
Initialize a client for the streaming endpoints
118 119 120 121 122 123 124 125 126 |
# File 'x-streaming/lib/x/streaming/streaming_client.rb', line 118 def initialize(client, read_timeout: DEFAULT_READ_TIMEOUT, max_reconnects: DEFAULT_MAX_RECONNECTS, on_reconnect: nil) @client = client @on_response = Stopper.guarding(client.on_response) @stream_client = client.with(read_timeout: Validator.read_timeout!(:read_timeout, read_timeout), on_response: CallbackError.tagging(@on_response)) @on_reconnect = Validator.callable!(:on_reconnect, on_reconnect) @reconnect_handler = ReconnectHandler.new(max_reconnects:, max_rate_limit_wait: client.max_rate_limit_wait, on_reconnect: Stopper.guarding(@on_reconnect)) @stream_parser = StreamParser.new @stopper = Stopper.new end |
Instance Attribute Details
#client ⇒ Client (readonly)
The client the stream authenticates and parses with
72 73 74 |
# File 'x-streaming/lib/x/streaming/streaming_client.rb', line 72 def client @client end |
#on_reconnect ⇒ #call? (readonly)
The callable passed the error that dropped a stream and its wait to reconnect
A stream that drops reconnects up to max_reconnects times in a row, without end by default, so one that can never connect, as with a host that does not resolve, or a proxy that refuses it, would reconnect in silence for as long as it runs; only a certificate that does not verify raises at once. on_reconnect is passed the error that dropped the stream, such as an X::NetworkError or an X::ServiceUnavailable, and the seconds the stream waits before it reconnects, before each wait, so it can report the reconnects, or give up on them by calling #stop, after which the stream returns nil rather than reconnect. An error it raises stops the stream, and reaches the caller as it was raised, as an error of the block of the stream does. A stop does not cut it short, as it does not the block of the stream.
It is called with these two arguments and no others, the error, which is never nil, since a stream reconnects only after one, and the wait, and every release of 1.x calls it so, so a lambda that takes exactly two, as ->(error, wait) does, serves each of them. What a later release of 1.x tells of a reconnect beside them, it passes to a hook of its own rather than to this one.
94 95 96 |
# File 'x-streaming/lib/x/streaming/streaming_client.rb', line 94 def on_reconnect @on_reconnect end |
Instance Method Details
#add_rules(rules, dry_run: false) {|problem| ... } ⇒ Array<StreamRule>
Add rules for the filtered stream to match posts against
A rule is a StreamRule, or a Hash of the value it matches and the tag it is labelled with, or a String, which is the value of a rule without a tag. A StreamRule is added by its value and tag, so the rules of one app, as #rules returns them, add themselves to another. No rules add none, and send no request. Anything else raises before a request, as it does for delete_rules.
The API adds the rules it can and reports the rest, such as a rule the app already has, as errors of a response that otherwise succeeds. The rules that were added are returned, and each rule that was not is yielded as the Problem the API reported, as a finder of x-objects yields the problems of a lookup. Without a block, a rule that was not added raises RulesRejected, which holds the rules that were as added, so that neither a rule the app already has nor one a dry run found invalid is passed over in silence.
312 313 314 315 316 317 318 |
# File 'x-streaming/lib/x/streaming/streaming_client.rb', line 312 def add_rules(rules, dry_run: false, &) rules = StreamRules.each_rule(rules) body = change_rules({add: rules.map { |rule| StreamRules.rule_to_add(rule) }}, dry_run:) unless rules.empty? added = StreamRules.rules_of(body) reporting(body, {added:}, &) added end |
#as_json ⇒ void
This method returns an undefined value.
Refuse to be read as JSON, which would write its credentials
ActiveSupport's Object#as_json reads every instance variable of an object that does not say how it is read, the client and its credentials among them, so it is refused as YAML is.
398 |
# File 'x-streaming/lib/x/streaming/streaming_client.rb', line 398 def as_json(*) = raise(TypeError, format(REFUSAL_MESSAGE, self.class, "JSON")) |
#delete_rules(rules, dry_run: false) {|problem| ... } ⇒ Integer
Delete rules of the filtered stream
A rule is deleted by its identifier, or by the value it matches: a StreamRule the API returned, a Hash that holds an id, an Integer, or an X::MatchingRule is deleted by identifier, so what #rules returned deletes itself, as do the rules a post of x-objects matched, which its matching_rules holds, and a StreamRule or a Hash that holds a value and no identifier, or a String, is deleted by value, so what add_rules was given deletes what it added. Anything else raises before a request, even something else with an id, such as a post, whose identifier would delete whichever rule shared it. No rules delete none, and send no request, since the API refuses a deletion that names no rule.
The API deletes the rules it can and reports the rest, such as a rule the app does not have, as errors of a response that otherwise succeeds. The number of rules that were deleted is returned, and each problem the API reported is yielded, as add_rules yields the rules it did not add, or, without a block, raises RulesRejected, which holds the number as deleted_count.
355 356 357 358 359 360 361 362 363 364 |
# File 'x-streaming/lib/x/streaming/streaming_client.rb', line 355 def delete_rules(rules, dry_run: false, &) rules = StreamRules.each_rule(rules) return 0 if rules.empty? ids, values = rules.partition { |rule| StreamRules.identifier_of(rule) } body = change_rules({delete: StreamRules.deletion(ids, values)}, dry_run:) deleted_count = body.to_h.dig("meta", "summary", "deleted").to_i reporting(body, {deleted_count:}, &) deleted_count end |
#encode_with(_coder) ⇒ void
This method returns an undefined value.
Refuse to be written as YAML, which would write its credentials
YAML reads no marshal_dump, and writes every instance variable of an object that does not say how it is written, the client and its credentials among them, so it is refused as Marshal is.
386 |
# File 'x-streaming/lib/x/streaming/streaming_client.rb', line 386 def encode_with(_coder) = raise(TypeError, format(REFUSAL_MESSAGE, self.class, "YAML")) |
#inspect ⇒ String
Summarize the streaming client for the console without revealing credentials
152 |
# File 'x-streaming/lib/x/streaming/streaming_client.rb', line 152 def inspect = "#<#{self.class} client=#{client.inspect}>" |
#marshal_dump ⇒ void
This method returns an undefined value.
Refuse to be written with Marshal, which would write its credentials
373 |
# File 'x-streaming/lib/x/streaming/streaming_client.rb', line 373 def marshal_dump = raise(TypeError, format(REFUSAL_MESSAGE, self.class, "Marshal")) |
#max_reconnects ⇒ Integer, Float
The maximum number of times in a row to reconnect a stream
A stream is reconnected when it drops, and the count starts over each time it delivers an object or reads the keep-alive X sends every 20 seconds, so that a stream that is quiet but connected never runs out of reconnects.
144 |
# File 'x-streaming/lib/x/streaming/streaming_client.rb', line 144 def max_reconnects = @reconnect_handler.max_reconnects |
#read_timeout ⇒ Integer, ...
The timeout for reading from a stream, in seconds
133 |
# File 'x-streaming/lib/x/streaming/streaming_client.rb', line 133 def read_timeout = @stream_client.read_timeout |
#rules(params: nil) ⇒ Array<StreamRule>
The rules the filtered stream matches posts against
The rules of an app are read, added, and deleted through a streaming client because they belong to the stream: they are what the filtered stream delivers, and they take the app-only authentication a stream takes. The API returns the rules a page at a time, and every page is read, so the rules are all of them. A page that names an empty next token, or the token of a page already read, is the last, as a page that names none is. The pages are read with while rather than Kernel#loop, which rescues StopIteration, so that a StopIteration the on_response of the client raises while a page is read reaches the caller, rather than end the reading in silence with the rules of the pages before it.
273 274 275 276 277 278 279 280 281 |
# File 'x-streaming/lib/x/streaming/streaming_client.rb', line 273 def rules(params: nil) rules = [] #: Array[StreamRule] spent = [] #: Array[String] while (token = read_rules(params, into: rules, spent:)) spent << token params = params.to_h.merge(pagination_token: token) end rules.freeze end |
#stop ⇒ nil
Stop every stream this streaming client runs, now and from then on
A stream waits on the API for most of its life, for the next object, the keep-alive X sends every 20 seconds, or the next reconnect, so its block, which runs only when an object arrives, cannot stop a stream that delivers nothing. This stops each stream running in another thread, or in this one, the next time it waits on the API, at once for a stream waiting now, closing its connection, and the stream returns nil. A block, or the on_response of the client, that is running when a stream is stopped runs to its end first, so that what it does with an object, or on_response with a failed response, is never cut short.
A streaming client that was stopped stays stopped: a stream it is asked to run later, as one a thread started just before stop may not yet have opened, returns nil at once, without a request. Client#streaming builds a new streaming client to stream with again, and builds another each time it is called, so stop is called on the streaming client the stream runs on, kept in a variable, not on one a second call builds. It may be called from the trap of a signal, and waits for no lock, so it may return before the streams it stops have ended.
245 |
# File 'x-streaming/lib/x/streaming/streaming_client.rb', line 245 def stop = @stopper.stop |
#stopped? ⇒ Boolean
Whether #stop was called, after which each stream returns nil at once
253 |
# File 'x-streaming/lib/x/streaming/streaming_client.rb', line 253 def stopped? = @stopper.stopped? |
#stream(endpoint, params: nil, headers: {}, array_class: client.default_array_class, object_class: client.default_object_class) {|Hash, Array| ... } ⇒ Object?
Stream data from the X API
The stream endpoints take app-only authentication, so a client that authenticates as a user streams with the bearer token its app_only client holds. A client that authenticates with OAuth 2.0 as a user and holds neither the app's bearer token nor its API key and secret has no app_only client, so it opens the stream as the user, which X refuses with 403 Forbidden, raising Forbidden, as any request X refuses the credentials of a client for does. A bearer token X rejects with 401 Unauthorized, as it does one that was invalidated, is fetched again with the API key and secret the app_only client holds, as it is for a request, and the stream is opened once more with it. A stream that drops, or that X disconnects with an operational-disconnect, reconnects, backing off as X recommends, up to max_reconnects times in a row. The API bills each object a stream delivers, so the client's on_response receives each one, as well as a failed response. An error on_response or the object_class raises stops the stream, and reaches the caller as it was raised, even an error of the X API, such as the X::ServiceUnavailable of a request on_response made, which a stream that dropped reconnects after.
Reconnects are unlimited by default, and made in silence but for on_reconnect, so a stream that cannot connect, as for a host that does not resolve or a network that is down, reconnects every 16 seconds, once its backoff has grown that long, for as long as it runs. Set max_reconnects to give up after that many in a row, when the stream raises the error of the last, or pass an on_reconnect that calls #stop, after which the stream returns nil. Only a certificate that does not verify, which will not verify the next time either, raises at once, as an X::NetworkError whose cause is the OpenSSL::SSL::SSLError.
A stream runs until its block stops it: break out of the block to stop the stream and return a value, throw to unwind to a catch further out, or raise, which stops the stream even where a drop would have reconnected, and reaches the caller unchanged, a StopIteration included. Another thread stops it with #stop, which ends a stream that delivers nothing as well, when it returns nil, as a stream of a streaming client that was stopped does at once, without a request.
205 206 207 208 209 210 211 212 213 214 215 216 |
# File 'x-streaming/lib/x/streaming/streaming_client.rb', line 205 def stream(endpoint, params: nil, headers: {}, array_class: client.default_array_class, object_class: client.default_object_class, &block) raise ArgumentError, NO_BLOCK_MESSAGE if block.nil? Validator.parsing_classes!(array_class:, object_class:) consumer = ->(object) { Stopper.guard { block.call(object) } } @stopper.run do @reconnect_handler.handle(consumer) do |deliver, alive| app_only(@stream_client).get_stream(endpoint, params:, headers:) { |response| read(response, array_class:, object_class:, alive:, &deliver) } end end end |
#to_json(_state = nil) ⇒ void
This method returns an undefined value.
Refuse to be written as JSON, which would write its credentials
It raises as #as_json does, for the reason that says, so that JSON.generate refuses a streaming client within what it writes as well.
411 |
# File 'x-streaming/lib/x/streaming/streaming_client.rb', line 411 def to_json(_state = nil) = raise(TypeError, format(REFUSAL_MESSAGE, self.class, "JSON")) |