Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -20,15 +20,15 @@ jobs:
build:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v6
- uses: actions/checkout@v7
with:
token: ${{ secrets.GITHUB_TOKEN }}
repository: ${{ github.event.pull_request.head.repo.full_name }}
ref: ${{ github.head_ref }}
fetch-depth: 0

- name: Set up Python
uses: actions/setup-python@v6
uses: actions/setup-python@v7
with:
python-version: "3.12"

Expand Down
6 changes: 3 additions & 3 deletions .github/workflows/release.yml
Original file line number Diff line number Diff line change
Expand Up @@ -12,14 +12,14 @@ jobs:
name: Create Release
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v6
- uses: actions/checkout@v7
- name: Set up Python
uses: actions/setup-python@v6
uses: actions/setup-python@v7
with:
python-version: '3.12'

- name: 'Get Previous tag'
uses: oprypin/find-latest-tag@v1.1.2
uses: oprypin/find-latest-tag@v1.1.3
with:
repository: ${{ github.repository }}
releases-only: true # We know that all relevant tags have a GitHub release for them.
Expand Down
8 changes: 8 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,13 @@
## Plugin Manager (dd-mm-yyyy)

### 1.1.11 (09-08-2026)

- Switched to babase.app.asyncio_loop and babase.app.threadpool
for concurrent execution
- A new threadpool to avoid blocking the babase.app.threadpool
which is for short parallel tasks
- Added a shutdown task to shutdown our threadpool

### 1.1.10 (12-06-2026)

- Fix for older bs versions using `EXPORT_CLASS_NAME_SHORTCUTS` import
Expand Down
6 changes: 6 additions & 0 deletions index.json
Original file line number Diff line number Diff line change
@@ -1,6 +1,12 @@
{
"plugin_manager_url": "https://github.com/bombsquad-community/plugin-manager/{content_type}/{tag}/plugin_manager.py",
"versions": {
"1.1.11": {
"api_version": 9,
"commit_sha": "a74246d",
"released_on": "09-08-2026",
"md5sum": "1dbf57b11602a22196711bdd3ad18fdc"
},
"1.1.10": {
"api_version": 9,
"commit_sha": "a1baa5f",
Expand Down
122 changes: 86 additions & 36 deletions plugin_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
from bauiv1lib import popup, confirm
from bauiv1lib.settings.allsettings import AllSettingsWindow

import urllib.error
import urllib.request
import http.client
import socket
Expand All @@ -12,20 +13,22 @@

import re
import os
import sys
import copy
import asyncio
import pathlib
import hashlib
import weakref
import threading
import contextlib
import concurrent.futures

from typing import override
from datetime import datetime

# Modules used for overriding AllSettingsWindow
import logging

PLUGIN_MANAGER_VERSION = "1.1.10"
PLUGIN_MANAGER_VERSION = "1.1.11"
REPOSITORY_URL = "https://github.com/bombsquad-community/plugin-manager"
# Current tag can be changed to "staging" or any other branch in
# plugin manager repo for testing purpose.
Expand All @@ -43,11 +46,39 @@
}
PLUGIN_DIRECTORY = _env["python_directory_user"]

# compatibility for older API versions.
if _env.get("build_number", 0) < 22714:
babase._asyncio._g_asyncio_event_loop = babase._asyncio._asyncio_event_loop
NETWORK_REQUEST_TIMEOUT = 10 # seconds

loop = babase.app.asyncio_loop
pool = babase.app.threadpool

# babase.app.threadpool (aka `pool`, also wired up as the asyncio loop's
# default executor) is reserved for short parallel work and warns/starves
# on long-running tasks. Network requests routinely run past that, so they
# get their own small dedicated pool instead.
_network_pool = concurrent.futures.ThreadPoolExecutor(
max_workers=4,
thread_name_prefix="PluginManagerNetwork",
)


async def _shutdown_network_pool() -> None:
"""Drain _network_pool's workers so none outlive app shutdown.

ThreadPoolExecutor.shutdown(wait=True) blocks, so it's run on a
throwaway thread and awaited from here rather than blocking the
event loop (other shutdown tasks run concurrently with this one).
In-flight requests are bounded by NETWORK_REQUEST_TIMEOUT, so this
settles quickly.
"""
done = asyncio.Event()

def _join() -> None:
_network_pool.shutdown(wait=True, cancel_futures=True)
loop.call_soon_threadsafe(done.set)

threading.Thread(target=_join, daemon=True).start()
await done.wait()

loop = babase._asyncio._g_asyncio_event_loop

open_popups = []

Expand Down Expand Up @@ -126,16 +157,16 @@ class CategoryMetadataParseError(Exception):


def send_network_request(request):
return urllib.request.urlopen(request)
return urllib.request.urlopen(request, timeout=NETWORK_REQUEST_TIMEOUT)


async def async_send_network_request(request):
response = await loop.run_in_executor(None, send_network_request, request)
response = await loop.run_in_executor(_network_pool, send_network_request, request)
return response


def stream_network_response_to_file(request, file, md5sum=None, retries=3):
response = urllib.request.urlopen(request)
response = urllib.request.urlopen(request, timeout=NETWORK_REQUEST_TIMEOUT)
chunk_size = 16 * 1024
content = b""
with open(file, "wb") as fout:
Expand All @@ -160,7 +191,7 @@ def stream_network_response_to_file(request, file, md5sum=None, retries=3):
async def async_stream_network_response_to_file(request, file, md5sum=None, retries=3):

content = await loop.run_in_executor(
None,
_network_pool,
stream_network_response_to_file,
request,
file,
Expand Down Expand Up @@ -201,40 +232,48 @@ class DNSBlockWorkaround:

_google_dns_cache = {}

def apply():
@classmethod
def apply(cls):
opener = urllib.request.build_opener(
DNSBlockWorkaround._HTTPHandler,
DNSBlockWorkaround._HTTPSHandler,
cls._HTTPHandler,
cls._HTTPSHandler,
)
urllib.request.install_opener(opener)

def _resolve_using_google_dns(hostname):
response = urllib.request.urlopen(f"https://dns.google/resolve?name={hostname}")
@classmethod
def _resolve_using_google_dns(cls, hostname):
response = urllib.request.urlopen(
f"https://dns.google/resolve?name={hostname}",
timeout=NETWORK_REQUEST_TIMEOUT,
)
response = response.read()
response = json.loads(response)
resolved_host = response["Answer"][0]["data"]
return resolved_host

def _resolve_using_system_dns(hostname):
@classmethod
def _resolve_using_system_dns(cls, hostname):
resolved_host = socket.gethostbyname(hostname)
return resolved_host

def _resolve_with_workaround(hostname):
resolved_host_from_cache = DNSBlockWorkaround._google_dns_cache.get(hostname)
@classmethod
def _resolve_with_workaround(cls, hostname):
resolved_host_from_cache = cls._google_dns_cache.get(hostname)
if resolved_host_from_cache:
return resolved_host_from_cache

resolved_host_by_system_dns = DNSBlockWorkaround._resolve_using_system_dns(hostname)
resolved_host_by_system_dns = cls._resolve_using_system_dns(hostname)

if DNSBlockWorkaround._is_blocked(hostname, resolved_host_by_system_dns):
resolved_host = DNSBlockWorkaround._resolve_using_google_dns(hostname)
DNSBlockWorkaround._google_dns_cache[hostname] = resolved_host
if cls._is_blocked(hostname, resolved_host_by_system_dns):
resolved_host = cls._resolve_using_google_dns(hostname)
cls._google_dns_cache[hostname] = resolved_host
else:
resolved_host = resolved_host_by_system_dns

return resolved_host

def _is_blocked(hostname, address):
@classmethod
def _is_blocked(cls, hostname, address):
is_blocked = False
if hostname == "raw.githubusercontent.com":
# Jio's DNS server may be blocking it.
Expand Down Expand Up @@ -578,7 +617,7 @@ async def get_content(self):
if not self.is_installed:
raise PluginNotInstalled("Plugin is not available locally.")

self._content = await loop.run_in_executor(None, self._get_content)
self._content = await loop.run_in_executor(pool, self._get_content)
return self._content

async def get_api_version(self):
Expand Down Expand Up @@ -659,7 +698,7 @@ async def enable(self):
self.save()

def load_plugin(self, entry_point):
plugin_class = babase._general.getclass(entry_point, babase.Plugin)
plugin_class = babase.getclass(entry_point, babase.Plugin)
loaded_plugin_instance = plugin_class()
loaded_plugin_instance.on_app_running()

Expand All @@ -683,7 +722,7 @@ def set_version(self, version):
async def set_content(self, content):
if not self._content:

await loop.run_in_executor(None, self._set_content, content)
await loop.run_in_executor(pool, self._set_content, content)
self._content = content
return self

Expand All @@ -705,14 +744,23 @@ def save(self):
class PluginVersion:
def __init__(self, plugin, version, tag=CURRENT_TAG):
self.number, info = version
self.plugin = plugin
# Plugin already owns its PluginVersions (via `versions`,
# `latest_version`, `latest_compatible_version`); holding a strong
# back-reference here would form a Plugin<->PluginVersion cycle
# that only the cyclic GC can free. A weakref avoids that so
# they're freed by refcounting alone.
self._plugin_ref = weakref.ref(plugin)
self.api_version = info["api_version"]
self.released_on = info["released_on"]
self.commit_sha = info["commit_sha"]
self.md5sum = info["md5sum"]

self.download_url = self.plugin.url.format(content_type="raw", tag=tag)
self.view_url = self.plugin.url.format(content_type="blob", tag=tag)
self.download_url = plugin.url.format(content_type="raw", tag=tag)
self.view_url = plugin.url.format(content_type="blob", tag=tag)

@property
def plugin(self):
return self._plugin_ref()

def __eq__(self, plugin_version):
return (self.number, self.plugin.name) == (plugin_version.number,
Expand Down Expand Up @@ -869,7 +917,7 @@ def __init__(self):
self._index = _CACHE.get("index", {})
self._changelog = _CACHE.get("changelog", {})
self.categories = {}
self.module_path = sys.modules[__name__].__file__
self.module_path = __file__
self._index_setup_in_progress = False
self._changelog_setup_in_progress = False

Expand Down Expand Up @@ -899,7 +947,7 @@ async def setup_index(self):
await self.setup_plugin_categories(index)
self._index_setup_in_progress = False

async def get_changelog(self) -> list[str, bool]:
async def get_changelog(self) -> tuple[str, bool]:
requested = False
if not self._changelog:
request = urllib.request.Request(CHANGELOG_META.format(
Expand All @@ -911,7 +959,7 @@ async def get_changelog(self) -> list[str, bool]:
response = await async_send_network_request(request)
self._changelog = response.read().decode()
requested = True
return [self._changelog, requested]
return self._changelog, requested

async def setup_changelog(self, version=None) -> None:
if version is None:
Expand All @@ -930,6 +978,7 @@ async def setup_changelog(self, version=None) -> None:
released_on = full_changelog[0].split(version)[1].split('\n')[0]
matches = re.findall(pattern, full_changelog[0], re.DOTALL)
else:
released_on = ' (Not Provided)'
matches = None

if matches:
Expand All @@ -938,7 +987,7 @@ async def setup_changelog(self, version=None) -> None:
'info': matches[0].strip()
}
else:
changelog = {'released_on': ' (Not Provided)',
changelog = {'released_on': released_on,
'info': f"Changelog entry for version {version} not found."}
else:
changelog = full_changelog[0]
Expand Down Expand Up @@ -1735,6 +1784,7 @@ def _cancel(self) -> None:
_remove_popup(self)
bui.containerwidget(edit=self._root_widget, transition='out_scale')

@staticmethod
def button(fn):
async def asyncio_handler(fn, self, *args, **kwargs):
await fn(self, *args, **kwargs)
Expand Down Expand Up @@ -2260,8 +2310,8 @@ def main_window_should_preserve_selection(self) -> bool:

def __init__(
self,
transition: str = "in_right",
origin_widget: bui.Widget = None
transition: str | None = "in_right",
origin_widget: bui.Widget | None = None
):
self.plugin_manager = PluginManager()
self.category_selection_button = None
Expand Down Expand Up @@ -3419,7 +3469,7 @@ def on_app_running(self) -> None:
from bauiv1lib.settings import allsettings
allsettings.AllSettingsWindow = NewAllSettingsWindow
DNSBlockWorkaround.apply()
asyncio.set_event_loop(babase._asyncio._g_asyncio_event_loop)
babase.app.add_shutdown_task(_shutdown_network_pool())
startup_tasks = StartupTasks()

loop.create_task(startup_tasks.execute())
Loading