From 533ba3c15dac34cbabdf3bd1d1f575a73a9ab95c Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Mon, 7 Sep 2026 16:10:33 +0900 Subject: [PATCH] slop Signed-off-by: Seongmin Lee --- Cargo.lock | 2 + crates/sh_org/Cargo.toml | 2 + crates/sh_org/src/ingest.rs | 7 +- crates/sh_org/src/legacy/sh_tangled.rs | 12 +- crates/sh_org/src/migrate.rs | 490 ++++++++++++++++++++++++- crates/sh_org/src/schema.sql | 23 +- crates/sh_org/src/store.rs | 206 +++++++++-- 7 files changed, 682 insertions(+), 60 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 0f37595f9..4f933058d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -8445,9 +8445,11 @@ dependencies = [ "chrono", "futures", "jacquard-common", + "jacquard-repo", "lexicons", "rusqlite", "serde", + "serde_ipld_dagcbor", "serde_json", "tokio", "tokio-tungstenite 0.29.0", diff --git a/crates/sh_org/Cargo.toml b/crates/sh_org/Cargo.toml index 334935c88..85104e07f 100644 --- a/crates/sh_org/Cargo.toml +++ b/crates/sh_org/Cargo.toml @@ -14,9 +14,11 @@ anyhow = { workspace = true } chrono = { workspace = true } futures = { workspace = true } jacquard-common = { workspace = true } +jacquard-repo = { workspace = true } lexicons = { workspace = true } rusqlite = { version = "0.38", features = ["bundled"] } serde = { workspace = true } +serde_ipld_dagcbor = { workspace = true } serde_json = { workspace = true } tokio = { workspace = true } tokio-tungstenite = { workspace = true } diff --git a/crates/sh_org/src/ingest.rs b/crates/sh_org/src/ingest.rs index 35449bf3b..bb9364ca8 100644 --- a/crates/sh_org/src/ingest.rs +++ b/crates/sh_org/src/ingest.rs @@ -228,7 +228,9 @@ pub fn ingest(frame: &RecordFrame, seq: u64) -> Option { let mut record = Record { seq: seq as i64, - sh_uri, + sh_did: frame.did.clone(), + sh_nsid: frame.collection.clone(), + sh_rkey: frame.rkey.clone(), sh_cid: frame.cid.clone(), sh_record: body.map(RawValue::get).unwrap_or_default().to_owned(), created_at: created_at(frame, body), @@ -243,7 +245,7 @@ pub fn ingest(frame: &RecordFrame, seq: u64) -> Option { repo, }), Err(error) => { - warn!(uri = %record.sh_uri, %error, "unusable record, storing it as failed"); + warn!(uri = %record.sh_uri(), %error, "unusable record, storing it as failed"); record.note = Some(error.to_string()); Op::Put(Ingested { record, @@ -342,6 +344,7 @@ fn parse_meta(frame: &RecordFrame, body: &RawValue) -> anyhow::Result<(Vec, - pub subject: StarSubject, + pub subject: Subject, #[serde(flatten, default, skip_serializing_if = "Option::is_none")] pub extra_data: Extra, } + + /// `subject` was an at-uri to the repo record before it became a union + #[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq)] + #[serde(untagged)] + pub enum Subject { + Uri(AtUri), + Union(StarSubject), + } } pub mod graph_vouch { diff --git a/crates/sh_org/src/migrate.rs b/crates/sh_org/src/migrate.rs index c9af79f8c..d853e84e3 100644 --- a/crates/sh_org/src/migrate.rs +++ b/crates/sh_org/src/migrate.rs @@ -6,9 +6,33 @@ use std::time::{Duration, Instant}; -use tracing::{error, info, warn}; +use anyhow::Context; +use jacquard_common::types::collection::Collection; +use jacquard_common::types::did::Did; +use jacquard_common::types::string::{AtStrError, AtUri, RecordKey, Rkey, Tid}; +use jacquard_common::types::uri::UriValue; +use lexicons::com_atproto::repo::strong_ref::StrongRef; +use lexicons::org_tangled; +use lexicons::sh_tangled; +use serde::Serialize; +use tracing::{debug, error, info, warn}; -use crate::store::{Migration, Record, Store}; +use crate::legacy::sh_tangled as legacy; +use crate::store::{Migrated, Record, Store}; + +macro_rules! info_extra { + ($row:expr, $extra:expr) => {{ + let mut extra = $extra; + if let Some(fields) = &mut extra { + fields.remove("$type"); + } + let extra = extra.filter(|fields| !fields.is_empty()); + if let Some(fields) = &extra { + info!(uri = %$row.sh_uri(), field = stringify!($extra), data = ?fields, "extra data"); + } + extra + }}; +} const STATS_INTERVAL: Duration = Duration::from_secs(10); /// how long to wait before looking for work again on an empty queue @@ -42,36 +66,466 @@ pub fn worker(mut store: Store) { std::process::exit(1); } }; - let Some(record) = claimed else { + let Some(row) = claimed else { std::thread::sleep(IDLE); continue; }; // nothing marks a row in flight: a crash mid-record leaves it pending and it simply // runs again, so migrate() has to tolerate being called twice for the same body - let verdict = match migrate(&store, &record) { - Ok(out) => store.finish(&record, out), + let verdict = match migrate(&store, &row) { + Ok(out) => store.finish(&row, out), Err(error) => { - warn!(uri = %record.sh_uri, %error, "migration failed"); - store.defer(&record) + warn!(uri = %row.sh_uri(), %error, "migration failed"); + store.defer(&row) } }; if let Err(error) = verdict { // the same row would come straight back and spin the loop - error!(uri = %record.sh_uri, %error, "recording the verdict failed, giving up"); + error!(uri = %row.sh_uri(), %error, "recording the verdict failed, giving up"); std::process::exit(1); } } } -/// `(records, depends_on, repos) -> (records, migrated)`. Read this record's at-uri edges -/// with `store.deps()` and repo identity with `store.repo()`. -/// -/// An `Err` pushes the row's next attempt out and eventually lands it in `failed`, which is -/// how a dependency that hasn't been migrated yet is meant to be reported. Returning an empty -/// `Vec` marks the record processed without producing anything. -pub fn migrate(_store: &Store, record: &Record) -> anyhow::Result> { - info!(uri = %record.sh_uri, "migrating record"); - Ok(Vec::new()) - // todo!("sh.tangled.* record -> org.tangled.* records") +pub fn migrate(store: &Store, row: &Record) -> anyhow::Result> { + match row.sh_nsid.as_str() { + "sh.tangled.actor.profile" => actor_profile(store, row), + "sh.tangled.feed.comment" => feed_comment(store, row), + "sh.tangled.feed.reaction" => feed_reaction(store, row), + "sh.tangled.feed.star" => feed_star(store, row), + "sh.tangled.graph.follow" => graph_follow(store, row), + "sh.tangled.graph.vouch" => graph_vouch(store, row), + "sh.tangled.publicKey" => public_key(store, row), + "sh.tangled.repo" => repo(store, row), + // "sh.tangled.repo.artifact" => repo_artifact(store, row), + // "sh.tangled.repo.collaborator" => repo_collaborator(store, row), + // "sh.tangled.repo.collaboratorAcceptance" => repo_collaboratorAcceptance(store, row), + // "sh.tangled.repo.issue" => repo_issue(store, row), + // "sh.tangled.repo.issue.state" => repo_issue_state(store, row), + // "sh.tangled.repo.pull" => repo_pull(store, row), + // "sh.tangled.repo.pull.status" => repo_pull_status(store, row), + // "sh.tangled.string" => string(store, row), + + // no-op + "sh.tangled.knot" + | "sh.tangled.knot.member" + | "sh.tangled.knot.memberAcceptance" + | "sh.tangled.label.definition" + | "sh.tangled.label.op" + | "sh.tangled.repo.issue.comment" + | "sh.tangled.repo.pull.comment" + | "sh.tangled.spindle" + | "sh.tangled.spindle.member" => { + debug!(uri = %row.sh_uri(), "no-op"); + Ok(Vec::new()) + } + _ => { + error!(uri = %row.sh_uri(), "unknown collection"); + Ok(Vec::new()) + } + } +} + +fn actor_profile(store: &Store, row: &Record) -> anyhow::Result> { + let profile: sh_tangled::actor::profile::Profile = serde_json::from_str(&row.sh_record)?; + + let pinned_repositories = profile + .pinned_repositories + .map(|pins| { + pins.iter() + .filter(|pin| !pin.is_empty()) + .map(|pin| { + if pin.starts_with("did:") { + return Ok(Did::new_owned(pin)?); + } + let uri = AtUri::new_owned(pin) + .with_context(|| format!("invalid pinned repository: {pin:?}"))?; + store + .repo_did(&uri)? + .with_context(|| format!("failed to lookup repo: {uri}")) + }) + .collect::>>() + }) + .transpose()? + .filter(|pins| !pins.is_empty()); + + let mut links: Vec<_> = profile + .links + .unwrap_or_default() + .into_iter() + .filter(|link| !link.as_str().is_empty()) + .collect(); + + if profile.bluesky && links.len() < 5 { + let bluesky_link = UriValue::new_owned(format!("https://bsky.app/profile/{}", row.sh_did))?; + if !links.contains(&bluesky_link) { + links.insert(0, bluesky_link); + } + } + + Ok(vec![basic( + row, + org_tangled::actor::profile::Profile { + avatar: profile.avatar, + description: profile.description, + is_organization: None, + labels: None, + links: (!links.is_empty()).then_some(links), + location: profile.location, + pinned_repositories, + preferred_handle: profile.preferred_handle, + pronouns: profile.pronouns, + stats: profile.stats, + extra_data: info_extra!(row, profile.extra_data), + }, + )?]) +} + +fn feed_comment(store: &Store, row: &Record) -> anyhow::Result> { + let comment: sh_tangled::feed::comment::Comment = serde_json::from_str(&row.sh_record)?; + + let subject = migrated_ref(store, &comment.subject.uri)?; + let reply_to = comment + .reply_to + .as_ref() + .map(|reply_to| migrated_ref(store, &reply_to.uri)) + .transpose()?; + let embed = comment + .embed + .map(|embed| org_tangled::embed::commit::Commit { + repo: embed.repo, + commit: org_tangled::git::oid::Oid { + oid: embed.commit.oid, + extra_data: None, + }, + extra_data: info_extra!(row, embed.extra_data), + }); + + // TODO: when subject is ticket/PR, follow their new owner. + let owner: Did = row.sh_did.clone(); + + let body_blobs = None; // todo!(); + + Ok(vec![with_cid( + AtUri::from_parts_owned( + &owner, + ::NSID, + // TODO: hmmm can we keep same TID? + &row.sh_rkey, + )?, + org_tangled::feed::comment::Comment { + subject, + reply_to, + body: org_tangled::markup::markdown::Markdown { + text: comment.body.text, + blobs: body_blobs, + extra_data: info_extra!(row, comment.body.extra_data), + }, + embed, + created_at: row.created_at.clone(), + extra_data: info_extra!(row, comment.extra_data), + }, + )?]) +} + +/// `sh.tangled.feed.reaction` -> `org.tangled.feed.reaction` +fn feed_reaction(store: &Store, row: &Record) -> anyhow::Result> { + let reaction: legacy::FeedReaction = serde_json::from_str(&row.sh_record)?; + + let subject = migrated_ref(store, &reaction.subject)?; + + // TODO: when subject is ticket,PR or comment to ticket/PR, follow their new owner. + let owner: Did = row.sh_did.clone(); + + Ok(vec![with_cid( + AtUri::from_parts_owned( + &owner, + ::NSID, + // TODO: hmmm can we keep same TID? + &row.sh_rkey, + )?, + org_tangled::feed::reaction::Reaction { + created_at: row.created_at.clone(), + reaction: reaction.reaction, + subject, + extra_data: info_extra!(row, reaction.extra_data), + }, + )?]) +} + +/// `sh.tangled.feed.star` -> `org.tangled.feed.star` +fn feed_star(store: &Store, row: &Record) -> anyhow::Result> { + use legacy::feed_star::{StarSubject, Subject}; + use org_tangled::feed::star as org_star; + + let star: legacy::FeedStar = serde_json::from_str(&row.sh_record)?; + + let subject = match star.subject { + Subject::Union(StarSubject::Repo(repo)) => { + org_star::StarSubject::Repo(Box::new(org_star::Repo { + identity: repo.did, + extra_data: info_extra!(row, repo.extra_data), + })) + } + Subject::Union(StarSubject::String(string)) => { + org_star::StarSubject::Record(Box::new(org_star::Record { + record: migrated_ref(store, &string.uri)?, + extra_data: info_extra!(row, string.extra_data), + })) + } + Subject::Uri(uri) + if uri + .collection() + .is_some_and(|n| n.as_str() == ::NSID) => + { + org_star::StarSubject::Repo(Box::new(org_star::Repo { + identity: store + .repo_did(&uri)? + .with_context(|| format!("failed to lookup repo: {uri}"))?, + extra_data: None, + })) + } + Subject::Uri(uri) => org_star::StarSubject::Record(Box::new(org_star::Record { + record: migrated_ref(store, &uri)?, + extra_data: None, + })), + }; + + Ok(vec![basic( + row, + org_star::Star { + created_at: row.created_at.clone(), + subject, + extra_data: info_extra!(row, star.extra_data), + }, + )?]) +} + +/// `sh.tangled.graph.follow` -> `org.tangled.graph.follow` +fn graph_follow(_store: &Store, row: &Record) -> anyhow::Result> { + let follow: sh_tangled::graph::follow::Follow = serde_json::from_str(&row.sh_record)?; + + Ok(vec![basic( + row, + org_tangled::graph::follow::Follow { + subject: follow.subject, + created_at: row.created_at.clone(), + extra_data: info_extra!(row, follow.extra_data), + }, + )?]) +} + +/// `sh.tangled.graph.vouch` -> `org.tangled.graph.vouch` +fn graph_vouch(store: &Store, row: &Record) -> anyhow::Result> { + let vouch: sh_tangled::graph::vouch::Vouch = serde_json::from_str(&row.sh_record)?; + + // legacy keyed the record by the vouchee's did; the org record carries it as a field + let subject = Did::new_owned(row.sh_rkey.to_string()) + .with_context(|| format!("vouch rkey {} is not a did", row.sh_rkey))?; + let evidences = vouch + .evidences + .as_deref() + .map(|evidences| { + evidences + .iter() + .map(|uri| migrated_ref(store, uri)) + .collect::>>() + }) + .transpose()?; + + Ok(vec![with_cid( + AtUri::from_parts_owned( + &row.sh_did, + ::NSID, + // org.tangled.graph.vouch is tid-keyed and the legacy rkey is a did, so the + // rkey comes from created_at: a re-run of the same row lands on the same one. + // ponytail: two vouches by one author in the same microsecond would collide; + // bump the clkid off the seq if that ever shows up + Tid::from_time(row.created_at.timestamp_micros() as u64, 0).as_str(), + )?, + org_tangled::graph::vouch::Vouch { + subject, + kind: vouch.kind, + evidences, + reason: vouch.reason, + created_at: row.created_at.clone(), + extra_data: info_extra!(row, vouch.extra_data), + }, + )?]) +} + +/// `sh.tangled.publicKey` -> `org.tangled.key.sshKey` +fn public_key(_store: &Store, row: &Record) -> anyhow::Result> { + let key: sh_tangled::public_key::PublicKey = serde_json::from_str(&row.sh_record)?; + + Ok(vec![basic( + row, + org_tangled::key::ssh_key::SshKey { + name: key.name, + payload: key.key, + usages: Some(vec!["auth".into(), "sign".into()]), + created_at: row.created_at.clone(), + extra_data: info_extra!(row, key.extra_data), + }, + )?]) +} + +/// `sh.tangled.repo` -> `org.tangled.repo.declaration` + `org.tangled.repo.manifest` + +/// `org.tangled.label.definition[]` +fn repo(store: &Store, record: &Record) -> anyhow::Result> { + let repo: legacy::Repo = serde_json::from_str(&record.sh_record)?; + + let repo_did = repo.repo_did.context("repo carries no repoDid")?; + + let slug = repo + .name + .and_then(|name| Rkey::new_owned(name).ok()) + .unwrap_or_else(|| record.sh_rkey.clone()); + let slug = RecordKey(slug); + + Ok(vec![ + with_cid( + AtUri::from_parts_owned( + &record.sh_did, + ::NSID, + &slug.0, + )?, + org_tangled::repo::declaration::Declaration { + repo: repo_did.clone(), + created_at: record.created_at.clone(), + extra_data: None, + }, + )?, + with_cid( + AtUri::from_parts_owned( + &repo_did, + ::NSID, + "self", + )?, + org_tangled::repo::manifest::Manifest { + created_at: record.created_at.clone(), + declaration: org_tangled::repo::manifest::Declaration { + owner: record.sh_did.clone(), + slug, + extra_data: None, + }, + description: repo.description, + spindle: repo.spindle.as_deref().map(did_from_hostname).transpose()?, + topics: repo.topics, + upstream: repo + .source + .map(|source| upstream_from_uri(store, source)) + .transpose()? + .flatten(), + website: repo.website, + extra_data: info_extra!(record, repo.extra_data), + }, + )?, + // org.tangled.label.definition[] + ]) +} + +fn did_from_hostname(hostname: &str) -> Result { + let authority = hostname.to_ascii_lowercase().replace(':', "%3A"); + Did::new_owned(format!("did:web:{authority}")) +} + +fn upstream_from_uri( + store: &Store, + uri: UriValue, +) -> anyhow::Result> { + use org_tangled::repo::manifest::ManifestUpstream; + use org_tangled::repo::{SourceExternal, SourceTangled}; + + let tangled = match &uri { + UriValue::Did(did) => did.clone(), + UriValue::At(sh_uri) => store + .repo_did(sh_uri)? + .with_context(|| format!("failed to lookup repo: {sh_uri}"))?, + UriValue::Any(raw) if raw.is_empty() => return Ok(None), + _ => { + return Ok(Some(ManifestUpstream::SourceExternal(Box::new( + SourceExternal { + uri, + extra_data: None, + }, + )))); + } + }; + + Ok(Some(ManifestUpstream::SourceTangled(Box::new( + SourceTangled { + did: tangled, + extra_data: None, + }, + )))) +} + +#[allow(dead_code)] +fn repo_artifact(_store: &Store, _row: &Record) -> anyhow::Result> { + // TODO: what we are going to do with artifacts? LFS? `org.tangled.track.release`? + todo!() +} + +#[allow(dead_code)] +fn repo_collaborator(_store: &Store, _row: &Record) -> anyhow::Result> { + todo!() +} + +#[allow(dead_code)] +fn repo_issue(_store: &Store, _row: &Record) -> anyhow::Result> { + todo!() +} + +#[allow(dead_code)] +fn repo_issue_state(_store: &Store, _row: &Record) -> anyhow::Result> { + // upsert `org.tangled.track.ticket` records `state` + todo!() +} + +#[allow(dead_code)] +fn repo_pull(_store: &Store, _row: &Record) -> anyhow::Result> { + todo!() +} + +#[allow(dead_code)] +fn repo_pull_status(_store: &Store, _row: &Record) -> anyhow::Result> { + // upsert `org.tangled.track.ticket` records `state` + todo!() +} + +#[allow(dead_code)] +fn string(_store: &Store, _row: &Record) -> anyhow::Result> { + // NOTE: this will create *new* blobs + todo!() +} + +/// A strong ref to the `org.tangled.*` record a legacy at-uri migrated into. +fn migrated_ref(store: &Store, sh_uri: &AtUri) -> anyhow::Result { + let migrated = store + .migrated(sh_uri)? + .with_context(|| format!("{sh_uri} is not migrated yet"))?; + Ok(StrongRef::builder() + .uri(migrated.org_uri) + .cid(migrated.org_cid) + .build()) +} + +fn with_cid(uri: AtUri, record: T) -> anyhow::Result { + Ok(Migrated { + org_uri: uri, + org_cid: jacquard_repo::mst::util::compute_cid(&serde_ipld_dagcbor::to_vec(&record)?) + .map_err(|error| anyhow::anyhow!("{error}"))? + .into(), + org_record: serde_json::to_string(&record)?, + }) +} + +/// Basic `sh.tangled.*` -> `org.tangled.*` migration. same authority, same rkey. +fn basic(row: &Record, org_record: T) -> anyhow::Result { + with_cid( + AtUri::from_parts_owned(&row.sh_did, T::NSID, &row.sh_rkey)?, + org_record, + ) } diff --git a/crates/sh_org/src/schema.sql b/crates/sh_org/src/schema.sql index 7c8cba523..2b7865e32 100644 --- a/crates/sh_org/src/schema.sql +++ b/crates/sh_org/src/schema.sql @@ -12,14 +12,20 @@ CREATE TABLE IF NOT EXISTS meta ( -- pending records to migrate CREATE TABLE IF NOT EXISTS records ( seq INTEGER NOT NULL, - sh_uri TEXT PRIMARY KEY, + sh_did TEXT NOT NULL, + sh_nsid TEXT NOT NULL, + sh_rkey TEXT NOT NULL, + -- rust never selects this; it is here so depends_on and deletes can still address a + -- record by its full uri + sh_uri TEXT GENERATED ALWAYS AS ('at://' || sh_did || '/' || sh_nsid || '/' || sh_rkey) VIRTUAL, sh_cid TEXT, sh_record TEXT NOT NULL, created_at TEXT NOT NULL, -- for ordering state TEXT NOT NULL, retry_count INTEGER NOT NULL, -- retry info used on phase 2. retry_after INTEGER NOT NULL, - note TEXT -- why a `failed` row failed + note TEXT, -- why a `failed` row failed + PRIMARY KEY (sh_did, sh_nsid, sh_rkey) -- a generated column may not be one ); -- M:M mapping between records @@ -31,10 +37,12 @@ CREATE TABLE IF NOT EXISTS depends_on ( -- repo identity metadata CREATE TABLE IF NOT EXISTS repos ( - did TEXT PRIMARY KEY, - owner TEXT NOT NULL, - slug TEXT NOT NULL, - note TEXT + did TEXT PRIMARY KEY, + owner TEXT NOT NULL, + slug TEXT NOT NULL, + sh_rkey TEXT NOT NULL, -- rkey of `sh.tangled.repo` record + sh_uri TEXT GENERATED ALWAYS AS ('at://' || owner || '/sh.tangled.repo/' || sh_rkey) VIRTUAL, + note TEXT ); -- we will `ingest()` function at this phase: @@ -72,5 +80,8 @@ CREATE TABLE IF NOT EXISTS migrated ( -- phase 2's claim, minus retry_after: a `<= now` predicate can't live in a partial index CREATE INDEX IF NOT EXISTS records_queue ON records(created_at) WHERE state = 'pending'; +-- required, not an optimisation: depends_on's FK needs a unique index on its parent key +CREATE UNIQUE INDEX IF NOT EXISTS records_sh_uri ON records(sh_uri); CREATE INDEX IF NOT EXISTS depends_on_source ON depends_on(source_sh_uri); CREATE INDEX IF NOT EXISTS migrated_sh_uri ON migrated(sh_uri); +CREATE INDEX IF NOT EXISTS repos_sh_uri ON repos(sh_uri); diff --git a/crates/sh_org/src/store.rs b/crates/sh_org/src/store.rs index 56d8a9f3c..0a04509b2 100644 --- a/crates/sh_org/src/store.rs +++ b/crates/sh_org/src/store.rs @@ -4,7 +4,7 @@ use std::path::Path; -use jacquard_common::types::string::{AtUri, Cid, Datetime, Did, Rkey}; +use jacquard_common::types::string::{AtUri, Cid, Datetime, Did, Nsid, Rkey}; use rusqlite::{Connection, OptionalExtension}; use crate::now; @@ -22,7 +22,9 @@ pub struct Store { #[derive(Debug)] pub struct Record { pub seq: i64, - pub sh_uri: AtUri, + pub sh_did: Did, + pub sh_nsid: Nsid, + pub sh_rkey: Rkey, pub sh_cid: Option, /// the record body, byte for byte as hydrant sent it pub sh_record: String, @@ -32,19 +34,27 @@ pub struct Record { pub note: Option, } +impl Record { + /// What the `sh_uri` generated column spells out. The parts are validated types, so + /// `raw` can't panic; this is the same construction as `RecordFrame::uri()`. + pub fn sh_uri(&self) -> AtUri { + AtUri::raw(format!("at://{}/{}/{}", self.sh_did, self.sh_nsid, self.sh_rkey).into()) + } +} + /// a row of `repos` #[derive(Debug)] pub struct Repo { pub did: Did, pub owner: Did, - /// the repo's rkey, which is its name pub slug: Rkey, + pub sh_rkey: Rkey, pub note: Option, } /// a row of `migrated` #[derive(Debug)] -pub struct Migration { +pub struct Migrated { pub org_uri: AtUri, pub org_cid: Cid, pub org_record: String, @@ -141,9 +151,9 @@ impl Store { }) => { // a record arriving again is new work, whatever became of the last attempt tx.prepare_cached( - "INSERT INTO records(seq, sh_uri, sh_cid, sh_record, created_at, state, retry_count, retry_after, note) - VALUES(?1, ?2, ?3, ?4, ?5, ?6, 0, 0, ?7) - ON CONFLICT(sh_uri) DO UPDATE SET + "INSERT INTO records(seq, sh_did, sh_nsid, sh_rkey, sh_cid, sh_record, created_at, state, retry_count, retry_after, note) + VALUES(?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, 0, 0, ?9) + ON CONFLICT(sh_did, sh_nsid, sh_rkey) DO UPDATE SET seq = excluded.seq, sh_cid = excluded.sh_cid, sh_record = excluded.sh_record, @@ -155,7 +165,9 @@ impl Store { )? .execute(( record.seq, - record.sh_uri.as_str(), + record.sh_did.as_str(), + record.sh_nsid.as_str(), + record.sh_rkey.as_str(), record.sh_cid.as_ref().map(Cid::as_str), &record.sh_record, record.created_at.as_str(), @@ -163,24 +175,29 @@ impl Store { &record.note, ))?; // the edge set is whatever the newest body says it is + let sh_uri = record.sh_uri(); tx.prepare_cached("DELETE FROM depends_on WHERE source_sh_uri = ?1")? - .execute((record.sh_uri.as_str(),))?; + .execute((sh_uri.as_str(),))?; for target in depends_on { tx.prepare_cached( "INSERT INTO depends_on(source_sh_uri, target_sh_uri) VALUES(?1, ?2)", )? - .execute((record.sh_uri.as_str(), target.as_str()))?; + .execute((sh_uri.as_str(), target.as_str()))?; } if let Some(repo) = repo { tx.prepare_cached( - "INSERT INTO repos(did, owner, slug, note) VALUES(?1, ?2, ?3, ?4) + "INSERT INTO repos(did, owner, slug, sh_rkey, note) VALUES(?1, ?2, ?3, ?4, ?5) ON CONFLICT(did) DO UPDATE SET - owner = excluded.owner, slug = excluded.slug, note = excluded.note", + owner = excluded.owner, + slug = excluded.slug, + sh_rkey = excluded.sh_rkey, + note = excluded.note", )? .execute(( repo.did.as_str(), repo.owner.as_str(), repo.slug.as_str(), + repo.sh_rkey.as_str(), &repo.note, ))?; } @@ -210,34 +227,53 @@ impl Store { Ok(written) } - /// The oldest record still waiting to be migrated. + /// The oldest record whose dependencies are all out of the way. + /// + /// A target still `pending` blocks; anything else does not. A target that is `failed`, + /// `deleted`, or absent from `records` entirely — 1,634 of the edges on the network + /// point at records this mirror never saw — is never going to become `migrated`, and + /// waiting on one would strand its source here forever instead of letting `migrate()` + /// say so and land in `failed`. + /// + /// Two pending records that point at each other would block each other for good. + /// There are none on the network today. The `task queue` line is what would show it: + /// a `pending` count that stops moving. pub fn claim(&self) -> anyhow::Result> { let row = self .conn .query_row( - "SELECT seq, sh_uri, sh_cid, sh_record, created_at, note FROM records - WHERE state = 'pending' AND retry_after <= ?1 - ORDER BY created_at LIMIT 1", + "SELECT r.seq, r.sh_did, r.sh_nsid, r.sh_rkey, r.sh_cid, r.sh_record, r.created_at, r.note FROM records r + WHERE r.state = 'pending' AND r.retry_after <= ?1 + AND NOT EXISTS ( + SELECT 1 FROM depends_on d + JOIN records t ON t.sh_uri = d.target_sh_uri + WHERE d.source_sh_uri = r.sh_uri AND t.state = 'pending' + ) + ORDER BY r.created_at LIMIT 1", (now(),), |row| { Ok(( row.get::<_, i64>(0)?, row.get::<_, String>(1)?, - row.get::<_, Option>(2)?, + row.get::<_, String>(2)?, row.get::<_, String>(3)?, - row.get::<_, String>(4)?, - row.get::<_, Option>(5)?, + row.get::<_, Option>(4)?, + row.get::<_, String>(5)?, + row.get::<_, String>(6)?, + row.get::<_, Option>(7)?, )) }, ) .optional()?; - let Some((seq, sh_uri, sh_cid, sh_record, created_at, note)) = row else { + let Some((seq, sh_did, sh_nsid, sh_rkey, sh_cid, sh_record, created_at, note)) = row else { return Ok(None); }; // parsed out here, not in the row mapper: these are not rusqlite's errors Ok(Some(Record { seq, - sh_uri: AtUri::new_owned(sh_uri)?, + sh_did: Did::new_owned(sh_did)?, + sh_nsid: Nsid::new_owned(sh_nsid)?, + sh_rkey: Rkey::new_owned(sh_rkey)?, sh_cid: sh_cid.map(|c| c.parse()).transpose()?, sh_record, created_at: created_at.parse()?, @@ -256,43 +292,87 @@ impl Store { .collect() } + /// What a `sh.tangled.*` record turned into. + pub fn migrated(&self, sh_uri: &AtUri) -> anyhow::Result> { + let row = self + .conn + .prepare_cached( + "SELECT org_uri, org_cid, org_record FROM migrated WHERE sh_uri = ?1 LIMIT 1", + )? + .query_row((sh_uri.as_str(),), |row| { + Ok(( + row.get::<_, String>(0)?, + row.get::<_, String>(1)?, + row.get::<_, String>(2)?, + )) + }) + .optional()?; + let Some((org_uri, org_cid, org_record)) = row else { + return Ok(None); + }; + Ok(Some(Migrated { + org_uri: AtUri::new_owned(org_uri)?, + org_cid: org_cid.parse()?, + org_record, + })) + } + /// Repo identity by repo did, for `migrate()` to resolve did-shaped references. pub fn repo(&self, did: &Did) -> anyhow::Result> { let row = self .conn .query_row( - "SELECT did, owner, slug, note FROM repos WHERE did = ?1", + "SELECT did, owner, slug, sh_rkey, note FROM repos WHERE did = ?1", (did.as_str(),), |row| { Ok(( row.get::<_, String>(0)?, row.get::<_, String>(1)?, row.get::<_, String>(2)?, - row.get::<_, Option>(3)?, + row.get::<_, String>(3)?, + row.get::<_, Option>(4)?, )) }, ) .optional()?; - let Some((did, owner, slug, note)) = row else { + let Some((did, owner, slug, sh_rkey, note)) = row else { return Ok(None); }; Ok(Some(Repo { did: Did::new_owned(did)?, owner: Did::new_owned(owner)?, slug: Rkey::new_owned(slug)?, + sh_rkey: Rkey::new_owned(sh_rkey)?, note, })) } + pub fn repo_did(&self, sh_uri: &AtUri) -> anyhow::Result> { + let did = self + .conn + .prepare_cached("SELECT did FROM repos WHERE sh_uri = ?1 LIMIT 1")? + .query_row((sh_uri.as_str(),), |row| row.get::<_, String>(0)) + .optional()?; + Ok(did.map(Did::new_owned).transpose()?) + } + /// Marks the record processed and stores whatever it turned into. `out` may be empty: /// `migrated` means we are done with it, not that it produced anything. - pub fn finish(&mut self, record: &Record, out: Vec) -> anyhow::Result<()> { + pub fn finish(&mut self, record: &Record, out: Vec) -> anyhow::Result<()> { let tx = self.conn.transaction()?; // guarded on seq: if the stream re-ingested this uri while migrate() ran, the row is // pending again on purpose and the old body's verdict must not land on it let claimed = tx - .prepare_cached("UPDATE records SET state = 'migrated' WHERE sh_uri = ?1 AND seq = ?2")? - .execute((record.sh_uri.as_str(), record.seq))?; + .prepare_cached( + "UPDATE records SET state = 'migrated' + WHERE sh_did = ?1 AND sh_nsid = ?2 AND sh_rkey = ?3 AND seq = ?4", + )? + .execute(( + record.sh_did.as_str(), + record.sh_nsid.as_str(), + record.sh_rkey.as_str(), + record.seq, + ))?; if claimed == 0 { return Ok(()); } @@ -307,7 +387,7 @@ impl Store { seeded = 0", )? .execute(( - record.sh_uri.as_str(), + record.sh_uri().as_str(), m.org_uri.as_str(), m.org_cid.as_str(), &m.org_record, @@ -324,12 +404,14 @@ impl Store { .prepare_cached( "UPDATE records SET retry_count = retry_count + 1, - retry_after = ?3 + ?4 * (1 << retry_count), - state = CASE WHEN retry_count + 1 > ?5 THEN 'failed' ELSE state END - WHERE sh_uri = ?1 AND seq = ?2", + retry_after = ?5 + ?6 * (1 << retry_count), + state = CASE WHEN retry_count + 1 > ?7 THEN 'failed' ELSE state END + WHERE sh_did = ?1 AND sh_nsid = ?2 AND sh_rkey = ?3 AND seq = ?4", )? .execute(( - record.sh_uri.as_str(), + record.sh_did.as_str(), + record.sh_nsid.as_str(), + record.sh_rkey.as_str(), record.seq, now(), RETRY_BASE, @@ -361,3 +443,63 @@ impl Store { )?) } } + +#[cfg(test)] +mod tests { + use super::*; + + /// The parts go in, the parts come back, and the generated column agrees with what + /// `Record::sh_uri()` computes. That last equality is what everything addressing a + /// record by uri — `depends_on`, `Op::Delete`, `migrated` — rests on. + #[test] + fn record_round_trip() { + let dir = std::env::temp_dir().join("sh_org_record_round_trip"); + let _ = std::fs::remove_dir_all(&dir); + std::fs::create_dir_all(&dir).unwrap(); + let mut store = Store::open(&dir.join("t.db")).unwrap(); + + let record = Record { + seq: 7, + sh_did: Did::new_owned("did:plc:abc234abc234abc234abc234").unwrap(), + sh_nsid: Nsid::new_owned("sh.tangled.feed.star").unwrap(), + sh_rkey: Rkey::new_owned("3jzfcijpj2z2a").unwrap(), + sh_cid: None, + sh_record: "{}".to_owned(), + created_at: "2024-01-01T00:00:00Z".parse().unwrap(), + note: None, + }; + let uri = record.sh_uri(); + let target = AtUri::new_owned("at://did:plc:xyz234xyz234xyz234xyz234/sh.tangled.repo/r") + .unwrap(); + store + .write(vec![Op::Put(Ingested { + record, + state: State::Pending, + depends_on: vec![target], + repo: None, + })]) + .unwrap(); + + let stored: String = store + .conn + .query_one("SELECT sh_uri FROM records", [], |row| row.get(0)) + .unwrap(); + assert_eq!(stored, uri.as_str()); + + // the edge landed, so the FK against the generated column holds + assert_eq!(store.dependencies(&uri).unwrap().len(), 1); + + let claimed = store.claim().unwrap().expect("the pending record"); + assert_eq!(claimed.seq, 7); + assert_eq!(claimed.sh_did.as_str(), "did:plc:abc234abc234abc234abc234"); + assert_eq!(claimed.sh_nsid.as_str(), "sh.tangled.feed.star"); + assert_eq!(claimed.sh_rkey.as_str(), "3jzfcijpj2z2a"); + assert_eq!(claimed.sh_uri().as_str(), uri.as_str()); + + // finish() addresses the row by its parts, so the queue has to go empty + store.finish(&claimed, Vec::new()).unwrap(); + assert!(store.claim().unwrap().is_none()); + + std::fs::remove_dir_all(&dir).unwrap(); + } +} -- 2.51.2