Something went wrong. Try again.
Monorepo for Tangled
Something went wrong. Try again.
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790package appview
import ( "context" "database/sql" "encoding/json" "errors" "fmt" "io" "log/slog" "maps" "net/http" "net/url" "slices" "strings" "sync"
"time"
"github.com/avast/retry-go/v4" "github.com/bluesky-social/indigo/atproto/syntax" jmodels "github.com/bluesky-social/jetstream/pkg/models" "github.com/go-git/go-git/v5/plumbing" "github.com/ipfs/go-cid" "golang.org/x/sync/errgroup" "tangled.org/core/api/tangled" "tangled.org/core/appview/cache" "tangled.org/core/appview/config" "tangled.org/core/appview/db" "tangled.org/core/appview/mentions" "tangled.org/core/appview/models" "tangled.org/core/appview/notify" "tangled.org/core/appview/repoverify" "tangled.org/core/appview/serververify" "tangled.org/core/appview/validator" "tangled.org/core/idresolver" "tangled.org/core/orm" "tangled.org/core/rbac")
type Ingester struct { Ctx context.Context Db *db.DB Enforcer *rbac.Enforcer IdResolver *idresolver.Resolver Cache *cache.Cache Config *config.Config Logger *slog.Logger Validator *validator.Validator MentionsResolver *mentions.Resolver Notifier notify.Notifier Verifier repoverify.Verifier}
type processFunc func(ctx context.Context, e *jmodels.Event) error
func (i *Ingester) Ingest() processFunc { return func(ctx context.Context, e *jmodels.Event) error { var err error
l := i.Logger.With("kind", e.Kind) switch e.Kind { case jmodels.EventKindAccount: // TODO: sync account state to db if e.Account.Active { break } // TODO: revoke sessions by DID if *e.Account.Status == "deactivated" { err = i.IdResolver.InvalidateIdent(ctx, e.Account.Did) } case jmodels.EventKindIdentity: err = i.IdResolver.InvalidateIdent(ctx, e.Identity.Did) case jmodels.EventKindCommit: switch e.Commit.Collection { case tangled.GraphFollowNSID: err = i.ingestFollow(e) case tangled.GraphVouchNSID: err = i.ingestVouch(ctx, e) case tangled.FeedStarNSID: err = i.ingestStar(ctx, e) case tangled.PublicKeyNSID: err = i.ingestPublicKey(e) case tangled.RepoArtifactNSID: err = i.ingestArtifact(ctx, e) case tangled.ActorProfileNSID: err = i.ingestProfile(ctx, e) case tangled.SpindleMemberNSID: err = i.ingestSpindleMember(ctx, e) case tangled.SpindleNSID: err = i.ingestSpindle(ctx, e) case tangled.KnotMemberNSID: err = i.ingestKnotMember(ctx, e) case tangled.KnotNSID: err = i.ingestKnot(ctx, e) case tangled.StringNSID: err = i.ingestString(e) case tangled.RepoIssueNSID: err = i.ingestIssue(ctx, e) case tangled.RepoPullNSID: err = i.ingestPull(ctx, e) case tangled.FeedCommentNSID: err = i.ingestComment(e) case tangled.RepoIssueCommentNSID: err = i.ingestIssueComment(e) case tangled.RepoPullCommentNSID: err = i.ingestPullComment(e) case tangled.LabelDefinitionNSID: err = i.ingestLabelDefinition(e) case tangled.LabelOpNSID: err = i.ingestLabelOp(e) case tangled.RepoNSID: err = i.ingestRepo(ctx, e) } l = i.Logger.With("nsid", e.Commit.Collection) }
if err != nil { l.Warn("failed to ingest record, skipping", "err", err) }
lastTimeUs := e.TimeUS + 1 if saveErr := i.Db.SaveLastTimeUs(lastTimeUs); saveErr != nil { l.Error("failed to save cursor", "err", saveErr) }
return nil }}
func (i *Ingester) resolveRepoRef(ref string) (*models.Repo, error) { if strings.HasPrefix(ref, "did:") { return db.GetRepoByDid(i.Db, ref) } return db.GetRepoByAtUri(i.Db, ref)}
func (i *Ingester) resolveOldFormatStar(raw json.RawMessage, star *models.Star, l *slog.Logger) (bool, error) { var legacy struct { Subject *string `json:"subject"` SubjectDid *string `json:"subjectDid"` } if err := json.Unmarshal(raw, &legacy); err != nil { return false, err }
switch { case legacy.SubjectDid != nil: repo, err := i.resolveRepoRef(*legacy.SubjectDid) if err != nil { l.Warn("skipping old-format star for unknown repo", "subjectDid", *legacy.SubjectDid) return false, nil } star.SubjectType = models.StarSubjectRepo star.Subject = repo.RepoDid return true, nil
case legacy.Subject != nil: uri, err := syntax.ParseATURI(*legacy.Subject) if err != nil { return false, fmt.Errorf("invalid old-format star subject: %w", err) } switch uri.Collection().String() { case tangled.RepoNSID: repo, err := db.GetRepoByAtUri(i.Db, uri.String()) if err != nil { l.Warn("skipping old-format star for unknown repo", "subject", *legacy.Subject) return false, nil } star.SubjectType = models.StarSubjectRepo star.Subject = repo.RepoDid return true, nil default: star.SubjectType = models.StarSubjectString star.Subject = *legacy.Subject return true, nil }
default: return false, fmt.Errorf("old-format star has neither subject nor subjectDid") }}
func (i *Ingester) ingestStar(ctx context.Context, e *jmodels.Event) error { var err error did := e.Did
l := i.Logger.With("handler", "ingestStar") l = l.With("nsid", e.Commit.Collection)
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) record := tangled.FeedStar{} unmarshalErr := json.Unmarshal(raw, &record)
star := &models.Star{ Did: did, Rkey: e.Commit.RKey, }
switch { case unmarshalErr != nil: resolved, resolveErr := i.resolveOldFormatStar(raw, star, l) if resolveErr != nil { l.Error("invalid record", "newFmtErr", unmarshalErr, "oldFmtErr", resolveErr) return unmarshalErr } if !resolved { return nil }
case record.Subject == nil: return fmt.Errorf("star record has nil subject")
case record.Subject.FeedStar_Repo != nil: repo, repoErr := i.resolveRepoRef(record.Subject.FeedStar_Repo.Did) if repoErr != nil { l.Warn("skipping star for unknown repo", "did", record.Subject.FeedStar_Repo.Did) return nil } star.SubjectType = models.StarSubjectRepo star.Subject = repo.RepoDid
case record.Subject.FeedStar_String != nil: star.SubjectType = models.StarSubjectString star.Subject = record.Subject.FeedStar_String.Uri
default: return fmt.Errorf("star record has empty subject union") }
err = db.AddStar(i.Db, star) case jmodels.CommitOperationDelete: err = db.DeleteStarByRkey(i.Db, did, e.Commit.RKey) }
if err != nil { return fmt.Errorf("failed to %s star record: %w", e.Commit.Operation, err) }
return nil}
func (i *Ingester) ingestFollow(e *jmodels.Event) error { var err error did := e.Did
l := i.Logger.With("handler", "ingestFollow") l = l.With("nsid", e.Commit.Collection)
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) record := tangled.GraphFollow{} err = json.Unmarshal(raw, &record) if err != nil { l.Error("invalid record", "err", err) return err }
err = db.AddFollow(i.Db, &models.Follow{ UserDid: did, SubjectDid: record.Subject, Rkey: e.Commit.RKey, }) case jmodels.CommitOperationDelete: err = db.DeleteFollowByRkey(i.Db, did, e.Commit.RKey) }
if err != nil { return fmt.Errorf("failed to %s follow record: %w", e.Commit.Operation, err) }
return nil}
func (i *Ingester) ingestVouch(ctx context.Context, e *jmodels.Event) error { var err error did := e.Did
l := i.Logger.With("handler", "ingestVouch") l = l.With("nsid", e.Commit.Collection) l.Info("ingesting vouch")
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) record := tangled.GraphVouch{} err = json.Unmarshal(raw, &record) if err != nil { l.Error("invalid record", "err", err) return err }
// rkey is the subject_did being vouched for/denounced subjectDID := e.Commit.RKey
_, err = syntax.ParseDID(subjectDID) if err != nil { l.Error("invalid subject_did in rkey", "err", err, "rkey", subjectDID) return fmt.Errorf("invalid subject_did: %w", err) }
if did == subjectDID { l.Warn("attempted self-vouch", "did", did) return fmt.Errorf("cannot vouch for self") }
subjectId, err := i.IdResolver.ResolveIdent(ctx, subjectDID) if err != nil { return err }
if subjectId.Handle.IsInvalidHandle() { return err }
kind, err := models.ParseVouchKind(record.Kind) if err != nil { l.Error("invalid kind", "kind", kind) return fmt.Errorf("invalid kind: %s", kind) }
recordCid, err := cid.Parse(e.Commit.CID) if err != nil { l.Error("invalid cid", "err", err, "cid", e.Commit.CID) return fmt.Errorf("invalid cid: %w", err) }
var evidences []syntax.ATURI for _, raw := range record.Evidences { uri, parseErr := syntax.ParseATURI(raw) if parseErr != nil { l.Warn("invalid evidence AT-URI, skipping", "uri", raw, "err", parseErr) continue } evidences = append(evidences, uri) }
tx, txErr := i.Db.Begin() if txErr != nil { return fmt.Errorf("failed to start transaction: %w", txErr) }
addErr := db.AddVouch(tx, &models.Vouch{ Did: syntax.DID(did), SubjectDid: subjectId.DID, Cid: recordCid, Kind: kind, Reason: record.Reason, Evidences: evidences, }) if addErr != nil { tx.Rollback() err = addErr } else { err = tx.Commit() }
case jmodels.CommitOperationDelete: err = db.DeleteVouchByRkey(i.Db, did, e.Commit.RKey) }
if err != nil { return fmt.Errorf("failed to %s vouch record: %w", e.Commit.Operation, err) }
return nil}
func (i *Ingester) ingestPublicKey(e *jmodels.Event) error { did := e.Did var err error
l := i.Logger.With("handler", "ingestPublicKey") l = l.With("nsid", e.Commit.Collection)
switch e.Commit.Operation { case jmodels.CommitOperationCreate: l.Debug("processing add of pubkey") raw := json.RawMessage(e.Commit.Record) record := tangled.PublicKey{} err = json.Unmarshal(raw, &record) if err != nil { l.Error("invalid record", "err", err) return err }
name := record.Name key := record.Key err = db.AddPublicKey(i.Db, did, name, key, e.Commit.RKey) case jmodels.CommitOperationUpdate: l.Debug("processing update of pubkey") raw := json.RawMessage(e.Commit.Record) record := tangled.PublicKey{} err = json.Unmarshal(raw, &record) if err != nil { l.Error("invalid record", "err", err) return err }
name := record.Name key := record.Key err = db.UpdatePublicKey(i.Db, did, name, key, e.Commit.RKey) case jmodels.CommitOperationDelete: l.Debug("processing delete of pubkey") err = db.DeletePublicKeyByRkey(i.Db, did, e.Commit.RKey) }
if err != nil { return fmt.Errorf("failed to %s pubkey record: %w", e.Commit.Operation, err) }
return nil}
func (i *Ingester) ingestArtifact(ctx context.Context, e *jmodels.Event) error { did := e.Did var err error
l := i.Logger.With("handler", "ingestArtifact") l = l.With("nsid", e.Commit.Collection)
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) record := tangled.RepoArtifact{} err = json.Unmarshal(raw, &record) if err != nil { l.Error("invalid record", "err", err) return err }
var repo *models.Repo if record.RepoDid != nil && *record.RepoDid != "" { repo, err = db.GetRepoByDid(i.Db, *record.RepoDid) if err != nil && !errors.Is(err, sql.ErrNoRows) { return fmt.Errorf("failed to look up repo by DID %s: %w", *record.RepoDid, err) } } if repo == nil && record.Repo != nil { repoAt, parseErr := syntax.ParseATURI(*record.Repo) if parseErr != nil { return parseErr } repo, err = db.GetRepoByAtUri(i.Db, repoAt.String()) if err != nil { return err } } if repo == nil { return fmt.Errorf("artifact record has neither valid repoDid nor repo field") }
ok, err := i.Enforcer.E.Enforce(did, repo.Knot, repo.RepoIdentifier(), "repo:push") if err != nil || !ok { return err }
repoDid := repo.RepoDid if repoDid == "" && record.RepoDid != nil { repoDid = *record.RepoDid } if repoDid != "" && (record.RepoDid == nil || *record.RepoDid == "") && record.Repo != nil { if enqErr := db.EnqueuePdsRecordMigration(ctx, i.Db, "add-repo-did", syntax.DID(did), syntax.NSID(tangled.RepoArtifactNSID), syntax.RecordKey(e.Commit.RKey)); enqErr != nil { l.Warn("failed to enqueue PDS rewrite for artifact", "err", enqErr, "did", did, "repoDid", repoDid) } }
createdAt, err := time.Parse(time.RFC3339, record.CreatedAt) if err != nil { createdAt = time.Now() }
artifact := models.Artifact{ Did: did, Rkey: e.Commit.RKey, RepoDid: syntax.DID(repo.RepoDid), Tag: plumbing.Hash(record.Tag), CreatedAt: createdAt, BlobCid: cid.Cid(record.Artifact.Ref), Name: record.Name, Size: uint64(record.Artifact.Size), MimeType: record.Artifact.MimeType, }
err = db.AddArtifact(i.Db, artifact) case jmodels.CommitOperationDelete: err = db.DeleteArtifact(i.Db, orm.FilterEq("did", did), orm.FilterEq("rkey", e.Commit.RKey)) }
if err != nil { return fmt.Errorf("failed to %s artifact record: %w", e.Commit.Operation, err) }
return nil}
func (i *Ingester) ingestProfile(ctx context.Context, e *jmodels.Event) error { did := e.Did var err error
l := i.Logger.With("handler", "ingestProfile") l = l.With("nsid", e.Commit.Collection)
if e.Commit.RKey != "self" { return fmt.Errorf("ingestProfile only ingests `self` record") }
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) record := tangled.ActorProfile{} err = json.Unmarshal(raw, &record) if err != nil { l.Error("invalid record", "err", err) return err }
avatar := "" if record.Avatar != nil { avatar = record.Avatar.Ref.String() }
description := "" if record.Description != nil { description = *record.Description }
includeBluesky := record.Bluesky
pronouns := "" if record.Pronouns != nil { pronouns = *record.Pronouns }
location := "" if record.Location != nil { location = *record.Location }
var links [5]string for i, l := range record.Links { if i < 5 { links[i] = l } }
var stats [2]models.VanityStat for i, s := range record.Stats { if i < 2 { stats[i].Kind = models.ParseVanityStatKind(s) } }
var pinned [6]string for i, r := range record.PinnedRepositories { if i < 6 { pinned[i] = r } }
var preferredHandle syntax.Handle if record.PreferredHandle != nil { if h, err := syntax.ParseHandle(*record.PreferredHandle); err == nil { ident, identErr := i.IdResolver.ResolveIdent(ctx, did) if identErr == nil && slices.Contains(ident.AlsoKnownAs, "at://"+string(h)) { preferredHandle = h } } }
profile := models.Profile{ Did: did, Avatar: avatar, Description: description, IncludeBluesky: includeBluesky, Location: location, Links: links, Stats: stats, PinnedRepos: pinned, Pronouns: pronouns, PreferredHandle: preferredHandle, }
tx, err := i.Db.Begin() if err != nil { return fmt.Errorf("failed to start transaction: %w", err) }
err = db.ValidateProfile(tx, &profile) if err != nil { return fmt.Errorf("invalid profile record") }
err = db.UpsertProfile(tx, &profile) if err == nil && i.Cache != nil { pipe := i.Cache.Pipeline() didKey := fmt.Sprintf(cache.PreferredHandleByDid, did) if preferredHandle != "" { pipe.Set(ctx, didKey, string(preferredHandle), cache.PreferredHandleTTL) pipe.Set(ctx, fmt.Sprintf(cache.PreferredHandleByHandle, string(preferredHandle)), did, cache.PreferredHandleTTL) } else { pipe.Del(ctx, didKey) } if _, execErr := pipe.Exec(ctx); execErr != nil { l.Warn("failed to update preferred handle cache", "err", execErr) } } case jmodels.CommitOperationDelete: tx, beginErr := i.Db.Begin() if beginErr != nil { return fmt.Errorf("failed to start transaction: %w", beginErr) }
priorHandle, phErr := db.GetPreferredHandle(tx, did) if phErr != nil && !errors.Is(phErr, sql.ErrNoRows) { l.Warn("failed to read prior preferred handle", "err", phErr) }
err = db.DeleteProfile(tx, did) if err == nil && i.Cache != nil { pipe := i.Cache.Pipeline() pipe.Del(ctx, fmt.Sprintf(cache.PreferredHandleByDid, did)) if priorHandle != "" { pipe.Del(ctx, fmt.Sprintf(cache.PreferredHandleByHandle, string(priorHandle))) } if _, execErr := pipe.Exec(ctx); execErr != nil { l.Warn("failed to evict preferred handle cache", "err", execErr) } } }
if err != nil { return fmt.Errorf("failed to %s profile record: %w", e.Commit.Operation, err) }
return nil}
func (i *Ingester) ingestSpindleMember(ctx context.Context, e *jmodels.Event) error { did := e.Did var err error
l := i.Logger.With("handler", "ingestSpindleMember") l = l.With("nsid", e.Commit.Collection)
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) record := tangled.SpindleMember{} err = json.Unmarshal(raw, &record) if err != nil { l.Error("invalid record", "err", err) return err }
// only spindle owner can invite to spindles ok, err := i.Enforcer.IsSpindleInviteAllowed(did, record.Instance) if err != nil { return fmt.Errorf("failed to check invite permission: %w", err) } if !ok { if verifyErr := i.verifySpindle(ctx, record.Instance, did); verifyErr != nil { return fmt.Errorf("invite denied and verify failed: %w", verifyErr) } ok, err = i.Enforcer.IsSpindleInviteAllowed(did, record.Instance) if err != nil { return fmt.Errorf("failed to re-check invite permission: %w", err) } if !ok { return fmt.Errorf("invite denied for did %s on spindle %s", did, record.Instance) } }
memberId, err := i.IdResolver.ResolveIdent(ctx, record.Subject) if err != nil { return err }
if memberId.Handle.IsInvalidHandle() { return fmt.Errorf("invalid handle for member %s", record.Subject) }
existing, err := db.GetSpindleMembers(i.Db, orm.FilterEq("did", did), orm.FilterEq("rkey", e.Commit.RKey), ) if err != nil { return fmt.Errorf("failed to look up existing member: %w", err) } if len(existing) > 1 { return fmt.Errorf("multiple spindle members with rkey %s", e.Commit.RKey) }
tx, err := i.Db.Begin() if err != nil { return fmt.Errorf("failed to start txn: %w", err) } committed := false defer func() { if committed { return } tx.Rollback() i.Enforcer.E.LoadPolicy() }()
if len(existing) == 1 { prev := existing[0] if prev.Instance != record.Instance || prev.Subject != memberId.DID { if err = db.RemoveSpindleMember(tx, orm.FilterEq("did", did), orm.FilterEq("rkey", e.Commit.RKey), ); err != nil { return fmt.Errorf("failed to remove stale row: %w", err) } if err = i.Enforcer.RemoveSpindleMember(prev.Instance, prev.Subject.String()); err != nil { return fmt.Errorf("failed to remove stale ACL: %w", err) } } }
if err = db.AddSpindleMember(tx, models.SpindleMember{ Did: syntax.DID(did), Rkey: e.Commit.RKey, Instance: record.Instance, Subject: memberId.DID, }); err != nil { return fmt.Errorf("failed to add to db: %w", err) }
if err = i.Enforcer.AddSpindleMember(record.Instance, memberId.DID.String()); err != nil { return fmt.Errorf("failed to update ACLs: %w", err) }
if err = tx.Commit(); err != nil { return fmt.Errorf("failed to commit txn: %w", err) }
if err = i.Enforcer.E.SavePolicy(); err != nil { return fmt.Errorf("failed to save ACLs: %w", err) } committed = true
l.Info("upserted spindle member") case jmodels.CommitOperationDelete: rkey := e.Commit.RKey
// get record from db first members, err := db.GetSpindleMembers( i.Db, orm.FilterEq("did", did), orm.FilterEq("rkey", rkey), ) if err != nil || len(members) != 1 { return fmt.Errorf("failed to get member: %w, len(members) = %d", err, len(members)) } member := members[0]
tx, err := i.Db.Begin() if err != nil { return fmt.Errorf("failed to start txn: %w", err) } committed := false defer func() { if committed { return } tx.Rollback() i.Enforcer.E.LoadPolicy() }()
// remove record by rkey && update enforcer if err = db.RemoveSpindleMember( tx, orm.FilterEq("did", did), orm.FilterEq("rkey", rkey), ); err != nil { return fmt.Errorf("failed to remove from db: %w", err) }
// update enforcer err = i.Enforcer.RemoveSpindleMember(member.Instance, member.Subject.String()) if err != nil { return fmt.Errorf("failed to update ACLs: %w", err) }
if err = tx.Commit(); err != nil { return fmt.Errorf("failed to commit txn: %w", err) }
if err = i.Enforcer.E.SavePolicy(); err != nil { return fmt.Errorf("failed to save ACLs: %w", err) } committed = true
l.Info("removed spindle member") }
return nil}
func (i *Ingester) ingestSpindle(ctx context.Context, e *jmodels.Event) error { did := e.Did var err error
l := i.Logger.With("handler", "ingestSpindle") l = l.With("nsid", e.Commit.Collection)
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) record := tangled.Spindle{} err = json.Unmarshal(raw, &record) if err != nil { l.Error("invalid record", "err", err) return err }
instance := e.Commit.RKey
err := db.AddSpindle(i.Db, models.Spindle{ Owner: syntax.DID(did), Instance: instance, }) if err != nil { l.Error("failed to add spindle to db", "err", err, "instance", instance) return err }
if err := i.verifySpindle(ctx, instance, did); err != nil { l.Warn("failed to verify spindle", "instance", instance, "did", did, "err", err) }
return nil
case jmodels.CommitOperationDelete: instance := e.Commit.RKey
// get record from db first spindles, err := db.GetSpindles( ctx, i.Db, orm.FilterEq("owner", did), orm.FilterEq("instance", instance), ) if err != nil || len(spindles) != 1 { return fmt.Errorf("failed to get spindles: %w, len(spindles) = %d", err, len(spindles)) } spindle := spindles[0]
tx, err := i.Db.Begin() if err != nil { return err } defer func() { tx.Rollback() i.Enforcer.E.LoadPolicy() }()
// remove spindle members first err = db.RemoveSpindleMember( tx, orm.FilterEq("owner", did), orm.FilterEq("instance", instance), ) if err != nil { return err }
err = db.DeleteSpindle( tx, orm.FilterEq("owner", did), orm.FilterEq("instance", instance), ) if err != nil { return err }
if spindle.Verified != nil { err = i.Enforcer.RemoveSpindle(instance) if err != nil { return err } }
err = tx.Commit() if err != nil { return err }
err = i.Enforcer.E.SavePolicy() if err != nil { return err } }
return nil}
func (i *Ingester) ingestString(e *jmodels.Event) error { did := e.Did rkey := e.Commit.RKey
var err error
l := i.Logger.With("handler", "ingestString", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) l.Info("ingesting record")
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) record := tangled.String{} err = json.Unmarshal(raw, &record) if err != nil { l.Error("invalid record", "err", err) return err }
string := models.StringFromRecord(did, rkey, record)
if err = i.Validator.ValidateString(&string); err != nil { l.Error("invalid record", "err", err) return err }
if err = db.AddString(i.Db, string); err != nil { l.Error("failed to add string", "err", err) return err }
return nil
case jmodels.CommitOperationDelete: if err := db.DeleteString( i.Db, orm.FilterEq("did", did), orm.FilterEq("rkey", rkey), ); err != nil { l.Error("failed to delete", "err", err) return fmt.Errorf("failed to delete string record: %w", err) }
return nil }
return nil}
func (i *Ingester) ingestKnotMember(ctx context.Context, e *jmodels.Event) error { did := e.Did var err error
l := i.Logger.With("handler", "ingestKnotMember") l = l.With("nsid", e.Commit.Collection)
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) record := tangled.KnotMember{} err = json.Unmarshal(raw, &record) if err != nil { l.Error("invalid record", "err", err) return err }
// only knot owner can invite to knots ok, err := i.Enforcer.IsKnotInviteAllowed(did, record.Domain) if err != nil { return fmt.Errorf("failed to check invite permission: %w", err) } if !ok { if verifyErr := i.verifyKnot(ctx, record.Domain, did); verifyErr != nil { return fmt.Errorf("invite denied and verify failed: %w", verifyErr) } ok, err = i.Enforcer.IsKnotInviteAllowed(did, record.Domain) if err != nil { return fmt.Errorf("failed to re-check invite permission: %w", err) } if !ok { return fmt.Errorf("invite denied for did %s on knot %s", did, record.Domain) } }
memberId, err := i.IdResolver.ResolveIdent(ctx, record.Subject) if err != nil { return err }
if memberId.Handle.IsInvalidHandle() { return fmt.Errorf("invalid handle for member %s", record.Subject) }
existing, err := db.GetKnotMembers(i.Db, orm.FilterEq("did", did), orm.FilterEq("rkey", e.Commit.RKey), ) if err != nil { return fmt.Errorf("failed to look up existing member: %w", err) } if len(existing) > 1 { return fmt.Errorf("multiple knot members with rkey %s", e.Commit.RKey) }
tx, err := i.Db.Begin() if err != nil { return fmt.Errorf("failed to start txn: %w", err) } committed := false defer func() { if committed { return } tx.Rollback() i.Enforcer.E.LoadPolicy() }()
if len(existing) == 1 { prev := existing[0] if prev.Domain != record.Domain || prev.Subject != memberId.DID { if err = db.RemoveKnotMember(tx, orm.FilterEq("did", did), orm.FilterEq("rkey", e.Commit.RKey), ); err != nil { return fmt.Errorf("failed to remove stale row: %w", err) } if err = i.Enforcer.RemoveKnotMember(prev.Domain, prev.Subject.String()); err != nil { return fmt.Errorf("failed to remove stale ACL: %w", err) } } }
if err = db.AddKnotMember(tx, models.KnotMember{ Did: syntax.DID(did), Rkey: e.Commit.RKey, Domain: record.Domain, Subject: memberId.DID, }); err != nil { return fmt.Errorf("failed to add to db: %w", err) }
if err = i.Enforcer.AddKnotMember(record.Domain, memberId.DID.String()); err != nil { return fmt.Errorf("failed to update ACLs: %w", err) }
if err = tx.Commit(); err != nil { return fmt.Errorf("failed to commit txn: %w", err) }
if err = i.Enforcer.E.SavePolicy(); err != nil { return fmt.Errorf("failed to save ACLs: %w", err) } committed = true
l.Info("upserted knot member") case jmodels.CommitOperationDelete: rkey := e.Commit.RKey
members, err := db.GetKnotMembers( i.Db, orm.FilterEq("did", did), orm.FilterEq("rkey", rkey), ) if err != nil { return fmt.Errorf("failed to look up knot member with rkey %s: %w", rkey, err) } if len(members) == 0 { l.Info("knot member already removed", "rkey", rkey) return nil } if len(members) > 1 { return fmt.Errorf("multiple knot members with rkey %s", rkey) } member := members[0]
tx, err := i.Db.Begin() if err != nil { return fmt.Errorf("failed to start txn: %w", err) } committed := false defer func() { if committed { return } tx.Rollback() i.Enforcer.E.LoadPolicy() }()
if err = db.RemoveKnotMember( tx, orm.FilterEq("did", did), orm.FilterEq("rkey", rkey), ); err != nil { return fmt.Errorf("failed to remove from db: %w", err) }
if err = i.Enforcer.RemoveKnotMember(member.Domain, member.Subject.String()); err != nil { return fmt.Errorf("failed to update ACLs: %w", err) }
if err = tx.Commit(); err != nil { return fmt.Errorf("failed to commit txn: %w", err) }
if err = i.Enforcer.E.SavePolicy(); err != nil { return fmt.Errorf("failed to save ACLs: %w", err) } committed = true
l.Info("removed knot member") }
return nil}
func (i *Ingester) ingestKnot(ctx context.Context, e *jmodels.Event) error { did := e.Did var err error
l := i.Logger.With("handler", "ingestKnot") l = l.With("nsid", e.Commit.Collection)
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) record := tangled.Knot{} err = json.Unmarshal(raw, &record) if err != nil { l.Error("invalid record", "err", err) return err }
domain := e.Commit.RKey
err := db.AddKnot(i.Db, domain, did) if err != nil { l.Error("failed to add knot to db", "err", err, "domain", domain) return err }
if err := i.verifyKnot(ctx, domain, did); err != nil { l.Warn("failed to verify knot", "domain", domain, "did", did, "err", err) }
return nil
case jmodels.CommitOperationDelete: domain := e.Commit.RKey
// get record from db first registrations, err := db.GetRegistrations( i.Db, orm.FilterEq("domain", domain), orm.FilterEq("did", did), ) if err != nil { return fmt.Errorf("failed to get registration: %w", err) } if len(registrations) != 1 { return fmt.Errorf("got incorrect number of registrations: %d, expected 1", len(registrations)) } registration := registrations[0]
tx, err := i.Db.Begin() if err != nil { return err } defer func() { tx.Rollback() i.Enforcer.E.LoadPolicy() }()
err = db.RemoveKnotMember( tx, orm.FilterEq("did", did), orm.FilterEq("domain", domain), ) if err != nil { return err }
err = db.DeleteKnot( tx, orm.FilterEq("did", did), orm.FilterEq("domain", domain), ) if err != nil { return err }
err = db.RemoveReposByKnot(tx, domain) if err != nil { return err }
if registration.Registered != nil { err = i.Enforcer.RemoveKnot(domain) if err != nil { return err } }
err = tx.Commit() if err != nil { return err }
err = i.Enforcer.E.SavePolicy() if err != nil { return err } }
return nil}
const ( verifyAttempts = 4 verifyMinDelay = 1 * time.Second verifyMaxDelay = 5 * time.Second)
func (i *Ingester) verifyKnot(ctx context.Context, domain, did string) error { regs, err := db.GetRegistrations(i.Db, orm.FilterEq("domain", domain), orm.FilterEq("did", did), ) if err != nil { return fmt.Errorf("look up registration: %w", err) } if len(regs) != 1 { return fmt.Errorf("no registration for %s by %s", domain, did) } if regs[0].Registered != nil { return nil }
err = retry.Do( func() error { return serververify.RunVerification(ctx, domain, did, i.Config.Core.Dev) }, retry.Context(ctx), retry.Attempts(verifyAttempts), retry.Delay(verifyMinDelay), retry.MaxDelay(verifyMaxDelay), retry.DelayType(retry.BackOffDelay), retry.LastErrorOnly(true), ) if err != nil { return fmt.Errorf("verify: %w", err) } return serververify.MarkKnotVerified(i.Db, i.Enforcer, domain, did)}
func (i *Ingester) verifySpindle(ctx context.Context, instance, did string) error { spindles, err := db.GetSpindles(ctx, i.Db, orm.FilterEq("instance", instance), orm.FilterEq("owner", did), ) if err != nil { return fmt.Errorf("look up spindle: %w", err) } if len(spindles) != 1 { return fmt.Errorf("no spindle for %s by %s", instance, did) } if spindles[0].Verified != nil { return nil }
err = retry.Do( func() error { return serververify.RunVerification(ctx, instance, did, i.Config.Core.Dev) }, retry.Context(ctx), retry.Attempts(verifyAttempts), retry.Delay(verifyMinDelay), retry.MaxDelay(verifyMaxDelay), retry.DelayType(retry.BackOffDelay), retry.LastErrorOnly(true), ) if err != nil { return fmt.Errorf("verify: %w", err) } _, err = serververify.MarkSpindleVerified(i.Db, i.Enforcer, instance, did) return err}
const sweepConcurrency = 4
func (i *Ingester) SweepPendingVerifications() { l := i.Logger.With("handler", "SweepPendingVerifications")
var g errgroup.Group g.SetLimit(sweepConcurrency)
regs, err := db.GetRegistrations(i.Db, orm.FilterIs("registered", nil)) if err != nil { l.Error("failed to list unverified knots", "err", err) } else { for _, reg := range regs { g.Go(func() error { if err := i.verifyKnot(i.Ctx, reg.Domain, reg.ByDid); err != nil { l.Warn("verify knot failed", "domain", reg.Domain, "did", reg.ByDid, "err", err) } return nil }) } }
spindles, err := db.GetSpindles(i.Ctx, i.Db, orm.FilterIs("verified", nil)) if err != nil { l.Error("failed to list unverified spindles", "err", err) g.Wait() return } for _, s := range spindles { g.Go(func() error { if err := i.verifySpindle(i.Ctx, s.Instance, s.Owner.String()); err != nil { l.Warn("verify spindle failed", "instance", s.Instance, "owner", s.Owner, "err", err) } return nil }) } g.Wait()}
func (i *Ingester) ingestIssue(ctx context.Context, e *jmodels.Event) error { did := e.Did rkey := e.Commit.RKey
var err error
l := i.Logger.With("handler", "ingestIssue", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) l.Info("ingesting record")
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) record := tangled.RepoIssue{} err = json.Unmarshal(raw, &record) if err != nil { l.Error("invalid record", "err", err) return err }
issue := models.IssueFromRecord(did, rkey, record)
if issue.RepoDid == "" { return fmt.Errorf("issue record has no repo field") } if _, err := syntax.ParseDID(string(issue.RepoDid)); err != nil { return fmt.Errorf("issue record repo field is not a valid DID: %w", err) }
if err := i.Validator.ValidateIssue(&issue); err != nil { return fmt.Errorf("failed to validate issue: %w", err) }
if record.Repo != "" && !strings.HasPrefix(record.Repo, "did:") { repo, repoErr := db.GetRepoByAtUri(i.Db, record.Repo) if repoErr == nil && repo.RepoDid != "" { if enqErr := db.EnqueuePdsRecordMigration(ctx, i.Db, "add-repo-did", syntax.DID(did), syntax.NSID(tangled.RepoIssueNSID), syntax.RecordKey(e.Commit.RKey)); enqErr != nil { l.Warn("failed to enqueue PDS rewrite for issue", "err", enqErr, "did", did, "repoDid", repo.RepoDid) } } }
tx, err := i.Db.BeginTx(ctx, nil) if err != nil { l.Error("failed to begin transaction", "err", err) return err } defer tx.Rollback()
err = db.PutIssue(tx, &issue) if err != nil { l.Error("failed to create issue", "err", err) return err }
err = tx.Commit() if err != nil { l.Error("failed to commit txn", "err", err) return err }
return nil
case jmodels.CommitOperationDelete: tx, err := i.Db.BeginTx(ctx, nil) if err != nil { l.Error("failed to begin transaction", "err", err) return err } defer tx.Rollback()
if err := db.DeleteIssues( tx, did, rkey, ); err != nil { l.Error("failed to delete", "err", err) return fmt.Errorf("failed to delete issue record: %w", err) } if err := tx.Commit(); err != nil { l.Error("failed to commit txn", "err", err) return err }
return nil }
return nil}
func (i *Ingester) ingestPull(ctx context.Context, e *jmodels.Event) error { did := e.Did rkey := e.Commit.RKey
var err error
l := i.Logger.With("handler", "ingestPull", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) l.Info("ingesting record")
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) record := tangled.RepoPull{} err = json.Unmarshal(raw, &record) if err != nil { l.Error("invalid record", "err", err) return err }
ownerId, err := i.IdResolver.ResolveIdent(ctx, did) if err != nil { l.Error("failed to resolve did") return err }
// go through and fetch all blobs in parallel readers := make([]*io.ReadCloser, len(record.Rounds)) var mu sync.Mutex
g, gctx := errgroup.WithContext(ctx)
for idx, b := range record.Rounds { g.Go(func() error { // for some reason, a blob is empty if b.PatchBlob == nil { return fmt.Errorf("missing patchBlob in round %d", idx) }
ownerPds := ownerId.PDSEndpoint() url, _ := url.Parse(fmt.Sprintf("%s/xrpc/com.atproto.sync.getBlob", ownerPds)) q := url.Query() q.Set("cid", b.PatchBlob.Ref.String()) q.Set("did", did) url.RawQuery = q.Encode()
req, err := http.NewRequestWithContext(gctx, http.MethodGet, url.String(), nil) if err != nil { l.Error("failed to create request") return err } req.Header.Set("Content-Type", "application/json")
resp, err := http.DefaultClient.Do(req) if err != nil { l.Error("failed to make request") return err }
mu.Lock() readers[idx] = &resp.Body mu.Unlock()
return nil }) }
if err := g.Wait(); err != nil { for _, r := range readers { if r != nil && *r != nil { (*r).Close() } } return err }
defer func() { for _, r := range readers { if r != nil && *r != nil { (*r).Close() } } }()
pull, err := models.PullFromRecord(did, rkey, record, readers) if err != nil { return fmt.Errorf("failed to parse pull from record: %w", err) } if err := i.Validator.ValidatePull(pull); err != nil { return fmt.Errorf("failed to validate pull: %w", err) }
tx, err := i.Db.BeginTx(ctx, nil) if err != nil { l.Error("failed to begin transaction", "err", err) return err } defer tx.Rollback()
err = db.PutPull(tx, pull) if err != nil { l.Error("failed to create pull", "err", err) return err }
err = tx.Commit() if err != nil { l.Error("failed to commit txn", "err", err) return err }
return nil
case jmodels.CommitOperationDelete: tx, err := i.Db.BeginTx(ctx, nil) if err != nil { l.Error("failed to begin transaction", "err", err) return err } defer tx.Rollback()
if err := db.AbandonPulls( tx, orm.FilterEq("owner_did", did), orm.FilterEq("rkey", rkey), ); err != nil { l.Error("failed to abandon", "err", err) return fmt.Errorf("failed to abandon pull record: %w", err) } if err := tx.Commit(); err != nil { l.Error("failed to commit txn", "err", err) return err }
return nil }
return nil}
// ingestIssueComment ingests legacy sh.tangled.repo.issue.comment deletionsfunc (i *Ingester) ingestIssueComment(e *jmodels.Event) error { l := i.Logger.With("handler", "ingestIssueComment", "nsid", e.Commit.Collection, "did", e.Did, "rkey", e.Commit.RKey) l.Info("ingesting record")
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: // no-op. sh.tangled.repo.issue.comment is deprecated
case jmodels.CommitOperationDelete: if err := db.PurgeComments( i.Db, orm.FilterEq("did", e.Did), orm.FilterEq("collection", e.Commit.Collection), orm.FilterEq("rkey", e.Commit.RKey), ); err != nil { return fmt.Errorf("failed to delete comment record: %w", err) } }
return nil}
// ingestPullComment ingests legacy sh.tangled.repo.pull.comment deletionsfunc (i *Ingester) ingestPullComment(e *jmodels.Event) error { l := i.Logger.With("handler", "ingestPullComment", "nsid", e.Commit.Collection, "did", e.Did, "rkey", e.Commit.RKey) l.Info("ingesting record")
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: // no-op. sh.tangled.repo.pull.comment is deprecated
case jmodels.CommitOperationDelete: if err := db.PurgeComments( i.Db, orm.FilterEq("did", e.Did), orm.FilterEq("collection", e.Commit.Collection), orm.FilterEq("rkey", e.Commit.RKey), ); err != nil { return fmt.Errorf("failed to delete comment record: %w", err) } }
return nil}
func (i *Ingester) ingestComment(e *jmodels.Event) error { did := e.Did rkey := e.Commit.RKey cid := e.Commit.CID
var err error
l := i.Logger.With("handler", "ingestComment", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) l.Info("ingesting record")
ctx := context.Background()
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) record := tangled.FeedComment{} err = json.Unmarshal(raw, &record) if err != nil { return fmt.Errorf("invalid record: %w", err) }
comment, err := models.CommentFromRecord(syntax.DID(did), syntax.RecordKey(rkey), syntax.CID(cid), record) if err != nil { return fmt.Errorf("failed to parse comment from record: %w", err) }
if err := comment.Validate(); err != nil { return fmt.Errorf("failed to validate comment: %w", err) }
var references []syntax.ATURI if comment.Body.Original != nil { _, references = i.MentionsResolver.Resolve(ctx, *comment.Body.Original) }
tx, err := i.Db.Begin() if err != nil { return fmt.Errorf("failed to start transaction: %w", err) } defer tx.Rollback()
_, err = db.PutComment(tx, comment, references) if err != nil { return fmt.Errorf("failed to create comment: %w", err) }
if err := tx.Commit(); err != nil { return err }
case jmodels.CommitOperationDelete: if err := db.DeleteComments( i.Db, orm.FilterEq("did", did), orm.FilterEq("collection", e.Commit.Collection), orm.FilterEq("rkey", rkey), ); err != nil { return fmt.Errorf("failed to delete comment record: %w", err) }
return nil }
return nil}
func (i *Ingester) ingestLabelDefinition(e *jmodels.Event) error { did := e.Did rkey := e.Commit.RKey
var err error
l := i.Logger.With("handler", "ingestLabelDefinition", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) l.Info("ingesting record")
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) record := tangled.LabelDefinition{} err = json.Unmarshal(raw, &record) if err != nil { return fmt.Errorf("invalid record: %w", err) }
def, err := models.LabelDefinitionFromRecord(did, rkey, record) if err != nil { return fmt.Errorf("failed to parse labeldef from record: %w", err) }
if err := i.Validator.ValidateLabelDefinition(def); err != nil { return fmt.Errorf("failed to validate labeldef: %w", err) }
_, err = db.AddLabelDefinition(i.Db, def) if err != nil { return fmt.Errorf("failed to create labeldef: %w", err) }
return nil
case jmodels.CommitOperationDelete: if err := db.DeleteLabelDefinition( i.Db, orm.FilterEq("did", did), orm.FilterEq("rkey", rkey), ); err != nil { return fmt.Errorf("failed to delete labeldef record: %w", err) }
return nil }
return nil}
func (i *Ingester) ingestLabelOp(e *jmodels.Event) error { did := e.Did rkey := e.Commit.RKey
var err error
l := i.Logger.With("handler", "ingestLabelOp", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) l.Info("ingesting record")
switch e.Commit.Operation { case jmodels.CommitOperationCreate: raw := json.RawMessage(e.Commit.Record) record := tangled.LabelOp{} err = json.Unmarshal(raw, &record) if err != nil { return fmt.Errorf("invalid record: %w", err) }
subject := syntax.ATURI(record.Subject) collection := subject.Collection()
var repo *models.Repo switch collection { case tangled.RepoIssueNSID: i, err := db.GetIssues(i.Db, orm.FilterEq("at_uri", subject)) if err != nil || len(i) != 1 { return fmt.Errorf("failed to find subject: %w || subject count %d", err, len(i)) } repo = i[0].Repo case tangled.RepoPullNSID: p, err := db.GetPulls(i.Db, orm.FilterEq("at_uri", subject)) if err != nil || len(p) != 1 { return fmt.Errorf("failed to find subject: %w || subject count %d", err, len(p)) } repo = p[0].Repo default: return fmt.Errorf("unsupported label subject: %s", collection) }
actx, err := db.NewLabelApplicationCtx(i.Db, orm.FilterIn("at_uri", repo.Labels)) if err != nil { return fmt.Errorf("failed to build label application ctx: %w", err) }
ops := models.LabelOpsFromRecord(did, rkey, record)
for _, o := range ops { def, ok := actx.Defs[o.OperandKey] if !ok { return fmt.Errorf("failed to find label def for key: %s, expected: %q", o.OperandKey, slices.Collect(maps.Keys(actx.Defs))) } if err := i.Validator.ValidateLabelOp(def, repo, &o); err != nil { return fmt.Errorf("failed to validate labelop: %w", err) } }
tx, err := i.Db.Begin() if err != nil { return err } defer tx.Rollback()
for _, o := range ops { _, err = db.AddLabelOp(tx, &o) if err != nil { return fmt.Errorf("failed to add labelop: %w", err) } }
if err = tx.Commit(); err != nil { return err } }
return nil}