mirror of
https://github.com/mattermost/mattermost.git
synced 2026-09-24 16:05:00 +08:00
MM-65661 - channel admin abac override previous jobs (#33872)
* MM-65661 - channel admin abac override previous jobs * more UX adjustments; always show the self-exclusion warning modal * use SubjectID parameter for more performant user lookup instead of fetching all matching users * improve validation of result based on PR feedback * performance optimization and DoS protection for access control sync jobs * refactor: rename context parameter from 'c' to 'rctx' in job-related functions for consistency * prevent duplicate save button clicks with immediate response and remove unnecessary debouncing time * remove dedicated endpoint and unify logic * improve filtering performance by including statuses * use a flag to use master directly to prevent db replication lags * adjust unit tests based on pr feedback --------- Co-authored-by: Mattermost Build <build@mattermost.com>
This commit is contained in:
co-authored by
Mattermost Build
parent
f5693467db
commit
de686b80bf
@@ -8,6 +8,7 @@ import (
|
||||
"net/http"
|
||||
|
||||
"github.com/mattermost/mattermost/server/public/model"
|
||||
"github.com/mattermost/mattermost/server/public/shared/mlog"
|
||||
"github.com/mattermost/mattermost/server/public/shared/request"
|
||||
"github.com/mattermost/mattermost/server/v8/channels/store"
|
||||
)
|
||||
@@ -74,7 +75,48 @@ func (a *App) GetJobsByTypesAndStatuses(rctx request.CTX, jobTypes []string, sta
|
||||
}
|
||||
|
||||
func (a *App) CreateJob(rctx request.CTX, job *model.Job) (*model.Job, *model.AppError) {
|
||||
return a.Srv().Jobs.CreateJob(rctx, job.Type, job.Data)
|
||||
switch job.Type {
|
||||
case model.JobTypeAccessControlSync:
|
||||
// Route ABAC jobs to specialized deduplication handler
|
||||
return a.CreateAccessControlSyncJob(rctx, job.Data)
|
||||
default:
|
||||
return a.Srv().Jobs.CreateJob(rctx, job.Type, job.Data)
|
||||
}
|
||||
}
|
||||
|
||||
func (a *App) CreateAccessControlSyncJob(rctx request.CTX, jobData map[string]string) (*model.Job, *model.AppError) {
|
||||
// Get the policy_id (channel ID) from job data to scope the deduplication
|
||||
policyID, exists := jobData["policy_id"]
|
||||
|
||||
// If policy_id is provided, this is a channel-specific job that needs deduplication
|
||||
if exists && policyID != "" {
|
||||
// Find existing pending or in-progress jobs for this specific policy/channel
|
||||
existingJobs, err := a.Srv().Store().Job().GetByTypeAndData(rctx, model.JobTypeAccessControlSync, map[string]string{
|
||||
"policy_id": policyID,
|
||||
}, true, model.JobStatusPending, model.JobStatusInProgress)
|
||||
if err != nil {
|
||||
return nil, model.NewAppError("CreateAccessControlSyncJob", "app.job.get_existing_jobs.error", nil, "", http.StatusInternalServerError).Wrap(err)
|
||||
}
|
||||
|
||||
// Cancel any existing active jobs for this policy (all returned jobs are already active)
|
||||
for _, job := range existingJobs {
|
||||
rctx.Logger().Info("Canceling existing access control sync job before creating new one",
|
||||
mlog.String("job_id", job.Id),
|
||||
mlog.String("policy_id", policyID),
|
||||
mlog.String("status", job.Status))
|
||||
|
||||
// directly cancel jobs for deduplication
|
||||
if err := a.Srv().Jobs.SetJobCanceled(job); err != nil {
|
||||
rctx.Logger().Warn("Failed to cancel existing access control sync job",
|
||||
mlog.String("job_id", job.Id),
|
||||
mlog.String("policy_id", policyID),
|
||||
mlog.Err(err))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Create the new job
|
||||
return a.Srv().Jobs.CreateJob(rctx, model.JobTypeAccessControlSync, jobData)
|
||||
}
|
||||
|
||||
func (a *App) CancelJob(rctx request.CTX, jobId string) *model.AppError {
|
||||
|
||||
@@ -219,6 +219,179 @@ func TestSessionHasPermissionToCreateAccessControlSyncJob(t *testing.T) {
|
||||
})
|
||||
}
|
||||
|
||||
func TestCreateAccessControlSyncJob(t *testing.T) {
|
||||
mainHelper.Parallel(t)
|
||||
th := Setup(t).InitBasic()
|
||||
defer th.TearDown()
|
||||
|
||||
t.Run("cancels pending job and creates new one", func(t *testing.T) {
|
||||
// Create an existing pending job manually in the store
|
||||
existingJob := &model.Job{
|
||||
Id: model.NewId(),
|
||||
Type: model.JobTypeAccessControlSync,
|
||||
Status: model.JobStatusPending,
|
||||
Data: map[string]string{
|
||||
"policy_id": "channel456",
|
||||
},
|
||||
}
|
||||
_, err := th.App.Srv().Store().Job().Save(existingJob)
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() {
|
||||
_, stErr := th.App.Srv().Store().Job().Delete(existingJob.Id)
|
||||
require.NoError(t, stErr)
|
||||
})
|
||||
|
||||
// Test the cancellation logic by calling the method directly
|
||||
existingJobs, storeErr := th.App.Srv().Store().Job().GetByTypeAndData(th.Context, model.JobTypeAccessControlSync, map[string]string{
|
||||
"policy_id": "channel456",
|
||||
}, false, model.JobStatusPending, model.JobStatusInProgress)
|
||||
require.NoError(t, storeErr)
|
||||
require.Len(t, existingJobs, 1)
|
||||
|
||||
// Verify that the store method finds the job
|
||||
assert.Equal(t, existingJob.Id, existingJobs[0].Id)
|
||||
assert.Equal(t, model.JobStatusPending, existingJobs[0].Status)
|
||||
|
||||
// Test the cancellation logic directly
|
||||
for _, job := range existingJobs {
|
||||
if job.Status == model.JobStatusPending || job.Status == model.JobStatusInProgress {
|
||||
appErr := th.App.CancelJob(th.Context, job.Id)
|
||||
require.Nil(t, appErr)
|
||||
}
|
||||
}
|
||||
|
||||
// Verify that the job was cancelled
|
||||
updatedJob, getErr := th.App.Srv().Store().Job().Get(th.Context, existingJob.Id)
|
||||
require.NoError(t, getErr)
|
||||
// Job should be either cancel_requested or canceled (async process)
|
||||
assert.Contains(t, []string{model.JobStatusCancelRequested, model.JobStatusCanceled}, updatedJob.Status)
|
||||
})
|
||||
|
||||
t.Run("cancels in-progress job and creates new one", func(t *testing.T) {
|
||||
// Create an existing in-progress job
|
||||
existingJob := &model.Job{
|
||||
Id: model.NewId(),
|
||||
Type: model.JobTypeAccessControlSync,
|
||||
Status: model.JobStatusInProgress,
|
||||
Data: map[string]string{
|
||||
"policy_id": "channel789",
|
||||
},
|
||||
}
|
||||
_, err := th.App.Srv().Store().Job().Save(existingJob)
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() {
|
||||
_, stErr := th.App.Srv().Store().Job().Delete(existingJob.Id)
|
||||
require.NoError(t, stErr)
|
||||
})
|
||||
|
||||
// Test that GetByTypeAndData finds the in-progress job
|
||||
existingJobs, storeErr := th.App.Srv().Store().Job().GetByTypeAndData(th.Context, model.JobTypeAccessControlSync, map[string]string{
|
||||
"policy_id": "channel789",
|
||||
}, false, model.JobStatusPending, model.JobStatusInProgress)
|
||||
require.NoError(t, storeErr)
|
||||
require.Len(t, existingJobs, 1)
|
||||
assert.Equal(t, model.JobStatusInProgress, existingJobs[0].Status)
|
||||
|
||||
// Test cancellation of in-progress job
|
||||
appErr := th.App.CancelJob(th.Context, existingJob.Id)
|
||||
require.Nil(t, appErr)
|
||||
|
||||
// Verify cancellation was requested (job cancellation is asynchronous)
|
||||
updatedJob, getErr := th.App.Srv().Store().Job().Get(th.Context, existingJob.Id)
|
||||
require.NoError(t, getErr)
|
||||
// Job should be either cancel_requested or canceled (async process)
|
||||
assert.Contains(t, []string{model.JobStatusCancelRequested, model.JobStatusCanceled}, updatedJob.Status)
|
||||
})
|
||||
|
||||
t.Run("leaves completed jobs alone", func(t *testing.T) {
|
||||
// Create an existing completed job
|
||||
existingJob := &model.Job{
|
||||
Id: model.NewId(),
|
||||
Type: model.JobTypeAccessControlSync,
|
||||
Status: model.JobStatusSuccess,
|
||||
Data: map[string]string{
|
||||
"policy_id": "channel101",
|
||||
},
|
||||
}
|
||||
_, err := th.App.Srv().Store().Job().Save(existingJob)
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() {
|
||||
_, stErr := th.App.Srv().Store().Job().Delete(existingJob.Id)
|
||||
require.NoError(t, stErr)
|
||||
})
|
||||
|
||||
// Test that GetByTypeAndData finds the completed job
|
||||
existingJobs, storeErr := th.App.Srv().Store().Job().GetByTypeAndData(th.Context, model.JobTypeAccessControlSync, map[string]string{
|
||||
"policy_id": "channel101",
|
||||
}, false)
|
||||
require.NoError(t, storeErr)
|
||||
require.Len(t, existingJobs, 1)
|
||||
assert.Equal(t, model.JobStatusSuccess, existingJobs[0].Status)
|
||||
|
||||
// Test that we don't cancel completed jobs (logic test)
|
||||
shouldCancel := existingJob.Status == model.JobStatusPending || existingJob.Status == model.JobStatusInProgress
|
||||
assert.False(t, shouldCancel, "Should not cancel completed jobs")
|
||||
|
||||
// Verify the job status is unchanged
|
||||
updatedJob, getErr := th.App.Srv().Store().Job().Get(th.Context, existingJob.Id)
|
||||
require.NoError(t, getErr)
|
||||
assert.Equal(t, model.JobStatusSuccess, updatedJob.Status)
|
||||
})
|
||||
|
||||
// Test deduplication logic with status filtering to ensure database optimization works correctly
|
||||
|
||||
t.Run("deduplication respects status filtering", func(t *testing.T) {
|
||||
// Create jobs with different statuses
|
||||
pendingJob := &model.Job{
|
||||
Id: model.NewId(),
|
||||
Type: model.JobTypeAccessControlSync,
|
||||
Status: model.JobStatusPending,
|
||||
Data: map[string]string{"policy_id": "channel999"},
|
||||
}
|
||||
|
||||
completedJob := &model.Job{
|
||||
Id: model.NewId(),
|
||||
Type: model.JobTypeAccessControlSync,
|
||||
Status: model.JobStatusSuccess,
|
||||
Data: map[string]string{"policy_id": "channel999"},
|
||||
}
|
||||
|
||||
for _, job := range []*model.Job{pendingJob, completedJob} {
|
||||
_, err := th.App.Srv().Store().Job().Save(job)
|
||||
require.NoError(t, err)
|
||||
|
||||
// Capture job ID to avoid closure variable capture issue
|
||||
jobID := job.Id
|
||||
t.Cleanup(func() {
|
||||
_, stErr := th.App.Srv().Store().Job().Delete(jobID)
|
||||
require.NoError(t, stErr)
|
||||
})
|
||||
}
|
||||
|
||||
// Verify status filtering returns only active jobs
|
||||
activeJobs, err := th.App.Srv().Store().Job().GetByTypeAndData(
|
||||
th.Context,
|
||||
model.JobTypeAccessControlSync,
|
||||
map[string]string{"policy_id": "channel999"},
|
||||
false,
|
||||
model.JobStatusPending, model.JobStatusInProgress, // Only active statuses
|
||||
)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, activeJobs, 1, "Should only find active jobs (pending/in-progress)")
|
||||
assert.Equal(t, pendingJob.Id, activeJobs[0].Id, "Should find the pending job")
|
||||
|
||||
// Verify all jobs are returned when no status filter is provided
|
||||
allJobs, err := th.App.Srv().Store().Job().GetByTypeAndData(
|
||||
th.Context,
|
||||
model.JobTypeAccessControlSync,
|
||||
map[string]string{"policy_id": "channel999"},
|
||||
false, // No status filter
|
||||
)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, allJobs, 2, "Should find all jobs when no status filter")
|
||||
})
|
||||
}
|
||||
|
||||
func TestSessionHasPermissionToReadJob(t *testing.T) {
|
||||
mainHelper.Parallel(t)
|
||||
th := Setup(t)
|
||||
|
||||
@@ -387,6 +387,38 @@ func (jss SqlJobStore) GetCountByStatusAndType(status string, jobType string) (i
|
||||
return count, nil
|
||||
}
|
||||
|
||||
func (jss SqlJobStore) GetByTypeAndData(rctx request.CTX, jobType string, data map[string]string, useMaster bool, statuses ...string) ([]*model.Job, error) {
|
||||
query := jss.jobQuery.Where(sq.Eq{"Type": jobType})
|
||||
|
||||
// Add status filtering if provided - enables full usage of idx_jobs_status_type index
|
||||
if len(statuses) > 0 {
|
||||
query = query.Where(sq.Eq{"Status": statuses})
|
||||
}
|
||||
|
||||
// Add JSON data filtering for each key-value pair
|
||||
for key, value := range data {
|
||||
query = query.Where(sq.Expr("Data->? = ?", key, fmt.Sprintf(`"%s"`, value)))
|
||||
}
|
||||
|
||||
queryString, args, err := query.ToSql()
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "get_by_type_and_data_tosql")
|
||||
}
|
||||
|
||||
var jobs []*model.Job
|
||||
// For consistency-critical operations (like job deduplication), use master
|
||||
db := jss.GetReplica()
|
||||
if useMaster {
|
||||
db = jss.GetMaster()
|
||||
}
|
||||
|
||||
if err := db.Select(&jobs, queryString, args...); err != nil {
|
||||
return nil, errors.Wrap(err, "failed to get Jobs by type and data")
|
||||
}
|
||||
|
||||
return jobs, nil
|
||||
}
|
||||
|
||||
func (jss SqlJobStore) Delete(id string) (string, error) {
|
||||
query, args, err := jss.getQueryBuilder().
|
||||
Delete("Jobs").
|
||||
|
||||
@@ -799,6 +799,7 @@ type JobStore interface {
|
||||
GetNewestJobByStatusAndType(status string, jobType string) (*model.Job, error)
|
||||
GetNewestJobByStatusesAndType(statuses []string, jobType string) (*model.Job, error)
|
||||
GetCountByStatusAndType(status string, jobType string) (int64, error)
|
||||
GetByTypeAndData(rctx request.CTX, jobType string, data map[string]string, useMaster bool, statuses ...string) ([]*model.Job, error)
|
||||
Delete(id string) (string, error)
|
||||
Cleanup(expiryTime int64, batchSize int) error
|
||||
}
|
||||
|
||||
@@ -32,6 +32,7 @@ func TestJobStore(t *testing.T, rctx request.CTX, ss store.Store) {
|
||||
t.Run("GetCountByStatusAndType", func(t *testing.T) { testJobStoreGetCountByStatusAndType(t, rctx, ss) })
|
||||
t.Run("JobUpdateOptimistically", func(t *testing.T) { testJobUpdateOptimistically(t, rctx, ss) })
|
||||
t.Run("JobUpdateStatusUpdateStatusOptimistically", func(t *testing.T) { testJobUpdateStatusUpdateStatusOptimistically(t, rctx, ss) })
|
||||
t.Run("JobGetByTypeAndData", func(t *testing.T) { testJobGetByTypeAndData(t, rctx, ss) })
|
||||
t.Run("JobDelete", func(t *testing.T) { testJobDelete(t, rctx, ss) })
|
||||
t.Run("JobCleanup", func(t *testing.T) { testJobCleanup(t, rctx, ss) })
|
||||
}
|
||||
@@ -792,3 +793,162 @@ func testJobCleanup(t *testing.T, rctx request.CTX, ss store.Store) {
|
||||
require.NoError(t, err)
|
||||
assert.Len(t, jobs, 0)
|
||||
}
|
||||
|
||||
func testJobGetByTypeAndData(t *testing.T, rctx request.CTX, ss store.Store) {
|
||||
// Test setup - create test jobs with different types and data
|
||||
jobType := model.JobTypeAccessControlSync
|
||||
otherJobType := model.JobTypeDataRetention
|
||||
|
||||
// Job 1: Access control sync job with policy_id = "channel1"
|
||||
job1 := &model.Job{
|
||||
Id: model.NewId(),
|
||||
Type: jobType,
|
||||
Status: model.JobStatusPending,
|
||||
Data: map[string]string{
|
||||
"policy_id": "channel1",
|
||||
"extra": "data1",
|
||||
},
|
||||
}
|
||||
|
||||
// Job 2: Access control sync job with policy_id = "channel2"
|
||||
job2 := &model.Job{
|
||||
Id: model.NewId(),
|
||||
Type: jobType,
|
||||
Status: model.JobStatusInProgress,
|
||||
Data: map[string]string{
|
||||
"policy_id": "channel2",
|
||||
"extra": "data2",
|
||||
},
|
||||
}
|
||||
|
||||
// Job 3: Access control sync job with policy_id = "channel1" (same as job1)
|
||||
job3 := &model.Job{
|
||||
Id: model.NewId(),
|
||||
Type: jobType,
|
||||
Status: model.JobStatusSuccess,
|
||||
Data: map[string]string{
|
||||
"policy_id": "channel1",
|
||||
"extra": "data3",
|
||||
},
|
||||
}
|
||||
|
||||
// Job 4: Different job type with same policy_id
|
||||
job4 := &model.Job{
|
||||
Id: model.NewId(),
|
||||
Type: otherJobType,
|
||||
Status: model.JobStatusPending,
|
||||
Data: map[string]string{
|
||||
"policy_id": "channel1",
|
||||
},
|
||||
}
|
||||
|
||||
// Save all jobs
|
||||
_, err := ss.Job().Save(job1)
|
||||
require.NoError(t, err)
|
||||
defer func() { _, _ = ss.Job().Delete(job1.Id) }()
|
||||
|
||||
_, err = ss.Job().Save(job2)
|
||||
require.NoError(t, err)
|
||||
defer func() { _, _ = ss.Job().Delete(job2.Id) }()
|
||||
|
||||
_, err = ss.Job().Save(job3)
|
||||
require.NoError(t, err)
|
||||
defer func() { _, _ = ss.Job().Delete(job3.Id) }()
|
||||
|
||||
_, err = ss.Job().Save(job4)
|
||||
require.NoError(t, err)
|
||||
defer func() { _, _ = ss.Job().Delete(job4.Id) }()
|
||||
|
||||
t.Run("finds jobs by type and single data field", func(t *testing.T) {
|
||||
// Should find job1 and job3 (both have policy_id = "channel1" and correct type)
|
||||
jobs, err := ss.Job().GetByTypeAndData(rctx, jobType, map[string]string{
|
||||
"policy_id": "channel1",
|
||||
}, false)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, jobs, 2)
|
||||
|
||||
// Should contain job1 and job3
|
||||
jobIds := []string{jobs[0].Id, jobs[1].Id}
|
||||
assert.Contains(t, jobIds, job1.Id)
|
||||
assert.Contains(t, jobIds, job3.Id)
|
||||
})
|
||||
|
||||
t.Run("finds jobs by type and multiple data fields", func(t *testing.T) {
|
||||
// Should find only job1 (has both policy_id = "channel1" AND extra = "data1")
|
||||
jobs, err := ss.Job().GetByTypeAndData(rctx, jobType, map[string]string{
|
||||
"policy_id": "channel1",
|
||||
"extra": "data1",
|
||||
}, false)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, jobs, 1)
|
||||
assert.Equal(t, job1.Id, jobs[0].Id)
|
||||
})
|
||||
|
||||
t.Run("returns empty slice when no matches", func(t *testing.T) {
|
||||
// Should find nothing (no jobs with policy_id = "nonexistent")
|
||||
jobs, err := ss.Job().GetByTypeAndData(rctx, jobType, map[string]string{
|
||||
"policy_id": "nonexistent",
|
||||
}, false)
|
||||
require.NoError(t, err)
|
||||
assert.Len(t, jobs, 0)
|
||||
})
|
||||
|
||||
t.Run("filters by job type correctly", func(t *testing.T) {
|
||||
// Should find only job4 (different job type with same policy_id)
|
||||
jobs, err := ss.Job().GetByTypeAndData(rctx, otherJobType, map[string]string{
|
||||
"policy_id": "channel1",
|
||||
}, false)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, jobs, 1)
|
||||
assert.Equal(t, job4.Id, jobs[0].Id)
|
||||
})
|
||||
|
||||
// Test status parameter filtering
|
||||
t.Run("filters by single status", func(t *testing.T) {
|
||||
// Filter by single status should return only matching jobs
|
||||
jobs, err := ss.Job().GetByTypeAndData(rctx, jobType, map[string]string{
|
||||
"policy_id": "channel1",
|
||||
}, false, model.JobStatusPending)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, jobs, 1)
|
||||
assert.Equal(t, job1.Id, jobs[0].Id)
|
||||
assert.Equal(t, model.JobStatusPending, jobs[0].Status)
|
||||
})
|
||||
|
||||
t.Run("filters by multiple statuses", func(t *testing.T) {
|
||||
// Filter by multiple statuses should return jobs matching any status
|
||||
jobs, err := ss.Job().GetByTypeAndData(rctx, jobType, map[string]string{
|
||||
"policy_id": "channel1",
|
||||
}, false, model.JobStatusPending, model.JobStatusSuccess)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, jobs, 2)
|
||||
|
||||
// Verify both statuses are represented
|
||||
statuses := []string{jobs[0].Status, jobs[1].Status}
|
||||
assert.Contains(t, statuses, model.JobStatusPending)
|
||||
assert.Contains(t, statuses, model.JobStatusSuccess)
|
||||
})
|
||||
|
||||
t.Run("no status filter returns all statuses", func(t *testing.T) {
|
||||
// No status filter should return all jobs regardless of status
|
||||
jobs, err := ss.Job().GetByTypeAndData(rctx, jobType, map[string]string{
|
||||
"policy_id": "channel1",
|
||||
}, false) // No status parameters
|
||||
require.NoError(t, err)
|
||||
require.Len(t, jobs, 2) // job1 (pending), job3 (success) - both have policy_id=channel1
|
||||
|
||||
// Verify both statuses are present
|
||||
statuses := []string{jobs[0].Status, jobs[1].Status}
|
||||
assert.Contains(t, statuses, model.JobStatusPending)
|
||||
assert.Contains(t, statuses, model.JobStatusSuccess)
|
||||
})
|
||||
|
||||
t.Run("filters by non-existent status returns empty", func(t *testing.T) {
|
||||
// Invalid status filter should return empty result
|
||||
jobs, err := ss.Job().GetByTypeAndData(rctx, jobType, map[string]string{
|
||||
"policy_id": "channel1",
|
||||
}, false, model.JobStatusError)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, jobs, 0)
|
||||
})
|
||||
}
|
||||
|
||||
@@ -271,6 +271,43 @@ func (_m *JobStore) GetAllByTypesPage(rctx request.CTX, jobTypes []string, offse
|
||||
return r0, r1
|
||||
}
|
||||
|
||||
// GetByTypeAndData provides a mock function with given fields: rctx, jobType, data, useMaster, statuses
|
||||
func (_m *JobStore) GetByTypeAndData(rctx request.CTX, jobType string, data map[string]string, useMaster bool, statuses ...string) ([]*model.Job, error) {
|
||||
_va := make([]interface{}, len(statuses))
|
||||
for _i := range statuses {
|
||||
_va[_i] = statuses[_i]
|
||||
}
|
||||
var _ca []interface{}
|
||||
_ca = append(_ca, rctx, jobType, data, useMaster)
|
||||
_ca = append(_ca, _va...)
|
||||
ret := _m.Called(_ca...)
|
||||
|
||||
if len(ret) == 0 {
|
||||
panic("no return value specified for GetByTypeAndData")
|
||||
}
|
||||
|
||||
var r0 []*model.Job
|
||||
var r1 error
|
||||
if rf, ok := ret.Get(0).(func(request.CTX, string, map[string]string, bool, ...string) ([]*model.Job, error)); ok {
|
||||
return rf(rctx, jobType, data, useMaster, statuses...)
|
||||
}
|
||||
if rf, ok := ret.Get(0).(func(request.CTX, string, map[string]string, bool, ...string) []*model.Job); ok {
|
||||
r0 = rf(rctx, jobType, data, useMaster, statuses...)
|
||||
} else {
|
||||
if ret.Get(0) != nil {
|
||||
r0 = ret.Get(0).([]*model.Job)
|
||||
}
|
||||
}
|
||||
|
||||
if rf, ok := ret.Get(1).(func(request.CTX, string, map[string]string, bool, ...string) error); ok {
|
||||
r1 = rf(rctx, jobType, data, useMaster, statuses...)
|
||||
} else {
|
||||
r1 = ret.Error(1)
|
||||
}
|
||||
|
||||
return r0, r1
|
||||
}
|
||||
|
||||
// GetCountByStatusAndType provides a mock function with given fields: status, jobType
|
||||
func (_m *JobStore) GetCountByStatusAndType(status string, jobType string) (int64, error) {
|
||||
ret := _m.Called(status, jobType)
|
||||
|
||||
@@ -6134,6 +6134,10 @@
|
||||
"id": "app.job.get_count_by_status_and_type.app_error",
|
||||
"translation": "Unable to get the job count by status and type."
|
||||
},
|
||||
{
|
||||
"id": "app.job.get_existing_jobs.error",
|
||||
"translation": "Unable to get existing jobs."
|
||||
},
|
||||
{
|
||||
"id": "app.job.get_newest_job_by_status_and_type.app_error",
|
||||
"translation": "Unable to get the newest job by status and type."
|
||||
|
||||
+4
@@ -100,4 +100,8 @@
|
||||
color: var(--error-text);
|
||||
}
|
||||
}
|
||||
|
||||
.canceled-status-content {
|
||||
padding: 32px;
|
||||
}
|
||||
}
|
||||
|
||||
+51
-21
@@ -15,6 +15,7 @@ import * as ChannelActions from 'mattermost-redux/actions/channels';
|
||||
import {getChannel} from 'mattermost-redux/selectors/entities/channels';
|
||||
import {getTeam} from 'mattermost-redux/selectors/entities/teams';
|
||||
|
||||
import AlertBanner from 'components/alert_banner';
|
||||
import CodeBlock from 'components/code_block/code_block';
|
||||
|
||||
import type {GlobalState} from 'types/store';
|
||||
@@ -184,19 +185,29 @@ export default function JobDetailsModal({job, onExited}: Props): JSX.Element {
|
||||
}
|
||||
modalSubheaderText={
|
||||
<div className='modal-subheader-text'>
|
||||
<FormattedMessage
|
||||
id='admin.access_control.jobTable.details.subheader'
|
||||
defaultMessage='Finished at {finishedAt}'
|
||||
values={{
|
||||
finishedAt: new Date(job.last_activity_at).toLocaleString(),
|
||||
}}
|
||||
/>
|
||||
{job.status === 'canceled' && job.type.includes('access_control_sync') ? (
|
||||
<FormattedMessage
|
||||
id='admin.access_control.jobTable.details.subheader.canceled'
|
||||
defaultMessage='Canceled at {canceledAt}'
|
||||
values={{
|
||||
canceledAt: new Date(job.last_activity_at).toLocaleString(),
|
||||
}}
|
||||
/>
|
||||
) : (
|
||||
<FormattedMessage
|
||||
id='admin.access_control.jobTable.details.subheader'
|
||||
defaultMessage='Finished at {finishedAt}'
|
||||
values={{
|
||||
finishedAt: new Date(job.last_activity_at).toLocaleString(),
|
||||
}}
|
||||
/>
|
||||
)}
|
||||
</div>
|
||||
}
|
||||
show={true}
|
||||
bodyPadding={false}
|
||||
>
|
||||
{job.status === 'error' ? (
|
||||
{job.status === 'error' && (
|
||||
<div className='error-status-content'>
|
||||
<div className='error-status-content__title'>
|
||||
<FormattedMessage
|
||||
@@ -209,20 +220,39 @@ export default function JobDetailsModal({job, onExited}: Props): JSX.Element {
|
||||
language='json'
|
||||
/>
|
||||
</div>
|
||||
) : (
|
||||
job.type.includes('access_control_sync') && syncResults && (
|
||||
<SearchableSyncJobChannelList
|
||||
channels={filteredChannels}
|
||||
teams={teamLookup}
|
||||
channelsPerPage={pageSize}
|
||||
nextPage={() => {}}
|
||||
isSearch={Boolean(searchTerm)}
|
||||
search={setSearchTerm}
|
||||
onViewDetails={handleViewDetails}
|
||||
noResultsText={noResultsText}
|
||||
syncResults={syncResults}
|
||||
)}
|
||||
{job.status === 'canceled' && job.type.includes('access_control_sync') && (
|
||||
<div className='canceled-status-content'>
|
||||
<AlertBanner
|
||||
mode='warning'
|
||||
variant='app'
|
||||
title={
|
||||
<FormattedMessage
|
||||
id='admin.access_control.jobTable.syncResults.canceled.title'
|
||||
defaultMessage='Job Canceled'
|
||||
/>
|
||||
}
|
||||
message={
|
||||
<FormattedMessage
|
||||
id='admin.access_control.jobTable.syncResults.canceled.message'
|
||||
defaultMessage='This sync job was canceled, likely because a newer sync job was started for the same channel. Channel members were not updated.'
|
||||
/>
|
||||
}
|
||||
/>
|
||||
)
|
||||
</div>
|
||||
)}
|
||||
{job.status !== 'error' && !(job.status === 'canceled' && job.type.includes('access_control_sync')) && job.type.includes('access_control_sync') && syncResults && (
|
||||
<SearchableSyncJobChannelList
|
||||
channels={filteredChannels}
|
||||
teams={teamLookup}
|
||||
channelsPerPage={pageSize}
|
||||
nextPage={() => {}}
|
||||
isSearch={Boolean(searchTerm)}
|
||||
search={setSearchTerm}
|
||||
onViewDetails={handleViewDetails}
|
||||
noResultsText={noResultsText}
|
||||
syncResults={syncResults}
|
||||
/>
|
||||
)}
|
||||
|
||||
{selectedChannel && selectedChannelResults && (
|
||||
|
||||
+2
@@ -97,6 +97,7 @@ exports[`components/admin_console/access_control/policy_details/PolicyDetails sh
|
||||
<TableEditor
|
||||
actions={
|
||||
Object {
|
||||
"createAccessControlSyncJob": [MockFunction],
|
||||
"createJob": [MockFunction],
|
||||
"deleteChannelPolicy": [MockFunction],
|
||||
"getAccessControlFields": [MockFunction],
|
||||
@@ -310,6 +311,7 @@ exports[`components/admin_console/access_control/policy_details/PolicyDetails sh
|
||||
<TableEditor
|
||||
actions={
|
||||
Object {
|
||||
"createAccessControlSyncJob": [MockFunction],
|
||||
"createJob": [MockFunction],
|
||||
"deleteChannelPolicy": [MockFunction],
|
||||
"getAccessControlFields": [MockFunction],
|
||||
|
||||
+1
@@ -93,6 +93,7 @@ describe('components/admin_console/access_control/policy_details/PolicyDetails',
|
||||
deleteChannelPolicy: jest.fn(),
|
||||
getChannelMembers: jest.fn(),
|
||||
createJob: jest.fn(),
|
||||
createAccessControlSyncJob: jest.fn(),
|
||||
updateAccessControlPolicyActive: jest.fn(),
|
||||
validateExpressionAgainstRequester: jest.fn(),
|
||||
});
|
||||
|
||||
+3
-6
@@ -26,7 +26,6 @@ import TextSetting from 'components/widgets/settings/text_setting';
|
||||
|
||||
import {useChannelAccessControlActions} from 'hooks/useChannelAccessControlActions';
|
||||
import {getHistory} from 'utils/browser_history';
|
||||
import {JobTypes} from 'utils/constants';
|
||||
|
||||
import ChannelList from './channel_list';
|
||||
|
||||
@@ -226,11 +225,9 @@ function PolicyDetails({
|
||||
// --- Step 4: Create Job if necessary ---
|
||||
if (apply) {
|
||||
try {
|
||||
const job: JobTypeBase & { data: any } = {
|
||||
type: JobTypes.ACCESS_CONTROL_SYNC,
|
||||
data: {policy_id: currentPolicyId},
|
||||
};
|
||||
await actions.createJob(job);
|
||||
await abacActions.createAccessControlSyncJob({
|
||||
policy_id: currentPolicyId,
|
||||
});
|
||||
} catch (error) {
|
||||
setServerError(formatMessage({
|
||||
id: 'admin.access_control.policy.edit_policy.error.create_job',
|
||||
|
||||
+52
@@ -62,6 +62,7 @@ describe('components/channel_settings_modal/ChannelSettingsAccessRulesTab', () =
|
||||
deleteChannelPolicy: jest.fn(),
|
||||
getChannelMembers: jest.fn(),
|
||||
createJob: jest.fn(),
|
||||
createAccessControlSyncJob: jest.fn(),
|
||||
updateAccessControlPolicyActive: jest.fn(),
|
||||
validateExpressionAgainstRequester: jest.fn(),
|
||||
};
|
||||
@@ -942,6 +943,57 @@ describe('components/channel_settings_modal/ChannelSettingsAccessRulesTab', () =
|
||||
});
|
||||
});
|
||||
|
||||
test('should prevent duplicate save button clicks', async () => {
|
||||
// Ensure the searchUsers mock returns current user for validation to pass
|
||||
mockActions.searchUsers.mockResolvedValue({
|
||||
data: {
|
||||
users: [{
|
||||
id: 'current_user_id',
|
||||
username: 'testuser',
|
||||
first_name: 'Test',
|
||||
last_name: 'User',
|
||||
}],
|
||||
},
|
||||
});
|
||||
mockActions.getChannelMembers.mockResolvedValue({data: []});
|
||||
|
||||
renderWithContext(
|
||||
<ChannelSettingsAccessRulesTab {...baseProps}/>,
|
||||
initialState,
|
||||
);
|
||||
|
||||
await waitFor(() => {
|
||||
expect(screen.getByTestId('table-editor')).toBeInTheDocument();
|
||||
});
|
||||
|
||||
// Change expression to trigger unsaved state
|
||||
const onChangeCallback = MockedTableEditor.mock.calls[0][0].onChange;
|
||||
onChangeCallback('user.attributes.department == "Engineering"');
|
||||
|
||||
// Wait for save button to appear
|
||||
await waitFor(() => {
|
||||
expect(screen.getByText('Save')).toBeInTheDocument();
|
||||
});
|
||||
|
||||
const saveButton = screen.getByText('Save');
|
||||
|
||||
// Clear mock calls before testing duplicate prevention
|
||||
mockActions.saveChannelPolicy.mockClear();
|
||||
|
||||
// Click save button multiple times rapidly
|
||||
await userEvent.click(saveButton);
|
||||
await userEvent.click(saveButton);
|
||||
await userEvent.click(saveButton);
|
||||
|
||||
// Wait for async operations to complete
|
||||
await waitFor(() => {
|
||||
expect(mockActions.saveChannelPolicy).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
// Should only have been called once due to duplicate prevention
|
||||
expect(mockActions.saveChannelPolicy).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
test('should reset changes when Reset button is clicked', async () => {
|
||||
renderWithContext(
|
||||
<ChannelSettingsAccessRulesTab {...baseProps}/>,
|
||||
|
||||
+43
-24
@@ -1,12 +1,11 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
import React, {useState, useEffect, useCallback, useMemo} from 'react';
|
||||
import React, {useState, useEffect, useCallback, useMemo, useRef} from 'react';
|
||||
import {FormattedMessage, useIntl} from 'react-intl';
|
||||
import {useSelector} from 'react-redux';
|
||||
|
||||
import type {Channel} from '@mattermost/types/channels';
|
||||
import type {JobTypeBase} from '@mattermost/types/jobs';
|
||||
import type {UserPropertyField} from '@mattermost/types/properties';
|
||||
|
||||
import {getAccessControlSettings} from 'mattermost-redux/selectors/entities/access_control';
|
||||
@@ -19,7 +18,6 @@ import SaveChangesPanel, {type SaveChangesPanelState} from 'components/widgets/m
|
||||
|
||||
import {useChannelAccessControlActions} from 'hooks/useChannelAccessControlActions';
|
||||
import {useChannelSystemPolicies} from 'hooks/useChannelSystemPolicies';
|
||||
import {JobTypes} from 'utils/constants';
|
||||
|
||||
import type {GlobalState} from 'types/store';
|
||||
|
||||
@@ -389,13 +387,9 @@ function ChannelSettingsAccessRulesTab({
|
||||
// This ensures both user removal (always) and addition (conditional) happen immediately
|
||||
if (expression.trim()) {
|
||||
try {
|
||||
const job: JobTypeBase & { data: {policy_id: string} } = {
|
||||
type: JobTypes.ACCESS_CONTROL_SYNC,
|
||||
data: {
|
||||
policy_id: channel.id, // Sync only this specific channel policy
|
||||
},
|
||||
};
|
||||
await actions.createJob(job);
|
||||
await actions.createAccessControlSyncJob({
|
||||
policy_id: channel.id, // Sync only this specific channel policy
|
||||
});
|
||||
} catch (jobError) {
|
||||
// Log job creation error but don't fail the save operation
|
||||
// eslint-disable-next-line no-console
|
||||
@@ -478,27 +472,52 @@ function ChannelSettingsAccessRulesTab({
|
||||
}
|
||||
}, [expression, autoSyncMembers, formatMessage, validateSelfExclusion, calculateMembershipChanges, performSave, isEmptyRulesState]);
|
||||
|
||||
// Handle confirmation modal confirm
|
||||
// Prevent duplicate saves with immediate response
|
||||
const saveInProgressRef = useRef(false);
|
||||
|
||||
// Handle confirmation modal confirm - immediate response, no debounce
|
||||
const handleConfirmSave = useCallback(async () => {
|
||||
const success = await performSave();
|
||||
if (success) {
|
||||
setSaveChangesPanelState('saved');
|
||||
} else {
|
||||
setSaveChangesPanelState('error');
|
||||
// Prevent duplicate clicks - immediate response
|
||||
if (saveInProgressRef.current) {
|
||||
return;
|
||||
}
|
||||
|
||||
saveInProgressRef.current = true;
|
||||
|
||||
try {
|
||||
const success = await performSave();
|
||||
if (success) {
|
||||
setSaveChangesPanelState('saved');
|
||||
} else {
|
||||
setSaveChangesPanelState('error');
|
||||
}
|
||||
} finally {
|
||||
saveInProgressRef.current = false;
|
||||
}
|
||||
}, [performSave]);
|
||||
|
||||
// Handle save changes panel actions
|
||||
// Handle save changes panel actions - immediate response, no debounce
|
||||
const handleSaveChanges = useCallback(async () => {
|
||||
const result = await handleSave();
|
||||
|
||||
if (result === 'saved') {
|
||||
setSaveChangesPanelState('saved');
|
||||
} else if (result === 'error') {
|
||||
setSaveChangesPanelState('error');
|
||||
// Prevent duplicate clicks - immediate response
|
||||
if (saveInProgressRef.current) {
|
||||
return;
|
||||
}
|
||||
|
||||
// If result is 'confirmation_required', do nothing to the panel state
|
||||
saveInProgressRef.current = true;
|
||||
|
||||
try {
|
||||
const result = await handleSave();
|
||||
|
||||
if (result === 'saved') {
|
||||
setSaveChangesPanelState('saved');
|
||||
} else if (result === 'error') {
|
||||
setSaveChangesPanelState('error');
|
||||
}
|
||||
|
||||
// If result is 'confirmation_required', do nothing to the panel state
|
||||
} finally {
|
||||
saveInProgressRef.current = false;
|
||||
}
|
||||
}, [handleSave]);
|
||||
|
||||
const handleCancel = useCallback(() => {
|
||||
|
||||
@@ -18,6 +18,7 @@ import {
|
||||
updateAccessControlPolicyActive,
|
||||
deleteAccessControlPolicy,
|
||||
validateExpressionAgainstRequester,
|
||||
createAccessControlSyncJob,
|
||||
} from 'mattermost-redux/actions/access_control';
|
||||
import {getChannelMembers} from 'mattermost-redux/actions/channels';
|
||||
import {createJob} from 'mattermost-redux/actions/jobs';
|
||||
@@ -34,6 +35,7 @@ export interface ChannelAccessControlActions {
|
||||
createJob: (job: JobTypeBase & { data: any }) => Promise<ActionResult>;
|
||||
updateAccessControlPolicyActive: (policyId: string, active: boolean) => Promise<ActionResult>;
|
||||
validateExpressionAgainstRequester: (expression: string) => Promise<ActionResult<{requester_matches: boolean}>>;
|
||||
createAccessControlSyncJob: (jobData: {policy_id: string}) => Promise<ActionResult>;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -124,6 +126,13 @@ export const useChannelAccessControlActions = (channelId?: string): ChannelAcces
|
||||
validateExpressionAgainstRequester: (expression: string) => {
|
||||
return dispatch(validateExpressionAgainstRequester(expression, channelId));
|
||||
},
|
||||
|
||||
/**
|
||||
* Create an access control sync job with deduplication
|
||||
*/
|
||||
createAccessControlSyncJob: (jobData: {policy_id: string}) => {
|
||||
return dispatch(createAccessControlSyncJob(jobData));
|
||||
},
|
||||
}), [dispatch, channelId]);
|
||||
};
|
||||
|
||||
|
||||
@@ -269,6 +269,9 @@
|
||||
"admin.access_control.edit_policy.save_policy": "Save policy",
|
||||
"admin.access_control.edit_policy.serverError": "There are errors in the form above: {serverError}",
|
||||
"admin.access_control.jobTable.details.subheader": "Finished at {finishedAt}",
|
||||
"admin.access_control.jobTable.details.subheader.canceled": "Canceled at {canceledAt}",
|
||||
"admin.access_control.jobTable.syncResults.canceled.message": "This sync job was canceled, likely because a newer sync job was started for the same channel. Channel members were not updated.",
|
||||
"admin.access_control.jobTable.syncResults.canceled.title": "Job Canceled",
|
||||
"admin.access_control.policies.add_policy": "Add policy",
|
||||
"admin.access_control.policies.applies_to": "Applies to",
|
||||
"admin.access_control.policies.description": "Create policies containing attribute based access rules and the resources they apply to.",
|
||||
|
||||
@@ -172,3 +172,16 @@ export function validateExpressionAgainstRequester(expression: string, channelId
|
||||
return {data};
|
||||
};
|
||||
}
|
||||
|
||||
export function createAccessControlSyncJob(jobData: {policy_id: string}): ActionFuncAsync<any> {
|
||||
return async (dispatch, getState) => {
|
||||
let data;
|
||||
try {
|
||||
data = await Client4.createAccessControlSyncJob(jobData);
|
||||
} catch (error) {
|
||||
forceLogoutIfNecessary(error as ServerError, dispatch, getState);
|
||||
return {error};
|
||||
}
|
||||
return {data};
|
||||
};
|
||||
}
|
||||
|
||||
@@ -10,7 +10,7 @@ import {bindClientFunc} from './helpers';
|
||||
|
||||
import {General} from '../constants';
|
||||
|
||||
export function createJob(job: JobTypeBase) {
|
||||
export function createJob(job: JobTypeBase & { data?: any }) {
|
||||
return bindClientFunc({
|
||||
clientFunc: Client4.createJob,
|
||||
onSuccess: JobTypes.RECEIVED_JOB,
|
||||
|
||||
@@ -96,7 +96,7 @@ import type {
|
||||
OutgoingWebhook,
|
||||
SubmitDialogResponse,
|
||||
} from '@mattermost/types/integrations';
|
||||
import type {Job, JobTypeBase} from '@mattermost/types/jobs';
|
||||
import type {Job, JobType, JobTypeBase} from '@mattermost/types/jobs';
|
||||
import type {ServerLimits} from '@mattermost/types/limits';
|
||||
import type {
|
||||
MarketplaceApp,
|
||||
@@ -3186,7 +3186,7 @@ export default class Client4 {
|
||||
);
|
||||
};
|
||||
|
||||
createJob = (job: JobTypeBase) => {
|
||||
createJob = (job: JobTypeBase & { data?: any }) => {
|
||||
return this.doFetch<Job>(
|
||||
`${this.getJobsRoute()}`,
|
||||
{method: 'post', body: JSON.stringify(job)},
|
||||
@@ -4537,6 +4537,14 @@ export default class Client4 {
|
||||
);
|
||||
};
|
||||
|
||||
createAccessControlSyncJob = (jobData: {[key: string]: string}) => {
|
||||
const job = {
|
||||
type: 'access_control_sync' as JobType,
|
||||
data: jobData,
|
||||
};
|
||||
return this.createJob(job);
|
||||
};
|
||||
|
||||
getAccessControlFields = (after: string, limit: number, channelId?: string) => {
|
||||
const params = new URLSearchParams({after, limit: limit.toString()});
|
||||
if (channelId) {
|
||||
|
||||
Reference in New Issue
Block a user