feat(plugin): route certified installs to shared workers

This commit is contained in:
RockChinQ
2026-09-25 12:48:55 +00:00
parent 9631eccaaf
commit 9f296bd57e
11 changed files with 320 additions and 90 deletions
@@ -6,6 +6,8 @@ import zipfile
import pytest
from langbot_plugin.entities.io.context import PluginExecutionMode
@pytest.mark.parametrize(
('deployment', 'certificate', 'force', 'expected_disposition', 'expected_code'),
@@ -150,6 +152,80 @@ def test_log_visibility_policy_only_scopes_valid_shared_certifications(
assert visibility.value == expected_visibility
@pytest.mark.parametrize(
('certification', 'expected_mode'),
[
(
{
'artifact_digest': 'a' * 64,
'verification': 'valid',
'certificate_runtime_profile': 'shared-runtime-v1',
'runtime_profile': 'shared-runtime-v1',
'admission_code': 'CERTIFIED_PLUGIN_SHARED_ELIGIBLE',
},
PluginExecutionMode.SHARED_CERTIFIED,
),
(None, PluginExecutionMode.DEDICATED),
({}, PluginExecutionMode.DEDICATED),
(
{
'artifact_digest': 'a' * 64,
'verification': 'invalid',
'certificate_runtime_profile': 'shared-runtime-v1',
'runtime_profile': 'shared-runtime-v1',
'admission_code': 'CERTIFIED_PLUGIN_SHARED_ELIGIBLE',
},
PluginExecutionMode.DEDICATED,
),
(
{
'artifact_digest': 'b' * 64,
'verification': 'valid',
'certificate_runtime_profile': 'shared-runtime-v1',
'runtime_profile': 'shared-runtime-v1',
'admission_code': 'CERTIFIED_PLUGIN_SHARED_ELIGIBLE',
},
PluginExecutionMode.DEDICATED,
),
(
{
'artifact_digest': 'a' * 64,
'verification': 'valid',
'certificate_runtime_profile': 'shared-runtime-v1',
'runtime_profile': 'dedicated',
'admission_code': 'CERTIFIED_PLUGIN_SHARED_ELIGIBLE',
},
PluginExecutionMode.DEDICATED,
),
(
{
'artifact_digest': 'a' * 64,
'verification': 'valid',
'certificate_runtime_profile': 'shared-runtime-v1',
'runtime_profile': 'shared-runtime-v1',
'admission_code': 'CERTIFIED_PLUGIN_OSS_FORCED_DEDICATED',
},
PluginExecutionMode.DEDICATED,
),
],
)
def test_persisted_certification_selects_shared_execution_only_for_exact_admitted_artifact(
certification: dict[str, str] | None,
expected_mode: PluginExecutionMode,
) -> None:
from langbot.pkg.plugin.certification import execution_mode_for_persisted_installation
install_info = {} if certification is None else {'_certification': certification}
assert (
execution_mode_for_persisted_installation(
artifact_digest='a' * 64,
install_info=install_info,
)
is expected_mode
)
def _archive_bytes(manifest: dict[str, object]) -> bytes:
buffer = io.BytesIO()
with zipfile.ZipFile(buffer, 'w') as archive:
@@ -17,7 +17,7 @@ from unittest.mock import AsyncMock, Mock
from importlib import import_module
from tests.factories import text_query
from langbot_plugin.entities.io.context import InstallationBinding
from langbot_plugin.entities.io.context import InstallationBinding, PluginExecutionMode
from langbot.pkg.api.http.context import ExecutionContext
from langbot.pkg.workspace.errors import WorkspaceNotFoundError
@@ -758,7 +758,16 @@ class TestSetPluginConfig:
runtime_revision=1,
artifact_digest=TEST_INSTALLATION_BINDING.artifact_digest,
enabled=True,
install_info={'_artifact_storage': 'tenant_binary_storage_v1'},
install_info={
'_artifact_storage': 'tenant_binary_storage_v1',
'_certification': {
'artifact_digest': TEST_INSTALLATION_BINDING.artifact_digest,
'verification': 'valid',
'certificate_runtime_profile': 'shared-runtime-v1',
'runtime_profile': 'shared-runtime-v1',
'admission_code': 'CERTIFIED_PLUGIN_SHARED_ELIGIBLE',
},
},
)
connector._setting_for_plugin = AsyncMock(return_value=(TEST_EXECUTION_CONTEXT, setting))
connector.ap.persistence_mgr.execute_async = AsyncMock(return_value=SimpleNamespace(rowcount=1))
@@ -777,6 +786,7 @@ class TestSetPluginConfig:
applied_binding,
artifact_package=None,
enabled=True,
execution_mode=PluginExecutionMode.SHARED_CERTIFIED,
)
@@ -8,7 +8,7 @@ from types import SimpleNamespace
from unittest.mock import AsyncMock, Mock
import pytest
from langbot_plugin.entities.io.context import InstallationBinding
from langbot_plugin.entities.io.context import InstallationBinding, PluginExecutionMode
from langbot_plugin.runtime.plugin.mgr import PluginInstallSource
from langbot.pkg.api.http.context import ExecutionContext
@@ -45,7 +45,17 @@ def mock_archive_admission(connector: PluginRuntimeConnector, digest: str) -> No
# is exercised by integration/plugin/test_certified_plugin_admission.py.
connector._admit_plugin_archive = Mock(
side_effect=lambda _package, info: (
{**info, '_certification': {'normalized_digest': digest}},
{
**info,
'_certification': {
'artifact_digest': digest,
'normalized_digest': digest,
'verification': 'valid',
'certificate_runtime_profile': 'shared-runtime-v1',
'runtime_profile': 'shared-runtime-v1',
'admission_code': 'CERTIFIED_PLUGIN_SHARED_ELIGIBLE',
},
},
SimpleNamespace(for_installation=lambda _uuid: SimpleNamespace(artifact_digest=digest)),
)
)
@@ -56,6 +66,7 @@ def plugin_setting(
artifact_digest: str,
*,
durable: bool = True,
certification: dict[str, str] | None = None,
) -> SimpleNamespace:
return SimpleNamespace(
plugin_author='author',
@@ -67,7 +78,10 @@ def plugin_setting(
priority=0,
created_at=datetime.datetime(2026, 1, 1),
install_source='local',
install_info={'_artifact_storage': 'tenant_binary_storage_v1'} if durable else {},
install_info={
**({'_artifact_storage': 'tenant_binary_storage_v1'} if durable else {}),
**({'_certification': certification} if certification is not None else {}),
},
)
@@ -163,6 +177,43 @@ async def test_shared_reconnect_replays_two_workspaces_and_removes_missing_proje
assert set(connector._known_desired_states) == {setting_a.installation_uuid}
@pytest.mark.asyncio
async def test_reconcile_reload_projects_certified_exact_artifact_to_shared_execution():
binding = execution_binding('workspace-a')
digest = 'a' * 64
setting = plugin_setting(
'01',
digest,
certification={
'artifact_digest': digest,
'verification': 'valid',
'certificate_runtime_profile': 'shared-runtime-v1',
'runtime_profile': 'shared-runtime-v1',
'admission_code': 'CERTIFIED_PLUGIN_SHARED_ELIGIBLE',
},
)
connector = shared_connector([[binding]], {'workspace-a': [setting]})
connector.handler = runtime_handler()
await connector._prepare_connected_runtime()
desired = connector.handler.reconcile_plugin_installations.await_args.args[0][0]
assert desired.execution_mode is PluginExecutionMode.SHARED_CERTIFIED
@pytest.mark.asyncio
async def test_reconcile_reload_defaults_legacy_installation_to_dedicated_execution():
binding = execution_binding('workspace-a')
setting = plugin_setting('01', 'a' * 64)
connector = shared_connector([[binding]], {'workspace-a': [setting]})
connector.handler = runtime_handler()
await connector._prepare_connected_runtime()
desired = connector.handler.reconcile_plugin_installations.await_args.args[0][0]
assert desired.execution_mode is PluginExecutionMode.DEDICATED
@pytest.mark.asyncio
async def test_empty_projected_workspaces_do_not_retain_installation_sets():
binding_a = execution_binding('workspace-a')
@@ -223,6 +274,7 @@ async def test_fresh_shared_runtime_cache_replays_persisted_local_package():
desired.binding,
artifact_package=package,
enabled=True,
execution_mode=PluginExecutionMode.DEDICATED,
)
@@ -304,6 +356,7 @@ async def test_local_install_persists_verified_package_before_runtime_apply():
binding,
artifact_package=package,
enabled=True,
execution_mode=PluginExecutionMode.DEDICATED,
)
@@ -396,6 +449,9 @@ async def test_marketplace_upgrade_reports_multistep_progress():
'download_current': 0,
'download_speed': 0,
}
assert connector.handler.apply_plugin_installation.await_args.kwargs['execution_mode'] is (
PluginExecutionMode.SHARED_CERTIFIED
)
@pytest.mark.asyncio
@@ -430,7 +486,13 @@ async def test_workspace_reads_do_not_wait_for_an_installation_apply():
connector._persist_installation_package = AsyncMock(return_value=(binding, None, False))
connector._wait_for_installed_plugin_ready = AsyncMock()
connector._load_workspace_desired_states = AsyncMock(
return_value=[PluginInstallationDesiredState(binding=binding, enabled=True)]
return_value=[
PluginInstallationDesiredState(
binding=binding,
enabled=True,
execution_mode=PluginExecutionMode.SHARED_CERTIFIED,
)
]
)
apply_started = asyncio.Event()
release_apply = asyncio.Event()
+25 -2
View File
@@ -10,7 +10,12 @@ from unittest.mock import AsyncMock, MagicMock, Mock
import pytest
from langbot_plugin.entities.io.actions.enums import LangBotToRuntimeAction, PluginToRuntimeAction
from langbot_plugin.entities.io.context import ActionContext, InstallationBinding, PluginInstallationDesiredState
from langbot_plugin.entities.io.context import (
ActionContext,
InstallationBinding,
PluginExecutionMode,
PluginInstallationDesiredState,
)
def make_handler(app):
@@ -90,7 +95,25 @@ async def test_reconcile_plugin_installations_accepts_configured_cold_start_time
await runtime_handler.reconcile_plugin_installations((desired,), timeout=900)
assert runtime_handler.call_action.await_args.kwargs["timeout"] == 900
assert runtime_handler.call_action.await_args.kwargs['timeout'] == 900
@pytest.mark.asyncio
async def test_apply_plugin_installation_serializes_certified_shared_execution_mode():
runtime_handler = make_handler(SimpleNamespace())
runtime_handler.send_file = AsyncMock(return_value='artifact-file')
runtime_handler.call_action = AsyncMock(return_value={'state': 'starting'})
binding = next(iter(runtime_handler._installation_bindings.values()))[0]
await runtime_handler.apply_plugin_installation(
binding,
artifact_package=b'package',
enabled=True,
execution_mode=PluginExecutionMode.SHARED_CERTIFIED,
)
assert runtime_handler.call_action.await_args.args[0] == LangBotToRuntimeAction.APPLY_PLUGIN_INSTALLATION
assert runtime_handler.call_action.await_args.args[1]['execution_mode'] == 'shared-runtime-v1'
class TestHandlerQueryVariables: