Add documentation

This commit is contained in:
Thorsten Sommer committed 2015-06-17 17:44:52 +02:00
1 parent 69d0bdabf4
commit 6dab89a1d4
144 files changed
+1107 -549

No files matched your search

@@ -1,72 +1,40 @@
package ICCC
import (
"fmt"
"github.com/SommerEngineering/Ocean/ICCC/Scheme"
"github.com/SommerEngineering/Ocean/Log"
LM "github.com/SommerEngineering/Ocean/Log/Meta"
"github.com/SommerEngineering/Ocean/Shutdown"
"gopkg.in/mgo.v2/bson"
"time"
)
func InitCacheNow() {
startCacheTimerLock.Lock()
defer startCacheTimerLock.Unlock()
if cacheTimerRunning {
return
}
cacheTimerLogic(false)
}
func StartCacheTimer() {
initCacheTimer()
}
func initCacheTimer() {
startCacheTimerLock.Lock()
defer startCacheTimerLock.Unlock()
if cacheTimerRunning {
return
} else {
cacheTimerRunning = true
}
go func() {
for {
cacheTimerLogic(true)
}
}()
}
func cacheTimerLogic(waiting bool) {
if Shutdown.IsDown() {
return
}
lastCount := cacheListenerDatabase.Len()
selection := bson.D{{`IsActive`, true}}
entriesIterator := collectionListener.Find(selection).Iter()
entry := Scheme.Listener{}
cacheListenerDatabaseLock.Lock()
cacheListenerDatabase.Init()
for entriesIterator.Next(&entry) {
cacheListenerDatabase.PushBack(entry)
}
cacheListenerDatabaseLock.Unlock()
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelINFO, LM.MessageNameEXECUTE, `The listener cache was refreshed with the values from the database.`, fmt.Sprintf(`last count=%d`, lastCount), fmt.Sprintf(`new count=%d`, cacheListenerDatabase.Len()))
if waiting {
nextDuration := time.Duration(5) * time.Minute
if cacheListenerDatabase.Len() == 0 {
nextDuration = time.Duration(10) * time.Second
}
time.Sleep(nextDuration)
}
}
package ICCC
import (
"fmt"
"github.com/SommerEngineering/Ocean/ICCC/Scheme"
"github.com/SommerEngineering/Ocean/Log"
LM "github.com/SommerEngineering/Ocean/Log/Meta"
"github.com/SommerEngineering/Ocean/Shutdown"
"gopkg.in/mgo.v2/bson"
"time"
)
func cacheTimerLogic(waiting bool) {
if Shutdown.IsDown() {
return
}
lastCount := cacheListenerDatabase.Len()
selection := bson.D{{`IsActive`, true}}
entriesIterator := collectionListener.Find(selection).Iter()
entry := Scheme.Listener{}
cacheListenerDatabaseLock.Lock()
cacheListenerDatabase.Init()
for entriesIterator.Next(&entry) {
cacheListenerDatabase.PushBack(entry)
}
cacheListenerDatabaseLock.Unlock()
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelINFO, LM.MessageNameEXECUTE, `The listener cache was refreshed with the values from the database.`, fmt.Sprintf(`last count=%d`, lastCount), fmt.Sprintf(`new count=%d`, cacheListenerDatabase.Len()))
if waiting {
nextDuration := time.Duration(5) * time.Minute
if cacheListenerDatabase.Len() == 0 {
nextDuration = time.Duration(10) * time.Second
}
time.Sleep(nextDuration)
}
}
+5
View File
@@ -6,6 +6,7 @@ import (
"strconv"
)
// Function to convert the HTTP data back to a message.
func Data2Message(target interface{}, data map[string][]string) (channel, command string, obj interface{}) {
if data == nil || len(data) == 0 {
channel = ``
@@ -14,12 +15,15 @@ func Data2Message(target interface{}, data map[string][]string) (channel, comman
return
}
// Use reflection for the target type:
element := reflect.ValueOf(target)
element = element.Elem()
elementType := element.Type()
channel = data[`channel`][0]
command = data[`command`][0]
// Use the order of the destination type's fields:
for i := 0; i < element.NumField(); i++ {
field := element.Field(i)
switch field.Kind().String() {
@@ -53,6 +57,7 @@ func Data2Message(target interface{}, data map[string][]string) (channel, comman
v, _ := strconv.ParseUint(mapValue, 16, 8)
field.SetUint(v)
// Case: Arrays...
case `slice`:
sliceInterface := field.Interface()
sliceKind := reflect.ValueOf(sliceInterface).Type().String()
+8 -12
View File
@@ -1,21 +1,15 @@
/*
This is the "[I]nter-[C]omponent [C]ommunication [C]hannel". It is a minimal
messaging service to connect different servers or even different parts of
huge systems across programming languages.
This is the "[I]nter-[C]omponent [C]ommunication [C]hannel". It is a minimal messaging service to connect different servers or even different parts of huge systems across programming languages.
The basis idea is to create such messaging service on top of HTTP, because
every programming language is able to process HTTP. Therefore, all messages
are transformed to HTTP form values (with URL encoding).
The basis idea is to create such messaging service on top of HTTP, because every programming language is able to process HTTP. Therefore, all messages are transformed to HTTP form values (with URL encoding).
To be able to marshal / parse the data back to objects, some additional
information is added:
To be able to marshal / parse the data back to objects, some additional information is added:
Example 01:
name=str:Surname
value=Sommer
The HTTP form name is 'str:Surname' and the value is 'Sommer'. The 'str' is
the indicator for the data type, in this case it is a string.
The HTTP form name is 'str:Surname' and the value is 'Sommer'. The 'str' is the indicator for the data type, in this case it is a string.
Known data types are:
* str := string
@@ -48,7 +42,9 @@ channel=CHANNEL
[any count of data tuples]
InternalCommPassword=[configured communication password e.g. an UUID etc.]
If you want to build a distributed system across the Internet, please use e.g. SSH tunnels
to keep things secret.
If you want to build a distributed system across the Internet, please use e.g. SSH tunnels to keep things secret.
Constrains to the environment:
The web server cannot reorder the fields of the request or response. The order of fields at the data object (message) must correspond with the order of fields inside the HTTP message. Therefore, a reorder is not possible at the moment.
*/
package ICCC
+13
View File
@@ -8,36 +8,49 @@ import (
"net/http"
)
// The HTTP handler for ICCC.
func ICCCHandler(response http.ResponseWriter, request *http.Request) {
// Cannot parse the form?
if errParse := request.ParseForm(); errParse != nil {
Log.LogFull(senderName, LM.CategorySYSTEM, LM.LevelERROR, LM.SeverityCritical, LM.ImpactCritical, LM.MessageNameNETWORK, `Was not able to parse the HTTP form data from an ICCC message!`)
http.NotFound(response, request)
return
}
// Read the data out of the request:
messageData := map[string][]string(request.PostForm)
// The data must contain at least three fields (command, channel & communication password)
if len(messageData) < 3 {
Log.LogFull(senderName, LM.CategorySYSTEM, LM.LevelERROR, LM.SeverityCritical, LM.ImpactCritical, LM.MessageNameNETWORK, `The ICCC message contains not enough data: At least the channel, command and password is required!`)
http.NotFound(response, request)
return
}
// Read the meta data:
channel := messageData[`channel`][0]
command := messageData[`command`][0]
password := messageData[`InternalCommPassword`][0]
// Check the password:
if password != Tools.InternalCommPassword() {
Log.LogFull(senderName, LM.CategorySYSTEM, LM.LevelSECURITY, LM.SeverityCritical, LM.ImpactNone, LM.MessageNamePASSWORD, `Received a ICCC message with wrong password!`, request.RemoteAddr)
http.NotFound(response, request)
return
}
// Build the key for the mapping of the listener cache:
key := fmt.Sprintf(`%s::%s`, channel, command)
// Get the matching listener
listener := listeners[key]
if listener == nil {
// Case: No such listener
Log.LogFull(senderName, LM.CategorySYSTEM, LM.LevelWARN, LM.SeverityCritical, LM.ImpactUnknown, LM.MessageNameCONFIGURATION, `Was not able to find the correct listener for these ICCC message.`, `channel=`+channel, `command`+command, `hostname=`+Tools.ThisHostname())
} else {
// Case: Everything is fine => deliver the message
listener(messageData)
}
}
+7
View File
@@ -7,16 +7,23 @@ import (
"github.com/SommerEngineering/Ocean/Tools"
)
// Init this package.
func init() {
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelINFO, LM.MessageNameINIT, `Start init of ICCC.`)
defer Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelINFO, LM.MessageNameINIT, `Done init ICCC.`)
// Create the list as cache for all global listener (not only listener from this server):
cacheListenerDatabase = list.New()
// Create a mapping as cache for all local listener end-points (functions):
listeners = make(map[string]func(data map[string][]string))
// Using the local IP address:
correctAddressWithPort = Tools.LocalIPAddressAndPort()
// Init the database:
initDB()
// Register this server to the listener (if not present):
registerHost2Database()
}
+13
View File
@@ -0,0 +1,13 @@
package ICCC
// Starts the timer cache once and exit it after (no thread, no endless loop).
func InitCacheNow() {
startCacheTimerLock.Lock()
defer startCacheTimerLock.Unlock()
if cacheTimerRunning {
return
}
cacheTimerLogic(false)
}
+21
View File
@@ -0,0 +1,21 @@
package ICCC
// Setup and starts the cache timer.
func initCacheTimer() {
startCacheTimerLock.Lock()
defer startCacheTimerLock.Unlock()
if cacheTimerRunning {
return
} else {
cacheTimerRunning = true
}
// Start another thread with the timer logic:
go func() {
// Endless loop:
for {
cacheTimerLogic(true)
}
}()
}
+6
View File
@@ -7,6 +7,7 @@ import (
"gopkg.in/mgo.v2"
)
// Init the database.
func initDB() {
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelINFO, LM.MessageNameINIT, `Start init of the ICCC collections.`)
defer Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelINFO, LM.MessageNameINIT, `Done init the ICCC collection.`)
@@ -14,6 +15,7 @@ func initDB() {
// Get the database:
dbSession, db = CustomerDB.DB()
// Case: Error?
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
@@ -23,7 +25,9 @@ func initDB() {
collectionListener = db.C(`ICCCListener`)
collectionHosts = db.C(`ICCCHosts`)
//
// Take care about the indexes for ICCCListener:
//
collectionListener.EnsureIndexKey(`Command`)
collectionListener.EnsureIndexKey(`Command`, `IsActive`)
@@ -45,7 +49,9 @@ func initDB() {
indexName1.Unique = true
collectionListener.EnsureIndex(indexName1)
//
// Index for hosts:
//
collectionHosts.EnsureIndexKey(`Hostname`, `IPAddressPort`)
indexName2 := mgo.Index{}
+9
View File
@@ -6,8 +6,13 @@ import (
"strconv"
)
// Function to convert an ICCC message to HTTP data.
func message2Data(channel, command string, message interface{}) (data map[string][]string) {
// Create the map:
data = make(map[string][]string)
// Add the meta information:
data[`command`] = []string{command}
data[`channel`] = []string{channel}
@@ -15,9 +20,12 @@ func message2Data(channel, command string, message interface{}) (data map[string
return
}
// Use reflection to determine the types:
element := reflect.ValueOf(message)
elementType := element.Type()
// Iterate over all fields of the data type.
// Transform the data regarding the type.
for i := 0; i < element.NumField(); i++ {
field := element.Field(i)
keyName := elementType.Field(i).Name
@@ -44,6 +52,7 @@ func message2Data(channel, command string, message interface{}) (data map[string
key := fmt.Sprintf(`ui8:%s`, keyName)
data[key] = []string{strconv.FormatUint(field.Uint(), 16)}
// Case: Arrays...
case `slice`:
sliceLen := field.Len()
if sliceLen > 0 {
+7
View File
@@ -7,10 +7,13 @@ import (
"gopkg.in/mgo.v2/bson"
)
// The internal function to register a command to ICCC.
func register2Database(channel, command string) {
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelINFO, LM.MessageNameSTARTUP, `Register this ICCC command in to the database.`, `channel=`+channel, `command=`+command)
//
// Case: Exist and active :)
//
emptyEntry := Scheme.Listener{}
selection := bson.D{{`Channel`, channel}, {`Command`, command}, {`IPAddressPort`, correctAddressWithPort}, {`IsActive`, true}}
count1, _ := collectionListener.Find(selection).Count()
@@ -20,7 +23,9 @@ func register2Database(channel, command string) {
return
}
//
// Case: Exist but not active
//
selection = bson.D{{`Channel`, channel}, {`Command`, command}, {`IPAddressPort`, correctAddressWithPort}, {`IsActive`, false}}
notActiveEntry := Scheme.Listener{}
collectionListener.Find(selection).One(&notActiveEntry)
@@ -32,7 +37,9 @@ func register2Database(channel, command string) {
return
}
//
// Case: Not exist
//
Log.LogFull(senderName, LM.CategorySYSTEM, LM.LevelWARN, LM.SeverityCritical, LM.ImpactNone, LM.MessageNameCONFIGURATION, `This ICCC command is not known.`, `Create now a new entry!`)
entry := Scheme.Listener{}
+11
View File
@@ -8,20 +8,31 @@ import (
"gopkg.in/mgo.v2/bson"
)
// Function to register this server to the ICCC.
func registerHost2Database() {
// Create the host entry:
host := Scheme.Host{}
host.Hostname = Tools.ThisHostname()
host.IPAddressPort = correctAddressWithPort
// The query to find already existing entries:
selection := bson.D{{`Hostname`, host.Hostname}, {`IPAddressPort`, host.IPAddressPort}}
// Count the already existing entries:
count, _ := collectionHosts.Find(selection).Count()
// Already exist?
if count == 1 {
// Case: Exists!
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelINFO, LM.MessageNameCONFIGURATION, `This host is already registered!`, `host=`+host.Hostname, `address=`+host.IPAddressPort)
} else {
// Case: Not exist.
if errInsert := collectionHosts.Insert(host); errInsert != nil {
// Case: Was not able insert in the database
Log.LogFull(senderName, LM.CategorySYSTEM, LM.LevelERROR, LM.SeverityCritical, LM.ImpactCritical, LM.MessageNameDATABASE, `Was not able to register this host.`, errInsert.Error(), `host=`+host.Hostname, `address=`+host.IPAddressPort)
} else {
// Case: Everything was fine.
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelINFO, LM.MessageNameSTARTUP, `This host is now registered.`, `host=`+host.Hostname, `address=`+host.IPAddressPort)
}
}
+4
View File
@@ -6,11 +6,15 @@ import (
LM "github.com/SommerEngineering/Ocean/Log/Meta"
)
// Register a command to ICCC for a specific channel.
func Registrar(channel, command string, callback func(data map[string][]string)) {
listenersLock.Lock()
defer listenersLock.Unlock()
// Write the command to the database:
register2Database(channel, command)
// Register the command at the local cache:
listeners[fmt.Sprintf(`%s::%s`, channel, command)] = callback
Log.LogShort(senderName, LM.CategorySYSTEM, LM.LevelINFO, LM.MessageNameCONFIGURATION, `The registrar has registered a new ICCC command.`, `channel=`+channel, `command=`+command)
+1
View File
@@ -1,5 +1,6 @@
package Scheme
// Scheme for the host database entry.
type Host struct {
Hostname string `bson:"Hostname"`
IPAddressPort string `bson:"IPAddressPort"`
+1
View File
@@ -1,5 +1,6 @@
package Scheme
// Type for the listener entries at the database.
type Listener struct {
Channel string `bson:"Channel"`
Command string `bson:"Command"`
+7 -1
View File
@@ -9,11 +9,17 @@ import (
"net/url"
)
// Send a message to a listener.
func sendMessage(listener Scheme.Listener, data map[string][]string) {
// Convert the data and encode it:
valuesHTTP := url.Values(data)
// Add the communication password:
valuesHTTP.Add(`InternalCommPassword`, Tools.InternalCommPassword())
// Try to deliver the message:
if _, err := http.PostForm(`http://`+listener.IPAddressPort+`/ICCC`, valuesHTTP); err != nil {
// Case: Was not possible to deliver.
Log.LogFull(senderName, LM.CategorySYSTEM, LM.LevelERROR, LM.SeverityCritical, LM.ImpactUnknown, LM.MessageNameNETWORK, `Was not able to send the ICCC message.`, err.Error())
}
+11 -6
View File
@@ -7,27 +7,32 @@ import (
"gopkg.in/mgo.v2/bson"
)
/*
Please do not use this type. It is an internal type of Ocean to provide a shutdown function!
*/
// Type to provide a shutdown function.
type ShutdownFunction struct {
}
/*
This function is called if the Ocean server is shutting down.
*/
// The shutdown function for ICCC.
func (a ShutdownFunction) Shutdown() {
Log.LogShort(senderName, LM.CategoryAPP, LM.LevelWARN, LM.MessageNameSHUTDOWN, `Shutting down now all ICCC listener for this host.`)
// Define the database query:
selection := bson.D{{`IPAddressPort`, correctAddressWithPort}}
// Reserve the memory for an answer:
entry := Scheme.Listener{}
// Execute the query and iterate over the results:
iterator := collectionListener.Find(selection).Iter()
for iterator.Next(&entry) {
// Update the entry and set it to active=false:
selectionUpdate := bson.D{{`Channel`, entry.Channel}, {`Command`, entry.Command}, {`IPAddressPort`, correctAddressWithPort}}
entry.IsActive = false
// Update the entry:
collectionListener.Update(selectionUpdate, entry)
}
// Disconnect the database:
db.Logout()
dbSession.Close()
Log.LogShort(senderName, LM.CategoryAPP, LM.LevelWARN, LM.MessageNameSHUTDOWN, `Done shutting down now all ICCC listener for this host.`)
+6
View File
@@ -0,0 +1,6 @@
package ICCC
// Starts the cache timer thread.
func StartCacheTimer() {
initCacheTimer()
}
+1
View File
@@ -1,5 +1,6 @@
package SystemMessages
// Message type for the startup message:
type ICCCStartUpMessage struct {
PublicIPAddressAndPort string
AdminIPAddressAndPort string
+19 -18
View File
@@ -7,26 +7,27 @@ import (
"sync"
)
// Some pre-defined channels:
const (
ChannelSYSTEM string = `System`
ChannelNUMGEN string = `System::NumGen`
ChannelSHUTDOWN string = `System::Shutdown`
ChannelSTARTUP string = `System::Startup`
ChannelICCC string = `System::ICCC`
ChannelSYSTEM string = `System` // The common system channel.
ChannelNUMGEN string = `System::NumGen` // A channel for the number generator.
ChannelSHUTDOWN string = `System::Shutdown` // A channel for system shutdown messages.
ChannelSTARTUP string = `System::Startup` // A channel for system startup messages.
ChannelICCC string = `System::ICCC` // A common ICCC channel.
)
var (
senderName LM.Sender = `ICCC`
db *mgo.Database = nil
dbSession *mgo.Session = nil
collectionListener *mgo.Collection = nil
collectionHosts *mgo.Collection = nil
reservedSystemChannels []string = []string{ChannelSYSTEM, ChannelNUMGEN, ChannelSHUTDOWN, ChannelSTARTUP, ChannelICCC}
listeners map[string]func(data map[string][]string) = nil
listenersLock sync.RWMutex = sync.RWMutex{}
cacheListenerDatabase *list.List = nil
cacheListenerDatabaseLock sync.RWMutex = sync.RWMutex{}
startCacheTimerLock sync.Mutex = sync.Mutex{}
cacheTimerRunning bool = false
correctAddressWithPort string = ``
senderName LM.Sender = `ICCC` // This is the name for logging event from this package
db *mgo.Database = nil // The database
dbSession *mgo.Session = nil // The database session
collectionListener *mgo.Collection = nil // The database collection for listeners
collectionHosts *mgo.Collection = nil // The database collection for hosts
reservedSystemChannels []string = []string{ChannelSYSTEM, ChannelNUMGEN, ChannelSHUTDOWN, ChannelSTARTUP, ChannelICCC} // The reserved and pre-defined system channels
listeners map[string]func(data map[string][]string) = nil // The listener cache for all local available listeners with local functions
listenersLock sync.RWMutex = sync.RWMutex{} // The mutex for the listener cache
cacheListenerDatabase *list.List = nil // The globally cache for all listeners from all servers
cacheListenerDatabaseLock sync.RWMutex = sync.RWMutex{} // The mutex for the globally cache
startCacheTimerLock sync.Mutex = sync.Mutex{} // Mutex for the start timer
cacheTimerRunning bool = false // Is the timer running?
correctAddressWithPort string = `` // The IP address and port of the this local server
)
+7
View File
@@ -6,20 +6,27 @@ import (
LM "github.com/SommerEngineering/Ocean/Log/Meta"
)
// Function to broadcast a message to all listeners.
func WriteMessage2All(channel, command string, message interface{}) {
cacheListenerDatabaseLock.RLock()
defer cacheListenerDatabaseLock.RUnlock()
// Convert the message to HTTP data:
data := message2Data(channel, command, message)
counter := 0
// Loop over all listeners which are currently available at the cache:
for entry := cacheListenerDatabase.Front(); entry != nil; entry = entry.Next() {
listener := entry.Value.(Scheme.Listener)
// If the channel and the command matches, deliver the message:
if listener.Channel == channel && listener.Command == command {
go sendMessage(listener, data)
counter++
}
}
// Was not able to deliver to any listener?
if counter == 0 {
Log.LogFull(senderName, LM.CategorySYSTEM, LM.LevelWARN, LM.SeverityCritical, LM.ImpactUnknown, LM.MessageNameCONFIGURATION, `It was not able to deliver this message, because no listener was found!`, `channel=`+channel, `command=`+command)
}
+8
View File
@@ -7,16 +7,22 @@ import (
"github.com/SommerEngineering/Ocean/Tools"
)
// Function to write a message to any listener.
func WriteMessage2Any(channel, command string, message interface{}) {
cacheListenerDatabaseLock.RLock()
defer cacheListenerDatabaseLock.RUnlock()
// Convert the message to HTTP data:
data := message2Data(channel, command, message)
maxCount := cacheListenerDatabase.Len()
entries := make([]Scheme.Listener, 0, maxCount)
counter := 0
// Loop over all listeners which are currently present at the cache:
for entry := cacheListenerDatabase.Front(); entry != nil; entry = entry.Next() {
listener := entry.Value.(Scheme.Listener)
// If the channel and the command matches, store the listener:
if listener.Channel == channel && listener.Command == command {
entries = entries[:len(entries)+1]
entries[counter] = listener
@@ -25,9 +31,11 @@ func WriteMessage2Any(channel, command string, message interface{}) {
count := len(entries)
if count > 0 {
// Case: Find at least one possible listener. Choose a random one and deliver:
listener := entries[Tools.RandomInteger(count)]
go sendMessage(listener, data)
} else {
// Case: Find no listener at all.
Log.LogFull(senderName, LM.CategorySYSTEM, LM.LevelWARN, LM.SeverityCritical, LM.ImpactUnknown, LM.MessageNameCONFIGURATION, `It was not able to deliver this message to any listener, because no listener was found!`, `channel=`+channel, `command=`+command)
}
}