class_name ConnectionLeaseClient extends Node const AssignmentState = preload("res://scripts/assignment_state.gd") signal reconciliation_failed(reason: String) const CLAIMED := "claimed" const UNAVAILABLE := "unavailable" const REJECTED := "rejected" var _base_url := "" var _workload_token := "" var _match_id := "" var _server_id := "" var _pending: Array[Dictionary] = [] var _processing := false func configure(base_url: String, workload_token: String, match_id: String, server_id: String) -> bool: base_url = base_url.strip_edges().trim_suffix("/") workload_token = workload_token.strip_edges() if not valid_configuration(base_url, workload_token, match_id, server_id): return false _base_url = base_url _workload_token = workload_token _match_id = match_id _server_id = server_id return true # Admission awaits one bounded request only. If the control plane is down, the # same process may continue using its local generation and this exact event is # retained ahead of every later disconnect/reconnect for ordered reconciliation. func claim(player_id: String, expected_generation: int) -> Dictionary: if not AssignmentState.is_valid_opaque_id(player_id) or expected_generation < 0: return {"status": REJECTED} var event := _connect_event(player_id, expected_generation) if _processing or not _pending.is_empty(): if expected_generation == 0: return {"status": REJECTED} _pending.append(event) _start_processing() return {"status": UNAVAILABLE, "generation": expected_generation + 1} var response := await _send(event, true) if String(response.get("status", "")) == UNAVAILABLE: if expected_generation == 0: return {"status": REJECTED} _pending.append(event) _start_processing() return {"status": UNAVAILABLE, "generation": expected_generation + 1} return response func record_disconnect(player_id: String, generation: int) -> void: if not AssignmentState.is_valid_opaque_id(player_id) or generation < 1: return _pending.append(_disconnect_event(player_id, generation)) _start_processing() func _start_processing() -> void: if _processing or _pending.is_empty() or not is_inside_tree(): return _process_pending() func _process_pending() -> void: _processing = true while not _pending.is_empty() and is_inside_tree(): var event := _pending[0] var response := await _send(event) var status := String(response.get("status", "")) if status == CLAIMED: _pending.pop_front() continue if status == REJECTED: reconciliation_failed.emit("durable connection lease conflict") _processing = false return await get_tree().create_timer(1.0).timeout _processing = false func _send(event: Dictionary, allow_recovery := false) -> Dictionary: var request := HTTPRequest.new() request.timeout = 1.0 add_child(request) var operation := String(event["operation"]) var endpoint := "%s/v1/servers/%s/%s" % [_base_url, _server_id.uri_encode(), operation] var start_error := request.request(endpoint, [ "Authorization: Bearer " + _workload_token, "Content-Type: application/json", "Idempotency-Key: " + String(event["key"]), ], HTTPClient.METHOD_POST, JSON.stringify(event["payload"])) if start_error != OK: request.queue_free() return {"status": UNAVAILABLE} var raw: Array = await request.request_completed request.queue_free() return classify_response(operation, int(event["generation"]), int(raw[0]), int(raw[1]), raw[3], allow_recovery) func _connect_event(player_id: String, expected_generation: int) -> Dictionary: return { "operation": "connect", "generation": expected_generation, "key": event_key(_match_id, player_id, "connect", expected_generation), "payload": {"player_id": player_id, "expected_generation": expected_generation}, } func _disconnect_event(player_id: String, generation: int) -> Dictionary: return { "operation": "disconnect", "generation": generation, "key": event_key(_match_id, player_id, "disconnect", generation), "payload": {"player_id": player_id, "generation": generation}, } static func classify_response(operation: String, generation: int, request_result: int, response_code: int, body: PackedByteArray, allow_recovery := false) -> Dictionary: if request_result != HTTPRequest.RESULT_SUCCESS or response_code == 0 or response_code == 429 or response_code >= 500: return {"status": UNAVAILABLE} if operation == "disconnect" and response_code == 204: return {"status": CLAIMED, "generation": generation} if operation == "connect" and response_code == 200: var decoded = JSON.parse_string(body.get_string_from_utf8()) if decoded is Dictionary and _valid_generation(decoded.get("generation")): var claimed_generation := int(decoded["generation"]) if claimed_generation == generation + 1 or (allow_recovery and generation == 0 and claimed_generation > 1): return {"status": CLAIMED, "generation": claimed_generation} return {"status": REJECTED} static func _valid_generation(value: Variant) -> bool: if value is int: return int(value) >= 1 if value is float: return is_finite(float(value)) and float(value) >= 1.0 and float(value) == floor(float(value)) and float(value) <= 9007199254740991.0 return false static func event_key(match_id: String, player_id: String, operation: String, generation: int) -> String: return "server-lease-" + (match_id + "\n" + player_id + "\n" + operation + "\n" + str(generation)).sha256_text() static func valid_configuration(base_url: String, workload_token: String, match_id: String, server_id: String) -> bool: if not (base_url.begins_with("http://") or base_url.begins_with("https://")) or base_url.contains("\n") or base_url.contains("\r") or base_url.contains("?") or base_url.contains("#") or base_url.contains("@"): return false if workload_token.is_empty() or workload_token.contains("\n") or workload_token.contains("\r"): return false return AssignmentState.is_valid_opaque_id(match_id) and AssignmentState.is_valid_opaque_id(server_id)