| 502 | 502 | | |
| 503 | 503 | | async fn open(&self, key: &str) -> Result<ArtifactsRepo> { |
| 504 | 504 | | let (namespace, name) = locate(key); |
| 505 | | − | let handle = invoke(&namespace, key, self.binding(&namespace)?, "get", &[name.as_str().into()], true).await?; |
| 506 | 505 | | Ok(ArtifactsRepo { |
| 507 | | − | handle, |
| 506 | + | handle: RefCell::new(None), |
| 507 | + | binding: self.binding(&namespace)?.clone(), |
| 508 | 508 | | key: key.to_owned(), |
| 509 | 509 | | name, |
| 510 | 510 | | namespace, |
| ⋯ |
| 525 | 525 | | /// A handle to one Artifacts repository. It is an RPC stub, so it is |
| 526 | 526 | | /// released when dropped. |
| 527 | 527 | | pub struct ArtifactsRepo { |
| 528 | | − | handle: JsValue, |
| 528 | + | /// The store's handle, asked for (`get`) on the first call that needs |
| 529 | + | /// the store: an answer from a cache, or a fetch over git with a kept |
| 530 | + | /// credential, never costs a `get`. |
| 531 | + | handle: RefCell<Option<JsValue>>, |
| 532 | + | /// The namespace's binding, for that `get`. |
| 533 | + | binding: JsValue, |
| 529 | 534 | | /// The repository's store key, which scopes its cached objects. |
| 530 | 535 | | key: String, |
| 531 | 536 | | /// Its name in its namespace. |
| ⋯ |
| 584 | 589 | | version.map(|version| CacheKey::Versioned(format!("branches/{version}"))) |
| 585 | 590 | | } |
| 586 | 591 | | |
| 592 | + | /// Objects named by their content, kept in the isolate ahead of the Cache |
| 593 | + | /// API, oldest out first past `MEMORY_CACHE_BYTES`. They never go stale, |
| 594 | + | /// and each is under its repository's key. |
| 595 | + | const MEMORY_CACHE_BYTES: usize = 16 * 1024 * 1024; |
| 596 | + | |
| 597 | + | #[derive(Default)] |
| 598 | + | struct MemoryCache { |
| 599 | + | entries: HashMap<String, Rc<Vec<u8>>>, |
| 600 | + | order: std::collections::VecDeque<String>, |
| 601 | + | bytes: usize, |
| 602 | + | } |
| 603 | + | |
| 604 | + | impl MemoryCache { |
| 605 | + | fn get(&self, url: &str) -> Option<Vec<u8>> { |
| 606 | + | self.entries.get(url).map(|bytes| bytes.as_ref().clone()) |
| 607 | + | } |
| 608 | + | |
| 609 | + | fn put(&mut self, url: String, bytes: &[u8]) { |
| 610 | + | if bytes.len() > MEMORY_CACHE_BYTES / 16 || self.entries.contains_key(&url) { |
| 611 | + | return; |
| 612 | + | } |
| 613 | + | self.bytes += bytes.len(); |
| 614 | + | self.entries.insert(url.clone(), Rc::new(bytes.to_vec())); |
| 615 | + | self.order.push_back(url); |
| 616 | + | while self.bytes > MEMORY_CACHE_BYTES { |
| 617 | + | let Some(oldest) = self.order.pop_front() else { break }; |
| 618 | + | if let Some(gone) = self.entries.remove(&oldest) { |
| 619 | + | self.bytes -= gone.len(); |
| 620 | + | } |
| 621 | + | } |
| 622 | + | } |
| 623 | + | } |
| 624 | + | |
| 625 | + | thread_local! { |
| 626 | + | static MEMORY: RefCell<MemoryCache> = RefCell::new(MemoryCache::default()); |
| 627 | + | } |
| 628 | + | |
| 587 | 629 | | impl ArtifactsRepo { |
| 588 | 630 | | fn cache_url(&self, path: &str) -> String { |
| 589 | 631 | | format!("{OBJECT_CACHE}{}/{path}", self.key) |
| 590 | 632 | | } |
| 591 | 633 | | |
| 592 | | − | async fn cached_at(&self, path: &str) -> Option<Vec<u8>> { |
| 593 | | − | let mut response = worker::Cache::default().get(self.cache_url(path), false).await.ok()??; |
| 594 | | − | response.bytes().await.ok() |
| 634 | + | /// A kept answer: from the isolate for one kept for good, else the |
| 635 | + | /// Cache API. Each look is metered (`cache.memory_hit`, `cache.edge_hit`, |
| 636 | + | /// `cache.miss`), so the usage check shows which cache answers. |
| 637 | + | async fn cached_at(&self, path: &str, forever: bool) -> Option<Vec<u8>> { |
| 638 | + | let url = self.cache_url(path); |
| 639 | + | if forever && let Some(bytes) = MEMORY.with(|memory| memory.borrow().get(&url)) { |
| 640 | + | meters::record_bytes("cache.memory_hit", &self.key, 0, bytes.len() as u64); |
| 641 | + | return Some(bytes); |
| 642 | + | } |
| 643 | + | let found = match worker::Cache::default().get(url.clone(), false).await { |
| 644 | + | Ok(Some(mut response)) => response.bytes().await.ok(), |
| 645 | + | _ => None, |
| 646 | + | }; |
| 647 | + | match &found { |
| 648 | + | Some(bytes) => { |
| 649 | + | meters::record_bytes("cache.edge_hit", &self.key, 0, bytes.len() as u64); |
| 650 | + | if forever { |
| 651 | + | MEMORY.with(|memory| memory.borrow_mut().put(url, bytes)); |
| 652 | + | } |
| 653 | + | } |
| 654 | + | None => meters::record("cache.miss", &self.key, 0, 0), |
| 655 | + | } |
| 656 | + | found |
| 595 | 657 | | } |
| 596 | 658 | | |
| 597 | 659 | | async fn cached(&self, kind: &str, hash: &str) -> Option<Vec<u8>> { |
| 598 | | − | self.cached_at(&format!("{kind}/{hash}")).await |
| 660 | + | self.cached_at(&format!("{kind}/{hash}"), true).await |
| 599 | 661 | | } |
| 600 | 662 | | |
| 601 | 663 | | async fn keep_at(&self, path: &str, bytes: Vec<u8>, max_age: &str) { |
| 664 | + | if max_age == OBJECT_MAX_AGE { |
| 665 | + | MEMORY.with(|memory| memory.borrow_mut().put(self.cache_url(path), &bytes)); |
| 666 | + | } |
| 602 | 667 | | let Ok(mut response) = worker::Response::from_bytes(bytes) else { |
| 603 | 668 | | return; |
| 604 | 669 | | }; |
| ⋯ |
| 613 | 678 | | |
| 614 | 679 | | async fn get_key(&self, key: &CacheKey) -> Option<Vec<u8>> { |
| 615 | 680 | | match key { |
| 616 | | − | CacheKey::Forever(path) | CacheKey::Versioned(path) => self.cached_at(path).await, |
| 681 | + | CacheKey::Forever(path) => self.cached_at(path, true).await, |
| 682 | + | CacheKey::Versioned(path) => self.cached_at(path, false).await, |
| 617 | 683 | | } |
| 618 | 684 | | } |
| 619 | 685 | | |
| ⋯ |
| 625 | 691 | | } |
| 626 | 692 | | |
| 627 | 693 | | async fn call(&self, method: &str, args: &[JsValue], retry: bool) -> std::result::Result<JsValue, StoreError> { |
| 628 | | − | invoke(&self.namespace, &self.key, &self.handle, method, args, retry).await |
| 694 | + | let handle = self.handle().await?; |
| 695 | + | invoke(&self.namespace, &self.key, &handle, method, args, retry).await |
| 629 | 696 | | } |
| 630 | 697 | | |
| 698 | + | /// The store's handle, asked for the first time it is needed. |
| 699 | + | async fn handle(&self) -> std::result::Result<JsValue, StoreError> { |
| 700 | + | if let Some(handle) = self.handle.borrow().as_ref() { |
| 701 | + | return Ok(handle.clone()); |
| 702 | + | } |
| 703 | + | let handle = invoke(&self.namespace, &self.key, &self.binding, "get", &[self.name.as_str().into()], true).await?; |
| 704 | + | *self.handle.borrow_mut() = Some(handle.clone()); |
| 705 | + | Ok(handle) |
| 706 | + | } |
| 707 | + | |
| 631 | 708 | | /// A new credential from the store. Its remote is worked out from the |
| 632 | 709 | | /// key once this isolate knows where the namespace's remotes start; |
| 633 | 710 | | /// until then the store is asked (`info()`) alongside the token. |
| ⋯ |
| 638 | 715 | | let token: RawToken = js::from_js(&self.call("createToken", &args, true).await?)?; |
| 639 | 716 | | return Ok(GitAccess { remote: remote_from(&prefix, &self.name), token: token.plaintext }); |
| 640 | 717 | | } |
| 718 | + | self.handle().await?; |
| 641 | 719 | | let (info, token) = futures_util::future::join(self.call("info", &[], true), self.call("createToken", &args, true)).await; |
| 642 | 720 | | let info: RawInfo = js::from_js(&info?)?; |
| 643 | 721 | | let token: RawToken = js::from_js(&token?)?; |
| ⋯ |
| 655 | 733 | | fn drop(&mut self) { |
| 656 | 734 | | let symbol = js::get(&worker::js_sys::global(), "Symbol"); |
| 657 | 735 | | let dispose = js::get(&symbol, "dispose"); |
| 658 | | − | if let Ok(function) = Reflect::get(&self.handle, &dispose) |
| 659 | | − | .and_then(|value| value.dyn_into::<worker::js_sys::Function>()) |
| 660 | | − | { |
| 661 | | − | let _ = function.call0(&self.handle); |
| 736 | + | let Some(handle) = self.handle.borrow_mut().take() else { return }; |
| 737 | + | if let Ok(function) = Reflect::get(&handle, &dispose).and_then(|value| value.dyn_into::<worker::js_sys::Function>()) { |
| 738 | + | let _ = function.call0(&handle); |
| 662 | 739 | | } |
| 663 | 740 | | } |
| 664 | 741 | | } |
| ⋯ |
| 981 | 1058 | | } |
| 982 | 1059 | | |
| 983 | 1060 | | #[test] |
| 1061 | + | fn the_memory_cache_drops_its_oldest_past_its_budget() { |
| 1062 | + | let mut cache = MemoryCache::default(); |
| 1063 | + | let chunk = vec![7u8; MEMORY_CACHE_BYTES / 16]; |
| 1064 | + | for n in 0..17 { |
| 1065 | + | cache.put(format!("k{n}"), &chunk); |
| 1066 | + | } |
| 1067 | + | // Sixteen chunks fit; the seventeenth pushed the first out. |
| 1068 | + | assert!(cache.get("k0").is_none()); |
| 1069 | + | assert_eq!(cache.get("k16").map(|b| b.len()), Some(chunk.len())); |
| 1070 | + | assert!(cache.bytes <= MEMORY_CACHE_BYTES); |
| 1071 | + | // Too large to keep at all. |
| 1072 | + | cache.put("big".into(), &vec![0u8; MEMORY_CACHE_BYTES / 16 + 1]); |
| 1073 | + | assert!(cache.get("big").is_none()); |
| 1074 | + | } |
| 1075 | + | |
| 1076 | + | #[test] |
| 984 | 1077 | | fn binding_methods_have_snake_case_meters() { |
| 985 | 1078 | | assert_eq!(meter_of("createToken"), "binding.create_token"); |
| 986 | 1079 | | assert_eq!(meter_of("readBlob"), "binding.read_blob"); |