From 44023a3d23e3c0c47c1aab6e2adc2d307989b0e5 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Thu, 3 Sep 2026 00:33:14 -0700 Subject: [PATCH] refactor(plugins/teams_pipeline,spotify,security-guidance): phase-split pipeline, dedupe Graph/Spotify clients, compact docs (3140->2631 LOC, --help/schema/pattern parity verified) --- plugins/platforms/teams/summary_writer.py | 92 ++---- plugins/security-guidance/__init__.py | 74 ++--- plugins/security-guidance/patterns.py | 53 ++-- plugins/spotify/__init__.py | 34 +-- plugins/spotify/client.py | 73 ++--- plugins/spotify/tools.py | 277 ++++++------------ plugins/teams_pipeline/cli.py | 86 ++---- plugins/teams_pipeline/meetings.py | 124 +++----- plugins/teams_pipeline/models.py | 104 +++---- plugins/teams_pipeline/pipeline.py | 330 ++++++++-------------- plugins/teams_pipeline/runtime.py | 31 +- plugins/teams_pipeline/store.py | 46 +-- plugins/teams_pipeline/subscriptions.py | 35 +-- 13 files changed, 425 insertions(+), 934 deletions(-) diff --git a/plugins/platforms/teams/summary_writer.py b/plugins/platforms/teams/summary_writer.py index 0de2f22d63..ae43fe7fa1 100644 --- a/plugins/platforms/teams/summary_writer.py +++ b/plugins/platforms/teams/summary_writer.py @@ -29,6 +29,9 @@ def _parse_bool(value: Any, *, default: bool = False) -> bool: _LIST_SECTIONS = (("Key decisions", "key_decisions"), ("Action items", "action_items"), ("Risks", "risks")) +# Env fallbacks for delivery config keys, applied only where nothing else set the key (access_token is a scoped secret). +_ENV_KEYS = {"delivery_mode": "TEAMS_DELIVERY_MODE", "incoming_webhook_url": "TEAMS_INCOMING_WEBHOOK_URL", + "access_token": "TEAMS_GRAPH_ACCESS_TOKEN", "team_id": "TEAMS_TEAM_ID", "channel_id": "TEAMS_CHANNEL_ID", "chat_id": "TEAMS_CHAT_ID"} class _StaticAccessTokenProvider: @@ -38,7 +41,6 @@ class _StaticAccessTokenProvider: self._access_token = str(access_token or "").strip() async def get_access_token(self, *, force_refresh: bool = False) -> str: - del force_refresh if not self._access_token: raise ValueError("TEAMS_GRAPH_ACCESS_TOKEN is required for graph delivery mode.") return self._access_token @@ -54,22 +56,17 @@ class TeamsSummaryWriter: self, platform_config: PlatformConfig | None = None, *, graph_client: Any | None = None, transport: httpx.AsyncBaseTransport | None = None, ) -> None: - self._platform_config = platform_config - self._graph_client = graph_client - self._transport = transport + self._platform_config, self._graph_client, self._transport = platform_config, graph_client, transport - async def write_summary( - self, payload: Any, config: dict[str, Any] | None, existing_record: Optional[dict[str, Any]] = None - ) -> dict[str, Any]: + async def write_summary(self, payload: Any, config: dict[str, Any] | None, existing_record: Optional[dict[str, Any]] = None) -> dict[str, Any]: merged = self._resolve_delivery_config(config) if existing_record and not _parse_bool(merged.get("force_resend"), default=False): return dict(existing_record) mode = str(merged.get("delivery_mode") or merged.get("mode") or "").strip().lower() - if not mode: - if merged.get("incoming_webhook_url"): - mode = "incoming_webhook" - elif merged.get("chat_id") or (merged.get("team_id") and merged.get("channel_id")): - mode = "graph" + if not mode and merged.get("incoming_webhook_url"): + mode = "incoming_webhook" + elif not mode and (merged.get("chat_id") or (merged.get("team_id") and merged.get("channel_id"))): + mode = "graph" if mode == "incoming_webhook": return await self._write_summary_via_incoming_webhook(payload, merged) if mode == "graph": @@ -86,22 +83,14 @@ class TeamsSummaryWriter: if platform_cfg.home_channel: merged.setdefault("channel_id", platform_cfg.home_channel.chat_id) merged.update(dict(config or {})) - env_defaults = { - "delivery_mode": os.getenv("TEAMS_DELIVERY_MODE", ""), - "incoming_webhook_url": os.getenv("TEAMS_INCOMING_WEBHOOK_URL", ""), - "access_token": _get_scoped_secret("TEAMS_GRAPH_ACCESS_TOKEN", ""), - "team_id": os.getenv("TEAMS_TEAM_ID", ""), - "channel_id": os.getenv("TEAMS_CHANNEL_ID", ""), - "chat_id": os.getenv("TEAMS_CHAT_ID", ""), - } - for key, value in env_defaults.items(): + for key, env in _ENV_KEYS.items(): + value = _get_scoped_secret(env, "") if key == "access_token" else os.getenv(env, "") if value and not merged.get(key): merged[key] = value return merged async def _write_summary_via_incoming_webhook(self, payload: Any, config: dict[str, Any]) -> dict[str, Any]: import httpx # lazy — see module docstring - webhook_url = str(config.get("incoming_webhook_url") or "").strip() if not webhook_url: raise ValueError("TEAMS_INCOMING_WEBHOOK_URL is required for incoming_webhook mode.") @@ -109,10 +98,7 @@ class TeamsSummaryWriter: async with httpx.AsyncClient(timeout=20.0, transport=self._transport) as client: response = await client.post(webhook_url, json=body) response.raise_for_status() - return { - "delivery_mode": "incoming_webhook", "webhook_url": webhook_url, - "status_code": response.status_code, "delivered": True, - } + return {"delivery_mode": "incoming_webhook", "webhook_url": webhook_url, "status_code": response.status_code, "delivered": True} async def _write_summary_via_graph(self, payload: Any, config: dict[str, Any]) -> dict[str, Any]: graph_client = self._build_graph_client(config) @@ -127,69 +113,43 @@ class TeamsSummaryWriter: raise ValueError("Graph delivery mode requires chat_id, or both team_id and channel_id.") path = f"/teams/{quote(team_id, safe='')}/channels/{quote(channel_id, safe='')}/messages" target = {"target_type": "channel", "team_id": team_id, "channel_id": channel_id} - response = await graph_client.post_json( - path, - json_body={"body": {"contentType": "html", "content": self._render_summary_html(payload)}}, - ) - return { - "delivery_mode": "graph", **target, - "message_id": (response or {}).get("id"), "web_url": (response or {}).get("webUrl"), - } + response = await graph_client.post_json(path, json_body={"body": {"contentType": "html", "content": self._render_summary_html(payload)}}) + return {"delivery_mode": "graph", **target, "message_id": (response or {}).get("id"), "web_url": (response or {}).get("webUrl")} def _build_graph_client(self, config: dict[str, Any]) -> Any: if self._graph_client is not None: return self._graph_client from tools.microsoft_graph_auth import MicrosoftGraphTokenProvider from tools.microsoft_graph_client import MicrosoftGraphClient - access_token = str(config.get("access_token") or "").strip() - if access_token: - return MicrosoftGraphClient(_StaticAccessTokenProvider(access_token), transport=self._transport) - return MicrosoftGraphClient(MicrosoftGraphTokenProvider.from_env(), transport=self._transport) + provider = _StaticAccessTokenProvider(access_token) if access_token else MicrosoftGraphTokenProvider.from_env() + return MicrosoftGraphClient(provider, transport=self._transport) def _render_summary_markdown(self, payload: Any) -> str: - lines = [ - f"**{self._title(payload)}**", - "", - f"Summary: {self._text(getattr(payload, 'summary', None), 'No summary available.')}", - ] + lines = [f"**{self._title(payload)}**", "", f"Summary: {self._text(getattr(payload, 'summary', None), 'No summary available.')}"] for heading, attr in _LIST_SECTIONS: lines += ["", f"{heading}:", *self._bullet_lines(getattr(payload, attr, None))] return "\n".join(lines) def _render_summary_html(self, payload: Any) -> str: - sections = [ - ("Summary", [self._text(getattr(payload, "summary", None), "No summary available.")]), - *((heading, list(getattr(payload, attr, None) or [])) for heading, attr in _LIST_SECTIONS), - ] - blocks = [f"
{html.escape(str(items[0]))}
") - continue - if items: - rendered = "".join(f"None
") - else: - blocks.append("None
") + summary = html.escape(self._text(getattr(payload, "summary", None), "No summary available.")) + blocks = [f"{summary}
"] + for heading, attr in _LIST_SECTIONS: + rendered = "".join(f"None
"] return "".join(blocks) @staticmethod def _title(payload: Any) -> str: - title = getattr(payload, "title", None) - if title: + if title := getattr(payload, "title", None): return str(title) meeting_ref = getattr(payload, "meeting_ref", None) - meeting_id = getattr(meeting_ref, "meeting_id", None) if meeting_ref else None - return f"Meeting {meeting_id or 'summary'}" + return f"Meeting {(getattr(meeting_ref, 'meeting_id', None) if meeting_ref else None) or 'summary'}" @staticmethod def _text(value: Any, default: str) -> str: - text = str(value or "").strip() - return text or default + return str(value or "").strip() or default @classmethod def _bullet_lines(cls, values: Any) -> list[str]: - items = [str(item).strip() for item in (values or []) if str(item).strip()] - return [f"- {item}" for item in items] or ["- None"] + return [f"- {str(item).strip()}" for item in (values or []) if str(item).strip()] or ["- None"] diff --git a/plugins/security-guidance/__init__.py b/plugins/security-guidance/__init__.py index 387a008524..999cffa6ec 100644 --- a/plugins/security-guidance/__init__.py +++ b/plugins/security-guidance/__init__.py @@ -1,11 +1,10 @@ """security-guidance plugin — fast pattern-matched security warnings on file writes. -Scans content written by ``write_file`` / ``patch`` / ``skill_manage`` for known dangerous -code patterns and appends a ``⚠️ Security guidance`` block to the tool result; the file is -still written and the model self-corrects next turn. Warn (not block) by default because -patterns have a real false-positive rate (``eval(`` in a tokenizer, ECB in a test fixture); -``SECURITY_GUIDANCE_BLOCK=1`` refuses the write instead, ``SECURITY_GUIDANCE_DISABLE=1`` is -a kill switch. Pattern data is ``patterns.py`` (Apache-2.0 fork, see LICENSE / NOTICE). +Scans content written by ``write_file`` / ``patch`` / ``skill_manage`` for known dangerous patterns +and appends a ``⚠️ Security guidance`` block to the tool result; the file is still written and the +model self-corrects next turn. Warn (not block) by default because patterns have a real false-positive +rate (``eval(`` in a tokenizer, ECB in a test fixture); ``SECURITY_GUIDANCE_BLOCK=1`` refuses the write +instead, ``SECURITY_GUIDANCE_DISABLE=1`` is a kill switch. Pattern data: ``patterns.py`` (Apache-2.0 fork). """ from __future__ import annotations @@ -42,24 +41,14 @@ def _compile_rules() -> List[Dict[str, Any]]: """Pre-compile regexes once; substrings stay plain (``in`` beats a literal regex).""" compiled: List[Dict[str, Any]] = [] for rule in _patterns.SECURITY_PATTERNS: - regex = None - re_src = rule.get("regex") - if re_src: - try: - regex = re.compile(re_src) - except re.error as err: - logger.warning( - "security-guidance: skipping rule %s — invalid regex %r: %s", - rule["ruleName"], re_src, err, - ) - continue + try: + regex = re.compile(rule["regex"]) if rule.get("regex") else None + except re.error as err: + logger.warning("security-guidance: skipping rule %s — invalid regex %r: %s", rule["ruleName"], rule["regex"], err) + continue compiled.append({ - "ruleName": rule["ruleName"], - "reminder": rule["reminder"], - "path_filter": rule.get("path_filter"), - "path_check": rule.get("path_check"), - "substrings": tuple(rule.get("substrings", ())), - "regex": regex, + "ruleName": rule["ruleName"], "reminder": rule["reminder"], "path_filter": rule.get("path_filter"), + "path_check": rule.get("path_check"), "substrings": tuple(rule.get("substrings", ())), "regex": regex, }) return compiled @@ -68,11 +57,8 @@ _COMPILED: List[Dict[str, Any]] = _compile_rules() def _rule_matches(entry: Dict[str, Any], path: str, content: str) -> bool: - """One rule against one write. Path predicates are best-effort: an exception is a non-match. - - path_check rules fire on the path ALONE (e.g. "you're editing a workflow file") and never - pattern-match content; path_filter gates content rules to relevant file types. - """ + """One rule against one write; a raising path predicate is a non-match. path_check rules fire on + the path ALONE and never scan content; path_filter gates content rules to relevant file types.""" try: if entry["path_check"] is not None: return bool(entry["path_check"](path)) @@ -80,9 +66,7 @@ def _rule_matches(entry: Dict[str, Any], path: str, content: str) -> bool: return False except Exception: return False - return any(sub in content for sub in entry["substrings"]) or ( - entry["regex"] is not None and bool(entry["regex"].search(content)) - ) + return any(sub in content for sub in entry["substrings"]) or (entry["regex"] is not None and bool(entry["regex"].search(content))) def _scan_content(path: str, content: str) -> List[Tuple[str, str]]: @@ -99,14 +83,8 @@ def _scan_args(tool_name: str, args: Any) -> List[Tuple[str, str]]: if _env_flag("SECURITY_GUIDANCE_DISABLE") or spec is None or not isinstance(args, dict): return [] path_key, content_keys = spec - path = args.get(path_key) or "" - if not isinstance(path, str): - path = "" - findings: List[Tuple[str, str]] = [] - for val in (args.get(ck) for ck in content_keys): - if isinstance(val, str) and val: - findings.extend(_scan_content(path, val)) - return findings + path = raw_path if isinstance(raw_path := args.get(path_key), str) else "" + return [finding for val in (args.get(ck) for ck in content_keys) if isinstance(val, str) and val for finding in _scan_content(path, val)] def _format_warning_block(findings: List[Tuple[str, str]]) -> str: @@ -125,29 +103,19 @@ def _format_warning_block(findings: List[Tuple[str, str]]) -> str: def _on_pre_tool_call(tool_name: str = "", args: Any = None, **_: Any) -> Optional[Dict[str, str]]: """Block mode only: refuse the write if any pattern matches (None = let it through).""" - if not _env_flag("SECURITY_GUIDANCE_BLOCK"): - return None - findings = _scan_args(tool_name, args) + findings = _scan_args(tool_name, args) if _env_flag("SECURITY_GUIDANCE_BLOCK") else [] if not findings: return None return { "action": "block", - "message": ( - "security-guidance refused this write: " - + _format_warning_block(findings) - + "\n\nTo override, unset SECURITY_GUIDANCE_BLOCK and retry." - ), + "message": "security-guidance refused this write: " + _format_warning_block(findings) + "\n\nTo override, unset SECURITY_GUIDANCE_BLOCK and retry.", } -def _on_transform_tool_result( - tool_name: str = "", args: Any = None, result: Any = None, **_: Any, -) -> Optional[str]: +def _on_transform_tool_result(tool_name: str = "", args: Any = None, result: Any = None, **_: Any) -> Optional[str]: """Warn mode: append the warning block to the result string (None = unchanged).""" # In block mode pre_tool_call already handled it — the tool didn't run, no result to wrap. - if _env_flag("SECURITY_GUIDANCE_BLOCK") or not isinstance(result, str): - return None - findings = _scan_args(tool_name, args) + findings = [] if _env_flag("SECURITY_GUIDANCE_BLOCK") or not isinstance(result, str) else _scan_args(tool_name, args) if not findings: return None # Don't decorate error results — the model already has bigger problems. diff --git a/plugins/security-guidance/patterns.py b/plugins/security-guidance/patterns.py index 296a947f99..287f8359fa 100644 --- a/plugins/security-guidance/patterns.py +++ b/plugins/security-guidance/patterns.py @@ -138,25 +138,18 @@ def _rule(name, reminder, **triggers): return {"ruleName": name, "reminder": reminder, **triggers} -# Security patterns configuration. Regex notes: -# - eval / exec lookbehinds exclude `.` so method calls (model.eval(), redis.eval()) don't match. -# - pickle matches deserialization only (load/loads/Unpickler); pickle.dump is not the RCE -# surface, and `pkl_load` needs a word boundary so similarly named safe loaders don't match. -# - script_src_without_sri: negative lookahead after `