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/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..b1162d6fc --- /dev/null +++ b/go/internal/comms/fabric_publish_pgtest_test.go @@ -0,0 +1,317 @@ +//go:build pgtest + +package comms + +import ( + "context" + "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" + + 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 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 + deadline time.Time + hasDeadline bool + calledAt time.Time +} + +// 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 + // 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() + } + 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, calledAt: calledAt, + }) + 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() + 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, 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 { + 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, 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) + } + wantSubject, err := fabric.CommsSubject(string(tenant), fabric.KindMessagePosted) + if err != nil { + t.Fatalf("CommsSubject: %v", err) + } + 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].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() + 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 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) + + // Detached from cancellation, the publish still needs its own bound, or a + // 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] + 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) + } +} + +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, fab, bootstrapTenant(st), 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, fab, bootstrapTenant(st), 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, fab, bootstrapTenant(st), string(answer.ID)) +} diff --git a/go/internal/comms/mapping.go b/go/internal/comms/mapping.go index 3e4e823c6..e8669e076 100644 --- a/go/internal/comms/mapping.go +++ b/go/internal/comms/mapping.go @@ -2,10 +2,13 @@ package comms import ( "context" + "log/slog" + "time" "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" ) @@ -495,12 +498,37 @@ 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{ 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 { + // 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. + 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/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) } 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)