Initial commit of Ocean's local development
This commit is contained in:
1 parent
9944a9e1df
commit
86451938ec
116 files changed
+3320
No files matched your search
@@ -0,0 +1,6 @@
|
||||
package NumGen
|
||||
|
||||
func BadNumber() (result int64) {
|
||||
result = badNumber64
|
||||
return
|
||||
}
|
||||
@@ -0,0 +1,30 @@
|
||||
package NumGen
|
||||
|
||||
import "net/http"
|
||||
import "net/url"
|
||||
import "strconv"
|
||||
import "github.com/SommerEngineering/Ocean/Shutdown"
|
||||
import "github.com/SommerEngineering/Ocean/Log"
|
||||
import LM "github.com/SommerEngineering/Ocean/Log/Meta"
|
||||
|
||||
func GetNextInt64(name string) (result int64) {
|
||||
result = badNumber64
|
||||
|
||||
if Shutdown.IsDown() {
|
||||
return
|
||||
}
|
||||
|
||||
if responseData, errRequest := http.PostForm(getHandler, url.Values{"name": {name}, "password": {correctPassword}}); errRequest != nil {
|
||||
Log.LogFull(senderName, LM.CategorySYSTEM, LM.LevelERROR, LM.SeverityCritical, LM.ImpactCritical, LM.MessageNameGENERATOR, `Requesting the next number was not possible.`, errRequest.Error())
|
||||
return
|
||||
} else {
|
||||
nextNumberText := responseData.Header.Get(`nextNumber`)
|
||||
if number, errAtio := strconv.Atoi(nextNumberText); errAtio != nil {
|
||||
Log.LogFull(senderName, LM.CategorySYSTEM, LM.LevelERROR, LM.SeverityCritical, LM.ImpactCritical, LM.MessageNameGENERATOR, `It was not possible to convert the answer into an int64.`, errAtio.Error())
|
||||
return
|
||||
} else {
|
||||
result = int64(number)
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,41 @@
|
||||
package NumGen
|
||||
|
||||
import "fmt"
|
||||
import "net/http"
|
||||
import "github.com/SommerEngineering/Ocean/Shutdown"
|
||||
import "github.com/SommerEngineering/Ocean/Log"
|
||||
import LM "github.com/SommerEngineering/Ocean/Log/Meta"
|
||||
|
||||
func HandlerGetNext(response http.ResponseWriter, request *http.Request) {
|
||||
if Shutdown.IsDown() {
|
||||
http.NotFound(response, request)
|
||||
return
|
||||
}
|
||||
|
||||
if !isActive {
|
||||
Log.LogFull(senderName, LM.CategorySYSTEM, LM.LevelWARN, LM.SeverityCritical, LM.ImpactNone, LM.MessageNameCONFIGURATION, `Called the get handler on an inactive host.`, `Wrong configuration?`)
|
||||
http.NotFound(response, request)
|
||||
return
|
||||
}
|
||||
|
||||
if correctPassword == `` {
|
||||
Log.LogFull(senderName, LM.CategorySYSTEM, LM.LevelERROR, LM.SeverityCritical, LM.ImpactCritical, LM.MessageNameSECURITY, `No communication password was set.`)
|
||||
http.NotFound(response, request)
|
||||
return
|
||||
}
|
||||
|
||||
name := request.FormValue(`name`)
|
||||
pwd := request.FormValue(`password`)
|
||||
|
||||
if pwd != correctPassword {
|
||||
Log.LogFull(senderName, LM.CategorySYSTEM, LM.LevelSECURITY, LM.SeverityCritical, LM.ImpactNone, LM.MessageNamePASSWORD, `A wrong password was used to access this system handler.`, `This should never happens: Is this a hacking attempt?`, `IP address of requester=`+request.RemoteAddr)
|
||||
http.NotFound(response, request)
|
||||
return
|
||||
}
|
||||
|
||||
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelDEBUG, LM.MessageNameANALYSIS, `Next number requested.`, name, pwd)
|
||||
channel := requestChannel4Name(name)
|
||||
nextNumber := <-channel
|
||||
|
||||
response.Header().Add(`nextNumber`, fmt.Sprintf(`%d`, nextNumber))
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
package NumGen
|
||||
|
||||
import "strings"
|
||||
import "strconv"
|
||||
import "github.com/SommerEngineering/Ocean/Tools"
|
||||
import "github.com/SommerEngineering/Ocean/ConfigurationDB"
|
||||
import "github.com/SommerEngineering/Ocean/Log"
|
||||
import LM "github.com/SommerEngineering/Ocean/Log/Meta"
|
||||
|
||||
func init() {
|
||||
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelINFO, LM.MessageNameSTARTUP, `Init the number generator.`)
|
||||
|
||||
channelListLock.Lock()
|
||||
defer channelListLock.Unlock()
|
||||
|
||||
correctPassword = ConfigurationDB.Read(`InternalCommPassword`)
|
||||
activeHost := ConfigurationDB.Read(`NumGenActiveHosts`)
|
||||
isActive = strings.Contains(activeHost, Tools.ThisHostname())
|
||||
getHandler = ConfigurationDB.Read(`NumGenGetHandler`)
|
||||
|
||||
if isActive {
|
||||
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelWARN, LM.MessageNameCONFIGURATION, `The number generator is active on this host.`, `This host is producer and consumer.`)
|
||||
|
||||
channelBufferSizeText := ConfigurationDB.Read(`NumGenBufferSize`)
|
||||
if bufferSizeNumber, errBufferSizeNumber := strconv.Atoi(channelBufferSizeText); errBufferSizeNumber != nil {
|
||||
Log.LogFull(senderName, LM.CategorySYSTEM, LM.LevelERROR, LM.SeverityCritical, LM.ImpactMiddle, LM.MessageNameCONFIGURATION, `Was not able to parse the configuration value of NumGenBufferSize.`, errBufferSizeNumber.Error(), `Use the default value now!`)
|
||||
} else {
|
||||
channelBufferSize = bufferSizeNumber
|
||||
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelINFO, LM.MessageNameCONFIGURATION, `The buffer size for the number generator was loaded.`, `Buffer size=`+channelBufferSizeText)
|
||||
}
|
||||
|
||||
channelList = make(map[string]chan int64)
|
||||
|
||||
initDB()
|
||||
} else {
|
||||
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelWARN, LM.MessageNameCONFIGURATION, `The number generator is not active on this host.`, `This host is just a consumer.`)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,28 @@
|
||||
package NumGen
|
||||
|
||||
import "labix.org/v2/mgo"
|
||||
import "github.com/SommerEngineering/Ocean/CustomerDB"
|
||||
import "github.com/SommerEngineering/Ocean/Log"
|
||||
import LM "github.com/SommerEngineering/Ocean/Log/Meta"
|
||||
|
||||
func initDB() {
|
||||
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelINFO, LM.MessageNameINIT, `Start init of number generator collection.`)
|
||||
defer Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelINFO, LM.MessageNameINIT, `Done init of number generator collection.`)
|
||||
|
||||
// Get the database:
|
||||
db = CustomerDB.DB()
|
||||
|
||||
if db == nil {
|
||||
Log.LogFull(senderName, LM.CategorySYSTEM, LM.LevelERROR, LM.SeverityCritical, LM.ImpactCritical, LM.MessageNameDATABASE, `Was not able to get the customer database.`)
|
||||
return
|
||||
}
|
||||
|
||||
// Get my collection:
|
||||
collectionNumGen = db.C(`NumGen`)
|
||||
|
||||
// Take care about the indexes:
|
||||
indexName := mgo.Index{}
|
||||
indexName.Key = []string{`Name`}
|
||||
indexName.Unique = true
|
||||
collectionNumGen.EnsureIndex(indexName)
|
||||
}
|
||||
@@ -0,0 +1,68 @@
|
||||
package NumGen
|
||||
|
||||
import "time"
|
||||
import "labix.org/v2/mgo/bson"
|
||||
import "github.com/SommerEngineering/Ocean/Shutdown"
|
||||
import "github.com/SommerEngineering/Ocean/Log"
|
||||
import LM "github.com/SommerEngineering/Ocean/Log/Meta"
|
||||
|
||||
func producer(name string) {
|
||||
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelINFO, LM.MessageNameSTARTUP, `The NumGen producer is now starting.`, `name=`+name)
|
||||
|
||||
// Get my channel:
|
||||
myChannel := requestChannel4Name(name)
|
||||
|
||||
// Read my next free number:
|
||||
currentNextFreeNumber := nextFreeNumberFromDatabase(name)
|
||||
|
||||
// Where is the next "reload"?
|
||||
nextReload := currentNextFreeNumber + int64(channelBufferSize)
|
||||
|
||||
// Set the next free number to the database:
|
||||
updateNextFreeNumber2Database(name, nextReload)
|
||||
|
||||
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelINFO, LM.MessageNameSTARTUP, `The NumGen producer is now running.`, `name=`+name)
|
||||
for nextNumber := currentNextFreeNumber; true; {
|
||||
if Shutdown.IsDown() {
|
||||
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelINFO, LM.MessageNameSHUTDOWN, `The NumGen producer is now down.`, `name=`+name)
|
||||
return
|
||||
}
|
||||
|
||||
if nextNumber > nextReload {
|
||||
nextReload = nextReload + int64(channelBufferSize)
|
||||
updateNextFreeNumber2Database(name, nextReload)
|
||||
|
||||
// Enables the administrator to monitor the frequence of chunks and is able to reconfigure the settings:
|
||||
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelDEBUG, LM.MessageNamePRODUCER, `The NumGen producer creates the next chunk.`, `name=`+name)
|
||||
}
|
||||
|
||||
// Enqueue the next number:
|
||||
select {
|
||||
case myChannel <- nextNumber:
|
||||
nextNumber++
|
||||
case <-time.After(time.Millisecond * 500):
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func nextFreeNumberFromDatabase(name string) (result int64) {
|
||||
selection := bson.D{{`Name`, name}}
|
||||
searchResult := NumberGenScheme{}
|
||||
|
||||
count, _ := collectionNumGen.Find(selection).Count()
|
||||
if count == 1 {
|
||||
collectionNumGen.Find(selection).One(&searchResult)
|
||||
result = searchResult.NextFreeNumber
|
||||
} else {
|
||||
searchResult.Name = name
|
||||
searchResult.NextFreeNumber = startValue64
|
||||
collectionNumGen.Insert(searchResult)
|
||||
result = searchResult.NextFreeNumber
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func updateNextFreeNumber2Database(name string, nextFreeNumber int64) {
|
||||
selection := bson.D{{`Name`, name}}
|
||||
collectionNumGen.Update(selection, NumberGenScheme{Name: name, NextFreeNumber: nextFreeNumber})
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
package NumGen
|
||||
|
||||
import "github.com/SommerEngineering/Ocean/Shutdown"
|
||||
import "github.com/SommerEngineering/Ocean/Log"
|
||||
import LM "github.com/SommerEngineering/Ocean/Log/Meta"
|
||||
|
||||
func requestChannel4Name(name string) (result chan int64) {
|
||||
|
||||
if Shutdown.IsDown() {
|
||||
return
|
||||
}
|
||||
|
||||
if !isActive {
|
||||
Log.LogFull(senderName, LM.CategorySYSTEM, LM.LevelWARN, LM.SeverityCritical, LM.ImpactNone, LM.MessageNameCONFIGURATION, `Called the requestChannel4Name() on an inactive host.`, `Wrong configuration?`)
|
||||
return
|
||||
}
|
||||
|
||||
channelListLock.RLock()
|
||||
channel, isPresent := channelList[name]
|
||||
channelListLock.RUnlock()
|
||||
|
||||
if isPresent {
|
||||
result = channel
|
||||
return
|
||||
}
|
||||
|
||||
// Create the entry:
|
||||
newChannel := make(chan int64, channelBufferSize)
|
||||
result = newChannel
|
||||
|
||||
channelListLock.Lock()
|
||||
channelList[name] = newChannel
|
||||
channelListLock.Unlock()
|
||||
|
||||
// Create the new producer:
|
||||
go producer(name)
|
||||
return
|
||||
}
|
||||
@@ -0,0 +1,6 @@
|
||||
package NumGen
|
||||
|
||||
type NumberGenScheme struct {
|
||||
Name string `bson:"Name"`
|
||||
NextFreeNumber int64 `bson:"NextFreeNumber"`
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
package NumGen
|
||||
|
||||
import "github.com/SommerEngineering/Ocean/Log"
|
||||
import LM "github.com/SommerEngineering/Ocean/Log/Meta"
|
||||
|
||||
type ShutdownFunction struct {
|
||||
}
|
||||
|
||||
func (a ShutdownFunction) Shutdown() {
|
||||
Log.LogShort(senderName, LM.CategoryAPP, LM.LevelWARN, LM.MessageNameSHUTDOWN, `Shutting down the number generator.`)
|
||||
}
|
||||
@@ -0,0 +1,22 @@
|
||||
package NumGen
|
||||
|
||||
import "sync"
|
||||
import "labix.org/v2/mgo"
|
||||
import LM "github.com/SommerEngineering/Ocean/Log/Meta"
|
||||
|
||||
var (
|
||||
correctPassword string = ``
|
||||
senderName LM.Sender = `System::NumGen::Producer`
|
||||
isActive bool = false
|
||||
getHandler string = ``
|
||||
db *mgo.Database = nil
|
||||
collectionNumGen *mgo.Collection = nil
|
||||
channelBufferSize int = 10
|
||||
channelList map[string]chan int64 = nil
|
||||
channelListLock sync.RWMutex = sync.RWMutex{}
|
||||
)
|
||||
|
||||
const (
|
||||
badNumber64 int64 = 9222222222222222222
|
||||
startValue64 int64 = -9223372036854775808
|
||||
)
|
||||
Reference in new issue
Block a user