diff --git a/Cargo.lock b/Cargo.lock index 744a340f..927cc4c8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1373,7 +1373,6 @@ dependencies = [ "axum", "base64", "bcrypt", - "chrono", "clap", "clap-markdown", "clap_complete", diff --git a/Cargo.toml b/Cargo.toml index 247e3637..70e52393 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -28,10 +28,7 @@ libc = "0.2" axum = { version = "0.8", features = ["ws"] } base64 = "0.22" bcrypt = "0.19" -chrono = { version = "0.4", default-features = false, features = [ - "serde", - "std", -] } +chrono = { version = "0.4", default-features = false, features = ["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 86502340..f97ef7b0 100644 --- a/hive-c0re/Cargo.toml +++ b/hive-c0re/Cargo.toml @@ -8,7 +8,6 @@ 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 6fb8dc6a..652df40e 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::from_secs(row.requested_at), + requested_at: row.requested_at, status: ApprovalStatus::Approved, - resolved_at: Some(hive_sh4re::wire_time::from_secs(resolved_at)), + resolved_at: Some(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::from_secs(row.requested_at), + requested_at: row.requested_at, status: ApprovalStatus::Cancelled, - resolved_at: Some(hive_sh4re::wire_time::from_secs(resolved_at)), + resolved_at: Some(resolved_at), note: Some(note), fetched_sha: row.fetched_sha, description: row.description, @@ -404,11 +404,9 @@ 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::from_secs(row.get(4)?), + requested_at: row.get(4)?, status, - resolved_at: row - .get::<_, Option>(6)? - .map(hive_sh4re::wire_time::from_secs), + resolved_at: row.get(6)?, 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 c145dde9..43f1772a 100644 --- a/hive-c0re/src/audit_log.rs +++ b/hive-c0re/src/audit_log.rs @@ -24,8 +24,6 @@ use std::sync::{Arc, Mutex, OnceLock}; use std::time::{SystemTime, UNIX_EPOCH}; use anyhow::{Context, Result}; - -use chrono::{DateTime, Utc}; use rusqlite::{Connection, params}; use serde::Serialize; @@ -91,7 +89,8 @@ impl AuditOutcome { #[derive(Debug, Clone, Serialize)] pub struct AuditEntry { pub id: i64, - pub ts_unix: DateTime, + #[serde(with = "hive_sh4re::wire_time::iso")] + pub ts_unix: i64, /// Agent on whose behalf the action was taken. pub agent: String, /// What was done (e.g. `restart_infra`). @@ -160,7 +159,7 @@ impl AuditLog { ) { Ok(_) => Some(AuditEntry { id: conn.last_insert_rowid(), - ts_unix: hive_sh4re::wire_time::from_secs(now), + ts_unix: now, agent: agent.to_owned(), action: action.to_owned(), target: target.to_owned(), @@ -253,7 +252,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::from_secs(r.get(1)?), + ts_unix: 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 269d7386..aedf314f 100644 --- a/hive-c0re/src/broker.rs +++ b/hive-c0re/src/broker.rs @@ -7,8 +7,6 @@ use std::sync::Mutex; use std::time::{SystemTime, UNIX_EPOCH}; use anyhow::{Context, Result}; - -use chrono::{DateTime, Utc}; use hive_sh4re::{InboxRow, Message}; use rusqlite::{Connection, OptionalExtension, params}; use serde::Serialize; @@ -76,8 +74,10 @@ pub struct PendingReminder { pub message: String, #[serde(skip_serializing_if = "Option::is_none")] pub file_path: Option, - pub due_at: DateTime, - pub created_at: DateTime, + #[serde(with = "hive_sh4re::wire_time::iso")] + pub due_at: i64, + #[serde(with = "hive_sh4re::wire_time::iso")] + pub created_at: i64, /// 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. @@ -926,8 +926,8 @@ impl Broker { agent: row.get(1)?, message: row.get(2)?, file_path: row.get(3)?, - due_at: hive_sh4re::wire_time::from_secs(row.get(4)?), - created_at: hive_sh4re::wire_time::from_secs(row.get(5)?), + due_at: row.get(4)?, + created_at: 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 b2e50f84..93068813 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::from_secs(resolved_at), + 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::from_secs(asked_at), - deadline_at: deadline_at.map(hive_sh4re::wire_time::from_secs), + asked_at, + deadline_at, 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::from_secs(answered_at), + 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 0f08c33d..5a0399e7 100644 --- a/hive-c0re/src/dashboard.rs +++ b/hive-c0re/src/dashboard.rs @@ -27,7 +27,6 @@ use tokio_stream::{Stream, StreamExt}; use crate::container_view::{ContainerView, claude_has_session}; use crate::coordinator::Coordinator; use crate::lifecycle::{self, MANAGER_NAME}; -use chrono::{DateTime, Utc}; mod approvals; mod build_logs; @@ -434,7 +433,8 @@ struct ApprovalHistoryView { /// `approved` / `denied` / `failed`. status: &'static str, /// RFC 3339 UTC. Renders as a relative time on the dashboard. - resolved_at: DateTime, + #[serde(with = "hive_sh4re::wire_time::iso")] + resolved_at: i64, /// Operator-supplied deny reason (for `denied`) or build error /// (for `failed`). None on `approved`. #[serde(skip_serializing_if = "Option::is_none")] @@ -475,7 +475,8 @@ 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: DateTime, + #[serde(with = "hive_sh4re::wire_time::iso")] + requested_at: i64, } /// Replace silent `.unwrap_or_default()` on the data sources behind @@ -907,7 +908,7 @@ fn history_view(a: Approval) -> ApprovalHistoryView { kind, sha_short, status, - resolved_at: a.resolved_at.unwrap_or_default(), + resolved_at: a.resolved_at.unwrap_or(0), note: a.note, } } @@ -1100,7 +1101,7 @@ async fn dashboard_history(State(state): State) -> Response { from, to, body, - at: hive_sh4re::wire_time::from_secs(at), + at, in_reply_to, file_refs, }) @@ -1120,7 +1121,7 @@ async fn dashboard_history(State(state): State) -> Response { from, to, body, - at: hive_sh4re::wire_time::from_secs(at), + at, in_reply_to, file_refs, }) @@ -1375,7 +1376,7 @@ async fn api_operator_inbox(State(state): State) -> Response { "id": id, "from": from, "body": body, - "at": hive_sh4re::wire_time::from_secs(at), + "at": hive_sh4re::wire_time::to_iso(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 21f5d820..48610473 100644 --- a/hive-c0re/src/dashboard_events.rs +++ b/hive-c0re/src/dashboard_events.rs @@ -9,7 +9,6 @@ use serde::Serialize; use crate::container_view::ContainerView; use crate::dashboard::{MetaInputView, TombstoneView}; use crate::rebuild_queue::QueueEntry; -use chrono::{DateTime, Utc}; #[derive(Debug, Clone, Serialize)] #[serde(rename_all = "snake_case", tag = "kind")] @@ -39,7 +38,8 @@ pub enum DashboardEvent { from: String, to: String, body: String, - at: DateTime, + #[serde(with = "hive_sh4re::wire_time::iso")] + at: i64, #[serde(default, skip_serializing_if = "Option::is_none")] in_reply_to: Option, #[serde(default, skip_serializing_if = "Vec::is_empty")] @@ -54,7 +54,8 @@ pub enum DashboardEvent { from: String, to: String, body: String, - at: DateTime, + #[serde(with = "hive_sh4re::wire_time::iso")] + at: i64, #[serde(default, skip_serializing_if = "Option::is_none")] in_reply_to: Option, #[serde(default, skip_serializing_if = "Vec::is_empty")] @@ -96,7 +97,8 @@ pub enum DashboardEvent { sha_short: Option, /// `"approved"` / `"denied"` / `"failed"`. status: &'static str, - resolved_at: DateTime, + #[serde(with = "hive_sh4re::wire_time::iso")] + resolved_at: i64, note: Option, description: Option, }, @@ -112,8 +114,10 @@ pub enum DashboardEvent { question: String, options: Vec, multi: bool, - asked_at: DateTime, - deadline_at: Option>, + #[serde(with = "hive_sh4re::wire_time::iso")] + asked_at: i64, + #[serde(with = "hive_sh4re::wire_time::iso_opt")] + deadline_at: Option, target: Option, /// Verified file-path tokens that appear in `question`. /// Same shape as broker `Sent`/`Delivered` events; the @@ -131,7 +135,8 @@ pub enum DashboardEvent { id: i64, answer: String, answerer: String, - answered_at: DateTime, + #[serde(with = "hive_sh4re::wire_time::iso")] + answered_at: i64, cancelled: bool, target: Option, /// Verified file-path tokens that appear in `answer`. @@ -332,7 +337,7 @@ mod tests { from: "a".into(), to: "b".into(), body: String::new(), - at: hive_sh4re::wire_time::from_secs(0), + at: 0, in_reply_to: None, file_refs: Vec::new(), }, @@ -342,7 +347,7 @@ mod tests { from: "a".into(), to: "b".into(), body: String::new(), - at: hive_sh4re::wire_time::from_secs(0), + at: 0, in_reply_to: None, file_refs: Vec::new(), }, @@ -363,7 +368,7 @@ mod tests { approval_kind: "apply_commit", sha_short: None, status: "approved", - resolved_at: hive_sh4re::wire_time::from_secs(0), + resolved_at: 0, note: None, description: None, }, @@ -374,7 +379,7 @@ mod tests { question: String::new(), options: Vec::new(), multi: false, - asked_at: hive_sh4re::wire_time::from_secs(0), + asked_at: 0, deadline_at: None, target: None, question_refs: Vec::new(), @@ -384,7 +389,7 @@ mod tests { id: 1, answer: String::new(), answerer: "a".into(), - answered_at: hive_sh4re::wire_time::from_secs(0), + answered_at: 0, cancelled: false, target: None, answer_refs: Vec::new(), @@ -447,7 +452,7 @@ mod tests { seq: 1, entry: crate::audit_log::AuditEntry { id: 1, - ts_unix: hive_sh4re::wire_time::from_secs(0), + ts_unix: 0, agent: "atlas".into(), action: "restart_infra".into(), target: "hive-ci".into(), @@ -476,7 +481,7 @@ mod tests { seq: 7, entry: crate::audit_log::AuditEntry { id: 42, - ts_unix: hive_sh4re::wire_time::from_secs(1_700_000_000), + ts_unix: 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 46f7e597..24eb34c2 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.timestamp()), + age_seconds: saturating_age(now, a.requested_at), }); } 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.timestamp()), + age_seconds: saturating_age(now, q.asked_at), }); } 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.timestamp()), + age_seconds: saturating_age(now, r.created_at), }); } 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.timestamp()), + age_seconds: saturating_age(now, a.requested_at), }); } 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.timestamp()), + age_seconds: saturating_age(now, q.asked_at), }); } 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.timestamp()), + age_seconds: saturating_age(now, r.created_at), }); } Ok(out) diff --git a/hive-c0re/src/main.rs b/hive-c0re/src/main.rs index a7546d63..5e5ac9e2 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::from_secs(at), + 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::from_secs(at), + at, in_reply_to, file_refs, }); diff --git a/hive-c0re/src/operator_questions.rs b/hive-c0re/src/operator_questions.rs index 604c43c8..ab51e9f0 100644 --- a/hive-c0re/src/operator_questions.rs +++ b/hive-c0re/src/operator_questions.rs @@ -14,8 +14,6 @@ use std::sync::Mutex; use std::time::{SystemTime, UNIX_EPOCH}; use anyhow::{Context, Result, bail}; - -use chrono::{DateTime, Utc}; use rusqlite::{Connection, OptionalExtension, params}; use serde::Serialize; @@ -76,12 +74,15 @@ pub struct OpQuestion { pub question: String, pub options: Vec, pub multi: bool, - pub asked_at: DateTime, + #[serde(with = "hive_sh4re::wire_time::iso")] + pub asked_at: i64, /// 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>, + #[serde(with = "hive_sh4re::wire_time::iso_opt")] + pub deadline_at: Option, + #[serde(with = "hive_sh4re::wire_time::iso_opt")] + pub answered_at: Option, pub answer: Option, /// Recipient of the question. `None` = the operator (dashboard /// path); `Some()` = a peer agent asked via @@ -289,14 +290,10 @@ fn row_to_question(row: &rusqlite::Row<'_>) -> rusqlite::Result { question: row.get(2)?, options, multi: multi != 0, - asked_at: hive_sh4re::wire_time::from_secs(row.get(5)?), - answered_at: row - .get::<_, Option>(6)? - .map(hive_sh4re::wire_time::from_secs), + asked_at: row.get(5)?, + answered_at: row.get(6)?, answer: row.get(7)?, - deadline_at: row - .get::<_, Option>(8)? - .map(hive_sh4re::wire_time::from_secs), + deadline_at: row.get(8)?, target: row.get(9)?, }) } diff --git a/hive-c0re/src/socket_server.rs b/hive-c0re/src/socket_server.rs index b4e4e14e..b21f60e3 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::from_secs(s.next_fire_at_unix), - created_at_unix: hive_sh4re::wire_time::from_secs(s.created_at_unix), + next_fire_at_unix: s.next_fire_at_unix, + created_at_unix: 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::from_secs), - paused_at_unix: s.paused_at_unix.map(hive_sh4re::wire_time::from_secs), + cancelled_at_unix: s.cancelled_at_unix, + paused_at_unix: s.paused_at_unix, 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::from_secs), - last_fired_at_unix: t.last_fired_at_unix.map(hive_sh4re::wire_time::from_secs), + cancelled_at_unix: t.cancelled_at_unix, + last_fired_at_unix: t.last_fired_at_unix, 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::from_secs(0), - created_at_unix: hive_sh4re::wire_time::from_secs(0), + next_fire_at_unix: 0, + created_at_unix: 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 d8193b91..09931b0d 100644 --- a/hive-sh4re/src/lib.rs +++ b/hive-sh4re/src/lib.rs @@ -1,6 +1,5 @@ //! Wire types shared between `hive-c0re` and the in-container harness. -use chrono::{DateTime, Utc}; use serde::{Deserialize, Serialize}; pub mod assets; @@ -210,10 +209,15 @@ 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: DateTime, + #[serde(with = "crate::wire_time::iso")] + pub requested_at: i64, pub status: ApprovalStatus, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub resolved_at: Option>, + #[serde( + default, + skip_serializing_if = "Option::is_none", + with = "crate::wire_time::iso_opt" + )] + 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 +453,8 @@ pub enum LooseEnd { id: i64, owner: String, message: String, - due_at: DateTime, + #[serde(with = "crate::wire_time::iso")] + due_at: i64, age_seconds: u64, }, /// Undelivered inbox messages waiting to be `recv`'d by this agent. @@ -1414,17 +1419,27 @@ pub struct WireSchedule { pub body: String, #[serde(default, skip_serializing_if = "Option::is_none")] pub interval_seconds: Option, - pub next_fire_at_unix: DateTime, - pub created_at_unix: DateTime, + #[serde(with = "crate::wire_time::iso")] + pub next_fire_at_unix: i64, + #[serde(with = "crate::wire_time::iso")] + pub created_at_unix: i64, pub source: WireScheduleSource, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub cancelled_at_unix: Option>, + #[serde( + default, + skip_serializing_if = "Option::is_none", + with = "crate::wire_time::iso_opt" + )] + 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>, + #[serde( + default, + skip_serializing_if = "Option::is_none", + with = "crate::wire_time::iso_opt" + )] + pub paused_at_unix: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub description: Option, pub targets: Vec, @@ -1440,10 +1455,18 @@ pub enum WireScheduleSource { #[derive(Debug, Clone, Serialize, Deserialize)] pub struct WireScheduleTarget { pub target: String, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub cancelled_at_unix: Option>, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub last_fired_at_unix: Option>, + #[serde( + default, + skip_serializing_if = "Option::is_none", + with = "crate::wire_time::iso_opt" + )] + pub cancelled_at_unix: Option, + #[serde( + default, + skip_serializing_if = "Option::is_none", + with = "crate::wire_time::iso_opt" + )] + 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 b246b1cf..cbc3dae7 100644 --- a/hive-sh4re/src/wire_time.rs +++ b/hive-sh4re/src/wire_time.rs @@ -1,38 +1,112 @@ -//! 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. +//! 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: `#[serde(with = "crate::wire_time::iso")]` on `i64` fields, +//! `#[serde(with = "crate::wire_time::iso_opt")]` on `Option` +//! (keep the usual `default` + `skip_serializing_if` attributes). -use chrono::{DateTime, Utc}; +use chrono::{DateTime, SecondsFormat, Utc}; +use serde::Deserialize; -/// 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. +/// 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. #[must_use] -pub fn from_secs(secs: i64) -> DateTime { - DateTime::::from_timestamp(secs, 0).unwrap_or_default() +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), + } + } +} + +/// Adaptor for required `i64` timestamp fields. +pub mod iso { + use serde::{Deserializer, Serializer}; + + use super::{Deserialize, EpochOrIso}; + + pub fn serialize(secs: &i64, ser: S) -> Result { + ser.serialize_str(&super::to_iso(*secs)) + } + + pub fn deserialize<'de, D: Deserializer<'de>>(de: D) -> Result { + EpochOrIso::deserialize(de)?.into_secs() + } +} + +/// Adaptor for `Option` timestamp fields. +pub mod iso_opt { + use serde::{Deserializer, Serializer}; + + use super::{Deserialize, EpochOrIso}; + + pub fn serialize(secs: &Option, ser: S) -> Result { + match secs { + Some(secs) => ser.serialize_str(&super::to_iso(*secs)), + None => ser.serialize_none(), + } + } + + pub fn deserialize<'de, D: Deserializer<'de>>(de: D) -> Result, D::Error> { + Option::::deserialize(de)? + .map(EpochOrIso::into_secs) + .transpose() + } } #[cfg(test)] mod tests { - use chrono::{DateTime, Utc}; use serde::{Deserialize, Serialize}; - use super::from_secs; - #[derive(Serialize, Deserialize, PartialEq, Debug)] struct Row { - at: DateTime, - #[serde(default, skip_serializing_if = "Option::is_none")] - maybe_at: Option>, + #[serde(with = "crate::wire_time::iso")] + at: i64, + #[serde( + default, + skip_serializing_if = "Option::is_none", + with = "crate::wire_time::iso_opt" + )] + maybe_at: Option, } #[test] - fn serializes_as_rfc3339_z() { + fn serializes_epoch_as_rfc3339_z() { let json = serde_json::to_string(&Row { - at: from_secs(1_751_480_000), + at: 1_751_480_000, maybe_at: None, }) .unwrap(); @@ -42,8 +116,8 @@ mod tests { #[test] fn round_trips_and_serializes_some() { let row = Row { - at: from_secs(0), - maybe_at: Some(from_secs(1_751_480_000)), + at: 0, + maybe_at: Some(1_751_480_000), }; let json = serde_json::to_string(&row).unwrap(); assert_eq!( @@ -54,13 +128,17 @@ mod tests { } #[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, from_secs(1_751_480_000)); + 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, 1_751_480_000); + assert_eq!(row.maybe_at, Some(1_751_480_000)); } #[test] - fn from_secs_clamps_out_of_range_to_epoch() { - assert_eq!(from_secs(i64::MAX), from_secs(0)); + 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, 1_751_480_000); } }