Documentation
¶
Index ¶
- Constants
- Variables
- func TenantMiddleware(limits Limits) middleware.Interface
- type ApplyStorageUpdatesFunc
- type Chunk
- type CompactorClient
- type DeleteRequestClientMetrics
- type DeleteRequestHandler
- func (dm *DeleteRequestHandler) AddDeleteRequestHandler(w http.ResponseWriter, r *http.Request)
- func (dm *DeleteRequestHandler) CancelDeleteRequestHandler(w http.ResponseWriter, r *http.Request)
- func (dm *DeleteRequestHandler) GetAllDeleteRequestsHandler(w http.ResponseWriter, r *http.Request)
- func (dm *DeleteRequestHandler) GetCacheGenerationNumberHandler(w http.ResponseWriter, r *http.Request)
- type DeleteRequestsClient
- type DeleteRequestsKind
- type DeleteRequestsManager
- func (d *DeleteRequestsManager) CanSkipSeries(userID []byte, lbls labels.Labels, seriesID []byte, _ model.Time, ...) bool
- func (d *DeleteRequestsManager) DropFromIndex(_ []byte, _ retention.Chunk, _ labels.Labels, _ model.Time, _ model.Time) bool
- func (d *DeleteRequestsManager) Expired(userID []byte, chk retention.Chunk, lbls labels.Labels, seriesID []byte, ...) (bool, filter.Func)
- func (d *DeleteRequestsManager) Init(tablesManager TablesManager, registerer prometheus.Registerer) error
- func (d *DeleteRequestsManager) IntervalMayHaveExpiredChunks(_ model.Interval, userID string) bool
- func (d *DeleteRequestsManager) JobBuilder() *JobBuilder
- func (d *DeleteRequestsManager) MarkPhaseFailed()
- func (d *DeleteRequestsManager) MarkPhaseFinished()
- func (d *DeleteRequestsManager) MarkPhaseStarted()
- func (d *DeleteRequestsManager) MarkPhaseTimedOut()
- func (d *DeleteRequestsManager) MarkSeriesAsProcessed(userID, seriesID []byte, lbls labels.Labels, tableName string) error
- func (d *DeleteRequestsManager) Start(ctx context.Context)
- type DeleteRequestsStore
- type DeleteRequestsStoreDBType
- type DeleteRequestsStoreOption
- type GRPCRequestHandler
- type GetChunkClientForTableFunc
- type JobBuilder
- type JobRunner
- type Limits
- type StorageUpdatesIterator
- type Table
- type TableIteratorFunc
- type TablesManager
- type TimeRange
Constants ¶
const (
DeleteRequestsTableName = "delete_requests"
)
const ForQuerytimeFilteringQueryParam = "for_querytime_filtering"
Variables ¶
var ErrDeleteRequestNotFound = errors.New("could not find matching delete requests")
var ErrNoChunksSelectedForDeletion = fmt.Errorf("no chunks selected for deletion")
var SupportedDeleteRequestsStoreDBTypes = []DeleteRequestsStoreDBType{DeleteRequestsStoreDBTypeBoltDB, DeleteRequestsStoreDBTypeSQLite}
Functions ¶
func TenantMiddleware ¶
func TenantMiddleware(limits Limits) middleware.Interface
Types ¶
type ApplyStorageUpdatesFunc ¶ added in v3.6.0
type ApplyStorageUpdatesFunc func(ctx context.Context, iterator StorageUpdatesIterator) error
type CompactorClient ¶
type DeleteRequestClientMetrics ¶
type DeleteRequestClientMetrics struct {
// contains filtered or unexported fields
}
func NewDeleteRequestClientMetrics ¶
func NewDeleteRequestClientMetrics(r prometheus.Registerer) *DeleteRequestClientMetrics
func (DeleteRequestClientMetrics) Unregister ¶ added in v3.1.0
func (m DeleteRequestClientMetrics) Unregister()
type DeleteRequestHandler ¶
type DeleteRequestHandler struct {
// contains filtered or unexported fields
}
DeleteRequestHandler provides handlers for delete requests
func NewDeleteRequestHandler ¶
func NewDeleteRequestHandler(deleteStore DeleteRequestsStore, maxInterval, deleteRequestCancelPeriod time.Duration, registerer prometheus.Registerer) *DeleteRequestHandler
NewDeleteRequestHandler creates a DeleteRequestHandler
func (*DeleteRequestHandler) AddDeleteRequestHandler ¶
func (dm *DeleteRequestHandler) AddDeleteRequestHandler(w http.ResponseWriter, r *http.Request)
AddDeleteRequestHandler handles addition of a new delete request
func (*DeleteRequestHandler) CancelDeleteRequestHandler ¶
func (dm *DeleteRequestHandler) CancelDeleteRequestHandler(w http.ResponseWriter, r *http.Request)
CancelDeleteRequestHandler handles delete request cancellation
func (*DeleteRequestHandler) GetAllDeleteRequestsHandler ¶
func (dm *DeleteRequestHandler) GetAllDeleteRequestsHandler(w http.ResponseWriter, r *http.Request)
GetAllDeleteRequestsHandler handles get all delete requests
func (*DeleteRequestHandler) GetCacheGenerationNumberHandler ¶
func (dm *DeleteRequestHandler) GetCacheGenerationNumberHandler(w http.ResponseWriter, r *http.Request)
GetCacheGenerationNumberHandler handles requests for a user's cache generation number
type DeleteRequestsClient ¶
type DeleteRequestsClient interface {
GetAllDeleteRequestsForUser(ctx context.Context, userID string, forQuerytimeFiltering bool, timeRange *TimeRange) ([]deletionproto.DeleteRequest, error)
Stop()
}
func NewDeleteRequestsClient ¶
func NewDeleteRequestsClient(compactorClient CompactorClient, deleteClientMetrics *DeleteRequestClientMetrics, clientType string, opts ...DeleteRequestsStoreOption) (DeleteRequestsClient, error)
func NewNoOpDeleteRequestsClient ¶ added in v3.5.0
func NewNoOpDeleteRequestsClient() DeleteRequestsClient
func NewPerTenantDeleteRequestsClient ¶
func NewPerTenantDeleteRequestsClient(c DeleteRequestsClient, l Limits) DeleteRequestsClient
type DeleteRequestsKind ¶ added in v3.6.0
type DeleteRequestsKind string
const ( DeleteRequestsWithLineFilters DeleteRequestsKind = "DeleteRequestsWithLineFilters" DeleteRequestsWithoutLineFilters DeleteRequestsKind = "DeleteRequestsWithoutLineFilters" DeleteRequestsAll DeleteRequestsKind = "DeleteRequestsAll" )
type DeleteRequestsManager ¶
type DeleteRequestsManager struct {
HSModeEnabled bool
// contains filtered or unexported fields
}
func NewDeleteRequestsManager ¶
func NewDeleteRequestsManager( workingDir string, store DeleteRequestsStore, deleteRequestCancelPeriod time.Duration, batchSize int, limits Limits, HSModeEnabled bool, deletionManifestStoreClient client.ObjectClient, registerer prometheus.Registerer, ) (*DeleteRequestsManager, error)
func (*DeleteRequestsManager) CanSkipSeries ¶ added in v3.5.0
func (*DeleteRequestsManager) DropFromIndex ¶
func (*DeleteRequestsManager) Init ¶ added in v3.6.0
func (d *DeleteRequestsManager) Init(tablesManager TablesManager, registerer prometheus.Registerer) error
func (*DeleteRequestsManager) IntervalMayHaveExpiredChunks ¶
func (d *DeleteRequestsManager) IntervalMayHaveExpiredChunks(_ model.Interval, userID string) bool
func (*DeleteRequestsManager) JobBuilder ¶ added in v3.6.0
func (d *DeleteRequestsManager) JobBuilder() *JobBuilder
func (*DeleteRequestsManager) MarkPhaseFailed ¶
func (d *DeleteRequestsManager) MarkPhaseFailed()
func (*DeleteRequestsManager) MarkPhaseFinished ¶
func (d *DeleteRequestsManager) MarkPhaseFinished()
func (*DeleteRequestsManager) MarkPhaseStarted ¶
func (d *DeleteRequestsManager) MarkPhaseStarted()
func (*DeleteRequestsManager) MarkPhaseTimedOut ¶
func (d *DeleteRequestsManager) MarkPhaseTimedOut()
func (*DeleteRequestsManager) MarkSeriesAsProcessed ¶ added in v3.5.0
func (d *DeleteRequestsManager) MarkSeriesAsProcessed(userID, seriesID []byte, lbls labels.Labels, tableName string) error
MarkSeriesAsProcessed marks a series as processed. It ignores the operation if the series progress file reference is nil.
func (*DeleteRequestsManager) Start ¶ added in v3.6.0
func (d *DeleteRequestsManager) Start(ctx context.Context)
Start starts the DeleteRequestsManager's background operations. It is a blocking call. To stop the background operations, cancel the passed context.
type DeleteRequestsStore ¶
type DeleteRequestsStore interface {
AddDeleteRequest(ctx context.Context, userID, query string, startTime, endTime model.Time, shardByInterval time.Duration) (string, error)
GetAllRequests(ctx context.Context) ([]deletionproto.DeleteRequest, error)
GetAllDeleteRequestsForUser(ctx context.Context, userID string, forQuerytimeFiltering bool, timeRange *TimeRange) ([]deletionproto.DeleteRequest, error)
RemoveDeleteRequest(ctx context.Context, userID string, requestID string) error
GetDeleteRequest(ctx context.Context, userID, requestID string) (deletionproto.DeleteRequest, error)
GetCacheGenerationNumber(ctx context.Context, userID string) (string, error)
MergeShardedRequests(ctx context.Context) error
// ToDo(Sandeep): To keep changeset smaller, below 2 methods treat a single shard as individual request. This can be refactored later in a separate PR.
MarkShardAsProcessed(ctx context.Context, req deletionproto.DeleteRequest) error
GetUnprocessedShards(ctx context.Context) ([]deletionproto.DeleteRequest, error)
Stop()
// contains filtered or unexported methods
}
func NewDeleteRequestsStore ¶ added in v3.5.0
func NewDeleteRequestsStore( deleteRequestsStoreDBType DeleteRequestsStoreDBType, workingDirectory string, indexStorageClient storage.Client, backupDeleteRequestStoreDBType DeleteRequestsStoreDBType, indexUpdatePropagationMaxDelay time.Duration, ) (DeleteRequestsStore, error)
type DeleteRequestsStoreDBType ¶ added in v3.5.0
type DeleteRequestsStoreDBType string
const ( DeleteRequestsStoreDBTypeBoltDB DeleteRequestsStoreDBType = "boltdb" DeleteRequestsStoreDBTypeSQLite DeleteRequestsStoreDBType = "sqlite" )
type DeleteRequestsStoreOption ¶
type DeleteRequestsStoreOption func(c *deleteRequestsClient)
func WithRequestClientCacheDuration ¶
func WithRequestClientCacheDuration(d time.Duration) DeleteRequestsStoreOption
type GRPCRequestHandler ¶
type GRPCRequestHandler struct {
// contains filtered or unexported fields
}
func NewGRPCRequestHandler ¶
func NewGRPCRequestHandler(deleteRequestsStore DeleteRequestsStore, limits Limits) *GRPCRequestHandler
func (*GRPCRequestHandler) GetCacheGenNumbers ¶
func (g *GRPCRequestHandler) GetCacheGenNumbers(ctx context.Context, _ *grpc.GetCacheGenNumbersRequest) (*grpc.GetCacheGenNumbersResponse, error)
func (*GRPCRequestHandler) GetDeleteRequests ¶
func (g *GRPCRequestHandler) GetDeleteRequests(ctx context.Context, req *grpc.GetDeleteRequestsRequest) (*grpc.GetDeleteRequestsResponse, error)
type GetChunkClientForTableFunc ¶ added in v3.6.0
type JobBuilder ¶ added in v3.6.0
type JobBuilder struct {
// contains filtered or unexported fields
}
func NewJobBuilder ¶ added in v3.6.0
func NewJobBuilder( deletionManifestStoreClient client.ObjectClient, applyStorageUpdatesFunc ApplyStorageUpdatesFunc, markRequestsAsProcessedFunc markRequestsAsProcessedFunc, r prometheus.Registerer, ) *JobBuilder
func (*JobBuilder) BuildJobs ¶ added in v3.6.0
func (b *JobBuilder) BuildJobs(ctx context.Context, jobsChan chan<- *grpc.Job)
BuildJobs implements jobqueue.Builder interface
func (*JobBuilder) JobsLeft ¶ added in v3.6.0
func (b *JobBuilder) JobsLeft() int
func (*JobBuilder) OnJobResponse ¶ added in v3.6.0
func (b *JobBuilder) OnJobResponse(response *grpc.JobResult) error
OnJobResponse implements jobqueue.Builder interface
type JobRunner ¶ added in v3.6.0
type JobRunner struct {
// contains filtered or unexported fields
}
func NewJobRunner ¶ added in v3.6.0
func NewJobRunner(chunkProcessingConcurrency int, getStorageClientForTableFunc GetChunkClientForTableFunc, r prometheus.Registerer) *JobRunner
type Limits ¶
type Limits interface {
DeletionMode(userID string) string
RetentionPeriod(userID string) time.Duration
StreamRetention(userID string) []validation.StreamRetention
}
type StorageUpdatesIterator ¶ added in v3.6.0
type Table ¶ added in v3.6.0
type Table interface {
GetUserIndex(userID string) (retention.SeriesIterator, error)
}
type TableIteratorFunc ¶ added in v3.6.0
type TablesManager ¶ added in v3.6.0
Source Files
¶
- delete_request.go
- delete_request_batch.go
- delete_requests_client.go
- delete_requests_db_boltdb.go
- delete_requests_db_sqlite.go
- delete_requests_manager.go
- delete_requests_store.go
- delete_requests_store_boltdb.go
- delete_requests_store_sqlite.go
- deletion_manifest_builder.go
- grpc_request_handler.go
- job_builder.go
- job_runner.go
- metrics.go
- request_handler.go
- tenant_delete_requests_client.go
- tenant_request_handler.go
- util.go