From 6c3cef29e291591f627a885e50813514c8d897f0 Mon Sep 17 00:00:00 2001 From: Paul Chobert Date: Thu, 24 Sep 2026 18:23:13 +0200 Subject: [PATCH 1/2] MSC4140: schedule delayed events through the dedicated endpoint MSC4140 schedules delayed events with `PUT /rooms/{roomId}/delayed_event/{eventType}/{txnId}`, taking `delay_ms`, `content` and, for state events, `state_key` in the body. Scheduling through the `org.matrix.msc4140.delay` query parameter on `/send` and `/state` is now one of the MSC's rejected alternatives. Schedule every delayed event through the unstable form of the endpoint, giving state events a transaction ID of their own. The existing "same txnID" subtest now checks that the endpoint is transactional, and the state event tests check that a state event scheduled through it lands as room state. A new test checks that a missing, zero or negative `delay_ms` is rejected with a 400. One test still schedules a message event and a state event through the query parameter, so that form stays covered until homeservers drop it. Signed-off-by: Paul Chobert --- tests/msc4140/delayed_event_test.go | 152 ++++++++++++++++++++++------ 1 file changed, 122 insertions(+), 30 deletions(-) diff --git a/tests/msc4140/delayed_event_test.go b/tests/msc4140/delayed_event_test.go index 6c93ffb53..47c1b9e38 100644 --- a/tests/msc4140/delayed_event_test.go +++ b/tests/msc4140/delayed_event_test.go @@ -62,6 +62,31 @@ func TestDelayedEvents(t *testing.T) { }) }) + t.Run("cannot schedule a delayed event without a positive delay", func(t *testing.T) { + for i, tc := range []struct { + name string + body map[string]interface{} + }{ + {"missing delay_ms", map[string]interface{}{"content": map[string]interface{}{}}}, + {"zero delay_ms", getDelayedEventBody(0, map[string]interface{}{})}, + {"negative delay_ms", getDelayedEventBody(-1, map[string]interface{}{})}, + } { + t.Run(tc.name, func(t *testing.T) { + res := user.Do( + t, + "PUT", + getPathForDelayedEvent(roomID, eventType, fmt.Sprintf("txn-delayed-invalid-delay-%d", i)), + client.WithJSONBody(t, tc.body), + ) + must.MatchResponse(t, res, match.HTTPResponse{ + StatusCode: 400, + }) + }) + } + + matchDelayedEvents(t, user, delayedEventsNumberEqual(0)) + }) + // FIXME: Too much mixing of tests that should be more independent t.Run("delayed message events are sent on timeout", func(t *testing.T) { var res *http.Response @@ -75,15 +100,15 @@ func TestDelayedEvents(t *testing.T) { countKey := "count" numEvents := 3 - for i, delayStr := range []string{"700", "800", "900"} { + for i, delayMs := range []int64{700, 800, 900} { + body := getDelayedEventBody(delayMs, map[string]interface{}{ + countKey: i + 1, + }) res = user.MustDo( t, "PUT", - getPathForSend(roomID, eventType, fmt.Sprintf(txnIdBase, i)), - client.WithJSONBody(t, map[string]interface{}{ - countKey: i + 1, - }), - getDelayQueryParam(delayStr), + getPathForDelayedEvent(roomID, eventType, fmt.Sprintf(txnIdBase, i)), + client.WithJSONBody(t, body), ) delayID := client.GetJSONFieldStr(t, client.ParseJSON(t, res), "delay_id") @@ -91,8 +116,8 @@ func TestDelayedEvents(t *testing.T) { res := user.MustDo( t, "PUT", - getPathForSend(roomID, eventType, fmt.Sprintf(txnIdBase, i)), - getDelayQueryParam(delayStr), + getPathForDelayedEvent(roomID, eventType, fmt.Sprintf(txnIdBase, i)), + client.WithJSONBody(t, body), ) must.MatchResponse(t, res, match.HTTPResponse{ JSON: []match.JSON{ @@ -146,11 +171,10 @@ func TestDelayedEvents(t *testing.T) { user.MustDo( t, "PUT", - getPathForState(roomID, eventType, stateKey), - client.WithJSONBody(t, map[string]interface{}{ + getPathForDelayedEvent(roomID, eventType, "txn-delayed-state-timeout"), + client.WithJSONBody(t, getDelayedStateEventBody(900, stateKey, map[string]interface{}{ setterKey: setterExpected, - }), - getDelayQueryParam("900"), + })), ) // Ensure that a delayed event is now scheduled @@ -199,6 +223,58 @@ func TestDelayedEvents(t *testing.T) { matchDelayedEvents(t, user, delayedEventsNumberEqual(0)) }) + // TODO: Remove once homeservers stop accepting the `org.matrix.msc4140.delay` query + // parameter on `/send` and `/state`, which MSC4140 lists as a rejected alternative to + // the dedicated endpoint. + t.Run("delayed events scheduled with the delay query parameter are sent on timeout", func(t *testing.T) { + var res *http.Response + + defer cleanupDelayedEvents(t, user) + + stateKey := "to_send_on_timeout_with_query_param" + + // Schedule a delayed message event and a delayed state event + setterKey := "setter" + setterExpected := "on_timeout_with_query_param" + user.MustDo( + t, + "PUT", + getPathForSend(roomID, eventType, "txn-delayed-msg-query-param"), + client.WithJSONBody(t, map[string]interface{}{ + setterKey: setterExpected, + }), + getDelayQueryParam("900"), + ) + user.MustDo( + t, + "PUT", + getPathForState(roomID, eventType, stateKey), + client.WithJSONBody(t, map[string]interface{}{ + setterKey: setterExpected, + }), + getDelayQueryParam("900"), + ) + matchDelayedEvents(t, user, delayedEventsNumberEqual(2)) + + // Check for both delayed events being sent (using `MustSyncUntil` to account for + // any processing or worker replication delays) + user.MustSyncUntil(t, client.SyncReq{}, client.SyncTimelineHas(roomID, func(ev gjson.Result) bool { + return ev.Get("type").Str == eventType && !ev.Get("state_key").Exists() && ev.Get("content."+setterKey).Str == setterExpected + })) + user.MustSyncUntil(t, client.SyncReq{UseStateAfter: true}, client.SyncStateAfterHas(roomID, func(ev gjson.Result) bool { + return ev.Get("type").Str == eventType && ev.Get("state_key").Str == stateKey + })) + // Make sure the state looks as expected after + res = user.MustDo(t, "GET", getPathForState(roomID, eventType, stateKey)) + must.MatchResponse(t, res, match.HTTPResponse{ + JSON: []match.JSON{ + match.JSONKeyEqual(setterKey, setterExpected), + }, + }) + // No more delayed events + matchDelayedEvents(t, user, delayedEventsNumberEqual(0)) + }) + t.Run("cannot update a delayed event without an action", func(t *testing.T) { res := unauthedClient.Do( t, @@ -254,11 +330,10 @@ func TestDelayedEvents(t *testing.T) { res = user.MustDo( t, "PUT", - getPathForState(roomID, eventType, stateKey), - client.WithJSONBody(t, map[string]interface{}{ + getPathForDelayedEvent(roomID, eventType, "txn-delayed-state-cancel"), + client.WithJSONBody(t, getDelayedStateEventBody(1500, stateKey, map[string]interface{}{ setterKey: setterExpected, - }), - getDelayQueryParam("1500"), + })), ) delayID := client.GetJSONFieldStr(t, client.ParseJSON(t, res), "delay_id") @@ -308,11 +383,10 @@ func TestDelayedEvents(t *testing.T) { res = user.MustDo( t, "PUT", - getPathForState(roomID, eventType, stateKey), - client.WithJSONBody(t, map[string]interface{}{ + getPathForDelayedEvent(roomID, eventType, "txn-delayed-state-send"), + client.WithJSONBody(t, getDelayedStateEventBody(100000, stateKey, map[string]interface{}{ setterKey: setterExpected, - }), - getDelayQueryParam("100000"), + })), ) delayID := client.GetJSONFieldStr(t, client.ParseJSON(t, res), "delay_id") @@ -364,11 +438,10 @@ func TestDelayedEvents(t *testing.T) { res = user.MustDo( t, "PUT", - getPathForState(roomID, eventType, stateKey), - client.WithJSONBody(t, map[string]interface{}{ + getPathForDelayedEvent(roomID, eventType, "txn-delayed-state-restart"), + client.WithJSONBody(t, getDelayedStateEventBody(1500, stateKey, map[string]interface{}{ setterKey: setterExpected, - }), - getDelayQueryParam("1500"), + })), ) delayID := client.GetJSONFieldStr(t, client.ParseJSON(t, res), "delay_id") @@ -437,9 +510,8 @@ func TestDelayedEvents(t *testing.T) { user.MustDo( t, "PUT", - getPathForState(roomID, eventType, stateKey1), - client.WithJSONBody(t, map[string]interface{}{}), - getDelayQueryParam("900"), + getPathForDelayedEvent(roomID, eventType, "txn-delayed-state-server-restart-1"), + client.WithJSONBody(t, getDelayedStateEventBody(900, stateKey1, map[string]interface{}{})), ) numberOfDelayedEvents++ @@ -471,11 +543,10 @@ func TestDelayedEvents(t *testing.T) { user.MustDo( t, "PUT", + getPathForDelayedEvent(roomID, eventType, fmt.Sprintf("txn-delayed-state-server-restart-%d", i+2)), // Avoid clashing state keys as that would cancel previous delayed events on the // same key (start at 2). - getPathForState(roomID, eventType, fmt.Sprintf("%d", i+2)), - client.WithJSONBody(t, map[string]interface{}{}), - getDelayQueryParam(fmt.Sprintf("%d", delay.Milliseconds())), + client.WithJSONBody(t, getDelayedStateEventBody(delay.Milliseconds(), fmt.Sprintf("%d", i+2), map[string]interface{}{})), ) numberOfDelayedEvents++ } @@ -532,6 +603,27 @@ func getPathForUpdateDelayedEvent(delayId string, action DelayedEventAction) []s return append(getPathForDelayedEvents(), delayId, string(action)) } +func getPathForDelayedEvent(roomID string, eventType string, txnID string) []string { + return []string{"_matrix", "client", "unstable", "org.matrix.msc4140", "rooms", roomID, "delayed_event", eventType, txnID} +} + +// getDelayedEventBody returns the body of a request to schedule a delayed message event +// through `getPathForDelayedEvent`. +func getDelayedEventBody(delayMs int64, content map[string]interface{}) map[string]interface{} { + return map[string]interface{}{ + "delay_ms": delayMs, + "content": content, + } +} + +// getDelayedStateEventBody returns the body of a request to schedule a delayed state event +// through `getPathForDelayedEvent`. +func getDelayedStateEventBody(delayMs int64, stateKey string, content map[string]interface{}) map[string]interface{} { + body := getDelayedEventBody(delayMs, content) + body["state_key"] = stateKey + return body +} + func getPathForSend(roomID string, eventType string, txnId string) []string { return []string{"_matrix", "client", "v3", "rooms", roomID, "send", eventType, txnId} } From 6125401c972e36b0acffea2a2aeca9e6ae5ec4d4 Mon Sep 17 00:00:00 2001 From: Paul Chobert Date: Thu, 24 Sep 2026 18:42:31 +0200 Subject: [PATCH 2/2] MSC4140: test finalised delayed events MSC4140 keeps delayed events once they are finalised: sent, cancelled by the user, or cancelled due to an error. `GET /delayed_events/{delay_id}` describes the outcome in a `finalised` object, with `finalised_ts`, plus `event_id` if the event was sent or `error` if sending it failed. A management action on a finalised delayed event succeeds if it matches the outcome, and answers 409 if it conflicts with it. Add tests for: - the lookup of a delayed event sent on timeout, cancelled, or that failed to be sent because its sender left the room - a 200 for a repeated `send` or `cancel`, including `cancel` on a delayed event that failed to be sent - a 409 for `cancel` on a sent delayed event, and for `send` or `restart` on a cancelled one - a 409 for `restart` on a sent delayed event, which MSC4140 leaves open and MSC4542 proposes - a 404 when another user looks up a finalised delayed event - the bulk `GET /delayed_events` still listing scheduled delayed events only Signed-off-by: Paul Chobert --- tests/msc4140/delayed_event_test.go | 275 ++++++++++++++++++++++++++++ 1 file changed, 275 insertions(+) diff --git a/tests/msc4140/delayed_event_test.go b/tests/msc4140/delayed_event_test.go index 47c1b9e38..1fb3aa4d3 100644 --- a/tests/msc4140/delayed_event_test.go +++ b/tests/msc4140/delayed_event_test.go @@ -494,6 +494,250 @@ func TestDelayedEvents(t *testing.T) { matchDelayedEvents(t, user, delayedEventsNumberEqual(0)) }) + t.Run("delayed events sent on timeout are finalised with their event ID", func(t *testing.T) { + defer cleanupDelayedEvents(t, user) + + markerKey := "marker" + markerExpected := "finalised_on_timeout" + content := map[string]interface{}{ + markerKey: markerExpected, + } + res := user.MustDo( + t, + "PUT", + getPathForDelayedEvent(roomID, eventType, "txn-delayed-msg-finalised-timeout"), + client.WithJSONBody(t, getDelayedEventBody(900, content)), + ) + delayID := client.GetJSONFieldStr(t, client.ParseJSON(t, res), "delay_id") + + // Wait for the delayed event to be sent, and capture the ID it got in the room + var eventID string + user.MustSyncUntil(t, client.SyncReq{}, client.SyncTimelineHas(roomID, func(ev gjson.Result) bool { + if ev.Get("type").Str != eventType || ev.Get("content."+markerKey).Str != markerExpected { + return false + } + eventID = ev.Get("event_id").Str + return true + })) + + matchFinalisedDelayedEvent(t, user, delayID, + match.JSONKeyEqual("room_id", roomID), + match.JSONKeyEqual("type", eventType), + match.JSONKeyMissing("state_key"), + match.JSONKeyEqual("delay_ms", 900), + match.JSONKeyTypeEqual("delayed_since_ts", gjson.Number), + match.JSONKeyEqual("content", content), + match.JSONKeyEqual("finalised.event_id", eventID), + match.JSONKeyMissing("finalised.error"), + ) + }) + + t.Run("cancelled delayed events are finalised without an event ID", func(t *testing.T) { + var res *http.Response + + defer cleanupDelayedEvents(t, user) + + res = user.MustDo( + t, + "PUT", + getPathForDelayedEvent(roomID, eventType, "txn-delayed-msg-finalised-cancel"), + client.WithJSONBody(t, getDelayedEventBody(100000, map[string]interface{}{})), + ) + delayID := client.GetJSONFieldStr(t, client.ParseJSON(t, res), "delay_id") + + // A delayed event that is still scheduled is not finalised + res = user.MustDo(t, "GET", getPathForLookupDelayedEvent(delayID)) + must.MatchResponse(t, res, match.HTTPResponse{ + JSON: []match.JSON{ + match.JSONKeyEqual("delay_id", delayID), + match.JSONKeyMissing("finalised"), + }, + }) + + user.MustDo( + t, + "POST", + getPathForUpdateDelayedEvent(delayID, DelayedEventActionCancel), + client.WithJSONBody(t, map[string]interface{}{}), + ) + + matchFinalisedDelayedEvent(t, user, delayID, + match.JSONKeyMissing("finalised.event_id"), + match.JSONKeyMissing("finalised.error"), + ) + + t.Run("cannot look up a finalised delayed event of another user", func(t *testing.T) { + res := user2.Do(t, "GET", getPathForLookupDelayedEvent(delayID)) + must.MatchResponse(t, res, match.HTTPResponse{ + StatusCode: 404, + JSON: []match.JSON{ + match.JSONKeyEqual("errcode", "M_NOT_FOUND"), + }, + }) + }) + }) + + t.Run("delayed events that fail to be sent are finalised with an error", func(t *testing.T) { + // Use a room of its own, as the sender leaves it before the delayed event is sent + failRoomID := user.MustCreateRoom(t, map[string]interface{}{ + "preset": "public_chat", + "power_level_content_override": map[string]interface{}{ + "events": map[string]int{ + eventType: 0, + }, + }, + }) + user2.MustJoinRoom(t, failRoomID, nil) + + stateKey := "to_fail_on_timeout" + res := user2.MustDo( + t, + "PUT", + getPathForDelayedEvent(failRoomID, eventType, "txn-delayed-state-finalised-failure"), + client.WithJSONBody(t, getDelayedStateEventBody(1500, stateKey, map[string]interface{}{})), + ) + delayID := client.GetJSONFieldStr(t, client.ParseJSON(t, res), "delay_id") + + // Leave the room so that the delayed event cannot be sent at its scheduled time + user2.MustLeaveRoom(t, failRoomID) + + matchFinalisedDelayedEvent(t, user2, delayID, + match.JSONKeyEqual("state_key", stateKey), + match.JSONKeyTypeEqual("finalised.error.errcode", gjson.String), + match.JSONKeyMissing("finalised.event_id"), + ) + + // Sanity check that the room state hasn't changed + res = user.Do(t, "GET", getPathForState(failRoomID, eventType, stateKey)) + must.MatchResponse(t, res, match.HTTPResponse{ + StatusCode: 404, + }) + + // A delayed event cancelled due to an error counts as cancelled + user2.MustDo( + t, + "POST", + getPathForUpdateDelayedEvent(delayID, DelayedEventActionCancel), + client.WithJSONBody(t, map[string]interface{}{}), + ) + }) + + for _, tc := range []struct { + name string + finalisedBy DelayedEventAction + statusByAction map[DelayedEventAction]int + finalisedChecks []match.JSON + }{ + { + name: "sent", + finalisedBy: DelayedEventActionSend, + statusByAction: map[DelayedEventAction]int{ + DelayedEventActionSend: 200, + DelayedEventActionCancel: 409, + // MSC4140 leaves this case open, MSC4542 proposes a 409 + DelayedEventActionRestart: 409, + }, + finalisedChecks: []match.JSON{ + match.JSONKeyTypeEqual("finalised.event_id", gjson.String), + }, + }, + { + name: "cancelled", + finalisedBy: DelayedEventActionCancel, + statusByAction: map[DelayedEventAction]int{ + DelayedEventActionCancel: 200, + DelayedEventActionSend: 409, + DelayedEventActionRestart: 409, + }, + finalisedChecks: []match.JSON{ + match.JSONKeyMissing("finalised.event_id"), + }, + }, + } { + t.Run(fmt.Sprintf("actions on %s delayed events succeed if repeated and conflict otherwise", tc.name), func(t *testing.T) { + defer cleanupDelayedEvents(t, user) + + res := user.MustDo( + t, + "PUT", + getPathForDelayedEvent(roomID, eventType, fmt.Sprintf("txn-delayed-msg-finalised-%s-actions", tc.name)), + client.WithJSONBody(t, getDelayedEventBody(100000, map[string]interface{}{})), + ) + delayID := client.GetJSONFieldStr(t, client.ParseJSON(t, res), "delay_id") + + user.MustDo( + t, + "POST", + getPathForUpdateDelayedEvent(delayID, tc.finalisedBy), + client.WithJSONBody(t, map[string]interface{}{}), + ) + matchFinalisedDelayedEvent(t, user, delayID, tc.finalisedChecks...) + + for _, action := range []DelayedEventAction{ + DelayedEventActionCancel, + DelayedEventActionRestart, + DelayedEventActionSend, + } { + statusCode := tc.statusByAction[action] + t.Run(fmt.Sprintf("%s returns %d", action, statusCode), func(t *testing.T) { + res := user.Do( + t, + "POST", + getPathForUpdateDelayedEvent(delayID, action), + client.WithJSONBody(t, map[string]interface{}{}), + ) + must.MatchResponse(t, res, match.HTTPResponse{ + StatusCode: statusCode, + }) + }) + } + + // The delayed event is still finalised the same way + matchFinalisedDelayedEvent(t, user, delayID, tc.finalisedChecks...) + }) + } + + t.Run("finalised delayed events are not listed", func(t *testing.T) { + var res *http.Response + + defer cleanupDelayedEvents(t, user) + + delayIDs := make([]string, 3) + for i := range delayIDs { + res = user.MustDo( + t, + "PUT", + getPathForDelayedEvent(roomID, eventType, fmt.Sprintf("txn-delayed-msg-finalised-list-%d", i)), + client.WithJSONBody(t, getDelayedEventBody(100000, map[string]interface{}{})), + ) + delayIDs[i] = client.GetJSONFieldStr(t, client.ParseJSON(t, res), "delay_id") + } + matchDelayedEvents(t, user, delayedEventsNumberEqual(len(delayIDs))) + + // Finalise all but the last delayed event + for i, action := range []DelayedEventAction{ + DelayedEventActionSend, + DelayedEventActionCancel, + } { + user.MustDo( + t, + "POST", + getPathForUpdateDelayedEvent(delayIDs[i], action), + client.WithJSONBody(t, map[string]interface{}{}), + ) + matchFinalisedDelayedEvent(t, user, delayIDs[i]) + } + + // Only the delayed event that is still scheduled is listed + matchDelayedEvents(t, user, delayedEventsNumberEqual(1)) + res = getDelayedEvents(t, user) + must.MatchResponse(t, res, match.HTTPResponse{ + JSON: []match.JSON{ + match.JSONKeyEqual("delayed_events.0.delay_id", delayIDs[2]), + }, + }) + }) + t.Run("delayed state events are kept on server restart", func(t *testing.T) { // Spec cannot enforce server restart behaviour runtime.SkipIf(t, runtime.Dendrite, runtime.Conduit, runtime.Conduwuit) @@ -603,6 +847,10 @@ func getPathForUpdateDelayedEvent(delayId string, action DelayedEventAction) []s return append(getPathForDelayedEvents(), delayId, string(action)) } +func getPathForLookupDelayedEvent(delayID string) []string { + return append(getPathForDelayedEvents(), delayID) +} + func getPathForDelayedEvent(roomID string, eventType string, txnID string) []string { return []string{"_matrix", "client", "unstable", "org.matrix.msc4140", "rooms", roomID, "delayed_event", eventType, txnID} } @@ -737,6 +985,33 @@ func matchDelayedEvents(t *testing.T, user *client.CSAPI, checks ...delayedEvent ) } +// matchFinalisedDelayedEvent looks up the given delayed event until it is finalised, then +// runs the given checks on it. This retries as the homeserver may still be sending the +// delayed event. +func matchFinalisedDelayedEvent(t *testing.T, user *client.CSAPI, delayID string, checks ...match.JSON) { + t.Helper() + + res := user.MustDo(t, "GET", getPathForLookupDelayedEvent(delayID), + client.WithRetryUntil( + 5*time.Second, + func(res *http.Response) bool { + body, err := io.ReadAll(res.Body) + if err != nil { + t.Log(err) + return false + } + return res.StatusCode == 200 && gjson.GetBytes(body, "finalised").Exists() + }, + ), + ) + must.MatchResponse(t, res, match.HTTPResponse{ + JSON: append([]match.JSON{ + match.JSONKeyEqual("delay_id", delayID), + match.JSONKeyTypeEqual("finalised.finalised_ts", gjson.Number), + }, checks...), + }) +} + // FIXME: Instead of using `cleanupDelayedEvents`, each test should just use their own // room func cleanupDelayedEvents(t *testing.T, user *client.CSAPI) {