Merge remote-tracking branch 'origin/main' into codex/bucket-metadata-merge-20260908

Signed-off-by: Feng Ruohang <rh@vonng.com>
This commit is contained in:
Feng Ruohang
2026-09-08 14:24:22 +08:00
5 changed files with 613 additions and 15 deletions
+346
View File
@@ -0,0 +1,346 @@
// Copyright (c) 2015-2025 MinIO, Inc.
// Copyright (c) 2025-2026 PGSTY
//
// This file is part of MinIO Object Storage stack
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// This program is distributed in the hope that it will be useful
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <http://www.gnu.org/licenses/>.
package cmd
import (
"bytes"
"encoding/json"
"encoding/xml"
"io"
"net/http"
"net/http/httptest"
"strings"
"sync"
"testing"
"github.com/minio/minio-go/v7/pkg/set"
"github.com/minio/minio/internal/auth"
"github.com/minio/minio/internal/config/dns"
"github.com/minio/minio/internal/hash"
xhttp "github.com/minio/minio/internal/http"
)
// federationRemoteCapture records the exact request headers the remote
// deployment receives on each forwarded request, so a test can assert what
// actually crossed the wire rather than what the destination ends up storing.
type federationRemoteCapture struct {
mu sync.Mutex
headers []http.Header
}
type federationResponseFilter struct {
http.ResponseWriter
filter func(http.Header)
}
func (w federationResponseFilter) WriteHeader(status int) {
w.filter(w.Header())
w.ResponseWriter.WriteHeader(status)
}
func (c *federationRemoteCapture) record(h http.Header) {
c.mu.Lock()
c.headers = append(c.headers, h.Clone())
c.mu.Unlock()
}
// reservedKeys returns every reserved-prefix header key seen across all
// forwarded requests.
func (c *federationRemoteCapture) reservedKeys() []string {
c.mu.Lock()
defer c.mu.Unlock()
var keys []string
for _, h := range c.headers {
for k := range h {
if stringsHasPrefixFold(k, ReservedMetadataPrefix) {
keys = append(keys, k)
}
}
}
return keys
}
// setupCopyObjectFederation makes the current process play both federation
// roles for a whole-object CopyObject. The destination bucket really exists in
// the shared backend, but the proxy's own bucket lookup is told it lives on a
// remote deployment, so CopyObjectHandler takes the legacy etcd federation
// branch and forwards the write through getRemoteInstanceClient and minio-go
// into a second HTTP endpoint that serves the real PutObjectHandler. That
// endpoint records the inbound headers first so a caller can inspect the wire.
//
// GetBucketLocation is deliberately left unregistered on that endpoint: the
// remoteBucketObjectLayer reports remoteBucket as missing (so the proxy takes
// the federation branch), which would also fail a real GetBucketLocation the
// minio-go client probes for. Leaving it unregistered makes the probe fall back
// to the default region, exactly as TestAPIFederatedCopyObjectPartChecksum does.
func setupCopyObjectFederation(t *testing.T, objectAPI ObjectLayer, apiRouter http.Handler,
instanceType, srcBucket string, responseFilters ...func(http.Header),
) (remoteBucket string, capture *federationRemoteCapture, cleanup func()) {
t.Helper()
remoteBucket = getRandomBucketName()
if err := objectAPI.MakeBucket(t.Context(), remoteBucket, MakeBucketOptions{}); err != nil {
t.Fatalf("%s: unable to create the remote bucket: %v", instanceType, err)
}
capture = &federationRemoteCapture{}
remote := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
capture.record(r.Header)
if len(responseFilters) != 0 && r.Method == http.MethodPut {
w = federationResponseFilter{ResponseWriter: w, filter: responseFilters[0]}
}
apiRouter.ServeHTTP(w, r)
}))
host, port, _ := strings.Cut(remote.Listener.Addr().String(), ":")
globalObjLayerMutex.Lock()
previousLayer := globalObjectAPI
globalObjectAPI = remoteBucketObjectLayer{ObjectLayer: previousLayer, remoteBucket: remoteBucket}
globalObjLayerMutex.Unlock()
previousDNS, previousFederation, previousIPs := globalDNSConfig, globalBucketFederation, globalDomainIPs
globalDNSConfig = federationTestDNS{records: map[string][]dns.SrvRecord{
srcBucket: {{Host: host, Port: json.Number(port)}},
remoteBucket: {{Host: host, Port: json.Number(port)}},
}}
// Every DNS record resolves to this process, so the bucket forwarding
// middleware always serves locally and only the handler proxies.
globalDomainIPs = set.CreateStringSet(remote.Listener.Addr().String())
globalBucketFederation = true
cleanup = func() {
remote.Close()
globalObjLayerMutex.Lock()
globalObjectAPI = previousLayer
globalObjLayerMutex.Unlock()
globalDNSConfig, globalBucketFederation, globalDomainIPs = previousDNS, previousFederation, previousIPs
}
return remoteBucket, capture, cleanup
}
// federatedCopyRequest drives a whole-object CopyObject at the proxy that
// forwards across deployments, and returns the recorded response.
func federatedCopyRequest(t *testing.T, apiRouter http.Handler, credentials auth.Credentials,
srcBucket, srcObject, dstBucket, dstObject string, headers map[string]string,
) *httptest.ResponseRecorder {
t.Helper()
req, err := newTestSignedRequestV4(http.MethodPut, getCopyObjectURL("", dstBucket, dstObject),
0, nil, credentials.AccessKey, credentials.SecretKey, headers)
if err != nil {
t.Fatalf("failed to build federated CopyObject request: %v", err)
}
req.Header.Set(xhttp.AmzCopySource, SlashSeparator+pathJoin(srcBucket, srcObject))
rec := httptest.NewRecorder()
apiRouter.ServeHTTP(rec, req)
return rec
}
// TestAPIFederatedCopyObjectInlineSource drives the legacy etcd federation
// branch of CopyObjectHandler end to end for a source object stored inline.
// Such a source carries x-minio-internal-inline-data in its stored metadata,
// which the remote deployment rejects as a reserved-prefix header. Before the
// fix the forwarded write failed with 400 InvalidArgument; the copy must now
// succeed and forward no reserved-prefix metadata at all.
func TestAPIFederatedCopyObjectInlineSource(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testAPIFederatedCopyObjectInlineSource,
endpoints: []string{"CopyObject", "PutObject", "HeadObject", "GetObject"},
})
}
func testAPIFederatedCopyObjectInlineSource(objectAPI ObjectLayer, instanceType, bucketName string,
apiRouter http.Handler, credentials auth.Credentials, t *testing.T,
) {
// A small object is stored inline, so its stored metadata carries
// x-minio-internal-inline-data. A user metadata key rides along to prove
// the fix strips only the reserved class, never ordinary metadata.
data := []byte("federated inline copy 26b!")
srcObject := "federation/inline-source"
putCopyChecksumSource(t, apiRouter, credentials, bucketName, srcObject, data,
map[string]string{"X-Amz-Meta-Origin": "inline-source"})
remoteBucket, capture, cleanup := setupCopyObjectFederation(t, objectAPI, apiRouter, instanceType, bucketName)
defer cleanup()
dstObject := "federation/inline-destination"
rec := federatedCopyRequest(t, apiRouter, credentials, bucketName, srcObject, remoteBucket, dstObject, nil)
if rec.Code != http.StatusOK {
t.Fatalf("%s: federated CopyObject of an inline source failed: %d %s",
instanceType, rec.Code, rec.Body.String())
}
// No reserved-prefix header may reach the remote deployment on any of the
// forwarded requests.
if leaked := capture.reservedKeys(); len(leaked) != 0 {
t.Fatalf("%s: forwarded reserved metadata to the remote: %v", instanceType, leaked)
}
// The destination object must be readable and byte-identical, and it must
// still carry the copied user metadata.
gr, err := objectAPI.GetObjectNInfo(t.Context(), remoteBucket, dstObject, nil, nil, ObjectOptions{})
if err != nil {
t.Fatalf("%s: unable to read the federated copy destination: %v", instanceType, err)
}
got, err := io.ReadAll(gr)
closeErr := gr.Close()
if err != nil {
t.Fatalf("%s: reading the federated copy destination failed: %v", instanceType, err)
}
if closeErr != nil {
t.Fatalf("%s: closing the federated copy destination failed: %v", instanceType, closeErr)
}
if !bytes.Equal(got, data) {
t.Fatalf("%s: federated copy destination body = %q, want %q", instanceType, got, data)
}
if origin, ok := gr.ObjInfo.UserDefined["X-Amz-Meta-Origin"]; !ok || origin != "inline-source" {
t.Fatalf("%s: destination lost copied user metadata: %v", instanceType, gr.ObjInfo.UserDefined)
}
}
// TestAPIFederatedCopyObjectRequestedChecksum drives the legacy etcd federation
// branch of CopyObjectHandler and verifies that a server-side checksum is both
// returned and persisted, matching the local CopyObject path. Before the fix
// the federated copy forwarded the write without asking for a checksum and
// discarded whatever the remote returned, so the response carried an empty
// checksum even when the client requested one (#99).
//
// The no-algorithm case is included deliberately: a checksum-less source gains
// the S3 default CRC-64NVME full-object checksum on the local path, so the
// federated path must return the same. "No requested algorithm" does not mean
// "no checksum".
func TestAPIFederatedCopyObjectRequestedChecksum(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testAPIFederatedCopyObjectRequestedChecksum,
endpoints: []string{"CopyObject", "PutObject", "HeadObject", "GetObject"},
})
}
func testAPIFederatedCopyObjectRequestedChecksum(objectAPI ObjectLayer, instanceType, bucketName string,
apiRouter http.Handler, credentials auth.Credentials, t *testing.T,
) {
data := []byte("federated copy checksum body")
srcObject := "federation/checksum-source"
putCopyChecksumSource(t, apiRouter, credentials, bucketName, srcObject, data, nil)
remoteBucket, _, cleanup := setupCopyObjectFederation(t, objectAPI, apiRouter, instanceType, bucketName)
defer cleanup()
cases := []struct {
name string
typ hash.ChecksumType
explicit bool
}{
{name: "CRC32", typ: hash.ChecksumCRC32, explicit: true},
{name: "CRC32C", typ: hash.ChecksumCRC32C, explicit: true},
{name: "SHA1", typ: hash.ChecksumSHA1, explicit: true},
{name: "SHA256", typ: hash.ChecksumSHA256, explicit: true},
{name: "CRC64NVME", typ: hash.ChecksumCRC64NVME, explicit: true},
// Default: no requested algorithm still yields the S3 CRC-64NVME.
{name: "default", typ: hash.ChecksumCRC64NVME},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
var headers map[string]string
if tc.explicit {
headers = map[string]string{xhttp.AmzChecksumAlgo: tc.typ.String()}
}
dstObject := "federation/checksum-destination-" + tc.name
rec := federatedCopyRequest(t, apiRouter, credentials, bucketName, srcObject, remoteBucket, dstObject, headers)
if rec.Code != http.StatusOK {
t.Fatalf("%s: federated CopyObject failed: %d %s", instanceType, rec.Code, rec.Body.String())
}
// The CopyObjectResult must carry the checksum of the copied bytes.
assertCopyChecksumResponse(t, rec, tc.typ, data)
// The remote must have persisted that same checksum.
assertCopyChecksum(t, objectAPI, remoteBucket, dstObject, tc.typ, data, false, nil)
})
}
}
func TestAPIFederatedCopyObjectRejectsInvalidRemoteChecksum(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: func(obj ObjectLayer, instanceType, bucket string, router http.Handler, credentials auth.Credentials, t *testing.T) {
data := []byte("remote checksum response fixture")
putCopyChecksumSource(t, router, credentials, bucket, "source", data, nil)
for _, value := range []string{"", "invalid-base64", "YQ=="} {
t.Run("checksum="+value, func(t *testing.T) {
remoteBucket, _, cleanup := setupCopyObjectFederation(t, obj, router, instanceType, bucket, func(header http.Header) {
header.Set(xhttp.AmzChecksumCRC32, value)
})
defer cleanup()
rec := federatedCopyRequest(t, router, credentials, bucket, "source", remoteBucket, "destination",
map[string]string{xhttp.AmzChecksumAlgo: "CRC32"})
if rec.Code < 500 {
t.Fatalf("invalid remote checksum must fail the copy, got %d: %s", rec.Code, rec.Body.String())
}
})
}
},
endpoints: []string{"CopyObject", "PutObject", "HeadObject", "GetObject"},
})
}
// TestAPIFederatedCopyObjectChecksumIsBoundToWrite guards the checksum
// representation: a federated copy that requests one algorithm must return only
// that algorithm, and a copy of a checksum-less source without a requested
// algorithm must never fabricate one other than the S3 default.
func TestAPIFederatedCopyObjectChecksumIsBoundToWrite(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testAPIFederatedCopyObjectChecksumIsBoundToWrite,
endpoints: []string{"CopyObject", "PutObject", "HeadObject", "GetObject"},
})
}
func testAPIFederatedCopyObjectChecksumIsBoundToWrite(objectAPI ObjectLayer, instanceType, bucketName string,
apiRouter http.Handler, credentials auth.Credentials, t *testing.T,
) {
data := []byte("federated copy single checksum body")
srcObject := "federation/single-checksum-source"
putCopyChecksumSource(t, apiRouter, credentials, bucketName, srcObject, data, nil)
remoteBucket, _, cleanup := setupCopyObjectFederation(t, objectAPI, apiRouter, instanceType, bucketName)
defer cleanup()
dstObject := "federation/single-checksum-destination"
rec := federatedCopyRequest(t, apiRouter, credentials, bucketName, srcObject, remoteBucket, dstObject,
map[string]string{xhttp.AmzChecksumAlgo: hash.ChecksumCRC32.String()})
if rec.Code != http.StatusOK {
t.Fatalf("%s: federated CopyObject failed: %d %s", instanceType, rec.Code, rec.Body.String())
}
var response CopyObjectResponse
if err := xml.Unmarshal(rec.Body.Bytes(), &response); err != nil {
t.Fatalf("%s: unable to decode CopyObjectResult: %v", instanceType, err)
}
if response.ChecksumCRC32 == "" {
t.Fatalf("%s: requested CRC32 checksum missing from response: %s", instanceType, rec.Body.String())
}
// Only the requested algorithm may be present.
if response.ChecksumCRC32C != "" || response.ChecksumSHA1 != "" ||
response.ChecksumSHA256 != "" || response.ChecksumCRC64NVME != "" {
t.Fatalf("%s: response carried checksums beyond the requested CRC32: %s", instanceType, rec.Body.String())
}
}
+69 -3
View File
@@ -1217,6 +1217,43 @@ var getRemoteInstanceClient = func(r *http.Request, host string) (*miniogo.Core,
return core, nil
}
// federatedChecksumType maps a requested server-side checksum type to the
// minio-go checksum the federation proxy asks the remote deployment to compute
// on the forwarded PutObject. Returns ChecksumNone for an unset/unknown type.
func federatedChecksumType(t hash.ChecksumType) miniogo.ChecksumType {
switch t.Base() {
case hash.ChecksumCRC32:
return miniogo.ChecksumCRC32
case hash.ChecksumCRC32C:
return miniogo.ChecksumCRC32C
case hash.ChecksumSHA1:
return miniogo.ChecksumSHA1
case hash.ChecksumSHA256:
return miniogo.ChecksumSHA256
case hash.ChecksumCRC64NVME:
return miniogo.ChecksumCRC64NVME
}
return miniogo.ChecksumNone
}
// federatedChecksumValue returns the base64 checksum the remote deployment
// reported for the requested type on the forwarded PutObject.
func federatedChecksumValue(t hash.ChecksumType, info miniogo.UploadInfo) string {
switch t.Base() {
case hash.ChecksumCRC32:
return info.ChecksumCRC32
case hash.ChecksumCRC32C:
return info.ChecksumCRC32C
case hash.ChecksumSHA1:
return info.ChecksumSHA1
case hash.ChecksumSHA256:
return info.ChecksumSHA256
case hash.ChecksumCRC64NVME:
return info.ChecksumCRC64NVME
}
return ""
}
// Check if the destination bucket is on a remote site, this code only gets executed
// when federation is enabled, ie when globalDNSConfig is non 'nil'.
//
@@ -1884,14 +1921,34 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
return
}
// Remove the metadata for remote calls.
delete(srcInfo.UserDefined, ReservedMetadataPrefix+"compression")
delete(srcInfo.UserDefined, ReservedMetadataPrefix+"actual-size")
// A plain federated PutObject must not carry any internal storage
// metadata. The remote rejects every reserved-prefix header as a class
// (containsReservedMetadata), so strip the whole class here rather than
// an enumerated subset: an inline source object also carries
// inline-data, and replication bookkeeping adds still more, so removing
// only compression/actual-size just defers the next rejected key. Match
// the remote's case-insensitive detection.
for k := range srcInfo.UserDefined {
if stringsHasPrefixFold(k, ReservedMetadataPrefix) {
delete(srcInfo.UserDefined, k)
}
}
opts := miniogo.PutObjectOptions{
UserMetadata: srcInfo.UserDefined,
ServerSideEncryption: dstOpts.ServerSideEncryption,
UserTags: tag.ToMap(),
}
// When a server-side checksum was requested (explicitly, inherited from
// the source, or the S3 default for a checksum-less object), the local
// path computes, persists and returns it; the federated path must do the
// same. Ask the remote to compute and persist that checksum by streaming
// it as a trailing checksum, so the forwarded write's response carries
// the value back to us. Without this the federated copy silently returns
// an empty checksum (#99).
wantChecksumType := dstOpts.WantServerSideChecksumType
if wantChecksumType.IsSet() {
opts.Checksum = federatedChecksumType(wantChecksumType)
}
remoteObjInfo, rerr := core.PutObject(ctx, dstBucket, dstObject, srcInfo.Reader,
srcInfo.Size, "", "", opts)
if rerr != nil {
@@ -1901,6 +1958,15 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re
objInfo.UserDefined = cloneMSS(opts.UserMetadata)
objInfo.ETag = remoteObjInfo.ETag
objInfo.ModTime = remoteObjInfo.LastModified
// Do not acknowledge a requested checksum the remote did not return.
if wantChecksumType.IsSet() {
cs := hash.NewChecksumWithType(wantChecksumType, federatedChecksumValue(wantChecksumType, remoteObjInfo))
if cs == nil {
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrInternalError), r.URL)
return
}
objInfo.Checksum = cs.AppendTo(nil, nil)
}
} else {
os = newObjSweeper(dstBucket, dstObject).WithVersioning(dstOpts.Versioned, dstOpts.VersionSuspended)
// Get appropriate object info to identify the remote object to delete
+4 -4
View File
@@ -395,9 +395,9 @@ func consoleLogIf(ctx context.Context, subsystem string, err error, errKind ...a
if err == nil {
return
}
if consoleTgt != nil {
if console := consoleTgt.Load(); console != nil {
entry := errToEntry(ctx, subsystem, err, errKind...)
consoleTgt.Send(ctx, entry)
(*console).Send(ctx, entry)
}
}
@@ -423,8 +423,8 @@ func sendLog(ctx context.Context, entry log.Entry) {
// Iterate over all logger targets to send the log entry
for _, t := range systemTgts {
if err := t.Send(ctx, entry); err != nil {
if consoleTgt != nil { // Sending to the console never fails
consoleTgt.Send(ctx, errToEntry(ctx, "logging", fmt.Errorf("unable to send log event to Logger target (%s): %v", t.String(), err), entry.Level))
if console := consoleTgt.Load(); console != nil { // Sending to the console never fails
(*console).Send(ctx, errToEntry(ctx, "logging", fmt.Errorf("unable to send log event to Logger target (%s): %v", t.String(), err), entry.Level))
}
}
}
+43 -8
View File
@@ -22,6 +22,7 @@ import (
"fmt"
"strings"
"sync"
"sync/atomic"
"github.com/minio/minio/internal/logger/target/http"
"github.com/minio/minio/internal/logger/target/kafka"
@@ -59,11 +60,35 @@ func (tl *targetsList) get() []Target {
return tl.list
}
func (tl *targetsList) add(t Target) {
// contains reports whether t is already registered.
func (tl *targetsList) contains(t Target) bool {
tl.mu.RLock()
defer tl.mu.RUnlock()
return tl.indexOf(t) >= 0
}
// addIfAbsent appends t unless it is already registered.
// Returns true if t was added.
func (tl *targetsList) addIfAbsent(t Target) bool {
tl.mu.Lock()
defer tl.mu.Unlock()
if tl.indexOf(t) >= 0 {
return false
}
tl.list = append(tl.list, t)
return true
}
// indexOf must be called with tl.mu held.
func (tl *targetsList) indexOf(t Target) int {
for i, existing := range tl.list {
if existing == t {
return i
}
}
return -1
}
func (tl *targetsList) set(tgts []Target) {
@@ -82,7 +107,10 @@ var (
auditTargets = newTargetsList()
// This is always set represent /dev/console target
consoleTgt Target
consoleTgt atomic.Pointer[Target]
// Init may emit logs, so serialize registration without holding the list lock.
systemTargetInitMu sync.Mutex
)
// SystemTargets returns active targets.
@@ -125,19 +153,26 @@ func CurrentStats() map[string]types.TargetStats {
}
// AddSystemTarget adds a new logger target to the
// list of enabled loggers
// list of enabled loggers. Adding a target that is already
// registered is a no-op, so callers may safely re-register
// long-lived targets such as the console logger.
func AddSystemTarget(ctx context.Context, t Target) error {
systemTargetInitMu.Lock()
defer systemTargetInitMu.Unlock()
if systemTargets.contains(t) {
return nil
}
if err := t.Init(ctx); err != nil {
return err
}
if consoleTgt == nil {
if t.Type() == types.TargetConsole {
consoleTgt = t
}
if t.Type() == types.TargetConsole {
consoleTgt.CompareAndSwap(nil, &t)
}
systemTargets.add(t)
systemTargets.addIfAbsent(t)
return nil
}
+151
View File
@@ -0,0 +1,151 @@
// Copyright (c) 2015-2026 MinIO, Inc.
//
// This file is part of MinIO Object Storage stack
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// This program is distributed in the hope that it will be useful
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <http://www.gnu.org/licenses/>.
package logger
import (
"context"
"errors"
"sync"
"sync/atomic"
"testing"
types "github.com/minio/minio/internal/logger/target/loggertypes"
)
// fakeTarget is a minimal Target used to exercise the registry.
type fakeTarget struct {
name string
kind types.TargetType
inits atomic.Int32
init func() error
}
func (f *fakeTarget) String() string { return f.name }
func (f *fakeTarget) Endpoint() string { return "" }
func (f *fakeTarget) Stats() types.TargetStats { return types.TargetStats{} }
func (f *fakeTarget) Init(context.Context) error {
f.inits.Add(1)
if f.init != nil {
return f.init()
}
return nil
}
func (f *fakeTarget) IsOnline(context.Context) bool { return true }
func (f *fakeTarget) Cancel() {}
func (f *fakeTarget) Send(context.Context, any) error { return nil }
func (f *fakeTarget) Type() types.TargetType { return f.kind }
// swapSystemTargets isolates the package-level registry for a test.
func swapSystemTargets(t *testing.T) {
t.Helper()
prevTargets, prevConsole := systemTargets, consoleTgt.Load()
systemTargets = newTargetsList()
consoleTgt.Store(nil)
t.Cleanup(func() {
systemTargets = prevTargets
consoleTgt.Store(prevConsole)
})
}
// AddSystemTarget must not register the same target twice: the console
// logger is re-added on every console-log subscription, and duplicates
// surface as identical series in the /logger/webhook metrics collector.
func TestAddSystemTargetIdempotent(t *testing.T) {
swapSystemTargets(t)
ctx := context.Background()
console := &fakeTarget{name: "console+http", kind: types.TargetConsole}
for range 3 {
if err := AddSystemTarget(ctx, console); err != nil {
t.Fatalf("AddSystemTarget: %v", err)
}
}
if got := len(SystemTargets()); got != 1 {
t.Fatalf("expected 1 system target after repeated add, got %d", got)
}
if got := console.inits.Load(); got != 1 {
t.Fatalf("expected Init to run once, ran %d times", got)
}
if got := consoleTgt.Load(); got == nil || *got != console {
t.Fatalf("consoleTgt not set to the console target")
}
other := &fakeTarget{name: "other", kind: types.TargetHTTP}
if err := AddSystemTarget(ctx, other); err != nil {
t.Fatalf("AddSystemTarget: %v", err)
}
if got := len(SystemTargets()); got != 2 {
t.Fatalf("expected 2 distinct system targets, got %d", got)
}
}
func TestAddSystemTargetConcurrent(t *testing.T) {
swapSystemTargets(t)
ctx := context.Background()
console := &fakeTarget{name: "console+http", kind: types.TargetConsole}
console.init = func() error {
// Initialization is allowed to inspect or write to the existing logger.
_ = SystemTargets()
consoleLogIf(ctx, "test", errors.New("initializing logger"))
return nil
}
var wg sync.WaitGroup
start := make(chan struct{})
for range 32 {
wg.Add(1)
go func() {
defer wg.Done()
<-start
if err := AddSystemTarget(ctx, console); err != nil {
t.Errorf("AddSystemTarget: %v", err)
}
consoleLogIf(ctx, "test", errors.New("concurrent logger"))
}()
}
close(start)
wg.Wait()
if got := len(SystemTargets()); got != 1 {
t.Fatalf("expected 1 system target after concurrent adds, got %d", got)
}
if got := console.inits.Load(); got != 1 {
t.Fatalf("expected one concurrent initialization, got %d", got)
}
}
func TestAddSystemTargetRetriesFailedInit(t *testing.T) {
swapSystemTargets(t)
ctx := context.Background()
initErr := errors.New("target unavailable")
target := &fakeTarget{name: "console", kind: types.TargetConsole}
target.init = func() error { return initErr }
if err := AddSystemTarget(ctx, target); !errors.Is(err, initErr) {
t.Fatalf("expected initialization error, got %v", err)
}
if len(SystemTargets()) != 0 || consoleTgt.Load() != nil {
t.Fatal("failed initialization published a target")
}
target.init = nil
if err := AddSystemTarget(ctx, target); err != nil {
t.Fatal(err)
}
if got := len(SystemTargets()); got != 1 || target.inits.Load() != 2 {
t.Fatalf("retry: targets=%d, initializations=%d", got, target.inits.Load())
}
}