From c7ae2d86ec7957ae38003d8bc9e2de539418aa1e Mon Sep 17 00:00:00 2001 From: Jerry Yan <792602257@qq.com> Date: Mon, 6 Jul 2026 18:07:53 +0800 Subject: [PATCH] feat: add skill learning runtime --- .../changes/skill-learning-runtime/tasks.md | 58 ++-- pyproject.toml | 2 + runtime/task.py | 70 ++++- skills_learning/__init__.py | 16 ++ skills_learning/config.py | 61 +++++ skills_learning/embeddings.py | 44 ++++ skills_learning/models.py | 186 +++++++++++++ skills_learning/retrieval.py | 55 ++++ skills_learning/store.py | 141 ++++++++++ skills_learning/synthesis.py | 248 ++++++++++++++++++ skills_learning/versioning.py | 72 +++++ tests/test_skill_retrieval.py | 91 +++++++ tests/test_skill_store.py | 54 ++++ tests/test_skill_synthesis.py | 113 ++++++++ tests/test_skill_task_runner.py | 217 +++++++++++++++ tests/test_skill_versioning.py | 72 +++++ tests/test_smoke.py | 1 + 17 files changed, 1470 insertions(+), 31 deletions(-) create mode 100644 skills_learning/__init__.py create mode 100644 skills_learning/config.py create mode 100644 skills_learning/embeddings.py create mode 100644 skills_learning/models.py create mode 100644 skills_learning/retrieval.py create mode 100644 skills_learning/store.py create mode 100644 skills_learning/synthesis.py create mode 100644 skills_learning/versioning.py create mode 100644 tests/test_skill_retrieval.py create mode 100644 tests/test_skill_store.py create mode 100644 tests/test_skill_synthesis.py create mode 100644 tests/test_skill_task_runner.py create mode 100644 tests/test_skill_versioning.py diff --git a/openspec/changes/skill-learning-runtime/tasks.md b/openspec/changes/skill-learning-runtime/tasks.md index 6dfd829..422853a 100644 --- a/openspec/changes/skill-learning-runtime/tasks.md +++ b/openspec/changes/skill-learning-runtime/tasks.md @@ -1,49 +1,49 @@ ## 1. Package scaffolding -- [ ] 1.1 Create `skills_learning/` package with `__init__.py`, `models.py`, `synthesis.py`, `versioning.py`, `embeddings.py`, `retrieval.py`, `config.py`, `store.py` -- [ ] 1.2 Add `skills_learning*` to `pyproject.toml`'s `[tool.setuptools.packages.find].include` list and add the embedding-provider SDK dependency -- [ ] 1.3 Add `skills_learning/config.py` with `SkillAuthoringConfig` (enable flag default `False`, divergence tolerance, embedding model name, default `top_k`) and a module-level accessor mirroring `semantic-scene-runtime`'s/`world-model-runtime`'s config-module pattern -- [ ] 1.4 Add `skills_learning/models.py` defining (or importing, if `skill-catalog-subscription` is already implemented) the shared `Skill`/`FlowTemplateSkill`/`SkillMetadata` dataclass shapes, plus this capability's own `source`, `version`, `parent_version_id` fields +- [x] 1.1 Create `skills_learning/` package with `__init__.py`, `models.py`, `synthesis.py`, `versioning.py`, `embeddings.py`, `retrieval.py`, `config.py`, `store.py` +- [x] 1.2 Add `skills_learning*` to `pyproject.toml`'s `[tool.setuptools.packages.find].include` list and add the embedding-provider SDK dependency +- [x] 1.3 Add `skills_learning/config.py` with `SkillAuthoringConfig` (enable flag default `False`, divergence tolerance, embedding model name, default `top_k`) and a module-level accessor mirroring `semantic-scene-runtime`'s/`world-model-runtime`'s config-module pattern +- [x] 1.4 Add `skills_learning/models.py` defining (or importing, if `skill-catalog-subscription` is already implemented) the shared `Skill`/`FlowTemplateSkill`/`SkillMetadata` dataclass shapes, plus this capability's own `source`, `version`, `parent_version_id` fields ## 2. Local skill store (skill-authoring, skill-versioning) -- [ ] 2.1 Implement `skills_learning/store.py`: a local store for locally-synthesized `FlowTemplateSkill` records, separate from `skill-catalog-subscription`'s synced catalog, with `create_version()`, `get_by_id()`, `get_latest_by_name()`, `list_versions(name)` -- [ ] 2.2 Enforce `source = "local-synthesis"` tagging on every record written by this store; add a guard/test that this store never writes into or imports a write-path of `skill-catalog-subscription`'s catalog module -- [ ] 2.3 Write unit tests for the store's version-chain semantics: creating a new version does not delete/modify prior versions, and `get_latest_by_name()` returns the highest `version` +- [x] 2.1 Implement `skills_learning/store.py`: a local store for locally-synthesized `FlowTemplateSkill` records, separate from `skill-catalog-subscription`'s synced catalog, with `create_version()`, `get_by_id()`, `get_latest_by_name()`, `list_versions(name)` +- [x] 2.2 Enforce `source = "local-synthesis"` tagging on every record written by this store; add a guard/test that this store never writes into or imports a write-path of `skill-catalog-subscription`'s catalog module +- [x] 2.3 Write unit tests for the store's version-chain semantics: creating a new version does not delete/modify prior versions, and `get_latest_by_name()` returns the highest `version` ## 3. Timeline extraction and parameter abstraction (skill-authoring) -- [ ] 3.1 Implement `skills_learning/synthesis.py`'s tool-call extraction: given a `task_id`, read `storage.timeline.Timeline.read(task_id)` and produce an ordered list of `(tool_name, args)` pairs, filtering out read-only tool names (`describe_screen`, `screenshot`, `ui_tree`, `find_text`, `find_icon`) -- [ ] 3.2 Implement skeleton matching: given an extracted tool-name sequence, look up any stored skill (via `store.py`) whose latest version has the identical tool-name sequence -- [ ] 3.3 Implement cross-execution argument diffing: compare extracted argument values position-by-position against a matched stored version's steps, and promote any differing value into a named `{param}` placeholder plus a corresponding entry in the skill's `parameters` schema -- [ ] 3.4 Implement parameter naming: prefer a name derived from the corresponding `SemanticScene.widgets[].purpose` label when available (optional dependency on `semantic/`'s output, degrading gracefully when absent), else fall back to a positional name (e.g. `param_2`) -- [ ] 3.5 Implement `synthesize_flow_skill(goal, timeline) -> FlowTemplateSkill`: orchestrates extraction → skeleton match → diffing → parameter promotion → returns a candidate skill record (not yet persisted) -- [ ] 3.6 Write unit tests for first-time synthesis (no prior match, zero parameters), second-execution parameter promotion, and identical-repeat synthesis (no spurious new parameters), using canned `TimelineRecord` fixtures +- [x] 3.1 Implement `skills_learning/synthesis.py`'s tool-call extraction: given a `task_id`, read `storage.timeline.Timeline.read(task_id)` and produce an ordered list of `(tool_name, args)` pairs, filtering out read-only tool names (`describe_screen`, `screenshot`, `ui_tree`, `find_text`, `find_icon`) +- [x] 3.2 Implement skeleton matching: given an extracted tool-name sequence, look up any stored skill (via `store.py`) whose latest version has the identical tool-name sequence +- [x] 3.3 Implement cross-execution argument diffing: compare extracted argument values position-by-position against a matched stored version's steps, and promote any differing value into a named `{param}` placeholder plus a corresponding entry in the skill's `parameters` schema +- [x] 3.4 Implement parameter naming: prefer a name derived from the corresponding `SemanticScene.widgets[].purpose` label when available (optional dependency on `semantic/`'s output, degrading gracefully when absent), else fall back to a positional name (e.g. `param_2`) +- [x] 3.5 Implement `synthesize_flow_skill(goal, timeline) -> FlowTemplateSkill`: orchestrates extraction → skeleton match → diffing → parameter promotion → returns a candidate skill record (not yet persisted) +- [x] 3.6 Write unit tests for first-time synthesis (no prior match, zero parameters), second-execution parameter promotion, and identical-repeat synthesis (no spurious new parameters), using canned `TimelineRecord` fixtures ## 4. Version divergence detection (skill-versioning) -- [ ] 4.1 Implement `skills_learning/versioning.py`'s `diff_flow_versions(stored_steps, executed_steps) -> VersionDiff`: detect tool-name-sequence insertion/deletion/reorder (structural divergence) versus argument-value-only differences -- [ ] 4.2 Implement version-bump logic: on structural divergence, construct a new `FlowTemplateSkill` version with incremented `version` and `parent_version_id` set to the prior version's id; on argument-only divergence, update the existing version's parameters in place (no bump) -- [ ] 4.3 Write unit tests: extra/missing/reordered step triggers a version bump; identical-sequence-different-values does not bump but does update parameters; assert prior version records remain retrievable and unmodified after a bump +- [x] 4.1 Implement `skills_learning/versioning.py`'s `diff_flow_versions(stored_steps, executed_steps) -> VersionDiff`: detect tool-name-sequence insertion/deletion/reorder (structural divergence) versus argument-value-only differences +- [x] 4.2 Implement version-bump logic: on structural divergence, construct a new `FlowTemplateSkill` version with incremented `version` and `parent_version_id` set to the prior version's id; on argument-only divergence, update the existing version's parameters in place (no bump) +- [x] 4.3 Write unit tests: extra/missing/reordered step triggers a version bump; identical-sequence-different-values does not bump but does update parameters; assert prior version records remain retrievable and unmodified after a bump ## 5. Post-task synthesis hook wiring -- [ ] 5.1 Add an optional `on_task_succeeded: Callable[[str, str, Timeline], None] | None = None` constructor argument to `TaskRunner` in `runtime/task.py`, invoked exactly once at the end of `run()` when the final status is `succeeded` -- [ ] 5.2 Wire a default hook (when `on_task_succeeded` is not explicitly passed and Skill Authoring is enabled in `skills_learning/config.py`) that calls `synthesis.synthesize_flow_skill()`, runs versioning via `versioning.py`, and persists the result via `store.py` -- [ ] 5.3 Verify that when Skill Authoring is disabled (default) or `on_task_succeeded` is left `None` and disabled, `TaskRunner.run()`'s behavior and return value are byte-for-byte identical to before this change -- [ ] 5.4 Write a unit test that runs a fake successful `TaskRunner` loop with Skill Authoring enabled and asserts a skill record is stored after completion, and a test that asserts no store write occurs when disabled +- [x] 5.1 Add an optional `on_task_succeeded: Callable[[str, str, Timeline], None] | None = None` constructor argument to `TaskRunner` in `runtime/task.py`, invoked exactly once at the end of `run()` when the final status is `succeeded` +- [x] 5.2 Wire a default hook (when `on_task_succeeded` is not explicitly passed and Skill Authoring is enabled in `skills_learning/config.py`) that calls `synthesis.synthesize_flow_skill()`, runs versioning via `versioning.py`, and persists the result via `store.py` +- [x] 5.3 Verify that when Skill Authoring is disabled (default) or `on_task_succeeded` is left `None` and disabled, `TaskRunner.run()`'s behavior and return value are byte-for-byte identical to before this change +- [x] 5.4 Write a unit test that runs a fake successful `TaskRunner` loop with Skill Authoring enabled and asserts a skill record is stored after completion, and a test that asserts no store write occurs when disabled ## 6. Embedding and retrieval (skill-embedding-retrieval) -- [ ] 6.1 Implement `skills_learning/embeddings.py`'s embedding client interface: `embed_skill_text(text) -> list[float] | None`, catching timeout/rate-limit/disabled-config/connection-error internally and returning `None` rather than raising, mirroring `semantic/llm_client.py`'s degrade-safe contract -- [ ] 6.2 Implement a local `skill_embeddings` index (skill id + version → vector, model name, `updated_at`) in `skills_learning/store.py` or a dedicated `skills_learning/embeddings_store.py` -- [ ] 6.3 Wire embedding computation into the post-synthesis/versioning path: call `embed_skill_text()` on `name + description + goal` for every newly stored skill version, storing the resulting vector (or leaving the skill un-embedded if the call returns `None`) -- [ ] 6.4 Implement `skills_learning/retrieval.py`'s `retrieve_candidate_skills(goal, top_k) -> list[ScoredSkill]`: embed the incoming goal, compute cosine similarity against every stored skill embedding, and return the top `top_k` ranked results, skipping skills with no stored embedding -- [ ] 6.5 Write unit tests: ranked ordering for a goal similar to a stored skill's originating goal (using a fake/deterministic embedding function), `top_k` truncation, empty-result case when no skill has an embedding, and a case where an embedding call returns `None` and the skill is stored but excluded from retrieval results +- [x] 6.1 Implement `skills_learning/embeddings.py`'s embedding client interface: `embed_skill_text(text) -> list[float] | None`, catching timeout/rate-limit/disabled-config/connection-error internally and returning `None` rather than raising, mirroring `semantic/llm_client.py`'s degrade-safe contract +- [x] 6.2 Implement a local `skill_embeddings` index (skill id + version → vector, model name, `updated_at`) in `skills_learning/store.py` or a dedicated `skills_learning/embeddings_store.py` +- [x] 6.3 Wire embedding computation into the post-synthesis/versioning path: call `embed_skill_text()` on `name + description + goal` for every newly stored skill version, storing the resulting vector (or leaving the skill un-embedded if the call returns `None`) +- [x] 6.4 Implement `skills_learning/retrieval.py`'s `retrieve_candidate_skills(goal, top_k) -> list[ScoredSkill]`: embed the incoming goal, compute cosine similarity against every stored skill embedding, and return the top `top_k` ranked results, skipping skills with no stored embedding +- [x] 6.5 Write unit tests: ranked ordering for a goal similar to a stored skill's originating goal (using a fake/deterministic embedding function), `top_k` truncation, empty-result case when no skill has an embedding, and a case where an embedding call returns `None` and the skill is stored but excluded from retrieval results ## 7. Integration tests and validation -- [ ] 7.1 Write an end-to-end test: run a fake successful task twice with slightly different goal text/argument values through `TaskRunner` (Skill Authoring enabled, embedding client mocked), asserting the second run produces a new skill version with a promoted parameter and its own embedding -- [ ] 7.2 Write an end-to-end test: run a fake successful task, then call `retrieve_candidate_skills()` with a new, semantically similar goal string, asserting the synthesized skill is returned -- [ ] 7.3 Run the full existing `pytest` suite and confirm zero existing test files require content changes (only new `tests/test_skill_*.py`-style files are added) -- [ ] 7.4 Add a smoke test importing `skills_learning` alongside existing `tests/` smoke coverage, confirming the package has no import-time dependency on `skill-catalog-subscription`'s sync client (only, optionally, its shared model shapes if already implemented) +- [x] 7.1 Write an end-to-end test: run a fake successful task twice with slightly different goal text/argument values through `TaskRunner` (Skill Authoring enabled, embedding client mocked), asserting the second run produces a new skill version with a promoted parameter and its own embedding +- [x] 7.2 Write an end-to-end test: run a fake successful task, then call `retrieve_candidate_skills()` with a new, semantically similar goal string, asserting the synthesized skill is returned +- [x] 7.3 Run the full existing `pytest` suite and confirm zero existing test files require content changes (only new `tests/test_skill_*.py`-style files are added) +- [x] 7.4 Add a smoke test importing `skills_learning` alongside existing `tests/` smoke coverage, confirming the package has no import-time dependency on `skill-catalog-subscription`'s sync client (only, optionally, its shared model shapes if already implemented) diff --git a/pyproject.toml b/pyproject.toml index 8afd8ed..dd9b6de 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -8,6 +8,7 @@ dependencies = [ "Appium-Python-Client>=5.1.1", "fastapi>=0.115.0", "mcp>=1.27,<2", + "openai>=1.0.0", "paddleocr>=3.0.0", "uvicorn[standard]>=0.30.0", ] @@ -31,6 +32,7 @@ include = [ "perception*", "runtime*", "semantic*", + "skills_learning*", "storage*", "tools*", "world*", diff --git a/runtime/task.py b/runtime/task.py index 246c515..916298f 100644 --- a/runtime/task.py +++ b/runtime/task.py @@ -1,5 +1,6 @@ from __future__ import annotations +import logging from collections.abc import Callable from dataclasses import dataclass, replace from inspect import Parameter, signature @@ -9,6 +10,15 @@ from runtime.context import TaskContext from runtime.executor import Executor from runtime.planner import PlannedStep, Planner from semantic.models import SemanticScene +from skills_learning.config import ( + SkillAuthoringConfig, + load_config as load_skill_authoring_config, +) +from skills_learning.embeddings import EmbeddingClient, embed_skill_text +from skills_learning.models import skill_embedding_text +from skills_learning.store import SkillStore, get_default_store +from skills_learning.synthesis import synthesize_flow_skill +from skills_learning.versioning import store_synthesized_skill from storage.task_metadata import TaskMetadataStore from storage.timeline import Timeline from tools.describe_screen import describe_screen @@ -16,6 +26,8 @@ from tools.screenshot import take_screenshot from world.config import WorldConfig, load_config as load_world_config from world.model import WorldModel +logger = logging.getLogger(__name__) + @dataclass class TaskRunnerConfig: @@ -24,6 +36,7 @@ class TaskRunnerConfig: Observer = Callable[[str], Scene] ScreenshotProvider = Callable[[str], bytes] +TaskSucceededHook = Callable[[str, str, Timeline], None] class TaskRunner: @@ -39,6 +52,10 @@ class TaskRunner: screenshot_provider: ScreenshotProvider | None = None, world_model: WorldModel | None = None, world_config: WorldConfig | None = None, + on_task_succeeded: TaskSucceededHook | None = None, + skill_authoring_config: SkillAuthoringConfig | None = None, + skill_store: SkillStore | None = None, + skill_embedding_client: EmbeddingClient | None = None, ) -> None: self.planner = planner or Planner() self.executor = executor or Executor() @@ -56,6 +73,17 @@ class TaskRunner: self.world_model = WorldModel(config=self.world_config) else: self.world_model = None + self.skill_authoring_config = ( + skill_authoring_config or load_skill_authoring_config() + ) + self.skill_store = skill_store + self.skill_embedding_client = skill_embedding_client + if on_task_succeeded is not None: + self.on_task_succeeded = on_task_succeeded + elif self.skill_authoring_config.enabled: + self.on_task_succeeded = self._default_task_succeeded_hook + else: + self.on_task_succeeded = None def run(self, task: Task) -> Task: context = TaskContext(task_id=task.id, goal=task.goal) @@ -72,8 +100,7 @@ class TaskRunner: scene=scene, context=context, ): - self._update_task(task, status="completed", completed=True) - return task + return self._complete_task(task) for step in steps: executable_step = self._step_for_device(step, task.device_id) @@ -101,6 +128,45 @@ class TaskRunner: ) return task + def _complete_task(self, task: Task) -> Task: + self._update_task(task, status="completed", completed=True) + self._notify_task_succeeded(task) + return task + + def _notify_task_succeeded(self, task: Task) -> None: + if self.on_task_succeeded is None or self.timeline is None: + return + try: + self.on_task_succeeded(task.id, task.goal, self.timeline) + except Exception as exc: + logger.info("task succeeded hook failed: %s", exc) + + def _default_task_succeeded_hook( + self, + task_id: str, + goal: str, + timeline: Timeline, + ) -> None: + store = self.skill_store or get_default_store() + candidate = synthesize_flow_skill( + goal, + timeline, + task_id=task_id, + store=store, + ) + stored = store_synthesized_skill(store, candidate).skill + vector = embed_skill_text( + skill_embedding_text(stored), + client=self.skill_embedding_client, + config=self.skill_authoring_config, + ) + if vector is not None: + store.store_embedding( + stored, + vector, + model_name=self.skill_authoring_config.embedding_model, + ) + def _plan( self, goal: str, diff --git a/skills_learning/__init__.py b/skills_learning/__init__.py new file mode 100644 index 0000000..dac905e --- /dev/null +++ b/skills_learning/__init__.py @@ -0,0 +1,16 @@ +"""Local skill learning from completed task timelines.""" + +from skills_learning.config import SkillAuthoringConfig, load_config +from skills_learning.models import FlowStep, FlowTemplateSkill, Skill, SkillMetadata +from skills_learning.store import SkillStore, get_default_store + +__all__ = [ + "FlowStep", + "FlowTemplateSkill", + "Skill", + "SkillAuthoringConfig", + "SkillMetadata", + "SkillStore", + "get_default_store", + "load_config", +] diff --git a/skills_learning/config.py b/skills_learning/config.py new file mode 100644 index 0000000..5f64a5f --- /dev/null +++ b/skills_learning/config.py @@ -0,0 +1,61 @@ +from __future__ import annotations + +import os +from collections.abc import Mapping +from dataclasses import dataclass + +DEFAULT_DIVERGENCE_TOLERANCE = 0.0 +DEFAULT_EMBEDDING_MODEL = "text-embedding-3-small" +DEFAULT_TOP_K = 5 + +ENABLED_ENV = "SKILL_AUTHORING_ENABLED" +DIVERGENCE_TOLERANCE_ENV = "SKILL_DIVERGENCE_TOLERANCE" +EMBEDDING_MODEL_ENV = "SKILL_EMBEDDING_MODEL" +TOP_K_ENV = "SKILL_RETRIEVAL_TOP_K" + + +@dataclass(frozen=True) +class SkillAuthoringConfig: + enabled: bool = False + divergence_tolerance: float = DEFAULT_DIVERGENCE_TOLERANCE + embedding_model: str = DEFAULT_EMBEDDING_MODEL + top_k: int = DEFAULT_TOP_K + + +def load_config(env: Mapping[str, str] | None = None) -> SkillAuthoringConfig: + values = env or os.environ + return SkillAuthoringConfig( + enabled=_parse_bool(values.get(ENABLED_ENV), default=False), + divergence_tolerance=_parse_float( + values.get(DIVERGENCE_TOLERANCE_ENV), + default=DEFAULT_DIVERGENCE_TOLERANCE, + ), + embedding_model=values.get(EMBEDDING_MODEL_ENV) or DEFAULT_EMBEDDING_MODEL, + top_k=_parse_int(values.get(TOP_K_ENV), default=DEFAULT_TOP_K), + ) + + +def _parse_bool(value: str | None, *, default: bool) -> bool: + if value is None: + return default + return value.strip().lower() in {"1", "true", "yes", "on", "enabled"} + + +def _parse_float(value: str | None, *, default: float) -> float: + if value is None: + return default + try: + parsed = float(value) + except ValueError: + return default + return parsed if parsed >= 0 else default + + +def _parse_int(value: str | None, *, default: int) -> int: + if value is None: + return default + try: + parsed = int(value) + except ValueError: + return default + return parsed if parsed > 0 else default diff --git a/skills_learning/embeddings.py b/skills_learning/embeddings.py new file mode 100644 index 0000000..9a7b36d --- /dev/null +++ b/skills_learning/embeddings.py @@ -0,0 +1,44 @@ +from __future__ import annotations + +from typing import Any, Protocol + +from skills_learning.config import SkillAuthoringConfig, load_config + + +class EmbeddingClient(Protocol): + def embed(self, text: str, *, model: str) -> list[float]: + ... + + +class OpenAIEmbeddingClient: + def __init__(self, *, transport: Any | None = None) -> None: + self._transport = transport + + def embed(self, text: str, *, model: str) -> list[float]: + client = self._client() + response = client.embeddings.create(model=model, input=text) + return [float(value) for value in response.data[0].embedding] + + def _client(self) -> Any: + if self._transport is not None: + return self._transport + from openai import OpenAI + + self._transport = OpenAI() + return self._transport + + +def embed_skill_text( + text: str, + *, + client: EmbeddingClient | None = None, + config: SkillAuthoringConfig | None = None, +) -> list[float] | None: + settings = config or load_config() + if not settings.enabled: + return None + try: + embedding_client = client or OpenAIEmbeddingClient() + return embedding_client.embed(text, model=settings.embedding_model) + except Exception: + return None diff --git a/skills_learning/models.py b/skills_learning/models.py new file mode 100644 index 0000000..ff420af --- /dev/null +++ b/skills_learning/models.py @@ -0,0 +1,186 @@ +from __future__ import annotations + +from dataclasses import dataclass, field, replace +from datetime import datetime +from typing import Any, Literal +from uuid import uuid4 + +from core.models import utc_now + +LOCAL_SYNTHESIS_SOURCE = "local-synthesis" +SkillKind = Literal["knowledge", "flow_template"] + + +@dataclass(frozen=True) +class SkillMetadata: + id: str = field(default_factory=lambda: uuid4().hex) + name: str = "" + description: str = "" + kind: SkillKind = "flow_template" + tags: list[str] = field(default_factory=list) + source: str = LOCAL_SYNTHESIS_SOURCE + version: int = 1 + parent_version_id: str | None = None + originating_goal: str | None = None + created_at: datetime = field(default_factory=utc_now) + updated_at: datetime = field(default_factory=utc_now) + + def to_dict(self) -> dict[str, Any]: + return { + "id": self.id, + "name": self.name, + "description": self.description, + "kind": self.kind, + "tags": list(self.tags), + "source": self.source, + "version": self.version, + "parent_version_id": self.parent_version_id, + "originating_goal": self.originating_goal, + "created_at": self.created_at.isoformat(), + "updated_at": self.updated_at.isoformat(), + } + + @classmethod + def from_dict(cls, data: dict[str, Any]) -> "SkillMetadata": + return cls( + id=str(data.get("id") or uuid4().hex), + name=str(data.get("name") or ""), + description=str(data.get("description") or ""), + kind=data.get("kind") or "flow_template", + tags=[str(tag) for tag in data.get("tags", [])], + source=str(data.get("source") or LOCAL_SYNTHESIS_SOURCE), + version=int(data.get("version") or 1), + parent_version_id=data.get("parent_version_id"), + originating_goal=data.get("originating_goal"), + created_at=_parse_datetime(data.get("created_at")), + updated_at=_parse_datetime(data.get("updated_at")), + ) + + +@dataclass(frozen=True) +class FlowStep: + tool_name: str + args: dict[str, Any] = field(default_factory=dict) + + def to_dict(self) -> dict[str, Any]: + return { + "tool_name": self.tool_name, + "args": dict(self.args), + } + + @classmethod + def from_dict(cls, data: dict[str, Any]) -> "FlowStep": + return cls( + tool_name=str(data.get("tool_name") or data.get("action") or ""), + args=dict(data.get("args") or {}), + ) + + +@dataclass(frozen=True) +class Skill: + metadata: SkillMetadata + + @property + def id(self) -> str: + return self.metadata.id + + @property + def name(self) -> str: + return self.metadata.name + + @property + def description(self) -> str: + return self.metadata.description + + @property + def source(self) -> str: + return self.metadata.source + + @property + def version(self) -> int: + return self.metadata.version + + @property + def parent_version_id(self) -> str | None: + return self.metadata.parent_version_id + + @property + def originating_goal(self) -> str | None: + return self.metadata.originating_goal + + +@dataclass(frozen=True) +class FlowTemplateSkill(Skill): + steps: list[FlowStep] = field(default_factory=list) + parameters: dict[str, dict[str, Any]] = field(default_factory=dict) + + def to_dict(self) -> dict[str, Any]: + return { + **self.metadata.to_dict(), + "steps": [step.to_dict() for step in self.steps], + "parameters": { + name: dict(schema) + for name, schema in self.parameters.items() + }, + } + + @classmethod + def from_dict(cls, data: dict[str, Any]) -> "FlowTemplateSkill": + return cls( + metadata=SkillMetadata.from_dict(data), + steps=[ + FlowStep.from_dict(step) + for step in data.get("steps", []) + ], + parameters={ + str(name): dict(schema) + for name, schema in (data.get("parameters") or {}).items() + }, + ) + + def with_metadata(self, **changes: Any) -> "FlowTemplateSkill": + return replace(self, metadata=replace(self.metadata, **changes)) + + def with_updates( + self, + *, + steps: list[FlowStep] | None = None, + parameters: dict[str, dict[str, Any]] | None = None, + **metadata_changes: Any, + ) -> "FlowTemplateSkill": + metadata = replace( + self.metadata, + updated_at=utc_now(), + **metadata_changes, + ) + return replace( + self, + metadata=metadata, + steps=list(steps) if steps is not None else list(self.steps), + parameters={ + name: dict(schema) + for name, schema in ( + parameters if parameters is not None else self.parameters + ).items() + }, + ) + + +def skill_embedding_text(skill: FlowTemplateSkill) -> str: + goal = skill.originating_goal or "" + return f"{skill.name}: {skill.description}\nOriginal goal: {goal}" + + +def clone_skill(skill: FlowTemplateSkill) -> FlowTemplateSkill: + return FlowTemplateSkill.from_dict(skill.to_dict()) + + +def _parse_datetime(value: Any) -> datetime: + if isinstance(value, datetime): + return value + if isinstance(value, str): + try: + return datetime.fromisoformat(value) + except ValueError: + pass + return utc_now() diff --git a/skills_learning/retrieval.py b/skills_learning/retrieval.py new file mode 100644 index 0000000..037a6b8 --- /dev/null +++ b/skills_learning/retrieval.py @@ -0,0 +1,55 @@ +from __future__ import annotations + +import math +from dataclasses import dataclass + +from skills_learning.config import SkillAuthoringConfig, load_config +from skills_learning.embeddings import EmbeddingClient, embed_skill_text +from skills_learning.models import FlowTemplateSkill +from skills_learning.store import SkillStore, get_default_store + + +@dataclass(frozen=True) +class ScoredSkill: + skill: FlowTemplateSkill + score: float + + +def retrieve_candidate_skills( + goal: str, + *, + store: SkillStore | None = None, + top_k: int | None = None, + embedding_client: EmbeddingClient | None = None, + config: SkillAuthoringConfig | None = None, +) -> list[ScoredSkill]: + settings = config or load_config() + skill_store = store or get_default_store() + query_vector = embed_skill_text( + goal, + client=embedding_client, + config=settings, + ) + if query_vector is None: + return [] + + scored: list[ScoredSkill] = [] + for record in skill_store.list_embeddings(): + skill = skill_store.get_by_id(record.skill_id) + if skill is None: + continue + score = _cosine_similarity(query_vector, record.vector) + scored.append(ScoredSkill(skill=skill, score=score)) + scored.sort(key=lambda item: item.score, reverse=True) + return scored[: top_k or settings.top_k] + + +def _cosine_similarity(left: list[float], right: list[float]) -> float: + if not left or not right or len(left) != len(right): + return 0.0 + dot = sum(a * b for a, b in zip(left, right, strict=True)) + left_norm = math.sqrt(sum(value * value for value in left)) + right_norm = math.sqrt(sum(value * value for value in right)) + if left_norm == 0 or right_norm == 0: + return 0.0 + return dot / (left_norm * right_norm) diff --git a/skills_learning/store.py b/skills_learning/store.py new file mode 100644 index 0000000..603e40c --- /dev/null +++ b/skills_learning/store.py @@ -0,0 +1,141 @@ +from __future__ import annotations + +from dataclasses import dataclass, field, replace +from datetime import datetime +from typing import Any +from uuid import uuid4 + +from core.models import utc_now +from skills_learning.models import ( + LOCAL_SYNTHESIS_SOURCE, + FlowTemplateSkill, + clone_skill, +) + + +@dataclass(frozen=True) +class SkillEmbeddingRecord: + skill_id: str + version: int + vector: list[float] + model_name: str + updated_at: datetime = field(default_factory=utc_now) + + def to_dict(self) -> dict[str, Any]: + return { + "skill_id": self.skill_id, + "version": self.version, + "vector": list(self.vector), + "model_name": self.model_name, + "updated_at": self.updated_at.isoformat(), + } + + +class SkillStore: + def __init__(self) -> None: + self._skills: dict[str, FlowTemplateSkill] = {} + self._embeddings: dict[tuple[str, int], SkillEmbeddingRecord] = {} + + def create_version( + self, + skill: FlowTemplateSkill, + *, + parent: FlowTemplateSkill | None = None, + ) -> FlowTemplateSkill: + version = parent.version + 1 if parent else self._next_version(skill.name) + stored = skill.with_metadata( + id=uuid4().hex, + source=LOCAL_SYNTHESIS_SOURCE, + version=version, + parent_version_id=parent.id if parent else skill.parent_version_id, + created_at=utc_now(), + updated_at=utc_now(), + ) + self._skills[stored.id] = clone_skill(stored) + return clone_skill(stored) + + def update_skill(self, skill: FlowTemplateSkill) -> FlowTemplateSkill: + if skill.id not in self._skills: + raise KeyError(f"unknown skill {skill.id}") + stored = skill.with_metadata( + source=LOCAL_SYNTHESIS_SOURCE, + updated_at=utc_now(), + ) + self._skills[stored.id] = clone_skill(stored) + return clone_skill(stored) + + def get_by_id(self, skill_id: str) -> FlowTemplateSkill | None: + skill = self._skills.get(skill_id) + return clone_skill(skill) if skill else None + + def get_latest_by_name(self, name: str) -> FlowTemplateSkill | None: + versions = self.list_versions(name) + return versions[-1] if versions else None + + def list_versions(self, name: str) -> list[FlowTemplateSkill]: + return sorted( + [ + clone_skill(skill) + for skill in self._skills.values() + if skill.name == name + ], + key=lambda skill: skill.version, + ) + + def list_latest(self) -> list[FlowTemplateSkill]: + latest: dict[str, FlowTemplateSkill] = {} + for skill in self._skills.values(): + current = latest.get(skill.name) + if current is None or skill.version > current.version: + latest[skill.name] = skill + return [clone_skill(skill) for skill in latest.values()] + + def list_all(self) -> list[FlowTemplateSkill]: + return [clone_skill(skill) for skill in self._skills.values()] + + def store_embedding( + self, + skill: FlowTemplateSkill, + vector: list[float], + *, + model_name: str, + ) -> SkillEmbeddingRecord: + record = SkillEmbeddingRecord( + skill_id=skill.id, + version=skill.version, + vector=[float(value) for value in vector], + model_name=model_name, + ) + self._embeddings[(record.skill_id, record.version)] = record + return replace(record, vector=list(record.vector)) + + def get_embedding( + self, + skill_id: str, + version: int, + ) -> SkillEmbeddingRecord | None: + record = self._embeddings.get((skill_id, version)) + return replace(record, vector=list(record.vector)) if record else None + + def list_embeddings(self) -> list[SkillEmbeddingRecord]: + return [ + replace(record, vector=list(record.vector)) + for record in self._embeddings.values() + ] + + def _next_version(self, name: str) -> int: + latest = self.get_latest_by_name(name) + return latest.version + 1 if latest else 1 + + +_DEFAULT_STORE = SkillStore() + + +def get_default_store() -> SkillStore: + return _DEFAULT_STORE + + +def reset_default_store() -> SkillStore: + global _DEFAULT_STORE + _DEFAULT_STORE = SkillStore() + return _DEFAULT_STORE diff --git a/skills_learning/synthesis.py b/skills_learning/synthesis.py new file mode 100644 index 0000000..6520cb1 --- /dev/null +++ b/skills_learning/synthesis.py @@ -0,0 +1,248 @@ +from __future__ import annotations + +import re +from collections.abc import Iterable +from typing import Any +from uuid import uuid4 + +from skills_learning.models import ( + LOCAL_SYNTHESIS_SOURCE, + FlowStep, + FlowTemplateSkill, + SkillMetadata, +) +from skills_learning.store import SkillStore + +READ_ONLY_TOOL_NAMES = { + "describe_screen", + "describe_screen_semantic", + "screenshot", + "take_screenshot", + "ui_tree", + "get_ui_tree", + "find_text", + "find_text_on_screen", + "find_icon", + "find_icon_on_screen", +} + + +def extract_tool_calls( + task_id: str, + timeline: Any, +) -> list[FlowStep]: + records = _timeline_records(timeline, task_id) + steps: list[FlowStep] = [] + for record in records: + tool_call = _record_value(record, "tool_call") + if not isinstance(tool_call, dict): + continue + tool_name = str(tool_call.get("action") or tool_call.get("tool_name") or "") + if not tool_name or tool_name in READ_ONLY_TOOL_NAMES: + continue + steps.append(FlowStep(tool_name=tool_name, args=dict(tool_call.get("args") or {}))) + return steps + + +def find_matching_skeleton( + steps: list[FlowStep], + store: SkillStore | None, +) -> FlowTemplateSkill | None: + if store is None: + return None + sequence = _tool_sequence(steps) + for skill in store.list_latest(): + if _tool_sequence(skill.steps) == sequence: + return skill + return None + + +def synthesize_flow_skill( + goal: str, + timeline: Any, + *, + task_id: str = "", + store: SkillStore | None = None, +) -> FlowTemplateSkill: + records = _timeline_records(timeline, task_id) + executed_steps = extract_tool_calls(task_id, records) + matched = find_matching_skeleton(executed_steps, store) + + if matched is None: + steps = executed_steps + parameters: dict[str, dict[str, Any]] = {} + name = _skill_name_from_goal(goal) + else: + steps, parameters = promote_parameters( + matched.steps, + executed_steps, + records=records, + existing_parameters=matched.parameters, + ) + name = matched.name + + return FlowTemplateSkill( + metadata=SkillMetadata( + id=uuid4().hex, + name=name, + description=f"Learned flow for: {goal}", + kind="flow_template", + source=LOCAL_SYNTHESIS_SOURCE, + version=matched.version if matched else 1, + parent_version_id=matched.parent_version_id if matched else None, + originating_goal=goal, + ), + steps=steps, + parameters=parameters, + ) + + +def promote_parameters( + stored_steps: list[FlowStep], + executed_steps: list[FlowStep], + *, + records: Iterable[Any] = (), + existing_parameters: dict[str, dict[str, Any]] | None = None, +) -> tuple[list[FlowStep], dict[str, dict[str, Any]]]: + parameters = { + name: dict(schema) + for name, schema in (existing_parameters or {}).items() + } + parameterized_steps = [ + FlowStep(step.tool_name, dict(step.args)) + for step in executed_steps + ] + used_names = set(parameters) + parameter_index = len(used_names) + 1 + + for step_index, (stored_step, executed_step) in enumerate( + zip(stored_steps, executed_steps, strict=False) + ): + for arg_name, executed_value in executed_step.args.items(): + stored_value = stored_step.args.get(arg_name) + if stored_value == executed_value: + continue + if _is_placeholder(stored_value): + parameterized_steps[step_index].args[arg_name] = stored_value + continue + if _is_placeholder(executed_value): + continue + + preferred_name = _parameter_name_from_semantic_record( + list(records), + step_index, + ) + parameter_name = _unique_parameter_name( + preferred_name or f"param_{parameter_index}", + used_names, + ) + used_names.add(parameter_name) + parameter_index += 1 + parameterized_steps[step_index].args[arg_name] = f"{{{parameter_name}}}" + parameters[parameter_name] = _parameter_schema( + arg_name, + executed_value, + preferred_name is not None, + ) + + return parameterized_steps, parameters + + +def _timeline_records(timeline: Any, task_id: str) -> list[Any]: + if isinstance(timeline, list): + return list(timeline) + if hasattr(timeline, "read"): + return list(timeline.read(task_id)) + return list(timeline) + + +def _record_value(record: Any, key: str) -> Any: + if isinstance(record, dict): + return record.get(key) + return getattr(record, key, None) + + +def _tool_sequence(steps: list[FlowStep]) -> list[str]: + return [step.tool_name for step in steps] + + +def _is_placeholder(value: Any) -> bool: + return isinstance(value, str) and value.startswith("{") and value.endswith("}") + + +def _parameter_schema( + arg_name: str, + value: Any, + from_semantic_label: bool, +) -> dict[str, Any]: + return { + "type": _json_type(value), + "description": ( + f"Value for {arg_name} inferred from a semantic widget label" + if from_semantic_label + else f"Value for {arg_name}" + ), + } + + +def _json_type(value: Any) -> str: + if isinstance(value, bool): + return "boolean" + if isinstance(value, int | float): + return "number" + if isinstance(value, list): + return "array" + if isinstance(value, dict): + return "object" + return "string" + + +def _parameter_name_from_semantic_record( + records: list[Any], + step_index: int, +) -> str | None: + if step_index >= len(records): + return None + result = _record_value(records[step_index], "result") + semantic_scene = _semantic_scene_payload(result) + if not isinstance(semantic_scene, dict): + return None + widgets = semantic_scene.get("widgets") + if not isinstance(widgets, list): + return None + for widget in widgets: + if isinstance(widget, dict) and widget.get("purpose"): + return _sanitize_name(str(widget["purpose"])) + return None + + +def _semantic_scene_payload(result: Any) -> Any: + if not isinstance(result, dict): + return None + if "semantic_scene" in result: + return result["semantic_scene"] + nested = result.get("result") + if isinstance(nested, dict): + return nested.get("semantic_scene") + return None + + +def _unique_parameter_name(name: str, used_names: set[str]) -> str: + candidate = _sanitize_name(name) or "param" + if candidate not in used_names: + return candidate + index = 2 + while f"{candidate}_{index}" in used_names: + index += 1 + return f"{candidate}_{index}" + + +def _sanitize_name(value: str) -> str: + sanitized = re.sub(r"[^0-9a-zA-Z]+", "_", value.strip().lower()).strip("_") + if sanitized and sanitized[0].isdigit(): + return f"param_{sanitized}" + return sanitized + + +def _skill_name_from_goal(goal: str) -> str: + return goal.strip() or "learned flow" diff --git a/skills_learning/versioning.py b/skills_learning/versioning.py new file mode 100644 index 0000000..e709498 --- /dev/null +++ b/skills_learning/versioning.py @@ -0,0 +1,72 @@ +from __future__ import annotations + +from dataclasses import dataclass +from typing import Any + +from skills_learning.models import FlowStep, FlowTemplateSkill +from skills_learning.store import SkillStore + + +@dataclass(frozen=True) +class VersionDiff: + structural_divergence: bool + stored_sequence: list[str] + executed_sequence: list[str] + argument_differences: list[tuple[int, str, Any, Any]] + + +@dataclass(frozen=True) +class VersioningResult: + skill: FlowTemplateSkill + created_new_version: bool + + +def diff_flow_versions( + stored_steps: list[FlowStep], + executed_steps: list[FlowStep], +) -> VersionDiff: + stored_sequence = [step.tool_name for step in stored_steps] + executed_sequence = [step.tool_name for step in executed_steps] + structural = stored_sequence != executed_sequence + differences: list[tuple[int, str, Any, Any]] = [] + if not structural: + for index, (stored_step, executed_step) in enumerate( + zip(stored_steps, executed_steps, strict=True) + ): + keys = set(stored_step.args) | set(executed_step.args) + for key in sorted(keys): + stored_value = stored_step.args.get(key) + executed_value = executed_step.args.get(key) + if stored_value != executed_value: + differences.append((index, key, stored_value, executed_value)) + return VersionDiff( + structural_divergence=structural, + stored_sequence=stored_sequence, + executed_sequence=executed_sequence, + argument_differences=differences, + ) + + +def store_synthesized_skill( + store: SkillStore, + candidate: FlowTemplateSkill, +) -> VersioningResult: + latest = store.get_latest_by_name(candidate.name) + if latest is None: + return VersioningResult(store.create_version(candidate), True) + + diff = diff_flow_versions(latest.steps, candidate.steps) + if diff.structural_divergence: + return VersioningResult(store.create_version(candidate, parent=latest), True) + + merged_parameters = { + **latest.parameters, + **candidate.parameters, + } + updated = latest.with_updates( + steps=candidate.steps, + parameters=merged_parameters, + description=candidate.description, + originating_goal=candidate.originating_goal, + ) + return VersioningResult(store.update_skill(updated), False) diff --git a/tests/test_skill_retrieval.py b/tests/test_skill_retrieval.py new file mode 100644 index 0000000..b4aeb4f --- /dev/null +++ b/tests/test_skill_retrieval.py @@ -0,0 +1,91 @@ +from __future__ import annotations + +from skills_learning.config import SkillAuthoringConfig +from skills_learning.embeddings import embed_skill_text +from skills_learning.models import FlowStep, FlowTemplateSkill, SkillMetadata +from skills_learning.retrieval import retrieve_candidate_skills +from skills_learning.store import SkillStore + + +class FakeEmbeddingClient: + def __init__(self, *, fail: bool = False) -> None: + self.fail = fail + + def embed(self, text: str, *, model: str) -> list[float]: + if self.fail: + raise TimeoutError("embedding timeout") + lowered = text.lower() + if "coffee" in lowered: + return [1.0, 0.0] + if "tea" in lowered: + return [0.0, 1.0] + return [0.5, 0.5] + + +def _skill(name: str, goal: str) -> FlowTemplateSkill: + return FlowTemplateSkill( + metadata=SkillMetadata( + name=name, + description=f"Learned flow for {goal}", + originating_goal=goal, + ), + steps=[FlowStep("input_text", {"text": goal})], + parameters={}, + ) + + +def test_embed_skill_text_degrades_to_none_on_provider_failure() -> None: + result = embed_skill_text( + "coffee", + client=FakeEmbeddingClient(fail=True), + config=SkillAuthoringConfig(enabled=True), + ) + + assert result is None + + +def test_retrieve_candidate_skills_ranks_by_similarity_and_truncates_top_k() -> None: + store = SkillStore() + coffee = store.create_version(_skill("search coffee", "search coffee")) + tea = store.create_version(_skill("search tea", "search tea")) + store.store_embedding(coffee, [1.0, 0.0], model_name="fake") + store.store_embedding(tea, [0.0, 1.0], model_name="fake") + + results = retrieve_candidate_skills( + "find coffee", + store=store, + top_k=1, + embedding_client=FakeEmbeddingClient(), + config=SkillAuthoringConfig(enabled=True), + ) + + assert [result.skill.name for result in results] == ["search coffee"] + + +def test_retrieve_candidate_skills_returns_empty_without_embeddings() -> None: + store = SkillStore() + store.create_version(_skill("search coffee", "search coffee")) + + results = retrieve_candidate_skills( + "find coffee", + store=store, + embedding_client=FakeEmbeddingClient(), + config=SkillAuthoringConfig(enabled=True), + ) + + assert results == [] + + +def test_skill_without_embedding_is_stored_but_excluded_from_retrieval() -> None: + store = SkillStore() + skill = store.create_version(_skill("search coffee", "search coffee")) + + results = retrieve_candidate_skills( + "find coffee", + store=store, + embedding_client=FakeEmbeddingClient(), + config=SkillAuthoringConfig(enabled=True), + ) + + assert store.get_by_id(skill.id) == skill + assert results == [] diff --git a/tests/test_skill_store.py b/tests/test_skill_store.py new file mode 100644 index 0000000..a4b1baf --- /dev/null +++ b/tests/test_skill_store.py @@ -0,0 +1,54 @@ +from __future__ import annotations + +import sys + +from skills_learning.models import ( + LOCAL_SYNTHESIS_SOURCE, + FlowStep, + FlowTemplateSkill, + SkillMetadata, +) +from skills_learning.store import SkillStore + + +def _skill( + *, + name: str = "search web", + steps: list[FlowStep] | None = None, + parameters: dict[str, dict[str, object]] | None = None, +) -> FlowTemplateSkill: + return FlowTemplateSkill( + metadata=SkillMetadata( + name=name, + description=f"Learned flow for {name}", + source="external-source", + originating_goal=name, + ), + steps=steps or [FlowStep("input_text", {"text": "coffee"})], + parameters=parameters or {}, + ) + + +def test_skill_store_enforces_local_synthesis_source_without_catalog_write_path() -> None: + store = SkillStore() + + stored = store.create_version(_skill()) + + assert stored.source == LOCAL_SYNTHESIS_SOURCE + assert "skills.catalog" not in sys.modules + + +def test_skill_store_keeps_version_chain_and_returns_latest_by_name() -> None: + store = SkillStore() + first = store.create_version(_skill()) + second = store.create_version( + _skill(steps=[FlowStep("input_text", {"text": "tea"})]), + parent=first, + ) + + assert first.version == 1 + assert second.version == 2 + assert second.parent_version_id == first.id + assert store.get_latest_by_name("search web") == second + assert store.get_by_id(first.id) == first + assert [skill.version for skill in store.list_versions("search web")] == [1, 2] diff --git a/tests/test_skill_synthesis.py b/tests/test_skill_synthesis.py new file mode 100644 index 0000000..9143742 --- /dev/null +++ b/tests/test_skill_synthesis.py @@ -0,0 +1,113 @@ +from __future__ import annotations + +from skills_learning.models import FlowStep, FlowTemplateSkill, SkillMetadata +from skills_learning.store import SkillStore +from skills_learning.synthesis import extract_tool_calls, synthesize_flow_skill + + +def _record( + action: str, + args: dict[str, object] | None = None, + *, + result: dict[str, object] | None = None, +) -> dict[str, object]: + return { + "tool_call": { + "action": action, + "args": args or {}, + }, + "result": result or {}, + } + + +def test_extract_tool_calls_filters_read_only_tools_in_order() -> None: + records = [ + _record("describe_screen"), + _record("launch_app", {"app_id": "com.example"}), + _record("screenshot"), + _record("tap", {"x": 1, "y": 2}), + _record("find_text", {"query": "Send"}), + _record("input_text", {"text": "coffee"}), + ] + + steps = extract_tool_calls("", records) + + assert steps == [ + FlowStep("launch_app", {"app_id": "com.example"}), + FlowStep("tap", {"x": 1, "y": 2}), + FlowStep("input_text", {"text": "coffee"}), + ] + + +def test_first_time_synthesis_has_literal_steps_and_no_parameters() -> None: + skill = synthesize_flow_skill( + "search coffee", + [_record("input_text", {"text": "coffee"})], + ) + + assert skill.name == "search coffee" + assert skill.steps == [FlowStep("input_text", {"text": "coffee"})] + assert skill.parameters == {} + + +def test_second_execution_promotes_differing_argument_to_parameter() -> None: + store = SkillStore() + store.create_version( + FlowTemplateSkill( + metadata=SkillMetadata( + name="search", + description="Learned search", + originating_goal="search coffee", + ), + steps=[FlowStep("input_text", {"text": "coffee"})], + parameters={}, + ) + ) + records = [ + _record( + "input_text", + {"text": "tea"}, + result={ + "semantic_scene": { + "page": "Search", + "intents": ["search"], + "widgets": [ + { + "element_id": "search-input", + "purpose": "search field", + } + ], + } + }, + ) + ] + + skill = synthesize_flow_skill("search tea", records, store=store) + + assert skill.name == "search" + assert skill.steps == [FlowStep("input_text", {"text": "{search_field}"})] + assert "search_field" in skill.parameters + + +def test_identical_repeat_does_not_add_parameters() -> None: + store = SkillStore() + store.create_version( + FlowTemplateSkill( + metadata=SkillMetadata( + name="search", + description="Learned search", + originating_goal="search coffee", + ), + steps=[FlowStep("input_text", {"text": "coffee"})], + parameters={}, + ) + ) + + skill = synthesize_flow_skill( + "search coffee again", + [_record("input_text", {"text": "coffee"})], + store=store, + ) + + assert skill.steps == [FlowStep("input_text", {"text": "coffee"})] + assert skill.parameters == {} diff --git a/tests/test_skill_task_runner.py b/tests/test_skill_task_runner.py new file mode 100644 index 0000000..dc9a4bd --- /dev/null +++ b/tests/test_skill_task_runner.py @@ -0,0 +1,217 @@ +from __future__ import annotations + +from core.models import Bounds, Scene, SceneElement, Task +from runtime.executor import Executor, ExecutorConfig +from runtime.planner import PlannedStep, Planner +from runtime.task import TaskRunner, TaskRunnerConfig +from skills_learning.config import SkillAuthoringConfig +from skills_learning.retrieval import retrieve_candidate_skills +from skills_learning.store import SkillStore +from storage.artifact_store import ArtifactStore +from storage.timeline import Timeline +from tests.fakes import PNG_10X20 + + +class ScriptedPlanner(Planner): + def __init__(self, steps: list[PlannedStep]) -> None: + self.steps = steps + + def plan(self, *, goal, scene, context): + if len(context.step_results) >= len(self.steps): + return [] + return [self.steps[len(context.step_results)]] + + def goal_reached(self, *, goal, scene, context): + return len(context.step_results) >= len(self.steps) and all( + result.success for result in context.step_results + ) + + +class FakeEmbeddingClient: + def __init__(self, *, fail: bool = False) -> None: + self.fail = fail + + def embed(self, text: str, *, model: str) -> list[float]: + if self.fail: + raise TimeoutError("embedding timeout") + lowered = text.lower() + if "coffee" in lowered: + return [1.0, 0.0] + if "tea" in lowered: + return [0.0, 1.0] + return [0.5, 0.5] + + +def _scene() -> Scene: + return Scene( + width=10, + height=20, + elements=[ + SceneElement( + id="input", + type="input", + text="Search", + bounds=Bounds(1, 2, 3, 4), + ) + ], + ) + + +def _runner( + *, + tmp_path, + planner: Planner, + timeline: Timeline | None = None, + store: SkillStore | None = None, + skill_config: SkillAuthoringConfig | None = None, + embedding_client: FakeEmbeddingClient | None = None, + on_task_succeeded=None, +) -> TaskRunner: + return TaskRunner( + planner=planner, + executor=Executor( + tools={"input_text": lambda **kwargs: {"ok": True, **kwargs}}, + config=ExecutorConfig(max_retries=1, backoff_seconds=0), + ), + timeline=timeline or Timeline(ArtifactStore(tmp_path / "history")), + config=TaskRunnerConfig(max_steps=5), + observer=lambda device_id: _scene(), + screenshot_provider=lambda device_id: PNG_10X20, + skill_store=store, + skill_authoring_config=skill_config, + skill_embedding_client=embedding_client, + on_task_succeeded=on_task_succeeded, + ) + + +def test_task_runner_calls_explicit_success_hook_once(tmp_path) -> None: + calls: list[tuple[str, str, Timeline]] = [] + timeline = Timeline(ArtifactStore(tmp_path / "history")) + runner = _runner( + tmp_path=tmp_path, + timeline=timeline, + planner=ScriptedPlanner( + [PlannedStep("input_text", "type", {"text": "coffee"})] + ), + on_task_succeeded=lambda task_id, goal, timeline: calls.append( + (task_id, goal, timeline) + ), + ) + task = Task(goal="search coffee", device_id="phone") + + result = runner.run(task) + + assert result.status == "completed" + assert calls == [(task.id, "search coffee", timeline)] + + +def test_task_runner_skill_authoring_disabled_by_default_writes_no_skill(tmp_path) -> None: + store = SkillStore() + runner = _runner( + tmp_path=tmp_path, + planner=ScriptedPlanner( + [PlannedStep("input_text", "type", {"text": "coffee"})] + ), + store=store, + ) + + result = runner.run(Task(goal="search coffee", device_id="phone")) + + assert result.status == "completed" + assert store.list_all() == [] + + +def test_task_runner_skill_authoring_enabled_stores_skill_and_embedding(tmp_path) -> None: + store = SkillStore() + runner = _runner( + tmp_path=tmp_path, + planner=ScriptedPlanner( + [PlannedStep("input_text", "type", {"text": "coffee"})] + ), + store=store, + skill_config=SkillAuthoringConfig(enabled=True, embedding_model="fake"), + embedding_client=FakeEmbeddingClient(), + ) + + result = runner.run(Task(goal="search coffee", device_id="phone")) + + assert result.status == "completed" + skill = store.get_latest_by_name("search coffee") + assert skill is not None + assert skill.steps[0].args == {"text": "coffee"} + assert store.get_embedding(skill.id, skill.version) is not None + + +def test_task_runner_embedding_failure_still_stores_skill_without_embedding(tmp_path) -> None: + store = SkillStore() + runner = _runner( + tmp_path=tmp_path, + planner=ScriptedPlanner( + [PlannedStep("input_text", "type", {"text": "coffee"})] + ), + store=store, + skill_config=SkillAuthoringConfig(enabled=True, embedding_model="fake"), + embedding_client=FakeEmbeddingClient(fail=True), + ) + + result = runner.run(Task(goal="search coffee", device_id="phone")) + + assert result.status == "completed" + skill = store.get_latest_by_name("search coffee") + assert skill is not None + assert store.get_embedding(skill.id, skill.version) is None + + +def test_two_successful_tasks_promote_parameter_and_store_embedding(tmp_path) -> None: + store = SkillStore() + config = SkillAuthoringConfig(enabled=True, embedding_model="fake") + + _runner( + tmp_path=tmp_path, + planner=ScriptedPlanner( + [PlannedStep("input_text", "type coffee", {"text": "coffee"})] + ), + store=store, + skill_config=config, + embedding_client=FakeEmbeddingClient(), + ).run(Task(goal="search coffee", device_id="phone")) + + _runner( + tmp_path=tmp_path, + planner=ScriptedPlanner( + [PlannedStep("input_text", "type tea", {"text": "tea"})] + ), + store=store, + skill_config=config, + embedding_client=FakeEmbeddingClient(), + ).run(Task(goal="search tea", device_id="phone")) + + skill = store.get_latest_by_name("search coffee") + assert skill is not None + assert skill.version == 1 + assert skill.steps[0].args == {"text": "{param_1}"} + assert "param_1" in skill.parameters + assert store.get_embedding(skill.id, skill.version) is not None + + +def test_retrieve_candidate_skills_after_successful_task(tmp_path) -> None: + store = SkillStore() + config = SkillAuthoringConfig(enabled=True, embedding_model="fake") + _runner( + tmp_path=tmp_path, + planner=ScriptedPlanner( + [PlannedStep("input_text", "type coffee", {"text": "coffee"})] + ), + store=store, + skill_config=config, + embedding_client=FakeEmbeddingClient(), + ).run(Task(goal="search coffee", device_id="phone")) + + results = retrieve_candidate_skills( + "find coffee", + store=store, + embedding_client=FakeEmbeddingClient(), + config=config, + ) + + assert [result.skill.name for result in results] == ["search coffee"] diff --git a/tests/test_skill_versioning.py b/tests/test_skill_versioning.py new file mode 100644 index 0000000..9621c18 --- /dev/null +++ b/tests/test_skill_versioning.py @@ -0,0 +1,72 @@ +from __future__ import annotations + +from skills_learning.models import FlowStep, FlowTemplateSkill, SkillMetadata +from skills_learning.store import SkillStore +from skills_learning.versioning import diff_flow_versions, store_synthesized_skill + + +def _skill( + *, + name: str = "search", + steps: list[FlowStep] | None = None, + parameters: dict[str, dict[str, object]] | None = None, +) -> FlowTemplateSkill: + return FlowTemplateSkill( + metadata=SkillMetadata( + name=name, + description="Learned search", + originating_goal="search coffee", + ), + steps=steps or [FlowStep("input_text", {"text": "coffee"})], + parameters=parameters or {}, + ) + + +def test_diff_flow_versions_detects_extra_missing_and_reordered_steps() -> None: + stored = [FlowStep("tap"), FlowStep("input_text")] + + assert diff_flow_versions(stored, [FlowStep("tap")]).structural_divergence + assert diff_flow_versions( + stored, + [FlowStep("tap"), FlowStep("input_text"), FlowStep("tap")], + ).structural_divergence + assert diff_flow_versions( + stored, + [FlowStep("input_text"), FlowStep("tap")], + ).structural_divergence + + +def test_argument_only_difference_updates_existing_version_without_bump() -> None: + store = SkillStore() + first = store.create_version(_skill()) + candidate = _skill( + steps=[FlowStep("input_text", {"text": "{search_query}"})], + parameters={"search_query": {"type": "string"}}, + ) + + result = store_synthesized_skill(store, candidate) + + assert result.created_new_version is False + assert result.skill.version == first.version + assert result.skill.id == first.id + assert result.skill.parameters == {"search_query": {"type": "string"}} + assert store.get_latest_by_name("search") == result.skill + + +def test_structural_divergence_creates_new_version_and_preserves_parent() -> None: + store = SkillStore() + first = store.create_version(_skill()) + candidate = _skill( + steps=[ + FlowStep("tap", {"x": 1}), + FlowStep("input_text", {"text": "coffee"}), + ] + ) + + result = store_synthesized_skill(store, candidate) + + assert result.created_new_version is True + assert result.skill.version == 2 + assert result.skill.parent_version_id == first.id + assert store.get_by_id(first.id) == first + assert store.get_latest_by_name("search") == result.skill diff --git a/tests/test_smoke.py b/tests/test_smoke.py index 4a6b329..d60e13e 100644 --- a/tests/test_smoke.py +++ b/tests/test_smoke.py @@ -12,6 +12,7 @@ def test_imports_new_packages() -> None: "perception", "runtime", "semantic", + "skills_learning", "storage", "tools", "world",