Add teabot daemon implementation
teabot watches Gitea repos and dispatches one-shot Claude Code sessions in Docker containers to work issues and review PRs, acting as configurable bot personalities. Claude-Session: https://claude.ai/code/session_015ur3i7D2azsMAWTSVABApv
This commit is contained in:
@@ -0,0 +1,219 @@
|
||||
// Package state persists which Gitea events teabot has already handled so a
|
||||
// restart does not re-trigger work. State is a single JSON file under the
|
||||
// configured state directory (default ~/.local/state/teabot/state.json).
|
||||
package state
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// StateFileName is the JSON file holding processed-event bookkeeping.
|
||||
const StateFileName = "state.json"
|
||||
|
||||
// RepoState tracks what has been handled for one repository.
|
||||
type RepoState struct {
|
||||
// ProcessedIssues holds issue indexes already dispatched.
|
||||
ProcessedIssues map[int64]bool `json:"processed_issues"`
|
||||
// ProcessedPulls holds pull-request indexes already dispatched.
|
||||
ProcessedPulls map[int64]bool `json:"processed_pulls"`
|
||||
// ProcessedComments holds comment IDs already dispatched.
|
||||
ProcessedComments map[int64]bool `json:"processed_comments"`
|
||||
// ActedIssues/ActedPulls record which issues/PRs teabot ran a session
|
||||
// for, so comment follow-ups only fire on threads the bot engaged with.
|
||||
ActedIssues map[int64]bool `json:"acted_issues"`
|
||||
ActedPulls map[int64]bool `json:"acted_pulls"`
|
||||
// Seeded is set the first time a repo is polled: existing open issues/PRs
|
||||
// and recent comments are recorded as processed WITHOUT dispatching, so a
|
||||
// fresh install does not stampede every open item.
|
||||
Seeded bool `json:"seeded"`
|
||||
// LastPoll is the time of the last completed poll, used to bound `since`
|
||||
// queries on subsequent cycles.
|
||||
LastPoll time.Time `json:"last_poll"`
|
||||
}
|
||||
|
||||
func newRepoState() *RepoState {
|
||||
return &RepoState{
|
||||
ProcessedIssues: map[int64]bool{},
|
||||
ProcessedPulls: map[int64]bool{},
|
||||
ProcessedComments: map[int64]bool{},
|
||||
ActedIssues: map[int64]bool{},
|
||||
ActedPulls: map[int64]bool{},
|
||||
}
|
||||
}
|
||||
|
||||
// data is the on-disk document.
|
||||
type data struct {
|
||||
Repos map[string]*RepoState `json:"repos"`
|
||||
}
|
||||
|
||||
// Store is a thread-safe, file-backed processed-event tracker.
|
||||
type Store struct {
|
||||
path string
|
||||
mu sync.Mutex
|
||||
d *data
|
||||
}
|
||||
|
||||
// New loads the store from dir, creating an empty one if the file is absent.
|
||||
func New(dir string) (*Store, error) {
|
||||
s := &Store{
|
||||
path: filepath.Join(dir, StateFileName),
|
||||
d: &data{Repos: map[string]*RepoState{}},
|
||||
}
|
||||
raw, err := os.ReadFile(s.path)
|
||||
if err != nil {
|
||||
if os.IsNotExist(err) {
|
||||
return s, nil
|
||||
}
|
||||
return nil, fmt.Errorf("reading state %s: %w", s.path, err)
|
||||
}
|
||||
if len(raw) == 0 {
|
||||
return s, nil
|
||||
}
|
||||
if err := json.Unmarshal(raw, s.d); err != nil {
|
||||
return nil, fmt.Errorf("parsing state %s: %w", s.path, err)
|
||||
}
|
||||
if s.d.Repos == nil {
|
||||
s.d.Repos = map[string]*RepoState{}
|
||||
}
|
||||
return s, nil
|
||||
}
|
||||
|
||||
// repo returns the RepoState for repo, creating it if needed. Caller holds mu.
|
||||
func (s *Store) repo(repo string) *RepoState {
|
||||
rs := s.d.Repos[repo]
|
||||
if rs == nil {
|
||||
rs = newRepoState()
|
||||
s.d.Repos[repo] = rs
|
||||
}
|
||||
// Guard against a partially-populated document loaded from disk.
|
||||
if rs.ProcessedIssues == nil {
|
||||
rs.ProcessedIssues = map[int64]bool{}
|
||||
}
|
||||
if rs.ProcessedPulls == nil {
|
||||
rs.ProcessedPulls = map[int64]bool{}
|
||||
}
|
||||
if rs.ProcessedComments == nil {
|
||||
rs.ProcessedComments = map[int64]bool{}
|
||||
}
|
||||
if rs.ActedIssues == nil {
|
||||
rs.ActedIssues = map[int64]bool{}
|
||||
}
|
||||
if rs.ActedPulls == nil {
|
||||
rs.ActedPulls = map[int64]bool{}
|
||||
}
|
||||
return rs
|
||||
}
|
||||
|
||||
// IssueProcessed reports whether an issue index was already handled.
|
||||
func (s *Store) IssueProcessed(repo string, index int64) bool {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return s.repo(repo).ProcessedIssues[index]
|
||||
}
|
||||
|
||||
// PullProcessed reports whether a PR index was already handled.
|
||||
func (s *Store) PullProcessed(repo string, index int64) bool {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return s.repo(repo).ProcessedPulls[index]
|
||||
}
|
||||
|
||||
// CommentProcessed reports whether a comment ID was already handled.
|
||||
func (s *Store) CommentProcessed(repo string, id int64) bool {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return s.repo(repo).ProcessedComments[id]
|
||||
}
|
||||
|
||||
// ActedOnIssue reports whether teabot ran a session for an issue.
|
||||
func (s *Store) ActedOnIssue(repo string, index int64) bool {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return s.repo(repo).ActedIssues[index]
|
||||
}
|
||||
|
||||
// ActedOnPull reports whether teabot ran a session for a PR.
|
||||
func (s *Store) ActedOnPull(repo string, index int64) bool {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return s.repo(repo).ActedPulls[index]
|
||||
}
|
||||
|
||||
// MarkIssue records an issue index as processed and acted-on.
|
||||
func (s *Store) MarkIssue(repo string, index int64) {
|
||||
s.mu.Lock()
|
||||
rs := s.repo(repo)
|
||||
rs.ProcessedIssues[index] = true
|
||||
rs.ActedIssues[index] = true
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
// MarkPull records a PR index as processed and acted-on.
|
||||
func (s *Store) MarkPull(repo string, index int64) {
|
||||
s.mu.Lock()
|
||||
rs := s.repo(repo)
|
||||
rs.ProcessedPulls[index] = true
|
||||
rs.ActedPulls[index] = true
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
// MarkComment records a comment ID as processed.
|
||||
func (s *Store) MarkComment(repo string, id int64) {
|
||||
s.mu.Lock()
|
||||
s.repo(repo).ProcessedComments[id] = true
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
// Seeded reports whether a repo has completed its baseline seeding poll.
|
||||
func (s *Store) Seeded(repo string) bool {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return s.repo(repo).Seeded
|
||||
}
|
||||
|
||||
// MarkSeeded records that a repo has completed baseline seeding.
|
||||
func (s *Store) MarkSeeded(repo string) {
|
||||
s.mu.Lock()
|
||||
s.repo(repo).Seeded = true
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
// LastPoll returns the time of the last completed poll for a repo.
|
||||
func (s *Store) LastPoll(repo string) time.Time {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return s.repo(repo).LastPoll
|
||||
}
|
||||
|
||||
// SetLastPoll records the time of the last completed poll for a repo.
|
||||
func (s *Store) SetLastPoll(repo string, t time.Time) {
|
||||
s.mu.Lock()
|
||||
s.repo(repo).LastPoll = t
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
// Save atomically writes the state document to disk.
|
||||
func (s *Store) Save() error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
if err := os.MkdirAll(filepath.Dir(s.path), 0o755); err != nil {
|
||||
return fmt.Errorf("creating state dir: %w", err)
|
||||
}
|
||||
raw, err := json.MarshalIndent(s.d, "", " ")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
tmp := s.path + ".tmp"
|
||||
if err := os.WriteFile(tmp, raw, 0o644); err != nil {
|
||||
return fmt.Errorf("writing state: %w", err)
|
||||
}
|
||||
if err := os.Rename(tmp, s.path); err != nil {
|
||||
return fmt.Errorf("committing state: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,122 @@
|
||||
package state
|
||||
|
||||
import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestMarkAndQuery(t *testing.T) {
|
||||
s, err := New(t.TempDir())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
const repo = "unkin/teabot"
|
||||
|
||||
if s.IssueProcessed(repo, 1) {
|
||||
t.Error("fresh store should not report issue 1 processed")
|
||||
}
|
||||
s.MarkIssue(repo, 1)
|
||||
if !s.IssueProcessed(repo, 1) {
|
||||
t.Error("issue 1 should be processed after MarkIssue")
|
||||
}
|
||||
if !s.ActedOnIssue(repo, 1) {
|
||||
t.Error("MarkIssue should also record acted-on")
|
||||
}
|
||||
if s.ActedOnPull(repo, 1) {
|
||||
t.Error("issue mark must not set acted-on-pull")
|
||||
}
|
||||
|
||||
s.MarkPull(repo, 2)
|
||||
if !s.PullProcessed(repo, 2) || !s.ActedOnPull(repo, 2) {
|
||||
t.Error("pull 2 should be processed and acted-on")
|
||||
}
|
||||
|
||||
s.MarkComment(repo, 99)
|
||||
if !s.CommentProcessed(repo, 99) {
|
||||
t.Error("comment 99 should be processed")
|
||||
}
|
||||
if s.CommentProcessed(repo, 100) {
|
||||
t.Error("comment 100 was never marked")
|
||||
}
|
||||
}
|
||||
|
||||
func TestPersistenceRoundTrip(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
s1, err := New(dir)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
const repo = "a/b"
|
||||
s1.MarkIssue(repo, 10)
|
||||
s1.MarkPull(repo, 11)
|
||||
s1.MarkComment(repo, 12)
|
||||
s1.MarkSeeded(repo)
|
||||
now := time.Now().Truncate(time.Second)
|
||||
s1.SetLastPoll(repo, now)
|
||||
if err := s1.Save(); err != nil {
|
||||
t.Fatalf("Save: %v", err)
|
||||
}
|
||||
|
||||
// A fresh Store loaded from the same dir must see the persisted state.
|
||||
s2, err := New(dir)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !s2.IssueProcessed(repo, 10) || !s2.PullProcessed(repo, 11) || !s2.CommentProcessed(repo, 12) {
|
||||
t.Error("processed sets did not survive reload")
|
||||
}
|
||||
if !s2.ActedOnIssue(repo, 10) || !s2.ActedOnPull(repo, 11) {
|
||||
t.Error("acted-on sets did not survive reload")
|
||||
}
|
||||
if !s2.Seeded(repo) {
|
||||
t.Error("seeded flag did not survive reload")
|
||||
}
|
||||
if !s2.LastPoll(repo).Equal(now) {
|
||||
t.Errorf("LastPoll = %v, want %v", s2.LastPoll(repo), now)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSaveIsAtomicFile(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
s, err := New(dir)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
s.MarkIssue("a/b", 1)
|
||||
if err := s.Save(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := os.Stat(filepath.Join(dir, StateFileName)); err != nil {
|
||||
t.Errorf("state file missing after Save: %v", err)
|
||||
}
|
||||
// No leftover temp file.
|
||||
if _, err := os.Stat(filepath.Join(dir, StateFileName+".tmp")); !os.IsNotExist(err) {
|
||||
t.Error("temp file should not remain after atomic rename")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSeededIndependentPerRepo(t *testing.T) {
|
||||
s, err := New(t.TempDir())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
s.MarkSeeded("a/b")
|
||||
if s.Seeded("c/d") {
|
||||
t.Error("seeding a/b must not seed c/d")
|
||||
}
|
||||
if !s.Seeded("a/b") {
|
||||
t.Error("a/b should be seeded")
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadCorruptStateFails(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
if err := os.WriteFile(filepath.Join(dir, StateFileName), []byte("{not json"), 0o644); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := New(dir); err == nil {
|
||||
t.Error("expected error loading corrupt state file")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user