diff --git a/docs/coverage/aws/s3.md b/docs/coverage/aws/s3.md
index 614ae2848..95dafa0a9 100644
--- a/docs/coverage/aws/s3.md
+++ b/docs/coverage/aws/s3.md
@@ -194,4 +194,5 @@ VersionedBucket is an optional extension a storage provider implements when
## Not in scope
-_Not documented yet. See the [emulator boundary](../../../README.md) for cloudemu-wide non-goals._
+- Azure: on the bare emulator host (such as `https://127.0.0.1:4568/`), Blob, Queue and Table share one endpoint, and the account-root calls (`GET /?comp=list`, `?restype=service`, `?restype=account`) look the same for each service. They go to Blob unless the User-Agent carries an Azure SDK Queue or Table product token (`azsdk-go-azqueue`, `azsdk-python-storage-queue`, `azsdk-java-azure-storage-queue`, `azsdk-js-storage-queue`, `azsdk-net-Storage.Queues`, and the matching `data-tables` clients). Other Queue and Table clients should use the `{account}.queue.core.windows.net` or `{account}.table.core.windows.net` host, or the path-style `/{account}/` form.
+- Azure: Queue and Table Set Service Properties validate the request and return 202, but do not store logging, metrics or CORS settings. Get Service Stats and Get User Delegation Key are not served.
diff --git a/docs/coverage/azure/blobstorage.md b/docs/coverage/azure/blobstorage.md
index 444f815c7..cfa31d3e7 100644
--- a/docs/coverage/azure/blobstorage.md
+++ b/docs/coverage/azure/blobstorage.md
@@ -199,4 +199,5 @@ StorageAccountKeys is an OPTIONAL Azure-specific capability, discovered by
## Not in scope
-_Not documented yet. See the [emulator boundary](../../../README.md) for cloudemu-wide non-goals._
+- Azure: on the bare emulator host (such as `https://127.0.0.1:4568/`), Blob, Queue and Table share one endpoint, and the account-root calls (`GET /?comp=list`, `?restype=service`, `?restype=account`) look the same for each service. They go to Blob unless the User-Agent carries an Azure SDK Queue or Table product token (`azsdk-go-azqueue`, `azsdk-python-storage-queue`, `azsdk-java-azure-storage-queue`, `azsdk-js-storage-queue`, `azsdk-net-Storage.Queues`, and the matching `data-tables` clients). Other Queue and Table clients should use the `{account}.queue.core.windows.net` or `{account}.table.core.windows.net` host, or the path-style `/{account}/` form.
+- Azure: Queue and Table Set Service Properties validate the request and return 202, but do not store logging, metrics or CORS settings. Get Service Stats and Get User Delegation Key are not served.
diff --git a/docs/coverage/gcp/gcs.md b/docs/coverage/gcp/gcs.md
index ca93b8b8c..43da702b0 100644
--- a/docs/coverage/gcp/gcs.md
+++ b/docs/coverage/gcp/gcs.md
@@ -72,4 +72,5 @@ GCSExtensions is an OPTIONAL GCS-specific capability, discovered by type
## Not in scope
-_Not documented yet. See the [emulator boundary](../../../README.md) for cloudemu-wide non-goals._
+- Azure: on the bare emulator host (such as `https://127.0.0.1:4568/`), Blob, Queue and Table share one endpoint, and the account-root calls (`GET /?comp=list`, `?restype=service`, `?restype=account`) look the same for each service. They go to Blob unless the User-Agent carries an Azure SDK Queue or Table product token (`azsdk-go-azqueue`, `azsdk-python-storage-queue`, `azsdk-java-azure-storage-queue`, `azsdk-js-storage-queue`, `azsdk-net-Storage.Queues`, and the matching `data-tables` clients). Other Queue and Table clients should use the `{account}.queue.core.windows.net` or `{account}.table.core.windows.net` host, or the path-style `/{account}/` form.
+- Azure: Queue and Table Set Service Properties validate the request and return 202, but do not store logging, metrics or CORS settings. Get Service Stats and Get User Delegation Key are not served.
diff --git a/docs/coverage/nongoals/storage.md b/docs/coverage/nongoals/storage.md
new file mode 100644
index 000000000..f24060a54
--- /dev/null
+++ b/docs/coverage/nongoals/storage.md
@@ -0,0 +1,2 @@
+- Azure: on the bare emulator host (such as `https://127.0.0.1:4568/`), Blob, Queue and Table share one endpoint, and the account-root calls (`GET /?comp=list`, `?restype=service`, `?restype=account`) look the same for each service. They go to Blob unless the User-Agent carries an Azure SDK Queue or Table product token (`azsdk-go-azqueue`, `azsdk-python-storage-queue`, `azsdk-java-azure-storage-queue`, `azsdk-js-storage-queue`, `azsdk-net-Storage.Queues`, and the matching `data-tables` clients). Other Queue and Table clients should use the `{account}.queue.core.windows.net` or `{account}.table.core.windows.net` host, or the path-style `/{account}/` form.
+- Azure: Queue and Table Set Service Properties validate the request and return 202, but do not store logging, metrics or CORS settings. Get Service Stats and Get User Delegation Key are not served.
diff --git a/providers/azure/eventgrid/delivery.go b/providers/azure/eventgrid/delivery.go
index 9e920441a..c2c6c70d7 100644
--- a/providers/azure/eventgrid/delivery.go
+++ b/providers/azure/eventgrid/delivery.go
@@ -217,9 +217,19 @@ func (m *Mock) dispatchStorageQueue(ctx context.Context, dest subscriptionDestin
return
}
- if dest.QueueName != "" {
- _ = m.storageQueue.DeliverExternal(ctx, dest.QueueName, string(body))
+ if dest.QueueName == "" {
+ return
}
+
+ // A queue in a storage account is keyed "{account}/{queue}" (the default
+ // account's queues by bare name), so try the account named by resourceId
+ // first and fall back to the default account.
+ if acct := resourceLeafName(dest.ResourceID); acct != "" &&
+ m.storageQueue.DeliverExternal(ctx, acct+"/"+dest.QueueName, string(body)) == nil {
+ return
+ }
+
+ _ = m.storageQueue.DeliverExternal(ctx, dest.QueueName, string(body))
}
// resourceLeafName returns the trailing path segment of an ARM resource id, the
diff --git a/providers/azure/servicebus/servicebus.go b/providers/azure/servicebus/servicebus.go
index 1dfa77913..3d41c18ac 100644
--- a/providers/azure/servicebus/servicebus.go
+++ b/providers/azure/servicebus/servicebus.go
@@ -527,6 +527,9 @@ func buildSendMessage(input *driver.SendMessageInput, sessionID string, now time
}
visibleAt := now.Add(time.Duration(delaySeconds) * time.Second)
+ if input.ScheduledEnqueueTime.After(visibleAt) {
+ visibleAt = input.ScheduledEnqueueTime
+ }
return &sbMessage{
ID: idgen.GenerateID("sb-msg-"),
diff --git a/server/azure/azure.go b/server/azure/azure.go
index 51b03a406..9363a7cd5 100644
--- a/server/azure/azure.go
+++ b/server/azure/azure.go
@@ -1237,8 +1237,12 @@ func New(d Drivers) http.Handler {
// that contain parentheses or a bare JSON POST, disjoint from Blob's
// container/blob paths and Queue's /messages surface. Registered before the
// permissive Blob fallback.
+ // Storage accounts scope the queue and table namespaces the same way they
+ // scope blob containers.
+ storageAccounts, _ := d.BlobStorage.(storagedriver.AzureStorageAccounts)
+
if d.TableStorage != nil {
- srv.Register(tablesrv.New(d.TableStorage))
+ srv.Register(tablesrv.New(d.TableStorage).WithAccounts(storageAccounts))
}
// Queue Storage matches the queue data-plane surface (/{queue}/messages,
@@ -1247,7 +1251,7 @@ func New(d Drivers) http.Handler {
// carries OData parentheses). Registered before the permissive Blob
// fallback.
if d.QueueStorage != nil {
- srv.Register(queue.New(d.QueueStorage))
+ srv.Register(queue.New(d.QueueStorage).WithAccounts(storageAccounts))
}
// Storage-account ARM control plane (Microsoft.Storage/storageAccounts).
diff --git a/server/azure/blobstorage/handler.go b/server/azure/blobstorage/handler.go
index 37b9a60b1..09fd1c7c9 100644
--- a/server/azure/blobstorage/handler.go
+++ b/server/azure/blobstorage/handler.go
@@ -131,6 +131,8 @@ func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
h.listContainers(w, r, account)
case container == "" && q.Get("comp") == compBlobs:
h.findBlobsByTags(w, r, account, "")
+ case container == "" && azurearm.IsStorageServiceOp(q):
+ h.serviceOp(w, r, account)
case container == "":
writeError(w, http.StatusNotImplemented, "NotImplemented", "operation not supported on root")
case blob == "" && q.Get("restype") == "container":
diff --git a/server/azure/blobstorage/service_properties.go b/server/azure/blobstorage/service_properties.go
new file mode 100644
index 000000000..3d523e88f
--- /dev/null
+++ b/server/azure/blobstorage/service_properties.go
@@ -0,0 +1,116 @@
+package blobstorage
+
+import (
+ "net/http"
+ "strings"
+
+ "github.com/stackshy/cloudemu/v2/server/wire/azurearm"
+ storagedriver "github.com/stackshy/cloudemu/v2/services/storage/driver"
+)
+
+// serviceOp serves the account-level Blob service operations: Get/Set Blob
+// Service Properties (?restype=service&comp=properties) and Get Account
+// Information (?restype=account&comp=properties). The delete retention policy
+// and CORS rules share their state with the ARM blobServices/default resource,
+// so a value set on either surface reads back on the other.
+func (h *Handler) serviceOp(w http.ResponseWriter, r *http.Request, account string) {
+ cfg, ok := h.bucket.(storagedriver.BlobServiceConfig)
+ if !ok || !azurearm.IsStorageServicePropertiesOp(r.URL.Query()) {
+ azurearm.ServeStorageServiceOp(w, r)
+ return
+ }
+
+ if account == "" {
+ account = storagedriver.AzureDefaultStorageAccount
+ }
+
+ current, err := cfg.BlobServiceProperties(r.Context(), account)
+ if err != nil {
+ writeErr(w, err)
+ return
+ }
+
+ switch r.Method {
+ case http.MethodGet:
+ props := toWireServiceProperties(¤t)
+ azurearm.WriteStorageServiceProperties(w, &props)
+ case http.MethodPut:
+ in, ok := azurearm.DecodeStorageServiceProperties(w, r)
+ if !ok {
+ return
+ }
+
+ mergeServiceProperties(¤t, &in)
+
+ if err := cfg.SetBlobServiceProperties(r.Context(), account, current); err != nil {
+ writeErr(w, err)
+ return
+ }
+
+ w.WriteHeader(http.StatusAccepted)
+ default:
+ writeError(w, http.StatusMethodNotAllowed, "UnsupportedHttpVerb", "method not allowed")
+ }
+}
+
+// toWireServiceProperties renders the stored Blob service properties as the
+// data-plane document, with the defaults for what cloudemu does not store.
+func toWireServiceProperties(p *storagedriver.BlobServiceProperties) azurearm.StorageServiceProperties {
+ out := azurearm.DefaultStorageServiceProperties()
+ out.DeleteRetentionPolicy = &azurearm.StorageRetentionPolicy{
+ Enabled: p.DeleteRetentionEnabled,
+ Days: p.DeleteRetentionDays,
+ }
+
+ for _, c := range p.CORS {
+ out.Cors.Rules = append(out.Cors.Rules, azurearm.StorageCorsRule{
+ AllowedOrigins: strings.Join(c.AllowedOrigins, ","),
+ AllowedMethods: strings.Join(c.AllowedMethods, ","),
+ AllowedHeaders: strings.Join(c.AllowedHeaders, ","),
+ ExposedHeaders: strings.Join(c.ExposeHeaders, ","),
+ MaxAgeInSeconds: c.MaxAgeSeconds,
+ })
+ }
+
+ return out
+}
+
+// mergeServiceProperties applies a Set Blob Service Properties body. As in
+// real Azure, an element the request omits keeps its current value.
+func mergeServiceProperties(p *storagedriver.BlobServiceProperties, in *azurearm.StorageServiceProperties) {
+ if in.DeleteRetentionPolicy != nil {
+ p.DeleteRetentionEnabled = in.DeleteRetentionPolicy.Enabled
+ p.DeleteRetentionDays = 0
+
+ if in.DeleteRetentionPolicy.Enabled {
+ p.DeleteRetentionDays = in.DeleteRetentionPolicy.Days
+ }
+ }
+
+ if in.Cors != nil {
+ p.CORS = make([]storagedriver.CORSRule, 0, len(in.Cors.Rules))
+
+ for _, c := range in.Cors.Rules {
+ p.CORS = append(p.CORS, storagedriver.CORSRule{
+ AllowedOrigins: splitList(c.AllowedOrigins),
+ AllowedMethods: splitList(c.AllowedMethods),
+ AllowedHeaders: splitList(c.AllowedHeaders),
+ ExposeHeaders: splitList(c.ExposedHeaders),
+ MaxAgeSeconds: c.MaxAgeInSeconds,
+ })
+ }
+ }
+}
+
+// splitList splits a comma-separated CORS list, dropping empty entries.
+func splitList(s string) []string {
+ var out []string
+
+ for _, v := range strings.Split(s, ",") {
+ if v = strings.TrimSpace(v); v != "" {
+ out = append(out, v)
+ }
+ }
+
+ return out
+}
diff --git a/server/azure/cosmosdb/handler.go b/server/azure/cosmosdb/handler.go
index 47b623720..c57b2a6a9 100644
--- a/server/azure/cosmosdb/handler.go
+++ b/server/azure/cosmosdb/handler.go
@@ -353,10 +353,12 @@ func (h *Handler) Matches(r *http.Request) bool {
return false
}
- // A root "GET /?comp=list" is a Storage service call (list queues / list
- // containers), not a Cosmos account probe (which carries no query). Decline
- // it so the Queue and Blob handlers, registered after this one, serve it.
- if rest == "/" && r.URL.Query().Get("comp") == "list" {
+ // A root request carrying comp= or restype= is a Storage account-level call
+ // (list queues or containers, service properties, account information,
+ // find blobs by tags), not a Cosmos account probe (which carries no query).
+ // Decline it so the Queue, Table and Blob handlers, registered after this
+ // one, serve it.
+ if q := r.URL.Query(); rest == "/" && (q.Has("comp") || q.Has("restype")) {
return false
}
diff --git a/server/azure/queue/account.go b/server/azure/queue/account.go
new file mode 100644
index 000000000..12c5c7754
--- /dev/null
+++ b/server/azure/queue/account.go
@@ -0,0 +1,102 @@
+package queue
+
+import (
+ "net/http"
+ "slices"
+ "strings"
+
+ "github.com/stackshy/cloudemu/v2/server/wire/azurearm"
+ storagedriver "github.com/stackshy/cloudemu/v2/services/storage/driver"
+)
+
+// WithAccounts scopes queues to the storage accounts in accounts: a queue
+// created through {account}.queue.core.windows.net (or the path-style
+// /{account}/ prefix) lives only in that account. Without it every request
+// uses the default account.
+func (h *Handler) WithAccounts(accounts storagedriver.AzureStorageAccounts) *Handler {
+ h.accounts = accounts
+ return h
+}
+
+// resolve returns the account a request targets and the path below it. The
+// {account}.queue host names the account when it is an existing storage
+// account; otherwise a path-style /{account}/ prefix does. The default account
+// "cloudemu" and the bare host map to the default namespace (""), which also
+// holds every queue from before accounts were modeled.
+func (h *Handler) resolve(r *http.Request) (account, path string) {
+ if acct, svc, ok := azurearm.StorageHost(r.Host); ok && svc == azurearm.StorageServiceQueue && h.accountExists(r, acct) {
+ return acct, r.URL.Path
+ }
+
+ q := r.URL.Query()
+ rootOp := q.Get("comp") == compList || azurearm.IsStorageServiceOp(q)
+
+ return azurearm.PeelStorageAccount(r.URL.Path, rootOp, func(name string) bool {
+ // A default-namespace queue of the same name keeps its URL.
+ return h.accountExists(r, name) && !h.queueExists(r, name)
+ })
+}
+
+func (h *Handler) accountExists(r *http.Request, name string) bool {
+ if h.accounts == nil || name == storagedriver.AzureDefaultStorageAccount {
+ return false
+ }
+
+ _, err := h.accounts.GetStorageAccount(r.Context(), name)
+
+ return err == nil
+}
+
+func (h *Handler) queueExists(r *http.Request, key string) bool {
+ _, err := h.resolveQueueURL(r, key)
+ return err == nil
+}
+
+// queueKey is the driver name of queue in account: the bare name in the
+// default account, "{account}/{queue}" in any other. Queue names cannot hold
+// "/", so keys never collide.
+func queueKey(account, queue string) string {
+ return storagedriver.AzureContainerKey(account, queue)
+}
+
+// inAccount reports whether the driver queue key belongs to account and, if
+// so, the queue's own name.
+func inAccount(key, account string) (string, bool) {
+ acct, name := storagedriver.SplitAzureContainerKey(key)
+ if acct == storagedriver.AzureDefaultStorageAccount {
+ acct = ""
+ }
+
+ return name, acct == account
+}
+
+// queueSDKProducts are the User-Agent product names of the Azure SDK Queue
+// clients, lowercased.
+//
+//nolint:gochecknoglobals // read-only lookup table, not mutable state
+var queueSDKProducts = []string{
+ "azsdk-go-azqueue",
+ "azsdk-python-storage-queue",
+ "azsdk-java-azure-storage-queue",
+ "azsdk-js-storage-queue",
+ "azsdk-net-storage.queues",
+}
+
+// isQueueClient reports whether a request on a host that does not name its
+// service comes from an Azure SDK Queue client. List Queues and the
+// account-level service calls have the same shape as their Blob counterparts,
+// so on a bare host only a product token of the User-Agent (such as
+// "azsdk-go-azqueue/v1.0.0") picks Queue. Free text such as an application id
+// is ignored, and anything else goes to Blob. Other Queue clients should use
+// the {account}.queue host or the path-style /{account}/ form.
+func isQueueClient(r *http.Request) bool {
+ for _, token := range strings.Fields(strings.ToLower(r.UserAgent())) {
+ product, _, _ := strings.Cut(token, "/")
+
+ if slices.Contains(queueSDKProducts, product) {
+ return true
+ }
+ }
+
+ return false
+}
diff --git a/server/azure/queue/handler.go b/server/azure/queue/handler.go
index a6e29b840..b32351eb2 100644
--- a/server/azure/queue/handler.go
+++ b/server/azure/queue/handler.go
@@ -32,6 +32,7 @@ import (
cerrors "github.com/stackshy/cloudemu/v2/errors"
"github.com/stackshy/cloudemu/v2/server/wire/azurearm"
mqdriver "github.com/stackshy/cloudemu/v2/services/messagequeue/driver"
+ storagedriver "github.com/stackshy/cloudemu/v2/services/storage/driver"
)
const (
@@ -43,6 +44,9 @@ const (
compList = "list"
compMetadata = "metadata"
+ // segMessages is the path segment of the queue message surface.
+ segMessages = "messages"
+
// maxUpdateVisibilityTimeout is the Azure ceiling for a message's visibility
// timeout (7 days, in seconds).
maxUpdateVisibilityTimeout = 7 * 24 * 60 * 60
@@ -79,7 +83,8 @@ const (
// Handler serves Azure Queue Storage REST requests against a messagequeue
// driver.
type Handler struct {
- mq mqdriver.MessageQueue
+ mq mqdriver.MessageQueue
+ accounts storagedriver.AzureStorageAccounts
}
// New returns a Queue handler backed by mq.
@@ -101,52 +106,68 @@ func New(mq mqdriver.MessageQueue) *Handler {
// - PUT|DELETE /{queue} with no restype=container query: Blob container ops
// always carry restype=container, so a bare PUT/DELETE on a single path
// segment is a queue create/delete. Disjoint from Blob container ops.
-// - GET /?comp=list (list queues): this shape is byte-for-byte identical to
-// Blob's list-containers; Azure disambiguates only by hostname. When both
-// handlers are registered, the Queue handler (registered first) owns it.
+// - GET /?comp=list (list queues) and the account-level service calls
+// (?restype=service|account): these shapes are byte-for-byte identical to
+// Blob's, and Azure tells them apart by hostname. The Queue handler claims
+// them on an {account}.queue host, and on a bare host only for a Queue
+// client (see isQueueClient). Every other root request goes to Blob.
//
-// The two shared-hostname shapes above are inherent to serving Queue and Blob
-// on one endpoint (real Azure uses distinct hostnames); documented, not fixable
-// without host-based routing.
+// A path-style /{account}/ prefix naming an existing storage account is
+// peeled before the shape checks (see resolve).
//
// Registered before the permissive Blob fallback so these shapes win.
-func (*Handler) Matches(r *http.Request) bool {
+func (h *Handler) Matches(r *http.Request) bool {
if strings.HasPrefix(r.URL.Path, "/subscriptions/") {
return false
}
// A storage host names its service, so a request to another service's
// host (such as {account}.blob.core.windows.net) is never a Queue call.
- if _, svc, ok := azurearm.StorageHost(r.Host); ok && svc != "queue" {
+ _, svc, storageHost := azurearm.StorageHost(r.Host)
+ if storageHost && svc != azurearm.StorageServiceQueue {
return false
}
- queue, sub, msgID := parseQueuePath(r.URL.Path)
+ _, path := h.resolve(r)
+ queue, sub, _ := parseQueuePath(path)
q := r.URL.Query()
- // /{queue}/messages[/{id}]: the unambiguous queue message surface.
- if sub == "messages" {
- _ = msgID
-
+ switch {
+ case sub == segMessages:
+ // /{queue}/messages[/{id}]: the unambiguous queue message surface.
return true
+ case queue == "":
+ return matchesRootOp(r, storageHost)
+ case sub != "" || q.Get("restype") != "" || r.Header.Get("X-Ms-Blob-Type") != "":
+ // Blob container ops carry restype=container and Put Blob carries
+ // x-ms-blob-type.
+ return false
}
- // GET /?comp=list: list queues (see Matches doc: shares Blob's shape).
- if queue == "" {
- return r.Method == http.MethodGet && q.Get("comp") == compList && q.Get("restype") == ""
- }
+ return matchesQueueOp(r.Method, q.Get("comp"))
+}
- // Bare /{queue} create/delete: PUT or DELETE with no container/blob query
- // markers. Blob container ops carry restype=container.
- if sub == "" {
- switch r.Method {
- case http.MethodPut, http.MethodDelete:
- return q.Get("restype") == ""
- case http.MethodGet, http.MethodHead:
- // GET|HEAD /{queue}?comp=metadata: queue properties. Blob container
- // metadata carries restype=container, so this is unambiguous.
- return q.Get("comp") == compMetadata && q.Get("restype") == ""
- }
+// matchesRootOp reports whether an account-root request is a Queue call: List
+// Queues or an account-level service call, on a queue host or from a Queue
+// client.
+func matchesRootOp(r *http.Request, queueHost bool) bool {
+ q := r.URL.Query()
+ listQueues := r.Method == http.MethodGet && q.Get("comp") == compList && q.Get("restype") == ""
+
+ return (listQueues || azurearm.IsStorageServiceOp(q)) && (queueHost || isQueueClient(r))
+}
+
+// matchesQueueOp reports whether a bare /{queue} request is a queue create,
+// delete or metadata call.
+func matchesQueueOp(method, comp string) bool {
+ switch method {
+ case http.MethodPut:
+ return comp == "" || comp == compMetadata
+ case http.MethodDelete:
+ return comp == ""
+ case http.MethodGet, http.MethodHead:
+ // GET|HEAD /{queue}?comp=metadata: queue properties.
+ return comp == compMetadata
}
return false
@@ -156,17 +177,24 @@ func (*Handler) Matches(r *http.Request) bool {
func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
w.Header().Set("X-Ms-Version", xmsVersion)
- queue, sub, msgID := parseQueuePath(r.URL.Path)
+ account, path := h.resolve(r)
+ queue, sub, msgID := parseQueuePath(path)
q := r.URL.Query()
+ if queue != "" {
+ queue = queueKey(account, queue)
+ }
+
switch {
case queue == "" && r.Method == http.MethodGet && q.Get("comp") == compList:
- h.listQueues(w, r)
+ h.listQueues(w, r, account)
+ case queue == "" && azurearm.IsStorageServiceOp(q):
+ azurearm.ServeStorageServiceOp(w, r)
case queue == "":
writeError(w, http.StatusNotImplemented, "NotImplemented", "operation not supported on root")
- case sub == "messages" && msgID == "":
+ case sub == segMessages && msgID == "":
h.messagesOp(w, r, queue)
- case sub == "messages":
+ case sub == segMessages:
h.messageIDOp(w, r, queue, msgID)
case sub == "":
h.queueOp(w, r, queue)
@@ -327,18 +355,21 @@ func (h *Handler) deleteQueue(w http.ResponseWriter, r *http.Request, queue stri
w.WriteHeader(http.StatusNoContent)
}
-func (h *Handler) listQueues(w http.ResponseWriter, r *http.Request) {
+func (h *Handler) listQueues(w http.ResponseWriter, r *http.Request, account string) {
prefix := r.URL.Query().Get("prefix")
- queues, err := h.mq.ListQueues(r.Context(), prefix)
+ queues, err := h.mq.ListQueues(r.Context(), "")
if err != nil {
writeErr(w, err)
return
}
out := listQueuesResult{Prefix: prefix}
+
for _, qi := range queues {
- out.Queues.Queues = append(out.Queues.Queues, queueXML{Name: qi.Name})
+ if name, ok := inAccount(qi.Name, account); ok && strings.HasPrefix(name, prefix) {
+ out.Queues.Queues = append(out.Queues.Queues, queueXML{Name: name})
+ }
}
writeXML(w, http.StatusOK, out)
diff --git a/server/azure/queue_servicebus_routing_test.go b/server/azure/queue_servicebus_routing_test.go
index f4d4de8da..c1e796073 100644
--- a/server/azure/queue_servicebus_routing_test.go
+++ b/server/azure/queue_servicebus_routing_test.go
@@ -390,9 +390,9 @@ func TestBlobTableRoutingUnaffected(t *testing.T) {
// TestStorageHostRoutesListToItsService: "GET /?comp=list" is List Containers
// on a blob host and List Queues on a queue host, as on real Azure where the
-// hostname picks the service. A bare host keeps the shared-endpoint rule
-// (Queue owns it). A blob named "messages" on a blob host is a blob, not a
-// queue message call.
+// hostname picks the service. On a bare host it is List Containers unless a
+// Queue client sends it (see storage_dataplane_isolation_test.go). A blob
+// named "messages" on a blob host is a blob, not a queue message call.
func TestStorageHostRoutesListToItsService(t *testing.T) {
ts := newFullAzureServer(t)
@@ -409,7 +409,7 @@ func TestStorageHostRoutesListToItsService(t *testing.T) {
}{
{"blob host lists containers", blobHost, "ctr-host", "q-host"},
{"queue host lists queues", queueHost, "q-host", "ctr-host"},
- {"bare host lists queues", "", "q-host", "ctr-host"},
+ {"bare host lists containers", "", "ctr-host", "q-host"},
}
for _, tt := range tests {
diff --git a/server/azure/servicebus/dataplane.go b/server/azure/servicebus/dataplane.go
index 230edb5de..08102ee89 100644
--- a/server/azure/servicebus/dataplane.go
+++ b/server/azure/servicebus/dataplane.go
@@ -604,15 +604,22 @@ func (b *wireBrokerProps) sendInput(url, body string, defaultTTLSecs int) mqdriv
in.MessageTTLSeconds = &ttl
}
- if b.ScheduledEnqueueTimeUtc != "" {
- if t, err := time.Parse(time.RFC3339, b.ScheduledEnqueueTimeUtc); err == nil {
- if delay := int(time.Until(t).Seconds()); delay > 0 {
- in.DelaySeconds = delay
- }
+ in.ScheduledEnqueueTime = parseScheduledEnqueueTime(b.ScheduledEnqueueTimeUtc)
+
+ return in
+}
+
+// parseScheduledEnqueueTime reads ScheduledEnqueueTimeUtc. The Service Bus REST
+// docs use the RFC 1123 form ("Sun, 06 Nov 1994 08:49:37 GMT"); ISO 8601 is
+// accepted too. An empty or unparseable value means no schedule.
+func parseScheduledEnqueueTime(raw string) time.Time {
+ for _, layout := range []string{time.RFC1123, time.RFC1123Z, time.RFC3339Nano, "2006-01-02T15:04:05"} {
+ if t, err := time.Parse(layout, raw); err == nil {
+ return t
}
}
- return in
+ return time.Time{}
}
// readSendBody reads a size-capped request body, writing an error on failure.
diff --git a/server/azure/servicebus/dataplane_message_test.go b/server/azure/servicebus/dataplane_message_test.go
index b4f4a7434..746fe2f4d 100644
--- a/server/azure/servicebus/dataplane_message_test.go
+++ b/server/azure/servicebus/dataplane_message_test.go
@@ -151,6 +151,54 @@ func TestDataPlaneScheduledMessageDelayed(t *testing.T) {
}
}
+// TestDataPlaneScheduledEnqueueTimeFormats is the regression for an RFC 1123
+// ScheduledEnqueueTimeUtc (the form the REST docs use) being ignored, so the
+// message was delivered at once. Each accepted form must hold the message
+// until the scheduled instant on the provider clock, then deliver it.
+func TestDataPlaneScheduledEnqueueTimeFormats(t *testing.T) {
+ tests := []struct {
+ name string
+ layout string
+ }{
+ {name: "rfc1123", layout: time.RFC1123},
+ {name: "rfc3339", layout: time.RFC3339},
+ {name: "iso8601-no-zone", layout: "2006-01-02T15:04:05"},
+ }
+
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ srv, clk := newClockServer(t)
+ seedNamespace(t, srv)
+
+ if r := doRequest(t, srv, http.MethodPut, queueURL("fmt")+apiVer, `{"properties":{}}`); r.StatusCode != http.StatusOK {
+ t.Fatalf("create queue = %d", r.StatusCode)
+ }
+
+ // A FakeClock far from the wall clock proves the schedule is measured
+ // on the provider clock.
+ clk.Advance(48 * time.Hour)
+ at := clk.Now().Add(time.Minute).UTC().Truncate(time.Second)
+ broker := fmt.Sprintf(`{"ScheduledEnqueueTimeUtc":%q}`, at.Format(tt.layout))
+
+ send := doRequest(t, srv, http.MethodPost, "/"+nsName+"/fmt/messages", "later",
+ map[string]string{"BrokerProperties": broker})
+ if send.StatusCode != http.StatusCreated {
+ t.Fatalf("send = %d, want 201", send.StatusCode)
+ }
+
+ if early := doRequest(t, srv, http.MethodDelete, "/"+nsName+"/fmt/messages/head", ""); early.StatusCode != http.StatusNoContent {
+ t.Fatalf("receive before schedule = %d, want 204", early.StatusCode)
+ }
+
+ clk.Advance(time.Minute + time.Second)
+
+ if got := doRequest(t, srv, http.MethodDelete, "/"+nsName+"/fmt/messages/head", ""); got.StatusCode != http.StatusOK {
+ t.Fatalf("receive after schedule = %d, want 200", got.StatusCode)
+ }
+ })
+ }
+}
+
// TestDataPlaneScheduledMessageTTLFromEnqueue is the regression for silent loss
// of a scheduled message whose TTL is shorter than the schedule delay: a message
// scheduled for +120s with a 30s TimeToLive must still be delivered once it
@@ -164,9 +212,7 @@ func TestDataPlaneScheduledMessageTTLFromEnqueue(t *testing.T) {
t.Fatalf("create queue = %d", r.StatusCode)
}
- // The dataplane derives the delivery delay from ScheduledEnqueueTimeUtc via the
- // wall clock, while TTL reaping runs on the FakeClock (both anchored at start).
- schedule := time.Now().Add(120 * time.Second).UTC().Format(time.RFC3339)
+ schedule := clk.Now().Add(120 * time.Second).UTC().Format(time.RFC3339)
broker := fmt.Sprintf(`{"ScheduledEnqueueTimeUtc":%q,"TimeToLive":30}`, schedule)
send := doRequest(t, srv, http.MethodPost, "/"+nsName+"/schedttl/messages", "future",
diff --git a/server/azure/storage_dataplane_isolation_test.go b/server/azure/storage_dataplane_isolation_test.go
new file mode 100644
index 000000000..f7b5d8647
--- /dev/null
+++ b/server/azure/storage_dataplane_isolation_test.go
@@ -0,0 +1,277 @@
+package azure_test
+
+import (
+ "context"
+ "io"
+ "net/http"
+ "net/http/httptest"
+ "strings"
+ "testing"
+
+ storagedriver "github.com/stackshy/cloudemu/v2/services/storage/driver"
+)
+
+const (
+ isoAcctA = "acctaa"
+ isoAcctB = "acctbb"
+ queueUA = "azsdk-go-azqueue/v1.0.0 (go1.25; darwin)"
+ tableUA = "azsdk-go-aztables/v1.3.0 (go1.25; darwin)"
+ xmsBlobTp = "x-ms-blob-type"
+)
+
+// newIsolationServer starts the full Azure server with two storage accounts.
+func newIsolationServer(t *testing.T) *httptest.Server {
+ t.Helper()
+
+ ts, p := newFullAzureServerWithProvider(t)
+
+ for _, name := range []string{isoAcctA, isoAcctB} {
+ ref := storagedriver.StorageAccountRef{Name: name, ResourceGroup: "rg"}
+ if _, err := p.BlobStorage.CreateStorageAccount(context.Background(), ref); err != nil {
+ t.Fatalf("create account %s: %v", name, err)
+ }
+ }
+
+ return ts
+}
+
+// storageDo sends one data-plane request with the given Host and User-Agent
+// and returns the status and body.
+func storageDo(t *testing.T, ts *httptest.Server, host, ua, method, path, body string,
+ hdr map[string]string,
+) (status int, respBody string, header http.Header) {
+ t.Helper()
+
+ req, err := http.NewRequestWithContext(context.Background(), method, ts.URL+path, strings.NewReader(body))
+ if err != nil {
+ t.Fatalf("new request: %v", err)
+ }
+
+ if host != "" {
+ req.Host = host
+ }
+
+ req.Header.Set("User-Agent", ua)
+
+ for k, v := range hdr {
+ req.Header.Set(k, v)
+ }
+
+ resp, err := ts.Client().Do(req)
+ if err != nil {
+ t.Fatalf("%s %s: %v", method, path, err)
+ }
+ defer resp.Body.Close()
+
+ data, _ := io.ReadAll(resp.Body)
+
+ return resp.StatusCode, string(data), resp.Header
+}
+
+func expectStatus(t *testing.T, what string, got, want int, body string) {
+ t.Helper()
+
+ if got != want {
+ t.Fatalf("%s = %d, want %d: %s", what, got, want, body)
+ }
+}
+
+// TestQueueAccountIsolation is the regression for AZSTO-06 on Queue: two
+// accounts each get their own "jobs" queue, a message in one is invisible to
+// the other, and listings only show the account's own queues. The default
+// account (bare host) is separate again, and the path-style /{account}/ form
+// reaches the same queue as the account host.
+func TestQueueAccountIsolation(t *testing.T) {
+ ts := newIsolationServer(t)
+ hostA := isoAcctA + ".queue.core.windows.net"
+ hostB := isoAcctB + ".queue.core.windows.net"
+
+ for _, h := range []string{hostA, hostB} {
+ st, body, _ := storageDo(t, ts, h, queueUA, http.MethodPut, "/jobs", "", nil)
+ expectStatus(t, "create jobs on "+h, st, http.StatusCreated, body)
+ }
+
+ msg := "hello"
+ st, body, _ := storageDo(t, ts, hostA, queueUA, http.MethodPost, "/jobs/messages", msg, nil)
+ expectStatus(t, "put message on A", st, http.StatusCreated, body)
+
+ _, body, _ = storageDo(t, ts, hostB, queueUA, http.MethodGet, "/jobs/messages?peekonly=true", "", nil)
+ if strings.Contains(body, "hello") {
+ t.Fatalf("account B sees account A's message: %s", body)
+ }
+
+ _, body, _ = storageDo(t, ts, "", queueUA, http.MethodGet, "/"+isoAcctA+"/jobs/messages?peekonly=true", "", nil)
+ if !strings.Contains(body, "hello") {
+ t.Fatalf("path-style peek on account A = %s, want the message", body)
+ }
+
+ _, body, _ = storageDo(t, ts, "", queueUA, http.MethodGet, "/?comp=list", "", nil)
+ if strings.Contains(body, "jobs") || strings.Contains(body, isoAcctA) {
+ t.Fatalf("default account lists another account's queue: %s", body)
+ }
+
+ _, body, _ = storageDo(t, ts, hostB, queueUA, http.MethodGet, "/?comp=list", "", nil)
+ if !strings.Contains(body, "jobs") {
+ t.Fatalf("account B list = %s, want jobs", body)
+ }
+
+ st, body, _ = storageDo(t, ts, hostA, queueUA, http.MethodDelete, "/jobs", "", nil)
+ expectStatus(t, "delete jobs on A", st, http.StatusNoContent, body)
+
+ st, body, _ = storageDo(t, ts, hostB, queueUA, http.MethodGet, "/jobs?comp=metadata", "", nil)
+ expectStatus(t, "account B jobs after A's delete", st, http.StatusOK, body)
+}
+
+// TestTableAccountIsolation is the regression for AZSTO-06 on Table: the same
+// table name in two accounts holds separate entities and lists separately.
+func TestTableAccountIsolation(t *testing.T) {
+ ts := newIsolationServer(t)
+ hostA := isoAcctA + ".table.core.windows.net"
+ hostB := isoAcctB + ".table.core.windows.net"
+ jsonHdr := map[string]string{"Content-Type": "application/json", "Accept": "application/json;odata=nometadata"}
+
+ for _, h := range []string{hostA, hostB} {
+ st, body, _ := storageDo(t, ts, h, tableUA, http.MethodPost, "/Tables", `{"TableName":"people"}`, jsonHdr)
+ expectStatus(t, "create people on "+h, st, http.StatusCreated, body)
+ }
+
+ st, body, _ := storageDo(t, ts, hostA, tableUA, http.MethodPost, "/people",
+ `{"PartitionKey":"org","RowKey":"alice"}`, jsonHdr)
+ expectStatus(t, "insert on A", st, http.StatusCreated, body)
+
+ st, body, _ = storageDo(t, ts, hostB, tableUA, http.MethodGet,
+ "/people(PartitionKey='org',RowKey='alice')", "", jsonHdr)
+ expectStatus(t, "get on B", st, http.StatusNotFound, body)
+
+ st, body, _ = storageDo(t, ts, "", tableUA, http.MethodGet,
+ "/"+isoAcctA+"/people(PartitionKey='org',RowKey='alice')", "", jsonHdr)
+ expectStatus(t, "path-style get on A", st, http.StatusOK, body)
+
+ _, body, _ = storageDo(t, ts, "", tableUA, http.MethodGet, "/Tables", "", jsonHdr)
+ if strings.Contains(body, "people") {
+ t.Fatalf("default account lists another account's table: %s", body)
+ }
+
+ _, body, _ = storageDo(t, ts, "", tableUA, http.MethodGet, "/"+isoAcctB+"/Tables", "", jsonHdr)
+ if !strings.Contains(body, `"TableName":"people"`) {
+ t.Fatalf("path-style list on B = %s, want people", body)
+ }
+}
+
+// TestStorageRootListAndServiceOps is the regression for AZSTO-07 and
+// AZSTO-08: a bare-host "GET /?comp=list" is List Containers unless a Queue
+// client sends it, and the account-level service calls reach Blob, Queue or
+// Table by host instead of the Cosmos account probe.
+func TestStorageRootListAndServiceOps(t *testing.T) {
+ ts := newIsolationServer(t)
+ blobA := isoAcctA + ".blob.core.windows.net"
+
+ st, body, _ := storageDo(t, ts, "", "", http.MethodPut, "/ctr1?restype=container", "", nil)
+ expectStatus(t, "create container", st, http.StatusCreated, body)
+
+ st, body, _ = storageDo(t, ts, "", "", http.MethodPut, "/q1", "", nil)
+ expectStatus(t, "create queue", st, http.StatusCreated, body)
+
+ _, body, _ = storageDo(t, ts, "", "azsdk-go-azblob/v1.6.0", http.MethodGet, "/?comp=list", "", nil)
+ if !strings.Contains(body, "ctr1") {
+ t.Fatalf("bare-host blob list = %s, want ctr1", body)
+ }
+
+ _, body, _ = storageDo(t, ts, "", queueUA, http.MethodGet, "/?comp=list", "", nil)
+ if !strings.Contains(body, "q1") {
+ t.Fatalf("bare-host queue list = %s, want q1", body)
+ }
+
+ const props = "/?restype=service&comp=properties"
+
+ for _, tc := range []struct{ host, ua string }{
+ {"", ""},
+ {blobA, ""},
+ {isoAcctA + ".queue.core.windows.net", ""},
+ {isoAcctA + ".table.core.windows.net", ""},
+ {"", queueUA},
+ {"", tableUA},
+ } {
+ st, body, _ := storageDo(t, ts, tc.host, tc.ua, http.MethodGet, props, "", nil)
+ if st != http.StatusOK || !strings.Contains(body, "") {
+ t.Fatalf("GET service properties on %q/%q = %d %s", tc.host, tc.ua, st, body)
+ }
+ }
+
+ set := `` +
+ `true7` +
+ `https://a.exampleGET,PUT` +
+ `**60` +
+ ``
+
+ st, body, _ = storageDo(t, ts, blobA, "", http.MethodPut, props, set, nil)
+ expectStatus(t, "set blob service properties", st, http.StatusAccepted, body)
+
+ _, body, _ = storageDo(t, ts, blobA, "", http.MethodGet, props, "", nil)
+ if !strings.Contains(body, "7") || !strings.Contains(body, "https://a.example") {
+ t.Fatalf("blob service properties after set = %s", body)
+ }
+
+ _, body, _ = storageDo(t, ts, isoAcctB+".blob.core.windows.net", "", http.MethodGet, props, "", nil)
+ if strings.Contains(body, "https://a.example") {
+ t.Fatalf("account B sees account A's blob service properties: %s", body)
+ }
+
+ st, body, hdr := storageDo(t, ts, "", "", http.MethodGet, "/?restype=account&comp=properties", "", nil)
+ if st != http.StatusOK || hdr.Get("x-ms-account-kind") != "StorageV2" {
+ t.Fatalf("get account info = %d kind %q: %s", st, hdr.Get("x-ms-account-kind"), body)
+ }
+}
+
+// TestBareHostRootListByUserAgent pins the bare-host rule for the shared
+// "GET /?comp=list" shape: only an Azure SDK Queue product token picks Queue.
+// An application id that merely contains "queue" stays with Blob, as does a
+// request with no User-Agent.
+func TestBareHostRootListByUserAgent(t *testing.T) {
+ ts := newIsolationServer(t)
+
+ st, body, _ := storageDo(t, ts, "", "", http.MethodPut, "/ctr1?restype=container", "", nil)
+ expectStatus(t, "create container", st, http.StatusCreated, body)
+
+ st, body, _ = storageDo(t, ts, "", "", http.MethodPut, "/q1", "", nil)
+ expectStatus(t, "create queue", st, http.StatusCreated, body)
+
+ tests := []struct {
+ name, ua, want string
+ }{
+ {"blob client with queue-like app id", "myqueue-worker azsdk-go-azblob/v1.6.0 (go1.25; darwin)", "ctr1"},
+ {"azqueue", queueUA, "q1"},
+ {"python queue sdk", "azsdk-python-storage-queue/12.10.0 Python/3.12", "q1"},
+ {"net queue sdk", "azsdk-net-Storage.Queues/12.19.0 (.NET 8.0)", "q1"},
+ {"no user agent", "", "ctr1"},
+ {"table client with queue-like app id", "queue-sync azsdk-go-aztables/v1.3.0", "ctr1"},
+ }
+
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ st, body, _ := storageDo(t, ts, "", tt.ua, http.MethodGet, "/?comp=list", "", nil)
+ if st != http.StatusOK || !strings.Contains(body, ""+tt.want+"") {
+ t.Fatalf("GET /?comp=list with UA %q = %d %s, want %s", tt.ua, st, body, tt.want)
+ }
+ })
+ }
+}
+
+// TestQueuePathStyleLeavesBlobPuts guards the path-style peel: a Put Blob into
+// a container of the default account whose name matches a storage account is
+// still a blob write, not a queue create.
+func TestQueuePathStyleLeavesBlobPuts(t *testing.T) {
+ ts := newIsolationServer(t)
+
+ st, body, _ := storageDo(t, ts, "", "", http.MethodPut, "/"+isoAcctA+"/ctr9?restype=container", "", nil)
+ expectStatus(t, "path-style create container", st, http.StatusCreated, body)
+
+ st, body, _ = storageDo(t, ts, "", "", http.MethodPut, "/"+isoAcctA+"/ctr9/b1", "data",
+ map[string]string{xmsBlobTp: "BlockBlob"})
+ expectStatus(t, "path-style put blob", st, http.StatusCreated, body)
+
+ st, body, _ = storageDo(t, ts, "", "", http.MethodGet, "/"+isoAcctA+"/ctr9/b1", "", nil)
+ if st != http.StatusOK || body != "data" {
+ t.Fatalf("path-style get blob = %d %q", st, body)
+ }
+}
diff --git a/server/azure/tablestorage/account.go b/server/azure/tablestorage/account.go
new file mode 100644
index 000000000..c0ea7f3fc
--- /dev/null
+++ b/server/azure/tablestorage/account.go
@@ -0,0 +1,88 @@
+package tablestorage
+
+import (
+ "net/http"
+ "slices"
+ "strings"
+
+ "github.com/stackshy/cloudemu/v2/server/wire/azurearm"
+ storagedriver "github.com/stackshy/cloudemu/v2/services/storage/driver"
+)
+
+// WithAccounts scopes tables to the storage accounts in accounts: a table
+// created through {account}.table.core.windows.net (or the path-style
+// /{account}/ prefix) lives only in that account. Without it every request
+// uses the default account.
+func (h *Handler) WithAccounts(accounts storagedriver.AzureStorageAccounts) *Handler {
+ h.accounts = accounts
+ return h
+}
+
+// resolve returns the account a request targets and the path below it. The
+// {account}.table host names the account when it is an existing storage
+// account; otherwise a path-style /{account}/ prefix does. The default account
+// "cloudemu" and the bare host map to the default namespace (""), which also
+// holds every table from before accounts were modeled.
+func (h *Handler) resolve(r *http.Request) (account, path string) {
+ if acct, svc, ok := azurearm.StorageHost(r.Host); ok && svc == azurearm.StorageServiceTable && h.accountExists(r, acct) {
+ return acct, r.URL.Path
+ }
+
+ rootOp := azurearm.IsStorageServiceOp(r.URL.Query())
+
+ return azurearm.PeelStorageAccount(r.URL.Path, rootOp, func(name string) bool {
+ return h.accountExists(r, name)
+ })
+}
+
+func (h *Handler) accountExists(r *http.Request, name string) bool {
+ if h.accounts == nil || name == storagedriver.AzureDefaultStorageAccount {
+ return false
+ }
+
+ _, err := h.accounts.GetStorageAccount(r.Context(), name)
+
+ return err == nil
+}
+
+// tableKey is the driver name of table in account: the bare name in the
+// default account, "{account}/{table}" in any other. Table names are
+// alphanumeric, so keys never collide.
+func tableKey(account, table string) string {
+ return storagedriver.AzureContainerKey(account, table)
+}
+
+// tableName is the table's own name within its account.
+func tableName(key string) string {
+ _, name := storagedriver.SplitAzureContainerKey(key)
+ return name
+}
+
+// tableSDKProducts are the User-Agent product names of the Azure SDK Table
+// clients, lowercased.
+//
+//nolint:gochecknoglobals // read-only lookup table, not mutable state
+var tableSDKProducts = []string{
+ "azsdk-go-aztables",
+ "azsdk-python-data-tables",
+ "azsdk-java-azure-data-tables",
+ "azsdk-js-data-tables",
+ "azsdk-net-data.tables",
+}
+
+// isTableClient reports whether a root request on a host that does not name
+// its service comes from an Azure SDK Table client. The account-level service
+// calls have the same shape for Blob, Queue and Table, so on a bare host only
+// a product token of the User-Agent (such as "azsdk-go-aztables/v1.3.0")
+// picks Table. Free text such as an application id is ignored.
+func isTableClient(r *http.Request) bool {
+ for _, token := range strings.Fields(strings.ToLower(r.UserAgent())) {
+ product, _, _ := strings.Cut(token, "/")
+
+ if slices.Contains(tableSDKProducts, product) {
+ return true
+ }
+ }
+
+ return false
+}
diff --git a/server/azure/tablestorage/batch.go b/server/azure/tablestorage/batch.go
index 82073ab2e..cfa7d707a 100644
--- a/server/azure/tablestorage/batch.go
+++ b/server/azure/tablestorage/batch.go
@@ -33,7 +33,7 @@ func batchErr(msg string) error {
// batch handles POST /$batch: an OData entity group transaction. It parses the
// multipart/mixed batch + change set, applies the operations atomically, and
// returns the multipart/mixed batch response the aztables client expects.
-func (h *Handler) batch(w http.ResponseWriter, r *http.Request) {
+func (h *Handler) batch(w http.ResponseWriter, r *http.Request, account string) {
if r.Method != http.MethodPost {
writeError(w, http.StatusMethodNotAllowed, "MethodNotAllowed", "method not allowed")
return
@@ -60,7 +60,12 @@ func (h *Handler) batch(w http.ResponseWriter, r *http.Request) {
return
}
- results, applyErr := h.ts.ApplyBatch(r.Context(), table, ops)
+ // A path-style change set addresses its entities as /{account}/{table}(…).
+ if account != "" {
+ table = strings.TrimPrefix(table, account+"/")
+ }
+
+ results, applyErr := h.ts.ApplyBatch(r.Context(), tableKey(account, table), ops)
if applyErr != nil {
writeBatchFailure(w, applyErr)
return
diff --git a/server/azure/tablestorage/handler.go b/server/azure/tablestorage/handler.go
index 980ab5c58..1fd445e44 100644
--- a/server/azure/tablestorage/handler.go
+++ b/server/azure/tablestorage/handler.go
@@ -30,6 +30,7 @@ import (
cerrors "github.com/stackshy/cloudemu/v2/errors"
"github.com/stackshy/cloudemu/v2/server/wire/azurearm"
+ storagedriver "github.com/stackshy/cloudemu/v2/services/storage/driver"
driver "github.com/stackshy/cloudemu/v2/services/tablestorage/driver"
)
@@ -47,12 +48,16 @@ const (
// pathBatch is the entity-group-transaction endpoint path segment.
pathBatch = "$batch"
+
+ // pathTables is the table lifecycle collection path.
+ pathTables = "Tables"
)
// Handler serves Azure Table Storage REST requests against a TableStorage
// driver.
type Handler struct {
- ts driver.TableStorage
+ ts driver.TableStorage
+ accounts storagedriver.AzureStorageAccounts
}
// New returns a Table handler backed by ts.
@@ -65,36 +70,50 @@ func New(ts driver.TableStorage) *Handler {
//
// - The path is /Tables or /Tables('name'): the table lifecycle surface,
// which no other service uses.
+//
// - The path's first segment carries an OData key predicate: it contains a
// "(": either "()" (query entities) or
// "(PartitionKey='…',RowKey='…')" (entity CRUD). Blob/Queue paths never
// contain parentheses, and ARM paths start with /subscriptions/, so this
// is unambiguous.
+//
// - POST /{table} (insert entity) is a bare single segment with a JSON body.
// Blob and Queue never use a bare POST on a single path segment, so the
// method+content-type discriminates it from their PUT/DELETE ops.
//
+// - The account-level service calls (?restype=service|account) on the root,
+// on an {account}.table host or, on a bare host, from a Table client.
+//
+// A path-style /{account}/ prefix naming an existing storage account is
+// peeled before the shape checks (see resolve).
+//
// Registered before the permissive Blob fallback so these shapes win.
-func (*Handler) Matches(r *http.Request) bool {
+func (h *Handler) Matches(r *http.Request) bool {
if strings.HasPrefix(r.URL.Path, "/subscriptions/") {
return false
}
// A storage host names its service, so a request to another service's
// host (such as {account}.blob.core.windows.net) is never a Table call.
- if _, svc, ok := azurearm.StorageHost(r.Host); ok && svc != "table" {
+ _, svc, storageHost := azurearm.StorageHost(r.Host)
+ if storageHost && svc != azurearm.StorageServiceTable {
return false
}
- path := strings.TrimPrefix(r.URL.Path, "/")
+ _, resolved := h.resolve(r)
+ path := strings.TrimPrefix(resolved, "/")
- // /$batch: entity group transactions.
- if path == pathBatch {
- return true
+ if path == "" {
+ return azurearm.IsStorageServiceOp(r.URL.Query()) && (storageHost || isTableClient(r))
}
- // /Tables and /Tables('name').
- if path == "Tables" || strings.HasPrefix(path, "Tables(") {
+ return matchesTablePath(r, path)
+}
+
+// matchesTablePath reports whether a path below the account is a Table call.
+func matchesTablePath(r *http.Request, path string) bool {
+ // /$batch (entity group transactions), /Tables and /Tables('name').
+ if path == pathBatch || path == pathTables || strings.HasPrefix(path, "Tables(") {
return true
}
@@ -125,36 +144,39 @@ func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
w.Header().Set("X-Ms-Version", xmsVersion)
w.Header().Set("Dataserviceversion", "3.0")
- path := strings.TrimPrefix(r.URL.Path, "/")
+ account, resolved := h.resolve(r)
+ path := strings.TrimPrefix(resolved, "/")
switch {
+ case path == "" && azurearm.IsStorageServiceOp(r.URL.Query()):
+ azurearm.ServeStorageServiceOp(w, r)
case path == pathBatch:
- h.batch(w, r)
- case path == "Tables":
- h.tablesCollectionOp(w, r)
+ h.batch(w, r, account)
+ case path == pathTables:
+ h.tablesCollectionOp(w, r, account)
case strings.HasPrefix(path, "Tables("):
- h.deleteTable(w, r, tableNameFromDelete(path))
+ h.deleteTable(w, r, tableKey(account, tableNameFromDelete(path)))
case r.Method == http.MethodPost && !strings.ContainsRune(path, '('):
// POST /{table}: insert entity into a bare table path.
- h.insertEntity(w, r, path)
+ h.insertEntity(w, r, tableKey(account, path))
default:
- h.entityOp(w, r, path)
+ h.entityOp(w, r, account, path)
}
}
// tablesCollectionOp handles POST (create) and GET (list) on /Tables.
-func (h *Handler) tablesCollectionOp(w http.ResponseWriter, r *http.Request) {
+func (h *Handler) tablesCollectionOp(w http.ResponseWriter, r *http.Request, account string) {
switch r.Method {
case http.MethodPost:
- h.createTable(w, r)
+ h.createTable(w, r, account)
case http.MethodGet:
- h.listTables(w, r)
+ h.listTables(w, r, account)
default:
writeError(w, http.StatusMethodNotAllowed, "MethodNotAllowed", "method not allowed")
}
}
-func (h *Handler) createTable(w http.ResponseWriter, r *http.Request) {
+func (h *Handler) createTable(w http.ResponseWriter, r *http.Request, account string) {
var body struct {
TableName string `json:"TableName"`
}
@@ -164,7 +186,7 @@ func (h *Handler) createTable(w http.ResponseWriter, r *http.Request) {
return
}
- if err := h.ts.CreateTable(r.Context(), body.TableName); err != nil {
+ if err := h.ts.CreateTable(r.Context(), tableKey(account, body.TableName)); err != nil {
// A duplicate table is TableAlreadyExists, distinct from the generic
// EntityAlreadyExists that a duplicate entity insert reports.
if cerrors.IsAlreadyExists(err) {
@@ -185,7 +207,7 @@ func (h *Handler) createTable(w http.ResponseWriter, r *http.Request) {
writeJSON(w, http.StatusCreated, resp)
}
-func (h *Handler) listTables(w http.ResponseWriter, r *http.Request) {
+func (h *Handler) listTables(w http.ResponseWriter, r *http.Request, account string) {
names, err := h.ts.ListTables(r.Context())
if err != nil {
writeErr(w, err)
@@ -193,8 +215,11 @@ func (h *Handler) listTables(w http.ResponseWriter, r *http.Request) {
}
value := make([]map[string]any, 0, len(names))
+
for _, n := range names {
- value = append(value, map[string]any{"TableName": n})
+ if acct, name := storagedriver.SplitAzureContainerKey(n); acct == account {
+ value = append(value, map[string]any{"TableName": name})
+ }
}
resp := map[string]any{
@@ -221,13 +246,15 @@ func (h *Handler) deleteTable(w http.ResponseWriter, r *http.Request, name strin
// entityOp routes entity-level requests: /{table}() (query) and
// /{table}(PartitionKey='p',RowKey='r') (CRUD).
-func (h *Handler) entityOp(w http.ResponseWriter, r *http.Request, path string) {
+func (h *Handler) entityOp(w http.ResponseWriter, r *http.Request, account, path string) {
table, predicate, ok := splitEntityPath(path)
if !ok {
writeError(w, http.StatusBadRequest, "InvalidUri", "unrecognized table path")
return
}
+ table = tableKey(account, table)
+
// Query: /{table}() with an empty predicate.
if strings.TrimSpace(predicate) == "" {
h.queryEntities(w, r, table)
@@ -287,7 +314,7 @@ func (h *Handler) queryEntities(w http.ResponseWriter, r *http.Request, table st
}
resp := map[string]any{
- "odata.metadata": fmt.Sprintf("%s://%s/$metadata#%s", scheme(r), r.Host, table),
+ "odata.metadata": fmt.Sprintf("%s://%s/$metadata#%s", scheme(r), r.Host, tableName(table)),
"value": value,
}
@@ -302,7 +329,7 @@ func (h *Handler) getEntity(w http.ResponseWriter, r *http.Request, table, pk, r
}
out := entityToJSON(ent)
- out["odata.metadata"] = fmt.Sprintf("%s://%s/$metadata#%s/@Element", scheme(r), r.Host, table)
+ out["odata.metadata"] = fmt.Sprintf("%s://%s/$metadata#%s/@Element", scheme(r), r.Host, tableName(table))
if etag := asString(ent[etagProp]); etag != "" {
w.Header().Set("ETag", etag)
@@ -329,7 +356,7 @@ func (h *Handler) insertEntity(w http.ResponseWriter, r *http.Request, table str
// Default (no Prefer header): return-content, echoing the entity with 201.
out := entityToJSON(ent)
- out["odata.metadata"] = fmt.Sprintf("%s://%s/$metadata#%s/@Element", scheme(r), r.Host, table)
+ out["odata.metadata"] = fmt.Sprintf("%s://%s/$metadata#%s/@Element", scheme(r), r.Host, tableName(table))
out[etagProp] = etag
w.Header().Set("ETag", etag)
diff --git a/server/wire/azurearm/storagehost.go b/server/wire/azurearm/storagehost.go
index 5d57a6f1e..75f5bbff4 100644
--- a/server/wire/azurearm/storagehost.go
+++ b/server/wire/azurearm/storagehost.go
@@ -1,6 +1,21 @@
package azurearm
-import "strings"
+import (
+ "net/url"
+ "strings"
+)
+
+// Storage service labels of an {account}.{service}.{suffix} host.
+const (
+ StorageServiceQueue = "queue"
+ StorageServiceTable = "table"
+)
+
+// restype values of the account-level storage service operations.
+const (
+ restypeService = "service"
+ restypeAccount = "account"
+)
// storageServices are the Azure Storage data-plane service labels that sit
// between the account name and the storage DNS suffix, as in
@@ -49,3 +64,31 @@ func StorageHost(host string) (account, service string, ok bool) {
return "", "", false
}
+
+// PeelStorageAccount reads a path-style storage request
+// (https://host:port/{account}/...), as the SDKs build it for an IP or
+// emulator endpoint. The leading segment is taken as the account, and the
+// rest of the path returned, when isAccount accepts it and either more path
+// follows or the request is an account-level operation (rootOp), such as
+// "GET /{account}?comp=list". Otherwise the request belongs to the default
+// account ("") and path is returned unchanged.
+func PeelStorageAccount(path string, rootOp bool, isAccount func(name string) bool) (account, rest string) {
+ name, after, more := strings.Cut(strings.TrimPrefix(path, "/"), "/")
+ if name == "" || (!more && !rootOp) || !isAccount(name) {
+ return "", path
+ }
+
+ return name, "/" + after
+}
+
+// IsStorageServiceOp reports whether q selects an account-level storage
+// service operation: Get/Set Service Properties, Get Service Stats
+// (restype=service) or Get Account Information (restype=account).
+func IsStorageServiceOp(q url.Values) bool {
+ switch q.Get("restype") {
+ case restypeService, restypeAccount:
+ return true
+ }
+
+ return false
+}
diff --git a/server/wire/azurearm/storageservice.go b/server/wire/azurearm/storageservice.go
new file mode 100644
index 000000000..5ef398724
--- /dev/null
+++ b/server/wire/azurearm/storageservice.go
@@ -0,0 +1,172 @@
+package azurearm
+
+import (
+ "encoding/xml"
+ "net/http"
+ "net/url"
+)
+
+// storageAnalyticsVersion is the version Storage Analytics reports for the
+// logging and metrics settings of a new account.
+const storageAnalyticsVersion = "1.0"
+
+// compProperties is the comp= value of the service and account property calls.
+const compProperties = "properties"
+
+// StorageRetentionPolicy is the RetentionPolicy / DeleteRetentionPolicy
+// element of the Storage service properties document.
+type StorageRetentionPolicy struct {
+ Enabled bool `xml:"Enabled"`
+ Days int `xml:"Days,omitempty"`
+}
+
+// StorageLogging is the Logging element of the service properties document.
+type StorageLogging struct {
+ Version string `xml:"Version"`
+ Delete bool `xml:"Delete"`
+ Read bool `xml:"Read"`
+ Write bool `xml:"Write"`
+ RetentionPolicy StorageRetentionPolicy `xml:"RetentionPolicy"`
+}
+
+// StorageMetrics is the HourMetrics / MinuteMetrics element.
+type StorageMetrics struct {
+ Version string `xml:"Version"`
+ Enabled bool `xml:"Enabled"`
+ IncludeAPIs *bool `xml:"IncludeAPIs,omitempty"`
+ RetentionPolicy StorageRetentionPolicy `xml:"RetentionPolicy"`
+}
+
+// StorageCorsRule is one CorsRule element. List-valued fields are
+// comma-separated, as on the wire.
+type StorageCorsRule struct {
+ AllowedOrigins string `xml:"AllowedOrigins"`
+ AllowedMethods string `xml:"AllowedMethods"`
+ AllowedHeaders string `xml:"AllowedHeaders"`
+ ExposedHeaders string `xml:"ExposedHeaders"`
+ MaxAgeInSeconds int `xml:"MaxAgeInSeconds"`
+}
+
+// StorageCors is the Cors element.
+type StorageCors struct {
+ Rules []StorageCorsRule `xml:"CorsRule"`
+}
+
+// StorageServiceProperties is the StorageServiceProperties document of Get
+// and Set Service Properties for the Blob, Queue and Table services. Optional
+// elements are pointers so a Set request can tell an omitted element (keep the
+// current value) from one that is present.
+type StorageServiceProperties struct {
+ XMLName xml.Name `xml:"StorageServiceProperties"`
+ Logging *StorageLogging `xml:"Logging,omitempty"`
+ HourMetrics *StorageMetrics `xml:"HourMetrics,omitempty"`
+ MinuteMetrics *StorageMetrics `xml:"MinuteMetrics,omitempty"`
+ Cors *StorageCors `xml:"Cors,omitempty"`
+ DefaultServiceVersion string `xml:"DefaultServiceVersion,omitempty"`
+ DeleteRetentionPolicy *StorageRetentionPolicy `xml:"DeleteRetentionPolicy,omitempty"`
+}
+
+// DefaultStorageServiceProperties returns the service properties real Azure
+// reports for a new account: logging and metrics off, no CORS rules.
+func DefaultStorageServiceProperties() StorageServiceProperties {
+ includeAPIs := false
+
+ return StorageServiceProperties{
+ Logging: &StorageLogging{Version: storageAnalyticsVersion},
+ HourMetrics: &StorageMetrics{Version: storageAnalyticsVersion, IncludeAPIs: &includeAPIs},
+ MinuteMetrics: &StorageMetrics{Version: storageAnalyticsVersion, IncludeAPIs: &includeAPIs},
+ Cors: &StorageCors{},
+ }
+}
+
+// DecodeStorageServiceProperties reads a Set Service Properties body, writing
+// a 400 InvalidXmlDocument and returning false when it does not parse.
+func DecodeStorageServiceProperties(w http.ResponseWriter, r *http.Request) (StorageServiceProperties, bool) {
+ var props StorageServiceProperties
+ if err := xml.NewDecoder(r.Body).Decode(&props); err != nil {
+ writeStorageXMLError(w, http.StatusBadRequest, "InvalidXmlDocument",
+ "XML specified is not syntactically valid.")
+
+ return props, false
+ }
+
+ return props, true
+}
+
+// WriteStorageServiceProperties writes a 200 Get Service Properties response.
+func WriteStorageServiceProperties(w http.ResponseWriter, props *StorageServiceProperties) {
+ out, err := xml.Marshal(props)
+ if err != nil {
+ writeStorageXMLError(w, http.StatusInternalServerError, "InternalError", err.Error())
+ return
+ }
+
+ w.Header().Set("Content-Type", "application/xml")
+ w.WriteHeader(http.StatusOK)
+ _, _ = w.Write([]byte(xml.Header))
+ _, _ = w.Write(out)
+}
+
+// WriteStorageAccountInfo answers Get Account Information
+// (?restype=account&comp=properties) with the SKU and kind of a default
+// StorageV2 account.
+func WriteStorageAccountInfo(w http.ResponseWriter) {
+ w.Header().Set("x-ms-sku-name", "Standard_LRS")
+ w.Header().Set("x-ms-account-kind", "StorageV2")
+ w.Header().Set("x-ms-is-hns-enabled", "false")
+ w.WriteHeader(http.StatusOK)
+}
+
+// storageXMLError is the Storage service error body.
+type storageXMLError struct {
+ XMLName xml.Name `xml:"Error"`
+ Code string `xml:"Code"`
+ Message string `xml:"Message"`
+}
+
+func writeStorageXMLError(w http.ResponseWriter, status int, code, msg string) {
+ out, _ := xml.Marshal(storageXMLError{Code: code, Message: msg})
+
+ w.Header().Set("Content-Type", "application/xml")
+ w.Header().Set("x-ms-error-code", code)
+ w.WriteHeader(status)
+ _, _ = w.Write([]byte(xml.Header))
+ _, _ = w.Write(out)
+}
+
+// ServeStorageServiceOp answers the account-level service operations that
+// keep no state in cloudemu's Queue and Table services: Get Service
+// Properties reports the defaults, Set Service Properties validates the body
+// and returns 202, and Get Account Information reports the account's SKU and
+// kind. Any other selector is a 400.
+func ServeStorageServiceOp(w http.ResponseWriter, r *http.Request) {
+ q := r.URL.Query()
+
+ if q.Get("comp") != compProperties {
+ writeStorageXMLError(w, http.StatusBadRequest, "InvalidQueryParameterValue",
+ "Value for one of the query parameters specified in the request URI is invalid.")
+
+ return
+ }
+
+ switch op := q.Get("restype") + " " + r.Method; op {
+ case restypeAccount + " " + http.MethodGet, restypeAccount + " " + http.MethodHead:
+ WriteStorageAccountInfo(w)
+ case restypeService + " " + http.MethodGet:
+ props := DefaultStorageServiceProperties()
+ WriteStorageServiceProperties(w, &props)
+ case restypeService + " " + http.MethodPut:
+ if _, ok := DecodeStorageServiceProperties(w, r); ok {
+ w.WriteHeader(http.StatusAccepted)
+ }
+ default:
+ writeStorageXMLError(w, http.StatusMethodNotAllowed, "UnsupportedHttpVerb",
+ "The resource doesn't support specified Http Verb.")
+ }
+}
+
+// IsStorageServicePropertiesOp reports whether q selects Get or Set Service
+// Properties.
+func IsStorageServicePropertiesOp(q url.Values) bool {
+ return q.Get("restype") == restypeService && q.Get("comp") == compProperties
+}
diff --git a/services/messagequeue/driver/driver.go b/services/messagequeue/driver/driver.go
index fbf503d3d..9b21f9ca0 100644
--- a/services/messagequeue/driver/driver.go
+++ b/services/messagequeue/driver/driver.go
@@ -153,6 +153,10 @@ type SendMessageInput struct {
// ReplyToSessionId, ContentType) so they survive a send/receive round-trip.
// Ignored by non-Azure providers.
SystemProperties map[string]string
+ // ScheduledEnqueueTime is Azure Service Bus's ScheduledEnqueueTimeUtc: the
+ // message stays invisible until this instant, measured on the provider's
+ // clock. The zero value means no schedule. Ignored by non-Azure providers.
+ ScheduledEnqueueTime time.Time
}
// SendMessageOutput is the result of sending a message.