Monorepo for Tangled forked from tangled.org/core
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304package spindle
import ( "context" "encoding/json" "errors" "fmt" "time"
"tangled.org/core/api/tangled" "tangled.org/core/eventconsumer" "tangled.org/core/rbac" "tangled.org/core/spindle/db"
comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/identity" "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/indigo/xrpc" "github.com/bluesky-social/jetstream/pkg/models" securejoin "github.com/cyphar/filepath-securejoin")
type Ingester func(ctx context.Context, e *models.Event) error
func (s *Spindle) ingest() Ingester { return func(ctx context.Context, e *models.Event) error { if e.Kind != models.EventKindCommit { return nil }
var err error switch e.Commit.Collection { case tangled.SpindleMemberNSID: err = s.ingestMember(ctx, e) case tangled.RepoNSID: err = s.ingestRepo(ctx, e) case tangled.RepoCollaboratorNSID: err = s.ingestCollaborator(ctx, e) }
if err != nil { s.l.Warn("failed to process message, skipping", "nsid", e.Commit.Collection, "err", err) }
lastTimeUs := e.TimeUS + 1 if saveErr := s.db.SaveLastTimeUs(lastTimeUs); saveErr != nil { s.l.Error("failed to save cursor", "err", saveErr) }
return nil }}
func (s *Spindle) ingestMember(_ context.Context, e *models.Event) error { var err error did := e.Did rkey := e.Commit.RKey
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 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 := db.AddSpindleMember(s.db, db.SpindleMember{ Did: syntax.DID(did), Rkey: rkey, Instance: recordInstance, Subject: syntax.DID(record.Subject), Created: time.Now(), }); err != nil { l.Error("failed to add member", "error", err) return fmt.Errorf("failed to add member: %w", err) }
if err := s.e.AddSpindleMember(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
case models.CommitOperationDelete: record, err := db.GetSpindleMember(s.db, did, rkey) if err != nil { l.Error("failed to find member", "error", err) return fmt.Errorf("failed to find member: %w", err) }
if err := db.RemoveSpindleMember(s.db, did, rkey); err != nil { l.Error("failed to remove member", "error", err) return fmt.Errorf("failed to remove member: %w", err) }
if err := s.e.RemoveSpindleMember(rbacDomain, record.Subject.String()); 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.RemoveDid(record.Subject.String()); err != nil { l.Error("failed to add did", "error", err) return fmt.Errorf("failed to add did: %w", err) } s.jc.RemoveDid(record.Subject.String())
} return nil}
func (s *Spindle) ingestRepo(ctx context.Context, e *models.Event) error { var err error did := e.Did
l := s.l.With("component", "ingester", "record", tangled.RepoNSID)
l.Info("ingesting repo record", "did", did)
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", "name", record.Name) return nil }
// this repo did not want this spindle if *record.Spindle != domain { l.Info("different spindle configured", "name", record.Name, "spindle", *record.Spindle, "domain", domain) return nil }
// add this repo to the watch list if err := s.db.AddRepo(record.Knot, did, record.Name); err != nil { l.Error("failed to add repo", "error", err) return fmt.Errorf("failed to add repo: %w", err) }
didSlashRepo, err := securejoin.SecureJoin(did, record.Name) if err != nil { return err }
// add repo to rbac if err := s.e.AddRepo(did, rbac.ThisServer, didSlashRepo); err != nil { l.Error("failed to add repo to enforcer", "error", err) return fmt.Errorf("failed to add repo: %w", err) }
// add collaborators to rbac owner, err := s.res.ResolveIdent(ctx, did) if err != nil || owner.Handle.IsInvalidHandle() { return err } if err := s.fetchAndAddCollaborators(ctx, owner, didSlashRepo); err != nil { return err }
// add this knot to the event consumer src := eventconsumer.NewKnotSource(record.Knot) s.ks.AddSource(context.Background(), src)
return nil
} return nil}
func (s *Spindle) ingestCollaborator(ctx context.Context, e *models.Event) error { var err error
l := s.l.With("component", "ingester", "record", tangled.RepoCollaboratorNSID, "did", e.Did)
l.Info("ingesting collaborator record")
switch e.Commit.Operation { case models.CommitOperationCreate, models.CommitOperationUpdate: raw := e.Commit.Record record := tangled.RepoCollaborator{} err = json.Unmarshal(raw, &record) if err != nil { l.Error("invalid record", "error", err) return err }
subjectId, err := s.res.ResolveIdent(ctx, record.Subject) if err != nil || subjectId.Handle.IsInvalidHandle() { return err }
var rbacResource string var ownerDid string switch { case record.Repo != nil: repoAt, parseErr := syntax.ParseATURI(*record.Repo) if parseErr != nil { l.Info("rejecting record, invalid repoAt", "repoAt", *record.Repo) return nil }
owner, resolveErr := s.res.ResolveIdent(ctx, repoAt.Authority().String()) if resolveErr != nil || owner.Handle.IsInvalidHandle() { return fmt.Errorf("failed to resolve handle: %w", resolveErr) }
xrpcc := xrpc.Client{ Host: owner.PDSEndpoint(), }
resp, getErr := comatproto.RepoGetRecord(ctx, &xrpcc, "", tangled.RepoNSID, repoAt.Authority().String(), repoAt.RecordKey().String()) if getErr != nil { return getErr }
repo := resp.Value.Val.(*tangled.Repo) rbacResource, _ = securejoin.SecureJoin(owner.DID.String(), repo.Name) ownerDid = owner.DID.String()
default: l.Info("rejecting collaborator record without repo at-uri (spindle RBAC keyed by owner/name)") return nil }
if ok, err := s.e.IsCollaboratorInviteAllowed(ownerDid, rbac.ThisServer, rbacResource); !ok || err != nil { return fmt.Errorf("insufficient permissions: %w", err) }
if err := s.e.AddCollaborator(record.Subject, rbac.ThisServer, rbacResource); err != nil { l.Error("failed to add collaborator to enforcer", "error", err) return fmt.Errorf("failed to add collaborator: %w", err) }
return nil } return nil}
func (s *Spindle) fetchAndAddCollaborators(ctx context.Context, owner *identity.Identity, didSlashRepo string) error { l := s.l.With("component", "ingester", "handler", "fetchAndAddCollaborators")
l.Info("fetching and adding existing collaborators")
xrpcc := xrpc.Client{ Host: owner.PDSEndpoint(), }
resp, err := comatproto.RepoListRecords(ctx, &xrpcc, tangled.RepoCollaboratorNSID, "", 50, owner.DID.String(), false) if err != nil { return err }
var errs error for _, r := range resp.Records { if r == nil { continue } record := r.Value.Val.(*tangled.RepoCollaborator)
if err := s.e.AddCollaborator(record.Subject, rbac.ThisServer, didSlashRepo); err != nil { l.Error("failed to add repo to enforcer", "error", err) errors.Join(errs, fmt.Errorf("failed to add repo: %w", err)) } }
return errs}