Files
seaweedfs/weed/admin/dash/admin_server.go
T
Mathieu ArnoldandGitHub e931cccc7b Manage bucket policies via the admin ui (#10895)
* admin: manage S3 bucket policies from the admin UI

Bucket policies were only manageable through the S3 PutBucketPolicy API;
the admin UI had no equivalent to the quota/owner/lifecycle editors it
already offers. Add GET/PUT/DELETE for a bucket's policy, sharing the
exact validation the S3 gateway uses.

- Extract validateBucketPolicy/validateResourceForBucket out of
  s3api_bucket_policy_handlers.go into policy_engine.ValidateBucketPolicy /
  ResourceMatchesBucket so both the S3 API and the admin UI enforce
  identical rules.
- weed/admin/dash/bucket_policy.go: Get/Set/DeleteBucketPolicy, writing
  through ObjectTransaction + PATCH_EXTENDED (the lifecycle pattern) so a
  concurrent owner/quota/lifecycle change on the same bucket entry isn't
  clobbered. Propagation to every S3 gateway is automatic via the existing
  filer metadata log subscription. The S3 gateway's IAM policy mirror is
  deliberately not replicated here (its delete path is already an
  unimplemented TODO on the S3 side).
- New GET/PUT/DELETE /api/s3/buckets/{bucket}/policy routes, CSRF-guarded
  on writes.
- Bucket list and details modal now show a statement-count badge, read
  from the entry already fetched (no extra RPC).
- UI: a JSON-textarea policy editor modal, matching the lifecycle modal's
  structure.

* admin: reuse the visual policy editor for bucket policies

Extract the structured policy editor (add/remove statement, action/
resource/principal rows with autocomplete, JSON tab kept in sync) out of
policies.templ's inline script into a shared
weed/admin/static/js/policy_editor.js, and wire the bucket policy modal
in s3_buckets.templ up to it instead of a bare JSON textarea.

- registerPolicyEditor(which, config) replaces the hardcoded create/edit
  id derivation with a per-instance config (textarea/tab/body ids,
  datalist ids, requirePrincipal, bucket). The IAM policies page keeps its
  exact pre-extraction ids via two registerPolicyEditor calls, so its
  markup is unchanged.
- New policy_datalists.templ exposes the three shared <datalist>s
  (actions/resources/principals) as @PolicyDatalists(), now rendered by
  both policies.templ and s3_buckets.templ.
- requirePrincipal seeds new bucket-policy statements with Principal: "*"
  and adds a client-side check before save (the server, via
  policy_engine.ValidateBucketPolicy, remains the actual authority); the
  bucket config pins the Resource autocomplete to the open bucket instead
  of fetching every bucket in the cluster.
- layout.templ loads policy_editor.js globally, after admin.js/
  modal-alerts.js (basePath/escapeHtml/showAlert) which it depends on.

3a (the extraction) is a byte-preserving move verified against the
unchanged policies.templ behavior before layering 3b's parameterization
and the bucket-policy wiring on top.

* admin: migrate S3 Tables bucket/table policy editors to the shared editor

Third consumer of the shared visual policy editor: the S3 Tables bucket
and table policy modals (a bare JSON textarea each) now get the same
structured Editor/JSON tabs as the bucket policy and IAM policy pages,
via registerPolicyEditor('s3tablesBucketPolicy'/'s3tablesTablePolicy',
{ textareaId: ... }). Storage and validation are untouched - S3 Tables
policies still go through their own s3tables.PolicyDocument type and the
s3tables.policy extended attribute, unrelated to policy_engine and
s3-bucket-policy; only the editor UI is shared.

Fix a real bug surfaced by adding this second load path: the bucket
policy modal (and the naive first draft of this s3tables port) called
commitPolicyTextareaToEditor() right after a GET and then force-switched
to the Editor tab. commitPolicyTextareaToEditor() is designed to leave
the current tab in place and the editor state untouched when a document
fails to parse (so an in-progress edit survives a bad tab switch), so
forcing the Editor tab afterwards could show empty/stale editor state
that a careless Save would then serialize over a perfectly valid but
structurally-unusual stored policy. Add
loadPolicyTextareaIntoEditor(which) to policy_editor.js, which has no
"current tab" to defer to and instead falls back to the JSON tab with an
alert on a document the structured editor can't represent - the same
safety editPolicy already had in policies.templ - and use it at all three
"populate the editor right after a GET" call sites (bucket policy,
S3 Tables bucket policy, S3 Tables table policy).

* admin: show policy statement count on the S3 Tables buckets page

Mirrors the "Policy" column already added to the classic S3 buckets
list: a clickable badge with the statement count when the table bucket
has a resource policy, "Not configured" otherwise. S3 Tables policies
are a separate mechanism (s3tables.PolicyDocument under the
s3tables.policy extended attribute) from the S3 bucket policy work
elsewhere in this branch (policy_engine.PolicyDocument /
s3-bucket-policy), so this is a parallel implementation of the same
pattern rather than shared code.

- S3TablesBucketSummary gains PolicyStatementCount, populated in
  GetS3TablesBucketsData from entry.Entry.Extended[s3tables.ExtendedKeyPolicy]
  via the new extractS3TablesPolicyStatementCountFromEntry - no extra RPC,
  the entry is already fetched for ExtendedKeyMetadata.
- The badge reuses the existing .s3tables-bucket-policy-btn class, so it
  opens the same policy modal as the row's action button with no JS
  changes.

* admin: don't let a failed policy GET open the door to an empty overwrite

loadS3TablesBucketPolicy/loadS3TablesTablePolicy cleared the textarea,
then unconditionally called loadPolicyTextareaIntoEditor() regardless of
whether the GET actually succeeded - including when fetch() rejected or
the response was not ok, silently logged to console only. That leaves
the structured editor holding a legitimate-looking empty policy
({version, statements: []}), with the Editor tab active by default.

If Save is then clicked, commitPolicyActiveTab() serializes that empty
state into the textarea as `{"Version":"2012-10-17","Statement":[]}` -
a non-empty string - before the "Policy JSON is required" guard ever
sees it, so the guard passes and the transient load failure gets
written over whatever policy was actually stored.

Add s3tablesBucketPolicyLoaded/s3tablesTablePolicyLoaded, set true only
once a GET has actually completed (ok, including a genuinely empty
policy) and false on any failure path (fetch rejection or a non-ok
response, which previously fell through silently). Both submit handlers
now check the flag before touching the editor at all, and a failed load
surfaces via alert() instead of only a console.error - the user
previously had no visible indication the load had failed.

Verified with a jsdom simulation driving the real rendered page against
a stubbed fetch: a failed GET followed by Save now sends no PUT at all
(previously it sent Statement: []); a successful GET followed by Save
still PUTs the loaded policy unchanged.

* admin: address code review findings on the policy editor

1. policy_editor.js: policyEditors is only pre-populated for 'create'/
   'edit'; every other `which` (bucket, s3tablesBucket, s3tablesTable)
   stays undefined until its first successful async load. Nothing in
   this file enforces that a page hide its Editor/JSON tabs and
   Add-statement button until that load completes - the S3 Tables policy
   modals don't - so a click in that window (e.g. Add statement, or
   switching to the JSON tab) threw "Cannot read properties of undefined
   (reading 'unparsed')". Add policyEditorState(which), which lazily
   initializes a default state, and route addPolicyStatement, the
   jsonTabBtn 'show.bs.tab' handler, commitPolicyActiveTab, and
   renderPolicyEditor through it. Verified with a jsdom simulation
   against a never-resolving fetch: the exact click threw on the
   pre-fix code and no longer does.

2. s3_buckets.templ: the bucket-policy Save handler checked the
   textarea for emptiness before calling commitPolicyActiveTab(), which
   is what actually serializes the structured Editor tab's fields into
   that textarea. A policy entered entirely through the Editor tab (the
   primary path - never touching the JSON tab) left the textarea at
   whatever it was at load time, so creating a new policy this way hit
   "Enter a policy document" and Save silently did nothing. Move the
   commit before the emptiness check, preserving the existing alert and
   early-return. Verified with a jsdom simulation: Add-statement then
   Save (no tab switch) now PUTs the entered statement; before the fix
   the same sequence never reached fetch().

3. s3tables_buckets.templ / s3tables_tables.templ: the policy Editor/
   JSON nav-tabs were missing the ARIA roles Bootstrap's own tab pattern
   expects (role="tab"/"tabpanel", aria-selected, aria-controls,
   aria-labelledby) - screen readers had no way to tell these were tabs
   or which pane went with which button. Added the standard Bootstrap 5
   tab markup to both.

* admin: guard policy load/save flows against overlapping requests

1. s3tables.js: loadS3TablesBucketPolicy/loadS3TablesTablePolicy had no
   protection against overlapping loads. Opening one bucket's (or
   table's) policy dialog and then another's before the first GET
   resolved let the late response write its document into the shared
   textarea and mark the dialog "loaded" while it was now targeting the
   second resource - a subsequent Save would then push the first
   resource's policy onto the second. Add a per-load monotonic sequence
   number (s3tablesBucketPolicyRequestSeq / s3tablesTablePolicyRequestSeq,
   the same pattern already used for the classic bucket-policy load in
   s3_buckets.templ); a response is only applied - textarea, loaded flag,
   editor state - if its captured sequence still matches the latest one
   issued.

   Verified with a jsdom simulation: bucket A's policy load (artificially
   slow) followed immediately by bucket B's (fast) previously left A's
   policy in the textarea once A's late response landed; it now correctly
   keeps B's.

2. s3_buckets.templ: the bucket-policy Save button lives outside the
   (initially hidden) editor wrapper, so it stays clickable while a load
   is still in flight - the existing policyRequestSeq guard only protects
   the *load* from a stale response, not Save from firing before any
   load for the current bucket has completed. Add bucketPolicyLoaded,
   reset before each GET and set only once the matching response lands,
   and check it at the top of the Save handler.

   Verified with a jsdom simulation: clicking Save immediately after
   opening the dialog, before a (deliberately never-resolving) GET
   settles, now sends no PUT; a normal load-then-save sequence still
   PUTs the loaded policy unchanged.

* admin: address further code review findings on the policy editor

1. s3tables.js: loadS3TablesBucketPolicy/loadS3TablesTablePolicy only
   reset the JSON textarea when a new load starts; the structured editor
   kept showing the previously loaded resource's statements (Editor tab
   is the default active one) until the new fetch resolved. Call
   loadPolicyTextareaIntoEditor() against the now-cleared textarea
   immediately, so switching resources visibly resets the editor right
   away instead of only once its own load completes. Verified with jsdom:
   opening bucket A (loads fully) then bucket B (GET never resolves) no
   longer leaves A's statements visible in B's editor.

2. s3tables.js: deleteS3TablesBucketPolicy/deleteS3TablesTablePolicy had
   no loaded-state check, so a failed GET (which already blocks Save)
   left Delete fully able to remove the resource's stored policy sight
   unseen. Add the same s3tablesBucketPolicyLoaded/s3tablesTablePolicyLoaded
   guard Save already uses. Verified with jsdom: delete after a failed
   load now sends no DELETE; delete after a successful load is unaffected.

3. s3_buckets.templ: the bucket-policy Editor/JSON nav-tabs were missing
   the same ARIA roles already added to the S3 Tables policy tabs in an
   earlier round (role="tab"/"tabpanel", aria-selected, aria-controls,
   aria-labelledby) - this instance was out of scope for that review
   comment but is the same gap. Bootstrap's own tab.js already manages
   aria-selected on tab switch once the attribute exists, so no extra JS
   was needed.

4. s3_buckets.templ: neither the bucket-policy Save nor Delete handler
   guarded against a double-click, or against firing while the other was
   still in flight - two overlapping PUT/DELETE requests for the same
   bucket could land in either order. Add a shared
   bucketPolicyMutationInFlight flag: set (and both buttons disabled)
   before each fetch, cleared (and buttons re-enabled) on failure so the
   user can retry, left set through the existing success hide-and-reload
   path, and also reset when a new bucket's dialog opens so an abandoned
   in-flight request from a closed dialog can't leave the buttons stuck
   disabled. Verified with jsdom: double-clicking Save now sends exactly
   one PUT, and a Delete click while that PUT is still pending sends no
   DELETE.

* admin: scope bucket-policy mutation completions to the bucket that started them

1. The previous round's fix reset bucketPolicyMutationInFlight whenever a
   new bucket's policy dialog opened, to avoid leaving Save/Delete stuck
   disabled if the modal was closed mid-request. That traded one bug for
   a worse one: if bucket A's PUT/DELETE was still in flight when the
   user opened bucket B's dialog, the reset let B's Save/Delete fire
   immediately, and A's completion handler - unaware anything had
   changed - would still hide the (now B's) modal and reload the page
   out from under whatever the user was doing with B, on success, or
   alert a message with no bucket context, on failure.

   Stop resetting on reopen, so a pending mutation for a previous bucket
   keeps this bucket's Save/Delete blocked until it settles (matches the
   "preventing overlapping mutations" the review comment describes).
   Instead, capture policyEditorBucket as targetBucket right before each
   fetch and compare it against policyEditorBucket again in the
   completion handler: the in-flight flag is always released so the
   buttons never get stuck, but the modal-hide/reload/alert only fire if
   this bucket is still the one showing; a stale completion for an
   abandoned bucket just logs to the console instead.

   Verified with a jsdom simulation: opening bucket B while bucket A's
   Save is still pending leaves B's Save button disabled and a click on
   it a no-op; once A's PUT resolves, B's button re-enables but no
   modal.hide()/reload() fires (previously both fired unconditionally).

2. bucketPolicyDeleteBtn had no bucketPolicyLoaded check, unlike Save -
   a failed GET blocked Save but left Delete free to remove a policy the
   client never actually saw (the same gap already fixed for the S3
   Tables policy modals in an earlier round). Added the same guard,
   ahead of the confirm() dialog. Verified with jsdom: Delete after a
   failed load now sends no DELETE request.

* admin: fix spelling mistake
2026-08-23 22:11:18 -07:00

2289 lines
74 KiB
Go

package dash
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"sort"
"strings"
"sync"
"time"
"github.com/seaweedfs/seaweedfs/weed/admin/maintenance"
adminplugin "github.com/seaweedfs/seaweedfs/weed/admin/plugin"
"github.com/seaweedfs/seaweedfs/weed/cluster"
clustermaintenance "github.com/seaweedfs/seaweedfs/weed/cluster/maintenance"
"github.com/seaweedfs/seaweedfs/weed/credential"
"github.com/seaweedfs/seaweedfs/weed/filer"
"github.com/seaweedfs/seaweedfs/weed/glog"
"github.com/seaweedfs/seaweedfs/weed/pb"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
"github.com/seaweedfs/seaweedfs/weed/pb/mq_pb"
"github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb"
"github.com/seaweedfs/seaweedfs/weed/pb/schema_pb"
"github.com/seaweedfs/seaweedfs/weed/security"
stats_collect "github.com/seaweedfs/seaweedfs/weed/stats"
"github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding"
"github.com/seaweedfs/seaweedfs/weed/storage/super_block"
"github.com/seaweedfs/seaweedfs/weed/util"
"github.com/seaweedfs/seaweedfs/weed/wdclient"
"google.golang.org/grpc"
"github.com/seaweedfs/seaweedfs/weed/s3api"
"github.com/seaweedfs/seaweedfs/weed/s3api/lifecycle_xml"
"github.com/seaweedfs/seaweedfs/weed/s3api/policy_engine"
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
"github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle"
"github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle/scheduler"
"github.com/seaweedfs/seaweedfs/weed/s3api/s3tables"
"github.com/seaweedfs/seaweedfs/weed/worker/tasks"
_ "github.com/seaweedfs/seaweedfs/weed/credential/grpc" // Register gRPC credential store
)
const (
defaultCacheTimeout = 10 * time.Second
defaultFilerCacheTimeout = 30 * time.Second
defaultStatsCacheTimeout = 30 * time.Second
)
// FilerConfig holds filer configuration needed for bucket operations
type FilerConfig struct {
BucketsPath string
FilerGroup string
}
// getFilerConfig retrieves the filer configuration (buckets path and filer group)
func (s *AdminServer) getFilerConfig() (*FilerConfig, error) {
config := &FilerConfig{
BucketsPath: s3_constants.DefaultBucketsPath,
FilerGroup: "",
}
err := s.WithFilerClient(func(client filer_pb.SeaweedFilerClient) error {
resp, err := client.GetFilerConfiguration(context.Background(), &filer_pb.GetFilerConfigurationRequest{})
if err != nil {
return fmt.Errorf("get filer configuration: %w", err)
}
if resp.DirBuckets != "" {
config.BucketsPath = resp.DirBuckets
}
config.FilerGroup = resp.FilerGroup
return nil
})
return config, err
}
// getFilerConf reads filer.conf so callers can inspect per-path rules such as
// the read-only flag that quota enforcement toggles. A missing filer.conf
// yields an empty (all-writable) config rather than an error.
func (s *AdminServer) getFilerConf() (*filer.FilerConf, error) {
fc := filer.NewFilerConf()
err := s.WithFilerClient(func(client filer_pb.SeaweedFilerClient) error {
content, err := filer.ReadInsideFiler(context.Background(), client, filer.DirectoryEtcSeaweedFS, filer.FilerConfName)
if err != nil {
if errors.Is(err, filer_pb.ErrNotFound) {
return nil
}
return err
}
if len(content) > 0 {
return fc.LoadFromBytes(content)
}
return nil
})
return fc, err
}
// getCollectionName returns the collection name for a bucket, prefixed with filer group if configured
func getCollectionName(filerGroup, bucketName string) string {
if filerGroup != "" {
return fmt.Sprintf("%s_%s", filerGroup, bucketName)
}
return bucketName
}
type AdminServer struct {
masterClient *wdclient.MasterClient
templateFS http.FileSystem
dataDir string
filerGroup string
grpcDialOption grpc.DialOption
cacheExpiration time.Duration
lastCacheUpdate time.Time
cachedTopology *ClusterTopology
// dashSamples is a bounded in-memory ring of recent cluster snapshots that
// powers the dashboard's at-a-glance sparklines (see dashboard_metrics.go),
// filled on the maintenance-metrics ticker. No Prometheus dependency.
dashSamples []dashSample
dashSamplesMu sync.Mutex
// Filer discovery and caching
cachedFilers []string
lastFilerUpdate time.Time
filerCacheExpiration time.Duration
// Credential management
credentialManager *credential.CredentialManager
// Configuration persistence
configPersistence *ConfigPersistence
// Maintenance system
maintenanceManager *maintenance.MaintenanceManager
plugin *adminplugin.Plugin
pluginLock *AdminLockManager
adminPresenceLock *adminPresenceLock
expireJobHandler func(jobID string, reason string) (*adminplugin.TrackedJob, bool, error)
// Topic retention purger
topicRetentionPurger *TopicRetentionPurger
// Worker gRPC server
workerGrpcServer *WorkerGrpcServer
// Background goroutine lifecycle
bgCancel context.CancelFunc
// Collection statistics caching
collectionStatsCache map[string]collectionStats
lastCollectionStatsUpdate time.Time
collectionStatsCacheThreshold time.Duration
s3TablesManager *s3tables.Manager
icebergPort int
lancePort int
}
// Type definitions moved to types.go
func NewAdminServer(masters string, filerGroup string, templateFS http.FileSystem, dataDir string, icebergPort, lancePort int) *AdminServer {
grpcDialOption := security.LoadClientTLS(util.GetViper(), "grpc.admin")
// Create master client with multiple master support
masterClient := wdclient.NewMasterClient(
grpcDialOption,
filerGroup,
"admin", // clientType
"", // clientHost - not needed for admin
"", // dataCenter - not needed for admin
"", // rack - not needed for admin
*pb.ServerAddresses(masters).ToServiceDiscovery(),
)
// Start master client connection process (like shell and filer do)
bgCtx, bgCancel := context.WithCancel(context.Background())
go masterClient.KeepConnectedToMaster(bgCtx)
lockManager := NewAdminLockManager(masterClient, adminLockClientName)
presenceLock := newAdminPresenceLock(masterClient)
if presenceLock != nil {
presenceLock.Start()
}
server := &AdminServer{
masterClient: masterClient,
templateFS: templateFS,
dataDir: dataDir,
filerGroup: filerGroup,
grpcDialOption: grpcDialOption,
cacheExpiration: defaultCacheTimeout,
filerCacheExpiration: defaultFilerCacheTimeout,
configPersistence: NewConfigPersistence(dataDir),
collectionStatsCacheThreshold: defaultStatsCacheTimeout,
s3TablesManager: newS3TablesManager(),
icebergPort: icebergPort,
lancePort: lancePort,
pluginLock: lockManager,
adminPresenceLock: presenceLock,
bgCancel: bgCancel,
}
// Initialize topic retention purger
server.topicRetentionPurger = NewTopicRetentionPurger(server)
// Initialize credential manager with defaults
credentialManager, err := credential.NewCredentialManagerWithDefaults(credential.StoreTypeGrpc)
if err != nil {
glog.Warningf("Failed to initialize credential manager: %v", err)
// Continue without credential manager - will fall back to legacy approach
} else {
server.credentialManager = credentialManager
glog.V(0).Infof("Credential manager initialized with store type: %s", credentialManager.GetStore().GetName())
// For stores that need filer address function, configure them
if store := credentialManager.GetStore(); store != nil {
if filerFuncSetter, ok := store.(interface {
SetFilerAddressFunc(func() pb.ServerAddress, grpc.DialOption)
}); ok {
// Configure the filer address function to dynamically return the current active filer
// This function will be called each time credentials need to be loaded/saved,
// so it will automatically use whatever filer is currently available (HA-aware)
filerFuncSetter.SetFilerAddressFunc(func() pb.ServerAddress {
return pb.ServerAddress(server.GetFilerAddress())
}, server.grpcDialOption)
glog.V(0).Infof("Credential store configured with dynamic filer address function")
} else {
glog.V(0).Infof("Credential store %s does not support filer address function", store.GetName())
}
// Mirror the filer's jwt.filer_signing.key so the admin UI's
// Users/Groups pages can present a valid Bearer token when the
// filer enforces IAM gRPC auth. When the key is empty, both
// sides run unauthenticated and no token is sent.
if signer, ok := store.(interface {
SetAdminSigning(security.SigningKey, int)
}); ok {
viper := util.GetViper()
key := security.SigningKey(viper.GetString("jwt.filer_signing.key"))
expires := viper.GetInt("jwt.filer_signing.expires_after_seconds")
signer.SetAdminSigning(key, expires)
if len(key) > 0 {
glog.V(0).Infof("Credential store configured with admin Bearer token signing")
}
}
}
}
// Initialize maintenance system - always initialize even without persistent storage
var maintenanceConfig *maintenance.MaintenanceConfig
if server.configPersistence.IsConfigured() {
var err error
maintenanceConfig, err = server.configPersistence.LoadMaintenanceConfig()
if err != nil {
glog.Errorf("Failed to load maintenance configuration: %v", err)
maintenanceConfig = maintenance.DefaultMaintenanceConfig()
}
// Apply new defaults to handle schema changes (like enabling by default)
schema := maintenance.GetMaintenanceConfigSchema()
if err := schema.ApplyDefaultsToProtobuf(maintenanceConfig); err != nil {
glog.Warningf("Failed to apply schema defaults to loaded config: %v", err)
}
// Force enable maintenance system for new default behavior
// This handles the case where old configs had Enabled=false as default
if !maintenanceConfig.Enabled {
glog.V(1).Infof("Enabling maintenance system (new default behavior)")
maintenanceConfig.Enabled = true
}
glog.V(1).Infof("Maintenance system initialized with persistent configuration (enabled: %v)", maintenanceConfig.Enabled)
} else {
maintenanceConfig = maintenance.DefaultMaintenanceConfig()
glog.V(1).Infof("No data directory configured, maintenance system will run in memory-only mode (enabled: %v)", maintenanceConfig.Enabled)
}
// Load saved task configurations from persistence. This has to run before the maintenance
// manager is created: creating it applies the maintenance policy to the registered
// detectors and schedulers, while this call replaces each task's whole config object, so
// running it afterwards would discard what the policy just applied. Both read the same
// persisted task config files, so the policy ends up as the last writer and stays
// authoritative for the task types it covers.
server.loadTaskConfigurationsFromPersistence()
// Always initialize maintenance manager
server.InitMaintenanceManager(maintenanceConfig)
// Start maintenance manager if enabled
if maintenanceConfig.Enabled {
go func() {
// Give master client a bit of time to connect before starting scans
time.Sleep(2 * time.Second)
if err := server.StartMaintenanceManager(); err != nil {
glog.Errorf("Failed to start maintenance manager: %v", err)
}
}()
}
pluginOpts := adminplugin.Options{
DataDir: dataDir,
ClusterContextProvider: func(_ context.Context) (*plugin_pb.ClusterContext, error) {
return server.buildDefaultPluginClusterContext(), nil
},
LockManager: lockManager,
ConfigDefaultsProvider: server.enrichConfigDefaults,
}
plugin, err := adminplugin.New(pluginOpts)
if err != nil && dataDir != "" {
glog.Warningf("Failed to initialize plugin with dataDir=%q: %v. Falling back to in-memory plugin state.", dataDir, err)
pluginOpts.DataDir = ""
plugin, err = adminplugin.New(pluginOpts)
}
if err != nil {
glog.Errorf("Failed to initialize plugin: %v", err)
} else {
server.plugin = plugin
glog.V(0).Infof("Plugin enabled")
go server.monitorVacuumWorker(bgCtx)
}
go server.publishMaintenanceMetrics(bgCtx)
return server
}
func (s *AdminServer) listClusterNodesRequest(clientType string) *master_pb.ListClusterNodesRequest {
return &master_pb.ListClusterNodesRequest{
ClientType: clientType,
FilerGroup: s.filerGroup,
}
}
// vacuumToggler abstracts the master's vacuum enable/disable for testing.
type vacuumToggler interface {
disableVacuum() error
enableVacuum() error
}
// masterVacuumToggler implements vacuumToggler via gRPC calls to the master.
type masterVacuumToggler struct {
server *AdminServer
}
func (m *masterVacuumToggler) disableVacuum() error {
return m.server.WithMasterClient(func(client master_pb.SeaweedClient) error {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
_, err := client.DisableVacuum(ctx, &master_pb.DisableVacuumRequest{ByPlugin: true})
return err
})
}
func (m *masterVacuumToggler) enableVacuum() error {
return m.server.WithMasterClient(func(client master_pb.SeaweedClient) error {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
_, err := client.EnableVacuum(ctx, &master_pb.EnableVacuumRequest{ByPlugin: true})
return err
})
}
// syncVacuumState performs a single sync step: checks if a vacuum-capable worker
// is present and calls disable/enable accordingly. Returns the updated state
// and whether the call failed (for log dedup on retries).
func syncVacuumState(hasWorker bool, previouslyActive bool, toggler vacuumToggler, retrying bool) (active bool, failed bool) {
if hasWorker == previouslyActive {
return previouslyActive, false
}
if hasWorker {
if !retrying {
glog.V(0).Infof("Vacuum plugin worker connected, disabling master automatic vacuum")
}
if err := toggler.disableVacuum(); err != nil {
glog.Warningf("Failed to disable vacuum on master: %v", err)
return false, true // retry next tick
}
return true, false
}
if !retrying {
glog.V(0).Infof("Vacuum plugin worker disconnected, re-enabling master automatic vacuum")
}
if err := toggler.enableVacuum(); err != nil {
glog.Warningf("Failed to enable vacuum on master: %v", err)
return true, true // retry next tick
}
return false, false
}
// monitorVacuumWorker polls the plugin registry for vacuum-capable workers and
// disables/enables the master's automatic scheduled vacuum accordingly.
func (s *AdminServer) monitorVacuumWorker(ctx context.Context) {
const pollInterval = 30 * time.Second
ticker := time.NewTicker(pollInterval)
defer ticker.Stop()
toggler := &masterVacuumToggler{server: s}
vacuumWorkerActive := false
retrying := false
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
if s.plugin == nil {
continue
}
hasWorker := s.plugin.HasCapableWorker("vacuum")
vacuumWorkerActive, retrying = syncVacuumState(hasWorker, vacuumWorkerActive, toggler, retrying)
}
}
}
// publishMaintenanceMetrics periodically snapshots the maintenance queue and
// worker fleet into Prometheus gauges. Counters and durations are recorded at
// their event sites; these gauges reflect current state at scrape resolution.
func (s *AdminServer) publishMaintenanceMetrics(ctx context.Context) {
const interval = 15 * time.Second
ticker := time.NewTicker(interval)
defer ticker.Stop()
// Seed one sample so the dashboard has something to draw before the first tick.
s.recordDashboardSample()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
s.collectMaintenanceMetrics()
s.recordDashboardSample()
}
}
}
// workerFleetTotals aggregates connected workers and their task slots across
// BOTH worker registries the admin server keeps: the legacy maintenance-worker
// registry (workers that register over the worker gRPC stream) and the plugin
// worker registry (workers started as `weed worker`). Reading only the legacy
// one reported zero workers on clusters that run the admin and the workers as
// separate components, where no legacy worker ever registers.
//
// A worker can appear in both registries: `weed mini` starts both runtimes from
// one working directory, so they share the persisted worker ID. Merging by ID
// keeps such a worker counted once, and its slots are taken from the legacy
// registry, which is where they were accounted for before.
func (s *AdminServer) workerFleetTotals() (workers, usedSlots, maxSlots int) {
var legacySlots map[string]maintenance.WorkerSlots
if s.maintenanceManager != nil {
legacySlots = s.maintenanceManager.GetWorkerSlots()
}
return mergeWorkerFleetTotals(legacySlots, s.GetPluginWorkers())
}
// mergeWorkerFleetTotals unions the legacy and plugin worker registries by
// worker ID. Plugin workers report their slots in the heartbeat, so one that
// has connected but not yet sent a heartbeat adds to the worker count with zero
// slots until its first heartbeat lands.
func mergeWorkerFleetTotals(legacySlots map[string]maintenance.WorkerSlots, pluginWorkers []*adminplugin.WorkerSession) (workers, usedSlots, maxSlots int) {
for _, slots := range legacySlots {
workers++
usedSlots += slots.Used
maxSlots += slots.Max
}
for _, session := range pluginWorkers {
if session == nil {
continue
}
if _, counted := legacySlots[session.WorkerID]; counted {
continue
}
workers++
if heartbeat := session.Heartbeat; heartbeat != nil {
used := int(heartbeat.DetectionSlotsUsed) + int(heartbeat.ExecutionSlotsUsed)
max := int(heartbeat.DetectionSlotsTotal) + int(heartbeat.ExecutionSlotsTotal)
// A worker's self-reported slots are untrusted input; a stale or
// misbehaving one should not be able to drive the aggregate gauge
// negative, matching the same defensiveness as registry.go's own
// slot arithmetic.
if used < 0 {
used = 0
}
if max < 0 {
max = 0
}
usedSlots += used
maxSlots += max
}
}
return
}
func (s *AdminServer) collectMaintenanceMetrics() {
// Published before the maintenanceManager guard below: plugin workers are
// tracked independently of the maintenance manager.
workers, usedSlots, maxSlots := s.workerFleetTotals()
stats_collect.AdminWorkersConnected.Set(float64(workers))
stats_collect.AdminWorkerSlots.WithLabelValues("used").Set(float64(usedSlots))
stats_collect.AdminWorkerSlots.WithLabelValues("max").Set(float64(maxSlots))
if s.maintenanceManager == nil {
return
}
stats := s.maintenanceManager.GetStats()
stats_collect.AdminMaintenanceTasksByStatus.Reset()
for status, count := range stats.TasksByStatus {
stats_collect.AdminMaintenanceTasksByStatus.WithLabelValues(string(status)).Set(float64(count))
}
stats_collect.AdminMaintenanceTasksByType.Reset()
for taskType, count := range stats.TasksByType {
stats_collect.AdminMaintenanceTasksByType.WithLabelValues(string(taskType)).Set(float64(count))
}
// NextScanTime is only meaningful while the scanner runs; GetStats computes
// it unconditionally, so clear the gauge when idle to avoid a stale value.
if s.maintenanceManager.IsRunning() && !stats.NextScanTime.IsZero() {
stats_collect.AdminMaintenanceNextScanTimestampSeconds.Set(float64(stats.NextScanTime.Unix()))
} else {
stats_collect.AdminMaintenanceNextScanTimestampSeconds.Set(0)
}
}
// loadTaskConfigurationsFromPersistence loads saved task configurations from protobuf files
func (s *AdminServer) loadTaskConfigurationsFromPersistence() {
if s.configPersistence == nil || !s.configPersistence.IsConfigured() {
glog.V(1).Infof("Config persistence not available, using default task configurations")
return
}
// Load task configurations dynamically using the config update registry
configUpdateRegistry := tasks.GetGlobalConfigUpdateRegistry()
configUpdateRegistry.UpdateAllConfigs(s.configPersistence)
}
// enrichConfigDefaults is called by the plugin when bootstrapping a job type's
// default config from its descriptor. It overlays admin.toml maintenance
// settings, and for admin_script fetches maintenance scripts from the master
// to use as the script default.
//
// MIGRATION: the admin_script part exists to help users migrate from
// master.toml [master.maintenance] to the admin script plugin worker.
// Remove after March 2027.
func (s *AdminServer) enrichConfigDefaults(cfg *plugin_pb.PersistedJobTypeConfig) *plugin_pb.PersistedJobTypeConfig {
applyPluginTomlDefaults(util.GetViper(), cfg)
if cfg.JobType != "admin_script" {
return cfg
}
var maintenanceScripts string
var sleepMinutes uint32
err := s.WithMasterClient(func(client master_pb.SeaweedClient) error {
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
resp, err := client.GetMasterConfiguration(ctx, &master_pb.GetMasterConfigurationRequest{})
if err != nil {
return err
}
maintenanceScripts = resp.MaintenanceScripts
sleepMinutes = resp.MaintenanceSleepMinutes
return nil
})
if err != nil {
glog.V(1).Infof("Could not fetch master configuration for admin_script defaults: %v", err)
return cfg
}
script := cleanMaintenanceScript(maintenanceScripts)
if script == "" {
return cfg
}
interval := int64(sleepMinutes)
if interval <= 0 {
interval = clustermaintenance.DefaultMaintenanceSleepMinutes
}
glog.V(0).Infof("Enriching admin_script defaults from master maintenance scripts (interval=%dm)", interval)
if cfg.AdminConfigValues == nil {
cfg.AdminConfigValues = make(map[string]*plugin_pb.ConfigValue)
}
cfg.AdminConfigValues["script"] = &plugin_pb.ConfigValue{
Kind: &plugin_pb.ConfigValue_StringValue{StringValue: script},
}
cfg.AdminConfigValues["run_interval_minutes"] = &plugin_pb.ConfigValue{
Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: interval},
}
cfg.UpdatedBy = "master_migration"
return cfg
}
// cleanMaintenanceScript strips lock/unlock commands and normalizes a
// maintenance script string for use with the admin script plugin worker.
//
// MIGRATION: Used by enrichConfigDefaults. Remove after March 2027.
func cleanMaintenanceScript(script string) string {
script = strings.ReplaceAll(script, "\r\n", "\n")
var lines []string
for _, line := range strings.Split(script, "\n") {
trimmed := strings.TrimSpace(line)
if trimmed == "" || strings.HasPrefix(trimmed, "#") {
continue
}
// Strip inline comments (e.g., "lock # migration note")
if idx := strings.Index(trimmed, "#"); idx >= 0 {
trimmed = strings.TrimSpace(trimmed[:idx])
if trimmed == "" {
continue
}
}
firstToken := strings.ToLower(strings.Fields(trimmed)[0])
if firstToken == "lock" || firstToken == "unlock" {
continue
}
lines = append(lines, trimmed)
}
return strings.Join(lines, "\n")
}
// GetCredentialManager returns the credential manager
func (s *AdminServer) GetCredentialManager() *credential.CredentialManager {
return s.credentialManager
}
// Filer discovery methods moved to client_management.go
// Client management methods moved to client_management.go
// WithFilerClient and WithVolumeServerClient methods moved to client_management.go
// Cluster topology methods moved to cluster_topology.go
// getTopologyViaGRPC method moved to cluster_topology.go
// InvalidateCache method moved to cluster_topology.go
// GetS3BucketsData retrieves Object Store buckets with pagination and sorting
func (s *AdminServer) GetS3BucketsData(page, pageSize int, sortBy, sortOrder string) (S3BucketsData, error) {
if page < 1 {
page = 1
}
if pageSize < 1 || pageSize > 1000 {
pageSize = 100
}
if sortBy == "" {
sortBy = "name"
}
if sortOrder == "" {
sortOrder = "asc"
}
buckets, err := s.GetS3Buckets()
if err != nil {
return S3BucketsData{}, err
}
var totalSize int64
for _, bucket := range buckets {
totalSize += bucket.PhysicalSize
}
totalBuckets := len(buckets)
// Sort buckets
s.sortBuckets(buckets, sortBy, sortOrder)
// Calculate pagination
totalPages := (totalBuckets + pageSize - 1) / pageSize
if totalPages == 0 {
totalPages = 1
}
if page > totalPages {
page = totalPages
}
startIndex := (page - 1) * pageSize
endIndex := startIndex + pageSize
if startIndex >= totalBuckets {
buckets = []S3Bucket{}
} else {
if endIndex > totalBuckets {
endIndex = totalBuckets
}
buckets = buckets[startIndex:endIndex]
}
return S3BucketsData{
Buckets: buckets,
TotalBuckets: totalBuckets,
TotalSize: totalSize,
LastUpdated: time.Now(),
CurrentPage: page,
TotalPages: totalPages,
PageSize: pageSize,
SortBy: sortBy,
SortOrder: sortOrder,
}, nil
}
// sortBuckets sorts the bucket slice in place by the given field and order
func (s *AdminServer) sortBuckets(buckets []S3Bucket, sortBy, sortOrder string) {
desc := sortOrder == "desc"
sort.Slice(buckets, func(i, j int) bool {
a, b := buckets[i], buckets[j]
switch sortBy {
case "owner":
if a.Owner != b.Owner {
if desc {
return a.Owner > b.Owner
}
return a.Owner < b.Owner
}
case "created":
if !a.CreatedAt.Equal(b.CreatedAt) {
if desc {
return a.CreatedAt.After(b.CreatedAt)
}
return a.CreatedAt.Before(b.CreatedAt)
}
case "logical_size":
if a.LogicalSize != b.LogicalSize {
if desc {
return a.LogicalSize > b.LogicalSize
}
return a.LogicalSize < b.LogicalSize
}
case "physical_size":
if a.PhysicalSize != b.PhysicalSize {
if desc {
return a.PhysicalSize > b.PhysicalSize
}
return a.PhysicalSize < b.PhysicalSize
}
}
// Tie-breaker: sort by name (also the default/primary for sortBy=="name")
if a.Name != b.Name {
if desc {
return a.Name > b.Name
}
return a.Name < b.Name
}
return false
})
}
// GetS3Buckets retrieves all Object Store buckets from the filer and collects size/object data from collections
func (s *AdminServer) GetS3Buckets() ([]S3Bucket, error) {
var buckets []S3Bucket
// Collect volume information by collection with caching
collectionMap, _ := s.getCollectionStats()
// Get filer configuration (buckets path and filer group)
filerConfig, err := s.getFilerConfig()
if err != nil {
glog.Warningf("Failed to get filer configuration, using defaults: %v", err)
}
// Read filer.conf so we can surface the read-only flag quota enforcement sets
fc, err := s.getFilerConf()
if err != nil {
glog.Warningf("Failed to read filer.conf for bucket read-only state: %v", err)
}
// Now list buckets from the filer and match with collection data
err = s.WithFilerClient(func(client filer_pb.SeaweedFilerClient) error {
// Paginate through all buckets in the buckets directory
const listPageSize = 1000
startFrom := ""
var snapshotTsNs int64
for {
stream, err := client.ListEntries(context.Background(), &filer_pb.ListEntriesRequest{
Directory: filerConfig.BucketsPath,
Prefix: "",
StartFromFileName: startFrom,
InclusiveStartFrom: false,
Limit: listPageSize,
SnapshotTsNs: snapshotTsNs,
})
if err != nil {
return err
}
pageCount := 0
lastName := ""
for {
resp, err := stream.Recv()
if err != nil {
if errors.Is(err, io.EOF) {
break
}
return err
}
if snapshotTsNs == 0 && resp.SnapshotTsNs != 0 {
snapshotTsNs = resp.SnapshotTsNs
}
if resp.Entry == nil {
continue
}
lastName = resp.Entry.Name
pageCount++
if !resp.Entry.IsDirectory {
continue
}
bucketName := resp.Entry.Name
if strings.HasPrefix(bucketName, ".") {
// Skip internal/system directories from Object Store bucket listing.
continue
}
if s3tables.IsTableBucketEntry(resp.Entry) || strings.HasSuffix(bucketName, "--table-s3") {
// Keep table buckets in the S3 Tables pages, not regular Object Store buckets.
continue
}
// Determine collection name for this bucket
collectionName := getCollectionName(filerConfig.FilerGroup, bucketName)
var physicalSize int64
var logicalSize int64
if collectionData, exists := collectionMap[collectionName]; exists {
physicalSize = collectionData.PhysicalSize
logicalSize = collectionData.LogicalSize
}
// Get quota information from entry
quota := resp.Entry.Quota
quotaEnabled := quota > 0
if quota < 0 {
// Negative quota means disabled
quota = -quota
quotaEnabled = false
}
// Get versioning, object lock, and owner information from extended attributes
versioningStatus := ""
objectLockEnabled := false
objectLockMode := ""
var objectLockDuration int32 = 0
var owner string
if resp.Entry.Extended != nil {
// Use shared utility to extract versioning information
versioningStatus = extractVersioningFromEntry(resp.Entry)
// Use shared utility to extract Object Lock information
objectLockEnabled, objectLockMode, objectLockDuration = extractObjectLockInfoFromEntry(resp.Entry)
// Extract owner information
if ownerBytes, ok := resp.Entry.Extended[s3_constants.AmzIdentityId]; ok {
owner = string(ownerBytes)
}
}
var createdAt, lastModified time.Time
if resp.Entry.Attributes != nil {
createdAt = time.Unix(resp.Entry.Attributes.Crtime, 0)
lastModified = time.Unix(resp.Entry.Attributes.Mtime, 0)
}
readOnly := fc.MatchStorageRule(filerConfig.BucketsPath + "/" + bucketName + "/").ReadOnly
lifecycleRuleCount, lifecycleEnabledCount := extractLifecycleCountsFromEntry(resp.Entry)
bucket := S3Bucket{
Name: bucketName,
CreatedAt: createdAt,
LogicalSize: logicalSize,
PhysicalSize: physicalSize,
LastModified: lastModified,
Quota: quota,
QuotaEnabled: quotaEnabled,
ReadOnly: readOnly,
VersioningStatus: versioningStatus,
ObjectLockEnabled: objectLockEnabled,
ObjectLockMode: objectLockMode,
ObjectLockDuration: objectLockDuration,
Owner: owner,
LifecycleRuleCount: lifecycleRuleCount,
LifecycleEnabledCount: lifecycleEnabledCount,
PolicyStatementCount: extractPolicyStatementCountFromEntry(resp.Entry),
}
buckets = append(buckets, bucket)
}
// If we received fewer entries than the page size, we've listed everything
if pageCount < listPageSize {
break
}
startFrom = lastName
}
return nil
})
if err != nil {
return nil, fmt.Errorf("failed to list Object Store buckets: %w", err)
}
return buckets, nil
}
// GetBucketDetails retrieves detailed information about a specific bucket
// Note: This no longer lists objects for performance reasons. Use GetS3Buckets for size/count data.
func (s *AdminServer) GetBucketDetails(bucketName string) (*BucketDetails, error) {
// Get filer configuration (buckets path)
filerConfig, err := s.getFilerConfig()
if err != nil {
glog.Warningf("Failed to get filer configuration, using defaults: %v", err)
}
details := &BucketDetails{
Bucket: S3Bucket{
Name: bucketName,
},
UpdatedAt: time.Now(),
}
// Get collection data for size and object count with caching
collectionName := getCollectionName(filerConfig.FilerGroup, bucketName)
stats, err := s.getCollectionStats()
if err != nil {
glog.Warningf("Failed to get collection data: %v", err)
// Continue without collection data - use zero values
} else if data, ok := stats[collectionName]; ok {
details.Bucket.LogicalSize = data.LogicalSize
details.Bucket.PhysicalSize = data.PhysicalSize
}
// Surface the read-only flag quota enforcement sets in filer.conf
if fc, err := s.getFilerConf(); err != nil {
glog.Warningf("Failed to read filer.conf for bucket read-only state: %v", err)
} else {
details.Bucket.ReadOnly = fc.MatchStorageRule(filerConfig.BucketsPath + "/" + bucketName + "/").ReadOnly
}
err = s.WithFilerClient(func(client filer_pb.SeaweedFilerClient) error {
// Get bucket info
bucketResp, err := client.LookupDirectoryEntry(context.Background(), &filer_pb.LookupDirectoryEntryRequest{
Directory: filerConfig.BucketsPath,
Name: bucketName,
})
if err != nil {
return fmt.Errorf("bucket not found: %w", err)
}
details.Bucket.CreatedAt = time.Unix(bucketResp.Entry.Attributes.Crtime, 0)
details.Bucket.LastModified = time.Unix(bucketResp.Entry.Attributes.Mtime, 0)
// Get quota information from entry
quota := bucketResp.Entry.Quota
quotaEnabled := quota > 0
if quota < 0 {
// Negative quota means disabled
quota = -quota
quotaEnabled = false
}
details.Bucket.Quota = quota
details.Bucket.QuotaEnabled = quotaEnabled
// Get versioning, object lock, and owner information from extended attributes
versioningStatus := ""
objectLockEnabled := false
objectLockMode := ""
var objectLockDuration int32 = 0
var owner string
if bucketResp.Entry.Extended != nil {
// Use shared utility to extract versioning information
versioningStatus = extractVersioningFromEntry(bucketResp.Entry)
// Use shared utility to extract Object Lock information
objectLockEnabled, objectLockMode, objectLockDuration = extractObjectLockInfoFromEntry(bucketResp.Entry)
// Extract owner information
if ownerBytes, ok := bucketResp.Entry.Extended[s3_constants.AmzIdentityId]; ok {
owner = string(ownerBytes)
}
}
details.Bucket.VersioningStatus = versioningStatus
details.Bucket.ObjectLockEnabled = objectLockEnabled
details.Bucket.ObjectLockMode = objectLockMode
details.Bucket.ObjectLockDuration = objectLockDuration
details.Bucket.Owner = owner
details.Bucket.LifecycleRuleCount, details.Bucket.LifecycleEnabledCount = extractLifecycleCountsFromEntry(bucketResp.Entry)
details.Bucket.PolicyStatementCount = extractPolicyStatementCountFromEntry(bucketResp.Entry)
return nil
})
if err != nil {
return nil, err
}
return details, nil
}
// GetBucketLifecycle returns the lifecycle configuration stored on a bucket's filer entry
func (s *AdminServer) GetBucketLifecycle(bucketName string) (*BucketLifecycle, error) {
filerConfig, err := s.getFilerConfig()
if err != nil {
glog.Warningf("Failed to get filer configuration, using defaults: %v", err)
}
lifecycle := &BucketLifecycle{
Bucket: bucketName,
Rules: []BucketLifecycleRule{},
}
err = s.WithFilerClient(func(client filer_pb.SeaweedFilerClient) error {
resp, err := client.LookupDirectoryEntry(context.Background(), &filer_pb.LookupDirectoryEntryRequest{
Directory: filerConfig.BucketsPath,
Name: bucketName,
})
if err != nil {
return fmt.Errorf("bucket not found: %w", err)
}
xmlBytes := resp.Entry.Extended[scheduler.BucketLifecycleConfigurationXMLKey]
if len(xmlBytes) == 0 {
return nil
}
rules, err := lifecycle_xml.ParseCanonical(xmlBytes)
if err != nil {
return fmt.Errorf("parse lifecycle configuration: %w", err)
}
lifecycle.XML = string(xmlBytes)
for _, rule := range rules {
lifecycle.Rules = append(lifecycle.Rules, toBucketLifecycleRule(rule))
}
return nil
})
if err != nil {
return nil, err
}
return lifecycle, nil
}
func toBucketLifecycleRule(rule *s3lifecycle.Rule) BucketLifecycleRule {
out := BucketLifecycleRule{
ID: rule.ID,
Status: rule.Status,
Prefix: rule.Prefix,
Tags: rule.FilterTags,
SizeGreaterThan: rule.FilterSizeGreaterThan,
SizeLessThan: rule.FilterSizeLessThan,
ExpirationDays: rule.ExpirationDays,
ExpiredObjectDeleteMarker: rule.ExpiredObjectDeleteMarker,
NoncurrentVersionExpirationDays: rule.NoncurrentVersionExpirationDays,
NewerNoncurrentVersions: rule.NewerNoncurrentVersions,
AbortMultipartDays: rule.AbortMPUDaysAfterInitiation,
}
if !rule.ExpirationDate.IsZero() {
out.ExpirationDate = rule.ExpirationDate.Format(time.DateOnly)
}
return out
}
// fromBucketLifecycleRule is the inverse of toBucketLifecycleRule, turning a
// rule edited in the admin UI back into the engine's canonical shape.
func fromBucketLifecycleRule(rule BucketLifecycleRule) (*s3lifecycle.Rule, error) {
out := &s3lifecycle.Rule{
ID: rule.ID,
Status: rule.Status,
Prefix: rule.Prefix,
FilterTags: rule.Tags,
FilterSizeGreaterThan: rule.SizeGreaterThan,
FilterSizeLessThan: rule.SizeLessThan,
ExpirationDays: rule.ExpirationDays,
ExpiredObjectDeleteMarker: rule.ExpiredObjectDeleteMarker,
NoncurrentVersionExpirationDays: rule.NoncurrentVersionExpirationDays,
NewerNoncurrentVersions: rule.NewerNoncurrentVersions,
AbortMPUDaysAfterInitiation: rule.AbortMultipartDays,
}
if rule.ExpirationDate != "" {
date, err := time.Parse(time.DateOnly, rule.ExpirationDate)
if err != nil {
return nil, fmt.Errorf("invalid expiration date %q: %w", rule.ExpirationDate, err)
}
out.ExpirationDate = date
}
return out, nil
}
// ErrBucketNotFound reports that the named bucket has no filer entry, so a
// handler can answer 404 rather than 500.
var ErrBucketNotFound = errors.New("bucket not found")
// SetBucketLifecycle replaces the lifecycle configuration stored on a
// bucket's filer entry. An empty rule list clears the configuration
// entirely, mirroring clearStoredBucketLifecycleConfiguration on the S3 API
// side. Callers must validate rules before calling this (see
// validateBucketLifecycleRules) — this only rejects what marshaling itself
// rejects.
func (s *AdminServer) SetBucketLifecycle(bucketName string, rules []BucketLifecycleRule) error {
canonicalRules := make([]*s3lifecycle.Rule, 0, len(rules))
for _, rule := range rules {
canonicalRule, err := fromBucketLifecycleRule(rule)
if err != nil {
return err
}
canonicalRules = append(canonicalRules, canonicalRule)
}
var lifecycleXML []byte
if len(canonicalRules) > 0 {
var err error
lifecycleXML, err = lifecycle_xml.MarshalCanonical(canonicalRules)
if err != nil {
return fmt.Errorf("marshal lifecycle configuration: %w", err)
}
if len(lifecycleXML) > scheduler.MaxBucketLifecycleConfigurationSize {
return fmt.Errorf("lifecycle configuration is %d bytes, which exceeds the %d byte limit", len(lifecycleXML), scheduler.MaxBucketLifecycleConfigurationSize)
}
}
filerConfig, err := s.getFilerConfig()
if err != nil {
return fmt.Errorf("get filer configuration: %w", err)
}
collection := getCollectionName(filerConfig.FilerGroup, bucketName)
return s.WithFilerClient(func(client filer_pb.SeaweedFilerClient) error {
// PATCH_EXTENDED is a no-op on a missing entry, so the existence
// check has to happen here rather than fall out of the write.
if _, err := filer_pb.LookupEntry(context.Background(), client, &filer_pb.LookupDirectoryEntryRequest{
Directory: filerConfig.BucketsPath,
Name: bucketName,
}); err != nil {
if errors.Is(err, filer_pb.ErrNotFound) {
return fmt.Errorf("%w: %s", ErrBucketNotFound, bucketName)
}
return fmt.Errorf("look up bucket %s: %w", bucketName, err)
}
// Migration: clear any legacy day-TTL filer.conf entries before
// writing the new XML, so a failure here leaves the bucket entry
// untouched instead of committing the new policy alongside a stale
// TTL rule. Same step and ordering as
// PutBucketLifecycleConfigurationHandler.
if err := filer.ClearBucketLifecycleDayTTLs(context.Background(), client, filerConfig.BucketsPath, bucketName, collection); err != nil {
return fmt.Errorf("failed to clear legacy lifecycle TTLs: %w", err)
}
bucketPath := filerConfig.BucketsPath + "/" + bucketName
resp, err := client.ObjectTransaction(context.Background(), &filer_pb.ObjectTransactionRequest{
LockKey: bucketPath,
RouteKey: s3_constants.ObjectWriteRouteKeyPrefix + bucketPath,
Mutations: []*filer_pb.ObjectMutation{bucketLifecycleMutation(filerConfig.BucketsPath, bucketName, lifecycleXML)},
})
if err != nil {
return fmt.Errorf("failed to update bucket lifecycle: %w", err)
}
if resp.Error != "" {
return fmt.Errorf("failed to update bucket lifecycle: %s", resp.Error)
}
return nil
})
}
// bucketLifecycleMutation patches the two lifecycle keys rather than writing
// the whole entry back: the filer re-reads and merges under the bucket path
// lock, so a concurrent owner/quota/versioning change is preserved instead of
// being reverted by a stale snapshot. Same mutation the S3 gateway uses for
// these keys (see patchBucketEntry in s3api_bucket_config.go). Empty XML
// clears the configuration, transition minimum size included.
func bucketLifecycleMutation(bucketsPath, bucketName string, lifecycleXML []byte) *filer_pb.ObjectMutation {
mutation := &filer_pb.ObjectMutation{
Type: filer_pb.ObjectMutation_PATCH_EXTENDED,
Directory: bucketsPath,
Name: bucketName,
}
if len(lifecycleXML) > 0 {
mutation.SetExtended = map[string][]byte{
scheduler.BucketLifecycleConfigurationXMLKey: lifecycleXML,
}
return mutation
}
mutation.DeleteExtended = []string{
scheduler.BucketLifecycleConfigurationXMLKey,
scheduler.BucketLifecycleTransitionMinimumObjectSizeKey,
}
return mutation
}
// CreateS3Bucket creates a new S3 bucket
func (s *AdminServer) CreateS3Bucket(bucketName string) error {
return s.CreateS3BucketWithQuota(bucketName, 0, false)
}
// DeleteS3Bucket deletes an S3 bucket and all its contents
func (s *AdminServer) DeleteS3Bucket(bucketName string) error {
ctx := context.Background()
// Get filer configuration (buckets path and filer group)
filerConfig, err := s.getFilerConfig()
if err != nil {
return fmt.Errorf("failed to get filer configuration: %w", err)
}
// Check if bucket has Object Lock enabled and if there are locked objects
err = s.WithFilerClient(func(client filer_pb.SeaweedFilerClient) error {
return s3api.CheckBucketForLockedObjects(ctx, client, filerConfig.BucketsPath, bucketName)
})
if err != nil {
return err
}
// Delete the collection first (same as s3.bucket.delete shell command)
// This ensures volume data is cleaned up properly
// Collection name must be prefixed with filer group if configured
collectionName := getCollectionName(filerConfig.FilerGroup, bucketName)
err = s.WithMasterClient(func(client master_pb.SeaweedClient) error {
_, err := client.CollectionDelete(ctx, &master_pb.CollectionDeleteRequest{
Name: collectionName,
})
return err
})
if err != nil {
return fmt.Errorf("failed to delete collection %s: %w", collectionName, err)
}
// Then delete bucket directory recursively from filer
// Use same parameters as s3.bucket.delete shell command and S3 API
return s.WithFilerClient(func(client filer_pb.SeaweedFilerClient) error {
_, err := client.DeleteEntry(ctx, &filer_pb.DeleteEntryRequest{
Directory: filerConfig.BucketsPath,
Name: bucketName,
IsDeleteData: false, // Collection already deleted, just remove metadata
IsRecursive: true,
IgnoreRecursiveError: true, // Same as S3 API and shell command
})
if err != nil {
return fmt.Errorf("failed to delete bucket: %w", err)
}
return nil
})
}
// IsStaticUser checks if a user is a static identity by loading the
// configuration from the credential manager and checking the IsStatic flag.
func (s *AdminServer) IsStaticUser(username string) bool {
if s.credentialManager == nil {
return false
}
s3cfg, err := s.credentialManager.LoadConfiguration(context.Background())
if err != nil {
return false
}
for _, ident := range s3cfg.Identities {
if ident.Name == username {
return ident.IsStatic
}
}
return false
}
// GetObjectStoreUsers retrieves object store users from identity.json
func (s *AdminServer) GetObjectStoreUsers(ctx context.Context) ([]ObjectStoreUser, error) {
if s.credentialManager == nil {
return []ObjectStoreUser{}, nil
}
s3cfg, err := s.credentialManager.LoadConfiguration(ctx)
if err != nil {
return nil, fmt.Errorf("failed to load IAM configuration: %w", err)
}
var users []ObjectStoreUser
// Convert IAM identities to ObjectStoreUser format
for _, identity := range s3cfg.Identities {
// Skip service accounts - they should not be parent users
if strings.HasPrefix(identity.Name, serviceAccountPrefix) {
continue
}
user := ObjectStoreUser{
Username: identity.Name,
Permissions: identity.Actions,
IsStatic: identity.IsStatic,
}
// Set email from account if available
if identity.Account != nil {
user.Email = identity.Account.EmailAddress
}
// Get first access key for display
if len(identity.Credentials) > 0 {
user.AccessKey = identity.Credentials[0].AccessKey
user.SecretKey = identity.Credentials[0].SecretKey
}
users = append(users, user)
}
return users, nil
}
// Volume server methods moved to volume_management.go
// Volume methods moved to volume_management.go
// sortVolumes method moved to volume_management.go
// GetClusterCollections method moved to collection_management.go
// GetClusterMasters retrieves cluster masters data
func (s *AdminServer) GetClusterMasters() (*ClusterMastersData, error) {
var masters []MasterInfo
var leaderCount int
// First, get master information from topology
topology, err := s.GetClusterTopology()
if err != nil {
return nil, err
}
// Create a map to merge topology and raft data
masterMap := make(map[string]*MasterInfo)
// Add masters from topology
for _, master := range topology.Masters {
masterInfo := &MasterInfo{
Address: pb.ServerAddress(master.Address).ToHttpAddress(),
IsLeader: master.IsLeader,
Suffrage: "",
}
if master.IsLeader {
leaderCount++
}
masterMap[masterInfo.Address] = masterInfo
}
// Then, get additional master information from Raft cluster
err = s.WithMasterClient(func(client master_pb.SeaweedClient) error {
resp, err := client.RaftListClusterServers(context.Background(), &master_pb.RaftListClusterServersRequest{})
if err != nil {
return err
}
// Process each raft server
for _, server := range resp.ClusterServers {
// Raft stores gRPC addresses, convert to HTTP address
httpAddress := pb.GrpcAddressToServerAddress(server.Address)
// Update existing master info or create new one
if masterInfo, exists := masterMap[httpAddress]; exists {
// Update existing master with raft data
masterInfo.IsLeader = server.IsLeader
masterInfo.Suffrage = server.Suffrage
} else {
// Create new master info from raft data
masterInfo := &MasterInfo{
Address: httpAddress,
IsLeader: server.IsLeader,
Suffrage: server.Suffrage,
}
masterMap[httpAddress] = masterInfo
}
if server.IsLeader {
// Update leader count based on raft data
leaderCount = 1 // There should only be one leader
}
}
return nil
})
if err != nil {
// If gRPC call fails, log the error but continue with topology data
currentMaster := s.masterClient.GetMaster(context.Background())
glog.Errorf("Failed to get raft cluster servers from master %s: %v", currentMaster, err)
}
// Convert map to slice
for _, masterInfo := range masterMap {
masters = append(masters, *masterInfo)
}
// Sort masters by address for consistent ordering on page refresh
sort.Slice(masters, func(i, j int) bool {
return masters[i].Address < masters[j].Address
})
// If no masters found at all, add the current master as fallback
if len(masters) == 0 {
currentMaster := s.masterClient.GetMaster(context.Background())
if currentMaster != "" {
masters = append(masters, MasterInfo{
Address: pb.ServerAddress(currentMaster).ToHttpAddress(),
IsLeader: true,
Suffrage: "Voter",
})
leaderCount = 1
}
}
return &ClusterMastersData{
Masters: masters,
TotalMasters: len(masters),
LeaderCount: leaderCount,
LastUpdated: time.Now(),
}, nil
}
// GetClusterFilers retrieves cluster filers data
func (s *AdminServer) GetClusterFilers() (*ClusterFilersData, error) {
var filers []FilerInfo
// Get filer information from master using ListClusterNodes
err := s.WithMasterClient(func(client master_pb.SeaweedClient) error {
resp, err := client.ListClusterNodes(context.Background(), s.listClusterNodesRequest(cluster.FilerType))
if err != nil {
return err
}
// Process each filer node
for _, node := range resp.ClusterNodes {
createdAt := time.Unix(0, node.CreatedAtNs)
filerInfo := FilerInfo{
Address: pb.ServerAddress(node.Address).ToHttpAddress(),
DataCenter: node.DataCenter,
Rack: node.Rack,
Version: node.Version,
CreatedAt: createdAt,
}
filers = append(filers, filerInfo)
}
return nil
})
if err != nil {
return nil, fmt.Errorf("failed to get filer nodes from master: %w", err)
}
// Sort filers by address for consistent ordering on page refresh
sort.Slice(filers, func(i, j int) bool {
return filers[i].Address < filers[j].Address
})
return &ClusterFilersData{
Filers: filers,
TotalFilers: len(filers),
LastUpdated: time.Now(),
}, nil
}
// GetClusterBrokers retrieves cluster message brokers data
func (s *AdminServer) GetClusterBrokers() (*ClusterBrokersData, error) {
var brokers []MessageBrokerInfo
// Get broker information from master using ListClusterNodes
err := s.WithMasterClient(func(client master_pb.SeaweedClient) error {
resp, err := client.ListClusterNodes(context.Background(), s.listClusterNodesRequest(cluster.BrokerType))
if err != nil {
return err
}
// Process each broker node
for _, node := range resp.ClusterNodes {
createdAt := time.Unix(0, node.CreatedAtNs)
brokerInfo := MessageBrokerInfo{
Address: pb.ServerAddress(node.Address).ToHttpAddress(),
DataCenter: node.DataCenter,
Rack: node.Rack,
Version: node.Version,
CreatedAt: createdAt,
}
brokers = append(brokers, brokerInfo)
}
return nil
})
if err != nil {
return nil, fmt.Errorf("failed to get broker nodes from master: %w", err)
}
// Sort brokers by address for consistent ordering on page refresh
sort.Slice(brokers, func(i, j int) bool {
return brokers[i].Address < brokers[j].Address
})
return &ClusterBrokersData{
Brokers: brokers,
TotalBrokers: len(brokers),
LastUpdated: time.Now(),
}, nil
}
// GetClusterS3Servers retrieves cluster S3 servers data
func (s *AdminServer) GetClusterS3Servers() (*ClusterS3ServersData, error) {
var s3Servers []S3ServerInfo
// Get S3 server information from master using ListClusterNodes
err := s.WithMasterClient(func(client master_pb.SeaweedClient) error {
resp, err := client.ListClusterNodes(context.Background(), s.listClusterNodesRequest(cluster.S3Type))
if err != nil {
return err
}
// Process each S3 server node
for _, node := range resp.ClusterNodes {
createdAt := time.Unix(0, node.CreatedAtNs)
s3ServerInfo := S3ServerInfo{
Address: pb.ServerAddress(node.Address).ToHttpAddress(),
DataCenter: node.DataCenter,
Version: node.Version,
CreatedAt: createdAt,
}
s3Servers = append(s3Servers, s3ServerInfo)
}
return nil
})
if err != nil {
return nil, fmt.Errorf("failed to get S3 server nodes from master: %w", err)
}
// Sort S3 servers by address for consistent ordering on page refresh
sort.Slice(s3Servers, func(i, j int) bool {
return s3Servers[i].Address < s3Servers[j].Address
})
return &ClusterS3ServersData{
S3Servers: s3Servers,
TotalS3Servers: len(s3Servers),
LastUpdated: time.Now(),
}, nil
}
// GetAllFilers method moved to client_management.go
// GetVolumeDetails method moved to volume_management.go
// VacuumVolume method moved to volume_management.go
// TriggerTopicRetentionPurgeAPI triggers topic retention purge via HTTP API
func (as *AdminServer) TriggerTopicRetentionPurgeAPI(w http.ResponseWriter, r *http.Request) {
err := as.TriggerTopicRetentionPurge()
if err != nil {
writeJSONError(w, http.StatusInternalServerError, err.Error())
return
}
writeJSON(w, http.StatusOK, map[string]interface{}{"message": "Topic retention purge triggered successfully"})
}
// GetConfigInfo returns information about the admin configuration
func (as *AdminServer) GetConfigInfo(w http.ResponseWriter, r *http.Request) {
configInfo := as.configPersistence.GetConfigInfo()
// Add additional admin server info
currentMaster := as.masterClient.GetMaster(context.Background())
configInfo["master_address"] = string(currentMaster)
configInfo["cache_expiration"] = as.cacheExpiration.String()
configInfo["filer_cache_expiration"] = as.filerCacheExpiration.String()
// Add maintenance system info
if as.maintenanceManager != nil {
configInfo["maintenance_enabled"] = true
configInfo["maintenance_running"] = as.maintenanceManager.IsRunning()
} else {
configInfo["maintenance_enabled"] = false
configInfo["maintenance_running"] = false
}
writeJSON(w, http.StatusOK, map[string]interface{}{
"config_info": configInfo,
"title": "Configuration Information",
})
}
// StartWorkerGrpcServer starts the worker gRPC server
func (s *AdminServer) StartWorkerGrpcServer(grpcPort int) error {
if s.workerGrpcServer != nil {
return fmt.Errorf("worker gRPC server is already running")
}
s.workerGrpcServer = NewWorkerGrpcServer(s)
return s.workerGrpcServer.StartWithTLS(grpcPort)
}
// StopWorkerGrpcServer stops the worker gRPC server
func (s *AdminServer) StopWorkerGrpcServer() error {
if s.workerGrpcServer != nil {
err := s.workerGrpcServer.Stop()
s.workerGrpcServer = nil
return err
}
return nil
}
// GetWorkerGrpcServer returns the worker gRPC server
func (s *AdminServer) GetWorkerGrpcServer() *WorkerGrpcServer {
return s.workerGrpcServer
}
// GetWorkerGrpcPort returns the worker gRPC listen port, or 0 when unavailable.
func (s *AdminServer) GetWorkerGrpcPort() int {
if s.workerGrpcServer == nil {
return 0
}
return s.workerGrpcServer.ListenPort()
}
// GetPlugin returns the plugin instance when enabled.
func (s *AdminServer) GetPlugin() *adminplugin.Plugin {
return s.plugin
}
func (s *AdminServer) acquirePluginLock(reason string) (func(), error) {
if s == nil || s.pluginLock == nil {
return func() {}, nil
}
return s.pluginLock.Acquire(reason)
}
// RequestPluginJobTypeDescriptor asks one worker for job type schema and returns the descriptor.
func (s *AdminServer) RequestPluginJobTypeDescriptor(ctx context.Context, jobType string, forceRefresh bool) (*plugin_pb.JobTypeDescriptor, error) {
if s.plugin == nil {
return nil, fmt.Errorf("plugin is not enabled")
}
return s.plugin.RequestConfigSchema(ctx, jobType, forceRefresh)
}
// LoadPluginJobTypeDescriptor loads persisted descriptor for one job type.
func (s *AdminServer) LoadPluginJobTypeDescriptor(jobType string) (*plugin_pb.JobTypeDescriptor, error) {
if s.plugin == nil {
return nil, fmt.Errorf("plugin is not enabled")
}
return s.plugin.LoadDescriptor(jobType)
}
// SavePluginJobTypeConfig persists plugin job type config in admin data dir.
func (s *AdminServer) SavePluginJobTypeConfig(config *plugin_pb.PersistedJobTypeConfig) error {
if s.plugin == nil {
return fmt.Errorf("plugin is not enabled")
}
return s.plugin.SaveJobTypeConfig(config)
}
// LoadPluginJobTypeConfig loads plugin job type config from persistence.
func (s *AdminServer) LoadPluginJobTypeConfig(jobType string) (*plugin_pb.PersistedJobTypeConfig, error) {
if s.plugin == nil {
return nil, fmt.Errorf("plugin is not enabled")
}
return s.plugin.LoadJobTypeConfig(jobType)
}
// RunPluginDetection triggers one detection pass for a job type and returns proposed jobs.
func (s *AdminServer) RunPluginDetection(
ctx context.Context,
jobType string,
clusterContext *plugin_pb.ClusterContext,
maxResults int32,
) ([]*plugin_pb.JobProposal, error) {
if s.plugin == nil {
return nil, fmt.Errorf("plugin is not enabled")
}
releaseLock, err := s.acquirePluginLock(fmt.Sprintf("plugin detection %s", jobType))
if err != nil {
return nil, err
}
if releaseLock != nil {
defer releaseLock()
}
return s.plugin.RunDetection(ctx, jobType, clusterContext, maxResults)
}
// FilterPluginProposalsWithActiveJobs drops proposals already represented by assigned/running jobs.
func (s *AdminServer) FilterPluginProposalsWithActiveJobs(
jobType string,
proposals []*plugin_pb.JobProposal,
) ([]*plugin_pb.JobProposal, int, error) {
if s.plugin == nil {
return nil, 0, fmt.Errorf("plugin is not enabled")
}
filtered, skipped := s.plugin.FilterProposalsWithActiveJobs(jobType, proposals)
return filtered, skipped, nil
}
// RunPluginDetectionWithReport triggers one detection pass and returns request metadata and proposals.
func (s *AdminServer) RunPluginDetectionWithReport(
ctx context.Context,
jobType string,
clusterContext *plugin_pb.ClusterContext,
maxResults int32,
) (*adminplugin.DetectionReport, error) {
if s.plugin == nil {
return nil, fmt.Errorf("plugin is not enabled")
}
releaseLock, err := s.acquirePluginLock(fmt.Sprintf("plugin detection %s", jobType))
if err != nil {
return nil, err
}
if releaseLock != nil {
defer releaseLock()
}
return s.plugin.RunDetectionWithReport(ctx, jobType, clusterContext, maxResults)
}
// DispatchPluginProposals dispatches a batch of proposals using the same
// capacity-aware dispatch logic as the scheduler loop (executor reservation with
// backoff, per-job retry on transient errors). The dispatch takes the cluster
// admin lock around each job itself; callers must not hold it.
func (s *AdminServer) DispatchPluginProposals(
ctx context.Context,
jobType string,
proposals []*plugin_pb.JobProposal,
clusterContext *plugin_pb.ClusterContext,
) (successCount, errorCount, canceledCount int, err error) {
if s.plugin == nil {
return 0, 0, 0, fmt.Errorf("plugin is not enabled")
}
sc, ec, cc := s.plugin.DispatchProposals(ctx, jobType, proposals, clusterContext)
return sc, ec, cc, nil
}
// ExecutePluginJob dispatches one job to a capable worker and waits for completion.
func (s *AdminServer) ExecutePluginJob(
ctx context.Context,
job *plugin_pb.JobSpec,
clusterContext *plugin_pb.ClusterContext,
attempt int32,
) (*plugin_pb.JobCompleted, error) {
if s.plugin == nil {
return nil, fmt.Errorf("plugin is not enabled")
}
jobType := ""
if job != nil {
jobType = strings.TrimSpace(job.JobType)
}
releaseLock, err := s.acquirePluginLock(fmt.Sprintf("plugin execution %s", jobType))
if err != nil {
return nil, err
}
if releaseLock != nil {
defer releaseLock()
}
return s.plugin.ExecuteJob(ctx, job, clusterContext, attempt)
}
// GetPluginRunHistory returns the bounded run history (last 10 success + last 10 error).
func (s *AdminServer) GetPluginRunHistory(jobType string) (*adminplugin.JobTypeRunHistory, error) {
if s.plugin == nil {
return nil, fmt.Errorf("plugin is not enabled")
}
return s.plugin.LoadRunHistory(jobType)
}
// ListPluginJobTypes returns known plugin job types from connected worker registry and persisted data.
func (s *AdminServer) ListPluginJobTypes() ([]adminplugin.JobTypeInfo, error) {
if s.plugin == nil {
return nil, fmt.Errorf("plugin is not enabled")
}
return s.plugin.ListKnownJobTypes()
}
// GetPluginWorkers returns currently connected plugin workers.
func (s *AdminServer) GetPluginWorkers() []*adminplugin.WorkerSession {
if s.plugin == nil {
return nil
}
return s.plugin.ListWorkers()
}
// ListPluginJobs returns tracked plugin jobs for monitoring.
func (s *AdminServer) ListPluginJobs(jobType, state string, limit int) []adminplugin.TrackedJob {
if s.plugin == nil {
return nil
}
return s.plugin.ListTrackedJobs(jobType, state, limit)
}
// GetPluginJob returns one tracked plugin job by ID.
func (s *AdminServer) GetPluginJob(jobID string) (*adminplugin.TrackedJob, bool) {
if s.plugin == nil {
return nil, false
}
return s.plugin.GetTrackedJob(jobID)
}
// GetPluginJobDetail returns detailed plugin job information with activity timeline.
func (s *AdminServer) GetPluginJobDetail(jobID string, activityLimit, relatedLimit int) (*adminplugin.JobDetail, bool, error) {
if s.plugin == nil {
return nil, false, fmt.Errorf("plugin is not enabled")
}
return s.plugin.BuildJobDetail(jobID, activityLimit, relatedLimit)
}
// ExpirePluginJob marks an active plugin job as failed so it no longer blocks scheduling.
func (s *AdminServer) ExpirePluginJob(jobID, reason string) (*adminplugin.TrackedJob, bool, error) {
if handler := s.expireJobHandler; handler != nil {
return handler(jobID, reason)
}
if s.plugin == nil {
return nil, false, fmt.Errorf("plugin is not enabled")
}
return s.plugin.ExpireJob(jobID, reason)
}
// ListPluginActivities returns plugin job activities for monitoring.
func (s *AdminServer) ListPluginActivities(jobType string, limit int) []adminplugin.JobActivity {
if s.plugin == nil {
return nil
}
return s.plugin.ListActivities(jobType, limit)
}
// ListPluginSchedulerStates returns per-job-type scheduler state.
func (s *AdminServer) ListPluginSchedulerStates() ([]adminplugin.SchedulerJobTypeState, error) {
if s.plugin == nil {
return nil, fmt.Errorf("plugin is not enabled")
}
return s.plugin.ListSchedulerStates()
}
// Maintenance system integration methods
// InitMaintenanceManager initializes the maintenance manager
func (s *AdminServer) InitMaintenanceManager(config *maintenance.MaintenanceConfig) {
// Hand the real config store to the manager so that, if it has to build the maintenance policy
// itself, it reads the persisted task configs instead of compiled-in defaults. Only pass it when
// a data directory is actually configured: an unconfigured store has nothing to read, and a typed
// nil pointer would satisfy the loaders' type assertion and then panic on use.
var configPersistence interface{}
if s.configPersistence != nil && s.configPersistence.IsConfigured() {
configPersistence = s.configPersistence
}
s.maintenanceManager = maintenance.NewMaintenanceManager(s, config, configPersistence)
// Set up task persistence if config persistence is available
if s.configPersistence != nil {
queue := s.maintenanceManager.GetQueue()
if queue != nil {
queue.SetPersistence(s.configPersistence)
// Load tasks from persistence on startup
if err := queue.LoadTasksFromPersistence(); err != nil {
glog.Errorf("Failed to load tasks from persistence: %v", err)
}
}
}
glog.V(1).Infof("Maintenance manager initialized (enabled: %v)", config.Enabled)
}
// GetMaintenanceManager returns the maintenance manager
func (s *AdminServer) GetMaintenanceManager() *maintenance.MaintenanceManager {
return s.maintenanceManager
}
// StartMaintenanceManager starts the maintenance manager
func (s *AdminServer) StartMaintenanceManager() error {
if s.maintenanceManager == nil {
return fmt.Errorf("maintenance manager not initialized")
}
return s.maintenanceManager.Start()
}
// StopMaintenanceManager stops the maintenance manager
func (s *AdminServer) StopMaintenanceManager() {
if s.maintenanceManager != nil {
s.maintenanceManager.Stop()
}
}
// TriggerTopicRetentionPurge triggers topic data purging based on retention policies
func (s *AdminServer) TriggerTopicRetentionPurge() error {
if s.topicRetentionPurger == nil {
return fmt.Errorf("topic retention purger not initialized")
}
glog.V(0).Infof("Triggering topic retention purge")
return s.topicRetentionPurger.PurgeExpiredTopicData()
}
// GetTopicRetentionPurger returns the topic retention purger
func (s *AdminServer) GetTopicRetentionPurger() *TopicRetentionPurger {
return s.topicRetentionPurger
}
// CreateTopicWithRetention creates a new topic with optional retention configuration
func (s *AdminServer) CreateTopicWithRetention(namespace, name string, partitionCount int32, retentionEnabled bool, retentionSeconds int64) error {
// Find broker leader to create the topic
brokerLeader, err := s.findBrokerLeader()
if err != nil {
return fmt.Errorf("failed to find broker leader: %w", err)
}
// Create retention configuration
var retention *mq_pb.TopicRetention
if retentionEnabled {
retention = &mq_pb.TopicRetention{
Enabled: true,
RetentionSeconds: retentionSeconds,
}
} else {
retention = &mq_pb.TopicRetention{
Enabled: false,
RetentionSeconds: 0,
}
}
// Create the topic via broker
err = s.withBrokerClient(brokerLeader, func(client mq_pb.SeaweedMessagingClient) error {
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
_, err := client.ConfigureTopic(ctx, &mq_pb.ConfigureTopicRequest{
Topic: &schema_pb.Topic{
Namespace: namespace,
Name: name,
},
PartitionCount: partitionCount,
Retention: retention,
})
return err
})
if err != nil {
return fmt.Errorf("failed to create topic: %w", err)
}
glog.V(0).Infof("Created topic %s.%s with %d partitions (retention: enabled=%v, seconds=%d)",
namespace, name, partitionCount, retentionEnabled, retentionSeconds)
return nil
}
// UpdateTopicRetention updates the retention configuration for an existing topic
func (s *AdminServer) UpdateTopicRetention(namespace, name string, enabled bool, retentionSeconds int64) error {
// Get broker information from master
var brokerAddress string
err := s.WithMasterClient(func(client master_pb.SeaweedClient) error {
resp, err := client.ListClusterNodes(context.Background(), s.listClusterNodesRequest(cluster.BrokerType))
if err != nil {
return err
}
// Find the first available broker
for _, node := range resp.ClusterNodes {
brokerAddress = node.Address
break
}
return nil
})
if err != nil {
return fmt.Errorf("failed to get broker nodes from master: %w", err)
}
if brokerAddress == "" {
return fmt.Errorf("no active brokers found")
}
// Create gRPC connection
conn, err := grpc.NewClient(brokerAddress, s.grpcDialOption)
if err != nil {
return fmt.Errorf("failed to connect to broker: %w", err)
}
defer conn.Close()
client := mq_pb.NewSeaweedMessagingClient(conn)
// First, get the current topic configuration to preserve existing settings
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
currentConfig, err := client.GetTopicConfiguration(ctx, &mq_pb.GetTopicConfigurationRequest{
Topic: &schema_pb.Topic{
Namespace: namespace,
Name: name,
},
})
if err != nil {
return fmt.Errorf("failed to get current topic configuration: %w", err)
}
// Create the topic configuration request, preserving all existing settings
configRequest := &mq_pb.ConfigureTopicRequest{
Topic: &schema_pb.Topic{
Namespace: namespace,
Name: name,
},
// Preserve existing partition count - this is critical!
PartitionCount: currentConfig.PartitionCount,
// Preserve existing schema if it exists
MessageRecordType: currentConfig.MessageRecordType,
KeyColumns: currentConfig.KeyColumns,
}
// Update only the retention configuration
if enabled {
configRequest.Retention = &mq_pb.TopicRetention{
RetentionSeconds: retentionSeconds,
Enabled: true,
}
} else {
// Set retention to disabled
configRequest.Retention = &mq_pb.TopicRetention{
RetentionSeconds: 0,
Enabled: false,
}
}
// Send the configuration request with preserved settings
_, err = client.ConfigureTopic(ctx, configRequest)
if err != nil {
return fmt.Errorf("failed to update topic retention: %w", err)
}
glog.V(0).Infof("Updated topic %s.%s retention (enabled: %v, seconds: %d) while preserving %d partitions",
namespace, name, enabled, retentionSeconds, currentConfig.PartitionCount)
return nil
}
// Shutdown gracefully shuts down the admin server
func (s *AdminServer) Shutdown() {
glog.V(1).Infof("Shutting down admin server...")
// Cancel background goroutines (vacuum monitor, etc.)
if s.bgCancel != nil {
s.bgCancel()
}
// Stop maintenance manager
s.StopMaintenanceManager()
if s.adminPresenceLock != nil {
s.adminPresenceLock.Stop()
}
if s.plugin != nil {
s.plugin.Shutdown()
}
// Stop worker gRPC server
if err := s.StopWorkerGrpcServer(); err != nil {
glog.Errorf("Failed to stop worker gRPC server: %v", err)
}
// Shutdown credential manager
if s.credentialManager != nil {
s.credentialManager.Shutdown()
}
glog.V(1).Infof("Admin server shutdown complete")
}
// Function to extract Object Lock information from bucket entry using shared utilities
func extractObjectLockInfoFromEntry(entry *filer_pb.Entry) (bool, string, int32) {
// Try to load Object Lock configuration using shared utility
if config, found := s3api.LoadObjectLockConfigurationFromExtended(entry); found {
return s3api.ExtractObjectLockInfoFromConfig(config)
}
return false, "", 0
}
// Function to extract versioning information from bucket entry using shared utilities
func extractVersioningFromEntry(entry *filer_pb.Entry) string {
return s3api.GetVersioningStatus(entry)
}
func extractLifecycleCountsFromEntry(entry *filer_pb.Entry) (ruleCount, enabledCount int) {
xmlBytes := entry.Extended[scheduler.BucketLifecycleConfigurationXMLKey]
if len(xmlBytes) == 0 {
return
}
rules, err := lifecycle_xml.ParseCanonical(xmlBytes)
if err != nil {
return
}
for _, rule := range rules {
ruleCount++
if rule.Status == s3lifecycle.StatusEnabled {
enabledCount++
}
}
return
}
// extractPolicyStatementCountFromEntry returns the number of statements in
// the bucket's policy, or 0 if it has none or the stored JSON can't be
// parsed. Forgiving on parse failure, same as extractLifecycleCountsFromEntry.
func extractPolicyStatementCountFromEntry(entry *filer_pb.Entry) int {
policyJSON := entry.Extended[s3api.BUCKET_POLICY_METADATA_KEY]
if len(policyJSON) == 0 {
return 0
}
var doc policy_engine.PolicyDocument
if err := json.Unmarshal(policyJSON, &doc); err != nil {
return 0
}
return len(doc.Statement)
}
// GetConfigPersistence returns the config persistence manager
func (as *AdminServer) GetConfigPersistence() *ConfigPersistence {
return as.configPersistence
}
type collectionStats struct {
PhysicalSize int64
LogicalSize int64
FileCount int64
}
// ecVolumeCounts combines EC volume counts reported by multiple nodes.
// Every node holding any shard of an EC volume reports the same file_count
// (total entries in the replicated .ecx), so we dedupe it per volume id by
// taking the max — a node that has not yet finished loading .ecx would
// otherwise pin the aggregate at 0 and zero out the bucket object count.
// In contrast, a needle delete is recorded locally on the shard holder
// that served it, so each node reports its own tombstone count and the
// true delete total is the sum across nodes.
type ecVolumeCounts struct {
collection string
fileCount uint64
deleteCount uint64
}
// volumeLiveCount is the live chunk count of one regular volume. Replicas
// mirror each other's needles and their deletes, so the fullest report is the
// volume's count — dividing each report by the copy count instead would lose
// a chunk to integer truncation and would halve a volume whose second replica
// has not reported yet.
type volumeLiveCount struct {
collection string
live uint64
}
func collectCollectionStats(topologyInfo *master_pb.TopologyInfo) map[string]collectionStats {
collectionMap := make(map[string]collectionStats)
ecVolumeAgg := make(map[uint32]*ecVolumeCounts)
volumeAgg := make(map[uint32]*volumeLiveCount)
for _, dc := range topologyInfo.DataCenterInfos {
for _, rack := range dc.RackInfos {
for _, node := range rack.DataNodeInfos {
for _, diskInfo := range node.DiskInfos {
for _, volInfo := range diskInfo.VolumeInfos {
collection := volInfo.Collection
if collection == "" {
collection = "default"
}
data := collectionMap[collection]
data.PhysicalSize += int64(volInfo.Size)
rp, _ := super_block.NewReplicaPlacementFromByte(byte(volInfo.ReplicaPlacement))
// NewReplicaPlacementFromByte never returns a nil rp. If there's an error,
// it returns a zero-valued ReplicaPlacement, for which GetCopyCount() is 1.
// This provides a safe fallback, so we can ignore the error.
replicaCount := int64(rp.GetCopyCount())
if volInfo.Size >= volInfo.DeletedByteCount {
data.LogicalSize += int64(volInfo.Size-volInfo.DeletedByteCount) / replicaCount
}
collectionMap[collection] = data
if volInfo.FileCount >= volInfo.DeleteCount {
agg, ok := volumeAgg[volInfo.Id]
if !ok {
agg = &volumeLiveCount{collection: collection}
volumeAgg[volInfo.Id] = agg
}
if live := volInfo.FileCount - volInfo.DeleteCount; live > agg.live {
agg.live = live
}
}
}
for _, ecShardInfo := range diskInfo.EcShardInfos {
collection := ecShardInfo.Collection
if collection == "" {
collection = "default"
}
shards := erasure_coding.ShardsInfoFromVolumeEcShardInformationMessage(ecShardInfo)
data := collectionMap[collection]
data.PhysicalSize += int64(shards.TotalSize())
data.LogicalSize += int64(shards.MinusParityShards(erasure_coding.DataShardsCount).TotalSize())
collectionMap[collection] = data
// fileCount is volume-wide (same .ecx on every shard
// holder) so take the max to dedupe — a node that has
// not yet finished loading .ecx reports 0 and must not
// pin the aggregate. deleteCount is node-local and is
// summed across shard holders.
agg, ok := ecVolumeAgg[ecShardInfo.Id]
if !ok {
agg = &ecVolumeCounts{collection: collection}
ecVolumeAgg[ecShardInfo.Id] = agg
}
if ecShardInfo.FileCount > agg.fileCount {
agg.fileCount = ecShardInfo.FileCount
}
agg.deleteCount += ecShardInfo.DeleteCount
}
}
}
}
}
// Fold the per-volume live counts in, one entry per volume id no matter
// how many replicas reported it.
for _, agg := range volumeAgg {
data := collectionMap[agg.collection]
data.FileCount += int64(agg.live)
collectionMap[agg.collection] = data
}
// Fold EC per-volume counts into the collection totals. fileCount is
// deduped via max across every node reporting shards for the volume;
// deleteCount is summed across the same nodes.
for vid, agg := range ecVolumeAgg {
data := collectionMap[agg.collection]
if agg.fileCount >= agg.deleteCount {
data.FileCount += int64(agg.fileCount - agg.deleteCount)
} else {
glog.Warningf("ec volume %d in collection %q: summed delete_count=%d exceeds file_count=%d; skipping object count",
vid, agg.collection, agg.deleteCount, agg.fileCount)
}
collectionMap[agg.collection] = data
}
return collectionMap
}
// totalCollectionFileCount is the cluster-wide live chunk count: the sum of
// every collection's deduped count. Volumes and EC volumes are reported by
// each replica or shard holder, so only this aggregation counts a chunk once.
func totalCollectionFileCount(topologyInfo *master_pb.TopologyInfo) int64 {
var total int64
for _, stats := range collectCollectionStats(topologyInfo) {
total += stats.FileCount
}
return total
}
// getCollectionStats returns current collection statistics with caching
func (s *AdminServer) getCollectionStats() (map[string]collectionStats, error) {
now := time.Now()
if s.collectionStatsCache != nil && now.Sub(s.lastCollectionStatsUpdate) < s.collectionStatsCacheThreshold {
return s.collectionStatsCache, nil
}
err := s.WithMasterClient(func(client master_pb.SeaweedClient) error {
resp, err := pb.CollectVolumeList(context.Background(), client, &master_pb.VolumeListRequest{})
if err != nil {
return err
}
if resp.TopologyInfo != nil {
s.collectionStatsCache = collectCollectionStats(resp.TopologyInfo)
s.lastCollectionStatsUpdate = now
}
return nil
})
return s.collectionStatsCache, err
}