diff --git a/addons/twitcher/lib/http/buffered_http_client.gd b/addons/twitcher/lib/http/buffered_http_client.gd index 5310b3f5..bdbed4bd 100644 --- a/addons/twitcher/lib/http/buffered_http_client.gd +++ b/addons/twitcher/lib/http/buffered_http_client.gd @@ -2,7 +2,9 @@ @tool extends Twitcher -## Http client that bufferes the requests and sends them sequentialy +## Http client that buffers the requests and sends at most [member max_parallel_requests] +## of them at the same time. Everything above that limit waits in a queue and is +## dispatched in the order it was requested as soon as a slot gets free. class_name BufferedHTTPClient @@ -17,7 +19,7 @@ signal request_done(response: ResponseData) class RequestData extends RefCounted: ## The client that the request belongs too var client: BufferedHTTPClient - ## The request node that is executing the request + ## The request node that is executing the request (null while the request waits in the queue) var http_request: HTTPRequest ## Path of the request var path: String @@ -29,10 +31,14 @@ class RequestData extends RefCounted: var body: String = "" ## Amount of retries var retry: int + ## `Time.get_ticks_msec()` when the request was put on the wire (first attempt or retry) + var started_at: int ## When you are done free the request func queue_free() -> void: - http_request.queue_free() + if http_request != null: + http_request.queue_free() + http_request = null ## Contains the response data @@ -56,19 +62,26 @@ class ResponseData extends RefCounted: ## When a request fails max_error_count then cancel that request -1 for endless amount of tries. @export var max_error_count : int = -1 +## How many requests may be on the wire at the same time. [code]1[/code] sends them strictly +## one after another, which is the safe default for API calls (see PR #130: several threaded +## requests started in the same frame can stall and time out). Raise it for clients that +## download many independent small files like emotes and badges. [code]0[/code] or a negative +## value removes the limit. +@export var max_parallel_requests : int = 1 @export var custom_header : Dictionary[String, String] = { "Accept": "*/*" } +## Seconds until a single attempt of a request is aborted with [constant HTTPRequest.RESULT_TIMEOUT] +@export var request_timeout : float = 30 +## Every request that was started and whose response wasn't consumed via `wait_for_request` yet. var requests : Array[RequestData] = [] -var current_request : RequestData -var current_response_data : PackedByteArray = PackedByteArray() +## Requests that wait for a free slot, in the order they were requested. +var queued_requests : Array[RequestData] = [] +## Requests that are currently on the wire (including retries). +var active_requests : Array[RequestData] = [] var responses : Dictionary = {} -var error_count : int - -## Only one poll at a time so block for all other tries to call it -var polling: bool var processing: bool: - get: return not requests.is_empty() || current_request != null + get: return not requests.is_empty() ## Starts a request that will be handled as soon as the client gets free. @@ -83,16 +96,11 @@ func request(path: String, method: int, headers: Dictionary, body: String) -> Re req.body = body req.headers = headers req.client = self - req.http_request = HTTPRequest.new() - req.http_request.use_threads = true - req.http_request.timeout = 30 - req.http_request.request_completed.connect(_on_request_completed.bind(req)) - add_child(req.http_request) - var err : Error = req.http_request.request(req.path, _pack_headers(req.headers), req.method, req.body) - if err != OK: logError("Problems with request to %s cause of %s" % [path, error_string(err)]) requests.append(req) + queued_requests.append(req) request_added.emit(req) - logDebug("[%s] request started " % [ path ]) + logDebug("[%s] request queued (queued: %s, active: %s)" % [ path, queued_requests.size(), active_requests.size() ]) + _dispatch() return req @@ -116,6 +124,39 @@ func wait_for_request(request_data: RequestData) -> ResponseData: return latest_response +## Sends queued requests as long as there are free slots +func _dispatch() -> void: + while not queued_requests.is_empty() and _has_free_slot(): + var req: RequestData = queued_requests.pop_front() + active_requests.append(req) + _send(req) + + +func _has_free_slot() -> bool: + return max_parallel_requests <= 0 or active_requests.size() < max_parallel_requests + + +## Puts a request on the wire. Used for the first attempt and for every retry. +func _send(request_data: RequestData) -> void: + if request_data.http_request != null: + request_data.http_request.queue_free() + var http_request: HTTPRequest = HTTPRequest.new() + http_request.use_threads = true + http_request.timeout = request_timeout + http_request.request_completed.connect(_on_request_completed.bind(request_data)) + add_child(http_request) + request_data.http_request = http_request + request_data.started_at = Time.get_ticks_msec() + var err : Error = http_request.request(request_data.path, _pack_headers(request_data.headers), request_data.method, request_data.body) + if err != OK: + logError("Problems with request to %s cause of %s" % [request_data.path, error_string(err)]) + # HTTPRequest doesn't emit request_completed when request() fails, finish it + # ourself otherwise the request blocks its slot forever and waiters hang. + _on_request_completed.call_deferred(HTTPRequest.Result.RESULT_REQUEST_FAILED, 0, PackedStringArray(), PackedByteArray(), request_data) + return + logDebug("[%s] request started " % [ request_data.path ]) + + func _on_request_completed(result: int, response_code: int, headers: PackedStringArray, body: PackedByteArray, request_data: RequestData) -> void: var response_data : ResponseData = ResponseData.new() if result != HTTPRequest.Result.RESULT_SUCCESS: @@ -124,17 +165,17 @@ func _on_request_completed(result: int, response_code: int, headers: PackedStrin if result == HTTPRequest.Result.RESULT_CONNECTION_ERROR || result == HTTPRequest.Result.RESULT_TLS_HANDSHAKE_ERROR: if request_data.retry == max_error_count: printerr("Maximum amount of retries for the request. Abort request: %s" % [request_data.path]) + # Fall through and deliver the failed response, so that waiters get an + # answer and the slot is released for the next queued request. + else: + var wait_time = pow(2, request_data.retry) + wait_time = min(wait_time, 30) + logDebug("Error happend during connection. Wait for %s" % wait_time) + await get_tree().create_timer(wait_time, true, false, true).timeout + request_data.retry += 1 + # The request keeps its slot while retrying + _send(request_data) return - var wait_time = pow(2, request_data.retry) - wait_time = min(wait_time, 30) - logDebug("Error happend during connection. Wait for %s" % wait_time) - await get_tree().create_timer(wait_time, true, false, true).timeout - var http_request: HTTPRequest = request_data.http_request.duplicate() - add_child(http_request) - request_data.http_request = http_request - request_data.retry += 1 - http_request.request(request_data.path, _pack_headers(request_data.headers), request_data.method, request_data.body) - http_request.request_completed.connect(_on_request_completed.bind(http_request)) response_data.result = result response_data.request_data = request_data @@ -142,7 +183,9 @@ func _on_request_completed(result: int, response_code: int, headers: PackedStrin response_data.response_code = response_code response_data.response_header = _get_response_headers_as_dictionary(headers) responses[request_data] = response_data - logInfo("[%s] request done with result HTTPRequest.Result[%s] " % [ request_data.path, result]) + logInfo("[%s] request done with result HTTPRequest.Result[%s] code %s after %sms (retries: %s)" % [ request_data.path, result, response_code, Time.get_ticks_msec() - request_data.started_at, request_data.retry ]) + active_requests.erase(request_data) + _dispatch() request_done.emit(response_data) @@ -170,12 +213,9 @@ func _pack_headers(headers: Dictionary) -> PackedStringArray: return result -## The amount of requests that are pending +## The amount of requests that are pending (waiting in the queue or on the wire) func queued_request_size() -> int: - var requests_size: int = requests.size() - if current_request != null: - requests_size += 1 - return requests_size + return queued_requests.size() + active_requests.size() func empty_response(request_data: RequestData) -> ResponseData: diff --git a/addons/twitcher/media/twitch_media_loader.gd b/addons/twitcher/media/twitch_media_loader.gd index 56450558..9520a0f8 100644 --- a/addons/twitcher/media/twitch_media_loader.gd +++ b/addons/twitcher/media/twitch_media_loader.gd @@ -23,6 +23,13 @@ const FALLBACK_PROFILE = preload("res://addons/twitcher/assets/no_profile.png") @export var fallback_texture: Texture2D = FALLBACK_TEXTURE @export var fallback_profile: Texture2D = FALLBACK_PROFILE @export var image_cdn_host: String = "https://static-cdn.jtvnw.net" +## How many images (emotes, badges, cheermotes, profiles) are downloaded at the same time. +## Unlike API calls, image downloads are many small independent files and get noticeably +## slow when they are fetched one after another. See [member BufferedHTTPClient.max_parallel_requests]. +@export var max_parallel_downloads: int = 8: + set(val): + max_parallel_downloads = val + if _client != null: _client.max_parallel_requests = val ## Will preload the whole badge and emote cache also to editor time (use it when you make a Editor Plugin with Twitch Support) @export var load_cache_in_editor: bool @@ -51,6 +58,7 @@ var _client: BufferedHTTPClient func _ready() -> void: _client = BufferedHTTPClient.new() _client.name = "TwitchMediaLoaderClient" + _client.max_parallel_requests = max_parallel_downloads add_child(_client) _load_cache() if api == null: api = TwitchAPI.instance diff --git a/test/unit/lib/test_buffered_http_client.gd b/test/unit/lib/test_buffered_http_client.gd new file mode 100644 index 00000000..0c73c395 --- /dev/null +++ b/test/unit/lib/test_buffered_http_client.gd @@ -0,0 +1,223 @@ +## Unit tests for [BufferedHTTPClient]. +## +## The client is exercised against a real [HTTPServer] on the loopback interface, +## because Godot doesn't allow overriding native methods and therefore +## [HTTPRequest] can't be doubled in a meaningful way. The server never answers +## on its own; each test decides when connections get a response (or get +## dropped), which is what makes the queue, the parallel limit and the retry path +## observable. +## +## Background: PR #130 — several threaded requests started in the same frame +## stalled and timed out, so API clients must be sequential by default, while +## emote and badge downloads still need parallelism to stay fast. +extends TwitcherTest + +const HttpServer := preload("res://addons/twitcher/lib/http/http_server.gd") +const Subject := preload("res://addons/twitcher/lib/http/buffered_http_client.gd") + +## Each test listens on its own port so a socket lingering from the previous +## test can never make the next one flaky. +static var _next_port := 48131 + +var _port: int +var _server: HTTPServer +## Connections that sent a request and now wait for an answer. +var _open: Array[HTTPServer.Client] = [] +## Connections we already counted, so a request arriving in two TCP chunks is +## counted once. +var _seen: Array[HTTPServer.Client] = [] +## Request paths in the order they hit the server. +var _paths: PackedStringArray = [] +## How many of the upcoming requests get their connection dropped without an answer. +var _drop_next := 0 + + +func before_each() -> void: + super() + _open = [] + _seen = [] + _paths = [] + _drop_next = 0 + _port = _next_port + _next_port += 1 + _server = HttpServer.create(_port, "127.0.0.1") + add_child_autofree(_server) + _server.request_received.connect(_on_request) + _server.start_listening() + + +func after_each() -> void: + _server.stop_listening() + super() + + +#region Helpers + +func _url(path: String) -> String: + return "http://127.0.0.1:%d/%s" % [_port, path] + + +func _make_client(max_parallel: int, max_errors: int = -1) -> BufferedHTTPClient: + var client: BufferedHTTPClient = Subject.new() + client.max_parallel_requests = max_parallel + client.max_error_count = max_errors + add_child_autofree(client) + return client + + +func _on_request(client: HTTPServer.Client) -> void: + var peer := client.peer + var data: Array = peer.get_data(peer.get_available_bytes()) + if _seen.has(client): + return + _seen.append(client) + var request_line: String = (data[1] as PackedByteArray).get_string_from_utf8().get_slice("\r\n", 0) + _paths.append(request_line.get_slice(" ", 1).trim_prefix("/")) + if _drop_next > 0: + _drop_next -= 1 + peer.disconnect_from_host() + return + _open.append(client) + + +func _received() -> int: + return _seen.size() + + +## Answers every connection that is waiting for a response with 200 OK. +func _respond_all() -> void: + for client: HTTPServer.Client in _open: + _server.send_response(client, "200 OK", "ok".to_utf8_buffer()) + _open = [] + + +## Waits until the server has seen [param count] requests, or gives up after [param timeout] seconds. +func _wait_until_received(count: int, timeout := 3.0) -> void: + var deadline := Time.get_ticks_msec() + int(timeout * 1000) + while _received() < count and Time.get_ticks_msec() < deadline: + await get_tree().process_frame + + +## Waits until [param request] has a response, or gives up after [param timeout] seconds. +func _wait_for_response(client: BufferedHTTPClient, request: BufferedHTTPClient.RequestData, timeout := 5.0) -> BufferedHTTPClient.ResponseData: + var deadline := Time.get_ticks_msec() + int(timeout * 1000) + while not client.responses.has(request) and Time.get_ticks_msec() < deadline: + await get_tree().process_frame + if not client.responses.has(request): + fail_test("No response for %s within %ss" % [request.path, timeout]) + return null + return await client.wait_for_request(request) + +#endregion + + +func test_default_limit_is_sequential() -> void: + var client: BufferedHTTPClient = Subject.new() + autofree(client) + assert_eq(client.max_parallel_requests, 1, "API clients must not fire requests in parallel by default") + + +func test_sequential_client_puts_only_one_request_on_the_wire() -> void: + var client := _make_client(1) + var a := client.request(_url("a"), HTTPClient.METHOD_GET, {}, "") + var b := client.request(_url("b"), HTTPClient.METHOD_GET, {}, "") + var c := client.request(_url("c"), HTTPClient.METHOD_GET, {}, "") + assert_eq(client.queued_request_size(), 3, "all three requests are pending") + + await _wait_until_received(1) + await wait_seconds(0.3) + assert_eq(_received(), 1, "only one request may be on the wire") + assert_eq(client.active_requests.size(), 1) + assert_eq(client.queued_requests.size(), 2) + + _respond_all() + await _wait_until_received(2) + await wait_seconds(0.3) + assert_eq(_received(), 2, "the second request starts once the first one is answered") + + _respond_all() + await _wait_until_received(3) + _respond_all() + + assert_eq(_paths, PackedStringArray(["a", "b", "c"]), "requests are sent in the order they were queued") + for request: BufferedHTTPClient.RequestData in [a, b, c]: + var response := await _wait_for_response(client, request) + if response != null: + assert_eq(response.response_code, 200, "%s answered" % request.path) + assert_eq(client.queued_request_size(), 0, "nothing pending after all responses were consumed") + assert_false(client.processing) + + +func test_parallel_limit_allows_that_many_requests_at_once() -> void: + var client := _make_client(2) + client.request(_url("a"), HTTPClient.METHOD_GET, {}, "") + client.request(_url("b"), HTTPClient.METHOD_GET, {}, "") + client.request(_url("c"), HTTPClient.METHOD_GET, {}, "") + + await _wait_until_received(2) + await wait_seconds(0.3) + assert_eq(_received(), 2, "two requests may be on the wire") + assert_eq(client.queued_requests.size(), 1) + + _respond_all() + await _wait_until_received(3) + assert_eq(_received(), 3, "the third request follows as soon as a slot is free") + _respond_all() + await wait_for_signal(client.request_done, 5) + + +func test_zero_limit_means_unlimited() -> void: + var client := _make_client(0) + for i: int in 5: + client.request(_url(str(i)), HTTPClient.METHOD_GET, {}, "") + + await _wait_until_received(5) + assert_eq(_received(), 5, "no limit: everything goes out immediately") + assert_eq(client.queued_requests.size(), 0) + _respond_all() + await wait_for_signal(client.request_done, 5) + + +## Before PR #130 an exhausted retry budget just returned, leaving `wait_for_request` +## hanging forever and (now) the slot occupied for good. +func test_exhausted_retries_deliver_an_error_response_and_free_the_slot() -> void: + var client := _make_client(1, 0) + _drop_next = 1 + var a := client.request(_url("a"), HTTPClient.METHOD_GET, {}, "") + var b := client.request(_url("b"), HTTPClient.METHOD_GET, {}, "") + + var response_a := await _wait_for_response(client, a) + if response_a == null: + return + assert_true(response_a.error, "dropped connection is reported as error") + assert_eq(response_a.result, HTTPRequest.RESULT_CONNECTION_ERROR) + assert_eq(a.retry, 0, "max_error_count 0 means no retry at all") + + await _wait_until_received(2) + assert_eq(_received(), 2, "the failed request released its slot for the next one") + _respond_all() + var response_b := await _wait_for_response(client, b) + if response_b != null: + assert_eq(response_b.response_code, 200) + + +## The retry used to reconnect `request_completed` bound to the new HTTPRequest node +## instead of the RequestData, so a retried request could never be completed. +func test_retry_completes_the_original_request() -> void: + var client := _make_client(1, 1) + _drop_next = 1 + var a := client.request(_url("a"), HTTPClient.METHOD_GET, {}, "") + + await _wait_until_received(2, 6.0) # first attempt fails, retry follows after a 1s backoff + assert_eq(_received(), 2, "the request was sent a second time") + _respond_all() + + var response := await _wait_for_response(client, a) + if response == null: + return + assert_eq(response.response_code, 200, "the retry answers the original request") + assert_false(response.error) + assert_eq(a.retry, 1) + assert_eq(client.active_requests.size(), 0, "the retried request released its slot") + await get_tree().process_frame # queue_free is deferred + assert_eq(client.get_child_count(), 0, "the HTTPRequest nodes of the failed attempt and of the retry were both freed") diff --git a/test/unit/lib/test_buffered_http_client.gd.uid b/test/unit/lib/test_buffered_http_client.gd.uid new file mode 100644 index 00000000..19909409 --- /dev/null +++ b/test/unit/lib/test_buffered_http_client.gd.uid @@ -0,0 +1 @@ +uid://unpuv8w0w1c4