Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
32 changes: 29 additions & 3 deletions loopx/control_plane/collaboration/source_grant_observation.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,9 @@

from ...agent_registry import registered_agent_ids_for_goal
from ..goals.activation import goal_is_stopped
from ..goals.goal_ref_validation import exact_goal_ref
from ..projects.registry_codec import (
SOURCE_SESSION_PROFILE_ID,
load_project_registry,
require_runtime_compatible_project_registry,
)
Expand Down Expand Up @@ -126,9 +128,33 @@ def source_context_authority(
registry = load_project_registry(registry_path)
if not isinstance(registry, dict):
raise ValueError("invalid registry")
require_runtime_compatible_project_registry(
registry, operation="context source recipient observation"
)
if registry.get("profile_id") == SOURCE_SESSION_PROFILE_ID:
goals = registry.get("goals")
if not isinstance(goals, list):
raise ValueError("source-session registry has no Goal list")
# This is observation for context handoff, not execution admission.
# Enumerate only instance-bound Goals; the handoff's own Goal scope
# still rechecks the selected exact GoalRef before it commits.
instantiated = []
for goal in goals:
if not isinstance(goal, dict):
continue
goal_id = goal.get("id")
instance_id = goal.get("goal_instance_id")
if not isinstance(goal_id, str) or not isinstance(instance_id, str):
continue
try:
exact_goal_ref(goal_id, instance_id)
except ValueError:
continue
instantiated.append(goal)
if not instantiated:
raise ValueError("source-session registry has no instantiated Goal")
registry = {**registry, "goals": instantiated}
else:
require_runtime_compatible_project_registry(
registry, operation="context source recipient observation"
)
except (OSError, ValueError, TypeError):
return {"mode": "unavailable", "targets": []}
observed = registered_context_recipients(registry)
Expand Down
2 changes: 1 addition & 1 deletion loopx/semantics/project_registry_io_manifest_v1.json
Original file line number Diff line number Diff line change
Expand Up @@ -1039,7 +1039,7 @@
},
{
"site": "loopx/control_plane/collaboration/source_grant_observation.py::<module>.source_context_authority::codec_read:load_project_registry#1",
"line": 126,
"line": 128,
"column": 20,
"kind": "codec_read",
"api": "load_project_registry",
Expand Down
42 changes: 42 additions & 0 deletions tests/test_manager_context_handoff.py
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,48 @@ def reject_enumeration(_registry):
assert not _root(root).exists()


@pytest.mark.parametrize("strict_envelope", [False, True])
@pytest.mark.parametrize("other_goal", [
{"goal_instance_id": "ginst_" + "b" * 32},
{},
{"goal_instance_id": "invalid"},
{"goal_instance_id": "ginst_" + "b" * 32, "id": "unsafe/alias"},
{"goal_instance_id": "ginst_" + "b" * 32, "activation_state": "stopped"},
{"goal_instance_id": "ginst_" + "b" * 32, "activation_state": "invalid"},
])
def test_source_session_catalog_keeps_only_current_instantiated_recipients(
fixture, strict_envelope, other_goal
):
from loopx.control_plane.projects import registry_codec

root, registry, session, turn, request = fixture
payload = json.loads(registry.read_text())
payload["profile_id"] = registry_codec.SOURCE_SESSION_PROFILE_ID
payload["goals"][0]["goal_instance_id"] = "ginst_" + "a" * 32
payload["goals"][1].update(other_goal)
source_registry = registry.with_name("source-session-registry.json")
if strict_envelope:
with registry_codec.source_session_registry_transaction(
source_registry,
operation="create context catalog fixture",
create=lambda: payload,
) as transaction:
transaction.commit(payload)
else:
source_registry.write_text(json.dumps(payload))
before = source_registry.read_bytes()
expected = [request]
if other_goal == {"goal_instance_id": "ginst_" + "b" * 32}:
expected.insert(0, {"goal_id": "other", "agent_id": "peer"})
result = authority(root, source_registry, session, turn)
assert result["mode"] == "context_only"
assert result["targets"] == expected
goal_session = {**session, "channel_id": "goal.research", "goal_id": "research"}
assert authority(root, source_registry, goal_session, turn)["targets"] == [request]
assert source_registry.read_bytes() == before
assert not _root(root).exists()


def test_stopped_goal_is_not_a_context_recipient_and_revokes_replay(fixture):
root, registry, session, turn, request = fixture
assert request in authority(root, registry, session, turn)["targets"]
Expand Down
Loading