Skip to content

Commit 2d71979

Browse files
fix: accept unknown originators in SubscribeTopics cursor (#1983)
Resolves #1982 ## Summary `SubscribeTopics` rejected requests with `CodeInvalidArgument` whenever a client's per-topic cursor referenced an originator the node had not yet indexed (i.e. absent from `gateway_envelopes_latest`). This breaks legitimate flows: - `libxmtp` sends cursors containing migrator originator IDs (10, 11) even in environments where no migrator has been configured. - A client may learn of a new originator before a given node has indexed any of its messages (new nodes, network partitions). - A client may still hold a cursor for an originator that was removed long ago and has no messages on any currently-live node. The originator set is node-local and lossy — treating a mismatch as a fatal client error is the wrong default. ## Change - Remove the unknown-originator rejection loop in `validateTopicFilters`. All other validation (filter count, topic length, vector clock length) stays in place. - `buildTopicCursors` already tolerated unknown originators: the client's entries pass through unchanged, and `FillMissingOriginators` still seeds known originators at sequence 0. The downstream `SelectGatewayEnvelopesByPerTopicCursors` is a `CROSS JOIN LATERAL` on the client's `(topic, node_id, seq_id)` tuples — entries that don't match any row contribute nothing. - `advanceTopicCursors` already handles previously-unknown originators starting to publish live, so that path needed no change. ## Test plan - [x] Replaced `TestSubscribeTopics_UnknownOriginatorInCursor` (previously asserted rejection) with `TestSubscribeTopics_AcceptsUnknownOriginatorInCursor`, which asserts the subscription opens and delivers envelopes from known originators. - [x] Added `TestSubscribeTopics_MixedKnownAndUnknownOriginators`: cursor mixes a partially caught-up known originator with an unknown one; catch-up delivers only the unseen messages for the known originator. - [x] Added `TestSubscribeTopics_AllUnknownOriginatorsInCursor`: all cursor entries are unknown; catch-up still delivers messages from known originators via `FillMissingOriginators`. - [x] `TestSubscribeTopics_Validation` unchanged — other validation errors still return `CodeInvalidArgument`. - [x] `go test ./pkg/api/message/... -count=1` - [x] `dev/lint-fix` <!-- Macroscope's pull request summary starts here --> <!-- Macroscope will only edit the content between these invisible markers, and the markers themselves will not be visible in the GitHub rendered markdown. --> <!-- If you delete either of the start / end markers from your PR's description, Macroscope will append its summary at the bottom of the description. --> > [!NOTE] > ### Accept unknown originator IDs in `SubscribeTopics` cursor without error > Previously, `validateTopicFilters` rejected cursors containing originator IDs not known to the node. This removes that check, since unknown originators simply match no rows downstream and should not block subscription processing. > > - Removes the `knownOriginators` parameter from [`validateTopicFilters`](https://github.com/xmtp/xmtpd/pull/1983/files#diff-ac39f1802be2ab7789d5ffc0d2344a0d6f82d8b6a85fde0dc434a0069097bcff), dropping the loop that rejected unknown originator IDs in `LastSeen` cursors. > - Moves `validateTopicFilters` to run before the known-originator lookup, simplifying the handler flow. > - Behavioral Change: cursors with unknown originator IDs are now accepted; catch-up queries for those originators return no results rather than an error. > > <!-- Macroscope's review summary starts here --> > > <sup><a href="https://app.macroscope.com">Macroscope</a> summarized 426f102.</sup> > <!-- Macroscope's review summary ends here --> > <!-- Macroscope's pull request summary ends here --> Co-authored-by: xmtp-coder-agent <>
1 parent 727266f commit 2d71979

2 files changed

Lines changed: 89 additions & 29 deletions

File tree

‎pkg/api/message/subscribe_topics.go‎

Lines changed: 9 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,10 @@ func (s *Service) SubscribeTopics(
6868
)
6969
}
7070

71+
if err := validateTopicFilters(filters); err != nil {
72+
return err
73+
}
74+
7175
knownOriginators, err := s.originatorList.GetOriginatorNodeIDs(ctx)
7276
if err != nil {
7377
return connect.NewError(
@@ -76,10 +80,6 @@ func (s *Service) SubscribeTopics(
7680
)
7781
}
7882

79-
if err := validateTopicFilters(filters, knownOriginators); err != nil {
80-
return err
81-
}
82-
8383
cursors, topics, topicKeys := buildTopicCursors(filters, knownOriginators)
8484

8585
envelopesCh := s.subscribeWorker.listen(ctx, &subscribeFilter{
@@ -146,9 +146,13 @@ func (s *Service) SubscribeTopics(
146146
}
147147

148148
// validateTopicFilters validates the topic filters in a SubscribeTopicsRequest.
149+
// Cursor entries for originators the node has not yet seen are allowed: a
150+
// client may learn of a new originator before this node has indexed any of
151+
// its messages, or may still hold cursors for originators removed long ago.
152+
// Unknown originators are harmless in the downstream LATERAL query — they
153+
// simply match no rows.
149154
func validateTopicFilters(
150155
filters []*message_api.SubscribeTopicsRequest_TopicFilter,
151-
knownOriginators []uint32,
152156
) error {
153157
if len(filters) == 0 {
154158
return connect.NewError(
@@ -164,24 +168,10 @@ func validateTopicFilters(
164168
)
165169
}
166170

167-
known := make(map[uint32]struct{}, len(knownOriginators))
168-
for _, id := range knownOriginators {
169-
known[id] = struct{}{}
170-
}
171-
172171
for _, f := range filters {
173172
if err := validateTopicFilter(f); err != nil {
174173
return connect.NewError(connect.CodeInvalidArgument, err)
175174
}
176-
177-
for origID := range f.GetLastSeen().GetNodeIdToSequenceId() {
178-
if _, ok := known[origID]; !ok {
179-
return connect.NewError(
180-
connect.CodeInvalidArgument,
181-
fmt.Errorf("unknown originator node ID in cursor: %d", origID),
182-
)
183-
}
184-
}
185175
}
186176

187177
return nil

‎pkg/api/message/subscribe_topics_test.go‎

Lines changed: 80 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -189,25 +189,95 @@ func TestSubscribeTopics_Validation(t *testing.T) {
189189
}
190190
}
191191

192-
func TestSubscribeTopics_UnknownOriginatorInCursor(t *testing.T) {
192+
func TestSubscribeTopics_AcceptsUnknownOriginatorInCursor(t *testing.T) {
193193
client, store, _ := setupTopicTest(t)
194194
payerID := db.NullInt32(testutils.CreatePayer(t, store))
195195

196196
insertAndWait(t, store, []queries.InsertGatewayEnvelopeV3Params{
197197
makeEnvRow(t, 100, 1, topicA, payerID),
198198
})
199199

200-
// Reference originator 999 which is not known.
201-
stream, err := client.SubscribeTopics(
200+
// Cursor references originator 999, which the node has never seen.
201+
// The subscription should open normally; FillMissingOriginators seeds
202+
// known originator 100 at sequence 0, so catch-up still delivers its
203+
// envelope. The unknown originator is a no-op in the LATERAL query.
204+
stream := subscribeTopics(
205+
t,
206+
client,
202207
t.Context(),
203-
connect.NewRequest(&message_api.SubscribeTopicsRequest{
204-
Filters: []*message_api.SubscribeTopicsRequest_TopicFilter{
205-
makeFilter(topicA, map[uint32]uint64{999: 0}),
206-
},
207-
}),
208+
[]*message_api.SubscribeTopicsRequest_TopicFilter{
209+
makeFilter(topicA, map[uint32]uint64{999: 0}),
210+
},
208211
)
209-
require.NoError(t, err)
210-
requireTopicStreamError(t, stream, connect.CodeInvalidArgument)
212+
213+
envs := collectTopicEnvelopes(t, stream, 1)
214+
require.Len(t, envs, 1)
215+
decoded := envelopeTestUtils.UnmarshalUnsignedOriginatorEnvelope(
216+
t, envs[0].GetUnsignedOriginatorEnvelope(),
217+
)
218+
require.EqualValues(t, 100, decoded.GetOriginatorNodeId())
219+
require.EqualValues(t, 1, decoded.GetOriginatorSequenceId())
220+
}
221+
222+
func TestSubscribeTopics_MixedKnownAndUnknownOriginators(t *testing.T) {
223+
client, store, _ := setupTopicTest(t)
224+
payerID := db.NullInt32(testutils.CreatePayer(t, store))
225+
226+
insertAndWait(t, store, []queries.InsertGatewayEnvelopeV3Params{
227+
makeEnvRow(t, 100, 1, topicA, payerID),
228+
makeEnvRow(t, 100, 2, topicA, payerID),
229+
makeEnvRow(t, 200, 1, topicA, payerID),
230+
})
231+
232+
// Cursor: known originator 100 caught up past seq 1, plus an unknown
233+
// originator 999. Expect (100, 2) and (200, 1); (100, 1) is skipped
234+
// and the 999 entry contributes nothing.
235+
stream := subscribeTopics(
236+
t,
237+
client,
238+
t.Context(),
239+
[]*message_api.SubscribeTopicsRequest_TopicFilter{
240+
makeFilter(topicA, map[uint32]uint64{100: 1, 999: 0}),
241+
},
242+
)
243+
244+
envs := collectTopicEnvelopes(t, stream, 2)
245+
require.Len(t, envs, 2)
246+
247+
seen := make(map[uint32]uint64)
248+
for _, env := range envs {
249+
decoded := envelopeTestUtils.UnmarshalUnsignedOriginatorEnvelope(
250+
t, env.GetUnsignedOriginatorEnvelope(),
251+
)
252+
seen[decoded.GetOriginatorNodeId()] = decoded.GetOriginatorSequenceId()
253+
}
254+
require.Equal(t, uint64(2), seen[100])
255+
require.Equal(t, uint64(1), seen[200])
256+
}
257+
258+
func TestSubscribeTopics_AllUnknownOriginatorsInCursor(t *testing.T) {
259+
client, store, _ := setupTopicTest(t)
260+
payerID := db.NullInt32(testutils.CreatePayer(t, store))
261+
262+
insertAndWait(t, store, []queries.InsertGatewayEnvelopeV3Params{
263+
makeEnvRow(t, 100, 1, topicA, payerID),
264+
makeEnvRow(t, 200, 1, topicA, payerID),
265+
})
266+
267+
// Cursor references only originators the node has never seen. Both
268+
// known originators are absent from the cursor, so FillMissingOriginators
269+
// seeds them at 0 and catch-up delivers everything.
270+
stream := subscribeTopics(
271+
t,
272+
client,
273+
t.Context(),
274+
[]*message_api.SubscribeTopicsRequest_TopicFilter{
275+
makeFilter(topicA, map[uint32]uint64{888: 5, 999: 10}),
276+
},
277+
)
278+
279+
envs := collectTopicEnvelopes(t, stream, 2)
280+
require.Len(t, envs, 2)
211281
}
212282

213283
// ---- Live-Only Tests (nil LastSeen) ----

0 commit comments

Comments
 (0)