Merge pull request #2546 from langbot-app/feat/marketplace-installed-state-and-search

feat(plugins): show installed state in marketplace and search install…
This commit is contained in:
RockChinQ
2026-09-22 00:43:46 +08:00
committed by GitHub
25 changed files with 878 additions and 174 deletions
+48
View File
@@ -94,8 +94,10 @@ async def _read_httpx_response_limited(
response: httpx.Response,
*,
max_bytes: int,
task_context: taskmgr.TaskContext | None = None,
) -> bytes:
content_length = response.headers.get('content-length')
declared_size: int | None = None
if content_length is not None:
try:
declared_size = int(content_length)
@@ -104,11 +106,25 @@ async def _read_httpx_response_limited(
if declared_size is not None and declared_size > max_bytes:
raise ValueError(f'Remote response exceeds the {max_bytes}-byte limit')
if task_context is not None and declared_size is not None:
# Publish the advertised size up-front so the UI can render a
# determinate bar even before the first chunk arrives.
task_context.metadata['download_total'] = declared_size
start_time = time.time()
body = bytearray()
async for chunk in response.aiter_bytes(chunk_size=64 * 1024):
body.extend(chunk)
if len(body) > max_bytes:
raise ValueError(f'Remote response exceeds the {max_bytes}-byte limit')
if task_context is not None:
elapsed = time.time() - start_time
task_context.metadata.update(
{
'download_current': len(body),
'download_speed': len(body) / elapsed if elapsed > 0 else 0,
}
)
return bytes(body)
@@ -118,6 +134,7 @@ async def _marketplace_get(
*,
max_bytes: int,
allow_not_found: bool = False,
task_context: taskmgr.TaskContext | None = None,
) -> tuple[int, bytes]:
async with client.stream('GET', url) as response:
if allow_not_found and response.status_code == 404:
@@ -126,6 +143,7 @@ async def _marketplace_get(
return response.status_code, await _read_httpx_response_limited(
response,
max_bytes=max_bytes,
task_context=task_context,
)
@@ -1711,6 +1729,7 @@ class PluginRuntimeConnector(ManagedRuntimeConnector):
client,
f'{space_url}/api/v1/marketplace/plugins/download/{plugin_author}/{plugin_name}/{version}',
max_bytes=_MARKETPLACE_PLUGIN_DOWNLOAD_MAX_BYTES,
task_context=task_context,
)
return plugin_package, version
@@ -1765,7 +1784,21 @@ class PluginRuntimeConnector(ManagedRuntimeConnector):
plugin_name = str(install_info.get('plugin_name') or '')
file_bytes: bytes | None
if task_context is not None:
# Reset the per-install counters so re-installing the same plugin
# cannot inherit stale progress metadata from a previous task.
task_context.set_current_action('preparing plugin install')
task_context.metadata.update(
{
'download_total': 0,
'download_current': 0,
'download_speed': 0,
}
)
if install_source == PluginInstallSource.MARKETPLACE:
if task_context is not None:
task_context.set_current_action('downloading plugin package')
file_bytes, version = await self._download_marketplace_package(
execution_context,
plugin_author,
@@ -1791,6 +1824,8 @@ class PluginRuntimeConnector(ManagedRuntimeConnector):
raise ValueError(f'Unsupported plugin install source: {install_source.value}')
install_info, verified_certificate = self._admit_plugin_archive(file_bytes, install_info)
if task_context is not None:
task_context.set_current_action('inspecting plugin package')
manifest_author, manifest_name = self._inspect_plugin_package(file_bytes, task_context)
if not manifest_author or not manifest_name:
raise ValueError('Plugin package manifest identity is missing')
@@ -1802,8 +1837,12 @@ class PluginRuntimeConnector(ManagedRuntimeConnector):
if task_context is not None:
task_context.metadata['plugin_name'] = f'{plugin_author}/{plugin_name}'
if task_context is not None:
task_context.set_current_action('storing plugin package')
artifact_digest = hashlib.sha256(file_bytes).hexdigest()
await self._store_artifact_package(execution_context, artifact_digest, file_bytes)
if task_context is not None:
task_context.set_current_action('persisting the installation')
try:
binding, previous_digest, previous_was_durable = await self._persist_installation_package(
execution_context,
@@ -1824,6 +1863,13 @@ class PluginRuntimeConnector(ManagedRuntimeConnector):
certification_facts = verified_certificate.for_installation(binding.installation_uuid)
if certification_facts.artifact_digest != install_info['_certification']['normalized_digest']:
raise RuntimeError('Plugin certification digest changed before Runtime apply')
if task_context is not None:
# The runtime installs the plugin's dependencies and starts it
# inside apply_plugin_installation. It does not stream
# per-dependency progress back to this task context, so this stage
# deliberately stays coarse instead of claiming a separate,
# unobservable "installing dependencies" step.
task_context.set_current_action('installing or starting plugin')
await self._apply_desired_state(
PluginInstallationDesiredState(binding=binding, enabled=True),
artifact_package=file_bytes,
@@ -1841,6 +1887,8 @@ class PluginRuntimeConnector(ManagedRuntimeConnector):
pass
except Exception as exc:
self.ap.logger.debug(f'Legacy OSS plugin cleanup skipped: {exc}')
if task_context is not None:
task_context.set_current_action('waiting for plugin to become ready')
await self._wait_for_installed_plugin_ready(plugin_author, plugin_name, task_context)
async def upgrade_plugin(