| 1 | use super::*; |
| 2 | use std::{future::Future, pin::Pin}; |
| 3 | |
| 4 | pub(super) enum Outcome { |
| 5 | Read(Vec<u8>), |
| 6 | Stamp(Box<Stamp>), |
| 7 | Published, |
| 8 | } |
| 9 | |
| 10 | pub(super) enum Pending { |
| 11 | Running(Pin<Box<dyn Future<Output = std::result::Result<Outcome, CommitError>>>>), |
| 12 | Ready(std::result::Result<Outcome, CommitError>), |
| 13 | } |
| 14 | |
| 15 | fn uncommitted(error: io::Error) -> CommitError { |
| 16 | CommitError { |
| 17 | state: CommitState::NotCommitted, |
| 18 | error, |
| 19 | } |
| 20 | } |
| 21 | |
| 22 | impl HostedRemote { |
| 23 | fn operation( |
| 24 | &mut self, |
| 25 | work: impl Future<Output = std::result::Result<Outcome, CommitError>> + 'static, |
| 26 | ) -> std::result::Result<Outcome, CommitError> { |
| 27 | match self.pending.take() { |
| 28 | Some(Pending::Ready(result)) => result, |
| 29 | pending => { |
| 30 | self.pending = Some(pending.unwrap_or_else(|| Pending::Running(Box::pin(work)))); |
| 31 | Err(uncommitted(io::ErrorKind::WouldBlock.into())) |
| 32 | } |
| 33 | } |
| 34 | } |
| 35 | |
| 36 | fn publication( |
| 37 | &mut self, |
| 38 | result: std::result::Result<Outcome, CommitError>, |
| 39 | ) -> std::result::Result<(), CommitError> { |
| 40 | match result { |
| 41 | Ok(Outcome::Published) => { |
| 42 | self.seen = None; |
| 43 | Ok(()) |
| 44 | } |
| 45 | Ok(_) => Err(uncommitted(io::ErrorKind::InvalidInput.into())), |
| 46 | Err(error) => { |
| 47 | if error.error.kind() != io::ErrorKind::WouldBlock |
| 48 | && error.state == CommitState::NotCommitted |
| 49 | { |
| 50 | self.rejected = true; |
| 51 | self.guest.inner.current.lock().unwrap().remove(&self.path); |
| 52 | } |
| 53 | Err(error) |
| 54 | } |
| 55 | } |
| 56 | } |
| 57 | } |
| 58 | |
| 59 | impl crate::Remote for HostedRemote { |
| 60 | fn pending(&mut self) -> Option<Pin<Box<dyn Future<Output = ()> + '_>>> { |
| 61 | let pending = self.pending.as_mut()?; |
| 62 | Some(Box::pin(async move { |
| 63 | if let Pending::Running(work) = pending { |
| 64 | *pending = Pending::Ready(work.await); |
| 65 | } |
| 66 | })) |
| 67 | } |
| 68 | |
| 69 | fn read(&mut self) -> io::Result<Vec<u8>> { |
| 70 | let (guest, path, seen) = (self.guest.clone(), self.path.clone(), self.seen.clone()); |
| 71 | let result = self |
| 72 | .operation(async move { |
| 73 | let image = match seen.as_ref().and_then(|stamp| guest.held(&path, stamp)) { |
| 74 | Some(image) => image, |
| 75 | None => guest |
| 76 | .read_async(kind::READ, &path, LIMIT) |
| 77 | .await |
| 78 | .map_err(uncommitted)?, |
| 79 | }; |
| 80 | Ok(Outcome::Read(image)) |
| 81 | }) |
| 82 | .map_err(|error| error.error)?; |
| 83 | match result { |
| 84 | Outcome::Read(image) => { |
| 85 | self.rejected = false; |
| 86 | Ok(image) |
| 87 | } |
| 88 | _ => Err(io::ErrorKind::InvalidInput.into()), |
| 89 | } |
| 90 | } |
| 91 | |
| 92 | fn stamp(&mut self) -> io::Result<Stamp> { |
| 93 | let (guest, path) = (self.guest.clone(), self.path.clone()); |
| 94 | let result = self |
| 95 | .operation(async move { |
| 96 | guest |
| 97 | .stamp_async(&path) |
| 98 | .await |
| 99 | .map(|stamp| Outcome::Stamp(Box::new(stamp))) |
| 100 | .map_err(uncommitted) |
| 101 | }) |
| 102 | .map_err(|error| error.error)?; |
| 103 | match result { |
| 104 | Outcome::Stamp(stamp) => { |
| 105 | self.seen = Some((*stamp).clone()); |
| 106 | Ok(*stamp) |
| 107 | } |
| 108 | _ => Err(io::ErrorKind::InvalidInput.into()), |
| 109 | } |
| 110 | } |
| 111 | |
| 112 | fn accepts_edits(&self) -> bool { |
| 113 | !self.rejected |
| 114 | && self |
| 115 | .guest |
| 116 | .host() |
| 117 | .is_some_and(|host| host.ops == Some(1) && host.kinds.contains(&kind::EDITS)) |
| 118 | } |
| 119 | |
| 120 | fn publish_edits( |
| 121 | &mut self, |
| 122 | transaction: &Transaction, |
| 123 | edits: &[crate::PendingEdit], |
| 124 | revisions: &BTreeMap<onestore::ExGuid, onestore::ExGuid>, |
| 125 | ) -> std::result::Result<(), CommitError> { |
| 126 | let (guest, path, transaction, edits, revisions) = ( |
| 127 | self.guest.clone(), |
| 128 | self.path.clone(), |
| 129 | transaction.clone(), |
| 130 | edits.to_vec(), |
| 131 | revisions.clone(), |
| 132 | ); |
| 133 | let result = self.operation(async move { |
| 134 | guest |
| 135 | .edits_async(&path, &transaction, &edits, &revisions) |
| 136 | .await?; |
| 137 | Ok(Outcome::Published) |
| 138 | }); |
| 139 | self.publication(result) |
| 140 | } |
| 141 | |
| 142 | fn publish(&mut self, transaction: &Transaction) -> std::result::Result<(), CommitError> { |
| 143 | let (guest, path, transaction) = |
| 144 | (self.guest.clone(), self.path.clone(), transaction.clone()); |
| 145 | let result = self.operation(async move { |
| 146 | guest.commit_async(&path, &transaction).await?; |
| 147 | guest.inner.published(&path, &transaction); |
| 148 | Ok(Outcome::Published) |
| 149 | }); |
| 150 | self.publication(result) |
| 151 | } |
| 152 | |
| 153 | fn confirm(&mut self, base: &Stamp) -> std::result::Result<(), CommitError> { |
| 154 | let (guest, path, base) = (self.guest.clone(), self.path.clone(), base.clone()); |
| 155 | let result = self.operation(async move { |
| 156 | guest.confirm_async(&path, &base).await?; |
| 157 | Ok(Outcome::Published) |
| 158 | }); |
| 159 | self.publication(result) |
| 160 | } |
| 161 | } |
| 162 | |
| 163 | struct Cached<'a>(&'a Hosted, Listings); |
| 164 | |
| 165 | impl Guest { |
| 166 | pub(super) async fn cache_entries( |
| 167 | &self, |
| 168 | folder: &str, |
| 169 | entries: &[discover::Entry], |
| 170 | ) -> io::Result<()> { |
| 171 | let Some(listed) = self.inner.catalog.lock().unwrap().clone() else { |
| 172 | return Ok(()); |
| 173 | }; |
| 174 | let files = listed.with_extension("files"); |
| 175 | let mut listings: Listings = crate::fs::read(&listed) |
| 176 | .ok() |
| 177 | .and_then(|bytes| serde_json::from_slice(&bytes).ok()) |
| 178 | .unwrap_or_default(); |
| 179 | let before = listings.get(folder).cloned().unwrap_or_default(); |
| 180 | let mut current = Vec::with_capacity(entries.len()); |
| 181 | crate::fs::create_dir_all(files.join(folder))?; |
| 182 | for entry in entries { |
| 183 | if entry.name.contains(['/', '\\', '\0']) || matches!(entry.name.as_str(), "." | "..") { |
| 184 | return Err(io::ErrorKind::InvalidData.into()); |
| 185 | } |
| 186 | let e = wire_entry(entry); |
| 187 | let listed = (e.name, e.kind, e.size, e.modified); |
| 188 | let path = if folder.is_empty() { |
| 189 | entry.name.clone() |
| 190 | } else { |
| 191 | format!("{folder}/{}", entry.name) |
| 192 | }; |
| 193 | if entry.kind == discover::EntryKind::File |
| 194 | && (!before.contains(&listed) || crate::fs::metadata(files.join(&path)).is_err()) |
| 195 | { |
| 196 | let kind = if path.ends_with(".one") || path.ends_with(".onetoc2") { |
| 197 | kind::READ |
| 198 | } else { |
| 199 | kind::READ_FILE |
| 200 | }; |
| 201 | let bytes = self.read_async(kind, &path, LIMIT).await?; |
| 202 | crate::fs::write(files.join(path), bytes)?; |
| 203 | } |
| 204 | current.push(listed); |
| 205 | } |
| 206 | // Other folders may have finished listing while these files were being read. |
| 207 | if let Ok(bytes) = crate::fs::read(&listed) { |
| 208 | listings = serde_json::from_slice(&bytes).map_err(io::Error::other)?; |
| 209 | } |
| 210 | if listings.get(folder) != Some(&current) { |
| 211 | listings.insert(folder.to_owned(), current); |
| 212 | crate::fs::write(&listed, serde_json::to_vec(&listings)?)?; |
| 213 | crate::fs::durable().await?; |
| 214 | if let Some(reports) = &*self.inner.watch.lock().unwrap() { |
| 215 | reports.catalog(folder); |
| 216 | } |
| 217 | (self.inner.events)(); |
| 218 | } |
| 219 | Ok(()) |
| 220 | } |
| 221 | } |
| 222 | |
| 223 | impl discover::Source for Cached<'_> { |
| 224 | fn entries(&mut self, path: &str, limit: usize) -> io::Result<Vec<discover::Entry>> { |
| 225 | let entries = self.1.get(path).ok_or(io::ErrorKind::NotConnected)?; |
| 226 | if entries.len() > limit { |
| 227 | return Err(io::ErrorKind::FileTooLarge.into()); |
| 228 | } |
| 229 | Ok(entries |
| 230 | .iter() |
| 231 | .map(|(name, kind, size, modified)| { |
| 232 | entry(&WireEntry { |
| 233 | name: name.clone(), |
| 234 | kind: *kind, |
| 235 | size: *size, |
| 236 | modified: *modified, |
| 237 | }) |
| 238 | }) |
| 239 | .collect()) |
| 240 | } |
| 241 | fn read(&mut self, path: &str, limit: usize) -> io::Result<Vec<u8>> { |
| 242 | let bytes = self |
| 243 | .0 |
| 244 | .guest |
| 245 | .inner |
| 246 | .image(path) |
| 247 | .map(|image| image.to_vec()) |
| 248 | .map(Ok) |
| 249 | .unwrap_or_else(|| crate::fs::read(self.0.files().join(path)))?; |
| 250 | if bytes.len() > limit { |
| 251 | return Err(io::ErrorKind::FileTooLarge.into()); |
| 252 | } |
| 253 | Ok(bytes) |
| 254 | } |
| 255 | fn read_asset(&mut self, path: &str, limit: usize) -> io::Result<Vec<u8>> { |
| 256 | self.read(path, limit) |
| 257 | } |
| 258 | } |
| 259 | |
| 260 | impl Hosted { |
| 261 | fn files(&self) -> PathBuf { |
| 262 | self.listed.with_extension("files") |
| 263 | } |
| 264 | |
| 265 | /// Loads a bounded catalog and its files before the synchronous editor opens it. |
| 266 | pub async fn prepare(guest: Arc<Guest>, cache: &std::path::Path) -> io::Result<()> { |
| 267 | let listed = |
| 268 | crate::session::listing(cache, &guest.location()).with_extension("entries.json"); |
| 269 | let hosted = Self::new(guest, listed); |
| 270 | let deadline = Instant::now() + Duration::from_secs(20); |
| 271 | let (_, waiting) = crate::task::channel(); |
| 272 | while hosted.guest.host().is_none() { |
| 273 | if Instant::now() >= deadline { |
| 274 | return Err(hosted.guest.offline()); |
| 275 | } |
| 276 | crate::task::wait(&waiting, Some(Duration::from_millis(50))).await; |
| 277 | } |
| 278 | let mut folders = vec![String::new()]; |
| 279 | let mut count = 0; |
| 280 | while let Some(folder) = folders.pop() { |
| 281 | if folder.split('/').count() > 32 { |
| 282 | return Err(io::ErrorKind::FileTooLarge.into()); |
| 283 | } |
| 284 | let entries = hosted.guest.entries_async(&folder).await?; |
| 285 | count += entries.len(); |
| 286 | if count > 10000 { |
| 287 | return Err(io::ErrorKind::FileTooLarge.into()); |
| 288 | } |
| 289 | crate::fs::create_dir_all(hosted.files().join(&folder))?; |
| 290 | for entry in &entries { |
| 291 | if entry.kind == discover::EntryKind::Directory { |
| 292 | folders.push(if folder.is_empty() { |
| 293 | entry.name.clone() |
| 294 | } else { |
| 295 | format!("{folder}/{}", entry.name) |
| 296 | }); |
| 297 | } |
| 298 | } |
| 299 | } |
| 300 | crate::fs::durable().await |
| 301 | } |
| 302 | } |
| 303 | |
| 304 | fn desktop() -> Error { |
| 305 | io::Error::new( |
| 306 | io::ErrorKind::Unsupported, |
| 307 | "Create and organize sections in desktop Snowbound.", |
| 308 | ) |
| 309 | .into() |
| 310 | } |
| 311 | |
| 312 | impl Storage for Hosted { |
| 313 | fn discover( |
| 314 | &self, |
| 315 | cache: &mut discover::Cache, |
| 316 | limits: discover::Limits, |
| 317 | ) -> Result<discover::Folder> { |
| 318 | let listings = |
| 319 | serde_json::from_slice(&crate::fs::read(&self.listed)?).map_err(io::Error::other)?; |
| 320 | Ok(cache.discover(&mut Cached(self, listings), limits)?) |
| 321 | } |
| 322 | fn location(&self) -> String { |
| 323 | self.guest.location() |
| 324 | } |
| 325 | fn entries(&self, path: &str) -> io::Result<Vec<discover::Entry>> { |
| 326 | let listings = |
| 327 | serde_json::from_slice(&crate::fs::read(&self.listed)?).map_err(io::Error::other)?; |
| 328 | discover::Source::entries(&mut Cached(self, listings), path, 10000) |
| 329 | } |
| 330 | fn exists(&self, path: &str) -> bool { |
| 331 | crate::fs::metadata(self.files().join(path)).is_ok() |
| 332 | } |
| 333 | fn stamp(&self, path: &str) -> io::Result<Stamp> { |
| 334 | Stamp::of(&self.read(path).map_err(io::Error::other)?).map_err(io::Error::other) |
| 335 | } |
| 336 | fn read(&self, path: &str) -> Result<Vec<u8>> { |
| 337 | self.read_file(path, LIMIT) |
| 338 | } |
| 339 | fn read_file(&self, path: &str, limit: usize) -> Result<Vec<u8>> { |
| 340 | Ok(discover::Source::read( |
| 341 | &mut Cached(self, Listings::new()), |
| 342 | path, |
| 343 | limit, |
| 344 | )?) |
| 345 | } |
| 346 | fn create(&self, _: &str, _: &[u8]) -> Result<()> { |
| 347 | Err(desktop()) |
| 348 | } |
| 349 | fn create_directory(&self, _: &str) -> Result<()> { |
| 350 | Err(desktop()) |
| 351 | } |
| 352 | fn hide(&self, _: &str) -> Result<()> { |
| 353 | Err(desktop()) |
| 354 | } |
| 355 | fn rename(&self, _: &str, _: &str) -> Result<()> { |
| 356 | Err(desktop()) |
| 357 | } |
| 358 | fn rename_root(&self, _: &str, _: &[String]) -> Result<String> { |
| 359 | Err(desktop()) |
| 360 | } |
| 361 | fn replace(&self, _: &str, _: &str) -> Result<()> { |
| 362 | Err(desktop()) |
| 363 | } |
| 364 | fn delete(&self, _: &str) -> Result<()> { |
| 365 | Err(desktop()) |
| 366 | } |
| 367 | fn place(&self, _: &str, _: [u8; 16], _: &str) -> Result<()> { |
| 368 | Err(desktop()) |
| 369 | } |
| 370 | fn commit(&self, _: &str, _: &Transaction) -> Result<()> { |
| 371 | Err(desktop()) |
| 372 | } |
| 373 | fn confirm(&self, _: &str, _: &Stamp) -> std::result::Result<(), CommitError> { |
| 374 | Err(uncommitted(io::Error::new( |
| 375 | io::ErrorKind::Unsupported, |
| 376 | "Section management requires desktop Snowbound", |
| 377 | ))) |
| 378 | } |
| 379 | fn supersede(&self, _: &str, _: &Stamp, _: &str) -> Result<()> { |
| 380 | Err(desktop()) |
| 381 | } |
| 382 | } |