487
}
488
489
>
namespaceEntry, err := e.namespaceRegistry.GetNamespaceByID(namespace.ID(partition.NamespaceId()))
matching_engine.go
490
>
if err != nil {
491
return nil, false, err
492
}
493
495
>
tqConfig := newTaskQueueConfig(partition.TaskQueue(), e.config, namespaceEntry.Name())
496
>
tqConfig.loadCause = loadCause
497
>
logger, throttledLogger, metricsHandler := e.loggerAndMetricsForPartition(namespaceEntry, partition, tqConfig)
498
>
onFatalErr := func(cause unloadCause) { newPM.unloadFromEngine(cause) }
499
>
onUserDataChanged := func(to *persistencespb.VersionedTaskQueueUserData) { newPM.userDataChanged(to) }
500
>
onEphemeralDataChanged := func(data *taskqueuespb.EphemeralData) { newPM.ephemeralDataChanged(data) }
501
>
userDataManager := newUserDataManager(
502
>
e.taskManager,
503
>
e.matchingRawClient,
504
>
onFatalErr,
505
>
onUserDataChanged,
506
>
onEphemeralDataChanged,
507
>
partition,
508
>
tqConfig,
509
>
logger,
510
>
e.namespaceRegistry,
511
>
)
512
>
newPM, err = newTaskQueuePartitionManager(
513
>
e,
514
>
namespaceEntry,
515
>
partition,
516
>
tqConfig,
517
>
logger,
518
>
throttledLogger,
519
>
metricsHandler,
520
>
userDataManager,
521
>
)
522
>
if err != nil {
523
return nil, false, err
524
}
525
526
// If it gets here, write lock and check again in case a task queue is created between the two locks
528
>
pm, ok = e.partitions[key]
529
>
if ok {
530
e.partitionsLock.Unlock()
531
// Lost the race with a concurrent load of the same partition. The unstarted