From 992e9e7971c3394963d000646d49a8076e22b12b Mon Sep 17 00:00:00 2001 From: mintaka Date: Tue, 6 Oct 2026 01:05:38 -0400 Subject: [PATCH] test(e2e): prove the owner peering lifecycle over the real transport (RIG-4365) Unpeered, one-sided, mutual, revoked and restored peering, each checked on OpenDM and on live delivery to the peer agent session. Co-authored-by: Matt Wilkinson --- go/e2e/agent_ops.go | 8 +- go/e2e/legcomms_tenant_test.go | 352 ++++++++++++++++++++++++++++++++- 2 files changed, 353 insertions(+), 7 deletions(-) diff --git a/go/e2e/agent_ops.go b/go/e2e/agent_ops.go index c937a9651..337b64da5 100644 --- a/go/e2e/agent_ops.go +++ b/go/e2e/agent_ops.go @@ -102,7 +102,13 @@ func (f *Fixture) Resume(ctx context.Context, containerName, resumeSessionID str // AwaitTurnSettled via its own derived settleTimeout, mirroring how // SubscribeComms opens under the caller's ctx and AwaitDelivery bounds the read. func (f *Fixture) OpenSessionTail(ctx context.Context, sessionID string) (*connect.ServerStreamForClient[compassv1.AgentSessionFrame], error) { - stream, err := f.Compass().SubscribeAgentSession(ctx, connect.NewRequest(&compassv1.SubscribeAgentSessionRequest{ + return f.OpenSessionTailAs(ctx, f.Compass(), sessionID) +} + +// OpenSessionTailAs opens the tail as client's account, which must be a member +// of the agent's home channel (another owner's agent is invisible to admin). +func (f *Fixture) OpenSessionTailAs(ctx context.Context, client compassServiceClient, sessionID string) (*connect.ServerStreamForClient[compassv1.AgentSessionFrame], error) { + stream, err := client.SubscribeAgentSession(ctx, connect.NewRequest(&compassv1.SubscribeAgentSessionRequest{ SessionId: sessionID, })) if err != nil { diff --git a/go/e2e/legcomms_tenant_test.go b/go/e2e/legcomms_tenant_test.go index e07e6b92c..1326fbbc9 100644 --- a/go/e2e/legcomms_tenant_test.go +++ b/go/e2e/legcomms_tenant_test.go @@ -58,6 +58,29 @@ const ( t4CanaryBody = "t4: globally-visible canary owner-2 IS entitled to" ) +const ( + t5Owner1Handle = "t5-owner-1" + t5Owner2Handle = "t5-owner-2" + t5Agent1Handle = "t5-agent-1" + t5Agent2Handle = "t5-agent-2" + t5Agent3Handle = "t5-agent-3" + t5GhostHandle = "t5-ghost-agent" + + t5SharedChannel = "t5-peer-room" + t5Topic = "general" + t5EventMarker = "t5-peer-lifecycle-event" + t5ModelMarker = "t5-peer-lifecycle-model" + t5CanaryMarker = "t5-peer-lifecycle-canary" +) + +func init() { + registerSharedFixtureOption( + WithCannedMarkerReply(t5EventMarker, "peer lifecycle message received"), + WithCannedMarkerReply(t5ModelMarker, "peer lifecycle message received"), + WithCannedMarkerReply(t5CanaryMarker, "peer lifecycle canary received"), + ) +} + // TestCommsTenantVisibilityTransport is the three-assertion multi-tenant // transport leg: // @@ -113,8 +136,8 @@ func TestCommsTenantVisibilityTransport(t *testing.T) { // about. So each agent is created over its OWN owner's observer client — the // same "call the generated client AsObserver returned" shape assertion 3 // uses for OpenDM, not a new fixture primitive. - agent1ID, agent1Owner := createAgentAs(ctx, t, owner1Comms, t4Agent1Handle, "T4 Agent One") - agent2ID, agent2Owner := createAgentAs(ctx, t, owner2Comms, t4Agent2Handle, "T4 Agent Two") + agent1ID, agent1Owner, _ := createAgentAs(ctx, t, owner1Comms, t4Agent1Handle, "T4 Agent One") + agent2ID, agent2Owner, _ := createAgentAs(ctx, t, owner2Comms, t4Agent2Handle, "T4 Agent Two") if agent1ID == agent2ID { t.Fatalf("the two owners' agents share account id %q; the per-owner agent namespaces are not distinct", agent1ID) } @@ -329,12 +352,12 @@ func TestCommsTenantVisibilityTransport(t *testing.T) { // createAgentAs creates an agent over an EXPLICIT comms client (an observer's), // so the agent is owned by that client's account rather than by the fixture's // bootstrap admin — CreateAgent places the new agent under the CALLER's resolved -// owner (comms/comms.go:109). Returns the new account id AND the owner the // server actually resolved, because per-owner ownership is an emergent property // of which client was passed: nothing in the request names an owner, so the -// caller cannot assume it landed where intended and must check. Fatal on +// caller cannot assume it landed where intended and must check. Also returns the +// agent's home channel. Fatal on // failure: a setup miss makes every assertion below meaningless. -func createAgentAs(ctx context.Context, t *testing.T, comms commsServiceClient, handle, displayName string) (accountID, ownerID string) { +func createAgentAs(ctx context.Context, t *testing.T, comms commsServiceClient, handle, displayName string) (accountID, ownerID, homeChannelID string) { t.Helper() rctx, cancel := context.WithTimeout(ctx, rpcTimeout) defer cancel() @@ -356,7 +379,7 @@ func createAgentAs(ctx context.Context, t *testing.T, comms commsServiceClient, if owner == "" { t.Fatalf("CreateAgent(%s) returned account %q with no agent owner id; cannot verify which tenant it landed under", handle, id) } - return id, owner + return id, owner, account.GetAgent().GetHomeChannelId() } // awaitSubscriptionLive drains the leading control frame a sinceSeq=0 @@ -456,3 +479,320 @@ func rejectionMessage(t *testing.T, peerHandle string, err error) string { } return cerr.Message() } + +func TestCommsPeerLifecycleTransport(t *testing.T) { + if !podmanUsable() { + t.Skip("rootless podman cannot run compass-agent:latest here; skipping the real-stack e2e") + } + ctx := context.Background() + f := sharedFixture(t) + + owner1ID, err := f.CreateUser(ctx, t5Owner1Handle, "T5 Owner One") + if err != nil { + t.Fatalf("CreateUser(owner 1): %v", err) + } + owner2ID, err := f.CreateUser(ctx, t5Owner2Handle, "T5 Owner Two") + if err != nil { + t.Fatalf("CreateUser(owner 2): %v", err) + } + owner1Compass, owner1Comms, err := f.AsObserver(ctx, t5Owner1Handle) + if err != nil { + t.Fatalf("AsObserver(owner 1): %v", err) + } + owner2Compass, owner2Comms, err := f.AsObserver(ctx, t5Owner2Handle) + if err != nil { + t.Fatalf("AsObserver(owner 2): %v", err) + } + agent1ID, agent1Owner, agent1Home := createAgentAs(ctx, t, owner1Comms, t5Agent1Handle, "T5 Agent One") + agent2ID, agent2Owner, agent2Home := createAgentAs(ctx, t, owner2Comms, t5Agent2Handle, "T5 Agent Two") + agent3ID, agent3Owner, _ := createAgentAs(ctx, t, owner2Comms, t5Agent3Handle, "T5 Agent Three") + if agent1ID == agent2ID || agent2ID == agent3ID || agent1Owner != owner1ID || agent2Owner != owner2ID || agent3Owner != owner2ID { + t.Fatalf("agent ownership = (%q,%q), (%q,%q), (%q,%q); want distinct agents under owners %q and %q", agent1ID, agent1Owner, agent2ID, agent2Owner, agent3ID, agent3Owner, owner1ID, owner2ID) + } + _, agent1Comms, err := f.AsObserver(ctx, t5Owner1Handle+"/"+t5Agent1Handle) + if err != nil { + t.Fatalf("AsObserver(agent 1): %v", err) + } + _, agent2Comms, err := f.AsObserver(ctx, t5Owner2Handle+"/"+t5Agent2Handle) + if err != nil { + t.Fatalf("AsObserver(agent 2): %v", err) + } + + roomID, err := f.CreateChannel(ctx, owner1ID, t5SharedChannel, false) + if err != nil { + t.Fatalf("CreateChannel(shared room): %v", err) + } + addMember := func(client commsServiceClient, handle string) { + t.Helper() + rctx, cancel := context.WithTimeout(ctx, rpcTimeout) + defer cancel() + _, err := client.UpdateChannelMembers(rctx, connect.NewRequest(&compassv1.UpdateChannelMembersRequest{ + ChannelId: roomID, AddMemberHandles: []string{handle}, SubscribeHandles: []string{handle}, + })) + if err != nil { + t.Fatalf("UpdateChannelMembers(add %s): %v", handle, err) + } + } + addMember(owner1Comms, t5Owner2Handle) + addMember(owner1Comms, t5Owner1Handle+"/"+t5Agent1Handle) + addMember(owner2Comms, t5Owner2Handle+"/"+t5Agent2Handle) + + owner1Stream, err := f.SubscribeCommsAsObserver(ctx, owner1Comms, 0) + if err != nil { + t.Fatalf("SubscribeCommsAsObserver(owner 1): %v", err) + } + defer owner1Stream.Close() + awaitSubscriptionLive(ctx, t, owner1Stream) + owner2Stream, err := f.SubscribeCommsAsObserver(ctx, owner2Comms, 0) + if err != nil { + t.Fatalf("SubscribeCommsAsObserver(owner 2): %v", err) + } + defer owner2Stream.Close() + awaitSubscriptionLive(ctx, t, owner2Stream) + + container1, err := f.Provision(ctx, agent1ID, "t5-agent-1-provision") + if err != nil { + t.Fatalf("Provision(agent 1): %v", err) + } + t.Cleanup(func() { + if err := f.RemoveWorkspace(ctx, container1, "t5-agent-1-teardown"); err != nil { + t.Errorf("RemoveWorkspace(agent 1): %v", err) + } + }) + session1, err := f.StartSession(ctx, container1) + if err != nil { + t.Fatalf("StartSession(agent 1): %v", err) + } + container2, err := f.Provision(ctx, agent2ID, "t5-agent-2-provision") + if err != nil { + t.Fatalf("Provision(agent 2): %v", err) + } + t.Cleanup(func() { + if err := f.RemoveWorkspace(ctx, container2, "t5-agent-2-teardown"); err != nil { + t.Errorf("RemoveWorkspace(agent 2): %v", err) + } + }) + session2, err := f.StartSession(ctx, container2) + if err != nil { + t.Fatalf("StartSession(agent 2): %v", err) + } + tail1, err := f.OpenSessionTailAs(ctx, owner1Compass, session1) + if err != nil { + t.Fatalf("OpenSessionTail(agent 1): %v", err) + } + defer tail1.Close() + tail2, err := f.OpenSessionTailAs(ctx, owner2Compass, session2) + if err != nil { + t.Fatalf("OpenSessionTail(agent 2): %v", err) + } + defer tail2.Close() + + peerHandle := t5Owner2Handle + "/" + t5Agent2Handle + unknownCode, unknownMsg := openDMRejection(ctx, t, agent1Comms, t5GhostHandle) + if unknownCode != connect.CodeNotFound { + t.Fatalf("OpenDM(unknown %q) = %v, want NOT_FOUND", t5GhostHandle, unknownCode) + } + assertRejectedLikeUnknown := func(name string, client commsServiceClient, handle string) { + t.Helper() + code, msg := openDMRejection(ctx, t, client, handle) + if code != connect.CodeNotFound { + t.Fatalf("%s OpenDM(%q) = %v, want NOT_FOUND", name, handle, code) + } + got := strings.ReplaceAll(msg, handle, "") + want := strings.ReplaceAll(unknownMsg, t5GhostHandle, "") + if got != want { + t.Fatalf("%s rejection %q is not redacted-identical to unknown rejection %q", name, got, want) + } + } + postAs := func(name string, client commsServiceClient, channelID, body string) string { + t.Helper() + id, err := f.PostMessageAsObserver(ctx, client, channelID, t5Topic, body) + if err != nil { + t.Fatalf("PostMessageAsObserver(%s): %v", name, err) + } + if id == "" { + t.Fatalf("PostMessageAsObserver(%s) returned an empty id", name) + } + return id + } + observe := func(name string, stream *connect.ServerStreamForClient[compassv1.SubscribeCommsResponse], id, author string) { + t.Helper() + msg, err := f.AwaitDelivery(ctx, stream, func(m *compassv1.Message) bool { return m.GetId() == id }) + if err != nil { + t.Fatalf("%s observer did not receive message %s: %v", name, id, err) + } + if msg.GetAuthorAccountId() != author { + t.Fatalf("%s message author = %q, want %q", name, msg.GetAuthorAccountId(), author) + } + } + assertDelivered := func(name string, tail *connect.ServerStreamForClient[compassv1.AgentSessionFrame], id string) { + t.Helper() + if _, err := f.AwaitControlDispatchOn(ctx, tail, func(kind, got string) bool { + return got == id && strings.Contains(kind, "DELIVER") + }); err != nil { + t.Fatalf("%s did not receive DELIVER for %s: %v", name, id, err) + } + if err := f.AwaitTurnSettled(ctx, tail); err != nil { + t.Fatalf("%s did not settle after %s: %v", name, id, err) + } + } + assertMentioned := func(name string, tail *connect.ServerStreamForClient[compassv1.AgentSessionFrame], id string) { + t.Helper() + kind, err := f.AwaitControlDispatchOn(ctx, tail, func(_, got string) bool { return got == id }) + if err != nil { + t.Fatalf("%s received no session injection for qualified mention %s: %v", name, id, err) + } + if !strings.Contains(kind, "STEER") { + t.Errorf("%s received %s for qualified mention %s, want STEER", name, kind, id) + } + if err := f.AwaitTurnSettled(ctx, tail); err != nil { + t.Fatalf("%s did not settle after mention %s: %v", name, id, err) + } + } + // A post by an agent with a live session is held until that session settles, + // so each agent-authored post is followed by an owner-driven turn of its author. + settleAuthor := func(name string, owner commsServiceClient, home string, tail *connect.ServerStreamForClient[compassv1.AgentSessionFrame]) { + t.Helper() + postAs(name+" settle trigger", owner, home, name+" "+t5EventMarker) + if err := f.AwaitTurnSettled(ctx, tail); err != nil { + t.Fatalf("%s author did not settle: %v", name, err) + } + } + settle1 := func(name string) { t.Helper(); settleAuthor(name, owner1Comms, agent1Home, tail1) } + settle2 := func(name string) { t.Helper(); settleAuthor(name, owner2Comms, agent2Home, tail2) } + + assertNoReach := func(name string, authorClient commsServiceClient, authorID string, settle func(string), deniedChannel, body string, + observer *connect.ServerStreamForClient[compassv1.SubscribeCommsResponse], recipientTail *connect.ServerStreamForClient[compassv1.AgentSessionFrame], + canaryClient commsServiceClient, canaryAuthor string, canaryChannel string) { + t.Helper() + deniedID := postAs(name+" denied", authorClient, deniedChannel, body) + observe(name+" denied event", observer, deniedID, authorID) + settle(name + " denied") + canaryID := postAs(name+" canary", canaryClient, canaryChannel, name+" "+t5CanaryMarker) + observe(name+" canary event", observer, canaryID, canaryAuthor) + var reachedID string + _, err := f.AwaitControlDispatchOn(ctx, recipientTail, func(kind, id string) bool { + if id == deniedID { + reachedID = id + return true + } + if id == canaryID && strings.Contains(kind, "DELIVER") { + reachedID = id + return true + } + return false + }) + if err != nil { + t.Fatalf("%s target session received neither the denied post nor its canary: %v", name, err) + } + if reachedID != canaryID { + t.Fatalf("%s target session received denied message %s before canary %s", name, deniedID, canaryID) + } + if err := f.AwaitTurnSettled(ctx, recipientTail); err != nil { + t.Fatalf("%s canary did not settle: %v", name, err) + } + } + + assertRejectedLikeUnknown("unpeered co-member", agent1Comms, peerHandle) + assertNoReach("unpeered A→B", agent1Comms, agent1ID, settle1, roomID, "t5-a-unpeered-post "+t5EventMarker, + owner2Stream, tail2, owner2Comms, owner2ID, roomID) + + rctx, cancel := context.WithTimeout(ctx, rpcTimeout) + oneSided, err := owner1Comms.ApprovePeer(rctx, connect.NewRequest(&compassv1.ApprovePeerRequest{PeerHandle: t5Owner2Handle})) + cancel() + if err != nil { + t.Fatalf("ApprovePeer(owner 1 → owner 2): %v", err) + } + if got := oneSided.Msg.GetPeering().GetState(); got != compassv1.PeeringState_PEERING_STATE_PENDING_OUTGOING { + t.Fatalf("one-sided peering state = %v, want PENDING_OUTGOING", got) + } + assertRejectedLikeUnknown("one-sided approval", agent1Comms, peerHandle) + rctx, cancel = context.WithTimeout(ctx, rpcTimeout) + mutual, err := owner2Comms.ApprovePeer(rctx, connect.NewRequest(&compassv1.ApprovePeerRequest{PeerHandle: t5Owner1Handle})) + cancel() + if err != nil { + t.Fatalf("ApprovePeer(owner 2 → owner 1): %v", err) + } + if got := mutual.Msg.GetPeering().GetState(); got != compassv1.PeeringState_PEERING_STATE_APPROVED { + t.Fatalf("mutual peering state = %v, want APPROVED", got) + } + + dm, err := agent1Comms.OpenDM(ctx, connect.NewRequest(&compassv1.OpenDMRequest{PeerHandle: peerHandle})) + if err != nil { + t.Fatalf("OpenDM(mutually peered): %v", err) + } + dmID := dm.Msg.GetChannel().GetId() + dmMessageID := postAs("agent 1 DM", agent1Comms, dmID, "t5-a-dm-positive "+t5EventMarker) + observe("owner 2 DM", owner2Stream, dmMessageID, agent1ID) + settle1("agent 1 DM") + assertDelivered("agent 2 DM", tail2, dmMessageID) + list, err := agent2Comms.ListMessages(ctx, connect.NewRequest(&compassv1.ListMessagesRequest{ + Container: &compassv1.ListMessagesRequest_ChannelId{ChannelId: dmID}, + Limit: 20, + })) + if err != nil { + t.Fatalf("agent 2 ListMessages(DM): %v", err) + } + foundDM := false + for _, msg := range list.Msg.GetMessages() { + if msg.GetId() == dmMessageID && firstBlockText(msg) == "t5-a-dm-positive "+t5EventMarker { + foundDM = true + } + } + if !foundDM { + t.Fatalf("agent 2's DM history does not contain message %s", dmMessageID) + } + + aPostID := postAs("agent 1 shared post", agent1Comms, roomID, "t5-a-approved-post "+t5EventMarker) + observe("owner 2 shared post", owner2Stream, aPostID, agent1ID) + settle1("agent 1 shared post") + assertDelivered("agent 2 shared post", tail2, aPostID) + bPostID := postAs("agent 2 shared post", agent2Comms, roomID, "t5-b-approved-post "+t5ModelMarker) + observe("owner 1 shared post", owner1Stream, bPostID, agent2ID) + settle2("agent 2 shared post") + assertDelivered("agent 1 shared post", tail1, bPostID) + aMentionID := postAs("agent 1 qualified mention", agent1Comms, roomID, + "t5-a-approved-mention @"+peerHandle+" "+t5EventMarker) + observe("owner 2 qualified mention", owner2Stream, aMentionID, agent1ID) + settle1("agent 1 qualified mention") + assertMentioned("agent 2", tail2, aMentionID) + bMentionID := postAs("agent 2 qualified mention", agent2Comms, roomID, + "t5-b-approved-mention @"+t5Owner1Handle+"/"+t5Agent1Handle+" "+t5ModelMarker) + observe("owner 1 qualified mention", owner1Stream, bMentionID, agent2ID) + settle2("agent 2 qualified mention") + assertMentioned("agent 1", tail1, bMentionID) + + rctx, cancel = context.WithTimeout(ctx, rpcTimeout) + revoked, err := owner1Comms.RevokePeer(rctx, connect.NewRequest(&compassv1.RevokePeerRequest{PeerHandle: t5Owner2Handle})) + cancel() + if err != nil { + t.Fatalf("RevokePeer(owner 2): %v", err) + } + if !revoked.Msg.GetDeleted() { + t.Fatal("RevokePeer(owner 2) returned deleted=false") + } + assertRejectedLikeUnknown("revoked co-member", agent1Comms, peerHandle) + assertRejectedLikeUnknown("revoked third agent", agent1Comms, t5Owner2Handle+"/"+t5Agent3Handle) + + assertNoReach("revoked A→B post", agent1Comms, agent1ID, settle1, roomID, "t5-a-revoked-post "+t5EventMarker, + owner2Stream, tail2, owner2Comms, owner2ID, roomID) + assertNoReach("revoked A→B mention", agent1Comms, agent1ID, settle1, roomID, "t5-a-revoked-mention @"+peerHandle+" "+t5EventMarker, + owner2Stream, tail2, owner2Comms, owner2ID, roomID) + assertNoReach("revoked B→A post", agent2Comms, agent2ID, settle2, roomID, "t5-b-revoked-post "+t5ModelMarker, + owner1Stream, tail1, owner1Comms, owner1ID, roomID) + assertNoReach("revoked B→A mention", agent2Comms, agent2ID, settle2, roomID, "t5-b-revoked-mention @"+t5Owner1Handle+"/"+t5Agent1Handle+" "+t5ModelMarker, + owner1Stream, tail1, owner1Comms, owner1ID, roomID) + assertNoReach("revoked existing DM post", agent1Comms, agent1ID, settle1, dmID, "t5-a-dm-revoked "+t5EventMarker, + owner2Stream, tail2, owner2Comms, owner2ID, roomID) + + rctx, cancel = context.WithTimeout(ctx, rpcTimeout) + if _, err := owner1Comms.ApprovePeer(rctx, connect.NewRequest(&compassv1.ApprovePeerRequest{PeerHandle: t5Owner2Handle})); err != nil { + cancel() + t.Fatalf("ApprovePeer(owner 2 after revoke): %v", err) + } + cancel() + restoredID := postAs("agent 1 restored shared post", agent1Comms, roomID, "t5-a-restored-post "+t5EventMarker) + observe("owner 2 restored shared post", owner2Stream, restoredID, agent1ID) + settle1("agent 1 restored shared post") + assertDelivered("agent 2 restored shared post", tail2, restoredID) +}