From 823913c50199d9dcfcc96c5c1a1000d9d08882bf Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Mon, 5 Feb 2024 17:53:06 -0800 Subject: [PATCH] automod: core engine metrics --- automod/engine/blobs.go | 8 ++++ automod/engine/engine.go | 44 ++++++++++++++++++++ automod/engine/fetch_account_meta.go | 3 ++ automod/engine/fetch_relationship.go | 1 + automod/engine/metrics.go | 61 ++++++++++++++++++++++++++++ automod/engine/persist.go | 18 ++++++++ automod/engine/persisthelpers.go | 4 +- 7 files changed, 138 insertions(+), 1 deletion(-) create mode 100644 automod/engine/metrics.go diff --git a/automod/engine/blobs.go b/automod/engine/blobs.go index 12432cef..7de7cb4e 100644 --- a/automod/engine/blobs.go +++ b/automod/engine/blobs.go @@ -5,6 +5,7 @@ import ( "io" "net/http" "strings" + "time" appbsky "github.com/bluesky-social/indigo/api/bsky" lexutil "github.com/bluesky-social/indigo/lex/util" @@ -94,6 +95,12 @@ func (c *RecordContext) Blobs() ([]lexutil.LexBlob, error) { func (c *RecordContext) fetchBlob(blob lexutil.LexBlob) ([]byte, error) { + start := time.Now() + defer func() { + duration := time.Since(start) + blobDownloadDuration.Observe(duration.Seconds()) + }() + var blobBytes []byte // TODO: better way to do this, eg a shared client? @@ -121,6 +128,7 @@ func (c *RecordContext) fetchBlob(blob lexutil.LexBlob) ([]byte, error) { } defer resp.Body.Close() + blobDownloadCount.WithLabelValues(fmt.Sprint(resp.StatusCode)).Inc() if resp.StatusCode != 200 { return nil, fmt.Errorf("failed to fetch blob from PDS. did=%s cid=%s statusCode=%d", c.Account.Identity.DID, blob.Ref, resp.StatusCode) } diff --git a/automod/engine/engine.go b/automod/engine/engine.go index 09cf6d47..fc6e73b0 100644 --- a/automod/engine/engine.go +++ b/automod/engine/engine.go @@ -44,10 +44,18 @@ type Engine struct { // // This method can be called concurrently, though cached state may end up inconsistent if multiple events for the same account (DID) are processed in parallel. func (eng *Engine) ProcessIdentityEvent(ctx context.Context, typ string, did syntax.DID) error { + eventProcessCount.WithLabelValues("identity").Inc() + start := time.Now() + defer func() { + duration := time.Since(start) + eventProcessDuration.WithLabelValues("identity").Observe(duration.Seconds()) + }() + // similar to an HTTP server, we want to recover any panics from rule execution defer func() { if r := recover(); r != nil { eng.Logger.Error("automod event execution exception", "err", r, "did", did, "type", typ) + eventErrorCount.WithLabelValues("identity").Inc() } }() ctx, cancel := context.WithTimeout(ctx, identityEventTimeout) @@ -59,25 +67,31 @@ func (eng *Engine) ProcessIdentityEvent(ctx context.Context, typ string, did syn } ident, err := eng.Directory.LookupDID(ctx, did) if err != nil { + eventErrorCount.WithLabelValues("identity").Inc() return fmt.Errorf("resolving identity: %w", err) } if ident == nil { + eventErrorCount.WithLabelValues("identity").Inc() return fmt.Errorf("identity not found for DID: %s", did.String()) } am, err := eng.GetAccountMeta(ctx, ident) if err != nil { + eventErrorCount.WithLabelValues("identity").Inc() return fmt.Errorf("failed to fetch account metadata: %w", err) } ac := NewAccountContext(ctx, eng, *am) if err := eng.Rules.CallIdentityRules(&ac); err != nil { + eventErrorCount.WithLabelValues("identity").Inc() return fmt.Errorf("rule execution failed: %w", err) } eng.CanonicalLogLineAccount(&ac) if err := eng.persistAccountModActions(&ac); err != nil { + eventErrorCount.WithLabelValues("identity").Inc() return fmt.Errorf("failed to persist actions for identity event: %w", err) } if err := eng.persistCounters(ctx, ac.effects); err != nil { + eventErrorCount.WithLabelValues("identity").Inc() return fmt.Errorf("failed to persist counters for identity event: %w", err) } return nil @@ -87,6 +101,13 @@ func (eng *Engine) ProcessIdentityEvent(ctx context.Context, typ string, did syn // // This method can be called concurrently, though cached state may end up inconsistent if multiple events for the same account (DID) are processed in parallel. func (eng *Engine) ProcessRecordOp(ctx context.Context, op RecordOp) error { + eventProcessCount.WithLabelValues("record").Inc() + start := time.Now() + defer func() { + duration := time.Since(start) + eventProcessDuration.WithLabelValues("record").Observe(duration.Seconds()) + }() + // similar to an HTTP server, we want to recover any panics from rule execution defer func() { if r := recover(); r != nil { @@ -97,18 +118,22 @@ func (eng *Engine) ProcessRecordOp(ctx context.Context, op RecordOp) error { defer cancel() if err := op.Validate(); err != nil { + eventErrorCount.WithLabelValues("record").Inc() return fmt.Errorf("bad record op: %w", err) } ident, err := eng.Directory.LookupDID(ctx, op.DID) if err != nil { + eventErrorCount.WithLabelValues("record").Inc() return fmt.Errorf("resolving identity: %w", err) } if ident == nil { + eventErrorCount.WithLabelValues("record").Inc() return fmt.Errorf("identity not found for DID: %s", op.DID) } am, err := eng.GetAccountMeta(ctx, ident) if err != nil { + eventErrorCount.WithLabelValues("record").Inc() return fmt.Errorf("failed to fetch account metadata: %w", err) } rc := NewRecordContext(ctx, eng, *am, op) @@ -116,13 +141,16 @@ func (eng *Engine) ProcessRecordOp(ctx context.Context, op RecordOp) error { switch op.Action { case CreateOp, UpdateOp: if err := eng.Rules.CallRecordRules(&rc); err != nil { + eventErrorCount.WithLabelValues("record").Inc() return fmt.Errorf("rule execution failed: %w", err) } case DeleteOp: if err := eng.Rules.CallRecordDeleteRules(&rc); err != nil { + eventErrorCount.WithLabelValues("record").Inc() return fmt.Errorf("rule execution failed: %w", err) } default: + eventErrorCount.WithLabelValues("record").Inc() return fmt.Errorf("unexpected op action: %s", op.Action) } eng.CanonicalLogLineRecord(&rc) @@ -133,9 +161,11 @@ func (eng *Engine) ProcessRecordOp(ctx context.Context, op RecordOp) error { } } if err := eng.persistRecordModActions(&rc); err != nil { + eventErrorCount.WithLabelValues("record").Inc() return fmt.Errorf("failed to persist actions for record event: %w", err) } if err := eng.persistCounters(ctx, rc.effects); err != nil { + eventErrorCount.WithLabelValues("record").Inc() return fmt.Errorf("failed to persist counts for record event: %w", err) } return nil @@ -143,6 +173,13 @@ func (eng *Engine) ProcessRecordOp(ctx context.Context, op RecordOp) error { // returns a boolean indicating "block the event" func (eng *Engine) ProcessNotificationEvent(ctx context.Context, senderDID, recipientDID syntax.DID, reason string, subject syntax.ATURI) (bool, error) { + eventProcessCount.WithLabelValues("notif").Inc() + start := time.Now() + defer func() { + duration := time.Since(start) + eventProcessDuration.WithLabelValues("notif").Observe(duration.Seconds()) + }() + // similar to an HTTP server, we want to recover any panics from rule execution defer func() { if r := recover(); r != nil { @@ -154,31 +191,38 @@ func (eng *Engine) ProcessNotificationEvent(ctx context.Context, senderDID, reci senderIdent, err := eng.Directory.LookupDID(ctx, senderDID) if err != nil { + eventErrorCount.WithLabelValues("notif").Inc() return false, fmt.Errorf("resolving identity: %w", err) } if senderIdent == nil { + eventErrorCount.WithLabelValues("notif").Inc() return false, fmt.Errorf("identity not found for sender DID: %s", senderDID.String()) } recipientIdent, err := eng.Directory.LookupDID(ctx, recipientDID) if err != nil { + eventErrorCount.WithLabelValues("notif").Inc() return false, fmt.Errorf("resolving identity: %w", err) } if recipientIdent == nil { + eventErrorCount.WithLabelValues("notif").Inc() return false, fmt.Errorf("identity not found for sender DID: %s", recipientDID.String()) } senderMeta, err := eng.GetAccountMeta(ctx, senderIdent) if err != nil { + eventErrorCount.WithLabelValues("notif").Inc() return false, fmt.Errorf("failed to fetch account metadata: %w", err) } recipientMeta, err := eng.GetAccountMeta(ctx, recipientIdent) if err != nil { + eventErrorCount.WithLabelValues("notif").Inc() return false, fmt.Errorf("failed to fetch account metadata: %w", err) } nc := NewNotificationContext(ctx, eng, *senderMeta, *recipientMeta, reason, subject) if err := eng.Rules.CallNotificationRules(&nc); err != nil { + eventErrorCount.WithLabelValues("notif").Inc() return false, fmt.Errorf("rule execution failed: %w", err) } eng.CanonicalLogLineNotification(&nc) diff --git a/automod/engine/fetch_account_meta.go b/automod/engine/fetch_account_meta.go index f2c1d7b3..3a2add8b 100644 --- a/automod/engine/fetch_account_meta.go +++ b/automod/engine/fetch_account_meta.go @@ -43,6 +43,9 @@ func (e *Engine) GetAccountMeta(ctx context.Context, ident *identity.Identity) ( return &am, nil } + // doing a "full" fetch from here on + accountMetaFetches.Inc() + flags, err := e.Flags.Get(ctx, ident.DID.String()) if err != nil { return nil, fmt.Errorf("failed checking account flag cache: %w", err) diff --git a/automod/engine/fetch_relationship.go b/automod/engine/fetch_relationship.go index 665a3e26..50226830 100644 --- a/automod/engine/fetch_relationship.go +++ b/automod/engine/fetch_relationship.go @@ -45,6 +45,7 @@ func (eng *Engine) GetAccountRelationship(ctx context.Context, primary, other sy } // fetch account relationship from AppView + accountRelationshipFetches.Inc() resp, err := appbsky.GraphGetRelationships(ctx, eng.BskyClient, primary.String(), []string{other.String()}) if err != nil || len(resp.Relationships) != 1 { logger.Warn("account relationship lookup failed", "err", err) diff --git a/automod/engine/metrics.go b/automod/engine/metrics.go new file mode 100644 index 00000000..bf71197d --- /dev/null +++ b/automod/engine/metrics.go @@ -0,0 +1,61 @@ +package engine + +import ( + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promauto" +) + +var eventProcessDuration = promauto.NewHistogramVec(prometheus.HistogramOpts{ + Name: "automod_event_duration_sec", + Help: "Total duration of automod event processing", +}, []string{"type"}) + +var eventProcessCount = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "automod_event_processed", + Help: "Number of events processed", +}, []string{"type"}) + +var eventErrorCount = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "automod_event_errors", + Help: "Number of events which failed processing", +}, []string{"type"}) + +var actionNewLabelCount = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "automod_new_action_labels", + Help: "Number of new labels persisted", +}, []string{"type", "val"}) + +var actionNewFlagCount = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "automod_new_action_flags", + Help: "Number of new flags persisted", +}, []string{"type", "val"}) + +var actionNewReportCount = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "automod_new_action_reports", + Help: "Number of new flags persisted", +}, []string{"type"}) + +var actionNewTakedownCount = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "automod_new_action_takedowns", + Help: "Number of new flags persisted", +}, []string{"type"}) + +var accountMetaFetches = promauto.NewCounter(prometheus.CounterOpts{ + Name: "automod_account_meta_fetches", + Help: "Number of account metadata reads (API calls)", +}) + +var accountRelationshipFetches = promauto.NewCounter(prometheus.CounterOpts{ + Name: "automod_account_relationship_fetches", + Help: "Number of account relationship reads (API calls)", +}) + +var blobDownloadCount = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "automod_blob_downloads", + Help: "Number of blobs downloaded, by HTTP status code", +}, []string{"status"}) + +var blobDownloadDuration = promauto.NewHistogram(prometheus.HistogramOpts{ + Name: "automod_blob_download_duration_sec", + Help: "Duration of blob download attempts", +}) diff --git a/automod/engine/persist.go b/automod/engine/persist.go index f6f8aff0..694444d1 100644 --- a/automod/engine/persist.go +++ b/automod/engine/persist.go @@ -68,6 +68,10 @@ func (eng *Engine) persistAccountModActions(c *AccountContext) error { // flags don't require admin auth if len(newFlags) > 0 { + for _, val := range newFlags { + // note: WithLabelValues is a prometheus label, not an atproto label + actionNewFlagCount.WithLabelValues("record", val).Inc() + } eng.Flags.Add(ctx, c.Account.Identity.DID.String(), newFlags) } @@ -83,6 +87,10 @@ func (eng *Engine) persistAccountModActions(c *AccountContext) error { if len(newLabels) > 0 { c.Logger.Info("labeling record", "newLabels", newLabels) + for _, val := range newLabels { + // note: WithLabelValues is a prometheus label, not an atproto label + actionNewLabelCount.WithLabelValues("account", val).Inc() + } comment := "[automod]: auto-labeling account" _, err := comatproto.AdminEmitModerationEvent(ctx, xrpcc, &comatproto.AdminEmitModerationEvent_Input{ CreatedBy: xrpcc.Auth.Did, @@ -118,6 +126,7 @@ func (eng *Engine) persistAccountModActions(c *AccountContext) error { if newTakedown { c.Logger.Warn("account-takedown") + actionNewTakedownCount.WithLabelValues("account").Inc() comment := "[automod]: auto account-takedown" _, err := comatproto.AdminEmitModerationEvent(ctx, xrpcc, &comatproto.AdminEmitModerationEvent_Input{ CreatedBy: xrpcc.Auth.Did, @@ -212,6 +221,10 @@ func (eng *Engine) persistRecordModActions(c *RecordContext) error { // flags don't require admin auth if len(newFlags) > 0 { + for _, val := range newFlags { + // note: WithLabelValues is a prometheus label, not an atproto label + actionNewFlagCount.WithLabelValues("record", val).Inc() + } eng.Flags.Add(ctx, atURI, newFlags) } @@ -238,6 +251,10 @@ func (eng *Engine) persistRecordModActions(c *RecordContext) error { xrpcc := eng.AdminClient if len(newLabels) > 0 { c.Logger.Info("labeling record", "newLabels", newLabels) + for _, val := range newLabels { + // note: WithLabelValues is a prometheus label, not an atproto label + actionNewLabelCount.WithLabelValues("record", val).Inc() + } comment := "[automod]: auto-labeling record" _, err := comatproto.AdminEmitModerationEvent(ctx, xrpcc, &comatproto.AdminEmitModerationEvent_Input{ CreatedBy: xrpcc.Auth.Did, @@ -266,6 +283,7 @@ func (eng *Engine) persistRecordModActions(c *RecordContext) error { if newTakedown { c.Logger.Warn("record-takedown") + actionNewTakedownCount.WithLabelValues("record").Inc() comment := "[automod]: automated record-takedown" _, err := comatproto.AdminEmitModerationEvent(ctx, xrpcc, &comatproto.AdminEmitModerationEvent_Input{ CreatedBy: xrpcc.Auth.Did, diff --git a/automod/engine/persisthelpers.go b/automod/engine/persisthelpers.go index acaddaee..f9aa31d8 100644 --- a/automod/engine/persisthelpers.go +++ b/automod/engine/persisthelpers.go @@ -142,6 +142,7 @@ func (eng *Engine) createReportIfFresh(ctx context.Context, xrpcc *xrpc.Client, } eng.Logger.Info("reporting account", "reasonType", mr.ReasonType, "comment", mr.Comment) + actionNewReportCount.WithLabelValues("account").Inc() comment := "[automod] " + mr.Comment _, err = comatproto.ModerationCreateReport(ctx, xrpcc, &comatproto.ModerationCreateReport_Input{ ReasonType: &mr.ReasonType, @@ -191,7 +192,8 @@ func (eng *Engine) createRecordReportIfFresh(ctx context.Context, xrpcc *xrpc.Cl return false, nil } - eng.Logger.Info("reporting account", "reasonType", mr.ReasonType, "comment", mr.Comment) + eng.Logger.Info("reporting record", "reasonType", mr.ReasonType, "comment", mr.Comment) + actionNewReportCount.WithLabelValues("record").Inc() comment := "[automod] " + mr.Comment _, err = comatproto.ModerationCreateReport(ctx, xrpcc, &comatproto.ModerationCreateReport_Input{ ReasonType: &mr.ReasonType, -- 2.51.2