Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123package spindle
import ( "context" "fmt"
"github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/jetstream/pkg/models" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/codes" "tangled.org/core/api/tangled" "tangled.org/core/spindle/observability" "tangled.org/core/tapc")
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.EventKindAccount { return s.ingestAccount(ctx, e) } if e.Kind != models.EventKindCommit { return nil }
ctx, span := observability.Tracer().Start(ctx, "jetstream.ingest") defer span.End()
if span.IsRecording() { var attrs []attribute.KeyValue if e.Did != "" { attrs = append(attrs, attribute.String(observability.UserDIDKey, e.Did)) } if e.Commit != nil { if e.Commit.Collection != "" { attrs = append(attrs, attribute.String(observability.CollectionKey, e.Commit.Collection)) } if e.Commit.RKey != "" { attrs = append(attrs, attribute.String(observability.RKeyKey, e.Commit.RKey)) } } span.SetAttributes(attrs...) } var err error switch e.Commit.Collection { case tangled.RepoNSID, tangled.RepoCollaboratorNSID: if evt, ok := jetstreamToTapEvent(e); ok { err = s.tap.processEvent(ctx, evt) } case tangled.RepoPullNSID: if evt, ok := jetstreamToTapEvent(e); ok { err = s.processPull(ctx, evt.Record) } case tangled.RepoPullStatusNSID: if evt, ok := jetstreamToTapEvent(e); ok { err = s.processPullStatus(ctx, evt.Record) } }
if err != nil { s.l.WarnContext(ctx, "failed to process message, skipping", "nsid", e.Commit.Collection, "did", e.Did, "rkey", e.Commit.RKey, "err", err) s.metrics.RecordEventIngestion("jetstream", "error") span.SetStatus(codes.Error, "failed to process jetstream event") } else { s.metrics.RecordEventIngestion("jetstream", "success") span.SetStatus(codes.Ok, "success") }
return nil }}
func jetstreamToTapEvent(e *models.Event) (tapc.Event, bool) { if e.Commit == nil { return tapc.Event{}, false } did, err := syntax.ParseDID(e.Did) if err != nil { return tapc.Event{}, false } var action tapc.RecordAction switch e.Commit.Operation { case models.CommitOperationCreate: action = tapc.RecordCreateAction case models.CommitOperationUpdate: action = tapc.RecordUpdateAction case models.CommitOperationDelete: action = tapc.RecordDeleteAction default: return tapc.Event{}, false } return tapc.Event{ Type: tapc.EvtRecord, Record: &tapc.RecordEventData{ Did: did, Rkey: syntax.RecordKey(e.Commit.RKey), Collection: syntax.NSID(e.Commit.Collection), Rev: e.Commit.Rev, Action: action, Record: e.Commit.Record, // jetstream is only used for live Live: true, }, }, true}
func (s *Spindle) ingestAccount(ctx context.Context, e *models.Event) error { if e.Account == nil || e.Account.Active { return nil } status := "unknown" if e.Account.Status != nil { status = *e.Account.Status } did, err := syntax.ParseDID(e.Account.Did) if err != nil { return fmt.Errorf("parsing account did: %w", err) } s.l.Warn("account went inactive, purging its repos", "did", did, "status", status) return s.WipeOwner(ctx, did, "account "+status)}