package knotstream import ( "sync" "tangled.org/core/knotfeed" "tangled.org/core/knotmirror/models" ) type subscription struct { hostname string mu sync.Mutex last knotfeed.Cursor applied knotfeed.Cursor scheduler *ParallelScheduler } func (s *subscription) Last() knotfeed.Cursor { s.mu.Lock() defer s.mu.Unlock() return s.last } func (s *subscription) Applied() knotfeed.Cursor { s.mu.Lock() defer s.mu.Unlock() return s.applied } func (s *subscription) Resume(cursor knotfeed.Cursor) { s.mu.Lock() defer s.mu.Unlock() s.last, s.applied = cursor, cursor } func (s *subscription) Seen(cursor knotfeed.Cursor) { s.mu.Lock() defer s.mu.Unlock() if cursor.Feed() != s.last.Feed() { s.applied = knotfeed.NewCursor(cursor.Feed(), 0) } s.last = cursor } func (s *subscription) MarkApplied(cursor knotfeed.Cursor) { s.mu.Lock() defer s.mu.Unlock() if cursor.Feed() == s.applied.Feed() { s.applied = knotfeed.NewCursor(s.applied.Feed(), max(s.applied.Seq(), cursor.Seq())) } } func (s *subscription) HostCursor() models.HostCursor { return models.HostCursor{ Hostname: s.hostname, Cursor: s.Applied(), } }