Leon4gr45's picture
Upload folder using huggingface_hub (part 6)
d6f631f verified
Raw
History Blame Contribute Delete
1.56 kB
package flushhandler
import (
"context"
"errors"
"github.com/openmeterio/openmeter/openmeter/sink/models"
)
type DrainCompleteFunc func()
var _ FlushEventHandler = (*FlushEventHandlers)(nil)
type FlushEventHandlers struct {
handlers []FlushEventHandler
onDrainComplete []DrainCompleteFunc
}
func NewFlushEventHandlers() *FlushEventHandlers {
return &FlushEventHandlers{}
}
func (f *FlushEventHandlers) AddHandler(handler FlushEventHandler) {
f.handlers = append(f.handlers, handler)
}
func (f *FlushEventHandlers) OnDrainComplete(fn DrainCompleteFunc) {
f.onDrainComplete = append(f.onDrainComplete, fn)
}
func (f *FlushEventHandlers) OnFlushSuccess(ctx context.Context, events []models.SinkMessage) error {
var errs []error
for _, handler := range f.handlers {
if err := handler.OnFlushSuccess(ctx, events); err != nil {
errs = append(errs, err)
}
}
return errors.Join(errs...)
}
func (f *FlushEventHandlers) Close() error {
var errs []error
for _, handler := range f.handlers {
if err := handler.Close(); err != nil {
errs = append(errs, err)
}
}
return errors.Join(errs...)
}
func (f *FlushEventHandlers) Start(ctx context.Context) error {
for _, handler := range f.handlers {
if err := handler.Start(ctx); err != nil {
return err
}
}
return nil
}
func (f *FlushEventHandlers) WaitForDrain(ctx context.Context) error {
for _, handler := range f.handlers {
if err := handler.WaitForDrain(ctx); err != nil {
return err
}
}
for _, fn := range f.onDrainComplete {
fn()
}
return nil
}