mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-08 15:45:50 +00:00
Compare commits
7
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1df165d514 | ||
|
|
0305e837fd | ||
|
|
8f80dac30f | ||
|
|
30cf53265b | ||
|
|
4d1f49c638 | ||
|
|
14fdd61aea | ||
|
|
49c25890ad |
@@ -10,3 +10,5 @@ updates:
|
||||
directory: "/"
|
||||
schedule:
|
||||
interval: weekly
|
||||
ignore:
|
||||
- dependency-name: "github.com/seaweedfs/goexif"
|
||||
|
||||
@@ -64,7 +64,7 @@ require (
|
||||
github.com/prometheus/procfs v0.22.0
|
||||
github.com/rcrowley/go-metrics v0.0.0-20201227073835-cf1acfcdf475 // indirect
|
||||
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect
|
||||
github.com/seaweedfs/goexif v2.0.0+incompatible
|
||||
github.com/seaweedfs/goexif v1.0.3
|
||||
github.com/seaweedfs/raft v1.2.1
|
||||
github.com/sirupsen/logrus v1.9.4 // indirect
|
||||
github.com/spf13/afero v1.15.0 // indirect
|
||||
@@ -258,7 +258,6 @@ require (
|
||||
github.com/rclone/Proton-API-Bridge v1.0.5 // indirect
|
||||
github.com/rclone/go-proton-api v1.0.4 // indirect
|
||||
github.com/rogpeppe/go-internal v1.15.0 // indirect
|
||||
github.com/rwcarlsen/goexif v0.0.0-20190401172101-9e8deecbddbd // indirect
|
||||
github.com/ryanuber/go-glob v1.0.0 // indirect
|
||||
github.com/sasha-s/go-deadlock v0.3.1 // indirect
|
||||
github.com/smarty/assertions v1.15.0 // indirect
|
||||
|
||||
@@ -1746,8 +1746,6 @@ github.com/rs/zerolog v1.34.0 h1:k43nTLIwcTVQAncfCw4KZ2VY6ukYoZaBPNOE8txlOeY=
|
||||
github.com/rs/zerolog v1.34.0/go.mod h1:bJsvje4Z08ROH4Nhs5iH600c3IkWhwp44iRc54W6wYQ=
|
||||
github.com/ruudk/golang-pdf417 v0.0.0-20181029194003-1af4ab5afa58/go.mod h1:6lfFZQK844Gfx8o5WFuvpxWRwnSoipWe/p622j1v06w=
|
||||
github.com/ruudk/golang-pdf417 v0.0.0-20201230142125-a7e3863a1245/go.mod h1:pQAZKsJ8yyVxGRWYNEm9oFB8ieLgKFnamEyDmSA0BRk=
|
||||
github.com/rwcarlsen/goexif v0.0.0-20190401172101-9e8deecbddbd h1:CmH9+J6ZSsIjUK3dcGsnCnO41eRBOnY12zwkn5qVwgc=
|
||||
github.com/rwcarlsen/goexif v0.0.0-20190401172101-9e8deecbddbd/go.mod h1:hPqNNc0+uJM6H+SuU8sEs5K5IQeKccPqeSjfgcKGgPk=
|
||||
github.com/ryanuber/go-glob v1.0.0 h1:iQh3xXAumdQ+4Ufa5b25cRpC5TYKlno6hsv6Cb3pkBk=
|
||||
github.com/ryanuber/go-glob v1.0.0/go.mod h1:807d1WSdnB0XRJzKNil9Om6lcp/3a0v4qIHxIXzX/Yc=
|
||||
github.com/sabhiram/go-gitignore v0.0.0-20210923224102-525f6e181f06 h1:OkMGxebDjyw0ULyrTYWeN0UNCCkmCWfjPnIA2W6oviI=
|
||||
@@ -1766,8 +1764,8 @@ github.com/seaweedfs/cockroachdb-parser v0.0.0-20260225204133-2f342c5ea564 h1:Tg
|
||||
github.com/seaweedfs/cockroachdb-parser v0.0.0-20260225204133-2f342c5ea564/go.mod h1:JSKCh6uCHBz91lQYFYHCyTrSVIPge4SUFVn28iwMNB0=
|
||||
github.com/seaweedfs/go-fuse/v2 v2.9.4 h1:ACyloiuopdhRSjdLLeSWbsVaemMPskORaRF01TY6GyM=
|
||||
github.com/seaweedfs/go-fuse/v2 v2.9.4/go.mod h1:zABdmWEa6A0bwaBeEOBUeUkGIZlxUhcdv+V1Dcc/U/I=
|
||||
github.com/seaweedfs/goexif v2.0.0+incompatible h1:x8pckiT12QQhifwhDQpeISgDfsqmQ6VR4LFPQ64JRps=
|
||||
github.com/seaweedfs/goexif v2.0.0+incompatible/go.mod h1:Oni780Z236sXpIQzk1XoJlTwqrJ02smEin9zQeff7Fk=
|
||||
github.com/seaweedfs/goexif v1.0.3 h1:ve/OjI7dxPW8X9YQsv3JuVMaxEyF9Rvfd04ouL+Bz30=
|
||||
github.com/seaweedfs/goexif v1.0.3/go.mod h1:Oni780Z236sXpIQzk1XoJlTwqrJ02smEin9zQeff7Fk=
|
||||
github.com/seaweedfs/raft v1.2.1 h1:QgFl/aaPnagpUxYB6Bx+fFss1NyetVVcJmraMmnaQ5Q=
|
||||
github.com/seaweedfs/raft v1.2.1/go.mod h1:fgs/rAVEzjQ7e04XMzG3eJhwZZRmBW+2uRtjakeCGeU=
|
||||
github.com/secure-systems-lab/go-securesystemslib v0.11.0 h1:iuCR9kcMFD4QurdKrGvPLoKZLv9YvwPYVr0473BdtFs=
|
||||
|
||||
+1042
-1049
File diff suppressed because it is too large
Load Diff
|
Before Width: | Height: | Size: 54 KiB After Width: | Height: | Size: 54 KiB |
@@ -1182,6 +1182,10 @@ pub struct Volume {
|
||||
|
||||
last_modified_ts_seconds: u64,
|
||||
last_append_at_ns: u64,
|
||||
last_write_append_at_ns: u64, // AppendAtNs of the newest write; tombstones don't move it
|
||||
last_write_needle_key: NeedleId, // the write behind the watermark
|
||||
last_write_deleted: bool, // that write was deleted, so recovery has to rescan
|
||||
keep_last_modified_ts_on_load: bool,
|
||||
pub last_disk_check_ns: Arc<std::sync::atomic::AtomicI64>, // for phantom volume detection cache
|
||||
|
||||
last_compact_index_offset: u64,
|
||||
@@ -1281,6 +1285,10 @@ impl Volume {
|
||||
location_disk_space_low: Arc::new(AtomicBool::new(false)),
|
||||
last_modified_ts_seconds: 0,
|
||||
last_append_at_ns: 0,
|
||||
last_write_append_at_ns: 0,
|
||||
last_write_needle_key: NeedleId(0),
|
||||
last_write_deleted: false,
|
||||
keep_last_modified_ts_on_load: false,
|
||||
last_disk_check_ns: Arc::new(std::sync::atomic::AtomicI64::new(0)),
|
||||
last_compact_index_offset: 0,
|
||||
last_compact_revision: 0,
|
||||
@@ -1327,6 +1335,10 @@ impl Volume {
|
||||
location_disk_space_low: Arc::new(AtomicBool::new(false)),
|
||||
last_modified_ts_seconds: 0,
|
||||
last_append_at_ns: 0,
|
||||
last_write_append_at_ns: 0,
|
||||
last_write_needle_key: NeedleId(0),
|
||||
last_write_deleted: false,
|
||||
keep_last_modified_ts_on_load: false,
|
||||
last_disk_check_ns: Arc::new(std::sync::atomic::AtomicI64::new(0)),
|
||||
last_compact_index_offset: 0,
|
||||
last_compact_revision: 0,
|
||||
@@ -1442,12 +1454,14 @@ impl Volume {
|
||||
Err(e) => return Err(e.into()),
|
||||
}
|
||||
|
||||
self.last_modified_ts_seconds = metadata
|
||||
.modified()
|
||||
.unwrap_or(SystemTime::UNIX_EPOCH)
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_secs();
|
||||
if !self.keep_last_modified_ts_on_load {
|
||||
self.last_modified_ts_seconds = metadata
|
||||
.modified()
|
||||
.unwrap_or(SystemTime::UNIX_EPOCH)
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_secs();
|
||||
}
|
||||
|
||||
if metadata.len() >= SUPER_BLOCK_SIZE as u64 {
|
||||
already_has_super_block = true;
|
||||
@@ -1549,7 +1563,9 @@ impl Volume {
|
||||
"volumeDataIntegrityChecking failed"
|
||||
);
|
||||
}
|
||||
self.recover_last_modified_ts();
|
||||
if !self.keep_last_modified_ts_on_load {
|
||||
self.recover_last_modified_ts();
|
||||
}
|
||||
|
||||
// Structural check: no .idx entry may reference bytes past the
|
||||
// end of .dat. The needle map's load walk above already
|
||||
@@ -2339,6 +2355,13 @@ impl Volume {
|
||||
.collect();
|
||||
}
|
||||
self.last_append_at_ns = last_append_at_ns;
|
||||
for ((n, _), r) in run.iter().zip(&staged) {
|
||||
if matches!(r, Ok(Some(_))) && n.append_at_ns > self.last_write_append_at_ns {
|
||||
self.last_write_append_at_ns = n.append_at_ns;
|
||||
self.last_write_needle_key = n.id;
|
||||
self.last_write_deleted = false;
|
||||
}
|
||||
}
|
||||
|
||||
// A durable entry that fails to publish stops the volume taking
|
||||
// writes, so the entries after it are refused the way a lone
|
||||
@@ -2521,6 +2544,9 @@ impl Volume {
|
||||
}
|
||||
|
||||
self.last_append_at_ns = n.append_at_ns;
|
||||
self.last_write_append_at_ns = n.append_at_ns;
|
||||
self.last_write_needle_key = n.id;
|
||||
self.last_write_deleted = false;
|
||||
|
||||
self.publish_write(n, offset, fsync)?;
|
||||
|
||||
@@ -2839,6 +2865,11 @@ impl Volume {
|
||||
if let Some(nm) = &mut self.nm {
|
||||
nm.delete(n.id, Offset::from_actual_offset(offset as i64))?;
|
||||
}
|
||||
if n.id == self.last_write_needle_key {
|
||||
self.last_write_append_at_ns = 0;
|
||||
self.last_write_needle_key = NeedleId(0);
|
||||
self.last_write_deleted = true;
|
||||
}
|
||||
let checkpoint_ok = self.maybe_checkpoint_index(false);
|
||||
|
||||
// Clear the EIO streak after a successful delete (tombstone append +
|
||||
@@ -3106,8 +3137,15 @@ impl Volume {
|
||||
return;
|
||||
}
|
||||
match self.find_last_write_append_at_ns() {
|
||||
Ok(0) => {}
|
||||
Ok(append_at_ns) => self.last_modified_ts_seconds = append_at_ns / 1_000_000_000,
|
||||
Ok((0, _)) => {}
|
||||
Ok((append_at_ns, key)) => {
|
||||
self.last_modified_ts_seconds = append_at_ns / 1_000_000_000;
|
||||
if append_at_ns > self.last_write_append_at_ns {
|
||||
self.last_write_append_at_ns = append_at_ns;
|
||||
self.last_write_needle_key = key;
|
||||
self.last_write_deleted = false;
|
||||
}
|
||||
}
|
||||
Err(e) => warn!(
|
||||
volume_id = self.id.0,
|
||||
error = %e,
|
||||
@@ -3117,7 +3155,8 @@ impl Volume {
|
||||
}
|
||||
|
||||
/// Scan the .idx backwards for the newest write — an entry that is not a
|
||||
/// deletion tombstone — and return that needle's append timestamp. The .idx
|
||||
/// deletion tombstone — and return that needle's append timestamp and key.
|
||||
/// The .idx
|
||||
/// and the .dat share an order, so an append-ordered volume answers with the
|
||||
/// first write the scan reaches. Vacuum rewrites both in key order, which
|
||||
/// tracks write order only because the master issues keys increasing: an
|
||||
@@ -3126,19 +3165,21 @@ impl Volume {
|
||||
/// nothing but tombstones, when a vacuumed volume holds more needles than
|
||||
/// the scan budget, or for a volume older than version 3, whose needles
|
||||
/// carry no append timestamp. Mirrors Go's findLastWriteAppendAtNs.
|
||||
fn find_last_write_append_at_ns(&self) -> Result<u64, VolumeError> {
|
||||
fn find_last_write_append_at_ns(&self) -> Result<(u64, NeedleId), VolumeError> {
|
||||
let version = self.version();
|
||||
if version != VERSION_3 {
|
||||
return Ok(0);
|
||||
return Ok((0, NeedleId(0)));
|
||||
}
|
||||
let idx_path = self.file_name(".idx");
|
||||
let idx_size = fs::metadata(&idx_path).map(|m| m.len()).unwrap_or(0) as i64;
|
||||
if idx_size == 0 || idx_size % NEEDLE_MAP_ENTRY_SIZE as i64 != 0 {
|
||||
return Ok(0);
|
||||
return Ok((0, NeedleId(0)));
|
||||
}
|
||||
let scan_every_write = self.super_block.compaction_revision > 0;
|
||||
let mut entry_budget = Self::VACUUMED_LAST_WRITE_SCAN_ENTRIES;
|
||||
let mut last_write_append_at_ns = 0u64;
|
||||
let mut last_write_key = NeedleId(0);
|
||||
let mut dead = HashSet::new();
|
||||
let mut idx_file = File::open(&idx_path)?;
|
||||
let mut block = vec![0u8; NEEDLE_MAP_ENTRY_SIZE * idx::ROWS_TO_READ];
|
||||
let mut end = idx_size;
|
||||
@@ -3149,7 +3190,13 @@ impl Volume {
|
||||
idx_file.read_exact(entries)?;
|
||||
for entry in entries.as_chunks::<NEEDLE_MAP_ENTRY_SIZE>().0.iter().rev() {
|
||||
let (key, offset, size) = idx_entry_from_bytes(entry);
|
||||
// The first row a key presents is its latest state: a tombstone
|
||||
// there retires the write rows beneath it.
|
||||
if dead.contains(&key) {
|
||||
continue;
|
||||
}
|
||||
if offset.is_zero() || size.is_deleted() {
|
||||
dead.insert(key);
|
||||
continue;
|
||||
}
|
||||
let Some(needle_offset) =
|
||||
@@ -3157,10 +3204,13 @@ impl Volume {
|
||||
else {
|
||||
continue;
|
||||
};
|
||||
last_write_append_at_ns = last_write_append_at_ns
|
||||
.max(self.read_needle_append_at_ns(needle_offset, size)?);
|
||||
let append_at_ns = self.read_needle_append_at_ns(needle_offset, size)?;
|
||||
if append_at_ns > last_write_append_at_ns {
|
||||
last_write_append_at_ns = append_at_ns;
|
||||
last_write_key = key;
|
||||
}
|
||||
if !scan_every_write {
|
||||
return Ok(last_write_append_at_ns);
|
||||
return Ok((last_write_append_at_ns, last_write_key));
|
||||
}
|
||||
entry_budget -= 1;
|
||||
if entry_budget == 0 {
|
||||
@@ -3169,12 +3219,12 @@ impl Volume {
|
||||
budget = Self::VACUUMED_LAST_WRITE_SCAN_ENTRIES,
|
||||
"too many needles to scan for the last write, keeping the .dat mtime"
|
||||
);
|
||||
return Ok(0);
|
||||
return Ok((0, NeedleId(0)));
|
||||
}
|
||||
}
|
||||
end = start;
|
||||
}
|
||||
Ok(last_write_append_at_ns)
|
||||
Ok((last_write_append_at_ns, last_write_key))
|
||||
}
|
||||
|
||||
/// The .dat offset holding the needle an .idx entry describes, or None when
|
||||
@@ -4294,6 +4344,9 @@ impl Volume {
|
||||
|
||||
// Update lastAppendAtNs (matches Go L352: v.lastAppendAtNs = appendAtNs)
|
||||
self.last_append_at_ns = append_at_ns;
|
||||
self.last_write_append_at_ns = append_at_ns;
|
||||
self.last_write_needle_key = needle_id;
|
||||
self.last_write_deleted = false;
|
||||
|
||||
// Update needle map index
|
||||
let offset = Offset::from_actual_offset(dat_size);
|
||||
@@ -4493,8 +4546,18 @@ impl Volume {
|
||||
self.write_compact_commit_marker()?;
|
||||
self.apply_compact_swap()?;
|
||||
|
||||
// Reload
|
||||
self.load(true, false, 0, self.version())?;
|
||||
// The write watermark already equals what recover_last_modified_ts
|
||||
// would rescan, so keep the clock instead of paying for the scan
|
||||
// under the lock. If its write was itself deleted, the reload
|
||||
// recovers the newest surviving write instead.
|
||||
if self.last_write_append_at_ns != 0 {
|
||||
self.last_modified_ts_seconds = self.last_write_append_at_ns / 1_000_000_000;
|
||||
}
|
||||
self.keep_last_modified_ts_on_load =
|
||||
self.last_modified_ts_seconds != 0 && !self.last_write_deleted;
|
||||
let load_result = self.load(true, false, 0, self.version());
|
||||
self.keep_last_modified_ts_on_load = false;
|
||||
load_result?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
@@ -6381,6 +6444,7 @@ mod tests {
|
||||
let prior = v.nm.as_ref().unwrap().get(NeedleId(1)).unwrap().unwrap();
|
||||
let dat_len_before = dat_len(&v);
|
||||
let last_append_before = v.last_append_at_ns;
|
||||
let last_write_append_before = v.last_write_append_at_ns;
|
||||
let last_modified_before = v.last_modified_ts_seconds;
|
||||
let file_count_before = v.file_count();
|
||||
|
||||
@@ -6402,6 +6466,7 @@ mod tests {
|
||||
);
|
||||
assert_eq!(dat_len(&v), dat_len_before, "the run is off the .dat");
|
||||
assert_eq!(v.last_append_at_ns, last_append_before);
|
||||
assert_eq!(v.last_write_append_at_ns, last_write_append_before);
|
||||
assert_eq!(v.last_modified_ts_seconds, last_modified_before);
|
||||
let now = v.nm.as_ref().unwrap().get(NeedleId(1)).unwrap().unwrap();
|
||||
assert_eq!((now.offset, now.size), (prior.offset, prior.size));
|
||||
@@ -7037,6 +7102,146 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
// Covers the reload that ends a vacuum commit: the clock must keep the
|
||||
// last write's append time rather than be re-derived, and the intervening
|
||||
// delete's tombstone must not freshen it.
|
||||
#[test]
|
||||
fn test_ttl_clock_carried_across_vacuum_commit() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let dir = tmp.path().to_str().unwrap();
|
||||
let ttl = crate::storage::needle::ttl::TTL::read("5m").unwrap();
|
||||
let last_write_ns = (SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.unwrap()
|
||||
.as_secs()
|
||||
- 2 * 60 * 60)
|
||||
* 1_000_000_000;
|
||||
|
||||
let mut v = make_ttl_volume(dir, ttl);
|
||||
let mut written = Vec::new();
|
||||
for i in 1..=3u64 {
|
||||
let data = format!("data {}", i);
|
||||
let mut n = Needle {
|
||||
id: NeedleId(i),
|
||||
cookie: Cookie(i as u32),
|
||||
data: data.as_bytes().to_vec(),
|
||||
data_size: data.len() as u32,
|
||||
..Needle::default()
|
||||
};
|
||||
let (offset, _, _) = v.write_needle(&mut n, true, false).unwrap();
|
||||
written.push((offset, n.size));
|
||||
}
|
||||
v.delete_needle(&mut Needle {
|
||||
id: NeedleId(2),
|
||||
cookie: Cookie(2),
|
||||
..Needle::default()
|
||||
})
|
||||
.unwrap();
|
||||
v.sync_to_disk().unwrap();
|
||||
for (offset, size) in written {
|
||||
backdate_append_at_ns(&v.dat_path(), offset, size, last_write_ns);
|
||||
}
|
||||
// Where a restart's recovery would have left the clock. The delete
|
||||
// above pushed last_append_at_ns to ~now; the commit must not use it.
|
||||
v.set_last_modified_ts_for_test(last_write_ns / 1_000_000_000);
|
||||
v.last_write_append_at_ns = last_write_ns;
|
||||
|
||||
v.compact_by_index(0, 0, |_| true).unwrap();
|
||||
v.commit_compact().unwrap();
|
||||
|
||||
assert_eq!(v.last_modified_ts(), last_write_ns / 1_000_000_000);
|
||||
assert!(
|
||||
v.is_expired(v.content_size(), 1024 * 1024),
|
||||
"a TTL volume whose last write is 2h old must stay expired across a vacuum commit"
|
||||
);
|
||||
}
|
||||
|
||||
// A write's client supplied modified time can lie ahead of or behind when
|
||||
// it was appended; the commit keeps the server-side write watermark the
|
||||
// recovery scan would recompute.
|
||||
#[test]
|
||||
fn test_ttl_clock_at_commit_uses_append_time() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let dir = tmp.path().to_str().unwrap();
|
||||
let ttl = crate::storage::needle::ttl::TTL::read("5m").unwrap();
|
||||
let future = SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.unwrap()
|
||||
.as_secs()
|
||||
+ 24 * 60 * 60;
|
||||
|
||||
let mut v = make_ttl_volume(dir, ttl);
|
||||
let data = b"data".to_vec();
|
||||
let mut n = Needle {
|
||||
id: NeedleId(1),
|
||||
cookie: Cookie(1),
|
||||
data,
|
||||
data_size: 4,
|
||||
last_modified: future,
|
||||
..Needle::default()
|
||||
};
|
||||
v.write_needle(&mut n, true, false).unwrap();
|
||||
assert!(v.last_modified_ts() >= future);
|
||||
let append_watermark_sec = v.last_write_append_at_ns / 1_000_000_000;
|
||||
|
||||
v.compact_by_index(0, 0, |_| true).unwrap();
|
||||
v.commit_compact().unwrap();
|
||||
|
||||
assert_eq!(v.last_modified_ts(), append_watermark_sec);
|
||||
}
|
||||
|
||||
// A vacuum commit whose newest write was deleted first: the carried
|
||||
// watermark belonged to that write, so the commit has to let the reload
|
||||
// rescan and land on the newest surviving write rather than keep the
|
||||
// volume alive on a deleted write's time.
|
||||
#[test]
|
||||
fn test_ttl_clock_at_commit_skips_deleted_write() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let dir = tmp.path().to_str().unwrap();
|
||||
let ttl = crate::storage::needle::ttl::TTL::read("5m").unwrap();
|
||||
let now = SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.unwrap()
|
||||
.as_secs();
|
||||
let old_write_ns = (now - 2 * 60 * 60) * 1_000_000_000;
|
||||
let new_write_ns = (now - 60 * 60) * 1_000_000_000;
|
||||
|
||||
let mut v = make_ttl_volume(dir, ttl);
|
||||
for (id, append_at_ns) in [(1u64, old_write_ns), (2, new_write_ns)] {
|
||||
let data = format!("data {}", id);
|
||||
let mut n = Needle {
|
||||
id: NeedleId(id),
|
||||
cookie: Cookie(id as u32),
|
||||
data: data.as_bytes().to_vec(),
|
||||
data_size: data.len() as u32,
|
||||
..Needle::default()
|
||||
};
|
||||
let (offset, _, _) = v.write_needle(&mut n, true, false).unwrap();
|
||||
v.sync_to_disk().unwrap();
|
||||
backdate_append_at_ns(&v.dat_path(), offset, n.size, append_at_ns);
|
||||
}
|
||||
// The delete lands inside the commit window: makeup_diff replays its
|
||||
// tombstone into the new .idx behind the write row the copy carried.
|
||||
v.compact_by_index(0, 0, |_| true).unwrap();
|
||||
v.delete_needle(&mut Needle {
|
||||
id: NeedleId(2),
|
||||
cookie: Cookie(2),
|
||||
..Needle::default()
|
||||
})
|
||||
.unwrap();
|
||||
assert!(
|
||||
v.last_write_deleted,
|
||||
"deleting the newest write must mark the watermark dead"
|
||||
);
|
||||
v.commit_compact().unwrap();
|
||||
|
||||
assert_eq!(v.last_modified_ts(), old_write_ns / 1_000_000_000);
|
||||
assert!(
|
||||
!v.last_write_deleted,
|
||||
"the reload's rescan must reseed the watermark off the surviving write"
|
||||
);
|
||||
}
|
||||
|
||||
// Guard the destroy time an EC volume is reclaimed on: it was recomputed as
|
||||
// now+TTL every time the .vif was written, so a read-only mark, a tier
|
||||
// upload or an EC encode handed an already expiring volume another full TTL.
|
||||
|
||||
+7
-10
@@ -7,20 +7,20 @@ require (
|
||||
github.com/linkedin/goavro/v2 v2.15.0
|
||||
github.com/seaweedfs/seaweedfs v0.0.0-00010101000000-000000000000
|
||||
github.com/segmentio/kafka-go v0.4.49
|
||||
github.com/stretchr/testify v1.11.1
|
||||
github.com/stretchr/testify v1.12.1
|
||||
google.golang.org/grpc v1.85.0-dev.0.20260915183914-4e49413dcab7
|
||||
)
|
||||
|
||||
replace github.com/seaweedfs/seaweedfs => ../../
|
||||
|
||||
require (
|
||||
github.com/andybalholm/brotli v1.2.2 // indirect
|
||||
github.com/andybalholm/brotli v1.2.3 // indirect
|
||||
github.com/aws/aws-sdk-go v1.55.8 // indirect
|
||||
github.com/beorn7/perks v1.0.1 // indirect
|
||||
github.com/cespare/xxhash/v2 v2.3.0 // indirect
|
||||
github.com/cognusion/imaging v1.0.4 // indirect
|
||||
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect
|
||||
github.com/dustin/go-humanize v1.0.1 // indirect
|
||||
github.com/dustin/go-humanize v1.1.0 // indirect
|
||||
github.com/eapache/go-resiliency v1.7.0 // indirect
|
||||
github.com/eapache/go-xerial-snappy v0.0.0-20230731223053-c322873962e3 // indirect
|
||||
github.com/eapache/queue v1.1.0 // indirect
|
||||
@@ -61,17 +61,15 @@ require (
|
||||
github.com/petermattis/goid v0.0.0-20260113132338-7c7de50cc741 // indirect
|
||||
github.com/pierrec/lz4/v4 v4.1.29 // indirect
|
||||
github.com/planetscale/vtprotobuf v0.6.1-0.20240319094008-0393e58bdf10 // indirect
|
||||
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect
|
||||
github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 // indirect
|
||||
github.com/prometheus/client_golang v1.24.1 // indirect
|
||||
github.com/prometheus/client_model v0.6.3 // indirect
|
||||
github.com/prometheus/common v0.70.1 // indirect
|
||||
github.com/prometheus/common v0.71.0 // indirect
|
||||
github.com/prometheus/procfs v0.22.0 // indirect
|
||||
github.com/rcrowley/go-metrics v0.0.0-20250401214520-65e299d6c5c9 // indirect
|
||||
github.com/rdleal/intervalst v1.5.0 // indirect
|
||||
github.com/rwcarlsen/goexif v0.0.0-20190401172101-9e8deecbddbd // indirect
|
||||
github.com/sagikazarmark/locafero v0.11.0 // indirect
|
||||
github.com/seaweedfs/goexif v2.0.0+incompatible // indirect
|
||||
github.com/seaweedfs/goexif v1.0.3 // indirect
|
||||
github.com/shirou/gopsutil/v4 v4.26.7 // indirect
|
||||
github.com/sourcegraph/conc v0.3.1-0.20240121214520-5f936abd7ae8 // indirect
|
||||
github.com/spf13/afero v1.15.0 // indirect
|
||||
@@ -91,8 +89,8 @@ require (
|
||||
github.com/xeipuuv/gojsonreference v0.0.0-20180127040603-bd5ef7bd5415 // indirect
|
||||
github.com/xeipuuv/gojsonschema v1.2.0 // indirect
|
||||
github.com/yusufpapurcu/wmi v1.2.4 // indirect
|
||||
go.yaml.in/yaml/v3 v3.0.4 // indirect
|
||||
golang.org/x/crypto v0.56.0 // indirect
|
||||
go.yaml.in/yaml/v3 v3.0.5 // indirect
|
||||
golang.org/x/crypto v0.57.0 // indirect
|
||||
golang.org/x/image v0.46.0 // indirect
|
||||
golang.org/x/net v0.58.0 // indirect
|
||||
golang.org/x/sync v0.23.0 // indirect
|
||||
@@ -101,5 +99,4 @@ require (
|
||||
google.golang.org/genproto/googleapis/rpc v0.0.0-20260819154853-08b0e4226688 // indirect
|
||||
google.golang.org/grpc/security/advancedtls v1.0.0 // indirect
|
||||
google.golang.org/protobuf v1.36.12 // indirect
|
||||
gopkg.in/yaml.v3 v3.0.1 // indirect
|
||||
)
|
||||
|
||||
+16
-20
@@ -10,8 +10,8 @@ github.com/alecthomas/assert/v2 v2.10.0 h1:jjRCHsj6hBJhkmhznrCzoNpbA3zqy0fYiUcYZ
|
||||
github.com/alecthomas/assert/v2 v2.10.0/go.mod h1:Bze95FyfUr7x34QZrjL+XP+0qgp/zg8yS+TtBj1WA3k=
|
||||
github.com/alecthomas/repr v0.4.0 h1:GhI2A8MACjfegCPVq9f1FLvIBS+DrQ2KQBFZP1iFzXc=
|
||||
github.com/alecthomas/repr v0.4.0/go.mod h1:Fr0507jx4eOXV7AlPV6AVZLYrLIuIeSOWtW57eE/O/4=
|
||||
github.com/andybalholm/brotli v1.2.2 h1:HzTuoo2ErYQqf5qvcJInB8uvqSVxRttzkFexPWtnceM=
|
||||
github.com/andybalholm/brotli v1.2.2/go.mod h1:rzTDkvFWvIrjDXZHkuS16NPggd91W3kUSvPlQ1pLaKY=
|
||||
github.com/andybalholm/brotli v1.2.3 h1:8H1qwOkl2LPfjf3YezB90JnCliZb6SInJ/OJkEbA5NQ=
|
||||
github.com/andybalholm/brotli v1.2.3/go.mod h1:rzTDkvFWvIrjDXZHkuS16NPggd91W3kUSvPlQ1pLaKY=
|
||||
github.com/aws/aws-sdk-go v1.55.8 h1:JRmEUbU52aJQZ2AjX4q4Wu7t4uZjOu71uyNmaWlUkJQ=
|
||||
github.com/aws/aws-sdk-go v1.55.8/go.mod h1:ZkViS9AqA6otK+JBBNH2++sx1sgxrPKcSzPPvQkUtXk=
|
||||
github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM=
|
||||
@@ -28,8 +28,8 @@ github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSs
|
||||
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
|
||||
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1VwoXQT9A3Wy9MM3WgvqSxFWenqJduM=
|
||||
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
|
||||
github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY=
|
||||
github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto=
|
||||
github.com/dustin/go-humanize v1.1.0 h1:dbKTrvD0klcbBV/h4AWJdMuZogJACoMlvWIWZ5b2xWg=
|
||||
github.com/dustin/go-humanize v1.1.0/go.mod h1:hc1CvRkJMsgxqjmjMQF3QNRAZBwY8AXBAzKYoSX9sFI=
|
||||
github.com/eapache/go-resiliency v1.7.0 h1:n3NRTnBn5N0Cbi/IeOHuQn9s2UwVUH7Ga0ZWcP+9JTA=
|
||||
github.com/eapache/go-resiliency v1.7.0/go.mod h1:5yPzW0MIvSe0JDsv0v+DvcjEv2FyD6iZYSs1ZI+iQho=
|
||||
github.com/eapache/go-xerial-snappy v0.0.0-20230731223053-c322873962e3 h1:Oy0F4ALJ04o5Qqpdz8XLIpNA3WM/iSIXqxtqo7UGVws=
|
||||
@@ -174,8 +174,8 @@ github.com/prometheus/client_golang v1.24.1 h1:JnJkREXzWxUdCuPFpIWZiPispT9xVV59u
|
||||
github.com/prometheus/client_golang v1.24.1/go.mod h1:F+oSRECHg4sse5ucfYpYDeIv/hu68Zo0uoHKetWnzcE=
|
||||
github.com/prometheus/client_model v0.6.3 h1:O0jaTVAYNxTHYInEPFJt5I3+sN8zqBtVMPTB1qyxiEo=
|
||||
github.com/prometheus/client_model v0.6.3/go.mod h1:gpN5P9S7Rr6Yr92PiQ+Ixvhf6JZEkF1dnxsYL2aPBEM=
|
||||
github.com/prometheus/common v0.70.1 h1:1HvjP4D5oL3t8RsPlwxA9onvvStjtIHYE5XuuwOi/PY=
|
||||
github.com/prometheus/common v0.70.1/go.mod h1:VdFUQDMZK3VLkurFUVhia6uys/0suUp86TJz5qbJRhc=
|
||||
github.com/prometheus/common v0.71.0 h1:9KDAKb7Mj3HEVKyFCK6Dc/HIwlBzZIN2l7/lrHl3KK8=
|
||||
github.com/prometheus/common v0.71.0/go.mod h1:CLJ5H8TEsGX8bl31BdMkfhIZ+QmZ9tBPPotUxUbfcmk=
|
||||
github.com/prometheus/procfs v0.22.0 h1:6q9+/JL9IKAPbCmBrv9n5O5Ty3NKnciV5X7YGw0oics=
|
||||
github.com/prometheus/procfs v0.22.0/go.mod h1:CvmFr/GVhIjIvWJZW3tgkODBQMRIf0EyWMQLHCHab58=
|
||||
github.com/rcrowley/go-metrics v0.0.0-20250401214520-65e299d6c5c9 h1:bsUq1dX0N8AOIL7EB/X911+m4EHsnWEHeJ0c+3TTBrg=
|
||||
@@ -184,12 +184,10 @@ github.com/rdleal/intervalst v1.5.0 h1:SEB9bCFz5IqD1yhfH1Wv8IBnY/JQxDplwkxHjT6ha
|
||||
github.com/rdleal/intervalst v1.5.0/go.mod h1:xO89Z6BC+LQDH+IPQQw/OESt5UADgFD41tYMUINGpxQ=
|
||||
github.com/rogpeppe/go-internal v1.15.0 h1:D0RCU5rMAp+SpgkiNdrjfJ+LX4J1M32V2NeCY7EJ6hc=
|
||||
github.com/rogpeppe/go-internal v1.15.0/go.mod h1:DrUVZyrJU+txYW5/1kwtXQSMFio52ZOxX7yM1VHvnxs=
|
||||
github.com/rwcarlsen/goexif v0.0.0-20190401172101-9e8deecbddbd h1:CmH9+J6ZSsIjUK3dcGsnCnO41eRBOnY12zwkn5qVwgc=
|
||||
github.com/rwcarlsen/goexif v0.0.0-20190401172101-9e8deecbddbd/go.mod h1:hPqNNc0+uJM6H+SuU8sEs5K5IQeKccPqeSjfgcKGgPk=
|
||||
github.com/sagikazarmark/locafero v0.11.0 h1:1iurJgmM9G3PA/I+wWYIOw/5SyBtxapeHDcg+AAIFXc=
|
||||
github.com/sagikazarmark/locafero v0.11.0/go.mod h1:nVIGvgyzw595SUSUE6tvCp3YYTeHs15MvlmU87WwIik=
|
||||
github.com/seaweedfs/goexif v2.0.0+incompatible h1:x8pckiT12QQhifwhDQpeISgDfsqmQ6VR4LFPQ64JRps=
|
||||
github.com/seaweedfs/goexif v2.0.0+incompatible/go.mod h1:Oni780Z236sXpIQzk1XoJlTwqrJ02smEin9zQeff7Fk=
|
||||
github.com/seaweedfs/goexif v1.0.3 h1:ve/OjI7dxPW8X9YQsv3JuVMaxEyF9Rvfd04ouL+Bz30=
|
||||
github.com/seaweedfs/goexif v1.0.3/go.mod h1:Oni780Z236sXpIQzk1XoJlTwqrJ02smEin9zQeff7Fk=
|
||||
github.com/segmentio/kafka-go v0.4.49 h1:GJiNX1d/g+kG6ljyJEoi9++PUMdXGAxb7JGPiDCuNmk=
|
||||
github.com/segmentio/kafka-go v0.4.49/go.mod h1:Y1gn60kzLEEaW28YshXyk2+VCUKbJ3Qr6DrnT3i4+9E=
|
||||
github.com/shirou/gopsutil/v4 v4.26.7 h1:IXzpHz/dkMRYAhKkOXr1HB6SuzWU3eoyyeWe7g3bNZc=
|
||||
@@ -216,8 +214,8 @@ github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/
|
||||
github.com/stretchr/testify v1.7.5/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU=
|
||||
github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU=
|
||||
github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4=
|
||||
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
|
||||
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
|
||||
github.com/stretchr/testify v1.12.1 h1:EuwCh5fleGS7H32xRwO3wRGT7DxrDhLAT6FF8MpWDWE=
|
||||
github.com/stretchr/testify v1.12.1/go.mod h1:MDEgiDPPsNp5cuIrHPPCyornHKgEVbtFUmoNlxoYthg=
|
||||
github.com/subosito/gotenv v1.6.0 h1:9NlTDc1FTs4qu0DDq7AEtTPNw6SVm7uBMsUCUjABIf8=
|
||||
github.com/subosito/gotenv v1.6.0/go.mod h1:Dk4QP5c2W3ibzajGcXpNraDfq2IrhjMIvMSWPKKo0FU=
|
||||
github.com/syndtr/goleveldb v1.0.1-0.20190318030020-c3a204f8e965 h1:1oFLiOyVl+W7bnBzGhf7BbIv9loSFQcieWWYIjLqcAw=
|
||||
@@ -263,15 +261,15 @@ go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto=
|
||||
go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE=
|
||||
go.yaml.in/yaml/v2 v2.4.4 h1:tuyd0P+2Ont/d6e2rl3be67goVK4R6deVxCUX5vyPaQ=
|
||||
go.yaml.in/yaml/v2 v2.4.4/go.mod h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ=
|
||||
go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc=
|
||||
go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg=
|
||||
go.yaml.in/yaml/v3 v3.0.5 h1:N6y/pJk8buWs9NY5ERU2HSMfm+IuD/OtfdAnq6kESPw=
|
||||
go.yaml.in/yaml/v3 v3.0.5/go.mod h1:HVTZu1O7/Vkt2N+BFy8Zza+lnLsABggaTM2ZpNIGuKg=
|
||||
golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
|
||||
golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
|
||||
golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto=
|
||||
golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc=
|
||||
golang.org/x/crypto v0.6.0/go.mod h1:OFC/31mSvZgRz0V1QTNCzfAI1aIRzbiufJtkMIlEp58=
|
||||
golang.org/x/crypto v0.56.0 h1:GUh5Ii4J5jtcseSMiRqr1jXCNHoxjeV9Fmekc2oLy6Y=
|
||||
golang.org/x/crypto v0.56.0/go.mod h1:OMW5y6CY9l38uPLmxU6l6pwcXp1obtLo3e6gT7gQR2I=
|
||||
golang.org/x/crypto v0.57.0 h1:3ZVCjf8Ggz7zneR/EHRVx68Ctf+2pmIMP2UFhh9cC6M=
|
||||
golang.org/x/crypto v0.57.0/go.mod h1:Fdz0i5U6CoizGwLda9DttjSk6qlZo25zYNtR+ycvuZA=
|
||||
golang.org/x/image v0.46.0 h1:b1+oYj0Jbp6K5MDT4i4/eZpYlk3V8SJhhDKh6LBHAyQ=
|
||||
golang.org/x/image v0.46.0/go.mod h1:3B3W05VGVQyuXucLINLjXKrqISASfi4Xj+iCVkLMwew=
|
||||
golang.org/x/mod v0.3.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA=
|
||||
@@ -288,8 +286,8 @@ golang.org/x/net v0.6.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs=
|
||||
golang.org/x/net v0.7.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs=
|
||||
golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To=
|
||||
golang.org/x/net v0.58.0/go.mod h1:YwCddHnFlT7eLQqVprV19OnhLGtc5xOKgE0RyqgfWAU=
|
||||
golang.org/x/oauth2 v0.36.0 h1:peZ/1z27fi9hUOFCAZaHyrpWG5lwe0RJEEEeH0ThlIs=
|
||||
golang.org/x/oauth2 v0.36.0/go.mod h1:YDBUJMTkDnJS+A4BP4eZBjCqtokkg1hODuPjwiGPO7Q=
|
||||
golang.org/x/oauth2 v0.37.0 h1:JUlcxA8oAtauLfiH8FX2/FkAWHAdi0QtGCGc+hofE98=
|
||||
golang.org/x/oauth2 v0.37.0/go.mod h1:IxwZNxUULJmpBFf9K/9NTMSIfZZuvuTy1gGxhigP/58=
|
||||
golang.org/x/sync v0.0.0-20180314180146-1d60e4601c6f/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
|
||||
golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
|
||||
golang.org/x/sync v0.0.0-20201020160332-67f06af15bc9/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
|
||||
@@ -351,8 +349,6 @@ google.golang.org/protobuf v1.23.0/go.mod h1:EGpADcykh3NcUnDUJcl1+ZksZNG86OlYog2
|
||||
google.golang.org/protobuf v1.36.12 h1:pJOKDDOyeXErUroCihFAd5LQuwXBSpVnKGrj5o/fwxc=
|
||||
google.golang.org/protobuf v1.36.12/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco=
|
||||
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk=
|
||||
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q=
|
||||
gopkg.in/fsnotify.v1 v1.4.7/go.mod h1:Tz8NjZHkW78fSQdbUxIjBTcgA1z1m8ZHf0WmKUhAMys=
|
||||
gopkg.in/tomb.v1 v1.0.0-20141024135613-dd632973f1e7 h1:uRGJdciOHaEIrze2W8Q3AKkepLTh2hOroT7a+7czfdQ=
|
||||
gopkg.in/tomb.v1 v1.0.0-20141024135613-dd632973f1e7/go.mod h1:dt/ZhP58zS4L8KSrWDmTeBkI65Dw0HsyUHuEVlX15mw=
|
||||
|
||||
@@ -98,6 +98,10 @@ func (cp *ConfigPersistence) ApplyMaintenanceConfigFromToml(v TomlConfig) error
|
||||
ecConf.ReplicaPlacement = v.GetString(k)
|
||||
ecChanged = true
|
||||
}
|
||||
if k := "maintenance.erasure_coding.strict_placement"; v.IsSet(k) {
|
||||
ecConf.StrictPlacement = v.GetBool(k)
|
||||
ecChanged = true
|
||||
}
|
||||
|
||||
if !maintenanceChanged && !vacuumChanged && !balanceChanged && !ecChanged {
|
||||
return nil
|
||||
@@ -224,6 +228,7 @@ var pluginConfigSections = []pluginConfigSection{
|
||||
"min_size_mb": int64Value,
|
||||
"preferred_tags": stringListValue,
|
||||
"replica_placement": stringValue,
|
||||
"strict_placement": boolValue,
|
||||
},
|
||||
// workers read collection_filter from the admin values, not the worker values
|
||||
adminKeys: map[string]func(v TomlConfig, key string) *plugin_pb.ConfigValue{
|
||||
@@ -240,6 +245,10 @@ func int64Value(v TomlConfig, key string) *plugin_pb.ConfigValue {
|
||||
return &plugin_pb.ConfigValue{Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: int64(v.GetInt(key))}}
|
||||
}
|
||||
|
||||
func boolValue(v TomlConfig, key string) *plugin_pb.ConfigValue {
|
||||
return &plugin_pb.ConfigValue{Kind: &plugin_pb.ConfigValue_BoolValue{BoolValue: v.GetBool(key)}}
|
||||
}
|
||||
|
||||
func stringValue(v TomlConfig, key string) *plugin_pb.ConfigValue {
|
||||
return &plugin_pb.ConfigValue{Kind: &plugin_pb.ConfigValue_StringValue{StringValue: v.GetString(key)}}
|
||||
}
|
||||
|
||||
@@ -174,17 +174,24 @@ func runFilerSynchronize(cmd *Command, args []string) bool {
|
||||
|
||||
// per-cluster HTTPS clients for volume server connections
|
||||
var httpClientA, httpClientB *util_http_client.HTTPClient
|
||||
var jwtForFilerA, jwtForFilerB security.FilerJwtProvider
|
||||
if *syncOptions.aSecurity != "" {
|
||||
var err error
|
||||
if httpClientA, err = security.LoadHTTPClientFromFile(*syncOptions.aSecurity); err != nil {
|
||||
glog.Fatalf("load HTTPS client config for filer A: %v", err)
|
||||
}
|
||||
if jwtForFilerA, err = security.LoadFilerJwtFromFile(*syncOptions.aSecurity); err != nil {
|
||||
glog.Fatalf("load filer JWT config for filer A: %v", err)
|
||||
}
|
||||
}
|
||||
if *syncOptions.bSecurity != "" {
|
||||
var err error
|
||||
if httpClientB, err = security.LoadHTTPClientFromFile(*syncOptions.bSecurity); err != nil {
|
||||
glog.Fatalf("load HTTPS client config for filer B: %v", err)
|
||||
}
|
||||
if jwtForFilerB, err = security.LoadFilerJwtFromFile(*syncOptions.bSecurity); err != nil {
|
||||
glog.Fatalf("load filer JWT config for filer B: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
grace.SetupProfiling(*syncCpuProfile, *syncMemProfile)
|
||||
@@ -263,7 +270,9 @@ func runFilerSynchronize(cmd *Command, args []string) bool {
|
||||
bFilerSignature,
|
||||
&syncStateA2B,
|
||||
httpClientA,
|
||||
httpClientB)
|
||||
httpClientB,
|
||||
jwtForFilerA,
|
||||
jwtForFilerB)
|
||||
if err != nil {
|
||||
glog.Errorf("sync from %s to %s: %v", *syncOptions.filerA, *syncOptions.filerB, err)
|
||||
time.Sleep(1747 * time.Millisecond)
|
||||
@@ -306,7 +315,9 @@ func runFilerSynchronize(cmd *Command, args []string) bool {
|
||||
aFilerSignature,
|
||||
&syncStateB2A,
|
||||
httpClientB,
|
||||
httpClientA)
|
||||
httpClientA,
|
||||
jwtForFilerB,
|
||||
jwtForFilerA)
|
||||
if err != nil {
|
||||
glog.Errorf("sync from %s to %s: %v", *syncOptions.filerB, *syncOptions.filerA, err)
|
||||
time.Sleep(2147 * time.Millisecond)
|
||||
@@ -336,7 +347,8 @@ func initOffsetFromTsMs(grpcDialOption grpc.DialOption, targetFiler pb.ServerAdd
|
||||
|
||||
func doSubscribeFilerMetaChanges(clientId int32, clientEpoch int32, sourceGrpcDialOption grpc.DialOption, sourceFiler pb.ServerAddress, sourcePath string, sourceExcludePaths []string, sourceReadChunkFromFiler bool, targetGrpcDialOption grpc.DialOption, targetFiler pb.ServerAddress, targetPath string,
|
||||
replicationStr, collection string, ttlSec int, sinkWriteChunkByFiler bool, diskType string, debug bool, concurrency int, chunkConcurrency int, doDeleteFiles bool, sourceFilerSignature int32, targetFilerSignature int32, statePtr *atomic.Pointer[syncState],
|
||||
sourceHttpClient *util_http_client.HTTPClient, sinkHttpClient *util_http_client.HTTPClient) error {
|
||||
sourceHttpClient *util_http_client.HTTPClient, sinkHttpClient *util_http_client.HTTPClient,
|
||||
sourceJwtProvider security.FilerJwtProvider, sinkJwtProvider security.FilerJwtProvider) error {
|
||||
|
||||
// if first time, start from now
|
||||
// if has previously synced, resume from that point of time
|
||||
@@ -357,12 +369,18 @@ func doSubscribeFilerMetaChanges(clientId int32, clientEpoch int32, sourceGrpcDi
|
||||
if sourceHttpClient != nil {
|
||||
filerSource.SetHttpClient(sourceHttpClient)
|
||||
}
|
||||
if sourceJwtProvider != nil {
|
||||
filerSource.SetFilerJwtProvider(sourceJwtProvider)
|
||||
}
|
||||
filerSink := &filersink.FilerSink{}
|
||||
filerSink.DoInitialize(targetFiler.ToHttpAddress(), targetFiler.ToGrpcAddress(), targetPath, replicationStr, collection, ttlSec, diskType, targetGrpcDialOption, sinkWriteChunkByFiler)
|
||||
filerSink.SetChunkConcurrency(chunkConcurrency)
|
||||
if sinkHttpClient != nil {
|
||||
filerSink.SetUploader(operation.NewUploaderWithHttpClient(sinkHttpClient))
|
||||
}
|
||||
if sinkJwtProvider != nil {
|
||||
filerSink.SetFilerJwtProvider(sinkJwtProvider)
|
||||
}
|
||||
filerSink.SetSourceFiler(filerSource)
|
||||
|
||||
persistEventFn := genProcessFunction(sourcePath, targetPath, sourceExcludePaths, nil, nil, nil, filerSink, doDeleteFiles, debug)
|
||||
|
||||
@@ -64,6 +64,8 @@
|
||||
# preferred_tags = ["fast", "ssd"]
|
||||
# EC shard placement constraint, e.g. "020"; empty uses the master default replication
|
||||
# replica_placement = ""
|
||||
# fail planning when the placement constraints cannot be met instead of relaxing them
|
||||
# strict_placement = false
|
||||
# max retry attempts for a failed erasure coding job
|
||||
# retry_limit = 1
|
||||
# seconds to wait between retry attempts
|
||||
|
||||
@@ -240,7 +240,7 @@ func (group *ChunkGroup) SetChunks(chunks []*filer_pb.FileChunk) error {
|
||||
continue
|
||||
}
|
||||
|
||||
resolvedChunks, err := resolveOneChunkManifest(context.Background(), group.lookupFn, chunk, group.cacheInvalidator, group.manifestCache)
|
||||
resolvedChunks, err := resolveOneChunkManifest(context.Background(), group.lookupFn, chunk, group.cacheInvalidator, group.manifestCache, nil)
|
||||
if err != nil {
|
||||
group.resolveErr = err
|
||||
return err
|
||||
|
||||
@@ -16,6 +16,7 @@ import (
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/security"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
util_http "github.com/seaweedfs/seaweedfs/weed/util/http"
|
||||
)
|
||||
@@ -55,7 +56,16 @@ func SeparateManifestChunks(chunks []*filer_pb.FileChunk) (manifestChunks, nonMa
|
||||
}
|
||||
|
||||
func ResolveChunkManifest(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, chunks []*filer_pb.FileChunk, startOffset, stopOffset int64, invalidator CacheInvalidator) (dataChunks, manifestChunks []*filer_pb.FileChunk, manifestResolveErr error) {
|
||||
resolver := newChunkManifestResolver(ctx, lookupFileIdFn, invalidator)
|
||||
resolver := newChunkManifestResolver(ctx, lookupFileIdFn, invalidator, nil)
|
||||
defer resolver.close()
|
||||
return resolver.resolve(chunks, startOffset, stopOffset)
|
||||
}
|
||||
|
||||
// ResolveChunkManifestWithFilerJwt is ResolveChunkManifest that signs proxied
|
||||
// manifest downloads with filerJwtFn instead of the process-wide filer read
|
||||
// key, for readers carrying per-source credentials.
|
||||
func ResolveChunkManifestWithFilerJwt(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, filerJwtFn security.FilerJwtProvider, chunks []*filer_pb.FileChunk, startOffset, stopOffset int64, invalidator CacheInvalidator) (dataChunks, manifestChunks []*filer_pb.FileChunk, manifestResolveErr error) {
|
||||
resolver := newChunkManifestResolver(ctx, lookupFileIdFn, invalidator, filerJwtFn)
|
||||
defer resolver.close()
|
||||
return resolver.resolve(chunks, startOffset, stopOffset)
|
||||
}
|
||||
@@ -81,6 +91,7 @@ type chunkManifestResolver struct {
|
||||
cancel context.CancelFunc
|
||||
lookupFileIdFn wdclient.LookupFileIdFunctionType
|
||||
invalidator CacheInvalidator
|
||||
filerJwtFn security.FilerJwtProvider
|
||||
jobs chan chunkManifestResolveJob
|
||||
overflowSem chan struct{}
|
||||
workers sync.WaitGroup
|
||||
@@ -88,7 +99,7 @@ type chunkManifestResolver struct {
|
||||
started bool
|
||||
}
|
||||
|
||||
func newChunkManifestResolver(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, invalidator CacheInvalidator) *chunkManifestResolver {
|
||||
func newChunkManifestResolver(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, invalidator CacheInvalidator, filerJwtFn security.FilerJwtProvider) *chunkManifestResolver {
|
||||
workCtx, cancel := context.WithCancel(ctx)
|
||||
resolver := &chunkManifestResolver{
|
||||
ctx: workCtx,
|
||||
@@ -96,6 +107,7 @@ func newChunkManifestResolver(ctx context.Context, lookupFileIdFn wdclient.Looku
|
||||
cancel: cancel,
|
||||
lookupFileIdFn: lookupFileIdFn,
|
||||
invalidator: invalidator,
|
||||
filerJwtFn: filerJwtFn,
|
||||
jobs: make(chan chunkManifestResolveJob, chunkManifestResolveJobBufferSize),
|
||||
overflowSem: make(chan struct{}, maxChunkManifestResolveWorkers),
|
||||
}
|
||||
@@ -103,7 +115,7 @@ func newChunkManifestResolver(ctx context.Context, lookupFileIdFn wdclient.Looku
|
||||
}
|
||||
|
||||
func (r *chunkManifestResolver) executeJob(job chunkManifestResolveJob) {
|
||||
job.result.chunks, job.result.err = ResolveOneChunkManifest(job.batchCtx, r.lookupFileIdFn, job.chunk, r.invalidator)
|
||||
job.result.chunks, job.result.err = resolveOneChunkManifest(job.batchCtx, r.lookupFileIdFn, job.chunk, r.invalidator, nil, r.filerJwtFn)
|
||||
if job.result.err != nil && r.parentCtx.Err() == nil {
|
||||
if job.batchCtx.Err() != nil && errors.Is(job.result.err, context.Canceled) {
|
||||
job.result.internalCancel = true
|
||||
@@ -293,14 +305,23 @@ func (r *chunkManifestResolver) resolve(chunks []*filer_pb.FileChunk, startOffse
|
||||
// Keeping this signature stable preserves the existing four-argument contract
|
||||
// for external callers.
|
||||
func ResolveOneChunkManifest(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, chunk *filer_pb.FileChunk, invalidator CacheInvalidator) (dataChunks []*filer_pb.FileChunk, manifestResolveErr error) {
|
||||
return resolveOneChunkManifest(ctx, lookupFileIdFn, chunk, invalidator, nil)
|
||||
return resolveOneChunkManifest(ctx, lookupFileIdFn, chunk, invalidator, nil, nil)
|
||||
}
|
||||
|
||||
// ResolveOneChunkManifestWithFilerJwt is ResolveOneChunkManifest that signs
|
||||
// proxied manifest downloads with filerJwtFn instead of the process-wide
|
||||
// filer read key, for readers carrying per-source credentials.
|
||||
func ResolveOneChunkManifestWithFilerJwt(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, filerJwtFn security.FilerJwtProvider, chunk *filer_pb.FileChunk, invalidator CacheInvalidator) (dataChunks []*filer_pb.FileChunk, manifestResolveErr error) {
|
||||
return resolveOneChunkManifest(ctx, lookupFileIdFn, chunk, invalidator, nil, filerJwtFn)
|
||||
}
|
||||
|
||||
// resolveOneChunkManifest is the cache-aware implementation. cache may be nil,
|
||||
// in which case the manifest is fetched and validated on every call, matching
|
||||
// the historical uncached behavior. A non-nil cache is owned by a single mount
|
||||
// (WFS) and coalesces concurrent cold misses via singleflight.
|
||||
func resolveOneChunkManifest(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, chunk *filer_pb.FileChunk, invalidator CacheInvalidator, cache *ChunkManifestCache) (dataChunks []*filer_pb.FileChunk, manifestResolveErr error) {
|
||||
// (WFS) and coalesces concurrent cold misses via singleflight. filerJwtFn
|
||||
// overrides the filer read credential for proxied downloads; nil means the
|
||||
// process-wide key.
|
||||
func resolveOneChunkManifest(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, chunk *filer_pb.FileChunk, invalidator CacheInvalidator, cache *ChunkManifestCache, filerJwtFn security.FilerJwtProvider) (dataChunks []*filer_pb.FileChunk, manifestResolveErr error) {
|
||||
if !chunk.IsChunkManifest {
|
||||
return
|
||||
}
|
||||
@@ -318,7 +339,7 @@ func resolveOneChunkManifest(ctx context.Context, lookupFileIdFn wdclient.Lookup
|
||||
bytesBuffer := bytesBufferPool.Get().(*bytes.Buffer)
|
||||
bytesBuffer.Reset()
|
||||
defer bytesBufferPool.Put(bytesBuffer)
|
||||
if err := fetchWholeChunk(ctx, bytesBuffer, lookupFileIdFn, key.fileID, chunk.CipherKey, chunk.IsCompressed, invalidator); err != nil {
|
||||
if err := fetchWholeChunk(ctx, bytesBuffer, lookupFileIdFn, key.fileID, chunk.CipherKey, chunk.IsCompressed, invalidator, filerJwtFn); err != nil {
|
||||
return nil, fmt.Errorf("fail to read manifest %s: %w", key.fileID, err)
|
||||
}
|
||||
// Copy before the buffer returns to the pool so concurrent callers
|
||||
@@ -357,13 +378,18 @@ func resolveOneChunkManifest(ctx context.Context, lookupFileIdFn wdclient.Lookup
|
||||
return m.Chunks, nil
|
||||
}
|
||||
|
||||
func fetchWholeChunk(ctx context.Context, bytesBuffer *bytes.Buffer, lookupFileIdFn wdclient.LookupFileIdFunctionType, fileId string, cipherKey []byte, isGzipped bool, invalidator CacheInvalidator) error {
|
||||
func fetchWholeChunk(ctx context.Context, bytesBuffer *bytes.Buffer, lookupFileIdFn wdclient.LookupFileIdFunctionType, fileId string, cipherKey []byte, isGzipped bool, invalidator CacheInvalidator, filerJwtFn security.FilerJwtProvider) error {
|
||||
urlStrings, err := lookupFileIdFn(ctx, fileId)
|
||||
if err != nil {
|
||||
glog.ErrorfCtx(ctx, "operation LookupFileId %s failed, err: %v", fileId, err)
|
||||
return err
|
||||
}
|
||||
jwt := ChunkReadJwt(urlStrings, fileId)
|
||||
var jwt string
|
||||
if filerJwtFn != nil && len(urlStrings) > 0 && util_http.IsProxyChunkUrl(urlStrings[0]) {
|
||||
jwt = string(filerJwtFn(false))
|
||||
} else {
|
||||
jwt = ChunkReadJwt(urlStrings, fileId)
|
||||
}
|
||||
if _, err = retriedStreamFetchChunkData(ctx, bytesBuffer, urlStrings, jwt, cipherKey, isGzipped, true, 0, 0, refreshUrls(ctx, invalidator, lookupFileIdFn, fileId)); err == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -23,7 +23,7 @@ func BenchmarkManifestResolutionRepeatedOpen(b *testing.B) {
|
||||
}
|
||||
chunk := newManifestCacheTestChunk("benchmark-cached")
|
||||
cache := NewChunkManifestCache(MaxMountChunkManifestCacheEntries, MaxMountChunkManifestCacheBytes)
|
||||
_, err := resolveOneChunkManifest(context.Background(), lookup, chunk, nil, cache)
|
||||
_, err := resolveOneChunkManifest(context.Background(), lookup, chunk, nil, cache, nil)
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
@@ -33,7 +33,7 @@ func BenchmarkManifestResolutionRepeatedOpen(b *testing.B) {
|
||||
b.ReportAllocs()
|
||||
b.ResetTimer()
|
||||
for i := 0; i < b.N; i++ {
|
||||
if _, err := resolveOneChunkManifest(context.Background(), lookup, chunk, nil, cache); err != nil {
|
||||
if _, err := resolveOneChunkManifest(context.Background(), lookup, chunk, nil, cache, nil); err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -258,7 +258,7 @@ func TestResolveOneChunkManifestHonorsCanceledContextOnCacheHit(t *testing.T) {
|
||||
return nil, errors.New("lookup should not be called")
|
||||
}
|
||||
|
||||
_, err := resolveOneChunkManifest(ctx, lookup, chunk, nil, cache)
|
||||
_, err := resolveOneChunkManifest(ctx, lookup, chunk, nil, cache, nil)
|
||||
require.ErrorIs(t, err, context.Canceled)
|
||||
require.False(t, lookupCalled, "a canceled cache hit must not issue a lookup")
|
||||
}
|
||||
@@ -278,7 +278,7 @@ func TestResolveOneChunkManifestCanceledWaiterReturnsDuringCoalescedMiss(t *test
|
||||
defer leaderCancel()
|
||||
leaderDone := make(chan error, 1)
|
||||
go func() {
|
||||
_, err := resolveOneChunkManifest(leaderCtx, fixture.lookup, chunk, nil, cache)
|
||||
_, err := resolveOneChunkManifest(leaderCtx, fixture.lookup, chunk, nil, cache, nil)
|
||||
leaderDone <- err
|
||||
}()
|
||||
|
||||
@@ -287,7 +287,7 @@ func TestResolveOneChunkManifestCanceledWaiterReturnsDuringCoalescedMiss(t *test
|
||||
// for the leader's result.
|
||||
waiterCtx, waiterCancel := context.WithCancel(context.Background())
|
||||
waiterCancel()
|
||||
_, err := resolveOneChunkManifest(waiterCtx, fixture.lookup, chunk, nil, cache)
|
||||
_, err := resolveOneChunkManifest(waiterCtx, fixture.lookup, chunk, nil, cache, nil)
|
||||
require.ErrorIs(t, err, context.Canceled)
|
||||
|
||||
// The leader must still complete successfully and populate the cache.
|
||||
|
||||
@@ -16,6 +16,8 @@ import (
|
||||
"google.golang.org/protobuf/proto"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/security"
|
||||
util_http "github.com/seaweedfs/seaweedfs/weed/util/http"
|
||||
)
|
||||
|
||||
func TestDoMaybeManifestize(t *testing.T) {
|
||||
@@ -621,7 +623,7 @@ func TestFetchWholeChunkRetriesFreshLocations(t *testing.T) {
|
||||
inv := &countingInvalidator{}
|
||||
bytesBuffer := fetchManifestBuffer(t)
|
||||
|
||||
assert.NoError(t, fetchWholeChunk(context.Background(), bytesBuffer, lookup.lookup, "5,stale", nil, false, inv))
|
||||
assert.NoError(t, fetchWholeChunk(context.Background(), bytesBuffer, lookup.lookup, "5,stale", nil, false, inv, nil))
|
||||
assert.Equal(t, int32(1), inv.invalidations.Load())
|
||||
assert.Equal(t, int32(2), lookup.calls.Load())
|
||||
|
||||
@@ -655,7 +657,7 @@ func TestFetchWholeChunkRefreshesLocationsAfterPartialFailure(t *testing.T) {
|
||||
inv := &countingInvalidator{}
|
||||
bytesBuffer := fetchManifestBuffer(t)
|
||||
|
||||
assert.NoError(t, fetchWholeChunk(context.Background(), bytesBuffer, lookup.lookup, "5,abc", nil, false, inv))
|
||||
assert.NoError(t, fetchWholeChunk(context.Background(), bytesBuffer, lookup.lookup, "5,abc", nil, false, inv, nil))
|
||||
assert.Equal(t, int32(1), inv.invalidations.Load())
|
||||
assert.Equal(t, int32(2), lookup.calls.Load())
|
||||
decoded := &filer_pb.FileChunkManifest{}
|
||||
@@ -676,7 +678,7 @@ func TestFetchWholeChunkWithoutInvalidator(t *testing.T) {
|
||||
freshUrls: []string{"http://unused:8080/5,abc"},
|
||||
}
|
||||
|
||||
assert.Error(t, fetchWholeChunk(context.Background(), fetchManifestBuffer(t), lookup.lookup, "5,abc", nil, false, nil))
|
||||
assert.Error(t, fetchWholeChunk(context.Background(), fetchManifestBuffer(t), lookup.lookup, "5,abc", nil, false, nil, nil))
|
||||
assert.Equal(t, int32(1), lookup.calls.Load())
|
||||
}
|
||||
|
||||
@@ -694,7 +696,7 @@ func TestFetchWholeChunkUnchangedLocations(t *testing.T) {
|
||||
}
|
||||
inv := &countingInvalidator{}
|
||||
|
||||
assert.Error(t, fetchWholeChunk(context.Background(), fetchManifestBuffer(t), lookup.lookup, "5,abc", nil, false, inv))
|
||||
assert.Error(t, fetchWholeChunk(context.Background(), fetchManifestBuffer(t), lookup.lookup, "5,abc", nil, false, inv, nil))
|
||||
assert.Equal(t, int32(2), lookup.calls.Load())
|
||||
assert.Equal(t, int32(1), inv.invalidations.Load())
|
||||
}
|
||||
@@ -711,7 +713,7 @@ func TestFetchWholeChunkCancelledKeepsLocations(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
cancel()
|
||||
|
||||
err := fetchWholeChunk(ctx, fetchManifestBuffer(t), lookup.lookup, "5,abc", nil, false, inv)
|
||||
err := fetchWholeChunk(ctx, fetchManifestBuffer(t), lookup.lookup, "5,abc", nil, false, inv, nil)
|
||||
assert.ErrorIs(t, err, context.Canceled)
|
||||
assert.Equal(t, int32(0), inv.invalidations.Load())
|
||||
assert.Equal(t, int32(1), lookup.calls.Load())
|
||||
@@ -730,3 +732,43 @@ func TestFetchWholeChunkCancelledKeepsLocations(t *testing.T) {
|
||||
})
|
||||
assert.ErrorIs(t, noInvalidator, context.Canceled)
|
||||
}
|
||||
|
||||
// TestFetchWholeChunkUsesProvidedFilerJwt covers a replicating reader whose
|
||||
// source filer authenticates proxied downloads with its own read key: the
|
||||
// supplied provider's token must reach the server, not the process-wide one.
|
||||
func TestFetchWholeChunkUsesProvidedFilerJwt(t *testing.T) {
|
||||
manifestBytes, err := proto.Marshal(&filer_pb.FileChunkManifest{
|
||||
Chunks: []*filer_pb.FileChunk{{FileId: "100,abc", Offset: 0, Size: 8}},
|
||||
})
|
||||
assert.NoError(t, err)
|
||||
|
||||
gotAuth := make(chan string, 1)
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
gotAuth <- r.Header.Get("Authorization")
|
||||
w.Header().Set("Content-Length", strconv.Itoa(len(manifestBytes)))
|
||||
w.Write(manifestBytes)
|
||||
}))
|
||||
t.Cleanup(srv.Close)
|
||||
|
||||
lookup := func(ctx context.Context, fileId string) ([]string, error) {
|
||||
return []string{srv.URL + "/?" + util_http.ProxyChunkIdParam + "=" + fileId}, nil
|
||||
}
|
||||
jwtFn := func(isWrite bool) security.EncodedJwt {
|
||||
assert.False(t, isWrite)
|
||||
return "side-read-jwt"
|
||||
}
|
||||
bytesBuffer := fetchManifestBuffer(t)
|
||||
assert.NoError(t, fetchWholeChunk(context.Background(), bytesBuffer, lookup, "5,abc", nil, false, nil, jwtFn))
|
||||
assert.Equal(t, security.BearerPrefix+"side-read-jwt", <-gotAuth)
|
||||
|
||||
// non-proxy URLs keep the volume-server credential and never call the provider
|
||||
volumeURL := manifestServer(t, manifestBytes).URL + "/5,abc"
|
||||
volumeLookup := func(ctx context.Context, fileId string) ([]string, error) {
|
||||
return []string{volumeURL}, nil
|
||||
}
|
||||
bytesBuffer.Reset()
|
||||
assert.NoError(t, fetchWholeChunk(context.Background(), bytesBuffer, volumeLookup, "5,abc", nil, false, nil, func(bool) security.EncodedJwt {
|
||||
t.Fatal("provider must not be consulted for a volume url")
|
||||
return ""
|
||||
}))
|
||||
}
|
||||
|
||||
@@ -190,7 +190,7 @@ func loadLogFileEntries(masterClient *wdclient.MasterClient, chunk *filer_pb.Fil
|
||||
lookupFileIdFn := func(ctx context.Context, fileId string) (targetUrls []string, err error) {
|
||||
return masterClient.LookupFileId(ctx, fileId)
|
||||
}
|
||||
if fetchErr := fetchWholeChunk(context.Background(), bytesBuffer, lookupFileIdFn, chunk.GetFileIdString(), chunk.CipherKey, chunk.IsCompressed, masterClient); fetchErr != nil {
|
||||
if fetchErr := fetchWholeChunk(context.Background(), bytesBuffer, lookupFileIdFn, chunk.GetFileIdString(), chunk.CipherKey, chunk.IsCompressed, masterClient, nil); fetchErr != nil {
|
||||
return nil, false, fetchErr
|
||||
}
|
||||
return decodeLogRecords(bytesBuffer.Bytes())
|
||||
|
||||
@@ -62,6 +62,7 @@ type UploadOption struct {
|
||||
SourceUrl string // optional: for logging when reading from a remote source
|
||||
MaxAttempts int // <=0 uses the default
|
||||
GenUploadUrl func(host, fileId string) string // if nil → fallback "http://{host}/{fileId}"
|
||||
FilerJwt security.FilerJwtProvider // credential for proxy chunk URLs; nil → process-wide jwt.filer_signing
|
||||
}
|
||||
|
||||
type UploadResult struct {
|
||||
@@ -217,7 +218,11 @@ func (uploader *Uploader) uploadWithRetryData(assignFn func() (fileId string, ho
|
||||
// The request addresses the filer, which authorizes it and mints the
|
||||
// volume credential itself. The AssignVolume token is not a filer
|
||||
// credential and gets the caller nowhere here.
|
||||
uploadOption.Jwt = security.EncodedJwt(util_http.JwtForFilerServer(true))
|
||||
if uploadOption.FilerJwt != nil {
|
||||
uploadOption.Jwt = uploadOption.FilerJwt(true)
|
||||
} else {
|
||||
uploadOption.Jwt = security.EncodedJwt(util_http.JwtForFilerServer(true))
|
||||
}
|
||||
}
|
||||
|
||||
uploadResult, err = uploader.retriedUploadData(context.Background(), data, uploadOption)
|
||||
|
||||
@@ -380,6 +380,7 @@ message ErasureCodingTaskConfig {
|
||||
string collection_filter = 4; // Only process volumes from specific collections
|
||||
repeated string preferred_tags = 5; // Disk tags to prioritize for EC shard placement
|
||||
string replica_placement = 6; // EC shard replica placement (e.g. "020"); empty falls back to master default replication
|
||||
bool strict_placement = 7; // fail planning instead of relaxing placement constraints
|
||||
}
|
||||
|
||||
// BalanceTaskConfig contains balance-specific configuration
|
||||
|
||||
@@ -2980,6 +2980,7 @@ type ErasureCodingTaskConfig struct {
|
||||
CollectionFilter string `protobuf:"bytes,4,opt,name=collection_filter,json=collectionFilter,proto3" json:"collection_filter,omitempty"` // Only process volumes from specific collections
|
||||
PreferredTags []string `protobuf:"bytes,5,rep,name=preferred_tags,json=preferredTags,proto3" json:"preferred_tags,omitempty"` // Disk tags to prioritize for EC shard placement
|
||||
ReplicaPlacement string `protobuf:"bytes,6,opt,name=replica_placement,json=replicaPlacement,proto3" json:"replica_placement,omitempty"` // EC shard replica placement (e.g. "020"); empty falls back to master default replication
|
||||
StrictPlacement bool `protobuf:"varint,7,opt,name=strict_placement,json=strictPlacement,proto3" json:"strict_placement,omitempty"` // fail planning instead of relaxing placement constraints
|
||||
unknownFields protoimpl.UnknownFields
|
||||
sizeCache protoimpl.SizeCache
|
||||
}
|
||||
@@ -3056,6 +3057,13 @@ func (x *ErasureCodingTaskConfig) GetReplicaPlacement() string {
|
||||
return ""
|
||||
}
|
||||
|
||||
func (x *ErasureCodingTaskConfig) GetStrictPlacement() bool {
|
||||
if x != nil {
|
||||
return x.StrictPlacement
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// BalanceTaskConfig contains balance-specific configuration
|
||||
type BalanceTaskConfig struct {
|
||||
state protoimpl.MessageState `protogen:"open.v1"`
|
||||
@@ -4285,14 +4293,15 @@ const file_worker_proto_rawDesc = "" +
|
||||
"\x10VacuumTaskConfig\x12+\n" +
|
||||
"\x11garbage_threshold\x18\x01 \x01(\x01R\x10garbageThreshold\x12/\n" +
|
||||
"\x14min_volume_age_hours\x18\x02 \x01(\x05R\x11minVolumeAgeHours\x120\n" +
|
||||
"\x14min_interval_seconds\x18\x03 \x01(\x05R\x12minIntervalSeconds\"\x9a\x02\n" +
|
||||
"\x14min_interval_seconds\x18\x03 \x01(\x05R\x12minIntervalSeconds\"\xc5\x02\n" +
|
||||
"\x17ErasureCodingTaskConfig\x12%\n" +
|
||||
"\x0efullness_ratio\x18\x01 \x01(\x01R\rfullnessRatio\x12*\n" +
|
||||
"\x11quiet_for_seconds\x18\x02 \x01(\x05R\x0fquietForSeconds\x12+\n" +
|
||||
"\x12min_volume_size_mb\x18\x03 \x01(\x05R\x0fminVolumeSizeMb\x12+\n" +
|
||||
"\x11collection_filter\x18\x04 \x01(\tR\x10collectionFilter\x12%\n" +
|
||||
"\x0epreferred_tags\x18\x05 \x03(\tR\rpreferredTags\x12+\n" +
|
||||
"\x11replica_placement\x18\x06 \x01(\tR\x10replicaPlacement\"\x9b\x01\n" +
|
||||
"\x11replica_placement\x18\x06 \x01(\tR\x10replicaPlacement\x12)\n" +
|
||||
"\x10strict_placement\x18\a \x01(\bR\x0fstrictPlacement\"\x9b\x01\n" +
|
||||
"\x11BalanceTaskConfig\x12/\n" +
|
||||
"\x13imbalance_threshold\x18\x01 \x01(\x01R\x12imbalanceThreshold\x12(\n" +
|
||||
"\x10min_server_count\x18\x02 \x01(\x05R\x0eminServerCount\x12+\n" +
|
||||
|
||||
@@ -125,6 +125,21 @@ func ReadIntConfig(values map[string]*plugin_pb.ConfigValue, field string, fallb
|
||||
return int(v)
|
||||
}
|
||||
|
||||
// ReadBoolConfig reads a bool-valued plugin config field.
|
||||
func ReadBoolConfig(values map[string]*plugin_pb.ConfigValue, field string, fallback bool) bool {
|
||||
if values == nil {
|
||||
return fallback
|
||||
}
|
||||
value := values[field]
|
||||
if value == nil {
|
||||
return fallback
|
||||
}
|
||||
if kind, ok := value.Kind.(*plugin_pb.ConfigValue_BoolValue); ok {
|
||||
return kind.BoolValue
|
||||
}
|
||||
return fallback
|
||||
}
|
||||
|
||||
// ReadBytesConfig reads a bytes-valued plugin config field, returning nil when
|
||||
// the value is missing or of a different kind.
|
||||
func ReadBytesConfig(values map[string]*plugin_pb.ConfigValue, field string) []byte {
|
||||
|
||||
@@ -195,7 +195,7 @@ func (fs *FilerSink) replicateOneManifestChunk(ctx context.Context, sourceChunk
|
||||
resolveName := fmt.Sprintf("resolve manifest %s", sourceChunk.GetFileIdString())
|
||||
missingGate := fs.newMissingSourceChunkGate(sourceChunk.GetFileIdString())
|
||||
err := util.RetryUntil(resolveName, func() error {
|
||||
rc, e := filer.ResolveOneChunkManifest(ctx, fs.filerSource.LookupFileId, sourceChunk, nil)
|
||||
rc, e := filer.ResolveOneChunkManifestWithFilerJwt(ctx, fs.filerSource.LookupFileId, fs.filerSource.FilerJwt(), sourceChunk, nil)
|
||||
if e != nil {
|
||||
return e
|
||||
}
|
||||
@@ -315,6 +315,7 @@ func (fs *FilerSink) uploadManifestChunk(path string, sourceMtimeNs int64, sourc
|
||||
}
|
||||
if fs.writeChunkByFiler {
|
||||
uploadOption.GenUploadUrl = operation.GenUploadUrlProxy(fs.address)
|
||||
uploadOption.FilerJwt = fs.jwtForFiler
|
||||
}
|
||||
currentFileId, uploadResult, uploadErr, _ := uploader.UploadWithRetry(
|
||||
fs,
|
||||
@@ -435,6 +436,7 @@ func (fs *FilerSink) fetchAndWrite(sourceChunk *filer_pb.FileChunk, path string,
|
||||
}
|
||||
if fs.writeChunkByFiler {
|
||||
uploadOption.GenUploadUrl = operation.GenUploadUrlProxy(fs.address)
|
||||
uploadOption.FilerJwt = fs.jwtForFiler
|
||||
}
|
||||
currentFileId, uploadResult, uploadErr, _ := uploader.UploadWithRetry(
|
||||
fs,
|
||||
|
||||
@@ -62,6 +62,7 @@ type FilerSink struct {
|
||||
signature int32
|
||||
activeTransfers sync.Map // chunkFileId -> *ChunkTransferStatus
|
||||
uploader *operation.Uploader
|
||||
jwtForFiler security.FilerJwtProvider
|
||||
// lastServedFileId is the most recent chunk the source did serve, the probe
|
||||
// sourceStillServesChunks re-checks before writing an entry off.
|
||||
lastServedFileId atomic.Pointer[string]
|
||||
@@ -116,6 +117,12 @@ func (fs *FilerSink) SetUploader(uploader *operation.Uploader) {
|
||||
fs.uploader = uploader
|
||||
}
|
||||
|
||||
// SetFilerJwtProvider sets the filer API credential for proxied chunk writes.
|
||||
// Must be called during initialization, before any replication goroutines start.
|
||||
func (fs *FilerSink) SetFilerJwtProvider(provider security.FilerJwtProvider) {
|
||||
fs.jwtForFiler = provider
|
||||
}
|
||||
|
||||
func (fs *FilerSink) getUploader() (*operation.Uploader, error) {
|
||||
if fs.uploader != nil {
|
||||
return fs.uploader, nil
|
||||
@@ -358,7 +365,7 @@ func (fs *FilerSink) UpdateEntry(key string, oldEntry *filer_pb.Entry, newParent
|
||||
existingEntry.RemoteEntry = newEntry.RemoteEntry
|
||||
default:
|
||||
// source-side chunks resolve via source filer; sink volume IDs may collide.
|
||||
deletedChunks, newChunks, err := compareChunks(context.Background(), filer.LookupFn(fs.filerSource), oldEntry, newEntry)
|
||||
deletedChunks, newChunks, err := compareChunks(context.Background(), filer.LookupFn(fs.filerSource), fs.filerSource.FilerJwt(), oldEntry, newEntry)
|
||||
if err != nil {
|
||||
return true, fmt.Errorf("replicate %s compare chunks error: %w", key, err)
|
||||
}
|
||||
@@ -404,12 +411,12 @@ func (fs *FilerSink) UpdateEntry(key string, oldEntry *filer_pb.Entry, newParent
|
||||
})
|
||||
|
||||
}
|
||||
func compareChunks(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, oldEntry, newEntry *filer_pb.Entry) (deletedChunks, newChunks []*filer_pb.FileChunk, err error) {
|
||||
aData, aMeta, aErr := filer.ResolveChunkManifest(ctx, lookupFileIdFn, oldEntry.GetChunks(), 0, math.MaxInt64, nil)
|
||||
func compareChunks(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, filerJwtFn security.FilerJwtProvider, oldEntry, newEntry *filer_pb.Entry) (deletedChunks, newChunks []*filer_pb.FileChunk, err error) {
|
||||
aData, aMeta, aErr := filer.ResolveChunkManifestWithFilerJwt(ctx, lookupFileIdFn, filerJwtFn, oldEntry.GetChunks(), 0, math.MaxInt64, nil)
|
||||
if aErr != nil {
|
||||
return nil, nil, aErr
|
||||
}
|
||||
bData, bMeta, bErr := filer.ResolveChunkManifest(ctx, lookupFileIdFn, newEntry.GetChunks(), 0, math.MaxInt64, nil)
|
||||
bData, bMeta, bErr := filer.ResolveChunkManifestWithFilerJwt(ctx, lookupFileIdFn, filerJwtFn, newEntry.GetChunks(), 0, math.MaxInt64, nil)
|
||||
if bErr != nil {
|
||||
return nil, nil, bErr
|
||||
}
|
||||
|
||||
@@ -34,6 +34,7 @@ type FilerSource struct {
|
||||
dataCenter string
|
||||
signature int32
|
||||
httpClient *util_http_client.HTTPClient
|
||||
jwtForFiler security.FilerJwtProvider
|
||||
}
|
||||
|
||||
func (fs *FilerSource) Initialize(configuration util.Configuration, prefix string) error {
|
||||
@@ -67,6 +68,16 @@ func (fs *FilerSource) SetHttpClient(client *util_http_client.HTTPClient) {
|
||||
fs.httpClient = client
|
||||
}
|
||||
|
||||
func (fs *FilerSource) SetFilerJwtProvider(provider security.FilerJwtProvider) {
|
||||
fs.jwtForFiler = provider
|
||||
}
|
||||
|
||||
// FilerJwt returns the side-specific filer API credential, or nil when the
|
||||
// process-wide jwt.filer_signing configuration applies.
|
||||
func (fs *FilerSource) FilerJwt() security.FilerJwtProvider {
|
||||
return fs.jwtForFiler
|
||||
}
|
||||
|
||||
func (fs *FilerSource) LookupFileId(ctx context.Context, part string) (fileUrls []string, err error) {
|
||||
|
||||
vid2Locations := make(map[string]*filer_pb.Locations)
|
||||
@@ -126,7 +137,11 @@ func (fs *FilerSource) ReadPart(fileId string, offset int64) (filename string, h
|
||||
|
||||
if fs.proxyByFiler {
|
||||
fileUrl := util_http.ProxyChunkUrl(fs.address, fileId)
|
||||
filename, header, resp, err = downloadFn(fileUrl, util_http.JwtForFilerServer(false), offset)
|
||||
jwt := util_http.JwtForFilerServer(false)
|
||||
if fs.jwtForFiler != nil {
|
||||
jwt = string(fs.jwtForFiler(false))
|
||||
}
|
||||
filename, header, resp, err = downloadFn(fileUrl, jwt, offset)
|
||||
if err == nil {
|
||||
err = readPartStatusError(fileUrl, resp)
|
||||
}
|
||||
|
||||
@@ -8,6 +8,8 @@ import (
|
||||
|
||||
jwt "github.com/golang-jwt/jwt/v5"
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
"github.com/spf13/viper"
|
||||
)
|
||||
|
||||
type EncodedJwt string
|
||||
@@ -151,3 +153,57 @@ func DecodeJwt(signingKey SigningKey, tokenString EncodedJwt, claims jwt.Claims)
|
||||
return []byte(signingKey), nil
|
||||
})
|
||||
}
|
||||
|
||||
// FilerJwtProvider signs the credential a filer's HTTP API expects. A nil
|
||||
// provider means the process-wide jwt.filer_signing configuration applies.
|
||||
type FilerJwtProvider func(isWrite bool) EncodedJwt
|
||||
|
||||
// LoadFilerJwtFromFile reads jwt.filer_signing from a security file the way
|
||||
// LoadClientTLSFromFile reads the TLS section, honoring the same WEED_
|
||||
// environment precedence. A nil provider means the file configures no filer
|
||||
// signing keys and the process-wide configuration applies. A file that sets
|
||||
// only one access level falls back to the process-wide key for the other.
|
||||
func LoadFilerJwtFromFile(configFile string) (FilerJwtProvider, error) {
|
||||
v := viper.New()
|
||||
v.SetConfigFile(configFile)
|
||||
v.AutomaticEnv()
|
||||
v.SetEnvPrefix("weed")
|
||||
v.SetEnvKeyReplacer(strings.NewReplacer(".", "_"))
|
||||
if err := v.ReadInConfig(); err != nil {
|
||||
return nil, fmt.Errorf("failed to read security config %s: %v", configFile, err)
|
||||
}
|
||||
|
||||
signingKey := SigningKey(v.GetString("jwt.filer_signing.key"))
|
||||
readSigningKey := SigningKey(v.GetString("jwt.filer_signing.read.key"))
|
||||
if len(signingKey) == 0 && len(readSigningKey) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
signingKeyExpires := v.GetInt("jwt.filer_signing.expires_after_seconds")
|
||||
readSigningKeyExpires := v.GetInt("jwt.filer_signing.read.expires_after_seconds")
|
||||
if len(signingKey) == 0 || len(readSigningKey) == 0 {
|
||||
gv := util.GetViper()
|
||||
if len(signingKey) == 0 {
|
||||
signingKey = SigningKey(gv.GetString("jwt.filer_signing.key"))
|
||||
signingKeyExpires = gv.GetInt("jwt.filer_signing.expires_after_seconds")
|
||||
}
|
||||
if len(readSigningKey) == 0 {
|
||||
readSigningKey = SigningKey(gv.GetString("jwt.filer_signing.read.key"))
|
||||
readSigningKeyExpires = gv.GetInt("jwt.filer_signing.read.expires_after_seconds")
|
||||
}
|
||||
}
|
||||
if signingKeyExpires < 0 || readSigningKeyExpires < 0 {
|
||||
return nil, fmt.Errorf("jwt.filer_signing lifetimes must not be negative")
|
||||
}
|
||||
if signingKeyExpires == 0 {
|
||||
signingKeyExpires = 10
|
||||
}
|
||||
if readSigningKeyExpires == 0 {
|
||||
readSigningKeyExpires = 60
|
||||
}
|
||||
return func(isWrite bool) EncodedJwt {
|
||||
if isWrite {
|
||||
return GenJwtForFilerServer(signingKey, signingKeyExpires)
|
||||
}
|
||||
return GenJwtForFilerServer(readSigningKey, readSigningKeyExpires)
|
||||
}, nil
|
||||
}
|
||||
|
||||
@@ -0,0 +1,146 @@
|
||||
package security
|
||||
|
||||
import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
)
|
||||
|
||||
func TestLoadFilerJwtFromFile(t *testing.T) {
|
||||
configFile := filepath.Join(t.TempDir(), "security.toml")
|
||||
config := `
|
||||
[jwt]
|
||||
[jwt.filer_signing]
|
||||
key = "side-write-key"
|
||||
expires_after_seconds = 30
|
||||
[jwt.filer_signing.read]
|
||||
key = "side-read-key"
|
||||
expires_after_seconds = 90
|
||||
`
|
||||
if err := os.WriteFile(configFile, []byte(config), 0644); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
provider, err := LoadFilerJwtFromFile(configFile)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if provider == nil {
|
||||
t.Fatal("config with filer signing keys gave a nil provider")
|
||||
}
|
||||
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
isWrite bool
|
||||
signedBy string
|
||||
otherKey string
|
||||
expires int64
|
||||
}{
|
||||
{"read", false, "side-read-key", "side-write-key", 90},
|
||||
{"write", true, "side-write-key", "side-read-key", 30},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
before := time.Now()
|
||||
token := provider(tc.isWrite)
|
||||
after := time.Now()
|
||||
claims := &SeaweedFilerClaims{}
|
||||
if _, err := DecodeJwt(SigningKey(tc.signedBy), token, claims); err != nil {
|
||||
t.Fatalf("token does not validate against the %s key: %v", tc.name, err)
|
||||
}
|
||||
if claims.ExpiresAt == nil {
|
||||
t.Fatal("token never expires")
|
||||
}
|
||||
expiresIn := claims.ExpiresAt.Time
|
||||
if expiresIn.Before(before.Add(time.Duration(tc.expires-1)*time.Second)) || expiresIn.After(after.Add(time.Duration(tc.expires+1)*time.Second)) {
|
||||
t.Fatalf("token expires at %v, want %ds after %v", expiresIn, tc.expires, before)
|
||||
}
|
||||
if _, err := DecodeJwt(SigningKey(tc.otherKey), token, &SeaweedFilerClaims{}); err == nil {
|
||||
t.Fatal("token also validates against the other access level's key")
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadFilerJwtFromFileWithoutKeys(t *testing.T) {
|
||||
configFile := filepath.Join(t.TempDir(), "security.toml")
|
||||
if err := os.WriteFile(configFile, []byte("[grpc.client]\n"), 0644); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
provider, err := LoadFilerJwtFromFile(configFile)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if provider != nil {
|
||||
t.Fatal("config without filer signing keys gave a provider")
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadFilerJwtFromFileMissing(t *testing.T) {
|
||||
if _, err := LoadFilerJwtFromFile(filepath.Join(t.TempDir(), "none.toml")); err == nil {
|
||||
t.Fatal("missing config file loaded without error")
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadFilerJwtFromFilePartialKeys(t *testing.T) {
|
||||
configFile := filepath.Join(t.TempDir(), "security.toml")
|
||||
config := `
|
||||
[jwt.filer_signing.read]
|
||||
key = "side-read-key"
|
||||
`
|
||||
if err := os.WriteFile(configFile, []byte(config), 0644); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
gv := util.GetViper()
|
||||
priorKey := gv.GetString("jwt.filer_signing.key")
|
||||
gv.Set("jwt.filer_signing.key", "global-write-key")
|
||||
t.Cleanup(func() { gv.Set("jwt.filer_signing.key", priorKey) })
|
||||
|
||||
provider, err := LoadFilerJwtFromFile(configFile)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if provider == nil {
|
||||
t.Fatal("config with a filer signing key gave a nil provider")
|
||||
}
|
||||
|
||||
if _, err := DecodeJwt(SigningKey("global-write-key"), provider(true), &SeaweedFilerClaims{}); err != nil {
|
||||
t.Fatalf("write token does not validate against the global write key: %v", err)
|
||||
}
|
||||
if _, err := DecodeJwt(SigningKey("side-read-key"), provider(false), &SeaweedFilerClaims{}); err != nil {
|
||||
t.Fatalf("read token does not validate against the side read key: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadFilerJwtFromFileEnvOverride(t *testing.T) {
|
||||
configFile := filepath.Join(t.TempDir(), "security.toml")
|
||||
config := `
|
||||
[jwt.filer_signing]
|
||||
key = "side-write-key"
|
||||
[jwt.filer_signing.read]
|
||||
key = "side-read-key"
|
||||
`
|
||||
if err := os.WriteFile(configFile, []byte(config), 0644); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Setenv("WEED_JWT_FILER_SIGNING_READ_KEY", "env-read-key")
|
||||
|
||||
provider, err := LoadFilerJwtFromFile(configFile)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if provider == nil {
|
||||
t.Fatal("config with filer signing keys gave a nil provider")
|
||||
}
|
||||
|
||||
if _, err := DecodeJwt(SigningKey("env-read-key"), provider(false), &SeaweedFilerClaims{}); err != nil {
|
||||
t.Fatalf("read token does not validate against the env override key: %v", err)
|
||||
}
|
||||
if _, err := DecodeJwt(SigningKey("side-write-key"), provider(true), &SeaweedFilerClaims{}); err != nil {
|
||||
t.Fatalf("write token does not validate against the side write key: %v", err)
|
||||
}
|
||||
}
|
||||
@@ -27,6 +27,11 @@ var (
|
||||
// (possibly-unflushed) gap, in case the flush notification is missed.
|
||||
unflushedGapRetryInterval = 2 * time.Second
|
||||
|
||||
// aggDiskReprobeInterval paces the aggregated persisted-log re-listing for
|
||||
// files no watermark signals: a peer past the flush low-watermark can land
|
||||
// a file without moving the minimum.
|
||||
aggDiskReprobeInterval = 2 * time.Second
|
||||
|
||||
// gapStallWarnInterval paces the warning for a subscriber that stays parked.
|
||||
gapStallWarnInterval = time.Minute
|
||||
|
||||
@@ -666,8 +671,9 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
|
||||
var lastHeartbeatNs int64
|
||||
baseEachLogEntryFn := eachLogEntryFn(req, sender, eachEventNotificationFn, &unsyncedEvents)
|
||||
// heldAtTsNs remembers the entry a read was held at (for the log line);
|
||||
// the rewind target is the last entry actually delivered.
|
||||
var heldAtTsNs int64
|
||||
// diskHeldAtTsNs is the same marker for the disk pass alone: a pending
|
||||
// disk hold keeps the pass re-reading until the entry is served.
|
||||
var heldAtTsNs, diskHeldAtTsNs int64
|
||||
// Each read path holds at its own watermark: persisted logs are complete
|
||||
// only up to every peer's flush watermark, the ring only up to every
|
||||
// peer's delivery watermark.
|
||||
@@ -698,7 +704,14 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
|
||||
// What the last disk pass proved covered: flushed on every peer AND inside
|
||||
// the pass's listing, so an empty pass proves (cursor, proven] empty.
|
||||
var diskPassProvenTsNs int64
|
||||
diskEachLogEntryFn := guardedEachLogEntryFn(func() int64 { return diskPassHoldTsNs })
|
||||
diskBaseEachLogEntryFn := guardedEachLogEntryFn(func() int64 { return diskPassHoldTsNs })
|
||||
diskEachLogEntryFn := func(logEntry *filer_pb.LogEntry) (bool, error) {
|
||||
isDone, err := diskBaseEachLogEntryFn(logEntry)
|
||||
if errors.Is(err, errHeldByPeerWatermark) {
|
||||
diskHeldAtTsNs = logEntry.TsNs
|
||||
}
|
||||
return isDone, err
|
||||
}
|
||||
memEachLogEntryFn := guardedEachLogEntryFn(holdMemTsNs)
|
||||
// waitHeld pauses a held read until a peer reports further progress, or
|
||||
// the retry interval elapses (a peer dropped past its grace, or a log file
|
||||
@@ -733,6 +746,11 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
|
||||
var readPersistedLogErr error
|
||||
var readInMemoryLogErr error
|
||||
var isDone bool
|
||||
var lastCheckedFlushTsNs int64 = -1 // Track the last local flush we read the disk under
|
||||
var lastCheckedFlushLowTsNs int64 = -1 // Track the last peer flush low-watermark we read the disk under
|
||||
var lastDiskReadTsNs int64 = -1 // Track the last read position we used for disk read
|
||||
var lastDiskRefsStopTsNs int64 = -1 // Track the last chunk listing bound we read the disk under
|
||||
var lastDiskPassAt time.Time // Paces re-probes for files no watermark signals
|
||||
sentRefs := make(map[string]sentRefState)
|
||||
|
||||
aggBuffer := fs.filer.MetaAggregator.MetaLogBuffer
|
||||
@@ -765,78 +783,110 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
|
||||
// diskPassHoldTsNs above).
|
||||
diskPassFlushLowTsNs = fs.filer.MetaAggregator.PeerLowFlushWatermarkTsNs()
|
||||
diskPassHoldTsNs = resolveAggReadHoldTsNs(diskPassFlushLowTsNs, time.Now().UnixNano(), metadataGapSettledHorizon)
|
||||
diskPassProvenTsNs = diskPassFlushLowTsNs
|
||||
|
||||
// Re-read the disk only when something changed it cannot miss: a local
|
||||
// flush landed, the peers' flush low-watermark moved in either
|
||||
// direction (a joining peer invalidates what an earlier pass proved),
|
||||
// the cursor moved, or a disk hold is pending (the ring read after an
|
||||
// empty pass parks internally, so skipping would strand the held
|
||||
// entry). A peer past the low-watermark can still land a file without
|
||||
// moving it, so the listing is also re-probed at a slow cadence;
|
||||
// chunk listings re-arm as soon as their read bound admits more files.
|
||||
currentFlushTsNs := fs.filer.LocalMetaLogBuffer.GetLastFlushTsNs()
|
||||
currentReadTsNs := lastReadTime.Time.UnixNano()
|
||||
var currentRefsStopTsNs int64
|
||||
if req.ClientSupportsMetadataChunks {
|
||||
refsStopTsNs := chunkRefsStopTsNs(diskPassHoldTsNs, req.UntilNs)
|
||||
// Nothing above the listing bound is proven by this pass.
|
||||
if refsStopTsNs < diskPassProvenTsNs {
|
||||
diskPassProvenTsNs = refsStopTsNs
|
||||
}
|
||||
if refsStopTsNs > lastReadTime.Time.UnixNano() {
|
||||
processedTsNs, isDone, readPersistedLogErr = fs.chunkDiskPass(ctx, sender, lastReadTime, refsStopTsNs, sentRefs, nil)
|
||||
} else {
|
||||
processedTsNs, isDone, readPersistedLogErr = 0, false, nil
|
||||
}
|
||||
} else {
|
||||
processedTsNs, isDone, readPersistedLogErr = fs.filer.ReadPersistedLogBuffer(ctx, lastReadTime, req.UntilNs, diskEachLogEntryFn)
|
||||
currentRefsStopTsNs = chunkRefsStopTsNs(diskPassHoldTsNs, req.UntilNs)
|
||||
}
|
||||
if errors.Is(readPersistedLogErr, errHeldByPeerWatermark) {
|
||||
// Stay at the last delivered entry; the held entry is re-read (and
|
||||
// re-checked) by the next pass.
|
||||
if processedTsNs > 0 {
|
||||
shouldReadFromDisk := lastCheckedFlushTsNs == -1 ||
|
||||
currentFlushTsNs > lastCheckedFlushTsNs ||
|
||||
diskPassFlushLowTsNs != lastCheckedFlushLowTsNs ||
|
||||
currentReadTsNs != lastDiskReadTsNs ||
|
||||
currentRefsStopTsNs > lastDiskRefsStopTsNs ||
|
||||
diskHeldAtTsNs != 0 ||
|
||||
time.Since(lastDiskPassAt) >= aggDiskReprobeInterval
|
||||
|
||||
diskAdvanced := false
|
||||
if shouldReadFromDisk {
|
||||
lastCheckedFlushTsNs = currentFlushTsNs
|
||||
lastCheckedFlushLowTsNs = diskPassFlushLowTsNs
|
||||
lastDiskReadTsNs = currentReadTsNs
|
||||
lastDiskRefsStopTsNs = currentRefsStopTsNs
|
||||
lastDiskPassAt = time.Now()
|
||||
diskHeldAtTsNs = 0
|
||||
diskPassProvenTsNs = diskPassFlushLowTsNs
|
||||
|
||||
if req.ClientSupportsMetadataChunks {
|
||||
refsStopTsNs := currentRefsStopTsNs
|
||||
// Nothing above the listing bound is proven by this pass.
|
||||
if refsStopTsNs < diskPassProvenTsNs {
|
||||
diskPassProvenTsNs = refsStopTsNs
|
||||
}
|
||||
if refsStopTsNs > lastReadTime.Time.UnixNano() {
|
||||
processedTsNs, isDone, readPersistedLogErr = fs.chunkDiskPass(ctx, sender, lastReadTime, refsStopTsNs, sentRefs, nil)
|
||||
} else {
|
||||
processedTsNs, isDone, readPersistedLogErr = 0, false, nil
|
||||
}
|
||||
} else {
|
||||
processedTsNs, isDone, readPersistedLogErr = fs.filer.ReadPersistedLogBuffer(ctx, lastReadTime, req.UntilNs, diskEachLogEntryFn)
|
||||
}
|
||||
if errors.Is(readPersistedLogErr, errHeldByPeerWatermark) {
|
||||
// Stay at the last delivered entry; the held entry is re-read (and
|
||||
// re-checked) by the next pass.
|
||||
if processedTsNs > 0 {
|
||||
lastReadTime = log_buffer.NewMessagePosition(processedTsNs, gapResumeCursorOffset)
|
||||
if processedTsNs > diskAnchorTsNs {
|
||||
diskAnchorTsNs = processedTsNs
|
||||
}
|
||||
}
|
||||
// A hold is not a gap: clear any stale ResumeFromDiskError so the
|
||||
// next pass's disk-miss handling cannot skip past the held entry.
|
||||
readInMemoryLogErr = nil
|
||||
if !waitHeld("disk", flushChan) {
|
||||
return nil
|
||||
}
|
||||
continue
|
||||
}
|
||||
if readPersistedLogErr != nil {
|
||||
return fmt.Errorf("reading from persisted logs: %w", readPersistedLogErr)
|
||||
}
|
||||
if isDone {
|
||||
return nil
|
||||
}
|
||||
|
||||
glog.V(4).Infof("processed to %v: %v", clientName, processedTsNs)
|
||||
diskAdvanced = diskReadAdvanced(processedTsNs, lastReadTime)
|
||||
// Read after the disk read (an eviction landing mid-read must count) and
|
||||
// in received-ts space: the ring's bumped stopTimes exceed anything on
|
||||
// any peer's disk, and gating disk cursors on them parks subscribers
|
||||
// that drained every peer's log.
|
||||
lastEvictedTsNs := fs.filer.MetaAggregator.MetaLogBuffer.GetLastEvictedOriginalTsNs()
|
||||
if diskAdvanced {
|
||||
gapStall.resumed()
|
||||
reportUnprovenAggregatedCrossing(cursorBeforeDiskTsNs, processedTsNs, lastEvictedTsNs, diskPassFlushLowTsNs, clientName, req.PathPrefix)
|
||||
lastReadTime = log_buffer.NewMessagePosition(processedTsNs, gapResumeCursorOffset)
|
||||
if processedTsNs > diskAnchorTsNs {
|
||||
diskAnchorTsNs = processedTsNs
|
||||
}
|
||||
}
|
||||
// A hold is not a gap: clear any stale ResumeFromDiskError so the
|
||||
// next pass's disk-miss handling cannot skip past the held entry.
|
||||
readInMemoryLogErr = nil
|
||||
if !waitHeld("disk", flushChan) {
|
||||
return nil
|
||||
}
|
||||
continue
|
||||
}
|
||||
if readPersistedLogErr != nil {
|
||||
return fmt.Errorf("reading from persisted logs: %w", readPersistedLogErr)
|
||||
}
|
||||
if isDone {
|
||||
return nil
|
||||
}
|
||||
|
||||
glog.V(4).Infof("processed to %v: %v", clientName, processedTsNs)
|
||||
diskAdvanced := diskReadAdvanced(processedTsNs, lastReadTime)
|
||||
// Read after the disk read (an eviction landing mid-read must count) and
|
||||
// in received-ts space: the ring's bumped stopTimes exceed anything on
|
||||
// any peer's disk, and gating disk cursors on them parks subscribers
|
||||
// that drained every peer's log.
|
||||
lastEvictedTsNs := fs.filer.MetaAggregator.MetaLogBuffer.GetLastEvictedOriginalTsNs()
|
||||
if diskAdvanced {
|
||||
gapStall.resumed()
|
||||
reportUnprovenAggregatedCrossing(cursorBeforeDiskTsNs, processedTsNs, lastEvictedTsNs, diskPassFlushLowTsNs, clientName, req.PathPrefix)
|
||||
lastReadTime = log_buffer.NewMessagePosition(processedTsNs, gapResumeCursorOffset)
|
||||
if processedTsNs > diskAnchorTsNs {
|
||||
diskAnchorTsNs = processedTsNs
|
||||
}
|
||||
} else if readInMemoryLogErr == nil {
|
||||
// Nothing on disk and memory never spoke: scan forward for the next
|
||||
// day that has logs.
|
||||
nextDayTs := util.GetNextDayTsNano(lastReadTime.Time.UnixNano())
|
||||
// The day jump delivers nothing; stay put until the hold point
|
||||
// covers the skipped range.
|
||||
if nextDayTs <= diskPassHoldTsNs {
|
||||
position := log_buffer.NewMessagePosition(nextDayTs, gapResumeCursorOffset)
|
||||
found, err := fs.filer.HasPersistedLogFiles(position)
|
||||
if err != nil {
|
||||
return fmt.Errorf("checking persisted log files: %w", err)
|
||||
}
|
||||
if found {
|
||||
gapStall.resumed()
|
||||
reportUnprovenAggregatedCrossing(cursorBeforeDiskTsNs, nextDayTs, lastEvictedTsNs, diskPassFlushLowTsNs, clientName, req.PathPrefix)
|
||||
lastReadTime = position
|
||||
if nextDayTs > diskAnchorTsNs {
|
||||
diskAnchorTsNs = nextDayTs
|
||||
} else if readInMemoryLogErr == nil {
|
||||
// Nothing on disk and memory never spoke: scan forward for the next
|
||||
// day that has logs.
|
||||
nextDayTs := util.GetNextDayTsNano(lastReadTime.Time.UnixNano())
|
||||
// The day jump delivers nothing; stay put until the hold point
|
||||
// covers the skipped range.
|
||||
if nextDayTs <= diskPassHoldTsNs {
|
||||
position := log_buffer.NewMessagePosition(nextDayTs, gapResumeCursorOffset)
|
||||
found, err := fs.filer.HasPersistedLogFiles(position)
|
||||
if err != nil {
|
||||
return fmt.Errorf("checking persisted log files: %w", err)
|
||||
}
|
||||
if found {
|
||||
gapStall.resumed()
|
||||
reportUnprovenAggregatedCrossing(cursorBeforeDiskTsNs, nextDayTs, lastEvictedTsNs, diskPassFlushLowTsNs, clientName, req.PathPrefix)
|
||||
lastReadTime = position
|
||||
if nextDayTs > diskAnchorTsNs {
|
||||
diskAnchorTsNs = nextDayTs
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -867,6 +917,7 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
|
||||
// every event whose original timestamp is at or below this.
|
||||
preMemDeliveryLowTsNs := fs.filer.MetaAggregator.PeerLowWatermarkTsNs()
|
||||
|
||||
diskReprobeDue := false
|
||||
lastReadTime, isDone, readInMemoryLogErr = fs.filer.MetaAggregator.MetaLogBuffer.LoopProcessLogData(aggReaderName, lastReadTime, req.UntilNs, func() bool {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
@@ -876,6 +927,13 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
|
||||
if !fs.hasClient(req.ClientId, req.ClientEpoch) {
|
||||
return false
|
||||
}
|
||||
// Caught-up readers park in the inner wait loop; the outer disk
|
||||
// gate never runs again unless this read returns, so unwind to
|
||||
// re-probe the persisted logs on the slow cadence.
|
||||
if time.Since(lastDiskPassAt) >= aggDiskReprobeInterval {
|
||||
diskReprobeDue = true
|
||||
return false
|
||||
}
|
||||
// Contiguous and caught up: advance the anchor to the delivery
|
||||
// low-watermark so long live tails keep eviction rewinds short.
|
||||
// Only once the run is connected to the ring - the empty-ring
|
||||
@@ -919,6 +977,9 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
|
||||
}
|
||||
}
|
||||
if isDone {
|
||||
if diskReprobeDue {
|
||||
continue
|
||||
}
|
||||
return nil
|
||||
}
|
||||
if !fs.hasClient(req.ClientId, req.ClientEpoch) {
|
||||
|
||||
@@ -0,0 +1,79 @@
|
||||
package weed_server
|
||||
|
||||
// The persisted log may end before a subscriber's start position (a filer
|
||||
// whose own journal is older than the position a backup client resumes from).
|
||||
// With nothing on disk and every ring entry held by a peer watermark, the
|
||||
// aggregated loop used to re-list and re-read the persisted log on every wake
|
||||
// - about a full CPU core per such subscriber. The disk pass may only re-run
|
||||
// when something it cannot miss has changed.
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/filer"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
)
|
||||
|
||||
type countingStore struct {
|
||||
filer.FilerStore
|
||||
logLists *atomic.Int64
|
||||
}
|
||||
|
||||
func (s *countingStore) ListDirectoryPrefixedEntries(ctx context.Context, dirPath util.FullPath, startFileName string, includeStartFile bool, limit int64, prefix string, eachEntryFunc filer.ListEachEntryFunc) (lastFileName string, err error) {
|
||||
if strings.HasPrefix(string(dirPath), filer.SystemLogDir) {
|
||||
s.logLists.Add(1)
|
||||
}
|
||||
return s.FilerStore.ListDirectoryPrefixedEntries(ctx, dirPath, startFileName, includeStartFile, limit, prefix, eachEntryFunc)
|
||||
}
|
||||
|
||||
func TestSubscribeLoop_AggregatedNoPersistedEntryAfterStart(t *testing.T) {
|
||||
h := newSubscribeHarness(t)
|
||||
|
||||
lists := &atomic.Int64{}
|
||||
h.f.SetStore(&countingStore{FilerStore: h.f.GetStore(), logLists: lists})
|
||||
|
||||
// Persisted log ends at T1: one flushed window, nothing after.
|
||||
h.append(h.tsAt(0, 0))
|
||||
h.append(h.tsAt(0, 1))
|
||||
h.f.LocalMetaLogBuffer.ForceFlush()
|
||||
waitForFlushedFiles(t, h, h.tsAt(0, 1))
|
||||
|
||||
// Client cursor sits just past the last persisted entry.
|
||||
cursor := h.tsAt(0, 1) + int64(time.Millisecond)
|
||||
|
||||
ma := h.startAggregator()
|
||||
|
||||
// Aggregated ring holds only much newer entries (peer events).
|
||||
recent := time.Now().UnixNano()
|
||||
h.appendAggregated(recent)
|
||||
h.appendAggregated(recent + int64(time.Millisecond))
|
||||
|
||||
// Peers' watermarks are stuck at the old log tail.
|
||||
reportPeersAt(ma, h.tsAt(0, 1), h.tsAt(0, 1))
|
||||
|
||||
r := h.subscribeAggregated(cursor)
|
||||
|
||||
// Warm up: the first pass plus the cursor-move re-read after the gap
|
||||
// machinery re-arms the cursor are legitimate.
|
||||
time.Sleep(150 * time.Millisecond)
|
||||
|
||||
before := lists.Load()
|
||||
time.Sleep(500 * time.Millisecond)
|
||||
rate := float64(lists.Load()-before) / 0.5
|
||||
if rate > 4 {
|
||||
t.Fatalf("%.1f persisted-log listings per second while parked; the disk pass re-ran on every wake", rate)
|
||||
}
|
||||
|
||||
// Peer progress through the held entries releases the read: the held
|
||||
// events are delivered without another disk pass.
|
||||
lists.Store(0)
|
||||
reportPeersAt(ma, recent+int64(time.Millisecond), h.tsAt(0, 1))
|
||||
waitForEvents(t, r, []int64{recent, recent + int64(time.Millisecond)}, 3*time.Second)
|
||||
if got := lists.Load(); got > 4 {
|
||||
t.Fatalf("%d listings while draining held entries; delivery should come from the ring", got)
|
||||
}
|
||||
}
|
||||
@@ -47,6 +47,15 @@ func (c *commandEcEncode) Help() string {
|
||||
If you only have less than 4 volume servers, with erasure coding, at least you can afford to
|
||||
have 4 corrupted shard files.
|
||||
|
||||
The guarantee follows from where shards land: a volume survives the loss of any
|
||||
nodes (or racks) that hold at most parityShards (4) shards between them. Spread
|
||||
is best-effort; -shardReplicaPlacement requests limits: its rack digit
|
||||
sets the requested shards per rack, and its node digit the requested shards
|
||||
per node (the data-center digit is not used for EC). For example
|
||||
-shardReplicaPlacement=021 requests at most 1 shard per node and 2 per rack.
|
||||
Only when the final placement meets these limits does losing one rack cost
|
||||
at most 2 shards; the command does not guarantee that the limits are met.
|
||||
|
||||
The -collection parameter is a comma-separated list of collection names, with
|
||||
"*" and "?" wildcards, and regex patterns:
|
||||
- One collection: ec.encode -collection="mybucket"
|
||||
|
||||
@@ -54,8 +54,12 @@ type Volume struct {
|
||||
asyncRequestsChan chan *needle.AsyncRequest
|
||||
asyncWorkerClosed bool
|
||||
|
||||
lastModifiedTsSeconds uint64 // unix time in seconds
|
||||
lastAppendAtNs uint64 // unix time in nanoseconds
|
||||
lastModifiedTsSeconds uint64 // unix time in seconds
|
||||
lastAppendAtNs uint64 // unix time in nanoseconds
|
||||
lastWriteAppendAtNs uint64 // AppendAtNs of the newest write; tombstones don't move it
|
||||
lastWriteNeedleKey types.NeedleId // the write behind the watermark
|
||||
lastWriteDeleted bool // that write was deleted, so recovery has to rescan
|
||||
keepLastModifiedTsOnLoad bool
|
||||
|
||||
lastCompactIndexOffset uint64
|
||||
lastCompactRevision uint16
|
||||
|
||||
@@ -264,7 +264,7 @@ func (v *Volume) recoverLastModifiedTs(indexFile *os.File) {
|
||||
if err != nil || indexSize == 0 {
|
||||
return
|
||||
}
|
||||
appendAtNs, err := findLastWriteAppendAtNs(v, indexFile, indexSize)
|
||||
appendAtNs, key, err := findLastWriteAppendAtNs(v, indexFile, indexSize)
|
||||
if err != nil {
|
||||
glog.Warningf("volume %d recover last write from %s: %v", v.Id, indexFile.Name(), err)
|
||||
return
|
||||
@@ -273,6 +273,11 @@ func (v *Volume) recoverLastModifiedTs(indexFile *os.File) {
|
||||
return
|
||||
}
|
||||
v.lastModifiedTsSeconds = appendAtNs / uint64(time.Second)
|
||||
if appendAtNs > v.lastWriteAppendAtNs {
|
||||
v.lastWriteAppendAtNs = appendAtNs
|
||||
v.lastWriteNeedleKey = key
|
||||
v.lastWriteDeleted = false
|
||||
}
|
||||
}
|
||||
|
||||
// vacuumedLastWriteScanEntries bounds the work a vacuumed volume's recovery
|
||||
@@ -285,25 +290,27 @@ var vacuumedLastWriteScanEntries = 1 << 16
|
||||
|
||||
// findLastWriteAppendAtNs scans the .idx backwards for the newest write -- an
|
||||
// entry that is not a deletion tombstone -- and returns that needle's append
|
||||
// timestamp. The .idx and the .dat share an order, so an append-ordered volume
|
||||
// answers with the first write the scan reaches. Vacuum rewrites both in key
|
||||
// order, which tracks write order only because the master issues keys
|
||||
// increasing: an overwrite keeps its original, lower key, so a vacuumed volume
|
||||
// has to take the maximum over every write it indexes. Returns 0 when the .idx
|
||||
// holds nothing but tombstones, when a vacuumed volume holds more needles than
|
||||
// the scan budget, or for a volume older than version 3, whose needles carry no
|
||||
// append timestamp.
|
||||
func findLastWriteAppendAtNs(v *Volume, indexFile *os.File, indexSize int64) (uint64, error) {
|
||||
// timestamp and key. The .idx and the .dat share an order, so an
|
||||
// append-ordered volume answers with the first write the scan reaches. Vacuum
|
||||
// rewrites both in key order, which tracks write order only because the master
|
||||
// issues keys increasing: an overwrite keeps its original, lower key, so a
|
||||
// vacuumed volume has to take the maximum over every write it indexes. Returns
|
||||
// 0 when the .idx holds nothing but tombstones, when a vacuumed volume holds
|
||||
// more needles than the scan budget, or for a volume older than version 3,
|
||||
// whose needles carry no append timestamp.
|
||||
func findLastWriteAppendAtNs(v *Volume, indexFile *os.File, indexSize int64) (uint64, types.NeedleId, error) {
|
||||
version := v.Version()
|
||||
if version != needle.Version3 {
|
||||
return 0, nil
|
||||
return 0, 0, nil
|
||||
}
|
||||
scanEveryWrite := v.SuperBlock.CompactionRevision > 0
|
||||
entryBudget := vacuumedLastWriteScanEntries
|
||||
if scanEveryWrite && !affordableVacuumedScan(indexFile, indexSize, v.Id, v.FileName(".dat")) {
|
||||
return 0, nil
|
||||
return 0, 0, nil
|
||||
}
|
||||
var lastWriteAppendAtNs uint64
|
||||
var lastWriteKey types.NeedleId
|
||||
dead := make(map[types.NeedleId]struct{})
|
||||
block := make([]byte, types.NeedleMapEntrySize*idx.RowsToRead)
|
||||
for end := indexSize; end > 0; {
|
||||
start := max(end-int64(len(block)), 0)
|
||||
@@ -313,11 +320,17 @@ func findLastWriteAppendAtNs(v *Volume, indexFile *os.File, indexSize int64) (ui
|
||||
err = nil
|
||||
}
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("read %s at %d: %v", indexFile.Name(), start, err)
|
||||
return 0, 0, fmt.Errorf("read %s at %d: %v", indexFile.Name(), start, err)
|
||||
}
|
||||
for i := len(entries) - types.NeedleMapEntrySize; i >= 0; i -= types.NeedleMapEntrySize {
|
||||
key, offset, size := idx.IdxFileEntry(entries[i : i+types.NeedleMapEntrySize])
|
||||
// The first row a key presents is its latest state: a tombstone
|
||||
// there retires the write rows beneath it.
|
||||
if _, gone := dead[key]; gone {
|
||||
continue
|
||||
}
|
||||
if offset.IsZero() || size.IsDeleted() {
|
||||
dead[key] = struct{}{}
|
||||
continue
|
||||
}
|
||||
needleOffset := findNeedleOffset(v.DataBackend, version, offset.ToActualOffset(), key, size)
|
||||
@@ -326,21 +339,24 @@ func findLastWriteAppendAtNs(v *Volume, indexFile *os.File, indexSize int64) (ui
|
||||
}
|
||||
appendAtNs, err := readNeedleAppendAtNs(v.DataBackend, needleOffset, size)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
return 0, 0, err
|
||||
}
|
||||
if appendAtNs > lastWriteAppendAtNs {
|
||||
lastWriteAppendAtNs = appendAtNs
|
||||
lastWriteKey = key
|
||||
}
|
||||
lastWriteAppendAtNs = max(lastWriteAppendAtNs, appendAtNs)
|
||||
if !scanEveryWrite {
|
||||
return lastWriteAppendAtNs, nil
|
||||
return lastWriteAppendAtNs, lastWriteKey, nil
|
||||
}
|
||||
if entryBudget--; entryBudget == 0 {
|
||||
glog.V(0).Infof("volume %d: more than %d needles to scan for its last write, keeping the %s mtime",
|
||||
v.Id, vacuumedLastWriteScanEntries, v.FileName(".dat"))
|
||||
return 0, nil
|
||||
return 0, 0, nil
|
||||
}
|
||||
}
|
||||
end = start
|
||||
}
|
||||
return lastWriteAppendAtNs, nil
|
||||
return lastWriteAppendAtNs, lastWriteKey, nil
|
||||
}
|
||||
|
||||
// affordableVacuumedScan reports whether a vacuumed volume's recovery scan fits
|
||||
|
||||
@@ -172,7 +172,7 @@ func (v *Volume) load(alsoLoadIndex bool, createDatIfMissing bool, needleMapKind
|
||||
return fmt.Errorf("load remote file %v: %w", v.volumeInfo, err)
|
||||
}
|
||||
// Set lastModifiedTsSeconds from remote file to prevent premature expiry on startup
|
||||
if len(v.volumeInfo.GetFiles()) > 0 {
|
||||
if len(v.volumeInfo.GetFiles()) > 0 && !v.keepLastModifiedTsOnLoad {
|
||||
remoteFileModifiedTime := v.volumeInfo.GetFiles()[0].GetModifiedTime()
|
||||
if remoteFileModifiedTime > 0 {
|
||||
v.lastModifiedTsSeconds = remoteFileModifiedTime
|
||||
@@ -201,7 +201,9 @@ func (v *Volume) load(alsoLoadIndex bool, createDatIfMissing bool, needleMapKind
|
||||
if err != nil {
|
||||
return datFileLoadError(v.FileName(".dat"), err)
|
||||
}
|
||||
v.lastModifiedTsSeconds = uint64(modifiedTime.Unix())
|
||||
if !v.keepLastModifiedTsOnLoad {
|
||||
v.lastModifiedTsSeconds = uint64(modifiedTime.Unix())
|
||||
}
|
||||
if fileSize >= super_block.SuperBlockSize {
|
||||
alreadyHasSuperBlock = true
|
||||
}
|
||||
@@ -292,7 +294,9 @@ func (v *Volume) load(alsoLoadIndex bool, createDatIfMissing bool, needleMapKind
|
||||
v.noWriteOrDelete = true
|
||||
glog.V(0).Infof("volumeDataIntegrityChecking failed %v", err)
|
||||
}
|
||||
v.recoverLastModifiedTs(indexFile)
|
||||
if !v.keepLastModifiedTsOnLoad {
|
||||
v.recoverLastModifiedTs(indexFile)
|
||||
}
|
||||
}
|
||||
|
||||
// The post-load structural check below uses the in-memory needle map
|
||||
|
||||
@@ -250,7 +250,7 @@ func TestVolumeTtlClockSkipsUnaffordableScanWithoutDatReads(t *testing.T) {
|
||||
defer func(budget int) { vacuumedLastWriteScanEntries = budget }(vacuumedLastWriteScanEntries)
|
||||
vacuumedLastWriteScanEntries = 2
|
||||
|
||||
appendAtNs, err := findLastWriteAppendAtNs(v, indexFile, indexStat.Size())
|
||||
appendAtNs, _, err := findLastWriteAppendAtNs(v, indexFile, indexStat.Size())
|
||||
if err != nil {
|
||||
t.Fatalf("recover last write: %v", err)
|
||||
}
|
||||
@@ -262,6 +262,152 @@ func TestVolumeTtlClockSkipsUnaffordableScanWithoutDatReads(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestVolumeTtlClockCarriedAcrossVacuumCommit covers the reload that ends a
|
||||
// vacuum commit: the clock must keep the last write's append time rather than
|
||||
// be re-derived, and the intervening delete's tombstone must not freshen it.
|
||||
func TestVolumeTtlClockCarriedAcrossVacuumCommit(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
ttl, err := needle.ReadTTL("5m")
|
||||
if err != nil {
|
||||
t.Fatalf("read ttl: %v", err)
|
||||
}
|
||||
|
||||
v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, ttl, 0, needle.GetCurrentVersion(), 0, 0)
|
||||
if err != nil {
|
||||
t.Fatalf("volume creation: %v", err)
|
||||
}
|
||||
defer v.Close()
|
||||
|
||||
lastWriteNs := uint64(time.Now().Add(-2 * time.Hour).UnixNano())
|
||||
for i := 1; i <= 3; i++ {
|
||||
n := newRandomNeedle(uint64(i))
|
||||
offset, _, _, err := v.writeNeedle2(n, true, false, false)
|
||||
if err != nil {
|
||||
t.Fatalf("write needle %d: %v", i, err)
|
||||
}
|
||||
backdateAppendAtNs(t, v, int64(offset), n.Size, lastWriteNs)
|
||||
}
|
||||
if _, err := v.doDeleteRequest(newEmptyNeedle(2)); err != nil {
|
||||
t.Fatalf("delete needle 2: %v", err)
|
||||
}
|
||||
// Where a restart's recovery would have left the clock. The delete above
|
||||
// pushed lastAppendAtNs to ~now; the commit must not consult it.
|
||||
v.lastModifiedTsSeconds = lastWriteNs / uint64(time.Second)
|
||||
v.lastWriteAppendAtNs = lastWriteNs
|
||||
|
||||
defer func(budget int) { vacuumedLastWriteScanEntries = budget }(vacuumedLastWriteScanEntries)
|
||||
vacuumedLastWriteScanEntries = 1
|
||||
|
||||
if err := v.CompactByIndex(nil); err != nil {
|
||||
t.Fatalf("compact: %v", err)
|
||||
}
|
||||
if err := v.CommitCompact(); err != nil {
|
||||
t.Fatalf("commit compact: %v", err)
|
||||
}
|
||||
if v.SuperBlock.CompactionRevision == 0 {
|
||||
t.Fatal("vacuum must bump CompactionRevision for this test to exercise the vacuumed path")
|
||||
}
|
||||
|
||||
if got, want := v.lastModifiedTsSeconds, lastWriteNs/uint64(time.Second); got != want {
|
||||
t.Errorf("TTL clock after commit is %d, want the last write at %d", got, want)
|
||||
}
|
||||
if !v.expired(v.ContentSize(), 1024*1024) {
|
||||
t.Error("a TTL volume whose last write is 2h old must stay expired across a vacuum commit")
|
||||
}
|
||||
}
|
||||
|
||||
// TestVolumeTtlClockAtCommitUsesAppendTime covers a write whose client
|
||||
// supplied modified time does not match when it was appended: the commit
|
||||
// keeps the server-side append watermark the recovery scan would recompute.
|
||||
func TestVolumeTtlClockAtCommitUsesAppendTime(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
ttl, err := needle.ReadTTL("5m")
|
||||
if err != nil {
|
||||
t.Fatalf("read ttl: %v", err)
|
||||
}
|
||||
|
||||
v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, ttl, 0, needle.GetCurrentVersion(), 0, 0)
|
||||
if err != nil {
|
||||
t.Fatalf("volume creation: %v", err)
|
||||
}
|
||||
defer v.Close()
|
||||
|
||||
future := uint64(time.Now().Add(24 * time.Hour).Unix())
|
||||
n := newRandomNeedle(1)
|
||||
n.LastModified = future
|
||||
if _, _, _, err := v.writeNeedle2(n, true, false, false); err != nil {
|
||||
t.Fatalf("write needle: %v", err)
|
||||
}
|
||||
if v.lastModifiedTsSeconds < future {
|
||||
t.Fatalf("clock %d did not follow the needle's modified time %d", v.lastModifiedTsSeconds, future)
|
||||
}
|
||||
appendWatermarkSec := v.lastWriteAppendAtNs / uint64(time.Second)
|
||||
|
||||
if err := v.CompactByIndex(nil); err != nil {
|
||||
t.Fatalf("compact: %v", err)
|
||||
}
|
||||
if err := v.CommitCompact(); err != nil {
|
||||
t.Fatalf("commit compact: %v", err)
|
||||
}
|
||||
|
||||
if got := v.lastModifiedTsSeconds; got != appendWatermarkSec {
|
||||
t.Errorf("TTL clock after commit is %d, want the append watermark %d", got, appendWatermarkSec)
|
||||
}
|
||||
}
|
||||
|
||||
// TestVolumeTtlClockAtCommitSkipsDeletedWrite covers a vacuum commit whose
|
||||
// newest write was deleted first: the carried watermark belonged to that
|
||||
// write, so the commit has to let the reload rescan and land on the newest
|
||||
// surviving write rather than keep the volume alive on a deleted write's time.
|
||||
func TestVolumeTtlClockAtCommitSkipsDeletedWrite(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
ttl, err := needle.ReadTTL("5m")
|
||||
if err != nil {
|
||||
t.Fatalf("read ttl: %v", err)
|
||||
}
|
||||
|
||||
v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, ttl, 0, needle.GetCurrentVersion(), 0, 0)
|
||||
if err != nil {
|
||||
t.Fatalf("volume creation: %v", err)
|
||||
}
|
||||
defer v.Close()
|
||||
|
||||
oldWriteNs := uint64(time.Now().Add(-2 * time.Hour).UnixNano())
|
||||
newWriteNs := uint64(time.Now().Add(-time.Hour).UnixNano())
|
||||
for _, w := range []struct {
|
||||
id uint64
|
||||
ns uint64
|
||||
}{{1, oldWriteNs}, {2, newWriteNs}} {
|
||||
n := newRandomNeedle(w.id)
|
||||
offset, _, _, err := v.writeNeedle2(n, true, false, false)
|
||||
if err != nil {
|
||||
t.Fatalf("write needle %d: %v", w.id, err)
|
||||
}
|
||||
backdateAppendAtNs(t, v, int64(offset), n.Size, w.ns)
|
||||
}
|
||||
// The delete lands inside the commit window: makeupDiff replays its
|
||||
// tombstone into the new .idx behind the write row the copy carried.
|
||||
if err := v.CompactByIndex(nil); err != nil {
|
||||
t.Fatalf("compact: %v", err)
|
||||
}
|
||||
if _, err := v.doDeleteRequest(newEmptyNeedle(2)); err != nil {
|
||||
t.Fatalf("delete needle 2: %v", err)
|
||||
}
|
||||
if !v.lastWriteDeleted {
|
||||
t.Fatal("deleting the newest write must mark the watermark dead")
|
||||
}
|
||||
if err := v.CommitCompact(); err != nil {
|
||||
t.Fatalf("commit compact: %v", err)
|
||||
}
|
||||
|
||||
if got, want := v.lastModifiedTsSeconds, oldWriteNs/uint64(time.Second); got != want {
|
||||
t.Errorf("TTL clock after commit is %d, want the surviving write at %d", got, want)
|
||||
}
|
||||
if v.lastWriteDeleted {
|
||||
t.Error("the reload's rescan must reseed the watermark off the surviving write")
|
||||
}
|
||||
}
|
||||
|
||||
// TestVolumeExpireAtSecCountsFromLastWrite guards the destroy time an EC volume
|
||||
// is reclaimed on (erasure_coding.EcVolume.IsTimeToDestroy). It was recomputed
|
||||
// as now+TTL on every .vif write, so a read-only mark, a tier upload or an EC
|
||||
|
||||
@@ -222,7 +222,17 @@ func (v *Volume) CommitCompact() error {
|
||||
//time.Sleep(20 * time.Second)
|
||||
|
||||
glog.V(3).Infof("Loading volume %d commit file...", v.Id)
|
||||
if e := v.load(true, false, v.needleMapKind, 0, v.Version()); e != nil {
|
||||
// The write watermark already equals what recoverLastModifiedTs would
|
||||
// rescan, so keep the clock instead of paying for the scan under the lock.
|
||||
// If its write was itself deleted, the reload recovers the newest
|
||||
// surviving write instead.
|
||||
if v.lastWriteAppendAtNs != 0 {
|
||||
v.lastModifiedTsSeconds = v.lastWriteAppendAtNs / uint64(time.Second)
|
||||
}
|
||||
v.keepLastModifiedTsOnLoad = v.lastModifiedTsSeconds != 0 && !v.lastWriteDeleted
|
||||
e := v.load(true, false, v.needleMapKind, 0, v.Version())
|
||||
v.keepLastModifiedTsOnLoad = false
|
||||
if e != nil {
|
||||
return e
|
||||
}
|
||||
glog.V(3).Infof("Finish committing volume %d", v.Id)
|
||||
|
||||
@@ -329,6 +329,10 @@ func (v *Volume) rollbackUnflushedWrite(n *needle.Needle, offset uint64, end int
|
||||
}
|
||||
}
|
||||
}
|
||||
if err == nil && n.Id == v.lastWriteNeedleKey {
|
||||
v.lastWriteAppendAtNs, v.lastWriteNeedleKey = 0, 0
|
||||
v.lastWriteDeleted = true
|
||||
}
|
||||
if err != nil {
|
||||
recoveryErr = errors.Join(recoveryErr,
|
||||
fmt.Errorf("roll back the index of needle %d in volume %d: %w", n.Id, v.Id, err))
|
||||
@@ -409,6 +413,9 @@ func (v *Volume) doWriteRequest(n *needle.Needle, checkCookie bool, fsync bool)
|
||||
return
|
||||
}
|
||||
v.lastAppendAtNs = n.AppendAtNs
|
||||
v.lastWriteAppendAtNs = n.AppendAtNs
|
||||
v.lastWriteNeedleKey = n.Id
|
||||
v.lastWriteDeleted = false
|
||||
|
||||
// add to needle map
|
||||
if !ok || uint64(nv.Offset.ToActualOffset()) < offset {
|
||||
@@ -490,6 +497,10 @@ func (v *Volume) doDeleteRequest(n *needle.Needle) (Size, error) {
|
||||
if err = v.nm.Delete(n.Id, ToOffset(int64(offset))); err != nil {
|
||||
return size, err
|
||||
}
|
||||
if n.Id == v.lastWriteNeedleKey {
|
||||
v.lastWriteAppendAtNs, v.lastWriteNeedleKey = 0, 0
|
||||
v.lastWriteDeleted = true
|
||||
}
|
||||
return size, err
|
||||
}
|
||||
return 0, nil
|
||||
@@ -524,6 +535,9 @@ func (v *Volume) processBatch(currentRequests []*needle.AsyncRequest) {
|
||||
}
|
||||
indexEnd := int64(v.nm.IndexFileSize())
|
||||
batchLastAppendAtNs := v.lastAppendAtNs
|
||||
batchLastWriteAppendAtNs := v.lastWriteAppendAtNs
|
||||
batchLastWriteNeedleKey := v.lastWriteNeedleKey
|
||||
batchLastWriteDeleted := v.lastWriteDeleted
|
||||
batchLastModifiedTsSeconds := v.lastModifiedTsSeconds
|
||||
for i := 0; i < len(currentRequests); i++ {
|
||||
needleID := currentRequests[i].N.Id
|
||||
@@ -563,6 +577,9 @@ func (v *Volume) processBatch(currentRequests []*needle.AsyncRequest) {
|
||||
if syncErr := v.DataBackend.Sync(); syncErr != nil {
|
||||
v.checkReadWriteError(syncErr)
|
||||
v.lastAppendAtNs = batchLastAppendAtNs
|
||||
v.lastWriteAppendAtNs = batchLastWriteAppendAtNs
|
||||
v.lastWriteNeedleKey = batchLastWriteNeedleKey
|
||||
v.lastWriteDeleted = batchLastWriteDeleted
|
||||
v.lastModifiedTsSeconds = batchLastModifiedTsSeconds
|
||||
batchErr := syncErr
|
||||
if recoveryErr := v.rollbackBatch(end, indexEnd, orderedSnapshots, metricRollbacker, batchMetrics); recoveryErr != nil {
|
||||
@@ -702,6 +719,9 @@ func (v *Volume) WriteNeedleBlob(needleId NeedleId, needleBlob []byte, size Size
|
||||
return err
|
||||
}
|
||||
v.lastAppendAtNs = appendAtNs
|
||||
v.lastWriteAppendAtNs = appendAtNs
|
||||
v.lastWriteNeedleKey = needleId
|
||||
v.lastWriteDeleted = false
|
||||
|
||||
// add to needle map
|
||||
if err = v.nm.Put(needleId, ToOffset(int64(offset)), size); err != nil {
|
||||
|
||||
@@ -18,6 +18,7 @@ type Config struct {
|
||||
MinSizeMB int `json:"min_size_mb"`
|
||||
PreferredTags []string `json:"preferred_tags"`
|
||||
ReplicaPlacement string `json:"replica_placement"` // e.g. "020"; empty falls back to the master default replication
|
||||
StrictPlacement bool `json:"strict_placement"` // fail planning instead of relaxing placement constraints
|
||||
}
|
||||
|
||||
// NewDefaultConfig creates a new default erasure coding configuration
|
||||
@@ -171,6 +172,18 @@ func GetConfigSpec() base.ConfigSpec {
|
||||
InputType: "text",
|
||||
CSSClasses: "form-control",
|
||||
},
|
||||
{
|
||||
Name: "strict_placement",
|
||||
JSONName: "strict_placement",
|
||||
Type: config.FieldTypeBool,
|
||||
DefaultValue: false,
|
||||
Required: false,
|
||||
DisplayName: "Strict Placement",
|
||||
Description: "Refuse to encode a volume when the placement constraints can't be satisfied",
|
||||
HelpText: "When enabled, a volume is only encoded if every shard can be placed within the per-disk, anti-affinity, replica-placement and per-rack caps, so the configured resilience is preserved. When disabled (default), unsatisfiable constraints are relaxed and noted in the log",
|
||||
InputType: "checkbox",
|
||||
CSSClasses: "form-check-input",
|
||||
},
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -192,6 +205,7 @@ func (c *Config) ToTaskPolicy() *worker_pb.TaskPolicy {
|
||||
CollectionFilter: c.CollectionFilter,
|
||||
PreferredTags: preferredTagsCopy,
|
||||
ReplicaPlacement: c.ReplicaPlacement,
|
||||
StrictPlacement: c.StrictPlacement,
|
||||
},
|
||||
},
|
||||
}
|
||||
@@ -216,6 +230,7 @@ func (c *Config) FromTaskPolicy(policy *worker_pb.TaskPolicy) error {
|
||||
c.CollectionFilter = ecConfig.CollectionFilter
|
||||
c.PreferredTags = append([]string(nil), ecConfig.PreferredTags...)
|
||||
c.ReplicaPlacement = ecConfig.ReplicaPlacement
|
||||
c.StrictPlacement = ecConfig.StrictPlacement
|
||||
}
|
||||
|
||||
return nil
|
||||
|
||||
@@ -487,8 +487,11 @@ func buildNodeAddressMap(at *topology.ActiveTopology) map[string]string {
|
||||
//
|
||||
// Encode is lenient (PlaceDurabilityFirst): it relaxes caps/anti-affinity/RP and,
|
||||
// last, the total-shards-per-rack cap as needed, failing only when no eligible
|
||||
// disk has room. It prefers the source disk type but spills if that type can't
|
||||
// hold every shard. rp is the resolved replica placement (may be nil).
|
||||
// disk has room. ecConfig.StrictPlacement instead fails the volume's planning
|
||||
// when the constraints cannot be satisfied (PlaceStrict), so a volume is only
|
||||
// encoded while its configured resilience is preserved. The task prefers the
|
||||
// source disk type but spills if that type can't hold every shard. rp is the
|
||||
// resolved replica placement (may be nil).
|
||||
func planECDestinations(snap *ecbalancer.Topology, nodeAddresses map[string]string, metric *types.VolumeHealthMetrics, ecConfig *Config, rp *super_block.ReplicaPlacement, dataShards, parityShards int) (*topology.MultiDestinationPlan, [][]uint32, error) {
|
||||
if snap == nil {
|
||||
return nil, nil, fmt.Errorf("EC placement snapshot not available")
|
||||
@@ -511,13 +514,17 @@ func planECDestinations(snap *ecbalancer.Topology, nodeAddresses map[string]stri
|
||||
for i := range need {
|
||||
need[i] = i
|
||||
}
|
||||
mode := ecbalancer.PlaceDurabilityFirst
|
||||
if ecConfig.StrictPlacement {
|
||||
mode = ecbalancer.PlaceStrict
|
||||
}
|
||||
res, err := snap.Place(metric.VolumeID, metric.Collection, need, ecbalancer.Constraints{
|
||||
DiskType: metric.DiskType,
|
||||
DiskTypePolicy: ecbalancer.DiskTypePrefer,
|
||||
PreferredTags: ecConfig.PreferredTags,
|
||||
ReplicaPlacement: rp,
|
||||
Ratio: func(string) (int, int) { return dataShards, parityShards },
|
||||
}, ecbalancer.PlaceDurabilityFirst)
|
||||
}, mode)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
@@ -0,0 +1,50 @@
|
||||
package erasure_coding
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding/ecbalancer"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/super_block"
|
||||
"github.com/seaweedfs/seaweedfs/weed/worker/types"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// StrictPlacement refuses a volume whose configured caps cannot be met,
|
||||
// rather than relaxing them: here one rack must hold all 14 shards under
|
||||
// a 2-shards-per-rack replica placement.
|
||||
func TestPlanECDestinationsStrictPlacement(t *testing.T) {
|
||||
activeTopology := buildActiveTopology(t, 7, []string{"hdd"}, 100, 0, "")
|
||||
metric := &types.VolumeHealthMetrics{
|
||||
VolumeID: 1,
|
||||
Server: "10.0.0.1:8080",
|
||||
Size: 100 * 1024 * 1024,
|
||||
}
|
||||
rp, err := super_block.NewReplicaPlacementFromString("020")
|
||||
require.NoError(t, err)
|
||||
nodeAddresses := buildNodeAddressMap(activeTopology)
|
||||
|
||||
snap := ecbalancer.FromActiveTopology(activeTopology, erasure_coding.DataShardsCount)
|
||||
cfg := NewDefaultConfig()
|
||||
plan, shardsPerPlan, err := planECDestinations(snap, nodeAddresses, metric, cfg, rp, erasure_coding.DataShardsCount, erasure_coding.ParityShardsCount)
|
||||
require.NoError(t, err, "lenient placement relaxes the unsatisfiable rack cap")
|
||||
requireAllShardsPlaced(t, plan, shardsPerPlan)
|
||||
|
||||
snap = ecbalancer.FromActiveTopology(activeTopology, erasure_coding.DataShardsCount)
|
||||
cfg.StrictPlacement = true
|
||||
_, _, err = planECDestinations(snap, nodeAddresses, metric, cfg, rp, erasure_coding.DataShardsCount, erasure_coding.ParityShardsCount)
|
||||
require.Error(t, err, "strict placement must refuse rather than weaken the rack cap")
|
||||
}
|
||||
|
||||
func TestStrictPlacementRoundTripsThroughTaskPolicy(t *testing.T) {
|
||||
cfg := NewDefaultConfig()
|
||||
cfg.StrictPlacement = true
|
||||
|
||||
restored := NewDefaultConfig()
|
||||
require.NoError(t, restored.FromTaskPolicy(cfg.ToTaskPolicy()))
|
||||
require.True(t, restored.StrictPlacement, "strict placement must survive the persisted policy round trip")
|
||||
|
||||
restored.StrictPlacement = false
|
||||
require.NoError(t, restored.FromTaskPolicy(restored.ToTaskPolicy()))
|
||||
require.False(t, restored.StrictPlacement)
|
||||
}
|
||||
@@ -146,6 +146,13 @@ func (h *ErasureCodingHandler) Descriptor() *plugin_pb.JobTypeDescriptor {
|
||||
FieldType: plugin_pb.ConfigFieldType_CONFIG_FIELD_TYPE_STRING,
|
||||
Widget: plugin_pb.ConfigWidget_CONFIG_WIDGET_TEXT,
|
||||
},
|
||||
{
|
||||
Name: "strict_placement",
|
||||
Label: "Strict Placement",
|
||||
Description: "Fail EC planning when the placement constraints cannot be met instead of relaxing them.",
|
||||
FieldType: plugin_pb.ConfigFieldType_CONFIG_FIELD_TYPE_BOOL,
|
||||
Widget: plugin_pb.ConfigWidget_CONFIG_WIDGET_TOGGLE,
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
@@ -165,6 +172,9 @@ func (h *ErasureCodingHandler) Descriptor() *plugin_pb.JobTypeDescriptor {
|
||||
"replica_placement": {
|
||||
Kind: &plugin_pb.ConfigValue_StringValue{StringValue: ""},
|
||||
},
|
||||
"strict_placement": {
|
||||
Kind: &plugin_pb.ConfigValue_BoolValue{BoolValue: false},
|
||||
},
|
||||
},
|
||||
},
|
||||
AdminRuntimeDefaults: &plugin_pb.AdminRuntimeDefaults{
|
||||
@@ -624,6 +634,8 @@ func deriveErasureCodingWorkerConfig(values map[string]*plugin_pb.ConfigValue) *
|
||||
|
||||
taskConfig.ReplicaPlacement = strings.TrimSpace(pluginworker.ReadStringConfig(values, "replica_placement", taskConfig.ReplicaPlacement))
|
||||
|
||||
taskConfig.StrictPlacement = pluginworker.ReadBoolConfig(values, "strict_placement", taskConfig.StrictPlacement)
|
||||
|
||||
return &erasureCodingWorkerConfig{
|
||||
TaskConfig: taskConfig,
|
||||
}
|
||||
|
||||
@@ -189,9 +189,11 @@ func (h *Handler) compactDataFiles(
|
||||
return "", nil, err
|
||||
}
|
||||
|
||||
// Build compaction bins: group small data files by partition.
|
||||
// Build compaction bins: group small data files by partition, each bin's
|
||||
// files in the order of their bounds so the merged file keeps the
|
||||
// ordering its inputs had.
|
||||
targetSize := compactionTargetSizeForPlan(config, rewritePlan)
|
||||
bins := buildCompactionBins(candidateEntries, targetSize, minInputFiles)
|
||||
bins := buildCompactionBins(candidateEntries, targetSize, minInputFiles, meta)
|
||||
initialBinCount := len(bins)
|
||||
bins = filterCompactionBinsByPlan(bins, config, rewritePlan)
|
||||
if len(bins) == 0 {
|
||||
@@ -269,10 +271,11 @@ func (h *Handler) compactDataFiles(
|
||||
|
||||
var mergedData []byte
|
||||
var recordCount int64
|
||||
rowGroupRows := rowsPerRowGroup(bin, config)
|
||||
if rewritePlan != nil && rewritePlan.strategy == "sort" {
|
||||
mergedData, recordCount, err = mergeParquetFilesSorted(ctx, filerClient, bucketName, dataPath, bin.Entries, positionDeletes, eqDeleteGroups, schema, rewritePlan)
|
||||
mergedData, recordCount, err = mergeParquetFilesSorted(ctx, filerClient, bucketName, dataPath, bin.Entries, positionDeletes, eqDeleteGroups, schema, rewritePlan, rowGroupRows)
|
||||
} else {
|
||||
mergedData, recordCount, err = mergeParquetFiles(ctx, filerClient, bucketName, dataPath, bin.Entries, positionDeletes, eqDeleteGroups, schema)
|
||||
mergedData, recordCount, err = mergeParquetFiles(ctx, filerClient, bucketName, dataPath, bin.Entries, positionDeletes, eqDeleteGroups, schema, rowGroupRows)
|
||||
}
|
||||
if err != nil {
|
||||
glog.Warningf("iceberg compact: failed to merge bin %d (%d files): %v", binIdx, len(bin.Entries), err)
|
||||
@@ -314,6 +317,14 @@ func (h *Handler) compactDataFiles(
|
||||
}
|
||||
writtenArtifacts = append(writtenArtifacts, artifact{dir: dataDir, fileName: mergedFileName})
|
||||
|
||||
// Record the column statistics readers prune files by; without
|
||||
// them every scan has to read every compacted file.
|
||||
if stats, statsErr := collectColumnStats(mergedData, schema); statsErr != nil {
|
||||
glog.Warningf("iceberg compact: no column statistics for %s: %v", mergedFileName, statsErr)
|
||||
} else {
|
||||
stats.applyTo(dfBuilder)
|
||||
}
|
||||
|
||||
mergedDataFile := dfBuilder.Build()
|
||||
summary.addFile(mergedDataFile)
|
||||
newEntry := iceberg.NewManifestEntry(
|
||||
@@ -527,7 +538,13 @@ func (h *Handler) compactDataFiles(
|
||||
// buildCompactionBins groups small data files by partition for bin-packing.
|
||||
// A file is "small" if it's below targetSize. A bin must have at least
|
||||
// minFiles entries to be worth compacting.
|
||||
func buildCompactionBins(entries []iceberg.ManifestEntry, targetSize int64, minFiles int) []compactionBin {
|
||||
//
|
||||
// With an order, every bin lists its files by their bounds on the ordering
|
||||
// column and an oversized partition is split into runs of consecutive files,
|
||||
// so each merged file covers one contiguous range of that column. Without one
|
||||
// the files keep manifest order and an oversized partition is packed
|
||||
// largest-first.
|
||||
func buildCompactionBins(entries []iceberg.ManifestEntry, targetSize int64, minFiles int, meta table.Metadata) []compactionBin {
|
||||
if minFiles < 2 {
|
||||
minFiles = 2
|
||||
}
|
||||
@@ -560,14 +577,37 @@ func buildCompactionBins(entries []iceberg.ManifestEntry, targetSize int64, minF
|
||||
bin.TotalSize += df.FileSizeBytes()
|
||||
}
|
||||
|
||||
// Filter to bins with enough files, splitting oversized bins
|
||||
// Filter to bins with enough files, splitting oversized bins. The merge
|
||||
// order resolves per group over the files that survived the eligibility
|
||||
// filter, so a file that can never participate (oversized, another
|
||||
// format) cannot veto ordering for the ones that can.
|
||||
var result []compactionBin
|
||||
for _, bin := range groups {
|
||||
if len(bin.Entries) < minFiles {
|
||||
continue
|
||||
}
|
||||
order := resolveCompactionOrder(meta, bin.Entries)
|
||||
order.sortEntries(bin.Entries)
|
||||
if bin.TotalSize <= targetSize {
|
||||
result = append(result, *bin)
|
||||
} else if order != nil {
|
||||
runs, leftover := splitOrderedBin(*bin, targetSize, minFiles)
|
||||
result = append(result, runs...)
|
||||
if order.bestEffort && len(leftover) > 0 {
|
||||
runtBin := compactionBin{PartitionKey: bin.PartitionKey, Partition: bin.Partition, SpecID: bin.SpecID, Entries: leftover}
|
||||
packed := splitOversizedBin(runtBin, targetSize, minFiles)
|
||||
if len(packed) == 0 {
|
||||
// Stragglers cannot merge even with each other; repacking
|
||||
// the whole bin lets them pair across the ordered runs.
|
||||
// Only replace the runs when the repack yields bins —
|
||||
// otherwise the runs stay valid compaction work.
|
||||
if full := splitOversizedBin(*bin, targetSize, minFiles); len(full) > 0 {
|
||||
result = result[:len(result)-len(runs)]
|
||||
packed = full
|
||||
}
|
||||
}
|
||||
result = append(result, packed...)
|
||||
}
|
||||
} else {
|
||||
result = append(result, splitOversizedBin(*bin, targetSize, minFiles)...)
|
||||
}
|
||||
@@ -918,6 +958,46 @@ func resolveEqualityColIndices(pqSchema *parquet.Schema, fieldIDs []int, iceberg
|
||||
return indices, nil
|
||||
}
|
||||
|
||||
const (
|
||||
// Iceberg's write.parquet.row-group-size-bytes default and PyIceberg's
|
||||
// write.parquet.row-group-limit default.
|
||||
defaultRowGroupSizeBytes = 128 * 1024 * 1024
|
||||
defaultRowGroupRowLimit = 1048576
|
||||
minRowGroupRows = 1024
|
||||
)
|
||||
|
||||
// rowsPerRowGroup is the most rows one row group of a compacted file holds:
|
||||
// the row limit, lowered so a row group stays near the byte size, estimated
|
||||
// from the bin's inputs' bytes per row. parquet-go's writers have no limit of
|
||||
// their own, so without one every compacted file is a single row group. The
|
||||
// floor only protects the estimate — a configured row limit always wins.
|
||||
func rowsPerRowGroup(bin compactionBin, config Config) int64 {
|
||||
rows, sizeBytes := config.RowGroupRowLimit, config.RowGroupSizeBytes
|
||||
explicit := rows > 0
|
||||
if !explicit {
|
||||
rows = defaultRowGroupRowLimit
|
||||
}
|
||||
if sizeBytes <= 0 {
|
||||
sizeBytes = defaultRowGroupSizeBytes
|
||||
}
|
||||
var totalRows int64
|
||||
for _, entry := range bin.Entries {
|
||||
totalRows += entry.DataFile().Count()
|
||||
}
|
||||
if totalRows > 0 && bin.TotalSize > 0 {
|
||||
bytesPerRow := max(bin.TotalSize/totalRows, 1)
|
||||
estimated := sizeBytes / bytesPerRow
|
||||
if !explicit {
|
||||
estimated = max(estimated, minRowGroupRows)
|
||||
}
|
||||
rows = min(rows, estimated)
|
||||
}
|
||||
if !explicit {
|
||||
rows = max(rows, minRowGroupRows)
|
||||
}
|
||||
return rows
|
||||
}
|
||||
|
||||
// mergeParquetFiles reads multiple small Parquet files and merges them into
|
||||
// a single Parquet file, optionally filtering out rows matching position or
|
||||
// equality deletes. Files are processed one at a time to keep memory usage
|
||||
@@ -930,6 +1010,7 @@ func mergeParquetFiles(
|
||||
positionDeletes map[string][]int64,
|
||||
eqDeleteGroups []equalityDeleteGroup,
|
||||
icebergSchema *iceberg.Schema,
|
||||
rowGroupRows int64,
|
||||
) ([]byte, int64, error) {
|
||||
if len(entries) == 0 {
|
||||
return nil, 0, fmt.Errorf("no entries to merge")
|
||||
@@ -954,7 +1035,7 @@ func mergeParquetFiles(
|
||||
}
|
||||
|
||||
var outputBuf bytes.Buffer
|
||||
writer := parquet.NewWriter(&outputBuf, parquetSchema)
|
||||
writer := parquet.NewWriter(&outputBuf, parquetSchema, parquet.MaxRowsPerRowGroup(rowGroupRows))
|
||||
|
||||
drainReader := func(reader *parquet.Reader, source string) (int64, error) {
|
||||
return visitFilteredParquetRows(ctx, reader, source, bucketName, dataPath, positionDeletes, resolvedEqGroups, func(filtered []parquet.Row) error {
|
||||
@@ -1119,6 +1200,7 @@ func mergeParquetFilesSorted(
|
||||
eqDeleteGroups []equalityDeleteGroup,
|
||||
icebergSchema *iceberg.Schema,
|
||||
rewritePlan *compactionRewritePlan,
|
||||
rowGroupRows int64,
|
||||
) ([]byte, int64, error) {
|
||||
if len(entries) == 0 {
|
||||
return nil, 0, fmt.Errorf("no entries to merge")
|
||||
@@ -1168,6 +1250,7 @@ func mergeParquetFilesSorted(
|
||||
|
||||
var outputBuf bytes.Buffer
|
||||
writer := parquet.NewSortingWriter[any](&outputBuf, sortBufferRows, parquetSchema,
|
||||
parquet.MaxRowsPerRowGroup(rowGroupRows),
|
||||
parquet.SortingWriterConfig(
|
||||
parquet.SortingColumns(sortingColumns...),
|
||||
parquet.SortingBuffers(parquet.NewFileBufferPool(spillDir, "seaweedfs-iceberg-sort-*")),
|
||||
|
||||
@@ -0,0 +1,201 @@
|
||||
package iceberg
|
||||
|
||||
import (
|
||||
"sort"
|
||||
|
||||
"github.com/apache/iceberg-go"
|
||||
"github.com/apache/iceberg-go/table"
|
||||
)
|
||||
|
||||
// compactionOrder names the column by whose bounds a bin's input files are
|
||||
// ordered before they are merged: the first identity field of the table's
|
||||
// sort order when it declares one, otherwise the first column every candidate
|
||||
// file carries bounds for. Rows are copied file by file, so a merged file is
|
||||
// as ordered as the sequence of its inputs.
|
||||
type compactionOrder struct {
|
||||
fieldID int
|
||||
typ iceberg.PrimitiveType
|
||||
descending bool
|
||||
// bestEffort marks an order inferred from column bounds rather than the
|
||||
// table's declared sort order; runs it cannot form may still merge by
|
||||
// size since contiguous ranges are a preference, not a requirement.
|
||||
bestEffort bool
|
||||
}
|
||||
|
||||
func resolveCompactionOrder(meta table.Metadata, entries []iceberg.ManifestEntry) *compactionOrder {
|
||||
if meta == nil || meta.CurrentSchema() == nil || len(entries) == 0 {
|
||||
return nil
|
||||
}
|
||||
schema := meta.CurrentSchema()
|
||||
if sortOrder := meta.SortOrder(); !sortOrder.IsUnsorted() {
|
||||
for _, sortField := range sortOrder.Fields() {
|
||||
if _, ok := sortField.Transform.(iceberg.IdentityTransform); !ok {
|
||||
continue
|
||||
}
|
||||
field, ok := schema.FindFieldByID(sortField.SourceID())
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
typ, ok := field.Type.(iceberg.PrimitiveType)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
candidate := &compactionOrder{fieldID: field.ID, typ: typ, descending: sortField.Direction == table.SortDESC}
|
||||
// The declared order only holds when every entry carries a bound
|
||||
// for it; otherwise merging would claim sortedness it cannot
|
||||
// verify. Fall back to inference instead of ordering by a later
|
||||
// sort field alone.
|
||||
if candidate.coversEntries(entries) {
|
||||
return candidate
|
||||
}
|
||||
break
|
||||
}
|
||||
}
|
||||
// A column only orders the merge when every entry in the group carries a
|
||||
// bound for it. The check is scoped to the entries handed in — a bin's
|
||||
// eligible files — so an unrelated file without bounds (oversized, a
|
||||
// different format, or another partition) cannot disable ordering here.
|
||||
for _, field := range schema.Fields() {
|
||||
typ, ok := field.Type.(iceberg.PrimitiveType)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
candidate := &compactionOrder{fieldID: field.ID, typ: typ, bestEffort: true}
|
||||
if candidate.coversEntries(entries) {
|
||||
return candidate
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (o *compactionOrder) coversEntries(entries []iceberg.ManifestEntry) bool {
|
||||
for _, entry := range entries {
|
||||
if _, ok := o.key(entry); !ok {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// key is the bound a file is ordered by: its lower bound on the column for an
|
||||
// ascending order, its upper bound for a descending one.
|
||||
func (o *compactionOrder) key(entry iceberg.ManifestEntry) (iceberg.Literal, bool) {
|
||||
df := entry.DataFile()
|
||||
raw := df.LowerBoundValues()[o.fieldID]
|
||||
if o.descending {
|
||||
raw = df.UpperBoundValues()[o.fieldID]
|
||||
}
|
||||
if raw == nil {
|
||||
return nil, false
|
||||
}
|
||||
lit, err := iceberg.LiteralFromBytes(o.typ, raw)
|
||||
return lit, err == nil
|
||||
}
|
||||
|
||||
// sortEntries orders a bin's files by their key. Files without a key sort
|
||||
// after the rest, in their original order.
|
||||
func (o *compactionOrder) sortEntries(entries []iceberg.ManifestEntry) {
|
||||
if o == nil || len(entries) < 2 {
|
||||
return
|
||||
}
|
||||
type keyed struct {
|
||||
key iceberg.Literal
|
||||
ok bool
|
||||
}
|
||||
keys := make([]keyed, len(entries))
|
||||
idx := make([]int, len(entries))
|
||||
for i, entry := range entries {
|
||||
k, ok := o.key(entry)
|
||||
keys[i], idx[i] = keyed{k, ok}, i
|
||||
}
|
||||
sort.SliceStable(idx, func(i, j int) bool {
|
||||
a, b := keys[idx[i]], keys[idx[j]]
|
||||
if a.ok != b.ok {
|
||||
return a.ok
|
||||
}
|
||||
if !a.ok {
|
||||
return false
|
||||
}
|
||||
c, ok := compareLiterals(a.key, b.key)
|
||||
if o.descending {
|
||||
c = -c
|
||||
}
|
||||
return ok && c < 0
|
||||
})
|
||||
sorted := make([]iceberg.ManifestEntry, len(entries))
|
||||
for i, j := range idx {
|
||||
sorted[i] = entries[j]
|
||||
}
|
||||
copy(entries, sorted)
|
||||
}
|
||||
|
||||
// splitOrderedBin splits a bin whose files are in bound order into runs of
|
||||
// consecutive files that stay under targetSize, so every output covers one
|
||||
// contiguous range of the ordering column. The returned entries are the runs
|
||||
// too short to reach minFiles: the caller can still merge them by size when
|
||||
// contiguous ranges are only a preference, or leave them for a later pass
|
||||
// when order must hold.
|
||||
func splitOrderedBin(bin compactionBin, targetSize int64, minFiles int) (valid []compactionBin, leftover []iceberg.ManifestEntry) {
|
||||
newBin := func() compactionBin {
|
||||
return compactionBin{PartitionKey: bin.PartitionKey, Partition: bin.Partition, SpecID: bin.SpecID}
|
||||
}
|
||||
current := newBin()
|
||||
flush := func() {
|
||||
if len(current.Entries) >= minFiles {
|
||||
valid = append(valid, current)
|
||||
} else {
|
||||
leftover = append(leftover, current.Entries...)
|
||||
}
|
||||
current = newBin()
|
||||
}
|
||||
for _, entry := range bin.Entries {
|
||||
size := entry.DataFile().FileSizeBytes()
|
||||
if current.TotalSize > 0 && current.TotalSize+size > targetSize {
|
||||
flush()
|
||||
}
|
||||
current.Entries = append(current.Entries, entry)
|
||||
current.TotalSize += size
|
||||
}
|
||||
flush()
|
||||
return valid, leftover
|
||||
}
|
||||
|
||||
// compareLiterals orders two literals of the same type; false when they are
|
||||
// not comparable.
|
||||
func compareLiterals(a, b iceberg.Literal) (int, bool) {
|
||||
switch x := a.(type) {
|
||||
case iceberg.TypedLiteral[bool]:
|
||||
return compareTyped(x, b)
|
||||
case iceberg.TypedLiteral[int32]:
|
||||
return compareTyped(x, b)
|
||||
case iceberg.TypedLiteral[int64]:
|
||||
return compareTyped(x, b)
|
||||
case iceberg.TypedLiteral[float32]:
|
||||
return compareTyped(x, b)
|
||||
case iceberg.TypedLiteral[float64]:
|
||||
return compareTyped(x, b)
|
||||
case iceberg.TypedLiteral[iceberg.Date]:
|
||||
return compareTyped(x, b)
|
||||
case iceberg.TypedLiteral[iceberg.Time]:
|
||||
return compareTyped(x, b)
|
||||
case iceberg.TypedLiteral[iceberg.Timestamp]:
|
||||
return compareTyped(x, b)
|
||||
case iceberg.TypedLiteral[iceberg.TimestampNano]:
|
||||
return compareTyped(x, b)
|
||||
case iceberg.TypedLiteral[string]:
|
||||
return compareTyped(x, b)
|
||||
case iceberg.TypedLiteral[[]byte]:
|
||||
return compareTyped(x, b)
|
||||
case iceberg.TypedLiteral[iceberg.Decimal]:
|
||||
return compareTyped(x, b)
|
||||
}
|
||||
return 0, false
|
||||
}
|
||||
|
||||
func compareTyped[T iceberg.LiteralType](a iceberg.TypedLiteral[T], b iceberg.Literal) (int, bool) {
|
||||
bt, ok := b.(iceberg.TypedLiteral[T])
|
||||
if !ok {
|
||||
return 0, false
|
||||
}
|
||||
return a.Comparator()(a.Value(), bt.Value()), true
|
||||
}
|
||||
@@ -0,0 +1,388 @@
|
||||
package iceberg
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"path"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/apache/iceberg-go"
|
||||
"github.com/apache/iceberg-go/table"
|
||||
|
||||
filer_pb "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3tables"
|
||||
)
|
||||
|
||||
type boundedFile struct {
|
||||
Name string
|
||||
IDs []int64
|
||||
Lower, Upper int64
|
||||
SizeBytes int64
|
||||
NoBounds bool
|
||||
}
|
||||
|
||||
// populateBoundedTable is populateTableWithDeleteFiles narrowed to data
|
||||
// files whose manifest entries carry bounds on the id column.
|
||||
func populateBoundedTable(t *testing.T, fs *fakeFilerServer, setup tableSetup, files []boundedFile) {
|
||||
t.Helper()
|
||||
populateBoundedTableSorted(t, fs, setup, files, table.UnsortedSortOrder)
|
||||
}
|
||||
|
||||
func populateBoundedTableSorted(t *testing.T, fs *fakeFilerServer, setup tableSetup, files []boundedFile, sortOrder table.SortOrder) {
|
||||
t.Helper()
|
||||
schema := newTestSchema()
|
||||
spec := *iceberg.UnpartitionedSpec
|
||||
|
||||
meta, err := table.NewMetadata(schema, &spec, sortOrder, "s3://"+setup.BucketName+"/"+setup.dataPath(), nil)
|
||||
if err != nil {
|
||||
t.Fatalf("create metadata: %v", err)
|
||||
}
|
||||
|
||||
bucketPath := path.Join(s3tables.TablesPath, setup.BucketName)
|
||||
nsPath := path.Join(bucketPath, setup.Namespace)
|
||||
tableFilerPath := path.Join(bucketPath, setup.dataPath())
|
||||
metaDir := path.Join(tableFilerPath, "metadata")
|
||||
dataDir := path.Join(tableFilerPath, "data")
|
||||
version := meta.Version()
|
||||
|
||||
var entries []iceberg.ManifestEntry
|
||||
for _, f := range files {
|
||||
rows := make([]struct {
|
||||
ID int64
|
||||
Name string
|
||||
}, len(f.IDs))
|
||||
for i, id := range f.IDs {
|
||||
rows[i] = struct {
|
||||
ID int64
|
||||
Name string
|
||||
}{id, fmt.Sprintf("r%d", id)}
|
||||
}
|
||||
data := writeTestParquetFile(t, fs, dataDir, f.Name, rows)
|
||||
size := int64(len(data))
|
||||
if f.SizeBytes > 0 {
|
||||
size = f.SizeBytes
|
||||
}
|
||||
dfb, err := iceberg.NewDataFileBuilder(spec, iceberg.EntryContentData, setup.fileRef("data", f.Name),
|
||||
iceberg.ParquetFile, map[int]any{}, nil, nil, int64(len(f.IDs)), size)
|
||||
if err != nil {
|
||||
t.Fatalf("build data file %s: %v", f.Name, err)
|
||||
}
|
||||
if !f.NoBounds {
|
||||
lo, _ := iceberg.Int64Literal(f.Lower).MarshalBinary()
|
||||
hi, _ := iceberg.Int64Literal(f.Upper).MarshalBinary()
|
||||
dfb.LowerBoundValues(map[int][]byte{1: lo}).UpperBoundValues(map[int][]byte{1: hi})
|
||||
}
|
||||
snapID := int64(1)
|
||||
entries = append(entries, iceberg.NewManifestEntry(iceberg.EntryStatusADDED, &snapID, nil, nil, dfb.Build()))
|
||||
}
|
||||
|
||||
var manifestBuf bytes.Buffer
|
||||
mf, err := iceberg.WriteManifest(setup.fileRef("metadata", "data-manifest-1.avro"), &manifestBuf,
|
||||
version, spec, schema, 1, entries)
|
||||
if err != nil {
|
||||
t.Fatalf("write manifest: %v", err)
|
||||
}
|
||||
fs.putEntry(metaDir, "data-manifest-1.avro", &filer_pb.Entry{
|
||||
Name: "data-manifest-1.avro", Content: manifestBuf.Bytes(),
|
||||
Attributes: &filer_pb.FuseAttributes{Mtime: time.Now().Unix(), FileSize: uint64(manifestBuf.Len())},
|
||||
})
|
||||
|
||||
var mlBuf bytes.Buffer
|
||||
seqNum := int64(1)
|
||||
if err := iceberg.WriteManifestList(version, &mlBuf, 1, nil, &seqNum, 0, []iceberg.ManifestFile{mf}); err != nil {
|
||||
t.Fatalf("write manifest list: %v", err)
|
||||
}
|
||||
fs.putEntry(metaDir, "snap-1.avro", &filer_pb.Entry{
|
||||
Name: "snap-1.avro", Content: mlBuf.Bytes(),
|
||||
Attributes: &filer_pb.FuseAttributes{Mtime: time.Now().Unix(), FileSize: uint64(mlBuf.Len())},
|
||||
})
|
||||
|
||||
snap := table.Snapshot{SnapshotID: 1, TimestampMs: time.Now().UnixMilli(),
|
||||
ManifestList: setup.fileRef("metadata", "snap-1.avro")}
|
||||
builder, err := table.MetadataBuilderFromBase(meta, "s3://"+setup.BucketName+"/"+setup.dataPath())
|
||||
if err != nil {
|
||||
t.Fatalf("metadata builder: %v", err)
|
||||
}
|
||||
if err := builder.AddSnapshot(&snap); err != nil {
|
||||
t.Fatalf("add snapshot: %v", err)
|
||||
}
|
||||
if err := builder.SetSnapshotRef(table.MainBranch, snap.SnapshotID, table.BranchRef); err != nil {
|
||||
t.Fatalf("set snapshot ref: %v", err)
|
||||
}
|
||||
meta, err = builder.Build()
|
||||
if err != nil {
|
||||
t.Fatalf("build metadata: %v", err)
|
||||
}
|
||||
|
||||
fullMetadataJSON, _ := json.Marshal(meta)
|
||||
xattr, _ := json.Marshal(map[string]interface{}{
|
||||
"metadataVersion": 1,
|
||||
"metadataLocation": setup.fileRef("metadata", "v1.metadata.json"),
|
||||
"metadata": map[string]interface{}{"fullMetadata": json.RawMessage(fullMetadataJSON)},
|
||||
})
|
||||
fs.putEntry(path.Join(s3tables.TablesPath), setup.BucketName, &filer_pb.Entry{
|
||||
Name: setup.BucketName, IsDirectory: true,
|
||||
Extended: map[string][]byte{s3tables.ExtendedKeyTableBucket: []byte("true")},
|
||||
})
|
||||
fs.putEntry(bucketPath, setup.Namespace, &filer_pb.Entry{Name: setup.Namespace, IsDirectory: true})
|
||||
fs.putEntry(nsPath, setup.TableName, &filer_pb.Entry{
|
||||
Name: setup.TableName, IsDirectory: true,
|
||||
Extended: map[string][]byte{s3tables.ExtendedKeyMetadata: xattr},
|
||||
})
|
||||
}
|
||||
|
||||
// A bin's files are merged in the order of their bounds on the
|
||||
// ordering column, not manifest order, so a merged file is as ordered
|
||||
// as the sequence of its inputs.
|
||||
func TestCompactDataFilesKeepsInputOrderByBounds(t *testing.T) {
|
||||
fs, client := startFakeFiler(t)
|
||||
setup := tableSetup{BucketName: "tb", Namespace: "ns", TableName: "tbl"}
|
||||
// manifest order is the reverse of bound order
|
||||
populateBoundedTable(t, fs, setup, []boundedFile{
|
||||
{Name: "hi.parquet", IDs: []int64{3, 4}, Lower: 3, Upper: 4},
|
||||
{Name: "lo.parquet", IDs: []int64{1, 2}, Lower: 1, Upper: 2},
|
||||
})
|
||||
|
||||
handler := NewHandler(nil)
|
||||
config := Config{TargetFileSizeBytes: 256 * 1024 * 1024, MinInputFiles: 2, MaxCommitRetries: 3}
|
||||
if _, _, err := handler.compactDataFiles(context.Background(), client, setup.BucketName, setup.tablePath(), config, nil); err != nil {
|
||||
t.Fatalf("compactDataFiles: %v", err)
|
||||
}
|
||||
|
||||
df := compactedDataFile(t, client, setup)
|
||||
ids := compactedIDs(t, client, setup, df.FilePath())
|
||||
want := []int64{1, 2, 3, 4}
|
||||
if fmt.Sprint(ids) != fmt.Sprint(want) {
|
||||
t.Errorf("merged rows = %v, want %v", ids, want)
|
||||
}
|
||||
}
|
||||
|
||||
// An oversized partition is split into runs of consecutive files in
|
||||
// bound order, so each output covers one contiguous range.
|
||||
func TestCompactDataFilesSplitsOrderedRuns(t *testing.T) {
|
||||
fs, client := startFakeFiler(t)
|
||||
setup := tableSetup{BucketName: "tb", Namespace: "ns", TableName: "tbl"}
|
||||
// four disjoint ranges, each "4 MB": one bin, split into two runs of two
|
||||
populateBoundedTable(t, fs, setup, []boundedFile{
|
||||
{Name: "q4.parquet", IDs: []int64{7, 8}, Lower: 7, Upper: 8, SizeBytes: 4 << 20},
|
||||
{Name: "q1.parquet", IDs: []int64{1, 2}, Lower: 1, Upper: 2, SizeBytes: 4 << 20},
|
||||
{Name: "q3.parquet", IDs: []int64{5, 6}, Lower: 5, Upper: 6, SizeBytes: 4 << 20},
|
||||
{Name: "q2.parquet", IDs: []int64{3, 4}, Lower: 3, Upper: 4, SizeBytes: 4 << 20},
|
||||
})
|
||||
|
||||
handler := NewHandler(nil)
|
||||
config := Config{TargetFileSizeBytes: 8 << 20, MinInputFiles: 2, MaxCommitRetries: 3}
|
||||
result, _, err := handler.compactDataFiles(context.Background(), client, setup.BucketName, setup.tablePath(), config, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("compactDataFiles: %v", err)
|
||||
}
|
||||
if !strings.Contains(result, "compacted 4 files into 2") {
|
||||
t.Fatalf("expected two ordered runs, got %q", result)
|
||||
}
|
||||
|
||||
// each output must cover one contiguous range
|
||||
state, err := loadCurrentMetadata(context.Background(), client, setup.BucketName, setup.tablePath())
|
||||
if err != nil {
|
||||
t.Fatalf("loadCurrentMetadata: %v", err)
|
||||
}
|
||||
manifests, err := loadCurrentManifests(context.Background(), client, setup.BucketName, state.DataPath, state.Metadata)
|
||||
if err != nil {
|
||||
t.Fatalf("loadCurrentManifests: %v", err)
|
||||
}
|
||||
var ranges [][2]int64
|
||||
for _, mf := range manifests {
|
||||
if mf.ManifestContent() != iceberg.ManifestContentData {
|
||||
continue
|
||||
}
|
||||
manifestData, err := loadFileByIcebergPath(context.Background(), client, setup.BucketName, state.DataPath, mf.FilePath())
|
||||
if err != nil {
|
||||
t.Fatalf("load manifest: %v", err)
|
||||
}
|
||||
entries, err := iceberg.ReadManifest(mf, bytes.NewReader(manifestData), true)
|
||||
if err != nil {
|
||||
t.Fatalf("read manifest: %v", err)
|
||||
}
|
||||
for _, entry := range entries {
|
||||
if entry.Status() == iceberg.EntryStatusDELETED {
|
||||
continue
|
||||
}
|
||||
lo, _ := iceberg.LiteralFromBytes(iceberg.PrimitiveTypes.Int64, entry.DataFile().LowerBoundValues()[1])
|
||||
hi, _ := iceberg.LiteralFromBytes(iceberg.PrimitiveTypes.Int64, entry.DataFile().UpperBoundValues()[1])
|
||||
ranges = append(ranges, [2]int64{int64(lo.(iceberg.Int64Literal)), int64(hi.(iceberg.Int64Literal))})
|
||||
}
|
||||
}
|
||||
if len(ranges) != 2 {
|
||||
t.Fatalf("expected 2 output files, got %d", len(ranges))
|
||||
}
|
||||
// ranges must be contiguous runs, not interleaved: {1..4} and {5..8}
|
||||
for _, r := range ranges {
|
||||
if r[1]-r[0] != 3 {
|
||||
t.Errorf("output covers non-contiguous range %v", r)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// An oversized file without bounds can never join a merge, so it must not
|
||||
// disable bounds ordering for the files that can.
|
||||
func TestCompactDataFilesOrdersPastIneligibleFile(t *testing.T) {
|
||||
fs, client := startFakeFiler(t)
|
||||
setup := tableSetup{BucketName: "tb", Namespace: "ns", TableName: "tbl"}
|
||||
populateBoundedTable(t, fs, setup, []boundedFile{
|
||||
{Name: "hi.parquet", IDs: []int64{3, 4}, Lower: 3, Upper: 4},
|
||||
{Name: "lo.parquet", IDs: []int64{1, 2}, Lower: 1, Upper: 2},
|
||||
{Name: "huge.parquet", IDs: []int64{9}, SizeBytes: 512 << 20, NoBounds: true},
|
||||
})
|
||||
|
||||
handler := NewHandler(nil)
|
||||
config := Config{TargetFileSizeBytes: 256 << 20, MinInputFiles: 2, MaxCommitRetries: 3}
|
||||
result, _, err := handler.compactDataFiles(context.Background(), client, setup.BucketName, setup.tablePath(), config, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("compactDataFiles: %v", err)
|
||||
}
|
||||
if !strings.Contains(result, "compacted 2 files into 1") {
|
||||
t.Fatalf("expected only the two eligible files to merge, got %q", result)
|
||||
}
|
||||
|
||||
var merged iceberg.DataFile
|
||||
for _, df := range liveDataFiles(t, client, setup) {
|
||||
if !strings.HasSuffix(df.FilePath(), "huge.parquet") {
|
||||
merged = df
|
||||
}
|
||||
}
|
||||
ids := compactedIDs(t, client, setup, merged.FilePath())
|
||||
want := []int64{1, 2, 3, 4}
|
||||
if fmt.Sprint(ids) != fmt.Sprint(want) {
|
||||
t.Errorf("merged rows = %v, want %v", ids, want)
|
||||
}
|
||||
}
|
||||
|
||||
// Runs too short to merge in bound order still merge by size, so files that
|
||||
// cannot form a contiguous run are not stranded for every later pass.
|
||||
func TestCompactDataFilesMergesStrandedRun(t *testing.T) {
|
||||
fs, client := startFakeFiler(t)
|
||||
setup := tableSetup{BucketName: "tb", Namespace: "ns", TableName: "tbl"}
|
||||
populateBoundedTable(t, fs, setup, []boundedFile{
|
||||
{Name: "r1.parquet", IDs: []int64{1}, Lower: 1, Upper: 2, SizeBytes: 4 << 20},
|
||||
{Name: "r2.parquet", IDs: []int64{3}, Lower: 3, Upper: 4, SizeBytes: 4 << 20},
|
||||
{Name: "r3.parquet", IDs: []int64{5}, Lower: 5, Upper: 6, SizeBytes: 6 << 20},
|
||||
{Name: "r4.parquet", IDs: []int64{7}, Lower: 7, Upper: 8, SizeBytes: 3 << 20},
|
||||
{Name: "r5.parquet", IDs: []int64{9}, Lower: 9, Upper: 10, SizeBytes: 3 << 20},
|
||||
})
|
||||
|
||||
handler := NewHandler(nil)
|
||||
config := Config{TargetFileSizeBytes: 10 << 20, MinInputFiles: 3, MaxCommitRetries: 3}
|
||||
result, _, err := handler.compactDataFiles(context.Background(), client, setup.BucketName, setup.tablePath(), config, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("compactDataFiles: %v", err)
|
||||
}
|
||||
if !strings.Contains(result, "compacted 3 files into 1") {
|
||||
t.Fatalf("expected stranded files to merge by size, got %q", result)
|
||||
}
|
||||
}
|
||||
|
||||
// A declared sort order is preserved at the cost of compaction progress:
|
||||
// runs too short to merge are left for later passes instead of being
|
||||
// repacked out of order.
|
||||
func TestCompactDataFilesKeepsStrandedRunForSortedTable(t *testing.T) {
|
||||
fs, client := startFakeFiler(t)
|
||||
setup := tableSetup{BucketName: "tb", Namespace: "ns", TableName: "tbl"}
|
||||
sortOrder, err := table.NewSortOrder(1, []table.SortField{{
|
||||
SourceIDs: []int{1},
|
||||
Transform: iceberg.IdentityTransform{},
|
||||
Direction: table.SortASC,
|
||||
NullOrder: table.NullsFirst,
|
||||
}})
|
||||
if err != nil {
|
||||
t.Fatalf("new sort order: %v", err)
|
||||
}
|
||||
populateBoundedTableSorted(t, fs, setup, []boundedFile{
|
||||
{Name: "r1.parquet", IDs: []int64{1}, Lower: 1, Upper: 2, SizeBytes: 4 << 20},
|
||||
{Name: "r2.parquet", IDs: []int64{3}, Lower: 3, Upper: 4, SizeBytes: 4 << 20},
|
||||
{Name: "r3.parquet", IDs: []int64{5}, Lower: 5, Upper: 6, SizeBytes: 6 << 20},
|
||||
{Name: "r4.parquet", IDs: []int64{7}, Lower: 7, Upper: 8, SizeBytes: 3 << 20},
|
||||
{Name: "r5.parquet", IDs: []int64{9}, Lower: 9, Upper: 10, SizeBytes: 3 << 20},
|
||||
}, sortOrder)
|
||||
|
||||
handler := NewHandler(nil)
|
||||
config := Config{TargetFileSizeBytes: 10 << 20, MinInputFiles: 3, MaxCommitRetries: 3}
|
||||
result, _, err := handler.compactDataFiles(context.Background(), client, setup.BucketName, setup.tablePath(), config, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("compactDataFiles: %v", err)
|
||||
}
|
||||
if !strings.Contains(result, "no files eligible for compaction") {
|
||||
t.Fatalf("expected underfilled ordered runs to be kept, got %q", result)
|
||||
}
|
||||
}
|
||||
|
||||
// A declared sort order only applies when every file carries a bound for its
|
||||
// column; a file without one falls back to inference or unordered packing.
|
||||
func TestResolveCompactionOrderSkipsUncoveredSortField(t *testing.T) {
|
||||
schema := newTestSchema()
|
||||
spec := *iceberg.UnpartitionedSpec
|
||||
sortOrder, err := table.NewSortOrder(1, []table.SortField{{
|
||||
SourceIDs: []int{1},
|
||||
Transform: iceberg.IdentityTransform{},
|
||||
Direction: table.SortASC,
|
||||
NullOrder: table.NullsFirst,
|
||||
}})
|
||||
if err != nil {
|
||||
t.Fatalf("new sort order: %v", err)
|
||||
}
|
||||
meta, err := table.NewMetadata(schema, &spec, sortOrder, "s3://b/p", nil)
|
||||
if err != nil {
|
||||
t.Fatalf("create metadata: %v", err)
|
||||
}
|
||||
|
||||
entry := func(name string, bounded bool) iceberg.ManifestEntry {
|
||||
dfb, err := iceberg.NewDataFileBuilder(spec, iceberg.EntryContentData,
|
||||
"s3://b/data/"+name, iceberg.ParquetFile, map[int]any{}, nil, nil, 2, 16)
|
||||
if err != nil {
|
||||
t.Fatalf("build data file: %v", err)
|
||||
}
|
||||
if bounded {
|
||||
lo, _ := iceberg.Int64Literal(1).MarshalBinary()
|
||||
hi, _ := iceberg.Int64Literal(2).MarshalBinary()
|
||||
dfb.LowerBoundValues(map[int][]byte{1: lo}).UpperBoundValues(map[int][]byte{1: hi})
|
||||
}
|
||||
snapID := int64(1)
|
||||
return iceberg.NewManifestEntry(iceberg.EntryStatusADDED, &snapID, nil, nil, dfb.Build())
|
||||
}
|
||||
|
||||
if order := resolveCompactionOrder(meta, []iceberg.ManifestEntry{
|
||||
entry("a.parquet", true), entry("b.parquet", true),
|
||||
}); order == nil || order.bestEffort {
|
||||
t.Fatalf("covered sort field: got %+v, want strict order", order)
|
||||
}
|
||||
if order := resolveCompactionOrder(meta, []iceberg.ManifestEntry{
|
||||
entry("a.parquet", true), entry("b.parquet", false),
|
||||
}); order != nil {
|
||||
t.Fatalf("uncovered sort field: got %+v, want nil", order)
|
||||
}
|
||||
}
|
||||
// A full-bin repack that produces nothing must not discard the ordered runs
|
||||
// it was meant to replace — they are still valid compaction work.
|
||||
func TestCompactDataFilesKeepsRunsWhenRepackFails(t *testing.T) {
|
||||
fs, client := startFakeFiler(t)
|
||||
setup := tableSetup{BucketName: "tb", Namespace: "ns", TableName: "tbl"}
|
||||
populateBoundedTable(t, fs, setup, []boundedFile{
|
||||
{Name: "r1.parquet", IDs: []int64{1}, Lower: 1, Upper: 2, SizeBytes: 4 << 20},
|
||||
{Name: "r2.parquet", IDs: []int64{3}, Lower: 3, Upper: 4, SizeBytes: 3 << 20},
|
||||
{Name: "r3.parquet", IDs: []int64{5}, Lower: 5, Upper: 6, SizeBytes: 3 << 20},
|
||||
{Name: "r4.parquet", IDs: []int64{7}, Lower: 7, Upper: 8, SizeBytes: 6 << 20},
|
||||
{Name: "r5.parquet", IDs: []int64{9}, Lower: 9, Upper: 10, SizeBytes: 6 << 20},
|
||||
})
|
||||
|
||||
handler := NewHandler(nil)
|
||||
config := Config{TargetFileSizeBytes: 10 << 20, MinInputFiles: 3, MaxCommitRetries: 3}
|
||||
result, _, err := handler.compactDataFiles(context.Background(), client, setup.BucketName, setup.tablePath(), config, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("compactDataFiles: %v", err)
|
||||
}
|
||||
if !strings.Contains(result, "compacted 3 files into 1") {
|
||||
t.Fatalf("expected the ordered run to survive the failed repack, got %q", result)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,97 @@
|
||||
package iceberg
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"testing"
|
||||
|
||||
"github.com/apache/iceberg-go"
|
||||
"github.com/parquet-go/parquet-go"
|
||||
)
|
||||
|
||||
// Compacted files get row groups from the table's
|
||||
// write.parquet.row-group-limit instead of one row group per file.
|
||||
func TestCompactDataFilesWritesSeveralRowGroups(t *testing.T) {
|
||||
fs, client := startFakeFiler(t)
|
||||
setup := tableSetup{BucketName: "tb", Namespace: "ns", TableName: "tbl"}
|
||||
|
||||
type rowsT = struct {
|
||||
ID int64
|
||||
Name string
|
||||
}
|
||||
makeRows := func(lo, hi int64) []rowsT {
|
||||
rows := make([]rowsT, 0, hi-lo)
|
||||
for i := lo; i < hi; i++ {
|
||||
rows = append(rows, rowsT{i, fmt.Sprintf("r%d", i)})
|
||||
}
|
||||
return rows
|
||||
}
|
||||
populateTableWithDeleteFiles(t, fs, setup,
|
||||
[]struct {
|
||||
Name string
|
||||
Rows []rowsT
|
||||
}{
|
||||
{"d1.parquet", makeRows(0, 1500)},
|
||||
{"d2.parquet", makeRows(1500, 3000)},
|
||||
},
|
||||
nil, nil,
|
||||
)
|
||||
|
||||
handler := NewHandler(nil)
|
||||
config := Config{
|
||||
TargetFileSizeBytes: 256 * 1024 * 1024,
|
||||
MinInputFiles: 2,
|
||||
MaxCommitRetries: 3,
|
||||
RowGroupRowLimit: 2000,
|
||||
}
|
||||
if _, _, err := handler.compactDataFiles(context.Background(), client, setup.BucketName, setup.tablePath(), config, nil); err != nil {
|
||||
t.Fatalf("compactDataFiles: %v", err)
|
||||
}
|
||||
|
||||
state, err := loadCurrentMetadata(context.Background(), client, setup.BucketName, setup.tablePath())
|
||||
if err != nil {
|
||||
t.Fatalf("loadCurrentMetadata: %v", err)
|
||||
}
|
||||
df := compactedDataFile(t, client, setup)
|
||||
data, err := loadFileByIcebergPath(context.Background(), client, setup.BucketName, state.DataPath, df.FilePath())
|
||||
if err != nil {
|
||||
t.Fatalf("load compacted file: %v", err)
|
||||
}
|
||||
f, err := parquet.OpenFile(bytes.NewReader(data), int64(len(data)))
|
||||
if err != nil {
|
||||
t.Fatalf("open compacted file: %v", err)
|
||||
}
|
||||
if got := len(f.Metadata().RowGroups); got != 2 {
|
||||
t.Errorf("row groups = %d, want 2 (3000 rows at row-group-limit 2000)", got)
|
||||
}
|
||||
}
|
||||
|
||||
// A configured row-group limit is a cap, not a floor: small explicit limits
|
||||
// are honored instead of being raised to the estimate floor.
|
||||
func TestRowsPerRowGroupHonorsConfiguredLimit(t *testing.T) {
|
||||
entries := make([]iceberg.ManifestEntry, 0, 2)
|
||||
spec := *iceberg.UnpartitionedSpec
|
||||
for _, name := range []string{"a.parquet", "b.parquet"} {
|
||||
dfb, err := iceberg.NewDataFileBuilder(spec, iceberg.EntryContentData,
|
||||
"s3://b/data/"+name, iceberg.ParquetFile, map[int]any{}, nil, nil, 300, 9000)
|
||||
if err != nil {
|
||||
t.Fatalf("build data file: %v", err)
|
||||
}
|
||||
snapID := int64(1)
|
||||
entries = append(entries, iceberg.NewManifestEntry(iceberg.EntryStatusADDED, &snapID, nil, nil, dfb.Build()))
|
||||
}
|
||||
bin := compactionBin{Entries: entries, TotalSize: 18000}
|
||||
|
||||
if got := rowsPerRowGroup(bin, Config{RowGroupRowLimit: 128}); got != 128 {
|
||||
t.Errorf("explicit limit: rowsPerRowGroup = %d, want 128", got)
|
||||
}
|
||||
// 30 bytes/row, byte cap 128MB: the estimate of ~4.4M rows exceeds the 2M
|
||||
// default, so the default holds; the floor lifts a degenerate estimate.
|
||||
if got := rowsPerRowGroup(bin, Config{RowGroupSizeBytes: 1024}); got != minRowGroupRows {
|
||||
t.Errorf("floored estimate: rowsPerRowGroup = %d, want %d", got, minRowGroupRows)
|
||||
}
|
||||
if got := rowsPerRowGroup(bin, Config{}); got != defaultRowGroupRowLimit {
|
||||
t.Errorf("default: rowsPerRowGroup = %d, want %d", got, defaultRowGroupRowLimit)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,231 @@
|
||||
package iceberg
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"math"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
"unicode/utf8"
|
||||
|
||||
"github.com/apache/iceberg-go"
|
||||
"github.com/parquet-go/parquet-go"
|
||||
"github.com/parquet-go/parquet-go/format"
|
||||
)
|
||||
|
||||
// boundTruncateLength is the default metrics mode of the Iceberg reference
|
||||
// implementation, truncate(16), applied to string and binary bounds.
|
||||
const boundTruncateLength = 16
|
||||
|
||||
// columnStats is what a data file's manifest entry records about its columns,
|
||||
// keyed by Iceberg field id. Readers plan scans from it.
|
||||
type columnStats struct {
|
||||
sizes, values, nulls map[int]int64
|
||||
lower, upper map[int][]byte
|
||||
splitOffsets []int64
|
||||
}
|
||||
|
||||
func (s *columnStats) applyTo(b *iceberg.DataFileBuilder) {
|
||||
b.ColumnSizes(s.sizes).ValueCounts(s.values).NullValueCounts(s.nulls).
|
||||
LowerBoundValues(s.lower).UpperBoundValues(s.upper).SplitOffsets(s.splitOffsets)
|
||||
}
|
||||
|
||||
// collectColumnStats reads back the footer of a Parquet file the compactor
|
||||
// wrote. Counts and sizes are summed over the row groups and the bounds are
|
||||
// the minimum and maximum of the row group statistics, encoded with the
|
||||
// spec's single-value serialization. A column whose bounds cannot be read
|
||||
// for every row group, or converted exactly, gets no bounds.
|
||||
func collectColumnStats(data []byte, schema *iceberg.Schema) (*columnStats, error) {
|
||||
file, err := parquet.OpenFile(bytes.NewReader(data), int64(len(data)),
|
||||
parquet.SkipPageIndex(true), parquet.SkipBloomFilters(true))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var leaves []*parquet.Column
|
||||
var walk func(*parquet.Column)
|
||||
walk = func(c *parquet.Column) {
|
||||
if c.Leaf() {
|
||||
leaves = append(leaves, c)
|
||||
return
|
||||
}
|
||||
for _, child := range c.Columns() {
|
||||
walk(child)
|
||||
}
|
||||
}
|
||||
walk(file.Root())
|
||||
|
||||
type bounds struct {
|
||||
field iceberg.NestedField
|
||||
typ parquet.Type
|
||||
min, max parquet.Value
|
||||
seen bool
|
||||
invalid bool
|
||||
}
|
||||
stats := &columnStats{sizes: map[int]int64{}, values: map[int]int64{}, nulls: map[int]int64{},
|
||||
lower: map[int][]byte{}, upper: map[int][]byte{}}
|
||||
agg := map[int]*bounds{}
|
||||
meta := file.Metadata()
|
||||
for rgIdx, rg := range file.RowGroups() {
|
||||
stats.splitOffsets = append(stats.splitOffsets, rowGroupOffset(&meta.RowGroups[rgIdx]))
|
||||
for col, cc := range rg.ColumnChunks() {
|
||||
field, ok := icebergFieldOf(leaves[col], schema)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
stats.sizes[field.ID] += meta.RowGroups[rgIdx].Columns[col].MetaData.TotalCompressedSize
|
||||
stats.values[field.ID] += cc.NumValues()
|
||||
b := agg[field.ID]
|
||||
if b == nil {
|
||||
b = &bounds{field: field, typ: leaves[col].Type()}
|
||||
agg[field.ID] = b
|
||||
}
|
||||
chunk, ok := cc.(*parquet.FileColumnChunk)
|
||||
if !ok {
|
||||
b.invalid = true
|
||||
continue
|
||||
}
|
||||
stats.nulls[field.ID] += chunk.NullCount()
|
||||
if chunk.NullCount() == chunk.NumValues() {
|
||||
continue // an all-null chunk constrains nothing
|
||||
}
|
||||
lo, hi, ok := chunk.Bounds()
|
||||
switch {
|
||||
case !ok:
|
||||
b.invalid = true
|
||||
case !b.seen:
|
||||
b.min, b.max, b.seen = lo.Clone(), hi.Clone(), true
|
||||
default:
|
||||
if b.typ.Compare(lo, b.min) < 0 {
|
||||
b.min = lo.Clone()
|
||||
}
|
||||
if b.typ.Compare(hi, b.max) > 0 {
|
||||
b.max = hi.Clone()
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
for id, b := range agg {
|
||||
if b.invalid || !b.seen {
|
||||
continue
|
||||
}
|
||||
lo, okLo := boundLiteral(b.min, b.typ, b.field.Type)
|
||||
hi, okHi := boundLiteral(b.max, b.typ, b.field.Type)
|
||||
if !okLo || !okHi {
|
||||
continue
|
||||
}
|
||||
lo, hi, okHi = truncateBounds(lo, hi)
|
||||
if raw, err := lo.MarshalBinary(); err == nil {
|
||||
stats.lower[id] = raw
|
||||
}
|
||||
if raw, err := hi.MarshalBinary(); err == nil && okHi {
|
||||
stats.upper[id] = raw
|
||||
}
|
||||
}
|
||||
sort.Slice(stats.splitOffsets, func(i, j int) bool { return stats.splitOffsets[i] < stats.splitOffsets[j] })
|
||||
return stats, nil
|
||||
}
|
||||
|
||||
// icebergFieldOf finds the Iceberg field of a Parquet leaf column by the field
|
||||
// id the writer stored, or by its dotted path in a file written without ids.
|
||||
// Only primitive fields outside lists and maps carry statistics.
|
||||
func icebergFieldOf(leaf *parquet.Column, schema *iceberg.Schema) (iceberg.NestedField, bool) {
|
||||
var field iceberg.NestedField
|
||||
var ok bool
|
||||
if leaf.MaxRepetitionLevel() > 0 {
|
||||
return field, false
|
||||
}
|
||||
if id := leaf.ID(); id > 0 {
|
||||
field, ok = schema.FindFieldByID(id)
|
||||
} else {
|
||||
field, ok = schema.FindFieldByName(strings.Join(leaf.Path(), "."))
|
||||
}
|
||||
if !ok {
|
||||
return field, false
|
||||
}
|
||||
_, primitive := field.Type.(iceberg.PrimitiveType)
|
||||
return field, primitive
|
||||
}
|
||||
|
||||
// rowGroupOffset is where a row group's bytes start: the split offset readers
|
||||
// plan one task per row group from.
|
||||
func rowGroupOffset(rg *format.RowGroup) int64 {
|
||||
if rg.FileOffset > 0 || len(rg.Columns) == 0 {
|
||||
return rg.FileOffset
|
||||
}
|
||||
first := rg.Columns[0].MetaData
|
||||
if first.DictionaryPageOffset > 0 && first.DictionaryPageOffset < first.DataPageOffset {
|
||||
return first.DictionaryPageOffset
|
||||
}
|
||||
return first.DataPageOffset
|
||||
}
|
||||
|
||||
// boundLiteral converts a Parquet statistics value into a literal of the
|
||||
// column's Iceberg type, when the conversion is exact. Other types (decimal,
|
||||
// time, uuid, fixed, timestamps not in the type's unit) get no bound.
|
||||
func boundLiteral(v parquet.Value, pt parquet.Type, t iceberg.Type) (iceberg.Literal, bool) {
|
||||
if v.IsNull() {
|
||||
return nil, false
|
||||
}
|
||||
switch t.(type) {
|
||||
case iceberg.BooleanType:
|
||||
return iceberg.BoolLiteral(v.Boolean()), v.Kind() == parquet.Boolean
|
||||
case iceberg.Int32Type:
|
||||
return iceberg.Int32Literal(v.Int32()), v.Kind() == parquet.Int32
|
||||
case iceberg.Int64Type:
|
||||
switch v.Kind() {
|
||||
case parquet.Int64:
|
||||
return iceberg.Int64Literal(v.Int64()), true
|
||||
case parquet.Int32:
|
||||
return iceberg.Int64Literal(int64(v.Int32())), true
|
||||
}
|
||||
case iceberg.Float32Type:
|
||||
f := v.Float()
|
||||
return iceberg.Float32Literal(f), v.Kind() == parquet.Float && !math.IsNaN(float64(f))
|
||||
case iceberg.Float64Type:
|
||||
f := v.Double()
|
||||
return iceberg.Float64Literal(f), v.Kind() == parquet.Double && !math.IsNaN(f)
|
||||
case iceberg.DateType:
|
||||
return iceberg.DateLiteral(iceberg.Date(v.Int32())), v.Kind() == parquet.Int32
|
||||
case iceberg.TimestampType, iceberg.TimestampTzType:
|
||||
return iceberg.TimestampLiteral(iceberg.Timestamp(v.Int64())), v.Kind() == parquet.Int64 && timeUnit(pt) == time.Microsecond
|
||||
case iceberg.TimestampNsType:
|
||||
return iceberg.TimestampNsLiteral(iceberg.TimestampNano(v.Int64())), v.Kind() == parquet.Int64 && timeUnit(pt) == time.Nanosecond
|
||||
case iceberg.StringType:
|
||||
b := v.ByteArray()
|
||||
return iceberg.StringLiteral(string(b)), v.Kind() == parquet.ByteArray && utf8.Valid(b)
|
||||
case iceberg.BinaryType:
|
||||
return iceberg.BinaryLiteral(bytes.Clone(v.ByteArray())), v.Kind() == parquet.ByteArray
|
||||
}
|
||||
return nil, false
|
||||
}
|
||||
|
||||
func timeUnit(pt parquet.Type) time.Duration {
|
||||
if pt == nil || pt.LogicalType() == nil {
|
||||
return 0
|
||||
}
|
||||
if ts, ok := pt.LogicalType().Value.(*format.TimestampType); ok && ts.Unit.Value != nil {
|
||||
return ts.Unit.Value.Duration()
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
// truncateBounds applies truncate(16) to string and binary bounds. A prefix is
|
||||
// still a lower bound; a shortened upper bound is not, so a value longer than
|
||||
// the limit keeps no upper bound (the reference implementation increments the
|
||||
// last character instead).
|
||||
func truncateBounds(lo, hi iceberg.Literal) (iceberg.Literal, iceberg.Literal, bool) {
|
||||
switch l := lo.(type) {
|
||||
case iceberg.StringLiteral:
|
||||
if r := []rune(string(l)); len(r) > boundTruncateLength {
|
||||
lo = iceberg.StringLiteral(string(r[:boundTruncateLength]))
|
||||
}
|
||||
return lo, hi, utf8.RuneCountInString(string(hi.(iceberg.StringLiteral))) <= boundTruncateLength
|
||||
case iceberg.BinaryLiteral:
|
||||
if len(l) > boundTruncateLength {
|
||||
lo = l[:boundTruncateLength]
|
||||
}
|
||||
return lo, hi, len(hi.(iceberg.BinaryLiteral)) <= boundTruncateLength
|
||||
}
|
||||
return lo, hi, true
|
||||
}
|
||||
@@ -0,0 +1,148 @@
|
||||
package iceberg
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"github.com/apache/iceberg-go"
|
||||
"github.com/parquet-go/parquet-go"
|
||||
|
||||
filer_pb "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
)
|
||||
|
||||
// compactedDataFile returns the data file the current snapshot's data
|
||||
// manifest references after a compaction run.
|
||||
func compactedDataFile(t *testing.T, client filer_pb.SeaweedFilerClient, setup tableSetup) iceberg.DataFile {
|
||||
t.Helper()
|
||||
found := liveDataFiles(t, client, setup)
|
||||
if len(found) != 1 {
|
||||
t.Fatalf("expected 1 live data file after compaction, got %d", len(found))
|
||||
}
|
||||
return found[0]
|
||||
}
|
||||
|
||||
func liveDataFiles(t *testing.T, client filer_pb.SeaweedFilerClient, setup tableSetup) []iceberg.DataFile {
|
||||
t.Helper()
|
||||
state, err := loadCurrentMetadata(context.Background(), client, setup.BucketName, setup.tablePath())
|
||||
if err != nil {
|
||||
t.Fatalf("loadCurrentMetadata: %v", err)
|
||||
}
|
||||
manifests, err := loadCurrentManifests(context.Background(), client, setup.BucketName, state.DataPath, state.Metadata)
|
||||
if err != nil {
|
||||
t.Fatalf("loadCurrentManifests: %v", err)
|
||||
}
|
||||
var found []iceberg.DataFile
|
||||
for _, mf := range manifests {
|
||||
if mf.ManifestContent() != iceberg.ManifestContentData {
|
||||
continue
|
||||
}
|
||||
manifestData, err := loadFileByIcebergPath(context.Background(), client, setup.BucketName, state.DataPath, mf.FilePath())
|
||||
if err != nil {
|
||||
t.Fatalf("load manifest: %v", err)
|
||||
}
|
||||
entries, err := iceberg.ReadManifest(mf, bytes.NewReader(manifestData), true)
|
||||
if err != nil {
|
||||
t.Fatalf("read manifest: %v", err)
|
||||
}
|
||||
for _, entry := range entries {
|
||||
if entry.Status() != iceberg.EntryStatusDELETED {
|
||||
found = append(found, entry.DataFile())
|
||||
}
|
||||
}
|
||||
}
|
||||
return found
|
||||
}
|
||||
|
||||
func compactedFileBytes(t *testing.T, client filer_pb.SeaweedFilerClient, setup tableSetup, filePath string) []byte {
|
||||
t.Helper()
|
||||
state, err := loadCurrentMetadata(context.Background(), client, setup.BucketName, setup.tablePath())
|
||||
if err != nil {
|
||||
t.Fatalf("loadCurrentMetadata: %v", err)
|
||||
}
|
||||
data, err := loadFileByIcebergPath(context.Background(), client, setup.BucketName, state.DataPath, filePath)
|
||||
if err != nil {
|
||||
t.Fatalf("load compacted file: %v", err)
|
||||
}
|
||||
return data
|
||||
}
|
||||
|
||||
func compactedIDs(t *testing.T, client filer_pb.SeaweedFilerClient, setup tableSetup, filePath string) []int64 {
|
||||
t.Helper()
|
||||
data := compactedFileBytes(t, client, setup, filePath)
|
||||
f, err := parquet.OpenFile(bytes.NewReader(data), int64(len(data)))
|
||||
if err != nil {
|
||||
t.Fatalf("open compacted file: %v", err)
|
||||
}
|
||||
var ids []int64
|
||||
for _, rg := range f.RowGroups() {
|
||||
rows := rg.Rows()
|
||||
buf := make([]parquet.Row, 8)
|
||||
for {
|
||||
n, err := rows.ReadRows(buf)
|
||||
for _, row := range buf[:n] {
|
||||
ids = append(ids, row[0].Int64())
|
||||
}
|
||||
if err != nil {
|
||||
break
|
||||
}
|
||||
}
|
||||
rows.Close()
|
||||
}
|
||||
return ids
|
||||
}
|
||||
|
||||
// A compacted file's manifest entry carries the column metrics readers
|
||||
// prune files by: counts, sizes, bounds and split offsets read back
|
||||
// from the file's own footer.
|
||||
func TestCompactDataFilesRecordsColumnStatistics(t *testing.T) {
|
||||
fs, client := startFakeFiler(t)
|
||||
setup := tableSetup{BucketName: "tb", Namespace: "ns", TableName: "tbl"}
|
||||
populateTableWithDeleteFiles(t, fs, setup,
|
||||
[]struct {
|
||||
Name string
|
||||
Rows []struct {
|
||||
ID int64
|
||||
Name string
|
||||
}
|
||||
}{
|
||||
{"d1.parquet", []struct {
|
||||
ID int64
|
||||
Name string
|
||||
}{{1, "a"}, {2, "b"}}},
|
||||
{"d2.parquet", []struct {
|
||||
ID int64
|
||||
Name string
|
||||
}{{3, "c"}}},
|
||||
},
|
||||
nil, nil,
|
||||
)
|
||||
|
||||
handler := NewHandler(nil)
|
||||
config := Config{TargetFileSizeBytes: 256 * 1024 * 1024, MinInputFiles: 2, MaxCommitRetries: 3}
|
||||
if _, _, err := handler.compactDataFiles(context.Background(), client, setup.BucketName, setup.tablePath(), config, nil); err != nil {
|
||||
t.Fatalf("compactDataFiles: %v", err)
|
||||
}
|
||||
|
||||
df := compactedDataFile(t, client, setup)
|
||||
if df.ValueCounts()[1] != 3 {
|
||||
t.Errorf("value_counts[1] = %d, want 3", df.ValueCounts()[1])
|
||||
}
|
||||
if len(df.LowerBoundValues()) == 0 || len(df.UpperBoundValues()) == 0 {
|
||||
t.Fatal("compacted file has no bounds")
|
||||
}
|
||||
lo, err := iceberg.LiteralFromBytes(iceberg.PrimitiveTypes.Int64, df.LowerBoundValues()[1])
|
||||
if err != nil {
|
||||
t.Fatalf("decode lower bound: %v", err)
|
||||
}
|
||||
hi, err := iceberg.LiteralFromBytes(iceberg.PrimitiveTypes.Int64, df.UpperBoundValues()[1])
|
||||
if err != nil {
|
||||
t.Fatalf("decode upper bound: %v", err)
|
||||
}
|
||||
if lo.(iceberg.Int64Literal) != 1 || hi.(iceberg.Int64Literal) != 3 {
|
||||
t.Errorf("id bounds = %v..%v, want 1..3", lo, hi)
|
||||
}
|
||||
if len(df.SplitOffsets()) == 0 {
|
||||
t.Error("compacted file has no split_offsets")
|
||||
}
|
||||
}
|
||||
@@ -130,6 +130,10 @@ type Config struct {
|
||||
SortMaxInputBytes int64
|
||||
SortBufferRows int64
|
||||
SortSpillDir string
|
||||
// Row groups of a compacted file, from the table's write.parquet.*
|
||||
// properties; zero means the default.
|
||||
RowGroupSizeBytes int64
|
||||
RowGroupRowLimit int64
|
||||
}
|
||||
|
||||
// ParseConfig extracts an iceberg maintenance Config from plugin config values.
|
||||
|
||||
@@ -477,7 +477,7 @@ func hasEligibleCompaction(
|
||||
}
|
||||
|
||||
targetSize := compactionTargetSizeForPlan(config, rewritePlan)
|
||||
bins := buildCompactionBins(candidateEntries, targetSize, minInputFiles)
|
||||
bins := buildCompactionBins(candidateEntries, targetSize, minInputFiles, meta)
|
||||
bins = filterCompactionBinsByPlan(bins, config, rewritePlan)
|
||||
return len(bins) > 0, nil
|
||||
}
|
||||
|
||||
@@ -701,7 +701,7 @@ func TestBuildCompactionBins(t *testing.T) {
|
||||
{path: "data/f3.parquet", size: 4096, partition: map[int]any{}},
|
||||
})
|
||||
|
||||
bins := buildCompactionBins(entries, targetSize, minFiles)
|
||||
bins := buildCompactionBins(entries, targetSize, minFiles, nil)
|
||||
if len(bins) != 1 {
|
||||
t.Fatalf("expected 1 bin, got %d", len(bins))
|
||||
}
|
||||
@@ -720,7 +720,7 @@ func TestBuildCompactionBinsFiltersLargeFiles(t *testing.T) {
|
||||
{path: "data/large.parquet", size: 5000, partition: map[int]any{}},
|
||||
})
|
||||
|
||||
bins := buildCompactionBins(entries, targetSize, minFiles)
|
||||
bins := buildCompactionBins(entries, targetSize, minFiles, nil)
|
||||
if len(bins) != 1 {
|
||||
t.Fatalf("expected 1 bin, got %d", len(bins))
|
||||
}
|
||||
@@ -756,7 +756,7 @@ func TestBuildCompactionBinsLowercaseParquetFormat(t *testing.T) {
|
||||
entries[0] = lowercaseFormatEntry{entries[0]}
|
||||
entries[1] = lowercaseFormatEntry{entries[1]}
|
||||
|
||||
bins := buildCompactionBins(entries, targetSize, minFiles)
|
||||
bins := buildCompactionBins(entries, targetSize, minFiles, nil)
|
||||
if len(bins) != 1 {
|
||||
t.Fatalf("expected 1 bin, got %d", len(bins))
|
||||
}
|
||||
@@ -774,7 +774,7 @@ func TestBuildCompactionBinsMinFilesThreshold(t *testing.T) {
|
||||
{path: "data/f2.parquet", size: 2048, partition: map[int]any{}},
|
||||
})
|
||||
|
||||
bins := buildCompactionBins(entries, targetSize, minFiles)
|
||||
bins := buildCompactionBins(entries, targetSize, minFiles, nil)
|
||||
if len(bins) != 0 {
|
||||
t.Errorf("expected 0 bins (below min threshold), got %d", len(bins))
|
||||
}
|
||||
@@ -801,7 +801,7 @@ func TestBuildCompactionBinsMultiplePartitions(t *testing.T) {
|
||||
{path: "data/b3.parquet", size: 4096, partition: partB, partitionSpec: &partitionSpec},
|
||||
})
|
||||
|
||||
bins := buildCompactionBins(entries, targetSize, minFiles)
|
||||
bins := buildCompactionBins(entries, targetSize, minFiles, nil)
|
||||
if len(bins) != 2 {
|
||||
t.Fatalf("expected 2 bins (one per partition), got %d", len(bins))
|
||||
}
|
||||
@@ -977,7 +977,7 @@ func TestBuildCompactionBinsMultipleSpecs(t *testing.T) {
|
||||
{path: "data/s1-f2.parquet", size: 2048, partition: map[int]any{}, specID: 1},
|
||||
}, partSpecs)
|
||||
|
||||
bins := buildCompactionBins(entries, targetSize, minFiles)
|
||||
bins := buildCompactionBins(entries, targetSize, minFiles, nil)
|
||||
if len(bins) != 2 {
|
||||
t.Fatalf("expected 2 bins (one per spec), got %d", len(bins))
|
||||
}
|
||||
@@ -1011,7 +1011,7 @@ func TestBuildCompactionBinsSingleSpec(t *testing.T) {
|
||||
{path: "data/f3.parquet", size: 4096, partition: map[int]any{}, specID: 0},
|
||||
}, partSpecs)
|
||||
|
||||
bins := buildCompactionBins(entries, targetSize, minFiles)
|
||||
bins := buildCompactionBins(entries, targetSize, minFiles, nil)
|
||||
if len(bins) != 1 {
|
||||
t.Fatalf("expected 1 bin, got %d", len(bins))
|
||||
}
|
||||
@@ -1303,7 +1303,7 @@ func TestMergeParquetFilesWithPositionDeletes(t *testing.T) {
|
||||
|
||||
merged, count, err := mergeParquetFiles(
|
||||
context.Background(), client, "test-bucket", "ns/tbl",
|
||||
entries, posDeletes, nil, nil,
|
||||
entries, posDeletes, nil, nil, 0,
|
||||
)
|
||||
if err != nil {
|
||||
t.Fatalf("mergeParquetFiles: %v", err)
|
||||
@@ -1419,7 +1419,7 @@ func TestMergeParquetFilesWithEqualityDeletes(t *testing.T) {
|
||||
|
||||
merged, count, err := mergeParquetFiles(
|
||||
context.Background(), client, "test-bucket", "ns/tbl",
|
||||
entries, nil, eqGroups, schema,
|
||||
entries, nil, eqGroups, schema, 0,
|
||||
)
|
||||
if err != nil {
|
||||
t.Fatalf("mergeParquetFiles: %v", err)
|
||||
@@ -1479,7 +1479,7 @@ func TestMergeParquetFilesDictionaryEncodedInput(t *testing.T) {
|
||||
|
||||
merged, count, err := mergeParquetFiles(
|
||||
context.Background(), client, "test-bucket", "ns/tbl",
|
||||
entries, nil, nil, nil,
|
||||
entries, nil, nil, nil, 0,
|
||||
)
|
||||
if err != nil {
|
||||
t.Fatalf("mergeParquetFiles: %v", err)
|
||||
|
||||
@@ -19,6 +19,8 @@ import (
|
||||
const (
|
||||
propTargetFileSize = "write.target-file-size-bytes"
|
||||
propDeleteTargetFileSize = "write.delete.target-file-size-bytes"
|
||||
propRowGroupSize = "write.parquet.row-group-size-bytes"
|
||||
propRowGroupRowLimit = "write.parquet.row-group-limit"
|
||||
propMaxSnapshotAgeMs = "history.expire.max-snapshot-age-ms"
|
||||
propMinSnapshotsToKeep = "history.expire.min-snapshots-to-keep"
|
||||
)
|
||||
@@ -50,6 +52,12 @@ func applyTableProperties(cfg Config, props iceberg.Properties) Config {
|
||||
if v, ok := propInt64(props, propDeleteTargetFileSize); ok {
|
||||
cfg.DeleteTargetFileSizeBytes = v
|
||||
}
|
||||
if v, ok := propInt64(props, propRowGroupSize); ok {
|
||||
cfg.RowGroupSizeBytes = v
|
||||
}
|
||||
if v, ok := propInt64(props, propRowGroupRowLimit); ok {
|
||||
cfg.RowGroupRowLimit = v
|
||||
}
|
||||
if v, ok := propInt64(props, propMaxSnapshotAgeMs); ok {
|
||||
cfg.SnapshotRetentionMs = v
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user