Skip to main content

stygian_graph/ports/
data_sink.rs

1//! DataSink port — outbound counterpart to [`DataSourcePort`](crate::ports::data_source::DataSourcePort).
2//!
3//! [`DataSinkPort`](crate::ports::data_sink::DataSinkPort) is the abstraction that lets pipeline nodes publish scraped
4//! records to an external system without being coupled to any particular backend
5//! (file system, webhook endpoint, message queue, Scrape Exchange, etc.).
6//!
7//! # Architecture
8//!
9//! Following the hexagonal architecture model:
10//!
11//! - This file lives in the **ports** layer — pure trait definitions, no I/O.
12//! - Concrete adapters (file sink, HTTP sink, …) implement this trait and live
13//!   under `adapters/`.
14//!
15//! # Example
16//!
17//! ```rust
18//! use stygian_graph::ports::data_sink::{DataSinkPort, SinkRecord};
19//!
20//! // Any adapter that implements DataSinkPort can be used here.
21//! async fn publish_one(sink: &dyn DataSinkPort, payload: serde_json::Value) {
22//!     let record = SinkRecord::new("my-schema", "https://example.com", payload);
23//!     match sink.publish(&record).await {
24//!         Ok(receipt) => println!("Published: {}", receipt.id),
25//!         Err(e) => eprintln!("Publish failed: {e}"),
26//!     }
27//! }
28//! ```
29
30use std::collections::HashMap;
31
32use async_trait::async_trait;
33use serde::{Deserialize, Serialize};
34use thiserror::Error;
35
36// ── Error type ────────────────────────────────────────────────────────────────
37
38/// Errors that a [`DataSinkPort`] implementation may return.
39#[derive(Debug, Error)]
40#[non_exhaustive]
41pub enum DataSinkError {
42    /// The record failed structural or semantic validation before being sent.
43    #[error("validation failed: {0}")]
44    ValidationFailed(String),
45
46    /// The underlying transport or API rejected the publish request.
47    #[error("publish failed: {0}")]
48    PublishFailed(String),
49
50    /// The sink is temporarily rate-limited; caller should back off.
51    #[error("rate limited: {0}")]
52    RateLimited(String),
53
54    /// Authentication or authorisation rejected the request.
55    #[error("unauthorized: {0}")]
56    Unauthorized(String),
57
58    /// The referenced schema identifier is not known to this sink.
59    #[error("schema not found: {0}")]
60    SchemaNotFound(String),
61}
62
63// ── Domain types ──────────────────────────────────────────────────────────────
64
65/// A single structured record to be published through a [`DataSinkPort`].
66///
67/// # Example
68///
69/// ```rust
70/// use stygian_graph::ports::data_sink::SinkRecord;
71/// use serde_json::json;
72///
73/// let record = SinkRecord::new(
74///     "product-v1",
75///     "https://shop.example.com/items/42",
76///     json!({ "sku": "ABC-42", "price": 9.99 }),
77/// );
78/// assert_eq!(record.schema_id, "product-v1");
79/// ```
80#[derive(Debug, Clone, Serialize, Deserialize)]
81pub struct SinkRecord {
82    /// The payload to publish. Any JSON value is accepted.
83    pub data: serde_json::Value,
84
85    /// Identifies the schema or data-contract version this record conforms to.
86    /// Sinks may use this for routing, validation, or schema-registry lookups.
87    pub schema_id: String,
88
89    /// The canonical URL the record was scraped from. Used for provenance and
90    /// deduplication. Stored as a `String` to avoid a `url` crate dependency
91    /// in the port layer.
92    pub source_url: String,
93
94    /// **T108 mandatory** — the wall-clock instant at which the source
95    /// was fetched, as observed by the transport (HTTP `Date` header
96    /// for HTTP-sourced records, browser `Date.now()` for
97    /// browser-sourced records).
98    ///
99    /// This field is required. Adapters and consumers cannot construct
100    /// a [`SinkRecord`] without supplying it — use
101    /// [`SinkRecord::with_fetched_at`] or
102    /// [`SinkRecord::fetched_at_or_default`] to provide it.
103    ///
104    /// The `fetched_at` value should be the *transport-supplied*
105    /// timestamp, not `Utc::now()` computed post-extraction. The
106    /// guard test in `tests/sink_invariants.rs` asserts every
107    /// adapter uses a transport-level source.
108    pub fetched_at: chrono::DateTime<chrono::Utc>,
109
110    /// Arbitrary string key-value metadata (content-type, run-id, tenant, …).
111    pub metadata: HashMap<String, String>,
112}
113
114impl SinkRecord {
115    /// Construct a new [`SinkRecord`] with empty metadata and
116    /// `fetched_at = chrono::Utc::now()` as the **fallback** for
117    /// callers that don't have a transport-level timestamp.
118    ///
119    /// **Prefer [`SinkRecord::with_fetched_at`]** — using
120    /// `Utc::now()` here means the record carries the time of
121    /// construction, not the time of the upstream fetch.
122    #[must_use]
123    pub fn new(
124        schema_id: impl Into<String>,
125        source_url: impl Into<String>,
126        data: serde_json::Value,
127    ) -> Self {
128        Self::fetched_at_or_default(schema_id, source_url, data, chrono::Utc::now())
129    }
130
131    /// Construct a [`SinkRecord`] with an explicit `fetched_at`
132    /// timestamp. This is the **preferred constructor** — the
133    /// caller is required to supply a transport-level timestamp.
134    ///
135    /// # Example
136    ///
137    /// ```rust
138    /// use stygian_graph::ports::data_sink::SinkRecord;
139    /// use chrono::{TimeZone, Utc};
140    ///
141    /// let fetched_at = Utc.with_ymd_and_hms(2026, 8, 22, 12, 0, 0).unwrap();
142    /// let r = SinkRecord::with_fetched_at(
143    ///     "schema-v1",
144    ///     "https://example.com/page",
145    ///     serde_json::Value::Null,
146    ///     fetched_at,
147    /// );
148    /// assert_eq!(r.fetched_at, fetched_at);
149    /// ```
150    #[must_use]
151    pub fn with_fetched_at(
152        schema_id: impl Into<String>,
153        source_url: impl Into<String>,
154        data: serde_json::Value,
155        fetched_at: chrono::DateTime<chrono::Utc>,
156    ) -> Self {
157        Self {
158            data,
159            schema_id: schema_id.into(),
160            source_url: source_url.into(),
161            fetched_at,
162            metadata: HashMap::new(),
163        }
164    }
165
166    /// Construct a [`SinkRecord`] using `default_fetched_at` only
167    /// when the caller does not have a transport-level timestamp.
168    /// The `default_fetched_at` value is recorded verbatim — use
169    /// [`chrono::Utc::now`] for a "now" fallback, or pass an HTTP
170    /// `Date` header value parsed into a [`chrono::DateTime<Utc>`].
171    #[must_use]
172    pub fn fetched_at_or_default(
173        schema_id: impl Into<String>,
174        source_url: impl Into<String>,
175        data: serde_json::Value,
176        default_fetched_at: chrono::DateTime<chrono::Utc>,
177    ) -> Self {
178        Self::with_fetched_at(schema_id, source_url, data, default_fetched_at)
179    }
180
181    /// Attach a metadata entry and return `self` for builder-style use.
182    ///
183    /// # Example
184    ///
185    /// ```rust
186    /// use stygian_graph::ports::data_sink::SinkRecord;
187    /// use chrono::Utc;
188    ///
189    /// let r = SinkRecord::new("s", "https://x.com", serde_json::Value::Null)
190    ///     .with_meta("run_id", "abc123");
191    /// assert_eq!(r.metadata["run_id"], "abc123");
192    /// ```
193    #[must_use]
194    pub fn with_meta(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
195        self.metadata.insert(key.into(), value.into());
196        self
197    }
198}
199
200/// Confirmation that a [`SinkRecord`] was successfully accepted by the sink.
201///
202/// # Example
203///
204/// ```rust
205/// use stygian_graph::ports::data_sink::SinkReceipt;
206///
207/// let receipt = SinkReceipt {
208///     id: "rec-001".to_string(),
209///     published_at: "2026-04-09T00:00:00Z".to_string(),
210///     platform: "file-sink".to_string(),
211/// };
212/// assert_eq!(receipt.platform, "file-sink");
213/// ```
214#[derive(Debug, Clone, Serialize, Deserialize)]
215pub struct SinkReceipt {
216    /// Platform-assigned identifier for this published record.
217    pub id: String,
218
219    /// ISO 8601 timestamp at which the sink accepted the record.
220    pub published_at: String,
221
222    /// Human-readable name of the sink platform (e.g. `"scrape-exchange"`, `"file"`).
223    pub platform: String,
224}
225
226// ── Port trait ────────────────────────────────────────────────────────────────
227
228/// Outbound data sink port — publish scraped records to an external system.
229///
230/// Implementations live in `adapters/` and are never imported by domain code.
231/// The port is always injected via `Arc<dyn DataSinkPort>`.
232///
233/// # Object safety
234///
235/// Native `async fn` in traits is not object-safe by itself. This trait uses
236/// `#[async_trait]`, which erases async methods into boxed futures and enables
237/// usage as `dyn DataSinkPort` through `Arc` in this workspace.
238///
239/// # Example
240///
241/// ```rust
242/// use stygian_graph::ports::data_sink::{DataSinkPort, SinkRecord, SinkReceipt, DataSinkError};
243///
244/// struct NoopSink;
245///
246/// #[async_trait::async_trait]
247/// impl DataSinkPort for NoopSink {
248///     async fn publish(&self, _record: &SinkRecord) -> Result<SinkReceipt, DataSinkError> {
249///         Ok(SinkReceipt {
250///             id: "noop".to_string(),
251///             published_at: "".to_string(),
252///             platform: "noop".to_string(),
253///         })
254///     }
255///
256///     async fn validate(&self, _record: &SinkRecord) -> Result<(), DataSinkError> {
257///         Ok(())
258///     }
259///
260///     async fn health_check(&self) -> Result<(), DataSinkError> {
261///         Ok(())
262///     }
263/// }
264/// ```
265#[async_trait]
266pub trait DataSinkPort: Send + Sync {
267    /// Validate and publish `record` to the sink.
268    ///
269    /// Implementations should validate the record before publishing; failing
270    /// fast with [`DataSinkError::ValidationFailed`] is preferred over sending
271    /// invalid data downstream.
272    ///
273    /// # Errors
274    ///
275    /// Returns [`DataSinkError`] on validation failure, transport error, or
276    /// rate-limit/auth rejection.
277    async fn publish(&self, record: &SinkRecord) -> Result<SinkReceipt, DataSinkError>;
278
279    /// Validate `record` without publishing it.
280    ///
281    /// Useful for preflight checks without side effects.
282    ///
283    /// # Errors
284    ///
285    /// Returns [`DataSinkError::ValidationFailed`] if the record is malformed
286    /// or violates schema constraints.
287    async fn validate(&self, record: &SinkRecord) -> Result<(), DataSinkError>;
288
289    /// Check that the sink backend is reachable and healthy.
290    ///
291    /// # Errors
292    ///
293    /// Returns [`DataSinkError::PublishFailed`] or [`DataSinkError::Unauthorized`]
294    /// if the backend is unreachable or misconfigured.
295    async fn health_check(&self) -> Result<(), DataSinkError>;
296}
297
298// ── Tests ─────────────────────────────────────────────────────────────────────
299
300#[cfg(test)]
301mod tests {
302    use super::*;
303    use serde_json::{Value, json};
304
305    #[test]
306    fn sink_record_construction_and_serde_roundtrip()
307    -> std::result::Result<(), Box<dyn std::error::Error>> {
308        let record = SinkRecord::new(
309            "product-v1",
310            "https://shop.example.com/items/42",
311            json!({ "sku": "ABC-42", "price": 9.99 }),
312        )
313        .with_meta("run_id", "abc123")
314        .with_meta("tenant", "acme");
315
316        assert_eq!(record.schema_id, "product-v1");
317        assert_eq!(record.source_url, "https://shop.example.com/items/42");
318        assert_eq!(
319            record.data.get("sku").and_then(Value::as_str),
320            Some("ABC-42")
321        );
322        assert_eq!(
323            record.metadata.get("run_id").map(String::as_str),
324            Some("abc123")
325        );
326        assert_eq!(
327            record.metadata.get("tenant").map(String::as_str),
328            Some("acme")
329        );
330
331        // Round-trip through JSON
332        let json_str = serde_json::to_string(&record)?;
333        let restored: SinkRecord = serde_json::from_str(&json_str)?;
334
335        assert_eq!(restored.schema_id, record.schema_id);
336        assert_eq!(restored.source_url, record.source_url);
337        assert_eq!(
338            restored.metadata.get("run_id").map(String::as_str),
339            Some("abc123")
340        );
341        Ok(())
342    }
343
344    #[test]
345    fn sink_receipt_serde_roundtrip() -> std::result::Result<(), Box<dyn std::error::Error>> {
346        let receipt = SinkReceipt {
347            id: "rec-001".to_string(),
348            published_at: "2026-04-09T00:00:00Z".to_string(),
349            platform: "test-sink".to_string(),
350        };
351
352        let json_str = serde_json::to_string(&receipt)?;
353        let restored: SinkReceipt = serde_json::from_str(&json_str)?;
354
355        assert_eq!(restored.id, receipt.id);
356        assert_eq!(restored.platform, receipt.platform);
357        Ok(())
358    }
359
360    #[test]
361    fn data_sink_error_display() {
362        assert_eq!(
363            DataSinkError::ValidationFailed("missing field".to_string()).to_string(),
364            "validation failed: missing field"
365        );
366        assert_eq!(
367            DataSinkError::PublishFailed("timeout".to_string()).to_string(),
368            "publish failed: timeout"
369        );
370        assert_eq!(
371            DataSinkError::RateLimited("429".to_string()).to_string(),
372            "rate limited: 429"
373        );
374        assert_eq!(
375            DataSinkError::Unauthorized("401".to_string()).to_string(),
376            "unauthorized: 401"
377        );
378        assert_eq!(
379            DataSinkError::SchemaNotFound("v99".to_string()).to_string(),
380            "schema not found: v99"
381        );
382    }
383
384    // ── T108 mandatory fetched_at tests ──────────────────────────────
385
386    #[test]
387    fn fetched_at_required_field_compiles_with_new_constructor() {
388        // The new constructor carries the explicit fetched_at
389        // timestamp — this is the happy path adapters must take.
390        use chrono::TimeZone;
391        let fetched_at = chrono::Utc
392            .with_ymd_and_hms(2026, 8, 22, 12, 0, 0)
393            .single()
394            .unwrap_or_else(chrono::Utc::now);
395        let r = SinkRecord::with_fetched_at(
396            "schema-v1",
397            "https://example.com",
398            json!({ "sku": "ABC-42" }),
399            fetched_at,
400        );
401        assert_eq!(r.fetched_at, fetched_at);
402    }
403
404    #[test]
405    fn fetched_at_or_default_records_verbatim() {
406        use chrono::TimeZone;
407        let supplied = chrono::Utc
408            .with_ymd_and_hms(2026, 1, 1, 0, 0, 0)
409            .single()
410            .unwrap_or_else(chrono::Utc::now);
411        let r = SinkRecord::fetched_at_or_default(
412            "schema-v1",
413            "https://example.com",
414            json!({}),
415            supplied,
416        );
417        assert_eq!(r.fetched_at, supplied);
418    }
419
420    #[test]
421    fn fetched_at_default_constructor_falls_back_to_now() {
422        use chrono::TimeZone;
423        let before = chrono::Utc
424            .with_ymd_and_hms(1970, 1, 1, 0, 0, 0)
425            .single()
426            .unwrap_or_else(chrono::Utc::now);
427        let r = SinkRecord::new("schema-v1", "https://example.com", json!({}));
428        let after = chrono::Utc::now() + chrono::Duration::seconds(1);
429        assert!(
430            r.fetched_at >= before && r.fetched_at <= after,
431            "fetched_at must fall within [before, after] window"
432        );
433    }
434
435    #[test]
436    fn fetched_at_round_trips_through_json() -> std::result::Result<(), Box<dyn std::error::Error>>
437    {
438        use chrono::TimeZone;
439        let fetched_at = chrono::Utc
440            .with_ymd_and_hms(2026, 8, 22, 12, 0, 0)
441            .single()
442            .unwrap_or_else(chrono::Utc::now);
443        let record = SinkRecord::with_fetched_at(
444            "schema-v1",
445            "https://example.com",
446            json!({ "x": 1 }),
447            fetched_at,
448        );
449        let json_str = serde_json::to_string(&record)?;
450        let restored: SinkRecord = serde_json::from_str(&json_str)?;
451        assert_eq!(restored.fetched_at, record.fetched_at);
452        Ok(())
453    }
454
455    #[test]
456    fn with_meta_is_chainable_after_new_constructor() {
457        let r = SinkRecord::new("s", "https://x.com", json!({}))
458            .with_meta("run_id", "abc123")
459            .with_meta("tenant", "acme");
460        assert_eq!(r.metadata.get("run_id").map(String::as_str), Some("abc123"));
461        assert_eq!(r.metadata.get("tenant").map(String::as_str), Some("acme"));
462    }
463}