This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / appview / state / knotstream.go
10 kB 377 lines
1package state 2 3import ( 4 "context" 5 "database/sql" 6 "encoding/json" 7 "errors" 8 "fmt" 9 "slices" 10 "time" 11 12 "tangled.org/core/appview/cloudflare" 13 "tangled.org/core/appview/notify" 14 15 "tangled.org/core/api/tangled" 16 "tangled.org/core/appview/cache" 17 "tangled.org/core/appview/config" 18 "tangled.org/core/appview/db" 19 "tangled.org/core/appview/models" 20 "tangled.org/core/appview/sites" 21 ec "tangled.org/core/eventconsumer" 22 "tangled.org/core/eventconsumer/cursor" 23 "tangled.org/core/log" 24 "tangled.org/core/orm" 25 "tangled.org/core/rbac" 26 "tangled.org/core/workflow" 27 28 "github.com/bluesky-social/indigo/atproto/syntax" 29 "github.com/go-git/go-git/v5/plumbing" 30 "github.com/posthog/posthog-go" 31) 32 33func Knotstream(ctx context.Context, c *config.Config, d *db.DB, enforcer *rbac.Enforcer, posthog posthog.Client, notifier notify.Notifier, cfClient *cloudflare.Client) (*ec.Consumer, error) { 34 logger := log.FromContext(ctx) 35 logger = log.SubLogger(logger, "knotstream") 36 37 knots, err := db.GetRegistrations( 38 d, 39 orm.FilterIsNot("registered", "null"), 40 ) 41 if err != nil { 42 return nil, err 43 } 44 45 srcs := make(map[ec.Source]struct{}) 46 for _, k := range knots { 47 s := ec.NewKnotSource(k.Domain) 48 srcs[s] = struct{}{} 49 } 50 51 cache := cache.New(c.Redis.Addr) 52 cursorStore := cursor.NewRedisCursorStore(cache) 53 54 cfg := ec.ConsumerConfig{ 55 Sources: srcs, 56 ProcessFunc: knotIngester(d, enforcer, posthog, notifier, c.Core.Dev, c, cfClient), 57 RetryInterval: c.Knotstream.RetryInterval, 58 MaxRetryInterval: c.Knotstream.MaxRetryInterval, 59 ConnectionTimeout: c.Knotstream.ConnectionTimeout, 60 WorkerCount: c.Knotstream.WorkerCount, 61 QueueSize: c.Knotstream.QueueSize, 62 Logger: logger, 63 Dev: c.Core.Dev, 64 CursorStore: &cursorStore, 65 } 66 67 return ec.NewConsumer(cfg), nil 68} 69 70func resolveRepo(d *db.DB, repoDid *string, ownerDid, repoName string) (*models.Repo, error) { 71 if repoDid != nil && *repoDid != "" { 72 return db.GetRepoByDid(d, *repoDid) 73 } 74 repos, err := db.GetRepos(d, orm.FilterEq("did", ownerDid), orm.FilterEq("name", repoName)) 75 if err != nil { 76 return nil, err 77 } 78 if len(repos) == 0 { 79 return nil, sql.ErrNoRows 80 } 81 return &repos[0], nil 82} 83 84func knotIngester(d *db.DB, enforcer *rbac.Enforcer, posthog posthog.Client, notifier notify.Notifier, dev bool, c *config.Config, cfClient *cloudflare.Client) ec.ProcessFunc { 85 return func(ctx context.Context, source ec.Source, msg ec.Message) error { 86 switch msg.Nsid { 87 case tangled.GitRefUpdateNSID: 88 return ingestRefUpdate(ctx, d, enforcer, posthog, notifier, dev, c, cfClient, source, msg) 89 case tangled.PipelineNSID: 90 return ingestPipeline(d, source, msg) 91 } 92 93 return nil 94 } 95} 96 97func ingestRefUpdate(ctx context.Context, d *db.DB, enforcer *rbac.Enforcer, pc posthog.Client, notifier notify.Notifier, dev bool, c *config.Config, cfClient *cloudflare.Client, source ec.Source, msg ec.Message) error { 98 logger := log.FromContext(ctx) 99 100 var record tangled.GitRefUpdate 101 err := json.Unmarshal(msg.EventJson, &record) 102 if err != nil { 103 return err 104 } 105 106 knownKnots, err := enforcer.GetKnotsForUser(record.CommitterDid) 107 if err != nil { 108 return err 109 } 110 if !slices.Contains(knownKnots, source.Key()) { 111 return fmt.Errorf("%s does not belong to %s, something is fishy", record.CommitterDid, source.Key()) 112 } 113 114 ownerDid := "" 115 if record.OwnerDid != nil { 116 ownerDid = *record.OwnerDid 117 } 118 119 repo, lookupErr := resolveRepo(d, record.RepoDid, ownerDid, record.RepoName) 120 if lookupErr != nil { 121 return fmt.Errorf("failed to look up repo: %w", lookupErr) 122 } 123 124 logger.Info("processing gitRefUpdate event", 125 "repo", repo.RepoIdentifier(), 126 "ref", record.Ref, 127 "old_sha", record.OldSha, 128 "new_sha", record.NewSha) 129 130 notifier.Push(ctx, repo, record.Ref, record.OldSha, record.NewSha, record.CommitterDid) 131 132 errPunchcard := populatePunchcard(d, record) 133 errLanguages := updateRepoLanguages(d, record) 134 135 var errPosthog error 136 if !dev && record.CommitterDid != "" { 137 errPosthog = pc.Enqueue(posthog.Capture{ 138 DistinctId: record.CommitterDid, 139 Event: "git_ref_update", 140 }) 141 } 142 143 // Trigger a sites redeploy if this push is to the configured sites branch. 144 if cfClient.Enabled() { 145 go triggerSitesDeployIfNeeded(ctx, d, cfClient, c, record, source) 146 } 147 148 return errors.Join(errPunchcard, errLanguages, errPosthog) 149} 150 151// triggerSitesDeployIfNeeded checks whether the pushed ref matches the sites 152// branch configured for this repo and, if so, syncs the site to R2 153func triggerSitesDeployIfNeeded(ctx context.Context, d *db.DB, cfClient *cloudflare.Client, c *config.Config, record tangled.GitRefUpdate, source ec.Source) { 154 logger := log.FromContext(ctx) 155 156 ref := plumbing.ReferenceName(record.Ref) 157 if !ref.IsBranch() { 158 return 159 } 160 pushedBranch := ref.Short() 161 162 ownerDid := "" 163 if record.OwnerDid != nil { 164 ownerDid = *record.OwnerDid 165 } 166 167 repo, err := resolveRepo(d, record.RepoDid, ownerDid, record.RepoName) 168 if err != nil { 169 return 170 } 171 172 siteConfig, err := db.GetRepoSiteConfig(d, repo.RepoAt().String()) 173 if err != nil || siteConfig == nil { 174 return 175 } 176 if siteConfig.Branch != pushedBranch { 177 return 178 } 179 180 scheme := "https" 181 if c.Core.Dev { 182 scheme = "http" 183 } 184 knotHost := fmt.Sprintf("%s://%s", scheme, source.Key()) 185 186 deploy := &models.SiteDeploy{ 187 RepoAt: repo.RepoAt().String(), 188 Branch: siteConfig.Branch, 189 Dir: siteConfig.Dir, 190 CommitSHA: record.NewSha, 191 Trigger: models.SiteDeployTriggerPush, 192 } 193 194 deployErr := sites.Deploy(ctx, cfClient, knotHost, repo.RepoIdentifier(), record.RepoName, siteConfig.Branch, siteConfig.Dir) 195 if deployErr != nil { 196 logger.Error("sites: R2 sync failed on push", "repo", repo.RepoIdentifier(), "err", deployErr) 197 deploy.Status = models.SiteDeployStatusFailure 198 deploy.Error = deployErr.Error() 199 } else { 200 deploy.Status = models.SiteDeployStatusSuccess 201 } 202 203 if err := db.AddSiteDeploy(d, deploy); err != nil { 204 logger.Error("sites: failed to record deploy", "repo", repo.RepoIdentifier(), "err", err) 205 } 206 207 if deployErr == nil { 208 logger.Info("site deployed to r2", "repo", repo.RepoIdentifier()) 209 } 210} 211 212func populatePunchcard(d *db.DB, record tangled.GitRefUpdate) error { 213 if record.CommitterDid == "" { 214 return nil 215 } 216 217 knownEmails, err := db.GetAllEmails(d, record.CommitterDid) 218 if err != nil { 219 return err 220 } 221 222 count := 0 223 for _, ke := range knownEmails { 224 if record.Meta == nil { 225 continue 226 } 227 if record.Meta.CommitCount == nil { 228 continue 229 } 230 for _, ce := range record.Meta.CommitCount.ByEmail { 231 if ce == nil { 232 continue 233 } 234 if ce.Email == ke.Address || ce.Email == record.CommitterDid { 235 count += int(ce.Count) 236 } 237 } 238 } 239 240 punch := models.Punch{ 241 Did: record.CommitterDid, 242 Date: time.Now(), 243 Count: count, 244 } 245 return db.AddPunch(d, punch) 246} 247 248func updateRepoLanguages(d *db.DB, record tangled.GitRefUpdate) error { 249 if record.Meta == nil || record.Meta.LangBreakdown == nil || record.Meta.LangBreakdown.Inputs == nil { 250 return fmt.Errorf("empty language data for repo: %v/%s", record.OwnerDid, record.RepoName) 251 } 252 253 ownerDid := "" 254 if record.OwnerDid != nil { 255 ownerDid = *record.OwnerDid 256 } 257 258 r, lookupErr := resolveRepo(d, record.RepoDid, ownerDid, record.RepoName) 259 if lookupErr != nil { 260 return fmt.Errorf("failed to look up repo: %w", lookupErr) 261 } 262 repo := *r 263 264 ref := plumbing.ReferenceName(record.Ref) 265 if !ref.IsBranch() { 266 return fmt.Errorf("%s is not a valid reference name", ref) 267 } 268 269 var langs []models.RepoLanguage 270 for _, l := range record.Meta.LangBreakdown.Inputs { 271 if l == nil { 272 continue 273 } 274 275 langs = append(langs, models.RepoLanguage{ 276 RepoAt: repo.RepoAt(), 277 Ref: ref.Short(), 278 IsDefaultRef: record.Meta.IsDefaultRef, 279 Language: l.Lang, 280 Bytes: l.Size, 281 }) 282 } 283 284 tx, err := d.Begin() 285 if err != nil { 286 return err 287 } 288 defer tx.Rollback() 289 290 // update appview's cache 291 err = db.UpdateRepoLanguages(tx, repo.RepoAt(), ref.Short(), langs) 292 if err != nil { 293 fmt.Printf("failed; %s\n", err) 294 // non-fatal 295 } 296 297 return tx.Commit() 298} 299 300func ingestPipeline(d *db.DB, source ec.Source, msg ec.Message) error { 301 var record tangled.Pipeline 302 err := json.Unmarshal(msg.EventJson, &record) 303 if err != nil { 304 return err 305 } 306 307 if record.TriggerMetadata == nil { 308 return fmt.Errorf("empty trigger metadata: nsid %s, rkey %s", msg.Nsid, msg.Rkey) 309 } 310 311 if record.TriggerMetadata.Repo == nil { 312 return fmt.Errorf("empty repo: nsid %s, rkey %s", msg.Nsid, msg.Rkey) 313 } 314 315 repoName := "" 316 if record.TriggerMetadata.Repo.Repo != nil { 317 repoName = *record.TriggerMetadata.Repo.Repo 318 } 319 320 repo, lookupErr := resolveRepo(d, record.TriggerMetadata.Repo.RepoDid, record.TriggerMetadata.Repo.Did, repoName) 321 if lookupErr != nil { 322 return fmt.Errorf("failed to look up repo: %w", lookupErr) 323 } 324 if repo.Spindle == "" { 325 return fmt.Errorf("repo does not have a spindle configured yet: nsid %s, rkey %s", msg.Nsid, msg.Rkey) 326 } 327 328 // trigger info 329 var trigger models.Trigger 330 var sha string 331 trigger.Kind = workflow.TriggerKind(record.TriggerMetadata.Kind) 332 switch trigger.Kind { 333 case workflow.TriggerKindPush: 334 trigger.PushRef = &record.TriggerMetadata.Push.Ref 335 trigger.PushNewSha = &record.TriggerMetadata.Push.NewSha 336 trigger.PushOldSha = &record.TriggerMetadata.Push.OldSha 337 sha = *trigger.PushNewSha 338 case workflow.TriggerKindPullRequest: 339 trigger.PRSourceBranch = &record.TriggerMetadata.PullRequest.SourceBranch 340 trigger.PRTargetBranch = &record.TriggerMetadata.PullRequest.TargetBranch 341 trigger.PRSourceSha = &record.TriggerMetadata.PullRequest.SourceSha 342 trigger.PRAction = &record.TriggerMetadata.PullRequest.Action 343 sha = *trigger.PRSourceSha 344 } 345 346 tx, err := d.Begin() 347 if err != nil { 348 return fmt.Errorf("failed to start txn: %w", err) 349 } 350 351 triggerId, err := db.AddTrigger(tx, trigger) 352 if err != nil { 353 return fmt.Errorf("failed to add trigger entry: %w", err) 354 } 355 356 pipeline := models.Pipeline{ 357 Rkey: msg.Rkey, 358 Knot: source.Key(), 359 RepoOwner: syntax.DID(record.TriggerMetadata.Repo.Did), 360 RepoName: repoName, 361 RepoDid: repo.RepoDid, 362 TriggerId: int(triggerId), 363 Sha: sha, 364 } 365 366 err = db.AddPipeline(tx, pipeline) 367 if err != nil { 368 return fmt.Errorf("failed to add pipeline: %w", err) 369 } 370 371 err = tx.Commit() 372 if err != nil { 373 return fmt.Errorf("failed to commit txn: %w", err) 374 } 375 376 return nil 377}