Files
2026-09-26 16:32:40 +03:00

602 lines
14 KiB
Go

package pipeline
import (
"crypto/sha256"
"encoding/hex"
"os"
"path/filepath"
"sort"
"strings"
"sync"
"time"
"aego/asset"
"aego/core/errs"
"aego/core/ids"
"aego/core/log"
"aego/project"
"aego/vfs"
)
const PipelineVersion = 1
const ServiceID = "asset.pipeline"
type Batch struct {
Imported []asset.Entry
Failed map[ids.AssetID]error
Removed []ids.AssetID
Retry []ids.AssetID
}
func (b *Batch) Empty() bool {
return len(b.Imported) == 0 && len(b.Failed) == 0 && len(b.Removed) == 0
}
type record struct {
id ids.AssetID
path string
meta *Meta
importer Importer
key string
group string
deps []ids.AssetID
missing bool
}
type Pipeline struct {
project *project.Project
imps *ImporterRegistry
log *log.ChannelLogger
watcher Watcher
index *asset.Index
artifacts vfs.FS
byID map[ids.AssetID]*record
byPath map[string]*record
groups map[string][]ids.AssetID
groupKey map[string]string
groupResults map[string]map[ids.AssetID]*Result
pages map[string]ids.AssetID
roots []string
mu sync.Mutex
}
func Open(p *project.Project, imps *ImporterRegistry, lg *log.Logger) (*Pipeline, error) {
if p == nil {
return nil, errs.New(errs.AssetNoProject, "pipeline needs a project")
}
if lg == nil {
lg = log.Nop()
}
cacheDir := p.CacheDir()
if err := os.MkdirAll(filepath.Join(cacheDir, "artifacts"), 0o755); err != nil {
return nil, errs.Wrap(err, errs.VFSIO, "cannot create the artifact cache")
}
files, err := vfs.New(vfs.Mount{Scheme: "cache", Root: cacheDir})
if err != nil {
return nil, err
}
pl := &Pipeline{
project: p,
imps: imps,
log: lg.With(log.Asset),
watcher: PollWatcher{},
index: asset.NewIndex(),
artifacts: files,
byID: map[ids.AssetID]*record{},
byPath: map[string]*record{},
groups: map[string][]ids.AssetID{},
groupKey: map[string]string{},
groupResults: map[string]map[ids.AssetID]*Result{},
pages: map[string]ids.AssetID{},
roots: []string{p.AssetsDir(), p.ScenesDir()},
}
if data, err := os.ReadFile(pl.indexPath()); err == nil {
if idx, err := asset.ReadIndex(data, "index.aetext"); err == nil {
pl.index = idx
} else {
pl.log.Warn("asset index is unreadable, rebuilding", log.F("err", err.Error()))
}
}
return pl, nil
}
func (p *Pipeline) SetWatcher(w Watcher) {
if w != nil {
p.watcher = w
}
}
func (p *Pipeline) Index() *asset.Index {
return p.index
}
func (p *Pipeline) Artifacts() vfs.FS {
return p.artifacts
}
func (p *Pipeline) indexPath() string {
return filepath.Join(p.project.CacheDir(), "index.aetext")
}
func (p *Pipeline) Lookup(path string) (ids.AssetID, bool) {
r, ok := p.byPath[filepath.ToSlash(path)]
if !ok {
return ids.AssetID{}, false
}
return r.id, true
}
func (p *Pipeline) PathOf(id ids.AssetID) (string, bool) {
r, ok := p.byID[id]
if !ok {
return "", false
}
return r.path, true
}
func (p *Pipeline) Dependents(id ids.AssetID) []ids.AssetID {
var out []ids.AssetID
for _, r := range p.byID {
for _, d := range r.deps {
if d == id {
out = append(out, r.id)
break
}
}
}
return out
}
func (p *Pipeline) Scan() (added, removed []ids.AssetID, err error) {
seen := map[string]bool{}
for _, root := range p.roots {
if err := p.scanRoot(root, seen, &added); err != nil {
return nil, nil, err
}
}
for path, r := range p.byPath {
if seen[path] {
continue
}
delete(p.byPath, path)
delete(p.byID, r.id)
removed = append(removed, r.id)
if e, ok := p.index.Lookup(r.id); ok {
e.Missing = true
}
p.log.Warn("asset source disappeared", log.F("path", path))
}
return added, removed, nil
}
func (p *Pipeline) scanRoot(root string, seen map[string]bool, added *[]ids.AssetID) error {
if _, err := os.Stat(root); os.IsNotExist(err) {
return nil
}
return filepath.WalkDir(root, func(path string, d os.DirEntry, err error) error {
if err != nil {
return nil
}
if d.IsDir() {
name := d.Name()
if path != root && (name == ".aego" || name[0] == '.') {
return filepath.SkipDir
}
return nil
}
if strings.HasSuffix(path, MetaSuffix) {
return nil
}
rel := p.rel(path)
imp, ok := p.imps.ForPath(path)
if !ok {
return nil
}
seen[rel] = true
if existing, ok := p.byPath[rel]; ok {
existing.importer = imp
return nil
}
meta, created, err := p.loadOrCreateMeta(path, imp)
if err != nil {
p.log.Error("cannot read the meta file", log.F("path", rel), log.F("err", err.Error()))
return nil
}
if other, dup := p.byID[meta.ID]; dup {
p.log.Error("duplicate asset id, regenerating",
log.F("path", rel), log.F("other", other.path))
meta.ID = ids.NewAssetID()
created = true
}
if created {
if err := os.WriteFile(path+MetaSuffix, WriteMeta(meta), 0o644); err != nil {
return nil
}
}
r := &record{id: meta.ID, path: rel, meta: meta, importer: imp}
p.byPath[rel] = r
p.byID[meta.ID] = r
*added = append(*added, meta.ID)
return nil
})
}
func (p *Pipeline) rel(path string) string {
r, err := filepath.Rel(p.project.Root, path)
if err != nil {
return filepath.ToSlash(path)
}
return filepath.ToSlash(r)
}
func (p *Pipeline) abs(rel string) string {
return filepath.Join(p.project.Root, filepath.FromSlash(rel))
}
func (p *Pipeline) loadOrCreateMeta(path string, imp Importer) (*Meta, bool, error) {
metaPath := path + MetaSuffix
data, err := os.ReadFile(metaPath)
if err == nil {
m, err := ParseMeta(data, p.rel(metaPath))
if err != nil {
return nil, false, err
}
if m.Importer == "" {
m.Importer = imp.ID()
}
for key, def := range imp.Defaults() {
if _, ok := m.Settings[key]; !ok {
m.Settings[key] = def
}
}
return m, false, nil
}
if !os.IsNotExist(err) {
return nil, false, err
}
return NewMeta(imp.ID(), imp.Version(), imp.Defaults()), true, nil
}
func (p *Pipeline) ImportAll() (*Batch, error) {
all := make([]ids.AssetID, 0, len(p.byID))
for id := range p.byID {
all = append(all, id)
}
sort.Slice(all, func(i, j int) bool {
return p.byID[all[i]].path < p.byID[all[j]].path
})
return p.Import(all...)
}
func (p *Pipeline) Import(list ...ids.AssetID) (*Batch, error) {
batch := &Batch{Failed: map[ids.AssetID]error{}}
groupsTouched := map[string]bool{}
for _, id := range list {
r, ok := p.byID[id]
if !ok {
continue
}
res, key, err := p.importOne(r)
if err != nil {
if errs.Has(err, errs.AssetSourceUnstable) {
batch.Retry = append(batch.Retry, id)
continue
}
batch.Failed[id] = err
p.log.Error("import failed", log.F("path", r.path), log.F("err", err.Error()))
continue
}
if res == nil {
continue
}
r.key = key
r.group = res.Group
r.deps = res.Deps
if res.Group != "" {
p.addToGroup(res.Group, id, res)
groupsTouched[res.Group] = true
continue
}
entry, err := p.writeArtifacts(r, res)
if err != nil {
batch.Failed[id] = err
continue
}
p.index.Put(entry)
batch.Imported = append(batch.Imported, entry)
}
for group := range groupsTouched {
entries, err := p.finalizeGroup(group)
if err != nil {
p.log.Error("atlas group failed", log.F("group", group), log.F("err", err.Error()))
for _, id := range p.groups[group] {
batch.Failed[id] = err
}
continue
}
batch.Imported = append(batch.Imported, entries...)
}
p.index.Sort()
if err := os.WriteFile(p.indexPath(), asset.WriteIndex(p.index), 0o644); err != nil {
return batch, errs.Wrap(err, errs.VFSIO, "cannot write the asset index")
}
return batch, nil
}
func (p *Pipeline) importOne(r *record) (*Result, string, error) {
full := p.abs(r.path)
before, err := os.Stat(full)
if err != nil {
return nil, "", errs.Wrap(err, errs.VFSIO, "cannot stat the asset source")
}
source, err := os.ReadFile(full)
if err != nil {
return nil, "", errs.Wrap(err, errs.VFSIO, "cannot read the asset source")
}
after, err := os.Stat(full)
if err != nil {
return nil, "", errs.Wrap(err, errs.VFSIO, "cannot stat the asset source")
}
if len(source) == 0 ||
int64(len(source)) != after.Size() ||
before.Size() != after.Size() ||
!before.ModTime().Equal(after.ModTime()) {
return nil, "", errs.New(errs.AssetSourceUnstable, "asset source is still being written",
log.F("path", r.path))
}
key := p.cacheKey(r, source)
if key == r.key && r.group == "" {
if _, ok := p.index.Lookup(r.id); ok {
return nil, key, nil
}
}
ctx := &ImportContext{
ID: r.id,
Path: r.path,
Source: source,
Meta: r.meta,
Log: p.log,
Lookup: p.Lookup,
}
res, err := r.importer.Import(ctx)
if err != nil {
return nil, "", err
}
return res, key, nil
}
func (p *Pipeline) cacheKey(r *record, source []byte) string {
h := sha256.New()
h.Write(source)
h.Write(canonicalSettings(r.meta.Settings))
writeUint(h, uint64(r.importer.Version()))
writeUint(h, PipelineVersion)
h.Write([]byte(r.importer.ID()))
for _, dep := range r.deps {
if e, ok := p.index.Lookup(dep); ok {
for _, a := range e.Artifacts {
h.Write([]byte(a.Hash))
}
}
}
return hex.EncodeToString(h.Sum(nil))
}
func writeUint(h interface{ Write([]byte) (int, error) }, v uint64) {
var b [8]byte
for i := 0; i < 8; i++ {
b[i] = byte(v >> (i * 8))
}
h.Write(b[:])
}
func (p *Pipeline) addToGroup(group string, id ids.AssetID, res *Result) {
members := p.groups[group]
for _, existing := range members {
if existing == id {
p.groupResults[group][id] = res
return
}
}
p.groups[group] = append(members, id)
if p.groupResults == nil {
p.groupResults = map[string]map[ids.AssetID]*Result{}
}
if p.groupResults[group] == nil {
p.groupResults[group] = map[ids.AssetID]*Result{}
}
p.groupResults[group][id] = res
}
func (p *Pipeline) writeArtifacts(r *record, res *Result) (asset.Entry, error) {
entry := asset.Entry{
ID: r.id,
Kind: res.Kind,
Name: assetName(r.path),
Source: r.path,
Deps: res.Deps,
}
for _, a := range res.Artifacts {
name := a.Name
if name == "" {
name = r.id.String() + artifactExt(a.Kind)
}
path := "cache://artifacts/" + name
if err := p.artifacts.WriteFile(path, a.Data); err != nil {
return asset.Entry{}, err
}
sum := sha256.Sum256(a.Data)
entry.Artifacts = append(entry.Artifacts, asset.Artifact{
Kind: a.Kind,
Path: path,
Hash: hex.EncodeToString(sum[:8]),
Size: int64(len(a.Data)),
})
}
return entry, nil
}
func artifactExt(kind asset.ArtifactKind) string {
switch kind {
case asset.ArtTexture, asset.ArtAtlasPage:
return ".aetex"
case asset.ArtSpriteRegion:
return ".aereg"
case asset.ArtScene:
return ".aescene"
case asset.ArtPrefab:
return ".aeprefab"
}
return ".bin"
}
func assetName(path string) string {
name := strings.TrimSuffix(filepath.Base(path), filepath.Ext(path))
dir := filepath.ToSlash(filepath.Dir(path))
if dir == "." {
return name
}
return dir + "/" + name
}
const maxRetries = 5
func (p *Pipeline) Watch(quiet time.Duration, fn func(*Batch)) (func(), error) {
if quiet <= 0 {
quiet = 150 * time.Millisecond
}
interval := quiet / 2
if interval < 50*time.Millisecond {
interval = 50 * time.Millisecond
}
pending := map[string]ChangeKind{}
attempts := map[string]int{}
var timer *time.Timer
var schedule func()
fire := func() {
p.mu.Lock()
paths := make(map[string]ChangeKind, len(pending))
for k, v := range pending {
paths[k] = v
}
pending = map[string]ChangeKind{}
p.mu.Unlock()
batch := p.applyChanges(paths)
if batch == nil {
return
}
if len(batch.Retry) > 0 {
p.mu.Lock()
for _, id := range batch.Retry {
r, ok := p.byID[id]
if !ok || attempts[r.path] >= maxRetries {
continue
}
attempts[r.path]++
pending[r.path] = ChangeWrite
}
p.mu.Unlock()
schedule()
}
for _, e := range batch.Imported {
p.mu.Lock()
delete(attempts, e.Source)
p.mu.Unlock()
}
if !batch.Empty() {
fn(batch)
}
}
schedule = func() {
p.mu.Lock()
defer p.mu.Unlock()
if len(pending) == 0 {
return
}
if timer != nil {
timer.Stop()
}
timer = time.AfterFunc(quiet, fire)
}
stop, err := p.watcher.Start(p.project.Root, []string{".aego", ".git", "exports", "native"}, interval, func(changes []Change) {
p.mu.Lock()
for _, c := range changes {
path := strings.TrimSuffix(c.Path, MetaSuffix)
pending[path] = c.Kind
}
p.mu.Unlock()
schedule()
})
if err != nil {
return nil, err
}
return func() {
stop()
p.mu.Lock()
if timer != nil {
timer.Stop()
}
p.mu.Unlock()
}, nil
}
func (p *Pipeline) applyChanges(paths map[string]ChangeKind) *Batch {
if _, _, err := p.Scan(); err != nil {
p.log.Error("rescan failed", log.F("err", err.Error()))
return nil
}
dirty := map[ids.AssetID]bool{}
var removed []ids.AssetID
for path, kind := range paths {
r, ok := p.byPath[path]
if !ok {
if kind == ChangeRemove {
continue
}
continue
}
if kind == ChangeRemove {
removed = append(removed, r.id)
continue
}
dirty[r.id] = true
for _, dep := range p.Dependents(r.id) {
dirty[dep] = true
}
if r.group != "" {
for _, member := range p.groups[r.group] {
dirty[member] = true
}
}
}
if len(dirty) == 0 && len(removed) == 0 {
return nil
}
list := make([]ids.AssetID, 0, len(dirty))
for id := range dirty {
list = append(list, id)
}
sort.Slice(list, func(i, j int) bool {
return p.byID[list[i]].path < p.byID[list[j]].path
})
batch, err := p.Import(list...)
if err != nil {
p.log.Error("import batch failed", log.F("err", err.Error()))
}
if batch != nil {
batch.Removed = removed
}
return batch
}
func (p *Pipeline) Close() {
_ = os.WriteFile(p.indexPath(), asset.WriteIndex(p.index), 0o644)
}