diff --git a/cmd/beemo/main.go b/cmd/beemo/main.go index cb9fbcfb..91ce4773 100644 --- a/cmd/beemo/main.go +++ b/cmd/beemo/main.go @@ -4,19 +4,7 @@ package main import ( - "bytes" - "context" - "encoding/json" - "fmt" - "net/http" "os" - "strings" - "time" - - comatproto "github.com/bluesky-social/indigo/api/atproto" - toolsozone "github.com/bluesky-social/indigo/api/ozone" - "github.com/bluesky-social/indigo/util" - "github.com/bluesky-social/indigo/xrpc" _ "github.com/joho/godotenv/autoload" _ "go.uber.org/automaxprocs" @@ -97,145 +85,3 @@ func run(args []string) error { return app.Run(args) } -func pollNewReports(cctx *cli.Context) error { - // record last-seen report timestamp - since := time.Now() - // NOTE: uncomment this for testing - //since = time.Now().Add(time.Duration(-12) * time.Hour) - period := time.Duration(cctx.Int("poll-period")) * time.Second - - // create a new session - xrpcc := &xrpc.Client{ - Client: util.RobustHTTPClient(), - Host: cctx.String("pds-host"), - Auth: &xrpc.AuthInfo{Handle: cctx.String("handle")}, - } - - auth, err := comatproto.ServerCreateSession(context.TODO(), xrpcc, &comatproto.ServerCreateSession_Input{ - Identifier: xrpcc.Auth.Handle, - Password: cctx.String("password"), - }) - if err != nil { - return err - } - xrpcc.Auth.AccessJwt = auth.AccessJwt - xrpcc.Auth.RefreshJwt = auth.RefreshJwt - xrpcc.Auth.Did = auth.Did - xrpcc.Auth.Handle = auth.Handle - - adminToken := cctx.String("admin-password") - if len(adminToken) > 0 { - xrpcc.AdminToken = &adminToken - } - log.Infof("report polling bot starting up...") - // can flip this bool to false to prevent spamming slack channel on startup - if true { - err := sendSlackMsg(cctx, fmt.Sprintf("restarted bot, monitoring for reports since `%s`...", since.Format(time.RFC3339))) - if err != nil { - return err - } - } - for { - // refresh session - xrpcc.Auth.AccessJwt = xrpcc.Auth.RefreshJwt - refresh, err := comatproto.ServerRefreshSession(context.TODO(), xrpcc) - if err != nil { - return err - } - xrpcc.Auth.AccessJwt = refresh.AccessJwt - xrpcc.Auth.RefreshJwt = refresh.RefreshJwt - - // query just new reports (regardless of resolution state) - // ModerationQueryEvents(ctx context.Context, c *xrpc.Client, createdBy string, cursor string, includeAllUserRecords bool, limit int64, sortDirection string, subject string, types []string) (*ModerationQueryEvents_Output, error) - var limit int64 = 50 - me, err := toolsozone.ModerationQueryEvents( - cctx.Context, - xrpcc, - nil, - nil, - "", - "", - "", - "", - "", - false, - true, - limit, - nil, - nil, - nil, - "", - "", - []string{"tools.ozone.moderation.defs#modEventReport"}, - ) - if err != nil { - return err - } - // this works out to iterate from newest to oldest, which is the behavior we want (report only newest, then break) - for _, evt := range me.Events { - report := evt.Event.ModerationDefs_ModEventReport - // TODO: filter out based on subject state? similar to old "report.ResolvedByActionIds" - createdAt, err := time.Parse(time.RFC3339, evt.CreatedAt) - if err != nil { - return fmt.Errorf("invalid time format for 'createdAt': %w", err) - } - if createdAt.After(since) { - shortType := "" - if report.ReportType != nil && strings.Contains(*report.ReportType, "#") { - shortType = strings.SplitN(*report.ReportType, "#", 2)[1] - } - // ok, we found a "new" report, need to notify - msg := fmt.Sprintf("⚠️ New report at `%s` ⚠️\n", evt.CreatedAt) - msg += fmt.Sprintf("report id: `%d`\t", evt.Id) - msg += fmt.Sprintf("instance: `%s`\n", cctx.String("pds-host")) - msg += fmt.Sprintf("reasonType: `%s`\t", shortType) - msg += fmt.Sprintf("Admin: %s/reports/%d\n", cctx.String("admin-host"), evt.Id) - //msg += fmt.Sprintf("reportedByDid: `%s`\n", report.ReportedByDid) - log.Infof("found new report, notifying slack: %s", report) - err := sendSlackMsg(cctx, msg) - if err != nil { - return fmt.Errorf("failed to send slack message: %w", err) - } - since = createdAt - break - } else { - log.Debugf("skipping report: %s", report) - } - } - log.Infof("... sleeping for %s", period) - time.Sleep(period) - } -} - -type SlackWebhookBody struct { - Text string `json:"text"` -} - -// sends a simple slack message to a channel via "incoming webhook" -// The slack incoming webhook must be already configured in the slack workplace. -func sendSlackMsg(cctx *cli.Context, msg string) error { - // loosely based on: https://golangcode.com/send-slack-messages-without-a-library/ - - webhookUrl := cctx.String("slack-webhook-url") - body, _ := json.Marshal(SlackWebhookBody{Text: msg}) - req, err := http.NewRequest(http.MethodPost, webhookUrl, bytes.NewBuffer(body)) - if err != nil { - return err - } - req.Header.Add("Content-Type", "application/json") - client := util.RobustHTTPClient() - resp, err := client.Do(req) - if err != nil { - return err - } - - defer resp.Body.Close() - - buf := new(bytes.Buffer) - buf.ReadFrom(resp.Body) - if resp.StatusCode != 200 || buf.String() != "ok" { - // TODO: in some cases print body? eg, if short and text - return fmt.Errorf("failed slack webhook POST request. status=%d", resp.StatusCode) - } - return nil -} diff --git a/cmd/beemo/notify_reports.go b/cmd/beemo/notify_reports.go new file mode 100644 index 00000000..4541c6dd --- /dev/null +++ b/cmd/beemo/notify_reports.go @@ -0,0 +1,125 @@ +package main + +import ( + "context" + "fmt" + "strings" + "time" + + comatproto "github.com/bluesky-social/indigo/api/atproto" + toolsozone "github.com/bluesky-social/indigo/api/ozone" + "github.com/bluesky-social/indigo/util" + "github.com/bluesky-social/indigo/xrpc" + + "github.com/urfave/cli/v2" +) + +func pollNewReports(cctx *cli.Context) error { + // record last-seen report timestamp + since := time.Now() + // NOTE: uncomment this for testing + //since = time.Now().Add(time.Duration(-12) * time.Hour) + period := time.Duration(cctx.Int("poll-period")) * time.Second + + // create a new session + xrpcc := &xrpc.Client{ + Client: util.RobustHTTPClient(), + Host: cctx.String("pds-host"), + Auth: &xrpc.AuthInfo{Handle: cctx.String("handle")}, + } + + auth, err := comatproto.ServerCreateSession(context.TODO(), xrpcc, &comatproto.ServerCreateSession_Input{ + Identifier: xrpcc.Auth.Handle, + Password: cctx.String("password"), + }) + if err != nil { + return err + } + xrpcc.Auth.AccessJwt = auth.AccessJwt + xrpcc.Auth.RefreshJwt = auth.RefreshJwt + xrpcc.Auth.Did = auth.Did + xrpcc.Auth.Handle = auth.Handle + + adminToken := cctx.String("admin-password") + if len(adminToken) > 0 { + xrpcc.AdminToken = &adminToken + } + log.Infof("report polling bot starting up...") + // can flip this bool to false to prevent spamming slack channel on startup + if true { + err := sendSlackMsg(cctx, fmt.Sprintf("restarted bot, monitoring for reports since `%s`...", since.Format(time.RFC3339))) + if err != nil { + return err + } + } + for { + // refresh session + xrpcc.Auth.AccessJwt = xrpcc.Auth.RefreshJwt + refresh, err := comatproto.ServerRefreshSession(context.TODO(), xrpcc) + if err != nil { + return err + } + xrpcc.Auth.AccessJwt = refresh.AccessJwt + xrpcc.Auth.RefreshJwt = refresh.RefreshJwt + + // query just new reports (regardless of resolution state) + // ModerationQueryEvents(ctx context.Context, c *xrpc.Client, createdBy string, cursor string, includeAllUserRecords bool, limit int64, sortDirection string, subject string, types []string) (*ModerationQueryEvents_Output, error) + var limit int64 = 50 + me, err := toolsozone.ModerationQueryEvents( + cctx.Context, + xrpcc, + nil, + nil, + "", + "", + "", + "", + "", + false, + true, + limit, + nil, + nil, + nil, + "", + "", + []string{"tools.ozone.moderation.defs#modEventReport"}, + ) + if err != nil { + return err + } + // this works out to iterate from newest to oldest, which is the behavior we want (report only newest, then break) + for _, evt := range me.Events { + report := evt.Event.ModerationDefs_ModEventReport + // TODO: filter out based on subject state? similar to old "report.ResolvedByActionIds" + createdAt, err := time.Parse(time.RFC3339, evt.CreatedAt) + if err != nil { + return fmt.Errorf("invalid time format for 'createdAt': %w", err) + } + if createdAt.After(since) { + shortType := "" + if report.ReportType != nil && strings.Contains(*report.ReportType, "#") { + shortType = strings.SplitN(*report.ReportType, "#", 2)[1] + } + // ok, we found a "new" report, need to notify + msg := fmt.Sprintf("⚠️ New report at `%s` ⚠️\n", evt.CreatedAt) + msg += fmt.Sprintf("report id: `%d`\t", evt.Id) + msg += fmt.Sprintf("instance: `%s`\n", cctx.String("pds-host")) + msg += fmt.Sprintf("reasonType: `%s`\t", shortType) + msg += fmt.Sprintf("Admin: %s/reports/%d\n", cctx.String("admin-host"), evt.Id) + //msg += fmt.Sprintf("reportedByDid: `%s`\n", report.ReportedByDid) + log.Infof("found new report, notifying slack: %s", report) + err := sendSlackMsg(cctx, msg) + if err != nil { + return fmt.Errorf("failed to send slack message: %w", err) + } + since = createdAt + break + } else { + log.Debugf("skipping report: %s", report) + } + } + log.Infof("... sleeping for %s", period) + time.Sleep(period) + } +} diff --git a/cmd/beemo/slack.go b/cmd/beemo/slack.go new file mode 100644 index 00000000..5d0aea2e --- /dev/null +++ b/cmd/beemo/slack.go @@ -0,0 +1,45 @@ +package main + +import ( + "bytes" + "encoding/json" + "fmt" + "net/http" + + "github.com/bluesky-social/indigo/util" + + "github.com/urfave/cli/v2" +) + +type SlackWebhookBody struct { + Text string `json:"text"` +} + +// sends a simple slack message to a channel via "incoming webhook" +// The slack incoming webhook must be already configured in the slack workplace. +func sendSlackMsg(cctx *cli.Context, msg string) error { + // loosely based on: https://golangcode.com/send-slack-messages-without-a-library/ + + webhookUrl := cctx.String("slack-webhook-url") + body, _ := json.Marshal(SlackWebhookBody{Text: msg}) + req, err := http.NewRequest(http.MethodPost, webhookUrl, bytes.NewBuffer(body)) + if err != nil { + return err + } + req.Header.Add("Content-Type", "application/json") + client := util.RobustHTTPClient() + resp, err := client.Do(req) + if err != nil { + return err + } + + defer resp.Body.Close() + + buf := new(bytes.Buffer) + buf.ReadFrom(resp.Body) + if resp.StatusCode != 200 || buf.String() != "ok" { + // TODO: in some cases print body? eg, if short and text + return fmt.Errorf("failed slack webhook POST request. status=%d", resp.StatusCode) + } + return nil +} -- 2.51.2 From b8f540ecfe89b8f49ce0671bf94a521b7cc8beac Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Fri, 24 May 2024 17:42:36 -0700 Subject: [PATCH 2/6] beemo: add mention notifier --- cmd/beemo/firehose_consumer.go | 139 +++++++++++++++++++++++++++++++++ cmd/beemo/main.go | 123 ++++++++++++++++++++--------- cmd/beemo/notify_mentions.go | 86 ++++++++++++++++++++ cmd/beemo/notify_reports.go | 17 ++-- cmd/beemo/slack.go | 8 +- 5 files changed, 324 insertions(+), 49 deletions(-) create mode 100644 cmd/beemo/firehose_consumer.go create mode 100644 cmd/beemo/notify_mentions.go diff --git a/cmd/beemo/firehose_consumer.go b/cmd/beemo/firehose_consumer.go new file mode 100644 index 00000000..a554df41 --- /dev/null +++ b/cmd/beemo/firehose_consumer.go @@ -0,0 +1,139 @@ +package main + +import ( + "bytes" + "context" + "fmt" + "log/slog" + "net/http" + "net/url" + "strings" + + comatproto "github.com/bluesky-social/indigo/api/atproto" + appbsky "github.com/bluesky-social/indigo/api/bsky" + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/bluesky-social/indigo/events/schedulers/autoscaling" + lexutil "github.com/bluesky-social/indigo/lex/util" + + "github.com/bluesky-social/indigo/events" + "github.com/bluesky-social/indigo/repo" + "github.com/bluesky-social/indigo/repomgr" + "github.com/carlmjohnson/versioninfo" + "github.com/gorilla/websocket" +) + +func RunFirehoseConsumer(ctx context.Context, logger *slog.Logger, relayHost string, postCallback func(context.Context, syntax.DID, syntax.RecordKey, appbsky.FeedPost) error) error { + + dialer := websocket.DefaultDialer + u, err := url.Parse(relayHost) + if err != nil { + return fmt.Errorf("invalid relayHost URI: %w", err) + } + // always continue at the current cursor offset (don't provide cursor query param) + u.Path = "xrpc/com.atproto.sync.subscribeRepos" + logger.Info("subscribing to repo event stream", "upstream", relayHost) + con, _, err := dialer.Dial(u.String(), http.Header{ + "User-Agent": []string{fmt.Sprintf("beemo/%s", versioninfo.Short())}, + }) + if err != nil { + return fmt.Errorf("subscribing to firehose failed (dialing): %w", err) + } + + rsc := &events.RepoStreamCallbacks{ + RepoCommit: func(evt *comatproto.SyncSubscribeRepos_Commit) error { + return HandleRepoCommit(ctx, logger, evt, postCallback) + }, + // NOTE: could add other callbacks as needed + } + + var scheduler events.Scheduler + // use auto-scaling scheduler + scaleSettings := autoscaling.DefaultAutoscaleSettings() + scheduler = autoscaling.NewScheduler(scaleSettings, relayHost, rsc.EventHandler) + logger.Info("beemo firehose scheduler configured", "scheduler", "autoscaling", "initial", scaleSettings.Concurrency, "max", scaleSettings.MaxConcurrency) + + return events.HandleRepoStream(ctx, con, scheduler) +} + +// TODO: move this to a "ParsePath" helper in syntax package? +func splitRepoPath(path string) (syntax.NSID, syntax.RecordKey, error) { + parts := strings.SplitN(path, "/", 3) + if len(parts) != 2 { + return "", "", fmt.Errorf("invalid record path: %s", path) + } + collection, err := syntax.ParseNSID(parts[0]) + if err != nil { + return "", "", err + } + rkey, err := syntax.ParseRecordKey(parts[1]) + if err != nil { + return "", "", err + } + return collection, rkey, nil +} + +// NOTE: for now, this function basically never errors, just logs and returns nil. Should think through error processing better. +func HandleRepoCommit(ctx context.Context, logger *slog.Logger, evt *comatproto.SyncSubscribeRepos_Commit, postCallback func(context.Context, syntax.DID, syntax.RecordKey, appbsky.FeedPost) error) error { + + logger = logger.With("event", "commit", "did", evt.Repo, "rev", evt.Rev, "seq", evt.Seq) + logger.Debug("received commit event") + + if evt.TooBig { + logger.Warn("skipping tooBig events for now") + return nil + } + + did, err := syntax.ParseDID(evt.Repo) + if err != nil { + logger.Error("bad DID syntax in event", "err", err) + return nil + } + + rr, err := repo.ReadRepoFromCar(ctx, bytes.NewReader(evt.Blocks)) + if err != nil { + logger.Error("failed to read repo from car", "err", err) + return nil + } + + for _, op := range evt.Ops { + logger = logger.With("eventKind", op.Action, "path", op.Path) + collection, rkey, err := splitRepoPath(op.Path) + if err != nil { + logger.Error("invalid path in repo op") + return nil + } + + ek := repomgr.EventKind(op.Action) + switch ek { + case repomgr.EvtKindCreateRecord, repomgr.EvtKindUpdateRecord: + // read the record bytes from blocks, and verify CID + rc, recordCBOR, err := rr.GetRecordBytes(ctx, op.Path) + if err != nil { + logger.Error("reading record from event blocks (CAR)", "err", err) + continue + } + if op.Cid == nil || lexutil.LexLink(rc) != *op.Cid { + logger.Error("mismatch between commit op CID and record block", "recordCID", rc, "opCID", op.Cid) + continue + } + + switch collection { + case "app.bsky.feed.post": + var post appbsky.FeedPost + if err := post.UnmarshalCBOR(bytes.NewReader(*recordCBOR)); err != nil { + logger.Error("failed to parse app.bsky.feed.post record", "err", err) + continue + } + if err := postCallback(ctx, did, rkey, post); err != nil { + logger.Error("failed to process post record", "err", err) + continue + } + } + + default: + // ignore other events + } + } + + return nil +} diff --git a/cmd/beemo/main.go b/cmd/beemo/main.go index 91ce4773..544ee3d8 100644 --- a/cmd/beemo/main.go +++ b/cmd/beemo/main.go @@ -4,21 +4,22 @@ package main import ( + "io" + "log/slog" "os" + "strings" _ "github.com/joho/godotenv/autoload" _ "go.uber.org/automaxprocs" "github.com/carlmjohnson/versioninfo" - logging "github.com/ipfs/go-log" "github.com/urfave/cli/v2" ) -var log = logging.Logger("beemo") - func main() { if err := run(os.Args); err != nil { - log.Fatal(err) + slog.Error("exiting", "err", err) + os.Exit(-1) } } @@ -32,34 +33,9 @@ func run(args []string) error { app.Flags = []cli.Flag{ &cli.StringFlag{ - Name: "pds-host", - Usage: "method, hostname, and port of PDS instance", - Value: "http://localhost:4849", - EnvVars: []string{"ATP_PDS_HOST"}, - }, - &cli.StringFlag{ - Name: "admin-host", - Usage: "method, hostname, and port of admin interface (eg, Ozone), for direct links", - Value: "http://localhost:3000", - EnvVars: []string{"ATP_ADMIN_HOST"}, - }, - &cli.StringFlag{ - Name: "handle", - Usage: "for PDS login", - Required: true, - EnvVars: []string{"ATP_AUTH_HANDLE"}, - }, - &cli.StringFlag{ - Name: "password", - Usage: "for PDS login", - Required: true, - EnvVars: []string{"ATP_AUTH_PASSWORD"}, - }, - &cli.StringFlag{ - Name: "admin-password", - Usage: "admin authentication password for PDS", - Required: true, - EnvVars: []string{"ATP_AUTH_ADMIN_PASSWORD"}, + Name: "log-level", + Usage: "log verbosity level (eg: warn, info, debug)", + EnvVars: []string{"BEEMO_LOG_LEVEL", "GO_LOG_LEVEL", "LOG_LEVEL"}, }, &cli.StringFlag{ Name: "slack-webhook-url", @@ -68,20 +44,91 @@ func run(args []string) error { Required: true, EnvVars: []string{"SLACK_WEBHOOK_URL"}, }, - &cli.IntFlag{ - Name: "poll-period", - Usage: "API poll period in seconds", - Value: 30, - EnvVars: []string{"POLL_PERIOD"}, - }, } app.Commands = []*cli.Command{ &cli.Command{ Name: "notify-reports", Usage: "watch for new moderation reports, notify in slack", Action: pollNewReports, + Flags: []cli.Flag{ + &cli.StringFlag{ + Name: "pds-host", + Usage: "method, hostname, and port of PDS instance", + Value: "http://localhost:4849", + EnvVars: []string{"ATP_PDS_HOST"}, + }, + &cli.StringFlag{ + Name: "admin-host", + Usage: "method, hostname, and port of admin interface (eg, Ozone), for direct links", + Value: "http://localhost:3000", + EnvVars: []string{"ATP_ADMIN_HOST"}, + }, + &cli.IntFlag{ + Name: "poll-period", + Usage: "API poll period in seconds", + Value: 30, + EnvVars: []string{"POLL_PERIOD"}, + }, + &cli.StringFlag{ + Name: "handle", + Usage: "for PDS login", + Required: true, + EnvVars: []string{"ATP_AUTH_HANDLE"}, + }, + &cli.StringFlag{ + Name: "password", + Usage: "for PDS login", + Required: true, + EnvVars: []string{"ATP_AUTH_PASSWORD"}, + }, + &cli.StringFlag{ + Name: "admin-password", + Usage: "admin authentication password for PDS", + Required: true, + EnvVars: []string{"ATP_AUTH_ADMIN_PASSWORD"}, + }, + }, + }, + &cli.Command{ + Name: "notify-mentions", + Usage: "watch firehose for posts mentioning specific accounts", + Action: notifyMentions, + Flags: []cli.Flag{ + &cli.StringFlag{ + Name: "relay-host", + Usage: "method, hostname, and port of Relay instance (websocket)", + Value: "wss://bsky.network", + EnvVars: []string{"ATP_RELAY_HOST"}, + }, + &cli.StringFlag{ + Name: "mention-dids", + Usage: "DIDs to look for in mentions (comma-separated)", + Required: true, + EnvVars: []string{"BEEMO_MENTION_DIDS"}, + }, + }, }, } return app.Run(args) } +func configLogger(cctx *cli.Context, writer io.Writer) *slog.Logger { + var level slog.Level + switch strings.ToLower(cctx.String("log-level")) { + case "error": + level = slog.LevelError + case "warn": + level = slog.LevelWarn + case "info": + level = slog.LevelInfo + case "debug": + level = slog.LevelDebug + default: + level = slog.LevelInfo + } + logger := slog.New(slog.NewJSONHandler(writer, &slog.HandlerOptions{ + Level: level, + })) + slog.SetDefault(logger) + return logger +} diff --git a/cmd/beemo/notify_mentions.go b/cmd/beemo/notify_mentions.go new file mode 100644 index 00000000..eaad36f8 --- /dev/null +++ b/cmd/beemo/notify_mentions.go @@ -0,0 +1,86 @@ +package main + +import ( + "context" + "fmt" + "log/slog" + "os" + "strings" + + appbsky "github.com/bluesky-social/indigo/api/bsky" + "github.com/bluesky-social/indigo/atproto/identity" + "github.com/bluesky-social/indigo/atproto/syntax" + + "github.com/urfave/cli/v2" +) + +type MentionChecker struct { + slackWebhookURL string + mentionDIDs []syntax.DID + logger *slog.Logger + directory identity.Directory +} + +func (mc *MentionChecker) ProcessPost(ctx context.Context, did syntax.DID, rkey syntax.RecordKey, post appbsky.FeedPost) error { + mc.logger.Debug("processing post record", "did", did, "rkey", rkey) + + for _, facet := range post.Facets { + for _, feature := range facet.Features { + mention := feature.RichtextFacet_Mention + if mention == nil { + continue + } + for _, d := range mc.mentionDIDs { + if mention.Did == d.String() { + mc.logger.Info("found mention", "target", d, "author", did, "rkey", rkey) + targetIdent, err := mc.directory.LookupDID(ctx, syntax.DID(mention.Did)) + if err != nil { + return err + } + authorIdent, err := mc.directory.LookupDID(ctx, did) + if err != nil { + return err + } + msg := fmt.Sprintf("Mention of `@%s` by `@%s` ():\n```%s```", targetIdent.Handle, authorIdent.Handle, did, rkey, post.Text) + return sendSlackMsg(ctx, msg, mc.slackWebhookURL) + } + } + } + } + return nil +} + +func notifyMentions(cctx *cli.Context) error { + ctx := context.Background() + logger := configLogger(cctx, os.Stdout) + relayHost := cctx.String("relay-host") + + mentionDIDs := []syntax.DID{} + for _, raw := range strings.Split(cctx.String("mention-dids"), ",") { + fmt.Println(raw) + did, err := syntax.ParseDID(raw) + if err != nil { + return err + } + mentionDIDs = append(mentionDIDs, did) + } + + checker := MentionChecker{ + slackWebhookURL: cctx.String("slack-webhook-url"), + mentionDIDs: mentionDIDs, + logger: logger, + directory: identity.DefaultDirectory(), + } + + logger.Info("beemo mention checker starting up...", "relayHost", relayHost, "mentionDIDs", mentionDIDs) + + // can flip this bool to false to prevent spamming slack channel on startup + if true { + err := sendSlackMsg(ctx, fmt.Sprintf("beemo booting, looking for account mentions: `%s`", mentionDIDs), checker.slackWebhookURL) + if err != nil { + return err + } + } + + return RunFirehoseConsumer(ctx, logger, relayHost, checker.ProcessPost) +} diff --git a/cmd/beemo/notify_reports.go b/cmd/beemo/notify_reports.go index 4541c6dd..58d52c7a 100644 --- a/cmd/beemo/notify_reports.go +++ b/cmd/beemo/notify_reports.go @@ -3,6 +3,7 @@ package main import ( "context" "fmt" + "os" "strings" "time" @@ -15,6 +16,10 @@ import ( ) func pollNewReports(cctx *cli.Context) error { + ctx := context.Background() + logger := configLogger(cctx, os.Stdout) + slackWebhookURL := cctx.String("slack-webhook-url") + // record last-seen report timestamp since := time.Now() // NOTE: uncomment this for testing @@ -44,10 +49,10 @@ func pollNewReports(cctx *cli.Context) error { if len(adminToken) > 0 { xrpcc.AdminToken = &adminToken } - log.Infof("report polling bot starting up...") + logger.Info("report polling bot starting up...") // can flip this bool to false to prevent spamming slack channel on startup if true { - err := sendSlackMsg(cctx, fmt.Sprintf("restarted bot, monitoring for reports since `%s`...", since.Format(time.RFC3339))) + err := sendSlackMsg(ctx, fmt.Sprintf("restarted bot, monitoring for reports since `%s`...", since.Format(time.RFC3339)), slackWebhookURL) if err != nil { return err } @@ -108,18 +113,18 @@ func pollNewReports(cctx *cli.Context) error { msg += fmt.Sprintf("reasonType: `%s`\t", shortType) msg += fmt.Sprintf("Admin: %s/reports/%d\n", cctx.String("admin-host"), evt.Id) //msg += fmt.Sprintf("reportedByDid: `%s`\n", report.ReportedByDid) - log.Infof("found new report, notifying slack: %s", report) - err := sendSlackMsg(cctx, msg) + logger.Info("found new report, notifying slack", "report", report) + err := sendSlackMsg(ctx, msg, slackWebhookURL) if err != nil { return fmt.Errorf("failed to send slack message: %w", err) } since = createdAt break } else { - log.Debugf("skipping report: %s", report) + logger.Debug("skipping report", "report", report) } } - log.Infof("... sleeping for %s", period) + logger.Info("... sleeping", "period", period) time.Sleep(period) } } diff --git a/cmd/beemo/slack.go b/cmd/beemo/slack.go index 5d0aea2e..6fca5af2 100644 --- a/cmd/beemo/slack.go +++ b/cmd/beemo/slack.go @@ -2,13 +2,12 @@ package main import ( "bytes" + "context" "encoding/json" "fmt" "net/http" "github.com/bluesky-social/indigo/util" - - "github.com/urfave/cli/v2" ) type SlackWebhookBody struct { @@ -17,12 +16,11 @@ type SlackWebhookBody struct { // sends a simple slack message to a channel via "incoming webhook" // The slack incoming webhook must be already configured in the slack workplace. -func sendSlackMsg(cctx *cli.Context, msg string) error { +func sendSlackMsg(ctx context.Context, msg, webhookURL string) error { // loosely based on: https://golangcode.com/send-slack-messages-without-a-library/ - webhookUrl := cctx.String("slack-webhook-url") body, _ := json.Marshal(SlackWebhookBody{Text: msg}) - req, err := http.NewRequest(http.MethodPost, webhookUrl, bytes.NewBuffer(body)) + req, err := http.NewRequestWithContext(ctx, http.MethodPost, webhookURL, bytes.NewBuffer(body)) if err != nil { return err } -- 2.51.2 From baba29cb4b457b477cf24ada12b6b5f495e204b5 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Fri, 24 May 2024 17:58:48 -0700 Subject: [PATCH 3/6] add a note about any quote/embed/media --- cmd/beemo/notify_mentions.go | 3 +++ 1 file changed, 3 insertions(+) diff --git a/cmd/beemo/notify_mentions.go b/cmd/beemo/notify_mentions.go index eaad36f8..b0b0fde8 100644 --- a/cmd/beemo/notify_mentions.go +++ b/cmd/beemo/notify_mentions.go @@ -42,6 +42,9 @@ func (mc *MentionChecker) ProcessPost(ctx context.Context, did syntax.DID, rkey return err } msg := fmt.Sprintf("Mention of `@%s` by `@%s` ():\n```%s```", targetIdent.Handle, authorIdent.Handle, did, rkey, post.Text) + if post.Embed.EmbedImages != nil || post.Embed.EmbedRecordWithMedia != nil || post.Embed.EmbedRecord != nil || post.Embed.EmbedExternal != nil { + msg += "\n(post also contains an embed/quote/media)" + } return sendSlackMsg(ctx, msg, mc.slackWebhookURL) } } -- 2.51.2 From ebf06e71a3e28a800bf0a4cd5f2a30c203ae45a2 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Fri, 24 May 2024 18:56:56 -0700 Subject: [PATCH 4/6] beemo: fix silly nil embed panic --- cmd/beemo/notify_mentions.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/cmd/beemo/notify_mentions.go b/cmd/beemo/notify_mentions.go index b0b0fde8..a4f09c33 100644 --- a/cmd/beemo/notify_mentions.go +++ b/cmd/beemo/notify_mentions.go @@ -42,7 +42,7 @@ func (mc *MentionChecker) ProcessPost(ctx context.Context, did syntax.DID, rkey return err } msg := fmt.Sprintf("Mention of `@%s` by `@%s` ():\n```%s```", targetIdent.Handle, authorIdent.Handle, did, rkey, post.Text) - if post.Embed.EmbedImages != nil || post.Embed.EmbedRecordWithMedia != nil || post.Embed.EmbedRecord != nil || post.Embed.EmbedExternal != nil { + if post.Embed != nil && (post.Embed.EmbedImages != nil || post.Embed.EmbedRecordWithMedia != nil || post.Embed.EmbedRecord != nil || post.Embed.EmbedExternal != nil) { msg += "\n(post also contains an embed/quote/media)" } return sendSlackMsg(ctx, msg, mc.slackWebhookURL) -- 2.51.2 From 1c226097194616a1d31c56a83d480e51b0eed393 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Fri, 24 May 2024 19:01:54 -0700 Subject: [PATCH 5/6] small cleanups --- cmd/beemo/notify_mentions.go | 1 - cmd/beemo/notify_reports.go | 4 ++-- cmd/beemo/slack.go | 1 - 3 files changed, 2 insertions(+), 4 deletions(-) diff --git a/cmd/beemo/notify_mentions.go b/cmd/beemo/notify_mentions.go index a4f09c33..96026ecf 100644 --- a/cmd/beemo/notify_mentions.go +++ b/cmd/beemo/notify_mentions.go @@ -60,7 +60,6 @@ func notifyMentions(cctx *cli.Context) error { mentionDIDs := []syntax.DID{} for _, raw := range strings.Split(cctx.String("mention-dids"), ",") { - fmt.Println(raw) did, err := syntax.ParseDID(raw) if err != nil { return err diff --git a/cmd/beemo/notify_reports.go b/cmd/beemo/notify_reports.go index 58d52c7a..5404619b 100644 --- a/cmd/beemo/notify_reports.go +++ b/cmd/beemo/notify_reports.go @@ -33,7 +33,7 @@ func pollNewReports(cctx *cli.Context) error { Auth: &xrpc.AuthInfo{Handle: cctx.String("handle")}, } - auth, err := comatproto.ServerCreateSession(context.TODO(), xrpcc, &comatproto.ServerCreateSession_Input{ + auth, err := comatproto.ServerCreateSession(ctx, xrpcc, &comatproto.ServerCreateSession_Input{ Identifier: xrpcc.Auth.Handle, Password: cctx.String("password"), }) @@ -60,7 +60,7 @@ func pollNewReports(cctx *cli.Context) error { for { // refresh session xrpcc.Auth.AccessJwt = xrpcc.Auth.RefreshJwt - refresh, err := comatproto.ServerRefreshSession(context.TODO(), xrpcc) + refresh, err := comatproto.ServerRefreshSession(ctx, xrpcc) if err != nil { return err } diff --git a/cmd/beemo/slack.go b/cmd/beemo/slack.go index 6fca5af2..58887e72 100644 --- a/cmd/beemo/slack.go +++ b/cmd/beemo/slack.go @@ -36,7 +36,6 @@ func sendSlackMsg(ctx context.Context, msg, webhookURL string) error { buf := new(bytes.Buffer) buf.ReadFrom(resp.Body) if resp.StatusCode != 200 || buf.String() != "ok" { - // TODO: in some cases print body? eg, if short and text return fmt.Errorf("failed slack webhook POST request. status=%d", resp.StatusCode) } return nil -- 2.51.2 From c13b161f5541be7a7365d7e06c32953e23a0094b Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Thu, 27 Jun 2024 12:11:09 -0700 Subject: [PATCH 6/6] tweaks from review (thanks jaz) --- cmd/beemo/firehose_consumer.go | 15 ++++++++++----- 1 file changed, 10 insertions(+), 5 deletions(-) diff --git a/cmd/beemo/firehose_consumer.go b/cmd/beemo/firehose_consumer.go index a554df41..55f2edc8 100644 --- a/cmd/beemo/firehose_consumer.go +++ b/cmd/beemo/firehose_consumer.go @@ -12,7 +12,7 @@ import ( comatproto "github.com/bluesky-social/indigo/api/atproto" appbsky "github.com/bluesky-social/indigo/api/bsky" "github.com/bluesky-social/indigo/atproto/syntax" - "github.com/bluesky-social/indigo/events/schedulers/autoscaling" + "github.com/bluesky-social/indigo/events/schedulers/parallel" lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/events" @@ -47,10 +47,15 @@ func RunFirehoseConsumer(ctx context.Context, logger *slog.Logger, relayHost str } var scheduler events.Scheduler - // use auto-scaling scheduler - scaleSettings := autoscaling.DefaultAutoscaleSettings() - scheduler = autoscaling.NewScheduler(scaleSettings, relayHost, rsc.EventHandler) - logger.Info("beemo firehose scheduler configured", "scheduler", "autoscaling", "initial", scaleSettings.Concurrency, "max", scaleSettings.MaxConcurrency) + // use parallel scheduler + parallelism := 4 + scheduler = parallel.NewScheduler( + parallelism, + 1000, + relayHost, + rsc.EventHandler, + ) + logger.Info("beemo firehose scheduler configured", "scheduler", "parallel", "workers", parallelism) return events.HandleRepoStream(ctx, con, scheduler) }