1use std::sync::Arc;
9
10use async_trait::async_trait;
11use bytes::Bytes;
12use futures::StreamExt;
13use lance_core::utils::tracing::{
14 AUDIT_MODE_CREATE, AUDIT_MODE_DELETE, AUDIT_TYPE_MANIFEST, TRACE_FILE_AUDIT,
15};
16use lance_core::{Error, Result};
17use lance_io::object_store::ObjectStore;
18use log::warn;
19use object_store::ObjectMeta;
20use object_store::ObjectStoreExt;
21use object_store::{Error as ObjectStoreError, ObjectStore as OSObjectStore, path::Path};
22use tracing::info;
23
24use super::{
25 MANIFEST_EXTENSION, ManifestLocation, ManifestNamingScheme, current_manifest_path,
26 default_resolve_version, make_staging_manifest_path, write_version_hint,
27};
28use crate::format::{IndexMetadata, Manifest, Transaction};
29use crate::io::commit::{CommitError, CommitHandler};
30
31#[async_trait]
44pub trait ExternalManifestStore: std::fmt::Debug + Send + Sync {
45 async fn get(&self, base_uri: &str, version: u64) -> Result<String>;
47
48 async fn get_manifest_location(
49 &self,
50 base_uri: &str,
51 version: u64,
52 ) -> Result<ManifestLocation> {
53 let path = self.get(base_uri, version).await?;
54 let path = Path::parse(&path).map_err(|e| Error::invalid_input(e.to_string()))?;
55 let naming_scheme = detect_naming_scheme_from_path(&path)?;
56 Ok(ManifestLocation {
57 version,
58 path,
59 size: None,
60 naming_scheme,
61 e_tag: None,
62 })
63 }
64
65 async fn get_latest_version(&self, base_uri: &str) -> Result<Option<(u64, String)>>;
69
70 async fn get_latest_manifest_location(
76 &self,
77 base_uri: &str,
78 ) -> Result<Option<ManifestLocation>> {
79 self.get_latest_version(base_uri).await.and_then(|res| {
80 res.map(|(version, uri)| {
81 let path = Path::parse(&uri).map_err(|e| Error::invalid_input(e.to_string()))?;
82 let naming_scheme = detect_naming_scheme_from_path(&path)?;
83 Ok(ManifestLocation {
84 version,
85 path,
86 size: None,
87 naming_scheme,
88 e_tag: None,
89 })
90 })
91 .transpose()
92 })
93 }
94
95 #[allow(clippy::too_many_arguments)]
104 async fn put(
105 &self,
106 base_path: &Path,
107 version: u64,
108 staging_path: &Path,
109 size: u64,
110 e_tag: Option<String>,
111 object_store: &dyn OSObjectStore,
112 naming_scheme: ManifestNamingScheme,
113 ) -> Result<ManifestLocation> {
114 self.put_if_not_exists(
118 base_path.as_ref(),
119 version,
120 staging_path.as_ref(),
121 size,
122 e_tag.clone(),
123 )
124 .await?;
125
126 let final_path = naming_scheme.manifest_path(base_path, version);
128 let copied = match copy_size_aware(object_store, staging_path, &final_path, size).await {
129 Ok(_) => true,
130 Err(ObjectStoreError::NotFound { .. }) => false,
131 Err(e) => return Err(e.into()),
132 };
133 if copied {
134 info!(target: TRACE_FILE_AUDIT, mode=AUDIT_MODE_CREATE, r#type=AUDIT_TYPE_MANIFEST, path = final_path.as_ref());
135 }
136
137 let final_meta = object_store.head(&final_path).await?;
140 let final_size = final_meta.size;
141 let final_e_tag = final_meta.e_tag;
142
143 let location = ManifestLocation {
144 version,
145 path: final_path.clone(),
146 size: Some(final_size),
147 naming_scheme,
148 e_tag: final_e_tag.clone(),
149 };
150
151 if !copied {
152 return Ok(location);
153 }
154
155 self.put_if_exists(
157 base_path.as_ref(),
158 version,
159 final_path.as_ref(),
160 final_size,
161 final_e_tag,
162 )
163 .await?;
164
165 match object_store.delete(staging_path).await {
167 Ok(_) => {}
168 Err(ObjectStoreError::NotFound { .. }) => {}
169 Err(e) => return Err(e.into()),
170 }
171 info!(target: TRACE_FILE_AUDIT, mode=AUDIT_MODE_DELETE, r#type=AUDIT_TYPE_MANIFEST, path = staging_path.as_ref());
172
173 Ok(location)
174 }
175
176 async fn put_if_not_exists(
178 &self,
179 base_uri: &str,
180 version: u64,
181 path: &str,
182 size: u64,
183 e_tag: Option<String>,
184 ) -> Result<()>;
185
186 async fn put_if_exists(
188 &self,
189 base_uri: &str,
190 version: u64,
191 path: &str,
192 size: u64,
193 e_tag: Option<String>,
194 ) -> Result<()>;
195
196 async fn delete(&self, _base_uri: &str) -> Result<()> {
198 Ok(())
199 }
200}
201
202pub(crate) fn detect_naming_scheme_from_path(path: &Path) -> Result<ManifestNamingScheme> {
203 path.filename()
204 .and_then(|name| {
205 ManifestNamingScheme::detect_scheme(name)
206 .or_else(|| Some(ManifestNamingScheme::detect_scheme_staging(name)))
207 })
208 .ok_or_else(|| {
209 Error::corrupt_file(
210 path.clone(),
211 "Path does not follow known manifest naming convention.",
212 )
213 })
214}
215
216const MAX_SERVER_SIDE_COPY_BYTES: u64 = 5 * 1024 * 1024 * 1024;
225
226const COPY_REWRITE_PART_SIZE: usize = 100 * 1024 * 1024;
232
233async fn copy_size_aware(
254 store: &dyn OSObjectStore,
255 from: &Path,
256 to: &Path,
257 size: u64,
258) -> std::result::Result<(), ObjectStoreError> {
259 if size < MAX_SERVER_SIDE_COPY_BYTES {
260 store.copy(from, to).await
261 } else {
262 copy_via_read_rewrite(store, from, to).await
263 }
264}
265
266async fn copy_via_read_rewrite(
273 store: &dyn OSObjectStore,
274 from: &Path,
275 to: &Path,
276) -> std::result::Result<(), ObjectStoreError> {
277 let mut stream = store.get(from).await?.into_stream();
279
280 let mut upload = store.put_multipart(to).await?;
290 let mut part_buf: Vec<u8> = Vec::with_capacity(COPY_REWRITE_PART_SIZE);
291
292 while let Some(chunk) = stream.next().await {
293 let chunk = match chunk {
294 Ok(b) => b,
295 Err(e) => {
296 let _ = upload.abort().await;
297 return Err(e);
298 }
299 };
300 let mut offset = 0;
306 while offset < chunk.len() {
307 let want = COPY_REWRITE_PART_SIZE - part_buf.len();
308 let take = want.min(chunk.len() - offset);
309 part_buf.extend_from_slice(&chunk[offset..offset + take]);
310 offset += take;
311
312 if part_buf.len() >= COPY_REWRITE_PART_SIZE {
313 let payload =
314 std::mem::replace(&mut part_buf, Vec::with_capacity(COPY_REWRITE_PART_SIZE));
315 if let Err(e) = upload.put_part(Bytes::from(payload).into()).await {
316 let _ = upload.abort().await;
317 return Err(e);
318 }
319 }
320 }
321 }
322
323 if !part_buf.is_empty()
326 && let Err(e) = upload.put_part(Bytes::from(part_buf).into()).await
327 {
328 let _ = upload.abort().await;
329 return Err(e);
330 }
331
332 if let Err(e) = upload.complete().await {
333 let _ = upload.abort().await;
334 return Err(e);
335 }
336 Ok(())
337}
338
339#[derive(Debug)]
343pub struct ExternalManifestCommitHandler {
344 pub external_manifest_store: Arc<dyn ExternalManifestStore>,
345}
346
347impl ExternalManifestCommitHandler {
348 async fn verify_finalized_manifest_location(
349 &self,
350 base_path: &Path,
351 location: ManifestLocation,
352 object_store: &dyn OSObjectStore,
353 ) -> std::result::Result<ManifestLocation, Error> {
354 match object_store.head(&location.path).await {
355 Ok(ObjectMeta { size, e_tag, .. }) => {
356 let ManifestLocation {
357 version,
358 path,
359 size: expected_size,
360 naming_scheme,
361 e_tag: expected_e_tag,
362 } = location;
363
364 let size = match expected_size {
365 Some(expected_size) if expected_size != size => {
366 return Err(Error::corrupt_file(
367 path,
368 format!(
369 "Manifest size mismatch for version {}: external store expected {}, object store returned {}",
370 version, expected_size, size
371 ),
372 ));
373 }
374 Some(expected_size) => Some(expected_size),
375 None => Some(size),
376 };
377
378 let e_tag = match expected_e_tag {
379 Some(expected_e_tag) => {
380 if e_tag.as_ref() != Some(&expected_e_tag) {
381 return Err(Error::corrupt_file(
382 path,
383 format!(
384 "Manifest e_tag mismatch for version {}: external store expected {:?}, object store returned {:?}",
385 version, expected_e_tag, e_tag
386 ),
387 ));
388 }
389 Some(expected_e_tag)
390 }
391 None => e_tag,
392 };
393
394 Ok(ManifestLocation {
395 version,
396 path,
397 size,
398 naming_scheme,
399 e_tag,
400 })
401 }
402 Err(ObjectStoreError::NotFound { .. }) => {
403 default_resolve_version(base_path, location.version, object_store).await
406 }
407 Err(e) => Err(e.into()),
408 }
409 }
410
411 #[allow(clippy::too_many_arguments)]
421 async fn finalize_manifest(
422 &self,
423 base_path: &Path,
424 staging_manifest_path: &Path,
425 version: u64,
426 size: u64,
427 store: &dyn OSObjectStore,
428 naming_scheme: ManifestNamingScheme,
429 ) -> std::result::Result<ManifestLocation, Error> {
430 let final_manifest_path = naming_scheme.manifest_path(base_path, version);
432
433 let copied =
434 match copy_size_aware(store, staging_manifest_path, &final_manifest_path, size).await {
435 Ok(_) => true,
436 Err(ObjectStoreError::NotFound { .. }) => false, Err(e) => return Err(e.into()),
438 };
439 if copied {
440 info!(target: TRACE_FILE_AUDIT, mode=AUDIT_MODE_CREATE, r#type=AUDIT_TYPE_MANIFEST, path = final_manifest_path.as_ref());
441 }
442
443 let final_meta = store.head(&final_manifest_path).await?;
446 let final_size = final_meta.size;
447 let final_e_tag = final_meta.e_tag;
448
449 let location = ManifestLocation {
450 version,
451 path: final_manifest_path,
452 size: Some(final_size),
453 naming_scheme,
454 e_tag: final_e_tag,
455 };
456
457 if !copied {
458 return Ok(location);
459 }
460
461 self.external_manifest_store
463 .put_if_exists(
464 base_path.as_ref(),
465 version,
466 location.path.as_ref(),
467 final_size,
468 location.e_tag.clone(),
469 )
470 .await?;
471
472 match store.delete(staging_manifest_path).await {
474 Ok(_) => {}
475 Err(ObjectStoreError::NotFound { .. }) => {}
476 Err(e) => return Err(e.into()),
477 }
478 info!(target: TRACE_FILE_AUDIT, mode=AUDIT_MODE_DELETE, r#type=AUDIT_TYPE_MANIFEST, path = staging_manifest_path.as_ref());
479
480 Ok(location)
481 }
482}
483
484#[async_trait]
485impl CommitHandler for ExternalManifestCommitHandler {
486 async fn resolve_latest_location(
487 &self,
488 base_path: &Path,
489 object_store: &ObjectStore,
490 ) -> std::result::Result<ManifestLocation, Error> {
491 let location = self
492 .external_manifest_store
493 .get_latest_manifest_location(base_path.as_ref())
494 .await?;
495
496 match location {
497 Some(location) => {
498 if location.path.extension() == Some(MANIFEST_EXTENSION) {
499 return self
500 .verify_finalized_manifest_location(
501 base_path,
502 location,
503 object_store.inner.as_ref(),
504 )
505 .await;
506 }
507
508 let ManifestLocation {
509 version,
510 path,
511 size,
512 naming_scheme,
513 e_tag: _,
514 } = location;
515
516 let size = if let Some(size) = size {
517 size
518 } else {
519 match object_store.inner.head(&path).await {
520 Ok(meta) => meta.size,
521 Err(ObjectStoreError::NotFound { .. }) => {
522 let new_location = self
524 .external_manifest_store
525 .get_manifest_location(base_path.as_ref(), version)
526 .await?;
527 return Ok(new_location);
528 }
529 Err(e) => return Err(e.into()),
530 }
531 };
532
533 let final_location = self
534 .finalize_manifest(
535 base_path,
536 &path,
537 version,
538 size,
539 &object_store.inner,
540 naming_scheme,
541 )
542 .await?;
543
544 Ok(final_location)
545 }
546 None => current_manifest_path(object_store, base_path).await,
549 }
550 }
551
552 async fn resolve_version_location(
553 &self,
554 base_path: &Path,
555 version: u64,
556 object_store: &dyn OSObjectStore,
557 ) -> std::result::Result<ManifestLocation, Error> {
558 let location_res = self
559 .external_manifest_store
560 .get_manifest_location(base_path.as_ref(), version)
561 .await;
562
563 let location = match location_res {
564 Ok(p) => p,
565 Err(Error::NotFound { .. }) => {
567 let path = default_resolve_version(base_path, version, object_store)
568 .await
569 .map_err(|_| Error::not_found(format!("{}@{}", base_path, version)))?
570 .path;
571 match object_store.head(&path).await {
572 Ok(ObjectMeta { size, e_tag, .. }) => {
573 let res = self
574 .external_manifest_store
575 .put_if_not_exists(
576 base_path.as_ref(),
577 version,
578 path.as_ref(),
579 size,
580 e_tag.clone(),
581 )
582 .await;
583 if let Err(e) = res {
584 warn!(
585 "could not update external manifest store during load, with error: {}",
586 e
587 );
588 }
589 let naming_scheme =
590 ManifestNamingScheme::detect_scheme_staging(path.filename().unwrap());
591 return Ok(ManifestLocation {
592 version,
593 path,
594 size: Some(size),
595 naming_scheme,
596 e_tag,
597 });
598 }
599 Err(ObjectStoreError::NotFound { .. }) => {
600 return Err(Error::not_found(path.to_string()));
601 }
602 Err(e) => return Err(e.into()),
603 }
604 }
605 Err(e) => return Err(e),
606 };
607
608 if location.path.extension() == Some(MANIFEST_EXTENSION) {
609 return self
610 .verify_finalized_manifest_location(base_path, location, object_store)
611 .await;
612 }
613
614 let naming_scheme =
615 ManifestNamingScheme::detect_scheme_staging(location.path.filename().unwrap());
616
617 let size = if let Some(size) = location.size {
618 size
619 } else {
620 let meta = object_store.head(&location.path).await?;
621 meta.size
622 };
623
624 self.finalize_manifest(
625 base_path,
626 &location.path,
627 version,
628 size,
629 object_store,
630 naming_scheme,
631 )
632 .await
633 }
634
635 async fn version_exists(
636 &self,
637 base_path: &Path,
638 version: u64,
639 object_store: &dyn OSObjectStore,
640 naming_scheme: ManifestNamingScheme,
641 ) -> Result<bool> {
642 match self
643 .external_manifest_store
644 .get_manifest_location(base_path.as_ref(), version)
645 .await
646 {
647 Ok(_) => Ok(true),
648 Err(Error::NotFound { .. }) => {
649 let path = naming_scheme.manifest_path(base_path, version);
650 match object_store.head(&path).await {
651 Ok(_) => Ok(true),
652 Err(ObjectStoreError::NotFound { .. }) => Ok(false),
653 Err(e) => Err(e.into()),
654 }
655 }
656 Err(e) => Err(e),
657 }
658 }
659
660 async fn commit(
661 &self,
662 manifest: &mut Manifest,
663 indices: Option<Vec<IndexMetadata>>,
664 base_path: &Path,
665 object_store: &ObjectStore,
666 manifest_writer: super::ManifestWriter,
667 naming_scheme: ManifestNamingScheme,
668 transaction: Option<Transaction>,
669 ) -> std::result::Result<ManifestLocation, CommitError> {
670 let path = naming_scheme.manifest_path(base_path, manifest.version);
675 let staging_path = make_staging_manifest_path(&path)?;
676 let write_res =
677 manifest_writer(object_store, manifest, indices, &staging_path, transaction).await?;
678
679 let result = self
681 .external_manifest_store
682 .put(
683 base_path,
684 manifest.version,
685 &staging_path,
686 write_res.size as u64,
687 write_res.e_tag,
688 &object_store.inner,
689 naming_scheme,
690 )
691 .await;
692
693 match result {
694 Ok(location) => {
695 write_version_hint(object_store, base_path, manifest.version).await;
696 Ok(location)
697 }
698 Err(_) => {
699 match object_store.inner.delete(&staging_path).await {
701 Ok(_) => {}
702 Err(ObjectStoreError::NotFound { .. }) => {}
703 Err(e) => return Err(CommitError::OtherError(e.into())),
704 }
705 info!(target: TRACE_FILE_AUDIT, mode=AUDIT_MODE_DELETE, r#type=AUDIT_TYPE_MANIFEST, path = staging_path.as_ref());
706 Err(CommitError::CommitConflict {})
707 }
708 }
709 }
710
711 async fn delete(&self, base_path: &Path) -> Result<()> {
712 self.external_manifest_store
713 .delete(base_path.as_ref())
714 .await
715 }
716}