job manager, jobm pkg
This commit is contained in:
parent
a17926f54b
commit
740fb94556
@ -2,6 +2,7 @@ package filter
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"io/ioutil"
|
"io/ioutil"
|
||||||
|
"os"
|
||||||
"path"
|
"path"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
@ -15,16 +16,14 @@ import (
|
|||||||
// RegisterFilters reads a directory and register filters from files within it
|
// RegisterFilters reads a directory and register filters from files within it
|
||||||
func RegisterFilters(dir string) {
|
func RegisterFilters(dir string) {
|
||||||
files, err := ioutil.ReadDir(dir)
|
files, err := ioutil.ReadDir(dir)
|
||||||
if err != nil {
|
logger.Eexit(err, "could not read from template filters dir '%s'", dir)
|
||||||
logger.Log.Panicf("could not read from template filters dir '%s': %s", dir, err)
|
|
||||||
}
|
|
||||||
for _, f := range files {
|
for _, f := range files {
|
||||||
if !f.IsDir() {
|
if !f.IsDir() {
|
||||||
switch path.Ext(f.Name()) {
|
switch path.Ext(f.Name()) {
|
||||||
case ".js":
|
case ".js":
|
||||||
fileBase := strings.TrimSuffix(f.Name(), ".js")
|
fileBase := strings.TrimSuffix(f.Name(), ".js")
|
||||||
jsFile := dir + "/" + f.Name()
|
jsFile := dir + "/" + f.Name()
|
||||||
logger.Log.Debugf("trying to register filter from: %s", jsFile)
|
logger.D("trying to register filter from: %s", jsFile)
|
||||||
/*
|
/*
|
||||||
jsStr, err := ioutil.ReadFile(jsFile)
|
jsStr, err := ioutil.ReadFile(jsFile)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@ -33,38 +32,31 @@ func RegisterFilters(dir string) {
|
|||||||
*/
|
*/
|
||||||
vm := motto.New()
|
vm := motto.New()
|
||||||
fn, err := vm.Run(jsFile)
|
fn, err := vm.Run(jsFile)
|
||||||
if err != nil {
|
logger.Eexit(err, "error in javascript vm for '%s'", jsFile)
|
||||||
logger.Log.Panicf("error in javascript vm for '%s': %s", jsFile, err)
|
|
||||||
}
|
|
||||||
if !fn.IsFunction() {
|
if !fn.IsFunction() {
|
||||||
logger.Log.Panicf("%s does not contain a function code", jsFile)
|
logger.E("%s does not contain a function code", jsFile)
|
||||||
|
os.Exit(1)
|
||||||
}
|
}
|
||||||
|
|
||||||
err = pongo2.RegisterFilter(
|
err = pongo2.RegisterFilter(
|
||||||
fileBase,
|
fileBase,
|
||||||
func(in, param *pongo2.Value) (out *pongo2.Value, erro *pongo2.Error) {
|
func(in, param *pongo2.Value) (out *pongo2.Value, erro *pongo2.Error) {
|
||||||
thisObj, _ := vm.Object("({})")
|
thisObj, _ := vm.Object("({})")
|
||||||
|
var err error
|
||||||
if mark2web.CurrentContext != nil {
|
if mark2web.CurrentContext != nil {
|
||||||
thisObj.Set("context", *mark2web.CurrentContext)
|
err = thisObj.Set("context", *mark2web.CurrentContext)
|
||||||
}
|
}
|
||||||
|
|
||||||
if err != nil {
|
logger.Perr(err, "could not set context in '%s': %s", jsFile)
|
||||||
logger.Log.Panicf("could not set context as in '%s': %s", jsFile, err)
|
|
||||||
}
|
|
||||||
ret, err := fn.Call(thisObj.Value(), in.Interface(), param.Interface())
|
ret, err := fn.Call(thisObj.Value(), in.Interface(), param.Interface())
|
||||||
if err != nil {
|
logger.Eexit(err, "error in javascript file '%s' while calling returned function", jsFile)
|
||||||
logger.Log.Panicf("error in javascript file '%s' while calling returned function: %s", jsFile, err)
|
|
||||||
}
|
|
||||||
retGo, err := ret.Export()
|
retGo, err := ret.Export()
|
||||||
if err != nil {
|
logger.Perr(err, "export error for '%s'", jsFile)
|
||||||
logger.Log.Panicf("export error for '%s': %s", jsFile, err)
|
|
||||||
}
|
|
||||||
return pongo2.AsValue(retGo), nil
|
return pongo2.AsValue(retGo), nil
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
if err != nil {
|
logger.Perr(err, "could not register filter from '%s'", jsFile)
|
||||||
logger.Log.Panicf("could not register filter from '%s': %s", jsFile, err)
|
|
||||||
}
|
|
||||||
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
@ -13,6 +13,7 @@ import (
|
|||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
"gitbase.de/apairon/mark2web/pkg/helper"
|
"gitbase.de/apairon/mark2web/pkg/helper"
|
||||||
|
"gitbase.de/apairon/mark2web/pkg/jobm"
|
||||||
"gitbase.de/apairon/mark2web/pkg/logger"
|
"gitbase.de/apairon/mark2web/pkg/logger"
|
||||||
"gitbase.de/apairon/mark2web/pkg/mark2web"
|
"gitbase.de/apairon/mark2web/pkg/mark2web"
|
||||||
"github.com/disintegration/imaging"
|
"github.com/disintegration/imaging"
|
||||||
@ -168,7 +169,8 @@ func ImageProcessFilter(in *pongo2.Value, param *pongo2.Value) (*pongo2.Value, *
|
|||||||
if f, err := os.Stat(imgTarget); err == nil && !f.IsDir() {
|
if f, err := os.Stat(imgTarget); err == nil && !f.IsDir() {
|
||||||
logger.Log.Noticef("skipped processing image from %s to %s, file already exists", imgSource, imgTarget)
|
logger.Log.Noticef("skipped processing image from %s to %s, file already exists", imgSource, imgTarget)
|
||||||
} else {
|
} else {
|
||||||
mark2web.ThreadStart(func() {
|
jobm.Enqueue(jobm.Job{
|
||||||
|
Function: func() {
|
||||||
logger.Log.Noticef("processing image from %s to %s", imgSource, imgTarget)
|
logger.Log.Noticef("processing image from %s to %s", imgSource, imgTarget)
|
||||||
if strings.HasPrefix(imgSource, "http://") || strings.HasPrefix(imgSource, "https://") {
|
if strings.HasPrefix(imgSource, "http://") || strings.HasPrefix(imgSource, "https://") {
|
||||||
// webrequest before finding target filename, because of file format in filename
|
// webrequest before finding target filename, because of file format in filename
|
||||||
@ -225,6 +227,9 @@ func ImageProcessFilter(in *pongo2.Value, param *pongo2.Value) (*pongo2.Value, *
|
|||||||
logger.Log.Panicf("filter:image_resize, could save image '%s': %s", imgTarget, err)
|
logger.Log.Panicf("filter:image_resize, could save image '%s': %s", imgTarget, err)
|
||||||
}
|
}
|
||||||
logger.Log.Noticef("finished image: %s", imgTarget)
|
logger.Log.Noticef("finished image: %s", imgTarget)
|
||||||
|
},
|
||||||
|
Description: "process image " + imgSource,
|
||||||
|
Category: "image process",
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
return pongo2.AsValue(mark2web.CurrentTreeNode.ResolveNavPath(p.Filename)), nil
|
return pongo2.AsValue(mark2web.CurrentTreeNode.ResolveNavPath(p.Filename)), nil
|
||||||
|
58
pkg/jobm/jobmanager.go
Normal file
58
pkg/jobm/jobmanager.go
Normal file
@ -0,0 +1,58 @@
|
|||||||
|
package jobm
|
||||||
|
|
||||||
|
import (
|
||||||
|
"runtime"
|
||||||
|
"sync"
|
||||||
|
|
||||||
|
"gitbase.de/apairon/mark2web/pkg/logger"
|
||||||
|
)
|
||||||
|
|
||||||
|
var wg sync.WaitGroup
|
||||||
|
var numCPU = runtime.NumCPU()
|
||||||
|
|
||||||
|
// Job is a wrapper to descripe a Job function
|
||||||
|
type Job struct {
|
||||||
|
Function func()
|
||||||
|
Description string
|
||||||
|
Category string
|
||||||
|
}
|
||||||
|
|
||||||
|
var jobChan = make(chan Job)
|
||||||
|
|
||||||
|
func worker(jobChan <-chan Job) {
|
||||||
|
defer wg.Done()
|
||||||
|
|
||||||
|
for job := range jobChan {
|
||||||
|
job.Function()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func init() {
|
||||||
|
logger.Log.Infof("number of CPU core: %d", numCPU)
|
||||||
|
// one core for main thread
|
||||||
|
for i := 0; i < (numCPU - 1); i++ {
|
||||||
|
wg.Add(1)
|
||||||
|
go worker(jobChan)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Enqueue enqueues a job to the job queue
|
||||||
|
func Enqueue(job Job) {
|
||||||
|
if numCPU <= 1 {
|
||||||
|
// no threading
|
||||||
|
job.Function()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
jobChan <- job
|
||||||
|
}
|
||||||
|
|
||||||
|
// Wait will wait for all jobs to finish
|
||||||
|
func Wait() {
|
||||||
|
close(jobChan)
|
||||||
|
wg.Wait()
|
||||||
|
}
|
||||||
|
|
||||||
|
// SetNumCPU is for testing package without threading
|
||||||
|
func SetNumCPU(i int) {
|
||||||
|
numCPU = i
|
||||||
|
}
|
@ -13,7 +13,7 @@ var logBackendLeveled logging.LeveledBackend
|
|||||||
|
|
||||||
// SetLogLevel sets log level for global logger (debug, info, notice, warning, error)
|
// SetLogLevel sets log level for global logger (debug, info, notice, warning, error)
|
||||||
func SetLogLevel(level string) {
|
func SetLogLevel(level string) {
|
||||||
logBackendLevel := logging.NOTICE
|
logBackendLevel := logging.INFO
|
||||||
switch level {
|
switch level {
|
||||||
case "debug":
|
case "debug":
|
||||||
logBackendLevel = logging.DEBUG
|
logBackendLevel = logging.DEBUG
|
||||||
@ -58,3 +58,67 @@ func init() {
|
|||||||
spew.Config.DisablePointerMethods = true
|
spew.Config.DisablePointerMethods = true
|
||||||
configureLogger()
|
configureLogger()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// D is shorthand for Debugf
|
||||||
|
func D(format string, args ...interface{}) {
|
||||||
|
Log.Debugf(format, args...)
|
||||||
|
}
|
||||||
|
|
||||||
|
// I is shorthand for Infof
|
||||||
|
func I(format string, args ...interface{}) {
|
||||||
|
Log.Infof(format, args...)
|
||||||
|
}
|
||||||
|
|
||||||
|
// N is shorthand for Noticef
|
||||||
|
func N(format string, args ...interface{}) {
|
||||||
|
Log.Noticef(format, args...)
|
||||||
|
}
|
||||||
|
|
||||||
|
// W is shorthand for Warningf
|
||||||
|
func W(format string, args ...interface{}) {
|
||||||
|
Log.Warningf(format, args...)
|
||||||
|
}
|
||||||
|
|
||||||
|
// E is shorthand for Errorf
|
||||||
|
func E(format string, args ...interface{}) {
|
||||||
|
Log.Errorf(format, args...)
|
||||||
|
}
|
||||||
|
|
||||||
|
// P is shorthand for Panicf
|
||||||
|
func P(format string, args ...interface{}) {
|
||||||
|
Log.Panicf(format, args...)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Eerr is shorthand for
|
||||||
|
// if err != nil {
|
||||||
|
// Log.Errorf(...)
|
||||||
|
// }
|
||||||
|
func Eerr(err error, format string, args ...interface{}) {
|
||||||
|
if err != nil {
|
||||||
|
args = append(args, err)
|
||||||
|
Log.Errorf(format+"\nError: %s", args...)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Eexit is shorthand for
|
||||||
|
// if err != nil {
|
||||||
|
// Log.Errorf(...)
|
||||||
|
// os.Exit(1)
|
||||||
|
// }
|
||||||
|
func Eexit(err error, format string, args ...interface{}) {
|
||||||
|
Eerr(err, format, args...)
|
||||||
|
if err != nil {
|
||||||
|
os.Exit(1)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Perr is shorthand for
|
||||||
|
// if err != nil {
|
||||||
|
// Log.Panicf(...)
|
||||||
|
// }
|
||||||
|
func Perr(err error, format string, args ...interface{}) {
|
||||||
|
if err != nil {
|
||||||
|
args = append(args, err)
|
||||||
|
Log.Panicf(format+"\nError: %s", args...)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
@ -7,11 +7,14 @@ import (
|
|||||||
"os"
|
"os"
|
||||||
"path"
|
"path"
|
||||||
|
|
||||||
|
"gitbase.de/apairon/mark2web/pkg/jobm"
|
||||||
|
|
||||||
"gitbase.de/apairon/mark2web/pkg/logger"
|
"gitbase.de/apairon/mark2web/pkg/logger"
|
||||||
)
|
)
|
||||||
|
|
||||||
func handleCompression(filename string, content []byte) {
|
func handleCompression(filename string, content []byte) {
|
||||||
ThreadStart(func() {
|
jobm.Enqueue(jobm.Job{
|
||||||
|
Function: func() {
|
||||||
if _, ok := Config.Compress.Extensions[path.Ext(filename)]; ok {
|
if _, ok := Config.Compress.Extensions[path.Ext(filename)]; ok {
|
||||||
|
|
||||||
if Config.Compress.Brotli {
|
if Config.Compress.Brotli {
|
||||||
@ -57,6 +60,9 @@ func handleCompression(filename string, content []byte) {
|
|||||||
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
},
|
||||||
|
Description: "compress " + filename,
|
||||||
|
Category: "compress",
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -1,5 +1,7 @@
|
|||||||
package mark2web
|
package mark2web
|
||||||
|
|
||||||
|
import "gitbase.de/apairon/mark2web/pkg/jobm"
|
||||||
|
|
||||||
// Run will do a complete run of mark2web
|
// Run will do a complete run of mark2web
|
||||||
func Run(inDir, outDir string, defaultPathConfig *PathConfig) {
|
func Run(inDir, outDir string, defaultPathConfig *PathConfig) {
|
||||||
SetTemplateDir(inDir + "/templates")
|
SetTemplateDir(inDir + "/templates")
|
||||||
@ -12,5 +14,5 @@ func Run(inDir, outDir string, defaultPathConfig *PathConfig) {
|
|||||||
|
|
||||||
tree.WriteWebserverConfig()
|
tree.WriteWebserverConfig()
|
||||||
|
|
||||||
Wait()
|
jobm.Wait()
|
||||||
}
|
}
|
||||||
|
@ -1,60 +0,0 @@
|
|||||||
package mark2web
|
|
||||||
|
|
||||||
import (
|
|
||||||
"runtime"
|
|
||||||
"sync"
|
|
||||||
|
|
||||||
"gitbase.de/apairon/mark2web/pkg/logger"
|
|
||||||
)
|
|
||||||
|
|
||||||
var wg sync.WaitGroup
|
|
||||||
var numCPU = runtime.NumCPU()
|
|
||||||
var curNumThreads = 1 // main thread is 1
|
|
||||||
|
|
||||||
func init() {
|
|
||||||
logger.Log.Infof("number of CPU core: %d", numCPU)
|
|
||||||
}
|
|
||||||
|
|
||||||
// Wait will wait for all our internal go threads
|
|
||||||
func Wait() {
|
|
||||||
wg.Wait()
|
|
||||||
}
|
|
||||||
|
|
||||||
// ThreadSetup adds 1 to wait group
|
|
||||||
func ThreadSetup() {
|
|
||||||
curNumThreads++
|
|
||||||
wg.Add(1)
|
|
||||||
}
|
|
||||||
|
|
||||||
// ThreadDone removes 1 from wait group
|
|
||||||
func ThreadDone() {
|
|
||||||
curNumThreads--
|
|
||||||
wg.Done()
|
|
||||||
}
|
|
||||||
|
|
||||||
// ThreadStart will start a thread an manages the wait group
|
|
||||||
func ThreadStart(f func(), forceNewThread ...bool) {
|
|
||||||
force := false
|
|
||||||
if len(forceNewThread) > 0 && forceNewThread[0] {
|
|
||||||
force = true
|
|
||||||
}
|
|
||||||
|
|
||||||
if numCPU > curNumThreads || force {
|
|
||||||
// only new thread if empty CPU core available or forced
|
|
||||||
threadF := func() {
|
|
||||||
f()
|
|
||||||
ThreadDone()
|
|
||||||
}
|
|
||||||
|
|
||||||
ThreadSetup()
|
|
||||||
go threadF()
|
|
||||||
} else {
|
|
||||||
logger.Log.Debugf("no more CPU core (%d used), staying in main thread", curNumThreads)
|
|
||||||
f()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// SetNumCPU is for testing package without threading
|
|
||||||
func SetNumCPU(i int) {
|
|
||||||
numCPU = i
|
|
||||||
}
|
|
Loading…
Reference in New Issue
Block a user