diff --git a/mnemosyne_ingest/src/dispatch.rs b/mnemosyne_ingest/src/dispatch.rs index 7eb7779..2eb5cbb 100644 --- a/mnemosyne_ingest/src/dispatch.rs +++ b/mnemosyne_ingest/src/dispatch.rs @@ -115,15 +115,39 @@ impl Dispatcher { /// and advances the cursor in the same transaction. If no handler /// is registered for the event's NSID, the cursor still advances — /// we want to make progress past events we don't care about. + /// + /// Wraps the whole event in one tracing span — handlers and the + /// cursor advance run inside it, and handlers attach + /// `skip_reason` / repo-specific fields via `Span::current().record()` + /// so the closing log line is the wide event for "what happened + /// with this firehose record." + #[tracing::instrument( + name = "ingest.event", + skip(self, event), + fields( + event.nsid = %event.collection, + event.author_did = %event.did, + event.rkey = %event.rkey, + event.operation = ?event.operation, + event.time_us = event.time_us, + outcome = tracing::field::Empty, + skip_reason = tracing::field::Empty, + repo.did = tracing::field::Empty, + subject.did = tracing::field::Empty, + ), + )] pub async fn apply( &self, event: &CommitEvent, ) -> Result { + let span = tracing::Span::current(); let mut tx = self.pool.begin().await?; let outcome = if let Some(handler) = self.handlers.get(&event.collection) { handler.apply(&mut tx, event).await?; + span.record("outcome", "applied"); HandlerOutcome::Applied } else { + span.record("outcome", "no_handler"); HandlerOutcome::Skipped }; mnemosyne_postgres::advance_ingest_cursor(&mut *tx, &self.cursor_id, event.time_us) diff --git a/mnemosyne_ingest/src/handlers/public_key.rs b/mnemosyne_ingest/src/handlers/public_key.rs index 73c3310..ff086c1 100644 --- a/mnemosyne_ingest/src/handlers/public_key.rs +++ b/mnemosyne_ingest/src/handlers/public_key.rs @@ -55,26 +55,19 @@ impl EventHandler for PublicKeyHandler { tx: &mut Transaction<'_, Postgres>, event: &CommitEvent, ) -> Result<(), DispatchError> { + let span = tracing::Span::current(); match event.operation { CommitOperation::Create | CommitOperation::Update => { let Some(record_value) = event.record.as_ref() else { - eprintln!( - "ingest/{NSID}: skipping {operation:?} (did={did} rkey={rkey}): missing record body", - operation = event.operation, - did = event.did, - rkey = event.rkey, - ); + span.record("skip_reason", "missing_record_body"); + tracing::warn!("publicKey record body missing"); return Ok(()); }; let record: PublicKeyRecord = match serde_json::from_value(record_value.clone()) { - Ok(r) => r, - Err(e) => { - eprintln!( - "ingest/{NSID}: skipping {operation:?} (did={did} rkey={rkey}): record decode failed: {e}", - operation = event.operation, - did = event.did, - rkey = event.rkey, - ); + Ok(record) => record, + Err(error) => { + span.record("skip_reason", "record_decode_failed"); + tracing::warn!(error = %error, "publicKey record decode failed"); return Ok(()); } }; diff --git a/mnemosyne_ingest/src/handlers/repo.rs b/mnemosyne_ingest/src/handlers/repo.rs index d31ae85..5ab3b0e 100644 --- a/mnemosyne_ingest/src/handlers/repo.rs +++ b/mnemosyne_ingest/src/handlers/repo.rs @@ -51,56 +51,41 @@ impl EventHandler for RepoHandler { tx: &mut Transaction<'_, Postgres>, event: &CommitEvent, ) -> Result<(), DispatchError> { + let span = tracing::Span::current(); match event.operation { CommitOperation::Create | CommitOperation::Update => { let Some(record_value) = event.record.as_ref() else { - eprintln!( - "ingest/{NSID}: skipping {operation:?} (did={did} rkey={rkey}): missing record body", - operation = event.operation, - did = event.did, - rkey = event.rkey, - ); + span.record("skip_reason", "missing_record_body"); + tracing::warn!("repo record body missing"); return Ok(()); }; let record: RepoRecord = match serde_json::from_value(record_value.clone()) { - Ok(r) => r, - Err(e) => { - eprintln!( - "ingest/{NSID}: skipping {operation:?} (did={did} rkey={rkey}): record decode failed: {e}", - operation = event.operation, - did = event.did, - rkey = event.rkey, - ); + Ok(record) => record, + Err(error) => { + span.record("skip_reason", "record_decode_failed"); + tracing::warn!(error = %error, "repo record decode failed"); return Ok(()); } }; let Some(repo_did) = record.repo_did else { - eprintln!( - "ingest/{NSID}: skipping {operation:?} (did={did} rkey={rkey}): record has no repoDid", - operation = event.operation, - did = event.did, - rkey = event.rkey, - ); + span.record("skip_reason", "missing_repo_did"); + tracing::warn!("repo record has no repoDid; cannot link to a knot row"); return Ok(()); }; + span.record("repo.did", repo_did.as_str()); let owner_did = match lookup_repo_owner(&mut **tx, &repo_did).await? { - Some(o) => o, + Some(owner_did) => owner_did, None => { - eprintln!( - "ingest/{NSID}: skipping {operation:?} (did={did} rkey={rkey}): repo {repo_did} not registered", - operation = event.operation, - did = event.did, - rkey = event.rkey, - ); + span.record("skip_reason", "unknown_repo"); + tracing::warn!("repo not registered on this knot"); return Ok(()); } }; if owner_did != event.did { - eprintln!( - "ingest/{NSID}: refusing {operation:?} (did={did} rkey={rkey}): repo {repo_did} owned by {owner_did}, not record author", - operation = event.operation, - did = event.did, - rkey = event.rkey, + span.record("skip_reason", "author_not_owner"); + tracing::warn!( + owner.did = %owner_did, + "repo record author does not own the referenced repo", ); return Ok(()); } @@ -112,11 +97,8 @@ impl EventHandler for RepoHandler { // state from this; the Go knot doesn't either. A // future explicit deletion API will handle teardown // with proper authority. - eprintln!( - "ingest/{NSID}: noting Delete (did={did} rkey={rkey}); repo teardown is out-of-band", - did = event.did, - rkey = event.rkey, - ); + span.record("skip_reason", "delete_is_out_of_band"); + tracing::info!("noting repo Delete; teardown is out-of-band, not firehose-driven"); } } Ok(()) diff --git a/mnemosyne_ingest/src/handlers/repo_collaborator.rs b/mnemosyne_ingest/src/handlers/repo_collaborator.rs index eb23873..51195ce 100644 --- a/mnemosyne_ingest/src/handlers/repo_collaborator.rs +++ b/mnemosyne_ingest/src/handlers/repo_collaborator.rs @@ -55,51 +55,45 @@ impl EventHandler for RepoCollaboratorHandler { tx: &mut Transaction<'_, Postgres>, event: &CommitEvent, ) -> Result<(), DispatchError> { + let span = tracing::Span::current(); match event.operation { CommitOperation::Create | CommitOperation::Update => { let Some(record_value) = event.record.as_ref() else { - eprintln!( - "ingest/{NSID}: skipping {operation:?} (did={did} rkey={rkey}): missing record body", - operation = event.operation, - did = event.did, - rkey = event.rkey, - ); + span.record("skip_reason", "missing_record_body"); + tracing::warn!("collaborator record body missing"); return Ok(()); }; let record: CollaboratorRecord = match serde_json::from_value(record_value.clone()) { - Ok(r) => r, - Err(e) => { - eprintln!( - "ingest/{NSID}: skipping {operation:?} (did={did} rkey={rkey}): record decode failed: {e}", - operation = event.operation, - did = event.did, - rkey = event.rkey, + Ok(record) => record, + Err(error) => { + span.record("skip_reason", "record_decode_failed"); + tracing::warn!( + error = %error, + "collaborator record decode failed", ); return Ok(()); } }; + span.record("repo.did", record.repo.as_str()); + span.record("subject.did", record.subject.as_str()); let owner_did = match lookup_repo_owner(&mut **tx, &record.repo).await? { - Some(o) => o, + Some(owner_did) => owner_did, None => { - eprintln!( - "ingest/{NSID}: revoking prior link for {operation:?} (did={did} rkey={rkey}): repo {repo} not registered", - operation = event.operation, - did = event.did, - rkey = event.rkey, - repo = record.repo, + span.record("skip_reason", "unknown_repo_implicit_revoke"); + tracing::warn!( + "collaborator record targets unknown repo; revoking any prior link", ); delete_collaborator(&mut **tx, &event.did, &event.rkey).await?; return Ok(()); } }; if owner_did != event.did { - eprintln!( - "ingest/{NSID}: revoking prior link for {operation:?} (did={did} rkey={rkey}): repo {repo} owned by {owner_did}, not record author — collaborators cannot add collaborators", - operation = event.operation, - did = event.did, - rkey = event.rkey, - repo = record.repo, + span.record("skip_reason", "author_not_owner_implicit_revoke"); + tracing::warn!( + owner.did = %owner_did, + "collaborator record author does not own the referenced repo; \ + collaborators cannot add collaborators — revoking any prior link", ); delete_collaborator(&mut **tx, &event.did, &event.rkey).await?; return Ok(()); diff --git a/mnemosyne_ingest/src/main.rs b/mnemosyne_ingest/src/main.rs index a4cb2e8..3903fbc 100644 --- a/mnemosyne_ingest/src/main.rs +++ b/mnemosyne_ingest/src/main.rs @@ -92,11 +92,11 @@ async fn main() -> anyhow::Result<()> { let cursor = read_ingest_cursor(&pool, &cfg.cursor_id) .await .context("read cursor")?; - eprintln!( - "ingest: starting from cursor {}", - cursor - .map(|c| c.to_string()) - .unwrap_or_else(|| "(live tail — no persisted cursor)".to_string()), + tracing::info!( + ingest.cursor_id = %cfg.cursor_id, + ingest.cursor = ?cursor, + ingest.jetstream_url = %cfg.jetstream_url, + "ingest starting", ); let mut dispatcher = Dispatcher::new(pool, cfg.cursor_id.clone()); @@ -117,7 +117,7 @@ async fn main() -> anyhow::Result<()> { let shutdown = Box::pin(async { let _ = tokio::signal::ctrl_c().await; - eprintln!("ingest: shutdown requested"); + tracing::info!("ingest shutdown requested"); }); match subscriber.run(shutdown).await { diff --git a/mnemosyne_ingest/src/subscriber.rs b/mnemosyne_ingest/src/subscriber.rs index 9aa773f..a33faa3 100644 --- a/mnemosyne_ingest/src/subscriber.rs +++ b/mnemosyne_ingest/src/subscriber.rs @@ -95,12 +95,20 @@ impl Subscriber { return outcome; } // Any other error: log + backoff + retry. - if let Err(e) = outcome { - eprintln!("ingest: stream disconnected: {e}"); + if let Err(error) = outcome { + tracing::warn!( + error = %error, + ingest.backoff_ms = backoff.as_millis() as u64, + "jetstream stream disconnected; will reconnect", + ); } } - Err(e) => { - eprintln!("ingest: connect failed: {e}"); + Err(error) => { + tracing::warn!( + error = %error, + ingest.backoff_ms = backoff.as_millis() as u64, + "jetstream connect failed; will retry", + ); } } } @@ -145,15 +153,21 @@ impl Subscriber { let event = match decode_commit_line(frame) { Ok(Some(event)) => event, Ok(None) => return, - Err(e) => { - eprintln!("ingest: malformed jetstream frame: {e}"); + Err(error) => { + tracing::warn!(error = %error, "malformed jetstream frame; skipping"); return; } }; - if let Err(e) = self.dispatcher.apply(&event).await { - eprintln!( - "ingest: handler/cursor failed for {} {} at {}: {e}", - event.collection, event.rkey, event.time_us, + if let Err(error) = self.dispatcher.apply(&event).await { + // Infra-level failure (db unreachable, schema mismatch). + // Distinct from per-event policy skips, which are logged + // inside the handler span without bubbling here. + tracing::error!( + error = %error, + event.nsid = %event.collection, + event.rkey = %event.rkey, + event.time_us = event.time_us, + "ingest dispatcher failure", ); } }