Compare commits

...
Author SHA1 Message Date
Chris Lu da39ee65a3 mount: add unit tests for refresh event buffering
Test all key scenarios for the cache refresh mechanism:

- Delete during refresh: subscription delete is buffered and replayed,
  preventing ghost entries from the stale snapshot
- Create during refresh: subscription create is buffered and replayed,
  preventing lost entries
- Stale entries removed: entries absent from the filer snapshot are
  cleaned from the cache
- Local delete buffered: direct DeleteEntry during refresh is both
  applied immediately and replayed after commit
- Local create buffered: direct InsertEntry during refresh is preserved
  after snapshot commit
- Cancel discards buffer: CancelRefresh drops buffered events and
  resumes normal event processing
- Concurrent deletes (issue #8442 scenario): 1000 concurrent delete
  events racing with a stale snapshot all resolve correctly
2026-03-06 00:47:37 -08:00
Chris Lu e91d43ef08 mount: use BeginRefresh/CommitRefresh in doEnsureVisited
Replace the incremental batch-insert approach in doEnsureVisited with
the new refresh lifecycle:

1. BeginRefresh - start buffering subscription events
2. ReadDirAllEntries - fetch full listing (no lock held)
3. CommitRefresh - atomically replace cache + replay buffered events

This ensures that any creates, deletes, or updates that arrive via the
subscription handler during the filer listing are not lost. The snapshot
replaces all stale entries, and buffered events are replayed on top to
bring the cache up to date.
2026-03-06 00:46:10 -08:00
Chris Lu 58cf7e9a89 mount: add refresh event buffering to MetaCache
Add infrastructure to buffer subscription events during directory cache
refresh. When a directory is being refreshed from the filer (BeginRefresh),
events from the subscription handler, local deletes, and local creates
are buffered. When the refresh completes (CommitRefresh), the filer
snapshot atomically replaces the cached entries, then buffered events
are replayed to ensure no mutations are lost.

This fixes a race condition where concurrent deletes or creates could be
overwritten by a stale filer snapshot during directory refresh, causing
ghost entries or lost files.

Fixes #8442
2026-03-06 00:45:23 -08:00
Chris Lu 338be16254 fix logs 2026-03-05 15:38:05 -08:00
Chris Lu 1b6e96614d s3api: cache parsed IAM policy engines for fallback auth
Previously, evaluateIAMPolicies created a new PolicyEngine and re-parsed
the JSON policy document for every policy on every request. This adds a
shared iamPolicyEngine field that caches compiled policies, kept in sync
by PutPolicy, DeletePolicy, and bulk config reload paths.

- PutPolicy deletes the old cache entry before setting the new one, so a
  parse failure on update does not leave a stale allow.
- Log warnings when policy compilation fails instead of silently
  discarding errors.
- Add test for valid-to-invalid policy update regression.
2026-03-05 14:27:48 -08:00
SrikanthBhandaryandGitHub 4eb45ecc5e s3api: add IAM policy fallback authorization tests (#8518)
* s3api: add IAM policy fallback auth with tests

* s3api: use policy engine for IAM fallback evaluation
2026-03-05 12:13:18 -08:00
Chris LuandGitHub 1f3df6e9ef admin: remove Alpha badge and unused Metrics/Logs menu items (#8525)
* admin: remove Alpha badge and unused Metrics/Logs menu items

* Update layout_templ.go
2026-03-05 11:51:11 -08:00
fcd5de9710 Fix YAML parse error in post-install-bucket-hook template (#8523)
The 'set -o pipefail' line was improperly indented outside the YAML block
scalar, causing a parse error when s3.enabled=true and s3.createBuckets
were populated. Moved the line to the beginning of the script block with
correct indentation (12 spaces).

Fixes #8520

Co-authored-by: Claude Haiku 4.5 <noreply@anthropic.com>
2026-03-05 10:01:31 -08:00
Steven CrespoandGitHub b6f6f0187e Add before-hook-creation delete policy to bucket-hook Job (#8519) 2026-03-05 06:28:03 -08:00
9 changed files with 717 additions and 63 deletions
@@ -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; }
-14
View File
@@ -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">
+11 -11
View File
@@ -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 {
+104 -8
View File
@@ -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()
+15 -26
View File
@@ -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
+340
View File
@@ -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))])
}
}
+94
View File
@@ -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()
+149
View File
@@ -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
+1 -1
View File
@@ -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{