feat: queue pool in thread

This commit is contained in:
aclist 2026-02-16 22:36:44 +09:00
parent 93a4e231ac
commit 7766b4c5c2
4 changed files with 79 additions and 36 deletions

View File

@ -43,6 +43,7 @@ from dzgui.config import update
from dzgui.config.query import lookup
from dzgui.config.userprefs import UserPrefs
from dzgui.controllers.emitter import Emitter
from dzgui.model.filtered_model import FilteredModelManager
from dzgui.model.misc_model import ModelManager
from dzgui.util import strings
from dzgui.util._json import read_json, write_json
@ -148,9 +149,9 @@ class Controller(GObject.GObject):
except AttributeError:
logger.critical(f"{attr} is not a valid AppNavigation attribute.")
def terminate_process(self) -> None:
# TODO: only used by server table multiprocessing queue
self.get_active_treeview().terminate_process()
# def terminate_process(self) -> None:
# # TODO: only used by server table multiprocessing queue
# self.get_active_treeview().terminate_process()
def get_prefs(self) -> UserPrefs:
return self.prefs
@ -410,7 +411,7 @@ class Controller(GObject.GObject):
def dump_ips(self, ips: list) -> None:
# NOTE: block malformed records (TODO: add github issue no.)
ips = [ip for ip in ips if len(ip.split(":")) == 3 and ip.split(":")[2] != "" ]
ips = [ip for ip in ips if len(ip.split(":")) == 3 and ip.split(":")[2] != ""]
with ThreadPoolExecutor() as executor:
futures = [
executor.submit(
@ -421,7 +422,7 @@ class Controller(GObject.GObject):
for ip in ips
]
# TODO: update dialog in main loop
#wait(futures)
# wait(futures)
serv = []
for future in as_completed(futures):
res = future.result()
@ -435,17 +436,34 @@ class Controller(GObject.GObject):
parsed = Servers.parse_json(serv)
self.push_data_success(parsed, FilterMode.INITIAL)
def dump_lan(self, port: int) -> None:
serv = []
with ThreadPoolExecutor() as executor:
futures = [executor.submit(Servers.test_ip, i, port) for i in range(1, 256)]
for future in as_completed(futures):
try:
res = future.result(timeout=3)
if res is None:
continue
serv.panned(res)
if len(serv) == 0:
self.cleanup_func = StoredFunc(self.cleanup_on_failure)
return
parsed = Servers.parse_json(serv)
self.push_data_success(parsed, FilterMode.INITIAL)
def dump_api(self) -> None:
key = self.query_config(Preferences.STEAM)
job = Servers.query_api
params = Servers.params
serv = []
i = 0
#i = 0
with ThreadPoolExecutor() as executor:
futures = [executor.submit(job, key, APPID_DAYZ, param) for param in params]
for future in as_completed(futures):
try:
i += 1
#i += 1
GLib.idle_add(lambda: self.wait_dialog.increment())
res = future.result(timeout=3)
if res.status != 200 or not res.parsed:
@ -621,8 +639,9 @@ class Controller(GObject.GObject):
# FIXME: signal should instead be emitted off of treeview when rows added/inserted
# self.update_mod_statusbar()
# call a2s on thread and update ephemeral model in situ
case ContextMenu.REFRESH_PLAYERS:
# get record
# call a2s on thread
pass
# update history model, update tab label, pop off of queue, write new list into file
@ -719,14 +738,13 @@ class Controller(GObject.GObject):
self.pending_jobs = 1
treeview = self.get_active_treeview()
treeview.set_loaded(True)
treeview.set_model(self.to_insert)
# TODO: this will allow history and saved tab to emit signals to statusbar
# CHORE: test if treeview's sort method inserts row at the correct index
# inserting a row serializes file on disk, updates control model for that tab, and
# reapplies filters to ephemeral model; since filters are applied, in-situ insertion might not be necessary
self.to_insert.connect("row-inserted", lambda: print("row inserted into model"))
map_man = treeview.get_map_man()
# TODO: signals or other approach to deferring map
# model insertion after thread closes
@ -738,6 +756,7 @@ class Controller(GObject.GObject):
# CHORE: this is placeholder logic
if self.first_iteration:
map_man = treeview.get_map_man()
map_man.set_unique_maps(self.new_maps)
self.emitter.emit("servers_loaded_init")
self.first_iteration = False
@ -773,14 +792,14 @@ class Controller(GObject.GObject):
dialog.run()
def push_data_success(self, data: tuple, mode: Optional[FilterMode]) -> None:
# FIXME: set outside of thread
treeview = self.get_active_treeview()
manager = treeview.get_filter_man()
# FIXME: calls treeview read methods in thread
# treeview = self.get_active_treeview()
# manager = treeview.get_filter_man()
manager = self.get_filter_man()
if data is None:
self.to_insert = None
else:
# TODO: consolidate into filter manager
if mode == FilterMode.INITIAL:
manager.set_control(data)
manager.filter(mode)
@ -789,7 +808,6 @@ class Controller(GObject.GObject):
u_maps = set([row[1] for row in data])
self.new_maps = sorted(u_maps)
treeview.set_loaded(True)
self.cleanup_func = StoredFunc(self.cleanup_on_success)
def highlight_stale_cleanup(self, stale_mods: list) -> None:
@ -892,23 +910,27 @@ class Controller(GObject.GObject):
func.call()
@call_on_thread(strings.dialog.filtering)
def filter_threaded(self, mode: FilterMode, label: str) -> None:
# FIXME: call outside of thread
tv = self.get_active_treeview()
filter_man = tv.get_filter_man()
filter_man.filter(mode, label)
def filter_threaded(
self, filter_man: "FilteredModelManager", mode: FilterMode, label: str
) -> None:
# TODO: why pushing empty data?
filter_man.filter(mode, label)
self.push_data_success("", mode)
# TODO: consolidate with method above
# TODO: call filter_man methods directly
def refilter_model(self, mode: FilterMode, label: Optional[str] = None) -> None:
tv = self.get_active_treeview()
# tv.freeze_child_notify()
filter_man = tv.get_filter_man()
if filter_man.get_control() is None:
return
self.filter_threaded(mode, label)
self.filter_threaded(filter_man, mode, label)
def get_filter_man(self) -> "FilteredModelManager":
return self.filter_man
def set_filter_man(self, filter_man: "FilteredModelManager") -> None:
self.filter_man = filter_man
def populate_model(self) -> None:
# NOTE: prepare GTK objects outside of thread
@ -930,6 +952,9 @@ class Controller(GObject.GObject):
self.first_iteration = True
self.pending_jobs = jobs
# TODO: get filter manager a priori and set it so that it can be accessed out of thread
filter_man = treeview.get_filter_man()
self.set_filter_man(filter_man)
self.run_query_func(func)
def get_favorite(self) -> tuple[str, str] | tuple[None, None]:

View File

@ -1,5 +1,4 @@
import logging
import multiprocessing
from math import radians, cos, sin, asin, sqrt
from typing import TYPE_CHECKING
@ -40,12 +39,12 @@ class Haversine:
return self.dist / 1609.344
class CalcDist(multiprocessing.Process):
class CalcDist:
def __init__(
self,
addr: str,
enum: "ServerTab",
result_queue: multiprocessing.Queue,
result_queue: "Queue",
controller: "Controller",
) -> None:
super().__init__()
@ -54,9 +53,9 @@ class CalcDist(multiprocessing.Process):
self.controller = controller
self.result_queue = result_queue
self.addr = addr
self.ip = addr.split(":")[0]
self.ip = self.addr #.split(":")[0]
def run(self) -> None:
#def run(self) -> None:
cache = self.controller.get_dist_cache()
if self.addr in cache:
logger.info(f"Address '{self.addr}' already in cache")

View File

@ -171,7 +171,7 @@ class OuterWindow(Gtk.Window):
self.halt_proc_and_quit()
def halt_proc_and_quit(self) -> None:
MainController.terminate_process()
#MainController.terminate_process()
MainController.save_res_and_quit()

View File

@ -64,7 +64,9 @@ class ServerTreeView(ContextMixin, TreeView):
self.seen_cache = []
self.current_proc = None
self.queue = multiprocessing.Queue()
from queue import Queue
self.queue = Queue()
#self.queue = multiprocessing.Queue()
prefs = self.controller.get_prefs()
columns = prefs.paths.columns
@ -129,6 +131,8 @@ class ServerTreeView(ContextMixin, TreeView):
self.connect("map", self._on_map)
self.connect("unmap", self._on_unmap)
self.thread = None
def _get_ping(
self,
column: Gtk.TreeViewColumn,
@ -235,12 +239,15 @@ class ServerTreeView(ContextMixin, TreeView):
# NOTE: get final width after drag action completes
GLib.idle_add(self.controller.propagate_column_width, col)
def terminate_process(self) -> None:
if self.current_proc and self.current_proc.is_alive():
self.current_proc.terminate()
#def terminate_process(self) -> None:
# if self.thread and self.thread.is_alive():
# print("thread still exists")
# if self.current_proc and self.current_proc.is_alive():
# self.current_proc.terminate()
def start_distcalc(self, emitter: Optional["Emitter"] = None):
self.terminate_process()
#self.terminate_process()
#self.queue_id = GLib.timeout_add(QUEUE_CHECK_DELAY, self._check_result_queue)
self.emitter.emit("distcalc_started")
record = self.get_record()
if record is None:
@ -255,10 +262,22 @@ class ServerTreeView(ContextMixin, TreeView):
self.controller.set_statusbar_dist(haversine, self.get_enum())
return
self.current_proc = CalcDist(
record.ip, self.get_enum(), self.queue, self.controller
#def t(ip, enum, queue) -> None:
# if queue.empty():
# queue.put(["LONG DISTANCE", 0])
enum = self.get_enum()
self.thread = threading.Thread(
daemon=True,
target=CalcDist,
args=(record.ip, enum, self.queue, self.controller),
)
self.current_proc.start()
self.thread.start()
#self.current_proc = CalcDist(
# record.ip, self.get_enum(), self.queue, self.controller
#)
#self.current_proc.start()
def _check_result_queue(self) -> Literal[True]:
latest_result = None