Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
95 changes: 75 additions & 20 deletions cpp/src/arrow/filesystem/azurefs.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -347,6 +347,22 @@ bool IsContainerNotFound(const Storage::StorageException& e) {
return false;
}

const auto kHierarchicalNamespaceIsDirectoryMetadataKey = "hdi_isFolder";
const auto kFlatNamespaceIsDirectoryMetadataKey = "is_directory";

bool MetadataIndicatesIsDirectory(const Storage::Metadata& metadata) {
// Inspired by
// https://github.com/Azure/azure-sdk-for-cpp/blob/12407e8bfcb9bc1aa43b253c1d0ec93bf795ae3b/sdk/storage/azure-storage-files-datalake/src/datalake_utilities.cpp#L86-L91
auto hierarchical_directory_metadata =
metadata.find(kHierarchicalNamespaceIsDirectoryMetadataKey);
if (hierarchical_directory_metadata != metadata.end()) {
return hierarchical_directory_metadata->second == "true";
}
auto flat_directory_metadata = metadata.find(kFlatNamespaceIsDirectoryMetadataKey);
return flat_directory_metadata != metadata.end() &&
flat_directory_metadata->second == "true";
}

template <typename ArrowType>
std::string FormatValue(typename TypeTraits<ArrowType>::CType value) {
struct StringAppender {
Expand DownExpand Up@@ -512,11 +528,18 @@ class ObjectInputFile final : public io::RandomAccessFile {

Status Init() {
if (content_length_ != kNoSize) {
// When the user provides the file size we don't validate that its a file. This is
// only a read so its not a big deal if the user makes a mistake.
DCHECK_GE(content_length_, 0);
return Status::OK();
}
try {
// To open an ObjectInputFile the Blob must exist and it must not represent
// a directory. Additionally we need to know the file size.
auto properties = blob_client_->GetProperties();
if (MetadataIndicatesIsDirectory(properties.Value.Metadata)) {
return NotAFile(location_);
}
content_length_ = properties.Value.BlobSize;
metadata_ = PropertiesToMetadata(properties.Value);
return Status::OK();
Expand DownExpand Up@@ -698,11 +721,10 @@ class ObjectAppendStream final : public io::OutputStream {
ObjectAppendStream(std::shared_ptr<Blobs::BlockBlobClient> block_blob_client,
const io::IOContext& io_context, const AzureLocation& location,
const std::shared_ptr<const KeyValueMetadata>& metadata,
const AzureOptions& options, int64_t size = kNoSize)
const AzureOptions& options)
: block_blob_client_(std::move(block_blob_client)),
io_context_(io_context),
location_(location),
content_length_(size) {
location_(location) {
if (metadata && metadata->size() != 0) {
metadata_ = ArrowMetadataToAzureMetadata(metadata);
} else if (options.default_metadata && options.default_metadata->size() != 0) {
Expand All@@ -716,17 +738,31 @@ class ObjectAppendStream final : public io::OutputStream {
io::internal::CloseFromDestructor(this);
}

Status Init() {
if (content_length_ != kNoSize) {
DCHECK_GE(content_length_, 0);
pos_ = content_length_;
Status Init(const bool truncate,
std::function<Status()> ensure_not_flat_namespace_directory) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

You can inject AzureFileSystem *azure_file_system here and not have to allocate a closure for this. You would call AzureFileSystem::Impl::EnsureNotFlatNamespaceDirectory(location) via azure_file_system->impl_ (accessible because the handles produced by the azure file system can be friends with the filesystem class).

@Tom-NewtonTom-NewtonFeb 20, 2024

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the extra info. I was planning to do this but I was struggling with the friends thing.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The trick is to forward-declare the class as well.

diff --git a/cpp/src/arrow/filesystem/azurefs.h b/cpp/src/arrow/filesystem/azurefs.h
index 2a131e40c..d48ef9dd7 100644
--- a/cpp/src/arrow/filesystem/azurefs.h
+++ b/cpp/src/arrow/filesystem/azurefs.h
@@ -44,6 +44,7 @@ classDataLakeServiceClient;
namespacearrow::fs {
+classObjectAppendStream;
classTestAzureFileSystem;
/// Options for the AzureFileSystem implementation.
@@ -180,6 +181,7 @@ classARROW_EXPORT AzureFileSystem : public FileSystem {
explicitAzureFileSystem(std::unique_ptr<Impl>&& impl);
+ friendclassObjectAppendStream;
friendclassTestAzureFileSystem;
voidForceCachedHierarchicalNamespaceSupport(int hns_support);

@Tom-NewtonTom-NewtonFeb 21, 2024

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think my main problem was that ObjectAppendStream is defined inside an anonymous namespace but I still haven't got it working as you describe.

Are you suggesting to use AzureFileSystem *azure_file_system or AzureFileSystem:Impl *azure_file_system as the argument to ObjectAppendStream::Impl. I don't know how I can get a AzureFileSystem pointer from inside AzureFileSystem::Impl and using AzureFileSystem::Impl as the argument leads to incomplete type errors which I don't think I can avoid.

Also if you wouldn't mind I would be interested to know what the disadvantage of a lambda function is compared to what you proposed.

Sorry about my lacking C++ knowledge here.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Also if you wouldn't mind I would be interested to know what the disadvantage of a lambda function is compared to what you proposed.

To create the std::function, you heap allocate an object with copies of the values in the capture list and generate a lot more extra code in the binary:

class function {
T valuesfromthecpapturelist;
RetType operator()(ArgsType ...) {...};
}

When you think about an std::function this way (a pair of context data and a function), you realize the class you already serves that purpose.

But hey, this is becoming challenging, so I won't hold the PR anymore because of this. Moving to Init() was a big step in the right direction.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for explaining

if (truncate) {
content_length_ = 0;
pos_ = 0;
// We need to create an empty file overwriting any existing file, but
// fail if there is an existing directory.
RETURN_NOT_OK(ensure_not_flat_namespace_directory());
// On hierarchical namespace CreateEmptyBlockBlob will fail if there is an existing
// directory so we don't need to check like we do on flat namespace.
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client_));
} else {
try {
auto properties = block_blob_client_->GetProperties();
if (MetadataIndicatesIsDirectory(properties.Value.Metadata)) {
return NotAFile(location_);
}
content_length_ = properties.Value.BlobSize;
pos_ = content_length_;
} catch (const Storage::StorageException& exception) {
if (exception.StatusCode == Http::HttpStatusCode::NotFound) {
// No file exists but on flat namespace its possible there is a directory
// marker or an implied directory. Ensure there is no directory before starting
// a new empty file.
RETURN_NOT_OK(ensure_not_flat_namespace_directory());
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client_));
} else {
return ExceptionToStatus(
Expand All@@ -743,6 +779,7 @@ class ObjectAppendStream final : public io::OutputStream {
block_ids_.push_back(block.Name);
}
}
initialised_ = true;
return Status::OK();
}

Expand DownExpand Up@@ -789,6 +826,11 @@ class ObjectAppendStream final : public io::OutputStream {

Status Flush() override {
RETURN_NOT_OK(CheckClosed("flush"));
if (!initialised_) {
// If the stream has not been successfully initialized then there is nothing to
// flush. This also avoids some unhandled errors when flushing in the destructor.
return Status::OK();
}
return CommitBlockList(block_blob_client_, block_ids_, metadata_);
}

Expand DownExpand Up@@ -840,10 +882,11 @@ class ObjectAppendStream final : public io::OutputStream {
std::shared_ptr<Blobs::BlockBlobClient> block_blob_client_;
const io::IOContext io_context_;
const AzureLocation location_;
int64_t content_length_ = kNoSize;

bool closed_ = false;
bool initialised_ = false;
int64_t pos_ = 0;
int64_t content_length_ = kNoSize;
std::vector<std::string> block_ids_;
Storage::Metadata metadata_;
};
Expand DownExpand Up@@ -1662,20 +1705,32 @@ class AzureFileSystem::Impl {
AzureFileSystem* fs) {
RETURN_NOT_OK(ValidateFileLocation(location));

const auto blob_container_client = GetBlobContainerClient(location.container);
auto block_blob_client = std::make_shared<Blobs::BlockBlobClient>(
blob_service_client_->GetBlobContainerClient(location.container)
.GetBlockBlobClient(location.path));
blob_container_client.GetBlockBlobClient(location.path));

auto ensure_not_flat_namespace_directory = [this, location,
blob_container_client]() -> Status {
ARROW_ASSIGN_OR_RAISE(
auto hns_support,
HierarchicalNamespaceSupport(GetFileSystemClient(location.container)));
if (hns_support == HNSSupport::kDisabled) {
// Flat namespace so we need to GetFileInfo in-case its a directory.
ARROW_ASSIGN_OR_RAISE(auto status, GetFileInfo(blob_container_client, location))
if (status.type() == FileType::Directory) {
return NotAFile(location);
}
}
// kContainerNotFound - it doesn't exist, so no need to check if its a directory.
// kEnabled - hierarchical namespace so Azure APIs will fail if its a directory. We
// don't need to explicitly check.
return Status::OK();
};

std::shared_ptr<ObjectAppendStream> stream;
if (truncate) {
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client));
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_, 0);
} else {
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_);
}
RETURN_NOT_OK(stream->Init());
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_);
RETURN_NOT_OK(stream->Init(truncate, ensure_not_flat_namespace_directory));
return stream;
}

Expand All@@ -1690,7 +1745,7 @@ class AzureFileSystem::Impl {
// on directory marker blobs.
// https://github.com/fsspec/adlfs/blob/32132c4094350fca2680155a5c236f2e9f991ba5/adlfs/spec.py#L855-L870
Blobs::UploadBlockBlobFromOptions blob_options;
blob_options.Metadata.emplace("is_directory", "true");
blob_options.Metadata.emplace(kFlatNamespaceIsDirectoryMetadataKey, "true");
block_blob_client.UploadFrom(nullptr, 0, blob_options);
}

Expand Down
60 changes: 60 additions & 0 deletions cpp/src/arrow/filesystem/azurefs_test.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -826,6 +826,41 @@ class TestAzureFileSystem : public ::testing::Test {
AssertFileInfo(fs(), subdir3, FileType::Directory);
}

void TestDisallowReadingOrWritingDirectoryMarkers() {
auto data = SetUpPreexistingData();
auto directory_path = data.Path("directory");

ASSERT_OK(fs()->CreateDir(directory_path));
ASSERT_RAISES(IOError, fs()->OpenInputFile(directory_path));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(directory_path));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(directory_path));

auto directory_path_with_slash = directory_path + "/";
ASSERT_RAISES(IOError, fs()->OpenInputFile(directory_path_with_slash));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(directory_path_with_slash));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(directory_path_with_slash));
}

void TestDisallowCreatingFileAndDirectoryWithTheSameName() {
auto data = SetUpPreexistingData();
auto path1 = data.Path("directory1");
ASSERT_OK(fs()->CreateDir(path1));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(path1));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(path1));
AssertFileInfo(fs(), path1, FileType::Directory);

auto path2 = data.Path("directory2");
ASSERT_OK(fs()->OpenOutputStream(path2));
// CreateDir returns OK even if there is already a file or directory at this
// location. Whether or not this is the desired behaviour is debatable.
ASSERT_OK(fs()->CreateDir(path2));
AssertFileInfo(fs(), path2, FileType::File);
}

void TestOpenOutputStreamWithMissingContainer() {
ASSERT_RAISES(IOError, fs()->OpenOutputStream("not-a-container/file", {}));
}

void TestDeleteDirSuccessEmpty() {
if (HasSubmitBatchBug()) {
GTEST_SKIP() << kSubmitBatchBugMessage;
Expand DownExpand Up@@ -1559,6 +1594,19 @@ TYPED_TEST(TestAzureFileSystemOnAllScenarios, CreateDirOnMissingContainer) {
this->TestCreateDirOnMissingContainer();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, DisallowReadingOrWritingDirectoryMarkers) {
this->TestDisallowReadingOrWritingDirectoryMarkers();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios,
DisallowCreatingFileAndDirectoryWithTheSameName) {
this->TestDisallowCreatingFileAndDirectoryWithTheSameName();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, OpenOutputStreamWithMissingContainer) {
this->TestOpenOutputStreamWithMissingContainer();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, DeleteDirSuccessEmpty) {
this->TestDeleteDirSuccessEmpty();
}
Expand DownExpand Up@@ -2162,6 +2210,18 @@ TEST_F(TestAzuriteFileSystem, WriteMetadata) {
.Value.Metadata;
// Defaults are overwritten and not merged.
EXPECT_EQ(Core::CaseInsensitiveMap{std::make_pair("bar", "foo")}, blob_metadata);

// Metadata can be written without writing any data.
ASSERT_OK_AND_ASSIGN(
output, fs_with_defaults->OpenAppendStream(
full_path, /*metadata=*/arrow::key_value_metadata({{"bar", "baz"}})));
ASSERT_OK(output->Close());
blob_metadata = blob_service_client_->GetBlobContainerClient(data.container_name)
.GetBlockBlobClient(blob_path)
.GetProperties()
.Value.Metadata;
// Defaults are overwritten and not merged.
EXPECT_EQ(Core::CaseInsensitiveMap{std::make_pair("bar", "baz")}, blob_metadata);
}

TEST_F(TestAzuriteFileSystem, OpenOutputStreamSmall) {
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all \u003cpre\u003e\u003ccode\u003e blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks"); } } catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); } })(); (function(){ try { var __m = "github.com"; var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
95 changes: 75 additions & 20 deletions cpp/src/arrow/filesystem/azurefs.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -347,6 +347,22 @@ bool IsContainerNotFound(const Storage::StorageException& e) {
return false;
}

const auto kHierarchicalNamespaceIsDirectoryMetadataKey = "hdi_isFolder";
const auto kFlatNamespaceIsDirectoryMetadataKey = "is_directory";

bool MetadataIndicatesIsDirectory(const Storage::Metadata& metadata) {
// Inspired by
// https://github.com/Azure/azure-sdk-for-cpp/blob/12407e8bfcb9bc1aa43b253c1d0ec93bf795ae3b/sdk/storage/azure-storage-files-datalake/src/datalake_utilities.cpp#L86-L91
auto hierarchical_directory_metadata =
metadata.find(kHierarchicalNamespaceIsDirectoryMetadataKey);
if (hierarchical_directory_metadata != metadata.end()) {
return hierarchical_directory_metadata->second == "true";
}
auto flat_directory_metadata = metadata.find(kFlatNamespaceIsDirectoryMetadataKey);
return flat_directory_metadata != metadata.end() &&
flat_directory_metadata->second == "true";
}

template <typename ArrowType>
std::string FormatValue(typename TypeTraits<ArrowType>::CType value) {
struct StringAppender {
Expand DownExpand Up@@ -512,11 +528,18 @@ class ObjectInputFile final : public io::RandomAccessFile {

Status Init() {
if (content_length_ != kNoSize) {
// When the user provides the file size we don't validate that its a file. This is
// only a read so its not a big deal if the user makes a mistake.
DCHECK_GE(content_length_, 0);
return Status::OK();
}
try {
// To open an ObjectInputFile the Blob must exist and it must not represent
// a directory. Additionally we need to know the file size.
auto properties = blob_client_->GetProperties();
if (MetadataIndicatesIsDirectory(properties.Value.Metadata)) {
return NotAFile(location_);
}
content_length_ = properties.Value.BlobSize;
metadata_ = PropertiesToMetadata(properties.Value);
return Status::OK();
Expand DownExpand Up@@ -698,11 +721,10 @@ class ObjectAppendStream final : public io::OutputStream {
ObjectAppendStream(std::shared_ptr<Blobs::BlockBlobClient> block_blob_client,
const io::IOContext& io_context, const AzureLocation& location,
const std::shared_ptr<const KeyValueMetadata>& metadata,
const AzureOptions& options, int64_t size = kNoSize)
const AzureOptions& options)
: block_blob_client_(std::move(block_blob_client)),
io_context_(io_context),
location_(location),
content_length_(size) {
location_(location) {
if (metadata && metadata->size() != 0) {
metadata_ = ArrowMetadataToAzureMetadata(metadata);
} else if (options.default_metadata && options.default_metadata->size() != 0) {
Expand All@@ -716,17 +738,31 @@ class ObjectAppendStream final : public io::OutputStream {
io::internal::CloseFromDestructor(this);
}

Status Init() {
if (content_length_ != kNoSize) {
DCHECK_GE(content_length_, 0);
pos_ = content_length_;
Status Init(const bool truncate,
std::function<Status()> ensure_not_flat_namespace_directory) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

You can inject AzureFileSystem *azure_file_system here and not have to allocate a closure for this. You would call AzureFileSystem::Impl::EnsureNotFlatNamespaceDirectory(location) via azure_file_system->impl_ (accessible because the handles produced by the azure file system can be friends with the filesystem class).

@Tom-NewtonTom-NewtonFeb 20, 2024

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the extra info. I was planning to do this but I was struggling with the friends thing.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The trick is to forward-declare the class as well.

diff --git a/cpp/src/arrow/filesystem/azurefs.h b/cpp/src/arrow/filesystem/azurefs.h
index 2a131e40c..d48ef9dd7 100644
--- a/cpp/src/arrow/filesystem/azurefs.h
+++ b/cpp/src/arrow/filesystem/azurefs.h
@@ -44,6 +44,7 @@ classDataLakeServiceClient;
namespacearrow::fs {
+classObjectAppendStream;
classTestAzureFileSystem;
/// Options for the AzureFileSystem implementation.
@@ -180,6 +181,7 @@ classARROW_EXPORT AzureFileSystem : public FileSystem {
explicitAzureFileSystem(std::unique_ptr<Impl>&& impl);
+ friendclassObjectAppendStream;
friendclassTestAzureFileSystem;
voidForceCachedHierarchicalNamespaceSupport(int hns_support);

@Tom-NewtonTom-NewtonFeb 21, 2024

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think my main problem was that ObjectAppendStream is defined inside an anonymous namespace but I still haven't got it working as you describe.

Are you suggesting to use AzureFileSystem *azure_file_system or AzureFileSystem:Impl *azure_file_system as the argument to ObjectAppendStream::Impl. I don't know how I can get a AzureFileSystem pointer from inside AzureFileSystem::Impl and using AzureFileSystem::Impl as the argument leads to incomplete type errors which I don't think I can avoid.

Also if you wouldn't mind I would be interested to know what the disadvantage of a lambda function is compared to what you proposed.

Sorry about my lacking C++ knowledge here.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Also if you wouldn't mind I would be interested to know what the disadvantage of a lambda function is compared to what you proposed.

To create the std::function, you heap allocate an object with copies of the values in the capture list and generate a lot more extra code in the binary:

class function {
T valuesfromthecpapturelist;
RetType operator()(ArgsType ...) {...};
}

When you think about an std::function this way (a pair of context data and a function), you realize the class you already serves that purpose.

But hey, this is becoming challenging, so I won't hold the PR anymore because of this. Moving to Init() was a big step in the right direction.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for explaining

if (truncate) {
content_length_ = 0;
pos_ = 0;
// We need to create an empty file overwriting any existing file, but
// fail if there is an existing directory.
RETURN_NOT_OK(ensure_not_flat_namespace_directory());
// On hierarchical namespace CreateEmptyBlockBlob will fail if there is an existing
// directory so we don't need to check like we do on flat namespace.
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client_));
} else {
try {
auto properties = block_blob_client_->GetProperties();
if (MetadataIndicatesIsDirectory(properties.Value.Metadata)) {
return NotAFile(location_);
}
content_length_ = properties.Value.BlobSize;
pos_ = content_length_;
} catch (const Storage::StorageException& exception) {
if (exception.StatusCode == Http::HttpStatusCode::NotFound) {
// No file exists but on flat namespace its possible there is a directory
// marker or an implied directory. Ensure there is no directory before starting
// a new empty file.
RETURN_NOT_OK(ensure_not_flat_namespace_directory());
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client_));
} else {
return ExceptionToStatus(
Expand All@@ -743,6 +779,7 @@ class ObjectAppendStream final : public io::OutputStream {
block_ids_.push_back(block.Name);
}
}
initialised_ = true;
return Status::OK();
}

Expand DownExpand Up@@ -789,6 +826,11 @@ class ObjectAppendStream final : public io::OutputStream {

Status Flush() override {
RETURN_NOT_OK(CheckClosed("flush"));
if (!initialised_) {
// If the stream has not been successfully initialized then there is nothing to
// flush. This also avoids some unhandled errors when flushing in the destructor.
return Status::OK();
}
return CommitBlockList(block_blob_client_, block_ids_, metadata_);
}

Expand DownExpand Up@@ -840,10 +882,11 @@ class ObjectAppendStream final : public io::OutputStream {
std::shared_ptr<Blobs::BlockBlobClient> block_blob_client_;
const io::IOContext io_context_;
const AzureLocation location_;
int64_t content_length_ = kNoSize;

bool closed_ = false;
bool initialised_ = false;
int64_t pos_ = 0;
int64_t content_length_ = kNoSize;
std::vector<std::string> block_ids_;
Storage::Metadata metadata_;
};
Expand DownExpand Up@@ -1662,20 +1705,32 @@ class AzureFileSystem::Impl {
AzureFileSystem* fs) {
RETURN_NOT_OK(ValidateFileLocation(location));

const auto blob_container_client = GetBlobContainerClient(location.container);
auto block_blob_client = std::make_shared<Blobs::BlockBlobClient>(
blob_service_client_->GetBlobContainerClient(location.container)
.GetBlockBlobClient(location.path));
blob_container_client.GetBlockBlobClient(location.path));

auto ensure_not_flat_namespace_directory = [this, location,
blob_container_client]() -> Status {
ARROW_ASSIGN_OR_RAISE(
auto hns_support,
HierarchicalNamespaceSupport(GetFileSystemClient(location.container)));
if (hns_support == HNSSupport::kDisabled) {
// Flat namespace so we need to GetFileInfo in-case its a directory.
ARROW_ASSIGN_OR_RAISE(auto status, GetFileInfo(blob_container_client, location))
if (status.type() == FileType::Directory) {
return NotAFile(location);
}
}
// kContainerNotFound - it doesn't exist, so no need to check if its a directory.
// kEnabled - hierarchical namespace so Azure APIs will fail if its a directory. We
// don't need to explicitly check.
return Status::OK();
};

std::shared_ptr<ObjectAppendStream> stream;
if (truncate) {
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client));
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_, 0);
} else {
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_);
}
RETURN_NOT_OK(stream->Init());
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_);
RETURN_NOT_OK(stream->Init(truncate, ensure_not_flat_namespace_directory));
return stream;
}

Expand All@@ -1690,7 +1745,7 @@ class AzureFileSystem::Impl {
// on directory marker blobs.
// https://github.com/fsspec/adlfs/blob/32132c4094350fca2680155a5c236f2e9f991ba5/adlfs/spec.py#L855-L870
Blobs::UploadBlockBlobFromOptions blob_options;
blob_options.Metadata.emplace("is_directory", "true");
blob_options.Metadata.emplace(kFlatNamespaceIsDirectoryMetadataKey, "true");
block_blob_client.UploadFrom(nullptr, 0, blob_options);
}

Expand Down
60 changes: 60 additions & 0 deletions cpp/src/arrow/filesystem/azurefs_test.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -826,6 +826,41 @@ class TestAzureFileSystem : public ::testing::Test {
AssertFileInfo(fs(), subdir3, FileType::Directory);
}

void TestDisallowReadingOrWritingDirectoryMarkers() {
auto data = SetUpPreexistingData();
auto directory_path = data.Path("directory");

ASSERT_OK(fs()->CreateDir(directory_path));
ASSERT_RAISES(IOError, fs()->OpenInputFile(directory_path));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(directory_path));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(directory_path));

auto directory_path_with_slash = directory_path + "/";
ASSERT_RAISES(IOError, fs()->OpenInputFile(directory_path_with_slash));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(directory_path_with_slash));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(directory_path_with_slash));
}

void TestDisallowCreatingFileAndDirectoryWithTheSameName() {
auto data = SetUpPreexistingData();
auto path1 = data.Path("directory1");
ASSERT_OK(fs()->CreateDir(path1));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(path1));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(path1));
AssertFileInfo(fs(), path1, FileType::Directory);

auto path2 = data.Path("directory2");
ASSERT_OK(fs()->OpenOutputStream(path2));
// CreateDir returns OK even if there is already a file or directory at this
// location. Whether or not this is the desired behaviour is debatable.
ASSERT_OK(fs()->CreateDir(path2));
AssertFileInfo(fs(), path2, FileType::File);
}

void TestOpenOutputStreamWithMissingContainer() {
ASSERT_RAISES(IOError, fs()->OpenOutputStream("not-a-container/file", {}));
}

void TestDeleteDirSuccessEmpty() {
if (HasSubmitBatchBug()) {
GTEST_SKIP() << kSubmitBatchBugMessage;
Expand DownExpand Up@@ -1559,6 +1594,19 @@ TYPED_TEST(TestAzureFileSystemOnAllScenarios, CreateDirOnMissingContainer) {
this->TestCreateDirOnMissingContainer();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, DisallowReadingOrWritingDirectoryMarkers) {
this->TestDisallowReadingOrWritingDirectoryMarkers();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios,
DisallowCreatingFileAndDirectoryWithTheSameName) {
this->TestDisallowCreatingFileAndDirectoryWithTheSameName();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, OpenOutputStreamWithMissingContainer) {
this->TestOpenOutputStreamWithMissingContainer();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, DeleteDirSuccessEmpty) {
this->TestDeleteDirSuccessEmpty();
}
Expand DownExpand Up@@ -2162,6 +2210,18 @@ TEST_F(TestAzuriteFileSystem, WriteMetadata) {
.Value.Metadata;
// Defaults are overwritten and not merged.
EXPECT_EQ(Core::CaseInsensitiveMap{std::make_pair("bar", "foo")}, blob_metadata);

// Metadata can be written without writing any data.
ASSERT_OK_AND_ASSIGN(
output, fs_with_defaults->OpenAppendStream(
full_path, /*metadata=*/arrow::key_value_metadata({{"bar", "baz"}})));
ASSERT_OK(output->Close());
blob_metadata = blob_service_client_->GetBlobContainerClient(data.container_name)
.GetBlockBlobClient(blob_path)
.GetProperties()
.Value.Metadata;
// Defaults are overwritten and not merged.
EXPECT_EQ(Core::CaseInsensitiveMap{std::make_pair("bar", "baz")}, blob_metadata);
}

TEST_F(TestAzuriteFileSystem, OpenOutputStreamSmall) {
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
95 changes: 75 additions & 20 deletions cpp/src/arrow/filesystem/azurefs.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -347,6 +347,22 @@ bool IsContainerNotFound(const Storage::StorageException& e) {
return false;
}

const auto kHierarchicalNamespaceIsDirectoryMetadataKey = "hdi_isFolder";
const auto kFlatNamespaceIsDirectoryMetadataKey = "is_directory";

bool MetadataIndicatesIsDirectory(const Storage::Metadata& metadata) {
// Inspired by
// https://github.com/Azure/azure-sdk-for-cpp/blob/12407e8bfcb9bc1aa43b253c1d0ec93bf795ae3b/sdk/storage/azure-storage-files-datalake/src/datalake_utilities.cpp#L86-L91
auto hierarchical_directory_metadata =
metadata.find(kHierarchicalNamespaceIsDirectoryMetadataKey);
if (hierarchical_directory_metadata != metadata.end()) {
return hierarchical_directory_metadata->second == "true";
}
auto flat_directory_metadata = metadata.find(kFlatNamespaceIsDirectoryMetadataKey);
return flat_directory_metadata != metadata.end() &&
flat_directory_metadata->second == "true";
}

template <typename ArrowType>
std::string FormatValue(typename TypeTraits<ArrowType>::CType value) {
struct StringAppender {
Expand DownExpand Up@@ -512,11 +528,18 @@ class ObjectInputFile final : public io::RandomAccessFile {

Status Init() {
if (content_length_ != kNoSize) {
// When the user provides the file size we don't validate that its a file. This is
// only a read so its not a big deal if the user makes a mistake.
DCHECK_GE(content_length_, 0);
return Status::OK();
}
try {
// To open an ObjectInputFile the Blob must exist and it must not represent
// a directory. Additionally we need to know the file size.
auto properties = blob_client_->GetProperties();
if (MetadataIndicatesIsDirectory(properties.Value.Metadata)) {
return NotAFile(location_);
}
content_length_ = properties.Value.BlobSize;
metadata_ = PropertiesToMetadata(properties.Value);
return Status::OK();
Expand DownExpand Up@@ -698,11 +721,10 @@ class ObjectAppendStream final : public io::OutputStream {
ObjectAppendStream(std::shared_ptr<Blobs::BlockBlobClient> block_blob_client,
const io::IOContext& io_context, const AzureLocation& location,
const std::shared_ptr<const KeyValueMetadata>& metadata,
const AzureOptions& options, int64_t size = kNoSize)
const AzureOptions& options)
: block_blob_client_(std::move(block_blob_client)),
io_context_(io_context),
location_(location),
content_length_(size) {
location_(location) {
if (metadata && metadata->size() != 0) {
metadata_ = ArrowMetadataToAzureMetadata(metadata);
} else if (options.default_metadata && options.default_metadata->size() != 0) {
Expand All@@ -716,17 +738,31 @@ class ObjectAppendStream final : public io::OutputStream {
io::internal::CloseFromDestructor(this);
}

Status Init() {
if (content_length_ != kNoSize) {
DCHECK_GE(content_length_, 0);
pos_ = content_length_;
Status Init(const bool truncate,
std::function<Status()> ensure_not_flat_namespace_directory) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

You can inject AzureFileSystem *azure_file_system here and not have to allocate a closure for this. You would call AzureFileSystem::Impl::EnsureNotFlatNamespaceDirectory(location) via azure_file_system->impl_ (accessible because the handles produced by the azure file system can be friends with the filesystem class).

@Tom-NewtonTom-NewtonFeb 20, 2024

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the extra info. I was planning to do this but I was struggling with the friends thing.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The trick is to forward-declare the class as well.

diff --git a/cpp/src/arrow/filesystem/azurefs.h b/cpp/src/arrow/filesystem/azurefs.h
index 2a131e40c..d48ef9dd7 100644
--- a/cpp/src/arrow/filesystem/azurefs.h
+++ b/cpp/src/arrow/filesystem/azurefs.h
@@ -44,6 +44,7 @@ classDataLakeServiceClient;
namespacearrow::fs {
+classObjectAppendStream;
classTestAzureFileSystem;
/// Options for the AzureFileSystem implementation.
@@ -180,6 +181,7 @@ classARROW_EXPORT AzureFileSystem : public FileSystem {
explicitAzureFileSystem(std::unique_ptr<Impl>&& impl);
+ friendclassObjectAppendStream;
friendclassTestAzureFileSystem;
voidForceCachedHierarchicalNamespaceSupport(int hns_support);

@Tom-NewtonTom-NewtonFeb 21, 2024

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think my main problem was that ObjectAppendStream is defined inside an anonymous namespace but I still haven't got it working as you describe.

Are you suggesting to use AzureFileSystem *azure_file_system or AzureFileSystem:Impl *azure_file_system as the argument to ObjectAppendStream::Impl. I don't know how I can get a AzureFileSystem pointer from inside AzureFileSystem::Impl and using AzureFileSystem::Impl as the argument leads to incomplete type errors which I don't think I can avoid.

Also if you wouldn't mind I would be interested to know what the disadvantage of a lambda function is compared to what you proposed.

Sorry about my lacking C++ knowledge here.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Also if you wouldn't mind I would be interested to know what the disadvantage of a lambda function is compared to what you proposed.

To create the std::function, you heap allocate an object with copies of the values in the capture list and generate a lot more extra code in the binary:

class function {
T valuesfromthecpapturelist;
RetType operator()(ArgsType ...) {...};
}

When you think about an std::function this way (a pair of context data and a function), you realize the class you already serves that purpose.

But hey, this is becoming challenging, so I won't hold the PR anymore because of this. Moving to Init() was a big step in the right direction.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for explaining

if (truncate) {
content_length_ = 0;
pos_ = 0;
// We need to create an empty file overwriting any existing file, but
// fail if there is an existing directory.
RETURN_NOT_OK(ensure_not_flat_namespace_directory());
// On hierarchical namespace CreateEmptyBlockBlob will fail if there is an existing
// directory so we don't need to check like we do on flat namespace.
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client_));
} else {
try {
auto properties = block_blob_client_->GetProperties();
if (MetadataIndicatesIsDirectory(properties.Value.Metadata)) {
return NotAFile(location_);
}
content_length_ = properties.Value.BlobSize;
pos_ = content_length_;
} catch (const Storage::StorageException& exception) {
if (exception.StatusCode == Http::HttpStatusCode::NotFound) {
// No file exists but on flat namespace its possible there is a directory
// marker or an implied directory. Ensure there is no directory before starting
// a new empty file.
RETURN_NOT_OK(ensure_not_flat_namespace_directory());
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client_));
} else {
return ExceptionToStatus(
Expand All@@ -743,6 +779,7 @@ class ObjectAppendStream final : public io::OutputStream {
block_ids_.push_back(block.Name);
}
}
initialised_ = true;
return Status::OK();
}

Expand DownExpand Up@@ -789,6 +826,11 @@ class ObjectAppendStream final : public io::OutputStream {

Status Flush() override {
RETURN_NOT_OK(CheckClosed("flush"));
if (!initialised_) {
// If the stream has not been successfully initialized then there is nothing to
// flush. This also avoids some unhandled errors when flushing in the destructor.
return Status::OK();
}
return CommitBlockList(block_blob_client_, block_ids_, metadata_);
}

Expand DownExpand Up@@ -840,10 +882,11 @@ class ObjectAppendStream final : public io::OutputStream {
std::shared_ptr<Blobs::BlockBlobClient> block_blob_client_;
const io::IOContext io_context_;
const AzureLocation location_;
int64_t content_length_ = kNoSize;

bool closed_ = false;
bool initialised_ = false;
int64_t pos_ = 0;
int64_t content_length_ = kNoSize;
std::vector<std::string> block_ids_;
Storage::Metadata metadata_;
};
Expand DownExpand Up@@ -1662,20 +1705,32 @@ class AzureFileSystem::Impl {
AzureFileSystem* fs) {
RETURN_NOT_OK(ValidateFileLocation(location));

const auto blob_container_client = GetBlobContainerClient(location.container);
auto block_blob_client = std::make_shared<Blobs::BlockBlobClient>(
blob_service_client_->GetBlobContainerClient(location.container)
.GetBlockBlobClient(location.path));
blob_container_client.GetBlockBlobClient(location.path));

auto ensure_not_flat_namespace_directory = [this, location,
blob_container_client]() -> Status {
ARROW_ASSIGN_OR_RAISE(
auto hns_support,
HierarchicalNamespaceSupport(GetFileSystemClient(location.container)));
if (hns_support == HNSSupport::kDisabled) {
// Flat namespace so we need to GetFileInfo in-case its a directory.
ARROW_ASSIGN_OR_RAISE(auto status, GetFileInfo(blob_container_client, location))
if (status.type() == FileType::Directory) {
return NotAFile(location);
}
}
// kContainerNotFound - it doesn't exist, so no need to check if its a directory.
// kEnabled - hierarchical namespace so Azure APIs will fail if its a directory. We
// don't need to explicitly check.
return Status::OK();
};

std::shared_ptr<ObjectAppendStream> stream;
if (truncate) {
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client));
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_, 0);
} else {
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_);
}
RETURN_NOT_OK(stream->Init());
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_);
RETURN_NOT_OK(stream->Init(truncate, ensure_not_flat_namespace_directory));
return stream;
}

Expand All@@ -1690,7 +1745,7 @@ class AzureFileSystem::Impl {
// on directory marker blobs.
// https://github.com/fsspec/adlfs/blob/32132c4094350fca2680155a5c236f2e9f991ba5/adlfs/spec.py#L855-L870
Blobs::UploadBlockBlobFromOptions blob_options;
blob_options.Metadata.emplace("is_directory", "true");
blob_options.Metadata.emplace(kFlatNamespaceIsDirectoryMetadataKey, "true");
block_blob_client.UploadFrom(nullptr, 0, blob_options);
}

Expand Down
60 changes: 60 additions & 0 deletions cpp/src/arrow/filesystem/azurefs_test.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -826,6 +826,41 @@ class TestAzureFileSystem : public ::testing::Test {
AssertFileInfo(fs(), subdir3, FileType::Directory);
}

void TestDisallowReadingOrWritingDirectoryMarkers() {
auto data = SetUpPreexistingData();
auto directory_path = data.Path("directory");

ASSERT_OK(fs()->CreateDir(directory_path));
ASSERT_RAISES(IOError, fs()->OpenInputFile(directory_path));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(directory_path));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(directory_path));

auto directory_path_with_slash = directory_path + "/";
ASSERT_RAISES(IOError, fs()->OpenInputFile(directory_path_with_slash));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(directory_path_with_slash));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(directory_path_with_slash));
}

void TestDisallowCreatingFileAndDirectoryWithTheSameName() {
auto data = SetUpPreexistingData();
auto path1 = data.Path("directory1");
ASSERT_OK(fs()->CreateDir(path1));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(path1));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(path1));
AssertFileInfo(fs(), path1, FileType::Directory);

auto path2 = data.Path("directory2");
ASSERT_OK(fs()->OpenOutputStream(path2));
// CreateDir returns OK even if there is already a file or directory at this
// location. Whether or not this is the desired behaviour is debatable.
ASSERT_OK(fs()->CreateDir(path2));
AssertFileInfo(fs(), path2, FileType::File);
}

void TestOpenOutputStreamWithMissingContainer() {
ASSERT_RAISES(IOError, fs()->OpenOutputStream("not-a-container/file", {}));
}

void TestDeleteDirSuccessEmpty() {
if (HasSubmitBatchBug()) {
GTEST_SKIP() << kSubmitBatchBugMessage;
Expand DownExpand Up@@ -1559,6 +1594,19 @@ TYPED_TEST(TestAzureFileSystemOnAllScenarios, CreateDirOnMissingContainer) {
this->TestCreateDirOnMissingContainer();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, DisallowReadingOrWritingDirectoryMarkers) {
this->TestDisallowReadingOrWritingDirectoryMarkers();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios,
DisallowCreatingFileAndDirectoryWithTheSameName) {
this->TestDisallowCreatingFileAndDirectoryWithTheSameName();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, OpenOutputStreamWithMissingContainer) {
this->TestOpenOutputStreamWithMissingContainer();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, DeleteDirSuccessEmpty) {
this->TestDeleteDirSuccessEmpty();
}
Expand DownExpand Up@@ -2162,6 +2210,18 @@ TEST_F(TestAzuriteFileSystem, WriteMetadata) {
.Value.Metadata;
// Defaults are overwritten and not merged.
EXPECT_EQ(Core::CaseInsensitiveMap{std::make_pair("bar", "foo")}, blob_metadata);

// Metadata can be written without writing any data.
ASSERT_OK_AND_ASSIGN(
output, fs_with_defaults->OpenAppendStream(
full_path, /*metadata=*/arrow::key_value_metadata({{"bar", "baz"}})));
ASSERT_OK(output->Close());
blob_metadata = blob_service_client_->GetBlobContainerClient(data.container_name)
.GetBlockBlobClient(blob_path)
.GetProperties()
.Value.Metadata;
// Defaults are overwritten and not merged.
EXPECT_EQ(Core::CaseInsensitiveMap{std::make_pair("bar", "baz")}, blob_metadata);
}

TEST_F(TestAzuriteFileSystem, OpenOutputStreamSmall) {
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length \u003e 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
95 changes: 75 additions & 20 deletions cpp/src/arrow/filesystem/azurefs.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -347,6 +347,22 @@ bool IsContainerNotFound(const Storage::StorageException& e) {
return false;
}

const auto kHierarchicalNamespaceIsDirectoryMetadataKey = "hdi_isFolder";
const auto kFlatNamespaceIsDirectoryMetadataKey = "is_directory";

bool MetadataIndicatesIsDirectory(const Storage::Metadata& metadata) {
// Inspired by
// https://github.com/Azure/azure-sdk-for-cpp/blob/12407e8bfcb9bc1aa43b253c1d0ec93bf795ae3b/sdk/storage/azure-storage-files-datalake/src/datalake_utilities.cpp#L86-L91
auto hierarchical_directory_metadata =
metadata.find(kHierarchicalNamespaceIsDirectoryMetadataKey);
if (hierarchical_directory_metadata != metadata.end()) {
return hierarchical_directory_metadata->second == "true";
}
auto flat_directory_metadata = metadata.find(kFlatNamespaceIsDirectoryMetadataKey);
return flat_directory_metadata != metadata.end() &&
flat_directory_metadata->second == "true";
}

template <typename ArrowType>
std::string FormatValue(typename TypeTraits<ArrowType>::CType value) {
struct StringAppender {
Expand DownExpand Up@@ -512,11 +528,18 @@ class ObjectInputFile final : public io::RandomAccessFile {

Status Init() {
if (content_length_ != kNoSize) {
// When the user provides the file size we don't validate that its a file. This is
// only a read so its not a big deal if the user makes a mistake.
DCHECK_GE(content_length_, 0);
return Status::OK();
}
try {
// To open an ObjectInputFile the Blob must exist and it must not represent
// a directory. Additionally we need to know the file size.
auto properties = blob_client_->GetProperties();
if (MetadataIndicatesIsDirectory(properties.Value.Metadata)) {
return NotAFile(location_);
}
content_length_ = properties.Value.BlobSize;
metadata_ = PropertiesToMetadata(properties.Value);
return Status::OK();
Expand DownExpand Up@@ -698,11 +721,10 @@ class ObjectAppendStream final : public io::OutputStream {
ObjectAppendStream(std::shared_ptr<Blobs::BlockBlobClient> block_blob_client,
const io::IOContext& io_context, const AzureLocation& location,
const std::shared_ptr<const KeyValueMetadata>& metadata,
const AzureOptions& options, int64_t size = kNoSize)
const AzureOptions& options)
: block_blob_client_(std::move(block_blob_client)),
io_context_(io_context),
location_(location),
content_length_(size) {
location_(location) {
if (metadata && metadata->size() != 0) {
metadata_ = ArrowMetadataToAzureMetadata(metadata);
} else if (options.default_metadata && options.default_metadata->size() != 0) {
Expand All@@ -716,17 +738,31 @@ class ObjectAppendStream final : public io::OutputStream {
io::internal::CloseFromDestructor(this);
}

Status Init() {
if (content_length_ != kNoSize) {
DCHECK_GE(content_length_, 0);
pos_ = content_length_;
Status Init(const bool truncate,
std::function<Status()> ensure_not_flat_namespace_directory) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

You can inject AzureFileSystem *azure_file_system here and not have to allocate a closure for this. You would call AzureFileSystem::Impl::EnsureNotFlatNamespaceDirectory(location) via azure_file_system->impl_ (accessible because the handles produced by the azure file system can be friends with the filesystem class).

@Tom-NewtonTom-NewtonFeb 20, 2024

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the extra info. I was planning to do this but I was struggling with the friends thing.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The trick is to forward-declare the class as well.

diff --git a/cpp/src/arrow/filesystem/azurefs.h b/cpp/src/arrow/filesystem/azurefs.h
index 2a131e40c..d48ef9dd7 100644
--- a/cpp/src/arrow/filesystem/azurefs.h
+++ b/cpp/src/arrow/filesystem/azurefs.h
@@ -44,6 +44,7 @@ classDataLakeServiceClient;
namespacearrow::fs {
+classObjectAppendStream;
classTestAzureFileSystem;
/// Options for the AzureFileSystem implementation.
@@ -180,6 +181,7 @@ classARROW_EXPORT AzureFileSystem : public FileSystem {
explicitAzureFileSystem(std::unique_ptr<Impl>&& impl);
+ friendclassObjectAppendStream;
friendclassTestAzureFileSystem;
voidForceCachedHierarchicalNamespaceSupport(int hns_support);

@Tom-NewtonTom-NewtonFeb 21, 2024

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think my main problem was that ObjectAppendStream is defined inside an anonymous namespace but I still haven't got it working as you describe.

Are you suggesting to use AzureFileSystem *azure_file_system or AzureFileSystem:Impl *azure_file_system as the argument to ObjectAppendStream::Impl. I don't know how I can get a AzureFileSystem pointer from inside AzureFileSystem::Impl and using AzureFileSystem::Impl as the argument leads to incomplete type errors which I don't think I can avoid.

Also if you wouldn't mind I would be interested to know what the disadvantage of a lambda function is compared to what you proposed.

Sorry about my lacking C++ knowledge here.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Also if you wouldn't mind I would be interested to know what the disadvantage of a lambda function is compared to what you proposed.

To create the std::function, you heap allocate an object with copies of the values in the capture list and generate a lot more extra code in the binary:

class function {
T valuesfromthecpapturelist;
RetType operator()(ArgsType ...) {...};
}

When you think about an std::function this way (a pair of context data and a function), you realize the class you already serves that purpose.

But hey, this is becoming challenging, so I won't hold the PR anymore because of this. Moving to Init() was a big step in the right direction.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for explaining

if (truncate) {
content_length_ = 0;
pos_ = 0;
// We need to create an empty file overwriting any existing file, but
// fail if there is an existing directory.
RETURN_NOT_OK(ensure_not_flat_namespace_directory());
// On hierarchical namespace CreateEmptyBlockBlob will fail if there is an existing
// directory so we don't need to check like we do on flat namespace.
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client_));
} else {
try {
auto properties = block_blob_client_->GetProperties();
if (MetadataIndicatesIsDirectory(properties.Value.Metadata)) {
return NotAFile(location_);
}
content_length_ = properties.Value.BlobSize;
pos_ = content_length_;
} catch (const Storage::StorageException& exception) {
if (exception.StatusCode == Http::HttpStatusCode::NotFound) {
// No file exists but on flat namespace its possible there is a directory
// marker or an implied directory. Ensure there is no directory before starting
// a new empty file.
RETURN_NOT_OK(ensure_not_flat_namespace_directory());
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client_));
} else {
return ExceptionToStatus(
Expand All@@ -743,6 +779,7 @@ class ObjectAppendStream final : public io::OutputStream {
block_ids_.push_back(block.Name);
}
}
initialised_ = true;
return Status::OK();
}

Expand DownExpand Up@@ -789,6 +826,11 @@ class ObjectAppendStream final : public io::OutputStream {

Status Flush() override {
RETURN_NOT_OK(CheckClosed("flush"));
if (!initialised_) {
// If the stream has not been successfully initialized then there is nothing to
// flush. This also avoids some unhandled errors when flushing in the destructor.
return Status::OK();
}
return CommitBlockList(block_blob_client_, block_ids_, metadata_);
}

Expand DownExpand Up@@ -840,10 +882,11 @@ class ObjectAppendStream final : public io::OutputStream {
std::shared_ptr<Blobs::BlockBlobClient> block_blob_client_;
const io::IOContext io_context_;
const AzureLocation location_;
int64_t content_length_ = kNoSize;

bool closed_ = false;
bool initialised_ = false;
int64_t pos_ = 0;
int64_t content_length_ = kNoSize;
std::vector<std::string> block_ids_;
Storage::Metadata metadata_;
};
Expand DownExpand Up@@ -1662,20 +1705,32 @@ class AzureFileSystem::Impl {
AzureFileSystem* fs) {
RETURN_NOT_OK(ValidateFileLocation(location));

const auto blob_container_client = GetBlobContainerClient(location.container);
auto block_blob_client = std::make_shared<Blobs::BlockBlobClient>(
blob_service_client_->GetBlobContainerClient(location.container)
.GetBlockBlobClient(location.path));
blob_container_client.GetBlockBlobClient(location.path));

auto ensure_not_flat_namespace_directory = [this, location,
blob_container_client]() -> Status {
ARROW_ASSIGN_OR_RAISE(
auto hns_support,
HierarchicalNamespaceSupport(GetFileSystemClient(location.container)));
if (hns_support == HNSSupport::kDisabled) {
// Flat namespace so we need to GetFileInfo in-case its a directory.
ARROW_ASSIGN_OR_RAISE(auto status, GetFileInfo(blob_container_client, location))
if (status.type() == FileType::Directory) {
return NotAFile(location);
}
}
// kContainerNotFound - it doesn't exist, so no need to check if its a directory.
// kEnabled - hierarchical namespace so Azure APIs will fail if its a directory. We
// don't need to explicitly check.
return Status::OK();
};

std::shared_ptr<ObjectAppendStream> stream;
if (truncate) {
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client));
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_, 0);
} else {
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_);
}
RETURN_NOT_OK(stream->Init());
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_);
RETURN_NOT_OK(stream->Init(truncate, ensure_not_flat_namespace_directory));
return stream;
}

Expand All@@ -1690,7 +1745,7 @@ class AzureFileSystem::Impl {
// on directory marker blobs.
// https://github.com/fsspec/adlfs/blob/32132c4094350fca2680155a5c236f2e9f991ba5/adlfs/spec.py#L855-L870
Blobs::UploadBlockBlobFromOptions blob_options;
blob_options.Metadata.emplace("is_directory", "true");
blob_options.Metadata.emplace(kFlatNamespaceIsDirectoryMetadataKey, "true");
block_blob_client.UploadFrom(nullptr, 0, blob_options);
}

Expand Down
60 changes: 60 additions & 0 deletions cpp/src/arrow/filesystem/azurefs_test.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -826,6 +826,41 @@ class TestAzureFileSystem : public ::testing::Test {
AssertFileInfo(fs(), subdir3, FileType::Directory);
}

void TestDisallowReadingOrWritingDirectoryMarkers() {
auto data = SetUpPreexistingData();
auto directory_path = data.Path("directory");

ASSERT_OK(fs()->CreateDir(directory_path));
ASSERT_RAISES(IOError, fs()->OpenInputFile(directory_path));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(directory_path));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(directory_path));

auto directory_path_with_slash = directory_path + "/";
ASSERT_RAISES(IOError, fs()->OpenInputFile(directory_path_with_slash));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(directory_path_with_slash));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(directory_path_with_slash));
}

void TestDisallowCreatingFileAndDirectoryWithTheSameName() {
auto data = SetUpPreexistingData();
auto path1 = data.Path("directory1");
ASSERT_OK(fs()->CreateDir(path1));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(path1));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(path1));
AssertFileInfo(fs(), path1, FileType::Directory);

auto path2 = data.Path("directory2");
ASSERT_OK(fs()->OpenOutputStream(path2));
// CreateDir returns OK even if there is already a file or directory at this
// location. Whether or not this is the desired behaviour is debatable.
ASSERT_OK(fs()->CreateDir(path2));
AssertFileInfo(fs(), path2, FileType::File);
}

void TestOpenOutputStreamWithMissingContainer() {
ASSERT_RAISES(IOError, fs()->OpenOutputStream("not-a-container/file", {}));
}

void TestDeleteDirSuccessEmpty() {
if (HasSubmitBatchBug()) {
GTEST_SKIP() << kSubmitBatchBugMessage;
Expand DownExpand Up@@ -1559,6 +1594,19 @@ TYPED_TEST(TestAzureFileSystemOnAllScenarios, CreateDirOnMissingContainer) {
this->TestCreateDirOnMissingContainer();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, DisallowReadingOrWritingDirectoryMarkers) {
this->TestDisallowReadingOrWritingDirectoryMarkers();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios,
DisallowCreatingFileAndDirectoryWithTheSameName) {
this->TestDisallowCreatingFileAndDirectoryWithTheSameName();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, OpenOutputStreamWithMissingContainer) {
this->TestOpenOutputStreamWithMissingContainer();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, DeleteDirSuccessEmpty) {
this->TestDeleteDirSuccessEmpty();
}
Expand DownExpand Up@@ -2162,6 +2210,18 @@ TEST_F(TestAzuriteFileSystem, WriteMetadata) {
.Value.Metadata;
// Defaults are overwritten and not merged.
EXPECT_EQ(Core::CaseInsensitiveMap{std::make_pair("bar", "foo")}, blob_metadata);

// Metadata can be written without writing any data.
ASSERT_OK_AND_ASSIGN(
output, fs_with_defaults->OpenAppendStream(
full_path, /*metadata=*/arrow::key_value_metadata({{"bar", "baz"}})));
ASSERT_OK(output->Close());
blob_metadata = blob_service_client_->GetBlobContainerClient(data.container_name)
.GetBlockBlobClient(blob_path)
.GetProperties()
.Value.Metadata;
// Defaults are overwritten and not merged.
EXPECT_EQ(Core::CaseInsensitiveMap{std::make_pair("bar", "baz")}, blob_metadata);
}

TEST_F(TestAzuriteFileSystem, OpenOutputStreamSmall) {
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
95 changes: 75 additions & 20 deletions cpp/src/arrow/filesystem/azurefs.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -347,6 +347,22 @@ bool IsContainerNotFound(const Storage::StorageException& e) {
return false;
}

const auto kHierarchicalNamespaceIsDirectoryMetadataKey = "hdi_isFolder";
const auto kFlatNamespaceIsDirectoryMetadataKey = "is_directory";

bool MetadataIndicatesIsDirectory(const Storage::Metadata& metadata) {
// Inspired by
// https://github.com/Azure/azure-sdk-for-cpp/blob/12407e8bfcb9bc1aa43b253c1d0ec93bf795ae3b/sdk/storage/azure-storage-files-datalake/src/datalake_utilities.cpp#L86-L91
auto hierarchical_directory_metadata =
metadata.find(kHierarchicalNamespaceIsDirectoryMetadataKey);
if (hierarchical_directory_metadata != metadata.end()) {
return hierarchical_directory_metadata->second == "true";
}
auto flat_directory_metadata = metadata.find(kFlatNamespaceIsDirectoryMetadataKey);
return flat_directory_metadata != metadata.end() &&
flat_directory_metadata->second == "true";
}

template <typename ArrowType>
std::string FormatValue(typename TypeTraits<ArrowType>::CType value) {
struct StringAppender {
Expand DownExpand Up@@ -512,11 +528,18 @@ class ObjectInputFile final : public io::RandomAccessFile {

Status Init() {
if (content_length_ != kNoSize) {
// When the user provides the file size we don't validate that its a file. This is
// only a read so its not a big deal if the user makes a mistake.
DCHECK_GE(content_length_, 0);
return Status::OK();
}
try {
// To open an ObjectInputFile the Blob must exist and it must not represent
// a directory. Additionally we need to know the file size.
auto properties = blob_client_->GetProperties();
if (MetadataIndicatesIsDirectory(properties.Value.Metadata)) {
return NotAFile(location_);
}
content_length_ = properties.Value.BlobSize;
metadata_ = PropertiesToMetadata(properties.Value);
return Status::OK();
Expand DownExpand Up@@ -698,11 +721,10 @@ class ObjectAppendStream final : public io::OutputStream {
ObjectAppendStream(std::shared_ptr<Blobs::BlockBlobClient> block_blob_client,
const io::IOContext& io_context, const AzureLocation& location,
const std::shared_ptr<const KeyValueMetadata>& metadata,
const AzureOptions& options, int64_t size = kNoSize)
const AzureOptions& options)
: block_blob_client_(std::move(block_blob_client)),
io_context_(io_context),
location_(location),
content_length_(size) {
location_(location) {
if (metadata && metadata->size() != 0) {
metadata_ = ArrowMetadataToAzureMetadata(metadata);
} else if (options.default_metadata && options.default_metadata->size() != 0) {
Expand All@@ -716,17 +738,31 @@ class ObjectAppendStream final : public io::OutputStream {
io::internal::CloseFromDestructor(this);
}

Status Init() {
if (content_length_ != kNoSize) {
DCHECK_GE(content_length_, 0);
pos_ = content_length_;
Status Init(const bool truncate,
std::function<Status()> ensure_not_flat_namespace_directory) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

You can inject AzureFileSystem *azure_file_system here and not have to allocate a closure for this. You would call AzureFileSystem::Impl::EnsureNotFlatNamespaceDirectory(location) via azure_file_system->impl_ (accessible because the handles produced by the azure file system can be friends with the filesystem class).

@Tom-NewtonTom-NewtonFeb 20, 2024

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the extra info. I was planning to do this but I was struggling with the friends thing.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The trick is to forward-declare the class as well.

diff --git a/cpp/src/arrow/filesystem/azurefs.h b/cpp/src/arrow/filesystem/azurefs.h
index 2a131e40c..d48ef9dd7 100644
--- a/cpp/src/arrow/filesystem/azurefs.h
+++ b/cpp/src/arrow/filesystem/azurefs.h
@@ -44,6 +44,7 @@ classDataLakeServiceClient;
namespacearrow::fs {
+classObjectAppendStream;
classTestAzureFileSystem;
/// Options for the AzureFileSystem implementation.
@@ -180,6 +181,7 @@ classARROW_EXPORT AzureFileSystem : public FileSystem {
explicitAzureFileSystem(std::unique_ptr<Impl>&& impl);
+ friendclassObjectAppendStream;
friendclassTestAzureFileSystem;
voidForceCachedHierarchicalNamespaceSupport(int hns_support);

@Tom-NewtonTom-NewtonFeb 21, 2024

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think my main problem was that ObjectAppendStream is defined inside an anonymous namespace but I still haven't got it working as you describe.

Are you suggesting to use AzureFileSystem *azure_file_system or AzureFileSystem:Impl *azure_file_system as the argument to ObjectAppendStream::Impl. I don't know how I can get a AzureFileSystem pointer from inside AzureFileSystem::Impl and using AzureFileSystem::Impl as the argument leads to incomplete type errors which I don't think I can avoid.

Also if you wouldn't mind I would be interested to know what the disadvantage of a lambda function is compared to what you proposed.

Sorry about my lacking C++ knowledge here.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Also if you wouldn't mind I would be interested to know what the disadvantage of a lambda function is compared to what you proposed.

To create the std::function, you heap allocate an object with copies of the values in the capture list and generate a lot more extra code in the binary:

class function {
T valuesfromthecpapturelist;
RetType operator()(ArgsType ...) {...};
}

When you think about an std::function this way (a pair of context data and a function), you realize the class you already serves that purpose.

But hey, this is becoming challenging, so I won't hold the PR anymore because of this. Moving to Init() was a big step in the right direction.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for explaining

if (truncate) {
content_length_ = 0;
pos_ = 0;
// We need to create an empty file overwriting any existing file, but
// fail if there is an existing directory.
RETURN_NOT_OK(ensure_not_flat_namespace_directory());
// On hierarchical namespace CreateEmptyBlockBlob will fail if there is an existing
// directory so we don't need to check like we do on flat namespace.
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client_));
} else {
try {
auto properties = block_blob_client_->GetProperties();
if (MetadataIndicatesIsDirectory(properties.Value.Metadata)) {
return NotAFile(location_);
}
content_length_ = properties.Value.BlobSize;
pos_ = content_length_;
} catch (const Storage::StorageException& exception) {
if (exception.StatusCode == Http::HttpStatusCode::NotFound) {
// No file exists but on flat namespace its possible there is a directory
// marker or an implied directory. Ensure there is no directory before starting
// a new empty file.
RETURN_NOT_OK(ensure_not_flat_namespace_directory());
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client_));
} else {
return ExceptionToStatus(
Expand All@@ -743,6 +779,7 @@ class ObjectAppendStream final : public io::OutputStream {
block_ids_.push_back(block.Name);
}
}
initialised_ = true;
return Status::OK();
}

Expand DownExpand Up@@ -789,6 +826,11 @@ class ObjectAppendStream final : public io::OutputStream {

Status Flush() override {
RETURN_NOT_OK(CheckClosed("flush"));
if (!initialised_) {
// If the stream has not been successfully initialized then there is nothing to
// flush. This also avoids some unhandled errors when flushing in the destructor.
return Status::OK();
}
return CommitBlockList(block_blob_client_, block_ids_, metadata_);
}

Expand DownExpand Up@@ -840,10 +882,11 @@ class ObjectAppendStream final : public io::OutputStream {
std::shared_ptr<Blobs::BlockBlobClient> block_blob_client_;
const io::IOContext io_context_;
const AzureLocation location_;
int64_t content_length_ = kNoSize;

bool closed_ = false;
bool initialised_ = false;
int64_t pos_ = 0;
int64_t content_length_ = kNoSize;
std::vector<std::string> block_ids_;
Storage::Metadata metadata_;
};
Expand DownExpand Up@@ -1662,20 +1705,32 @@ class AzureFileSystem::Impl {
AzureFileSystem* fs) {
RETURN_NOT_OK(ValidateFileLocation(location));

const auto blob_container_client = GetBlobContainerClient(location.container);
auto block_blob_client = std::make_shared<Blobs::BlockBlobClient>(
blob_service_client_->GetBlobContainerClient(location.container)
.GetBlockBlobClient(location.path));
blob_container_client.GetBlockBlobClient(location.path));

auto ensure_not_flat_namespace_directory = [this, location,
blob_container_client]() -> Status {
ARROW_ASSIGN_OR_RAISE(
auto hns_support,
HierarchicalNamespaceSupport(GetFileSystemClient(location.container)));
if (hns_support == HNSSupport::kDisabled) {
// Flat namespace so we need to GetFileInfo in-case its a directory.
ARROW_ASSIGN_OR_RAISE(auto status, GetFileInfo(blob_container_client, location))
if (status.type() == FileType::Directory) {
return NotAFile(location);
}
}
// kContainerNotFound - it doesn't exist, so no need to check if its a directory.
// kEnabled - hierarchical namespace so Azure APIs will fail if its a directory. We
// don't need to explicitly check.
return Status::OK();
};

std::shared_ptr<ObjectAppendStream> stream;
if (truncate) {
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client));
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_, 0);
} else {
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_);
}
RETURN_NOT_OK(stream->Init());
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_);
RETURN_NOT_OK(stream->Init(truncate, ensure_not_flat_namespace_directory));
return stream;
}

Expand All@@ -1690,7 +1745,7 @@ class AzureFileSystem::Impl {
// on directory marker blobs.
// https://github.com/fsspec/adlfs/blob/32132c4094350fca2680155a5c236f2e9f991ba5/adlfs/spec.py#L855-L870
Blobs::UploadBlockBlobFromOptions blob_options;
blob_options.Metadata.emplace("is_directory", "true");
blob_options.Metadata.emplace(kFlatNamespaceIsDirectoryMetadataKey, "true");
block_blob_client.UploadFrom(nullptr, 0, blob_options);
}

Expand Down
60 changes: 60 additions & 0 deletions cpp/src/arrow/filesystem/azurefs_test.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -826,6 +826,41 @@ class TestAzureFileSystem : public ::testing::Test {
AssertFileInfo(fs(), subdir3, FileType::Directory);
}

void TestDisallowReadingOrWritingDirectoryMarkers() {
auto data = SetUpPreexistingData();
auto directory_path = data.Path("directory");

ASSERT_OK(fs()->CreateDir(directory_path));
ASSERT_RAISES(IOError, fs()->OpenInputFile(directory_path));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(directory_path));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(directory_path));

auto directory_path_with_slash = directory_path + "/";
ASSERT_RAISES(IOError, fs()->OpenInputFile(directory_path_with_slash));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(directory_path_with_slash));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(directory_path_with_slash));
}

void TestDisallowCreatingFileAndDirectoryWithTheSameName() {
auto data = SetUpPreexistingData();
auto path1 = data.Path("directory1");
ASSERT_OK(fs()->CreateDir(path1));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(path1));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(path1));
AssertFileInfo(fs(), path1, FileType::Directory);

auto path2 = data.Path("directory2");
ASSERT_OK(fs()->OpenOutputStream(path2));
// CreateDir returns OK even if there is already a file or directory at this
// location. Whether or not this is the desired behaviour is debatable.
ASSERT_OK(fs()->CreateDir(path2));
AssertFileInfo(fs(), path2, FileType::File);
}

void TestOpenOutputStreamWithMissingContainer() {
ASSERT_RAISES(IOError, fs()->OpenOutputStream("not-a-container/file", {}));
}

void TestDeleteDirSuccessEmpty() {
if (HasSubmitBatchBug()) {
GTEST_SKIP() << kSubmitBatchBugMessage;
Expand DownExpand Up@@ -1559,6 +1594,19 @@ TYPED_TEST(TestAzureFileSystemOnAllScenarios, CreateDirOnMissingContainer) {
this->TestCreateDirOnMissingContainer();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, DisallowReadingOrWritingDirectoryMarkers) {
this->TestDisallowReadingOrWritingDirectoryMarkers();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios,
DisallowCreatingFileAndDirectoryWithTheSameName) {
this->TestDisallowCreatingFileAndDirectoryWithTheSameName();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, OpenOutputStreamWithMissingContainer) {
this->TestOpenOutputStreamWithMissingContainer();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, DeleteDirSuccessEmpty) {
this->TestDeleteDirSuccessEmpty();
}
Expand DownExpand Up@@ -2162,6 +2210,18 @@ TEST_F(TestAzuriteFileSystem, WriteMetadata) {
.Value.Metadata;
// Defaults are overwritten and not merged.
EXPECT_EQ(Core::CaseInsensitiveMap{std::make_pair("bar", "foo")}, blob_metadata);

// Metadata can be written without writing any data.
ASSERT_OK_AND_ASSIGN(
output, fs_with_defaults->OpenAppendStream(
full_path, /*metadata=*/arrow::key_value_metadata({{"bar", "baz"}})));
ASSERT_OK(output->Close());
blob_metadata = blob_service_client_->GetBlobContainerClient(data.container_name)
.GetBlockBlobClient(blob_path)
.GetProperties()
.Value.Metadata;
// Defaults are overwritten and not merged.
EXPECT_EQ(Core::CaseInsensitiveMap{std::make_pair("bar", "baz")}, blob_metadata);
}

TEST_F(TestAzuriteFileSystem, OpenOutputStreamSmall) {
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
95 changes: 75 additions & 20 deletions cpp/src/arrow/filesystem/azurefs.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -347,6 +347,22 @@ bool IsContainerNotFound(const Storage::StorageException& e) {
return false;
}

const auto kHierarchicalNamespaceIsDirectoryMetadataKey = "hdi_isFolder";
const auto kFlatNamespaceIsDirectoryMetadataKey = "is_directory";

bool MetadataIndicatesIsDirectory(const Storage::Metadata& metadata) {
// Inspired by
// https://github.com/Azure/azure-sdk-for-cpp/blob/12407e8bfcb9bc1aa43b253c1d0ec93bf795ae3b/sdk/storage/azure-storage-files-datalake/src/datalake_utilities.cpp#L86-L91
auto hierarchical_directory_metadata =
metadata.find(kHierarchicalNamespaceIsDirectoryMetadataKey);
if (hierarchical_directory_metadata != metadata.end()) {
return hierarchical_directory_metadata->second == "true";
}
auto flat_directory_metadata = metadata.find(kFlatNamespaceIsDirectoryMetadataKey);
return flat_directory_metadata != metadata.end() &&
flat_directory_metadata->second == "true";
}

template <typename ArrowType>
std::string FormatValue(typename TypeTraits<ArrowType>::CType value) {
struct StringAppender {
Expand DownExpand Up@@ -512,11 +528,18 @@ class ObjectInputFile final : public io::RandomAccessFile {

Status Init() {
if (content_length_ != kNoSize) {
// When the user provides the file size we don't validate that its a file. This is
// only a read so its not a big deal if the user makes a mistake.
DCHECK_GE(content_length_, 0);
return Status::OK();
}
try {
// To open an ObjectInputFile the Blob must exist and it must not represent
// a directory. Additionally we need to know the file size.
auto properties = blob_client_->GetProperties();
if (MetadataIndicatesIsDirectory(properties.Value.Metadata)) {
return NotAFile(location_);
}
content_length_ = properties.Value.BlobSize;
metadata_ = PropertiesToMetadata(properties.Value);
return Status::OK();
Expand DownExpand Up@@ -698,11 +721,10 @@ class ObjectAppendStream final : public io::OutputStream {
ObjectAppendStream(std::shared_ptr<Blobs::BlockBlobClient> block_blob_client,
const io::IOContext& io_context, const AzureLocation& location,
const std::shared_ptr<const KeyValueMetadata>& metadata,
const AzureOptions& options, int64_t size = kNoSize)
const AzureOptions& options)
: block_blob_client_(std::move(block_blob_client)),
io_context_(io_context),
location_(location),
content_length_(size) {
location_(location) {
if (metadata && metadata->size() != 0) {
metadata_ = ArrowMetadataToAzureMetadata(metadata);
} else if (options.default_metadata && options.default_metadata->size() != 0) {
Expand All@@ -716,17 +738,31 @@ class ObjectAppendStream final : public io::OutputStream {
io::internal::CloseFromDestructor(this);
}

Status Init() {
if (content_length_ != kNoSize) {
DCHECK_GE(content_length_, 0);
pos_ = content_length_;
Status Init(const bool truncate,
std::function<Status()> ensure_not_flat_namespace_directory) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

You can inject AzureFileSystem *azure_file_system here and not have to allocate a closure for this. You would call AzureFileSystem::Impl::EnsureNotFlatNamespaceDirectory(location) via azure_file_system->impl_ (accessible because the handles produced by the azure file system can be friends with the filesystem class).

@Tom-NewtonTom-NewtonFeb 20, 2024

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the extra info. I was planning to do this but I was struggling with the friends thing.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The trick is to forward-declare the class as well.

diff --git a/cpp/src/arrow/filesystem/azurefs.h b/cpp/src/arrow/filesystem/azurefs.h
index 2a131e40c..d48ef9dd7 100644
--- a/cpp/src/arrow/filesystem/azurefs.h
+++ b/cpp/src/arrow/filesystem/azurefs.h
@@ -44,6 +44,7 @@ classDataLakeServiceClient;
namespacearrow::fs {
+classObjectAppendStream;
classTestAzureFileSystem;
/// Options for the AzureFileSystem implementation.
@@ -180,6 +181,7 @@ classARROW_EXPORT AzureFileSystem : public FileSystem {
explicitAzureFileSystem(std::unique_ptr<Impl>&& impl);
+ friendclassObjectAppendStream;
friendclassTestAzureFileSystem;
voidForceCachedHierarchicalNamespaceSupport(int hns_support);

@Tom-NewtonTom-NewtonFeb 21, 2024

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think my main problem was that ObjectAppendStream is defined inside an anonymous namespace but I still haven't got it working as you describe.

Are you suggesting to use AzureFileSystem *azure_file_system or AzureFileSystem:Impl *azure_file_system as the argument to ObjectAppendStream::Impl. I don't know how I can get a AzureFileSystem pointer from inside AzureFileSystem::Impl and using AzureFileSystem::Impl as the argument leads to incomplete type errors which I don't think I can avoid.

Also if you wouldn't mind I would be interested to know what the disadvantage of a lambda function is compared to what you proposed.

Sorry about my lacking C++ knowledge here.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Also if you wouldn't mind I would be interested to know what the disadvantage of a lambda function is compared to what you proposed.

To create the std::function, you heap allocate an object with copies of the values in the capture list and generate a lot more extra code in the binary:

class function {
T valuesfromthecpapturelist;
RetType operator()(ArgsType ...) {...};
}

When you think about an std::function this way (a pair of context data and a function), you realize the class you already serves that purpose.

But hey, this is becoming challenging, so I won't hold the PR anymore because of this. Moving to Init() was a big step in the right direction.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for explaining

if (truncate) {
content_length_ = 0;
pos_ = 0;
// We need to create an empty file overwriting any existing file, but
// fail if there is an existing directory.
RETURN_NOT_OK(ensure_not_flat_namespace_directory());
// On hierarchical namespace CreateEmptyBlockBlob will fail if there is an existing
// directory so we don't need to check like we do on flat namespace.
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client_));
} else {
try {
auto properties = block_blob_client_->GetProperties();
if (MetadataIndicatesIsDirectory(properties.Value.Metadata)) {
return NotAFile(location_);
}
content_length_ = properties.Value.BlobSize;
pos_ = content_length_;
} catch (const Storage::StorageException& exception) {
if (exception.StatusCode == Http::HttpStatusCode::NotFound) {
// No file exists but on flat namespace its possible there is a directory
// marker or an implied directory. Ensure there is no directory before starting
// a new empty file.
RETURN_NOT_OK(ensure_not_flat_namespace_directory());
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client_));
} else {
return ExceptionToStatus(
Expand All@@ -743,6 +779,7 @@ class ObjectAppendStream final : public io::OutputStream {
block_ids_.push_back(block.Name);
}
}
initialised_ = true;
return Status::OK();
}

Expand DownExpand Up@@ -789,6 +826,11 @@ class ObjectAppendStream final : public io::OutputStream {

Status Flush() override {
RETURN_NOT_OK(CheckClosed("flush"));
if (!initialised_) {
// If the stream has not been successfully initialized then there is nothing to
// flush. This also avoids some unhandled errors when flushing in the destructor.
return Status::OK();
}
return CommitBlockList(block_blob_client_, block_ids_, metadata_);
}

Expand DownExpand Up@@ -840,10 +882,11 @@ class ObjectAppendStream final : public io::OutputStream {
std::shared_ptr<Blobs::BlockBlobClient> block_blob_client_;
const io::IOContext io_context_;
const AzureLocation location_;
int64_t content_length_ = kNoSize;

bool closed_ = false;
bool initialised_ = false;
int64_t pos_ = 0;
int64_t content_length_ = kNoSize;
std::vector<std::string> block_ids_;
Storage::Metadata metadata_;
};
Expand DownExpand Up@@ -1662,20 +1705,32 @@ class AzureFileSystem::Impl {
AzureFileSystem* fs) {
RETURN_NOT_OK(ValidateFileLocation(location));

const auto blob_container_client = GetBlobContainerClient(location.container);
auto block_blob_client = std::make_shared<Blobs::BlockBlobClient>(
blob_service_client_->GetBlobContainerClient(location.container)
.GetBlockBlobClient(location.path));
blob_container_client.GetBlockBlobClient(location.path));

auto ensure_not_flat_namespace_directory = [this, location,
blob_container_client]() -> Status {
ARROW_ASSIGN_OR_RAISE(
auto hns_support,
HierarchicalNamespaceSupport(GetFileSystemClient(location.container)));
if (hns_support == HNSSupport::kDisabled) {
// Flat namespace so we need to GetFileInfo in-case its a directory.
ARROW_ASSIGN_OR_RAISE(auto status, GetFileInfo(blob_container_client, location))
if (status.type() == FileType::Directory) {
return NotAFile(location);
}
}
// kContainerNotFound - it doesn't exist, so no need to check if its a directory.
// kEnabled - hierarchical namespace so Azure APIs will fail if its a directory. We
// don't need to explicitly check.
return Status::OK();
};

std::shared_ptr<ObjectAppendStream> stream;
if (truncate) {
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client));
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_, 0);
} else {
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_);
}
RETURN_NOT_OK(stream->Init());
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_);
RETURN_NOT_OK(stream->Init(truncate, ensure_not_flat_namespace_directory));
return stream;
}

Expand All@@ -1690,7 +1745,7 @@ class AzureFileSystem::Impl {
// on directory marker blobs.
// https://github.com/fsspec/adlfs/blob/32132c4094350fca2680155a5c236f2e9f991ba5/adlfs/spec.py#L855-L870
Blobs::UploadBlockBlobFromOptions blob_options;
blob_options.Metadata.emplace("is_directory", "true");
blob_options.Metadata.emplace(kFlatNamespaceIsDirectoryMetadataKey, "true");
block_blob_client.UploadFrom(nullptr, 0, blob_options);
}

Expand Down
60 changes: 60 additions & 0 deletions cpp/src/arrow/filesystem/azurefs_test.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -826,6 +826,41 @@ class TestAzureFileSystem : public ::testing::Test {
AssertFileInfo(fs(), subdir3, FileType::Directory);
}

void TestDisallowReadingOrWritingDirectoryMarkers() {
auto data = SetUpPreexistingData();
auto directory_path = data.Path("directory");

ASSERT_OK(fs()->CreateDir(directory_path));
ASSERT_RAISES(IOError, fs()->OpenInputFile(directory_path));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(directory_path));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(directory_path));

auto directory_path_with_slash = directory_path + "/";
ASSERT_RAISES(IOError, fs()->OpenInputFile(directory_path_with_slash));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(directory_path_with_slash));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(directory_path_with_slash));
}

void TestDisallowCreatingFileAndDirectoryWithTheSameName() {
auto data = SetUpPreexistingData();
auto path1 = data.Path("directory1");
ASSERT_OK(fs()->CreateDir(path1));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(path1));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(path1));
AssertFileInfo(fs(), path1, FileType::Directory);

auto path2 = data.Path("directory2");
ASSERT_OK(fs()->OpenOutputStream(path2));
// CreateDir returns OK even if there is already a file or directory at this
// location. Whether or not this is the desired behaviour is debatable.
ASSERT_OK(fs()->CreateDir(path2));
AssertFileInfo(fs(), path2, FileType::File);
}

void TestOpenOutputStreamWithMissingContainer() {
ASSERT_RAISES(IOError, fs()->OpenOutputStream("not-a-container/file", {}));
}

void TestDeleteDirSuccessEmpty() {
if (HasSubmitBatchBug()) {
GTEST_SKIP() << kSubmitBatchBugMessage;
Expand DownExpand Up@@ -1559,6 +1594,19 @@ TYPED_TEST(TestAzureFileSystemOnAllScenarios, CreateDirOnMissingContainer) {
this->TestCreateDirOnMissingContainer();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, DisallowReadingOrWritingDirectoryMarkers) {
this->TestDisallowReadingOrWritingDirectoryMarkers();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios,
DisallowCreatingFileAndDirectoryWithTheSameName) {
this->TestDisallowCreatingFileAndDirectoryWithTheSameName();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, OpenOutputStreamWithMissingContainer) {
this->TestOpenOutputStreamWithMissingContainer();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, DeleteDirSuccessEmpty) {
this->TestDeleteDirSuccessEmpty();
}
Expand DownExpand Up@@ -2162,6 +2210,18 @@ TEST_F(TestAzuriteFileSystem, WriteMetadata) {
.Value.Metadata;
// Defaults are overwritten and not merged.
EXPECT_EQ(Core::CaseInsensitiveMap{std::make_pair("bar", "foo")}, blob_metadata);

// Metadata can be written without writing any data.
ASSERT_OK_AND_ASSIGN(
output, fs_with_defaults->OpenAppendStream(
full_path, /*metadata=*/arrow::key_value_metadata({{"bar", "baz"}})));
ASSERT_OK(output->Close());
blob_metadata = blob_service_client_->GetBlobContainerClient(data.container_name)
.GetBlockBlobClient(blob_path)
.GetProperties()
.Value.Metadata;
// Defaults are overwritten and not merged.
EXPECT_EQ(Core::CaseInsensitiveMap{std::make_pair("bar", "baz")}, blob_metadata);
}

TEST_F(TestAzuriteFileSystem, OpenOutputStreamSmall) {
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
95 changes: 75 additions & 20 deletions cpp/src/arrow/filesystem/azurefs.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -347,6 +347,22 @@ bool IsContainerNotFound(const Storage::StorageException& e) {
return false;
}

const auto kHierarchicalNamespaceIsDirectoryMetadataKey = "hdi_isFolder";
const auto kFlatNamespaceIsDirectoryMetadataKey = "is_directory";

bool MetadataIndicatesIsDirectory(const Storage::Metadata& metadata) {
// Inspired by
// https://github.com/Azure/azure-sdk-for-cpp/blob/12407e8bfcb9bc1aa43b253c1d0ec93bf795ae3b/sdk/storage/azure-storage-files-datalake/src/datalake_utilities.cpp#L86-L91
auto hierarchical_directory_metadata =
metadata.find(kHierarchicalNamespaceIsDirectoryMetadataKey);
if (hierarchical_directory_metadata != metadata.end()) {
return hierarchical_directory_metadata->second == "true";
}
auto flat_directory_metadata = metadata.find(kFlatNamespaceIsDirectoryMetadataKey);
return flat_directory_metadata != metadata.end() &&
flat_directory_metadata->second == "true";
}

template <typename ArrowType>
std::string FormatValue(typename TypeTraits<ArrowType>::CType value) {
struct StringAppender {
Expand DownExpand Up@@ -512,11 +528,18 @@ class ObjectInputFile final : public io::RandomAccessFile {

Status Init() {
if (content_length_ != kNoSize) {
// When the user provides the file size we don't validate that its a file. This is
// only a read so its not a big deal if the user makes a mistake.
DCHECK_GE(content_length_, 0);
return Status::OK();
}
try {
// To open an ObjectInputFile the Blob must exist and it must not represent
// a directory. Additionally we need to know the file size.
auto properties = blob_client_->GetProperties();
if (MetadataIndicatesIsDirectory(properties.Value.Metadata)) {
return NotAFile(location_);
}
content_length_ = properties.Value.BlobSize;
metadata_ = PropertiesToMetadata(properties.Value);
return Status::OK();
Expand DownExpand Up@@ -698,11 +721,10 @@ class ObjectAppendStream final : public io::OutputStream {
ObjectAppendStream(std::shared_ptr<Blobs::BlockBlobClient> block_blob_client,
const io::IOContext& io_context, const AzureLocation& location,
const std::shared_ptr<const KeyValueMetadata>& metadata,
const AzureOptions& options, int64_t size = kNoSize)
const AzureOptions& options)
: block_blob_client_(std::move(block_blob_client)),
io_context_(io_context),
location_(location),
content_length_(size) {
location_(location) {
if (metadata && metadata->size() != 0) {
metadata_ = ArrowMetadataToAzureMetadata(metadata);
} else if (options.default_metadata && options.default_metadata->size() != 0) {
Expand All@@ -716,17 +738,31 @@ class ObjectAppendStream final : public io::OutputStream {
io::internal::CloseFromDestructor(this);
}

Status Init() {
if (content_length_ != kNoSize) {
DCHECK_GE(content_length_, 0);
pos_ = content_length_;
Status Init(const bool truncate,
std::function<Status()> ensure_not_flat_namespace_directory) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

You can inject AzureFileSystem *azure_file_system here and not have to allocate a closure for this. You would call AzureFileSystem::Impl::EnsureNotFlatNamespaceDirectory(location) via azure_file_system->impl_ (accessible because the handles produced by the azure file system can be friends with the filesystem class).

@Tom-NewtonTom-NewtonFeb 20, 2024

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the extra info. I was planning to do this but I was struggling with the friends thing.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The trick is to forward-declare the class as well.

diff --git a/cpp/src/arrow/filesystem/azurefs.h b/cpp/src/arrow/filesystem/azurefs.h
index 2a131e40c..d48ef9dd7 100644
--- a/cpp/src/arrow/filesystem/azurefs.h
+++ b/cpp/src/arrow/filesystem/azurefs.h
@@ -44,6 +44,7 @@ classDataLakeServiceClient;
namespacearrow::fs {
+classObjectAppendStream;
classTestAzureFileSystem;
/// Options for the AzureFileSystem implementation.
@@ -180,6 +181,7 @@ classARROW_EXPORT AzureFileSystem : public FileSystem {
explicitAzureFileSystem(std::unique_ptr<Impl>&& impl);
+ friendclassObjectAppendStream;
friendclassTestAzureFileSystem;
voidForceCachedHierarchicalNamespaceSupport(int hns_support);

@Tom-NewtonTom-NewtonFeb 21, 2024

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think my main problem was that ObjectAppendStream is defined inside an anonymous namespace but I still haven't got it working as you describe.

Are you suggesting to use AzureFileSystem *azure_file_system or AzureFileSystem:Impl *azure_file_system as the argument to ObjectAppendStream::Impl. I don't know how I can get a AzureFileSystem pointer from inside AzureFileSystem::Impl and using AzureFileSystem::Impl as the argument leads to incomplete type errors which I don't think I can avoid.

Also if you wouldn't mind I would be interested to know what the disadvantage of a lambda function is compared to what you proposed.

Sorry about my lacking C++ knowledge here.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Also if you wouldn't mind I would be interested to know what the disadvantage of a lambda function is compared to what you proposed.

To create the std::function, you heap allocate an object with copies of the values in the capture list and generate a lot more extra code in the binary:

class function {
T valuesfromthecpapturelist;
RetType operator()(ArgsType ...) {...};
}

When you think about an std::function this way (a pair of context data and a function), you realize the class you already serves that purpose.

But hey, this is becoming challenging, so I won't hold the PR anymore because of this. Moving to Init() was a big step in the right direction.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for explaining

if (truncate) {
content_length_ = 0;
pos_ = 0;
// We need to create an empty file overwriting any existing file, but
// fail if there is an existing directory.
RETURN_NOT_OK(ensure_not_flat_namespace_directory());
// On hierarchical namespace CreateEmptyBlockBlob will fail if there is an existing
// directory so we don't need to check like we do on flat namespace.
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client_));
} else {
try {
auto properties = block_blob_client_->GetProperties();
if (MetadataIndicatesIsDirectory(properties.Value.Metadata)) {
return NotAFile(location_);
}
content_length_ = properties.Value.BlobSize;
pos_ = content_length_;
} catch (const Storage::StorageException& exception) {
if (exception.StatusCode == Http::HttpStatusCode::NotFound) {
// No file exists but on flat namespace its possible there is a directory
// marker or an implied directory. Ensure there is no directory before starting
// a new empty file.
RETURN_NOT_OK(ensure_not_flat_namespace_directory());
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client_));
} else {
return ExceptionToStatus(
Expand All@@ -743,6 +779,7 @@ class ObjectAppendStream final : public io::OutputStream {
block_ids_.push_back(block.Name);
}
}
initialised_ = true;
return Status::OK();
}

Expand DownExpand Up@@ -789,6 +826,11 @@ class ObjectAppendStream final : public io::OutputStream {

Status Flush() override {
RETURN_NOT_OK(CheckClosed("flush"));
if (!initialised_) {
// If the stream has not been successfully initialized then there is nothing to
// flush. This also avoids some unhandled errors when flushing in the destructor.
return Status::OK();
}
return CommitBlockList(block_blob_client_, block_ids_, metadata_);
}

Expand DownExpand Up@@ -840,10 +882,11 @@ class ObjectAppendStream final : public io::OutputStream {
std::shared_ptr<Blobs::BlockBlobClient> block_blob_client_;
const io::IOContext io_context_;
const AzureLocation location_;
int64_t content_length_ = kNoSize;

bool closed_ = false;
bool initialised_ = false;
int64_t pos_ = 0;
int64_t content_length_ = kNoSize;
std::vector<std::string> block_ids_;
Storage::Metadata metadata_;
};
Expand DownExpand Up@@ -1662,20 +1705,32 @@ class AzureFileSystem::Impl {
AzureFileSystem* fs) {
RETURN_NOT_OK(ValidateFileLocation(location));

const auto blob_container_client = GetBlobContainerClient(location.container);
auto block_blob_client = std::make_shared<Blobs::BlockBlobClient>(
blob_service_client_->GetBlobContainerClient(location.container)
.GetBlockBlobClient(location.path));
blob_container_client.GetBlockBlobClient(location.path));

auto ensure_not_flat_namespace_directory = [this, location,
blob_container_client]() -> Status {
ARROW_ASSIGN_OR_RAISE(
auto hns_support,
HierarchicalNamespaceSupport(GetFileSystemClient(location.container)));
if (hns_support == HNSSupport::kDisabled) {
// Flat namespace so we need to GetFileInfo in-case its a directory.
ARROW_ASSIGN_OR_RAISE(auto status, GetFileInfo(blob_container_client, location))
if (status.type() == FileType::Directory) {
return NotAFile(location);
}
}
// kContainerNotFound - it doesn't exist, so no need to check if its a directory.
// kEnabled - hierarchical namespace so Azure APIs will fail if its a directory. We
// don't need to explicitly check.
return Status::OK();
};

std::shared_ptr<ObjectAppendStream> stream;
if (truncate) {
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client));
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_, 0);
} else {
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_);
}
RETURN_NOT_OK(stream->Init());
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_);
RETURN_NOT_OK(stream->Init(truncate, ensure_not_flat_namespace_directory));
return stream;
}

Expand All@@ -1690,7 +1745,7 @@ class AzureFileSystem::Impl {
// on directory marker blobs.
// https://github.com/fsspec/adlfs/blob/32132c4094350fca2680155a5c236f2e9f991ba5/adlfs/spec.py#L855-L870
Blobs::UploadBlockBlobFromOptions blob_options;
blob_options.Metadata.emplace("is_directory", "true");
blob_options.Metadata.emplace(kFlatNamespaceIsDirectoryMetadataKey, "true");
block_blob_client.UploadFrom(nullptr, 0, blob_options);
}

Expand Down
60 changes: 60 additions & 0 deletions cpp/src/arrow/filesystem/azurefs_test.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -826,6 +826,41 @@ class TestAzureFileSystem : public ::testing::Test {
AssertFileInfo(fs(), subdir3, FileType::Directory);
}

void TestDisallowReadingOrWritingDirectoryMarkers() {
auto data = SetUpPreexistingData();
auto directory_path = data.Path("directory");

ASSERT_OK(fs()->CreateDir(directory_path));
ASSERT_RAISES(IOError, fs()->OpenInputFile(directory_path));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(directory_path));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(directory_path));

auto directory_path_with_slash = directory_path + "/";
ASSERT_RAISES(IOError, fs()->OpenInputFile(directory_path_with_slash));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(directory_path_with_slash));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(directory_path_with_slash));
}

void TestDisallowCreatingFileAndDirectoryWithTheSameName() {
auto data = SetUpPreexistingData();
auto path1 = data.Path("directory1");
ASSERT_OK(fs()->CreateDir(path1));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(path1));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(path1));
AssertFileInfo(fs(), path1, FileType::Directory);

auto path2 = data.Path("directory2");
ASSERT_OK(fs()->OpenOutputStream(path2));
// CreateDir returns OK even if there is already a file or directory at this
// location. Whether or not this is the desired behaviour is debatable.
ASSERT_OK(fs()->CreateDir(path2));
AssertFileInfo(fs(), path2, FileType::File);
}

void TestOpenOutputStreamWithMissingContainer() {
ASSERT_RAISES(IOError, fs()->OpenOutputStream("not-a-container/file", {}));
}

void TestDeleteDirSuccessEmpty() {
if (HasSubmitBatchBug()) {
GTEST_SKIP() << kSubmitBatchBugMessage;
Expand DownExpand Up@@ -1559,6 +1594,19 @@ TYPED_TEST(TestAzureFileSystemOnAllScenarios, CreateDirOnMissingContainer) {
this->TestCreateDirOnMissingContainer();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, DisallowReadingOrWritingDirectoryMarkers) {
this->TestDisallowReadingOrWritingDirectoryMarkers();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios,
DisallowCreatingFileAndDirectoryWithTheSameName) {
this->TestDisallowCreatingFileAndDirectoryWithTheSameName();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, OpenOutputStreamWithMissingContainer) {
this->TestOpenOutputStreamWithMissingContainer();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, DeleteDirSuccessEmpty) {
this->TestDeleteDirSuccessEmpty();
}
Expand DownExpand Up@@ -2162,6 +2210,18 @@ TEST_F(TestAzuriteFileSystem, WriteMetadata) {
.Value.Metadata;
// Defaults are overwritten and not merged.
EXPECT_EQ(Core::CaseInsensitiveMap{std::make_pair("bar", "foo")}, blob_metadata);

// Metadata can be written without writing any data.
ASSERT_OK_AND_ASSIGN(
output, fs_with_defaults->OpenAppendStream(
full_path, /*metadata=*/arrow::key_value_metadata({{"bar", "baz"}})));
ASSERT_OK(output->Close());
blob_metadata = blob_service_client_->GetBlobContainerClient(data.container_name)
.GetBlockBlobClient(blob_path)
.GetProperties()
.Value.Metadata;
// Defaults are overwritten and not merged.
EXPECT_EQ(Core::CaseInsensitiveMap{std::make_pair("bar", "baz")}, blob_metadata);
}

TEST_F(TestAzuriteFileSystem, OpenOutputStreamSmall) {
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
95 changes: 75 additions & 20 deletions cpp/src/arrow/filesystem/azurefs.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -347,6 +347,22 @@ bool IsContainerNotFound(const Storage::StorageException& e) {
return false;
}

const auto kHierarchicalNamespaceIsDirectoryMetadataKey = "hdi_isFolder";
const auto kFlatNamespaceIsDirectoryMetadataKey = "is_directory";

bool MetadataIndicatesIsDirectory(const Storage::Metadata& metadata) {
// Inspired by
// https://github.com/Azure/azure-sdk-for-cpp/blob/12407e8bfcb9bc1aa43b253c1d0ec93bf795ae3b/sdk/storage/azure-storage-files-datalake/src/datalake_utilities.cpp#L86-L91
auto hierarchical_directory_metadata =
metadata.find(kHierarchicalNamespaceIsDirectoryMetadataKey);
if (hierarchical_directory_metadata != metadata.end()) {
return hierarchical_directory_metadata->second == "true";
}
auto flat_directory_metadata = metadata.find(kFlatNamespaceIsDirectoryMetadataKey);
return flat_directory_metadata != metadata.end() &&
flat_directory_metadata->second == "true";
}

template <typename ArrowType>
std::string FormatValue(typename TypeTraits<ArrowType>::CType value) {
struct StringAppender {
Expand DownExpand Up@@ -512,11 +528,18 @@ class ObjectInputFile final : public io::RandomAccessFile {

Status Init() {
if (content_length_ != kNoSize) {
// When the user provides the file size we don't validate that its a file. This is
// only a read so its not a big deal if the user makes a mistake.
DCHECK_GE(content_length_, 0);
return Status::OK();
}
try {
// To open an ObjectInputFile the Blob must exist and it must not represent
// a directory. Additionally we need to know the file size.
auto properties = blob_client_->GetProperties();
if (MetadataIndicatesIsDirectory(properties.Value.Metadata)) {
return NotAFile(location_);
}
content_length_ = properties.Value.BlobSize;
metadata_ = PropertiesToMetadata(properties.Value);
return Status::OK();
Expand DownExpand Up@@ -698,11 +721,10 @@ class ObjectAppendStream final : public io::OutputStream {
ObjectAppendStream(std::shared_ptr<Blobs::BlockBlobClient> block_blob_client,
const io::IOContext& io_context, const AzureLocation& location,
const std::shared_ptr<const KeyValueMetadata>& metadata,
const AzureOptions& options, int64_t size = kNoSize)
const AzureOptions& options)
: block_blob_client_(std::move(block_blob_client)),
io_context_(io_context),
location_(location),
content_length_(size) {
location_(location) {
if (metadata && metadata->size() != 0) {
metadata_ = ArrowMetadataToAzureMetadata(metadata);
} else if (options.default_metadata && options.default_metadata->size() != 0) {
Expand All@@ -716,17 +738,31 @@ class ObjectAppendStream final : public io::OutputStream {
io::internal::CloseFromDestructor(this);
}

Status Init() {
if (content_length_ != kNoSize) {
DCHECK_GE(content_length_, 0);
pos_ = content_length_;
Status Init(const bool truncate,
std::function<Status()> ensure_not_flat_namespace_directory) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

You can inject AzureFileSystem *azure_file_system here and not have to allocate a closure for this. You would call AzureFileSystem::Impl::EnsureNotFlatNamespaceDirectory(location) via azure_file_system->impl_ (accessible because the handles produced by the azure file system can be friends with the filesystem class).

@Tom-NewtonTom-NewtonFeb 20, 2024

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the extra info. I was planning to do this but I was struggling with the friends thing.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The trick is to forward-declare the class as well.

diff --git a/cpp/src/arrow/filesystem/azurefs.h b/cpp/src/arrow/filesystem/azurefs.h
index 2a131e40c..d48ef9dd7 100644
--- a/cpp/src/arrow/filesystem/azurefs.h
+++ b/cpp/src/arrow/filesystem/azurefs.h
@@ -44,6 +44,7 @@ classDataLakeServiceClient;
namespacearrow::fs {
+classObjectAppendStream;
classTestAzureFileSystem;
/// Options for the AzureFileSystem implementation.
@@ -180,6 +181,7 @@ classARROW_EXPORT AzureFileSystem : public FileSystem {
explicitAzureFileSystem(std::unique_ptr<Impl>&& impl);
+ friendclassObjectAppendStream;
friendclassTestAzureFileSystem;
voidForceCachedHierarchicalNamespaceSupport(int hns_support);

@Tom-NewtonTom-NewtonFeb 21, 2024

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think my main problem was that ObjectAppendStream is defined inside an anonymous namespace but I still haven't got it working as you describe.

Are you suggesting to use AzureFileSystem *azure_file_system or AzureFileSystem:Impl *azure_file_system as the argument to ObjectAppendStream::Impl. I don't know how I can get a AzureFileSystem pointer from inside AzureFileSystem::Impl and using AzureFileSystem::Impl as the argument leads to incomplete type errors which I don't think I can avoid.

Also if you wouldn't mind I would be interested to know what the disadvantage of a lambda function is compared to what you proposed.

Sorry about my lacking C++ knowledge here.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Also if you wouldn't mind I would be interested to know what the disadvantage of a lambda function is compared to what you proposed.

To create the std::function, you heap allocate an object with copies of the values in the capture list and generate a lot more extra code in the binary:

class function {
T valuesfromthecpapturelist;
RetType operator()(ArgsType ...) {...};
}

When you think about an std::function this way (a pair of context data and a function), you realize the class you already serves that purpose.

But hey, this is becoming challenging, so I won't hold the PR anymore because of this. Moving to Init() was a big step in the right direction.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for explaining

if (truncate) {
content_length_ = 0;
pos_ = 0;
// We need to create an empty file overwriting any existing file, but
// fail if there is an existing directory.
RETURN_NOT_OK(ensure_not_flat_namespace_directory());
// On hierarchical namespace CreateEmptyBlockBlob will fail if there is an existing
// directory so we don't need to check like we do on flat namespace.
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client_));
} else {
try {
auto properties = block_blob_client_->GetProperties();
if (MetadataIndicatesIsDirectory(properties.Value.Metadata)) {
return NotAFile(location_);
}
content_length_ = properties.Value.BlobSize;
pos_ = content_length_;
} catch (const Storage::StorageException& exception) {
if (exception.StatusCode == Http::HttpStatusCode::NotFound) {
// No file exists but on flat namespace its possible there is a directory
// marker or an implied directory. Ensure there is no directory before starting
// a new empty file.
RETURN_NOT_OK(ensure_not_flat_namespace_directory());
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client_));
} else {
return ExceptionToStatus(
Expand All@@ -743,6 +779,7 @@ class ObjectAppendStream final : public io::OutputStream {
block_ids_.push_back(block.Name);
}
}
initialised_ = true;
return Status::OK();
}

Expand DownExpand Up@@ -789,6 +826,11 @@ class ObjectAppendStream final : public io::OutputStream {

Status Flush() override {
RETURN_NOT_OK(CheckClosed("flush"));
if (!initialised_) {
// If the stream has not been successfully initialized then there is nothing to
// flush. This also avoids some unhandled errors when flushing in the destructor.
return Status::OK();
}
return CommitBlockList(block_blob_client_, block_ids_, metadata_);
}

Expand DownExpand Up@@ -840,10 +882,11 @@ class ObjectAppendStream final : public io::OutputStream {
std::shared_ptr<Blobs::BlockBlobClient> block_blob_client_;
const io::IOContext io_context_;
const AzureLocation location_;
int64_t content_length_ = kNoSize;

bool closed_ = false;
bool initialised_ = false;
int64_t pos_ = 0;
int64_t content_length_ = kNoSize;
std::vector<std::string> block_ids_;
Storage::Metadata metadata_;
};
Expand DownExpand Up@@ -1662,20 +1705,32 @@ class AzureFileSystem::Impl {
AzureFileSystem* fs) {
RETURN_NOT_OK(ValidateFileLocation(location));

const auto blob_container_client = GetBlobContainerClient(location.container);
auto block_blob_client = std::make_shared<Blobs::BlockBlobClient>(
blob_service_client_->GetBlobContainerClient(location.container)
.GetBlockBlobClient(location.path));
blob_container_client.GetBlockBlobClient(location.path));

auto ensure_not_flat_namespace_directory = [this, location,
blob_container_client]() -> Status {
ARROW_ASSIGN_OR_RAISE(
auto hns_support,
HierarchicalNamespaceSupport(GetFileSystemClient(location.container)));
if (hns_support == HNSSupport::kDisabled) {
// Flat namespace so we need to GetFileInfo in-case its a directory.
ARROW_ASSIGN_OR_RAISE(auto status, GetFileInfo(blob_container_client, location))
if (status.type() == FileType::Directory) {
return NotAFile(location);
}
}
// kContainerNotFound - it doesn't exist, so no need to check if its a directory.
// kEnabled - hierarchical namespace so Azure APIs will fail if its a directory. We
// don't need to explicitly check.
return Status::OK();
};

std::shared_ptr<ObjectAppendStream> stream;
if (truncate) {
RETURN_NOT_OK(CreateEmptyBlockBlob(*block_blob_client));
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_, 0);
} else {
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_);
}
RETURN_NOT_OK(stream->Init());
stream = std::make_shared<ObjectAppendStream>(block_blob_client, fs->io_context(),
location, metadata, options_);
RETURN_NOT_OK(stream->Init(truncate, ensure_not_flat_namespace_directory));
return stream;
}

Expand All@@ -1690,7 +1745,7 @@ class AzureFileSystem::Impl {
// on directory marker blobs.
// https://github.com/fsspec/adlfs/blob/32132c4094350fca2680155a5c236f2e9f991ba5/adlfs/spec.py#L855-L870
Blobs::UploadBlockBlobFromOptions blob_options;
blob_options.Metadata.emplace("is_directory", "true");
blob_options.Metadata.emplace(kFlatNamespaceIsDirectoryMetadataKey, "true");
block_blob_client.UploadFrom(nullptr, 0, blob_options);
}

Expand Down
60 changes: 60 additions & 0 deletions cpp/src/arrow/filesystem/azurefs_test.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -826,6 +826,41 @@ class TestAzureFileSystem : public ::testing::Test {
AssertFileInfo(fs(), subdir3, FileType::Directory);
}

void TestDisallowReadingOrWritingDirectoryMarkers() {
auto data = SetUpPreexistingData();
auto directory_path = data.Path("directory");

ASSERT_OK(fs()->CreateDir(directory_path));
ASSERT_RAISES(IOError, fs()->OpenInputFile(directory_path));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(directory_path));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(directory_path));

auto directory_path_with_slash = directory_path + "/";
ASSERT_RAISES(IOError, fs()->OpenInputFile(directory_path_with_slash));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(directory_path_with_slash));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(directory_path_with_slash));
}

void TestDisallowCreatingFileAndDirectoryWithTheSameName() {
auto data = SetUpPreexistingData();
auto path1 = data.Path("directory1");
ASSERT_OK(fs()->CreateDir(path1));
ASSERT_RAISES(IOError, fs()->OpenOutputStream(path1));
ASSERT_RAISES(IOError, fs()->OpenAppendStream(path1));
AssertFileInfo(fs(), path1, FileType::Directory);

auto path2 = data.Path("directory2");
ASSERT_OK(fs()->OpenOutputStream(path2));
// CreateDir returns OK even if there is already a file or directory at this
// location. Whether or not this is the desired behaviour is debatable.
ASSERT_OK(fs()->CreateDir(path2));
AssertFileInfo(fs(), path2, FileType::File);
}

void TestOpenOutputStreamWithMissingContainer() {
ASSERT_RAISES(IOError, fs()->OpenOutputStream("not-a-container/file", {}));
}

void TestDeleteDirSuccessEmpty() {
if (HasSubmitBatchBug()) {
GTEST_SKIP() << kSubmitBatchBugMessage;
Expand DownExpand Up@@ -1559,6 +1594,19 @@ TYPED_TEST(TestAzureFileSystemOnAllScenarios, CreateDirOnMissingContainer) {
this->TestCreateDirOnMissingContainer();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, DisallowReadingOrWritingDirectoryMarkers) {
this->TestDisallowReadingOrWritingDirectoryMarkers();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios,
DisallowCreatingFileAndDirectoryWithTheSameName) {
this->TestDisallowCreatingFileAndDirectoryWithTheSameName();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, OpenOutputStreamWithMissingContainer) {
this->TestOpenOutputStreamWithMissingContainer();
}

TYPED_TEST(TestAzureFileSystemOnAllScenarios, DeleteDirSuccessEmpty) {
this->TestDeleteDirSuccessEmpty();
}
Expand DownExpand Up@@ -2162,6 +2210,18 @@ TEST_F(TestAzuriteFileSystem, WriteMetadata) {
.Value.Metadata;
// Defaults are overwritten and not merged.
EXPECT_EQ(Core::CaseInsensitiveMap{std::make_pair("bar", "foo")}, blob_metadata);

// Metadata can be written without writing any data.
ASSERT_OK_AND_ASSIGN(
output, fs_with_defaults->OpenAppendStream(
full_path, /*metadata=*/arrow::key_value_metadata({{"bar", "baz"}})));
ASSERT_OK(output->Close());
blob_metadata = blob_service_client_->GetBlobContainerClient(data.container_name)
.GetBlockBlobClient(blob_path)
.GetProperties()
.Value.Metadata;
// Defaults are overwritten and not merged.
EXPECT_EQ(Core::CaseInsensitiveMap{std::make_pair("bar", "baz")}, blob_metadata);
}

TEST_F(TestAzuriteFileSystem, OpenOutputStreamSmall) {
Expand Down