Merge pull request #232 from Dvorinka/feat/content-audit

feat(storage): policy-level content audit on share creation
pull/3593/head
Tomáš Dvořák 2 weeks ago committed by GitHub
commit dbac8db105
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

@ -790,7 +790,8 @@
"report_abuse": "Report abuse",
"oauth_grant_create": "OAuth App Authorization",
"oauth_token_exchange": "OAuth Token Exchange",
"oauth_grant_revoke": "OAuth App Authorization Revoke"
"oauth_grant_revoke": "OAuth App Authorization Revoke",
"content_audit_blocked": "Share Blocked by Content Audit"
},
"server": "Server",
"tempPath": "Temporary path",
@ -1097,6 +1098,10 @@
"maxTotalSizeDes": "Maximum total bytes stored under this policy. Uploads and copies that would exceed it are rejected. 0 means unlimited.",
"overflowPolicy": "Overflow policy",
"overflowPolicyDes": "When this policy is out of capacity, new uploads spill into the selected policy. Chains of policies are supported.",
"auditEndpoint": "Content audit endpoint",
"auditEndpointDes": "HTTP endpoint for content audit (e.g. image moderation). When set, sharing a file stored on this policy POSTs its content to the endpoint; a {\"flagged\":true} response blocks the share. Leave empty to disable.",
"auditMaxSize": "Audit size limit",
"auditMaxSizeDes": "Maximum entity size sent to the audit endpoint. Larger files are skipped. 0 means no limit.",
"overflowNone": "None",
"enterFileExt": "Separated by semi-colon commas, leave blank to allow all file extensions.",
"extList": "File extension restrictions",

@ -790,7 +790,8 @@
"report_abuse": "举报滥用",
"oauth_grant_create": "OAuth 应用授权",
"oauth_token_exchange": "OAuth 令牌交换",
"oauth_grant_revoke": "OAuth 应用授权撤销"
"oauth_grant_revoke": "OAuth 应用授权撤销",
"content_audit_blocked": "分享被内容审计拦截"
},
"server": "服务器设置",
"tempPath": "临时路径",
@ -1097,6 +1098,10 @@
"maxTotalSizeDes": "此存储策略下可存储的文件总大小上限,超出后上传和复制将被拒绝。输入 0 表示不限制。",
"overflowPolicy": "溢出存储策略",
"overflowPolicyDes": "当此存储策略容量用尽时,新上传将自动切换到所选存储策略。支持多级链式溢出。",
"auditEndpoint": "内容审计端点",
"auditEndpointDes": "用于内容审计(鉴黄)的 HTTP 端点。设置后,分享此策略下的文件时会将文件内容 POST 到该端点;返回 {\"flagged\":true} 时阻止分享。留空表示不启用。",
"auditMaxSize": "审计大小上限",
"auditMaxSizeDes": "发送至审计端点的文件大小上限,超过此大小的文件跳过审计。输入 0 表示不限制。",
"overflowNone": "无",
"enterFileExt": "留空表示不限制文件扩展名,多个请以半角逗号 , 隔开。",
"extList": "文件扩展名限制",

@ -256,6 +256,8 @@ export interface PolicySetting {
thumb_max_size?: number;
max_total_size?: number;
overflow_policy_id?: number;
audit_endpoint?: string;
audit_max_size?: number;
relay?: boolean;
pre_allocate?: boolean;
media_meta_exts?: string[];

@ -487,6 +487,7 @@ export const AuditLogType = {
oauth_grant_create: 59,
oauth_token_exchange: 60,
oauth_grant_revoke: 61,
content_audit_blocked: 62,
};
export interface MultipleUriService {

@ -74,7 +74,13 @@ export const eventCategories = {
share: {
title: "settings.shareEvents",
description: "settings.shareEventsDes",
events: [AuditLogType.share, AuditLogType.share_link_viewed, AuditLogType.edit_share, AuditLogType.delete_share],
events: [
AuditLogType.share,
AuditLogType.share_link_viewed,
AuditLogType.edit_share,
AuditLogType.delete_share,
AuditLogType.content_audit_blocked,
],
},
version: {
title: "settings.versionEvents",

@ -145,6 +145,27 @@ const StorageAndUploadSection = () => {
[setPolicy],
);
const onAuditEndpointChange = useCallback(
(e: React.ChangeEvent<HTMLInputElement>) => {
const v = e.target.value.trim();
setPolicy((p: StoragePolicy) => ({
...p,
settings: { ...p.settings, audit_endpoint: v === "" ? undefined : v },
}));
},
[setPolicy],
);
const onAuditMaxSizeChange = useCallback(
(e: number) => {
setPolicy((p: StoragePolicy) => ({
...p,
settings: { ...p.settings, audit_max_size: e === 0 ? undefined : e },
}));
},
[setPolicy],
);
const fileExts = useMemo(() => {
return values.settings?.file_type?.join() ?? "";
}, [values.settings?.file_type]);
@ -335,6 +356,22 @@ const StorageAndUploadSection = () => {
<NoMarginHelperText>{t("policy.overflowPolicyDes")}</NoMarginHelperText>
</FormControl>
</SettingForm>
<SettingForm title={t("policy.auditEndpoint")} lgWidth={5}>
<FormControl fullWidth>
<DenseFilledTextField
placeholder="https://"
value={values.settings?.audit_endpoint ?? ""}
onChange={onAuditEndpointChange}
/>
<NoMarginHelperText>{t("policy.auditEndpointDes")}</NoMarginHelperText>
</FormControl>
</SettingForm>
<SettingForm title={t("policy.auditMaxSize")} lgWidth={5}>
<FormControl fullWidth>
<SizeInput variant={"outlined"} value={values.settings?.audit_max_size ?? 0} onChange={onAuditMaxSizeChange} />
<NoMarginHelperText>{t("policy.auditMaxSizeDes")}</NoMarginHelperText>
</FormControl>
</SettingForm>
<SettingForm title={t("policy.extList")} lgWidth={5}>
<FormControl fullWidth>
<DenseFilledTextField

@ -65,4 +65,5 @@ const (
EventOAuthGrantCreate = 59
EventOAuthTokenExchange = 60
EventOAuthGrantRevoke = 61
EventContentAuditBlocked = 62
)

@ -155,6 +155,12 @@ type (
// stored — they ignore the source file's policy and land here. 0 keeps
// thumbnails on the source policy.
ThumbStoragePolicyID int `json:"thumb_storage_policy_id,omitempty"`
// AuditEndpoint enables content audit: when a share is created over
// files on this policy, each image entity is POSTed to this URL and the
// share is rejected when the endpoint reports it flagged. Empty disables.
AuditEndpoint string `json:"audit_endpoint,omitempty"`
// AuditMaxSize caps the entity size sent for audit. 0 means no limit.
AuditMaxSize int64 `json:"audit_max_size,omitempty"`
// NativeMediaProcessing whether to use native media processing API from storage provider.
NativeMediaProcessing bool `json:"native_media_processing"`
// S3DeleteBatchSize the number of objects to delete in each batch.

@ -0,0 +1,125 @@
package manager
import (
"context"
"encoding/json"
"io"
"net/http"
"net/url"
"strconv"
"strings"
"time"
"github.com/cloudreve/Cloudreve/v4/ent"
"github.com/cloudreve/Cloudreve/v4/inventory/types"
"github.com/cloudreve/Cloudreve/v4/pkg/activity"
"github.com/cloudreve/Cloudreve/v4/pkg/filemanager/fs"
"github.com/cloudreve/Cloudreve/v4/pkg/request"
"github.com/cloudreve/Cloudreve/v4/pkg/serializer"
)
// auditTimeout bounds a single audit request so a stuck endpoint cannot hang
// share creation.
const auditTimeout = 30 * time.Second
type auditResponse struct {
Flagged bool `json:"flagged"`
Reason string `json:"reason"`
}
// auditShareEntities submits every auditable entity covered by a share to the
// owning policy's audit endpoint (PolicySetting.AuditEndpoint, 鉴黄). Only
// image entities within AuditMaxSize are sent; a flagged result aborts share
// creation. Endpoint failures are fail-closed — an unavailable audit service
// must not silently bypass the check.
func (l *manager) auditShareEntities(ctx context.Context, files []fs.File) error {
policyClient := l.dep.StoragePolicyClient()
mime := l.dep.MimeDetector(ctx)
seen := map[int]struct{}{}
policies := map[int]*ent.StoragePolicy{}
for _, file := range files {
// The extension lives on the file name; entity sources are blob keys.
fileMime := mime.TypeByName(file.Name())
if !strings.HasPrefix(fileMime, "image/") {
continue
}
for _, e := range file.Entities() {
if _, ok := seen[e.ID()]; ok {
continue
}
seen[e.ID()] = struct{}{}
if e.Type() != types.EntityTypeVersion || e.Size() == 0 {
continue
}
policy, ok := policies[e.PolicyID()]
if !ok {
p, err := policyClient.GetPolicyByID(ctx, e.PolicyID())
if err != nil {
return serializer.NewError(serializer.CodeContentAuditFailed, "failed to load entity storage policy", err)
}
policy = p
policies[e.PolicyID()] = p
}
endpoint := policy.Settings.AuditEndpoint
if endpoint == "" {
continue
}
if policy.Settings.AuditMaxSize > 0 && e.Size() > policy.Settings.AuditMaxSize {
continue
}
src, err := l.GetEntitySource(ctx, e.ID(), fs.WithEntity(e))
if err != nil {
return serializer.NewError(serializer.CodeContentAuditFailed, "failed to open entity for audit", err)
}
flagged, reason, err := callAuditEndpoint(l.dep.RequestClient(request.WithContext(ctx), request.WithTimeout(auditTimeout)),
endpoint, file.Name(), e.Size(), fileMime, src)
src.Close()
if err != nil {
return err
}
if flagged {
activity.Record(ctx, l.settings, l.dep.ActivityClient(), types.EventContentAuditBlocked,
activity.File(file.ID()), activity.Extra(map[string]any{
"entity_id": e.ID(),
"reason": reason,
}))
return serializer.NewError(serializer.CodeContentAuditRejected, "content rejected by audit", nil)
}
}
}
return nil
}
// callAuditEndpoint POSTs the entity body to the audit endpoint and returns the
// flagged verdict. Endpoint must be plain http(s); the endpoint answers
// {"flagged": bool, "reason": string}.
func callAuditEndpoint(client request.Client, endpoint, name string, size int64, mimeType string, body io.Reader) (bool, string, error) {
u, err := url.Parse(endpoint)
if err != nil || (u.Scheme != "http" && u.Scheme != "https") || u.Host == "" {
return false, "", serializer.NewError(serializer.CodeContentAuditFailed, "invalid audit endpoint", nil)
}
resp := client.Request(http.MethodPost, endpoint, body,
request.WithContentLength(size),
request.WithHeader(http.Header{
"Content-Type": {mimeType},
"X-Entity-Name": {name},
"X-Entity-Size": {strconv.FormatInt(size, 10)},
}),
)
raw, err := resp.CheckHTTPResponse(http.StatusOK).GetResponse()
if err != nil {
return false, "", serializer.NewError(serializer.CodeContentAuditFailed, "audit endpoint error", err)
}
var result auditResponse
if err := json.Unmarshal([]byte(raw), &result); err != nil {
return false, "", serializer.NewError(serializer.CodeContentAuditFailed, "invalid audit response", err)
}
return result.Flagged, result.Reason, nil
}

@ -0,0 +1,82 @@
package manager
import (
"io"
"net/http"
"net/http/httptest"
"strings"
"testing"
"github.com/cloudreve/Cloudreve/v4/pkg/request"
"github.com/cloudreve/Cloudreve/v4/pkg/serializer"
"github.com/stretchr/testify/require"
)
func newAuditTestClient(t *testing.T) request.Client {
return request.NewClient(nil, request.WithTimeout(0))
}
func TestCallAuditEndpoint(t *testing.T) {
t.Run("flagged", func(t *testing.T) {
var gotBody []byte
var gotName, gotSize, gotMime string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
gotBody, _ = io.ReadAll(r.Body)
gotName = r.Header.Get("X-Entity-Name")
gotSize = r.Header.Get("X-Entity-Size")
gotMime = r.Header.Get("Content-Type")
w.Write([]byte(`{"flagged":true,"reason":"nsfw"}`))
}))
t.Cleanup(srv.Close)
flagged, reason, err := callAuditEndpoint(newAuditTestClient(t), srv.URL, "pic.png", 4, "image/png", strings.NewReader("data"))
require.NoError(t, err)
require.True(t, flagged)
require.Equal(t, "nsfw", reason)
require.Equal(t, "data", string(gotBody))
require.Equal(t, "pic.png", gotName)
require.Equal(t, "4", gotSize)
require.Equal(t, "image/png", gotMime)
})
t.Run("clean", func(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Write([]byte(`{"flagged":false}`))
}))
t.Cleanup(srv.Close)
flagged, _, err := callAuditEndpoint(newAuditTestClient(t), srv.URL, "pic.png", 4, "image/png", strings.NewReader("data"))
require.NoError(t, err)
require.False(t, flagged)
})
t.Run("fail closed on endpoint error", func(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusInternalServerError)
}))
t.Cleanup(srv.Close)
_, _, err := callAuditEndpoint(newAuditTestClient(t), srv.URL, "pic.png", 4, "image/png", strings.NewReader("data"))
require.Error(t, err)
require.Equal(t, serializer.CodeContentAuditFailed, err.(serializer.AppError).Code)
})
t.Run("fail closed on bad json", func(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Write([]byte(`not json`))
}))
t.Cleanup(srv.Close)
_, _, err := callAuditEndpoint(newAuditTestClient(t), srv.URL, "pic.png", 4, "image/png", strings.NewReader("data"))
require.Error(t, err)
require.Equal(t, serializer.CodeContentAuditFailed, err.(serializer.AppError).Code)
})
t.Run("rejects non-http schemes", func(t *testing.T) {
for _, endpoint := range []string{"file:///etc/passwd", "ftp://x/", "://bad", ""} {
_, _, err := callAuditEndpoint(newAuditTestClient(t), endpoint, "pic.png", 4, "image/png", strings.NewReader("data"))
require.Error(t, err, endpoint)
require.Equal(t, serializer.CodeContentAuditFailed, err.(serializer.AppError).Code)
}
})
}

@ -319,7 +319,7 @@ func (l *manager) CreateOrUpdateShare(ctx context.Context, paths []*fs.URI, args
files := make([]fs.File, 0, len(paths))
seen := make(map[int]struct{}, len(paths))
for _, path := range paths {
file, err := l.fs.Get(ctx, path, dbfs.WithRequiredCapabilities(dbfs.NavigatorCapabilityShare), dbfs.WithNotRoot())
file, err := l.fs.Get(ctx, path, dbfs.WithRequiredCapabilities(dbfs.NavigatorCapabilityShare), dbfs.WithNotRoot(), dbfs.WithFileEntities())
if err != nil {
return nil, serializer.NewError(serializer.CodeNotFound, "src file not found", err)
}
@ -389,6 +389,12 @@ func (l *manager) CreateOrUpdateShare(ctx context.Context, paths []*fs.URI, args
Note: args.Note,
}
// Content audit: entities on policies with an audit endpoint are checked
// before the share record is created.
if err := l.auditShareEntities(ctx, files); err != nil {
return nil, err
}
fileIDs := lo.Map(files, func(f fs.File, _ int) int { return f.ID() })
share, err := shareClient.Upsert(ctx, &inventory.CreateShareParams{
OwnerID: file.OwnerID(),

@ -272,6 +272,10 @@ const (
CodeSmsCodeErr = 40095
// CodeFailedSendSms 短信发送失败
CodeFailedSendSms = 40096
// CodeContentAuditFailed 内容审核服务不可用或响应异常
CodeContentAuditFailed = 40097
// CodeContentAuditRejected 内容审核未通过
CodeContentAuditRejected = 40098
// CodeDBError 数据库操作失败
CodeDBError = 50001
// CodeEncryptError 加密失败

Loading…
Cancel
Save