diff --git a/Cargo.lock b/Cargo.lock index 927cc4c8..744a340f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1373,6 +1373,7 @@ dependencies = [ "axum", "base64", "bcrypt", + "chrono", "clap", "clap-markdown", "clap_complete", diff --git a/Cargo.toml b/Cargo.toml index 70e52393..247e3637 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -28,7 +28,10 @@ libc = "0.2" axum = { version = "0.8", features = ["ws"] } base64 = "0.22" bcrypt = "0.19" -chrono = { version = "0.4", default-features = false, features = ["std"] } +chrono = { version = "0.4", default-features = false, features = [ + "serde", + "std", +] } clap = { version = "4", features = ["derive"] } clap_complete = "4" hive-sh4re = { path = "hive-sh4re" } diff --git a/hive-c0re/Cargo.toml b/hive-c0re/Cargo.toml index f97ef7b0..86502340 100644 --- a/hive-c0re/Cargo.toml +++ b/hive-c0re/Cargo.toml @@ -8,6 +8,7 @@ workspace = true [dependencies] anyhow.workspace = true axum.workspace = true +chrono.workspace = true base64.workspace = true bcrypt.workspace = true reqwest.workspace = true diff --git a/hive-c0re/src/approvals.rs b/hive-c0re/src/approvals.rs index 5254075f..6fb8dc6a 100644 --- a/hive-c0re/src/approvals.rs +++ b/hive-c0re/src/approvals.rs @@ -231,9 +231,9 @@ impl Approvals { agent: row.agent, kind: kind_from_str(&row.kind)?, commit_ref: row.commit_ref, - requested_at: hive_sh4re::wire_time::WireTime(row.requested_at), + requested_at: hive_sh4re::wire_time::from_secs(row.requested_at), status: ApprovalStatus::Approved, - resolved_at: Some(hive_sh4re::wire_time::WireTime(resolved_at)), + resolved_at: Some(hive_sh4re::wire_time::from_secs(resolved_at)), note: None, fetched_sha: row.fetched_sha, description: row.description, @@ -294,9 +294,9 @@ impl Approvals { agent: row.agent, kind: kind_from_str(&row.kind)?, commit_ref: row.commit_ref, - requested_at: hive_sh4re::wire_time::WireTime(row.requested_at), + requested_at: hive_sh4re::wire_time::from_secs(row.requested_at), status: ApprovalStatus::Cancelled, - resolved_at: Some(hive_sh4re::wire_time::WireTime(resolved_at)), + resolved_at: Some(hive_sh4re::wire_time::from_secs(resolved_at)), note: Some(note), fetched_sha: row.fetched_sha, description: row.description, @@ -404,11 +404,11 @@ fn row_to_approval(row: &rusqlite::Row<'_>) -> rusqlite::Result { agent: row.get(1)?, kind, commit_ref: row.get(3)?, - requested_at: hive_sh4re::wire_time::WireTime(row.get(4)?), + requested_at: hive_sh4re::wire_time::from_secs(row.get(4)?), status, resolved_at: row .get::<_, Option>(6)? - .map(hive_sh4re::wire_time::WireTime), + .map(hive_sh4re::wire_time::from_secs), note: row.get(7)?, fetched_sha: row.get(8)?, description: row.get(9)?, diff --git a/hive-c0re/src/audit_log.rs b/hive-c0re/src/audit_log.rs index b8e7e509..c145dde9 100644 --- a/hive-c0re/src/audit_log.rs +++ b/hive-c0re/src/audit_log.rs @@ -24,7 +24,8 @@ use std::sync::{Arc, Mutex, OnceLock}; use std::time::{SystemTime, UNIX_EPOCH}; use anyhow::{Context, Result}; -use hive_sh4re::wire_time::WireTime; + +use chrono::{DateTime, Utc}; use rusqlite::{Connection, params}; use serde::Serialize; @@ -90,7 +91,7 @@ impl AuditOutcome { #[derive(Debug, Clone, Serialize)] pub struct AuditEntry { pub id: i64, - pub ts_unix: WireTime, + pub ts_unix: DateTime, /// Agent on whose behalf the action was taken. pub agent: String, /// What was done (e.g. `restart_infra`). @@ -159,7 +160,7 @@ impl AuditLog { ) { Ok(_) => Some(AuditEntry { id: conn.last_insert_rowid(), - ts_unix: hive_sh4re::wire_time::WireTime(now), + ts_unix: hive_sh4re::wire_time::from_secs(now), agent: agent.to_owned(), action: action.to_owned(), target: target.to_owned(), @@ -252,7 +253,7 @@ pub fn spawn_vacuum(coord: &Arc) { fn row_to_entry(r: &rusqlite::Row) -> rusqlite::Result { Ok(AuditEntry { id: r.get(0)?, - ts_unix: hive_sh4re::wire_time::WireTime(r.get(1)?), + ts_unix: hive_sh4re::wire_time::from_secs(r.get(1)?), agent: r.get(2)?, action: r.get(3)?, target: r.get(4)?, diff --git a/hive-c0re/src/broker.rs b/hive-c0re/src/broker.rs index 09c3215c..269d7386 100644 --- a/hive-c0re/src/broker.rs +++ b/hive-c0re/src/broker.rs @@ -7,7 +7,8 @@ use std::sync::Mutex; use std::time::{SystemTime, UNIX_EPOCH}; use anyhow::{Context, Result}; -use hive_sh4re::wire_time::WireTime; + +use chrono::{DateTime, Utc}; use hive_sh4re::{InboxRow, Message}; use rusqlite::{Connection, OptionalExtension, params}; use serde::Serialize; @@ -75,8 +76,8 @@ pub struct PendingReminder { pub message: String, #[serde(skip_serializing_if = "Option::is_none")] pub file_path: Option, - pub due_at: WireTime, - pub created_at: WireTime, + pub due_at: DateTime, + pub created_at: DateTime, /// Most recent delivery failure for this row, if any. Cleared /// to NULL on operator retry. Surfaced inline in the dashboard /// so a stuck reminder doesn't just silently retry forever. @@ -925,8 +926,8 @@ impl Broker { agent: row.get(1)?, message: row.get(2)?, file_path: row.get(3)?, - due_at: hive_sh4re::wire_time::WireTime(row.get(4)?), - created_at: hive_sh4re::wire_time::WireTime(row.get(5)?), + due_at: hive_sh4re::wire_time::from_secs(row.get(4)?), + created_at: hive_sh4re::wire_time::from_secs(row.get(5)?), last_error: row.get(6)?, attempt_count: u32::try_from(attempts).unwrap_or(0), }) diff --git a/hive-c0re/src/coordinator.rs b/hive-c0re/src/coordinator.rs index fbb1bbda..b2e50f84 100644 --- a/hive-c0re/src/coordinator.rs +++ b/hive-c0re/src/coordinator.rs @@ -805,7 +805,7 @@ impl Coordinator { approval_kind, sha_short, status, - resolved_at: hive_sh4re::wire_time::WireTime(resolved_at), + resolved_at: hive_sh4re::wire_time::from_secs(resolved_at), note, description, }); @@ -838,8 +838,8 @@ impl Coordinator { question: question.to_owned(), options: options.to_vec(), multi, - asked_at: hive_sh4re::wire_time::WireTime(asked_at), - deadline_at: deadline_at.map(hive_sh4re::wire_time::WireTime), + asked_at: hive_sh4re::wire_time::from_secs(asked_at), + deadline_at: deadline_at.map(hive_sh4re::wire_time::from_secs), target: target.map(str::to_owned), question_refs, }); @@ -869,7 +869,7 @@ impl Coordinator { id, answer: answer.to_owned(), answerer: answerer.to_owned(), - answered_at: hive_sh4re::wire_time::WireTime(answered_at), + answered_at: hive_sh4re::wire_time::from_secs(answered_at), cancelled, target: target.map(str::to_owned), answer_refs, diff --git a/hive-c0re/src/dashboard.rs b/hive-c0re/src/dashboard.rs index e4f0709a..0f08c33d 100644 --- a/hive-c0re/src/dashboard.rs +++ b/hive-c0re/src/dashboard.rs @@ -27,7 +27,7 @@ use tokio_stream::{Stream, StreamExt}; use crate::container_view::{ContainerView, claude_has_session}; use crate::coordinator::Coordinator; use crate::lifecycle::{self, MANAGER_NAME}; -use hive_sh4re::wire_time::WireTime; +use chrono::{DateTime, Utc}; mod approvals; mod build_logs; @@ -434,7 +434,7 @@ struct ApprovalHistoryView { /// `approved` / `denied` / `failed`. status: &'static str, /// RFC 3339 UTC. Renders as a relative time on the dashboard. - resolved_at: WireTime, + resolved_at: DateTime, /// Operator-supplied deny reason (for `denied`) or build error /// (for `failed`). None on `approved`. #[serde(skip_serializing_if = "Option::is_none")] @@ -475,7 +475,7 @@ struct ApprovalView { /// RFC 3339 UTC time the approval was queued. Rendered as a /// relative time on the card so the operator can spot a stale /// request. - requested_at: WireTime, + requested_at: DateTime, } /// Replace silent `.unwrap_or_default()` on the data sources behind @@ -1100,7 +1100,7 @@ async fn dashboard_history(State(state): State) -> Response { from, to, body, - at: hive_sh4re::wire_time::WireTime(at), + at: hive_sh4re::wire_time::from_secs(at), in_reply_to, file_refs, }) @@ -1120,7 +1120,7 @@ async fn dashboard_history(State(state): State) -> Response { from, to, body, - at: hive_sh4re::wire_time::WireTime(at), + at: hive_sh4re::wire_time::from_secs(at), in_reply_to, file_refs, }) @@ -1375,7 +1375,7 @@ async fn api_operator_inbox(State(state): State) -> Response { "id": id, "from": from, "body": body, - "at": hive_sh4re::wire_time::WireTime(at), + "at": hive_sh4re::wire_time::from_secs(at), "in_reply_to": in_reply_to, "file_refs": file_refs, })) diff --git a/hive-c0re/src/dashboard_events.rs b/hive-c0re/src/dashboard_events.rs index 92e1c58e..21f5d820 100644 --- a/hive-c0re/src/dashboard_events.rs +++ b/hive-c0re/src/dashboard_events.rs @@ -9,7 +9,7 @@ use serde::Serialize; use crate::container_view::ContainerView; use crate::dashboard::{MetaInputView, TombstoneView}; use crate::rebuild_queue::QueueEntry; -use hive_sh4re::wire_time::WireTime; +use chrono::{DateTime, Utc}; #[derive(Debug, Clone, Serialize)] #[serde(rename_all = "snake_case", tag = "kind")] @@ -39,7 +39,7 @@ pub enum DashboardEvent { from: String, to: String, body: String, - at: WireTime, + at: DateTime, #[serde(default, skip_serializing_if = "Option::is_none")] in_reply_to: Option, #[serde(default, skip_serializing_if = "Vec::is_empty")] @@ -54,7 +54,7 @@ pub enum DashboardEvent { from: String, to: String, body: String, - at: WireTime, + at: DateTime, #[serde(default, skip_serializing_if = "Option::is_none")] in_reply_to: Option, #[serde(default, skip_serializing_if = "Vec::is_empty")] @@ -96,7 +96,7 @@ pub enum DashboardEvent { sha_short: Option, /// `"approved"` / `"denied"` / `"failed"`. status: &'static str, - resolved_at: WireTime, + resolved_at: DateTime, note: Option, description: Option, }, @@ -112,8 +112,8 @@ pub enum DashboardEvent { question: String, options: Vec, multi: bool, - asked_at: WireTime, - deadline_at: Option, + asked_at: DateTime, + deadline_at: Option>, target: Option, /// Verified file-path tokens that appear in `question`. /// Same shape as broker `Sent`/`Delivered` events; the @@ -131,7 +131,7 @@ pub enum DashboardEvent { id: i64, answer: String, answerer: String, - answered_at: WireTime, + answered_at: DateTime, cancelled: bool, target: Option, /// Verified file-path tokens that appear in `answer`. @@ -332,7 +332,7 @@ mod tests { from: "a".into(), to: "b".into(), body: String::new(), - at: hive_sh4re::wire_time::WireTime(0), + at: hive_sh4re::wire_time::from_secs(0), in_reply_to: None, file_refs: Vec::new(), }, @@ -342,7 +342,7 @@ mod tests { from: "a".into(), to: "b".into(), body: String::new(), - at: hive_sh4re::wire_time::WireTime(0), + at: hive_sh4re::wire_time::from_secs(0), in_reply_to: None, file_refs: Vec::new(), }, @@ -363,7 +363,7 @@ mod tests { approval_kind: "apply_commit", sha_short: None, status: "approved", - resolved_at: hive_sh4re::wire_time::WireTime(0), + resolved_at: hive_sh4re::wire_time::from_secs(0), note: None, description: None, }, @@ -374,7 +374,7 @@ mod tests { question: String::new(), options: Vec::new(), multi: false, - asked_at: hive_sh4re::wire_time::WireTime(0), + asked_at: hive_sh4re::wire_time::from_secs(0), deadline_at: None, target: None, question_refs: Vec::new(), @@ -384,7 +384,7 @@ mod tests { id: 1, answer: String::new(), answerer: "a".into(), - answered_at: hive_sh4re::wire_time::WireTime(0), + answered_at: hive_sh4re::wire_time::from_secs(0), cancelled: false, target: None, answer_refs: Vec::new(), @@ -447,7 +447,7 @@ mod tests { seq: 1, entry: crate::audit_log::AuditEntry { id: 1, - ts_unix: hive_sh4re::wire_time::WireTime(0), + ts_unix: hive_sh4re::wire_time::from_secs(0), agent: "atlas".into(), action: "restart_infra".into(), target: "hive-ci".into(), @@ -476,7 +476,7 @@ mod tests { seq: 7, entry: crate::audit_log::AuditEntry { id: 42, - ts_unix: hive_sh4re::wire_time::WireTime(1_700_000_000), + ts_unix: hive_sh4re::wire_time::from_secs(1_700_000_000), agent: "atlas".into(), action: "restart_infra".into(), target: "hive-gateway".into(), diff --git a/hive-c0re/src/loose_ends.rs b/hive-c0re/src/loose_ends.rs index 70399329..46f7e597 100644 --- a/hive-c0re/src/loose_ends.rs +++ b/hive-c0re/src/loose_ends.rs @@ -70,7 +70,7 @@ pub fn for_agent(coord: &Coordinator, agent: &str) -> Result> { agent: a.agent, commit_ref: a.commit_ref, description: a.description, - age_seconds: saturating_age(now, a.requested_at.secs()), + age_seconds: saturating_age(now, a.requested_at.timestamp()), }); } for q in coord.questions.pending_all()? { @@ -83,7 +83,7 @@ pub fn for_agent(coord: &Coordinator, agent: &str) -> Result> { asker: q.asker, target: q.target, question: q.question, - age_seconds: saturating_age(now, q.asked_at.secs()), + age_seconds: saturating_age(now, q.asked_at.timestamp()), }); } for r in coord.broker.list_pending_reminders()? { @@ -95,7 +95,7 @@ pub fn for_agent(coord: &Coordinator, agent: &str) -> Result> { owner: r.agent, message: r.message, due_at: r.due_at, - age_seconds: saturating_age(now, r.created_at.secs()), + age_seconds: saturating_age(now, r.created_at.timestamp()), }); } Ok(out) @@ -114,7 +114,7 @@ pub fn hive_wide(coord: &Coordinator) -> Result> { agent: a.agent, commit_ref: a.commit_ref, description: a.description, - age_seconds: saturating_age(now, a.requested_at.secs()), + age_seconds: saturating_age(now, a.requested_at.timestamp()), }); } for q in coord.questions.pending_all()? { @@ -123,7 +123,7 @@ pub fn hive_wide(coord: &Coordinator) -> Result> { asker: q.asker, target: q.target, question: q.question, - age_seconds: saturating_age(now, q.asked_at.secs()), + age_seconds: saturating_age(now, q.asked_at.timestamp()), }); } for r in coord.broker.list_pending_reminders()? { @@ -132,7 +132,7 @@ pub fn hive_wide(coord: &Coordinator) -> Result> { owner: r.agent, message: r.message, due_at: r.due_at, - age_seconds: saturating_age(now, r.created_at.secs()), + age_seconds: saturating_age(now, r.created_at.timestamp()), }); } Ok(out) diff --git a/hive-c0re/src/main.rs b/hive-c0re/src/main.rs index 46b123fb..a7546d63 100644 --- a/hive-c0re/src/main.rs +++ b/hive-c0re/src/main.rs @@ -497,7 +497,7 @@ fn spawn_broker_to_dashboard_forwarder(coord: Arc) { from, to, body, - at: hive_sh4re::wire_time::WireTime(at), + at: hive_sh4re::wire_time::from_secs(at), in_reply_to, file_refs, }); @@ -517,7 +517,7 @@ fn spawn_broker_to_dashboard_forwarder(coord: Arc) { from, to, body, - at: hive_sh4re::wire_time::WireTime(at), + at: hive_sh4re::wire_time::from_secs(at), in_reply_to, file_refs, }); diff --git a/hive-c0re/src/operator_questions.rs b/hive-c0re/src/operator_questions.rs index bf243846..604c43c8 100644 --- a/hive-c0re/src/operator_questions.rs +++ b/hive-c0re/src/operator_questions.rs @@ -14,7 +14,8 @@ use std::sync::Mutex; use std::time::{SystemTime, UNIX_EPOCH}; use anyhow::{Context, Result, bail}; -use hive_sh4re::wire_time::WireTime; + +use chrono::{DateTime, Utc}; use rusqlite::{Connection, OptionalExtension, params}; use serde::Serialize; @@ -75,12 +76,12 @@ pub struct OpQuestion { pub question: String, pub options: Vec, pub multi: bool, - pub asked_at: WireTime, + pub asked_at: DateTime, /// Deadline after which a watchdog auto-resolves the question with /// answer `[expired]`. `None` = no expiry. Surfaced on the /// dashboard as a remaining-time chip. - pub deadline_at: Option, - pub answered_at: Option, + pub deadline_at: Option>, + pub answered_at: Option>, pub answer: Option, /// Recipient of the question. `None` = the operator (dashboard /// path); `Some()` = a peer agent asked via @@ -288,14 +289,14 @@ fn row_to_question(row: &rusqlite::Row<'_>) -> rusqlite::Result { question: row.get(2)?, options, multi: multi != 0, - asked_at: hive_sh4re::wire_time::WireTime(row.get(5)?), + asked_at: hive_sh4re::wire_time::from_secs(row.get(5)?), answered_at: row .get::<_, Option>(6)? - .map(hive_sh4re::wire_time::WireTime), + .map(hive_sh4re::wire_time::from_secs), answer: row.get(7)?, deadline_at: row .get::<_, Option>(8)? - .map(hive_sh4re::wire_time::WireTime), + .map(hive_sh4re::wire_time::from_secs), target: row.get(9)?, }) } diff --git a/hive-c0re/src/socket_server.rs b/hive-c0re/src/socket_server.rs index 09af25d6..b4e4e14e 100644 --- a/hive-c0re/src/socket_server.rs +++ b/hive-c0re/src/socket_server.rs @@ -2018,8 +2018,8 @@ fn schedule_to_wire(s: crate::scheduled_prompts::Schedule) -> hive_sh4re::WireSc owner: s.owner, body: s.body, interval_seconds: s.interval_seconds, - next_fire_at_unix: hive_sh4re::wire_time::WireTime(s.next_fire_at_unix), - created_at_unix: hive_sh4re::wire_time::WireTime(s.created_at_unix), + next_fire_at_unix: hive_sh4re::wire_time::from_secs(s.next_fire_at_unix), + created_at_unix: hive_sh4re::wire_time::from_secs(s.created_at_unix), source: match s.source { crate::scheduled_prompts::ScheduleSource::Operator => { hive_sh4re::WireScheduleSource::Operator @@ -2028,16 +2028,16 @@ fn schedule_to_wire(s: crate::scheduled_prompts::Schedule) -> hive_sh4re::WireSc hive_sh4re::WireScheduleSource::Approval { id } } }, - cancelled_at_unix: s.cancelled_at_unix.map(hive_sh4re::wire_time::WireTime), - paused_at_unix: s.paused_at_unix.map(hive_sh4re::wire_time::WireTime), + cancelled_at_unix: s.cancelled_at_unix.map(hive_sh4re::wire_time::from_secs), + paused_at_unix: s.paused_at_unix.map(hive_sh4re::wire_time::from_secs), description: s.description, targets: s .targets .into_iter() .map(|t| hive_sh4re::WireScheduleTarget { target: t.target, - cancelled_at_unix: t.cancelled_at_unix.map(hive_sh4re::wire_time::WireTime), - last_fired_at_unix: t.last_fired_at_unix.map(hive_sh4re::wire_time::WireTime), + cancelled_at_unix: t.cancelled_at_unix.map(hive_sh4re::wire_time::from_secs), + last_fired_at_unix: t.last_fired_at_unix.map(hive_sh4re::wire_time::from_secs), last_result: t.last_result, }) .collect(), @@ -2131,8 +2131,8 @@ mod tests { owner: "operator".to_owned(), body: "ping".to_owned(), interval_seconds: None, - next_fire_at_unix: hive_sh4re::wire_time::WireTime(0), - created_at_unix: hive_sh4re::wire_time::WireTime(0), + next_fire_at_unix: hive_sh4re::wire_time::from_secs(0), + created_at_unix: hive_sh4re::wire_time::from_secs(0), source: hive_sh4re::WireScheduleSource::Operator, cancelled_at_unix: None, paused_at_unix: None, diff --git a/hive-sh4re/src/lib.rs b/hive-sh4re/src/lib.rs index 42c42c52..d8193b91 100644 --- a/hive-sh4re/src/lib.rs +++ b/hive-sh4re/src/lib.rs @@ -1,6 +1,6 @@ //! Wire types shared between `hive-c0re` and the in-container harness. -use crate::wire_time::WireTime; +use chrono::{DateTime, Utc}; use serde::{Deserialize, Serialize}; pub mod assets; @@ -210,10 +210,10 @@ pub struct Approval { /// hive-c0re refreshes this + re-renders the card for re-review. #[serde(default, skip_serializing_if = "Option::is_none")] pub fetched_sha: Option, - pub requested_at: WireTime, + pub requested_at: DateTime, pub status: ApprovalStatus, #[serde(default, skip_serializing_if = "Option::is_none")] - pub resolved_at: Option, + pub resolved_at: Option>, #[serde(default, skip_serializing_if = "Option::is_none")] pub note: Option, /// Free-text description the manager attached at submission time; @@ -449,7 +449,7 @@ pub enum LooseEnd { id: i64, owner: String, message: String, - due_at: WireTime, + due_at: DateTime, age_seconds: u64, }, /// Undelivered inbox messages waiting to be `recv`'d by this agent. @@ -1414,17 +1414,17 @@ pub struct WireSchedule { pub body: String, #[serde(default, skip_serializing_if = "Option::is_none")] pub interval_seconds: Option, - pub next_fire_at_unix: WireTime, - pub created_at_unix: WireTime, + pub next_fire_at_unix: DateTime, + pub created_at_unix: DateTime, pub source: WireScheduleSource, #[serde(default, skip_serializing_if = "Option::is_none")] - pub cancelled_at_unix: Option, + pub cancelled_at_unix: Option>, /// Set while the schedule is paused. Worker skips paused rows; /// they keep their `next_fire_at_unix` so resuming at any time /// fires at the next intended instant (no catch-up clamp needed /// — a paused schedule simply slips its next fire). #[serde(default, skip_serializing_if = "Option::is_none")] - pub paused_at_unix: Option, + pub paused_at_unix: Option>, #[serde(default, skip_serializing_if = "Option::is_none")] pub description: Option, pub targets: Vec, @@ -1441,9 +1441,9 @@ pub enum WireScheduleSource { pub struct WireScheduleTarget { pub target: String, #[serde(default, skip_serializing_if = "Option::is_none")] - pub cancelled_at_unix: Option, + pub cancelled_at_unix: Option>, #[serde(default, skip_serializing_if = "Option::is_none")] - pub last_fired_at_unix: Option, + pub last_fired_at_unix: Option>, #[serde(default, skip_serializing_if = "Option::is_none")] pub last_result: Option, } diff --git a/hive-sh4re/src/wire_time.rs b/hive-sh4re/src/wire_time.rs index 1c5422c7..b246b1cf 100644 --- a/hive-sh4re/src/wire_time.rs +++ b/hive-sh4re/src/wire_time.rs @@ -1,124 +1,38 @@ -//! Serde adaptors for timestamp fields: `i64` unix-epoch seconds in -//! Rust, RFC 3339 UTC strings (`2026-07-02T18:30:00Z`) in JSON. -//! -//! Rust code keeps doing plain integer arithmetic on these fields — -//! only the serialized representation changes, so the dashboard (and -//! any other JSON consumer) can feed the value straight into -//! `new Date(s)` without the `* 1000` epoch dance. -//! -//! Deserialization is lenient: both the RFC 3339 string form and the -//! legacy bare-integer form are accepted. That keeps a rolling deploy -//! safe (an old peer emitting epoch ints into a new reader) and lets -//! previously persisted JSON blobs re-load unchanged. -//! -//! Usage: type timestamp fields as [`WireTime`] / `Option` -//! (keep the usual `default` + `skip_serializing_if` attributes on the -//! optional form). The type IS the adaptor — no `#[serde(with = …)]` -//! needed. +//! Timestamp conventions for the wire types: fields are +//! `chrono::DateTime` (serde serializes them as RFC 3339 UTC `Z` +//! strings, e.g. `2026-07-02T18:30:00Z`), while sqlite storage and +//! agent-facing *input* args stay unix-epoch seconds (`i64`). This +//! module owns the two conversions at those boundaries. -use chrono::{DateTime, SecondsFormat, Utc}; -use serde::Deserialize; +use chrono::{DateTime, Utc}; -/// A timestamp on the wire: unix-epoch seconds in Rust, RFC 3339 UTC -/// string in JSON. Carries "this is a timestamp" in the type system -/// instead of a bare `i64` — the module docs above describe the wire -/// behaviour (ISO out, lenient epoch-or-ISO in). -/// -/// The inner value is public: arithmetic like `now + delay` stays -/// plain integer math (`WireTime(now_secs + delay)`), no chrono types -/// leak into call sites. -#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, PartialOrd, Ord, Hash)] -pub struct WireTime(pub i64); - -impl WireTime { - /// The wrapped unix-epoch seconds. - #[must_use] - pub fn secs(self) -> i64 { - self.0 - } - - /// RFC 3339 UTC `Z` string form (same as the serialized shape). - #[must_use] - pub fn to_iso(self) -> String { - to_iso(self.0) - } -} - -impl From for WireTime { - fn from(secs: i64) -> Self { - Self(secs) - } -} - -impl std::fmt::Display for WireTime { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - f.write_str(&to_iso(self.0)) - } -} - -impl serde::Serialize for WireTime { - fn serialize(&self, ser: S) -> Result { - ser.serialize_str(&to_iso(self.0)) - } -} - -impl<'de> Deserialize<'de> for WireTime { - fn deserialize>(de: D) -> Result { - EpochOrIso::deserialize(de)?.into_secs().map(Self) - } -} - -/// Format unix-epoch seconds as an RFC 3339 UTC string with a `Z` -/// suffix. Out-of-range values (never produced by our clocks) clamp to -/// the epoch rather than erroring — serialization must not fail. +/// Convert unix-epoch seconds (the sqlite column / input-arg form) +/// into the wire timestamp type. Out-of-range values (never produced +/// by our clocks) clamp to the epoch rather than erroring — the db +/// read path must not fail on a weird row. #[must_use] -pub fn to_iso(secs: i64) -> String { - DateTime::::from_timestamp(secs, 0) - .unwrap_or_default() - .to_rfc3339_opts(SecondsFormat::Secs, true) -} - -/// Parse an RFC 3339 string back to unix-epoch seconds. Any UTC offset -/// is accepted and normalized. -pub fn from_iso(s: &str) -> Result { - Ok(DateTime::parse_from_rfc3339(s)?.timestamp()) -} - -/// Lenient wire form: either the legacy epoch integer or the RFC 3339 -/// string. `untagged` tries the integer first (cheap), then the string. -#[derive(Deserialize)] -#[serde(untagged)] -enum EpochOrIso { - Epoch(i64), - Iso(String), -} - -impl EpochOrIso { - fn into_secs(self) -> Result { - match self { - Self::Epoch(secs) => Ok(secs), - Self::Iso(s) => from_iso(&s).map_err(E::custom), - } - } +pub fn from_secs(secs: i64) -> DateTime { + DateTime::::from_timestamp(secs, 0).unwrap_or_default() } #[cfg(test)] mod tests { + use chrono::{DateTime, Utc}; use serde::{Deserialize, Serialize}; - use super::WireTime; + use super::from_secs; #[derive(Serialize, Deserialize, PartialEq, Debug)] struct Row { - at: WireTime, + at: DateTime, #[serde(default, skip_serializing_if = "Option::is_none")] - maybe_at: Option, + maybe_at: Option>, } #[test] - fn serializes_epoch_as_rfc3339_z() { + fn serializes_as_rfc3339_z() { let json = serde_json::to_string(&Row { - at: WireTime(1_751_480_000), + at: from_secs(1_751_480_000), maybe_at: None, }) .unwrap(); @@ -128,8 +42,8 @@ mod tests { #[test] fn round_trips_and_serializes_some() { let row = Row { - at: WireTime(0), - maybe_at: Some(WireTime(1_751_480_000)), + at: from_secs(0), + maybe_at: Some(from_secs(1_751_480_000)), }; let json = serde_json::to_string(&row).unwrap(); assert_eq!( @@ -140,17 +54,13 @@ mod tests { } #[test] - fn deserializes_legacy_epoch_ints() { - // Rolling-deploy skew: an old writer still emits bare epoch - // integers — the lenient reader must accept them. - let row: Row = serde_json::from_str(r#"{"at":1751480000,"maybe_at":1751480000}"#).unwrap(); - assert_eq!(row.at, WireTime(1_751_480_000)); - assert_eq!(row.maybe_at, Some(WireTime(1_751_480_000))); + fn deserializes_offset_form_normalized_to_utc() { + let row: Row = serde_json::from_str(r#"{"at":"2025-07-02T20:13:20+02:00"}"#).unwrap(); + assert_eq!(row.at, from_secs(1_751_480_000)); } #[test] - fn deserializes_offset_form_normalized_to_utc() { - let row: Row = serde_json::from_str(r#"{"at":"2025-07-02T20:13:20+02:00"}"#).unwrap(); - assert_eq!(row.at, WireTime(1_751_480_000)); + fn from_secs_clamps_out_of_range_to_epoch() { + assert_eq!(from_secs(i64::MAX), from_secs(0)); } }