From bba8ad0154c9b43aeb3eab2d463f183017ac7581 Mon Sep 17 00:00:00 2001 From: mintaka Date: Fri, 25 Sep 2026 12:13:04 -0400 Subject: [PATCH 1/4] feat(comms): publish message_posted on the event fabric (RIG-3107) Adds the publish side of the delivery cutover in `docs/designs/infra/runtime/compass-managed-delivery-cutover/design.md` (Plan, T1). - `NewComms` takes a `fabric.EventFabric`. A nil fabric keeps `message_posted` on the bus only. Every call site passes `nil` until server assembly constructs the fabric. - After the bus publish, `publishMessagePosted` publishes `EventRef{tenant, message_posted, row id}` on `CommsSubject(tenant, message_posted)`. The tenant comes from `Store.EffectiveTenant`, the exported set-or-bootstrap resolver the hub already uses. The record names it `ResolveTenant`; this reuses the existing export instead of adding a second one. - The auth registry guard now names the internal compass.v1 service files (runner, agent gateway, guest control) as ungated. Linking the fabric into comms put them in the auth test binary, and they are never mounted behind AdminGate. - A publish failure after commit is logged and counted on `compass.delivery.fabric_publish_failures`, never returned to the poster. Tests (pgtest, fake fabric): a genuine post publishes exactly one ref after the row commits; an idempotent retry publishes nothing; a publish failure still succeeds and increments the counter; `RespondToAsk` publishes the answer message. Spec-impact: none. Refs RIG-3107 Co-authored-by: Matt Wilkinson --- go/internal/auth/classify_exhaustive_test.go | 20 +- go/internal/auth/interceptor_pgtest_test.go | 2 +- go/internal/comms/agent_caller_pgtest_test.go | 2 +- go/internal/comms/comms.go | 28 ++- go/internal/comms/comms_test.go | 2 +- .../comms/fabric_publish_pgtest_test.go | 222 ++++++++++++++++++ go/internal/comms/mapping.go | 19 ++ go/internal/comms/roster_pgtest_test.go | 2 +- .../comms/subscribe_failclosed_test.go | 2 +- go/internal/comms/subscribe_test.go | 2 +- .../runnerhub/integration_pgtest_test.go | 2 +- go/server/cors_pgtest_test.go | 2 +- go/server/dm_e2e_pgtest_test.go | 2 +- go/server/lifecycle_e2e_pgtest_test.go | 2 +- go/server/lifecycle_pgtest_test.go | 2 +- go/server/offline_mention_e2e_pgtest_test.go | 2 +- go/server/otel_emission_pgtest_test.go | 6 +- go/server/serve.go | 2 +- go/server/serve_forge_pgtest_test.go | 2 +- go/server/serve_seed_pgtest_test.go | 2 +- 20 files changed, 301 insertions(+), 24 deletions(-) create mode 100644 go/internal/comms/fabric_publish_pgtest_test.go diff --git a/go/internal/auth/classify_exhaustive_test.go b/go/internal/auth/classify_exhaustive_test.go index 1e6a23cab..9901db71c 100644 --- a/go/internal/auth/classify_exhaustive_test.go +++ b/go/internal/auth/classify_exhaustive_test.go @@ -13,6 +13,7 @@ import ( "google.golang.org/protobuf/reflect/protoregistry" compassv1 "github.com/RigelBuild/compass/go/gen/compass/v1" + compassv1internal "github.com/RigelBuild/compass/go/internal/gen/compass/v1" ) // procedurePath reconstructs the connect procedure path for a method descriptor. @@ -36,6 +37,17 @@ func gatedFileDescriptors() []protoreflect.FileDescriptor { } } +// ungatedFileDescriptors are compass.v1 service files never mounted behind +// AdminGate: Runner, agent-socket and guest-vsock surfaces with their own authz. +// Importing them keeps the registry guard below independent of the link set. +func ungatedFileDescriptors() []protoreflect.FileDescriptor { + return []protoreflect.FileDescriptor{ + compassv1internal.File_compass_v1_runner_proto, + compassv1internal.File_compass_v1_agent_gateway_proto, + compassv1internal.File_compass_v1_guest_control_proto, + } +} + // TestClassifyProcedureCoversEveryGeneratedProcedure fails if any generated // CompassService or CommsService procedure is not explicitly classified by // classifyProcedure. This is the build-time gate the doc comment promises: a new @@ -127,6 +139,12 @@ func TestClassificationGateCoversEveryRegisteredCompassService(t *testing.T) { covered[services.Get(si).FullName()] = true } } + for _, file := range ungatedFileDescriptors() { + services := file.Services() + for si := range services.Len() { + covered[services.Get(si).FullName()] = true + } + } if len(packages) == 0 { t.Fatal("gatedFileDescriptors is empty — the classification gate covers nothing") } @@ -139,7 +157,7 @@ func TestClassificationGateCoversEveryRegisteredCompassService(t *testing.T) { svc := services.Get(si) checked++ if !covered[svc.FullName()] { - t.Errorf("service %q is registered in proto package %q but its file is not in gatedFileDescriptors — add its File_..._proto to the slice so classifyProcedure's exhaustiveness gate covers its RPCs (otherwise they silently fail-closed to adminOnly on the network door)", svc.FullName(), pkg) + t.Errorf("service %q is registered in proto package %q but its file is in neither gatedFileDescriptors nor ungatedFileDescriptors — add it to the gated slice so classifyProcedure's exhaustiveness gate covers its RPCs (otherwise they silently fail-closed to adminOnly on the network door), or to the ungated slice if it is never mounted behind AdminGate", svc.FullName(), pkg) } } return true diff --git a/go/internal/auth/interceptor_pgtest_test.go b/go/internal/auth/interceptor_pgtest_test.go index 87ccfd0ee..53233b044 100644 --- a/go/internal/auth/interceptor_pgtest_test.go +++ b/go/internal/auth/interceptor_pgtest_test.go @@ -74,7 +74,7 @@ func TestBearerInterceptorSetsCommsActorNotAdminFallback(t *testing.T) { // diverts attribution to the real caller. commsBus := events.NewBus[*compassv1.SubscribeCommsResponse]() t.Cleanup(commsBus.Close) - commsSvc := comms.NewComms(st, commsBus, admin) + commsSvc := comms.NewComms(st, commsBus, nil, admin) // Drive CreateChannelGroup through the bearer interceptor: the interceptor // resolves the member token and threads the caller into the ctx it hands the diff --git a/go/internal/comms/agent_caller_pgtest_test.go b/go/internal/comms/agent_caller_pgtest_test.go index 411a0f127..cbe95ecac 100644 --- a/go/internal/comms/agent_caller_pgtest_test.go +++ b/go/internal/comms/agent_caller_pgtest_test.go @@ -160,7 +160,7 @@ func TestPostAsAccountEmptyAccountFailsClosedNoAdminWrite(t *testing.T) { if err != nil { t.Fatalf("BootstrapAdmin: %v", err) } - svc := NewComms(st, bus, admin.ID) + svc := NewComms(st, bus, nil, admin.ID) ctx := context.Background() // A channel the admin founds (so admin is a member) — the write target that diff --git a/go/internal/comms/comms.go b/go/internal/comms/comms.go index fcfcd32db..8da4ef414 100644 --- a/go/internal/comms/comms.go +++ b/go/internal/comms/comms.go @@ -18,14 +18,18 @@ package comms import ( "context" "errors" + "log/slog" "connectrpc.com/connect" + "go.opentelemetry.io/otel" "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/metric" "go.opentelemetry.io/otel/trace" "github.com/RigelBuild/compass/go/events" compassv1 "github.com/RigelBuild/compass/go/gen/compass/v1" "github.com/RigelBuild/compass/go/gen/compass/v1/compassv1connect" + "github.com/RigelBuild/compass/go/internal/fabric" "github.com/RigelBuild/compass/go/internal/store" ) @@ -37,13 +41,17 @@ import ( // edge stamps Seq/AtUnixMs/InstanceEpoch onto a copy. type commsBus = *events.Bus[*compassv1.SubscribeCommsResponse] +const instrumentationScope = "github.com/RigelBuild/compass/go/internal/comms" + // Comms implements compassv1connect.CommsServiceHandler over the store and the // comms event bus. Cheap to share by pointer; the store and bus are each safe // for concurrent use, and Comms holds no mutable state of its own — the store is // the source of truth, so there is no in-memory account/channel/group state. type Comms struct { - store *store.Store - bus commsBus + store *store.Store + bus commsBus + fabric fabric.EventFabric + fabricPublishFailures metric.Int64Counter // adminID attributes every RPC on the local-socket door (the door has no // interceptor yet). The T3 interceptor overrides this per-request by setting // a caller on the context; adminID is the fallback when none is set. @@ -59,9 +67,19 @@ type Comms struct { // NewComms constructs the CommsService handler over store and bus. adminID is the // bootstrap-admin account the local-socket door attributes callers to until the -// T3 interceptor sets a real identity (design.md:1219-1222). -func NewComms(st *store.Store, bus commsBus, adminID store.AccountID) *Comms { - return &Comms{store: st, bus: bus, adminID: adminID} +// T3 interceptor sets a real identity (design.md:1219-1222). A nil fab keeps +// message_posted on the bus only. +func NewComms(st *store.Store, bus commsBus, fab fabric.EventFabric, adminID store.AccountID) *Comms { + // A metric miss must never fail construction, matching the delivery counter. + failures, err := otel.Meter(instrumentationScope).Int64Counter( + "compass.delivery.fabric_publish_failures", + metric.WithDescription("Count of message_posted fabric publishes that failed after commit."), + ) + if err != nil { + slog.Warn("comms: failed to create fabric publish failure counter; metric disabled", "err", err) + failures = nil + } + return &Comms{store: st, bus: bus, fabric: fab, fabricPublishFailures: failures, adminID: adminID} } // SetPresenceSource wires the in-memory presence enum source GetRoster joins diff --git a/go/internal/comms/comms_test.go b/go/internal/comms/comms_test.go index caeddb695..646ad0dba 100644 --- a/go/internal/comms/comms_test.go +++ b/go/internal/comms/comms_test.go @@ -59,7 +59,7 @@ func newHandler(t *testing.T) (*Comms, *store.Store) { if err != nil { t.Fatalf("BootstrapAdmin: %v", err) } - return NewComms(st, bus, admin.ID), st + return NewComms(st, bus, nil, admin.ID), st } func TestCreateChannelEmitsChannelChanged(t *testing.T) { diff --git a/go/internal/comms/fabric_publish_pgtest_test.go b/go/internal/comms/fabric_publish_pgtest_test.go new file mode 100644 index 000000000..beb86f395 --- /dev/null +++ b/go/internal/comms/fabric_publish_pgtest_test.go @@ -0,0 +1,222 @@ +//go:build pgtest + +package comms + +import ( + "context" + "errors" + "sync" + "testing" + + "connectrpc.com/connect" + "go.opentelemetry.io/otel" + sdkmetric "go.opentelemetry.io/otel/sdk/metric" + "go.opentelemetry.io/otel/sdk/metric/metricdata" + + compassv1 "github.com/RigelBuild/compass/go/gen/compass/v1" + "github.com/RigelBuild/compass/go/internal/fabric" + "github.com/RigelBuild/compass/go/internal/store" +) + +// publishedRef is one Publish call the fake fabric recorded, plus whether the +// named row was already committed when it arrived. +type publishedRef struct { + subject string + ref fabric.EventRef + committed bool +} + +// fakeEventFabric records publishes and, for each one, re-reads the row so a +// test can prove the publish happened after the commit, as the consumer needs. +type fakeEventFabric struct { + st *store.Store + err error + + mu sync.Mutex + got []publishedRef +} + +func (f *fakeEventFabric) Publish(ctx context.Context, subject string, ref fabric.EventRef) error { + _, readErr := f.st.MessageByID(store.WithSystemRole(ctx), ref.RowID) + f.mu.Lock() + f.got = append(f.got, publishedRef{subject: subject, ref: ref, committed: readErr == nil}) + f.mu.Unlock() + return f.err +} + +func (f *fakeEventFabric) Subscribe(context.Context, string, func(fabric.EventRef)) (fabric.Unsubscribe, error) { + return nil, errors.New("fakeEventFabric: Subscribe not used by comms") +} + +func (f *fakeEventFabric) SubscribeKind(context.Context, fabric.EventKind, func(fabric.EventRef)) (fabric.Unsubscribe, error) { + return nil, errors.New("fakeEventFabric: SubscribeKind not used by comms") +} + +func (f *fakeEventFabric) published() []publishedRef { + f.mu.Lock() + defer f.mu.Unlock() + return append([]publishedRef(nil), f.got...) +} + +// newFabricHandler builds a handler over a real store with fab wired in. +func newFabricHandler(t *testing.T, fabErr error) (*Comms, *store.Store, *fakeEventFabric) { + t.Helper() + st := newTestStore(t) + admin, err := st.BootstrapAdmin(context.Background(), store.NewUser{Handle: "root", DisplayName: "Root"}) + if err != nil { + t.Fatalf("BootstrapAdmin: %v", err) + } + fab := &fakeEventFabric{st: st, err: fabErr} + return NewComms(st, newBus(t), fab, admin.ID), st, fab +} + +func postText(ctx context.Context, t *testing.T, svc *Comms, actor store.AccountID, ch store.ChannelID, body, requestID string) string { + t.Helper() + resp, err := svc.PostMessage(WithActor(ctx, actor), connect.NewRequest(&compassv1.PostMessageRequest{ + Container: &compassv1.PostMessageRequest_ChannelId{ChannelId: string(ch)}, + Topic: &compassv1.PostMessageRequest_TopicName{TopicName: "general"}, + CreateTopic: true, + Blocks: []*compassv1.MessageBlock{{Block: &compassv1.MessageBlock_Text{Text: body}}}, + ClientRequestId: requestID, + })) + if err != nil { + t.Fatalf("PostMessage(%q): %v", body, err) + } + return resp.Msg.GetMessage().GetId() +} + +func assertOnePostedRef(t *testing.T, st *store.Store, fab *fakeEventFabric, wantRowID string) { + t.Helper() + got := fab.published() + if len(got) != 1 { + t.Fatalf("fabric received %d publishes, want exactly 1: %+v", len(got), got) + } + tenant := string(st.EffectiveTenant(context.Background())) + wantSubject, err := fabric.CommsSubject(tenant, fabric.KindMessagePosted) + if err != nil { + t.Fatalf("CommsSubject: %v", err) + } + want := fabric.EventRef{Tenant: tenant, Kind: fabric.KindMessagePosted, RowID: wantRowID} + if got[0].ref != want || got[0].subject != wantSubject { + t.Fatalf("published %q %+v, want %q %+v", got[0].subject, got[0].ref, wantSubject, want) + } + if !got[0].committed { + t.Fatal("fabric publish arrived before the message row was committed") + } +} + +func TestPostMessagePublishesPostedRefOnFabric(t *testing.T) { + svc, st, fab := newFabricHandler(t, nil) + ctx := context.Background() + poster := mustUser(t, st, "poster") + ch, err := st.CreateChannel(ctx, poster.ID, store.NewChannel{Name: "room", Kind: store.ChannelKindChannel}) + if err != nil { + t.Fatalf("CreateChannel: %v", err) + } + + id := postText(ctx, t, svc, poster.ID, ch.ID, "hello", "") + assertOnePostedRef(t, st, fab, id) +} + +func TestPostMessageIdempotentRetryDoesNotRepublishOnFabric(t *testing.T) { + svc, st, fab := newFabricHandler(t, nil) + ctx := context.Background() + poster := mustUser(t, st, "poster") + ch, err := st.CreateChannel(ctx, poster.ID, store.NewChannel{Name: "room", Kind: store.ChannelKindChannel}) + if err != nil { + t.Fatalf("CreateChannel: %v", err) + } + + id := postText(ctx, t, svc, poster.ID, ch.ID, "once", "req-dup") + if retry := postText(ctx, t, svc, poster.ID, ch.ID, "again", "req-dup"); retry != id { + t.Fatalf("retry returned %q, want the stored %q", retry, id) + } + assertOnePostedRef(t, st, fab, id) +} + +func TestPostMessageSurvivesFabricPublishFailureAndCountsIt(t *testing.T) { + reader := sdkmetric.NewManualReader() + mp := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader)) + t.Cleanup(func() { + if err := mp.Shutdown(context.Background()); err != nil { + t.Errorf("meter provider shutdown: %v", err) + } + }) + prev := otel.GetMeterProvider() + t.Cleanup(func() { otel.SetMeterProvider(prev) }) + otel.SetMeterProvider(mp) + + svc, st, fab := newFabricHandler(t, errors.New("nats unavailable")) + ctx := context.Background() + poster := mustUser(t, st, "poster") + ch, err := st.CreateChannel(ctx, poster.ID, store.NewChannel{Name: "room", Kind: store.ChannelKindChannel}) + if err != nil { + t.Fatalf("CreateChannel: %v", err) + } + + // postText fails the test if the RPC returns an error, which is the point. + id := postText(ctx, t, svc, poster.ID, ch.ID, "committed anyway", "") + assertOnePostedRef(t, st, fab, id) + + var rm metricdata.ResourceMetrics + if err := reader.Collect(ctx, &rm); err != nil { + t.Fatalf("collect metrics: %v", err) + } + var total int64 + for _, sm := range rm.ScopeMetrics { + for _, m := range sm.Metrics { + if m.Name != "compass.delivery.fabric_publish_failures" { + continue + } + sum, ok := m.Data.(metricdata.Sum[int64]) + if !ok { + t.Fatalf("fabric_publish_failures data = %T, want Sum[int64]", m.Data) + } + for _, dp := range sum.DataPoints { + total += dp.Value + } + } + } + if total != 1 { + t.Fatalf("compass.delivery.fabric_publish_failures = %d, want 1", total) + } +} + +func TestRespondToAskPublishesAnswerRefOnFabric(t *testing.T) { + svc, st, fab := newFabricHandler(t, nil) + ctx := context.Background() + agent := mustUser(t, st, "agent") + ch, err := st.CreateChannel(ctx, agent.ID, store.NewChannel{Name: "room", Kind: store.ChannelKindChannel}) + if err != nil { + t.Fatalf("CreateChannel: %v", err) + } + askMsg, _, err := st.AppendMessage(ctx, store.Message{AuthorAccountID: agent.ID, Blocks: []store.MessageBlock{pendingAskStore("ask-1")}}, string(ch.ID), store.TopicRef{Name: "general", Create: true}, "") + if err != nil { + t.Fatalf("AppendMessage: %v", err) + } + if got := fab.published(); len(got) != 0 { + t.Fatalf("a store-level append published %d refs, want 0", len(got)) + } + + if _, err := svc.RespondToAsk(WithActor(ctx, agent.ID), connect.NewRequest(&compassv1.RespondToAskRequest{ + AskId: "ask-1", + Answers: []*compassv1.AskQuestionAnswer{{QuestionId: "q1", ChosenOptionIds: []string{"opt-a"}}}, + })); err != nil { + t.Fatalf("RespondToAsk: %v", err) + } + got := fab.published() + if len(got) != 1 { + t.Fatalf("fabric received %d publishes, want exactly 1 for the answer", len(got)) + } + if got[0].ref.RowID == string(askMsg.ID) { + t.Fatal("published the ask message, want the new answer message") + } + answer, err := st.MessageByID(store.WithSystemRole(ctx), got[0].ref.RowID) + if err != nil { + t.Fatalf("published row %q is not a stored message: %v", got[0].ref.RowID, err) + } + if len(answer.Blocks) != 1 || answer.Blocks[0].AskAnswer == nil { + t.Fatalf("published row %q is not an ask_answer message: %+v", answer.ID, answer.Blocks) + } + assertOnePostedRef(t, st, fab, string(answer.ID)) +} diff --git a/go/internal/comms/mapping.go b/go/internal/comms/mapping.go index 3e4e823c6..5699a0eb8 100644 --- a/go/internal/comms/mapping.go +++ b/go/internal/comms/mapping.go @@ -2,10 +2,12 @@ package comms import ( "context" + "log/slog" "connectrpc.com/connect" compassv1 "github.com/RigelBuild/compass/go/gen/compass/v1" + "github.com/RigelBuild/compass/go/internal/fabric" "github.com/RigelBuild/compass/go/internal/store" ) @@ -501,6 +503,23 @@ func (c *Comms) publishMessagePosted(ctx context.Context, m store.Message) { MessagePosted: &compassv1.MessagePosted{Message: MessageToWire(m)}, }, }) + if c.fabric == nil { + return + } + tenant := string(c.store.EffectiveTenant(ctx)) + ref := fabric.EventRef{Tenant: tenant, Kind: fabric.KindMessagePosted, RowID: string(m.ID)} + subject, err := fabric.CommsSubject(tenant, fabric.KindMessagePosted) + if err == nil { + err = c.fabric.Publish(ctx, subject, ref) + } + // The row is already committed, so failing the RPC would report a persisted + // write as lost; the delivery recovery sweep owns redelivery instead. + if err != nil { + slog.ErrorContext(ctx, "comms: publishing message_posted to fabric failed", "error", err, "message_id", string(m.ID)) + if c.fabricPublishFailures != nil { + c.fabricPublishFailures.Add(ctx, 1) + } + } } func (c *Comms) publishMessageUpdated(m store.Message) { diff --git a/go/internal/comms/roster_pgtest_test.go b/go/internal/comms/roster_pgtest_test.go index 2676d3bf8..9d0675b02 100644 --- a/go/internal/comms/roster_pgtest_test.go +++ b/go/internal/comms/roster_pgtest_test.go @@ -295,7 +295,7 @@ func TestGetRosterActivitySurvivesSimulatedRestart(t *testing.T) { // Simulate the restart: a brand-new handler over the SAME store, its in-memory // presence projection empty (nothing re-enrolled yet). - fresh := NewComms(st, newBus(t), owner.ID) + fresh := NewComms(st, newBus(t), nil, owner.ID) resp, err := fresh.GetRoster(WithActor(ctx, owner.ID), connect.NewRequest(&compassv1.GetRosterRequest{ Scope: compassv1.RosterScope_ROSTER_SCOPE_OWNER, diff --git a/go/internal/comms/subscribe_failclosed_test.go b/go/internal/comms/subscribe_failclosed_test.go index 029286146..159e00194 100644 --- a/go/internal/comms/subscribe_failclosed_test.go +++ b/go/internal/comms/subscribe_failclosed_test.go @@ -455,7 +455,7 @@ func TestForwardCommsLiveTailOverrunEmitsTerminalResync(t *testing.T) { // surfaces as a whole-suite timeout instead of this test failing. func driveUnderflowResync(t *testing.T, bus *events.Bus[*compassv1.SubscribeCommsResponse], req *compassv1.SubscribeCommsRequest) forwardResult { t.Helper() - svc := NewComms(nil, bus, testActor) + svc := NewComms(nil, bus, nil, testActor) path, handler := compassv1connect.NewCommsServiceHandler(svc) mux := http.NewServeMux() diff --git a/go/internal/comms/subscribe_test.go b/go/internal/comms/subscribe_test.go index ae0a01ea4..735a6d4f9 100644 --- a/go/internal/comms/subscribe_test.go +++ b/go/internal/comms/subscribe_test.go @@ -70,7 +70,7 @@ func newStreamHarness(t *testing.T) streamHarness { if err != nil { t.Fatalf("BootstrapAdmin: %v", err) } - svc := NewComms(st, bus, admin.ID) + svc := NewComms(st, bus, nil, admin.ID) path, handler := compassv1connect.NewCommsServiceHandler(svc) mux := http.NewServeMux() diff --git a/go/internal/runnerhub/integration_pgtest_test.go b/go/internal/runnerhub/integration_pgtest_test.go index c6b6040ff..f67469b59 100644 --- a/go/internal/runnerhub/integration_pgtest_test.go +++ b/go/internal/runnerhub/integration_pgtest_test.go @@ -320,7 +320,7 @@ func openStoreFixture(t *testing.T, ctx context.Context, dsn string) (*store.Sto // adminID is the comms handler's ambient fallback; the agent-initiated leg // (PostAsAccount) overrides it per-call with the resolved agent account, so // the post attributes to the agent, never the admin. - commsSvc := comms.NewComms(st, bus, admin.ID) + commsSvc := comms.NewComms(st, bus, nil, admin.ID) sub, err := bus.Subscribe(0, 0) if err != nil { t.Fatalf("bus.Subscribe: %v", err) diff --git a/go/server/cors_pgtest_test.go b/go/server/cors_pgtest_test.go index 19fb19a05..568cdf3cd 100644 --- a/go/server/cors_pgtest_test.go +++ b/go/server/cors_pgtest_test.go @@ -47,7 +47,7 @@ func buildDoorHandler(t *testing.T, corsOrigin string) http.Handler { commsBus := events.NewBus[*compassv1.SubscribeCommsResponse]() t.Cleanup(commsBus.Close) - commsSvc := comms.NewComms(st, commsBus, admin) + commsSvc := comms.NewComms(st, commsBus, nil, admin) cfg := ServeConfig{ StateDir: t.TempDir(), // bootstrap-admin token file lands here (0600) diff --git a/go/server/dm_e2e_pgtest_test.go b/go/server/dm_e2e_pgtest_test.go index b201ad466..063835003 100644 --- a/go/server/dm_e2e_pgtest_test.go +++ b/go/server/dm_e2e_pgtest_test.go @@ -76,7 +76,7 @@ func newDME2EWire(t *testing.T) *dmE2EWire { commsBus := events.NewBus[*compassv1.SubscribeCommsResponse]() t.Cleanup(commsBus.Close) - commsSvc := comms.NewComms(st, commsBus, admin.ID) + commsSvc := comms.NewComms(st, commsBus, nil, admin.ID) // The production delivery wire (sinks.go:142-155), assembled inline with the // REAL resume-based waker whose dm opener is the comms service — so the spawn diff --git a/go/server/lifecycle_e2e_pgtest_test.go b/go/server/lifecycle_e2e_pgtest_test.go index 86a78157a..1d608d055 100644 --- a/go/server/lifecycle_e2e_pgtest_test.go +++ b/go/server/lifecycle_e2e_pgtest_test.go @@ -477,7 +477,7 @@ func newE2EWire(t *testing.T) *e2eWire { brd := board.NewProjection(bus) commsBus := events.NewBus[*compassv1.SubscribeCommsResponse]() t.Cleanup(commsBus.Close) - commsSvc := comms.NewComms(st, commsBus, admin.ID) + commsSvc := comms.NewComms(st, commsBus, nil, admin.ID) hub := newRunnerHub(st, brd, newSessionTail(), commsSvc, discardLogE2E()) // Wire the lifecycleService as the hub's LifecycleCaller — the serve.go:250 diff --git a/go/server/lifecycle_pgtest_test.go b/go/server/lifecycle_pgtest_test.go index 1d1c17ad3..e2c7e6d35 100644 --- a/go/server/lifecycle_pgtest_test.go +++ b/go/server/lifecycle_pgtest_test.go @@ -47,7 +47,7 @@ func newLifecycleFixture(t *testing.T) lifecycleFixture { // is unused — OpenDMAsAccount always sets the caller explicitly). commsBus := events.NewBus[*compassv1.SubscribeCommsResponse]() t.Cleanup(commsBus.Close) - commsSvc := comms.NewComms(pf.store, commsBus, owner) + commsSvc := comms.NewComms(pf.store, commsBus, nil, owner) return lifecycleFixture{ placementFixture: pf, lc: newLifecycleService(pf.store, pf.hub, commsSvc), diff --git a/go/server/offline_mention_e2e_pgtest_test.go b/go/server/offline_mention_e2e_pgtest_test.go index c516d699f..0249a20cc 100644 --- a/go/server/offline_mention_e2e_pgtest_test.go +++ b/go/server/offline_mention_e2e_pgtest_test.go @@ -77,7 +77,7 @@ func newMentionE2EWire(t *testing.T) *mentionE2EWire { // A real caller (not nil) makes the agent-authored leg drivable. commsBus := events.NewBus[*compassv1.SubscribeCommsResponse]() t.Cleanup(commsBus.Close) - commsSvc := comms.NewComms(st, commsBus, admin.ID) + commsSvc := comms.NewComms(st, commsBus, nil, admin.ID) // The hub over a discard board + tail — otherwise the same shape // newPlacementFixtureWith builds. diff --git a/go/server/otel_emission_pgtest_test.go b/go/server/otel_emission_pgtest_test.go index 8ffbf253d..0f547fb8a 100644 --- a/go/server/otel_emission_pgtest_test.go +++ b/go/server/otel_emission_pgtest_test.go @@ -192,7 +192,7 @@ func TestNetworkDoorExposesTraceResponseHeader(t *testing.T) { svc := newService("otel-cors-test", bus, st, nil, nil, nil, nil) commsBus := events.NewBus[*compassv1.SubscribeCommsResponse]() t.Cleanup(commsBus.Close) - commsSvc := comms.NewComms(st, commsBus, admin) + commsSvc := comms.NewComms(st, commsBus, nil, admin) secretsSvc := newSecretsService(st, nil, nil, nil) otelIC, err := otelconnect.NewInterceptor() if err != nil { @@ -241,7 +241,7 @@ func TestNetworkDoorAllowsPostHogSessionRequestHeader(t *testing.T) { svc := newService("otel-cors-session-test", bus, st, nil, nil, nil, nil) commsBus := events.NewBus[*compassv1.SubscribeCommsResponse]() t.Cleanup(commsBus.Close) - commsSvc := comms.NewComms(st, commsBus, admin) + commsSvc := comms.NewComms(st, commsBus, nil, admin) secretsSvc := newSecretsService(st, nil, nil, nil) otelIC, err := otelconnect.NewInterceptor() if err != nil { @@ -309,7 +309,7 @@ func TestNetworkDoorStampsPostHogSessionIDOnTheSpan(t *testing.T) { svc := newService("otel-session-test", bus, st, nil, nil, nil, nil) commsBus := events.NewBus[*compassv1.SubscribeCommsResponse]() t.Cleanup(commsBus.Close) - commsSvc := comms.NewComms(st, commsBus, admin) + commsSvc := comms.NewComms(st, commsBus, nil, admin) secretsSvc := newSecretsService(st, nil, nil, nil) otelIC, err := otelconnect.NewInterceptor() if err != nil { diff --git a/go/server/serve.go b/go/server/serve.go index 0faa508f9..bd215d823 100644 --- a/go/server/serve.go +++ b/go/server/serve.go @@ -788,7 +788,7 @@ func Serve(ctx context.Context, cfg ServeConfig) error { // agent-initiated comms calls through this handler (the CommsCaller). commsBus := events.NewBus[*compassv1.SubscribeCommsResponse]() defer commsBus.Close() - commsSvc := comms.NewComms(st, commsBus, admin.ID) + commsSvc := comms.NewComms(st, commsBus, nil, admin.ID) // Register the coordination-channel reconcile as the store's in-tx hook, so // the two parent-edge writers auto-provision/reconcile a manager's // coordination channel atomically with the tree edge (RIG-1722 T5). Wired here diff --git a/go/server/serve_forge_pgtest_test.go b/go/server/serve_forge_pgtest_test.go index 5dbe654d2..2ea37c02b 100644 --- a/go/server/serve_forge_pgtest_test.go +++ b/go/server/serve_forge_pgtest_test.go @@ -518,7 +518,7 @@ func TestBuildDoorsRoutesTheResolverInstancesOverTheRealCallGraph(t *testing.T) if err != nil { t.Fatalf("CreateUser: %v", err) } - commsSvc := comms.NewComms(st, commsBus, admin.ID) + commsSvc := comms.NewComms(st, commsBus, nil, admin.ID) brd := board.NewProjection(bus) issueBrd := board.NewIssueProjection(bus, st) tail := newSessionTail() diff --git a/go/server/serve_seed_pgtest_test.go b/go/server/serve_seed_pgtest_test.go index e3d8d5542..ec2cd4207 100644 --- a/go/server/serve_seed_pgtest_test.go +++ b/go/server/serve_seed_pgtest_test.go @@ -83,7 +83,7 @@ func newSeedHarness(t *testing.T) *seedHarness { tail := newSessionTail() commsBus := events.NewBus[*compassv1.SubscribeCommsResponse]() t.Cleanup(commsBus.Close) - commsSvc := comms.NewComms(st, commsBus, admin.ID) + commsSvc := comms.NewComms(st, commsBus, nil, admin.ID) hub := newRunnerHub(st, brd, tail, commsSvc, slog.New(slog.DiscardHandler)) svc := newService("test", bus, st, hub, brd, nil, tail) From d62502fe7000be20cf3c37bb537fd8cc43225cef Mon Sep 17 00:00:00 2001 From: mintaka Date: Fri, 25 Sep 2026 15:03:13 -0400 Subject: [PATCH 2/4] fix(comms): detach the fabric publish from request cancellation (RIG-3107) Review of the T1 publish found a lost-trigger window: - The publish ran on the request ctx. A client hanging up after commit canceled it, and an idempotent retry never republishes, so the post waited for the floor sweep. It now runs on a WithoutCancel ctx bounded by a 5s timeout. - New tests pin a non-bootstrap request tenant to its subject, and pin the publish surviving a cancel mid-ack. The fake keeps its row-read error instead of folding it into a bool. - The design record's T1 interface named Store.ResolveTenant; it now names Store.EffectiveTenant, which landed with the same semantics. Spec-impact: none. Refs RIG-3107 Co-authored-by: Matt Wilkinson --- .../design.md | 10 +- .../comms/fabric_publish_pgtest_test.go | 114 +++++++++++++++--- go/internal/comms/mapping.go | 11 +- go/internal/store/tenant.go | 10 +- 4 files changed, 115 insertions(+), 30 deletions(-) diff --git a/docs/designs/infra/runtime/compass-managed-delivery-cutover/design.md b/docs/designs/infra/runtime/compass-managed-delivery-cutover/design.md index 664d25049..11ffc6f2a 100644 --- a/docs/designs/infra/runtime/compass-managed-delivery-cutover/design.md +++ b/docs/designs/infra/runtime/compass-managed-delivery-cutover/design.md @@ -274,10 +274,10 @@ resolve the tenant from ctx with the bootstrap fallback. and `CommsSubject(tenant string, kind EventKind) (string, error)` + `fabric.KindMessagePosted` (PR3); `store.TenantFromContext(ctx context.Context) (store.TenantID, bool)` - (`go/internal/store/context.go:23`). Produces: a new exported - `func (s *Store) ResolveTenant(ctx context.Context) TenantID` (promoting the - unexported `resolveTenant`, `go/internal/store/tenant.go:54-58`, so comms - gets the same set-or-bootstrap-fallback semantics the store's writes use); + (`go/internal/store/context.go:23`). Produces: comms resolves the tenant + through the exported `func (s *Store) EffectiveTenant(ctx context.Context) + TenantID` (which RIG-3108 added after this record froze, with the same + set-or-bootstrap-fallback semantics as the unexported `resolveTenant`); `comms.NewComms` gains a `fabric fabric.EventFabric` parameter (nil-safe: nil ⇒ bus-only, so unit tests and any not-yet-wired assembly keep working); `publishMessagePosted(ctx, m)` extended with the fabric publish + the @@ -448,7 +448,7 @@ exists anymore) and re-derive the no-loss argument from JetStream durability. ## Tasks -- [ ] T1: fabric publish in `publishMessagePosted` + `Store.ResolveTenant` + +- [ ] T1: fabric publish in `publishMessagePosted` + `Store.EffectiveTenant` + failure counter (tests a–d) - [ ] T2: consumer trigger cutover — `SubscribeKind` in, bus tail out, per OQ-1/OQ-2/OQ-3 rulings; fabric serial-callback contract doc + no-overlap diff --git a/go/internal/comms/fabric_publish_pgtest_test.go b/go/internal/comms/fabric_publish_pgtest_test.go index beb86f395..fcf4a0873 100644 --- a/go/internal/comms/fabric_publish_pgtest_test.go +++ b/go/internal/comms/fabric_publish_pgtest_test.go @@ -7,8 +7,10 @@ import ( "errors" "sync" "testing" + "time" "connectrpc.com/connect" + "github.com/jackc/pgx/v5" "go.opentelemetry.io/otel" sdkmetric "go.opentelemetry.io/otel/sdk/metric" "go.opentelemetry.io/otel/sdk/metric/metricdata" @@ -18,12 +20,13 @@ import ( "github.com/RigelBuild/compass/go/internal/store" ) -// publishedRef is one Publish call the fake fabric recorded, plus whether the -// named row was already committed when it arrived. +// publishedRef is one Publish call the fake fabric recorded, plus the row read +// it made during the call and the call's context state. type publishedRef struct { - subject string - ref fabric.EventRef - committed bool + subject string + ref fabric.EventRef + readErr error + ctxErr error } // fakeEventFabric records publishes and, for each one, re-reads the row so a @@ -31,15 +34,21 @@ type publishedRef struct { type fakeEventFabric struct { st *store.Store err error + // onPublish runs first in Publish, standing in for the client hanging up + // while the JetStream ack is in flight. + onPublish func() mu sync.Mutex got []publishedRef } func (f *fakeEventFabric) Publish(ctx context.Context, subject string, ref fabric.EventRef) error { + if f.onPublish != nil { + f.onPublish() + } _, readErr := f.st.MessageByID(store.WithSystemRole(ctx), ref.RowID) f.mu.Lock() - f.got = append(f.got, publishedRef{subject: subject, ref: ref, committed: readErr == nil}) + f.got = append(f.got, publishedRef{subject: subject, ref: ref, readErr: readErr, ctxErr: ctx.Err()}) f.mu.Unlock() return f.err } @@ -61,13 +70,43 @@ func (f *fakeEventFabric) published() []publishedRef { // newFabricHandler builds a handler over a real store with fab wired in. func newFabricHandler(t *testing.T, fabErr error) (*Comms, *store.Store, *fakeEventFabric) { t.Helper() - st := newTestStore(t) + svc, st, fab, _ := newFabricHandlerDSN(t, fabErr) + return svc, st, fab +} + +func newFabricHandlerDSN(t *testing.T, fabErr error) (*Comms, *store.Store, *fakeEventFabric, string) { + t.Helper() + st, dsn := newTestStoreDSN(t) admin, err := st.BootstrapAdmin(context.Background(), store.NewUser{Handle: "root", DisplayName: "Root"}) if err != nil { t.Fatalf("BootstrapAdmin: %v", err) } fab := &fakeEventFabric{st: st, err: fabErr} - return NewComms(st, newBus(t), fab, admin.ID), st, fab + return NewComms(st, newBus(t), fab, admin.ID), st, fab, dsn +} + +// seedTenant inserts a non-bootstrap tenant row. The store exposes no tenant +// create, and the tenants table is RLS-exempt, so a direct insert is enough. +func seedTenant(t *testing.T, dsn, slug string) store.TenantID { + t.Helper() + ctx := context.Background() + conn, err := pgx.Connect(ctx, dsn) + if err != nil { + t.Fatalf("connect to seed tenant: %v", err) + } + defer func() { + if err := conn.Close(ctx); err != nil { + t.Errorf("close seed conn: %v", err) + } + }() + id := "tenant-" + slug + if _, err := conn.Exec(ctx, + "INSERT INTO tenants (id, slug, display_name, created_at_unix_ms) VALUES ($1, $2, $3, $4)", + id, slug, slug, time.Now().UnixMilli(), + ); err != nil { + t.Fatalf("seed tenant %q: %v", slug, err) + } + return store.TenantID(id) } func postText(ctx context.Context, t *testing.T, svc *Comms, actor store.AccountID, ch store.ChannelID, body, requestID string) string { @@ -85,26 +124,33 @@ func postText(ctx context.Context, t *testing.T, svc *Comms, actor store.Account return resp.Msg.GetMessage().GetId() } -func assertOnePostedRef(t *testing.T, st *store.Store, fab *fakeEventFabric, wantRowID string) { +func assertOnePostedRef(t *testing.T, fab *fakeEventFabric, tenant store.TenantID, wantRowID string) { t.Helper() got := fab.published() if len(got) != 1 { t.Fatalf("fabric received %d publishes, want exactly 1: %+v", len(got), got) } - tenant := string(st.EffectiveTenant(context.Background())) - wantSubject, err := fabric.CommsSubject(tenant, fabric.KindMessagePosted) + wantSubject, err := fabric.CommsSubject(string(tenant), fabric.KindMessagePosted) if err != nil { t.Fatalf("CommsSubject: %v", err) } - want := fabric.EventRef{Tenant: tenant, Kind: fabric.KindMessagePosted, RowID: wantRowID} + want := fabric.EventRef{Tenant: string(tenant), Kind: fabric.KindMessagePosted, RowID: wantRowID} if got[0].ref != want || got[0].subject != wantSubject { t.Fatalf("published %q %+v, want %q %+v", got[0].subject, got[0].ref, wantSubject, want) } - if !got[0].committed { - t.Fatal("fabric publish arrived before the message row was committed") + if got[0].readErr != nil { + t.Fatalf("reading the published row during Publish (was it committed?): %v", got[0].readErr) + } + if got[0].ctxErr != nil { + t.Fatalf("fabric publish ran on a done context: %v", got[0].ctxErr) } } +// bootstrapTenant is the tenant a request without a tenant context resolves to. +func bootstrapTenant(st *store.Store) store.TenantID { + return st.EffectiveTenant(context.Background()) +} + func TestPostMessagePublishesPostedRefOnFabric(t *testing.T) { svc, st, fab := newFabricHandler(t, nil) ctx := context.Background() @@ -115,7 +161,39 @@ func TestPostMessagePublishesPostedRefOnFabric(t *testing.T) { } id := postText(ctx, t, svc, poster.ID, ch.ID, "hello", "") - assertOnePostedRef(t, st, fab, id) + assertOnePostedRef(t, fab, bootstrapTenant(st), id) +} + +func TestPostMessagePublishesOnTheRequestTenantSubject(t *testing.T) { + svc, st, fab, dsn := newFabricHandlerDSN(t, nil) + tenant := seedTenant(t, dsn, "other") + ctx := store.WithTenant(context.Background(), tenant) + poster, err := st.CreateUser(ctx, store.NewUser{Handle: "poster", DisplayName: "poster"}) + if err != nil { + t.Fatalf("CreateUser: %v", err) + } + ch, err := st.CreateChannel(ctx, poster.ID, store.NewChannel{Name: "room", Kind: store.ChannelKindChannel}) + if err != nil { + t.Fatalf("CreateChannel: %v", err) + } + + id := postText(ctx, t, svc, poster.ID, ch.ID, "hello", "") + assertOnePostedRef(t, fab, tenant, id) +} + +func TestPostMessageFabricPublishOutlivesRequestCancellation(t *testing.T) { + svc, st, fab := newFabricHandler(t, nil) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + fab.onPublish = cancel + poster := mustUser(t, st, "poster") + ch, err := st.CreateChannel(ctx, poster.ID, store.NewChannel{Name: "room", Kind: store.ChannelKindChannel}) + if err != nil { + t.Fatalf("CreateChannel: %v", err) + } + + id := postText(ctx, t, svc, poster.ID, ch.ID, "hello", "") + assertOnePostedRef(t, fab, bootstrapTenant(st), id) } func TestPostMessageIdempotentRetryDoesNotRepublishOnFabric(t *testing.T) { @@ -131,7 +209,7 @@ func TestPostMessageIdempotentRetryDoesNotRepublishOnFabric(t *testing.T) { if retry := postText(ctx, t, svc, poster.ID, ch.ID, "again", "req-dup"); retry != id { t.Fatalf("retry returned %q, want the stored %q", retry, id) } - assertOnePostedRef(t, st, fab, id) + assertOnePostedRef(t, fab, bootstrapTenant(st), id) } func TestPostMessageSurvivesFabricPublishFailureAndCountsIt(t *testing.T) { @@ -156,7 +234,7 @@ func TestPostMessageSurvivesFabricPublishFailureAndCountsIt(t *testing.T) { // postText fails the test if the RPC returns an error, which is the point. id := postText(ctx, t, svc, poster.ID, ch.ID, "committed anyway", "") - assertOnePostedRef(t, st, fab, id) + assertOnePostedRef(t, fab, bootstrapTenant(st), id) var rm metricdata.ResourceMetrics if err := reader.Collect(ctx, &rm); err != nil { @@ -218,5 +296,5 @@ func TestRespondToAskPublishesAnswerRefOnFabric(t *testing.T) { if len(answer.Blocks) != 1 || answer.Blocks[0].AskAnswer == nil { t.Fatalf("published row %q is not an ask_answer message: %+v", answer.ID, answer.Blocks) } - assertOnePostedRef(t, st, fab, string(answer.ID)) + assertOnePostedRef(t, fab, bootstrapTenant(st), string(answer.ID)) } diff --git a/go/internal/comms/mapping.go b/go/internal/comms/mapping.go index 5699a0eb8..e8669e076 100644 --- a/go/internal/comms/mapping.go +++ b/go/internal/comms/mapping.go @@ -3,6 +3,7 @@ package comms import ( "context" "log/slog" + "time" "connectrpc.com/connect" @@ -497,6 +498,10 @@ func (c *Comms) publishAgentWorkspaceChanged(w store.AgentWorkspace) { }) } +// fabricPublishTimeout bounds the post-commit publish once it no longer follows +// the request's cancellation, so a stalled JetStream ack cannot pin the RPC. +const fabricPublishTimeout = 5 * time.Second + func (c *Comms) publishMessagePosted(ctx context.Context, m store.Message) { c.bus.PublishCtx(ctx, &compassv1.SubscribeCommsResponse{ Payload: &compassv1.SubscribeCommsResponse_MessagePosted{ @@ -510,7 +515,11 @@ func (c *Comms) publishMessagePosted(ctx context.Context, m store.Message) { ref := fabric.EventRef{Tenant: tenant, Kind: fabric.KindMessagePosted, RowID: string(m.ID)} subject, err := fabric.CommsSubject(tenant, fabric.KindMessagePosted) if err == nil { - err = c.fabric.Publish(ctx, subject, ref) + // A retry of this post is idempotent and never republishes, so a client + // hanging up after commit must not cancel the only publish attempt. + pubCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), fabricPublishTimeout) + err = c.fabric.Publish(pubCtx, subject, ref) + cancel() } // The row is already committed, so failing the RPC would report a persisted // write as lost; the delivery recovery sweep owns redelivery instead. diff --git a/go/internal/store/tenant.go b/go/internal/store/tenant.go index 09c520cf8..d5a3460ab 100644 --- a/go/internal/store/tenant.go +++ b/go/internal/store/tenant.go @@ -61,12 +61,10 @@ func (s *Store) resolveTenant(ctx context.Context) TenantID { // EffectiveTenant returns the tenant a request-scoped call resolves against: // the context tenant if the auth layer set one, else the bootstrap tenant (the // OSS single-tenant degenerate path). It is the exported form of resolveTenant, -// for a caller that must name the tenant a binding was written under — the -// RIG-3108 hub, which publishes a BindingChange on the per-tenant routing -// subject after recording a binding on the request ctx. A system-role ctx -// carries no tenant, so the bootstrap fallback applies there too; the hub never -// records a binding under the system role, so that degenerate value is never -// published. +// for a caller outside the store that must name the tenant a write landed under: +// the runner hub's BindingChange and the comms fabric publish, both keyed on a +// per-tenant subject. A system-role ctx carries no tenant, so it resolves to the +// bootstrap tenant; neither caller writes under the system role. func (s *Store) EffectiveTenant(ctx context.Context) TenantID { return s.resolveTenant(ctx) } From f7c9fa4713cf63639efb6bd0923a71d58946f010 Mon Sep 17 00:00:00 2001 From: mintaka Date: Fri, 25 Sep 2026 15:39:26 -0400 Subject: [PATCH 3/4] test(comms): pin the fabric publish timeout bound (RIG-3107) The cancellation test proved the publish is detached from the request, but removing the 5s timeout left it green. The fake now records the publish deadline, and the test fails when none is set. Spec-impact: none. Refs RIG-3107 Co-authored-by: Matt Wilkinson --- .../comms/fabric_publish_pgtest_test.go | 24 +++++++++++++++---- 1 file changed, 19 insertions(+), 5 deletions(-) diff --git a/go/internal/comms/fabric_publish_pgtest_test.go b/go/internal/comms/fabric_publish_pgtest_test.go index fcf4a0873..50df37703 100644 --- a/go/internal/comms/fabric_publish_pgtest_test.go +++ b/go/internal/comms/fabric_publish_pgtest_test.go @@ -23,10 +23,12 @@ import ( // publishedRef is one Publish call the fake fabric recorded, plus the row read // it made during the call and the call's context state. type publishedRef struct { - subject string - ref fabric.EventRef - readErr error - ctxErr error + subject string + ref fabric.EventRef + readErr error + ctxErr error + deadline time.Time + hasDeadline bool } // fakeEventFabric records publishes and, for each one, re-reads the row so a @@ -47,8 +49,12 @@ func (f *fakeEventFabric) Publish(ctx context.Context, subject string, ref fabri f.onPublish() } _, readErr := f.st.MessageByID(store.WithSystemRole(ctx), ref.RowID) + deadline, hasDeadline := ctx.Deadline() f.mu.Lock() - f.got = append(f.got, publishedRef{subject: subject, ref: ref, readErr: readErr, ctxErr: ctx.Err()}) + f.got = append(f.got, publishedRef{ + subject: subject, ref: ref, readErr: readErr, ctxErr: ctx.Err(), + deadline: deadline, hasDeadline: hasDeadline, + }) f.mu.Unlock() return f.err } @@ -192,8 +198,16 @@ func TestPostMessageFabricPublishOutlivesRequestCancellation(t *testing.T) { t.Fatalf("CreateChannel: %v", err) } + before := time.Now() id := postText(ctx, t, svc, poster.ID, ch.ID, "hello", "") assertOnePostedRef(t, fab, bootstrapTenant(st), id) + + // Detached from cancellation, the publish still needs its own bound, or a + // stalled JetStream ack would pin the RPC. + got := fab.published()[0] + if !got.hasDeadline || got.deadline.After(time.Now().Add(fabricPublishTimeout)) || got.deadline.Before(before) { + t.Fatalf("publish deadline = %v (set %v), want within %v of the call", got.deadline, got.hasDeadline, fabricPublishTimeout) + } } func TestPostMessageIdempotentRetryDoesNotRepublishOnFabric(t *testing.T) { From e6c0ca2b30f4479beb14da653d645df55e380c9e Mon Sep 17 00:00:00 2001 From: mintaka Date: Fri, 25 Sep 2026 15:46:06 -0400 Subject: [PATCH 4/4] test(comms): bound the publish deadline from both sides (RIG-3107) The deadline check caught a missing or longer timeout but passed a shorter one. It now measures from Publish entry and fails on a 2s or 10s bound. Spec-impact: none. Refs RIG-3107 Co-authored-by: Matt Wilkinson --- go/internal/comms/fabric_publish_pgtest_test.go | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/go/internal/comms/fabric_publish_pgtest_test.go b/go/internal/comms/fabric_publish_pgtest_test.go index 50df37703..b1162d6fc 100644 --- a/go/internal/comms/fabric_publish_pgtest_test.go +++ b/go/internal/comms/fabric_publish_pgtest_test.go @@ -29,6 +29,7 @@ type publishedRef struct { ctxErr error deadline time.Time hasDeadline bool + calledAt time.Time } // fakeEventFabric records publishes and, for each one, re-reads the row so a @@ -48,12 +49,13 @@ func (f *fakeEventFabric) Publish(ctx context.Context, subject string, ref fabri if f.onPublish != nil { f.onPublish() } + calledAt := time.Now() _, readErr := f.st.MessageByID(store.WithSystemRole(ctx), ref.RowID) deadline, hasDeadline := ctx.Deadline() f.mu.Lock() f.got = append(f.got, publishedRef{ subject: subject, ref: ref, readErr: readErr, ctxErr: ctx.Err(), - deadline: deadline, hasDeadline: hasDeadline, + deadline: deadline, hasDeadline: hasDeadline, calledAt: calledAt, }) f.mu.Unlock() return f.err @@ -198,15 +200,16 @@ func TestPostMessageFabricPublishOutlivesRequestCancellation(t *testing.T) { t.Fatalf("CreateChannel: %v", err) } - before := time.Now() id := postText(ctx, t, svc, poster.ID, ch.ID, "hello", "") assertOnePostedRef(t, fab, bootstrapTenant(st), id) // Detached from cancellation, the publish still needs its own bound, or a - // stalled JetStream ack would pin the RPC. + // stalled JetStream ack would pin the RPC. The deadline is set just before + // Publish is entered, so it lands within a second under calledAt+timeout. got := fab.published()[0] - if !got.hasDeadline || got.deadline.After(time.Now().Add(fabricPublishTimeout)) || got.deadline.Before(before) { - t.Fatalf("publish deadline = %v (set %v), want within %v of the call", got.deadline, got.hasDeadline, fabricPublishTimeout) + limit := got.calledAt.Add(fabricPublishTimeout) + if !got.hasDeadline || got.deadline.After(limit) || got.deadline.Before(limit.Add(-time.Second)) { + t.Fatalf("publish deadline = %v (set %v), want within 1s under %v", got.deadline, got.hasDeadline, limit) } }