From 8f2adf2138f916b6d58bcd2f0ca7398a87d5cf68 Mon Sep 17 00:00:00 2001 From: Andre Duffeck Date: Mon, 5 Oct 2026 14:28:58 +0200 Subject: [PATCH] Reindex disabled spaces once they get enabled again (#3579) Spaces that are disabled while reindexing are skipped. They are now remembered in a NATS KV bucket (including whether a forced rescan was requested) and reindexed when the SpaceEnabled event comes in. The entry is removed when the space gets deleted. --- .../opencloud/services/search/v0/search.pb.go | 241 ++++++++++++------ .../services/search/v0/search.swagger.json | 43 +++- .../opencloud/services/search/v0/search.proto | 13 + services/search/pkg/command/index.go | 37 ++- services/search/pkg/command/server.go | 51 +++- services/search/pkg/search/skippedspaces.go | 136 ++++++++++ .../search/pkg/search/skippedspaces_test.go | 138 ++++++++++ services/search/pkg/server/grpc/option.go | 8 + services/search/pkg/server/grpc/server.go | 1 + .../search/pkg/service/event/debouncer.go | 17 +- .../pkg/service/event/debouncer_test.go | 88 +++++-- services/search/pkg/service/event/service.go | 69 +++-- .../search/pkg/service/event/service_test.go | 3 +- services/search/pkg/service/grpc/v0/option.go | 8 + .../search/pkg/service/grpc/v0/service.go | 59 +++-- 15 files changed, 728 insertions(+), 184 deletions(-) create mode 100644 services/search/pkg/search/skippedspaces.go create mode 100644 services/search/pkg/search/skippedspaces_test.go diff --git a/protogen/gen/opencloud/services/search/v0/search.pb.go b/protogen/gen/opencloud/services/search/v0/search.pb.go index 1346c98197..82489194fa 100644 --- a/protogen/gen/opencloud/services/search/v0/search.pb.go +++ b/protogen/gen/opencloud/services/search/v0/search.pb.go @@ -25,6 +25,62 @@ const ( _ = protoimpl.EnforceVersion(protoimpl.MaxVersion - 20) ) +type IndexSpaceResponse_Status int32 + +const ( + IndexSpaceResponse_STATUS_UNKNOWN IndexSpaceResponse_Status = 0 + // The space has been indexed successfully. + IndexSpaceResponse_STATUS_SUCCESS IndexSpaceResponse_Status = 1 + // The space has been skipped because it is disabled. It will be indexed + // once it gets enabled again. + IndexSpaceResponse_STATUS_SKIPPED IndexSpaceResponse_Status = 2 + // Indexing the space failed, see the error field for details. + IndexSpaceResponse_STATUS_ERROR IndexSpaceResponse_Status = 3 +) + +// Enum value maps for IndexSpaceResponse_Status. +var ( + IndexSpaceResponse_Status_name = map[int32]string{ + 0: "STATUS_UNKNOWN", + 1: "STATUS_SUCCESS", + 2: "STATUS_SKIPPED", + 3: "STATUS_ERROR", + } + IndexSpaceResponse_Status_value = map[string]int32{ + "STATUS_UNKNOWN": 0, + "STATUS_SUCCESS": 1, + "STATUS_SKIPPED": 2, + "STATUS_ERROR": 3, + } +) + +func (x IndexSpaceResponse_Status) Enum() *IndexSpaceResponse_Status { + p := new(IndexSpaceResponse_Status) + *p = x + return p +} + +func (x IndexSpaceResponse_Status) String() string { + return protoimpl.X.EnumStringOf(x.Descriptor(), protoreflect.EnumNumber(x)) +} + +func (IndexSpaceResponse_Status) Descriptor() protoreflect.EnumDescriptor { + return file_opencloud_services_search_v0_search_proto_enumTypes[0].Descriptor() +} + +func (IndexSpaceResponse_Status) Type() protoreflect.EnumType { + return &file_opencloud_services_search_v0_search_proto_enumTypes[0] +} + +func (x IndexSpaceResponse_Status) Number() protoreflect.EnumNumber { + return protoreflect.EnumNumber(x) +} + +// Deprecated: Use IndexSpaceResponse_Status.Descriptor instead. +func (IndexSpaceResponse_Status) EnumDescriptor() ([]byte, []int) { + return file_opencloud_services_search_v0_search_proto_rawDescGZIP(), []int{5, 0} +} + type SearchRequest struct { state protoimpl.MessageState sizeCache protoimpl.SizeCache @@ -389,6 +445,8 @@ type IndexSpaceResponse struct { TotalSpaces int64 `protobuf:"varint,4,opt,name=total_spaces,json=totalSpaces,proto3" json:"total_spaces,omitempty"` // Contains an error message in case indexing this particular space failed. Error string `protobuf:"bytes,5,opt,name=error,proto3" json:"error,omitempty"` + // The outcome of indexing this particular space. + Status IndexSpaceResponse_Status `protobuf:"varint,6,opt,name=status,proto3,enum=opencloud.services.search.v0.IndexSpaceResponse_Status" json:"status,omitempty"` } func (x *IndexSpaceResponse) Reset() { @@ -458,6 +516,13 @@ func (x *IndexSpaceResponse) GetError() string { return "" } +func (x *IndexSpaceResponse) GetStatus() IndexSpaceResponse_Status { + if x != nil { + return x.Status + } + return IndexSpaceResponse_STATUS_UNKNOWN +} + var File_opencloud_services_search_v0_search_proto protoreflect.FileDescriptor var file_opencloud_services_search_v0_search_proto_rawDesc = []byte{ @@ -531,7 +596,7 @@ var file_opencloud_services_search_v0_search_proto_rawDesc = []byte{ 0x52, 0x0c, 0x66, 0x6f, 0x72, 0x63, 0x65, 0x52, 0x65, 0x69, 0x6e, 0x64, 0x65, 0x78, 0x12, 0x26, 0x0a, 0x0b, 0x63, 0x6f, 0x6e, 0x63, 0x75, 0x72, 0x72, 0x65, 0x6e, 0x63, 0x79, 0x18, 0x04, 0x20, 0x01, 0x28, 0x05, 0x42, 0x04, 0xe2, 0x41, 0x01, 0x01, 0x52, 0x0b, 0x63, 0x6f, 0x6e, 0x63, 0x75, - 0x72, 0x72, 0x65, 0x6e, 0x63, 0x79, 0x22, 0xd1, 0x01, 0x0a, 0x12, 0x49, 0x6e, 0x64, 0x65, 0x78, + 0x72, 0x72, 0x65, 0x6e, 0x63, 0x79, 0x22, 0xfa, 0x02, 0x0a, 0x12, 0x49, 0x6e, 0x64, 0x65, 0x78, 0x53, 0x70, 0x61, 0x63, 0x65, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x12, 0x19, 0x0a, 0x08, 0x73, 0x70, 0x61, 0x63, 0x65, 0x5f, 0x69, 0x64, 0x18, 0x01, 0x20, 0x01, 0x28, 0x09, 0x52, 0x07, 0x73, 0x70, 0x61, 0x63, 0x65, 0x49, 0x64, 0x12, 0x40, 0x0a, 0x0e, 0x73, 0x70, 0x61, 0x63, @@ -544,61 +609,71 @@ var file_opencloud_services_search_v0_search_proto_rawDesc = []byte{ 0x73, 0x12, 0x21, 0x0a, 0x0c, 0x74, 0x6f, 0x74, 0x61, 0x6c, 0x5f, 0x73, 0x70, 0x61, 0x63, 0x65, 0x73, 0x18, 0x04, 0x20, 0x01, 0x28, 0x03, 0x52, 0x0b, 0x74, 0x6f, 0x74, 0x61, 0x6c, 0x53, 0x70, 0x61, 0x63, 0x65, 0x73, 0x12, 0x14, 0x0a, 0x05, 0x65, 0x72, 0x72, 0x6f, 0x72, 0x18, 0x05, 0x20, - 0x01, 0x28, 0x09, 0x52, 0x05, 0x65, 0x72, 0x72, 0x6f, 0x72, 0x32, 0xb3, 0x02, 0x0a, 0x0e, 0x53, - 0x65, 0x61, 0x72, 0x63, 0x68, 0x50, 0x72, 0x6f, 0x76, 0x69, 0x64, 0x65, 0x72, 0x12, 0x85, 0x01, - 0x0a, 0x06, 0x53, 0x65, 0x61, 0x72, 0x63, 0x68, 0x12, 0x2b, 0x2e, 0x6f, 0x70, 0x65, 0x6e, 0x63, - 0x6c, 0x6f, 0x75, 0x64, 0x2e, 0x73, 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x73, 0x2e, 0x73, 0x65, - 0x61, 0x72, 0x63, 0x68, 0x2e, 0x76, 0x30, 0x2e, 0x53, 0x65, 0x61, 0x72, 0x63, 0x68, 0x52, 0x65, - 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, 0x2c, 0x2e, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, - 0x64, 0x2e, 0x73, 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x73, 0x2e, 0x73, 0x65, 0x61, 0x72, 0x63, - 0x68, 0x2e, 0x76, 0x30, 0x2e, 0x53, 0x65, 0x61, 0x72, 0x63, 0x68, 0x52, 0x65, 0x73, 0x70, 0x6f, - 0x6e, 0x73, 0x65, 0x22, 0x20, 0x82, 0xd3, 0xe4, 0x93, 0x02, 0x1a, 0x3a, 0x01, 0x2a, 0x22, 0x15, - 0x2f, 0x61, 0x70, 0x69, 0x2f, 0x76, 0x30, 0x2f, 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x2f, 0x73, - 0x65, 0x61, 0x72, 0x63, 0x68, 0x12, 0x98, 0x01, 0x0a, 0x0a, 0x49, 0x6e, 0x64, 0x65, 0x78, 0x53, - 0x70, 0x61, 0x63, 0x65, 0x12, 0x2f, 0x2e, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, - 0x2e, 0x73, 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x73, 0x2e, 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, - 0x2e, 0x76, 0x30, 0x2e, 0x49, 0x6e, 0x64, 0x65, 0x78, 0x53, 0x70, 0x61, 0x63, 0x65, 0x52, 0x65, - 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, 0x30, 0x2e, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, - 0x64, 0x2e, 0x73, 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x73, 0x2e, 0x73, 0x65, 0x61, 0x72, 0x63, - 0x68, 0x2e, 0x76, 0x30, 0x2e, 0x49, 0x6e, 0x64, 0x65, 0x78, 0x53, 0x70, 0x61, 0x63, 0x65, 0x52, - 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x22, 0x25, 0x82, 0xd3, 0xe4, 0x93, 0x02, 0x1f, 0x3a, - 0x01, 0x2a, 0x22, 0x1a, 0x2f, 0x61, 0x70, 0x69, 0x2f, 0x76, 0x30, 0x2f, 0x73, 0x65, 0x61, 0x72, - 0x63, 0x68, 0x2f, 0x69, 0x6e, 0x64, 0x65, 0x78, 0x2d, 0x73, 0x70, 0x61, 0x63, 0x65, 0x30, 0x01, - 0x32, 0xa7, 0x01, 0x0a, 0x0d, 0x49, 0x6e, 0x64, 0x65, 0x78, 0x50, 0x72, 0x6f, 0x76, 0x69, 0x64, - 0x65, 0x72, 0x12, 0x95, 0x01, 0x0a, 0x06, 0x53, 0x65, 0x61, 0x72, 0x63, 0x68, 0x12, 0x30, 0x2e, + 0x01, 0x28, 0x09, 0x52, 0x05, 0x65, 0x72, 0x72, 0x6f, 0x72, 0x12, 0x4f, 0x0a, 0x06, 0x73, 0x74, + 0x61, 0x74, 0x75, 0x73, 0x18, 0x06, 0x20, 0x01, 0x28, 0x0e, 0x32, 0x37, 0x2e, 0x6f, 0x70, 0x65, + 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2e, 0x73, 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x73, 0x2e, + 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x2e, 0x76, 0x30, 0x2e, 0x49, 0x6e, 0x64, 0x65, 0x78, 0x53, + 0x70, 0x61, 0x63, 0x65, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x2e, 0x53, 0x74, 0x61, + 0x74, 0x75, 0x73, 0x52, 0x06, 0x73, 0x74, 0x61, 0x74, 0x75, 0x73, 0x22, 0x56, 0x0a, 0x06, 0x53, + 0x74, 0x61, 0x74, 0x75, 0x73, 0x12, 0x12, 0x0a, 0x0e, 0x53, 0x54, 0x41, 0x54, 0x55, 0x53, 0x5f, + 0x55, 0x4e, 0x4b, 0x4e, 0x4f, 0x57, 0x4e, 0x10, 0x00, 0x12, 0x12, 0x0a, 0x0e, 0x53, 0x54, 0x41, + 0x54, 0x55, 0x53, 0x5f, 0x53, 0x55, 0x43, 0x43, 0x45, 0x53, 0x53, 0x10, 0x01, 0x12, 0x12, 0x0a, + 0x0e, 0x53, 0x54, 0x41, 0x54, 0x55, 0x53, 0x5f, 0x53, 0x4b, 0x49, 0x50, 0x50, 0x45, 0x44, 0x10, + 0x02, 0x12, 0x10, 0x0a, 0x0c, 0x53, 0x54, 0x41, 0x54, 0x55, 0x53, 0x5f, 0x45, 0x52, 0x52, 0x4f, + 0x52, 0x10, 0x03, 0x32, 0xb3, 0x02, 0x0a, 0x0e, 0x53, 0x65, 0x61, 0x72, 0x63, 0x68, 0x50, 0x72, + 0x6f, 0x76, 0x69, 0x64, 0x65, 0x72, 0x12, 0x85, 0x01, 0x0a, 0x06, 0x53, 0x65, 0x61, 0x72, 0x63, + 0x68, 0x12, 0x2b, 0x2e, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2e, 0x73, 0x65, + 0x72, 0x76, 0x69, 0x63, 0x65, 0x73, 0x2e, 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x2e, 0x76, 0x30, + 0x2e, 0x53, 0x65, 0x61, 0x72, 0x63, 0x68, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, 0x2c, + 0x2e, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2e, 0x73, 0x65, 0x72, 0x76, 0x69, + 0x63, 0x65, 0x73, 0x2e, 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x2e, 0x76, 0x30, 0x2e, 0x53, 0x65, + 0x61, 0x72, 0x63, 0x68, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x22, 0x20, 0x82, 0xd3, + 0xe4, 0x93, 0x02, 0x1a, 0x3a, 0x01, 0x2a, 0x22, 0x15, 0x2f, 0x61, 0x70, 0x69, 0x2f, 0x76, 0x30, + 0x2f, 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x2f, 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x12, 0x98, + 0x01, 0x0a, 0x0a, 0x49, 0x6e, 0x64, 0x65, 0x78, 0x53, 0x70, 0x61, 0x63, 0x65, 0x12, 0x2f, 0x2e, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2e, 0x73, 0x65, 0x72, 0x76, 0x69, 0x63, - 0x65, 0x73, 0x2e, 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x2e, 0x76, 0x30, 0x2e, 0x53, 0x65, 0x61, - 0x72, 0x63, 0x68, 0x49, 0x6e, 0x64, 0x65, 0x78, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, - 0x31, 0x2e, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2e, 0x73, 0x65, 0x72, 0x76, - 0x69, 0x63, 0x65, 0x73, 0x2e, 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x2e, 0x76, 0x30, 0x2e, 0x53, - 0x65, 0x61, 0x72, 0x63, 0x68, 0x49, 0x6e, 0x64, 0x65, 0x78, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, - 0x73, 0x65, 0x22, 0x26, 0x82, 0xd3, 0xe4, 0x93, 0x02, 0x20, 0x3a, 0x01, 0x2a, 0x22, 0x1b, 0x2f, - 0x61, 0x70, 0x69, 0x2f, 0x76, 0x30, 0x2f, 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x2f, 0x69, 0x6e, - 0x64, 0x65, 0x78, 0x2f, 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x42, 0xf2, 0x02, 0x5a, 0x4a, 0x67, - 0x69, 0x74, 0x68, 0x75, 0x62, 0x2e, 0x63, 0x6f, 0x6d, 0x2f, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, - 0x6f, 0x75, 0x64, 0x2d, 0x65, 0x75, 0x2f, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, - 0x2f, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x67, 0x65, 0x6e, 0x2f, 0x67, 0x65, 0x6e, 0x2f, 0x6f, 0x70, - 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2f, 0x73, 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x2f, - 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x2f, 0x76, 0x30, 0x92, 0x41, 0xa2, 0x02, 0x12, 0xb7, 0x01, - 0x0a, 0x10, 0x4f, 0x70, 0x65, 0x6e, 0x43, 0x6c, 0x6f, 0x75, 0x64, 0x20, 0x73, 0x65, 0x61, 0x72, - 0x63, 0x68, 0x22, 0x51, 0x0a, 0x0e, 0x4f, 0x70, 0x65, 0x6e, 0x43, 0x6c, 0x6f, 0x75, 0x64, 0x20, - 0x47, 0x6d, 0x62, 0x48, 0x12, 0x29, 0x68, 0x74, 0x74, 0x70, 0x73, 0x3a, 0x2f, 0x2f, 0x67, 0x69, - 0x74, 0x68, 0x75, 0x62, 0x2e, 0x63, 0x6f, 0x6d, 0x2f, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, - 0x75, 0x64, 0x2d, 0x65, 0x75, 0x2f, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x1a, - 0x14, 0x73, 0x75, 0x70, 0x70, 0x6f, 0x72, 0x74, 0x40, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, - 0x75, 0x64, 0x2e, 0x65, 0x75, 0x2a, 0x49, 0x0a, 0x0a, 0x41, 0x70, 0x61, 0x63, 0x68, 0x65, 0x2d, - 0x32, 0x2e, 0x30, 0x12, 0x3b, 0x68, 0x74, 0x74, 0x70, 0x73, 0x3a, 0x2f, 0x2f, 0x67, 0x69, 0x74, - 0x68, 0x75, 0x62, 0x2e, 0x63, 0x6f, 0x6d, 0x2f, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, - 0x64, 0x2d, 0x65, 0x75, 0x2f, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2f, 0x62, - 0x6c, 0x6f, 0x62, 0x2f, 0x6d, 0x61, 0x69, 0x6e, 0x2f, 0x4c, 0x49, 0x43, 0x45, 0x4e, 0x53, 0x45, - 0x32, 0x05, 0x31, 0x2e, 0x30, 0x2e, 0x30, 0x2a, 0x02, 0x01, 0x02, 0x32, 0x10, 0x61, 0x70, 0x70, - 0x6c, 0x69, 0x63, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x2f, 0x6a, 0x73, 0x6f, 0x6e, 0x3a, 0x10, 0x61, - 0x70, 0x70, 0x6c, 0x69, 0x63, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x2f, 0x6a, 0x73, 0x6f, 0x6e, 0x72, - 0x3e, 0x0a, 0x10, 0x44, 0x65, 0x76, 0x65, 0x6c, 0x6f, 0x70, 0x65, 0x72, 0x20, 0x4d, 0x61, 0x6e, - 0x75, 0x61, 0x6c, 0x12, 0x2a, 0x68, 0x74, 0x74, 0x70, 0x73, 0x3a, 0x2f, 0x2f, 0x64, 0x6f, 0x63, - 0x73, 0x2e, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2e, 0x65, 0x75, 0x2f, 0x73, - 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x73, 0x2f, 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x2f, 0x62, - 0x06, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x33, + 0x65, 0x73, 0x2e, 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x2e, 0x76, 0x30, 0x2e, 0x49, 0x6e, 0x64, + 0x65, 0x78, 0x53, 0x70, 0x61, 0x63, 0x65, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, 0x30, + 0x2e, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2e, 0x73, 0x65, 0x72, 0x76, 0x69, + 0x63, 0x65, 0x73, 0x2e, 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x2e, 0x76, 0x30, 0x2e, 0x49, 0x6e, + 0x64, 0x65, 0x78, 0x53, 0x70, 0x61, 0x63, 0x65, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, + 0x22, 0x25, 0x82, 0xd3, 0xe4, 0x93, 0x02, 0x1f, 0x3a, 0x01, 0x2a, 0x22, 0x1a, 0x2f, 0x61, 0x70, + 0x69, 0x2f, 0x76, 0x30, 0x2f, 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x2f, 0x69, 0x6e, 0x64, 0x65, + 0x78, 0x2d, 0x73, 0x70, 0x61, 0x63, 0x65, 0x30, 0x01, 0x32, 0xa7, 0x01, 0x0a, 0x0d, 0x49, 0x6e, + 0x64, 0x65, 0x78, 0x50, 0x72, 0x6f, 0x76, 0x69, 0x64, 0x65, 0x72, 0x12, 0x95, 0x01, 0x0a, 0x06, + 0x53, 0x65, 0x61, 0x72, 0x63, 0x68, 0x12, 0x30, 0x2e, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, + 0x75, 0x64, 0x2e, 0x73, 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x73, 0x2e, 0x73, 0x65, 0x61, 0x72, + 0x63, 0x68, 0x2e, 0x76, 0x30, 0x2e, 0x53, 0x65, 0x61, 0x72, 0x63, 0x68, 0x49, 0x6e, 0x64, 0x65, + 0x78, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, 0x31, 0x2e, 0x6f, 0x70, 0x65, 0x6e, 0x63, + 0x6c, 0x6f, 0x75, 0x64, 0x2e, 0x73, 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x73, 0x2e, 0x73, 0x65, + 0x61, 0x72, 0x63, 0x68, 0x2e, 0x76, 0x30, 0x2e, 0x53, 0x65, 0x61, 0x72, 0x63, 0x68, 0x49, 0x6e, + 0x64, 0x65, 0x78, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x22, 0x26, 0x82, 0xd3, 0xe4, + 0x93, 0x02, 0x20, 0x3a, 0x01, 0x2a, 0x22, 0x1b, 0x2f, 0x61, 0x70, 0x69, 0x2f, 0x76, 0x30, 0x2f, + 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x2f, 0x69, 0x6e, 0x64, 0x65, 0x78, 0x2f, 0x73, 0x65, 0x61, + 0x72, 0x63, 0x68, 0x42, 0xf2, 0x02, 0x5a, 0x4a, 0x67, 0x69, 0x74, 0x68, 0x75, 0x62, 0x2e, 0x63, + 0x6f, 0x6d, 0x2f, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2d, 0x65, 0x75, 0x2f, + 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2f, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x67, + 0x65, 0x6e, 0x2f, 0x67, 0x65, 0x6e, 0x2f, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, + 0x2f, 0x73, 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x2f, 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x2f, + 0x76, 0x30, 0x92, 0x41, 0xa2, 0x02, 0x12, 0xb7, 0x01, 0x0a, 0x10, 0x4f, 0x70, 0x65, 0x6e, 0x43, + 0x6c, 0x6f, 0x75, 0x64, 0x20, 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x22, 0x51, 0x0a, 0x0e, 0x4f, + 0x70, 0x65, 0x6e, 0x43, 0x6c, 0x6f, 0x75, 0x64, 0x20, 0x47, 0x6d, 0x62, 0x48, 0x12, 0x29, 0x68, + 0x74, 0x74, 0x70, 0x73, 0x3a, 0x2f, 0x2f, 0x67, 0x69, 0x74, 0x68, 0x75, 0x62, 0x2e, 0x63, 0x6f, + 0x6d, 0x2f, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2d, 0x65, 0x75, 0x2f, 0x6f, + 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x1a, 0x14, 0x73, 0x75, 0x70, 0x70, 0x6f, 0x72, + 0x74, 0x40, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2e, 0x65, 0x75, 0x2a, 0x49, + 0x0a, 0x0a, 0x41, 0x70, 0x61, 0x63, 0x68, 0x65, 0x2d, 0x32, 0x2e, 0x30, 0x12, 0x3b, 0x68, 0x74, + 0x74, 0x70, 0x73, 0x3a, 0x2f, 0x2f, 0x67, 0x69, 0x74, 0x68, 0x75, 0x62, 0x2e, 0x63, 0x6f, 0x6d, + 0x2f, 0x6f, 0x70, 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2d, 0x65, 0x75, 0x2f, 0x6f, 0x70, + 0x65, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2f, 0x62, 0x6c, 0x6f, 0x62, 0x2f, 0x6d, 0x61, 0x69, + 0x6e, 0x2f, 0x4c, 0x49, 0x43, 0x45, 0x4e, 0x53, 0x45, 0x32, 0x05, 0x31, 0x2e, 0x30, 0x2e, 0x30, + 0x2a, 0x02, 0x01, 0x02, 0x32, 0x10, 0x61, 0x70, 0x70, 0x6c, 0x69, 0x63, 0x61, 0x74, 0x69, 0x6f, + 0x6e, 0x2f, 0x6a, 0x73, 0x6f, 0x6e, 0x3a, 0x10, 0x61, 0x70, 0x70, 0x6c, 0x69, 0x63, 0x61, 0x74, + 0x69, 0x6f, 0x6e, 0x2f, 0x6a, 0x73, 0x6f, 0x6e, 0x72, 0x3e, 0x0a, 0x10, 0x44, 0x65, 0x76, 0x65, + 0x6c, 0x6f, 0x70, 0x65, 0x72, 0x20, 0x4d, 0x61, 0x6e, 0x75, 0x61, 0x6c, 0x12, 0x2a, 0x68, 0x74, + 0x74, 0x70, 0x73, 0x3a, 0x2f, 0x2f, 0x64, 0x6f, 0x63, 0x73, 0x2e, 0x6f, 0x70, 0x65, 0x6e, 0x63, + 0x6c, 0x6f, 0x75, 0x64, 0x2e, 0x65, 0x75, 0x2f, 0x73, 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x73, + 0x2f, 0x73, 0x65, 0x61, 0x72, 0x63, 0x68, 0x2f, 0x62, 0x06, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x33, } var ( @@ -613,35 +688,38 @@ func file_opencloud_services_search_v0_search_proto_rawDescGZIP() []byte { return file_opencloud_services_search_v0_search_proto_rawDescData } +var file_opencloud_services_search_v0_search_proto_enumTypes = make([]protoimpl.EnumInfo, 1) var file_opencloud_services_search_v0_search_proto_msgTypes = make([]protoimpl.MessageInfo, 6) var file_opencloud_services_search_v0_search_proto_goTypes = []interface{}{ - (*SearchRequest)(nil), // 0: opencloud.services.search.v0.SearchRequest - (*SearchResponse)(nil), // 1: opencloud.services.search.v0.SearchResponse - (*SearchIndexRequest)(nil), // 2: opencloud.services.search.v0.SearchIndexRequest - (*SearchIndexResponse)(nil), // 3: opencloud.services.search.v0.SearchIndexResponse - (*IndexSpaceRequest)(nil), // 4: opencloud.services.search.v0.IndexSpaceRequest - (*IndexSpaceResponse)(nil), // 5: opencloud.services.search.v0.IndexSpaceResponse - (*v0.Reference)(nil), // 6: opencloud.messages.search.v0.Reference - (*v0.Match)(nil), // 7: opencloud.messages.search.v0.Match - (*durationpb.Duration)(nil), // 8: google.protobuf.Duration + (IndexSpaceResponse_Status)(0), // 0: opencloud.services.search.v0.IndexSpaceResponse.Status + (*SearchRequest)(nil), // 1: opencloud.services.search.v0.SearchRequest + (*SearchResponse)(nil), // 2: opencloud.services.search.v0.SearchResponse + (*SearchIndexRequest)(nil), // 3: opencloud.services.search.v0.SearchIndexRequest + (*SearchIndexResponse)(nil), // 4: opencloud.services.search.v0.SearchIndexResponse + (*IndexSpaceRequest)(nil), // 5: opencloud.services.search.v0.IndexSpaceRequest + (*IndexSpaceResponse)(nil), // 6: opencloud.services.search.v0.IndexSpaceResponse + (*v0.Reference)(nil), // 7: opencloud.messages.search.v0.Reference + (*v0.Match)(nil), // 8: opencloud.messages.search.v0.Match + (*durationpb.Duration)(nil), // 9: google.protobuf.Duration } var file_opencloud_services_search_v0_search_proto_depIdxs = []int32{ - 6, // 0: opencloud.services.search.v0.SearchRequest.ref:type_name -> opencloud.messages.search.v0.Reference - 7, // 1: opencloud.services.search.v0.SearchResponse.matches:type_name -> opencloud.messages.search.v0.Match - 6, // 2: opencloud.services.search.v0.SearchIndexRequest.ref:type_name -> opencloud.messages.search.v0.Reference - 7, // 3: opencloud.services.search.v0.SearchIndexResponse.matches:type_name -> opencloud.messages.search.v0.Match - 8, // 4: opencloud.services.search.v0.IndexSpaceResponse.space_duration:type_name -> google.protobuf.Duration - 0, // 5: opencloud.services.search.v0.SearchProvider.Search:input_type -> opencloud.services.search.v0.SearchRequest - 4, // 6: opencloud.services.search.v0.SearchProvider.IndexSpace:input_type -> opencloud.services.search.v0.IndexSpaceRequest - 2, // 7: opencloud.services.search.v0.IndexProvider.Search:input_type -> opencloud.services.search.v0.SearchIndexRequest - 1, // 8: opencloud.services.search.v0.SearchProvider.Search:output_type -> opencloud.services.search.v0.SearchResponse - 5, // 9: opencloud.services.search.v0.SearchProvider.IndexSpace:output_type -> opencloud.services.search.v0.IndexSpaceResponse - 3, // 10: opencloud.services.search.v0.IndexProvider.Search:output_type -> opencloud.services.search.v0.SearchIndexResponse - 8, // [8:11] is the sub-list for method output_type - 5, // [5:8] 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 + 7, // 0: opencloud.services.search.v0.SearchRequest.ref:type_name -> opencloud.messages.search.v0.Reference + 8, // 1: opencloud.services.search.v0.SearchResponse.matches:type_name -> opencloud.messages.search.v0.Match + 7, // 2: opencloud.services.search.v0.SearchIndexRequest.ref:type_name -> opencloud.messages.search.v0.Reference + 8, // 3: opencloud.services.search.v0.SearchIndexResponse.matches:type_name -> opencloud.messages.search.v0.Match + 9, // 4: opencloud.services.search.v0.IndexSpaceResponse.space_duration:type_name -> google.protobuf.Duration + 0, // 5: opencloud.services.search.v0.IndexSpaceResponse.status:type_name -> opencloud.services.search.v0.IndexSpaceResponse.Status + 1, // 6: opencloud.services.search.v0.SearchProvider.Search:input_type -> opencloud.services.search.v0.SearchRequest + 5, // 7: opencloud.services.search.v0.SearchProvider.IndexSpace:input_type -> opencloud.services.search.v0.IndexSpaceRequest + 3, // 8: opencloud.services.search.v0.IndexProvider.Search:input_type -> opencloud.services.search.v0.SearchIndexRequest + 2, // 9: opencloud.services.search.v0.SearchProvider.Search:output_type -> opencloud.services.search.v0.SearchResponse + 6, // 10: opencloud.services.search.v0.SearchProvider.IndexSpace:output_type -> opencloud.services.search.v0.IndexSpaceResponse + 4, // 11: opencloud.services.search.v0.IndexProvider.Search:output_type -> opencloud.services.search.v0.SearchIndexResponse + 9, // [9:12] is the sub-list for method output_type + 6, // [6:9] is the sub-list for method input_type + 6, // [6:6] is the sub-list for extension type_name + 6, // [6:6] is the sub-list for extension extendee + 0, // [0:6] is the sub-list for field type_name } func init() { file_opencloud_services_search_v0_search_proto_init() } @@ -728,13 +806,14 @@ func file_opencloud_services_search_v0_search_proto_init() { File: protoimpl.DescBuilder{ GoPackagePath: reflect.TypeOf(x{}).PkgPath(), RawDescriptor: file_opencloud_services_search_v0_search_proto_rawDesc, - NumEnums: 0, + NumEnums: 1, NumMessages: 6, NumExtensions: 0, NumServices: 2, }, GoTypes: file_opencloud_services_search_v0_search_proto_goTypes, DependencyIndexes: file_opencloud_services_search_v0_search_proto_depIdxs, + EnumInfos: file_opencloud_services_search_v0_search_proto_enumTypes, MessageInfos: file_opencloud_services_search_v0_search_proto_msgTypes, }.Build() File_opencloud_services_search_v0_search_proto = out.File diff --git a/protogen/gen/opencloud/services/search/v0/search.swagger.json b/protogen/gen/opencloud/services/search/v0/search.swagger.json index 43ff1e2c79..db24bd32b5 100644 --- a/protogen/gen/opencloud/services/search/v0/search.swagger.json +++ b/protogen/gen/opencloud/services/search/v0/search.swagger.json @@ -46,7 +46,7 @@ "$ref": "#/definitions/v0IndexSpaceResponse" }, "error": { - "$ref": "#/definitions/rpcStatus" + "$ref": "#/definitions/googlerpcStatus" } }, "title": "Stream result of v0IndexSpaceResponse" @@ -55,7 +55,7 @@ "default": { "description": "An unexpected error response.", "schema": { - "$ref": "#/definitions/rpcStatus" + "$ref": "#/definitions/googlerpcStatus" } } }, @@ -87,7 +87,7 @@ "default": { "description": "An unexpected error response.", "schema": { - "$ref": "#/definitions/rpcStatus" + "$ref": "#/definitions/googlerpcStatus" } } }, @@ -119,7 +119,7 @@ "default": { "description": "An unexpected error response.", "schema": { - "$ref": "#/definitions/rpcStatus" + "$ref": "#/definitions/googlerpcStatus" } } }, @@ -140,16 +140,7 @@ } }, "definitions": { - "protobufAny": { - "type": "object", - "properties": { - "@type": { - "type": "string" - } - }, - "additionalProperties": {} - }, - "rpcStatus": { + "googlerpcStatus": { "type": "object", "properties": { "code": { @@ -167,6 +158,15 @@ } } }, + "protobufAny": { + "type": "object", + "properties": { + "@type": { + "type": "string" + } + }, + "additionalProperties": {} + }, "v0Audio": { "type": "object", "properties": { @@ -382,9 +382,24 @@ "error": { "type": "string", "description": "Contains an error message in case indexing this particular space failed." + }, + "status": { + "$ref": "#/definitions/v0IndexSpaceResponseStatus", + "description": "The outcome of indexing this particular space." } } }, + "v0IndexSpaceResponseStatus": { + "type": "string", + "enum": [ + "STATUS_UNKNOWN", + "STATUS_SUCCESS", + "STATUS_SKIPPED", + "STATUS_ERROR" + ], + "default": "STATUS_UNKNOWN", + "description": " - STATUS_SUCCESS: The space has been indexed successfully.\n - STATUS_SKIPPED: The space has been skipped because it is disabled. It will be indexed\nonce it gets enabled again.\n - STATUS_ERROR: Indexing the space failed, see the error field for details." + }, "v0LivePhoto": { "type": "object", "properties": { diff --git a/protogen/proto/opencloud/services/search/v0/search.proto b/protogen/proto/opencloud/services/search/v0/search.proto index 62b4b31436..806245d95e 100644 --- a/protogen/proto/opencloud/services/search/v0/search.proto +++ b/protogen/proto/opencloud/services/search/v0/search.proto @@ -112,6 +112,17 @@ message IndexSpaceRequest { } message IndexSpaceResponse { + enum Status { + STATUS_UNKNOWN = 0; + // The space has been indexed successfully. + STATUS_SUCCESS = 1; + // The space has been skipped because it is disabled. It will be indexed + // once it gets enabled again. + STATUS_SKIPPED = 2; + // Indexing the space failed, see the error field for details. + STATUS_ERROR = 3; + } + // The id of the space that has just been indexed. string space_id = 1; // The duration it took to index this space. @@ -122,4 +133,6 @@ message IndexSpaceResponse { int64 total_spaces = 4; // Contains an error message in case indexing this particular space failed. string error = 5; + // The outcome of indexing this particular space. + Status status = 6; } diff --git a/services/search/pkg/command/index.go b/services/search/pkg/command/index.go index c9b27814c1..86a8ae0e1e 100644 --- a/services/search/pkg/command/index.go +++ b/services/search/pkg/command/index.go @@ -8,6 +8,7 @@ import ( "io" "os" "os/signal" + "strconv" "syscall" "github.com/opencloud-eu/opencloud/pkg/config/configlog" @@ -92,14 +93,7 @@ func Index(cfg *config.Config) *cobra.Command { return err } - if progress.GetError() != "" { - fmt.Printf("[%d/%d] failed to index space %s: %s\n", - progress.GetIndexedSpaces(), progress.GetTotalSpaces(), progress.GetSpaceId(), progress.GetError()) - continue - } - - fmt.Printf("[%d/%d] indexed space %s in %s\n", - progress.GetIndexedSpaces(), progress.GetTotalSpaces(), progress.GetSpaceId(), progress.GetSpaceDuration().AsDuration()) + printProgress(progress) } return nil }, @@ -138,3 +132,30 @@ func Index(cfg *config.Config) *cobra.Command { return indexCmd } + +// printProgress prints a single progress line, e.g. +// +// [ 1/12 SKIPPED] is disabled, it will be indexed once it is enabled again +// [ 2/12 SUCCESS] indexed in 40.9ms +// [ 3/12 ERROR ] failed: +// +// The counter is padded to the width of the total and the status to the +// longest status so that the messages line up. +func printProgress(progress *searchsvc.IndexSpaceResponse) { + var status, msg string + switch { + case progress.GetStatus() == searchsvc.IndexSpaceResponse_STATUS_SKIPPED: + status = "SKIPPED" + msg = progress.GetSpaceId() + " is disabled, it will be indexed once it is enabled again" + // servers not setting a status only report failures via the error field + case progress.GetStatus() == searchsvc.IndexSpaceResponse_STATUS_ERROR || progress.GetError() != "": + status = "ERROR" + msg = progress.GetSpaceId() + " failed: " + progress.GetError() + default: + status = "SUCCESS" + msg = fmt.Sprintf("%s indexed in %s", progress.GetSpaceId(), progress.GetSpaceDuration().AsDuration()) + } + + width := len(strconv.FormatInt(progress.GetTotalSpaces(), 10)) + fmt.Printf("[%*d/%d %-7s] %s\n", width, progress.GetIndexedSpaces(), progress.GetTotalSpaces(), status, msg) +} diff --git a/services/search/pkg/command/server.go b/services/search/pkg/command/server.go index dbfe28af9b..269292e6e2 100644 --- a/services/search/pkg/command/server.go +++ b/services/search/pkg/command/server.go @@ -29,6 +29,8 @@ import ( "github.com/opencloud-eu/reva/v2/pkg/events/raw" "github.com/opencloud-eu/reva/v2/pkg/rgrpc/todo/pool" "github.com/spf13/cobra" + + "github.com/nats-io/nats.go/jetstream" ) // Server is the entrypoint for the server command. @@ -122,6 +124,40 @@ func Server(cfg *config.Config) *cobra.Command { ss := search.NewService(selector, eng, extractor, mtrcs, logger, cfg) + // The KV bucket is used to remember spaces that were skipped during + // (re)indexing because they were disabled. It is shared between the gRPC + // reindex handler (which fills it) and the event consumer (which drains it + // again once a space gets enabled). + rawEventsCfg := raw.Config{ + Endpoint: cfg.Events.Endpoint, + Cluster: cfg.Events.Cluster, + EnableTLS: cfg.Events.EnableTLS, + TLSInsecure: cfg.Events.TLSInsecure, + TLSRootCACertificate: cfg.Events.TLSRootCACertificate, + AuthUsername: cfg.Events.AuthUsername, + AuthPassword: cfg.Events.AuthPassword, + MaxAckPending: cfg.Events.MaxAckPending, + AckWait: cfg.Events.AckWait, + } + skippedSpaces := search.NewSkippedSpaces(nil) + if !cfg.Events.Disabled { + kvConnName := generators.GenerateConnectionName(cfg.Service.Name, generators.NTypeKeyValue) + js, err := raw.JetStream(ctx, kvConnName, rawEventsCfg) + if err != nil { + logger.Error().Err(err).Msg("Failed to connect to NATS jetstream") + return err + } + kv, err := js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{ + Bucket: search.SkippedSpacesBucket, + Description: "Spaces that were skipped during (re)indexing because they were disabled", + }) + if err != nil { + logger.Error().Err(err).Msg("Failed to create the skipped spaces KV bucket") + return err + } + skippedSpaces = search.NewSkippedSpaces(kv) + } + // setup the servers gr := runner.NewGroup() @@ -136,6 +172,7 @@ func Server(cfg *config.Config) *cobra.Command { grpc.TraceProvider(traceProvider), grpc.GatewaySelector(selector), grpc.Searcher(ss), + grpc.SkippedSpaces(skippedSpaces), ) if err != nil { logger.Error().Err(err).Str("transport", "grpc").Msg("Failed to initialize server") @@ -149,23 +186,13 @@ func Server(cfg *config.Config) *cobra.Command { if !cfg.Events.Disabled { connName := generators.GenerateConnectionName(cfg.Service.Name, generators.NTypeBus) - bus, err := raw.FromConfig(context.Background(), connName, raw.Config{ - Endpoint: cfg.Events.Endpoint, - Cluster: cfg.Events.Cluster, - EnableTLS: cfg.Events.EnableTLS, - TLSInsecure: cfg.Events.TLSInsecure, - TLSRootCACertificate: cfg.Events.TLSRootCACertificate, - AuthUsername: cfg.Events.AuthUsername, - AuthPassword: cfg.Events.AuthPassword, - MaxAckPending: cfg.Events.MaxAckPending, - AckWait: cfg.Events.AckWait, - }) + bus, err := raw.FromConfig(context.Background(), connName, rawEventsCfg) if err != nil { logger.Error().Err(err).Msg("Failed to create event bus client") return err } - eventSvc, err := svcEvent.New(ctx, bus, logger, traceProvider, mtrcs, ss, cfg.Events.DebounceDuration, cfg.Events.NumConsumers, cfg.Events.AsyncUploads) + eventSvc, err := svcEvent.New(ctx, bus, logger, traceProvider, mtrcs, ss, skippedSpaces, cfg.Events.DebounceDuration, cfg.Events.NumConsumers, cfg.Events.AsyncUploads) if err != nil { logger.Error().Err(err).Str("transport", "event").Msg("Failed to initialize server") return err diff --git a/services/search/pkg/search/skippedspaces.go b/services/search/pkg/search/skippedspaces.go new file mode 100644 index 0000000000..5fb880d4b0 --- /dev/null +++ b/services/search/pkg/search/skippedspaces.go @@ -0,0 +1,136 @@ +package search + +import ( + "context" + "encoding/base64" + "encoding/json" + "errors" + "time" + + provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1" + "github.com/nats-io/nats.go/jetstream" + "github.com/opencloud-eu/reva/v2/pkg/storagespace" +) + +const ( + SkippedSpacesBucket = "search-skipped-spaces" + skippedSpacesTimeout = 10 * time.Second +) + +// SkippedSpaces keeps track of spaces that were skipped during (re)indexing +// because they were disabled. When such a space gets enabled again it can be +// looked up here and reindexed. +type SkippedSpaces struct { + kv jetstream.KeyValue +} + +// skippedSpace is the value stored in the KV bucket for a skipped space. +type skippedSpace struct { + SpaceID string `json:"spaceId"` + ForceReindex bool `json:"forceReindex"` +} + +// NewSkippedSpaces creates a new SkippedSpaces tracker backed by the given +// NATS JS KV bucket. A nil bucket turns all operations into no-ops, which is +// handy for tests and for deployments that run without an event bus. +func NewSkippedSpaces(kv jetstream.KeyValue) *SkippedSpaces { + return &SkippedSpaces{kv: kv} +} + +// key returns the canonical, KV-safe key for the given space id. NATS KV keys +// may only contain a limited set of characters, so the (possibly `$`/`!` +// separated) space id is normalized and base64url encoded. +func (s *SkippedSpaces) key(spaceID *provider.StorageSpaceId) (string, error) { + rid, err := storagespace.ParseID(spaceID.GetOpaqueId()) + if err != nil { + return "", err + } + canonical := storagespace.FormatResourceID(&provider.ResourceId{ + StorageId: rid.GetStorageId(), + SpaceId: rid.GetSpaceId(), + }) + return base64.RawURLEncoding.EncodeToString([]byte(canonical)), nil +} + +// Mark records that the given space was skipped while it was disabled. +// forceReindex indicates whether the skipped scan was a forced reindex. An +// existing forced mark is never downgraded to a shallow one. +func (s *SkippedSpaces) Mark(spaceID *provider.StorageSpaceId, forceReindex bool) error { + if s == nil || s.kv == nil { + return nil + } + key, err := s.key(spaceID) + if err != nil { + return err + } + ctx, cancel := context.WithTimeout(context.Background(), skippedSpacesTimeout) + defer cancel() + + if !forceReindex { + entry, err := s.kv.Get(ctx, key) + switch { + case err == nil: + if decodeSkippedSpace(entry.Value()).ForceReindex { + return nil + } + case !errors.Is(err, jetstream.ErrKeyNotFound): + return err + } + } + + value, err := json.Marshal(skippedSpace{SpaceID: spaceID.GetOpaqueId(), ForceReindex: forceReindex}) + if err != nil { + return err + } + _, err = s.kv.Put(ctx, key, value) + return err +} + +// decodeSkippedSpace decodes a KV value. Values that can not be decoded are +// treated as shallow (non-forced) marks. +func decodeSkippedSpace(value []byte) skippedSpace { + var ss skippedSpace + _ = json.Unmarshal(value, &ss) + return ss +} + +// Unmark removes any record for the given space. It is a no-op if the space +// was not tracked. +func (s *SkippedSpaces) Unmark(spaceID *provider.StorageSpaceId) error { + if s == nil || s.kv == nil { + return nil + } + key, err := s.key(spaceID) + if err != nil { + return err + } + ctx, cancel := context.WithTimeout(context.Background(), skippedSpacesTimeout) + defer cancel() + if err := s.kv.Delete(ctx, key); err != nil && !errors.Is(err, jetstream.ErrKeyNotFound) { + return err + } + return nil +} + +// IsMarked reports whether the given space was previously skipped while it was +// disabled, and whether the skipped scan was a forced reindex. +func (s *SkippedSpaces) IsMarked(spaceID *provider.StorageSpaceId) (marked bool, forceReindex bool, err error) { + if s == nil || s.kv == nil { + return false, false, nil + } + key, err := s.key(spaceID) + if err != nil { + return false, false, err + } + ctx, cancel := context.WithTimeout(context.Background(), skippedSpacesTimeout) + defer cancel() + entry, err := s.kv.Get(ctx, key) + switch { + case err == nil: + return true, decodeSkippedSpace(entry.Value()).ForceReindex, nil + case errors.Is(err, jetstream.ErrKeyNotFound): + return false, false, nil + default: + return false, false, err + } +} diff --git a/services/search/pkg/search/skippedspaces_test.go b/services/search/pkg/search/skippedspaces_test.go new file mode 100644 index 0000000000..c0d6830404 --- /dev/null +++ b/services/search/pkg/search/skippedspaces_test.go @@ -0,0 +1,138 @@ +package search_test + +import ( + "context" + "time" + + provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1" + nserver "github.com/nats-io/nats-server/v2/server" + "github.com/nats-io/nats.go" + "github.com/nats-io/nats.go/jetstream" + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" + + "github.com/opencloud-eu/opencloud/services/search/pkg/search" +) + +var _ = Describe("SkippedSpaces", func() { + var ( + ctx context.Context + kv jetstream.KeyValue + skipped *search.SkippedSpaces + + spaceID = &provider.StorageSpaceId{OpaqueId: "storage$space"} + ) + + BeforeEach(func() { + ctx = context.Background() + + srv, err := nserver.NewServer(&nserver.Options{ + DontListen: true, + JetStream: true, + StoreDir: GinkgoT().TempDir(), + NoLog: true, + NoSigs: true, + }) + Expect(err).ToNot(HaveOccurred()) + go srv.Start() + Expect(srv.ReadyForConnections(5 * time.Second)).To(BeTrue()) + DeferCleanup(srv.Shutdown) + + nc, err := nats.Connect("", nats.InProcessServer(srv)) + Expect(err).ToNot(HaveOccurred()) + DeferCleanup(nc.Close) + + js, err := jetstream.New(nc) + Expect(err).ToNot(HaveOccurred()) + kv, err = js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{Bucket: search.SkippedSpacesBucket}) + Expect(err).ToNot(HaveOccurred()) + + skipped = search.NewSkippedSpaces(kv) + }) + + isMarked := func(id *provider.StorageSpaceId) (bool, bool) { + marked, force, err := skipped.IsMarked(id) + Expect(err).ToNot(HaveOccurred()) + return marked, force + } + + It("reports unknown spaces as not marked", func() { + marked, force := isMarked(spaceID) + Expect(marked).To(BeFalse()) + Expect(force).To(BeFalse()) + }) + + It("marks a space for a shallow reindex", func() { + Expect(skipped.Mark(spaceID, false)).To(Succeed()) + + marked, force := isMarked(spaceID) + Expect(marked).To(BeTrue()) + Expect(force).To(BeFalse()) + }) + + It("marks a space for a forced reindex", func() { + Expect(skipped.Mark(spaceID, true)).To(Succeed()) + + marked, force := isMarked(spaceID) + Expect(marked).To(BeTrue()) + Expect(force).To(BeTrue()) + }) + + It("does not downgrade a forced mark to a shallow one", func() { + Expect(skipped.Mark(spaceID, true)).To(Succeed()) + Expect(skipped.Mark(spaceID, false)).To(Succeed()) + + _, force := isMarked(spaceID) + Expect(force).To(BeTrue()) + }) + + It("upgrades a shallow mark to a forced one", func() { + Expect(skipped.Mark(spaceID, false)).To(Succeed()) + Expect(skipped.Mark(spaceID, true)).To(Succeed()) + + _, force := isMarked(spaceID) + Expect(force).To(BeTrue()) + }) + + It("unmarks a space", func() { + Expect(skipped.Mark(spaceID, true)).To(Succeed()) + Expect(skipped.Unmark(spaceID)).To(Succeed()) + + marked, force := isMarked(spaceID) + Expect(marked).To(BeFalse()) + Expect(force).To(BeFalse()) + }) + + It("does not fail when unmarking an unknown space", func() { + Expect(skipped.Unmark(spaceID)).To(Succeed()) + }) + + It("keeps spaces apart", func() { + other := &provider.StorageSpaceId{OpaqueId: "storage$other"} + Expect(skipped.Mark(spaceID, true)).To(Succeed()) + + marked, _ := isMarked(other) + Expect(marked).To(BeFalse()) + }) + + It("returns an error for invalid space ids", func() { + invalid := &provider.StorageSpaceId{} + Expect(skipped.Mark(invalid, false)).ToNot(Succeed()) + Expect(skipped.Unmark(invalid)).ToNot(Succeed()) + _, _, err := skipped.IsMarked(invalid) + Expect(err).To(HaveOccurred()) + }) + + DescribeTable("is a no-op without a bucket", + func(s *search.SkippedSpaces) { + Expect(s.Mark(spaceID, true)).To(Succeed()) + Expect(s.Unmark(spaceID)).To(Succeed()) + marked, force, err := s.IsMarked(spaceID) + Expect(err).ToNot(HaveOccurred()) + Expect(marked).To(BeFalse()) + Expect(force).To(BeFalse()) + }, + Entry("nil bucket", search.NewSkippedSpaces(nil)), + Entry("nil tracker", (*search.SkippedSpaces)(nil)), + ) +}) diff --git a/services/search/pkg/server/grpc/option.go b/services/search/pkg/server/grpc/option.go index 1686dce219..77378992a1 100644 --- a/services/search/pkg/server/grpc/option.go +++ b/services/search/pkg/server/grpc/option.go @@ -29,6 +29,7 @@ type Options struct { TraceProvider trace.TracerProvider GatewaySelector *pool.Selector[gateway.GatewayAPIClient] Searcher search.Searcher + SkippedSpaces *search.SkippedSpaces } // newOptions initializes the available default options. @@ -111,3 +112,10 @@ func Searcher(val search.Searcher) Option { o.Searcher = val } } + +// SkippedSpaces provides a function to set the SkippedSpaces option. +func SkippedSpaces(val *search.SkippedSpaces) Option { + return func(o *Options) { + o.SkippedSpaces = val + } +} diff --git a/services/search/pkg/server/grpc/server.go b/services/search/pkg/server/grpc/server.go index b6e8f1868e..174a9dd586 100644 --- a/services/search/pkg/server/grpc/server.go +++ b/services/search/pkg/server/grpc/server.go @@ -39,6 +39,7 @@ func Server(opts ...Option) (grpc.Service, error) { svc.Metrics(options.Metrics), svc.GatewaySelector(options.GatewaySelector), svc.Searcher(options.Searcher), + svc.SkippedSpaces(options.SkippedSpaces), ) if err != nil { options.Logger.Error(). diff --git a/services/search/pkg/service/event/debouncer.go b/services/search/pkg/service/event/debouncer.go index 1444b9dbd7..9d2a10b08c 100644 --- a/services/search/pkg/service/event/debouncer.go +++ b/services/search/pkg/service/event/debouncer.go @@ -13,7 +13,7 @@ import ( type SpaceDebouncer struct { after time.Duration timeout time.Duration - f func(id *provider.StorageSpaceId) + f func(id *provider.StorageSpaceId, force bool) pending map[string]*workItem inProgress sync.Map @@ -24,6 +24,7 @@ type SpaceDebouncer struct { type workItem struct { t *time.Timer timeout *time.Timer + force bool work func() } @@ -31,7 +32,7 @@ type workItem struct { type AckFunc func() error // NewSpaceDebouncer returns a new SpaceDebouncer instance -func NewSpaceDebouncer(d time.Duration, timeout time.Duration, f func(id *provider.StorageSpaceId), logger log.Logger) *SpaceDebouncer { +func NewSpaceDebouncer(d time.Duration, timeout time.Duration, f func(id *provider.StorageSpaceId, force bool), logger log.Logger) *SpaceDebouncer { return &SpaceDebouncer{ after: d, timeout: timeout, @@ -42,12 +43,15 @@ func NewSpaceDebouncer(d time.Duration, timeout time.Duration, f func(id *provid } } -// Debounce restars the debounce timer for the given space -func (d *SpaceDebouncer) Debounce(id *provider.StorageSpaceId, ack AckFunc) { +// Debounce restarts the debounce timer for the given space. If force is set, +// the scheduled run is a forced one. A pending run stays forced even if it is +// debounced again without force. +func (d *SpaceDebouncer) Debounce(id *provider.StorageSpaceId, ack AckFunc, force bool) { d.mutex.Lock() defer d.mutex.Unlock() if wi := d.pending[id.OpaqueId]; wi != nil { + wi.force = wi.force || force if ack != nil { go ack() // Acknowledge the event immediately, the according space is already scheduled for indexing } @@ -55,7 +59,7 @@ func (d *SpaceDebouncer) Debounce(id *provider.StorageSpaceId, ack AckFunc) { return } - wi := &workItem{} + wi := &workItem{force: force} wi.work = func() { if _, ok := d.inProgress.Load(id.OpaqueId); ok { // Reschedule this run for when the previous run has finished @@ -70,13 +74,14 @@ func (d *SpaceDebouncer) Debounce(id *provider.StorageSpaceId, ack AckFunc) { d.mutex.Lock() wi.timeout.Stop() // stop the timeout timer if it is running delete(d.pending, id.OpaqueId) + force := wi.force d.inProgress.Store(id.OpaqueId, true) defer func() { d.inProgress.Delete(id.OpaqueId) }() d.mutex.Unlock() // release the lock early to allow other goroutines to debounce - d.f(id) + d.f(id, force) go func() { if ack != nil { if err := ack(); err != nil { diff --git a/services/search/pkg/service/event/debouncer_test.go b/services/search/pkg/service/event/debouncer_test.go index 50b9ccfb34..8169b36625 100644 --- a/services/search/pkg/service/event/debouncer_test.go +++ b/services/search/pkg/service/event/debouncer_test.go @@ -25,7 +25,7 @@ var _ = Describe("SpaceDebouncer", func() { BeforeEach(func() { callCount = atomic.Int32{} - debouncer = event.NewSpaceDebouncer(50*time.Millisecond, 10*time.Second, func(id *sprovider.StorageSpaceId) { + debouncer = event.NewSpaceDebouncer(50*time.Millisecond, 10*time.Second, func(id *sprovider.StorageSpaceId, _ bool) { if id.OpaqueId == "spaceid" { callCount.Add(1) } @@ -33,22 +33,22 @@ var _ = Describe("SpaceDebouncer", func() { }) It("debounces", func() { - debouncer.Debounce(spaceid, nil) - debouncer.Debounce(spaceid, nil) - debouncer.Debounce(spaceid, nil) + debouncer.Debounce(spaceid, nil, false) + debouncer.Debounce(spaceid, nil, false) + debouncer.Debounce(spaceid, nil, false) Eventually(func() int { return int(callCount.Load()) }, "200ms").Should(Equal(1)) }) It("works multiple times", func() { - debouncer.Debounce(spaceid, nil) - debouncer.Debounce(spaceid, nil) - debouncer.Debounce(spaceid, nil) + debouncer.Debounce(spaceid, nil, false) + debouncer.Debounce(spaceid, nil, false) + debouncer.Debounce(spaceid, nil, false) time.Sleep(100 * time.Millisecond) - debouncer.Debounce(spaceid, nil) - debouncer.Debounce(spaceid, nil) + debouncer.Debounce(spaceid, nil, false) + debouncer.Debounce(spaceid, nil, false) Eventually(func() int { return int(callCount.Load()) @@ -56,16 +56,16 @@ var _ = Describe("SpaceDebouncer", func() { }) It("doesn't trigger twice simultaneously", func() { - debouncer = event.NewSpaceDebouncer(50*time.Millisecond, 5*time.Second, func(id *sprovider.StorageSpaceId) { + debouncer = event.NewSpaceDebouncer(50*time.Millisecond, 5*time.Second, func(id *sprovider.StorageSpaceId, _ bool) { if id.OpaqueId == "spaceid" { callCount.Add(1) } time.Sleep(300 * time.Millisecond) }, log.NewLogger()) - debouncer.Debounce(spaceid, nil) + debouncer.Debounce(spaceid, nil, false) time.Sleep(100 * time.Millisecond) // Let it trigger once - debouncer.Debounce(spaceid, nil) + debouncer.Debounce(spaceid, nil, false) time.Sleep(100 * time.Millisecond) // shouldn't trigger as the other run is still in progress Expect(int(callCount.Load())).To(Equal(1)) @@ -75,7 +75,7 @@ var _ = Describe("SpaceDebouncer", func() { }) It("fires at the timeout even when continuously debounced", func() { - debouncer = event.NewSpaceDebouncer(100*time.Millisecond, 250*time.Millisecond, func(id *sprovider.StorageSpaceId) { + debouncer = event.NewSpaceDebouncer(100*time.Millisecond, 250*time.Millisecond, func(id *sprovider.StorageSpaceId, _ bool) { if id.OpaqueId == "spaceid" { callCount.Add(1) } @@ -87,11 +87,11 @@ var _ = Describe("SpaceDebouncer", func() { // a Debounce call arriving right after the timeout fires would find // pending empty and schedule a second workItem, breaking the assertion // below that the work function is invoked exactly once. - debouncer.Debounce(spaceid, nil) + debouncer.Debounce(spaceid, nil, false) for i := 0; i < 4 && callCount.Load() == 0; i++ { time.Sleep(50 * time.Millisecond) if callCount.Load() == 0 { - debouncer.Debounce(spaceid, nil) + debouncer.Debounce(spaceid, nil, false) } } @@ -108,14 +108,14 @@ var _ = Describe("SpaceDebouncer", func() { }) It("doesn't run the timeout function if the work function has been called", func() { - debouncer = event.NewSpaceDebouncer(100*time.Millisecond, 250*time.Millisecond, func(id *sprovider.StorageSpaceId) { + debouncer = event.NewSpaceDebouncer(100*time.Millisecond, 250*time.Millisecond, func(id *sprovider.StorageSpaceId, _ bool) { if id.OpaqueId == "spaceid" { callCount.Add(1) } }, log.NewLogger()) // Initial call to start the timers - debouncer.Debounce(spaceid, nil) + debouncer.Debounce(spaceid, nil, false) // Wait for the debounce timer to fire Eventually(func() int { @@ -134,7 +134,7 @@ var _ = Describe("SpaceDebouncer", func() { return nil } - debouncer.Debounce(spaceid, ackFunc) + debouncer.Debounce(spaceid, ackFunc, false) Eventually(func() int { return int(callCount.Load()) @@ -158,12 +158,12 @@ var _ = Describe("SpaceDebouncer", func() { } // First call, sets up the trigger - debouncer.Debounce(spaceid, firstAckFunc) + debouncer.Debounce(spaceid, firstAckFunc, false) Expect(firstAckCalled.Load()).To(BeFalse()) Expect(secondAckCalled.Load()).To(BeFalse()) // Second call, should call its ack immediately - debouncer.Debounce(spaceid, secondAckFunc) + debouncer.Debounce(spaceid, secondAckFunc, false) Eventually(func() bool { return secondAckCalled.Load() }, "50ms").Should(BeTrue()) @@ -178,4 +178,52 @@ var _ = Describe("SpaceDebouncer", func() { return firstAckCalled.Load() }, "200ms").Should(BeTrue()) }) + + Describe("force", func() { + var forced atomic.Bool + + BeforeEach(func() { + forced = atomic.Bool{} + debouncer = event.NewSpaceDebouncer(50*time.Millisecond, 10*time.Second, func(id *sprovider.StorageSpaceId, force bool) { + if id.OpaqueId == "spaceid" { + forced.Store(force) + callCount.Add(1) + } + }, log.NewLogger()) + }) + + It("is not forced by default", func() { + debouncer.Debounce(spaceid, nil, false) + Eventually(func() int { + return int(callCount.Load()) + }, "200ms").Should(Equal(1)) + Expect(forced.Load()).To(BeFalse()) + }) + + It("passes the force flag", func() { + debouncer.Debounce(spaceid, nil, true) + Eventually(func() int { + return int(callCount.Load()) + }, "200ms").Should(Equal(1)) + Expect(forced.Load()).To(BeTrue()) + }) + + It("keeps a pending run forced when debounced again without force", func() { + debouncer.Debounce(spaceid, nil, true) + debouncer.Debounce(spaceid, nil, false) + Eventually(func() int { + return int(callCount.Load()) + }, "200ms").Should(Equal(1)) + Expect(forced.Load()).To(BeTrue()) + }) + + It("upgrades a pending run to forced", func() { + debouncer.Debounce(spaceid, nil, false) + debouncer.Debounce(spaceid, nil, true) + Eventually(func() int { + return int(callCount.Load()) + }, "200ms").Should(Equal(1)) + Expect(forced.Load()).To(BeTrue()) + }) + }) }) diff --git a/services/search/pkg/service/event/service.go b/services/search/pkg/service/event/service.go index c4eff9bfb7..da1e79a715 100644 --- a/services/search/pkg/service/event/service.go +++ b/services/search/pkg/service/event/service.go @@ -32,6 +32,7 @@ type Service struct { tp trace.TracerProvider m *metrics.Metrics index search.Searcher + skippedSpaces *search.SkippedSpaces events []events.Unmarshaller stream raw.Stream indexSpaceDebouncer *SpaceDebouncer @@ -41,16 +42,17 @@ type Service struct { } // New returns a service implementation for Service. -func New(ctx context.Context, stream raw.Stream, logger log.Logger, tp trace.TracerProvider, m *metrics.Metrics, index search.Searcher, debounceDuration int, numConsumers int, asyncUploads bool) (Service, error) { +func New(ctx context.Context, stream raw.Stream, logger log.Logger, tp trace.TracerProvider, m *metrics.Metrics, index search.Searcher, skippedSpaces *search.SkippedSpaces, debounceDuration int, numConsumers int, asyncUploads bool) (Service, error) { svc := Service{ - ctx: ctx, - log: logger, - tp: tp, - m: m, - index: index, - stream: stream, - stopCh: make(chan struct{}, 1), - stopped: new(atomic.Bool), + ctx: ctx, + log: logger, + tp: tp, + m: m, + index: index, + skippedSpaces: skippedSpaces, + stream: stream, + stopCh: make(chan struct{}, 1), + stopped: new(atomic.Bool), events: []events.Unmarshaller{ events.ItemTrashed{}, events.ItemPurged{}, @@ -64,6 +66,7 @@ func New(ctx context.Context, stream raw.Stream, logger log.Logger, tp trace.Tra events.TagsRemoved{}, events.SpaceRenamed{}, events.SpaceDeleted{}, + events.SpaceEnabled{}, events.LabelAdded{}, events.LabelRemoved{}, }, @@ -76,8 +79,8 @@ func New(ctx context.Context, stream raw.Stream, logger log.Logger, tp trace.Tra svc.events = append(svc.events, events.FileUploaded{}) } - svc.indexSpaceDebouncer = NewSpaceDebouncer(time.Duration(debounceDuration)*time.Millisecond, 30*time.Second, func(id *provider.StorageSpaceId) { - if err := svc.index.IndexSpace(id, false); err != nil { + svc.indexSpaceDebouncer = NewSpaceDebouncer(time.Duration(debounceDuration)*time.Millisecond, 30*time.Second, func(id *provider.StorageSpaceId, force bool) { + if err := svc.index.IndexSpace(id, force); err != nil { svc.log.Error().Err(err).Interface("spaceID", id).Msg("error while indexing a space") } }, svc.log) @@ -174,9 +177,9 @@ func (s Service) processEvent(e raw.Event) error { ack := e.Ack - debounce := func(id *provider.StorageSpaceId) { + debounce := func(id *provider.StorageSpaceId, force bool) { ack = func() error { return nil } - s.indexSpaceDebouncer.Debounce(id, e.Ack) + s.indexSpaceDebouncer.Debounce(id, e.Ack, force) } var err error @@ -184,37 +187,55 @@ func (s Service) processEvent(e raw.Event) error { switch ev := e.Event.Event.(type) { case events.ItemTrashed: s.index.TrashItem(ev.ID) - debounce(getSpaceID(ev.Ref)) + debounce(getSpaceID(ev.Ref), false) case events.ItemPurged: s.index.PurgeItem(ev.Ref) case events.TrashbinPurged: err = s.index.PurgeDeleted(getSpaceID(ev.Ref)) case events.ItemMoved: s.index.MoveItem(ev.Ref) - debounce(getSpaceID(ev.Ref)) + debounce(getSpaceID(ev.Ref), false) case events.ItemRestored: s.index.RestoreItem(ev.Ref) - debounce(getSpaceID(ev.Ref)) + debounce(getSpaceID(ev.Ref), false) case events.ContainerCreated: - debounce(getSpaceID(ev.Ref)) + debounce(getSpaceID(ev.Ref), false) case events.FileTouched: - debounce(getSpaceID(ev.Ref)) + debounce(getSpaceID(ev.Ref), false) case events.FileVersionRestored: - debounce(getSpaceID(ev.Ref)) + debounce(getSpaceID(ev.Ref), false) case events.TagsAdded: s.index.UpsertItem(ev.Ref) - debounce(getSpaceID(ev.Ref)) + debounce(getSpaceID(ev.Ref), false) case events.TagsRemoved: s.index.UpsertItem(ev.Ref) - debounce(getSpaceID(ev.Ref)) + debounce(getSpaceID(ev.Ref), false) case events.FileUploaded: - debounce(getSpaceID(ev.Ref)) + debounce(getSpaceID(ev.Ref), false) case events.UploadReady: - debounce(getSpaceID(ev.FileRef)) + debounce(getSpaceID(ev.FileRef), false) case events.SpaceRenamed: - debounce(ev.ID) + debounce(ev.ID, false) + case events.SpaceEnabled: + // A space that was skipped during (re)indexing while it was disabled + // needs to be reindexed now that it is enabled again. + marked, force, err := s.skippedSpaces.IsMarked(ev.ID) + if err != nil { + s.log.Error().Err(err).Interface("spaceID", ev.ID).Msg("failed to check whether space was skipped while disabled") + break + } + if !marked { + break + } + if err = s.skippedSpaces.Unmark(ev.ID); err != nil { + s.log.Error().Err(err).Interface("spaceID", ev.ID).Msg("failed to remove space from the skipped spaces bucket") + } + debounce(ev.ID, force) case events.SpaceDeleted: err = s.index.PurgeSpace(ev.ID) + if err = s.skippedSpaces.Unmark(ev.ID); err != nil { + s.log.Error().Err(err).Interface("spaceID", ev.ID).Msg("failed to remove space from the skipped spaces bucket") + } case events.LabelAdded: s.index.UpsertItem(ev.Ref) case events.LabelRemoved: diff --git a/services/search/pkg/service/event/service_test.go b/services/search/pkg/service/event/service_test.go index d0afa2fea9..04984e1a0a 100644 --- a/services/search/pkg/service/event/service_test.go +++ b/services/search/pkg/service/event/service_test.go @@ -8,6 +8,7 @@ import ( . "github.com/onsi/ginkgo/v2" . "github.com/onsi/gomega" "github.com/opencloud-eu/opencloud/pkg/log" + "github.com/opencloud-eu/opencloud/services/search/pkg/search" searchMocks "github.com/opencloud-eu/opencloud/services/search/pkg/search/mocks" "github.com/opencloud-eu/opencloud/services/search/pkg/service/event" "github.com/opencloud-eu/reva/v2/pkg/events" @@ -27,7 +28,7 @@ var _ = DescribeTable("event", ch := make(chan raw.Event, 1) stream.EXPECT().Consume(mock.Anything, mock.Anything).Return((<-chan raw.Event)(ch), nil) - event, err := event.New(context.Background(), stream, log.NewLogger(), nil, nil, s, 50, 1, asyncUploads) + event, err := event.New(context.Background(), stream, log.NewLogger(), nil, nil, s, search.NewSkippedSpaces(nil), 50, 1, asyncUploads) Expect(err).NotTo(HaveOccurred()) go func() { diff --git a/services/search/pkg/service/grpc/v0/option.go b/services/search/pkg/service/grpc/v0/option.go index e284651895..23507a7a5e 100644 --- a/services/search/pkg/service/grpc/v0/option.go +++ b/services/search/pkg/service/grpc/v0/option.go @@ -23,6 +23,7 @@ type Options struct { Metrics *metrics.Metrics GatewaySelector *pool.Selector[gateway.GatewayAPIClient] Searcher search.Searcher + SkippedSpaces *search.SkippedSpaces } func newOptions(opts ...Option) Options { @@ -85,3 +86,10 @@ func Searcher(val search.Searcher) Option { o.Searcher = val } } + +// SkippedSpaces provides a function to set the SkippedSpaces option. +func SkippedSpaces(val *search.SkippedSpaces) Option { + return func(o *Options) { + o.SkippedSpaces = val + } +} diff --git a/services/search/pkg/service/grpc/v0/service.go b/services/search/pkg/service/grpc/v0/service.go index 6cb258ee43..48f14ec75c 100644 --- a/services/search/pkg/service/grpc/v0/service.go +++ b/services/search/pkg/service/grpc/v0/service.go @@ -31,6 +31,10 @@ import ( "github.com/opencloud-eu/opencloud/services/search/pkg/search" ) +// _spaceStateTrashed is the value the storage provider sets in the space's +// opaque "trashed" field when the space is disabled. +const _spaceStateTrashed = "trashed" + // NewHandler returns a service implementation for Service. func NewHandler(opts ...Option) (searchsvc.SearchProviderHandler, error) { options := newOptions(opts...) @@ -56,25 +60,27 @@ func NewHandler(opts ...Option) (searchsvc.SearchProviderHandler, error) { } return &Service{ - id: cfg.GRPC.Namespace + "." + cfg.Service.Name, - log: &options.Logger, - searcher: options.Searcher, - cache: cache, - tokenManager: tokenManager, - gws: options.GatewaySelector, - cfg: cfg, + id: cfg.GRPC.Namespace + "." + cfg.Service.Name, + log: &options.Logger, + searcher: options.Searcher, + skippedSpaces: options.SkippedSpaces, + cache: cache, + tokenManager: tokenManager, + gws: options.GatewaySelector, + cfg: cfg, }, nil } // Service implements the searchServiceHandler interface type Service struct { - id string - log *log.Logger - searcher search.Searcher - cache *ttlcache.Cache - tokenManager token.Manager - gws *pool.Selector[gateway.GatewayAPIClient] - cfg *config.Config + id string + log *log.Logger + searcher search.Searcher + skippedSpaces *search.SkippedSpaces + cache *ttlcache.Cache + tokenManager token.Manager + gws *pool.Selector[gateway.GatewayAPIClient] + cfg *config.Config } // Search handles the search @@ -134,9 +140,11 @@ func (s Service) IndexSpace(_ context.Context, in *searchsvc.IndexSpaceRequest, SpaceId: in.GetSpaceId(), IndexedSpaces: 1, TotalSpaces: 1, + Status: searchsvc.IndexSpaceResponse_STATUS_SUCCESS, } if err != nil { resp.Error = err.Error() + resp.Status = searchsvc.IndexSpaceResponse_STATUS_ERROR } if sendErr := stream.Send(resp); sendErr != nil { return sendErr @@ -193,11 +201,24 @@ func (s Service) IndexSpace(_ context.Context, in *searchsvc.IndexSpaceRequest, s.log.Info().Str("space_id", space.GetId().GetOpaqueId()).Msg("indexing space") t := time.Now() - indexErr := s.searcher.IndexSpace(space.GetId(), in.GetForceReindex()) - if indexErr != nil { - s.log.Error().Err(indexErr).Str("space_id", space.GetId().GetOpaqueId()).Msg("failed to index space") + var indexErr error + status := searchsvc.IndexSpaceResponse_STATUS_SUCCESS + if utils.ReadPlainFromOpaque(space.GetOpaque(), "trashed") == _spaceStateTrashed { + status = searchsvc.IndexSpaceResponse_STATUS_SKIPPED + // The space is disabled and cannot be indexed right now. Remember it + // so it gets reindexed once it is enabled again (see the SpaceEnabled + // event handler in the event service). + s.log.Info().Str("space_id", space.GetId().GetOpaqueId()).Msg("space is disabled, skipping and marking it for reindexing") + if markErr := s.skippedSpaces.Mark(space.GetId(), in.GetForceReindex()); markErr != nil { + s.log.Error().Err(markErr).Str("space_id", space.GetId().GetOpaqueId()).Msg("failed to mark disabled space for reindexing") + } } else { - s.log.Info().Str("space_id", space.GetId().GetOpaqueId()).Msg("finished indexing space") + indexErr = s.searcher.IndexSpace(space.GetId(), in.GetForceReindex()) + if indexErr != nil { + s.log.Error().Err(indexErr).Str("space_id", space.GetId().GetOpaqueId()).Msg("failed to index space") + } else { + s.log.Info().Str("space_id", space.GetId().GetOpaqueId()).Msg("finished indexing space") + } } mu.Lock() @@ -214,9 +235,11 @@ func (s Service) IndexSpace(_ context.Context, in *searchsvc.IndexSpaceRequest, IndexedSpaces: indexedCount, TotalSpaces: totalSpaces, SpaceDuration: durationpb.New(time.Since(t)), + Status: status, } if indexErr != nil { progress.Error = indexErr.Error() + progress.Status = searchsvc.IndexSpaceResponse_STATUS_ERROR } return stream.Send(progress) })