diff --git a/core/cmd/admin/queue_backup.go b/core/cmd/admin/queue_backup.go index 7886060b..ad8924f0 100644 --- a/core/cmd/admin/queue_backup.go +++ b/core/cmd/admin/queue_backup.go @@ -101,10 +101,50 @@ var listFailedBackupCmd = &cobra.Command{ }, } +//nolint:gochecknoglobals +var retryBackupCmd = &cobra.Command{ + Use: "retry [cluster-name]", + Short: "Retry failed backup tasks in the queue", + Long: "Retry failed backup tasks in the queue.\n\n" + + "With no arguments, all failed backup tasks are retried. If a cluster name " + + "is given, all failed backup tasks for that cluster are retried.", + Args: cobra.MaximumNArgs(1), + RunE: func(cmd *cobra.Command, args []string) error { + socketPath, err := cmd.Flags().GetString("socket-path") + if err != nil { + return fmt.Errorf("while getting the socketPath flag: %w", err) + } + + conn, err := connectToAdminServer(socketPath) + if err != nil { + return err + } + defer func() { + _ = conn.Close() + }() + + var request klioGRPC.QueueRetryBackupsRequest + if len(args) > 0 { + clusterName := args[0] + request.ClusterName = &clusterName + } + + adminClient := klioGRPC.NewAdminClient(conn) + _, err = adminClient.QueueRetryBackups(cmd.Context(), &request) + if err != nil { + return fmt.Errorf("while calling queue retry backups entrypoint: %w", err) + } + + return nil + }, +} + //nolint:gochecknoinits func init() { queueCmd.AddCommand(queueBackupCmd) - queueBackupCmd.AddCommand(listFailedBackupCmd) + queueBackupCmd.AddCommand(listFailedBackupCmd) listFailedBackupCmd.Flags().String("cluster-name", "", "Cluster name to filter failed backup tasks (optional)") + + queueBackupCmd.AddCommand(retryBackupCmd) } diff --git a/core/cmd/admin/queue_wal.go b/core/cmd/admin/queue_wal.go index 2fac56c5..83dc6d23 100644 --- a/core/cmd/admin/queue_wal.go +++ b/core/cmd/admin/queue_wal.go @@ -22,7 +22,6 @@ package admin import ( "fmt" "os" - "strconv" "time" "github.com/spf13/cobra" @@ -89,7 +88,6 @@ var listFailedWALCmd = &cobra.Command{ rows := make([][]string, 0, len(response.GetWals())) for _, wal := range response.GetWals() { rows = append(rows, []string{ - strconv.FormatUint(wal.GetSequence(), 10), wal.GetClusterName(), wal.GetWalName(), wal.GetLastAttemptTime().AsTime().Format(time.RFC3339), @@ -97,7 +95,7 @@ var listFailedWALCmd = &cobra.Command{ } if err := writeTable( os.Stdout, - []string{"SEQUENCE", "CLUSTER", "WAL NAME", "LAST ATTEMPT"}, + []string{"CLUSTER", "WAL NAME", "LAST ATTEMPT"}, rows, ); err != nil { return fmt.Errorf("while writing table output: %w", err) @@ -108,10 +106,54 @@ var listFailedWALCmd = &cobra.Command{ }, } +//nolint:gochecknoglobals +var retryWALCmd = &cobra.Command{ + Use: "retry [cluster-name] [WAL1 WAL2 ...]", + Short: "Retry failed WAL tasks in the queue", + Long: "Retry failed WAL tasks in the queue.\n\n" + + "With no arguments, all failed WAL tasks are retried. If a cluster name " + + "is given, all failed WAL tasks for that cluster are retried. If WAL " + + "files are also given, only those are retried.", + Args: cobra.ArbitraryArgs, + RunE: func(cmd *cobra.Command, args []string) error { + socketPath, err := cmd.Flags().GetString("socket-path") + if err != nil { + return fmt.Errorf("while getting the socketPath flag: %w", err) + } + + conn, err := connectToAdminServer(socketPath) + if err != nil { + return err + } + defer func() { + _ = conn.Close() + }() + + var request klioGRPC.QueueRetryWALsRequest + if len(args) > 0 { + clusterName := args[0] + request.ClusterName = &clusterName + } + if len(args) > 1 { + request.WalNames = args[1:] + } + + adminClient := klioGRPC.NewAdminClient(conn) + _, err = adminClient.QueueRetryWALs(cmd.Context(), &request) + if err != nil { + return fmt.Errorf("while calling queue retry wals entrypoint: %w", err) + } + + return nil + }, +} + //nolint:gochecknoinits func init() { queueCmd.AddCommand(queueWALCmd) - queueWALCmd.AddCommand(listFailedWALCmd) + queueWALCmd.AddCommand(listFailedWALCmd) listFailedWALCmd.Flags().String("cluster-name", "", "Cluster name to filter failed WAL tasks (optional)") + + queueWALCmd.AddCommand(retryWALCmd) } diff --git a/core/internal/grpc/klio_admin.pb.go b/core/internal/grpc/klio_admin.pb.go index 8397605b..68ed0d2a 100644 --- a/core/internal/grpc/klio_admin.pb.go +++ b/core/internal/grpc/klio_admin.pb.go @@ -552,6 +552,138 @@ func (x *FailedWAL) GetLastAttemptTime() *timestamppb.Timestamp { return nil } +type QueueRetryWALsRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + ClusterName *string `protobuf:"bytes,1,opt,name=cluster_name,json=clusterName,proto3,oneof" json:"cluster_name,omitempty"` + WalNames []string `protobuf:"bytes,2,rep,name=wal_names,json=walNames,proto3" json:"wal_names,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *QueueRetryWALsRequest) Reset() { + *x = QueueRetryWALsRequest{} + mi := &file_proto_klio_admin_proto_msgTypes[10] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *QueueRetryWALsRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*QueueRetryWALsRequest) ProtoMessage() {} + +func (x *QueueRetryWALsRequest) ProtoReflect() protoreflect.Message { + mi := &file_proto_klio_admin_proto_msgTypes[10] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use QueueRetryWALsRequest.ProtoReflect.Descriptor instead. +func (*QueueRetryWALsRequest) Descriptor() ([]byte, []int) { + return file_proto_klio_admin_proto_rawDescGZIP(), []int{10} +} + +func (x *QueueRetryWALsRequest) GetClusterName() string { + if x != nil && x.ClusterName != nil { + return *x.ClusterName + } + return "" +} + +func (x *QueueRetryWALsRequest) GetWalNames() []string { + if x != nil { + return x.WalNames + } + return nil +} + +type QueueRetryBackupsRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + ClusterName *string `protobuf:"bytes,1,opt,name=cluster_name,json=clusterName,proto3,oneof" json:"cluster_name,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *QueueRetryBackupsRequest) Reset() { + *x = QueueRetryBackupsRequest{} + mi := &file_proto_klio_admin_proto_msgTypes[11] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *QueueRetryBackupsRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*QueueRetryBackupsRequest) ProtoMessage() {} + +func (x *QueueRetryBackupsRequest) ProtoReflect() protoreflect.Message { + mi := &file_proto_klio_admin_proto_msgTypes[11] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use QueueRetryBackupsRequest.ProtoReflect.Descriptor instead. +func (*QueueRetryBackupsRequest) Descriptor() ([]byte, []int) { + return file_proto_klio_admin_proto_rawDescGZIP(), []int{11} +} + +func (x *QueueRetryBackupsRequest) GetClusterName() string { + if x != nil && x.ClusterName != nil { + return *x.ClusterName + } + return "" +} + +type QueueRetryResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *QueueRetryResponse) Reset() { + *x = QueueRetryResponse{} + mi := &file_proto_klio_admin_proto_msgTypes[12] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *QueueRetryResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*QueueRetryResponse) ProtoMessage() {} + +func (x *QueueRetryResponse) ProtoReflect() protoreflect.Message { + mi := &file_proto_klio_admin_proto_msgTypes[12] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use QueueRetryResponse.ProtoReflect.Descriptor instead. +func (*QueueRetryResponse) Descriptor() ([]byte, []int) { + return file_proto_klio_admin_proto_rawDescGZIP(), []int{12} +} + type QueueStatusRequest struct { state protoimpl.MessageState `protogen:"open.v1"` unknownFields protoimpl.UnknownFields @@ -560,7 +692,7 @@ type QueueStatusRequest struct { func (x *QueueStatusRequest) Reset() { *x = QueueStatusRequest{} - mi := &file_proto_klio_admin_proto_msgTypes[10] + mi := &file_proto_klio_admin_proto_msgTypes[13] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -572,7 +704,7 @@ func (x *QueueStatusRequest) String() string { func (*QueueStatusRequest) ProtoMessage() {} func (x *QueueStatusRequest) ProtoReflect() protoreflect.Message { - mi := &file_proto_klio_admin_proto_msgTypes[10] + mi := &file_proto_klio_admin_proto_msgTypes[13] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -585,7 +717,7 @@ func (x *QueueStatusRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use QueueStatusRequest.ProtoReflect.Descriptor instead. func (*QueueStatusRequest) Descriptor() ([]byte, []int) { - return file_proto_klio_admin_proto_rawDescGZIP(), []int{10} + return file_proto_klio_admin_proto_rawDescGZIP(), []int{13} } type QueueStatusResponse struct { @@ -600,7 +732,7 @@ type QueueStatusResponse struct { func (x *QueueStatusResponse) Reset() { *x = QueueStatusResponse{} - mi := &file_proto_klio_admin_proto_msgTypes[11] + mi := &file_proto_klio_admin_proto_msgTypes[14] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -612,7 +744,7 @@ func (x *QueueStatusResponse) String() string { func (*QueueStatusResponse) ProtoMessage() {} func (x *QueueStatusResponse) ProtoReflect() protoreflect.Message { - mi := &file_proto_klio_admin_proto_msgTypes[11] + mi := &file_proto_klio_admin_proto_msgTypes[14] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -625,7 +757,7 @@ func (x *QueueStatusResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use QueueStatusResponse.ProtoReflect.Descriptor instead. func (*QueueStatusResponse) Descriptor() ([]byte, []int) { - return file_proto_klio_admin_proto_rawDescGZIP(), []int{11} + return file_proto_klio_admin_proto_rawDescGZIP(), []int{14} } func (x *QueueStatusResponse) GetPendingBackups() uint64 { @@ -658,7 +790,7 @@ type DeleteBackupRequest struct { func (x *DeleteBackupRequest) Reset() { *x = DeleteBackupRequest{} - mi := &file_proto_klio_admin_proto_msgTypes[12] + mi := &file_proto_klio_admin_proto_msgTypes[15] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -670,7 +802,7 @@ func (x *DeleteBackupRequest) String() string { func (*DeleteBackupRequest) ProtoMessage() {} func (x *DeleteBackupRequest) ProtoReflect() protoreflect.Message { - mi := &file_proto_klio_admin_proto_msgTypes[12] + mi := &file_proto_klio_admin_proto_msgTypes[15] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -683,7 +815,7 @@ func (x *DeleteBackupRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use DeleteBackupRequest.ProtoReflect.Descriptor instead. func (*DeleteBackupRequest) Descriptor() ([]byte, []int) { - return file_proto_klio_admin_proto_rawDescGZIP(), []int{12} + return file_proto_klio_admin_proto_rawDescGZIP(), []int{15} } func (x *DeleteBackupRequest) GetBackupName() string { @@ -716,7 +848,7 @@ type DeleteBackupResponse struct { func (x *DeleteBackupResponse) Reset() { *x = DeleteBackupResponse{} - mi := &file_proto_klio_admin_proto_msgTypes[13] + mi := &file_proto_klio_admin_proto_msgTypes[16] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -728,7 +860,7 @@ func (x *DeleteBackupResponse) String() string { func (*DeleteBackupResponse) ProtoMessage() {} func (x *DeleteBackupResponse) ProtoReflect() protoreflect.Message { - mi := &file_proto_klio_admin_proto_msgTypes[13] + mi := &file_proto_klio_admin_proto_msgTypes[16] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -741,7 +873,7 @@ func (x *DeleteBackupResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use DeleteBackupResponse.ProtoReflect.Descriptor instead. func (*DeleteBackupResponse) Descriptor() ([]byte, []int) { - return file_proto_klio_admin_proto_rawDescGZIP(), []int{13} + return file_proto_klio_admin_proto_rawDescGZIP(), []int{16} } var File_proto_klio_admin_proto protoreflect.FileDescriptor @@ -771,7 +903,15 @@ const file_proto_klio_admin_proto_rawDesc = "" + "\fcluster_name\x18\x01 \x01(\tR\vclusterName\x12\x19\n" + "\bwal_name\x18\x02 \x01(\tR\awalName\x12\x1a\n" + "\bsequence\x18\x03 \x01(\x04R\bsequence\x12F\n" + - "\x11last_attempt_time\x18\x04 \x01(\v2\x1a.google.protobuf.TimestampR\x0flastAttemptTime\"\x14\n" + + "\x11last_attempt_time\x18\x04 \x01(\v2\x1a.google.protobuf.TimestampR\x0flastAttemptTime\"m\n" + + "\x15QueueRetryWALsRequest\x12&\n" + + "\fcluster_name\x18\x01 \x01(\tH\x00R\vclusterName\x88\x01\x01\x12\x1b\n" + + "\twal_names\x18\x02 \x03(\tR\bwalNamesB\x0f\n" + + "\r_cluster_name\"S\n" + + "\x18QueueRetryBackupsRequest\x12&\n" + + "\fcluster_name\x18\x01 \x01(\tH\x00R\vclusterName\x88\x01\x01B\x0f\n" + + "\r_cluster_name\"\x14\n" + + "\x12QueueRetryResponse\"\x14\n" + "\x12QueueStatusRequest\"a\n" + "\x13QueueStatusResponse\x12'\n" + "\x0fpending_backups\x18\x01 \x01(\x04R\x0ependingBackups\x12!\n" + @@ -787,12 +927,14 @@ const file_proto_klio_admin_proto_rawDesc = "" + "\n" + "\x06TIER_1\x10\x01\x12\n" + "\n" + - "\x06TIER_2\x10\x022\xab\x04\n" + + "\x06TIER_2\x10\x022\xe3\x05\n" + "\x05Admin\x12D\n" + "\aRefresh\x12\x1b.klio.wal.v1.RefreshRequest\x1a\x1a.klio.wal.v1.RefreshResult\"\x00\x12P\n" + "\vListBackups\x12\x1f.klio.wal.v1.ListBackupsRequest\x1a\x1e.klio.wal.v1.ListBackupsResult\"\x00\x12s\n" + "\x16QueueListFailedBackups\x12*.klio.wal.v1.QueueListFailedBackupsRequest\x1a+.klio.wal.v1.QueueListFailedBackupsResponse\"\x00\x12j\n" + - "\x13QueueListFailedWALs\x12'.klio.wal.v1.QueueListFailedWALsRequest\x1a(.klio.wal.v1.QueueListFailedWALsResponse\"\x00\x12R\n" + + "\x13QueueListFailedWALs\x12'.klio.wal.v1.QueueListFailedWALsRequest\x1a(.klio.wal.v1.QueueListFailedWALsResponse\"\x00\x12W\n" + + "\x0eQueueRetryWALs\x12\".klio.wal.v1.QueueRetryWALsRequest\x1a\x1f.klio.wal.v1.QueueRetryResponse\"\x00\x12]\n" + + "\x11QueueRetryBackups\x12%.klio.wal.v1.QueueRetryBackupsRequest\x1a\x1f.klio.wal.v1.QueueRetryResponse\"\x00\x12R\n" + "\vQueueStatus\x12\x1f.klio.wal.v1.QueueStatusRequest\x1a .klio.wal.v1.QueueStatusResponse\"\x00\x12U\n" + "\fDeleteBackup\x12 .klio.wal.v1.DeleteBackupRequest\x1a!.klio.wal.v1.DeleteBackupResponse\"\x00B3Z1github.com/cloudnative-pg/klio/core/internal/grpcb\x06proto3" @@ -809,7 +951,7 @@ func file_proto_klio_admin_proto_rawDescGZIP() []byte { } var file_proto_klio_admin_proto_enumTypes = make([]protoimpl.EnumInfo, 1) -var file_proto_klio_admin_proto_msgTypes = make([]protoimpl.MessageInfo, 14) +var file_proto_klio_admin_proto_msgTypes = make([]protoimpl.MessageInfo, 17) var file_proto_klio_admin_proto_goTypes = []any{ (Tier)(0), // 0: klio.wal.v1.Tier (*RefreshRequest)(nil), // 1: klio.wal.v1.RefreshRequest @@ -822,32 +964,39 @@ var file_proto_klio_admin_proto_goTypes = []any{ (*QueueListFailedWALsResponse)(nil), // 8: klio.wal.v1.QueueListFailedWALsResponse (*FailedBackup)(nil), // 9: klio.wal.v1.FailedBackup (*FailedWAL)(nil), // 10: klio.wal.v1.FailedWAL - (*QueueStatusRequest)(nil), // 11: klio.wal.v1.QueueStatusRequest - (*QueueStatusResponse)(nil), // 12: klio.wal.v1.QueueStatusResponse - (*DeleteBackupRequest)(nil), // 13: klio.wal.v1.DeleteBackupRequest - (*DeleteBackupResponse)(nil), // 14: klio.wal.v1.DeleteBackupResponse - (*timestamppb.Timestamp)(nil), // 15: google.protobuf.Timestamp + (*QueueRetryWALsRequest)(nil), // 11: klio.wal.v1.QueueRetryWALsRequest + (*QueueRetryBackupsRequest)(nil), // 12: klio.wal.v1.QueueRetryBackupsRequest + (*QueueRetryResponse)(nil), // 13: klio.wal.v1.QueueRetryResponse + (*QueueStatusRequest)(nil), // 14: klio.wal.v1.QueueStatusRequest + (*QueueStatusResponse)(nil), // 15: klio.wal.v1.QueueStatusResponse + (*DeleteBackupRequest)(nil), // 16: klio.wal.v1.DeleteBackupRequest + (*DeleteBackupResponse)(nil), // 17: klio.wal.v1.DeleteBackupResponse + (*timestamppb.Timestamp)(nil), // 18: google.protobuf.Timestamp } var file_proto_klio_admin_proto_depIdxs = []int32{ 9, // 0: klio.wal.v1.QueueListFailedBackupsResponse.backups:type_name -> klio.wal.v1.FailedBackup 10, // 1: klio.wal.v1.QueueListFailedWALsResponse.wals:type_name -> klio.wal.v1.FailedWAL - 15, // 2: klio.wal.v1.FailedBackup.last_attempt_time:type_name -> google.protobuf.Timestamp - 15, // 3: klio.wal.v1.FailedWAL.last_attempt_time:type_name -> google.protobuf.Timestamp + 18, // 2: klio.wal.v1.FailedBackup.last_attempt_time:type_name -> google.protobuf.Timestamp + 18, // 3: klio.wal.v1.FailedWAL.last_attempt_time:type_name -> google.protobuf.Timestamp 0, // 4: klio.wal.v1.DeleteBackupRequest.tiers:type_name -> klio.wal.v1.Tier 1, // 5: klio.wal.v1.Admin.Refresh:input_type -> klio.wal.v1.RefreshRequest 3, // 6: klio.wal.v1.Admin.ListBackups:input_type -> klio.wal.v1.ListBackupsRequest 5, // 7: klio.wal.v1.Admin.QueueListFailedBackups:input_type -> klio.wal.v1.QueueListFailedBackupsRequest 7, // 8: klio.wal.v1.Admin.QueueListFailedWALs:input_type -> klio.wal.v1.QueueListFailedWALsRequest - 11, // 9: klio.wal.v1.Admin.QueueStatus:input_type -> klio.wal.v1.QueueStatusRequest - 13, // 10: klio.wal.v1.Admin.DeleteBackup:input_type -> klio.wal.v1.DeleteBackupRequest - 2, // 11: klio.wal.v1.Admin.Refresh:output_type -> klio.wal.v1.RefreshResult - 4, // 12: klio.wal.v1.Admin.ListBackups:output_type -> klio.wal.v1.ListBackupsResult - 6, // 13: klio.wal.v1.Admin.QueueListFailedBackups:output_type -> klio.wal.v1.QueueListFailedBackupsResponse - 8, // 14: klio.wal.v1.Admin.QueueListFailedWALs:output_type -> klio.wal.v1.QueueListFailedWALsResponse - 12, // 15: klio.wal.v1.Admin.QueueStatus:output_type -> klio.wal.v1.QueueStatusResponse - 14, // 16: klio.wal.v1.Admin.DeleteBackup:output_type -> klio.wal.v1.DeleteBackupResponse - 11, // [11:17] is the sub-list for method output_type - 5, // [5:11] is the sub-list for method input_type + 11, // 9: klio.wal.v1.Admin.QueueRetryWALs:input_type -> klio.wal.v1.QueueRetryWALsRequest + 12, // 10: klio.wal.v1.Admin.QueueRetryBackups:input_type -> klio.wal.v1.QueueRetryBackupsRequest + 14, // 11: klio.wal.v1.Admin.QueueStatus:input_type -> klio.wal.v1.QueueStatusRequest + 16, // 12: klio.wal.v1.Admin.DeleteBackup:input_type -> klio.wal.v1.DeleteBackupRequest + 2, // 13: klio.wal.v1.Admin.Refresh:output_type -> klio.wal.v1.RefreshResult + 4, // 14: klio.wal.v1.Admin.ListBackups:output_type -> klio.wal.v1.ListBackupsResult + 6, // 15: klio.wal.v1.Admin.QueueListFailedBackups:output_type -> klio.wal.v1.QueueListFailedBackupsResponse + 8, // 16: klio.wal.v1.Admin.QueueListFailedWALs:output_type -> klio.wal.v1.QueueListFailedWALsResponse + 13, // 17: klio.wal.v1.Admin.QueueRetryWALs:output_type -> klio.wal.v1.QueueRetryResponse + 13, // 18: klio.wal.v1.Admin.QueueRetryBackups:output_type -> klio.wal.v1.QueueRetryResponse + 15, // 19: klio.wal.v1.Admin.QueueStatus:output_type -> klio.wal.v1.QueueStatusResponse + 17, // 20: klio.wal.v1.Admin.DeleteBackup:output_type -> klio.wal.v1.DeleteBackupResponse + 13, // [13:21] is the sub-list for method output_type + 5, // [5:13] is the sub-list for method input_type 5, // [5:5] is the sub-list for extension type_name 5, // [5:5] is the sub-list for extension extendee 0, // [0:5] is the sub-list for field type_name @@ -860,13 +1009,15 @@ func file_proto_klio_admin_proto_init() { } file_proto_klio_admin_proto_msgTypes[4].OneofWrappers = []any{} file_proto_klio_admin_proto_msgTypes[6].OneofWrappers = []any{} + file_proto_klio_admin_proto_msgTypes[10].OneofWrappers = []any{} + file_proto_klio_admin_proto_msgTypes[11].OneofWrappers = []any{} type x struct{} out := protoimpl.TypeBuilder{ File: protoimpl.DescBuilder{ GoPackagePath: reflect.TypeOf(x{}).PkgPath(), RawDescriptor: unsafe.Slice(unsafe.StringData(file_proto_klio_admin_proto_rawDesc), len(file_proto_klio_admin_proto_rawDesc)), NumEnums: 1, - NumMessages: 14, + NumMessages: 17, NumExtensions: 0, NumServices: 1, }, diff --git a/core/internal/grpc/klio_admin_grpc.pb.go b/core/internal/grpc/klio_admin_grpc.pb.go index 5568990a..60c25f27 100644 --- a/core/internal/grpc/klio_admin_grpc.pb.go +++ b/core/internal/grpc/klio_admin_grpc.pb.go @@ -42,6 +42,8 @@ const ( Admin_ListBackups_FullMethodName = "/klio.wal.v1.Admin/ListBackups" Admin_QueueListFailedBackups_FullMethodName = "/klio.wal.v1.Admin/QueueListFailedBackups" Admin_QueueListFailedWALs_FullMethodName = "/klio.wal.v1.Admin/QueueListFailedWALs" + Admin_QueueRetryWALs_FullMethodName = "/klio.wal.v1.Admin/QueueRetryWALs" + Admin_QueueRetryBackups_FullMethodName = "/klio.wal.v1.Admin/QueueRetryBackups" Admin_QueueStatus_FullMethodName = "/klio.wal.v1.Admin/QueueStatus" Admin_DeleteBackup_FullMethodName = "/klio.wal.v1.Admin/DeleteBackup" ) @@ -58,6 +60,10 @@ type AdminClient interface { QueueListFailedBackups(ctx context.Context, in *QueueListFailedBackupsRequest, opts ...grpc.CallOption) (*QueueListFailedBackupsResponse, error) // List WAL files failed to be processed from the queue QueueListFailedWALs(ctx context.Context, in *QueueListFailedWALsRequest, opts ...grpc.CallOption) (*QueueListFailedWALsResponse, error) + // Retry WAL files that failed to be processed from the queue + QueueRetryWALs(ctx context.Context, in *QueueRetryWALsRequest, opts ...grpc.CallOption) (*QueueRetryResponse, error) + // Retry Backups that failed to be processed from the queue + QueueRetryBackups(ctx context.Context, in *QueueRetryBackupsRequest, opts ...grpc.CallOption) (*QueueRetryResponse, error) // Get the status of the task queue (pending backups and WALs) QueueStatus(ctx context.Context, in *QueueStatusRequest, opts ...grpc.CallOption) (*QueueStatusResponse, error) // Delete a backup from the server @@ -112,6 +118,26 @@ func (c *adminClient) QueueListFailedWALs(ctx context.Context, in *QueueListFail return out, nil } +func (c *adminClient) QueueRetryWALs(ctx context.Context, in *QueueRetryWALsRequest, opts ...grpc.CallOption) (*QueueRetryResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(QueueRetryResponse) + err := c.cc.Invoke(ctx, Admin_QueueRetryWALs_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *adminClient) QueueRetryBackups(ctx context.Context, in *QueueRetryBackupsRequest, opts ...grpc.CallOption) (*QueueRetryResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(QueueRetryResponse) + err := c.cc.Invoke(ctx, Admin_QueueRetryBackups_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + func (c *adminClient) QueueStatus(ctx context.Context, in *QueueStatusRequest, opts ...grpc.CallOption) (*QueueStatusResponse, error) { cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) out := new(QueueStatusResponse) @@ -144,6 +170,10 @@ type AdminServer interface { QueueListFailedBackups(context.Context, *QueueListFailedBackupsRequest) (*QueueListFailedBackupsResponse, error) // List WAL files failed to be processed from the queue QueueListFailedWALs(context.Context, *QueueListFailedWALsRequest) (*QueueListFailedWALsResponse, error) + // Retry WAL files that failed to be processed from the queue + QueueRetryWALs(context.Context, *QueueRetryWALsRequest) (*QueueRetryResponse, error) + // Retry Backups that failed to be processed from the queue + QueueRetryBackups(context.Context, *QueueRetryBackupsRequest) (*QueueRetryResponse, error) // Get the status of the task queue (pending backups and WALs) QueueStatus(context.Context, *QueueStatusRequest) (*QueueStatusResponse, error) // Delete a backup from the server @@ -170,6 +200,12 @@ func (UnimplementedAdminServer) QueueListFailedBackups(context.Context, *QueueLi func (UnimplementedAdminServer) QueueListFailedWALs(context.Context, *QueueListFailedWALsRequest) (*QueueListFailedWALsResponse, error) { return nil, status.Error(codes.Unimplemented, "method QueueListFailedWALs not implemented") } +func (UnimplementedAdminServer) QueueRetryWALs(context.Context, *QueueRetryWALsRequest) (*QueueRetryResponse, error) { + return nil, status.Error(codes.Unimplemented, "method QueueRetryWALs not implemented") +} +func (UnimplementedAdminServer) QueueRetryBackups(context.Context, *QueueRetryBackupsRequest) (*QueueRetryResponse, error) { + return nil, status.Error(codes.Unimplemented, "method QueueRetryBackups not implemented") +} func (UnimplementedAdminServer) QueueStatus(context.Context, *QueueStatusRequest) (*QueueStatusResponse, error) { return nil, status.Error(codes.Unimplemented, "method QueueStatus not implemented") } @@ -269,6 +305,42 @@ func _Admin_QueueListFailedWALs_Handler(srv interface{}, ctx context.Context, de return interceptor(ctx, in, info, handler) } +func _Admin_QueueRetryWALs_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(QueueRetryWALsRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(AdminServer).QueueRetryWALs(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: Admin_QueueRetryWALs_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(AdminServer).QueueRetryWALs(ctx, req.(*QueueRetryWALsRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _Admin_QueueRetryBackups_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(QueueRetryBackupsRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(AdminServer).QueueRetryBackups(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: Admin_QueueRetryBackups_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(AdminServer).QueueRetryBackups(ctx, req.(*QueueRetryBackupsRequest)) + } + return interceptor(ctx, in, info, handler) +} + func _Admin_QueueStatus_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { in := new(QueueStatusRequest) if err := dec(in); err != nil { @@ -328,6 +400,14 @@ var Admin_ServiceDesc = grpc.ServiceDesc{ MethodName: "QueueListFailedWALs", Handler: _Admin_QueueListFailedWALs_Handler, }, + { + MethodName: "QueueRetryWALs", + Handler: _Admin_QueueRetryWALs_Handler, + }, + { + MethodName: "QueueRetryBackups", + Handler: _Admin_QueueRetryBackups_Handler, + }, { MethodName: "QueueStatus", Handler: _Admin_QueueStatus_Handler, diff --git a/core/internal/queue/backup.go b/core/internal/queue/backup.go index 49688bda..fa3e7db4 100644 --- a/core/internal/queue/backup.go +++ b/core/internal/queue/backup.go @@ -51,7 +51,7 @@ func (t BackupTask) Cluster() string { // NotifyBackupReceived is called to notify the consumers that a new backup // has been uploaded. func (q *Conn) NotifyBackupReceived(ctx context.Context, task *BackupTask) error { - return q.notifyMessage(ctx, backupSubject(task.ClusterName), task) + return q.notifyMessage(ctx, backupSubject(task.ClusterName), task, nil) } // BackupTaskHandler is called for every backup task message that should be handled. @@ -62,7 +62,16 @@ type BackupTaskHandler func(ctx context.Context, t *BackupTask) error // when the context is canceled. After a successful handler run, all dead-letter // queue entries for the backed-up cluster are purged. func (q *Conn) ConsumeBackupReceivedMessages(ctx context.Context, handler BackupTaskHandler) error { - wrapped := func(ctx context.Context, t *BackupTask, _ nats.Header) error { + logger := log.FromContext(ctx).WithName("backup-consumer") + + wrapped := func(ctx context.Context, t *BackupTask, headers nats.Header) error { + if isDLQRetry(headers) { + logger.Info( + "Retrying backup task re-enqueued from the dead-letter queue", + "cluster", t.ClusterName, + ) + } + if err := handler(ctx, t); err != nil { return err } diff --git a/core/internal/queue/dlq_test.go b/core/internal/queue/dlq_test.go index a3b2c83a..038a4d03 100644 --- a/core/internal/queue/dlq_test.go +++ b/core/internal/queue/dlq_test.go @@ -23,7 +23,6 @@ import ( "context" "encoding/json" "fmt" - "strconv" "testing" "time" @@ -38,18 +37,18 @@ import ( const testWALName = "000000010000000000000001" // publishWALMessage publishes a WAL task to the WAL work-queue stream and -// returns its assigned stream sequence. When dlqRetrySeq is non-zero, the -// message carries the CLI retry marker pointing at that DLQ sequence. -func publishWALMessage(t *testing.T, js jetstream.JetStream, clusterName string, dlqRetrySeq uint64) uint64 { +// returns its assigned stream sequence. When retried is true, the message +// carries the dead-letter queue retry origin marker. +func publishWALMessage(t *testing.T, js jetstream.JetStream, clusterName string, retried bool) uint64 { t.Helper() data, err := json.Marshal(WALTask{ClusterName: clusterName, WALName: testWALName}) require.NoError(t, err) msg := &nats.Msg{Subject: walSubject(clusterName), Data: data} - if dlqRetrySeq != 0 { + if retried { msg.Header = nats.Header{} - msg.Header.Set(DLQAdvisorySequenceHeader, strconv.FormatUint(dlqRetrySeq, 10)) + msg.Header.Set(TaskOriginHeaderKey, TaskOriginDLQRetry) } ack, err := js.PublishMsg(t.Context(), msg) @@ -74,9 +73,8 @@ func publishBackupMessage(t *testing.T, js jetstream.JetStream, clusterName stri // seedDLQAdvisory publishes a synthetic max-deliveries advisory onto the // dead-letter queue subject for the given stream/consumer, pointing at the -// original message sequence streamSeq, and returns the advisory's sequence in -// the DLQ stream. -func seedDLQAdvisory(t *testing.T, js jetstream.JetStream, streamName, consumerName string, streamSeq uint64) uint64 { +// original message sequence streamSeq. +func seedDLQAdvisory(t *testing.T, js jetstream.JetStream, streamName, consumerName string, streamSeq uint64) { t.Helper() advisory := server.JSConsumerDeliveryExceededAdvisory{ @@ -89,10 +87,8 @@ func seedDLQAdvisory(t *testing.T, js jetstream.JetStream, streamName, consumerN subject := fmt.Sprintf("%s.%s.%s", server.JSAdvisoryConsumerMaxDeliveryExceedPre, streamName, consumerName) - ack, err := js.Publish(t.Context(), subject, data) + _, err = js.Publish(t.Context(), subject, data) require.NoError(t, err) - - return ack.Sequence } // dlqMsgCount returns the number of messages currently stored in stream. @@ -105,7 +101,7 @@ func dlqMsgCount(t *testing.T, stream jetstream.Stream) uint64 { return info.State.Msgs } -func TestPurgeWALDLQEntryRemovesEntryAndReleasesOriginal(t *testing.T) { +func TestPurgeWALDLQEntriesRemovesEntryAndReleasesOriginal(t *testing.T) { ns, url := startNATSServer(t) defer ns.Shutdown() @@ -122,20 +118,20 @@ func TestPurgeWALDLQEntryRemovesEntryAndReleasesOriginal(t *testing.T) { // The original failed WAL message stays in the work queue after exhausting // its delivery budget; its DLQ advisory points at that sequence. - poisonSeq := publishWALMessage(t, js, "purge-cluster", 0) - dlqSeq := seedDLQAdvisory(t, js, klioWalStreamName, klioWalConsumerName, poisonSeq) + poisonSeq := publishWALMessage(t, js, "purge-cluster", false) + seedDLQAdvisory(t, js, klioWalStreamName, klioWalConsumerName, poisonSeq) require.Equal(t, uint64(1), dlqMsgCount(t, streamHandle(ctx, t, conn.conn, klioDLQWalStreamName))) require.Equal(t, uint64(1), dlqMsgCount(t, streamHandle(ctx, t, conn.conn, klioWalStreamName))) - require.NoError(t, conn.purgeWALDLQEntry(ctx, dlqSeq)) + require.NoError(t, conn.purgeWALDLQEntries(ctx, "purge-cluster", testWALName)) assert.Equal(t, uint64(0), dlqMsgCount(t, streamHandle(ctx, t, conn.conn, klioDLQWalStreamName)), - "the dead-letter queue entry must be purged by sequence") + "the dead-letter queue entry for the cluster and WAL must be purged") assert.Equal(t, uint64(0), dlqMsgCount(t, streamHandle(ctx, t, conn.conn, klioWalStreamName)), "the original failed WAL message must be released from the work queue") } -func TestPurgeWALDLQEntryToleratesMissingMessages(t *testing.T) { +func TestPurgeWALDLQEntriesLeavesOtherWALsUntouched(t *testing.T) { ns, url := startNATSServer(t) defer ns.Shutdown() @@ -150,18 +146,41 @@ func TestPurgeWALDLQEntryToleratesMissingMessages(t *testing.T) { js, err := jetstream.New(nc) require.NoError(t, err) - // A purge for a DLQ sequence that does not exist must be a no-op. - require.NoError(t, conn.purgeWALDLQEntry(ctx, 999)) + // Two failed WALs for the same cluster: only the one matching the requested + // WAL name must be purged. + targetSeq := publishWALMessage(t, js, "multi-cluster", false) + seedDLQAdvisory(t, js, klioWalStreamName, klioWalConsumerName, targetSeq) - // A purge whose original message is already gone must still remove the DLQ - // entry without error. - poisonSeq := publishWALMessage(t, js, "gone-cluster", 0) - dlqSeq := seedDLQAdvisory(t, js, klioWalStreamName, klioWalConsumerName, poisonSeq) - require.NoError(t, streamHandle(ctx, t, conn.conn, klioWalStreamName).DeleteMsg(ctx, poisonSeq)) + otherData, err := json.Marshal(WALTask{ClusterName: "multi-cluster", WALName: "000000010000000000000002"}) + require.NoError(t, err) + otherAck, err := js.PublishMsg(ctx, &nats.Msg{Subject: walSubject("multi-cluster"), Data: otherData}) + require.NoError(t, err) + seedDLQAdvisory(t, js, klioWalStreamName, klioWalConsumerName, otherAck.Sequence) - require.NoError(t, conn.purgeWALDLQEntry(ctx, dlqSeq)) - assert.Equal(t, uint64(0), dlqMsgCount(t, streamHandle(ctx, t, conn.conn, klioDLQWalStreamName)), - "the dead-letter queue entry must be purged even when its original is gone") + require.Equal(t, uint64(2), dlqMsgCount(t, streamHandle(ctx, t, conn.conn, klioDLQWalStreamName))) + + require.NoError(t, conn.purgeWALDLQEntries(ctx, "multi-cluster", testWALName)) + + assert.Equal(t, uint64(1), dlqMsgCount(t, streamHandle(ctx, t, conn.conn, klioDLQWalStreamName)), + "only the entry matching the requested WAL name must be purged") + assert.Equal(t, uint64(1), dlqMsgCount(t, streamHandle(ctx, t, conn.conn, klioWalStreamName)), + "the non-matching WAL's original message must be retained") +} + +func TestPurgeWALDLQEntriesWithNoMatchingEntriesIsNoOp(t *testing.T) { + ns, url := startNATSServer(t) + defer ns.Shutdown() + + nc, err := nats.Connect(url) + require.NoError(t, err) + defer nc.Close() + + ctx := context.Background() + conn, err := New(ctx, nc) + require.NoError(t, err) + + // A purge for a cluster/WAL with no dead-letter queue entries is a no-op. + require.NoError(t, conn.purgeWALDLQEntries(ctx, "empty-cluster", testWALName)) } func TestPurgeBackupDLQEntriesRemovesClusterEntries(t *testing.T) { @@ -213,14 +232,14 @@ func TestWALConsumerPurgesDLQOnRetrySuccess(t *testing.T) { js, err := jetstream.New(nc) require.NoError(t, err) - // The original failed WAL message stays in the work queue; the retry - // republishes the same task carrying the DLQ sequence in its marker. - poisonSeq := publishWALMessage(t, js, "retry-cluster", 0) - dlqSeq := seedDLQAdvisory(t, js, klioWalStreamName, klioWalConsumerName, poisonSeq) + // The retried task carries the DLQ retry origin marker. Its advisory points + // at the task being processed; the purge runs inside the handler before the + // message is acked, so its original is still present and is released + // together with the advisory. + retrySeq := publishWALMessage(t, js, "retry-cluster", true) + seedDLQAdvisory(t, js, klioWalStreamName, klioWalConsumerName, retrySeq) require.Equal(t, uint64(1), dlqMsgCount(t, streamHandle(ctx, t, conn.conn, klioDLQWalStreamName))) - publishWALMessage(t, js, "retry-cluster", dlqSeq) - handler := func(_ context.Context, _ *WALTask) error { return nil } go func() { _ = conn.ConsumeWALReceivedMessages(ctx, handler) @@ -230,7 +249,7 @@ func TestWALConsumerPurgesDLQOnRetrySuccess(t *testing.T) { info, infoErr := streamHandle(ctx, t, conn.conn, klioDLQWalStreamName).Info(ctx) return infoErr == nil && info.State.Msgs == 0 }, 5*time.Second, 50*time.Millisecond, - "a successful CLI retry must purge the referenced WAL dead-letter queue entry") + "a successful retry must purge the dead-letter queue entry for the cluster and WAL") } func TestWALConsumerSkipsDLQWithoutRetryMarker(t *testing.T) { @@ -248,7 +267,7 @@ func TestWALConsumerSkipsDLQWithoutRetryMarker(t *testing.T) { js, err := jetstream.New(nc) require.NoError(t, err) - poisonSeq := publishWALMessage(t, js, "normal-cluster", 0) + poisonSeq := publishWALMessage(t, js, "normal-cluster", false) seedDLQAdvisory(t, js, klioWalStreamName, klioWalConsumerName, poisonSeq) handler := func(_ context.Context, _ *WALTask) error { return nil } diff --git a/core/internal/queue/manager.go b/core/internal/queue/manager.go index eb8d58c8..122543be 100644 --- a/core/internal/queue/manager.go +++ b/core/internal/queue/manager.go @@ -48,6 +48,10 @@ var errIncompleteDLQListing = errors.New("incomplete DLQ listing") // empty result without one is ambiguous. var errAmbiguousSourceRead = errors.New("ambiguous source stream read: no message and no error") +// errWALFilterUnsupported indicates WithWALs was passed to an operation on a +// task type that has no WAL name to filter on (e.g. backups). +var errWALFilterUnsupported = errors.New("WAL name filtering is only supported for WAL tasks") + // FailedTask represents a task that has failed and has been sent to the Dead Letter Queue (DLQ) stream. type FailedTask[T clusterTask] struct { // Sequence is the sequence number of the message in the DLQ stream. @@ -63,26 +67,36 @@ type clusterTask interface { Cluster() string } -// ListOption configures a DLQ listing call. Pass options to -// ListFailedWALTasks / ListFailedBackupTasks via functional options -// (e.g., WithCluster("foo")). -type ListOption func(*listConfig) +// Option configures a queue operation on failed tasks. +// Pass options via functional options (e.g., WithCluster("foo")). Each +// operation uses only the subset of fields relevant to it. +type Option func(*optionConfig) -type listConfig struct { +type optionConfig struct { cluster string + wals []string } -// WithCluster restricts the returned DLQ entries to those whose original -// task belongs to the given cluster. An empty cluster name is a no-op. -func WithCluster(name string) ListOption { - return func(c *listConfig) { +// WithCluster restricts the operation to failed tasks whose original task +// belongs to the given cluster. +func WithCluster(name string) Option { + return func(c *optionConfig) { c.cluster = name } } +// WithWALs restricts the operation to failed tasks for the given WAL file +// names. +func WithWALs(wals ...string) Option { + return func(c *optionConfig) { + c.wals = wals + } +} + // StreamManager provides methods to interact with NATS streams. type StreamManager struct { mgr *jsm.Manager + js jetstream.JetStream mu sync.Mutex streams map[string]*jsm.Stream @@ -95,8 +109,14 @@ func NewStreamManager(conn *nats.Conn) (*StreamManager, error) { return nil, err } + js, err := jetstream.New(conn) + if err != nil { + return nil, fmt.Errorf("while creating JetStream instance: %w", err) + } + return &StreamManager{ mgr: mgr, + js: js, streams: make(map[string]*jsm.Stream), }, nil } @@ -121,7 +141,7 @@ func (m *StreamManager) GetStatus() (*Status, error) { } // ListFailedWALTasks retrieves a list of failed WAL tasks from the Dead Letter Queue (DLQ) stream. -func (m *StreamManager) ListFailedWALTasks(ctx context.Context, opts ...ListOption) ([]FailedTask[WALTask], error) { +func (m *StreamManager) ListFailedWALTasks(ctx context.Context, opts ...Option) ([]FailedTask[WALTask], error) { walStream, err := m.loadStreamOrNil(klioWalStreamName) if err != nil { return nil, err @@ -136,14 +156,37 @@ func (m *StreamManager) ListFailedWALTasks(ctx context.Context, opts ...ListOpti return nil, nil } - return listFailedTasks[WALTask](ctx, dlqWALStream, walStream, opts...) + tasks, err := listFailedTasks[WALTask](ctx, dlqWALStream, walStream, opts...) + if err != nil { + return nil, err + } + + var cfg optionConfig + for _, opt := range opts { + opt(&cfg) + } + if len(cfg.wals) > 0 { + tasks = slices.DeleteFunc(tasks, func(task FailedTask[WALTask]) bool { + return !slices.Contains(cfg.wals, task.Task.WALName) + }) + } + + return tasks, nil } // ListFailedBackupTasks retrieves a list of failed backup tasks from the Dead Letter Queue (DLQ) stream. func (m *StreamManager) ListFailedBackupTasks( ctx context.Context, - opts ...ListOption, + opts ...Option, ) ([]FailedTask[BackupTask], error) { + var cfg optionConfig + for _, opt := range opts { + opt(&cfg) + } + if len(cfg.wals) > 0 { + return nil, errWALFilterUnsupported + } + backupStream, err := m.loadStreamOrNil(klioBackupStreamName) if err != nil { return nil, err @@ -161,6 +204,83 @@ func (m *StreamManager) ListFailedBackupTasks( return listFailedTasks[BackupTask](ctx, dlqBackupStream, backupStream, opts...) } +// RetryFailedWALTasks re-enqueues failed WAL tasks from the dead-letter queue. +func (m *StreamManager) RetryFailedWALTasks( + ctx context.Context, + opts ...Option, +) error { + failedTasks, err := m.ListFailedWALTasks(ctx, opts...) + if err != nil { + return fmt.Errorf("while listing failed WAL tasks: %w", err) + } + + return m.enqueueWALTasks(ctx, failedTasks) +} + +// RetryFailedBackupTasks re-enqueues failed backup tasks from the dead-letter queue. +func (m *StreamManager) RetryFailedBackupTasks( + ctx context.Context, + opts ...Option, +) error { + failedTasks, err := m.ListFailedBackupTasks(ctx, opts...) + if err != nil { + return fmt.Errorf("while listing failed backup tasks: %w", err) + } + + return m.enqueueBackupTasks(ctx, failedTasks) +} + +// enqueueWALTasks re-publishes the given failed WAL tasks onto the work queue +// carrying the DLQ retry origin marker, skipping duplicate tasks. +func (m *StreamManager) enqueueWALTasks(ctx context.Context, tasks []FailedTask[WALTask]) error { + attempted := make(map[WALTask]struct{}, len(tasks)) + for _, task := range tasks { + if _, ok := attempted[task.Task]; ok { + continue + } + if err := m.notifyMessage( + ctx, + walSubject(task.Task.Cluster()), + task.Task, + nats.Header{ + TaskOriginHeaderKey: []string{TaskOriginDLQRetry}, + }, + ); err != nil { + return fmt.Errorf("while retrying failed WAL task for cluster %s, wal %s: %w", + task.Task.ClusterName, + task.Task.WALName, + err) + } + attempted[task.Task] = struct{}{} + } + + return nil +} + +// enqueueBackupTasks re-publishes the given failed backup tasks onto the work +// queue carrying the DLQ retry origin marker, skipping duplicate tasks. +func (m *StreamManager) enqueueBackupTasks(ctx context.Context, tasks []FailedTask[BackupTask]) error { + attempted := make(map[string]struct{}, len(tasks)) + for _, task := range tasks { + if _, ok := attempted[task.Task.ClusterName]; ok { + continue + } + if err := m.notifyMessage( + ctx, + backupSubject(task.Task.Cluster()), + task.Task, + nats.Header{ + TaskOriginHeaderKey: []string{TaskOriginDLQRetry}, + }, + ); err != nil { + return fmt.Errorf("while retrying failed backup task for cluster %s: %w", task.Task.ClusterName, err) + } + attempted[task.Task.ClusterName] = struct{}{} + } + + return nil +} + // configureStreams creates or updates all JetStream streams required by Klio. func (m *StreamManager) configureStreams(ctx context.Context, js jetstream.JetStream) error { configs := []jetstream.StreamConfig{ @@ -221,10 +341,7 @@ func (m *StreamManager) configureStreams(ctx context.Context, js jetstream.JetSt return nil } -// purgeWALDLQEntry removes the WAL dead-letter queue entry at the given stream -// sequence and releases the original message it references from the WAL -// work-queue stream. -func (m *StreamManager) purgeWALDLQEntry(_ context.Context, dlqSequence uint64) error { +func (m *StreamManager) purgeWALDLQEntries(ctx context.Context, clusterName, walName string) error { dlqStream, err := m.loadStreamOrNil(klioDLQWalStreamName) if err != nil { return err @@ -237,7 +354,22 @@ func (m *StreamManager) purgeWALDLQEntry(_ context.Context, dlqSequence uint64) return nil } - return m.purgeDLQEntryBySequence(dlqStream, sourceStream, dlqSequence) + failed, err := listFailedTasks[WALTask](ctx, dlqStream, sourceStream, WithCluster(clusterName)) + if err != nil { + return err + } + + var errs []error + for _, task := range failed { + if task.Task.WALName != walName { + continue + } + if err := m.purgeDLQEntryBySequence(dlqStream, sourceStream, task.Sequence); err != nil { + errs = append(errs, err) + } + } + + return errors.Join(errs...) } // purgeBackupDLQEntries removes every backup dead-letter queue entry belonging to the @@ -367,12 +499,36 @@ func (m *StreamManager) purgeDLQEntryBySequence( return nil } +// notifyMessage is called to send a message on the queue. +func (m *StreamManager) notifyMessage(ctx context.Context, subject string, task any, headers nats.Header) error { + contextLogger := log.FromContext(ctx) + contextLogger.Info("Sending message", "subject", subject, "task", task) + + rawContent, err := json.Marshal(task) + if err != nil { + return fmt.Errorf("while marshalling task to JSON: %w", err) + } + + msg := &nats.Msg{ + Subject: subject, + Data: rawContent, + Header: headers, + } + + _, err = m.js.PublishMsg(ctx, msg) + if err != nil { + return fmt.Errorf("while pushing message to the queue: %w", err) + } + + return nil +} + func listFailedTasks[T clusterTask]( ctx context.Context, dlqStream, taskStream *jsm.Stream, - opts ...ListOption, + opts ...Option, ) ([]FailedTask[T], error) { - var cfg listConfig + var cfg optionConfig for _, opt := range opts { opt(&cfg) } @@ -412,7 +568,7 @@ func listPager[T clusterTask](ctx context.Context, pgr *jsm.StreamPager, readMessage func(seq uint64) (*api.StoredMsg, error), expected uint64, - cfg listConfig, + cfg optionConfig, ) ([]FailedTask[T], error) { // The DLQ pager and the source-stream reads share the same NATS // connection, and the pager's reply inbox shares the connection's diff --git a/core/internal/queue/manager_test.go b/core/internal/queue/manager_test.go index 6e5ece80..47c9d0ee 100644 --- a/core/internal/queue/manager_test.go +++ b/core/internal/queue/manager_test.go @@ -232,7 +232,7 @@ func TestListPagerShortRead(t *testing.T) { // Claim the stream holds two entries while the pager will only ever deliver // the single seeded one, simulating an early pager termination. - _, err = listPager[WALTask](t.Context(), pgr, source.ReadMessage, 2, listConfig{}) + _, err = listPager[WALTask](t.Context(), pgr, source.ReadMessage, 2, optionConfig{}) require.ErrorIs(t, err, errIncompleteDLQListing) } @@ -250,7 +250,7 @@ func TestListPagerContextCancelled(t *testing.T) { ctx, cancel := context.WithCancel(t.Context()) cancel() - _, err = listPager[WALTask](ctx, pgr, source.ReadMessage, 1, listConfig{}) + _, err = listPager[WALTask](ctx, pgr, source.ReadMessage, 1, optionConfig{}) require.ErrorIs(t, err, context.Canceled) } diff --git a/core/internal/queue/queue.go b/core/internal/queue/queue.go index b2b5524d..fa408e9a 100644 --- a/core/internal/queue/queue.go +++ b/core/internal/queue/queue.go @@ -24,7 +24,6 @@ import ( "encoding/json" "errors" "fmt" - "strconv" "strings" "sync" "time" @@ -69,23 +68,19 @@ const ( // AckWait and redelivers a message that is still being processed. const heartbeatInterval = 15 * time.Second -// DLQAdvisorySequenceHeader is the NATS message header carrying the DLQ advisory sequence for a CLI-driven retry. -const DLQAdvisorySequenceHeader = "Klio-Dlq-Advisory-Sequence" - -// dlqRetrySequence returns the dead-letter queue sequence carried by the -// Klio-Dlq-Retry-Sequence header, if present. -func dlqRetrySequence(headers nats.Header) (uint64, bool) { - raw := headers.Get(DLQAdvisorySequenceHeader) - if raw == "" { - return 0, false - } - - seq, err := strconv.ParseUint(raw, 10, 64) - if err != nil { - return 0, false - } +const ( + // TaskOriginHeaderKey is the NATS message header describing the provenance + // of a task. + TaskOriginHeaderKey = "Klio-Task-Origin" + // TaskOriginDLQRetry marks a task that was manually re-enqueued from the + // dead-letter queue. + TaskOriginDLQRetry = "dlq-retry" +) - return seq, true +// isDLQRetry reports whether the message was manually re-enqueued from the +// dead-letter queue. +func isDLQRetry(headers nats.Header) bool { + return headers.Get(TaskOriginHeaderKey) == TaskOriginDLQRetry } func backupSubject(clusterName string) string { @@ -355,29 +350,6 @@ type Status struct { PendingWALs uint64 } -// notifyMessage is called to send a message on the queue. -func (q *Conn) notifyMessage(ctx context.Context, subject string, task any) error { - contextLogger := log.FromContext(ctx) - contextLogger.Info("Sending message", "subject", subject, "task", task) - - js, err := jetstream.New(q.conn) - if err != nil { - return fmt.Errorf("while creating JetStream instance: %w", err) - } - - rawContent, err := json.Marshal(task) - if err != nil { - return fmt.Errorf("while marshalling task to JSON: %w", err) - } - - _, err = js.Publish(ctx, subject, rawContent) - if err != nil { - return fmt.Errorf("while pushing message to the queue: %w", err) - } - - return nil -} - // internalConsumeMessages starts consuming messages and ends // when the context is canceled. func internalConsumeMessages[T any]( diff --git a/core/internal/queue/retry_test.go b/core/internal/queue/retry_test.go new file mode 100644 index 00000000..b94ac591 --- /dev/null +++ b/core/internal/queue/retry_test.go @@ -0,0 +1,305 @@ +package queue + +import ( + "context" + "encoding/json" + "testing" + + "github.com/nats-io/nats.go" + "github.com/nats-io/nats.go/jetstream" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// seedFailedWAL publishes an original WAL task to the WAL work-queue stream and +// a matching dead-letter queue advisory, simulating a WAL that has exhausted +// its delivery budget. +func seedFailedWAL(t *testing.T, js jetstream.JetStream, clusterName, walName string) { + t.Helper() + + data, err := json.Marshal(WALTask{ClusterName: clusterName, WALName: walName}) + require.NoError(t, err) + + ack, err := js.PublishMsg(t.Context(), &nats.Msg{Subject: walSubject(clusterName), Data: data}) + require.NoError(t, err) + + seedDLQAdvisory(t, js, klioWalStreamName, klioWalConsumerName, ack.Sequence) +} + +// retriedWALs returns the set of WAL tasks re-enqueued onto the WAL work-queue +// stream, identified by the DLQ retry origin marker. +func retriedWALs(t *testing.T, stream jetstream.Stream) map[WALTask]struct{} { + t.Helper() + + info, err := stream.Info(t.Context()) + require.NoError(t, err) + + out := make(map[WALTask]struct{}) + for seq := info.State.FirstSeq; seq <= info.State.LastSeq && seq != 0; seq++ { + msg, err := stream.GetMsg(t.Context(), seq) + if err != nil { + // Sequences may be absent (e.g. deleted); skip them. + continue + } + if msg.Header.Get(TaskOriginHeaderKey) != TaskOriginDLQRetry { + continue + } + + var task WALTask + require.NoError(t, json.Unmarshal(msg.Data, &task)) + out[task] = struct{}{} + } + + return out +} + +func TestRetryFailedWALTasksRetriesAllClusters(t *testing.T) { + ns, url := startNATSServer(t) + defer ns.Shutdown() + + nc, err := nats.Connect(url) + require.NoError(t, err) + defer nc.Close() + + ctx := context.Background() + conn, err := New(ctx, nc) + require.NoError(t, err) + + js, err := jetstream.New(nc) + require.NoError(t, err) + + seedFailedWAL(t, js, "cluster-a", "000000010000000000000001") + seedFailedWAL(t, js, "cluster-b", "000000010000000000000002") + + require.NoError(t, conn.RetryFailedWALTasks(ctx)) + + retried := retriedWALs(t, streamHandle(ctx, t, conn.conn, klioWalStreamName)) + assert.Equal(t, map[WALTask]struct{}{ + {ClusterName: "cluster-a", WALName: "000000010000000000000001"}: {}, + {ClusterName: "cluster-b", WALName: "000000010000000000000002"}: {}, + }, retried) +} + +func TestRetryFailedWALTasksRetriesSingleCluster(t *testing.T) { + ns, url := startNATSServer(t) + defer ns.Shutdown() + + nc, err := nats.Connect(url) + require.NoError(t, err) + defer nc.Close() + + ctx := context.Background() + conn, err := New(ctx, nc) + require.NoError(t, err) + + js, err := jetstream.New(nc) + require.NoError(t, err) + + seedFailedWAL(t, js, "cluster-a", "000000010000000000000001") + seedFailedWAL(t, js, "cluster-b", "000000010000000000000002") + + require.NoError(t, conn.RetryFailedWALTasks(ctx, WithCluster("cluster-a"))) + + retried := retriedWALs(t, streamHandle(ctx, t, conn.conn, klioWalStreamName)) + assert.Equal(t, map[WALTask]struct{}{ + {ClusterName: "cluster-a", WALName: "000000010000000000000001"}: {}, + }, retried, "only the requested cluster's failed WAL must be retried") +} + +func TestRetryFailedWALTasksRetriesSpecificWALs(t *testing.T) { + ns, url := startNATSServer(t) + defer ns.Shutdown() + + nc, err := nats.Connect(url) + require.NoError(t, err) + defer nc.Close() + + ctx := context.Background() + conn, err := New(ctx, nc) + require.NoError(t, err) + + js, err := jetstream.New(nc) + require.NoError(t, err) + + seedFailedWAL(t, js, "cluster-a", "000000010000000000000001") + seedFailedWAL(t, js, "cluster-a", "000000010000000000000002") + seedFailedWAL(t, js, "cluster-a", "000000010000000000000003") + + require.NoError(t, conn.RetryFailedWALTasks( + ctx, + WithCluster("cluster-a"), + WithWALs("000000010000000000000001", "000000010000000000000003"), + )) + + retried := retriedWALs(t, streamHandle(ctx, t, conn.conn, klioWalStreamName)) + assert.Equal(t, map[WALTask]struct{}{ + {ClusterName: "cluster-a", WALName: "000000010000000000000001"}: {}, + {ClusterName: "cluster-a", WALName: "000000010000000000000003"}: {}, + }, retried, "only the requested WAL names must be retried") +} + +func TestRetryFailedWALTasksSkipsUnknownWALs(t *testing.T) { + ns, url := startNATSServer(t) + defer ns.Shutdown() + + nc, err := nats.Connect(url) + require.NoError(t, err) + defer nc.Close() + + ctx := context.Background() + conn, err := New(ctx, nc) + require.NoError(t, err) + + js, err := jetstream.New(nc) + require.NoError(t, err) + + seedFailedWAL(t, js, "cluster-a", "000000010000000000000001") + + // An unknown WAL name is silently ignored; the known one is still retried. + require.NoError(t, conn.RetryFailedWALTasks( + ctx, + WithCluster("cluster-a"), + WithWALs("000000010000000000000001", "000000019999999999999999"), + )) + + retried := retriedWALs(t, streamHandle(ctx, t, conn.conn, klioWalStreamName)) + assert.Equal(t, map[WALTask]struct{}{ + {ClusterName: "cluster-a", WALName: "000000010000000000000001"}: {}, + }, retried, "only the WAL names that matched a failed task must be retried") +} + +// seedFailedBackup publishes an original backup task to the backup work-queue +// stream and a matching dead-letter queue advisory, simulating a backup that +// has exhausted its delivery budget. +func seedFailedBackup(t *testing.T, js jetstream.JetStream, clusterName string) { + t.Helper() + + seq := publishBackupMessage(t, js, clusterName) + seedDLQAdvisory(t, js, klioBackupStreamName, klioBackupConsumerName, seq) +} + +// retriedBackupClusters returns, per cluster, the number of backup tasks +// re-enqueued onto the backup work-queue stream, identified by the DLQ retry +// origin marker. +func retriedBackupClusters(t *testing.T, stream jetstream.Stream) map[string]int { + t.Helper() + + info, err := stream.Info(t.Context()) + require.NoError(t, err) + + out := make(map[string]int) + for seq := info.State.FirstSeq; seq <= info.State.LastSeq && seq != 0; seq++ { + msg, err := stream.GetMsg(t.Context(), seq) + if err != nil { + // Sequences may be absent (e.g. deleted); skip them. + continue + } + if msg.Header.Get(TaskOriginHeaderKey) != TaskOriginDLQRetry { + continue + } + + var task BackupTask + require.NoError(t, json.Unmarshal(msg.Data, &task)) + out[task.ClusterName]++ + } + + return out +} + +func TestRetryFailedBackupTasksRetriesAllClusters(t *testing.T) { + ns, url := startNATSServer(t) + defer ns.Shutdown() + + nc, err := nats.Connect(url) + require.NoError(t, err) + defer nc.Close() + + ctx := context.Background() + conn, err := New(ctx, nc) + require.NoError(t, err) + + js, err := jetstream.New(nc) + require.NoError(t, err) + + seedFailedBackup(t, js, "cluster-a") + seedFailedBackup(t, js, "cluster-b") + + require.NoError(t, conn.RetryFailedBackupTasks(ctx)) + + retried := retriedBackupClusters(t, streamHandle(ctx, t, conn.conn, klioBackupStreamName)) + assert.Equal(t, map[string]int{"cluster-a": 1, "cluster-b": 1}, retried) +} + +func TestRetryFailedBackupTasksRetriesSingleCluster(t *testing.T) { + ns, url := startNATSServer(t) + defer ns.Shutdown() + + nc, err := nats.Connect(url) + require.NoError(t, err) + defer nc.Close() + + ctx := context.Background() + conn, err := New(ctx, nc) + require.NoError(t, err) + + js, err := jetstream.New(nc) + require.NoError(t, err) + + seedFailedBackup(t, js, "cluster-a") + seedFailedBackup(t, js, "cluster-b") + + require.NoError(t, conn.RetryFailedBackupTasks(ctx, WithCluster("cluster-a"))) + + retried := retriedBackupClusters(t, streamHandle(ctx, t, conn.conn, klioBackupStreamName)) + assert.Equal(t, map[string]int{"cluster-a": 1}, retried, + "only the requested cluster's failed backup must be retried") +} + +func TestRetryFailedBackupTasksDeduplicatesByCluster(t *testing.T) { + ns, url := startNATSServer(t) + defer ns.Shutdown() + + nc, err := nats.Connect(url) + require.NoError(t, err) + defer nc.Close() + + ctx := context.Background() + conn, err := New(ctx, nc) + require.NoError(t, err) + + js, err := jetstream.New(nc) + require.NoError(t, err) + + // Two failed backups for the same cluster must collapse into a single retry. + seedFailedBackup(t, js, "cluster-a") + seedFailedBackup(t, js, "cluster-a") + + require.NoError(t, conn.RetryFailedBackupTasks(ctx)) + + retried := retriedBackupClusters(t, streamHandle(ctx, t, conn.conn, klioBackupStreamName)) + assert.Equal(t, map[string]int{"cluster-a": 1}, retried, + "multiple failed backups for one cluster must be retried only once") +} + +func TestRetryFailedBackupTasksRejectsWALFilter(t *testing.T) { + ns, url := startNATSServer(t) + defer ns.Shutdown() + + nc, err := nats.Connect(url) + require.NoError(t, err) + defer nc.Close() + + ctx := context.Background() + conn, err := New(ctx, nc) + require.NoError(t, err) + + js, err := jetstream.New(nc) + require.NoError(t, err) + + seedFailedBackup(t, js, "cluster-a") + + // BackupTask has no individual WAL name to filter on: WithWALs must be + // rejected rather than silently ignored. + err = conn.RetryFailedBackupTasks(ctx, WithWALs("000000010000000000000001")) + require.ErrorIs(t, err, errWALFilterUnsupported) +} diff --git a/core/internal/queue/wal.go b/core/internal/queue/wal.go index a58d9224..dc82956b 100644 --- a/core/internal/queue/wal.go +++ b/core/internal/queue/wal.go @@ -50,7 +50,7 @@ func (t WALTask) Cluster() string { // NotifyWALReceived is called to notify the consumers that a new WAL // is available in the Klio repository. func (q *Conn) NotifyWALReceived(ctx context.Context, task *WALTask) error { - return q.notifyMessage(ctx, walSubject(task.ClusterName), task) + return q.notifyMessage(ctx, walSubject(task.ClusterName), task, nil) } // WALTaskHandler is called for every WAL task message that should be handled. @@ -61,20 +61,27 @@ type WALTaskHandler func(ctx context.Context, t *WALTask) error // when the context is canceled. After a successful handler run, the WAL is // recorded as the latest uploaded WAL for its cluster. func (q *Conn) ConsumeWALReceivedMessages(ctx context.Context, handler WALTaskHandler) error { + logger := log.FromContext(ctx).WithName("wal-consumer") + wrapped := func(ctx context.Context, t *WALTask, headers nats.Header) error { + isRetried := isDLQRetry(headers) + if isRetried { + logger.Info( + "Retrying WAL task re-enqueued from the dead-letter queue", + "cluster", t.ClusterName, "wal", t.WALName, + ) + } + if err := handler(ctx, t); err != nil { return err } - // retried messages carry the marker identifying the exact dead-letter queue - // entry to remove. - if dlqSequence, ok := dlqRetrySequence(headers); ok { - if err := q.purgeWALDLQEntry(ctx, dlqSequence); err != nil { - log.FromContext(ctx).Error( + if isRetried { + if err := q.purgeWALDLQEntries(ctx, t.ClusterName, t.WALName); err != nil { + logger.Error( err, - "Failed to purge WAL dead-letter queue entry after successful retry", - "task", t, - "dlqSequence", dlqSequence, + "Failed to purge WAL dead-letter queue entries after successful retry", + "cluster", t.ClusterName, "wal", t.WALName, ) } } @@ -83,10 +90,14 @@ func (q *Conn) ConsumeWALReceivedMessages(ctx context.Context, handler WALTaskHa return nil } + // TODO: now that we allow holes in the WAL sequence, the logic of gating the deletion of WALs based on the + // latest uploaded WAL is not correct anymore. If WAL X and X+1 fail to upload, but X+2 is uploaded, the + // post-backup clean-up logic may delete WAL X and X+1 even if they have not been uploaded to tier2 yet. if err := q.notifyMessage( ctx, latestUploadedWalSubject(t.ClusterName), t, + nil, ); err != nil { log.FromContext(ctx).Error( err, diff --git a/core/internal/server/admin/admin.go b/core/internal/server/admin/admin.go index 3d01ef29..018de5fc 100644 --- a/core/internal/server/admin/admin.go +++ b/core/internal/server/admin/admin.go @@ -245,7 +245,7 @@ func (s *Server) QueueListFailedBackups( ) } - opts := make([]queue.ListOption, 0) + opts := make([]queue.Option, 0) if name := req.GetClusterName(); name != "" { opts = append(opts, queue.WithCluster(name)) @@ -281,7 +281,7 @@ func (s *Server) QueueListFailedWALs( "failed WALs not available: server not configured with Stream Manager", ) } - opts := make([]queue.ListOption, 0) + opts := make([]queue.Option, 0) if name := req.GetClusterName(); name != "" { opts = append(opts, queue.WithCluster(name)) @@ -308,6 +308,66 @@ func (s *Server) QueueListFailedWALs( }, nil } +// QueueRetryWALs implements [grpc.AdminServer]. +func (s *Server) QueueRetryWALs( + ctx context.Context, + req *klioGRPC.QueueRetryWALsRequest, +) (*klioGRPC.QueueRetryResponse, error) { + if s.streamMgr == nil { + return nil, status.Errorf( + codes.Unavailable, + "failed WALs not available: server not configured with Stream Manager", + ) + } + + clusterName := req.GetClusterName() + wals := req.GetWalNames() + + if len(wals) > 0 && clusterName == "" { + return nil, status.Errorf(codes.InvalidArgument, "WAL names require a cluster name") + } + + var retryOpts []queue.Option + if clusterName != "" { + retryOpts = append(retryOpts, queue.WithCluster(clusterName)) + } + if len(wals) > 0 { + retryOpts = append(retryOpts, queue.WithWALs(wals...)) + } + + if err := s.streamMgr.RetryFailedWALTasks(ctx, retryOpts...); err != nil { + return nil, status.Errorf(codes.Internal, "while retrying failed WALs: %s", err.Error()) + } + + return &klioGRPC.QueueRetryResponse{}, nil +} + +// QueueRetryBackups implements [grpc.AdminServer]. +func (s *Server) QueueRetryBackups( + ctx context.Context, + req *klioGRPC.QueueRetryBackupsRequest, +) (*klioGRPC.QueueRetryResponse, error) { + if s.streamMgr == nil { + return nil, status.Errorf( + codes.Unavailable, + "failed backups not available: server not configured with Stream Manager", + ) + } + + clusterName := req.GetClusterName() + + var retryOpts []queue.Option + if clusterName != "" { + retryOpts = append(retryOpts, queue.WithCluster(clusterName)) + } + + if err := s.streamMgr.RetryFailedBackupTasks(ctx, retryOpts...); err != nil { + return nil, status.Errorf(codes.Internal, "while retrying failed backups: %s", err.Error()) + } + + return &klioGRPC.QueueRetryResponse{}, nil +} + // DeleteBackup implements [grpc.AdminServer]. func (s *Server) DeleteBackup( ctx context.Context, diff --git a/core/proto/klio_admin.proto b/core/proto/klio_admin.proto index f29aa24d..0097e065 100644 --- a/core/proto/klio_admin.proto +++ b/core/proto/klio_admin.proto @@ -37,6 +37,12 @@ service Admin { // List WAL files failed to be processed from the queue rpc QueueListFailedWALs(QueueListFailedWALsRequest) returns (QueueListFailedWALsResponse) {} + // Retry WAL files that failed to be processed from the queue + rpc QueueRetryWALs(QueueRetryWALsRequest) returns (QueueRetryResponse) {} + + // Retry Backups that failed to be processed from the queue + rpc QueueRetryBackups(QueueRetryBackupsRequest) returns (QueueRetryResponse) {} + // Get the status of the task queue (pending backups and WALs) rpc QueueStatus(QueueStatusRequest) returns (QueueStatusResponse) {} @@ -94,6 +100,18 @@ message FailedWAL { google.protobuf.Timestamp last_attempt_time = 4; } +message QueueRetryWALsRequest { + optional string cluster_name = 1; + repeated string wal_names = 2; +} + +message QueueRetryBackupsRequest { + optional string cluster_name = 1; +} + +message QueueRetryResponse { +} + message QueueStatusRequest { } diff --git a/documentation/.wordlist.txt b/documentation/.wordlist.txt index 8d9adc9f..48c43b2b 100644 --- a/documentation/.wordlist.txt +++ b/documentation/.wordlist.txt @@ -135,6 +135,11 @@ QueueListFailedBackupsResponse QueueListFailedWALs QueueListFailedWALsRequest QueueListFailedWALsResponse +QueueRetryBackups +QueueRetryBackupsRequest +QueueRetryResponse +QueueRetryWALs +QueueRetryWALsRequest QueueStatus QueueStatusRequest QueueStatusResponse diff --git a/documentation/web/docs/developer/_protocol.md b/documentation/web/docs/developer/_protocol.md index 4a1665e2..3a3d1c2a 100644 --- a/documentation/web/docs/developer/_protocol.md +++ b/documentation/web/docs/developer/_protocol.md @@ -14,6 +14,9 @@ - [QueueListFailedBackupsResponse](#klio-wal-v1-QueueListFailedBackupsResponse) - [QueueListFailedWALsRequest](#klio-wal-v1-QueueListFailedWALsRequest) - [QueueListFailedWALsResponse](#klio-wal-v1-QueueListFailedWALsResponse) + - [QueueRetryBackupsRequest](#klio-wal-v1-QueueRetryBackupsRequest) + - [QueueRetryResponse](#klio-wal-v1-QueueRetryResponse) + - [QueueRetryWALsRequest](#klio-wal-v1-QueueRetryWALsRequest) - [QueueStatusRequest](#klio-wal-v1-QueueStatusRequest) - [QueueStatusResponse](#klio-wal-v1-QueueStatusResponse) - [RefreshRequest](#klio-wal-v1-RefreshRequest) @@ -198,6 +201,47 @@ DeleteBackupResponse is the response to a backup deletion request. + + +### QueueRetryBackupsRequest + + + +| Field | Type | Label | Description | +| ----- | ---- | ----- | ----------- | +| cluster_name | [string](#string) | optional | | + + + + + + + + +### QueueRetryResponse + + + + + + + + + +### QueueRetryWALsRequest + + + +| Field | Type | Label | Description | +| ----- | ---- | ----- | ----------- | +| cluster_name | [string](#string) | optional | | +| wal_names | [string](#string) | repeated | | + + + + + + ### QueueStatusRequest @@ -274,6 +318,8 @@ Tier represents a storage tier in the backup system. | ListBackups | [ListBackupsRequest](#klio-wal-v1-ListBackupsRequest) | [ListBackupsResult](#klio-wal-v1-ListBackupsResult) | List every backup on the server | | QueueListFailedBackups | [QueueListFailedBackupsRequest](#klio-wal-v1-QueueListFailedBackupsRequest) | [QueueListFailedBackupsResponse](#klio-wal-v1-QueueListFailedBackupsResponse) | List backups failed to be processed from the queue | | QueueListFailedWALs | [QueueListFailedWALsRequest](#klio-wal-v1-QueueListFailedWALsRequest) | [QueueListFailedWALsResponse](#klio-wal-v1-QueueListFailedWALsResponse) | List WAL files failed to be processed from the queue | +| QueueRetryWALs | [QueueRetryWALsRequest](#klio-wal-v1-QueueRetryWALsRequest) | [QueueRetryResponse](#klio-wal-v1-QueueRetryResponse) | Retry WAL files that failed to be processed from the queue | +| QueueRetryBackups | [QueueRetryBackupsRequest](#klio-wal-v1-QueueRetryBackupsRequest) | [QueueRetryResponse](#klio-wal-v1-QueueRetryResponse) | Retry Backups that failed to be processed from the queue | | QueueStatus | [QueueStatusRequest](#klio-wal-v1-QueueStatusRequest) | [QueueStatusResponse](#klio-wal-v1-QueueStatusResponse) | Get the status of the task queue (pending backups and WALs) | | DeleteBackup | [DeleteBackupRequest](#klio-wal-v1-DeleteBackupRequest) | [DeleteBackupResponse](#klio-wal-v1-DeleteBackupResponse) | Delete a backup from the server | diff --git a/documentation/web/docs/user/cli/klio_admin_queue_backup.md b/documentation/web/docs/user/cli/klio_admin_queue_backup.md index a1937143..a2018527 100644 --- a/documentation/web/docs/user/cli/klio_admin_queue_backup.md +++ b/documentation/web/docs/user/cli/klio_admin_queue_backup.md @@ -35,4 +35,5 @@ Manage the queue backup tasks * [klio admin queue](klio_admin_queue.md) - Manage the queue tasks * [klio admin queue backup list-failed](klio_admin_queue_backup_list-failed.md) - List failed backup tasks in the queue +* [klio admin queue backup retry](klio_admin_queue_backup_retry.md) - Retry failed backup tasks in the queue diff --git a/documentation/web/docs/user/cli/klio_admin_queue_backup_retry.md b/documentation/web/docs/user/cli/klio_admin_queue_backup_retry.md new file mode 100644 index 00000000..ed5e8cb8 --- /dev/null +++ b/documentation/web/docs/user/cli/klio_admin_queue_backup_retry.md @@ -0,0 +1,47 @@ +--- +title: klio admin queue backup retry +--- + +## klio admin queue backup retry + +Retry failed backup tasks in the queue + +### Synopsis + +Retry failed backup tasks in the queue. + +With no arguments, all failed backup tasks are retried. If a cluster name is given, all failed backup tasks for that cluster are retried. + +``` +klio admin queue backup retry [cluster-name] [flags] +``` + +### Options + +``` + -h, --help help for retry +``` + +### Options inherited from parent commands + +``` + --config string config file (default is $HOME/.klio.yaml) + --debug enable debug logging + --json Output in JSON format + --log-destination string where the log stream will be written + --log-field-level string JSON log field to report severity in (default: level) + --log-field-timestamp string JSON log field to report timestamp in (default: ts) + --log-level string the desired log level, one of error, info, debug and trace (default "info") + --pprof-server string enable the PPROF server using the specified address + --socket-path string Unix socket used by the administration server (default "/tmp/.klio-admin") + --zap-devel Development Mode defaults(encoder=consoleEncoder,logLevel=Debug,stackTraceLevel=Warn). Production Mode defaults(encoder=jsonEncoder,logLevel=Info,stackTraceLevel=Error) + --zap-encoder encoder Zap log encoding (one of 'json' or 'console') + --zap-log-level level Zap Level to configure the verbosity of logging. Can be one of 'debug', 'info', 'error', 'panic' or any integer value > 0 which corresponds to custom debug levels of increasing verbosity + --zap-stacktrace-level level Zap Level at and above which stacktraces are captured (one of 'info', 'error', 'panic'). + --zap-time-encoding time-encoding Zap time encoding (one of 'epoch', 'millis', 'nano', 'iso8601', 'rfc3339' or 'rfc3339nano'). Defaults to 'epoch'. +``` + +### SEE ALSO + +* [klio admin queue backup](klio_admin_queue_backup.md) - Manage the queue backup tasks + diff --git a/documentation/web/docs/user/cli/klio_admin_queue_wal.md b/documentation/web/docs/user/cli/klio_admin_queue_wal.md index 562a303b..4575f86d 100644 --- a/documentation/web/docs/user/cli/klio_admin_queue_wal.md +++ b/documentation/web/docs/user/cli/klio_admin_queue_wal.md @@ -35,4 +35,5 @@ Manage the queue WAL tasks * [klio admin queue](klio_admin_queue.md) - Manage the queue tasks * [klio admin queue wal list-failed](klio_admin_queue_wal_list-failed.md) - List failed WAL tasks in the queue +* [klio admin queue wal retry](klio_admin_queue_wal_retry.md) - Retry failed WAL tasks in the queue diff --git a/documentation/web/docs/user/cli/klio_admin_queue_wal_retry.md b/documentation/web/docs/user/cli/klio_admin_queue_wal_retry.md new file mode 100644 index 00000000..2287ccfa --- /dev/null +++ b/documentation/web/docs/user/cli/klio_admin_queue_wal_retry.md @@ -0,0 +1,47 @@ +--- +title: klio admin queue wal retry +--- + +## klio admin queue wal retry + +Retry failed WAL tasks in the queue + +### Synopsis + +Retry failed WAL tasks in the queue. + +With no arguments, all failed WAL tasks are retried. If a cluster name is given, all failed WAL tasks for that cluster are retried. If WAL files are also given, only those are retried. + +``` +klio admin queue wal retry [cluster-name] [WAL1 WAL2 ...] [flags] +``` + +### Options + +``` + -h, --help help for retry +``` + +### Options inherited from parent commands + +``` + --config string config file (default is $HOME/.klio.yaml) + --debug enable debug logging + --json Output in JSON format + --log-destination string where the log stream will be written + --log-field-level string JSON log field to report severity in (default: level) + --log-field-timestamp string JSON log field to report timestamp in (default: ts) + --log-level string the desired log level, one of error, info, debug and trace (default "info") + --pprof-server string enable the PPROF server using the specified address + --socket-path string Unix socket used by the administration server (default "/tmp/.klio-admin") + --zap-devel Development Mode defaults(encoder=consoleEncoder,logLevel=Debug,stackTraceLevel=Warn). Production Mode defaults(encoder=jsonEncoder,logLevel=Info,stackTraceLevel=Error) + --zap-encoder encoder Zap log encoding (one of 'json' or 'console') + --zap-log-level level Zap Level to configure the verbosity of logging. Can be one of 'debug', 'info', 'error', 'panic' or any integer value > 0 which corresponds to custom debug levels of increasing verbosity + --zap-stacktrace-level level Zap Level at and above which stacktraces are captured (one of 'info', 'error', 'panic'). + --zap-time-encoding time-encoding Zap time encoding (one of 'epoch', 'millis', 'nano', 'iso8601', 'rfc3339' or 'rfc3339nano'). Defaults to 'epoch'. +``` + +### SEE ALSO + +* [klio admin queue wal](klio_admin_queue_wal.md) - Manage the queue WAL tasks +