From 22878f5d891aeb3a15254f9293e1c355e216dde5 Mon Sep 17 00:00:00 2001 From: witbrock Date: Tue, 15 Sep 2026 11:29:42 -0400 Subject: [PATCH] Clarify task source categories and reconcile corrected creates --- .../integrations/internal_mcp/catalogue.py | 13 +- src/backend/services/adaptive_turn_service.py | 13 +- src/backend/services/task_ontology_service.py | 9 +- tests/backend/test_adaptive_turn_service.py | 26 +++- .../test_coding_task_creation_recovery.py | 135 ++++++++++++++++++ 5 files changed, 190 insertions(+), 6 deletions(-) diff --git a/src/backend/integrations/internal_mcp/catalogue.py b/src/backend/integrations/internal_mcp/catalogue.py index 146eca393..cb2d54cb9 100644 --- a/src/backend/integrations/internal_mcp/catalogue.py +++ b/src/backend/integrations/internal_mcp/catalogue.py @@ -45115,7 +45115,13 @@ def _build_default_catalogue_task_and_workflow_definitions() -> List[MethodDefin enum_values={ "priority": ("low", "medium", "high", "critical"), }, - description="Create a new task in Vontology.", + description=( + "Create a new task in Vontology. task_source_id and its alias " + "source_id identify a registered source category (von_native " + "or jira_imported), not an external message ID or URL. For " + "a task created from a message, omit the category to default " + "to von_native; retain provenance in notes or reference_code." + ), ), output_schema=task_create_output_schema, category="write", @@ -45153,6 +45159,11 @@ def _build_default_catalogue_task_and_workflow_definitions() -> List[MethodDefin "Priority: low, medium, high, critical. Tasks start in 'pending' status and can include " "planning metadata (components, fix versions, sprint values, backlog rank), canonical " "task categories, source semantics, and richer task-detail fields. " + "task_source_id (alias source_id) is a registered source category: " + "von_native or jira_imported, also accepted as canonical concept IDs. " + "It is not an external message ID or URL. A task created from Gmail " + "is von_native by default: omit the category and retain the message " + "reference in notes or reference_code. " "For coding-agent work, preserve a user-specified model and reasoning level in " "requested_model (exact model ID, e.g. gpt-6-astra) and " "requested_reasoning_effort (e.g. high or xhigh; extra-high is xhigh). " diff --git a/src/backend/services/adaptive_turn_service.py b/src/backend/services/adaptive_turn_service.py index a95dd9747..7ccf2ee6b 100644 --- a/src/backend/services/adaptive_turn_service.py +++ b/src/backend/services/adaptive_turn_service.py @@ -4896,7 +4896,7 @@ def _reconcile_task_creation_attempts( ) -> None: """Reconcile a rejected create with a verified record for the same request. - A corrected invalid type changes the handler fingerprint, so neither an + A corrected invalid category can change the handler fingerprint, so neither an exact effect ID nor a task ID exists to join the rejected attempt. Compare the remaining creation arguments instead; assignment may be completed by a later verified field update on the returned task. This is an evidence @@ -4913,6 +4913,11 @@ def arguments_for(invocation: Mapping[str, Any]) -> dict[str, Any]: ): if alias in arguments: arguments.setdefault(field, arguments.pop(alias)) + # Match the handler's source-category alias precedence, including calls + # that supply both spellings or an empty canonical field. + source_id = arguments.pop("source_id", None) + if source_id or arguments.get("task_source_id"): + arguments["task_source_id"] = arguments.get("task_source_id") or source_id if not arguments.get("assignee_concept_id") and arguments.get( "created_by_concept_id" ): @@ -4937,6 +4942,10 @@ def arguments_for(invocation: Mapping[str, Any]) -> dict[str, Any]: error_code == "INVALID_DATA" and "unknown task type" in str(error or "").lower() ) + corrected_invalid_source = ( + error_code == "INVALID_DATA" + and "unknown task source" in str(error or "").lower() + ) for later_index in range(index + 1, len(tool_invocations)): later = tool_invocations[later_index] if later.get("execution_method", later.get("tool")) != "task_create": @@ -4954,6 +4963,8 @@ def arguments_for(invocation: Mapping[str, Any]) -> dict[str, Any]: excluded = set(assignment_fields) if corrected_invalid_type: excluded.add("task_type_ids") + if corrected_invalid_source: + excluded.add("task_source_id") if {k: v for k, v in expected.items() if k not in excluded} != { k: v for k, v in observed_arguments.items() if k not in excluded }: diff --git a/src/backend/services/task_ontology_service.py b/src/backend/services/task_ontology_service.py index 76bfa709b..2ccafd625 100644 --- a/src/backend/services/task_ontology_service.py +++ b/src/backend/services/task_ontology_service.py @@ -511,7 +511,14 @@ def normalise_task_source_id( resolved = _TASK_SOURCE_BY_SLUG.get(cleaned.lower()) or cleaned if resolved not in _TASK_SOURCE_IDS: - raise ValueError(f"Unknown task source: {cleaned}") + raise ValueError( + f"Unknown task source: {cleaned}. task_source_id (alias source_id) " + "must identify a registered source category: " + + ", ".join(item["slug"] for item in TASK_SOURCE_DEFINITIONS) + + ". It is not an external message ID or URL. For a task created " + "from a message, omit the source category to use von_native and " + "retain the external reference in notes or reference_code." + ) return resolved diff --git a/tests/backend/test_adaptive_turn_service.py b/tests/backend/test_adaptive_turn_service.py index 7d7a8976b..a6e8e5f34 100644 --- a/tests/backend/test_adaptive_turn_service.py +++ b/tests/backend/test_adaptive_turn_service.py @@ -16712,15 +16712,30 @@ def call(name: str, arguments: dict[str, Any]) -> None: "description", "organisation_concept_id", "requested_model", + "requested_reasoning_effort", + "priority", + "due_date", + "notes", "assignee_concept_id", "unverified", "changed", + "partial", "other_type_error", "later_reassignment", ], ) +@pytest.mark.parametrize( + "invalid_field,error", + [ + ("task_type_id", "Unknown task type"), + ("task_source_id", "Unknown task source"), + ("source_id", "Unknown task source"), + ], +) def test_task_create_reconciliation_requires_matching_verified_requested_object( difference: str, + invalid_field: str, + error: str, ) -> None: arguments = { "title": "Requested work", @@ -16728,6 +16743,11 @@ def test_task_create_reconciliation_requires_matching_verified_requested_object( "organisation_concept_id": "#V#org", "assignee_concept_id": "#V#codex_dgx", "requested_model": "test-model", + "requested_reasoning_effort": "high", + "priority": "high", + "due_date": "2026-10-01T12:00:00Z", + "notes": "Source: fixture-message-001", + "idempotency_key": "same-key-is-not-enough", } later_arguments = dict(arguments) fields = {"assignee_concept_id": "#V#codex_dgx"} @@ -16741,7 +16761,7 @@ def test_task_create_reconciliation_requires_matching_verified_requested_object( "tool": "task_create", "effective_arguments": { **arguments, - "task_type_id": "#V#task_specification", + invalid_field: "#V#invalid_category", }, }, { @@ -16752,14 +16772,14 @@ def test_task_create_reconciliation_requires_matching_verified_requested_object( ] snapshot = { "failed": { - "effect_status": "failed", + "effect_status": "partial" if difference == "partial" else "failed", "changed": difference == "changed", "failure_fact": { "error_code": "INVALID_DATA", "error": ( "Other validation failure" if difference == "other_type_error" - else "Unknown task type" + else error ), }, }, diff --git a/tests/backend/test_coding_task_creation_recovery.py b/tests/backend/test_coding_task_creation_recovery.py index d78ea5c5e..1cb3e839a 100644 --- a/tests/backend/test_coding_task_creation_recovery.py +++ b/tests/backend/test_coding_task_creation_recovery.py @@ -272,6 +272,141 @@ def test_create_schema_exposes_preferences_without_exposing_trusted_creator(gate assert "created_by_concept_id" not in schema["properties"] +def test_create_contract_distinguishes_source_category_from_external_reference(gateway): + definition = gateway.get_method_definition("task_create") + schema = _model_visible_input_schema(definition) + assert {"task_source_id", "source_id", "notes", "reference_code"} <= schema[ + "properties" + ].keys() + for description in (definition.description, definition.input_schema.description): + assert "source_id" in description + assert "category" in description + assert "message ID" in description + assert "notes" in description and "reference_code" in description + + +@pytest.mark.parametrize("omit_priority", [False, True]) +def test_message_source_correction_retains_provenance_and_reconciles_final_answer( + native_store, gateway, monkeypatch, omit_priority +): + from src.backend.languagemodels.structured_tool_calling.types import ( + LLMResponse, + ToolCall, + ) + from src.backend.services.adaptive_turn_service import execute_adaptive_turn + from src.backend.services import turn_execution_record_service as records + + observations = [] + + def record_observation(**kwargs): + observations.append(deepcopy(kwargs)) + return {"updated": True, "duplicate": False} + + monkeypatch.setattr(records, "record_effect_observation_phase", record_observation) + requested = { + "title": "Review the seminar agenda", + "description": "Check the agenda and prepare discussion points.", + "priority": "high", + "due_date": "2026-10-01T12:00:00Z", + "notes": "Source: Gmail message fixture-message-001; seminar agenda.", + "reference_code": "fixture-message-001", + "idempotency_key": "fixture-seminar-agenda", + } + + class MessageSourceReplay: + call = 0 + answer = "" + + def generate_with_tools(self, prompt, available_tools, **kwargs): + self.call += 1 + assert self.call <= 4, "unexpected replay model call" + if self.call == 4: + task_id = native_store.concepts.find_one({})["concept_id"] + self.answer = f"Created task {task_id}, retaining the Gmail reference." + return LLMResponse(text_response=self.answer) + assert native_store.concepts.count_documents({}) == 0 + arguments = dict(requested) + if self.call == 1: + arguments["task_source_id"] = "fixture-message-001" + if self.call <= 2: + arguments["source_id"] = "fixture-message-001" + elif omit_priority: + arguments.pop("priority") + return LLMResponse( + text_response="", + tool_calls=[ + ToolCall( + tool_name="turn_invoke_capability", + call_id=f"source-fixture-{self.call}", + payload={"name": "task_create", "arguments": arguments}, + ) + ], + ) + + client = MessageSourceReplay() + result = execute_adaptive_turn( + gateway=gateway, + prompt="Create a high-priority task from the seminar email, due 1 October at noon UTC. Retain the Gmail reference.", + context=[], + llm_client=client, + model="fixture-no-model-request", + user_namespace="#V#alice@lab", + user_concept_id="#V#alice", + org_concept_id="#V#lab", + conversation_id="fixture-source", + turn_id="fixture-message-source-recovery", + turn_budget_seconds=20, + final_synthesis_reserve_seconds=2, + ) + assert len(result.tool_invocations) == 3 + first, second, successful = result.tool_invocations + for rejected in (first, second): + assert rejected["effect_status"] == "failed" + assert rejected["changed"] is False + assert rejected["error_code"] == "INVALID_DATA" + error = rejected["failure_fact"]["error"] + assert "Unknown task source: fixture-message-001" in error + assert "registered source category" in error + assert "alias source_id" in error + assert "notes or reference_code" in error + assert (rejected.get("recovery_status") == "succeeded") is not omit_priority + if not omit_priority: + assert rejected["recovered_by_effect_id"] == successful["effect_id"] + assert first["effective_arguments"]["task_source_id"] == "fixture-message-001" + assert first["effective_arguments"]["source_id"] == "fixture-message-001" + assert "task_source_id" not in second["effective_arguments"] + assert second["effective_arguments"]["source_id"] == "fixture-message-001" + # Finality is a projection; durable failure observations remain unchanged. + failures = [ + item["observation"] + for item in observations + if item.get("observation", {}).get("effect_status") == "failed" + ] + assert len(failures) >= 2 + assert all("recovery_status" not in item for item in failures) + assert successful["canonical_readback"]["verified"] is True + task_id = successful["canonical_readback"]["task_concept_id"] + with override_current_actor("#V#alice", "#V#lab"): + task = tasks.get_task(task_id) + assert native_store.concepts.count_documents({}) == 1 + assert task["notes"] == requested["notes"] + assert task["reference_code"] == requested["reference_code"] + assert task["task_source_id"] == "#V#von_native_task_source" + assert task["assignee_concept_id"] == "#V#alice" + assert task["organisation_concept_id"] == "#V#lab" + assert task["priority"] == ("medium" if omit_priority else "high") + assert result.terminal_status == ( + "effect_partially_completed" if omit_priority else "completed" + ), result.response_text + if not omit_priority: + assert result.response_text == client.answer + assert result.response_authority == "model" + effects = records._build_tool_authored_mutation_effects( + tool_invocations=result.tool_invocations + ) + assert effects and all(effect["status"] == "satisfied" for effect in effects) + + @pytest.mark.parametrize( "field,value,hint", [