From c804940ec5bcdcfa0ffe90e03c6866d3bb651416 Mon Sep 17 00:00:00 2001
From: zhangmeng <775834166@qq.com>
Date: 星期二, 14 一月 2020 17:13:21 +0800
Subject: [PATCH] add chan cache for async

---
 run.go |  266 ++++++++++++++++++++++++++++++++++------------------
 1 files changed, 174 insertions(+), 92 deletions(-)

diff --git a/run.go b/run.go
index e1f3cbd..bba0c82 100644
--- a/run.go
+++ b/run.go
@@ -2,6 +2,7 @@
 
 import (
 	"context"
+	"os"
 	"sync"
 	"time"
 
@@ -12,9 +13,10 @@
 	"basic.com/valib/gogpu.git"
 )
 
+const maxTryBeforeReboot = 10
+
 type face struct {
 	handle *SDKFace
-	list   *sdkhelper.LockList
 
 	maxChannel      int
 	ftrackChans     map[string]chan protomsg.SdkMessage
@@ -30,6 +32,40 @@
 	ipc2Rule            string
 	ruleMsgMaxCacheSize int
 	reserved            map[string]interface{}
+
+	running     bool
+	rebootUntil int
+	mtxRunning  sync.Mutex
+}
+
+func (f *face) maybeReboot(ctx context.Context) {
+	for {
+		select {
+		case <-ctx.Done():
+			return
+		default:
+			f.mtxRunning.Lock()
+			running := f.running
+			f.mtxRunning.Unlock()
+
+			if running {
+				f.rebootUntil = 0
+
+				f.mtxRunning.Lock()
+				f.running = false
+				f.mtxRunning.Unlock()
+
+			} else {
+				f.rebootUntil++
+				f.fnLogger("Face No Running: ", f.rebootUntil)
+				if f.rebootUntil > maxTryBeforeReboot {
+					f.fnLogger("Face Too Long Running, Reboot")
+					os.Exit(0)
+				}
+			}
+			time.Sleep(time.Second)
+		}
+	}
 }
 
 // Create create sdk
@@ -101,7 +137,6 @@
 
 	return &face{
 		handle: handle,
-		list:   sdkhelper.NewLockList(maxChan + maxChan/2),
 
 		maxChannel:      maxChan,
 		ftrackChans:     make(map[string]chan protomsg.SdkMessage, maxChan),
@@ -116,6 +151,9 @@
 		ipc2Rule:            ipc2Rule,
 		ruleMsgMaxCacheSize: ruleMaxSize,
 		reserved:            reserved,
+
+		running:     true,
+		rebootUntil: maxTryBeforeReboot,
 	}
 }
 
@@ -129,101 +167,28 @@
 func Run(ctx context.Context, i interface{}) {
 	s := i.(*face)
 
-	chRcv, chSnd := sdkhelper.FlowCreate(ctx, s.id, s.shm, s.ipc2Rule, s.ruleMsgMaxCacheSize, s.fnLogger)
-	go sdkhelper.FlowBatch(ctx, chRcv, chSnd, s.typ, s.list.Push, s.list.Drain, s.run, s.release, s.fnLogger)
-}
+	const (
+		postPull = `_1`
+		postPush = `_2`
+	)
+	ipcRcv := sdkhelper.GetIpcAddress(s.shm, s.id+postPull)
+	ipcSnd := sdkhelper.GetIpcAddress(s.shm, s.id+postPush)
+	chRcv := make(chan []byte, s.maxChannel)
+	chSnd := make(chan sdkstruct.MsgSDK, s.maxChannel)
 
-func (f *face) run(msgs []protomsg.SdkMessage, out chan<- sdkstruct.MsgSDK, typ string) {
+	rcver := sdkhelper.NewReciever(ipcRcv, chRcv, s.shm, s.fnLogger)
+	snder := sdkhelper.NewSender(ipcSnd, chSnd, s.shm, s.fnLogger)
+	torule := sdkhelper.NewToRule(s.ipc2Rule, s.ruleMsgMaxCacheSize, s.fnLogger)
 
-	for _, msg := range msgs {
-		if _, ok := f.ftrackChans[msg.Cid]; ok {
-			f.ftrackChans[msg.Cid] <- msg
-		} else {
+	snder.ApplyCallbackFunc(torule.Push)
 
-			f.ftrackChans[msg.Cid] = make(chan protomsg.SdkMessage, cacheFrameNum)
-			chn := f.getAvailableChn()
-			if chn < 0 {
-				f.fnLogger("TOO MUCH CHANNEL")
-				sdkhelper.EjectResult(nil, msg, out)
-				continue
-			}
-			f.ftrackChannels[msg.Cid] = chn
+	go rcver.Run(ctx)
+	go snder.Run(ctx)
+	go torule.Run(ctx)
 
-			i := sdkhelper.UnpackImage(msg, f.typ, f.fnLogger)
-			if i == nil {
-				sdkhelper.EjectResult(nil, msg, out)
-				continue
-			}
-			// conv to bgr24 and resize
-			imgW, imgH := int(i.Width), int(i.Height)
-			ret := f.handle.TrackerResize(imgW, imgH, chn)
-			f.fnLogger("ResizeFaceTracker: cid: ", msg.Cid, " chan: ", chn, " wXh: ", imgW, "x", imgH, " result:", ret)
-			go f.detectTrackOneChn(f.ftrackChans[msg.Cid], out, chn)
-			f.ftrackChans[msg.Cid] <- msg
-		}
-	}
-}
+	go s.run(ctx, chRcv, chSnd)
 
-func (f *face) detectTrackOneChn(in <-chan protomsg.SdkMessage, out chan<- sdkstruct.MsgSDK, dtchn int) {
-	tm := time.Now()
-	sc := 0
-	f.fnLogger("DETECTTRACKONECHN DTCHN: ", dtchn)
-	var curCid string
-
-	for {
-		select {
-		case rMsg := <-in:
-
-			if !sdkhelper.ValidRemoteMessage(rMsg, f.typ, f.fnLogger) {
-				sdkhelper.EjectResult(nil, rMsg, out)
-				f.fnLogger("Face!!!!!!SkdMessage Invalid")
-
-				continue
-			}
-
-			i := sdkhelper.UnpackImage(rMsg, f.typ, f.fnLogger)
-			if i == nil || i.Data == nil || i.Width <= 0 || i.Height <= 0 {
-				sdkhelper.EjectResult(nil, rMsg, out)
-				f.fnLogger("Face!!!!!!Unpack Image From SkdMessage Failed")
-
-				continue
-			}
-
-			curCid = i.Cid
-
-			// conv to bgr24 and resize
-			imgW, imgH := int(i.Width), int(i.Height)
-
-			count, data, _ := f.handle.Run(i.Data, imgW, imgH, 3, dtchn)
-
-			sdkhelper.EjectResult(data, rMsg, out)
-			var id, name string
-			if rMsg.Tasklab != nil {
-				id, name = rMsg.Tasklab.Taskid, rMsg.Tasklab.Taskname
-			}
-			f.fnLogger("CAMERAID: ", rMsg.Cid, " TASKID: ", id, " TASKNAME: ", name, " DETECT ", f.typ, " COUNT: ", count)
-
-			sc++
-			if sc == 25 {
-				f.fnLogger("CHAN:%d, FACE RUN 25 FRAME USE TIME: ", dtchn, time.Since(tm))
-				sc = 0
-				tm = time.Now()
-			}
-
-			if time.Since(tm) > time.Second {
-				f.fnLogger("CHAN: ", dtchn, " FACE RUN ", sc, " FRAME USE TIME: ", time.Since(tm))
-				sc = 0
-				tm = time.Now()
-			}
-		case <-time.After(trackChnTimeout * time.Second):
-			f.fnLogger("Timeout to get image, curCid:", curCid)
-			if curCid != "" {
-				delete(f.ftrackChans, curCid)
-				f.releaseChn(dtchn)
-			}
-			return
-		}
-	}
+	go s.maybeReboot(ctx)
 }
 
 //////////////////////////////////////////////////////////////////
@@ -258,3 +223,120 @@
 	f.ftrackChanStats[chn] = false
 	f.chnLock.Unlock()
 }
+
+func (f *face) run(ctx context.Context, in <-chan []byte, out chan<- sdkstruct.MsgSDK) {
+
+	chMsg := make(chan protomsg.SdkMessage)
+	go sdkhelper.UnserilizeProto(ctx, in, chMsg, f.fnLogger)
+
+	for {
+		select {
+		case <-ctx.Done():
+			f.handle.Free()
+			return
+		case rMsg := <-chMsg:
+			if !sdkhelper.ValidRemoteMessage(rMsg, f.typ, f.fnLogger) {
+				f.fnLogger("FACE TRACK VALIDREMOTEMESSAGE INVALID")
+				sdkhelper.EjectResult(nil, rMsg, out)
+				continue
+			}
+
+			if _, ok := f.ftrackChans[rMsg.Cid]; ok {
+				f.fnLogger("Face Cache Size: ", len(f.ftrackChans))
+				f.ftrackChans[rMsg.Cid] <- rMsg
+			} else {
+
+				f.ftrackChans[rMsg.Cid] = make(chan protomsg.SdkMessage, cacheFrameNum)
+				chn := f.getAvailableChn()
+				if chn < 0 {
+					f.fnLogger("TOO MUCH CHANNEL")
+					sdkhelper.EjectResult(nil, rMsg, out)
+					continue
+				}
+				f.ftrackChannels[rMsg.Cid] = chn
+
+				i := sdkhelper.UnpackImage(rMsg, f.typ, f.fnLogger)
+				if i == nil {
+					sdkhelper.EjectResult(nil, rMsg, out)
+					continue
+				}
+				// conv to bgr24 and resize
+				imgW, imgH := int(i.Width), int(i.Height)
+				ret := f.handle.TrackerResize(imgW, imgH, chn)
+				f.fnLogger("ResizeFaceTracker: cid: ", rMsg.Cid, " chan: ", chn, " wXh: ", imgW, "x", imgH, " result:", ret)
+				go f.detectTrackOneChn(ctx, f.ftrackChans[rMsg.Cid], out, chn)
+				f.ftrackChans[rMsg.Cid] <- rMsg
+			}
+		default:
+			time.Sleep(time.Millisecond * 100)
+		}
+	}
+}
+
+func (f *face) detectTrackOneChn(ctx context.Context, in <-chan protomsg.SdkMessage, out chan<- sdkstruct.MsgSDK, dtchn int) {
+	tm := time.Now()
+	sc := 0
+	f.fnLogger("DETECTTRACKONECHN DTCHN: ", dtchn)
+	var curCid string
+
+	for {
+		select {
+		case <-ctx.Done():
+			return
+
+		case rMsg := <-in:
+
+			if !sdkhelper.ValidRemoteMessage(rMsg, f.typ, f.fnLogger) {
+				sdkhelper.EjectResult(nil, rMsg, out)
+				continue
+			}
+
+			i := sdkhelper.UnpackImage(rMsg, f.typ, f.fnLogger)
+			if i == nil || i.Data == nil || i.Width <= 0 || i.Height <= 0 {
+				sdkhelper.EjectResult(nil, rMsg, out)
+				continue
+			}
+
+			curCid = i.Cid
+
+			// conv to bgr24 and resize
+			imgW, imgH := int(i.Width), int(i.Height)
+
+			f.fnLogger("Face Start Run:", dtchn, "CAMERAID: ", rMsg.Cid)
+
+			count, data, _ := f.handle.Run(i.Data, imgW, imgH, 3, dtchn)
+
+			sdkhelper.EjectResult(data, rMsg, out)
+			f.mtxRunning.Lock()
+			f.running = true
+			f.mtxRunning.Unlock()
+
+			var id, name string
+			if rMsg.Tasklab != nil {
+				id, name = rMsg.Tasklab.Taskid, rMsg.Tasklab.Taskname
+			}
+			f.fnLogger("Chan:", dtchn, "CAMERAID: ", rMsg.Cid, " TASKID: ", id, " TASKNAME: ", name, " DETECT ", f.typ, " COUNT: ", count)
+
+			sc++
+			if sc == 25 {
+				f.fnLogger("CHAN:%d, FACE RUN 25 FRAME USE TIME: ", dtchn, time.Since(tm))
+				sc = 0
+				tm = time.Now()
+			}
+
+			if time.Since(tm) > time.Second {
+				f.fnLogger("CHAN: ", dtchn, " FACE RUN ", sc, " FRAME USE TIME: ", time.Since(tm))
+				sc = 0
+				tm = time.Now()
+			}
+		case <-time.After(trackChnTimeout * time.Second):
+			f.fnLogger("Timeout to get image, curCid:", curCid)
+			if curCid != "" {
+				delete(f.ftrackChans, curCid)
+				f.releaseChn(dtchn)
+			}
+			return
+
+		}
+	}
+}

--
Gitblit v1.8.0