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
8 changes: 8 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

- **Asynchronous storage path visitors** — `StoragePathVisitor::visit_path` is now asynchronous so storage drivers can flush bounded path batches without full-list buffering; external driver implementations must update their visitor method to `async`.

### Fixed

- **Storage migration multipart memory bound** — Storage-policy Blob migration now
plans provider part limits separately from the local heap budget, uses bounded
reader uploads with reopenable source ranges for retries, and exposes multipart
capability results during dry-run preflight. Existing hash, verification, abort,
checkpoint, and Blob CAS semantics remain unchanged.

## [v0.6.0] - 2026-09-12

### Release Highlights
Expand Down
15 changes: 8 additions & 7 deletions crates/aster_drive_storage/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -77,11 +77,12 @@ pub use traits::driver::{
};
pub use traits::{
DirectDownloadStorageDriver, ExactSizeReader, ListStorageDriver, LocalPathStorageDriver,
MultipartStorageDriver, NativeMediaMetadataRequest, NativeMediaMetadataResult,
NativeMediaMetadataStorageDriver, NativeThumbnailRequest, NativeThumbnailStorageDriver,
PresignedUploadStorageDriver, ProviderResumableUploadCapabilities,
ProviderResumableUploadDriver, ProviderResumableUploadFragmentOutcome,
ProviderResumableUploadSession, ProviderResumableUploadStatus, StorageCapacityInfo,
StorageCapacityStatus, StorageDriverExtensions, StreamUploadAttempt, StreamUploadCleanup,
StreamUploadDriver, UploadedMultipartPart, checked_upload_size, exact_size_error_kind,
MultipartStorageCapabilities, MultipartStorageDriver, MultipartUploadMode,
NativeMediaMetadataRequest, NativeMediaMetadataResult, NativeMediaMetadataStorageDriver,
NativeThumbnailRequest, NativeThumbnailStorageDriver, PresignedUploadStorageDriver,
ProviderResumableUploadCapabilities, ProviderResumableUploadDriver,
ProviderResumableUploadFragmentOutcome, ProviderResumableUploadSession,
ProviderResumableUploadStatus, StorageCapacityInfo, StorageCapacityStatus,
StorageDriverExtensions, StreamUploadAttempt, StreamUploadCleanup, StreamUploadDriver,
UploadedMultipartPart, checked_upload_size, exact_size_error_kind,
};
6 changes: 6 additions & 0 deletions crates/aster_drive_storage/src/traits/driver.rs
Original file line number Diff line number Diff line change
Expand Up @@ -196,6 +196,12 @@ pub trait StorageDriver: Send + Sync {
false
}

/// Optional maximum size for a single non-multipart PUT accepted by the
/// provider. This is distinct from multipart part/object limits.
fn max_single_put_size(&self) -> Option<u64> {
None
}

/// 删除文件
async fn delete(&self, path: &str) -> Result<()>;

Expand Down
5 changes: 4 additions & 1 deletion crates/aster_drive_storage/src/traits/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,4 +18,7 @@ pub use extensions::{
StorageDriverExtensions, StreamUploadAttempt, StreamUploadCleanup, StreamUploadDriver,
checked_upload_size, exact_size_error_kind,
};
pub use multipart::{MultipartStorageDriver, UploadedMultipartPart};
pub use multipart::{
MultipartStorageCapabilities, MultipartStorageDriver, MultipartUploadMode,
UploadedMultipartPart,
};
59 changes: 59 additions & 0 deletions crates/aster_drive_storage/src/traits/multipart.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,41 @@ use tokio::io::{AsyncRead, AsyncReadExt};
const DEFAULT_MULTIPART_READER_BUFFER_SIZE: usize = 64 * 1024;
const MAX_DEFAULT_MULTIPART_READER_SIZE: usize = 64 * 1024 * 1024;

/// 描述 provider 对 numbered multipart/block upload 的限制以及 reader 上传能力。
///
/// `max_part_size` 是 provider 的协议限制,`upload_mode` 描述驱动在本地如何
/// 消费 part。两者必须分开:一个 1 GiB 的 provider part 可以通过固定小 buffer
/// 流式发送,并不意味着进程需要分配 1 GiB 的 `Vec`。
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct MultipartStorageCapabilities {
pub min_part_size: u64,
pub max_part_size: Option<u64>,
pub max_parts: u64,
pub max_object_size: Option<u64>,
pub upload_mode: MultipartUploadMode,
}

#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MultipartUploadMode {
NativeStreaming,
Buffered { max_size: u64 },
}

impl MultipartStorageCapabilities {
pub const fn conservative_default() -> Self {
Self {
min_part_size: 5 * 1024 * 1024,
max_part_size: None,
max_parts: 10_000,
max_object_size: None,
upload_mode: MultipartUploadMode::Buffered {
max_size: (MAX_DEFAULT_MULTIPART_READER_SIZE - DEFAULT_MULTIPART_READER_BUFFER_SIZE)
as u64,
},
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
}
}

/// Provider 端已经接收的 multipart part 明细。
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct UploadedMultipartPart {
Expand All @@ -38,6 +73,14 @@ pub struct UploadedMultipartPart {
/// 协议勉强模拟。
#[async_trait]
pub trait MultipartStorageDriver: Send + Sync {
/// 返回 provider part 限制与 reader 上传的本地内存语义。
///
/// 未覆写的驱动使用保守 buffered fallback,迁移 preflight 会在 logical
/// part 超过该 fallback 上限时提前阻塞任务,而不是让 worker 在运行中才失败。
fn capabilities(&self) -> MultipartStorageCapabilities {
MultipartStorageCapabilities::conservative_default()
}

/// 创建 multipart upload,返回 provider 端的 upload_id
async fn create_multipart_upload(&self, path: &str) -> Result<String>;

Expand Down Expand Up @@ -421,4 +464,20 @@ mod tests {

assert_eq!(error.kind(), StorageErrorKind::Unsupported);
}

#[test]
fn default_capabilities_describe_the_bounded_buffered_fallback() {
let driver = CapturingMultipartDriver::new();
let capabilities = driver.capabilities();
assert_eq!(capabilities.min_part_size, 5 * 1024 * 1024);
assert_eq!(capabilities.max_parts, 10_000);
assert_eq!(capabilities.max_object_size, None);
assert_eq!(
capabilities.upload_mode,
MultipartUploadMode::Buffered {
max_size: (MAX_DEFAULT_MULTIPART_READER_SIZE - DEFAULT_MULTIPART_READER_BUFFER_SIZE)
as u64
}
);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,192 @@
import { render, screen } from "@testing-library/react";
import type { ReactNode } from "react";
import { describe, expect, it, vi } from "vitest";
import { StoragePolicyMigrationDialog } from "@/components/admin/admin-policies-page/StoragePolicyMigrationDialog";
import type { StoragePolicy, StoragePolicyMigrationDryRun } from "@/types/api";

vi.mock("react-i18next", () => ({
useTranslation: () => ({
t: (key: string, values?: Record<string, unknown>) =>
values ? `${key} ${Object.values(values).join(" ")}` : key,
}),
}));
vi.mock("@/components/ui/dialog", () => ({
Dialog: ({ children, open }: { children: ReactNode; open: boolean }) =>
open ? <div>{children}</div> : null,
DialogContent: ({ children }: { children: ReactNode }) => (
<div>{children}</div>
),
DialogDescription: ({ children }: { children: ReactNode }) => (
<p>{children}</p>
),
DialogFooter: ({ children }: { children: ReactNode }) => (
<footer>{children}</footer>
),
DialogHeader: ({ children }: { children: ReactNode }) => (
<header>{children}</header>
),
DialogTitle: ({ children }: { children: ReactNode }) => <h2>{children}</h2>,
}));
vi.mock("@/components/ui/button", () => ({
Button: ({
children,
...props
}: {
children?: ReactNode;
[key: string]: unknown;
}) => <button {...props}>{children}</button>,
}));
vi.mock("@/components/ui/icon", () => ({
Icon: ({ name }: { name: string }) => <i>{name}</i>,
}));
vi.mock("@/components/ui/label", () => ({
Label: ({ children }: { children: ReactNode }) => <span>{children}</span>,
}));
vi.mock("@/components/ui/select", () => ({
Select: ({ children }: { children: ReactNode }) => <div>{children}</div>,
SelectContent: ({ children }: { children: ReactNode }) => (
<div>{children}</div>
),
SelectItem: ({ children }: { children: ReactNode }) => (
<span>{children}</span>
),
SelectTrigger: ({ children }: { children: ReactNode }) => (
<div>{children}</div>
),
SelectValue: ({ children }: { children?: ReactNode }) => (
<span>{children}</span>
),
}));

const policies = [
{ id: 1, name: "Source" },
{ id: 2, name: "Target" },
] as StoragePolicy[];

function dryRun(
reason: "provider_limits" | "buffered_heap_budget" | null,
uploadMode: "native_streaming" | "buffered" = "native_streaming",
) {
return {
can_start: reason === null,
content_sha256_blob_count: 0,
estimated_copy_blob_count: 1,
opaque_blob_count: 0,
opaque_key_conflict_count: 0,
source_blob_count: 1,
source_policy_id: 1,
source_total_bytes: 1024,
target_capacity: null,
target_capacity_check: "sufficient",
target_connection_ok: true,
target_matching_blob_count: 0,
target_policy_id: 2,
target_supports_stream_upload: true,
warnings: [],
multipart_plan: {
blob_size: 1024,
can_start: reason === null,
heap_budget: 64 * 1024 * 1024,
part_count: 1,
part_size: 1024,
provider_max_parts: 10_000,
provider_max_part_size: null,
reason,
upload_mode: uploadMode,
},
} as StoragePolicyMigrationDryRun;
}

function renderDialog(dryRunValue: StoragePolicyMigrationDryRun) {
return render(
<StoragePolicyMigrationDialog
dryRun={dryRunValue}
dryRunLoading={false}
open
policies={policies}
sourcePolicyId="1"
submitting={false}
targetPolicyId="2"
onOpenChange={vi.fn()}
onDryRun={vi.fn()}
onSourcePolicyChange={vi.fn()}
onSubmit={vi.fn()}
onTargetPolicyChange={vi.fn()}
/>,
);
}

function renderWithoutDryRun() {
return render(
<StoragePolicyMigrationDialog
dryRun={null}
dryRunLoading={false}
open
policies={policies}
sourcePolicyId="1"
submitting={false}
targetPolicyId="2"
onOpenChange={vi.fn()}
onDryRun={vi.fn()}
onSourcePolicyChange={vi.fn()}
onSubmit={vi.fn()}
onTargetPolicyChange={vi.fn()}
/>,
);
}

describe("StoragePolicyMigrationDialog", () => {
it("renders native streaming multipart plan details", () => {
renderDialog(dryRun(null));
expect(
screen.getByText("policy_migration_multipart_plan"),
).toBeInTheDocument();
expect(
screen.getByText("policy_migration_multipart_mode_native_streaming"),
).toBeInTheDocument();
});

it("renders the translated multipart block reason", () => {
renderDialog(dryRun("buffered_heap_budget"));
expect(
screen.getByText(
"policy_migration_multipart_reason_buffered_heap_budget",
),
).toBeInTheDocument();
});

it("renders the buffered multipart mode", () => {
renderDialog(dryRun("provider_limits", "buffered"));
expect(
screen.getByText("policy_migration_multipart_mode_buffered"),
).toBeInTheDocument();
});

it("renders a verified target capacity detail", () => {
const value = dryRun(null);
value.target_capacity = {
status: "supported",
total_bytes: 100,
available_bytes: 40,
used_bytes: 60,
source: "test",
observed_at: "2026-01-01T00:00:00Z",
};
renderDialog(value);
expect(
screen.getByText(/policy_migration_capacity_available_of_total.*40.*100/),
).toBeInTheDocument();
Comment thread
coderabbitai[bot] marked this conversation as resolved.
});

it("renders the dialog before a dry run exists", () => {
renderWithoutDryRun();
expect(screen.getByText("policy_migration_dry_run")).toBeInTheDocument();
});

it("renders a dry run without a multipart plan", () => {
const value = dryRun(null);
value.multipart_plan = null;
renderDialog(value);
expect(screen.queryByText("policy_migration_multipart_plan")).toBeNull();
});
});
Original file line number Diff line number Diff line change
Expand Up @@ -82,8 +82,8 @@ export function StoragePolicyMigrationDialog({
sourceId !== targetId &&
!dryRunLoading &&
!submitting;
const targetAvailableBytes = dryRun?.target_capacity.available_bytes;
const targetTotalBytes = dryRun?.target_capacity.total_bytes;
const targetAvailableBytes = dryRun?.target_capacity?.available_bytes;
const targetTotalBytes = dryRun?.target_capacity?.total_bytes;
const targetCapacityDetail =
typeof targetAvailableBytes === "number" &&
typeof targetTotalBytes === "number"
Expand Down Expand Up @@ -277,6 +277,33 @@ export function StoragePolicyMigrationDialog({
})}
</div>
</div>
{dryRun.multipart_plan ? (
<div className="rounded-md border bg-background/70 px-2.5 py-2 text-xs text-muted-foreground">
<div className="font-medium text-foreground">
{t("policy_migration_multipart_plan")}
</div>
<div className="mt-1 flex flex-wrap gap-x-3 gap-y-1">
<span>
{t("policy_migration_multipart_part_size", {
size: formatBytes(dryRun.multipart_plan.part_size),
count: dryRun.multipart_plan.part_count,
})}
</span>
<span>
{dryRun.multipart_plan.upload_mode === "native_streaming"
? t("policy_migration_multipart_mode_native_streaming")
: t("policy_migration_multipart_mode_buffered")}
</span>
</div>
{dryRun.multipart_plan.reason ? (
<div className="mt-1 text-destructive">
{t(
`policy_migration_multipart_reason_${dryRun.multipart_plan.reason}`,
)}
</div>
Comment thread
coderabbitai[bot] marked this conversation as resolved.
) : null}
</div>
) : null}
{dryRun.warnings.length > 0 ? (
<div className="space-y-1 rounded-md border border-amber-200 bg-amber-50 px-2.5 py-2 text-xs text-amber-800 dark:border-amber-900 dark:bg-amber-950/40 dark:text-amber-200">
{dryRun.warnings.map((warning) => (
Expand Down
8 changes: 8 additions & 0 deletions frontend-panel/src/i18n/locales/en/admin/policies.json
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,14 @@
"policy_migration_capacity_unavailable": "Unavailable",
"policy_migration_capacity_available_of_total": "({{available}} available / {{total}} total)",
"policy_migration_warning_target_capacity_unavailable": "Target free capacity cannot be verified by this storage driver. Confirm the target has enough capacity before creating the migration task.",
"policy_migration_warning_multipart_capability_unavailable": "The target multipart capability does not cover the largest Blob part plan, so migration is blocked.",
"policy_migration_multipart_plan": "Multipart Plan",
"policy_migration_multipart_part_size": "Logical part {{size}}, {{count}} parts",
"policy_migration_multipart_mode_native_streaming": "Native streaming",
"policy_migration_multipart_mode_buffered": "Bounded buffering",
"policy_migration_multipart_reason_provider_limits": "The target part size or count limits cannot represent this Blob.",
"policy_migration_multipart_reason_buffered_heap_budget": "The target driver would buffer a full part beyond the migration heap budget.",
"policy_migration_multipart_reason_provider_object_size": "The Blob exceeds the target storage object size limit.",
"policy_capacity_title": "AsterDrive Usage",
"policy_capacity_loading": "Checking storage capacity...",
"policy_capacity_checking": "Checking",
Expand Down
Loading
Loading