shared project
This commit is contained in:
		
						commit
						24aefd98aa
					
				
							
								
								
									
										61
									
								
								.gitignore
									
									
									
									
										vendored
									
									
										Normal file
									
								
							
							
						
						
									
										61
									
								
								.gitignore
									
									
									
									
										vendored
									
									
										Normal file
									
								
							| @ -0,0 +1,61 @@ | ||||
| # Created by .ignore support plugin (hsz.mobi) | ||||
| ### JetBrains template | ||||
| # Covers JetBrains IDEs: IntelliJ, RubyMine, PhpStorm, AppCode, PyCharm, CLion, Android Studio and Webstorm | ||||
| # Reference: https://intellij-support.jetbrains.com/hc/en-us/articles/206544839 | ||||
| 
 | ||||
| # User-specific stuff: | ||||
| .idea/**/workspace.xml | ||||
| .idea/**/tasks.xml | ||||
| .idea/dictionaries | ||||
| 
 | ||||
| # Sensitive or high-churn files: | ||||
| .idea/**/dataSources/ | ||||
| .idea/**/dataSources.ids | ||||
| .idea/**/dataSources.xml | ||||
| .idea/**/dataSources.local.xml | ||||
| .idea/**/sqlDataSources.xml | ||||
| .idea/**/dynamic.xml | ||||
| .idea/**/uiDesigner.xml | ||||
| 
 | ||||
| # Gradle: | ||||
| .idea/**/gradle.xml | ||||
| .idea/**/libraries | ||||
| 
 | ||||
| # Mongo Explorer plugin: | ||||
| .idea/**/mongoSettings.xml | ||||
| 
 | ||||
| ## File-based project format: | ||||
| *.iws | ||||
| 
 | ||||
| ## Plugin-specific files: | ||||
| 
 | ||||
| # IntelliJ | ||||
| /out/ | ||||
| 
 | ||||
| # mpeltonen/sbt-idea plugin | ||||
| .idea_modules/ | ||||
| 
 | ||||
| # JIRA plugin | ||||
| atlassian-ide-plugin.xml | ||||
| 
 | ||||
| # Crashlytics plugin (for Android Studio and IntelliJ) | ||||
| com_crashlytics_export_strings.xml | ||||
| crashlytics.properties | ||||
| crashlytics-build.properties | ||||
| fabric.properties | ||||
| ### Go template | ||||
| # Binaries for programs and plugins | ||||
| *.exe | ||||
| *.dll | ||||
| *.so | ||||
| *.dylib | ||||
| 
 | ||||
| # Test binary, build with `go test -c` | ||||
| *.test | ||||
| 
 | ||||
| # Output of the go coverage tool, specifically when used with LiteIDE | ||||
| *.out | ||||
| 
 | ||||
| # Project-local glide cache, RE: https://github.com/Masterminds/glide/issues/736 | ||||
| .glide/ | ||||
| 
 | ||||
							
								
								
									
										111
									
								
								queue.go
									
									
									
									
									
										Normal file
									
								
							
							
						
						
									
										111
									
								
								queue.go
									
									
									
									
									
										Normal file
									
								
							| @ -0,0 +1,111 @@ | ||||
| package queue | ||||
| 
 | ||||
| import ( | ||||
| 	"container/heap" | ||||
| 	"time" | ||||
| 	"loafle.com/overflow/agent_api/observer" | ||||
| 	"github.com/prometheus/common/log" | ||||
| ) | ||||
| 
 | ||||
| type EventType int | ||||
| 
 | ||||
| const ( | ||||
| 
 | ||||
| 	DATA_TYPE EventType  = 0 | ||||
| 	EVENT_TYPE EventType = 1 | ||||
| 
 | ||||
| ) | ||||
| type item struct { | ||||
| 	value interface{} | ||||
| 	priority int | ||||
| } | ||||
| 
 | ||||
| type LoafleQueue struct { | ||||
| 	items []*item | ||||
| 	queueType string | ||||
| 	interval time.Duration | ||||
| 
 | ||||
| 	eventChanel chan interface{} | ||||
| } | ||||
| 
 | ||||
| func (lq LoafleQueue) Len() int { | ||||
| 	return len(lq.items) | ||||
| } | ||||
| 
 | ||||
| func (lq LoafleQueue) Less(i, j int) bool { | ||||
| 
 | ||||
| 	return lq.items[i].priority < lq.items[j].priority | ||||
| } | ||||
| 
 | ||||
| func (lq LoafleQueue) Swap(i, j int) { | ||||
| 	lq.items[i], lq.items[j] = lq.items[j], lq.items[i] | ||||
| } | ||||
| 
 | ||||
| func (lq *LoafleQueue) GetItems() []*item { | ||||
| 
 | ||||
| 	time.Sleep(time.Second * lq.interval) | ||||
| 	resultItems := make([]*item, 0) | ||||
| 	for lq.Len() > 0 { | ||||
| 		item := heap.Pop(lq).(*item) | ||||
| 		resultItems = append(resultItems, item) | ||||
| 	} | ||||
| 	return resultItems | ||||
| } | ||||
| 
 | ||||
| func (lq *LoafleQueue) Push(i interface{}) { | ||||
| 	n := len(lq.items) | ||||
| 	nItem := i.(*item) | ||||
| 	nItem.priority = n | ||||
| 	lq.items = append(lq.items, nItem) | ||||
| } | ||||
| 
 | ||||
| func (lq *LoafleQueue) Pop() interface{} { | ||||
| 	old := lq.items | ||||
| 	n := len(old) | ||||
| 	nItem := old[n-1] | ||||
| 	lq.items = old[0 : n-1] | ||||
| 	return nItem | ||||
| } | ||||
| func (lq LoafleQueue) newItem(value interface{}) *item { | ||||
| 	return &item{ | ||||
| 		value:value, | ||||
| 	} | ||||
| } | ||||
| 
 | ||||
| func (lq *LoafleQueue) notifyEventHandler(c chan interface{}) { | ||||
| 	for data := range c { | ||||
| 		it := lq.newItem(data) | ||||
| 		heap.Push(lq, it) | ||||
| 	} | ||||
| } | ||||
| 
 | ||||
| func (lq *LoafleQueue) Close() { | ||||
| 	if lq.eventChanel != nil { | ||||
| 		observer.Remove(lq.queueType, lq.eventChanel) | ||||
| 	} | ||||
| } | ||||
| func NewLoafleQueue(eventType EventType, interval time.Duration) *LoafleQueue { | ||||
| 	items := make([]*item, 0) | ||||
| 	event := make(chan interface{},0) | ||||
| 	var tempType string | ||||
| 
 | ||||
| 	if eventType == DATA_TYPE { | ||||
| 		tempType = "DATA_TYPE" | ||||
| 	}else if eventType == EVENT_TYPE{ | ||||
| 		tempType = "EVENT_TYPE" | ||||
| 	}else { | ||||
| 		log.Fatal("Event Type Error") | ||||
| 	} | ||||
| 
 | ||||
| 	lq := &LoafleQueue{ | ||||
| 		items:items, | ||||
| 		queueType:tempType, | ||||
| 		interval:interval, | ||||
| 		eventChanel:event, | ||||
| 	} | ||||
| 
 | ||||
| 	observer.Add(lq.queueType, lq.eventChanel) | ||||
| 	go lq.notifyEventHandler(event) | ||||
| 
 | ||||
| 	return lq | ||||
| } | ||||
		Loading…
	
	
			
			x
			
			
		
	
		Reference in New Issue
	
	Block a user