README
¶
âïž Go Job Kit
Cloud Tasks ãžæå ¥ããGCS ãžææç©ãæžãåºãéåæãžã§ããæ±ããµãŒãã¹åãã®å ±éåºç€ã§ãã
ããžã§ããæå
¥ãã â é²è¡ç¶æ³ãèšé²ãã â å±¥æŽãããŒãžã³ã°ããŠäžèЧããããšããéªšæ Œã¯ãçæç©ã鳿¥œã§ããåç»ã§ããæŒ«ç»ã§ããåã圢ã«ãªããŸãã
ã«ããããããåãµãŒãã¹ãããããã® internal/ ã«å®è£
ãæ±ãããšãåãã³ãŒããå°ããã€é£ãéã£ããŸãŸå¢ããŠãããŸãã
æ¬ã©ã€ãã©ãªã¯ãã®å ±ééšåã ããæãåºãããã®ã§ãææç©ãã®ãã®ã®ãã¡ã€ã³ã«ã¯äžåèžã¿èŸŒãŸãªãããšãèšèšæ¹éãšããŠããŸãã
ðº å šäœå (Overview)
ãžã§ãã®äžçãšãããããã®å±é¢ãæ åœããããã±ãŒãžã®å¯Ÿå¿ã§ãã
| å±é¢ | èµ·ããããš | æ åœ |
|---|---|---|
| æå ¥ | HTTP ãã³ãã©ãŒã Cloud Tasks ãžç©ã¿ãqueued ãèšé²ãã |
jobstatus.Store |
| å®è¡ | ã¯ãŒã«ãŒã running â succeeded / failed ãèšé²ãããåé
ä¿¡ãããã¿ã¹ã¯ã¯ããã§æã¡åã |
jobstatus.Recorder |
| åç § | UIã»M2M ã¯ã©ã€ã¢ã³ããé²è¡ç¶æ³ãããŒãªã³ã°ãã | jobstatus.Store |
| äžèЧ | å±¥æŽç»é¢ããžã§ã ID ãéããããŒãžãåãåºããã¡ã¿ããŒã¿ãèªã | joblist + paging + cache |
ð ãããžã§ã¯ãã¬ã€ã¢ãŠã (Project Layout)
go-job-kit/
âââ jobstatus/ # ãžã§ãé²è¡ç¶æ³ã®èšé²ã»ååŸïŒStatus / Store / RecorderïŒ
âââ joblist/ # ã¹ãã¬ãŒãžã®ç䌌ãã£ã¬ã¯ããªèµ°æ»ãããžã§ã ID ãåéïŒCollectïŒ
âââ paging/ # ãžã§ã ID äžèŠ§ã® 1 å§ãŸãããŒãžã³ã°ïŒSelectIDs / LoadPage / PageMetaïŒ
âââ cache/ # ã€ã³ã¡ã¢ãªãã£ãã·ã¥ïŒTTL / IDListïŒ
ðŠ jobstatus â ãžã§ãç¶æ
ã®èšé²
çæã®æåŠããããŸã§ Slack éç¥ã«ããæ®ããã倱æãããžã§ãã UI ããå®å šã«æ¶ããŠããåé¡ãè§£æ¶ããããã®èšé²å±€ã§ããããããŠãCloud Tasks ã® at-least-once é ä¿¡ã«å¯Ÿããåå®è¡ã¬ãŒãã®æ ¹æ ã«ããªããŸãã
1. ç¶æ ã®åãå®çŸ©ãã
ææç©ã®ä¿åå
ã¯ãµãŒãã¹ããšã«åœ¢ãéãïŒåäž URIã»åºåãã£ã¬ã¯ããªã»è€æ° URIïŒãããå
±éãã£ãŒã«ãã ãã jobstatus.Status ãæã¡ãŸãããµãŒãã¹åºæã®ãã£ãŒã«ãã¯ããããåã蟌ãã æ§é äœã«è¶³ããŠãã ããã
type JobStatus struct {
jobstatus.Status
OutputDir string `json:"output_dir,omitempty"`
}
Go ã¯åãèŸŒã¿æ§é äœã JSON ã§ãã©ããã«å±éãããããä¿åãããåœ¢ã¯æ¬¡ã®ããã«ãªããŸããæ¢åã® status.json ããã®ãŸãŸèªã¿æžãã§ããã®ã¯ãã®ããã§ãã
{
"job_id": "c20260726-120000-abcd1234",
"command": "generate",
"state": "running",
"title": "äœåå",
"attempts": 2,
"queued_at": "2026-07-26T12:00:00Z",
"updated_at": "2026-07-26T12:00:31Z",
"output_dir": "gs://bucket/jobs/c20260726-120000-abcd1234"
}
å ¥ãåã®ãã€ããŒããå°å ¥ããªãã§ãã ããã åã蟌ã¿ããããæç¹ã§ JSON ã®åœ¢ãå€ãããä¿åæžã¿ã®ç¶æ ãã¡ã€ã«ãèªããªããªããŸãã
2. Store ã§èªã¿æžããã
// storage ã«ã¯ remoteio.Store ããã®ãŸãŸæž¡ããŸãã
store := jobstatus.NewStore[JobStatus](
storage,
jobstatus.UnderJobDir("gs://"+bucket+"/jobs"), // â .../jobs/{jobID}/status.json
)
// æå
¥çŽåŸã« queued ãèšé²ãããJobID ãš UpdatedAt 㯠Save ãæå»ããŸãã
// åã蟌ãã Status ã®ãã£ãŒã«ãã¯ãGo 1.27 ãããã®ãŸãŸè€åãªãã©ã«ãžæžããŸãã
err := store.Save(ctx, jobID, JobStatus{State: jobstatus.StateQueued, Command: "generate"})
// UIã»M2M ã¯ã©ã€ã¢ã³ãããé²è¡ç¶æ³ã远ã
status, err := store.Get(ctx, jobID)
switch {
case errors.Is(err, jobstatus.ErrNotFound):
// æªèšé²ãèšé²åã®æå
¥ãããã®æ©èœããåã«äœããããžã§ãã§ãèµ·ããæ£åžžç³» â 404
case errors.Is(err, jobstatus.ErrUnavailable):
// ååšããã¯ããªã®ã«èªããªãïŒæš©éäžè¶³ã»GCS é害ïŒâ 503ãåå 㯠err ã«å
ãŸããŠãã
case errors.Is(err, jobstatus.ErrInvalidJobID):
// URL ã«çŽã蟌ãã äžæ£ãª ID â 400ãå詊è¡ããŠãçŽããªãã®ã§ 5xx ã«ããªã
case err != nil:
// å£ãã JSON ãªã© â 500
}
ãæªèšé²ãïŒErrNotFoundïŒãšãããã¯ããªã®ã«èªããªããïŒErrUnavailableïŒãå¥ã®ãšã©ãŒã«ããŠããã®ã¯ãäž¡è
ã§åãã¹ãå€æãæ£å察ã ããã§ããèšé²ãç¡ãã®ã¯æ£åžžãªç¶æ
ãªã®ã§å
ãžé²ãã§ãããèªããªãã£ãã ãã®å Žåããç¡ãããšã¿ãªããšãå®äºæžã¿ã®ãžã§ããæªå®äºãšèª€èªããŠçæããŸãããšããçŽããŸããåãåãã«ã¯ remoteio ãæªååšã os.ErrNotExist ã«å
ãã§è¿ãããšã䜿ã£ãŠããŠã远å ã®ã¹ãã¬ãŒãžåŸåŸ©ã¯ãããŸããã
ã©ã®ãšã©ãŒãåå ãå
ãã ãŸãŸè¿ãã®ã§ãerrors.Is ã®å€å®ã¯ãã®ãŸãŸã«ããã°ã«ã¯å
ã®å€±æçç±ãæ®ããŸãã
èªã¿æžãã«ã¯ encoding/json/v2ïŒGo 1.27ïŒã䜿ããŸããJSON å€ã®ããšã«æ®ã£ããã€ãåãšãéè€ããããŒãæåŠããããã§ããåŸæ¥ã® json.Decoder ã¯ã©ã¡ããé»ã£ãŠåãåãããšãã«éè€ããŒã¯åŸåã¡ã ã£ãã®ã§ãéäžã§åããæžã蟌ã¿ã«æ¬¡ã®æžã蟌ã¿ãç¶ãã status.json ã {"state":"succeeded",âŠ,"state":"running"} ã®åœ¢ã«ãªããš running ãšããŠèªããŠããŸãããåå®è¡ã¬ãŒããé²ãã§ããã¯ãã®å·»ãæ»ãããèšé²ã§ã¯ãªãèªã¿åãã®åŽããå
¥ã£ãŠããçµè·¯ã§ããæžã蟌ã¿åŽã v2 ã«æããŠãããã & < > ã¯ãšã¹ã±ãŒããããŸããããJSON ãšããŠã¯åå€ã§ãv1 ãæžããæ¢åã®ãã¡ã€ã«ããã®ãŸãŸèªããŸãã
é
眮ãèªåã§æ±ºãããå Žå㯠Locator ãæž¡ããŸããUnderJobDir ã¯ãææç©ãšåããžã§ããã£ã¬ã¯ããªé
äžã«çœ®ãããšããæ¢å®ã®é
眮ã§ãå±¥æŽåé€ïŒãã¬ãã£ãã¯ã¹ã®äžæ¬åé€ïŒã§ç¶æ
ãã¡ã€ã«ãèªåçã«çä»ããŸãã
locate := func(jobID string) (string, error) {
return remoteio.BuildURI(remoteio.SchemeGCS, bucket, layout.JobStatusPath(jobID)), nil
}
ãžã§ã ID ã®æ£èŠå㯠Store ã®å
éšã§å¿
ãè¡ãããŸãïŒLocator ãåãåãã®ã¯æ£èŠåæžã¿ã® ID ã§ãïŒãåŒã³åºãåŽã§ jobid.Sanitize ãéãå¿
èŠã¯ãããŸããã
3. Recorder ã§ã¯ãŒã«ãŒããèšé²ãã
ã¯ãŒã«ãŒã¯ç¶æ
ãå€ãããã³ã«ã¿ã¹ã¯ããç¶æ
ãçµã¿ç«ãŠçŽããŸããçŽ æŽã«æžããšãå詊è¡ã®ãã³ã«è©Šè¡åæ°ãšæå
¥æå»ã倱ãããŸããRecorder ã¯ååã®èšé²ããå
±éãã£ãŒã«ããåŒãç¶ãã ããã§ä¿åããŸãã
rec := jobstatus.NewRecorder(store) // store ã nil ãªãèšé²ã¯è¡ãããªã
// ââ ã¯ãŒã«ãŒã®å
¥å£ïŒåå®è¡ã¬ãŒã + åŠçéå§ã®èšé² ââ
// Cloud Tasks 㯠at-least-once é
ä¿¡ã§ããéç¥ã®å€±æãªã©ã§ã¯ãŒã«ãŒããšã©ãŒãè¿ããš
// åãã¿ã¹ã¯ãåé
ä¿¡ãããçæã³ã¹ãããã®ãŸãŸäºéã«çºçããŸãã
//
// Attemptsã»QueuedAt ã¯ååã®èšé²ããåŒãç¶ãããŸãïŒapply ã¯ãã®åŸã«åŒã°ããŸãïŒã
// å€å®ãšåŒãç¶ãã¯åãèšé²ãèŠãã®ã§ãstatus.json ã®ååŸã¯ 1 åã ãã§ãã
done, err := rec.Begin(ctx, task.JobID, newStatus(task, jobstatus.StateRunning),
func(next, prev *JobStatus) {
next.Attempts++
if prev != nil {
next.OutputDir = prev.OutputDir // ãµãŒãã¹åºæãã£ãŒã«ãã®åŒãç¶ã
}
})
if err != nil {
return err // ç¶æ
ãèªããå€å®ã§ããªãããšã©ãŒãè¿ããŠåé
ä¿¡ã«å§ãã
}
if done {
return nil // å®äºæžã¿ã®åé
ä¿¡ãèšé²ãããªã
}
// ââ 以éã®èšé² ââ
rec.Record(ctx, task.JobID, newStatus(task, jobstatus.StateSucceeded))
Begin ã¯ãAlreadySucceeded ã§æã¡åããå€å®ããŠãã Record ã§ running ãæžã 2 段ã 1 ã€ã«ãŸãšãããã®ã§ããRecord ã¯åŒãç¶ãã®ããã«ååã®èšé²ãèªã¿çŽãã®ã§ã以åã¯åã status.json ã 1 ã¿ã¹ã¯ã§ 2 床ååŸããŠããŸãããèšé²ã䌎ããªãæã¡åãå€å®ã ããèŠããšã㯠AlreadySucceeded ã䜿ã£ãŠãã ããã
åå®è¡ã¬ãŒãã¯ãæªèšé²ïŒErrNotFoundïŒããå®äºããŠããªãããšã㊠false ã«ããŸããäžæ¹ãç¶æ
ãèªããªãã£ãå ŽåïŒErrUnavailableïŒã¯ãšã©ãŒãè¿ããŸãããæªå®äºãã«åããšå®äºæžã¿ã®ãžã§ããäœãçŽããŠã¬ãŒããé²ãã¯ãã®ã³ã¹ããèªåã§çºçããããå®äºæžã¿ãã«åããšæªå®äºã®ãžã§ããã¿ã¹ã¯ããš ACK ãããŠäºåºŠãšå®è¡ãããªãããã§ãã©ã¡ããžãåãããåŒã³åºãåŽãåé
ä¿¡ã«å§ããããããã«ããŠããŸãã
éåãã®åãããŒããå¡ãã§ãããŸããRecord ã¯ãå®äºæžã¿ã®èšé²ãž running / failed ãæžãããšãããšãã¯ä¿åããŸãããåé
ä¿¡ãããã¿ã¹ã¯ã¯ç¶æ
ãçµã¿ç«ãŠçŽãããããã®ãŸãŸæžããšå®äºãããžã§ããããŒãªã³ã°äžã®ç»é¢ã«ãåå®è¡ã¬ãŒãã«ãããŸã çµãã£ãŠããªãããšæ ãããã§ããqueued ã¯äŸå€ã§ããã®ãŸãŸä¿åããŸãã åé
ä¿¡ãããã¿ã¹ã¯ãæžãã®ã¯ running ã failed ã§ãããqueued ãæžãã®ã¯æ°ããäŸé Œã ããªã®ã§ãåããžã§ã ID ã§ã®äœãçŽããããã§æ¢ããŠããŸããªãããã§ãã
åŒãç¶ãã®èŠå㯠Status.CarryOver ã«ãŸãšãŸã£ãŠããŸãã
| ãã£ãŒã«ã | åŒãç¶ã | çç± |
|---|---|---|
Attempts / QueuedAt |
ãã | çµã¿ç«ãŠçŽãã®ãã³ã«å€±ãããªã |
Title / Command |
ä»åã空ã®ãšãã ã | çæã®éäžã§ç¢ºå®ããé¡ç®ããå€ãå€ã§äžæžãããªã |
State / Error / UpdatedAt |
ããªã | ãä»åã®èšé²ãã衚ãå€ãæååŸã«å€ã倱æçç±ãæ®ã£ãŠããŸã |
Recorder ãåãåãã®ã¯ StatusStore[T] ã€ã³ã¿ãŒãã§ãŒã¹ã§ã*jobstatus.Store[T] ã¯ãã®ãŸãŸæž¡ããŸãã
type StatusStore[T any] interface {
Get(ctx context.Context, jobID string) (T, error)
Save(ctx context.Context, jobID string, status T) error
}
ã¢ããªåŽã«ç¶æ
ä¿åã® port ã眮ããŠããå Žåã¯ããã®åœ¢ã«æããŠããããšãããããããŸããæããŠããã° *jobstatus.Store[T] ããã®ãŸãŸ port ã®å®è£
ã«ãªããéã«äœãæãŸãã«æžã¿ãŸãã
æããªãã£ãå Žåã«å¿
èŠã«ãªãã®ã¯èãã¢ããã¿ 1 ã€ã§ãããå©çšåŽãå¢ãããšãã®æ°ã ãåããã®ãå¢ããŸãïŒå®éããžã§ã ID ãç¶æ
ã«å«ãã圢 Save(ctx, status) ãä¿ã£ãŠãã 3 ãµãŒãã¹ããåãã¢ããã¿ãããããæã€ããšã«ãªããŸããïŒã
API äžèЧ
| çš®å¥ | ã·ã³ãã« | 説æ |
|---|---|---|
| å | State / StateQueued StateRunning StateSucceeded StateFailed |
ã©ã€ããµã€ã¯ã«äžã®ç¶æ |
| å | Status |
å ±éãã£ãŒã«ãããµãŒãã¹åºæã®åãžåã蟌ãã§äœ¿ã |
| ã¡ãœãã | Status.IsTerminal() |
ããŒãªã³ã°ãæ¢ããŠãããïŒsucceeded ã®ã¿ trueïŒ |
| ã¡ãœãã | Status.CarryOver(prev) / Status.Common() |
ååèšé²ããã®åŒãç¶ãïŒRecorder ã䜿ãïŒ |
| å | Store[T] / NewStore |
ä¿å (Save)ã»ååŸ (Get)ã»åé€ (Delete) |
| å | Locator / UnderJobDir |
ç¶æ ãã¡ã€ã«ã®é 眮 |
| å | Recorder[T] / NewRecorder |
ã¯ãŒã«ãŒå
¥å£ (Begin)ã»åå®è¡ã¬ãŒã (AlreadySucceeded)ã»åŒãç¶ãä»ãèšé² (Record) |
| å | StatusStore[T] |
Recorder ãèŠæ±ããèªã¿æžãã*Store[T] ã¯ãã®ãŸãŸæºãã |
| ãšã©ãŒ | ErrNotFound |
æªèšé²ïŒèšé²åã®æå ¥ã»ãã®æ©èœããåã®ãžã§ãïŒã404 ãž |
| ãšã©ãŒ | ErrInvalidJobID |
ãžã§ã ID ãæ£èŠåãéããªãïŒåŒã³åºãåŽã®å
¥åïŒã400 ãžãåå 㯠jobid.ErrEmpty ãªã©ãŸã§ errors.Is ã§èŸ¿ãã |
| ãšã©ãŒ | ErrUnavailable |
ååšããã¯ãã®ç¶æ ãèªããªãïŒæš©éäžè¶³ã»é害ïŒã503 ãžãåå®è¡ã¬ãŒãã¯ããããšã©ãŒãšããŠè¿ã |
ð joblist â ãžã§ã ID ã®åé
äžèЧã®å ¥å£ã§ããããžã§ã 1 ä»¶ = ãã¬ãã£ãã¯ã¹çŽäžã®ãã£ã¬ã¯ã㪠1 ã€ããšããé 眮ãããç䌌ãã£ã¬ã¯ããªåããžã§ã ID ãšããŠéããŸãã
jobIDs, err := jobIDCache.Load(ctx, prefix, func(ctx context.Context) ([]string, error) {
return joblist.Collect(ctx, storage, prefix) // storage ã«ã¯ remoteio.Store ããã®ãŸãŸæž¡ããŸã
})
åºåãæåãæå®ããŠèµ°æ»ããããããžã§ã 1 ä»¶ã 1 ãšã³ããªãšããŠåãåããŸããæå®ããªããšé äžã®ææç©ãå šä»¶è¿ãã1 ãžã§ãã«ã€ãææç©ã®æ°ã ãçµæãåãåã£ãããã§ãåŒã³åºãåŽã§éè€ã朰ãããšã«ãªããŸãããã¬ãã£ãã¯ã¹çŽäžã«çŽæ¥çœ®ããããªããžã§ã¯ãã¯ããžã§ãã§ã¯ãªããã察象å€ã§ãã
æŸããã®ã¯ãã£ã¬ã¯ããªåã ããªã®ã§ãã¡ã¿ããŒã¿ã®ä¿ååã«èœã¡ããžã§ãã ID ãšããŠã¯çŸããŸããäžèЧã«èŠãããé€ããã¯èªã¿èŸŒã¿åŽã®å€æã§ãïŒLoadPage ã®ãã©ãŒã«ããã¯ã®é
ãåç
§ïŒã
äœæ¥çšãžã§ãã®é€å€ã ID 圢åŒã®æ€èšŒã¯ããªãã·ã§ã³ã§å·®ã蟌ã¿ãŸãã
jobIDs, err := joblist.Collect(ctx, reader, prefix,
joblist.WithValidIDsOnly(), // jobid.Validate ãéã ID ã ã
joblist.WithKeep(func(id string) bool { return !strings.HasPrefix(id, "regen-") }),
)
è¿ã䞊ã³é ã¯ã¹ãã¬ãŒãžã®åæé ã®ãŸãŸã§ãããæ°ããé ããžã®äžŠã¹æ¿ããšããŒãžåãåºãã¯ã次㮠paging ãæ
ããŸãããã±ããå
šäœã®èµ°æ»ã«ãªããããåŒã³åºã㯠cache.IDList è¶ãã«è¡ãããšãå§ããŸãã
API äžèЧ
| çš®å¥ | ã·ã³ãã« | 説æ |
|---|---|---|
| 颿° | Collect(ctx, reader, prefix, opts...) |
ç䌌ãã£ã¬ã¯ããªåããžã§ã ID ãšããŠåéãã |
| ãªãã·ã§ã³ | WithKeep(fn) |
éãã ID ãçµã蟌ãïŒè€æ°æå®ã¯ ANDïŒ |
| ãªãã·ã§ã³ | WithValidIDsOnly() |
jobid.Validate ãéã ID ã ããéãã |
ð paging â å±¥æŽã®ããŒãžã³ã°
ãžã§ã ID ã®äžèŠ§ãæ°ããé ã«åãåºããç»é¢è¡šç€ºã«å¿
èŠãªã¡ã¿ããŒã¿ãçµã¿ç«ãŠãŸããããŒãžçªå·ã¯ 1 å§ãŸããperPage ã 0 以äžã®ãšãã¯ããŒãžã³ã°ããå
šä»¶ãè¿ããŸãã
ããŒãžãåãåºã
ids, meta := paging.SelectIDs(jobIDs, page, perPage, jobid.SortKey) // åŒæ°ã®ã¹ã©ã€ã¹ã¯å€æŽããŸãã
items := loadHistories(ctx, ids) // äžéšã®èªã¿èŸŒã¿å€±æã¯ã¹ããã
meta = paging.AdjustItemCount(meta, len(items)) // å®ä»¶æ°ã«åãã㊠From/To ãè£æ£
AdjustItemCount ãããã®ã¯ãã1ã10 ä»¶ç®ã衚瀺ããšåºããªãã 8 ä»¶ãã䞊ã°ãªãããšãããºã¬ãé²ãããã§ããäžèŠ§ã¯ ID ã䞊ã¹ãŠããã¡ã¿ããŒã¿æ¬äœãèªã¿ã«ãããããäžéšã®èªã¿èŸŒã¿ã«å€±æãããšè¡šç€ºä»¶æ°ã SelectIDs ã®æ³å®ããå°ãªããªããŸãã
PageMeta ã® JSON ã¿ã°ã¯ãç§»æ€å
ã®ãµãŒãã¹ãè¿ããŠããæ¢åã®ã¬ã¹ãã³ã¹ãšåã圢ã§ããç»é¢ãš M2M ã¯ã©ã€ã¢ã³ãã®åæ¹ãäŸåããŠããããã倿Žãããšãã¯å©çšåŽã®è¿œéãèŠããŸãã
type PageMeta struct {
Page int `json:"page"` // 1 å§ãŸããç¯å²å€ã¯æçµããŒãžãžäžžãããã
PerPage int `json:"per_page"`
Total int `json:"total"` // äžèЧå
šäœã®ä»¶æ°
TotalPages int `json:"total_pages"`
HasPrev bool `json:"has_prev"`
HasNext bool `json:"has_next"`
PrevPage int `json:"prev_page"`
NextPage int `json:"next_page"`
From int `json:"from"` // ãnãm ä»¶ç®ã衚瀺ãã® n
To int `json:"to"` // åãã m
}
䞊ã³é
äžŠã¹æ¿ãã®ããŒã¯åŒæ°ã§å¿ ãéžã³ãŸããæ¢å®ã眮ããšãééã£ãæ¢å®ã®ãŸãŸåŒã¹ãŠããŸãããã§ããããã§èª€ã£ãŠãäžèŠ§ã¯æ£åžžã«è¿ããé åºã ããéãã«åŽ©ããŸãã
nil ãæž¡ããš ID ãã®ãã®ã®éé ã«ãªããŸããããããæç³»åãšäžèŽããã®ã¯ ID ã®å
é ãçææå»ã§å§ãŸãå Žåââãã¬ãã£ãã¯ã¹ãä»ããªãããå
šä»¶ã§åäžã®ãã¬ãã£ãã¯ã¹ã䜿ãå Žåââã«éãããŸããçšéããšã«ç°ãªããã¬ãã£ãã¯ã¹ãæ··åšããããæ¡çªã®åœ¢åŒãéäžã§å€ãã£ãããããšãæå忝èŒã¯æå»ã§ã¯ãªããã¬ãã£ãã¯ã¹é ã«ãªããŸãïŒæå»éšåããåã«å·®ãåºãããïŒã
ãã®ããé垞㯠go-utils/jobid ã® SortKey ãæž¡ããŠãã ãããjobid.New ãçæãã圢åŒã«å ããŠåãµãŒãã¹ãç¬èªæ¡çªããŠãã圢åŒãèªãããããçºè¡å
ã®éã ID ãæ··åšããäžèЧã§ã䞊ã³é ã厩ããŸãããæå»ãåãåºããªã ID ã§ã¯ç©ºæåãè¿ããéé ã§ã¯æ«å°Ÿã«åããŸãã
ããŒãžåã䞊è¡ã«èªã¿èŸŒã
1 ããŒãžåã®ã¡ã¿ããŒã¿ååŸã¯ãçŽåã«ãããšã¹ãã¬ãŒãžèªã¿åããæ°åå䞊ã³ãŸããããšãã£ãŠäžŠè¡ã«ãããšé åºã厩ããŸããLoadPage ã¯ãã®çµã¿ç«ãŠïŒåãåºã â 䞊è¡èªã¿èŸŒã¿ â ä»¶æ°è£æ£ïŒããŸãšããŠè¡ããSelectIDs ãè¿ãã䞊ã³é ããã®ãŸãŸä¿ã¡ãŸãã
items, meta, err := paging.LoadPage(ctx, jobIDs, page, perPage, jobid.SortKey, repo.loadHistory,
paging.WithConcurrency(10), // æ¢å®ã¯ 10
)
èªã¿èŸŒã¿ã«å€±æãã ID ã¯èŠåãã°ãæ®ããŠäžèЧããåãé€ãããŸãã代ããã«ãžã§ã ID ã ãã®è¡ãæ®ãããå Žåã¯ãload ã®äžã§ãã©ãŒã«ããã¯å€ãè¿ããŠãã ããïŒã¡ãã»ãŒãžã代æ¿å€ã®åœ¢ã¯ãµãŒãã¹ããšã«éããããã©ã€ãã©ãªåŽã§ã¯æã¡ãŸããïŒã
load := func(ctx context.Context, jobID string) (History, error) {
h, err := repo.load(ctx, jobID)
if err != nil {
return History{JobID: jobID, Title: jobID}, nil // äžèЧã«ã¯æ®ã
}
return h, nil
}
API äžèЧ
| çš®å¥ | ã·ã³ãã« | 説æ |
|---|---|---|
| 颿° | SelectIDs(jobIDs, page, perPage, sortKey) |
æ°ããé ã«äžŠã¹æ¿ããŠããŒãžãåãåºã |
| 颿° | LoadPage(ctx, jobIDs, page, perPage, sortKey, load, opts...) |
åãåºã + 䞊è¡èªã¿èŸŒã¿ + ä»¶æ°è£æ£ |
| 颿° | AdjustItemCount(meta, itemCount) |
å®éã«èªããä»¶æ°ãž From/To ãè£æ£ |
| å | PageMeta |
ç»é¢ãããŒãžããŒã·ã§ã³ãæç»ããããã®ã¡ã¿ããŒã¿ |
| å | SortKeyFunc |
äžŠã¹æ¿ãããŒã®åãåºãæ¹ãåŒæ°ã§å¿
ãéžã¶ïŒnil 㯠ID ãã®ãã®ïŒ |
| ãªãã·ã§ã³ | WithConcurrency(n) / WithLogger(l) |
åæå®è¡æ°ã»ãã¬ãŒïŒLoadPage ã§ã®ã¿æå¹ïŒ |
â¡ cache â ã€ã³ã¡ã¢ãªãã£ãã·ã¥
å±¥æŽäžèЧã¯ãžã§ã ID ã䞊ã¹ãŠãã 1 ä»¶ãã€ã¡ã¿ããŒã¿ãèªã¿ã«ãããããããŒãžãããããã³ã«åããªããžã§ã¯ããèªã¿çŽããŸãããããé¿ããããã®èããã£ãã·ã¥ã§ãã
ã¡ã¿ããŒã¿ã®ãã£ãã·ã¥ (TTL)
ããŒã¯ãžã§ã ID ã§ãã"abc" ãš "dir/abc" ãå¥ãšã³ããªã«ãªããšæŽæ°ããå€ãèªãŸããªããã£ãã·ã¥ãã¹ã«ãªããããããŒã¯ã¹ãã¬ãŒãžãã¹ãšåãæ£èŠåãéããŠæããããŸãã
histories := cache.NewTTL[HistoryItem](10 * time.Minute)
defer histories.Close()
if h, ok := histories.Get(jobID); ok {
return h, nil
}
NewTTL ã¯æéåããšã³ããªã®ååãŸã§éå§ããŸããéå§ãå©çšåŽã®æé ã«ãããšãåŒã³å¿ãããã®ãŸãŸã¡ã¢ãªã®æ»çã«ãªãããã§ãã䜿ãçµãã£ãã Close ãåŒãã§ãã ããïŒäœåºŠåŒãã§ããè€æ°ã®ãŽã«ãŒãã³ããåæã«åŒãã§ãå®å
šã§ãïŒã
äžèŠ§èµ°æ»ã®ãã£ãã·ã¥ (IDList)
äžèŠ§èµ°æ»ãã®ãã®ïŒãã¬ãã£ãã¯ã¹é
äžå
šäœã® ListïŒã¯ããããç¡ããšå±¥æŽç»é¢ãéããã³ã«èµ°ããŸããä¿ææéã¯ã¡ã¿ããŒã¿æ¬äœãã倧å¹
ã«çãåããæ°ãã宿ãããžã§ããäžèЧã«çŸãããŸã§ã®é
å»¶ãæå°ã«ããŸããåé€ã»è¿œå 㯠Invalidate ã§å³åº§ã«åæ ãããŸãã
TTL ãåããç¬éã«åãããŒãžã®ãªã¯ãšã¹ããéãªã£ãŠããèµ°æ»ã¯ 1 åã ãå®è¡ãããå šå¡ããã®çµæãå ±æããŸããåŸ ã€ã®ã¯èªåã® ctx ã®ç¯å²ã ãã§ãèµ°æ»ãå§ããåŒã³åºãããã£ã³ã»ã«ãããŠããåŸ ã£ãŠããåŽã¯å·»ã蟌ãŸããã«ããçŽããŸãã
jobIDCache := cache.NewIDList(time.Minute)
jobIDs, err := jobIDCache.Load(ctx, prefix, r.collectJobIDs) // è¿ãã®ã¯åžžã«è€è£œ
...
jobIDCache.Invalidate(prefix)
TTL ãšéãããã¡ãã¯ååãéå§ãã Close ããããŸãããä¿æããã®ã¯äžèЧã®åäœããšã« 1 ãšã³ããªã ãã§ãåãããŒãžã®æžã蟌ã¿ã§äžæžããããŠããããã§ããéã«èšãã°ãããŒã«ã¯äžèЧã®åäœïŒèµ°æ»å¯Ÿè±¡ã®ãã¬ãã£ãã¯ã¹ãªã©ïŒã ãã䜿ã£ãŠãã ããã ãžã§ã ID ã®ããã«å¢ãç¶ããå€ãããŒã«ãããšãæéåãã®ãšã³ããªãååãããªããŸãŸæºãŸããŸãã
API äžèЧ
| çš®å¥ | ã·ã³ãã« | 説æ |
|---|---|---|
| å | TTL[T] / NewTTL(ttl) |
ãžã§ã ID ãããŒãšããã¡ã¿ããŒã¿ã®ãã£ãã·ã¥ |
| ã¡ãœãã | Get / Set / Delete / Len / Close |
ååŸã»ä¿åã»åé€ã»ä»¶æ°ã»ååã®åæ¢ |
| å | IDList / NewIDList(ttl) |
ãžã§ã ID äžèЧã®çæãã£ãã·ã¥ |
| ã¡ãœãã | Load(ctx, key, collect) / Invalidate(key) / Len |
èµ°æ»çµæã®ååŸã»ç Žæ£ã»ä»¶æ° |
| 宿° | DefaultTTLïŒ10 åïŒ/ DefaultIDListTTLïŒ1 åïŒ |
ttl ã« 0 以äžãæž¡ãããšãã®æ¢å®å€ |
ð§ èšèšäžã®çŽæ (Invariants)
ãã®ã©ã€ãã©ãªãåŒãåãããåŒã³åºãåŽãæèããªããŠããåæã§ãã
-
ãžã§ã ID ã¯å¿ ãæ£èŠåããŠãã䜿ã â ãžã§ã ID 㯠URL ãã¹ãšã¹ãã¬ãŒãžãã¹ã®åæ¹ã«çŸãããããæ€èšŒã¯ã»ãã¥ãªãã£å¢çãå ŒããŸããã¹ãã¬ãŒãžãã¹ã®çµã¿ç«ãŠããã£ãã·ã¥ããŒã®çæããã©ã€ãã©ãªå éšã§
go-utils/jobidã«ããæ£èŠåãéããŸããïŒç§»æ€å ã§ã¯ããã®æ£èŠåãéãç®æãšéããªãç®æãæ··åšããŠããŸããïŒ -
䞊ã³é ã¯ãæ°ããé ãã«çµ±äžãã â
SelectIDs/LoadPageã¯ãœãŒãããŒãåŒæ°ã§å¿ ãéžã°ããæ¢å®ã眮ããŸããã誀ã£ãŠãäžèŠ§ã¯æ£åžžã«è¿ããé åºã ããéãã«åŽ©ããããã§ãïŒã䞊ã³é ããåç §ïŒã -
ååã®éå§ã¯åŒã³åºãåŽã®æé ã«ããªã â
cache.TTLã¯çæãšåæã«æéåããšã³ããªã®ååãå§ããŸãã忢ã¯Closeã«éçŽããŠããŸãã -
ç¶æ ãã¡ã€ã«ã¯åžžã«ææ°ã® 1 äžä»£ã®ã¿ â ç¶æ ã¯äžæžãã§æŽæ°ããå±¥æŽã¯æ®ããŸãããCDNã»ãã©ãŠã¶ã«ãã£ãã·ã¥ãããªããã
no-storeãä»äžããŸãã -
ç¶æ ã®èšé²ã«å€±æããŠãçæã¯æ¢ããªã â
Recorderã¯ä¿åã®å€±æãèŠåãã°ã«çããåŒã³åºãåŽãžã¯äŒããŸãããç¶æ ã¯ãããŸã§èŠ³æž¬ã®ããã®èšé²ã§ãããæžããªãã£ãããšãçç±ã«çæãäžæããã»ãã害ã倧ããããã§ããäžæ¹ãç¶æ ãèªããªãã£ããšãã®åå®è¡ã¬ãŒãïŒBegin/AlreadySucceededïŒã¯ãæªå®äºãã«ããå®äºæžã¿ãã«ãåãããšã©ãŒãè¿ããåŒã³åºãåŽã Cloud Tasks ã®åé ä¿¡ã«å§ããããããã«ããŸãïŒãæªå®äºããšã¿ãªãã®ã¯æªèšé²ErrNotFoundã ãã§ãïŒã -
äžèЧã®äžéšãèªããªããŠãããŒãžå šäœã¯è¿ã â
LoadPageã¯èªã¿èŸŒã¿ã«å€±æãã ID ãäžèЧããåãé€ããPageMetaãå®ä»¶æ°ãžè£æ£ããŸãããã ãctxã®ãã£ã³ã»ã«ã»æéåãã¯ãšã©ãŒãšããŠè¿ããŸããèªã¿èŸŒã¿ãè»äžŠã¿å€±æããçµæã®ã0 ä»¶ãããæ£åžžãªç©ºäžèЧãšåãéããããªãããã§ãã
ð¥ åé²åºæº (What belongs here)
ããžã§ã管çããäœã§ãåãå ¥ããŠããŸãååã®ãããåé²ã®å¯åŠã¯ä»¥äžã§å€æããŸãã
- 2 ã€ä»¥äžã®ãµãŒãã¹ãã䜿ããã â åäžãµãŒãã¹ã§ãã䜿ããªããã®ã¯ããã®ãµãŒãã¹ã®
internal/ã«çœ®ããŠãã ããã - ãµãŒãã¹åºæã®ãã¡ã€ã³ãæã¡èŸŒãŸãªã â çæç©ãã®ãã®ã衚ãåã¯ããã«ã¯çœ®ããŸãããç¶æ ãš ID ã ããæ±ããŸãã
- ãžã§ãã®ã©ã€ããµã€ã¯ã«ã«é¢ãã â èªèšŒã»HTTPã»éç¥ã¯ãããã
gcp-kit/go-http-kit/go-notifyã®æ åœã§ãã
ð äž»ãªäŸåé¢ä¿ (Dependencies)
| ããã±ãŒãž | çšé |
|---|---|
| shouni/go-remote-io | ç¶æ ãã¡ã€ã«ã®èªã¿æžããšäžèŠ§èµ°æ»ïŒGCS / S3 / ããŒã«ã«ãééçã«æ±ãïŒ |
| shouni/go-utils | ãžã§ã ID ã®æ€èšŒã»æ£èŠå (jobid) |
| jellydator/ttlcache | TTL ä»ãã€ã³ã¡ã¢ãªãã£ãã·ã¥ |
| golang.org/x/sync | åæã®ãã£ãã·ã¥ãã¹ã§äžèŠ§èµ°æ»ã 1 åã«ãŸãšããïŒcache.IDListïŒ |
paging ã¯æšæºã©ã€ãã©ãªã®ã¿ã«äŸåããŸãã
ð ã©ã€ã»ã³ã¹ (License)
ãã®ãããžã§ã¯ã㯠MIT License ã®äžã§å ¬éãããŠããŸãã
Directories
¶
| Path | Synopsis |
|---|---|
|
Package cache ã¯ãå±¥æŽäžèЧã®ããã® TTL ä»ãã€ã³ã¡ã¢ãªãã£ãã·ã¥ãæäŸããŸãã
|
Package cache ã¯ãå±¥æŽäžèЧã®ããã® TTL ä»ãã€ã³ã¡ã¢ãªãã£ãã·ã¥ãæäŸããŸãã |
|
Package joblist ã¯ãã¹ãã¬ãŒãžã®ç䌌ãã£ã¬ã¯ããªèµ°æ»ãããžã§ã ID ã®äžèЧãéããŸãã
|
Package joblist ã¯ãã¹ãã¬ãŒãžã®ç䌌ãã£ã¬ã¯ããªèµ°æ»ãããžã§ã ID ã®äžèЧãéããŸãã |
|
Package jobstatus ã¯ãéåæãžã§ãã®é²è¡ç¶æ³ããªã¢ãŒãã¹ãã¬ãŒãžäžã® JSON ãšã㊠èªã¿æžãããŸãã
|
Package jobstatus ã¯ãéåæãžã§ãã®é²è¡ç¶æ³ããªã¢ãŒãã¹ãã¬ãŒãžäžã® JSON ãšã㊠èªã¿æžãããŸãã |
|
Package paging ã¯ããžã§ã ID ã®äžèŠ§ãæ°ããé ã«åãåºããç»é¢è¡šç€ºã«å¿
èŠãª ããŒãžã¡ã¿ããŒã¿ãçµã¿ç«ãŠãŸãã
|
Package paging ã¯ããžã§ã ID ã®äžèŠ§ãæ°ããé ã«åãåºããç»é¢è¡šç€ºã«å¿ èŠãª ããŒãžã¡ã¿ããŒã¿ãçµã¿ç«ãŠãŸãã |