diff --git a/.github/RELEASE.md b/.github/RELEASE.md new file mode 100644 index 0000000..44cf8b1 --- /dev/null +++ b/.github/RELEASE.md @@ -0,0 +1,59 @@ +# Release publication and catalog notification + +Publish only an accepted plugin package bound to its reviewed source/tag and exact archive digest. +Keep an existing public tag and archive immutable. A candidate, draft or notification receipt does +not establish authenticated vendor or installed-host acceptance. + +The Validate and package workflow produces candidate artifacts; it does not make a release public. After the accepted release becomes public, +`notify-catalog.yml` requests a complete catalog reconciliation. It also observes public edits, +channel promotion, unpublishing and deletion; those events never authorize catalog withdrawal by +themselves. The central publisher retains verified history and applies its reviewed withdrawal +policy. It verifies actual GitHub release sources rather than trusting an event payload. + +## Notification authority + +The workflow pins the website's central notification action to a reviewed full commit. Publish that +central commit before enabling a plugin workflow that references it. Review and update this pin +when adopting changes to the notification contract. The caller checks its immutable repository ID, +does not check out package code, and grants its own job token no repository permissions. + +Supply `CATALOG_DISPATCH_TOKEN` using existing reviewed authority with Actions write access to +`computer-mcp/computer-mcp.github.io` only. Website Contents write access is unnecessary. The action +can also receive an existing temporary token directly from a publishing job. Neither workflow +creates or persists credentials. Missing or rejected authority fails visibly; the +publisher's independent schedule still reconciles missed notifications. + +## Publication and retry + +A manual public release emits the release event. Publication performed with a repository's +`GITHUB_TOKEN` does not trigger ordinary release-event workflows. After that publication succeeds, +its automation must explicitly call this reusable workflow as a dependent job: + +```yaml +notify-catalog: + needs: publish + uses: ./.github/workflows/notify-catalog.yml + secrets: + CATALOG_DISPATCH_TOKEN: ${{ secrets.CATALOG_DISPATCH_TOKEN }} +``` + +Here `publish` is the job that actually makes the accepted release public, not the candidate-build +or draft-upload job. When using an existing short-lived token within that publishing job, invoke +the same pinned central action directly after publication instead. Keep token values out of command +arguments, printed output and release metadata. + +For an operator-driven publication or a missed/failed notification, explicitly dispatch: + +```sh +gh workflow run notify-catalog.yml --repo computer-mcp/plugin-cursor --ref main +``` + +This schedules notification using its configured authority; it does not publish or rewrite a +release. Inspect the notification run and its returned central `run_url`. A successful dispatch +proves request acceptance only. Verify the central run completed successfully and the public index +contains the exact expected release identities and generation. If the release is already public and +notification fails, retry notification without changing or republishing the release. Complete +reconciliation is idempotent and repairs duplicate/missed events. + +See the central [catalog publication and notification contract](https://github.com/computer-mcp/computer-mcp.github.io/blob/main/docs/plugin-catalog.md) +for provenance, credentials, retry bounds and deployment semantics. diff --git a/.github/workflows/notify-catalog.yml b/.github/workflows/notify-catalog.yml new file mode 100644 index 0000000..60a5b17 --- /dev/null +++ b/.github/workflows/notify-catalog.yml @@ -0,0 +1,28 @@ +name: Notify official plugin catalog +on: + release: + types: [published, edited, released, unpublished, deleted] + workflow_dispatch: + workflow_call: + secrets: + CATALOG_DISPATCH_TOKEN: + description: Existing receiver-scoped Actions write authority + required: true + outputs: + run_url: + description: Accepted central run; verify its deployment separately + value: ${{ jobs.notify.outputs.run_url }} +permissions: {} +jobs: + notify: + if: github.repository_id == '1384745559' + runs-on: ubuntu-latest + timeout-minutes: 3 + outputs: + run_url: ${{ steps.catalog.outputs.run-url }} + steps: + - name: Request complete catalog reconciliation + id: catalog + uses: computer-mcp/computer-mcp.github.io/.github/actions/notify-catalog@fc27dd0f370d028a3e5021e3585274891f696578 + with: + token: ${{ secrets.CATALOG_DISPATCH_TOKEN }} diff --git a/Documentation/Architecture/Package.md b/Documentation/Architecture/Package.md index 8bc7de2..ec61102 100644 --- a/Documentation/Architecture/Package.md +++ b/Documentation/Architecture/Package.md @@ -6,6 +6,13 @@ The package is independently maintained and distributed. Computer MCP owns regis `bin/cursor-mcp-adapter` owns native sessions, one active operation per session, bounded events and pending permission/question/plan responses. `bin/plugin_runtime.py` is a package-private standard-library implementation of bounded MCP framing/validation, process ownership and retention; it is distributed with this plugin, not installed into or imported from Computer MCP. Each package can update independently. Python 3.13+ is a runtime dependency selected by the launch PATH. +The standard MCP work resource projects these same session owners. The adapter +binds each session to its creation request's host correlation and supplies live +or uncertain state; the MCP server bounds and versions complete snapshots. +Background threads and unconfirmed cleanup remain owned even after the creating +tool returns. Host configuration generations and continuation routing remain +host responsibilities. The resource does not introduce another task lifecycle. + The supervisor lifeline and cleanup receipt are private implementation channels, not plugin-host protocol extensions. Vendor subprocesses inherit neither channel nor host private descriptors. Explicit cleanup acknowledgement is separate from exit status. Runtime errors do not authorize retries or privilege escalation. Scripts provide deterministic packaging, non-model native interface checks and isolated unchanged-host integration checks. Tests exercise deterministic peers and lifecycle failures. The distribution excludes build/test/evidence files. No full native GUI, Windows backend or hosted API is claimed by this package. diff --git a/Documentation/Reference/Installation.md b/Documentation/Reference/Installation.md index a0537ad..56e2019 100644 --- a/Documentation/Reference/Installation.md +++ b/Documentation/Reference/Installation.md @@ -14,6 +14,8 @@ The example exposes this plugin's complete MCP tool catalog with no extra prefix On Computer MCP 1.2.2, plugin activation/selection changes require idle Gateway client admission. Finish or safely pause clients before production installation changes. No host binary replacement or host release is required. Do not restart the active development control connection merely to test installation. +Computer MCP 1.3.0 publishes plugin configuration changes to connected clients atomically. New calls use the current configuration; existing work keeps its owning runtime until release. Installation does not grant access, and later calls use current authorization. + ## Build ```sh @@ -28,16 +30,25 @@ CI runs deterministic fixture tests and packaging; it does not install a vendor ## Isolated host interoperability -After packaging, validate the exact ZIP with an unchanged installed Computer MCP 1.2.2 binary: +After packaging, validate the exact ZIP against the reviewed Computer MCP binary. Select its release version explicitly; a mismatch fails before package extraction: ```sh python3 Scripts/validate_host.py \ - --host "/Applications/Computer MCP.app/Contents/Resources/computer-mcp" \ + --host "/absolute/path/to/candidate/computer-mcp" \ + --expected-host-version 1.3.0 \ --archive /output/PLUGIN.zip \ --output /new/evidence/directory ``` -Replace `PLUGIN.zip` with the package's actual archive name. This uses a temporary directory and the installed host's archive worker and standalone MCP entrypoint. It does not connect to the production App's control socket or database. Vendor tool execution is replaced with inert fixtures; the native version/help check is a separate command. A new evidence directory is required to avoid overwriting an earlier run. +Replace `PLUGIN.zip` with the package's actual archive name. This uses a temporary directory and the selected host's archive worker and standalone MCP entrypoint. It does not connect to the production App's control socket or database. Vendor tool execution is replaced with inert fixtures; the native version/help check is a separate command. A new evidence directory is required to avoid overwriting an earlier run. + +For a candidate host implementing the provider-work contract, add +`--require-work-ownership`. This gate verifies the same connection retains an +opened session, a detached prompt waiting for input and an idle session, then +observes no owners after confirmed session close. It also checks that the +gateway does not reexport the downstream-only resource declaration. Hosts +without this observation surface can still run the ordinary interoperability +gate without that option. A real ACP handshake can be checked independently, without authenticating, opening a conversation or invoking a model: diff --git a/Documentation/Reference/Interface.md b/Documentation/Reference/Interface.md index 0589bac..8196e53 100644 --- a/Documentation/Reference/Interface.md +++ b/Documentation/Reference/Interface.md @@ -8,6 +8,54 @@ The CLI contribution preserves the declared argv order, optional native flags an The MCP adapter implements standard newline-delimited JSON-RPC over stdio, with MCP initialization, tools/list, tools/call, ping and cancellation notifications. Supported MCP dates are 2024-11-05, 2025-03-26 and 2025-06-18. An unsupported proposed date receives the adapter's supported date rather than an unimplemented echo. Tool results use `structuredContent.result`; tool failures set `isError` and include an error code. Protocol errors remain JSON-RPC errors. Tool schemas are the callable source of truth. +## Host risk metadata + +Every MCP tool declares `_meta["io.github.computer-mcp/risk"]`. Model execution +and continuation declare `full-shell`: native permission defaults are not a +host-enforced sandbox. Catalog, result, event and pending-request inspection +declare `read-only`. Cancellation and owned-process retirement declare +`destructive`. The host applies these as minimum classifications, intersects +its own grants, and retains approval authority. Standard MCP annotations remain +hints rather than permissions. + +Session open/load, mode changes and interactive responses also declare +`full-shell`; they can initialize native tools or continue executable work. + +## Runtime work resource + +The adapter advertises ordinary MCP resources and the version-1 declaration in +`_meta["io.github.computer-mcp/work"]` on its tool definitions. `resources/list` +lists `computer-mcp://runtime/work/v1`; `resources/read` returns one JSON text +content entry. The report has `format_version`, a connection-local UUID +`instance_id`, a nonnegative exact integer `revision`, and the complete +`resources` array. Unchanged resource sets keep their revision; changed sets +advance it. At most 1,024 resources and 512 KiB of report text are allowed. + +Each live resource has kind `cursor.session`, the adapter `session` handle as +`id`, the opening call's host-supplied +`_meta["io.github.computer-mcp/work-invocation"]` UUID as `acquired_by`, and +state `active` or `uncertain`. This reference binds lifecycle observation; it +grants no authority and is not a native session ID. It is never read from tool +arguments or forwarded to the vendor process. + +A session is owned from pending startup through idle periods, repeated or +background prompts and interactive requests. Completing a prompt does not close +its session. Removal requires confirmed process cleanup, settled startup and +completion of the owned background thread. Closing with unconfirmed cleanup +reports `uncertain` and keeps the session against capacity. A one-shot prompt +releases its session after the same cleanup boundary. No report can substitute +for permission, native execution success or authenticated model evidence. + +Ordinary clients may continue to call tools without work-invocation metadata. +If such a client creates a live session, the work resource returns unavailable +evidence until unbound work is released; it never reports a falsely empty +snapshot. Malformed invocation metadata is rejected before execution. Report +reads do not launch a vendor, terminate a session, change permissions or replay +work. The adapter does not require resource subscriptions; hosts may poll. + +See the host's [provider work contract](https://github.com/computer-mcp/computer-mcp/blob/master/Documentation/Reference/MCPProtocol.md#downstream-provider-work) +for acquisition expiry, snapshot validation and host-side uncertainty. + ## Tools | Native MCP tool | Behavior | @@ -39,15 +87,15 @@ The default permission policy is `reject-once`. Explicit `allow-once`/`allow-alw {"session":"ADAPTER_HANDLE","request_id":"PENDING_REQUEST","response":{"outcome":{"outcome":"selected","optionId":"EXACT_OFFERED_OPTION"}}} ``` -Question answers must use the offered question and option IDs and obey single/multiple selection. Plans accept `accepted`, `rejected` or `cancelled` and their documented optional fields. Duplicate, stale and unoffered replies fail. Unknown vendor requests receive method-not-found, not fabricated success. +Question answers must use the offered question and option IDs and obey single/multiple selection. Plans accept `accepted`, `rejected` or `cancelled` and their documented optional fields. Duplicate, stale and unoffered replies fail. Each request is bound to its active native operation and session; response delivery and operation retirement are serialized. Unknown vendor requests receive method-not-found, not fabricated success. ## Bounds, cancellation and failures -An ACP frame is bounded to 1 MiB before waiting for a newline. Event retention is bounded by both 256 events and 256 KiB. Text accumulation is bounded to 128 KiB; omissions/truncation are reported. Pages use absolute cursors, `next_cursor`, `has_more` and `missed_events`. An event larger than the requested page becomes explicit omission metadata so pagination can progress. Responses retain the bounded native prompt result; absent stopReason is an error, and native cancellation is not successful task completion. +An ACP frame is bounded to 1 MiB before waiting for a newline. Event retention is bounded by both 256 events and 256 KiB. Text accumulation and its serialized JSON value are each bounded to 128 KiB; omissions/truncation are reported. Pages use absolute cursors, `next_cursor`, `has_more` and `missed_events`. An event larger than the requested page becomes explicit omission metadata so pagination can progress. Responses retain the bounded native prompt result; absent stopReason is an error, and native cancellation is not successful task completion. Prompts default to 300 seconds and accept up to 1800 seconds. Session setup defaults to 45 seconds. Transport startup and cleanup can add bounded latency. Pending requests expire with their native operation. Cancellation sends ACP session/cancel, waits for native settlement, and retires the owned process if it cannot settle within its grace period. A cancelled background request is not automatically replayed. Cancelling the start call after it returned does not identify the background task: use the returned session/prompt handle. -The private supervisor observes adapter EOF/termination and owns the vendor process group. It retains the leader until termination and reaping, then sends a separate bounded cleanup acknowledgement. Missing acknowledgement is `cleanup_unconfirmed`, never a clean success inferred from exit alone. Escaped, independently reparented processes are not claimed as owned. Host callback descriptors/COMPUTER_MCP metadata are not forwarded to the vendor. The process working directory is not an OS sandbox. +The private supervisor observes adapter EOF/termination and owns the vendor process group. It retains the leader until termination and reaping, then sends a separate bounded cleanup acknowledgement. Missing acknowledgement is `cleanup_unconfirmed`, never a clean success inferred from exit alone. A session with unconfirmed cleanup remains retained and counts against admission capacity; repeating close cannot erase the failure. Escaped, independently reparented processes are not claimed as owned. Host callback descriptors/COMPUTER_MCP metadata are not forwarded to the vendor. The process working directory is not an OS sandbox. Representative errors include invalid_arguments, unknown_session, busy, incompatible_vendor, vendor_failed, invalid_vendor_response, frame_too_large, result_too_large, timeout, cancelled and cleanup_unconfirmed. Failed/unknown writes are not automatically retried. Native output may contain sensitive user content; callers must handle it accordingly. @@ -58,3 +106,16 @@ Representative errors include invalid_arguments, unknown_session, busy, incompat - The installed native `agent --help`, `agent --version` and `agent acp --help` used to maintain the pinned CLI tree. Vendor protocol observations, fixture tests, host interoperability and authenticated backend execution are separate evidence classes. + +## Continuation binding + +Tools that accept an existing adapter handle declare +`_meta["io.github.computer-mcp/continuation"]` with format version 1. The selector +matches kind `cursor.session` and primary resource `id` against argument +`session` using JSON Pointer `/session`. This identifies the actual +connection-owned lifetime; it does not rebind acquisition or grant permissions. +New work and unscoped listings do not claim an existing owner. The declaration +uses ordinary MCP metadata and requires no private Host Services. Hosts validate +and retain it on its originating connection; gateway reexports strip it. Runtime +generation selection remains host-owned, and this declaration alone does not +enable live configuration changes. diff --git a/README.md b/README.md index 7263dfb..9d071d8 100644 --- a/README.md +++ b/README.md @@ -8,6 +8,12 @@ The CLI contribution describes verified non-interactive commands. The native bas The adapter provides ACP session open/load, repeated prompts, background prompt start/result, session listing/mode/cancel/close, cursor-paginated events and explicit responses to permission/question/plan requests. `cursor.acp.prompt` remains the one-shot convenience path. Twelve tools are discovered through MCP; discovery does not start a vendor process. +The ordinary MCP work resource reports connection-owned sessions through confirmed +cleanup, including idle sessions, background prompts and interactive requests. +Hosts that support this resource can account for work after a tool reply. It +requires no private Host Services permission and does not itself enable host +configuration changes. + `permission_policy` defaults to `reject-once`. `allow-once` and `allow-always` are explicit native decisions and only select offered options. `manual` exposes pending permission requests for an explicit response. Questions and plans always require an explicit response; no answer is invented. Use `session.open` plus `session.prompt.start` for these workflows. ## Use and verify diff --git a/Scripts/validate_host.py b/Scripts/validate_host.py index ce4ad19..c4d4009 100644 --- a/Scripts/validate_host.py +++ b/Scripts/validate_host.py @@ -108,7 +108,7 @@ def configuration(package, workspace, cli_fixture, acp_fixture, vendor, readonly ''' -def validate(host, archive, output): +def validate(host, archive, output, expected_host_version, require_work_ownership=False): host = host.resolve(strict=True) archive = archive.resolve(strict=True) output.mkdir(parents=True, exist_ok=False) @@ -126,7 +126,8 @@ def validate(host, archive, output): environment = {'PATH': os.pathsep.join([str(Path(sys.executable).parent), '/usr/bin', '/bin', '/usr/sbin', '/sbin']), 'HOME': str(home), 'TMPDIR': str(work), 'PYTHONDONTWRITEBYTECODE': '1', 'LANG': 'en_US.UTF-8'} version = capture([str(host), '--version'], work, environment).decode().strip() - require(version.startswith('1.2.2 '), 'This validator targets Computer MCP 1.2.2; review another host before use') + require(bool(expected_host_version) and version.startswith(expected_host_version + ' '), + 'Host version differs from the explicitly selected acceptance version') inputs = work / 'inputs' inputs.mkdir() (inputs / f'{vendor}.zip').write_bytes(archive.read_bytes()) @@ -166,6 +167,19 @@ def validate(host, archive, output): native_names = [tool['name'] for tool in catalog if tool['name'].startswith(vendor + '.')] require(len(native_names) == (12 if vendor == 'cursor' else 6), 'Adapter catalog missing required tools') checks['catalog'] = {'status':'passed', 'cli_tools':len(projected), 'mcp_tools':len(native_names)} + if require_work_ownership: + require(vendor == 'cursor', 'Work ownership acceptance requires the Cursor session contract') + require(all(key not in tool.get('_meta',{}) for tool in catalog + for key in ('io.github.computer-mcp/work','io.github.computer-mcp/continuation')), + 'Gateway exports advertise downstream-only ownership metadata') + def work_status(count): + def observe(): + servers = checked(client,'mcp.servers.status',{'server':'fixture-adapter'})['servers'] + value = servers[0]['connection'].get('provider_work') + require(isinstance(value,dict), 'Candidate host does not expose provider-work observation') + return value + return wait_for(observe, lambda value:value['resource_count']==count + and value['unsettled_invocation_count']==0 and not value['observation_pending']) prompt = "--leading 'quotes' 中文\nnot-a-shell-command" arguments = {'prompt': prompt} if vendor == 'claude': @@ -187,15 +201,23 @@ def validate(host, archive, output): if vendor == 'cursor': session = checked(client, 'cursor.acp.session.open', {'permission_policy':'manual'})['session'] + work_evidence = {'opened':work_status(1)} if require_work_ownership else None run = checked(client, 'cursor.acp.session.prompt.start', {'session':session,'prompt':'permission'})['prompt_id'] pending = wait_for(lambda: checked(client, 'cursor.acp.requests.list', {'session':session}), lambda v: bool(v['requests']))['requests'][0] + if work_evidence is not None: work_evidence['waiting_for_input'] = work_status(1) checked(client, 'cursor.acp.requests.respond', {'session':session,'request_id':pending['request_id'], 'response':{'outcome':{'outcome':'selected','optionId':'opaque-no'}}}) completed = wait_for(lambda: checked(client, 'cursor.acp.session.prompt.result', {'session':session,'prompt_id':run}), lambda v:v.get('completed')) require(not completed.get('is_error'), 'Background ACP fixture failed') + if work_evidence is not None: work_evidence['idle_session'] = work_status(1) checked(client, 'cursor.acp.events.read', {'session':session,'max_bytes':2048}) checked(client, 'cursor.acp.session.close', {'session':session}) require(not checked(client, 'cursor.acp.session.list')['sessions'], 'Closed ACP session remains live') + if work_evidence is not None: + work_evidence['released'] = work_status(0) + require(len({state['instance_id'] for state in work_evidence.values()})==1, + 'Provider connection changed during session ownership acceptance') + checks['provider_work'] = work_evidence else: run = checked(client, 'claude.run.start', {'prompt':'hello','permission_mode':'plan'})['run_id'] completed = wait_for(lambda: checked(client, 'claude.run.result', {'run_id':run}), lambda v:v.get('completed')) @@ -216,9 +238,30 @@ def validate(host, archive, output): finally: client.close() checks['read_only_profile'] = 'passed' + restricted = configuration(package, workspace, cli_fixture, ROOT / 'Tests/Fixtures/vendor.py', vendor) + restricted = restricted.replace('mode = "local-full-access"', 'mode = "workspace-operations"') + restricted = restricted.replace('full_shell_enabled = true', 'full_shell_enabled = false') + execution_tool = 'cursor.acp.prompt' if vendor == 'cursor' else 'claude.run' + # An explicit low host risk must not bypass the publisher's execution floor. + restricted += '\n[mcp.servers.tool_risks]\n' + json.dumps(execution_tool) + ' = "read-only"\n' + config.write_text(restricted) + client = Client(str(workspace), environment, [str(host), 'serve', 'stdio', '--config', str(config)]) + try: + tools = client.request('tools/list')['result']['tools'] + require(execution_tool not in {tool['name'] for tool in tools}, 'Restricted profile exposed arbitrary vendor execution') + inspection = 'cursor.acp.session.list' if vendor == 'cursor' else 'claude.run.list' + checked(client, inspection) + for name, arguments in [(execution_tool, {'prompt':'must-not-execute'}), + ('mcp.tools.call', {'server':'fixture-adapter','tool':execution_tool,'arguments':{'prompt':'must-not-execute'}})]: + denied = client.call(name, arguments) + require('error' in denied or denied['result'].get('isError'), 'Restricted profile admitted vendor execution') + finally: + client.close() + checks['publisher_floor_under_restricted_profile'] = 'passed' require(digest(host) == host_hash and digest(archive) == archive_hash, 'Host or archive changed during acceptance') report = {'status':'passed', 'observed_at':datetime.datetime.now(datetime.timezone.utc).isoformat(), 'plugin_id':vendor, 'plugin_version':manifest['version'], 'host_version':version, + 'expected_host_version':expected_host_version, 'host_sha256':host_hash, 'archive_sha256':archive_hash, 'checks':checks, 'scope':'host archive validation and ordinary standalone registrations with inert vendor fixtures', 'production_installation':False, 'authenticated_model_execution':False, @@ -230,10 +273,14 @@ def validate(host, archive, output): def main(): parser = argparse.ArgumentParser(description=__doc__) parser.add_argument('--host', required=True, type=Path) + parser.add_argument('--expected-host-version', required=True, + help='Exact reviewed host release version, for example 1.3.0') parser.add_argument('--archive', required=True, type=Path) parser.add_argument('--output', required=True, type=Path) + parser.add_argument('--require-work-ownership', action='store_true', + help='Require a candidate host to observe session ownership through final release') args = parser.parse_args() - validate(args.host, args.archive, args.output) + validate(args.host, args.archive, args.output, args.expected_host_version, args.require_work_ownership) if __name__ == '__main__': diff --git a/Tests/Fixtures/vendor.py b/Tests/Fixtures/vendor.py index 5048c01..cddefd7 100755 --- a/Tests/Fixtures/vendor.py +++ b/Tests/Fixtures/vendor.py @@ -14,6 +14,7 @@ session='fixture-session' waiting={} active=None +ignore_cancel=False def emit(value):print(json.dumps(value),flush=True) def reply(identifier,value):emit({'jsonrpc':'2.0','id':identifier,'result':value}) @@ -31,6 +32,7 @@ def finish(identifier,reason='end_turn'):reply(identifier,{'stopReason':reason}) elif method=='session/set_mode':reply(identifier,{'observed_mode':params['modeId']}) elif method=='session/prompt': prompt=params['prompt'][0]['text'];active=identifier + ignore_cancel=prompt=='ignore-cancel' if prompt in ('permission','question','plan'): if prompt=='permission': name='session/request_permission';value={'sessionId':session,'toolCall':{'toolCallId':'fixture-tool','title':'Fixture permission'},'options':[{'optionId':'opaque-yes','kind':'allow_once','name':'Allow once'},{'optionId':'opaque-no','kind':'reject_once','name':'Reject once'}]} @@ -40,7 +42,7 @@ def finish(identifier,reason='end_turn'):reply(identifier,{'stopReason':reason}) name='cursor/create_plan';value={'toolCallId':'fixture-plan','plan':'Fixture plan','todos':[]} waiting[900]=identifier emit({'jsonrpc':'2.0','id':900,'method':name,'params':value}) - elif prompt=='slow':update('started') + elif prompt in ('slow','ignore-cancel'):update('started') elif prompt=='fork': child=subprocess.Popen([sys.executable,'-c','import time;time.sleep(120)']) Path(os.environ['FIXTURE_MARKER']).write_text(json.dumps([os.getpid(),child.pid])) @@ -56,7 +58,7 @@ def finish(identifier,reason='end_turn'):reply(identifier,{'stopReason':reason}) finish(identifier) else:update('hello');finish(identifier) elif method=='session/cancel': - if active is not None:finish(active,'cancelled');active=None + if active is not None and not ignore_cancel:finish(active,'cancelled');active=None elif method is None and identifier in waiting: update(json.dumps({'received_response':q.get('result')})) finish(waiting.pop(identifier)) diff --git a/Tests/support.py b/Tests/support.py index 63c2aae..2b16b5f 100644 --- a/Tests/support.py +++ b/Tests/support.py @@ -82,6 +82,8 @@ def wait_file(path, seconds=8): def process_alive(pid): r=subprocess.run(['/bin/ps','-p',str(pid),'-o','stat='],capture_output=True,text=True,timeout=1) + if r.stderr.strip() or r.returncode not in (0,1): + raise AssertionError('Cannot verify owned process state: '+r.stderr.strip()) return bool(r.stdout.strip()) and not r.stdout.strip().startswith('Z') diff --git a/Tests/test_adapter.py b/Tests/test_adapter.py index ddedada..43cfd82 100644 --- a/Tests/test_adapter.py +++ b/Tests/test_adapter.py @@ -14,6 +14,14 @@ def value(self,response): self.assertFalse(response['result']['isError'],response) return Client.value(response) def open(self,**args):return self.value(self.call('session.open',args))['session'] + def test_catalog_declares_execution_and_inspection_risk(self): + expected = {'cursor.acp.prompt': 'full-shell', 'cursor.acp.session.open': 'full-shell', 'cursor.acp.session.prompt': 'full-shell', 'cursor.acp.session.prompt.start': 'full-shell', 'cursor.acp.session.prompt.result': 'read-only', 'cursor.acp.session.list': 'read-only', 'cursor.acp.session.mode': 'full-shell', 'cursor.acp.session.cancel': 'destructive', 'cursor.acp.session.close': 'destructive', 'cursor.acp.events.read': 'read-only', 'cursor.acp.requests.list': 'read-only', 'cursor.acp.requests.respond': 'full-shell'} + tools = self.client.request('tools/list')['result']['tools'] + self.assertEqual({tool['name']:tool['_meta']['io.github.computer-mcp/risk'] for tool in tools},expected) + for tool in tools: + risk = expected[tool['name']] + self.assertEqual(tool['annotations']['readOnlyHint'],risk=='read-only') + self.assertEqual(tool['annotations']['destructiveHint'],risk in {'destructive','full-shell'}) def test_acp_prompt_through_mcp(self): r=self.value(self.call('prompt',{'prompt':'hello'})) self.assertEqual(r['text'],'hello');self.assertEqual(r['stop_reason'],'end_turn') @@ -64,6 +72,20 @@ def test_native_cancel_settles_prompt_and_session_can_be_reused(self): self.value(self.call('session.cancel',{'session':session})) self.assertEqual(Client.value(self.client.wait(request))['stop_reason'],'cancelled') self.value(self.call('session.prompt',{'session':session,'prompt':'after-cancel'})) + def test_unsettled_native_cancellation_retires_owned_session(self): + session=self.open() + request=self.client.begin('cursor.acp.session.prompt',{'session':session,'prompt':'ignore-cancel'}) + for _ in range(80): + if self.value(self.call('events.read',{'session':session}))['events']:break + time.sleep(.02) + else:self.fail('Native prompt did not start') + started=time.monotonic() + self.value(self.call('session.cancel',{'session':session})) + result=Client.value(self.client.wait(request)) + self.assertEqual(result['error']['code'],'cancelled') + self.assertLess(time.monotonic()-started,5) + self.assertTrue(self.call('session.prompt',{'session':session,'prompt':'after'})['result']['isError']) + def test_early_exit_does_not_drop_buffered_result(self): self.assertEqual(self.value(self.call('prompt',{'prompt':'exit-fast'}))['stop_reason'],'end_turn') def test_aggregate_events_are_bounded_and_cursor_progresses(self): diff --git a/Tests/test_protocol_admission.py b/Tests/test_protocol_admission.py new file mode 100644 index 0000000..29c7dc7 --- /dev/null +++ b/Tests/test_protocol_admission.py @@ -0,0 +1,76 @@ +"""Malformed peer input must not terminate unrelated MCP work.""" +import json +import sys +import threading +import unittest +from unittest.mock import patch +from support import ROOT +sys.path.insert(0,str(ROOT/'bin')) +import plugin_runtime as runtime + +class ProtocolAdmissionTests(unittest.TestCase): + def test_invalid_unicode_is_rejected_in_all_json_string_positions(self): + for raw in (b'{"id":"\\ud800"}',b'{"\\udfff":1}',b'["\\ud800"]'): + with self.subTest(raw=raw),self.assertRaises(ValueError): + runtime.decoded(raw) + self.assertEqual(runtime.decoded(b'"\\ud83d\\ude00"'),'\U0001f600') + def test_bad_id_does_not_end_active_tool_or_connection(self): + entered,release=threading.Event(),threading.Event() + def handler(_args,job): + entered.set() + if not release.wait(2): raise AssertionError('fixture not released') + job.check() + return {'finished':True} + server=runtime.MCPServer('fixture','1',[runtime.Tool('hold','fixture',runtime.schema({}),handler,risk='read-only')],lambda:None) + server.initialized=True + messages=[] + def write(_fd,data,*_): messages.append(json.loads(data)) + with patch.object(runtime,'write_bytes',side_effect=write): + server.dispatch({'jsonrpc':'2.0','id':1,'method':'tools/call','params':{'name':'hold'}}) + try: + self.assertTrue(entered.wait(1)) + server.dispatch({'jsonrpc':'2.0','id':'\ud800','method':'ping'}) + self.assertEqual(messages[-1]['error']['code'],-32600) + self.assertIsNone(messages[-1]['id']) + server.dispatch({'jsonrpc':'2.0','id':2,'method':'ping'}) + self.assertEqual(messages[-1]['result'],{}) + self.assertFalse(server.stop.is_set()) + finally: + with server.lock: threads=list(server.threads) + release.set() + for thread in threads: thread.join(2) + completion=next(m for m in messages if m['id']==1) + self.assertFalse(completion['result']['isError']) + def test_active_request_id_cannot_be_reused_by_ping(self): + server=runtime.MCPServer('fixture','1',[],lambda:None) + server.jobs[server.id_key(5)]=runtime.Job() + with patch.object(server,'emit') as emit: + server.dispatch({'jsonrpc':'2.0','id':5,'method':'ping'}) + self.assertIn('error',emit.call_args.args[0]) + self.assertIn(server.id_key(5),server.jobs) + + def test_denied_process_inspection_cannot_confirm_cleanup(self): + child=unittest.mock.Mock(pid=42,returncode=0) + writes=[] + report=unittest.mock.Mock(returncode=1,stdout='',stderr='permission denied') + with patch.object(runtime.subprocess,'Popen',return_value=child), patch.object(runtime.os,'waitid',return_value=object()), patch.object(runtime.os,'killpg',side_effect=PermissionError()), patch.object(runtime.subprocess,'run',return_value=report), patch.object(runtime.os,'write',side_effect=lambda fd,data:writes.append((fd,data))), patch.object(runtime.os,'close'), patch.object(runtime.signal,'signal'): + result=runtime.supervise(10,11,['/unused']) + self.assertEqual(result,70) + receipt=runtime.decoded(next(data for fd,data in writes if fd==11)) + self.assertFalse(receipt['cleanup_confirmed']) + + def test_protocol_negotiation_uses_supported_dates_and_typed_fields(self): + for date in (*runtime.SUPPORTED_MCP,'unsupported'): + server=runtime.MCPServer('fixture','1',[],lambda:None) + with patch.object(server,'emit') as emit: + server.dispatch({'jsonrpc':'2.0','id':1,'method':'initialize','params':{'protocolVersion':date,'capabilities':{},'clientInfo':{'name':'test','version':'1'}}}) + self.assertIn(emit.call_args.args[0]['result']['protocolVersion'],runtime.SUPPORTED_MCP) + for fields in ({'protocolVersion':True},{'protocolVersion':'2025-06-18','capabilities':[]}, + {'protocolVersion':'2025-06-18','capabilities':{},'clientInfo':{'name':4,'version':'1'}}): + server=runtime.MCPServer('fixture','1',[],lambda:None) + with patch.object(server,'emit') as emit: + server.dispatch({'jsonrpc':'2.0','id':1,'method':'initialize','params':fields}) + self.assertIn('error',emit.call_args.args[0]) + self.assertFalse(server.initialized) + +if __name__=='__main__': unittest.main() diff --git a/Tests/test_review_regressions.py b/Tests/test_review_regressions.py new file mode 100644 index 0000000..57016e2 --- /dev/null +++ b/Tests/test_review_regressions.py @@ -0,0 +1,186 @@ +"""ACP request identity, operation ownership and terminal-state regressions.""" +import queue +import threading +import time +import unittest +from unittest.mock import patch +from test_completion import adapter +from plugin_runtime import Failure, Job, encoded + +PERMISSION = {'jsonrpc':'2.0','id':900,'method':'session/request_permission', + 'params':{'sessionId':'native-session','options':[{'optionId':'yes','kind':'allow_once'}]}} +ALLOW = {'outcome':{'outcome':'selected','optionId':'yes'}} + +class Peer: + def __init__(self): + self.frames=queue.Queue() + self.sent=[] + self.requested=threading.Event() + self.closed=False + self.on_send=lambda _:None + def send(self,message,*_): + self.on_send(message) + self.sent.append(message) + if 'method' in message and 'id' in message:self.requested.set() + def read(self,deadline,*_): + try:return encoded(self.frames.get(timeout=max(0,deadline-time.monotonic()))) + except queue.Empty:raise Failure('timeout','fixture read deadline') + def close(self):self.closed=True + +class ReviewRegressions(unittest.TestCase): + def begin(self,policy='manual',method='session/prompt'): + session=adapter.Session('handle','/unused',policy) + session.vendor_id='native-session' + peer=Peer();session.process=peer + result=[] + def execute(): + try:result.append(session.request(method,{},time.monotonic()+3,Job())) + except Exception as error:result.append(error) + thread=threading.Thread(target=execute) + thread.start() + self.assertTrue(peer.requested.wait(1)) + def cleanup(): + peer.frames.put({'jsonrpc':'2.0','id':session.sequence,'result':{'stopReason':'end_turn'}}) + thread.join(4) + self.assertFalse(thread.is_alive()) + self.addCleanup(cleanup) + return session,peer,thread,result + def pending(self,session): + deadline=time.monotonic()+1 + while time.monotonic()32768: - self.respond_raw(message['id'],error={'code':-32602,'message':'Request parameters exceed the adapter contract'}) - return - if method == 'session/request_permission' and self.policy != 'manual': - response = permission_choice(params,self.policy) - self.respond_raw(message['id'],response) - self.events.append({'method':method,'params':params,'response':response}) - elif method in ('session/request_permission','cursor/ask_question','cursor/create_plan'): - with self.lock: + with self.lock: + try: key = MCPServer.id_key(message.get('id')) + except ValueError as error: + raise Failure('invalid_vendor_request','ACP request has an invalid identity') from error + operation = self.active_operation + if operation is None or self.closed.is_set(): + raise Failure('stale_request','No operation owns this native request') + if any(MCPServer.id_key(item['id']) == key for item in self.pending.values()): + raise Failure('invalid_vendor_request','ACP request identity is already active') + method, params = message['method'], message.get('params',{}) + self.check_session(method,params) + if len(encoded(params))>32768: + self.respond_raw(message['id'],error={'code':-32602,'message':'Request parameters exceed the adapter contract'},deadline=deadline) + return + if operation['job'].cancelled.is_set() or (operation['method']=='session/prompt' and self.cancelled.is_set()) or deadline<=time.monotonic(): + self.respond_raw(message['id'],{'outcome':{'outcome':'cancelled'}},deadline=deadline) + elif method == 'session/request_permission' and self.policy != 'manual': + response = permission_choice(params,self.policy) + self.respond_raw(message['id'],response,deadline=deadline) + self.events.append({'method':method,'params':params,'response':response}) + elif method in ('session/request_permission','cursor/ask_question','cursor/create_plan'): if len(self.pending)>=8: - self.respond_raw(message['id'],{'outcome':{'outcome':'cancelled'}}) + self.respond_raw(message['id'],{'outcome':{'outcome':'cancelled'}},deadline=deadline) return token = uuid.uuid4().hex - self.pending[token] = {'id':message['id'],'method':method,'params':params,'deadline':deadline} - self.events.append({'method':method,'request_id':token,'params':params}) - else: - self.respond_raw(message['id'],error={'code':-32601,'message':'Unsupported ACP client request'}) + self.pending[token] = {'id':message['id'],'method':method,'params':params,'deadline':deadline,'operation':operation['id']} + self.events.append({'method':method,'request_id':token,'params':params}) + else: + self.respond_raw(message['id'],error={'code':-32601,'message':'Unsupported ACP client request'},deadline=deadline) def request(self, method, params, deadline, job): if not self.operation.acquire(blocking=False): raise Failure('busy','Another request owns this Cursor session') try: if self.closed.is_set(): raise Failure('session_closed','Cursor session is closed') - self.sequence += 1 - request_id = self.sequence + with self.lock: + self.sequence += 1 + request_id = self.sequence + self.active_operation = {'id':request_id,'deadline':deadline,'job':job,'method':method} self.process.send({'jsonrpc':'2.0','id':request_id,'method':method,'params':params},deadline,job) cancellation_sent = False while True: @@ -180,7 +235,7 @@ class Session: self.process.send({'jsonrpc':'2.0','method':'session/cancel','params':{'sessionId':self.vendor_id}},min(deadline,time.monotonic()+1)) cancellation_sent = True deadline = min(deadline,time.monotonic()+2) - self.cancel_pending() + self.cancel_pending(deadline) try: raw = self.process.read(min(deadline,time.monotonic()+.1)) except Failure as error: @@ -197,6 +252,10 @@ class Session: if not isinstance(message,dict) or message.get('jsonrpc')!='2.0': raise Failure('invalid_vendor_response','Cursor emitted an invalid ACP envelope') if message.get('id')==request_id and type(message.get('id')) is int and ('result' in message or 'error' in message): + if ('result' in message) == ('error' in message) or 'method' in message: + raise Failure('invalid_vendor_response','ACP response must contain exactly one result or error') + if 'error' in message and (not isinstance(message['error'],dict) or type(message['error'].get('code')) is not int or not isinstance(message['error'].get('message'),str)): + raise Failure('invalid_vendor_response','ACP error has invalid fields') if cancellation_sent: # The matching response consumes the outstanding RPC; this session can be reused. return {'stopReason':'cancelled','cancelled':True} @@ -209,8 +268,11 @@ class Session: raise Failure('invalid_vendor_response','ACP result must be an object') return value if 'id' in message and isinstance(message.get('method'),str): + if 'result' in message or 'error' in message: + raise Failure('invalid_vendor_response','ACP request contains response fields') self.vendor_request(message,deadline) elif isinstance(message.get('method'),str): + self.check_session(message['method'],message.get('params',{})) self.events.append(message) update = message.get('params',{}).get('update',{}) if isinstance(message.get('params',{}),dict) else {} if isinstance(update,dict) and update.get('sessionUpdate')=='agent_message_chunk': @@ -226,7 +288,14 @@ class Session: self.close() raise finally: - self.operation.release() + try: + self.cancel_pending(deadline) + except Exception: + self.close() + raise + finally: + with self.lock: self.active_operation = None + self.operation.release() def set_mode(self, mode, deadline, job): available = self.session_info.get('modes',{}).get('availableModes',[]) if not any(isinstance(x,dict) and x.get('id')==mode for x in available): @@ -250,10 +319,19 @@ class Session: if len(encoded(result)) > 131072: self.close() raise Failure('result_too_large', 'Cursor prompt result exceeds its byte bound') + text = self.text.decode('utf-8',errors='ignore') + # JSON escaping counts against the text budget as well as UTF-8 bytes. + if len(encoded(text))>131072: + low, high = 0, len(text) + while low=8: raise Failure('capacity','Cursor session capacity reached or adapter is shutting down') self.sessions[session.handle] = session @@ -340,8 +440,10 @@ class Cursor: raise Failure('unknown_session','No session with this handle belongs to this adapter connection') return session def close(self,args,_job): - with self.lock: session = self.sessions.pop(args['session'],None) - if session: session.close() + with self.lock: session = self.sessions.get(args['session']) + if session: + session.close() + with self.lock: self.sessions.pop(args['session'],None) return {'closed':True,'already_closed':session is None} def once(self,args,job): setup = {k:v for k,v in args.items() if k in OPEN_FIELDS} @@ -375,35 +477,40 @@ class Cursor: def listing(self,_args,_job): with self.lock: sessions = list(self.sessions.values()) return {'sessions':[s.info() for s in sessions]} + def work_resources(self): + with self.lock: sessions = list(self.sessions.values()) + return [resource for session in sessions if (resource:=session.work_resource()) is not None] def mode(self,args,job): session = self.get(args['session']) return session.set_mode(args['mode'],time.monotonic()+15,job) def shutdown(self): with self.lock: self.stopping = True - sessions, self.sessions = list(self.sessions.values()), {} + sessions = list(self.sessions.values()) errors = [] for session in sessions: try: session.close() + if session.released(): + with self.lock: self.sessions.pop(session.handle,None) except Exception as error: errors.append(error) if errors: raise Failure('cleanup_unconfirmed', 'One or more Cursor sessions could not confirm cleanup') def tools(self): return [ - Tool('cursor.acp.prompt','Run one Cursor ACP prompt and retire its owned process. Use explicit sessions for interactive input.',schema({**OPEN_FIELDS,**PROMPT_FIELDS},['prompt']),self.once), - Tool('cursor.acp.session.open','Initialize Cursor ACP, authenticate existing vendor credentials, and create or load a session in the host-selected working directory.',schema(OPEN_FIELDS),self.open), - Tool('cursor.acp.session.prompt','Prompt an existing Cursor ACP session; read events and answer pending requests concurrently.',schema({'session':SESSION,**PROMPT_FIELDS},['session','prompt']),self.prompt), - Tool('cursor.acp.session.prompt.start','Start a prompt without holding a long MCP request; use its prompt ID to read completion.',schema({'session':SESSION,**PROMPT_FIELDS},['session','prompt']),self.start_prompt), - Tool('cursor.acp.session.prompt.result','Read the most recent background prompt result without replaying it.',schema({'session':SESSION,'prompt_id':text_schema(64)},['session','prompt_id']),self.prompt_result), - Tool('cursor.acp.session.list','List only this adapter connection\'s Cursor sessions.',schema({}),self.listing), - Tool('cursor.acp.session.mode','Select a mode advertised by the current Cursor session.',schema({'session':SESSION,'mode':MODE},['session','mode']),self.mode), - Tool('cursor.acp.session.cancel','Request native ACP prompt cancellation without guessing a permission response.',schema({'session':SESSION},['session']),self.cancel), - Tool('cursor.acp.session.close','Close this Cursor session and retire its process group.',schema({'session':SESSION},['session']),self.close), - Tool('cursor.acp.events.read','Read byte-bounded Cursor events using a session-local cursor.',schema({'session':SESSION,'after_cursor':integer_schema(0,2**53-1,0),'limit':integer_schema(1,256,128),'max_bytes':integer_schema(1024,196608,196608)},['session']),self.events), - Tool('cursor.acp.requests.list','Read pending Cursor permission, question and plan requests; this does not approve them.',schema({'session':SESSION},['session']),self.pending), - Tool('cursor.acp.requests.respond','Respond once to an exact pending vendor request using its native outcome schema.',schema({'session':SESSION,'request_id':text_schema(64),'response':{'type':'object'}},['session','request_id','response']),self.answer), + Tool('cursor.acp.prompt','Run one Cursor ACP prompt and retire its owned process. Use explicit sessions for interactive input.',schema({**OPEN_FIELDS,**PROMPT_FIELDS},['prompt']),self.once,risk='full-shell'), + Tool('cursor.acp.session.open','Initialize Cursor ACP, authenticate existing vendor credentials, and create or load a session in the host-selected working directory.',schema(OPEN_FIELDS),self.open,risk='full-shell'), + Tool('cursor.acp.session.prompt','Prompt an existing Cursor ACP session; read events and answer pending requests concurrently.',schema({'session':SESSION,**PROMPT_FIELDS},['session','prompt']),self.prompt,risk='full-shell',continuation=('cursor.session','session')), + Tool('cursor.acp.session.prompt.start','Start a prompt without holding a long MCP request; use its prompt ID to read completion.',schema({'session':SESSION,**PROMPT_FIELDS},['session','prompt']),self.start_prompt,risk='full-shell',continuation=('cursor.session','session')), + Tool('cursor.acp.session.prompt.result','Read the most recent background prompt result without replaying it.',schema({'session':SESSION,'prompt_id':text_schema(64)},['session','prompt_id']),self.prompt_result,risk='read-only',continuation=('cursor.session','session')), + Tool('cursor.acp.session.list','List only this adapter connection\'s Cursor sessions.',schema({}),self.listing,risk='read-only'), + Tool('cursor.acp.session.mode','Select a mode advertised by the current Cursor session.',schema({'session':SESSION,'mode':MODE},['session','mode']),self.mode,risk='full-shell',continuation=('cursor.session','session')), + Tool('cursor.acp.session.cancel','Request native ACP prompt cancellation without guessing a permission response.',schema({'session':SESSION},['session']),self.cancel,risk='destructive',continuation=('cursor.session','session')), + Tool('cursor.acp.session.close','Close this Cursor session and retire its process group.',schema({'session':SESSION},['session']),self.close,risk='destructive',continuation=('cursor.session','session')), + Tool('cursor.acp.events.read','Read byte-bounded Cursor events using a session-local cursor.',schema({'session':SESSION,'after_cursor':integer_schema(0,2**53-1,0),'limit':integer_schema(1,256,128),'max_bytes':integer_schema(1024,196608,196608)},['session']),self.events,risk='read-only',continuation=('cursor.session','session')), + Tool('cursor.acp.requests.list','Read pending Cursor permission, question and plan requests; this does not approve them.',schema({'session':SESSION},['session']),self.pending,risk='read-only',continuation=('cursor.session','session')), + Tool('cursor.acp.requests.respond','Respond once to an exact pending vendor request using its native outcome schema.',schema({'session':SESSION,'request_id':text_schema(64),'response':{'type':'object'}},['session','request_id','response']),self.answer,risk='full-shell',continuation=('cursor.session','session')), ] @@ -413,7 +520,7 @@ def main(): parser.add_argument('--executable',help='Absolute user-owned Cursor Agent path; otherwise CURSOR_AGENT_EXECUTABLE or PATH') args = parser.parse_args() cursor = Cursor(args.executable) - MCPServer('cursor-mcp-adapter',VERSION,cursor.tools(),cursor.shutdown).serve() + MCPServer('cursor-mcp-adapter',VERSION,cursor.tools(),cursor.shutdown,work=cursor.work_resources).serve() if __name__=='__main__': main() diff --git a/bin/plugin_runtime.py b/bin/plugin_runtime.py index 0c64892..33e056e 100644 --- a/bin/plugin_runtime.py +++ b/bin/plugin_runtime.py @@ -16,10 +16,15 @@ import threading import time import tomllib +import uuid MAX_FRAME = 1_048_576 MAX_RESULT = 196_608 SUPPORTED_MCP = ('2024-11-05', '2025-03-26', '2025-06-18') +WORK_URI = 'computer-mcp://runtime/work/v1' +WORK_METADATA = 'io.github.computer-mcp/work' +CONTINUATION_METADATA = 'io.github.computer-mcp/continuation' +WORK_INVOCATION = 'io.github.computer-mcp/work-invocation' class Failure(Exception): @@ -47,7 +52,16 @@ def finite_float(text): if not math.isfinite(value): raise ValueError('JSON floating-point number is not finite') return value - return json.loads(raw.decode('utf-8'), object_pairs_hook=unique, parse_constant=invalid_number, parse_float=finite_float) + value = json.loads(raw.decode('utf-8'), object_pairs_hook=unique, parse_constant=invalid_number, parse_float=finite_float) + pending = [value] + while pending: + item = pending.pop() + if isinstance(item,str): item.encode('utf-8') + elif isinstance(item,dict): + pending.extend(item.keys()) + pending.extend(item.values()) + elif isinstance(item,list): pending.extend(item) + return value def schema(properties, required=()): @@ -107,8 +121,9 @@ def validate(value, spec, location='arguments', depth=0): class Job: - def __init__(self): + def __init__(self, work_invocation=None): self.cancelled = threading.Event() + self.work_invocation = work_invocation def cancel(self): self.cancelled.set() def check(self): @@ -237,7 +252,7 @@ def supervise(read_fd, receipt_fd, command): ['/bin/ps', '-g', str(child.pid), '-o', 'pid=,stat='], capture_output=True, text=True, timeout=1) rows = [line.split() for line in report.stdout.splitlines() if line.strip()] - if report.returncode not in (0, 1) or any( + if report.stderr.strip() or report.returncode not in (0, 1) or any( len(row) != 2 or not row[1].startswith('Z') for row in rows ): raise Failure('cleanup_unconfirmed', 'Owned group could not be retired') @@ -285,10 +300,15 @@ def __init__(self, command, cwd): self.stderr = bytearray() self.stderr_truncated = False self.stop_stderr = threading.Event() - os.set_blocking(self.child.stdin.fileno(), False) - self.output = LineReader(self.child.stdout) - self.stderr_thread = threading.Thread(target=self._stderr, daemon=True) - self.stderr_thread.start() + self.output, self.stderr_thread = None, None + try: + os.set_blocking(self.child.stdin.fileno(), False) + self.output = LineReader(self.child.stdout) + self.stderr_thread = threading.Thread(target=self._stderr, daemon=True) + self.stderr_thread.start() + except BaseException: + self.close() + raise def _stderr(self): try: with selectors.DefaultSelector() as selector: @@ -346,30 +366,33 @@ def close(self): raise self.cleanup_failure return self.closed = True - os.close(self.life) try: try: - self.child.wait(timeout=3) - except subprocess.TimeoutExpired: - self.child.terminate() + os.close(self.life) try: - self.child.wait(timeout=2) - except subprocess.TimeoutExpired as error: - self.child.kill() - self.child.wait(timeout=2) - raise Failure('cleanup_unconfirmed', 'Vendor supervisor did not confirm shutdown') from error - self.confirm_cleanup(self.receipt) - except Failure as error: - self.cleanup_failure = error - raise - finally: - os.close(self.receipt) - self.output.close() - self.stop_stderr.set() - self.stderr_thread.join(timeout=.5) - with self.writer_lock: - for stream in (self.child.stdin, self.child.stdout, self.child.stderr): - stream.close() + self.child.wait(timeout=3) + except subprocess.TimeoutExpired: + self.child.terminate() + try: + self.child.wait(timeout=2) + except subprocess.TimeoutExpired as error: + self.child.kill() + self.child.wait(timeout=2) + raise Failure('cleanup_unconfirmed', 'Vendor supervisor did not confirm shutdown') from error + self.confirm_cleanup(self.receipt) + finally: + os.close(self.receipt) + if self.output is not None: self.output.close() + self.stop_stderr.set() + if self.stderr_thread is not None and self.stderr_thread.ident is not None: + self.stderr_thread.join(timeout=.5) + with self.writer_lock: + for stream in (self.child.stdin, self.child.stdout, self.child.stderr): + stream.close() + except Exception as error: + self.cleanup_failure = error if isinstance(error,Failure) else Failure('cleanup_unconfirmed','Owned process cleanup could not be confirmed') + if self.cleanup_failure is error: raise + raise self.cleanup_failure from error def resolve_executable(explicit, environment_key, default): @@ -391,7 +414,7 @@ def load_package(entry): return manifest, tree -def verify_version(executable, tree, job, cwd): +def verify_version(executable, tree, job, cwd, retain_process=None): expected = next((x.get('stdout') for x in tree.get('executable_checks', []) if x.get('args') == ['--version']), None) if expected is None: raise Failure('invalid_configuration', 'Package lacks its native version assertion') @@ -399,6 +422,7 @@ def verify_version(executable, tree, job, cwd): deadline = time.monotonic()+5 output = bytearray() try: + if retain_process is not None: retain_process(process) process.end_input() while True: raw = process.read(deadline, job) @@ -459,12 +483,36 @@ class Tool: description: str input_schema: dict handler: object + risk: str + continuation: tuple[str,str] | None = None + def __post_init__(self): + if self.risk not in {'read-only','workspace-write','external-write','destructive','full-shell'}: + raise ValueError('Every tool requires a known publisher risk classification') + if self.continuation is not None: + kind, argument = self.continuation + if (not isinstance(kind,str) or not kind or len(kind.encode('utf-8'))>1024 + or any(ord(c)<32 or 127<=ord(c)<=159 for c in kind) + or not isinstance(argument,str) or not argument + or len(('/'+argument.replace('~','~0').replace('/','~1')).encode('utf-8'))>1024 + or any(ord(c)<32 or 127<=ord(c)<=159 for c in argument) + or argument not in self.input_schema.get('required',[]) + or self.input_schema.get('properties',{}).get(argument,{}).get('type') not in {'string','integer'}): + raise ValueError('Continuation requires a resource kind and required scalar handle') def definition(self): - return {'name':self.name,'description':self.description,'inputSchema':self.input_schema} + read_only = self.risk == 'read-only' + metadata = {'io.github.computer-mcp/risk':self.risk} + if self.continuation is not None: + kind, argument = self.continuation + pointer = '/' + argument.replace('~','~0').replace('/','~1') + metadata[CONTINUATION_METADATA] = {'format_version':1,'selectors':[{'kind':kind,'handles':{'id':pointer}}]} + return {'name':self.name,'description':self.description,'inputSchema':self.input_schema, + '_meta':metadata, + 'annotations':{'readOnlyHint':read_only,'destructiveHint':self.risk in {'destructive','full-shell'}, + 'idempotentHint':read_only,'openWorldHint':not read_only}} class MCPServer: - def __init__(self, name, version, tools, shutdown): + def __init__(self, name, version, tools, shutdown, work=None): self.name, self.version = name, version self.tools = {t.name:t for t in tools} self.shutdown = shutdown @@ -472,8 +520,54 @@ def __init__(self, name, version, tools, shutdown): self.jobs, self.threads = {}, set() self.lock, self.write_lock = threading.Lock(), threading.Lock() self.initialized = False + self.work = work + self.work_instance = str(uuid.uuid4()) + self.work_revision = 0 + self.work_last = None + def work_resource(self): + resources = self.work() + if not isinstance(resources,list) or len(resources)>1024: + raise Failure('work_unavailable','Complete work observation is unavailable') + keys = set() + rows = [] + def identifier(value): + return isinstance(value,str) and 0=2**63: + raise Failure('work_unavailable','Work observation revision is exhausted') + snapshot = {'format_version':1,'instance_id':self.work_instance,'revision':revision,'resources':rows} + body = encoded(snapshot) + if len(body)>524288: + raise Failure('work_unavailable','Complete work observation exceeds its byte bound') + self.work_revision, self.work_last = revision, rows + return {'contents':[{'uri':WORK_URI,'mimeType':'application/json','text':body.decode('utf-8')}]} def emit(self, message): - raw = encoded(message)+b'\n' + try: + raw = encoded(message)+b'\n' + except (ValueError,UnicodeError,RecursionError): + try: + self.id_key(message.get('id')) + identifier = message['id'] + except ValueError: + identifier = None + raw = encoded({'jsonrpc':'2.0','id':identifier,'error':{'code':-32603,'message':'Response cannot be encoded'}})+b'\n' if len(raw) > MAX_FRAME: raw = encoded({'jsonrpc':'2.0','id':message.get('id'),'error':{'code':-32603,'message':'Response exceeds the protocol bound'}})+b'\n' try: @@ -507,6 +601,7 @@ def call(self, request_id, key, tool, arguments, job): def id_key(value): if type(value) not in (int,str) or (isinstance(value,str) and (not value or len(value)>256)): raise ValueError('Invalid request ID') + if isinstance(value,str): value.encode('utf-8') return type(value).__name__, value def dispatch(self, message): if not isinstance(message, dict) or message.get('jsonrpc') != '2.0' or not isinstance(message.get('method'),str): @@ -524,16 +619,22 @@ def dispatch(self, message): try: key = self.id_key(request_id) except ValueError: self.rpc_error(None,-32600,'Invalid request ID'); return + with self.lock: active = key in self.jobs + if active: + self.rpc_error(request_id,-32600,'Request ID is already active'); return if not isinstance(params,dict): self.rpc_error(request_id,-32602,'Parameters must be an object'); return if method == 'initialize': if self.initialized: self.rpc_error(request_id,-32600,'Connection is already initialized'); return requested = params.get('protocolVersion') - if not isinstance(requested,str): - self.rpc_error(request_id,-32602,'protocolVersion is required'); return + client = params.get('clientInfo') + if not isinstance(requested,str) or not isinstance(params.get('capabilities'),dict) or not isinstance(client,dict) or not all(isinstance(client.get(k),str) and client[k] for k in ('name','version')): + self.rpc_error(request_id,-32602,'Initialize requires protocolVersion, capabilities and typed clientInfo'); return self.initialized = True - self.emit({'jsonrpc':'2.0','id':request_id,'result':{'protocolVersion':requested if requested in SUPPORTED_MCP else SUPPORTED_MCP[-1],'capabilities':{'tools':{'listChanged':False}},'serverInfo':{'name':self.name,'version':self.version}}}) + capabilities = {'tools':{'listChanged':False}} + if self.work is not None: capabilities['resources'] = {'subscribe':False,'listChanged':False} + self.emit({'jsonrpc':'2.0','id':request_id,'result':{'protocolVersion':requested if requested in SUPPORTED_MCP else SUPPORTED_MCP[-1],'capabilities':capabilities,'serverInfo':{'name':self.name,'version':self.version}}}) elif method == 'ping': self.emit({'jsonrpc':'2.0','id':request_id,'result':{}}) elif not self.initialized: @@ -541,18 +642,44 @@ def dispatch(self, message): elif method == 'tools/list': if params.get('cursor'): self.rpc_error(request_id,-32602,'This finite catalog has no continuation cursor'); return - self.emit({'jsonrpc':'2.0','id':request_id,'result':{'tools':[t.definition() for t in self.tools.values()]}}) + definitions = [t.definition() for t in self.tools.values()] + if self.work is not None: + for definition in definitions: + definition['_meta'][WORK_METADATA] = {'format_version':1,'uri':WORK_URI} + self.emit({'jsonrpc':'2.0','id':request_id,'result':{'tools':definitions}}) + elif method == 'resources/list' and self.work is not None: + if params.get('cursor'): + self.rpc_error(request_id,-32602,'This finite resource catalog has no continuation cursor'); return + self.emit({'jsonrpc':'2.0','id':request_id,'result':{'resources':[{'uri':WORK_URI,'name':'Runtime work','mimeType':'application/json','description':'Complete connection-owned live work and cleanup observation.'}]}}) + elif method == 'resources/read' and self.work is not None: + if params.get('uri') != WORK_URI: + self.rpc_error(request_id,-32602,'Unknown resource'); return + try: result = self.work_resource() + except Failure as error: + self.rpc_error(request_id,-32000,str(error)); return + except Exception: + self.rpc_error(request_id,-32000,'Complete work observation is unavailable'); return + self.emit({'jsonrpc':'2.0','id':request_id,'result':result}) elif method == 'tools/call': name = params.get('name') if not isinstance(name,str) or name not in self.tools: self.rpc_error(request_id,-32602,'Unknown tool'); return args = params.get('arguments', {}) + origin = None + if self.work is not None: + meta = params.get('_meta',{}) + if not isinstance(meta,dict): + self.rpc_error(request_id,-32602,'Tool metadata must be an object'); return + if WORK_INVOCATION in meta: + try: origin = str(uuid.UUID(meta[WORK_INVOCATION])) + except (ValueError,TypeError,AttributeError): + self.rpc_error(request_id,-32602,'Work invocation reference must be a UUID'); return with self.lock: if key in self.jobs: self.rpc_error(request_id,-32600,'Request ID is already active'); return if len(self.jobs)>=16: self.rpc_error(request_id,-32000,'Concurrent request capacity reached'); return - job = Job() + job = Job(work_invocation=origin) thread = threading.Thread(target=self.call,args=(request_id,key,self.tools[name],args,job)) self.jobs[key] = job self.threads.add(thread) diff --git a/computer-mcp-plugin.toml b/computer-mcp-plugin.toml index c9dba12..f7c427d 100644 --- a/computer-mcp-plugin.toml +++ b/computer-mcp-plugin.toml @@ -1,6 +1,6 @@ id = "cursor" name = "Cursor" -version = "0.1.0" +version = "0.1.1" repository = "https://github.com/computer-mcp/plugin-cursor" description = "Cursor Agent CLI automation plus an ACP-to-MCP adapter." @@ -27,7 +27,7 @@ id = "acp" transport = "stdio" executable = { path = "bin/cursor-mcp-adapter" } prefix = "cursor" -capabilities = ["tools"] +capabilities = ["tools", "resources"] [[skills]] id = "cursor-agent"