implement scan and index queue
This commit is contained in:
parent
118cfd91de
commit
e0314dcf56
4 changed files with 81 additions and 33 deletions
|
|
@ -166,7 +166,7 @@ func indexHandler(ctx *Context) {
|
||||||
if fpath, ok = hasTrimPrefix(path, core.IndexPrefix); ok {
|
if fpath, ok = hasTrimPrefix(path, core.IndexPrefix); ok {
|
||||||
log = log.WithValues("file", fpath)
|
log = log.WithValues("file", fpath)
|
||||||
log.V(2).Info("started indexing file")
|
log.V(2).Info("started indexing file")
|
||||||
err = scanFile(ctx.Srv, fpath)
|
err = ctx.Srv.SQ.Add(fpath)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error(err, "error indexing file")
|
log.Error(err, "error indexing file")
|
||||||
}
|
}
|
||||||
|
|
@ -214,7 +214,7 @@ func indexHandler(ctx *Context) {
|
||||||
ctx.Srv.DB.IndexMutex.Set(uint64(len(files)))
|
ctx.Srv.DB.IndexMutex.Set(uint64(len(files)))
|
||||||
for _, f := range files {
|
for _, f := range files {
|
||||||
if !ctx.Srv.DB.IsIndexed(f) {
|
if !ctx.Srv.DB.IsIndexed(f) {
|
||||||
err := scanFile(ctx.Srv, f)
|
err := ctx.Srv.SQ.Add(f)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error(err, "error indexing file", "file", f)
|
log.Error(err, "error indexing file", "file", f)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -28,37 +28,6 @@ func getStats(i *bleve.Index) *core.Stats {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func scanFile(srv *HTTPServer, path string) (err error) {
|
|
||||||
err = srv.FS.AddScan(path)
|
|
||||||
if err != nil {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
go func() {
|
|
||||||
defer srv.FS.RemoveScan(path)
|
|
||||||
log := srv.Log.WithValues("file", path)
|
|
||||||
file, txt, err := srv.Scanner.Scan(path)
|
|
||||||
if err != nil {
|
|
||||||
log.Error(err, "error scanning file")
|
|
||||||
return
|
|
||||||
}
|
|
||||||
tags, err := srv.DB.GetAllRTags()
|
|
||||||
if err != nil {
|
|
||||||
log.Error(err, "error getting tags")
|
|
||||||
}
|
|
||||||
found := tags.Match(txt)
|
|
||||||
id, err := srv.DB.Index.Add(txt, found)
|
|
||||||
if err != nil {
|
|
||||||
log.Error(err, "error adding file to index")
|
|
||||||
return
|
|
||||||
}
|
|
||||||
err = srv.DB.NewFile(id, file, found)
|
|
||||||
if err != nil {
|
|
||||||
log.Error(err, "error adding file to DB")
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
// Parts taken from strings.TrimPrefix
|
// Parts taken from strings.TrimPrefix
|
||||||
func hasTrimPrefix(s, prefix string) (string, bool) {
|
func hasTrimPrefix(s, prefix string) (string, bool) {
|
||||||
if len(s) >= len(prefix) && s[0:len(prefix)] == prefix {
|
if len(s) >= len(prefix) && s[0:len(prefix)] == prefix {
|
||||||
|
|
|
||||||
71
pkg/server/queue.go
Normal file
71
pkg/server/queue.go
Normal file
|
|
@ -0,0 +1,71 @@
|
||||||
|
// Copyright (C) 2022 Marius Schellenberger
|
||||||
|
|
||||||
|
package server
|
||||||
|
|
||||||
|
import "context"
|
||||||
|
|
||||||
|
type ScanQueue struct {
|
||||||
|
ctx context.Context
|
||||||
|
srv *HTTPServer
|
||||||
|
pchan chan string
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewScanQueue(ctx context.Context, srv *HTTPServer, buffer int) *ScanQueue {
|
||||||
|
return &ScanQueue{
|
||||||
|
ctx: ctx,
|
||||||
|
srv: srv,
|
||||||
|
pchan: make(chan string, buffer),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (sq *ScanQueue) Add(path string) (err error) {
|
||||||
|
err = sq.srv.FS.AddScan(path)
|
||||||
|
if err != nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
go func() {
|
||||||
|
select {
|
||||||
|
case <-sq.ctx.Done():
|
||||||
|
return
|
||||||
|
case sq.pchan <- path:
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
func (sq *ScanQueue) Scan() {
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-sq.ctx.Done():
|
||||||
|
//close(sq.pchan)
|
||||||
|
return
|
||||||
|
case p := <-sq.pchan:
|
||||||
|
sq.scanFile(p)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (sq *ScanQueue) scanFile(path string) {
|
||||||
|
srv := sq.srv
|
||||||
|
defer srv.FS.RemoveScan(path)
|
||||||
|
log := srv.Log.WithValues("file", path)
|
||||||
|
file, txt, err := srv.Scanner.Scan(path)
|
||||||
|
if err != nil {
|
||||||
|
log.Error(err, "error scanning file")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
tags, err := srv.DB.GetAllRTags()
|
||||||
|
if err != nil {
|
||||||
|
log.Error(err, "error getting tags")
|
||||||
|
}
|
||||||
|
found := tags.Match(txt)
|
||||||
|
id, err := srv.DB.Index.Add(txt, found)
|
||||||
|
if err != nil {
|
||||||
|
log.Error(err, "error adding file to index")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
err = srv.DB.NewFile(id, file, found)
|
||||||
|
if err != nil {
|
||||||
|
log.Error(err, "error adding file to DB")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -39,21 +39,27 @@ type HTTPServer struct {
|
||||||
JWT *jwt.JWT
|
JWT *jwt.JWT
|
||||||
Scanner *scan.Scanner
|
Scanner *scan.Scanner
|
||||||
FS *fs.Filesystem
|
FS *fs.Filesystem
|
||||||
|
SQ *ScanQueue
|
||||||
Log logr.Logger
|
Log logr.Logger
|
||||||
plain logr.Logger
|
plain logr.Logger
|
||||||
|
|
||||||
|
cancel func()
|
||||||
|
|
||||||
templ map[string]*template.Template
|
templ map[string]*template.Template
|
||||||
res map[string][]byte
|
res map[string][]byte
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewHTTPServer(log logr.Logger, cfg *core.Config, l net.Listener, s *scan.Scanner) (srv *HTTPServer) {
|
func NewHTTPServer(log logr.Logger, cfg *core.Config, l net.Listener, s *scan.Scanner) (srv *HTTPServer) {
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
srv = &HTTPServer{
|
srv = &HTTPServer{
|
||||||
Config: cfg,
|
Config: cfg,
|
||||||
listener: l,
|
listener: l,
|
||||||
Scanner: s,
|
Scanner: s,
|
||||||
Log: log.WithName("server"),
|
Log: log.WithName("server"),
|
||||||
plain: log,
|
plain: log,
|
||||||
|
cancel: cancel,
|
||||||
}
|
}
|
||||||
|
srv.SQ = NewScanQueue(ctx, srv, 100)
|
||||||
srv.loadTemplates()
|
srv.loadTemplates()
|
||||||
srv.srv = &http.Server{}
|
srv.srv = &http.Server{}
|
||||||
return
|
return
|
||||||
|
|
@ -88,6 +94,7 @@ func (s *HTTPServer) Start() (err error) {
|
||||||
} else {
|
} else {
|
||||||
s.srv.Handler = s.handler
|
s.srv.Handler = s.handler
|
||||||
}
|
}
|
||||||
|
go s.SQ.Scan()
|
||||||
go func() {
|
go func() {
|
||||||
err = s.srv.Serve(s.listener)
|
err = s.srv.Serve(s.listener)
|
||||||
if err != nil && err != http.ErrServerClosed {
|
if err != nil && err != http.ErrServerClosed {
|
||||||
|
|
@ -98,6 +105,7 @@ func (s *HTTPServer) Start() (err error) {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *HTTPServer) Stop() error {
|
func (s *HTTPServer) Stop() error {
|
||||||
|
s.cancel()
|
||||||
s.DB.Close()
|
s.DB.Close()
|
||||||
s.JWT.Stop()
|
s.JWT.Stop()
|
||||||
ctx, cancel := context.WithTimeout(context.Background(), time.Second*5)
|
ctx, cancel := context.WithTimeout(context.Background(), time.Second*5)
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue