diff --git a/server/channels/app/job.go b/server/channels/app/job.go index d3c23325bb8..abbe2b38e6b 100644 --- a/server/channels/app/job.go +++ b/server/channels/app/job.go @@ -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 { diff --git a/server/channels/app/job_test.go b/server/channels/app/job_test.go index 3f4a6ab1d9d..a7235de6c0b 100644 --- a/server/channels/app/job_test.go +++ b/server/channels/app/job_test.go @@ -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) diff --git a/server/channels/store/sqlstore/job_store.go b/server/channels/store/sqlstore/job_store.go index bc82a79f4dd..5db507bbbc5 100644 --- a/server/channels/store/sqlstore/job_store.go +++ b/server/channels/store/sqlstore/job_store.go @@ -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"). diff --git a/server/channels/store/store.go b/server/channels/store/store.go index a9942111678..1ed65e65d25 100644 --- a/server/channels/store/store.go +++ b/server/channels/store/store.go @@ -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 } diff --git a/server/channels/store/storetest/job_store.go b/server/channels/store/storetest/job_store.go index 4285b46446b..7ece32288b7 100644 --- a/server/channels/store/storetest/job_store.go +++ b/server/channels/store/storetest/job_store.go @@ -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) + }) +} diff --git a/server/channels/store/storetest/mocks/JobStore.go b/server/channels/store/storetest/mocks/JobStore.go index cfa6577a27e..0ae2fdffdd3 100644 --- a/server/channels/store/storetest/mocks/JobStore.go +++ b/server/channels/store/storetest/mocks/JobStore.go @@ -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) diff --git a/server/i18n/en.json b/server/i18n/en.json index 729c15ffc6b..b1b31087f9a 100644 --- a/server/i18n/en.json +++ b/server/i18n/en.json @@ -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." diff --git a/webapp/channels/src/components/admin_console/access_control/modals/job_details/job_details_modal.scss b/webapp/channels/src/components/admin_console/access_control/modals/job_details/job_details_modal.scss index 6b90b79ac56..768e9de8e4d 100644 --- a/webapp/channels/src/components/admin_console/access_control/modals/job_details/job_details_modal.scss +++ b/webapp/channels/src/components/admin_console/access_control/modals/job_details/job_details_modal.scss @@ -100,4 +100,8 @@ color: var(--error-text); } } + + .canceled-status-content { + padding: 32px; + } } diff --git a/webapp/channels/src/components/admin_console/access_control/modals/job_details/job_details_modal.tsx b/webapp/channels/src/components/admin_console/access_control/modals/job_details/job_details_modal.tsx index 90b494460ba..c46a9ae0f01 100644 --- a/webapp/channels/src/components/admin_console/access_control/modals/job_details/job_details_modal.tsx +++ b/webapp/channels/src/components/admin_console/access_control/modals/job_details/job_details_modal.tsx @@ -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={