/
githubmirror
/
XX-Net
Обзор
Документация
Войти
/
githubmirror
/
XX-Net
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
code/default/lib/noarch/front_base/http1.py
306 строк
11 KB
micheal
5.16.4 improve performance
25 авг 2025, 02:10
25 авг 2025, 02:10
758928d
Код
Авторство
О чём код?
from queue import Queue import threading from .http_common import * import simple_http_client import utils def pack_headers(headers): out_list = [] for k, v in headers.items(): if isinstance(v, int): out_list.append(b'%s: %d\r\n' % (utils.to_bytes(k), v)) else: out_list.append(b'%s: %s\r\n' % (utils.to_bytes(k), utils.to_bytes(v))) return b''.join(out_list) class Http1Worker(HttpWorker): def __init__(self, logger, ip_manager, config, ssl_sock, close_cb, retry_task_cb, idle_cb, log_debug_data): super(Http1Worker, self).__init__(logger, ip_manager, config, ssl_sock, close_cb, retry_task_cb, idle_cb, log_debug_data) self.version = "1.1" self.task = None self.transfered_size = 0 self.trace_time = [] self.trace_time.append([ssl_sock.create_time, "connect"]) self.record_active("init") self.task_queue = Queue() threading.Thread(target=self.work_loop, name="%s_http1_work_loop" % self.logger.name).start() self.idle_cb() if self.config.http1_first_ping_wait or self.config.http1_ping_interval: threading.Thread(target=self.keep_alive_thread, name="%s_http1_keep_alive" % self.logger.name).start() def record_active(self, active=""): self.trace_time.append([time.time(), active]) if len(self.trace_time) > self.config.http1_trace_size: self.trace_time.pop(0) # self.logger.debug("%s stat:%s", self.ip, active) def get_trace(self): out_list = [] last_time = self.trace_time[0][0] for t, stat in self.trace_time: time_diff = int((t - last_time) * 1000) last_time = t out_list.append(" %d:%s" % (time_diff, stat)) out_list.append(":%d" % ((time.time() - last_time) * 1000)) out_list.append(" processed:%d" % self.processed_tasks) out_list.append(" transfered:%d" % self.transfered_size) out_list.append(" sni:%s" % self.ssl_sock.sni) return ",".join(out_list) def request(self, task): self.accept_task = False self.task = task self.task_queue.put(task) def keep_alive_thread(self): while time.time() - self.ssl_sock.create_time < self.config.http1_first_ping_wait: if not self.keep_running: self.close("exit") return time.sleep(3) if self.config.http1_first_ping_wait and self.processed_tasks == 0: self.task_queue.put("ping") if self.config.http1_ping_interval: while self.keep_running: time_to_ping = max(self.config.http1_ping_interval - (time.time() - self.last_recv_time), 3) time.sleep(time_to_ping) if not self.request_onway and \ time.time() - self.last_recv_time > self.config.http1_ping_interval - 3: self.task_queue.put("ping") time.sleep(3) elif self.config.http1_idle_time: while self.keep_running: time_to_sleep = max(self.config.http1_idle_time - (time.time() - self.last_recv_time), 3) time.sleep(time_to_sleep) if not self.request_onway and time.time() - self.last_recv_time > self.config.http1_idle_time: self.close("idle timeout") return def work_loop(self): while self.keep_running: try: task = self.task_queue.get(block=True) except: task = None if not task: # None task means exit self.accept_task = False self.keep_running = False return if task == "ping": if not self.head_request(): self.ip_manager.recheck_ip(self.ssl_sock.ip_str) self.close("keep alive") return continue # self.logger.debug("http1 get task") time_now = time.time() if time_now - self.last_recv_time > self.config.http1_idle_time: self.logger.warn("get task but inactive time:%d", time_now - self.last_recv_time) self.task = task self.close("inactive timeout %d" % (time_now - self.last_recv_time)) return self.request_task(task) self.request_onway = False self.last_send_time = time_now life_end_reason = self.is_life_end() if life_end_reason: self.close("life_end:" + life_end_reason) return def request_task(self, task): timeout = task.timeout self.request_onway = True start_time = time.time() self.record_active("request") task.set_state("h1_req") task.headers[b'Host'] = self.get_host(task.host) request_len = len(task.body) task.headers[b"Content-Length"] = request_len request_data = b'%s %s HTTP/1.1\r\n' % (task.method, task.path) request_data += pack_headers(task.headers) request_data += b'\r\n' try: self.ssl_sock.send(request_data) payload_len = len(task.body) start = 0 while start < payload_len: send_size = min(payload_len - start, 65535) sended = self.ssl_sock.send(task.body[start:start + send_size]) start += sended task.set_state("h1_req_sent") except Exception as e: self.logger.warn("%s %s h1_request send:%r inactive_time:%d task.timeout:%d", self.ip_str, self.ssl_sock.getsockname(), e, time.time() - self.last_recv_time, task.timeout) self.logger.warn('%s trace:%s', self.ip_str, self.get_trace()) self.retry_task_cb(task) self.task = None self.close("send fail") return try: response = simple_http_client.Response(self.ssl_sock) response.begin(timeout=timeout) task.set_state("response_begin") self.last_recv_time = time.time() except Exception as e: self.logger.warn("%s h1_request recv:%r inactive_time:%d task.timeout:%d", self.ip_str, e, time.time() - self.last_recv_time, task.timeout) self.logger.warn('%s trace:%s', self.ip_str, self.get_trace()) self.retry_task_cb(task) self.task = None self.close("recv fail") return task.set_state("h1_get_head") time_left = timeout - (time.time() - start_time) if task.method == b"HEAD" or response.status in [204, 304]: response.content_length = 0 response.ssl_sock = self.ssl_sock response.task = task response.worker = self task.content_length = response.content_length task.responsed = True if task.queue: task.queue.put(response) if self.config.http2_show_debug: self.logger.debug("got res for %s status:%d", self.ip_str, response.status) try: read_target = int(response.content_length) except: read_target = 0 data_len = 0 while True: try: data = response.read(timeout=time_left) if not data: break except Exception as e: self.logger.warn("read fail, ip:%s, chunk:%d url:%s task.timeout:%d e:%r", self.ip_str, response.chunked, task.url, task.timeout, e) self.logger.warn('%s trace:%s', self.ip_str, self.get_trace()) self.close("read fail") return task.put_data(data) length = len(data) data_len += length if read_target and data_len >= read_target: break if read_target > data_len: self.logger.warn("read fail, ip:%s, chunk:%d url:%s task.timeout:%d ", self.ip_str, response.chunked, task.url, task.timeout) self.ip_manager.recheck_ip(self.ssl_sock.ip_str) self.close("down fail") task.finish() self.ssl_sock.received_size += data_len time_now = time.time() time_cost = (time_now - start_time) xcost = float(response.headers.get(b"X-Cost", -1)) if isinstance(xcost, list): xcost = float(xcost[0]) road_time = time_cost - xcost if xcost != -1 and road_time > 0: self.update_speed(road_time, request_len, data_len) task.set_state("h1_finish[RTT:%d]" % (road_time * 1000)) self.transfered_size += len(request_data) + data_len self.task = None self.accept_task = True self.idle_cb() self.processed_tasks += 1 self.last_recv_time = time.time() self.record_active("Res") def head_request(self): if not self.ssl_sock.host: # self.logger.warn("try head but no host set") return True # for keep alive, not work now. self.request_onway = True self.record_active("head") # start_time = time.time() # self.logger.debug("head request %s", self.ip) request_data = b'GET / HTTP/1.1\r\nHost: %s\r\n\r\n' % utils.to_bytes(self.ssl_sock.host) try: data = request_data ret = self.ssl_sock.send(data) if ret != len(data): self.logger.warn("h1 head send len:%r %d %s", ret, len(data), self.ip_str) self.logger.warn('%s trace:%s', self.ip_str, self.get_trace()) return False response = simple_http_client.Response(self.ssl_sock) response.begin(timeout=5) status = response.status if status not in [200, 404]: self.logger.warn("%s host:%s head fail status:%d", self.ip_str, self.ssl_sock.host, status) return False content = response.readall(timeout=5) self.record_active("head end") self.last_recv_time = time.time() return True except Exception as e: self.logger.warn("h1 %s HEAD keep alive request fail:%r", self.ssl_sock.ip_str, e) self.logger.warn('%s trace:%s', self.ip_str, self.get_trace()) return False finally: self.request_onway = False def close(self, reason=""): # Notify loop to exit # This function may be call by outside http2 # When gae_proxy found the appid or ip is wrong self.task_queue.put(None) if self.task is not None: if self.task.responsed: self.task.finish() else: self.retry_task_cb(self.task) self.task = None super(Http1Worker, self).close(reason)