mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-08-03 05:46:33 +00:00
* fix(s3/audit): emit audit log for successful GET/HEAD Successful GET/HEAD object requests never produced a fluent audit entry because those handlers write the response directly (streaming for GET, WriteHeader for HEAD) and never reach a PostLog call site. The wiki advertises GET as an audited verb, so the asymmetry surprises operators who rely on the log for read-access auditing. Move the safety net into the track() middleware: tag each request with an audit-tracking flag, let PostLog/PostAccessLog (delete path) mark it, and emit a single fallback entry after the handler returns when nothing fired. The recorder's status flows into the fallback so the audit row still reflects 200/206 vs 404 etc. No double logging for handlers that already emit (write helpers, error paths, bulk delete). Refs #9463 * fix(s3/audit): defensive nil checks on audit-tracking helpers Address PR review: guard against nil request and nil *atomic.Bool stored under the audit-tracking key. The conditions are unreachable today (the key is private and we only ever store new(atomic.Bool)), but the checks are free and keep the helpers safe if a future caller misbehaves. * test(s3/audit): track() audit fallback coverage + stale comment cleanup (#9469) test(s3/audit): cover track() fallback wiring + cleanup Adds two unit tests in weed/s3api/stats_test.go that exercise the audit-tracking flag set up by track(): one verifies the fallback path fires when a handler writes the response directly (the GET/HEAD object regression in #9463), the other verifies the flag is set when a handler emits PostLog itself so the fallback is skipped. To make the wiring observable without standing up fluent, PostLog now marks the audit flag before short-circuiting on a nil Logger; production behavior is unchanged (no logger, no posting) but the flag stays consistent. Also drops two stale comments in s3api_object_handlers.go that still referenced proxyToFiler — that helper was removed when GET/HEAD started streaming from volume servers directly. Stacks on #9467.
257 lines
8.2 KiB
Go
257 lines
8.2 KiB
Go
package s3err
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"net"
|
|
"net/http"
|
|
"os"
|
|
"strings"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
"github.com/fluent/fluent-logger-golang/fluent"
|
|
"github.com/seaweedfs/seaweedfs/weed/glog"
|
|
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
|
|
"github.com/seaweedfs/seaweedfs/weed/util/request_id"
|
|
)
|
|
|
|
type AccessLogExtend struct {
|
|
AccessLog
|
|
AccessLogHTTP
|
|
}
|
|
|
|
type AccessLog struct {
|
|
Bucket string `msg:"bucket" json:"bucket"` // awsexamplebucket1
|
|
Time int64 `msg:"time" json:"time"` // [06/Feb/2019:00:00:38 +0000]
|
|
RemoteIP string `msg:"remote_ip" json:"remote_ip,omitempty"` // 192.0.2.3
|
|
Requester string `msg:"requester" json:"requester,omitempty"` // IAM user id
|
|
RequestID string `msg:"request_id" json:"request_id,omitempty"` // 3E57427F33A59F07
|
|
Operation string `msg:"operation" json:"operation,omitempty"` // REST.HTTP_method.resource_type REST.PUT.OBJECT
|
|
Key string `msg:"key" json:"key,omitempty"` // /photos/2019/08/puppy.jpg
|
|
ErrorCode string `msg:"error_code" json:"error_code,omitempty"`
|
|
HostId string `msg:"host_id" json:"host_id,omitempty"`
|
|
HostHeader string `msg:"host_header" json:"host_header,omitempty"` // s3.us-west-2.amazonaws.com
|
|
UserAgent string `msg:"user_agent" json:"user_agent,omitempty"`
|
|
HTTPStatus int `msg:"status" json:"status,omitempty"`
|
|
SignatureVersion string `msg:"signature_version" json:"signature_version,omitempty"`
|
|
}
|
|
|
|
type AccessLogHTTP struct {
|
|
RequestURI string `json:"request_uri,omitempty"` // "GET /awsexamplebucket1/photos/2019/08/puppy.jpg?x-foo=bar HTTP/1.1"
|
|
BytesSent string `json:"bytes_sent,omitempty"`
|
|
ObjectSize string `json:"object_size,omitempty"`
|
|
TotalTime int `json:"total_time,omitempty"`
|
|
TurnAroundTime int `json:"turn_around_time,omitempty"`
|
|
Referer string `json:"Referer,omitempty"`
|
|
VersionId string `json:"version_id,omitempty"`
|
|
CipherSuite string `json:"cipher_suite,omitempty"`
|
|
AuthenticationType string `json:"auth_type,omitempty"`
|
|
TLSVersion string `json:"TLS_version,omitempty"`
|
|
}
|
|
|
|
const tag = "s3.access"
|
|
|
|
var (
|
|
Logger *fluent.Fluent
|
|
hostname = os.Getenv("HOSTNAME")
|
|
environment = os.Getenv("ENVIRONMENT")
|
|
)
|
|
|
|
func InitAuditLog(config string) {
|
|
configContent, readErr := os.ReadFile(config)
|
|
if readErr != nil {
|
|
glog.Errorf("fail to read fluent config %s : %v", config, readErr)
|
|
return
|
|
}
|
|
fluentConfig := &fluent.Config{}
|
|
if err := json.Unmarshal(configContent, fluentConfig); err != nil {
|
|
glog.Errorf("fail to parse fluent config %s : %v", string(configContent), err)
|
|
return
|
|
}
|
|
if len(fluentConfig.TagPrefix) == 0 && len(environment) > 0 {
|
|
fluentConfig.TagPrefix = environment
|
|
}
|
|
fluentConfig.Async = true
|
|
fluentConfig.AsyncResultCallback = func(data []byte, err error) {
|
|
if err != nil {
|
|
glog.Warning("Error while posting log: ", err)
|
|
}
|
|
}
|
|
var err error
|
|
Logger, err = fluent.New(*fluentConfig)
|
|
if err != nil {
|
|
glog.Errorf("fail to load fluent config: %v", err)
|
|
}
|
|
}
|
|
|
|
// getRemoteIP returns the client IP for the audit log, honoring forwarding
|
|
// headers set by reverse proxies. Preference order: X-Forwarded-For (first
|
|
// non-empty entry), X-Real-IP, then r.RemoteAddr. Headers are trusted as-is;
|
|
// operators who expose the S3 endpoint directly to untrusted networks should
|
|
// strip these headers at the proxy boundary so clients cannot spoof them.
|
|
func getRemoteIP(r *http.Request) string {
|
|
if forwardedFor := r.Header.Get("X-Forwarded-For"); forwardedFor != "" {
|
|
for _, entry := range strings.Split(forwardedFor, ",") {
|
|
if ip := strings.TrimSpace(entry); ip != "" {
|
|
return ip
|
|
}
|
|
}
|
|
}
|
|
if realIP := strings.TrimSpace(r.Header.Get("X-Real-IP")); realIP != "" {
|
|
return realIP
|
|
}
|
|
if host, _, err := net.SplitHostPort(r.RemoteAddr); err == nil {
|
|
return host
|
|
}
|
|
return r.RemoteAddr
|
|
}
|
|
|
|
func getREST(httpMetod string, resourceType string) string {
|
|
return fmt.Sprintf("REST.%s.%s", httpMetod, resourceType)
|
|
}
|
|
|
|
func getResourceType(object string, query_key string, metod string) (string, bool) {
|
|
if object == "/" {
|
|
switch query_key {
|
|
case "delete":
|
|
return "BATCH.DELETE.OBJECT", true
|
|
case "tagging":
|
|
return getREST(metod, "OBJECTTAGGING"), true
|
|
case "lifecycle":
|
|
return getREST(metod, "LIFECYCLECONFIGURATION"), true
|
|
case "acl":
|
|
return getREST(metod, "ACCESSCONTROLPOLICY"), true
|
|
case "policy":
|
|
return getREST(metod, "BUCKETPOLICY"), true
|
|
default:
|
|
return getREST(metod, "BUCKET"), false
|
|
}
|
|
} else {
|
|
switch query_key {
|
|
case "tagging":
|
|
return getREST(metod, "OBJECTTAGGING"), true
|
|
default:
|
|
return getREST(metod, "OBJECT"), false
|
|
}
|
|
}
|
|
}
|
|
|
|
func getOperation(object string, r *http.Request) string {
|
|
queries := r.URL.Query()
|
|
var operation string
|
|
var queryFound bool
|
|
for key, _ := range queries {
|
|
operation, queryFound = getResourceType(object, key, r.Method)
|
|
if queryFound {
|
|
return operation
|
|
}
|
|
}
|
|
if len(queries) == 0 {
|
|
operation, _ = getResourceType(object, "", r.Method)
|
|
}
|
|
return operation
|
|
}
|
|
|
|
func GetAccessLog(r *http.Request, HTTPStatusCode int, s3errCode ErrorCode) *AccessLog {
|
|
bucket, key := s3_constants.GetBucketAndObject(r)
|
|
var errorCode string
|
|
if s3errCode != ErrNone {
|
|
errorCode = GetAPIError(s3errCode).Code
|
|
}
|
|
remoteIP := getRemoteIP(r)
|
|
hostHeader := r.Header.Get("X-Forwarded-Host")
|
|
if len(hostHeader) == 0 {
|
|
hostHeader = r.Host
|
|
}
|
|
return &AccessLog{
|
|
HostHeader: hostHeader,
|
|
RequestID: request_id.GetFromRequest(r),
|
|
RemoteIP: remoteIP,
|
|
Requester: s3_constants.GetIdentityNameFromContext(r), // Get from context, not header (secure)
|
|
SignatureVersion: r.Header.Get(s3_constants.AmzAuthType),
|
|
UserAgent: r.Header.Get("user-agent"),
|
|
HostId: hostname,
|
|
Bucket: bucket,
|
|
HTTPStatus: HTTPStatusCode,
|
|
Time: time.Now().Unix(),
|
|
Key: key,
|
|
Operation: getOperation(key, r),
|
|
ErrorCode: errorCode,
|
|
}
|
|
}
|
|
|
|
func PostLog(r *http.Request, HTTPStatusCode int, errorCode ErrorCode) {
|
|
if r == nil {
|
|
return
|
|
}
|
|
// Mark before the Logger nil-check so the middleware fallback in track()
|
|
// still sees that audit was handled by the caller in deployments that
|
|
// haven't configured fluent — keeps the flag behavior consistent and
|
|
// makes wiring testable without standing up a fluent server.
|
|
markAuditLogged(r)
|
|
if Logger == nil {
|
|
return
|
|
}
|
|
if err := Logger.Post(tag, *GetAccessLog(r, HTTPStatusCode, errorCode)); err != nil {
|
|
glog.Warning("Error while posting log: ", err)
|
|
}
|
|
}
|
|
|
|
func PostAccessLog(log AccessLog) {
|
|
if Logger == nil || len(log.Key) == 0 {
|
|
return
|
|
}
|
|
if err := Logger.Post(tag, log); err != nil {
|
|
glog.Warning("Error while posting log: ", err)
|
|
}
|
|
}
|
|
|
|
// auditLogCtxKey identifies the per-request flag used to detect whether an
|
|
// audit entry has already been emitted, so middleware can supply a fallback
|
|
// log for handlers that don't call PostLog themselves without double-logging
|
|
// the ones that do.
|
|
type auditLogCtxKey struct{}
|
|
|
|
// EnsureAuditTracking attaches an audit-tracking flag to the request context
|
|
// if one is not already present. Safe to call when no fluent logger is
|
|
// configured; the flag is harmless in that case.
|
|
func EnsureAuditTracking(r *http.Request) *http.Request {
|
|
if r == nil {
|
|
return nil
|
|
}
|
|
if v, ok := r.Context().Value(auditLogCtxKey{}).(*atomic.Bool); ok && v != nil {
|
|
return r
|
|
}
|
|
flag := new(atomic.Bool)
|
|
return r.WithContext(context.WithValue(r.Context(), auditLogCtxKey{}, flag))
|
|
}
|
|
|
|
// MarkAuditLogged signals that an audit entry has already been emitted for r.
|
|
// Callers that emit logs via paths other than PostLog (e.g. PostAccessLog in
|
|
// batch delete) should invoke this so the middleware fallback skips r.
|
|
func MarkAuditLogged(r *http.Request) {
|
|
markAuditLogged(r)
|
|
}
|
|
|
|
func markAuditLogged(r *http.Request) {
|
|
if r == nil {
|
|
return
|
|
}
|
|
if v, ok := r.Context().Value(auditLogCtxKey{}).(*atomic.Bool); ok && v != nil {
|
|
v.Store(true)
|
|
}
|
|
}
|
|
|
|
// AuditAlreadyLogged reports whether PostLog (or MarkAuditLogged) has run for r.
|
|
func AuditAlreadyLogged(r *http.Request) bool {
|
|
if r == nil {
|
|
return false
|
|
}
|
|
if v, ok := r.Context().Value(auditLogCtxKey{}).(*atomic.Bool); ok && v != nil {
|
|
return v.Load()
|
|
}
|
|
return false
|
|
}
|