use std::fmt;
use std::pin::Pin;
use std::str::FromStr;
use async_trait::async_trait;
use futures::Stream;
use serde::{Deserialize, Serialize};
use crate::kb::ObjectMeta;
use crate::object_path::ObjectPath;
use crate::slug::KbSlug;
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
#[serde(transparent)]
pub struct EventId(pub u64);
impl fmt::Display for EventId {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "{}", self.0)
}
}
impl FromStr for EventId {
type Err = std::num::ParseIntError;
fn from_str(s: &str) -> Result<Self, Self::Err> {
s.parse::<u64>().map(Self)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum EventSource {
Http,
Webdav,
Mcp,
FsWatch,
Reconcile,
Indexer,
}
impl EventSource {
#[must_use]
pub fn as_str(self) -> &'static str {
match self {
Self::Http => "http",
Self::Webdav => "webdav",
Self::Mcp => "mcp",
Self::FsWatch => "fs-watch",
Self::Reconcile => "reconcile",
Self::Indexer => "indexer",
}
}
}
impl fmt::Display for EventSource {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "event")]
pub enum ObjectEventKind {
#[serde(rename = "object.written")]
Written {
etag: String,
size: u64,
mime: String,
mtime: i64,
},
#[serde(rename = "object.deleted")]
Deleted,
#[serde(rename = "object.indexed")]
Indexed {
etag: String,
mime: String,
chunks: u32,
},
#[serde(rename = "object.index_failed")]
IndexFailed {
#[serde(default, skip_serializing_if = "Option::is_none")]
etag: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
mime: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
summary: Option<String>,
},
}
impl ObjectEventKind {
#[must_use]
pub fn name(&self) -> &'static str {
match self {
Self::Written { .. } => "object.written",
Self::Deleted => "object.deleted",
Self::Indexed { .. } => "object.indexed",
Self::IndexFailed { .. } => "object.index_failed",
}
}
#[must_use]
pub fn mime(&self) -> Option<&str> {
match self {
Self::Written { mime, .. } | Self::Indexed { mime, .. } => Some(mime),
Self::IndexFailed { mime, .. } => mime.as_deref(),
Self::Deleted => None,
}
}
#[must_use]
pub fn etag(&self) -> Option<&str> {
match self {
Self::Written { etag, .. } | Self::Indexed { etag, .. } => Some(etag),
Self::IndexFailed { etag, .. } => etag.as_deref(),
Self::Deleted => None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ObjectEvent {
pub kb: KbSlug,
pub object_key: ObjectPath,
#[serde(flatten)]
pub kind: ObjectEventKind,
pub source: EventSource,
pub occurred_at: String,
}
impl ObjectEvent {
#[must_use]
pub fn written(
kb: KbSlug,
object_key: ObjectPath,
etag: String,
size: u64,
content_type: String,
last_modified: i64,
source: EventSource,
) -> Self {
Self {
kb,
object_key,
kind: ObjectEventKind::Written {
etag,
size,
mime: content_type,
mtime: last_modified,
},
source,
occurred_at: now_rfc3339(),
}
}
#[must_use]
pub fn deleted(kb: KbSlug, object_key: ObjectPath, source: EventSource) -> Self {
Self {
kb,
object_key,
kind: ObjectEventKind::Deleted,
source,
occurred_at: now_rfc3339(),
}
}
#[must_use]
pub fn indexed(
kb: KbSlug,
object_key: ObjectPath,
etag: String,
mime: String,
chunks: u32,
) -> Self {
Self {
kb,
object_key,
kind: ObjectEventKind::Indexed { etag, mime, chunks },
source: EventSource::Indexer,
occurred_at: now_rfc3339(),
}
}
#[must_use]
pub fn index_failed(
kb: KbSlug,
object_key: ObjectPath,
etag: Option<String>,
mime: Option<String>,
summary: String,
) -> Self {
Self {
kb,
object_key,
kind: ObjectEventKind::IndexFailed {
etag,
mime,
summary: Some(summary),
},
source: EventSource::Indexer,
occurred_at: now_rfc3339(),
}
}
#[must_use]
pub fn from_meta(
kb: KbSlug,
object_key: ObjectPath,
meta: &ObjectMeta,
source: EventSource,
) -> Self {
Self::written(
kb,
object_key,
meta.etag.clone().unwrap_or_default(),
meta.size,
meta.content_type.clone().unwrap_or_default(),
meta.last_modified.unwrap_or(0),
source,
)
}
}
#[derive(Debug, thiserror::Error)]
pub enum PublishError {
#[error("event backend unavailable: {message}")]
Unavailable {
message: String,
},
}
#[derive(Debug, thiserror::Error)]
pub enum SubscribeError {
#[error("events after {requested} are no longer retained; the oldest retained id is {oldest}")]
Gone {
requested: EventId,
oldest: EventId,
},
#[error("event backend unavailable: {message}")]
Unavailable {
message: String,
},
}
#[derive(Debug, thiserror::Error)]
pub enum StreamError {
#[error("subscriber lagged behind the event log; resume after {resume_after}")]
Lagged {
resume_after: EventId,
},
#[error("event backend failed mid-stream: {message}")]
Broker {
message: String,
},
}
pub type EventStream =
Pin<Box<dyn Stream<Item = Result<(EventId, ObjectEvent), StreamError>> + Send>>;
#[async_trait]
pub trait EventPublisher: Send + Sync {
async fn publish(&self, event: ObjectEvent) -> Result<EventId, PublishError>;
async fn subscribe(
&self,
kb: &KbSlug,
after: Option<EventId>,
) -> Result<EventStream, SubscribeError>;
fn ready(&self) -> bool;
fn backend_name(&self) -> &'static str;
}
fn now_rfc3339() -> String {
let secs = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| i64::try_from(d.as_secs()).unwrap_or(i64::MAX));
unix_to_rfc3339(secs)
}
#[must_use]
pub fn unix_to_rfc3339(secs: i64) -> String {
let secs = secs.max(0);
let days = secs.div_euclid(86_400);
let rem = secs.rem_euclid(86_400);
let (hh, mm, ss) = (rem / 3600, (rem % 3600) / 60, rem % 60);
let z = days + 719_468;
let era = z.div_euclid(146_097);
let doe = z.rem_euclid(146_097);
let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365;
let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
let mp = (5 * doy + 2) / 153;
let day = doy - (153 * mp + 2) / 5 + 1;
let month = if mp < 10 { mp + 3 } else { mp - 9 };
let year = yoe + era * 400 + i64::from(month <= 2);
format!("{year:04}-{month:02}-{day:02}T{hh:02}:{mm:02}:{ss:02}Z")
}
#[cfg(test)]
mod tests {
use super::*;
fn kb() -> KbSlug {
KbSlug::try_new("notes").expect("valid slug")
}
fn key(s: &str) -> ObjectPath {
ObjectPath::try_from_str(s).expect("valid path")
}
#[test]
fn a_write_serialises_to_the_documented_shape() {
let mut event = ObjectEvent::written(
kb(),
key("inbox/memo.mp3"),
"\"9a3f\"".into(),
48_213_011,
"audio/mpeg".into(),
1_757_950_000,
EventSource::Http,
);
event.occurred_at = "2026-09-15T14:33:20Z".into();
let json = serde_json::to_value(&event).expect("serialises");
assert_eq!(
json,
serde_json::json!({
"event": "object.written",
"kb": "notes",
"object_key": "inbox/memo.mp3",
"etag": "\"9a3f\"",
"size": 48_213_011,
"mime": "audio/mpeg",
"mtime": 1_757_950_000,
"source": "http",
"occurred_at": "2026-09-15T14:33:20Z",
})
);
let back: ObjectEvent = serde_json::from_value(json).expect("round trips");
assert_eq!(back, event);
}
#[test]
fn a_delete_carries_no_stamp_and_round_trips() {
let event = ObjectEvent::deleted(kb(), key("old.md"), EventSource::Webdav);
let json = serde_json::to_value(&event).expect("serialises");
assert_eq!(json["event"], "object.deleted");
assert_eq!(json["source"], "webdav");
assert!(json.get("etag").is_none());
assert!(json.get("mime").is_none());
let back: ObjectEvent = serde_json::from_value(json).expect("round trips");
assert_eq!(back, event);
assert_eq!(event.kind.name(), "object.deleted");
assert_eq!(event.kind.mime(), None);
}
#[test]
fn an_indexed_event_serialises_to_the_documented_shape() {
let mut event = ObjectEvent::indexed(
kb(),
key("inbox/memo.md"),
"\"b71c\"".into(),
"text/markdown".into(),
3,
);
event.occurred_at = "2026-09-15T14:33:21Z".into();
let json = serde_json::to_value(&event).expect("serialises");
assert_eq!(
json,
serde_json::json!({
"event": "object.indexed",
"kb": "notes",
"object_key": "inbox/memo.md",
"etag": "\"b71c\"",
"mime": "text/markdown",
"chunks": 3,
"source": "indexer",
"occurred_at": "2026-09-15T14:33:21Z",
})
);
let back: ObjectEvent = serde_json::from_value(json).expect("round trips");
assert_eq!(back, event);
assert_eq!(event.kind.name(), "object.indexed");
assert_eq!(event.kind.mime(), Some("text/markdown"));
assert_eq!(event.kind.etag(), Some("\"b71c\""));
}
#[test]
fn an_index_failure_omits_what_it_does_not_know() {
let blind = ObjectEvent::index_failed(
kb(),
key("a.md"),
None,
None,
"storage.head_object failed: backend unavailable".into(),
);
let json = serde_json::to_value(&blind).expect("serialises");
assert_eq!(json["event"], "object.index_failed");
assert_eq!(json["source"], "indexer");
assert!(json.get("etag").is_none());
assert!(json.get("mime").is_none());
assert_eq!(
json["summary"],
"storage.head_object failed: backend unavailable"
);
assert_eq!(blind.kind.mime(), None);
assert_eq!(blind.kind.etag(), None);
let mut stamped = ObjectEvent::index_failed(
kb(),
key("a.md"),
Some("\"e1\"".into()),
Some("text/plain".into()),
"embedder.embed failed: connection refused".into(),
);
stamped.occurred_at = "2026-09-15T14:33:22Z".into();
let json = serde_json::to_value(&stamped).expect("serialises");
assert_eq!(
json,
serde_json::json!({
"event": "object.index_failed",
"kb": "notes",
"object_key": "a.md",
"etag": "\"e1\"",
"mime": "text/plain",
"summary": "embedder.embed failed: connection refused",
"source": "indexer",
"occurred_at": "2026-09-15T14:33:22Z",
})
);
let back: ObjectEvent = serde_json::from_value(json).expect("round trips");
assert_eq!(back, stamped);
assert_eq!(stamped.kind.name(), "object.index_failed");
assert_eq!(stamped.kind.mime(), Some("text/plain"));
assert_eq!(stamped.kind.etag(), Some("\"e1\""));
}
#[test]
fn a_withheld_summary_reads_back_as_absent() {
let json = serde_json::json!({
"event": "object.index_failed",
"kb": "notes",
"object_key": "a.md",
"etag": "\"e1\"",
"mime": "text/plain",
"source": "indexer",
"occurred_at": "2026-09-15T14:33:22Z",
});
let event: ObjectEvent = serde_json::from_value(json).expect("deserialises");
assert_eq!(
event.kind,
ObjectEventKind::IndexFailed {
etag: Some("\"e1\"".into()),
mime: Some("text/plain".into()),
summary: None,
}
);
}
#[test]
fn sources_use_the_documented_spellings() {
for (source, wire) in [
(EventSource::Http, "http"),
(EventSource::Webdav, "webdav"),
(EventSource::Mcp, "mcp"),
(EventSource::FsWatch, "fs-watch"),
(EventSource::Reconcile, "reconcile"),
(EventSource::Indexer, "indexer"),
] {
assert_eq!(source.as_str(), wire);
assert_eq!(serde_json::to_value(source).unwrap(), wire);
}
}
#[test]
fn from_meta_takes_the_head_stamp() {
let meta = ObjectMeta {
key: "a.md".into(),
size: 12,
last_modified: Some(1_700_000_000),
content_type: Some("text/markdown".into()),
etag: Some("\"e1\"".into()),
};
let event = ObjectEvent::from_meta(kb(), key("a.md"), &meta, EventSource::FsWatch);
assert_eq!(
event.kind,
ObjectEventKind::Written {
etag: "\"e1\"".into(),
size: 12,
mime: "text/markdown".into(),
mtime: 1_700_000_000,
}
);
assert_eq!(event.source, EventSource::FsWatch);
}
#[test]
fn event_ids_render_and_parse_as_decimals() {
assert_eq!(EventId(4812).to_string(), "4812");
assert_eq!("4812".parse::<EventId>().unwrap(), EventId(4812));
assert!("".parse::<EventId>().is_err());
assert!("-1".parse::<EventId>().is_err());
assert!("abc".parse::<EventId>().is_err());
}
#[test]
fn rfc3339_rendering_matches_known_instants() {
assert_eq!(unix_to_rfc3339(0), "1970-01-01T00:00:00Z");
assert_eq!(unix_to_rfc3339(1_757_950_000), "2025-09-15T15:26:40Z");
assert_eq!(unix_to_rfc3339(951_782_400), "2000-02-29T00:00:00Z");
assert_eq!(unix_to_rfc3339(4_107_542_399), "2100-02-28T23:59:59Z");
assert_eq!(unix_to_rfc3339(-5), "1970-01-01T00:00:00Z");
}
}