pub struct Stream { /* private fields */ }Expand description
Priority data stream entity
Implementations§
Source§impl Stream
impl Stream
Sourcepub fn new(
session_id: Id<SessionMarker>,
source_data: JsonData,
config: StreamConfig,
) -> Stream
pub fn new( session_id: Id<SessionMarker>, source_data: JsonData, config: StreamConfig, ) -> Stream
Create new stream
Sourcepub fn id(&self) -> Id<StreamMarker>
pub fn id(&self) -> Id<StreamMarker>
Get stream ID
Sourcepub fn session_id(&self) -> Id<SessionMarker>
pub fn session_id(&self) -> Id<SessionMarker>
Get session ID
Sourcepub fn state(&self) -> &StreamState
pub fn state(&self) -> &StreamState
Get current state
Sourcepub fn config(&self) -> &StreamConfig
pub fn config(&self) -> &StreamConfig
Get configuration
Sourcepub fn stats(&self) -> &StreamStats
pub fn stats(&self) -> &StreamStats
Get statistics
Sourcepub fn created_at(&self) -> DateTime<Utc>
pub fn created_at(&self) -> DateTime<Utc>
Get creation timestamp
Sourcepub fn updated_at(&self) -> DateTime<Utc>
pub fn updated_at(&self) -> DateTime<Utc>
Get last update timestamp
Sourcepub fn completed_at(&self) -> Option<DateTime<Utc>>
pub fn completed_at(&self) -> Option<DateTime<Utc>>
Get completion timestamp
Sourcepub fn source_data(&self) -> Option<&JsonData>
pub fn source_data(&self) -> Option<&JsonData>
Get source data
Sourcepub fn add_metadata(&mut self, key: String, value: String)
pub fn add_metadata(&mut self, key: String, value: String)
Add metadata
Sourcepub fn start_streaming(&mut self) -> Result<(), DomainError>
pub fn start_streaming(&mut self) -> Result<(), DomainError>
Start streaming (transition to Streaming state)
Sourcepub fn complete(&mut self) -> Result<(), DomainError>
pub fn complete(&mut self) -> Result<(), DomainError>
Complete stream successfully
Sourcepub fn cancel(&mut self) -> Result<(), DomainError>
pub fn cancel(&mut self) -> Result<(), DomainError>
Cancel stream
Sourcepub fn create_skeleton_frame(&mut self) -> Result<Frame, DomainError>
pub fn create_skeleton_frame(&mut self) -> Result<Frame, DomainError>
Generate skeleton frame for the stream
Sourcepub fn create_patch_frames(
&mut self,
priority_threshold: Priority,
max_frames: usize,
) -> Result<Vec<Frame>, DomainError>
pub fn create_patch_frames( &mut self, priority_threshold: Priority, max_frames: usize, ) -> Result<Vec<Frame>, DomainError>
Create batch of patch frames based on priority
Sourcepub fn extract_prioritized_patches(
&self,
priority_threshold: Priority,
) -> Result<Vec<(FramePatch, Priority)>, DomainError>
pub fn extract_prioritized_patches( &self, priority_threshold: Priority, ) -> Result<Vec<(FramePatch, Priority)>, DomainError>
Compute prioritized patches for this stream’s current source_data,
without mutating any state — the expensive half of
Self::create_patch_frames (a full traversal of source_data,
computing a priority for every leaf value).
Split out from Self::commit_patch_frames so a caller that needs
its own concurrency control around the mutating half (e.g. a
repository holding a per-session lock for the read-modify-write) can
run this traversal lock-free and only take the lock for the commit
step — which is itself O(patches.len()), not cheap or
max_frames-bounded; see Self::commit_patch_frames’s docs. Safe
to call without holding any lock: source_data is set once at
stream creation and never mutated afterward.
Sourcepub fn commit_patch_frames(
&mut self,
patches: Vec<(FramePatch, Priority)>,
max_frames: usize,
) -> Result<Vec<Frame>, DomainError>
pub fn commit_patch_frames( &mut self, patches: Vec<(FramePatch, Priority)>, max_frames: usize, ) -> Result<Vec<Frame>, DomainError>
Turn already-extracted prioritized patches (from
Self::extract_prioritized_patches) into frames and commit their
bookkeeping (next_sequence, stats) — the other half of
Self::create_patch_frames.
max_frames only bounds the number of frames returned, not the work
done to produce them: every patch in patches is still cloned into
some frame, so this call’s cost is O(patches.len()), proportional to
the total number of patches extracted, not to max_frames.
Re-checks the streaming-state precondition itself: if the stream
transitioned out of Streaming between the caller’s earlier
Self::extract_prioritized_patches call and this one, this fails
cleanly instead of committing frames for a stream that can no longer
accept them.
Not strictly all-or-nothing: chunks are finalized one at a time via
Self::finalize_patch_frame, so a construction failure partway
through would leave earlier chunks already committed. Currently
unobservable — Frame::patch only errors on an empty patch vector,
and chunking never produces an empty chunk here.
Sourcepub fn finalize_patch_frame(
&mut self,
priority: Priority,
frame_patches: Vec<FramePatch>,
) -> Result<Frame, DomainError>
pub fn finalize_patch_frame( &mut self, priority: Priority, frame_patches: Vec<FramePatch>, ) -> Result<Frame, DomainError>
Build and commit one patch frame from an already-chunked group of
patches (see Self::chunk_patches_for_commit).
This is the only place patch-frame next_sequence/stats bookkeeping
happens: assigns the next sequence number, constructs the Frame,
and records it. Callers that decide
which chunked candidates survive before committing (e.g. a
cross-stream priority truncation across multiple streams) should call
this only for the surviving subset, so discarded candidates never
consume a sequence number or inflate stats.
§Examples
use pjson_rs_domain::entities::Stream;
use pjson_rs_domain::entities::frame::FramePatch;
use pjson_rs_domain::value_objects::{JsonData, JsonPath, Priority, SessionId};
let mut stream = Stream::new(
SessionId::new(),
JsonData::Object(Default::default()),
Default::default(),
);
stream.start_streaming().unwrap();
let patch = FramePatch::set(JsonPath::root(), JsonData::Bool(true));
let frame = stream
.finalize_patch_frame(Priority::HIGH, vec![patch])
.unwrap();
assert_eq!(frame.sequence(), 1);
assert_eq!(stream.stats().total_frames, 1);Sourcepub fn create_completion_frame(
&mut self,
checksum: Option<String>,
) -> Result<Frame, DomainError>
pub fn create_completion_frame( &mut self, checksum: Option<String>, ) -> Result<Frame, DomainError>
Create completion frame
Sourcepub fn is_finished(&self) -> bool
pub fn is_finished(&self) -> bool
Check if stream is finished
Sourcepub fn update_config(&mut self, config: StreamConfig) -> Result<(), DomainError>
pub fn update_config(&mut self, config: StreamConfig) -> Result<(), DomainError>
Update configuration
Sourcepub fn chunk_patches_for_commit(
patches: Vec<(FramePatch, Priority)>,
max_frames: usize,
) -> Vec<(Priority, Vec<FramePatch>)>
pub fn chunk_patches_for_commit( patches: Vec<(FramePatch, Priority)>, max_frames: usize, ) -> Vec<(Priority, Vec<FramePatch>)>
Group prioritized patches into per-frame chunks without constructing
any Frame or mutating stream state.
Each returned chunk’s priority is the maximum priority of the patches
it contains, so per-frame ordering downstream reflects the most
important content the frame will carry. Pure and side-effect-free so a
caller (e.g. a cross-stream priority batch) can chunk candidates from
several streams, decide which survive truncation, and only then
commit the survivors via Self::finalize_patch_frame — see that
method’s docs.
§Examples
use pjson_rs_domain::entities::Stream;
use pjson_rs_domain::entities::frame::FramePatch;
use pjson_rs_domain::value_objects::{JsonData, JsonPath, Priority, SessionId};
let mut stream = Stream::new(
SessionId::new(),
JsonData::Object(Default::default()),
Default::default(),
);
stream.start_streaming().unwrap();
let patches = vec![
(
FramePatch::set(JsonPath::root(), JsonData::Bool(true)),
Priority::LOW,
),
(
FramePatch::set(JsonPath::root(), JsonData::Bool(false)),
Priority::HIGH,
),
];
// Chunk without mutating the stream — a caller can inspect priorities
// and pick survivors before any sequence number is spent.
let mut chunks = Stream::chunk_patches_for_commit(patches, 2);
assert_eq!(chunks.len(), 2);
chunks.sort_by_key(|(priority, _)| std::cmp::Reverse(*priority));
// Only commit the highest-priority chunk.
let (priority, frame_patches) = chunks.remove(0);
let frame = stream
.finalize_patch_frame(priority, frame_patches)
.unwrap();
assert_eq!(frame.priority(), Priority::HIGH);
assert_eq!(stream.stats().total_frames, 1);Trait Implementations§
Source§impl<'de> Deserialize<'de> for Stream
impl<'de> Deserialize<'de> for Stream
Source§fn deserialize<__D>(
__deserializer: __D,
) -> Result<Stream, <__D as Deserializer<'de>>::Error>where
__D: Deserializer<'de>,
fn deserialize<__D>(
__deserializer: __D,
) -> Result<Stream, <__D as Deserializer<'de>>::Error>where
__D: Deserializer<'de>,
Source§impl Serialize for Stream
impl Serialize for Stream
Source§fn serialize<__S>(
&self,
__serializer: __S,
) -> Result<<__S as Serializer>::Ok, <__S as Serializer>::Error>where
__S: Serializer,
fn serialize<__S>(
&self,
__serializer: __S,
) -> Result<<__S as Serializer>::Ok, <__S as Serializer>::Error>where
__S: Serializer,
Auto Trait Implementations§
impl Freeze for Stream
impl RefUnwindSafe for Stream
impl Send for Stream
impl Sync for Stream
impl Unpin for Stream
impl UnsafeUnpin for Stream
impl UnwindSafe for Stream
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> DeserializeOwned for Twhere
T: for<'de> Deserialize<'de>,
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more