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}