mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-08 23:55:51 +00:00
Compare commits
9
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
da39ee65a3 | ||
|
|
e91d43ef08 | ||
|
|
58cf7e9a89 | ||
|
|
338be16254 | ||
|
|
1b6e96614d | ||
|
|
4eb45ecc5e | ||
|
|
1f3df6e9ef | ||
|
|
fcd5de9710 | ||
|
|
b6f6f0187e |
@@ -51,7 +51,7 @@ metadata:
|
||||
annotations:
|
||||
"helm.sh/hook": post-install,post-upgrade
|
||||
"helm.sh/hook-weight": "-5"
|
||||
"helm.sh/hook-delete-policy": hook-succeeded
|
||||
"helm.sh/hook-delete-policy": before-hook-creation,hook-succeeded
|
||||
spec:
|
||||
template:
|
||||
metadata:
|
||||
@@ -92,6 +92,7 @@ spec:
|
||||
- "/bin/sh"
|
||||
- "-ec"
|
||||
- |
|
||||
set -o pipefail
|
||||
wait_for_service() {
|
||||
local url=$1
|
||||
local max_attempts=60 # 5 minutes total (5s * 60)
|
||||
@@ -117,8 +118,7 @@ spec:
|
||||
wait_for_service "http://$WEED_CLUSTER_SW_MASTER{{ .Values.master.readinessProbe.httpGet.path }}"
|
||||
wait_for_service "http://$WEED_CLUSTER_SW_FILER{{ .Values.filer.readinessProbe.httpGet.path }}"
|
||||
{{- end }}
|
||||
set -o pipefail
|
||||
{{- range $createBuckets }}
|
||||
{{- range $createBuckets }}
|
||||
{{- $bucketName := .name }}
|
||||
{{- $bucketLock := or .lock .objectLock .withLock }}
|
||||
bucket_list=$(/bin/echo 's3.bucket.list' | /usr/bin/weed shell) || { echo "Error listing s3 buckets"; exit 1; }
|
||||
|
||||
@@ -58,7 +58,6 @@ templ Layout(view ViewContext, content templ.Component) {
|
||||
<a class="navbar-brand fw-bold" href="/admin">
|
||||
<i class="fas fa-server me-2"></i>
|
||||
SeaweedFS Admin
|
||||
<span class="badge bg-warning text-dark ms-2">ALPHA</span>
|
||||
</a>
|
||||
|
||||
<button class="navbar-toggler" type="button" data-bs-toggle="collapse" data-bs-target="#navbarNav">
|
||||
@@ -245,19 +244,6 @@ templ Layout(view ViewContext, content templ.Component) {
|
||||
</div>
|
||||
}
|
||||
</li>
|
||||
<!-- Commented out for later -->
|
||||
<!--
|
||||
<li class="nav-item">
|
||||
<a class="nav-link" href="/metrics">
|
||||
<i class="fas fa-chart-line me-2"></i>Metrics
|
||||
</a>
|
||||
</li>
|
||||
<li class="nav-item">
|
||||
<a class="nav-link" href="/logs">
|
||||
<i class="fas fa-file-alt me-2"></i>Logs
|
||||
</a>
|
||||
</li>
|
||||
-->
|
||||
</ul>
|
||||
|
||||
<h6 class="sidebar-heading px-3 mt-4 mb-1 text-muted">
|
||||
|
||||
@@ -71,14 +71,14 @@ func Layout(view ViewContext, content templ.Component) templ.Component {
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 2, "\"><link rel=\"icon\" href=\"/static/favicon.ico\" type=\"image/x-icon\"><!-- Bootstrap CSS --><link href=\"/static/css/bootstrap.min.css\" rel=\"stylesheet\"><!-- Font Awesome CSS --><link href=\"/static/css/fontawesome.min.css\" rel=\"stylesheet\"><!-- HTMX --><script src=\"/static/js/htmx.min.js\"></script><!-- Custom CSS --><link rel=\"stylesheet\" href=\"/static/css/admin.css\"></head><body><div class=\"container-fluid p-0\"><!-- Header --><header class=\"navbar navbar-expand-lg navbar-dark bg-primary sticky-top\"><div class=\"container-fluid\"><a class=\"navbar-brand fw-bold\" href=\"/admin\"><i class=\"fas fa-server me-2\"></i> SeaweedFS Admin <span class=\"badge bg-warning text-dark ms-2\">ALPHA</span></a> <button class=\"navbar-toggler\" type=\"button\" data-bs-toggle=\"collapse\" data-bs-target=\"#navbarNav\"><span class=\"navbar-toggler-icon\"></span></button><div class=\"collapse navbar-collapse\" id=\"navbarNav\"><ul class=\"navbar-nav ms-auto\"><li class=\"nav-item dropdown\"><a class=\"nav-link dropdown-toggle\" href=\"#\" role=\"button\" data-bs-toggle=\"dropdown\"><i class=\"fas fa-user me-1\"></i>")
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 2, "\"><link rel=\"icon\" href=\"/static/favicon.ico\" type=\"image/x-icon\"><!-- Bootstrap CSS --><link href=\"/static/css/bootstrap.min.css\" rel=\"stylesheet\"><!-- Font Awesome CSS --><link href=\"/static/css/fontawesome.min.css\" rel=\"stylesheet\"><!-- HTMX --><script src=\"/static/js/htmx.min.js\"></script><!-- Custom CSS --><link rel=\"stylesheet\" href=\"/static/css/admin.css\"></head><body><div class=\"container-fluid p-0\"><!-- Header --><header class=\"navbar navbar-expand-lg navbar-dark bg-primary sticky-top\"><div class=\"container-fluid\"><a class=\"navbar-brand fw-bold\" href=\"/admin\"><i class=\"fas fa-server me-2\"></i> SeaweedFS Admin</a> <button class=\"navbar-toggler\" type=\"button\" data-bs-toggle=\"collapse\" data-bs-target=\"#navbarNav\"><span class=\"navbar-toggler-icon\"></span></button><div class=\"collapse navbar-collapse\" id=\"navbarNav\"><ul class=\"navbar-nav ms-auto\"><li class=\"nav-item dropdown\"><a class=\"nav-link dropdown-toggle\" href=\"#\" role=\"button\" data-bs-toggle=\"dropdown\"><i class=\"fas fa-user me-1\"></i>")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
var templ_7745c5c3_Var3 string
|
||||
templ_7745c5c3_Var3, templ_7745c5c3_Err = templ.JoinStringErrs(username)
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 72, Col: 73}
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 71, Col: 73}
|
||||
}
|
||||
_, templ_7745c5c3_Err = templ_7745c5c3_Buffer.WriteString(templ.EscapeString(templ_7745c5c3_Var3))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
@@ -113,7 +113,7 @@ func Layout(view ViewContext, content templ.Component) templ.Component {
|
||||
var templ_7745c5c3_Var6 string
|
||||
templ_7745c5c3_Var6, templ_7745c5c3_Err = templ.JoinStringErrs(fmt.Sprintf("%t", isClusterPage))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 99, Col: 207}
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 98, Col: 207}
|
||||
}
|
||||
_, templ_7745c5c3_Err = templ_7745c5c3_Buffer.WriteString(templ.EscapeString(templ_7745c5c3_Var6))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
@@ -170,7 +170,7 @@ func Layout(view ViewContext, content templ.Component) templ.Component {
|
||||
var templ_7745c5c3_Var11 string
|
||||
templ_7745c5c3_Var11, templ_7745c5c3_Err = templ.JoinStringErrs(fmt.Sprintf("%t", isStoragePage))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 124, Col: 207}
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 123, Col: 207}
|
||||
}
|
||||
_, templ_7745c5c3_Err = templ_7745c5c3_Buffer.WriteString(templ.EscapeString(templ_7745c5c3_Var11))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
@@ -254,7 +254,7 @@ func Layout(view ViewContext, content templ.Component) templ.Component {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
}
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 24, "</li><!-- Commented out for later --><!--\n <li class=\"nav-item\">\n <a class=\"nav-link\" href=\"/metrics\">\n <i class=\"fas fa-chart-line me-2\"></i>Metrics\n </a>\n </li>\n <li class=\"nav-item\">\n <a class=\"nav-link\" href=\"/logs\">\n <i class=\"fas fa-file-alt me-2\"></i>Logs\n </a>\n </li>\n --></ul><h6 class=\"sidebar-heading px-3 mt-4 mb-1 text-muted\"><span>WORKERS</span></h6><ul class=\"nav flex-column\"><li class=\"nav-item\">")
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 24, "</li></ul><h6 class=\"sidebar-heading px-3 mt-4 mb-1 text-muted\"><span>WORKERS</span></h6><ul class=\"nav flex-column\"><li class=\"nav-item\">")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
@@ -344,7 +344,7 @@ func Layout(view ViewContext, content templ.Component) templ.Component {
|
||||
var templ_7745c5c3_Var14 string
|
||||
templ_7745c5c3_Var14, templ_7745c5c3_Err = templ.JoinStringErrs(fmt.Sprintf("%d", time.Now().Year()))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 339, Col: 60}
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 325, Col: 60}
|
||||
}
|
||||
_, templ_7745c5c3_Err = templ_7745c5c3_Buffer.WriteString(templ.EscapeString(templ_7745c5c3_Var14))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
@@ -357,7 +357,7 @@ func Layout(view ViewContext, content templ.Component) templ.Component {
|
||||
var templ_7745c5c3_Var15 string
|
||||
templ_7745c5c3_Var15, templ_7745c5c3_Err = templ.JoinStringErrs(version.VERSION_NUMBER)
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 339, Col: 102}
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 325, Col: 102}
|
||||
}
|
||||
_, templ_7745c5c3_Err = templ_7745c5c3_Buffer.WriteString(templ.EscapeString(templ_7745c5c3_Var15))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
@@ -409,7 +409,7 @@ func LoginForm(title string, errorMessage string, csrfToken string) templ.Compon
|
||||
var templ_7745c5c3_Var17 string
|
||||
templ_7745c5c3_Var17, templ_7745c5c3_Err = templ.JoinStringErrs(title)
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 367, Col: 17}
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 353, Col: 17}
|
||||
}
|
||||
_, templ_7745c5c3_Err = templ_7745c5c3_Buffer.WriteString(templ.EscapeString(templ_7745c5c3_Var17))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
@@ -422,7 +422,7 @@ func LoginForm(title string, errorMessage string, csrfToken string) templ.Compon
|
||||
var templ_7745c5c3_Var18 string
|
||||
templ_7745c5c3_Var18, templ_7745c5c3_Err = templ.JoinStringErrs(title)
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 381, Col: 57}
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 367, Col: 57}
|
||||
}
|
||||
_, templ_7745c5c3_Err = templ_7745c5c3_Buffer.WriteString(templ.EscapeString(templ_7745c5c3_Var18))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
@@ -440,7 +440,7 @@ func LoginForm(title string, errorMessage string, csrfToken string) templ.Compon
|
||||
var templ_7745c5c3_Var19 string
|
||||
templ_7745c5c3_Var19, templ_7745c5c3_Err = templ.JoinStringErrs(errorMessage)
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 388, Col: 45}
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 374, Col: 45}
|
||||
}
|
||||
_, templ_7745c5c3_Err = templ_7745c5c3_Buffer.WriteString(templ.EscapeString(templ_7745c5c3_Var19))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
@@ -458,7 +458,7 @@ func LoginForm(title string, errorMessage string, csrfToken string) templ.Compon
|
||||
var templ_7745c5c3_Var20 string
|
||||
templ_7745c5c3_Var20, templ_7745c5c3_Err = templ.JoinStringErrs(csrfToken)
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 393, Col: 84}
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `view/layout/layout.templ`, Line: 379, Col: 84}
|
||||
}
|
||||
_, templ_7745c5c3_Err = templ_7745c5c3_Buffer.WriteString(templ.EscapeString(templ_7745c5c3_Var20))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
|
||||
@@ -18,6 +18,20 @@ import (
|
||||
// need to have logic similar to FilerStoreWrapper
|
||||
// e.g. fill fileId field for chunks
|
||||
|
||||
// bufferedEvent represents a subscription event captured during a directory refresh.
|
||||
type bufferedEvent struct {
|
||||
oldPath util.FullPath
|
||||
newEntry *filer.Entry
|
||||
}
|
||||
|
||||
// refreshState tracks events that arrive while a directory is being refreshed
|
||||
// from the filer. After the refresh snapshot is applied, buffered events are
|
||||
// replayed so that creates, deletes, and updates that raced with the snapshot
|
||||
// are not lost.
|
||||
type refreshState struct {
|
||||
events []bufferedEvent
|
||||
}
|
||||
|
||||
type MetaCache struct {
|
||||
root util.FullPath
|
||||
localStore filer.VirtualFilerStore
|
||||
@@ -29,6 +43,7 @@ type MetaCache struct {
|
||||
invalidateFunc func(fullpath util.FullPath, entry *filer_pb.Entry)
|
||||
onDirectoryUpdate func(dir util.FullPath)
|
||||
visitGroup singleflight.Group // deduplicates concurrent EnsureVisited calls for the same path
|
||||
refreshing map[util.FullPath]*refreshState
|
||||
}
|
||||
|
||||
func NewMetaCache(dbFolder string, uidGidMapper *UidGidMapper, root util.FullPath,
|
||||
@@ -45,6 +60,7 @@ func NewMetaCache(dbFolder string, uidGidMapper *UidGidMapper, root util.FullPat
|
||||
invalidateFunc: func(fullpath util.FullPath, entry *filer_pb.Entry) {
|
||||
invalidateFunc(fullpath, entry)
|
||||
},
|
||||
refreshing: make(map[util.FullPath]*refreshState),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -69,6 +85,11 @@ func openMetaStore(dbFolder string) (*leveldb.LevelDBStore, filer.VirtualFilerSt
|
||||
func (mc *MetaCache) InsertEntry(ctx context.Context, entry *filer.Entry) error {
|
||||
mc.Lock()
|
||||
defer mc.Unlock()
|
||||
// Buffer the insert if the parent directory is being refreshed
|
||||
dir, _ := entry.DirAndName()
|
||||
if state := mc.isRefreshingDir(util.FullPath(dir)); state != nil {
|
||||
state.events = append(state.events, bufferedEvent{newEntry: entry})
|
||||
}
|
||||
return mc.doInsertEntry(ctx, entry)
|
||||
}
|
||||
|
||||
@@ -76,16 +97,33 @@ func (mc *MetaCache) doInsertEntry(ctx context.Context, entry *filer.Entry) erro
|
||||
return mc.localStore.InsertEntry(ctx, entry)
|
||||
}
|
||||
|
||||
// doBatchInsertEntries inserts multiple entries using LevelDB's batch write.
|
||||
// This is more efficient than inserting entries one by one.
|
||||
func (mc *MetaCache) doBatchInsertEntries(ctx context.Context, entries []*filer.Entry) error {
|
||||
return mc.leveldbStore.BatchInsertEntries(ctx, entries)
|
||||
}
|
||||
|
||||
func (mc *MetaCache) AtomicUpdateEntryFromFiler(ctx context.Context, oldPath util.FullPath, newEntry *filer.Entry) error {
|
||||
mc.Lock()
|
||||
defer mc.Unlock()
|
||||
|
||||
// If the affected directory is being refreshed, buffer the event
|
||||
// instead of applying it. It will be replayed after the snapshot commits.
|
||||
if oldPath != "" {
|
||||
dir, _ := oldPath.DirAndName()
|
||||
if state := mc.isRefreshingDir(util.FullPath(dir)); state != nil {
|
||||
state.events = append(state.events, bufferedEvent{oldPath, newEntry})
|
||||
return nil
|
||||
}
|
||||
}
|
||||
if newEntry != nil {
|
||||
newDir, _ := newEntry.DirAndName()
|
||||
if state := mc.isRefreshingDir(util.FullPath(newDir)); state != nil {
|
||||
state.events = append(state.events, bufferedEvent{oldPath, newEntry})
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
return mc.doAtomicUpdateEntryFromFiler(ctx, oldPath, newEntry)
|
||||
}
|
||||
|
||||
// doAtomicUpdateEntryFromFiler is the core logic for applying a filer event
|
||||
// to the local cache. Caller must hold mc.Lock().
|
||||
func (mc *MetaCache) doAtomicUpdateEntryFromFiler(ctx context.Context, oldPath util.FullPath, newEntry *filer.Entry) error {
|
||||
entry, err := mc.localStore.FindEntry(ctx, oldPath)
|
||||
if err != nil && err != filer_pb.ErrNotFound {
|
||||
glog.Errorf("Metacache: find entry error: %v", err)
|
||||
@@ -104,8 +142,6 @@ func (mc *MetaCache) AtomicUpdateEntryFromFiler(ctx context.Context, oldPath uti
|
||||
}
|
||||
}
|
||||
}
|
||||
} else {
|
||||
// println("unknown old directory:", oldDir)
|
||||
}
|
||||
|
||||
if newEntry != nil {
|
||||
@@ -143,6 +179,14 @@ func (mc *MetaCache) FindEntry(ctx context.Context, fp util.FullPath) (entry *fi
|
||||
func (mc *MetaCache) DeleteEntry(ctx context.Context, fp util.FullPath) (err error) {
|
||||
mc.Lock()
|
||||
defer mc.Unlock()
|
||||
// Buffer the delete if the parent directory is being refreshed
|
||||
dir, _ := fp.DirAndName()
|
||||
if state := mc.isRefreshingDir(util.FullPath(dir)); state != nil {
|
||||
state.events = append(state.events, bufferedEvent{oldPath: fp})
|
||||
}
|
||||
// Always apply the delete directly as well, so the entry is removed
|
||||
// immediately for the current node's view. CommitRefresh's snapshot
|
||||
// may re-insert it, but the buffered event replay will re-delete it.
|
||||
return mc.localStore.DeleteEntry(ctx, fp)
|
||||
}
|
||||
func (mc *MetaCache) DeleteFolderChildren(ctx context.Context, fp util.FullPath) (err error) {
|
||||
@@ -151,6 +195,58 @@ func (mc *MetaCache) DeleteFolderChildren(ctx context.Context, fp util.FullPath)
|
||||
return mc.localStore.DeleteFolderChildren(ctx, fp)
|
||||
}
|
||||
|
||||
// BeginRefresh starts buffering subscription events for dirPath.
|
||||
// While a refresh is active, AtomicUpdateEntryFromFiler will buffer events
|
||||
// instead of applying them, so they can be replayed after the snapshot is committed.
|
||||
func (mc *MetaCache) BeginRefresh(dirPath util.FullPath) {
|
||||
mc.Lock()
|
||||
defer mc.Unlock()
|
||||
mc.refreshing[dirPath] = &refreshState{}
|
||||
}
|
||||
|
||||
// CommitRefresh atomically replaces a directory's cached entries with the
|
||||
// filer snapshot, then replays any subscription events that were buffered
|
||||
// during the refresh. This ensures no creates, deletes, or updates are lost.
|
||||
func (mc *MetaCache) CommitRefresh(ctx context.Context, dirPath util.FullPath, entries []*filer.Entry) error {
|
||||
mc.Lock()
|
||||
defer mc.Unlock()
|
||||
|
||||
// Clear stale entries and insert the fresh snapshot
|
||||
if err := mc.localStore.DeleteFolderChildren(ctx, dirPath); err != nil {
|
||||
return err
|
||||
}
|
||||
if len(entries) > 0 {
|
||||
if err := mc.leveldbStore.BatchInsertEntries(ctx, entries); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
// Replay buffered events so mutations that raced with the snapshot are applied
|
||||
state := mc.refreshing[dirPath]
|
||||
delete(mc.refreshing, dirPath)
|
||||
if state != nil {
|
||||
for _, ev := range state.events {
|
||||
if err := mc.doAtomicUpdateEntryFromFiler(ctx, ev.oldPath, ev.newEntry); err != nil {
|
||||
glog.Warningf("replay buffered event for %s: %v", dirPath, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// CancelRefresh discards the refresh state without replaying buffered events.
|
||||
func (mc *MetaCache) CancelRefresh(dirPath util.FullPath) {
|
||||
mc.Lock()
|
||||
defer mc.Unlock()
|
||||
delete(mc.refreshing, dirPath)
|
||||
}
|
||||
|
||||
// isRefreshing returns the refresh state for the directory containing fp,
|
||||
// or nil if no refresh is active. Caller must hold mc.Lock().
|
||||
func (mc *MetaCache) isRefreshingDir(dirPath util.FullPath) *refreshState {
|
||||
return mc.refreshing[dirPath]
|
||||
}
|
||||
|
||||
func (mc *MetaCache) ListDirectoryEntries(ctx context.Context, dirPath util.FullPath, startFileName string, includeStartFile bool, limit int64, eachEntryFunc filer.ListEachEntryFunc) error {
|
||||
mc.RLock()
|
||||
defer mc.RUnlock()
|
||||
|
||||
@@ -49,11 +49,6 @@ func EnsureVisited(mc *MetaCache, client filer_pb.FilerClient, dirPath util.Full
|
||||
return g.Wait()
|
||||
}
|
||||
|
||||
// batchInsertSize is the number of entries to accumulate before flushing to LevelDB.
|
||||
// 100 provides a balance between memory usage (~100 Entry pointers) and write efficiency
|
||||
// (fewer disk syncs). Larger values reduce I/O overhead but increase memory and latency.
|
||||
const batchInsertSize = 100
|
||||
|
||||
func doEnsureVisited(ctx context.Context, mc *MetaCache, client filer_pb.FilerClient, path util.FullPath) error {
|
||||
// Use singleflight to deduplicate concurrent requests for the same path
|
||||
_, err, _ := mc.visitGroup.Do(string(path), func() (interface{}, error) {
|
||||
@@ -69,42 +64,36 @@ func doEnsureVisited(ctx context.Context, mc *MetaCache, client filer_pb.FilerCl
|
||||
|
||||
glog.V(4).Infof("ReadDirAllEntries %s ...", path)
|
||||
|
||||
// Collect entries in batches for efficient LevelDB writes
|
||||
var batch []*filer.Entry
|
||||
// Start buffering subscription events for this directory.
|
||||
// Any events that arrive during the fetch will be replayed
|
||||
// after the snapshot is committed, preventing lost mutations.
|
||||
mc.BeginRefresh(path)
|
||||
|
||||
// Collect all entries from the filer. No lock is held during
|
||||
// network I/O so subscription events can still be buffered.
|
||||
var allEntries []*filer.Entry
|
||||
|
||||
fetchErr := util.Retry("ReadDirAllEntries", func() error {
|
||||
batch = nil // Reset batch on retry, allow GC of previous entries
|
||||
allEntries = nil // Reset on retry
|
||||
return filer_pb.ReadDirAllEntries(ctx, client, path, "", func(pbEntry *filer_pb.Entry, isLast bool) error {
|
||||
entry := filer.FromPbEntry(string(path), pbEntry)
|
||||
if IsHiddenSystemEntry(string(path), entry.Name()) {
|
||||
return nil
|
||||
}
|
||||
|
||||
batch = append(batch, entry)
|
||||
|
||||
// Flush batch when it reaches the threshold
|
||||
// Don't rely on isLast here - hidden entries may cause early return
|
||||
if len(batch) >= batchInsertSize {
|
||||
// No lock needed - LevelDB Write() is thread-safe
|
||||
if err := mc.doBatchInsertEntries(ctx, batch); err != nil {
|
||||
return fmt.Errorf("batch insert for %s: %w", path, err)
|
||||
}
|
||||
// Create new slice to allow GC of flushed entries
|
||||
batch = make([]*filer.Entry, 0, batchInsertSize)
|
||||
}
|
||||
allEntries = append(allEntries, entry)
|
||||
return nil
|
||||
})
|
||||
})
|
||||
|
||||
if fetchErr != nil {
|
||||
mc.CancelRefresh(path)
|
||||
return nil, fmt.Errorf("list %s: %w", path, fetchErr)
|
||||
}
|
||||
|
||||
// Flush any remaining entries in the batch
|
||||
if len(batch) > 0 {
|
||||
if err := mc.doBatchInsertEntries(ctx, batch); err != nil {
|
||||
return nil, fmt.Errorf("batch insert remaining for %s: %w", path, err)
|
||||
}
|
||||
// Atomically replace cached entries with the snapshot and replay
|
||||
// any buffered events that arrived during the fetch.
|
||||
if err := mc.CommitRefresh(ctx, path, allEntries); err != nil {
|
||||
return nil, fmt.Errorf("commit refresh for %s: %w", path, err)
|
||||
}
|
||||
mc.markCachedFn(path)
|
||||
return nil, nil
|
||||
|
||||
@@ -0,0 +1,340 @@
|
||||
package meta_cache
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/filer"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
)
|
||||
|
||||
func newTestMetaCache(t *testing.T) (*MetaCache, map[util.FullPath]bool) {
|
||||
t.Helper()
|
||||
uidGidMapper, err := NewUidGidMapper("", "")
|
||||
if err != nil {
|
||||
t.Fatalf("create uid/gid mapper: %v", err)
|
||||
}
|
||||
cached := make(map[util.FullPath]bool)
|
||||
mc := NewMetaCache(
|
||||
filepath.Join(t.TempDir(), "meta"),
|
||||
uidGidMapper,
|
||||
util.FullPath("/"),
|
||||
func(path util.FullPath) { cached[path] = true },
|
||||
func(path util.FullPath) bool { return cached[path] },
|
||||
func(util.FullPath, *filer_pb.Entry) {},
|
||||
nil,
|
||||
)
|
||||
t.Cleanup(func() { mc.Shutdown() })
|
||||
return mc, cached
|
||||
}
|
||||
|
||||
func makeEntry(dir, name string) *filer.Entry {
|
||||
return &filer.Entry{
|
||||
FullPath: util.NewFullPath(dir, name),
|
||||
Attr: filer.Attr{Mode: 0644},
|
||||
}
|
||||
}
|
||||
|
||||
func listEntries(t *testing.T, mc *MetaCache, dir string) []string {
|
||||
t.Helper()
|
||||
var names []string
|
||||
err := mc.ListDirectoryEntries(context.Background(), util.FullPath(dir), "", false, 10000, func(entry *filer.Entry) (bool, error) {
|
||||
names = append(names, entry.Name())
|
||||
return true, nil
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("list %s: %v", dir, err)
|
||||
}
|
||||
return names
|
||||
}
|
||||
|
||||
// TestCommitRefreshDeleteDuringRefresh verifies that a delete event arriving
|
||||
// during a directory refresh is not overwritten by the stale snapshot.
|
||||
func TestCommitRefreshDeleteDuringRefresh(t *testing.T) {
|
||||
mc, cached := newTestMetaCache(t)
|
||||
ctx := context.Background()
|
||||
dir := util.FullPath("/testdir")
|
||||
cached[dir] = true
|
||||
|
||||
// Pre-populate the cache with 3 entries
|
||||
for _, name := range []string{"a", "b", "c"} {
|
||||
if err := mc.InsertEntry(ctx, makeEntry("/testdir", name)); err != nil {
|
||||
t.Fatalf("insert %s: %v", name, err)
|
||||
}
|
||||
}
|
||||
|
||||
// Start a refresh — simulate fetching a snapshot that includes all 3
|
||||
mc.BeginRefresh(dir)
|
||||
|
||||
// While refresh is in progress, a subscription event deletes "b"
|
||||
if err := mc.AtomicUpdateEntryFromFiler(ctx, util.NewFullPath("/testdir", "b"), nil); err != nil {
|
||||
t.Fatalf("atomic delete b: %v", err)
|
||||
}
|
||||
|
||||
// Commit the snapshot (which still includes "b") — the buffered delete
|
||||
// should be replayed, removing "b" from the final result
|
||||
snapshot := []*filer.Entry{
|
||||
makeEntry("/testdir", "a"),
|
||||
makeEntry("/testdir", "b"),
|
||||
makeEntry("/testdir", "c"),
|
||||
}
|
||||
if err := mc.CommitRefresh(ctx, dir, snapshot); err != nil {
|
||||
t.Fatalf("commit refresh: %v", err)
|
||||
}
|
||||
|
||||
names := listEntries(t, mc, "/testdir")
|
||||
expected := map[string]bool{"a": true, "c": true}
|
||||
if len(names) != len(expected) {
|
||||
t.Fatalf("expected %v, got %v", expected, names)
|
||||
}
|
||||
for _, n := range names {
|
||||
if !expected[n] {
|
||||
t.Errorf("unexpected entry %q in listing", n)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestCommitRefreshCreateDuringRefresh verifies that a create event arriving
|
||||
// during a directory refresh is preserved after the snapshot is applied.
|
||||
func TestCommitRefreshCreateDuringRefresh(t *testing.T) {
|
||||
mc, cached := newTestMetaCache(t)
|
||||
ctx := context.Background()
|
||||
dir := util.FullPath("/testdir")
|
||||
cached[dir] = true
|
||||
|
||||
// Pre-populate with "a"
|
||||
if err := mc.InsertEntry(ctx, makeEntry("/testdir", "a")); err != nil {
|
||||
t.Fatalf("insert a: %v", err)
|
||||
}
|
||||
|
||||
// Start refresh — snapshot will have "a" only
|
||||
mc.BeginRefresh(dir)
|
||||
|
||||
// During refresh, a subscription event creates "d"
|
||||
newEntry := makeEntry("/testdir", "d")
|
||||
if err := mc.AtomicUpdateEntryFromFiler(ctx, "", newEntry); err != nil {
|
||||
t.Fatalf("atomic create d: %v", err)
|
||||
}
|
||||
|
||||
// Commit snapshot (only "a") — buffered create of "d" should be replayed
|
||||
snapshot := []*filer.Entry{makeEntry("/testdir", "a")}
|
||||
if err := mc.CommitRefresh(ctx, dir, snapshot); err != nil {
|
||||
t.Fatalf("commit refresh: %v", err)
|
||||
}
|
||||
|
||||
names := listEntries(t, mc, "/testdir")
|
||||
expected := map[string]bool{"a": true, "d": true}
|
||||
if len(names) != len(expected) {
|
||||
t.Fatalf("expected %v, got %v", expected, names)
|
||||
}
|
||||
for _, n := range names {
|
||||
if !expected[n] {
|
||||
t.Errorf("unexpected entry %q", n)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestCommitRefreshStaleEntriesRemoved verifies that entries deleted on the
|
||||
// filer (absent from the snapshot) are removed from the local cache.
|
||||
func TestCommitRefreshStaleEntriesRemoved(t *testing.T) {
|
||||
mc, cached := newTestMetaCache(t)
|
||||
ctx := context.Background()
|
||||
dir := util.FullPath("/testdir")
|
||||
cached[dir] = true
|
||||
|
||||
// Pre-populate with a, b, c
|
||||
for _, name := range []string{"a", "b", "c"} {
|
||||
if err := mc.InsertEntry(ctx, makeEntry("/testdir", name)); err != nil {
|
||||
t.Fatalf("insert %s: %v", name, err)
|
||||
}
|
||||
}
|
||||
|
||||
// Refresh with snapshot that only has "a" and "c" — "b" was deleted on filer
|
||||
mc.BeginRefresh(dir)
|
||||
snapshot := []*filer.Entry{
|
||||
makeEntry("/testdir", "a"),
|
||||
makeEntry("/testdir", "c"),
|
||||
}
|
||||
if err := mc.CommitRefresh(ctx, dir, snapshot); err != nil {
|
||||
t.Fatalf("commit refresh: %v", err)
|
||||
}
|
||||
|
||||
names := listEntries(t, mc, "/testdir")
|
||||
for _, n := range names {
|
||||
if n == "b" {
|
||||
t.Errorf("stale entry 'b' should have been removed by refresh")
|
||||
}
|
||||
}
|
||||
if len(names) != 2 {
|
||||
t.Fatalf("expected [a, c], got %v", names)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCommitRefreshLocalDeleteBuffered verifies that a local DeleteEntry
|
||||
// during refresh is both applied immediately and replayed after commit.
|
||||
func TestCommitRefreshLocalDeleteBuffered(t *testing.T) {
|
||||
mc, cached := newTestMetaCache(t)
|
||||
ctx := context.Background()
|
||||
dir := util.FullPath("/testdir")
|
||||
cached[dir] = true
|
||||
|
||||
// Pre-populate
|
||||
for _, name := range []string{"a", "b"} {
|
||||
if err := mc.InsertEntry(ctx, makeEntry("/testdir", name)); err != nil {
|
||||
t.Fatalf("insert %s: %v", name, err)
|
||||
}
|
||||
}
|
||||
|
||||
mc.BeginRefresh(dir)
|
||||
|
||||
// Local delete of "a" — simulates Unlink calling DeleteEntry directly
|
||||
if err := mc.DeleteEntry(ctx, util.NewFullPath("/testdir", "a")); err != nil {
|
||||
t.Fatalf("delete a: %v", err)
|
||||
}
|
||||
|
||||
// Verify "a" is gone immediately (before commit)
|
||||
_, err := mc.FindEntry(ctx, util.NewFullPath("/testdir", "a"))
|
||||
if err != filer_pb.ErrNotFound {
|
||||
t.Fatalf("expected ErrNotFound for deleted entry, got: %v", err)
|
||||
}
|
||||
|
||||
// Commit with snapshot that includes "a" — the buffered delete replays
|
||||
snapshot := []*filer.Entry{
|
||||
makeEntry("/testdir", "a"),
|
||||
makeEntry("/testdir", "b"),
|
||||
}
|
||||
if err := mc.CommitRefresh(ctx, dir, snapshot); err != nil {
|
||||
t.Fatalf("commit refresh: %v", err)
|
||||
}
|
||||
|
||||
// "a" should still be gone after commit
|
||||
names := listEntries(t, mc, "/testdir")
|
||||
for _, n := range names {
|
||||
if n == "a" {
|
||||
t.Errorf("locally deleted entry 'a' reappeared after commit")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestCommitRefreshLocalCreateBuffered verifies that a local InsertEntry
|
||||
// during refresh is preserved after the snapshot is applied.
|
||||
func TestCommitRefreshLocalCreateBuffered(t *testing.T) {
|
||||
mc, cached := newTestMetaCache(t)
|
||||
ctx := context.Background()
|
||||
dir := util.FullPath("/testdir")
|
||||
cached[dir] = true
|
||||
|
||||
mc.BeginRefresh(dir)
|
||||
|
||||
// Local create during refresh
|
||||
if err := mc.InsertEntry(ctx, makeEntry("/testdir", "new")); err != nil {
|
||||
t.Fatalf("insert new: %v", err)
|
||||
}
|
||||
|
||||
// Commit with empty snapshot — the buffered create should still appear
|
||||
if err := mc.CommitRefresh(ctx, dir, nil); err != nil {
|
||||
t.Fatalf("commit refresh: %v", err)
|
||||
}
|
||||
|
||||
names := listEntries(t, mc, "/testdir")
|
||||
if len(names) != 1 || names[0] != "new" {
|
||||
t.Fatalf("expected [new], got %v", names)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCancelRefreshDiscardsBuffer verifies that CancelRefresh discards
|
||||
// buffered events and resumes normal event processing.
|
||||
func TestCancelRefreshDiscardsBuffer(t *testing.T) {
|
||||
mc, cached := newTestMetaCache(t)
|
||||
ctx := context.Background()
|
||||
dir := util.FullPath("/testdir")
|
||||
cached[dir] = true
|
||||
|
||||
if err := mc.InsertEntry(ctx, makeEntry("/testdir", "a")); err != nil {
|
||||
t.Fatalf("insert a: %v", err)
|
||||
}
|
||||
|
||||
mc.BeginRefresh(dir)
|
||||
|
||||
// Buffer a delete event
|
||||
if err := mc.AtomicUpdateEntryFromFiler(ctx, util.NewFullPath("/testdir", "a"), nil); err != nil {
|
||||
t.Fatalf("atomic delete: %v", err)
|
||||
}
|
||||
|
||||
// Cancel — buffered events are discarded
|
||||
mc.CancelRefresh(dir)
|
||||
|
||||
// "a" should still be in cache since the buffered delete was discarded
|
||||
entry, err := mc.FindEntry(ctx, util.NewFullPath("/testdir", "a"))
|
||||
if err != nil {
|
||||
t.Fatalf("find a: %v", err)
|
||||
}
|
||||
if entry == nil {
|
||||
t.Fatal("entry 'a' should still exist after cancel")
|
||||
}
|
||||
|
||||
// After cancel, events should apply normally (not buffered)
|
||||
if err := mc.AtomicUpdateEntryFromFiler(ctx, util.NewFullPath("/testdir", "a"), nil); err != nil {
|
||||
t.Fatalf("atomic delete after cancel: %v", err)
|
||||
}
|
||||
_, err = mc.FindEntry(ctx, util.NewFullPath("/testdir", "a"))
|
||||
if err != filer_pb.ErrNotFound {
|
||||
t.Fatalf("expected ErrNotFound after direct delete, got: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// TestConcurrentDeletesDuringRefresh simulates the scenario from issue #8442:
|
||||
// multiple concurrent deletes racing with a directory refresh.
|
||||
func TestConcurrentDeletesDuringRefresh(t *testing.T) {
|
||||
mc, cached := newTestMetaCache(t)
|
||||
ctx := context.Background()
|
||||
dir := util.FullPath("/testdir")
|
||||
cached[dir] = true
|
||||
|
||||
const numFiles = 1000
|
||||
|
||||
// Pre-populate
|
||||
for i := 0; i < numFiles; i++ {
|
||||
name := fmt.Sprintf("file_%04d", i)
|
||||
if err := mc.InsertEntry(ctx, makeEntry("/testdir", name)); err != nil {
|
||||
t.Fatalf("insert %s: %v", name, err)
|
||||
}
|
||||
}
|
||||
|
||||
// Build snapshot (taken before deletes happen)
|
||||
var snapshot []*filer.Entry
|
||||
for i := 0; i < numFiles; i++ {
|
||||
snapshot = append(snapshot, makeEntry("/testdir", fmt.Sprintf("file_%04d", i)))
|
||||
}
|
||||
|
||||
mc.BeginRefresh(dir)
|
||||
|
||||
// Concurrently delete all files via subscription events (simulates
|
||||
// deletes from two mount nodes)
|
||||
var wg sync.WaitGroup
|
||||
for i := 0; i < numFiles; i++ {
|
||||
wg.Add(1)
|
||||
go func(idx int) {
|
||||
defer wg.Done()
|
||||
name := fmt.Sprintf("file_%04d", idx)
|
||||
fp := util.NewFullPath("/testdir", name)
|
||||
mc.AtomicUpdateEntryFromFiler(ctx, fp, nil)
|
||||
}(i)
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
// Commit with the stale snapshot — all deletes should be replayed
|
||||
if err := mc.CommitRefresh(ctx, dir, snapshot); err != nil {
|
||||
t.Fatalf("commit refresh: %v", err)
|
||||
}
|
||||
|
||||
names := listEntries(t, mc, "/testdir")
|
||||
if len(names) != 0 {
|
||||
t.Errorf("expected 0 entries after all deletes, got %d: first few: %v",
|
||||
len(names), names[:min(5, len(names))])
|
||||
}
|
||||
}
|
||||
@@ -21,6 +21,7 @@ import (
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/iam_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/policy_engine"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3err"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util/wildcard"
|
||||
@@ -67,6 +68,10 @@ type IdentityAccessManagement struct {
|
||||
// Bucket policy engine for evaluating bucket policies
|
||||
policyEngine *BucketPolicyEngine
|
||||
|
||||
// Cached policy engine for IAM policy fallback evaluation.
|
||||
// Keyed by policy name, kept in sync by PutPolicy/DeletePolicy.
|
||||
iamPolicyEngine *policy_engine.PolicyEngine
|
||||
|
||||
// background polling
|
||||
stopChan chan struct{}
|
||||
shutdownOnce sync.Once
|
||||
@@ -658,6 +663,7 @@ func (iam *IdentityAccessManagement) ReplaceS3ApiConfiguration(config *iam_pb.S3
|
||||
iam.nameToIdentity = nameToIdentity
|
||||
iam.accessKeyIdent = accessKeyIdent
|
||||
iam.policies = policies
|
||||
iam.rebuildIAMPolicyEngineLocked()
|
||||
|
||||
// Re-add environment-based identities that were preserved
|
||||
for _, envIdent := range envIdentities {
|
||||
@@ -914,6 +920,7 @@ func (iam *IdentityAccessManagement) MergeS3ApiConfiguration(config *iam_pb.S3Ap
|
||||
iam.nameToIdentity = nameToIdentity
|
||||
iam.accessKeyIdent = accessKeyIdent
|
||||
iam.policies = policies
|
||||
iam.rebuildIAMPolicyEngineLocked()
|
||||
// Update authentication state based on whether identities exist
|
||||
// Once enabled, keep it enabled (one-way toggle)
|
||||
authJustEnabled := iam.updateAuthenticationState(len(identities))
|
||||
@@ -1658,6 +1665,51 @@ func determineIAMAuthPath(sessionToken, principal, principalArn string) iamAuthP
|
||||
return iamAuthPathNone
|
||||
}
|
||||
|
||||
// evaluateIAMPolicies evaluates attached IAM policies for a user identity.
|
||||
// Returns true if any matching statement explicitly allows the action.
|
||||
// Uses the cached iamPolicyEngine to avoid re-parsing policy JSON on every request.
|
||||
func (iam *IdentityAccessManagement) evaluateIAMPolicies(r *http.Request, identity *Identity, action Action, bucket, object string) bool {
|
||||
if identity == nil || len(identity.PolicyNames) == 0 {
|
||||
return false
|
||||
}
|
||||
|
||||
iam.m.RLock()
|
||||
engine := iam.iamPolicyEngine
|
||||
iam.m.RUnlock()
|
||||
|
||||
if engine == nil {
|
||||
return false
|
||||
}
|
||||
|
||||
resource := buildResourceARN(bucket, object)
|
||||
principal := buildPrincipalARN(identity, r)
|
||||
s3Action := ResolveS3Action(r, string(action), bucket, object)
|
||||
explicitAllow := false
|
||||
conditions := policy_engine.ExtractConditionValuesFromRequest(r)
|
||||
for k, v := range policy_engine.ExtractPrincipalVariables(principal) {
|
||||
conditions[k] = v
|
||||
}
|
||||
|
||||
for _, policyName := range identity.PolicyNames {
|
||||
result := engine.EvaluatePolicy(policyName, &policy_engine.PolicyEvaluationArgs{
|
||||
Action: s3Action,
|
||||
Resource: resource,
|
||||
Principal: principal,
|
||||
Conditions: conditions,
|
||||
Claims: identity.Claims,
|
||||
})
|
||||
|
||||
if result == policy_engine.PolicyResultDeny {
|
||||
return false
|
||||
}
|
||||
if result == policy_engine.PolicyResultAllow {
|
||||
explicitAllow = true
|
||||
}
|
||||
}
|
||||
|
||||
return explicitAllow
|
||||
}
|
||||
|
||||
// VerifyActionPermission checks if the identity is allowed to perform the action on the resource.
|
||||
// It handles both traditional identities (via Actions) and IAM/STS identities (via Policy).
|
||||
func (iam *IdentityAccessManagement) VerifyActionPermission(r *http.Request, identity *Identity, action Action, bucket, object string) s3err.ErrorCode {
|
||||
@@ -1679,6 +1731,7 @@ func (iam *IdentityAccessManagement) VerifyActionPermission(r *http.Request, ide
|
||||
return iam.authorizeWithIAM(r, identity, action, bucket, object)
|
||||
}
|
||||
|
||||
// Traditional actions-based authorization from static S3 config.
|
||||
if len(identity.Actions) > 0 {
|
||||
if !identity.CanDo(action, bucket, object) {
|
||||
return s3err.ErrAccessDenied
|
||||
@@ -1686,6 +1739,14 @@ func (iam *IdentityAccessManagement) VerifyActionPermission(r *http.Request, ide
|
||||
return s3err.ErrNone
|
||||
}
|
||||
|
||||
// IAM policy fallback for identities with attached policies but without IAM integration.
|
||||
if len(identity.PolicyNames) > 0 {
|
||||
if iam.evaluateIAMPolicies(r, identity, action, bucket, object) {
|
||||
return s3err.ErrNone
|
||||
}
|
||||
return s3err.ErrAccessDenied
|
||||
}
|
||||
|
||||
return s3err.ErrAccessDenied
|
||||
}
|
||||
|
||||
@@ -1752,6 +1813,12 @@ func (iam *IdentityAccessManagement) PutPolicy(name string, content string) erro
|
||||
iam.policies = make(map[string]*iam_pb.Policy)
|
||||
}
|
||||
iam.policies[name] = &iam_pb.Policy{Name: name, Content: content}
|
||||
iam.ensureIAMPolicyEngine()
|
||||
// Remove old entry first so that a parse failure doesn't leave a stale allow.
|
||||
_ = iam.iamPolicyEngine.DeleteBucketPolicy(name)
|
||||
if err := iam.iamPolicyEngine.SetBucketPolicy(name, content); err != nil {
|
||||
glog.Warningf("IAM policy %q is stored but could not be compiled for cache: %v", name, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -1770,9 +1837,36 @@ func (iam *IdentityAccessManagement) DeletePolicy(name string) error {
|
||||
iam.m.Lock()
|
||||
defer iam.m.Unlock()
|
||||
delete(iam.policies, name)
|
||||
if iam.iamPolicyEngine != nil {
|
||||
_ = iam.iamPolicyEngine.DeleteBucketPolicy(name)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// ensureIAMPolicyEngine lazily initializes the shared IAM policy engine.
|
||||
// Must be called with iam.m held.
|
||||
func (iam *IdentityAccessManagement) ensureIAMPolicyEngine() {
|
||||
if iam.iamPolicyEngine == nil {
|
||||
iam.iamPolicyEngine = policy_engine.NewPolicyEngine()
|
||||
}
|
||||
}
|
||||
|
||||
// rebuildIAMPolicyEngineLocked rebuilds the entire IAM policy engine cache
|
||||
// from the current policies map. Must be called with iam.m held.
|
||||
func (iam *IdentityAccessManagement) rebuildIAMPolicyEngineLocked() {
|
||||
if len(iam.policies) == 0 {
|
||||
iam.iamPolicyEngine = nil
|
||||
return
|
||||
}
|
||||
engine := policy_engine.NewPolicyEngine()
|
||||
for name, p := range iam.policies {
|
||||
if err := engine.SetBucketPolicy(name, p.Content); err != nil {
|
||||
glog.Warningf("IAM policy cache rebuild: skipping invalid policy %q: %v", name, err)
|
||||
}
|
||||
}
|
||||
iam.iamPolicyEngine = engine
|
||||
}
|
||||
|
||||
// ListPolicies lists all policies
|
||||
func (iam *IdentityAccessManagement) ListPolicies() []*iam_pb.Policy {
|
||||
iam.m.RLock()
|
||||
|
||||
@@ -1,7 +1,9 @@
|
||||
package s3api
|
||||
|
||||
import (
|
||||
"crypto/tls"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"os"
|
||||
"reflect"
|
||||
"sync"
|
||||
@@ -9,6 +11,7 @@ import (
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/credential"
|
||||
. "github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3err"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util/wildcard"
|
||||
"github.com/stretchr/testify/assert"
|
||||
|
||||
@@ -260,6 +263,152 @@ func TestMatchWildcardPattern(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestVerifyActionPermissionPolicyFallback(t *testing.T) {
|
||||
buildRequest := func(t *testing.T, method string) *http.Request {
|
||||
t.Helper()
|
||||
req, err := http.NewRequest(method, "http://s3.amazonaws.com/test-bucket/test-object", nil)
|
||||
assert.NoError(t, err)
|
||||
return req
|
||||
}
|
||||
|
||||
t.Run("policy allow grants access", func(t *testing.T) {
|
||||
iam := &IdentityAccessManagement{}
|
||||
err := iam.PutPolicy("allowGet", `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Action":"s3:GetObject","Resource":"arn:aws:s3:::test-bucket/*"}]}`)
|
||||
assert.NoError(t, err)
|
||||
|
||||
identity := &Identity{
|
||||
Name: "policy-user",
|
||||
Account: &AccountAdmin,
|
||||
PolicyNames: []string{"allowGet"},
|
||||
}
|
||||
|
||||
errCode := iam.VerifyActionPermission(buildRequest(t, http.MethodGet), identity, Action(ACTION_READ), "test-bucket", "test-object")
|
||||
assert.Equal(t, s3err.ErrNone, errCode)
|
||||
})
|
||||
|
||||
t.Run("explicit deny overrides allow", func(t *testing.T) {
|
||||
iam := &IdentityAccessManagement{}
|
||||
err := iam.PutPolicy("allowAllGet", `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Action":"s3:GetObject","Resource":"arn:aws:s3:::test-bucket/*"}]}`)
|
||||
assert.NoError(t, err)
|
||||
err = iam.PutPolicy("denySecret", `{"Version":"2012-10-17","Statement":[{"Effect":"Deny","Action":"s3:GetObject","Resource":"arn:aws:s3:::test-bucket/secret.txt"}]}`)
|
||||
assert.NoError(t, err)
|
||||
|
||||
identity := &Identity{
|
||||
Name: "policy-user",
|
||||
Account: &AccountAdmin,
|
||||
PolicyNames: []string{"allowAllGet", "denySecret"},
|
||||
}
|
||||
|
||||
errCode := iam.VerifyActionPermission(buildRequest(t, http.MethodGet), identity, Action(ACTION_READ), "test-bucket", "secret.txt")
|
||||
assert.Equal(t, s3err.ErrAccessDenied, errCode)
|
||||
})
|
||||
|
||||
t.Run("implicit deny when no statement matches", func(t *testing.T) {
|
||||
iam := &IdentityAccessManagement{}
|
||||
err := iam.PutPolicy("allowOtherBucket", `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Action":"s3:GetObject","Resource":"arn:aws:s3:::other-bucket/*"}]}`)
|
||||
assert.NoError(t, err)
|
||||
|
||||
identity := &Identity{
|
||||
Name: "policy-user",
|
||||
Account: &AccountAdmin,
|
||||
PolicyNames: []string{"allowOtherBucket"},
|
||||
}
|
||||
|
||||
errCode := iam.VerifyActionPermission(buildRequest(t, http.MethodGet), identity, Action(ACTION_READ), "test-bucket", "test-object")
|
||||
assert.Equal(t, s3err.ErrAccessDenied, errCode)
|
||||
})
|
||||
|
||||
t.Run("invalid policy document does not allow", func(t *testing.T) {
|
||||
iam := &IdentityAccessManagement{}
|
||||
err := iam.PutPolicy("invalidPolicy", "{not-json")
|
||||
assert.NoError(t, err)
|
||||
|
||||
identity := &Identity{
|
||||
Name: "policy-user",
|
||||
Account: &AccountAdmin,
|
||||
PolicyNames: []string{"invalidPolicy"},
|
||||
}
|
||||
|
||||
errCode := iam.VerifyActionPermission(buildRequest(t, http.MethodGet), identity, Action(ACTION_READ), "test-bucket", "test-object")
|
||||
assert.Equal(t, s3err.ErrAccessDenied, errCode)
|
||||
})
|
||||
|
||||
t.Run("notresource excludes denied object", func(t *testing.T) {
|
||||
iam := &IdentityAccessManagement{}
|
||||
err := iam.PutPolicy("denyNotResource", `{"Version":"2012-10-17","Statement":[{"Effect":"Deny","Action":"s3:GetObject","NotResource":"arn:aws:s3:::test-bucket/public/*"}]}`)
|
||||
assert.NoError(t, err)
|
||||
err = iam.PutPolicy("allowAllGet", `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Action":"s3:GetObject","Resource":"arn:aws:s3:::test-bucket/*"}]}`)
|
||||
assert.NoError(t, err)
|
||||
|
||||
identity := &Identity{
|
||||
Name: "policy-user",
|
||||
Account: &AccountAdmin,
|
||||
PolicyNames: []string{"allowAllGet", "denyNotResource"},
|
||||
}
|
||||
|
||||
errCode := iam.VerifyActionPermission(buildRequest(t, http.MethodGet), identity, Action(ACTION_READ), "test-bucket", "private/secret.txt")
|
||||
assert.Equal(t, s3err.ErrAccessDenied, errCode)
|
||||
|
||||
errCode = iam.VerifyActionPermission(buildRequest(t, http.MethodGet), identity, Action(ACTION_READ), "test-bucket", "public/readme.txt")
|
||||
assert.Equal(t, s3err.ErrNone, errCode)
|
||||
})
|
||||
|
||||
t.Run("condition securetransport enforced", func(t *testing.T) {
|
||||
iam := &IdentityAccessManagement{}
|
||||
err := iam.PutPolicy("allowTLSOnly", `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Action":"s3:GetObject","Resource":"arn:aws:s3:::test-bucket/*","Condition":{"Bool":{"aws:SecureTransport":"true"}}}]}`)
|
||||
assert.NoError(t, err)
|
||||
|
||||
identity := &Identity{
|
||||
Name: "policy-user",
|
||||
Account: &AccountAdmin,
|
||||
PolicyNames: []string{"allowTLSOnly"},
|
||||
}
|
||||
|
||||
httpReq := buildRequest(t, http.MethodGet)
|
||||
errCode := iam.VerifyActionPermission(httpReq, identity, Action(ACTION_READ), "test-bucket", "test-object")
|
||||
assert.Equal(t, s3err.ErrAccessDenied, errCode)
|
||||
|
||||
httpsReq := buildRequest(t, http.MethodGet)
|
||||
httpsReq.TLS = &tls.ConnectionState{}
|
||||
errCode = iam.VerifyActionPermission(httpsReq, identity, Action(ACTION_READ), "test-bucket", "test-object")
|
||||
assert.Equal(t, s3err.ErrNone, errCode)
|
||||
})
|
||||
|
||||
t.Run("valid policy updated to invalid denies access", func(t *testing.T) {
|
||||
iam := &IdentityAccessManagement{}
|
||||
err := iam.PutPolicy("myPolicy", `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Action":"s3:GetObject","Resource":"arn:aws:s3:::test-bucket/*"}]}`)
|
||||
assert.NoError(t, err)
|
||||
|
||||
identity := &Identity{
|
||||
Name: "policy-user",
|
||||
Account: &AccountAdmin,
|
||||
PolicyNames: []string{"myPolicy"},
|
||||
}
|
||||
|
||||
errCode := iam.VerifyActionPermission(buildRequest(t, http.MethodGet), identity, Action(ACTION_READ), "test-bucket", "test-object")
|
||||
assert.Equal(t, s3err.ErrNone, errCode)
|
||||
|
||||
// Update to invalid JSON — should revoke access.
|
||||
err = iam.PutPolicy("myPolicy", "{broken")
|
||||
assert.NoError(t, err)
|
||||
|
||||
errCode = iam.VerifyActionPermission(buildRequest(t, http.MethodGet), identity, Action(ACTION_READ), "test-bucket", "test-object")
|
||||
assert.Equal(t, s3err.ErrAccessDenied, errCode)
|
||||
})
|
||||
|
||||
t.Run("actions based path still works", func(t *testing.T) {
|
||||
iam := &IdentityAccessManagement{}
|
||||
identity := &Identity{
|
||||
Name: "legacy-user",
|
||||
Account: &AccountAdmin,
|
||||
Actions: []Action{"Read:test-bucket"},
|
||||
}
|
||||
|
||||
errCode := iam.VerifyActionPermission(buildRequest(t, http.MethodGet), identity, Action(ACTION_READ), "test-bucket", "any-object")
|
||||
assert.Equal(t, s3err.ErrNone, errCode)
|
||||
})
|
||||
}
|
||||
|
||||
type LoadS3ApiConfigurationTestCase struct {
|
||||
pbAccount *iam_pb.Account
|
||||
pbIdent *iam_pb.Identity
|
||||
|
||||
@@ -113,7 +113,7 @@ func (ms *MasterServer) SendHeartbeat(stream master_pb.Seaweed_SendHeartbeatServ
|
||||
// tell the volume servers about the leader
|
||||
newLeader, err := ms.Topo.MaybeLeader()
|
||||
if err != nil || newLeader == "" {
|
||||
glog.Warningf("SendHeartbeat find leader: %v", err)
|
||||
glog.Warningf("SendHeartbeat find leader: %v, %v", newLeader, err)
|
||||
return raft.NotLeaderError
|
||||
}
|
||||
if err := stream.Send(&master_pb.HeartbeatResponse{
|
||||
|
||||
Reference in New Issue
Block a user