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) Chunk(key, digest string) (ChunkCheckpoint, CheckpointDecision) { l.mu.Lock() defer l.mu.Unlock() return l.inner.Chunk(key, digest) } 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) ChunkRunning(key, digest string) error { return r.call(func() error { return r.inner.ChunkRunning(key, digest) }) } func (r *lockedCheckpointRecorder) ChunkSucceeded(key, digest string, chunks []source.Chunk, warnings []contracts.Warning) error { return r.call(func() error { return r.inner.ChunkSucceeded(key, digest, chunks, warnings) }) } func (r *lockedCheckpointRecorder) ChunkRejected(key, digest string, rejected contracts.RejectedOutput) error { return r.call(func() error { return r.inner.ChunkRejected(key, digest, rejected) }) } func (r *lockedCheckpointRecorder) ChunkFailed(key, digest string, err error) error { return r.call(func() error { return r.inner.ChunkFailed(key, digest, 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) }) }