diff --git a/events/diskpersist.go b/events/diskpersist.go index 686b1025..44a2bfc4 100644 --- a/events/diskpersist.go +++ b/events/diskpersist.go @@ -184,6 +184,8 @@ func (dp *DiskPersistence) initLogFile() error { return nil } +// swapLog swaps the current log file out for a new empty one +// must only be called while holding dp.lk func (dp *DiskPersistence) swapLog(ctx context.Context) error { if err := dp.logfi.Close(); err != nil { return fmt.Errorf("failed to close current log file: %w", err) @@ -259,26 +261,6 @@ const ( var emptyHeader = make([]byte, headerSize) -func (p *DiskPersistence) addJobsToQueue(jobs []persistJob) error { - p.lk.Lock() - defer p.lk.Unlock() - - for _, job := range jobs { - if err := p.doPersist(job); err != nil { - return err - } - - // TODO: for some reason replacing this constant with p.writeBufferSize dramatically reduces perf... - if len(p.evtbuf) > 400 { - if err := p.flushLog(context.TODO()); err != nil { - return fmt.Errorf("failed to flush disk log: %w", err) - } - } - } - - return nil -} - func (p *DiskPersistence) addJobToQueue(job persistJob) error { p.lk.Lock() defer p.lk.Unlock() diff --git a/events/events.go b/events/events.go index c00d8a4d..5e2e928d 100644 --- a/events/events.go +++ b/events/events.go @@ -78,17 +78,6 @@ func (em *EventManager) persistAndSendEvent(ctx context.Context, evt *XRPCStream } } -type batchPersister interface { - PersistMany(ctx context.Context, evts []*XRPCStreamEvent) error -} - -func (em *EventManager) persistAndSendEvents(ctx context.Context, evts []*XRPCStreamEvent) { - pm := em.persister.(batchPersister) - if err := pm.PersistMany(ctx, evts); err != nil { - log.Errorf("failed to persist outbound events: %s", err) - } -} - type Subscriber struct { outgoing chan *XRPCStreamEvent @@ -136,14 +125,6 @@ func (em *EventManager) AddEvent(ctx context.Context, ev *XRPCStreamEvent) error return nil } -func (em *EventManager) AddEventBatch(ctx context.Context, evs []*XRPCStreamEvent) error { - ctx, span := otel.Tracer("events").Start(ctx, "AddEventBatch") - defer span.End() - - em.persistAndSendEvents(ctx, evs) - return nil -} - var ErrPlaybackShutdown = fmt.Errorf("playback shutting down") func (em *EventManager) Subscribe(ctx context.Context, filter func(*XRPCStreamEvent) bool, since *int64) (<-chan *XRPCStreamEvent, func(), error) {