fix: reduce deadlocks merging PRs w/ async label stat recalcs (#9868)
The intent of this change is to reduce the scope of deadlock issues identified in #9785. I've identified other deadlock issues from synthetic testing, so this is not a complete fix, but it's a partial fix. This design was discussed in #9785 and this is the most basic implementation, with a very small scope of work converted to use it. Introduces a new `forgejo.org/services/stats` module which allows for the queuing and routing of recalc requests for object stats; in this case, the "number of issues" that are assigned to a label, and the number of closed issues that are assigned to a label. The reasons that these calculations are performed asynchronously through a queue are: - User operations that are common and performance-sensitive don't have to wait for recalculations that don't need to be exactly up-to-date at all times. For example, merging a pull request will be a faster operation; as it closes an issue, it needs to recalculate `label.num_closed_issues` for every label attached to the PR. - Database deadlocks that can occur between concurrent operations -- for example, if you were holding a lock on an issue while recalculating a label's count of open issues -- can be broken by making the recalculation occur outside of the transaction. ## Checklist The [contributor guide](https://forgejo.org/docs/next/contributor/) contains information that will be helpful to first time contributors. There also are a few [conditions for merging Pull Requests in Forgejo repositories](https://codeberg.org/forgejo/governance/src/branch/main/PullRequestsAgreement.md). You are also welcome to join the [Forgejo development chatroom](https://matrix.to/#/#forgejo-development:matrix.org). ### Tests - I added test coverage for Go changes... - [x] in their respective `*_test.go` for unit tests. - [ ] in the `tests/integration` directory if it involves interactions with a live Forgejo server. - I added test coverage for JavaScript changes... - [ ] in `web_src/js/*.test.js` if it can be unit tested. - [ ] in `tests/e2e/*.test.e2e.js` if it requires interactions with a live Forgejo server (see also the [developer guide for JavaScript testing](https://codeberg.org/forgejo/forgejo/src/branch/forgejo/tests/e2e/README.md#end-to-end-tests)). ### Documentation - [ ] I created a pull request [to the documentation](https://codeberg.org/forgejo/docs) to explain to Forgejo users how to use this change. - [x] I did not document these changes and I do not expect someone else to do it. - Internal developer documentation is present. ### Release notes - [ ] I do not want this change to show in the release notes. - [ ] I want the title to show in the release notes with a link to this pull request. - [x] I want the content of the `release-notes/<pull request number>.md` to be be used for the release notes instead of the title. Reviewed-on: https://codeberg.org/forgejo/forgejo/pulls/9868 Reviewed-by: Gusted <gusted@noreply.codeberg.org> Co-authored-by: Mathieu Fenniak <mathieu@fenniak.net> Co-committed-by: Mathieu Fenniak <mathieu@fenniak.net>
This commit is contained in:
committed by
Mathieu Fenniak
parent
2a3d852e46
commit
9e07bb07be
@@ -11,6 +11,7 @@ import (
|
||||
"forgejo.org/models/db"
|
||||
access_model "forgejo.org/models/perm/access"
|
||||
user_model "forgejo.org/models/user"
|
||||
"forgejo.org/services/stats"
|
||||
|
||||
"xorm.io/builder"
|
||||
)
|
||||
@@ -56,7 +57,7 @@ func newIssueLabel(ctx context.Context, issue *Issue, label *Label, doer *user_m
|
||||
|
||||
issue.Labels = append(issue.Labels, label)
|
||||
|
||||
return updateLabelCols(ctx, label, "num_issues", "num_closed_issue")
|
||||
return stats.QueueRecalcLabelByID(label.ID)
|
||||
}
|
||||
|
||||
// Remove all issue labels in the given exclusive scope
|
||||
@@ -191,7 +192,7 @@ func deleteIssueLabel(ctx context.Context, issue *Issue, label *Label, doer *use
|
||||
return err
|
||||
}
|
||||
|
||||
return updateLabelCols(ctx, label, "num_issues", "num_closed_issue")
|
||||
return stats.QueueRecalcLabelByID(label.ID)
|
||||
}
|
||||
|
||||
// DeleteIssueLabel deletes issue-label relation.
|
||||
|
||||
@@ -24,6 +24,7 @@ import (
|
||||
api "forgejo.org/modules/structs"
|
||||
"forgejo.org/modules/timeutil"
|
||||
"forgejo.org/modules/util"
|
||||
"forgejo.org/services/stats"
|
||||
|
||||
"xorm.io/builder"
|
||||
)
|
||||
@@ -101,8 +102,8 @@ func doChangeIssueStatus(ctx context.Context, issue *Issue, doer *user_model.Use
|
||||
if err := issue.LoadLabels(ctx); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for idx := range issue.Labels {
|
||||
if err := updateLabelCols(ctx, issue.Labels[idx], "num_issues", "num_closed_issue"); err != nil {
|
||||
for _, label := range issue.Labels {
|
||||
if err := stats.QueueRecalcLabelByID(label.ID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
+26
-6
@@ -16,6 +16,7 @@ import (
|
||||
"forgejo.org/modules/optional"
|
||||
"forgejo.org/modules/timeutil"
|
||||
"forgejo.org/modules/util"
|
||||
"forgejo.org/services/stats"
|
||||
|
||||
"xorm.io/builder"
|
||||
)
|
||||
@@ -240,7 +241,12 @@ func UpdateLabel(ctx context.Context, l *Label) error {
|
||||
}
|
||||
l.Color = color
|
||||
|
||||
return updateLabelCols(ctx, l, "name", "description", "color", "exclusive", "archived_unix")
|
||||
_, err = db.GetEngine(ctx).Cols("name", "description", "color", "exclusive", "archived_unix").ID(l.ID).Update(l)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return stats.QueueRecalcLabelByID(l.ID)
|
||||
}
|
||||
|
||||
// DeleteLabel delete a label
|
||||
@@ -510,20 +516,34 @@ func CountLabelsByOrgID(ctx context.Context, orgID int64) (int64, error) {
|
||||
return db.GetEngine(ctx).Where("org_id = ?", orgID).Count(&Label{})
|
||||
}
|
||||
|
||||
func updateLabelCols(ctx context.Context, l *Label, cols ...string) error {
|
||||
_, err := db.GetEngine(ctx).ID(l.ID).
|
||||
func init() {
|
||||
stats.RegisterRecalc(stats.LabelByLabelID, doRecalcLabelByID)
|
||||
stats.RegisterRecalc(stats.LabelByRepoID, doRecalcLabelByRepoID)
|
||||
}
|
||||
|
||||
func doRecalcLabelByID(ctx context.Context, labelID int64) error {
|
||||
return doRecalcLabel(ctx, builder.Eq{"id": labelID})
|
||||
}
|
||||
|
||||
func doRecalcLabelByRepoID(ctx context.Context, repoID int64) error {
|
||||
return doRecalcLabel(ctx, builder.Eq{"repo_id": repoID})
|
||||
}
|
||||
|
||||
func doRecalcLabel(ctx context.Context, cond builder.Cond) error {
|
||||
_, err := db.GetEngine(ctx).
|
||||
SetExpr("num_issues",
|
||||
builder.Select("count(*)").From("issue_label").
|
||||
Where(builder.Eq{"label_id": l.ID}),
|
||||
Where(builder.Eq{"label_id": builder.Expr("label.id")}),
|
||||
).
|
||||
SetExpr("num_closed_issues",
|
||||
builder.Select("count(*)").From("issue_label").
|
||||
InnerJoin("issue", "issue_label.issue_id = issue.id").
|
||||
Where(builder.Eq{
|
||||
"issue_label.label_id": l.ID,
|
||||
"issue_label.label_id": builder.Expr("label.id"),
|
||||
"issue.is_closed": true,
|
||||
}),
|
||||
).
|
||||
Cols(cols...).Update(l)
|
||||
Where(cond).
|
||||
Update(&Label{})
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -0,0 +1,71 @@
|
||||
// Copyright 2025 The Forgejo Authors. All rights reserved.
|
||||
// SPDX-License-Identifier: GPL-3.0-or-later
|
||||
|
||||
package issues
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"forgejo.org/models/db"
|
||||
"forgejo.org/models/unittest"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestRecalcLabelByLabelID(t *testing.T) {
|
||||
require.NoError(t, unittest.PrepareTestDatabase())
|
||||
|
||||
// Verify no error on recalc of a deleted/non-existent object; important because async recalcs can be queued and
|
||||
// then occur later after more state changes have happened.
|
||||
err := doRecalcLabelByID(t.Context(), -1000)
|
||||
require.NoError(t, err)
|
||||
|
||||
// Intentionally corrupt counts from fixture, then recalc them
|
||||
label := unittest.AssertExistsAndLoadBean(t, &Label{ID: 1})
|
||||
updated, err := db.GetEngine(t.Context()).
|
||||
Table(&Label{}).
|
||||
Where("id = ?", label.ID).
|
||||
Update(map[string]any{"num_issues": 1000, "num_closed_issues": 1001})
|
||||
require.NoError(t, err)
|
||||
require.EqualValues(t, 1, updated)
|
||||
err = doRecalcLabelByID(t.Context(), label.ID)
|
||||
require.NoError(t, err)
|
||||
label = unittest.AssertExistsAndLoadBean(t, &Label{ID: 1})
|
||||
assert.Equal(t, 2, label.NumIssues)
|
||||
assert.Equal(t, 0, label.NumClosedIssues)
|
||||
}
|
||||
|
||||
func TestRecalcLabelByRepoID(t *testing.T) {
|
||||
require.NoError(t, unittest.PrepareTestDatabase())
|
||||
|
||||
// Verify no error on recalc of a deleted/non-existent object; important because async recalcs can be queued and
|
||||
// then occur later after more state changes have happened.
|
||||
err := doRecalcLabelByRepoID(t.Context(), -1000)
|
||||
require.NoError(t, err)
|
||||
|
||||
// Intentionally corrupt counts from fixture, then recalc them
|
||||
label1 := unittest.AssertExistsAndLoadBean(t, &Label{ID: 1})
|
||||
updated, err := db.GetEngine(t.Context()).
|
||||
Table(&Label{}).
|
||||
Where("id = ?", label1.ID).
|
||||
Update(map[string]any{"num_issues": 1000, "num_closed_issues": 1001})
|
||||
require.NoError(t, err)
|
||||
require.EqualValues(t, 1, updated)
|
||||
label2 := unittest.AssertExistsAndLoadBean(t, &Label{ID: 2})
|
||||
require.Equal(t, label1.RepoID, label2.RepoID) // sanity check
|
||||
updated, err = db.GetEngine(t.Context()).
|
||||
Table(&Label{}).
|
||||
Where("id = ?", label2.ID).
|
||||
Update(map[string]any{"num_issues": 1000, "num_closed_issues": 1001})
|
||||
require.NoError(t, err)
|
||||
require.EqualValues(t, 1, updated)
|
||||
err = doRecalcLabelByRepoID(t.Context(), label1.RepoID)
|
||||
require.NoError(t, err)
|
||||
label1 = unittest.AssertExistsAndLoadBean(t, &Label{ID: 1})
|
||||
label2 = unittest.AssertExistsAndLoadBean(t, &Label{ID: 2})
|
||||
assert.Equal(t, 2, label1.NumIssues)
|
||||
assert.Equal(t, 0, label1.NumClosedIssues)
|
||||
assert.Equal(t, 1, label2.NumIssues)
|
||||
assert.Equal(t, 1, label2.NumClosedIssues)
|
||||
}
|
||||
@@ -314,6 +314,7 @@ func TestNewIssueLabel(t *testing.T) {
|
||||
LabelID: label.ID,
|
||||
Content: "1",
|
||||
})
|
||||
unittest.FlushAsyncCalcs(t)
|
||||
label = unittest.AssertExistsAndLoadBean(t, &issues_model.Label{ID: 2})
|
||||
assert.Equal(t, prevNumIssues+1, label.NumIssues)
|
||||
|
||||
@@ -366,6 +367,7 @@ func TestNewIssueLabels(t *testing.T) {
|
||||
LabelID: label1.ID,
|
||||
Content: "1",
|
||||
})
|
||||
unittest.FlushAsyncCalcs(t)
|
||||
unittest.AssertExistsAndLoadBean(t, &issues_model.IssueLabel{IssueID: issue.ID, LabelID: label1.ID})
|
||||
label1 = unittest.AssertExistsAndLoadBean(t, &issues_model.Label{ID: 1})
|
||||
assert.Equal(t, 3, label1.NumIssues)
|
||||
@@ -409,6 +411,7 @@ func TestDeleteIssueLabel(t *testing.T) {
|
||||
IssueID: issueID,
|
||||
LabelID: labelID,
|
||||
}, `content=""`)
|
||||
unittest.FlushAsyncCalcs(t)
|
||||
label = unittest.AssertExistsAndLoadBean(t, &issues_model.Label{ID: labelID})
|
||||
assert.Equal(t, expectedNumIssues, label.NumIssues)
|
||||
assert.Equal(t, expectedNumClosedIssues, label.NumClosedIssues)
|
||||
|
||||
+3
-4
@@ -19,6 +19,7 @@ import (
|
||||
"forgejo.org/models/unit"
|
||||
user_model "forgejo.org/models/user"
|
||||
"forgejo.org/modules/log"
|
||||
"forgejo.org/services/stats"
|
||||
|
||||
"xorm.io/builder"
|
||||
)
|
||||
@@ -80,13 +81,11 @@ func labelStatsCorrectNumIssuesRepo(ctx context.Context, id int64) error {
|
||||
}
|
||||
|
||||
func labelStatsCorrectNumClosedIssues(ctx context.Context, id int64) error {
|
||||
_, err := db.GetEngine(ctx).Exec("UPDATE `label` SET num_closed_issues=(SELECT COUNT(*) FROM `issue_label`,`issue` WHERE `issue_label`.label_id=`label`.id AND `issue_label`.issue_id=`issue`.id AND `issue`.is_closed=?) WHERE `label`.id=?", true, id)
|
||||
return err
|
||||
return stats.QueueRecalcLabelByID(id)
|
||||
}
|
||||
|
||||
func labelStatsCorrectNumClosedIssuesRepo(ctx context.Context, id int64) error {
|
||||
_, err := db.GetEngine(ctx).Exec("UPDATE `label` SET num_closed_issues=(SELECT COUNT(*) FROM `issue_label`,`issue` WHERE `issue_label`.label_id=`label`.id AND `issue_label`.issue_id=`issue`.id AND `issue`.is_closed=?) WHERE `label`.repo_id=?", true, id)
|
||||
return err
|
||||
return stats.QueueRecalcLabelByRepoID(id)
|
||||
}
|
||||
|
||||
var milestoneStatsQueryNumIssues = "SELECT `milestone`.id FROM `milestone` WHERE `milestone`.num_closed_issues!=(SELECT COUNT(*) FROM `issue` WHERE `issue`.milestone_id=`milestone`.id AND `issue`.is_closed=?) OR `milestone`.num_issues!=(SELECT COUNT(*) FROM `issue` WHERE `issue`.milestone_id=`milestone`.id)"
|
||||
|
||||
@@ -27,6 +27,7 @@ var consistencyCheckMap = make(map[string]func(t *testing.T, bean any))
|
||||
|
||||
// CheckConsistencyFor test that all matching database entries are consistent
|
||||
func CheckConsistencyFor(t *testing.T, beansToCheck ...any) {
|
||||
FlushAsyncCalcs(t)
|
||||
for _, bean := range beansToCheck {
|
||||
sliceType := reflect.SliceOf(reflect.TypeOf(bean))
|
||||
sliceValue := reflect.MakeSlice(sliceType, 0, 10)
|
||||
|
||||
@@ -20,7 +20,9 @@ import (
|
||||
"forgejo.org/modules/setting"
|
||||
"forgejo.org/modules/setting/config"
|
||||
"forgejo.org/modules/storage"
|
||||
"forgejo.org/modules/test"
|
||||
"forgejo.org/modules/util"
|
||||
"forgejo.org/services/stats"
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
"xorm.io/xorm"
|
||||
@@ -158,6 +160,7 @@ func MainTest(m *testing.M, testOpts ...*TestOptions) {
|
||||
if err = storage.Init(); err != nil {
|
||||
fatalTestError("storage.Init: %v\n", err)
|
||||
}
|
||||
initStats()
|
||||
if err = util.RemoveAll(repoRootPath); err != nil {
|
||||
fatalTestError("util.RemoveAll: %v\n", err)
|
||||
}
|
||||
@@ -211,6 +214,22 @@ func MainTest(m *testing.M, testOpts ...*TestOptions) {
|
||||
os.Exit(exitStatus)
|
||||
}
|
||||
|
||||
func initStats() {
|
||||
// Use an in-memory queue for the `stats` module during testing. This queue will collect requests for recalc during
|
||||
// tests which can be performed by invoking `unittest.FlushAsyncCalcs(t)`.
|
||||
cfg, err := setting.NewConfigProviderFromData(`
|
||||
[queue.stats_recalc]
|
||||
TYPE = channel
|
||||
`)
|
||||
if err != nil {
|
||||
fatalTestError("NewConfigProviderFromData: %v\n", err)
|
||||
}
|
||||
defer test.MockVariableValue(&setting.CfgProvider, cfg)()
|
||||
if err := stats.Init(); err != nil {
|
||||
fatalTestError("stats.Init: %v\n", err)
|
||||
}
|
||||
}
|
||||
|
||||
// FixturesOptions fixtures needs to be loaded options
|
||||
type FixturesOptions struct {
|
||||
Dir string
|
||||
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
"testing"
|
||||
|
||||
"forgejo.org/models/db"
|
||||
"forgejo.org/services/stats"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
@@ -162,3 +163,7 @@ func AssertCountByCond(t testing.TB, tableName string, cond builder.Cond, expect
|
||||
return assert.EqualValues(t, expected, GetCountByCond(t, tableName, cond),
|
||||
"Failed consistency test, the counted bean (of table %s) was %+v", tableName, cond)
|
||||
}
|
||||
|
||||
func FlushAsyncCalcs(t testing.TB) {
|
||||
require.NoError(t, stats.Flush(t.Context()))
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user