Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120package feed
import ( "context" "log/slog" "sync" "time"
"tangled.org/core/hostutil" knotfeed "tangled.org/core/knotfeed")
type Hooks struct { LoadCursor func(ctx context.Context, knot string) (knotfeed.Cursor, error) StoreCursor func(ctx context.Context, knot string, cursor knotfeed.Cursor) error Handle func(ctx context.Context, knot string, msg knotfeed.Message) error OutdatedReplay func(ctx context.Context, knot string, feed knotfeed.Feed) knotfeed.Cursor OnConnectError func(knot string, err error)}
type Feed struct { logger *slog.Logger hooks Hooks
mu sync.Mutex subs map[string]*subscriber}
type subscriber struct { cancel context.CancelFunc}
func New(logger *slog.Logger, hooks Hooks) *Feed { return &Feed{ logger: logger, hooks: hooks, subs: make(map[string]*subscriber), }}
func (f *Feed) Start(ctx context.Context, knots []string) { for _, knot := range knots { f.Subscribe(ctx, knot) }}
func (f *Feed) Subscribe(ctx context.Context, knot string) { host, noTLS, err := hostutil.ParseHostname(knot) if err != nil { f.logger.Error("unsubscribable knot host", "knot", knot, "err", err) return }
f.mu.Lock() if _, ok := f.subs[knot]; ok { f.mu.Unlock() return } runCtx, cancel := context.WithCancel(ctx) sub := &subscriber{cancel: cancel} f.subs[knot] = sub f.mu.Unlock()
consumer := &knotfeed.Consumer{ Host: host, NoTLS: noTLS, Logger: f.logger, LoadCursor: func(ctx context.Context) (knotfeed.Cursor, error) { return f.hooks.LoadCursor(ctx, knot) }, StoreCursor: func(ctx context.Context, cursor knotfeed.Cursor) error { return f.hooks.StoreCursor(ctx, knot, cursor) }, Handle: func(ctx context.Context, msg knotfeed.Message) error { return f.hooks.Handle(ctx, knot, msg) }, OutdatedReplay: func(ctx context.Context, feed knotfeed.Feed) knotfeed.Cursor { if f.hooks.OutdatedReplay == nil { return feed.Live(time.Now()) } return f.hooks.OutdatedReplay(ctx, knot, feed) }, } if f.hooks.OnConnectError != nil { consumer.OnConnectError = func(err error) { f.hooks.OnConnectError(knot, err) } }
go f.run(runCtx, knot, sub, consumer)}
func (f *Feed) Unsubscribe(knot string) { f.mu.Lock() sub, ok := f.subs[knot] if !ok { f.mu.Unlock() return } delete(f.subs, knot) f.mu.Unlock()
sub.cancel()}
func (f *Feed) run(ctx context.Context, knot string, sub *subscriber, consumer *knotfeed.Consumer) { defer f.drop(knot, sub) if err := consumer.Run(ctx); err != nil && ctx.Err() == nil { f.logger.Error("knot feed exited", "knot", knot, "err", err) }}
func (f *Feed) drop(knot string, sub *subscriber) { f.mu.Lock() defer f.mu.Unlock() if f.subs[knot] == sub { delete(f.subs, knot) }}