// ============================================================================= // Active-GO Framework // Copyright (c) 2025 ActiveThing (https://activething.com) // Author: Juan V. Navarro juanvnl@activething.com // ============================================================================= // // Permission is hereby granted, free of charge, to any person obtaining a copy // of this software and associated documentation files (the "Software"), to deal // in the Software without restriction, including without limitation the rights // to use, copy, modify, merge, publish, distribute, sublicense, and/or sell // copies of the Software, and to permit persons to whom the Software is // furnished to do so, subject to the following conditions: // // The above copyright notice and this permission notice shall be included in // all copies or substantial portions of the Software. // // THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR // IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, // FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE // AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER // LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, // OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN // THE SOFTWARE. // // ============================================================================= // Date Create 19/05/2025 // ============================================================================= package logs import ( "context" "core/errs" "fmt" "reflect" "sync" ) type ( Dispatcher struct { mu sync.RWMutex processors map[string]*ProcessorEntry formatters map[string]errs.Formatter bufferSize int cacheMu sync.Mutex cacheEnabled map[string]struct{} formatCache map[string]map[uintptr][]byte inputChan chan error ctx context.Context cancel context.CancelFunc wg sync.WaitGroup } DispatcherFnc func (dispatcher *Dispatcher) error ) func NewDispatcher(options ...DispatcherFnc) (*Dispatcher,error) { dp := &Dispatcher{ processors : make(map[string]*ProcessorEntry), formatters : make(map[string]errs.Formatter), formatCache: make(map[string]map[uintptr][]byte), cacheEnabled: make(map[string]struct{}), inputChan : make(chan error,100), } for _,o := range options { if e := o(dp); e != nil { return nil,e } } return dp,nil } func (d *Dispatcher) Start () { d.mu.Lock() defer d.mu.Unlock() if d.ctx != nil { //todo //panic dispatcher already started } d.ctx,d.cancel = context.WithCancel(context.Background()) d.wg.Add(1) go d.run() for _,entry := range d.processors { d.wg.Add(1) go d.processErrors(d.ctx, entry) } } func (d *Dispatcher) RegisterProcessor(processor Processor, options ...ProcessorEntryFnc) error { d.mu.Lock() defer d.mu.Unlock() if processor == nil { return errs.CreateParamNilErr("processor") } if _,ex := d.formatters[processor.FormatterID()]; !ex { // //todo create newerror formatter_not_registered } if _,ex := d.processors[processor.ID()]; ex { // //todo create newerror processor_registered } pe,er := NewProcessorEntry(options...) if er != nil { return er } d.processors[processor.ID()] = pe return nil } func (d *Dispatcher) processErrors (entry *ProcessorEntry) { defer d.wg.Done() for { select { case <- d.ctx.Done(): // todo existing // case er,ok := <-entry.channel: if !ok { // channel closed return } fi := entry.Processor.FormatterID() } } } // getFormattedError retrieves the formatted error from the cache or formats it if not found. // It uses the dispatcher's formatters and cache. // This function is called by the processor's dedicated goroutine. func (d *Dispatcher) getFormattedError(err error, formatterID string) ([]byte, error) { // Use the error's underlying value's pointer address as a unique identifier for this specific error instance. // This is the correct and idiomatic way to get a stable ID for caching error instances. errID := reflect.ValueOf(err).Pointer() // Using reflect.ValueOf().Pointer() for errID // Check if caching is enabled for this formatter ID. d.mu.RLock() // Use RLock to check cacheEnabled map _, cachingEnabled := d.cacheEnabled[formatterID] d.mu.RUnlock() // Release RLock if cachingEnabled { d.cacheMu.Lock() // Acquire cache lock to read from cache // Check if the formatted error is already in the cache for this formatter ID and error. if formatterCache, ok := d.formatCache[formatterID]; ok { if formatted, ok := formatterCache[errID]; ok { d.cacheMu.Unlock() // Release cache lock before returning return formatted, nil // Found in cache } } d.cacheMu.Unlock() // Release cache lock before accessing formatters (which require d.mu) } // Formatted error not in cache or caching is disabled, find the formatter and format it. d.mu.RLock() // Need RLock to access the formatters map formatter, exists := d.formatters[formatterID] d.mu.RUnlock() // Release RLock if !exists { // This should ideally not happen if RegisterProcessor checks were sufficient, // but handle defensively. return nil, fmt.Errorf("formatter with ID '%s' not found", formatterID) } // Format the error using the found formatter. formatted, formatErr := formatter.Format(err) if formatErr != nil { return nil, fmt.Errorf("failed to format error using formatter '%s': %w", formatterID, formatErr) } // Store the formatted error in the cache only if caching is enabled for this formatter. if cachingEnabled { d.cacheMu.Lock() // Acquire cache lock to write to cache if _, ok := d.formatCache[formatterID]; !ok { d.formatCache[formatterID] = make(map[uintptr][]byte) } d.formatCache[formatterID][errID] = formatted d.cacheMu.Unlock() // Release cache lock after writing to cache } return formatted, nil }