From 0720a0ff2680107079bd6e3830f270b7ffc37d23 Mon Sep 17 00:00:00 2001 From: LeonardoTrapani Date: Tue, 12 Aug 2025 01:33:37 +0200 Subject: [PATCH] change pipeline handling --- internal/daemon/daemon.go | 170 +++++++++++++++++++++------------- internal/notify/notify.go | 32 +++++-- internal/pipeline/pipeline.go | 88 ++++++++++++++++++ 3 files changed, 218 insertions(+), 72 deletions(-) create mode 100644 internal/pipeline/pipeline.go diff --git a/internal/daemon/daemon.go b/internal/daemon/daemon.go index 9229931..538e0fb 100644 --- a/internal/daemon/daemon.go +++ b/internal/daemon/daemon.go @@ -13,24 +13,29 @@ import ( "github.com/leonardotrapani/hyprvoice/internal/bus" "github.com/leonardotrapani/hyprvoice/internal/notify" + "github.com/leonardotrapani/hyprvoice/internal/pipeline" ) -type Status string +type Status = pipeline.Status const ( - Idle Status = "idle" - Recording Status = "recording" - Transcribing Status = "transcribing" - Injecting Status = "injecting" - Completed Status = "completed" + Idle = pipeline.Idle + Recording = pipeline.Recording + Transcribing = pipeline.Transcribing + Injecting = pipeline.Injecting ) type Daemon struct { - mu sync.Mutex + mu sync.RWMutex status Status notifier notify.Notifier - ctx context.Context - cancel context.CancelFunc + + ctx context.Context + cancel context.CancelFunc + + pipeline pipeline.Pipeline + pipelineCancel context.CancelFunc + statusCh <-chan Status } func New(n notify.Notifier) *Daemon { @@ -38,20 +43,57 @@ func New(n notify.Notifier) *Daemon { n = notify.Desktop{} } ctx, cancel := context.WithCancel(context.Background()) - return &Daemon{ + d := &Daemon{ notifier: n, ctx: ctx, cancel: cancel, status: Idle, } + + return d } func (d *Daemon) Status() Status { - d.mu.Lock() - defer d.mu.Unlock() + d.mu.RLock() + defer d.mu.RUnlock() return d.status } +func (d *Daemon) startStatusReader(ctx context.Context, statusCh <-chan Status) { + go func() { + defer func() { + // Always clean up when this goroutine exits + d.mu.Lock() + d.status = Idle + d.statusCh = nil + d.pipeline = nil + d.pipelineCancel = nil + d.mu.Unlock() + }() + + for { + select { + case status, ok := <-statusCh: + if !ok { + return // Channel closed + } + + d.mu.Lock() + oldStatus := d.status + d.status = status + d.mu.Unlock() + + if oldStatus != status { + log.Printf("Status changed: %s -> %s", oldStatus, status) + } + + case <-ctx.Done(): + return // Context cancelled + } + } + }() +} + func (d *Daemon) Run() error { if err := bus.CheckExistingDaemon(); err != nil { return err @@ -80,36 +122,18 @@ func (d *Daemon) Run() error { log.Printf("Daemon started, listening on socket") - // Accept connections in a goroutine - connCh := make(chan net.Conn) - errCh := make(chan error) - - go func() { - for { - c, err := ln.Accept() - if err != nil { - errCh <- err - return - } - connCh <- c - } - }() - for { - select { - case <-d.ctx.Done(): - log.Printf("Shutdown requested, exiting") - return nil - case c := <-connCh: - go d.handle(c) - case err := <-errCh: - // If context is cancelled, this is expected + c, err := ln.Accept() + if err != nil { if d.ctx.Err() != nil { + log.Printf("Shutdown requested") return nil } log.Printf("Accept error: %v", err) return fmt.Errorf("accept failed: %w", err) } + + go d.handle(c) } } @@ -129,35 +153,14 @@ func (d *Daemon) handle(c net.Conn) { cmd := line[0] switch cmd { - case 't': // toggle - d.mu.Lock() - defer d.mu.Unlock() - - switch d.status { - case Idle: - d.status = Recording - - d.notifier.RecordingChanged(true) - log.Printf("Recording toggled: true") - fmt.Fprintf(c, "STATUS recording=%s\n", d.status) - default: - d.status = Idle - - // TODO: trigger transcription - - d.notifier.RecordingChanged(false) - log.Printf("Recording toggled: false") - fmt.Fprintf(c, "STATUS recording=%s\n", d.status) - } - case 's': // status - d.mu.Lock() - status := d.status - d.mu.Unlock() - + case 't': + d.toggle() + case 's': + status := d.Status() fmt.Fprintf(c, "STATUS recording=%s\n", status) - case 'v': // protocol version + case 'v': fmt.Fprintf(c, "STATUS proto=%s\n", bus.ProtoVer) - case 'q': // quit daemon + case 'q': log.Printf("Shutdown requested") fmt.Fprint(c, "OK quitting\n") d.cancel() @@ -166,3 +169,46 @@ func (d *Daemon) handle(c net.Conn) { fmt.Fprintf(c, "ERR unknown=%q\n", cmd) } } + +func (d *Daemon) toggle() { + d.mu.Lock() + defer d.mu.Unlock() + + var notification func() + + switch d.status { + case Idle: + ctx, cancel := context.WithCancel(d.ctx) + p := pipeline.New() + d.pipeline = p + d.pipelineCancel = cancel + d.statusCh = p.Run(ctx) + notification = d.notifier.RecordingStarted + + // Start status reader for this pipeline + d.startStatusReader(ctx, d.statusCh) + + case Recording: + if d.pipelineCancel != nil { + d.pipelineCancel() // Context cleanup handles the rest + } + notification = d.notifier.RecordingEnded + + case Transcribing: + if d.pipeline != nil { + d.pipeline.Inject() + } + // No notification for injection start + + case Injecting: + if d.pipelineCancel != nil { + d.pipelineCancel() // Context cleanup handles the rest + } + notification = d.notifier.RecordingEnded + } + + // Send notification after releasing lock + if notification != nil { + go notification() + } +} diff --git a/internal/notify/notify.go b/internal/notify/notify.go index bee7b87..2f98f93 100644 --- a/internal/notify/notify.go +++ b/internal/notify/notify.go @@ -1,25 +1,35 @@ package notify import ( - "fmt" "log" "os/exec" ) type Notifier interface { - RecordingChanged(on bool) + RecordingStarted() + RecordingEnded() + Transcribing() Error(msg string) } type Desktop struct{} -func (Desktop) RecordingChanged(on bool) { - state := "Stopped" - if on { - state = "Started" +func (Desktop) RecordingStarted() { + cmd := exec.Command("notify-send", "-a", "Hyprvoice", "Hyprvoice: Recording Started") + if err := cmd.Run(); err != nil { + log.Printf("Failed to send notification: %v", err) } - cmd := exec.Command("notify-send", "-a", "Hyprvoice", - fmt.Sprintf("Hyprvoice: %s Recording", state)) +} + +func (Desktop) RecordingEnded() { + cmd := exec.Command("notify-send", "-a", "Hyprvoice", "Hyprvoice: Recording Ended") + if err := cmd.Run(); err != nil { + log.Printf("Failed to send notification: %v", err) + } +} + +func (Desktop) Transcribing() { + cmd := exec.Command("notify-send", "-a", "Hyprvoice", "Hyprvoice: Transcribing...") if err := cmd.Run(); err != nil { log.Printf("Failed to send notification: %v", err) } @@ -36,5 +46,7 @@ func (Desktop) Error(msg string) { // Useful in unit tests or headless builds. type Nop struct{} -func (Nop) RecordingChanged(on bool) {} -func (Nop) Error(msg string) {} +func (Nop) RecordingStarted() {} +func (Nop) RecordingEnded() {} +func (Nop) Transcribing() {} +func (Nop) Error(msg string) {} diff --git a/internal/pipeline/pipeline.go b/internal/pipeline/pipeline.go new file mode 100644 index 0000000..e05cb26 --- /dev/null +++ b/internal/pipeline/pipeline.go @@ -0,0 +1,88 @@ +package pipeline + +import ( + "context" + "log" + "time" +) + +type Status string + +const ( + Idle Status = "idle" + Recording Status = "recording" + Transcribing Status = "transcribing" + Injecting Status = "injecting" +) + +type Pipeline interface { + Run(ctx context.Context) <-chan Status + Inject() +} + +type pipeline struct { + injectCh chan struct{} +} + +func New() Pipeline { + return &pipeline{ + injectCh: make(chan struct{}, 1), + } +} + +func (p *pipeline) Run(ctx context.Context) <-chan Status { + statusCh := make(chan Status, 1) + go p.run(ctx, statusCh) + return statusCh +} + + +func (p *pipeline) Inject() { + select { + case p.injectCh <- struct{}{}: + default: + } +} + +func (p *pipeline) run(ctx context.Context, statusCh chan<- Status) { + defer close(statusCh) + + // Start recording + log.Printf("Pipeline: Starting recording") + statusCh <- Recording + + // Recording phase + select { + case <-time.After(2 * time.Second): + log.Printf("Pipeline: Recording complete, waiting for injection") + statusCh <- Transcribing + case <-ctx.Done(): + log.Printf("Pipeline: Stopped during recording") + return + } + + // Wait for injection or timeout + select { + case <-p.injectCh: + log.Printf("Pipeline: Injection started") + statusCh <- Injecting + + // Injection work + select { + case <-time.After(1 * time.Second): + log.Printf("Pipeline: Injection complete") + statusCh <- Idle // Instead of Completed + case <-ctx.Done(): + log.Printf("Pipeline: Stopped during injection") + return + } + + case <-time.After(10 * time.Second): // use context timeout + log.Printf("Pipeline: Auto-timeout, completing") + statusCh <- Idle // Instead of Completed + + case <-ctx.Done(): + log.Printf("Pipeline: Stopped during transcription wait") + return + } +}