From 10cbf5f216055b54b889050c81f93351e5c732d1 Mon Sep 17 00:00:00 2001 From: Tom Softreck Date: Sat, 19 Sep 2026 19:07:31 +0200 Subject: [PATCH 1/2] fix(mcp): share validated observation contract across interfaces --- project/ticket-078/README.md | 26 ++++ project/ticket-078/intent.json | 110 ++++++++++++++++ src/monag/cli.py | 18 ++- src/monag/dsl.py | 105 +++++++++------ src/monag/dsl_llm.py | 5 +- src/monag/mcp.py | 153 +++++++++------------- src/monag/nl_contract.py | 151 +++++++++++++++++++++ src/monag/panel.py | 29 ++++ tests/test_nl_contract.py | 233 +++++++++++++++++++++++++++++++++ 9 files changed, 700 insertions(+), 130 deletions(-) create mode 100644 project/ticket-078/README.md create mode 100644 project/ticket-078/intent.json create mode 100644 src/monag/nl_contract.py create mode 100644 tests/test_nl_contract.py diff --git a/project/ticket-078/README.md b/project/ticket-078/README.md new file mode 100644 index 0000000..e07b384 --- /dev/null +++ b/project/ticket-078/README.md @@ -0,0 +1,26 @@ +# Ticket 078: MCP observation contract + +- **ID**: ticket-078 +- **Owner**: codex-monag-mcp-20260919 +- **Status**: IN_PROGRESS +- **Workflow state**: EDIT +- **Created**: 2026-09-19 +- **Planfile**: PLF-028 +- **Issue**: https://github.com/semcod/monag/issues/95 + +## Goal and scope + +SESSION_EXECUTION_AUTHORIZATION: user requested continuation after the installed MCP assessment, under the existing fix/publication authorization. Adopt the read-only observation interface from wellmanifest/nl-dsl-llm at 2040efe37b9eb898350f3fec0285a2d4e69d4e34. Implement execute_dsl, nl_ask, describe_grammar, schema://current, standard structured results and correct error propagation, with one validated observation executor including advisory queries. Scope is disjoint from ticket-076 and the concurrent governance adoption ticket-077. No claim of full mutation-interface conformance or client MCP registration. + +## Acceptance criteria + +- [x] AC-01: Standard MCP tools/resource and shared CLI/REST observation entrypoints preserve provenance, validate commands, reject malformed DSL and emit correct errors; legacy tool names remain usable. +- [ ] AC-02: Regression/full suite and governance pass; independent protected publication and installed-runtime canaries succeed. + +## Standard source + +The immutable source, supported schema profile and invocation examples are returned by describe_grammar and schema://current; source and tests are the material deliverables. + +## Validation + +Full application run: 361 tests and 30 subtests passed. Twelve focused conformance tests pass, including an additional CLI compatibility case. The pinned normative command/result schemas independently validate real success and error envelopes; the advertised observation schema is valid JSON Schema. Ruff on all changed Python files and governance pass. Protected review and installed-runtime validation remain external delivery receipts. diff --git a/project/ticket-078/intent.json b/project/ticket-078/intent.json new file mode 100644 index 0000000..729a5ea --- /dev/null +++ b/project/ticket-078/intent.json @@ -0,0 +1,110 @@ +{ + "schema": "new-project.intent/v3", + "ticket": "ticket-078", + "summary": "Conform MCP observation tools and shared CLI REST contract to NL-DSL-LLM", + "workstream": "application", + "classification": { + "kind": "BUG", + "priority": "P1", + "origin": "requested" + }, + "allowedPaths": [ + "project/ticket-078/**", + "src/monag/cli.py", + "src/monag/dsl.py", + "src/monag/dsl_llm.py", + "src/monag/mcp.py", + "src/monag/nl_contract.py", + "src/monag/panel.py", + "tests/test_nl_contract.py" + ], + "forbiddenPaths": [ + "project/ticket-*/user-*.md" + ], + "stacks": [], + "dependsOn": [], + "conflictsWith": [], + "integrationTicket": null, + "delivery": { + "acceptedBaseSha": "fb5259da3fb739cb0d138bd9008d422b59232381", + "targetBranch": "main", + "outcome": "Expose standard MCP observation tools/resource, strict validated DSL execution and consistent structured results through MCP, CLI and REST; publish and deploy verified merged package", + "nonGoals": [ + "No mutation command redesign or remote MCP registration", + "No edits to concurrent governance adoption or historical branches", + "No new runtime dependencies" + ], + "complexity": "L", + "estimatedMinutes": 75, + "budgets": { + "maxImplementationFiles": 9, + "maxAffectedComponents": 4, + "maxPublicInterfaceChanges": 3, + "maxRuntimeDependencies": 0 + }, + "architecture": { + "status": "accepted", + "decision": "Share a closed JSON-schema observation profile, its small validator, the existing DSL executor and standard result envelope across adapters; preserve legacy monag_* text alongside machine-readable results", + "components": [ + { + "name": "observation-contract", + "paths": [ + "src/monag/dsl.py", + "src/monag/dsl_llm.py", + "src/monag/nl_contract.py" + ] + }, + { + "name": "mcp", + "paths": [ + "src/monag/mcp.py" + ] + }, + { + "name": "cli-rest", + "paths": [ + "src/monag/cli.py", + "src/monag/panel.py" + ] + }, + { + "name": "conformance", + "paths": [ + "tests/test_nl_contract.py" + ] + } + ], + "responsibilityChanges": false, + "interfaceChanges": [ + "MCP standard tools and schema resource", + "CLI dsl and structured ask", + "REST versioned observation routes" + ], + "dataChanges": [], + "ui": { + "impact": "none", + "states": [], + "evidence": [] + }, + "rollback": "Revert the protected merge; switch installed launcher back to its recorded previous runtime" + }, + "runtimeDependencies": [], + "validation": [ + { + "criterion": "AC-01", + "commands": [ + "python -m pytest -q tests/test_nl_contract.py tests/test_mcp.py tests/test_dsl.py tests/test_dsl_llm.py" + ], + "evidence": "Strict invalid-input rejection, real stdio/HTTP/CLI parity, provider fallback and legacy tool compatibility" + }, + { + "criterion": "AC-02", + "commands": [ + "python -m pytest -q", + "./project/governance-check.sh --actor agent" + ], + "evidence": "Full tests and protected exact-head checks pass; installed merged MCP handshake/tool call verified" + } + ] + } +} diff --git a/src/monag/cli.py b/src/monag/cli.py index d43a029..305e3c5 100644 --- a/src/monag/cli.py +++ b/src/monag/cli.py @@ -277,6 +277,8 @@ def main(argv=None): query_parser = sub.add_parser('query', aliases=['ask'], help='execute a natural language query or OBSERVE DSL command') query_parser.add_argument('query', nargs='+', help='natural language query phrase or OBSERVE DSL command') + dsl_parser = sub.add_parser('dsl', help='execute one validated OBSERVE statement') + dsl_parser.add_argument('query', nargs='+', help='canonical OBSERVE DSL') sub.add_parser('catalog', help='read-only, local-only catalog of what each repository under --root ' 'declares itself to be (description, stack, entry points)') export_parser = sub.add_parser('export', help='read-only staging list of candidate work items ' @@ -594,7 +596,19 @@ def display_report(document): from . import mcp mcp.run_stdio_server(root, depth=args.depth) return 0 - if args.mode in ('query', 'ask'): + if args.mode in ('ask', 'dsl'): + from . import nl_contract + result = nl_contract.execute(' '.join(args.query), root, + direct=args.mode == 'dsl', depth=args.depth, + registry=registry) + if output_format == 'json': + print(json.dumps(result, ensure_ascii=True)) + elif result['success']: + display_report(result['meta']['markdown']) + else: + print(result['errors'][0]['message'], file=sys.stderr) + return 0 if result['success'] else 1 + if args.mode == 'query': from . import dsl_llm query_str = ' '.join(args.query) res = dsl_llm.execute(query_str, root, depth=args.depth, registry=registry) @@ -606,7 +620,7 @@ def display_report(document): else: print(f"Error: {res.get('error')}", file=sys.stderr) return 1 - return 0 + return 0 if res.get('status') == 'ok' else 1 if args.mode == 'catalog': from . import catalog if output_format != 'json' and sys.stderr.isatty(): diff --git a/src/monag/dsl.py b/src/monag/dsl.py index 4544699..1852495 100644 --- a/src/monag/dsl.py +++ b/src/monag/dsl.py @@ -20,14 +20,15 @@ SCHEMA = 'monag.dsl/v1' -VALID_DOMAINS = {'prs', 'pr', 'audit', 'status', 'agents', 'resume', 'usage', 'catalog'} +VALID_DOMAINS = {'prs', 'pr', 'audit', 'status', 'agents', 'resume', 'usage', 'catalog', 'advise'} class Query: """Structured representation of a Monag observation query.""" def __init__(self, target, hours=24.0, state='all', limit=20, unpushed_only=False, - worktrees_only=False, raw_input=None): + worktrees_only=False, raw_input=None, issue_limit=None, + radar=False, tier="all", emit_planfile=False): # Normalize target if target in {'pr', 'prs'}: self.target = 'prs' @@ -35,12 +36,18 @@ def __init__(self, target, hours=24.0, state='all', limit=20, unpushed_only=Fals self.target = 'status' else: self.target = target - self.hours = float(hours) if hours is not None else 24.0 - self.state = state if state in {'open', 'merged', 'all'} else 'all' - self.limit = int(limit) if limit is not None else 20 - self.unpushed_only = bool(unpushed_only) - self.worktrees_only = bool(worktrees_only) + self.hours = hours if hours is not None else 24.0 + self.state = state + self.limit = limit if limit is not None else 20 + self.unpushed_only = unpushed_only + self.worktrees_only = worktrees_only + self.issue_limit = issue_limit + self.radar = radar + self.tier = tier + self.emit_planfile = emit_planfile self.raw_input = raw_input or '' + from .nl_contract import validate_query + validate_query(self) def to_dsl(self): """Serialize query to canonical DSL string.""" @@ -55,6 +62,14 @@ def to_dsl(self): parts.append('UNPUSHED_ONLY') if self.worktrees_only: parts.append('WORKTREES_ONLY') + if self.issue_limit is not None: + parts.extend(['ISSUE_LIMIT', str(self.issue_limit)]) + if self.radar: + parts.append('RADAR') + if self.tier != 'all': + parts.extend(['TIER', self.tier]) + if self.emit_planfile: + parts.append('EMIT_PLANFILE') return ' '.join(parts) def to_dict(self): @@ -65,6 +80,10 @@ def to_dict(self): 'limit': self.limit, 'unpushed_only': self.unpushed_only, 'worktrees_only': self.worktrees_only, + 'issue_limit': self.issue_limit, + 'radar': self.radar, + 'tier': self.tier, + 'emit_planfile': self.emit_planfile, 'dsl': self.to_dsl(), 'raw_input': self.raw_input, } @@ -97,34 +116,28 @@ def parse_dsl(text): return None kwargs = {'target': target, 'raw_input': text} + converters = {'HOURS': float, 'LIMIT': int, 'ISSUE_LIMIT': int, + 'STATE': str.lower, 'TIER': str.lower} + flags = {'UNPUSHED_ONLY', 'WORKTREES_ONLY', 'RADAR', 'EMIT_PLANFILE'} + seen = set() i = 2 - while i < len(tokens): - key = tokens[i].upper() - if key == 'HOURS' and i + 1 < len(tokens): - try: - kwargs['hours'] = float(tokens[i + 1]) - except ValueError: - pass - i += 2 - elif key == 'STATE' and i + 1 < len(tokens): - kwargs['state'] = tokens[i + 1].lower() - i += 2 - elif key == 'LIMIT' and i + 1 < len(tokens): - try: - kwargs['limit'] = int(tokens[i + 1]) - except ValueError: - pass - i += 2 - elif key == 'UNPUSHED_ONLY': - kwargs['unpushed_only'] = True - i += 1 - elif key == 'WORKTREES_ONLY': - kwargs['worktrees_only'] = True - i += 1 - else: - i += 1 - - return Query(**kwargs) + try: + while i < len(tokens): + key = tokens[i].upper() + if key in seen: + return None + seen.add(key) + if key in converters and i + 1 < len(tokens): + kwargs[key.lower()] = converters[key](tokens[i + 1]) + i += 2 + elif key in flags: + kwargs[key.lower()] = True + i += 1 + else: + return None + return Query(**kwargs) + except (ValueError, TypeError, OverflowError): + return None def parse_natural_language(text): @@ -142,6 +155,9 @@ def parse_natural_language(text): elif re.search(r'(scalon|zmergowan|merged)', low): state = 'merged' + if re.search(r'(advise|advice|porad|zaleceni|rekomendac)', low): + return Query('advise', raw_input=raw) + # Domain 1: PRs & branches if re.search(r'(\bprs?\b|\bpr-[a-z0-9]+\b|pull\s*requests?|ga[łl][ęe]z|branch|\bmerg|scal|nieprzepchni)', low): return Query('prs', hours=hours if hours is not None else 24.0, @@ -182,10 +198,8 @@ def parse(text): if not text or not text.strip(): return None raw = text.strip() - if raw.upper().startswith('OBSERVE '): - dsl_query = parse_dsl(raw) - if dsl_query: - return dsl_query + if raw.upper().startswith('OBSERVE'): + return parse_dsl(raw) return parse_natural_language(raw) @@ -204,6 +218,8 @@ def execute(query_or_text, root, depth=2, pr_limit=200, issue_limit=200, registr else: query = query_or_text + from .nl_contract import validate_query + validate_query(query) target = query.target data = None doc = '' @@ -216,7 +232,7 @@ def execute(query_or_text, root, depth=2, pr_limit=200, issue_limit=200, registr elif target == 'audit': from . import audit - data = audit.scan(root, depth=depth, issue_limit=issue_limit, + data = audit.scan(root, depth=depth, issue_limit=query.issue_limit or issue_limit, recent_hours=query.hours if query.hours != 24.0 else None, worktrees_hours=query.hours, worktrees_only=query.worktrees_only) @@ -246,6 +262,17 @@ def execute(query_or_text, root, depth=2, pr_limit=200, issue_limit=200, registr data = catalog.scan(root, depth=depth) doc = catalog.markdown(data, limit=query.limit) + elif target == 'advise': + from . import advise + data = advise.advise(root, depth=depth, limit=query.limit, + radar=query.radar, tier=query.tier) + if query.emit_planfile: + import json + data = advise.export_planfile_tickets(data, tier=query.tier) + doc = json.dumps(data, ensure_ascii=False, indent=2) + else: + doc = advise.markdown(data) + return { 'schema': SCHEMA, 'status': 'ok', diff --git a/src/monag/dsl_llm.py b/src/monag/dsl_llm.py index 0402e73..5bd8083 100644 --- a/src/monag/dsl_llm.py +++ b/src/monag/dsl_llm.py @@ -24,7 +24,7 @@ Grammar (exactly one line, conforming or nothing): OBSERVE [HOURS ] [STATE ] [LIMIT ] [UNPUSHED_ONLY] [WORKTREES_ONLY] -domain := prs | audit | status | resume | usage | catalog +domain := prs | audit | status | resume | usage | catalog | advise Domains: prs = pull requests and branches, audit = planfile/GitHub coverage, status = agent processes and checkout snapshot, resume = worktree checkouts @@ -94,6 +94,9 @@ def resolve(text, command=None, timeout=None): if rule_query is not None: return rule_query, provenance('rule', rule_query.to_dsl(), raw) + if raw.upper().startswith('OBSERVE'): + return None, provenance('none', None, raw) + if command is None: command = os.environ.get('MONAG_LLM_COMMAND', '') if timeout is None: diff --git a/src/monag/mcp.py b/src/monag/mcp.py index 8dadcef..60a375c 100644 --- a/src/monag/mcp.py +++ b/src/monag/mcp.py @@ -19,7 +19,7 @@ from pathlib import Path import sys -from . import dsl +from . import dsl, nl_contract MCP_PROTOCOL_VERSION = '2024-11-05' SERVER_NAME = 'monag-mcp' @@ -120,7 +120,24 @@ }, ] +TOOLS.extend([ + {'name': 'execute_dsl', 'description': 'Execute one strictly validated MONAG OBSERVE statement.', + 'inputSchema': {'type': 'object', 'required': ['dsl_command'], + 'properties': {'dsl_command': {'type': 'string'}}}}, + {'name': 'nl_ask', 'description': 'Resolve Polish or English observation queries; LLM translation is optional.', + 'inputSchema': {'type': 'object', 'required': ['query'], 'properties': { + 'query': {'type': 'string'}, 'allow_llm_fallback': {'type': 'boolean', 'default': True}, + 'locale': {'type': 'string', 'enum': ['pl', 'en']}}}}, + {'name': 'describe_grammar', 'description': 'Describe the pinned observation schema, domains and interfaces.', + 'inputSchema': {'type': 'object', 'properties': {}}}, +]) +for tool in TOOLS: + tool['inputSchema']['additionalProperties'] = False + tool['annotations'] = {'readOnlyHint': True, 'destructiveHint': False} + + RESOURCES = [ + {'uri': 'schema://current', 'name': 'Observation command JSON Schema', 'mimeType': 'application/json'}, {'uri': 'monag://snapshot', 'name': 'Workspace Snapshot', 'mimeType': 'text/markdown'}, {'uri': 'monag://prs', 'name': 'Pull Requests & Unpushed Branches', 'mimeType': 'text/markdown'}, {'uri': 'monag://audit', 'name': 'Planfile / GitHub Coverage Audit', 'mimeType': 'text/markdown'}, @@ -130,78 +147,51 @@ ] -def handle_tool_call(name, arguments, root, depth=2): - """Execute tool and return MCP formatted content array.""" - arguments = arguments or {} - - if name == 'monag_prs': - query = dsl.Query( - target='prs', - hours=float(arguments.get('hours', 24.0)), - state=arguments.get('state', 'all'), - unpushed_only=bool(arguments.get('unpushed_only', False)), - ) - res = dsl.execute(query, root, depth=depth) - return [{'type': 'text', 'text': res['markdown']}] - - elif name == 'monag_audit': - query = dsl.Query( - target='audit', - issue_limit=int(arguments.get('issue_limit', 200)), - worktrees_only=bool(arguments.get('worktrees_only', False)), - hours=float(arguments.get('hours', 24.0)), - ) - res = dsl.execute(query, root, depth=depth) - return [{'type': 'text', 'text': res['markdown']}] - - elif name == 'monag_status': - query = dsl.Query( - target='status', - hours=float(arguments.get('hours', 24.0)), - limit=int(arguments.get('limit', 20)), - ) - res = dsl.execute(query, root, depth=depth) - return [{'type': 'text', 'text': res['markdown']}] - - elif name == 'monag_resume': - query = dsl.Query(target='resume') - res = dsl.execute(query, root, depth=depth) - return [{'type': 'text', 'text': res['markdown']}] - - elif name == 'monag_usage': - query = dsl.Query(target='usage') - res = dsl.execute(query, root, depth=depth) - return [{'type': 'text', 'text': res['markdown']}] - - elif name == 'monag_catalog': - query = dsl.Query(target='catalog') - res = dsl.execute(query, root, depth=depth) - return [{'type': 'text', 'text': res['markdown']}] - - elif name == 'monag_query': - from . import dsl_llm - q_str = arguments.get('query', '') - res = dsl_llm.execute(q_str, root, depth=depth) - if res.get('status') == 'error': - return [{'type': 'text', 'text': f"Error: {res.get('error')}"}] - return [{'type': 'text', 'text': res['markdown']}] - - elif name == 'monag_advise': - from . import advise - data = advise.advise(root, depth=depth, limit=int(arguments.get('limit', 10)), - radar=bool(arguments.get('radar', False)), - tier=arguments.get('tier')) - if arguments.get('emit_planfile'): - payload = advise.export_planfile_tickets(data, tier=arguments.get('tier')) - return [{'type': 'text', 'text': json.dumps(payload, ensure_ascii=False, indent=2)}] - return [{'type': 'text', 'text': advise.markdown(data)}] - +def execute_tool(name, arguments, root, depth=2): + """Keep legacy text while returning the same machine result as CLI/REST.""" + try: + tool = next((item for item in TOOLS if item['name'] == name), None) + if tool is None: + raise ValueError(f'Unknown tool: {name}') + nl_contract.validate(arguments, tool['inputSchema']) + if name == 'describe_grammar': + result = nl_contract.grammar() + return {'content': [{'type': 'text', 'text': json.dumps(result)}], + 'structuredContent': result, 'isError': False} + if name == 'execute_dsl': + result = nl_contract.execute(arguments['dsl_command'], root, direct=True, depth=depth) + elif name in {'nl_ask', 'monag_query'}: + result = nl_contract.execute(arguments['query'], root, depth=depth, + allow_llm_fallback=arguments.get('allow_llm_fallback', True), + locale=arguments.get('locale')) + else: + target = name.removeprefix('monag_') + values = dict(arguments) + if target == 'advise': + values.setdefault('limit', 10) + result = nl_contract.execute({'entity': target, 'operation': 'query', + 'arguments': values}, root, direct=True, depth=depth) + except (ValueError, TypeError, OverflowError) as error: + result = nl_contract.envelope(success=False, status='VALIDATION_ERROR', error=error) + if isinstance(name, str) and name in {'execute_dsl', 'nl_ask'}: + text = json.dumps(result, ensure_ascii=False) + elif result['success']: + text = result['meta']['markdown'] else: - raise ValueError(f'Unknown tool: {name}') + text = 'Error: ' + result['errors'][0]['message'] + return {'content': [{'type': 'text', 'text': text}], + 'structuredContent': result, 'isError': not result['success']} + + +def handle_tool_call(name, arguments, root, depth=2): + """Compatibility helper for callers expecting just the text content array.""" + return execute_tool(name, arguments, root, depth)['content'] def read_resource(uri, root, depth=2): """Read resource content by URI and return markdown text.""" + if uri == 'schema://current': + return json.dumps(nl_contract.grammar()['commandSchema']) target = uri.replace('monag://', '').strip('/') if target == 'snapshot': res = dsl.execute('OBSERVE status', root, depth=depth) @@ -214,9 +204,7 @@ def read_resource(uri, root, depth=2): elif target == 'usage': res = dsl.execute('OBSERVE usage', root, depth=depth) elif target == 'advise': - from . import advise - data = advise.advise(root, depth=depth, limit=10) - return advise.markdown(data) + res = dsl.execute('OBSERVE advise LIMIT 10', root, depth=depth) else: raise ValueError(f'Unknown resource: {uri}') return res['markdown'] @@ -234,6 +222,9 @@ def process_message(msg, root, depth=2): method = msg.get('method') params = msg.get('params', {}) + if not isinstance(params, dict): + return {'jsonrpc': '2.0', 'id': msg_id, 'error': {'code': -32602, 'message': 'params must be an object'}} + if not method: return {'jsonrpc': '2.0', 'id': msg_id, 'error': {'code': -32600, 'message': 'Method required'}} @@ -270,22 +261,8 @@ def process_message(msg, root, depth=2): if method == 'tools/call': tool_name = params.get('name') arguments = params.get('arguments', {}) - try: - content = handle_tool_call(tool_name, arguments, root, depth=depth) - return { - 'jsonrpc': '2.0', - 'id': msg_id, - 'result': {'content': content, 'isError': False}, - } - except Exception as error: - return { - 'jsonrpc': '2.0', - 'id': msg_id, - 'result': { - 'content': [{'type': 'text', 'text': f'Tool execution error: {error}'}], - 'isError': True, - }, - } + return {'jsonrpc': '2.0', 'id': msg_id, + 'result': execute_tool(tool_name, arguments, root, depth)} if method == 'resources/list': return { @@ -302,7 +279,7 @@ def process_message(msg, root, depth=2): 'jsonrpc': '2.0', 'id': msg_id, 'result': { - 'contents': [{'uri': uri, 'mimeType': 'text/markdown', 'text': text}], + 'contents': [{'uri': uri, 'mimeType': 'application/json' if uri == 'schema://current' else 'text/markdown', 'text': text}], }, } except Exception as error: diff --git a/src/monag/nl_contract.py b/src/monag/nl_contract.py new file mode 100644 index 0000000..7606ba9 --- /dev/null +++ b/src/monag/nl_contract.py @@ -0,0 +1,151 @@ +"""Closed observation profile of wellmanifest/nl-dsl-llm (2040efe37b9e). + +All executable inputs validate against COMMAND_SCHEMA before DSL dispatch. +The validator implements only the keywords used by this fixed, local schema; +it is not a general-purpose JSON Schema library or a remote schema resolver. +Mutating MONAG commands and MCP client registration are outside this profile. +""" +import copy +import math +import time +import uuid + +STANDARD_REVISION = '2040efe37b9eb898350f3fec0285a2d4e69d4e34' +DOMAINS = ['prs', 'audit', 'status', 'resume', 'usage', 'catalog', 'advise'] +PARAMETERS = { + 'hours': {'type': 'number', 'minimum': 0}, + 'state': {'type': 'string', 'enum': ['all', 'open', 'merged']}, + 'limit': {'type': 'integer', 'minimum': 1}, + 'unpushed_only': {'type': 'boolean'}, + 'worktrees_only': {'type': 'boolean'}, + 'issue_limit': {'type': 'integer', 'minimum': 1}, + 'radar': {'type': 'boolean'}, + 'tier': {'type': 'string', 'enum': ['all', 'floor', 'mission', 'hygiene', 'backlog']}, + 'emit_planfile': {'type': 'boolean'}, +} +COMMAND_SCHEMA = { + '$schema': 'https://json-schema.org/draft/2020-12/schema', + 'title': 'MONAG read-only observation command', + 'type': 'object', 'additionalProperties': False, + 'required': ['entity', 'operation'], + 'properties': { + 'entity': {'type': 'string', 'enum': DOMAINS}, + 'operation': {'type': 'string', 'enum': ['query']}, + 'arguments': {'type': 'object', 'additionalProperties': False, + 'properties': PARAMETERS}, + }, +} + + +def validate(value, schema, field='input'): + """Validate this profile's closed object/scalar schemas without coercion.""" + kind = schema['type'] + matches = { + 'object': lambda: isinstance(value, dict), + 'string': lambda: isinstance(value, str), + 'boolean': lambda: type(value) is bool, + 'integer': lambda: type(value) is int, + 'number': lambda: type(value) in (int, float) and math.isfinite(value), + } + if kind not in matches or not matches[kind](): + raise ValueError(f'{field} must be {kind}') + if 'enum' in schema and value not in schema['enum']: + raise ValueError(f'{field} must be one of {schema["enum"]}') + if 'minimum' in schema and value < schema['minimum']: + raise ValueError(f'{field} must be >= {schema["minimum"]}') + if kind == 'object': + properties = schema.get('properties', {}) + if schema.get('additionalProperties') is False and set(value) - set(properties): + raise ValueError(f'{field} has unsupported fields') + if set(schema.get('required', [])) - set(value): + raise ValueError(f'{field} is missing required fields') + for name, item in value.items(): + if name in properties: + validate(item, properties[name], f'{field}.{name}') + + +def command(query): + arguments = {name: getattr(query, name) for name in PARAMETERS} + if arguments['issue_limit'] is None: + del arguments['issue_limit'] + return {'entity': query.target, 'operation': 'query', 'arguments': arguments} + + +def validate_query(query): + validate(command(query), COMMAND_SCHEMA) + + +def grammar(): + """Return a fresh JSON schema and discoverable invocation contract.""" + return { + 'standard': 'wellmanifest/nl-dsl-llm', 'sourceRevision': STANDARD_REVISION, + 'profile': 'read-only observations; mutation commands are not covered', + 'commandSchema': copy.deepcopy(COMMAND_SCHEMA), + 'dsl': 'OBSERVE [HOURS n] [STATE all|open|merged] [LIMIT n] ' + '[ISSUE_LIMIT n] [UNPUSHED_ONLY] [WORKTREES_ONLY] [RADAR] ' + '[TIER all|floor|mission|hygiene|backlog] [EMIT_PLANFILE]', + 'domains': list(DOMAINS), + 'examples': ['OBSERVE catalog', 'OBSERVE prs STATE open', 'OBSERVE advise LIMIT 10'], + 'interfaces': {'mcp': ['nl_ask', 'execute_dsl', 'describe_grammar'], + 'cli': ['monag --json ask "pokaż katalog"', + 'monag --json dsl "OBSERVE catalog"'], + 'rest': ['/api/v1/query', '/api/v1/dsl', '/api/v1/schema']}, + } + + +def envelope(*, success, status, data=None, error=None, source='direct_dsl', + canonical=None, provenance=None, markdown='', started=None): + return { + 'success': success, 'status': status, 'data': data, + 'errors': [] if success else [{'code': status, 'message': str(error)}], + 'meta': {'executionTimeMs': max(0, (time.perf_counter() - started) * 1000) if started else 0, + 'sourceLayer': source, 'canonicalDsl': canonical, + 'traceId': str(uuid.uuid4()), 'provenance': provenance, + 'markdown': markdown, 'standardRevision': STANDARD_REVISION}, + } + + +def execute(value, root, *, direct=False, allow_llm_fallback=True, locale=None, **options): + """Resolve and validate once, then dispatch exclusively through the DSL.""" + from . import dsl, dsl_llm + started = time.perf_counter() + source = 'direct_dsl' if direct else 'nl_fast_path' + canonical = None + prov = None + try: + if type(allow_llm_fallback) is not bool: + raise ValueError('allow_llm_fallback must be boolean') + if locale not in (None, 'pl', 'en'): + raise ValueError('locale must be pl or en') + if direct and isinstance(value, dict): + validate(value, COMMAND_SCHEMA) + query = dsl.Query(value['entity'], **value.get('arguments', {})) + prov = dsl_llm.provenance('rule', query.to_dsl()) + else: + if not isinstance(value, str) or not value.strip(): + raise ValueError('A nonempty query or DSL statement is required') + canonical_input = value.lstrip().upper().startswith('OBSERVE') + if direct or canonical_input: + source = 'direct_dsl' + query = dsl.parse_dsl(value) + prov = dsl_llm.provenance('rule', query.to_dsl() if query else None, value) + else: + query, prov = dsl_llm.resolve(value, command=None if allow_llm_fallback else '') + source = 'llm_fallback' if prov['engine'] == 'llm' else 'nl_fast_path' + if query is None: + raise ValueError('Unrecognized or invalid observation command') + validate_query(query) + canonical = query.to_dsl() + except (ValueError, TypeError, OverflowError) as error: + return envelope(success=False, status='VALIDATION_ERROR', error=error, + source=source, canonical=canonical, provenance=prov, started=started) + try: + result = dsl.execute(query, root, **options) + if result.get('status') == 'error': + raise ValueError(result.get('error', 'Observation failed')) + return envelope(success=True, status='OK', data=result.get('data'), + source=source, canonical=canonical, provenance=prov, + markdown=result.get('markdown', ''), started=started) + except Exception as error: + return envelope(success=False, status='EXECUTION_ERROR', error=error, + source=source, canonical=canonical, provenance=prov, started=started) diff --git a/src/monag/panel.py b/src/monag/panel.py index 5443465..ed858df 100644 --- a/src/monag/panel.py +++ b/src/monag/panel.py @@ -383,6 +383,11 @@ def do_GET(self): self._send(json.dumps({'status': 'ok', 'updated': updated, 'config': cfg}, ensure_ascii=False).encode(), 'application/json; charset=utf-8') return + if parsed.path == '/api/v1/schema': + from . import nl_contract + self._send(json.dumps(nl_contract.grammar()).encode(), + 'application/json; charset=utf-8') + return if parsed.path in {'/api/query', '/api/query.json'}: params = parse_qs(parsed.query) q = (params.get('q') or params.get('query') or params.get('nl') or [''])[0] @@ -442,6 +447,30 @@ def do_POST(self): self._send(json.dumps({'status': 'ok', 'updated': updated, 'config': cfg}, ensure_ascii=False).encode(), 'application/json; charset=utf-8') return + if parsed.path in {'/api/v1/query', '/api/v1/dsl'}: + from . import nl_contract + try: + length = int(self.headers.get('Content-Length', 0)) + if not 0 < length <= 65536: + raise ValueError('JSON request size must be 1..65536 bytes') + body = json.loads(self.rfile.read(length)) + if not isinstance(body, dict): + raise ValueError('JSON body must be an object') + direct = parsed.path.endswith('/dsl') + allowed = {'dsl_command'} if direct else {'query', 'locale', 'allow_llm_fallback'} + if set(body) - allowed: + raise ValueError('Unsupported request fields') + value = body.get('dsl_command' if direct else 'query') + result = nl_contract.execute(value, state.root, direct=direct, + allow_llm_fallback=body.get('allow_llm_fallback', True), + locale=body.get('locale'), depth=state.depth, + registry=state.registry) + except (ValueError, TypeError) as error: + result = nl_contract.envelope(success=False, status='VALIDATION_ERROR', error=error) + code = 200 if result['success'] else (400 if result['status'] == 'VALIDATION_ERROR' else 500) + self._send(json.dumps(result, ensure_ascii=True).encode(), + 'application/json; charset=utf-8', status=code) + return if parsed.path in {'/api/query', '/api/query.json'}: try: length = int(self.headers.get('Content-Length', 0)) diff --git a/tests/test_nl_contract.py b/tests/test_nl_contract.py new file mode 100644 index 0000000..e47c322 --- /dev/null +++ b/tests/test_nl_contract.py @@ -0,0 +1,233 @@ +"""Behavioral conformance of the shared read-only NL/DSL/MCP interface.""" +import http.client +import json +import os +from pathlib import Path +import subprocess +import sys +import tempfile +import threading +import unittest +from unittest.mock import patch + +from monag import dsl, dsl_llm, mcp, nl_contract, panel + + +class ContractTests(unittest.TestCase): + def setUp(self): + self.temporary = tempfile.TemporaryDirectory() + self.addCleanup(self.temporary.cleanup) + self.root = Path(self.temporary.name) + self.addCleanup(patch.stopall) + patch.dict(os.environ, {'MONAG_LLM_COMMAND': ''}).start() + + def tool(self, name, arguments): + return mcp.process_message({'jsonrpc': '2.0', 'id': 1, 'method': 'tools/call', + 'params': {'name': name, 'arguments': arguments}}, self.root)['result'] + + def assert_envelope(self, result, success=True): + self.assertEqual(set(result), {'success', 'status', 'data', 'errors', 'meta'}) + self.assertEqual(result['success'], success) + self.assertEqual(bool(result['errors']), not success) + self.assertGreaterEqual(result['meta']['executionTimeMs'], 0) + self.assertIn(result['meta']['sourceLayer'], {'direct_dsl', 'nl_fast_path', 'llm_fallback'}) + self.assertTrue(result['meta']['traceId']) + + def test_tools_and_resource_publish_same_closed_schema(self): + self.assertTrue({'execute_dsl', 'nl_ask', 'describe_grammar'} <= {t['name'] for t in mcp.TOOLS}) + grammar = self.tool('describe_grammar', {})['structuredContent'] + schema = json.loads(mcp.read_resource('schema://current', self.root)) + self.assertEqual(schema, grammar['commandSchema']) + self.assertEqual(grammar['sourceRevision'], nl_contract.STANDARD_REVISION) + grammar['commandSchema']['properties'].clear() + self.assertIn('entity', nl_contract.grammar()['commandSchema']['properties']) + + def test_direct_and_multilingual_fast_paths_never_invoke_provider(self): + with patch.object(dsl_llm, 'run_provider', side_effect=AssertionError('unexpected LLM')): + for tool, key, value, layer in [ + ('execute_dsl', 'dsl_command', 'OBSERVE catalog', 'direct_dsl'), + ('nl_ask', 'query', 'pokaż katalog', 'nl_fast_path'), + ('nl_ask', 'query', 'show catalog', 'nl_fast_path'), + ]: + with self.subTest(tool=tool, value=value): + result = self.tool(tool, {key: value}) + self.assertFalse(result['isError']) + envelope = result['structuredContent'] + self.assert_envelope(envelope) + self.assertEqual(json.loads(result['content'][0]['text']), envelope) + self.assertEqual(envelope['meta']['sourceLayer'], layer) + self.assertEqual(envelope['meta']['canonicalDsl'], 'OBSERVE catalog') + self.assertEqual(envelope['meta']['provenance']['engine'], 'rule') + + def test_invalid_dsl_is_rejected_before_provider_or_scanner(self): + invalid = ['OBSERVE catalog SURPRISE', 'OBSERVE catalog LIMIT', + 'OBSERVE catalog LIMIT xyz', 'OBSERVE catalog LIMIT -1', + 'OBSERVE catalog STATE unexpected', 'OBSERVE catalog HOURS NaN', + 'OBSERVE catalog HOURS inf', 'OBSERVE catalog HOURS -1', + 'OBSERVE catalog LIMIT 1 LIMIT 2', + 'OBSERVE catalog; OBSERVE prs', 'OBSERVE catalog\nOBSERVE prs', + 'OBSERVE unknown', 'show catalog'] + with patch.object(dsl_llm, 'run_provider', side_effect=AssertionError('unexpected LLM')): + with patch.object(dsl, 'execute', side_effect=AssertionError('unexpected scan')): + for value in invalid: + with self.subTest(value=value): + result = self.tool('execute_dsl', {'dsl_command': value}) + self.assertTrue(result['isError']) + self.assert_envelope(result['structuredContent'], False) + self.assertIsNone(dsl.parse('OBSERVE catalog LIMIT garbage')) + self.assertIsNone(dsl_llm.resolve('OBSERVE catalog LIMIT garbage', command='fixture')[0]) + + def test_tool_argument_validation(self): + for name, args in [('nl_ask', {}), ('nl_ask', []), ('nl_ask', None), + ('nl_ask', {'query': 5}), ('nl_ask', {'query': 'catalog', 'extra': True}), + ('nl_ask', {'query': 'catalog', 'allow_llm_fallback': 'false'}), + ('monag_status', {'limit': True}), ('monag_audit', {'issue_limit': 0}), + ('monag_advise', {'tier': 'invalid'}), (['invalid'], {})]: + with self.subTest(name=name, args=args): + result = self.tool(name, args) + self.assertTrue(result['isError']) + self.assert_envelope(result['structuredContent'], False) + + def test_unknown_query_error_and_provenance_survive_legacy_adapter(self): + result = self.tool('monag_query', {'query': 'qqzz_unknown_123'}) + self.assertTrue(result['isError']) + self.assertTrue(result['content'][0]['text'].startswith('Error:')) + self.assert_envelope(result['structuredContent'], False) + success = self.tool('monag_query', {'query': 'OBSERVE catalog'}) + self.assertIn('# MONAG', success['content'][0]['text']) + self.assertEqual(success['structuredContent']['meta']['provenance']['engine'], 'rule') + + def test_provider_translation_validates_and_can_be_disabled_per_call(self): + with patch.dict(os.environ, {'MONAG_LLM_COMMAND': 'fixture-provider'}): + with patch.object(dsl_llm, 'run_provider', return_value='OBSERVE catalog') as provider: + disabled = self.tool('nl_ask', {'query': 'qqzz_unknown_123', 'allow_llm_fallback': False}) + self.assertTrue(disabled['isError']) + provider.assert_not_called() + enabled = self.tool('nl_ask', {'query': 'qqzz_unknown_123'}) + self.assertFalse(enabled['isError']) + self.assertEqual(enabled['structuredContent']['meta']['sourceLayer'], 'llm_fallback') + provider.assert_called_once() + with patch.object(dsl_llm, 'run_provider', return_value='OBSERVE catalog EXTRA'): + self.assertTrue(self.tool('nl_ask', {'query': 'qqzz_unknown_123'})['isError']) + + def test_observer_exception_is_execution_error(self): + with patch('monag.catalog.scan', side_effect=OSError('fixture failure')): + result = self.tool('execute_dsl', {'dsl_command': 'OBSERVE catalog'}) + self.assertTrue(result['isError']) + self.assertEqual(result['structuredContent']['status'], 'EXECUTION_ERROR') + + def test_advisory_and_issue_limit_use_validated_dsl(self): + with patch('monag.advise.advise', return_value={'recommendations': []}) as observer: + with patch('monag.advise.markdown', return_value='advice'): + with patch.object(dsl, 'execute', wraps=dsl.execute) as executor: + result = self.tool('monag_advise', {'limit': 3, 'radar': True, 'tier': 'floor'}) + self.assertFalse(result['isError']) + executor.assert_called_once() + self.assertEqual(executor.call_args.args[0].target, 'advise') + self.assertEqual(observer.call_args.kwargs['limit'], 3) + self.assertTrue(observer.call_args.kwargs['radar']) + with patch('monag.audit.scan', return_value={}) as scan: + with patch('monag.audit.markdown', return_value='audit'): + self.assertFalse(self.tool('monag_audit', {'issue_limit': 7})['isError']) + self.assertEqual(scan.call_args.kwargs['issue_limit'], 7) + + def test_structured_commands_and_mutated_query_objects_are_validated(self): + result = nl_contract.execute({'entity': 'catalog', 'operation': 'query'}, self.root, direct=True) + self.assert_envelope(result) + for value in [{'entity': 'catalog', 'operation': 'delete'}, + {'entity': 'catalog', 'operation': 'query', 'arguments': {'limit': True}}, + {'entity': 'catalog', 'operation': 'query', 'extra': 1}]: + self.assert_envelope(nl_contract.execute(value, self.root, direct=True), False) + query = dsl.Query('catalog') + query.limit = -1 + with self.assertRaises(ValueError): + dsl.execute(query, self.root) + + def start_http(self): + server, _ = panel.build_server(self.root, self.root / 'state', port=0) + thread = threading.Thread(target=server.serve_forever, kwargs={'poll_interval': .02}, daemon=True) + thread.start() + self.addCleanup(thread.join, 2) + self.addCleanup(server.server_close) + self.addCleanup(server.shutdown) + return server + + def request(self, server, path, body=None): + connection = http.client.HTTPConnection(*server.server_address, timeout=5) + self.addCleanup(connection.close) + connection.request('GET' if body is None else 'POST', path, + body=None if body is None else json.dumps(body), + headers={'Content-Type': 'application/json'}) + response = connection.getresponse() + return response.status, json.loads(response.read()) + + def test_cli_http_and_mcp_share_success_and_error_contract(self): + server = self.start_http() + code, schema = self.request(server, '/api/v1/schema') + self.assertEqual(code, 200) + self.assertEqual(schema, nl_contract.grammar()) + for valid in (True, False): + text = 'OBSERVE catalog' if valid else 'OBSERVE catalog BAD' + code, http = self.request(server, '/api/v1/dsl', {'dsl_command': text}) + rpc = self.tool('execute_dsl', {'dsl_command': text})['structuredContent'] + cli = subprocess.run([sys.executable, '-m', 'monag', '--root', str(self.root), + '--json', 'dsl', text], capture_output=True, text=True, timeout=20) + self.assertEqual(cli.returncode, 0 if valid else 1, cli.stderr) + cli_result = json.loads(cli.stdout) + for result in (http, rpc, cli_result): + self.assert_envelope(result, valid) + self.assertEqual(result['status'], rpc['status']) + self.assertEqual(result['meta']['canonicalDsl'], rpc['meta']['canonicalDsl']) + self.assertEqual(code, 200 if valid else 400) + code, result = self.request(server, '/api/v1/query', {'query': 'pokaż katalog', 'locale': 'pl'}) + self.assertEqual(code, 200) + self.assertEqual(result['meta']['sourceLayer'], 'nl_fast_path') + for body in [[], {'query': 4}, {'query': 'catalog', 'extra': 1}]: + code, result = self.request(server, '/api/v1/query', body) + self.assertEqual(code, 400) + self.assert_envelope(result, False) + with patch('monag.catalog.scan', side_effect=OSError('fixture failure')): + code, result = self.request(server, '/api/v1/dsl', {'dsl_command': 'OBSERVE catalog'}) + self.assertEqual(code, 500) + self.assert_envelope(result, False) + + def test_cli_ask_uses_standard_envelope_and_legacy_query_keeps_format(self): + for mode in ('ask', 'query'): + process = subprocess.run([sys.executable, '-m', 'monag', '--root', str(self.root), + '--json', mode, 'show catalog'], + capture_output=True, text=True, timeout=20) + self.assertEqual(process.returncode, 0, process.stderr) + result = json.loads(process.stdout) + if mode == 'ask': + self.assert_envelope(result) + self.assertEqual(result['meta']['sourceLayer'], 'nl_fast_path') + else: + self.assertEqual(result['status'], 'ok') + process = subprocess.run([sys.executable, '-m', 'monag', '--root', str(self.root), + '--json', 'query', 'qqzz_unknown_123'], + capture_output=True, text=True, timeout=20) + self.assertEqual(process.returncode, 1) + + def test_real_stdio_transports_schema_result_and_error(self): + messages = [ + {'jsonrpc': '2.0', 'id': 1, 'method': 'initialize'}, + {'jsonrpc': '2.0', 'method': 'notifications/initialized'}, + {'jsonrpc': '2.0', 'id': 2, 'method': 'resources/read', 'params': {'uri': 'schema://current'}}, + {'jsonrpc': '2.0', 'id': 3, 'method': 'tools/call', + 'params': {'name': 'execute_dsl', 'arguments': {'dsl_command': 'OBSERVE catalog'}}}, + {'jsonrpc': '2.0', 'id': 4, 'method': 'tools/call', + 'params': {'name': 'monag_query', 'arguments': {'query': 'qqzz_unknown_123'}}}, + ] + process = subprocess.run([sys.executable, '-m', 'monag', '--root', str(self.root), 'mcp'], + input='\n'.join(map(json.dumps, messages)) + '\n', + capture_output=True, text=True, timeout=20) + self.assertEqual(process.returncode, 0, process.stderr) + replies = [json.loads(line) for line in process.stdout.splitlines()] + self.assertEqual([r['id'] for r in replies], [1, 2, 3, 4]) + self.assertEqual(replies[1]['result']['contents'][0]['mimeType'], 'application/json') + self.assertFalse(replies[2]['result']['isError']) + self.assertTrue(replies[3]['result']['isError']) + + +if __name__ == '__main__': + unittest.main() From 5cb7cd6c175c4d7c69ef8a30737ce7a678ab45e0 Mon Sep 17 00:00:00 2001 From: Tom Softreck Date: Sat, 19 Sep 2026 19:11:01 +0200 Subject: [PATCH 2/2] fix(mcp): retain provenance when LLM translation is rejected --- project/ticket-078/README.md | 2 +- src/monag/dsl_llm.py | 4 ++-- src/monag/nl_contract.py | 3 ++- tests/test_nl_contract.py | 5 ++++- 4 files changed, 9 insertions(+), 5 deletions(-) diff --git a/project/ticket-078/README.md b/project/ticket-078/README.md index e07b384..2d2d004 100644 --- a/project/ticket-078/README.md +++ b/project/ticket-078/README.md @@ -23,4 +23,4 @@ The immutable source, supported schema profile and invocation examples are retur ## Validation -Full application run: 361 tests and 30 subtests passed. Twelve focused conformance tests pass, including an additional CLI compatibility case. The pinned normative command/result schemas independently validate real success and error envelopes; the advertised observation schema is valid JSON Schema. Ruff on all changed Python files and governance pass. Protected review and installed-runtime validation remain external delivery receipts. +Final full application run: 362 tests and 30 subtests passed, including twelve conformance tests covering CLI compatibility and rejected LLM translation provenance. The pinned normative command/result schemas independently validate real success and error envelopes; the advertised observation schema is valid JSON Schema. Ruff on all changed Python files and governance pass. Protected review and installed-runtime validation remain external delivery receipts. diff --git a/src/monag/dsl_llm.py b/src/monag/dsl_llm.py index 5bd8083..c461c3d 100644 --- a/src/monag/dsl_llm.py +++ b/src/monag/dsl_llm.py @@ -113,10 +113,10 @@ def resolve(text, command=None, timeout=None): answer = run_provider(SYSTEM_PROMPT + '\nRequest: ' + raw, command, timeout) line = extract_observe_line(answer) if line is None or line.lower() == 'observe none': - return None, provenance('none', None, raw) + return None, dict(provenance('none', None, raw), attempted_engine='llm') query = dsl.parse_dsl(line) if query is None: - return None, provenance('none', None, raw) + return None, dict(provenance('none', None, raw), attempted_engine='llm') query.raw_input = raw return query, provenance('llm', query.to_dsl(), raw) diff --git a/src/monag/nl_contract.py b/src/monag/nl_contract.py index 7606ba9..091b194 100644 --- a/src/monag/nl_contract.py +++ b/src/monag/nl_contract.py @@ -131,7 +131,8 @@ def execute(value, root, *, direct=False, allow_llm_fallback=True, locale=None, prov = dsl_llm.provenance('rule', query.to_dsl() if query else None, value) else: query, prov = dsl_llm.resolve(value, command=None if allow_llm_fallback else '') - source = 'llm_fallback' if prov['engine'] == 'llm' else 'nl_fast_path' + source = ('llm_fallback' if prov['engine'] == 'llm' or + prov.get('attempted_engine') == 'llm' else 'nl_fast_path') if query is None: raise ValueError('Unrecognized or invalid observation command') validate_query(query) diff --git a/tests/test_nl_contract.py b/tests/test_nl_contract.py index e47c322..5a366cc 100644 --- a/tests/test_nl_contract.py +++ b/tests/test_nl_contract.py @@ -108,7 +108,10 @@ def test_provider_translation_validates_and_can_be_disabled_per_call(self): self.assertEqual(enabled['structuredContent']['meta']['sourceLayer'], 'llm_fallback') provider.assert_called_once() with patch.object(dsl_llm, 'run_provider', return_value='OBSERVE catalog EXTRA'): - self.assertTrue(self.tool('nl_ask', {'query': 'qqzz_unknown_123'})['isError']) + rejected = self.tool('nl_ask', {'query': 'qqzz_unknown_123'}) + self.assertTrue(rejected['isError']) + self.assertEqual(rejected['structuredContent']['meta']['sourceLayer'], 'llm_fallback') + self.assertEqual(rejected['structuredContent']['meta']['provenance']['engine'], 'none') def test_observer_exception_is_execution_error(self): with patch('monag.catalog.scan', side_effect=OSError('fixture failure')):