Implement e2e, add container with server
This commit is contained in:
@@ -8,11 +8,12 @@ use std::{
|
||||
|
||||
use fuser::INodeNo;
|
||||
use sea_orm::{DatabaseConnection, EntityTrait};
|
||||
use tokio::sync::Notify;
|
||||
use tokio_stream::StreamExt;
|
||||
|
||||
use crate::item::Item;
|
||||
use crate::origins::FileWatcher;
|
||||
use crate::origins::network::transport::NetworkTransport;
|
||||
use crate::origins::{FileWatcher, WatcherHandle};
|
||||
use crate::proto::ManifestEntry as ProtoManifestEntry;
|
||||
use tracing::{info, warn};
|
||||
|
||||
@@ -51,8 +52,9 @@ impl NetworkOriginFileWatcher {
|
||||
/// The trait signature still requires it for parity with LocalOriginFileWatcher;
|
||||
/// a future refactor can rebuild the map in place here.
|
||||
impl FileWatcher for NetworkOriginFileWatcher {
|
||||
fn watch(&self, files: Arc<std::sync::Mutex<BTreeMap<INodeNo, Item>>>) {
|
||||
fn watch(&self, files: Arc<std::sync::Mutex<BTreeMap<INodeNo, Item>>>) -> WatcherHandle {
|
||||
let runtime_handle = self.runtime_handle.clone();
|
||||
let shutdown = Arc::new(Notify::new());
|
||||
let state = WatcherState {
|
||||
transport: self.transport.clone(),
|
||||
client: self.client.clone(),
|
||||
@@ -60,8 +62,13 @@ impl FileWatcher for NetworkOriginFileWatcher {
|
||||
latest_manifest: self.latest_manifest.clone(),
|
||||
files,
|
||||
};
|
||||
thread::spawn(move || {
|
||||
runtime_handle.block_on(state.run_loop());
|
||||
let loop_shutdown = shutdown.clone();
|
||||
let join = thread::spawn(move || {
|
||||
runtime_handle.block_on(state.run_loop(loop_shutdown));
|
||||
});
|
||||
return WatcherHandle::new(move || {
|
||||
shutdown.notify_one();
|
||||
let _ = join.join();
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -76,22 +83,33 @@ struct WatcherState {
|
||||
}
|
||||
|
||||
impl WatcherState {
|
||||
async fn run_loop(self) {
|
||||
async fn run_loop(self, shutdown: Arc<Notify>) {
|
||||
loop {
|
||||
match self.transport.subscribe_events().await {
|
||||
let subscribed = tokio::select! {
|
||||
biased;
|
||||
_ = shutdown.notified() => return,
|
||||
s = self.transport.subscribe_events() => s,
|
||||
};
|
||||
match subscribed {
|
||||
Ok(mut stream) => {
|
||||
info!("network watcher: subscribed to /events");
|
||||
while let Some(item) = stream.next().await {
|
||||
loop {
|
||||
let item = tokio::select! {
|
||||
biased;
|
||||
_ = shutdown.notified() => return,
|
||||
item = stream.next() => item,
|
||||
};
|
||||
match item {
|
||||
Ok(_event) => {
|
||||
Some(Ok(_event)) => {
|
||||
if let Err(e) = self.reconcile_once().await {
|
||||
warn!(error = %e, "network watcher: reconcile after event failed");
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
Some(Err(e)) => {
|
||||
warn!(error = %e, "network watcher: stream error; reconnecting");
|
||||
break;
|
||||
}
|
||||
None => break,
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -99,7 +117,11 @@ impl WatcherState {
|
||||
warn!(error = %e, "network watcher: subscribe failed; will retry");
|
||||
}
|
||||
}
|
||||
tokio::time::sleep(POLL_FALLBACK_INTERVAL).await;
|
||||
tokio::select! {
|
||||
biased;
|
||||
_ = shutdown.notified() => return,
|
||||
_ = tokio::time::sleep(POLL_FALLBACK_INTERVAL) => {}
|
||||
}
|
||||
if let Err(e) = self.reconcile_once().await {
|
||||
warn!(error = %e, "network watcher: poll reconcile failed");
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user