From 338b6f2fdf5dc4513f6154de0aa14c4c325da59f Mon Sep 17 00:00:00 2001 From: Gareth George Date: Wed, 22 Nov 2023 23:46:03 -0800 Subject: [PATCH] chore: add indexing for snapshot field --- gen/go/v1/operations.pb.go | 124 +++---- internal/api/server.go | 13 +- .../indexutil/indexutil.go | 55 ++- .../indexutil/indexutil_test.go | 0 internal/{database => }/oplog/oplog.go | 312 +++++++++--------- internal/{database => }/oplog/oplog_test.go | 151 ++++++--- .../serializationutil/serializationutil.go | 9 +- .../serializationutil_test.go | 0 internal/orchestrator/orchestrator.go | 36 +- internal/orchestrator/tasks.go | 93 +++--- pkg/restic/outputs.go | 129 ++++---- proto/v1/operations.proto | 1 + resticui.go | 17 +- webui/gen/ts/v1/operations.pb.ts | 1 + 14 files changed, 522 insertions(+), 419 deletions(-) rename internal/{database => oplog}/indexutil/indexutil.go (54%) rename internal/{database => oplog}/indexutil/indexutil_test.go (100%) rename internal/{database => }/oplog/oplog.go (57%) rename internal/{database => }/oplog/oplog_test.go (57%) rename internal/{database => oplog}/serializationutil/serializationutil.go (85%) rename internal/{database => oplog}/serializationutil/serializationutil_test.go (100%) diff --git a/gen/go/v1/operations.pb.go b/gen/go/v1/operations.pb.go index 8552ad93..641c1cba 100644 --- a/gen/go/v1/operations.pb.go +++ b/gen/go/v1/operations.pb.go @@ -178,8 +178,9 @@ type Operation struct { unknownFields protoimpl.UnknownFields Id int64 `protobuf:"varint,1,opt,name=id,proto3" json:"id,omitempty"` - RepoId string `protobuf:"bytes,2,opt,name=repo_id,json=repoId,proto3" json:"repo_id,omitempty"` // repo id if associated with a repo (always true) - PlanId string `protobuf:"bytes,3,opt,name=plan_id,json=planId,proto3" json:"plan_id,omitempty"` // plan id if associated with a plan (always true) + RepoId string `protobuf:"bytes,2,opt,name=repo_id,json=repoId,proto3" json:"repo_id,omitempty"` // repo id if associated with a repo (always true) + PlanId string `protobuf:"bytes,3,opt,name=plan_id,json=planId,proto3" json:"plan_id,omitempty"` // plan id if associated with a plan (always true) + SnapshotId string `protobuf:"bytes,8,opt,name=snapshot_id,json=snapshotId,proto3" json:"snapshot_id,omitempty"` // snapshot id if associated with a snapshot. Status OperationStatus `protobuf:"varint,4,opt,name=status,proto3,enum=v1.OperationStatus" json:"status,omitempty"` UnixTimeStartMs int64 `protobuf:"varint,5,opt,name=unix_time_start_ms,json=unixTimeStartMs,proto3" json:"unix_time_start_ms,omitempty"` UnixTimeEndMs int64 `protobuf:"varint,6,opt,name=unix_time_end_ms,json=unixTimeEndMs,proto3" json:"unix_time_end_ms,omitempty"` @@ -244,6 +245,13 @@ func (x *Operation) GetPlanId() string { return "" } +func (x *Operation) GetSnapshotId() string { + if x != nil { + return x.SnapshotId + } + return "" +} + func (x *Operation) GetStatus() OperationStatus { if x != nil { return x.Status @@ -469,65 +477,67 @@ var file_v1_operations_proto_rawDesc = []byte{ 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x4c, 0x69, 0x73, 0x74, 0x12, 0x2d, 0x0a, 0x0a, 0x6f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x73, 0x18, 0x01, 0x20, 0x03, 0x28, 0x0b, 0x32, 0x0d, 0x2e, 0x76, 0x31, 0x2e, 0x4f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x52, 0x0a, - 0x6f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x73, 0x22, 0x99, 0x03, 0x0a, 0x09, 0x4f, + 0x6f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x73, 0x22, 0xba, 0x03, 0x0a, 0x09, 0x4f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x12, 0x0e, 0x0a, 0x02, 0x69, 0x64, 0x18, 0x01, 0x20, 0x01, 0x28, 0x03, 0x52, 0x02, 0x69, 0x64, 0x12, 0x17, 0x0a, 0x07, 0x72, 0x65, 0x70, 0x6f, 0x5f, 0x69, 0x64, 0x18, 0x02, 0x20, 0x01, 0x28, 0x09, 0x52, 0x06, 0x72, 0x65, 0x70, 0x6f, 0x49, 0x64, 0x12, 0x17, 0x0a, 0x07, 0x70, 0x6c, 0x61, 0x6e, 0x5f, 0x69, 0x64, 0x18, 0x03, 0x20, 0x01, - 0x28, 0x09, 0x52, 0x06, 0x70, 0x6c, 0x61, 0x6e, 0x49, 0x64, 0x12, 0x2b, 0x0a, 0x06, 0x73, 0x74, - 0x61, 0x74, 0x75, 0x73, 0x18, 0x04, 0x20, 0x01, 0x28, 0x0e, 0x32, 0x13, 0x2e, 0x76, 0x31, 0x2e, - 0x4f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x53, 0x74, 0x61, 0x74, 0x75, 0x73, 0x52, - 0x06, 0x73, 0x74, 0x61, 0x74, 0x75, 0x73, 0x12, 0x2b, 0x0a, 0x12, 0x75, 0x6e, 0x69, 0x78, 0x5f, - 0x74, 0x69, 0x6d, 0x65, 0x5f, 0x73, 0x74, 0x61, 0x72, 0x74, 0x5f, 0x6d, 0x73, 0x18, 0x05, 0x20, - 0x01, 0x28, 0x03, 0x52, 0x0f, 0x75, 0x6e, 0x69, 0x78, 0x54, 0x69, 0x6d, 0x65, 0x53, 0x74, 0x61, - 0x72, 0x74, 0x4d, 0x73, 0x12, 0x27, 0x0a, 0x10, 0x75, 0x6e, 0x69, 0x78, 0x5f, 0x74, 0x69, 0x6d, - 0x65, 0x5f, 0x65, 0x6e, 0x64, 0x5f, 0x6d, 0x73, 0x18, 0x06, 0x20, 0x01, 0x28, 0x03, 0x52, 0x0d, - 0x75, 0x6e, 0x69, 0x78, 0x54, 0x69, 0x6d, 0x65, 0x45, 0x6e, 0x64, 0x4d, 0x73, 0x12, 0x27, 0x0a, - 0x0f, 0x64, 0x69, 0x73, 0x70, 0x6c, 0x61, 0x79, 0x5f, 0x6d, 0x65, 0x73, 0x73, 0x61, 0x67, 0x65, - 0x18, 0x07, 0x20, 0x01, 0x28, 0x09, 0x52, 0x0e, 0x64, 0x69, 0x73, 0x70, 0x6c, 0x61, 0x79, 0x4d, - 0x65, 0x73, 0x73, 0x61, 0x67, 0x65, 0x12, 0x40, 0x0a, 0x10, 0x6f, 0x70, 0x65, 0x72, 0x61, 0x74, - 0x69, 0x6f, 0x6e, 0x5f, 0x62, 0x61, 0x63, 0x6b, 0x75, 0x70, 0x18, 0x64, 0x20, 0x01, 0x28, 0x0b, - 0x32, 0x13, 0x2e, 0x76, 0x31, 0x2e, 0x4f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x42, - 0x61, 0x63, 0x6b, 0x75, 0x70, 0x48, 0x00, 0x52, 0x0f, 0x6f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, - 0x6f, 0x6e, 0x42, 0x61, 0x63, 0x6b, 0x75, 0x70, 0x12, 0x56, 0x0a, 0x18, 0x6f, 0x70, 0x65, 0x72, - 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x5f, 0x69, 0x6e, 0x64, 0x65, 0x78, 0x5f, 0x73, 0x6e, 0x61, 0x70, - 0x73, 0x68, 0x6f, 0x74, 0x18, 0x65, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x1a, 0x2e, 0x76, 0x31, 0x2e, - 0x4f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x49, 0x6e, 0x64, 0x65, 0x78, 0x53, 0x6e, - 0x61, 0x70, 0x73, 0x68, 0x6f, 0x74, 0x48, 0x00, 0x52, 0x16, 0x6f, 0x70, 0x65, 0x72, 0x61, 0x74, - 0x69, 0x6f, 0x6e, 0x49, 0x6e, 0x64, 0x65, 0x78, 0x53, 0x6e, 0x61, 0x70, 0x73, 0x68, 0x6f, 0x74, - 0x42, 0x04, 0x0a, 0x02, 0x6f, 0x70, 0x22, 0x69, 0x0a, 0x0e, 0x4f, 0x70, 0x65, 0x72, 0x61, 0x74, - 0x69, 0x6f, 0x6e, 0x45, 0x76, 0x65, 0x6e, 0x74, 0x12, 0x2a, 0x0a, 0x04, 0x74, 0x79, 0x70, 0x65, - 0x18, 0x01, 0x20, 0x01, 0x28, 0x0e, 0x32, 0x16, 0x2e, 0x76, 0x31, 0x2e, 0x4f, 0x70, 0x65, 0x72, - 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x45, 0x76, 0x65, 0x6e, 0x74, 0x54, 0x79, 0x70, 0x65, 0x52, 0x04, - 0x74, 0x79, 0x70, 0x65, 0x12, 0x2b, 0x0a, 0x09, 0x6f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, - 0x6e, 0x18, 0x02, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x0d, 0x2e, 0x76, 0x31, 0x2e, 0x4f, 0x70, 0x65, - 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x52, 0x09, 0x6f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, - 0x6e, 0x22, 0x4b, 0x0a, 0x0f, 0x4f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x42, 0x61, - 0x63, 0x6b, 0x75, 0x70, 0x12, 0x38, 0x0a, 0x0b, 0x6c, 0x61, 0x73, 0x74, 0x5f, 0x73, 0x74, 0x61, - 0x74, 0x75, 0x73, 0x18, 0x03, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x17, 0x2e, 0x76, 0x31, 0x2e, 0x42, - 0x61, 0x63, 0x6b, 0x75, 0x70, 0x50, 0x72, 0x6f, 0x67, 0x72, 0x65, 0x73, 0x73, 0x45, 0x6e, 0x74, - 0x72, 0x79, 0x52, 0x0a, 0x6c, 0x61, 0x73, 0x74, 0x53, 0x74, 0x61, 0x74, 0x75, 0x73, 0x22, 0x48, - 0x0a, 0x16, 0x4f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x49, 0x6e, 0x64, 0x65, 0x78, - 0x53, 0x6e, 0x61, 0x70, 0x73, 0x68, 0x6f, 0x74, 0x12, 0x2e, 0x0a, 0x08, 0x73, 0x6e, 0x61, 0x70, - 0x73, 0x68, 0x6f, 0x74, 0x18, 0x02, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x12, 0x2e, 0x76, 0x31, 0x2e, - 0x52, 0x65, 0x73, 0x74, 0x69, 0x63, 0x53, 0x6e, 0x61, 0x70, 0x73, 0x68, 0x6f, 0x74, 0x52, 0x08, - 0x73, 0x6e, 0x61, 0x70, 0x73, 0x68, 0x6f, 0x74, 0x2a, 0x4d, 0x0a, 0x12, 0x4f, 0x70, 0x65, 0x72, - 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x45, 0x76, 0x65, 0x6e, 0x74, 0x54, 0x79, 0x70, 0x65, 0x12, 0x11, - 0x0a, 0x0d, 0x45, 0x56, 0x45, 0x4e, 0x54, 0x5f, 0x55, 0x4e, 0x4b, 0x4e, 0x4f, 0x57, 0x4e, 0x10, - 0x00, 0x12, 0x11, 0x0a, 0x0d, 0x45, 0x56, 0x45, 0x4e, 0x54, 0x5f, 0x43, 0x52, 0x45, 0x41, 0x54, - 0x45, 0x44, 0x10, 0x01, 0x12, 0x11, 0x0a, 0x0d, 0x45, 0x56, 0x45, 0x4e, 0x54, 0x5f, 0x55, 0x50, - 0x44, 0x41, 0x54, 0x45, 0x44, 0x10, 0x02, 0x2a, 0x76, 0x0a, 0x0f, 0x4f, 0x70, 0x65, 0x72, 0x61, - 0x74, 0x69, 0x6f, 0x6e, 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, 0x50, 0x45, 0x4e, 0x44, 0x49, 0x4e, 0x47, - 0x10, 0x01, 0x12, 0x15, 0x0a, 0x11, 0x53, 0x54, 0x41, 0x54, 0x55, 0x53, 0x5f, 0x49, 0x4e, 0x50, - 0x52, 0x4f, 0x47, 0x52, 0x45, 0x53, 0x53, 0x10, 0x02, 0x12, 0x12, 0x0a, 0x0e, 0x53, 0x54, 0x41, - 0x54, 0x55, 0x53, 0x5f, 0x53, 0x55, 0x43, 0x43, 0x45, 0x53, 0x53, 0x10, 0x03, 0x12, 0x10, 0x0a, - 0x0c, 0x53, 0x54, 0x41, 0x54, 0x55, 0x53, 0x5f, 0x45, 0x52, 0x52, 0x4f, 0x52, 0x10, 0x04, 0x42, - 0x2e, 0x5a, 0x2c, 0x67, 0x69, 0x74, 0x68, 0x75, 0x62, 0x2e, 0x63, 0x6f, 0x6d, 0x2f, 0x67, 0x61, - 0x72, 0x65, 0x74, 0x68, 0x67, 0x65, 0x6f, 0x72, 0x67, 0x65, 0x2f, 0x72, 0x65, 0x73, 0x74, 0x69, - 0x63, 0x75, 0x69, 0x2f, 0x67, 0x6f, 0x2f, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x2f, 0x76, 0x31, 0x62, - 0x06, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x33, + 0x28, 0x09, 0x52, 0x06, 0x70, 0x6c, 0x61, 0x6e, 0x49, 0x64, 0x12, 0x1f, 0x0a, 0x0b, 0x73, 0x6e, + 0x61, 0x70, 0x73, 0x68, 0x6f, 0x74, 0x5f, 0x69, 0x64, 0x18, 0x08, 0x20, 0x01, 0x28, 0x09, 0x52, + 0x0a, 0x73, 0x6e, 0x61, 0x70, 0x73, 0x68, 0x6f, 0x74, 0x49, 0x64, 0x12, 0x2b, 0x0a, 0x06, 0x73, + 0x74, 0x61, 0x74, 0x75, 0x73, 0x18, 0x04, 0x20, 0x01, 0x28, 0x0e, 0x32, 0x13, 0x2e, 0x76, 0x31, + 0x2e, 0x4f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x53, 0x74, 0x61, 0x74, 0x75, 0x73, + 0x52, 0x06, 0x73, 0x74, 0x61, 0x74, 0x75, 0x73, 0x12, 0x2b, 0x0a, 0x12, 0x75, 0x6e, 0x69, 0x78, + 0x5f, 0x74, 0x69, 0x6d, 0x65, 0x5f, 0x73, 0x74, 0x61, 0x72, 0x74, 0x5f, 0x6d, 0x73, 0x18, 0x05, + 0x20, 0x01, 0x28, 0x03, 0x52, 0x0f, 0x75, 0x6e, 0x69, 0x78, 0x54, 0x69, 0x6d, 0x65, 0x53, 0x74, + 0x61, 0x72, 0x74, 0x4d, 0x73, 0x12, 0x27, 0x0a, 0x10, 0x75, 0x6e, 0x69, 0x78, 0x5f, 0x74, 0x69, + 0x6d, 0x65, 0x5f, 0x65, 0x6e, 0x64, 0x5f, 0x6d, 0x73, 0x18, 0x06, 0x20, 0x01, 0x28, 0x03, 0x52, + 0x0d, 0x75, 0x6e, 0x69, 0x78, 0x54, 0x69, 0x6d, 0x65, 0x45, 0x6e, 0x64, 0x4d, 0x73, 0x12, 0x27, + 0x0a, 0x0f, 0x64, 0x69, 0x73, 0x70, 0x6c, 0x61, 0x79, 0x5f, 0x6d, 0x65, 0x73, 0x73, 0x61, 0x67, + 0x65, 0x18, 0x07, 0x20, 0x01, 0x28, 0x09, 0x52, 0x0e, 0x64, 0x69, 0x73, 0x70, 0x6c, 0x61, 0x79, + 0x4d, 0x65, 0x73, 0x73, 0x61, 0x67, 0x65, 0x12, 0x40, 0x0a, 0x10, 0x6f, 0x70, 0x65, 0x72, 0x61, + 0x74, 0x69, 0x6f, 0x6e, 0x5f, 0x62, 0x61, 0x63, 0x6b, 0x75, 0x70, 0x18, 0x64, 0x20, 0x01, 0x28, + 0x0b, 0x32, 0x13, 0x2e, 0x76, 0x31, 0x2e, 0x4f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, + 0x42, 0x61, 0x63, 0x6b, 0x75, 0x70, 0x48, 0x00, 0x52, 0x0f, 0x6f, 0x70, 0x65, 0x72, 0x61, 0x74, + 0x69, 0x6f, 0x6e, 0x42, 0x61, 0x63, 0x6b, 0x75, 0x70, 0x12, 0x56, 0x0a, 0x18, 0x6f, 0x70, 0x65, + 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x5f, 0x69, 0x6e, 0x64, 0x65, 0x78, 0x5f, 0x73, 0x6e, 0x61, + 0x70, 0x73, 0x68, 0x6f, 0x74, 0x18, 0x65, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x1a, 0x2e, 0x76, 0x31, + 0x2e, 0x4f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x49, 0x6e, 0x64, 0x65, 0x78, 0x53, + 0x6e, 0x61, 0x70, 0x73, 0x68, 0x6f, 0x74, 0x48, 0x00, 0x52, 0x16, 0x6f, 0x70, 0x65, 0x72, 0x61, + 0x74, 0x69, 0x6f, 0x6e, 0x49, 0x6e, 0x64, 0x65, 0x78, 0x53, 0x6e, 0x61, 0x70, 0x73, 0x68, 0x6f, + 0x74, 0x42, 0x04, 0x0a, 0x02, 0x6f, 0x70, 0x22, 0x69, 0x0a, 0x0e, 0x4f, 0x70, 0x65, 0x72, 0x61, + 0x74, 0x69, 0x6f, 0x6e, 0x45, 0x76, 0x65, 0x6e, 0x74, 0x12, 0x2a, 0x0a, 0x04, 0x74, 0x79, 0x70, + 0x65, 0x18, 0x01, 0x20, 0x01, 0x28, 0x0e, 0x32, 0x16, 0x2e, 0x76, 0x31, 0x2e, 0x4f, 0x70, 0x65, + 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x45, 0x76, 0x65, 0x6e, 0x74, 0x54, 0x79, 0x70, 0x65, 0x52, + 0x04, 0x74, 0x79, 0x70, 0x65, 0x12, 0x2b, 0x0a, 0x09, 0x6f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, + 0x6f, 0x6e, 0x18, 0x02, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x0d, 0x2e, 0x76, 0x31, 0x2e, 0x4f, 0x70, + 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x52, 0x09, 0x6f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, + 0x6f, 0x6e, 0x22, 0x4b, 0x0a, 0x0f, 0x4f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x42, + 0x61, 0x63, 0x6b, 0x75, 0x70, 0x12, 0x38, 0x0a, 0x0b, 0x6c, 0x61, 0x73, 0x74, 0x5f, 0x73, 0x74, + 0x61, 0x74, 0x75, 0x73, 0x18, 0x03, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x17, 0x2e, 0x76, 0x31, 0x2e, + 0x42, 0x61, 0x63, 0x6b, 0x75, 0x70, 0x50, 0x72, 0x6f, 0x67, 0x72, 0x65, 0x73, 0x73, 0x45, 0x6e, + 0x74, 0x72, 0x79, 0x52, 0x0a, 0x6c, 0x61, 0x73, 0x74, 0x53, 0x74, 0x61, 0x74, 0x75, 0x73, 0x22, + 0x48, 0x0a, 0x16, 0x4f, 0x70, 0x65, 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x49, 0x6e, 0x64, 0x65, + 0x78, 0x53, 0x6e, 0x61, 0x70, 0x73, 0x68, 0x6f, 0x74, 0x12, 0x2e, 0x0a, 0x08, 0x73, 0x6e, 0x61, + 0x70, 0x73, 0x68, 0x6f, 0x74, 0x18, 0x02, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x12, 0x2e, 0x76, 0x31, + 0x2e, 0x52, 0x65, 0x73, 0x74, 0x69, 0x63, 0x53, 0x6e, 0x61, 0x70, 0x73, 0x68, 0x6f, 0x74, 0x52, + 0x08, 0x73, 0x6e, 0x61, 0x70, 0x73, 0x68, 0x6f, 0x74, 0x2a, 0x4d, 0x0a, 0x12, 0x4f, 0x70, 0x65, + 0x72, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x45, 0x76, 0x65, 0x6e, 0x74, 0x54, 0x79, 0x70, 0x65, 0x12, + 0x11, 0x0a, 0x0d, 0x45, 0x56, 0x45, 0x4e, 0x54, 0x5f, 0x55, 0x4e, 0x4b, 0x4e, 0x4f, 0x57, 0x4e, + 0x10, 0x00, 0x12, 0x11, 0x0a, 0x0d, 0x45, 0x56, 0x45, 0x4e, 0x54, 0x5f, 0x43, 0x52, 0x45, 0x41, + 0x54, 0x45, 0x44, 0x10, 0x01, 0x12, 0x11, 0x0a, 0x0d, 0x45, 0x56, 0x45, 0x4e, 0x54, 0x5f, 0x55, + 0x50, 0x44, 0x41, 0x54, 0x45, 0x44, 0x10, 0x02, 0x2a, 0x76, 0x0a, 0x0f, 0x4f, 0x70, 0x65, 0x72, + 0x61, 0x74, 0x69, 0x6f, 0x6e, 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, 0x50, 0x45, 0x4e, 0x44, 0x49, 0x4e, + 0x47, 0x10, 0x01, 0x12, 0x15, 0x0a, 0x11, 0x53, 0x54, 0x41, 0x54, 0x55, 0x53, 0x5f, 0x49, 0x4e, + 0x50, 0x52, 0x4f, 0x47, 0x52, 0x45, 0x53, 0x53, 0x10, 0x02, 0x12, 0x12, 0x0a, 0x0e, 0x53, 0x54, + 0x41, 0x54, 0x55, 0x53, 0x5f, 0x53, 0x55, 0x43, 0x43, 0x45, 0x53, 0x53, 0x10, 0x03, 0x12, 0x10, + 0x0a, 0x0c, 0x53, 0x54, 0x41, 0x54, 0x55, 0x53, 0x5f, 0x45, 0x52, 0x52, 0x4f, 0x52, 0x10, 0x04, + 0x42, 0x2e, 0x5a, 0x2c, 0x67, 0x69, 0x74, 0x68, 0x75, 0x62, 0x2e, 0x63, 0x6f, 0x6d, 0x2f, 0x67, + 0x61, 0x72, 0x65, 0x74, 0x68, 0x67, 0x65, 0x6f, 0x72, 0x67, 0x65, 0x2f, 0x72, 0x65, 0x73, 0x74, + 0x69, 0x63, 0x75, 0x69, 0x2f, 0x67, 0x6f, 0x2f, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x2f, 0x76, 0x31, + 0x62, 0x06, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x33, } var ( diff --git a/internal/api/server.go b/internal/api/server.go index 9027daa7..0acfb58e 100644 --- a/internal/api/server.go +++ b/internal/api/server.go @@ -10,7 +10,8 @@ import ( "github.com/garethgeorge/resticui/gen/go/types" v1 "github.com/garethgeorge/resticui/gen/go/v1" "github.com/garethgeorge/resticui/internal/config" - "github.com/garethgeorge/resticui/internal/database/oplog" + "github.com/garethgeorge/resticui/internal/oplog" + "github.com/garethgeorge/resticui/internal/oplog/indexutil" "github.com/garethgeorge/resticui/internal/orchestrator" "github.com/garethgeorge/resticui/pkg/restic" "go.uber.org/zap" @@ -185,10 +186,10 @@ func (s *Server) GetOperationEvents(_ *emptypb.Empty, stream v1.ResticUI_GetOper } func (s *Server) GetOperations(ctx context.Context, req *v1.GetOperationsRequest) (*v1.OperationList, error) { - filter := oplog.FilterKeepAll() + collector := indexutil.CollectAll() if req.LastN != 0 { - filter = oplog.FilterLastN(req.LastN) + collector = indexutil.CollectLastN(int(req.LastN)) } var err error @@ -196,11 +197,11 @@ func (s *Server) GetOperations(ctx context.Context, req *v1.GetOperationsRequest if req.RepoId != "" && req.PlanId != "" { return nil, errors.New("cannot specify both repoId and planId") } else if req.PlanId != "" { - ops, err = s.oplog.GetByPlan(req.PlanId, filter) + ops, err = s.oplog.GetByPlan(req.PlanId, collector) } else if req.RepoId != "" { - ops, err = s.oplog.GetByRepo(req.RepoId, filter) + ops, err = s.oplog.GetByRepo(req.RepoId, collector) } else { - ops, err = s.oplog.GetAll(filter) + ops, err = s.oplog.GetAll() } if err != nil { return nil, fmt.Errorf("failed to get operations: %w", err) diff --git a/internal/database/indexutil/indexutil.go b/internal/oplog/indexutil/indexutil.go similarity index 54% rename from internal/database/indexutil/indexutil.go rename to internal/oplog/indexutil/indexutil.go index 4cfdd0da..12dbd962 100644 --- a/internal/database/indexutil/indexutil.go +++ b/internal/oplog/indexutil/indexutil.go @@ -2,8 +2,9 @@ package indexutil import ( "bytes" + "sort" - "github.com/garethgeorge/resticui/internal/database/serializationutil" + "github.com/garethgeorge/resticui/internal/oplog/serializationutil" bolt "go.etcd.io/bbolt" ) @@ -14,14 +15,20 @@ func IndexByteValue(b *bolt.Bucket, value []byte, recordId int64) error { return b.Put(key, []byte{}) } +func IndexRemoveByteValue(b *bolt.Bucket, value []byte, recordId int64) error { + key := serializationutil.BytesToKey(value) + key = append(key, serializationutil.Itob(recordId)...) + return b.Delete(key) +} + // IndexSearchByteValue searches the index given a value and returns an iterator over the associated recordIds. func IndexSearchByteValue(b *bolt.Bucket, value []byte) *IndexSearchIterator { return newSearchIterator(b, serializationutil.BytesToKey(value)) } type IndexSearchIterator struct { - c *bolt.Cursor - k []byte + c *bolt.Cursor + k []byte prefix []byte } @@ -29,8 +36,8 @@ func newSearchIterator(b *bolt.Bucket, prefix []byte) *IndexSearchIterator { c := b.Cursor() k, _ := c.Seek(prefix) return &IndexSearchIterator{ - c: c, - k: k, + c: c, + k: k, prefix: prefix, } } @@ -54,4 +61,40 @@ func (i *IndexSearchIterator) ToSlice() []int64 { ids = append(ids, id) } return ids -} \ No newline at end of file +} + +type Collector func(*IndexSearchIterator) []int64 + +func CollectAll() Collector { + return func(iter *IndexSearchIterator) []int64 { + return iter.ToSlice() + } +} + +func CollectFirstN(firstN int) Collector { + return func(iter *IndexSearchIterator) []int64 { + ids := make([]int64, 0, firstN) + for id, ok := iter.Next(); ok && len(ids) < firstN; id, ok = iter.Next() { + ids = append(ids, id) + } + sort.Slice(ids, func(i, j int) bool { + return ids[i] < ids[j] + }) + return ids + } +} + +func CollectLastN(lastN int) Collector { + return func(iter *IndexSearchIterator) []int64 { + ids := make([]int64, lastN) + count := 0 + for id, ok := iter.Next(); ok; id, ok = iter.Next() { + ids[count%lastN] = id + count += 1 + } + if count < lastN { + return ids[:count] + } + return ids + } +} diff --git a/internal/database/indexutil/indexutil_test.go b/internal/oplog/indexutil/indexutil_test.go similarity index 100% rename from internal/database/indexutil/indexutil_test.go rename to internal/oplog/indexutil/indexutil_test.go diff --git a/internal/database/oplog/oplog.go b/internal/oplog/oplog.go similarity index 57% rename from internal/database/oplog/oplog.go rename to internal/oplog/oplog.go index 1459750c..d5843df6 100644 --- a/internal/database/oplog/oplog.go +++ b/internal/oplog/oplog.go @@ -6,11 +6,12 @@ import ( "os" "path" "sync" + "sync/atomic" "time" v1 "github.com/garethgeorge/resticui/gen/go/v1" - "github.com/garethgeorge/resticui/internal/database/indexutil" - "github.com/garethgeorge/resticui/internal/database/serializationutil" + "github.com/garethgeorge/resticui/internal/oplog/indexutil" + "github.com/garethgeorge/resticui/internal/oplog/serializationutil" bolt "go.etcd.io/bbolt" "go.uber.org/zap" "google.golang.org/protobuf/proto" @@ -24,12 +25,16 @@ const ( EventTypeOpUpdated = EventType(iota) ) +const ( + schemaVersion int64 = 1 +) + var ( - SystemBucket = []byte("oplog.system") // system stores metadata - OpLogBucket = []byte("oplog.log") // oplog stores the operations themselves - RepoIndexBucket = []byte("oplog.repo_idx") // repo_index tracks IDs of operations affecting a given repo - PlanIndexBucket = []byte("oplog.plan_idx") // plan_index tracks IDs of operations affecting a given plan - IndexedSnapshotsSetBucket = []byte("oplog.indexed_snapshots") // indexed_snapshots is a set of snapshot IDs that have been indexed + SystemBucket = []byte("oplog.system") // system stores metadata + OpLogBucket = []byte("oplog.log") // oplog stores the operations themselves + RepoIndexBucket = []byte("oplog.repo_idx") // repo_index tracks IDs of operations affecting a given repo + PlanIndexBucket = []byte("oplog.plan_idx") // plan_index tracks IDs of operations affecting a given plan + SnapshotIndexBucket = []byte("oplog.snapshot_idx") // snapshot_index tracks IDs of operations affecting a given snapshot ) // OpLog represents a log of operations performed. @@ -39,6 +44,7 @@ type OpLog struct { subscribersMu sync.RWMutex subscribers []*func(EventType, *v1.Operation) + nextId atomic.Int64 } func NewOpLog(databasePath string) (*OpLog, error) { @@ -51,18 +57,21 @@ func NewOpLog(databasePath string) (*OpLog, error) { return nil, fmt.Errorf("error opening database: %s", err) } + o := &OpLog{db: db} + if err := db.Update(func(tx *bolt.Tx) error { // Create the buckets if they don't exist for _, bucket := range [][]byte{ - SystemBucket, OpLogBucket, RepoIndexBucket, PlanIndexBucket, IndexedSnapshotsSetBucket, + SystemBucket, OpLogBucket, RepoIndexBucket, PlanIndexBucket, SnapshotIndexBucket, } { if _, err := tx.CreateBucketIfNotExists(bucket); err != nil { return fmt.Errorf("creating bucket %s: %s", string(bucket), err) } } - // Validate the operation log on startup. sysBucket := tx.Bucket(SystemBucket) + + // Validate the operation log on startup. opLogBucket := tx.Bucket(OpLogBucket) c := opLogBucket.Cursor() if lastValidated := sysBucket.Get([]byte("last_validated")); lastValidated != nil { @@ -77,13 +86,11 @@ func NewOpLog(databasePath string) (*OpLog, error) { if op.Status == v1.OperationStatus_STATUS_INPROGRESS { op.Status = v1.OperationStatus_STATUS_ERROR op.DisplayMessage = "Operation timeout." - bytes, err := proto.Marshal(op) - if err != nil { - return fmt.Errorf("marshalling operation: %w", err) - } - if err := opLogBucket.Put(k, bytes); err != nil { - return fmt.Errorf("putting operation into bucket: %w", err) - } + } + + if err := o.addOperationHelper(tx, op); err != nil { + zap.L().Error("error re-adding operation, there may be corruption in the oplog", zap.Error(err)) + continue } } if lastValidated, _ := c.Last(); lastValidated != nil { @@ -97,7 +104,7 @@ func NewOpLog(databasePath string) (*OpLog, error) { return nil, err } - return &OpLog{db: db}, nil + return o, nil } func (o *OpLog) Close() error { @@ -126,6 +133,9 @@ func (o *OpLog) Add(op *v1.Operation) error { func (o *OpLog) BulkAdd(ops []*v1.Operation) error { err := o.db.Update(func(tx *bolt.Tx) error { for _, op := range ops { + if op.Id != 0 { + return errors.New("operation already has an ID, OpLog.BulkAdd is expected to set the ID") + } if err := o.addOperationHelper(tx, op); err != nil { return err } @@ -140,15 +150,62 @@ func (o *OpLog) BulkAdd(ops []*v1.Operation) error { return err } +func (o *OpLog) Update(op *v1.Operation) error { + if op.Id == 0 { + return errors.New("operation does not have an ID, OpLog.Update expects operation with an ID") + } + + err := o.db.Update(func(tx *bolt.Tx) error { + if err := o.deleteOperationHelper(tx, op.Id); err != nil { + return fmt.Errorf("deleting existing value prior to update: %w", err) + } + if err := o.addOperationHelper(tx, op); err != nil { + return fmt.Errorf("adding updated value: %w", err) + } + return nil + }) + if err == nil { + o.notifyHelper(EventTypeOpUpdated, op) + } + return err +} + +func (o *OpLog) notifyHelper(eventType EventType, op *v1.Operation) { + o.subscribersMu.RLock() + defer o.subscribersMu.RUnlock() + for _, sub := range o.subscribers { + (*sub)(eventType, op) + } +} + +func (o *OpLog) getOperationHelper(b *bolt.Bucket, id int64) (*v1.Operation, error) { + bytes := b.Get(serializationutil.Itob(id)) + if bytes == nil { + return nil, fmt.Errorf("operation with ID %d does not exist", id) + } + + var op v1.Operation + if err := proto.Unmarshal(bytes, &op); err != nil { + return nil, fmt.Errorf("error unmarshalling operation: %w", err) + } + + return &op, nil +} + func (o *OpLog) addOperationHelper(tx *bolt.Tx, op *v1.Operation) error { b := tx.Bucket(OpLogBucket) - id, err := b.NextSequence() - if err != nil { - return fmt.Errorf("error getting next sequence: %w", err) + if op.Id == 0 { + // Create a unique ID sorted based on the start time in milliseconds and + // a counter to ensure uniqueness in the case of multiple operations + // starting at the same time. + op.Id = op.UnixTimeStartMs<<20 | (o.nextId.Add(1) & (1<<20 - 1)) + if op.Id < 0 { + return fmt.Errorf("overflow in operation ID generation") + } } - op.Id = int64(id) + op.SnapshotId = NormalizeSnapshotId(op.SnapshotId) bytes, err := proto.Marshal(op) if err != nil { @@ -170,105 +227,53 @@ func (o *OpLog) addOperationHelper(tx *bolt.Tx, op *v1.Operation) error { return fmt.Errorf("error adding operation to repo index: %w", err) } } - - // Update operation type dependent indices. - switch wrappedOp := op.Op.(type) { - case *v1.Operation_OperationBackup: - // Nothing extra to be done. - case *v1.Operation_OperationIndexSnapshot: - if wrappedOp.OperationIndexSnapshot == nil || wrappedOp.OperationIndexSnapshot.Snapshot == nil { - return errors.New("op.OperationIndexSnapshot or op.OperationIndexSnapshot.Snapshot is nil") + if op.SnapshotId != "" { + if err := indexutil.IndexByteValue(tx.Bucket(SnapshotIndexBucket), []byte(op.SnapshotId), op.Id); err != nil { + return fmt.Errorf("error adding operation to snapshot index: %w", err) } - snapshotId := serializationutil.NormalizeSnapshotId(wrappedOp.OperationIndexSnapshot.Snapshot.Id) - key := serializationutil.BytesToKey([]byte(snapshotId)) - if err := tx.Bucket(IndexedSnapshotsSetBucket).Put(key, serializationutil.Itob(op.Id)); err != nil { - return fmt.Errorf("error adding OperationIndexSnapshot to indexed snapshots set: %w", err) - } - default: - return fmt.Errorf("unknown operation type: %T", wrappedOp) } return nil } -func (o *OpLog) HasIndexedSnapshot(snapshotId string) (int64, error) { - var id int64 - if err := o.db.View(func(tx *bolt.Tx) error { - snapshotId := serializationutil.NormalizeSnapshotId(snapshotId) - key := serializationutil.BytesToKey([]byte(snapshotId)) - idBytes := tx.Bucket(IndexedSnapshotsSetBucket).Get(key) - if idBytes == nil { - id = -1 - } else { - var err error - id, err = serializationutil.Btoi(idBytes) - if err != nil { - return fmt.Errorf("database corrupt, couldn't convert ID bytes to int: %w", err) - } +func (o *OpLog) deleteOperationHelper(tx *bolt.Tx, id int64) error { + b := tx.Bucket(OpLogBucket) + + prevValue, err := o.getOperationHelper(b, id) + if err != nil { + return fmt.Errorf("getting operation %v: %w", id, err) + } + + if prevValue.PlanId != "" { + if err := indexutil.IndexRemoveByteValue(tx.Bucket(PlanIndexBucket), []byte(prevValue.PlanId), id); err != nil { + return fmt.Errorf("removing operation %v from plan index: %w", id, err) } - return nil - }); err != nil { - return 0, err - } - return id, nil -} - -func (o *OpLog) Update(op *v1.Operation) error { - if op.Id == 0 { - return errors.New("operation does not have an ID, OpLog.Update expects operation with an ID") } - err := o.db.Update(func(tx *bolt.Tx) error { - b := tx.Bucket(OpLogBucket) - - if b.Get(serializationutil.Itob(op.Id)) == nil { - return fmt.Errorf("operation with ID %d does not exist", op.Id) + if prevValue.RepoId != "" { + if err := indexutil.IndexRemoveByteValue(tx.Bucket(RepoIndexBucket), []byte(prevValue.RepoId), id); err != nil { + return fmt.Errorf("removing operation %v from repo index: %w", id, err) } + } - bytes, err := proto.Marshal(op) - if err != nil { - return fmt.Errorf("error marshalling operation: %w", err) + if prevValue.SnapshotId != "" { + if err := indexutil.IndexRemoveByteValue(tx.Bucket(SnapshotIndexBucket), []byte(prevValue.SnapshotId), id); err != nil { + return fmt.Errorf("removing operation %v from snapshot index: %w", id, err) } - - if err := b.Put(serializationutil.Itob(op.Id), bytes); err != nil { - return fmt.Errorf("error putting operation into bucket: %w", err) - } - - return nil - }) - if err == nil { - o.notifyHelper(EventTypeOpUpdated, op) - } - return err -} - -func (o *OpLog) notifyHelper(eventType EventType, op *v1.Operation) { - o.subscribersMu.RLock() - defer o.subscribersMu.RUnlock() - for _, sub := range o.subscribers { - (*sub)(eventType, op) - } -} - -func (o *OpLog) getHelper(b *bolt.Bucket, id int64) (*v1.Operation, error) { - bytes := b.Get(serializationutil.Itob(id)) - if bytes == nil { - return nil, fmt.Errorf("operation with ID %d does not exist", id) } - var op v1.Operation - if err := proto.Unmarshal(bytes, &op); err != nil { - return nil, fmt.Errorf("error unmarshalling operation: %w", err) + if err := b.Delete(serializationutil.Itob(id)); err != nil { + return fmt.Errorf("deleting operation %v from bucket: %w", id, err) } - return &op, nil + return nil } func (o *OpLog) Get(id int64) (*v1.Operation, error) { var op *v1.Operation if err := o.db.View(func(tx *bolt.Tx) error { var err error - op, err = o.getHelper(tx.Bucket(OpLogBucket), id) + op, err = o.getOperationHelper(tx.Bucket(OpLogBucket), id) return err }); err != nil { return nil, err @@ -276,51 +281,54 @@ func (o *OpLog) Get(id int64) (*v1.Operation, error) { return op, nil } -func (o *OpLog) GetByRepo(repoId string, filter Filter) ([]*v1.Operation, error) { +func (o *OpLog) GetByRepo(repoId string, collector indexutil.Collector) ([]*v1.Operation, error) { + var err error var ops []*v1.Operation - if err := o.db.View(func(tx *bolt.Tx) error { - ids := indexutil.IndexSearchByteValue(tx.Bucket(RepoIndexBucket), []byte(repoId)).ToSlice() - ids = filter(ids) - - b := tx.Bucket(OpLogBucket) - for _, id := range ids { - op, err := o.getHelper(b, id) - if err != nil { - return err - } - ops = append(ops, op) - } - + o.db.View(func(tx *bolt.Tx) error { + ids := collector(indexutil.IndexSearchByteValue(tx.Bucket(RepoIndexBucket), []byte(repoId))) + ops, err = o.getOpsByIds(tx, ids) return nil - }); err != nil { - return nil, err + }) + return ops, err +} + +func (o *OpLog) GetByPlan(planId string, collector indexutil.Collector) ([]*v1.Operation, error) { + var err error + var ops []*v1.Operation + o.db.View(func(tx *bolt.Tx) error { + ids := collector(indexutil.IndexSearchByteValue(tx.Bucket(PlanIndexBucket), []byte(planId))) + ops, err = o.getOpsByIds(tx, ids) + return nil + }) + return ops, err +} + +func (o *OpLog) GetBySnapshotId(snapshotId string, collector indexutil.Collector) ([]*v1.Operation, error) { + snapshotId = NormalizeSnapshotId(snapshotId) + var err error + var ops []*v1.Operation + o.db.View(func(tx *bolt.Tx) error { + ids := collector(indexutil.IndexSearchByteValue(tx.Bucket(SnapshotIndexBucket), []byte(snapshotId))) + ops, err = o.getOpsByIds(tx, ids) + return nil + }) + return ops, err +} + +func (o *OpLog) getOpsByIds(tx *bolt.Tx, ids []int64) ([]*v1.Operation, error) { + b := tx.Bucket(OpLogBucket) + ops := make([]*v1.Operation, 0, len(ids)) + for _, id := range ids { + op, err := o.getOperationHelper(b, id) + if err != nil { + return nil, err + } + ops = append(ops, op) } return ops, nil } -func (o *OpLog) GetByPlan(planId string, filter Filter) ([]*v1.Operation, error) { - var ops []*v1.Operation - if err := o.db.View(func(tx *bolt.Tx) error { - ids := indexutil.IndexSearchByteValue(tx.Bucket(PlanIndexBucket), []byte(planId)).ToSlice() - ids = filter(ids) - - b := tx.Bucket(OpLogBucket) - for _, id := range ids { - op, err := o.getHelper(b, id) - if err != nil { - return err - } - ops = append(ops, op) - } - - return nil - }); err != nil { - return nil, err - } - return ops, nil -} - -func (o *OpLog) GetAll(filter Filter) ([]*v1.Operation, error) { +func (o *OpLog) GetAll() ([]*v1.Operation, error) { var ops []*v1.Operation if err := o.db.View(func(tx *bolt.Tx) error { c := tx.Bucket(OpLogBucket).Cursor() @@ -356,31 +364,9 @@ func (o *OpLog) Unsubscribe(callback *func(EventType, *v1.Operation)) { } } -type Filter func([]int64) []int64 - -func FilterKeepAll() Filter { - return func(ids []int64) []int64 { - return ids - } -} - -func FilterLastN(n int64) Filter { - return func(ids []int64) []int64 { - if len(ids) > int(n) { - ids = ids[len(ids)-int(n):] - } - return ids - } -} - -func FilterLimitOffset(limit, offset int64) Filter { - return func(ids []int64) []int64 { - if len(ids) > int(offset) { - ids = ids[offset:] - } - if len(ids) > int(limit) { - ids = ids[:limit] - } - return ids +func NormalizeSnapshotId(id string) string { + if len(id) < 8 { + return id } + return id[:8] } diff --git a/internal/database/oplog/oplog_test.go b/internal/oplog/oplog_test.go similarity index 57% rename from internal/database/oplog/oplog_test.go rename to internal/oplog/oplog_test.go index 60d60d6e..d1141737 100644 --- a/internal/database/oplog/oplog_test.go +++ b/internal/oplog/oplog_test.go @@ -5,6 +5,7 @@ import ( "testing" v1 "github.com/garethgeorge/resticui/gen/go/v1" + "github.com/garethgeorge/resticui/internal/oplog/indexutil" ) func TestCreate(t *testing.T) { @@ -33,16 +34,13 @@ func TestAddOperation(t *testing.T) { wantErr bool }{ { - name: "no operation", - op: &v1.Operation{ - Id: 0, - }, - wantErr: true, + name: "basic operation", + op: &v1.Operation{}, + wantErr: false, }, { name: "basic backup operation", op: &v1.Operation{ - Id: 0, Op: &v1.Operation_OperationBackup{}, }, wantErr: false, @@ -61,16 +59,6 @@ func TestAddOperation(t *testing.T) { }, wantErr: false, }, - { - name: "basic snapshot operation with no snapshot", - op: &v1.Operation{ - Id: 0, - Op: &v1.Operation_OperationIndexSnapshot{ - OperationIndexSnapshot: &v1.OperationIndexSnapshot{}, - }, - }, - wantErr: true, - }, { name: "operation with ID", op: &v1.Operation{ @@ -187,9 +175,9 @@ func TestListOperation(t *testing.T) { var ops []*v1.Operation var err error if tc.byPlan { - ops, err = log.GetByPlan(tc.id, FilterKeepAll()) + ops, err = log.GetByPlan(tc.id, indexutil.CollectAll()) } else if tc.byRepo { - ops, err = log.GetByRepo(tc.id, FilterKeepAll()) + ops, err = log.GetByRepo(tc.id, indexutil.CollectAll()) } else { t.Fatalf("must specify byPlan or byRepo") } @@ -205,14 +193,17 @@ func TestListOperation(t *testing.T) { } func TestBigIO(t *testing.T) { - // t.Parallel() + t.Parallel() + + count := 10 + log, err := NewOpLog(t.TempDir() + "/test.boltdb") if err != nil { t.Fatalf("error creating oplog: %s", err) } t.Cleanup(func() { log.Close() }) - for i := 0; i < 100; i++ { + for i := 0; i < count; i++ { if err := log.Add(&v1.Operation{ PlanId: "plan1", RepoId: "repo1", @@ -222,20 +213,20 @@ func TestBigIO(t *testing.T) { } } - ops, err := log.GetByPlan("plan1", FilterKeepAll()) + ops, err := log.GetByPlan("plan1", indexutil.CollectAll()) if err != nil { t.Fatalf("error listing operations: %s", err) } - if len(ops) != 100 { - t.Errorf("want 100 operations, got %d", len(ops)) + if len(ops) != count { + t.Errorf("want %v operations, got %d", count, len(ops)) } - ops, err = log.GetByRepo("repo1", FilterKeepAll()) + ops, err = log.GetByRepo("repo1", indexutil.CollectAll()) if err != nil { t.Fatalf("error listing operations: %s", err) } - if len(ops) != 100 { - t.Errorf("want 100 operations, got %d", len(ops)) + if len(ops) != count { + t.Errorf("want %v operations, got %d", count, len(ops)) } } @@ -248,34 +239,108 @@ func TestIndexSnapshot(t *testing.T) { t.Cleanup(func() { log.Close() }) op := &v1.Operation{ - PlanId: "plan1", - RepoId: "repo1", - Op: &v1.Operation_OperationIndexSnapshot{ - OperationIndexSnapshot: &v1.OperationIndexSnapshot{ - Snapshot: &v1.ResticSnapshot{ - Id: "abcdefghijklmnop", - }, - }, - }, + PlanId: "plan1", + RepoId: "repo1", + SnapshotId: "abcdefgh", + Op: &v1.Operation_OperationIndexSnapshot{}, } if err := log.Add(op); err != nil { t.Fatalf("error adding operation: %s", err) } - id, err := log.HasIndexedSnapshot("abcdefgh") + ops, err := log.GetBySnapshotId("abcdefgh", indexutil.CollectAll()) if err != nil { t.Fatalf("error checking for snapshot: %s", err) } - if id != op.Id { - t.Fatalf("want id %d, got %d", op.Id, id) + if len(ops) != 1 { + t.Fatalf("want 1 operation, got %d", len(ops)) + } + if ops[0].Id != op.Id { + t.Errorf("want operation ID %d, got %d", op.Id, ops[0].Id) + } +} + +func TestUpdateOperation(t *testing.T) { + t.Parallel() + log, err := NewOpLog(t.TempDir() + "/test.boltdb") + if err != nil { + t.Fatalf("error creating oplog: %s", err) + } + t.Cleanup(func() { log.Close() }) + + // Insert initial operation + op := &v1.Operation{ + PlanId: "oldplan", + RepoId: "oldrepo", + SnapshotId: "12345678", + } + if err := log.Add(op); err != nil { + t.Fatalf("error adding operation: %s", err) + } + opId := op.Id + + // Validate initial values are indexed + if ops, err := log.GetByPlan("oldplan", indexutil.CollectAll()); err != nil { + t.Fatalf("error checking for plan: %s", err) + } else if len(ops) != 1 { + t.Fatalf("want 1 operation, got %d", len(ops)) + } + if ops, err := log.GetByRepo("oldrepo", indexutil.CollectAll()); err != nil { + t.Fatalf("error checking for repo: %s", err) + } else if len(ops) != 1 { + t.Fatalf("want 1 operation, got %d", len(ops)) + } + if ops, err := log.GetBySnapshotId("12345678", indexutil.CollectAll()); err != nil { + t.Fatalf("error checking for snapshot: %s", err) + } else if len(ops) != 1 { + t.Fatalf("want 1 operation, got %d", len(ops)) } - id, err = log.HasIndexedSnapshot("notfound") - if err != nil { - t.Fatalf("error checking for snapshot: %s", err) + // Update indexed values + op.SnapshotId = "abcdefgh" + op.PlanId = "myplan" + op.RepoId = "myrepo" + if err := log.Update(op); err != nil { + t.Fatalf("error updating operation: %s", err) } - if id != -1 { - t.Fatalf("want id -1, got %d", id) + + // Validate updated values are indexed + if opId != op.Id { + t.Errorf("want operation ID %d, got %d", opId, op.Id) + } + if ops, err := log.GetBySnapshotId("abcdefgh", indexutil.CollectAll()); err != nil { + t.Fatalf("error checking for snapshot: %s", err) + } else if len(ops) != 1 { + t.Fatalf("want 1 operation, got %d", len(ops)) + } + + if ops, err := log.GetByPlan("myplan", indexutil.CollectAll()); err != nil { + t.Fatalf("error checking for plan: %s", err) + } else if len(ops) != 1 { + t.Fatalf("want 1 operation, got %d", len(ops)) + } + + if ops, err := log.GetByRepo("myrepo", indexutil.CollectAll()); err != nil { + t.Fatalf("error checking for repo: %s", err) + } else if len(ops) != 1 { + t.Fatalf("want 1 operation, got %d", len(ops)) + } + + // Validate prior values are gone + if ops, err := log.GetByPlan("oldplan", indexutil.CollectAll()); err != nil { + t.Fatalf("error checking for plan: %s", err) + } else if len(ops) != 0 { + t.Fatalf("want 0 operations, got %d", len(ops)) + } + if ops, err := log.GetByRepo("oldrepo", indexutil.CollectAll()); err != nil { + t.Fatalf("error checking for repo: %s", err) + } else if len(ops) != 0 { + t.Fatalf("want 0 operations, got %d", len(ops)) + } + if ops, err := log.GetBySnapshotId("12345678", indexutil.CollectAll()); err != nil { + t.Fatalf("error checking for snapshot: %s", err) + } else if len(ops) != 0 { + t.Fatalf("want 0 operations, got %d", len(ops)) } } diff --git a/internal/database/serializationutil/serializationutil.go b/internal/oplog/serializationutil/serializationutil.go similarity index 85% rename from internal/database/serializationutil/serializationutil.go rename to internal/oplog/serializationutil/serializationutil.go index 62b5db90..d1961895 100644 --- a/internal/database/serializationutil/serializationutil.go +++ b/internal/oplog/serializationutil/serializationutil.go @@ -35,7 +35,7 @@ func Btos(b []byte) (string, int64, error) { if int64(len(b)) < 8+length { return "", 0, ErrInvalidLength } - return string(b[8:8+length]), 8+length, nil + return string(b[8 : 8+length]), 8 + length, nil } func BytesToKey(b []byte) []byte { @@ -44,10 +44,3 @@ func BytesToKey(b []byte) []byte { key = append(key, b...) return key } - -func NormalizeSnapshotId(id string) string { - if len(id) < 8 { - return id - } - return id[:8] -} \ No newline at end of file diff --git a/internal/database/serializationutil/serializationutil_test.go b/internal/oplog/serializationutil/serializationutil_test.go similarity index 100% rename from internal/database/serializationutil/serializationutil_test.go rename to internal/oplog/serializationutil/serializationutil_test.go diff --git a/internal/orchestrator/orchestrator.go b/internal/orchestrator/orchestrator.go index 3aa1d617..220a060e 100644 --- a/internal/orchestrator/orchestrator.go +++ b/internal/orchestrator/orchestrator.go @@ -9,7 +9,7 @@ import ( v1 "github.com/garethgeorge/resticui/gen/go/v1" "github.com/garethgeorge/resticui/internal/config" - "github.com/garethgeorge/resticui/internal/database/oplog" + "github.com/garethgeorge/resticui/internal/oplog" "github.com/garethgeorge/resticui/pkg/restic" "go.uber.org/zap" "google.golang.org/protobuf/proto" @@ -21,14 +21,13 @@ var ErrPlanNotFound = errors.New("plan not found") // Orchestrator is responsible for managing repos and backups. type Orchestrator struct { - mu sync.Mutex - config *v1.Config - oplog *oplog.OpLog + mu sync.Mutex + config *v1.Config + OpLog *oplog.OpLog repoPool *resticRepoPool - configUpdates chan *v1.Config // configUpdates chan makes config changes available to Run() - externTasks chan Task // externTasks is a channel that externally added tasks can be added to, they will be consumed by Run() + externTasks chan Task // externTasks is a channel that externally added tasks can be added to, they will be consumed by Run() } func NewOrchestrator(configProvider config.ConfigStore, oplog *oplog.OpLog) (*Orchestrator, error) { @@ -38,9 +37,9 @@ func NewOrchestrator(configProvider config.ConfigStore, oplog *oplog.OpLog) (*Or } return &Orchestrator{ - config: cfg, - oplog: oplog, - repoPool: newResticRepoPool(&config.MemoryStore{Config: cfg}), + config: cfg, + OpLog: oplog, + repoPool: newResticRepoPool(&config.MemoryStore{Config: cfg}), externTasks: make(chan Task, 2), }, nil } @@ -51,7 +50,7 @@ func (o *Orchestrator) ApplyConfig(cfg *v1.Config) error { o.config = cfg zap.L().Debug("Applying config to orchestrator", zap.Any("config", cfg)) - + // Update the config provided to the repo pool. if err := o.repoPool.configProvider.Update(cfg); err != nil { return fmt.Errorf("failed to update repo pool config: %w", err) @@ -70,7 +69,7 @@ func (o *Orchestrator) GetRepo(repoId string) (repo *RepoOrchestrator, err error defer o.mu.Unlock() r, err := o.repoPool.GetRepo(repoId) - if err != nil { + if err != nil { return nil, fmt.Errorf("failed to get repo %q: %w", repoId, err) } return r, nil @@ -119,7 +118,7 @@ func (o *Orchestrator) Run(mainCtx context.Context) error { func (o *Orchestrator) runVersion(mainCtx context.Context, config *v1.Config) bool { var lock sync.Mutex ctx, cancel := context.WithCancel(mainCtx) - + var wg sync.WaitGroup var execTask func(t Task) @@ -149,7 +148,7 @@ func (o *Orchestrator) runVersion(mainCtx context.Context, config *v1.Config) bo lock.Lock() defer lock.Unlock() zap.L().Info("running task", zap.String("task", t.Name())) - + // Task execution runs with mainCtx meaning config changes do not interrupt it, but cancelling the orchestration loop will. if err := t.Run(mainCtx); err != nil { zap.L().Error("task failed", zap.String("task", t.Name()), zap.Error(err)) @@ -159,7 +158,7 @@ func (o *Orchestrator) runVersion(mainCtx context.Context, config *v1.Config) bo if ctx.Err() != nil { zap.L().Debug("not attempting to reschedule task, orchestrator context is cancelled.", zap.String("task", t.Name())) - return + return } execTask(t) @@ -202,15 +201,14 @@ func (o *Orchestrator) EnqueueTask(t Task) { // resticRepoPool caches restic repos. type resticRepoPool struct { - mu sync.Mutex - repos map[string]*RepoOrchestrator + mu sync.Mutex + repos map[string]*RepoOrchestrator configProvider config.ConfigStore } - func newResticRepoPool(configProvider config.ConfigStore) *resticRepoPool { return &resticRepoPool{ - repos: make(map[string]*RepoOrchestrator), + repos: make(map[string]*RepoOrchestrator), configProvider: configProvider, } } @@ -244,7 +242,7 @@ func (rp *resticRepoPool) GetRepo(repoId string) (repo *RepoOrchestrator, err er if ok && proto.Equal(repo.repoConfig, repoProto) { return repo, nil } - delete(rp.repos, repoId); + delete(rp.repos, repoId) var opts []restic.GenericOption opts = append(opts, restic.WithPropagatedEnvVars(restic.EnvToPropagate...)) diff --git a/internal/orchestrator/tasks.go b/internal/orchestrator/tasks.go index 3c693401..2e291fde 100644 --- a/internal/orchestrator/tasks.go +++ b/internal/orchestrator/tasks.go @@ -6,16 +6,16 @@ import ( "time" v1 "github.com/garethgeorge/resticui/gen/go/v1" - "github.com/garethgeorge/resticui/internal/database/oplog" + "github.com/garethgeorge/resticui/internal/oplog" + "github.com/garethgeorge/resticui/internal/oplog/indexutil" "github.com/garethgeorge/resticui/pkg/restic" "github.com/gitploy-io/cronexpr" "github.com/hashicorp/go-multierror" "go.uber.org/zap" ) - type Task interface { - Name() string // huamn readable name for this task. + Name() string // huamn readable name for this task. Next(now time.Time) *time.Time // when this task would like to be run. Run(ctx context.Context) error // run the task. } @@ -23,8 +23,8 @@ type Task interface { // BackupTask is a scheduled backup operation. type ScheduledBackupTask struct { orchestrator *Orchestrator // owning orchestrator - plan *v1.Plan - schedule *cronexpr.Schedule + plan *v1.Plan + schedule *cronexpr.Schedule } var _ Task = &ScheduledBackupTask{} @@ -37,8 +37,8 @@ func NewScheduledBackupTask(orchestrator *Orchestrator, plan *v1.Plan) (*Schedul return &ScheduledBackupTask{ orchestrator: orchestrator, - plan: plan, - schedule: sched, + plan: plan, + schedule: sched, }, nil } @@ -58,15 +58,15 @@ func (t *ScheduledBackupTask) Run(ctx context.Context) error { // OnetimeBackupTask is a single backup operation. type OnetimeBackupTask struct { orchestrator *Orchestrator - plan *v1.Plan - time *time.Time + plan *v1.Plan + time *time.Time } func NewOneofBackupTask(orchestrator *Orchestrator, plan *v1.Plan, at time.Time) *OnetimeBackupTask { return &OnetimeBackupTask{ orchestrator: orchestrator, - plan: plan, - time: &at, + plan: plan, + time: &at, } } @@ -75,9 +75,9 @@ func (t *OnetimeBackupTask) Name() string { } func (t *OnetimeBackupTask) Next(now time.Time) *time.Time { - ret := t.time + ret := t.time t.time = nil - return ret + return ret } func (t *OnetimeBackupTask) Run(ctx context.Context) error { @@ -91,16 +91,16 @@ func backupHelper(ctx context.Context, orchestrator *Orchestrator, plan *v1.Plan } op := &v1.Operation{ - PlanId: plan.Id, - RepoId: plan.Repo, + PlanId: plan.Id, + RepoId: plan.Repo, UnixTimeStartMs: curTimeMillis(), - Status: v1.OperationStatus_STATUS_INPROGRESS, - Op: backupOp, + Status: v1.OperationStatus_STATUS_INPROGRESS, + Op: backupOp, } startTime := time.Now() - err := WithOperation(orchestrator.oplog, op, func() error { + err := WithOperation(orchestrator.OpLog, op, func() error { zap.L().Info("Starting backup", zap.String("plan", plan.Id)) repo, err := orchestrator.GetRepo(plan.Repo) if err != nil { @@ -109,13 +109,13 @@ func backupHelper(ctx context.Context, orchestrator *Orchestrator, plan *v1.Plan lastSent := time.Now() // debounce progress updates, these can endup being very frequent. summary, err := repo.Backup(ctx, plan, func(entry *restic.BackupProgressEntry) { - if time.Since(lastSent) < 200 * time.Millisecond { + if time.Since(lastSent) < 200*time.Millisecond { return } lastSent = time.Now() backupOp.OperationBackup.LastStatus = entry.ToProto() - if err := orchestrator.oplog.Update(op); err != nil { + if err := orchestrator.OpLog.Update(op); err != nil { zap.S().Errorf("failed to update oplog with progress for backup: %v", err) } zap.L().Debug("backup progress", zap.Float64("progress", entry.PercentDone)) @@ -124,10 +124,8 @@ func backupHelper(ctx context.Context, orchestrator *Orchestrator, plan *v1.Plan return fmt.Errorf("repo.Backup for repo %q: %w", plan.Repo, err) } + op.SnapshotId = summary.SnapshotId backupOp.OperationBackup.LastStatus = summary.ToProto() - if err := orchestrator.oplog.Update(op); err != nil { - return fmt.Errorf("update oplog with summary for backup: %v", err) - } zap.L().Info("backup complete", zap.String("plan", plan.Id), zap.Duration("duration", time.Since(startTime))) return nil @@ -136,7 +134,7 @@ func backupHelper(ctx context.Context, orchestrator *Orchestrator, plan *v1.Plan return fmt.Errorf("backup operation: %w", err) } - // this could alternatively be a separate operation, but it probably makes sense to index snapshots immediately after a backup. + // this could alternatively be scheduled as a separate task, but it probably makes sense to index snapshots immediately after a backup. if err := indexSnapshotsHelper(ctx, orchestrator, plan); err != nil { return fmt.Errorf("reindexing snapshots after backup operation: %w", err) } @@ -157,47 +155,56 @@ func indexSnapshotsHelper(ctx context.Context, orchestrator *Orchestrator, plan startTime := time.Now() alreadyIndexed := 0 - opTime := curTimeMillis() var indexOps []*v1.Operation for _, snapshot := range snapshots { - opid, err := orchestrator.oplog.HasIndexedSnapshot(snapshot.Id) + ops, err := orchestrator.OpLog.GetBySnapshotId(snapshot.Id, indexutil.CollectAll()) if err != nil { return fmt.Errorf("HasIndexSnapshot for snapshot %q: %w", snapshot.Id, err) } - if opid >= 0 { + if containsSnapshotOperation(ops) { alreadyIndexed += 1 continue } - + + snapshotProto := snapshot.ToProto() indexOps = append(indexOps, &v1.Operation{ - RepoId: plan.Repo, - PlanId: plan.Id, - UnixTimeStartMs: opTime, - UnixTimeEndMs: opTime, - Status: v1.OperationStatus_STATUS_SUCCESS, + RepoId: plan.Repo, + PlanId: plan.Id, + UnixTimeStartMs: snapshotProto.UnixTimeMs, + UnixTimeEndMs: snapshotProto.UnixTimeMs, + Status: v1.OperationStatus_STATUS_SUCCESS, Op: &v1.Operation_OperationIndexSnapshot{ OperationIndexSnapshot: &v1.OperationIndexSnapshot{ - Snapshot: snapshot.ToProto(), + Snapshot: snapshotProto, }, }, }) } - if err := orchestrator.oplog.BulkAdd(indexOps); err != nil { + if err := orchestrator.OpLog.BulkAdd(indexOps); err != nil { return fmt.Errorf("BulkAdd snapshot operations: %w", err) } - zap.L().Debug("Indexed snapshots", - zap.String("plan", plan.Id), - zap.Duration("duration", time.Since(startTime)), - zap.Int("alreadyIndexed", alreadyIndexed), - zap.Int("newlyAdded", len(snapshots) - alreadyIndexed), + zap.L().Debug("Indexed snapshots", + zap.String("plan", plan.Id), + zap.Duration("duration", time.Since(startTime)), + zap.Int("alreadyIndexed", alreadyIndexed), + zap.Int("newlyAdded", len(snapshots)-alreadyIndexed), ) return err } +func containsSnapshotOperation(ops []*v1.Operation) bool { + for _, op := range ops { + if _, ok := op.Op.(*v1.Operation_OperationIndexSnapshot); ok { + return true + } + } + return false +} + // WithOperation is a utility that creates an operation to track the function's execution. // timestamps are automatically added and the status is automatically updated if an error occurs. func WithOperation(oplog *oplog.OpLog, op *v1.Operation, do func() error) error { @@ -209,7 +216,7 @@ func WithOperation(oplog *oplog.OpLog, op *v1.Operation, do func() error) error } err := do() if err != nil { - op.Status = v1.OperationStatus_STATUS_ERROR + op.Status = v1.OperationStatus_STATUS_ERROR op.DisplayMessage = err.Error() } op.UnixTimeEndMs = curTimeMillis() @@ -224,5 +231,5 @@ func WithOperation(oplog *oplog.OpLog, op *v1.Operation, do func() error) error func curTimeMillis() int64 { t := time.Now() - return t.Unix() * 1000 + int64(t.Nanosecond() / 1000000) -} \ No newline at end of file + return t.Unix()*1000 + int64(t.Nanosecond()/1000000) +} diff --git a/pkg/restic/outputs.go b/pkg/restic/outputs.go index 53e2ba87..c13512db 100644 --- a/pkg/restic/outputs.go +++ b/pkg/restic/outputs.go @@ -14,13 +14,13 @@ import ( ) type LsEntry struct { - Name string `json:"name"` - Type string `json:"type"` - Path string `json:"path"` - Uid int `json:"uid"` - Gid int `json:"gid"` - Size int `json:"size"` - Mode int `json:"mode"` + Name string `json:"name"` + Type string `json:"type"` + Path string `json:"path"` + Uid int `json:"uid"` + Gid int `json:"gid"` + Size int `json:"size"` + Mode int `json:"mode"` Mtime string `json:"mtime"` Atime string `json:"atime"` Ctime string `json:"ctime"` @@ -28,13 +28,13 @@ type LsEntry struct { func (e *LsEntry) ToProto() *v1.LsEntry { return &v1.LsEntry{ - Name: e.Name, - Type: e.Type, - Path: e.Path, - Uid: int64(e.Uid), - Gid: int64(e.Gid), - Size: int64(e.Size), - Mode: int64(e.Mode), + Name: e.Name, + Type: e.Type, + Path: e.Path, + Uid: int64(e.Uid), + Gid: int64(e.Gid), + Size: int64(e.Size), + Mode: int64(e.Mode), Mtime: e.Mtime, Atime: e.Atime, Ctime: e.Ctime, @@ -42,28 +42,27 @@ func (e *LsEntry) ToProto() *v1.LsEntry { } type Snapshot struct { - Id string `json:"id"` - Time string `json:"time"` - Tree string `json:"tree"` - Paths []string `json:"paths"` - Hostname string `json:"hostname"` - Username string `json:"username"` - Tags []string `json:"tags"` - Parent string `json:"parent"` - unixTimeMs int64 `json:"-"` + Id string `json:"id"` + Time string `json:"time"` + Tree string `json:"tree"` + Paths []string `json:"paths"` + Hostname string `json:"hostname"` + Username string `json:"username"` + Tags []string `json:"tags"` + Parent string `json:"parent"` + unixTimeMs int64 `json:"-"` } func (s *Snapshot) ToProto() *v1.ResticSnapshot { - return &v1.ResticSnapshot{ - Id: s.Id, + Id: s.Id, UnixTimeMs: s.UnixTimeMs(), - Tree: s.Tree, - Paths: s.Paths, - Hostname: s.Hostname, - Username: s.Username, - Tags: s.Tags, - Parent: s.Parent, + Tree: s.Tree, + Paths: s.Paths, + Hostname: s.Hostname, + Username: s.Username, + Tags: s.Tags, + Parent: s.Parent, } } @@ -84,26 +83,26 @@ type BackupProgressEntry struct { MessageType string `json:"message_type"` // "summary" or "status" // Summary fields - FilesNew int `json:"files_new"` - FilesChanged int `json:"files_changed"` - FilesUnmodified int `json:"files_unmodified"` - DirsNew int `json:"dirs_new"` - DirsChanged int `json:"dirs_changed"` - DirsUnmodified int `json:"dirs_unmodified"` - DataBlobs int `json:"data_blobs"` - TreeBlobs int `json:"tree_blobs"` - DataAdded int `json:"data_added"` - TotalFilesProcessed int `json:"total_files_processed"` - TotalBytesProcessed int `json:"total_bytes_processed"` - TotalDuration float64 `json:"total_duration"` - SnapshotId string `json:"snapshot_id"` + FilesNew int `json:"files_new"` + FilesChanged int `json:"files_changed"` + FilesUnmodified int `json:"files_unmodified"` + DirsNew int `json:"dirs_new"` + DirsChanged int `json:"dirs_changed"` + DirsUnmodified int `json:"dirs_unmodified"` + DataBlobs int `json:"data_blobs"` + TreeBlobs int `json:"tree_blobs"` + DataAdded int `json:"data_added"` + TotalFilesProcessed int `json:"total_files_processed"` + TotalBytesProcessed int `json:"total_bytes_processed"` + TotalDuration float64 `json:"total_duration"` + SnapshotId string `json:"snapshot_id"` // Status fields PercentDone float64 `json:"percent_done"` - TotalFiles int `json:"total_files"` - FilesDone int `json:"files_done"` - TotalBytes int `json:"total_bytes"` - BytesDone int `json:"bytes_done"` + TotalFiles int `json:"total_files"` + FilesDone int `json:"files_done"` + TotalBytes int `json:"total_bytes"` + BytesDone int `json:"bytes_done"` } func (b *BackupProgressEntry) ToProto() *v1.BackupProgressEntry { @@ -112,19 +111,19 @@ func (b *BackupProgressEntry) ToProto() *v1.BackupProgressEntry { return &v1.BackupProgressEntry{ Entry: &v1.BackupProgressEntry_Summary{ Summary: &v1.BackupProgressSummary{ - FilesNew: int64(b.FilesNew), - FilesChanged: int64(b.FilesChanged), - FilesUnmodified: int64(b.FilesUnmodified), - DirsNew: int64(b.DirsNew), - DirsChanged: int64(b.DirsChanged), - DirsUnmodified: int64(b.DirsUnmodified), - DataBlobs: int64(b.DataBlobs), - TreeBlobs: int64(b.TreeBlobs), - DataAdded: int64(b.DataAdded), + FilesNew: int64(b.FilesNew), + FilesChanged: int64(b.FilesChanged), + FilesUnmodified: int64(b.FilesUnmodified), + DirsNew: int64(b.DirsNew), + DirsChanged: int64(b.DirsChanged), + DirsUnmodified: int64(b.DirsUnmodified), + DataBlobs: int64(b.DataBlobs), + TreeBlobs: int64(b.TreeBlobs), + DataAdded: int64(b.DataAdded), TotalFilesProcessed: int64(b.TotalFilesProcessed), TotalBytesProcessed: int64(b.TotalBytesProcessed), - TotalDuration: float64(b.TotalDuration), - SnapshotId: b.SnapshotId, + TotalDuration: float64(b.TotalDuration), + SnapshotId: b.SnapshotId, }, }, } @@ -133,16 +132,16 @@ func (b *BackupProgressEntry) ToProto() *v1.BackupProgressEntry { Entry: &v1.BackupProgressEntry_Status{ Status: &v1.BackupProgressStatusEntry{ PercentDone: b.PercentDone, - TotalFiles: int64(b.TotalFiles), - FilesDone: int64(b.FilesDone), - TotalBytes: int64(b.TotalBytes), - BytesDone: int64(b.BytesDone), + TotalFiles: int64(b.TotalFiles), + FilesDone: int64(b.FilesDone), + TotalBytes: int64(b.TotalBytes), + BytesDone: int64(b.BytesDone), }, }, } default: log.Fatalf("unknown message type: %s", b.MessageType) - return nil + return nil } } @@ -212,4 +211,4 @@ func readLs(output io.Reader) (*Snapshot, []*LsEntry, error) { entries = append(entries, entry) } return snapshot, entries, nil -} \ No newline at end of file +} diff --git a/proto/v1/operations.proto b/proto/v1/operations.proto index 7cc36d67..39c0dcf4 100644 --- a/proto/v1/operations.proto +++ b/proto/v1/operations.proto @@ -14,6 +14,7 @@ message Operation { int64 id = 1; string repo_id = 2; // repo id if associated with a repo (always true) string plan_id = 3; // plan id if associated with a plan (always true) + string snapshot_id = 8; // snapshot id if associated with a snapshot. OperationStatus status = 4; int64 unix_time_start_ms = 5; int64 unix_time_end_ms = 6; diff --git a/resticui.go b/resticui.go index 3f93a801..82dc5887 100644 --- a/resticui.go +++ b/resticui.go @@ -14,7 +14,7 @@ import ( "github.com/garethgeorge/resticui/internal/api" "github.com/garethgeorge/resticui/internal/config" - "github.com/garethgeorge/resticui/internal/database/oplog" + "github.com/garethgeorge/resticui/internal/oplog" "github.com/garethgeorge/resticui/internal/orchestrator" static "github.com/garethgeorge/resticui/webui" "github.com/mattn/go-colorable" @@ -33,18 +33,17 @@ func main() { ctx := context.Background() ctx, cancel := context.WithCancel(ctx) go onterm(cancel) - + if _, err := config.Default.Get(); err != nil { zap.S().Fatalf("Error loading config: %v", err) } - + var wg sync.WaitGroup // Configure the HTTP mux mux := http.NewServeMux() mux.Handle("/", http.FileServer(http.FS(&SubdirFilesystem{FS: static.FS, subdir: "dist"}))) - // Create and serve API server oplogFile := path.Join(dataPath(), "oplog.boltdb") oplog, err := oplog.NewOpLog(oplogFile) @@ -115,9 +114,9 @@ func init() { c.EncodeLevel = zapcore.CapitalColorLevelEncoder c.EncodeTime = zapcore.ISO8601TimeEncoder l := zap.New(zapcore.NewCore( - zapcore.NewConsoleEncoder(c), - zapcore.AddSync(colorable.NewColorableStdout()), - zapcore.DebugLevel, + zapcore.NewConsoleEncoder(c), + zapcore.AddSync(colorable.NewColorableStdout()), + zapcore.DebugLevel, )) zap.ReplaceGlobals(l) } @@ -142,7 +141,7 @@ func (s *SubdirFilesystem) Open(name string) (fs.File, error) { return s.FS.Open(path.Join(s.subdir, name)) } -func (s *SubdirFilesystem) ReadDir(name string) ([]fs.DirEntry, error) { +func (s *SubdirFilesystem) ReadDir(name string) ([]fs.DirEntry, error) { readDirFS := s.FS.(fs.ReadDirFS) if readDirFS == nil { return nil, &fs.PathError{Op: "readdir", Path: name, Err: errors.New("not implemented")} @@ -156,4 +155,4 @@ func dataPath() string { datahome = path.Join(os.Getenv("HOME") + "/.local/share") } return path.Join(datahome, "resticui") -} \ No newline at end of file +} diff --git a/webui/gen/ts/v1/operations.pb.ts b/webui/gen/ts/v1/operations.pb.ts index 7b1711f0..474b5cbf 100644 --- a/webui/gen/ts/v1/operations.pb.ts +++ b/webui/gen/ts/v1/operations.pb.ts @@ -38,6 +38,7 @@ type BaseOperation = { id?: string repoId?: string planId?: string + snapshotId?: string status?: OperationStatus unixTimeStartMs?: string unixTimeEndMs?: string