mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-20 03:22:18 +00:00
Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 4d8262e58b |
Generated
+93
-124
@@ -1198,9 +1198,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "aws-smithy-http-client"
|
||||
version = "1.4.0"
|
||||
version = "1.3.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "ebfd138fac0337cee7516c352757ea73b9f2266e57d0bcb5bc70e9547e45aef1"
|
||||
checksum = "3c1c8a04cb31ba74d0115af5a890bb8c0d48fba64b52812fa13929a6ef0cc83c"
|
||||
dependencies = [
|
||||
"aws-smithy-async",
|
||||
"aws-smithy-protocol-test",
|
||||
@@ -1280,9 +1280,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "aws-smithy-runtime"
|
||||
version = "1.14.0"
|
||||
version = "1.13.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "b82e438d30e02a825d363bd639a9efaed68a8089d86101054b0081e7e0d3e606"
|
||||
checksum = "483b858ff67522011c4786310c5cd8fd88d0be7ea3d5f1a48328446300c4269e"
|
||||
dependencies = [
|
||||
"aws-smithy-async",
|
||||
"aws-smithy-http",
|
||||
@@ -1306,9 +1306,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "aws-smithy-runtime-api"
|
||||
version = "1.15.0"
|
||||
version = "1.14.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "954c563ce84507722d2679f07a35d21b9c6466b3872d513020d0281fc8112ac9"
|
||||
checksum = "3b98f2e1fd67ec06618f9c291e5e495a468e60519e44c9c1979cd0521f3affdb"
|
||||
dependencies = [
|
||||
"aws-smithy-async",
|
||||
"aws-smithy-runtime-api-macros",
|
||||
@@ -1598,7 +1598,7 @@ version = "0.10.4"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "3078c7629b62d3f0439517fa394996acacc5cbc91c5a20d8c658e77abd503a71"
|
||||
dependencies = [
|
||||
"generic-array 0.14.9",
|
||||
"generic-array 0.14.7",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
@@ -1617,7 +1617,7 @@ version = "0.3.3"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "a8894febbff9f758034a5b8e12d87918f56dfc64a8e1fe757d65e29041538d93"
|
||||
dependencies = [
|
||||
"generic-array 0.14.9",
|
||||
"generic-array 0.14.7",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
@@ -1968,7 +1968,7 @@ version = "0.4.4"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "773f3b9af64447d2ce9850330c473515014aa235e6a783b02db81ff39e4a3dad"
|
||||
dependencies = [
|
||||
"crypto-common 0.1.6",
|
||||
"crypto-common 0.1.7",
|
||||
"inout 0.1.4",
|
||||
]
|
||||
|
||||
@@ -2428,7 +2428,7 @@ version = "0.5.5"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "0dc92fb57ca44df6db8059111ab3af99a63d5d0f8375d9972e319a379c6bab76"
|
||||
dependencies = [
|
||||
"generic-array 0.14.9",
|
||||
"generic-array 0.14.7",
|
||||
"rand_core 0.6.4",
|
||||
"subtle",
|
||||
"zeroize",
|
||||
@@ -2453,11 +2453,11 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "crypto-common"
|
||||
version = "0.1.6"
|
||||
version = "0.1.7"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "1bfb12502f3fc46cca1bb51ac28df9d618d813cdc3d2f25b9fe775a34af26bb3"
|
||||
checksum = "78c8292055d1c1df0cce5d180393dc8cce0abec0a7102adb6c7b1eef6016d60a"
|
||||
dependencies = [
|
||||
"generic-array 0.14.9",
|
||||
"generic-array 0.14.7",
|
||||
"typenum",
|
||||
]
|
||||
|
||||
@@ -2703,9 +2703,8 @@ checksum = "4583a4551df46e2792f82ceeac45e850d2e2d5debba0b91f102385cda5b11f06"
|
||||
|
||||
[[package]]
|
||||
name = "datafusion"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "96f76f0167ed0842b29a3d1e41be3c034c0a46409a3a703cc4cc84ee8c24abf4"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-schema",
|
||||
@@ -2752,9 +2751,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-catalog"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "d79ec3460f6ed5c58f9b3f2d873fbc77748b82653bff1b4cdaf06de33bb4e05f"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"async-trait",
|
||||
@@ -2777,9 +2775,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-catalog-listing"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "b48cef241e2efcfd496fe05ae4d0d5de20793451862faefe406c397a467e12d4"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"async-trait",
|
||||
@@ -2801,9 +2798,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-common"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "3f72810485975c258f1b4d00baab31728470676c60c5546f366ebd0d99f05ab6"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-ipc",
|
||||
@@ -2828,9 +2824,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-common-runtime"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "533c28e75dba52f41bde187d23a1cb24ab91c7c097966824fa471e67b60320ea"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"futures",
|
||||
"log",
|
||||
@@ -2839,9 +2834,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-datasource"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "5b00a1fa0da26f6087136a82fea7f13c76a672cbab452d4086952a7cf770a19b"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"async-trait",
|
||||
@@ -2869,9 +2863,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-datasource-arrow"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "5ad17ec881bff2ed7768b4bfe971d3efbf3473f2fd1f9d365447bccbdf908678"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-ipc",
|
||||
@@ -2893,9 +2886,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-datasource-csv"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "b5345285b0c3eaab412e7539b706973c083bd7e5bce575de5e0a3da488d08d1d"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"async-trait",
|
||||
@@ -2916,9 +2908,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-datasource-json"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "da02fb9324f56bd8c53f1ee2e949547425cb66f76adc6832b10d44f80a1221d2"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"async-trait",
|
||||
@@ -2939,9 +2930,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-datasource-parquet"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "3c0b0dc1453952952fd5c69ad1c7f6042176e69ed233011d47e07cf74ed0949e"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-schema",
|
||||
@@ -2971,15 +2961,13 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-doc"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "a88fd985bc0550c36f557db69543cc9d6393b1509783520b30e902f23c555da6"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-execution"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "a98f1052f91b4991f0bf2ce1e4e36dfbdcda454a956b8c8d562c7c845e8fce1d"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-buffer",
|
||||
@@ -3003,9 +2991,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-expr"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "464625a1f0e4b9df552d894fafcc8aac953ebbc8b0fa0acdaf20975fd615040e"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-schema",
|
||||
@@ -3026,9 +3013,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-expr-common"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "2604994999d5aeca1d1df645ffc98bc787447aaff05dde27aad0342b48fc1fe0"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"datafusion-common",
|
||||
@@ -3038,9 +3024,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-functions"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "051e97533e6af53e4aa0a0667cadc886abcaf36c4a5925019c55c0aa4c218fde"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-buffer",
|
||||
@@ -3066,9 +3051,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-functions-aggregate"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "2d0f1bb166d3572b6ed40e1afb2faaacade962abc08c2fcf04babee74681c56b"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"datafusion-common",
|
||||
@@ -3087,9 +3071,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-functions-aggregate-common"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "7ed756770f5f98369e181d692fd5ee6b1127ffd7322caba92f3730f9f5c92333"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"datafusion-common",
|
||||
@@ -3099,9 +3082,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-functions-nested"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "91173fdb5c0ff2a41169a8ffa1b385b8844f18728747bb0a37e35ad7d5772a4f"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-ord",
|
||||
@@ -3124,9 +3106,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-functions-table"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "b1bcdfb286a745461b126719c32700777e83df4f17cc44db5d71ebce5731e840"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"async-trait",
|
||||
@@ -3140,9 +3121,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-functions-window"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "9ec4b508f1f93f00038ba3e737e894ec6c775528b4369413386655ae6125f0fc"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"datafusion-common",
|
||||
@@ -3157,9 +3137,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-functions-window-common"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "0b352020834140073fbf5b46ee0ceb926e5074a9d0bcae1dbd91d0586d999cde"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"datafusion-common",
|
||||
"datafusion-physical-expr-common",
|
||||
@@ -3167,9 +3146,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-macros"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "15192effab05d38cce10e92a6fb48c967b5f166b27b7195a165a72b232569c58"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"datafusion-doc",
|
||||
"quote",
|
||||
@@ -3178,9 +3156,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-optimizer"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "854445d9f7847e1e46089cf61b8d341a64382f14484e912c83a0f23b31216896"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"chrono",
|
||||
@@ -3198,9 +3175,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-physical-expr"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "671558dad1d2aa253c39c0a4c52515958b99eb91abf649f4b88d5e69cc55282f"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"datafusion-common",
|
||||
@@ -3220,9 +3196,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-physical-expr-adapter"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "ffae3d78c2da80ecc829cb58536cc5aca2e99cf1365eda694fc75bfe288861e0"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"datafusion-common",
|
||||
@@ -3235,9 +3210,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-physical-expr-common"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "3d9092ed15e7203fbd0903215172f7c9d18f10d94cba35137f3b3836f7c46f16"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"chrono",
|
||||
@@ -3252,9 +3226,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-physical-optimizer"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "9005b6cf50b57b72d476c6ed4662b04be7ca6be5320ba9127c6d0b7e4218095b"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"datafusion-common",
|
||||
@@ -3272,9 +3245,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-physical-plan"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "5787e4fcff4adc4fce8948441103a99705018b49c8dff0720b650bd7a15da112"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-data",
|
||||
@@ -3307,9 +3279,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-pruning"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "9e651c8df0b90daed6a7be5921ec0ee379e6909705f063eeff70fd4e35010e4c"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"datafusion-common",
|
||||
@@ -3323,9 +3294,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-session"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "fb56667ee38217efab19b895d9a936052cfb47ed438a19663351bdc42a6214a1"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow-schema",
|
||||
"async-trait",
|
||||
@@ -3338,9 +3308,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "datafusion-sql"
|
||||
version = "55.0.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "9c29067cb9d32f8e603c45e15d61ea18f1069f96ceafeceb4e18466b8e5b31d9"
|
||||
version = "54.1.0"
|
||||
source = "git+https://github.com/apache/datafusion.git?rev=e08aed1e5de41dcf81d529140dae07723b942a5e#e08aed1e5de41dcf81d529140dae07723b942a5e"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"bigdecimal",
|
||||
@@ -3695,7 +3664,7 @@ checksum = "9ed9a281f7bc9b7576e61468ba615a66a5c8cfdff42420a70aa82701a3b1e292"
|
||||
dependencies = [
|
||||
"block-buffer 0.10.4",
|
||||
"const-oid 0.9.6",
|
||||
"crypto-common 0.1.6",
|
||||
"crypto-common 0.1.7",
|
||||
"subtle",
|
||||
]
|
||||
|
||||
@@ -3955,7 +3924,7 @@ dependencies = [
|
||||
"crypto-bigint 0.5.5",
|
||||
"digest 0.10.7",
|
||||
"ff 0.13.1",
|
||||
"generic-array 0.14.9",
|
||||
"generic-array 0.14.7",
|
||||
"group 0.13.0",
|
||||
"hkdf 0.12.4",
|
||||
"pem-rfc7468 0.7.0",
|
||||
@@ -4400,9 +4369,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "generic-array"
|
||||
version = "0.14.9"
|
||||
version = "0.14.7"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "4bb6743198531e02858aeaea5398fcc883e71851fcbcb5a2f773e2fb6cb1edf2"
|
||||
checksum = "85649ca51fd72272d7821adaf274ad91c288277713d9c18820d8499a7ff69e9a"
|
||||
dependencies = [
|
||||
"typenum",
|
||||
"version_check",
|
||||
@@ -4415,7 +4384,7 @@ version = "1.4.5"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "337d46834ee672ab3e48caca2cb0c78cc174fb12b3a68d0d88f99a0519a5e36e"
|
||||
dependencies = [
|
||||
"generic-array 0.14.9",
|
||||
"generic-array 0.14.7",
|
||||
"rustversion",
|
||||
"typenum",
|
||||
]
|
||||
@@ -4757,9 +4726,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "h2"
|
||||
version = "0.4.17"
|
||||
version = "0.4.16"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "9f877e75f39e9827ec50a572dd592684ac28c029578726c85f1b2aa6ab807449"
|
||||
checksum = "a9f37a958b41b3b19ee2707c06439c0e9e547e847223eb791ecb0cb821c65e27"
|
||||
dependencies = [
|
||||
"atomic-waker",
|
||||
"bytes",
|
||||
@@ -5449,7 +5418,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "879f10e63c20629ecabbb64a8010319738c66a5cd0c29b02d63d272b03751d01"
|
||||
dependencies = [
|
||||
"block-padding 0.3.3",
|
||||
"generic-array 0.14.9",
|
||||
"generic-array 0.14.7",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
@@ -8659,18 +8628,18 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "ref-cast"
|
||||
version = "1.0.27"
|
||||
version = "1.0.26"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "7e440fb4e4b4147295338efb76001ab9e4efc0e5839df2c47fc5ac2381d365c3"
|
||||
checksum = "216e8f773d7923bcba9ceb86a86c93cabb3903a11872fc3f138c49630e50b96d"
|
||||
dependencies = [
|
||||
"ref-cast-impl",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "ref-cast-impl"
|
||||
version = "1.0.27"
|
||||
version = "1.0.26"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "92ecd8964f8453721699a1ed72037b0db49ce2f5a5138486ee89bed6f67cdf3a"
|
||||
checksum = "2c9283685feec7d69af75fb0e858d5e7378f33fe4fc699383b2916ab9273e03c"
|
||||
dependencies = [
|
||||
"proc-macro2",
|
||||
"quote",
|
||||
@@ -10864,7 +10833,7 @@ checksum = "d3e97a565f76233a6003f9f5c54be1d9c5bdfa3eccfb189469f11ec4901c47dc"
|
||||
dependencies = [
|
||||
"base16ct 0.2.0",
|
||||
"der 0.7.10",
|
||||
"generic-array 0.14.9",
|
||||
"generic-array 0.14.7",
|
||||
"pkcs8 0.10.2",
|
||||
"subtle",
|
||||
"zeroize",
|
||||
@@ -11649,9 +11618,9 @@ checksum = "13c2bddecc57b384dee18652358fb23172facb8a2c51ccc10d74c157bdea3292"
|
||||
|
||||
[[package]]
|
||||
name = "suppaftp"
|
||||
version = "10.0.2"
|
||||
version = "10.0.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "821001051ea3d12a60fb790b8c7cb9a6f5f8698dcfdca4cd533a025fefb0b5b8"
|
||||
checksum = "9c890e698eaf58526b6e7105d74c5d91ebe76a4da1faac2e20ff10e8e5c8bcff"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"chrono",
|
||||
@@ -11842,7 +11811,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd"
|
||||
dependencies = [
|
||||
"fastrand",
|
||||
"getrandom 0.4.3",
|
||||
"getrandom 0.3.4",
|
||||
"once_cell",
|
||||
"rustix",
|
||||
"windows-sys 0.61.2",
|
||||
@@ -13392,9 +13361,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "zerovec"
|
||||
version = "0.11.8"
|
||||
version = "0.11.7"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "bb0464e17806c1d976d5cba29399c7f08e516e279e2ba493f63123b5fca67dd8"
|
||||
checksum = "94b5c6b5976d66c1d703c4fd17d3f5e43c8cedaacf604961b171adc7130896d8"
|
||||
dependencies = [
|
||||
"yoke",
|
||||
"zerofrom",
|
||||
|
||||
+5
-4
@@ -231,8 +231,8 @@ aws-credential-types = { version = "1.3.0" }
|
||||
aws-sdk-kms = { default-features = false, version = "1.115.0" }
|
||||
aws-sdk-s3 = { default-features = false, version = "1.142.0" }
|
||||
aws-sdk-sts = { default-features = false, version = "1.111.0" }
|
||||
aws-smithy-http-client = { default-features = false, version = "1.4.0" }
|
||||
aws-smithy-runtime-api = { version = "1.15.0" }
|
||||
aws-smithy-http-client = { default-features = false, version = "1.3.0" }
|
||||
aws-smithy-runtime-api = { version = "1.14.0" }
|
||||
aws-smithy-types = { version = "1.6.2" }
|
||||
base64 = "0.23.1"
|
||||
base64-simd = "0.8.0"
|
||||
@@ -245,7 +245,8 @@ crossbeam-queue = "0.3.13"
|
||||
crossbeam-channel = "0.5.16"
|
||||
crossbeam-deque = "0.8.7"
|
||||
crossbeam-utils = "0.8.22"
|
||||
datafusion = { default-features = false, version = "55.0.0" }
|
||||
datafusion = { default-features = false, git = "https://github.com/apache/datafusion.git", rev = "e08aed1e5de41dcf81d529140dae07723b942a5e" }
|
||||
#datafusion = { default-features = false, version = "54.1.0" }
|
||||
derive_builder = "0.20.2"
|
||||
enumset = "1.1.14"
|
||||
faster-hex = "0.10.0"
|
||||
@@ -340,7 +341,7 @@ pyroscope = { version = "2.1.1" }
|
||||
# FTP and SFTP
|
||||
libunftp = { version = "0.23.0" }
|
||||
unftp-core = "0.1.0"
|
||||
suppaftp = { version = "10.0.2" }
|
||||
suppaftp = { version = "10.0.1" }
|
||||
rcgen = { version = "0.14.9", default-features = false, features = ["aws_lc_rs", "crypto", "pem"] }
|
||||
russh = { version = "0.62.7" }
|
||||
russh-sftp = "2.4.0"
|
||||
|
||||
+14
-3231
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,351 @@
|
||||
// Copyright 2024 RustFS Team
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
use crate::{Error, Result};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::collections::HashSet;
|
||||
use std::path::Path;
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
use tokio::sync::RwLock;
|
||||
use tracing::{debug, warn};
|
||||
|
||||
use super::super::{BUCKET_META_PREFIX, DiskStore, HealDiskExt as _, RUSTFS_META_BUCKET};
|
||||
use super::{
|
||||
LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, PersistThrottle, RESUME_CHECKPOINT_FILE, delete_resume_file, path_to_str,
|
||||
validate_resume_task_id,
|
||||
};
|
||||
|
||||
const EVENT_HEAL_CHECKPOINT_STATE: &str = "heal_checkpoint_state";
|
||||
|
||||
/// Current on-disk schema version for `ResumeCheckpoint`. Same rationale as
|
||||
/// `CURRENT_RESUME_SCHEMA`: pre-per-version dedup identities are not comparable
|
||||
/// to the new `compose_key` identities, so a stale checkpoint is discarded.
|
||||
pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 5;
|
||||
|
||||
/// resume checkpoint
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct ResumeCheckpoint {
|
||||
/// on-disk schema version; absent in legacy snapshots (defaults to 0)
|
||||
#[serde(default)]
|
||||
pub schema_version: u32,
|
||||
/// task id
|
||||
pub task_id: String,
|
||||
/// checkpoint time
|
||||
pub checkpoint_time: u64,
|
||||
/// current bucket index
|
||||
pub current_bucket_index: usize,
|
||||
/// current object index
|
||||
pub current_object_index: usize,
|
||||
/// Objects healed since the last completed page. HashSet: with the
|
||||
/// previous Vec the per-object `contains` was O(n) and made large-bucket
|
||||
/// heals O(N²). Only spans the in-flight page — completed pages are
|
||||
/// covered by `current_object_index`, so `complete_page` prunes the sets.
|
||||
pub processed_objects: HashSet<String>,
|
||||
/// failed objects
|
||||
pub failed_objects: HashSet<String>,
|
||||
/// skipped objects
|
||||
pub skipped_objects: HashSet<String>,
|
||||
}
|
||||
|
||||
impl ResumeCheckpoint {
|
||||
pub fn new(task_id: String) -> Self {
|
||||
Self {
|
||||
schema_version: CURRENT_CHECKPOINT_SCHEMA,
|
||||
task_id,
|
||||
checkpoint_time: SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs(),
|
||||
current_bucket_index: 0,
|
||||
current_object_index: 0,
|
||||
processed_objects: HashSet::new(),
|
||||
failed_objects: HashSet::new(),
|
||||
skipped_objects: HashSet::new(),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn update_position(&mut self, bucket_index: usize, object_index: usize) {
|
||||
self.current_bucket_index = bucket_index;
|
||||
self.current_object_index = object_index;
|
||||
self.checkpoint_time = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
|
||||
}
|
||||
|
||||
pub fn add_processed_object(&mut self, object: String) {
|
||||
self.processed_objects.insert(object);
|
||||
}
|
||||
|
||||
pub fn add_failed_object(&mut self, object: String) {
|
||||
self.failed_objects.insert(object);
|
||||
}
|
||||
|
||||
pub fn add_skipped_object(&mut self, object: String) {
|
||||
self.skipped_objects.insert(object);
|
||||
}
|
||||
|
||||
/// Advance past a fully-processed page: objects below `object_index` are
|
||||
/// skipped by position on resume, so the per-object sets no longer need
|
||||
/// their entries and would otherwise grow with the whole bucket.
|
||||
pub fn complete_page(&mut self, bucket_index: usize, object_index: usize) {
|
||||
self.update_position(bucket_index, object_index);
|
||||
self.processed_objects.clear();
|
||||
self.skipped_objects.clear();
|
||||
self.failed_objects.clear();
|
||||
}
|
||||
|
||||
/// Reset the scan to the start and clear the per-object sets so a retry
|
||||
/// re-scans the whole set.
|
||||
pub fn reset_for_retry(&mut self) {
|
||||
self.update_position(0, 0);
|
||||
self.processed_objects.clear();
|
||||
self.skipped_objects.clear();
|
||||
self.failed_objects.clear();
|
||||
}
|
||||
}
|
||||
|
||||
/// resume checkpoint manager
|
||||
pub struct CheckpointManager {
|
||||
disk: DiskStore,
|
||||
checkpoint: Arc<RwLock<ResumeCheckpoint>>,
|
||||
throttle: Mutex<PersistThrottle>,
|
||||
}
|
||||
|
||||
impl CheckpointManager {
|
||||
/// create new checkpoint manager
|
||||
pub async fn new(disk: DiskStore, task_id: String) -> Result<Self> {
|
||||
validate_resume_task_id(&task_id)?;
|
||||
let checkpoint = ResumeCheckpoint::new(task_id);
|
||||
let manager = Self {
|
||||
disk,
|
||||
checkpoint: Arc::new(RwLock::new(checkpoint)),
|
||||
throttle: Mutex::new(PersistThrottle::new()),
|
||||
};
|
||||
|
||||
// save initial checkpoint
|
||||
if let Err(e) = manager.save_checkpoint().await {
|
||||
warn!(
|
||||
target: "rustfs::heal::resume",
|
||||
event = EVENT_HEAL_CHECKPOINT_STATE,
|
||||
component = LOG_COMPONENT_HEAL,
|
||||
subsystem = LOG_SUBSYSTEM_RESUME,
|
||||
state = "initial_save_failed",
|
||||
error = %e,
|
||||
"Heal checkpoint persistence failed"
|
||||
);
|
||||
}
|
||||
Ok(manager)
|
||||
}
|
||||
|
||||
/// load checkpoint from disk
|
||||
pub async fn load_from_disk(disk: DiskStore, task_id: &str) -> Result<Self> {
|
||||
validate_resume_task_id(task_id)?;
|
||||
let checkpoint_data = Self::read_checkpoint_file(&disk, task_id).await?;
|
||||
let mut checkpoint: ResumeCheckpoint =
|
||||
serde_json::from_slice(&checkpoint_data).map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to deserialize checkpoint: {e}"),
|
||||
})?;
|
||||
if checkpoint.task_id != task_id {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: "Resume checkpoint task id does not match filename".to_string(),
|
||||
});
|
||||
}
|
||||
|
||||
// A checkpoint from an older schema stored latest-only dedup identities
|
||||
// that are not comparable to the new per-version `compose_key`
|
||||
// identities. Discard the stale sets and position, then stamp the
|
||||
// current schema so the scan restarts cleanly.
|
||||
if checkpoint.schema_version > CURRENT_CHECKPOINT_SCHEMA {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!(
|
||||
"Checkpoint schema {} is newer than supported schema {CURRENT_CHECKPOINT_SCHEMA}",
|
||||
checkpoint.schema_version
|
||||
),
|
||||
});
|
||||
}
|
||||
if checkpoint.schema_version < CURRENT_CHECKPOINT_SCHEMA {
|
||||
warn!(
|
||||
target: "rustfs::heal::resume",
|
||||
event = EVENT_HEAL_CHECKPOINT_STATE,
|
||||
component = LOG_COMPONENT_HEAL,
|
||||
subsystem = LOG_SUBSYSTEM_RESUME,
|
||||
task_id,
|
||||
found_schema = checkpoint.schema_version,
|
||||
current_schema = CURRENT_CHECKPOINT_SCHEMA,
|
||||
state = "schema_discarded",
|
||||
"Heal checkpoint schema is stale; discarding dedup sets and position"
|
||||
);
|
||||
checkpoint.processed_objects.clear();
|
||||
checkpoint.failed_objects.clear();
|
||||
checkpoint.skipped_objects.clear();
|
||||
checkpoint.current_bucket_index = 0;
|
||||
checkpoint.current_object_index = 0;
|
||||
checkpoint.schema_version = CURRENT_CHECKPOINT_SCHEMA;
|
||||
}
|
||||
|
||||
Ok(Self {
|
||||
disk,
|
||||
checkpoint: Arc::new(RwLock::new(checkpoint)),
|
||||
throttle: Mutex::new(PersistThrottle::new()),
|
||||
})
|
||||
}
|
||||
|
||||
/// check if checkpoint exists
|
||||
pub async fn has_checkpoint(disk: &DiskStore, task_id: &str) -> bool {
|
||||
if validate_resume_task_id(task_id).is_err() {
|
||||
return false;
|
||||
}
|
||||
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
|
||||
match path_to_str(&file_path) {
|
||||
Ok(path_str) => match disk.read_all(RUSTFS_META_BUCKET, path_str).await {
|
||||
Ok(data) => !data.is_empty(),
|
||||
Err(_) => false,
|
||||
},
|
||||
Err(_) => false,
|
||||
}
|
||||
}
|
||||
|
||||
/// get current checkpoint
|
||||
pub async fn get_checkpoint(&self) -> ResumeCheckpoint {
|
||||
self.checkpoint.read().await.clone()
|
||||
}
|
||||
|
||||
/// update position
|
||||
pub async fn update_position(&self, bucket_index: usize, object_index: usize) -> Result<()> {
|
||||
let mut checkpoint = self.checkpoint.write().await;
|
||||
checkpoint.update_position(bucket_index, object_index);
|
||||
drop(checkpoint);
|
||||
self.save_checkpoint_throttled().await
|
||||
}
|
||||
|
||||
/// Advance past a completed page and prune the per-object sets, then persist.
|
||||
pub async fn complete_page(&self, bucket_index: usize, object_index: usize) -> Result<()> {
|
||||
let mut checkpoint = self.checkpoint.write().await;
|
||||
checkpoint.complete_page(bucket_index, object_index);
|
||||
drop(checkpoint);
|
||||
self.save_checkpoint_throttled().await
|
||||
}
|
||||
|
||||
/// Reset the checkpoint to the start of the scan for a retry, then persist.
|
||||
pub async fn reset_for_retry(&self) -> Result<()> {
|
||||
let mut checkpoint = self.checkpoint.write().await;
|
||||
checkpoint.reset_for_retry();
|
||||
drop(checkpoint);
|
||||
self.save_checkpoint_throttled().await
|
||||
}
|
||||
|
||||
/// Add a processed object. Called once per healed object, so persistence
|
||||
/// is batched (`PERSIST_EVERY_MUTATIONS` / `PERSIST_INTERVAL`); positions
|
||||
/// and page boundaries still persist unconditionally.
|
||||
pub async fn add_processed_object(&self, object: String) -> Result<()> {
|
||||
let mut checkpoint = self.checkpoint.write().await;
|
||||
checkpoint.add_processed_object(object);
|
||||
drop(checkpoint);
|
||||
self.save_checkpoint_if_due().await
|
||||
}
|
||||
|
||||
/// add failed object (batched, see `add_processed_object`)
|
||||
pub async fn add_failed_object(&self, object: String) -> Result<()> {
|
||||
let mut checkpoint = self.checkpoint.write().await;
|
||||
checkpoint.add_failed_object(object);
|
||||
drop(checkpoint);
|
||||
self.save_checkpoint_if_due().await
|
||||
}
|
||||
|
||||
/// add skipped object (batched, see `add_processed_object`)
|
||||
pub async fn add_skipped_object(&self, object: String) -> Result<()> {
|
||||
let mut checkpoint = self.checkpoint.write().await;
|
||||
checkpoint.add_skipped_object(object);
|
||||
drop(checkpoint);
|
||||
self.save_checkpoint_if_due().await
|
||||
}
|
||||
|
||||
async fn save_checkpoint_if_due(&self) -> Result<()> {
|
||||
let should_save = self.throttle.lock().map(|mut throttle| throttle.record()).unwrap_or(true);
|
||||
if !should_save {
|
||||
return Ok(());
|
||||
}
|
||||
self.save_checkpoint_throttled().await
|
||||
}
|
||||
|
||||
async fn save_checkpoint_throttled(&self) -> Result<()> {
|
||||
let result = self.save_checkpoint().await;
|
||||
if result.is_ok()
|
||||
&& let Ok(mut throttle) = self.throttle.lock()
|
||||
{
|
||||
throttle.mark_saved();
|
||||
}
|
||||
result
|
||||
}
|
||||
|
||||
/// cleanup checkpoint
|
||||
pub async fn cleanup(&self) -> Result<()> {
|
||||
let task_id = self.checkpoint.read().await.task_id.clone();
|
||||
validate_resume_task_id(&task_id)?;
|
||||
|
||||
let checkpoint_file = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
|
||||
delete_resume_file(&self.disk, &checkpoint_file).await?;
|
||||
|
||||
debug!(
|
||||
target: "rustfs::heal::resume",
|
||||
event = EVENT_HEAL_CHECKPOINT_STATE,
|
||||
component = LOG_COMPONENT_HEAL,
|
||||
subsystem = LOG_SUBSYSTEM_RESUME,
|
||||
task_id,
|
||||
state = "cleaned",
|
||||
"Heal checkpoint cleaned"
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// save checkpoint to disk
|
||||
async fn save_checkpoint(&self) -> Result<()> {
|
||||
let checkpoint = self.checkpoint.read().await;
|
||||
validate_resume_task_id(&checkpoint.task_id)?;
|
||||
let checkpoint_data = serde_json::to_vec(&*checkpoint).map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to serialize checkpoint: {e}"),
|
||||
})?;
|
||||
|
||||
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{}_{}", checkpoint.task_id, RESUME_CHECKPOINT_FILE));
|
||||
|
||||
let path_str = path_to_str(&file_path)?;
|
||||
self.disk
|
||||
.write_all(RUSTFS_META_BUCKET, path_str, checkpoint_data.into())
|
||||
.await
|
||||
.map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to save checkpoint: {e}"),
|
||||
})?;
|
||||
|
||||
debug!(
|
||||
target: "rustfs::heal::resume",
|
||||
event = EVENT_HEAL_CHECKPOINT_STATE,
|
||||
component = LOG_COMPONENT_HEAL,
|
||||
subsystem = LOG_SUBSYSTEM_RESUME,
|
||||
task_id = %checkpoint.task_id,
|
||||
state = "saved",
|
||||
"Heal checkpoint persisted"
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// read checkpoint file from disk
|
||||
async fn read_checkpoint_file(disk: &DiskStore, task_id: &str) -> Result<Vec<u8>> {
|
||||
validate_resume_task_id(task_id)?;
|
||||
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
|
||||
|
||||
let path_str = path_to_str(&file_path)?;
|
||||
disk.read_all(RUSTFS_META_BUCKET, path_str)
|
||||
.await
|
||||
.map(|bytes| bytes.to_vec())
|
||||
.map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to read checkpoint file: {e}"),
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,688 @@
|
||||
// Copyright 2024 RustFS Team
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
use crate::{Error, Result};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::collections::HashSet;
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
|
||||
use super::super::HealDiskExt as _;
|
||||
|
||||
use super::super::storage_api::owner::{EcstoreConditionalFileUpdate, EcstoreDiskBytes};
|
||||
use super::{
|
||||
DiskError, DiskStore, RUSTFS_META_BUCKET, ResumeManager, ResumeState, delete_resume_file, ensure_replacement_recovery_dir,
|
||||
injected_replacement_proof_write_error, is_replacement_intent, legacy_replacement_completion_proof_path, path_to_str,
|
||||
replacement_completion_proof_path, replacement_intent_seal_path, replacement_recovery_conflict,
|
||||
replacement_recovery_corruption, validate_resume_task_id,
|
||||
};
|
||||
|
||||
/// Durable-proof schema version.
|
||||
const CURRENT_REPLACEMENT_COMPLETION_PROOF_SCHEMA: u32 = 1;
|
||||
|
||||
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum ReplacementPhase {
|
||||
#[default]
|
||||
None,
|
||||
Intent,
|
||||
Rebuilding,
|
||||
Verified,
|
||||
CleanupPending,
|
||||
Abandoned,
|
||||
}
|
||||
|
||||
/// Target-specific state for a durable automatic replacement generation.
|
||||
///
|
||||
/// This is deliberately separate from the legacy background-heal status
|
||||
/// contract. Consumers must treat [`Self::Unknown`] as non-definitive.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum ReplacementRecoveryState {
|
||||
WaitingForReplacement,
|
||||
Running,
|
||||
Incomplete,
|
||||
Unrecoverable,
|
||||
CleanupPending,
|
||||
Completed,
|
||||
Unknown,
|
||||
}
|
||||
|
||||
/// Read-only status derived from one durable replacement generation.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct ReplacementRecoveryRecord {
|
||||
pub task_id: String,
|
||||
pub state: ReplacementRecoveryState,
|
||||
pub generation: Option<String>,
|
||||
pub set_disk_id: Option<String>,
|
||||
pub target_slots: Vec<String>,
|
||||
pub reason: Option<String>,
|
||||
pub verified_at: Option<u64>,
|
||||
}
|
||||
|
||||
impl ReplacementRecoveryRecord {
|
||||
pub(super) fn from_state(state: ResumeState) -> Option<Self> {
|
||||
if !is_replacement_intent(&state) {
|
||||
return None;
|
||||
}
|
||||
|
||||
let invariant_holds = state.replacement_generation.as_deref() == Some(state.task_id.as_str())
|
||||
&& replacement_targets_match_identities(&state.replacement_targets, &state.replacement_target_identities);
|
||||
if !invariant_holds {
|
||||
return Some(Self::unknown(
|
||||
state.task_id,
|
||||
"durable replacement state violates its generation or target identity binding",
|
||||
));
|
||||
}
|
||||
|
||||
let (state_kind, reason) = if !state.completed && state.retry_count >= state.max_retries {
|
||||
(
|
||||
ReplacementRecoveryState::Unrecoverable,
|
||||
Some("replacement retry budget exhausted".to_string()),
|
||||
)
|
||||
} else if let Some(reason) = state.error_message.clone() {
|
||||
(ReplacementRecoveryState::Incomplete, Some(reason))
|
||||
} else {
|
||||
match state.replacement_phase {
|
||||
ReplacementPhase::Intent => (ReplacementRecoveryState::WaitingForReplacement, None),
|
||||
ReplacementPhase::Rebuilding => (ReplacementRecoveryState::Running, None),
|
||||
ReplacementPhase::Verified | ReplacementPhase::CleanupPending => (ReplacementRecoveryState::CleanupPending, None),
|
||||
ReplacementPhase::Abandoned => (
|
||||
ReplacementRecoveryState::Unrecoverable,
|
||||
Some("replacement generation was abandoned".to_string()),
|
||||
),
|
||||
ReplacementPhase::None => (ReplacementRecoveryState::Unknown, Some("replacement phase is missing".to_string())),
|
||||
}
|
||||
};
|
||||
|
||||
Some(Self {
|
||||
task_id: state.task_id,
|
||||
state: state_kind,
|
||||
generation: state.replacement_generation,
|
||||
set_disk_id: Some(state.set_disk_id),
|
||||
target_slots: state.replacement_targets,
|
||||
reason,
|
||||
verified_at: None,
|
||||
})
|
||||
}
|
||||
|
||||
pub(super) fn from_completion_proof(proof: &ReplacementCompletionProof) -> Self {
|
||||
Self {
|
||||
task_id: proof.task_id.clone(),
|
||||
state: ReplacementRecoveryState::Completed,
|
||||
generation: Some(proof.replacement_generation.clone()),
|
||||
set_disk_id: Some(proof.set_disk_id.clone()),
|
||||
target_slots: proof.replacement_targets.clone(),
|
||||
reason: None,
|
||||
verified_at: Some(proof.verified_at),
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn unknown(task_id: String, reason: &str) -> Self {
|
||||
Self {
|
||||
task_id,
|
||||
state: ReplacementRecoveryState::Unknown,
|
||||
generation: None,
|
||||
set_disk_id: None,
|
||||
target_slots: Vec::new(),
|
||||
reason: Some(reason.to_string()),
|
||||
verified_at: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn replacement_targets_match_identities(targets: &[String], identities: &[ReplacementTargetIdentity]) -> bool {
|
||||
!targets.is_empty()
|
||||
&& targets.len() == identities.len()
|
||||
&& targets.iter().collect::<HashSet<_>>().len() == targets.len()
|
||||
&& identities.iter().map(|identity| &identity.endpoint).eq(targets.iter())
|
||||
}
|
||||
|
||||
/// Stable evidence for the mounted replacement instance that owns a repair
|
||||
/// generation. Endpoint text alone is not sufficient because a later disk can
|
||||
/// be mounted at the same configured path.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub struct ReplacementTargetIdentity {
|
||||
pub endpoint: String,
|
||||
pub canonical_path: String,
|
||||
pub physical_device_ids: Vec<String>,
|
||||
pub filesystem_identity: String,
|
||||
}
|
||||
|
||||
/// Durable terminal evidence for one automatic replacement generation. This
|
||||
/// lives on the healthy non-target anchor rather than in the resumable state,
|
||||
/// because resume cleanup must not erase proof that the generation completed.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub(crate) struct ReplacementCompletionProof {
|
||||
pub schema_version: u32,
|
||||
pub task_id: String,
|
||||
pub replacement_generation: String,
|
||||
pub set_disk_id: String,
|
||||
pub replacement_targets: Vec<String>,
|
||||
pub replacement_target_identities: Vec<ReplacementTargetIdentity>,
|
||||
pub verified_at: u64,
|
||||
}
|
||||
|
||||
impl ReplacementCompletionProof {
|
||||
pub(super) fn from_state(state: &ResumeState, verified_at: u64) -> Result<Self> {
|
||||
let replacement_generation = state
|
||||
.replacement_generation
|
||||
.clone()
|
||||
.ok_or_else(|| Error::TaskExecutionFailed {
|
||||
message: format!("Replacement completion has no generation for task {}", state.task_id),
|
||||
})?;
|
||||
if replacement_generation != state.task_id
|
||||
|| state.replacement_targets.is_empty()
|
||||
|| state
|
||||
.replacement_target_identities
|
||||
.iter()
|
||||
.map(|identity| &identity.endpoint)
|
||||
.collect::<Vec<_>>()
|
||||
!= state.replacement_targets.iter().collect::<Vec<_>>()
|
||||
{
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Replacement completion identity does not match task {}", state.task_id),
|
||||
});
|
||||
}
|
||||
|
||||
Ok(Self {
|
||||
schema_version: CURRENT_REPLACEMENT_COMPLETION_PROOF_SCHEMA,
|
||||
task_id: state.task_id.clone(),
|
||||
replacement_generation,
|
||||
set_disk_id: state.set_disk_id.clone(),
|
||||
replacement_targets: state.replacement_targets.clone(),
|
||||
replacement_target_identities: state.replacement_target_identities.clone(),
|
||||
verified_at,
|
||||
})
|
||||
}
|
||||
|
||||
fn matches_state(&self, state: &ResumeState) -> bool {
|
||||
self.schema_version == CURRENT_REPLACEMENT_COMPLETION_PROOF_SCHEMA
|
||||
&& self.task_id == state.task_id
|
||||
&& state.replacement_generation.as_deref() == Some(self.replacement_generation.as_str())
|
||||
&& self.set_disk_id == state.set_disk_id
|
||||
&& self.replacement_targets == state.replacement_targets
|
||||
&& self.replacement_target_identities == state.replacement_target_identities
|
||||
}
|
||||
|
||||
fn validate(&self, expected_task_id: &str) -> Result<()> {
|
||||
if self.schema_version != CURRENT_REPLACEMENT_COMPLETION_PROOF_SCHEMA {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Replacement completion proof schema {} is unsupported", self.schema_version),
|
||||
});
|
||||
}
|
||||
validate_resume_task_id(expected_task_id)?;
|
||||
if self.task_id != expected_task_id
|
||||
|| self.replacement_generation != self.task_id
|
||||
|| self.set_disk_id.is_empty()
|
||||
|| self.verified_at == 0
|
||||
|| !replacement_targets_match_identities(&self.replacement_targets, &self.replacement_target_identities)
|
||||
{
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Replacement completion proof does not match task {expected_task_id}"),
|
||||
});
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn replacement_target_identities_match(
|
||||
expected: &[ReplacementTargetIdentity],
|
||||
actual: &[ReplacementTargetIdentity],
|
||||
) -> bool {
|
||||
let mut expected = expected.to_vec();
|
||||
let mut actual = actual.to_vec();
|
||||
expected.sort_by(|left, right| left.endpoint.cmp(&right.endpoint));
|
||||
actual.sort_by(|left, right| left.endpoint.cmp(&right.endpoint));
|
||||
expected == actual
|
||||
}
|
||||
|
||||
/// Build the canonical, provably-injective dedup identity for an object
|
||||
/// version. Length-prefixing the object key makes the encoding injective: no
|
||||
/// two distinct `(object, version_id)` pairs can collide, even for adversarial
|
||||
/// keys containing `:` or embedded null bytes. This is the single source of
|
||||
/// truth for per-version dedup across the heal loop and the checkpoint sets.
|
||||
pub fn compose_key(object: &str, version_id: Option<&str>) -> String {
|
||||
format!("{}:{}{}", object.len(), object, version_id.unwrap_or(""))
|
||||
}
|
||||
|
||||
impl ResumeManager {
|
||||
/// Seal a durably published intent before the caller may format a target.
|
||||
/// A torn intent without this seal is known to have failed before its
|
||||
/// creator returned and can be atomically recreated on retry.
|
||||
pub(super) async fn ensure_replacement_intent_seal(&self) -> Result<()> {
|
||||
let task_id = self.state.read().await.task_id.clone();
|
||||
validate_resume_task_id(&task_id)?;
|
||||
let path = replacement_intent_seal_path(&task_id);
|
||||
let path = path_to_str(&path)?;
|
||||
match self.disk.read_all(RUSTFS_META_BUCKET, path).await {
|
||||
Ok(_) => return Ok(()),
|
||||
Err(DiskError::FileNotFound) => {}
|
||||
Err(error) => {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to read replacement intent seal: {error}"),
|
||||
});
|
||||
}
|
||||
}
|
||||
self.disk
|
||||
.write_all(RUSTFS_META_BUCKET, path, b"sealed".as_slice().into())
|
||||
.await
|
||||
.map_err(|error| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to save replacement intent seal: {error}"),
|
||||
})
|
||||
}
|
||||
|
||||
pub async fn mark_replacement_rebuilding(
|
||||
&self,
|
||||
mut replacement_target_identities: Vec<ReplacementTargetIdentity>,
|
||||
) -> Result<()> {
|
||||
replacement_target_identities.sort_by(|left, right| left.endpoint.cmp(&right.endpoint));
|
||||
replacement_target_identities.dedup_by(|left, right| left.endpoint == right.endpoint);
|
||||
let mut state = self.state.write().await;
|
||||
if !matches!(state.replacement_phase, ReplacementPhase::Intent | ReplacementPhase::Rebuilding) {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Replacement intent is not active for task {}", state.task_id),
|
||||
});
|
||||
}
|
||||
if replacement_target_identities
|
||||
.iter()
|
||||
.map(|identity| &identity.endpoint)
|
||||
.collect::<Vec<_>>()
|
||||
!= state.replacement_targets.iter().collect::<Vec<_>>()
|
||||
{
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Replacement identities do not match targets for task {}", state.task_id),
|
||||
});
|
||||
}
|
||||
if !replacement_target_identities_match(&state.replacement_target_identities, &replacement_target_identities) {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Replacement target changed after format for task {}", state.task_id),
|
||||
});
|
||||
}
|
||||
state.replacement_phase = ReplacementPhase::Rebuilding;
|
||||
state.last_update = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
|
||||
drop(state);
|
||||
self.save_state_strict().await
|
||||
}
|
||||
|
||||
/// Persist survivor-anchor completion proof before transitioning this
|
||||
/// resumable state to `Verified`. If proof persistence fails, this state
|
||||
/// stays rebuildable and the caller must retain the healing marker.
|
||||
pub async fn mark_replacement_completed_and_verified(&self) -> Result<()> {
|
||||
let state = self.state.read().await.clone();
|
||||
if !matches!(state.replacement_phase, ReplacementPhase::Intent | ReplacementPhase::Rebuilding) {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Replacement verification is not active for task {}", state.task_id),
|
||||
});
|
||||
}
|
||||
let proof = self.write_replacement_completion_proof(&state, None).await?;
|
||||
|
||||
let mut state = self.state.write().await;
|
||||
if !matches!(state.replacement_phase, ReplacementPhase::Intent | ReplacementPhase::Rebuilding) {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Replacement verification changed for task {}", state.task_id),
|
||||
});
|
||||
}
|
||||
state.mark_completed();
|
||||
state.replacement_phase = ReplacementPhase::Verified;
|
||||
state.last_update = proof.verified_at;
|
||||
drop(state);
|
||||
self.save_state_strict().await
|
||||
}
|
||||
|
||||
/// Verify or backfill the terminal proof before marker removal or resume
|
||||
/// cleanup. This supports restart recovery from a `Verified` state written
|
||||
/// by a prior binary that did not yet have a separate proof record.
|
||||
pub(crate) async fn ensure_replacement_completion_proof(&self) -> Result<ReplacementCompletionProof> {
|
||||
let state = self.state.read().await.clone();
|
||||
if !state.completed || !matches!(state.replacement_phase, ReplacementPhase::Verified | ReplacementPhase::CleanupPending) {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Replacement completion is not verified for task {}", state.task_id),
|
||||
});
|
||||
}
|
||||
self.write_replacement_completion_proof(&state, Some(state.last_update)).await
|
||||
}
|
||||
|
||||
/// Record that the healing markers have been removed, so a later retry can
|
||||
/// safely delete the remaining resume artifacts without touching markers.
|
||||
pub async fn mark_replacement_cleanup_pending(&self) -> Result<()> {
|
||||
let mut state = self.state.write().await;
|
||||
if !state.completed || !matches!(state.replacement_phase, ReplacementPhase::Verified | ReplacementPhase::CleanupPending) {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Replacement cleanup is not ready for task {}", state.task_id),
|
||||
});
|
||||
}
|
||||
state.replacement_phase = ReplacementPhase::CleanupPending;
|
||||
state.last_update = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
|
||||
drop(state);
|
||||
self.save_state_strict().await
|
||||
}
|
||||
|
||||
/// Load the durable terminal proof from the healthy survivor anchor.
|
||||
pub(crate) async fn load_replacement_completion_proof(disk: DiskStore, task_id: &str) -> Result<ReplacementCompletionProof> {
|
||||
Self::replacement_completion_proof_if_present(disk, task_id)
|
||||
.await?
|
||||
.ok_or_else(|| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to read replacement completion proof: proof is missing for task {task_id}"),
|
||||
})
|
||||
}
|
||||
|
||||
async fn replacement_completion_proof_if_present(
|
||||
disk: DiskStore,
|
||||
task_id: &str,
|
||||
) -> Result<Option<ReplacementCompletionProof>> {
|
||||
validate_resume_task_id(task_id)?;
|
||||
let mut proofs = Vec::new();
|
||||
for path in [
|
||||
replacement_completion_proof_path(task_id),
|
||||
legacy_replacement_completion_proof_path(task_id),
|
||||
] {
|
||||
let path_str = path_to_str(&path)?;
|
||||
let bytes = match disk.read_all(RUSTFS_META_BUCKET, path_str).await {
|
||||
Ok(bytes) => bytes,
|
||||
Err(DiskError::FileNotFound) => continue,
|
||||
Err(error) => {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to read replacement completion proof: {error}"),
|
||||
});
|
||||
}
|
||||
};
|
||||
let proof: ReplacementCompletionProof =
|
||||
serde_json::from_slice(&bytes).map_err(|error| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to deserialize replacement completion proof: {error}"),
|
||||
})?;
|
||||
proof.validate(task_id)?;
|
||||
proofs.push(proof);
|
||||
}
|
||||
|
||||
match proofs.as_slice() {
|
||||
[] => Ok(None),
|
||||
[proof] => Ok(Some(proof.clone())),
|
||||
[proof, legacy_proof] if proof == legacy_proof => Ok(Some(proof.clone())),
|
||||
_ => Err(replacement_recovery_conflict(format!(
|
||||
"Replacement completion proof conflicts with legacy proof for task {task_id}"
|
||||
))),
|
||||
}
|
||||
}
|
||||
|
||||
/// Reconcile the proof-first publication order after a crash. A matching
|
||||
/// proof is durable evidence that rebuilding finished, so it must win over
|
||||
/// an older active state before a retry may format the target again.
|
||||
pub(super) async fn reconcile_replacement_completion_proof(&self) -> Result<()> {
|
||||
let task_id = self.state.read().await.task_id.clone();
|
||||
let Some(proof) = Self::replacement_completion_proof_if_present(self.disk.clone(), &task_id).await? else {
|
||||
return Ok(());
|
||||
};
|
||||
|
||||
let mut state = self.state.write().await;
|
||||
if !proof.matches_state(&state) {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Replacement completion proof does not match active intent for task {}", state.task_id),
|
||||
});
|
||||
}
|
||||
if state.completed && matches!(state.replacement_phase, ReplacementPhase::Verified | ReplacementPhase::CleanupPending) {
|
||||
return Ok(());
|
||||
}
|
||||
if state.completed || !matches!(state.replacement_phase, ReplacementPhase::Intent | ReplacementPhase::Rebuilding) {
|
||||
return Err(replacement_recovery_conflict(format!(
|
||||
"Replacement completion proof conflicts with state for task {}",
|
||||
state.task_id
|
||||
)));
|
||||
}
|
||||
|
||||
state.mark_completed();
|
||||
state.replacement_phase = ReplacementPhase::Verified;
|
||||
state.last_update = proof.verified_at;
|
||||
drop(state);
|
||||
self.save_state_strict().await
|
||||
}
|
||||
|
||||
pub(super) async fn migrate_legacy_replacement_completion_proof(disk: &DiskStore, task_id: &str) -> Result<bool> {
|
||||
validate_resume_task_id(task_id)?;
|
||||
let legacy_path = legacy_replacement_completion_proof_path(task_id);
|
||||
let legacy_path_str = path_to_str(&legacy_path)?;
|
||||
let legacy_bytes = match disk.read_all(RUSTFS_META_BUCKET, legacy_path_str).await {
|
||||
Ok(bytes) => bytes,
|
||||
Err(DiskError::FileNotFound) => return Ok(false),
|
||||
Err(error) => {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to read legacy replacement completion proof: {error}"),
|
||||
});
|
||||
}
|
||||
};
|
||||
let legacy_proof: ReplacementCompletionProof = serde_json::from_slice(&legacy_bytes).map_err(|error| {
|
||||
replacement_recovery_corruption(format!("Failed to deserialize legacy replacement completion proof: {error}"))
|
||||
})?;
|
||||
legacy_proof
|
||||
.validate(task_id)
|
||||
.map_err(|error| replacement_recovery_corruption(format!("Invalid legacy replacement completion proof: {error}")))?;
|
||||
|
||||
ensure_replacement_recovery_dir(disk)
|
||||
.await
|
||||
.map_err(|error| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to create replacement recovery directory: {error}"),
|
||||
})?;
|
||||
let path = replacement_completion_proof_path(task_id);
|
||||
let path_str = path_to_str(&path)?;
|
||||
for _ in 0..2 {
|
||||
match disk.read_all(RUSTFS_META_BUCKET, path_str).await {
|
||||
Ok(bytes) => {
|
||||
let proof: ReplacementCompletionProof =
|
||||
serde_json::from_slice(&bytes).map_err(|error| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to deserialize replacement completion proof: {error}"),
|
||||
})?;
|
||||
proof.validate(task_id).map_err(|error| {
|
||||
replacement_recovery_corruption(format!("Invalid replacement completion proof: {error}"))
|
||||
})?;
|
||||
if proof != legacy_proof {
|
||||
return Err(replacement_recovery_conflict(format!(
|
||||
"Replacement completion proof conflicts with legacy proof for task {task_id}"
|
||||
)));
|
||||
}
|
||||
delete_resume_file(disk, &legacy_path).await?;
|
||||
return Ok(true);
|
||||
}
|
||||
Err(DiskError::FileNotFound) => {}
|
||||
Err(error) => {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to read replacement completion proof: {error}"),
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
match super::super::storage_api::owner::EcstoreDiskAPI::compare_and_update_file(
|
||||
disk.as_ref(),
|
||||
RUSTFS_META_BUCKET,
|
||||
path_str,
|
||||
None,
|
||||
Some(legacy_bytes.clone()),
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(EcstoreConditionalFileUpdate::Updated) => {
|
||||
delete_resume_file(disk, &legacy_path).await?;
|
||||
return Ok(true);
|
||||
}
|
||||
Ok(EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch) => continue,
|
||||
Err(error) => {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to migrate replacement completion proof: {error}"),
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Err(Error::TaskExecutionFailed {
|
||||
message: format!("Replacement completion proof changed while migrating task {task_id}"),
|
||||
})
|
||||
}
|
||||
|
||||
pub async fn abandon_replacement_intent(&self) -> Result<()> {
|
||||
let mut state = self.state.write().await;
|
||||
if matches!(state.replacement_phase, ReplacementPhase::Abandoned) {
|
||||
return Ok(());
|
||||
}
|
||||
state.replacement_phase = ReplacementPhase::Abandoned;
|
||||
state.last_update = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
|
||||
drop(state);
|
||||
self.save_state_strict().await
|
||||
}
|
||||
|
||||
pub async fn set_replacement_targets(&self, replacement_targets: Vec<String>) -> Result<()> {
|
||||
{
|
||||
let mut state = self.state.write().await;
|
||||
state.replacement_targets = replacement_targets;
|
||||
}
|
||||
self.save_state().await
|
||||
}
|
||||
|
||||
pub(super) async fn publish_new_replacement_intent(&self, expected: Option<EcstoreDiskBytes>) -> Result<()> {
|
||||
let state = self.state.read().await.clone();
|
||||
validate_resume_task_id(&state.task_id)?;
|
||||
let state_data = EcstoreDiskBytes::from(serde_json::to_vec(&state).map_err(|error| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to serialize resume state: {error}"),
|
||||
})?);
|
||||
let path = self.state_file.path(&state.task_id);
|
||||
let path = path_to_str(&path)?;
|
||||
|
||||
ensure_replacement_recovery_dir(&self.disk)
|
||||
.await
|
||||
.map_err(|error| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to create replacement recovery directory: {error}"),
|
||||
})?;
|
||||
match super::super::storage_api::owner::EcstoreDiskAPI::compare_and_update_file(
|
||||
self.disk.as_ref(),
|
||||
RUSTFS_META_BUCKET,
|
||||
path,
|
||||
expected,
|
||||
Some(state_data),
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(EcstoreConditionalFileUpdate::Updated) => Ok(()),
|
||||
Ok(EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch) => {
|
||||
Err(Error::TaskExecutionFailed {
|
||||
message: format!("Replacement intent changed before publication for task {}", state.task_id),
|
||||
})
|
||||
}
|
||||
Err(error) => Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to save resume state: {error}"),
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
async fn write_replacement_completion_proof(
|
||||
&self,
|
||||
state: &ResumeState,
|
||||
verified_at: Option<u64>,
|
||||
) -> Result<ReplacementCompletionProof> {
|
||||
ensure_replacement_recovery_dir(&self.disk)
|
||||
.await
|
||||
.map_err(|error| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to create replacement recovery directory: {error}"),
|
||||
})?;
|
||||
let path = replacement_completion_proof_path(&state.task_id);
|
||||
let path_str = path_to_str(&path)?;
|
||||
let proof = ReplacementCompletionProof::from_state(
|
||||
state,
|
||||
verified_at.unwrap_or_else(|| SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs()),
|
||||
)?;
|
||||
let proof_data = EcstoreDiskBytes::from(serde_json::to_vec(&proof).map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to serialize replacement completion proof: {e}"),
|
||||
})?);
|
||||
if let Some(error) = injected_replacement_proof_write_error(path_str) {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to save replacement completion proof: {error}"),
|
||||
});
|
||||
}
|
||||
|
||||
// Publish through the disk CAS primitive: `write_all` can expose a
|
||||
// partially written proof to a crash/restart reader. If a prior
|
||||
// version left torn bytes behind, replace exactly the observed bytes;
|
||||
// a concurrently published valid proof is never overwritten.
|
||||
for _ in 0..2 {
|
||||
let expected = match self.disk.read_all(RUSTFS_META_BUCKET, path_str).await {
|
||||
Ok(existing) => match serde_json::from_slice::<ReplacementCompletionProof>(&existing) {
|
||||
Ok(existing_proof) => {
|
||||
existing_proof.validate(&state.task_id)?;
|
||||
if existing_proof.matches_state(state) {
|
||||
return Ok(existing_proof);
|
||||
}
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Replacement completion proof does not match task {}", state.task_id),
|
||||
});
|
||||
}
|
||||
Err(_) => Some(existing),
|
||||
},
|
||||
Err(DiskError::FileNotFound) => None,
|
||||
Err(error) => {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to read replacement completion proof: {error}"),
|
||||
});
|
||||
}
|
||||
};
|
||||
|
||||
match super::super::storage_api::owner::EcstoreDiskAPI::compare_and_update_file(
|
||||
self.disk.as_ref(),
|
||||
RUSTFS_META_BUCKET,
|
||||
path_str,
|
||||
expected,
|
||||
Some(proof_data.clone()),
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(EcstoreConditionalFileUpdate::Updated) => return Ok(proof),
|
||||
Ok(EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch) => continue,
|
||||
Err(error) => {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to save replacement completion proof: {error}"),
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Err(Error::TaskExecutionFailed {
|
||||
message: format!("Replacement completion proof changed while publishing task {}", state.task_id),
|
||||
})
|
||||
}
|
||||
|
||||
pub(super) async fn write_replacement_intent_state(
|
||||
&self,
|
||||
path: &str,
|
||||
state_data: EcstoreDiskBytes,
|
||||
) -> std::result::Result<(), DiskError> {
|
||||
ensure_replacement_recovery_dir(&self.disk).await?;
|
||||
for _ in 0..2 {
|
||||
let expected = match self.disk.read_all(RUSTFS_META_BUCKET, path).await {
|
||||
Ok(existing) => Some(existing),
|
||||
Err(DiskError::FileNotFound) => None,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
match super::super::storage_api::owner::EcstoreDiskAPI::compare_and_update_file(
|
||||
self.disk.as_ref(),
|
||||
RUSTFS_META_BUCKET,
|
||||
path,
|
||||
expected,
|
||||
Some(state_data.clone()),
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(EcstoreConditionalFileUpdate::Updated) => return Ok(()),
|
||||
Ok(EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch) => continue,
|
||||
Err(error) => return Err(error),
|
||||
}
|
||||
}
|
||||
Err(DiskError::other("replacement intent changed while publishing"))
|
||||
}
|
||||
}
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,311 @@
|
||||
// Copyright 2024 RustFS Team
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
use crate::{Error, Result};
|
||||
use std::collections::HashSet;
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
use tracing::{debug, warn};
|
||||
use uuid::Uuid;
|
||||
|
||||
use super::super::{BUCKET_META_PREFIX, DiskError, DiskStore, HealDiskExt as _, RUSTFS_META_BUCKET};
|
||||
use super::replacement::{ReplacementPhase, ReplacementRecoveryRecord};
|
||||
use super::{
|
||||
EVENT_HEAL_RESUME_STATE, LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, REPLACEMENT_COMPLETION_PROOF_FILE,
|
||||
REPLACEMENT_INTENT_FILE, RESUME_STATE_FILE, ResumeManager, ResumeStateFile, is_replacement_intent, path_to_str,
|
||||
replacement_recovery_corruption_for_state_load, replacement_recovery_dir, validate_resume_task_id,
|
||||
};
|
||||
|
||||
/// resume utils
|
||||
pub struct ResumeUtils;
|
||||
|
||||
impl ResumeUtils {
|
||||
/// generate unique task id
|
||||
pub fn generate_task_id() -> String {
|
||||
Uuid::new_v4().to_string()
|
||||
}
|
||||
|
||||
/// check if task can be resumed
|
||||
pub async fn can_resume_task(disk: &DiskStore, task_id: &str) -> bool {
|
||||
ResumeManager::has_resume_state(disk, task_id).await
|
||||
}
|
||||
|
||||
/// get all resumable task ids
|
||||
pub async fn get_resumable_tasks(disk: &DiskStore) -> Result<Vec<String>> {
|
||||
// List all files in the buckets metadata directory
|
||||
let entries = match disk.list_dir("", RUSTFS_META_BUCKET, BUCKET_META_PREFIX, -1).await {
|
||||
Ok(entries) => entries,
|
||||
Err(e) => {
|
||||
debug!(
|
||||
target: "rustfs::heal::resume",
|
||||
event = EVENT_HEAL_RESUME_STATE,
|
||||
component = LOG_COMPONENT_HEAL,
|
||||
subsystem = LOG_SUBSYSTEM_RESUME,
|
||||
state = "list_failed",
|
||||
error = %e,
|
||||
"Heal resume state listing failed"
|
||||
);
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
};
|
||||
|
||||
let mut task_ids = Vec::new();
|
||||
|
||||
// Filter files that end with ahm_resume_state.json and extract task IDs
|
||||
for entry in entries {
|
||||
if entry.ends_with(&format!("_{RESUME_STATE_FILE}")) {
|
||||
// Extract task ID from filename: {task_id}_ahm_resume_state.json
|
||||
if let Some(task_id) = entry.strip_suffix(&format!("_{RESUME_STATE_FILE}"))
|
||||
&& validate_resume_task_id(task_id).is_ok()
|
||||
{
|
||||
task_ids.push(task_id.to_string());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
debug!(
|
||||
target: "rustfs::heal::resume",
|
||||
event = EVENT_HEAL_RESUME_STATE,
|
||||
component = LOG_COMPONENT_HEAL,
|
||||
subsystem = LOG_SUBSYSTEM_RESUME,
|
||||
task_count = task_ids.len(),
|
||||
state = "listed",
|
||||
"Heal resume states listed"
|
||||
);
|
||||
Ok(task_ids)
|
||||
}
|
||||
|
||||
/// Return replacement intent task IDs from the dedicated recovery
|
||||
/// directory. Periodic recovery must never enumerate the ordinary resume
|
||||
/// directory, whose cardinality is unrelated to replacement work.
|
||||
pub async fn get_replacement_intent_tasks(disk: &DiskStore) -> Result<Vec<String>> {
|
||||
let entries = Self::replacement_recovery_entries(disk).await?;
|
||||
let suffix = format!("_{REPLACEMENT_INTENT_FILE}");
|
||||
let mut task_ids = HashSet::new();
|
||||
|
||||
for entry in entries {
|
||||
if let Some(task_id) = entry.strip_suffix(&suffix)
|
||||
&& validate_resume_task_id(task_id).is_ok()
|
||||
{
|
||||
task_ids.insert(task_id.to_string());
|
||||
continue;
|
||||
}
|
||||
}
|
||||
|
||||
let mut task_ids = task_ids.into_iter().collect::<Vec<_>>();
|
||||
task_ids.sort_unstable();
|
||||
Ok(task_ids)
|
||||
}
|
||||
|
||||
async fn replacement_recovery_entries(disk: &DiskStore) -> Result<Vec<String>> {
|
||||
let recovery_dir = replacement_recovery_dir();
|
||||
let recovery_dir = path_to_str(&recovery_dir)?;
|
||||
match disk.list_dir("", RUSTFS_META_BUCKET, recovery_dir, -1).await {
|
||||
Ok(entries) => Ok(entries),
|
||||
Err(DiskError::FileNotFound) => Ok(Vec::new()),
|
||||
Err(error @ DiskError::UnformattedDisk) => Err(error.into()),
|
||||
Err(error) => Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to list replacement recovery records: {error}"),
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
/// Migrate flat replacement artifacts from earlier builds exactly once at
|
||||
/// manager startup. The normal scanner only uses the dedicated directory;
|
||||
/// ordinary resume JSON is never read on its periodic path.
|
||||
pub async fn migrate_legacy_replacement_records(disk: &DiskStore) -> Result<()> {
|
||||
let entries = disk
|
||||
.list_dir("", RUSTFS_META_BUCKET, BUCKET_META_PREFIX, -1)
|
||||
.await
|
||||
.map_err(|error| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to list legacy replacement records: {error}"),
|
||||
})?;
|
||||
let ordinary_suffix = format!("_{RESUME_STATE_FILE}");
|
||||
let intent_suffix = format!("_{REPLACEMENT_INTENT_FILE}");
|
||||
let proof_suffix = format!("_{REPLACEMENT_COMPLETION_PROOF_FILE}");
|
||||
let mut ordinary_task_ids = HashSet::new();
|
||||
let mut intent_task_ids = HashSet::new();
|
||||
let mut proof_task_ids = HashSet::new();
|
||||
|
||||
for entry in entries {
|
||||
if let Some(task_id) = entry.strip_suffix(&intent_suffix)
|
||||
&& validate_resume_task_id(task_id).is_ok()
|
||||
{
|
||||
intent_task_ids.insert(task_id.to_string());
|
||||
continue;
|
||||
}
|
||||
if let Some(task_id) = entry.strip_suffix(&ordinary_suffix)
|
||||
&& validate_resume_task_id(task_id).is_ok()
|
||||
{
|
||||
ordinary_task_ids.insert(task_id.to_string());
|
||||
continue;
|
||||
}
|
||||
if let Some(task_id) = entry.strip_suffix(&proof_suffix)
|
||||
&& validate_resume_task_id(task_id).is_ok()
|
||||
{
|
||||
proof_task_ids.insert(task_id.to_string());
|
||||
}
|
||||
}
|
||||
|
||||
let mut state_task_ids = intent_task_ids.into_iter().collect::<Vec<_>>();
|
||||
state_task_ids.extend(ordinary_task_ids);
|
||||
state_task_ids.sort_unstable();
|
||||
state_task_ids.dedup();
|
||||
for task_id in state_task_ids {
|
||||
let has_flat_intent = ResumeManager::has_state_file(disk, &task_id, ResumeStateFile::LegacyReplacementIntent).await;
|
||||
if !has_flat_intent {
|
||||
let manager = ResumeManager::load_from_disk(disk.clone(), &task_id).await.map_err(|error| {
|
||||
replacement_recovery_corruption_for_state_load(
|
||||
format!("Failed to load legacy replacement recovery candidate {task_id}"),
|
||||
error,
|
||||
)
|
||||
})?;
|
||||
if !is_replacement_intent(&manager.get_state().await) {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
ResumeManager::load_replacement_intent(disk.clone(), &task_id).await?;
|
||||
}
|
||||
|
||||
for task_id in proof_task_ids {
|
||||
ResumeManager::migrate_legacy_replacement_completion_proof(disk, &task_id).await?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Return all durable replacement states and completion proofs stored on
|
||||
/// one survivor disk. Unlike the legacy resumable-task helper, listing
|
||||
/// failures are returned to the caller so an observability surface cannot
|
||||
/// silently turn an unreadable durable record into a green result.
|
||||
pub async fn get_replacement_recovery_records(disk: &DiskStore) -> Result<Vec<ReplacementRecoveryRecord>> {
|
||||
let entries = Self::replacement_recovery_entries(disk).await?;
|
||||
let proof_suffix = format!("_{REPLACEMENT_COMPLETION_PROOF_FILE}");
|
||||
let mut records = Vec::new();
|
||||
let mut intent_task_ids = HashSet::new();
|
||||
|
||||
for task_id in Self::get_replacement_intent_tasks(disk).await? {
|
||||
let state = ResumeManager::load_replacement_intent(disk.clone(), &task_id)
|
||||
.await?
|
||||
.get_state()
|
||||
.await;
|
||||
intent_task_ids.insert(task_id.clone());
|
||||
records.push(ReplacementRecoveryRecord::from_state(state).unwrap_or_else(|| {
|
||||
ReplacementRecoveryRecord::unknown(
|
||||
task_id,
|
||||
"isolated replacement intent violates its generation or target identity binding",
|
||||
)
|
||||
}));
|
||||
}
|
||||
|
||||
for entry in entries {
|
||||
let Some(task_id) = entry.strip_suffix(&proof_suffix) else {
|
||||
continue;
|
||||
};
|
||||
if validate_resume_task_id(task_id).is_err() {
|
||||
continue;
|
||||
}
|
||||
if intent_task_ids.contains(task_id) {
|
||||
continue;
|
||||
}
|
||||
let proof = ResumeManager::load_replacement_completion_proof(disk.clone(), task_id).await?;
|
||||
records.push(ReplacementRecoveryRecord::from_completion_proof(&proof));
|
||||
}
|
||||
|
||||
records.sort_by(|left, right| left.task_id.cmp(&right.task_id).then(left.state.cmp(&right.state)));
|
||||
Ok(records)
|
||||
}
|
||||
|
||||
/// cleanup expired resume states
|
||||
pub async fn cleanup_expired_states(disk: &DiskStore, max_age_hours: u64) -> Result<()> {
|
||||
let task_ids = Self::get_resumable_tasks(disk).await?;
|
||||
let current_time = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_secs();
|
||||
|
||||
for task_id in task_ids {
|
||||
if let Ok(resume_manager) = ResumeManager::load_from_disk(disk.clone(), &task_id).await {
|
||||
let state = resume_manager.get_state().await;
|
||||
let age_hours = current_time.saturating_sub(state.last_update) / 3600;
|
||||
|
||||
if !state.completed && matches!(state.replacement_phase, ReplacementPhase::Intent | ReplacementPhase::Rebuilding)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
if state.completed
|
||||
&& matches!(state.replacement_phase, ReplacementPhase::Verified | ReplacementPhase::CleanupPending)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
|
||||
if age_hours > max_age_hours {
|
||||
debug!(
|
||||
target: "rustfs::heal::resume",
|
||||
event = EVENT_HEAL_RESUME_STATE,
|
||||
component = LOG_COMPONENT_HEAL,
|
||||
subsystem = LOG_SUBSYSTEM_RESUME,
|
||||
task_id,
|
||||
age_hours,
|
||||
state = "expired_cleanup_started",
|
||||
"Heal resume cleanup started"
|
||||
);
|
||||
if let Err(e) = resume_manager.cleanup().await {
|
||||
warn!(
|
||||
target: "rustfs::heal::resume",
|
||||
event = EVENT_HEAL_RESUME_STATE,
|
||||
component = LOG_COMPONENT_HEAL,
|
||||
subsystem = LOG_SUBSYSTEM_RESUME,
|
||||
task_id,
|
||||
age_hours,
|
||||
state = "expired_cleanup_failed",
|
||||
error = %e,
|
||||
"Heal resume state cleanup failed"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
for task_id in Self::get_replacement_intent_tasks(disk).await? {
|
||||
if let Ok(resume_manager) = ResumeManager::load_replacement_intent(disk.clone(), &task_id).await {
|
||||
let state = resume_manager.get_state().await;
|
||||
let age_hours = current_time.saturating_sub(state.last_update) / 3600;
|
||||
|
||||
if !state.completed && matches!(state.replacement_phase, ReplacementPhase::Intent | ReplacementPhase::Rebuilding)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
if state.completed
|
||||
&& matches!(state.replacement_phase, ReplacementPhase::Verified | ReplacementPhase::CleanupPending)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
|
||||
if age_hours > max_age_hours
|
||||
&& let Err(e) = resume_manager.cleanup().await
|
||||
{
|
||||
warn!(
|
||||
target: "rustfs::heal::resume",
|
||||
event = EVENT_HEAL_RESUME_STATE,
|
||||
component = LOG_COMPONENT_HEAL,
|
||||
subsystem = LOG_SUBSYSTEM_RESUME,
|
||||
task_id,
|
||||
age_hours,
|
||||
state = "expired_cleanup_failed",
|
||||
error = %e,
|
||||
"Replacement intent cleanup failed"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user