diff --git a/.github/copilot-instructions.md b/.github/copilot-instructions.md index e72667258..642372820 100644 --- a/.github/copilot-instructions.md +++ b/.github/copilot-instructions.md @@ -31,11 +31,13 @@ For everything that isn't an attempt to use the kit (general questions, code exp - `/menu` - `/troubleshoot` - `/flightcheck` + - `/org-announcements` 2. **Intent hint — natural-language equivalent.** The user isn't typing a slash-command but is unambiguously asking to *run* the kit from this workspace. Examples: - "How do I set up the kit?" / "How do I run setup?" / "Start the ESS Maker Kit" - "Run flightcheck" / "Run the readiness check on my agent" - "Create a topic" / "Connect ServiceNow" / "Scan my agent for errors" — when phrased as a request to *do it now* in this workspace, not as a general "how does this work?" question. + - "Create an organization announcement" / "Post an announcement" / "Manage organization announcements" — when phrased as a request to act in this workspace. When in doubt, prefer the default behavior (answer normally) over firing the redirect. A user asking "what does /flightcheck do?" is asking a documentation question — answer it from the README and `solutions/ess-maker-skills/` files; do **not** redirect. @@ -51,7 +53,7 @@ When (and only when) the trigger conditions above are met, respond with **only** > 2. Navigate **inside** this folder, then **into** `solutions`, and select `ess-maker-skills` > 3. Click `Select Folder` > 4. VS Code will reopen with the kit loaded -> 5. Type your command again (for example `/setup` or `/run`) — it will work this time +> 5. Type your command again (for example `/setup`, `/org-announcements`, or `/run`) — it will work this time > > See the [README](README.md) for the full getting-started walkthrough. > diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 1501e9e4b..84d50a3ed 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -102,6 +102,7 @@ jobs: run: >- python -m pytest tests/mcp/agentconfig_core + tests/mcp/test_import_isolation.py tests/scripts/test_mcp_config.py -q @@ -137,6 +138,31 @@ jobs: node extension.test.js npm run validate + org-announcements: + name: Org Announcements configuration + runs-on: ubuntu-latest + + steps: + - uses: actions/checkout@v6 + + - name: Set up Python + uses: actions/setup-python@v6 + with: + python-version: '3.11' + + - name: Install Python dependencies + run: >- + pip install + -r requirements-dev.txt + -r solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/requirements.txt + + - name: Run Org Announcements tests + run: >- + python -m pytest + tests/mcp/agentconfig_org_announcements + tests/setup/test_org_announcements_integration.py + -q + flightcheck-tests: name: FlightCheck offline test suite runs-on: ubuntu-latest diff --git a/setup/Install-EssAdk.Tests.ps1 b/setup/Install-EssAdk.Tests.ps1 index af7d66ff6..36ee217a8 100644 --- a/setup/Install-EssAdk.Tests.ps1 +++ b/setup/Install-EssAdk.Tests.ps1 @@ -422,11 +422,21 @@ Test 'maker installer identity is instrumented (post-consolidation)' { # now a compat shim into the unified installer that passes -InstallMode maker, # the guard would drop all shim events. It must be gone. $old = $env:ESS_ADK_TELEMETRY + $oldSend = (Get-Command Send-EssTelEvent).ScriptBlock try { - $env:ESS_ADK_TELEMETRY = '' + $env:ESS_ADK_TELEMETRY = 'on' + Set-Item -Path Function:\Send-EssTelEvent -Value { + param([hashtable]$Envelope) + return 200 + } Initialize-EssInstallTelemetry -Installer 'lite' -InstallMode 'maker' if (-not $script:EssTel.Ready) { throw 'legacy lite installer identity should be telemetry-ready after consolidation' } - } finally { $env:ESS_ADK_TELEMETRY = $old; $script:EssTel.Ready = $false; $script:EssTel.Completed = $false } + } finally { + Set-Item -Path Function:\Send-EssTelEvent -Value $oldSend + $env:ESS_ADK_TELEMETRY = $old + $script:EssTel.Ready = $false + $script:EssTel.Completed = $false + } } Test 'PowerShell emitter no longer guards out the legacy lite installer' { $emitterSrc = Get-Content $psEmitter -Raw diff --git a/setup/README.md b/setup/README.md index 8141ab9e6..bcea5ead2 100644 --- a/setup/README.md +++ b/setup/README.md @@ -48,7 +48,7 @@ iex (irm https://raw.githubusercontent.com/microsoft/Employee-Self-Service-Agent Maker mode is the same install as Developer mode plus the **ESS Maker Profile** extension applying: - A chat-only layout with all developer surfaces hidden -- Big-button "Quick Actions" rail for common tasks (Connect, Customize landing page, Create, Scan, FlightCheck, Push) +- Big-button "Quick Actions" rail for common tasks (Setup, Customize landing page, Post an announcement, Create, Scan, FlightCheck, Push) - A built-in tutorial explaining each button You can switch between Maker mode and Developer mode at any time using the toggle buttons in the Quick Actions panel. diff --git a/solutions/ess-maker-skills/.github/copilot-instructions.md b/solutions/ess-maker-skills/.github/copilot-instructions.md index d34962891..c2d5d745c 100644 --- a/solutions/ess-maker-skills/.github/copilot-instructions.md +++ b/solutions/ess-maker-skills/.github/copilot-instructions.md @@ -278,8 +278,8 @@ Order of grounding sources (highest to lowest): microsoft/CopilotStudioSamples Employee Self-Service Agent samples. 3. `src/skills/` - kit-shipped skill instructions for /create, /update, /delete, /test, /scan, /evaluate, /push, /flightcheck, - /backup-template-configs, /restore-template-configs, and landing-page - configuration. + /backup-template-configs, /restore-template-configs, /org-announcements, + and landing-page configuration. 4. `src/reference/` (other subfolders) - additional kit-shipped guidance. 5. Web fetch / general knowledge - only when none of the above answer the question and only after telling the user you're falling back. @@ -452,6 +452,8 @@ pushed. Run the push pipeline when the maker asks to push local changes. | Restore or re-apply hybrid Workday HCM template configs | `src/skills/restore-template-configs/SKILL.md` | | View or configure ESS landing-page branding, quick links, starter prompts, insight cards, name, or icon | `src/skills/landing-page-config/SKILL.md` | | Invoke any tool from the `ess-landing-page-config` MCP server | `src/skills/landing-page-config/SKILL.md` | +| Create, edit, republish, archive, or manage organization announcements or bulletins | `src/skills/org-announcements/SKILL.md` | +| Invoke any tool from the `ess-org-announcements` MCP server | `src/skills/org-announcements/SKILL.md` | **Trigger phrases for connect:** "connect ServiceNow", "set up ServiceNow", "integrate ServiceNow", "connect Workday", "set up Workday", "add ServiceNow", @@ -475,6 +477,18 @@ links, starter prompts, Stay Up to Date, Quick Access, the agent name, or the agent icon, or asks what any landing-page setting controls for employees. Do not call an AgentConfiguration MCP tool from a generic flow. +**Org Announcements invocation:** Before invoking ANY tool from the +`ess-org-announcements` MCP server, read and follow +`src/skills/org-announcements/SKILL.md`. Its own `list_agent_configs` and +`search_agents` tools resolve missing deployed titleIds; do not start or call +the landing-page server for announcement discovery. This applies whether the user asks to +see, create, edit, republish, archive, or delete an announcement, mentions +announcements, org announcements, bulletins, or alerts, or asks who an +announcement reaches. Org Announcements are scoped to the authenticated tenant +and selected deployed agent's required `titleId`. The tenant is token-derived; +the title is not an audience group or author permission. Do not call an Org +Announcements MCP tool from a generic flow. + **FlightCheck results rendering:** When presenting `/flightcheck` results (Step 3 of `src/skills/flightcheck/SKILL.md`), read `workspace/flightcheck/results.json` with your file-reading tool and format the summary banner and tables **yourself, diff --git a/solutions/ess-maker-skills/.github/prompts/menu.prompt.md b/solutions/ess-maker-skills/.github/prompts/menu.prompt.md index 61f50ab91..5ef79a7de 100644 --- a/solutions/ess-maker-skills/.github/prompts/menu.prompt.md +++ b/solutions/ess-maker-skills/.github/prompts/menu.prompt.md @@ -14,6 +14,7 @@ Here's what I can help you with: | Command | What it does | |---------|-------------| | `/landing-page` | Configure the branding and content employees see when they open the ESS agent | +| `/org-announcements` | Create and manage announcements for the selected deployed ESS agent | | `/connect-workday` | Connect the active ESS HR agent to Workday | | `/connect` | Choose an available integration | | `/create` | Create a topic, workflow, or evaluation test set locally | diff --git a/solutions/ess-maker-skills/.github/prompts/org-announcements.prompt.md b/solutions/ess-maker-skills/.github/prompts/org-announcements.prompt.md new file mode 100644 index 000000000..c3ca03cc7 --- /dev/null +++ b/solutions/ess-maker-skills/.github/prompts/org-announcements.prompt.md @@ -0,0 +1,8 @@ +--- +mode: agent +description: "Create and manage announcements for the selected ESS agent" +--- + +# Org Announcements + +Read `src/skills/org-announcements/SKILL.md` and follow it. diff --git a/solutions/ess-maker-skills/.vscode/mcp.defaults.json b/solutions/ess-maker-skills/.vscode/mcp.defaults.json index 13dbccae3..b72a37cb2 100644 --- a/solutions/ess-maker-skills/.vscode/mcp.defaults.json +++ b/solutions/ess-maker-skills/.vscode/mcp.defaults.json @@ -4,6 +4,11 @@ "command": "{pythonExecutable}", "args": ["server.py"], "cwd": "${workspaceFolder}/src/mcp/agentconfig_landing_page" + }, + "ess-org-announcements": { + "command": "{pythonExecutable}", + "args": ["server.py"], + "cwd": "${workspaceFolder}/src/mcp/agentconfig_org_announcements" } } } diff --git a/solutions/ess-maker-skills/README.md b/solutions/ess-maker-skills/README.md index 82a19e3f1..5042d5e98 100644 --- a/solutions/ess-maker-skills/README.md +++ b/solutions/ess-maker-skills/README.md @@ -48,6 +48,48 @@ Copilot Studio and deployed to the organization. `/setup` installs and extracts the Power Platform agent; publication, admin approval, and Integrated apps deployment are separate steps. +### 📢 Post Organization Announcements + +Publish announcements for the selected deployed ESS agent and its audiences. +Run `/org-announcements`, ask `Create an announcement`, or use the **Post an +announcement** Quick Action. + +- **Standard announcements** carry a title, description, priority, and up to + two actions. +- **Alerts** carry a single link action for time-sensitive notices. +- **Audiences** are security groups, mail-enabled security groups, or classic + distribution groups, searched by name or email in one combined query. +- **Scheduling** publishes an announcement for a start/end window, and expired + announcements can be published again through the normal editor after + reviewing and updating their schedule. +- **Lifecycle** actions archive, unarchive, move back to draft, duplicate, or + delete an announcement. + +Describe the announcement in chat and the kit opens a pre-filled editor for +you to review — nothing is saved until you publish or save a draft in that +editor. + +Org Announcements are **scoped to the authenticated tenant and selected +agent's `titleId`**, not shared across agents. The current 100 limit and latest +50 archive window apply per tenant-and-agent pair. There is no tenant-wide +fallback. The title is resolved using `list_agent_configs` and `search_agents` +on the `ess-org-announcements` provider. Discovery shares neutral Python code +with the landing-page provider, but does not require its MCP process or +initialize its configuration. + +Announcement authoring requires the Org Announcements feature to be enabled +for your tenant, and audience search requires the `Directory.Read.All` +Microsoft Graph permission to be consented in your tenant. Graph uses a +separate resource token for the same authoring tenant and account. The current +account-context check requires readable `tid` and `oid` claims; opaque tokens +or credentials missing those claims return an explicit authentication failure +rather than using a different account. The API still validates tokens and +authorizes every request. + +This development surface requires the matching agent-qualified v1.1 backend +and scoped widget. The MCP rejects unscoped canonical responses instead of +silently consuming records from an older backend. + ### 📖 Pre-Loaded ESS Documentation, Samples & Best Practices The kit ships with a complete reference library that the AI agent reads at task time — you don't need to look anything up yourself. @@ -375,6 +417,7 @@ Then **run `/setup`** in GitHub Copilot Chat to configure your environment. |---------|-------------| | `/setup` | Connect this workspace to an existing editable DA Dev agent | | `/landing-page` | Configure landing-page branding and content | +| `/org-announcements` | Create and manage announcements for the selected ESS agent | | `/connect` | Explain the DA-GA product extension requirement | | `/create` | Create a topic, workflow, or evaluation test set locally | | `/update` | Update a topic, workflow, or evaluation test set locally | diff --git a/solutions/ess-maker-skills/src/mcp/agentconfig_core/_odata.py b/solutions/ess-maker-skills/src/mcp/agentconfig_core/_odata.py index 51379c5b1..41134c855 100644 --- a/solutions/ess-maker-skills/src/mcp/agentconfig_core/_odata.py +++ b/solutions/ess-maker-skills/src/mcp/agentconfig_core/_odata.py @@ -70,6 +70,14 @@ def _escape_odata_literal(value: str, name: str) -> str: return _validate_odata_string(value, name).replace("'", "''") +def _validate_title_id(title_id: str) -> str: + """Validate the opaque EmployeeAgents key shared by authoring surfaces.""" + _validate_odata_string(title_id, "titleId") + if len(title_id) > 256: + raise ValueError("titleId must not exceed 256 characters") + return title_id + + def _require_odata_id(value: str, name: str) -> str: """Validate a non-empty, control-char-free id and encode it as an OData key.""" return urllib.parse.quote(_escape_odata_literal(value, name), safe="") diff --git a/solutions/ess-maker-skills/src/mcp/agentconfig_core/agent_discovery.py b/solutions/ess-maker-skills/src/mcp/agentconfig_core/agent_discovery.py new file mode 100644 index 000000000..366114914 --- /dev/null +++ b/solutions/ess-maker-skills/src/mcp/agentconfig_core/agent_discovery.py @@ -0,0 +1,81 @@ +# Copyright (c) Microsoft Corporation. +# Licensed under the MIT License. + +"""Read-only deployed-agent discovery shared by feature-owned MCP providers.""" + +from __future__ import annotations + +from typing import Any + +from base_client import AgentConfigApiError, AgentConfigBaseClient + + +_MAX_SEARCH_LENGTH = 256 + + +def _convert_key_case(value: Any, *, upper: bool) -> Any: + if isinstance(value, list): + return [_convert_key_case(item, upper=upper) for item in value] + if not isinstance(value, dict): + return value + converted: dict[str, Any] = {} + for key, item in value.items(): + if key and key[0].isalpha(): + first = key[0].upper() if upper else key[0].lower() + converted_key = first + key[1:] + else: + converted_key = key + converted[converted_key] = _convert_key_case(item, upper=upper) + return converted + + +def _to_api_payload(value: Any) -> Any: + return _convert_key_case(value, upper=True) + + +def _to_tool_payload(value: Any) -> Any: + return _convert_key_case(value, upper=False) + + +def _unwrap_agent_collection(payload: Any) -> list[dict[str, Any]]: + if isinstance(payload, list): + return payload + if isinstance(payload, dict) and isinstance(payload.get("value"), list): + return payload["value"] + raise AgentConfigApiError( + "AgentConfiguration API returned an invalid collection response" + ) + + +class AgentDiscoveryClient(AgentConfigBaseClient): + """Discover deployed titleIds without initializing any feature configuration. + + Feature clients retain their own base URL, transport, and authentication. + Explicit response conversion keeps discovery independent of each feature's + canonical payload casing and collection routes. + """ + + def _agent_collection_path(self) -> str: + return f"tenants('{self.tenant_id}')/EmployeeAgents" + + async def list_agent_configs(self) -> list[dict[str, Any]]: + payload = await self._request( + "GET", self._agent_collection_path(), transform_payload=False + ) + return _unwrap_agent_collection(_to_tool_payload(payload)) + + async def search_agents(self, search_string: str) -> list[dict[str, Any]]: + if not isinstance(search_string, str) or not search_string.strip(): + raise ValueError("searchString must be a non-empty string") + normalized = search_string.strip() + if len(normalized) > _MAX_SEARCH_LENGTH: + raise ValueError( + f"searchString must not exceed {_MAX_SEARCH_LENGTH} characters" + ) + payload = await self._request( + "POST", + f"{self._agent_collection_path()}/SearchAgents", + json={"SearchString": normalized}, + transform_payload=False, + ) + return _unwrap_agent_collection(_to_tool_payload(payload)) diff --git a/solutions/ess-maker-skills/src/mcp/agentconfig_core/base_client.py b/solutions/ess-maker-skills/src/mcp/agentconfig_core/base_client.py index 9d47ac4ef..8fd7c6f07 100644 --- a/solutions/ess-maker-skills/src/mcp/agentconfig_core/base_client.py +++ b/solutions/ess-maker-skills/src/mcp/agentconfig_core/base_client.py @@ -28,7 +28,6 @@ import asyncio import base64 -import binascii import http.server import json import logging @@ -39,9 +38,11 @@ import urllib.parse import uuid import webbrowser +from enum import Enum from typing import Any, Optional import httpx +from portalocker.exceptions import LockException from _tenant_context import configured_tenant_id from _token_cache import create_token_cache @@ -68,6 +69,41 @@ def __init__(self, message: str, *, http_status: int | None = None): self.http_status = http_status +class _CredentialFailure(Enum): + FILE_MISSING = ( + "AGENTCONFIG_ACCESS_TOKEN_FILE does not exist. Restore the configured " + "token file for the original account and retry." + ) + FILE_UNREADABLE = ( + "AGENTCONFIG_ACCESS_TOKEN_FILE could not be read. " + "Check token-file access and retry." + ) + FILE_EMPTY = ( + "AGENTCONFIG_ACCESS_TOKEN_FILE is empty. " + "Provide a valid token for the original account and retry." + ) + INVALID_TOKEN = ( + "The replacement access token is malformed or has no valid tenant id. " + "Provide a valid token for the original account and retry." + ) + TENANT_MISMATCH = ( + "The replacement access token does not match the original tenant. " + "Provide a token for the original account and tenant and retry." + ) + ACCOUNT_MISMATCH = ( + "The replacement access token does not identify the original account. " + "Provide a matching token and retry." + ) + + +class LocalCredentialError(ValueError): + """Local validation failure with an allowlisted renewal diagnostic.""" + + def __init__(self, message: str, reason: _CredentialFailure): + super().__init__(message) + self.reason = reason + + def _resolve_token() -> str: """Resolve a token without writing it to logs or MCP configuration.""" expected_tenant_id = configured_tenant_id() @@ -85,14 +121,19 @@ def _resolve_token() -> str: def _read_token_file(token_file: str) -> str: if not os.path.isfile(token_file): - raise ValueError( - f"AGENTCONFIG_ACCESS_TOKEN_FILE={token_file!r} does not exist" + raise LocalCredentialError( + "AGENTCONFIG_ACCESS_TOKEN_FILE does not exist", _CredentialFailure.FILE_MISSING ) - with open(token_file, "r", encoding="utf-8") as handle: - token = handle.read().strip() + try: + with open(token_file, "r", encoding="utf-8") as handle: + token = handle.read().strip() + except (OSError, UnicodeError) as error: + raise LocalCredentialError( + "AGENTCONFIG_ACCESS_TOKEN_FILE could not be read", _CredentialFailure.FILE_UNREADABLE + ) from error if not token: - raise ValueError( - f"AGENTCONFIG_ACCESS_TOKEN_FILE={token_file!r} is empty" + raise LocalCredentialError( + "AGENTCONFIG_ACCESS_TOKEN_FILE is empty", _CredentialFailure.FILE_EMPTY ) return token @@ -102,9 +143,10 @@ def _validate_token_tenant(token: str, expected_tenant_id: str | None) -> str: expected_tenant_id is not None and _decode_tenant_id_from_jwt(token) != expected_tenant_id ): - raise ValueError( + raise LocalCredentialError( "The AgentConfiguration token tenant does not match the configured " - "Dataverse environment. Sign in to that tenant or provide a matching token." + "Dataverse environment. Sign in to that tenant or provide a matching token.", + _CredentialFailure.TENANT_MISMATCH, ) return token @@ -157,18 +199,19 @@ def acquire_token_msal_interactive(expected_tenant_id: str | None = None) -> str result = _acquire_token_interactive_form_post(app) if "access_token" not in result: - error = result.get("error", "unknown_error") - description = result.get("error_description", "") - raise ValueError(f"MSAL sign-in failed ({error}): {description}") + raise ValueError("AgentConfiguration sign-in failed. Sign in again and retry.") return result["access_token"] -def _refresh_msal_token(tenant_id: str, object_id: str) -> str: +def account_for_identity(app: Any, tenant_id: str, object_id: str) -> Any | None: + """Find one MSAL home account through its tenant-local profile. + + get_accounts() groups tenant profiles by home account. Its exposed local id + may belong to another tenant (notably for guests), so match the original + tenant profile in the shared cache before selecting a grouped account. + """ import msal - app = _create_msal_app(tenant_id) - # get_accounts() groups profiles by home account and can expose another - # tenant's local id. Resolve the original tenant profile through the cache. home_account_ids = { account["home_account_id"] for account in app.token_cache.search(msal.TokenCache.CredentialType.ACCOUNT) @@ -182,13 +225,19 @@ def _refresh_msal_token(tenant_id: str, object_id: str) -> str: account for account in app.get_accounts() if account.get("home_account_id") in home_account_ids ] - if len(accounts) != 1: + return accounts[0] if len(accounts) == 1 else None + + +def _refresh_msal_token(tenant_id: str, object_id: str) -> str: + app = _create_msal_app(tenant_id) + account = account_for_identity(app, tenant_id, object_id) + if account is None: raise AgentConfigApiError( "The original account is unavailable for token refresh. " "Sign in again as that account and retry.", http_status=401, ) - result = app.acquire_token_silent(_SCOPE, account=accounts[0], force_refresh=True) + result = app.acquire_token_silent(_SCOPE, account=account, force_refresh=True) if not result or not result.get("access_token"): raise AgentConfigApiError( "Could not refresh the access token for the original account. " @@ -223,10 +272,10 @@ def _acquire_token_interactive_form_post(app: Any) -> dict[str, Any]: server.server_close() if not _FormPostCaptureHandler.captured: - return { - "error": "timeout", - "error_description": "No sign-in callback received within 300 seconds.", - } + raise ValueError( + "AgentConfiguration sign-in timed out waiting for the browser callback. " + "Sign in again and retry." + ) return app.acquire_token_by_auth_code_flow( flow, @@ -243,30 +292,39 @@ def _decode_jwt_payload(token: str) -> dict[str, Any]: """ parts = token.split(".") if len(parts) != 3: - raise ValueError( + raise LocalCredentialError( "AGENTCONFIG_ACCESS_TOKEN does not look like a JWT " - "(expected three dot-separated segments)" + "(expected three dot-separated segments)", + _CredentialFailure.INVALID_TOKEN, ) payload_segment = parts[1] padded = payload_segment + "=" * (-len(payload_segment) % 4) try: - return json.loads(base64.urlsafe_b64decode(padded)) - except (binascii.Error, json.JSONDecodeError, UnicodeDecodeError) as error: - raise ValueError( - f"Could not decode AGENTCONFIG_ACCESS_TOKEN payload: {error}" + payload = json.loads(base64.urlsafe_b64decode(padded)) + except ValueError as error: + raise LocalCredentialError( + "Could not decode AGENTCONFIG_ACCESS_TOKEN payload", + _CredentialFailure.INVALID_TOKEN, ) from error + if not isinstance(payload, dict): + raise LocalCredentialError( + "AGENTCONFIG_ACCESS_TOKEN payload must be an object", _CredentialFailure.INVALID_TOKEN + ) + return payload def _decode_tenant_id_from_jwt(token: str) -> str: """Decode and validate the tenant ID (``tid``) used to address the route.""" tenant_id = _decode_jwt_payload(token).get("tid") if not isinstance(tenant_id, str) or not tenant_id: - raise ValueError("AGENTCONFIG_ACCESS_TOKEN payload has no 'tid' claim") + raise LocalCredentialError( + "AGENTCONFIG_ACCESS_TOKEN payload has no 'tid' claim", _CredentialFailure.INVALID_TOKEN + ) try: return str(uuid.UUID(tenant_id)) except ValueError as error: - raise ValueError( - "AGENTCONFIG_ACCESS_TOKEN payload has an invalid 'tid' claim" + raise LocalCredentialError( + "AGENTCONFIG_ACCESS_TOKEN payload has an invalid 'tid' claim", _CredentialFailure.INVALID_TOKEN ) from error @@ -287,6 +345,19 @@ def _decode_object_id_from_jwt(token: str) -> Optional[str]: return None +def validate_token_identity(token: str, tenant_id: str, object_id: str | None) -> str: + """Check resource-token context; services still validate and authorize tokens.""" + _validate_token_tenant(token, tenant_id) + if not object_id: + raise ValueError("The authoring account cannot be identified without an oid claim") + actual = _decode_object_id_from_jwt(token) + if actual is None or actual.casefold() != object_id.casefold(): + raise LocalCredentialError( + "The token does not identify the intended authoring account", _CredentialFailure.ACCOUNT_MISMATCH + ) + return token + + class AgentConfigBaseClient: """Neutral AgentConfiguration client core shared by the landing-page and planner MCPs. @@ -326,6 +397,11 @@ def __repr__(self) -> str: f"tenant_id={self.tenant_id!r}>" ) + @property + def object_id(self) -> str | None: + """Captured tenant-local principal context, never a tool argument.""" + return self._object_id + def _transform_response(self, payload: Any) -> Any: """Surface-specific response key transform; identity in the neutral core. @@ -360,7 +436,7 @@ async def aclose(self) -> None: self._client = None def _acquire_replacement_token(self) -> str: - if self._object_id is None: + if self.object_id is None: raise AgentConfigApiError( "The original account cannot be identified for token refresh. " "Provide a token with an object id and recreate the client.", @@ -377,23 +453,33 @@ def _acquire_replacement_token(self) -> str: http_status=401, ) else: - token = _refresh_msal_token(self.tenant_id, self._object_id) - - _validate_token_tenant(token, self.tenant_id) - object_id = _decode_object_id_from_jwt(token) - if object_id is None or object_id.casefold() != self._object_id.casefold(): - raise AgentConfigApiError( - "The replacement token does not identify the original account. " - "Provide a matching token and retry.", - http_status=401, - ) - return token + token = _refresh_msal_token(self.tenant_id, self.object_id) + return validate_token_identity(token, self.tenant_id, self.object_id) async def _refresh_token(self, rejected_token: str) -> None: async with self._token_lock: if self._token != rejected_token: return - token = await asyncio.to_thread(self._acquire_replacement_token) + try: + token = await asyncio.to_thread(self._acquire_replacement_token) + except (ValueError, LockException, OSError) as error: + message = ( + "Could not renew credentials for the original account. " + "Provide matching credentials and retry." + ) + if isinstance(error, LocalCredentialError): + message = error.reason.value + elif isinstance(error, LockException): + message = ( + "The credential cache could not be locked. Retry shortly; " + "if this persists, check local cache access." + ) + elif isinstance(error, (PermissionError, FileNotFoundError)): + message = ( + "Local credential access failed. Check credential-file and " + "cache availability and permissions, then retry." + ) + raise AgentConfigApiError(message, http_status=401) from error if token == rejected_token: raise AgentConfigApiError( "No replacement access token is available. " diff --git a/solutions/ess-maker-skills/src/mcp/agentconfig_landing_page/client.py b/solutions/ess-maker-skills/src/mcp/agentconfig_landing_page/client.py index 616ac67cb..f16cc9190 100644 --- a/solutions/ess-maker-skills/src/mcp/agentconfig_landing_page/client.py +++ b/solutions/ess-maker-skills/src/mcp/agentconfig_landing_page/client.py @@ -28,52 +28,20 @@ from _odata import ( # noqa: E402 _require_odata_id, _validate_https_base_url, - _validate_odata_string, + _validate_title_id, ) -from base_client import AgentConfigApiError, AgentConfigBaseClient # noqa: E402 +from agent_discovery import ( # noqa: E402 + AgentDiscoveryClient, + _to_api_payload, + _to_tool_payload, +) +from base_client import AgentConfigApiError as AgentConfigApiError # noqa: E402 DEFAULT_AGENTCONFIG_BASE_URL = "https://substrate.office.com/weveb2/api/v1.1" -_MAX_TITLE_ID_LENGTH = 256 -_MAX_SEARCH_LENGTH = 256 - - -def _validate_title_id(title_id: str) -> str: - _validate_odata_string(title_id, "titleId") - if len(title_id) > _MAX_TITLE_ID_LENGTH: - raise ValueError( - f"titleId must not exceed {_MAX_TITLE_ID_LENGTH} characters" - ) - return title_id - - -def _convert_key_case(value: Any, *, upper: bool) -> Any: - """Recursively convert the first character of JSON object keys.""" - if isinstance(value, list): - return [_convert_key_case(item, upper=upper) for item in value] - if not isinstance(value, dict): - return value - - converted: dict[str, Any] = {} - for key, item in value.items(): - if key and key[0].isalpha(): - first = key[0].upper() if upper else key[0].lower() - converted_key = first + key[1:] - else: - converted_key = key - converted[converted_key] = _convert_key_case(item, upper=upper) - return converted - -def _to_api_payload(value: Any) -> Any: - return _convert_key_case(value, upper=True) - -def _to_tool_payload(value: Any) -> Any: - return _convert_key_case(value, upper=False) - - -class AgentConfigClient(AgentConfigBaseClient): +class AgentConfigClient(AgentDiscoveryClient): """Async client for production EmployeeAgents list/search/create/get/PATCH. Inherits auth, token decode, the httpx session, and the retrying @@ -97,41 +65,12 @@ def _transform_response(self, payload: Any) -> Any: return _to_tool_payload(payload) def _collection_path(self) -> str: - return f"tenants('{self.tenant_id}')/EmployeeAgents" + return self._agent_collection_path() def _agent_path(self, title_id: str) -> str: encoded = _require_odata_id(_validate_title_id(title_id), "titleId") return f"{self._collection_path()}('{encoded}')" - @staticmethod - def _unwrap_collection(payload: Any) -> list[dict[str, Any]]: - if isinstance(payload, list): - return payload - if isinstance(payload, dict) and isinstance(payload.get("value"), list): - return payload["value"] - raise AgentConfigApiError( - "AgentConfiguration API returned an invalid collection response" - ) - - async def list_agent_configs(self) -> list[dict[str, Any]]: - payload = await self._request("GET", self._collection_path()) - return self._unwrap_collection(payload) - - async def search_agents(self, search_string: str) -> list[dict[str, Any]]: - if not isinstance(search_string, str) or not search_string.strip(): - raise ValueError("searchString must be a non-empty string") - normalized = search_string.strip() - if len(normalized) > _MAX_SEARCH_LENGTH: - raise ValueError( - f"searchString must not exceed {_MAX_SEARCH_LENGTH} characters" - ) - payload = await self._request( - "POST", - f"{self._collection_path()}/SearchAgents", - json={"SearchString": normalized}, - ) - return self._unwrap_collection(payload) - async def create_agent_config(self, title_id: str) -> dict[str, Any]: title_id = _validate_title_id(title_id) return await self._request( diff --git a/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/.gitignore b/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/.gitignore new file mode 100644 index 000000000..57006c3d8 --- /dev/null +++ b/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/.gitignore @@ -0,0 +1,12 @@ +# Compiled Python +__pycache__/ +*.pyc + +# This server owns NO token cache of its own. It uses exactly the two shared +# delegated caches: +# - Microsoft Graph -> solutions/ess-maker-skills/.local/.token_cache.bin +# - AgentConfiguration -> src/mcp/agentconfig_core/.local/msal_token_cache.bin +# both of which are ignored by their own directories' rules. This entry is +# belt-and-braces only: if anything ever drops local state here it must not be +# committed. +.local/ diff --git a/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/client.py b/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/client.py new file mode 100644 index 000000000..4869e2f9c --- /dev/null +++ b/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/client.py @@ -0,0 +1,659 @@ +# Copyright (c) Microsoft Corporation. +# Licensed under the MIT License. + +"""WeveNova Org Announcements (EssBulletin) authoring client. + +The shared client core — bearer-token acquisition, the JWT ``tid`` decode, the +httpx session, and the retrying ``_request`` — lives in the neutral +``agentconfig_core`` core (``base_client.AgentConfigBaseClient``). This module +keeps only what is specific to the Org Announcements authoring surface: the +v1.1 base URL, the three agent-qualified ``EssBulletins`` routes, and the +manager-state classification. + +The collection is keyed by the authenticated tenant and deployed ESS titleId. +This new backend surface must never fall back to tenant-only routes. +""" + +from __future__ import annotations + +import asyncio +import os +import sys +from datetime import datetime, timezone +from typing import Any, Optional + +import httpx + +# The AgentConfiguration MCP family lives at the ``src/mcp`` root as sibling +# folders sharing the neutral ``agentconfig_core`` client core. There is no +# package __init__.py, and each server launches with cwd set to its own folder +# on a flat sys.path, so make the sibling ``agentconfig_core`` folder importable. +sys.path.insert( + 0, + os.path.join( + os.path.dirname(os.path.abspath(__file__)), "..", "agentconfig_core" + ), +) + +from _odata import ( # noqa: E402 + _require_odata_id, + _validate_https_base_url, +) +from agent_discovery import AgentDiscoveryClient # noqa: E402 +from base_client import AgentConfigApiError # noqa: E402 +from validation import validate_bulletin_id, validate_title_id # noqa: E402 + + +DEFAULT_ORG_ANNOUNCEMENTS_BASE_URL = "https://substrate.office.com/weveb2/api/v1.1" + +# WeveNova returns every current item followed by at most this many archived +# items, with no envelope metadata. Hitting the cap means the archived window is +# truncated; it does not reveal an exact archived total. +ARCHIVED_WINDOW_SIZE = 50 +STATUS_DELETED = "deleted" + +_VALID_STATUSES = frozenset({"draft", "published", "retired", STATUS_DELETED}) +_VALID_BULLETIN_TYPES = frozenset({"standard", "alert"}) +_VALID_ACTION_TYPES = frozenset({"externalLink", "copilotChat"}) +_VALID_PRIORITIES = frozenset({0, 1}) +_COMMITTED_RELOAD_404_DELAYS = (0.05, 0.15) + + +class IndeterminateWriteError(AgentConfigApiError): + """An unkeyed create may have committed but was never acknowledged. + + Raised only for writes that cannot be safely replayed: an unkeyed create or + a duplicate that failed with an ambiguous network error or a 502/503/504 + gateway response. The caller must refresh canonical state before retrying so + a committed-but-unacknowledged POST is never duplicated. + """ + + +class BulletinValidationError(AgentConfigApiError): + """The OData ``Save`` action returned HTTP 200 with field errors. + + WeveNova's ``EssBulletinSaveResult`` reports validation failure *inside* a + success response: ``Errors`` is non-empty and ``Id`` is absent. Every + ``{code, field, message}`` entry is preserved verbatim so the widget can + render its localized copy against the exact backend code and attach the + message to the exact field. Flattening them into one generic message would + silently destroy that mapping and leave the maker with no way to know which + field to fix. + """ + + def __init__(self, errors: list[dict[str, Any]]): + super().__init__( + "The Org Announcements service rejected the announcement.", + http_status=200, + ) + self.errors = errors + + +class CommittedCanonicalReloadError(Exception): + """A Save committed, but its canonical keyed reload failed.""" + + def __init__(self, cause: AgentConfigApiError | httpx.RequestError): + super().__init__( + "The announcement was saved, but its canonical state could not be reloaded." + ) + self.cause = cause + + +_BULLETIN_FIELD_MAP = { + "Type": "type", + "Priority": "priority", + "Title": "title", + "Description": "description", + "PrimaryAction": "primaryAction", + "SecondaryAction": "secondaryAction", + "StartDate": "startDate", + "EndDate": "endDate", +} + +_ACTION_FIELD_MAP = { + "ActionType": "actionType", + "Label": "label", + "Url": "url", + "Prompt": "prompt", +} + + +def _require_odata_enum( + value: Any, field: str, allowed: frozenset[Any] +) -> Any: + if not any( + type(value) is type(candidate) and value == candidate + for candidate in allowed + ): + raise AgentConfigApiError( + f"Org Announcements API returned an invalid {field}" + ) + return value + + +def _require_input_enum( + value: Any, field: str, allowed: frozenset[Any] +) -> Any: + if not any( + type(value) is type(candidate) and value == candidate + for candidate in allowed + ): + raise ValueError(f"{field} has an unsupported value") + return value + + +def _from_odata_action(payload: Any) -> Optional[dict[str, Any]]: + if payload is None: + return None + if not isinstance(payload, dict): + raise AgentConfigApiError( + "Org Announcements API returned an invalid bulletin action" + ) + _require_odata_enum( + payload.get("ActionType"), "action type", _VALID_ACTION_TYPES + ) + return { + target: payload[source] + for source, target in _ACTION_FIELD_MAP.items() + if source in payload + } + + +def _from_odata_bulletin(payload: Any) -> dict[str, Any]: + if not isinstance(payload, dict): + raise AgentConfigApiError( + "Org Announcements API returned invalid bulletin content" + ) + if "Id" in payload or "id" in payload: + raise AgentConfigApiError( + "Org Announcements API returned a duplicate id inside bulletin content" + ) + + result: dict[str, Any] = {} + for source, target in _BULLETIN_FIELD_MAP.items(): + if source not in payload: + continue + value = payload[source] + if source in ("PrimaryAction", "SecondaryAction"): + value = _from_odata_action(value) + elif source == "Type": + value = _require_odata_enum( + value, "bulletin type", _VALID_BULLETIN_TYPES + ) + elif source == "Priority": + value = _require_odata_enum( + value, "bulletin priority", _VALID_PRIORITIES + ) + result[target] = value + return result + + +def _to_odata_action(payload: Any) -> Optional[dict[str, Any]]: + if payload is None: + return None + if not isinstance(payload, dict): + raise ValueError("bulletin actions must be objects") + _require_input_enum( + payload.get("actionType"), "actionType", _VALID_ACTION_TYPES + ) + return { + source: payload[target] + for source, target in _ACTION_FIELD_MAP.items() + if target in payload + } + + +def _to_odata_bulletin(payload: Any) -> dict[str, Any]: + if not isinstance(payload, dict): + raise ValueError("bulletin must be an object") + result: dict[str, Any] = {} + for source, target in _BULLETIN_FIELD_MAP.items(): + if target not in payload: + continue + value = payload[target] + if target in ("primaryAction", "secondaryAction"): + value = _to_odata_action(value) + elif target == "type": + value = _require_input_enum( + value, "type", _VALID_BULLETIN_TYPES + ) + elif target == "priority": + value = _require_input_enum( + value, "priority", _VALID_PRIORITIES + ) + result[source] = value + return result + + +def _to_odata_save_input(payload: dict[str, Any]) -> dict[str, Any]: + result: dict[str, Any] = {} + if "id" in payload: + result["Id"] = validate_bulletin_id(payload["id"]) + if "bulletin" in payload: + result["Bulletin"] = _to_odata_bulletin(payload["bulletin"]) + if "audience" in payload: + result["Audience"] = payload["audience"] + if "status" not in payload: + raise ValueError("status is required") + result["Status"] = _require_input_enum( + payload["status"], "status", _VALID_STATUSES + ) + return result + + +def _parse_instant(value: Any) -> Optional[datetime]: + """Parse a UTC ISO instant, returning ``None`` for absent or unparseable text. + + ``None`` means "no boundary", which the classifier treats as open-ended + rather than as an error: an unscheduled published item is current. + """ + if not isinstance(value, str) or not value: + return None + text = value.strip() + if not text: + return None + if text.endswith("Z") or text.endswith("z"): + text = f"{text[:-1]}+00:00" + try: + parsed = datetime.fromisoformat(text) + except ValueError: + return None + if parsed.tzinfo is None: + return parsed.replace(tzinfo=timezone.utc) + return parsed.astimezone(timezone.utc) + + +def is_deleted_item(config: dict[str, Any]) -> bool: + """Classify one stored configuration as soft-deleted. + + ``delete`` is a lifecycle transition, not a hard removal, so the list + endpoint can still return the row. A deleted announcement is not a working + item and is not an archived item the maker can restore, so it is excluded + from the manager entirely rather than being counted in either bucket. + """ + return config.get("status") == "deleted" + + +def is_archived_item(config: dict[str, Any], now: datetime) -> bool: + """Classify one stored configuration as archived. + + Archived means Retired, or Published with an end instant already elapsed. + Everything else — Draft, and Published without an elapsed end — is current. + Position in the API response is not used, because the list carries no + envelope metadata that would make position authoritative. + """ + status = config.get("status") + if status == "retired": + return True + if status != "published": + return False + + bulletin = config.get("bulletin") + end = _parse_instant(bulletin.get("endDate")) if isinstance(bulletin, dict) else None + return end is not None and end < now + + +class OrgAnnouncementsClient(AgentDiscoveryClient): + """Async client for the tenant-and-agent-scoped ``EssBulletins`` routes. + + Inherits auth, the token decode, the httpx session, and the retrying + ``_request`` from ``AgentConfigBaseClient``; adds only the v1.1 base URL and + the three authoring routes. WeveNova's OData properties are PascalCase; + this client explicitly adapts them to the existing lower-camel MCP/widget + contract at the HTTP boundary. + """ + + def __init__(self, *, transport: Optional[httpx.AsyncBaseTransport] = None): + base_url = _validate_https_base_url( + os.environ.get( + "ORG_ANNOUNCEMENTS_BASE_URL", DEFAULT_ORG_ANNOUNCEMENTS_BASE_URL + ), + "ORG_ANNOUNCEMENTS_BASE_URL", + ) + super().__init__( + base_url=base_url, + logger_name="ess-org-announcements", + transport=transport, + ) + + def _collection_path(self, title_id: str) -> str: + encoded = _require_odata_id(validate_title_id(title_id), "titleId") + return f"tenants('{self.tenant_id}')/EmployeeAgents('{encoded}')/EssBulletins" + + @staticmethod + def _require_config( + payload: Any, + title_id: str, + *, + expected_bulletin_id: Optional[str] = None, + ) -> dict[str, Any]: + """Reject a success-shaped response that is not a canonical record. + + A malformed body must not become an empty default, because the widget + would then render a blank editor over real stored content. + """ + if not isinstance(payload, dict): + raise AgentConfigApiError( + "Org Announcements API returned an invalid bulletin configuration" + ) + if payload.get("TitleId") != title_id: + raise AgentConfigApiError( + "Org Announcements API returned a configuration with missing " + "or mismatched titleId" + ) + raw_id = payload.get("Id") + if not isinstance(raw_id, str): + raise AgentConfigApiError( + "Org Announcements API returned a configuration without a valid id" + ) + try: + bulletin_id = validate_bulletin_id(raw_id) + except ValueError as error: + raise AgentConfigApiError( + "Org Announcements API returned a configuration without a valid id" + ) from error + if expected_bulletin_id is not None and bulletin_id != expected_bulletin_id: + raise AgentConfigApiError( + "Org Announcements API returned a configuration for a different " + "bulletin id" + ) + audience = payload.get("Audience") + if not isinstance(audience, list) or not all( + isinstance(group_id, str) for group_id in audience + ): + raise AgentConfigApiError( + "Org Announcements API returned an invalid audience" + ) + status = _require_odata_enum( + payload.get("Status"), "status", _VALID_STATUSES + ) + + config = { + "id": bulletin_id, + "titleId": title_id, + "bulletin": _from_odata_bulletin(payload.get("Bulletin")), + "audience": audience, + "status": status, + } + for source, target in ( + ("CreatedBy", "createdBy"), + ("CreatedOn", "createdOn"), + ("ModifiedDate", "modifiedDate"), + ): + if source in payload: + config[target] = payload[source] + return config + + @classmethod + def _unwrap_collection(cls, payload: Any, title_id: str) -> list[dict[str, Any]]: + if isinstance(payload, dict) and isinstance(payload.get("value"), list): + items = payload["value"] + else: + raise AgentConfigApiError( + "Org Announcements API returned an invalid collection response" + ) + return [cls._require_config(item, title_id) for item in items] + + @staticmethod + def _normalize_save_errors(raw: Any) -> list[dict[str, Any]]: + """Project the wrapper's ``errors`` into ``{code, field, message}`` rows. + + Every entry is preserved — none are collapsed, deduplicated, or dropped + — because the widget maps each backend ``code`` to localized copy and + binds each ``field`` to an input. An entry that is not an object at all + is still represented (with the raw text as its message) rather than + discarded, so a schema drift surfaces as a visible error instead of a + silently successful save. + """ + normalized: list[dict[str, Any]] = [] + for entry in raw: + if isinstance(entry, dict): + code = entry.get("Code") + field = entry.get("Field") + message = entry.get("Message") + normalized.append( + { + "code": ( + code + if isinstance(code, str) and code + else "InvalidRequest" + ), + "field": field if isinstance(field, str) and field else None, + "message": ( + message + if isinstance(message, str) and message + else "The Org Announcements service rejected this value." + ), + } + ) + else: + normalized.append( + { + "code": "InvalidRequest", + "field": None, + "message": str(entry), + } + ) + return normalized + + @classmethod + def _unwrap_save_result( + cls, payload: Any, *, requested_id: Optional[str] + ) -> str: + """Validate and unwrap an ``EssBulletinSaveResult``. + + The OData action answers HTTP 200 with ``{Id, Errors}`` and reports + validation failure inside a success status. Three outcomes are + distinguished: + + * ``Errors`` non-empty → :class:`BulletinValidationError` carrying every + entry, so field-level backend codes survive to the widget. + * ``Errors`` empty and ``Id`` valid → the affected canonical key, after + checking it agrees with the ID the caller asked to update. + * anything else → :class:`AgentConfigApiError`, never an empty default: + a blank record would render an empty editor over real stored content. + """ + if not isinstance(payload, dict): + raise AgentConfigApiError( + "Org Announcements API returned an invalid save response" + ) + + errors = payload.get("Errors") + if not isinstance(errors, list): + raise AgentConfigApiError( + "Org Announcements API returned an invalid save error list" + ) + if errors: + raise BulletinValidationError(cls._normalize_save_errors(errors)) + + raw_id = payload.get("Id") + if not isinstance(raw_id, str): + raise AgentConfigApiError( + "Org Announcements API returned a saved announcement without an id" + ) + try: + canonical_id = validate_bulletin_id(raw_id) + except ValueError as error: + raise AgentConfigApiError( + "Org Announcements API returned a saved announcement without a valid id" + ) from error + if requested_id is not None and canonical_id != requested_id: + raise AgentConfigApiError( + "Org Announcements API returned a different announcement " + "than the one that was updated" + ) + return canonical_id + + async def list_bulletins(self, title_id: str) -> list[dict[str, Any]]: + """List every current item plus the most recent archived window.""" + payload = await self._request( + "GET", + f"{self._collection_path(title_id)}/ManagementView()", + transform_payload=False, + ) + return self._unwrap_collection(payload, title_id) + + async def get_bulletin(self, title_id: str, bulletin_id: str) -> dict[str, Any]: + """Load one canonical stored configuration.""" + validated_id = validate_bulletin_id(bulletin_id) + path = f"{self._collection_path(title_id)}({validated_id})" + return self._require_config( + await self._request("GET", path, transform_payload=False), + title_id, + expected_bulletin_id=validated_id, + ) + + async def save_bulletin( + self, title_id: str, payload: dict[str, Any] + ) -> dict[str, Any]: + """Create or update through the authoring OData ``Save`` action. + + Returns the canonical configuration from a keyed GET after validating + the ``EssBulletinSaveResult`` receipt. A payload without ``id`` is an + unkeyed create; replaying one after an ambiguous failure could duplicate + a committed record, so ambiguous retries are disabled and surfaced as + :class:`IndeterminateWriteError`. + """ + requested_id = ( + validate_bulletin_id(payload["id"]) if payload.get("id") else None + ) + is_create = not requested_id + try: + result = await self._request( + "POST", + f"{self._collection_path(title_id)}/Save", + json={"input": _to_odata_save_input(payload)}, + transform_payload=False, + idempotent=not is_create, + ) + except (AgentConfigApiError, httpx.RequestError) as error: + if is_create and _is_ambiguous_write_failure(error): + raise IndeterminateWriteError( + "The announcement may have been created but the response was " + "never received. Refresh before retrying." + ) from error + raise + saved_id = self._unwrap_save_result( + result, + requested_id=requested_id if not is_create else None, + ) + return await self._reload_committed_bulletin(title_id, saved_id) + + async def _reload_committed_bulletin( + self, title_id: str, bulletin_id: str + ) -> dict[str, Any]: + """Reload a committed row, retrying only short-lived keyed GET 404s.""" + for attempt in range(len(_COMMITTED_RELOAD_404_DELAYS) + 1): + try: + return await self.get_bulletin(title_id, bulletin_id) + except (AgentConfigApiError, httpx.RequestError) as error: + is_retryable_404 = ( + isinstance(error, AgentConfigApiError) + and error.http_status == 404 + and attempt < len(_COMMITTED_RELOAD_404_DELAYS) + ) + if not is_retryable_404: + raise CommittedCanonicalReloadError(error) from error + await asyncio.sleep(_COMMITTED_RELOAD_404_DELAYS[attempt]) + raise AssertionError("committed reload retry loop exhausted") + + async def transition_bulletin( + self, title_id: str, bulletin_id: str, status: str + ) -> Optional[dict[str, Any]]: + """Apply a minimal lifecycle status change. + + WeveNova performs the transition inside ``SaveAsync``: it loads the + canonical record server-side, preserves authored content and audience, + validates the transition, and writes the new status. Sending only + ``{id, status}`` therefore avoids a separate read/merge/write and cannot + clobber content with stale client state. + + The response is the same ``EssBulletinSaveResult`` receipt the content + save returns. Non-delete transitions are followed by a keyed GET; + delete is a tombstone and therefore has no readable canonical resource. + """ + validated_id = validate_bulletin_id(bulletin_id) + validated_status = _require_input_enum( + status, "status", _VALID_STATUSES + ) + saved_id = self._unwrap_save_result( + await self._request( + "POST", + f"{self._collection_path(title_id)}/Save", + json={ + "input": { + "Id": validated_id, + "Status": validated_status, + } + }, + transform_payload=False, + idempotent=True, + ), + requested_id=validated_id, + ) + if validated_status == STATUS_DELETED: + return None + return await self._reload_committed_bulletin(title_id, saved_id) + + +def _is_ambiguous_write_failure( + error: AgentConfigApiError | httpx.RequestError, +) -> bool: + """Decide whether a write failure leaves the commit outcome unknown. + + A transport error never reached a response, and a 502/503/504 came from an + intermediary that may have forwarded the request. A 4xx or a 500 from the + service itself is a definite rejection, so it stays a normal error. + """ + if isinstance(error, httpx.RequestError): + return True + return error.http_status in (502, 503, 504) + + +def build_manager_state( + items: list[dict[str, Any]], + audience_metadata: dict[str, list[dict[str, Any]]], + now: datetime, + *, + tenant_id: str, + title_id: str, +) -> dict[str, Any]: + """Adapt the API list into the widget's manager state. + + Soft-deleted rows are dropped before any counting: they are neither a + working item nor a restorable archived item, so including them would inflate + ``workingSetCount`` or push ``archivedTruncated`` true off records the maker + cannot see or act on. + + ``workingSetCount`` is the exact number of current items, computed from + status and schedule rather than from list position. ``archivedTruncated`` + reports only that the archived window filled, matching the widget copy + "Showing the 50 most recently archived announcements"; it never claims an + exact archived total. API item order is preserved. + """ + archived_count = 0 + working_set_count = 0 + view_models: list[dict[str, Any]] = [] + + for config in items: + if is_deleted_item(config): + continue + if is_archived_item(config, now): + archived_count += 1 + else: + working_set_count += 1 + bulletin_id = config.get("id", "") + view_models.append( + { + "config": config, + "audienceMetadata": audience_metadata.get(bulletin_id, []), + } + ) + + return { + "tenantId": tenant_id, + "titleId": title_id, + "items": view_models, + "workingSetCount": working_set_count, + "archivedTruncated": archived_count >= ARCHIVED_WINDOW_SIZE, + } diff --git a/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/drafts.py b/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/drafts.py new file mode 100644 index 000000000..a40b3ff54 --- /dev/null +++ b/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/drafts.py @@ -0,0 +1,508 @@ +# Copyright (c) Microsoft Corporation. +# Licensed under the MIT License. + +"""Typed Org Announcements draft and payload models. + +Two families live here: + +``SuggestedBulletinDraft`` + The *input* contract for a pre-hydrated create. It mirrors Vorpal's + ``suggestedBulletinDraftSchema`` exactly and is deliberately narrower than + the canonical record: it cannot carry a bulletin ID, a persisted lifecycle + status, creator/modifier identity, audit timestamps, or any backend version + or storage field. It carries both proposed content and editable copies, + including actions that need repair. Opening that unpublished client state + is never an authorization to write or proof of save/publish validity. + +``AnnouncementEditorDraft`` and the ``open_org_announcements`` payloads + The *output* contract consumed by the widget. ``build_create_draft`` + overlays a validated suggestion onto the editor defaults so an omitted + field keeps the exact value the empty editor would have shown. + +Field mapping between the two families is intentional and one-directional: +suggested ``priority`` becomes ``standardPriority`` and suggested +``secondaryAction`` becomes ``standardSecondaryAction``, matching the widget's +Standard-only editor fields. +""" + +from __future__ import annotations + +import re +from copy import deepcopy +from datetime import date, datetime, timezone +from typing import Any, Literal, Optional + +from pydantic import ( + BaseModel, + ConfigDict, + Field, + ModelWrapValidatorHandler, + PrivateAttr, + model_validator, +) + + +AnnouncementType = Literal["standard", "alert"] +AnnouncementPriority = Literal[0, 1] +AnnouncementStatus = Literal["draft", "published", "retired", "deleted"] +EditorMode = Literal["create", "edit", "duplicate"] +BulletinActionType = Literal["externalLink", "copilotChat"] + +# The values the empty Vorpal editor shows before a maker types anything. An +# omitted suggested field must land on exactly these, so they are named rather +# than repeated inline. +DEFAULT_TYPE: AnnouncementType = "standard" +DEFAULT_TITLE = "" +DEFAULT_DESCRIPTION = "" +DEFAULT_PRIMARY_ACTION = None +DEFAULT_START_DATE = "" +DEFAULT_END_DATE = "" +DEFAULT_STANDARD_PRIORITY: AnnouncementPriority = 1 +DEFAULT_STANDARD_SECONDARY_ACTION = None + +# ``EssBulletinInput.startDate``/``endDate`` are nullable ``DateTimeOffset`` +# UTC instants. The widget's DatePicker exchanges instants as +# ``YYYY-MM-DDTHH:MM:SS.mmmZ``, so every accepted suggested value is normalized +# into exactly that shape before it reaches the editor draft. +_INSTANT_FORMAT = "%Y-%m-%dT%H:%M:%S.%f" +_DATE_ONLY_PATTERN = re.compile(r"^\d{4}-\d{2}-\d{2}$") + + +def _format_instant(moment: datetime) -> str: + """Render an aware datetime as a UTC instant with millisecond precision.""" + utc = moment.astimezone(timezone.utc) + # ``%f`` is microseconds; truncate (never round) so a normalized instant is + # always at or before the value the maker supplied. + return f"{utc.strftime(_INSTANT_FORMAT)[:-3]}Z" + + +def normalize_suggested_instant(value: str, *, boundary: str) -> str: + """Normalize one suggested schedule boundary to a UTC millisecond instant. + + Three inputs are accepted, and nothing else: + + * ``""`` — an explicit "unset" that the DatePicker renders as empty. It is + preserved rather than defaulted so the widget shows its own validation. + * ``YYYY-MM-DD`` — a date-only value. A start becomes the first instant of + that UTC day and an end becomes the last, so a single-day announcement + spans the whole day instead of collapsing to midnight-to-midnight. + * A timezone-aware ISO-8601 instant — converted to UTC. + + A timezone-naive datetime is rejected rather than assumed to be UTC: the + maker's local day boundary is not knowable here, and silently guessing would + schedule an announcement at the wrong time. Natural language ("next + Monday") and any other unparseable text are rejected for the same reason — + a suggestion the model invented must never become a schedule nobody + reviewed. + """ + if not isinstance(value, str): + raise ValueError(f"{boundary} must be a string") + text = value.strip() + if not text: + return "" + + if _DATE_ONLY_PATTERN.match(text): + try: + day = date.fromisoformat(text) + except ValueError as error: + raise ValueError( + f"{boundary} must be a calendar date (YYYY-MM-DD) or a UTC instant" + ) from error + moment = ( + datetime(day.year, day.month, day.day, 0, 0, 0, 0, tzinfo=timezone.utc) + if boundary == "startDate" + else datetime( + day.year, day.month, day.day, 23, 59, 59, 999000, tzinfo=timezone.utc + ) + ) + return _format_instant(moment) + + candidate = f"{text[:-1]}+00:00" if text.endswith(("Z", "z")) else text + try: + parsed = datetime.fromisoformat(candidate) + except ValueError as error: + raise ValueError( + f"{boundary} must be an empty string, a calendar date " + f"(YYYY-MM-DD), or a timezone-aware ISO-8601 instant" + ) from error + if parsed.tzinfo is None or parsed.utcoffset() is None: + raise ValueError( + f"{boundary} must carry a timezone offset; a local datetime is " + f"ambiguous and is not assumed to be UTC" + ) + return _format_instant(parsed) + + +class StrictModel(BaseModel): + """Reject any field the contract does not name.""" + + model_config = ConfigDict(extra="forbid", hide_input_in_errors=True) + + +class BulletinAction(StrictModel): + """A primary or secondary announcement action. + + ``url`` belongs to ``externalLink`` and ``prompt`` to ``copilotChat``. A + payload carrying the *wrong* target for its type is malformed at the + contract level, so it is rejected here rather than forwarded. + + A *missing or blank* target is structurally representable so read-only + editors can retain content for repair. Mutation input also uses this model, + but WeveNova validates every present action on both Draft save and Publish. + Parsing a request here does not establish backend acceptance; service-owned + validation failures retain their field-level error codes. + """ + + actionType: BulletinActionType + label: str + url: Optional[str] = None + prompt: Optional[str] = None + + @model_validator(mode="after") + def require_matching_target(self) -> "BulletinAction": + if self.actionType == "externalLink": + if self.prompt is not None: + raise ValueError("externalLink actions must not carry a prompt") + elif self.url is not None: + raise ValueError("copilotChat actions must not carry a url") + return self + + +class SuggestedBulletinAction(BulletinAction): + """A structurally typed action in read-only proposed or copied content. + + A missing own target remains available for explicit repair in the editor. + Discriminators, field types, unknown fields and mismatched target members + still use :class:`BulletinAction` validation. Action readiness belongs to + save/publish validation, not to opening an unsaved working copy. + """ + + +class SuggestedBulletinDraft(StrictModel): + """Partial proposed or copied content for an unsaved create editor. + + Every field is optional; omitted fields fall back to the editor defaults in + :func:`build_create_draft`. An explicit empty string is a real draft value + and is preserved so the widget surfaces its normal validation instead of + silently substituting a default. Primary actions are retained for repair + even when incompatible with the selected announcement type. + """ + + type: Optional[AnnouncementType] = None + priority: Optional[AnnouncementPriority] = Field( + default=None, + description=( + "Standard announcement priority: 0 = Important; 1 = Informational. " + "Omitting priority defaults to Informational (1). Omit for Alert announcements." + ), + ) + title: Optional[str] = None + description: Optional[str] = None + primaryAction: Optional[SuggestedBulletinAction] = None + secondaryAction: Optional[SuggestedBulletinAction] = None + startDate: Optional[str] = None + endDate: Optional[str] = None + audience: Optional[list[str]] = None + + _retry_input: dict[str, Any] = PrivateAttr(default_factory=dict) + + @model_validator(mode="wrap") + @classmethod + def preserve_retry_input( + cls, value: Any, handler: ModelWrapValidatorHandler["SuggestedBulletinDraft"] + ) -> "SuggestedBulletinDraft": + # Keep the validated original proposal for retries; date normalization + # belongs to editor hydration, not to the maker's original request. + original = deepcopy(value) if isinstance(value, dict) else None + draft = handler(value) + if original is not None: + draft._retry_input = original + return draft + + def retry_payload(self) -> dict[str, Any]: + return deepcopy(self._retry_input) + + @model_validator(mode="after") + def normalize_schedule(self) -> "SuggestedBulletinDraft": + """Coerce every accepted schedule value into a UTC millisecond instant. + + Normalizing here — rather than in ``build_create_draft`` — means the + opener rejects an unusable suggestion before any Graph audience lookup + is issued, and the editor only ever receives instants the DatePicker + can render. + """ + if self.startDate is not None: + self.startDate = normalize_suggested_instant( + self.startDate, boundary="startDate" + ) + if self.endDate is not None: + self.endDate = normalize_suggested_instant( + self.endDate, boundary="endDate" + ) + return self + + @model_validator(mode="after") + def reject_alert_standard_only_fields(self) -> "SuggestedBulletinDraft": + """Keep the frontend's type-aware projection of Standard-only fields. + + Only an *explicit* incompatible value is rejected. Omitted Standard + fields still receive their required editor defaults, because the widget + keeps them in the draft even while they are hidden. A primary action, + unlike these Standard-only fields, stays visible for explicit repair. + """ + if self.type != "alert": + return self + if self.priority is not None: + raise ValueError("alert announcements do not support priority") + if self.secondaryAction is not None: + raise ValueError( + "alert announcements do not support a secondary action" + ) + return self + + @model_validator(mode="after") + def reject_blank_audience_ids(self) -> "SuggestedBulletinDraft": + if self.audience is None: + return self + for group_id in self.audience: + if not group_id or not group_id.strip(): + raise ValueError("audience group IDs must be non-empty") + return self + + +class OpenAnnouncementsRequest(StrictModel): + """Validate flat opener arguments before they can become widget retry state.""" + + titleId: str + view: Literal["manager", "editor"] + mode: Optional[Literal["create", "edit"]] = None + bulletinId: Optional[str] = None + suggestedDraft: Optional[SuggestedBulletinDraft] = None + + @model_validator(mode="after") + def require_valid_editor_intent(self) -> "OpenAnnouncementsRequest": + if self.view == "manager": + if any( + value is not None + for value in (self.mode, self.bulletinId, self.suggestedDraft) + ): + raise ValueError("manager does not take editor arguments") + elif self.mode is None: + raise ValueError("editor requires explicit create or edit mode") + elif self.mode == "create": + if self.bulletinId is not None: + raise ValueError("create must not supply bulletinId") + elif self.bulletinId is None: + raise ValueError("edit requires bulletinId") + elif self.suggestedDraft is not None: + raise ValueError("suggestedDraft applies only to create") + return self + + +class AudienceGroup(StrictModel): + """Display metadata for one canonical audience group ID. + + ``isValid`` is false when the ID could not be resolved through Graph or + resolved to a group category this kit does not author. The ID is always + retained so the widget can require removal or replacement instead of losing + stored state. + """ + + id: str + displayName: str + mail: Optional[str] = None + isValid: bool = True + + +class AnnouncementEditorDraft(StrictModel): + """The complete editor draft the widget renders.""" + + id: Optional[str] = None + type: AnnouncementType + title: str + description: str + primaryAction: Optional[BulletinAction] = None + startDate: str + endDate: str + audience: list[AudienceGroup] + standardPriority: AnnouncementPriority + standardSecondaryAction: Optional[BulletinAction] = None + + +def build_create_draft( + suggestion: Optional[SuggestedBulletinDraft], + audience_metadata: Optional[list[AudienceGroup]] = None, +) -> AnnouncementEditorDraft: + """Overlay a validated suggestion onto the empty-editor defaults. + + ``audience_metadata`` is the resolved, order-preserving metadata for + ``suggestion.audience``; the caller resolves it because resolution needs the + Graph client. A create draft never carries an ID. + """ + if suggestion is None: + suggestion = SuggestedBulletinDraft() + + return AnnouncementEditorDraft( + type=suggestion.type if suggestion.type is not None else DEFAULT_TYPE, + title=suggestion.title if suggestion.title is not None else DEFAULT_TITLE, + description=( + suggestion.description + if suggestion.description is not None + else DEFAULT_DESCRIPTION + ), + primaryAction=( + suggestion.primaryAction + if suggestion.primaryAction is not None + else DEFAULT_PRIMARY_ACTION + ), + startDate=( + suggestion.startDate + if suggestion.startDate is not None + else DEFAULT_START_DATE + ), + endDate=( + suggestion.endDate + if suggestion.endDate is not None + else DEFAULT_END_DATE + ), + audience=list(audience_metadata or []), + standardPriority=( + suggestion.priority + if suggestion.priority is not None + else DEFAULT_STANDARD_PRIORITY + ), + standardSecondaryAction=( + suggestion.secondaryAction + if suggestion.secondaryAction is not None + else DEFAULT_STANDARD_SECONDARY_ACTION + ), + ) + + +def build_editor_draft_from_config( + config: dict[str, Any], + audience_metadata: list[AudienceGroup], +) -> AnnouncementEditorDraft: + """Project canonical content into the normal editor, preserving its schedule.""" + bulletin = config.get("bulletin") + if not isinstance(bulletin, dict): + raise ValueError("bulletin configuration is missing its bulletin content") + + bulletin_id = config.get("id") + if not isinstance(bulletin_id, str) or not bulletin_id: + raise ValueError("bulletin configuration is missing its id") + + primary_action = bulletin.get("primaryAction") + secondary_action = bulletin.get("secondaryAction") + priority = bulletin.get("priority") + + return AnnouncementEditorDraft( + id=bulletin_id, + type=bulletin.get("type") or DEFAULT_TYPE, + title=bulletin.get("title") or DEFAULT_TITLE, + description=bulletin.get("description") or DEFAULT_DESCRIPTION, + primaryAction=( + BulletinAction.model_validate(primary_action) + if isinstance(primary_action, dict) + else DEFAULT_PRIMARY_ACTION + ), + startDate=bulletin.get("startDate") or DEFAULT_START_DATE, + endDate=bulletin.get("endDate") or DEFAULT_END_DATE, + audience=list(audience_metadata), + standardPriority=( + priority if priority in (0, 1) else DEFAULT_STANDARD_PRIORITY + ), + standardSecondaryAction=( + BulletinAction.model_validate(secondary_action) + if isinstance(secondary_action, dict) + else DEFAULT_STANDARD_SECONDARY_ACTION + ), + ) + + +_SCHEDULE_FIELDS = ("startDate", "endDate") + + +def is_blank_schedule_value(value: Any) -> bool: + """Report whether a schedule value is the widget's "unset" sentinel. + + The single definition of the rule, shared by the validated save path and the + duplicate path, so the two cannot drift into disagreeing about what "unset" + means. + """ + return isinstance(value, str) and not value.strip() + + +def without_blank_schedule(bulletin: dict[str, Any]) -> dict[str, Any]: + """Drop schedule keys holding a blank sentinel from a raw content dict. + + Used by the duplicate path, which forwards *stored* content verbatim and so + never passes through :class:`BulletinInput`. Without this, duplicating a + source whose schedule is blank would put ``""`` on the wire and fail the + backend's model binding, even though the equivalent ordinary save is fine. + """ + return { + key: value + for key, value in bulletin.items() + if not (key in _SCHEDULE_FIELDS and is_blank_schedule_value(value)) + } + + +class BulletinInput(StrictModel): + """The authored content a mutation sends to the authoring API. + + ``startDate``/``endDate`` map to nullable ``DateTimeOffset`` on + ``EssBulletinInput``. The widget's DatePicker represents "unset" as the + empty string, so a blank arriving here is coerced to ``None`` and then + dropped entirely by ``exclude_none=True`` at serialization. Sending ``""`` + instead would fail the backend's *model binding* — a JSON string is not a + DateTimeOffset — which surfaces as an opaque 400 rather than as the + field-level validation error the widget knows how to render, and would make + a perfectly ordinary unscheduled Draft unsaveable. + """ + + type: AnnouncementType + priority: Optional[AnnouncementPriority] = None + title: str + description: str + primaryAction: Optional[BulletinAction] = None + secondaryAction: Optional[BulletinAction] = None + startDate: Optional[str] = None + endDate: Optional[str] = None + + @model_validator(mode="after") + def blank_schedule_means_unset(self) -> "BulletinInput": + """Treat a blank or whitespace-only boundary as absent, not as a value. + + Only blanks are touched. A real instant is passed through byte-for-byte: + the backend owns the canonical schedule, and silently rewriting an + instant here could shift a schedule the maker already reviewed. + """ + if self.startDate is not None and is_blank_schedule_value(self.startDate): + self.startDate = None + if self.endDate is not None and is_blank_schedule_value(self.endDate): + self.endDate = None + return self + + +class SaveBulletinRequest(StrictModel): + """A complete create-or-update of authored content and audience. + + ``id`` present means update, ``id`` absent means create. The widget always + sends the whole authored state, so this model never merges with server + state. Agent identity belongs to the route, never to this HTTP body. + """ + + id: Optional[str] = None + bulletin: BulletinInput + audience: list[str] + status: Literal["draft", "published"] + + @model_validator(mode="after") + def reject_blank_identifiers(self) -> "SaveBulletinRequest": + if self.id is not None and not self.id.strip(): + raise ValueError("id must be a non-empty string when provided") + for group_id in self.audience: + if not group_id or not group_id.strip(): + raise ValueError("audience group IDs must be non-empty") + return self diff --git a/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/graph_directory_client.py b/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/graph_directory_client.py new file mode 100644 index 000000000..baa592def --- /dev/null +++ b/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/graph_directory_client.py @@ -0,0 +1,699 @@ +# Copyright (c) Microsoft Corporation. +# Licensed under the MIT License. + +"""Microsoft Graph directory client for audience discovery and metadata. + +Graph is a *discovery and metadata* source only. WeveNova remains authoritative +for tenant-local existence and save validity; this client decides which group +categories the kit is willing to author against, which is deliberately narrower +than everything the backend might accept. + +Authentication is delegated MSAL using the established Graph Command Line Tools +public client and the ADK ``.local/.token_cache.bin`` cache, anchored to the +solution root so it is the *same* cache `/setup` populates. It is a separate +delegated resource from the AgentConfiguration/WeveNova token: no Graph scope is +added to that token, and no Azure CLI credential is used. + +Each instance is bound to the captured authoring tenant and account, not to an +ESS titleId. Resource tokens remain separate; the service owns authorization. +""" + +from __future__ import annotations + +import asyncio +import logging +import os +import sys +import threading +import time +from typing import Any, Iterable, Optional +from uuid import UUID + +import httpx +from portalocker.exceptions import LockException + +sys.path.insert( + 0, + os.path.join(os.path.dirname(os.path.abspath(__file__)), "..", "agentconfig_core"), +) + +from base_client import account_for_identity, validate_token_identity # noqa: E402 +from _token_cache import create_token_cache # noqa: E402 + + +logging.getLogger("httpx").setLevel(logging.WARNING) +logging.getLogger("httpcore").setLevel(logging.WARNING) + +GRAPH_BASE_URL = "https://graph.microsoft.com/v1.0" + +# The established Microsoft Graph Command Line Tools public client, reused so +# audience discovery needs no new app registration or preauthorization. This is +# the same first-party client the kit's `/setup` and FlightCheck Graph paths +# use, which is what makes a silent (prompt-free) acquisition possible here. +GRAPH_CLIENT_ID = "14d82eec-204b-4c2f-b7e8-296a70dab67e" +GRAPH_SCOPES = ["https://graph.microsoft.com/Directory.Read.All"] + +# The one delegated Graph/Dataverse cache the whole ADK shares. It is anchored +# to the *solution root* derived from this module's own location, never to the +# process working directory: this server launches with cwd set to its own MCP +# folder, so a cwd-relative ".local/.token_cache.bin" would mint a third, +# private cache and force a second interactive sign-in for a maker who already +# signed in through `/setup`. +# .../solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/ +# parents: [3] [2] [1] [0] +_SOLUTION_ROOT = os.path.abspath( + os.path.join(os.path.dirname(os.path.abspath(__file__)), "..", "..", "..") +) +_LOCAL_STATE_DIR = os.path.join(_SOLUTION_ROOT, ".local") +GRAPH_TOKEN_CACHE_PATH = os.path.join(_LOCAL_STATE_DIR, ".token_cache.bin") + +MAX_ELIGIBLE_RESULTS = 20 +MAX_SEARCH_PAGES = 3 +GRAPH_PAGE_SIZE = 50 +MAX_QUERY_LENGTH = 256 + +# Graph caps ``$filter`` id lists; batch lookups stay well inside that bound. +_METADATA_BATCH_SIZE = 15 + +# Every mutation refreshes manager state, which re-resolves the audience of +# every listed bulletin. Without a cache that is one Graph round trip per batch +# per mutation for metadata that changes rarely. The TTL is short so a renamed +# or deleted group self-corrects quickly, and WeveNova — not this cache — +# remains authoritative for save validity. +GROUP_METADATA_TTL_SECONDS = 300.0 +GROUP_METADATA_CACHE_MAX_ENTRIES = 512 + +# Shown when a group carries neither a display name nor a mail address. A raw +# directory ID is never used as a name: it is not human-meaningful and it would +# put a tenant group identifier into the widget's visible surface. +UNNAMED_GROUP_LABEL = "Unnamed group" + +_GROUP_SELECT = "id,displayName,mail,mailEnabled,securityEnabled,groupTypes" + +# Identifier syntax alone cannot distinguish a provider code from private +# account details. Only these fixed diagnostic categories may reach the log. +_SAFE_MSAL_ERROR_CODES = frozenset( + {"access_denied", "interaction_required", "consent_required"} +) + + +class GraphDirectoryError(RuntimeError): + """Raised when a Graph directory call fails in a way the caller must surface.""" + + def __init__( + self, + message: str, + *, + code: str, + retryable: bool = False, + http_status: Optional[int] = None, + ): + super().__init__(message) + self.code = code + self.retryable = retryable + self.http_status = http_status + + +class GraphAuthenticationError(GraphDirectoryError): + """Raised when no delegated Graph token could be acquired.""" + + def __init__(self, message: str): + super().__init__(message, code="AuthenticationRequired", retryable=False) + + +def escape_search_value(query: str) -> str: + """Escape a maker-supplied value for a Graph ``$search`` phrase. + + The value is embedded inside a double-quoted ``displayName:`` phrase, + so a raw backslash or double quote would terminate or corrupt the phrase. + Both are backslash-escaped; control characters are rejected outright rather + than silently stripped, because a stripped query would search for something + the maker did not type. + """ + if not isinstance(query, str): + raise ValueError("query must be a string") + normalized = query.strip() + if not normalized: + raise ValueError("query must be a non-empty string") + if len(normalized) > MAX_QUERY_LENGTH: + raise ValueError(f"query must not exceed {MAX_QUERY_LENGTH} characters") + if any( + ord(character) < 0x20 or ord(character) == 0x7F for character in normalized + ): + raise ValueError("query must not contain control characters") + return normalized.replace("\\", "\\\\").replace('"', '\\"') + + +def is_eligible_group(group: Any) -> bool: + """Apply the MVP audience eligibility policy. + + Included: non-dynamic security groups, mail-enabled security groups, and + classic distribution groups. Excluded: Microsoft 365/Unified groups and any + group whose ``groupTypes`` contains ``DynamicMembership``. Exchange dynamic + distribution lists are not returned by the Graph groups collection at all, + so they need no predicate branch. + + This is the single eligibility decision in the server; hydration and search + both call it so a stored group and a searched group are judged identically. + """ + if not isinstance(group, dict): + return False + if not isinstance(group.get("id"), str) or not group["id"]: + return False + + group_types = group.get("groupTypes") + if not isinstance(group_types, list): + return False + lowered = {value.lower() for value in group_types if isinstance(value, str)} + if "unified" in lowered or "dynamicmembership" in lowered: + return False + + security_enabled = group.get("securityEnabled") is True + mail_enabled = group.get("mailEnabled") is True + # Security group (with or without mail), or a classic distribution group. + return security_enabled or mail_enabled + + +def to_audience_group(group: dict[str, Any]) -> dict[str, Any]: + """Project a Graph group into the widget's audience metadata shape. + + The display name falls back through ``displayName`` → ``mail`` → + :data:`UNNAMED_GROUP_LABEL`. The raw directory ID is deliberately *not* a + fallback: rendering it as a name shows an opaque GUID where the maker + expects a group, and it surfaces a tenant directory identifier in the + widget and in anything that copies the label. + """ + mail = group.get("mail") + mail = mail.strip() if isinstance(mail, str) else "" + display_name = group.get("displayName") + display_name = display_name.strip() if isinstance(display_name, str) else "" + return { + "id": group["id"], + "displayName": display_name or mail or UNNAMED_GROUP_LABEL, + "mail": mail or None, + "isValid": True, + } + + +def _validate_graph_token(token: str, tenant_id: str, object_id: Optional[str]) -> str: + try: + return validate_token_identity(token, tenant_id, object_id) + except ValueError as error: + raise GraphAuthenticationError( + "Microsoft Graph sign-in must identify the current authoring tenant " + "and account. Use matching delegated credentials with tid and oid claims." + ) from error + + +def acquire_graph_token(tenant_id: str, object_id: Optional[str]) -> str: + """Acquire a delegated Graph token from the shared ADK cache. + + Silent first, so a maker who already signed in for `/setup` is not prompted + again; the Graph CLI public client is first-party and FOCI, so a prior + Dataverse sign-in usually redeems silently. Only when that fails does this + fall back to MSAL's interactive browser flow. + + Device code is deliberately *not* used. An MCP server speaks JSON-RPC over + stdio, so printing a device code to stdout corrupts the protocol stream, and + the maker is already in a desktop session where a browser can open. + + This is blocking (MSAL is synchronous and the interactive flow waits on a + human). Callers must run it off the event loop. + """ + tenant_id = str(UUID(tenant_id)) + if not object_id: + raise GraphAuthenticationError( + "The authoring account cannot be identified. Audience lookup requires " + "an AgentConfiguration delegated token with an oid claim." + ) + token = os.environ.get("GRAPH_ACCESS_TOKEN", "").strip() + if token: + return _validate_graph_token(token, tenant_id, object_id) + + import msal + + cache = create_token_cache(GRAPH_TOKEN_CACHE_PATH) + app = msal.PublicClientApplication( + GRAPH_CLIENT_ID, + authority=f"https://login.microsoftonline.com/{tenant_id}", + token_cache=cache, + ) + + account = account_for_identity(app, tenant_id, object_id) + result = ( + app.acquire_token_silent(GRAPH_SCOPES, account=account) + if account is not None + else None + ) + if not result or "access_token" not in result: + # Logged, never printed: stdout is the MCP transport. + logging.getLogger("ess-org-announcements.graph").info( + "Opening a browser for Microsoft Graph sign-in (audience discovery)." + ) + result = app.acquire_token_interactive( + GRAPH_SCOPES, prompt="select_account" + ) + + if not isinstance(result, dict) or "access_token" not in result: + error_code = result.get("error") if isinstance(result, dict) else None + if not isinstance(error_code, str) or error_code not in _SAFE_MSAL_ERROR_CODES: + error_code = "unknown_error" + logging.getLogger("ess-org-announcements.graph").warning( + "Microsoft Graph sign-in failed (%s).", error_code + ) + raise GraphAuthenticationError( + "Microsoft Graph sign-in failed. Sign in with the authoring account." + ) + return _validate_graph_token(result["access_token"], tenant_id, object_id) + + +class _GroupMetadataCache: + """A bounded, TTL'd in-process cache of resolved group metadata. + + Only *successful* resolutions are cached. A negative entry would pin a group + that was just created, just consented to, or just made eligible as invalid + for the whole TTL, which the maker would experience as an audience they + cannot add; re-querying the small set of unresolved IDs is the cheaper + mistake. + + Insertion-ordered with FIFO eviction so a long-running server cannot grow + without bound. Guarded by a lock because ``resolve_groups`` may be awaited + concurrently. + """ + + def __init__( + self, + *, + ttl_seconds: float = GROUP_METADATA_TTL_SECONDS, + max_entries: int = GROUP_METADATA_CACHE_MAX_ENTRIES, + ): + self._ttl = ttl_seconds + self._max_entries = max_entries + self._entries: dict[str, tuple[float, dict[str, Any]]] = {} + self._lock = threading.Lock() + + def _now(self) -> float: + # Monotonic: a wall-clock adjustment must not resurrect or expire + # entries early. + return time.monotonic() + + def get(self, group_id: str) -> Optional[dict[str, Any]]: + with self._lock: + entry = self._entries.get(group_id) + if entry is None: + return None + expires_at, metadata = entry + if expires_at <= self._now(): + del self._entries[group_id] + return None + return dict(metadata) + + def put(self, group_id: str, metadata: dict[str, Any]) -> None: + with self._lock: + expires_at = self._now() + self._ttl + # Re-insert so a refreshed entry moves to the newest position and + # FIFO eviction stays meaningful. + self._entries.pop(group_id, None) + self._entries[group_id] = (expires_at, dict(metadata)) + while len(self._entries) > self._max_entries: + self._entries.pop(next(iter(self._entries))) + + def clear(self) -> None: + with self._lock: + self._entries.clear() + + +class GraphDirectoryClient: + """Bounded, deterministic Graph group search and metadata hydration.""" + + def __init__( + self, + *, + tenant_id: str, + object_id: Optional[str], + transport: Optional[httpx.AsyncBaseTransport] = None, + access_token: Optional[str] = None, + metadata_cache: Optional["_GroupMetadataCache"] = None, + ): + self._logger = logging.getLogger("ess-org-announcements.graph") + self._tenant_id = str(UUID(tenant_id)) + self._object_id = object_id + # Deliberately NOT acquired here. Constructing the client must stay a + # pure, non-blocking operation: the server builds it lazily from a tool + # call, and acquiring a token in __init__ would run a blocking MSAL + # sign-in (and possibly open a browser) inside the asyncio event loop, + # stalling every other in-flight request and the MCP stdio transport. + self._token = ( + _validate_graph_token(access_token.strip(), self._tenant_id, self.object_id) + if access_token else "" + ) + self._token_lock = asyncio.Lock() + self._transport = transport + self._client: Optional[httpx.AsyncClient] = None + self._metadata_cache = ( + metadata_cache if metadata_cache is not None else _GroupMetadataCache() + ) + self.timeout = 30.0 + + def __repr__(self) -> str: + return f"<{type(self).__name__} base_url={GRAPH_BASE_URL!r}>" + + @property + def tenant_id(self) -> str: + """The captured authoring tenant; a tenant change requires a new instance.""" + return self._tenant_id + + @property + def object_id(self) -> Optional[str]: + """The captured account; changing identity requires a separate cache/client.""" + return self._object_id + + async def _access_token(self) -> str: + """Return the delegated token, acquiring it once, off the event loop. + + ``acquire_graph_token`` is synchronous and can block for as long as a + human takes to complete an interactive sign-in, so it runs in a worker + thread. The lock makes concurrent first calls share one sign-in instead + of racing two browser prompts. + """ + if self._token: + return self._token + async with self._token_lock: + if not self._token: + try: + token = await asyncio.to_thread( + acquire_graph_token, self.tenant_id, self.object_id + ) + except (LockException, OSError) as error: + # Keep persistence failures inside the feature's normal error + # contract, especially when hydration follows a committed write. + raise GraphAuthenticationError( + "Microsoft Graph sign-in could not be completed. Try again." + ) from error + self._token = _validate_graph_token(token, self.tenant_id, self.object_id) + return self._token + + async def _ensure_client(self) -> httpx.AsyncClient: + if self._client is None or self._client.is_closed: + token = await self._access_token() + # Another first caller may have constructed it while we awaited auth. + if self._client is None or self._client.is_closed: + self._client = httpx.AsyncClient( + base_url=GRAPH_BASE_URL, + headers={ + "Authorization": f"Bearer {token}", + "Accept": "application/json", + "ConsistencyLevel": "eventual", + }, + timeout=self.timeout, + verify=True, + transport=self._transport, + follow_redirects=False, + ) + return self._client + + async def aclose(self) -> None: + async with self._token_lock: + client = self._discard_credentials() + if client is not None and not client.is_closed: + await client.aclose() + + def _discard_credentials(self) -> Optional[httpx.AsyncClient]: + self._token = "" + client, self._client = self._client, None + # A late request may still hold this old cache. Replace, rather than + # merely clear it, so its completion cannot populate fresh credentials. + self._metadata_cache.clear() + self._metadata_cache = _GroupMetadataCache() + return client + + async def _invalidate_credentials(self, failed_token: str) -> None: + """Drop an access token Graph has rejected, and the client holding it. + + A cached access token eventually expires. The MCP server keeps one + process-global Graph client, so without this a single 401 would poison + every later audience lookup until the server was restarted — the maker's + only recovery would be reloading the whole MCP host. + + The clear is a compare-and-swap against ``failed_token``: a concurrent + caller may already have completed a fresh acquisition between the + request going out and the 401 coming back, and wiping *that* token would + throw away a good credential and force a needless second sign-in. + + The ``httpx.AsyncClient`` is dropped too, not just the token, because + the bearer is baked into its default headers at construction: keeping it + would keep sending the dead token no matter what ``_token`` says. + """ + async with self._token_lock: + if failed_token and self._token != failed_token: + # Someone already refreshed; their token is still good. + return + client = self._discard_credentials() + + # Closed outside the lock: aclose() awaits connection teardown and must + # not hold up a concurrent re-acquisition. + if client is not None and not client.is_closed: + await client.aclose() + + async def _get( + self, url: str, params: Optional[dict[str, Any]] = None, + *, client: Optional[httpx.AsyncClient] = None, + ) -> dict[str, Any]: + if client is None: + client = await self._ensure_client() + if client.is_closed: + raise GraphAuthenticationError( + "Microsoft Graph sign-in was reset. Try the request again to sign in." + ) + request_token = client.headers["Authorization"].removeprefix("Bearer ") + try: + response = await client.get(url, params=params) + response.raise_for_status() + payload = response.json() + except httpx.HTTPStatusError as error: + status = error.response.status_code + if status == 401: + # Expired or revoked credential. Invalidate so the *next* call + # re-runs silent-first acquisition, which usually refreshes from + # the shared cache with no prompt at all. + # + # The failed request is deliberately NOT replayed here. These + # are idempotent GETs so a replay would be safe, but + # re-acquisition can fall through to an interactive browser + # sign-in, and blocking an in-flight tool call on a surprise + # prompt is worse than returning a clear, actionable + # AuthenticationRequired the maker can act on. The next call + # succeeds silently in the common case. + await self._invalidate_credentials(request_token) + raise GraphDirectoryError( + "Microsoft Graph rejected the cached sign-in. Try the " + "request again to sign in.", + code="AuthenticationRequired", + retryable=False, + http_status=status, + ) from error + if status == 403: + raise GraphDirectoryError( + "Microsoft Graph denied the directory request. " + "Directory.Read.All requires tenant consent and may be " + "blocked by Conditional Access.", + code="AuthorizationDenied", + retryable=False, + http_status=status, + ) from error + raise GraphDirectoryError( + f"Microsoft Graph returned HTTP {status}.", + code="SearchUnavailable", + retryable=status in (429, 502, 503, 504), + http_status=status, + ) from error + except httpx.RequestError as error: + raise GraphDirectoryError( + "Microsoft Graph could not be reached.", + code="NetworkError", + retryable=True, + ) from error + + if not isinstance(payload, dict) or not isinstance(payload.get("value"), list): + raise GraphDirectoryError( + "Microsoft Graph returned an unexpected directory response.", + code="SearchUnavailable", + retryable=False, + ) + return payload + + async def search_groups(self, query: str) -> dict[str, Any]: + """Search groups, returning at most 20 eligible, deduplicated results. + + Paging stops at 20 eligible results, at result-set exhaustion, or after + three pages, whichever comes first. ``exhausted`` distinguishes an + exhausted result set from a search that stopped at the page cap, so the + caller never reports "no more matches" for a capped search. Once the + result limit is reached, the remainder of that already-fetched page is + inspected only to distinguish exact exhaustion from omitted matches. + """ + escaped = escape_search_value(query) + client = await self._ensure_client() + params: Optional[dict[str, Any]] = { + "$search": f'"displayName:{escaped}" OR "mail:{escaped}"', + "$select": _GROUP_SELECT, + "$count": "true", + "$top": str(GRAPH_PAGE_SIZE), + } + url = "/groups" + + groups: list[dict[str, Any]] = [] + seen: set[str] = set() + pages = 0 + exhausted = False + + while pages < MAX_SEARCH_PAGES: + payload = await self._get(url, params, client=client) + pages += 1 + has_unreturned_eligible_group = False + for item in payload["value"]: + if not is_eligible_group(item): + continue + if item["id"] in seen: + continue + seen.add(item["id"]) + # Graph's relevance order is preserved; dedupe keeps the first + # occurrence so repeated runs return the same ordering. + if len(groups) < MAX_ELIGIBLE_RESULTS: + groups.append(to_audience_group(item)) + else: + has_unreturned_eligible_group = True + + next_link = payload.get("@odata.nextLink") + has_next_page = isinstance(next_link, str) and bool(next_link) + if len(groups) >= MAX_ELIGIBLE_RESULTS: + return { + "groups": groups, + "exhausted": ( + not has_unreturned_eligible_group and not has_next_page + ), + "pagesExamined": pages, + } + if not has_next_page: + exhausted = True + break + # The nextLink already carries every query option; re-sending params + # would duplicate them. + url = next_link + params = None + + return { + "groups": groups, + "exhausted": exhausted, + "pagesExamined": pages, + } + + async def resolve_groups( + self, group_ids: Iterable[str] + ) -> dict[str, dict[str, Any]]: + """Resolve display metadata for a deduplicated set of group IDs. + + Deduplication applies only to the lookup batch. The caller re-expands + the result against each bulletin's canonical audience list so per-item + order and multiplicity stay exact. + + Results already held in the bounded TTL cache are served without a Graph + round trip; only the remaining IDs are batched. A Graph failure + propagates and leaves the cache untouched, so a transient outage cannot + poison later lookups. + + Any ID that Graph does not return, or that resolves to an ineligible + group category, is omitted from the result. The caller marks it + ``isValid: false`` and keeps the ID. + """ + unique = [ + group_id + for group_id in dict.fromkeys(group_ids) + if isinstance(group_id, str) and group_id + ] + if not unique: + return {} + + client = await self._ensure_client() + metadata_cache = self._metadata_cache + + resolved: dict[str, dict[str, Any]] = {} + pending: list[str] = [] + for group_id in unique: + cached = metadata_cache.get(group_id) + if cached is None: + pending.append(group_id) + else: + resolved[group_id] = cached + + for start in range(0, len(pending), _METADATA_BATCH_SIZE): + batch = pending[start : start + _METADATA_BATCH_SIZE] + clause = " or ".join( + f"id eq '{group_id}'" + for group_id in batch + if _is_safe_filter_id(group_id) + ) + if not clause: + continue + url = "/groups" + params: Optional[dict[str, Any]] = { + "$filter": clause, + "$select": _GROUP_SELECT, + "$top": str(len(batch)), + } + while True: + payload = await self._get(url, params, client=client) + for item in payload["value"]: + if is_eligible_group(item): + metadata = to_audience_group(item) + resolved[item["id"]] = metadata + metadata_cache.put(item["id"], metadata) + + next_link = payload.get("@odata.nextLink") + if not isinstance(next_link, str) or not next_link: + break + url = next_link + params = None + + return resolved + + +def _is_safe_filter_id(group_id: str) -> bool: + """Reject any ID that could break out of an OData string literal. + + Directory object IDs are GUIDs, so a value containing a quote, a control + character, or a backslash is not a real ID; skipping it keeps the filter + well-formed and the ID is still reported as invalid downstream. + """ + if "'" in group_id or "\\" in group_id: + return False + return not any( + ord(character) < 0x20 or ord(character) == 0x7F for character in group_id + ) + + +def build_audience_metadata( + audience_ids: list[str], resolved: dict[str, dict[str, Any]] +) -> list[dict[str, Any]]: + """Expand canonical audience IDs into order-preserving display metadata. + + The output is one-to-one with ``audience_ids`` and in the same order, which + the widget's schema enforces. An unresolved or ineligible ID keeps its + position and is marked invalid; its display name falls back to a neutral + placeholder so a raw group ID is never rendered as a name. + """ + metadata: list[dict[str, Any]] = [] + for group_id in audience_ids: + group = resolved.get(group_id) + if group is None: + metadata.append( + { + "id": group_id, + "displayName": "Unavailable group", + "mail": None, + "isValid": False, + } + ) + else: + metadata.append(dict(group)) + return metadata diff --git a/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/requirements.txt b/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/requirements.txt new file mode 100644 index 000000000..dcc23755e --- /dev/null +++ b/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/requirements.txt @@ -0,0 +1,6 @@ +mcp>=1.29.0,<2.0.0 +anyio>=4.0,<5.0 +httpx>=0.27.0,<1.0 +msal>=1.35.0 +pydantic>=2.0,<3.0 +-r ../agentconfig_core/requirements.txt diff --git a/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/server.py b/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/server.py new file mode 100644 index 000000000..05f5c9c09 --- /dev/null +++ b/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/server.py @@ -0,0 +1,1622 @@ +# Copyright (c) Microsoft Corporation. +# Licensed under the MIT License. + +"""ESS Org Announcements MCP server. + +Announcements are scoped to the authenticated tenant and deployed ESS titleId. +The tenant comes exclusively from the AgentConfiguration access token; each +request carries its agent identity through every read, write, and refresh. + +Tool visibility is deliberate: + +* ``open_org_announcements`` is read-only and model/app-visible. It is the single + entry point a maker turn may use, and the widget uses it for scoped navigation. +* ``save_bulletin``, ``transition_bulletin``, and ``duplicate_bulletin`` are + app-visible only. Once the widget is open it owns the editing session, so the + model must not issue a duplicate write. +* ``search_audience_groups`` is read-only and visible to both, because the skill + needs it to resolve explicit audience names before an opener call. + +Nothing here logs announcement content, group identifiers, group names, group +mail, tokens, claims, or the opener request payload. +""" + +from __future__ import annotations + +import asyncio +import html +import json +import logging +import os +import time +from collections.abc import AsyncIterator +from contextlib import asynccontextmanager +from datetime import datetime, timezone +from typing import Any, Literal, Optional +from urllib.parse import urlsplit + +import anyio +import httpx +from mcp.server.fastmcp import FastMCP +from mcp.server.fastmcp.exceptions import ToolError +from mcp.types import CallToolResult, TextContent, ToolAnnotations +from portalocker.exceptions import LockException +from pydantic import ValidationError + +from client import ( + AgentConfigApiError, + BulletinValidationError, + CommittedCanonicalReloadError, + IndeterminateWriteError, + OrgAnnouncementsClient, + STATUS_DELETED, + build_manager_state, + is_deleted_item, +) +from base_client import LocalCredentialError +from drafts import ( + AnnouncementEditorDraft, + AudienceGroup, + EditorMode, + SaveBulletinRequest, + SuggestedBulletinDraft, + OpenAnnouncementsRequest, + build_create_draft, + build_editor_draft_from_config, + without_blank_schedule, +) +from graph_directory_client import ( + GraphDirectoryClient, + GraphDirectoryError, + build_audience_metadata, + escape_search_value, +) +from telemetry import ( + SOURCE_BACKEND, + SOURCE_GRAPH, + SOURCE_MCP, + normalize_error_code, + record_operation, +) +from validation import validate_bulletin_id, validate_title_id + + +DEFAULT_WIDGET_ORIGIN = "https://workforceinsights.m365.cloud.microsoft" +ALLOWED_WIDGET_ORIGINS = frozenset( + { + "https://workforceinsights.m365.cloud.dev.microsoft", + "https://df.workforceinsights.m365.cloud.microsoft", + DEFAULT_WIDGET_ORIGIN, + } +) +DEVELOPMENT_WIDGET_ORIGIN_ENV = "VORPAL_WIDGET_ALLOW_DEVELOPMENT_ORIGIN" +WIDGET_MIME_TYPE = "text/html;profile=mcp-app" +ORG_ANNOUNCEMENTS_RESOURCE_URI = ( + "ui://widget/org-announcements/OrgAnnouncements.html" +) + +_LOGGER = logging.getLogger("ess-org-announcements") + +_READ_ONLY_ANNOTATIONS = ToolAnnotations( + readOnlyHint=True, + destructiveHint=False, + idempotentHint=True, + openWorldHint=False, +) +_MUTATION_ANNOTATIONS = ToolAnnotations( + readOnlyHint=False, + destructiveHint=False, + idempotentHint=False, + openWorldHint=False, +) +_DESTRUCTIVE_ANNOTATIONS = ToolAnnotations( + readOnlyHint=False, + destructiveHint=True, + idempotentHint=True, + openWorldHint=False, +) + +# Widget transition -> stored status. Publish now is intentionally absent: the +# scheduled row already holds canonical content and audience in the widget, so +# Vorpal calls save_bulletin with that complete state and only moves startDate +# to the current UTC instant. Routing it here would need a server-side +# GET-then-POST and would race a concurrent edit. +TRANSITION_STATUS = { + "archive": "retired", + "unarchive": "draft", + "moveToDraft": "draft", + "delete": STATUS_DELETED, +} +TransitionName = Literal["archive", "unarchive", "moveToDraft", "delete"] + +# Backend codes that mean the tenant's OrgAnnouncementsSettings gate is off or +# the authoring surface is not deployed. These are reported as a recoverable +# unavailable state; the server never falls back to mock persistence. +# +# The list is intentionally exact-match and conservative. A substring or prefix +# rule would swallow a genuine field-level validation code that merely mentions +# a feature, turning a fixable "this value is wrong" into a dead-end "your +# tenant is not enabled" that the maker cannot act on. +_FEATURE_DISABLED_CODES = frozenset( + { + "FeatureDisabled", + "FeatureNotEnabled", + "OrgAnnouncementsDisabled", + "NotSupported", + } +) + +# The neutral client core synthesizes this code when a failing response carried +# no backend ``Code`` at all. It is a transport placeholder, NOT a backend +# validation code, so it must never be surfaced as one: the widget would look up +# localized copy for a code the backend never emitted, and a 500 would be +# reported to the maker as if they had typed something wrong. +_FALLBACK_BACKEND_CODE = "HttpError" + +# Backend failure copy safe for the model-visible opener. The opener rebuilds +# even recognized failures from this local map rather than trusting the message +# paired with a familiar backend code. Every other backend code/detail remains +# available only to app-visible mutation tools. +_MODEL_VISIBLE_BACKEND_MESSAGES = { + "AuthenticationRequired": "Sign in again to author organization announcements.", + "AuthorizationDenied": ( + "This account is not authorized to author organization " + "announcements in this tenant." + ), + "FeatureUnavailable": "Organization announcements are not enabled for this tenant.", + "NetworkError": "The Org Announcements service is temporarily unavailable.", + "NotFound": "The announcement was not found.", + "ServiceError": "The Org Announcements service could not complete the request.", +} +_BACKEND_DIAGNOSTIC_CODES = frozenset( + {*_MODEL_VISIBLE_BACKEND_MESSAGES, "CommittedRefreshFailed"} +) + +# Closed authored-content projection accepted by WeveNova's EssBulletinContent +# DTO. A duplicate never forwards provider-added fields back into Save. +_AUTHORED_BULLETIN_FIELDS = frozenset( + { + "type", + "priority", + "title", + "description", + "primaryAction", + "secondaryAction", + "startDate", + "endDate", + } +) + + +def _resolve_widget_origin( + value: Optional[str] = None, + *, + allow_development: Optional[bool] = None, +) -> str: + origin = ( + value or os.environ.get("VORPAL_WIDGET_ORIGIN") or DEFAULT_WIDGET_ORIGIN + ).rstrip("/") + parsed = urlsplit(origin) + if ( + parsed.scheme != "https" + or not parsed.hostname + or parsed.username is not None + or parsed.password is not None + or parsed.path + or parsed.query + or parsed.fragment + ): + raise ValueError( + "VORPAL_WIDGET_ORIGIN must be an HTTPS origin without credentials, " + "a path, a query, or a fragment." + ) + + normalized_origin = f"https://{parsed.netloc.lower()}" + if normalized_origin in ALLOWED_WIDGET_ORIGINS: + return normalized_origin + + if allow_development is None: + allow_development = ( + os.environ.get(DEVELOPMENT_WIDGET_ORIGIN_ENV, "").strip() == "1" + ) + hostname = parsed.hostname.lower() + is_local_development_origin = ( + hostname in {"localhost", "127.0.0.1", "::1"} + or hostname.endswith(".devtunnels.ms") + ) + if allow_development and is_local_development_origin: + return normalized_origin + + raise ValueError( + "VORPAL_WIDGET_ORIGIN must use an approved Vorpal deployment origin. " + f"Set {DEVELOPMENT_WIDGET_ORIGIN_ENV}=1 for localhost, loopback, or " + "*.devtunnels.ms development origins." + ) + + +WIDGET_ORIGIN = _resolve_widget_origin() + + +def _widget_tool_meta() -> dict[str, Any]: + return { + "ui": { + "resourceUri": ORG_ANNOUNCEMENTS_RESOURCE_URI, + "visibility": ["model", "app"], + } + } + + +def _app_only_tool_meta() -> dict[str, Any]: + """Hide a mutation from the model so only the open widget can call it.""" + return {"ui": {"visibility": ["app"]}} + + +def _shared_tool_meta() -> dict[str, Any]: + return {"ui": {"visibility": ["model", "app"]}} + + +def _widget_resource_meta() -> dict[str, Any]: + return { + "ui": { + "domain": WIDGET_ORIGIN, + "csp": {"resourceDomains": [WIDGET_ORIGIN], "connectDomains": []}, + } + } + + +def _widget_shell() -> str: + script_url = html.escape( + f"{WIDGET_ORIGIN}/mcp-widget/org-announcements/widget.js", quote=True + ) + return ( + "\n" + '\n' + " \n" + ' \n' + ' \n' + " \n" + " \n" + '
\n' + f' \n' + " \n" + "\n" + ) + + +mcp = FastMCP( + "ess-org-announcements", + instructions=( + "Author ESS announcements for one deployed agent in the authenticated " + "tenant. A required titleId selects the agent, not an author permission " + "or an audience group. The current100/latest50 limits apply per tenant " + "and agent. Open the management view or the editor " + "through open_org_announcements; the widget owns every save, lifecycle " + "action, and audience search after it opens." + ), +) + +_client: Optional[OrgAnnouncementsClient] = None +_client_construction_task: Optional[asyncio.Task[None]] = None +_client_users: dict[OrgAnnouncementsClient, int] = {} +_retired_clients: set[OrgAnnouncementsClient] = set() +_graph_client: Optional[GraphDirectoryClient] = None +_graph_client_users: dict[GraphDirectoryClient, int] = {} + +# Guards construction and invalidation of the process-global authoring client so +# concurrent tool calls share one sign-in instead of racing two browser prompts. +_client_lock = asyncio.Lock() + + +def _observe_task_failure(task: asyncio.Task[None]) -> None: + """Retrieve an unobserved construction failure after all waiters cancel.""" + if not task.cancelled(): + task.exception() + + +async def _construct_and_publish_client() -> None: + """Build one authoring client outside the lock and publish or dispose it.""" + global _client, _client_construction_task + current_task = asyncio.current_task() + client = None + published = False + try: + try: + client = await asyncio.to_thread(OrgAnnouncementsClient) + except (LocalCredentialError, LockException, OSError) as error: + raise _FailureResult( + "AuthenticationRequired", + "Organization announcement sign-in could not be completed. Try again.", + source=SOURCE_MCP, + ) from error + + with anyio.CancelScope(shield=True): + close_client = None + async with _client_lock: + if _client is None: + _client = client + published = True + else: + close_client = client + if close_client is not None: + await close_client.aclose() + except asyncio.CancelledError: + if client is not None and not published: + with anyio.CancelScope(shield=True): + await client.aclose() + raise + finally: + with anyio.CancelScope(shield=True): + async with _client_lock: + if _client_construction_task is current_task: + _client_construction_task = None + + +async def get_client() -> OrgAnnouncementsClient: + """Lease the authoring client, constructing it lazily off the event loop. + + ``OrgAnnouncementsClient.__init__`` resolves a delegated token through the + shared core, which is synchronous and — on a cold cache — blocks for as long + as a human takes to finish an interactive browser sign-in. Constructing it + inline would run all of that *inside* the asyncio event loop, freezing every + other in-flight request and the MCP stdio transport itself for the duration. + + One shared construction task is registered behind the lock, but the + potentially human-duration authentication itself runs outside it. Lease + registration remains lock-protected. Callers must release the returned + client, normally through ``get_client_lease``, so a 401 reset can retire the + client without closing it during another request. + """ + global _client_construction_task + while True: + async with _client_lock: + if _client is not None: + client = _client + _client_users[client] = _client_users.get(client, 0) + 1 + return client + construction = _client_construction_task + if construction is None: + construction = asyncio.create_task(_construct_and_publish_client()) + construction.add_done_callback(_observe_task_failure) + _client_construction_task = construction + try: + await asyncio.shield(construction) + except asyncio.CancelledError: + current = asyncio.current_task() + if construction.cancelled() and ( + current is None or not current.cancelling() + ): + continue + raise + + +async def release_client(client: OrgAnnouncementsClient) -> None: + """Release one authoring-client lease and close a retired final user.""" + with anyio.CancelScope(shield=True): + close_client = None + async with _client_lock: + users = _client_users.get(client) + if users is None: + return + if users > 1: + _client_users[client] = users - 1 + return + del _client_users[client] + if client in _retired_clients: + _retired_clients.remove(client) + close_client = client + if close_client is not None: + await close_client.aclose() + + +@asynccontextmanager +async def get_client_lease() -> AsyncIterator[OrgAnnouncementsClient]: + """Keep one authoring client alive for the complete tool operation.""" + client = await get_client() + try: + yield client + finally: + await release_client(client) + + +async def reset_client( + failed_client: Optional[OrgAnnouncementsClient] = None, +) -> bool: + """Drop the authoring client so the next call reauthenticates. + + The shared core resolves its token once, in ``__init__``, and holds it for + the object's lifetime; there is no lazy re-acquisition inside it. Because + this server keeps one process-global client, an expired token would + otherwise fail every subsequent tool call until the MCP host was restarted. + Dropping the object is therefore the reauthentication mechanism: the next + ``get_client()`` builds a fresh one and re-runs silent-first MSAL, which + normally refreshes from the shared cache with no prompt. + + ``failed_client`` makes invalidation compare-and-swap: a delayed 401 from a + retired client cannot evict a healthy replacement. Active users retain a + lease, so the retired client closes only after its final operation exits. + """ + global _client + close_client = None + async with _client_lock: + client = _client + if client is None or ( + failed_client is not None and failed_client is not client + ): + return False + _client = None + if _client_users.get(client): + _retired_clients.add(client) + else: + close_client = client + if close_client is not None: + await close_client.aclose() + return True + + +@asynccontextmanager +async def get_graph_client( + tenant_id: str, object_id: Optional[str] +) -> AsyncIterator[GraphDirectoryClient]: + """Capture one tenant-and-account-bound directory client for an operation. + + Keep only the current tenant's idle client. Replaced clients are closed + after their captured users finish, not while another request is using them. + The account comes from the captured authoring client, never from tool input. + """ + global _graph_client + previous = _graph_client + client = previous + if ( + client is None + or client.tenant_id != tenant_id + or client.object_id != object_id + ): + client = GraphDirectoryClient(tenant_id=tenant_id, object_id=object_id) + _graph_client = client + _graph_client_users[client] = _graph_client_users.get(client, 0) + 1 + try: + if ( + previous is not None + and previous is not client + and not _graph_client_users.get(previous) + ): + with anyio.CancelScope(shield=True): + await previous.aclose() + yield client + finally: + with anyio.CancelScope(shield=True): + remaining = _graph_client_users[client] - 1 + if remaining: + _graph_client_users[client] = remaining + else: + del _graph_client_users[client] + if client is not _graph_client: + await client.aclose() + + +def _now() -> datetime: + return datetime.now(timezone.utc) + + +def _elapsed_ms(started: float) -> int: + return int((time.monotonic() - started) * 1000) + + +def _text_result(payload: dict[str, Any], message: str) -> CallToolResult: + return CallToolResult( + content=[TextContent(type="text", text=message)], + structuredContent=payload, + ) + + +class _FailureResult(Exception): + """Carries a stable discriminated failure back to the tool boundary. + + Every MCP-only failure is one of the documented codes so the widget can + branch on ``code`` and ``retryable`` instead of parsing message text. + ``errors`` carries a *structured backend validation* payload when there is + one, so field-level codes survive intact instead of being flattened into a + single message. + """ + + def __init__( + self, + code: str, + message: str, + *, + retryable: bool = False, + field: Optional[str] = None, + source: str = SOURCE_MCP, + errors: Optional[list[dict[str, Any]]] = None, + ): + super().__init__(message) + self.code = code + self.message = message + self.retryable = retryable + self.field = field + self.source = source + self.errors = errors + + def as_error(self) -> dict[str, Any]: + return { + "code": self.code, + "field": self.field, + "message": self.message, + "retryable": self.retryable, + } + + def as_error_list(self) -> list[dict[str, Any]]: + """Every error to report, preserving a structured backend payload.""" + if self.errors: + return [ + { + "code": entry["code"], + "field": entry["field"], + "message": entry["message"], + "retryable": False, + } + for entry in self.errors + ] + return [self.as_error()] + + +def _diagnostic_code(failure: _FailureResult) -> str: + """Return a content-free code for logs and telemetry.""" + if ( + failure.source == SOURCE_BACKEND + and failure.code not in _BACKEND_DIAGNOSTIC_CODES + ): + return "BackendValidationError" + return normalize_error_code(failure.code) + + +def _validation_failure(error: BulletinValidationError) -> _FailureResult: + """Adapt an HTTP-200 ``EssBulletinSaveResult`` rejection. + + The wrapper's ``errors`` are carried through verbatim. The summary ``code`` + is the *first* backend code so a caller that only reads ``code`` still gets + a real backend code rather than a generic one, and ``field`` is that entry's + field so single-error cases keep their input binding. + """ + first = error.errors[0] + return _FailureResult( + first["code"], + first["message"], + field=first["field"], + retryable=False, + source=SOURCE_BACKEND, + errors=error.errors, + ) + + +def _classify_api_error(error: AgentConfigApiError) -> _FailureResult: + """Map a backend failure onto the documented discriminated results. + + Structured backend validation errors keep their own code and message: the + widget renders localized copy for ``AudienceRequired``, + ``AudienceGroupInvalid``, ``BulletinLimitExceeded``, and the field-level + title/description/action/date/priority/lifecycle codes, so replacing them + with a generic code would lose that mapping. + + The backend's own ``Code`` is inspected before the HTTP status so a + feature-gated tenant is reported as ``FeatureUnavailable`` rather than being + flattened into an authorization or not-found result — but only when the code + is a *real* backend code. ``HttpError`` is the core's placeholder for "the + response carried no code", so it is discarded before any code-based + decision. + """ + if isinstance(error, IndeterminateWriteError): + return _FailureResult("IndeterminateWrite", str(error), retryable=False) + if isinstance(error, BulletinValidationError): + return _validation_failure(error) + + raw_code, _, detail = str(error).partition(": ") + backend_code = "" if raw_code == _FALLBACK_BACKEND_CODE else raw_code + status = error.http_status + + if backend_code in _FEATURE_DISABLED_CODES and backend_code: + return _FailureResult( + "FeatureUnavailable", + "Organization announcements are not enabled for this tenant.", + source=SOURCE_BACKEND, + ) + if status == 401: + return _FailureResult( + "AuthenticationRequired", + "Sign in again to author organization announcements.", + source=SOURCE_BACKEND, + ) + if status == 403: + return _FailureResult( + "AuthorizationDenied", + "This account is not authorized to author organization " + "announcements in this tenant.", + source=SOURCE_BACKEND, + ) + if status == 404: + return _FailureResult( + "NotFound", + "The announcement was not found.", + source=SOURCE_BACKEND, + ) + if status in (400, 409, 422): + # A real backend validation failure. Without a backend code there is + # nothing to map, so report a generic invalid request rather than + # inventing one; the detail is the service's own message. + return _FailureResult( + backend_code or "InvalidRequest", + detail or str(error), + source=SOURCE_BACKEND, + ) + if status in (429, 502, 503, 504): + return _FailureResult( + "NetworkError", + "The Org Announcements service is temporarily unavailable.", + retryable=True, + source=SOURCE_BACKEND, + ) + if status in (405, 501): + return _FailureResult( + "FeatureUnavailable", + "Organization announcements are not enabled for this tenant.", + source=SOURCE_BACKEND, + ) + # Everything left is a server-side failure (500) or a response this client + # could not interpret. It is emphatically NOT a validation error the maker + # can fix, and the raw text can echo internal detail, so it becomes one + # stable, privacy-safe code. Non-retryable: the request reached the service + # and was processed, so an automatic replay would risk a duplicate write + # without any expectation of a different answer. + return _FailureResult( + "ServiceError", + "The Org Announcements service could not complete the request.", + retryable=False, + source=SOURCE_BACKEND, + ) + + +def _classify_graph_error(error: GraphDirectoryError) -> _FailureResult: + return _FailureResult( + error.code, + str(error), + retryable=error.retryable, + source=SOURCE_GRAPH, + ) + + +class _CommittedRefreshError(Exception): + """The write committed; only the follow-up refresh failed. + + Raised *after* the authoring API has definitely acknowledged a save, so the + announcement exists regardless of what happens next. It is deliberately a + distinct type from :class:`_FailureResult` because the two demand opposite + handling: a normal failure means "nothing was written, you may retry", and + this means "it was written, do not retry". + """ + + def __init__(self, cause: _FailureResult): + super().__init__(cause.message) + self.cause = cause + + +def _committed_refresh_failure(cause: _FailureResult) -> _FailureResult: + """Build the explicit partial-success contract for a committed write. + + Retrying here is not merely wasteful, it is unsafe: the create path is + unkeyed, so a replay would create a *second* announcement. The result is + therefore a stable, non-retryable ``CommittedRefreshFailed`` that states + plainly that the write succeeded and only the refresh did not, and the + original refresh failure is preserved as a second entry so the widget can + still tell a Graph outage from a service outage. + """ + summary_message = ( + "The change was saved, but the updated list could not be loaded. " + "Refresh to see the current state — do not repeat the action." + ) + return _FailureResult( + "CommittedRefreshFailed", + summary_message, + retryable=False, + source=cause.source, + errors=[ + { + "code": "CommittedRefreshFailed", + "field": None, + "message": summary_message, + }, + { + "code": cause.code, + "field": cause.field, + "message": cause.message, + }, + ], + ) + + +async def _resolve_audience_metadata( + graph_client: GraphDirectoryClient, + audience_ids: list[str], +) -> list[AudienceGroup]: + """Resolve one canonical audience list into ordered display metadata.""" + if not audience_ids: + return [] + try: + resolved = await graph_client.resolve_groups(audience_ids) + except GraphDirectoryError as error: + raise _FailureResult( + "AudienceMetadataUnavailable", + f"Audience group names could not be loaded. {error}", + retryable=error.retryable, + source=SOURCE_GRAPH, + ) from error + return [ + AudienceGroup.model_validate(item) + for item in build_audience_metadata(audience_ids, resolved) + ] + + +async def _resolve_manager_metadata( + graph_client: GraphDirectoryClient, + items: list[dict[str, Any]], +) -> dict[str, list[dict[str, Any]]]: + """Resolve audience metadata for every listed bulletin in one Graph batch. + + Soft-deleted rows are skipped: the manager drops them, so resolving their + audience would issue Graph lookups for groups nobody will see. + + Deduplication applies only to the lookup batch. Each bulletin's canonical + audience list is re-expanded in its own order, so per-item multiplicity and + ordering stay exact. + """ + visible = [config for config in items if not is_deleted_item(config)] + + all_ids: list[str] = [] + for config in visible: + audience = config.get("audience") + if isinstance(audience, list): + all_ids.extend( + group_id for group_id in audience if isinstance(group_id, str) + ) + + if not all_ids: + return {} + + try: + resolved = await graph_client.resolve_groups(all_ids) + except GraphDirectoryError as error: + raise _FailureResult( + "AudienceMetadataUnavailable", + f"Audience group names could not be loaded. {error}", + retryable=error.retryable, + source=SOURCE_GRAPH, + ) from error + + metadata: dict[str, list[dict[str, Any]]] = {} + for config in visible: + bulletin_id = config.get("id") + audience = config.get("audience") + metadata[bulletin_id or ""] = build_audience_metadata( + [group_id for group_id in audience if isinstance(group_id, str)] + if isinstance(audience, list) + else [], + resolved, + ) + return metadata + + +async def _manager_state( + client: OrgAnnouncementsClient, scope: dict[str, str], + graph_client: GraphDirectoryClient, +) -> dict[str, Any]: + items = await client.list_bulletins(scope["titleId"]) + return build_manager_state( + items, + await _resolve_manager_metadata(graph_client, items), + _now(), + tenant_id=scope["tenantId"], + title_id=scope["titleId"], + ) + + +def _editor_payload( + mode: EditorMode, + config: Optional[dict[str, Any]], + draft: AnnouncementEditorDraft, + scope: dict[str, str], +) -> dict[str, Any]: + return { + **scope, + "view": "editor", + "mode": mode, + "config": config, + "draft": draft.model_dump( + mode="json", exclude={"id"} if draft.id is None else set() + ), + } + + +def _model_visible_failure(failure: _FailureResult) -> _FailureResult: + """Replace backend-controlled detail before exposing a failure to the model.""" + if failure.source == SOURCE_BACKEND: + safe_message = _MODEL_VISIBLE_BACKEND_MESSAGES.get(failure.code) + if safe_message is None: + return _FailureResult( + "InvalidRequest", + "The Org Announcements service rejected the request.", + source=SOURCE_BACKEND, + ) + return _FailureResult( + failure.code, + safe_message, + retryable=failure.retryable, + source=SOURCE_BACKEND, + ) + return failure + + +def _open_error_payload( + request: dict[str, Any], failure: _FailureResult, scope: dict[str, str] +) -> CallToolResult: + """Return a recoverable open error instead of a half-rendered editor. + + The request is response-only retry state, including the original suggestion. + It is never logged or sent to telemetry. + """ + visible_failure = _model_visible_failure(failure) + + payload = { + **scope, + "view": "error", + "request": request, + "code": visible_failure.code, + "message": visible_failure.message, + "retryable": visible_failure.retryable, + } + return CallToolResult( + content=[TextContent(type="text", text=visible_failure.message)], + structuredContent=payload, + isError=True, + ) + + +def _open_success( + payload: dict[str, Any], message: str, started: float +) -> CallToolResult: + record_operation( + "open_org_announcements", + outcome="success", + latency_ms=_elapsed_ms(started), + ) + return _text_result(payload, message) + + +def _open_failure( + request: dict[str, Any], + failure: _FailureResult, + started: float, + scope: dict[str, str], +) -> CallToolResult: + diagnostic_code = _diagnostic_code(failure) + _LOGGER.warning("open_org_announcements failed: %s", diagnostic_code) + record_operation( + "open_org_announcements", + outcome="failure", + latency_ms=_elapsed_ms(started), + error_code=diagnostic_code, + error_source=failure.source, + ) + return _open_error_payload(request, failure, scope) + + +async def _discovery_tool_error( + operation: str, + error: Exception, + started: float, + client: Optional[OrgAnnouncementsClient], +) -> ToolError: + failure = _model_visible_failure(await _failure_from(error, client)) + diagnostic_code = _diagnostic_code(failure) + _LOGGER.warning("%s failed: %s", operation, diagnostic_code) + record_operation( + operation, + outcome="failure", + latency_ms=_elapsed_ms(started), + error_code=diagnostic_code, + error_source=failure.source, + ) + return ToolError(failure.message) + + +@mcp.resource( + ORG_ANNOUNCEMENTS_RESOURCE_URI, + name="Org announcements", + title="Org announcements", + description="Announcement manager and editor for the selected ESS agent.", + mime_type=WIDGET_MIME_TYPE, + meta=_widget_resource_meta(), +) +def org_announcements_widget() -> str: + return _widget_shell() + + +@mcp.tool( + annotations=_READ_ONLY_ANNOTATIONS, +) +async def list_agent_configs() -> str: + """List configured deployed ESS agents; never initialize configuration.""" + started = time.monotonic() + client = None + try: + async with get_client_lease() as client: + result = await client.list_agent_configs() + except (_FailureResult, AgentConfigApiError, httpx.RequestError) as error: + raise await _discovery_tool_error( + "list_agent_configs", error, started, client + ) from None + record_operation( + "list_agent_configs", + outcome="success", + latency_ms=_elapsed_ms(started), + ) + return json.dumps(result, indent=2) + + +@mcp.tool( + annotations=_READ_ONLY_ANNOTATIONS, +) +async def search_agents(searchString: str) -> str: + """Find deployed ESS agents by name to resolve their titleId.""" + started = time.monotonic() + client = None + try: + async with get_client_lease() as client: + result = await client.search_agents(searchString) + except (_FailureResult, AgentConfigApiError, httpx.RequestError) as error: + raise await _discovery_tool_error( + "search_agents", error, started, client + ) from None + record_operation( + "search_agents", + outcome="success", + latency_ms=_elapsed_ms(started), + ) + return json.dumps(result, indent=2) + + +@mcp.tool( + meta=_widget_tool_meta(), + annotations=_READ_ONLY_ANNOTATIONS, +) +async def open_org_announcements( + titleId: str, + view: Literal["manager", "editor"], + mode: Optional[Literal["create", "edit"]] = None, + bulletinId: Optional[str] = None, + suggestedDraft: Optional[SuggestedBulletinDraft] = None, +) -> CallToolResult: + """Open announcements for the selected agent in the manager or editor. + + This tool only reads: it never creates, updates, or publishes. The model + calls it at most once per maker turn; the widget uses it for navigation + within the same tenant-and-agent scope. Manager takes no editor arguments. + Editor requires mode=create without bulletinId, or mode=edit with bulletinId. + suggestedDraft is supported only for create, including editable copies + whose actions need repair. Opening preserves working content; it does not + establish that the content is valid to save or publish. + """ + # The flat tool signature stays compatible with MCP callers. Validate the + # combination before entering the recoverable widget-error path: an invalid + # request is not safe retry state and cannot be rendered as an error view. + try: + validate_title_id(titleId) + if bulletinId is not None: + validate_bulletin_id(bulletinId) + validated = OpenAnnouncementsRequest( + titleId=titleId, view=view, mode=mode, + bulletinId=bulletinId, suggestedDraft=suggestedDraft, + ) + except ValidationError as error: + messages = "; ".join( + entry["msg"] for entry in error.errors(include_input=False) + ) + raise ToolError(f"Invalid announcement opener: {messages}") from None + except ValueError as error: + raise ToolError(f"Invalid announcement opener: {error}") from None + + started = time.monotonic() + scope = {"titleId": titleId} + client = None + request = validated.model_dump(exclude_none=True, exclude={"suggestedDraft"}) + if suggestedDraft is not None: + request["suggestedDraft"] = suggestedDraft.retry_payload() + + try: + # Capture once even for an empty create. Subsequent reads and post-save + # refreshes must never adopt a different global client's tenant. + async with get_client_lease() as client: + scope = {"tenantId": client.tenant_id, "titleId": titleId} + + if view == "manager": + async with get_graph_client( + scope["tenantId"], client.object_id + ) as graph_client: + state = await _manager_state(client, scope, graph_client) + return _open_success( + {"view": "manager", **state}, + "Opened announcements for the selected ESS agent.", + started, + ) + + if mode == "create": + async with get_graph_client( + scope["tenantId"], client.object_id + ) as graph_client: + audience_metadata = await _resolve_audience_metadata( + graph_client, + list(suggestedDraft.audience) + if suggestedDraft is not None and suggestedDraft.audience + else [], + ) + draft = build_create_draft(suggestedDraft, audience_metadata) + return _open_success( + _editor_payload("create", None, draft, scope), + "Opened a new organization announcement for review. Nothing is " + "saved until you publish or save a draft in the editor.", + started, + ) + + config = await client.get_bulletin(titleId, bulletinId) + audience = config.get("audience") + async with get_graph_client( + scope["tenantId"], client.object_id + ) as graph_client: + audience_metadata = await _resolve_audience_metadata( + graph_client, + [ + group_id + for group_id in audience + if isinstance(group_id, str) + ] + if isinstance(audience, list) + else [], + ) + draft = build_editor_draft_from_config(config, audience_metadata) + message = "Opened the announcement for editing." + return _open_success( + _editor_payload("edit", config, draft, scope), message, started + ) + + except _FailureResult as failure: + return _open_failure(request, failure, started, scope) + except AgentConfigApiError as error: + return _open_failure( + request, + await _failure_from(error, client), + started, + scope, + ) + except httpx.RequestError: + return _open_failure( + request, + _FailureResult( + "NetworkError", + "The Org Announcements service could not be reached.", + retryable=True, + source=SOURCE_BACKEND, + ), + started, + scope, + ) + except (ValidationError, ValueError) as error: + return _open_failure( + request, _FailureResult("InvalidRequest", str(error)), started, scope + ) + + +def _fail( + operation: str, failure: _FailureResult, started: float, scope: dict[str, str] +) -> CallToolResult: + """Emit the content-free failure event and build the tool result.""" + diagnostic_code = _diagnostic_code(failure) + record_operation( + operation, + outcome="failure", + latency_ms=_elapsed_ms(started), + error_code=diagnostic_code, + error_source=failure.source, + ) + return _mutation_failure(failure, scope) + + +def _mutation_failure( + failure: _FailureResult, scope: dict[str, str] +) -> CallToolResult: + return CallToolResult( + content=[TextContent(type="text", text=failure.message)], + structuredContent={ + **scope, "status": "failure", "errors": failure.as_error_list() + }, + isError=True, + ) + + +async def _failure_from( + error: Exception, + authoring_client: Optional[OrgAnnouncementsClient], +) -> _FailureResult: + """Classify a failure and reauthenticate the authoring client on a 401. + + A 401 from the authoring API means the cached delegated token is expired or + revoked. The shared core holds its token for the client object's lifetime, + so the stale credential is discarded by discarding the client; the next tool + call then rebuilds it and re-runs silent-first MSAL. + + Only an *authoring-side* 401 resets the authoring client. The Graph client + reports the same ``AuthenticationRequired`` code for its own expiry and + handles that itself, so the source bucket is what distinguishes them — + resetting on a Graph 401 would throw away a perfectly good WeveNova token. + """ + failure = _to_failure(error) + if failure.code == "AuthenticationRequired" and failure.source == SOURCE_BACKEND: + await reset_client(authoring_client) + return failure + + +def _to_failure(error: Exception) -> _FailureResult: + if isinstance(error, _FailureResult): + return error + if isinstance(error, AgentConfigApiError): + return _classify_api_error(error) + if isinstance(error, GraphDirectoryError): + return _classify_graph_error(error) + if isinstance(error, httpx.RequestError): + return _FailureResult( + "NetworkError", + "The Org Announcements service could not be reached.", + retryable=True, + source=SOURCE_BACKEND, + ) + return _FailureResult("InvalidRequest", str(error)) + + +# Everything a mutation can raise. Listed once so the write step and the +# post-commit refresh step cannot drift apart and let an exception escape one +# but not the other. +_MUTATION_ERRORS = ( + _FailureResult, + AgentConfigApiError, + GraphDirectoryError, + httpx.RequestError, + ValueError, +) + + +async def _saved_item_result( + client: OrgAnnouncementsClient, + scope: dict[str, str], + config: dict[str, Any], + message: str, +) -> CallToolResult: + """Return the canonical saved item plus refreshed manager state. + + Called only after the write has been acknowledged, so any error raised here + is a *refresh* failure over a committed record. It is wrapped in + :class:`_CommittedRefreshError` so the caller cannot mistake it for a failed + write and offer a retry that would duplicate the announcement. + """ + try: + audience = config.get("audience") + async with get_graph_client(scope["tenantId"], client.object_id) as graph_client: + audience_metadata = await _resolve_audience_metadata( + graph_client, + [group_id for group_id in audience if isinstance(group_id, str)] + if isinstance(audience, list) + else [], + ) + manager = await _manager_state(client, scope, graph_client) + except _MUTATION_ERRORS as error: + raise _CommittedRefreshError( + await _failure_from(error, client) + ) from error + + return _text_result( + { + **scope, + "status": "success", + "item": { + "config": config, + "audienceMetadata": [ + group.model_dump(mode="json") for group in audience_metadata + ], + }, + "manager": manager, + }, + message, + ) + + +@mcp.tool( + meta=_app_only_tool_meta(), + annotations=_MUTATION_ANNOTATIONS, +) +async def save_bulletin( + titleId: str, + bulletin: dict[str, Any], + audience: list[str], + status: Literal["draft", "published"], + id: Optional[str] = None, # noqa: A002 — the widget's wire field name +) -> CallToolResult: + """Create or update an announcement with its complete authored state.""" + started = time.monotonic() + scope = {"titleId": titleId} + try: + validate_title_id(titleId) + if id is not None: + validate_bulletin_id(id) + request = SaveBulletinRequest.model_validate( + { + "id": id, + "bulletin": bulletin, + "audience": audience, + "status": status, + } + ) + except ValueError as error: + return _fail( + "save_bulletin", + _FailureResult("InvalidRequest", str(error)), + started, + scope, + ) + + payload = request.model_dump(mode="json", exclude_none=True) + client = None + try: + async with get_client_lease() as client: + scope = {"tenantId": client.tenant_id, "titleId": titleId} + try: + saved = await client.save_bulletin(titleId, payload) + except CommittedCanonicalReloadError as error: + failure = await _failure_from(error.cause, client) + _LOGGER.warning( + "save_bulletin canonical reload failed after commit: %s", + _diagnostic_code(failure), + ) + return _fail( + "save_bulletin", + _committed_refresh_failure(failure), + started, + scope, + ) + + try: + result = await _saved_item_result( + client, + scope, + saved, + "Published the organization announcement." + if status == "published" + else "Saved the organization announcement draft.", + ) + except _CommittedRefreshError as error: + # The write is committed. Report the explicit partial success so the + # widget tells the maker to refresh instead of offering a retry that + # would create a second announcement. + _LOGGER.warning( + "save_bulletin refresh failed after commit: %s", + _diagnostic_code(error.cause), + ) + return _fail( + "save_bulletin", + _committed_refresh_failure(error.cause), + started, + scope, + ) + + record_operation( + "save_bulletin", + outcome="success", + latency_ms=_elapsed_ms(started), + ) + return result + except _MUTATION_ERRORS as error: + failure = await _failure_from(error, client) + _LOGGER.warning( + "save_bulletin failed: %s (create=%s)", + _diagnostic_code(failure), + id is None, + ) + return _fail("save_bulletin", failure, started, scope) + + +@mcp.tool( + meta=_app_only_tool_meta(), + annotations=_DESTRUCTIVE_ANNOTATIONS, +) +async def transition_bulletin( + titleId: str, + id: str, # noqa: A002 — the widget's wire field name + transition: TransitionName, +) -> CallToolResult: + """Archive, unarchive, move to draft, or delete an announcement. + + The request carries only the identifier and the new status. The service + loads the canonical record, preserves its authored content and audience, and + validates the lifecycle change, so no client-side merge is performed. + + ``TransitionName`` lists exactly the four supported operations. Publish now + is deliberately absent and is rejected at argument validation: it must go + through ``save_bulletin`` with the row's complete canonical state so the new + start instant cannot race a concurrent edit. + + KNOWN LIMITATION (Vorpal compatibility). A legacy client sending + ``transition: "publishNow"`` is rejected by FastMCP's *schema* validation, + before this function runs, so it surfaces as a protocol ``ToolError`` rather + than this server's structured ``InvalidRequest`` result. Converting it would + require widening ``transition`` to a free string, which would advertise + ``publishNow`` as acceptable in the production schema and re-open the race + that removing it closed — so the rollout prerequisite stands: ship a Vorpal + build that routes publish-now through ``save_bulletin``. The current + boundary is pinned by tests so it stays a known contract. + """ + started = time.monotonic() + scope = {"titleId": titleId} + client = None + try: + validate_title_id(titleId) + validate_bulletin_id(id) + async with get_client_lease() as client: + scope = {"tenantId": client.tenant_id, "titleId": titleId} + try: + changed = await client.transition_bulletin( + titleId, id, TRANSITION_STATUS[transition] + ) + except CommittedCanonicalReloadError as error: + cause = await _failure_from(error.cause, client) + _LOGGER.warning( + "transition_bulletin canonical reload failed after commit: %s (%s)", + _diagnostic_code(cause), + transition, + ) + return _fail( + "transition_bulletin", + _committed_refresh_failure(cause), + started, + scope, + ) + + # The transition is committed from here on. A refresh failure must never be + # reported as a retryable normal failure, because the lifecycle change has + # already been applied and re-issuing it could fail validation or move the + # record again. + try: + async with get_graph_client( + scope["tenantId"], client.object_id + ) as graph_client: + manager = await _manager_state(client, scope, graph_client) + except _MUTATION_ERRORS as error: + cause = await _failure_from(error, client) + _LOGGER.warning( + "transition_bulletin refresh failed after commit: %s (%s)", + _diagnostic_code(cause), + transition, + ) + return _fail( + "transition_bulletin", + _committed_refresh_failure(cause), + started, + scope, + ) + + record_operation( + "transition_bulletin", + outcome="success", + latency_ms=_elapsed_ms(started), + ) + # Non-delete transitions include the canonical changed row alongside + # the manager state. Delete returns only the refreshed manager because + # the tombstoned resource is no longer readable. + payload = {**scope, "status": "success", "manager": manager} + if changed is not None: + payload["item"] = {"config": changed} + return _text_result( + payload, + "Updated the organization announcement.", + ) + except _MUTATION_ERRORS as error: + failure = await _failure_from(error, client) + _LOGGER.warning( + "transition_bulletin failed: %s (%s)", + _diagnostic_code(failure), + transition, + ) + return _fail("transition_bulletin", failure, started, scope) + + +@mcp.tool( + meta=_app_only_tool_meta(), + annotations=_MUTATION_ANNOTATIONS, +) +async def duplicate_bulletin( + titleId: str, + id: str, # noqa: A002 — the widget's wire field name +) -> CallToolResult: + """Copy an existing announcement into a new Draft. + + Loads the canonical source, projects only WeveNova-authored content fields, + and creates a new Draft. A missing source is a not-found failure, never a + create: duplicating something that no longer exists must not invent a record. + """ + started = time.monotonic() + scope = {"titleId": titleId} + client = None + try: + validate_title_id(titleId) + validate_bulletin_id(id) + async with get_client_lease() as client: + scope = {"tenantId": client.tenant_id, "titleId": titleId} + source = await client.get_bulletin(titleId, id) + + # Stored content does not pass through BulletinInput. Project the + # provider's closed DTO and strip blank schedule sentinels before Save. + bulletin = without_blank_schedule( + { + key: value + for key, value in source["bulletin"].items() + if key in _AUTHORED_BULLETIN_FIELDS + } + ) + audience = source.get("audience") + payload = { + "bulletin": bulletin, + "audience": [ + group_id + for group_id in ( + audience if isinstance(audience, list) else [] + ) + if isinstance(group_id, str) + ], + "status": "draft", + } + + try: + created = await client.save_bulletin(titleId, payload) + except CommittedCanonicalReloadError as error: + failure = await _failure_from(error.cause, client) + _LOGGER.warning( + "duplicate_bulletin canonical reload failed after commit: %s", + _diagnostic_code(failure), + ) + return _fail( + "duplicate_bulletin", + _committed_refresh_failure(failure), + started, + scope, + ) + except _MUTATION_ERRORS as error: + failure = await _failure_from(error, client) + _LOGGER.warning( + "duplicate_bulletin failed: %s", + _diagnostic_code(failure), + ) + return _fail("duplicate_bulletin", failure, started, scope) + + try: + result = await _saved_item_result( + client, + scope, + created, + "Created a draft copy of the organization announcement.", + ) + except _CommittedRefreshError as error: + # The copy exists. This is an unkeyed create, so a retry would produce a + # second copy; report the committed partial success instead. + _LOGGER.warning( + "duplicate_bulletin refresh failed after commit: %s", + _diagnostic_code(error.cause), + ) + return _fail( + "duplicate_bulletin", + _committed_refresh_failure(error.cause), + started, + scope, + ) + + record_operation( + "duplicate_bulletin", + outcome="success", + latency_ms=_elapsed_ms(started), + ) + return result + except _MUTATION_ERRORS as error: + failure = await _failure_from(error, client) + _LOGGER.warning( + "duplicate_bulletin source load failed: %s", + _diagnostic_code(failure), + ) + return _fail("duplicate_bulletin", failure, started, scope) + + +@mcp.tool( + meta=_shared_tool_meta(), + annotations=_READ_ONLY_ANNOTATIONS, +) +async def search_audience_groups(query: str) -> CallToolResult: + """Search eligible audience groups by display name. + + Returns at most 20 security groups, mail-enabled security groups, or classic + distribution groups, deduplicated by ID and in a deterministic order. + Microsoft 365 groups and dynamic-membership groups are not offered. + + ``exhausted`` and ``pagesExamined`` are reported so the caller can tell "no + more matches exist" from "the page budget ran out". Without them a capped + search looks identical to an exhausted one and the maker would be told a + group does not exist when it simply was not reached. + """ + started = time.monotonic() + authoring_client = None + try: + escape_search_value(query) + async with get_client_lease() as authoring_client: + async with get_graph_client( + authoring_client.tenant_id, authoring_client.object_id + ) as graph_client: + result = await graph_client.search_groups(query) + except (_FailureResult, GraphDirectoryError, AgentConfigApiError, httpx.RequestError) as error: + failure = await _failure_from(error, authoring_client) + diagnostic_code = _diagnostic_code(failure) + _LOGGER.warning("search_audience_groups failed: %s", diagnostic_code) + record_operation( + "search_audience_groups", + outcome="failure", + latency_ms=_elapsed_ms(started), + error_code=diagnostic_code, + error_source=failure.source, + ) + return CallToolResult( + content=[TextContent(type="text", text=failure.message)], + structuredContent={ + "status": "failure", + "code": failure.code, + "retryable": failure.retryable, + }, + isError=True, + ) + except ValueError as error: + _LOGGER.warning("search_audience_groups rejected an invalid query") + record_operation( + "search_audience_groups", + outcome="failure", + latency_ms=_elapsed_ms(started), + error_code="InvalidRequest", + error_source=SOURCE_MCP, + ) + return CallToolResult( + content=[TextContent(type="text", text=str(error))], + structuredContent={ + "status": "failure", + "code": "InvalidRequest", + "retryable": False, + }, + isError=True, + ) + + groups = result["groups"] + record_operation( + "search_audience_groups", + outcome="success", + latency_ms=_elapsed_ms(started), + ) + return _text_result( + { + "status": "success", + "groups": groups, + "exhausted": result["exhausted"], + "pagesExamined": result["pagesExamined"], + }, + f"Found {len(groups)} eligible audience group(s).", + ) + + +if __name__ == "__main__": + mcp.run() diff --git a/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/telemetry.py b/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/telemetry.py new file mode 100644 index 000000000..65d12b210 --- /dev/null +++ b/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/telemetry.py @@ -0,0 +1,178 @@ +# Copyright (c) Microsoft Corporation. +# Licensed under the MIT License. + +"""Content-free telemetry for the Org Announcements MCP server. + +This is a thin, deliberately narrow adapter over the kit's existing +``scripts/adk_telemetry.py`` conventions. It emits one ``adk.api.call`` event +per tool invocation carrying exactly four things: + +* ``operation`` — the MCP tool name, from a fixed allowlist; +* ``outcome`` — ``success`` or ``failure``; +* ``latency_ms`` — wall-clock duration of the operation; +* ``error_code`` — one of this server's stable discriminated codes; +* ``error_category`` — a broad source bucket (``backend``/``graph``/``mcp``). + +Nothing else. In particular this module never emits, and has no parameter that +could carry, announcement content (title, description, action label, URL, or +prompt), audience group identifiers or names, the tenant's API endpoint, tokens, +claims, or the opener request payload. ``error_message`` is deliberately never +populated: backend messages can echo authored content back, and the stable code +is what a dashboard or a support engineer actually needs. + +Every path fails open. Telemetry must never turn a working save into a failed +tool call, so import errors, missing config, and emit errors are all swallowed. +""" + +from __future__ import annotations + +import logging +import os +import sys +from typing import Any, Optional + + +_LOGGER = logging.getLogger("ess-org-announcements.telemetry") + +# The MCP tool names this server exposes. An operation outside the allowlist is +# bucketed rather than emitted verbatim, so a future tool cannot silently mint a +# new dimension value (and cannot smuggle caller-controlled text into Aria). +_OPERATIONS = frozenset( + { + "list_agent_configs", + "open_org_announcements", + "save_bulletin", + "search_agents", + "transition_bulletin", + "duplicate_bulletin", + "search_audience_groups", + } +) +OPERATION_UNKNOWN = "unknown" + +# Broad source buckets. Anything narrower would start describing the tenant's +# configuration. +SOURCE_BACKEND = "backend" +SOURCE_GRAPH = "graph" +SOURCE_MCP = "mcp" +_SOURCES = frozenset({SOURCE_BACKEND, SOURCE_GRAPH, SOURCE_MCP}) + +# Stable codes are short identifiers, never free text. This bound is a +# belt-and-braces guard so a malformed code can never carry a payload. +_MAX_CODE_LENGTH = 64 +_ERROR_CODES = frozenset( + { + "AudienceMetadataUnavailable", + "AuthenticationRequired", + "AuthorizationDenied", + "BackendValidationError", + "CommittedRefreshFailed", + "FeatureUnavailable", + "IndeterminateWrite", + "InvalidRequest", + "NetworkError", + "NotFound", + "SearchUnavailable", + "ServiceError", + } +) + +# Emitting is opt-out through the same switch the rest of the ADK honours; the +# module-level import is resolved lazily so a server started outside the kit +# layout still runs. +_ADK_TELEMETRY: Optional[Any] = None +_ADK_TELEMETRY_RESOLVED = False + +_SCRIPTS_DIR = os.path.abspath( + os.path.join( + os.path.dirname(os.path.abspath(__file__)), "..", "..", "..", "scripts" + ) +) + + +def normalize_operation(operation: str) -> str: + """Clamp an operation name to the allowlist.""" + if not isinstance(operation, str): + return OPERATION_UNKNOWN + return operation if operation in _OPERATIONS else OPERATION_UNKNOWN + + +def normalize_source(source: str) -> str: + """Clamp an error source to the broad bucket allowlist.""" + if not isinstance(source, str): + return SOURCE_MCP + return source if source in _SOURCES else SOURCE_MCP + + +def normalize_error_code(error_code: str) -> str: + """Clamp an error code to the fixed diagnostic allowlist. + + Backend values can be caller-controlled even when they look like short + identifiers, so shape validation alone is not a privacy boundary. + """ + if not isinstance(error_code, str): + return "" + candidate = error_code.strip() + if not candidate: + return "" + if ( + len(candidate) > _MAX_CODE_LENGTH + or not candidate.isascii() + or not candidate.replace("_", "").isalnum() + or candidate not in _ERROR_CODES + ): + return "UnknownError" + return candidate + + +def _adk_telemetry() -> Optional[Any]: + """Resolve ``scripts/adk_telemetry`` once, tolerating its absence.""" + global _ADK_TELEMETRY, _ADK_TELEMETRY_RESOLVED + if _ADK_TELEMETRY_RESOLVED: + return _ADK_TELEMETRY + _ADK_TELEMETRY_RESOLVED = True + try: + if _SCRIPTS_DIR not in sys.path: + sys.path.append(_SCRIPTS_DIR) + import adk_telemetry # noqa: PLC0415 — resolved lazily and optionally + + _ADK_TELEMETRY = adk_telemetry + except Exception: # noqa: BLE001 — telemetry must never break a tool call + _ADK_TELEMETRY = None + return _ADK_TELEMETRY + + +def record_operation( + operation: str, + *, + outcome: str, + latency_ms: int, + error_code: str = "", + error_source: str = SOURCE_MCP, +) -> None: + """Emit one content-free ``adk.api.call`` event for a tool invocation. + + ``api_endpoint`` carries the *tool name*, not a URL: the tenant's API host + is environment-specific and is never reported. Failure adds only the stable + code and the broad source bucket; ``error_message`` is left empty on purpose. + """ + telemetry = _adk_telemetry() + if telemetry is None: + return + + normalized_outcome = "success" if outcome == "success" else "failure" + fields: dict[str, Any] = { + "api_endpoint": normalize_operation(operation), + "outcome": normalized_outcome, + "latency_ms": max(0, int(latency_ms)), + } + if normalized_outcome == "failure": + fields["error_code"] = normalize_error_code(error_code) + fields["error_category"] = normalize_source(error_source) + # Never a message: backend text can echo the announcement back. + fields["error_message"] = "" + + try: + telemetry.emit_api_call(**fields) + except Exception: # noqa: BLE001 — fail open, never break the tool call + _LOGGER.debug("Org Announcements telemetry emit failed", exc_info=False) diff --git a/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/validation.py b/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/validation.py new file mode 100644 index 000000000..eef49bf36 --- /dev/null +++ b/solutions/ess-maker-skills/src/mcp/agentconfig_org_announcements/validation.py @@ -0,0 +1,58 @@ +# Copyright (c) Microsoft Corporation. +# Licensed under the MIT License. + +"""Request-boundary validation shared by the announcements client and server.""" + +from __future__ import annotations + +import os +import sys +from uuid import UUID + +# The AgentConfiguration MCP family lives at the ``src/mcp`` root as sibling +# folders sharing the neutral ``agentconfig_core`` client core. There is no +# package __init__.py, and each server launches with cwd set to its own folder +# on a flat sys.path, so make the sibling ``agentconfig_core`` folder importable. +sys.path.insert( + 0, + os.path.join(os.path.dirname(os.path.abspath(__file__)), "..", "agentconfig_core"), +) + +from _odata import _validate_title_id # noqa: E402 + + +_MAX_BULLETIN_ID_LENGTH = 256 + + +def validate_title_id(title_id: str) -> str: + """Validate the opaque Employee Agent route key.""" + return _validate_title_id(title_id) + + +def validate_bulletin_id(bulletin_id: str) -> str: + """Validate and canonicalize the backend-assigned OData Guid key.""" + if not isinstance(bulletin_id, str) or not bulletin_id: + raise ValueError("bulletinId must be a non-empty GUID string") + if bulletin_id != bulletin_id.strip(): + raise ValueError("bulletinId must not have surrounding whitespace") + if bulletin_id in {".", ".."}: + raise ValueError("bulletinId must not be a URL dot-segment") + if len(bulletin_id) > _MAX_BULLETIN_ID_LENGTH: + raise ValueError( + f"bulletinId must not exceed {_MAX_BULLETIN_ID_LENGTH} characters" + ) + if any( + ord(character) < 0x20 or ord(character) == 0x7F for character in bulletin_id + ): + raise ValueError("bulletinId must not contain control characters") + if any(separator in bulletin_id for separator in ("/", "\\", "?", "#", "%")): + raise ValueError( + "bulletinId must not contain path, query, fragment, or escape separators" + ) + try: + parsed = UUID(bulletin_id) + except (ValueError, AttributeError) as error: + raise ValueError("bulletinId must be a valid GUID") from error + if parsed.int == 0: + raise ValueError("bulletinId must not be the empty GUID") + return str(parsed) diff --git a/solutions/ess-maker-skills/src/skills/foundation-setup/SKILL.md b/solutions/ess-maker-skills/src/skills/foundation-setup/SKILL.md index 1f857a483..d3959ec22 100644 --- a/solutions/ess-maker-skills/src/skills/foundation-setup/SKILL.md +++ b/solutions/ess-maker-skills/src/skills/foundation-setup/SKILL.md @@ -241,6 +241,7 @@ another setup completion summary. After the maker selects it, show: The {agent display name} agent is now active. - Run `/landing-page` to configure branding and the content employees see. +- Run `/org-announcements` to create or manage announcements for this agent. - Run `/connect` to choose an integration. - Type `/menu` to see all available capabilities. diff --git a/solutions/ess-maker-skills/src/skills/org-announcements/SKILL.md b/solutions/ess-maker-skills/src/skills/org-announcements/SKILL.md new file mode 100644 index 000000000..93d0ee701 --- /dev/null +++ b/solutions/ess-maker-skills/src/skills/org-announcements/SKILL.md @@ -0,0 +1,504 @@ +--- +name: org-announcements +description: >- + Create, edit, republish, and manage announcements for a deployed ESS agent + through the Org Announcements MCP server. Use for announcements, org + announcements, bulletins, alerts, announcement audiences, archiving or + republishing an announcement, and any call to the ess-org-announcements + MCP server. +--- + +# Org Announcements + +Orchestrate organization announcements through the Org Announcements MCP +server. Open the right view once, let the widget own the editing session, and +never claim a save the tools did not return. + +## Tenant-and-agent scope + +Org Announcements belong to **one deployed ESS agent in the authenticated +tenant**, selected by its required `titleId`. They are not shared across every +agent. `titleId` is an agent identifier, not an author permission, a Dataverse +`botId`, or a Graph audience group. Tenant identity comes only from the +authoring sign-in; never supply `tenantId` to a tool. + +The 100-current-item limit and latest-50 archived window apply separately to +each tenant-and-agent pair. There is no migration or tenant-wide fallback. +Tell the maker which selected agent they are managing: + +> These announcements belong to **{agent name}**. Their selected audiences apply +> within that agent, not across your other ESS agents. + +## Setup-state check + +For every request that reads or authors an announcement, read +`.local/config.json` before calling any MCP tool. + +If the file does not exist, or its `setup` value is not `"complete"`, show: + +> Welcome to the ESS Maker Kit. Before using `/org-announcements`, type `/setup` to set up your environment. + +and STOP. + +Requests that only ask what announcements are, or what this skill can do, do not +require local setup. Answer them directly. + +## MCP availability check + +Before any request that requires an announcement tool, inspect the tools +available in the current conversation for the `ess-org-announcements` server. + +When its tools are available, continue to **Resolve the target**. + +When its tools are unavailable: + +1. Run: + + ```text + python scripts/mcp_config.py validate --server ess-org-announcements + ``` + +2. Parse `MCP_CONFIG_STATUS_JSON:`: + - `configured`: follow **Start the announcements MCP server**. + - `missing-file` or `missing-server`: run: + + ```text + python scripts/mcp_config.py materialize-defaults + ``` + + Parse `MCP_CONFIG_RESULT_JSON:` and confirm `ess-org-announcements` + appears in `addedServers`, or run `validate` again and confirm its status + is `configured`. Then follow **Start the announcements MCP server**. + - command failure or any other result: show the exact error and stop. Do not + replace malformed JSON or overwrite an existing configuration. + +### Start the announcements MCP server + +Show: + +> The announcements server is configured, but its tools are not available in +> this chat yet. +> +> 1. Press `Ctrl+Shift+P`. +> 2. Run `MCP: List Servers`. +> 3. Select `ess-org-announcements`. +> 4. Choose `Start`. +> +> Type `done` when the server shows `Running`. + +Wait for the maker. When they confirm, inspect the available tools again. If the +tools are available, continue the original request. If they remain unavailable, +tell the maker to reload the VS Code window, rerun `/org-announcements`, and +stop. + +## Resolve the target + +Reuse the setup configuration already loaded. Do not guess an identifier or +use announcement content to decide which agent owns it. + +1. Select the active agent from the backward-compatible `agent` object. For + another configured agent, match an `agents` entry by `slug`, `botId`, or + unambiguous `name`. +2. Use the maker's explicit `titleId`, otherwise the selected entry's stored + `titleId`. A verified target can be reused for this conversation; do not + force `get_agent_config` before each opener. +3. If the title is missing, use this provider's read-only discovery tools: + `list_agent_configs`, then `search_agents` with a distinctive agent-name + substring when the list has no unambiguous match. These discover deployed, + tenant-visible agents; a `botId` is never a substitute for `titleId`. + Both tools belong to `ess-org-announcements`. No landing-page MCP process + is needed; follow this skill's availability check if they are unavailable. +4. Ask the maker to choose among ambiguous candidates. If no candidate matches, + stop and ask them to confirm the agent name and have the published agent + approved and deployed to the organization. Never fall back to tenant-wide + announcements. +5. Persist a discovered, verified `titleId` using the existing local-agent + convention: reread the complete config, find the selected `agents` entry by + `botId` then `slug` (name only if unambiguous), and change only its `titleId`. + If the active agent exists only as `agent`, copy that complete object into + `agents` first. Also update `agent.titleId` when the target matches + `activeAgent` or the active object's botId/slug. Preserve all other fields + and agents. Reread to verify both copies before calling another tool. + With no matching local entry, use the discovered title for this request + without fabricating a partial local agent. + +Discovery does **not** initialize landing-page configuration. Never call +`create_agent_config` or `update_agent_config` as part of this flow, including +when a target was found only through search. No landing-page existence check +or creation is needed to open announcements. + +## Hard rules + +1. Route every call to the `ess-org-announcements` MCP server through this + skill. +2. Call `open_org_announcements` **at most once per maker turn**. It is the only + announcement tool you may call to open a view. +3. `open_org_announcements` reads only. It never creates, updates, publishes, + archives, or deletes. +4. After the widget opens, **stop**. The widget owns every save, publish, + lifecycle action, and audience search for that editing session. Do not call + `save_bulletin`, `transition_bulletin`, or `duplicate_bulletin` — they are not + available to you, and asking for them is a bug. +5. Do not issue any further `search_audience_groups` call once the widget is + open. The widget performs its own searches. +6. Treat all suggested content as reviewable draft state. A suggestion is never + authorization to publish. Say what will open, not what was saved. +7. Never manufacture a bulletin ID, a status, an audit field, or a version. Only + the tools produce canonical state. +8. Never infer that an announcement was created, saved, published, archived, or + deleted from the request or opener alone. The widget reports normal success, + and chat must not duplicate its success message. The only exceptions are the + explicit `IndeterminateWrite` and `CommittedRefreshFailed` results described + below; report those outcomes exactly as instructed. +9. Use the announcement tools for server access. Do not call the backing REST + API directly. +10. Pass the resolved `titleId` on every opener. Only the widget may call + mutations, and it retains this title throughout navigation and retries. + On an agent change, open the new scope; never reuse the previous agent's + IDs, drafts, manager state, or retry request. +11. Keep `titleId` outside `suggestedDraft` and bulletin/editor content. + Audience search remains `{query}` in the tenant directory, not per-agent. + +## Classify the request + +Classify every announcement request into exactly one of these, then follow the +matching flow. + +| Maker intent | Flow | +|---|---| +| "Show me our announcements", "manage announcements", "what's published" | **Open management** | +| "Create an announcement", "new announcement", "let me write one" | **Empty create** | +| "Announce the benefits deadline to Finance", any request with real content | **Pre-hydrated create** | +| "Edit the all-hands announcement", "change the end date on X" | **Edit** | +| "Repost the parking notice", "that one expired, run it again" | **Edit** (review the schedule before publishing again) | + +When the intent is ambiguous, ask one short clarifying question before opening +anything. Opening the wrong view costs the maker a turn. + +### Open management + +Call `open_org_announcements` with: + +```json +{ "titleId": "", "view": "manager" } +``` + +Then stop and let the maker work in the widget. + +### Empty create + +Use this when the maker wants to write the announcement themselves. + +Call `open_org_announcements` with: + +```json +{ "titleId": "", "view": "editor", "mode": "create" } +``` + +Do not invent a title or description to "help". An empty create means empty. + +### Pre-hydrated create + +Use this when the maker's message already carries real announcement content. + +1. Build a `suggestedDraft` from what the maker actually said. Include only the + fields they supplied or clearly implied. Every omitted field falls back to + the editor default, which is what the maker would have seen anyway. +2. Resolve any named audience through `search_audience_groups` **before** the + opener call. See **Resolve suggested audiences**. +3. Call `open_org_announcements` once: + + ```json + { + "titleId": "", + "view": "editor", + "mode": "create", + "suggestedDraft": { + "type": "standard", + "priority": 1, + "title": "", + "description": "", + "primaryAction": { + "actionType": "externalLink", + "label": "