diff --git a/cmd/recordcollector/main.go b/cmd/recordcollector/main.go index 253125c..0601fee 100644 --- a/cmd/recordcollector/main.go +++ b/cmd/recordcollector/main.go @@ -23,6 +23,7 @@ import ( "github.com/edavis/recordcollector/pkg" _ "github.com/joho/godotenv/autoload" _ "github.com/mattn/go-sqlite3" + "github.com/redis/go-redis/v9" ) var RecordCollectorDid = os.Getenv("DID") @@ -50,6 +51,7 @@ type handler struct { logger *slog.Logger client *xrpc.Client // client that makes Ozone requests db *sql.DB // sqlite3 DB for tracking labels (to avoid dupes) + redis *redis.Client lk sync.Mutex count int } @@ -95,10 +97,24 @@ func main() { log.Fatalf("could not init database: %v", err) } + redisClient := redis.NewClient(&redis.Options{ + Addr: os.Getenv("REDIS_ADDR"), + }) + if err := redisClient.Ping(ctx).Err(); err != nil { + log.Fatalf("could not connect to redis: %v", err) + } + defer func() { + if err := redisClient.Close(); err != nil { + logger.Error("failed to close redis", "err", err) + } + logger.Info("redis closed") + }() + h := &handler{ logger: logger, client: &ozoneClient, db: dbCnx, + redis: redisClient, } h.runRefreshSession(ctx) @@ -175,6 +191,19 @@ func (h *handler) runRefreshSession(ctx context.Context) { } func (h *handler) addLabel(ctx context.Context, did, label string) error { + key := fmt.Sprintf("recordcollector_%s_%s", did, label) + + // Check if already labeled recently + exists, err := h.redis.Exists(ctx, key).Result() + if err != nil { + h.logger.Error("redis error checking key", "err", err, "key", key) + // Continue anyway - better to potentially duplicate than miss labels + } + if exists > 0 { + h.logger.Info("skipping label, already emitted recently", "did", did, "label", label) + return nil + } + var labelDuration int64 = 24 * 30 input := &toolsozone.ModerationEmitEvent_Input{ CreatedBy: RecordCollectorDid, @@ -199,6 +228,12 @@ func (h *handler) addLabel(ctx context.Context, did, label string) error { } h.logger.Info("added label", "view", view) + // Set key with 24h TTL after successful emission + if err := h.redis.Set(ctx, key, "1", 24*time.Hour).Err(); err != nil { + h.logger.Error("failed to set redis key", "err", err, "key", key) + // Don't return error - label was emitted successfully + } + return nil } diff --git a/compose.yaml b/compose.yaml index 49497d8..7431467 100644 --- a/compose.yaml +++ b/compose.yaml @@ -10,3 +10,14 @@ services: - type: bind source: ./labels.db target: /labels.db + depends_on: + - redis + + redis: + image: redis:7-alpine + container_name: recordcollector-redis + volumes: + - redis-data:/data + +volumes: + redis-data: diff --git a/go.mod b/go.mod index e100877..0f7b6ee 100644 --- a/go.mod +++ b/go.mod @@ -9,6 +9,7 @@ require ( github.com/bluesky-social/jetstream v0.0.0-20241210005130-ea96859b93d1 github.com/joho/godotenv v1.5.1 github.com/mattn/go-sqlite3 v1.14.22 + github.com/redis/go-redis/v9 v9.17.3 golang.org/x/text v0.33.0 ) @@ -16,6 +17,7 @@ require ( github.com/beorn7/perks v1.0.1 // indirect github.com/carlmjohnson/versioninfo v0.22.5 // indirect github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f // indirect github.com/felixge/httpsnoop v1.0.4 // indirect github.com/go-logr/logr v1.4.1 // indirect github.com/go-logr/stdr v1.2.2 // indirect diff --git a/go.sum b/go.sum index 4565d6d..c953cfb 100644 --- a/go.sum +++ b/go.sum @@ -6,6 +6,10 @@ github.com/bluesky-social/indigo v0.0.0-20250320052052-4873aceeabf4 h1:kKsrsvvMG github.com/bluesky-social/indigo v0.0.0-20250320052052-4873aceeabf4/go.mod h1:NVBwZvbBSa93kfyweAmKwOLYawdVHdwZ9s+GZtBBVLA= github.com/bluesky-social/jetstream v0.0.0-20241210005130-ea96859b93d1 h1:CFvRtYNSnWRAi/98M3O466t9dYuwtesNbu6FVPymRrA= github.com/bluesky-social/jetstream v0.0.0-20241210005130-ea96859b93d1/go.mod h1:WiYEeyJSdUwqoaZ71KJSpTblemUCpwJfh5oVXplK6T4= +github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs= +github.com/bsm/ginkgo/v2 v2.12.0/go.mod h1:SwYbGRRDovPVboqFv0tPTcG1sN61LM1Z4ARdbAV9g4c= +github.com/bsm/gomega v1.27.10 h1:yeMWxP2pV2fG3FgAODIY8EiRE3dy0aeFYt4l7wh6yKA= +github.com/bsm/gomega v1.27.10/go.mod h1:JyEr/xRbxbtgWNi8tIEVPUYZ5Dzef52k01W3YH0H+O0= github.com/carlmjohnson/versioninfo v0.22.5 h1:O00sjOLUAFxYQjlN/bzYTuZiS0y6fWDQjMRvwtKgwwc= github.com/carlmjohnson/versioninfo v0.22.5/go.mod h1:QT9mph3wcVfISUKd0i9sZfVrPviHuSF+cUtLjm2WSf8= github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= @@ -14,6 +18,8 @@ github.com/cpuguy83/go-md2man/v2 v2.0.0-20190314233015-f79a8a8ca69d/go.mod h1:ma github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f h1:lO4WD4F/rVNCu3HqELle0jiPLLBs70cWOduZpkS1E78= +github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f/go.mod h1:cuUVRXasLTGF7a8hSLbxyZXjz+1KgoB3wDUb6vlszIc= github.com/felixge/httpsnoop v1.0.4 h1:NFTV2Zj1bL4mc9sqWACXbQFVBBg2W3GPvqp8/ESS2Wg= github.com/felixge/httpsnoop v1.0.4/go.mod h1:m8KPJKqk1gH5J9DgRY2ASl2lWCfGKXixSwevea8zH2U= github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A= @@ -124,6 +130,8 @@ github.com/prometheus/common v0.54.0 h1:ZlZy0BgJhTwVZUn7dLOkwCZHUkrAqd3WYtcFCWnM github.com/prometheus/common v0.54.0/go.mod h1:/TQgMJP5CuVYveyT7n/0Ix8yLNNXy9yRSkhnLTHPDIQ= github.com/prometheus/procfs v0.15.1 h1:YagwOFzUgYfKKHX6Dr+sHT7km/hxC76UB0learggepc= github.com/prometheus/procfs v0.15.1/go.mod h1:fB45yRUv8NstnjriLhBQLuOUt+WW4BsoGhij/e3PBqk= +github.com/redis/go-redis/v9 v9.17.3 h1:fN29NdNrE17KttK5Ndf20buqfDZwGNgoUr9qjl1DQx4= +github.com/redis/go-redis/v9 v9.17.3/go.mod h1:u410H11HMLoB+TP67dz8rL9s6QW2j76l0//kSOd3370= github.com/rogpeppe/go-internal v1.3.0/go.mod h1:M8bDsm7K2OlrFYOpmOWEs/qY81heoFRclV5y23lUDJ4= github.com/rogpeppe/go-internal v1.12.0 h1:exVL4IDcn6na9z1rAb56Vxr+CgyK3nn3O+epU5NdKM8= github.com/rogpeppe/go-internal v1.12.0/go.mod h1:E+RYuTGaKKdloAfM02xzb0FW3Paa99yedzYV+kq4uf4=