Write a scan's findings with prepared statements, and its jobs in the same commit
`persist` runs after every scan, for every photograph the scan listed. On a settled library that is the folders whose ETag changed -- a sidecar written there by a rating is enough -- so one relisted folder of 1,600 images is an ordinary pass, and a first scan is all 24,000. Per photograph it prepared four statements from their SQL (a folder lookup, the image upsert, the id read-back, the remote upsert) and then, after the commit, found the image again by path and enqueued its thumbnail job as an autocommitting statement of its own -- a commit per photograph, for rows that were almost all already queued. Now the statements are prepared once per pass, a folder's id is looked up once per folder rather than once per photograph in it, and the job is enqueued inside the transaction with the id already in hand. That also makes the job atomic with the row it points at, which is what the old ordering after the commit was trying to guarantee. `jobs::enqueue` uses a cached statement for the same reason. persist_bench on a copy of the reference catalog, CPU, best of runs: largest folder (1,589 images) 102-118 ms -> 10-13 ms whole library (23,582 images) 1.55-2.19 s -> 188-192 ms The fingerprint of images, remote, jobs and folders after the run is the same for both builds.
This commit is contained in:
@@ -185,7 +185,8 @@ pub fn enqueue(
|
|||||||
priority: Priority,
|
priority: Priority,
|
||||||
payload: Option<&str>,
|
payload: Option<&str>,
|
||||||
) -> Result<(), CatalogError> {
|
) -> Result<(), CatalogError> {
|
||||||
conn.execute(
|
// Cached: a scan enqueues one per photograph it lists.
|
||||||
|
conn.prepare_cached(
|
||||||
"INSERT INTO jobs(kind, subject_id, priority, state, payload)
|
"INSERT INTO jobs(kind, subject_id, priority, state, payload)
|
||||||
VALUES (?1, ?2, ?3, 0, ?4)
|
VALUES (?1, ?2, ?3, 0, ?4)
|
||||||
ON CONFLICT(kind, subject_id) DO UPDATE SET
|
ON CONFLICT(kind, subject_id) DO UPDATE SET
|
||||||
@@ -195,8 +196,13 @@ pub fn enqueue(
|
|||||||
state = CASE WHEN jobs.state = 2 THEN 0 ELSE jobs.state END,
|
state = CASE WHEN jobs.state = 2 THEN 0 ELSE jobs.state END,
|
||||||
attempts = CASE WHEN jobs.state = 2 THEN 0 ELSE jobs.attempts END,
|
attempts = CASE WHEN jobs.state = 2 THEN 0 ELSE jobs.attempts END,
|
||||||
not_before = CASE WHEN jobs.state = 2 THEN 0 ELSE jobs.not_before END",
|
not_before = CASE WHEN jobs.state = 2 THEN 0 ELSE jobs.not_before END",
|
||||||
rusqlite::params![kind as i64, subject_id, priority as i64, payload],
|
)?
|
||||||
)?;
|
.execute(rusqlite::params![
|
||||||
|
kind as i64,
|
||||||
|
subject_id,
|
||||||
|
priority as i64,
|
||||||
|
payload
|
||||||
|
])?;
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -603,24 +603,28 @@ pub(super) fn persist(
|
|||||||
|
|
||||||
// Folder ETags first — without these persisted, the next scan prunes
|
// Folder ETags first — without these persisted, the next scan prunes
|
||||||
// nothing and walks the whole tree again (ARCH §6.6).
|
// nothing and walks the whole tree again (ARCH §6.6).
|
||||||
for (path, validator) in &result.directories {
|
{
|
||||||
tx.execute(
|
let mut folder = tx.prepare_cached(
|
||||||
"INSERT INTO folders(root_id, path, etag) VALUES (?1, ?2, ?3)
|
"INSERT INTO folders(root_id, path, etag) VALUES (?1, ?2, ?3)
|
||||||
ON CONFLICT(root_id, path) DO UPDATE SET etag = excluded.etag",
|
ON CONFLICT(root_id, path) DO UPDATE SET etag = excluded.etag",
|
||||||
rusqlite::params![root_id, path.as_str(), validator.as_str()],
|
|
||||||
)?;
|
)?;
|
||||||
|
for (path, validator) in &result.directories {
|
||||||
|
folder.execute(rusqlite::params![
|
||||||
|
root_id,
|
||||||
|
path.as_str(),
|
||||||
|
validator.as_str()
|
||||||
|
])?;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
for entry in &result.images {
|
// Every statement below runs once per photograph listed -- 1,600 for one
|
||||||
let folder_id: Option<i64> = entry.path.parent().and_then(|p| {
|
// folder whose sidecar changed, 24,000 for a first scan -- so each is
|
||||||
tx.query_row(
|
// prepared once, and a folder's id is looked up once per folder rather
|
||||||
"SELECT id FROM folders WHERE root_id = ?1 AND path = ?2",
|
// than once per photograph in it.
|
||||||
rusqlite::params![root_id, p.as_str()],
|
let mut folder_of =
|
||||||
|r| r.get(0),
|
tx.prepare_cached("SELECT id FROM folders WHERE root_id = ?1 AND path = ?2")?;
|
||||||
)
|
let mut folder_ids: std::collections::HashMap<String, Option<i64>> =
|
||||||
.ok()
|
std::collections::HashMap::new();
|
||||||
});
|
|
||||||
|
|
||||||
// TRACES: FR-PLAT-AND-2 | FR-CAT-9
|
// TRACES: FR-PLAT-AND-2 | FR-CAT-9
|
||||||
// The `availability` arm is what ends an offline library, and it does
|
// The `availability` arm is what ends an offline library, and it does
|
||||||
// it one photograph at a time. 3 is `Availability::Offline` and 0 is
|
// it one photograph at a time. 3 is `Availability::Offline` and 0 is
|
||||||
@@ -633,7 +637,7 @@ pub(super) fn persist(
|
|||||||
// `dr_catalog::walk` restores per file rather than per root: the only
|
// `dr_catalog::walk` restores per file rather than per root: the only
|
||||||
// thing that may clear "I could not reach this" is having reached it,
|
// thing that may clear "I could not reach this" is having reached it,
|
||||||
// and this statement runs precisely once per file the scan listed.
|
// and this statement runs precisely once per file the scan listed.
|
||||||
tx.execute(
|
let mut upsert = tx.prepare_cached(
|
||||||
"INSERT INTO images(root_id, folder_id, source_ref, format, file_size,
|
"INSERT INTO images(root_id, folder_id, source_ref, format, file_size,
|
||||||
availability, metadata_state, added_at)
|
availability, metadata_state, added_at)
|
||||||
VALUES (?1, ?2, ?3, ?4, ?5, 0, 1, ?6)
|
VALUES (?1, ?2, ?3, ?4, ?5, 0, 1, ?6)
|
||||||
@@ -642,7 +646,32 @@ pub(super) fn persist(
|
|||||||
folder_id = excluded.folder_id,
|
folder_id = excluded.folder_id,
|
||||||
availability = CASE WHEN images.availability = 3
|
availability = CASE WHEN images.availability = 3
|
||||||
THEN 0 ELSE images.availability END",
|
THEN 0 ELSE images.availability END",
|
||||||
rusqlite::params![
|
)?;
|
||||||
|
let mut image_of =
|
||||||
|
tx.prepare_cached("SELECT id FROM images WHERE root_id = ?1 AND source_ref = ?2")?;
|
||||||
|
let mut remote = tx.prepare_cached(
|
||||||
|
"INSERT INTO remote(image_id, file_id, etag, remote_path)
|
||||||
|
VALUES (?1, ?2, ?3, ?4)
|
||||||
|
ON CONFLICT(image_id) DO UPDATE SET
|
||||||
|
etag = excluded.etag, remote_path = excluded.remote_path",
|
||||||
|
)?;
|
||||||
|
|
||||||
|
for entry in &result.images {
|
||||||
|
let folder_id: Option<i64> = match entry.path.parent() {
|
||||||
|
None => None,
|
||||||
|
Some(p) => match folder_ids.get(p.as_str()) {
|
||||||
|
Some(id) => *id,
|
||||||
|
None => {
|
||||||
|
let id = folder_of
|
||||||
|
.query_row(rusqlite::params![root_id, p.as_str()], |r| r.get(0))
|
||||||
|
.ok();
|
||||||
|
folder_ids.insert(p.as_str().to_string(), id);
|
||||||
|
id
|
||||||
|
}
|
||||||
|
},
|
||||||
|
};
|
||||||
|
|
||||||
|
upsert.execute(rusqlite::params![
|
||||||
root_id,
|
root_id,
|
||||||
folder_id,
|
folder_id,
|
||||||
entry.path.as_str(),
|
entry.path.as_str(),
|
||||||
@@ -653,53 +682,41 @@ pub(super) fn persist(
|
|||||||
.map(|(_, e)| e.to_ascii_lowercase()),
|
.map(|(_, e)| e.to_ascii_lowercase()),
|
||||||
entry.size as i64,
|
entry.size as i64,
|
||||||
now_secs(),
|
now_secs(),
|
||||||
],
|
])?;
|
||||||
)?;
|
|
||||||
|
|
||||||
let image_id: i64 = tx.query_row(
|
let image_id: i64 = image_of
|
||||||
"SELECT id FROM images WHERE root_id = ?1 AND source_ref = ?2",
|
.query_row(rusqlite::params![root_id, entry.path.as_str()], |r| {
|
||||||
rusqlite::params![root_id, entry.path.as_str()],
|
r.get(0)
|
||||||
|r| r.get(0),
|
})?;
|
||||||
)?;
|
|
||||||
|
|
||||||
// Remote identity, keyed on oc:fileid so a server-side move is a move
|
// Remote identity, keyed on oc:fileid so a server-side move is a move
|
||||||
// rather than a re-download (FR-NC-5).
|
// rather than a re-download (FR-NC-5).
|
||||||
if let RemoteId::Stable(file_id) = entry.id {
|
if let RemoteId::Stable(file_id) = entry.id {
|
||||||
tx.execute(
|
remote.execute(rusqlite::params![
|
||||||
"INSERT INTO remote(image_id, file_id, etag, remote_path)
|
|
||||||
VALUES (?1, ?2, ?3, ?4)
|
|
||||||
ON CONFLICT(image_id) DO UPDATE SET
|
|
||||||
etag = excluded.etag, remote_path = excluded.remote_path",
|
|
||||||
rusqlite::params![
|
|
||||||
image_id,
|
image_id,
|
||||||
file_id as i64,
|
file_id as i64,
|
||||||
entry.validator.as_str(),
|
entry.validator.as_str(),
|
||||||
entry.path.as_str()
|
entry.path.as_str()
|
||||||
],
|
])?;
|
||||||
)?;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
tx.commit()?;
|
// Its thumbnail job, in the same transaction as the row it points at,
|
||||||
|
// so a failure mid-insert cannot leave one pointing at a row that
|
||||||
// Thumbnail jobs after the commit, so a failure mid-insert does not leave
|
// never landed. It used to follow the commit, one autocommitting
|
||||||
// jobs pointing at rows that never landed.
|
// statement per photograph, each finding its image again by path: a
|
||||||
for entry in &result.images {
|
// relisted folder of 1,600 paid 1,600 commits for rows that are
|
||||||
if let Ok(image_id) = conn.query_row(
|
// almost all already queued.
|
||||||
"SELECT id FROM images WHERE root_id = ?1 AND source_ref = ?2",
|
|
||||||
rusqlite::params![root_id, entry.path.as_str()],
|
|
||||||
|r| r.get::<_, i64>(0),
|
|
||||||
) {
|
|
||||||
let _ = dr_catalog::jobs::enqueue(
|
let _ = dr_catalog::jobs::enqueue(
|
||||||
conn,
|
&tx,
|
||||||
JobKind::Thumbnail,
|
JobKind::Thumbnail,
|
||||||
Some(image_id),
|
Some(image_id),
|
||||||
Priority::Background,
|
Priority::Background,
|
||||||
None,
|
None,
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
|
drop((folder_of, upsert, image_of, remote));
|
||||||
|
tx.commit()?;
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user