Complete local-first content migration slice
This commit is contained in:
@@ -85,35 +85,6 @@ pub fn initialize(
|
||||
}
|
||||
}
|
||||
|
||||
// Remove sqlite database as part of Logout v0.
|
||||
// TODO: Implement per user scoping of sqlite.
|
||||
#[cfg_attr(not(feature = "local_fs"), allow(unused_variables))]
|
||||
pub fn remove(sender: &Option<SyncSender<ModelEvent>>) {
|
||||
cfg_if::cfg_if! {
|
||||
if #[cfg(feature = "local_fs")] {
|
||||
if let Some(sender) = sender.clone() {
|
||||
sqlite::remove(sender);
|
||||
}
|
||||
} else {
|
||||
log::info!("Local filesystem persistence is not enabled.");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Reconstruct sqlite database as part of Logout v0.
|
||||
#[cfg_attr(not(feature = "local_fs"), allow(unused_variables))]
|
||||
pub fn reconstruct(sender: &Option<SyncSender<ModelEvent>>) {
|
||||
cfg_if::cfg_if! {
|
||||
if #[cfg(feature = "local_fs")] {
|
||||
if let Some(sender) = sender.clone() {
|
||||
sqlite::reconstruct(sender);
|
||||
}
|
||||
} else {
|
||||
log::info!("Local filesystem persistence is not enabled.");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Holds interfaces to the writer thread.
|
||||
pub struct WriterHandles {
|
||||
pub handle: JoinHandle<()>,
|
||||
@@ -175,12 +146,10 @@ impl Entity for PersistenceWriter {
|
||||
|
||||
impl SingletonEntity for PersistenceWriter {}
|
||||
|
||||
/// TODO: all of this data should eventually be indexed by user_id so that
|
||||
/// the logged in user sees the data for their user (and if another user logs in,
|
||||
/// they see their respective data). To do this, we can simply return a mapping
|
||||
/// of user ID->SqliteData and get the respective AppState after the user logs in.
|
||||
/// Data restored from Galaxy's local application database.
|
||||
///
|
||||
/// For now, to address the global scoping here, we clear all persisted data on logout.
|
||||
/// This data belongs to the local installation rather than an inherited Warp
|
||||
/// account, so logging out of a legacy account must not clear it.
|
||||
pub struct PersistedData {
|
||||
/// Session restoration data
|
||||
pub app_state: AppState,
|
||||
@@ -299,12 +268,6 @@ pub enum ModelEvent {
|
||||
SaveExperiments {
|
||||
experiments: Vec<ServerExperiment>,
|
||||
},
|
||||
// `PauseAndRemoveDatabase` and `ReconstructAndResume` are used to pause and resume the writer thread.
|
||||
// These are employed as part of Logout v0 to ensure that the writer thread
|
||||
// does not continue writing to the DB after the user has logged out and the DB is deleted.
|
||||
PauseAndRemoveDatabase,
|
||||
#[cfg(feature = "local_fs")]
|
||||
ReconstructAndResume,
|
||||
InsertObjectAction {
|
||||
object_action: ObjectAction,
|
||||
},
|
||||
|
||||
@@ -460,41 +460,12 @@ fn ensure_owner_only_file(_path: &Path) -> Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(super) fn remove(sender: SyncSender<ModelEvent>) {
|
||||
// Instruct the writer thread to remove the database and pause processing
|
||||
// events.
|
||||
// Ideally, we'd drop any other events in the channel, but it's not worth the complexity right
|
||||
// now. Having the writer thread remove the database file prevents race conditions if the
|
||||
// thread is in the middle of another update.
|
||||
report_if_error!(sender
|
||||
.send(ModelEvent::PauseAndRemoveDatabase)
|
||||
.context("Error requesting database deletion"));
|
||||
}
|
||||
|
||||
pub(super) fn reconstruct(sender: SyncSender<ModelEvent>) {
|
||||
report_if_error!(sender
|
||||
.send(ModelEvent::ReconstructAndResume)
|
||||
.context("Error resuming SQLite thread"));
|
||||
}
|
||||
|
||||
fn reconstruct_database(path: &Path) -> Result<SqliteConnection> {
|
||||
// If the DB still exists, logout might have failed. However, it's more likely that something
|
||||
// else wrote to it before the user logged back in.
|
||||
if std::fs::metadata(path).is_ok() {
|
||||
log::info!("Reconstructing database, but it already exists");
|
||||
}
|
||||
|
||||
// Always reinitialize DB - setup_database will only create it if it doesn't exist.
|
||||
setup_database(path)
|
||||
}
|
||||
|
||||
fn start_writer(conn: SqliteConnection, database_path: PathBuf) -> Result<WriterHandles> {
|
||||
let (tx, rx) = std::sync::mpsc::sync_channel(CHANNEL_SIZE);
|
||||
let mut current_conn = conn;
|
||||
let handle = thread::Builder::new()
|
||||
.name("SQLite Writer".into())
|
||||
.spawn(move || {
|
||||
let mut paused = false;
|
||||
loop {
|
||||
let events = match rx.recv() {
|
||||
Ok(event) => {
|
||||
@@ -515,38 +486,11 @@ fn start_writer(conn: SqliteConnection, database_path: PathBuf) -> Result<Writer
|
||||
|
||||
for event in events {
|
||||
match event {
|
||||
ModelEvent::ReconstructAndResume => {
|
||||
match reconstruct_database(&database_path) {
|
||||
Ok(conn) => {
|
||||
current_conn = conn;
|
||||
paused = false;
|
||||
log::info!("SQLite Writer is resumed");
|
||||
}
|
||||
Err(err) => {
|
||||
report_db_error("reconstruction", err, &database_path);
|
||||
}
|
||||
}
|
||||
}
|
||||
ModelEvent::PauseAndRemoveDatabase => {
|
||||
paused = true;
|
||||
log::info!("SQLite Writer is paused");
|
||||
|
||||
if let Err(err) = std::fs::remove_file(&database_path) {
|
||||
report_error!(anyhow::Error::new(err)
|
||||
.context("Error removing SQLite database"));
|
||||
} else {
|
||||
log::info!("Removed SQLite database");
|
||||
}
|
||||
}
|
||||
ModelEvent::Terminate => {
|
||||
log::info!("Shutting down SQLite writer thread");
|
||||
return;
|
||||
}
|
||||
event => {
|
||||
if paused {
|
||||
log::info!("Ignoring event as SQLite Writer is on pause");
|
||||
continue;
|
||||
}
|
||||
if let Err(err) = handle_model_event(event, &mut current_conn) {
|
||||
report_db_error("Model", err, &database_path);
|
||||
}
|
||||
@@ -560,14 +504,10 @@ fn start_writer(conn: SqliteConnection, database_path: PathBuf) -> Result<Writer
|
||||
|
||||
/// Handles a single [`ModelEvent`] by dispatching to an event-specific function.
|
||||
/// Events which affect the SQLite writer event loop _must_ instead be handled by the event loop itself:
|
||||
/// * [`ModelEvent::PauseAndRemoveDatabase`]
|
||||
/// * [`ModelEvent::ReconstructAndResume`]
|
||||
/// * [`ModelEvent::Terminate`]
|
||||
fn handle_model_event(event: ModelEvent, connection: &mut SqliteConnection) -> anyhow::Result<()> {
|
||||
match event {
|
||||
ModelEvent::PauseAndRemoveDatabase
|
||||
| ModelEvent::ReconstructAndResume
|
||||
| ModelEvent::Terminate => {
|
||||
ModelEvent::Terminate => {
|
||||
panic!("Unhandled control-flow event {event:?}");
|
||||
}
|
||||
ModelEvent::SaveBlock(BlockCompleted {
|
||||
|
||||
@@ -15,6 +15,7 @@ use super::{
|
||||
encode_path, get_all_codebase_index_metadata, read_sqlite_data, save_app_state,
|
||||
save_codebase_index_metadata, setup_database, start_writer, GALAXY_SQLITE_FILE_NAME,
|
||||
};
|
||||
use crate::ai::facts::{AIFact, AIMemory, CloudAIFact};
|
||||
use crate::app_state::{
|
||||
AppState, CodePaneSnapShot, CodePaneTabSnapshot, LeafContents, LeafSnapshot, PaneNodeSnapshot,
|
||||
TabGroupSnapshot, TabSnapshot, TerminalPaneSnapshot, WindowSnapshot,
|
||||
@@ -24,11 +25,13 @@ use crate::code::editor_management::CodeSource;
|
||||
use crate::notebooks::{CloudNotebook, CloudNotebookModel};
|
||||
use crate::persistence::model::ObjectPermissions;
|
||||
use crate::persistence::{BlockCompleted, ModelEvent, PersistenceScope};
|
||||
use crate::server::ids::ClientId;
|
||||
use crate::server::ids::{ClientId, SyncId};
|
||||
use crate::tab::SelectedTabColor;
|
||||
use crate::terminal::model::block::SerializedBlock;
|
||||
use crate::terminal::ShellLaunchData;
|
||||
use crate::themes::theme::AnsiColorIdentifier;
|
||||
use crate::workflows::workflow::Workflow;
|
||||
use crate::workflows::CloudWorkflow;
|
||||
use crate::workspace::tab_group::TabGroupId;
|
||||
|
||||
#[test]
|
||||
@@ -197,6 +200,157 @@ fn sqlite_writer_reuses_codebase_index_metadata_events() {
|
||||
let restored = get_all_codebase_index_metadata(&mut conn).expect("metadata should load");
|
||||
assert!(restored.is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn sqlite_writer_restores_and_deletes_local_rules() {
|
||||
let tempdir = tempfile::tempdir().expect("tempdir should be created");
|
||||
let database_path = tempdir.path().join("warp.sqlite");
|
||||
let conn = setup_database(&database_path).expect("database should initialize");
|
||||
let id = SyncId::ClientId(ClientId::new());
|
||||
let fact = AIFact::Memory(AIMemory {
|
||||
name: Some("Rust".to_string()),
|
||||
content: "Never unwrap".to_string(),
|
||||
is_autogenerated: false,
|
||||
suggested_logging_id: None,
|
||||
});
|
||||
let rule = crate::local_object_repository::new_local_rule(id, fact.clone());
|
||||
|
||||
let writer = start_writer(conn, database_path.clone()).expect("writer should start");
|
||||
writer
|
||||
.sender
|
||||
.send(ModelEvent::UpsertGenericStringObject {
|
||||
object: Box::new(rule),
|
||||
})
|
||||
.expect("rule upsert should send");
|
||||
writer
|
||||
.sender
|
||||
.send(ModelEvent::Terminate)
|
||||
.expect("terminate event should send");
|
||||
writer.handle.join().expect("writer should terminate");
|
||||
|
||||
let mut conn = setup_database(&database_path).expect("database should reopen");
|
||||
let restored = read_sqlite_data(&mut conn, None).expect("persisted data should load");
|
||||
let restored_rule = restored
|
||||
.cloud_objects
|
||||
.iter()
|
||||
.find_map(|object| {
|
||||
let rule: Option<&CloudAIFact> = object.into();
|
||||
rule
|
||||
})
|
||||
.expect("local rule should be restored");
|
||||
assert_eq!(restored_rule.id, id);
|
||||
assert_eq!(restored_rule.model().string_model, fact);
|
||||
assert!(!restored_rule.metadata.has_pending_content_changes());
|
||||
|
||||
let writer = start_writer(conn, database_path.clone()).expect("writer should restart");
|
||||
writer
|
||||
.sender
|
||||
.send(ModelEvent::DeleteObjects {
|
||||
ids: vec![(id, crate::cloud_object::ObjectIdType::GenericStringObject)],
|
||||
})
|
||||
.expect("rule deletion should send");
|
||||
writer
|
||||
.sender
|
||||
.send(ModelEvent::Terminate)
|
||||
.expect("terminate event should send");
|
||||
writer.handle.join().expect("writer should terminate");
|
||||
|
||||
let mut conn = setup_database(&database_path).expect("database should reopen");
|
||||
let restored = read_sqlite_data(&mut conn, None).expect("persisted data should load");
|
||||
assert!(restored.cloud_objects.iter().all(|object| {
|
||||
let rule: Option<&CloudAIFact> = object.into();
|
||||
rule.is_none()
|
||||
}));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn sqlite_writer_restores_and_deletes_local_notebooks_and_workflows() {
|
||||
let tempdir = tempfile::tempdir().expect("tempdir should be created");
|
||||
let database_path = tempdir.path().join("warp.sqlite");
|
||||
let conn = setup_database(&database_path).expect("database should initialize");
|
||||
let notebook_id = SyncId::ClientId(ClientId::new());
|
||||
let workflow_id = SyncId::ClientId(ClientId::new());
|
||||
let notebook = crate::local_object_repository::new_local_notebook(
|
||||
notebook_id,
|
||||
None,
|
||||
CloudNotebookModel {
|
||||
title: "Local notebook".to_string(),
|
||||
data: "echo local".to_string(),
|
||||
ai_document_id: None,
|
||||
conversation_id: None,
|
||||
},
|
||||
);
|
||||
let workflow = crate::local_object_repository::new_local_workflow(
|
||||
workflow_id,
|
||||
None,
|
||||
Workflow::new("Local workflow", "cargo test"),
|
||||
);
|
||||
|
||||
let writer = start_writer(conn, database_path.clone()).expect("writer should start");
|
||||
writer
|
||||
.sender
|
||||
.send(ModelEvent::UpsertNotebook { notebook })
|
||||
.expect("notebook upsert should send");
|
||||
writer
|
||||
.sender
|
||||
.send(ModelEvent::UpsertWorkflow { workflow })
|
||||
.expect("workflow upsert should send");
|
||||
writer
|
||||
.sender
|
||||
.send(ModelEvent::Terminate)
|
||||
.expect("terminate event should send");
|
||||
writer.handle.join().expect("writer should terminate");
|
||||
|
||||
let mut conn = setup_database(&database_path).expect("database should reopen");
|
||||
let restored = read_sqlite_data(&mut conn, None).expect("persisted data should load");
|
||||
let restored_notebook = restored
|
||||
.cloud_objects
|
||||
.iter()
|
||||
.find_map(|object| {
|
||||
let notebook: Option<&CloudNotebook> = object.into();
|
||||
notebook
|
||||
})
|
||||
.expect("local notebook should be restored");
|
||||
let restored_workflow = restored
|
||||
.cloud_objects
|
||||
.iter()
|
||||
.find_map(|object| {
|
||||
let workflow: Option<&CloudWorkflow> = object.into();
|
||||
workflow
|
||||
})
|
||||
.expect("local workflow should be restored");
|
||||
assert_eq!(restored_notebook.id, notebook_id);
|
||||
assert_eq!(restored_notebook.model().title, "Local notebook");
|
||||
assert!(!restored_notebook.metadata.has_pending_content_changes());
|
||||
assert_eq!(restored_workflow.id, workflow_id);
|
||||
assert_eq!(restored_workflow.model().data.name(), "Local workflow");
|
||||
assert!(!restored_workflow.metadata.has_pending_content_changes());
|
||||
|
||||
let writer = start_writer(conn, database_path.clone()).expect("writer should restart");
|
||||
writer
|
||||
.sender
|
||||
.send(ModelEvent::DeleteObjects {
|
||||
ids: vec![
|
||||
(notebook_id, crate::cloud_object::ObjectIdType::Notebook),
|
||||
(workflow_id, crate::cloud_object::ObjectIdType::Workflow),
|
||||
],
|
||||
})
|
||||
.expect("local object deletion should send");
|
||||
writer
|
||||
.sender
|
||||
.send(ModelEvent::Terminate)
|
||||
.expect("terminate event should send");
|
||||
writer.handle.join().expect("writer should terminate");
|
||||
|
||||
let mut conn = setup_database(&database_path).expect("database should reopen");
|
||||
let restored = read_sqlite_data(&mut conn, None).expect("persisted data should load");
|
||||
assert!(restored.cloud_objects.iter().all(|object| {
|
||||
let notebook: Option<&CloudNotebook> = object.into();
|
||||
let workflow: Option<&CloudWorkflow> = object.into();
|
||||
notebook.is_none() && workflow.is_none()
|
||||
}));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_deduplicate_snapshots() {
|
||||
let local_notebook = CloudNotebook::new_local(
|
||||
|
||||
Reference in New Issue
Block a user