Monorepo for Tangled
Something went wrong. Try again.
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251package appview
import ( "bytes" "context" "database/sql" "encoding/json" "errors" "fmt" "io" "log/slog" "net/http" "net/url" "slices" "strings"
"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/knotacl" "tangled.org/core/appview/mentions" "tangled.org/core/appview/models" "tangled.org/core/appview/notify" "tangled.org/core/appview/serververify" "tangled.org/core/consts" "tangled.org/core/idresolver" "tangled.org/core/orm" "tangled.org/core/rbac" "tangled.org/core/repoverify")
type RepoPermissionChecker interface { HasRepoPermissionErr(ctx context.Context, repo *models.Repo, userDid, perm string) (bool, error)}
type Ingester struct { Ctx context.Context Db *db.DB Enforcer *rbac.Enforcer Acl RepoPermissionChecker IdResolver *idresolver.Resolver Cache *cache.Cache Config *config.Config Logger *slog.Logger 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: l = l.With( "nsid", e.Commit.Collection, "did", e.Did, "rkey", e.Commit.RKey, "op", e.Commit.Operation, ) switch e.Commit.Collection { case tangled.GraphFollowNSID: err = i.ingestFollow(e, l) case tangled.GraphVouchNSID: err = i.ingestVouch(ctx, e, l) case tangled.FeedStarNSID: err = i.ingestStar(ctx, e, l) case tangled.FeedReactionNSID: err = i.ingestReaction(e, l) case tangled.PublicKeyNSID: err = i.ingestPublicKey(e, l) case tangled.RepoArtifactNSID: err = i.ingestArtifact(ctx, e, l) case tangled.ActorProfileNSID: err = i.ingestProfile(ctx, e, l) case tangled.SpindleMemberNSID: err = i.ingestSpindleMember(ctx, e, l) case tangled.SpindleNSID: err = i.ingestSpindle(ctx, e, l) case tangled.KnotMemberNSID: err = i.ingestKnotMember(ctx, e, l) case tangled.KnotNSID: err = i.ingestKnot(ctx, e, l) case tangled.StringNSID: err = i.ingestString(e, l) case tangled.RepoIssueNSID: err = i.ingestIssue(ctx, e, l) case tangled.RepoIssueStateNSID: err = i.ingestState(ctx, e, l, issueStateSpec) case tangled.RepoPullNSID: err = i.ingestPull(ctx, e, l) case tangled.RepoPullStatusNSID: err = i.ingestState(ctx, e, l, pullStatusSpec) case tangled.FeedCommentNSID: err = i.ingestComment(e, l) case tangled.RepoIssueCommentNSID: err = i.ingestIssueComment(e, l) case tangled.RepoPullCommentNSID: err = i.ingestPullComment(e, l) case tangled.LabelDefinitionNSID: err = i.ingestLabelDefinition(e, l) case tangled.LabelOpNSID: err = i.ingestLabelOp(ctx, e, l) case tangled.RepoNSID: err = i.ingestRepo(ctx, e, l) } }
if err != nil { l.Warn("failed to ingest record, skipping", "err", err) }
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, l *slog.Logger) error { var err error did := e.Did
l = l.With("handler", "ingestStar")
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) record := tangled.FeedStar{} unmarshalErr := json.Unmarshal(raw, &record)
createdAt, parseErr := time.Parse(time.RFC3339, record.CreatedAt) if parseErr != nil { createdAt = time.Now() }
star := models.Star{ Did: did, Created: createdAt, }
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.UpsertStar(i.Db, e.Commit.RKey, 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) } l.Info("processed star", "operation", e.Commit.Operation, "rkey", e.Commit.RKey)
l.Info("ingested record") return nil}
func (i *Ingester) ingestFollow(e *jmodels.Event, l *slog.Logger) error { var err error did := e.Did
l = l.With("handler", "ingestFollow")
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 := syntax.ParseDID(record.Subject) if err != nil { l.Error("invalid record. subject is invalid DID", "err", err) return err }
followedAt, err := time.Parse(time.RFC3339, record.CreatedAt) if err != nil { err = fmt.Errorf("createdAt is invalid datetime: %w", err) l.Error("invalid record", "err", err) return err }
err = db.UpsertFollow(i.Db, e.Commit.RKey, models.Follow{ UserDid: did, SubjectDid: record.Subject, FollowedAt: followedAt, }) 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) } l.Info("processed follow", "operation", e.Commit.Operation, "rkey", e.Commit.RKey)
l.Info("ingested record") return nil}
func (i *Ingester) ingestVouch(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { var err error did := e.Did
l = l.With("handler", "ingestVouch")
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) }
l.Info("ingested record") return nil}
func (i *Ingester) ingestPublicKey(e *jmodels.Event, l *slog.Logger) error { did := e.Did var err error
l = l.With("handler", "ingestPublicKey")
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 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 } pubKey, err := models.PublicKeyFromRecord(syntax.DID(did), syntax.RecordKey(e.Commit.RKey), record) if err != nil { l.Error("invalid record", "err", err) return err } if err := pubKey.Validate(); err != nil { l.Error("invalid record", "err", err) return err }
err = db.UpsertPublicKey(i.Db, pubKey) 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) } l.Info("processed pubkey", "operation", e.Commit.Operation, "rkey", e.Commit.RKey)
l.Info("ingested record") return nil}
func (i *Ingester) ingestArtifact(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { did := e.Did var err error
l = l.With("handler", "ingestArtifact")
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") }
allowed, permErr := i.Acl.HasRepoPermissionErr(ctx, repo, did, "repo:push") if permErr != nil { l.Warn("ingesting artifact without permission check", "did", did, "repo", repo.RepoIdentifier(), "err", permErr) } else if !allowed { l.Info("skipping unauthorized artifact", "did", did, "repo", repo.RepoIdentifier()) return nil }
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, parseErr := time.Parse(time.RFC3339, record.CreatedAt) if parseErr != 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) }
l.Info("ingested record") return nil}
func (i *Ingester) ingestProfile(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { did := e.Did var err error
l = l.With("handler", "ingestProfile")
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, }
err = db.ValidateProfile(i.Db, &profile) if err != nil { return fmt.Errorf("invalid profile record: %w", err) }
err = db.UpsertProfile(i.Db, &profile) if err != nil { return fmt.Errorf("upserting profile: %w", err) }
if 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) }
l.Info("ingested record") return nil}
func (i *Ingester) ingestSpindleMember(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { did := e.Did var err error
l = l.With("handler", "ingestSpindleMember")
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, l *slog.Logger) error { did := e.Did var err error
l = l.With("handler", "ingestSpindle")
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) }
l.Info("ingested record", "instance", instance) 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 fmt.Errorf("failed to start txn: %w", 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 fmt.Errorf("failed to remove spindle members: %w", err) }
err = db.DeleteSpindle( tx, orm.FilterEq("owner", did), orm.FilterEq("instance", instance), ) if err != nil { return fmt.Errorf("failed to delete spindle: %w", err) }
if spindle.Verified != nil { err = i.Enforcer.RemoveSpindle(instance) if err != nil { return fmt.Errorf("failed to remove spindle from enforcer: %w", err) } }
err = tx.Commit() if err != nil { return fmt.Errorf("failed to commit txn: %w", err) }
err = i.Enforcer.E.SavePolicy() if err != nil { return fmt.Errorf("failed to save ACLs: %w", err) }
l.Info("ingested record", "instance", instance) }
return nil}
func (i *Ingester) ingestString(e *jmodels.Event, l *slog.Logger) error { did := e.Did rkey := e.Commit.RKey
var err error
l = l.With("handler", "ingestString")
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 = string.Validate(); 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 }
l.Info("ingested record") 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) }
l.Info("ingested record") return nil }
return nil}
func (i *Ingester) ingestKnotMember(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { did := e.Did var err error
l = l.With("handler", "ingestKnotMember")
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, l *slog.Logger) error { did := e.Did var err error
l = l.With("handler", "ingestKnot")
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) }
l.Info("ingested record", "domain", domain) 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 fmt.Errorf("failed to start txn: %w", 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 fmt.Errorf("failed to remove knot members: %w", err) }
err = db.DeleteKnot( tx, orm.FilterEq("did", did), orm.FilterEq("domain", domain), ) if err != nil { return fmt.Errorf("failed to delete knot: %w", err) }
l.Error("attempt to delete repos by knot", "knot", domain) // err = db.RemoveReposByKnot(tx, domain) // if err != nil { // return fmt.Errorf("failed to remove repos by knot: %w", err) // }
if registration.Registered != nil { err = i.Enforcer.RemoveKnot(domain) if err != nil { return fmt.Errorf("failed to remove knot from enforcer: %w", err) } }
err = tx.Commit() if err != nil { return fmt.Errorf("failed to commit txn: %w", err) }
err = i.Enforcer.E.SavePolicy() if err != nil { return fmt.Errorf("failed to save ACLs: %w", err) }
l.Info("ingested record", "domain", domain) }
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, l *slog.Logger) error { did := e.Did rkey := e.Commit.RKey
var err error
l = l.With("handler", "ingestIssue")
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 := issue.Validate(); 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 }
if err := db.ResolveIssueState(tx, issue.AtUri()); err != nil { l.Error("failed to resolve issue state", "err", err) return err }
err = tx.Commit() if err != nil { l.Error("failed to commit txn", "err", err) return err }
i.drainPendingState(ctx, issue.AtUri(), l)
l.Info("ingested record") 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 }
l.Info("ingested record") return nil }
return nil}
func (i *Ingester) ingestPull(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { did := e.Did rkey := e.Commit.RKey
var err error
l = l.With("handler", "ingestPull")
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", "err", err) return err }
// go through and fetch all blobs in parallel blobs := make([]io.Reader, len(record.Rounds))
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 } defer resp.Body.Close()
var buf bytes.Buffer if _, err := io.Copy(&buf, io.LimitReader(resp.Body, 16<<20)); err != nil { return fmt.Errorf("failed to read blob in round %d: %w", idx, err) } blobs[idx] = &buf
return nil }) }
if err := g.Wait(); err != nil { return err }
pull, err := models.PullFromRecord(did, rkey, record, blobs) if err != nil { return fmt.Errorf("failed to parse pull from record: %w", err) } if err := pull.Validate(); err != nil { return fmt.Errorf("failed to validate pull: %w", err) } if pull.DependentOn != nil { if err := func() error { dependentPull, err := db.GetPull( i.Db, orm.FilterEq("dependent_on", pull.DependentOn.String()), ) if errors.Is(err, sql.ErrNoRows) { return nil } if err != nil { return fmt.Errorf("failed to fetch pulls with same dependency: %w", err) } if dependentPull.AtUri() == pull.AtUri() { return nil } return fmt.Errorf("another pull already depends on %s, which would form a DAG, this is presently disallowed", pull.DependentOn.String()) }(); err != nil { return fmt.Errorf("failed to validate pull stack: %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 }
if err := db.ResolvePullStatus(tx, pull.AtUri()); err != nil { l.Error("failed to resolve pull status", "err", err) return err }
err = tx.Commit() if err != nil { l.Error("failed to commit txn", "err", err) return err }
i.drainPendingState(ctx, pull.AtUri(), l)
l.Info("ingested record") 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 }
l.Info("ingested record") return nil }
return nil}
func (i *Ingester) authorizeStateRecord(ctx context.Context, repo *models.Repo, subjectAuthorDid, recordAuthorDid string, l *slog.Logger) (bool, error) { if recordAuthorDid == subjectAuthorDid { return true, nil } if recordAuthorDid == consts.TangledDid { return true, nil }
ok, err := i.Acl.HasRepoPermissionErr(ctx, repo, recordAuthorDid, "repo:push") if err != nil { if errors.Is(err, knotacl.ErrKnotUnreachable) { l.Warn("ingesting state record without permission check", "did", recordAuthorDid, "err", err) return true, nil } return false, err } return ok, nil}
type stateIngestSpec struct { subjectNSID string parse func(did, rkey string, raw json.RawMessage) (models.StateRecord, error) findSubject func(e db.Execer, subject syntax.ATURI) (repo *models.Repo, authorDid string, found bool, err error) put func(tx *sql.Tx, rec models.StateRecord) (syntax.ATURI, error) resolve func(tx *sql.Tx, subject syntax.ATURI) error recompute func(tx *sql.Tx, subject syntax.ATURI) error del func(tx *sql.Tx, did, rkey string) (syntax.ATURI, error)}
var issueStateSpec = stateIngestSpec{ subjectNSID: tangled.RepoIssueNSID, parse: func(did, rkey string, raw json.RawMessage) (models.StateRecord, error) { record := tangled.RepoIssueState{} if err := json.Unmarshal(raw, &record); err != nil { return models.StateRecord{}, fmt.Errorf("invalid record: %w", err) } return models.IssueStateFromRecord(did, rkey, record) }, findSubject: func(e db.Execer, subject syntax.ATURI) (*models.Repo, string, bool, error) { issues, err := db.GetIssues(e, orm.FilterEq("at_uri", subject)) if err != nil { return nil, "", false, err } if len(issues) != 1 || issues[0].Repo == nil { return nil, "", false, nil } return issues[0].Repo, issues[0].Did, true, nil }, put: db.PutIssueState, resolve: db.ResolveIssueState, recompute: db.RecomputeIssueState, del: db.DeleteIssueState,}
var pullStatusSpec = stateIngestSpec{ subjectNSID: tangled.RepoPullNSID, parse: func(did, rkey string, raw json.RawMessage) (models.StateRecord, error) { record := tangled.RepoPullStatus{} if err := json.Unmarshal(raw, &record); err != nil { return models.StateRecord{}, fmt.Errorf("invalid record: %w", err) } return models.PullStatusFromRecord(did, rkey, record) }, findSubject: func(e db.Execer, subject syntax.ATURI) (*models.Repo, string, bool, error) { pulls, err := db.GetPulls(e, orm.FilterEq("at_uri", subject)) if err != nil { return nil, "", false, err } if len(pulls) != 1 || pulls[0].Repo == nil { return nil, "", false, nil } return pulls[0].Repo, pulls[0].OwnerDid, true, nil }, put: db.PutPullStatus, resolve: db.ResolvePullStatus, recompute: db.RecomputePullStatus, del: db.DeletePullStatus,}
func (i *Ingester) ingestState(ctx context.Context, e *jmodels.Event, l *slog.Logger, spec stateIngestSpec) error { did := e.Did rkey := e.Commit.RKey nsid := e.Commit.Collection
l = l.With("handler", "ingestState", "nsid", nsid)
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: return i.applyStateRecord(ctx, did, rkey, nsid, e.Commit.Record, spec, l) case jmodels.CommitOperationDelete: return i.deleteStateRecord(ctx, did, rkey, nsid, spec, l) }
return nil}
func (i *Ingester) applyStateRecord(ctx context.Context, did, rkey, nsid string, raw []byte, spec stateIngestSpec, l *slog.Logger) error { rec, err := spec.parse(did, rkey, json.RawMessage(raw)) if err != nil { return err } if string(rec.Subject.Collection()) != spec.subjectNSID { return fmt.Errorf("state subject is not %s: %s", spec.subjectNSID, rec.Subject) }
repo, authorDid, found, err := spec.findSubject(i.Db, rec.Subject) if err != nil { return fmt.Errorf("failed to look up state subject: %w", err) } if !found { return i.parkStateRecord(ctx, did, rkey, nsid, rec.Subject, raw, l) }
authorized, err := i.authorizeStateRecord(ctx, repo, authorDid, did, l) if err != nil { return i.parkStateRecord(ctx, did, rkey, nsid, rec.Subject, raw, l) }
tx, err := i.Db.BeginTx(ctx, nil) if err != nil { return err } defer tx.Rollback()
if err := db.UnparkStateRecord(tx, did, rkey, nsid); err != nil { return fmt.Errorf("failed to unpark state record: %w", err) }
if !authorized { if err := tx.Commit(); err != nil { return err } l.Warn("dropped unauthorized state record", "did", did, "rkey", rkey, "subject", rec.Subject) return nil }
priorSubject, err := spec.put(tx, rec) if err != nil { return fmt.Errorf("failed to put state record: %w", err) } if err := spec.resolve(tx, rec.Subject); err != nil { return fmt.Errorf("failed to resolve state: %w", err) } if priorSubject != "" { if err := spec.recompute(tx, priorSubject); err != nil { return fmt.Errorf("failed to recompute prior subject state: %w", err) } }
if err := tx.Commit(); err != nil { return err }
l.Info("ingested record") return nil}
func (i *Ingester) deleteStateRecord(ctx context.Context, did, rkey, nsid string, spec stateIngestSpec, l *slog.Logger) error { tx, err := i.Db.BeginTx(ctx, nil) if err != nil { return err } defer tx.Rollback()
if err := db.UnparkStateRecord(tx, did, rkey, nsid); err != nil { return fmt.Errorf("failed to unpark state record: %w", err) } subject, err := spec.del(tx, did, rkey) if err != nil { return fmt.Errorf("failed to delete state record: %w", err) } if subject != "" { if err := spec.recompute(tx, subject); err != nil { return fmt.Errorf("failed to recompute state: %w", err) } }
if err := tx.Commit(); err != nil { return err }
l.Info("ingested record") return nil}
func (i *Ingester) parkStateRecord(ctx context.Context, did, rkey, nsid string, subject syntax.ATURI, raw []byte, l *slog.Logger) error { tx, err := i.Db.BeginTx(ctx, nil) if err != nil { return err } defer tx.Rollback()
if err := db.ParkStateRecord(tx, db.PendingStateRecord{ Did: did, Rkey: rkey, Nsid: nsid, Subject: subject, Record: raw, }); err != nil { return fmt.Errorf("failed to park state record: %w", err) }
if err := tx.Commit(); err != nil { return err }
l.Info("parked record for retry", "subject", subject, "nsid", nsid) return nil}
func (i *Ingester) drainPendingState(ctx context.Context, subject syntax.ATURI, l *slog.Logger) { pending, err := db.PendingStateRecordsForSubject(i.Db, subject) if err != nil { l.Error("failed to load pending records", "err", err, "subject", subject) return } for _, p := range pending { if err := i.reapplyPendingRecord(ctx, p, l); err != nil { l.Error("failed to drain pending record", "err", err, "did", p.Did, "rkey", p.Rkey, "nsid", p.Nsid) } }}
func (i *Ingester) drainPendingLabelOps(l *slog.Logger) { subjects, err := db.PendingStateSubjectsForNsid(i.Db, tangled.LabelOpNSID) if err != nil { l.Error("failed to list pending label op subjects", "err", err) return } for _, subject := range subjects { i.drainPendingState(i.Ctx, subject, l) }}
func (i *Ingester) reapplyPendingRecord(ctx context.Context, p db.PendingStateRecord, l *slog.Logger) error { switch p.Nsid { case tangled.RepoIssueStateNSID: return i.applyStateRecord(ctx, p.Did, p.Rkey, p.Nsid, p.Record, issueStateSpec, l) case tangled.RepoPullStatusNSID: return i.applyStateRecord(ctx, p.Did, p.Rkey, p.Nsid, p.Record, pullStatusSpec, l) case tangled.LabelOpNSID: return i.applyLabelOpRecord(ctx, p.Did, p.Rkey, p.Record, l) default: return fmt.Errorf("no reapply handler for parked nsid: %s", p.Nsid) }}
const ( pendingStateReconcileInterval = time.Hour pendingStateRecordTTL = 7 * 24 * time.Hour)
func (i *Ingester) StartPendingStateReconciler() { i.ReconcilePendingState()
ticker := time.NewTicker(pendingStateReconcileInterval) defer ticker.Stop() for { select { case <-i.Ctx.Done(): return case <-ticker.C: i.ReconcilePendingState() } }}
func (i *Ingester) ReconcilePendingState() { l := i.Logger.With("handler", "reconcilePendingState")
subjects, err := db.DistinctPendingStateSubjects(i.Db) if err != nil { l.Error("failed to list pending state subjects", "err", err) } for _, subject := range subjects { i.drainPendingState(i.Ctx, subject, l) }
cutoff := time.Now().Add(-pendingStateRecordTTL).UTC().Format(time.RFC3339) evicted, err := db.EvictStalePendingStateRecords(i.Db, cutoff) if err != nil { l.Error("failed to evict stale pending state records", "err", err) return } if evicted > 0 { l.Warn("evicted stale pending state records", "count", evicted, "olderThan", cutoff) }}
// ingestIssueComment ingests legacy sh.tangled.repo.issue.comment deletionsfunc (i *Ingester) ingestIssueComment(e *jmodels.Event, l *slog.Logger) error { l = l.With("handler", "ingestIssueComment")
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) } }
l.Info("ingested record") return nil}
// ingestPullComment ingests legacy sh.tangled.repo.pull.comment deletionsfunc (i *Ingester) ingestPullComment(e *jmodels.Event, l *slog.Logger) error { l = l.With("handler", "ingestPullComment")
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) } }
l.Info("ingested record") return nil}
func (i *Ingester) ingestComment(e *jmodels.Event, l *slog.Logger) error { did := e.Did rkey := e.Commit.RKey cid := e.Commit.CID
var err error
l = l.With("handler", "ingestComment")
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) } }
l.Info("ingested record") return nil}
func (i *Ingester) ingestReaction(e *jmodels.Event, l *slog.Logger) error { did := e.Did rkey := e.Commit.RKey
l = l.With("handler", "ingestReaction")
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) record := tangled.FeedReaction{} if err := json.Unmarshal(raw, &record); err != nil { return fmt.Errorf("invalid record: %w", err) }
subjectUri, err := syntax.ParseATURI(record.Subject) if err != nil { return fmt.Errorf("invalid reaction subject %q: %w", record.Subject, err) } subjectUri = models.NormalizeReactionSubject(subjectUri)
kind, ok := models.ParseReactionKind(record.Reaction) if !ok { return fmt.Errorf("invalid reaction kind: %q", record.Reaction) }
created, parseErr := time.Parse(time.RFC3339, record.CreatedAt) if parseErr != nil { created = time.Now() }
reaction := models.Reaction{ ReactedByDid: did, Rkey: rkey, ThreadAt: subjectUri, Kind: kind, Created: created, } if err := db.UpsertReaction(i.Db, reaction); err != nil { return fmt.Errorf("failed to upsert reaction: %w", err) }
case jmodels.CommitOperationDelete: if err := db.DeleteReactionByRkey(i.Db, did, rkey); err != nil { return fmt.Errorf("failed to delete reaction record: %w", err) } }
l.Info("ingested record") return nil}
func (i *Ingester) ingestLabelDefinition(e *jmodels.Event, l *slog.Logger) error { did := e.Did rkey := e.Commit.RKey
var err error
l = l.With("handler", "ingestLabelDefinition")
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 := def.Validate(); 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) }
if e.Commit.Operation == jmodels.CommitOperationCreate { i.drainPendingLabelOps(l) }
l.Info("ingested record") 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) }
l.Info("ingested record") return nil }
return nil}
func (i *Ingester) ingestLabelOp(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { l = l.With("handler", "ingestLabelOp") did := e.Did rkey := e.Commit.RKey
switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: return i.applyLabelOpRecord(ctx, did, rkey, e.Commit.Record, l) case jmodels.CommitOperationDelete: return i.deleteLabelOpRecord(ctx, did, rkey, l) }
return nil}
func (i *Ingester) findLabelSubjectRepo(subject syntax.ATURI) (*models.Repo, bool, error) { var spec stateIngestSpec switch subject.Collection() { case tangled.RepoIssueNSID: spec = issueStateSpec case tangled.RepoPullNSID: spec = pullStatusSpec default: return nil, false, fmt.Errorf("unsupported label subject: %s", subject.Collection()) } repo, _, found, err := spec.findSubject(i.Db, subject) return repo, found, err}
func (i *Ingester) applyLabelOpRecord(ctx context.Context, did, rkey string, raw []byte, l *slog.Logger) error { record := tangled.LabelOp{} if err := json.Unmarshal(raw, &record); err != nil { return fmt.Errorf("invalid record: %w", err) }
subject := syntax.ATURI(record.Subject) park := func() error { return i.parkStateRecord(ctx, did, rkey, tangled.LabelOpNSID, subject, raw, l) }
repo, found, err := i.findLabelSubjectRepo(subject) if err != nil { return err } if !found { return park() }
// validate permissions: only collaborators can apply labels currently // // TODO: introduce a repo:triage permission allowed, permErr := i.Acl.HasRepoPermissionErr(ctx, repo, did, "repo:push") if permErr != nil { if !errors.Is(permErr, knotacl.ErrKnotUnreachable) { return park() } l.Warn("ingesting labelop without permission check", "did", did, "err", permErr) allowed = true }
if !allowed { if err := i.unparkLabelOp(ctx, did, rkey); err != nil { return err } l.Warn("dropped unauthorized label op", "did", did, "rkey", rkey, "subject", subject) return nil }
ops := models.LabelOpsFromRecord(did, rkey, record)
operandKeys := make([]string, 0, len(ops)) for idx := range ops { operandKeys = append(operandKeys, ops[idx].OperandKey) }
actx, err := db.NewLabelApplicationCtx(i.Db, orm.FilterIn("at_uri", operandKeys)) if err != nil { return fmt.Errorf("failed to build label application ctx: %w", err) }
for idx := range ops { def, ok := actx.Defs[ops[idx].OperandKey] if !ok { return park() } if err := def.ValidateOperandValue(&ops[idx]); err != nil { return fmt.Errorf("failed to validate labelop: %w", err) } }
if err := i.materializeLabelOps(ctx, did, rkey, ops); err != nil { return err }
l.Info("ingested record") return nil}
func (i *Ingester) unparkLabelOp(ctx context.Context, did, rkey string) error { tx, err := i.Db.BeginTx(ctx, nil) if err != nil { return err } defer tx.Rollback()
if err := db.UnparkStateRecord(tx, did, rkey, tangled.LabelOpNSID); err != nil { return fmt.Errorf("failed to unpark label op: %w", err) }
return tx.Commit()}
func (i *Ingester) materializeLabelOps(ctx context.Context, did, rkey string, ops []models.LabelOp) error { tx, err := i.Db.BeginTx(ctx, nil) if err != nil { return err } defer tx.Rollback()
if err := db.UnparkStateRecord(tx, did, rkey, tangled.LabelOpNSID); err != nil { return fmt.Errorf("failed to unpark label op: %w", err) } if err := db.DeleteLabelOps(tx, orm.FilterEq("did", did), orm.FilterEq("rkey", rkey)); err != nil { return fmt.Errorf("failed to clear prior label ops: %w", err) } for idx := range ops { if _, err := db.AddLabelOp(tx, &ops[idx]); err != nil { return fmt.Errorf("failed to add labelop: %w", err) } } return tx.Commit()}
func (i *Ingester) deleteLabelOpRecord(ctx context.Context, did, rkey string, l *slog.Logger) error { if err := i.materializeLabelOps(ctx, did, rkey, nil); err != nil { return err } l.Info("ingested record") return nil}