Compare commits

..
6 Commits
Author SHA1 Message Date
Ultradesu 2fc5fd7960 Improved DHT publication. Adjust batching, added validation per key
Build and Publish / Build and Publish Docker Image (push) Successful in 4m9s
2026-07-23 16:44:35 +03:00
Ultradesu 63506e3af2 Aupdated FRID to v4 protocol
Build and Publish / Build and Publish Docker Image (push) Successful in 5m4s
2026-07-23 15:52:42 +03:00
Ultradesu 5b339aa921 Fixed content_id persistence
Build and Publish / Build and Publish Docker Image (push) Successful in 7m30s
2026-07-20 18:40:32 +03:00
Ultradesu 42c772f735 Fixed blake3 computation
Build and Publish / Build and Publish Docker Image (push) Successful in 5m5s
2026-07-20 18:24:39 +03:00
Ultradesu e738086573 Extend DHT scheme with content_id
Build and Publish / Build and Publish Docker Image (push) Successful in 5m6s
2026-07-20 18:05:31 +03:00
Ultradesu 4b7756c36e Extend DHT scheme with content_id 2026-07-20 18:05:16 +03:00
5 changed files with 315 additions and 185 deletions
Generated
+76 -119
View File
@@ -23,8 +23,8 @@ version = "0.5.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d122413f284cf2d62fb1b7db97e02edb8cda96d769b16e443a4f6195e35662b0"
dependencies = [
"crypto-common 0.1.7",
"generic-array 0.14.7",
"crypto-common 0.1.6",
"generic-array 0.14.9",
]
[[package]]
@@ -290,7 +290,7 @@ checksum = "ae36dc4177970ef04fde5178d3e2429882def40e57a451f919c098f72baa6cec"
dependencies = [
"proc-macro2",
"quote",
"syn 3.0.2",
"syn 3.0.3",
]
[[package]]
@@ -544,7 +544,7 @@ version = "0.10.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3078c7629b62d3f0439517fa394996acacc5cbc91c5a20d8c658e77abd503a71"
dependencies = [
"generic-array 0.14.7",
"generic-array 0.14.9",
]
[[package]]
@@ -684,15 +684,15 @@ version = "0.4.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "773f3b9af64447d2ce9850330c473515014aa235e6a783b02db81ff39e4a3dad"
dependencies = [
"crypto-common 0.1.7",
"crypto-common 0.1.6",
"inout",
]
[[package]]
name = "clap"
version = "4.6.2"
version = "4.6.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "dd059f9da4f5c36b3787f65d38ccaab1cc315f07b01f89abc8359ee6a8205011"
checksum = "d91e0c145792ef73a6ad36d27c75ac09f1832222a3c209689d90f534685ee5b7"
dependencies = [
"clap_builder",
"clap_derive",
@@ -712,14 +712,14 @@ dependencies = [
[[package]]
name = "clap_derive"
version = "4.6.1"
version = "4.6.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f2ce8604710f6733aa641a2b3731eaa1e8b3d9973d5e3565da11800813f997a9"
checksum = "d012d2b9d65aca7f18f4d9878a045bc17899bba951561ba5ec3c2ba1eed9a061"
dependencies = [
"heck",
"proc-macro2",
"quote",
"syn 2.0.119",
"syn 3.0.3",
]
[[package]]
@@ -1108,7 +1108,7 @@ version = "0.5.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0dc92fb57ca44df6db8059111ab3af99a63d5d0f8375d9972e319a379c6bab76"
dependencies = [
"generic-array 0.14.7",
"generic-array 0.14.9",
"rand_core 0.6.4",
"subtle",
"zeroize",
@@ -1116,11 +1116,11 @@ dependencies = [
[[package]]
name = "crypto-common"
version = "0.1.7"
version = "0.1.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "78c8292055d1c1df0cce5d180393dc8cce0abec0a7102adb6c7b1eef6016d60a"
checksum = "1bfb12502f3fc46cca1bb51ac28df9d618d813cdc3d2f25b9fe775a34af26bb3"
dependencies = [
"generic-array 0.14.7",
"generic-array 0.14.9",
"typenum",
]
@@ -1131,6 +1131,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ce6e4c961d6cd6c9a86db418387425e8bdeaf05b3c8bc1411e6dca4c252f1453"
dependencies = [
"hybrid-array",
"rand_core 0.10.1",
]
[[package]]
@@ -1181,9 +1182,9 @@ dependencies = [
[[package]]
name = "curve25519-dalek"
version = "5.0.0-rc.0"
version = "5.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "4f359e08ca85e7bd759e1fd933ff2bccd81864c60a8fba0e259c7f822b0924bf"
checksum = "b5eed333089e2e1c1ac8c6c0398e5e2497b4c9926ca6d0365ed1e099afa5bc23"
dependencies = [
"cfg-if",
"cpufeatures 0.3.0",
@@ -1417,7 +1418,7 @@ checksum = "9ed9a281f7bc9b7576e61468ba615a66a5c8cfdff42420a70aa82701a3b1e292"
dependencies = [
"block-buffer 0.10.4",
"const-oid 0.9.6",
"crypto-common 0.1.7",
"crypto-common 0.1.6",
"subtle",
]
@@ -1558,11 +1559,11 @@ dependencies = [
[[package]]
name = "ed25519-dalek"
version = "3.0.0-rc.0"
version = "3.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b011170fe4f04665565b4110afef66774fe9ffff278f3eb5b81cc73d26e27d60"
checksum = "6ebaa1a2bf1290ab3bfe5a7b771d050ebffab2711c19a81691c683a5144a25de"
dependencies = [
"curve25519-dalek 5.0.0-rc.0",
"curve25519-dalek 5.0.0",
"ed25519 3.0.0",
"rand_core 0.10.1",
"serde",
@@ -1591,7 +1592,7 @@ dependencies = [
"crypto-bigint",
"digest 0.10.7",
"ff",
"generic-array 0.14.7",
"generic-array 0.14.9",
"group",
"hkdf",
"pem-rfc7468 0.7.0",
@@ -1717,7 +1718,7 @@ dependencies = [
[[package]]
name = "federation-net"
version = "0.1.0"
source = "git+https://gt.hexor.cy/ab/frid.git#a897737978d476c223b66e65b960b70376084f2a"
source = "git+https://gt.hexor.cy/ab/frid.git#512a818a6a52ec713678e9a4e1cf0f50bb1e34ab"
dependencies = [
"blake3",
"data-encoding",
@@ -1844,11 +1845,12 @@ checksum = "e6d5a32815ae3f33302d95fdcb2ce17862f8c65363dcfd29360480ba1001fc9c"
[[package]]
name = "furumusic"
version = "0.6.2-fd"
version = "0.6.7-fd"
dependencies = [
"anyhow",
"async-trait",
"base64 0.22.1",
"blake3",
"chrono",
"cot",
"croner",
@@ -2032,9 +2034,9 @@ dependencies = [
[[package]]
name = "generic-array"
version = "0.14.7"
version = "0.14.9"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "85649ca51fd72272d7821adaf274ad91c288277713d9c18820d8499a7ff69e9a"
checksum = "4bb6743198531e02858aeaea5398fcc883e71851fcbcb5a2f773e2fb6cb1edf2"
dependencies = [
"typenum",
"version_check",
@@ -2470,9 +2472,9 @@ dependencies = [
[[package]]
name = "hyper"
version = "1.10.1"
version = "1.11.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "55281c53a1894c864990125767da440a4e630446785086f52523b20033b74498"
checksum = "d22053281f852e11534f5198498373cbb59295120a20771d90f7ed1897490a72"
dependencies = [
"atomic-waker",
"bytes",
@@ -2774,7 +2776,7 @@ checksum = "bee2c455ca60511a054699102d40ce7153621cb2429558f9eaa600e4499b5984"
dependencies = [
"proc-macro2",
"quote",
"syn 3.0.2",
"syn 3.0.3",
]
[[package]]
@@ -2783,7 +2785,7 @@ version = "0.1.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "879f10e63c20629ecabbb64a8010319738c66a5cd0c29b02d63d272b03751d01"
dependencies = [
"generic-array 0.14.7",
"generic-array 0.14.9",
]
[[package]]
@@ -2828,9 +2830,9 @@ dependencies = [
[[package]]
name = "iroh"
version = "1.0.2"
version = "1.0.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5fca9b4b462c343ff88fc0af4096c186f939b602a0bc08723536ef2c31c93971"
checksum = "460de6bc52163b41b1646931f2897e5ab986f0966ade444467fec25024751a72"
dependencies = [
"backon",
"blake3",
@@ -2839,7 +2841,7 @@ dependencies = [
"ctutils",
"data-encoding",
"derive_more",
"ed25519-dalek 3.0.0-rc.0",
"ed25519-dalek 3.0.0",
"futures-util",
"getrandom 0.4.3",
"hickory-resolver",
@@ -2879,15 +2881,15 @@ dependencies = [
[[package]]
name = "iroh-base"
version = "1.0.2"
version = "1.0.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "830a582cd54410dc1aa71d4786a82c3297d7b0165accd8b6dbbb3b240b48140d"
checksum = "6be73e16ee21c923aca9b3121aaa0db936f7c7ecc156ff47b8dac944c68d59a8"
dependencies = [
"curve25519-dalek 5.0.0-rc.0",
"curve25519-dalek 5.0.0",
"data-encoding",
"data-encoding-macro",
"derive_more",
"ed25519-dalek 3.0.0-rc.0",
"ed25519-dalek 3.0.0",
"getrandom 0.4.3",
"n0-error",
"rand 0.10.2",
@@ -2898,9 +2900,9 @@ dependencies = [
[[package]]
name = "iroh-dns"
version = "1.0.2"
version = "1.0.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "516e4eedc38e33ab69a6bd325520332dc3d67b25454e2d590ebb84a25240dd9a"
checksum = "46f6a9b39d18e6345f5c151afd299f2488e2cb5c520fe41b107b6bd3dc4c3349"
dependencies = [
"arc-swap",
"cfg_aliases",
@@ -2949,9 +2951,9 @@ dependencies = [
[[package]]
name = "iroh-relay"
version = "1.0.2"
version = "1.0.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8149bb6a57126225a07d6928846d82dcedfd24ea0f863ef7b2eb475e1d726354"
checksum = "24bd586cf927f7b700f56ec3639b53cb5fa901ce284784051ff71092bfbf8193"
dependencies = [
"blake3",
"bytes",
@@ -2988,7 +2990,6 @@ dependencies = [
"tokio-websockets",
"tracing",
"url",
"vergen-gitcl",
"webpki-roots 1.0.9",
"ws_stream_wasm",
]
@@ -3153,9 +3154,9 @@ dependencies = [
[[package]]
name = "libc"
version = "0.2.186"
version = "0.2.189"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "68ab91017fe16c622486840e4c83c9a37afeff978bd239b5293d61ece587de66"
checksum = "3eaf3ede3fee6db1a4c2ee091bf8a8b4dccdc6d17f656fb07896ee72867612f2"
[[package]]
name = "libm"
@@ -3484,7 +3485,7 @@ dependencies = [
"crc",
"document-features",
"dyn-clone",
"ed25519-dalek 3.0.0-rc.0",
"ed25519-dalek 3.0.0",
"flume 0.12.0",
"futures-lite",
"getrandom 0.4.3",
@@ -3621,7 +3622,7 @@ dependencies = [
[[package]]
name = "music-dht"
version = "0.1.0"
source = "git+https://gt.hexor.cy/ab/frid.git#a897737978d476c223b66e65b960b70376084f2a"
source = "git+https://gt.hexor.cy/ab/frid.git#512a818a6a52ec713678e9a4e1cf0f50bb1e34ab"
dependencies = [
"async-trait",
"blake3",
@@ -3846,9 +3847,9 @@ checksum = "38bf9645c8b145698bb0b18a4637dcacbc421ea49bef2317e4fd8065a387cf21"
[[package]]
name = "noq"
version = "1.0.1"
version = "1.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "4bf95190af1bd4a00a10e8255ca0c8ddd9e9a9f5e79151d7a7eb6d56aff5dc89"
checksum = "e11803df44ac03a30988d61585ea50885d5428e42da944fe1e498799da7886a2"
dependencies = [
"bytes",
"cfg_aliases",
@@ -3868,9 +3869,9 @@ dependencies = [
[[package]]
name = "noq-proto"
version = "1.0.1"
version = "1.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "aa6c890013591e709a3e45dd53501351b7e27e7ff3c7e9fc3dce43e300e7e9d3"
checksum = "334c3c9833f7b2c573cceb9896ddc7aaeb58c8807cbb63211b24d1fe88bf866e"
dependencies = [
"aes-gcm",
"bytes",
@@ -3895,9 +3896,9 @@ dependencies = [
[[package]]
name = "noq-udp"
version = "1.0.1"
version = "1.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3137a52df66c20090a889828d1c655f21f52294cba64e5c4fbb04fc83eee7c8e"
checksum = "bde7a5d5102f1cff03d482240f0ed20551661f63663620f4b26112ed751165e9"
dependencies = [
"cfg_aliases",
"libc",
@@ -4033,15 +4034,6 @@ dependencies = [
"syn 2.0.119",
]
[[package]]
name = "num_threads"
version = "0.1.7"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5c7398b9c8b70908f6371f47ed36737907c87c52af34c268fed0bf0ceb92ead9"
dependencies = [
"libc",
]
[[package]]
name = "oauth2"
version = "5.0.0"
@@ -4959,7 +4951,7 @@ checksum = "2c9283685feec7d69af75fb0e858d5e7378f33fe4fc699383b2916ab9273e03c"
dependencies = [
"proc-macro2",
"quote",
"syn 3.0.2",
"syn 3.0.3",
]
[[package]]
@@ -5208,9 +5200,9 @@ dependencies = [
[[package]]
name = "rustls-pki-types"
version = "1.15.0"
version = "1.15.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "764899a24af3980067ee14bc143654f297b22eaebfe3c7b6b211920a5a59b046"
checksum = "2f4925028c7eb5d1fcdaf196971378ed9d2c1c4efc7dc5d011256f76c99c0a96"
dependencies = [
"web-time",
"zeroize",
@@ -5363,7 +5355,7 @@ checksum = "d3e97a565f76233a6003f9f5c54be1d9c5bdfa3eccfb189469f11ec4901c47dc"
dependencies = [
"base16ct 0.2.0",
"der 0.7.10",
"generic-array 0.14.7",
"generic-array 0.14.9",
"pkcs8 0.10.2",
"subtle",
"zeroize",
@@ -5471,7 +5463,7 @@ checksum = "e7a5d71263a5a7d47b41f6b3f06ba276f10cc18b0931f1799f710578e2309348"
dependencies = [
"proc-macro2",
"quote",
"syn 3.0.2",
"syn 3.0.3",
]
[[package]]
@@ -5671,6 +5663,9 @@ name = "signature"
version = "3.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "28d567dcbaf0049cb8ac2608a76cd95ff9e4412e1899d389ee400918ca7537f5"
dependencies = [
"rand_core 0.10.1",
]
[[package]]
name = "simd-adler32"
@@ -5913,7 +5908,7 @@ dependencies = [
"futures-core",
"futures-io",
"futures-util",
"generic-array 0.14.7",
"generic-array 0.14.9",
"hex 0.4.3",
"hkdf",
"hmac",
@@ -6276,9 +6271,9 @@ dependencies = [
[[package]]
name = "syn"
version = "3.0.2"
version = "3.0.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a207d6d6a2b7fc470b80443726053f18a2481b7e1eee970597051596567987a3"
checksum = "53e9bae58849f64dfa4f5d5ae372c8341f7305f82a3868709269343628b659a3"
dependencies = [
"proc-macro2",
"quote",
@@ -6388,7 +6383,7 @@ checksum = "43cbfe0cf76104d42a574802844187e84a305e531ed54455f11fbde0f10541cd"
dependencies = [
"proc-macro2",
"quote",
"syn 3.0.2",
"syn 3.0.3",
]
[[package]]
@@ -6408,9 +6403,7 @@ checksum = "3e1d5e639ff6bab73cb6885cc7e7b1de96c3f32c68ec55f3952614bec1092244"
dependencies = [
"deranged",
"js-sys",
"libc",
"num-conv",
"num_threads",
"powerfmt",
"serde_core",
"time-core",
@@ -6460,9 +6453,9 @@ checksum = "1f3ccbac311fea05f86f61904b462b55fb3df8837a366dfc601a0161d0532f20"
[[package]]
name = "tokio"
version = "1.53.0"
version = "1.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d988bcd52dbe076d3d46903332f58c912b87a2c49b1428419a5845154762ffee"
checksum = "202caea871b69668250d242070849eb495be178ed697a3e98aebce5bc81a0bed"
dependencies = [
"bytes",
"libc",
@@ -6535,9 +6528,9 @@ dependencies = [
[[package]]
name = "tokio-stream"
version = "0.1.18"
version = "0.1.19"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "32da49809aab5c3bc678af03902d4ccddea2a87d028d86392a4b1560c6906c70"
checksum = "a3d06f0b082ba57c26b79407372e57cf2a1e28124f78e9479fe80322cf53420b"
dependencies = [
"futures-core",
"pin-project-lite",
@@ -6547,14 +6540,15 @@ dependencies = [
[[package]]
name = "tokio-util"
version = "0.7.18"
version = "0.7.19"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9ae9cec805b01e8fc3fd2fe289f89149a9b66dd16786abd8b19cfa7b48cb0098"
checksum = "494815d09bf52b5548659851081238f0ca39ff638363907596da739561c62c52"
dependencies = [
"bytes",
"futures-core",
"futures-sink",
"futures-util",
"libc",
"pin-project-lite",
"tokio",
]
@@ -6867,7 +6861,7 @@ version = "0.5.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "fc1de2c688dc15305988b563c3854064043356019f97a4b46276fe734c4f07ea"
dependencies = [
"crypto-common 0.1.7",
"crypto-common 0.1.6",
"subtle",
]
@@ -6937,43 +6931,6 @@ version = "0.2.15"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "accd4ea62f7bb7a82fe23066fb0957d48ef677f6eeb8215f372f52e48bb32426"
[[package]]
name = "vergen"
version = "9.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b849a1f6d8639e8de261e81ee0fc881e3e3620db1af9f2e0da015d4382ceaf75"
dependencies = [
"anyhow",
"derive_builder",
"rustversion",
"vergen-lib",
]
[[package]]
name = "vergen-gitcl"
version = "9.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "77ff3b5300a085d6bcd8fc96a507f706a28ae3814693236c9b409db71a1d15b9"
dependencies = [
"anyhow",
"derive_builder",
"rustversion",
"time",
"vergen",
"vergen-lib",
]
[[package]]
name = "vergen-lib"
version = "9.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b34a29ba7e9c59e62f229ae1932fb1b8fb8a6fdcc99215a641913f5f5a59a569"
dependencies = [
"anyhow",
"derive_builder",
"rustversion",
]
[[package]]
name = "version_check"
version = "0.9.5"
@@ -7646,18 +7603,18 @@ dependencies = [
[[package]]
name = "zerocopy"
version = "0.8.54"
version = "0.8.55"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b7cbbc0a705a0fd05cc3676525980d2bf5a9bc4adac6d6475209a7887cf59d19"
checksum = "b5a105cd7b140f6eeec8acff2ea38135d3cab283ada58540f629fe51e46696eb"
dependencies = [
"zerocopy-derive",
]
[[package]]
name = "zerocopy-derive"
version = "0.8.54"
version = "0.8.55"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e2e817b7b52d0c7358d3246da9d69935ebb18116b2b102b4230dac079b4862f5"
checksum = "0fe976fb70c78cd64cccfe3a6fc142244e8a77b70959b30faf9d0ac37ee228eb"
dependencies = [
"proc-macro2",
"quote",
+2 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "furumusic"
version = "0.6.3-fd"
version = "0.7.0"
edition = "2024"
description = "Reusable web-app boilerplate: auth, OIDC/SSO, admin panel, user management, i18n, PostgreSQL"
@@ -16,6 +16,7 @@ reqwest = { version = "0.12", default-features = false, features = ["rustls-tls"
tokio = { version = "1", features = ["sync", "fs", "io-util"] }
tower = "0.5"
base64 = "0.22"
blake3 = "1"
serde_json = "1"
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
+136 -3
View File
@@ -15,6 +15,7 @@
mod serve;
mod storage;
use std::collections::{HashMap, HashSet};
use std::path::PathBuf;
use std::sync::{Arc, OnceLock};
use std::time::Duration;
@@ -42,10 +43,19 @@ struct Running {
tasks: Vec<tokio::task::JoinHandle<()>>,
}
struct ContentHashJob {
media_file_id: i64,
sha256_hash: String,
file_path: String,
}
pub struct Federation {
/// Transport data directory; server-side DHT state and identity live in PostgreSQL.
data_dir: PathBuf,
database_url: std::sync::Mutex<String>,
storage_dir: std::sync::Mutex<String>,
content_cache: std::sync::Mutex<HashMap<i64, (String, String)>>,
content_pending: std::sync::Mutex<HashSet<i64>>,
pool: tokio::sync::OnceCell<PgPool>,
running: tokio::sync::Mutex<Option<Running>>,
last_sync: std::sync::Mutex<Option<String>>,
@@ -69,6 +79,9 @@ pub fn handle() -> Arc<Federation> {
Arc::new(Federation {
data_dir: PathBuf::from(crate::media_paths::resolve_config_path("federation")),
database_url: std::sync::Mutex::new(String::new()),
storage_dir: std::sync::Mutex::new(String::new()),
content_cache: std::sync::Mutex::new(Default::default()),
content_pending: std::sync::Mutex::new(Default::default()),
pool: tokio::sync::OnceCell::new(),
running: tokio::sync::Mutex::new(None),
last_sync: std::sync::Mutex::new(None),
@@ -149,6 +162,7 @@ impl Federation {
/// node. Called at boot and every time the admin settings are saved.
pub async fn apply(self: &Arc<Self>, config: &AppConfig) {
*lock(&self.database_url) = config.database_url.clone();
*lock(&self.storage_dir) = config.agent_storage_dir.clone();
let network = config.federation_network_id.trim().to_string();
if config.federation_enabled && !network.is_empty() {
if let Err(err) = self.start(network, config.agent_storage_dir.clone()).await {
@@ -281,7 +295,7 @@ impl Federation {
Ok(())
}
async fn sync_once(&self, service: &MusicDhtService) -> Result<SyncStats> {
async fn sync_once(self: &Arc<Self>, service: &MusicDhtService) -> Result<SyncStats> {
let specs = match self.collect_specs().await {
Ok(specs) => specs,
Err(err) => {
@@ -341,7 +355,7 @@ impl Federation {
/// Everything the regular player shows, as DHT item specs: non-hidden
/// artists, releases and tracks (a track also hides with its release).
async fn collect_specs(&self) -> Result<Vec<ItemSpec>> {
async fn collect_specs(self: &Arc<Self>) -> Result<Vec<ItemSpec>> {
let pool = self.pool().await?;
let mut specs = Vec::new();
@@ -362,6 +376,7 @@ impl Federation {
track_number: None,
disc_number: None,
duration_seconds: None,
content_id: None,
});
}
@@ -400,6 +415,7 @@ impl Federation {
track_number: None,
disc_number: None,
duration_seconds: None,
content_id: None,
});
}
@@ -426,16 +442,38 @@ impl Federation {
}
let tracks = sqlx::query(
"SELECT t.id, t.title, COALESCE(t.year, r.year), t.duration_seconds,
r.title, r.release_type, t.track_number, t.disc_number
r.title, r.release_type, t.track_number, t.disc_number,
t.audio_file_id, m.file_path, m.sha256_hash, c.content_id
FROM furumusic__track t
JOIN furumusic__release r ON r.id = t.release_id
JOIN furumusic__media_file m ON m.id = t.audio_file_id
LEFT JOIN furumusic__federation_content_id_cache c
ON c.media_file_id = m.id AND c.sha256_hash = m.sha256_hash
WHERE t.is_hidden = false AND r.is_hidden = false",
)
.fetch_all(&pool)
.await?;
let storage_dir = lock(&self.storage_dir).clone();
let mut content_hash_jobs = Vec::new();
for row in &tracks {
let id: i64 = row.get(0);
let duration: f64 = row.get(3);
let media_file_id: i64 = row.get(8);
let file_path: String = row.get(9);
let sha256_hash: String = row.get(10);
let cached_content_id: Option<String> = row.get(11);
let content_id = cached_content_id
.or_else(|| self.cached_content_id_for_media(media_file_id, &sha256_hash));
if content_id.is_none()
&& !storage_dir.trim().is_empty()
&& self.mark_content_hash_pending(media_file_id)
{
content_hash_jobs.push(ContentHashJob {
media_file_id,
sha256_hash,
file_path,
});
}
specs.push(ItemSpec {
local_key: format!("track:{id}"),
kind: ItemKind::Track,
@@ -448,12 +486,72 @@ impl Federation {
track_number: row.get(6),
disc_number: row.get(7),
duration_seconds: (duration > 0.0).then_some(duration),
content_id,
});
}
self.spawn_content_warmer(pool.clone(), storage_dir, content_hash_jobs);
Ok(specs)
}
fn cached_content_id_for_media(&self, media_file_id: i64, sha256_hash: &str) -> Option<String> {
if let Some((cached_hash, content_id)) = lock(&self.content_cache).get(&media_file_id)
&& cached_hash == sha256_hash
{
return Some(content_id.clone());
}
None
}
fn mark_content_hash_pending(&self, media_file_id: i64) -> bool {
lock(&self.content_pending).insert(media_file_id)
}
fn spawn_content_warmer(
self: &Arc<Self>,
pool: PgPool,
storage_dir: String,
jobs: Vec<ContentHashJob>,
) {
if jobs.is_empty() {
return;
}
let fed = Arc::clone(self);
tokio::spawn(async move {
let total = jobs.len();
let mut stored = 0usize;
for job in jobs {
let job_storage_dir = storage_dir.clone();
let job_file_path = job.file_path.clone();
let content_id = tokio::task::spawn_blocking(move || {
audio_content_id(&job_storage_dir, &job_file_path)
})
.await
.ok()
.flatten();
lock(&fed.content_pending).remove(&job.media_file_id);
if let Some(content_id) = content_id {
lock(&fed.content_cache).insert(
job.media_file_id,
(job.sha256_hash.clone(), content_id.clone()),
);
if let Err(err) =
persist_content_id(&pool, job.media_file_id, &job.sha256_hash, &content_id)
.await
{
tracing::warn!(
media_file_id = job.media_file_id,
"federation content-id cache write failed: {err:#}"
);
} else {
stored += 1;
}
}
}
tracing::info!(total, stored, "federation content-id cache warm finished");
});
}
/// Live status for the admin page.
pub async fn status(&self) -> Value {
let guard = self.running.lock().await;
@@ -511,6 +609,41 @@ impl Federation {
}
}
async fn persist_content_id(
pool: &PgPool,
media_file_id: i64,
sha256_hash: &str,
content_id: &str,
) -> Result<()> {
sqlx::query(
"INSERT INTO furumusic__federation_content_id_cache
(media_file_id, sha256_hash, content_id, updated_at)
VALUES ($1, $2, $3, $4)
ON CONFLICT (media_file_id) DO UPDATE SET
sha256_hash = EXCLUDED.sha256_hash,
content_id = EXCLUDED.content_id,
updated_at = EXCLUDED.updated_at",
)
.bind(media_file_id)
.bind(sha256_hash)
.bind(content_id)
.bind(now_iso())
.execute(pool)
.await?;
Ok(())
}
fn audio_content_id(storage_dir: &str, file_path: &str) -> Option<String> {
if storage_dir.trim().is_empty() {
return None;
}
let path = crate::media_paths::resolve_media_file_path(storage_dir, file_path);
let mut file = std::fs::File::open(path).ok()?;
let mut hasher = blake3::Hasher::new();
std::io::copy(&mut file, &mut hasher).ok()?;
Some(format!("b3:{}", hasher.finalize().to_hex()))
}
async fn stop_running(running: Option<Running>) {
let Some(running) = running else { return };
for task in &running.tasks {
+24 -18
View File
@@ -101,6 +101,8 @@ struct CatalogTrack {
track_number: Option<i32>,
disc_number: Option<i32>,
duration_seconds: Option<f64>,
#[serde(skip_serializing_if = "Option::is_none")]
content_id: Option<String>,
item_id: String,
}
@@ -612,32 +614,36 @@ async fn build_catalog(pool: &PgPool, own: &EndpointId, artist: &str) -> Result<
for release_row in release_rows {
let release_id: i64 = release_row.get(0);
let track_rows = sqlx::query(
"SELECT id, title, track_number, disc_number, duration_seconds
FROM furumusic__track
WHERE release_id = $1 AND is_hidden = false
ORDER BY disc_number NULLS FIRST, track_number NULLS LAST, title",
"SELECT t.id, t.title, t.track_number, t.disc_number, t.duration_seconds,
c.content_id
FROM furumusic__track t
JOIN furumusic__media_file m ON m.id = t.audio_file_id
LEFT JOIN furumusic__federation_content_id_cache c
ON c.media_file_id = m.id AND c.sha256_hash = m.sha256_hash
WHERE t.release_id = $1 AND t.is_hidden = false
ORDER BY t.disc_number NULLS FIRST, t.track_number NULLS LAST, t.title",
)
.bind(release_id)
.fetch_all(pool)
.await?;
let mut tracks = Vec::with_capacity(track_rows.len());
for row in track_rows {
let track_id: i64 = row.get(0);
let duration: f64 = row.get(4);
tracks.push(CatalogTrack {
title: row.get(1),
track_number: row.get(2),
disc_number: row.get(3),
duration_seconds: (duration > 0.0).then_some(duration),
content_id: row.get(5),
item_id: item_id_of(own, track_id),
});
}
releases.push(CatalogRelease {
title: release_row.get(1),
release_type: release_row.get(2),
year: release_row.get(3),
tracks: track_rows
.into_iter()
.map(|row| {
let track_id: i64 = row.get(0);
let duration: f64 = row.get(4);
CatalogTrack {
title: row.get(1),
track_number: row.get(2),
disc_number: row.get(3),
duration_seconds: (duration > 0.0).then_some(duration),
item_id: item_id_of(own, track_id),
}
})
.collect(),
tracks,
});
}
+77 -44
View File
@@ -44,6 +44,14 @@ const SCHEMA: &[&str] = &[
ticket TEXT NOT NULL,
last_seen_ms BIGINT NOT NULL
)",
"CREATE TABLE IF NOT EXISTS furumusic__federation_content_id_cache (
media_file_id BIGINT PRIMARY KEY,
sha256_hash TEXT NOT NULL,
content_id TEXT NOT NULL,
updated_at TEXT NOT NULL
)",
"CREATE INDEX IF NOT EXISTS idx_furumusic_federation_content_id_cache_content_id
ON furumusic__federation_content_id_cache(content_id)",
];
#[derive(Debug, Clone)]
@@ -157,52 +165,21 @@ impl MusicDhtStorage for PostgresFederationStorage {
}
async fn store_dht_record(&self, key: DhtKey, record: StoredRecord) -> music_dht::Result<bool> {
let existing = sqlx::query(
"SELECT revision, deleted, expires_at_ms
FROM furumusic__federation_dht_record
WHERE dht_key = $1 AND item_id = $2 AND owner_peer_id = $3",
)
.bind(key.as_bytes().as_slice())
.bind(record.item.id.as_bytes().as_slice())
.bind(record.item.owner.to_string())
.fetch_optional(&self.pool)
.await
.map_err(db_error)?
.map(|row| {
(
row.get::<i64, _>(0) as u64,
row.get::<bool, _>(1),
row.get::<i64, _>(2) as u64,
)
});
let mut conn = self.pool.acquire().await.map_err(db_error)?;
store_record_in_conn(&mut conn, &key, &record).await
}
match decide_store(existing, &record) {
StoreDecision::Ignore => return Ok(false),
StoreDecision::Write | StoreDecision::RefreshExpiry(_) => {}
async fn store_dht_records(
&self,
entries: Vec<(DhtKey, StoredRecord)>,
) -> music_dht::Result<Vec<bool>> {
let mut tx = self.pool.begin().await.map_err(db_error)?;
let mut stored = Vec::with_capacity(entries.len());
for (key, record) in &entries {
stored.push(store_record_in_conn(&mut tx, key, record).await?);
}
let payload = postcard::to_stdvec(&record).map_err(db_error)?;
sqlx::query(
"INSERT INTO furumusic__federation_dht_record
(dht_key, item_id, owner_peer_id, payload, revision, deleted, expires_at_ms)
VALUES ($1, $2, $3, $4, $5, $6, $7)
ON CONFLICT (dht_key, item_id, owner_peer_id) DO UPDATE SET
payload = EXCLUDED.payload,
revision = EXCLUDED.revision,
deleted = EXCLUDED.deleted,
expires_at_ms = EXCLUDED.expires_at_ms",
)
.bind(key.as_bytes().as_slice())
.bind(record.item.id.as_bytes().as_slice())
.bind(record.item.owner.to_string())
.bind(payload)
.bind(record.item.revision as i64)
.bind(record.item.deleted)
.bind(record.expires_at_ms as i64)
.execute(&self.pool)
.await
.map_err(db_error)?;
Ok(true)
tx.commit().await.map_err(db_error)?;
Ok(stored)
}
async fn dht_records_by_key(
@@ -311,6 +288,62 @@ fn secret_from_bytes(bytes: Vec<u8>) -> music_dht::Result<SecretKey> {
Ok(SecretKey::from_bytes(&bytes))
}
/// Applies one validated record following the revision/tombstone rules.
/// Returns `true` if the record was written or refreshed. Runs against a
/// pooled connection or an open transaction.
async fn store_record_in_conn(
conn: &mut sqlx::PgConnection,
key: &DhtKey,
record: &StoredRecord,
) -> music_dht::Result<bool> {
let existing = sqlx::query(
"SELECT revision, deleted, expires_at_ms
FROM furumusic__federation_dht_record
WHERE dht_key = $1 AND item_id = $2 AND owner_peer_id = $3",
)
.bind(key.as_bytes().as_slice())
.bind(record.item.id.as_bytes().as_slice())
.bind(record.item.owner.to_string())
.fetch_optional(&mut *conn)
.await
.map_err(db_error)?
.map(|row| {
(
row.get::<i64, _>(0) as u64,
row.get::<bool, _>(1),
row.get::<i64, _>(2) as u64,
)
});
match decide_store(existing, record) {
StoreDecision::Ignore => return Ok(false),
StoreDecision::Write | StoreDecision::RefreshExpiry(_) => {}
}
let payload = postcard::to_stdvec(record).map_err(db_error)?;
sqlx::query(
"INSERT INTO furumusic__federation_dht_record
(dht_key, item_id, owner_peer_id, payload, revision, deleted, expires_at_ms)
VALUES ($1, $2, $3, $4, $5, $6, $7)
ON CONFLICT (dht_key, item_id, owner_peer_id) DO UPDATE SET
payload = EXCLUDED.payload,
revision = EXCLUDED.revision,
deleted = EXCLUDED.deleted,
expires_at_ms = EXCLUDED.expires_at_ms",
)
.bind(key.as_bytes().as_slice())
.bind(record.item.id.as_bytes().as_slice())
.bind(record.item.owner.to_string())
.bind(payload)
.bind(record.item.revision as i64)
.bind(record.item.deleted)
.bind(record.expires_at_ms as i64)
.execute(&mut *conn)
.await
.map_err(db_error)?;
Ok(true)
}
fn db_error(err: impl std::fmt::Display) -> MusicDhtError {
MusicDhtError::Database(err.to_string())
}