MCPcopy Create free account
hub / github.com/comnik/declarative-dataflow / register_source

Method register_source

src/server/mod.rs:433–479  ·  view source on GitHub ↗

Handle a RegisterSource request.

(
        &mut self,
        source: Box<dyn Sourceable<S>>,
        scope: &mut S,
    )

Source from the content-addressed store, hash-verified

431
432 /// Handle a RegisterSource request.
433 pub fn register_source<S: Scope<Timestamp = T>>(
434 &mut self,
435 source: Box<dyn Sourceable<S>>,
436 scope: &mut S,
437 ) -> Result<(), Error> {
438 // use timely::logging::Logger;
439 // let timely_logger = scope.log_register().remove("timely");
440
441 // let differential_logger = scope.log_register().remove("differential/arrange");
442
443 let context = SourcingContext {
444 t0: self.t0,
445 scheduler: Rc::downgrade(&self.scheduler),
446 domain_probe: self.context.internal.domain_probe().clone(),
447 timely_events: self.timely_events.clone().unwrap(),
448 differential_events: self.differential_events.clone().unwrap(),
449 };
450
451 // self.timely_events = None;
452 // self.differential_events = None;
453
454 let mut attribute_streams = source.source(scope, context);
455
456 for (aid, config, datoms) in attribute_streams.drain(..) {
457 self.context
458 .internal
459 .create_sourced_attribute(&aid, config, &datoms)?;
460 }
461
462 // if let Some(logger) = timely_logger {
463 // if let Ok(logger) = logger.downcast::<Logger<TimelyEvent>>() {
464 // scope
465 // .log_register()
466 // .insert_logger::<TimelyEvent>("timely", *logger);
467 // }
468 // }
469
470 // if let Some(logger) = differential_logger {
471 // if let Ok(logger) = logger.downcast::<Logger<DifferentialEvent>>() {
472 // scope
473 // .log_register()
474 // .insert_logger::<DifferentialEvent>("differential/arrange", *logger);
475 // }
476 // }
477
478 Ok(())
479 }
480
481 /// Handle an AdvanceDomain request.
482 pub fn advance_domain(&mut self, name: Option<String>, next: T) -> Result<(), Error> {

Callers 4

mainFunction · 0.80
mainFunction · 0.80
mainFunction · 0.80
mainFunction · 0.80

Calls 3

domain_probeMethod · 0.80
sourceMethod · 0.45

Tested by

no test coverage detected