Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
Show all changes
37 commits
Select commit Hold shift + click to select a range
796def9
Add initial version of Cached
StephenWakely May 25, 2023
b205354
Remove the deadlock from writing to the cache after reading
StephenWakely May 26, 2023
d0c36ed
Insert source and service tag into events
StephenWakely May 26, 2023
ec666b0
Add bytesize count grouped by source and service
StephenWakely May 30, 2023
3e06791
Add EventCountTags trait
StephenWakely May 30, 2023
3a5fe02
Merge remote-tracking branch 'origin' into stephen/cached_events
StephenWakely May 31, 2023
15ea497
Fix merge
StephenWakely May 31, 2023
d66654c
Fix compile errors
StephenWakely May 31, 2023
d1c7df6
Merge remote-tracking branch 'origin' into stephen/cached_events
StephenWakely Jun 1, 2023
c895653
Set source and service tags for most other sinks
StephenWakely Jun 1, 2023
d28a809
Register the events with a trait rather than a Fn
StephenWakely Jun 2, 2023
4ade9be
These tests are round trip
StephenWakely Jun 2, 2023
3632fa2
Merge remote-tracking branch 'origin' into stephen/cached_events
StephenWakely Jun 5, 2023
4d9c5df
Clippy
StephenWakely Jun 5, 2023
c9640d9
Add event count tags for loki sink
StephenWakely Jun 5, 2023
a758704
RegisterEvent inherits form RegisterInternalEvent
StephenWakely Jun 6, 2023
5cbfee4
Merge remote-tracking branch 'origin' into stephen/cached_events
StephenWakely Jun 6, 2023
9eab386
Added telemetry options
StephenWakely Jun 7, 2023
ee50bf5
Only collect configured tags
StephenWakely Jun 8, 2023
9fd7854
Added tests and clippy.
StephenWakely Jun 8, 2023
01d8ad1
Little tidy
StephenWakely Jun 9, 2023
87ca474
TaggedEventsSent doesn't need Output
StephenWakely Jun 9, 2023
ebd10c8
Spelling
StephenWakely Jun 9, 2023
b30a4bb
Remove default impl of take_metadata
StephenWakely Jun 9, 2023
8b74470
Outer event should get tags from inner event type
StephenWakely Jun 9, 2023
b11c7a6
Driver should not consume the metadata
StephenWakely Jun 9, 2023
3327904
Merge remote-tracking branch 'origin' into stephen/cached_events
StephenWakely Jun 9, 2023
7a43045
Tags should be an associated type of RegisterEvent
StephenWakely Jun 9, 2023
f43af1b
Feedback from Bruce
StephenWakely Jun 15, 2023
b141f30
Merge remote-tracking branch 'origin' into stephen/cached_events
StephenWakely Jun 19, 2023
8ac7020
Set source tag to be Arc<ComponentKey>
StephenWakely Jun 19, 2023
243ec8b
Replace take_metadata with metadata_mut
StephenWakely Jun 21, 2023
574cbae
Add test for telemetry tags to Kafka sink
StephenWakely Jun 21, 2023
73131e6
Fix datadog integration test
StephenWakely Jun 21, 2023
4dee27c
Use Derivative to replace clone with no bounds
StephenWakely Jun 23, 2023
f1af663
Spelling
StephenWakely Jun 23, 2023
8f0cc1b
Feedback from Bruce
StephenWakely Jun 26, 2023
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 14 additions & 1 deletion lib/vector-common/src/internal_event/events_sent.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
use metrics::{register_counter, Counter};
use tracing::trace;

use super::{CountByteSize, Output, SharedString};
use super::{CountByteSize, Output, RegisterEvent, SharedString};

pub const DEFAULT_OUTPUT: &str = "_default";

Expand Down Expand Up @@ -96,6 +96,19 @@ crate::registered_event!(
}
);

/// TODO: This can probably become a part of the previous macro.
impl RegisterEvent<(Option<String>, Option<String>), TaggedEventsSent> for TaggedEventsSent {
fn register(
tags: &(Option<String>, Option<String>),
) -> <TaggedEventsSent as super::RegisterInternalEvent>::Handle {
super::register(TaggedEventsSent::new(
tags.0.clone(),
tags.1.clone(),
Output(None),
))
}
}

impl TaggedEventsSent {
#[must_use]
pub fn new(source: Option<String>, service: Option<String>, output: Output) -> Self {
Expand Down
35 changes: 25 additions & 10 deletions lib/vector-common/src/internal_event/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -248,23 +248,38 @@ macro_rules! registered_event {
};
}

#[derive(Clone)]
pub struct Cached<Tags, Event, Register> {
cache: Arc<RwLock<BTreeMap<Tags, Event>>>,
register: Register,
pub trait RegisterEvent<Tags, Event>
where
Event: RegisterInternalEvent,
{
fn register(tags: &Tags) -> <Event as RegisterInternalEvent>::Handle;
}
Comment thread
StephenWakely marked this conversation as resolved.
Outdated

pub struct Cached<Tags, Event: RegisterInternalEvent> {
cache: Arc<RwLock<BTreeMap<Tags, <Event as RegisterInternalEvent>::Handle>>>,
}

/// Deriving `Clone` for `Cached` doesn't work since the `Event` type is not clone,
/// we can happily implement our own `clone` however since we are just cloning
/// the `Arc`.
impl<Tags, Event: RegisterInternalEvent> Clone for Cached<Tags, Event> {
fn clone(&self) -> Self {
Self {
cache: Arc::clone(&self.cache),
}
}
}

impl<Tags, Event, Register, Data> Cached<Tags, Event, Register>
impl<Tags, Event, EventHandle, Data> Cached<Tags, Event>
where
Data: Sized,
Register: Fn(&Tags) -> Event,
Event: InternalEventHandle<Data = Data>,
EventHandle: InternalEventHandle<Data = Data>,
Tags: Ord + Clone,
Event: RegisterInternalEvent<Handle = EventHandle> + RegisterEvent<Tags, Event>,
{
pub fn new(register: Register) -> Self {
pub fn new() -> Self {
Self {
cache: Arc::new(RwLock::new(BTreeMap::new())),
register,
}
}

Expand All @@ -273,7 +288,7 @@ where
if let Some(event) = read.get(tags) {
event.emit(value);
} else {
let event = (self.register)(tags);
let event = <Event as RegisterEvent<Tags, Event>>::register(tags);
event.emit(value);

// Ensure the read lock is dropped so we can write.
Expand Down
20 changes: 6 additions & 14 deletions lib/vector-core/src/stream/driver.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,8 @@ use tokio::{pin, select};
use tower::Service;
use tracing::Instrument;
use vector_common::internal_event::{
register, ByteSize, BytesSent, Cached, CallError, InternalEventHandle as _, Output,
PollReadyError, Registered, SharedString, TaggedEventsSent,
register, ByteSize, BytesSent, Cached, CallError, InternalEventHandle as _, PollReadyError,
Registered, SharedString, TaggedEventsSent,
};
use vector_common::request_metadata::{MetaDescriptive, RequestCountByteSize, RequestMetadata};

Expand Down Expand Up @@ -99,13 +99,7 @@ where
pin!(batched_input);

let bytes_sent = protocol.map(|protocol| register(BytesSent { protocol }));
let events_sent = Cached::new(|tags: &(Option<String>, Option<String>)| {
register(TaggedEventsSent::new(
tags.0.clone(),
tags.1.clone(),
Output(None),
))
});
let events_sent = Cached::new();

loop {
// Core behavior of the loop:
Expand Down Expand Up @@ -205,16 +199,14 @@ where
Ok(())
}

fn handle_response<T>(
fn handle_response(
result: Result<Svc::Response, Svc::Error>,
request_id: usize,
finalizers: EventFinalizers,
metadata: &RequestMetadata,
bytes_sent: &Option<Registered<BytesSent>>,
events_sent: &Cached<(Option<String>, Option<String>), Registered<TaggedEventsSent>, T>,
) where
T: Fn(&(Option<String>, Option<String>)) -> Registered<TaggedEventsSent>,
{
events_sent: &Cached<(Option<String>, Option<String>), TaggedEventsSent>,
) {
match result {
Err(error) => {
Self::emit_call_error(Some(error), request_id, metadata.event_count());
Expand Down