package 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) } }