From ee4a51c1f1939cd3db4fcb0e1acfc6dbdd3944c4 Mon Sep 17 00:00:00 2001 From: showxu <10173746+showxu@users.noreply.github.com> Date: Sun, 27 Sep 2026 12:29:46 +0800 Subject: [PATCH 1/7] fix: preserve Cursor request and cleanup ownership --- Documentation/Reference/Interface.md | 6 +- Tests/Fixtures/vendor.py | 6 +- Tests/support.py | 2 + Tests/test_adapter.py | 14 ++ Tests/test_protocol_admission.py | 76 +++++++++ Tests/test_review_regressions.py | 186 ++++++++++++++++++++++ bin/cursor-mcp-adapter | 221 ++++++++++++++++++--------- bin/plugin_runtime.py | 32 +++- computer-mcp-plugin.toml | 2 +- 9 files changed, 463 insertions(+), 82 deletions(-) create mode 100644 Tests/test_protocol_admission.py create mode 100644 Tests/test_review_regressions.py diff --git a/Documentation/Reference/Interface.md b/Documentation/Reference/Interface.md index 0589bac..b8e2f70 100644 --- a/Documentation/Reference/Interface.md +++ b/Documentation/Reference/Interface.md @@ -39,15 +39,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. 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..53c6b6f 100644 --- a/Tests/test_adapter.py +++ b/Tests/test_adapter.py @@ -64,6 +64,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..daf54fb --- /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)],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 +217,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 +234,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 +250,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 +270,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 +301,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 +417,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} diff --git a/bin/plugin_runtime.py b/bin/plugin_runtime.py index 0c64892..f5477e4 100644 --- a/bin/plugin_runtime.py +++ b/bin/plugin_runtime.py @@ -47,7 +47,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=()): @@ -237,7 +246,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') @@ -473,7 +482,15 @@ def __init__(self, name, version, tools, shutdown): self.lock, self.write_lock = threading.Lock(), threading.Lock() self.initialized = False 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 +524,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,14 +542,18 @@ 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}}}) elif method == 'ping': diff --git a/computer-mcp-plugin.toml b/computer-mcp-plugin.toml index c9dba12..c595afc 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." From cf8a5fa8e58347049ad3a0640b277600b1a0258c Mon Sep 17 00:00:00 2001 From: showxu <10173746+showxu@users.noreply.github.com> Date: Sun, 27 Sep 2026 15:51:28 +0800 Subject: [PATCH 2/7] fix(permissions): declare execution risk on MCP tools --- Documentation/Reference/Interface.md | 13 +++++++++++++ Scripts/validate_host.py | 20 ++++++++++++++++++++ Tests/test_adapter.py | 8 ++++++++ Tests/test_protocol_admission.py | 2 +- bin/cursor-mcp-adapter | 24 ++++++++++++------------ bin/plugin_runtime.py | 10 +++++++++- 6 files changed, 63 insertions(+), 14 deletions(-) diff --git a/Documentation/Reference/Interface.md b/Documentation/Reference/Interface.md index b8e2f70..e8c3dba 100644 --- a/Documentation/Reference/Interface.md +++ b/Documentation/Reference/Interface.md @@ -8,6 +8,19 @@ 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. + ## Tools | Native MCP tool | Behavior | diff --git a/Scripts/validate_host.py b/Scripts/validate_host.py index ce4ad19..cdab113 100644 --- a/Scripts/validate_host.py +++ b/Scripts/validate_host.py @@ -216,6 +216,26 @@ 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, diff --git a/Tests/test_adapter.py b/Tests/test_adapter.py index 53c6b6f..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') diff --git a/Tests/test_protocol_admission.py b/Tests/test_protocol_admission.py index daf54fb..29c7dc7 100644 --- a/Tests/test_protocol_admission.py +++ b/Tests/test_protocol_admission.py @@ -21,7 +21,7 @@ def handler(_args,job): 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)],lambda:None) + 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)) diff --git a/bin/cursor-mcp-adapter b/bin/cursor-mcp-adapter index b42bdc7..986382e 100755 --- a/bin/cursor-mcp-adapter +++ b/bin/cursor-mcp-adapter @@ -471,18 +471,18 @@ class Cursor: 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'), + 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'), + 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'), + 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'), + Tool('cursor.acp.session.cancel','Request native ACP prompt cancellation without guessing a permission response.',schema({'session':SESSION},['session']),self.cancel,risk='destructive'), + Tool('cursor.acp.session.close','Close this Cursor session and retire its process group.',schema({'session':SESSION},['session']),self.close,risk='destructive'), + 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'), + 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'), + 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'), ] diff --git a/bin/plugin_runtime.py b/bin/plugin_runtime.py index f5477e4..af17d21 100644 --- a/bin/plugin_runtime.py +++ b/bin/plugin_runtime.py @@ -468,8 +468,16 @@ class Tool: description: str input_schema: dict handler: object + risk: str + 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') def definition(self): - return {'name':self.name,'description':self.description,'inputSchema':self.input_schema} + read_only = self.risk == 'read-only' + return {'name':self.name,'description':self.description,'inputSchema':self.input_schema, + '_meta':{'io.github.computer-mcp/risk':self.risk}, + 'annotations':{'readOnlyHint':read_only,'destructiveHint':self.risk in {'destructive','full-shell'}, + 'idempotentHint':read_only,'openWorldHint':not read_only}} class MCPServer: From 23c9f0bf16029ebb7fa513b603c67f29bb2f9ed7 Mon Sep 17 00:00:00 2001 From: showxu <10173746+showxu@users.noreply.github.com> Date: Sun, 27 Sep 2026 17:21:58 +0800 Subject: [PATCH 3/7] feat(mcp): report Cursor session ownership --- Documentation/Architecture/Package.md | 7 + Documentation/Reference/Installation.md | 8 + Documentation/Reference/Interface.md | 35 +++ README.md | 6 + Scripts/validate_host.py | 26 ++- Tests/test_work_resources.py | 281 ++++++++++++++++++++++++ bin/cursor-mcp-adapter | 42 +++- bin/plugin_runtime.py | 142 +++++++++--- computer-mcp-plugin.toml | 2 +- 9 files changed, 508 insertions(+), 41 deletions(-) create mode 100644 Tests/test_work_resources.py 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..707c1e4 100644 --- a/Documentation/Reference/Installation.md +++ b/Documentation/Reference/Installation.md @@ -39,6 +39,14 @@ python3 Scripts/validate_host.py \ 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. +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: ```sh diff --git a/Documentation/Reference/Interface.md b/Documentation/Reference/Interface.md index e8c3dba..828f9d6 100644 --- a/Documentation/Reference/Interface.md +++ b/Documentation/Reference/Interface.md @@ -21,6 +21,41 @@ 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 | 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 cdab113..424bcfd 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, require_work_ownership=False): host = host.resolve(strict=True) archive = archive.resolve(strict=True) output.mkdir(parents=True, exist_ok=False) @@ -166,6 +166,18 @@ 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('io.github.computer-mcp/work' not in tool.get('_meta',{}) for tool in catalog), + 'Gateway exports advertise a downstream-only work resource') + 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 +199,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')) @@ -252,8 +272,10 @@ def main(): parser.add_argument('--host', required=True, type=Path) 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.require_work_ownership) if __name__ == '__main__': diff --git a/Tests/test_work_resources.py b/Tests/test_work_resources.py new file mode 100644 index 0000000..ade16ff --- /dev/null +++ b/Tests/test_work_resources.py @@ -0,0 +1,281 @@ +"""Connection-owned Cursor sessions outlive their MCP creation responses.""" +import json +import os +import tempfile +import threading +import time +import unittest +import uuid +from unittest.mock import Mock, patch +from support import Client +from test_completion import adapter +from plugin_runtime import Failure, Job, MCPServer, Process, WORK_INVOCATION, WORK_METADATA, WORK_URI + + +class WorkResourceTests(unittest.TestCase): + def call(self, client, name, arguments=None, origin=None): + params = {'name':name,'arguments':arguments or {}} + if origin is not None: params['_meta'] = {WORK_INVOCATION:origin} + return client.request('tools/call',params) + + def snapshot(self, client): + response = client.request('resources/read',{'uri':WORK_URI}) + self.assertNotIn('error',response,response) + contents = response['result']['contents'] + self.assertEqual(len(contents),1) + self.assertEqual(contents[0]['uri'],WORK_URI) + self.assertEqual(contents[0]['mimeType'],'application/json') + return json.loads(contents[0]['text']) + + def test_native_session_report_spans_background_prompt_and_interactive_work(self): + with tempfile.TemporaryDirectory() as root: + client = Client(root) + try: + definitions = client.request('tools/list')['result']['tools'] + self.assertEqual(len(definitions),12) + for definition in definitions: + self.assertEqual(definition['_meta'][WORK_METADATA],{'format_version':1,'uri':WORK_URI}) + self.assertIn('io.github.computer-mcp/risk',definition['_meta']) + self.assertEqual(client.request('resources/list')['result']['resources'][0]['uri'],WORK_URI) + initial = self.snapshot(client) + self.assertEqual(initial['resources'],[]) + origin = str(uuid.uuid4()) + opened = self.call(client,'cursor.acp.session.open',{'permission_policy':'manual'},origin) + self.assertFalse(opened['result']['isError'],opened) + session = Client.value(opened)['session'] + report = self.snapshot(client) + self.assertEqual(report['instance_id'],initial['instance_id']) + self.assertGreater(report['revision'],initial['revision']) + self.assertEqual(report['resources'],[{'kind':'cursor.session','id':session,'acquired_by':origin,'state':'active'}]) + self.assertEqual(self.snapshot(client),report) + started = self.call(client,'cursor.acp.session.prompt.start',{'session':session,'prompt':'permission'},str(uuid.uuid4())) + prompt = Client.value(started)['prompt_id'] + for _ in range(100): + pending = Client.value(client.call('cursor.acp.requests.list',{'session':session}))['requests'] + if pending: break + time.sleep(.02) + else: self.fail('Native permission request did not become pending') + self.assertEqual(self.snapshot(client),report) + response = client.call('cursor.acp.requests.respond',{'session':session,'request_id':pending[0]['request_id'],'response':{'outcome':{'outcome':'selected','optionId':'opaque-no'}}}) + self.assertFalse(response['result']['isError'],response) + for _ in range(100): + result = Client.value(client.call('cursor.acp.session.prompt.result',{'session':session,'prompt_id':prompt})) + if result.get('completed'): break + time.sleep(.02) + else: self.fail('Native prompt did not complete') + self.assertEqual(self.snapshot(client),report) + closed = client.call('cursor.acp.session.close',{'session':session}) + self.assertFalse(closed['result']['isError'],closed) + released = self.snapshot(client) + self.assertEqual(released['resources'],[]) + self.assertGreater(released['revision'],report['revision']) + finally: client.close() + + def test_unbound_live_sessions_fail_observation_without_breaking_ordinary_calls(self): + with tempfile.TemporaryDirectory() as root: + client = Client(root) + try: + self.assertEqual(self.snapshot(client)['resources'],[]) + opened = client.call('cursor.acp.session.open') + self.assertFalse(opened['result']['isError'],opened) + session = Client.value(opened)['session'] + missing = client.request('resources/read',{'uri':WORK_URI}) + self.assertIn('error',missing) + self.assertIn('no bound host acquisition',missing['error']['message']) + prompted = client.call('cursor.acp.session.prompt',{'session':session,'prompt':'hello'}) + self.assertFalse(prompted['result']['isError'],prompted) + self.assertFalse(client.call('cursor.acp.session.close',{'session':session})['result']['isError']) + self.assertEqual(self.snapshot(client)['resources'],[]) + finally: client.close() + + def test_concurrent_identical_open_calls_keep_distinct_acquisitions(self): + with tempfile.TemporaryDirectory() as root: + client = Client(root) + try: + requests = [] + for _ in range(2): + origin = str(uuid.uuid4()) + client.serial += 1 + requests.append((client.serial,origin)) + client.send({'jsonrpc':'2.0','id':client.serial,'method':'tools/call', + 'params':{'name':'cursor.acp.session.open','arguments':{}, + '_meta':{WORK_INVOCATION:origin}}}) + expected = {} + for request, origin in requests: + opened = client.wait(request) + self.assertFalse(opened['result']['isError'],opened) + expected[Client.value(opened)['session']] = origin + self.assertEqual(len(expected),2) + report = self.snapshot(client) + self.assertEqual({row['id']:row['acquired_by'] for row in report['resources']},expected) + for session in expected: + self.assertFalse(client.call('cursor.acp.session.close',{'session':session})['result']['isError']) + self.assertEqual(self.snapshot(client)['resources'],[]) + finally: client.close() + + def test_invalid_metadata_and_resource_uris_do_not_start_work(self): + with tempfile.TemporaryDirectory() as root: + client = Client(root) + try: + for value in (None,True,7,{},[], 'bad-reference'): + response = client.request('tools/call',{'name':'cursor.acp.session.open','_meta':{WORK_INVOCATION:value}}) + self.assertEqual(response['error']['code'],-32602) + forged = self.call(client,'cursor.acp.session.open',{'acquired_by':str(uuid.uuid4())},str(uuid.uuid4())) + self.assertTrue(forged['result']['isError']) + self.assertEqual(self.snapshot(client)['resources'],[]) + self.assertIn('error',client.request('resources/read',{'uri':'https://unrelated.invalid'})) + self.assertEqual(client.request('ping')['result'],{}) + finally: client.close() + + def test_unknown_cleanup_is_retained_through_close_and_shutdown(self): + origin = str(uuid.uuid4()) + cursor = adapter.Cursor('/unused') + session = adapter.Session('owned','/unused','manual',origin) + session.process = Mock() + session.process.close.side_effect = Failure('cleanup_unconfirmed','Fixture has no cleanup acknowledgement') + cursor.sessions[session.handle] = session + for close in (lambda:cursor.close({'session':'owned'},Job()),cursor.shutdown): + with self.assertRaises(Failure): close() + self.assertIn('owned',cursor.sessions) + self.assertEqual(cursor.work_resources(),[{'kind':'cursor.session','id':'owned','acquired_by':origin,'state':'uncertain'}]) + + def test_background_thread_must_drain_before_session_release(self): + session = adapter.Session('owned','/unused','manual',str(uuid.uuid4())) + session.process = Mock() + session.prompt_thread = Mock() + session.prompt_thread.is_alive.return_value = True + with self.assertRaises(Failure): session.close() + self.assertFalse(session.released()) + self.assertEqual(session.work_resource()['state'],'uncertain') + session.prompt_thread.is_alive.return_value = False + session.close() + self.assertTrue(session.released()) + self.assertIsNone(session.work_resource()) + + def test_self_closing_background_work_is_retained_until_its_thread_finishes(self): + entered, release = threading.Event(), threading.Event() + session = adapter.Session('owned','/unused','manual',str(uuid.uuid4())) + def close_then_finish(): + session.close() + entered.set() + release.wait(2) + session.prompt_thread = threading.Thread(target=close_then_finish) + session.prompt_thread.start() + try: + self.assertTrue(entered.wait(1)) + self.assertTrue(session.cleanup_confirmed) + self.assertFalse(session.released()) + self.assertIsNotNone(session.work_resource()) + finally: + release.set() + session.prompt_thread.join(2) + self.assertTrue(session.released()) + + def test_pending_startup_is_owned_before_native_execution(self): + origin = str(uuid.uuid4()) + cursor = adapter.Cursor('/unused') + entered, release = threading.Event(), threading.Event() + errors = [] + def verify(*_): + entered.set() + if not release.wait(2): raise AssertionError('Probe not released') + raise Failure('incompatible_vendor','No native session was launched') + def run(): + try: cursor.open({},Job(work_invocation=origin)) + except Exception as error: errors.append(error) + with patch.object(adapter,'resolve_executable',return_value='/unused'),patch.object(adapter,'verify_version',side_effect=verify),patch.object(adapter,'Process') as process: + thread = threading.Thread(target=run) + thread.start() + try: + self.assertTrue(entered.wait(1)) + rows = cursor.work_resources() + self.assertEqual(len(rows),1) + self.assertEqual(rows[0]['acquired_by'],origin) + self.assertEqual(rows[0]['state'],'active') + finally: + release.set() + thread.join(3) + self.assertFalse(thread.is_alive()) + process.assert_not_called() + self.assertEqual(len(errors),1) + self.assertEqual(cursor.work_resources(),[]) + + def test_version_probe_cleanup_uncertainty_retains_its_session(self): + origin = str(uuid.uuid4()) + cursor = adapter.Cursor('/unused') + probe = Mock() + expected = next(row['stdout'] for row in adapter.TREE['executable_checks'] if row['args']==['--version']) + probe.read.side_effect = [expected.encode(),None] + probe.finish.return_value = 0 + probe.close.side_effect = Failure('cleanup_unconfirmed','Probe supervisor supplied no acknowledgement') + with patch.object(adapter,'resolve_executable',return_value='/unused'),patch('plugin_runtime.Process',return_value=probe),patch.object(adapter,'Process') as native: + with self.assertRaises(Failure): cursor.open({},Job(work_invocation=origin)) + native.assert_not_called() + self.assertEqual(len(cursor.sessions),1) + self.assertEqual(cursor.work_resources()[0]['state'],'uncertain') + self.assertEqual(cursor.work_resources()[0]['acquired_by'],origin) + with self.assertRaises(Failure): cursor.shutdown() + self.assertEqual(len(cursor.sessions),1) + + def test_cleanup_io_failure_cannot_become_success_on_repeated_close(self): + process = Process.__new__(Process) + process.close_lock, process.writer_lock = threading.Lock(), threading.Lock() + process.closed, process.cleanup_failure = False, None + process.life, process.receipt = 100, 101 + process.child, process.output, process.stderr_thread = Mock(), Mock(), Mock() + process.stop_stderr = threading.Event() + def close(fd): + if fd==process.life: raise OSError('Fixture lifeline close failed') + with patch('plugin_runtime.os.close',side_effect=close),patch.object(process,'confirm_cleanup') as confirmed: + for _ in range(2): + with self.assertRaises(Failure) as error: process.close() + self.assertEqual(error.exception.code,'cleanup_unconfirmed') + confirmed.assert_not_called() + + def test_partial_process_initialization_retains_unknown_cleanup(self): + for confirmed in (True,False): + with self.subTest(confirmed=confirmed): + cursor = adapter.Cursor('/unused') + owners = [] + def close(process): + owners.append(process) + if not confirmed: raise Failure('cleanup_unconfirmed','Partial startup cleanup unavailable') + try: + with patch.object(adapter,'resolve_executable',return_value='/unused'),patch('plugin_runtime.subprocess.Popen',return_value=Mock()),patch('plugin_runtime.os.set_blocking',side_effect=OSError('Fixture pipe setup failed')),patch.object(Process,'close',autospec=True,side_effect=close): + with self.assertRaises(OSError if confirmed else Failure): + cursor.open({},Job(work_invocation=str(uuid.uuid4()))) + self.assertEqual(len(owners),1) + if confirmed: + self.assertEqual(cursor.work_resources(),[]) + else: + self.assertEqual(len(cursor.sessions),1) + self.assertEqual(cursor.work_resources()[0]['state'],'uncertain') + with self.assertRaises(Failure): cursor.shutdown() + finally: + # These descriptors belong to the inert, mocked process owner. + for owner in owners: + os.close(owner.life) + os.close(owner.receipt) + + def test_failed_snapshot_never_advances_revision_or_discards_valid_state(self): + origin = str(uuid.uuid4()) + row = {'kind':'session','id':1,'acquired_by':origin,'state':'active'} + rows = [row] + server = MCPServer('fixture','1',[],lambda:None,work=lambda:rows) + def snapshot(): return json.loads(server.work_resource()['contents'][0]['text']) + first = snapshot() + rows.append(dict(row)) + with self.assertRaises(Failure): snapshot() + rows.pop() + self.assertEqual(snapshot(),first) + rows.append(dict(row,id='1')) + second = snapshot() + self.assertGreater(second['revision'],first['revision']) + self.assertEqual(len(second['resources']),2) + rows.reverse() + self.assertEqual(snapshot(),second) + rows[0]['state']='uncertain' + self.assertGreater(snapshot()['revision'],second['revision']) + + +if __name__=='__main__': unittest.main() diff --git a/bin/cursor-mcp-adapter b/bin/cursor-mcp-adapter index 986382e..4bf0e5e 100755 --- a/bin/cursor-mcp-adapter +++ b/bin/cursor-mcp-adapter @@ -84,10 +84,12 @@ def checked_response(method, params, response): class Session: - def __init__(self, handle, executable, policy): + def __init__(self, handle, executable, policy, work_invocation=None): self.handle, self.executable, self.policy = handle, executable, policy + self.work_invocation = work_invocation self.vendor_id = None self.process = None + self.version_process = None self.lock = threading.RLock() self.operation = threading.Lock() self.closed = threading.Event() @@ -98,6 +100,7 @@ class Session: self.active_operation = None self.cleanup_confirmed = False self.startup_job = None + self.startup_cleanup_failure = None self.startup_thread = None self.startup_done = threading.Event() self.startup_done.set() @@ -113,6 +116,16 @@ class Session: def info(self): with self.lock: return {'session':self.handle,'session_id':self.vendor_id,'closed':self.closed.is_set(),'busy':self.busy,'pending_requests':len(self.pending),'permission_policy':self.policy} + def work_resource(self): + with self.lock: + if self.released(): return None + return {'kind':'cursor.session','id':self.handle,'acquired_by':self.work_invocation, + 'state':'uncertain' if self.closed.is_set() and not self.cleanup_confirmed else 'active'} + def released(self): + with self.lock: + return self.cleanup_confirmed and self.startup_done.is_set() and not (self.prompt_thread and self.prompt_thread.is_alive()) + def retain_version_process(self, process): + with self.lock: self.version_process = process def open(self, args, job): with self.lock: if self.closed.is_set(): @@ -122,7 +135,8 @@ class Session: self.startup_done.clear() try: deadline = time.monotonic()+args.get('timeout_seconds',45) - verify_version(self.executable,TREE,job,os.getcwd()) + verify_version(self.executable,TREE,job,os.getcwd(),self.retain_version_process) + with self.lock: self.version_process = None job.check() with self.lock: if self.closed.is_set(): @@ -150,6 +164,10 @@ class Session: if args.get('mode'): self.set_mode(args['mode'],deadline,job) return {**self.info(),'agent_capabilities':self.initialization.get('agentCapabilities',{}),'session_info':self.session_info} + except Failure as error: + if error.code=='cleanup_unconfirmed': + with self.lock: self.startup_cleanup_failure = error + raise finally: with self.lock: self.startup_job = None @@ -384,11 +402,16 @@ class Session: if startup_thread is not None and startup_thread != threading.get_ident(): if not self.startup_done.wait(10): raise Failure('cleanup_unconfirmed','Session startup has not confirmed cleanup') + if self.startup_cleanup_failure is not None: raise self.startup_cleanup_failure if self.process is not None: self.process.close() - with self.lock: self.cleanup_confirmed = self.startup_done.is_set() + if self.version_process is not None: + self.version_process.close() if self.prompt_thread and self.prompt_thread is not threading.current_thread(): self.prompt_thread.join(timeout=4) + if self.prompt_thread.is_alive(): + raise Failure('cleanup_unconfirmed','Background prompt has not confirmed completion') + with self.lock: self.cleanup_confirmed = self.startup_done.is_set() class Cursor: @@ -399,9 +422,9 @@ class Cursor: self.stopping = False def open(self,args,job): executable = resolve_executable(self.executable,'CURSOR_AGENT_EXECUTABLE','agent') - session = Session(uuid.uuid4().hex,executable,args.get('permission_policy','reject-once')) + session = Session(uuid.uuid4().hex,executable,args.get('permission_policy','reject-once'),job.work_invocation) with self.lock: - self.sessions = {k:v for k,v in self.sessions.items() if not v.cleanup_confirmed} + self.sessions = {k:v for k,v in self.sessions.items() if not v.released()} if self.stopping or len(self.sessions)>=8: raise Failure('capacity','Cursor session capacity reached or adapter is shutting down') self.sessions[session.handle] = session @@ -454,17 +477,22 @@ 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: @@ -492,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 af17d21..8a11e85 100644 --- a/bin/plugin_runtime.py +++ b/bin/plugin_runtime.py @@ -16,10 +16,14 @@ 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' +WORK_INVOCATION = 'io.github.computer-mcp/work-invocation' class Failure(Exception): @@ -116,8 +120,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): @@ -294,10 +299,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: @@ -355,30 +365,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): @@ -400,7 +413,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') @@ -408,6 +421,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) @@ -481,7 +495,7 @@ def definition(self): 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 @@ -489,6 +503,44 @@ 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): try: raw = encoded(message)+b'\n' @@ -563,7 +615,9 @@ def dispatch(self, message): 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: @@ -571,18 +625,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 c595afc..f7c427d 100644 --- a/computer-mcp-plugin.toml +++ b/computer-mcp-plugin.toml @@ -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" From 90fb11a82cb7d98121347b4ab7d03261d9765e1a Mon Sep 17 00:00:00 2001 From: showxu <10173746+showxu@users.noreply.github.com> Date: Sun, 27 Sep 2026 21:45:09 +0800 Subject: [PATCH 4/7] feat: declare ownership of handle continuations --- Documentation/Reference/Interface.md | 13 +++++++++++++ Scripts/validate_host.py | 5 +++-- Tests/test_work_resources.py | 14 +++++++++++++- bin/cursor-mcp-adapter | 18 +++++++++--------- bin/plugin_runtime.py | 19 ++++++++++++++++++- 5 files changed, 56 insertions(+), 13 deletions(-) diff --git a/Documentation/Reference/Interface.md b/Documentation/Reference/Interface.md index 828f9d6..8196e53 100644 --- a/Documentation/Reference/Interface.md +++ b/Documentation/Reference/Interface.md @@ -106,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/Scripts/validate_host.py b/Scripts/validate_host.py index 424bcfd..174f0b8 100644 --- a/Scripts/validate_host.py +++ b/Scripts/validate_host.py @@ -168,8 +168,9 @@ def validate(host, archive, output, require_work_ownership=False): 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('io.github.computer-mcp/work' not in tool.get('_meta',{}) for tool in catalog), - 'Gateway exports advertise a downstream-only work resource') + 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'] diff --git a/Tests/test_work_resources.py b/Tests/test_work_resources.py index ade16ff..02dbe20 100644 --- a/Tests/test_work_resources.py +++ b/Tests/test_work_resources.py @@ -9,7 +9,7 @@ from unittest.mock import Mock, patch from support import Client from test_completion import adapter -from plugin_runtime import Failure, Job, MCPServer, Process, WORK_INVOCATION, WORK_METADATA, WORK_URI +from plugin_runtime import CONTINUATION_METADATA, Failure, Job, MCPServer, Process, WORK_INVOCATION, WORK_METADATA, WORK_URI class WorkResourceTests(unittest.TestCase): @@ -36,6 +36,18 @@ def test_native_session_report_spans_background_prompt_and_interactive_work(self for definition in definitions: self.assertEqual(definition['_meta'][WORK_METADATA],{'format_version':1,'uri':WORK_URI}) self.assertIn('io.github.computer-mcp/risk',definition['_meta']) + continuations = { + 'cursor.acp.session.prompt','cursor.acp.session.prompt.start', + 'cursor.acp.session.prompt.result','cursor.acp.session.mode', + 'cursor.acp.session.cancel','cursor.acp.session.close', + 'cursor.acp.events.read','cursor.acp.requests.list','cursor.acp.requests.respond', + } + declared = {tool['name']:tool['_meta'][CONTINUATION_METADATA] + for tool in definitions if CONTINUATION_METADATA in tool['_meta']} + self.assertEqual(set(declared),continuations) + for selector in declared.values(): + self.assertEqual(selector,{'format_version':1,'selectors':[ + {'kind':'cursor.session','handles':{'id':'/session'}}]}) self.assertEqual(client.request('resources/list')['result']['resources'][0]['uri'],WORK_URI) initial = self.snapshot(client) self.assertEqual(initial['resources'],[]) diff --git a/bin/cursor-mcp-adapter b/bin/cursor-mcp-adapter index 4bf0e5e..493b5c2 100755 --- a/bin/cursor-mcp-adapter +++ b/bin/cursor-mcp-adapter @@ -501,16 +501,16 @@ class Cursor: 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,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'), - 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'), - 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'), + 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'), - Tool('cursor.acp.session.cancel','Request native ACP prompt cancellation without guessing a permission response.',schema({'session':SESSION},['session']),self.cancel,risk='destructive'), - Tool('cursor.acp.session.close','Close this Cursor session and retire its process group.',schema({'session':SESSION},['session']),self.close,risk='destructive'), - 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'), - 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'), - 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'), + 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')), ] diff --git a/bin/plugin_runtime.py b/bin/plugin_runtime.py index 8a11e85..33e056e 100644 --- a/bin/plugin_runtime.py +++ b/bin/plugin_runtime.py @@ -23,6 +23,7 @@ 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' @@ -483,13 +484,29 @@ class Tool: 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): 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':{'io.github.computer-mcp/risk':self.risk}, + '_meta':metadata, 'annotations':{'readOnlyHint':read_only,'destructiveHint':self.risk in {'destructive','full-shell'}, 'idempotentHint':read_only,'openWorldHint':not read_only}} From 8a5ee55384d15e59c4a712b9a1cd81a1b2a93bec Mon Sep 17 00:00:00 2001 From: showxu <10173746+showxu@users.noreply.github.com> Date: Mon, 28 Sep 2026 08:07:50 +0800 Subject: [PATCH 5/7] ci: notify the official catalog after plugin releases --- .github/RELEASE.md | 59 ++++++++++++++++++++++++++++ .github/workflows/notify-catalog.yml | 28 +++++++++++++ 2 files changed, 87 insertions(+) create mode 100644 .github/RELEASE.md create mode 100644 .github/workflows/notify-catalog.yml 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..1ed04ef --- /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@e2a150570bbec9d2e4297b050defe05185ea56b4 + with: + token: ${{ secrets.CATALOG_DISPATCH_TOKEN }} From 52c98112c2a0a0c0438ea98919c0e441650b30e3 Mon Sep 17 00:00:00 2001 From: showxu <10173746+showxu@users.noreply.github.com> Date: Mon, 28 Sep 2026 08:18:17 +0800 Subject: [PATCH 6/7] ci: adopt catalog notification cooldown handling --- .github/workflows/notify-catalog.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/notify-catalog.yml b/.github/workflows/notify-catalog.yml index 1ed04ef..60a5b17 100644 --- a/.github/workflows/notify-catalog.yml +++ b/.github/workflows/notify-catalog.yml @@ -23,6 +23,6 @@ jobs: steps: - name: Request complete catalog reconciliation id: catalog - uses: computer-mcp/computer-mcp.github.io/.github/actions/notify-catalog@e2a150570bbec9d2e4297b050defe05185ea56b4 + uses: computer-mcp/computer-mcp.github.io/.github/actions/notify-catalog@fc27dd0f370d028a3e5021e3585274891f696578 with: token: ${{ secrets.CATALOG_DISPATCH_TOKEN }} From c8a2a3be29c3628f0d1fcb6bf3c11bb05ad8a47e Mon Sep 17 00:00:00 2001 From: showxu <10173746+showxu@users.noreply.github.com> Date: Tue, 29 Sep 2026 13:00:05 +0800 Subject: [PATCH 7/7] test: validate explicitly selected host versions --- Documentation/Reference/Installation.md | 9 ++++++--- Scripts/validate_host.py | 10 +++++++--- 2 files changed, 13 insertions(+), 6 deletions(-) diff --git a/Documentation/Reference/Installation.md b/Documentation/Reference/Installation.md index 707c1e4..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,17 @@ 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 diff --git a/Scripts/validate_host.py b/Scripts/validate_host.py index 174f0b8..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, require_work_ownership=False): +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, require_work_ownership=False): 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()) @@ -260,6 +261,7 @@ def observe(): 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, @@ -271,12 +273,14 @@ def observe(): 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, args.require_work_ownership) + validate(args.host, args.archive, args.output, args.expected_host_version, args.require_work_ownership) if __name__ == '__main__':