From 041ccd9bad1c94551aa2a328bc300c6c3898fc37 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Tue, 29 Jul 2025 17:57:15 -0700 Subject: [PATCH] spxrpc: add report creation with clipping --- Makefile | 3 + go.mod | 6 +- go.sum | 8 +- .../content/docs/lex-reference/openapi.json | 118 ++++++++++++++++++ .../com/atproto/moderation/createReport.json | 90 +++++++++++++ pkg/api/api_internal.go | 2 +- pkg/media/clip_user.go | 8 +- pkg/spxrpc/com_atproto_moderation.go | 108 ++++++++++++++++ pkg/spxrpc/spxrpc.go | 21 ++++ pkg/spxrpc/stubs.go | 19 +++ 10 files changed, 371 insertions(+), 12 deletions(-) create mode 100644 lexicons/com/atproto/moderation/createReport.json create mode 100644 pkg/spxrpc/com_atproto_moderation.go diff --git a/Makefile b/Makefile index 51110720..e5f9a1dd 100644 --- a/Makefile +++ b/Makefile @@ -130,6 +130,9 @@ js-lexicons: && sed -i.bak 's/PlaceStreamChatProfile\.Main/PlaceStreamChatProfile\.Record/' $$(find ./js/streamplace/src/lexicons/types/place/stream -type f) \ && sed -i.bak "s/import\ \*\ as\ AppBskyFeedDefs\ from\ '.\/defs'/import \{ AppBskyFeedDefs } from '@atproto\/api'/" $$(find ./js/streamplace/src/lexicons/types -type f) \ && sed -i.bak "s/import\ \*\ as\ AppBskyActorDefs\ from\ '.\/defs'/import \{ AppBskyActorDefs } from '@atproto\/api'/" $$(find ./js/streamplace/src/lexicons -type f) \ + && sed -i.bak "s/import\ \*\ as\ ComAtprotoAdminDefs\ from\ .*$$/import \{ ComAtprotoAdminDefs } from '@atproto\/api'/" $$(find ./js/streamplace/src/lexicons -type f) \ + && sed -i.bak "s/import\ \*\ as\ ComAtprotoRepoStrongRef\ from\ .*$$/import \{ ComAtprotoRepoStrongRef } from '@atproto\/api'/" $$(find ./js/streamplace/src/lexicons -type f) \ + && sed -i.bak "s/import\ \*\ as\ ComAtprotoModerationDefs\ from\ .*$$/import \{ ComAtprotoModerationDefs } from '@atproto\/api'/" $$(find ./js/streamplace/src/lexicons -type f) \ && npx prettier --write $$(find ./js/streamplace/src/lexicons -type f -name '*.ts') \ && find . | grep bak$$ | xargs rm diff --git a/go.mod b/go.mod index 8d80ee36..aaf13b9a 100644 --- a/go.mod +++ b/go.mod @@ -10,8 +10,6 @@ replace github.com/gocql/gocql => github.com/scylladb/gocql v1.14.4 replace github.com/AxisCommunications/go-dpop => github.com/streamplace/go-dpop v0.0.0-20250510031900-c897158a8ad4 -replace github.com/bluesky-social/indigo => github.com/streamplace/indigo v0.0.0-20250729000314-d9cf6fb86611 - require ( firebase.google.com/go/v4 v4.14.1 git.stream.place/streamplace/c2pa-go v0.7.0 @@ -19,7 +17,7 @@ require ( github.com/NYTimes/gziphandler v1.1.1 github.com/ThalesGroup/crypto11 v0.0.0-00010101000000-000000000000 github.com/acarl005/stripansi v0.0.0-20180116102854-5a71ef0e047d - github.com/bluesky-social/indigo v0.0.0-20250617211950-336ebe49427b + github.com/bluesky-social/indigo v0.0.0-20250729223159-573ae927246a github.com/cenkalti/backoff/v5 v5.0.2 github.com/decred/dcrd/dcrec/secp256k1 v1.0.4 github.com/dunglas/httpsfv v1.0.2 @@ -56,7 +54,7 @@ require ( github.com/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e github.com/slok/go-http-metrics v0.13.0 github.com/streamplace/atproto-oauth-golang v0.0.0-20250619231223-a9c04fb888ac - github.com/streamplace/oatproxy v0.0.0-20250619231549-b15df1b82a3a + github.com/streamplace/oatproxy v0.0.0-20250722215839-1c878921d185 github.com/stretchr/testify v1.10.0 github.com/whyrusleeping/cbor-gen v0.3.1 github.com/whyrusleeping/go-did v0.0.0-20230824162731-404d1707d5d6 diff --git a/go.sum b/go.sum index e274a3ed..e44a2eb4 100644 --- a/go.sum +++ b/go.sum @@ -119,6 +119,8 @@ github.com/bkielbasa/cyclop v1.2.3 h1:faIVMIGDIANuGPWH031CZJTi2ymOQBULs9H21HSMa5 github.com/bkielbasa/cyclop v1.2.3/go.mod h1:kHTwA9Q0uZqOADdupvcFJQtp/ksSnytRMe8ztxG8Fuo= github.com/blizzy78/varnamelen v0.8.0 h1:oqSblyuQvFsW1hbBHh1zfwrKe3kcSj0rnXkKzsQ089M= github.com/blizzy78/varnamelen v0.8.0/go.mod h1:V9TzQZ4fLJ1DSrjVDfl89H7aMnTvKkApdHeyESmyR7k= +github.com/bluesky-social/indigo v0.0.0-20250729223159-573ae927246a h1:S12KN45uIkRglMHC8PqD/Vsz0+u3KbIaBF/6rit8/Pg= +github.com/bluesky-social/indigo v0.0.0-20250729223159-573ae927246a/go.mod h1:0XUyOCRtL4/OiyeqMTmr6RlVHQMDgw3LS7CfibuZR5Q= github.com/bmizerany/assert v0.0.0-20160611221934-b7ed37b82869 h1:DDGfHa7BWjL4YnC6+E63dPcxHo2sUxDIu8g3QgEJdRY= github.com/bmizerany/assert v0.0.0-20160611221934-b7ed37b82869/go.mod h1:Ekp36dRnpXw/yCqJaO+ZrUyxD+3VXMFFr56k5XYrpB4= github.com/bombsimon/wsl/v4 v4.7.0 h1:1Ilm9JBPRczjyUs6hvOPKvd7VL1Q++PL8M0SXBDf+jQ= @@ -934,10 +936,8 @@ github.com/streamplace/atproto-oauth-golang v0.0.0-20250619231223-a9c04fb888ac h github.com/streamplace/atproto-oauth-golang v0.0.0-20250619231223-a9c04fb888ac/go.mod h1:9LlKkqciiO5lRfbX0n4Wn5KNY9nvFb4R3by8FdW2TWc= github.com/streamplace/go-dpop v0.0.0-20250510031900-c897158a8ad4 h1:L1fS4HJSaAyNnkwfuZubgfeZy8rkWmA0cMtH5Z0HqNc= github.com/streamplace/go-dpop v0.0.0-20250510031900-c897158a8ad4/go.mod h1:bGUXY9Wd4mnd+XUrOYZr358J2f6z9QO/dLhL1SsiD+0= -github.com/streamplace/indigo v0.0.0-20250729000314-d9cf6fb86611 h1:xk4uA9CyEg9HaEHnIpnBAGq/zU72YA7kiFx08XsxJQk= -github.com/streamplace/indigo v0.0.0-20250729000314-d9cf6fb86611/go.mod h1:0XUyOCRtL4/OiyeqMTmr6RlVHQMDgw3LS7CfibuZR5Q= -github.com/streamplace/oatproxy v0.0.0-20250619231549-b15df1b82a3a h1:rbfSEm9IgK6jhXXtOExexL+ACNDOE/zJ5VYp6ujTpCA= -github.com/streamplace/oatproxy v0.0.0-20250619231549-b15df1b82a3a/go.mod h1:pXi24hA7xBHj8eEywX6wGqJOR9FaEYlGwQ/72rN6okw= +github.com/streamplace/oatproxy v0.0.0-20250722215839-1c878921d185 h1:knTi1I/zbppb+CmbDRmPKNgqO4QPOKdZBds63mMWTms= +github.com/streamplace/oatproxy v0.0.0-20250722215839-1c878921d185/go.mod h1:pXi24hA7xBHj8eEywX6wGqJOR9FaEYlGwQ/72rN6okw= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw= github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpEOglKo= diff --git a/js/docs/src/content/docs/lex-reference/openapi.json b/js/docs/src/content/docs/lex-reference/openapi.json index caad90af..73b72dd1 100644 --- a/js/docs/src/content/docs/lex-reference/openapi.json +++ b/js/docs/src/content/docs/lex-reference/openapi.json @@ -705,6 +705,97 @@ } } }, + "/xrpc/com.atproto.moderation.createReport": { + "post": { + "summary": "Submit a moderation report regarding an atproto account or record. Implemented by moderation services (with PDS proxying), and requires auth.", + "operationId": "com.atproto.moderation.createReport", + "tags": ["com.atproto.moderation"], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "id": { + "type": "integer" + }, + "reasonType": { + "$ref": "#/components/schemas/com.atproto.moderation.defs_reasonType" + }, + "reason": { + "type": "string", + "maxLength": 20000 + }, + "subject": { + "oneOf": [ + { + "$ref": "#/components/schemas/com.atproto.admin.defs_repoRef" + }, + { + "$ref": "#/components/schemas/com.atproto.repo.strongRef" + } + ] + }, + "reportedBy": { + "type": "string", + "format": "did" + }, + "createdAt": { + "type": "string", + "format": "date-time" + } + }, + "required": [ + "id", + "reasonType", + "subject", + "reportedBy", + "createdAt" + ] + } + } + } + } + }, + "requestBody": { + "required": true, + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "reasonType": { + "$ref": "#/components/schemas/com.atproto.moderation.defs_reasonType", + "description": "Indicates the broad category of violation the report is for." + }, + "reason": { + "type": "string", + "description": "Additional context about the content and violation.", + "maxLength": 20000 + }, + "subject": { + "oneOf": [ + { + "$ref": "#/components/schemas/com.atproto.admin.defs_repoRef" + }, + { + "$ref": "#/components/schemas/com.atproto.repo.strongRef" + } + ] + }, + "modTool": { + "$ref": "#/components/schemas/com.atproto.moderation.createReport_modTool" + } + }, + "required": ["reasonType", "subject"] + } + } + } + } + } + }, "/xrpc/com.atproto.identity.resolveHandle": { "get": { "summary": "Resolves an atproto handle (hostname) to a DID. Does not necessarily bi-directionally verify against the the DID document.", @@ -1550,6 +1641,33 @@ }, "required": ["uri", "cid", "value"] }, + "com.atproto.moderation.defs_reasonType": { + "type": "string" + }, + "com.atproto.admin.defs_repoRef": { + "type": "object", + "properties": { + "did": { + "type": "string", + "format": "did" + } + }, + "required": ["did"] + }, + "com.atproto.moderation.createReport_modTool": { + "type": "object", + "description": "Moderation tool information for tracing the source of the action", + "properties": { + "name": { + "type": "string", + "description": "Name/identifier of the source (e.g., 'bsky-app/android', 'bsky-web/chrome')" + }, + "meta": { + "description": "Additional arbitrary metadata about the source" + } + }, + "required": ["name"] + }, "app.bsky.feed.defs_skeletonFeedPost": { "type": "object", "properties": { diff --git a/lexicons/com/atproto/moderation/createReport.json b/lexicons/com/atproto/moderation/createReport.json new file mode 100644 index 00000000..2b8fc740 --- /dev/null +++ b/lexicons/com/atproto/moderation/createReport.json @@ -0,0 +1,90 @@ +{ + "lexicon": 1, + "id": "com.atproto.moderation.createReport", + "defs": { + "main": { + "type": "procedure", + "description": "Submit a moderation report regarding an atproto account or record. Implemented by moderation services (with PDS proxying), and requires auth.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["reasonType", "subject"], + "properties": { + "reasonType": { + "type": "ref", + "description": "Indicates the broad category of violation the report is for.", + "ref": "com.atproto.moderation.defs#reasonType" + }, + "reason": { + "type": "string", + "maxGraphemes": 2000, + "maxLength": 20000, + "description": "Additional context about the content and violation." + }, + "subject": { + "type": "union", + "refs": [ + "com.atproto.admin.defs#repoRef", + "com.atproto.repo.strongRef" + ] + }, + "modTool": { + "type": "ref", + "ref": "#modTool" + } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": [ + "id", + "reasonType", + "subject", + "reportedBy", + "createdAt" + ], + "properties": { + "id": { "type": "integer" }, + "reasonType": { + "type": "ref", + "ref": "com.atproto.moderation.defs#reasonType" + }, + "reason": { + "type": "string", + "maxGraphemes": 2000, + "maxLength": 20000 + }, + "subject": { + "type": "union", + "refs": [ + "com.atproto.admin.defs#repoRef", + "com.atproto.repo.strongRef" + ] + }, + "reportedBy": { "type": "string", "format": "did" }, + "createdAt": { "type": "string", "format": "datetime" } + } + } + } + }, + "modTool": { + "type": "object", + "description": "Moderation tool information for tracing the source of the action", + "required": ["name"], + "properties": { + "name": { + "type": "string", + "description": "Name/identifier of the source (e.g., 'bsky-app/android', 'bsky-web/chrome')" + }, + "meta": { + "type": "unknown", + "description": "Additional arbitrary metadata about the source" + } + } + } + } +} diff --git a/pkg/api/api_internal.go b/pkg/api/api_internal.go index f7d56e9b..db142739 100644 --- a/pkg/api/api_internal.go +++ b/pkg/api/api_internal.go @@ -606,7 +606,7 @@ func (a *StreamplaceAPI) InternalHandler(ctx context.Context) (http.Handler, err } after := time.Now().Add(-time.Duration(secs) * time.Second) w.Header().Set("Content-Type", "video/mp4") - err = a.MediaManager.ClipUser(ctx, user, w, nil, &after) + err = media.ClipUser(ctx, a.Model, a.CLI, user, w, nil, &after) if err != nil { errors.WriteHTTPInternalServerError(w, "unable to clip user", err) return diff --git a/pkg/media/clip_user.go b/pkg/media/clip_user.go index 4d630de2..3ffb3894 100644 --- a/pkg/media/clip_user.go +++ b/pkg/media/clip_user.go @@ -8,10 +8,12 @@ import ( "time" "stream.place/streamplace/pkg/aqtime" + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/model" ) -func (mm *MediaManager) ClipUser(ctx context.Context, user string, writer io.Writer, before *time.Time, after *time.Time) error { - segments, err := mm.model.LatestSegmentsForUser(user, -1, before, after) +func ClipUser(ctx context.Context, mod model.Model, cli *config.CLI, user string, writer io.Writer, before *time.Time, after *time.Time) error { + segments, err := mod.LatestSegmentsForUser(user, -1, before, after) if err != nil { return fmt.Errorf("unable to get segments: %w", err) } @@ -25,7 +27,7 @@ func (mm *MediaManager) ClipUser(ctx context.Context, user string, writer io.Wri segmentFiles := []string{} for _, segment := range segments { aqt := aqtime.FromTime(segment.StartTime) - fpath, err := mm.cli.SegmentFilePath(user, fmt.Sprintf("%s.%s", aqt.FileSafeString(), "mp4")) + fpath, err := cli.SegmentFilePath(user, fmt.Sprintf("%s.%s", aqt.FileSafeString(), "mp4")) if err != nil { return fmt.Errorf("unable to get segment file path: %w", err) } diff --git a/pkg/spxrpc/com_atproto_moderation.go b/pkg/spxrpc/com_atproto_moderation.go new file mode 100644 index 00000000..2c2bbb91 --- /dev/null +++ b/pkg/spxrpc/com_atproto_moderation.go @@ -0,0 +1,108 @@ +package spxrpc + +import ( + "context" + "fmt" + "net/http" + "time" + + comatprototypes "github.com/bluesky-social/indigo/api/atproto" + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/bluesky-social/indigo/xrpc" + "github.com/google/uuid" + "github.com/labstack/echo/v4" + "github.com/streamplace/oatproxy/pkg/oatproxy" + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/media" + "stream.place/streamplace/pkg/model" +) + +func (s *Server) handleComAtprotoModerationCreateReport(ctx context.Context, body *comatprototypes.ModerationCreateReport_Input) (*comatprototypes.ModerationCreateReport_Output, error) { + c, ok := ctx.Value(echoContextKey).(echo.Context) + if !ok { + return nil, echo.NewHTTPError(http.StatusInternalServerError, "echo context not found") + } + + atprotoProxy := c.Request().Header.Get("Atproto-Proxy") + if atprotoProxy == "" { + return nil, echo.NewHTTPError(http.StatusBadRequest, "Atproto-Proxy header is required (where are you sending this report?)") + } + + log.Log(ctx, "handleComAtprotoModerationCreateReport", "body", body) + + session, client := oatproxy.GetOAuthSession(ctx) + if session == nil { + return nil, echo.NewHTTPError(http.StatusUnauthorized, "oauth session not found") + } + + if body.Reason == nil { + empty := "" + body.Reason = &empty + } + + if body.Subject == nil { + return nil, echo.NewHTTPError(http.StatusBadRequest, "subject is required") + } + + var did string + + if body.Subject.AdminDefs_RepoRef != nil { + d, err := syntax.ParseDID(body.Subject.AdminDefs_RepoRef.Did) + if err != nil { + return nil, echo.NewHTTPError(http.StatusBadRequest, "invalid subject did") + } + did = d.String() + } else if body.Subject.RepoStrongRef != nil { + aturi, err := syntax.ParseATURI(body.Subject.RepoStrongRef.Uri) + if err != nil { + return nil, echo.NewHTTPError(http.StatusBadRequest, "invalid subject uri") + } + did = aturi.Authority().String() + } else { + return nil, echo.NewHTTPError(http.StatusBadRequest, "invalid subject") + } + + clipID, err := makeClip(ctx, s.cli, s.model, did) + if err != nil { + // we still want the report to go through! + log.Error(ctx, "failed to make clip for report", "error", err) + } else { + clipURL := fmt.Sprintf("https://%s/data/%s/clips/%s.mp4", s.cli.PublicHost, did, clipID) + newReason := fmt.Sprintf("%s\n\nClip: %s", *body.Reason, clipURL) + body.Reason = &newReason + } + + client.SetHeaders(map[string]string{ + "Atproto-Proxy": c.Request().Header.Get("Atproto-Proxy"), + }) + + var output comatprototypes.ModerationCreateReport_Output + err = client.Do(ctx, xrpc.Procedure, "application/json", "com.atproto.moderation.createReport", nil, body, &output) + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) + } + + return &output, nil +} + +func makeClip(ctx context.Context, cli *config.CLI, mod model.Model, did string) (string, error) { + after := time.Now().Add(-time.Duration(60) * time.Second) + + uu, err := uuid.NewV7() + if err != nil { + return "", echo.NewHTTPError(http.StatusInternalServerError, "failed to generate uuid") + } + + fd, err := cli.DataFileCreate([]string{did, "clips", fmt.Sprintf("%s.mp4", uu.String())}, false) + if err != nil { + return "", echo.NewHTTPError(http.StatusInternalServerError, "failed to create data file") + } + defer fd.Close() + + err = media.ClipUser(ctx, mod, cli, did, fd, nil, &after) + if err != nil { + return "", echo.NewHTTPError(http.StatusInternalServerError, "failed to clip user") + } + return uu.String(), nil +} diff --git a/pkg/spxrpc/spxrpc.go b/pkg/spxrpc/spxrpc.go index 071e678e..34f8bf8e 100644 --- a/pkg/spxrpc/spxrpc.go +++ b/pkg/spxrpc/spxrpc.go @@ -27,6 +27,7 @@ func NewServer(ctx context.Context, cli *config.CLI, model model.Model, op *oatp model: model, } e.Use(s.ErrorHandlingMiddleware()) + e.Use(s.ContextPreservingMiddleware()) e.Use(echomiddleware.Handler("", mdlw)) e.Use(op.OAuthMiddleware) err := s.RegisterHandlersPlaceStream(e) @@ -71,3 +72,23 @@ func (s *Server) ErrorHandlingMiddleware() echo.MiddlewareFunc { } } } + +// unique type to prevent assignment. +type echoContextKeyType struct{} + +// singleton value to identify our logging metadata in context +var echoContextKey = echoContextKeyType{} + +func (s *Server) ContextPreservingMiddleware() echo.MiddlewareFunc { + return func(next echo.HandlerFunc) echo.HandlerFunc { + return func(c echo.Context) error { + ctx := c.Request().Context() + if ctx == nil { + ctx = context.Background() + } + ctx = context.WithValue(ctx, echoContextKey, c) + c.SetRequest(c.Request().WithContext(ctx)) + return next(c) + } + } +} diff --git a/pkg/spxrpc/stubs.go b/pkg/spxrpc/stubs.go index 8b1410bf..368a8101 100644 --- a/pkg/spxrpc/stubs.go +++ b/pkg/spxrpc/stubs.go @@ -63,6 +63,7 @@ func (s *Server) RegisterHandlersChatBsky(e *echo.Echo) error { func (s *Server) RegisterHandlersComAtproto(e *echo.Echo) error { e.GET("/xrpc/com.atproto.identity.resolveHandle", s.HandleComAtprotoIdentityResolveHandle) + e.POST("/xrpc/com.atproto.moderation.createReport", s.HandleComAtprotoModerationCreateReport) e.GET("/xrpc/com.atproto.repo.describeRepo", s.HandleComAtprotoRepoDescribeRepo) e.GET("/xrpc/com.atproto.repo.getRecord", s.HandleComAtprotoRepoGetRecord) e.GET("/xrpc/com.atproto.repo.listRecords", s.HandleComAtprotoRepoListRecords) @@ -87,6 +88,24 @@ func (s *Server) HandleComAtprotoIdentityResolveHandle(c echo.Context) error { return c.JSON(200, out) } +func (s *Server) HandleComAtprotoModerationCreateReport(c echo.Context) error { + ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandleComAtprotoModerationCreateReport") + defer span.End() + + var body comatprototypes.ModerationCreateReport_Input + if err := c.Bind(&body); err != nil { + return err + } + var out *comatprototypes.ModerationCreateReport_Output + var handleErr error + // func (s *Server) handleComAtprotoModerationCreateReport(ctx context.Context,body *comatprototypes.ModerationCreateReport_Input) (*comatprototypes.ModerationCreateReport_Output, error) + out, handleErr = s.handleComAtprotoModerationCreateReport(ctx, &body) + if handleErr != nil { + return handleErr + } + return c.JSON(200, out) +} + func (s *Server) HandleComAtprotoRepoDescribeRepo(c echo.Context) error { ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandleComAtprotoRepoDescribeRepo") defer span.End() -- 2.51.2