Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
1.2 kB · 60 lines
Go
at next-dev
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061package 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(), }}