From 9187c772b6ba7f005320aae85f2f33c906b4299c Mon Sep 17 00:00:00 2001 From: Ben Cooke Date: Mon, 17 Jun 2024 12:07:05 -0400 Subject: [PATCH] [MM-56074] mmctl job commands (#26855) * add job list and update job status command to mmctl --- api/v4/source/jobs.yaml | 60 +++++- e2e-tests/cypress/tests/support/api/role.js | 4 +- server/channels/api4/job.go | 105 ++++++++- server/channels/api4/job_local.go | 1 + server/channels/api4/job_test.go | 120 ++++++++++- server/channels/app/app_iface.go | 3 + server/channels/app/job.go | 59 +++++ .../app/opentracing/opentracing_layer.go | 61 ++++++ server/channels/app/permissions_migrations.go | 35 +++ .../opentracinglayer/opentracinglayer.go | 18 ++ .../channels/store/retrylayer/retrylayer.go | 21 ++ server/channels/store/sqlstore/job_store.go | 20 ++ server/channels/store/store.go | 1 + server/channels/store/storetest/job_store.go | 57 +++++ .../store/storetest/mocks/JobStore.go | 30 +++ .../channels/store/timerlayer/timerlayer.go | 16 ++ server/channels/testlib/store.go | 1 + server/cmd/mmctl/client/client.go | 3 +- server/cmd/mmctl/commands/export.go | 2 +- server/cmd/mmctl/commands/extract.go | 2 +- server/cmd/mmctl/commands/import.go | 65 +----- server/cmd/mmctl/commands/import_test.go | 4 +- server/cmd/mmctl/commands/job.go | 202 +++++++++++++++++ server/cmd/mmctl/commands/job_test.go | 204 ++++++++++++++++++ server/cmd/mmctl/commands/ldap.go | 2 +- server/cmd/mmctl/commands/ldap_test.go | 4 +- server/cmd/mmctl/docs/mmctl.rst | 1 + server/cmd/mmctl/docs/mmctl_job.rst | 42 ++++ server/cmd/mmctl/docs/mmctl_job_list.rst | 60 ++++++ server/cmd/mmctl/docs/mmctl_job_update.rst | 59 +++++ server/cmd/mmctl/mocks/client_mock.go | 23 +- server/i18n/en.json | 12 ++ server/public/model/client4.go | 21 +- server/public/model/job.go | 40 +++- server/public/model/job_test.go | 110 ++++++++++ server/public/model/migration.go | 1 + server/public/model/permission.go | 49 +++++ server/public/model/role.go | 6 + 38 files changed, 1423 insertions(+), 101 deletions(-) create mode 100644 server/cmd/mmctl/commands/job.go create mode 100644 server/cmd/mmctl/commands/job_test.go create mode 100644 server/cmd/mmctl/docs/mmctl_job.rst create mode 100644 server/cmd/mmctl/docs/mmctl_job_list.rst create mode 100644 server/cmd/mmctl/docs/mmctl_job_update.rst diff --git a/api/v4/source/jobs.yaml b/api/v4/source/jobs.yaml index eb7baca517..208dd93e7a 100644 --- a/api/v4/source/jobs.yaml +++ b/api/v4/source/jobs.yaml @@ -25,7 +25,17 @@ description: The number of jobs per page. schema: type: integer - default: 60 + default: 5 + - name: job_type + in: query + description: The type of jobs to fetch. + schema: + type: string + - name: status + in: query + description: The status of jobs to fetch. + schema: + type: string responses: "200": description: Job list retrieval successful @@ -223,3 +233,51 @@ $ref: "#/components/responses/Unauthorized" "403": $ref: "#/components/responses/Forbidden" + "/api/v4/jobs/{job_id}/status": + patch: + tags: + - jobs + summary: Update the status of a job + description: > + Update the status of a job. Valid status updates: + - 'in_progress' -> 'pending' + - 'in_progress' | 'pending' -> 'cancel_requested' + - 'cancel_requested' -> 'canceled' + + Add force to the body of the PATCH request to bypass the given rules, the only statuses you can go to are: pending, cancel_requested and canceled. This can have unexpected consequences and should be used with caution. + operationId: UpdateJobStatus + parameters: + - name: job_id + in: path + description: Job GUID + required: true + schema: + type: string + requestBody: + required: true + content: + application/json: + schema: + type: object + required: + - status + properties: + status: + type: string + description: The status you want to set + force: + type: boolean + description: Set this to true to bypass status restrictions + responses: + "200": + description: Status successfully set. + content: + application/json: + schema: + $ref: "#/components/schemas/StatusOK" + "400": + $ref: "#/components/responses/BadRequest" + "401": + $ref: "#/components/responses/Unauthorized" + "403": + $ref: "#/components/responses/Forbidden" diff --git a/e2e-tests/cypress/tests/support/api/role.js b/e2e-tests/cypress/tests/support/api/role.js index 4837f5053e..9d6e047a99 100644 --- a/e2e-tests/cypress/tests/support/api/role.js +++ b/e2e-tests/cypress/tests/support/api/role.js @@ -17,10 +17,10 @@ export const defaultRolesPermissions = { playbook_member: 'playbook_public_view playbook_public_manage_members playbook_public_manage_properties playbook_private_view playbook_private_manage_members playbook_private_manage_properties run_create', run_admin: 'run_manage_properties run_manage_members', run_member: 'run_view', - system_admin: 'sysconsole_write_environment_elasticsearch playbook_public_manage_properties sysconsole_write_authentication_ldap run_view manage_jobs manage_roles playbook_public_create manage_public_channel_properties sysconsole_read_plugins delete_post purge_elasticsearch_indexes sysconsole_read_integrations_bot_accounts read_data_retention_job manage_private_channel_members create_elasticsearch_post_indexing_job sysconsole_read_authentication_guest_access create_elasticsearch_post_aggregation_job join_public_teams sysconsole_read_site_public_links add_saml_idp_cert sysconsole_write_site_announcement_banner sysconsole_write_site_notices sysconsole_read_experimental_feature_flags sysconsole_read_site_users_and_teams manage_slash_commands sysconsole_read_authentication_ldap read_channel read_channel_content sysconsole_write_authentication_password list_users_without_team sysconsole_read_authentication_email add_saml_public_cert playbook_private_create promote_guest sysconsole_read_user_management_system_roles manage_public_channel_members create_data_retention_job add_saml_private_cert sysconsole_write_user_management_users sysconsole_read_compliance_compliance_monitoring playbook_public_manage_members sysconsole_write_environment_database sysconsole_write_user_management_teams playbook_private_manage_roles read_public_channel sysconsole_write_plugins sysconsole_read_authentication_openid sysconsole_write_user_management_groups sysconsole_write_site_file_sharing_and_downloads playbook_private_manage_properties sysconsole_read_site_customization join_public_channels add_user_to_team restore_custom_group download_compliance_export_result sysconsole_write_user_management_system_roles sysconsole_write_environment_session_lengths create_custom_group manage_private_channel_properties create_post_public remove_ldap_private_cert sysconsole_write_site_public_links import_team sysconsole_read_environment_developer sysconsole_read_environment_database sysconsole_read_environment_web_server use_channel_mentions view_team remove_others_reactions sysconsole_read_environment_session_lengths sysconsole_write_integrations_bot_accounts playbook_public_view use_group_mentions sysconsole_write_environment_web_server add_ldap_private_cert read_public_channel_groups invite_guest sysconsole_read_environment_smtp create_post sysconsole_read_about_edition_and_license sysconsole_read_authentication_signup sysconsole_read_authentication_saml sysconsole_read_environment_file_storage sysconsole_write_experimental_feature_flags sysconsole_write_site_localization sysconsole_write_environment_rate_limiting sysconsole_read_environment_rate_limiting sysconsole_read_products_boards get_saml_cert_status sysconsole_read_environment_high_availability manage_secure_connections read_compliance_export_job sysconsole_write_compliance_custom_terms_of_service read_user_access_token edit_post sysconsole_write_environment_logging sysconsole_read_environment_push_notification_server sysconsole_write_site_customization read_other_users_teams read_elasticsearch_post_aggregation_job sysconsole_write_compliance_data_retention_policy sysconsole_read_user_management_permissions sysconsole_read_site_emoji sysconsole_read_compliance_data_retention_policy read_license_information sysconsole_read_experimental_features read_deleted_posts sysconsole_read_environment_logging sysconsole_read_reporting_site_statistics test_elasticsearch sysconsole_read_site_posts add_reaction sysconsole_write_authentication_signup manage_outgoing_webhooks create_post_ephemeral sysconsole_read_environment_image_proxy invite_user manage_others_outgoing_webhooks create_user_access_token sysconsole_write_environment_image_proxy sysconsole_write_products_boards read_elasticsearch_post_indexing_job purge_bleve_indexes sysconsole_write_environment_performance_monitoring sysconsole_write_authentication_guest_access sysconsole_read_compliance_custom_terms_of_service edit_others_posts sysconsole_write_billing get_saml_metadata_from_idp sysconsole_write_authentication_saml create_post_bleve_indexes_job invalidate_caches sysconsole_write_experimental_bleve view_members manage_others_bots run_create join_private_teams convert_private_channel_to_public read_audits assign_bot read_jobs remove_user_from_team revoke_user_access_token manage_team sysconsole_read_reporting_server_logs get_public_link manage_others_slash_commands manage_system delete_public_channel read_private_channel_groups sysconsole_read_authentication_mfa delete_emojis list_private_teams create_emojis sysconsole_read_billing sysconsole_write_site_emoji invalidate_email_invite sysconsole_write_environment_file_storage sysconsole_write_compliance_compliance_monitoring remove_saml_public_cert sysconsole_read_compliance_compliance_export sysconsole_read_site_localization manage_team_roles list_public_teams get_logs sysconsole_write_integrations_integration_management sysconsole_read_integrations_cors manage_oauth manage_outgoing_oauth_connections delete_others_emojis sysconsole_write_integrations_gif manage_incoming_webhooks sysconsole_write_authentication_email create_private_channel playbook_private_make_public manage_bots add_ldap_public_cert remove_ldap_public_cert sysconsole_write_site_notifications sysconsole_write_environment_developer playbook_private_manage_members sysconsole_read_user_management_teams edit_custom_group remove_reaction playbook_public_manage_roles sysconsole_write_reporting_server_logs read_others_bots sysconsole_write_site_posts sysconsole_read_site_notifications sysconsole_read_authentication_password playbook_private_view manage_system_wide_oauth get_analytics list_team_channels sysconsole_write_user_management_channels delete_private_channel manage_custom_group_members test_s3 create_ldap_sync_job sysconsole_read_integrations_integration_management test_site_url recycle_database_connections sysconsole_read_site_announcement_banner test_email manage_shared_channels read_bots sysconsole_write_environment_smtp sysconsole_read_experimental_bleve sysconsole_write_environment_push_notification_server sysconsole_write_user_management_permissions sysconsole_read_environment_elasticsearch sysconsole_write_reporting_site_statistics sysconsole_write_site_users_and_teams demote_to_guest create_team test_ldap remove_saml_idp_cert delete_others_posts edit_other_users sysconsole_write_reporting_team_statistics sysconsole_read_integrations_gif sysconsole_read_site_notices sysconsole_write_about_edition_and_license manage_others_incoming_webhooks run_manage_members create_bot sysconsole_write_authentication_mfa sysconsole_read_user_management_users assign_system_admin_role sysconsole_write_experimental_features edit_brand create_group_channel sysconsole_write_authentication_openid create_direct_channel manage_license_information reload_config manage_channel_roles sysconsole_read_user_management_groups create_compliance_export_job read_ldap_sync_job upload_file sysconsole_read_site_file_sharing_and_downloads delete_custom_group sysconsole_read_user_management_channels sysconsole_write_compliance_compliance_export remove_saml_private_cert sysconsole_read_environment_performance_monitoring create_public_channel sysconsole_write_integrations_cors sysconsole_write_environment_high_availability playbook_public_make_private run_manage_properties sysconsole_read_reporting_team_statistics convert_public_channel_to_private add_bookmark_public_channel edit_bookmark_public_channel delete_bookmark_public_channel order_bookmark_public_channel add_bookmark_private_channel edit_bookmark_private_channel delete_bookmark_private_channel order_bookmark_private_channel', + system_admin: 'sysconsole_write_environment_elasticsearch playbook_public_manage_properties sysconsole_write_authentication_ldap run_view manage_jobs manage_roles playbook_public_create manage_public_channel_properties sysconsole_read_plugins delete_post purge_elasticsearch_indexes sysconsole_read_integrations_bot_accounts read_data_retention_job manage_private_channel_members create_elasticsearch_post_indexing_job manage_elasticsearch_post_indexing_job sysconsole_read_authentication_guest_access create_elasticsearch_post_aggregation_job manage_elasticsearch_post_aggregation_job join_public_teams sysconsole_read_site_public_links add_saml_idp_cert sysconsole_write_site_announcement_banner sysconsole_write_site_notices sysconsole_read_experimental_feature_flags sysconsole_read_site_users_and_teams manage_slash_commands sysconsole_read_authentication_ldap read_channel read_channel_content sysconsole_write_authentication_password list_users_without_team sysconsole_read_authentication_email add_saml_public_cert playbook_private_create promote_guest sysconsole_read_user_management_system_roles manage_public_channel_members create_data_retention_job manage_data_retention_job add_saml_private_cert sysconsole_write_user_management_users sysconsole_read_compliance_compliance_monitoring playbook_public_manage_members sysconsole_write_environment_database sysconsole_write_user_management_teams playbook_private_manage_roles read_public_channel sysconsole_write_plugins sysconsole_read_authentication_openid sysconsole_write_user_management_groups sysconsole_write_site_file_sharing_and_downloads playbook_private_manage_properties sysconsole_read_site_customization join_public_channels add_user_to_team restore_custom_group download_compliance_export_result sysconsole_write_user_management_system_roles sysconsole_write_environment_session_lengths create_custom_group manage_private_channel_properties create_post_public remove_ldap_private_cert sysconsole_write_site_public_links import_team sysconsole_read_environment_developer sysconsole_read_environment_database sysconsole_read_environment_web_server use_channel_mentions view_team remove_others_reactions sysconsole_read_environment_session_lengths sysconsole_write_integrations_bot_accounts playbook_public_view use_group_mentions sysconsole_write_environment_web_server add_ldap_private_cert read_public_channel_groups invite_guest sysconsole_read_environment_smtp create_post sysconsole_read_about_edition_and_license sysconsole_read_authentication_signup sysconsole_read_authentication_saml sysconsole_read_environment_file_storage sysconsole_write_experimental_feature_flags sysconsole_write_site_localization sysconsole_write_environment_rate_limiting sysconsole_read_environment_rate_limiting sysconsole_read_products_boards get_saml_cert_status sysconsole_read_environment_high_availability manage_secure_connections read_compliance_export_job sysconsole_write_compliance_custom_terms_of_service read_user_access_token edit_post sysconsole_write_environment_logging sysconsole_read_environment_push_notification_server sysconsole_write_site_customization read_other_users_teams read_elasticsearch_post_aggregation_job sysconsole_write_compliance_data_retention_policy sysconsole_read_user_management_permissions sysconsole_read_site_emoji sysconsole_read_compliance_data_retention_policy read_license_information sysconsole_read_experimental_features read_deleted_posts sysconsole_read_environment_logging sysconsole_read_reporting_site_statistics test_elasticsearch sysconsole_read_site_posts add_reaction sysconsole_write_authentication_signup manage_outgoing_webhooks create_post_ephemeral sysconsole_read_environment_image_proxy invite_user manage_others_outgoing_webhooks create_user_access_token sysconsole_write_environment_image_proxy sysconsole_write_products_boards read_elasticsearch_post_indexing_job purge_bleve_indexes sysconsole_write_environment_performance_monitoring sysconsole_write_authentication_guest_access sysconsole_read_compliance_custom_terms_of_service edit_others_posts sysconsole_write_billing get_saml_metadata_from_idp sysconsole_write_authentication_saml create_post_bleve_indexes_job manage_post_bleve_indexes_job invalidate_caches sysconsole_write_experimental_bleve view_members manage_others_bots run_create join_private_teams convert_private_channel_to_public read_audits assign_bot read_jobs remove_user_from_team revoke_user_access_token manage_team sysconsole_read_reporting_server_logs get_public_link manage_others_slash_commands manage_system delete_public_channel read_private_channel_groups sysconsole_read_authentication_mfa delete_emojis list_private_teams create_emojis sysconsole_read_billing sysconsole_write_site_emoji invalidate_email_invite sysconsole_write_environment_file_storage sysconsole_write_compliance_compliance_monitoring remove_saml_public_cert sysconsole_read_compliance_compliance_export sysconsole_read_site_localization manage_team_roles list_public_teams get_logs sysconsole_write_integrations_integration_management sysconsole_read_integrations_cors manage_oauth manage_outgoing_oauth_connections delete_others_emojis sysconsole_write_integrations_gif manage_incoming_webhooks sysconsole_write_authentication_email create_private_channel playbook_private_make_public manage_bots add_ldap_public_cert remove_ldap_public_cert sysconsole_write_site_notifications sysconsole_write_environment_developer playbook_private_manage_members sysconsole_read_user_management_teams edit_custom_group remove_reaction playbook_public_manage_roles sysconsole_write_reporting_server_logs read_others_bots sysconsole_write_site_posts sysconsole_read_site_notifications sysconsole_read_authentication_password playbook_private_view manage_system_wide_oauth get_analytics list_team_channels sysconsole_write_user_management_channels delete_private_channel manage_custom_group_members test_s3 create_ldap_sync_job manage_ldap_sync_job sysconsole_read_integrations_integration_management test_site_url recycle_database_connections sysconsole_read_site_announcement_banner test_email manage_shared_channels read_bots sysconsole_write_environment_smtp sysconsole_read_experimental_bleve sysconsole_write_environment_push_notification_server sysconsole_write_user_management_permissions sysconsole_read_environment_elasticsearch sysconsole_write_reporting_site_statistics sysconsole_write_site_users_and_teams demote_to_guest create_team test_ldap remove_saml_idp_cert delete_others_posts edit_other_users sysconsole_write_reporting_team_statistics sysconsole_read_integrations_gif sysconsole_read_site_notices sysconsole_write_about_edition_and_license manage_others_incoming_webhooks run_manage_members create_bot sysconsole_write_authentication_mfa sysconsole_read_user_management_users assign_system_admin_role sysconsole_write_experimental_features edit_brand create_group_channel sysconsole_write_authentication_openid create_direct_channel manage_license_information reload_config manage_channel_roles sysconsole_read_user_management_groups create_compliance_export_job manage_compliance_export_job read_ldap_sync_job upload_file sysconsole_read_site_file_sharing_and_downloads delete_custom_group sysconsole_read_user_management_channels sysconsole_write_compliance_compliance_export remove_saml_private_cert sysconsole_read_environment_performance_monitoring create_public_channel sysconsole_write_integrations_cors sysconsole_write_environment_high_availability playbook_public_make_private run_manage_properties sysconsole_read_reporting_team_statistics convert_public_channel_to_private add_bookmark_public_channel edit_bookmark_public_channel delete_bookmark_public_channel order_bookmark_public_channel add_bookmark_private_channel edit_bookmark_private_channel delete_bookmark_private_channel order_bookmark_private_channel', system_custom_group_admin: 'create_custom_group edit_custom_group delete_custom_group restore_custom_group manage_custom_group_members', system_guest: 'create_group_channel create_direct_channel', - system_manager: 'sysconsole_read_site_announcement_banner manage_private_channel_properties edit_brand read_private_channel_groups manage_private_channel_members manage_team_roles sysconsole_write_environment_session_lengths sysconsole_read_site_emoji sysconsole_write_environment_developer sysconsole_read_user_management_groups sysconsole_write_user_management_groups sysconsole_write_environment_rate_limiting delete_private_channel sysconsole_read_environment_performance_monitoring sysconsole_read_environment_rate_limiting sysconsole_write_user_management_teams sysconsole_write_integrations_integration_management sysconsole_write_site_public_links sysconsole_read_authentication_ldap sysconsole_write_integrations_cors reload_config sysconsole_write_user_management_channels sysconsole_read_environment_high_availability sysconsole_read_site_users_and_teams sysconsole_read_user_management_teams sysconsole_write_site_users_and_teams sysconsole_read_site_customization sysconsole_write_environment_high_availability sysconsole_read_integrations_bot_accounts sysconsole_read_authentication_guest_access sysconsole_read_site_public_links read_elasticsearch_post_indexing_job sysconsole_read_user_management_channels sysconsole_read_reporting_team_statistics invalidate_caches sysconsole_read_authentication_signup read_elasticsearch_post_aggregation_job sysconsole_write_environment_smtp manage_public_channel_members list_public_teams add_user_to_team sysconsole_read_environment_web_server sysconsole_read_site_localization get_logs sysconsole_write_site_posts sysconsole_write_integrations_bot_accounts sysconsole_write_user_management_permissions sysconsole_read_environment_elasticsearch sysconsole_read_environment_smtp list_private_teams read_public_channel_groups sysconsole_write_environment_file_storage sysconsole_write_integrations_gif manage_public_channel_properties sysconsole_write_environment_performance_monitoring sysconsole_write_site_notifications sysconsole_read_site_notifications sysconsole_read_environment_image_proxy sysconsole_write_site_announcement_banner sysconsole_write_site_emoji test_site_url sysconsole_read_integrations_gif sysconsole_write_environment_logging convert_public_channel_to_private get_analytics sysconsole_read_user_management_permissions sysconsole_write_environment_image_proxy test_elasticsearch recycle_database_connections sysconsole_write_site_localization sysconsole_read_reporting_server_logs create_elasticsearch_post_indexing_job sysconsole_read_reporting_site_statistics test_ldap delete_public_channel sysconsole_write_environment_push_notification_server read_license_information sysconsole_write_products_boards sysconsole_read_about_edition_and_license convert_private_channel_to_public sysconsole_read_integrations_integration_management create_elasticsearch_post_aggregation_job purge_elasticsearch_indexes sysconsole_read_environment_database join_public_teams sysconsole_read_authentication_email sysconsole_read_environment_push_notification_server view_team read_channel sysconsole_read_authentication_password read_ldap_sync_job sysconsole_read_integrations_cors sysconsole_read_environment_logging manage_team sysconsole_read_authentication_openid read_public_channel sysconsole_write_environment_elasticsearch sysconsole_read_plugins manage_channel_roles remove_user_from_team test_email sysconsole_write_site_file_sharing_and_downloads test_s3 sysconsole_read_site_file_sharing_and_downloads sysconsole_read_site_notices sysconsole_read_environment_file_storage join_private_teams sysconsole_read_products_boards sysconsole_read_environment_session_lengths sysconsole_write_environment_database sysconsole_read_authentication_saml sysconsole_read_authentication_mfa sysconsole_write_site_notices sysconsole_write_environment_web_server sysconsole_read_site_posts sysconsole_read_environment_developer sysconsole_write_site_customization manage_outgoing_oauth_connections', + system_manager: 'sysconsole_read_site_announcement_banner manage_private_channel_properties edit_brand read_private_channel_groups manage_private_channel_members manage_team_roles sysconsole_write_environment_session_lengths sysconsole_read_site_emoji sysconsole_write_environment_developer sysconsole_read_user_management_groups sysconsole_write_user_management_groups sysconsole_write_environment_rate_limiting delete_private_channel sysconsole_read_environment_performance_monitoring sysconsole_read_environment_rate_limiting sysconsole_write_user_management_teams sysconsole_write_integrations_integration_management sysconsole_write_site_public_links sysconsole_read_authentication_ldap sysconsole_write_integrations_cors reload_config sysconsole_write_user_management_channels sysconsole_read_environment_high_availability sysconsole_read_site_users_and_teams sysconsole_read_user_management_teams sysconsole_write_site_users_and_teams sysconsole_read_site_customization sysconsole_write_environment_high_availability sysconsole_read_integrations_bot_accounts sysconsole_read_authentication_guest_access sysconsole_read_site_public_links read_elasticsearch_post_indexing_job sysconsole_read_user_management_channels sysconsole_read_reporting_team_statistics invalidate_caches sysconsole_read_authentication_signup read_elasticsearch_post_aggregation_job sysconsole_write_environment_smtp manage_public_channel_members list_public_teams add_user_to_team sysconsole_read_environment_web_server sysconsole_read_site_localization get_logs sysconsole_write_site_posts sysconsole_write_integrations_bot_accounts sysconsole_write_user_management_permissions sysconsole_read_environment_elasticsearch sysconsole_read_environment_smtp list_private_teams read_public_channel_groups sysconsole_write_environment_file_storage sysconsole_write_integrations_gif manage_public_channel_properties sysconsole_write_environment_performance_monitoring sysconsole_write_site_notifications sysconsole_read_site_notifications sysconsole_read_environment_image_proxy sysconsole_write_site_announcement_banner sysconsole_write_site_emoji test_site_url sysconsole_read_integrations_gif sysconsole_write_environment_logging convert_public_channel_to_private get_analytics sysconsole_read_user_management_permissions sysconsole_write_environment_image_proxy test_elasticsearch recycle_database_connections sysconsole_write_site_localization sysconsole_read_reporting_server_logs create_elasticsearch_post_indexing_job manage_elasticsearch_post_indexing_job sysconsole_read_reporting_site_statistics test_ldap delete_public_channel sysconsole_write_environment_push_notification_server read_license_information sysconsole_write_products_boards sysconsole_read_about_edition_and_license convert_private_channel_to_public sysconsole_read_integrations_integration_management create_elasticsearch_post_aggregation_job manage_elasticsearch_post_aggregation_job purge_elasticsearch_indexes sysconsole_read_environment_database join_public_teams sysconsole_read_authentication_email sysconsole_read_environment_push_notification_server view_team read_channel sysconsole_read_authentication_password read_ldap_sync_job sysconsole_read_integrations_cors sysconsole_read_environment_logging manage_team sysconsole_read_authentication_openid read_public_channel sysconsole_write_environment_elasticsearch sysconsole_read_plugins manage_channel_roles remove_user_from_team test_email sysconsole_write_site_file_sharing_and_downloads test_s3 sysconsole_read_site_file_sharing_and_downloads sysconsole_read_site_notices sysconsole_read_environment_file_storage join_private_teams sysconsole_read_products_boards sysconsole_read_environment_session_lengths sysconsole_write_environment_database sysconsole_read_authentication_saml sysconsole_read_authentication_mfa sysconsole_write_site_notices sysconsole_write_environment_web_server sysconsole_read_site_posts sysconsole_read_environment_developer sysconsole_write_site_customization manage_outgoing_oauth_connections', system_post_all: 'use_group_mentions use_channel_mentions create_post', system_post_all_public: 'use_group_mentions use_channel_mentions create_post_public', system_read_only_admin: 'sysconsole_read_authentication_guest_access download_compliance_export_result sysconsole_read_compliance_data_retention_policy get_logs sysconsole_read_environment_file_storage read_channel sysconsole_read_integrations_integration_management sysconsole_read_compliance_custom_terms_of_service sysconsole_read_site_notices sysconsole_read_environment_rate_limiting sysconsole_read_about_edition_and_license read_public_channel sysconsole_read_experimental_features test_ldap sysconsole_read_user_management_permissions read_elasticsearch_post_aggregation_job sysconsole_read_environment_image_proxy sysconsole_read_compliance_compliance_export sysconsole_read_integrations_bot_accounts sysconsole_read_authentication_openid sysconsole_read_site_posts sysconsole_read_user_management_users sysconsole_read_experimental_feature_flags sysconsole_read_reporting_team_statistics sysconsole_read_site_localization read_private_channel_groups sysconsole_read_site_file_sharing_and_downloads sysconsole_read_user_management_channels sysconsole_read_authentication_email read_data_retention_job read_audits sysconsole_read_plugins view_team get_analytics sysconsole_read_user_management_groups sysconsole_read_experimental_bleve sysconsole_read_products_boards read_compliance_export_job sysconsole_read_environment_logging sysconsole_read_authentication_signup sysconsole_read_environment_smtp sysconsole_read_environment_session_lengths sysconsole_read_environment_developer sysconsole_read_environment_high_availability read_ldap_sync_job sysconsole_read_environment_performance_monitoring sysconsole_read_authentication_saml read_public_channel_groups sysconsole_read_integrations_gif sysconsole_read_authentication_mfa list_public_teams sysconsole_read_environment_database list_private_teams sysconsole_read_authentication_ldap sysconsole_read_compliance_compliance_monitoring sysconsole_read_site_notifications sysconsole_read_site_announcement_banner read_other_users_teams sysconsole_read_authentication_password sysconsole_read_environment_push_notification_server sysconsole_read_site_users_and_teams sysconsole_read_site_public_links sysconsole_read_site_emoji sysconsole_read_environment_elasticsearch read_license_information sysconsole_read_integrations_cors sysconsole_read_user_management_teams sysconsole_read_reporting_server_logs sysconsole_read_site_customization sysconsole_read_reporting_site_statistics sysconsole_read_environment_web_server read_elasticsearch_post_indexing_job', diff --git a/server/channels/api4/job.go b/server/channels/api4/job.go index eeca37c3f4..6797601735 100644 --- a/server/channels/api4/job.go +++ b/server/channels/api4/job.go @@ -23,6 +23,7 @@ func (api *API) InitJob() { api.BaseRoutes.Jobs.Handle("/{job_id:[A-Za-z0-9]+}/download", api.APISessionRequiredTrustRequester(downloadJob)).Methods("GET") api.BaseRoutes.Jobs.Handle("/{job_id:[A-Za-z0-9]+}/cancel", api.APISessionRequired(cancelJob)).Methods("POST") api.BaseRoutes.Jobs.Handle("/type/{job_type:[A-Za-z0-9_-]+}", api.APISessionRequired(getJobsByType)).Methods("GET") + api.BaseRoutes.Jobs.Handle("/{job_id:[A-Za-z0-9]+}/status", api.APISessionRequired(updateJobStatus)).Methods("PATCH") } func getJob(c *Context, w http.ResponseWriter, r *http.Request) { @@ -147,23 +148,58 @@ func getJobs(c *Context, w http.ResponseWriter, r *http.Request) { return } + jobType := r.URL.Query().Get("job_type") var validJobTypes []string - for _, jobType := range model.AllJobTypes { + + if jobType != "" { + isValidJobType := model.IsValidJobType(jobType) + if !isValidJobType { + c.SetInvalidURLParam("job_type") + return + } hasPermission, permissionRequired := c.App.SessionHasPermissionToReadJob(*c.AppContext.Session(), jobType) if permissionRequired == nil { - c.Logger.Warn("The job types of a job you are trying to retrieve does not contain permissions", mlog.String("jobType", jobType)) - continue + c.Err = model.NewAppError("getJobsByType", "api.job.retrieve.nopermissions", nil, "", http.StatusBadRequest) + return } - if hasPermission { - validJobTypes = append(validJobTypes, jobType) + if !hasPermission { + c.SetPermissionError(permissionRequired) + return + } + validJobTypes = append(validJobTypes, jobType) + } else { + for _, jType := range model.AllJobTypes { + hasPermission, permissionRequired := c.App.SessionHasPermissionToReadJob(*c.AppContext.Session(), jType) + if permissionRequired == nil { + c.Logger.Warn("The job types of a job you are trying to retrieve does not contain permissions", mlog.String("jobType", jType)) + continue + } + if hasPermission { + validJobTypes = append(validJobTypes, jType) + } } } + if len(validJobTypes) == 0 { c.SetPermissionError() return } - jobs, appErr := c.App.GetJobsByTypesPage(c.AppContext, validJobTypes, c.Params.Page, c.Params.PerPage) + status := r.URL.Query().Get("status") + isValidStatus := model.IsValidJobStatus(status) + if status != "" && !isValidStatus { + c.Err = model.NewAppError("getJobs", "api.job.status.invalid", nil, "", http.StatusBadRequest) + } + + var jobs []*model.Job + var appErr *model.AppError + + if status == "" { + jobs, appErr = c.App.GetJobsByTypesPage(c.AppContext, validJobTypes, c.Params.Page, c.Params.PerPage) + } else { + jobs, appErr = c.App.GetJobsByTypeAndStatus(c.AppContext, validJobTypes, status, c.Params.Page, c.Params.PerPage) + } + if appErr != nil { c.Err = appErr return @@ -248,3 +284,60 @@ func cancelJob(c *Context, w http.ResponseWriter, r *http.Request) { ReturnStatusOK(w) } + +func updateJobStatus(c *Context, w http.ResponseWriter, r *http.Request) { + c.RequireJobId() + if c.Err != nil { + return + } + + auditRec := c.MakeAuditRecord("updateJobStatus", audit.Fail) + defer c.LogAuditRec(auditRec) + audit.AddEventParameter(auditRec, "job_id", c.Params.JobId) + + props := model.StringInterfaceFromJSON(r.Body) + status, ok := props["status"].(string) + if !ok { + c.SetInvalidParam("status") + return + } + + force, ok := props["force"].(bool) + if !ok { + force = false + } + + job, err := c.App.GetJob(c.AppContext, c.Params.JobId) + if err != nil { + c.Err = err + return + } + + auditRec.AddEventPriorState(job) + auditRec.AddEventObjectType("job") + + hasPermission, permissionRequired := c.App.SessionHasPermissionToManageJob(*c.AppContext.Session(), job) + if permissionRequired == nil { + c.Err = model.NewAppError("updateJobStatus", "api.job.unable_to_manage_job.incorrect_job_type", nil, "", http.StatusBadRequest) + return + } + + if !hasPermission { + c.SetPermissionError(permissionRequired) + return + } + + if !force && !job.IsValidStatusChange(status) { + c.Err = model.NewAppError("updateJobStatus", "api.job.status.invalid", nil, "", http.StatusBadRequest) + return + } + + if err := c.App.UpdateJobStatus(c.AppContext, job, status); err != nil { + c.Err = err + return + } + + auditRec.Success() + + ReturnStatusOK(w) +} diff --git a/server/channels/api4/job_local.go b/server/channels/api4/job_local.go index fe5ed10aa0..0d14dc4499 100644 --- a/server/channels/api4/job_local.go +++ b/server/channels/api4/job_local.go @@ -9,4 +9,5 @@ func (api *API) InitJobLocal() { api.BaseRoutes.Jobs.Handle("/{job_id:[A-Za-z0-9]+}", api.APILocal(getJob)).Methods("GET") api.BaseRoutes.Jobs.Handle("/{job_id:[A-Za-z0-9]+}/cancel", api.APILocal(cancelJob)).Methods("POST") api.BaseRoutes.Jobs.Handle("/type/{job_type:[A-Za-z0-9_-]+}", api.APILocal(getJobsByType)).Methods("GET") + api.BaseRoutes.Jobs.Handle("/{job_id:[A-Za-z0-9]+}/status", api.APILocal(updateJobStatus)).Methods("PATCH") } diff --git a/server/channels/api4/job_test.go b/server/channels/api4/job_test.go index 8dada07306..b96358fda4 100644 --- a/server/channels/api4/job_test.go +++ b/server/channels/api4/job_test.go @@ -102,6 +102,12 @@ func TestGetJobs(t *testing.T) { Type: jobType, CreateAt: t0 + 2, }, + { + Id: model.NewId(), + Type: model.JobTypeLdapSync, + CreateAt: t0 + 3, + Status: model.JobStatusPending, + }, } for _, job := range jobs { @@ -110,21 +116,47 @@ func TestGetJobs(t *testing.T) { defer th.App.Srv().Store().Job().Delete(job.Id) } - received, _, err := th.SystemAdminClient.GetJobs(context.Background(), 0, 2) - require.NoError(t, err) + t.Run("Get 2 jobs", func(t *testing.T) { + received, _, err := th.SystemAdminClient.GetJobs(context.Background(), "", "", 0, 2) + require.NoError(t, err) - require.Len(t, received, 2, "received wrong number of jobs") - require.Equal(t, jobs[2].Id, received[0].Id, "should've received newest job first") - require.Equal(t, jobs[0].Id, received[1].Id, "should've received second newest job second") + require.Len(t, received, 2, "received wrong number of jobs") + require.Equal(t, jobs[3].Id, received[0].Id, "should've received newest job first") + require.Equal(t, jobs[2].Id, received[1].Id, "should've received second newest job second") + }) - received, _, err = th.SystemAdminClient.GetJobs(context.Background(), 1, 2) - require.NoError(t, err) + t.Run("Get oldest job using paging", func(t *testing.T) { + received, _, err := th.SystemAdminClient.GetJobs(context.Background(), "", "", 1, 3) + require.NoError(t, err) + require.Equal(t, jobs[1].Id, received[0].Id, "should've received oldest job last") + }) - require.Equal(t, jobs[1].Id, received[0].Id, "should've received oldest job last") + t.Run("Return error fetching job without permissions", func(t *testing.T) { + _, resp, err := th.Client.GetJobs(context.Background(), "", "", 0, 60) + require.Error(t, err) + CheckForbiddenStatus(t, resp) + }) - _, resp, err := th.Client.GetJobs(context.Background(), 0, 60) - require.Error(t, err) - CheckForbiddenStatus(t, resp) + t.Run("Get job by type", func(t *testing.T) { + received, _, err := th.SystemAdminClient.GetJobs(context.Background(), model.JobTypeLdapSync, "", 0, 3) + require.NoError(t, err) + require.Len(t, received, 1, "received wrong number of jobs") + require.Equal(t, jobs[3].Id, received[0].Id, "should've received the ldap sync job") + }) + + t.Run("Get job by status", func(t *testing.T) { + received, _, err := th.SystemAdminClient.GetJobs(context.Background(), "", model.JobStatusPending, 0, 3) + require.NoError(t, err) + require.Len(t, received, 1, "received wrong number of jobs") + require.Equal(t, jobs[3].Id, received[0].Id, "should've received the ldap sync job") + }) + + t.Run("Get job by type and status", func(t *testing.T) { + received, _, err := th.SystemAdminClient.GetJobs(context.Background(), model.JobTypeLdapSync, model.JobStatusPending, 0, 3) + require.NoError(t, err) + require.Len(t, received, 1, "received wrong number of jobs") + require.Equal(t, jobs[3].Id, received[0].Id, "should've received the ldap sync job") + }) } func TestGetJobsByType(t *testing.T) { @@ -336,3 +368,69 @@ func TestCancelJob(t *testing.T) { require.Error(t, err) CheckNotFoundStatus(t, resp) } + +func TestUpdateJobStatus(t *testing.T) { + th := Setup(t) + defer th.TearDown() + + jobType := model.JobTypeDataRetention + jobs := []*model.Job{ + { + Id: model.NewId(), + Type: jobType, + Status: model.JobStatusPending, + }, + { + Id: model.NewId(), + Type: jobType, + Status: model.JobStatusInProgress, + }, + { + Id: model.NewId(), + Type: jobType, + Status: model.JobStatusSuccess, + }, + { + Id: model.NewId(), + Type: jobType, + Status: model.JobStatusPending, + }, + } + + for _, job := range jobs { + _, err := th.App.Srv().Store().Job().Save(job) + require.NoError(t, err) + defer th.App.Srv().Store().Job().Delete(job.Id) + } + + t.Run("Fail to update job status without permission", func(t *testing.T) { + resp, err := th.Client.UpdateJobStatus(context.Background(), jobs[0].Id, model.JobStatusCancelRequested, false) + require.Error(t, err) + CheckForbiddenStatus(t, resp) + }) + + t.Run("Change a pending job to cancel requested without force with sysadmin client", func(t *testing.T) { + _, err := th.SystemAdminClient.UpdateJobStatus(context.Background(), jobs[0].Id, model.JobStatusCancelRequested, false) + require.NoError(t, err) + }) + + t.Run("Change a pending job to cancel requested without force with local client", func(t *testing.T) { + _, err := th.LocalClient.UpdateJobStatus(context.Background(), jobs[3].Id, model.JobStatusCancelRequested, false) + require.NoError(t, err) + }) + + t.Run("Fail to change a pending job to canceled without force", func(t *testing.T) { + th.TestForSystemAdminAndLocal(t, func(t *testing.T, client *model.Client4) { + resp, err := client.UpdateJobStatus(context.Background(), jobs[0].Id, model.JobStatusCanceled, false) + require.Error(t, err) + CheckBadRequestStatus(t, resp) + }) + }) + + t.Run("Change a pending job to canceled with force", func(t *testing.T) { + th.TestForSystemAdminAndLocal(t, func(t *testing.T, client *model.Client4) { + _, err := client.UpdateJobStatus(context.Background(), jobs[0].Id, model.JobStatusCanceled, true) + require.NoError(t, err) + }) + }) +} diff --git a/server/channels/app/app_iface.go b/server/channels/app/app_iface.go index 5bae578084..83ea25b1bd 100644 --- a/server/channels/app/app_iface.go +++ b/server/channels/app/app_iface.go @@ -720,6 +720,7 @@ type AppIface interface { GetIncomingWebhooksPageByUser(userID string, page, perPage int) ([]*model.IncomingWebhook, *model.AppError) GetJob(c request.CTX, id string) (*model.Job, *model.AppError) GetJobsByType(c request.CTX, jobType string, offset int, limit int) ([]*model.Job, *model.AppError) + GetJobsByTypeAndStatus(c request.CTX, jobTypes []string, status string, page int, perPage int) ([]*model.Job, *model.AppError) GetJobsByTypePage(c request.CTX, jobType string, page int, perPage int) ([]*model.Job, *model.AppError) GetJobsByTypes(c request.CTX, jobTypes []string, offset int, limit int) ([]*model.Job, *model.AppError) GetJobsByTypesPage(c request.CTX, jobType []string, page int, perPage int) ([]*model.Job, *model.AppError) @@ -1099,6 +1100,7 @@ type AppIface interface { SessionHasPermissionToChannelByPost(session model.Session, postID string, permission *model.Permission) bool SessionHasPermissionToCreateJob(session model.Session, job *model.Job) (bool, *model.Permission) SessionHasPermissionToGroup(session model.Session, groupID string, permission *model.Permission) bool + SessionHasPermissionToManageJob(session model.Session, job *model.Job) (bool, *model.Permission) SessionHasPermissionToReadJob(session model.Session, jobType string) (bool, *model.Permission) SessionHasPermissionToTeam(session model.Session, teamID string, permission *model.Permission) bool SessionHasPermissionToUser(session model.Session, userID string) bool @@ -1173,6 +1175,7 @@ type AppIface interface { UpdateHashedPassword(user *model.User, newHashedPassword string) *model.AppError UpdateHashedPasswordByUserId(userID, newHashedPassword string) *model.AppError UpdateIncomingWebhook(oldHook, updatedHook *model.IncomingWebhook) (*model.IncomingWebhook, *model.AppError) + UpdateJobStatus(c request.CTX, job *model.Job, newStatus string) *model.AppError UpdateMfa(c request.CTX, activate bool, userID, token string) *model.AppError UpdateMobileAppBadge(userID string) UpdateOAuthApp(oldApp, updatedApp *model.OAuthApp) (*model.OAuthApp, *model.AppError) diff --git a/server/channels/app/job.go b/server/channels/app/job.go index 42871d380e..a921d88a76 100644 --- a/server/channels/app/job.go +++ b/server/channels/app/job.go @@ -52,6 +52,14 @@ func (a *App) GetJobsByTypes(c request.CTX, jobTypes []string, offset int, limit return jobs, nil } +func (a *App) GetJobsByTypeAndStatus(c request.CTX, jobTypes []string, status string, page int, perPage int) ([]*model.Job, *model.AppError) { + jobs, err := a.Srv().Store().Job().GetAllByTypeAndStatusPage(c, jobTypes, status, page*perPage, perPage) + if err != nil { + return nil, model.NewAppError("GetAllByTypeAndStatusPage", "app.job.get_all.app_error", nil, "", http.StatusInternalServerError).Wrap(err) + } + return jobs, nil +} + func (a *App) CreateJob(c request.CTX, job *model.Job) (*model.Job, *model.AppError) { return a.Srv().Jobs.CreateJob(c, job.Type, job.Data) } @@ -60,6 +68,19 @@ func (a *App) CancelJob(c request.CTX, jobId string) *model.AppError { return a.Srv().Jobs.RequestCancellation(c, jobId) } +func (a *App) UpdateJobStatus(c request.CTX, job *model.Job, newStatus string) *model.AppError { + switch newStatus { + case model.JobStatusPending: + return a.Srv().Jobs.SetJobPending(job) + case model.JobStatusCancelRequested: + return a.Srv().Jobs.RequestCancellation(c, job.Id) + case model.JobStatusCanceled: + return a.Srv().Jobs.SetJobCanceled(job) + default: + return model.NewAppError("UpdateJobStatus", "app.job.update_status.app_error", nil, "", http.StatusInternalServerError) + } +} + func (a *App) SessionHasPermissionToCreateJob(session model.Session, job *model.Job) (bool, *model.Permission) { switch job.Type { case model.JobTypeBlevePostIndexing: @@ -92,6 +113,44 @@ func (a *App) SessionHasPermissionToCreateJob(session model.Session, job *model. return false, nil } +func (a *App) SessionHasPermissionToManageJob(session model.Session, job *model.Job) (bool, *model.Permission) { + var permission *model.Permission + + switch job.Type { + case model.JobTypeBlevePostIndexing: + permission = model.PermissionManagePostBleveIndexesJob + case model.JobTypeDataRetention: + permission = model.PermissionManageDataRetentionJob + case model.JobTypeMessageExport: + permission = model.PermissionManageComplianceExportJob + case model.JobTypeElasticsearchPostIndexing: + permission = model.PermissionManageElasticsearchPostIndexingJob + case model.JobTypeElasticsearchPostAggregation: + permission = model.PermissionManageElasticsearchPostAggregationJob + case model.JobTypeLdapSync: + permission = model.PermissionManageLdapSyncJob + case + model.JobTypeMigrations, + model.JobTypePlugins, + model.JobTypeProductNotices, + model.JobTypeExpiryNotify, + model.JobTypeActiveUsers, + model.JobTypeImportProcess, + model.JobTypeImportDelete, + model.JobTypeExportProcess, + model.JobTypeExportDelete, + model.JobTypeCloud, + model.JobTypeExtractContent: + permission = model.PermissionManageJobs + } + + if permission == nil { + return false, nil + } + + return a.SessionHasPermissionTo(session, permission), permission +} + func (a *App) SessionHasPermissionToReadJob(session model.Session, jobType string) (bool, *model.Permission) { switch jobType { case model.JobTypeDataRetention: diff --git a/server/channels/app/opentracing/opentracing_layer.go b/server/channels/app/opentracing/opentracing_layer.go index 776457db7e..4665d158eb 100644 --- a/server/channels/app/opentracing/opentracing_layer.go +++ b/server/channels/app/opentracing/opentracing_layer.go @@ -7308,6 +7308,28 @@ func (a *OpenTracingAppLayer) GetJobsByType(c request.CTX, jobType string, offse return resultVar0, resultVar1 } +func (a *OpenTracingAppLayer) GetJobsByTypeAndStatus(c request.CTX, jobTypes []string, status string, page int, perPage int) ([]*model.Job, *model.AppError) { + origCtx := a.ctx + span, newCtx := tracing.StartSpanWithParentByContext(a.ctx, "app.GetJobsByTypeAndStatus") + + a.ctx = newCtx + a.app.Srv().Store().SetContext(newCtx) + defer func() { + a.app.Srv().Store().SetContext(origCtx) + a.ctx = origCtx + }() + + defer span.Finish() + resultVar0, resultVar1 := a.app.GetJobsByTypeAndStatus(c, jobTypes, status, page, perPage) + + if resultVar1 != nil { + span.LogFields(spanlog.Error(resultVar1)) + ext.Error.Set(span, true) + } + + return resultVar0, resultVar1 +} + func (a *OpenTracingAppLayer) GetJobsByTypePage(c request.CTX, jobType string, page int, perPage int) ([]*model.Job, *model.AppError) { origCtx := a.ctx span, newCtx := tracing.StartSpanWithParentByContext(a.ctx, "app.GetJobsByTypePage") @@ -16316,6 +16338,23 @@ func (a *OpenTracingAppLayer) SessionHasPermissionToManageBot(rctx request.CTX, return resultVar0 } +func (a *OpenTracingAppLayer) SessionHasPermissionToManageJob(session model.Session, job *model.Job) (bool, *model.Permission) { + origCtx := a.ctx + span, newCtx := tracing.StartSpanWithParentByContext(a.ctx, "app.SessionHasPermissionToManageJob") + + a.ctx = newCtx + a.app.Srv().Store().SetContext(newCtx) + defer func() { + a.app.Srv().Store().SetContext(origCtx) + a.ctx = origCtx + }() + + defer span.Finish() + resultVar0, resultVar1 := a.app.SessionHasPermissionToManageJob(session, job) + + return resultVar0, resultVar1 +} + func (a *OpenTracingAppLayer) SessionHasPermissionToReadJob(session model.Session, jobType string) (bool, *model.Permission) { origCtx := a.ctx span, newCtx := tracing.StartSpanWithParentByContext(a.ctx, "app.SessionHasPermissionToReadJob") @@ -18083,6 +18122,28 @@ func (a *OpenTracingAppLayer) UpdateIncomingWebhook(oldHook *model.IncomingWebho return resultVar0, resultVar1 } +func (a *OpenTracingAppLayer) UpdateJobStatus(c request.CTX, job *model.Job, newStatus string) *model.AppError { + origCtx := a.ctx + span, newCtx := tracing.StartSpanWithParentByContext(a.ctx, "app.UpdateJobStatus") + + a.ctx = newCtx + a.app.Srv().Store().SetContext(newCtx) + defer func() { + a.app.Srv().Store().SetContext(origCtx) + a.ctx = origCtx + }() + + defer span.Finish() + resultVar0 := a.app.UpdateJobStatus(c, job, newStatus) + + if resultVar0 != nil { + span.LogFields(spanlog.Error(resultVar0)) + ext.Error.Set(span, true) + } + + return resultVar0 +} + func (a *OpenTracingAppLayer) UpdateMfa(c request.CTX, activate bool, userID string, token string) *model.AppError { origCtx := a.ctx span, newCtx := tracing.StartSpanWithParentByContext(a.ctx, "app.UpdateMfa") diff --git a/server/channels/app/permissions_migrations.go b/server/channels/app/permissions_migrations.go index c68dea55ee..b51763b91b 100644 --- a/server/channels/app/permissions_migrations.go +++ b/server/channels/app/permissions_migrations.go @@ -1181,6 +1181,40 @@ func (a *App) getAddChannelBookmarksPermissionsMigration() (permissionsMap, erro return transformations, nil } +func (a *App) getAddManageJobAncillaryPermissionsMigration() (permissionsMap, error) { + transformations := []permissionTransformation{} + + transformations = append(transformations, permissionTransformation{ + On: permissionExists(model.PermissionSysconsoleWriteAuthenticationLdap.Id), + Add: []string{model.PermissionManageLdapSyncJob.Id}, + }) + + transformations = append(transformations, permissionTransformation{ + On: permissionExists(model.PermissionSysconsoleWriteComplianceDataRetentionPolicy.Id), + Add: []string{model.PermissionManageDataRetentionJob.Id}, + }) + + transformations = append(transformations, permissionTransformation{ + On: permissionExists(model.PermissionSysconsoleWriteExperimentalBleve.Id), + Add: []string{model.PermissionManagePostBleveIndexesJob.Id}, + }) + + transformations = append(transformations, permissionTransformation{ + On: permissionExists(model.PermissionSysconsoleWriteComplianceComplianceExport.Id), + Add: []string{model.PermissionManageComplianceExportJob.Id}, + }) + + transformations = append(transformations, permissionTransformation{ + On: permissionExists(model.PermissionSysconsoleWriteEnvironmentElasticsearch.Id), + Add: []string{ + model.PermissionManageElasticsearchPostIndexingJob.Id, + model.PermissionManageElasticsearchPostAggregationJob.Id, + }, + }) + + return transformations, nil +} + // DoPermissionsMigrations execute all the permissions migrations need by the current version. func (a *App) DoPermissionsMigrations() error { return a.Srv().doPermissionsMigrations() @@ -1228,6 +1262,7 @@ func (s *Server) doPermissionsMigrations() error { {Key: model.MigrationKeyAddIPFilteringPermissions, Migration: a.getAddIPFilterPermissionsMigration}, {Key: model.MigrationKeyAddOutgoingOAuthConnectionsPermissions, Migration: a.getAddOutgoingOAuthConnectionsPermissions}, {Key: model.MigrationKeyAddChannelBookmarksPermissions, Migration: a.getAddChannelBookmarksPermissionsMigration}, + {Key: model.MigrationKeyAddManageJobAncillaryPermissions, Migration: a.getAddManageJobAncillaryPermissionsMigration}, } roles, err := s.Store().Role().GetAll() diff --git a/server/channels/store/opentracinglayer/opentracinglayer.go b/server/channels/store/opentracinglayer/opentracinglayer.go index 100741dc64..f145f49fea 100644 --- a/server/channels/store/opentracinglayer/opentracinglayer.go +++ b/server/channels/store/opentracinglayer/opentracinglayer.go @@ -5212,6 +5212,24 @@ func (s *OpenTracingLayerJobStore) GetAllByTypeAndStatus(c request.CTX, jobType return result, err } +func (s *OpenTracingLayerJobStore) GetAllByTypeAndStatusPage(c request.CTX, jobType []string, status string, offset int, limit int) ([]*model.Job, error) { + origCtx := s.Root.Store.Context() + span, newCtx := tracing.StartSpanWithParentByContext(s.Root.Store.Context(), "JobStore.GetAllByTypeAndStatusPage") + s.Root.Store.SetContext(newCtx) + defer func() { + s.Root.Store.SetContext(origCtx) + }() + + defer span.Finish() + result, err := s.JobStore.GetAllByTypeAndStatusPage(c, jobType, status, offset, limit) + if err != nil { + span.LogFields(spanlog.Error(err)) + ext.Error.Set(span, true) + } + + return result, err +} + func (s *OpenTracingLayerJobStore) GetAllByTypePage(c request.CTX, jobType string, offset int, limit int) ([]*model.Job, error) { origCtx := s.Root.Store.Context() span, newCtx := tracing.StartSpanWithParentByContext(s.Root.Store.Context(), "JobStore.GetAllByTypePage") diff --git a/server/channels/store/retrylayer/retrylayer.go b/server/channels/store/retrylayer/retrylayer.go index a8f63cee00..8d28f84aa2 100644 --- a/server/channels/store/retrylayer/retrylayer.go +++ b/server/channels/store/retrylayer/retrylayer.go @@ -5903,6 +5903,27 @@ func (s *RetryLayerJobStore) GetAllByTypeAndStatus(c request.CTX, jobType string } +func (s *RetryLayerJobStore) GetAllByTypeAndStatusPage(c request.CTX, jobType []string, status string, offset int, limit int) ([]*model.Job, error) { + + tries := 0 + for { + result, err := s.JobStore.GetAllByTypeAndStatusPage(c, jobType, status, offset, limit) + if err == nil { + return result, nil + } + if !isRepeatableError(err) { + return result, err + } + tries++ + if tries >= 3 { + err = errors.Wrap(err, "giving up after 3 consecutive repeatable transaction failures") + return result, err + } + timepkg.Sleep(100 * timepkg.Millisecond) + } + +} + func (s *RetryLayerJobStore) GetAllByTypePage(c request.CTX, jobType string, offset int, limit int) ([]*model.Job, error) { tries := 0 diff --git a/server/channels/store/sqlstore/job_store.go b/server/channels/store/sqlstore/job_store.go index d087e44091..7e0f3a8017 100644 --- a/server/channels/store/sqlstore/job_store.go +++ b/server/channels/store/sqlstore/job_store.go @@ -316,6 +316,26 @@ func (jss SqlJobStore) GetAllByStatus(c request.CTX, status string) ([]*model.Jo return statuses, nil } +func (jss SqlJobStore) GetAllByTypeAndStatusPage(c request.CTX, jobType []string, status string, offset int, limit int) ([]*model.Job, error) { + query, args, err := jss.getQueryBuilder(). + Select("*"). + From("Jobs"). + Where(sq.Eq{"Type": jobType, "Status": status}). + OrderBy("CreateAt DESC"). + Limit(uint64(limit)). + Offset(uint64(offset)).ToSql() + if err != nil { + return nil, errors.Wrap(err, "job_tosql") + } + + jobs := []*model.Job{} + if err = jss.GetReplicaX().Select(&jobs, query, args...); err != nil { + return nil, errors.Wrapf(err, "failed to find Jobs with type=%s and status=%s", strings.Join(jobType, ","), status) + } + + return jobs, nil +} + func (jss SqlJobStore) GetNewestJobByStatusAndType(status string, jobType string) (*model.Job, error) { return jss.GetNewestJobByStatusesAndType([]string{status}, jobType) } diff --git a/server/channels/store/store.go b/server/channels/store/store.go index 9703180794..bbb9862a1a 100644 --- a/server/channels/store/store.go +++ b/server/channels/store/store.go @@ -761,6 +761,7 @@ type JobStore interface { GetAllByTypePage(c request.CTX, jobType string, offset int, limit int) ([]*model.Job, error) GetAllByTypesPage(c request.CTX, jobTypes []string, offset int, limit int) ([]*model.Job, error) GetAllByStatus(c request.CTX, status string) ([]*model.Job, error) + GetAllByTypeAndStatusPage(c request.CTX, jobType []string, status string, offset int, limit int) ([]*model.Job, error) GetNewestJobByStatusAndType(status string, jobType string) (*model.Job, error) GetNewestJobByStatusesAndType(statuses []string, jobType string) (*model.Job, error) GetCountByStatusAndType(status string, jobType string) (int64, error) diff --git a/server/channels/store/storetest/job_store.go b/server/channels/store/storetest/job_store.go index 6f9988bb61..54fc07477a 100644 --- a/server/channels/store/storetest/job_store.go +++ b/server/channels/store/storetest/job_store.go @@ -24,6 +24,7 @@ func TestJobStore(t *testing.T, rctx request.CTX, ss store.Store) { t.Run("JobGetAllByTypeAndStatus", func(t *testing.T) { testJobGetAllByTypeAndStatus(t, rctx, ss) }) t.Run("JobGetAllByTypePage", func(t *testing.T) { testJobGetAllByTypePage(t, rctx, ss) }) t.Run("JobGetAllByTypesPage", func(t *testing.T) { testJobGetAllByTypesPage(t, rctx, ss) }) + t.Run("JobGetAllByTypeAndStatusPage", func(t *testing.T) { testJobGetAllByTypeAndStatusPage(t, rctx, ss) }) t.Run("JobGetAllByStatus", func(t *testing.T) { testJobGetAllByStatus(t, rctx, ss) }) t.Run("GetNewestJobByStatusAndType", func(t *testing.T) { testJobStoreGetNewestJobByStatusAndType(t, rctx, ss) }) t.Run("GetNewestJobByStatusesAndType", func(t *testing.T) { testJobStoreGetNewestJobByStatusesAndType(t, rctx, ss) }) @@ -259,6 +260,62 @@ func testJobGetAllByTypesPage(t *testing.T, rctx request.CTX, ss store.Store) { require.Equal(t, received[0].Id, jobs[1].Id, "should've received oldest job last") } +func testJobGetAllByTypeAndStatusPage(t *testing.T, rctx request.CTX, ss store.Store) { + jobType := model.NewId() + jobType2 := model.NewId() + t0 := model.GetMillis() + + jobs := []*model.Job{ + { + Id: model.NewId(), + Type: jobType, + Status: model.JobStatusPending, + CreateAt: t0, + }, + { + Id: model.NewId(), + Type: jobType, + Status: model.JobStatusPending, + CreateAt: t0 + 1, + }, + { + Id: model.NewId(), + Type: jobType2, + Status: model.JobStatusCanceled, + CreateAt: t0 + 2, + }, + { + Id: model.NewId(), + Type: jobType2, + Status: model.JobStatusCanceled, + CreateAt: t0 + 3, + }, + } + + for _, job := range jobs { + _, err := ss.Job().Save(job) + require.NoError(t, err) + defer ss.Job().Delete(job.Id) + } + + jobTypes := []string{jobType, jobType2} + received, err := ss.Job().GetAllByTypeAndStatusPage(rctx, jobTypes, model.JobStatusPending, 0, 4) + require.NoError(t, err) + require.Len(t, received, 2) + require.Equal(t, received[0].Id, jobs[1].Id, "should've received newest job first") + require.Equal(t, received[1].Id, jobs[0].Id, "should've received oldest job last") + + received, err = ss.Job().GetAllByTypeAndStatusPage(rctx, jobTypes, model.JobStatusPending, 1, 1) + require.NoError(t, err) + require.Len(t, received, 1) + require.Equal(t, received[0].Id, jobs[0].Id, "should've received the oldest pending job") + + received, err = ss.Job().GetAllByTypeAndStatusPage(rctx, []string{jobType2}, model.JobStatusCanceled, 1, 1) + require.NoError(t, err) + require.Len(t, received, 1) + require.Equal(t, received[0].Id, jobs[2].Id, "should've received the oldest canceled job") +} + func testJobGetAllByStatus(t *testing.T, rctx request.CTX, ss store.Store) { jobType := model.NewId() status := model.NewId() diff --git a/server/channels/store/storetest/mocks/JobStore.go b/server/channels/store/storetest/mocks/JobStore.go index 3f05f44b4e..3c7411e326 100644 --- a/server/channels/store/storetest/mocks/JobStore.go +++ b/server/channels/store/storetest/mocks/JobStore.go @@ -181,6 +181,36 @@ func (_m *JobStore) GetAllByTypeAndStatus(c request.CTX, jobType string, status return r0, r1 } +// GetAllByTypeAndStatusPage provides a mock function with given fields: c, jobType, status, offset, limit +func (_m *JobStore) GetAllByTypeAndStatusPage(c request.CTX, jobType []string, status string, offset int, limit int) ([]*model.Job, error) { + ret := _m.Called(c, jobType, status, offset, limit) + + if len(ret) == 0 { + panic("no return value specified for GetAllByTypeAndStatusPage") + } + + var r0 []*model.Job + var r1 error + if rf, ok := ret.Get(0).(func(request.CTX, []string, string, int, int) ([]*model.Job, error)); ok { + return rf(c, jobType, status, offset, limit) + } + if rf, ok := ret.Get(0).(func(request.CTX, []string, string, int, int) []*model.Job); ok { + r0 = rf(c, jobType, status, offset, limit) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).([]*model.Job) + } + } + + if rf, ok := ret.Get(1).(func(request.CTX, []string, string, int, int) error); ok { + r1 = rf(c, jobType, status, offset, limit) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + // GetAllByTypePage provides a mock function with given fields: c, jobType, offset, limit func (_m *JobStore) GetAllByTypePage(c request.CTX, jobType string, offset int, limit int) ([]*model.Job, error) { ret := _m.Called(c, jobType, offset, limit) diff --git a/server/channels/store/timerlayer/timerlayer.go b/server/channels/store/timerlayer/timerlayer.go index db9128b8bf..50879eec60 100644 --- a/server/channels/store/timerlayer/timerlayer.go +++ b/server/channels/store/timerlayer/timerlayer.go @@ -4731,6 +4731,22 @@ func (s *TimerLayerJobStore) GetAllByTypeAndStatus(c request.CTX, jobType string return result, err } +func (s *TimerLayerJobStore) GetAllByTypeAndStatusPage(c request.CTX, jobType []string, status string, offset int, limit int) ([]*model.Job, error) { + start := time.Now() + + result, err := s.JobStore.GetAllByTypeAndStatusPage(c, jobType, status, offset, limit) + + elapsed := float64(time.Since(start)) / float64(time.Second) + if s.Root.Metrics != nil { + success := "false" + if err == nil { + success = "true" + } + s.Root.Metrics.ObserveStoreMethodDuration("JobStore.GetAllByTypeAndStatusPage", success, elapsed) + } + return result, err +} + func (s *TimerLayerJobStore) GetAllByTypePage(c request.CTX, jobType string, offset int, limit int) ([]*model.Job, error) { start := time.Now() diff --git a/server/channels/testlib/store.go b/server/channels/testlib/store.go index 466d0aa302..630c9434e9 100644 --- a/server/channels/testlib/store.go +++ b/server/channels/testlib/store.go @@ -76,6 +76,7 @@ func GetMockStoreForSetupFunctions() *mocks.Store { systemStore.On("GetByName", model.MigrationKeyAddIPFilteringPermissions).Return(&model.System{Name: model.MigrationKeyAddIPFilteringPermissions, Value: "true"}, nil) systemStore.On("GetByName", model.MigrationKeyAddOutgoingOAuthConnectionsPermissions).Return(&model.System{Name: model.MigrationKeyAddOutgoingOAuthConnectionsPermissions, Value: "true"}, nil) systemStore.On("GetByName", model.MigrationKeyAddChannelBookmarksPermissions).Return(&model.System{Name: model.MigrationKeyAddChannelBookmarksPermissions, Value: "true"}, nil) + systemStore.On("GetByName", model.MigrationKeyAddManageJobAncillaryPermissions).Return(&model.System{Name: model.MigrationKeyAddManageJobAncillaryPermissions, Value: "true"}, nil) systemStore.On("GetByName", "CustomGroupAdminRoleCreationMigrationComplete").Return(&model.System{Name: model.MigrationKeyAddPlayboosksManageRolesPermissions, Value: "true"}, nil) systemStore.On("GetByName", "products_boards").Return(&model.System{Name: "products_boards", Value: "true"}, nil) systemStore.On("GetByName", "elasticsearch_fix_channel_index_migration").Return(&model.System{Name: "elasticsearch_fix_channel_index_migration", Value: "true"}, nil) diff --git a/server/cmd/mmctl/client/client.go b/server/cmd/mmctl/client/client.go index a045277650..fb7372d492 100644 --- a/server/cmd/mmctl/client/client.go +++ b/server/cmd/mmctl/client/client.go @@ -128,10 +128,11 @@ type Client interface { UploadData(ctx context.Context, uploadID string, data io.Reader) (*model.FileInfo, *model.Response, error) ListImports(ctx context.Context) ([]string, *model.Response, error) GetJob(ctx context.Context, id string) (*model.Job, *model.Response, error) - GetJobs(ctx context.Context, page int, perPage int) ([]*model.Job, *model.Response, error) + GetJobs(ctx context.Context, jobType string, status string, page int, perPage int) ([]*model.Job, *model.Response, error) GetJobsByType(ctx context.Context, jobType string, page int, perPage int) ([]*model.Job, *model.Response, error) CreateJob(ctx context.Context, job *model.Job) (*model.Job, *model.Response, error) CancelJob(ctx context.Context, jobID string) (*model.Response, error) + UpdateJobStatus(ctx context.Context, jobId string, status string, force bool) (*model.Response, error) CreateIncomingWebhook(ctx context.Context, hook *model.IncomingWebhook) (*model.IncomingWebhook, *model.Response, error) UpdateIncomingWebhook(ctx context.Context, hook *model.IncomingWebhook) (*model.IncomingWebhook, *model.Response, error) GetIncomingWebhooks(ctx context.Context, page int, perPage int, etag string) ([]*model.IncomingWebhook, *model.Response, error) diff --git a/server/cmd/mmctl/commands/export.go b/server/cmd/mmctl/commands/export.go index 2634bcd0ee..2d30f26006 100644 --- a/server/cmd/mmctl/commands/export.go +++ b/server/cmd/mmctl/commands/export.go @@ -270,7 +270,7 @@ func exportDownloadCmdF(c client.Client, command *cobra.Command, args []string) } func exportJobListCmdF(c client.Client, command *cobra.Command, args []string) error { - return jobListCmdF(c, command, model.JobTypeExportProcess) + return jobListCmdF(c, command, model.JobTypeExportProcess, "") } func exportJobShowCmdF(c client.Client, command *cobra.Command, args []string) error { diff --git a/server/cmd/mmctl/commands/extract.go b/server/cmd/mmctl/commands/extract.go index f39430896b..d83fefaaf2 100644 --- a/server/cmd/mmctl/commands/extract.go +++ b/server/cmd/mmctl/commands/extract.go @@ -107,7 +107,7 @@ func extractJobShowCmdF(c client.Client, command *cobra.Command, args []string) } func extractJobListCmdF(c client.Client, command *cobra.Command, args []string) error { - return jobListCmdF(c, command, model.JobTypeExtractContent) + return jobListCmdF(c, command, model.JobTypeExtractContent, "") } func printExtractContentJob(job *model.Job) { diff --git a/server/cmd/mmctl/commands/import.go b/server/cmd/mmctl/commands/import.go index 23907a6d9b..94455619c5 100644 --- a/server/cmd/mmctl/commands/import.go +++ b/server/cmd/mmctl/commands/import.go @@ -308,24 +308,6 @@ func importProcessCmdF(c client.Client, command *cobra.Command, args []string) e return nil } -func printJob(job *model.Job) { - if job.StartAt > 0 { - printer.PrintT(fmt.Sprintf(` ID: {{.Id}} - Status: {{.Status}} - Created: %s - Started: %s - Data: {{.Data}} -`, - time.Unix(job.CreateAt/1000, 0), time.Unix(job.StartAt/1000, 0)), job) - } else { - printer.PrintT(fmt.Sprintf(` ID: {{.Id}} - Status: {{.Status}} - Created: %s -`, - time.Unix(job.CreateAt/1000, 0)), job) - } -} - func importJobShowCmdF(c client.Client, command *cobra.Command, args []string) error { job, _, err := c.GetJob(context.TODO(), args[0]) if err != nil { @@ -337,53 +319,8 @@ func importJobShowCmdF(c client.Client, command *cobra.Command, args []string) e return nil } -func jobListCmdF(c client.Client, command *cobra.Command, jobType string) error { - page, err := command.Flags().GetInt("page") - if err != nil { - return err - } - perPage, err := command.Flags().GetInt("per-page") - if err != nil { - return err - } - showAll, err := command.Flags().GetBool("all") - if err != nil { - return err - } - - if showAll { - page = 0 - } - - for { - jobs, _, err := c.GetJobsByType(context.TODO(), jobType, page, perPage) - if err != nil { - return fmt.Errorf("failed to get jobs: %w", err) - } - - if len(jobs) == 0 { - if !showAll || page == 0 { - printer.Print("No jobs found") - } - return nil - } - - for _, job := range jobs { - printJob(job) - } - - if !showAll { - break - } - - page++ - } - - return nil -} - func importJobListCmdF(c client.Client, command *cobra.Command, args []string) error { - return jobListCmdF(c, command, model.JobTypeImportProcess) + return jobListCmdF(c, command, model.JobTypeImportProcess, "") } type Statistics struct { diff --git a/server/cmd/mmctl/commands/import_test.go b/server/cmd/mmctl/commands/import_test.go index 4f36b6a3ff..1c7907e81d 100644 --- a/server/cmd/mmctl/commands/import_test.go +++ b/server/cmd/mmctl/commands/import_test.go @@ -163,7 +163,7 @@ func (s *MmctlUnitTestSuite) TestImportJobListCmdF() { s.client. EXPECT(). - GetJobsByType(context.TODO(), model.JobTypeImportProcess, 0, perPage). + GetJobs(context.TODO(), model.JobTypeImportProcess, "", 0, perPage). Return(mockJobs, &model.Response{}, nil). Times(1) @@ -196,7 +196,7 @@ func (s *MmctlUnitTestSuite) TestImportJobListCmdF() { s.client. EXPECT(). - GetJobsByType(context.TODO(), model.JobTypeImportProcess, 0, perPage). + GetJobs(context.TODO(), model.JobTypeImportProcess, "", 0, perPage). Return(mockJobs, &model.Response{}, nil). Times(1) diff --git a/server/cmd/mmctl/commands/job.go b/server/cmd/mmctl/commands/job.go new file mode 100644 index 0000000000..0a861c854f --- /dev/null +++ b/server/cmd/mmctl/commands/job.go @@ -0,0 +1,202 @@ +// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved. +// See LICENSE.txt for license information. + +package commands + +import ( + "context" + "fmt" + "time" + + "github.com/hashicorp/go-multierror" + "github.com/mattermost/mattermost/server/public/model" + "github.com/mattermost/mattermost/server/v8/cmd/mmctl/client" + "github.com/mattermost/mattermost/server/v8/cmd/mmctl/printer" + + "github.com/spf13/cobra" +) + +var JobCmd = &cobra.Command{ + Use: "job", + Short: "Management of jobs", +} + +var listJobsCmd = &cobra.Command{ + Use: "list", + Short: "List the latest jobs", + Example: ` job list + job list --ids jobID1,jobID2 + job list --type ldap_sync --status success + job list --type ldap_sync --status success --page 0 --per-page 10`, + Args: cobra.NoArgs, + RunE: withClient(listJobsCmdF), +} + +var updateJobCmd = &cobra.Command{ + Use: "update [job] [status]", + Short: "Update the status of a job", + Long: `Update the status of a job. The following restrictions are in place: + - in_progress -> pending + - in_progress | pending -> cancel_requested + - cancel_requested -> canceled + + Those restriction can be bypassed with --force=true but the only statuses you can go to are: pending, cancel_requested and canceled. This can have unexpected consequences and should be used with caution.`, + Example: ` job update myJobID pending + job update myJobID pending --force true + job update myJobID canceled --force true`, + Args: cobra.MinimumNArgs(2), + RunE: withClient(updateJobCmdF), +} + +func init() { + listJobsCmd.Flags().Int("page", 0, "Page number to fetch for the list of import jobs") + listJobsCmd.Flags().Int("per-page", 5, "Number of import jobs to be fetched") + listJobsCmd.Flags().Bool("all", false, "Fetch all import jobs. --page flag will be ignored if provided") + listJobsCmd.Flags().StringSlice("ids", nil, "Comma-separated list of job IDs to which the operation will be applied. All other flags are ignored") + listJobsCmd.Flags().String("status", "", "Filter by job status") + listJobsCmd.Flags().String("type", "", "Filter by job type") + + updateJobCmd.Flags().Bool("force", false, "Setting a job status is restricted to certain statuses. You can overwrite these restrictions by using --force. This might cause unexpected behaviour on your Mattermost Server. Use this option with caution.") + + JobCmd.AddCommand( + listJobsCmd, + updateJobCmd, + ) + + RootCmd.AddCommand(JobCmd) +} + +func listJobsCmdF(c client.Client, cmd *cobra.Command, args []string) error { + ids, err := cmd.Flags().GetStringSlice("ids") + if err != nil { + return err + } + jobType, err := cmd.Flags().GetString("type") + if err != nil { + return err + } + status, err := cmd.Flags().GetString("status") + if err != nil { + return err + } + + if len(ids) > 0 { + jobs := make([]*model.Job, 0, len(ids)) + var result *multierror.Error + for _, id := range ids { + isValidId := model.IsValidId(id) + if !isValidId { + result = multierror.Append(result, fmt.Errorf("invalid job ID: %s", id)) + continue + } + + job, _, err := c.GetJob(context.TODO(), id) + if err != nil { + result = multierror.Append(result, err) + continue + } + jobs = append(jobs, job) + } + for _, job := range jobs { + printJob(job) + } + return result.ErrorOrNil() + } + + return jobListCmdF(c, cmd, jobType, status) +} + +func updateJobCmdF(c client.Client, cmd *cobra.Command, args []string) error { + force, err := cmd.Flags().GetBool("force") + if err != nil { + return err + } + + jobId := args[0] + if !model.IsValidId(jobId) { + return fmt.Errorf("invalid job ID: %s", jobId) + } + status := args[1] + if !model.IsValidJobStatus(status) { + return fmt.Errorf("invalid job status: %s", status) + } + + _, err = c.UpdateJobStatus(context.TODO(), jobId, status, force) + if err != nil { + return err + } + + return nil +} + +func jobListCmdF(c client.Client, command *cobra.Command, jobType string, status string) error { + page, err := command.Flags().GetInt("page") + if err != nil { + return err + } + perPage, err := command.Flags().GetInt("per-page") + if err != nil { + return err + } + showAll, err := command.Flags().GetBool("all") + if err != nil { + return err + } + + if showAll { + page = 0 + } + + if jobType != "" && !model.IsValidJobType(jobType) { + return fmt.Errorf("invalid job type: %s", jobType) + } + + if status != "" && !model.IsValidJobStatus(status) { + return fmt.Errorf("invalid job status: %s", status) + } + + for { + jobs, _, err := c.GetJobs(context.TODO(), jobType, status, page, perPage) + if err != nil { + return fmt.Errorf("failed to get jobs: %w", err) + } + + if len(jobs) == 0 { + if !showAll || page == 0 { + printer.Print("No jobs found") + } + return nil + } + + for _, job := range jobs { + printJob(job) + } + + if !showAll { + break + } + + page++ + } + + return nil +} + +func printJob(job *model.Job) { + if job.StartAt > 0 { + printer.PrintT(fmt.Sprintf(` ID: {{.Id}} + Type: {{.Type}} + Status: {{.Status}} + Created: %s + Started: %s + Data: {{.Data}} +`, + time.Unix(job.CreateAt/1000, 0), time.Unix(job.StartAt/1000, 0)), job) + } else { + printer.PrintT(fmt.Sprintf(` ID: {{.Id}} + Status: {{.Status}} + Created: %s +`, + time.Unix(job.CreateAt/1000, 0)), job) + } +} diff --git a/server/cmd/mmctl/commands/job_test.go b/server/cmd/mmctl/commands/job_test.go new file mode 100644 index 0000000000..b36563fa30 --- /dev/null +++ b/server/cmd/mmctl/commands/job_test.go @@ -0,0 +1,204 @@ +// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved. +// See LICENSE.txt for license information. + +package commands + +import ( + "context" + + "github.com/mattermost/mattermost/server/public/model" + + "github.com/mattermost/mattermost/server/v8/cmd/mmctl/printer" + + "github.com/spf13/cobra" +) + +func (s *MmctlUnitTestSuite) TestListJobsCmdF() { + s.Run("no jobs found", func() { + printer.Clean() + var mockJobs []*model.Job + + cmd := &cobra.Command{} + perPage := 10 + cmd.Flags().Int("page", 0, "") + cmd.Flags().Int("per-page", perPage, "") + cmd.Flags().Bool("all", false, "") + cmd.Flags().StringSlice("ids", []string{}, "") + cmd.Flags().String("status", "", "") + cmd.Flags().String("type", "", "") + + s.client. + EXPECT(). + GetJobs(context.TODO(), "", "", 0, perPage). + Return(mockJobs, &model.Response{}, nil). + Times(1) + + err := listJobsCmdF(s.client, cmd, nil) + s.Require().Nil(err) + s.Len(printer.GetLines(), 1) + s.Empty(printer.GetErrorLines()) + s.Equal("No jobs found", printer.GetLines()[0]) + }) + + s.Run("3 jobs found", func() { + printer.Clean() + mockJobs := []*model.Job{ + { + Id: model.NewId(), + }, + { + Id: model.NewId(), + }, + { + Id: model.NewId(), + }, + } + + cmd := &cobra.Command{} + perPage := 3 + cmd.Flags().Int("page", 0, "") + cmd.Flags().Int("per-page", perPage, "") + cmd.Flags().Bool("all", false, "") + cmd.Flags().StringSlice("ids", []string{}, "") + cmd.Flags().String("status", "", "") + cmd.Flags().String("type", "", "") + + s.client. + EXPECT(). + GetJobs(context.TODO(), "", "", 0, perPage). + Return(mockJobs, &model.Response{}, nil). + Times(1) + + err := listJobsCmdF(s.client, cmd, nil) + s.Require().Nil(err) + s.Len(printer.GetLines(), len(mockJobs)) + s.Empty(printer.GetErrorLines()) + for i, line := range printer.GetLines() { + s.Equal(mockJobs[i], line.(*model.Job)) + } + }) + + s.Run("return 1 job using ids flag", func() { + printer.Clean() + id := model.NewId() + mockJob := &model.Job{ + Id: id, + } + + cmd := &cobra.Command{} + perPage := 3 + cmd.Flags().Int("page", 0, "") + cmd.Flags().Int("per-page", perPage, "") + cmd.Flags().Bool("all", false, "") + cmd.Flags().StringSlice("ids", []string{id}, "") + cmd.Flags().String("status", "", "") + cmd.Flags().String("type", "", "") + + s.client. + EXPECT(). + GetJob(context.TODO(), id). + Return(mockJob, &model.Response{}, nil). + Times(1) + + err := listJobsCmdF(s.client, cmd, nil) + s.Require().Nil(err) + s.Len(printer.GetLines(), 1) + s.Empty(printer.GetErrorLines()) + for _, line := range printer.GetLines() { + s.Equal(mockJob, line.(*model.Job)) + } + }) + + s.Run("return 2 jobs by status", func() { + printer.Clean() + mockJobs := []*model.Job{ + { + Id: model.NewId(), + Status: model.JobStatusSuccess, + }, + { + Id: model.NewId(), + Status: model.JobStatusSuccess, + }, + } + + cmd := &cobra.Command{} + perPage := 2 + cmd.Flags().Int("page", 0, "") + cmd.Flags().Int("per-page", perPage, "") + cmd.Flags().Bool("all", false, "") + cmd.Flags().String("status", model.JobStatusSuccess, "") + cmd.Flags().StringSlice("ids", []string{}, "") + cmd.Flags().String("type", "", "") + + s.client. + EXPECT(). + GetJobs(context.TODO(), "", model.JobStatusSuccess, 0, perPage). + Return(mockJobs, &model.Response{}, nil). + Times(1) + + err := listJobsCmdF(s.client, cmd, nil) + s.Require().Nil(err) + s.Len(printer.GetLines(), len(mockJobs)) + s.Empty(printer.GetErrorLines()) + for i, line := range printer.GetLines() { + s.Equal(mockJobs[i], line.(*model.Job)) + } + }) + + s.Run("return 2 jobs by type", func() { + printer.Clean() + mockJobs := []*model.Job{ + { + Id: model.NewId(), + Type: model.JobTypeDataRetention, + }, + { + Id: model.NewId(), + Type: model.JobTypeDataRetention, + }, + } + + cmd := &cobra.Command{} + perPage := 2 + cmd.Flags().Int("page", 0, "") + cmd.Flags().Int("per-page", perPage, "") + cmd.Flags().Bool("all", false, "") + cmd.Flags().String("type", model.JobTypeDataRetention, "") + cmd.Flags().StringSlice("ids", []string{}, "") + cmd.Flags().String("status", "", "") + + s.client. + EXPECT(). + GetJobs(context.TODO(), model.JobTypeDataRetention, "", 0, perPage). + Return(mockJobs, &model.Response{}, nil). + Times(1) + + err := listJobsCmdF(s.client, cmd, nil) + s.Require().Nil(err) + s.Len(printer.GetLines(), len(mockJobs)) + s.Empty(printer.GetErrorLines()) + for i, line := range printer.GetLines() { + s.Equal(mockJobs[i], line.(*model.Job)) + } + }) +} + +func (s *MmctlUnitTestSuite) TestUpdateJobCmdF() { + s.Run("update job status", func() { + printer.Clean() + id := model.NewId() + + cmd := &cobra.Command{} + cmd.Flags().Bool("force", true, "") + + s.client. + EXPECT(). + UpdateJobStatus(context.TODO(), id, model.JobStatusPending, true). + Return(&model.Response{}, nil). + Times(1) + + err := updateJobCmdF(s.client, cmd, []string{id, model.JobStatusPending}) + s.Require().Nil(err) + }) +} diff --git a/server/cmd/mmctl/commands/ldap.go b/server/cmd/mmctl/commands/ldap.go index bd504f6ddb..c260668314 100644 --- a/server/cmd/mmctl/commands/ldap.go +++ b/server/cmd/mmctl/commands/ldap.go @@ -121,7 +121,7 @@ func ldapIDMigrateCmdF(c client.Client, cmd *cobra.Command, args []string) error } func ldapJobListCmdF(c client.Client, command *cobra.Command, args []string) error { - return jobListCmdF(c, command, model.JobTypeLdapSync) + return jobListCmdF(c, command, model.JobTypeLdapSync, "") } func ldapJobShowCmdF(c client.Client, command *cobra.Command, args []string) error { diff --git a/server/cmd/mmctl/commands/ldap_test.go b/server/cmd/mmctl/commands/ldap_test.go index a853dd1149..cda0cbd6e2 100644 --- a/server/cmd/mmctl/commands/ldap_test.go +++ b/server/cmd/mmctl/commands/ldap_test.go @@ -128,7 +128,7 @@ func (s *MmctlUnitTestSuite) TestLdapJobListCmdF() { s.client. EXPECT(). - GetJobsByType(context.TODO(), model.JobTypeLdapSync, 0, perPage). + GetJobs(context.TODO(), model.JobTypeLdapSync, "", 0, perPage). Return(mockJobs, &model.Response{}, nil). Times(1) @@ -161,7 +161,7 @@ func (s *MmctlUnitTestSuite) TestLdapJobListCmdF() { s.client. EXPECT(). - GetJobsByType(context.TODO(), model.JobTypeLdapSync, 0, perPage). + GetJobs(context.TODO(), model.JobTypeLdapSync, "", 0, perPage). Return(mockJobs, &model.Response{}, nil). Times(1) diff --git a/server/cmd/mmctl/docs/mmctl.rst b/server/cmd/mmctl/docs/mmctl.rst index c5532b3359..c2915fd5fa 100644 --- a/server/cmd/mmctl/docs/mmctl.rst +++ b/server/cmd/mmctl/docs/mmctl.rst @@ -42,6 +42,7 @@ SEE ALSO * `mmctl group `_ - Management of groups * `mmctl import `_ - Management of imports * `mmctl integrity `_ - Check database records integrity. +* `mmctl job `_ - Management of jobs * `mmctl ldap `_ - LDAP related utilities * `mmctl license `_ - Licensing commands * `mmctl logs `_ - Display logs in a human-readable format diff --git a/server/cmd/mmctl/docs/mmctl_job.rst b/server/cmd/mmctl/docs/mmctl_job.rst new file mode 100644 index 0000000000..59d84dd0bb --- /dev/null +++ b/server/cmd/mmctl/docs/mmctl_job.rst @@ -0,0 +1,42 @@ +.. _mmctl_job: + +mmctl job +--------- + +Management of jobs + +Synopsis +~~~~~~~~ + + +Management of jobs + +Options +~~~~~~~ + +:: + + -h, --help help for job + +Options inherited from parent commands +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +:: + + --config string path to the configuration file (default "$XDG_CONFIG_HOME/mmctl/config") + --disable-pager disables paged output + --insecure-sha1-intermediate allows to use insecure TLS protocols, such as SHA-1 + --insecure-tls-version allows to use TLS versions 1.0 and 1.1 + --json the output format will be in json format + --local allows communicating with the server through a unix socket + --quiet prevent mmctl to generate output for the commands + --strict will only run commands if the mmctl version matches the server one + --suppress-warnings disables printing warning messages + +SEE ALSO +~~~~~~~~ + +* `mmctl `_ - Remote client for the Open Source, self-hosted Slack-alternative +* `mmctl job list `_ - List the latest jobs +* `mmctl job update `_ - Update the status of a job + diff --git a/server/cmd/mmctl/docs/mmctl_job_list.rst b/server/cmd/mmctl/docs/mmctl_job_list.rst new file mode 100644 index 0000000000..acb6aa3cf7 --- /dev/null +++ b/server/cmd/mmctl/docs/mmctl_job_list.rst @@ -0,0 +1,60 @@ +.. _mmctl_job_list: + +mmctl job list +-------------- + +List the latest jobs + +Synopsis +~~~~~~~~ + + +List the latest jobs + +:: + + mmctl job list [flags] + +Examples +~~~~~~~~ + +:: + + job list + job list --ids jobID1,jobID2 + job list --type ldap_sync --status success + job list --type ldap_sync --status success --page 0 --per-page 10 + +Options +~~~~~~~ + +:: + + --all Fetch all import jobs. --page flag will be ignored if provided + -h, --help help for list + --ids strings Comma-separated list of job IDs to which the operation will be applied. All other flags are ignored + --page int Page number to fetch for the list of import jobs + --per-page int Number of import jobs to be fetched (default 5) + --status string Filter by job status + --type string Filter by job type + +Options inherited from parent commands +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +:: + + --config string path to the configuration file (default "$XDG_CONFIG_HOME/mmctl/config") + --disable-pager disables paged output + --insecure-sha1-intermediate allows to use insecure TLS protocols, such as SHA-1 + --insecure-tls-version allows to use TLS versions 1.0 and 1.1 + --json the output format will be in json format + --local allows communicating with the server through a unix socket + --quiet prevent mmctl to generate output for the commands + --strict will only run commands if the mmctl version matches the server one + --suppress-warnings disables printing warning messages + +SEE ALSO +~~~~~~~~ + +* `mmctl job `_ - Management of jobs + diff --git a/server/cmd/mmctl/docs/mmctl_job_update.rst b/server/cmd/mmctl/docs/mmctl_job_update.rst new file mode 100644 index 0000000000..b9b4aa280e --- /dev/null +++ b/server/cmd/mmctl/docs/mmctl_job_update.rst @@ -0,0 +1,59 @@ +.. _mmctl_job_update: + +mmctl job update +---------------- + +Update the status of a job + +Synopsis +~~~~~~~~ + + +Update the status of a job. The following restrictions are in place: + - in_progress -> pending + - in_progress | pending -> cancel_requested + - cancel_requested -> canceled + + Those restriction can be bypassed with --force=true but the only statuses you can go to are: pending, cancel_requested and canceled. This can have unexpected consequences and should be used with caution. + +:: + + mmctl job update [job] [status] [flags] + +Examples +~~~~~~~~ + +:: + + job update myJobID pending + job update myJobID pending --force true + job update myJobID canceled --force true + +Options +~~~~~~~ + +:: + + --force Setting a job status is restricted to certain statuses. You can overwrite these restrictions by using --force. This might cause unexpected behaviour on your Mattermost Server. Use this option with caution. + -h, --help help for update + +Options inherited from parent commands +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +:: + + --config string path to the configuration file (default "$XDG_CONFIG_HOME/mmctl/config") + --disable-pager disables paged output + --insecure-sha1-intermediate allows to use insecure TLS protocols, such as SHA-1 + --insecure-tls-version allows to use TLS versions 1.0 and 1.1 + --json the output format will be in json format + --local allows communicating with the server through a unix socket + --quiet prevent mmctl to generate output for the commands + --strict will only run commands if the mmctl version matches the server one + --suppress-warnings disables printing warning messages + +SEE ALSO +~~~~~~~~ + +* `mmctl job `_ - Management of jobs + diff --git a/server/cmd/mmctl/mocks/client_mock.go b/server/cmd/mmctl/mocks/client_mock.go index 59b3cc2fb8..6fe8d254cf 100644 --- a/server/cmd/mmctl/mocks/client_mock.go +++ b/server/cmd/mmctl/mocks/client_mock.go @@ -863,9 +863,9 @@ func (mr *MockClientMockRecorder) GetJob(arg0, arg1 interface{}) *gomock.Call { } // GetJobs mocks base method. -func (m *MockClient) GetJobs(arg0 context.Context, arg1, arg2 int) ([]*model.Job, *model.Response, error) { +func (m *MockClient) GetJobs(arg0 context.Context, arg1, arg2 string, arg3, arg4 int) ([]*model.Job, *model.Response, error) { m.ctrl.T.Helper() - ret := m.ctrl.Call(m, "GetJobs", arg0, arg1, arg2) + ret := m.ctrl.Call(m, "GetJobs", arg0, arg1, arg2, arg3, arg4) ret0, _ := ret[0].([]*model.Job) ret1, _ := ret[1].(*model.Response) ret2, _ := ret[2].(error) @@ -873,9 +873,9 @@ func (m *MockClient) GetJobs(arg0 context.Context, arg1, arg2 int) ([]*model.Job } // GetJobs indicates an expected call of GetJobs. -func (mr *MockClientMockRecorder) GetJobs(arg0, arg1, arg2 interface{}) *gomock.Call { +func (mr *MockClientMockRecorder) GetJobs(arg0, arg1, arg2, arg3, arg4 interface{}) *gomock.Call { mr.mock.ctrl.T.Helper() - return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetJobs", reflect.TypeOf((*MockClient)(nil).GetJobs), arg0, arg1, arg2) + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetJobs", reflect.TypeOf((*MockClient)(nil).GetJobs), arg0, arg1, arg2, arg3, arg4) } // GetJobsByType mocks base method. @@ -2105,6 +2105,21 @@ func (mr *MockClientMockRecorder) UpdateIncomingWebhook(arg0, arg1 interface{}) return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "UpdateIncomingWebhook", reflect.TypeOf((*MockClient)(nil).UpdateIncomingWebhook), arg0, arg1) } +// UpdateJobStatus mocks base method. +func (m *MockClient) UpdateJobStatus(arg0 context.Context, arg1, arg2 string, arg3 bool) (*model.Response, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "UpdateJobStatus", arg0, arg1, arg2, arg3) + ret0, _ := ret[0].(*model.Response) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// UpdateJobStatus indicates an expected call of UpdateJobStatus. +func (mr *MockClientMockRecorder) UpdateJobStatus(arg0, arg1, arg2, arg3 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "UpdateJobStatus", reflect.TypeOf((*MockClient)(nil).UpdateJobStatus), arg0, arg1, arg2, arg3) +} + // UpdateOutgoingWebhook mocks base method. func (m *MockClient) UpdateOutgoingWebhook(arg0 context.Context, arg1 *model.OutgoingWebhook) (*model.OutgoingWebhook, *model.Response, error) { m.ctrl.T.Helper() diff --git a/server/i18n/en.json b/server/i18n/en.json index b3b62c7165..3b25194742 100644 --- a/server/i18n/en.json +++ b/server/i18n/en.json @@ -2176,6 +2176,10 @@ "id": "api.job.retrieve.nopermissions", "translation": "The job types of a job you are trying to retrieve does not contain permissions" }, + { + "id": "api.job.status.invalid", + "translation": "Invalid status set" + }, { "id": "api.job.unable_to_create_job.incorrect_job_type", "translation": "The job type of the job you are trying to create is invalid" @@ -2188,6 +2192,10 @@ "id": "api.job.unable_to_download_job.incorrect_job_type", "translation": "The job type you are trying to download is not supported at the moment" }, + { + "id": "api.job.unable_to_manage_job.incorrect_job_type", + "translation": "You do not have permission to manage this job type" + }, { "id": "api.ldap_group.not_found", "translation": "ldap group not found" @@ -5670,6 +5678,10 @@ "id": "app.job.update.app_error", "translation": "Unable to update the job." }, + { + "id": "app.job.update_status.app_error", + "translation": "Unable to update job status. Invalid status set" + }, { "id": "app.last_accessible_file.app_error", "translation": "Error fetching last accessible file" diff --git a/server/public/model/client4.go b/server/public/model/client4.go index fc229a74c7..45bfae3c3c 100644 --- a/server/public/model/client4.go +++ b/server/public/model/client4.go @@ -7018,8 +7018,8 @@ func (c *Client4) GetJob(ctx context.Context, id string) (*Job, *Response, error } // GetJobs gets all jobs, sorted with the job that was created most recently first. -func (c *Client4) GetJobs(ctx context.Context, page int, perPage int) ([]*Job, *Response, error) { - r, err := c.DoAPIGet(ctx, c.jobsRoute()+fmt.Sprintf("?page=%v&per_page=%v", page, perPage), "") +func (c *Client4) GetJobs(ctx context.Context, jobType string, status string, page int, perPage int) ([]*Job, *Response, error) { + r, err := c.DoAPIGet(ctx, c.jobsRoute()+fmt.Sprintf("?page=%v&per_page=%v&job_type=%v&status=%v", page, perPage, jobType, status), "") if err != nil { return nil, BuildResponse(r), err } @@ -7088,6 +7088,23 @@ func (c *Client4) DownloadJob(ctx context.Context, jobId string) ([]byte, *Respo return data, BuildResponse(r), nil } +// UpdateJobStatus updates the status of a job +func (c *Client4) UpdateJobStatus(ctx context.Context, jobId string, status string, force bool) (*Response, error) { + buf, err := json.Marshal(map[string]any{ + "status": status, + "force": force, + }) + if err != nil { + return nil, NewAppError("UpdateJobStatus", "api.marshal_error", nil, "", http.StatusInternalServerError).Wrap(err) + } + r, err := c.DoAPIPatchBytes(ctx, c.jobsRoute()+fmt.Sprintf("/%v/status", jobId), buf) + if err != nil { + return BuildResponse(r), err + } + defer closeBody(r) + return BuildResponse(r), nil +} + // Roles Section // GetAllRoles returns a list of all the roles. diff --git a/server/public/model/job.go b/server/public/model/job.go index 603d19009b..14f00f3594 100644 --- a/server/public/model/job.go +++ b/server/public/model/job.go @@ -108,7 +108,31 @@ func (j *Job) IsValid() *AppError { return NewAppError("Job.IsValid", "model.job.is_valid.create_at.app_error", nil, "id="+j.Id, http.StatusBadRequest) } - switch j.Status { + validStatus := IsValidJobStatus(j.Status) + if !validStatus { + return NewAppError("Job.IsValid", "model.job.is_valid.status.app_error", nil, "id="+j.Id, http.StatusBadRequest) + } + + return nil +} + +func (j *Job) IsValidStatusChange(newStatus string) bool { + currentStatus := j.Status + + switch currentStatus { + case JobStatusInProgress: + return newStatus == JobStatusPending || newStatus == JobStatusCancelRequested + case JobStatusPending: + return newStatus == JobStatusCancelRequested + case JobStatusCancelRequested: + return newStatus == JobStatusCanceled + } + + return false +} + +func IsValidJobStatus(status string) bool { + switch status { case JobStatusPending, JobStatusInProgress, JobStatusSuccess, @@ -117,10 +141,20 @@ func (j *Job) IsValid() *AppError { JobStatusCancelRequested, JobStatusCanceled: default: - return NewAppError("Job.IsValid", "model.job.is_valid.status.app_error", nil, "id="+j.Id, http.StatusBadRequest) + return false } - return nil + return true +} + +func IsValidJobType(jobType string) bool { + for _, t := range AllJobTypes { + if t == jobType { + return true + } + } + + return false } func (j *Job) LogClone() any { diff --git a/server/public/model/job_test.go b/server/public/model/job_test.go index 360b533088..40e5e93323 100644 --- a/server/public/model/job_test.go +++ b/server/public/model/job_test.go @@ -121,3 +121,113 @@ func TestJobIsValid(t *testing.T) { } }) } + +func TestJobIsValidStatusChange(t *testing.T) { + t.Run("invalid status change", func(t *testing.T) { + job := &Job{ + Id: "arandomstring0123456789012", + Type: JobTypeExportProcess, + Priority: 42, + CreateAt: 1336, + StartAt: 1337, + LastActivityAt: 1666609360813, + Status: JobStatusInProgress, + Progress: 32, + Data: StringMap{"Hello": "World"}, + } + + require.False(t, job.IsValidStatusChange("invalid!")) + }) + + t.Run("valid status change from in_progress", func(t *testing.T) { + job := &Job{ + Id: "arandomstring0123456789012", + Type: JobTypeExportProcess, + Priority: 42, + CreateAt: 1336, + StartAt: 1337, + LastActivityAt: 1666609360813, + Status: JobStatusInProgress, + Progress: 32, + Data: StringMap{"Hello": "World"}, + } + + require.True(t, job.IsValidStatusChange(JobStatusPending)) + require.True(t, job.IsValidStatusChange(JobStatusCancelRequested)) + require.False(t, job.IsValidStatusChange(JobStatusCanceled)) + }) + + t.Run("valid status change from pending", func(t *testing.T) { + job := &Job{ + Id: "arandomstring0123456789012", + Type: JobTypeExportProcess, + Priority: 42, + CreateAt: 1336, + StartAt: 1337, + LastActivityAt: 1666609360813, + Status: JobStatusPending, + Progress: 32, + Data: StringMap{"Hello": "World"}, + } + + require.True(t, job.IsValidStatusChange(JobStatusCancelRequested)) + require.False(t, job.IsValidStatusChange(JobStatusInProgress)) + }) + + t.Run("valid status change from cancel_requested", func(t *testing.T) { + job := &Job{ + Id: "arandomstring0123456789012", + Type: JobTypeExportProcess, + Priority: 42, + CreateAt: 1336, + StartAt: 1337, + LastActivityAt: 1666609360813, + Status: JobStatusCancelRequested, + Progress: 32, + Data: StringMap{"Hello": "World"}, + } + + require.True(t, job.IsValidStatusChange(JobStatusCanceled)) + require.False(t, job.IsValidStatusChange(JobStatusPending)) + }) +} + +func TestIsValidJobType(t *testing.T) { + t.Run("valid", func(t *testing.T) { + validTypes := []string{JobTypeExportProcess, JobTypeImportProcess} + for _, jobType := range validTypes { + t.Run(jobType, func(t *testing.T) { + require.True(t, IsValidJobType(jobType)) + }) + } + }) + + t.Run("invalid", func(t *testing.T) { + invalidTypes := []string{"invalid!", ""} + for _, jobType := range invalidTypes { + t.Run(jobType, func(t *testing.T) { + require.False(t, IsValidJobType(jobType)) + }) + } + }) +} + +func TestIsValidJobStatus(t *testing.T) { + t.Run("valid", func(t *testing.T) { + validStatuses := []string{JobStatusCancelRequested, JobStatusCanceled, JobStatusError, JobStatusInProgress, JobStatusPending, JobStatusSuccess, JobStatusWarning} + for _, status := range validStatuses { + t.Run(status, func(t *testing.T) { + require.True(t, IsValidJobStatus(status)) + }) + } + }) + + t.Run("invalid", func(t *testing.T) { + invalidStatuses := []string{"invalid!", ""} + for _, status := range invalidStatuses { + t.Run(status, func(t *testing.T) { + require.False(t, IsValidJobStatus(status)) + }) + } + }) +} diff --git a/server/public/model/migration.go b/server/public/model/migration.go index bb48685049..bdd0531bd6 100644 --- a/server/public/model/migration.go +++ b/server/public/model/migration.go @@ -47,4 +47,5 @@ const ( MigrationKeyAddIPFilteringPermissions = "add_ip_filtering_permissions" MigrationKeyAddOutgoingOAuthConnectionsPermissions = "add_outgoing_oauth_connections_permissions" MigrationKeyAddChannelBookmarksPermissions = "add_channel_bookmarks_permissions" + MigrationKeyAddManageJobAncillaryPermissions = "add_manage_jobs_ancillary_permissions" ) diff --git a/server/public/model/permission.go b/server/public/model/permission.go index 39e3594a40..a6864f1186 100644 --- a/server/public/model/permission.go +++ b/server/public/model/permission.go @@ -123,8 +123,10 @@ var PermissionManageSharedChannels *Permission var PermissionManageSecureConnections *Permission var PermissionDownloadComplianceExportResult *Permission var PermissionCreateDataRetentionJob *Permission +var PermissionManageDataRetentionJob *Permission var PermissionReadDataRetentionJob *Permission var PermissionCreateComplianceExportJob *Permission +var PermissionManageComplianceExportJob *Permission var PermissionReadComplianceExportJob *Permission var PermissionReadAudits *Permission var PermissionTestElasticsearch *Permission @@ -136,12 +138,16 @@ var PermissionRecycleDatabaseConnections *Permission var PermissionPurgeElasticsearchIndexes *Permission var PermissionTestEmail *Permission var PermissionCreateElasticsearchPostIndexingJob *Permission +var PermissionManageElasticsearchPostIndexingJob *Permission var PermissionCreateElasticsearchPostAggregationJob *Permission +var PermissionManageElasticsearchPostAggregationJob *Permission var PermissionReadElasticsearchPostIndexingJob *Permission var PermissionReadElasticsearchPostAggregationJob *Permission var PermissionPurgeBleveIndexes *Permission var PermissionCreatePostBleveIndexesJob *Permission +var PermissionManagePostBleveIndexesJob *Permission var PermissionCreateLdapSyncJob *Permission +var PermissionManageLdapSyncJob *Permission var PermissionReadLdapSyncJob *Permission var PermissionTestLdap *Permission var PermissionInvalidateEmailInvite *Permission @@ -790,6 +796,12 @@ func initializePermissions() { "", PermissionScopeSystem, } + PermissionManageDataRetentionJob = &Permission{ + "manage_data_retention_job", + "", + "", + PermissionScopeSystem, + } PermissionReadDataRetentionJob = &Permission{ "read_data_retention_job", "", @@ -803,6 +815,12 @@ func initializePermissions() { "", PermissionScopeSystem, } + PermissionManageComplianceExportJob = &Permission{ + "manage_compliance_export_job", + "", + "", + PermissionScopeSystem, + } PermissionReadComplianceExportJob = &Permission{ "read_compliance_export_job", "", @@ -831,12 +849,25 @@ func initializePermissions() { PermissionScopeSystem, } + PermissionManagePostBleveIndexesJob = &Permission{ + "manage_post_bleve_indexes_job", + "", + "", + PermissionScopeSystem, + } + PermissionCreateLdapSyncJob = &Permission{ "create_ldap_sync_job", "", "", PermissionScopeSystem, } + PermissionManageLdapSyncJob = &Permission{ + "manage_ldap_sync_job", + "", + "", + PermissionScopeSystem, + } PermissionReadLdapSyncJob = &Permission{ "read_ldap_sync_job", "", @@ -1029,12 +1060,24 @@ func initializePermissions() { "", PermissionScopeSystem, } + PermissionManageElasticsearchPostIndexingJob = &Permission{ + "manage_elasticsearch_post_indexing_job", + "", + "", + PermissionScopeSystem, + } PermissionCreateElasticsearchPostAggregationJob = &Permission{ "create_elasticsearch_post_aggregation_job", "", "", PermissionScopeSystem, } + PermissionManageElasticsearchPostAggregationJob = &Permission{ + "manage_elasticsearch_post_aggregation_job", + "", + "", + PermissionScopeSystem, + } PermissionReadElasticsearchPostIndexingJob = &Permission{ "read_elasticsearch_post_indexing_job", "", @@ -2347,8 +2390,10 @@ func initializePermissions() { PermissionManageSecureConnections, PermissionDownloadComplianceExportResult, PermissionCreateDataRetentionJob, + PermissionManageDataRetentionJob, PermissionReadDataRetentionJob, PermissionCreateComplianceExportJob, + PermissionManageComplianceExportJob, PermissionReadComplianceExportJob, PermissionReadAudits, PermissionTestSiteURL, @@ -2360,12 +2405,16 @@ func initializePermissions() { PermissionPurgeElasticsearchIndexes, PermissionTestEmail, PermissionCreateElasticsearchPostIndexingJob, + PermissionManageElasticsearchPostIndexingJob, PermissionCreateElasticsearchPostAggregationJob, + PermissionManageElasticsearchPostAggregationJob, PermissionReadElasticsearchPostIndexingJob, PermissionReadElasticsearchPostAggregationJob, PermissionPurgeBleveIndexes, PermissionCreatePostBleveIndexesJob, + PermissionManagePostBleveIndexesJob, PermissionCreateLdapSyncJob, + PermissionManageLdapSyncJob, PermissionReadLdapSyncJob, PermissionTestLdap, PermissionInvalidateEmailInvite, diff --git a/server/public/model/role.go b/server/public/model/role.go index 8b8619c539..759e87c8d6 100644 --- a/server/public/model/role.go +++ b/server/public/model/role.go @@ -90,7 +90,9 @@ func init() { PermissionSysconsoleWriteEnvironmentElasticsearch.Id: { PermissionTestElasticsearch, PermissionCreateElasticsearchPostIndexingJob, + PermissionManageElasticsearchPostIndexingJob, PermissionCreateElasticsearchPostAggregationJob, + PermissionManageElasticsearchPostAggregationJob, PermissionPurgeElasticsearchIndexes, }, PermissionSysconsoleWriteEnvironmentFileStorage.Id: { @@ -145,12 +147,14 @@ func init() { }, PermissionSysconsoleWriteComplianceDataRetentionPolicy.Id: { PermissionCreateDataRetentionJob, + PermissionManageDataRetentionJob, }, PermissionSysconsoleReadComplianceDataRetentionPolicy.Id: { PermissionReadDataRetentionJob, }, PermissionSysconsoleWriteComplianceComplianceExport.Id: { PermissionCreateComplianceExportJob, + PermissionManageComplianceExportJob, PermissionDownloadComplianceExportResult, }, PermissionSysconsoleReadComplianceComplianceExport.Id: { @@ -163,9 +167,11 @@ func init() { PermissionSysconsoleWriteExperimentalBleve.Id: { PermissionCreatePostBleveIndexesJob, PermissionPurgeBleveIndexes, + PermissionManagePostBleveIndexesJob, }, PermissionSysconsoleWriteAuthenticationLdap.Id: { PermissionCreateLdapSyncJob, + PermissionManageLdapSyncJob, PermissionAddLdapPublicCert, PermissionRemoveLdapPublicCert, PermissionAddLdapPrivateCert,