mirror of
https://github.com/opencloud-eu/opencloud.git
synced 2026-10-06 10:52:08 -04:00
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.
This commit is contained in:
1 parent
b955c7c4f0
commit
8f2adf2138
15 files changed
+728
-184
No files matched your search
@@ -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
|
||||
|
||||
@@ -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": {
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
@@ -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] <space id> is disabled, it will be indexed once it is enabled again
|
||||
// [ 2/12 SUCCESS] <space id> indexed in 40.9ms
|
||||
// [ 3/12 ERROR ] <space id> failed: <error>
|
||||
//
|
||||
// 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)
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
@@ -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)),
|
||||
)
|
||||
})
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
@@ -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().
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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())
|
||||
})
|
||||
})
|
||||
})
|
||||
@@ -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:
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
})
|
||||
|
||||
Reference in new issue
Block a user