Handle a RegisterSource request.
(
&mut self,
source: Box<dyn Sourceable<S>>,
scope: &mut S,
)
| 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> { |
no test coverage detected