138 lines
5.5 KiB
Go
138 lines
5.5 KiB
Go
package pipeline
|
|
|
|
import (
|
|
"sync"
|
|
|
|
"gitea.maximumdirect.net/eric/notarius/internal/core/source"
|
|
"gitea.maximumdirect.net/eric/notarius/internal/framework/contracts"
|
|
)
|
|
|
|
type lockedDebugRecorder struct {
|
|
inner DebugRecorder
|
|
mu sync.Mutex
|
|
}
|
|
|
|
func synchronizedDebugRecorder(inner DebugRecorder) DebugRecorder {
|
|
if _, ok := inner.(*lockedDebugRecorder); ok {
|
|
return inner
|
|
}
|
|
return &lockedDebugRecorder{inner: inner}
|
|
}
|
|
|
|
// SynchronizedDebugRecorder serializes access when a recorder is shared by
|
|
// pipeline operations and construction-injected LLM clients.
|
|
func SynchronizedDebugRecorder(inner DebugRecorder) DebugRecorder {
|
|
return synchronizedDebugRecorder(inner)
|
|
}
|
|
func (r *lockedDebugRecorder) Enabled() bool {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
return r.inner.Enabled()
|
|
}
|
|
func (r *lockedDebugRecorder) WriteJSON(name string, payload any) error {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
return r.inner.WriteJSON(name, payload)
|
|
}
|
|
func (r *lockedDebugRecorder) WriteBytes(name string, data []byte) error {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
return r.inner.WriteBytes(name, data)
|
|
}
|
|
|
|
type lockedCheckpointLoader struct {
|
|
inner CheckpointLoader
|
|
mu sync.Mutex
|
|
}
|
|
|
|
func synchronizedCheckpointLoader(inner CheckpointLoader) CheckpointLoader {
|
|
if _, ok := inner.(*lockedCheckpointLoader); ok {
|
|
return inner
|
|
}
|
|
return &lockedCheckpointLoader{inner: inner}
|
|
}
|
|
func (l *lockedCheckpointLoader) Enabled() bool {
|
|
l.mu.Lock()
|
|
defer l.mu.Unlock()
|
|
return l.inner.Enabled()
|
|
}
|
|
func (l *lockedCheckpointLoader) Source(key string) (SourceCheckpoint, CheckpointDecision) {
|
|
l.mu.Lock()
|
|
defer l.mu.Unlock()
|
|
return l.inner.Source(key)
|
|
}
|
|
func (l *lockedCheckpointLoader) Extract(lane, key string, deps []CheckpointFingerprint) (ExtractCheckpoint, CheckpointDecision) {
|
|
l.mu.Lock()
|
|
defer l.mu.Unlock()
|
|
return l.inner.Extract(lane, key, deps)
|
|
}
|
|
func (l *lockedCheckpointLoader) Merge(lane, key string, deps []CheckpointFingerprint) (MergeCheckpoint, CheckpointDecision) {
|
|
l.mu.Lock()
|
|
defer l.mu.Unlock()
|
|
return l.inner.Merge(lane, key, deps)
|
|
}
|
|
func (l *lockedCheckpointLoader) Normalize(lane, key string, deps []CheckpointFingerprint) (NormalizeCheckpoint, CheckpointDecision) {
|
|
l.mu.Lock()
|
|
defer l.mu.Unlock()
|
|
return l.inner.Normalize(lane, key, deps)
|
|
}
|
|
|
|
type lockedCheckpointRecorder struct {
|
|
inner CheckpointRecorder
|
|
mu sync.Mutex
|
|
}
|
|
|
|
func synchronizedCheckpointRecorder(inner CheckpointRecorder) CheckpointRecorder {
|
|
if _, ok := inner.(*lockedCheckpointRecorder); ok {
|
|
return inner
|
|
}
|
|
return &lockedCheckpointRecorder{inner: inner}
|
|
}
|
|
func (r *lockedCheckpointRecorder) call(fn func() error) error {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
return fn()
|
|
}
|
|
func (r *lockedCheckpointRecorder) SourceRunning(key string) error {
|
|
return r.call(func() error { return r.inner.SourceRunning(key) })
|
|
}
|
|
func (r *lockedCheckpointRecorder) SourceSucceeded(key string, doc *source.SourceDocument) error {
|
|
return r.call(func() error { return r.inner.SourceSucceeded(key, doc) })
|
|
}
|
|
func (r *lockedCheckpointRecorder) SourceFailed(key string, err error) error {
|
|
return r.call(func() error { return r.inner.SourceFailed(key, err) })
|
|
}
|
|
func (r *lockedCheckpointRecorder) ExtractRunning(lane, key string, deps []CheckpointFingerprint) error {
|
|
return r.call(func() error { return r.inner.ExtractRunning(lane, key, deps) })
|
|
}
|
|
func (r *lockedCheckpointRecorder) ExtractSucceeded(lane, key string, deps []CheckpointFingerprint, outputs []CheckpointArtifact, rejected []contracts.RejectedOutput, warnings []contracts.Warning) error {
|
|
return r.call(func() error { return r.inner.ExtractSucceeded(lane, key, deps, outputs, rejected, warnings) })
|
|
}
|
|
func (r *lockedCheckpointRecorder) ExtractFailed(lane, key string, deps []CheckpointFingerprint, err error) error {
|
|
return r.call(func() error { return r.inner.ExtractFailed(lane, key, deps, err) })
|
|
}
|
|
func (r *lockedCheckpointRecorder) MergeRunning(lane, key string, deps []CheckpointFingerprint) error {
|
|
return r.call(func() error { return r.inner.MergeRunning(lane, key, deps) })
|
|
}
|
|
func (r *lockedCheckpointRecorder) MergeSucceeded(lane, key string, deps []CheckpointFingerprint, output CheckpointArtifact, warnings []contracts.Warning) error {
|
|
return r.call(func() error { return r.inner.MergeSucceeded(lane, key, deps, output, warnings) })
|
|
}
|
|
func (r *lockedCheckpointRecorder) MergeRejected(lane, key string, deps []CheckpointFingerprint, rejected contracts.RejectedOutput) error {
|
|
return r.call(func() error { return r.inner.MergeRejected(lane, key, deps, rejected) })
|
|
}
|
|
func (r *lockedCheckpointRecorder) MergeFailed(lane, key string, deps []CheckpointFingerprint, err error) error {
|
|
return r.call(func() error { return r.inner.MergeFailed(lane, key, deps, err) })
|
|
}
|
|
func (r *lockedCheckpointRecorder) NormalizeRunning(lane, key string, deps []CheckpointFingerprint) error {
|
|
return r.call(func() error { return r.inner.NormalizeRunning(lane, key, deps) })
|
|
}
|
|
func (r *lockedCheckpointRecorder) NormalizeSucceeded(lane, key string, deps []CheckpointFingerprint, output CheckpointArtifact, warnings []contracts.Warning) error {
|
|
return r.call(func() error { return r.inner.NormalizeSucceeded(lane, key, deps, output, warnings) })
|
|
}
|
|
func (r *lockedCheckpointRecorder) NormalizeRejected(lane, key string, deps []CheckpointFingerprint, rejected contracts.RejectedOutput) error {
|
|
return r.call(func() error { return r.inner.NormalizeRejected(lane, key, deps, rejected) })
|
|
}
|
|
func (r *lockedCheckpointRecorder) NormalizeFailed(lane, key string, deps []CheckpointFingerprint, err error) error {
|
|
return r.call(func() error { return r.inner.NormalizeFailed(lane, key, deps, err) })
|
|
}
|