98 lines
2.6 KiB
Go
98 lines
2.6 KiB
Go
package appendorder
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
|
|
"gitea.maximumdirect.net/eric/notarius/internal/core/artifacts"
|
|
"gitea.maximumdirect.net/eric/notarius/internal/core/source"
|
|
"gitea.maximumdirect.net/eric/notarius/internal/framework/contracts"
|
|
"gitea.maximumdirect.net/eric/notarius/internal/framework/pipeline"
|
|
)
|
|
|
|
const Key = "appendorder"
|
|
|
|
var _ contracts.Merger = (*Merger)(nil)
|
|
|
|
type Merger struct{}
|
|
|
|
func New() *Merger {
|
|
return &Merger{}
|
|
}
|
|
|
|
func (m *Merger) Key() string {
|
|
return Key
|
|
}
|
|
|
|
func (m *Merger) Merge(ctx context.Context, req contracts.MergeRequest) (contracts.MergeResult, error) {
|
|
if m == nil {
|
|
return contracts.MergeResult{}, mergerErrorf("merger must not be nil")
|
|
}
|
|
if ctx == nil {
|
|
return contracts.MergeResult{}, mergerErrorf("context must not be nil")
|
|
}
|
|
if err := ctx.Err(); err != nil {
|
|
return contracts.MergeResult{}, mergerErrorf("context error before merge: %w", err)
|
|
}
|
|
|
|
var candidates []artifacts.ArtifactCandidate
|
|
for _, chunkArtifacts := range req.ChunkArtifacts {
|
|
candidates = append(candidates, cloneCandidates(chunkArtifacts.Candidates)...)
|
|
}
|
|
return contracts.MergeResult{Candidates: candidates}, nil
|
|
}
|
|
|
|
func ModuleSpec() pipeline.ModuleSpec {
|
|
return pipeline.ModuleSpec{
|
|
Key: Key,
|
|
Stage: pipeline.StageMerge,
|
|
Provides: []string{"merged"},
|
|
}
|
|
}
|
|
|
|
func Register(registry *pipeline.MergerRegistry) error {
|
|
return registry.RegisterWithSpec(ModuleSpec(), func() (contracts.Merger, error) {
|
|
return New(), nil
|
|
})
|
|
}
|
|
|
|
func cloneCandidates(candidates []artifacts.ArtifactCandidate) []artifacts.ArtifactCandidate {
|
|
if len(candidates) == 0 {
|
|
return nil
|
|
}
|
|
|
|
out := make([]artifacts.ArtifactCandidate, 0, len(candidates))
|
|
for _, candidate := range candidates {
|
|
out = append(out, cloneCandidate(candidate))
|
|
}
|
|
return out
|
|
}
|
|
|
|
func cloneCandidate(candidate artifacts.ArtifactCandidate) artifacts.ArtifactCandidate {
|
|
return artifacts.ArtifactCandidate{
|
|
Index: candidate.Index,
|
|
ExtractorKey: candidate.ExtractorKey,
|
|
ArtifactType: candidate.ArtifactType,
|
|
SchemaVersion: candidate.SchemaVersion,
|
|
Payload: append(json.RawMessage(nil), candidate.Payload...),
|
|
SourceRefs: append([]source.SourceRef(nil), candidate.SourceRefs...),
|
|
Metadata: cloneMetadata(candidate.Metadata),
|
|
}
|
|
}
|
|
|
|
func cloneMetadata(metadata map[string]any) map[string]any {
|
|
if len(metadata) == 0 {
|
|
return nil
|
|
}
|
|
out := make(map[string]any, len(metadata))
|
|
for key, value := range metadata {
|
|
out[key] = value
|
|
}
|
|
return out
|
|
}
|
|
|
|
func mergerErrorf(format string, args ...any) error {
|
|
return fmt.Errorf("appendorder merger: "+format, args...)
|
|
}
|