Something went wrong. Try again.
Monorepo for Tangled — https://tangled.org forked from tangled.org/core
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139package spindle
import ( "context" "encoding/json" "fmt"
"tangled.sh/tangled.sh/core/api/tangled" "tangled.sh/tangled.sh/core/eventconsumer"
"github.com/bluesky-social/jetstream/pkg/models")
type Ingester func(ctx context.Context, e *models.Event) error
func (s *Spindle) ingest() Ingester { return func(ctx context.Context, e *models.Event) error { var err error defer func() { eventTime := e.TimeUS lastTimeUs := eventTime + 1 if err := s.db.SaveLastTimeUs(lastTimeUs); err != nil { err = fmt.Errorf("(deferred) failed to save last time us: %w", err) } }()
if e.Kind != models.EventKindCommit { return nil }
switch e.Commit.Collection { case tangled.SpindleMemberNSID: s.ingestMember(ctx, e) case tangled.RepoNSID: s.ingestRepo(ctx, e) }
return err }}
func (s *Spindle) ingestMember(_ context.Context, e *models.Event) error { did := e.Did var err error
l := s.l.With("component", "ingester", "record", tangled.SpindleMemberNSID)
switch e.Commit.Operation { case models.CommitOperationCreate, models.CommitOperationUpdate: raw := e.Commit.Record record := tangled.SpindleMember{} err = json.Unmarshal(raw, &record) if err != nil { l.Error("invalid record", "error", err) return err }
domain := s.cfg.Server.Hostname if s.cfg.Server.Dev { domain = s.cfg.Server.ListenAddr } recordInstance := record.Instance
if recordInstance != domain { l.Error("domain mismatch", "domain", recordInstance, "expected", domain) return fmt.Errorf("domain mismatch: %s != %s", record.Instance, domain) }
ok, err := s.e.IsSpindleInviteAllowed(did, rbacDomain) if err != nil || !ok { l.Error("failed to add member", "did", did, "error", err) return fmt.Errorf("failed to enforce permissions: %w", err) }
if err := s.e.AddKnotMember(rbacDomain, record.Subject); err != nil { l.Error("failed to add member", "error", err) return fmt.Errorf("failed to add member: %w", err) } l.Info("added member from firehose", "member", record.Subject)
if err := s.db.AddDid(record.Subject); err != nil { l.Error("failed to add did", "error", err) return fmt.Errorf("failed to add did: %w", err) } s.jc.AddDid(record.Subject)
return nil
} return nil}
func (s *Spindle) ingestRepo(_ context.Context, e *models.Event) error { var err error
l := s.l.With("component", "ingester", "record", tangled.RepoNSID)
l.Info("ingesting repo record")
switch e.Commit.Operation { case models.CommitOperationCreate, models.CommitOperationUpdate: raw := e.Commit.Record record := tangled.Repo{} err = json.Unmarshal(raw, &record) if err != nil { l.Error("invalid record", "error", err) return err }
domain := s.cfg.Server.Hostname
// no spindle configured for this repo if record.Spindle == nil { l.Info("no spindle configured", "did", record.Owner, "name", record.Name) return nil }
// this repo did not want this spindle if *record.Spindle != domain { l.Info("different spindle configured", "did", record.Owner, "name", record.Name, "spindle", *record.Spindle, "domain", domain) return nil }
// add this repo to the watch list if err := s.db.AddRepo(record.Knot, record.Owner, record.Name); err != nil { l.Error("failed to add repo", "error", err) return fmt.Errorf("failed to add repo: %w", err) }
// add this knot to the event consumer src := eventconsumer.NewKnotSource(record.Knot) s.ks.AddSource(context.Background(), src)
return nil
} return nil}