mirror of
https://github.com/agentscope-ai/ReMe.git
synced 2026-10-05 02:41:43 +00:00
feat(jobs): add a unified enabled switch for all job types (#580)
* feat(jobs): allow disabling cron schedules independently * refactor(jobs): unify job activation with enabled flag
This commit is contained in:
parent
67135cfc57
commit
81c16d3e2c
19 changed files with 184 additions and 18 deletions
|
|
@ -74,6 +74,20 @@ Keep secrets in `.env` or the process environment, never in committed configurat
|
|||
|
||||
`session_dir` must remain workspace-relative.
|
||||
|
||||
## Enabling jobs
|
||||
|
||||
`jobs.<name>.enabled` defaults to `true` for all Job types. Disabled jobs retain their configuration but do not
|
||||
start, expose service interfaces, or accept calls through `Application.run_job()` / `run_stream_job()`.
|
||||
For example, disable ReMe's Dream cron when a host plugin owns the schedule:
|
||||
|
||||
```bash
|
||||
reme start jobs.dream_cron.enabled=false
|
||||
```
|
||||
|
||||
Restart the service to apply the override. The separate `auto_dream` API remains available, and other jobs continue running.
|
||||
`enable_serve` independently controls service exposure: `enabled=true, enable_serve=false` keeps a Job available for
|
||||
local calls. Background and cron jobs are never service-exposed.
|
||||
|
||||
## LLM
|
||||
|
||||
The default LLM uses an OpenAI-compatible interface:
|
||||
|
|
|
|||
|
|
@ -79,6 +79,19 @@ base_url: ${LLM_BASE_URL:-https://example.com/v1}
|
|||
|
||||
`session_dir` 必须保持 workspace-relative。其他 workspace 子目录也应使用清晰、稳定的相对名称。
|
||||
|
||||
## Job 启停
|
||||
|
||||
`jobs.<name>.enabled` 默认为 `true`,适用于所有 Job 类型。禁用的 Job 保留配置,但不启动、不暴露接口,
|
||||
通过 `Application.run_job()` / `run_stream_job()` 调用时会报错。由宿主插件负责调度时,可以关闭 ReMe 的 Dream cron:
|
||||
|
||||
```bash
|
||||
reme start jobs.dream_cron.enabled=false
|
||||
```
|
||||
|
||||
重启服务后生效。独立的 `auto_dream` 接口仍可调用,其他 Job 继续运行。
|
||||
`enable_serve` 独立控制接口暴露:`enabled=true, enable_serve=false` 允许本地调用。
|
||||
Background 和 cron Job 始终不暴露服务接口。
|
||||
|
||||
## LLM 配置
|
||||
|
||||
默认 LLM 使用 OpenAI-compatible 接口:
|
||||
|
|
|
|||
|
|
@ -41,9 +41,14 @@ host-specific plugin.
|
|||
|
||||
```bash
|
||||
reme start workspace_dir=/absolute/path/to/your/reme-workspace \
|
||||
service.host=127.0.0.1 service.port=3457
|
||||
service.host=127.0.0.1 service.port=3457 \
|
||||
jobs.dream_cron.enabled=false
|
||||
```
|
||||
|
||||
The plugin owns the daily Dream schedule by default, so this disables ReMe's `dream_cron` while keeping the
|
||||
`auto_dream` API available. To use ReMe scheduling instead, omit this override and set the plugin's
|
||||
`autoDreamEnabled` to `false`. Shared services should have only one scheduler.
|
||||
|
||||
For development and screenshots, use an isolated directory outside the repository, such as `/tmp/reme-dsh-demo`. Do not write runtime memory into the repository's `.reme/` directory.
|
||||
|
||||
### 3.2 Install the DSH bundle
|
||||
|
|
|
|||
|
|
@ -40,9 +40,13 @@ ReMe HTTP 服务默认监听 `http://127.0.0.1:2333`,不使用 API Key 认证
|
|||
|
||||
```bash
|
||||
reme start workspace_dir=/absolute/path/to/your/reme-workspace \
|
||||
service.host=127.0.0.1 service.port=3457
|
||||
service.host=127.0.0.1 service.port=3457 \
|
||||
jobs.dream_cron.enabled=false
|
||||
```
|
||||
|
||||
插件默认负责每日 Dream 调度,因此这里关闭 ReMe 的 `dream_cron`,保留 `auto_dream` 接口。
|
||||
若改由 ReMe 调度,省略该覆盖并将插件的 `autoDreamEnabled` 设为 `false`。多个宿主共享服务时,只启用一个调度方。
|
||||
|
||||
开发或截图测试时建议使用仓库外的独立目录,例如 `/tmp/reme-dsh-demo`,不要把运行时记忆写入 ReMe 仓库自身的 `.reme/`。
|
||||
|
||||
### 3.2 安装插件
|
||||
|
|
|
|||
|
|
@ -40,9 +40,14 @@ pip install "reme-ai[core]"
|
|||
reme start \
|
||||
workspace_dir=/absolute/path/to/reme-workspace \
|
||||
service.host=127.0.0.1 \
|
||||
service.port=3458
|
||||
service.port=3458 \
|
||||
jobs.dream_cron.enabled=false
|
||||
```
|
||||
|
||||
The plugin owns the daily Dream schedule by default, so this disables ReMe's `dream_cron` while keeping the
|
||||
`auto_dream` API available. To use ReMe scheduling instead, omit this override and set the plugin's
|
||||
`autoDreamEnabled` to `false`. Shared services should have only one scheduler.
|
||||
|
||||
Verify the service without exposing configuration or credentials:
|
||||
|
||||
```bash
|
||||
|
|
|
|||
|
|
@ -37,9 +37,13 @@ pip install "reme-ai[core]"
|
|||
reme start \
|
||||
workspace_dir=/absolute/path/to/reme-workspace \
|
||||
service.host=127.0.0.1 \
|
||||
service.port=3458
|
||||
service.port=3458 \
|
||||
jobs.dream_cron.enabled=false
|
||||
```
|
||||
|
||||
插件默认负责每日 Dream 调度,因此这里关闭 ReMe 的 `dream_cron`,保留 `auto_dream` 接口。
|
||||
若改由 ReMe 调度,省略该覆盖并将插件的 `autoDreamEnabled` 设为 `false`。多个宿主共享服务时,只启用一个调度方。
|
||||
|
||||
不读取或输出模型配置即可验证服务:
|
||||
|
||||
```bash
|
||||
|
|
|
|||
|
|
@ -206,6 +206,9 @@ class Application(BaseComponent):
|
|||
|
||||
async def _start_one(self, c: BaseComponent) -> None:
|
||||
"""Start one component and record it for ordered shutdown."""
|
||||
if isinstance(c, BaseJob) and not c.enabled:
|
||||
self.logger.info(f"Skipping disabled job: {c.name}")
|
||||
return
|
||||
try:
|
||||
if isinstance(c, BackgroundJob):
|
||||
self.logger.info(f"Starting background job: {c.name}")
|
||||
|
|
@ -367,18 +370,23 @@ class Application(BaseComponent):
|
|||
|
||||
# ----- Job execution -------------------------------------------------
|
||||
|
||||
async def run_job(self, name: str, /, **kwargs) -> Response:
|
||||
"""Execute a registered job by name and return its final Response."""
|
||||
def _get_enabled_job(self, name: str) -> BaseJob:
|
||||
"""Resolve a job and reject execution when disabled."""
|
||||
if name not in self.context.jobs:
|
||||
raise KeyError(f"Job '{name}' not found")
|
||||
return await self.context.jobs[name](**kwargs)
|
||||
job = self.context.jobs[name]
|
||||
job.check_enabled()
|
||||
return job
|
||||
|
||||
async def run_job(self, name: str, /, **kwargs) -> Response:
|
||||
"""Execute a registered job by name and return its final Response."""
|
||||
return await self._get_enabled_job(name)(**kwargs)
|
||||
|
||||
async def run_stream_job(self, name: str, /, **kwargs) -> AsyncGenerator[StreamChunk, None]:
|
||||
"""Execute a streaming job, yielding chunks as they are produced."""
|
||||
if name not in self.context.jobs:
|
||||
raise KeyError(f"Job '{name}' not found")
|
||||
job = self._get_enabled_job(name)
|
||||
stream_queue: asyncio.Queue = asyncio.Queue()
|
||||
task = asyncio.create_task(self.context.jobs[name](stream_queue=stream_queue, **kwargs))
|
||||
task = asyncio.create_task(job(stream_queue=stream_queue, **kwargs))
|
||||
async for chunk in execute_stream_task(
|
||||
stream_queue=stream_queue,
|
||||
task=task,
|
||||
|
|
|
|||
|
|
@ -123,6 +123,7 @@ class BackgroundJob(BaseJob):
|
|||
|
||||
async def __call__(self, **kwargs) -> Response:
|
||||
"""Default body: run steps in order; errors propagate to supervisor."""
|
||||
self.check_enabled()
|
||||
self._record_call()
|
||||
merged = {**self.kwargs, **kwargs}
|
||||
context = RuntimeContext(stop_event=self._stop_event, **merged)
|
||||
|
|
|
|||
|
|
@ -44,12 +44,14 @@ class BaseJob(BaseComponent):
|
|||
parameters: dict | None = None,
|
||||
steps: list[ComponentConfig | dict] | None = None,
|
||||
enable_serve: bool = True,
|
||||
enabled: bool = True,
|
||||
**kwargs,
|
||||
):
|
||||
super().__init__(**kwargs)
|
||||
self.description = description
|
||||
self.parameters = parameters or {}
|
||||
self.step_configs = steps or []
|
||||
self.enabled = enabled
|
||||
self.enable_serve = enable_serve
|
||||
self.step_specs: list[tuple[type["BaseStep"], dict]] = []
|
||||
|
||||
|
|
@ -83,8 +85,14 @@ class BaseJob(BaseComponent):
|
|||
if isinstance(metadata, dict):
|
||||
global_counter_inc(metadata, ["__job_counter", self.name])
|
||||
|
||||
def check_enabled(self) -> None:
|
||||
"""Reject execution of a disabled job."""
|
||||
if not self.enabled:
|
||||
raise ValueError(f"Job '{self.name}' is disabled")
|
||||
|
||||
async def __call__(self, **kwargs) -> Response:
|
||||
"""Run all steps in order, capturing any failure into the response."""
|
||||
self.check_enabled()
|
||||
self._record_call()
|
||||
merged = {**self.kwargs, **kwargs}
|
||||
context = RuntimeContext(**merged)
|
||||
|
|
|
|||
|
|
@ -41,6 +41,7 @@ class CronJob(BackgroundJob):
|
|||
return context.response
|
||||
|
||||
async def __call__(self, **kwargs) -> Response:
|
||||
self.check_enabled()
|
||||
assert self._stop_event is not None
|
||||
while not self._stop_event.is_set():
|
||||
await self._wait_or_stop(self._next_fire_delay())
|
||||
|
|
|
|||
|
|
@ -12,6 +12,7 @@ class StreamJob(BaseJob):
|
|||
|
||||
async def __call__(self, **kwargs) -> None:
|
||||
"""Run steps; emit failures as ERROR chunks, then a terminal DONE marker."""
|
||||
self.check_enabled()
|
||||
self._record_call()
|
||||
merged = {**self.kwargs, **kwargs}
|
||||
context = RuntimeContext(**merged)
|
||||
|
|
|
|||
|
|
@ -76,12 +76,16 @@ class BaseService(BaseComponent):
|
|||
missing = sorted(self.jobs.difference(app.context.jobs))
|
||||
if missing:
|
||||
raise KeyError(f"Service jobs not found: {', '.join(missing)}")
|
||||
disabled = sorted(name for name in self.jobs if not app.context.jobs[name].enable_serve)
|
||||
disabled = sorted(
|
||||
name
|
||||
for name in self.jobs
|
||||
if not app.context.jobs[name].enabled or not app.context.jobs[name].enable_serve
|
||||
)
|
||||
if disabled:
|
||||
raise ValueError(f"Service jobs are not enabled for serving: {', '.join(disabled)}")
|
||||
|
||||
for name, job in app.context.jobs.items():
|
||||
if not job.enable_serve or (self.jobs is not None and name not in self.jobs):
|
||||
if not job.enabled or not job.enable_serve or (self.jobs is not None and name not in self.jobs):
|
||||
continue
|
||||
try:
|
||||
added = self.add_job(job)
|
||||
|
|
|
|||
|
|
@ -110,7 +110,10 @@ class HttpService(BaseService):
|
|||
conflicts = sorted(
|
||||
job.name
|
||||
for name, job in app.context.jobs.items()
|
||||
if job.enable_serve and (self.jobs is None or name in self.jobs) and f"/{job.name}" == self.mcp_path
|
||||
if job.enabled
|
||||
and job.enable_serve
|
||||
and (self.jobs is None or name in self.jobs)
|
||||
and f"/{job.name}" == self.mcp_path
|
||||
)
|
||||
if conflicts:
|
||||
names = ", ".join(conflicts)
|
||||
|
|
|
|||
|
|
@ -24,6 +24,7 @@ class JobConfig(ComponentConfig):
|
|||
parameters: dict = Field(default_factory=dict, description="Job-level parameters")
|
||||
steps: list[ComponentConfig] = Field(default_factory=list, description="Ordered step configs")
|
||||
enable_serve: bool = Field(default=True, description="Whether to expose this job through the service layer")
|
||||
enabled: bool = Field(default=True, description="Whether to start, serve, and allow execution of this job")
|
||||
|
||||
|
||||
class ApplicationConfig(BaseModel):
|
||||
|
|
|
|||
|
|
@ -31,7 +31,7 @@ class HelpStep(BaseStep):
|
|||
lines = []
|
||||
if self.app_context is not None:
|
||||
for name, job in self.app_context.jobs.items():
|
||||
if name == "help" or not getattr(job, "enable_serve", True):
|
||||
if name == "help" or not getattr(job, "enabled", True) or not getattr(job, "enable_serve", True):
|
||||
continue
|
||||
lines.append(f"🛠️ `{name}` — {job.description} 📥 {_format_params(job.parameters)}")
|
||||
|
||||
|
|
|
|||
|
|
@ -10,8 +10,24 @@ from reme.config.config_parser import (
|
|||
_read_config_file,
|
||||
parse_args,
|
||||
parse_dot_notation,
|
||||
parse_kwargs,
|
||||
resolve_app_config,
|
||||
)
|
||||
from reme.schema import ApplicationConfig
|
||||
|
||||
|
||||
def test_cli_disables_dream_schedule_without_disabling_manual_job():
|
||||
"""A CLI schedule override preserves callable jobs and other default schedules."""
|
||||
cfg = ApplicationConfig(
|
||||
**resolve_app_config(log_config=False, **parse_kwargs("jobs.dream_cron.enabled=false")),
|
||||
)
|
||||
|
||||
assert cfg.jobs["dream_cron"].enabled is False
|
||||
assert cfg.jobs["dream_cron"].backend == "cron"
|
||||
assert cfg.jobs["auto_dream"].enable_serve is True
|
||||
assert cfg.jobs["auto_dream"].enabled is True
|
||||
assert cfg.jobs["auto_dream"].steps
|
||||
assert cfg.jobs["optimize_index_cron"].enabled is True
|
||||
|
||||
|
||||
def test_load_builtin_config_by_filename_with_suffix():
|
||||
|
|
|
|||
|
|
@ -222,6 +222,15 @@ def test_http_service_fails_startup_preflight_for_mcp_job_conflict() -> None:
|
|||
service.add_jobs(app)
|
||||
|
||||
|
||||
def test_http_service_ignores_disabled_mcp_job_conflict() -> None:
|
||||
"""A disabled job does not reserve a service route."""
|
||||
service = HttpService(web_enabled=False)
|
||||
service.build_service(_FakeApplication()) # type: ignore[arg-type]
|
||||
app = SimpleNamespace(context=SimpleNamespace(jobs={"mcp": BaseJob(name="mcp", enabled=False)}))
|
||||
|
||||
service.add_jobs(app)
|
||||
|
||||
|
||||
def test_http_service_does_not_serve_symlinks_outside_static_dir(
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
|
|
|
|||
|
|
@ -384,6 +384,73 @@ def test_application_starts_jobs_base_stream_background_cron():
|
|||
asyncio.run(run())
|
||||
|
||||
|
||||
@pytest.mark.parametrize("cls", [BaseJob, StreamJob, BackgroundJob, CronJob])
|
||||
@pytest.mark.parametrize("enabled", [True, False])
|
||||
def test_application_respects_job_enabled_switch(cls, enabled):
|
||||
async def run():
|
||||
app = object.__new__(Application)
|
||||
app._started_components = []
|
||||
app.logger = MagicMock()
|
||||
app.context = SimpleNamespace(thread_pool=None)
|
||||
context = SimpleNamespace(registry=ComponentRegistry(), metadata={})
|
||||
options = {"cron": "0 23 * * *"} if cls is CronJob else {}
|
||||
job = cls(
|
||||
enabled=enabled,
|
||||
app_context=context,
|
||||
steps=[] if enabled else [{"backend": "missing_step"}],
|
||||
**options,
|
||||
)
|
||||
|
||||
try:
|
||||
await app._start_one(job)
|
||||
assert job.is_started is enabled
|
||||
assert app._started_components == ([job] if enabled else [])
|
||||
assert "enabled" not in job.kwargs
|
||||
if isinstance(job, BackgroundJob):
|
||||
assert (job._task is not None) is enabled
|
||||
assert job.enable_serve is False
|
||||
if not enabled:
|
||||
assert job.step_specs == []
|
||||
finally:
|
||||
await app._close_started_components()
|
||||
|
||||
assert not job.is_started
|
||||
if isinstance(job, BackgroundJob):
|
||||
assert job._task is None
|
||||
|
||||
asyncio.run(run())
|
||||
|
||||
|
||||
@pytest.mark.parametrize("cls", [BaseJob, StreamJob, BackgroundJob, CronJob])
|
||||
def test_application_rejects_disabled_job_calls(cls):
|
||||
async def run():
|
||||
app = object.__new__(Application)
|
||||
options = {"cron": "0 23 * * *"} if cls is CronJob else {}
|
||||
job = cls(name="disabled", enabled=False, **options)
|
||||
app.context = SimpleNamespace(jobs={"disabled": job})
|
||||
|
||||
with pytest.raises(ValueError, match="Job 'disabled' is disabled"):
|
||||
await job()
|
||||
with pytest.raises(ValueError, match="Job 'disabled' is disabled"):
|
||||
await app.run_job("disabled")
|
||||
with pytest.raises(ValueError, match="Job 'disabled' is disabled"):
|
||||
async for _ in app.run_stream_job("disabled"):
|
||||
pytest.fail("Disabled job emitted a chunk")
|
||||
|
||||
asyncio.run(run())
|
||||
|
||||
|
||||
def test_application_can_call_non_servable_job():
|
||||
async def run():
|
||||
app = object.__new__(Application)
|
||||
job = BaseJob(name="manual", enable_serve=False)
|
||||
app.context = SimpleNamespace(jobs={"manual": job})
|
||||
response = await app.run_job("manual")
|
||||
assert response.success is True
|
||||
|
||||
asyncio.run(run())
|
||||
|
||||
|
||||
def test_application_start_failure_propagates_and_closes_started_components():
|
||||
async def run():
|
||||
class GoodComponent(BaseComponent):
|
||||
|
|
|
|||
|
|
@ -36,12 +36,13 @@ def _app_with_jobs(**jobs):
|
|||
return SimpleNamespace(context=SimpleNamespace(jobs=jobs))
|
||||
|
||||
|
||||
def test_service_registers_all_enabled_jobs_by_default():
|
||||
@pytest.mark.parametrize("options", [{"enable_serve": False}, {"enabled": False}])
|
||||
def test_service_registers_all_enabled_jobs_by_default(options):
|
||||
"""Omitting service.jobs preserves registration of every service-enabled job."""
|
||||
service = MCPService()
|
||||
service.add_job = Mock(return_value=True)
|
||||
enabled = BaseJob(name="enabled")
|
||||
disabled = BaseJob(name="disabled", enable_serve=False)
|
||||
disabled = BaseJob(name="disabled", **options)
|
||||
|
||||
service.add_jobs(_app_with_jobs(enabled=enabled, disabled=disabled))
|
||||
|
||||
|
|
@ -73,7 +74,8 @@ def test_empty_service_jobs_disables_job_registration():
|
|||
service.add_job.assert_not_called()
|
||||
|
||||
|
||||
def test_explicit_service_jobs_reject_missing_disabled_and_unsupported_jobs():
|
||||
@pytest.mark.parametrize("options", [{"enable_serve": False}, {"enabled": False}])
|
||||
def test_explicit_service_jobs_reject_missing_disabled_and_unsupported_jobs(options):
|
||||
"""An explicit service.jobs list fails instead of starting an incomplete service."""
|
||||
missing_service = MCPService(jobs=["missing"])
|
||||
with pytest.raises(KeyError, match="missing"):
|
||||
|
|
@ -81,7 +83,7 @@ def test_explicit_service_jobs_reject_missing_disabled_and_unsupported_jobs():
|
|||
|
||||
disabled_service = MCPService(jobs=["disabled"])
|
||||
with pytest.raises(ValueError, match="disabled"):
|
||||
disabled_service.add_jobs(_app_with_jobs(disabled=BaseJob(name="disabled", enable_serve=False)))
|
||||
disabled_service.add_jobs(_app_with_jobs(disabled=BaseJob(name="disabled", **options)))
|
||||
|
||||
stream_service = MCPService(jobs=["stream"])
|
||||
stream_service.add_job = Mock(return_value=False)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue